jason810496 commented on code in PR #73847:
URL: https://github.com/apache/airflow/pull/73847#discussion_r4163184711


##########
dev/breeze/src/airflow_breeze/commands/kubernetes_commands.py:
##########
@@ -2702,6 +2710,93 @@ def _lang_sdk_build_go_bundle(
         shutil.copy(output_bin, go_dir / LANG_SDK_GO_BUNDLE_NAME)
 
 
+def _lang_sdk_build_ts_bundle(staging: Path, output: Output | None, *, native: 
bool = False) -> None:
+    """Pack ``ts-sdk/example`` into ``staging/ts-artifacts/bundle.min.mjs``.
+
+    The example depends on the in-repo ts-sdk by a workspace link, so the pack 
runs against the
+    checkout rather than a scratch copy: unlike Go's module replace, pnpm 
resolves the link from
+    the workspace root and a copied tree would lose it.
+
+    By default the build runs in an ephemeral Node toolchain container so the 
host needs no Node
+    install. In ``native`` mode (used in CI, where ``actions/setup-node`` has 
already provisioned
+    and cached one) it invokes the host toolchain directly.
+    """
+    ts_dir = staging / "ts-artifacts"
+    ts_dir.mkdir(parents=True, exist_ok=True)
+    pack = ["pnpm", "--filter", "apache-airflow-ts-sdk-example", "run", 
"build"]

Review Comment:
   Done in b3c7933b05. Both paths install and build the SDK before the example, 
and the container enables corepack in a writable home. The docstring no longer 
says CI provisions Node.



##########
ts-sdk/example/README.md:
##########
@@ -25,7 +25,14 @@ This example shows the coordinator-mode shape for TypeScript 
task handlers:
 - `src/main.ts` and `src/taskflow.ts` register a `TaskHandler` per stub task 
and start the coordinator runtime.
   One bundle provides for both Dags, and both declare a task called 
`build_message`.
   A handler binds the `(dag_id, task_id)` pair, so the two are different tasks 
with different bodies.
-- `dist/bundle.min.mjs` is the generated Node.js bundle that Airflow launches.
+- `src/native.ts` declares a third Dag, `typescript_native_example`, 
**entirely in TypeScript** — no

Review Comment:
   Done in bb98968fdc. The setup no longer uses `bundles_root`: #74042 adds 
`dist/` as the `ts-example` Dag bundle, and `ts`, the README's only Node 
coordinator, parses the native Dag from `bundle.min.mjs` there. I added that 
the Dag processor needs `node`, a link to the parsing docs, and the trigger 
line.



##########
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, where the ``ts-native`` coordinator claims 
it. Every task, the
+trigger included, runs in the TypeScript runtime.
+"""
+
+from __future__ import annotations
+
+from dataclasses import dataclass
+from datetime import datetime, timezone
+
+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"
+        )
+
+
[email protected](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
+
+
[email protected](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(timezone.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_bundle_entry_module(parsed_dag: AirflowClient):
+    """The Code view shows the TypeScript entry module the bundle embeds, not 
the bundle itself."""
+    content = parsed_dag.get_dag_source(_DAG_ID)["content"]
+
+    assert "new Bundle()" 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_a_cleared_skipped_branch_stays_skipped(completed_run: _CompletedRun):

Review Comment:
   Done in ceb3d9a259. Renamed to 
`test_the_skipped_branches_are_recorded_in_xcom`.



##########
dev/breeze/src/airflow_breeze/utils/selective_checks.py:
##########
@@ -181,6 +181,16 @@ def __hash__(self):
         return hash(frozenset(self))
 
 
+# Core and Task SDK sources on the native Lang-SDK Dag path: discovery, 
parsing through a

Review Comment:
   Done in 19e07411de. The jobs no longer key on `dagbag.py`, `manager.py` or 
`processor.py`, which are left to the canary run, and now also key on 
`serialized_objects.py`.



##########
ts-sdk/tsconfig.json:
##########
@@ -23,10 +23,18 @@
     // narrows `include` back to src, so none of this reaches the package.
     "allowJs": true,
 
+    // `example/` imports the package by name, which would otherwise resolve to
+    // the built `dist/`. A test that reached into both would hold two copies 
of
+    // the SDK, whose classes are not each other's. vitest.config.ts carries 
the
+    // matching alias so the compiler and the runner agree.
+    "paths": {

Review Comment:
   Done in 448b13ffbb. `example/tsconfig.json` sets `"paths": {}`, so the pack 
uses `dist/`, and the example typecheck no longer reports TS6059.



##########
kubernetes-tests/lang_sdk/README.md:
##########
@@ -24,6 +24,12 @@ End-to-end test that one Dag mixing **Python + Go + Java** 
tasks runs to success
 `[sdk] coordinators` config. The test lives at
 
`kubernetes-tests/tests/kubernetes_tests/test_lang_sdk_coordinator_executor.py`.
 
+The same file also holds `TestNativeTypeScriptDagOnKubernetes`, which runs a 
Dag with **no Python
+file at all**: the Dag processor asks the packed TypeScript bundle to parse 
itself. It is gated on

Review Comment:
   Done in 64af0954bd and 13b448d786. The README names the real skip variable, 
and both it and the test comment name the two gaps left: `bundle.min.mjs` 
reaches only worker pods, and the Dag processor image has no Node. The 
coordinator is no longer a gap, since the chart's only Node coordinator now 
parses packed `*.min.mjs` bundles in every Dag bundle.



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to