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

jason810496 pushed a commit to branch jason/lang-sdk-e2e/08-native-e2e
in repository https://gitbox.apache.org/repos/asf/airflow.git

commit d72a4e218d053a8ac24227a3ae83c4b10770f6ee
Author: ZHE YOU LIU <[email protected]>
AuthorDate: Mon Sep 28 08:01:45 2026 +0000

    Add a native Java Dag to the compose e2e
    
    A dedicated bundle declares two Dags in Java, one with the interface API
    and one with annotations, and routes every task to a "java-native" queue.
    Its JAR goes into the Dag bundle, where a root-less Java coordinator parses
    it in the Dag processor and runs its tasks. The test checks the edges, the
    queues, the run, the XComs and the Dag source the JAR embeds.
---
 airflow-e2e-tests/docker/java.yml                  |  10 +-
 airflow-e2e-tests/java-native-bundle/.gitignore    |   2 +
 airflow-e2e-tests/java-native-bundle/build.gradle  |  50 +++++++
 .../java-native-bundle/gradle.properties           |  20 +++
 .../java-native-bundle/settings.gradle             |  33 +++++
 .../apache/airflow/e2e/NativeBundleBuilder.java    |  84 +++++++++++
 .../airflow/e2e/nativedag/AnnotationDag.java       |  62 +++++++++
 .../tests/airflow_e2e_tests/conftest.py            |  23 +++-
 .../tests/airflow_e2e_tests/constants.py           |   5 +
 .../java_sdk_tests/test_java_sdk_native_dag.py     | 153 +++++++++++++++++++++
 10 files changed, 434 insertions(+), 8 deletions(-)

diff --git a/airflow-e2e-tests/docker/java.yml 
b/airflow-e2e-tests/docker/java.yml
index db4897f44c8..e91ef21f030 100644
--- a/airflow-e2e-tests/docker/java.yml
+++ b/airflow-e2e-tests/docker/java.yml
@@ -22,14 +22,18 @@
 # the pre-built bundle JARs (the Java example under /opt/airflow/java-jars, the
 # Scala Spark example under /opt/airflow/scala-jars, and the runner-behaviour
 # test fixtures under /opt/airflow/java-test-jars), and configures the worker 
to
-# consume the "java", "scala", and "java-test" Celery queues where @task.stub
-# tasks are routed.
+# consume the "java", "scala", "java-test", and "java-native" Celery queues
+# where Java tasks are routed. The native-Dag JAR is in the Dag bundle, so the
+# Dag processor runs on the same image to parse it.
 ---
 services:
+  airflow-dag-processor:
+    image: airflow-java-worker
+
   airflow-worker:
     image: airflow-java-worker
     volumes:
       - ./java-jars:/opt/airflow/java-jars:ro
       - ./scala-jars:/opt/airflow/scala-jars:ro
       - ./java-test-jars:/opt/airflow/java-test-jars:ro
