jason810496 commented on code in PR #74341: URL: https://github.com/apache/airflow/pull/74341#discussion_r4209235905
########## 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 Review Comment: Fixed in fc67f865848, with a smaller change than resolving the bundle first. `ExecutableDagImporter.can_handle` now checks only the name, the same as the Java and Node importers: a bundle binary has no file extension. The trailer check moved to `might_contain_dag`, which runs only when a bundle is scanned. So `execute_task` and `manager.py:813` work with the relative path unchanged. go.rst and the spec now state the no-extension rule. `TestExecuteTaskNativeBundle` covers `bin/orders` with a cwd that does not contain it. ########## 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: We'll leave this out. The Go SDK releases so far are betas, so we don't add version gates for them. Rebuilding the binary with the current Go SDK fixes the parse. ########## 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: Fixed in 985f2f5970e: `_open_checked_bundle` raises with the reason, and the parse import error includes it, for example a binary SHA-256 mismatch after `strip` or `codesign`. The bundle scan keeps the non-raising wrapper. ########## 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: Fixed in 51a902f850d: the manifest and source index are read and verified once per file and cached by file identity. A read hashes only the source it returns. `test_reads_the_bundle_index_once_for_many_dags` covers it. ########## 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: Fixed in e215043bfef: the spec now describes the notice instead of quoting it, so the two cannot drift. ########## 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: Added in b8f59bb34cb, matching java.rst and typescript.rst. ########## 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: Done in fcc01093c08. ########## 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: Fixed in 81613013b2d. fc67f865848 then moved the trailer check into `might_contain_dag`, so a file that cannot be read is now dropped there. -- 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]
