Skip to content
Closed
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
15 changes: 15 additions & 0 deletions .pre-commit-config.yaml
Original file line number Diff line number Diff line change
Expand Up @@ -337,6 +337,21 @@ repos:
(?x)
^airflow-core/src/airflow/serialization/schema\.json$|
^ts-sdk/schema/dag-schema\.json$
# check-ts-sdk-serialization-conformance in ts-sdk/ covers the SDK's own files; this runs the same
# check when the shared harness or Airflow's serializer changes, which that project cannot see.
- id: check-ts-sdk-serialization-conformance-shared
name: Check the TS SDK serializes Dags the way Airflow does, after a shared change
description: "Serialize the shared test Dags with the TS SDK and with Airflow, and compare the two"
entry: ./ts-sdk/scripts/ci/prek/check_serialization_conformance.py
language: node
additional_dependencies: ['pnpm@10.28.1']
pass_filenames: false
require_serial: true
files: >
(?x)
^airflow-core/src/airflow/serialization/serialized_objects\.py$|
^airflow-core/src/airflow/serialization/schema\.json$|
^scripts/ci/lang_sdk_serialization/.*$
- id: check-go-version-in-sync
name: Check Go toolchain version is consistent across build files
entry: ./scripts/ci/prek/check_go_version_in_sync.py
Expand Down
10 changes: 8 additions & 2 deletions airflow-core/adr/lang-sdk/0008-control-flow-constructs.md
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,8 @@ Proposed.
2. **A branch selects a task, not a string.** A case is the task reference the SDK already handed back, and the value on the wire is that task's task_id, so the compiler checks the candidate exists and no label has to be kept in step with it.
3. **No default case.** `BranchPythonOperator` has none to serialize, and a one-sided `If` whose condition is false follows nothing at all.
A decider returning a ref that is not one of the declared cases is a run-time error the SDK raises, a narrower check than `skip_all_except`, which only rejects a task_id missing from the whole Dag.
4. **Triggering a Dag run is an ordinary DSL task**, it is pure DSL purpose instead of a new runtime.
4. **Triggering a Dag run is an ordinary task that the SDK's own runtime runs.** It has no handler.
The runtime sends the `TriggerDagRun` request itself and follows `TriggerDagRunOperator` for every option it offers, including `wait_for_completion` and `deferrable`.
5. **Grouping keeps Python's semantics**: a scope offering the same task and nesting methods as the Dag, prefixing each task_id with the group id (`prefix_group_id`),
and can be ordered against a task or another group, as `group1 >> group2` does in Python.

Expand Down Expand Up @@ -89,7 +90,12 @@ That operator defaults to skipping every task in its downstream closure and igno
a single-ref return cannot express that, and no Lang SDK offers it for now.
An author needing several paths together puts them behind one task, or gates each with its own condition.
This limitation is accepted rather than open.
- **No Lang SDK needs a deferral mechanism for now** to offer `deferrable` or `wait_for_completion`, because the trigger task runs in Python.
- **Each Lang SDK runtime implements the trigger task**, so a native Dag needs no Python worker.
`deferrable` needs no trigger written in that language: the runtime defers to the standard provider's `DagStateTrigger`, which runs in the Python triggerer, and the task keeps its queue, so it resumes in the same runtime with `next_method` set.
- **The runtime reads the config the trigger task needs from its environment**, since it cannot read `airflow.cfg`.
It reads `AIRFLOW__API__BASE_URL`, `AIRFLOW__OPERATORS__DEFAULT_DEFERRABLE` and `AIRFLOW__TRIGGERER__QUEUES_ENABLED`, and an unset one takes Python's fallback.
The coordinator exports them from Airflow's config (#73842).
- **The trigger's arguments are not Jinja-rendered**, because no Python code runs the task. A value is sent as the author wrote it.
- **Nothing extra reaches the Dag JSON**, which carries no branch-candidate field at all. A ref is a task_id by the time the decision is sent, so each SDK stores its cases and nothing else.
- **A group edge needs one base type per SDK** that both a task and a group satisfy, since either can sit at the end of an edge.
Python already has it: `TaskGroup(TaskGroupMixin, DAGNode)` (`task-sdk/src/airflow/sdk/definitions/taskgroup.py:96`) and every operator inherit `DependencyMixin` (`.../definitions/_internal/mixins.py:35`), where `set_upstream` and `set_downstream` live.
Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -311,6 +311,153 @@ already exists changes nothing. Each returns the reference it was called on, so
Pass a value as an argument when the downstream task needs it, and use ``before`` or ``after`` when
it only needs to run in order.

