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