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