This is an automated email from the ASF dual-hosted git repository.
jason810496 pushed a commit to branch main
in repository https://gitbox.apache.org/repos/asf/airflow.git
The following commit(s) were added to refs/heads/main by this push:
new 46d3c59dfec Parse native Go Dags and show each Dag's own source
(#74341)
46d3c59dfec is described below
commit 46d3c59dfeccf7b57626c6f88678fa9694d99599
Author: Jason(Zhe-You) Liu <[email protected]>
AuthorDate: Thu Oct 8 12:02:58 2026 +0800
Parse native Go Dags and show each Dag's own source (#74341)
* Read per-Dag sources from executable bundles
The source region now holds the embedded files back to back, and the
metadata indexes them through `sources`, `dag_source_paths` and
`entrypoint_path`, replacing `source`. Add readers that verify each
region and resolve a Dag to its own file, falling back to the
entrypoint. Open and verify a bundle once so the metadata and source
reads share the same file handle.
* Add ExecutableDagImporter and the native Dag parse command
The importer claims any file that ends with the AFBNDL01 magic, so a
bundle needs no extension, and shows each Dag's embedded source file in
the Code view. ExecutableCoordinator now builds the parse command for
such a file from its verified metadata and supervisor schema version.
* Document native Dag parsing on the Go SDK page
Explain that the Dag processor parses every bundle binary once an
ExecutableCoordinator is configured, that the Code view shows each Dag's
own source file, and that handler-only binaries should be listed in
.airflowignore until the Go runtime answers the Dag-parse request.
* Say why an executable bundle is rejected
A bundle that fails verification now raises with the reason, such as a
binary SHA-256 mismatch or an unknown footer version.
* Read an executable bundle's manifest once per parse
The bundle index is cached per file identity and only the requested source
is read and hashed, so Dag count no longer multiplies the work.
* Fix the ExecutableDagImporter might_contain_dag docstring
A file that cannot be opened is not claimed by can_handle, so only a bundle
that fails verification is kept.
* Document dag_bundle_to_coordinator for native Go Dags
Mirror the Java and TypeScript pages: with several ExecutableCoordinators,
each Dag bundle needs a mapping.
* Describe the missing-source notice in the bundle spec
The spec quoted a notice text that the UI does not show, so it now
describes the behaviour.
* Autospec the _ensure_executable patch
The patch now checks the call signature of the function it replaces.
* Recognize an executable bundle by its name when a task runs
A task runs a native Dag file by its bundle-relative path, so the importer
claims a file by its extensionless name and checks the bundle trailer only
while listing a bundle.
---
.../authoring-and-scheduling/language-sdks/go.rst | 59 ++++-
task-sdk/docs/airflow-metadata.schema.json | 49 +++-
task-sdk/docs/executable-bundle-spec.rst | 86 +++++--
.../airflow/sdk/coordinators/_bundle_metadata.py | 16 ++
.../src/airflow/sdk/coordinators/_dag_importer.py | 1 +
.../sdk/coordinators/executable/_bundle_reader.py | 195 ++++++++++++++
.../sdk/coordinators/executable/_dag_importer.py | 105 ++++++++
.../sdk/coordinators/executable/coordinator.py | 108 ++++++--
.../coordinators/executable/_bundle_test_utils.py | 98 +++++++
.../coordinators/executable/test_bundle_reader.py | 217 ++++++++++++++++
.../coordinators/executable/test_coordinator.py | 168 +++++++++++-
.../coordinators/executable/test_dag_importer.py | 283 +++++++++++++++++++++
12 files changed, 1321 insertions(+), 64 deletions(-)
diff --git a/airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst
b/airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst
index f2d7613e356..fac17348007 100644
--- a/airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst
+++ b/airflow-core/docs/authoring-and-scheduling/language-sdks/go.rst
@@ -62,6 +62,8 @@ is the same coordinator mechanism the Java SDK uses. Because
the mature Python s
Airflow-facing concerns, Go tasks inherit remote task logs (S3/GCS), the full
range of task states, and
alternate XCom backends rather than implementing them again in Go.
+.. _go-sdk/quick-start:
+
Quick start
-----------
@@ -221,6 +223,33 @@ There is no separate Go worker to run: the Airflow worker
forks the bundle binar
and the Dag processor resolve ``task_handler_bundle_name`` through it, and
wherever the ``[sdk]`` config
is read it is rejected if the name is missing there.
+A Dag processor with this ``[sdk]`` configuration also parses the bundle
binaries of every Dag bundle (see
+:ref:`go-sdk/native-dag-parsing`). Keep it from parsing the binaries that only
register task handlers with
+that Dag bundle's ``.airflowignore``, whose pattern syntax follows ``[core]
dag_ignore_file_syntax``. Tasks do
+not read ``.airflowignore``, so the binaries still run.
+
+With ``task_handler_bundle_name`` set to its own Dag bundle, as above,
``go-task-handlers`` holds only handler
+binaries, so ignore everything in it:
+
+.. code-block:: bash
+
+ # [core] dag_ignore_file_syntax = glob (the default)
+ echo '*' > /opt/airflow/go-task-handlers/.airflowignore
+
+ # [core] dag_ignore_file_syntax = regexp; a bare * is not a valid pattern
and is dropped
+ echo '.' > /opt/airflow/go-task-handlers/.airflowignore
+
+With ``task_handler_bundle_name`` unset, the binaries sit in the same Dag
bundle as the Python Dag file. Go
+binaries have no extension, so put them in a folder such as ``bin/`` and
ignore that folder:
+
+.. code-block:: bash
+
+ # [core] dag_ignore_file_syntax = glob (the default)
+ echo 'bin/*' >> /opt/airflow/dags/.airflowignore
+
+ # [core] dag_ignore_file_syntax = regexp
+ echo '^bin/' >> /opt/airflow/dags/.airflowignore
+
Writing tasks
-------------
@@ -596,13 +625,39 @@ Deploying
Copy or mount the packed bundle into the Dag bundle named by the coordinator's
``task_handler_bundle_name``.
The :class:`~airflow.sdk.coordinators.executable.ExecutableCoordinator` scans
that Dag bundle recursively,
matches the incoming ``dag_id`` against each bundle's manifest, verifies the
bundle's integrity hash, and
-launches the matching bundle. Bundles are identified by the trailer magic, not
by filename (no extension on
-Linux/macOS, ``.exe`` on Windows), so the file name on the worker is
irrelevant.
+launches the matching bundle. This scan identifies bundles by the trailer
magic, not by file name. The Dag
+processor also requires a file name without an extension, see
:ref:`go-sdk/native-dag-parsing`.
The matching bundle is marked executable before it is launched, so any Dag
bundle works, including an
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
+parses every bundle binary in each Dag bundle and runs it to collect the Dags
it defines. The Dag processor
+recognizes a bundle binary only when its file name has no extension, such as
``bin/orders``, not
+``bin/orders.bin``. A binary that only registers task handlers is parsed too:
each parse runs it and finds
+no Dags, and a binary built with a Go SDK that cannot answer the Dag-parse
request records an import error.
+Keep those out of the Dag processor with ``.airflowignore``, as described in
:ref:`Quick start <go-sdk/quick-start>`.
+
+With one :class:`~airflow.sdk.coordinators.executable.ExecutableCoordinator`,
it parses the binaries of every Dag
+bundle. With several, for example a second one for another
``task_handler_bundle_name``, map each Dag bundle that
+holds native Go Dags to one of them in ``[sdk] dag_bundle_to_coordinator``. A
binary in a Dag bundle that has no
+entry, or whose entry names no
:class:`~airflow.sdk.coordinators.executable.ExecutableCoordinator`, fails to
parse
+with an import error:
+
+.. code-block:: ini
+
+ [sdk]
+ dag_bundle_to_coordinator = {"dags-folder": "go"}
+
+* The Code view shows each native Dag's own source file, taken from the
sources the bundle embeds.
+* A task of a native Go Dag runs the bundle binary its Dag was parsed from.
+
.. _go-sdk/coordinator-config:
:class:`~airflow.sdk.coordinators.executable.ExecutableCoordinator`
configuration
diff --git a/task-sdk/docs/airflow-metadata.schema.json
b/task-sdk/docs/airflow-metadata.schema.json
index d4a5cd240e7..c9c699ba029 100644
--- a/task-sdk/docs/airflow-metadata.schema.json
+++ b/task-sdk/docs/airflow-metadata.schema.json
@@ -4,7 +4,7 @@
"title": "Airflow Executable SDK Bundle Metadata",
"description": "Build-time manifest declaring DAG and task identifiers
exposed by an Airflow native-executable SDK bundle. See the Executable Bundle
Spec documentation in the Airflow Task SDK.",
"type": "object",
- "required": ["airflow_bundle_metadata_version", "sdk", "source", "dags"],
+ "required": ["airflow_bundle_metadata_version", "sdk", "dags"],
"additionalProperties": true,
"properties": {
"airflow_bundle_metadata_version": {
@@ -35,11 +35,26 @@
}
}
},
- "source": {
+ "entrypoint_path": {
"type": "string",
- "description": "Original filename of the primary Dag source file (e.g.
'example.go'). Bundle formats with an embedded source region use this as its
display name; other formats treat it as the logical authoring name.",
+ "description": "Path of the entrypoint source file (e.g.
'example/bundle/main.go'). It must be one of the sources paths. Shown for a Dag
that dag_source_paths does not map.",
"minLength": 1
},
+ "dag_source_paths": {
+ "type": "object",
+ "description": "Mapping of dag_id to the path of the file that defines
the Dag. Every value must be one of the sources paths.",
+ "additionalProperties": {
+ "type": "string",
+ "minLength": 1
+ }
+ },
+ "sources": {
+ "type": "array",
+ "description": "The embedded source files, in the order they appear in
the source region. A bundle without it embeds no source. Paths must be unique.",
+ "items": {
+ "$ref": "#/$defs/sourceEntry"
+ }
+ },
"dags": {
"type": "object",
"description": "Mapping of dag_id to DAG entry. Every dag_id the bundle
exposes must appear here.",
@@ -50,6 +65,34 @@
}
},
"$defs": {
+ "sourceEntry": {
+ "type": "object",
+ "description": "One embedded source file inside the bundle's source
region.",
+ "required": ["path", "offset", "length", "sha256"],
+ "additionalProperties": true,
+ "properties": {
+ "path": {
+ "type": "string",
+ "description": "The file's path as the author knows it (e.g.
'example/bundle/main.go').",
+ "minLength": 1
+ },
+ "offset": {
+ "type": "integer",
+ "description": "The file's start, in bytes from the start of the
source region.",
+ "minimum": 0
+ },
+ "length": {
+ "type": "integer",
+ "description": "The file's length in bytes.",
+ "minimum": 0
+ },
+ "sha256": {
+ "type": "string",
+ "description": "Lower-case hexadecimal SHA-256 of the file's bytes.",
+ "pattern": "^[0-9a-f]{64}$"
+ }
+ }
+ },
"dagEntry": {
"type": "object",
"description": "Static description of a single DAG declared in the
bundle.",
diff --git a/task-sdk/docs/executable-bundle-spec.rst
b/task-sdk/docs/executable-bundle-spec.rst
index afb0969a2cb..b8041b6357a 100644
--- a/task-sdk/docs/executable-bundle-spec.rst
+++ b/task-sdk/docs/executable-bundle-spec.rst
@@ -42,29 +42,37 @@ A bundle file therefore has three regions, in order from
offset 0:
1. The native executable (ELF / Mach-O / PE), including any code-signing
structures the platform appends.
-2. The primary DAG source file, embedded verbatim (UTF-8). MAY have length 0.
+2. The embedded source files, each verbatim (UTF-8), back to back. MAY have
length 0.
+ The metadata indexes them. See :ref:`the source region
<bundle-source-region>`.
3. The build-time manifest (``airflow-metadata.yaml`` content, UTF-8).
The file ends with a fixed 64-byte trailer that locates regions (2) and (3),
carries an integrity hash of the binary region, and identifies the file as a
bundle. See :ref:`bundle-trailer-layout`.
-Filenames follow OS conventions for executables: no extension on Linux/macOS,
-``.exe`` on Windows. The scanner identifies bundles by the trailer's magic,
-not by the filename.
+A bundle file has no file extension, and the Dag processor does not treat a
+file with an extension as a bundle. The scanner identifies bundles by the
+trailer's magic.
The complete bundle file regions are:
.. code-block:: text
[0, source_start) native binary (must be non-empty)
- [source_start, metadata_start) embedded source (may be zero length)
+ [source_start, metadata_start) embedded source files (may be zero length)
[metadata_start, file_size-64) build-time manifest
[file_size-64, file_size) 64-byte trailer
where ``metadata_start = file_size - 64 - metadata_len`` and
``source_start = metadata_start - source_len``.
+.. _bundle-source-region:
+
+The source region holds the source files of the Dags the bundle defines, one
entry per file. Files
+are concatenated with no separator, and the ``sources`` list in the manifest
gives each file's
+``path``, ``offset`` and ``length`` within the region and its ``sha256``. A
file's ``offset`` is
+relative to ``source_start``, so the first file has offset ``0``.
+
Reference Implementation
------------------------
@@ -84,7 +92,7 @@ for SDK users. Go SDK's ``airflow-go-pack`` is a good example.
BINARY = pathlib.Path(...) # Path to the compiled executable.
OUTPUT = pathlib.Path(...) # Where to put the processed executable.
- SOURCE = b"..." # Source code to embed.
+ SOURCES = [b"..."] # Source files to embed, in the order the manifest
lists them.
METADATA = b"..." # UTF-8-encoded YAML metadata.
# SHA-256 covers the binary region only: bytes [0, source_start).
@@ -92,7 +100,7 @@ for SDK users. Go SDK's ``airflow-go-pack`` is a good
example.
trailer = struct.pack(
"<III 32s 12s 8s",
- len(SOURCE), # source_len
+ sum(map(len, SOURCES)), # source_len
len(METADATA), # metadata_len
1, # footer_ver
binary_sha256,
@@ -103,7 +111,7 @@ for SDK users. Go SDK's ``airflow-go-pack`` is a good
example.
shutil.copy(BINARY, OUTPUT)
with OUTPUT.open("ab") as fh:
- fh.write(SOURCE) # Embedded source region.
+ fh.writelines(SOURCES) # Embedded source region.
fh.write(METADATA) # Metadata region.
fh.write(trailer)
OUTPUT.chmod(0o755)
@@ -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``.
+ To read a file, 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, and the UI shows a notice in place of the source.
+
+ A Dag's source file is the ``dag_source_paths`` entry for its ``dag_id``. A
Dag with no entry,
+ such as one built dynamically, shows ``entrypoint_path``. A Dag owned by
another language, such as
+ the Python Dag that a bundle's task handlers run for, shows its own source
and not an embedded
+ file. ``entrypoint_path`` and every ``dag_source_paths`` value MUST be one
of the ``sources``
+ paths.
Source comes *before* metadata so a future ``footer_ver`` MAY introduce
additional trailing blobs (e.g. signed checksums, compressed deps) by
@@ -182,7 +198,19 @@ and editors.
language: go
version: "0.1.0"
supervisor_schema_version: "2026-06-16"
- source: example.go
+ entrypoint_path: example/bundle/main.go
+ dag_source_paths:
+ example_dag: example/bundle/main.go
+ another_dag: example/bundle/dags/another.go
+ sources:
+ - path: example/bundle/main.go
+ offset: 0
+ length: 1532
+ sha256: 0f3a...e91c
+ - path: example/bundle/dags/another.go
+ offset: 1532
+ length: 811
+ sha256: 7b21...04d8
dags:
example_dag:
tasks:
@@ -214,12 +242,24 @@ Top-level keys:
task-execution time, and an unknown version causes that bundle to be
skipped.
-``source`` (string, required)
- Original filename of the primary DAG source file (e.g. ``example.go``).
- The file's bytes live in the source region of the bundle, not at this
- path; this field is a display name the Airflow UI uses to label the
- source-view panel and pick a syntax-highlighting mode from the
- extension.
+``entrypoint_path`` (string, optional)
+ Path of the entrypoint source file, such as the Go ``main`` package file.
It MUST be one of
+ the ``sources`` paths. The Airflow UI shows it for a Dag that
``dag_source_paths`` does not map.
+
+``dag_source_paths`` (mapping, optional)
+ Mapping of ``dag_id`` to the path of the file that defines the Dag. Every
value MUST be one of
+ the ``sources`` paths.
+
+``sources`` (list, optional)
+ The embedded source files, in the order they appear in the source region.
A bundle without it
+ embeds no source. Each entry has:
+
+ - ``path`` (string, required): the file's path as the author knows it,
such as
+ ``example/bundle/main.go``. Paths MUST be unique. The Airflow UI picks a
syntax-highlighting
+ mode from the extension.
+ - ``offset`` (integer, required): the file's start, in bytes from
``source_start``.
+ - ``length`` (integer, required): the file's length in bytes.
+ - ``sha256`` (string, required): the lower-case hexadecimal SHA-256 of the
file's bytes.
``dags`` (mapping, required)
Mapping of ``dag_id`` to a *DAG entry*. Every ``dag_id`` the bundle
@@ -243,16 +283,16 @@ Go bundle::
example
├── ELF/Mach-O/PE executable
- ├── source region: contents of example.go
- ├── metadata region: airflow-metadata.yaml (source: example.go)
+ ├── source region: example/bundle/main.go, example/bundle/dags/another.go
+ ├── metadata region: airflow-metadata.yaml (entrypoint_path,
dag_source_paths, sources)
└── trailer (64 B): lengths + binary_sha256 + AFBNDL01 magic
Rust bundle::
pipeline
├── ELF/Mach-O/PE executable
- ├── source region: contents of main.rs
- ├── metadata region: airflow-metadata.yaml (source: main.rs)
+ ├── source region: src/main.rs
+ ├── metadata region: airflow-metadata.yaml (entrypoint_path: src/main.rs)
└── trailer (64 B): lengths + binary_sha256 + AFBNDL01 magic
The bundle is one file. ``./example`` runs the binary; the appended data
@@ -272,7 +312,7 @@ that perform additional post-build steps MUST observe the
following order:
binary region; nothing has been written past its OS-defined end yet, so
the digest matches what the reader will recompute over
``[0, source_start)`` after the append.
-- **Append** ``<source><metadata><trailer>`` in a single write so a
+- **Append** ``<sources><metadata><trailer>`` in a single write so a
partially written file fails the magic or hash check rather than
appearing as a half-valid bundle.
diff --git a/task-sdk/src/airflow/sdk/coordinators/_bundle_metadata.py
b/task-sdk/src/airflow/sdk/coordinators/_bundle_metadata.py
index b97ed80e8d8..076cc90b107 100644
--- a/task-sdk/src/airflow/sdk/coordinators/_bundle_metadata.py
+++ b/task-sdk/src/airflow/sdk/coordinators/_bundle_metadata.py
@@ -124,3 +124,19 @@ def extract_supervisor_schema_version(metadata: dict[str,
Any]) -> str:
if not isinstance(value, str) or not value:
raise ValueError("missing or invalid sdk.supervisor_schema_version")
return value
+
+
+def resolve_source_path(metadata: dict[str, Any], dag_id: str | None) -> str |
None:
+ """
+ Return the embedded source path to show for *dag_id*, or ``None`` when
there is none.
+
+ A Dag mapped in ``dag_source_paths`` resolves to its own file. Any other
Dag, such as one built
+ dynamically, and a *dag_id* of ``None`` resolve to ``entrypoint_path``.
+ """
+ dag_source_paths = metadata.get("dag_source_paths")
+ if dag_id is not None and isinstance(dag_source_paths, dict):
+ mapped = dag_source_paths.get(dag_id)
+ if isinstance(mapped, str):
+ return mapped
+ entrypoint_path = metadata.get("entrypoint_path")
+ return entrypoint_path if isinstance(entrypoint_path, str) else None
diff --git a/task-sdk/src/airflow/sdk/coordinators/_dag_importer.py
b/task-sdk/src/airflow/sdk/coordinators/_dag_importer.py
index 4a72755d11c..4f4746211bf 100644
--- a/task-sdk/src/airflow/sdk/coordinators/_dag_importer.py
+++ b/task-sdk/src/airflow/sdk/coordinators/_dag_importer.py
@@ -44,6 +44,7 @@ if TYPE_CHECKING:
COORDINATOR_DAG_IMPORTERS: Final[tuple[str, ...]] = (
"airflow.sdk.coordinators.java._dag_importer.JavaDagImporter",
"airflow.sdk.coordinators.node._dag_importer.NodeDagImporter",
+ "airflow.sdk.coordinators.executable._dag_importer.ExecutableDagImporter",
)
"""
The classpaths of the :class:`CoordinatorDagImporter` subclasses a Dag
bundle's registry may hold.
diff --git a/task-sdk/src/airflow/sdk/coordinators/executable/_bundle_reader.py
b/task-sdk/src/airflow/sdk/coordinators/executable/_bundle_reader.py
new file mode 100644
index 00000000000..1443ca5e91f
--- /dev/null
+++ b/task-sdk/src/airflow/sdk/coordinators/executable/_bundle_reader.py
@@ -0,0 +1,195 @@
+#
+# 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.
+"""Read the per-Dag source files an executable bundle embeds."""
+
+from __future__ import annotations
+
+import functools
+import hashlib
+import pathlib
+import re
+from typing import Any
+
+import attrs
+
+from airflow.sdk.coordinators._bundle_metadata import resolve_source_path
+from airflow.sdk.coordinators.executable.coordinator import (
+ _VERIFY_CACHE_MAXSIZE,
+ _open_checked_bundle,
+ _VerifiedBundle,
+)
+
+_SHA256_HEX = re.compile(r"[0-9a-f]{64}")
+
+
[email protected](frozen=True)
+class _SourceRegion:
+ """One embedded file's path and its byte range inside the source region."""
+
+ path: str
+ offset: int
+ length: int
+ sha256: str
+
+
[email protected](frozen=True)
+class _BundleIndex:
+ """
+ A verified bundle's manifest and source index, shared by every read of the
same file.
+
+ The mappings are shared between callers, so they MUST NOT be mutated.
+ """
+
+ metadata: dict[str, Any]
+ source_start: int
+ regions: dict[str, _SourceRegion] | None
+
+
+def read_bundle_source(bundle_path: pathlib.Path, dag_id: str | None = None)
-> str | None:
+ """
+ Return the source file the bundle embeds for *dag_id*, or ``None``.
+
+ Pass *dag_id* only for a Dag the bundle defines natively. Omit it for a
Dag owned by another
+ language (Python), which shows its own source; the result is then ``None``
rather than an
+ unrelated file. A Dag mapped in ``dag_source_paths`` returns its own file,
and one that is not
+ (such as a Dag built dynamically) falls back to the entrypoint source.
+
+ :raises ValueError: when the file is not a valid bundle or its embedded
sources are invalid.
+ """
+ if dag_id is None:
+ return None
+ return _read_source(bundle_path, dag_id)
+
+
+def read_bundle_entrypoint_source(bundle_path: pathlib.Path) -> str | None:
+ """
+ Return the entrypoint source the bundle embeds, or ``None`` when it embeds
none.
+
+ :raises ValueError: when the file is not a valid bundle or its embedded
sources are invalid.
+ """
+ return _read_source(bundle_path, None)
+
+
+def read_bundle_language(bundle_path: pathlib.Path) -> str | None:
+ """Return the ``sdk.language`` of the bundle's metadata, or ``None`` when
it has none."""
+ try:
+ sdk = _read_index(bundle_path).metadata.get("sdk")
+ except ValueError:
+ return None
+ language = sdk.get("language") if isinstance(sdk, dict) else None
+ return language if isinstance(language, str) and language else None
+
+
+def _read_index(bundle_path: pathlib.Path) -> _BundleIndex:
+ try:
+ st = bundle_path.stat()
+ except OSError as exc:
+ raise ValueError(f"Cannot stat bundle file {bundle_path}: {exc}") from
exc
+ return _load_index(str(bundle_path), st.st_ino, st.st_mtime_ns, st.st_size)
+
+
[email protected]_cache(maxsize=_VERIFY_CACHE_MAXSIZE)
+def _load_index(path: str, ino: int, mtime_ns: int, size: int) -> _BundleIndex:
+ # The file identity is part of the key, so a replaced bundle misses the
cache and is read again.
+ with _open_checked_bundle(pathlib.Path(path)) as bundle:
+ return _BundleIndex(bundle.metadata, bundle.footer.source_start,
_parse_source_regions(bundle))
+
+
+def _read_source(bundle_path: pathlib.Path, dag_id: str | None) -> str | None:
+ index = _read_index(bundle_path)
+ if index.regions is None:
+ return None
+ source_path = resolve_source_path(index.metadata, dag_id)
+ if source_path is None:
+ return None
+ return _decode_source(
+ _read_region(bundle_path, index.source_start,
index.regions[source_path]), source_path
+ )
+
+
+def _parse_source_regions(bundle: _VerifiedBundle) -> dict[str, _SourceRegion]
| None:
+ """Validate the ``sources`` index and the paths that refer to it, or
return ``None`` without one."""
+ metadata = bundle.metadata
+ sources = metadata.get("sources")
+ if sources is None:
+ return None
+ if not isinstance(sources, list):
+ raise ValueError("bundle metadata sources must be a list")
+
+ regions: dict[str, _SourceRegion] = {}
+ for entry in sources:
+ region = _parse_source_region(entry, bundle.footer.source_len)
+ if region.path in regions:
+ raise ValueError(f"bundle metadata declares duplicate source path
{region.path!r}")
+ regions[region.path] = region
+
+ entrypoint_path = metadata.get("entrypoint_path")
+ if entrypoint_path is not None and entrypoint_path not in regions:
+ raise ValueError(f"bundle entrypoint_path {entrypoint_path!r} is not
one of its sources")
+
+ dag_source_paths = metadata.get("dag_source_paths")
+ if dag_source_paths is not None:
+ if not isinstance(dag_source_paths, dict):
+ raise ValueError("bundle metadata dag_source_paths must be a
mapping")
+ for dag_id, source_path in dag_source_paths.items():
+ if source_path not in regions:
+ raise ValueError(f"bundle dag_source_paths maps {dag_id!r} to
{source_path!r}, not a source")
+ return regions
+
+
+def _parse_source_region(entry: Any, source_len: int) -> _SourceRegion:
+ if not isinstance(entry, dict):
+ raise ValueError("bundle metadata sources entries must be mappings")
+ path = entry.get("path")
+ if not isinstance(path, str) or not path:
+ raise ValueError("bundle metadata source path must be a non-empty
string")
+ offset = _non_negative_int(entry.get("offset"), path, "offset")
+ length = _non_negative_int(entry.get("length"), path, "length")
+ sha256 = entry.get("sha256")
+ if not isinstance(sha256, str) or _SHA256_HEX.fullmatch(sha256) is None:
+ raise ValueError(f"bundle source {path!r} sha256 must be 64 lowercase
hexadecimal digits")
+ if offset + length > source_len:
+ raise ValueError(f"bundle source {path!r} extends past the source
region")
+ return _SourceRegion(path=path, offset=offset, length=length,
sha256=sha256)
+
+
+def _non_negative_int(value: Any, source_path: str, name: str) -> int:
+ if not isinstance(value, int) or isinstance(value, bool) or value < 0:
+ raise ValueError(f"bundle source {source_path!r} {name} must be a
non-negative integer")
+ return value
+
+
+def _read_region(bundle_path: pathlib.Path, source_start: int, region:
_SourceRegion) -> bytes:
+ try:
+ with open(bundle_path, "rb") as f:
+ f.seek(source_start + region.offset)
+ payload = f.read(region.length)
+ except OSError as exc:
+ raise ValueError(f"Cannot read source {region.path!r} of
{bundle_path}: {exc}") from exc
+ if len(payload) != region.length:
+ raise ValueError(f"{bundle_path.name} was truncated while reading
source {region.path!r}")
+ if hashlib.sha256(payload).hexdigest() != region.sha256:
+ raise ValueError(f"{bundle_path.name} source {region.path!r} SHA-256
mismatch")
+ return payload
+
+
+def _decode_source(payload: bytes, source_path: str) -> str:
+ try:
+ return payload.decode("utf-8")
+ except UnicodeDecodeError as exc:
+ raise ValueError(f"embedded airflow source {source_path!r} is not
valid UTF-8: {exc}") from exc
diff --git a/task-sdk/src/airflow/sdk/coordinators/executable/_dag_importer.py
b/task-sdk/src/airflow/sdk/coordinators/executable/_dag_importer.py
new file mode 100644
index 00000000000..4a6d2054557
--- /dev/null
+++ b/task-sdk/src/airflow/sdk/coordinators/executable/_dag_importer.py
@@ -0,0 +1,105 @@
+#
+# 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,
get_file_suffix
+
+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 binary has
+ no file extension, and the ``AFBNDL01`` magic that ends the file
identifies it when the bundle is scanned.
+ """
+
+ coordinator_classpath: ClassVar[str] =
"airflow.sdk.coordinators.executable.ExecutableCoordinator"
+ artifact_suffix = ""
+ supported_extensions: list[str] = []
+
+ def can_handle(self, definition: DagDefinition | str | Path) -> bool:
+ return get_file_suffix(definition) == ""
+
+ 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 whether the file ends with the bundle trailer.
+
+ ``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 bundle that
fails verification is kept,
+ so that parsing it reports why.
+ """
+ try:
+ with definition.as_file() as path, open(path, "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 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:
+ source = (
+ read_bundle_source(path, dag_id)
+ if dag_id is not None
+ else read_bundle_entrypoint_source(path)
+ )
+ language = read_bundle_language(path)
+ return DagSourceCode(source_code=source or _NO_SOURCE,
language=language or _DEFAULT_LANGUAGE)
diff --git a/task-sdk/src/airflow/sdk/coordinators/executable/coordinator.py
b/task-sdk/src/airflow/sdk/coordinators/executable/coordinator.py
index 8e52490385c..9b1b7519686 100644
--- a/task-sdk/src/airflow/sdk/coordinators/executable/coordinator.py
+++ b/task-sdk/src/airflow/sdk/coordinators/executable/coordinator.py
@@ -19,6 +19,7 @@
from __future__ import annotations
+import contextlib
import hashlib
import os
import pathlib
@@ -185,7 +186,25 @@ class _BinaryDigestCache:
_digest_cache = _BinaryDigestCache(maxsize=_VERIFY_CACHE_MAXSIZE)
-def _read_bundle_metadata(path: pathlib.Path) -> dict[str, Any] | None:
[email protected]
+class _VerifiedBundle:
+ """A bundle whose trailer and binary digest check out, with its metadata
and the open file."""
+
+ file: BinaryIO
+ footer: _Footer
+ metadata: dict[str, Any]
+
+
[email protected]
+def _open_checked_bundle(path: pathlib.Path) -> Iterator[_VerifiedBundle]:
+ """
+ Open *path* and yield it as a verified bundle.
+
+ The file stays open for the ``with`` block, so a caller can read more
regions from the
+ bundle it verified.
+
+ :raises ValueError: with the reason *path* is not a usable bundle.
+ """
# One open per bundle: trailer-parse, hash (on cache miss), and
# metadata-read all share the same fd, and the stat that keys the
# digest cache comes from that fd too. This both halves the syscall
@@ -194,55 +213,65 @@ def _read_bundle_metadata(path: pathlib.Path) ->
dict[str, Any] | None:
try:
f = open(path, "rb")
except OSError as exc:
- log.debug("Cannot open bundle file; skipping", path=str(path),
error=str(exc))
- return None
+ raise ValueError(f"Cannot open bundle file {path}: {exc}") from exc
with f:
try:
st = os.fstat(f.fileno())
except OSError as exc:
- log.debug("Cannot stat bundle file; skipping", path=str(path),
error=str(exc))
- return None
+ raise ValueError(f"Cannot stat bundle file {path}: {exc}") from exc
try:
footer = _Footer.read(f, path, st.st_size)
- except (OSError, ValueError) as exc:
- log.debug("Invalid bundle trailer; skipping", path=str(path),
error=str(exc))
- return None
+ except OSError as exc:
+ raise ValueError(f"Cannot read bundle trailer of {path}: {exc}")
from exc
if footer is None:
- return None
+ raise ValueError(f"{path} has no bundle trailer")
cache_key: _DigestKey = (str(path), footer.source_start, st.st_ino,
st.st_mtime_ns, st.st_size)
actual_digest = _digest_cache.get(cache_key)
if actual_digest is None:
try:
actual_digest = _hash_open_file(f, footer.source_start, path)
- except (OSError, ValueError) as exc:
- log.debug("Failed to hash bundle binary region",
path=str(path), error=str(exc))
- return None
+ except OSError as exc:
+ raise ValueError(f"Cannot hash the binary region of {path}:
{exc}") from exc
_digest_cache.put(cache_key, actual_digest)
if actual_digest != footer.binary_sha256:
- log.debug(
- "Bundle binary_sha256 mismatch; skipping",
- path=str(path),
- expected=footer.binary_sha256.hex(),
- actual=actual_digest.hex(),
+ raise ValueError(
+ f"{path} binary SHA-256 does not match its trailer; "
+ "was it changed after packing, for example by strip or
codesign?"
)
- return None
try:
f.seek(footer.metadata_start)
metadata_bytes = f.read(footer.metadata_len)
except OSError as exc:
- log.debug("Cannot read bundle metadata; skipping", path=str(path),
error=str(exc))
- return None
+ raise ValueError(f"Cannot read the metadata of {path}: {exc}")
from exc
+
+ try:
+ metadata = parse_metadata_mapping(metadata_bytes, source="bundle
metadata")
+ except ValueError as exc:
+ raise ValueError(f"Cannot decode the metadata of {path}: {exc}")
from exc
+
+ yield _VerifiedBundle(file=f, footer=footer, metadata=metadata)
- try:
- return parse_metadata_mapping(metadata_bytes, source="bundle metadata")
- except ValueError as exc:
- log.debug("Cannot decode bundle metadata; skipping", path=str(path),
error=str(exc))
- return None
+
[email protected]
+def _open_verified_bundle(path: pathlib.Path) -> Iterator[_VerifiedBundle |
None]:
+ """Like :func:`_open_checked_bundle`, but yield ``None`` instead of
raising for an unusable bundle."""
+ with contextlib.ExitStack() as stack:
+ try:
+ bundle = stack.enter_context(_open_checked_bundle(path))
+ except ValueError as exc:
+ log.debug("Not a usable bundle; skipping", path=str(path),
error=str(exc))
+ bundle = None
+ yield bundle
+
+
+def _read_bundle_metadata(path: pathlib.Path) -> dict[str, Any] | None:
+ with _open_verified_bundle(path) as bundle:
+ return None if bundle is None else bundle.metadata
def _dag_ids(metadata: dict[str, Any]) -> set[str]:
@@ -352,7 +381,11 @@ class _Bundle(ResolvedBundle):
@attrs.define(kw_only=True)
class ExecutableCoordinator(SubprocessCoordinator):
"""
- Coordinator that launches a native executable subprocess for task
execution.
+ Coordinator that launches a native executable subprocess for task
execution and Dag parsing.
+
+ It runs the bundle that holds a Python Dag's task handlers, and it parses
the native Dags of
+ every bundle in a Dag bundle: the Dag processor runs each bundle binary to
collect its Dags.
+ A task of a native Dag runs the bundle its Dag was parsed from.
Configuration is taken from the ``[sdk] coordinators`` entry that
constructs
this instance::
@@ -377,3 +410,26 @@ class ExecutableCoordinator(SubprocessCoordinator):
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."""
+ try:
+ with _open_checked_bundle(path) as checked:
+ metadata = checked.metadata
+ except ValueError as exc:
+ raise ValueError(f"{path} is not a valid executable bundle:
{exc}") from exc
+ 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)
diff --git
a/task-sdk/tests/task_sdk/coordinators/executable/_bundle_test_utils.py
b/task-sdk/tests/task_sdk/coordinators/executable/_bundle_test_utils.py
new file mode 100644
index 00000000000..6ec93466ead
--- /dev/null
+++ b/task-sdk/tests/task_sdk/coordinators/executable/_bundle_test_utils.py
@@ -0,0 +1,98 @@
+#
+# 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.
+from __future__ import annotations
+
+import hashlib
+import stat
+import struct
+from typing import TYPE_CHECKING, Any
+
+import yaml
+
+if TYPE_CHECKING:
+ from pathlib import Path
+
+SCHEMA_VERSION = "2026-06-16"
+FOOTER_MAGIC = b"AFBNDL01"
+BINARY_PAYLOAD = b"\x7fELF" + b"binary-stub-payload"
+ENTRYPOINT_PATH = "example/bundle/main.go"
+DEFAULT_SOURCES = {ENTRYPOINT_PATH: b"package main\n\nfunc main() {}\n"}
+
+
+def write_bundle(
+ path: Path,
+ *dag_ids: str,
+ sources: dict[str, bytes] | None = None,
+ dag_source_paths: dict[str, str] | None = None,
+ entrypoint_path: str | None = ENTRYPOINT_PATH,
+ language: str | None = "go",
+ schema_version: str | None = SCHEMA_VERSION,
+ index: list[dict[str, Any]] | None = None,
+ omit_sources: bool = False,
+ source_region: bytes | None = None,
+ binary: bytes = BINARY_PAYLOAD,
+) -> Path:
+ """
+ Write a synthetic executable bundle to *path* and return it.
+
+ *sources* are embedded back to back, indexed from their real offsets and
digests. Pass *index* to
+ replace the ``sources`` entries and *source_region* to replace the
embedded bytes, to build an
+ invalid bundle. Without *sources*, the default entrypoint is embedded.
*omit_sources* leaves the
+ ``sources`` key out.
+ """
+ embedded = DEFAULT_SOURCES if sources is None else sources
+ region = b""
+ entries = []
+ for source_path, content in embedded.items():
+ entries.append(
+ {
+ "path": source_path,
+ "offset": len(region),
+ "length": len(content),
+ "sha256": hashlib.sha256(content).hexdigest(),
+ }
+ )
+ region += content
+ if source_region is not None:
+ region = source_region
+
+ sdk: dict[str, str] = {"version": "0.1.0"}
+ if language is not None:
+ sdk["language"] = language
+ if schema_version is not None:
+ sdk["supervisor_schema_version"] = schema_version
+ metadata: dict[str, Any] = {"airflow_bundle_metadata_version": "1.0",
"sdk": sdk}
+ if entrypoint_path is not None:
+ metadata["entrypoint_path"] = entrypoint_path
+ metadata["dag_source_paths"] = (
+ {dag_id: ENTRYPOINT_PATH for dag_id in dag_ids} if dag_source_paths is
None else dag_source_paths
+ )
+ if not omit_sources:
+ metadata["sources"] = entries if index is None else index
+ metadata["dags"] = {dag_id: {"tasks": ["task1"]} for dag_id in dag_ids}
+ metadata_bytes = yaml.safe_dump(metadata, sort_keys=True).encode("utf-8")
+
+ trailer = (
+ struct.pack("<III", len(region), len(metadata_bytes), 1)
+ + hashlib.sha256(binary).digest()
+ + bytes(12)
+ + FOOTER_MAGIC
+ )
+ path.write_bytes(binary + region + metadata_bytes + trailer)
+ path.chmod(path.stat().st_mode | stat.S_IEXEC | stat.S_IXGRP |
stat.S_IXOTH)
+ return path
diff --git
a/task-sdk/tests/task_sdk/coordinators/executable/test_bundle_reader.py
b/task-sdk/tests/task_sdk/coordinators/executable/test_bundle_reader.py
new file mode 100644
index 00000000000..bb6a7f97d30
--- /dev/null
+++ b/task-sdk/tests/task_sdk/coordinators/executable/test_bundle_reader.py
@@ -0,0 +1,217 @@
+#
+# 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.
+from __future__ import annotations
+
+import hashlib
+
+import pytest
+from task_sdk.coordinators.executable._bundle_test_utils import
ENTRYPOINT_PATH, write_bundle
+
+from airflow.sdk.coordinators.executable._bundle_reader import (
+ _load_index,
+ read_bundle_entrypoint_source,
+ read_bundle_language,
+ read_bundle_source,
+)
+from airflow.sdk.coordinators.executable.coordinator import _digest_cache
+
+from tests_common.test_utils.version_compat import AIRFLOW_V_3_3_PLUS
+
+if not AIRFLOW_V_3_3_PLUS:
+ pytest.skip("Coordinator is only compatible with Airflow >= 3.3.0",
allow_module_level=True)
+
+ORDERS = b'package main\n\nfunc orders() { dag.New("orders") }\n'
+REPORTS = b'package dags\n\nfunc reports() { dag.New("reports") }\n'
+MAIN = b"package main\n\nfunc main() {}\n"
+REPORTS_PATH = "example/bundle/dags/reports.go"
+TWO_FILES = {ENTRYPOINT_PATH: ORDERS, REPORTS_PATH: REPORTS}
+
+
[email protected](autouse=True)
+def _clear_caches():
+ _digest_cache.clear()
+ _load_index.cache_clear()
+
+
+def _entry(source_path: str, content: bytes, offset: int = 0, **overrides) ->
dict:
+ return {
+ "path": source_path,
+ "offset": offset,
+ "length": len(content),
+ "sha256": hashlib.sha256(content).hexdigest(),
+ **overrides,
+ }
+
+
+class TestReadBundleSource:
+ def test_returns_the_file_of_a_mapped_dag(self, tmp_path):
+ bundle = write_bundle(
+ tmp_path / "b",
+ "orders",
+ "reports",
+ sources=TWO_FILES,
+ dag_source_paths={"orders": ENTRYPOINT_PATH, "reports":
REPORTS_PATH},
+ )
+
+ assert read_bundle_source(bundle, "orders") == ORDERS.decode()
+ assert read_bundle_source(bundle, "reports") == REPORTS.decode()
+
+ def test_falls_back_to_the_entrypoint_for_an_unmapped_dag(self, tmp_path):
+ bundle = write_bundle(tmp_path / "b", "orders", sources=TWO_FILES,
dag_source_paths={})
+
+ assert read_bundle_source(bundle, "dynamic") == ORDERS.decode()
+
+ def test_returns_none_without_a_dag_id(self, tmp_path):
+ bundle = write_bundle(tmp_path / "b", "orders")
+
+ assert read_bundle_source(bundle) is None
+
+ def test_returns_none_without_sources(self, tmp_path):
+ bundle = write_bundle(tmp_path / "b", "orders", omit_sources=True)
+
+ assert read_bundle_source(bundle, "orders") is None
+ assert read_bundle_entrypoint_source(bundle) is None
+
+ def test_returns_none_without_an_entrypoint(self, tmp_path):
+ bundle = write_bundle(tmp_path / "b", "orders", sources={},
entrypoint_path=None, dag_source_paths={})
+
+ assert read_bundle_source(bundle, "orders") is None
+
+ def test_decodes_utf8(self, tmp_path):
+ bundle = write_bundle(tmp_path / "b", "orders",
sources={ENTRYPOINT_PATH: "// café\n".encode()})
+
+ assert read_bundle_entrypoint_source(bundle) == "// café\n"
+
+ def test_rereads_a_replaced_bundle(self, tmp_path):
+ path = tmp_path / "b"
+ write_bundle(path, "orders", sources={ENTRYPOINT_PATH: ORDERS})
+ assert read_bundle_entrypoint_source(path) == ORDERS.decode()
+
+ write_bundle(path, "orders", sources={ENTRYPOINT_PATH: REPORTS + b"//
more\n"})
+
+ assert read_bundle_entrypoint_source(path) == (REPORTS + b"//
more\n").decode()
+
+ def test_only_the_requested_source_is_hashed(self, tmp_path):
+ bundle = write_bundle(
+ tmp_path / "b",
+ "orders",
+ sources=TWO_FILES,
+ dag_source_paths={"orders": ENTRYPOINT_PATH, "reports":
REPORTS_PATH},
+ source_region=ORDERS + bytes(len(REPORTS)),
+ )
+
+ assert read_bundle_entrypoint_source(bundle) == ORDERS.decode()
+ with pytest.raises(ValueError, match="SHA-256 mismatch"):
+ read_bundle_source(bundle, "reports")
+
+
+class TestReadBundleEntrypointSource:
+ def test_returns_the_entrypoint_not_a_dag_file(self, tmp_path):
+ bundle = write_bundle(
+ tmp_path / "b",
+ "reports",
+ sources={ENTRYPOINT_PATH: MAIN, REPORTS_PATH: REPORTS},
+ dag_source_paths={"reports": REPORTS_PATH},
+ )
+
+ assert read_bundle_entrypoint_source(bundle) == MAIN.decode()
+
+
+class TestReadBundleLanguage:
+ def test_returns_the_sdk_language(self, tmp_path):
+ assert read_bundle_language(write_bundle(tmp_path / "b", "orders",
language="rust")) == "rust"
+
+ def test_returns_none_without_one(self, tmp_path):
+ assert read_bundle_language(write_bundle(tmp_path / "b", "orders",
language=None)) is None
+
+ def test_returns_none_for_a_file_that_is_not_a_bundle(self, tmp_path):
+ path = tmp_path / "plain"
+ path.write_bytes(b"not a bundle")
+
+ assert read_bundle_language(path) is None
+
+
+class TestInvalidSources:
+ def test_rejects_a_file_that_is_not_a_bundle(self, tmp_path):
+ path = tmp_path / "plain"
+ path.write_bytes(b"not a bundle")
+
+ with pytest.raises(ValueError, match="has no bundle trailer"):
+ read_bundle_entrypoint_source(path)
+
+ def test_rejects_a_missing_file(self, tmp_path):
+ with pytest.raises(ValueError, match="Cannot stat bundle file"):
+ read_bundle_entrypoint_source(tmp_path / "gone")
+
+ def test_rejects_duplicate_paths(self, tmp_path):
+ index = [_entry(ENTRYPOINT_PATH, MAIN), _entry(ENTRYPOINT_PATH, MAIN)]
+ bundle = write_bundle(tmp_path / "b", "orders", index=index)
+
+ with pytest.raises(ValueError, match="duplicate source path"):
+ read_bundle_entrypoint_source(bundle)
+
+ @pytest.mark.parametrize(
+ ("overrides", "message"),
+ [
+ ({"offset": 1}, "extends past the source region"),
+ ({"length": len(MAIN) + 1}, "extends past the source region"),
+ ({"offset": -1}, "offset must be a non-negative integer"),
+ ({"length": -1}, "length must be a non-negative integer"),
+ ({"offset": "0"}, "offset must be a non-negative integer"),
+ ({"length": True}, "length must be a non-negative integer"),
+ ({"sha256": hashlib.sha256(MAIN).hexdigest().upper()}, "64
lowercase hexadecimal digits"),
+ ({"path": ""}, "path must be a non-empty string"),
+ ],
+ )
+ def test_rejects_a_malformed_region(self, tmp_path, overrides, message):
+ bundle = write_bundle(tmp_path / "b", "orders",
index=[_entry(ENTRYPOINT_PATH, MAIN, **overrides)])
+
+ with pytest.raises(ValueError, match=message):
+ read_bundle_entrypoint_source(bundle)
+
+ def test_rejects_a_digest_mismatch(self, tmp_path):
+ bundle = write_bundle(
+ tmp_path / "b", "orders", index=[_entry(ENTRYPOINT_PATH, MAIN,
sha256="0" * 64)]
+ )
+
+ with pytest.raises(ValueError, match="SHA-256 mismatch"):
+ read_bundle_entrypoint_source(bundle)
+
+ def test_rejects_an_entrypoint_missing_from_sources(self, tmp_path):
+ bundle = write_bundle(tmp_path / "b", "orders",
entrypoint_path="other/main.go")
+
+ with pytest.raises(ValueError, match="entrypoint_path 'other/main.go'
is not one of its sources"):
+ read_bundle_entrypoint_source(bundle)
+
+ def test_rejects_a_dag_mapped_to_a_file_missing_from_sources(self,
tmp_path):
+ bundle = write_bundle(tmp_path / "b", "orders",
dag_source_paths={"orders": "missing.go"})
+
+ with pytest.raises(ValueError, match="maps 'orders' to 'missing.go',
not a source"):
+ read_bundle_source(bundle, "orders")
+
+ def test_rejects_invalid_utf8(self, tmp_path):
+ bundle = write_bundle(tmp_path / "b", "orders",
sources={ENTRYPOINT_PATH: b"\xff\xfe"})
+
+ with pytest.raises(ValueError, match="not valid UTF-8"):
+ read_bundle_entrypoint_source(bundle)
+
+ def test_rejects_a_non_list_sources(self, tmp_path):
+ bundle = write_bundle(tmp_path / "b", "orders", index={"a": 1}) #
type: ignore[arg-type]
+
+ with pytest.raises(ValueError, match="sources must be a list"):
+ read_bundle_entrypoint_source(bundle)
diff --git
a/task-sdk/tests/task_sdk/coordinators/executable/test_coordinator.py
b/task-sdk/tests/task_sdk/coordinators/executable/test_coordinator.py
index 69483d39796..37ac3fe9f73 100644
--- a/task-sdk/tests/task_sdk/coordinators/executable/test_coordinator.py
+++ b/task-sdk/tests/task_sdk/coordinators/executable/test_coordinator.py
@@ -32,7 +32,9 @@ import pytest
import yaml
from uuid6 import uuid7
-from airflow.sdk.api.datamodels._generated import TaskInstance
+from airflow.dag_processing.bundles.base import BaseDagBundle
+from airflow.sdk.api.datamodels._generated import BundleInfo, TaskInstance
+from airflow.sdk.coordinators._subprocess import _PopenActivitySubprocess
from airflow.sdk.coordinators.executable.coordinator import (
FOOTER_MAGIC,
FOOTER_SIZE,
@@ -44,6 +46,7 @@ from airflow.sdk.coordinators.executable.coordinator import (
)
from airflow.sdk.execution_time.coordinator import BaseCoordinator
from airflow.sdk.execution_time.supervisor import ActivitySubprocess
+from airflow.sdk.importers import reset_importer_registry
from tests_common.test_utils.config import conf_vars
from tests_common.test_utils.version_compat import AIRFLOW_V_3_3_PLUS
@@ -62,7 +65,7 @@ def _make_metadata(dag_ids, source_filename: str =
"example.go") -> dict:
"version": "0.1.0",
"supervisor_schema_version": "2026-06-16",
},
- "source": source_filename,
+ "entrypoint_path": source_filename,
"dags": {dag_id: {"tasks": ["task1"]} for dag_id in dag_ids},
}
@@ -338,12 +341,7 @@ class TestBundleFind:
with pytest.raises(FileNotFoundError, match="cannot find
executable bundle"):
_Bundle.find([tmp_path], "tutorial_dag")
- mock_log.debug.assert_any_call(
- "Bundle binary_sha256 mismatch; skipping",
- path=str(bundle_path),
- expected=mock.ANY,
- actual=mock.ANY,
- )
+ mock_log.debug.assert_any_call("Not a usable bundle; skipping",
path=str(bundle_path), error=mock.ANY)
def test_captures_schema_version_from_metadata(self, tmp_path):
_build_bundle(tmp_path / "with_schema", dag_ids=["tutorial_dag"])
@@ -412,7 +410,7 @@ class TestBundleFind:
_Bundle.find([tmp_path], "tutorial_dag")
mock_log.debug.assert_any_call(
- "Cannot decode bundle metadata; skipping",
+ "Not a usable bundle; skipping",
path=str(bundle_path),
error=mock.ANY,
)
@@ -429,7 +427,7 @@ class TestBundleFind:
_Bundle.find([tmp_path], "tutorial_dag")
mock_log.debug.assert_any_call(
- "Cannot decode bundle metadata; skipping",
+ "Not a usable bundle; skipping",
path=str(bundle_path),
error=mock.ANY,
)
@@ -447,6 +445,110 @@ class TestExecutableCoordinatorAttributes:
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="SHA-256 does not match"):
+ ExecutableCoordinator()._build_parse_dag_command(path=binary)
+
+ def test_raises_with_the_reason_for_an_unknown_footer_version(self,
tmp_path):
+ binary = _build_bundle(tmp_path / "future", footer_ver=2)
+
+ with pytest.raises(ValueError, match="is not a valid executable
bundle: .*footer_ver=2"):
+ 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(
+
"airflow.sdk.coordinators.executable.coordinator._ensure_executable",
+ autospec=True,
+ return_value="denied",
+ ),
+ pytest.raises(ValueError, match="Cannot run bundle .*: denied"),
+ ):
+ ExecutableCoordinator()._build_parse_dag_command(path=binary)
+
+
+class TestBuildDagFileCommand:
+ def test_runs_the_bundle_the_dag_was_parsed_from(self, tmp_path):
+ binary = _build_bundle(tmp_path / "my_bundle", dag_ids=["native_dag"])
+
+ command, schema_version =
ExecutableCoordinator()._build_dag_file_command(
+ what=_make_ti(dag_id="native_dag"), 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", dag_ids=["native_dag"])
+ binary.chmod(0o644)
+
+
ExecutableCoordinator()._build_dag_file_command(what=_make_ti(dag_id="native_dag"),
path=binary)
+
+ assert os.access(binary, os.X_OK)
+
+ 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_dag_file_command(what=_make_ti(dag_id="native_dag"),
path=plain)
+
+
class TestBuildExecuteTaskCommand:
def test_returns_resolved_executable_and_schema_version(self, tmp_path):
binary = _build_bundle(tmp_path / "my_bundle",
dag_ids=["tutorial_dag"])
@@ -590,3 +692,49 @@ class TestExecutableCoordinatorExecuteTask:
assert isinstance(result, BaseCoordinator.ExecutionResult)
assert result.exit_code == 0
+
+
+class TestExecuteTaskNativeBundle:
+ def test_runs_a_bundle_found_by_its_bundle_relative_name(self, tmp_path,
monkeypatch, mock_client):
+ dags = tmp_path / "dags"
+ (dags / "bin").mkdir(parents=True)
+ binary = _build_bundle(dags / "bin" / "orders", dag_ids=["orders"])
+ elsewhere = tmp_path / "elsewhere"
+ elsewhere.mkdir()
+ monkeypatch.chdir(elsewhere)
+ bundle = MagicMock(spec=BaseDagBundle, path=dags, version="v1")
+ bundle.name = "dags"
+ coordinators = {
+ ("sdk", "coordinators"): json.dumps(
+ {"go": {"classpath":
"airflow.sdk.coordinators.executable.ExecutableCoordinator"}}
+ )
+ }
+
+ reset_importer_registry()
+ try:
+ with (
+ conf_vars(coordinators),
+
patch("airflow.sdk.coordinators._subprocess.initialize_ti_bundle",
autospec=True) as init,
+
patch("airflow.sdk.coordinators._subprocess.BundleVersionLock", autospec=True),
+ patch.object(_PopenActivitySubprocess, "start", autospec=True)
as mock_start,
+ patch.object(
+ ExecutableCoordinator,
+ "_build_execute_task_command",
+ autospec=True,
+
side_effect=ExecutableCoordinator._build_execute_task_command,
+ ) as mock_scan,
+ ):
+ init.return_value = bundle
+ mock_start.return_value.wait.return_value = 0
+ ExecutableCoordinator().execute_task(
+ what=_make_ti(dag_id="orders"),
+ dag_rel_path="bin/orders",
+ bundle_info=BundleInfo(name="dags", version="v1"),
+ client=mock_client,
+ subprocess_logs_to_stdout=False,
+ )
+ finally:
+ reset_importer_registry()
+
+ mock_scan.assert_not_called()
+ assert mock_start.call_args.kwargs["command"][0] ==
str(binary.resolve())
diff --git
a/task-sdk/tests/task_sdk/coordinators/executable/test_dag_importer.py
b/task-sdk/tests/task_sdk/coordinators/executable/test_dag_importer.py
new file mode 100644
index 00000000000..b323318a229
--- /dev/null
+++ b/task-sdk/tests/task_sdk/coordinators/executable/test_dag_importer.py
@@ -0,0 +1,283 @@
+#
+# 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.
+from __future__ import annotations
+
+import json
+from pathlib import Path
+from types import SimpleNamespace
+from unittest import mock
+
+import pytest
+from task_sdk.coordinators.executable._bundle_test_utils import
ENTRYPOINT_PATH, write_bundle
+
+from airflow.sdk.coordinators._dag_importer import find_claiming_importer
+from airflow.sdk.coordinators.executable._bundle_reader import _load_index,
_open_checked_bundle
+from airflow.sdk.coordinators.executable._dag_importer import
ExecutableDagImporter
+from airflow.sdk.coordinators.executable.coordinator import
ExecutableCoordinator, _digest_cache
+from airflow.sdk.importers import (
+ DagSourceCode,
+ FilesystemDagDefinition,
+ get_importer_registry,
+ reset_importer_registry,
+)
+
+from tests_common.test_utils.config import conf_vars
+from tests_common.test_utils.version_compat import AIRFLOW_V_3_3_PLUS
+
+if not AIRFLOW_V_3_3_PLUS:
+ pytest.skip("Coordinator is only compatible with Airflow >= 3.3.0",
allow_module_level=True)
+
+EXECUTABLE_COORDINATOR =
"airflow.sdk.coordinators.executable.ExecutableCoordinator"
+ORDERS = b'package main\n\nfunc orders() { dag.New("orders") }\n'
+REPORTS = b'package dags\n\nfunc reports() { dag.New("reports") }\n'
+REPORTS_PATH = "example/bundle/dags/reports.go"
+COORDINATORS = {("sdk", "coordinators"): json.dumps({"go": {"classpath":
EXECUTABLE_COORDINATOR}})}
+
+
[email protected](autouse=True)
+def _clean_state():
+ _digest_cache.clear()
+ _load_index.cache_clear()
+ reset_importer_registry()
+ yield
+ reset_importer_registry()
+
+
[email protected]
+def importer() -> ExecutableDagImporter:
+ return ExecutableDagImporter(bundle_name="dags-folder")
+
+
+def _bundle(root: SimpleNamespace | Path) -> SimpleNamespace:
+ return SimpleNamespace(name="dags-folder", path=root)
+
+
+class TestCanHandle:
+ @pytest.mark.parametrize(
+ "name",
+ ["bin/orders", "orders", "/opt/airflow/dags/bin/orders"],
+ ids=["relative", "bare-name", "absolute"],
+ )
+ def test_claims_a_name_without_an_extension(self, importer, name):
+ assert importer.can_handle(name) is True
+ assert importer.can_handle(Path(name)) is True
+ assert importer.can_handle(FilesystemDagDefinition(Path(name))) is True
+
+ @pytest.mark.parametrize("name", ["bin/orders.bin", "dags/etl.py",
"x.jar", "orders.exe"])
+ def test_ignores_a_name_with_an_extension(self, importer, name):
+ assert importer.can_handle(name) is False
+ assert importer.can_handle(Path(name)) is False
+ assert importer.can_handle(FilesystemDagDefinition(Path(name))) is
False
+
+ def test_ignores_a_bundle_with_an_extension(self, importer, tmp_path):
+ path = write_bundle(tmp_path / "orders.bin", "orders")
+
+ assert importer.can_handle(path) is False
+
+ def test_does_not_read_the_file(self, importer, tmp_path, monkeypatch):
+ monkeypatch.chdir(tmp_path)
+
+ assert importer.can_handle("bin/orders") is True
+
+ def test_declares_no_extension(self, importer):
+ assert importer.artifact_suffix == ""
+ assert importer.supported_extensions == []
+
+
+class TestListDagDefinitions:
+ def test_lists_every_bundle_under_the_root(self, importer, tmp_path):
+ top = write_bundle(tmp_path / "top", "orders")
+ (tmp_path / "team").mkdir()
+ nested = write_bundle(tmp_path / "team" / "nested", "reports")
+ (tmp_path / "README.md").write_text("docs")
+ (tmp_path / "dag.py").write_text("x = 1\n")
+
+ definitions = list(importer.list_dag_definitions(_bundle(tmp_path)))
+
+ assert sorted(d.path for d in definitions) == sorted([top, nested])
+
+ def test_lists_only_the_bundles_without_an_extension(self, importer,
tmp_path):
+ (tmp_path / "bin").mkdir()
+ bundle = write_bundle(tmp_path / "bin" / "orders", "orders")
+ write_bundle(tmp_path / "bin" / "orders.bin", "orders")
+ (tmp_path / "README").write_text("docs")
+
+ definitions = list(importer.list_dag_definitions(_bundle(tmp_path)))
+
+ assert [d.path for d in definitions] == [bundle]
+
+ def test_honors_airflowignore(self, importer, tmp_path):
+ kept = write_bundle(tmp_path / "kept", "orders")
+ write_bundle(tmp_path / "handlers-only", "reports")
+ (tmp_path / ".airflowignore").write_text("handlers-only\n")
+
+ definitions = list(importer.list_dag_definitions(_bundle(tmp_path)))
+
+ assert [d.path for d in definitions] == [kept]
+
+ def test_a_single_file_root_scopes_the_listing_to_it(self, importer,
tmp_path):
+ target = write_bundle(tmp_path / "target", "orders")
+ write_bundle(tmp_path / "other", "reports")
+
+ assert [d.path for d in
importer.list_dag_definitions(_bundle(target))] == [target]
+
+ def test_a_single_file_root_that_is_not_a_bundle_lists_nothing(self,
importer, tmp_path):
+ plain = tmp_path / "plain"
+ plain.write_bytes(b"not a bundle")
+
+ assert list(importer.list_dag_definitions(_bundle(plain))) == []
+
+
+class TestMightContainDag:
+ @pytest.mark.parametrize("safe_mode", [True, False])
+ def test_keeps_every_bundle(self, importer, tmp_path, safe_mode):
+ handlers_only = write_bundle(tmp_path / "handlers", "orders",
dag_source_paths={}, sources={})
+
+ assert
importer.might_contain_dag(FilesystemDagDefinition(handlers_only), safe_mode)
is True
+
+ @pytest.mark.parametrize(
+ "content",
+ [b"", b"short", b"x" * 100, b"AFBNDL01 and then more bytes"],
+ ids=["empty", "shorter-than-magic", "plain", "magic-not-at-the-end"],
+ )
+ def test_drops_a_file_that_does_not_end_with_the_magic(self, importer,
tmp_path, content):
+ path = tmp_path / "plain"
+ path.write_bytes(content)
+
+ assert importer.might_contain_dag(FilesystemDagDefinition(path), True)
is False
+
+ @pytest.mark.parametrize("name", ["gone", "."])
+ def test_drops_a_file_it_cannot_read(self, importer, tmp_path, name):
+ assert importer.might_contain_dag(FilesystemDagDefinition(tmp_path /
name), True) is False
+
+
+class TestGetSourceCode:
+ @pytest.fixture
+ def bundle(self, tmp_path) -> FilesystemDagDefinition:
+ path = write_bundle(
+ tmp_path / "bundle",
+ "orders",
+ "reports",
+ sources={ENTRYPOINT_PATH: ORDERS, REPORTS_PATH: REPORTS},
+ dag_source_paths={"orders": ENTRYPOINT_PATH, "reports":
REPORTS_PATH},
+ )
+ return FilesystemDagDefinition(path)
+
+ def test_returns_each_dags_own_file(self, importer, bundle):
+ assert importer.get_source_code(bundle, "orders") ==
DagSourceCode(ORDERS.decode(), "go")
+ assert importer.get_source_code(bundle, "reports") ==
DagSourceCode(REPORTS.decode(), "go")
+
+ def test_returns_the_entrypoint_without_a_dag_id(self, importer, bundle):
+ assert importer.get_source_code(bundle) ==
DagSourceCode(ORDERS.decode(), "go")
+
+ def test_returns_the_entrypoint_for_an_unmapped_dag(self, importer,
bundle):
+ assert importer.get_source_code(bundle, "dynamic") ==
DagSourceCode(ORDERS.decode(), "go")
+
+ @mock.patch(
+
"airflow.sdk.coordinators.executable._bundle_reader._open_checked_bundle",
+ wraps=_open_checked_bundle,
+ )
+ def test_reads_the_bundle_index_once_for_many_dags(self, mock_open,
importer, bundle):
+ for dag_id in ("orders", "reports", "dynamic"):
+ importer.get_source_code(bundle, dag_id)
+
+ mock_open.assert_called_once()
+
+ def test_returns_a_notice_when_the_bundle_embeds_no_source(self, importer,
tmp_path):
+ path = write_bundle(tmp_path / "bundle", "orders", omit_sources=True)
+
+ source_code = importer.get_source_code(FilesystemDagDefinition(path),
"orders")
+
+ assert source_code == DagSourceCode(
+ "// Source code is not available: the bundle embeds no source.\n",
"go"
+ )
+
+ def test_falls_back_to_text_without_a_language(self, importer, tmp_path):
+ path = write_bundle(tmp_path / "bundle", "orders", language=None)
+
+ assert
importer.get_source_code(FilesystemDagDefinition(path)).language == "text"
+
+ def test_raises_for_an_invalid_bundle(self, importer, tmp_path):
+ path = write_bundle(tmp_path / "bundle", "orders",
entrypoint_path="missing.go")
+
+ with pytest.raises(ValueError, match="not one of its sources"):
+ importer.get_source_code(FilesystemDagDefinition(path))
+
+
+class TestRegistry:
+ def test_routes_a_bundle_to_the_importer(self, tmp_path):
+ path = write_bundle(tmp_path / "go-bundle", "orders")
+ with conf_vars(COORDINATORS):
+ routed = get_importer_registry("dags-folder").get_importer(path)
+ claiming = find_claiming_importer(path, "dags-folder")
+
+ assert isinstance(routed, ExecutableDagImporter)
+ assert isinstance(claiming, ExecutableDagImporter)
+ assert claiming.bundle_name == "dags-folder"
+
+ def test_does_not_route_a_bundle_with_an_extension(self, tmp_path):
+ path = write_bundle(tmp_path / "go-bundle.bin", "orders")
+ with conf_vars(COORDINATORS):
+ assert get_importer_registry("dags-folder").get_importer(path) is
None
+ assert find_claiming_importer(path, "dags-folder") is None
+
+ def test_leaves_other_files_to_the_python_importers(self, tmp_path):
+ python_file = tmp_path / "dag.py"
+ python_file.write_text("x = 1\n")
+ plain = tmp_path / "notes.txt"
+ plain.write_text("notes")
+ with conf_vars(COORDINATORS):
+ assert find_claiming_importer(python_file, "dags-folder") is None
+ assert find_claiming_importer(plain, "dags-folder") is None
+
+ def
test_a_python_file_ending_with_the_magic_still_belongs_to_the_python_importer(self,
tmp_path):
+ path = tmp_path / "dag.py"
+ path.write_bytes(b"from airflow.sdk import DAG\n" + b"AFBNDL01")
+ with conf_vars(COORDINATORS):
+ assert find_claiming_importer(path, "dags-folder") is None
+ listed = [
+ (type(importer).__name__, item.path)
+ for importer, item in
get_importer_registry("dags-folder").list_dag_definitions(
+ _bundle(tmp_path)
+ )
+ ]
+
+ assert [name for name, _ in listed] == ["PythonDagImporter"]
+
+ def
test_registry_listing_keeps_only_the_bundles_without_an_extension(self,
tmp_path):
+ suffixless = write_bundle(tmp_path / "go-bundle", "orders")
+ write_bundle(tmp_path / "other.bin", "reports")
+ with conf_vars(COORDINATORS):
+ listed = [
+ item.path
+ for importer, item in
get_importer_registry("dags-folder").list_dag_definitions(
+ _bundle(tmp_path)
+ )
+ if isinstance(importer, ExecutableDagImporter)
+ ]
+
+ assert listed == [suffixless]
+
+ def test_registers_nothing_without_an_executable_coordinator(self,
tmp_path):
+ path = write_bundle(tmp_path / "go-bundle", "orders")
+ with conf_vars({("sdk", "coordinators"): "{}"}):
+ assert find_claiming_importer(path, "dags-folder") is None
+
+ def test_parsing_coordinator_is_the_executable_coordinator(self, importer):
+ with conf_vars(COORDINATORS):
+ assert isinstance(importer.get_parsing_coordinator(),
ExecutableCoordinator)