kaxil commented on code in PR #74037:
URL: https://github.com/apache/airflow/pull/74037#discussion_r4159730581
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst:
##########
@@ -514,6 +516,42 @@ instance's Dag. The artifact's name does not matter beyond
that suffix, so one r
and a Dag is routed to whichever declares it. If multiple bundles declare the
same Dag, the first configured
root wins, and within a root the first in sorted path order.
+.. _typescript-sdk/native-parsing:
+
+Parsing native Dags
+-------------------
+
+A Dag declared in TypeScript ships in its packed bundle, so the bundle goes
into a Dag bundle,
+next to any Python Dag files. A coordinator configured without
``bundles_root`` parses the ``*.min.mjs``
+bundles of the Dag bundles it serves: every Dag bundle, or only the one
``dag_bundle_name`` names.
Review Comment:
In `dag_bundle_name` mode, `_init_root_source` resolves
`BundleInfo(name=dag_bundle_name)` with no version, so tasks run whatever
version that bundle is on when they start, not the one their run was created
from. Before this layer that bundle only held handlers for a Python Dag; now it
can be the native Dag's own bundle, so a task renamed in v2 fails mid-run for a
v1 run. Should this use the TI's `bundle_info` when its name matches, or should
this section say runs aren't version-pinned in that mode? Java has the same
exposure.
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst:
##########
@@ -514,6 +516,42 @@ instance's Dag. The artifact's name does not matter beyond
that suffix, so one r
and a Dag is routed to whichever declares it. If multiple bundles declare the
same Dag, the first configured
root wins, and within a root the first in sorted path order.
+.. _typescript-sdk/native-parsing:
+
+Parsing native Dags
+-------------------
+
+A Dag declared in TypeScript ships in its packed bundle, so the bundle goes
into a Dag bundle,
+next to any Python Dag files. A coordinator configured without
``bundles_root`` parses the ``*.min.mjs``
+bundles of the Dag bundles it serves: every Dag bundle, or only the one
``dag_bundle_name`` names.
+The same entry runs the tasks:
+
+.. code-block:: ini
+
+ [sdk]
+ coordinators = {
+ "ts": {
+ "classpath": "airflow.sdk.coordinators.node.NodeCoordinator",
+ "kwargs": {"node_executable": "/usr/local/bin/node"}
+ }
+ }
+ queue_to_coordinator = {"typescript": "ts"}
+
+At most one such coordinator can serve a Dag bundle, and a coordinator with
``bundles_root`` parses nothing.
+
+The Dag processor runs each bundle with ``node`` to parse it, so it needs
Node.js and this ``[sdk]``
Review Comment:
At this layer the manager still queues only `.py` files and zip archives, so
a `bundle.min.mjs` never reaches `NodeDagImporter` until the importer-driven
discovery lands (#73841 here, #74020 upstream). Could the merge note name it
next to #73442 and #73445? Also, the prerequisite at line 51 reads as if Node
is only needed to parse TypeScript Dags, but a coordinator without
`bundles_root` runs `node` for every headered `*.min.mjs` in a served bundle,
handler-only ones included, and each gets an import error if `node` is missing.
##########
task-sdk/src/airflow/sdk/coordinators/node/coordinator.py:
##########
@@ -134,3 +134,12 @@ def _build_execute_task_command(self, *, what:
TaskInstance) -> tuple[list[str],
roots = self._get_scan_roots()
bundle = _Bundle.find(roots, what.dag_id)
return [self.node_executable, os.fspath(bundle.path)],
bundle.schema_version
+
+ @classmethod
+ def get_dag_importer_class(cls) -> type[NodeDagImporter]:
+ from airflow.sdk.coordinators.node._dag_importer import NodeDagImporter
Review Comment:
`_dag_importer` only imports `coordinator` under `TYPE_CHECKING`, so this
can be a top-level import, as `JavaCoordinator` does with `JavaDagImporter`,
and the `TYPE_CHECKING` copy at line 45 can go.
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst:
##########
@@ -514,6 +516,42 @@ instance's Dag. The artifact's name does not matter beyond
that suffix, so one r
and a Dag is routed to whichever declares it. If multiple bundles declare the
same Dag, the first configured
root wins, and within a root the first in sorted path order.
+.. _typescript-sdk/native-parsing:
+
+Parsing native Dags
+-------------------
+
+A Dag declared in TypeScript ships in its packed bundle, so the bundle goes
into a Dag bundle,
+next to any Python Dag files. A coordinator configured without
``bundles_root`` parses the ``*.min.mjs``
+bundles of the Dag bundles it serves: every Dag bundle, or only the one
``dag_bundle_name`` names.
+The same entry runs the tasks:
+
+.. code-block:: ini
+
+ [sdk]
+ coordinators = {
+ "ts": {
+ "classpath": "airflow.sdk.coordinators.node.NodeCoordinator",
+ "kwargs": {"node_executable": "/usr/local/bin/node"}
+ }
+ }
+ queue_to_coordinator = {"typescript": "ts"}
+
+At most one such coordinator can serve a Dag bundle, and a coordinator with
``bundles_root`` parses nothing.
+
+The Dag processor runs each bundle with ``node`` to parse it, so it needs
Node.js and this ``[sdk]``
+configuration. It parses only files that end in ``.min.mjs`` and start with
the header ``airflow-ts-pack``
+writes. A bundle that fails its integrity check, or whose Dags cannot be
serialized or fail validation,
Review Comment:
java.rst carries "Do not declare a Dag in Java that a Python file in the
same bundle also defines", and this page needs it more: the Quick start's
Python stub `typescript_example` binds the same `dagId`, so moving that Dag to
`new Dag("typescript_example")` and keeping the stub gives two files that
overwrite each other's Dag on every parse.
##########
task-sdk/src/airflow/sdk/coordinators/node/_dag_importer.py:
##########
@@ -0,0 +1,71 @@
+#
+# 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 that parses native TypeScript Dags from packed bundles."""
+
+from __future__ import annotations
+
+from typing import TYPE_CHECKING, Final
+
+import structlog
+
+from airflow.sdk.coordinators._dag_importer import CoordinatorDagImporter
+from airflow.sdk.coordinators.node._bundle_reader import (
+ BUNDLE_SUFFIX,
+ has_bundle_layout_prefix,
+ read_bundle_entrypoint_source,
+)
+from airflow.sdk.importers.base import DagSourceCode
+
+if TYPE_CHECKING:
+ from structlog.typing import FilteringBoundLogger
+
+ from airflow.sdk.coordinators.node.coordinator import NodeCoordinator
+ from airflow.sdk.importers.base import DagDefinition
+
+log: FilteringBoundLogger =
structlog.get_logger(logger_name="coordinators.node")
+
+_NO_SOURCE: Final = "// Source code is not available: the bundle embeds no
entrypoint source.\n"
+
+
+class NodeDagImporter(CoordinatorDagImporter):
+ """Parse the native Dags of packed ``*.min.mjs`` TypeScript bundles with a
:class:`NodeCoordinator`."""
+
+ artifact_suffix = BUNDLE_SUFFIX
+ supported_extensions = [".mjs"]
+ coordinator: NodeCoordinator
+
+ def might_contain_dag(self, definition: DagDefinition, safe_mode: bool) ->
bool:
+ # The header identifies a packed bundle rather than guessing at
content, so safe_mode keeps it.
+ with definition.as_file() as path:
+ try:
+ return has_bundle_layout_prefix(path)
+ except OSError:
+ # Keep the file, so parsing it records why it cannot be read.
+ return True
+
+ def get_source_code(self, definition: DagDefinition) -> DagSourceCode:
+ """Return the embedded entrypoint source of the bundle, or a notice
when there is none."""
+ with definition.as_file() as path:
+ try:
+ source = read_bundle_entrypoint_source(path)
+ except (OSError, ValueError) as e:
Review Comment:
Both callers of `get_source_code` already catch any exception and store a
placeholder (`_read_dag_source_codes` and the Dag bag), and the Java and Python
importers let errors through. Dropping this `except` would report a broken
bundle in one place, with one wording.
##########
task-sdk/tests/task_sdk/coordinators/node/test_dag_importer.py:
##########
@@ -0,0 +1,139 @@
+#
+# 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 pathlib
+from types import SimpleNamespace
+
+import pytest
+from task_sdk.coordinators.node._bundle_test_utils import (
+ BUNDLE_NAME,
+ LAYOUT_PREFIX,
+ SOURCE_OPEN,
+ mutate_byte,
+ mutate_section,
+ write_bundle,
+)
+
+from airflow.sdk.coordinators.node._bundle_reader import _digest_cache
+from airflow.sdk.coordinators.node._dag_importer import NodeDagImporter
+from airflow.sdk.coordinators.node.coordinator import NodeCoordinator
+from airflow.sdk.importers import DagSourceCode, FilesystemDagDefinition
+
+
[email protected](autouse=True)
+def clear_digest_cache():
+ _digest_cache.clear()
+
+
[email protected]
+def importer() -> NodeDagImporter:
+ return NodeDagImporter(coordinator=NodeCoordinator())
+
+
+def test_lists_only_packed_bundles(importer, tmp_path):
+ nested = write_bundle(tmp_path / "team", "sales")
+ tampered = write_bundle(tmp_path, "inventory", name="tampered.min.mjs")
+ mutate_section(tampered, "code")
+ write_bundle(tmp_path, "orders", name="plain.mjs")
+ (tmp_path / "vendor.min.mjs").write_bytes(b"export {};\n")
+
+ definitions =
importer.list_dag_definitions(SimpleNamespace(name="dags-folder",
path=tmp_path))
+
+ # A tampered bundle keeps its header, so parsing it reports the integrity
failure.
+ assert sorted(d.path for d in definitions) == sorted([nested, tampered])
+
+
+class TestMightContainDag:
+ @pytest.mark.parametrize("safe_mode", [True, False])
+ @pytest.mark.parametrize(
+ ("content", "expected"),
+ [
+ (None, True),
+ (b"export {};\n", False),
+ (b"", False),
+ (LAYOUT_PREFIX[:-1], False),
+ ],
+ ids=["bundle", "plain-module", "empty", "truncated-prefix"],
+ )
+ def test_checks_the_layout_header(self, importer, tmp_path, content,
expected, safe_mode):
+ path = write_bundle(tmp_path, "sales")
+ if content is not None:
+ path.write_bytes(content)
+
+ assert importer.might_contain_dag(FilesystemDagDefinition(path),
safe_mode) is expected
+
+ def test_keeps_an_unreadable_file(self, importer, tmp_path, monkeypatch):
+ path = write_bundle(tmp_path, "sales")
+ original_open = pathlib.Path.open
+
+ def raise_permission_error(self, *args, **kwargs):
+ if self.name == BUNDLE_NAME:
+ raise PermissionError("denied")
+ return original_open(self, *args, **kwargs)
+
+ monkeypatch.setattr(pathlib.Path, "open", raise_permission_error)
+
+ assert importer.might_contain_dag(FilesystemDagDefinition(path), True)
is True
+
+
+class TestGetSourceCode:
+ def test_returns_the_entrypoint_source_for_every_dag(self, importer,
tmp_path):
+ sales = '/* the */ export const sales = new Dag({ dagId: "sales" });\n'
+ inventory = 'export const inventory = new Dag({ dagId: "inventory"
});\n'
+ path = write_bundle(
+ tmp_path,
+ "sales",
+ "inventory",
+ sources=[("sales.ts", sales.encode()), ("inventory.ts",
inventory.encode())],
+ dag_source_paths={"sales": "sales.ts", "inventory":
"inventory.ts"},
+ entrypoint_path="sales.ts",
Review Comment:
`entrypoint_path` is `sales.ts` here, which is also the first Dag's mapped
file, so a resolver that returned the first `dag_source_paths` entry instead of
the entry module still passes (same in test_bundle_reader.py). A separate
`main.ts` entrypoint would pin the Code view choice the description makes.
##########
airflow-core/docs/authoring-and-scheduling/language-sdks/typescript.rst:
##########
@@ -514,6 +516,42 @@ instance's Dag. The artifact's name does not matter beyond
that suffix, so one r
and a Dag is routed to whichever declares it. If multiple bundles declare the
same Dag, the first configured
root wins, and within a root the first in sorted path order.
+.. _typescript-sdk/native-parsing:
+
+Parsing native Dags
+-------------------
+
+A Dag declared in TypeScript ships in its packed bundle, so the bundle goes
into a Dag bundle,
+next to any Python Dag files. A coordinator configured without
``bundles_root`` parses the ``*.min.mjs``
+bundles of the Dag bundles it serves: every Dag bundle, or only the one
``dag_bundle_name`` names.
+The same entry runs the tasks:
+
+.. code-block:: ini
+
+ [sdk]
+ coordinators = {
+ "ts": {
+ "classpath": "airflow.sdk.coordinators.node.NodeCoordinator",
+ "kwargs": {"node_executable": "/usr/local/bin/node"}
+ }
+ }
+ queue_to_coordinator = {"typescript": "ts"}
+
+At most one such coordinator can serve a Dag bundle, and a coordinator with
``bundles_root`` parses nothing.
+
+The Dag processor runs each bundle with ``node`` to parse it, so it needs
Node.js and this ``[sdk]``
+configuration. It parses only files that end in ``.min.mjs`` and start with
the header ``airflow-ts-pack``
+writes. A bundle that fails its integrity check, or whose Dags cannot be
serialized or fail validation,
+such as a cycle drawn with ``before`` and ``after``, is reported as an import
error.
+
+The Code view shows the bundle's entry module for each of its Dags, as
``airflow-ts-pack`` embeds it
Review Comment:
The build section above still says the packer embeds each Dag's own file "so
Airflow has something readable to display for each Dag", but this shows the
entry module for every Dag, so a Dag declared in `src/dags/etl.ts` and imported
from `src/main.ts` shows `main.ts`. The bundle already carries
`dag_source_paths`, and `read_bundle_source(path, dag_id)` resolves it but has
no production caller. Is the entry module the end state, or a stopgap until
`_read_dag_source_codes` can ask for a Dag's own source? `dag_source_codes` is
new in the 2026-10-30 schema, which no release has shipped yet, so a per-Dag
override there is cheap now and a migration later. Either way the two
paragraphs should agree.
##########
task-sdk/src/airflow/sdk/coordinators/node/coordinator.py:
##########
@@ -134,3 +134,12 @@ def _build_execute_task_command(self, *, what:
TaskInstance) -> tuple[list[str],
roots = self._get_scan_roots()
bundle = _Bundle.find(roots, what.dag_id)
Review Comment:
The parse runs the exact file the Dag came from, but a task still takes the
first sorted `*.min.mjs` that declares its dag_id, and `walk_files` doesn't
read `.airflowignore`. So with `archive/sales.min.mjs` ignored and
`sales.min.mjs` at the root, the Dag and its Code view come from the root file
while every task runs the archived copy; with two unignored bundles sharing a
dag_id, the last-parsed graph is stored while tasks run the first-sorted one.
For a native Dag, could the task run the TI's `dag_rel_path` when it names a
`.min.mjs` (checking it still declares the dag_id)? At minimum, "The same entry
runs the tasks" in typescript.rst overstates it.
--
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]