diff --git a/.github/workflows/k8s-tests.yml b/.github/workflows/k8s-tests.yml index ecad4e15e02ec..e4e3a6147ae7f 100644 --- a/.github/workflows/k8s-tests.yml +++ b/.github/workflows/k8s-tests.yml @@ -106,7 +106,7 @@ jobs: # preparing k8s environment with uv takes < 15 seconds with `uv` - there is no point in caching it. # The lang-SDK coordinator system test is intentionally NOT run in this matrix. It has its own # dedicated single-combo job (tests-kubernetes-lang-sdk below), so it no longer adds the ~6 minutes - # of Go/Java build + coordinator run to each of these system-test jobs. + # of Go/Java/TypeScript build + coordinator run to each of these system-test jobs. - name: "\ Run complete K8S tests ${{ matrix.executor }}-${{ env.PYTHON_MAJOR_MINOR_VERSION }}-\ ${{env.KUBERNETES_VERSION}}-${{ matrix.use-standard-naming }}" @@ -140,11 +140,12 @@ jobs: run: breeze k8s delete-cluster --all if: always() - # The Multi-Lang KubernetesExecutor coordinator system test builds a Go bundle and a Java jar and runs - # them as workers under KubernetesExecutor. It is executor/naming/version-independent, so instead of - # bolting it onto all six ``KubernetesExecutor``/``use-standard-naming == false`` system-test jobs above - # (which ran it redundantly and added ~6 minutes each), it runs once here on a single default combo and - # executes only the lang-SDK test (``-k test_lang_sdk_combined_dag_succeeds``) - not the full suite. + # The lang-SDK KubernetesExecutor coordinator system tests build a Go bundle, Java jars and + # TypeScript bundles and run them as workers under KubernetesExecutor. They are + # executor/naming/version-independent, so instead of bolting them onto all six + # ``KubernetesExecutor``/``use-standard-naming == false`` system-test jobs above (which ran them + # redundantly and added several minutes each), they run once here on a single default combo and + # execute only the lang-SDK tests (``-k TestLangSdk``) - not the full suite. tests-kubernetes-lang-sdk: timeout-minutes: 60 name: "K8S Lang-SDK:${{ inputs.lang-sdk-kubernetes-combo }}" @@ -182,7 +183,7 @@ jobs: use-uv: ${{ inputs.use-uv }} make-mnt-writeable-and-cleanup: true id: breeze - # Provision the lang-SDK Go/Java toolchains on the host so breeze builds the artifacts natively + # Provision the lang-SDK Go/Java/Node toolchains on the host so breeze builds the artifacts natively # (LANG_SDK_NATIVE_TOOLCHAIN=true) instead of in throwaway toolchain containers, skipping the image # pulls and reusing the caches restored below. The module/build and Gradle caches are keyed # explicitly (rather than via the setup-* built-in caching) so the key carries an ``-vN-`` salt: @@ -213,15 +214,28 @@ jobs: path: | ~/.gradle/caches ~/.gradle/wrapper - key: lang-sdk-gradle-v1-${{ runner.os }}-${{ runner.arch }}-${{ hashFiles('java-sdk/**/*.gradle*', 'java-sdk/**/gradle-wrapper.properties', 'kubernetes-tests/lang_sdk/java_example/**/*.gradle*') }} # yamllint disable-line rule:line-length + key: lang-sdk-gradle-v1-${{ runner.os }}-${{ runner.arch }}-${{ hashFiles('java-sdk/**/*.gradle*', 'java-sdk/**/gradle-wrapper.properties', 'kubernetes-tests/lang_sdk/java_example/**/*.gradle*', 'airflow-e2e-tests/java-native-bundle/**/*.gradle*') }} # yamllint disable-line rule:line-length restore-keys: | lang-sdk-gradle-v1-${{ runner.os }}-${{ runner.arch }}- + - name: "Setup pnpm for the lang-SDK TypeScript build" + uses: pnpm/action-setup@ea17c68df8912ef543352723c149a84f56e3d413 # v6.1.0 + with: + package_json_file: ts-sdk/package.json + run_install: false + - name: "Setup Node for the lang-SDK TypeScript build" + uses: actions/setup-node@820762786026740c76f36085b0efc47a31fe5020 # v7.0.0 + with: + node-version: 22 + cache: 'pnpm' + cache-dependency-path: | + ts-sdk/pnpm-lock.yaml + kubernetes-tests/lang_sdk/ts_example/pnpm-lock.yaml # No --upgrade here: the redeploy/upgrade path is already exercised across the main matrix above; - # this job only needs a single deploy to validate the Multi-Lang coordinator, saving ~1.5 minutes. - - name: "Run lang-SDK K8S test ${{ inputs.lang-sdk-kubernetes-combo }}" + # this job only needs a single deploy to validate the lang-SDK coordinators, saving ~1.5 minutes. + - name: "Run lang-SDK K8S tests ${{ inputs.lang-sdk-kubernetes-combo }}" run: >- breeze k8s run-complete-tests --no-copy-local-sources - -- -k test_lang_sdk_combined_dag_succeeds + -- -k TestLangSdk env: EXECUTOR: "KubernetesExecutor" USE_STANDARD_NAMING: "false" diff --git a/.rat-excludes b/.rat-excludes index dc4e0ffcd6a16..b41f42959cfb1 100644 --- a/.rat-excludes +++ b/.rat-excludes @@ -215,6 +215,7 @@ www-hash.txt **/.prettierrc **/.airflowignore **/.airflowignore_glob +**/airflowignore_all # Vendor includes diff --git a/airflow-e2e-tests/docker/java.yml b/airflow-e2e-tests/docker/java.yml index db4897f44c818..3204e395132a0 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 Dags folder, 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 0000000000000..7f6823bcc0f50 --- /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 0000000000000..0e572c073fce5 --- /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/kubernetes-tests/lang_sdk/Dockerfile.java b/airflow-e2e-tests/java-native-bundle/gradle.properties similarity index 57% rename from kubernetes-tests/lang_sdk/Dockerfile.java rename to airflow-e2e-tests/java-native-bundle/gradle.properties index 59a91f8699f9f..c3c94e0057f3f 100644 --- a/kubernetes-tests/lang_sdk/Dockerfile.java +++ b/airflow-e2e-tests/java-native-bundle/gradle.properties @@ -14,18 +14,7 @@ # KIND, either express or implied. See the License for the # specific language governing permissions and limitations # under the License. -# -# Java worker image for the "java" queue: the stock prod image plus a headless -# JRE so the JavaCoordinator can exec `java`. The Go queue runs on the plain -# prod image (a Go bundle is a self-contained static binary), so keeping this a -# separate image demonstrates routing each coordinator's queue to its own -# pod_template_file with a queue-specific base image. -ARG BASE_IMAGE -FROM ${BASE_IMAGE} -USER root -RUN apt-get update \ - && apt-get install --no-install-recommends -y default-jre-headless \ - && apt-get clean \ - && rm -rf /var/lib/apt/lists/* -USER airflow +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 0000000000000..75016bb20bdcf --- /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 0000000000000..b204371c37286 --- /dev/null +++ b/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/NativeBundleBuilder.java @@ -0,0 +1,150 @@ +/* + * 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.e2e.nativedag.TargetDag; +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. + * + *
This is the bundle's main class, so it is the Dag source the Airflow UI shows. Every task sets + * {@code queue}, which {@code queue_to_coordinator} routes to the {@code java-jdk} coordinator. + * The {@code java-native} coordinator only parses this JAR. + */ +public class NativeBundleBuilder { + public 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); + } + } + } + + /** Its boolean picks one of {@link ReportMany} and {@link ReportFew}; the other is skipped. */ + public static class HasRows implements ConditionTask { + @Override + public boolean decide(Context context, Client client) { + return ((Number) client.getXCom("transform")).longValue() > 0; + } + } + + public static class ReportMany implements Task { + @Override + public void execute(Context context, Client client) {} + } + + public static class ReportFew implements Task { + @Override + public void execute(Context context, Client client) {} + } + + /** Names the one of {@link ReportLong} and {@link ReportShort} that runs; the other is skipped. */ + public static class PickReport implements SwitchTask { + @Override + public Class extends Task> choose(Context context, Client client) { + return ((Number) client.getXCom("transform")).longValue() > 100 ? ReportLong.class : ReportShort.class; + } + } + + public static class ReportLong implements Task { + @Override + public void execute(Context context, Client client) {} + } + + public static class ReportShort implements Task { + @Override + public void execute(Context context, Client client) {} + } + + public static class Audit implements Task { + @Override + public void execute(Context context, Client client) {} + } + + 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); + + var reportMany = dag.task("report_many", ReportMany.class).config("queue", QUEUE); + var reportFew = dag.task("report_few", ReportFew.class).config("queue", QUEUE); + dag.If("has_rows", HasRows.class).after(transform).config("queue", QUEUE).Then(reportMany).Else(reportFew); + + var reportLong = dag.task("report_long", ReportLong.class).config("queue", QUEUE); + var reportShort = dag.task("report_short", ReportShort.class).config("queue", QUEUE); + dag.Switch("pick_report", PickReport.class) + .after(transform) + .config("queue", QUEUE) + .Case(reportLong) + .Case(reportShort); + + // A task that starts a run of another Dag; it runs no Java code. + var trigger = + dag.task("trigger_downstream", new TriggerDagRun("java_native_target_e2e")).config("queue", QUEUE); + load.before(trigger); + + // Ordering-only edge: the checks group runs after extract, with no data flowing. + var checks = dag.taskGroup("checks"); + checks.task("audit", Audit.class).config("queue", QUEUE); + extract.before(checks); + + return dag; + } + + public static Bundle build() { + return new Bundle().register(buildDag()).register(AnnotationDag.class).register(TargetDag.build()); + } + + 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 0000000000000..e4d8fc71c5528 --- /dev/null +++ b/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/nativedag/AnnotationDag.java @@ -0,0 +1,106 @@ +/* + * 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 static org.apache.airflow.e2e.NativeBundleBuilder.QUEUE; + +import org.apache.airflow.sdk.*; + +/** A native Java Dag declared with annotations; its tasks use the interface Dag's queue. */ +@Builder.Dag( + 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 = QUEUE) + public long extract() { + return 42L; + } + + @Builder.Task(id = "transform", queue = QUEUE) + 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 = QUEUE) + public void load(long transformed) { + if (transformed != 63L) { + throw new IllegalStateException("load expected 63 from transform, got " + transformed); + } + } + + @Builder.Task(id = "report_many", queue = QUEUE) + public void reportMany() {} + + @Builder.Task(id = "report_few", queue = QUEUE) + public void reportFew() {} + + /** Picks one of {@link #reportMany} and {@link #reportFew}; the other is skipped. */ + @Builder.If(id = "has_rows", queue = QUEUE) + public boolean hasRows(long transformed) { + return transformed > 0; + } + + @Builder.Task(id = "report_long", queue = QUEUE) + public void reportLong() {} + + @Builder.Task(id = "report_short", queue = QUEUE) + public void reportShort() {} + + /** Names the one of {@link #reportLong} and {@link #reportShort} that runs; the other is skipped. */ + @Builder.Switch(id = "pick_report", queue = QUEUE) + public Class extends Task> pickReport(long transformed) { + return transformed > 100 ? AnnotationDagBuilder.ReportLong.class : AnnotationDagBuilder.ReportShort.class; + } + + // A task that starts a run of another Dag. The method runs when this Dag is + // built, not when the task runs, so it takes no arguments. + @Builder.Task(id = "trigger_downstream", queue = QUEUE) + public TriggerDagRun triggerDownstream() { + return new TriggerDagRun("java_native_target_e2e"); + } + + // A task group: everything it declares is prefixed with its id, so this is + // the task "checks.audit". + @Builder.TaskGroup(id = "checks") + static class Checks { + @Builder.Task(id = "audit", queue = QUEUE) + public void audit() {} + } + + @Builder.Deps + static class Wiring implements AnnotationDagDeps { + void depends() { + var extracted = extract(); + var transformed = transform(extracted, lit(1.5)); + var loaded = load(transformed); + hasRows(transformed).Then(reportMany()).Else(reportFew()); + pickReport(transformed).Case(reportLong()).Case(reportShort()); + loaded.before(triggerDownstream()); + // Ordering-only edge: the checks group runs after extract, with no data + // flowing. + extracted.before(checks()); + checks().audit(); + } + } +} diff --git a/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/nativedag/TargetDag.java b/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/nativedag/TargetDag.java new file mode 100644 index 0000000000000..2e1be90d106aa --- /dev/null +++ b/airflow-e2e-tests/java-native-bundle/src/java/org/apache/airflow/e2e/nativedag/TargetDag.java @@ -0,0 +1,46 @@ +/* + * 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 static org.apache.airflow.e2e.NativeBundleBuilder.QUEUE; + +import java.util.List; +import org.apache.airflow.sdk.*; + +/** The Dag the e2e's other native Dags trigger, so they can't trigger each other in a loop. */ +public class TargetDag { + public static class Receive implements Task { + @Override + public void execute(Context context, Client client) { + client.setXCom("triggered"); + } + } + + public static DagDef build() { + var dag = + new DagDef("java_native_target_e2e") + .config("description", "Native Java Dag the e2e's other native Dags trigger") + .config("catchup", false) + .config("tags", List.of("java-sdk", "native", "e2e")); + dag.task("receive", Receive.class).config("queue", QUEUE); + return dag; + } +} diff --git a/airflow-e2e-tests/tests/airflow_e2e_tests/conftest.py b/airflow-e2e-tests/tests/airflow_e2e_tests/conftest.py index 6c48ac89bf497..beb1376e7ea3d 100644 --- a/airflow-e2e-tests/tests/airflow_e2e_tests/conftest.py +++ b/airflow-e2e-tests/tests/airflow_e2e_tests/conftest.py @@ -50,6 +50,8 @@ GO_SDK_STATE_STORE_RETENTION_DAYS, 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, @@ -338,6 +340,9 @@ def _run_java_sdk_gradle(workdir, *gradle_argv, capture_output=False, native=Fal image's HOME (/root) which the non-root process cannot write to. * files/m2 is mounted directly as ~/.m2 so publishToMavenLocal writes there without nesting, and its contents are visible on the host. + maven.repo.local names it explicitly: the JVM takes its home from + /etc/passwd, not HOME, so a host UID that the image does define (such + as 1000, its "ubuntu" user) would otherwise publish into the container. """ if native: cwd = workdir @@ -370,6 +375,7 @@ def _run_java_sdk_gradle(workdir, *gradle_argv, capture_output=False, native=Fal "eclipse-temurin:17-jdk", "/repo/java-sdk/gradlew", "--no-daemon", + "-Dmaven.repo.local=/workspace-home/.m2/repository", *gradle_argv, ] return subprocess.run(argv, cwd=cwd, check=True, capture_output=capture_output, text=True) @@ -398,8 +404,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 @@ -414,14 +420,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 = [ @@ -439,6 +448,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") @@ -452,7 +463,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 @@ -474,11 +485,11 @@ def _setup_java_sdk_integration(dot_env_file, tmp_dir): check=True, ) - # One JavaCoordinator per queue on the same worker image, each serving its - # own artifact bundle (one bundle is one classpath). The scala-jdk entry pins - # main_class (Spark's large classpath makes Main-Class discovery ambiguous) - # and carries Spark's Java 17 module openings, a small driver heap, and a - # longer startup timeout for its large dependency classpath. + # Four JavaCoordinators on the same worker image. The first three each serve their + # own artifact bundle for the stub tasks of their queue (one bundle is one classpath). + # The scala-jdk entry pins main_class (Spark's large classpath makes Main-Class + # discovery ambiguous) and carries Spark's Java 17 module openings, a small driver + # heap, and a longer startup timeout for its large dependency classpath. dag_bundle_config = _build_dag_bundle_config( { "java-task-handlers": "/opt/airflow/java-jars", @@ -505,11 +516,19 @@ def _setup_java_sdk_integration(dot_env_file, tmp_dir): "classpath": "airflow.sdk.coordinators.java.JavaCoordinator", "kwargs": {"task_handler_bundle_name": "java-test-task-handlers"}, }, + # Only parses: dag_bundle_to_coordinator picks it for the Dags folder, since + # there are four JavaCoordinators. No queue routes to it, so the native + # Dag's tasks run on java-jdk, from the Dag's own bundle. + "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-jdk"} ) + dag_bundle_to_coordinator = json.dumps({"dags-folder": "java-native"}) # Connection expected by the Java example bundle tasks. The JSON form # covers all connection fields, in particular the port: wire integers @@ -531,6 +550,7 @@ def _setup_java_sdk_integration(dot_env_file, tmp_dir): f"AIRFLOW__DAG_PROCESSOR__DAG_BUNDLE_CONFIG_LIST='{dag_bundle_config}'\n" f"AIRFLOW__SDK__COORDINATORS='{coordinator_config}'\n" f"AIRFLOW__SDK__QUEUE_TO_COORDINATOR='{queue_to_coordinator}'\n" + f"AIRFLOW__SDK__DAG_BUNDLE_TO_COORDINATOR='{dag_bundle_to_coordinator}'\n" f"AIRFLOW_CONN_TEST_HTTP='{test_http_conn}'\n" # Variable expected by the Java example bundle tasks. "AIRFLOW_VAR_MY_VARIABLE=test_value\n" @@ -779,6 +799,11 @@ def _setup_ts_sdk_integration(dot_env_file, tmp_dir): copyfile(TS_SDK_EXAMPLE_PATH / "dist" / "bundle.min.mjs", ts_bundles_dir / "example.min.mjs") # A native Dag: the Dag processor parses it, so `ts_hitl_ai_approval` has no Python Dag file. copyfile(TS_SDK_EXAMPLE_PATH / "dist" / "ai-approval.min.mjs", ts_bundles_dir / "ai-approval.min.mjs") + # The same artifact inside the Dag bundle, where the Dag processor finds its native Dag. + (tmp_dir / "dags" / "typescript").mkdir() + copyfile( + TS_SDK_EXAMPLE_PATH / "dist" / "bundle.min.mjs", tmp_dir / "dags" / "typescript" / "example.min.mjs" + ) # Both of the example bundle's Dags: one bundle.mjs provides for two dag_ids, # and the tests check that dispatch tells their same-named tasks apart. @@ -786,6 +811,9 @@ def _setup_ts_sdk_integration(dot_env_file, tmp_dir): copyfile(TS_SDK_EXAMPLE_PATH / "dags" / dag_file, tmp_dir / "dags" / dag_file) dag_bundle_config = _build_dag_bundle_config({"ts-task-handlers": "/opt/airflow/ts-bundles"}) + # "ts" is the only NodeCoordinator, so it parses the native Dag in every Dag bundle + # that holds one, the Dags folder here. It runs the stub tasks from the ts-task-handlers + # Dag bundle, and the native Dag's tasks from the Dag's own bundle. coordinator_config = json.dumps( { "ts": { @@ -794,7 +822,7 @@ def _setup_ts_sdk_integration(dot_env_file, tmp_dir): "task_handler_bundle_name": "ts-task-handlers", "node_executable": "/opt/nodejs/node", }, - } + }, } ) queue_to_coordinator = json.dumps({"typescript": "ts"}) diff --git a/airflow-e2e-tests/tests/airflow_e2e_tests/constants.py b/airflow-e2e-tests/tests/airflow_e2e_tests/constants.py index 87e7fa5bdf6f0..024a90fb01bec 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_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/e2e_test_utils/clients.py b/airflow-e2e-tests/tests/airflow_e2e_tests/e2e_test_utils/clients.py index 7f85f07332210..1f1d1ad435a71 100644 --- a/airflow-e2e-tests/tests/airflow_e2e_tests/e2e_test_utils/clients.py +++ b/airflow-e2e-tests/tests/airflow_e2e_tests/e2e_test_utils/clients.py @@ -184,6 +184,12 @@ def get_tasks(self, dag_id: str): """List a Dag's tasks, with the edges each one carries.""" return self._make_request(method="GET", endpoint=f"dags/{dag_id}/tasks") + def get_task_instance_links(self, dag_id: str, run_id: str, task_id: str): + """Get the extra links of a task instance, keyed by link name.""" + return self._make_request( + method="GET", endpoint=f"dags/{dag_id}/dagRuns/{run_id}/taskInstances/{task_id}/links" + ) + def get_dag_source(self, dag_id: str): """Get the source code stored for a Dag's latest version.""" return self._make_request(method="GET", endpoint=f"dagSources/{dag_id}") @@ -253,6 +259,13 @@ def get_task_logs( endpoint=endpoint, ) + def get_event_logs(self, dag_id: str, run_id: str, task_id: str | None = None) -> list[dict]: + """List the audit log events of a Dag run, or of one of its tasks, oldest first.""" + params = {"dag_id": dag_id, "run_id": run_id, "order_by": "event_log_id"} + if task_id is not None: + params["task_id"] = task_id + return self._make_request(method="GET", endpoint="eventLogs", params=params)["event_logs"] + class TaskSDKClient: """Client for interacting with the Task SDK API.""" 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 0000000000000..41219f4017018 --- /dev/null +++ b/airflow-e2e-tests/tests/airflow_e2e_tests/java_sdk_tests/test_java_sdk_native_dag.py @@ -0,0 +1,215 @@ +# 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. There are four ``JavaCoordinator``s, so ``[sdk] dag_bundle_to_coordinator`` picks the +``java-native`` one to parse it: the Dag processor runs the JAR through it. Each task sets +``queue="java-native"``, which routes to ``java-jdk``. That coordinator runs the same JAR from the Dag's +own bundle, whatever its ``task_handler_bundle_name`` says. ``java_native_e2e`` is declared with the interface +API, ``java_native_annotation_e2e`` with annotations. Both also exercise a task group, an ``If``, a +``Switch`` and a ``TriggerDagRun`` of ``java_native_target_e2e``, the third Dag the JAR declares. +""" + +from __future__ import annotations + +from dataclasses import dataclass +from datetime import UTC, datetime + +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 + +# The queue the Dag's tasks set, which ``queue_to_coordinator`` routes to ``java-jdk``. +_QUEUE = "java-native" + +# The Dag ``trigger_downstream`` starts a run of; not triggered through the client itself. +_TARGET_DAG_ID = "java_native_target_e2e" + +# Both native Dags wire an identical If/Switch/task-group/trigger shape onto their own +# extract -> transform -> load chain, so their graphs and skip sets match exactly; only +# the transform value (84 vs 63) differs, and both land on the same side of each decider. +_CONTROL_FLOW_DOWNSTREAM: dict[str, set[str]] = { + "transform": {"load", "has_rows", "pick_report"}, + "load": {"trigger_downstream"}, + "has_rows": {"report_many", "report_few"}, + "report_many": set(), + "report_few": set(), + "pick_report": {"report_long", "report_short"}, + "report_long": set(), + "report_short": set(), + "trigger_downstream": set(), + "checks.audit": set(), +} +# The Else side of the If, and the Case the Switch does not choose. +_CONTROL_FLOW_SKIPPED = frozenset({"report_few", "report_long"}) + + +@dataclass(frozen=True) +class _NativeDag: + dag_id: str + downstream: dict[str, set[str]] + return_values: dict[str, int | bool | str] + # Marks the Java source of the file that declares this Dag; each Dag's own source file, + # not always the bundle's main class, since the bundle embeds one source per Dag. + source_marker: str + skipped: frozenset[str] = frozenset() + + +_NATIVE_DAGS = [ + _NativeDag( + dag_id="java_native_e2e", + downstream={"extract": {"transform", "checks.audit"}, **_CONTROL_FLOW_DOWNSTREAM}, + return_values={"extract": 42, "transform": 84, "has_rows": True, "pick_report": "report_short"}, + source_marker="public class NativeBundleBuilder", + skipped=_CONTROL_FLOW_SKIPPED, + ), + _NativeDag( + dag_id="java_native_annotation_e2e", + downstream={"extract": {"transform", "checks.audit"}, **_CONTROL_FLOW_DOWNSTREAM}, + # transform(extracted, lit(1.5)). + return_values={"extract": 42, "transform": 63, "has_rows": True, "pick_report": "report_short"}, + source_marker="public class AnnotationDag", + skipped=_CONTROL_FLOW_SKIPPED, + ), +] + +_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] + + +@pytest.fixture(scope="module") +def parsed_dags() -> AirflowClient: + """A client that has waited for the Dag processor to register every Dag from the JAR.""" + client = AirflowClient() + for dag_id in (*[native_dag.dag_id for native_dag in _NATIVE_DAGS], _TARGET_DAG_ID): + client.wait_for_dag(dag_id, timeout=_JAVA_TASK_TIMEOUT) + # trigger_downstream never goes through the client, so nothing else un-pauses its + # target; a run of a still-paused Dag can leave its tasks queued instead of running. + client.un_pause_dag(_TARGET_DAG_ID) + return client + + +@pytest.fixture(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(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_sets_the_native_queue(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. + + The bundle embeds one source file per Dag, so the Code view shows the file that actually + declares each Dag, not always the bundle's main class. + """ + content = parsed_dags.get_dag_source(native_dag.dag_id)["content"] + + assert native_dag.source_marker 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 that hop. + expected = { + task_id: "skipped" if task_id in native_dag.skipped else "success" + for task_id in native_dag.downstream + } + assert run.ti_states == expected + + +@_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}" + + +@_by_dag_id +def test_trigger_downstream_run_succeeded( + parsed_dags: AirflowClient, completed_runs: dict[str, _CompletedRun], native_dag: _NativeDag +): + """``trigger_downstream`` pushes the triggered run's ID; follow it and check it too succeeded.""" + run_id = completed_runs[native_dag.dag_id].run_id + triggered_run_id = parsed_dags.get_xcom_value( + dag_id=native_dag.dag_id, task_id="trigger_downstream", run_id=run_id, key="trigger_run_id" + ).get("value") + assert triggered_run_id, f"{native_dag.dag_id}.trigger_downstream pushed no run ID" + + state = parsed_dags.wait_for_dag_run( + dag_id=_TARGET_DAG_ID, run_id=triggered_run_id, timeout=_JAVA_TASK_TIMEOUT + ) + assert state == "success", ( + f"expected the {_TARGET_DAG_ID} run {triggered_run_id} to succeed; got {state!r}" + ) diff --git a/airflow-e2e-tests/tests/airflow_e2e_tests/ts_sdk_tests/test_ts_sdk_native_dag.py b/airflow-e2e-tests/tests/airflow_e2e_tests/ts_sdk_tests/test_ts_sdk_native_dag.py new file mode 100644 index 0000000000000..b7f10ca89e41e --- /dev/null +++ b/airflow-e2e-tests/tests/airflow_e2e_tests/ts_sdk_tests/test_ts_sdk_native_dag.py @@ -0,0 +1,253 @@ +# 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 a Dag declared entirely in TypeScript. + +Run with:: + + E2E_TEST_MODE=ts_sdk uv run --project airflow-e2e-tests pytest \\ + tests/airflow_e2e_tests/ts_sdk_tests/test_ts_sdk_native_dag.py -xvs + +Unlike ``test_ts_sdk_dag.py``, no Python file declares this Dag: the Dag processor asks the +``airflow-ts-pack`` bundle to parse itself, and the bundle answers with the serialized Dag that +``ts-sdk/example/src/native.ts`` built. The graph is a graph rather than a chain -- a task group, a +named fan-in, order-only edges, a conditional, a multi-way branch, and a task that triggers another +Dag's run. + +The bundle sits in the Dag bundle. The ``ts`` coordinator, the only ``NodeCoordinator``, parses it +there. Every task, the trigger included, runs in the TypeScript runtime. +""" + +from __future__ import annotations + +from dataclasses import dataclass +from datetime import UTC, datetime + +import pytest + +from airflow_e2e_tests.e2e_test_utils.clients import AirflowClient + +# Parsing the bundle launches node before the first task is even scheduled, so +# allow the same headroom the mixed-language suite does. +_TS_TASK_TIMEOUT = 600 + +_DAG_ID = "typescript_native_example" +# The Dag trigger_downstream starts; see the triggerDagRun call in ts-sdk/example/src/native.ts. +_DOWNSTREAM_DAG_ID = "typescript_example" + +# Read by the handlers; see ts-sdk/example/src/native.ts. +_NORTH_ROWS_VARIABLE = "typescript_native_north_rows" +_SOUTH_ROWS_VARIABLE = "typescript_native_south_rows" +_CADENCE_VARIABLE = "typescript_native_cadence" + + +@dataclass +class _CompletedRun: + client: AirflowClient + run_id: str + state: str + ti_states: dict[str, str] + + def xcom(self, task_id: str, key: str = "return_value"): + return self.client.get_xcom_value(dag_id=_DAG_ID, task_id=task_id, run_id=self.run_id, key=key).get( + "value" + ) + + +@pytest.fixture(scope="module") +def parsed_dag() -> AirflowClient: + """A client that has waited for the bundle's Dag to be parsed and registered.""" + client = AirflowClient() + # The Dag processor spawns node to parse the bundle, so the Dag appears some + # time after the deployment is up; every read below would 404 until then. + client.wait_for_dag(_DAG_ID, timeout=_TS_TASK_TIMEOUT) + return client + + +@pytest.fixture(scope="module") +def completed_run(parsed_dag: AirflowClient) -> _CompletedRun: + """Trigger the native Dag once, with the inputs every test below reads back.""" + client = parsed_dag + # Both regions non-empty, so the conditional takes its `then` branch, and a + # weekly cadence so the branch's choice is known rather than guessed. + for key, value in ( + (_NORTH_ROWS_VARIABLE, "3"), + (_SOUTH_ROWS_VARIABLE, "2"), + (_CADENCE_VARIABLE, "weekly"), + ): + client.set_variable(key, value) + + # Dags are paused at creation here, so the run trigger_downstream starts would stay queued. + client.un_pause_dag(_DOWNSTREAM_DAG_ID) + client.un_pause_dag(_DAG_ID) + resp = client.trigger_dag(_DAG_ID, json={"logical_date": datetime.now(UTC).isoformat()}) + run_id = resp["dag_run_id"] + state = client.wait_for_dag_run(dag_id=_DAG_ID, run_id=run_id, timeout=_TS_TASK_TIMEOUT) + ti_resp = client.get_task_instances(dag_id=_DAG_ID, run_id=run_id) + return _CompletedRun( + client=client, + run_id=run_id, + state=state, + ti_states={ti["task_id"]: ti.get("state") for ti in ti_resp.get("task_instances", [])}, + ) + + +def test_the_dag_the_bundle_parsed_is_registered(parsed_dag: AirflowClient): + """The Dag exists without any Python file declaring it.""" + dag = parsed_dag.get_dag(_DAG_ID) + + assert dag["dag_id"] == _DAG_ID + # The serializer expands a cron preset, so "@daily" is recorded as its expression. + assert dag.get("timetable_summary") == "0 0 * * *" + # `tags` is a list of objects, each naming one tag. + assert {tag["name"] for tag in dag.get("tags") or []} >= {"typescript", "native"} + + +def test_the_dag_source_is_the_file_that_declares_it(parsed_dag: AirflowClient): + """The Code view shows native.ts, the file that declares this Dag, not the bundle itself.""" + content = parsed_dag.get_dag_source(_DAG_ID)["content"] + + assert 'new Dag("typescript_native_example"' in content + assert "airflow-ts-pack" not in content + + +def test_the_graph_carries_every_construct(parsed_dag: AirflowClient): + """Group prefixes, the fan-in, the branches and the trigger all survive parsing.""" + tasks = parsed_dag.get_tasks(_DAG_ID).get("tasks", []) + downstream = {task["task_id"]: set(task["downstream_task_ids"]) for task in tasks} + + # The task group prefixed its members. + assert {"extract.north", "extract.south"} <= set(downstream) + # The named fan-in put both extract tasks upstream of summarize. The API + # reports downstream edges only, so upstream is read by inverting them. + assert {"extract.north", "extract.south"} <= { + task_id for task_id, down in downstream.items() if "summarize" in down + } + # The conditional and the branch reach their candidates. + assert downstream["has_rows"] >= {"load_rows", "report_empty"} + assert downstream["pick_cadence"] >= {"publish_daily", "publish_weekly"} + # The group edge, expanded onto the tasks the group leaves from. + assert {"extract.north", "extract.south"} <= { + task_id for task_id, down in downstream.items() if "pick_cadence" in down + } + # Cleanup sits behind every branch outcome, which is what its trigger rule is for. + assert {"load_rows", "report_empty", "publish_daily", "publish_weekly"} <= { + task_id for task_id, down in downstream.items() if "cleanup" in down + } + # The trigger reads as the operator it mirrors, though the Node coordinator runs it. + by_id = {task["task_id"]: task for task in tasks} + assert by_id["trigger_downstream"]["operator_name"] == "TriggerDagRunOperator" + + +def test_dag_run_succeeded(completed_run: _CompletedRun): + assert completed_run.state == "success", ( + f"expected the run to succeed; got {completed_run.state!r}. task states: {completed_run.ti_states}" + ) + + +def test_the_taken_branches_ran_and_the_others_skipped(completed_run: _CompletedRun): + """A branch is a run-time skip, so the states are what prove it worked.""" + states = completed_run.ti_states + + # Both regions had rows, so the conditional followed `then`. + assert states.get("load_rows") == "success" + assert states.get("report_empty") == "skipped" + + # The cadence Variable said weekly, so that is the case the branch chose. + assert states.get("publish_weekly") == "success" + assert states.get("publish_daily") == "skipped" + + +def test_every_other_task_succeeded(completed_run: _CompletedRun): + always_run = [ + "extract.north", + "extract.south", + "summarize", + "has_rows", + "pick_cadence", + "cleanup", + "trigger_downstream", + ] + for task_id in always_run: + assert completed_run.ti_states.get(task_id) == "success", ( + f"{task_id!r} did not succeed. all task states: {completed_run.ti_states}" + ) + + +def test_the_trigger_started_the_downstream_run(completed_run: _CompletedRun): + """The TypeScript runtime ran ``trigger_downstream``, and it triggered the Dag.""" + run_id = completed_run.xcom("trigger_downstream", key="trigger_run_id") + client = completed_run.client + + state = client.wait_for_dag_run(dag_id=_DOWNSTREAM_DAG_ID, run_id=run_id, timeout=_TS_TASK_TIMEOUT) + run = client.get_dag_run(_DOWNSTREAM_DAG_ID, run_id) + + assert state == "success", f"expected the downstream run to succeed; got {run!r}" + assert run["run_type"] == "operator_triggered" + # Sent as the TypeScript Dag wrote it: there is no Jinja rendering. + assert run["conf"] == {"triggered_by": _DAG_ID} + + links = client.get_task_instance_links(_DAG_ID, completed_run.run_id, "trigger_downstream") + assert links["extra_links"]["Triggered DAG"].endswith(f"/dags/{_DOWNSTREAM_DAG_ID}/runs/{run_id}") + + +def test_xcoms_flow_between_typescript_tasks(completed_run: _CompletedRun): + """The fan-in's arguments reached the handler, and its own push came back.""" + assert completed_run.xcom("extract.north") == {"region": "north", "rows": 3} + assert completed_run.xcom("extract.south") == {"region": "south", "rows": 2} + assert completed_run.xcom("summarize") == {"total": 5, "regions": 2} + # Pushed under its own key by the summarize handler, and read back by load_rows. + assert completed_run.xcom("summarize", key="region_total") == 5 + assert completed_run.xcom("load_rows") == {"loaded": 5} + + +def test_the_decisions_are_recorded(completed_run: _CompletedRun): + """A condition returns its boolean, and a branch the task id it chose.""" + assert completed_run.xcom("has_rows") is True + assert completed_run.xcom("pick_cadence") == "publish_weekly" + + +def test_the_skipped_branches_are_recorded_in_xcom(completed_run: _CompletedRun): + """The skipmixin XCom a later clear reads to keep these branches skipped.""" + assert completed_run.xcom("has_rows", key="skipmixin_key") == {"skipped": ["report_empty"]} + assert completed_run.xcom("pick_cadence", key="skipmixin_key") == {"skipped": ["publish_daily"]} + + +def test_the_trigger_deferred_until_the_downstream_run_finished(completed_run: _CompletedRun): + """It deferred to ``DagStateTrigger``, which the Python triggerer runs, and resumed in TypeScript.""" + client = completed_run.client + events = client.get_event_logs(dag_id=_DAG_ID, run_id=completed_run.run_id, task_id="trigger_downstream") + lifecycle = [event for event in events if event["event"] in {"running", "deferred", "success"}] + + # Two starts of one try: the second is the resume, not a retry. + assert [(event["event"], event["try_number"]) for event in lifecycle] == [ + ("running", 1), + ("deferred", 1), + ("running", 1), + ("success", 1), + ], f"trigger_downstream events: {events}" + + # The queue the TypeScript coordinator serves, which the resumed task instance kept. + tis = client.get_task_instances(dag_id=_DAG_ID, run_id=completed_run.run_id)["task_instances"] + assert next(ti for ti in tis if ti["task_id"] == "trigger_downstream")["queue"] == "typescript" + + # Event ids grow in the order the events are written, so the resume came after every event of + # the run it waited on. + downstream_run_id = completed_run.xcom("trigger_downstream", key="trigger_run_id") + downstream = client.get_event_logs(dag_id=_DOWNSTREAM_DAG_ID, run_id=downstream_run_id) + assert downstream, f"no events for {_DOWNSTREAM_DAG_ID} run {downstream_run_id}" + assert max(event["event_log_id"] for event in downstream) < lifecycle[2]["event_log_id"] diff --git a/dev/breeze/doc/05_test_commands.rst b/dev/breeze/doc/05_test_commands.rst index 111344c11ea80..980caf471f05b 100644 --- a/dev/breeze/doc/05_test_commands.rst +++ b/dev/breeze/doc/05_test_commands.rst @@ -591,11 +591,11 @@ Setting up the lang-SDK coordinator system test ............................................... ``breeze k8s setup-lang-sdk-test`` provisions a cluster for the lang-SDK coordinator -system test: it builds the Go and Java example bundles, deploys an in-cluster localstack -S3, uploads the artifacts and the Python stub Dag to their buckets, renders the -coordinator pod-template image placeholders, and installs the Helm release configured for -the ``golang`` and ``java`` queues. After it completes, run the test with -``breeze k8s tests``. +system tests: it builds the Go, Java and TypeScript example bundles (mixed-language and +native), deploys an in-cluster localstack S3, uploads the artifacts and the Dag files to +their buckets, builds and loads the shared runtime image (prod + JRE + Node), and installs +the Helm release configured for the ``golang``, ``java``, ``java-native`` and ``typescript`` +queues. After it completes, run the tests with ``breeze k8s tests``. All parameters of the command are here: diff --git a/dev/breeze/doc/ci/04_selective_checks.md b/dev/breeze/doc/ci/04_selective_checks.md index f75ab8d37cfa4..8940f5a289fac 100644 --- a/dev/breeze/doc/ci/04_selective_checks.md +++ b/dev/breeze/doc/ci/04_selective_checks.md @@ -470,11 +470,18 @@ together using `pytest-xdist` (pytest-xdist distributes the tests among parallel canary and the changed tests must pass after it * `Java SDK E2E tests` (the `java_sdk` mode of the deployed-stack tests, exposed as the `run-java-sdk-e2e-tests` output) run when the Java SDK sources (`java-sdk/`, excluding `.md`), the - Java test-fixture bundle (`airflow-e2e-tests/java-test-bundle/`), the Java e2e suite or its Docker + Java test-fixture and native-Dag bundles (`airflow-e2e-tests/java-test-bundle/`, + `airflow-e2e-tests/java-native-bundle/`), the Java e2e suite or its Docker files (`airflow-e2e-tests/tests/airflow_e2e_tests/java_sdk_tests/`, - `airflow-e2e-tests/docker/java.yml`, `airflow-e2e-tests/docker/Dockerfile.java`), or the Java - coordinator (`task-sdk/src/airflow/sdk/coordinators/java/`, `_subprocess.py`) change. Like the - other deployed e2e suites, enabling them forces `PROD Image building`. + `airflow-e2e-tests/docker/java.yml`, `airflow-e2e-tests/docker/Dockerfile.java`), the Java + coordinator (`task-sdk/src/airflow/sdk/coordinators/java/`, `_subprocess.py`), or the native + Lang-SDK Dag parsing sources change. Those are the Dag processing files `importer_routing.py` and + `lang_sdk_processor.py`, `models/dagcode.py`, `serialization/serialized_objects.py`, + the Task SDK importers (`task-sdk/src/airflow/sdk/importers/`), the coordinator registry + (`task-sdk/src/airflow/sdk/execution_time/coordinator.py`), and `coordinators/_dag_importer.py`; + they trigger the TypeScript SDK E2E tests too. Shared Dag processing files such as `manager.py` + are left to `canary` runs. Like the other + deployed e2e suites, enabling them forces `PROD Image building`. * `OpenLineage E2E tests` (the `openlineage` mode of the deployed-stack tests under `airflow-e2e-tests/tests/airflow_e2e_tests/openlineage_tests`, exposed as the `run-openlineage-e2e-tests` output) run when the `openlineage` or `common` providers or the @@ -640,7 +647,7 @@ GitHub Actions to pass the list of parameters to a command to execute | run-system-tests | Whether system tests should be run ("true"/"false") | true | | | run-task-sdk-tests | Whether Task SDK tests should be run ("true"/"false") | true | | | run-ts-sdk-docs | Whether the TypeScript SDK API reference should be built — on `ts-sdk/api-docs/`, `ts-sdk/docs/`, or `ts-sdk/src/` changes, including Markdown ("true"/"false") | true | | -| run-ts-sdk-e2e-tests | Whether TypeScript SDK e2e tests should be run — on runtime-affecting `ts-sdk/`, TS e2e test, or Node coordinator changes ("true"/"false") | true | | +| run-ts-sdk-e2e-tests | Whether TypeScript SDK e2e tests should be run — on runtime-affecting `ts-sdk/`, TS e2e test, Node coordinator, or native Lang-SDK Dag parsing changes ("true"/"false") | true | | | run-ui-tests | Whether UI tests should be run ("true"/"false") | true | | | run-unit-tests | Whether unit tests should be run ("true"/"false") | true | | | run-www-tests | Whether Legacy WWW tests should be run ("true"/"false") | true | | diff --git a/dev/breeze/doc/images/output_k8s.svg b/dev/breeze/doc/images/output_k8s.svg index 085b233e8aa3e..6ba5d750acf29 100644 --- a/dev/breeze/doc/images/output_k8s.svg +++ b/dev/breeze/doc/images/output_k8s.svg @@ -1,4 +1,4 @@ -