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


##########
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:
   Confirmed all five against croniter before changing anything. The check now 
matches each comma-separated element on its own, and lets the croniter-only 
aliases through:
   
   ```kotlin
   private const val CRON_VALUE = "(\\d+|\\*|\\?|[A-Z]{3})"
   
   private val CRON_ELEMENT =
     Regex("$CRON_VALUE(-$CRON_VALUE)?([/#]\\d+)?[LW]*|L(-\\d+)?|LW", 
RegexOption.IGNORE_CASE)
   
   private val CRON_ALIASES = setOf("@midnight", "@annually")
   ```
   
   `@midnight` and `@annually` serialize unexpanded, which is what 
`cron_presets.get(cron, cron)` leaves them as. Tests cover the five you named 
plus `0 0 15W * *` and `*/5 1-5/2 * * MON-FRI`, and still reject `every 
tuesday`, `@bogus` and `0 0 * * tuesday`, so the prose guard you wanted kept is 
intact.
   
   Fixed in b8938f81f2a.
   



##########
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:
   The code was the part that was right and the docs were wrong. `bindInput` 
checks `wiredInputs(context)` first, so a native Java Dag binds from its own 
wiring and the serialized spec is what Airflow records and shows. Go reading 
the delivered bindings is a legitimate difference rather than a drift, since 
both emit the same spec, and for Java it is also what keeps `TaskInput` 
coherent: a wired `TaskInput` decodes one upstream value into the whole object, 
where the spec is positional and named.
   
   So: the `java.rst` note now says wiring wins and that runtime bindings are 
what a `@Builder.TaskHandler` reads; ADR-0007 G gained a paragraph recording 
that a runtime holding its own Dag may resolve locally while the spec stays the 
recorded form; and a new `ArgValuesTest` case pairs a client that carries 
bindings with a wired context and asserts the wired value is what the task 
receives. Nothing pinned that before, which is how three statements managed to 
disagree.
   
   Fixed in b8938f81f2a.
   



##########
.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:
   Agreed on the cost, but the gate had to be wider than `JAVA_SDK_FILES`. I 
checked the claim before writing it: a change under 
`airflow-core/src/airflow/serialization/` or 
`scripts/ci/lang_sdk_serialization/` does **not** force `full_tests_needed` 
(those paths only select `SelectiveCoreTestType.SERIALIZATION`), so a skip 
gated on java-sdk files alone would have stopped the hook running on exactly 
the PRs it exists for.
   
   Added `FileGroupForCi.JAVA_SDK_CONFORMANCE_FILES`, matching `java-sdk/` plus 
`serialized_objects.py`, `schema.json` and 
`scripts/ci/lang_sdk_serialization/`; the hook is skipped unless that group 
matches. `04_selective_checks.md` and `test_selective_checks.py` are updated 
alongside it, with cases asserting a serializer, schema or harness change does 
not skip it.
   
   `internal/Fields.kt` is in the hook's own filter now too.
   
   Fixed in b8938f81f2a.
   



##########
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:
   I checked Python, Go, TS before touching this, and it drops them as well.
   
   Noted in b8938f81f2a.
   



##########
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:
   Covered both ways. The `create_cron_data_intervals` TODO now names `[core] 
default_timezone` and `_extract_tz` as the same cannot-read-the-config gap, and 
`java.rst` says a cron Dag with no `startDate` is scheduled in UTC, so set 
`startDate` to pin the zone.
   
   Fixed in b8938f81f2a.
   



##########
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:
   Right, and the net effect was that no Java task in the run carried a spec at 
all. `SerializeJava` now wires each case's `upstream` handles followed by its 
`literals` as the task's call arguments, and the hook passes `--supports 
literal_inputs`.
   
   Binding names are positional, `arg0..argN`, which is the convention the Go 
SDK's serializer already emits (`serialize.go`) because reflect cannot read Go 
parameter names either; the Java interface API has the same gap, so no change 
to the shared `test_dags.yaml` was needed.
   
   The run now produces bindings on 9 tasks, nested literals included:
   
   ```json
   
{"name":"arg2","kind":"literal","value":{"limit":10,"ratio":0.5,"tags":["a",null],"nested":{"enabled":true}}}
   ```
   
   Fixed in b8938f81f2a.
   



-- 
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