This is an automated email from the ASF dual-hosted git repository.

jason810496 pushed a commit to branch jason/lang-sdk-e2e/03e-lang-sdk-dag-bag
in repository https://gitbox.apache.org/repos/asf/airflow.git

commit 1610ac78b417540affce715bac6a8fcc4359f222
Author: ZHE YOU LIU <[email protected]>
AuthorDate: Wed Sep 30 03:03:03 2026 +0000

    Report a native Dag's task that reaches a Python worker
    
    A Python worker handed a task from a Dag file that a coordinator's Dag
    importer claims now exits with an error naming the queue to route. It
    checks before building a Dag bag, so it never starts the runtime.
---
 .../src/airflow/sdk/execution_time/task_runner.py  | 23 ++++++
 .../task_sdk/execution_time/test_task_runner.py    | 92 ++++++++++++++++++++++
 2 files changed, 115 insertions(+)

diff --git a/task-sdk/src/airflow/sdk/execution_time/task_runner.py 
b/task-sdk/src/airflow/sdk/execution_time/task_runner.py
index ed622023308..c4bf7f9f5b2 100644
--- a/task-sdk/src/airflow/sdk/execution_time/task_runner.py
+++ b/task-sdk/src/airflow/sdk/execution_time/task_runner.py
@@ -60,6 +60,7 @@ from airflow.sdk.bases.operator import BaseOperator, 
ExecutorSafeguard
 from airflow.sdk.bases.skipmixin import XCOM_SKIPMIXIN_KEY
 from airflow.sdk.bases.xcom import BaseXCom
 from airflow.sdk.configuration import conf
+from airflow.sdk.coordinators._dag_importer import CoordinatorDagImporter
 from airflow.sdk.definitions._internal.dag_parsing_context import 
_airflow_parsing_context_manager
 from airflow.sdk.definitions._internal.types import NOTSET, ArgNotSet, 
