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 ea1f4723e4f TS SDK: generate Dag and task options from the
serialization schema (#73436)
ea1f4723e4f is described below
commit ea1f4723e4f8fe969c2b73b7fd979531132484f0
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Thu Sep 24 13:07:51 2026 +0800
TS SDK: generate Dag and task options from the serialization schema (#73436)
A Dag authored in TypeScript has to describe itself with the same options
Airflow's scheduler reads, and that vocabulary lives in airflow-core's Dag
serialization schema. Hand-listing those options in the SDK would leave two
copies to keep aligned by hand, and the copy the SDK ships would quietly rot
every time core adds, renames, or drops a key.
Deriving them instead makes the drift visible: a new scalar key shows up in
the
regenerated diff or fails generation until someone says why the SDK omits
it,
and a curated entry that outlives the key it names fails the same way. The
serialization version is the one thing the schema cannot supply, so a core
bump
now fails a test in airflow-core rather than producing Dags stamped with a
version core no longer writes.
This follows the policy the Go SDK's TaskSpec generator set and ADR-0008
records, so the three language SDKs expose the same Dag surface for the same
documented reasons.
---
.pre-commit-config.yaml | 10 +
.../language-sdks/typescript.rst | 9 +-
.../test_ts_sdk_serialization_version.py | 53 +++
ts-sdk/.pre-commit-config.yaml | 12 +
ts-sdk/package.json | 3 +-
ts-sdk/schema/dag-schema.json | 474 +++++++++++++++++++
ts-sdk/scripts/ci/prek/check_dag_schema.py | 65 +++
ts-sdk/scripts/ci/prek/sync_dag_schema.py | 65 +++
ts-sdk/scripts/generate-dag-schema.mjs | 512 +++++++++++++++++++++
ts-sdk/src/generated/dag-schema-fields.ts | 215 +++++++++
ts-sdk/src/sdk/dag.ts | 163 +++++--
ts-sdk/tests/public-api.test.ts | 48 +-
ts-sdk/tests/scripts/generate-dag-schema.test.ts | 419 +++++++++++++++++
ts-sdk/tests/sdk/dag.test.ts | 150 ++++--
ts-sdk/tsconfig.json | 7 +-
15 files changed, 2094 insertions(+), 111 deletions(-)
diff --git a/.pre-commit-config.yaml b/.pre-commit-config.yaml
index 9b4b05b804a..e4c5472d7b1 100644
--- a/.pre-commit-config.yaml
+++ b/.pre-commit-config.yaml
@@ -308,6 +308,16 @@ repos:
additional_dependencies: ['PyYAML>=6.0', 'rich>=13.6.0']
pass_filenames: false
require_serial: true
+ - id: sync-ts-sdk-dag-schema
+ name: Sync TypeScript SDK Dag serialization schema with airflow-core
+ description: "Copy airflow-core's serialization schema when TS SDK's
vendored dag-schema.json drifts"
+ entry: ./ts-sdk/scripts/ci/prek/sync_dag_schema.py
+ language: python
+ pass_filenames: false
+ files: >
+ (?x)
+ ^airflow-core/src/airflow/serialization/schema\.json$|
+ ^ts-sdk/schema/dag-schema\.json$
- 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 aad721a1a4f..d3ddca0c9b4 100644
--- a/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
+++ b/airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst
@@ -243,9 +243,9 @@ task instance.
Declaring a Dag in TypeScript
-----------------------------
-A ``Dag`` is declared on this side rather than in Python: its tasks and the
edges between them are
-written in TypeScript. The surface is still growing, so a Dag declared this
way is not served to
-Airflow yet.
+A ``Dag`` is declared on this side rather than in Python: its schedule, its
tasks, their options and
+the edges between them are all written in TypeScript. The surface is still
growing, so a Dag declared
+this way is not served to Airflow yet.
``dag.task(taskId, handler)`` returns a *factory*. Calling it places the task
in the Dag and supplies the
handler's arguments, so the call graph is the task graph:
@@ -277,7 +277,8 @@ argument itself: one buried inside an array or an object is
a literal, and draws
Every task has to be called exactly once. An uncalled task fails when the Dag
is read, so none can be
left out of the graph by accident.
-``new Dag`` and ``dag.task`` also take a trailing ``spec`` object that is not
used yet; do not set it.
+``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.
Writing tasks
-------------
diff --git
a/airflow-core/tests/unit/serialization/test_ts_sdk_serialization_version.py
b/airflow-core/tests/unit/serialization/test_ts_sdk_serialization_version.py
new file mode 100644
index 00000000000..87834b72e88
--- /dev/null
+++ b/airflow-core/tests/unit/serialization/test_ts_sdk_serialization_version.py
@@ -0,0 +1,53 @@
+# 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 re
+
+import pytest
+
+from airflow.serialization.serialized_objects import DagSerialization
+
+from tests_common.pytest_plugin import AIRFLOW_ROOT_PATH
+
+GENERATED_FIELDS_PATH = AIRFLOW_ROOT_PATH / "ts-sdk" / "src" / "generated" /
"dag-schema-fields.ts"
+
+SERIALIZATION_VERSION_PATTERN = re.compile(
+ r"^export const SERIALIZATION_VERSION = (?P<version>\d+);$", re.MULTILINE
+)
+
+
[email protected](
+ not GENERATED_FIELDS_PATH.is_file(),
+ reason="TypeScript SDK sources are absent (source distribution build)",
+)
+def test_ts_sdk_pins_the_current_serialization_version():
+ """
+ The schema constrains ``__version`` only to a positive integer, so the
TypeScript SDK
+ pins it by hand. Without this check, bumping the serializer would leave
the SDK emitting
+ Dags stamped with a version core no longer writes.
+ """
+ match =
SERIALIZATION_VERSION_PATTERN.search(GENERATED_FIELDS_PATH.read_text(encoding="utf-8"))
+ assert match is not None, f"No SERIALIZATION_VERSION export found in
{GENERATED_FIELDS_PATH}"
+
+ assert int(match.group("version")) == DagSerialization.SERIALIZER_VERSION,
(
+ f"{GENERATED_FIELDS_PATH.relative_to(AIRFLOW_ROOT_PATH)} pins
serialization version "
+ f"{match.group('version')}, but DagSerialization.SERIALIZER_VERSION is
"
+ f"{DagSerialization.SERIALIZER_VERSION}. Update SERIALIZATION_VERSION
in "
+ "ts-sdk/scripts/generate-dag-schema.mjs, re-run `pnpm run
generate:dag-schema`, and make "
+ "sure the SDK still emits Dags core can read."
+ )
diff --git a/ts-sdk/.pre-commit-config.yaml b/ts-sdk/.pre-commit-config.yaml
index 508a913a7ee..ad2bcb99c17 100644
--- a/ts-sdk/.pre-commit-config.yaml
+++ b/ts-sdk/.pre-commit-config.yaml
@@ -49,6 +49,18 @@ repos:
additional_dependencies: ['[email protected]']
pass_filenames: false
require_serial: true
+ - id: check-ts-sdk-dag-schema
+ name: Check TypeScript SDK Dag schema fields are up to date
+ entry: ./scripts/ci/prek/check_dag_schema.py
+ language: node
+ files: >
+ (?x)
+ ^src/generated/dag-schema-fields\.ts$|
+ ^schema/dag-schema\.json$|
+ ^scripts/generate-dag-schema\.mjs$
+ 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 b09a714106a..9ad8c80b044 100644
--- a/ts-sdk/package.json
+++ b/ts-sdk/package.json
@@ -51,7 +51,8 @@
"build": "pnpm run clean && tsc -p tsconfig.build.json",
"verify:package": "node scripts/verify-package.mjs",
"prepack": "pnpm run build",
- "generate:supervisor": "node scripts/generate-supervisor.mjs"
+ "generate:supervisor": "node scripts/generate-supervisor.mjs",
+ "generate:dag-schema": "node scripts/generate-dag-schema.mjs"
},
"keywords": [
"airflow",
diff --git a/ts-sdk/schema/dag-schema.json b/ts-sdk/schema/dag-schema.json
new file mode 100644
index 00000000000..fa149734b7b
--- /dev/null
+++ b/ts-sdk/schema/dag-schema.json
@@ -0,0 +1,474 @@
+{
+ "$schema": "http://json-schema.org/draft-07/schema#",
+ "$id": "https://airflow.apache.com/schemas/serialized-dags.json",
+ "definitions": {
+ "datetime": {
+ "description": "A date time, stored as fractional seconds since the
epoch",
+ "type": "number"
+ },
+ "timedelta": {
+ "type": "number",
+ "minimum": 0
+ },
+ "typed_timedelta": {
+ "type": "object",
+ "properties": {
+ "__type": {
+ "type": "string",
+ "const": "timedelta"
+ },
+ "__var": { "$ref": "#/definitions/timedelta" }
+ },
+ "required": [
+ "__type",
+ "__var"
+ ],
+ "additionalProperties": false
+ },
+ "typed_relativedelta": {
+ "type": "object",
+ "description": "A dateutil.relativedelta.relativedelta object",
+ "properties": {
+ "__type": {
+ "type": "string",
+ "const": "relativedelta"
+ },
+ "__var": {
+ "type": "object",
+ "properties": {
+ "weekday": {
+ "type": "array",
+ "items": { "type": "integer" },
+ "minItems": 1,
+ "maxItems": 2
+ }
+ },
+ "additionalProperties": { "type": "integer" }
+ }
+ }
+ },
+ "timezone": {
+ "anyOf": [
+ { "type": "string" },
+ { "type": "integer" }
+ ]
+ },
+ "asset_definition": {
+ "type": "object",
+ "properties": {
+ "uri": { "type": "string" },
+ "name": { "type": "string" },
+ "group": { "type": "string" },
+ "extra": {
+ "anyOf": [
+ {"type": "null"},
+ { "$ref": "#/definitions/dict" }
+ ]
+ },
+ "watchers": {
+ "type": "array",
+ "items": { "$ref": "#/definitions/trigger" }
+ }
+ },
+ "required": [ "uri", "extra" ]
+ },
+ "asset": {
+ "type": "object",
+ "properties": {
+ "uri": { "type": "string" },
+ "extra": {
+ "anyOf": [
+ {"type": "null"},
+ { "$ref": "#/definitions/dict" }
+ ]
+ }
+ },
+ "required": [ "uri", "extra" ]
+ },
+ "typed_asset": {
+ "type": "object",
+ "properties": {
+ "__type": {
+ "type": "string",
+ "constant": "asset"
+ },
+ "__var": { "$ref": "#/definitions/asset" }
+ },
+ "required": [
+ "__type",
+ "__var"
+ ],
+ "additionalProperties": false
+ },
+ "typed_asset_cond": {
+ "type": "object",
+ "properties": {
+ "__type": {
+ "anyOf": [{
+ "type": "string",
+ "constant": "asset_or"
+ },
+ {
+ "type": "string",
+ "constant": "asset_and"
+ }
+ ]
+ },
+ "__var": {
+ "type": "array",
+ "items": {
+ "anyOf": [
+ {"$ref": "#/definitions/typed_asset"},
+ { "$ref": "#/definitions/typed_asset_cond"}
+ ]
+ }
+ }
+ },
+ "required": [
+ "__type",
+ "__var"
+ ],
+ "additionalProperties": false
+ },
+ "trigger": {
+ "type": "object",
+ "properties": {
+ "classpath": { "type": "string" },
+ "kwargs": { "$ref": "#/definitions/dict" }
+ },
+ "required": [ "classpath", "kwargs" ]
+ },
+ "dict": {
+ "description": "A python dictionary containing values of any type",
+ "type": "object"
+ },
+ "arg_binding": {
+ "$comment": "One captured TaskFlow call argument of a @task.stub task.
Materialized directly by _serialize_node, so it stays plain JSON with no
{__type, __var} encoding. The object stays open so future binding fields keep
validating on older cores",
+ "type": "object",
+ "properties": {
+ "name": { "type": "string" },
+ "kind": { "type": "string", "enum": [ "xcom", "literal" ] },
+ "value_schema": { "type": "object" },
+ "task_id": { "type": "string" },
+ "value": {},
+ "from_default": { "type": "boolean" }
+ },
+ "required": [ "name", "kind" ]
+ },
+ "color": {
+ "type": "string",
+ "pattern": "^#[a-fA-F0-9]{3,6}$"
+ },
+ "extra_links": {
+ "type": "array",
+ "items": {
+ "type": "object",
+ "minProperties": 1,
+ "maxProperties": 1
+ }
+ },
+ "dag_dependencies": {
+ "type": "array",
+ "items": {
+ "type": "object"
+ }
+ },
+ "dag": {
+ "type": "object",
+ "properties": {
+ "params": { "$ref": "#/definitions/params" },
+ "dag_id": { "type": "string" },
+ "tasks": { "$ref": "#/definitions/tasks" },
+ "timezone": { "$ref": "#/definitions/timezone" },
+ "owner_links": { "type": "object" },
+ "timetable": {
+ "type": "object",
+ "properties": {
+ "type": { "type": "string" },
+ "value": { "$ref": "#/definitions/dict" }
+ }
+ },
+ "catchup": { "type": "boolean" },
+ "allowed_run_types": {
+ "anyOf": [
+ { "type": "array", "items": { "type": "string" } },
+ { "type": "null" }
+ ]
+ },
+ "fail_fast": { "type": "boolean", "default": false },
+ "fileloc": { "type" : "string"},
+ "relative_fileloc": { "type" : "string"},
+ "bundle_name": { "anyOf": [{ "type": "null" }, { "type": "string" }] },
+ "_processor_dags_folder": {
+ "anyOf": [
+ { "type": "null" },
+ { "type": "string" }
+ ]
+ },
+ "dag_display_name": { "type" : "string"},
+ "description": { "type" : "string"},
+ "deadline": {
+ "anyOf": [
+ { "$ref": "#/definitions/dict" },
+ {
+ "type": "array",
+ "items": { "$ref": "#/definitions/dict" }
+ },
+ {
+ "$comment": "Once persisted, a Dag's deadline alerts live
as rows in the deadline_alert table and the serialized Dag keeps only a list of
UUID strings referencing them (see
SerializedDagModel._generate_deadline_uuids). This branch lets the stored form
validate at any lifecycle stage, not only before the dict->UUID rewrite.",
+ "type": "array",
+ "items": { "type": "string",
+ "pattern":
"^[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}$"
+ }
+ },
+ { "type": "null" }
+ ]
+ },
+ "_concurrency": { "type" : "number"},
+ "max_active_tasks": { "type" : "number" },
+ "max_active_runs": { "type" : "number" },
+ "max_consecutive_failed_dag_runs": { "type" : "number" },
+ "default_args": { "$ref": "#/definitions/dict" },
+ "start_date": { "$ref": "#/definitions/datetime" },
+ "end_date": { "$ref": "#/definitions/datetime" },
+ "dagrun_timeout": { "$ref": "#/definitions/timedelta" },
+ "doc_md": { "type" : "string"},
+ "access_control": {"$ref": "#/definitions/dict" },
+ "is_paused_upon_creation": { "type": "boolean" },
+ "has_on_success_callback": { "type": "boolean", "default": false },
+ "has_on_failure_callback": { "type": "boolean", "default": false },
+ "render_template_as_native_obj": { "type": "boolean", "default":
false },
+ "tags": { "type": "array" },
+ "task_group": {"anyOf": [
+ { "type": "null" },
+ { "$ref": "#/definitions/task_group" }
+ ]},
+ "edge_info": { "$ref": "#/definitions/edge_info" },
+ "dag_dependencies": { "$ref": "#/definitions/dag_dependencies" },
+ "disable_bundle_versioning": {"type": "boolean" },
+ "rerun_with_latest_version": {"type": ["boolean", "null"], "default":
null}
+ },
+ "required": [
+ "dag_id",
+ "fileloc",
+ "tasks"
+ ],
+ "additionalProperties": false
+ },
+ "tasks": {
+ "type": "array",
+ "additionalProperties": { "$ref": "#/definitions/operator" }
+ },
+ "params": {
+ "type": "array",
+ "prefixItems": [
+ { "type": "string" },
+ { "$ref": "#/definitions/param" }
+ ],
+ "unevaluatedItems": false
+ },
+ "param": {
+ "$comment": "A param for a dag / operator",
+ "type": "object",
+ "required": [
+ "__class",
+ "default"
+ ],
+ "properties": {
+ "__class": { "type": "string" },
+ "default": {},
+ "description": {"anyOf": [{"type":"string"}, {"type":"null"}]},
+ "schema": { "$ref": "#/definitions/dict" }
+ }
+ },
+ "operator": {
+ "$comment": "A task/operator in a DAG",
+ "type": "object",
+ "required": [
+ "task_type",
+ "_task_module",
+ "task_id",
+ "ui_color",
+ "ui_fgcolor",
+ "template_fields"
+ ],
+ "properties": {
+ "task_type": { "type": "string", "default": "BaseOperator"},
+ "_task_module": { "type": "string" },
+ "_operator_extra_links": { "$ref": "#/definitions/extra_links" },
+ "task_id": { "type": "string" },
+ "_task_display_name": { "type": "string" },
+ "owner": { "type": "string", "default": "airflow" },
+ "start_date": { "$ref": "#/definitions/datetime" },
+ "end_date": { "$ref": "#/definitions/datetime" },
+ "trigger_rule": { "type": "string", "default": "all_success" },
+ "depends_on_past": { "type": "boolean", "default": false },
+ "ignore_first_depends_on_past": { "type": "boolean", "default": false
},
+ "wait_for_past_depends_before_skipping": { "type": "boolean",
"default": false },
+ "wait_for_downstream": { "type": "boolean", "default": false },
+ "retries": { "type": "number", "default": 0 },
+ "queue": { "type": "string", "default": "default" },
+ "pool": { "type": "string", "default": "default_pool" },
+ "pool_slots": { "type": "number", "default": 1 },
+ "execution_timeout": { "$ref": "#/definitions/timedelta" },
+ "retry_delay": { "$ref": "#/definitions/timedelta", "default": 300.0 },
+ "retry_exponential_backoff": { "type": "number", "default": 0 },
+ "max_retry_delay": { "$ref": "#/definitions/timedelta" },
+ "params": { "$ref": "#/definitions/params" },
+ "priority_weight": { "type": "number", "default": 1 },
+ "weight_rule": { "type": "string", "default": "downstream" },
+ "executor": { "type": "string" },
+ "executor_config": { "$ref": "#/definitions/dict" },
+ "do_xcom_push": { "type": "boolean", "default": true },
+ "email_on_failure": { "type": "boolean", "default": true },
+ "email_on_retry": { "type": "boolean", "default": true },
+ "ui_color": { "type": "string", "default": "#fff" },
+ "ui_fgcolor": { "type": "string", "default": "#000" },
+ "template_fields": {
+ "type": "array",
+ "items": { "type": "string" },
+ "default": []
+ },
+ "template_ext": {"type": "array", "default": []},
+ "template_fields_renderers": {"$ref": "#/definitions/dict", "default":
{}},
+ "downstream_task_ids": {
+ "type": "array",
+ "items": { "type": "string" },
+ "default": []
+ },
+ "doc": { "type": "string" },
+ "doc_md": { "type": "string" },
+ "doc_json": { "type": "string" },
+ "doc_yaml": { "type": "string" },
+ "doc_rst": { "type": "string" },
+ "_logger_name": { "type": "string" },
+ "_needs_expansion": { "type": "boolean"},
+ "_is_mapped": { "const": true, "$comment": "only present when True",
"default": false },
+ "_is_sensor": { "const": true, "$comment": "only present when True",
"default": false },
+ "partial_kwargs": { "type": "object" },
+ "_disallow_kwargs_override": { "type": "boolean"},
+ "_expand_input_attr": { "type": "string" },
+ "map_index_template": { "type": "string" },
+ "allow_nested_operators": { "type": "boolean", "default": true },
+ "render_template_as_native_obj": { "anyOf": [{"type": "boolean"},
{"type": "null"}], "default": null },
+ "inlets": {"type": "array", "default": []},
+ "outlets": {"type": "array", "default": []},
+ "has_on_execute_callback": {"type": "boolean", "default": false},
+ "has_on_failure_callback": {"type": "boolean", "default": false},
+ "has_on_skipped_callback": {"type": "boolean", "default": false},
+ "has_on_success_callback": {"type": "boolean", "default": false},
+ "has_on_retry_callback": {"type": "boolean", "default": false},
+ "multiple_outputs": {"type": "boolean", "default": false},
+ "start_from_trigger": {"type": "boolean", "default": false},
+ "start_trigger_args": {"type": "object", "default": null},
+ "is_setup": {"type": "boolean", "default": false},
+ "is_teardown": {"type": "boolean", "default": false},
+ "on_failure_fail_dagrun": {"type": "boolean", "default": false},
+ "max_active_tis_per_dag": {"type": "integer"},
+ "max_active_tis_per_dagrun": {"type": "integer"},
+ "_arg_bindings": {
+ "$comment": "Only present on @task.stub tasks called with TaskFlow
arguments",
+ "type": "array",
+ "items": { "$ref": "#/definitions/arg_binding" }
+ }
+ },
+ "dependencies": {
+ "expand_input": ["partial_kwargs", "_is_mapped"],
+ "partial_kwargs": ["expand_input", "_is_mapped"],
+ "_is_mapped": ["expand_input", "partial_kwargs"]
+ },
+ "additionalProperties": true
+ },
+ "task_group": {
+ "$comment": "A TaskGroup containing tasks",
+ "type": "object",
+ "required": [
+ "_group_id",
+ "group_display_name",
+ "prefix_group_id",
+ "children",
+ "tooltip",
+ "ui_color",
+ "ui_fgcolor",
+ "upstream_group_ids",
+ "downstream_group_ids",
+ "upstream_task_ids",
+ "downstream_task_ids"
+ ],
+ "properties": {
+ "_group_id": {"anyOf": [{"type": "null"}, { "type": "string" }]},
+ "group_display_name": {"type": "string" },
+ "is_mapped": { "type": "boolean" },
+ "prefix_group_id": { "type": "boolean" },
+ "children": { "$ref": "#/definitions/dict" },
+ "tooltip": { "type": "string" },
+ "doc_md": {
+ "anyOf": [
+ { "type": "string" },
+ { "type": "null" }
+ ]},
+ "ui_color": { "type": "string" },
+ "ui_fgcolor": { "type": "string" },
+ "upstream_group_ids": {
+ "type": "array",
+ "items": { "type": "string" }
+ },
+ "downstream_group_ids": {
+ "type": "array",
+ "items": { "type": "string" }
+ },
+ "upstream_task_ids": {
+ "type": "array",
+ "items": { "type": "string" }
+ },
+ "downstream_task_ids": {
+ "type": "array",
+ "items": { "type": "string" }
+ }
+ },
+ "additionalProperties": false
+ },
+ "edge_info": {
+ "$comment": "Metadata about DAG edges",
+ "type": "object",
+ "additionalProperties": {
+ "type": "object",
+ "additionalProperties": {
+ "type": "object",
+ "properties": {
+ "label": { "type": "string" }
+ },
+ "required": ["label"],
+ "additionalProperties": false
+ }
+ }
+ }
+ },
+
+ "type": "object",
+ "allOf": [
+ {
+ "type": "object",
+ "properties": {
+ "__version": {
+ "type": "integer",
+ "exclusiveMinimum": 0
+ },
+ "dag": { "$ref": "#/definitions/dag" },
+ "client_defaults": {
+ "type": "object",
+ "description": "SDK-specific default values that differ from schema
defaults",
+ "properties": {
+ "tasks": {
+ "type": "object",
+ "description": "Task-level default overrides"
+ }
+ },
+ "additionalProperties": false
+ }
+ },
+ "additionalProperties": false,
+ "required": [ "__version", "dag" ]
+ }
+ ]
+}
diff --git a/ts-sdk/scripts/ci/prek/check_dag_schema.py
b/ts-sdk/scripts/ci/prek/check_dag_schema.py
new file mode 100755
index 00000000000..df214e841ac
--- /dev/null
+++ b/ts-sdk/scripts/ci/prek/check_dag_schema.py
@@ -0,0 +1,65 @@
+#!/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.
+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, console, 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"
+ )
+
+# Path of the generated file, relative to the repo root (for git diff /
messages).
+GENERATED = "ts-sdk/src/generated/dag-schema-fields.ts"
+
+if __name__ == "__main__":
+ directory = AIRFLOW_ROOT_PATH / "ts-sdk"
+ run_command(["pnpm", "config", "set", "store-dir", ".pnpm-store"],
cwd=directory)
+ run_command(["pnpm", "install", "--frozen-lockfile",
"--config.confirmModulesPurge=false"], cwd=directory)
+ # Regenerate, then format exactly as `pnpm run format` would, so the diff
+ # reflects a stale schema and not raw-vs-prettier formatting noise.
+ run_command(["pnpm", "run", "generate:dag-schema"], cwd=directory)
+ run_command(["pnpm", "exec", "prettier", "--write",
"src/generated/dag-schema-fields.ts"], cwd=directory)
+
+ diff = subprocess.run(
+ ["git", "diff", "--", GENERATED],
+ cwd=AIRFLOW_ROOT_PATH,
+ capture_output=True,
+ text=True,
+ check=False,
+ )
+ if diff.stdout.strip():
+ message = (
+ f"{GENERATED} is out of date with the vendored Dag serialization
schema "
+ "(ts-sdk/schema/dag-schema.json).\n"
+ "Regenerate it with `pnpm run generate:dag-schema` in ts-sdk/ and
commit the result."
+ )
+ if console:
+ console.print(f"[red]{message}[/]")
+ console.print(diff.stdout)
+ else:
+ print(message, file=sys.stderr)
+ print(diff.stdout, file=sys.stderr)
+ raise SystemExit(1)
diff --git a/ts-sdk/scripts/ci/prek/sync_dag_schema.py
b/ts-sdk/scripts/ci/prek/sync_dag_schema.py
new file mode 100755
index 00000000000..ab2f158ae38
--- /dev/null
+++ b/ts-sdk/scripts/ci/prek/sync_dag_schema.py
@@ -0,0 +1,65 @@
+#!/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.
+"""
+Refresh the TypeScript SDK's vendored copy of the Dag serialization schema.
+
+The SDK vendors the schema so a published npm package can be built without the
+monorepo, the same way the Java SDK vendors it into ``java-sdk/sdk/schema/``.
+This keeps the copy honest; ``check-ts-sdk-dag-schema`` then keeps the
+generated TypeScript honest with respect to the copy.
+"""
+
+from __future__ import annotations
+
+import shutil
+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, console
+
+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"
+ )
+
+SOURCE = Path("airflow-core/src/airflow/serialization/schema.json")
+VENDORED = Path("ts-sdk/schema/dag-schema.json")
+
+if __name__ == "__main__":
+ source_path = AIRFLOW_ROOT_PATH / SOURCE
+ vendored_path = AIRFLOW_ROOT_PATH / VENDORED
+ if not source_path.is_file():
+ raise SystemExit(f"{SOURCE} is missing; cannot refresh {VENDORED}")
+
+ if vendored_path.is_file() and vendored_path.read_bytes() ==
source_path.read_bytes():
+ sys.exit(0)
+
+ vendored_path.parent.mkdir(parents=True, exist_ok=True)
+ shutil.copyfile(source_path, vendored_path)
+ message = (
+ f"Refreshed {VENDORED} from {SOURCE}.\n"
+ "Review the diff, re-run `pnpm run generate:dag-schema` in ts-sdk/,
and commit both."
+ )
+ if console:
+ console.print(f"[yellow]{message}[/]")
+ else:
+ print(message, file=sys.stderr)
+ raise SystemExit(1)
diff --git a/ts-sdk/scripts/generate-dag-schema.mjs
b/ts-sdk/scripts/generate-dag-schema.mjs
new file mode 100644
index 00000000000..11b2daf33b2
--- /dev/null
+++ b/ts-sdk/scripts/generate-dag-schema.mjs
@@ -0,0 +1,512 @@
+#!/usr/bin/env node
+/*!
+ * 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.
+ */
+
+// Codegen for the Dag authoring fields of `DagSpec` and `TaskSpec`.
+//
+// Reads the vendored Dag serialization schema (`schema/dag-schema.json`, kept
+// in sync with `airflow-core/src/airflow/serialization/schema.json` by the
+// `sync-ts-sdk-dag-schema` prek hook) and emits
+// `src/generated/dag-schema-fields.ts`: the two field interfaces plus the
+// table of schema keys and defaults the serializer needs.
+//
+// Field selection mirrors the Go SDK's `TaskSpec` generator
+// (`go-sdk/bundle/bundlev1/gen`), as decision 12 of
+// `go-sdk/adr/0008-native-dag-interface.md` requires of every Lang SDK. Every
scalar
+// property of the "operator" definition becomes a task field, in schema
+// order, unless one of these rules skips it:
+//
+// - "_"-prefixed keys are serializer internals (_task_module, _is_mapped,
...)
+// - keys in the definition's "required" list are written by the serializer
+// itself (task_type, ui_color, ...)
+// - "has_on_*" keys are flags derived from Python callbacks, not settable
+// - non-scalar keys (arrays, objects, anyOf, refs other than timedelta and
+// datetime) cannot be expressed as a scalar authoring field
+// - EXCLUDED_TASK_FIELDS entries are Python-only concerns, documented per
key
+//
+// A new scalar key added to the schema therefore shows up in the regenerated
+// file — visible in its diff — or must gain an EXCLUDED_TASK_FIELDS entry,
+// instead of going silently missing. Dag-level fields come from the curated
+// DAG_FIELDS allowlist instead, matching Go's hand-written `DagSpec`, but
+// their types and defaults are still read from the schema so an allowlist
+// entry cannot outlive the key it names, and every other eligible "dag" key
+// has to be named in EXCLUDED_DAG_FIELDS or UNEXPRESSIBLE_DAG_FIELDS so an
+// allowlist cannot quietly skip a new one either.
+//
+// Three things differ from the Go generator, because TypeScript expresses
+// them directly and Go cannot:
+//
+// - Optionality is `?`, so a boolean whose schema default is true needs no
+// pointer to tell "unset" from an explicit false.
+// - `number` covers both integer and floating-point schema types, so
+// retry_exponential_backoff needs no per-key type override.
+// - Identity is positional — `new Dag(dagId)` and `dag.task(taskId, ...)` —
+// so dag_id and task_id stay serializer-owned here rather than being
+// re-exposed as spec fields the way Go's `TaskId` is.
+//
+// Names are the schema key in camelCase, matching the SDK's existing `dagId` /
+// `taskId` spelling. Unlike Go, no initialism table is needed: TypeScript
+// camelCase lowercases an acronym tail (`docMd`, `doXcomPush`).
+
+import console from "node:console";
+import { readFileSync, writeFileSync, mkdirSync } from "node:fs";
+import { dirname, join } from "node:path";
+import process from "node:process";
+import { fileURLToPath, pathToFileURL } from "node:url";
+
+const HERE = dirname(fileURLToPath(import.meta.url));
+const ROOT = join(HERE, "..");
+const SCHEMA_PATH = join(ROOT, "schema/dag-schema.json");
+const OUT_PATH = join(ROOT, "src/generated/dag-schema-fields.ts");
+
+/**
+ * Version the serializer stamps into a serialized Dag's `__version`.
+ *
+ * The schema cannot supply it: its root `allOf` only constrains `__version` to
+ * a positive integer. A core bump therefore has to be mirrored here by hand,
+ * and
`airflow-core/tests/unit/serialization/test_ts_sdk_serialization_version.py`
+ * fails until it is.
+ */
+export const SERIALIZATION_VERSION = 3;
+
+/**
+ * Scalar "operator" keys deliberately kept off `TaskSpec`, each with the
+ * reason. Every other eligible scalar key becomes a field, so an entry that no
+ * longer matches an eligible schema key fails generation and the list cannot
+ * go stale.
+ */
+export const EXCLUDED_TASK_FIELDS = {
+ doc: "legacy doc attribute; the UI renders Markdown, so only doc_md is
exposed",
+ doc_json: "legacy doc attribute; the UI renders Markdown, so only doc_md is
exposed",
+ doc_yaml: "legacy doc attribute; the UI renders Markdown, so only doc_md is
exposed",
+ doc_rst: "legacy doc attribute; the UI renders Markdown, so only doc_md is
exposed",
+ allow_nested_operators:
+ "Python runtime concern: warns when an operator executes inside another
operator",
+ multiple_outputs: "TaskFlow (@task) dict-unpacking concern, meaningless
outside Python",
+ start_from_trigger:
+ "deferrable-operator machinery; its start_trigger_args counterpart is an
object the spec cannot express",
+ is_setup: "setup/teardown flags carry trigger-rule invariants the SDK does
not model yet",
+ is_teardown: "setup/teardown flags carry trigger-rule invariants the SDK
does not model yet",
+ on_failure_fail_dagrun: "only valid on teardown tasks, which the SDK does
not model yet",
+};
+
+/**
+ * Eligible "operator" keys the spec cannot express, each with the shape that
+ * stops it. Checked the same way as {@link EXCLUDED_TASK_FIELDS}: an entry
that
+ * no longer matches fails generation, and so does an eligible key that is in
+ * neither list.
+ *
+ * Without this, a key changing shape — `{"type": "boolean"}` becoming
+ * `{"anyOf": [...]}` — would delete its `TaskSpec` field, and every
+ * `dag.task(..., { thatField })` already written would start failing as an
+ * unknown option.
+ */
+export const UNEXPRESSIBLE_TASK_FIELDS = {
+ params: "a params object of arbitrary shape",
+ executor_config: "a dict of arbitrary shape",
+ template_ext: "an array; an operator's templated file extensions are its
own",
+ template_fields_renderers: "a dict; an operator's field renderers are its
own",
+ downstream_task_ids: "an array holding the graph, which the Dag builds from
its wiring",
+ partial_kwargs: "an object; dynamic task mapping, which the SDK does not
model yet",
+ render_template_as_native_obj: "an anyOf the generator cannot narrow to one
type",
+ inlets: "an array of assets, which the SDK does not model yet",
+ outlets: "an array of assets, which the SDK does not model yet",
+ start_trigger_args: "an object; deferrable-operator machinery",
+};
+
+/**
+ * The Dag-level allowlist, in the order the fields are emitted.
+ *
+ * Hand-curated rather than derived, matching Go's `DagSpec`: the "dag"
+ * definition mixes authoring options with bundle bookkeeping (fileloc,
+ * bundle_name, ...) and structure the SDK builds itself (tasks, task_group,
+ * edge_info), and no mechanical rule separates the two. Types and defaults
+ * still come from the schema, so an entry naming a key the schema dropped
+ * fails generation.
+ *
+ * `tsType` overrides the type for a key the schema types too loosely to map.
+ * `mapsTo` marks a virtual field: one the authoring surface offers and the
+ * serializer lowers onto a schema key of a different shape.
+ */
+export const DAG_FIELDS = [
+ {
+ name: "schedule",
+ mapsTo: "timetable",
+ tsType: "string",
+ doc: 'When the Dag runs: `"@once"`, `"@continuous"`, a cron expression, or
unset for no schedule. The serializer lowers it onto the schema\'s `timetable`
object.',
+ },
+ { key: "description" },
+ { key: "start_date" },
+ { key: "end_date" },
+ // The schema declares a bare array with no item type, so it cannot say this.
+ { key: "tags", tsType: "readonly string[]", schemaType: "string[]" },
+ { key: "dag_display_name" },
+ { key: "doc_md" },
+ { key: "max_active_tasks" },
+ { key: "max_active_runs" },
+ { key: "max_consecutive_failed_dag_runs" },
+ { key: "dagrun_timeout" },
+ { key: "catchup" },
+ { key: "fail_fast" },
+ { key: "render_template_as_native_obj" },
+ { key: "disable_bundle_versioning" },
+ { key: "is_paused_upon_creation" },
+];
+
+/**
+ * Scalar "dag" keys deliberately kept off `DagSpec`, each with the reason.
+ * Checked like {@link EXCLUDED_TASK_FIELDS}.
+ */
+export const EXCLUDED_DAG_FIELDS = {
+ relative_fileloc: "bundle bookkeeping: where the Dag file sits in its
bundle, set when packing",
+};
+
+/**
+ * Eligible "dag" keys the spec cannot express, each with the shape that stops
+ * it. Checked like {@link UNEXPRESSIBLE_TASK_FIELDS}.
+ */
+export const UNEXPRESSIBLE_DAG_FIELDS = {
+ params: "a params object of arbitrary shape",
+ timezone: "a timezone object; the schedule and the dates carry their own",
+ owner_links: "a dict of owner names to links",
+ allowed_run_types: "an anyOf the generator cannot narrow to one type",
+ bundle_name: "an anyOf; the bundle a Dag ships in is chosen when packing",
+ deadline: "an anyOf; Dag-run deadlines, which the SDK does not model yet",
+ default_args: "a dict of arbitrary shape",
+ access_control: "a dict of arbitrary shape",
+ task_group: "an anyOf holding the group tree, which the Dag builds from its
wiring",
+ edge_info: "an object holding the edges, which the Dag builds from its
wiring",
+ dag_dependencies: "an object the serializer derives from the Dag's tasks",
+ rerun_with_latest_version: "a nullable boolean the generator cannot narrow
to one type",
+};
+
+/** Whether the serializer, not the author, owns `key` in `definition`. */
+function isSerializerOwned(key, required) {
+ return key.startsWith("_") || required.has(key) || key.startsWith("has_on_");
+}
+
+/**
+ * Map one schema property onto an authoring field, or return undefined when it
+ * is not a scalar the spec can express (array, object, anyOf, or a `$ref`
+ * other than timedelta/datetime).
+ */
+export function resolveField(key, prop) {
+ const ref =
+ typeof prop.$ref === "string" && prop.$ref.startsWith("#/definitions/")
+ ? prop.$ref.slice("#/definitions/".length)
+ : undefined;
+ // A `"default": null` says "absent by default", which is what leaving the
+ // field unset already means.
+ const hasDefault = Object.hasOwn(prop, "default") && prop.default !== null;
+ const base = { key, name: toCamelCase(key), ...(hasDefault && { default:
prop.default }) };
+
+ if (ref === "timedelta") return { ...base, tsType: "number", schemaType:
"timedelta" };
+ if (ref === "datetime") return { ...base, tsType: "Date", schemaType:
"datetime" };
+ if (ref !== undefined) return undefined;
+
+ switch (prop.type) {
+ case "string":
+ return { ...base, tsType: "string", schemaType: "string" };
+ case "integer":
+ case "number":
+ return { ...base, tsType: "number", schemaType: "number" };
+ case "boolean":
+ return { ...base, tsType: "boolean", schemaType: "boolean" };
+ default:
+ return undefined;
+ }
+}
+
+function getDefinition(schema, name) {
+ const definition = schema.definitions?.[name];
+ if (!definition?.properties) {
+ throw new Error(`schema has no "${name}" definition with properties`);
+ }
+ return definition;
+}
+
+/** Derive the `TaskSpec` fields from the "operator" definition, in schema
order. */
+export function selectTaskFields(schema) {
+ const operator = getDefinition(schema, "operator");
+ const required = new Set(operator.required ?? []);
+ const matchedExclusions = new Set();
+ const matchedUnexpressible = new Set();
+ const fields = [];
+
+ for (const [key, prop] of Object.entries(operator.properties)) {
+ if (isSerializerOwned(key, required)) continue;
+ if (Object.hasOwn(EXCLUDED_TASK_FIELDS, key)) {
+ matchedExclusions.add(key);
+ continue;
+ }
+ const field = resolveField(key, prop);
+ if (field) {
+ fields.push(field);
+ continue;
+ }
+ // A key the spec cannot express has to be named, so a key that changes
+ // shape breaks generation instead of quietly deleting its field.
+ if (Object.hasOwn(UNEXPRESSIBLE_TASK_FIELDS, key)) {
+ matchedUnexpressible.add(key);
+ continue;
+ }
+ throw new Error(
+ `schema property "${key}" is not a scalar TaskSpec can express; add it
to ` +
+ "UNEXPRESSIBLE_TASK_FIELDS with the shape that stops it, or to
EXCLUDED_TASK_FIELDS",
+ );
+ }
+
+ checkNoStaleEntries([
+ [EXCLUDED_TASK_FIELDS, matchedExclusions, "EXCLUDED_TASK_FIELDS"],
+ [UNEXPRESSIBLE_TASK_FIELDS, matchedUnexpressible,
"UNEXPRESSIBLE_TASK_FIELDS"],
+ ]);
+ return fields;
+}
+
+/** Fail on a list entry the schema no longer has, so no list can go stale. */
+function checkNoStaleEntries(lists) {
+ for (const [list, matched, name] of lists) {
+ for (const key of Object.keys(list)) {
+ if (!matched.has(key)) {
+ throw new Error(
+ `${name} entry "${key}" matches no eligible schema property; remove
or fix it`,
+ );
+ }
+ }
+ }
+}
+
+/** Resolve the `DAG_FIELDS` allowlist against the "dag" definition, in
allowlist order. */
+export function selectDagFields(schema) {
+ const dag = getDefinition(schema, "dag");
+ const required = new Set(dag.required ?? []);
+ const fields = DAG_FIELDS.map((entry) => {
+ if (entry.mapsTo !== undefined) {
+ if (!Object.hasOwn(dag.properties, entry.mapsTo)) {
+ throw new Error(
+ `virtual Dag field "${entry.name}" lowers onto schema key
"${entry.mapsTo}", which the "dag" definition no longer declares`,
+ );
+ }
+ return {
+ key: entry.mapsTo,
+ name: entry.name,
+ tsType: entry.tsType,
+ schemaType: entry.tsType,
+ doc: entry.doc,
+ virtual: true,
+ };
+ }
+
+ const prop = dag.properties[entry.key];
+ if (prop === undefined) {
+ throw new Error(
+ `DAG_FIELDS entry "${entry.key}" is not a property of the "dag"
definition; remove or fix it`,
+ );
+ }
+ if (isSerializerOwned(entry.key, required)) {
+ throw new Error(`DAG_FIELDS entry "${entry.key}" is serializer-owned and
cannot be authored`);
+ }
+
+ const field = resolveField(entry.key, prop);
+ if (entry.tsType) {
+ return {
+ ...(field ?? { key: entry.key, name: toCamelCase(entry.key) }),
+ tsType: entry.tsType,
+ schemaType: entry.schemaType ?? entry.tsType,
+ };
+ }
+ if (!field) {
+ throw new Error(
+ `DAG_FIELDS entry "${entry.key}" is not a scalar the spec can express;
give it a tsType override or remove it`,
+ );
+ }
+ return field;
+ });
+
+ // The allowlist says which keys are authored; this says every other eligible
+ // key was looked at, so a key the schema gains is a generation failure
rather
+ // than an option nobody notices is missing.
+ const authored = new Set(DAG_FIELDS.map((entry) => entry.mapsTo ??
entry.key));
+ const matchedExclusions = new Set();
+ const matchedUnexpressible = new Set();
+ for (const [key, prop] of Object.entries(dag.properties)) {
+ if (isSerializerOwned(key, required) || authored.has(key)) continue;
+ if (Object.hasOwn(EXCLUDED_DAG_FIELDS, key)) {
+ matchedExclusions.add(key);
+ continue;
+ }
+ if (!resolveField(key, prop) && Object.hasOwn(UNEXPRESSIBLE_DAG_FIELDS,
key)) {
+ matchedUnexpressible.add(key);
+ continue;
+ }
+ throw new Error(
+ `schema property "${key}" of the "dag" definition is in no list; add it
to DAG_FIELDS to ` +
+ "author it, to UNEXPRESSIBLE_DAG_FIELDS with the shape that stops it,
or to EXCLUDED_DAG_FIELDS",
+ );
+ }
+ checkNoStaleEntries([
+ [EXCLUDED_DAG_FIELDS, matchedExclusions, "EXCLUDED_DAG_FIELDS"],
+ [UNEXPRESSIBLE_DAG_FIELDS, matchedUnexpressible,
"UNEXPRESSIBLE_DAG_FIELDS"],
+ ]);
+ return fields;
+}
+
+/** `max_active_tis_per_dag` -> `maxActiveTisPerDag`. */
+export function toCamelCase(key) {
+ return key.replace(/_(.)/g, (_match, char) => char.toUpperCase());
+}
+
+const LICENSE_HEADER = `/*!
+ * 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.
+ */
+
+// AUTO-GENERATED by scripts/generate-dag-schema.mjs — do not edit by hand.
+// Source: schema/dag-schema.json, vendored from
+// airflow-core/src/airflow/serialization/schema.json.
+//
+// Re-run with: pnpm run generate:dag-schema
+`;
+
+/** Render one field's doc comment: what it maps to, and its schema default. */
+function renderFieldDoc(field) {
+ if (field.doc) return ` /** ${field.doc} */`;
+ const unit = field.schemaType === "timedelta" ? ", in seconds" : "";
+ const suffix =
+ field.default === undefined ? "" : ` (schema default
\`${JSON.stringify(field.default)}\`)`;
+ return ` /** Maps to the schema key \`${field.key}\`${unit}${suffix}. */`;
+}
+
+function renderInterface(name, doc, fields) {
+ const members = fields.flatMap((field) => [
+ renderFieldDoc(field),
+ ` readonly ${field.name}?: ${field.tsType};`,
+ ]);
+ return `${doc}\nexport interface ${name} {\n${members.join("\n")}\n}\n`;
+}
+
+function renderTable(name, doc, typeName, fields) {
+ const entries = fields.map((field) => {
+ const parts = [
+ `key: ${JSON.stringify(field.key)}`,
+ `type: ${JSON.stringify(field.schemaType)}`,
+ ];
+ if (field.default !== undefined) parts.push(`default:
${JSON.stringify(field.default)}`);
+ if (field.virtual) parts.push("virtual: true");
+ return ` ${field.name}: { ${parts.join(", ")} },`;
+ });
+ return `${doc}\nexport const ${name} = {\n${entries.join(
+ "\n",
+ )}\n} as const satisfies Readonly<Record<keyof ${typeName},
SchemaField>>;\n`;
+}
+
+export function renderModule(dagFields, taskFields) {
+ return [
+ LICENSE_HEADER,
+ renderInterface(
+ "GeneratedDagFields",
+ `/**
+ * Dag-level authoring fields, from the curated allowlist in
+ * scripts/generate-dag-schema.mjs resolved against the "dag" definition.
+ *
+ * Every field is optional: leaving one unset means the scheduler applies its
+ * own default, so \`{}\` stays a valid spec and a field added here later
cannot
+ * break an existing call site.
+ */`,
+ dagFields,
+ ),
+ renderInterface(
+ "GeneratedTaskFields",
+ `/**
+ * Task-level authoring fields, from the scalar properties of the "operator"
+ * definition that the author rather than the serializer owns.
+ *
+ * Optional on the same terms as {@link GeneratedDagFields}.
+ */`,
+ taskFields,
+ ),
+ `/** How a field is carried in the serialized Dag. \`datetime\` is
fractional
+ * seconds since the epoch and \`timedelta\` is seconds; both are written as
+ * numbers. */
+export type SchemaFieldType = "string" | "number" | "boolean" | "string[]" |
"datetime" | "timedelta";
+
+/** What the serializer needs to know about one authoring field. */
+export interface SchemaField {
+ /** The key this field takes in the serialized Dag. */
+ readonly key: string;
+ readonly type: SchemaFieldType;
+ /** The schema's default, absent when it declares none. A value equal to it
+ * can be omitted, as Python's BaseSerialization omits what the scheduler
+ * re-derives. */
+ readonly default?: string | number | boolean;
+ /** Set when the serializer derives \`key\` from this field rather than
+ * writing the value straight through. */
+ readonly virtual?: boolean;
+}
+`,
+ renderTable(
+ "DAG_SCHEMA_FIELDS",
+ "/** Schema key, wire type and default of every {@link
GeneratedDagFields} field. */",
+ "GeneratedDagFields",
+ dagFields,
+ ),
+ renderTable(
+ "TASK_SCHEMA_FIELDS",
+ "/** Schema key, wire type and default of every {@link
GeneratedTaskFields} field. */",
+ "GeneratedTaskFields",
+ taskFields,
+ ),
+ `/** The \`__version\` the serializer stamps into a serialized Dag. The
schema
+ * constrains it only to a positive integer, so it is pinned here and guarded
+ * by a test in airflow-core. */
+export const SERIALIZATION_VERSION = ${SERIALIZATION_VERSION};
+`,
+ ].join("\n");
+}
+
+function main() {
+ const schema = JSON.parse(readFileSync(SCHEMA_PATH, "utf8"));
+ const dagFields = selectDagFields(schema);
+ const taskFields = selectTaskFields(schema);
+
+ mkdirSync(dirname(OUT_PATH), { recursive: true });
+ writeFileSync(OUT_PATH, renderModule(dagFields, taskFields), "utf8");
+
+ console.log(`wrote ${OUT_PATH}`);
+ console.log(` dag fields=${dagFields.length}, task
fields=${taskFields.length}`);
+}
+
+// Importable for tests; runs only when invoked as a script.
+if (process.argv[1] && pathToFileURL(process.argv[1]).href ===
import.meta.url) {
+ main();
+}
diff --git a/ts-sdk/src/generated/dag-schema-fields.ts
b/ts-sdk/src/generated/dag-schema-fields.ts
new file mode 100644
index 00000000000..de155670fda
--- /dev/null
+++ b/ts-sdk/src/generated/dag-schema-fields.ts
@@ -0,0 +1,215 @@
+/*!
+ * 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.
+ */
+
+// AUTO-GENERATED by scripts/generate-dag-schema.mjs — do not edit by hand.
+// Source: schema/dag-schema.json, vendored from
+// airflow-core/src/airflow/serialization/schema.json.
+//
+// Re-run with: pnpm run generate:dag-schema
+
+/**
+ * Dag-level authoring fields, from the curated allowlist in
+ * scripts/generate-dag-schema.mjs resolved against the "dag" definition.
+ *
+ * Every field is optional: leaving one unset means the scheduler applies its
+ * own default, so `{}` stays a valid spec and a field added here later cannot
+ * break an existing call site.
+ */
+export interface GeneratedDagFields {
+ /** When the Dag runs: `"@once"`, `"@continuous"`, a cron expression, or
unset for no schedule. The serializer lowers it onto the schema's `timetable`
object. */
+ readonly schedule?: string;
+ /** Maps to the schema key `description`. */
+ readonly description?: string;
+ /** Maps to the schema key `start_date`. */
+ readonly startDate?: Date;
+ /** Maps to the schema key `end_date`. */
+ readonly endDate?: Date;
+ /** Maps to the schema key `tags`. */
+ readonly tags?: readonly string[];
+ /** Maps to the schema key `dag_display_name`. */
+ readonly dagDisplayName?: string;
+ /** Maps to the schema key `doc_md`. */
+ readonly docMd?: string;
+ /** Maps to the schema key `max_active_tasks`. */
+ readonly maxActiveTasks?: number;
+ /** Maps to the schema key `max_active_runs`. */
+ readonly maxActiveRuns?: number;
+ /** Maps to the schema key `max_consecutive_failed_dag_runs`. */
+ readonly maxConsecutiveFailedDagRuns?: number;
+ /** Maps to the schema key `dagrun_timeout`, in seconds. */
+ readonly dagrunTimeout?: number;
+ /** Maps to the schema key `catchup`. */
+ readonly catchup?: boolean;
+ /** Maps to the schema key `fail_fast` (schema default `false`). */
+ readonly failFast?: boolean;
+ /** Maps to the schema key `render_template_as_native_obj` (schema default
`false`). */
+ readonly renderTemplateAsNativeObj?: boolean;
+ /** Maps to the schema key `disable_bundle_versioning`. */
+ readonly disableBundleVersioning?: boolean;
+ /** Maps to the schema key `is_paused_upon_creation`. */
+ readonly isPausedUponCreation?: boolean;
+}
+
+/**
+ * Task-level authoring fields, from the scalar properties of the "operator"
+ * definition that the author rather than the serializer owns.
+ *
+ * Optional on the same terms as {@link GeneratedDagFields}.
+ */
+export interface GeneratedTaskFields {
+ /** Maps to the schema key `owner` (schema default `"airflow"`). */
+ readonly owner?: string;
+ /** Maps to the schema key `start_date`. */
+ readonly startDate?: Date;
+ /** Maps to the schema key `end_date`. */
+ readonly endDate?: Date;
+ /** Maps to the schema key `trigger_rule` (schema default `"all_success"`).
*/
+ readonly triggerRule?: string;
+ /** Maps to the schema key `depends_on_past` (schema default `false`). */
+ readonly dependsOnPast?: boolean;
+ /** Maps to the schema key `ignore_first_depends_on_past` (schema default
`false`). */
+ readonly ignoreFirstDependsOnPast?: boolean;
+ /** Maps to the schema key `wait_for_past_depends_before_skipping` (schema
default `false`). */
+ readonly waitForPastDependsBeforeSkipping?: boolean;
+ /** Maps to the schema key `wait_for_downstream` (schema default `false`). */
+ readonly waitForDownstream?: boolean;
+ /** Maps to the schema key `retries` (schema default `0`). */
+ readonly retries?: number;
+ /** Maps to the schema key `queue` (schema default `"default"`). */
+ readonly queue?: string;
+ /** Maps to the schema key `pool` (schema default `"default_pool"`). */
+ readonly pool?: string;
+ /** Maps to the schema key `pool_slots` (schema default `1`). */
+ readonly poolSlots?: number;
+ /** Maps to the schema key `execution_timeout`, in seconds. */
+ readonly executionTimeout?: number;
+ /** Maps to the schema key `retry_delay`, in seconds (schema default `300`).
*/
+ readonly retryDelay?: number;
+ /** Maps to the schema key `retry_exponential_backoff` (schema default `0`).
*/
+ readonly retryExponentialBackoff?: number;
+ /** Maps to the schema key `max_retry_delay`, in seconds. */
+ readonly maxRetryDelay?: number;
+ /** Maps to the schema key `priority_weight` (schema default `1`). */
+ readonly priorityWeight?: number;
+ /** Maps to the schema key `weight_rule` (schema default `"downstream"`). */
+ readonly weightRule?: string;
+ /** Maps to the schema key `executor`. */
+ readonly executor?: string;
+ /** Maps to the schema key `do_xcom_push` (schema default `true`). */
+ readonly doXcomPush?: boolean;
+ /** Maps to the schema key `email_on_failure` (schema default `true`). */
+ readonly emailOnFailure?: boolean;
+ /** Maps to the schema key `email_on_retry` (schema default `true`). */
+ readonly emailOnRetry?: boolean;
+ /** Maps to the schema key `doc_md`. */
+ readonly docMd?: string;
+ /** Maps to the schema key `map_index_template`. */
+ readonly mapIndexTemplate?: string;
+ /** Maps to the schema key `max_active_tis_per_dag`. */
+ readonly maxActiveTisPerDag?: number;
+ /** Maps to the schema key `max_active_tis_per_dagrun`. */
+ readonly maxActiveTisPerDagrun?: number;
+}
+
+/** How a field is carried in the serialized Dag. `datetime` is fractional
+ * seconds since the epoch and `timedelta` is seconds; both are written as
+ * numbers. */
+export type SchemaFieldType =
+ "string" | "number" | "boolean" | "string[]" | "datetime" | "timedelta";
+
+/** What the serializer needs to know about one authoring field. */
+export interface SchemaField {
+ /** The key this field takes in the serialized Dag. */
+ readonly key: string;
+ readonly type: SchemaFieldType;
+ /** The schema's default, absent when it declares none. A value equal to it
+ * can be omitted, as Python's BaseSerialization omits what the scheduler
+ * re-derives. */
+ readonly default?: string | number | boolean;
+ /** Set when the serializer derives `key` from this field rather than
+ * writing the value straight through. */
+ readonly virtual?: boolean;
+}
+
+/** Schema key, wire type and default of every {@link GeneratedDagFields}
field. */
+export const DAG_SCHEMA_FIELDS = {
+ schedule: { key: "timetable", type: "string", virtual: true },
+ description: { key: "description", type: "string" },
+ startDate: { key: "start_date", type: "datetime" },
+ endDate: { key: "end_date", type: "datetime" },
+ tags: { key: "tags", type: "string[]" },
+ dagDisplayName: { key: "dag_display_name", type: "string" },
+ docMd: { key: "doc_md", type: "string" },
+ maxActiveTasks: { key: "max_active_tasks", type: "number" },
+ maxActiveRuns: { key: "max_active_runs", type: "number" },
+ maxConsecutiveFailedDagRuns: { key: "max_consecutive_failed_dag_runs", type:
"number" },
+ dagrunTimeout: { key: "dagrun_timeout", type: "timedelta" },
+ catchup: { key: "catchup", type: "boolean" },
+ failFast: { key: "fail_fast", type: "boolean", default: false },
+ renderTemplateAsNativeObj: {
+ key: "render_template_as_native_obj",
+ type: "boolean",
+ default: false,
+ },
+ disableBundleVersioning: { key: "disable_bundle_versioning", type: "boolean"
},
+ isPausedUponCreation: { key: "is_paused_upon_creation", type: "boolean" },
+} as const satisfies Readonly<Record<keyof GeneratedDagFields, SchemaField>>;
+
+/** Schema key, wire type and default of every {@link GeneratedTaskFields}
field. */
+export const TASK_SCHEMA_FIELDS = {
+ owner: { key: "owner", type: "string", default: "airflow" },
+ startDate: { key: "start_date", type: "datetime" },
+ endDate: { key: "end_date", type: "datetime" },
+ triggerRule: { key: "trigger_rule", type: "string", default: "all_success" },
+ dependsOnPast: { key: "depends_on_past", type: "boolean", default: false },
+ ignoreFirstDependsOnPast: {
+ key: "ignore_first_depends_on_past",
+ type: "boolean",
+ default: false,
+ },
+ waitForPastDependsBeforeSkipping: {
+ key: "wait_for_past_depends_before_skipping",
+ type: "boolean",
+ default: false,
+ },
+ waitForDownstream: { key: "wait_for_downstream", type: "boolean", default:
false },
+ retries: { key: "retries", type: "number", default: 0 },
+ queue: { key: "queue", type: "string", default: "default" },
+ pool: { key: "pool", type: "string", default: "default_pool" },
+ poolSlots: { key: "pool_slots", type: "number", default: 1 },
+ executionTimeout: { key: "execution_timeout", type: "timedelta" },
+ retryDelay: { key: "retry_delay", type: "timedelta", default: 300 },
+ retryExponentialBackoff: { key: "retry_exponential_backoff", type: "number",
default: 0 },
+ maxRetryDelay: { key: "max_retry_delay", type: "timedelta" },
+ priorityWeight: { key: "priority_weight", type: "number", default: 1 },
+ weightRule: { key: "weight_rule", type: "string", default: "downstream" },
+ executor: { key: "executor", type: "string" },
+ doXcomPush: { key: "do_xcom_push", type: "boolean", default: true },
+ emailOnFailure: { key: "email_on_failure", type: "boolean", default: true },
+ emailOnRetry: { key: "email_on_retry", type: "boolean", default: true },
+ docMd: { key: "doc_md", type: "string" },
+ mapIndexTemplate: { key: "map_index_template", type: "string" },
+ maxActiveTisPerDag: { key: "max_active_tis_per_dag", type: "number" },
+ maxActiveTisPerDagrun: { key: "max_active_tis_per_dagrun", type: "number" },
+} as const satisfies Readonly<Record<keyof GeneratedTaskFields, SchemaField>>;
+
+/** The `__version` the serializer stamps into a serialized Dag. The schema
+ * constrains it only to a positive integer, so it is pinned here and guarded
+ * by a test in airflow-core. */
+export const SERIALIZATION_VERSION = 3;
diff --git a/ts-sdk/src/sdk/dag.ts b/ts-sdk/src/sdk/dag.ts
index 94e2a45ec63..a264d306104 100644
--- a/ts-sdk/src/sdk/dag.ts
+++ b/ts-sdk/src/sdk/dag.ts
@@ -22,6 +22,12 @@
// Dag and supplies its arguments, the way calling a TaskFlow function does in
// Python.
+import {
+ DAG_SCHEMA_FIELDS,
+ TASK_SCHEMA_FIELDS,
+ type GeneratedDagFields,
+ type GeneratedTaskFields,
+} from "../generated/dag-schema-fields.js";
import { brand, DUPLICATE_COPY_HINT, hasBrand } from "./brand.js";
import type { JsonValue } from "./client-types.js";
import type { TaskFunction } from "./task.js";
@@ -32,31 +38,74 @@ function isPlainRecord(value: unknown): value is
Record<string, unknown> {
return prototype === Object.prototype || prototype === null;
}
-function validateEmptySpec(name: string, value: unknown): void {
- if (!isPlainRecord(value) || Reflect.ownKeys(value).length > 0) {
- throw new Error(`${name} must be an empty object`);
+/**
+ * A copy of `spec` that nothing can change afterwards.
+ *
+ * Nothing reads a spec until the Dag is packed, long after the author's module
+ * has run, so an edit of the object they passed would silently change what
+ * ships. A shallow freeze is not enough: `tags` is an array the author still
+ * holds a reference to.
+ */
+function freezeSpec<TSpec extends object>(spec: TSpec, describe: () =>
string): TSpec {
+ return deepFreeze(spec, describe, new WeakSet()) as TSpec;
+}
+
+function deepFreeze(value: unknown, describe: () => string, seen:
WeakSet<object>): unknown {
+ if (typeof value !== "object" || value === null) return value;
+ // `startDate` and `endDate` are Dates, and a Date's setters change it in
+ // place, so the recorded spec takes a copy. Freezing would not help:
+ // Object.freeze does not stop setFullYear.
+ if (value instanceof Date) return new Date(value.getTime());
+ if (seen.has(value)) {
+ throw new Error(`${describe()} refers back to itself, so it cannot be
recorded`);
+ }
+ seen.add(value);
+ if (Array.isArray(value)) {
+ return Object.freeze(value.map((element) => deepFreeze(element, describe,
seen)));
+ }
+ if (!isPlainRecord(value)) {
+ // Passed through, this would be neither copied nor frozen, so the author
+ // could still change what ships.
+ throw new Error(
+ `${describe()} holds a ${kindOf(value)}, which cannot be recorded; a
spec field takes a ` +
+ "string, number, boolean, Date, or an array of them",
+ );
}
+ const copy: Record<string, unknown> = {};
+ for (const [key, nested] of Object.entries(value)) {
+ copy[key] = deepFreeze(nested, describe, seen);
+ }
+ return Object.freeze(copy);
+}
+
+function kindOf(value: object): string {
+ const prototype = Object.getPrototypeOf(value) as { constructor?: { name?:
string } } | null;
+ return prototype?.constructor?.name ?? "value";
}
+const DAG_SPEC_KEYS: ReadonlySet<string> = new
Set(Object.keys(DAG_SCHEMA_FIELDS));
+const TASK_SPEC_KEYS: ReadonlySet<string> = new
Set(Object.keys(TASK_SCHEMA_FIELDS));
+
/**
- * Dag-level options.
+ * Dag-level options: the schedule, the tags, how many runs may be active, and
+ * the rest of what `DAG(...)` takes in Python.
*
- * No fields yet, so only `{}` is accepted: a field that would be silently
- * dropped, such as `new Dag("d", { schedule: "@daily" })`, is a compile error.
+ * Every field is optional, so `{}` stays valid and a field added later cannot
+ * break a call site. An unknown key is rejected, so a misspelled field is an
+ * error rather than a Dag that quietly ignores it.
*
- * Native Dag declaration will add optional fields here, generated from the
- * serialized-Dag JSON schema as `src/generated/supervisor.ts` is.
+ * Setting a field records it. A Dag declared in TypeScript is not served to
+ * Airflow yet, so nothing reads it.
*/
-export type DagSpec = Record<string, never>;
+export type DagSpec = GeneratedDagFields;
/**
- * Task-level options.
- *
- * No fields yet, so only `{}` is accepted, as with {@link DagSpec}.
+ * Task-level options: the retries, the pool, the trigger rule, and the rest of
+ * what an operator takes in Python.
*
- * Future task fields (retries, ...) will land here.
+ * Optional and record-only on the same terms as {@link DagSpec}.
*/
-export type TaskSpec = Record<string, never>;
+export type TaskSpec = GeneratedTaskFields;
// Carries a reference's return type without carrying a value. Not exported, so
// the property cannot be read or written from outside; it exists only so
@@ -200,14 +249,13 @@ export type TaskFactory<TParams extends readonly
unknown[], TReturn = unknown> =
: (...inputs: PositionalInputs<TParams>) => TaskRef<TReturn>;
/**
- * Named options for `dag.task()`.
+ * The trailing argument of `dag.task()`: the task's own {@link TaskSpec}, plus
+ * the names of the handler's positional arguments.
*
- * Keyword-only so future fields can be added without a new parameter. Unknown
- * keys are rejected, so a typo fails at import time rather than being ignored.
+ * Every field is optional, and an unknown key is rejected, so a misspelled
+ * field is an error rather than a task that quietly ignores it.
*/
-export interface TaskOptions {
- /** Task-level options. Stored, but not used yet — see {@link TaskSpec}. */
- readonly spec?: TaskSpec;
+export type TaskOptions = TaskSpec & {
/**
* Names for the handler's positional arguments, in declaration order.
*
@@ -216,8 +264,8 @@ export interface TaskOptions {
* order, so a name left out only costs the label: `arg0`, `arg1` and so on
* stand in for it.
*/
- readonly argNames?: readonly string[];
-}
+ readonly argBindings?: readonly string[];
+};
/** Per-task record a Dag retains: the reference, the handler, and its spec. */
export interface TaskRecord {
@@ -288,14 +336,10 @@ export class Dag {
}
constructor(dagId: string, spec: DagSpec = {}) {
- validateEmptySpec(`spec for Dag "${dagId}"`, spec);
+ validateDagSpec(dagId, spec);
brand(this, "Dag");
this.dagId = dagId;
- // Copied and frozen, as task specs are: nothing reads a spec until the
- // bundle manifest is built, long after the user's module has run, so a
- // later mutation of their object would silently change what is packed.
- // Shallow, so a nested value in a future generated spec stays mutable.
- this.spec = Object.freeze({ ...spec });
+ this.spec = freezeSpec(spec, () => `The spec for Dag "${dagId}"`);
}
/** Task IDs attached to this Dag, in attachment order. */
@@ -308,7 +352,8 @@ export class Dag {
*
* Every argument the handler declares becomes an input of the returned
* {@link TaskFactory}; `getContext()` and `getClient()` reach the runtime
- * from inside the call, so neither is an argument.
+ * from inside the call, so neither is an argument. The trailing options
object
+ * carries this task's own {@link TaskSpec}.
*/
task<TParams extends readonly unknown[] = [], TReturn = unknown>(
taskId: string,
@@ -329,43 +374,49 @@ export class Dag {
"declare every task while the module is loading",
);
}
- this.#validateOptions(taskId, options);
- const { spec = {} } = options;
- validateEmptySpec(`spec for Dag "${this.dagId}" task "${taskId}"`, spec);
- const argNames = this.#validateArgNames(taskId, options.argNames);
+ const spec = this.#taskSpecOf(taskId, options);
+ const argBindings = this.#validateArgBindings(taskId, options.argBindings);
const task = createTaskRef(this.dagId, taskId);
this.#tasks.set(taskId, {
task,
// The runtime dispatches every handler through one instantiation, as it
// does a registered TaskHandler; a positional one is wrapped at wiring.
fn: handler as unknown as TaskFunction,
- spec: Object.freeze({ ...spec }),
+ spec: freezeSpec(spec, () => `The spec for Dag "${this.dagId}" task
"${taskId}"`),
});
return ((...inputs: unknown[]) => {
- this.#wire(taskId, inputs, argNames);
+ this.#wire(taskId, inputs, argBindings);
return task;
}) as TaskFactory<TParams, TReturn>;
}
- // TypeScript is bypassable — from plain JavaScript, or an `as TaskOptions`
- // cast — so an unknown key is rejected rather than silently ignored.
- #validateOptions(taskId: string, options: TaskOptions): void {
+ // TypeScript is bypassable — from plain JavaScript, or an `as TaskSpec` cast
+ // — so an unknown key is rejected rather than silently ignored.
`argBindings`
+ // names the handler's arguments rather than configuring the task, so it is
+ // taken out here instead of reaching the spec.
+ #taskSpecOf(taskId: string, options: TaskOptions): TaskSpec {
const value: unknown = options;
if (!isPlainRecord(value)) {
- throw new Error(`options for Dag "${this.dagId}" task "${taskId}" must
be an object`);
+ throw new Error(`spec for Dag "${this.dagId}" task "${taskId}" must be
an object`);
}
- for (const key of Object.keys(value)) {
- if (key !== "spec" && key !== "argNames") {
- throw new Error(`Unknown option "${key}" for Dag "${this.dagId}" task
"${taskId}"`);
+ const spec: Record<string, unknown> = {};
+ for (const key of Reflect.ownKeys(value)) {
+ if (key === "argBindings") continue;
+ if (typeof key !== "string" || !TASK_SPEC_KEYS.has(key)) {
+ throw new Error(
+ `Unknown option "${String(key)}" in the spec for Dag "${this.dagId}"
task "${taskId}"`,
+ );
}
+ spec[key] = value[key];
}
+ return spec as TaskSpec;
}
// TypeScript is bypassable, and these names become the keys the arguments
are
// recorded under, so an integer-like one would reorder what it labels.
- #validateArgNames(taskId: string, names: unknown): readonly string[] |
undefined {
+ #validateArgBindings(taskId: string, names: unknown): readonly string[] |
undefined {
if (names === undefined) return undefined;
- const describe = `argNames for Dag "${this.dagId}" task "${taskId}"`;
+ const describe = `argBindings for Dag "${this.dagId}" task "${taskId}"`;
if (!Array.isArray(names)) throw new Error(`${describe} must be an array
of names`);
const seen = new Set<string>();
for (const name of names as unknown[]) {
@@ -381,7 +432,11 @@ export class Dag {
return Object.freeze([...(names as string[])]);
}
- #wire(taskId: string, inputs: readonly unknown[], argNames: readonly
string[] | undefined): void {
+ #wire(
+ taskId: string,
+ inputs: readonly unknown[],
+ argBindings: readonly string[] | undefined,
+ ): void {
if (this.#finalized) {
throw new Error(
`Task "${taskId}" of Dag "${this.dagId}" was called after the Dag was
read; ` +
@@ -397,7 +452,7 @@ export class Dag {
const positional = !isNamedCall(inputs);
const recorded = this.#checkInputs(
taskId,
- positional ? positionalInputs(inputs, argNames) : (inputs[0] as
Record<string, unknown>),
+ positional ? positionalInputs(inputs, argBindings) : (inputs[0] as
Record<string, unknown>),
);
this.#inputs.set(taskId, recorded);
if (positional && inputs.length > 0) {
@@ -483,11 +538,11 @@ function isNamedCall(inputs: readonly unknown[]): boolean
{
function positionalInputs(
inputs: readonly unknown[],
- argNames: readonly string[] | undefined,
+ argBindings: readonly string[] | undefined,
): Record<string, unknown> {
const byName: Record<string, unknown> = {};
inputs.forEach((value, index) => {
- byName[argNames?.[index] ?? `arg${index}`] = value;
+ byName[argBindings?.[index] ?? `arg${index}`] = value;
});
return byName;
}
@@ -511,6 +566,18 @@ function createTaskRef(dagId: string, taskId: string):
TaskRef {
return Object.freeze(task);
}
+function validateDagSpec(dagId: string, spec: DagSpec): void {
+ const value: unknown = spec;
+ if (!isPlainRecord(value)) {
+ throw new Error(`spec for Dag "${dagId}" must be an object`);
+ }
+ for (const key of Reflect.ownKeys(value)) {
+ if (typeof key !== "string" || !DAG_SPEC_KEYS.has(key)) {
+ throw new Error(`Unknown option "${String(key)}" in the spec for Dag
"${dagId}"`);
+ }
+ }
+}
+
/**
* Internal: the task records of a Dag, for bundle lookups.
*
diff --git a/ts-sdk/tests/public-api.test.ts b/ts-sdk/tests/public-api.test.ts
index b62075579f1..b3a55e54f58 100644
--- a/ts-sdk/tests/public-api.test.ts
+++ b/ts-sdk/tests/public-api.test.ts
@@ -291,12 +291,11 @@ describe("public API", () => {
// a wider one is, and not the other way round.
expectTypeOf<TaskRef<boolean>>().toMatchTypeOf<TaskRef>();
expectTypeOf<TaskRef>().not.toMatchTypeOf<TaskRef<boolean>>();
- // Wiring moved to the factory call, so `inputs` is no longer an option and
- // the only remaining one is the spec.
- expectTypeOf<TaskOptions>().toEqualTypeOf<{
- readonly spec?: TaskSpec;
- readonly argNames?: readonly string[];
- }>();
+ // Wiring moved to the factory call and the spec is the trailing argument
+ // itself, alongside the argument names the packer fills in.
+ expectTypeOf<TaskOptions>().toEqualTypeOf<
+ TaskSpec & { readonly argBindings?: readonly string[] }
+ >();
// Each named argument takes any upstream reference or a literal of its
own type.
expectTypeOf<TaskInputs<{ rows: number }>>().toEqualTypeOf<{ rows: TaskRef
| number }>();
// A positional argument takes a literal or a reference of the argument's
own
@@ -321,11 +320,22 @@ describe("public API", () => {
) => TaskFactory<TParams, TReturn>
>();
expectTypeOf<Dag["taskIds"]>().toEqualTypeOf<readonly string[]>();
- // Reserved with no fields yet, so only `{}` is expressible. Generated
specs
- // will be all-optional (weak) types, and `{}` stays assignable to those,
so
- // filling these in later cannot break a call site.
- expectTypeOf<DagSpec>().toEqualTypeOf<Record<string, never>>();
- expectTypeOf<TaskSpec>().toEqualTypeOf<Record<string, never>>();
+ // Both specs are all-optional, so `{}` stays assignable and a field the
+ // schema gains later cannot break a call site.
+ const emptyDagSpec: DagSpec = {};
+ const emptyTaskSpec: TaskSpec = {};
+ expect([emptyDagSpec, emptyTaskSpec]).toEqual([{}, {}]);
+ expectTypeOf<DagSpec["schedule"]>().toEqualTypeOf<string | undefined>();
+ expectTypeOf<DagSpec["tags"]>().toEqualTypeOf<readonly string[] |
undefined>();
+ expectTypeOf<DagSpec["startDate"]>().toEqualTypeOf<Date | undefined>();
+ expectTypeOf<TaskSpec["retries"]>().toEqualTypeOf<number | undefined>();
+ // A timedelta is seconds, not a Date, and a bool defaulting to true is
+ // still a plain optional: `undefined` already means "unset".
+ expectTypeOf<TaskSpec["retryDelay"]>().toEqualTypeOf<number | undefined>();
+ expectTypeOf<TaskSpec["doXcomPush"]>().toEqualTypeOf<boolean |
undefined>();
+ // Identity is positional, so it is not restated in either spec.
+ expectTypeOf<DagSpec>().not.toHaveProperty("dagId");
+ expectTypeOf<TaskSpec>().not.toHaveProperty("taskId");
});
it("uses idiomatic TypeScript names for public client types", () => {
@@ -417,16 +427,20 @@ describe("public API", () => {
new Dag("example").task("extract");
const dag = new Dag("example");
const extract = dag.task("extract", async () => undefined);
- // @ts-expect-error wiring belongs to the factory call, not the options.
+ // @ts-expect-error wiring belongs to the factory call, not the spec.
dag.task("transform", async () => undefined, { inputs: { count: 1 } });
- // @ts-expect-error the spec is keyword-only, not positional.
+ // @ts-expect-error the spec holds task options, not task references.
dag.task("transform2", async () => undefined, { extract });
// @ts-expect-error a Dag spec is an options object, not a primitive.
new Dag("spec_dag", 42);
- // @ts-expect-error DagSpec has no fields yet, so a schedule cannot be
declared here.
- new Dag("spec_dag", { schedule: "@daily" });
- // @ts-expect-error TaskSpec has no fields yet, so retries cannot be
declared here.
- dag.task("transform3", async () => undefined, { spec: { retries: 2 } });
+ new Dag("scheduled_dag", { schedule: "@daily", tags: ["team-a"] });
+ dag.task("transform3", async () => undefined, { retries: 2 });
+ // @ts-expect-error specs use the camelCased field name, not the schema
key.
+ new Dag("snake_case_dag", { dag_display_name: "Example" });
+ // @ts-expect-error a field the schema does not define is a typo.
+ new Dag("typo_dag", { scheduled: "@daily" });
+ // @ts-expect-error a timedelta field is seconds, not a Date.
+ dag.task("retry_dag", async () => undefined, { retryDelay: new Date() });
// @ts-expect-error a handler with no arguments is called with none.
extract({ rows: 1 });
const transform = dag.task("transform4", async (_: { rows: number }) =>
undefined);
diff --git a/ts-sdk/tests/scripts/generate-dag-schema.test.ts
b/ts-sdk/tests/scripts/generate-dag-schema.test.ts
new file mode 100644
index 00000000000..035a0e91707
--- /dev/null
+++ b/ts-sdk/tests/scripts/generate-dag-schema.test.ts
@@ -0,0 +1,419 @@
+/*!
+ * 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 { readFileSync } from "node:fs";
+import { dirname, join } from "node:path";
+import { fileURLToPath } from "node:url";
+import { describe, expect, it } from "vitest";
+import {
+ DAG_SCHEMA_FIELDS,
+ SERIALIZATION_VERSION,
+ TASK_SCHEMA_FIELDS,
+} from "../../src/generated/dag-schema-fields.js";
+import {
+ DAG_FIELDS,
+ EXCLUDED_DAG_FIELDS,
+ EXCLUDED_TASK_FIELDS,
+ UNEXPRESSIBLE_DAG_FIELDS,
+ UNEXPRESSIBLE_TASK_FIELDS,
+ renderModule,
+ resolveField,
+ selectDagFields,
+ selectTaskFields,
+ toCamelCase,
+} from "../../scripts/generate-dag-schema.mjs";
+
+const SCHEMA_PATH = join(dirname(fileURLToPath(import.meta.url)),
"../../schema/dag-schema.json");
+
+interface SchemaDefinition {
+ required?: string[];
+ properties: Record<string, Record<string, unknown>>;
+}
+
+/** The vendored schema the checked-in `src/generated/dag-schema-fields.ts`
was built from. */
+function readSchema(): { definitions: { dag: SchemaDefinition; operator:
SchemaDefinition } } {
+ return JSON.parse(readFileSync(SCHEMA_PATH, "utf8"));
+}
+
+/**
+ * A dag definition holding every allowlisted key with its real shape, so a
+ * drift case only has to add the one key it is about.
+ */
+function dagSchema(properties: Record<string, Record<string, unknown>>) {
+ const declared = readSchema().definitions.dag.properties;
+ return {
+ definitions: {
+ operator: { required: ["task_id"], properties: { task_id: { type:
"string" } } },
+ dag: {
+ required: [],
+ properties: {
+ ...Object.fromEntries(
+ DAG_FIELDS.map((entry) => {
+ const key = entry.mapsTo ?? entry.key;
+ return [key, declared[key]];
+ }),
+ ),
+ ...properties,
+ },
+ },
+ },
+ };
+}
+
+/** A minimal schema shaped like the real one, for the drift cases. */
+function buildSchema(overrides: {
+ operator?: Record<string, unknown>;
+ dag?: Record<string, unknown>;
+}) {
+ return {
+ definitions: {
+ operator: { required: ["task_id"], properties: { task_id: { type:
"string" } } },
+ dag: { required: ["dag_id"], properties: { timetable: { type: "object" }
} },
+ ...overrides,
+ },
+ };
+}
+
+describe("task field selection", () => {
+ it("accounts for every property of the operator definition", () => {
+ // Nothing here names the fields: an eligible key in neither exclusion list
+ // fails selection, so a key the schema gains has to be triaged.
+ const selected = selectTaskFields(readSchema());
+ const operator = readSchema().definitions.operator;
+ const accounted = new Set([
+ ...selected.map((field) => field.key),
+ ...Object.keys(EXCLUDED_TASK_FIELDS),
+ ...Object.keys(UNEXPRESSIBLE_TASK_FIELDS),
+ ...(operator.required ?? []),
+ ]);
+ const unaccounted = Object.keys(operator.properties).filter(
+ (key) => !accounted.has(key) && !key.startsWith("_") &&
!key.startsWith("has_on_"),
+ );
+ expect(unaccounted).toEqual([]);
+ expect(selected.length).toBeGreaterThan(0);
+ });
+
+ it.each([
+ ["a schema-required key the serializer writes", "task_type"],
+ ["the positionally supplied task id", "task_id"],
+ ["a serializer internal", "_task_module"],
+ ["a callback-derived flag", "has_on_failure_callback"],
+ ["an array", "inlets"],
+ ["an object ref", "executor_config"],
+ ["an anyOf", "render_template_as_native_obj"],
+ ["an excluded Python-only concern", "is_setup"],
+ ])("skips %s", (_label, schemaKey) => {
+ const keys = selectTaskFields(readSchema()).map((field) => field.key);
+ expect(keys).not.toContain(schemaKey);
+ });
+
+ it("fails when an exclusion no longer matches an eligible schema key", () =>
{
+ const [excluded] = Object.keys(EXCLUDED_TASK_FIELDS);
+ const schema = buildSchema({
+ operator: {
+ required: [],
+ properties: Object.fromEntries(
+ Object.keys(EXCLUDED_TASK_FIELDS)
+ .filter((key) => key !== excluded)
+ .map((key) => [key, { type: "boolean" }]),
+ ),
+ },
+ });
+ expect(() => selectTaskFields(schema)).toThrowError(
+ new RegExp(`EXCLUDED_TASK_FIELDS entry "${excluded}" matches no eligible
schema property`),
+ );
+ });
+
+ it("fails when an unexpressible entry no longer matches an eligible schema
key", () => {
+ const [unexpressible] = Object.keys(UNEXPRESSIBLE_TASK_FIELDS);
+ const schema = buildSchema({
+ operator: {
+ required: [],
+ properties: {
+ ...Object.fromEntries(
+ Object.keys(EXCLUDED_TASK_FIELDS).map((key) => [key, { type:
"boolean" }]),
+ ),
+ ...Object.fromEntries(
+ Object.keys(UNEXPRESSIBLE_TASK_FIELDS)
+ .filter((key) => key !== unexpressible)
+ .map((key) => [key, { type: "object" }]),
+ ),
+ },
+ },
+ });
+ expect(() => selectTaskFields(schema)).toThrowError(
+ new RegExp(`UNEXPRESSIBLE_TASK_FIELDS entry "${unexpressible}" matches
no eligible schema`),
+ );
+ });
+
+ it("fails when an eligible key is in neither list", () => {
+ // A key that changes shape would otherwise delete its TaskSpec field, and
+ // every spec already setting it would start failing as an unknown option.
+ const schema = buildSchema({
+ operator: {
+ required: [],
+ properties: {
+ ...Object.fromEntries(
+ Object.keys(EXCLUDED_TASK_FIELDS).map((key) => [key, { type:
"boolean" }]),
+ ),
+ ...Object.fromEntries(
+ Object.keys(UNEXPRESSIBLE_TASK_FIELDS).map((key) => [key, { type:
"object" }]),
+ ),
+ retries: { anyOf: [{ type: "integer" }, { type: "null" }] },
+ },
+ },
+ });
+ expect(() => selectTaskFields(schema)).toThrowError(
+ /schema property "retries" is not a scalar TaskSpec can express/,
+ );
+ });
+
+ it("fails when the operator definition is gone", () => {
+ expect(() => selectTaskFields({ definitions: {} })).toThrowError(
+ /schema has no "operator" definition/,
+ );
+ });
+});
+
+describe("Dag field selection", () => {
+ it("resolves the allowlist in its own order", () => {
+ const names = selectDagFields(readSchema()).map((field) => field.name);
+
+ expect(names).toEqual(DAG_FIELDS.map((entry) => toCamelCase(entry.name ??
entry.key)));
+ });
+
+ it("accounts for every property of the dag definition", () => {
+ const dag = readSchema().definitions.dag;
+ const accounted = new Set([
+ ...DAG_FIELDS.map((entry) => entry.mapsTo ?? entry.key),
+ ...Object.keys(EXCLUDED_DAG_FIELDS),
+ ...Object.keys(UNEXPRESSIBLE_DAG_FIELDS),
+ ...(dag.required ?? []),
+ ]);
+ const unaccounted = Object.keys(dag.properties).filter(
+ (key) => !accounted.has(key) && !key.startsWith("_") &&
!key.startsWith("has_on_"),
+ );
+
+ expect(unaccounted).toEqual([]);
+ });
+
+ it("fails when an eligible dag key is in neither list", () => {
+ // Without this an allowlist would quietly skip a new authoring option; the
+ // task side is checked the same way.
+ const schema = dagSchema({ newly_added: { type: "string" } });
+
+ expect(() => selectDagFields(schema)).toThrowError(
+ /schema property "newly_added" of the "dag" definition is in no list/,
+ );
+ });
+
+ it.each([
+ ["exclusion", "EXCLUDED_DAG_FIELDS"],
+ ["unexpressible entry", "UNEXPRESSIBLE_DAG_FIELDS"],
+ ])("fails when a dag %s no longer matches a schema key", (_label, name) => {
+ // A scalar key belongs to the exclusions and a non-scalar one to the
+ // unexpressible list, so each entry is declared with the shape its own
list
+ // expects and the dropped one is simply left out.
+ const shapes: Record<string, Record<string, unknown>> = {
+ ...Object.fromEntries(
+ Object.keys(EXCLUDED_DAG_FIELDS).map((key) => [key, { type: "string"
}]),
+ ),
+ ...Object.fromEntries(
+ Object.keys(UNEXPRESSIBLE_DAG_FIELDS).map((key) => [key, { type:
"array" }]),
+ ),
+ };
+ const list = name === "EXCLUDED_DAG_FIELDS" ? EXCLUDED_DAG_FIELDS :
UNEXPRESSIBLE_DAG_FIELDS;
+ const dropped = Object.keys(list)[0]!;
+ const schema = dagSchema(
+ Object.fromEntries(Object.entries(shapes).filter(([key]) => key !==
dropped)),
+ );
+
+ expect(() => selectDagFields(schema)).toThrowError(
+ new RegExp(`${name} entry "${dropped}" matches no eligible schema
property`),
+ );
+ });
+
+ it("keeps the allowlist within what the dag definition actually declares",
() => {
+ const declared = Object.keys(readSchema().definitions.dag.properties);
+ for (const entry of DAG_FIELDS) {
+ expect(declared).toContain(entry.mapsTo ?? entry.key);
+ }
+ });
+
+ it("lowers the virtual schedule field onto the timetable key", () => {
+ const schedule = selectDagFields(readSchema()).find((field) => field.name
=== "schedule");
+ expect(schedule).toMatchObject({ key: "timetable", tsType: "string",
virtual: true });
+ });
+
+ it("types tags from the override, since the schema declares a bare array",
() => {
+ const tags = selectDagFields(readSchema()).find((field) => field.name ===
"tags");
+ expect(tags).toMatchObject({ tsType: "readonly string[]", schemaType:
"string[]" });
+ });
+
+ it.each([
+ [
+ "the virtual field's target is gone",
+ { dag: { properties: { catchup: { type: "boolean" } } } },
+ /virtual Dag field "schedule" lowers onto schema key "timetable"/,
+ ],
+ [
+ "an allowlisted key is gone",
+ { dag: { properties: { timetable: { type: "object" } } } },
+ /DAG_FIELDS entry "description" is not a property of the "dag"
definition/,
+ ],
+ [
+ "an allowlisted key became serializer-owned",
+ {
+ dag: {
+ required: ["description"],
+ properties: { timetable: { type: "object" }, description: { type:
"string" } },
+ },
+ },
+ /DAG_FIELDS entry "description" is serializer-owned/,
+ ],
+ [
+ "an allowlisted key stopped being expressible and has no override",
+ {
+ dag: {
+ properties: { timetable: { type: "object" }, description: { type:
"array" } },
+ },
+ },
+ /DAG_FIELDS entry "description" is not a scalar the spec can express/,
+ ],
+ ])("fails when %s", (_label, overrides, expected) => {
+ expect(() =>
selectDagFields(buildSchema(overrides))).toThrowError(expected);
+ });
+
+ it("fails when the dag definition is gone", () => {
+ expect(() => selectDagFields({ definitions: {} })).toThrowError(
+ /schema has no "dag" definition/,
+ );
+ });
+});
+
+describe("property resolution", () => {
+ it.each([
+ ["a string", { type: "string" }, "string", "string"],
+ ["an integer", { type: "integer" }, "number", "number"],
+ ["a number", { type: "number" }, "number", "number"],
+ ["a boolean", { type: "boolean" }, "boolean", "boolean"],
+ ["a timedelta ref", { $ref: "#/definitions/timedelta" }, "number",
"timedelta"],
+ ["a datetime ref", { $ref: "#/definitions/datetime" }, "Date", "datetime"],
+ ])("maps %s", (_label, prop, tsType, schemaType) => {
+ expect(resolveField("some_key", prop)).toMatchObject({ tsType, schemaType
});
+ });
+
+ it.each([
+ ["an array", { type: "array" }],
+ ["an object", { type: "object" }],
+ ["an anyOf", { anyOf: [{ type: "string" }] }],
+ ["a ref the spec cannot express", { $ref: "#/definitions/dict" }],
+ ])("cannot express %s", (_label, prop) => {
+ expect(resolveField("some_key", prop)).toBeUndefined();
+ });
+
+ it("carries a declared default, but reads an explicit null as no default",
() => {
+ expect(resolveField("k", { type: "string", default: "x"
})).toHaveProperty("default", "x");
+ expect(resolveField("k", { type: "string", default: null
})).not.toHaveProperty("default");
+ expect(resolveField("k", { type: "string"
})).not.toHaveProperty("default");
+ });
+
+ it("camel-cases schema keys without special-casing acronyms", () => {
+ expect(toCamelCase("max_active_tis_per_dag")).toBe("maxActiveTisPerDag");
+ expect(toCamelCase("do_xcom_push")).toBe("doXcomPush");
+ expect(toCamelCase("owner")).toBe("owner");
+ });
+});
+
+describe("the emitted module", () => {
+ it("documents the schema key, the seconds unit and the default of each
field", () => {
+ const rendered = renderModule(
+ [{ key: "d", name: "d", tsType: "string", schemaType: "string", doc:
"Custom." }],
+ [
+ {
+ key: "retry_delay",
+ name: "retryDelay",
+ tsType: "number",
+ schemaType: "timedelta",
+ default: 300,
+ },
+ { key: "executor", name: "executor", tsType: "string", schemaType:
"string" },
+ ],
+ );
+ expect(rendered).toContain("/** Custom. */");
+ expect(rendered).toContain(
+ "/** Maps to the schema key `retry_delay`, in seconds (schema default
`300`). */",
+ );
+ expect(rendered).toContain("/** Maps to the schema key `executor`. */");
+ expect(rendered).toContain("readonly executor?: string;");
+ });
+
+ it("pins the serialization version the schema cannot express", () => {
+ expect(SERIALIZATION_VERSION).toBe(3);
+ });
+});
+
+describe("the field table", () => {
+ it("gives every authoring field its schema key and wire type", () => {
+ for (const table of [DAG_SCHEMA_FIELDS, TASK_SCHEMA_FIELDS]) {
+ for (const [name, field] of Object.entries(table)) {
+ expect(field.key, name).toBeTruthy();
+ expect(
+ ["string", "number", "boolean", "string[]", "datetime", "timedelta"],
+ name,
+ ).toContain(field.type);
+ }
+ }
+ });
+
+ it("stays in step with the generator's own selection", () => {
+ const schema = readSchema();
+
expect(Object.keys(DAG_SCHEMA_FIELDS)).toEqual(selectDagFields(schema).map((f)
=> f.name));
+
expect(Object.keys(TASK_SCHEMA_FIELDS)).toEqual(selectTaskFields(schema).map((f)
=> f.name));
+ });
+
+ it("carries the schema defaults the serializer omits against", () => {
+ expect(TASK_SCHEMA_FIELDS.owner).toEqual({ key: "owner", type: "string",
default: "airflow" });
+ expect(TASK_SCHEMA_FIELDS.retryDelay).toEqual({
+ key: "retry_delay",
+ type: "timedelta",
+ default: 300,
+ });
+ expect(TASK_SCHEMA_FIELDS.doXcomPush).toEqual({
+ key: "do_xcom_push",
+ type: "boolean",
+ default: true,
+ });
+ expect(DAG_SCHEMA_FIELDS.failFast).toEqual({
+ key: "fail_fast",
+ type: "boolean",
+ default: false,
+ });
+ // No schema default, so the scheduler's own default stands.
+ expect(DAG_SCHEMA_FIELDS.catchup).toEqual({ key: "catchup", type:
"boolean" });
+ });
+
+ it("marks schedule as derived rather than written straight through", () => {
+ expect(DAG_SCHEMA_FIELDS.schedule).toEqual({
+ key: "timetable",
+ type: "string",
+ virtual: true,
+ });
+ });
+});
diff --git a/ts-sdk/tests/sdk/dag.test.ts b/ts-sdk/tests/sdk/dag.test.ts
index 5c772eb57c8..098a2064b21 100644
--- a/ts-sdk/tests/sdk/dag.test.ts
+++ b/ts-sdk/tests/sdk/dag.test.ts
@@ -23,7 +23,9 @@ import {
finalizeDag,
getDagTaskInputs,
getDagTaskRecords,
+ type DagSpec,
type TaskRef,
+ type TaskSpec,
} from "../../src/sdk/dag.js";
import { Bundle, finalizeBundleDags } from "../../src/sdk/bundle.js";
@@ -45,7 +47,7 @@ describe("Dag", () => {
"transform",
async (_: { extracted: { rows: number } }) => undefined,
);
- const load = dag.task("load", async (_: { transformed: undefined }) =>
undefined, { spec: {} });
+ const load = dag.task("load", async (_: { transformed: undefined }) =>
undefined, {});
const extracted = extract();
const transformed = transform({ extracted });
@@ -203,11 +205,11 @@ describe("Dag", () => {
expect(Object.keys(getDagTaskInputs(dag).get("transform")!)).toEqual(["arg0",
"arg1"]);
});
- it("names positional arguments from argNames", () => {
+ it("names positional arguments from argBindings", () => {
const dag = new Dag("example_dag");
const extract = dag.task("extract", async (): Promise<number> => 1);
const transform = dag.task("transform", async (rows: number, region:
string) => `${region}`, {
- argNames: ["rows", "region"],
+ argBindings: ["rows", "region"],
});
const extracted = extract();
@@ -216,10 +218,10 @@ describe("Dag", () => {
expect(getDagTaskInputs(dag).get("transform")).toEqual({ rows: extracted,
region: "us" });
});
- it("labels the arguments argNames does not reach", () => {
+ it("labels the arguments argBindings does not reach", () => {
const dag = new Dag("example_dag");
const transform = dag.task("transform", async (rows: number, region:
string) => `${region}`, {
- argNames: ["rows"],
+ argBindings: ["rows"],
});
transform(1, "us");
@@ -228,12 +230,20 @@ describe("Dag", () => {
});
it.each([
- ["not an array", { argNames: 1 }, /argNames for Dag "d" task "t" must be
an array of names/],
- ["not a string", { argNames: [1] }, /holds 1; each name must be a
non-empty string/],
- ["empty", { argNames: [""] }, /holds ""; each name must be a non-empty
string/],
- ["a number", { argNames: ["0"] }, /holds "0"; each name must be a
non-empty string/],
- ["a duplicate", { argNames: ["a", "a"] }, /argNames for Dag "d" task "t"
names "a" twice/],
- ])("rejects argNames that are %s", (_label, options, expected) => {
+ [
+ "not an array",
+ { argBindings: 1 },
+ /argBindings for Dag "d" task "t" must be an array of names/,
+ ],
+ ["not a string", { argBindings: [1] }, /holds 1; each name must be a
non-empty string/],
+ ["empty", { argBindings: [""] }, /holds ""; each name must be a non-empty
string/],
+ ["a number", { argBindings: ["0"] }, /holds "0"; each name must be a
non-empty string/],
+ [
+ "a duplicate",
+ { argBindings: ["a", "a"] },
+ /argBindings for Dag "d" task "t" names "a" twice/,
+ ],
+ ])("rejects argBindings that are %s", (_label, options, expected) => {
const dag = new Dag("d");
expect(() => dag.task("t", async (a: number) => a, options as
never)).toThrowError(expected);
@@ -268,7 +278,7 @@ describe("Dag", () => {
async (rows: number, region: string) => {
seen.push(rows, region);
},
- { argNames: ["rows", "region"] },
+ { argBindings: ["rows", "region"] },
);
transform(1, "us");
@@ -361,7 +371,7 @@ describe("Dag", () => {
const taskSpec = {};
const handler = async () => "hello";
const dag = new Dag("example_dag", dagSpec);
- dag.task("my_task", handler, { spec: taskSpec })();
+ dag.task("my_task", handler, taskSpec)();
expect(dag.dagId).toBe("example_dag");
expect(dag.spec).toEqual(dagSpec);
@@ -372,29 +382,101 @@ describe("Dag", () => {
expect(Object.isFrozen(record!.spec)).toBe(true);
});
+ it("copies a spec deeply, so editing the array afterwards cannot change what
ships", () => {
+ // Nothing reads a spec until the Dag is packed, long after the author's
+ // module has run, and `tags` is an array they still hold.
+ const tags = ["etl"];
+ const dag = new Dag("deep_spec_dag", { tags });
+
+ tags.push("injected");
+
+ expect(dag.spec.tags).toEqual(["etl"]);
+ expect(Object.isFrozen(dag.spec.tags)).toBe(true);
+ });
+
+ it("copies a Date in a spec, so a setter afterwards cannot change what
ships", () => {
+ const startDate = new Date("2026-01-01T00:00:00Z");
+ const dag = new Dag("dated_dag", { startDate });
+ const task = dag.task("extract", async () => undefined, { startDate });
+
+ startDate.setFullYear(2030);
+ task();
+
+ expect(dag.spec.startDate?.toISOString()).toBe("2026-01-01T00:00:00.000Z");
+
expect(getDagTaskRecords(dag).get("extract")?.spec.startDate?.toISOString()).toBe(
+ "2026-01-01T00:00:00.000Z",
+ );
+ });
+
+ it("rejects a spec value that is neither JSON nor a Date", () => {
+ expect(() => new Dag("map_dag", { tags: [new Map()] as unknown as string[]
})).toThrowError(
+ /holds a Map, which cannot be recorded/,
+ );
+ });
+
+ it("rejects a spec that refers back to itself", () => {
+ const spec: Record<string, unknown> = {};
+ spec["tags"] = spec;
+
+ expect(() => new Dag("cyclic_dag", spec as never)).toThrowError(
+ /The spec for Dag "cyclic_dag" refers back to itself/,
+ );
+ });
+
+ it("accepts the generated Dag and task fields", () => {
+ const dag = new Dag("specced_dag", { schedule: "@daily", tags: ["etl"],
catchup: false });
+ dag.task("extract", async () => undefined, { retries: 2, retryDelay: 30
})();
+
+ expect(dag.spec).toEqual({ schedule: "@daily", tags: ["etl"], catchup:
false });
+ expect(getDagTaskRecords(dag).get("extract")?.spec).toEqual({ retries: 2,
retryDelay: 30 });
+ });
+
+ it.each([
+ ["a misspelling", "scheduled"],
+ // The generated fields are camelCase, so the schema's own spelling is a
+ // typo here rather than a second accepted name.
+ ["the raw schema key", "dag_display_name"],
+ // Identity is positional, so it is not a spec field.
+ ["the positional dag_id", "dagId"],
+ ])("rejects %s in the Dag spec", (_label, key) => {
+ expect(() => new Dag("example_dag", { [key]: "x" } as unknown as
DagSpec)).toThrowError(
+ new RegExp(`Unknown option "${key}" in the spec for Dag "example_dag"`),
+ );
+ });
+
+ it.each([
+ ["a misspelling", "retry"],
+ ["the raw schema key", "retry_delay"],
+ ["the positional task_id", "taskId"],
+ ])("rejects %s in the task spec", (_label, key) => {
+ const dag = new Dag("example_dag");
+ expect(() =>
+ dag.task("transform", async () => undefined, { [key]: 1 } as unknown as
TaskSpec),
+ ).toThrowError(
+ new RegExp(`Unknown option "${key}" in the spec for Dag "example_dag"
task "transform"`),
+ );
+ expect(dag.taskIds).toEqual([]);
+ });
+
it.each([
- ["a populated object", { schedule: "@daily" }],
["null", null],
["an array", []],
["a non-plain object", new Date()],
- ])("rejects a Dag spec that is not an empty object: %s", (_label, spec) => {
- expect(() => new Dag("example_dag", spec as unknown as Record<string,
never>)).toThrowError(
- /spec for Dag "example_dag" must be an empty object/,
+ ])("rejects a Dag spec that is not an options object: %s", (_label, spec) =>
{
+ expect(() => new Dag("example_dag", spec as unknown as
DagSpec)).toThrowError(
+ /spec for Dag "example_dag" must be an object/,
);
});
it.each([
- ["a populated object", { retries: 2 }],
["null", null],
["an array", []],
["a non-plain object", new Date()],
- ])("rejects a task spec that is not an empty object: %s", (_label, spec) => {
+ ])("rejects a task spec that is not an options object: %s", (_label, spec)
=> {
const dag = new Dag("example_dag");
expect(() =>
- dag.task("transform", async () => undefined, {
- spec: spec as unknown as Record<string, never>,
- }),
- ).toThrowError(/spec for Dag "example_dag" task "transform" must be an
empty object/);
+ dag.task("transform", async () => undefined, spec as unknown as
TaskSpec),
+ ).toThrowError(/spec for Dag "example_dag" task "transform" must be an
object/);
expect(dag.taskIds).toEqual([]);
});
@@ -408,28 +490,16 @@ describe("Dag", () => {
it.each([
["inputs, which the factory call now carries", { inputs: {} }],
- ["a misspelled spec key", { specs: {} }],
- ["an upstream reference passed as an option", { upstream: { dagId: "d",
taskId: "t" } }],
- ])("rejects %s in the task options", (_label, options) => {
+ ["a nested spec, left over from the old options object", { spec: {} }],
+ ["an upstream reference", { upstream: { dagId: "d", taskId: "t" } }],
+ ])("rejects %s in the task spec", (_label, spec) => {
const dag = new Dag("example_dag");
expect(() =>
- dag.task("transform", async () => undefined, options as unknown as
Record<string, never>),
- ).toThrowError(/Unknown option ".+" for Dag "example_dag" task
"transform"/);
+ dag.task("transform", async () => undefined, spec as unknown as
TaskSpec),
+ ).toThrowError(/Unknown option ".+" in the spec for Dag "example_dag" task
"transform"/);
expect(dag.taskIds).toEqual([]);
});
- it.each([
- ["null", null],
- ["an array", []],
- ["a string", "spec"],
- ["a non-plain object", new Date()],
- ])("rejects task options that are not an options object: %s", (_label,
options) => {
- const dag = new Dag("example_dag");
- expect(() =>
- dag.task("transform", async () => undefined, options as unknown as
Record<string, never>),
- ).toThrowError(/options for Dag "example_dag" task "transform" must be an
object/);
- });
-
it("rejects duplicate taskIds within a Dag", () => {
const dag = new Dag("example_dag");
dag.task("dup", async () => undefined);
diff --git a/ts-sdk/tsconfig.json b/ts-sdk/tsconfig.json
index d1c294a5f23..a47f167e378 100644
--- a/ts-sdk/tsconfig.json
+++ b/ts-sdk/tsconfig.json
@@ -18,10 +18,15 @@
"resolveJsonModule": true,
"skipLibCheck": true,
+ // The codegen scripts are plain ESM. Including them lets their tests infer
+ // real signatures rather than import an untyped module;
tsconfig.build.json
+ // narrows `include` back to src, so none of this reaches the package.
+ "allowJs": true,
+
"noEmit": true,
"declaration": true,
"declarationMap": true,
"sourceMap": true
},
- "include": ["src/**/*.ts", "tests/**/*.ts"]
+ "include": ["src/**/*.ts", "tests/**/*.ts", "scripts/**/*.mjs"]
}