jason810496 opened a new pull request, #74035: URL: https://github.com/apache/airflow/pull/74035
Replaces #73842, which GitHub closed as merged into a stack branch when the stack was reordered. 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](https://github.com/apache/airflow/compare/jason/lang-sdk-e2e/02b-shared-dag-cycle-detection...jason/lang-sdk-e2e/03-native-dag-parse)). ## 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 #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: ```python 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 ``` ```ini [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`). ```bash 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 - Java `Serde.kt` (#71190) still writes the stock values of the five settings above, so a native Java Dag does not get the config values yet. The Go SDK serializer (#67155) writes them too. - The Go SDK does not parse native Dags yet: its coordinator hands out no Dag importer, and a packed Go bundle has no file extension to claim. - Decision 2 of `airflow-core/adr/lang-sdk/0009-provider-operators-as-generated-dsl.md` says a provider DSL task runs on a Python worker, which a task of a native Dag cannot do. It needs its own amendment. --- ##### Was generative AI tooling used to co-author this PR? - [x] Yes, with help of Claude Code Opus 5.5 following [the guidelines](https://github.com/apache/airflow/blob/main/contributing-docs/05_pull_requests.rst#gen-ai-assisted-contributions) -- This is an automated message from the Apache Git Service. To respond to the message, please log on to GitHub and use the URL above to go to the specific comment. To unsubscribe, e-mail: [email protected] For queries about this service, please contact Infrastructure at: [email protected]
