kaxil commented on code in PR #71190:
URL: https://github.com/apache/airflow/pull/71190#discussion_r4206221067


##########
java-sdk/scripts/ci/prek/check_serialization_conformance.py:
##########
@@ -0,0 +1,54 @@
+#!/usr/bin/env python3
+# Licensed to the Apache Software Foundation (ASF) under one
+# or more contributor license agreements.  See the NOTICE file
+# distributed with this work for additional information
+# regarding copyright ownership.  The ASF licenses this file
+# to you under the Apache License, Version 2.0 (the
+# "License"); you may not use this file except in compliance
+# with the License.  You may obtain a copy of the License at
+#
+#   http://www.apache.org/licenses/LICENSE-2.0
+#
+# Unless required by applicable law or agreed to in writing,
+# software distributed under the License is distributed on an
+# "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+# KIND, either express or implied.  See the License for the
+# specific language governing permissions and limitations
+# under the License.
+"""Check the Java SDK serializes Dags as Airflow does, with 
scripts/ci/lang_sdk_serialization/compare.py."""
+
+from __future__ import annotations
+
+import subprocess
+import sys
+from pathlib import Path
+
+sys.path.insert(0, str(Path(__file__).resolve().parents[4] / "scripts" / "ci" 
/ "prek"))
+
+from common_prek_utils import AIRFLOW_ROOT_PATH
+
+if __name__ not in ("__main__", "__mp_main__"):
+    raise SystemExit(
+        "This file is intended to be executed as an executable program. You 
cannot use it as a module."
+        f"To run this script, run the ./{__file__} command"
+    )
+
+if __name__ == "__main__":
+    java_sdk = AIRFLOW_ROOT_PATH / "java-sdk"
+    # Gradle prints its own progress on stdout, so the classpath is the last 
line.
+    printed = subprocess.run(
+        [str(java_sdk / "gradlew"), "-p", str(java_sdk), "-q", 
":sdk:printConformanceClasspath"],
+        check=False,
+        capture_output=True,
+        text=True,
+    )
+    lines = printed.stdout.strip().splitlines()
+    if printed.returncode or not lines:
+        sys.stderr.write(printed.stdout)
+        sys.stderr.write(printed.stderr)
+        raise SystemExit("Could not build the Java SDK conformance classpath; 
see the Gradle output above")
+    classpath = lines[-1]
+    compare = AIRFLOW_ROOT_PATH / "scripts" / "ci" / "lang_sdk_serialization" 
/ "compare.py"
+    serializer = ["java", "-cp", classpath, 
"org.apache.airflow.sdk.conformance.SerializeJavaKt"]
+    command = [sys.executable, str(compare), "--sdk", "java", "--", 
*serializer]

Review Comment:
   This passes no `--supports`, so compare.py drops 
`conformance_literal_inputs` (it `requires: [literal_inputs]`), and 
`SerializeJava.kt` wires every task through `ref.after(...)` without reading 
`literals`. So no Java task in the run carries `_arg_bindings`. compare.py 
skips `_arg_bindings` on the grounds that loading the SDK's output through 
Airflow checks it, but for Java that load never sees one. Could SerializeJava 
wire a case's `upstream` and `literals` as named arguments, and this pass 
`--supports literal_inputs` (the Go hook passes `--supports all`)?



##########
.pre-commit-config.yaml:
##########
@@ -430,6 +430,25 @@ repos:
           ^go-sdk/schema/dag-schema\.json$|
           ^scripts/ci/lang_sdk_serialization/.*$|
           ^scripts/ci/prek/check_go_sdk_serialization_conformance\.py$
+      - id: check-java-sdk-serialization-conformance

Review Comment:
   CI runs `prek --all-files` (`ci-image-checks.yml`), and this filter always 
matches `serialized_objects.py` and `schema.json`, so the hook runs in every 
PR's static checks. Each run compiles the SDK with Gradle and resolves its 
dependencies from Maven Central. `selective_checks.py` skips `ktlint` and 
`regenerate-java-sdk-verification-metadata` when no java-sdk file changed for 
exactly that reason. Could this hook go in that skip set too, gated on java-sdk 
files, the two serialization files or `scripts/ci/lang_sdk_serialization/` 
changing? Separately, `internal/Fields.kt` isn't in the filter, though its 
`checkConfigValue` decides the form a config value is stored in before Serde 
writes it.



##########
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/ArgValues.kt:
##########
@@ -232,14 +231,15 @@ object ArgValues {
   ): Any? = decode(value, type)
 
   /**
-   * The inputs the Dag wired for this task, or null when the run's arguments
-   * come from the stub call site. A task with no wired inputs reads the
-   * bindings, so a stub call that bound nothing keeps its own diagnostics.
+   * The inputs the Dag wired for this task, or null when it wired none and the
+   * arguments come from the stub call site instead.
+   *
+   * A natively authored Dag always answers here, so its serialized binding
+   * spec is what Airflow records and shows rather than something read back at
+   * run time. A task handler wires nothing, so it reads the bindings and keeps
+   * its own diagnostics.
    */
-  internal fun wiredInputs(
-    context: Context,
-    client: Client,
-  ): List<Arg<*>>? = if (client.argBindings.isEmpty()) 
context.taskDef?.inputs?.takeIf { it.isNotEmpty() } else null
+  internal fun wiredInputs(context: Context): List<Arg<*>>? = 
context.taskDef?.inputs?.takeIf { it.isNotEmpty() }

Review Comment:
   With the bindings check gone, a native Java task always reads its wired 
inputs, even though `is_stub` now makes the API server send it `_arg_bindings`. 
Two places still describe the old order: the note in `java.rst` (around line 
695) says "Runtime argument bindings win over Java-declared wiring", and 
ADR-0007 decision G expects the spec to be what native Lang-SDK tasks bind 
from, which is what the Go task runner does 
(`convertArgBindings(details.TIContext.ArgBindings)`). Should the ADR and that 
note say Java treats the spec as informational, or should Java read the 
bindings the way Go does? Either way, none of the tests pair a client that has 
bindings with a wired context, so nothing pins which one wins.



##########
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt:
##########
@@ -0,0 +1,468 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.airflow.sdk.execution
+
+import com.fasterxml.jackson.databind.ObjectMapper
+import org.apache.airflow.sdk.Bundle
+import org.apache.airflow.sdk.DagDef
+import org.apache.airflow.sdk.GroupEdges
+import org.apache.airflow.sdk.GroupExpansion
+import org.apache.airflow.sdk.LiteralArg
+import org.apache.airflow.sdk.TaskDef
+import org.apache.airflow.sdk.TaskGroupRef
+import org.apache.airflow.sdk.TaskRef
+import org.apache.airflow.sdk.execution.comm.DagFileParseRequest
+import org.apache.airflow.sdk.internal.Field
+import org.apache.airflow.sdk.internal.SchemaFields
+import java.nio.file.InvalidPathException
+import java.nio.file.Paths
+import java.time.Duration
+import java.time.Instant
+import java.time.OffsetDateTime
+
+// Serializes Dags to Airflow DagSerialization v3 JSON, mirroring the 
TypeScript
+// SDK's serde (ts-sdk/src/coordinator/serde.ts), which in turn matches 
Python's
+// DagSerialization output.
+
+private val defaultsMapper = ObjectMapper()
+
+// Python drops both unless the operator names an email recipient. A Java task
+// has no email field to name one, so they are never written.
+private val OMITTED_TASK_KEYS = setOf("email_on_failure", "email_on_retry")
+
+/**
+ * Processes a [DagFileParseRequest] by serialising every Dag registered on
+ * [bundle] to DagSerialization v3 and returning the result as a
+ * DagFileParsingResult body.
+ *
+ * A Dag that cannot be serialized becomes an import error rather than taking
+ * the rest of the bundle with it. Airflow keys an import error by the
+ * bundle-relative path and holds one row per file, so every failure here is
+ * reported under that one key with its Dag named in the message.
+ */
+internal fun parseDags(
+  bundle: Bundle,
+  request: DagFileParseRequest,
+): Map<String, Any?> {
+  val fileloc = request.file ?: ""
+  val relativeFileloc = computeRelativeFileloc(fileloc, request.bundlePath)
+  val serializedDags = mutableListOf<Map<String, Any?>>()
+  val failures = mutableListOf<String>()
+  bundle.dags.values.forEach { dag ->
+    runCatching { serializeDag(dag, fileloc, relativeFileloc) }
+      .onSuccess { serializedDags += mapOf("data" to mapOf("__version" to 3, 
"dag" to it)) }
+      .onFailure { failures += "Dag \"${dag.id}\": ${it.message ?: 
it.javaClass.name}" }
+  }
+  return linkedMapOf<String, Any?>(
+    "type" to "DagFileParsingResult",
+    "fileloc" to fileloc,
+    "serialized_dags" to serializedDags,
+  ).apply {
+    if (failures.isNotEmpty()) this["import_errors"] = mapOf(relativeFileloc 
to failures.joinToString("\n"))
+  }
+}
+
+/**
+ * Converts a [DagDef] to Airflow DagSerialization v3 format. Required fields 
are
+ * always present; config-driven fields follow the rules in [applyDagConfig]
+ * (some always emitted, some only when set).
+ */
+internal fun serializeDag(
+  dag: DagDef,
+  fileloc: String,
+  relativeFileloc: String,
+): Map<String, Any?> {
+  // Group edges mean tasks only once the Dag is complete, so they are worked
+  // out here rather than carried on the groups themselves.
+  val expansion = dag.expandGroupEdges()
+  val downstream = linkedMapOf<String, MutableList<String>>()
+  dag.tasks.forEach { (taskId, def) ->
+    expansion.upstreamsOf(def).forEach { upstream ->
+      downstream.getOrPut(upstream) { mutableListOf() } += taskId
+    }
+  }
+
+  val result =
+    linkedMapOf<String, Any?>(
+      "dag_id" to dag.id,
+      "fileloc" to fileloc,
+      "relative_fileloc" to relativeFileloc,
+      "timezone" to dagTimezone(dag.dagConfig),
+      "timetable" to serializeTimetable(dag.id, dag.dagConfig),
+      "tasks" to dag.tasks.map { (taskId, def) -> serializeTask(taskId, def, 
downstream[taskId]) },
+      "dag_dependencies" to emptyList<Any?>(),
+      "task_group" to serializeTaskGroups(dag, expansion),
+      "edge_info" to emptyMap<String, Any?>(),
+      "params" to emptyList<Any?>(),
+      "deadline" to null,
+      "allowed_run_types" to null,
+    )
+  applyDagConfig(result, dag.dagConfig)
+  return result
+}
+
+/**
+ * Converts one task to the Airflow serialization format. `downstream` is the
+ * inverted view of the Dag's upstream edges, sorted for stable JSON.
+ */
+private fun serializeTask(
+  taskId: String,
+  def: TaskDef,
+  downstream: List<String>?,
+): Map<String, Any?> {
+  val data =
+    linkedMapOf<String, Any?>(
+      "task_id" to taskId,
+      "task_type" to def.definition.simpleName,
+      "_task_module" to def.definition.packageName,
+      "language" to "java",
+      // Python's operator serializer always emits template_fields (its list
+      // value never matches the tuple default it is compared against), so it
+      // is unconditional here too. Java tasks have no template fields.
+      "template_fields" to emptyList<Any?>(),
+      // What marks a task whose arguments Airflow resolves per instance for a
+      // runtime outside Python, as `@task.stub` does on the Python side.
+      // `get_arg_bindings` reads nothing without it.
+      "is_stub" to true,
+    )
+  argBindings(taskId, def)?.let { data["_arg_bindings"] = it }
+  // Emit only config entries that differ from their schema default, mirroring
+  // Python BaseSerialization's "omit hard-coded default" behavior. Operator
+  // fields are stored unwrapped, so the __type encoding is stripped.
+  def.configValues.forEach { (key, value) ->
+    if (key !in OMITTED_TASK_KEYS && 
!matchesSchemaDefault(SchemaFields.TASK[key], value)) {
+      data[key] = unwrapTypeEncoding(serializeValue(value))
+    }
+  }
+  if (!downstream.isNullOrEmpty()) {
+    data["downstream_task_ids"] = downstream.sorted()
+  }
+  return mapOf(
+    "__type" to "operator",
+    "__var" to data,
+  )
+}
+
+/**
+ * The task's arguments as the binding spec Airflow records, one entry per
+ * argument in the order the Dag's call passed them, as
+ * [ADR-0007](../../adr/lang-sdk/0007-taskflow-across-language-boundary.md)
+ * defines it.
+ *
+ * An upstream's handle becomes an `xcom` binding naming that task, and
+ * anything else a `literal` carrying the value. `value_schema` is left out: it
+ * constrains the decode side, and a Java task decodes into the type its own
+ * parameter declares.
+ *
+ * A Java task reads its arguments from the Dag in its own bundle rather than
+ * from this spec, so what it carries is what Airflow shows and what a change
+ * to an argument is seen in.
+ *
+ * Null for a task the Dag called with no arguments, which needs no spec.
+ */
+private fun argBindings(
+  taskId: String,
+  def: TaskDef,
+): List<Map<String, Any?>>? {
+  // A Dag wired by hand through Refs names no argument, and a spec without
+  // names binds nothing, so it is left out rather than written half-filled.
+  if (def.inputNames.size != def.inputs.size || def.inputs.isEmpty()) return 
null
+  return def.inputNames.zip(def.inputs) { name, input ->
+    when (input) {
+      is TaskRef<*> -> mapOf("name" to name, "kind" to "xcom", "task_id" to 
input.def.id)
+      is LiteralArg<*> ->
+        mapOf("name" to name, "kind" to "literal", "value" to 
plainJson(input.value, name, taskId))
+    }
+  }
+}
+
+/**
+ * [value] as the JSON the binding spec travels as, rejecting anything that has
+ * no JSON form.
+ *
+ * The spec is part of the serialized Dag, so a literal Airflow cannot store is
+ * refused where the Dag is written rather than where the task reads it.
+ */
+private fun plainJson(
+  value: Any?,
+  name: String,
+  taskId: String,
+): Any? =
+  when (value) {
+    null, is String, is Boolean -> value
+    is Double ->
+      value.takeIf { it.isFinite() }
+        ?: throw IllegalArgumentException(
+          "Argument '$name' of task '$taskId' is $value, which JSON has no 
form for; pass it as a string",
+        )
+    is Float -> plainJson(value.toDouble(), name, taskId)
+    is Int, is Long, is Short, is Byte -> value
+    is Collection<*> -> value.map { plainJson(it, name, taskId) }
+    is Array<*> -> value.map { plainJson(it, name, taskId) }
+    is Map<*, *> ->
+      value.entries.associate { (key, entry) ->
+        require(key is String) { "Argument '$name' of task '$taskId' has a map 
key that is not a string" }
+        key to plainJson(entry, name, taskId)
+      }
+    else ->
+      throw IllegalArgumentException(
+        "Argument '$name' of task '$taskId' is a ${value.javaClass.name}, 
which has no JSON form; the " +
+          "Dag's call arguments travel as JSON, so pass a string, number, 
boolean, list, or map",
+      )
+  }
+
+/**
+ * Writes Dag-level config onto [data], leaving out every field the Dag did
+ * not set. That includes the fields Python reads from Airflow's config
+ * (max_active_tasks, max_active_runs, max_consecutive_failed_dag_runs,
+ * catchup, disable_bundle_versioning): Airflow fills those in from its own
+ * config when it receives the Dag.
+ */
+private fun applyDagConfig(
+  data: MutableMap<String, Any?>,
+  config: Map<String, Any>,
+) {
+  listOf("description", "dag_display_name", "doc_md", "start_date", 
"end_date", "dagrun_timeout").forEach { key ->
+    config[key]?.let { data[key] = unwrapTypeEncoding(serializeValue(it)) }
+  }
+  (config["tags"] as? List<*>)?.let { tags ->
+    // Python stores tags in a set and serializes them sorted (for a stable
+    // dag_hash); mirror that regardless of registration order.
+    data["tags"] = tags.map { it.toString() }.distinct().sorted()
+  }
+  listOf(
+    "max_active_tasks",
+    "max_active_runs",
+    "max_consecutive_failed_dag_runs",
+    "catchup",
+    "disable_bundle_versioning",
+  ).forEach { key -> config[key]?.let { data[key] = it } }
+  // fail_fast and render_template_as_native_obj have schema default false, so
+  // Python omits them when false; keep that behavior.
+  if (config["fail_fast"] == true) data["fail_fast"] = true
+  if (config["render_template_as_native_obj"] == true) 
data["render_template_as_native_obj"] = true
+  config["is_paused_upon_creation"]?.let { data["is_paused_upon_creation"] = 
it }
+}
+
+// TODO: respect [scheduler] create_cron_data_intervals like Python's
+// _create_timetable; the JVM bundle cannot read airflow.cfg, so the
+// supervisor must send those flags over the coordinator protocol first.
+// The TypeScript SDK waits on the same flag; tracked at
+// https://github.com/apache/airflow/issues/67938
+private fun serializeTimetable(
+  dagId: String,
+  config: Map<String, Any>,
+): Map<String, Any?> =
+  when (val schedule = config["schedule"] as String?) {
+    null -> mapOf("__type" to "airflow.timetables.simple.NullTimetable", 
"__var" to emptyMap<String, Any?>())
+    "@once" -> mapOf("__type" to "airflow.timetables.simple.OnceTimetable", 
"__var" to emptyMap<String, Any?>())
+    "@continuous" ->
+      mapOf("__type" to "airflow.timetables.simple.ContinuousTimetable", 
"__var" to emptyMap<String, Any?>())
+    else -> {
+      val expression = CRON_PRESETS[schedule] ?: schedule
+      require(isCronExpression(expression)) {
+        "Schedule '$schedule' of Dag '$dagId' is not a cron expression or a 
preset " +
+          "(${CRON_PRESETS.keys.joinToString()}, @once, @continuous); a 
schedule the scheduler cannot " +
+          "parse would leave the Dag unschedulable"
+      }
+      mapOf(
+        "__type" to "airflow.timetables.trigger.CronTriggerTimetable",
+        "__var" to
+          mapOf(
+            "expression" to expression,
+            "timezone" to dagTimezone(config),
+            "interval" to 0.0,
+            "run_immediately" to false,
+          ),
+      )
+    }
+  }
+
+/**
+ * The Dag's timezone, as Python's `encode_timezone` writes it: `"UTC"` for a
+ * zero offset, otherwise the offset in seconds.
+ *
+ * Python takes it from `start_date`, so a cron schedule runs in the zone the
+ * Dag's start date was written in. A Dag with no start date runs in UTC.
+ */
+private fun dagTimezone(config: Map<String, Any>): Any =

Review Comment:
   Python falls back to `[core] default_timezone` when there's no start date 
(`DAG._extract_tz`), so with `default_timezone = America/New_York` a Python Dag 
with a cron schedule and no start date runs on New York time, while the same 
Java Dag runs on UTC. It's the same can't-read-the-config gap as the 
`create_cron_data_intervals` TODO above. Could that TODO cover this too, or the 
Java docs say a cron Dag without a start date is scheduled in UTC?



##########
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt:
##########
@@ -0,0 +1,468 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.airflow.sdk.execution
+
+import com.fasterxml.jackson.databind.ObjectMapper
+import org.apache.airflow.sdk.Bundle
+import org.apache.airflow.sdk.DagDef
+import org.apache.airflow.sdk.GroupEdges
+import org.apache.airflow.sdk.GroupExpansion
+import org.apache.airflow.sdk.LiteralArg
+import org.apache.airflow.sdk.TaskDef
+import org.apache.airflow.sdk.TaskGroupRef
+import org.apache.airflow.sdk.TaskRef
+import org.apache.airflow.sdk.execution.comm.DagFileParseRequest
+import org.apache.airflow.sdk.internal.Field
+import org.apache.airflow.sdk.internal.SchemaFields
+import java.nio.file.InvalidPathException
+import java.nio.file.Paths
+import java.time.Duration
+import java.time.Instant
+import java.time.OffsetDateTime
+
+// Serializes Dags to Airflow DagSerialization v3 JSON, mirroring the 
TypeScript
+// SDK's serde (ts-sdk/src/coordinator/serde.ts), which in turn matches 
Python's
+// DagSerialization output.
+
+private val defaultsMapper = ObjectMapper()
+
+// Python drops both unless the operator names an email recipient. A Java task
+// has no email field to name one, so they are never written.
+private val OMITTED_TASK_KEYS = setOf("email_on_failure", "email_on_retry")
+
+/**
+ * Processes a [DagFileParseRequest] by serialising every Dag registered on
+ * [bundle] to DagSerialization v3 and returning the result as a
+ * DagFileParsingResult body.
+ *
+ * A Dag that cannot be serialized becomes an import error rather than taking
+ * the rest of the bundle with it. Airflow keys an import error by the
+ * bundle-relative path and holds one row per file, so every failure here is
+ * reported under that one key with its Dag named in the message.
+ */
+internal fun parseDags(
+  bundle: Bundle,
+  request: DagFileParseRequest,
+): Map<String, Any?> {
+  val fileloc = request.file ?: ""
+  val relativeFileloc = computeRelativeFileloc(fileloc, request.bundlePath)
+  val serializedDags = mutableListOf<Map<String, Any?>>()
+  val failures = mutableListOf<String>()
+  bundle.dags.values.forEach { dag ->
+    runCatching { serializeDag(dag, fileloc, relativeFileloc) }
+      .onSuccess { serializedDags += mapOf("data" to mapOf("__version" to 3, 
"dag" to it)) }
+      .onFailure { failures += "Dag \"${dag.id}\": ${it.message ?: 
it.javaClass.name}" }
+  }
+  return linkedMapOf<String, Any?>(
+    "type" to "DagFileParsingResult",
+    "fileloc" to fileloc,
+    "serialized_dags" to serializedDags,
+  ).apply {
+    if (failures.isNotEmpty()) this["import_errors"] = mapOf(relativeFileloc 
to failures.joinToString("\n"))
+  }
+}
+
+/**
+ * Converts a [DagDef] to Airflow DagSerialization v3 format. Required fields 
are
+ * always present; config-driven fields follow the rules in [applyDagConfig]
+ * (some always emitted, some only when set).
+ */
+internal fun serializeDag(
+  dag: DagDef,
+  fileloc: String,
+  relativeFileloc: String,
+): Map<String, Any?> {
+  // Group edges mean tasks only once the Dag is complete, so they are worked
+  // out here rather than carried on the groups themselves.
+  val expansion = dag.expandGroupEdges()
+  val downstream = linkedMapOf<String, MutableList<String>>()
+  dag.tasks.forEach { (taskId, def) ->
+    expansion.upstreamsOf(def).forEach { upstream ->
+      downstream.getOrPut(upstream) { mutableListOf() } += taskId
+    }
+  }
+
+  val result =
+    linkedMapOf<String, Any?>(
+      "dag_id" to dag.id,
+      "fileloc" to fileloc,
+      "relative_fileloc" to relativeFileloc,
+      "timezone" to dagTimezone(dag.dagConfig),
+      "timetable" to serializeTimetable(dag.id, dag.dagConfig),
+      "tasks" to dag.tasks.map { (taskId, def) -> serializeTask(taskId, def, 
downstream[taskId]) },
+      "dag_dependencies" to emptyList<Any?>(),
+      "task_group" to serializeTaskGroups(dag, expansion),
+      "edge_info" to emptyMap<String, Any?>(),
+      "params" to emptyList<Any?>(),
+      "deadline" to null,
+      "allowed_run_types" to null,
+    )
+  applyDagConfig(result, dag.dagConfig)
+  return result
+}
+
+/**
+ * Converts one task to the Airflow serialization format. `downstream` is the
+ * inverted view of the Dag's upstream edges, sorted for stable JSON.
+ */
+private fun serializeTask(
+  taskId: String,
+  def: TaskDef,
+  downstream: List<String>?,
+): Map<String, Any?> {
+  val data =
+    linkedMapOf<String, Any?>(
+      "task_id" to taskId,
+      "task_type" to def.definition.simpleName,
+      "_task_module" to def.definition.packageName,
+      "language" to "java",
+      // Python's operator serializer always emits template_fields (its list
+      // value never matches the tuple default it is compared against), so it
+      // is unconditional here too. Java tasks have no template fields.
+      "template_fields" to emptyList<Any?>(),
+      // What marks a task whose arguments Airflow resolves per instance for a
+      // runtime outside Python, as `@task.stub` does on the Python side.
+      // `get_arg_bindings` reads nothing without it.
+      "is_stub" to true,
+    )
+  argBindings(taskId, def)?.let { data["_arg_bindings"] = it }
+  // Emit only config entries that differ from their schema default, mirroring
+  // Python BaseSerialization's "omit hard-coded default" behavior. Operator
+  // fields are stored unwrapped, so the __type encoding is stripped.
+  def.configValues.forEach { (key, value) ->
+    if (key !in OMITTED_TASK_KEYS && 
!matchesSchemaDefault(SchemaFields.TASK[key], value)) {
+      data[key] = unwrapTypeEncoding(serializeValue(value))
+    }
+  }
+  if (!downstream.isNullOrEmpty()) {
+    data["downstream_task_ids"] = downstream.sorted()
+  }
+  return mapOf(
+    "__type" to "operator",
+    "__var" to data,
+  )
+}
+
+/**
+ * The task's arguments as the binding spec Airflow records, one entry per
+ * argument in the order the Dag's call passed them, as
+ * [ADR-0007](../../adr/lang-sdk/0007-taskflow-across-language-boundary.md)
+ * defines it.
+ *
+ * An upstream's handle becomes an `xcom` binding naming that task, and
+ * anything else a `literal` carrying the value. `value_schema` is left out: it
+ * constrains the decode side, and a Java task decodes into the type its own
+ * parameter declares.
+ *
+ * A Java task reads its arguments from the Dag in its own bundle rather than
+ * from this spec, so what it carries is what Airflow shows and what a change
+ * to an argument is seen in.
+ *
+ * Null for a task the Dag called with no arguments, which needs no spec.
+ */
+private fun argBindings(
+  taskId: String,
+  def: TaskDef,
+): List<Map<String, Any?>>? {
+  // A Dag wired by hand through Refs names no argument, and a spec without
+  // names binds nothing, so it is left out rather than written half-filled.
+  if (def.inputNames.size != def.inputs.size || def.inputs.isEmpty()) return 
null
+  return def.inputNames.zip(def.inputs) { name, input ->
+    when (input) {
+      is TaskRef<*> -> mapOf("name" to name, "kind" to "xcom", "task_id" to 
input.def.id)
+      is LiteralArg<*> ->
+        mapOf("name" to name, "kind" to "literal", "value" to 
plainJson(input.value, name, taskId))
+    }
+  }
+}
+
+/**
+ * [value] as the JSON the binding spec travels as, rejecting anything that has
+ * no JSON form.
+ *
+ * The spec is part of the serialized Dag, so a literal Airflow cannot store is
+ * refused where the Dag is written rather than where the task reads it.
+ */
+private fun plainJson(
+  value: Any?,
+  name: String,
+  taskId: String,
+): Any? =
+  when (value) {
+    null, is String, is Boolean -> value
+    is Double ->
+      value.takeIf { it.isFinite() }
+        ?: throw IllegalArgumentException(
+          "Argument '$name' of task '$taskId' is $value, which JSON has no 
form for; pass it as a string",
+        )
+    is Float -> plainJson(value.toDouble(), name, taskId)
+    is Int, is Long, is Short, is Byte -> value
+    is Collection<*> -> value.map { plainJson(it, name, taskId) }
+    is Array<*> -> value.map { plainJson(it, name, taskId) }
+    is Map<*, *> ->
+      value.entries.associate { (key, entry) ->
+        require(key is String) { "Argument '$name' of task '$taskId' has a map 
key that is not a string" }
+        key to plainJson(entry, name, taskId)
+      }
+    else ->
+      throw IllegalArgumentException(
+        "Argument '$name' of task '$taskId' is a ${value.javaClass.name}, 
which has no JSON form; the " +
+          "Dag's call arguments travel as JSON, so pass a string, number, 
boolean, list, or map",
+      )
+  }
+
+/**
+ * Writes Dag-level config onto [data], leaving out every field the Dag did
+ * not set. That includes the fields Python reads from Airflow's config
+ * (max_active_tasks, max_active_runs, max_consecutive_failed_dag_runs,
+ * catchup, disable_bundle_versioning): Airflow fills those in from its own
+ * config when it receives the Dag.
+ */
+private fun applyDagConfig(
+  data: MutableMap<String, Any?>,
+  config: Map<String, Any>,
+) {
+  listOf("description", "dag_display_name", "doc_md", "start_date", 
"end_date", "dagrun_timeout").forEach { key ->
+    config[key]?.let { data[key] = unwrapTypeEncoding(serializeValue(it)) }
+  }
+  (config["tags"] as? List<*>)?.let { tags ->
+    // Python stores tags in a set and serializes them sorted (for a stable
+    // dag_hash); mirror that regardless of registration order.
+    data["tags"] = tags.map { it.toString() }.distinct().sorted()
+  }
+  listOf(
+    "max_active_tasks",
+    "max_active_runs",
+    "max_consecutive_failed_dag_runs",
+    "catchup",
+    "disable_bundle_versioning",
+  ).forEach { key -> config[key]?.let { data[key] = it } }
+  // fail_fast and render_template_as_native_obj have schema default false, so
+  // Python omits them when false; keep that behavior.
+  if (config["fail_fast"] == true) data["fail_fast"] = true
+  if (config["render_template_as_native_obj"] == true) 
data["render_template_as_native_obj"] = true
+  config["is_paused_upon_creation"]?.let { data["is_paused_upon_creation"] = 
it }
+}
+
+// TODO: respect [scheduler] create_cron_data_intervals like Python's
+// _create_timetable; the JVM bundle cannot read airflow.cfg, so the
+// supervisor must send those flags over the coordinator protocol first.
+// The TypeScript SDK waits on the same flag; tracked at
+// https://github.com/apache/airflow/issues/67938
+private fun serializeTimetable(
+  dagId: String,
+  config: Map<String, Any>,
+): Map<String, Any?> =
+  when (val schedule = config["schedule"] as String?) {
+    null -> mapOf("__type" to "airflow.timetables.simple.NullTimetable", 
"__var" to emptyMap<String, Any?>())
+    "@once" -> mapOf("__type" to "airflow.timetables.simple.OnceTimetable", 
"__var" to emptyMap<String, Any?>())
+    "@continuous" ->
+      mapOf("__type" to "airflow.timetables.simple.ContinuousTimetable", 
"__var" to emptyMap<String, Any?>())
+    else -> {
+      val expression = CRON_PRESETS[schedule] ?: schedule
+      require(isCronExpression(expression)) {
+        "Schedule '$schedule' of Dag '$dagId' is not a cron expression or a 
preset " +
+          "(${CRON_PRESETS.keys.joinToString()}, @once, @continuous); a 
schedule the scheduler cannot " +
+          "parse would leave the Dag unschedulable"
+      }
+      mapOf(
+        "__type" to "airflow.timetables.trigger.CronTriggerTimetable",
+        "__var" to
+          mapOf(
+            "expression" to expression,
+            "timezone" to dagTimezone(config),
+            "interval" to 0.0,
+            "run_immediately" to false,
+          ),
+      )
+    }
+  }
+
+/**
+ * The Dag's timezone, as Python's `encode_timezone` writes it: `"UTC"` for a
+ * zero offset, otherwise the offset in seconds.
+ *
+ * Python takes it from `start_date`, so a cron schedule runs in the zone the
+ * Dag's start date was written in. A Dag with no start date runs in UTC.
+ */
+private fun dagTimezone(config: Map<String, Any>): Any =
+  (config["start_date"] as? OffsetDateTime)
+    ?.offset
+    ?.totalSeconds
+    ?.takeIf { it != 0 }
+    ?: "UTC"
+
+/**
+ * Presets expanded the way `CronMixin.__init__` expands them, so the
+ * serialized expression is the one Python records, which the Dag's summary and
+ * its hash are both taken from. Mirrors `airflow.utils.dates.cron_presets`.
+ */
+private val CRON_PRESETS =
+  mapOf(
+    "@hourly" to "0 * * * *",
+    "@daily" to "0 0 * * *",
+    "@weekly" to "0 0 * * 0",
+    "@monthly" to "0 0 1 * *",
+    "@quarterly" to "0 0 1 */3 *",
+    "@yearly" to "0 0 1 1 *",
+  )
+
+private val CRON_FIELD = Regex("[\\d*,\\-/?LW#]+|[A-Z]{3}(-[A-Z]{3})?", 
RegexOption.IGNORE_CASE)

Review Comment:
   This rejects some ordinary schedules. The comma is only allowed in the 
numeric alternative, so `0 9 * * MON,WED,FRI`, `0 0 * JAN,JUL *` and `0 0 * * 
MON#2` all fail the shape check, as do `@midnight` and `@annually` (Python 
passes both to croniter unexpanded since they aren't in `cron_presets`, and 
croniter accepts them). I ran this regex next to croniter to check. The Dag 
then becomes an import error. Matching each comma-separated element on its own, 
and letting the croniter-only `@` aliases through, would fix it. Core's 
`validate_serialized_dag` doesn't call `timetable.validate()`, so this check is 
the only thing catching prose like `every tuesday`, which is worth keeping.



##########
java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt:
##########
@@ -0,0 +1,468 @@
+/*
+ * Licensed to the Apache Software Foundation (ASF) under one
+ * or more contributor license agreements.  See the NOTICE file
+ * distributed with this work for additional information
+ * regarding copyright ownership.  The ASF licenses this file
+ * to you under the Apache License, Version 2.0 (the
+ * "License"); you may not use this file except in compliance
+ * with the License.  You may obtain a copy of the License at
+ *
+ *   http://www.apache.org/licenses/LICENSE-2.0
+ *
+ * Unless required by applicable law or agreed to in writing,
+ * software distributed under the License is distributed on an
+ * "AS IS" BASIS, WITHOUT WARRANTIES OR CONDITIONS OF ANY
+ * KIND, either express or implied.  See the License for the
+ * specific language governing permissions and limitations
+ * under the License.
+ */
+
+package org.apache.airflow.sdk.execution
+
+import com.fasterxml.jackson.databind.ObjectMapper
+import org.apache.airflow.sdk.Bundle
+import org.apache.airflow.sdk.DagDef
+import org.apache.airflow.sdk.GroupEdges
+import org.apache.airflow.sdk.GroupExpansion
+import org.apache.airflow.sdk.LiteralArg
+import org.apache.airflow.sdk.TaskDef
+import org.apache.airflow.sdk.TaskGroupRef
+import org.apache.airflow.sdk.TaskRef
+import org.apache.airflow.sdk.execution.comm.DagFileParseRequest
+import org.apache.airflow.sdk.internal.Field
+import org.apache.airflow.sdk.internal.SchemaFields
+import java.nio.file.InvalidPathException
+import java.nio.file.Paths
+import java.time.Duration
+import java.time.Instant
+import java.time.OffsetDateTime
+
+// Serializes Dags to Airflow DagSerialization v3 JSON, mirroring the 
TypeScript
+// SDK's serde (ts-sdk/src/coordinator/serde.ts), which in turn matches 
Python's
+// DagSerialization output.
+
+private val defaultsMapper = ObjectMapper()
+
+// Python drops both unless the operator names an email recipient. A Java task
+// has no email field to name one, so they are never written.
+private val OMITTED_TASK_KEYS = setOf("email_on_failure", "email_on_retry")
+
+/**
+ * Processes a [DagFileParseRequest] by serialising every Dag registered on
+ * [bundle] to DagSerialization v3 and returning the result as a
+ * DagFileParsingResult body.
+ *
+ * A Dag that cannot be serialized becomes an import error rather than taking
+ * the rest of the bundle with it. Airflow keys an import error by the
+ * bundle-relative path and holds one row per file, so every failure here is
+ * reported under that one key with its Dag named in the message.
+ */
+internal fun parseDags(
+  bundle: Bundle,
+  request: DagFileParseRequest,
+): Map<String, Any?> {
+  val fileloc = request.file ?: ""
+  val relativeFileloc = computeRelativeFileloc(fileloc, request.bundlePath)
+  val serializedDags = mutableListOf<Map<String, Any?>>()
+  val failures = mutableListOf<String>()
+  bundle.dags.values.forEach { dag ->
+    runCatching { serializeDag(dag, fileloc, relativeFileloc) }
+      .onSuccess { serializedDags += mapOf("data" to mapOf("__version" to 3, 
"dag" to it)) }
+      .onFailure { failures += "Dag \"${dag.id}\": ${it.message ?: 
it.javaClass.name}" }
+  }
+  return linkedMapOf<String, Any?>(
+    "type" to "DagFileParsingResult",
+    "fileloc" to fileloc,
+    "serialized_dags" to serializedDags,
+  ).apply {
+    if (failures.isNotEmpty()) this["import_errors"] = mapOf(relativeFileloc 
to failures.joinToString("\n"))
+  }
+}
+
+/**
+ * Converts a [DagDef] to Airflow DagSerialization v3 format. Required fields 
are
+ * always present; config-driven fields follow the rules in [applyDagConfig]
+ * (some always emitted, some only when set).
+ */
+internal fun serializeDag(
+  dag: DagDef,
+  fileloc: String,
+  relativeFileloc: String,
+): Map<String, Any?> {
+  // Group edges mean tasks only once the Dag is complete, so they are worked
+  // out here rather than carried on the groups themselves.
+  val expansion = dag.expandGroupEdges()
+  val downstream = linkedMapOf<String, MutableList<String>>()
+  dag.tasks.forEach { (taskId, def) ->
+    expansion.upstreamsOf(def).forEach { upstream ->
+      downstream.getOrPut(upstream) { mutableListOf() } += taskId
+    }
+  }
+
+  val result =
+    linkedMapOf<String, Any?>(
+      "dag_id" to dag.id,
+      "fileloc" to fileloc,
+      "relative_fileloc" to relativeFileloc,
+      "timezone" to dagTimezone(dag.dagConfig),
+      "timetable" to serializeTimetable(dag.id, dag.dagConfig),
+      "tasks" to dag.tasks.map { (taskId, def) -> serializeTask(taskId, def, 
downstream[taskId]) },
+      "dag_dependencies" to emptyList<Any?>(),
+      "task_group" to serializeTaskGroups(dag, expansion),
+      "edge_info" to emptyMap<String, Any?>(),
+      "params" to emptyList<Any?>(),
+      "deadline" to null,
+      "allowed_run_types" to null,
+    )
+  applyDagConfig(result, dag.dagConfig)
+  return result
+}
+
+/**
+ * Converts one task to the Airflow serialization format. `downstream` is the
+ * inverted view of the Dag's upstream edges, sorted for stable JSON.
+ */
+private fun serializeTask(
+  taskId: String,
+  def: TaskDef,
+  downstream: List<String>?,
+): Map<String, Any?> {
+  val data =
+    linkedMapOf<String, Any?>(
+      "task_id" to taskId,
+      "task_type" to def.definition.simpleName,
+      "_task_module" to def.definition.packageName,
+      "language" to "java",
+      // Python's operator serializer always emits template_fields (its list
+      // value never matches the tuple default it is compared against), so it
+      // is unconditional here too. Java tasks have no template fields.
+      "template_fields" to emptyList<Any?>(),
+      // What marks a task whose arguments Airflow resolves per instance for a
+      // runtime outside Python, as `@task.stub` does on the Python side.
+      // `get_arg_bindings` reads nothing without it.
+      "is_stub" to true,
+    )
+  argBindings(taskId, def)?.let { data["_arg_bindings"] = it }
+  // Emit only config entries that differ from their schema default, mirroring
+  // Python BaseSerialization's "omit hard-coded default" behavior. Operator
+  // fields are stored unwrapped, so the __type encoding is stripped.
+  def.configValues.forEach { (key, value) ->
+    if (key !in OMITTED_TASK_KEYS && 
!matchesSchemaDefault(SchemaFields.TASK[key], value)) {

Review Comment:
   With no `client_defaults` in the payload, and `fill_config_defaults` only 
filling Dag-level keys, a Java task that leaves `retries`, `queue` or `owner` 
unset loads with the `SerializedBaseOperator` class defaults rather than 
`[core] default_task_retries` or `[operators] default_queue`. Go and TS don't 
send `client_defaults` either, so filling task defaults probably belongs in 
core next to `fill_config_defaults`. But dropping an explicit value that equals 
the schema default here would get in the way of that: once core fills task 
defaults from config, an explicit `retries=0` would look unset and pick up the 
configured value. Could explicitly set task values always be written?



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to