This is an automated email from the ASF dual-hosted git repository.
jason810496 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new a91216d5ea9 TS SDK: serialize a native Dag to Dag JSON (#73441)
a91216d5ea9 is described below
commit a91216d5ea9fe03819a60b9d5b5b3b82390de863
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Thu Oct 1 13:16:41 2026 +0800
TS SDK: serialize a native Dag to Dag JSON (#73441)
---
.pre-commit-config.yaml | 15 +
.../language-sdks/typescript.rst | 23 +
scripts/ci/lang_sdk_serialization/__init__.py | 16 +
scripts/ci/lang_sdk_serialization/compare.py | 233 ++++++
.../ci/lang_sdk_serialization/serialize_python.py | 140 ++++
scripts/ci/lang_sdk_serialization/test_dags.yaml | 166 ++++
.../tests/ci/lang_sdk_serialization/__init__.py | 16 +
.../ci/lang_sdk_serialization/test_compare.py | 245 ++++++
ts-sdk/.pre-commit-config.yaml | 14 +
ts-sdk/package.json | 3 +-
ts-sdk/pnpm-lock.yaml | 27 +-
.../ci/prek/check_serialization_conformance.py | 44 ++
ts-sdk/src/coordinator/serde.ts | 751 ++++++++++++++++++
ts-sdk/src/sdk/dag.ts | 25 +-
ts-sdk/tests/conformance/serialize_typescript.ts | 163 ++++
ts-sdk/tests/coordinator/arg-binding.test.ts | 19 +
ts-sdk/tests/coordinator/serde.test.ts | 846 +++++++++++++++++++++
ts-sdk/tests/public-api.test.ts | 4 +
18 files changed, 2734 insertions(+), 16 deletions(-)
diff --git a/.pre-commit-config.yaml b/.pre-commit-config.yaml
index 3b6fe5200e1..9a38259abcf 100644
--- a/.pre-commit-config.yaml
+++ b/.pre-commit-config.yaml
@@ -365,6 +365,21 @@ repos:
^scripts/ci/prek/check_go_sdk_spec_drift\.py$
pass_filenames: false
require_serial: true
+ # check-ts-sdk-serialization-conformance in ts-sdk/ covers the SDK's own
files; this runs the same
+ # check when the shared harness or Airflow's serializer changes, which
that project cannot see.
+ - id: check-ts-sdk-serialization-conformance-shared
+ name: Check the TS SDK serializes Dags the way Airflow does, after a
shared change
+ description: "Serialize the shared test Dags with the TS SDK and with
Airflow, and compare the two"
+ entry: ./ts-sdk/scripts/ci/prek/check_serialization_conformance.py
+ language: node
+ additional_dependencies: ['[email protected]']
+ pass_filenames: false
+ require_serial: true
+ files: >
+ (?x)
+ ^airflow-core/src/airflow/serialization/serialized_objects\.py$|
+ ^airflow-core/src/airflow/serialization/schema\.json$|
+ ^scripts/ci/lang_sdk_serialization/.*$
- id: check-go-version-in-sync
name: Check Go toolchain version is consistent across build files
entry: ./scripts/ci/prek/check_go_version_in_sync.py
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 aea8077c1ee..47070c581da 100644
--- a/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
+++ b/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
@@ -336,6 +336,29 @@ Pass ``{ prefixGroupId: false }`` to keep the ids declared
in a group as written
does in Python; they then have to be unique across the Dag. A group id is made
of letters, digits, dashes and
underscores, and is at most 200 characters.
+Serialization
+~~~~~~~~~~~~~
+
+A native Dag serializes into the same Dag JSON a Python Dag produces, so the
scheduler reads it
+without knowing which language declared it.
+
+``schedule`` accepts what maps to a stock timetable: unset, ``@once``,
``@continuous``, or a cron
+expression. A cron preset such as ``@daily`` is recorded as the expression it
stands for. Anything
+else names a Python object a TypeScript bundle cannot point at, and is
rejected.
+
+Every task of a native Dag runs on the Node coordinator, so it needs the queue
the deployment routes
+there. Set it once on the Dag and each task inherits it:
+
+.. code-block:: typescript
+
+ const dag = new Dag("ts_etl", { schedule: "@daily", queue: "typescript" });
+
+ // ...and one task that needs its own.
+ dag.task("heavy", heavyHandler, { queue: "typescript_large" })();
+
+``queue`` on a task wins over the Dag's. See
:ref:`typescript-sdk/coordinator-config` for the
+``queue_to_coordinator`` entry that sends that queue to the coordinator.
+
``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/scripts/ci/lang_sdk_serialization/__init__.py
b/scripts/ci/lang_sdk_serialization/__init__.py
new file mode 100644
index 00000000000..13a83393a91
--- /dev/null
+++ b/scripts/ci/lang_sdk_serialization/__init__.py
@@ -0,0 +1,16 @@
+# 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.
diff --git a/scripts/ci/lang_sdk_serialization/compare.py
b/scripts/ci/lang_sdk_serialization/compare.py
new file mode 100644
index 00000000000..35f488a5c35
--- /dev/null
+++ b/scripts/ci/lang_sdk_serialization/compare.py
@@ -0,0 +1,233 @@
+# 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.
+r"""
+Check that a language SDK serializes Dags the way Airflow does.
+
+Runs the SDK's serializer and serialize_python.py over test_dags.yaml, each
writing its serialization
+to a JSON file in a temporary directory. The Python side also takes the SDK's
output as Airflow
+receives it, filling in the fields the SDK leaves to Airflow's config, and
loads it through Airflow's
+deserializer. The two are then compared field by field, so this fails both
when the serializers drift
+apart and when Airflow cannot read what the SDK writes.
+
+The SDK's serializer is the command after ``--``. It runs from the repository
root, with the paths of
+test_dags.yaml and of the JSON file to write appended, and writes each Dag as
+``DagSerialization.to_dict`` would, keyed by Dag id. The Python side needs
``uv``::
+
+ python3 scripts/ci/lang_sdk_serialization/compare.py --sdk typescript -- \
+ pnpm --dir ts-sdk exec tsx tests/conformance/serialize_typescript.ts
+
+On a failure the files are kept, and their directory is printed.
+"""
+
+from __future__ import annotations
+
+import argparse
+import json
+import shutil
+import subprocess
+import sys
+import tempfile
+from pathlib import Path
+from typing import Any
+
+HERE = Path(__file__).resolve().parent
+REPO_ROOT = HERE.parents[2]
+TEST_DAGS = HERE / "test_dags.yaml"
+SCHEMA = REPO_ROOT / "airflow-core" / "src" / "airflow" / "serialization" /
"schema.json"
+
+# Name the file a Dag was declared in: Python's own Dag file, or the SDK's
bundle.
+DAG_KEYS_NOT_COMPARED = frozenset({"fileloc", "relative_fileloc",
"_processor_dags_folder"})
+
+TASK_KEYS_NOT_COMPARED = frozenset(
+ {
+ # Operator identity: Python names the Python class that ran, a
language SDK a fixed pair.
+ # Neither is ever imported on the Airflow side.
+ "task_type",
+ "_task_module",
+ # The SDK's marker for a task in its language, which Python has no
equivalent for.
+ "language",
+ # Python bookkeeping for mapped tasks and retry policies, neither of
which a language SDK's
+ # Dag can declare.
+ "_needs_expansion",
+ "has_retry_policy",
+ # The stub contract every language SDK task carries, which the Python
operator built here
+ # has no counterpart for. Loading the SDK's output through Airflow
checks them.
+ "is_stub",
+ "_arg_bindings",
+ }
+)
+
+
+def run(command: list[str]) -> None:
+ if subprocess.run(command, cwd=REPO_ROOT, check=False).returncode:
+ raise SystemExit(f"`{' '.join(command)}` failed")
+
+
+def get_task_defaults() -> dict[str, Any]:
+ """Map each task field the Dag schema gives a default to that default."""
+ fields =
json.loads(SCHEMA.read_text())["definitions"]["operator"]["properties"]
+ return {key: field["default"] for key, field in fields.items() if
field.get("default") is not None}
+
+
+def normalize_as_javascript(value: Any) -> Any:
+ """Read a JSON value as JavaScript does: one number type, and a bool that
is not a number."""
+ if isinstance(value, bool):
+ return ("bool", value)
+ if isinstance(value, (int, float)):
+ return float(value)
+ if isinstance(value, list):
+ return [normalize_as_javascript(item) for item in value]
+ if isinstance(value, dict):
+ return {key: normalize_as_javascript(item) for key, item in
value.items()}
+ return value
+
+
+def is_same_json(python: Any, sdk: Any) -> bool:
+ return normalize_as_javascript(python) == normalize_as_javascript(sdk)
+
+
+def find_differences(path: str, python: Any, sdk: Any) -> list[str]:
+ """List where two JSON values differ, down to the innermost key or
index."""
+ if isinstance(python, dict) and isinstance(sdk, dict):
+ problems = []
+ for key in sorted(python.keys() | sdk.keys()):
+ if key not in sdk:
+ problems.append(f"{path}.{key} is missing, Python writes
{python[key]!r}")
+ elif key not in python:
+ problems.append(f"{path}.{key} is {sdk[key]!r}, which Python
does not write")
+ else:
+ problems.extend(find_differences(f"{path}.{key}", python[key],
sdk[key]))
+ return problems
+ if isinstance(python, list) and isinstance(sdk, list) and len(python) ==
len(sdk):
+ return [
+ problem
+ for index, (python_item, sdk_item) in enumerate(zip(python, sdk))
+ for problem in find_differences(f"{path}[{index}]", python_item,
sdk_item)
+ ]
+ if not is_same_json(python, sdk):
+ return [f"{path} is {sdk!r}, Python writes {python!r}"]
+ return []
+
+
+def compare_fields(
+ python: dict[str, Any], sdk: dict[str, Any], not_compared: frozenset[str],
defaults: dict[str, Any]
+) -> list[str]:
+ """
+ Compare two serialized objects key by key.
+
+ A key the SDK leaves out is fine when Python wrote its schema default.
Python keeps such a value
+ when its ``client_defaults`` table disagrees with the schema, a table a
language SDK is never sent,
+ and Airflow reads a missing field as its default anyway.
+ """
+ problems = []
+ for key in sorted((python.keys() | sdk.keys()) - not_compared):
+ if key in python and key in sdk:
+ problems.extend(find_differences(key, python[key], sdk[key]))
+ elif key in sdk:
+ problems.append(f"{key} is {sdk[key]!r}, which Python does not
write")
+ elif key not in defaults or not is_same_json(python[key],
defaults[key]):
+ problems.append(f"{key} is missing, Python writes {python[key]!r}")
+ return problems
+
+
+def compare_dag(python: dict[str, Any], sdk: dict[str, Any], defaults:
dict[str, Any]) -> list[str]:
+ problems = compare_fields(python, sdk, DAG_KEYS_NOT_COMPARED | {"tasks"},
{})
+ python_ids = [task["__var"]["task_id"] for task in python["tasks"]]
+ sdk_ids = [task["__var"]["task_id"] for task in sdk["tasks"]]
+ if sdk_ids != python_ids:
+ return [*problems, f"tasks are {sdk_ids}, Python writes {python_ids}"]
+ for python_task, sdk_task in zip(python["tasks"], sdk["tasks"]):
+ task_id = python_task["__var"]["task_id"]
+ if sdk_task["__type"] != python_task["__type"]:
+ problems.append(f"task {task_id} is a {sdk_task['__type']!r}, not
a {python_task['__type']!r}")
+ problems.extend(
+ f"task {task_id}: {problem}"
+ for problem in compare_fields(
+ python_task["__var"], sdk_task["__var"],
TASK_KEYS_NOT_COMPARED, defaults
+ )
+ )
+ return problems
+
+
+def compare(python: dict[str, Any], sdk: dict[str, Any], task_defaults:
dict[str, Any]) -> list[str]:
+ """List every way the SDK's serialization differs from Python's."""
+ if sdk.keys() != python.keys():
+ return [f"the Dags are {sorted(sdk)}, Python writes {sorted(python)}"]
+ problems = []
+ for dag_id, python_dag in python.items():
+ sdk_dag = sdk[dag_id]
+ if sdk_dag["__version"] != python_dag["__version"]:
+ problems.append(
+ f"{dag_id}: __version is {sdk_dag['__version']}, Python writes
{python_dag['__version']}"
+ )
+ problems.extend(
+ f"{dag_id}: {problem}"
+ for problem in compare_dag(python_dag["dag"], sdk_dag["dag"],
task_defaults)
+ )
+ return problems
+
+
+def main(argv: list[str] | None = None) -> int:
+ parser = argparse.ArgumentParser(
+ description=__doc__,
formatter_class=argparse.RawDescriptionHelpFormatter
+ )
+ parser.add_argument("--sdk", required=True, help="the SDK's name, as in
serialized_<sdk>.json")
+ parser.add_argument("command", nargs="+", help="the SDK's serializer,
after --")
+ args = parser.parse_args(argv)
+
+ directory = Path(tempfile.mkdtemp(prefix=f"{args.sdk}-serialization-"))
+ python_output = directory / "serialized_python.json"
+ sdk_output = directory / f"serialized_{args.sdk}.json"
+ received_output = directory / f"received_{args.sdk}.json"
+ run([*args.command, str(TEST_DAGS), str(sdk_output)])
+ run(
+ [
+ "uv",
+ "run",
+ "--project",
+ "airflow-core",
+ # airflow-core's dev group pulls in providers and extras that
build native code; the
+ # serializer only needs airflow-core itself.
+ "--no-dev",
+ "python",
+ str(HERE / "serialize_python.py"),
+ str(TEST_DAGS),
+ str(python_output),
+ "--receive",
+ str(sdk_output),
+ str(received_output),
+ ]
+ )
+
+ python = json.loads(python_output.read_text())
+ problems = compare(python, json.loads(received_output.read_text()),
get_task_defaults())
+ if problems:
+ print(
+ f"The {args.sdk} serialization differs from Python's in
{len(problems)} place(s):",
+ file=sys.stderr,
+ )
+ for problem in problems:
+ print(f" {problem}", file=sys.stderr)
+ print(f"The serializations are kept in {directory}", file=sys.stderr)
+ return 1
+ shutil.rmtree(directory)
+ print(f"The {args.sdk} SDK serializes all {len(python)} Dags of
{TEST_DAGS.name} as Python does")
+ return 0
+
+
+if __name__ == "__main__":
+ sys.exit(main())
diff --git a/scripts/ci/lang_sdk_serialization/serialize_python.py
b/scripts/ci/lang_sdk_serialization/serialize_python.py
new file mode 100644
index 00000000000..3bb1406cb0e
--- /dev/null
+++ b/scripts/ci/lang_sdk_serialization/serialize_python.py
@@ -0,0 +1,140 @@
+# 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.
+r"""
+Serialize the Dags of test_dags.yaml with Airflow's own serializer.
+
+Each Dag is built with the Python authoring API and written as
``DagSerialization.to_dict`` returns it,
+keyed by Dag id. ``--receive`` also takes a language SDK's output as Airflow
receives it: it fills in the
+Dag fields the SDK leaves to Airflow's config, writes the result, and loads
every Dag through
+``DagSerialization.validate_schema`` and ``from_dict``. compare.py runs it as::
+
+ uv run --project airflow-core --no-dev python
scripts/ci/lang_sdk_serialization/serialize_python.py \
+ scripts/ci/lang_sdk_serialization/test_dags.yaml
serialized_python.json \
+ --receive serialized_typescript.json received_typescript.json
+"""
+
+from __future__ import annotations
+
+import argparse
+import copy
+import datetime
+import json
+import sys
+from collections.abc import Callable
+from pathlib import Path
+from typing import Any
+
+import yaml
+
+from airflow.configuration import conf
+from airflow.sdk import DAG, BaseOperator, TaskGroup
+from airflow.serialization.serialized_objects import DagSerialization
+
+# The Dag fields the Python DAG reads from config when they are left unset. A
language SDK cannot read
+# Airflow's config, so it leaves them out, and Airflow fills them in when it
receives the SDK's Dags.
+CONFIG_BACKED_DAG_FIELDS: dict[str, Callable[[], Any]] = {
+ "max_active_tasks": lambda: conf.getint("core",
"max_active_tasks_per_dag"),
+ "max_active_runs": lambda: conf.getint("core", "max_active_runs_per_dag"),
+ "max_consecutive_failed_dag_runs": lambda: conf.getint("core",
"max_consecutive_failed_dag_runs_per_dag"),
+ "catchup": lambda: conf.getboolean("scheduler", "catchup_by_default"),
+ "disable_bundle_versioning": lambda: conf.getboolean("dag_processor",
"disable_bundle_versioning"),
+}
+
+
+class NoopOperator(BaseOperator):
+ """Stands in for a language SDK's task: a task with no Python behaviour."""
+
+ def execute(self, context):
+ return None
+
+
+def construct_datetime(loader: yaml.SafeLoader, node: yaml.ScalarNode) ->
datetime.datetime:
+ return datetime.datetime.fromisoformat(loader.construct_scalar(node))
+
+
+def construct_timedelta(loader: yaml.SafeLoader, node: yaml.ScalarNode) ->
datetime.timedelta:
+ return datetime.timedelta(seconds=float(loader.construct_scalar(node)))
+
+
+yaml.SafeLoader.add_constructor("!datetime", construct_datetime)
+yaml.SafeLoader.add_constructor("!timedelta", construct_timedelta)
+
+
+def build_dag(case: dict[str, Any]) -> DAG:
+ dag = DAG(case["dag_id"], **case.get("spec", {}))
+ groups: dict[str, TaskGroup] = {}
+ for group_id in case.get("groups", []):
+ # A group id is fully qualified, so its parent is whatever comes
before the last dot.
+ parent_id, _, local_id = group_id.rpartition(".")
+ groups[group_id] = TaskGroup(local_id, dag=dag,
parent_group=groups[parent_id] if parent_id else None)
+ for task in case["tasks"]:
+ NoopOperator(
+ task_id=task["task_id"], dag=dag,
task_group=groups.get(task.get("group")), **task.get("spec", {})
+ )
+ for task in case["tasks"]:
+ task_id = f"{task['group']}.{task['task_id']}" if "group" in task else
task["task_id"]
+ for upstream in task.get("upstream", []):
+ dag.get_task(upstream) >> dag.get_task(task_id)
+ for upstream, downstream in case.get("order_edges", []):
+ get_node(dag, groups, upstream) >> get_node(dag, groups, downstream)
+ return dag
+
+
+def get_node(dag: DAG, groups: dict[str, TaskGroup], node_id: str):
+ """Resolve an edge endpoint: the task group with that id if there is one,
else the task."""
+ return groups[node_id] if node_id in groups else dag.get_task(node_id)
+
+
+def receive(sdk_output: Path, received_output: Path) -> None:
+ received = json.loads(sdk_output.read_text())
+ for data in received.values():
+ for key, read_config in CONFIG_BACKED_DAG_FIELDS.items():
+ data["dag"].setdefault(key, read_config())
+ received_output.write_text(json.dumps(received, indent=2) + "\n")
+ for dag_id, data in received.items():
+ try:
+ DagSerialization.validate_schema(data)
+ DagSerialization.from_dict(copy.deepcopy(data))
+ except Exception:
+ print(f"Airflow cannot load Dag {dag_id!r} as the SDK wrote it",
file=sys.stderr)
+ raise
+
+
+def main() -> None:
+ parser = argparse.ArgumentParser(
+ description=__doc__,
formatter_class=argparse.RawDescriptionHelpFormatter
+ )
+ parser.add_argument("test_dags", type=Path, help="the test cases,
test_dags.yaml")
+ parser.add_argument("output", type=Path, help="the JSON file to write")
+ parser.add_argument(
+ "--receive",
+ nargs=2,
+ type=Path,
+ metavar=("SDK_OUTPUT", "RECEIVED_OUTPUT"),
+ help="a JSON file a language SDK wrote, and where to write it as
Airflow receives it",
+ )
+ args = parser.parse_args()
+
+ cases = yaml.safe_load(args.test_dags.read_text())["dags"]
+ serialized = {case["dag_id"]: DagSerialization.to_dict(build_dag(case))
for case in cases}
+ args.output.write_text(json.dumps(serialized, indent=2) + "\n")
+ if args.receive:
+ receive(*args.receive)
+
+
+if __name__ == "__main__":
+ main()
diff --git a/scripts/ci/lang_sdk_serialization/test_dags.yaml
b/scripts/ci/lang_sdk_serialization/test_dags.yaml
new file mode 100644
index 00000000000..d80c90bbb40
--- /dev/null
+++ b/scripts/ci/lang_sdk_serialization/test_dags.yaml
@@ -0,0 +1,166 @@
+# 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 Dags compare.py checks a language SDK's serializer on.
serialize_python.py and the SDK's own
+# serializer each build every Dag here with their SDK's authoring API, and the
two serializations
+# have to agree.
+#
+# `spec` keys are the serialized Dag schema's snake_case names, which each
side maps onto
+# its authoring API. A task is called with the tasks named in its `upstream`,
which come
+# before it. Task and group ids are fully qualified, and a group comes after
its parent.
+# `order_edges` are `[upstream, downstream]` pairs of task or group ids.
+#
+# `!datetime` takes an ISO 8601 timestamp, and `!timedelta` a number of
seconds.
+dags:
+ # No schedule and no options: NullTimetable plus the fields Python always
emits.
+ - dag_id: conformance_minimal
+ tasks:
+ - task_id: solo
+
+ # A cron schedule with every Dag option set, and one fully configured task
followed by
+ # one left entirely at its schema defaults.
+ - dag_id: conformance_cron
+ spec:
+ schedule: "0 3 * * *"
+ description: a Dag that sets everything
+ dag_display_name: Conformance Cron
+ doc_md: "# notes"
+ tags: [gamma, alpha, gamma]
+ start_date: !datetime 2026-01-01T00:00:00+00:00
+ end_date: !datetime 2026-12-31T23:30:15+00:00
+ dagrun_timeout: !timedelta 300
+ catchup: true
+ render_template_as_native_obj: true
+ disable_bundle_versioning: true
+ is_paused_upon_creation: true
+ max_active_tasks: 8
+ max_active_runs: 3
+ max_consecutive_failed_dag_runs: 2
+ tasks:
+ - task_id: extract
+ spec:
+ owner: data-team
+ start_date: !datetime 2026-02-01T00:00:00+00:00
+ end_date: !datetime 2026-11-30T00:00:00+00:00
+ trigger_rule: all_done
+ depends_on_past: true
+ ignore_first_depends_on_past: true
+ wait_for_past_depends_before_skipping: true
+ wait_for_downstream: true
+ retries: 2
+ queue: typescript
+ pool: tiny
+ pool_slots: 2
+ execution_timeout: !timedelta 120
+ retry_delay: !timedelta 600
+ retry_exponential_backoff: 2
+ max_retry_delay: !timedelta 900
+ priority_weight: 5
+ weight_rule: upstream
+ executor: LocalExecutor
+ do_xcom_push: false
+ email_on_failure: false
+ email_on_retry: false
+ doc_md: extracts things
+ map_index_template: "{{ task.task_id }}"
+ max_active_tis_per_dag: 3
+ max_active_tis_per_dagrun: 4
+ # Every value explicitly at its schema default, so all of them are
omitted.
+ - task_id: transform
+ upstream: [extract]
+ spec:
+ owner: airflow
+ trigger_rule: all_success
+ depends_on_past: false
+ ignore_first_depends_on_past: false
+ wait_for_past_depends_before_skipping: false
+ wait_for_downstream: false
+ retries: 0
+ queue: default
+ pool: default_pool
+ pool_slots: 1
+ retry_delay: !timedelta 300
+ retry_exponential_backoff: 0
+ priority_weight: 1
+ weight_rule: downstream
+ do_xcom_push: true
+ email_on_failure: true
+ email_on_retry: true
+
+ # The @once preset, fail_fast, and a task feeding two downstreams. fail_fast
is set here
+ # rather than on conformance_cron, since Airflow only allows it when every
task uses the
+ # all_success trigger rule.
+ - dag_id: conformance_once
+ spec:
+ schedule: "@once"
+ fail_fast: true
+ tasks:
+ - task_id: seed
+ - task_id: beta
+ upstream: [seed]
+ - task_id: alpha
+ upstream: [seed]
+
+ # The @continuous preset. Airflow caps a continuous Dag at one active run, so
+ # max_active_runs is set rather than left at the default of 16.
+ - dag_id: conformance_continuous
+ spec:
+ schedule: "@continuous"
+ max_active_runs: 1
+ tasks:
+ - task_id: watch
+
+ # Fan-out and fan-in, so downstream_task_ids has to be sorted rather than
kept in wiring
+ # order. The schedule is not written "@daily" because Python expands a
preset before
+ # serializing it, which serde.test.ts covers separately.
+ - dag_id: conformance_diamond
+ spec:
+ schedule: "0 0 * * *"
+ tasks:
+ - task_id: root
+ - task_id: right
+ upstream: [root]
+ - task_id: left
+ upstream: [root]
+ - task_id: join
+ upstream: [right, left]
+
+ # Nested task groups and order-only edges. A group edge is recorded on the
group, and is
+ # also expanded onto the task graph, where it reaches the group's roots and
leaves
+ # rather than every task it holds: staging.stage feeds staging.checks.nulls,
so only
+ # the latter is a leaf. Both sides have to agree on that expansion, and on
the local
+ # _group_id a nested group carries.
+ - dag_id: conformance_groups
+ spec:
+ schedule: "@once"
+ groups: [staging, staging.checks, publish]
+ tasks:
+ - task_id: extract
+ - task_id: stage
+ group: staging
+ - task_id: nulls
+ group: staging.checks
+ upstream: [staging.stage]
+ - task_id: push
+ group: publish
+ - task_id: load
+ - task_id: notify
+ order_edges:
+ - [extract, staging]
+ - [staging, publish]
+ - [publish, load]
+ - [load, notify]
diff --git a/scripts/tests/ci/lang_sdk_serialization/__init__.py
b/scripts/tests/ci/lang_sdk_serialization/__init__.py
new file mode 100644
index 00000000000..13a83393a91
--- /dev/null
+++ b/scripts/tests/ci/lang_sdk_serialization/__init__.py
@@ -0,0 +1,16 @@
+# 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.
diff --git a/scripts/tests/ci/lang_sdk_serialization/test_compare.py
b/scripts/tests/ci/lang_sdk_serialization/test_compare.py
new file mode 100644
index 00000000000..603b7a8f4df
--- /dev/null
+++ b/scripts/tests/ci/lang_sdk_serialization/test_compare.py
@@ -0,0 +1,245 @@
+# 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.
+from __future__ import annotations
+
+import copy
+import json
+import subprocess
+from pathlib import Path
+from unittest import mock
+
+import pytest
+from ci.lang_sdk_serialization.compare import TEST_DAGS, compare,
get_task_defaults, main
+
+DEFAULTS = {"retry_delay": 300.0}
+
+
+def build_serialized(fileloc: str, tasks: list[dict]) -> dict:
+ return {
+ "d": {
+ "__version": 3,
+ "dag": {
+ "dag_id": "d",
+ "fileloc": fileloc,
+ "timezone": "UTC",
+ "catchup": False,
+ "tags": ["a"],
+ "task_group": {"prefix_group_id": True, "children":
{"extract": ["operator", "extract"]}},
+ "tasks": [{"__type": "operator", "__var": task} for task in
tasks],
+ },
+ }
+ }
+
+
+PYTHON = build_serialized(
+ "/dags/d.py",
+ [
+ {
+ "task_id": "extract",
+ "retries": 2,
+ "pool_slots": 1,
+ "retry_delay": 300.0,
+ "task_type": "NoopOperator",
+ },
+ {"task_id": "load", "retry_delay": 300.0, "task_type": "NoopOperator"},
+ ],
+)
+
+# Differs from PYTHON only where an SDK may: where it was declared, what
stands for its
+# tasks, a number JavaScript writes without a fraction, and a field left at
its default.
+SDK = build_serialized(
+ "/bundles/app/bundle.mjs",
+ [
+ {"task_id": "extract", "retries": 2, "pool_slots": 1.0, "task_type":
"Task", "is_stub": True},
+ {"task_id": "load", "task_type": "Task", "language": "typescript",
"_arg_bindings": []},
+ ],
+)
+
+
+def get_task(serialized: dict, index: int = 0) -> dict:
+ return serialized["d"]["dag"]["tasks"][index]["__var"]
+
+
+def test_accepts_the_differences_an_sdk_is_allowed():
+ assert compare(PYTHON, SDK, DEFAULTS) == []
+
+
[email protected](
+ ("change", "expected"),
+ [
+ pytest.param(
+ lambda sdk: get_task(sdk).update(retries=3),
+ ["d: task extract: retries is 3, Python writes 2"],
+ id="task-value",
+ ),
+ pytest.param(
+ lambda sdk: get_task(sdk).update(pool_slots=True),
+ ["d: task extract: pool_slots is True, Python writes 1"],
+ id="bool-is-not-a-number",
+ ),
+ pytest.param(
+ lambda sdk: get_task(sdk).pop("retries"),
+ ["d: task extract: retries is missing, Python writes 2"],
+ id="task-key-missing",
+ ),
+ pytest.param(
+ lambda sdk: get_task(sdk).update(owner="me"),
+ ["d: task extract: owner is 'me', which Python does not write"],
+ id="task-key-extra",
+ ),
+ pytest.param(
+ lambda sdk: sdk["d"]["dag"].pop("timezone"),
+ ["d: timezone is missing, Python writes 'UTC'"],
+ id="dag-key-missing",
+ ),
+ pytest.param(
+ lambda sdk: sdk["d"]["dag"].update(tags=["b"]),
+ ["d: tags[0] is 'b', Python writes 'a'"],
+ id="nested-list",
+ ),
+ pytest.param(
+ lambda sdk:
sdk["d"]["dag"]["task_group"].update(prefix_group_id=False),
+ ["d: task_group.prefix_group_id is False, Python writes True"],
+ id="nested-dict",
+ ),
+ pytest.param(
+ lambda sdk: sdk["d"]["dag"]["tasks"].reverse(),
+ ["d: tasks are ['load', 'extract'], Python writes ['extract',
'load']"],
+ id="task-order",
+ ),
+ pytest.param(
+ lambda sdk: sdk["d"]["dag"]["tasks"][0].update(__type="taskgroup"),
+ ["d: task extract is a 'taskgroup', not a 'operator'"],
+ id="task-encoding",
+ ),
+ pytest.param(
+ lambda sdk: sdk["d"].update(__version=4),
+ ["d: __version is 4, Python writes 3"],
+ id="version",
+ ),
+ pytest.param(
+ lambda sdk: sdk.update(e=sdk["d"]),
+ ["the Dags are ['d', 'e'], Python writes ['d']"],
+ id="dag-ids",
+ ),
+ ],
+)
+def test_reports_each_difference(change, expected):
+ sdk = copy.deepcopy(SDK)
+ change(sdk)
+
+ assert compare(PYTHON, sdk, DEFAULTS) == expected
+
+
+def test_accepts_a_left_out_task_key_only_at_its_schema_default():
+ python = copy.deepcopy(PYTHON)
+ get_task(python)["retry_delay"] = 600.0
+
+ assert compare(python, SDK, DEFAULTS) == ["d: task extract: retry_delay is
missing, Python writes 600.0"]
+
+
+def test_reads_the_task_defaults_from_the_dag_schema():
+ defaults = get_task_defaults()
+
+ assert defaults["retry_delay"] == 300.0
+ # A null default and no default both mean the field has none.
+ assert "render_template_as_native_obj" not in defaults
+ assert "task_id" not in defaults
+
+
+def write_outputs(received: dict):
+ """
+ Stand in for both serializers, each writing its output to the paths
compare.py hands it.
+
+ The SDK leaves catchup out and the received output has it, as when Airflow
fills it in from its
+ config, so only the received output agrees with Python.
+ """
+ sdk = copy.deepcopy(received)
+ del sdk["d"]["dag"]["catchup"]
+
+ def run(command, **kwargs):
+ if command[0] == "uv":
+ python_output, _, _, received_output = command[-4:]
+ Path(python_output).write_text(json.dumps(PYTHON))
+ Path(received_output).write_text(json.dumps(received))
+ else:
+ Path(command[-1]).write_text(json.dumps(sdk))
+ return subprocess.CompletedProcess(command, 0)
+
+ return run
+
+
[email protected]("ci.lang_sdk_serialization.compare.get_task_defaults",
autospec=True, return_value=DEFAULTS)
[email protected]("ci.lang_sdk_serialization.compare.tempfile.mkdtemp",
autospec=True)
[email protected]("ci.lang_sdk_serialization.compare.subprocess.run", autospec=True)
+def test_main_runs_both_serializers_and_removes_their_output_when_they_agree(
+ mock_run, mock_mkdtemp, mock_get_task_defaults, tmp_path, capsys
+):
+ mock_mkdtemp.return_value = str(tmp_path)
+ mock_run.side_effect = write_outputs(SDK)
+
+ assert main(["--sdk", "typescript", "--", "pnpm", "exec", "tsx",
"serialize.ts"]) == 0
+
+ sdk_output = str(tmp_path / "serialized_typescript.json")
+ assert [call.args[0] for call in mock_run.call_args_list] == [
+ ["pnpm", "exec", "tsx", "serialize.ts", str(TEST_DAGS), sdk_output],
+ [
+ "uv",
+ "run",
+ "--project",
+ "airflow-core",
+ "--no-dev",
+ "python",
+ str(TEST_DAGS.parent / "serialize_python.py"),
+ str(TEST_DAGS),
+ str(tmp_path / "serialized_python.json"),
+ "--receive",
+ sdk_output,
+ str(tmp_path / "received_typescript.json"),
+ ],
+ ]
+ assert not tmp_path.exists()
+ assert "serializes all 1 Dags" in capsys.readouterr().out
+
+
[email protected]("ci.lang_sdk_serialization.compare.get_task_defaults",
autospec=True, return_value=DEFAULTS)
[email protected]("ci.lang_sdk_serialization.compare.tempfile.mkdtemp",
autospec=True)
[email protected]("ci.lang_sdk_serialization.compare.subprocess.run", autospec=True)
+def test_main_reports_the_differences_and_keeps_the_output(
+ mock_run, mock_mkdtemp, mock_get_task_defaults, tmp_path, capsys
+):
+ mock_mkdtemp.return_value = str(tmp_path)
+ sdk = copy.deepcopy(SDK)
+ get_task(sdk)["retries"] = 3
+ mock_run.side_effect = write_outputs(sdk)
+
+ assert main(["--sdk", "typescript", "--", "serialize"]) == 1
+
+ assert (tmp_path / "serialized_typescript.json").exists()
+ assert "d: task extract: retries is 3, Python writes 2" in
capsys.readouterr().err
+
+
[email protected]("ci.lang_sdk_serialization.compare.tempfile.mkdtemp",
autospec=True)
[email protected]("ci.lang_sdk_serialization.compare.subprocess.run", autospec=True)
+def test_main_stops_when_a_serializer_fails(mock_run, mock_mkdtemp, tmp_path):
+ mock_mkdtemp.return_value = str(tmp_path)
+ mock_run.return_value = subprocess.CompletedProcess([], 1)
+
+ with pytest.raises(SystemExit, match="`serialize .*` failed"):
+ main(["--sdk", "typescript", "--", "serialize"])
+
+ assert mock_run.call_count == 1
diff --git a/ts-sdk/.pre-commit-config.yaml b/ts-sdk/.pre-commit-config.yaml
index ad2bcb99c17..4346cce9f4a 100644
--- a/ts-sdk/.pre-commit-config.yaml
+++ b/ts-sdk/.pre-commit-config.yaml
@@ -61,6 +61,20 @@ repos:
additional_dependencies: ['[email protected]']
pass_filenames: false
require_serial: true
+ - id: check-ts-sdk-serialization-conformance
+ name: Check the TS SDK serializes Dags the way Airflow does
+ description: "Serialize the shared test Dags with the TS SDK and with
Airflow, and compare the two"
+ entry: ./scripts/ci/prek/check_serialization_conformance.py
+ language: node
+ files: >
+ (?x)
+ ^schema/dag-schema\.json$|
+ ^src/coordinator/serde\.ts$|
+ ^src/sdk/dag\.ts$|
+ ^tests/conformance/.*$
+ additional_dependencies: ['[email protected]']
+ pass_filenames: false
+ require_serial: true
- id: check-ts-sdk-docs-package-version-in-sync
name: Check ts-sdk/docs package.json versions match ts-sdk/package.json
entry: ../scripts/ci/prek/check_ts_sdk_docs_package_version_in_sync.py
diff --git a/ts-sdk/package.json b/ts-sdk/package.json
index bd91883d55e..ab205567bee 100644
--- a/ts-sdk/package.json
+++ b/ts-sdk/package.json
@@ -86,6 +86,7 @@
"tsx": "^4.21.0",
"typescript": "^6.0.2",
"typescript-eslint": "^8.60.0",
- "vitest": "^4.1.7"
+ "vitest": "^4.1.7",
+ "yaml": "^2.9.1"
}
}
diff --git a/ts-sdk/pnpm-lock.yaml b/ts-sdk/pnpm-lock.yaml
index 543171cdcff..f1cd0251f6f 100644
--- a/ts-sdk/pnpm-lock.yaml
+++ b/ts-sdk/pnpm-lock.yaml
@@ -41,7 +41,10 @@ importers:
version: 8.67.0([email protected])([email protected])
vitest:
specifier: ^4.1.7
- version:
4.1.7(@types/[email protected])(@vitest/[email protected])([email protected](@types/[email protected])([email protected])([email protected]))
+ version:
4.1.7(@types/[email protected])(@vitest/[email protected])([email protected](@types/[email protected])([email protected])([email protected])([email protected]))
+ yaml:
+ specifier: ^2.9.1
+ version: 2.9.1
example:
dependencies:
@@ -1183,6 +1186,11 @@ packages:
resolution: {integrity:
sha512-BN22B5eaMMI9UMtjrGd5g5eCYPpCPDUy0FJXbYsaT5zYxjFOckS53SQDE3pWkVoWpHXVb3BrYcEN4Twa55B5cA==}
engines: {node: '>=0.10.0'}
+ [email protected]:
+ resolution: {integrity:
sha512-3NxN8+78OdzbT7C/WjGsyfPAtJaN3FNDsWxv7Y7mcDsT/oOmgW8BpyQQFFBnvZE3j9Y2Sdz1ULFLezL7Eb2yFw==}
+ engines: {node: '>= 14.6'}
+ hasBin: true
+
[email protected]:
resolution: {integrity:
sha512-rVksvsnNCdJ/ohGc6xgPwyN8eheCxsiLM8mxuE/t/mOVqJewPuO1miLpTHQiRgTKCLexL4MeAFVagts7HmNZ2Q==}
engines: {node: '>=10'}
@@ -1571,7 +1579,7 @@ snapshots:
obug: 2.1.4
std-env: 4.2.0
tinyrainbow: 3.1.1
- vitest:
4.1.7(@types/[email protected])(@vitest/[email protected])([email protected](@types/[email protected])([email protected])([email protected]))
+ vitest:
4.1.7(@types/[email protected])(@vitest/[email protected])([email protected](@types/[email protected])([email protected])([email protected])([email protected]))
optional: true
'@vitest/[email protected]':
@@ -1583,13 +1591,13 @@ snapshots:
chai: 6.2.2
tinyrainbow: 3.1.0
-
'@vitest/[email protected]([email protected](@types/[email protected])([email protected])([email protected]))':
+
'@vitest/[email protected]([email protected](@types/[email protected])([email protected])([email protected])([email protected]))':
dependencies:
'@vitest/spy': 4.1.7
estree-walker: 3.0.3
magic-string: 0.30.21
optionalDependencies:
- vite: 8.0.10(@types/[email protected])([email protected])([email protected])
+ vite: 8.0.10(@types/[email protected])([email protected])([email protected])([email protected])
'@vitest/[email protected]':
dependencies:
@@ -2115,7 +2123,7 @@ snapshots:
dependencies:
punycode: 2.3.1
- [email protected](@types/[email protected])([email protected])([email protected]):
+ [email protected](@types/[email protected])([email protected])([email protected])([email protected]):
dependencies:
lightningcss: 1.33.0
picomatch: 4.0.7
@@ -2127,11 +2135,12 @@ snapshots:
esbuild: 0.28.2
fsevents: 2.3.3
tsx: 4.23.12
+ yaml: 2.9.1
-
[email protected](@types/[email protected])(@vitest/[email protected])([email protected](@types/[email protected])([email protected])([email protected])):
+
[email protected](@types/[email protected])(@vitest/[email protected])([email protected](@types/[email protected])([email protected])([email protected])([email protected])):
dependencies:
'@vitest/expect': 4.1.7
- '@vitest/mocker':
4.1.7([email protected](@types/[email protected])([email protected])([email protected]))
+ '@vitest/mocker':
4.1.7([email protected](@types/[email protected])([email protected])([email protected])([email protected]))
'@vitest/pretty-format': 4.1.7
'@vitest/runner': 4.1.7
'@vitest/snapshot': 4.1.7
@@ -2148,7 +2157,7 @@ snapshots:
tinyexec: 1.1.1
tinyglobby: 0.2.16
tinyrainbow: 3.1.0
- vite: 8.0.10(@types/[email protected])([email protected])([email protected])
+ vite: 8.0.10(@types/[email protected])([email protected])([email protected])([email protected])
why-is-node-running: 2.3.0
optionalDependencies:
'@types/node': 26.2.0
@@ -2167,4 +2176,6 @@ snapshots:
[email protected]: {}
+ [email protected]: {}
+
[email protected]: {}
diff --git a/ts-sdk/scripts/ci/prek/check_serialization_conformance.py
b/ts-sdk/scripts/ci/prek/check_serialization_conformance.py
new file mode 100755
index 00000000000..c51dfffe183
--- /dev/null
+++ b/ts-sdk/scripts/ci/prek/check_serialization_conformance.py
@@ -0,0 +1,44 @@
+#!/usr/bin/env python3
+# 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.
+"""Check the TS SDK serializes Dags as Airflow does, with
scripts/ci/lang_sdk_serialization/compare.py."""
+
+from __future__ import annotations
+
+import subprocess
+import sys
+from pathlib import Path
+
+sys.path.insert(0, str(Path(__file__).resolve().parents[4] / "scripts" / "ci"
/ "prek"))
+
+from common_prek_utils import AIRFLOW_ROOT_PATH, run_command
+
+if __name__ not in ("__main__", "__mp_main__"):
+ raise SystemExit(
+ "This file is intended to be executed as an executable program. You
cannot use it as a module."
+ f"To run this script, run the ./{__file__} command"
+ )
+
+if __name__ == "__main__":
+ run_command(
+ ["pnpm", "install", "--frozen-lockfile",
"--config.confirmModulesPurge=false"],
+ cwd=AIRFLOW_ROOT_PATH / "ts-sdk",
+ )
+ compare = AIRFLOW_ROOT_PATH / "scripts" / "ci" / "lang_sdk_serialization"
/ "compare.py"
+ serializer = ["pnpm", "--dir", "ts-sdk", "exec", "tsx",
"tests/conformance/serialize_typescript.ts"]
+ command = [sys.executable, str(compare), "--sdk", "typescript", "--",
*serializer]
+ sys.exit(subprocess.run(command, check=False).returncode)
diff --git a/ts-sdk/src/coordinator/serde.ts b/ts-sdk/src/coordinator/serde.ts
new file mode 100644
index 00000000000..900b41b6887
--- /dev/null
+++ b/ts-sdk/src/coordinator/serde.ts
@@ -0,0 +1,751 @@
+/*!
+ * 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.
+ */
+
+// Turns a Dag declared in TypeScript into Airflow's DagSerialization v3 JSON —
+// what the Dag processor stores and the scheduler reads. The format is
+// Airflow-internal rather than an SDK schema, so it is reimplemented per
+// language against `airflow-core/src/airflow/serialization/schema.json`; see
+// `airflow-core/adr/lang-sdk/0004-dag-parsing.md` for the field table this
+// follows, and the Java SDK's `Serde.kt` for the same job in another language.
+//
+// Byte-parity with Python's serializer is not the goal — Python omits fields
+// against a `client_defaults` table this SDK does not receive. What has to
hold
+// is that `DagSerialization.from_dict` rebuilds the same Dag, which
+// scripts/ci/lang_sdk_serialization/compare.py checks against Python's own
serialization.
+//
+// This module only produces the payload. Answering a Dag-parsing request with
+// it is the bundle's job, once the coordinator has a parse request to answer.
+
+import { relative as relativePath } from "node:path";
+
+import {
+ DAG_SCHEMA_FIELDS,
+ TASK_SCHEMA_FIELDS,
+ type SchemaField,
+} from "../generated/dag-schema-fields.js";
+import type { JsonValue } from "../sdk/client-types.js";
+import {
+ getDagOrderEdges,
+ getDagTaskGroups,
+ getDagTaskInputs,
+ getDagTaskRecords,
+ isPlainRecord,
+ isTaskRef,
+ type Dag,
+ type RecordedInputs,
+ type TaskGroupRecord,
+} from "../sdk/dag.js";
+
+/** A serialized Dag: JSON, by the time it reaches the supervisor as msgpack.
*/
+type SerializedValue = JsonValue;
+
+/** Airflow's type/var encoding, as `BaseSerialization.serialize()` emits it.
*/
+interface TypeEncoded {
+ readonly __type: string;
+ readonly __var: SerializedValue;
+}
+
+/**
+ * Identity every TypeScript task carries, in place of the Python operator
class
+ * a Python Dag would name.
+ *
+ * Fixed rather than derived: nothing on the Airflow side imports
`_task_module`
+ * — `SerializedBaseOperator.populate_operator` only compares the pair as
+ * strings when matching plugin extra links — so the pair is free to name the
+ * coordinator that actually runs the task, which makes every TypeScript task
+ * greppable in the UI and the metadata DB.
+ */
+const TASK_TYPE = "TypeScriptOperator";
+const TASK_MODULE = "airflow.sdk.coordinators.node";
+
+/**
+ * Marks the tasks this SDK serialized, as the Java SDK marks its own.
+ *
+ * Nothing in airflow-core reads it today; the `operator` schema definition
+ * allows additional properties, so it rides along as a marker for tooling that
+ * wants to tell language-native tasks apart without parsing `_task_module`.
+ */
+const TASK_LANGUAGE = "typescript";
+
+/** How one set of authoring fields is written into a serialized object. */
+interface FieldRules {
+ readonly fields: Readonly<Record<string, SchemaField>>;
+ /**
+ * Fields whose {__type, __var} wrapper survives; every other field is
+ * unwrapped to the bare __var, as Python's `serialize_to_json` does.
+ *
+ * Neither set overlaps the generated authoring fields today, so in practice
+ * everything is unwrapped. They are named so that a decorated field added to
+ * the schema later takes the right path rather than silently losing its
+ * wrapper. Source: `DagSerialization._decorated_fields` and
+ * `OperatorSerialization._decorated_fields`.
+ */
+ readonly decorated: ReadonlySet<string>;
+ /** Fields never written, whatever the author set. */
+ readonly omitted: ReadonlySet<string>;
+}
+
+const DAG_FIELD_RULES: FieldRules = {
+ fields: DAG_SCHEMA_FIELDS,
+ decorated: new Set(["default_args", "access_control"]),
+ omitted: new Set(),
+};
+
+const TASK_FIELD_RULES: FieldRules = {
+ fields: TASK_SCHEMA_FIELDS,
+ decorated: new Set(["executor_config"]),
+ // Python drops both unless the operator names an email recipient
+ // (`OperatorSerialization._serialize_node`). A TypeScript task has no
`email`
+ // field to name one, so writing them would describe a notification that can
+ // never be sent.
+ omitted: new Set(["email_on_failure", "email_on_retry"]),
+};
+
+const NULL_TIMETABLE = "airflow.timetables.simple.NullTimetable";
+const ONCE_TIMETABLE = "airflow.timetables.simple.OnceTimetable";
+const CONTINUOUS_TIMETABLE = "airflow.timetables.simple.ContinuousTimetable";
+const CRON_TIMETABLE = "airflow.timetables.trigger.CronTriggerTimetable";
+
+/** Serialize one Dag to the `dag` object of a DagSerialization v3 payload. */
+export function serializeDag(
+ dag: Dag,
+ fileloc: string,
+ relativeFileloc: string,
+): Record<string, SerializedValue> {
+ const graph = buildDagGraph(dag);
+ const inputs = getDagTaskInputs(dag);
+ const data: Record<string, SerializedValue> = {
+ dag_id: dag.dagId,
+ fileloc,
+ relative_fileloc: relativeFileloc,
+ timezone: "UTC",
+ timetable: serializeTimetable(dag.spec.schedule, dag.dagId),
+ tasks: [...getDagTaskRecords(dag)].map(([taskId, record]) =>
+ serializeTask(
+ dag.dagId,
+ taskId,
+ withDagQueue(record.spec, dag.spec.queue),
+ graph.downstreamTaskIds.get(taskId),
+ inputs.get(taskId),
+ ),
+ ),
+ dag_dependencies: [],
+ task_group: serializeTaskGroups(dag, graph),
+ edge_info: {},
+ params: [],
+ // Always written by Python's serializer, so a Dag without either still
+ // round-trips to the same object.
+ deadline: null,
+ allowed_run_types: null,
+ };
+ // A field the Dag leaves unset stays out, including the ones Python reads
from
+ // config, such as max_active_tasks and catchup: Airflow fills those in from
its
+ // own config when it receives the Dag.
+ applySchemaFields(data, dag.spec, DAG_FIELD_RULES, `Dag "${dag.dagId}"`);
+ return data;
+}
+
+/**
+ * The task's spec with the Dag's queue filled in, when the task named none.
+ *
+ * Merged before the fields are written rather than after, so a queue that
+ * happens to equal the schema default is omitted the way any other defaulted
+ * field is.
+ */
+function withDagQueue(spec: object, dagQueue: string | undefined): object {
+ if (dagQueue === undefined || "queue" in spec) return spec;
+ return { ...spec, queue: dagQueue };
+}
+
+/** Serialize one task, with its downstream edges sorted for a stable payload.
*/
+function serializeTask(
+ dagId: string,
+ taskId: string,
+ spec: object,
+ downstream: ReadonlySet<string> | undefined,
+ inputs: RecordedInputs | undefined,
+): SerializedValue {
+ const data: Record<string, SerializedValue> = {
+ task_id: taskId,
+ task_type: TASK_TYPE,
+ _task_module: TASK_MODULE,
+ language: TASK_LANGUAGE,
+ // Python's operator serializer always emits this — its list value never
+ // matches the tuple default it is compared against. A TypeScript task has
+ // no Jinja templating, so the list is empty rather than absent.
+ template_fields: [],
+ // What marks a task whose arguments the API server resolves per instance
+ // and sends to a foreign runtime, as `@task.stub` does on the Python side.
+ // `get_arg_bindings` reads nothing without it.
+ is_stub: true,
+ };
+ const label = `task "${taskId}" of Dag "${dagId}"`;
+ const bindings = serializeArgBindings(inputs, label);
+ if (bindings) data["_arg_bindings"] = bindings;
+ applySchemaFields(data, spec, TASK_FIELD_RULES, label);
+ if (downstream?.size) {
+ data["downstream_task_ids"] = [...downstream].sort();
+ }
+ return { __type: "operator", __var: data };
+}
+
+/**
+ * The task's arguments as the binding spec the API server hands back at run
+ * time, one entry per argument in the order the call named them.
+ *
+ * A reference becomes an `xcom` binding naming the upstream task, and anything
+ * else a `literal` carrying the value, matching the `TaskArgBinding` union the
+ * execution API declares. `value_schema` is left out: it constrains the decode
+ * side, and nothing records a TypeScript argument's type at pack time, so
+ * omitting it says "unconstrained" rather than asserting a wrong type.
+ *
+ * `undefined` for a task called with no arguments, which needs no spec.
+ */
+function serializeArgBindings(
+ inputs: RecordedInputs | undefined,
+ label: string,
+): SerializedValue | undefined {
+ if (inputs === undefined) return undefined;
+ const entries = Object.entries(inputs);
+ if (entries.length === 0) return undefined;
+ return entries.map(([name, value]): SerializedValue =>
+ isTaskRef(value)
+ ? { name, kind: "xcom", task_id: value.taskId }
+ : { name, kind: "literal", value: toPlainJson(value, `Input "${name}" of
${label}`) },
+ );
+}
+
+/**
+ * Copy a literal argument as plain JSON, without {@link serializeValue}'s
+ * `{__type, __var}` wrapper: Python writes a literal as it is, and the runtime
+ * hands it to the handler undecoded.
+ *
+ * A value JSON cannot carry, such as a `Date` or a `Map`, is rejected rather
+ * than reaching the handler as something else.
+ */
+function toPlainJson(value: unknown, label: string): SerializedValue {
+ if (value === null || value === undefined) return null;
+ if (typeof value === "string" || typeof value === "boolean") return value;
+ if (typeof value === "number" && Number.isFinite(value)) return value;
+ if (Array.isArray(value)) return value.map((item) => toPlainJson(item,
label));
+ if (isPlainRecord(value)) {
+ const copy: Record<string, SerializedValue> = {};
+ for (const [key, item] of Object.entries(value)) copy[key] =
toPlainJson(item, label);
+ return copy;
+ }
+ throw new Error(
+ `${label} holds ${describeType(value)}, which JSON cannot carry; pass a
string, a finite ` +
+ "number, a boolean, null, an array or a plain object",
+ );
+}
+
+/**
+ * The Dag's task-group tree, rooted at the group Python builds for every Dag.
+ *
+ * `children[label] = [kind, value]`, as `serialize_for_task_group()` writes
it:
+ * a task is `["operator", task_id]` and a nested group is `["taskgroup", ...]`
+ * holding that group's own object, so the tree nests by embedding rather than
+ * by reference. `_group_id` is the group's *local* segment, not its qualified
+ * id, matching what Python records.
+ */
+function serializeTaskGroups(dag: Dag, graph: DagGraph): SerializedValue {
+ const groups = getDagTaskGroups(dag);
+ const grouped = new Set<string>();
+ for (const group of groups.values()) {
+ for (const taskId of group.taskIds) grouped.add(taskId);
+ }
+ const rootTaskIds = [...getDagTaskRecords(dag).keys()].filter((id) =>
!grouped.has(id));
+ const rootGroupIds = [...groups.values()]
+ .filter((group) => group.parentGroupId === undefined)
+ .map((group) => group.groupId);
+
+ return taskGroupObject(null, rootTaskIds, rootGroupIds, groups, graph);
+}
+
+/** One group object: the root when `groupId` is null, otherwise a nested one.
*/
+function taskGroupObject(
+ groupId: string | null,
+ taskIds: readonly string[],
+ childGroupIds: readonly string[],
+ groups: ReadonlyMap<string, TaskGroupRecord>,
+ graph: DagGraph,
+): SerializedValue {
+ const children: Record<string, SerializedValue> = {};
+ for (const taskId of taskIds) {
+ children[taskId] = ["operator", taskId];
+ }
+ for (const childId of childGroupIds) {
+ const child = groups.get(childId)!;
+ children[childId] = [
+ "taskgroup",
+ taskGroupObject(childId, child.taskIds, child.childGroupIds, groups,
graph),
+ ];
+ }
+ const own = groupId === null ? undefined : graph.groupEdges.get(groupId);
+ return {
+ // The local segment: Python's TaskGroup stores the id it was given, and
+ // rebuilds the qualified one from where the group sits in the tree.
+ _group_id: groupId === null ? null : localGroupId(groupId),
+ group_display_name: "",
+ prefix_group_id: groupId === null || groups.get(groupId)!.prefixGroupId,
+ tooltip: "",
+ ui_color: "CornflowerBlue",
+ ui_fgcolor: "#000",
+ children,
+ upstream_group_ids: sorted(own?.upstreamGroups),
+ downstream_group_ids: sorted(own?.downstreamGroups),
+ upstream_task_ids: sorted(own?.upstreamTasks),
+ downstream_task_ids: sorted(own?.downstreamTasks),
+ };
+}
+
+function localGroupId(groupId: string): string {
+ const cut = groupId.lastIndexOf(GROUP_SEPARATOR);
+ return cut === -1 ? groupId : groupId.slice(cut + 1);
+}
+
+function sorted(values: ReadonlySet<string> | undefined): string[] {
+ return [...(values ?? [])].sort();
+}
+
+const GROUP_SEPARATOR = ".";
+
+interface GroupEdgeSets {
+ readonly upstreamGroups: Set<string>;
+ readonly downstreamGroups: Set<string>;
+ readonly upstreamTasks: Set<string>;
+ readonly downstreamTasks: Set<string>;
+}
+
+/** Both views of a Dag's edges: the task graph, and what each group records.
*/
+interface DagGraph {
+ /** Each task's downstream task IDs, with every group endpoint expanded. */
+ readonly downstreamTaskIds: Map<string, Set<string>>;
+ readonly groupEdges: Map<string, GroupEdgeSets>;
+}
+
+/**
+ * Resolve a Dag's two kinds of edge into the two views a serialized Dag holds.
+ *
+ * An order-only edge with a group at either end lands in both: the group
object
+ * records it for the UI, and the task graph records it expanded, because the
+ * scheduler only ever reads task-to-task edges. A group expands to its *roots*
+ * when it is downstream and its *leaves* when it is upstream — an edge into a
+ * group reaches the tasks that start it, and one out of a group leaves from
the
+ * tasks that finish it — which is what `TaskGroup.set_upstream` does in
Python.
+ */
+function buildDagGraph(dag: Dag): DagGraph {
+ const groups = getDagTaskGroups(dag);
+ const downstreamTaskIds = new Map<string, Set<string>>();
+ const groupEdges = new Map<string, GroupEdgeSets>();
+ const link = (upstream: string, downstream: string): void => {
+ // Two arguments fed by the same upstream are one edge, as is an order-only
+ // edge redeclaring one the wiring already drew.
+ const edges = downstreamTaskIds.get(upstream) ?? new Set<string>();
+ edges.add(downstream);
+ downstreamTaskIds.set(upstream, edges);
+ };
+
+ for (const [taskId, inputs] of getDagTaskInputs(dag)) {
+ for (const value of Object.values(inputs)) {
+ if (isTaskRef(value)) link(value.taskId, taskId);
+ }
+ }
+ const orderEdges = getDagOrderEdges(dag);
+ for (const { upstream, downstream } of orderEdges) {
+ if (!groups.has(upstream) && !groups.has(downstream)) link(upstream,
downstream);
+ }
+
+ // Roots and leaves are read off the task-to-task graph, which holds every
+ // edge that can sit inside a group by now: wiring, and any order-only edge
+ // between two tasks. Python resolves them at `>>` time and so is equally
+ // order-sensitive, which is what keeps the two in step.
+ const ends = new GroupEnds(groups, downstreamTaskIds);
+
+ // Which group edges each endpoint has, for stepping over a group that holds
+ // no tasks.
+ const upstreamsOf = new Map<string, Set<string>>();
+ const downstreamsOf = new Map<string, Set<string>>();
+ for (const { upstream, downstream } of orderEdges) {
+ addTo(downstreamsOf, upstream, downstream);
+ addTo(upstreamsOf, downstream, upstream);
+ }
+
+ /**
+ * The tasks an edge endpoint stands for: the task itself, or a group's
leaves
+ * when it is upstream and its roots when it is downstream.
+ *
+ * A group holding no tasks has neither, so the edge steps over it and
+ * continues along the group edges beyond — `x >> empty >> y` still runs `y`
+ * after `x`, as Python's `find_leaves` walk does.
+ */
+ const tasksAt = (id: string, side: "upstream" | "downstream"): string[] => {
+ if (!groups.has(id)) return [id];
+ const own = side === "upstream" ? ends.leaves(id) : ends.roots(id);
+ return own.length > 0 ? own : tasksBeyond(id, side, new Set());
+ };
+ const tasksBeyond = (
+ id: string,
+ side: "upstream" | "downstream",
+ seen: Set<string>,
+ ): string[] => {
+ if (seen.has(id)) return [];
+ seen.add(id);
+ const next = side === "upstream" ? upstreamsOf.get(id) :
downstreamsOf.get(id);
+ return [...(next ?? [])].flatMap((other) => {
+ if (!groups.has(other)) return [other];
+ const own = side === "upstream" ? ends.leaves(other) : ends.roots(other);
+ return own.length > 0 ? own : tasksBeyond(other, side, seen);
+ });
+ };
+ const setsFor = (groupId: string): GroupEdgeSets => {
+ let sets = groupEdges.get(groupId);
+ if (!sets) {
+ sets = {
+ upstreamGroups: new Set(),
+ downstreamGroups: new Set(),
+ upstreamTasks: new Set(),
+ downstreamTasks: new Set(),
+ };
+ groupEdges.set(groupId, sets);
+ }
+ return sets;
+ };
+
+ for (const { upstream, downstream } of orderEdges) {
+ const upstreamIsGroup = groups.has(upstream);
+ const downstreamIsGroup = groups.has(downstream);
+ if (!upstreamIsGroup && !downstreamIsGroup) continue;
+
+ const from = tasksAt(upstream, "upstream");
+ const to = tasksAt(downstream, "downstream");
+ for (const tail of from) {
+ for (const head of to) link(tail, head);
+ }
+
+ if (downstreamIsGroup) {
+ const sets = setsFor(downstream);
+ for (const tail of from) sets.upstreamTasks.add(tail);
+ if (upstreamIsGroup) sets.upstreamGroups.add(upstream);
+ }
+ // Only a group whose downstream is a plain task records it as a task; when
+ // both ends are groups the pair is recorded as a group edge on this side
+ // and as the expanded tasks on the other, which is how Python leaves it.
+ if (upstreamIsGroup) {
+ const sets = setsFor(upstream);
+ if (downstreamIsGroup) sets.downstreamGroups.add(downstream);
+ else sets.downstreamTasks.add(downstream);
+ }
+ }
+ return { downstreamTaskIds, groupEdges };
+}
+
+function addTo(index: Map<string, Set<string>>, key: string, value: string):
void {
+ const existing = index.get(key) ?? new Set<string>();
+ existing.add(value);
+ index.set(key, existing);
+}
+
+/** The tasks an edge reaches when it points at a group, cached per group. */
+class GroupEnds {
+ readonly #groups: ReadonlyMap<string, TaskGroupRecord>;
+ readonly #downstream: ReadonlyMap<string, ReadonlySet<string>>;
+ readonly #members = new Map<string, Set<string>>();
+
+ constructor(
+ groups: ReadonlyMap<string, TaskGroupRecord>,
+ downstream: ReadonlyMap<string, ReadonlySet<string>>,
+ ) {
+ this.#groups = groups;
+ this.#downstream = downstream;
+ }
+
+ /** Tasks in the group with no upstream inside it: where an edge in arrives.
*/
+ roots(groupId: string): string[] {
+ const members = this.#membersOf(groupId);
+ const hasInternalUpstream = new Set<string>();
+ for (const [upstream, downstream] of this.#downstream) {
+ if (!members.has(upstream)) continue;
+ for (const task of downstream) {
+ if (members.has(task)) hasInternalUpstream.add(task);
+ }
+ }
+ return [...members].filter((task) => !hasInternalUpstream.has(task));
+ }
+
+ /** Tasks in the group with no downstream inside it: where an edge out
leaves. */
+ leaves(groupId: string): string[] {
+ const members = this.#membersOf(groupId);
+ return [...members].filter(
+ (task) => ![...(this.#downstream.get(task) ?? [])].some((other) =>
members.has(other)),
+ );
+ }
+
+ /** Every task the group holds, nested groups included. */
+ #membersOf(groupId: string): Set<string> {
+ const cached = this.#members.get(groupId);
+ if (cached) return cached;
+ const members = new Set<string>();
+ const pending = [groupId];
+ for (let i = 0; i < pending.length; i += 1) {
+ const group = this.#groups.get(pending[i]!);
+ if (!group) continue;
+ for (const taskId of group.taskIds) members.add(taskId);
+ pending.push(...group.childGroupIds);
+ }
+ this.#members.set(groupId, members);
+ return members;
+ }
+}
+
+/**
+ * Lower a `schedule` onto the timetable the scheduler reconstructs.
+ *
+ * Only the four schedules that map to a stock timetable are accepted. Anything
+ * else — an asset expression, a custom timetable — is a Python object the
+ * scheduler has to import, which a TypeScript bundle cannot name, so it is
+ * rejected here rather than serialized into a Dag that fails to deserialize.
+ */
+function serializeTimetable(schedule: unknown, dagId: string): SerializedValue
{
+ if (schedule === undefined || schedule === null) {
+ return simpleTimetable(NULL_TIMETABLE);
+ }
+ if (typeof schedule !== "string") {
+ throw new Error(
+ `schedule for Dag "${dagId}" must be "@once", "@continuous", or a cron
expression; ` +
+ `${describeType(schedule)} schedule names a Python object this SDK
cannot serialize`,
+ );
+ }
+ if (schedule.trim() === "") {
+ throw new Error(
+ `schedule for Dag "${dagId}" is empty; leave it unset for a Dag with no
schedule`,
+ );
+ }
+ if (schedule === "@once") return simpleTimetable(ONCE_TIMETABLE);
+ if (schedule === "@continuous") return simpleTimetable(CONTINUOUS_TIMETABLE);
+ const expression = CRON_PRESETS[schedule] ?? schedule;
+ if (!isCronExpression(expression)) {
+ throw new Error(
+ `schedule ${JSON.stringify(schedule)} for Dag "${dagId}" is not a cron
expression or a ` +
+ `preset (${Object.keys(CRON_PRESETS).join(", ")}, @once, @continuous);
a schedule the ` +
+ "scheduler cannot parse would leave the Dag unschedulable",
+ );
+ }
+ // TODO: honour [scheduler] create_cron_data_intervals, which switches Python
+ // to CronDataIntervalTimetable. A bundle cannot read airflow.cfg, so the
+ // supervisor has to send the flag first; tracked at
+ // https://github.com/apache/airflow/issues/67938
+ return {
+ __type: CRON_TIMETABLE,
+ __var: { expression, timezone: "UTC", interval: 0, run_immediately: false
},
+ };
+}
+
+/**
+ * Presets expanded the way `CronMixin.__init__` expands them, so the
serialized
+ * expression is the one Python records — which the Dag's summary and its hash
+ * are both taken from. Mirrors `airflow.utils.dates.cron_presets`.
+ */
+const CRON_PRESETS: Readonly<Record<string, string>> = {
+ "@hourly": "0 * * * *",
+ "@daily": "0 0 * * *",
+ "@weekly": "0 0 * * 0",
+ "@monthly": "0 0 1 * *",
+ "@quarterly": "0 0 1 */3 *",
+ "@yearly": "0 0 1 1 *",
+};
+
+/**
+ * Whether `expression` has the shape croniter accepts: five or six
+ * space-separated fields of cron characters.
+ *
+ * A shape check, not a parse: croniter validates the ranges, and repeating
that
+ * here would be a second implementation to keep in step. What it does catch is
+ * prose — `"every tuesday"` — which would otherwise be written into a Dag that
+ * the scheduler then fails to build a timetable for.
+ */
+function isCronExpression(expression: string): boolean {
+ const fields = expression.trim().split(/\s+/);
+ if (fields.length !== 5 && fields.length !== 6) return false;
+ return fields.every(
+ (field) => /^[\d*,\-/?LW#]+$/i.test(field) ||
/^[A-Z]{3}(-[A-Z]{3})?$/i.test(field),
+ );
+}
+
+function simpleTimetable(type: string): SerializedValue {
+ return { __type: type, __var: {} };
+}
+
+/**
+ * Write the fields a spec set onto `data`, skipping any left at its schema
+ * default — Python's serializer omits what the scheduler re-derives.
+ */
+function applySchemaFields(
+ data: Record<string, SerializedValue>,
+ spec: object,
+ rules: FieldRules,
+ label: string,
+): void {
+ const values = spec as Record<string, unknown>;
+ for (const [name, field] of Object.entries(rules.fields)) {
+ // A virtual field names a key the serializer derives rather than writes;
+ // `schedule` becomes `timetable`.
+ if (field.virtual || rules.omitted.has(field.key)) continue;
+ const value = values[name];
+ if (value === undefined || value === field.default) continue;
+ const encoded = encodeField(field, value, `${name} for ${label}`);
+ data[field.key] = rules.decorated.has(field.key) ? encoded :
unwrapTypeEncoding(encoded);
+ }
+}
+
+/** Encode one authoring value as the schema's type for that field. Rejects a
+ * value of the wrong type: this is the last point before the scheduler, and a
+ * mistyped field would otherwise surface as an unreadable Dag. */
+function encodeField(field: SchemaField, value: unknown, label: string):
SerializedValue {
+ switch (field.type) {
+ case "string":
+ if (typeof value !== "string") throw typeError(label, "a string", value);
+ return value;
+ case "boolean":
+ if (typeof value !== "boolean") throw typeError(label, "a boolean",
value);
+ return value;
+ case "number":
+ if (typeof value !== "number" || !Number.isFinite(value)) {
+ throw typeError(label, "a finite number", value);
+ }
+ return value;
+ case "timedelta":
+ if (typeof value !== "number" || !Number.isFinite(value)) {
+ throw typeError(label, "a duration in seconds", value);
+ }
+ return { __type: "timedelta", __var: value };
+ case "datetime":
+ if (!(value instanceof Date) || Number.isNaN(value.getTime())) {
+ throw typeError(label, "a valid Date", value);
+ }
+ return serializeValue(value);
+ case "string[]": {
+ if (!Array.isArray(value) || value.some((item) => typeof item !==
"string")) {
+ throw typeError(label, "an array of strings", value);
+ }
+ // Python holds these in a set, so duplicates collapse and the order is
+ // the sorted one that keeps a Dag's hash stable across runs.
+ return serializeValue(new Set(value as string[]));
+ }
+ }
+}
+
+function typeError(label: string, expected: string, value: unknown): Error {
+ return new Error(`${label} must be ${expected}, not ${describeType(value)}`);
+}
+
+function describeType(value: unknown): string {
+ if (value === null) return "null";
+ if (Array.isArray(value)) return "an array";
+ if (typeof value === "number" && !Number.isFinite(value)) return
String(value);
+ const noun =
+ typeof value === "object" && !isPlainRecord(value) ? getClassName(value) :
typeof value;
+ return `${/^[aeiou]/i.test(noun) ? "an" : "a"} ${noun}`;
+}
+
+function getClassName(value: object): string {
+ const prototype = Object.getPrototypeOf(value) as { constructor?: { name?:
string } } | null;
+ return prototype?.constructor?.name || "object";
+}
+
+/**
+ * Encode a value the way `BaseSerialization.serialize()` does.
+ *
+ * A duration has no distinct runtime type in TypeScript — it is a number of
+ * seconds — so `timedelta` is applied by {@link encodeField} from the schema
+ * rather than inferred here.
+ */
+export function serializeValue(value: unknown): SerializedValue {
+ if (value === null || value === undefined) return null;
+ if (typeof value === "string" || typeof value === "boolean") return value;
+ if (typeof value === "number") {
+ if (!Number.isFinite(value)) {
+ throw new Error(`Cannot serialize the non-finite number
${String(value)}`);
+ }
+ return value;
+ }
+ if (value instanceof Date) {
+ if (Number.isNaN(value.getTime())) throw new Error("Cannot serialize an
invalid Date");
+ return { __type: "datetime", __var: value.getTime() / 1000 };
+ }
+ if (value instanceof Set) {
+ return {
+ __type: "set",
+ __var: [...value].map(serializeValue).sort(compareSerialized),
+ };
+ }
+ if (Array.isArray(value)) return value.map(serializeValue);
+ if (value instanceof Map) {
+ return { __type: "dict", __var: serializeEntries(value.entries()) };
+ }
+ if (typeof value === "object") {
+ return { __type: "dict", __var: serializeEntries(Object.entries(value)) };
+ }
+ throw new Error(`Cannot serialize a ${typeof value}`);
+}
+
+function serializeEntries(entries: Iterable<[unknown, unknown]>):
Record<string, SerializedValue> {
+ const encoded: Record<string, SerializedValue> = {};
+ for (const [key, item] of entries) {
+ encoded[String(key)] = serializeValue(item);
+ }
+ return encoded;
+}
+
+// Python sorts a set's members before writing them; JSON's default sort is
+// lexicographic on the string form, which matches for the string sets this
+// SDK produces and stays total for anything else.
+function compareSerialized(left: SerializedValue, right: SerializedValue):
number {
+ const a = typeof left === "string" ? left : JSON.stringify(left);
+ const b = typeof right === "string" ? right : JSON.stringify(right);
+ return a < b ? -1 : a > b ? 1 : 0;
+}
+
+/**
+ * Strip the type encoding from a non-decorated field, as Python's
+ * `serialize_to_json` does: it serializes every field, then keeps only the
+ * `__var` of the ones outside its decorated set.
+ */
+export function unwrapTypeEncoding(value: SerializedValue): SerializedValue {
+ if (!isTypeEncoded(value)) return value;
+ return value.__var;
+}
+
+function isTypeEncoded(value: SerializedValue): value is TypeEncoded &
SerializedValue {
+ return (
+ typeof value === "object" &&
+ value !== null &&
+ !Array.isArray(value) &&
+ "__type" in value &&
+ "__var" in value
+ );
+}
+
+/** Where the Dag file sits inside its bundle, as Airflow records it. */
+export function computeRelativeFileloc(fileloc: string, bundlePath: string):
string {
+ if (!fileloc) return "";
+ if (!bundlePath) return ".";
+ const result = relativePath(bundlePath, fileloc);
+ return result === "" ? "." : result;
+}
diff --git a/ts-sdk/src/sdk/dag.ts b/ts-sdk/src/sdk/dag.ts
index 8f7a2042dea..f643e290f45 100644
--- a/ts-sdk/src/sdk/dag.ts
+++ b/ts-sdk/src/sdk/dag.ts
@@ -33,7 +33,8 @@ import type { JsonValue } from "./client-types.js";
import { getCurrentModuleSource } from "./module-source.js";
import type { TaskFunction } from "./task.js";
-function isPlainRecord(value: unknown): value is Record<string, unknown> {
+/** Internal: whether `value` is an object literal, not an array or a class
instance. */
+export function isPlainRecord(value: unknown): value is Record<string,
unknown> {
if (typeof value !== "object" || value === null || Array.isArray(value))
return false;
const prototype = Object.getPrototypeOf(value);
return prototype === Object.prototype || prototype === null;
@@ -84,7 +85,7 @@ function kindOf(value: object): string {
return prototype?.constructor?.name ?? "value";
}
-const DAG_SPEC_KEYS: ReadonlySet<string> = new
Set(Object.keys(DAG_SCHEMA_FIELDS));
+const DAG_SPEC_KEYS: ReadonlySet<string> = new
Set([...Object.keys(DAG_SCHEMA_FIELDS), "queue"]);
// `taskId` is hand-written rather than generated: the schema's task_id is
// serializer-owned, and this is the authoring surface's own way to set it.
const TASK_SPEC_KEYS: ReadonlySet<string> = new
Set([...Object.keys(TASK_SCHEMA_FIELDS), "taskId"]);
@@ -97,10 +98,20 @@ const TASK_SPEC_KEYS: ReadonlySet<string> = new
Set([...Object.keys(TASK_SCHEMA_
* break a call site. An unknown key is rejected, so a misspelled field is an
* error rather than a Dag that quietly ignores it.
*
- * Setting a field records it. A Dag declared in TypeScript is not served to
- * Airflow yet, so nothing reads it.
+ * `queue` is the one hand-written field: Airflow's 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 rather than on each task.
*/
-export type DagSpec = GeneratedDagFields;
+export interface DagSpec extends GeneratedDagFields {
+ /**
+ * Queue the Dag's tasks run on, unless a task names its own.
+ *
+ * A native Dag's tasks are executed by the Node coordinator, which the
+ * deployment's `queue_to_coordinator` maps a queue to, so this is what
+ * routes them there. `queue` on a {@link TaskSpec} wins for that task.
+ */
+ readonly queue?: string;
+}
/**
* Task-level options: the retries, the pool, the trigger rule, and the rest of
@@ -287,8 +298,8 @@ function nodeId(node: Node): string | undefined {
return undefined;
}
-/** Whether `value` is a TaskRef returned by any copy of this package. */
-function isTaskRef(value: unknown): value is TaskRef {
+/** Internal: whether `value` is a TaskRef returned by any copy of this
package. */
+export function isTaskRef(value: unknown): value is TaskRef {
return hasBrand(value, "TaskRef");
}
diff --git a/ts-sdk/tests/conformance/serialize_typescript.ts
b/ts-sdk/tests/conformance/serialize_typescript.ts
new file mode 100644
index 00000000000..948d12a1688
--- /dev/null
+++ b/ts-sdk/tests/conformance/serialize_typescript.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.
+ */
+
+// Serializes the Dags of scripts/ci/lang_sdk_serialization/test_dags.yaml
with this SDK, as
+// serialize_python.py there does with Airflow's serializer, and writes them
keyed by Dag id. The
+// check-ts-sdk-serialization-conformance hook has that directory's compare.py
run it as:
+//
+// pnpm --dir ts-sdk exec tsx tests/conformance/serialize_typescript.ts
<test_dags.yaml> <output.json>
+
+import { readFileSync, writeFileSync } from "node:fs";
+import { argv } from "node:process";
+
+import { parse, type ScalarTag } from "yaml";
+
+import { serializeDag } from "../../src/coordinator/serde.js";
+import {
+ DAG_SCHEMA_FIELDS,
+ SERIALIZATION_VERSION,
+ TASK_SCHEMA_FIELDS,
+ type SchemaField,
+} from "../../src/generated/dag-schema-fields.js";
+import {
+ Dag,
+ type DagSpec,
+ type TaskGroupRef,
+ type TaskRef,
+ type TaskSpec,
+} from "../../src/index.js";
+
+interface TaskCase {
+ readonly task_id: string;
+ readonly group?: string;
+ readonly upstream?: readonly string[];
+ readonly spec?: Readonly<Record<string, unknown>>;
+}
+
+interface DagCase {
+ readonly dag_id: string;
+ readonly spec?: Readonly<Record<string, unknown>>;
+ readonly groups?: readonly string[];
+ readonly tasks: readonly TaskCase[];
+ readonly order_edges?: readonly (readonly [string, string])[];
+}
+
+// A duration is already a number of seconds in this SDK's authoring API.
+const TAGS: ScalarTag[] = [
+ { tag: "!datetime", resolve: (value) => new Date(value) },
+ { tag: "!timedelta", resolve: (value) => Number(value) },
+];
+
+/**
+ * Map each schema key of a generated table to the authoring name it is set by.
+ *
+ * A virtual field is keyed by its authoring name instead, since its schema
key names what the
+ * serializer derives: `schedule` is spelled that way in both Python's `DAG()`
and this SDK, but
+ * its schema key is `timetable`.
+ */
+function authoringNames(table: Readonly<Record<string, SchemaField>>):
Map<string, string> {
+ return new Map(
+ Object.entries(table).map(([name, field]) => [field.virtual ? name :
field.key, name]),
+ );
+}
+
+const DAG_NAMES = authoringNames(DAG_SCHEMA_FIELDS);
+const TASK_NAMES = authoringNames(TASK_SCHEMA_FIELDS);
+
+function toAuthoringSpec(
+ spec: Readonly<Record<string, unknown>>,
+ names: Map<string, string>,
+ label: string,
+): Record<string, unknown> {
+ const authoring: Record<string, unknown> = {};
+ for (const [key, value] of Object.entries(spec)) {
+ const name = names.get(key);
+ if (name === undefined) throw new Error(`${label}: "${key}" is not an
authoring field`);
+ authoring[name] = value;
+ }
+ return authoring;
+}
+
+function qualifiedTaskId(task: TaskCase): string {
+ return task.group === undefined ? task.task_id :
`${task.group}.${task.task_id}`;
+}
+
+function buildDag(dagCase: DagCase): Dag {
+ const dag = new Dag(
+ dagCase.dag_id,
+ toAuthoringSpec(dagCase.spec ?? {}, DAG_NAMES, dagCase.dag_id) as DagSpec,
+ );
+ // A group id is fully qualified, so its parent is whatever comes before the
last dot.
+ const groups = new Map<string, TaskGroupRef>();
+ for (const groupId of dagCase.groups ?? []) {
+ const cut = groupId.lastIndexOf(".");
+ const scope = cut === -1 ? dag : groups.get(groupId.slice(0, cut))!;
+ groups.set(groupId, scope.taskGroup(groupId.slice(cut + 1)));
+ }
+
+ const factories = new Map<string, (inputs: Record<string, TaskRef>) =>
TaskRef>();
+ for (const task of dagCase.tasks) {
+ const label = `${dagCase.dag_id}.${qualifiedTaskId(task)}`;
+ const spec = toAuthoringSpec(task.spec ?? {}, TASK_NAMES, label) as
TaskSpec;
+ const scope = task.group === undefined ? dag : groups.get(task.group)!;
+ factories.set(
+ qualifiedTaskId(task),
+ scope.task(task.task_id, async (_inputs: Record<string, TaskRef>) =>
undefined, spec),
+ );
+ }
+ const refs = new Map<string, TaskRef>();
+ for (const task of dagCase.tasks) {
+ const inputs: Record<string, TaskRef> = {};
+ for (const upstream of task.upstream ?? []) {
+ const ref = refs.get(upstream);
+ if (!ref)
+ throw new Error(`${dagCase.dag_id}: "${upstream}" has to come before
its downstream`);
+ inputs[`from_${upstream}`] = ref;
+ }
+ refs.set(qualifiedTaskId(task),
factories.get(qualifiedTaskId(task))!(inputs));
+ }
+ for (const [upstream, downstream] of dagCase.order_edges ?? []) {
+ const from = groups.get(upstream) ?? refs.get(upstream);
+ const to = groups.get(downstream) ?? refs.get(downstream);
+ if (!from || !to)
+ throw new Error(`${dagCase.dag_id}: no node "${from ? downstream :
upstream}"`);
+ from.before(to);
+ }
+ return dag;
+}
+
+const [testDags, output] = argv.slice(2);
+if (testDags === undefined || output === undefined) {
+ throw new Error("Usage: serialize_typescript.ts <test_dags.yaml>
<output.json>");
+}
+const { dags } = parse(readFileSync(testDags, "utf-8"), { customTags: TAGS })
as {
+ dags: DagCase[];
+};
+// compare.py leaves fileloc out, as it names the file a Dag was declared in,
so any bundle path
+// does. Airflow still needs one to load the Dag.
+const serialized = Object.fromEntries(
+ dags.map((dagCase) => [
+ dagCase.dag_id,
+ {
+ __version: SERIALIZATION_VERSION,
+ dag: serializeDag(buildDag(dagCase), "/bundles/app/bundle.mjs",
"bundle.mjs"),
+ },
+ ]),
+);
+writeFileSync(output, `${JSON.stringify(serialized, null, 2)}\n`);
diff --git a/ts-sdk/tests/coordinator/arg-binding.test.ts
b/ts-sdk/tests/coordinator/arg-binding.test.ts
index 1e3f1161d34..663462613e8 100644
--- a/ts-sdk/tests/coordinator/arg-binding.test.ts
+++ b/ts-sdk/tests/coordinator/arg-binding.test.ts
@@ -23,7 +23,9 @@ import { foldArgName, resolveArgs, type BoundArgs } from
"../../src/coordinator/
import type { CoordinatorClient, XComEntry } from
"../../src/coordinator/client.js";
import type { LogChannel } from "../../src/coordinator/log-channel.js";
import type { ArgBindings } from "../../src/generated/supervisor.js";
+import { Bundle } from "../../src/sdk/bundle.js";
import type { GetXComOpts } from "../../src/sdk/client-types.js";
+import { Dag } from "../../src/sdk/dag.js";
function literal(name: string, value: unknown, extra: Record<string, unknown>
= {}) {
return { name, kind: "literal" as const, value, ...extra };
@@ -76,6 +78,23 @@ async function bind(
return { ...bound, warning, pulls };
}
+describe("a native task's bound arguments", () => {
+ it("reach a named handler as the object it destructures", async () => {
+ const dag = new Dag("d");
+ const seen: unknown[] = [];
+ const store = dag.task("store", async ({ rows }: { rows: number }) => {
+ seen.push(rows);
+ });
+ store({ rows: 1 });
+ const handler = new Bundle(dag).getTaskHandler("d", "store")!;
+
+ const { args } = await bind([literal("rows", 7)]);
+ await handler(args as never);
+
+ expect(seen).toEqual([7]);
+ });
+});
+
describe("foldArgName", () => {
it.each([
["region_code", "regioncode"],
diff --git a/ts-sdk/tests/coordinator/serde.test.ts
b/ts-sdk/tests/coordinator/serde.test.ts
new file mode 100644
index 00000000000..f2ddc7ad22c
--- /dev/null
+++ b/ts-sdk/tests/coordinator/serde.test.ts
@@ -0,0 +1,846 @@
+/*!
+ * 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.
+ */
+
+import { describe, expect, it } from "vitest";
+
+import {
+ computeRelativeFileloc,
+ serializeDag,
+ serializeValue,
+ unwrapTypeEncoding,
+} from "../../src/coordinator/serde.js";
+import {
+ Dag,
+ type DagSpec,
+ type TaskGroupRef,
+ type TaskRef,
+ type TaskSpec,
+} from "../../src/sdk/dag.js";
+
+type Json = Record<string, unknown>;
+type WiredArgs = Record<string, TaskRef>;
+
+/** Declare a task and place it, optionally behind some upstreams. */
+function place(
+ scope: Dag | TaskGroupRef,
+ taskId: string,
+ upstream: readonly TaskRef[] = [],
+ spec?: TaskSpec,
+) {
+ const factory = scope.task(taskId, async (_args: WiredArgs) => undefined,
spec);
+ const inputs: Record<string, TaskRef> = {};
+ upstream.forEach((ref, index) => {
+ inputs[`in_${index}`] = ref;
+ });
+ return factory(inputs);
+}
+
+/** A one-task Dag, serialized. */
+function serializeWith(spec: DagSpec, taskSpec?: TaskSpec): Json {
+ const dag = new Dag("d", spec);
+ place(dag, "t", [], taskSpec);
+ return serializeDag(dag, "/bundles/app/bundle.mjs", "bundle.mjs") as Json;
+}
+
+function taskVar(serialized: Json, index = 0): Json {
+ const tasks = serialized["tasks"] as { __type: string; __var: Json }[];
+ expect(tasks[index]!.__type).toBe("operator");
+ return tasks[index]!.__var;
+}
+
+/** Serialized tasks keyed by task id. */
+function taskMap(serialized: Json): Map<string, Json> {
+ const tasks = serialized["tasks"] as { __var: Json }[];
+ return new Map(tasks.map(({ __var }) => [__var["task_id"] as string,
__var]));
+}
+
+describe("serializeDag", () => {
+ it("writes the fields Airflow always expects", () => {
+ const serialized = serializeWith({});
+
+ expect(serialized).toMatchObject({
+ dag_id: "d",
+ fileloc: "/bundles/app/bundle.mjs",
+ relative_fileloc: "bundle.mjs",
+ timezone: "UTC",
+ dag_dependencies: [],
+ edge_info: {},
+ params: [],
+ deadline: null,
+ allowed_run_types: null,
+ });
+ });
+
+ describe("a field Python reads from config when the Dag leaves it unset", ()
=> {
+ const keys = [
+ "max_active_tasks",
+ "max_active_runs",
+ "max_consecutive_failed_dag_runs",
+ "catchup",
+ "disable_bundle_versioning",
+ ];
+
+ it("is left out, for Airflow to fill in from its own config", () => {
+ const serialized = serializeWith({});
+
+ expect(keys.filter((key) => key in serialized)).toEqual([]);
+ });
+
+ it("is written when the Dag sets it, even to the stock default", () => {
+ const serialized = serializeWith({
+ maxActiveTasks: 16,
+ maxActiveRuns: 16,
+ maxConsecutiveFailedDagRuns: 0,
+ catchup: false,
+ disableBundleVersioning: false,
+ });
+
+ expect(serialized).toMatchObject({
+ max_active_tasks: 16,
+ max_active_runs: 16,
+ max_consecutive_failed_dag_runs: 0,
+ catchup: false,
+ disable_bundle_versioning: false,
+ });
+ });
+ });
+
+ it("puts every task in the flat root group", () => {
+ const dag = new Dag("d");
+ const first = place(dag, "first");
+ place(dag, "second", [first]);
+
+ expect(serializeDag(dag, "", ".")["task_group"]).toEqual({
+ _group_id: null,
+ group_display_name: "",
+ prefix_group_id: true,
+ tooltip: "",
+ ui_color: "CornflowerBlue",
+ ui_fgcolor: "#000",
+ children: { first: ["operator", "first"], second: ["operator", "second"]
},
+ upstream_group_ids: [],
+ downstream_group_ids: [],
+ upstream_task_ids: [],
+ downstream_task_ids: [],
+ });
+ });
+
+ it("identifies a task as TypeScript rather than as a Python operator", () =>
{
+ expect(taskVar(serializeWith({}))).toEqual({
+ task_id: "t",
+ task_type: "TypeScriptOperator",
+ _task_module: "airflow.sdk.coordinators.node",
+ language: "typescript",
+ template_fields: [],
+ is_stub: true,
+ });
+ });
+
+ describe("queue", () => {
+ it("gives every task the Dag's queue, which is what routes them here", ()
=> {
+ const dag = new Dag("d", { queue: "typescript" });
+ place(dag, "one");
+ place(dag, "two");
+
+ const tasks = taskMap(serializeDag(dag, "", ".") as Json);
+ expect([...tasks.values()].map((task) => task["queue"])).toEqual([
+ "typescript",
+ "typescript",
+ ]);
+ });
+
+ it("lets a task name its own queue instead", () => {
+ const dag = new Dag("d", { queue: "typescript" });
+ place(dag, "light");
+ place(dag, "heavy", [], { queue: "typescript_large" });
+
+ const tasks = taskMap(serializeDag(dag, "", ".") as Json);
+ expect(tasks.get("light")?.["queue"]).toBe("typescript");
+ expect(tasks.get("heavy")?.["queue"]).toBe("typescript_large");
+ });
+
+ it("writes no queue when neither the Dag nor the task names one", () => {
+ expect(taskVar(serializeWith({}))).not.toHaveProperty("queue");
+ });
+
+ it("omits a queue that is already the schema default", () => {
+ // The scheduler re-derives it, as it does any other defaulted field.
+ const dag = new Dag("d", { queue: "default" });
+ place(dag, "one");
+
+ expect(taskMap(serializeDag(dag, "", ".") as
Json).get("one")).not.toHaveProperty("queue");
+ });
+
+ it("is not written onto the Dag itself, which has no queue field", () => {
+ expect(serializeWith({ queue: "typescript"
})).not.toHaveProperty("queue");
+ });
+ });
+
+ describe("arg bindings", () => {
+ it("binds an upstream reference as the xcom the API server resolves", ()
=> {
+ const dag = new Dag("d");
+ const extracted = place(dag, "extract");
+ const transform = dag.task("transform", async (_: { extracted: unknown
}) => undefined);
+ transform({ extracted });
+
+ expect(taskMap(serializeDag(dag, "", ".") as
Json).get("transform")).toMatchObject({
+ is_stub: true,
+ _arg_bindings: [{ name: "extracted", kind: "xcom", task_id: "extract"
}],
+ });
+ });
+
+ it("binds a literal as plain JSON, as Python writes it", () => {
+ const dag = new Dag("d");
+ const transform = dag.task(
+ "transform",
+ async (_: {
+ regionCode: string;
+ limits: number[];
+ options: { retry: { count: number } };
+ rows: { n: number }[];
+ cursor: string | null;
+ }) => undefined,
+ );
+ transform({
+ regionCode: "us",
+ limits: [1, 2],
+ options: { retry: { count: 2 } },
+ rows: [{ n: 1 }],
+ cursor: null,
+ });
+
+ expect(
+ taskMap(serializeDag(dag, "", ".") as
Json).get("transform")!["_arg_bindings"],
+ ).toEqual([
+ { name: "regionCode", kind: "literal", value: "us" },
+ { name: "limits", kind: "literal", value: [1, 2] },
+ { name: "options", kind: "literal", value: { retry: { count: 2 } } },
+ { name: "rows", kind: "literal", value: [{ n: 1 }] },
+ { name: "cursor", kind: "literal", value: null },
+ ]);
+ });
+
+ it.each([
+ ["a Date", new Date(0)],
+ ["a Map", new Map([["a", 1]])],
+ ["a Set", new Set([1])],
+ ["a Point", new (class Point {})()],
+ ["an Int8Array", new Int8Array(1)],
+ ["Infinity", Number.POSITIVE_INFINITY],
+ ])("rejects a literal holding %s, which JSON cannot carry", (description,
value) => {
+ const dag = new Dag("d");
+ const transform = dag.task("transform", async (_: { options: unknown })
=> undefined);
+ transform({ options: { items: [value] } } as never);
+
+ expect(() => serializeDag(dag, "", ".")).toThrowError(
+ `Input "options" of task "transform" of Dag "d" holds ${description},
which JSON cannot carry`,
+ );
+ });
+
+ it("keeps the order the call named the arguments in", () => {
+ const dag = new Dag("d");
+ const north = place(dag, "north");
+ const south = place(dag, "south");
+ const summarize = dag.task(
+ "summarize",
+ async (_: { south: unknown; north: unknown; label: string }) =>
undefined,
+ );
+ summarize({ south, north, label: "both" });
+
+ const bindings = taskMap(serializeDag(dag, "", ".") as
Json).get("summarize")![
+ "_arg_bindings"
+ ] as { name: string }[];
+ expect(bindings.map(({ name }) => name)).toEqual(["south", "north",
"label"]);
+ });
+
+ it("leaves the key out for a task called with no arguments", () => {
+ expect(taskVar(serializeWith({}))).not.toHaveProperty("_arg_bindings");
+ });
+ });
+
+ describe("timetable", () => {
+ it.each([
+ ["unset", undefined, { __type:
"airflow.timetables.simple.NullTimetable", __var: {} }],
+ ["@once", "@once", { __type: "airflow.timetables.simple.OnceTimetable",
__var: {} }],
+ [
+ "@continuous",
+ "@continuous",
+ { __type: "airflow.timetables.simple.ContinuousTimetable", __var: {} },
+ ],
+ [
+ "a cron expression",
+ "0 3 * * *",
+ {
+ __type: "airflow.timetables.trigger.CronTriggerTimetable",
+ __var: {
+ expression: "0 3 * * *",
+ timezone: "UTC",
+ interval: 0,
+ run_immediately: false,
+ },
+ },
+ ],
+ ])("maps %s onto the matching timetable", (_name, schedule, expected) => {
+ expect(serializeWith({ schedule })["timetable"]).toEqual(expected);
+ });
+
+ it.each([
+ ["@hourly", "0 * * * *"],
+ ["@daily", "0 0 * * *"],
+ ["@weekly", "0 0 * * 0"],
+ ["@monthly", "0 0 1 * *"],
+ ["@quarterly", "0 0 1 */3 *"],
+ ["@yearly", "0 0 1 1 *"],
+ ])("expands the preset %s the way Python records it", (schedule,
expression) => {
+ // Both spellings rebuild the same timetable, but the expression is what
+ // the Dag's summary and its hash are taken from, so writing the preset
+ // verbatim would disagree with the same Dag written in Python.
+ const timetable = serializeWith({ schedule })["timetable"] as Json;
+
expect(timetable["__type"]).toBe("airflow.timetables.trigger.CronTriggerTimetable");
+ expect((timetable["__var"] as Json)["expression"]).toBe(expression);
+ });
+
+ it.each([
+ ["prose", "every tuesday"],
+ ["too few fields", "0 0 * *"],
+ ["too many fields", "0 0 * * * * *"],
+ ["an unknown preset", "@fortnightly"],
+ ])("rejects %s rather than writing a schedule nothing can parse", (_label,
schedule) => {
+ expect(() => serializeWith({ schedule })).toThrowError(
+ /is not a cron expression or a preset/,
+ );
+ });
+
+ it.each([
+ ["a six-field expression", "0 0 0 * * *"],
+ ["named weekdays", "0 0 * * MON-FRI"],
+ ["a step", "*/15 * * * *"],
+ ])("accepts %s", (_label, schedule) => {
+ expect(() => serializeWith({ schedule })).not.toThrow();
+ });
+
+ it.each([
+ ["an asset expression", { assets: ["s3://bucket/key"] }, /an object
schedule names a Python/],
+ ["a number", 86400, /a number schedule names a Python/],
+ ["an empty string", "", /schedule for Dag "d" is empty/],
+ ["a blank string", " ", /schedule for Dag "d" is empty/],
+ ])("rejects %s", (_name, schedule, expected) => {
+ expect(() => serializeWith({ schedule } as
DagSpec)).toThrowError(expected);
+ });
+ });
+
+ describe("non-decorated fields", () => {
+ it("writes them as bare values, without the type encoding", () => {
+ const serialized = serializeWith({
+ startDate: new Date("2026-01-01T00:00:00Z"),
+ endDate: new Date("2026-12-31T23:30:15Z"),
+ dagrunTimeout: 300,
+ tags: ["gamma", "alpha"],
+ description: "demo",
+ });
+
+ expect(serialized).toMatchObject({
+ start_date: 1767225600,
+ end_date: 1798759815,
+ dagrun_timeout: 300,
+ // Python holds tags in a set and writes them sorted.
+ tags: ["alpha", "gamma"],
+ description: "demo",
+ });
+ });
+
+ it("keeps only what a decorated field would keep", () => {
+ // The two halves of the split: everything above went through
+ // serializeValue and then lost its wrapper, because no authoring field
is
+ // in Python's decorated set. A decorated field would stop at the first.
+ const wrapped = serializeValue(new Date("2026-01-01T00:00:00Z"));
+ expect(wrapped).toEqual({ __type: "datetime", __var: 1767225600 });
+ expect(unwrapTypeEncoding(wrapped)).toBe(1767225600);
+ });
+
+ it("collapses duplicate tags, as the set Python holds them in does", () =>
{
+ expect(serializeWith({ tags: ["b", "a", "b"] })["tags"]).toEqual(["a",
"b"]);
+ });
+ });
+
+ describe("omit-if-default", () => {
+ it("omits a Dag field left at its schema default", () => {
+ const serialized = serializeWith({ failFast: false,
renderTemplateAsNativeObj: false });
+ expect(serialized).not.toHaveProperty("fail_fast");
+ expect(serialized).not.toHaveProperty("render_template_as_native_obj");
+ });
+
+ it("writes a Dag field that differs from its schema default", () => {
+ expect(serializeWith({ failFast: true })).toMatchObject({ fail_fast:
true });
+ });
+
+ it("omits task fields left at their schema defaults", () => {
+ const task = taskVar(
+ serializeWith(
+ {},
+ { retries: 0, queue: "default", pool: "default_pool", retryDelay:
300, owner: "airflow" },
+ ),
+ );
+ for (const key of ["retries", "queue", "pool", "retry_delay", "owner"]) {
+ expect(task).not.toHaveProperty(key);
+ }
+ });
+
+ it("writes task fields that differ from their schema defaults", () => {
+ const task = taskVar(
+ serializeWith(
+ {},
+ { retries: 2, queue: "typescript", retryDelay: 600,
executionTimeout: 5 },
+ ),
+ );
+ expect(task).toMatchObject({
+ retries: 2,
+ queue: "typescript",
+ retry_delay: 600,
+ execution_timeout: 5,
+ });
+ });
+
+ it("never writes the email flags, which have no recipient to reach", () =>
{
+ const task = taskVar(serializeWith({}, { emailOnFailure: false,
emailOnRetry: false }));
+ expect(task).not.toHaveProperty("email_on_failure");
+ expect(task).not.toHaveProperty("email_on_retry");
+ });
+ });
+
+ describe("downstream_task_ids", () => {
+ it("inverts the recorded wiring, sorted", () => {
+ const dag = new Dag("d");
+ const root = place(dag, "root");
+ const right = place(dag, "right", [root]);
+ const left = place(dag, "left", [root]);
+ place(dag, "join", [right, left]);
+
+ const serialized = serializeDag(dag, "", ".") as Json;
+ expect(taskVar(serialized, 0)["downstream_task_ids"]).toEqual(["left",
"right"]);
+ expect(taskVar(serialized, 1)["downstream_task_ids"]).toEqual(["join"]);
+ expect(taskVar(serialized, 3)).not.toHaveProperty("downstream_task_ids");
+ });
+
+ it("counts one edge when two arguments come from the same upstream", () =>
{
+ const dag = new Dag("d");
+ const upstream = place(dag, "up");
+ place(dag, "down", [upstream, upstream]);
+
+ expect(taskVar(serializeDag(dag, "", ".") as Json,
0)["downstream_task_ids"]).toEqual([
+ "down",
+ ]);
+ });
+
+ it("ignores a literal argument that looks like a reference", () => {
+ const dag = new Dag("d");
+ place(dag, "up");
+ const factory = dag.task("down", async (_args: WiredArgs) => undefined);
+ factory({ config: { dagId: "d", taskId: "up" } as unknown as TaskRef });
+
+ expect(taskVar(serializeDag(dag, "", ".") as Json,
0)).not.toHaveProperty(
+ "downstream_task_ids",
+ );
+ });
+ });
+
+ describe("order-only edges", () => {
+ it("writes an edge between two tasks onto the upstream task", () => {
+ const dag = new Dag("d");
+ const loaded = place(dag, "load");
+ const cleaned = place(dag, "cleanup");
+ loaded.before(cleaned);
+
+ expect(taskVar(serializeDag(dag, "", ".") as Json,
0)["downstream_task_ids"]).toEqual([
+ "cleanup",
+ ]);
+ });
+
+ it("joins the wiring on one graph, since the serialized Dag has only one",
() => {
+ const dag = new Dag("d");
+ const extracted = place(dag, "extract");
+ place(dag, "transform", [extracted]);
+ const cleaned = place(dag, "cleanup");
+ extracted.before(cleaned);
+
+ expect(taskVar(serializeDag(dag, "", ".") as Json,
0)["downstream_task_ids"]).toEqual([
+ "cleanup",
+ "transform",
+ ]);
+ });
+
+ it("counts one edge when the wiring already drew it", () => {
+ const dag = new Dag("d");
+ const extracted = place(dag, "extract");
+ const transformed = place(dag, "transform", [extracted]);
+ extracted.before(transformed);
+
+ expect(taskVar(serializeDag(dag, "", ".") as Json,
0)["downstream_task_ids"]).toEqual([
+ "transform",
+ ]);
+ });
+ });
+
+ describe("task_group", () => {
+ /** The Dag's root task group, as the serialized payload carries it. */
+ function rootGroup(dag: Dag): Json {
+ return (serializeDag(dag, "", ".") as Json)["task_group"] as Json;
+ }
+
+ it("puts a Dag's top-level tasks in the root group", () => {
+ const dag = new Dag("d");
+ place(dag, "alpha");
+ place(dag, "beta");
+
+ expect(rootGroup(dag)).toMatchObject({
+ _group_id: null,
+ prefix_group_id: true,
+ children: { alpha: ["operator", "alpha"], beta: ["operator", "beta"] },
+ });
+ });
+
+ it("nests a group by embedding its own object, as Python does", () => {
+ const dag = new Dag("d");
+ place(dag, "top");
+ const staging = dag.taskGroup("staging");
+ place(staging, "stage");
+
+ const root = rootGroup(dag);
+ expect(root["children"]).toMatchObject({ top: ["operator", "top"] });
+ const [kind, nested] = (root["children"] as Json)["staging"] as [string,
Json];
+ expect(kind).toBe("taskgroup");
+ expect(nested).toMatchObject({
+ _group_id: "staging",
+ children: { "staging.stage": ["operator", "staging.stage"] },
+ });
+ });
+
+ it("carries prefixGroupId off, with the ids it holds as written", () => {
+ const dag = new Dag("d");
+ place(dag.taskGroup("checks", { prefixGroupId: false }), "nulls");
+
+ const [, checks] = (rootGroup(dag)["children"] as Json)["checks"] as
[string, Json];
+ expect(checks).toMatchObject({
+ _group_id: "checks",
+ prefix_group_id: false,
+ children: { nulls: ["operator", "nulls"] },
+ });
+ });
+
+ it("keeps a grouped task out of the root group's children", () => {
+ const dag = new Dag("d");
+ place(dag.taskGroup("staging"), "stage");
+
+ expect(Object.keys(rootGroup(dag)["children"] as
Json)).toEqual(["staging"]);
+ });
+
+ it("nests to any depth", () => {
+ const dag = new Dag("d");
+ const outer = dag.taskGroup("outer");
+ place(outer.taskGroup("inner"), "deep");
+
+ const root = rootGroup(dag);
+ const [, outerGroup] = (root["children"] as Json)["outer"] as [string,
Json];
+ const [, innerGroup] = (outerGroup["children"] as Json)["outer.inner"]
as [string, Json];
+ // The local segment, as Python records it: the qualified id is rebuilt
+ // from where the group sits in the tree.
+ expect(innerGroup).toMatchObject({
+ _group_id: "inner",
+ children: { "outer.inner.deep": ["operator", "outer.inner.deep"] },
+ });
+ });
+
+ it("records an edge between a group and a task on the group", () => {
+ const dag = new Dag("d");
+ const staging = dag.taskGroup("staging");
+ place(staging, "stage");
+ const loaded = place(dag, "load");
+ staging.before(loaded);
+
+ const [, group] = (rootGroup(dag)["children"] as Json)["staging"] as
[string, Json];
+ expect(group).toMatchObject({
+ downstream_task_ids: ["load"],
+ upstream_task_ids: [],
+ downstream_group_ids: [],
+ upstream_group_ids: [],
+ });
+ // And expanded onto the task graph, which is the only one the scheduler
+ // reads: the group's leaves carry the edge to the task it points at.
+ expect(taskVar(serializeDag(dag, "", ".") as Json,
0)["downstream_task_ids"]).toEqual([
+ "load",
+ ]);
+ });
+
+ it("records an edge between two groups on both of them", () => {
+ const dag = new Dag("d");
+ const first = dag.taskGroup("first");
+ place(first, "a");
+ const second = dag.taskGroup("second");
+ place(second, "b");
+ first.before(second);
+
+ const children = rootGroup(dag)["children"] as Json;
+ const [, firstGroup] = children["first"] as [string, Json];
+ const [, secondGroup] = children["second"] as [string, Json];
+ expect(firstGroup).toMatchObject({
+ downstream_group_ids: ["second"],
+ upstream_group_ids: [],
+ });
+ expect(secondGroup).toMatchObject({
+ upstream_group_ids: ["first"],
+ downstream_group_ids: [],
+ });
+ });
+
+ it("records an upstream task on the group it points at", () => {
+ const dag = new Dag("d");
+ const extracted = place(dag, "extract");
+ const staging = dag.taskGroup("staging");
+ place(staging, "stage");
+ extracted.before(staging);
+
+ const [, group] = (rootGroup(dag)["children"] as Json)["staging"] as
[string, Json];
+ expect(group).toMatchObject({ upstream_task_ids: ["extract"],
downstream_task_ids: [] });
+ });
+
+ it("expands an edge into a group onto its roots, not every task it holds",
() => {
+ const dag = new Dag("d");
+ const extracted = place(dag, "extract");
+ const staging = dag.taskGroup("staging");
+ const staged = place(staging, "stage");
+ place(staging, "check", [staged]);
+ extracted.before(staging);
+
+ // "staging.check" already runs after "staging.stage", so the edge in
+ // reaches only the task that starts the group.
+ expect(taskVar(serializeDag(dag, "", ".") as Json,
0)["downstream_task_ids"]).toEqual([
+ "staging.stage",
+ ]);
+ });
+
+ it("expands an edge out of a group from its leaves", () => {
+ const dag = new Dag("d");
+ const staging = dag.taskGroup("staging");
+ const staged = place(staging, "stage");
+ place(staging, "check", [staged]);
+ const loaded = place(dag, "load");
+ staging.before(loaded);
+
+ const serialized = serializeDag(dag, "", ".") as Json;
+ // "staging.stage" keeps its own intra-group edge and nothing more: the
+ // edge out of the group leaves from the task that finishes it.
+ expect(taskVar(serialized,
0)["downstream_task_ids"]).toEqual(["staging.check"]);
+ expect(taskVar(serialized, 1)["downstream_task_ids"]).toEqual(["load"]);
+ });
+
+ it("joins one group's leaves to the next group's roots", () => {
+ const dag = new Dag("d");
+ const first = dag.taskGroup("first");
+ const firstHead = place(first, "head");
+ place(first, "tail", [firstHead]);
+ const second = dag.taskGroup("second");
+ const secondHead = place(second, "head");
+ place(second, "tail", [secondHead]);
+ first.before(second);
+
+ const serialized = serializeDag(dag, "", ".") as Json;
+ expect(taskVar(serialized,
1)["downstream_task_ids"]).toEqual(["second.head"]);
+ const children = rootGroup(dag)["children"] as Json;
+ const [, secondGroup] = children["second"] as [string, Json];
+ // The downstream group records the upstream's leaves as well as the
+ // group itself; the upstream group records only the group edge.
+ expect(secondGroup).toMatchObject({
+ upstream_group_ids: ["first"],
+ upstream_task_ids: ["first.tail"],
+ });
+ const [, firstGroup] = children["first"] as [string, Json];
+ expect(firstGroup).toMatchObject({
+ downstream_group_ids: ["second"],
+ downstream_task_ids: [],
+ });
+ });
+
+ it("reaches a task held by a nested group", () => {
+ const dag = new Dag("d");
+ const extracted = place(dag, "extract");
+ const outer = dag.taskGroup("outer");
+ place(outer.taskGroup("inner"), "deep");
+ extracted.before(outer);
+
+ expect(taskVar(serializeDag(dag, "", ".") as Json,
0)["downstream_task_ids"]).toEqual([
+ "outer.inner.deep",
+ ]);
+ });
+
+ it("carries an empty group, which holds nothing and constrains nothing",
() => {
+ const dag = new Dag("d");
+ place(dag, "solo");
+ dag.taskGroup("empty");
+
+ const [, group] = (rootGroup(dag)["children"] as Json)["empty"] as
[string, Json];
+ expect(group).toMatchObject({ _group_id: "empty", children: {} });
+ });
+
+ it("steps over an empty group so a chain through it still orders its
ends", () => {
+ // An empty group has no roots and no leaves, so the two edges used to
+ // expand into nothing at all and `after` ran beside `before`. Python
+ // bridges the gap by walking up through the group's own upstreams.
+ const dag = new Dag("d");
+ const before = place(dag, "before");
+ const after = place(dag, "after");
+ const empty = dag.taskGroup("empty");
+ before.before(empty);
+ empty.before(after);
+
+ const tasks = taskMap(serializeDag(dag, "", ".") as Json);
+ expect(tasks.get("before")?.["downstream_task_ids"]).toEqual(["after"]);
+ });
+
+ it("steps over a chain of empty groups", () => {
+ const dag = new Dag("d");
+ const before = place(dag, "before");
+ const after = place(dag, "after");
+ const first = dag.taskGroup("first");
+ const second = dag.taskGroup("second");
+ before.before(first);
+ first.before(second);
+ second.before(after);
+
+ const tasks = taskMap(serializeDag(dag, "", ".") as Json);
+ expect(tasks.get("before")?.["downstream_task_ids"]).toEqual(["after"]);
+ });
+
+ it("draws nothing when an empty group has no other side to reach", () => {
+ const dag = new Dag("d");
+ const before = place(dag, "before");
+ before.before(dag.taskGroup("empty"));
+
+ const tasks = taskMap(serializeDag(dag, "", ".") as Json);
+ expect(tasks.get("before")).not.toHaveProperty("downstream_task_ids");
+ });
+ });
+
+ describe("rejects a value the schema cannot carry", () => {
+ it.each([
+ ["startDate", { startDate: "2026-01-01" }, /startDate for Dag "d" must
be a valid Date/],
+ ["an invalid Date", { startDate: new Date("nope") }, /must be a valid
Date/],
+ ["description", { description: 7 }, /description for Dag "d" must be a
string/],
+ ["catchup", { catchup: "yes" }, /catchup for Dag "d" must be a boolean/],
+ [
+ "maxActiveRuns",
+ { maxActiveRuns: "3" },
+ /maxActiveRuns for Dag "d" must be a finite number/,
+ ],
+ ["dagrunTimeout", { dagrunTimeout: Infinity }, /must be a duration in
seconds/],
+ ["tags", { tags: ["a", 2] }, /tags for Dag "d" must be an array of
strings/],
+ ])("on %s", (_name, spec, expected) => {
+ expect(() => serializeWith(spec as DagSpec)).toThrowError(expected);
+ });
+
+ it("names the task a bad task field belongs to", () => {
+ expect(() => serializeWith({}, { retries: "two" } as unknown as
TaskSpec)).toThrowError(
+ /retries for task "t" of Dag "d" must be a finite number/,
+ );
+ });
+ });
+});
+
+describe("serializeValue", () => {
+ it.each([
+ ["a string", "x", "x"],
+ ["a boolean", true, true],
+ ["a number", 1.5, 1.5],
+ ["null", null, null],
+ ["undefined", undefined, null],
+ ["a list, without a wrapper", [1, "a"], [1, "a"]],
+ ])("passes %s through", (_name, value, expected) => {
+ expect(serializeValue(value)).toEqual(expected);
+ });
+
+ it("encodes a Date as fractional epoch seconds", () => {
+ expect(serializeValue(new Date("2026-01-01T00:00:00.500Z"))).toEqual({
+ __type: "datetime",
+ __var: 1767225600.5,
+ });
+ });
+
+ it("encodes a Set as a sorted list", () => {
+ expect(serializeValue(new Set(["gamma", "alpha", "beta"]))).toEqual({
+ __type: "set",
+ __var: ["alpha", "beta", "gamma"],
+ });
+ });
+
+ it.each([
+ ["an object", { b: 1, a: "x" }],
+ [
+ "a Map",
+ new Map<string, unknown>([
+ ["b", 1],
+ ["a", "x"],
+ ]),
+ ],
+ ])("encodes %s as a dict", (_name, value) => {
+ expect(serializeValue(value)).toEqual({ __type: "dict", __var: { b: 1, a:
"x" } });
+ });
+
+ it("recurses into nested values", () => {
+ expect(serializeValue({ when: new Date("2026-01-01T00:00:00Z"), items: [{
n: 1 }] })).toEqual({
+ __type: "dict",
+ __var: {
+ when: { __type: "datetime", __var: 1767225600 },
+ items: [{ __type: "dict", __var: { n: 1 } }],
+ },
+ });
+ });
+
+ it.each([
+ ["a non-finite number", Number.NaN, /non-finite number/],
+ ["an invalid Date", new Date("nope"), /invalid Date/],
+ ["a function", () => undefined, /Cannot serialize a function/],
+ ])("rejects %s", (_name, value, expected) => {
+ expect(() => serializeValue(value)).toThrowError(expected);
+ });
+});
+
+describe("unwrapTypeEncoding", () => {
+ it("takes the __var of an encoded value", () => {
+ expect(unwrapTypeEncoding({ __type: "timedelta", __var: 300 })).toBe(300);
+ });
+
+ it.each([
+ ["a primitive", 5],
+ ["a list", [1, 2]],
+ ["an object that is not encoded", { __var: 1 }],
+ ])("leaves %s alone", (_name, value) => {
+ expect(unwrapTypeEncoding(value)).toEqual(value);
+ });
+});
+
+describe("computeRelativeFileloc", () => {
+ it.each([
+ ["a file inside the bundle", "/bundles/app/dags/bundle.mjs",
"/bundles/app", "dags/bundle.mjs"],
+ ["a file at the bundle root", "/bundles/app/bundle.mjs", "/bundles/app",
"bundle.mjs"],
+ ["a file that is the bundle", "/bundles/app", "/bundles/app", "."],
+ ["an unknown bundle path", "/bundles/app/bundle.mjs", "", "."],
+ ["an unknown file", "", "/bundles/app", ""],
+ ])("resolves %s", (_name, fileloc, bundlePath, expected) => {
+ expect(computeRelativeFileloc(fileloc, bundlePath)).toBe(expected);
+ });
+});
diff --git a/ts-sdk/tests/public-api.test.ts b/ts-sdk/tests/public-api.test.ts
index 769a8342b0e..b6cb9541627 100644
--- a/ts-sdk/tests/public-api.test.ts
+++ b/ts-sdk/tests/public-api.test.ts
@@ -360,6 +360,10 @@ describe("public API", () => {
expectTypeOf<Dag["taskIds"]>().toEqualTypeOf<readonly string[]>();
// Both specs are all-optional, so `{}` stays assignable and a field the
// schema gains later cannot break a call site.
+ // `queue` is the one Dag field Airflow's schema does not have: a native
+ // Dag's tasks all run on the same coordinator, so the queue that routes
+ // them there belongs on the Dag.
+ expectTypeOf<DagSpec["queue"]>().toEqualTypeOf<string | undefined>();
const emptyDagSpec: DagSpec = {};
const emptyTaskSpec: TaskSpec = {};
expect([emptyDagSpec, emptyTaskSpec]).toEqual([{}, {}]);