Skip to content

Parse native Lang-SDK Dags with their coordinator's runtime - #73842

Merged
jason810496 merged 7 commits into
jason/lang-sdk-e2e/02-importer-registryfrom
jason/lang-sdk-e2e/03-native-dag-parse
Oct 1, 2026
Merged

jason810496 merged 7 commits into
jason/lang-sdk-e2e/02-importer-registryfrom
jason/lang-sdk-e2e/03-native-dag-parse

Conversation

@jason810496

@jason810496 jason810496 commented Sep 28, 2026 •

Copy link
Copy Markdown
Member

Stack (bottom to top): #73841, #74004, #73842, #73843, #73844, #73845, #73846, #73847

related: #71929

Part of the native Dag e2e stack. Builds on #74004, the shared cycle check (layer diff).

Why

Nothing can parse a Dag written entirely in a Lang SDK yet. This is the parse_dag half of #71929.

What changes

  • Routing. A coordinator's get_dag_importer() serves its dag_bundle_name bundle, or every bundle if it sets neither that nor an explicit root. Two coordinators that claim one extension in a bundle are rejected when [sdk] coordinators is loaded, before either is built.
  • Dag processor. LangSDKDagFileProcessorProcess, in the new dag_processing/lang_sdk_processor.py, parses a file a CoordinatorDagImporter claims. Its child finds the coordinator and execs the runtime, which connects back to the Dag processor over TCP and answers the parse request. It shares the new BaseDagFileProcessorProcess with the Python processor, so fork, the macOS exec path and log forwarding work as for a Python file. A runtime still running 5s after its result is killed. Callbacks for the file are dropped.
  • Config defaults. A runtime cannot read the Airflow config, so a native Dag may leave out max_active_tasks, max_active_runs, max_consecutive_failed_dag_runs, catchup and disable_bundle_versioning. The new DagSerialization.fill_config_defaults fills each unset one from [core] max_active_tasks_per_dag, [core] max_active_runs_per_dag, [core] max_consecutive_failed_dag_runs_per_dag, [scheduler] catchup_by_default and [dag_processor] disable_bundle_versioning, as a Python Dag does. A value the Dag sets is kept, and the stored Dag carries the filled values. The field table of airflow-core/adr/lang-sdk/0004-dag-parsing.md says so.
  • Validation. After the fill, each returned Dag must pass the new DagSerialization.validate_serialized_dag: the JSON schema, a load, and no cycle in the task graph, checked with detect_cycle from Share Dag cycle detection between the Task SDK and core #74004. A failed start, a missing result, an invalid frame or message, or an invalid Dag gives an import error.
  • Import timeout. The parse child resolves get_dagbag_import_timeout for the file, as for a Python Dag file. A parse past it is killed and gives an import error.
  • Runtime environment. Like the log levels, a runtime now gets the resolved [api] base_url, [operators] default_deferrable and [triggerer] queues_enabled, with the fallbacks Python uses.
  • Dag bag. CoordinatorDagImporter.import_definition runs the same process, with the same fill and validation, and returns each Dag as the SerializedDAG the scheduler loads. A Dag bag has no API client, so a runtime request other than the parse result or MaskSecret gets an error. The Dag bag keeps its duplicate id check and skips the SDK-only checks and cluster policies. sync_bag_to_db, which airflow dags reserialize uses, leaves a coordinator's files to the Dag processor.
  • Python cannot run it. airflow dags test, tasks test, tasks render and tasks list refuse a native Dag. A Python worker handed one of its tasks exits with an error naming the queue to route, before it builds a Dag bag.

Commits

  1. Let a coordinator parse the Dag files of the bundles it serves
  2. Check that a serialized Dag can be stored and loaded
  3. Move the Dag file processor's shared plumbing into a base class (no behavior change)
  4. Parse coordinator-claimed Dag files with their runtime
  5. Fill a native Dag's unset settings from the Airflow config
  6. Report a Lang-SDK parse past its import timeout as an import error
  7. Report a native Dag's task that reaches a Python worker

Limitations

  • Cluster policies do not run on a native Dag.
  • Only the checks above run on a native Dag. A rule the Python DAG class enforces, such as max_active_runs <= 1 for a @continuous schedule, is not checked.
  • A task of a native Dag runs only on its coordinator, so a Python operator task inside one cannot run on a Python worker.

