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:

Reply via email to