This is an automated email from the ASF dual-hosted git repository.

jason810496 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git


The following commit(s) were added to refs/heads/main by this push:
     new 01b977568e2 Java SDK: Serialize native Dags to DagSerialization v3 
(#71190)
01b977568e2 is described below

commit 01b977568e25844a87c56f28e489e5e735d68f70
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Thu Oct 8 11:56:57 2026 +0800

    Java SDK: Serialize native Dags to DagSerialization v3 (#71190)
    
    * Java SDK: Serialize native Dags to DagSerialization v3 on parse requests
    
    A Java-authored Dag could not reach the scheduler on its own: the runtime
    answered task-execution requests only, so a Python stub file still had to
    exist purely to describe the Dag's structure. With dependency edges and
    schema-keyed configuration now recorded on the Java Dag model, the runtime
    has everything it needs to answer the coordinator's DagFileParseRequest the
    same way the Go SDK does, and the nativedag examples become real schedulable
    Dags with no Python counterpart.
    
    Native Java tasks deliberately emit no _arg_bindings: the execution API
    delivers bindings only for Python _StubOperator tasks, and a Java task
    always runs inside the JVM bundle that already holds its wired inputs, so
    the runtime resolves them locally.
    
    Cron schedules map to CronTriggerTimetable only, mirroring the Go SDK until
    the supervisor forwards the [scheduler] timetable flags over the coordinator
    protocol (see the TODO at the timetable serializer).
    
    * Test that the conformance check compares Dag params
    
    * Java SDK: Leave config-backed Dag fields for Airflow to fill in
    
    * Java SDK: Leave out the email flags a Java task cannot act on
    
    * Java SDK: Serialize a native Dag's task groups
    
    * Java SDK: Write each Dag tag once, as Python does
    
    * Java SDK: Check native Dag serialization against Airflow's
    
    * Java SDK: Run the serialization conformance check on every file that 
steers it
    
    `Bundle.register` expands the group edges, `Refs` decides which group each
    task lands in, and `sdk/build.gradle.kts` generates the schema fields the
    serializer writes from. All three change the serialized output, so add
    them to the hook's trigger set.
    
    Report a failed classpath build with the Gradle output instead of an
    `IndexError` on an empty line list or a bare `CalledProcessError`.
    
    * Java SDK: Build a native Dag's timetable the way Python does
    
    Expand a cron preset before writing it, so `@daily` serializes as
    `0 0 * * *` and the Dag hashes the same as the Python one. Reject a
    schedule that is neither a preset nor a cron expression where the Dag is
    written, rather than letting it reach a scheduler that cannot build a
    timetable from it.
    
    Take the timezone from the Dag's start date, as Python's `DAG.timezone`
    does, instead of always writing UTC: a Dag started at a non-UTC offset
    now fires at that offset rather than at the same wall clock in UTC.
    
    * Java SDK: Carry a native Dag's call arguments as argument bindings
    
    A native Dag's serialized form recorded its task edges but nothing about
    what each task was called with, so changing a literal argument left the
    Dag byte-identical: no new version, and nothing in the version diff.
    
    Write the same `is_stub` flag and `_arg_bindings` spec the TypeScript SDK
    writes, as ADR-0007 decision G says a natively authored Dag should. The
    wiring view now passes its parameter names along with the arguments, so
    each binding can name the parameter it feeds.
    
    A Java task still resolves its arguments from the Dag in its own bundle
    rather than reading the spec back: the bundle that runs the task also
    holds the call that wired it, so the values never have to travel, and a
    `TaskInput` keeps binding as one whole input. The spec is what Airflow
    records and shows.
    
    Since the spec travels as JSON, a literal with no JSON form is now
    rejected where the Dag is parsed instead of reaching a task that cannot
    receive it.
    
    * Java SDK: Keep one unserializable Dag from failing a whole bundle's parse
    
    A Dag that could not be serialized threw out of the parse response, so
    every other Dag in the same JAR vanished with it and `import_errors` was
    never filled. Report the failure under the bundle-relative path the way
    the TypeScript SDK does, naming the Dag, and serialize the rest.
    
    * Java SDK: Mark the native Dag capabilities as shipping in 3.4
    
    `native-dag-authoring`, `task-args`, `taskflow-dependencies` and
    `task-group` are new in this release, not in 3.3, which is what the
    matrix's "supported since" column says.
    
    * Java SDK: Accept every cron schedule croniter does and bind wired inputs 
first
    
    The cron shape check now accepts comma lists of names, month and weekday
    names, the `#`, `L` and `W` qualifiers and the `@midnight` and `@annually`
    aliases, which Python stores unexpanded. Each of these previously turned
    the Dag into an import error. The schema-default comment records why
    explicit default values are dropped, and the timezone TODO and the Java
    docs note that a cron Dag with no start date is scheduled in UTC.
    
    A wired input wins over a runtime binding in ArgValues. The Java docs and
    ADR 0007 now say so, and a test pins it.
    
    The conformance run now passes `--supports literal_inputs` and the Java
    serializer wires each task's upstream handles and literals as call
    arguments, so the run exercises `_arg_bindings`. The hook's file filter
    covers `internal/Fields.kt`. Selective checks skip the hook unless a
    java-sdk file, Airflow's serializer or schema, or the shared harness
    changed, through a new file group because serializer changes do not force
    full tests.
---
 .pre-commit-config.yaml                            |  19 +
 .../0007-taskflow-across-language-boundary.md      |   5 +
 .../language-sdks/java.rst                         |  13 +-
 dev/breeze/doc/ci/04_selective_checks.md           |   6 +
 .../src/airflow_breeze/utils/selective_checks.py   |  14 +
 dev/breeze/tests/test_selective_checks.py          |  88 +++-
 java-sdk/README.md                                 |  22 +-
 java-sdk/capabilities.yaml                         |  17 +-
 .../org/apache/airflow/sdk/BuilderProcessor.kt     |   6 +-
 .../kotlin/org/apache/airflow/sdk/BuilderTest.kt   |   5 +-
 .../ci/prek/check_serialization_conformance.py     |  63 +++
 java-sdk/sdk/build.gradle.kts                      |   9 +
 java-sdk/sdk/module.md                             |  22 +-
 .../main/kotlin/org/apache/airflow/sdk/DagDef.kt   |   3 +
 .../src/main/kotlin/org/apache/airflow/sdk/Deps.kt |   4 +
 .../main/kotlin/org/apache/airflow/sdk/Server.kt   |  11 +
 .../org/apache/airflow/sdk/execution/Serde.kt      | 490 +++++++++++++++++++++
 .../org/apache/airflow/sdk/internal/ArgValues.kt   |  24 +-
 .../kotlin/org/apache/airflow/sdk/internal/Refs.kt |  15 +-
 .../org/apache/airflow/sdk/internal/TaskArgs.kt    |   2 +-
 .../org/apache/airflow/sdk/ArgTestSupport.kt       |   2 +-
 .../kotlin/org/apache/airflow/sdk/ArgValuesTest.kt |  11 +
 .../kotlin/org/apache/airflow/sdk/ServerTest.kt    |  63 +++
 .../airflow/sdk/conformance/SerializeJava.kt       | 140 ++++++
 .../org/apache/airflow/sdk/execution/SerdeTest.kt  | 428 ++++++++++++++++++
 .../apache/airflow/sdk/internal/ArgValuesTest.kt   |   2 +-
 .../org/apache/airflow/sdk/internal/RefsTest.kt    |   6 +-
 scripts/ci/lang_sdk_serialization/compare.py       |  17 +-
 scripts/ci/lang_sdk_serialization/test_dags.yaml   |   4 +-
 .../ci/lang_sdk_serialization/test_compare.py      |  11 +
 30 files changed, 1433 insertions(+), 89 deletions(-)

diff --git a/.pre-commit-config.yaml b/.pre-commit-config.yaml
index 7ccab9e4dc9..ec3dd995779 100644
--- a/.pre-commit-config.yaml
+++ b/.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
+        name: Check the Java SDK serializes Dags the way Airflow does
+        description: "Serialize the shared test Dags with the Java SDK and 
with Airflow, and compare the two"
+        entry: ./java-sdk/scripts/ci/prek/check_serialization_conformance.py
+        language: system
+        pass_filenames: false
+        require_serial: true
+        files: >
+          (?x)
+          ^airflow-core/src/airflow/serialization/serialized_objects\.py$|
+          ^airflow-core/src/airflow/serialization/schema\.json$|
+          ^java-sdk/scripts/ci/prek/check_serialization_conformance\.py$|
+          ^java-sdk/sdk/schema/dag-schema\.json$|
+          ^java-sdk/sdk/build\.gradle\.kts$|
+          
^java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/(Bundle|DagDef|Deps|Arg|TaskGroupRef)\.kt$|
+          
^java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde\.kt$|
+          
^java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/(Refs|Fields)\.kt$|
+          ^java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/conformance/.*$|
+          ^scripts/ci/lang_sdk_serialization/.*$
       - id: check-go-version-in-sync
         name: Check Go toolchain version is consistent across build files
         entry: ./scripts/ci/prek/check_go_version_in_sync.py
diff --git 
a/airflow-core/adr/lang-sdk/0007-taskflow-across-language-boundary.md 
b/airflow-core/adr/lang-sdk/0007-taskflow-across-language-boundary.md
index c1f9d06fcf6..48588b35a77 100644
--- a/airflow-core/adr/lang-sdk/0007-taskflow-across-language-boundary.md
+++ b/airflow-core/adr/lang-sdk/0007-taskflow-across-language-boundary.md
@@ -256,6 +256,11 @@ modes. It is also why the spec is deliberately not 
`@task.stub`-shaped — no fi
 natively authored call. (Its docstrings still say "stub", reflecting the only 
producer that
 exists today; that is wording to revisit, not a constraint in the format.)
 
+A runtime that holds its own Dag may resolve call arguments from that Dag 
rather than from the
+delivered spec. The spec stays the recorded and displayed form, which keeps 
one wire format across
+both authoring modes. The Go SDK reads the delivered bindings and the Java SDK 
resolves locally;
+both emit the same spec.
+
 ### H. Scope: what does not cross the boundary
 
 The following raise when the Dag is serialized, rather than being silently 
dropped or deferred to
diff --git a/airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst 
b/airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst
index 446557db8ea..12b65d85580 100644
--- a/airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst
+++ b/airflow-core/docs/authoring-and-scheduling/language-sdks/java.rst
@@ -675,7 +675,9 @@ generated view, with a no-argument ``depends()`` method:
 
 Every ``@Builder.Task`` method must be called in the wiring class; a task the 
wiring missed fails at
 Dag-parse time.  ``lit(...)`` wires an inline constant where no upstream feeds 
a parameter.  A bare
-``double`` cannot be an ``Arg``, so a constant is wrapped.  A view method that 
takes no arguments
+``double`` cannot be an ``Arg``, so a constant is wrapped.  Airflow records 
what each task is called
+with, and that record travels as JSON, so a constant has to be a string, 
number, boolean, list or
+map.  A view method that takes no arguments
 returns the same handle every time, so it names one node wherever it appears; 
one that takes
 arguments is called once, and the wiring fails if it is called again with 
arguments, so hold its
 handle in a local and reuse that.
@@ -690,9 +692,9 @@ class that supplies only task bodies, for a Dag a Python 
file declares, carries
 
 .. note::
 
-   Runtime argument bindings win over Java-declared wiring.  When the 
supervisor delivers bindings
-   for a run (see :ref:`java-sdk/arg-binding`), the binding at a parameter's 
position is what the
-   task receives.  Wired inputs are the fallback, which is what a native Java 
Dag always uses.
+   A native Java Dag binds its task arguments from its own wiring, and the 
``_arg_bindings`` it
+   serializes are what Airflow records and shows.  Runtime bindings (see 
:ref:`java-sdk/arg-binding`)
+   are what a ``@Builder.TaskHandler`` class reads, for a task whose Dag a 
Python file declares.
 
 Task groups
 ~~~~~~~~~~~
@@ -772,6 +774,9 @@ Durations and date-times are ISO-8601 strings in 
annotations (``retryDelay = "PT
 ``java.time.OffsetDateTime`` values in ``config`` calls.  An unknown key or a 
mismatched value type
 fails the build for an annotation, and the ``config`` call itself for an 
object.
 
+A Dag with a cron ``schedule`` runs in the time zone of its ``startDate``.  
With no ``startDate`` it
+is scheduled in UTC, so set ``startDate`` to pin the zone.
+
 .. _java-sdk/task-state-store:
 
 Task state store
diff --git a/dev/breeze/doc/ci/04_selective_checks.md 
b/dev/breeze/doc/ci/04_selective_checks.md
index ec177bd87a8..9e16809dde8 100644
--- a/dev/breeze/doc/ci/04_selective_checks.md
+++ b/dev/breeze/doc/ci/04_selective_checks.md
@@ -548,6 +548,12 @@ when some files are not changed. Those are the rules 
implemented:
     Gradle wrapper, which downloads the Gradle distribution, and the latter 
additionally
     resolves the whole Java SDK dependency graph from Maven Central, so we 
avoid those
     downloads on PRs that do not touch `java-sdk/`)
+  * if no `Java SDK conformance files` changed - 
`check-java-sdk-serialization-conformance`
+    is skipped (it compiles the Java SDK with Gradle and resolves its 
dependencies from Maven
+    Central). The group is wider than `Java SDK files`: it also covers 
Airflow's serializer
+    (`serialized_objects.py`), `schema.json` and the shared harness under
+    `scripts/ci/lang_sdk_serialization/`, because the check compares the SDK 
against those and
+    none of them forces `full_tests_needed`
   * if no `TS SDK files` (`ts-sdk/`) changed - 
`check-ts-sdk-supervisor-schema` check is
     skipped (it regenerates and diffs the generated ts-sdk file; a change to 
the supervisor
     wire schema alone deliberately does not trigger it - regenerating the 
ts-sdk types is
diff --git a/dev/breeze/src/airflow_breeze/utils/selective_checks.py 
b/dev/breeze/src/airflow_breeze/utils/selective_checks.py
index 0388797404e..39334c54135 100644
--- a/dev/breeze/src/airflow_breeze/utils/selective_checks.py
+++ b/dev/breeze/src/airflow_breeze/utils/selective_checks.py
@@ -140,6 +140,7 @@ class FileGroupForCi(Enum):
     AGENT_FRAMEWORK_FILES = auto()
     GO_SDK_FILES = auto()
     JAVA_SDK_FILES = auto()
+    JAVA_SDK_CONFORMANCE_FILES = auto()
     TS_SDK_FILES = auto()
     TS_SDK_DOCS_FILES = auto()
     AIRFLOW_CTL_FILES = auto()
@@ -517,6 +518,14 @@ CI_FILE_GROUP_MATCHES: HashableDict[FileGroupForCi] = 
HashableDict(
             # `.md` excluded — doc-only edits do not affect the Gradle build.
             r"^java-sdk/(?!.*\.md$).*",
         ],
+        FileGroupForCi.JAVA_SDK_CONFORMANCE_FILES: [
+            # The Java SDK, plus what its serialization conformance check 
compares it against:
+            # Airflow's serializer and schema, and the shared harness.
+            r"^java-sdk/(?!.*\.md$).*",
+            r"^airflow-core/src/airflow/serialization/serialized_objects\.py$",
+            r"^airflow-core/src/airflow/serialization/schema\.json$",
+            r"^scripts/ci/lang_sdk_serialization/.*",
+        ],
         FileGroupForCi.TS_SDK_DOCS_FILES: [
             # TypeDoc renders the reference from the SDK sources and category 
entry points,
             # and the landing page is authored in ts-sdk/docs — unlike 
TS_SDK_FILES, `.md`
@@ -1942,6 +1951,11 @@ class SelectiveChecks:
             # from Maven Central. Skip it when no java-sdk files changed so 
unrelated PRs do not
             # depend on that resolution.
             prek_hooks_to_skip.add("regenerate-java-sdk-verification-metadata")
+        if not self._matching_files(FileGroupForCi.JAVA_SDK_CONFORMANCE_FILES, 
CI_FILE_GROUP_MATCHES):
+            # This hook compiles the Java SDK with Gradle and resolves its 
dependencies from Maven
+            # Central. Skip it unless a java-sdk file, Airflow's serializer or 
schema, or the shared
+            # harness changed. Those last do not force full_tests_needed, so 
they need their own group.
+            prek_hooks_to_skip.add("check-java-sdk-serialization-conformance")
         if not self._matching_files(FileGroupForCi.TS_SDK_FILES, 
CI_FILE_GROUP_MATCHES):
             # This hook regenerates ts-sdk/src/generated/supervisor.ts from 
the wire schema and
             # diffs it. Schema-only changes deliberately do not trigger it: 
regenerating the
diff --git a/dev/breeze/tests/test_selective_checks.py 
b/dev/breeze/tests/test_selective_checks.py
index 7e8d6e8d712..28664c92413 100644
--- a/dev/breeze/tests/test_selective_checks.py
+++ b/dev/breeze/tests/test_selective_checks.py
@@ -113,7 +113,7 @@ LIST_OF_ALL_PROVIDER_TESTS_AS_JSON = json.dumps(
 
 
 ALL_SKIPPED_COMMITS_ON_NO_CI_IMAGE = (
-    
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,flynt,identity,ktlint,lint-helm-chart,"
+    
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,flynt,identity,ktlint,lint-helm-chart,"
     
"mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
     "mypy-scripts,"
@@ -127,7 +127,7 @@ ALL_SKIPPED_COMMITS_ON_NO_CI_IMAGE = (
 ALL_SKIPPED_COMMITS_BY_DEFAULT_ON_ALL_TESTS_NEEDED = "identity,update-uv-lock"
 
 ALL_SKIPPED_COMMITS_IF_ONLY_UI_OPENAPI_CHANGED = (
-    
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,flynt,identity,ktlint,"
+    
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,flynt,identity,ktlint,"
     
"lint-helm-chart,mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,"
     
"mypy-airflow-e2e-tests,mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,"
     
"mypy-kubernetes-tests,mypy-scripts,mypy-shared-configuration,mypy-shared-dagnode,"
@@ -139,7 +139,7 @@ ALL_SKIPPED_COMMITS_IF_ONLY_UI_OPENAPI_CHANGED = (
 )
 
 ALL_SKIPPED_COMMITS_IF_NO_UI = (
-    
"check-ts-sdk-supervisor-schema,identity,ktlint,mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
+    
"check-java-sdk-serialization-conformance,check-ts-sdk-supervisor-schema,identity,ktlint,mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
     "mypy-scripts,"
     
"mypy-shared-configuration,mypy-shared-dagnode,mypy-shared-listeners,mypy-shared-logging,"
@@ -149,7 +149,7 @@ ALL_SKIPPED_COMMITS_IF_NO_UI = (
     
"regenerate-java-sdk-verification-metadata,ts-compile-lint-simple-auth-manager-ui,ts-compile-lint-ui,update-uv-lock"
 )
 ALL_SKIPPED_COMMITS_IF_NO_HELM_TESTS = (
-    "check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+    
"check-java-sdk-serialization-conformance,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
     
"mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
     "mypy-scripts,"
@@ -160,7 +160,7 @@ ALL_SKIPPED_COMMITS_IF_NO_HELM_TESTS = (
 )
 
 ALL_SKIPPED_COMMITS_IF_NO_UI_AND_HELM_TESTS = (
-    "check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+    
"check-java-sdk-serialization-conformance,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
     
"mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
     "mypy-scripts,"
@@ -175,7 +175,7 @@ ALL_SKIPPED_COMMITS_IF_NO_UI_AND_HELM_TESTS = (
 # forced. airflow-core Python changed (so mypy-airflow-core + flynt run); no
 # provider.yaml, helm, or UI files changed, so those checks stay skipped.
 ALL_SKIPPED_COMMITS_IF_ONLY_API_SOURCE_CHANGED = (
-    
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+    
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
     "mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
     "mypy-scripts,"
@@ -188,7 +188,7 @@ ALL_SKIPPED_COMMITS_IF_ONLY_API_SOURCE_CHANGED = (
 )
 
 ALL_SKIPPED_COMMITS_IF_NO_PROVIDERS_AND_UI = (
-    "check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,"
+    
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,"
     
"mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
     "mypy-scripts,"
@@ -200,7 +200,7 @@ ALL_SKIPPED_COMMITS_IF_NO_PROVIDERS_AND_UI = (
 )
 
 ALL_SKIPPED_COMMITS_IF_NO_PROVIDERS = (
-    
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+    
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
     
"mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
     "mypy-scripts,"
@@ -213,7 +213,7 @@ ALL_SKIPPED_COMMITS_IF_NO_PROVIDERS = (
 
 
 ALL_SKIPPED_COMMITS_IF_NO_PROVIDERS_UI_AND_HELM_TESTS = (
-    
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+    
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
     
"mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
     "mypy-scripts,"
@@ -225,7 +225,7 @@ ALL_SKIPPED_COMMITS_IF_NO_PROVIDERS_UI_AND_HELM_TESTS = (
 )
 
 ALL_SKIPPED_COMMITS_IF_NO_CODE_PROVIDERS_AND_HELM_TESTS = (
-    
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,flynt,identity,ktlint,lint-helm-chart,"
+    
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,flynt,identity,ktlint,lint-helm-chart,"
     
"mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
     "mypy-scripts,"
@@ -236,7 +236,7 @@ ALL_SKIPPED_COMMITS_IF_NO_CODE_PROVIDERS_AND_HELM_TESTS = (
 )
 
 ALL_SKIPPED_COMMITS_IF_NOT_IMPORTANT_FILES_CHANGED = (
-    
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,flynt,identity,ktlint,lint-helm-chart,"
+    
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,flynt,identity,ktlint,lint-helm-chart,"
     
"mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
     "mypy-scripts,"
@@ -481,7 +481,7 @@ def assert_outputs_are_printed(expected_outputs: dict[str, 
str], stderr: str):
                     "run-amazon-tests": "false",
                     "docs-build": "true",
                     "skip-prek-hooks": (
-                        
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+                        
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
                         
"mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
                         
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
                         "mypy-scripts,"
@@ -527,7 +527,7 @@ def assert_outputs_are_printed(expected_outputs: dict[str, 
str], stderr: str):
                     "run-api-tests": "true",
                     "docs-build": "true",
                     "skip-prek-hooks": (
-                        
"check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+                        
"check-java-sdk-serialization-conformance,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
                         
"mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
                         
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
                         "mypy-scripts,"
@@ -782,7 +782,7 @@ def assert_outputs_are_printed(expected_outputs: dict[str, 
str], stderr: str):
                     "docs-build": "true",
                     "full-tests-needed": "false",
                     "skip-prek-hooks": (
-                        
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+                        
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
                         
"mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
                         
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
                         "mypy-scripts,"
@@ -820,7 +820,7 @@ def assert_outputs_are_printed(expected_outputs: dict[str, 
str], stderr: str):
                     "docs-build": "false",
                     "full-tests-needed": "false",
                     "skip-prek-hooks": (
-                        
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+                        
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
                         
"mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
                         
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
                         "mypy-scripts,"
@@ -883,7 +883,7 @@ def assert_outputs_are_printed(expected_outputs: dict[str, 
str], stderr: str):
                     "docs-build": "true",
                     "full-tests-needed": "false",
                     "skip-prek-hooks": (
-                        
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+                        
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
                         
"mypy-airflow-core,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
                         
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
                         "mypy-scripts,"
@@ -919,7 +919,7 @@ def assert_outputs_are_printed(expected_outputs: dict[str, 
str], stderr: str):
                     "docs-build": "false",
                     "full-tests-needed": "false",
                     "skip-prek-hooks": (
-                        
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+                        
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
                         
"mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-e2e-tests,"
                         
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
                         "mypy-scripts,"
@@ -1268,7 +1268,7 @@ def assert_outputs_are_printed(expected_outputs: 
dict[str, str], stderr: str):
                 "docs-build": "false",
                 "run-kubernetes-tests": "false",
                 "skip-prek-hooks": (
-                    
"check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+                    
"check-java-sdk-serialization-conformance,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
                     
"mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
                     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
                     "mypy-scripts,"
@@ -1475,7 +1475,7 @@ def assert_outputs_are_printed(expected_outputs: 
dict[str, str], stderr: str):
                 "run-amazon-tests": "false",
                 "docs-build": "true",
                 "skip-prek-hooks": (
-                    
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,flynt,identity,ktlint,"
+                    
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,flynt,identity,ktlint,"
                     
"mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
                     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
                     "mypy-scripts,"
@@ -1879,7 +1879,7 @@ def assert_outputs_are_printed(expected_outputs: 
dict[str, str], stderr: str):
                 ("shared/logging/src/airflow_shared/logging/remote.py",),
                 {
                     "skip-prek-hooks": (
-                        
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+                        
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
                         
"mypy-airflow-core,mypy-airflow-ctl,mypy-airflow-ctl-tests,"
                         
"mypy-airflow-e2e-tests,mypy-dev,mypy-devel-common,mypy-docker-tests,"
                         "mypy-helm-tests,mypy-kubernetes-tests,mypy-scripts,"
@@ -1911,7 +1911,10 @@ def test_expected_output_pull_request_main(
     assert_outputs_are_printed(expected_outputs, str(stderr))
 
 
[email protected]("hook", ["ktlint", 
"regenerate-java-sdk-verification-metadata"])
[email protected](
+    "hook",
+    ["ktlint", "regenerate-java-sdk-verification-metadata", 
"check-java-sdk-serialization-conformance"],
+)
 @pytest.mark.parametrize(
     ("files", "hook_skipped"),
     [
@@ -1948,6 +1951,43 @@ def 
test_java_sdk_gradle_hooks_only_run_for_java_sdk_changes(
     assert (hook in skipped_hooks) is hook_skipped
 
 
[email protected](
+    ("files", "hook_skipped"),
+    [
+        pytest.param(
+            ("airflow-core/src/airflow/serialization/serialized_objects.py",),
+            False,
+            id="runs when Airflow's serializer changes",
+        ),
+        pytest.param(
+            ("airflow-core/src/airflow/serialization/schema.json",),
+            False,
+            id="runs when Airflow's Dag schema changes",
+        ),
+        pytest.param(
+            ("scripts/ci/lang_sdk_serialization/compare.py",),
+            False,
+            id="runs when the shared conformance harness changes",
+        ),
+        pytest.param(
+            ("SECURITY.md",),
+            True,
+            id="skipped when none of its inputs change",
+        ),
+    ],
+)
+def test_java_sdk_conformance_hook_runs_for_serializer_changes(files: 
tuple[str, ...], hook_skipped: bool):
+    stderr = SelectiveChecks(
+        files=files,
+        commit_ref=NEUTRAL_COMMIT,
+        github_event=GithubEvents.PULL_REQUEST,
+        pr_labels=tuple(),
+        default_branch="main",
+    )
+    skipped_hooks = 
get_outputs_from_stderr(str(stderr))["skip-prek-hooks"].split(",")
+    assert ("check-java-sdk-serialization-conformance" in skipped_hooks) is 
hook_skipped
+
+
 @pytest.mark.parametrize(
     ("files", "hook_skipped"),
     [
@@ -2821,7 +2861,7 @@ def test_expected_output_push(
                 "docs-build": "true",
                 "docs-list-as-string": ALL_DOCS_SELECTED_FOR_BUILD,
                 "skip-prek-hooks": (
-                    
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+                    
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
                     
"mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
                     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
                     "mypy-scripts,"
@@ -2863,7 +2903,7 @@ def test_expected_output_push(
                 "microsoft.mssql mongo mysql openlineage oracle postgres "
                 "presto salesforce samba sftp ssh standard trino",
                 "skip-prek-hooks": (
-                    
"check-ts-sdk-supervisor-schema,identity,ktlint,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
+                    
"check-java-sdk-serialization-conformance,check-ts-sdk-supervisor-schema,identity,ktlint,mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
                     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
                     "mypy-scripts,"
                     
"mypy-shared-configuration,mypy-shared-dagnode,mypy-shared-listeners,mypy-shared-logging,"
@@ -2908,7 +2948,7 @@ def test_expected_output_push(
                 "docs-build": "true",
                 "docs-list-as-string": "apache-airflow",
                 "skip-prek-hooks": (
-                    
"check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
+                    
"check-java-sdk-serialization-conformance,check-provider-yaml-valid,check-ts-sdk-supervisor-schema,identity,ktlint,lint-helm-chart,"
                     
"mypy-airflow-ctl,mypy-airflow-ctl-tests,mypy-airflow-e2e-tests,"
                     
"mypy-dev,mypy-devel-common,mypy-docker-tests,mypy-helm-tests,mypy-kubernetes-tests,"
                     "mypy-scripts,"
diff --git a/java-sdk/README.md b/java-sdk/README.md
index 33a3c560f7b..6b4ba3522fd 100644
--- a/java-sdk/README.md
+++ b/java-sdk/README.md
@@ -635,17 +635,17 @@ prek hook regenerate it.
 | capability: `asset-event-emit` | MAY | ✗ | – | runtime does not emit asset 
events yet |
 | capability: `asset-event-read` | MAY | ✗ | – | no task-facing asset-event 
API yet |
 | **Native-Dag authoring** |  |  |  |  |
-| capability: `native-dag-authoring` | SHOULD | ✗ | – | native Dag authoring 
not implemented yet |
-| capability: `task-args` | MUST † | n/a | – |  |
-| capability: `dag-params` | MUST † | n/a | – |  |
-| capability: `taskflow-dependencies` | MUST † | n/a | – |  |
-| capability: `branching` | SHOULD † | n/a | – |  |
-| capability: `dag-test` | SHOULD † | n/a | – |  |
-| capability: `task-group` | MAY † | n/a | – |  |
-| capability: `dynamic-task-mapping` | MAY † | n/a | – |  |
-| capability: `asset-inlets-outlets` | MAY † | n/a | – |  |
-| capability: `asset-scheduling` | MAY † | n/a | – |  |
-| capability: `object-store` | MAY † | n/a | – |  |
+| capability: `native-dag-authoring` | SHOULD | ✓ | 3.4 |  |
+| capability: `task-args` | MUST † | ✓ | 3.4 |  |
+| capability: `dag-params` | MUST † | ✗ | – |  |
+| capability: `taskflow-dependencies` | MUST † | ✓ | 3.4 |  |
+| capability: `branching` | SHOULD † | ✗ | – |  |
+| capability: `dag-test` | SHOULD † | ✗ | – |  |
+| capability: `task-group` | MAY † | ✓ | 3.4 |  |
+| capability: `dynamic-task-mapping` | MAY † | ✗ | – |  |
+| capability: `asset-inlets-outlets` | MAY † | ✗ | – |  |
+| capability: `asset-scheduling` | MAY † | ✗ | – |  |
+| capability: `object-store` | MAY † | ✗ | – |  |
 
 *Marks: ✓ supported · ✗ not supported · n/a not applicable. A tier marked † 
applies only when `native-dag-authoring` is supported.*
 
diff --git a/java-sdk/capabilities.yaml b/java-sdk/capabilities.yaml
index 6f61f999240..75c6053fc89 100644
--- a/java-sdk/capabilities.yaml
+++ b/java-sdk/capabilities.yaml
@@ -52,8 +52,8 @@ states:
     supported: true
     since: "3.3"
 
-# Runtime capabilities reflect the task-facing Client surface; native-Dag 
authoring is not
-# implemented yet, so every native capability is unsupported.
+# Runtime capabilities reflect the task-facing Client surface. Native-Dag 
capabilities describe
+# what a Java-authored Dag can declare and deliver through Dag serialization.
 capabilities:
   mixed-lang-stub-target:
     supported: true
@@ -96,20 +96,23 @@ capabilities:
     supported: false
     note: "no task-facing asset-event API yet"
   native-dag-authoring:
-    supported: false
-    note: "native Dag authoring not implemented yet"
+    supported: true
+    since: "3.4"
   task-args:
-    supported: false
+    supported: true
+    since: "3.4"
   dag-params:
     supported: false
   taskflow-dependencies:
-    supported: false
+    supported: true
+    since: "3.4"
   branching:
     supported: false
   dag-test:
     supported: false
   task-group:
-    supported: false
+    supported: true
+    since: "3.4"
   dynamic-task-mapping:
     supported: false
   asset-inlets-outlets:
diff --git 
a/java-sdk/processor/src/main/kotlin/org/apache/airflow/sdk/BuilderProcessor.kt 
b/java-sdk/processor/src/main/kotlin/org/apache/airflow/sdk/BuilderProcessor.kt
index 1101dbb9084..cc0d9c594f8 100644
--- 
a/java-sdk/processor/src/main/kotlin/org/apache/airflow/sdk/BuilderProcessor.kt
+++ 
b/java-sdk/processor/src/main/kotlin/org/apache/airflow/sdk/BuilderProcessor.kt
@@ -362,11 +362,15 @@ class BuilderProcessor : AbstractProcessor() {
     if (decl.dataParams.isEmpty()) {
       method.addStatement($$"return $T.node($L, $L)", REFS_TYPE, group, def)
     } else {
+      // The parameter names ride along so the serialized Dag can name each
+      // argument, as the binding spec ADR-0007 defines requires.
       method.addStatement(
-        $$"return $T.call($L, $L, $L)",
+        $$"return $T.call($L, $L, $T.of($L), $L)",
         REFS_TYPE,
         group,
         def,
+        LIST_TYPE,
+        CodeBlock.join(decl.dataParams.map { CodeBlock.of($$"$S", it.name) }, 
", "),
         decl.dataParams.joinToString { it.name },
       )
     }
diff --git 
a/java-sdk/processor/src/test/kotlin/org/apache/airflow/sdk/BuilderTest.kt 
b/java-sdk/processor/src/test/kotlin/org/apache/airflow/sdk/BuilderTest.kt
index d75c33aa45d..c3759fa1998 100644
--- a/java-sdk/processor/src/test/kotlin/org/apache/airflow/sdk/BuilderTest.kt
+++ b/java-sdk/processor/src/test/kotlin/org/apache/airflow/sdk/BuilderTest.kt
@@ -136,6 +136,7 @@ class BuilderTest {
 
          import java.lang.Integer;
          import java.lang.Void;
+         import java.util.List;
          import org.apache.airflow.sdk.Arg;
          import org.apache.airflow.sdk.Deps;
          import org.apache.airflow.sdk.TaskDef;
@@ -159,7 +160,7 @@ class BuilderTest {
            }
 
            default TaskRef<Void> t3(Arg<? extends Integer> value) {
-             return Refs.call("", new TaskDef("t3", 
TestExampleBuilder.T3.class), value);
+             return Refs.call("", new TaskDef("t3", 
TestExampleBuilder.T3.class), List.of("value"), value);
            }
          }
         """,
@@ -381,7 +382,7 @@ class BuilderTest {
 
            default TaskRef<Void> t(Arg<? extends String> text, Arg<?> anything,
                Arg<? extends List<String>> items, Arg<? extends Long> boxed) {
-             return Refs.call("", new TaskDef("t", 
TestExampleBuilder.T.class), text, anything, items, boxed);
+             return Refs.call("", new TaskDef("t", 
TestExampleBuilder.T.class), List.of("text", "anything", "items", "boxed"), 
text, anything, items, boxed);
            }
          }
         """,
diff --git a/java-sdk/scripts/ci/prek/check_serialization_conformance.py 
b/java-sdk/scripts/ci/prek/check_serialization_conformance.py
new file mode 100755
index 00000000000..9869f4ccc5c
--- /dev/null
+++ b/java-sdk/scripts/ci/prek/check_serialization_conformance.py
@@ -0,0 +1,63 @@
+#!/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",
+        "--supports",
+        "literal_inputs",
+        "--",
+        *serializer,
+    ]
+    sys.exit(subprocess.run(command, check=False).returncode)
diff --git a/java-sdk/sdk/build.gradle.kts b/java-sdk/sdk/build.gradle.kts
index 9ef581634bd..0159e1d2803 100644
--- a/java-sdk/sdk/build.gradle.kts
+++ b/java-sdk/sdk/build.gradle.kts
@@ -865,3 +865,12 @@ publishing {
         }
     }
 }
+
+// Prints the classpath that runs the conformance serializer, for
+// java-sdk/scripts/ci/prek/check_serialization_conformance.py.
+tasks.register("printConformanceClasspath") {
+    dependsOn("testClasses")
+    // Capture early to keep compatibility to the Gradle configuration cache.
+    val classpath = sourceSets.test.get().runtimeClasspath
+    doLast { println(classpath.asPath) }
+}
diff --git a/java-sdk/sdk/module.md b/java-sdk/sdk/module.md
index 07e272d1431..3d11a42da5e 100644
--- a/java-sdk/sdk/module.md
+++ b/java-sdk/sdk/module.md
@@ -56,17 +56,17 @@ meaning of each dimension is defined in the
 | capability: `asset-event-emit` | MAY | ✗ | – | runtime does not emit asset 
events yet |
 | capability: `asset-event-read` | MAY | ✗ | – | no task-facing asset-event 
API yet |
 | **Native-Dag authoring** |  |  |  |  |
-| capability: `native-dag-authoring` | SHOULD | ✗ | – | native Dag authoring 
not implemented yet |
-| capability: `task-args` | MUST † | n/a | – |  |
-| capability: `dag-params` | MUST † | n/a | – |  |
-| capability: `taskflow-dependencies` | MUST † | n/a | – |  |
-| capability: `branching` | SHOULD † | n/a | – |  |
-| capability: `dag-test` | SHOULD † | n/a | – |  |
-| capability: `task-group` | MAY † | n/a | – |  |
-| capability: `dynamic-task-mapping` | MAY † | n/a | – |  |
-| capability: `asset-inlets-outlets` | MAY † | n/a | – |  |
-| capability: `asset-scheduling` | MAY † | n/a | – |  |
-| capability: `object-store` | MAY † | n/a | – |  |
+| capability: `native-dag-authoring` | SHOULD | ✓ | 3.4 |  |
+| capability: `task-args` | MUST † | ✓ | 3.4 |  |
+| capability: `dag-params` | MUST † | ✗ | – |  |
+| capability: `taskflow-dependencies` | MUST † | ✓ | 3.4 |  |
+| capability: `branching` | SHOULD † | ✗ | – |  |
+| capability: `dag-test` | SHOULD † | ✗ | – |  |
+| capability: `task-group` | MAY † | ✓ | 3.4 |  |
+| capability: `dynamic-task-mapping` | MAY † | ✗ | – |  |
+| capability: `asset-inlets-outlets` | MAY † | ✗ | – |  |
+| capability: `asset-scheduling` | MAY † | ✗ | – |  |
+| capability: `object-store` | MAY † | ✗ | – |  |
 
 *Marks: ✓ supported · ✗ not supported · n/a not applicable. A tier marked † 
applies only when `native-dag-authoring` is supported.*
 
diff --git a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/DagDef.kt 
b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/DagDef.kt
index d5ba757beab..ff6f740af18 100644
--- a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/DagDef.kt
+++ b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/DagDef.kt
@@ -316,6 +316,9 @@ class TaskDef(
 
   internal val configValues = linkedMapOf<String, Any>()
   internal val inputs = mutableListOf<Arg<*>>()
+
+  /** Name of the task parameter each of [inputs] feeds, in the same order. */
+  internal val inputNames = mutableListOf<String>()
   internal val upstreams = linkedSetOf<TaskDef>()
   internal var owner: DagDef? = null
 
diff --git a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Deps.kt 
b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Deps.kt
index 0dab9c978f7..01b26177cd6 100644
--- a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Deps.kt
+++ b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Deps.kt
@@ -125,6 +125,10 @@ interface Deps {
    * `transform(extract(), lit(0.9))`. It is passed to the task as a constant
    * and creates no dependency edge.
    *
+   * The Dag's call arguments travel to Airflow as JSON, so the value has to be
+   * a string, number, boolean, list, or map; anything else fails when the Dag
+   * is parsed.
+   *
    * @param value Constant to bind; may be null for a nullable parameter.
    */
   fun <T> lit(value: T?): Arg<T> = Arg.lit(value)
diff --git a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.kt 
b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.kt
index 4b42eb68d4c..47b44662b99 100644
--- a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.kt
+++ b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/Server.kt
@@ -33,8 +33,10 @@ import kotlinx.coroutines.runBlocking
 import org.apache.airflow.sdk.execution.CoordinatorComm
 import org.apache.airflow.sdk.execution.LogSender
 import org.apache.airflow.sdk.execution.Logger
+import org.apache.airflow.sdk.execution.comm.DagFileParseRequest
 import org.apache.airflow.sdk.execution.comm.ErrorResponse
 import org.apache.airflow.sdk.execution.comm.StartupDetails
+import org.apache.airflow.sdk.execution.parseDags
 import org.apache.airflow.sdk.execution.runTask
 import kotlin.text.substringAfterLast
 import kotlin.text.substringBeforeLast
@@ -181,6 +183,7 @@ class Server(
     val frame = coordinator.readMessage()
     when (val body = frame.body) {
       is StartupDetails -> runTaskAndReport(bundle, body, coordinator)
+      is DagFileParseRequest -> parseDagsAndReport(bundle, body, coordinator)
       is ErrorResponse -> throw ApiError("[${body.error}] ${body.detail}")
       else -> throw ApiError("Unexpected initial frame (id=${frame.id})")
     }
@@ -194,4 +197,12 @@ class Server(
     val result = runTask(bundle, startup, coordinator)
     coordinator.communicate<Unit>(result)
   }
+
+  private suspend fun parseDagsAndReport(
+    bundle: Bundle,
+    request: DagFileParseRequest,
+    coordinator: CoordinatorComm,
+  ) {
+    coordinator.communicate<Unit>(parseDags(bundle, request))
+  }
 }
diff --git 
a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt 
b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt
new file mode 100644
index 00000000000..89245cf6946
--- /dev/null
+++ b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/execution/Serde.kt
@@ -0,0 +1,490 @@
+/*
+ * 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, which the 
Go
+  // and TypeScript SDKs mirror too. Operator fields are stored unwrapped, so 
the
+  // __type encoding is stripped. If core grows a task-level 
fill_config_defaults,
+  // every SDK has to keep explicitly set values instead, or an explicit 
retries=0
+  // reads as unset and picks up the configured default.
+  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 same gap applies to [core] default_timezone, which Python's
+// _extract_tz uses for a Dag with no start date; this uses UTC.
+// 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 + CRON_ALIASES).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 *",
+  )
+
+/** One value in a cron field: a number, `*`, `?`, or a three-letter month or 
weekday name. */
+private const val CRON_VALUE = "(\\d+|\\*|\\?|[A-Z]{3})"
+
+/**
+ * One comma-separated element of a cron field: a value or a range, either
+ * stepped, with croniter's `L`, `W` and `#` qualifiers.
+ */
+private val CRON_ELEMENT =
+  Regex("$CRON_VALUE(-$CRON_VALUE)?([/#]\\d+)?[LW]*|L(-\\d+)?|LW", 
RegexOption.IGNORE_CASE)
+
+/**
+ * Aliases croniter accepts that `cron_presets` does not expand, so Python
+ * stores them unexpanded and so does this.
+ */
+private val CRON_ALIASES = setOf("@midnight", "@annually")
+
+/**
+ * Whether [expression] has the shape croniter accepts: five or six
+ * space-separated fields of cron characters, or an `@` alias.
+ *
+ * A shape check, not a parse: croniter validates the ranges, and repeating
+ * that here would be a second implementation to keep in step. What it catches
+ * is prose, such as `"every tuesday"`, which would otherwise be written into a
+ * Dag the scheduler then fails to build a timetable for.
+ */
+private fun isCronExpression(expression: String): Boolean {
+  val trimmed = expression.trim()
+  if (trimmed.startsWith("@")) return trimmed in CRON_PRESETS || trimmed in 
CRON_ALIASES
+  val fields = trimmed.split(Regex("\\s+"))
+  return fields.size in 5..6 &&
+    fields.all { field -> field.split(',').all { CRON_ELEMENT.matches(it) } }
+}
+
+/**
+ * Serializes the Dag's task groups as Python's `TaskGroupSerialization` does:
+ * a root group holding the tasks in no group and the top-level groups, each
+ * group nesting its own tasks and groups.
+ */
+private fun serializeTaskGroups(
+  dag: DagDef,
+  expansion: GroupExpansion,
+): Map<String, Any?> {
+  val grouped = dag.groups.values.flatMapTo(mutableSetOf()) { it.taskIds }
+  return taskGroupObject(
+    null,
+    dag.tasks.keys.filterNot { it in grouped },
+    dag.groups.values.filterNot { '.' in it.id },
+    expansion,
+  )
+}
+
+/** One group object: the root when [group] is null, otherwise a nested one. */
+private fun taskGroupObject(
+  group: TaskGroupRef?,
+  taskIds: List<String>,
+  children: List<TaskGroupRef>,
+  expansion: GroupExpansion,
+): Map<String, Any?> =
+  mapOf(
+    // The local segment: Python's TaskGroup stores the ID it was given, and
+    // rebuilds the full one from where the group sits in the tree.
+    "_group_id" to group?.id?.substringAfterLast('.'),
+    "group_display_name" to "",
+    "prefix_group_id" to true,
+    "tooltip" to "",
+    "ui_color" to "CornflowerBlue",
+    "ui_fgcolor" to "#000",
+    "children" to
+      taskIds.associateWith { listOf("operator", it) } +
+      children.associate {
+        it.id to listOf("taskgroup", taskGroupObject(it, it.taskIds, 
it.children, expansion))
+      },
+    "upstream_group_ids" to group.edges(expansion).upstreamGroupIds.sorted(),
+    "downstream_group_ids" to 
group.edges(expansion).downstreamGroupIds.sorted(),
+    "upstream_task_ids" to group.edges(expansion).upstreamTaskIds.sorted(),
+    "downstream_task_ids" to group.edges(expansion).downstreamTaskIds.sorted(),
+  )
+
+/** The edges this group records for itself; the root group records none. */
+private fun TaskGroupRef?.edges(expansion: GroupExpansion): GroupEdges = 
this?.let { expansion.edgesOf(it.id) } ?: GroupEdges()
+
+/**
+ * Recursively serializes a value with Airflow's type/var encoding, matching
+ * Python's `BaseSerialization.serialize()` output: primitives pass through,
+ * date-times become `{"__type": "datetime", "__var": epoch_seconds}`,
+ * durations `{"__type": "timedelta", "__var": total_seconds}`, and maps
+ * `{"__type": "dict", "__var": {...}}`.
+ */
+internal fun serializeValue(value: Any?): Any? =
+  when (value) {
+    null -> null
+    is String, is Boolean, is Int, is Long, is Double -> value
+    is Byte, is Short -> (value as Number).toInt()
+    is Float -> value.toDouble()
+    is OffsetDateTime -> serializeValue(value.toInstant())
+    is Instant ->
+      mapOf(
+        "__type" to "datetime",
+        "__var" to value.epochSecond + value.nano / 1e9,
+      )
+    is Duration ->
+      mapOf(
+        "__type" to "timedelta",
+        "__var" to value.toNanos() / 1e9,
+      )
+    is Map<*, *> ->
+      mapOf(
+        "__type" to "dict",
+        "__var" to value.entries.associate { (k, v) -> k.toString() to 
serializeValue(v) },
+      )
+    is List<*> -> value.map(::serializeValue)
+    is Array<*> -> value.map(::serializeValue)
+    else -> value
+  }
+
+/**
+ * Extracts the `__var` part from a type-encoded value: in Python's
+ * `serialize_to_json`, non-decorated fields are serialized then unwrapped.
+ */
+internal fun unwrapTypeEncoding(value: Any?): Any? {
+  val map = value as? Map<*, *> ?: return value
+  if ("__type" !in map) return value
+  return if ("__var" in map) map["__var"] else value
+}
+
+/** Whether a config value equals the schema default and can be omitted. */
+private fun matchesSchemaDefault(
+  field: Field?,
+  value: Any,
+): Boolean {
+  val defaultJson = field?.defaultJson ?: return false
+  val node = defaultsMapper.readTree(defaultJson)
+  return when (value) {
+    is String -> node.isTextual && node.asText() == value
+    is Boolean -> node.isBoolean && node.asBoolean() == value
+    is Number -> node.isNumber && node.asDouble() == value.toDouble()
+    is Duration -> node.isNumber && node.asDouble() == value.toNanos() / 1e9
+    else -> false
+  }
+}
+
+private fun computeRelativeFileloc(
+  fileloc: String,
+  bundlePath: String?,
+): String {
+  if (fileloc.isEmpty()) return ""
+  if (bundlePath.isNullOrEmpty()) return "."
+  return try {
+    Paths
+      .get(bundlePath)
+      .relativize(Paths.get(fileloc))
+      .toString()
+      .ifEmpty { "." }
+  } catch (e: InvalidPathException) {
+    "."
+  } catch (e: IllegalArgumentException) {
+    "."
+  }
+}
diff --git 
a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/ArgValues.kt 
b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/ArgValues.kt
index 559e975645b..e3920f2b67f 100644
--- a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/ArgValues.kt
+++ b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/ArgValues.kt
@@ -51,10 +51,9 @@ import java.lang.reflect.Type
  * binding at their position (through [TaskArgs]); [TaskInput] fields resolve
  * bindings by name.
  *
- * A natively authored Dag has no stub call site, so the supervisor sends no
- * bindings for it and the inputs the Dag itself wired stand in. When the
- * supervisor sends bindings they are used for every parameter; the Dag's own
- * inputs are read only when it sends none.
+ * A natively authored Dag resolves its own arguments instead: the bundle that
+ * holds the task also holds the call that wired it, so a task that wired any
+ * input reads those and never the bindings.
  *
  * A count that does not match is fatal for flat parameters and a warning for a
  * [TaskInput]: a position has no name to fall back on, while a field does, so
@@ -92,7 +91,7 @@ object ArgValues {
     // Runtime bindings carry argument names to match fields against. A wired
     // input carries none, so it decodes into the whole input at once -- which
     // is well defined because a TaskInput is a task's only data parameter.
-    wiredInputs(context, client)?.let { wired ->
+    wiredInputs(context)?.let { wired ->
       warnWiredArity(client, type, wired.size)
       return type.cast(decode(resolveWiredAll(wired.take(1), client).single(), 
type))
         ?: throw missingInput(wired[0], type.simpleName)
@@ -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() }
 
   /**
    * The failure for a wired argument that resolved to nothing where a value is
diff --git 
a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/Refs.kt 
b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/Refs.kt
index 387e6fec365..761d5b8f282 100644
--- a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/Refs.kt
+++ b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/Refs.kt
@@ -92,14 +92,16 @@ object Refs {
   fun <T> node(
     groupId: String,
     def: TaskDef,
-  ): TaskRef<T> = call(groupId, def)
+  ): TaskRef<T> = call(groupId, def, emptyList())
 
   /**
-   * Records a task and the data edge for every [TaskRef] among [args]; a
-   * literal argument records a baked value and no edge.
+   * Records a task called with named arguments, as a generated wiring view
+   * does. [names] is the task parameter each argument feeds, in the same
+   * order, and is what the serialized Dag carries as the task's binding spec.
    *
    * @param groupId Full ID of the task group holding it, empty when it sits in
-   *    none.
+   *    none. The generated wiring view knows which it is.
+   *
    * @return The handle representing this task, memoized by [TaskDef.id] so a 
result
    *    held in a local and reused refers to one node.
    * @throws IllegalArgumentException if an argument is a raw Java `null`
@@ -110,8 +112,12 @@ object Refs {
   fun <T> call(
     groupId: String,
     def: TaskDef,
+    names: List<String>,
     vararg args: Arg<*>?,
   ): TaskRef<T> {
+    require(names.isEmpty() || names.size == args.size) {
+      "Task '${def.id}' was wired with ${args.size} argument(s) under 
${names.size} name(s)"
+    }
     val inputs =
       args.mapIndexed { i, arg ->
         requireNotNull(arg) {
@@ -131,6 +137,7 @@ object Refs {
     }
     inputs.filterIsInstance<TaskRef<*>>().forEach { def.dependsOn(it.def) }
     def.inputs += inputs
+    def.inputNames += names
     if (groupId.isEmpty()) {
       active.dag.addTask(def)
     } else {
diff --git 
a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/TaskArgs.kt 
b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/TaskArgs.kt
index 14c6404758b..c4bd53a5e14 100644
--- a/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/TaskArgs.kt
+++ b/java-sdk/sdk/src/main/kotlin/org/apache/airflow/sdk/internal/TaskArgs.kt
@@ -75,7 +75,7 @@ class TaskArgs private constructor(
       client: Client,
       declared: Int,
     ): TaskArgs {
-      ArgValues.wiredInputs(context, client)?.let { wired ->
+      ArgValues.wiredInputs(context)?.let { wired ->
         // Positional binding is strict in both directions: a parameter has no
         // name to fall back on, so a count that does not match cannot be
         // resolved and the run fails here rather than mid-task.
diff --git 
a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/ArgTestSupport.kt 
b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/ArgTestSupport.kt
index 50b0a8ef3ae..0929dfe4e4e 100644
--- a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/ArgTestSupport.kt
+++ b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/ArgTestSupport.kt
@@ -129,6 +129,6 @@ internal class NoopTask : Task {
 internal fun contextWiredWith(inputs: List<Arg<*>>): Context {
   val dag = DagDef("d")
   val def = TaskDef("t", NoopTask::class.java)
-  Refs.record(dag, listOf("t"), emptyList()) { Refs.call<Unit>("", def, 
*inputs.toTypedArray()) }
+  Refs.record(dag, listOf("t"), emptyList()) { Refs.call<Unit>("", def, 
emptyList(), *inputs.toTypedArray()) }
   return taskContext().also { it.taskDef = def }
 }
diff --git 
a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/ArgValuesTest.kt 
b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/ArgValuesTest.kt
index b6b5347ba4e..3c0a6fed279 100644
--- a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/ArgValuesTest.kt
+++ b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/ArgValuesTest.kt
@@ -121,6 +121,17 @@ internal class ArgValuesTest {
     assertEquals(0.5, input.threshold)
   }
 
+  @Test
+  @DisplayName("Should bind the Dag's wired input over the runtime bindings")
+  fun shouldPreferWiredInputOverRuntimeBindings() {
+    val (client, _) = clientWith(listOf(literal("threshold", 9.0)))
+    val context = contextWiredWith(listOf(Arg.lit(mapOf("threshold" to 0.5))))
+
+    val input = ArgValues.bindInput(context, client, 
ThresholdInput::class.java)
+
+    assertEquals(0.5, input.threshold)
+  }
+
   @Test
   @DisplayName("Should match a camelCase field to a snake_case argument 
through the fold")
   fun shouldFoldSnakeCaseArgument() {
diff --git a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/ServerTest.kt 
b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/ServerTest.kt
index e436dd2fc61..d6e93fe2114 100644
--- a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/ServerTest.kt
+++ b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/ServerTest.kt
@@ -27,16 +27,25 @@ import kotlinx.coroutines.runBlocking
 import org.apache.airflow.sdk.execution.CoordinatorComm
 import org.apache.airflow.sdk.execution.Frame
 import org.apache.airflow.sdk.execution.IncomingFrame
+import org.apache.airflow.sdk.execution.RawFrame
 import org.apache.airflow.sdk.execution.comm.TaskState
 import org.junit.jupiter.api.Assertions
 import org.junit.jupiter.api.DisplayName
 import org.junit.jupiter.api.Test
 import org.junit.jupiter.api.Timeout
 import org.msgpack.core.MessagePack
+import org.msgpack.core.buffer.ArrayBufferInput
 import java.io.ByteArrayOutputStream
 import java.util.concurrent.ArrayBlockingQueue
 import java.util.concurrent.TimeUnit
 
+private class ServerNoopTask : Task {
+  override fun execute(
+    context: Context,
+    client: Client,
+  ) = Unit
+}
+
 class ServerTest {
   private fun hexToBytes(hex: String): ByteArray =
     hex
@@ -112,6 +121,60 @@ class ServerTest {
     comm.close()
   }
 
+  @Test
+  @DisplayName("Should serialize the bundle's dags when the initial frame is a 
parse request")
+  @Timeout(value = 30, unit = TimeUnit.SECONDS)
+  fun parsesDagsAndReportsResult() {
+    val toServer = ByteChannel(autoFlush = true)
+    val fromServer = ByteChannel(autoFlush = true)
+    val comm = CoordinatorComm(toServer, fromServer)
+    val server = Server(InetSocketAddress("localhost", 0), 
InetSocketAddress("localhost", 0))
+    val bundle = Bundle(listOf(DagDef("parsed_dag").addTask(TaskDef("t", 
ServerNoopTask::class.java))))
+
+    val reported = ArrayBlockingQueue<RawFrame>(1)
+    val supervisor =
+      Thread {
+        runBlocking {
+          toServer.writeFrame(parseRequestFrame(3, "/bundle/dags/java.py", 
"/bundle"))
+          val prefix = fromServer.readByteArray(4)
+          val payload = 
fromServer.readByteArray(Frame.parseLengthPrefix(prefix).toInt())
+          val raw = Frame.decodeRaw(ArrayBufferInput(payload))
+          reported.put(raw)
+          toServer.writeFrame(ackFrame(raw.id))
+        }
+      }
+    supervisor.start()
+
+    runBlocking { server.dispatchTask(bundle, comm) }
+    supervisor.join()
+
+    val body = reported.take().rawBody as Map<*, *>
+    Assertions.assertEquals("DagFileParsingResult", body["type"])
+    Assertions.assertEquals("/bundle/dags/java.py", body["fileloc"])
+    val dags = body["serialized_dags"] as List<*>
+    Assertions.assertEquals(1, dags.size)
+    val dag = ((dags[0] as Map<*, *>)["data"] as Map<*, *>)["dag"] as Map<*, *>
+    Assertions.assertEquals("parsed_dag", dag["dag_id"])
+    comm.close()
+  }
+
+  private fun parseRequestFrame(
+    id: Int,
+    file: String,
+    bundlePath: String,
+  ): ByteArray {
+    val out = ByteArrayOutputStream()
+    MessagePack.newDefaultPacker(out).use { packer ->
+      packer.packArrayHeader(2)
+      packer.packInt(id)
+      packer.packMapHeader(3)
+      packer.packString("type").packString("DagFileParseRequest")
+      packer.packString("file").packString(file)
+      packer.packString("bundle_path").packString(bundlePath)
+    }
+    return out.toByteArray()
+  }
+
   private companion object {
     // [2, msg, null] with msg coming from
     // 
https://github.com/astronomer/airflow/blob/f39c8da8/task-sdk/tests/task_sdk/execution_time/test_comms.py#L73-L108
diff --git 
a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/conformance/SerializeJava.kt
 
b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/conformance/SerializeJava.kt
new file mode 100644
index 00000000000..7d9aa937010
--- /dev/null
+++ 
b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/conformance/SerializeJava.kt
@@ -0,0 +1,140 @@
+/*
+ * 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.
+ */
+
+// Serializes the shared test Dags with this SDK, for
+// scripts/ci/lang_sdk_serialization/compare.py:
+//
+//   java -cp <sdk test runtime classpath> 
org.apache.airflow.sdk.conformance.SerializeJavaKt \
+//       scripts/ci/lang_sdk_serialization/test_dags.yaml serialized_java.json
+//
+// Each Dag is built with the interface API and written as the runtime answers
+// a parse request, keyed by Dag ID. A task's `upstream` becomes an ordering
+// edge, which a Java task serializes the same way as a data edge.
+package org.apache.airflow.sdk.conformance
+
+import com.fasterxml.jackson.databind.JsonNode
+import com.fasterxml.jackson.databind.ObjectMapper
+import com.fasterxml.jackson.dataformat.yaml.YAMLMapper
+import org.apache.airflow.sdk.Arg
+import org.apache.airflow.sdk.Bundle
+import org.apache.airflow.sdk.Client
+import org.apache.airflow.sdk.Context
+import org.apache.airflow.sdk.DagDef
+import org.apache.airflow.sdk.Deps
+import org.apache.airflow.sdk.Task
+import org.apache.airflow.sdk.TaskGroupRef
+import org.apache.airflow.sdk.TaskRef
+import org.apache.airflow.sdk.execution.serializeDag
+import org.apache.airflow.sdk.internal.Field
+import org.apache.airflow.sdk.internal.FieldType
+import org.apache.airflow.sdk.internal.SchemaFields
+import java.io.File
+import java.time.Duration
+import java.time.OffsetDateTime
+
+class ConformanceTask : Task {
+  override fun execute(
+    context: Context,
+    client: Client,
+  ) = Unit
+}
+
+fun main(args: Array<String>) {
+  require(args.size == 2) { "usage: SerializeJava <test_dags.yaml> 
<output.json>" }
+  val cases = YAMLMapper().readTree(File(args[0])).path("dags")
+  val bundle = Bundle()
+  cases.forEach { bundle.register(buildDag(it)) }
+  val serialized =
+    bundle.dags.values.associate { dag ->
+      dag.id to mapOf("__version" to 3, "dag" to serializeDag(dag, "", "."))
+    }
+  ObjectMapper().writerWithDefaultPrettyPrinter().writeValue(File(args[1]), 
serialized)
+}
+
+private fun buildDag(case: JsonNode): DagDef {
+  val dag = DagDef(case.path("dag_id").asText())
+  case.path("spec").fields().forEach { (key, value) -> dag.config(key, 
toValue(SchemaFields.DAG, key, value)) }
+
+  // A group ID is fully qualified, so its parent is whatever comes before the 
last dot.
+  val groups = linkedMapOf<String, TaskGroupRef>()
+  case.path("groups").forEach { node ->
+    val groupId = node.asText()
+    val parentId = groupId.substringBeforeLast('.', "")
+    val localId = groupId.substringAfterLast('.')
+    groups[groupId] = if (parentId.isEmpty()) dag.taskGroup(localId) else 
groups.getValue(parentId).taskGroup(localId)
+  }
+
+  val tasks = linkedMapOf<String, TaskRef<Unit>>()
+  case.path("tasks").forEach { task ->
+    val groupId = task.path("group").asText("")
+    val localId = task.path("task_id").asText()
+    val ref =
+      if (groupId.isEmpty()) {
+        dag.task<Unit>(localId, ConformanceTask::class.java)
+      } else {
+        groups.getValue(groupId).task(localId, ConformanceTask::class.java)
+      }
+    task.path("spec").fields().forEach { (key, value) -> ref.config(key, 
toValue(SchemaFields.TASK, key, value)) }
+    // A task's `upstream` handles and its `literals` are its call arguments, 
in that order, so the
+    // Dag carries the binding spec a stub call would. Names are positional, 
as the Go SDK names
+    // them: the interface API has no signature to read parameter names from.
+    val inputs: List<Arg<*>> =
+      task.path("upstream").map { tasks.getValue(it.asText()) } +
+        task.path("literals").map { Arg.lit(toJsonValue(it)) }
+    inputs.filterIsInstance<TaskRef<*>>().forEach { ref.after(it) }
+    ref.def.inputs += inputs
+    ref.def.inputNames += inputs.indices.map { "arg$it" }
+    tasks[ref.def.id] = ref
+  }
+
+  case.path("order_edges").forEach { edge ->
+    val node = { id: String -> groups[id] as Deps.Flow? ?: tasks.getValue(id) }
+    node(edge[0].asText()).before(node(edge[1].asText()))
+  }
+  return dag
+}
+
+/** Reads a YAML value as the Java type the config key takes. */
+private fun toValue(
+  table: Map<String, Field>,
+  key: String,
+  node: JsonNode,
+): Any =
+  when (requireNotNull(table[key]) { "Unknown config key: '$key'" }.type) {
+    FieldType.STRING -> node.asText()
+    FieldType.BOOLEAN -> node.asBoolean()
+    FieldType.INTEGER -> node.asInt()
+    FieldType.NUMBER -> node.numberValue()
+    // `!datetime` is an ISO 8601 timestamp, and `!timedelta` a number of 
seconds.
+    FieldType.DATETIME -> OffsetDateTime.parse(node.asText())
+    FieldType.TIMEDELTA -> Duration.ofNanos((node.asText().toDouble() * 
1e9).toLong())
+    FieldType.STRING_ARRAY -> node.map { it.asText() }
+  }
+
+/** Reads a YAML literal as the plain value `Serde` writes out. */
+private fun toJsonValue(node: JsonNode): Any? =
+  when {
+    node.isNull -> null
+    node.isTextual -> node.asText()
+    node.isBoolean -> node.asBoolean()
+    node.isIntegralNumber -> node.numberValue()
+    node.isNumber -> node.asDouble()
+    node.isArray -> node.map { toJsonValue(it) }
+    else -> node.fields().asSequence().associate { (key, value) -> key to 
toJsonValue(value) }
+  }
diff --git 
a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/execution/SerdeTest.kt 
b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/execution/SerdeTest.kt
new file mode 100644
index 00000000000..e025d9a5b15
--- /dev/null
+++ b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/execution/SerdeTest.kt
@@ -0,0 +1,428 @@
+/*
+ * 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 org.apache.airflow.sdk.Arg
+import org.apache.airflow.sdk.Bundle
+import org.apache.airflow.sdk.Client
+import org.apache.airflow.sdk.Context
+import org.apache.airflow.sdk.DagDef
+import org.apache.airflow.sdk.Task
+import org.apache.airflow.sdk.TaskDef
+import org.apache.airflow.sdk.execution.comm.DagFileParseRequest
+import org.apache.airflow.sdk.internal.Refs
+import org.junit.jupiter.api.Assertions.assertEquals
+import org.junit.jupiter.api.Assertions.assertFalse
+import org.junit.jupiter.api.Assertions.assertNull
+import org.junit.jupiter.api.Assertions.assertThrows
+import org.junit.jupiter.api.DisplayName
+import org.junit.jupiter.api.Test
+import java.time.Duration
+import java.time.OffsetDateTime
+
+private class SerdeNoopTask : Task {
+  override fun execute(
+    context: Context,
+    client: Client,
+  ) = Unit
+}
+
+@Suppress("UNCHECKED_CAST")
+private fun taskData(
+  serialized: Map<String, Any?>,
+  index: Int,
+): Map<String, Any?> {
+  val tasks = serialized["tasks"] as List<Map<String, Any?>>
+  assertEquals("operator", tasks[index]["__type"])
+  return tasks[index]["__var"] as Map<String, Any?>
+}
+
+internal class SerdeTest {
+  @Test
+  @DisplayName("Should emit required dag fields and leave unset config-backed 
fields out")
+  fun shouldEmitRequiredDagFields() {
+    val serialized = serializeDag(DagDef("d"), "/bundles/app/dags.jar", 
"app/dags.jar")
+
+    assertEquals("d", serialized["dag_id"])
+    assertEquals("/bundles/app/dags.jar", serialized["fileloc"])
+    assertEquals("app/dags.jar", serialized["relative_fileloc"])
+    assertEquals("UTC", serialized["timezone"])
+    assertEquals(
+      mapOf("__type" to "airflow.timetables.simple.NullTimetable", "__var" to 
emptyMap<String, Any?>()),
+      serialized["timetable"],
+    )
+    assertEquals(emptyList<Any?>(), serialized["tasks"])
+    assertEquals(emptyList<Any?>(), serialized["dag_dependencies"])
+    assertEquals(emptyMap<String, Any?>(), serialized["edge_info"])
+    assertEquals(emptyList<Any?>(), serialized["params"])
+    assertNull(serialized["deadline"])
+    assertNull(serialized["allowed_run_types"])
+    assertFalse("max_active_tasks" in serialized)
+    assertFalse("max_active_runs" in serialized)
+    assertFalse("max_consecutive_failed_dag_runs" in serialized)
+    assertFalse("catchup" in serialized)
+    assertFalse("disable_bundle_versioning" in serialized)
+    assertFalse("description" in serialized)
+    assertFalse("fail_fast" in serialized)
+    assertFalse("tags" in serialized)
+  }
+
+  @Test
+  @DisplayName("Should map schedule strings to the matching timetable")
+  fun shouldMapScheduleToTimetable() {
+    val cron = serializeDag(DagDef("d").config("schedule", "@daily"), "", ".")
+    assertEquals(
+      mapOf(
+        "__type" to "airflow.timetables.trigger.CronTriggerTimetable",
+        "__var" to
+          mapOf(
+            "expression" to "0 0 * * *",
+            "timezone" to "UTC",
+            "interval" to 0.0,
+            "run_immediately" to false,
+          ),
+      ),
+      cron["timetable"],
+    )
+
+    val once = serializeDag(DagDef("d").config("schedule", "@once"), "", ".")
+    assertEquals(
+      mapOf("__type" to "airflow.timetables.simple.OnceTimetable", "__var" to 
emptyMap<String, Any?>()),
+      once["timetable"],
+    )
+
+    val continuous = serializeDag(DagDef("d").config("schedule", 
"@continuous"), "", ".")
+    assertEquals(
+      mapOf("__type" to "airflow.timetables.simple.ContinuousTimetable", 
"__var" to emptyMap<String, Any?>()),
+      continuous["timetable"],
+    )
+  }
+
+  @Test
+  @DisplayName("Should apply dag config values with Python's emit rules")
+  fun shouldApplyDagConfig() {
+    val dag =
+      DagDef("d")
+        .config("description", "demo")
+        .config("tags", listOf("b", "a", "b"))
+        .config("catchup", true)
+        .config("fail_fast", true)
+        .config("max_active_runs", 3)
+        .config("dagrun_timeout", Duration.ofMinutes(5))
+        .config("start_date", OffsetDateTime.parse("2026-01-01T00:00:00Z"))
+
+    val serialized = serializeDag(dag, "", ".")
+
+    assertEquals("demo", serialized["description"])
+    assertEquals(listOf("a", "b"), serialized["tags"])
+    assertEquals(true, serialized["catchup"])
+    assertEquals(true, serialized["fail_fast"])
+    assertEquals(3, serialized["max_active_runs"])
+    assertEquals(300.0, serialized["dagrun_timeout"])
+    assertEquals(1.7672256E9, serialized["start_date"])
+  }
+
+  @Test
+  @DisplayName("Should serialize tasks with identity fields, config, and 
sorted downstream ids")
+  fun shouldSerializeTasks() {
+    val extractDef =
+      TaskDef("extract", SerdeNoopTask::class.java)
+        .config("retries", 2)
+        .config("queue", "q")
+        .config("retry_delay", Duration.ofMinutes(10))
+        .config("email_on_failure", false)
+        .config("email_on_retry", false)
+    val transformDef =
+      TaskDef("transform", SerdeNoopTask::class.java)
+        .dependsOn(extractDef)
+        // Explicitly at schema defaults: omitted from the serialized form.
+        .config("retries", 0)
+        .config("queue", "default")
+        .config("retry_delay", Duration.ofMinutes(5))
+    val dag = DagDef("d").addTask(extractDef).addTask(transformDef)
+
+    val serialized = serializeDag(dag, "", ".")
+
+    val extract = taskData(serialized, 0)
+    assertEquals("extract", extract["task_id"])
+    assertEquals("SerdeNoopTask", extract["task_type"])
+    assertEquals("org.apache.airflow.sdk.execution", extract["_task_module"])
+    assertEquals("java", extract["language"])
+    assertEquals(emptyList<Any?>(), extract["template_fields"])
+    assertEquals(2, extract["retries"])
+    assertEquals("q", extract["queue"])
+    assertEquals(600.0, extract["retry_delay"])
+    assertEquals(listOf("transform"), extract["downstream_task_ids"])
+    assertFalse("email_on_failure" in extract)
+    assertFalse("email_on_retry" in extract)
+
+    val transform = taskData(serialized, 1)
+    assertEquals("transform", transform["task_id"])
+    assertFalse("retries" in transform)
+    assertFalse("queue" in transform)
+    assertFalse("retry_delay" in transform)
+    assertFalse("downstream_task_ids" in transform)
+
+    assertEquals(
+      mapOf("extract" to listOf("operator", "extract"), "transform" to 
listOf("operator", "transform")),
+      (serialized["task_group"] as Map<*, *>)["children"],
+    )
+  }
+
+  @Test
+  @DisplayName("Should serialize nested task groups with their own edges")
+  fun shouldSerializeTaskGroups() {
+    val dag = DagDef("d")
+    val extract = dag.task<Unit>("extract", SerdeNoopTask::class.java)
+    val staging = dag.taskGroup("staging")
+    staging.task<Unit>("stage", SerdeNoopTask::class.java)
+    staging.taskGroup("checks").task<Unit>("nulls", SerdeNoopTask::class.java)
+    extract.before(staging)
+
+    val root = serializeDag(dag, "", ".")["task_group"] as Map<*, *>
+
+    val group = { name: String, children: Map<String, Any?>, upstreamTasks: 
List<String> ->
+      mapOf(
+        "_group_id" to name,
+        "group_display_name" to "",
+        "prefix_group_id" to true,
+        "tooltip" to "",
+        "ui_color" to "CornflowerBlue",
+        "ui_fgcolor" to "#000",
+        "children" to children,
+        "upstream_group_ids" to emptyList<String>(),
+        "downstream_group_ids" to emptyList<String>(),
+        "upstream_task_ids" to upstreamTasks,
+        "downstream_task_ids" to emptyList<String>(),
+      )
+    }
+    val checks = group("checks", mapOf("staging.checks.nulls" to 
listOf("operator", "staging.checks.nulls")), emptyList())
+    assertEquals(
+      mapOf(
+        "extract" to listOf("operator", "extract"),
+        "staging" to
+          listOf(
+            "taskgroup",
+            group(
+              "staging",
+              mapOf(
+                "staging.stage" to listOf("operator", "staging.stage"),
+                "staging.checks" to listOf("taskgroup", checks),
+              ),
+              listOf("extract"),
+            ),
+          ),
+      ),
+      root["children"],
+    )
+    assertEquals(null, root["_group_id"])
+  }
+
+  @Test
+  @DisplayName("Should serialize wiring-registered dags with their data-flow 
edges")
+  fun shouldSerializeWiredDag() {
+    val dag = DagDef("d")
+    Refs.record(dag, listOf("extract", "transform"), emptyList()) {
+      val extracted = Refs.node<Long>("", TaskDef("extract", 
SerdeNoopTask::class.java))
+      Refs.call<Unit>("", TaskDef("transform", SerdeNoopTask::class.java), 
listOf("rows"), extracted)
+    }
+
+    val serialized = serializeDag(dag, "", ".")
+
+    assertEquals(listOf("transform"), taskData(serialized, 
0)["downstream_task_ids"])
+    assertEquals(
+      listOf(mapOf("name" to "rows", "kind" to "xcom", "task_id" to 
"extract")),
+      taskData(serialized, 1)["_arg_bindings"],
+    )
+    assertEquals(true, taskData(serialized, 1)["is_stub"])
+  }
+
+  @Test
+  @DisplayName("Should bind a wired literal as a literal argument binding")
+  fun shouldSerializeLiteralArgBinding() {
+    val dag = DagDef("d")
+    Refs.record(dag, listOf("transform"), emptyList()) {
+      Refs.call<Unit>("", TaskDef("transform", SerdeNoopTask::class.java), 
listOf("region"), Arg.lit("uk"))
+    }
+
+    val serialized = serializeDag(dag, "", ".")
+
+    assertEquals(
+      listOf(mapOf("name" to "region", "kind" to "literal", "value" to "uk")),
+      taskData(serialized, 0)["_arg_bindings"],
+    )
+  }
+
+  @Test
+  @DisplayName("Should leave out the binding spec of a task called with no 
arguments")
+  fun shouldOmitArgBindingsWithoutArguments() {
+    val serialized =
+      serializeDag(DagDef("d").addTask(TaskDef("t", 
SerdeNoopTask::class.java)), "", ".")
+
+    assertFalse("_arg_bindings" in taskData(serialized, 0))
+    assertEquals(true, taskData(serialized, 0)["is_stub"])
+  }
+
+  @Test
+  @DisplayName("Should reject a wired literal that has no JSON form")
+  fun shouldRejectNonJsonLiteral() {
+    val dag = DagDef("d")
+    Refs.record(dag, listOf("t"), emptyList()) {
+      Refs.call<Unit>("", TaskDef("t", SerdeNoopTask::class.java), 
listOf("at"), Arg.lit(Duration.ofSeconds(5)))
+    }
+
+    val error = assertThrows(IllegalArgumentException::class.java) { 
serializeDag(dag, "", ".") }
+
+    assertEquals(
+      "Argument 'at' of task 't' is a java.time.Duration, which has no JSON 
form; the Dag's call " +
+        "arguments travel as JSON, so pass a string, number, boolean, list, or 
map",
+      error.message,
+    )
+  }
+
+  @Test
+  @DisplayName("Should take the dag timezone from the start date, as Python 
does")
+  fun shouldTakeTimezoneFromStartDate() {
+    val dag =
+      DagDef("d")
+        .config("schedule", "0 3 * * *")
+        .config("start_date", 
OffsetDateTime.parse("2026-01-01T00:00:00+05:30"))
+
+    val serialized = serializeDag(dag, "", ".")
+
+    assertEquals(19800, serialized["timezone"])
+    assertEquals(19800, (serialized["timetable"] as Map<*, *>)["__var"].let { 
(it as Map<*, *>)["timezone"] })
+  }
+
+  private fun cronExpression(schedule: String): Any? {
+    val timetable = serializeDag(DagDef("d").config("schedule", schedule), "", 
".")["timetable"] as Map<*, *>
+    return (timetable["__var"] as Map<*, *>)["expression"]
+  }
+
+  @Test
+  @DisplayName("Should accept every cron schedule croniter does")
+  fun shouldAcceptCroniterSchedules() {
+    listOf(
+      "0 9 * * MON,WED,FRI",
+      "0 0 * JAN,JUL *",
+      "0 0 * * MON#2",
+      "0 0 15W * *",
+      "*/5 1-5/2 * * MON-FRI",
+    ).forEach { assertEquals(it, cronExpression(it)) }
+  }
+
+  @Test
+  @DisplayName("Should serialize croniter's @midnight and @annually aliases 
unexpanded")
+  fun shouldKeepCronAliasesUnexpanded() {
+    assertEquals("@midnight", cronExpression("@midnight"))
+    assertEquals("@annually", cronExpression("@annually"))
+    assertEquals("0 0 * * *", cronExpression("@daily"))
+  }
+
+  @Test
+  @DisplayName("Should reject schedules that are not cron expressions")
+  fun shouldRejectProseAndUnknownAliases() {
+    listOf("every tuesday", "@bogus", "0 0 * * tuesday").forEach {
+      assertThrows(IllegalArgumentException::class.java) { cronExpression(it) }
+    }
+  }
+
+  @Test
+  @DisplayName("Should reject a schedule the scheduler cannot build a 
timetable from")
+  fun shouldRejectNonCronSchedule() {
+    val dag = DagDef("d").config("schedule", "every monday")
+
+    val error = assertThrows(IllegalArgumentException::class.java) { 
serializeDag(dag, "", ".") }
+
+    assertEquals(
+      "Schedule 'every monday' of Dag 'd' is not a cron expression or a preset 
" +
+        "(@hourly, @daily, @weekly, @monthly, @quarterly, @yearly, @midnight, 
@annually, @once, @continuous); a schedule the " +
+        "scheduler cannot parse would leave the Dag unschedulable",
+      error.message,
+    )
+  }
+
+  @Test
+  @DisplayName("Should report a dag that cannot be serialized as an import 
error")
+  fun shouldReportUnserializableDagAsImportError() {
+    val broken = DagDef("broken").config("schedule", "every monday")
+    val healthy = DagDef("healthy").addTask(TaskDef("t", 
SerdeNoopTask::class.java))
+    val request =
+      DagFileParseRequest().also {
+        it.file = "/bundles/app/dags.jar"
+        it.bundlePath = "/bundles"
+      }
+
+    val result = parseDags(Bundle(listOf(broken, healthy)), request)
+
+    assertEquals(1, (result["serialized_dags"] as List<*>).size)
+    assertEquals(
+      mapOf(
+        "app/dags.jar" to
+          "Dag \"broken\": Schedule 'every monday' of Dag 'broken' is not a 
cron expression or a preset " +
+          "(@hourly, @daily, @weekly, @monthly, @quarterly, @yearly, 
@midnight, @annually, @once, @continuous); a schedule the " +
+          "scheduler cannot parse would leave the Dag unschedulable",
+      ),
+      result["import_errors"],
+    )
+  }
+
+  @Test
+  @DisplayName("Should wrap parsed dags in a DagFileParsingResult body")
+  fun shouldBuildParsingResult() {
+    val bundle = Bundle(listOf(DagDef("d").addTask(TaskDef("t", 
SerdeNoopTask::class.java))))
+    val request =
+      DagFileParseRequest().also {
+        it.file = "/bundles/app/dags.jar"
+        it.bundlePath = "/bundles"
+      }
+
+    val result = parseDags(bundle, request)
+
+    assertEquals("DagFileParsingResult", result["type"])
+    assertEquals("/bundles/app/dags.jar", result["fileloc"])
+    val dags = result["serialized_dags"] as List<*>
+    assertEquals(1, dags.size)
+    val data = (dags[0] as Map<*, *>)["data"] as Map<*, *>
+    assertEquals(3, data["__version"])
+    val dag = data["dag"] as Map<*, *>
+    assertEquals("d", dag["dag_id"])
+    assertEquals("app/dags.jar", dag["relative_fileloc"])
+  }
+
+  @Test
+  @DisplayName("Should encode temporals and nested maps with the type/var 
envelope")
+  fun shouldEncodeValuesWithTypeEnvelope() {
+    assertEquals(
+      mapOf("__type" to "timedelta", "__var" to 90.0),
+      serializeValue(Duration.ofSeconds(90)),
+    )
+    assertEquals(
+      mapOf("__type" to "datetime", "__var" to 1.7672256E9),
+      serializeValue(OffsetDateTime.parse("2026-01-01T00:00:00Z")),
+    )
+    assertEquals(
+      mapOf("__type" to "dict", "__var" to mapOf("k" to listOf(1, 2))),
+      serializeValue(mapOf("k" to listOf(1, 2))),
+    )
+    assertEquals(42, unwrapTypeEncoding(mapOf("__type" to "timedelta", "__var" 
to 42)))
+    assertEquals(mapOf("plain" to 1), unwrapTypeEncoding(mapOf("plain" to 1)))
+  }
+}
diff --git 
a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/internal/ArgValuesTest.kt 
b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/internal/ArgValuesTest.kt
index fad223f605f..b9fd480e10c 100644
--- 
a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/internal/ArgValuesTest.kt
+++ 
b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/internal/ArgValuesTest.kt
@@ -141,7 +141,7 @@ internal class ArgValuesTest {
       .distinct()
       .forEach { dag.addTask(it) }
     val def = TaskDef("consumer", NoopArgTask::class.java)
-    Refs.record(dag, listOf("consumer"), emptyList()) { Refs.call<Unit>("", 
def, *inputs.toTypedArray()) }
+    Refs.record(dag, listOf("consumer"), emptyList()) { Refs.call<Unit>("", 
def, emptyList(), *inputs.toTypedArray()) }
     return contextWithoutTaskDef().also { it.taskDef = def }
   }
 
diff --git 
a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/internal/RefsTest.kt 
b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/internal/RefsTest.kt
index 975dd8e449c..3a06f7c472b 100644
--- a/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/internal/RefsTest.kt
+++ b/java-sdk/sdk/src/test/kotlin/org/apache/airflow/sdk/internal/RefsTest.kt
@@ -54,7 +54,7 @@ internal class RefsTest {
     val dag = DagDef("d")
     Refs.record(dag, listOf("p", "c"), emptyList()) {
       val producer = Refs.node<Long>("", TaskDef("p", NoopRefTask::class.java))
-      Refs.call<Unit>("", TaskDef("c", NoopRefTask::class.java), producer, 
Arg.lit(5))
+      Refs.call<Unit>("", TaskDef("c", NoopRefTask::class.java), emptyList(), 
producer, Arg.lit(5))
     }
 
     val consumerDef = dag.tasks.getValue("c")
@@ -126,7 +126,7 @@ internal class RefsTest {
     val error =
       assertThrows(IllegalArgumentException::class.java) {
         Refs.record(DagDef("d"), listOf("t"), emptyList()) {
-          Refs.call<Unit>("", TaskDef("t", NoopRefTask::class.java), 
Arg.lit(1), null)
+          Refs.call<Unit>("", TaskDef("t", NoopRefTask::class.java), 
emptyList(), Arg.lit(1), null)
         }
       }
 
@@ -140,7 +140,7 @@ internal class RefsTest {
       assertThrows(IllegalArgumentException::class.java) {
         Refs.record(DagDef("d"), listOf("t"), emptyList()) {
           Refs.node<Unit>("", TaskDef("t", NoopRefTask::class.java))
-          Refs.call<Unit>("", TaskDef("t", NoopRefTask::class.java), 
Arg.lit(1))
+          Refs.call<Unit>("", TaskDef("t", NoopRefTask::class.java), 
emptyList(), Arg.lit(1))
         }
       }
 
diff --git a/scripts/ci/lang_sdk_serialization/compare.py 
b/scripts/ci/lang_sdk_serialization/compare.py
index 0d1dc19bf8f..e6fbab441c8 100644
--- a/scripts/ci/lang_sdk_serialization/compare.py
+++ b/scripts/ci/lang_sdk_serialization/compare.py
@@ -35,6 +35,9 @@ features the SDK has, as a comma-separated list or ``all``, 
and the copy of test
 the Dags whose features are all among them. A Dag that requires nothing is 
always in it, and the SDK
 has to write exactly the Dags of the copy.
 
+Each SDK's own prek hook runs it this way, such as
+java-sdk/scripts/ci/prek/check_serialization_conformance.py for the Java SDK.
+
 On a failure the files are kept, and their directory is printed.
 """
 
@@ -147,21 +150,25 @@ def get_task_defaults() -> dict[str, Any]:
     return {key: field["default"] for key, field in fields.items() if 
field.get("default") is not None}
 
 
-def normalize_as_javascript(value: Any) -> Any:
-    """Read a JSON value as JavaScript does: one number type, and a bool that 
is not a number."""
+def normalize_numbers(value: Any) -> Any:
+    """
+    Read a JSON value with one number type, and a bool that is not a number.
+
+    An SDK may write ``2`` where Python writes ``2.0``, as JavaScript does, 
which has one number type.
+    """
     if isinstance(value, bool):
         return ("bool", value)
     if isinstance(value, (int, float)):
         return float(value)
     if isinstance(value, list):
-        return [normalize_as_javascript(item) for item in value]
+        return [normalize_numbers(item) for item in value]
     if isinstance(value, dict):
-        return {key: normalize_as_javascript(item) for key, item in 
value.items()}
+        return {key: normalize_numbers(item) for key, item in value.items()}
     return value
 
 
 def is_same_json(python: Any, sdk: Any) -> bool:
-    return normalize_as_javascript(python) == normalize_as_javascript(sdk)
+    return normalize_numbers(python) == normalize_numbers(sdk)
 
 
 def find_differences(path: str, python: Any, sdk: Any) -> list[str]:
diff --git a/scripts/ci/lang_sdk_serialization/test_dags.yaml 
b/scripts/ci/lang_sdk_serialization/test_dags.yaml
index a9591031bb8..ae5a756426d 100644
--- a/scripts/ci/lang_sdk_serialization/test_dags.yaml
+++ b/scripts/ci/lang_sdk_serialization/test_dags.yaml
@@ -79,7 +79,7 @@ dags:
           wait_for_past_depends_before_skipping: true
           wait_for_downstream: true
           retries: 2
-          queue: typescript
+          queue: workers
           pool: tiny
           pool_slots: 2
           execution_timeout: !timedelta 120
@@ -143,7 +143,7 @@ dags:
 
   # Fan-out and fan-in, so downstream_task_ids has to be sorted rather than 
kept in wiring
   # order. The schedule is not written "@daily" because Python expands a 
preset before
-  # serializing it, which serde.test.ts covers separately.
+  # serializing it, which each SDK's own serializer tests cover separately.
   - dag_id: conformance_diamond
     spec:
       schedule: "0 0 * * *"
diff --git a/scripts/tests/ci/lang_sdk_serialization/test_compare.py 
b/scripts/tests/ci/lang_sdk_serialization/test_compare.py
index a5efc191657..6bba7d0abbc 100644
--- a/scripts/tests/ci/lang_sdk_serialization/test_compare.py
+++ b/scripts/tests/ci/lang_sdk_serialization/test_compare.py
@@ -48,6 +48,7 @@ def build_serialized(fileloc: str, tasks: list[dict]) -> dict:
                 "timezone": "UTC",
                 "catchup": False,
                 "tags": ["a"],
+                "params": [],
                 "task_group": {"prefix_group_id": True, "children": 
{"extract": ["operator", "extract"]}},
                 "tasks": [{"__type": "operator", "__var": task} for task in 
tasks],
             },
@@ -121,6 +122,16 @@ def test_accepts_the_differences_an_sdk_is_allowed():
             ["d: tags[0] is 'b', Python writes 'a'"],
             id="nested-list",
         ),
+        pytest.param(
+            lambda sdk: sdk["d"]["dag"].update(
+                params=[["limit", {"__class": 
"airflow.sdk.definitions.param.Param", "default": 5}]]
+            ),
+            [
+                "d: params is [['limit', {'__class': 
'airflow.sdk.definitions.param.Param', 'default': 5}]], "
+                "Python writes []"
+            ],
+            id="dag-params",
+        ),
         pytest.param(
             lambda sdk: 
sdk["d"]["dag"]["task_group"].update(prefix_group_id=False),
             ["d: task_group.prefix_group_id is False, Python writes True"],

Reply via email to