-    command: celery worker -q java,scala,java-test,default
+    command: celery worker -q java,scala,java-test,java-native,default
diff --git a/airflow-e2e-tests/java-native-bundle/.gitignore 
b/airflow-e2e-tests/java-native-bundle/.gitignore
new file mode 100644
index 00000000000..7f6823bcc0f
--- /dev/null
+++ b/airflow-e2e-tests/java-native-bundle/.gitignore
@@ -0,0 +1,2 @@
+.gradle
+build/
diff --git a/airflow-e2e-tests/java-native-bundle/build.gradle 
b/airflow-e2e-tests/java-native-bundle/build.gradle
new file mode 100644
index 00000000000..0e572c073fc
--- /dev/null
+++ b/airflow-e2e-tests/java-native-bundle/build.gradle
@@ -0,0 +1,50 @@
+/*
+ * 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.
+ */
+
+plugins {
+    id("org.apache.airflow.sdk") version "${projectVersion}"
+}
+
+repositories {
+    mavenLocal()
+    mavenCentral()
+}
+
+dependencies {
+    
annotationProcessor("org.apache.airflow:airflow-sdk-processor:${projectVersion}")
+    implementation("org.apache.airflow:airflow-sdk:${projectVersion}")
+    implementation("org.apache.airflow:airflow-sdk-jpl:${projectVersion}")
+}
+
+java {
+    toolchain {
+        languageVersion.set(JavaLanguageVersion.of(11))
+    }
+    sourceCompatibility = JavaVersion.VERSION_11
+}
+
+sourceSets {
+    main {
+        java.srcDir("src/java")
+    }
+}
+
+airflowBundle {
+    mainClass = "org.apache.airflow.e2e.NativeBundleBuilder"
+}
diff --git a/airflow-e2e-tests/java-native-bundle/gradle.properties 
b/airflow-e2e-tests/java-native-bundle/gradle.properties
new file mode 100644
index 00000000000..c3c94e0057f
--- /dev/null
+++ b/airflow-e2e-tests/java-native-bundle/gradle.properties
@@ -0,0 +1,20 @@
+# 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.
+
+org.gradle.configuration-cache=true
+
+projectVersion=1.0.0-SNAPSHOT
diff --git a/airflow-e2e-tests/java-native-bundle/settings.gradle 
b/airflow-e2e-tests/java-native-bundle/settings.gradle
new file mode 100644
index 00000000000..75016bb20bd
--- /dev/null
+++ b/airflow-e2e-tests/java-native-bundle/settings.gradle
@@ -0,0 +1,33 @@
+/*
+ * 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.
+ */
+
+// Route the plugin lookup to the SDK build published to the local Maven
+// repository by conftest._setup_java_sdk_integration.
+pluginManagement {
+    repositories {
+        mavenLocal()
+        gradlePluginPortal()
+    }
+}
+
+plugins {
+    id("org.gradle.toolchains.foojay-resolver-convention") version "0.10.0"
+}
+
+rootProject.name = "airflow-e2e-java-native-bundle"
diff --git 
a/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/NativeBundleBuilder.java
 
