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 3801de79ad614b660541975264d0e76e18722e7d Author: ZHE YOU LIU <[email protected]> AuthorDate: Mon Sep 21 09:11:04 2026 +0000 TS SDK: cover native TypeScript Dags end to end (cherry picked from commit 2b76bb3fae1a4c1c5f08ae47d46ec5b0171f4bf6) --- .../language-sdks/typescript.rst | 6 + .../airflow_e2e_tests/e2e_test_utils/clients.py | 16 ++ .../ts_sdk_tests/test_ts_sdk_native_dag.py | 209 +++++++++++++++++++++ dev/breeze/doc/images/output_k8s.svg | 8 +- dev/breeze/doc/images/output_k8s.txt | 2 +- .../doc/images/output_k8s_setup-lang-sdk-test.svg | 38 ++-- .../doc/images/output_k8s_setup-lang-sdk-test.txt | 2 +- .../airflow_breeze/commands/kubernetes_commands.py | 165 ++++++++++++++-- .../commands/kubernetes_commands_config.py | 1 + kubernetes-tests/lang_sdk/Dockerfile.typescript | 41 ++++ kubernetes-tests/lang_sdk/README.md | 17 +- kubernetes-tests/lang_sdk/config/values.yaml | 12 +- .../pod_templates/lang_sdk_typescript.yaml | 115 ++++++++++++ .../test_lang_sdk_coordinator_executor.py | 80 ++++++++ ts-sdk/adr/0002-native-dag-interface.md | 18 +- ts-sdk/example/README.md | 8 +- ts-sdk/example/src/main.ts | 12 +- ts-sdk/example/src/native.ts | 163 ++++++++++++++++ ts-sdk/tests/sdk/native-dag-example.test.ts | 207 ++++++++++++++++++++ ts-sdk/tsconfig.json | 10 +- ts-sdk/vitest.config.ts | 11 ++ 21 files changed, 1089 insertions(+), 52 deletions(-) diff --git a/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst b/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst index 23f60986224..77726061b5d 100644 --- a/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst +++ b/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst @@ -445,6 +445,12 @@ in the task spec. Templated arguments pass through untouched, so ``{{ ds }}`` in a ``conf`` value is rendered server-side where rendering already happens. +A worked example +~~~~~~~~~~~~~~~~ + +``ts-sdk/example/src/native.ts`` puts the constructs above into one Dag, registered on the same +bundle as the mixed-language handlers beside it, so a single artifact serves both authoring modes. + ``new Dag`` and ``dag.task`` both take a trailing spec of Airflow options: ``{ schedule: "@daily", tags: ["etl"] }`` for the Dag, ``{ retries: 2, retryDelay: 30 }`` for a task. 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 d56cecae3d7..803edf515e7 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 @@ -19,6 +19,7 @@ from __future__ import annotations import time from datetime import datetime, timezone from functools import cached_property +from http import HTTPStatus import boto3 import requests @@ -163,6 +164,21 @@ class AirflowClient: """Get an Airflow Variable via API.""" return self._make_request(method="GET", endpoint=f"variables/{key}") + def set_variable(self, key: str, value: str, description: str | None = None): + """Create or replace an Airflow Variable via API.""" + body = {"key": key, "value": value, "description": description} + try: + return self._make_request(method="POST", endpoint="variables", json=body) + except requests.HTTPError as exc: + # 409 == it already exists, from an earlier run of the same suite. + if exc.response is None or exc.response.status_code != HTTPStatus.CONFLICT: + raise + return self._make_request(method="PATCH", endpoint=f"variables/{key}", json=body) + + def get_tasks(self, dag_id: str): + """List a Dag's tasks, with the edges each one carries.""" + return self._make_request(method="GET", endpoint=f"dags/{dag_id}/tasks") + 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 new file mode 100644 index 00000000000..f23684d7a5b --- /dev/null +++ b/airflow-e2e-tests/tests/airflow_e2e_tests/ts_sdk_tests/test_ts_sdk_native_dag.py @@ -0,0 +1,209 @@ +# 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 RUN_TS_SDK_NATIVE_DAG_TESTS=true \\ + 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 ``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. +""" + +from __future__ import annotations + +import os +from dataclasses import dataclass +from datetime import datetime, timezone + +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 + +_DAG_ID = "typescript_native_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" + +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: + 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) + + 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_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 + } + # And the one task Python runs rather than the Node coordinator. + 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_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): + """What `_can_skip_downstream` plus the skipmixin XCom are for.""" + assert completed_run.xcom("has_rows", key="skipmixin_key") == {"skipped": ["report_empty"]} + assert completed_run.xcom("pick_cadence", key="skipmixin_key") == {"skipped": ["publish_daily"]} diff --git a/dev/breeze/doc/images/output_k8s.svg b/dev/breeze/doc/images/output_k8s.svg index 085b233e8aa..92f835a53e8 100644 --- a/dev/breeze/doc/images/output_k8s.svg +++ b/dev/breeze/doc/images/output_k8s.svg @@ -206,10 +206,10 @@ </text><text class="breeze-k8s-r5" x="0" y="727.6" textLength="12.2" clip-path="url(#breeze-k8s-line-29)">│</text><text class="breeze-k8s-r1" x="280.6" y="727.6" textLength="1159" clip-path="url(#breeze-k8s-line-29)">resources, and run the optional per-overlay pytest module.                              &# [...] </text><text class="breeze-k8s-r5" x="0" y="752" textLength="12.2" clip-path="url(#breeze-k8s-line-30)">│</text><text class="breeze-k8s-r4" x="24.4" y="752" textLength="231.8" clip-path="url(#breeze-k8s-line-30)">run-complete-tests </text><text class="breeze-k8s-r1" x="280.6" y="752" textLength="1159" clip-path="url(#breeze-k8s-line-30)">Run complete k8s tests consisting of: creating cluster, building and uploading image, d [...] </text><text class="breeze-k8s-r5" x="0" y="776.4" textLength="12.2" clip-path="url(#breeze-k8s-line-31)">│</text><text class="breeze-k8s-r1" x="280.6" y="776.4" textLength="1159" clip-path="url(#breeze-k8s-line-31)">airflow, running tests and deleting clusters (optionally for all clusters in parallel).        </text><text class="breeze-k8s-r5" x="1451.8" y="776.4" textLength="12.2" clip-path=" [...] -</text><text class="breeze-k8s-r5" x="0" y="800.8" textLength="12.2" clip-path="url(#breeze-k8s-line-32)">│</text><text class="breeze-k8s-r4" x="24.4" y="800.8" textLength="231.8" clip-path="url(#breeze-k8s-line-32)">setup-lang-sdk-test</text><text class="breeze-k8s-r1" x="280.6" y="800.8" textLength="1159" clip-path="url(#breeze-k8s-line-32)">Provision the lang-SDK (Go + Java) coordinator system test on an already-deployed  [...] -</text><text class="breeze-k8s-r5" x="0" y="825.2" textLength="12.2" clip-path="url(#breeze-k8s-line-33)">│</text><text class="breeze-k8s-r1" x="280.6" y="825.2" textLength="1159" clip-path="url(#breeze-k8s-line-33)">KubernetesExecutor cluster: build artifacts, build + load the Java worker image, deploy        </text><text class="breeze-k8s-r5" x="1451.8" y="825.2" textLength="12.2" clip-path=" [...] -</text><text class="breeze-k8s-r5" x="0" y="849.6" textLength="12.2" clip-path="url(#breeze-k8s-line-34)">│</text><text class="breeze-k8s-r1" x="280.6" y="849.6" textLength="1159" clip-path="url(#breeze-k8s-line-34)">localstack S3, upload artifacts + stub Dag, create config, and upgrade the Helm release. Run   </text><text class="breeze-k8s-r5" x="1451.8" y="849.6" textLength="12.2" clip-path="url(#breez [...] -</text><text class="breeze-k8s-r5" x="0" y="874" textLength="12.2" clip-path="url(#breeze-k8s-line-35)">│</text><text class="breeze-k8s-r1" x="280.6" y="874" textLength="866.2" clip-path="url(#breeze-k8s-line-35)">the test afterwards with `RUN_LANG_SDK_K8S_TESTS=true breeze k8s tests </text><text class="breeze-k8s-r4" x="1146.8" y="874" textLength="122" clip-path="url(#breeze-k8s-line-35)">--executor</text><text class="breeze-k8s-r5" x="1451.8" y=" [...] +</text><text class="breeze-k8s-r5" x="0" y="800.8" textLength="12.2" clip-path="url(#breeze-k8s-line-32)">│</text><text class="breeze-k8s-r4" x="24.4" y="800.8" textLength="231.8" clip-path="url(#breeze-k8s-line-32)">setup-lang-sdk-test</text><text class="breeze-k8s-r1" x="280.6" y="800.8" textLength="1159" clip-path="url(#breeze-k8s-line-32)">Provision the lang-SDK (Go + Java + TypeScript) coordinator system test on an alr [...] +</text><text class="breeze-k8s-r5" x="0" y="825.2" textLength="12.2" clip-path="url(#breeze-k8s-line-33)">│</text><text class="breeze-k8s-r1" x="280.6" y="825.2" textLength="1159" clip-path="url(#breeze-k8s-line-33)">KubernetesExecutor cluster: build artifacts, build + load the Java and TypeScript worker       </text><text class="breeze-k8s-r5" x="1451.8" y="825.2" textLength="12.2" clip-path="url(# [...] +</text><text class="breeze-k8s-r5" x="0" y="849.6" textLength="12.2" clip-path="url(#breeze-k8s-line-34)">│</text><text class="breeze-k8s-r1" x="280.6" y="849.6" textLength="1159" clip-path="url(#breeze-k8s-line-34)">images, deploy localstack S3, upload artifacts + stub Dag, create config, and upgrade the Helm </text><text class="breeze-k8s-r5" x="1451.8" y="849.6" textLength="12.2" clip-path="url(#breeze-k8s-line [...] +</text><text class="breeze-k8s-r5" x="0" y="874" textLength="12.2" clip-path="url(#breeze-k8s-line-35)">│</text><text class="breeze-k8s-r1" x="280.6" y="874" textLength="1024.8" clip-path="url(#breeze-k8s-line-35)">release. Run the test afterwards with `RUN_LANG_SDK_K8S_TESTS=true breeze k8s tests </text><text class="breeze-k8s-r4" x="1305.4" y="874" textLength="122" clip-path="url(#breeze-k8s-line-35)">--executor</text><text class="breez [...] </text><text class="breeze-k8s-r5" x="0" y="898.4" textLength="12.2" clip-path="url(#breeze-k8s-line-36)">│</text><text class="breeze-k8s-r1" x="280.6" y="898.4" textLength="268.4" clip-path="url(#breeze-k8s-line-36)">KubernetesExecutor -- </text><text class="breeze-k8s-r6" x="549" y="898.4" textLength="24.4" clip-path="url(#breeze-k8s-line-36)">-k</text><text class="breeze-k8s-r1" x="573.4" y="898.4" textLength="866.2" clip-path="url(#breeze-k8s-line-36)"> test_lang_sdk_c [...] </text><text class="breeze-k8s-r5" x="0" y="922.8" textLength="12.2" clip-path="url(#breeze-k8s-line-37)">│</text><text class="breeze-k8s-r4" x="24.4" y="922.8" textLength="231.8" clip-path="url(#breeze-k8s-line-37)">shell              </text><text class="breeze-k8s-r1" x="280.6" y="922.8" textLength="1159" clip-path="url(#breeze-k8s-line-37)">Run shell environment for the current KinD [...] </text><text class="breeze-k8s-r5" x="0" y="947.2" textLength="12.2" clip-path="url(#breeze-k8s-line-38)">│</text><text class="breeze-k8s-r4" x="24.4" y="947.2" textLength="231.8" clip-path="url(#breeze-k8s-line-38)">k9s                </text><text class="breeze-k8s-r1" x="280.6" y="947.2" textLength="1159" clip-path="url(#breeze-k8s-line-38)">Run k9s tool. You can pass any  [...] diff --git a/dev/breeze/doc/images/output_k8s.txt b/dev/breeze/doc/images/output_k8s.txt index d0c38089093..b983d60ab98 100644 --- a/dev/breeze/doc/images/output_k8s.txt +++ b/dev/breeze/doc/images/output_k8s.txt @@ -1 +1 @@ -09b8773fcafec9f720cc6512fa12f73d +e6bce3a1cc0d3660f16669fe5c4b1338 diff --git a/dev/breeze/doc/images/output_k8s_setup-lang-sdk-test.svg b/dev/breeze/doc/images/output_k8s_setup-lang-sdk-test.svg index 55b1d78afae..a9e646da6c4 100644 --- a/dev/breeze/doc/images/output_k8s_setup-lang-sdk-test.svg +++ b/dev/breeze/doc/images/output_k8s_setup-lang-sdk-test.svg @@ -1,4 +1,4 @@ -<svg class="rich-terminal" viewBox="0 0 1482 611.1999999999999" xmlns="http://www.w3.org/2000/svg"> +<svg class="rich-terminal" viewBox="0 0 1482 684.4" xmlns="http://www.w3.org/2000/svg"> <!-- Generated with Rich https://www.textualize.io --> <style> @@ -43,7 +43,7 @@ <defs> <clipPath id="breeze-k8s-setup-lang-sdk-test-clip-terminal"> - <rect x="0" y="0" width="1463.0" height="560.1999999999999" /> + <rect x="0" y="0" width="1463.0" height="633.4" /> </clipPath> <clipPath id="breeze-k8s-setup-lang-sdk-test-line-0"> <rect x="0" y="1.5" width="1464" height="24.65"/> @@ -111,9 +111,18 @@ <clipPath id="breeze-k8s-setup-lang-sdk-test-line-21"> <rect x="0" y="513.9" width="1464" height="24.65"/> </clipPath> +<clipPath id="breeze-k8s-setup-lang-sdk-test-line-22"> + <rect x="0" y="538.3" width="1464" height="24.65"/> + </clipPath> +<clipPath id="breeze-k8s-setup-lang-sdk-test-line-23"> + <rect x="0" y="562.7" width="1464" height="24.65"/> + </clipPath> +<clipPath id="breeze-k8s-setup-lang-sdk-test-line-24"> + <rect x="0" y="587.1" width="1464" height="24.65"/> + </clipPath> </defs> - <rect fill="#292929" stroke="rgba(255,255,255,0.35)" stroke-width="1" x="1" y="1" width="1480" height="609.2" rx="8"/><text class="breeze-k8s-setup-lang-sdk-test-title" fill="#c5c8c6" text-anchor="middle" x="740" y="27">Command: k8s setup-lang-sdk-test</text> + <rect fill="#292929" stroke="rgba(255,255,255,0.35)" stroke-width="1" x="1" y="1" width="1480" height="682.4" rx="8"/><text class="breeze-k8s-setup-lang-sdk-test-title" fill="#c5c8c6" text-anchor="middle" x="740" y="27">Command: k8s setup-lang-sdk-test</text> <g transform="translate(26,22)"> <circle cx="0" cy="0" r="7" fill="#ff5f57"/> <circle cx="22" cy="0" r="7" fill="#febc2e"/> @@ -126,10 +135,10 @@ <text class="breeze-k8s-setup-lang-sdk-test-r1" x="1464" y="20" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-0)"> </text><text class="breeze-k8s-setup-lang-sdk-test-r2" x="12.2" y="44.4" textLength="73.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-1)">Usage:</text><text class="breeze-k8s-setup-lang-sdk-test-r3" x="97.6" y="44.4" textLength="366" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-1)">breeze k8s setup-lang-sdk-test</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="475.8" y="44.4" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-1)">[</te [...] </text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="1464" y="68.8" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-2)"> -</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="12.2" y="93.2" textLength="1415.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-3)">Provision the lang-SDK (Go + Java) coordinator system test on an already-deployed KubernetesExecutor cluster: build </text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="1464" y="93.2" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-3)"> -</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="12.2" y="117.6" textLength="1427.4" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-4)">artifacts, build + load the Java worker image, deploy localstack S3, upload artifacts + stub Dag, create config, and </text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="1464" y="117.6" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang- [...] -</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="12.2" y="142" textLength="1232.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-5)">upgrade the Helm release. Run the test afterwards with `RUN_LANG_SDK_K8S_TESTS=true breeze k8s tests </text><text class="breeze-k8s-setup-lang-sdk-test-r4" x="1244.4" y="142" textLength="122" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-5)">--executor</text><text class="b [...] -</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="12.2" y="166.4" textLength="268.4" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-6)">KubernetesExecutor -- </text><text class="breeze-k8s-setup-lang-sdk-test-r5" x="280.6" y="166.4" textLength="24.4" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-6)">-k</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="305" y="166.4" textLength="463.6" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-6)"> test_l [...] +</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="12.2" y="93.2" textLength="1390.8" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-3)">Provision the lang-SDK (Go + Java + TypeScript) coordinator system test on an already-deployed KubernetesExecutor </text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="1464" y="93.2" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-3)"> +</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="12.2" y="117.6" textLength="1439.6" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-4)">cluster: build artifacts, build + load the Java and TypeScript worker images, deploy localstack S3, upload artifacts +</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="1464" y="117.6" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test- [...] +</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="12.2" y="142" textLength="1378.6" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-5)">stub Dag, create config, and upgrade the Helm release. Run the test afterwards with `RUN_LANG_SDK_K8S_TESTS=true </text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="1464" y="142" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-5)"> +</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="12.2" y="166.4" textLength="207.4" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-6)">breeze k8s tests </text><text class="breeze-k8s-setup-lang-sdk-test-r4" x="219.6" y="166.4" textLength="122" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-6)">--executor</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="341.6" y="166.4" textLength="280.6" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-6)"> [...] </text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="1464" y="190.8" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-7)"> </text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="215.2" textLength="24.4" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-8)">╭─</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="24.4" y="215.2" textLength="305" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-8)"> K8S lang-SDK test flags </text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="329.4" y="215.2" textLength="1110.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-8 [...] </text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="239.6" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-9)">│</text><text class="breeze-k8s-setup-lang-sdk-test-r4" x="24.4" y="239.6" textLength="244" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-9)">--python            </text><text class="breeze-k8s-setup-lang-sdk-test-r5" x="292.8" y="239.6" textLength="24.4" clip-path="url(#breeze-k8s [...] @@ -140,12 +149,15 @@ </text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="361.6" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-14)">│</text><text class="breeze-k8s-setup-lang-sdk-test-r4" x="24.4" y="361.6" textLength="244" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-14)">--java-image        </text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="341.6" y="361.6" textLength="1098" clip-path="url(#breeze-k8s-setup-lang-sdk-te [...] </text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="386" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-15)">│</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="341.6" y="386" textLength="1098" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-15)">the prod image plus a headless JRE (Dockerfile.java) and loading it into the kind cluster.</text><text class="breeze-k8s-setup-la [...] </text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="410.4" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-16)">│</text><text class="breeze-k8s-setup-lang-sdk-test-r7" x="341.6" y="410.4" textLength="73.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-16)">(TEXT)</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="1451.8" y="410.4" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-16)">│</text><text class="breeze-k8s-setup- [...] -</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="434.8" textLength="1464" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-17)">╰──────────────────────────────────────────────────────────────────────────────────────────────────────────────────────╯</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="1464" y="434.8" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-17)"> -</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="459.2" textLength="24.4" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-18)">╭─</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="24.4" y="459.2" textLength="195.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-18)"> Common options </text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="219.6" y="459.2" textLength="1220" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-18)">───────────── [...] -</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="483.6" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-19)">│</text><text class="breeze-k8s-setup-lang-sdk-test-r4" x="24.4" y="483.6" textLength="109.8" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-19)">--verbose</text><text class="breeze-k8s-setup-lang-sdk-test-r5" x="158.6" y="483.6" textLength="24.4" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-19)">-v</text><text class="breeze-k8s-set [...] -</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="508" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-20)">│</text><text class="breeze-k8s-setup-lang-sdk-test-r4" x="24.4" y="508" textLength="109.8" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-20)">--dry-run</text><text class="breeze-k8s-setup-lang-sdk-test-r5" x="158.6" y="508" textLength="24.4" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-20)">-D</text><text class="breeze-k8s-setup-lan [...] -</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="532.4" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-21)">│</text><text class="breeze-k8s-setup-lang-sdk-test-r4" x="24.4" y="532.4" textLength="109.8" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-21)">--help   </text><text class="breeze-k8s-setup-lang-sdk-test-r5" x="158.6" y="532.4" textLength="24.4" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-21)">-h</text><text class= [...] -</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="556.8" textLength="1464" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-22)">╰──────────────────────────────────────────────────────────────────────────────────────────────────────────────────────╯</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="1464" y="556.8" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-22)"> +</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="434.8" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-17)">│</text><text class="breeze-k8s-setup-lang-sdk-test-r4" x="24.4" y="434.8" textLength="244" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-17)">--ts-image          </text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="341.6" y="434.8" textLength="1098" clip-path="url(#breeze-k8s-setup-l [...] +</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="459.2" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-18)">│</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="341.6" y="459.2" textLength="1098" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-18)">building the prod image plus Node (Dockerfile.typescript) and loading it into the kind    </text><text class="breez [...] +</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="483.6" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-19)">│</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="341.6" y="483.6" textLength="109.8" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-19)">cluster. </text><text class="breeze-k8s-setup-lang-sdk-test-r7" x="451.4" y="483.6" textLength="73.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-19)">(TEXT)</text><text class="bree [...] +</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="508" textLength="1464" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-20)">╰──────────────────────────────────────────────────────────────────────────────────────────────────────────────────────╯</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="1464" y="508" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-20)"> +</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="532.4" textLength="24.4" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-21)">╭─</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="24.4" y="532.4" textLength="195.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-21)"> Common options </text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="219.6" y="532.4" textLength="1220" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-21)">───────────── [...] +</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="556.8" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-22)">│</text><text class="breeze-k8s-setup-lang-sdk-test-r4" x="24.4" y="556.8" textLength="109.8" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-22)">--verbose</text><text class="breeze-k8s-setup-lang-sdk-test-r5" x="158.6" y="556.8" textLength="24.4" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-22)">-v</text><text class="breeze-k8s-set [...] +</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="581.2" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-23)">│</text><text class="breeze-k8s-setup-lang-sdk-test-r4" x="24.4" y="581.2" textLength="109.8" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-23)">--dry-run</text><text class="breeze-k8s-setup-lang-sdk-test-r5" x="158.6" y="581.2" textLength="24.4" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-23)">-D</text><text class="breeze-k8s-set [...] +</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="605.6" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-24)">│</text><text class="breeze-k8s-setup-lang-sdk-test-r4" x="24.4" y="605.6" textLength="109.8" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-24)">--help   </text><text class="breeze-k8s-setup-lang-sdk-test-r5" x="158.6" y="605.6" textLength="24.4" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-24)">-h</text><text class= [...] +</text><text class="breeze-k8s-setup-lang-sdk-test-r6" x="0" y="630" textLength="1464" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-25)">╰──────────────────────────────────────────────────────────────────────────────────────────────────────────────────────╯</text><text class="breeze-k8s-setup-lang-sdk-test-r1" x="1464" y="630" textLength="12.2" clip-path="url(#breeze-k8s-setup-lang-sdk-test-line-25)"> </text> </g> </g> diff --git a/dev/breeze/doc/images/output_k8s_setup-lang-sdk-test.txt b/dev/breeze/doc/images/output_k8s_setup-lang-sdk-test.txt index 732dfe90120..c9d173a905d 100644 --- a/dev/breeze/doc/images/output_k8s_setup-lang-sdk-test.txt +++ b/dev/breeze/doc/images/output_k8s_setup-lang-sdk-test.txt @@ -1 +1 @@ -e79130d4518f85b8b1783b9a730f78fa +6a9a0568a2b313fb5a6780431efd2937 diff --git a/dev/breeze/src/airflow_breeze/commands/kubernetes_commands.py b/dev/breeze/src/airflow_breeze/commands/kubernetes_commands.py index 564c8fcc1c0..6b2eb800c5b 100644 --- a/dev/breeze/src/airflow_breeze/commands/kubernetes_commands.py +++ b/dev/breeze/src/airflow_breeze/commands/kubernetes_commands.py @@ -2505,6 +2505,14 @@ LANG_SDK_GRADLE_CACHE_PATH = AIRFLOW_ROOT_PATH / "files" / "gradle" # pod_template_file base image. LANG_SDK_JAVA_WORKER_IMAGE = "lang-sdk-java-worker:latest" LANG_SDK_JAVA_DOCKERFILE = LANG_SDK_PATH / "Dockerfile.java" +# The typescript queue needs Node, which the prod image does not carry either. +# The packed bundle is `ts-sdk/example`, already built by that workspace, so +# there is no example dir of its own here. +LANG_SDK_TS_WORKER_IMAGE = "lang-sdk-ts-worker:latest" +LANG_SDK_TS_DOCKERFILE = LANG_SDK_PATH / "Dockerfile.typescript" +LANG_SDK_TS_EXAMPLE_PATH = AIRFLOW_ROOT_PATH / "ts-sdk" / "example" +LANG_SDK_TS_BUILDER_IMAGE = os.environ.get("NODE_BUILDER_IMAGE", "node:22-alpine") +LANG_SDK_TS_BUNDLE_NAME = "bundle.min.mjs" LANG_SDK_AWS_CONN_URI = ( "aws://test:test@/?region_name=us-east-1&" "endpoint_url=http%3A%2F%2Flocalstack.airflow.svc.cluster.local%3A4566" @@ -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"] + if native: + get_console(output=output).print("[info]Packing the TypeScript bundle with the host Node toolchain") + run_command(pack, cwd=AIRFLOW_ROOT_PATH / "ts-sdk", output=output, check=True) + else: + uid_gid = f"{os.getuid()}:{os.getgid()}" + get_console(output=output).print( + f"[info]Packing the TypeScript bundle in {LANG_SDK_TS_BUILDER_IMAGE}" + ) + run_command( + [ + "docker", + "run", + "--rm", + "-u", + uid_gid, + "-v", + f"{AIRFLOW_ROOT_PATH}:{AIRFLOW_ROOT_PATH}", + "-w", + str(AIRFLOW_ROOT_PATH / "ts-sdk"), + "-e", + "npm_config_cache=/tmp/.npm", + "-e", + "COREPACK_ENABLE_DOWNLOAD_PROMPT=0", + LANG_SDK_TS_BUILDER_IMAGE, + "sh", + "-c", + "corepack enable pnpm && pnpm install --frozen-lockfile && " + " ".join(pack), + ], + output=output, + check=True, + ) + if get_dry_run(): + return + built = LANG_SDK_TS_EXAMPLE_PATH / "dist" / LANG_SDK_TS_BUNDLE_NAME + if not built.is_file(): + raise SystemExit(f"{built} was not produced; check the pack output above") + shutil.copyfile(built, ts_dir / LANG_SDK_TS_BUNDLE_NAME) + + +def _lang_sdk_build_ts_worker_image( + base_image: str, python: str, kubernetes_version: str, output: Output | None +) -> str: + """Build the prod+Node TypeScript worker image and load it into the kind cluster.""" + get_console(output=output).print( + f"[info]Building TypeScript worker image {LANG_SDK_TS_WORKER_IMAGE} (Node on top of {base_image})" + ) + run_command( + [ + "docker", + "build", + "--build-arg", + f"BASE_IMAGE={base_image}", + "-t", + LANG_SDK_TS_WORKER_IMAGE, + "-f", + str(LANG_SDK_TS_DOCKERFILE), + str(LANG_SDK_PATH), + ], + output=output, + check=True, + ) + cluster_name = get_kind_cluster_name(python=python, kubernetes_version=kubernetes_version) + get_console(output=output).print(f"[info]Loading {LANG_SDK_TS_WORKER_IMAGE} into {cluster_name}") + run_command_with_k8s_env( + ["kind", "load", "docker-image", "--name", cluster_name, LANG_SDK_TS_WORKER_IMAGE], + python=python, + kubernetes_version=kubernetes_version, + output=output, + check=True, + ) + return LANG_SDK_TS_WORKER_IMAGE + + def _lang_sdk_build_java_jar( staging: Path, java_sdk_source: Path, output: Output | None, *, native: bool = False ) -> None: @@ -2862,23 +2957,30 @@ def _lang_sdk_upload_artifacts( ).stdout.strip() go_bundle = staging / "go-artifacts" / "lang_sdk_combined" + ts_bundle = staging / "ts-artifacts" / LANG_SDK_TS_BUNDLE_NAME if get_dry_run(): # The dry-run build steps produce no jar; use the placeholder name so the commands still print. java_jar = staging / "java-artifacts" / "app.jar" else: java_jar = next((staging / "java-artifacts").glob("*.jar")) stub_dag = LANG_SDK_PATH / "dags" / "lang_sdk_combined.py" + # The native TypeScript Dag triggers this one, so it has to be in the + # cluster for its trigger task to have a target. It is the same Python Dag + # the mixed-language half of `ts-sdk/example` is written against. + ts_stub_dag = LANG_SDK_TS_EXAMPLE_PATH / "dags" / "typescript_example.py" for src, dest in ( (go_bundle, "/tmp/go_bundle"), (java_jar, "/tmp/app.jar"), + (ts_bundle, f"/tmp/{LANG_SDK_TS_BUNDLE_NAME}"), (stub_dag, "/tmp/lang_sdk_combined.py"), + (ts_stub_dag, "/tmp/typescript_example.py"), ): _lang_sdk_kubectl( ["cp", str(src), f"{HELM_AIRFLOW_NAMESPACE}/{pod}:{dest}"], python, kubernetes_version, output ) - for bucket in ("go-artifacts", "java-artifacts", "dags"): + for bucket in ("go-artifacts", "java-artifacts", "ts-artifacts", "dags"): _lang_sdk_kubectl( ["exec", "-n", HELM_AIRFLOW_NAMESPACE, pod, "--", "awslocal", "s3", "mb", f"s3://{bucket}"], python, @@ -2889,7 +2991,9 @@ def _lang_sdk_upload_artifacts( uploads = ( ("/tmp/go_bundle", "s3://go-artifacts/lang_sdk_combined"), ("/tmp/app.jar", "s3://java-artifacts/app.jar"), + (f"/tmp/{LANG_SDK_TS_BUNDLE_NAME}", f"s3://ts-artifacts/{LANG_SDK_TS_BUNDLE_NAME}"), ("/tmp/lang_sdk_combined.py", "s3://dags/lang_sdk_combined.py"), + ("/tmp/typescript_example.py", "s3://dags/typescript_example.py"), ) for src, dest in uploads: _lang_sdk_kubectl( @@ -2901,14 +3005,21 @@ def _lang_sdk_upload_artifacts( def _lang_sdk_apply_configmaps_and_secret( - python: str, kubernetes_version: str, go_image: str, java_image: str, output: Output | None + python: str, + kubernetes_version: str, + go_image: str, + java_image: str, + ts_image: str, + output: Output | None, ) -> None: with tempfile.TemporaryDirectory(prefix="lang_sdk_pt_") as tmp: rendered = Path(tmp) - for name in ("lang_sdk_golang.yaml", "lang_sdk_java.yaml"): + for name in ("lang_sdk_golang.yaml", "lang_sdk_java.yaml", "lang_sdk_typescript.yaml"): text = (LANG_SDK_PATH / "pod_templates" / name).read_text() - text = text.replace("__LANG_SDK_GO_IMAGE__", go_image).replace( - "__LANG_SDK_JAVA_IMAGE__", java_image + text = ( + text.replace("__LANG_SDK_GO_IMAGE__", go_image) + .replace("__LANG_SDK_JAVA_IMAGE__", java_image) + .replace("__LANG_SDK_TS_IMAGE__", ts_image) ) (rendered / name).write_text(text) # Idempotent configmap/secret application via `--dry-run | apply`. @@ -3067,20 +3178,24 @@ def _setup_lang_sdk_test( kubernetes_version: str, go_image: str | None = None, java_image: str | None = None, + ts_image: str | None = None, output: Output | None = None, ) -> None: """Provision the lang-SDK coordinator env on an already-deployed KubernetesExecutor cluster. - Resolves the go-sdk/java-sdk sources, then builds the Go/Java artifacts, the Java worker - image and deploys localstack in parallel, then serially uploads the artifacts, applies the - config + secret, and helm-upgrades Airflow with the lang-SDK values. + Resolves the go-sdk/java-sdk sources, then builds the Go/Java/TypeScript artifacts, the Java + and TypeScript worker images and deploys localstack in parallel, then serially uploads the + artifacts, applies the config + secret, and helm-upgrades Airflow with the lang-SDK values. """ go_image = go_image or f"{BuildProdParams(python=python).airflow_image_kubernetes}:latest" build_java_image = java_image is None + build_ts_image = ts_image is None if java_image is None: # The worker-image build below produces this fixed tag; resolve it up-front so the config # rendering (which needs the tag, not the build result) does not depend on the parallel run. java_image = LANG_SDK_JAVA_WORKER_IMAGE + if ts_image is None: + ts_image = LANG_SDK_TS_WORKER_IMAGE # In CI the Go/Java toolchains are provisioned + cached on the host (actions/setup-go, setup-java), # so building the artifacts natively skips the toolchain-image pulls and reuses the runner caches. native = os.environ.get("LANG_SDK_NATIVE_TOOLCHAIN", "").lower() == "true" @@ -3096,6 +3211,10 @@ def _setup_lang_sdk_test( "Build Java jar", lambda o: _lang_sdk_build_java_jar(staging, java_sdk_source, o, native=native), ), + ( + "Build TypeScript bundle", + lambda o: _lang_sdk_build_ts_bundle(staging, o, native=native), + ), ("Deploy localstack", lambda o: _lang_sdk_deploy_localstack(python, kubernetes_version, o)), ] if build_java_image: @@ -3105,17 +3224,25 @@ def _setup_lang_sdk_test( lambda o: _lang_sdk_build_java_worker_image(go_image, python, kubernetes_version, o), ) ) + if build_ts_image: + steps.append( + ( + "Build TypeScript worker image", + lambda o: _lang_sdk_build_ts_worker_image(go_image, python, kubernetes_version, o), + ) + ) _run_lang_sdk_parallel(steps, output=output) _lang_sdk_upload_artifacts(staging, python, kubernetes_version, output) - _lang_sdk_apply_configmaps_and_secret(python, kubernetes_version, go_image, java_image, output) + _lang_sdk_apply_configmaps_and_secret(python, kubernetes_version, go_image, java_image, ts_image, output) _lang_sdk_deploy_airflow(python, kubernetes_version, output) @kubernetes_group.command( name="setup-lang-sdk-test", - help="Provision the lang-SDK (Go + Java) coordinator system test on an already-deployed " - "KubernetesExecutor cluster: build artifacts, build + load the Java worker image, deploy " - "localstack S3, upload artifacts + stub Dag, create config, and upgrade the Helm release. " + help="Provision the lang-SDK (Go + Java + TypeScript) coordinator system test on an " + "already-deployed KubernetesExecutor cluster: build artifacts, build + load the Java and " + "TypeScript worker images, deploy localstack S3, upload artifacts + stub Dag, create config, " + "and upgrade the Helm release. " "Run the test afterwards with `RUN_LANG_SDK_K8S_TESTS=true breeze k8s tests " "--executor KubernetesExecutor -- -k test_lang_sdk_combined_dag_succeeds`.", ) @@ -3130,9 +3257,20 @@ def _setup_lang_sdk_test( help="Image for the Java (JavaCoordinator) worker pod. Must include a JRE. Defaults to building " "the prod image plus a headless JRE (Dockerfile.java) and loading it into the kind cluster.", ) [email protected]( + "--ts-image", + help="Image for the TypeScript (NodeCoordinator) worker pod. Must include Node. Defaults to " + "building the prod image plus Node (Dockerfile.typescript) and loading it into the kind cluster.", +) @option_verbose @option_dry_run -def setup_lang_sdk_test(python: str, kubernetes_version: str, go_image: str | None, java_image: str | None): +def setup_lang_sdk_test( + python: str, + kubernetes_version: str, + go_image: str | None, + java_image: str | None, + ts_image: str | None, +): result = sync_virtualenv(force_venv_setup=False) if result.returncode != 0: sys.exit(result.returncode) @@ -3142,6 +3280,7 @@ def setup_lang_sdk_test(python: str, kubernetes_version: str, go_image: str | No kubernetes_version=kubernetes_version, go_image=go_image, java_image=java_image, + ts_image=ts_image, output=None, ) console_print( diff --git a/dev/breeze/src/airflow_breeze/commands/kubernetes_commands_config.py b/dev/breeze/src/airflow_breeze/commands/kubernetes_commands_config.py index 919934cc754..b72db432868 100644 --- a/dev/breeze/src/airflow_breeze/commands/kubernetes_commands_config.py +++ b/dev/breeze/src/airflow_breeze/commands/kubernetes_commands_config.py @@ -56,6 +56,7 @@ KUBERNETES_PARAMETERS: dict[str, list[dict[str, str | list[str]]]] = { "--kubernetes-version", "--go-image", "--java-image", + "--ts-image", ], } ], diff --git a/kubernetes-tests/lang_sdk/Dockerfile.typescript b/kubernetes-tests/lang_sdk/Dockerfile.typescript new file mode 100644 index 00000000000..1dfda850443 --- /dev/null +++ b/kubernetes-tests/lang_sdk/Dockerfile.typescript @@ -0,0 +1,41 @@ +# 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. +# +# TypeScript worker image for the "typescript" queue: the stock prod image plus +# Node.js, so NodeCoordinator can exec `node` on the packed bundle. The Go +# queue runs on the plain prod image and Java on a JRE image, so each +# coordinator's queue gets the base image its runtime needs. +ARG BASE_IMAGE +FROM ${BASE_IMAGE} + +USER root +# The keyring step pipes curl into gpg, so a failure mid-pipe has to fail the build. +SHELL ["/bin/bash", "-o", "pipefail", "-c"] +RUN apt-get update \ + && apt-get install --no-install-recommends -y ca-certificates curl gnupg \ + && mkdir -p /etc/apt/keyrings \ + && curl -fsSL https://deb.nodesource.com/gpgkey/nodesource-repo.gpg.key \ + | gpg --dearmor -o /etc/apt/keyrings/nodesource.gpg \ + && echo "deb [signed-by=/etc/apt/keyrings/nodesource.gpg] https://deb.nodesource.com/node_22.x nodistro main" \ + > /etc/apt/sources.list.d/nodesource.list \ + && apt-get update \ + && apt-get install --no-install-recommends -y nodejs \ + && apt-get purge -y --auto-remove gnupg \ + && apt-get clean \ + && rm -rf /var/lib/apt/lists/* +SHELL ["/bin/sh", "-c"] +USER airflow diff --git a/kubernetes-tests/lang_sdk/README.md b/kubernetes-tests/lang_sdk/README.md index 93583ecbd8b..a9f032a76ed 100644 --- a/kubernetes-tests/lang_sdk/README.md +++ b/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 +`RUN_TS_SDK_NATIVE_DAG_K8S_TESTS` as well as `RUN_LANG_SDK_K8S_TESTS`, because dispatching a parse +request to a language coordinator is not in Airflow yet — the provisioning below is in place, so the +test becomes real as soon as it is. + ## How it fits together ``` @@ -57,11 +63,18 @@ coordinator scans. | `stage_artifacts.py` | Init-container entrypoint; stages an artifact bucket via DagBundle. | | `pod_templates/lang_sdk_golang.yaml` | `golang` queue worker pod: prod image + go-artifacts init container. | | `pod_templates/lang_sdk_java.yaml` | `java` queue worker pod: JVM image + java-artifacts init container. | +| `pod_templates/lang_sdk_typescript.yaml` | `typescript` queue worker pod: Node image + ts-artifacts init container. | +| `Dockerfile.typescript` | Prod image plus Node.js, which `NodeCoordinator` execs. | | `manifests/localstack.yaml` | In-cluster S3 (localstack). | | `config/values.yaml` | Helm overrides: KubernetesExecutor, coordinators (+extra.pod_template_file), queue routing, stub-Dag S3 bundle, AWS conn, scheduler pod-template mount. | -The Go binary, Java jar, and stub Dag share one object store (localstack) but live in -**separate buckets** (`go-artifacts`, `java-artifacts`, `dags`). +The Go binary, Java jar, TypeScript bundle and the Dag files share one object store (localstack) +but live in **separate buckets** (`go-artifacts`, `java-artifacts`, `ts-artifacts`, `dags`). + +The TypeScript bundle is `ts-sdk/example`, packed by `airflow-ts-pack`. It carries both halves of +that example: the handlers for the Python-declared `typescript_example` Dag, and the natively +declared `typescript_native_example`. `typescript_example.py` is uploaded to the `dags` bucket too, +because the native Dag's trigger task targets it. ## Which SDK sources get built diff --git a/kubernetes-tests/lang_sdk/config/values.yaml b/kubernetes-tests/lang_sdk/config/values.yaml index a4b5b4004b4..8d782c478f3 100644 --- a/kubernetes-tests/lang_sdk/config/values.yaml +++ b/kubernetes-tests/lang_sdk/config/values.yaml @@ -62,13 +62,13 @@ dagProcessor: aws_conn_id: "aws_localstack" config: - # Route the golang/java queues to their coordinators, each carrying an - # extra.pod_template_file (the worktree-1 feature) so the worker pod is - # launched from a coordinator-specific template that stages the artifact and, - # for Java, runs on the JVM-bearing image. + # Route the golang/java/typescript queues to their coordinators, each carrying + # an extra.pod_template_file so the worker pod is launched from a + # coordinator-specific template that stages the artifact and, for Java and + # TypeScript, runs on the image carrying that runtime. sdk: # Single-line JSON: airflow.cfg is an ini file, so a value with embedded # newlines would break the rendered configmap. # yamllint disable-line rule:line-length - coordinators: '{"go-sdk": {"classpath": "airflow.sdk.coordinators.executable.ExecutableCoordinator", "kwargs": {"executables_root": ["/opt/airflow/artifacts/go-artifacts"]}, "extra": {"pod_template_file": "/opt/airflow/lang_sdk/pod_templates/lang_sdk_golang.yaml"}}, "java-sdk": {"classpath": "airflow.sdk.coordinators.java.JavaCoordinator", "kwargs": {"jars_root": ["/opt/airflow/artifacts/java-artifacts"]}, "extra": {"pod_template_file": "/opt/airflow/lang_sdk/pod_templates/lang_sdk_j [...] - queue_to_coordinator: '{"golang": "go-sdk", "java": "java-sdk"}' + coordinators: '{"go-sdk": {"classpath": "airflow.sdk.coordinators.executable.ExecutableCoordinator", "kwargs": {"executables_root": ["/opt/airflow/artifacts/go-artifacts"]}, "extra": {"pod_template_file": "/opt/airflow/lang_sdk/pod_templates/lang_sdk_golang.yaml"}}, "java-sdk": {"classpath": "airflow.sdk.coordinators.java.JavaCoordinator", "kwargs": {"jars_root": ["/opt/airflow/artifacts/java-artifacts"]}, "extra": {"pod_template_file": "/opt/airflow/lang_sdk/pod_templates/lang_sdk_j [...] + queue_to_coordinator: '{"golang": "go-sdk", "java": "java-sdk", "typescript": "ts-sdk"}' diff --git a/kubernetes-tests/lang_sdk/pod_templates/lang_sdk_typescript.yaml b/kubernetes-tests/lang_sdk/pod_templates/lang_sdk_typescript.yaml new file mode 100644 index 00000000000..1a316f2e7e7 --- /dev/null +++ b/kubernetes-tests/lang_sdk/pod_templates/lang_sdk_typescript.yaml @@ -0,0 +1,115 @@ +# 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. +# +# Worker pod template for the "typescript" queue (NodeCoordinator). +# +# An init container stages the packed ``bundle.min.mjs`` from the localstack +# ``ts-artifacts`` bucket into a shared emptyDir at /opt/airflow/artifacts via +# the DagBundle interface (stage_artifacts.py). That path is the coordinator's +# bundles_root. The base image carries Node.js (Dockerfile.typescript), which the plain +# prod image does not. Image placeholders are rendered by the +# ``breeze k8s setup-lang-sdk-test`` command. +--- +kind: Pod +apiVersion: v1 +metadata: + name: placeholder-name-dont-delete + namespace: placeholder-name-dont-delete +spec: + initContainers: + - name: stage-ts-bundle + image: __LANG_SDK_TS_IMAGE__ + imagePullPolicy: IfNotPresent + command: ["python", "/opt/airflow/lang_sdk/stage_artifacts.py"] + env: + - name: AIRFLOW_HOME + value: /opt/airflow + - name: ARTIFACT_BUNDLE_NAME + value: ts-artifacts + - name: AIRFLOW__DAG_PROCESSOR__DAG_BUNDLE_STORAGE_PATH + value: /opt/airflow/artifacts + - name: AIRFLOW__DAG_PROCESSOR__DAG_BUNDLE_CONFIG_LIST + value: >- + [{"name": "ts-artifacts", + "classpath": "airflow.providers.amazon.aws.bundles.s3.S3DagBundle", + "kwargs": {"bucket_name": "ts-artifacts", "aws_conn_id": "aws_localstack"}}] + - name: AIRFLOW_CONN_AWS_LOCALSTACK + valueFrom: + secretKeyRef: + name: lang-sdk-aws-conn + key: uri + volumeMounts: + - name: lang-sdk-artifacts + mountPath: /opt/airflow/artifacts + - name: lang-sdk-scripts + readOnly: true + mountPath: /opt/airflow/lang_sdk + containers: + - name: base + image: __LANG_SDK_TS_IMAGE__ + env: + - name: AIRFLOW__CORE__EXECUTOR + value: LocalExecutor + - name: AIRFLOW_HOME + value: /opt/airflow + - name: AIRFLOW__CORE__DAGS_FOLDER + value: /opt/airflow/dags + - name: AIRFLOW__CORE__FERNET_KEY + valueFrom: + secretKeyRef: + name: airflow-fernet-key + key: fernet-key + - name: AIRFLOW__DATABASE__SQL_ALCHEMY_CONN + valueFrom: + secretKeyRef: + name: airflow-metadata + key: connection + volumeMounts: + - name: airflow-logs + mountPath: /opt/airflow/logs + - name: lang-sdk-artifacts + mountPath: /opt/airflow/artifacts + - name: airflow-config + readOnly: true + mountPath: /opt/airflow/airflow.cfg + subPath: airflow.cfg + - name: airflow-config + readOnly: true + mountPath: /opt/airflow/config/airflow_local_settings.py + subPath: airflow_local_settings.py + imagePullPolicy: IfNotPresent + volumes: + - name: airflow-logs + emptyDir: {} + - name: lang-sdk-artifacts + emptyDir: {} + - name: lang-sdk-scripts + configMap: + name: lang-sdk-scripts + defaultMode: 420 + - name: airflow-config + configMap: + name: airflow-config + defaultMode: 420 + restartPolicy: Never + terminationGracePeriodSeconds: 30 + serviceAccountName: airflow-worker-kubernetes + serviceAccount: airflow-worker-kubernetes + securityContext: + runAsUser: 50000 + fsGroup: 50000 + schedulerName: default-scheduler diff --git a/kubernetes-tests/tests/kubernetes_tests/test_lang_sdk_coordinator_executor.py b/kubernetes-tests/tests/kubernetes_tests/test_lang_sdk_coordinator_executor.py index 813329d4b8a..7bde81fbf32 100644 --- a/kubernetes-tests/tests/kubernetes_tests/test_lang_sdk_coordinator_executor.py +++ b/kubernetes-tests/tests/kubernetes_tests/test_lang_sdk_coordinator_executor.py @@ -38,6 +38,20 @@ from kubernetes_tests.test_base import EXECUTOR, BaseK8STest _RUN_LANG_SDK = os.environ.get("RUN_LANG_SDK_K8S_TESTS", "").lower() in ("true", "1") DAG_ID = "lang_sdk_combined" +# Declared entirely in TypeScript: the Dag processor asks the bundle to parse +# itself, so no Python file in the bundles folder mentions it. +NATIVE_TS_DAG_ID = "typescript_native_example" +NATIVE_TS_TASK_IDS = [ + "extract.north", + "extract.south", + "summarize", + "has_rows", + "load_rows", + "pick_cadence", + "publish_weekly", + "cleanup", + "trigger_downstream", +] TASK_IDS = [ "python_task_1", "go_extract", @@ -50,6 +64,10 @@ TASK_IDS = [ # artifact + start a coordinator subprocess, so allow generous headroom. _TIMEOUT = 600 +# The native Dag additionally needs a Dag processor that dispatches a parse +# request to the Node coordinator; until that lands the Dag never appears. +_RUN_NATIVE_TS = os.environ.get("RUN_TS_SDK_NATIVE_DAG_K8S_TESTS", "").lower() in ("true", "1") + @pytest.mark.skipif( EXECUTOR != "KubernetesExecutor" or not _RUN_LANG_SDK, @@ -86,3 +104,65 @@ class TestLangSdkCoordinatorExecutor(BaseK8STest): expected_final_state="success", timeout=_TIMEOUT, ) + + [email protected]( + EXECUTOR != "KubernetesExecutor" or not _RUN_NATIVE_TS, + reason="Runs only on KubernetesExecutor with a Dag processor that parses Lang SDK bundles " + "(RUN_TS_SDK_NATIVE_DAG_K8S_TESTS)", +) +class TestNativeTypeScriptDagOnKubernetes(BaseK8STest): + """The same Dag the SDK and Airflow e2e suites run, under the Node coordinator. + + Covers the path a Dag with no Python file takes: the Dag processor parses the + ``airflow-ts-pack`` bundle through the coordinator, the scheduler reads the serialized Dag, + and each task runs in its own pod. The graph is a real one -- a task group, a named fan-in, + order-only edges, a conditional and a multi-way branch -- so pod-per-task scheduling is + exercised against branches that skip rather than a linear chain. + """ + + def _ensure_variable(self, key: str, value: str) -> None: + resp = self.session.post(f"http://{self.host}/variables", json={"key": key, "value": value}) + assert resp.status_code in (200, 201, 409), f"Could not create variable {key}: {resp.text}" + + @pytest.mark.execution_timeout(1200) + def test_native_typescript_dag_succeeds(self): + # Both regions non-empty, so the conditional takes its `then` branch and + # `report_empty` is the side that skips. The cadence names the case the + # multi-way branch chooses, so `publish_daily` is the side that skips. + self._ensure_variable("typescript_native_north_rows", "3") + self._ensure_variable("typescript_native_south_rows", "2") + self._ensure_variable("typescript_native_cadence", "weekly") + + dag_run_id, logical_date = self.start_job_in_kubernetes(NATIVE_TS_DAG_ID, self.host) + print(f"Triggered {NATIVE_TS_DAG_ID} run {dag_run_id} (logical_date={logical_date})") + + for task_id in NATIVE_TS_TASK_IDS: + self.monitor_task( + host=self.host, + dag_run_id=dag_run_id, + dag_id=NATIVE_TS_DAG_ID, + task_id=task_id, + expected_final_state="success", + timeout=_TIMEOUT, + ) + + # The branch that was not taken is skipped, not failed, which is what + # proves the skip reached the supervisor rather than the task erroring. + for skipped in ("report_empty", "publish_daily"): + self.monitor_task( + host=self.host, + dag_run_id=dag_run_id, + dag_id=NATIVE_TS_DAG_ID, + task_id=skipped, + expected_final_state="skipped", + timeout=_TIMEOUT, + ) + + self.ensure_dag_expected_state( + host=self.host, + logical_date=logical_date, + dag_id=NATIVE_TS_DAG_ID, + expected_final_state="success", + timeout=_TIMEOUT, + ) diff --git a/ts-sdk/adr/0002-native-dag-interface.md b/ts-sdk/adr/0002-native-dag-interface.md index 7c08f4f6009..94c798bcb4d 100644 --- a/ts-sdk/adr/0002-native-dag-interface.md +++ b/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. ## Decision @@ -206,13 +206,17 @@ Dag is read. - One authoring surface (`dag.task()` plus its factory) covers the graph and each task's arguments, and `before`/`after` cover edges that carry nothing. - Handlers are unit-testable as plain functions of their data, with no SDK fixture to construct. -- `DagSpec` and `TaskSpec` are empty placeholders today — `Record<string, never>` - (`ts-sdk/src/sdk/dag.ts`), so `new Dag("d", { schedule: "@daily" })` is currently a compile error - by design. Native declaration is what fills them, generated from the serialized-Dag JSON schema the - way `src/generated/supervisor.ts` is. This ADR does not choose those fields; it fixes where an - author writes them. +- `DagSpec` and `TaskSpec` carry the fields of Airflow's own serialized-Dag JSON schema, generated + from it the way `src/generated/supervisor.ts` is, so every field is optional and a misspelled one is + a compile error. This ADR does not choose those fields; it fixes where an author writes them. `queue` + is the one hand-written Dag field: the schema has no Dag-level queue, but every task of a native Dag + runs on the same coordinator, so the queue that routes them there belongs on the Dag. - `TaskOptions` carries the task's spec and nothing else: the names on the wire are the keys of the - call itself. With wiring moved to the factory call, `inputs` is no longer an option. + call itself. With wiring moved to the factory call, `inputs` is no longer an option, so the trailing + argument is the spec's own fields. +- A native Dag's tasks serialize as stub tasks (`is_stub`, with `_arg_bindings`), which is how the API + server resolves each argument per task instance and hands it to a foreign runtime. So the wiring an + author writes is the same mechanism a `@task.stub` call already uses, rather than a second one. - `TaskHandlerArgs` is removed from the public API, `DagRegistry` becomes `Bundle`, and `serveDags(registry)` becomes `bundle.serve()`, which breaks 0.1.0-beta1 authors; see [ADR-0001](0001-mixed-lang-dag-interface.md) for the shipped call sites diff --git a/ts-sdk/example/README.md b/ts-sdk/example/README.md index 94c2ead3951..a51e089cdb5 100644 --- a/ts-sdk/example/README.md +++ b/ts-sdk/example/README.md @@ -25,7 +25,13 @@ 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 + Python file declares its graph. Its schedule, task options and edges are all written on this side, + and the bundle answers the Dag processor's parse request with the serialized Dag. 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 `triggerDagRun` task a Python worker executes. +- `dist/bundle.min.mjs` is the generated Node.js bundle that Airflow launches. One artifact serves both + authoring modes. The build uses the SDK's `airflow-ts-pack` tool, which bundles the entrypoint with esbuild and embeds the Airflow metadata generated from the bundle's diff --git a/ts-sdk/example/src/main.ts b/ts-sdk/example/src/main.ts index fa97d3f9899..14e68ab2b9e 100644 --- a/ts-sdk/example/src/main.ts +++ b/ts-sdk/example/src/main.ts @@ -17,13 +17,15 @@ * under the License. */ -// The bundle entry point: one bundle, two Python-owned Dags. +// The bundle entry point: one bundle, three Dags, both authoring modes. // -// Both Dags in `dags/` are declared in Python with `@task.stub` tasks routed to the Node -// coordinator, so this side only supplies the task bodies. +// The two Dags in `dags/` are declared in Python with `@task.stub` tasks routed to the Node +// coordinator, so this side only supplies their task bodies. The third, in `native.ts`, is +// declared in TypeScript outright — graph, schedule and all. import { Bundle, getClient, getContext, TaskHandler } from "apache-airflow-ts-sdk"; +import { dag as nativeDag } from "./native.js"; import { buildSummaryMessage, report, summarize } from "./taskflow.js"; export async function buildMessage() { @@ -72,8 +74,12 @@ export async function readConnection() { // One register call lists everything this bundle provides. // `build_message` appears under both Dags: two different handlers, told apart by the dag_id each // is bound to and never by the task_id alone. +// +// One bundle serves both authoring modes: handlers for the Dags a Python file declares, and the +// natively declared Dag from `native.ts`, whose graph this side owns outright. const bundle = new Bundle(); bundle.register( + nativeDag, new TaskHandler("typescript_example", "build_message", buildMessage), new TaskHandler("typescript_example", "read_connection", readConnection), new TaskHandler("typescript_example", "write_and_delete_variable", writeAndDeleteVariable), diff --git a/ts-sdk/example/src/native.ts b/ts-sdk/example/src/native.ts new file mode 100644 index 00000000000..64b8efb2192 --- /dev/null +++ b/ts-sdk/example/src/native.ts @@ -0,0 +1,163 @@ +/*! + * 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. + */ + +// A Dag declared entirely in TypeScript: no Python file declares its graph. +// +// The counterpart of `main.ts`, which supplies handlers for a Dag that a Python +// file declares. Here the schedule, the tasks, their options and the edges +// between them are all written on this side, and the bundle answers the Dag +// processor's parse request with the serialized Dag. +// +// The graph is a graph rather than a chain, so the e2e suites exercise the +// constructs against each other: 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. +// +// `main.ts` registers this Dag on the same bundle as the mixed-language +// handlers, so one artifact covers both authoring modes. + +import { Dag, getClient } from "apache-airflow-ts-sdk"; + +/** The regions this example works over. */ +const REGIONS = ["north", "south"] as const; + +interface RegionRows { + readonly region: string; + readonly rows: number; +} + +export const dag = new Dag("typescript_native_example", { + schedule: "@daily", + catchup: false, + tags: ["typescript", "native"], + description: "A Dag whose graph, schedule and task options are all declared in TypeScript.", + // Routes every task of this Dag to the Node coordinator; see the + // `queue_to_coordinator` entry in the TypeScript SDK docs. + queue: "typescript", +}); + +// --- Extraction, inside a task group ------------------------------------- +// +// A group prefixes the ids of everything in it, so these become +// `extract.north` and `extract.south`. + +const extract = dag.taskGroup("extract"); + +const extractNorth = extract.task("north", async (): Promise<RegionRows> => { + const configured = await getClient().getVariable("typescript_native_north_rows"); + return { region: "north", rows: Number(configured ?? 3) }; +}); + +const extractSouth = extract.task("south", async (): Promise<RegionRows> => { + const configured = await getClient().getVariable("typescript_native_south_rows"); + return { region: "south", rows: Number(configured ?? 0) }; +}); + +// --- A named fan-in ------------------------------------------------------- +// +// Two upstreams feed one task under the argument names its handler declares, +// which is what the wiring object is for. + +const summarize = dag.task( + "summarize", + async ({ north, south }: { north: RegionRows; south: RegionRows }) => { + const total = north.rows + south.rows; + await getClient().setXCom({ key: "region_total", value: total }); + return { total, regions: REGIONS.length }; + }, +); + +const summarized = summarize({ north: extractNorth(), south: extractSouth() }); + +// --- A conditional -------------------------------------------------------- +// +// The guarded tasks take no parameter for the control edge: the condition's +// boolean is a run-time signal rather than data, so it is never an argument. + +const loadRows = dag.task("load_rows", async () => { + const total = await getClient().getXCom<number>({ key: "region_total", taskId: "summarize" }); + return { loaded: total ?? 0 }; +}); + +const reportEmpty = dag.task("report_empty", async () => ({ loaded: 0 })); + +const loaded = loadRows(); +const reportedEmpty = reportEmpty(); + +// The upstream is an input like any other, named after the argument it +// supplies. +const hasRows = dag.task( + "has_rows", + async ({ summary }: { summary: { total: number } }) => summary.total > 0, +); +const gated = hasRows({ summary: summarized }); + +dag.if(gated).then(loaded).else(reportedEmpty); + +// --- A multi-way branch --------------------------------------------------- +// +// A case is the reference the SDK handed back, so renaming a handler cannot +// silently rewire the Dag, and the compiler checks the candidate exists. + +const publishDaily = dag.task("publish_daily", async () => ({ cadence: "daily" })); +const publishWeekly = dag.task("publish_weekly", async () => ({ cadence: "weekly" })); + +const daily = publishDaily(); +const weekly = publishWeekly(); + +// Decided from a Variable rather than from an upstream value, so the decider +// takes no argument and the extract group is what orders it. +const pickCadence = dag.task("pick_cadence", async () => { + const cadence = await getClient().getVariable("typescript_native_cadence"); + return cadence === "weekly" ? weekly : daily; +}); +const picked = pickCadence(); + +dag.switch(picked).case(daily).case(weekly); + +// --- Order-only edges ----------------------------------------------------- +// +// `cleanup` follows whichever side of each branch ran, but takes nothing from +// it, so the edges are drawn between the references rather than through an +// argument. Its `all_done` rule is what lets it run when a branch skipped the +// other side. +// +// `extract.before(picked)` is the same kind of edge with a group at one end: +// the decider reads no upstream value, so this is what puts it after every +// task the group holds. + +const cleanup = dag.task("cleanup", async () => undefined, { triggerRule: "all_done" }); +const cleaned = cleanup(); + +cleaned.after(loaded, reportedEmpty, daily, weekly); +extract.before(picked); + +// --- A task Python runs --------------------------------------------------- +// +// A trigger wraps no TypeScript function: it serializes as +// `TriggerDagRunOperator` and a Python worker executes it, so it inherits no +// queue from this Dag and needs nothing to call. + +dag + .triggerDagRun({ + taskId: "trigger_downstream", + dagId: "typescript_example", + conf: { triggered_by: "{{ dag.dag_id }}" }, + }) + .after(cleaned); diff --git a/ts-sdk/tests/sdk/native-dag-example.test.ts b/ts-sdk/tests/sdk/native-dag-example.test.ts new file mode 100644 index 00000000000..72b0353eb83 --- /dev/null +++ b/ts-sdk/tests/sdk/native-dag-example.test.ts @@ -0,0 +1,207 @@ +/*! + * 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. + */ + +// The shipped native Dag example, built and serialized. +// +// The unit suites cover each construct on a Dag written for that construct +// alone; this reads the example an author is pointed at, so the graph the docs +// describe is the graph the SDK produces. The Airflow and Kubernetes suites +// run the same Dag end to end. + +import { describe, expect, it } from "vitest"; +import { serializeDag } from "../../src/coordinator/serde.js"; +import { Bundle, listBundleNativeDags } from "../../src/sdk/bundle.js"; +import { finalizeDag, getDagOrderEdges, type Dag } from "../../src/sdk/dag.js"; +import { dag } from "../../example/src/native.js"; + +type Json = Record<string, unknown>; + +const FILELOC = "/bundles/example/bundle.min.mjs"; + +/** The serialized Dag, and its tasks keyed by task id. */ +function serialized(target: Dag) { + const json = serializeDag(target, FILELOC, "bundle.min.mjs") as Json; + const tasks = json["tasks"] as { __var: Json }[]; + return { json, tasks: new Map(tasks.map(({ __var }) => [__var["task_id"] as string, __var])) }; +} + +describe("the native Dag example", () => { + it("declares its schedule, options and queue in TypeScript", () => { + expect(dag.dagId).toBe("typescript_native_example"); + expect(dag.spec).toMatchObject({ + schedule: "@daily", + catchup: false, + tags: ["typescript", "native"], + queue: "typescript", + }); + }); + + it("builds a graph, not a chain", () => { + expect([...dag.taskIds].sort()).toEqual([ + "cleanup", + "extract.north", + "extract.south", + "has_rows", + "load_rows", + "pick_cadence", + "publish_daily", + "publish_weekly", + "report_empty", + "summarize", + "trigger_downstream", + ]); + }); + + it("is fully laid out, so reading it raises nothing", () => { + expect(() => finalizeDag(dag)).not.toThrow(); + }); + + it("is served by the same bundle as the mixed-language handlers", () => { + // What `main.ts` builds: importing it would start the runtime, so the + // registration is rebuilt here rather than imported. + const bundle = new Bundle(dag); + + expect(listBundleNativeDags(bundle).map((served) => served.dagId)).toEqual([ + "typescript_native_example", + ]); + }); + + describe("serializes", () => { + it("a named fan-in as an edge from each upstream into one task", () => { + const { tasks } = serialized(dag); + + expect(tasks.get("extract.north")?.["downstream_task_ids"]).toContain("summarize"); + expect(tasks.get("extract.south")?.["downstream_task_ids"]).toContain("summarize"); + expect(tasks.get("summarize")?.["downstream_task_ids"]).toEqual(["has_rows"]); + }); + + it("the fan-in's arguments as the bindings the API server resolves", () => { + expect(serialized(dag).tasks.get("summarize")?.["_arg_bindings"]).toEqual([ + { name: "north", kind: "xcom", task_id: "extract.north" }, + { name: "south", kind: "xcom", task_id: "extract.south" }, + ]); + }); + + it("the positional call's argument as one binding on the upstream", () => { + // The label is whatever named the parameter, which only the packer can + // read; the order and the upstream are what bind. + expect(serialized(dag).tasks.get("has_rows")?.["_arg_bindings"]).toMatchObject([ + { kind: "xcom", task_id: "summarize" }, + ]); + }); + + it("the task group, with its tasks prefixed and nested under it", () => { + const root = serialized(dag).json["task_group"] as Json; + const [kind, group] = (root["children"] as Json)["extract"] as [string, Json]; + + expect(kind).toBe("taskgroup"); + expect(group).toMatchObject({ + _group_id: "extract", + children: { + "extract.north": ["operator", "extract.north"], + "extract.south": ["operator", "extract.south"], + }, + }); + }); + + it("the group edge, expanded onto the tasks the group leaves from", () => { + // `extract.before(picked)` is drawn against the group, and reaches the + // task graph through the group's leaves — both extract tasks, since + // neither runs after the other inside it. + expect(getDagOrderEdges(dag)).toContainEqual({ + upstream: "extract", + downstream: "pick_cadence", + }); + + const { tasks } = serialized(dag); + expect(tasks.get("extract.north")?.["downstream_task_ids"]).toContain("pick_cadence"); + expect(tasks.get("extract.south")?.["downstream_task_ids"]).toContain("pick_cadence"); + }); + + it("the conditional's control edges as ordinary edges, and its skip marker", () => { + // Nothing names the branch in the Dag JSON: the decision is a run-time + // skip, and `_can_skip_downstream` is what makes a cleared branch stay + // skipped. + const { tasks } = serialized(dag); + const hasRows = tasks.get("has_rows")!; + + expect(hasRows["downstream_task_ids"]).toEqual(["load_rows", "report_empty"]); + expect(hasRows["_can_skip_downstream"]).toBe(true); + expect(tasks.get("load_rows")?.["downstream_task_ids"]).toEqual(["cleanup"]); + }); + + it("the multi-way branch's cases as ordinary edges, and its skip marker", () => { + const pick = serialized(dag).tasks.get("pick_cadence")!; + + expect(pick["downstream_task_ids"]).toEqual(["publish_daily", "publish_weekly"]); + expect(pick["_can_skip_downstream"]).toBe(true); + }); + + it("cleanup behind every branch outcome, which is what its all_done rule is for", () => { + const { tasks } = serialized(dag); + + for (const outcome of ["load_rows", "report_empty", "publish_daily", "publish_weekly"]) { + expect(tasks.get(outcome)?.["downstream_task_ids"]).toContain("cleanup"); + } + expect(tasks.get("cleanup")?.["trigger_rule"]).toBe("all_done"); + expect(tasks.get("cleanup")?.["downstream_task_ids"]).toEqual(["trigger_downstream"]); + }); + + it("the trigger task as the Python operator a Python worker runs", () => { + const trigger = serialized(dag).tasks.get("trigger_downstream")!; + + expect(trigger).toMatchObject({ + task_type: "TriggerDagRunOperator", + _task_module: "airflow.providers.standard.operators.trigger_dagrun", + trigger_dag_id: "typescript_example", + }); + // No TypeScript marker, and no queue of this Dag's, so a Python worker + // can pick it up. + expect(trigger).not.toHaveProperty("language"); + expect(trigger).not.toHaveProperty("queue"); + }); + + it("the triggered Dag as a dependency of this one", () => { + expect(serialized(dag).json["dag_dependencies"]).toEqual([ + { + source: "typescript_native_example", + target: "typescript_example", + label: "trigger_downstream", + dependency_type: "trigger", + dependency_id: "trigger_downstream", + }, + ]); + }); + + it("every other task as one this runtime executes, on the Dag's queue", () => { + const { tasks } = serialized(dag); + const typescript = [...tasks.values()].filter((task) => task["language"] === "typescript"); + + expect(typescript).toHaveLength(dag.taskIds.length - 1); + expect(typescript.every((task) => task["queue"] === "typescript")).toBe(true); + }); + + it("the schedule as the expanded cron the scheduler rebuilds", () => { + expect(serialized(dag).json["timetable"]).toMatchObject({ + __type: "airflow.timetables.trigger.CronTriggerTimetable", + __var: { expression: "0 0 * * *" }, + }); + }); + }); +}); diff --git a/ts-sdk/tsconfig.json b/ts-sdk/tsconfig.json index a47f167e378..14bba01027d 100644 --- a/ts-sdk/tsconfig.json +++ b/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": { + "apache-airflow-ts-sdk": ["./src/index.ts"] + }, + "noEmit": true, "declaration": true, "declarationMap": true, "sourceMap": true }, - "include": ["src/**/*.ts", "tests/**/*.ts", "scripts/**/*.mjs"] + "include": ["src/**/*.ts", "tests/**/*.ts", "scripts/**/*.mjs", "example/src/**/*.ts"] } diff --git a/ts-sdk/vitest.config.ts b/ts-sdk/vitest.config.ts index 723796e6d29..c1ea185f655 100644 --- a/ts-sdk/vitest.config.ts +++ b/ts-sdk/vitest.config.ts @@ -17,9 +17,20 @@ * under the License. */ +import { fileURLToPath } from "node:url"; + import { defineConfig } from "vitest/config"; export default defineConfig({ + resolve: { + alias: { + // `example/` imports the package by name, which resolves to the built + // `dist/`. A test that reaches into both would hold two copies of the + // SDK, and a Dag from one is not an instance of the other's class. + // Tests run against the sources, so the example does too. + "apache-airflow-ts-sdk": fileURLToPath(new URL("./src/index.ts", import.meta.url)), + }, + }, test: { include: ["tests/**/*.test.ts", "scripts/**/*.test.mjs"], },
