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}
