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]

Reply via email to