Task groups
~~~~~~~~~~~

``dag.taskGroup(groupId)`` opens a scope with the same ``task`` and ``taskGroup`` methods as the Dag,
prefixing the id of everything declared in it, as Python's ``prefix_group_id`` does:

.. code-block:: typescript

const staging = dag.taskGroup("staging");
staging.task("stage_rows", stageRows)(); // task id "staging.stage_rows"
staging.taskGroup("checks").task("nulls", checkNulls)(); // "staging.checks.nulls"

staging.before(loaded); // staging >> loaded

A group is an edge endpoint in its own right, so ``before`` and ``after`` order a whole group against
a task or against another group.

Tasks and groups share one id namespace, as they do in Python, so a Dag cannot hold both a task and a
group called ``staging``. A ``.`` is what separates a group from what it holds, so it cannot appear in
an id of either.

Pass ``{ prefixGroupId: false }`` to keep the ids declared in a group as written, as ``prefix_group_id=False``
does in Python; they then have to be unique across the Dag. A group id is made of letters, digits, dashes and
underscores, and is at most 200 characters.

Serialization
~~~~~~~~~~~~~

A native Dag serializes into the same Dag JSON a Python Dag produces, so the scheduler reads it
without knowing which language declared it.

``schedule`` accepts what maps to a stock timetable: unset, ``@once``, ``@continuous``, or a cron
expression. A cron preset such as ``@daily`` is recorded as the expression it stands for. Anything
else names a Python object a TypeScript bundle cannot point at, and is rejected.

Every task of a native Dag runs on the Node coordinator, so it needs the queue the deployment routes
there. Set it once on the Dag and each task inherits it:

.. code-block:: typescript

const dag = new Dag("ts_etl", { schedule: "@daily", queue: "typescript" });

// ...and one task that needs its own.
dag.task("heavy", heavyHandler, { queue: "typescript_large" })();

``queue`` on a task wins over the Dag's. See :ref:`typescript-sdk/coordinator-config` for the
``queue_to_coordinator`` entry that sends that queue to the coordinator.

Conditional branching
~~~~~~~~~~~~~~~~~~~~~

``dag.if`` takes a task whose handler returns a boolean, and names the task each outcome runs:

.. code-block:: typescript

const condition = dag.task("has_rows", async ({ rows }: { rows: number }) => rows > 0);
const gated = condition({ rows: extracted });

dag.if(gated).then(loaded).else(reportedEmpty);

The condition is an ordinary task, so it is declared, typed and wired like any other, and the
compiler checks that its handler really returns a boolean. ``else`` is optional: a one-sided
condition skips its own branch when the condition fails and follows nothing.

A guarded task takes no argument for the control edge, because a condition's boolean decides whether
the task runs rather than what it runs on. Read a value from the condition with
``getClient().getXCom``.

The side not taken is skipped when the run reaches it, and stays skipped if you clear it later. Only
the branches named here are skipped, so a task that several branches converge on still runs — unlike
Python's ``@task.branch``, which skips every immediate downstream it did not follow.

Multi-way branching
~~~~~~~~~~~~~~~~~~~

``dag.switch`` is the multi-way form: a task whose handler returns one of the cases it is given.

.. code-block:: typescript

const decider = dag.task("pick_path", async ({ rows }: { rows: number }) =>
rows > 1000 ? handleLong : handleShort,
);
const picked = decider({ rows: extracted });

dag.switch(picked).case(handleLong).case(handleShort);

A case is the task reference itself, so the compiler checks the candidate exists and renaming a
handler cannot silently rewire a Dag. The task's own value is the chosen task's id, which a
downstream task can read from its XCom.

There is no default case. A decider that returns anything outside its cases fails the task, naming
what it chose and what it could have chosen.

Exactly one case is selected. Python's branch callable may return a list of task ids, and no language
SDK offers that yet: put the paths that run together behind one task, or gate each with its own
condition.

Triggering another Dag
~~~~~~~~~~~~~~~~~~~~~~