is_arg_set
 from airflow.sdk.definitions.asset import (
@@ -150,6 +151,7 @@ from airflow.sdk.execution_time.email_backend import (
 from airflow.sdk.execution_time.sentry import Sentry
 from airflow.sdk.execution_time.tracing import detail_span
 from airflow.sdk.execution_time.xcom import XCom
+from airflow.sdk.importers import get_importer_registry
 from airflow.sdk.listener import get_listener_manager
 from airflow.sdk.observability.metrics import stats_utils
 from airflow.sdk.serde import allow_class, iter_pydantic_models
@@ -1018,6 +1020,16 @@ def _register_deserialization_allowed_classes(dag, log: 
Logger) -> None:
                     )
 
 
+def _is_lang_sdk_dag_file(path: str, bundle_name: str) -> bool:
+    """Return whether a coordinator's Dag importer claims *path*, so that a 
Lang-SDK runtime parses it."""
+    try:
+        importer = get_importer_registry(bundle_name).get_importer(path)
+    except Exception:
+        # Building the Dag bag reports a broken importer configuration.
+        return False
+    return isinstance(importer, CoordinatorDagImporter)
+
+
 @detail_span("parse")
 def parse(what: StartupDetails, log: Logger) -> RuntimeTaskInstance:
     # TODO: Task-SDK:
@@ -1030,6 +1042,17 @@ def parse(what: StartupDetails, log: Logger) -> 
RuntimeTaskInstance:
     bundle_prepare_ms = int((time.monotonic() - bundle_prepare_start) * 1000)
 
     dag_absolute_path = os.fspath(Path(bundle_instance.path, 
what.dag_rel_path))
+    if _is_lang_sdk_dag_file(dag_absolute_path, bundle_info.name):
+        log.error(
+            "A task of a native Lang-SDK Dag cannot run in Python. Route its 
queue to the coordinator "
+            "that parses the Dag, with [sdk] queue_to_coordinator",
+            dag_id=what.ti.dag_id,
+            task_id=what.ti.task_id,
+            queue=what.ti.queue,
+            path=what.dag_rel_path,
+        )
+        sys.exit(1)
+
     dag_file_parse_start = time.monotonic()
     bag = BundleDagBag(
         dag_folder=dag_absolute_path,
diff --git a/task-sdk/tests/task_sdk/execution_time/test_task_runner.py 
b/task-sdk/tests/task_sdk/execution_time/test_task_runner.py
index 9cdb75446b1..26ae41ac01e 100644
--- a/task-sdk/tests/task_sdk/execution_time/test_task_runner.py
+++ b/task-sdk/tests/task_sdk/execution_time/test_task_runner.py
@@ -31,6 +31,7 @@ from typing import TYPE_CHECKING, Any
 from unittest import mock
 from unittest.mock import call, patch
 
+import attrs
 import pandas as pd
 import pytest
 import structlog
@@ -78,6 +79,8 @@ from airflow.sdk.api.datamodels._generated import (
 )
 from airflow.sdk.bases.operator import ExecutorSafeguard
 from airflow.sdk.bases.xcom import BaseXCom
+from airflow.sdk.coordinators._dag_importer import CoordinatorDagImporter
+from airflow.sdk.coordinators._subprocess import SubprocessCoordinator
 from airflow.sdk.definitions._internal.types import NOTSET, 
SET_DURING_EXECUTION, is_arg_set
 from airflow.sdk.definitions.asset import Asset, AssetAlias, AssetUniqueKey, 
AssetUriRef, Dataset, Model
 from airflow.sdk.definitions.param import DagParam
@@ -176,6 +179,7 @@ from airflow.sdk.execution_time.task_runner import (
     TaskRunnerMarker,
     _defer_task,
     _execute_task,
+    _is_lang_sdk_dag_file,
     _make_task_span,
     _push_xcom_if_needed,
     _register_deserialization_allowed_classes,
@@ -189,6 +193,7 @@ from airflow.sdk.execution_time.task_runner import (
     startup,
 )
 from airflow.sdk.execution_time.xcom import XCom
+from airflow.sdk.importers import DagSourceCode, reset_importer_registry
 from airflow.sdk.serde import deserialize
 from airflow.triggers.base import BaseEventTrigger, BaseTrigger, TriggerEvent
 from airflow.triggers.callback import CallbackTrigger
@@ -902,6 +907,93 @@ def test_parse_module_in_bundle_root(tmp_path: Path, 
make_ti_context):
     assert ti.task.dag.dag_id == "dag_name"
 
 
+class NativeDagImporter(CoordinatorDagImporter):
+    artifact_suffix = ".native"
+    supported_extensions = [".native"]
+
+    def get_source_code(self, definition):
+        return DagSourceCode(source_code=definition.read_text(), 
language="native")
+
+
[email protected](kw_only=True)
+class NativeCoordinator(SubprocessCoordinator):
+    """A coordinator whose Dag importer claims ``.native`` files in every 
bundle."""
+
+    def get_dag_importer(self):
+        return NativeDagImporter(coordinator=self)
+
+
+@patch("airflow.dag_processing.dagbag.BundleDagBag", autospec=True)
+def test_parse_rejects_a_task_of_a_native_dag(mock_bag, tmp_path: Path, 
make_ti_context):
+    tmp_path.joinpath("dag.native").write_text("{}")
+    what = StartupDetails(
+        ti=TaskInstance(
+            id=uuid7(),
+            task_id="a",
+            dag_id="native_dag",
+            run_id="c",
+            try_number=1,
+            dag_version_id=uuid7(),
+            queue="default",
+        ),
+        dag_rel_path="dag.native",
+        bundle_info=BundleInfo(name="my-bundle", version=None),
+        ti_context=make_ti_context(),
+        start_date=timezone.utcnow(),
+        sentry_integration="",
+    )
+    bundle_config = [
+        {
+            "name": "my-bundle",
+            "classpath": "airflow.dag_processing.bundles.local.LocalDagBundle",
+            "kwargs": {"path": str(tmp_path), "refresh_interval": 1},
+        }
+    ]
+    coordinators = {"native": {"classpath": f"{__name__}.NativeCoordinator", 
"kwargs": {}}}
+    log = mock.Mock()
+
+    reset_importer_registry()
+    try:
+        with (
+            patch.dict(
+                os.environ,
+                {
+                    "AIRFLOW__DAG_PROCESSOR__DAG_BUNDLE_CONFIG_LIST": 
json.dumps(bundle_config),
+                    "AIRFLOW__SDK__COORDINATORS": json.dumps(coordinators),
+                },
+            ),
+            pytest.raises(SystemExit, match="1"),
+        ):
+            parse(what, log)
+    finally:
+        reset_importer_registry()
+
+    mock_bag.assert_not_called()
+    log.error.assert_called_once_with(
+        "A task of a native Lang-SDK Dag cannot run in Python. Route its queue 
to the coordinator "
+        "that parses the Dag, with [sdk] queue_to_coordinator",
+        dag_id="native_dag",
+        task_id="a",
+        queue="default",
+        path="dag.native",
+    )
+
+
[email protected](
+    ("file_name", "expected"),
+    [("dag.native", True), ("dag.py", False), ("dag.pyc", False), ("dags.zip", 
False)],
+)
+def test_is_lang_sdk_dag_file(file_name, expected):
+    coordinators = {"native": {"classpath": f"{__name__}.NativeCoordinator", 
"kwargs": {}}}
+
+    reset_importer_registry()
+    try:
+        with patch.dict(os.environ, {"AIRFLOW__SDK__COORDINATORS": 
json.dumps(coordinators)}):
+            assert _is_lang_sdk_dag_file(f"/bundle/{file_name}", "my-bundle") 
is expected
+    finally:
+        reset_importer_registry()
+
+
 @pytest.mark.parametrize("use_queues", [False, True])
 def test_run_deferred_basic(time_machine, create_runtime_ti, 
mock_supervisor_comms, use_queues: bool):
     """Test that a task can transition to a deferred state."""

Reply via email to