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 e07e905051510d7a32c92438829dde17bcb53d99 Author: ZHE YOU LIU <[email protected]> AuthorDate: Wed Sep 30 02:58:52 2026 +0000 Report a Lang-SDK parse past its import timeout as an import error The parse child resolves the get_dagbag_import_timeout policy for the file and sends the timeout with the runtime's schema version. The policy thus runs where it does for a Python Dag file, and a failure is only that file's import error. The timeout counts from the start of the parse. Past it, with no result yet, the runtime is killed and the file gets an import error. LangSDKDagFileProcessorProcess.run applies the same timeout, so it no longer takes one. Until the child reports it, dag_file_processor_timeout bounds the parse, as it does in the Dag processor. At the deadline run stops waiting even if an exited runtime left a process that holds its output open. --- .../airflow/dag_processing/lang_sdk_processor.py | 58 ++++++++++---- .../unit/dag_processing/test_lang_sdk_processor.py | 88 ++++++++++++++++++---- .../tests/unit/dag_processing/test_manager.py | 57 ++++++++++++++ 3 files changed, 174 insertions(+), 29 deletions(-) diff --git a/airflow-core/src/airflow/dag_processing/lang_sdk_processor.py b/airflow-core/src/airflow/dag_processing/lang_sdk_processor.py index 877af756ed6..dc7d2fca7bc 100644 --- a/airflow-core/src/airflow/dag_processing/lang_sdk_processor.py +++ b/airflow-core/src/airflow/dag_processing/lang_sdk_processor.py @@ -32,6 +32,8 @@ import msgspec from pydantic import BaseModel, Field, TypeAdapter from uuid6 import uuid7 +from airflow import settings +from airflow.configuration import conf from airflow.dag_processing.importer_routing import get_claiming_coordinator from airflow.dag_processing.processor import ( BaseDagFileProcessorProcess, @@ -83,12 +85,22 @@ class StartLangSDKRuntime(BaseModel): class LangSDKRuntimeSchemaVersion(BaseModel): - """The supervisor schema version of the runtime the parse child is about to exec.""" + """The schema version and the import timeout of the runtime the parse child is about to exec.""" schema_version: str | None + import_timeout: float | None = None + """Seconds from the start of the parse; ``None`` means no timeout.""" type: Literal["LangSDKRuntimeSchemaVersion"] = "LangSDKRuntimeSchemaVersion" +def _get_import_timeout(path: str) -> float | None: + """Return the ``get_dagbag_import_timeout`` policy's timeout for *path*; ``None`` means none.""" + timeout = settings.get_dagbag_import_timeout(path) + if not isinstance(timeout, (int, float)): + raise TypeError(f"Value ({timeout}) from get_dagbag_import_timeout must be int or float") + return timeout if timeout > 0 else None + + def _start_runtime_entrypoint() -> None: """Exec the runtime that parses the file named by the start request, or report why it cannot start.""" os.environ["_AIRFLOW_PROCESS_CONTEXT"] = "client" @@ -106,9 +118,11 @@ def _start_runtime_entrypoint() -> None: raise RuntimeError(f"Required first message to be a StartLangSDKRuntime, it was {msg}") def report_schema_version(schema_version: str | None) -> None: - comms.send(LangSDKRuntimeSchemaVersion(schema_version=schema_version)) + comms.send(LangSDKRuntimeSchemaVersion(schema_version=schema_version, import_timeout=import_timeout)) try: + # The policy is user code: it runs in this child, where a failure is only this file's import error. + import_timeout = _get_import_timeout(msg.file) coordinator = get_claiming_coordinator(msg.file, msg.bundle_name) if coordinator is None: raise RuntimeError(f"No coordinator's Dag importer claims {msg.file}") @@ -154,6 +168,7 @@ class LangSDKDagFileProcessorProcess(BaseDagFileProcessorProcess): _listeners: dict[_Channel, socket] _parse_request: DagFileParseRequest _runtime_schema_version: str | None = attrs.field(default=None, init=False) + _import_timeout: float | None = attrs.field(default=None, init=False) _schema_version_reported: bool = attrs.field(default=False, init=False) _parsing_result_monotonic: float | None = attrs.field(default=None, init=False) _unverified_connections: list[tuple[socket, _Channel]] = attrs.field(factory=list, init=False) @@ -217,17 +232,16 @@ class LangSDKDagFileProcessorProcess(BaseDagFileProcessorProcess): bundle_path: Path, bundle_name: str, dag_file_rel_path: str, - timeout: float | None, logger: FilteringBoundLogger, ) -> DagFileParsingResult: """ Parse *path* outside the Dag processor and wait for the result. - There is no API client, so each request of the runtime that needs one gets an error. - - :raises TimeoutError: if the parse does not finish within *timeout* seconds. The runtime is - killed. + There is no API client, so each request of the runtime that needs one gets an error. The file's import + timeout bounds the parse, and ``[dag_processor] dag_file_processor_timeout`` until the parse child + reports it. """ + processor_timeout = conf.getfloat("dag_processor", "dag_file_processor_timeout") with selectors.DefaultSelector() as selector: proc = cls.start( id=uuid7(), @@ -240,14 +254,13 @@ class LangSDKDagFileProcessorProcess(BaseDagFileProcessorProcess): ) try: while not proc.is_ready: - wait = 0.1 - if timeout is not None: - if (remaining := proc.start_time + timeout - time.monotonic()) <= 0: - raise TimeoutError( - f"The Lang-SDK runtime did not parse {os.fspath(path)} within {timeout}s" - ) - wait = min(wait, remaining) - proc._service_subprocess(max_wait_time=wait) + timeout = proc._import_timeout if proc._schema_version_reported else processor_timeout + if timeout is not None and time.monotonic() - proc.start_time > timeout: + # Unlike is_ready, this does not wait for an exited runtime's leftover processes, + # which can hold its sockets open. close() closes them. + proc._time_out(timeout) + break + proc._service_subprocess(max_wait_time=0.1) except BaseException: proc._kill_runtime() raise @@ -381,6 +394,7 @@ class LangSDKDagFileProcessorProcess(BaseDagFileProcessorProcess): self._reject_request(msg, log, req_id) return ResponseSent.ALREADY_SENT self._runtime_schema_version = msg.schema_version + self._import_timeout = msg.import_timeout self._schema_version_reported = True return None, {} @@ -466,6 +480,13 @@ class LangSDKDagFileProcessorProcess(BaseDagFileProcessorProcess): ): self.process_log.warning("The Lang-SDK runtime did not exit after its parse result; killing it") self._kill_runtime() + if ( + self._import_timeout is not None + and self.parsing_result is None + and self._exit_code is None + and time.monotonic() - self.start_time > self._import_timeout + ): + self._time_out(self._import_timeout) if self._check_subprocess_exit() is None: return False self._close_listeners() @@ -477,6 +498,13 @@ class LangSDKDagFileProcessorProcess(BaseDagFileProcessorProcess): ) return True + def _time_out(self, timeout: float) -> None: + if self.parsing_result is None: + self._set_import_error( + f"The Lang-SDK runtime did not parse {self._parse_request.file} within {timeout}s" + ) + self._kill_runtime() + def _kill_runtime(self) -> None: """Kill the runtime and wait for it, without servicing its sockets, whose handler may have failed.""" if self._exit_code is not None: diff --git a/airflow-core/tests/unit/dag_processing/test_lang_sdk_processor.py b/airflow-core/tests/unit/dag_processing/test_lang_sdk_processor.py index 4a5fd758681..dd3075a7c33 100644 --- a/airflow-core/tests/unit/dag_processing/test_lang_sdk_processor.py +++ b/airflow-core/tests/unit/dag_processing/test_lang_sdk_processor.py @@ -36,6 +36,7 @@ from airflow.configuration import conf from airflow.dag_processing.lang_sdk_processor import ( LangSDKDagFileProcessorProcess, LangSDKRuntimeSchemaVersion, + _get_import_timeout, ) from airflow.dag_processing.processor import DagFileParseRequest, DagFileParsingResult from airflow.sdk import DAG, BaseOperator @@ -336,6 +337,27 @@ class TestLangSDKDagFileProcessorProcess: assert proc._exit_code == -signal.SIGKILL assert "The Lang-SDK runtime did not exit after its parse result; killing it" in cap_structlog + @pytest.mark.parametrize( + ("policy", "error"), + [ + pytest.param( + {"side_effect": RuntimeError("policy bug")}, "RuntimeError: policy bug", id="raises" + ), + pytest.param( + {"return_value": "30"}, + "TypeError: Value (30) from get_dagbag_import_timeout must be int or float", + id="not-a-number", + ), + ], + ) + def test_a_failing_import_timeout_policy_is_an_import_error(self, parse, policy, error): + with patch("airflow.settings.get_dagbag_import_timeout", autospec=True, **policy): + proc = parse() + + assert proc.parsing_result.import_errors == { + "dag.native": f"Cannot start the Lang-SDK runtime: {error}" + } + @pytest.mark.skipif(not Path("/proc/self/fd").is_dir(), reason="reads /proc") @pytest.mark.parametrize("use_exec", [False, True], ids=["fork", "spawn"]) def test_the_runtime_inherits_only_its_standard_streams(self, monkeypatch, tmp_path, use_exec): @@ -361,13 +383,12 @@ class TestLangSDKDagFileProcessorProcess: class TestRun: @staticmethod - def _run(tmp_path, *, timeout: float | None = 30) -> DagFileParsingResult: + def _run(tmp_path, **spec) -> DagFileParsingResult: return LangSDKDagFileProcessorProcess.run( - path=write_native_file(tmp_path / "dag.native"), + path=write_native_file(tmp_path / "dag.native", **spec), bundle_path=tmp_path, bundle_name="testing", dag_file_rel_path="dag.native", - timeout=timeout, logger=structlog.get_logger(), ) @@ -402,9 +423,9 @@ class TestRun: mock_mask_secret.assert_called_once_with("native-secret", "native_conn") @pytest.mark.parametrize("connected", [False, True], ids=["before-connecting", "after-connecting"]) - def test_a_parse_past_its_timeout_is_killed(self, tmp_path, connected): + @patch("airflow.settings.get_dagbag_import_timeout", autospec=True, return_value=1) + def test_a_parse_past_the_import_timeout_is_killed(self, mock_timeout, tmp_path, connected): # A runtime that never connects leaves both listeners open when it is killed. - write_native_file(tmp_path / "dag.native", argv=["/bin/sh", "-c", "exec sleep 60"]) fds_before = _get_open_fds() with ( @@ -422,22 +443,61 @@ class TestRun: autospec=True, side_effect=LangSDKDagFileProcessorProcess.close, ) as mock_close, - pytest.raises(TimeoutError, match=r"did not parse .*dag\.native within 1s"), ): - LangSDKDagFileProcessorProcess.run( - path=tmp_path / "dag.native", - bundle_path=tmp_path, - bundle_name="testing", - dag_file_rel_path="dag.native", - timeout=1, - logger=structlog.get_logger(), - ) + result = self._run(tmp_path, argv=["/bin/sh", "-c", "exec sleep 60"]) + assert result.import_errors == { + "dag.native": f"The Lang-SDK runtime did not parse {tmp_path / 'dag.native'} within 1.0s" + } [proc] = [c.args[0] for c in mock_close.call_args_list] assert proc._exit_code == -9 assert not proc._open_sockets assert _get_open_fds() <= fds_before + @patch("airflow.settings.get_dagbag_import_timeout", autospec=True, return_value=1) + def test_the_import_timeout_holds_after_the_runtime_exits(self, mock_timeout, tmp_path): + with patch.object( + LangSDKDagFileProcessorProcess, + "close", + autospec=True, + side_effect=LangSDKDagFileProcessorProcess.close, + ) as mock_close: + # The runtime exits, and the process it leaves behind keeps its output open. + result = self._run(tmp_path, argv=["/bin/sh", "-c", "sleep 30 & exit 0"]) + [proc] = [c.args[0] for c in mock_close.call_args_list] + os.killpg(proc.pid, signal.SIGKILL) + + assert result.import_errors == { + "dag.native": f"The Lang-SDK runtime did not parse {tmp_path / 'dag.native'} within 1.0s" + } + assert proc._exit_code == 0 + assert not proc._open_sockets + + @conf_vars({("dag_processor", "dag_file_processor_timeout"): "1"}) + @patch.object( + FakeCoordinator, + "_build_parse_dag_command", + autospec=True, + side_effect=lambda self, *, path: time.sleep(60), + ) + def test_the_dag_file_processor_timeout_applies_until_the_import_timeout_is_reported( + self, mock_build_parse_dag_command, tmp_path + ): + result = self._run(tmp_path) + + assert result.import_errors == { + "dag.native": f"The Lang-SDK runtime did not parse {tmp_path / 'dag.native'} within 1.0s" + } + + [email protected](("configured", "expected"), [(30, 30), (0.5, 0.5), (0, None), (-1, None)]) +@patch("airflow.settings.get_dagbag_import_timeout", autospec=True) +def test_only_a_positive_import_timeout_applies(mock_timeout, configured, expected): + mock_timeout.return_value = configured + + assert _get_import_timeout("/b/dag.native") == expected + mock_timeout.assert_called_once_with("/b/dag.native") + def _make_process(**kwargs) -> LangSDKDagFileProcessorProcess: return LangSDKDagFileProcessorProcess( diff --git a/airflow-core/tests/unit/dag_processing/test_manager.py b/airflow-core/tests/unit/dag_processing/test_manager.py index 2ea33b386a0..b92f93493db 100644 --- a/airflow-core/tests/unit/dag_processing/test_manager.py +++ b/airflow-core/tests/unit/dag_processing/test_manager.py @@ -1682,6 +1682,63 @@ class TestDagFileProcessorManager: assert import_error.stacktrace.startswith("The Lang-SDK runtime sent an invalid frame: ") assert manager.selector.get_map() == {} + @mock.patch("airflow.settings.get_dagbag_import_timeout", autospec=True, return_value=1) + @mock.patch.object(FakeCoordinator, "parse_dag", autospec=True) + @mock.patch.object( + DagFileProcessorManager, "_find_files_in_bundle", autospec=True, return_value=[Path("slow.native")] + ) + def test_a_coordinator_file_past_the_import_timeout_is_an_import_error( + self, mock_find_files, mock_parse_dag, mock_timeout, tmp_path, configure_testing_dag_bundle + ): + mock_parse_dag.side_effect = play_runtime(lambda request, comms: time.sleep(60)) + slow = write_native_file(tmp_path / "slow.native") + + with fake_coordinator(), configure_testing_dag_bundle(tmp_path): + manager = DagFileProcessorManager(max_runs=1, processor_timeout=60) + manager.run() + + with create_session() as session: + [import_error] = session.scalars(select(ParseImportError)).all() + + assert import_error.filename == "slow.native" + assert import_error.stacktrace == f"The Lang-SDK runtime did not parse {slow} within 1.0s" + assert manager.selector.get_map() == {} + + @mock.patch("airflow.settings.get_dagbag_import_timeout", autospec=True) + @mock.patch.object( + DagFileProcessorManager, + "_find_files_in_bundle", + autospec=True, + return_value=[Path("broken.native"), Path("python_dag.py")], + ) + def test_a_failing_import_timeout_policy_fails_only_its_file( + self, mock_find_files, mock_timeout, tmp_path, configure_testing_dag_bundle + ): + def get_dagbag_import_timeout(dag_file_path): + if dag_file_path.endswith(".native"): + raise RuntimeError("policy bug") + return 30 + + mock_timeout.side_effect = get_dagbag_import_timeout + write_native_file(tmp_path / "broken.native") + (tmp_path / "python_dag.py").write_text( + "from airflow.sdk import DAG\nfrom airflow.sdk.bases.operator import BaseOperator\n\n" + 'with DAG("python_dag", schedule=None):\n BaseOperator(task_id="task")\n' + ) + + with fake_coordinator(), configure_testing_dag_bundle(tmp_path): + DagFileProcessorManager(max_runs=1, processor_timeout=60).run() + + with create_session() as session: + dag_ids = session.scalars(select(SerializedDagModel.dag_id)).all() + [import_error] = session.scalars(select(ParseImportError)).all() + + assert dag_ids == ["python_dag"] + assert (import_error.filename, import_error.stacktrace) == ( + "broken.native", + "Cannot start the Lang-SDK runtime: RuntimeError: policy bug", + ) + def test_terminate_orphan_processes_kills_then_closes_processor(self): manager = DagFileProcessorManager(max_runs=1) processor, _ = self.mock_processor()
