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]

Reply via email to