``dag.triggerDagRun`` starts another Dag's run, as a task of this one:

.. code-block:: typescript

const trigger = dag.triggerDagRun({
taskId: "trigger_downstream",
dagId: "downstream_etl",
waitForCompletion: true,
});

trigger.after(loaded);

There is nothing to call: the task takes no TypeScript arguments, so it hands back the reference
directly. A second options object carries the task's own spec, as ``dag.task`` takes one.

The task runs in the TypeScript runtime, like the Dag's other tasks, and inherits the Dag's queue. It
behaves as ``TriggerDagRunOperator`` does. ``waitForCompletion`` polls the run every ``pokeInterval``
seconds until it reaches one of ``allowedStates`` or ``failedStates``. With ``deferrable`` as well, the
task defers to ``DagStateTrigger`` instead of holding a worker slot, and resumes in this runtime when
the run finishes. ``DagStateTrigger`` runs in the Python triggerer, so the triggerer needs the standard
provider installed.

The task pushes the ``trigger_run_id`` XCom, and the task's "Triggered DAG" link opens the run it
started.

The task follows these config options, with Python's fallback when one is unset:

- ``[api] base_url`` is the base of the "Triggered DAG" link. It falls back to ``/``.
- ``[operators] default_deferrable`` is the default of ``deferrable``. It falls back to ``false``.
- ``[triggerer] queues_enabled`` decides whether the deferred trigger gets the task's queue. It falls
back to ``false``, so the trigger gets none.

Some of what ``TriggerDagRunOperator`` does is not offered:

- ``logical_date`` and ``run_after``. The triggered run's logical date is the time the task triggers
it, as in Python when neither is set.
- Jinja. Values are sent as written, so ``{{ ds }}`` in ``conf`` reaches the triggered run as that
literal string.
- OpenLineage parent injection (``openlineage_inject_parent_info``). The runtime does not add the
parent task's OpenLineage details to ``conf``.

A worked example
~~~~~~~~~~~~~~~~

``ts-sdk/example/src/native.ts`` puts the constructs above into one Dag, registered on the same
bundle as the mixed-language handlers beside it, so a single artifact serves both authoring modes.

``new Dag`` and ``dag.task`` both take a trailing spec of Airflow options:
``{ schedule: "@daily", tags: ["etl"] }`` for the Dag, ``{ retries: 2, retryDelay: 30 }`` for a task.

Expand Down
Original file line number Diff line number Diff line change
Expand Up @@ -19,6 +19,7 @@
import time
from datetime import datetime, timezone
from functools import cached_property
from http import HTTPStatus

import boto3
import requests
Expand Down Expand Up @@ -163,6 +164,21 @@ def get_variable(self, key: str):
"""Get an Airflow Variable via API."""
return self._make_request(method="GET", endpoint=f"variables/{key}")

def set_variable(self, key: str, value: str, description: str | None = None):
"""Create or replace an Airflow Variable via API."""
body = {"key": key, "value": value, "description": description}
try:
return self._make_request(method="POST", endpoint="variables", json=body)
except requests.HTTPError as exc:
# 409 == it already exists, from an earlier run of the same suite.
if exc.response is None or exc.response.status_code != HTTPStatus.CONFLICT:
raise
return self._make_request(method="PATCH", endpoint=f"variables/{key}", json=body)

def get_tasks(self, dag_id: str):
"""List a Dag's tasks, with the edges each one carries."""
return self._make_request(method="GET", endpoint=f"dags/{dag_id}/tasks")

def trigger_dag_and_wait(self, dag_id: str, json=None):
"""Trigger a DAG and wait for it to complete."""
self.un_pause_dag(dag_id)
Expand Down Expand Up @@ -201,6 +217,13 @@ def get_task_logs(
endpoint=endpoint,
)

def get_event_logs(self, dag_id: str, run_id: str, task_id: str | None = None) -> list[dict]:
"""List the audit log events of a Dag run, or of one of its tasks, oldest first."""
params = {"dag_id": dag_id, "run_id": run_id, "order_by": "event_log_id"}
if task_id is not None:
params["task_id"] = task_id
return self._make_request(method="GET", endpoint="eventLogs", params=params)["event_logs"]


class TaskSDKClient:
"""Client for interacting with the Task SDK API."""
Expand Down
Loading