Deviations from #71929

  • get_dag_importer() returns None by default, so a coordinator opts in; Java and TypeScript do so in the next layers.
  • The Python parse child execs the runtime and hands it TCP addresses, so no fd-0 bridge stays in between. parse_dag is on SubprocessCoordinator only.
  • A Dag bag holds a native Dag as a SerializedDAG, not an airflow.sdk.DAG.

Example

After the next layers, a coordinator opts in like this:

class NodeDagImporter(CoordinatorDagImporter):
    artifact_suffix = ".min.mjs"
    supported_extensions = [".mjs"]  # a registry routes by the last suffix

    def get_source_code(self, definition): ...


class NodeCoordinator(SubprocessCoordinator):
    @classmethod
    def get_dag_importer_class(cls):
        return NodeDagImporter

    def _build_parse_dag_command(self, *, path):
        # parse_dag appends --comm and --logs
        return [self.node_executable, os.fspath(path)], read_bundle(path).supervisor_schema_version
[sdk]
coordinators = {
    "ts-native": {"classpath": "airflow.sdk.coordinators.node.NodeCoordinator", "kwargs": {"node_executable": "/usr/bin/node"}},
    "java-native": {"classpath": "airflow.sdk.coordinators.java.JavaCoordinator", "kwargs": {"dag_bundle_name": "java-dags"}}
  }

How to test

The core tests replace the coordinator's parse_dag, so the forked parse child plays the runtime (fake_lang_sdk.play_runtime).

uv run --project airflow-core pytest airflow-core/tests/unit/dag_processing/test_lang_sdk_processor.py airflow-core/tests/unit/dag_processing/test_processor.py airflow-core/tests/unit/dag_processing/test_manager.py airflow-core/tests/unit/dag_processing/test_dagbag.py airflow-core/tests/unit/dag_processing/test_importer_routing.py airflow-core/tests/unit/utils/test_cli_util.py -xvs
uv run --project airflow-core pytest airflow-core/tests/unit/serialization/test_dag_serialization.py -k "ValidateSerializedDag or FillConfigDefaults" -xvs
uv run --project task-sdk pytest task-sdk/tests/task_sdk/coordinators task-sdk/tests/task_sdk/importers task-sdk/tests/task_sdk/execution_time/test_coordinator.py -xvs
uv run --project task-sdk pytest task-sdk/tests/task_sdk/execution_time/test_task_runner.py -k native -xvs

Follow-ups


Was generative AI tooling used to co-author this PR?

@boring-cyborg boring-cyborg Bot added area:coordinator Coordinator: The interface to spawn Lang-SDK subprocesses area:DAG-processing area:task-sdk labels Sep 28, 2026
@jason810496
jason810496 added this pull request to stack #73848 September 28, 2026 13:20
@jason810496 jason810496 changed the title jason/lang sdk e2e/03 native dag parse Parse native Lang-SDK Dags with their coordinator's runtime Sep 28, 2026
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/02-importer-registry branch from 82ee9ae to 40bf9b3 Compare September 28, 2026 15:57
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/03-native-dag-parse branch from 6242f83 to 8e8c982 Compare September 28, 2026 15:57
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/02-importer-registry branch from 40bf9b3 to a73513d Compare September 29, 2026 01:25
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/03-native-dag-parse branch from 8e8c982 to 0080efb Compare September 29, 2026 01:26
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/02-importer-registry branch from a73513d to 98bfa5b Compare September 29, 2026 01:29
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/03-native-dag-parse branch from 0080efb to c4b3431 Compare September 29, 2026 01:29
@jason810496
jason810496 removed this pull request from stack #73848 October 1, 2026 03:02
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/02-importer-registry branch from 98bfa5b to 57c9a34 Compare October 1, 2026 03:03
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/03-native-dag-parse branch from c4b3431 to e6de2df Compare October 1, 2026 03:03
@jason810496
jason810496 changed the base branch from jason/lang-sdk-e2e/02-importer-registry to jason/lang-sdk-e2e/02b-shared-dag-cycle-detection October 1, 2026 03:03
@jason810496
jason810496 added this pull request to stack #74005 October 1, 2026 03:04
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/03-native-dag-parse branch 2 times, most recently from ec6f22e to c35472f Compare October 1, 2026 12:41
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/02b-shared-dag-cycle-detection branch from e14a584 to 7caeacb Compare October 1, 2026 12:41
A coordinator opts in to parsing native Dags by returning a Dag importer
from get_dag_importer(). It serves the Dag bundle its dag_bundle_name
names, or every bundle when it sets neither that nor an explicit root.
Two coordinators claiming one extension in a bundle is a config error.

SubprocessCoordinator.parse_dag execs the runtime that parses a file.
The runtime inherits only the standard streams and connects back to the
given addresses.

Like the log levels, a runtime now gets the resolved [api] base_url,
[operators] default_deferrable and [triggerer] queues_enabled in its
environment, with the fallbacks Python uses.
DagSerialization.validate_serialized_dag checks the JSON schema,
deserializes a copy and looks for a cycle in the task graph. It raises
DeserializationError naming the Dag, and for a cycle a task on it. It is
meant for a serialized Dag that no Python code built, such as one a
Lang-SDK runtime returns, so the Python SDK's cycle check never ran on
it. The cycle check runs the shared detect_cycle over the serialized
downstream edges.
BaseDagFileProcessorProcess holds the parse log forwarding, the request
handlers, readiness and cleanup. DagFileProcessorProcess keeps only how
the Python parse child starts. No behavior change.
The Dag processor parses a file a coordinator's Dag importer claims with
LangSDKDagFileProcessorProcess, in the new lang_sdk_processor module.
Its Python child finds the coordinator, reports the runtime's schema
version over fd 0 and execs the runtime, which connects back to two
listeners the manager owns and answers the parse request. Each Dag must
pass validate_serialized_dag. A failed start, a missing result, an
invalid frame or message, or an invalid Dag is an import error. A
runtime still running 5s after its result is killed, and the result is
kept. Callbacks for the file are dropped.

CoordinatorDagImporter.import_definition runs the same process for a Dag
bag and returns each Dag as a SerializedDAG. A Dag bag has no API
client, so the runtime's requests get an error there. The Dag bag skips
the SDK-only checks and cluster policies for such a Dag, sync_bag_to_db
does not store the Dags or import errors of such a file, and CLI
commands that run a Dag refuse it.
A Lang-SDK runtime cannot read the Airflow config, so it leaves out the
Dag settings that a Python Dag reads from it when unset. For each Dag
the runtime returns, the Dag processor now fills max_active_tasks,
max_active_runs, max_consecutive_failed_dag_runs, catchup and
disable_bundle_versioning from the config with
DagSerialization.fill_config_defaults, before it validates the Dag. A
value the Dag sets is kept, and the stored Dag carries the filled
values. A Dag bag gets the same.
The parse child resolves the get_dagbag_import_timeout policy for the
file and sends the timeout with the runtime's schema version. The policy
thus runs where it does for a Python Dag file, and a failure is only
that file's import error. The timeout counts from the start of the
parse. Past it, with no result yet, the runtime is killed and the file
gets an import error.

LangSDKDagFileProcessorProcess.run applies the same timeout, so it no
longer takes one. Until the child reports it, dag_file_processor_timeout
bounds the parse, as it does in the Dag processor. At the deadline run
stops waiting even if an exited runtime left a process that holds its
output open.
A Python worker handed a task from a Dag file that a coordinator's Dag
importer claims now exits with an error naming the queue to route. It
checks before building a Dag bag, so it never starts the runtime.
@jason810496
jason810496 removed this pull request from stack #74005 October 1, 2026 13:17
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/03-native-dag-parse branch from c35472f to 6f91c3e Compare October 1, 2026 13:18
@jason810496
jason810496 force-pushed the jason/lang-sdk-e2e/02b-shared-dag-cycle-detection branch from 7caeacb to 5c20182 Compare October 1, 2026 13:18
Base automatically changed from jason/lang-sdk-e2e/02b-shared-dag-cycle-detection to jason/lang-sdk-e2e/02-importer-registry October 1, 2026 13:18
@jason810496
jason810496 merged commit 6f91c3e into jason/lang-sdk-e2e/02-importer-registry Oct 1, 2026
@jason810496
jason810496 deleted the jason/lang-sdk-e2e/03-native-dag-parse branch October 1, 2026 13:18
@jason810496
jason810496 restored the jason/lang-sdk-e2e/03-native-dag-parse branch October 1, 2026 13:36
@jason810496

Copy link
Copy Markdown
Member Author

Replaced by #74035. GitHub marked this PR as merged into a stack branch when the stack was reordered; nothing was merged into main.

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

area:coordinator Coordinator: The interface to spawn Lang-SDK subprocesses area:DAG-processing area:task-sdk

Projects

None yet

Development

Successfully merging this pull request may close these issues.

1 participant