b/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/NativeBundleBuilder.java
new file mode 100644
index 00000000000..17ac91e4cd3
--- /dev/null
+++ 
b/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/NativeBundleBuilder.java
@@ -0,0 +1,84 @@
+/*
+ * 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.e2e;
+
+import java.util.List;
+import org.apache.airflow.e2e.nativedag.AnnotationDag;
+import org.apache.airflow.sdk.*;
+
+/**
+ * Bundle for the native Dag E2E tests: Dags declared entirely in Java, parsed 
by the Dag processor
+ * from this JAR.
+ *
+ * <p>This is the bundle's main class, so it is the Dag source the Airflow UI 
shows. Every task sets
+ * {@code queue}, which routes it to the {@code java-native} coordinator.
+ */
+public class NativeBundleBuilder {
+  private static final String QUEUE = "java-native";
+
+  public static class Extract implements Task {
+    @Override
+    public void execute(Context context, Client client) {
+      client.setXCom(42L);
+    }
+  }
+
+  public static class Transform implements Task {
+    @Override
+    public void execute(Context context, Client client) {
+      var extracted = ((Number) client.getXCom("extract")).longValue();
+      client.setXCom(extracted * 2);
+    }
+  }
+
+  /** Fails the run unless the value reached it, so a successful run proves 
the XCom flow. */
+  public static class Load implements Task {
+    @Override
+    public void execute(Context context, Client client) {
+      var transformed = ((Number) client.getXCom("transform")).longValue();
+      if (transformed != 84L) {
+        throw new IllegalStateException("load expected 84 from transform, got 
" + transformed);
+      }
+    }
+  }
+
+  public static DagDef buildDag() {
+    var dag =
+        new DagDef("java_native_e2e")
+            .config("description", "Native Java Dag of the Airflow E2E tests")
+            .config("catchup", false)
+            .config("tags", List.of("java-sdk", "native", "e2e"));
+
+    var extract = dag.task("extract", Extract.class).config("queue", QUEUE);
+    var transform = dag.task("transform", Transform.class).config("queue", 
QUEUE);
+    var load = dag.task("load", Load.class).config("queue", QUEUE);
+
+    transform.after(extract).before(load);
+    return dag;
+  }
+
+  public static Bundle build() {
+    return new Bundle().register(buildDag()).register(AnnotationDag.class);
+  }
+
+  public static void main(String[] args) {
+    Server.create(args).serve(build());
+  }
+}
diff --git 
a/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/nativedag/AnnotationDag.java
 
b/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/nativedag/AnnotationDag.java
new file mode 100644
index 00000000000..e00dccb27b5
--- /dev/null
+++ 
b/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/nativedag/AnnotationDag.java
@@ -0,0 +1,62 @@
+/*
+ * 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.
+ */
+
+// "native" is a Java keyword, so the native Dags live in "nativedag".
+package org.apache.airflow.e2e.nativedag;
+
+import org.apache.airflow.sdk.*;
+
+/** A native Java Dag declared with annotations; its tasks set {@code queue} 
like the interface Dag's. */
[email protected](
+    id = "java_native_annotation_e2e",
+    description = "Native Java Dag of the Airflow E2E tests, declared with 
annotations",
+    catchup = false,
+    tags = {"java-sdk", "native", "e2e"})
+public class AnnotationDag {
+  @Builder.Task(id = "extract", queue = "java-native")
+  public long extract() {
+    return 42L;
+  }
+
+  @Builder.Task(id = "transform", queue = "java-native")
+  public long transform(long extracted, double factor) {
+    return (long) (extracted * factor);
+  }
+
+  /** Fails the run unless the value reached it, so a successful run proves 
the XCom flow. */
+  @Builder.Task(id = "load", queue = "java-native")
+  public void load(long transformed) {
+    if (transformed != 63L) {
+      throw new IllegalStateException("load expected 63 from transform, got " 
+ transformed);
+    }
+  }
+
+  @Builder.Task(id = "audit", queue = "java-native")
+  public void audit() {}
+
+  @Builder.Deps
+  static class Wiring implements AnnotationDagDeps {
+    void depends() {
+      var extracted = extract();
+      load(transform(extracted, lit(1.5)));
+      // Ordering-only edge: audit runs after extract, with no data flowing.
+      extracted.before(audit());
+    }
+  }
+}
diff --git a/airflow-e2e-tests/tests/airflow_e2e_tests/conftest.py 
b/airflow-e2e-tests/tests/airflow_e2e_tests/conftest.py
index 3fd6f9a6b6e..aa56af778f8 100644
--- a/airflow-e2e-tests/tests/airflow_e2e_tests/conftest.py
+++ b/airflow-e2e-tests/tests/airflow_e2e_tests/conftest.py
@@ -47,6 +47,8 @@ from airflow_e2e_tests.constants import (
     GO_SDK_ROOT_PATH,
     JAVA_COMPOSE_PATH,
     JAVA_DOCKERFILE_PATH,
+    JAVA_NATIVE_BUNDLE_LIBS_PATH,
+    JAVA_NATIVE_BUNDLE_ROOT_PATH,
     JAVA_SDK_EXAMPLE_DAGS_PATH,
     JAVA_SDK_EXAMPLE_LIBS_PATH,
     JAVA_SDK_MAVEN_CACHE_PATH,
@@ -383,8 +385,8 @@ def _setup_java_sdk_integration(dot_env_file, tmp_dir):
     console.print("[yellow]Publishing Java SDK artifacts to local Maven 
repository...")
     _run_java_sdk_gradle(JAVA_SDK_ROOT_PATH, "publishToMavenLocal", 
"-PskipSigning=true", native=native)
 
-    # The example, scala_spark_example, and java-test-bundle are independent
-    # Gradle builds that all consume the SDK artifact published above, so build
+    # The example, scala_spark_example, java-test-bundle, and 
java-native-bundle
+    # are independent Gradle builds that all consume the SDK artifact 
published above, so build
     # them concurrently. Sharing a writable Gradle user home between concurrent
     # builds is safe because each build can ping the other's lock-owner port 
over
     # one shared loopback - the host's own in native mode, --network=host in 
the
@@ -399,14 +401,17 @@ def _setup_java_sdk_integration(dot_env_file, tmp_dir):
     rmtree(JAVA_SDK_EXAMPLE_LIBS_PATH, ignore_errors=True)
     rmtree(SCALA_SPARK_EXAMPLE_LIBS_PATH, ignore_errors=True)
     rmtree(JAVA_TEST_BUNDLE_LIBS_PATH, ignore_errors=True)
+    rmtree(JAVA_NATIVE_BUNDLE_LIBS_PATH, ignore_errors=True)
     toolchain = "host toolchain" if native else "eclipse-temurin:17-jdk"
     console.print(
-        f"[yellow]Building Java SDK, Scala Spark, and test-fixture bundles 
concurrently ({toolchain})..."
+        "[yellow]Building Java SDK, Scala Spark, test-fixture, and native-Dag 
bundles concurrently "
+        f"({toolchain})..."
     )
     example_bundle_workdirs = [
         JAVA_SDK_ROOT_PATH / "example",
         JAVA_SDK_ROOT_PATH / "scala_spark_example",
         JAVA_TEST_BUNDLE_ROOT_PATH,
+        JAVA_NATIVE_BUNDLE_ROOT_PATH,
     ]
     with ThreadPoolExecutor(max_workers=len(example_bundle_workdirs)) as pool:
         bundle_builds = [
@@ -424,6 +429,8 @@ def _setup_java_sdk_integration(dot_env_file, tmp_dir):
     copytree(JAVA_SDK_EXAMPLE_LIBS_PATH, tmp_dir / "java-jars")
     copytree(SCALA_SPARK_EXAMPLE_LIBS_PATH, tmp_dir / "scala-jars")
     copytree(JAVA_TEST_BUNDLE_LIBS_PATH, tmp_dir / "java-test-jars")
+    # The native-Dag bundle goes into the Dag bundle, where the Dag processor 
parses it.
+    copytree(JAVA_NATIVE_BUNDLE_LIBS_PATH, tmp_dir / "dags" / "java-native")
 
     # Copy the Java SDK example Dag files so Airflow can discover them.
     copyfile(JAVA_SDK_EXAMPLE_DAGS_PATH / "java_examples.py", tmp_dir / "dags" 
/ "java_examples.py")
@@ -437,7 +444,7 @@ def _setup_java_sdk_integration(dot_env_file, tmp_dir):
     # JRE and copies nothing from the context, so without this docker build 
would
     # tar and stream the bundles (hundreds of MB of Spark JARs) to the daemon 
for
     # nothing. The JARs reach the worker via the compose bind-mounts, not the 
image.
-    (tmp_dir / 
".dockerignore").write_text("java-jars/\nscala-jars/\njava-test-jars/\n")
+    (tmp_dir / 
".dockerignore").write_text("java-jars/\nscala-jars/\njava-test-jars/\ndags/\n")
 
     # Build a local Docker image that extends DOCKER_IMAGE with a JRE.
     # We do this explicitly so testcontainers' DockerCompose.start() does not
@@ -483,10 +490,16 @@ def _setup_java_sdk_integration(dot_env_file, tmp_dir):
                 "classpath": "airflow.sdk.coordinators.java.JavaCoordinator",
                 "kwargs": {"jars_root": ["/opt/airflow/java-test-jars"]},
             },
+            # No root, so it serves the Dag bundle: it parses the native-Dag 
JAR
+            # there and runs that JAR's tasks.
+            "java-native": {
+                "classpath": "airflow.sdk.coordinators.java.JavaCoordinator",
+                "kwargs": {},
+            },
         }
     )
     queue_to_coordinator = json.dumps(
-        {"java": "java-jdk", "scala": "scala-jdk", "java-test": 
"java-test-jdk"}
+        {"java": "java-jdk", "scala": "scala-jdk", "java-test": 
"java-test-jdk", "java-native": "java-native"}
     )
 
     # Connection expected by the Java example bundle tasks. The JSON form
diff --git a/airflow-e2e-tests/tests/airflow_e2e_tests/constants.py 
b/airflow-e2e-tests/tests/airflow_e2e_tests/constants.py
index 4a1a2f7be38..a9d1eab2c77 100644
--- a/airflow-e2e-tests/tests/airflow_e2e_tests/constants.py
+++ b/airflow-e2e-tests/tests/airflow_e2e_tests/constants.py
@@ -84,6 +84,11 @@ JAVA_TEST_BUNDLE_ROOT_PATH = AIRFLOW_ROOT_PATH / 
"airflow-e2e-tests" / "java-tes
 JAVA_TEST_BUNDLE_DAGS_PATH = JAVA_TEST_BUNDLE_ROOT_PATH / "src" / "resources" 
/ "dags"
 JAVA_TEST_BUNDLE_LIBS_PATH = JAVA_TEST_BUNDLE_ROOT_PATH / "build" / "bundle"
 
+# Java native-Dag bundle paths (Dags declared in Java, parsed from the JAR by 
the
+# Dag processor; its JAR goes into the Dag bundle rather than a coordinator 
root).
+JAVA_NATIVE_BUNDLE_ROOT_PATH = AIRFLOW_ROOT_PATH / "airflow-e2e-tests" / 
"java-native-bundle"
+JAVA_NATIVE_BUNDLE_LIBS_PATH = JAVA_NATIVE_BUNDLE_ROOT_PATH / "build" / 
"bundle"
+
 # Go SDK E2E test paths
 GO_SDK_ROOT_PATH = AIRFLOW_ROOT_PATH / "go-sdk"
 GO_SDK_DAGS_PATH = GO_SDK_ROOT_PATH / "dags"
