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 75ac609e9bc561d97820dbf13957d549a97e4790 Author: ZHE YOU LIU <[email protected]> AuthorDate: Mon Sep 28 07:44:42 2026 +0000 Parse the native TypeScript Dag in the compose e2e Put the example bundle in the Dag bundle as well, add a root-less "ts-native" Node coordinator that claims it there, and give the Dag processor node. The existing "ts" coordinator still runs every TypeScript task, so the native Dag test needs no gate anymore. --- airflow-e2e-tests/docker/ts.yml | 13 +++++++++++-- airflow-e2e-tests/tests/airflow_e2e_tests/conftest.py | 15 ++++++++++++++- .../ts_sdk_tests/test_ts_sdk_native_dag.py | 17 +++-------------- 3 files changed, 28 insertions(+), 17 deletions(-) diff --git a/airflow-e2e-tests/docker/ts.yml b/airflow-e2e-tests/docker/ts.yml index 056106ff9d5..902ce7562c9 100644 --- a/airflow-e2e-tests/docker/ts.yml +++ b/airflow-e2e-tests/docker/ts.yml @@ -17,10 +17,12 @@ # Docker Compose override for ts_sdk E2E test mode. # -# The stock worker image ships no Node.js runtime, so node-provider copies the +# The stock Airflow image ships no Node.js runtime, so node-provider copies the # node binary from the same image conftest builds the bundle with into a # shared volume. The bundle is bind-mounted where NodeCoordinator scans, and -# the worker consumes the "typescript" queue where @task.stub tasks are routed. +# the worker consumes the "typescript" queue where TypeScript tasks are routed. +# The Dag processor gets node too, to parse the bundle's native Dag from the +# Dag bundle. --- services: node-provider: @@ -29,6 +31,13 @@ services: volumes: - nodejs-bin:/opt/nodejs + airflow-dag-processor: + volumes: + - nodejs-bin:/opt/nodejs:ro + depends_on: + node-provider: + condition: service_completed_successfully + airflow-worker: volumes: - ./ts-bundles:/opt/airflow/ts-bundles:ro diff --git a/airflow-e2e-tests/tests/airflow_e2e_tests/conftest.py b/airflow-e2e-tests/tests/airflow_e2e_tests/conftest.py index 39ff3fd0a64..4b441bb549f 100644 --- a/airflow-e2e-tests/tests/airflow_e2e_tests/conftest.py +++ b/airflow-e2e-tests/tests/airflow_e2e_tests/conftest.py @@ -724,12 +724,21 @@ def _setup_ts_sdk_integration(dot_env_file, tmp_dir): ts_bundles_dir.mkdir() # Deliberately renamed: the coordinator routes on embedded metadata, not on a fixed name. copyfile(TS_SDK_EXAMPLE_PATH / "dist" / "bundle.min.mjs", ts_bundles_dir / "example.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. for dag_file in ("typescript_example.py", "typescript_taskflow_example.py"): copyfile(TS_SDK_EXAMPLE_PATH / "dags" / dag_file, tmp_dir / "dags" / dag_file) + # "ts" runs every TypeScript task, stub or native, from /opt/airflow/ts-bundles. + # "ts-native" has no root, so it serves the Dag bundle: it parses the bundle's + # native Dag, in the Dag processor and on the worker that runs the Dag's + # Python task (trigger_downstream). coordinator_config = json.dumps( { "ts": { @@ -738,7 +747,11 @@ def _setup_ts_sdk_integration(dot_env_file, tmp_dir): "bundles_root": ["/opt/airflow/ts-bundles"], "node_executable": "/opt/nodejs/node", }, - } + }, + "ts-native": { + "classpath": "airflow.sdk.coordinators.node.NodeCoordinator", + "kwargs": {"node_executable": "/opt/nodejs/node"}, + }, } ) queue_to_coordinator = json.dumps({"typescript": "ts"}) 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 index f23684d7a5b..f8aa1f48a72 100644 --- 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 @@ -19,8 +19,7 @@ End-to-end test of a Dag declared entirely in TypeScript. Run with:: - E2E_TEST_MODE=ts_sdk RUN_TS_SDK_NATIVE_DAG_TESTS=true \\ - uv run --project airflow-e2e-tests pytest \\ + 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 @@ -29,14 +28,12 @@ Unlike ``test_ts_sdk_dag.py``, no Python file declares this Dag: the Dag process named fan-in, order-only edges, a conditional, a multi-way branch, and a ``TriggerDagRunOperator`` a Python worker runs. -Gated behind ``RUN_TS_SDK_NATIVE_DAG_TESTS`` because it needs a Dag processor that can dispatch a -parse request to a language coordinator. Until that lands, the bundle answers a request nothing -sends, and the Dag never appears. +The bundle sits in the Dag bundle, where the ``ts-native`` coordinator claims it. The Python worker +that runs ``trigger_downstream`` parses the bundle the same way to rebuild the operator. """ from __future__ import annotations -import os from dataclasses import dataclass from datetime import datetime, timezone @@ -44,8 +41,6 @@ import pytest from airflow_e2e_tests.e2e_test_utils.clients import AirflowClient -_RUN_NATIVE = os.environ.get("RUN_TS_SDK_NATIVE_DAG_TESTS", "").lower() in ("true", "1") - # 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 @@ -57,12 +52,6 @@ _NORTH_ROWS_VARIABLE = "typescript_native_north_rows" _SOUTH_ROWS_VARIABLE = "typescript_native_south_rows" _CADENCE_VARIABLE = "typescript_native_cadence" -pytestmark = pytest.mark.skipif( - not _RUN_NATIVE, - reason="Needs a Dag processor that dispatches parse requests to a language coordinator " - "(RUN_TS_SDK_NATIVE_DAG_TESTS)", -) - @dataclass class _CompletedRun:
