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 bac79752c037ebebb20d8831a944256a1dba6fc3
Author: ZHE YOU LIU <[email protected]>
AuthorDate: Mon Sep 28 07:46:10 2026 +0000

    Check the native TypeScript Dag's source and triggered run
    
    The Code view must show the bundle's entry module, and the Python task of
    the native Dag must trigger typescript_example with its rendered conf.
---
 .../airflow_e2e_tests/e2e_test_utils/clients.py    |  7 ++++++
 .../ts_sdk_tests/test_ts_sdk_native_dag.py         | 26 ++++++++++++++++++++++
 2 files changed, 33 insertions(+)

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 803edf515e7..ee541f97f7d 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
@@ -179,6 +179,13 @@ class AirflowClient:
         """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_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}")
+
+    def get_dag_run(self, dag_id: str, run_id: str):
+        return self._make_request(method="GET", 
endpoint=f"dags/{dag_id}/dagRuns/{run_id}")
+
     def trigger_dag_and_wait(self, dag_id: str, json=None):
         """Trigger a DAG and wait for it to complete."""
         self.un_pause_dag(dag_id)
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 f8aa1f48a72..f4d304e17e8 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
@@ -46,6 +46,8 @@ from airflow_e2e_tests.e2e_test_utils.clients import 
AirflowClient
 _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"
@@ -89,6 +91,8 @@ def completed_run(parsed_dag: AirflowClient) -> _CompletedRun:
     ):
         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(timezone.utc).isoformat()})
     run_id = resp["dag_run_id"]
@@ -113,6 +117,14 @@ def 
test_the_dag_the_bundle_parsed_is_registered(parsed_dag: AirflowClient):
     assert {tag["name"] for tag in dag.get("tags") or []} >= {"typescript", 
"native"}
 
 
+def test_the_dag_source_is_the_bundle_entry_module(parsed_dag: AirflowClient):
+    """The Code view shows the TypeScript the bundle was packed from, not the 
minified bundle."""
+    content = parsed_dag.get_dag_source(_DAG_ID)["content"]
+
+    assert "new Bundle()" in content
+    assert 'from "./native.js"' 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", [])
@@ -176,6 +188,20 @@ def test_every_other_task_succeeded(completed_run: 
_CompletedRun):
         )
 
 
+def test_the_trigger_started_the_downstream_run(completed_run: _CompletedRun):
+    """A Python worker rebuilt ``trigger_downstream`` from the bundle 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"
+    # The conf template is rendered on the worker, from the rebuilt operator.
+    assert run["conf"] == {"triggered_by": _DAG_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}

Reply via email to