diff --git 
a/airflow-e2e-tests/tests/airflow_e2e_tests/java_sdk_tests/test_java_sdk_native_dag.py
 
b/airflow-e2e-tests/tests/airflow_e2e_tests/java_sdk_tests/test_java_sdk_native_dag.py
new file mode 100644
index 00000000000..f8c00b31fb3
--- /dev/null
+++ 
b/airflow-e2e-tests/tests/airflow_e2e_tests/java_sdk_tests/test_java_sdk_native_dag.py
@@ -0,0 +1,153 @@
+# 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.
+"""
+End-to-end test of Dags declared entirely in Java.
+
+Run with::
+
+    E2E_TEST_MODE=java_sdk uv run --project airflow-e2e-tests pytest \\
+        tests/airflow_e2e_tests/java_sdk_tests/test_java_sdk_native_dag.py -xvs
+
+No Python file declares these Dags. The 
``airflow-e2e-tests/java-native-bundle`` JAR sits in the Dag
+bundle, where the root-less ``java-native`` coordinator claims it: the Dag 
processor runs the JAR to
+parse it, and the worker runs the same JAR for each task, which every task 
routes to with its
+``queue``. ``java_native_e2e`` is declared with the interface API, 
``java_native_annotation_e2e``
+with annotations.
+"""
+
+from __future__ import annotations
+
+from dataclasses import dataclass
+from datetime import datetime, timezone
+
+import pytest
+
+from airflow_e2e_tests.e2e_test_utils.clients import AirflowClient
+
+# The Dag processor starts a JVM to parse the JAR, and each task starts 
another.
+_JAVA_TASK_TIMEOUT = 600
+
+_QUEUE = "java-native"
+
+
+@dataclass(frozen=True)
+class _NativeDag:
+    dag_id: str
+    downstream: dict[str, set[str]]
+    return_values: dict[str, int]
+
+
+_NATIVE_DAGS = [
+    _NativeDag(
+        dag_id="java_native_e2e",
+        downstream={"extract": {"transform"}, "transform": {"load"}, "load": 
set()},
+        return_values={"extract": 42, "transform": 84},
+    ),
+    _NativeDag(
+        dag_id="java_native_annotation_e2e",
+        downstream={"extract": {"transform", "audit"}, "transform": {"load"}, 
"load": set(), "audit": set()},
+        # transform(extracted, lit(1.5)).
+        return_values={"extract": 42, "transform": 63},
+    ),
+]
+
+_by_dag_id = pytest.mark.parametrize("native_dag", _NATIVE_DAGS, ids=lambda d: 
d.dag_id)
+
+
+@dataclass
+class _CompletedRun:
+    run_id: str
+    state: str
+    ti_states: dict[str, str]
+
+
[email protected](scope="module")
+def parsed_dags() -> AirflowClient:
+    """A client that has waited for the Dag processor to register both Dags 
from the JAR."""
+    client = AirflowClient()
+    for native_dag in _NATIVE_DAGS:
+        client.wait_for_dag(native_dag.dag_id, timeout=_JAVA_TASK_TIMEOUT)
+    return client
+
+
[email protected](scope="module")
+def completed_runs(parsed_dags: AirflowClient) -> dict[str, _CompletedRun]:
+    """Trigger both Dags at once, then wait for both runs."""
+    client = parsed_dags
+    run_ids = {
+        native_dag.dag_id: client.trigger_dag(
+            native_dag.dag_id, json={"logical_date": 
datetime.now(timezone.utc).isoformat()}
+        )["dag_run_id"]
+        for native_dag in _NATIVE_DAGS
+    }
+    runs = {}
+    for dag_id, run_id in run_ids.items():
+        state = client.wait_for_dag_run(dag_id=dag_id, run_id=run_id, 
timeout=_JAVA_TASK_TIMEOUT)
+        ti_resp = client.get_task_instances(dag_id=dag_id, run_id=run_id)
+        runs[dag_id] = _CompletedRun(
+            run_id=run_id,
+            state=state,
+            ti_states={ti["task_id"]: ti.get("state") for ti in 
ti_resp.get("task_instances", [])},
+        )
+    return runs
+
+
+@_by_dag_id
+def test_the_graph_is_the_one_java_declared(parsed_dags: AirflowClient, 
native_dag: _NativeDag):
+    tasks = parsed_dags.get_tasks(native_dag.dag_id).get("tasks", [])
+
+    assert {task["task_id"]: set(task["downstream_task_ids"]) for task in 
tasks} == native_dag.downstream
+
+
+@_by_dag_id
+def test_every_task_is_routed_to_the_native_coordinator(parsed_dags: 
AirflowClient, native_dag: _NativeDag):
+    """The Java DSL has no Dag-level queue, so each task sets its own."""
+    tasks = parsed_dags.get_tasks(native_dag.dag_id).get("tasks", [])
+
+    assert {task["task_id"]: task["queue"] for task in tasks} == 
dict.fromkeys(native_dag.downstream, _QUEUE)
+
+
+@_by_dag_id
+def test_the_dag_source_is_the_bundle_main_class(parsed_dags: AirflowClient, 
native_dag: _NativeDag):
+    """The Code view shows the Java source the JAR embeds, not the JAR read as 
text."""
+    content = parsed_dags.get_dag_source(native_dag.dag_id)["content"]
+
+    assert "public class NativeBundleBuilder" in content
+    assert 'new DagDef("java_native_e2e")' in content
+
+
+@_by_dag_id
+def test_dag_run_succeeded(completed_runs: dict[str, _CompletedRun], 
native_dag: _NativeDag):
+    run = completed_runs[native_dag.dag_id]
+
+    assert run.state == "success", (
+        f"expected the run to succeed; got {run.state!r}. task states: 
{run.ti_states}"
+    )
+    # load throws unless transform's value reached it, so its success proves 
the last hop.
+    assert run.ti_states == dict.fromkeys(native_dag.downstream, "success")
+
+
+@_by_dag_id
+def test_xcoms_flow_between_java_tasks(
+    parsed_dags: AirflowClient, completed_runs: dict[str, _CompletedRun], 
native_dag: _NativeDag
+):
+    run_id = completed_runs[native_dag.dag_id].run_id
+    for task_id, expected in native_dag.return_values.items():
+        value = parsed_dags.get_xcom_value(
+            dag_id=native_dag.dag_id, task_id=task_id, run_id=run_id, 
key="return_value"
+        ).get("value")
+        assert value == expected, f"{native_dag.dag_id}.{task_id} returned 
{value!r}"

Reply via email to