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]

Reply via email to