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]