kaxil commented on code in PR #73847:
URL: https://github.com/apache/airflow/pull/73847#discussion_r4159731160
##########
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:
On a clean checkout this fails with `sh: airflow-ts-pack: not found`, even
once the corepack shim is sorted: the example's `airflow-ts-pack` bin points at
the SDK's `dist/cli/main.js`, and nothing here builds the SDK, so `pnpm
install` skips the bin link. The compose helper's sequence (install, `pnpm run
build` in ts-sdk, then install and build the example) works here too. The
native branch has the same gap and runs no install at all, and the docstring's
"CI provisions Node with `actions/setup-node`" doesn't match k8s-tests.yml,
which has no such step.
##########
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:
The stack below does dispatch the parse request now, so "not in Airflow yet"
(and the matching comment in test_lang_sdk_coordinator_executor.py) isn't the
reason. What's missing is provisioning: `bundle.min.mjs` only goes to the
`ts-artifacts` bucket, which only worker pods stage, the chart's only Node
coordinator has `bundles_root` so it parses nothing, and the Dag processor has
no Node. Could this name those three? The skipif also checks only
`RUN_TS_SDK_NATIVE_DAG_K8S_TESTS`, not `RUN_LANG_SDK_K8S_TESTS` as this says.
##########
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:
This keys the Java and TS e2e jobs, and a PROD image build, on `manager.py`,
`processor.py` and `dagbag.py`, the most-edited files in Dag processing.
Replaying apache/main since April, it would have added both jobs to 54 of the
62 commits that touch these paths and a PROD image to 37, about nine PRs a
month, mostly for changes off the Lang-SDK path. Would keying it on the
Lang-SDK modules (`lang_sdk_processor.py`, `importer_routing.py`,
`coordinators/_dag_importer.py`, `execution_time/coordinator.py`) and leaving
the shared files to the canary run be enough? It also misses
`serialization/serialized_objects.py`, where `fill_config_defaults` and
`validate_serialized_dag` live, so a change to the checks every native Dag goes
through runs neither job.
##########
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:
The configuration further down only sets up `ts` with `bundles_root`, and a
coordinator with a root parses nothing, so following these steps gives
`typescript_example` and `typescript_taskflow_example` but never
`typescript_native_example`, with no import error to say why. Could the setup
add the root-less coordinator the compose e2e uses, say `bundle.min.mjs` goes
in the Dags folder, and note the Dag processor needs `node`? Linking to
"Parsing native Dags" in typescript.rst would also do.
##########
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:
esbuild reads `paths` through the `extends` chain too, so
example/tsconfig.json inherits this and `airflow-ts-pack` now bundles the SDK
from `../src/*.ts` instead of the built package. The compose and k8s e2e then
stop covering the `dist/` that ships (a `tsconfig.build.json` slip would no
longer fail them), and every SDK module gets a source tag. Setting `"paths":
{}` in example/tsconfig.json keeps the alias for the root typecheck and vitest,
and should also clear the TS6059 listed under known issues.
##########
ts-sdk/adr/0002-native-dag-interface.md:
##########
@@ -21,7 +21,7 @@
## Status
-Proposed. Revised after the review on #72047.
+Accepted. Revised after the review on #72047, and again as it was implemented.
Review Comment:
Is there a vote or review that accepted this ADR? The native-Dag ADRs for Go
(0008) and Java (0002), and the TS 0001, are all still Proposed, and flipping
one to Accepted in an e2e PR without a link makes it hard to tell later when
the decision was made.
##########
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:
This doesn't clear anything; it checks the `skipmixin_key` XCom that
`has_rows` pushed. Could it be renamed to what it checks, or actually clear
`report_empty` and assert it ends skipped again? As named it reads as if the
clear path were covered.
--
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]