kaxil commented on code in PR #74341: URL: https://github.com/apache/airflow/pull/74341#discussion_r4206220147
########## task-sdk/src/airflow/sdk/coordinators/executable/_dag_importer.py: ########## @@ -0,0 +1,106 @@ +# +# 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. +"""The Dag importer of :class:`~airflow.sdk.coordinators.executable.ExecutableCoordinator`.""" + +from __future__ import annotations + +import os +from pathlib import Path +from typing import TYPE_CHECKING, ClassVar, Final + +from airflow.sdk._shared.module_loading.file_discovery import find_path_from_directory +from airflow.sdk.configuration import conf +from airflow.sdk.coordinators._dag_importer import CoordinatorDagImporter +from airflow.sdk.coordinators.executable._bundle_reader import ( + read_bundle_entrypoint_source, + read_bundle_language, + read_bundle_source, +) +from airflow.sdk.coordinators.executable.coordinator import FOOTER_MAGIC +from airflow.sdk.importers.base import DagSourceCode, FilesystemDagDefinition + +if TYPE_CHECKING: + from collections.abc import Iterator + + from airflow.dag_processing.bundles.base import BaseDagBundle # noqa: SDK002 + from airflow.sdk.importers.base import DagDefinition + +_NO_SOURCE: Final = "// Source code is not available: the bundle embeds no source.\n" +_DEFAULT_LANGUAGE: Final = "text" + + +class ExecutableDagImporter(CoordinatorDagImporter): + """ + Claim the native Dags of executable bundles, such as the ones the Go SDK packs. + + An :class:`~airflow.sdk.coordinators.executable.ExecutableCoordinator` parses them. A bundle is + identified by the ``AFBNDL01`` magic that ends the file, not by its name, so it can have any extension + or none. + """ + + coordinator_classpath: ClassVar[str] = "airflow.sdk.coordinators.executable.ExecutableCoordinator" + artifact_suffix = "" + supported_extensions: list[str] = [] + + def can_handle(self, definition: DagDefinition | str | Path) -> bool: + try: + with open(str(definition), "rb") as bundle_file: + bundle_file.seek(-len(FOOTER_MAGIC), os.SEEK_END) + return bundle_file.read(len(FOOTER_MAGIC)) == FOOTER_MAGIC + except OSError: + return False + + def list_dag_definitions( + self, bundle: BaseDagBundle, *, safe_mode: bool = True + ) -> Iterator[FilesystemDagDefinition]: + root = Path(bundle.path) + if root.is_file(): + paths: Iterator[Path] = iter([root]) + else: + ignore_file_syntax = conf.get_mandatory_value("core", "DAG_IGNORE_FILE_SYNTAX", fallback="glob") + paths = (Path(p) for p in find_path_from_directory(root, ".airflowignore", ignore_file_syntax)) + for path in paths: + if path.is_file() and self.can_handle(path): + definition = FilesystemDagDefinition(path=path) + if self.might_contain_dag(definition, safe_mode): + yield definition + + def might_contain_dag(self, definition: DagDefinition, safe_mode: bool) -> bool: + """ + Return ``True`` for every bundle. + + ``safe_mode`` does not apply: whether a bundle defines Dags is known only by running it, and a + bundle that only registers task handlers is parsed too. A file that cannot be read is kept, so + that parsing it reports why. + """ + return True + + def get_source_code(self, definition: DagDefinition, dag_id: str | None = None) -> DagSourceCode: + """ + Return the embedded source of *dag_id*'s own file, or a notice when the bundle embeds none. + + Without *dag_id*, or for one the bundle maps to no file, this is the entrypoint source. + """ + with definition.as_file() as path: Review Comment: `_read_dag_source_codes` calls this once per Dag from the parse-result handler, which runs in the Dag processor manager. Each call opens and verifies the bundle twice, `yaml.safe_load`s the whole manifest twice, and re-hashes every embedded source region to return one of them. The manifest grows with the Dag count, so this is quadratic. Using `write_bundle` from this PR's test utils with 8 KB per source, 100 Dags took about 5 s and 300 Dags about 45 s per parse result locally, all of it blocking the manager loop. Could the bundle be read and verified once per file and reused across Dags (cached the way `_digest_cache` keys the binary hash), or read only the requested region? ########## task-sdk/src/airflow/sdk/coordinators/executable/coordinator.py: ########## @@ -377,3 +413,23 @@ def _build_execute_task_command(self, *, what: TaskInstance) -> tuple[list[str], roots = self._get_scan_roots() bundle = _Bundle.find(roots, what.dag_id) return [str(bundle.path)], bundle.schema_version + + def _build_bundle_command(self, path: pathlib.Path) -> tuple[list[str], str | None]: + """Return the command that runs the verified bundle at *path*, and its supervisor schema version.""" + if (metadata := _read_bundle_metadata(path)) is None: + raise ValueError(f"{path} is not a valid executable bundle") Review Comment: When `_open_verified_bundle` rejects the file, the reason (binary SHA-256 mismatch, unknown `footer_ver`, undecodable manifest) only goes to a debug log, so the parse import error says just "is not a valid executable bundle". A binary that was `strip`ped or `codesign`ed after packing hits the hash mismatch, and nothing in the UI says so. Could the parse path raise with the reason, e.g. a verify function that raises, wrapped by the catching version `_Bundle.find` keeps for scanning? ########## task-sdk/src/airflow/sdk/coordinators/executable/_dag_importer.py: ########## @@ -0,0 +1,106 @@ +# +# 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. +"""The Dag importer of :class:`~airflow.sdk.coordinators.executable.ExecutableCoordinator`.""" + +from __future__ import annotations + +import os +from pathlib import Path +from typing import TYPE_CHECKING, ClassVar, Final + +from airflow.sdk._shared.module_loading.file_discovery import find_path_from_directory +from airflow.sdk.configuration import conf +from airflow.sdk.coordinators._dag_importer import CoordinatorDagImporter +from airflow.sdk.coordinators.executable._bundle_reader import ( + read_bundle_entrypoint_source, + read_bundle_language, + read_bundle_source, +) +from airflow.sdk.coordinators.executable.coordinator import FOOTER_MAGIC +from airflow.sdk.importers.base import DagSourceCode, FilesystemDagDefinition + +if TYPE_CHECKING: + from collections.abc import Iterator + + from airflow.dag_processing.bundles.base import BaseDagBundle # noqa: SDK002 + from airflow.sdk.importers.base import DagDefinition + +_NO_SOURCE: Final = "// Source code is not available: the bundle embeds no source.\n" +_DEFAULT_LANGUAGE: Final = "text" + + +class ExecutableDagImporter(CoordinatorDagImporter): + """ + Claim the native Dags of executable bundles, such as the ones the Go SDK packs. + + An :class:`~airflow.sdk.coordinators.executable.ExecutableCoordinator` parses them. A bundle is + identified by the ``AFBNDL01`` magic that ends the file, not by its name, so it can have any extension + or none. + """ + + coordinator_classpath: ClassVar[str] = "airflow.sdk.coordinators.executable.ExecutableCoordinator" + artifact_suffix = "" + supported_extensions: list[str] = [] + + def can_handle(self, definition: DagDefinition | str | Path) -> bool: + try: + with open(str(definition), "rb") as bundle_file: + bundle_file.seek(-len(FOOTER_MAGIC), os.SEEK_END) + return bundle_file.read(len(FOOTER_MAGIC)) == FOOTER_MAGIC + except OSError: + return False + + def list_dag_definitions( + self, bundle: BaseDagBundle, *, safe_mode: bool = True + ) -> Iterator[FilesystemDagDefinition]: + root = Path(bundle.path) + if root.is_file(): + paths: Iterator[Path] = iter([root]) + else: + ignore_file_syntax = conf.get_mandatory_value("core", "DAG_IGNORE_FILE_SYNTAX", fallback="glob") + paths = (Path(p) for p in find_path_from_directory(root, ".airflowignore", ignore_file_syntax)) + for path in paths: + if path.is_file() and self.can_handle(path): + definition = FilesystemDagDefinition(path=path) + if self.might_contain_dag(definition, safe_mode): + yield definition + + def might_contain_dag(self, definition: DagDefinition, safe_mode: bool) -> bool: + """ + Return ``True`` for every bundle. + + ``safe_mode`` does not apply: whether a bundle defines Dags is known only by running it, and a + bundle that only registers task handlers is parsed too. A file that cannot be read is kept, so Review Comment: `can_handle` already returns False for a file it can't open (line 65), so `list_dag_definitions` never yields one and the "a file that cannot be read is kept" sentence doesn't hold. ########## task-sdk/src/airflow/sdk/coordinators/executable/coordinator.py: ########## @@ -377,3 +413,23 @@ def _build_execute_task_command(self, *, what: TaskInstance) -> tuple[list[str], roots = self._get_scan_roots() bundle = _Bundle.find(roots, what.dag_id) return [str(bundle.path)], bundle.schema_version + + def _build_bundle_command(self, path: pathlib.Path) -> tuple[list[str], str | None]: + """Return the command that runs the verified bundle at *path*, and its supervisor schema version.""" + if (metadata := _read_bundle_metadata(path)) is None: + raise ValueError(f"{path} is not a valid executable bundle") + try: + bundle = _Bundle(path=path.resolve(), schema_version=extract_supervisor_schema_version(metadata)) + except (TypeError, ValueError) as exc: + raise ValueError(f"Bundle {path} has no usable supervisor schema version: {exc}") from exc + if (reason := _ensure_executable(path)) is not None: + raise ValueError(f"Cannot run bundle {path}: {reason}") + return [str(bundle.path)], bundle.schema_version + + def _build_dag_file_command( + self, *, what: TaskInstance, path: pathlib.Path + ) -> tuple[list[str], str | None]: + return self._build_bundle_command(path) + + def _build_parse_dag_command(self, *, path: pathlib.Path) -> tuple[list[str], str | None]: + return self._build_bundle_command(path) Review Comment: Java's `_build_parse_dag_command` refuses a JAR whose supervisor schema predates Dag parsing (`_DAG_PARSING_SCHEMA_VERSION`) with "rebuild it with a newer Java SDK or list it in .airflowignore". Should this do the same? Every released Go SDK (v1.0.0-beta1 to beta3) declares `2026-06-16`, so after an upgrade each existing handler binary that isn't in `.airflowignore` is spawned on every parse and ends in a generic Lang-SDK import error that doesn't say what to do. Sharing the constant with Java and checking it here would skip the spawn and give the actionable message. ########## task-sdk/tests/task_sdk/coordinators/executable/test_coordinator.py: ########## @@ -447,6 +447,102 @@ def test_build_command_scans_passed_roots_in_colocated_mode(self, tmp_path): assert schema_version == "2026-06-16" +class TestBuildParseDagCommand: + def test_returns_the_bundle_and_its_schema_version(self, tmp_path): + binary = _build_bundle(tmp_path / "my_bundle", dag_ids=["native_dag"]) + + command, schema_version = ExecutableCoordinator()._build_parse_dag_command(path=binary) + + assert command == [str(binary.resolve())] + assert schema_version == "2026-06-16" + + def test_marks_the_bundle_executable(self, tmp_path): + binary = _build_bundle(tmp_path / "my_bundle") + binary.chmod(0o644) + + ExecutableCoordinator()._build_parse_dag_command(path=binary) + + assert os.access(binary, os.X_OK) + + def test_parses_a_bundle_that_registers_no_dag(self, tmp_path): + binary = _build_bundle(tmp_path / "handlers_only", dag_ids=[]) + + command, _ = ExecutableCoordinator()._build_parse_dag_command(path=binary) + + assert command == [str(binary.resolve())] + + def test_raises_for_a_file_that_is_not_a_bundle(self, tmp_path): + plain = tmp_path / "plain" + plain.write_bytes(b"not a bundle") + + with pytest.raises(ValueError, match="is not a valid executable bundle"): + ExecutableCoordinator()._build_parse_dag_command(path=plain) + + def test_raises_for_a_tampered_bundle(self, tmp_path): + binary = _build_bundle(tmp_path / "tampered") + data = bytearray(binary.read_bytes()) + data[0] ^= 0xFF + binary.write_bytes(bytes(data)) + _digest_cache.clear() + + with pytest.raises(ValueError, match="is not a valid executable bundle"): + ExecutableCoordinator()._build_parse_dag_command(path=binary) + + def test_raises_when_the_bundle_omits_the_schema_version(self, tmp_path): + metadata = _make_metadata(["native_dag"]) + del metadata["sdk"]["supervisor_schema_version"] + binary = _build_bundle(tmp_path / "no_schema", metadata=metadata) + + with pytest.raises(ValueError, match="no usable supervisor schema version"): + ExecutableCoordinator()._build_parse_dag_command(path=binary) + + def test_raises_for_an_unknown_schema_version(self, tmp_path): + metadata = _make_metadata(["native_dag"]) + metadata["sdk"]["supervisor_schema_version"] = "1999-01-01" + binary = _build_bundle(tmp_path / "unknown_schema", metadata=metadata) + + with pytest.raises(ValueError, match="no usable supervisor schema version"): + ExecutableCoordinator()._build_parse_dag_command(path=binary) + + def test_raises_when_the_bundle_cannot_be_made_executable(self, tmp_path): + binary = _build_bundle(tmp_path / "locked") + + with ( + patch( Review Comment: Nit: `autospec=True` on this patch, so a signature change to `_ensure_executable` fails the test. ########## task-sdk/docs/executable-bundle-spec.rst: ########## @@ -156,9 +164,17 @@ Reader algorithm: re-hash on every exec; a cache miss (file replaced, mtime bumped) triggers re-verification. 7. Read ``metadata_len`` bytes from ``metadata_start`` for the manifest. -8. Read ``source_len`` bytes from ``source_start`` for the source view. - If ``source_len == 0``, no source is embedded; the UI displays - "(source not available)". +8. Read the source files through the manifest's ``sources`` list. For each entry, check that + ``offset`` and ``length`` are non-negative integers and that ``offset + length <= source_len``, + read ``length`` bytes from ``source_start + offset``, and compare their SHA-256 to ``sha256``. A + duplicate ``path`` or a digest mismatch is an error. Without a ``sources`` key, no source is + embedded; the UI displays "(source not available)". Review Comment: The notice the importer shows is `// Source code is not available: the bundle embeds no source.` (`_NO_SOURCE`), not "(source not available)". ########## airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst: ########## @@ -602,6 +631,20 @@ The matching bundle is marked executable before it is launched, so any Dag bundl object-store one such as ``S3DagBundle`` that has no concept of file permissions and so cannot preserve the execute bit the build produced. +.. _go-sdk/native-dag-parsing: + +Parsing native Dags +~~~~~~~~~~~~~~~~~~~ + +Once an :class:`~airflow.sdk.coordinators.executable.ExecutableCoordinator` is configured, the Dag processor Review Comment: java.rst and typescript.rst both cover running more than one coordinator of the class: map each Dag bundle with native Dags to one of them in `[sdk] dag_bundle_to_coordinator`. That applies here too (say a second `ExecutableCoordinator` for another `task_handler_bundle_name`), and without the mapping `get_dag_parsing_coordinator_key` raises `InvalidCoordinatorError` for every bundle binary the Dag processor picks up. Could this section get the same paragraph? -- 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]
