jason810496 commented on code in PR #74037:
URL: https://github.com/apache/airflow/pull/74037#discussion_r4163183828


##########
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:
   Done. The description now names #74020, and 71f9d62ddf says the Dag 
processor runs `node` on every packed bundle it serves, including handler-only 
ones.



##########
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:
   For now it's a stopgap; showing each Dag's own source is a follow-up. 
574c46a7ab makes both paragraphs say the Code view shows the entry module.



##########
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:
   Done in 6a997106e7. A native Dag's task runs the bundle named by the task 
instance's `dag_rel_path`, if that bundle still declares the dag_id. Otherwise 
it falls back to the sorted search. The docs now say so.



##########
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:
   Done in 25b8a2d858.



##########
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:
   Done in fed3ceadbc.



##########
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:
   Done in 44bef1e376.



##########
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:
   Done in 4e5d37bf8f. Both tests now use a separate `main.ts` entry point.



##########
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:
   Done in f4596dd224 on #74036. When `dag_bundle_name` names the task's own 
bundle, the task uses the version its run was created with. Both pages say so.



-- 
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