From 04b028b7023a57ededce12b442b1fb91513338b0 Mon Sep 17 00:00:00 2001 From: Venkat Saddala Date: Fri, 10 Apr 2026 12:49:26 +0530 Subject: [PATCH 01/10] Fix DagProcessorJob crash with ValueError: write to closed file on orphan kill After SIGKILL, stale socket registrations on the shared selector cause the next select() callback to write to an already-closed log handle. Add _deregister_processor_sockets() to clean up sockets before closing the handle, and guard target.log() in the supervisor against ValueError. Closes: #64959 --- airflow-core/newsfragments/64959.bugfix.rst | 1 + airflow-core/src/airflow/dag_processing/manager.py | 11 +++++++++++ task-sdk/src/airflow/sdk/execution_time/supervisor.py | 6 +++++- 3 files changed, 17 insertions(+), 1 deletion(-) create mode 100644 airflow-core/newsfragments/64959.bugfix.rst diff --git a/airflow-core/newsfragments/64959.bugfix.rst b/airflow-core/newsfragments/64959.bugfix.rst new file mode 100644 index 0000000000000..886d9df511898 --- /dev/null +++ b/airflow-core/newsfragments/64959.bugfix.rst @@ -0,0 +1 @@ +Fix ``DagProcessorJob`` crash with ``ValueError: write to closed file`` when terminating orphan subprocess processors. diff --git a/airflow-core/src/airflow/dag_processing/manager.py b/airflow-core/src/airflow/dag_processing/manager.py index b2428985a6506..f30cb7787e5fa 100644 --- a/airflow-core/src/airflow/dag_processing/manager.py +++ b/airflow-core/src/airflow/dag_processing/manager.py @@ -1157,6 +1157,15 @@ def remove_orphaned_file_stats(self, present: set[DagFileInfo]): for file in stats_to_remove: del self._file_stats[file] + def _deregister_processor_sockets(self, processor: DagFileProcessorProcess) -> None: + """Unregister and close all sockets for a killed processor before closing its log handle.""" + for sock in list(processor._open_sockets.keys()): + with contextlib.suppress(KeyError): + self.selector.unregister(sock) + with contextlib.suppress(Exception): + sock.close() + processor._open_sockets.clear() + def terminate_orphan_processes(self, present: set[DagFileInfo]): """Stop processors that are working on deleted files.""" present_keys = {file.presence_key for file in present} @@ -1181,6 +1190,7 @@ def terminate_orphan_processes(self, present: set[DagFileInfo]): ), ) processor.kill(signal.SIGKILL) + self._deregister_processor_sockets(processor) processor.logger_filehandle.close() self._file_stats.pop(file, None) @@ -1618,6 +1628,7 @@ def _kill_timed_out_processors(self): # Clean up `self._processors` after iterating over it for proc in processors_to_remove: processor = self._processors.pop(proc) + self._deregister_processor_sockets(processor) processor.logger_filehandle.close() def _add_files_to_queue( diff --git a/task-sdk/src/airflow/sdk/execution_time/supervisor.py b/task-sdk/src/airflow/sdk/execution_time/supervisor.py index 5a5e3852fd735..be824750bf0c1 100644 --- a/task-sdk/src/airflow/sdk/execution_time/supervisor.py +++ b/task-sdk/src/airflow/sdk/execution_time/supervisor.py @@ -2340,7 +2340,11 @@ def process_log_messages_from_subprocess( if level := NAME_TO_LEVEL.get(event.pop("level")): msg = event.pop("event", None) for target in loggers: - target.log(level, msg, **event) + try: + target.log(level, msg, **event) + except ValueError: + # Log handle already closed; discard remaining output from killed subprocess. + return def forward_to_log( From 450748b060b45b4dadaa078f2b5c0c967be2886c Mon Sep 17 00:00:00 2001 From: Venkat Saddala Date: Sat, 18 Apr 2026 14:56:12 +0530 Subject: [PATCH 02/10] Addressing Copilot review: narrow exception handling and add socket tests. --- .../src/airflow/dag_processing/manager.py | 2 +- .../tests/unit/dag_processing/test_manager.py | 46 ++++++++++++++++++ .../airflow/sdk/execution_time/supervisor.py | 7 +-- .../execution_time/test_supervisor.py | 48 +++++++++++++++++++ 4 files changed, 99 insertions(+), 4 deletions(-) diff --git a/airflow-core/src/airflow/dag_processing/manager.py b/airflow-core/src/airflow/dag_processing/manager.py index f30cb7787e5fa..ea24b1628b927 100644 --- a/airflow-core/src/airflow/dag_processing/manager.py +++ b/airflow-core/src/airflow/dag_processing/manager.py @@ -1162,7 +1162,7 @@ def _deregister_processor_sockets(self, processor: DagFileProcessorProcess) -> N for sock in list(processor._open_sockets.keys()): with contextlib.suppress(KeyError): self.selector.unregister(sock) - with contextlib.suppress(Exception): + with contextlib.suppress(OSError, ValueError): sock.close() processor._open_sockets.clear() diff --git a/airflow-core/tests/unit/dag_processing/test_manager.py b/airflow-core/tests/unit/dag_processing/test_manager.py index 4b55b6b2f1ee9..a5fa476b53a2b 100644 --- a/airflow-core/tests/unit/dag_processing/test_manager.py +++ b/airflow-core/tests/unit/dag_processing/test_manager.py @@ -1218,6 +1218,52 @@ def test_terminate_normalizes_file_path_stats_tag(self): ) processor.kill.assert_called_once_with(signal.SIGTERM, escalation_delay=5.0) + def test_terminate_orphan_processes_deregisters_sockets(self): + manager = DagFileProcessorManager(max_runs=1) + processor, _ = self.mock_processor() + + sock1, sock2 = mock.Mock(spec=socket), mock.Mock(spec=socket) + processor._open_sockets[sock1] = "log" + processor._open_sockets[sock2] = "log" + + file_info = DagFileInfo( + bundle_name="testing", rel_path=Path("removed.py"), bundle_path=TEST_DAGS_FOLDER + ) + manager._processors = {file_info: processor} + manager.selector = MagicMock() + + with mock.patch.object(type(processor), "kill"): + manager.terminate_orphan_processes(present=set()) + + manager.selector.unregister.assert_any_call(sock1) + manager.selector.unregister.assert_any_call(sock2) + sock1.close.assert_called_once() + sock2.close.assert_called_once() + assert not processor._open_sockets + processor.logger_filehandle.close.assert_called_once() + + def test_kill_timed_out_processors_deregisters_sockets(self): + manager = DagFileProcessorManager(max_runs=1, processor_timeout=5) + start_time = time.monotonic() - manager.processor_timeout - 1 + processor, _ = self.mock_processor(start_time=start_time) + + sock = mock.Mock(spec=socket) + processor._open_sockets[sock] = "log" + + file_info = DagFileInfo( + bundle_name="testing", rel_path=Path("abc.txt"), bundle_path=TEST_DAGS_FOLDER + ) + manager._processors = {file_info: processor} + manager.selector = MagicMock() + + with mock.patch.object(type(processor), "kill"): + manager._kill_timed_out_processors() + + manager.selector.unregister.assert_called_once_with(sock) + sock.close.assert_called_once() + assert not processor._open_sockets + processor.logger_filehandle.close.assert_called_once() + def test_handle_parsing_result_provides_its_own_session_when_caller_omits(self): """``handle_parsing_result`` is wrapped in ``@provide_session`` so subclasses overriding it can run without a caller-supplied session.""" manager = DagFileProcessorManager(max_runs=1) diff --git a/task-sdk/src/airflow/sdk/execution_time/supervisor.py b/task-sdk/src/airflow/sdk/execution_time/supervisor.py index be824750bf0c1..d9a3f8a582e6c 100644 --- a/task-sdk/src/airflow/sdk/execution_time/supervisor.py +++ b/task-sdk/src/airflow/sdk/execution_time/supervisor.py @@ -2342,9 +2342,10 @@ def process_log_messages_from_subprocess( for target in loggers: try: target.log(level, msg, **event) - except ValueError: - # Log handle already closed; discard remaining output from killed subprocess. - return + except ValueError as e: + if "write to closed file" not in str(e): + raise + # This logger's handle is closed (subprocess was killed); skip it. def forward_to_log( diff --git a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py index 491a1b0e3fd95..a2edf3692184a 100644 --- a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py +++ b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py @@ -3819,6 +3819,54 @@ def test_process_log_messages_from_subprocess(monkeypatch, caplog): ] +def test_process_log_messages_closed_logger_is_skipped(): + """A logger whose file handle is closed is skipped; other loggers still receive messages + and the generator continues processing subsequent lines.""" + from structlog.typing import FilteringBoundLogger + + closed_logger = mock.Mock(spec=FilteringBoundLogger) + closed_logger.log.side_effect = ValueError("write to closed file") + + good_logger = mock.Mock(spec=FilteringBoundLogger) + + def fake_reconfigure(log, *args, **kwargs): + return log + + with mock.patch( + "airflow.sdk.execution_time.supervisor.reconfigure_logger", + side_effect=fake_reconfigure, + ): + gen = process_log_messages_from_subprocess(loggers=(closed_logger, good_logger)) + next(gen) + + # closed_logger raises; good_logger should still receive both messages + gen.send(b'{"level": "info", "event": "hello"}\n') + gen.send(b'{"level": "info", "event": "world"}\n') + + assert good_logger.log.call_count == 2 + + +def test_process_log_messages_unexpected_value_error_is_reraised(): + """A ValueError unrelated to a closed file handle must propagate, not be silently swallowed.""" + from structlog.typing import FilteringBoundLogger + + buggy_logger = mock.Mock(spec=FilteringBoundLogger) + buggy_logger.log.side_effect = ValueError("unexpected formatting bug") + + def fake_reconfigure(log, *args, **kwargs): + return log + + with mock.patch( + "airflow.sdk.execution_time.supervisor.reconfigure_logger", + side_effect=fake_reconfigure, + ): + gen = process_log_messages_from_subprocess(loggers=(buggy_logger,)) + next(gen) + + with pytest.raises(ValueError, match="unexpected formatting bug"): + gen.send(b'{"level": "info", "event": "test"}\n') + + def test_reinit_supervisor_comms(monkeypatch, client_with_ti_start, caplog): def subprocess_main(): # This is run in the subprocess! From cd5a6c82cae4120af635ecc4e9c9a543be416860 Mon Sep 17 00:00:00 2001 From: Venkat Saddala Date: Sat, 18 Apr 2026 15:40:28 +0530 Subject: [PATCH 03/10] Restore suppress tests lost during rebase Re-add test_deregister_processor_sockets_suppresses_already_unregistered and test_deregister_processor_sockets_suppresses_close_error which directly exercise the KeyError and OSError suppression paths in _deregister_processor_sockets. --- .../tests/unit/dag_processing/test_manager.py | 27 +++++++++++++++++++ 1 file changed, 27 insertions(+) diff --git a/airflow-core/tests/unit/dag_processing/test_manager.py b/airflow-core/tests/unit/dag_processing/test_manager.py index a5fa476b53a2b..8a97a7b402a06 100644 --- a/airflow-core/tests/unit/dag_processing/test_manager.py +++ b/airflow-core/tests/unit/dag_processing/test_manager.py @@ -1264,6 +1264,33 @@ def test_kill_timed_out_processors_deregisters_sockets(self): assert not processor._open_sockets processor.logger_filehandle.close.assert_called_once() + def test_deregister_processor_sockets_suppresses_already_unregistered(self): + manager = DagFileProcessorManager(max_runs=1) + processor, _ = self.mock_processor() + + sock = mock.Mock(spec=socket) + processor._open_sockets[sock] = "log" + manager.selector = MagicMock() + manager.selector.unregister.side_effect = KeyError + + manager._deregister_processor_sockets(processor) + + sock.close.assert_called_once() + assert not processor._open_sockets + + def test_deregister_processor_sockets_suppresses_close_error(self): + manager = DagFileProcessorManager(max_runs=1) + processor, _ = self.mock_processor() + + sock = mock.Mock(spec=socket) + sock.close.side_effect = OSError("fd already closed") + processor._open_sockets[sock] = "log" + manager.selector = MagicMock() + + manager._deregister_processor_sockets(processor) + + assert not processor._open_sockets + def test_handle_parsing_result_provides_its_own_session_when_caller_omits(self): """``handle_parsing_result`` is wrapped in ``@provide_session`` so subclasses overriding it can run without a caller-supplied session.""" manager = DagFileProcessorManager(max_runs=1) From 8f8a630c788ee22a41108133e0e018acf30beef8 Mon Sep 17 00:00:00 2001 From: Venkat Saddala Date: Thu, 23 Apr 2026 12:12:00 +0530 Subject: [PATCH 04/10] Fix ruff formatting and rename newsfragment to PR number --- .../{64959.bugfix.rst => 65002.bugfix.rst} | 0 .../src/airflow/dag_processing/manager.py | 37 ++- .../tests/unit/dag_processing/test_manager.py | 281 +++++++++++++++++- 3 files changed, 306 insertions(+), 12 deletions(-) rename airflow-core/newsfragments/{64959.bugfix.rst => 65002.bugfix.rst} (100%) diff --git a/airflow-core/newsfragments/64959.bugfix.rst b/airflow-core/newsfragments/65002.bugfix.rst similarity index 100% rename from airflow-core/newsfragments/64959.bugfix.rst rename to airflow-core/newsfragments/65002.bugfix.rst diff --git a/airflow-core/src/airflow/dag_processing/manager.py b/airflow-core/src/airflow/dag_processing/manager.py index ea24b1628b927..4d5d2f3830b10 100644 --- a/airflow-core/src/airflow/dag_processing/manager.py +++ b/airflow-core/src/airflow/dag_processing/manager.py @@ -1158,12 +1158,39 @@ def remove_orphaned_file_stats(self, present: set[DagFileInfo]): del self._file_stats[file] def _deregister_processor_sockets(self, processor: DagFileProcessorProcess) -> None: - """Unregister and close all sockets for a killed processor before closing its log handle.""" + """ + Drain buffered log data from all sockets of a killed processor, then close them. + + After SIGKILL the OS closes the subprocess write ends; kernel buffers may still hold + unread data. Reading and processing that data here prevents log lines from being lost. + """ for sock in list(processor._open_sockets.keys()): - with contextlib.suppress(KeyError): - self.selector.unregister(sock) - with contextlib.suppress(OSError, ValueError): - sock.close() + try: + key = self.selector.get_key(sock) + except KeyError: + pass + else: + try: + socket_handler, on_close_callback = key.data + sock.setblocking(False) + while True: + try: + if not socket_handler(sock): + break + except (BlockingIOError, InterruptedError, OSError): + break + if on_close_callback is not None: + on_close_callback(sock) + else: + with contextlib.suppress(KeyError): + self.selector.unregister(sock) + except Exception: + self.log.exception("Failed to drain log buffer for killed processor socket.") + with contextlib.suppress(KeyError): + self.selector.unregister(sock) + finally: + with contextlib.suppress(OSError, ValueError): + sock.close() processor._open_sockets.clear() def terminate_orphan_processes(self, present: set[DagFileInfo]): diff --git a/airflow-core/tests/unit/dag_processing/test_manager.py b/airflow-core/tests/unit/dag_processing/test_manager.py index 8a97a7b402a06..19d2c80e431b9 100644 --- a/airflow-core/tests/unit/dag_processing/test_manager.py +++ b/airflow-core/tests/unit/dag_processing/test_manager.py @@ -1156,6 +1156,122 @@ def test_create_process_subprocess_logs_to_stdout( _, kwargs = mock_start.call_args assert kwargs["subprocess_logs_to_stdout"] is expected_subprocess_logs_to_stdout + def test_terminate_orphan_processes_deregisters_sockets_before_closing_log_handle(self): + """Sockets must be removed from the selector BEFORE the log handle is closed. + + This is the exact ordering that caused the original crash: stale socket events + were delivered after logger_filehandle.close(), causing ValueError in structlog. + See: https://github.com/apache/airflow/issues/64959 + """ + manager = DagFileProcessorManager(max_runs=1) + processor, _ = self.mock_processor() + + sock = mock.Mock(spec=socket) + processor._open_sockets[sock] = "log" + + file_info = DagFileInfo( + bundle_name="testing", rel_path=Path("removed.py"), bundle_path=TEST_DAGS_FOLDER + ) + manager._processors = {file_info: processor} + manager.selector = MagicMock() + key, _, _ = self._make_selector_key() + manager.selector.get_key.return_value = key + + call_order: list[str] = [] + original_deregister = manager._deregister_processor_sockets + + def tracking_deregister(proc): + call_order.append("deregister") + original_deregister(proc) + + processor.logger_filehandle.close.side_effect = lambda: call_order.append("close_log") + + with mock.patch.object(type(processor), "kill"): + with mock.patch.object(manager, "_deregister_processor_sockets", side_effect=tracking_deregister): + manager.terminate_orphan_processes(present=set()) + + assert call_order == ["deregister", "close_log"], ( + "Sockets must be deregistered before the log handle is closed to avoid " + "stale selector events writing to a closed BytesLogger" + ) + + def test_kill_timed_out_processors_deregisters_sockets_before_closing_log_handle(self): + """Sockets must be removed from the selector BEFORE the log handle is closed. + + Same ordering requirement as terminate_orphan_processes. + See: https://github.com/apache/airflow/issues/64959 + """ + manager = DagFileProcessorManager(max_runs=1, processor_timeout=5) + start_time = time.monotonic() - manager.processor_timeout - 1 + processor, _ = self.mock_processor(start_time=start_time) + + sock = mock.Mock(spec=socket) + processor._open_sockets[sock] = "log" + + file_info = DagFileInfo(bundle_name="testing", rel_path=Path("abc.txt"), bundle_path=TEST_DAGS_FOLDER) + manager._processors = {file_info: processor} + manager.selector = MagicMock() + key, _, _ = self._make_selector_key() + manager.selector.get_key.return_value = key + + call_order: list[str] = [] + original_deregister = manager._deregister_processor_sockets + + def tracking_deregister(proc): + call_order.append("deregister") + original_deregister(proc) + + processor.logger_filehandle.close.side_effect = lambda: call_order.append("close_log") + + with mock.patch.object(type(processor), "kill"): + with mock.patch.object(manager, "_deregister_processor_sockets", side_effect=tracking_deregister): + manager._kill_timed_out_processors() + + assert call_order == ["deregister", "close_log"], ( + "Sockets must be deregistered before the log handle is closed to avoid " + "stale selector events writing to a closed BytesLogger" + ) + + def test_stale_socket_after_kill_does_not_deliver_to_closed_log_handle(self): + """After a processor is killed, its sockets must not remain in the shared selector. + + Regression test for https://github.com/apache/airflow/issues/64959: + stale socket events from a killed processor were being delivered to a closed + BytesLogger, raising ValueError and crashing the entire DagProcessorJob. + """ + import selectors + + manager = DagFileProcessorManager(max_runs=1) + processor, _ = self.mock_processor() + + # Use a real socketpair so the selector can register them (needs a real fd) + sock, peer = socketpair() + real_selector = selectors.DefaultSelector() + try: + processor._open_sockets[sock] = "log" + + file_info = DagFileInfo( + bundle_name="testing", rel_path=Path("removed.py"), bundle_path=TEST_DAGS_FOLDER + ) + manager._processors = {file_info: processor} + manager.selector = real_selector + + handler = mock.Mock(return_value=False) + + def on_close(s): + real_selector.unregister(s) + + real_selector.register(sock, selectors.EVENT_READ, (handler, on_close)) + + with mock.patch.object(type(processor), "kill"): + manager.terminate_orphan_processes(present=set()) + + with pytest.raises((KeyError, ValueError)): + real_selector.get_key(sock) + finally: + real_selector.close() + peer.close() + def test_kill_timed_out_processors_kill(self): manager = DagFileProcessorManager(max_runs=1, processor_timeout=5) # Set start_time to ensure timeout occurs: start_time = current_time - (timeout + 1) = always (timeout + 1) seconds @@ -1218,6 +1334,17 @@ def test_terminate_normalizes_file_path_stats_tag(self): ) processor.kill.assert_called_once_with(signal.SIGTERM, escalation_delay=5.0) + def _make_selector_key(self, handler_side_effect=None): + """Return a (key, handler, on_close) triple whose key.data is (handler, on_close).""" + if handler_side_effect is not None: + handler = mock.Mock(side_effect=handler_side_effect) + else: + handler = mock.Mock(return_value=False) + on_close = mock.Mock() + key = mock.Mock() + key.data = (handler, on_close) + return key, handler, on_close + def test_terminate_orphan_processes_deregisters_sockets(self): manager = DagFileProcessorManager(max_runs=1) processor, _ = self.mock_processor() @@ -1231,12 +1358,13 @@ def test_terminate_orphan_processes_deregisters_sockets(self): ) manager._processors = {file_info: processor} manager.selector = MagicMock() + key, _, _ = self._make_selector_key() + manager.selector.get_key.return_value = key with mock.patch.object(type(processor), "kill"): manager.terminate_orphan_processes(present=set()) - manager.selector.unregister.assert_any_call(sock1) - manager.selector.unregister.assert_any_call(sock2) + manager.selector.unregister.assert_not_called() sock1.close.assert_called_once() sock2.close.assert_called_once() assert not processor._open_sockets @@ -1250,35 +1378,38 @@ def test_kill_timed_out_processors_deregisters_sockets(self): sock = mock.Mock(spec=socket) processor._open_sockets[sock] = "log" - file_info = DagFileInfo( - bundle_name="testing", rel_path=Path("abc.txt"), bundle_path=TEST_DAGS_FOLDER - ) + file_info = DagFileInfo(bundle_name="testing", rel_path=Path("abc.txt"), bundle_path=TEST_DAGS_FOLDER) manager._processors = {file_info: processor} manager.selector = MagicMock() + key, _, _ = self._make_selector_key() + manager.selector.get_key.return_value = key with mock.patch.object(type(processor), "kill"): manager._kill_timed_out_processors() - manager.selector.unregister.assert_called_once_with(sock) + manager.selector.unregister.assert_not_called() sock.close.assert_called_once() assert not processor._open_sockets processor.logger_filehandle.close.assert_called_once() def test_deregister_processor_sockets_suppresses_already_unregistered(self): + """Socket not in selector (get_key raises KeyError) must not crash; socket is still closed.""" manager = DagFileProcessorManager(max_runs=1) processor, _ = self.mock_processor() sock = mock.Mock(spec=socket) processor._open_sockets[sock] = "log" manager.selector = MagicMock() - manager.selector.unregister.side_effect = KeyError + manager.selector.get_key.side_effect = KeyError manager._deregister_processor_sockets(processor) + manager.selector.unregister.assert_not_called() sock.close.assert_called_once() assert not processor._open_sockets def test_deregister_processor_sockets_suppresses_close_error(self): + """OSError from sock.close() must not crash; _open_sockets is still cleared.""" manager = DagFileProcessorManager(max_runs=1) processor, _ = self.mock_processor() @@ -1286,9 +1417,145 @@ def test_deregister_processor_sockets_suppresses_close_error(self): sock.close.side_effect = OSError("fd already closed") processor._open_sockets[sock] = "log" manager.selector = MagicMock() + key, _, on_close = self._make_selector_key() + manager.selector.get_key.return_value = key + + manager._deregister_processor_sockets(processor) + + on_close.assert_called_once_with(sock) + assert not processor._open_sockets + + def test_deregister_processor_sockets_drains_all_data_before_closing(self): + """Handler is called repeatedly until EOF (returns False) before socket is closed.""" + manager = DagFileProcessorManager(max_runs=1) + processor, _ = self.mock_processor() + + sock = mock.Mock(spec=socket) + processor._open_sockets[sock] = "log" + manager.selector = MagicMock() + key, handler, on_close = self._make_selector_key(handler_side_effect=[True, True, False]) + manager.selector.get_key.return_value = key + + manager._deregister_processor_sockets(processor) + + assert handler.call_count == 3 + on_close.assert_called_once_with(sock) + manager.selector.unregister.assert_not_called() + sock.close.assert_called_once() + assert not processor._open_sockets + + def test_deregister_processor_sockets_stops_drain_on_blocking_io_error(self): + """BlockingIOError (no data left on non-blocking socket) ends drain; on_close and cleanup still run.""" + manager = DagFileProcessorManager(max_runs=1) + processor, _ = self.mock_processor() + + sock = mock.Mock(spec=socket) + processor._open_sockets[sock] = "log" + manager.selector = MagicMock() + key, handler, on_close = self._make_selector_key( + handler_side_effect=[True, BlockingIOError("no data")] + ) + manager.selector.get_key.return_value = key + + manager._deregister_processor_sockets(processor) + + assert handler.call_count == 2 + on_close.assert_called_once_with(sock) + manager.selector.unregister.assert_not_called() + sock.close.assert_called_once() + assert not processor._open_sockets + + def test_deregister_processor_sockets_stops_drain_on_oserror(self): + """OSError during drain (e.g. broken pipe after SIGKILL) ends drain; on_close and cleanup still run.""" + manager = DagFileProcessorManager(max_runs=1) + processor, _ = self.mock_processor() + + sock = mock.Mock(spec=socket) + processor._open_sockets[sock] = "log" + manager.selector = MagicMock() + key, handler, on_close = self._make_selector_key(handler_side_effect=OSError("connection reset")) + manager.selector.get_key.return_value = key manager._deregister_processor_sockets(processor) + handler.assert_called_once_with(sock) + on_close.assert_called_once_with(sock) + manager.selector.unregister.assert_not_called() + sock.close.assert_called_once() + assert not processor._open_sockets + + def test_deregister_processor_sockets_sets_socket_nonblocking_before_drain(self): + """Socket is switched to non-blocking mode before the drain loop.""" + manager = DagFileProcessorManager(max_runs=1) + processor, _ = self.mock_processor() + + sock = mock.Mock(spec=socket) + processor._open_sockets[sock] = "log" + manager.selector = MagicMock() + call_order = [] + sock.setblocking.side_effect = lambda _: call_order.append("setblocking") + key, handler, on_close = self._make_selector_key() + handler.side_effect = lambda _: call_order.append("handler") or False + manager.selector.get_key.return_value = key + + manager._deregister_processor_sockets(processor) + + assert call_order == ["setblocking", "handler"] + sock.setblocking.assert_called_once_with(False) + on_close.assert_called_once_with(sock) + + def test_deregister_processor_sockets_drains_multiple_sockets_independently(self): + """A drain failure on one socket does not prevent the other sockets from being drained.""" + manager = DagFileProcessorManager(max_runs=1) + processor, _ = self.mock_processor() + + sock1, sock2 = mock.Mock(spec=socket), mock.Mock(spec=socket) + processor._open_sockets[sock1] = "log" + processor._open_sockets[sock2] = "log" + manager.selector = MagicMock() + + bad_handler = mock.Mock(side_effect=RuntimeError("unexpected")) + good_handler = mock.Mock(return_value=False) + bad_on_close = mock.Mock() + good_on_close = mock.Mock() + bad_key = mock.Mock() + bad_key.data = (bad_handler, bad_on_close) + good_key = mock.Mock() + good_key.data = (good_handler, good_on_close) + manager.selector.get_key.side_effect = lambda s: bad_key if s is sock1 else good_key + + manager._deregister_processor_sockets(processor) + + sock1.close.assert_called_once() + sock2.close.assert_called_once() + manager.selector.unregister.assert_called_once_with(sock1) + good_handler.assert_called_once_with(sock2) + good_on_close.assert_called_once_with(sock2) + bad_on_close.assert_not_called() + assert not processor._open_sockets + + def test_deregister_processor_sockets_logs_unexpected_exception(self, caplog): + """Unexpected exception during drain is logged (not silently swallowed) and socket is still closed.""" + import logging + + manager = DagFileProcessorManager(max_runs=1) + processor, _ = self.mock_processor() + + sock = mock.Mock(spec=socket) + processor._open_sockets[sock] = "log" + manager.selector = MagicMock() + on_close = mock.Mock() + key = mock.Mock() + key.data = (mock.Mock(side_effect=RuntimeError("unexpected")), on_close) + manager.selector.get_key.return_value = key + + with caplog.at_level(logging.ERROR): + manager._deregister_processor_sockets(processor) + + assert "Failed to drain log buffer for killed processor socket." in caplog.text + on_close.assert_not_called() + manager.selector.unregister.assert_called_once_with(sock) + sock.close.assert_called_once() assert not processor._open_sockets def test_handle_parsing_result_provides_its_own_session_when_caller_omits(self): From 951fc95921ba4f91c451207ae14e87ea0032be1b Mon Sep 17 00:00:00 2001 From: Hemkumar Chheda Date: Fri, 3 Jul 2026 12:28:24 +0530 Subject: [PATCH 05/10] Add missing `contextlib` import in dag_processing/manager.py `_deregister_processor_sockets` uses `contextlib.suppress` in three places but the module never imported `contextlib`, which `ruff` (F821) flags as an undefined name. Caught while rebasing PR #65002 onto current main; the PR's own CI run predates this code path. --- airflow-core/src/airflow/dag_processing/manager.py | 1 + 1 file changed, 1 insertion(+) diff --git a/airflow-core/src/airflow/dag_processing/manager.py b/airflow-core/src/airflow/dag_processing/manager.py index 4d5d2f3830b10..70b82b2087ea8 100644 --- a/airflow-core/src/airflow/dag_processing/manager.py +++ b/airflow-core/src/airflow/dag_processing/manager.py @@ -19,6 +19,7 @@ from __future__ import annotations +import contextlib import functools import gc import inspect From d2a972069b82a18e91991788145892b09169aa32 Mon Sep 17 00:00:00 2001 From: Hemkumar Chheda Date: Tue, 7 Jul 2026 11:10:35 +0530 Subject: [PATCH 06/10] Align orphan-kill regression tests with review guidance --- airflow-core/tests/unit/dag_processing/test_manager.py | 10 +++------- .../tests/task_sdk/execution_time/test_supervisor.py | 5 +---- 2 files changed, 4 insertions(+), 11 deletions(-) diff --git a/airflow-core/tests/unit/dag_processing/test_manager.py b/airflow-core/tests/unit/dag_processing/test_manager.py index 19d2c80e431b9..8353de0ca0bfc 100644 --- a/airflow-core/tests/unit/dag_processing/test_manager.py +++ b/airflow-core/tests/unit/dag_processing/test_manager.py @@ -23,6 +23,7 @@ import os import random import re +import selectors import shutil import signal import textwrap @@ -1161,7 +1162,6 @@ def test_terminate_orphan_processes_deregisters_sockets_before_closing_log_handl This is the exact ordering that caused the original crash: stale socket events were delivered after logger_filehandle.close(), causing ValueError in structlog. - See: https://github.com/apache/airflow/issues/64959 """ manager = DagFileProcessorManager(max_runs=1) processor, _ = self.mock_processor() @@ -1199,7 +1199,6 @@ def test_kill_timed_out_processors_deregisters_sockets_before_closing_log_handle """Sockets must be removed from the selector BEFORE the log handle is closed. Same ordering requirement as terminate_orphan_processes. - See: https://github.com/apache/airflow/issues/64959 """ manager = DagFileProcessorManager(max_runs=1, processor_timeout=5) start_time = time.monotonic() - manager.processor_timeout - 1 @@ -1235,12 +1234,9 @@ def tracking_deregister(proc): def test_stale_socket_after_kill_does_not_deliver_to_closed_log_handle(self): """After a processor is killed, its sockets must not remain in the shared selector. - Regression test for https://github.com/apache/airflow/issues/64959: - stale socket events from a killed processor were being delivered to a closed - BytesLogger, raising ValueError and crashing the entire DagProcessorJob. + Stale socket events from a killed processor must not reach a closed + BytesLogger and crash the DagProcessorJob. """ - import selectors - manager = DagFileProcessorManager(max_runs=1) processor, _ = self.mock_processor() diff --git a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py index a2edf3692184a..a8eec74dcedf8 100644 --- a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py +++ b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py @@ -49,6 +49,7 @@ from opentelemetry.sdk.trace.export.in_memory_span_exporter import InMemorySpanExporter from opentelemetry.trace import get_current_span from pytest_unordered import unordered +from structlog.typing import FilteringBoundLogger from task_sdk import FAKE_BUNDLE, make_client from uuid6 import uuid7 @@ -3822,8 +3823,6 @@ def test_process_log_messages_from_subprocess(monkeypatch, caplog): def test_process_log_messages_closed_logger_is_skipped(): """A logger whose file handle is closed is skipped; other loggers still receive messages and the generator continues processing subsequent lines.""" - from structlog.typing import FilteringBoundLogger - closed_logger = mock.Mock(spec=FilteringBoundLogger) closed_logger.log.side_effect = ValueError("write to closed file") @@ -3848,8 +3847,6 @@ def fake_reconfigure(log, *args, **kwargs): def test_process_log_messages_unexpected_value_error_is_reraised(): """A ValueError unrelated to a closed file handle must propagate, not be silently swallowed.""" - from structlog.typing import FilteringBoundLogger - buggy_logger = mock.Mock(spec=FilteringBoundLogger) buggy_logger.log.side_effect = ValueError("unexpected formatting bug") From 04374a35a2f9c0ae3e8a97750e37cfe439a70359 Mon Sep 17 00:00:00 2001 From: Hemkumar Chheda Date: Tue, 7 Jul 2026 11:48:07 +0530 Subject: [PATCH 07/10] Remove newsfragment from DagProcessorJob fix --- airflow-core/newsfragments/65002.bugfix.rst | 1 - 1 file changed, 1 deletion(-) delete mode 100644 airflow-core/newsfragments/65002.bugfix.rst diff --git a/airflow-core/newsfragments/65002.bugfix.rst b/airflow-core/newsfragments/65002.bugfix.rst deleted file mode 100644 index 886d9df511898..0000000000000 --- a/airflow-core/newsfragments/65002.bugfix.rst +++ /dev/null @@ -1 +0,0 @@ -Fix ``DagProcessorJob`` crash with ``ValueError: write to closed file`` when terminating orphan subprocess processors. From 79980bcfdfb960ae645ec54335f5d2ab9baa2642 Mon Sep 17 00:00:00 2001 From: Hemkumar Chheda Date: Wed, 15 Jul 2026 19:28:18 +0530 Subject: [PATCH 08/10] Fix duplicate dag processor log handle close --- airflow-core/src/airflow/dag_processing/manager.py | 2 -- 1 file changed, 2 deletions(-) diff --git a/airflow-core/src/airflow/dag_processing/manager.py b/airflow-core/src/airflow/dag_processing/manager.py index 19c8ea731a114..dfd486f74460a 100644 --- a/airflow-core/src/airflow/dag_processing/manager.py +++ b/airflow-core/src/airflow/dag_processing/manager.py @@ -1223,7 +1223,6 @@ def terminate_orphan_processes(self, present: set[DagFileInfo]): ) processor.kill(signal.SIGKILL) self._deregister_processor_sockets(processor) - processor.logger_filehandle.close() processor.close() self._file_stats.pop(file, None) @@ -1662,7 +1661,6 @@ def _kill_timed_out_processors(self): for proc in processors_to_remove: processor = self._processors.pop(proc) self._deregister_processor_sockets(processor) - processor.logger_filehandle.close() processor.close() def _add_files_to_queue( From da59ecaea662a51d3e311b02d3e810fdc496eb70 Mon Sep 17 00:00:00 2001 From: Hemkumar Chheda Date: Sun, 2 Aug 2026 23:49:34 +0530 Subject: [PATCH 09/10] Address orphan-kill socket cleanup review feedback --- .../src/airflow/dag_processing/manager.py | 39 -- .../src/airflow/dag_processing/processor.py | 1 + .../tests/unit/dag_processing/test_manager.py | 339 ++---------------- .../sdk/execution_time/schema/schema.json | 2 +- .../schema/versions/__init__.py | 5 + .../schema/versions/v2026_08_02.py | 26 ++ .../airflow/sdk/execution_time/supervisor.py | 59 ++- .../execution_time/test_supervisor.py | 91 ++++- 8 files changed, 201 insertions(+), 361 deletions(-) create mode 100644 task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_08_02.py diff --git a/airflow-core/src/airflow/dag_processing/manager.py b/airflow-core/src/airflow/dag_processing/manager.py index 4b8c7f076efee..3f60cb3f4c4de 100644 --- a/airflow-core/src/airflow/dag_processing/manager.py +++ b/airflow-core/src/airflow/dag_processing/manager.py @@ -19,7 +19,6 @@ from __future__ import annotations -import contextlib import functools import gc import inspect @@ -1164,42 +1163,6 @@ def remove_orphaned_file_stats(self, present: set[DagFileInfo]): for file in stats_to_remove: del self._file_stats[file] - def _deregister_processor_sockets(self, processor: DagFileProcessorProcess) -> None: - """ - Drain buffered log data from all sockets of a killed processor, then close them. - - After SIGKILL the OS closes the subprocess write ends; kernel buffers may still hold - unread data. Reading and processing that data here prevents log lines from being lost. - """ - for sock in list(processor._open_sockets.keys()): - try: - key = self.selector.get_key(sock) - except KeyError: - pass - else: - try: - socket_handler, on_close_callback = key.data - sock.setblocking(False) - while True: - try: - if not socket_handler(sock): - break - except (BlockingIOError, InterruptedError, OSError): - break - if on_close_callback is not None: - on_close_callback(sock) - else: - with contextlib.suppress(KeyError): - self.selector.unregister(sock) - except Exception: - self.log.exception("Failed to drain log buffer for killed processor socket.") - with contextlib.suppress(KeyError): - self.selector.unregister(sock) - finally: - with contextlib.suppress(OSError, ValueError): - sock.close() - processor._open_sockets.clear() - def terminate_orphan_processes(self, present: set[DagFileInfo]): """Stop processors that are working on deleted files.""" present_keys = {file.presence_key for file in present} @@ -1224,7 +1187,6 @@ def terminate_orphan_processes(self, present: set[DagFileInfo]): ), ) processor.kill(signal.SIGKILL) - self._deregister_processor_sockets(processor) processor.close() self._file_stats.pop(file, None) @@ -1662,7 +1624,6 @@ def _kill_timed_out_processors(self): # Clean up `self._processors` after iterating over it for proc in processors_to_remove: processor = self._processors.pop(proc) - self._deregister_processor_sockets(processor) processor.close() def _add_files_to_queue( diff --git a/airflow-core/src/airflow/dag_processing/processor.py b/airflow-core/src/airflow/dag_processing/processor.py index ca7539131f60e..497af0d9509e2 100644 --- a/airflow-core/src/airflow/dag_processing/processor.py +++ b/airflow-core/src/airflow/dag_processing/processor.py @@ -733,6 +733,7 @@ def wait(self) -> int: raise NotImplementedError(f"Don't call wait on {type(self).__name__} objects") def close(self): + self.cleanup_sockets_after_kill() try: self.logger_filehandle.close() except OSError: diff --git a/airflow-core/tests/unit/dag_processing/test_manager.py b/airflow-core/tests/unit/dag_processing/test_manager.py index 25f32099c3c28..139ccf430cb8d 100644 --- a/airflow-core/tests/unit/dag_processing/test_manager.py +++ b/airflow-core/tests/unit/dag_processing/test_manager.py @@ -1188,116 +1188,71 @@ def test_create_process_subprocess_logs_to_stdout( _, kwargs = mock_start.call_args assert kwargs["subprocess_logs_to_stdout"] is expected_subprocess_logs_to_stdout - def test_terminate_orphan_processes_deregisters_sockets_before_closing_log_handle(self): - """Sockets must be removed from the selector BEFORE the log handle is closed. - - This is the exact ordering that caused the original crash: stale socket events - were delivered after logger_filehandle.close(), causing ValueError in structlog. - """ + def test_terminate_orphan_processes_kills_then_closes_processor(self): manager = DagFileProcessorManager(max_runs=1) processor, _ = self.mock_processor() - - sock = mock.Mock(spec=socket) - processor._open_sockets[sock] = "log" - file_info = DagFileInfo( bundle_name="testing", rel_path=Path("removed.py"), bundle_path=TEST_DAGS_FOLDER ) manager._processors = {file_info: processor} - manager.selector = MagicMock() - key, _, _ = self._make_selector_key() - manager.selector.get_key.return_value = key call_order: list[str] = [] - original_deregister = manager._deregister_processor_sockets - - def tracking_deregister(proc): - call_order.append("deregister") - original_deregister(proc) - - processor.logger_filehandle.close.side_effect = lambda: call_order.append("close_log") - - with mock.patch.object(type(processor), "kill"): - with mock.patch.object(manager, "_deregister_processor_sockets", side_effect=tracking_deregister): - manager.terminate_orphan_processes(present=set()) - - assert call_order == ["deregister", "close_log"], ( - "Sockets must be deregistered before the log handle is closed to avoid " - "stale selector events writing to a closed BytesLogger" - ) - - def test_kill_timed_out_processors_deregisters_sockets_before_closing_log_handle(self): - """Sockets must be removed from the selector BEFORE the log handle is closed. + processor.close = mock.Mock(side_effect=lambda: call_order.append("close")) - Same ordering requirement as terminate_orphan_processes. - """ - manager = DagFileProcessorManager(max_runs=1, processor_timeout=5) - start_time = time.monotonic() - manager.processor_timeout - 1 - processor, _ = self.mock_processor(start_time=start_time) - - sock = mock.Mock(spec=socket) - processor._open_sockets[sock] = "log" - - file_info = DagFileInfo(bundle_name="testing", rel_path=Path("abc.txt"), bundle_path=TEST_DAGS_FOLDER) - manager._processors = {file_info: processor} - manager.selector = MagicMock() - key, _, _ = self._make_selector_key() - manager.selector.get_key.return_value = key - - call_order: list[str] = [] - original_deregister = manager._deregister_processor_sockets - - def tracking_deregister(proc): - call_order.append("deregister") - original_deregister(proc) - - processor.logger_filehandle.close.side_effect = lambda: call_order.append("close_log") - - with mock.patch.object(type(processor), "kill"): - with mock.patch.object(manager, "_deregister_processor_sockets", side_effect=tracking_deregister): - manager._kill_timed_out_processors() - - assert call_order == ["deregister", "close_log"], ( - "Sockets must be deregistered before the log handle is closed to avoid " - "stale selector events writing to a closed BytesLogger" - ) + with mock.patch.object( + type(processor), "kill", side_effect=lambda *_args, **_kwargs: call_order.append("kill") + ): + manager.terminate_orphan_processes(present=set()) - def test_stale_socket_after_kill_does_not_deliver_to_closed_log_handle(self): - """After a processor is killed, its sockets must not remain in the shared selector. + assert call_order == ["kill", "close"] - Stale socket events from a killed processor must not reach a closed - BytesLogger and crash the DagProcessorJob. - """ + def test_terminate_orphan_processes_does_not_dispatch_request_frames_after_kill(self): manager = DagFileProcessorManager(max_runs=1) processor, _ = self.mock_processor() - - # Use a real socketpair so the selector can register them (needs a real fd) - sock, peer = socketpair() + request_sock, request_peer = socketpair() real_selector = selectors.DefaultSelector() try: - processor._open_sockets[sock] = "log" + processor.selector = real_selector + processor._open_sockets[request_sock] = "requests" file_info = DagFileInfo( bundle_name="testing", rel_path=Path("removed.py"), bundle_path=TEST_DAGS_FOLDER ) manager._processors = {file_info: processor} - manager.selector = real_selector - handler = mock.Mock(return_value=False) + request_handler = mock.Mock(return_value=False) - def on_close(s): - real_selector.unregister(s) + def on_close(sock): + real_selector.unregister(sock) - real_selector.register(sock, selectors.EVENT_READ, (handler, on_close)) + real_selector.register(request_sock, selectors.EVENT_READ, (request_handler, on_close)) with mock.patch.object(type(processor), "kill"): manager.terminate_orphan_processes(present=set()) + request_handler.assert_not_called() with pytest.raises((KeyError, ValueError)): - real_selector.get_key(sock) + real_selector.get_key(request_sock) finally: real_selector.close() - peer.close() + request_peer.close() + + def test_kill_timed_out_processors_kills_then_closes_processor(self): + manager = DagFileProcessorManager(max_runs=1, processor_timeout=5) + start_time = time.monotonic() - manager.processor_timeout - 1 + processor, _ = self.mock_processor(start_time=start_time) + file_info = DagFileInfo(bundle_name="testing", rel_path=Path("abc.txt"), bundle_path=TEST_DAGS_FOLDER) + manager._processors = {file_info: processor} + + call_order: list[str] = [] + processor.close = mock.Mock(side_effect=lambda: call_order.append("close")) + + with mock.patch.object( + type(processor), "kill", side_effect=lambda *_args, **_kwargs: call_order.append("kill") + ): + manager._kill_timed_out_processors() + + assert call_order == ["kill", "close"] def test_kill_timed_out_processors_kill(self): manager = DagFileProcessorManager(max_runs=1, processor_timeout=5) @@ -1381,230 +1336,6 @@ def test_terminate_normalizes_file_path_stats_tag(self): ) processor.kill.assert_called_once_with(signal.SIGTERM, escalation_delay=5.0) - def _make_selector_key(self, handler_side_effect=None): - """Return a (key, handler, on_close) triple whose key.data is (handler, on_close).""" - if handler_side_effect is not None: - handler = mock.Mock(side_effect=handler_side_effect) - else: - handler = mock.Mock(return_value=False) - on_close = mock.Mock() - key = mock.Mock() - key.data = (handler, on_close) - return key, handler, on_close - - def test_terminate_orphan_processes_deregisters_sockets(self): - manager = DagFileProcessorManager(max_runs=1) - processor, _ = self.mock_processor() - - sock1, sock2 = mock.Mock(spec=socket), mock.Mock(spec=socket) - processor._open_sockets[sock1] = "log" - processor._open_sockets[sock2] = "log" - - file_info = DagFileInfo( - bundle_name="testing", rel_path=Path("removed.py"), bundle_path=TEST_DAGS_FOLDER - ) - manager._processors = {file_info: processor} - manager.selector = MagicMock() - key, _, _ = self._make_selector_key() - manager.selector.get_key.return_value = key - - with mock.patch.object(type(processor), "kill"): - manager.terminate_orphan_processes(present=set()) - - manager.selector.unregister.assert_not_called() - sock1.close.assert_called_once() - sock2.close.assert_called_once() - assert not processor._open_sockets - processor.logger_filehandle.close.assert_called_once() - - def test_kill_timed_out_processors_deregisters_sockets(self): - manager = DagFileProcessorManager(max_runs=1, processor_timeout=5) - start_time = time.monotonic() - manager.processor_timeout - 1 - processor, _ = self.mock_processor(start_time=start_time) - - sock = mock.Mock(spec=socket) - processor._open_sockets[sock] = "log" - - file_info = DagFileInfo(bundle_name="testing", rel_path=Path("abc.txt"), bundle_path=TEST_DAGS_FOLDER) - manager._processors = {file_info: processor} - manager.selector = MagicMock() - key, _, _ = self._make_selector_key() - manager.selector.get_key.return_value = key - - with mock.patch.object(type(processor), "kill"): - manager._kill_timed_out_processors() - - manager.selector.unregister.assert_not_called() - sock.close.assert_called_once() - assert not processor._open_sockets - processor.logger_filehandle.close.assert_called_once() - - def test_deregister_processor_sockets_suppresses_already_unregistered(self): - """Socket not in selector (get_key raises KeyError) must not crash; socket is still closed.""" - manager = DagFileProcessorManager(max_runs=1) - processor, _ = self.mock_processor() - - sock = mock.Mock(spec=socket) - processor._open_sockets[sock] = "log" - manager.selector = MagicMock() - manager.selector.get_key.side_effect = KeyError - - manager._deregister_processor_sockets(processor) - - manager.selector.unregister.assert_not_called() - sock.close.assert_called_once() - assert not processor._open_sockets - - def test_deregister_processor_sockets_suppresses_close_error(self): - """OSError from sock.close() must not crash; _open_sockets is still cleared.""" - manager = DagFileProcessorManager(max_runs=1) - processor, _ = self.mock_processor() - - sock = mock.Mock(spec=socket) - sock.close.side_effect = OSError("fd already closed") - processor._open_sockets[sock] = "log" - manager.selector = MagicMock() - key, _, on_close = self._make_selector_key() - manager.selector.get_key.return_value = key - - manager._deregister_processor_sockets(processor) - - on_close.assert_called_once_with(sock) - assert not processor._open_sockets - - def test_deregister_processor_sockets_drains_all_data_before_closing(self): - """Handler is called repeatedly until EOF (returns False) before socket is closed.""" - manager = DagFileProcessorManager(max_runs=1) - processor, _ = self.mock_processor() - - sock = mock.Mock(spec=socket) - processor._open_sockets[sock] = "log" - manager.selector = MagicMock() - key, handler, on_close = self._make_selector_key(handler_side_effect=[True, True, False]) - manager.selector.get_key.return_value = key - - manager._deregister_processor_sockets(processor) - - assert handler.call_count == 3 - on_close.assert_called_once_with(sock) - manager.selector.unregister.assert_not_called() - sock.close.assert_called_once() - assert not processor._open_sockets - - def test_deregister_processor_sockets_stops_drain_on_blocking_io_error(self): - """BlockingIOError (no data left on non-blocking socket) ends drain; on_close and cleanup still run.""" - manager = DagFileProcessorManager(max_runs=1) - processor, _ = self.mock_processor() - - sock = mock.Mock(spec=socket) - processor._open_sockets[sock] = "log" - manager.selector = MagicMock() - key, handler, on_close = self._make_selector_key( - handler_side_effect=[True, BlockingIOError("no data")] - ) - manager.selector.get_key.return_value = key - - manager._deregister_processor_sockets(processor) - - assert handler.call_count == 2 - on_close.assert_called_once_with(sock) - manager.selector.unregister.assert_not_called() - sock.close.assert_called_once() - assert not processor._open_sockets - - def test_deregister_processor_sockets_stops_drain_on_oserror(self): - """OSError during drain (e.g. broken pipe after SIGKILL) ends drain; on_close and cleanup still run.""" - manager = DagFileProcessorManager(max_runs=1) - processor, _ = self.mock_processor() - - sock = mock.Mock(spec=socket) - processor._open_sockets[sock] = "log" - manager.selector = MagicMock() - key, handler, on_close = self._make_selector_key(handler_side_effect=OSError("connection reset")) - manager.selector.get_key.return_value = key - - manager._deregister_processor_sockets(processor) - - handler.assert_called_once_with(sock) - on_close.assert_called_once_with(sock) - manager.selector.unregister.assert_not_called() - sock.close.assert_called_once() - assert not processor._open_sockets - - def test_deregister_processor_sockets_sets_socket_nonblocking_before_drain(self): - """Socket is switched to non-blocking mode before the drain loop.""" - manager = DagFileProcessorManager(max_runs=1) - processor, _ = self.mock_processor() - - sock = mock.Mock(spec=socket) - processor._open_sockets[sock] = "log" - manager.selector = MagicMock() - call_order = [] - sock.setblocking.side_effect = lambda _: call_order.append("setblocking") - key, handler, on_close = self._make_selector_key() - handler.side_effect = lambda _: call_order.append("handler") or False - manager.selector.get_key.return_value = key - - manager._deregister_processor_sockets(processor) - - assert call_order == ["setblocking", "handler"] - sock.setblocking.assert_called_once_with(False) - on_close.assert_called_once_with(sock) - - def test_deregister_processor_sockets_drains_multiple_sockets_independently(self): - """A drain failure on one socket does not prevent the other sockets from being drained.""" - manager = DagFileProcessorManager(max_runs=1) - processor, _ = self.mock_processor() - - sock1, sock2 = mock.Mock(spec=socket), mock.Mock(spec=socket) - processor._open_sockets[sock1] = "log" - processor._open_sockets[sock2] = "log" - manager.selector = MagicMock() - - bad_handler = mock.Mock(side_effect=RuntimeError("unexpected")) - good_handler = mock.Mock(return_value=False) - bad_on_close = mock.Mock() - good_on_close = mock.Mock() - bad_key = mock.Mock() - bad_key.data = (bad_handler, bad_on_close) - good_key = mock.Mock() - good_key.data = (good_handler, good_on_close) - manager.selector.get_key.side_effect = lambda s: bad_key if s is sock1 else good_key - - manager._deregister_processor_sockets(processor) - - sock1.close.assert_called_once() - sock2.close.assert_called_once() - manager.selector.unregister.assert_called_once_with(sock1) - good_handler.assert_called_once_with(sock2) - good_on_close.assert_called_once_with(sock2) - bad_on_close.assert_not_called() - assert not processor._open_sockets - - def test_deregister_processor_sockets_logs_unexpected_exception(self, caplog): - """Unexpected exception during drain is logged (not silently swallowed) and socket is still closed.""" - import logging - - manager = DagFileProcessorManager(max_runs=1) - processor, _ = self.mock_processor() - - sock = mock.Mock(spec=socket) - processor._open_sockets[sock] = "log" - manager.selector = MagicMock() - on_close = mock.Mock() - key = mock.Mock() - key.data = (mock.Mock(side_effect=RuntimeError("unexpected")), on_close) - manager.selector.get_key.return_value = key - - with caplog.at_level(logging.ERROR): - manager._deregister_processor_sockets(processor) - - assert "Failed to drain log buffer for killed processor socket." in caplog.text - on_close.assert_not_called() - manager.selector.unregister.assert_called_once_with(sock) - sock.close.assert_called_once() - assert not processor._open_sockets - def test_handle_parsing_result_provides_its_own_session_when_caller_omits(self): """``handle_parsing_result`` is wrapped in ``@provide_session`` so subclasses overriding it can run without a caller-supplied session.""" manager = DagFileProcessorManager(max_runs=1) diff --git a/task-sdk/src/airflow/sdk/execution_time/schema/schema.json b/task-sdk/src/airflow/sdk/execution_time/schema/schema.json index 01e3c8970b5cb..b65201a2333bf 100644 --- a/task-sdk/src/airflow/sdk/execution_time/schema/schema.json +++ b/task-sdk/src/airflow/sdk/execution_time/schema/schema.json @@ -1,6 +1,6 @@ { "$schema": "https://json-schema.org/draft/2020-12/schema", - "api_version": "2026-06-16", + "api_version": "2026-08-02", "description": "Apache Airflow SDK Supervisor Schema", "$defs": { "AssetAliasReferenceAssetEventDagRun": { diff --git a/task-sdk/src/airflow/sdk/execution_time/schema/versions/__init__.py b/task-sdk/src/airflow/sdk/execution_time/schema/versions/__init__.py index 9491a8993fdc3..ad1c4df4f6be6 100644 --- a/task-sdk/src/airflow/sdk/execution_time/schema/versions/__init__.py +++ b/task-sdk/src/airflow/sdk/execution_time/schema/versions/__init__.py @@ -37,8 +37,13 @@ def get_bundle() -> VersionBundle: """ from cadwyn import HeadVersion, Version, VersionBundle + from airflow.sdk.execution_time.schema.versions.v2026_08_02 import ( + TrackDagProcessorSocketCleanupOwnership, + ) + return VersionBundle( HeadVersion(), + Version("2026-08-02", TrackDagProcessorSocketCleanupOwnership), Version("2026-06-16"), ) diff --git a/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_08_02.py b/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_08_02.py new file mode 100644 index 0000000000000..c54321b77f2a3 --- /dev/null +++ b/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_08_02.py @@ -0,0 +1,26 @@ +# 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 + +from cadwyn import VersionChange + + +class TrackDagProcessorSocketCleanupOwnership(VersionChange): + """Track the parser supervisor contract snapshot used by orphan-kill cleanup.""" + + description = __doc__ + instructions_to_migrate_to_previous_version = () diff --git a/task-sdk/src/airflow/sdk/execution_time/supervisor.py b/task-sdk/src/airflow/sdk/execution_time/supervisor.py index dce16fd243611..908f66dcec52c 100644 --- a/task-sdk/src/airflow/sdk/execution_time/supervisor.py +++ b/task-sdk/src/airflow/sdk/execution_time/supervisor.py @@ -1007,6 +1007,46 @@ def _cleanup_open_sockets(self): self.selector.close() self.stdin.close() + def cleanup_sockets_after_kill(self) -> None: + """Drain log-bearing sockets, then close every remaining socket after a forced kill.""" + for sock, socket_type in list(self._open_sockets.items()): + try: + key = self.selector.get_key(sock) + except KeyError: + key = None + + if key is not None: + socket_handler, on_close = key.data + try: + if socket_type != "requests": + sock.setblocking(False) + while True: + try: + if not socket_handler(sock): + break + except (BlockingIOError, InterruptedError, OSError): + break + + if on_close is not None: + on_close(sock) + else: + with suppress(KeyError): + self.selector.unregister(sock) + self._open_sockets.pop(sock, None) + except Exception: + log.exception( + "Failed to clean up killed subprocess socket", + pid=self.pid, + socket_type=socket_type, + ) + with suppress(KeyError): + self.selector.unregister(sock) + self._open_sockets.pop(sock, None) + with suppress(OSError, ValueError): + sock.close() + + self._open_sockets.clear() + def kill( self, signal_to_send: signal.Signals = signal.SIGINT, @@ -2350,12 +2390,17 @@ def process_log_messages_from_subprocess( if level := NAME_TO_LEVEL.get(event.pop("level")): msg = event.pop("event", None) for target in loggers: - try: - target.log(level, msg, **event) - except ValueError as e: - if "write to closed file" not in str(e): - raise - # This logger's handle is closed (subprocess was killed); skip it. + _log_to_target(target, level, msg, **event) + + +def _log_to_target(target: FilteringBoundLogger, level: int, msg: str | None, **event) -> None: + rendered_msg = msg if msg is not None else "" + try: + target.log(level, rendered_msg, **event) + except ValueError as e: + if "closed file" not in str(e): + raise + log.debug("Dropped log line for closed logger handle", level=level, logger=event.get("logger")) def forward_to_log( @@ -2370,7 +2415,7 @@ def forward_to_log( except UnicodeDecodeError: msg = line.decode("ascii", errors="replace") for log in target_loggers: - log.log(level, msg, logger=logger) + _log_to_target(log, level, msg, logger=logger) def ensure_secrets_backend_loaded() -> list[BaseSecretsBackend]: diff --git a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py index 221cc62b873df..4be9e3b2e9b38 100644 --- a/task-sdk/tests/task_sdk/execution_time/test_supervisor.py +++ b/task-sdk/tests/task_sdk/execution_time/test_supervisor.py @@ -167,6 +167,7 @@ ProcessTracker, _make_process_nondumpable, _remote_logging_conn, + forward_to_log, in_process_api_server, make_buffered_socket_reader, process_log_messages_from_subprocess, @@ -3947,29 +3948,51 @@ def test_process_log_messages_from_subprocess(monkeypatch, caplog): ] -def test_process_log_messages_closed_logger_is_skipped(): - """A logger whose file handle is closed is skipped; other loggers still receive messages - and the generator continues processing subsequent lines.""" +@pytest.mark.parametrize( + "error_message", + ["write to closed file", "I/O operation on closed file"], +) +def test_process_log_messages_closed_logger_is_skipped(error_message): closed_logger = mock.Mock(spec=FilteringBoundLogger) - closed_logger.log.side_effect = ValueError("write to closed file") + closed_logger.log.side_effect = ValueError(error_message) good_logger = mock.Mock(spec=FilteringBoundLogger) - def fake_reconfigure(log, *args, **kwargs): - return log + def fake_reconfigure(logger, *args, **kwargs): + return logger - with mock.patch( - "airflow.sdk.execution_time.supervisor.reconfigure_logger", - side_effect=fake_reconfigure, + with ( + mock.patch( + "airflow.sdk.execution_time.supervisor.reconfigure_logger", + side_effect=fake_reconfigure, + ), + mock.patch.object(supervisor.log, "debug") as mock_debug, ): gen = process_log_messages_from_subprocess(loggers=(closed_logger, good_logger)) next(gen) - # closed_logger raises; good_logger should still receive both messages gen.send(b'{"level": "info", "event": "hello"}\n') gen.send(b'{"level": "info", "event": "world"}\n') assert good_logger.log.call_count == 2 + assert mock_debug.call_count == 2 + + +def test_forward_to_log_closed_logger_is_skipped(): + closed_logger = mock.Mock(spec=FilteringBoundLogger) + closed_logger.log.side_effect = ValueError("I/O operation on closed file") + good_logger = mock.Mock(spec=FilteringBoundLogger) + + with mock.patch.object(supervisor.log, "debug") as mock_debug: + gen = forward_to_log((closed_logger, good_logger), logger="task.stdout", level=logging.INFO) + next(gen) + gen.send(b"hello\n") + gen.send(b"world\n") + + assert good_logger.log.call_count == 2 + good_logger.log.assert_any_call(logging.INFO, "hello", logger="task.stdout") + good_logger.log.assert_any_call(logging.INFO, "world", logger="task.stdout") + assert mock_debug.call_count == 2 def test_process_log_messages_unexpected_value_error_is_reraised(): @@ -3991,6 +4014,54 @@ def fake_reconfigure(log, *args, **kwargs): gen.send(b'{"level": "info", "event": "test"}\n') +def test_cleanup_sockets_after_kill_drains_logs_but_not_requests(mocker): + request_read, request_write = socket.socketpair() + stdout_read, stdout_write = socket.socketpair() + log_read, log_write = socket.socketpair() + + subprocess = ActivitySubprocess( + process_log=mocker.MagicMock(), + id=TI_ID, + pid=12345, + stdin=stdout_write, + client=mocker.Mock(), + process=mocker.Mock(), + ) + selector = selectors.DefaultSelector() + subprocess.selector = selector + + request_handler = mock.Mock(return_value=False) + stdout_handler = mock.Mock(return_value=False) + log_handler = mock.Mock(return_value=False) + + def on_close(sock): + selector.unregister(sock) + subprocess._open_sockets.pop(sock, None) + + try: + subprocess._open_sockets[request_read] = "requests" + subprocess._open_sockets[stdout_read] = "stdout" + subprocess._open_sockets[log_read] = "logs" + + selector.register(request_read, selectors.EVENT_READ, (request_handler, on_close)) + selector.register(stdout_read, selectors.EVENT_READ, (stdout_handler, on_close)) + selector.register(log_read, selectors.EVENT_READ, (log_handler, on_close)) + + subprocess.cleanup_sockets_after_kill() + + request_handler.assert_not_called() + stdout_handler.assert_called_once_with(stdout_read) + log_handler.assert_called_once_with(log_read) + assert not subprocess._open_sockets + with pytest.raises((KeyError, ValueError)): + selector.get_key(request_read) + finally: + selector.close() + request_write.close() + stdout_write.close() + log_write.close() + + def test_reinit_supervisor_comms(monkeypatch, client_with_ti_start, caplog): def subprocess_main(): # This is run in the subprocess! From 1d79e80579cb88ea991c4e07a2703522d663eb8b Mon Sep 17 00:00:00 2001 From: Hemkumar Chheda Date: Tue, 4 Aug 2026 22:41:03 +0530 Subject: [PATCH 10/10] Remove no-op supervisor schema version bump --- .../sdk/execution_time/schema/schema.json | 2 +- .../schema/versions/__init__.py | 5 ---- .../schema/versions/v2026_08_02.py | 26 ------------------- 3 files changed, 1 insertion(+), 32 deletions(-) delete mode 100644 task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_08_02.py diff --git a/task-sdk/src/airflow/sdk/execution_time/schema/schema.json b/task-sdk/src/airflow/sdk/execution_time/schema/schema.json index b65201a2333bf..01e3c8970b5cb 100644 --- a/task-sdk/src/airflow/sdk/execution_time/schema/schema.json +++ b/task-sdk/src/airflow/sdk/execution_time/schema/schema.json @@ -1,6 +1,6 @@ { "$schema": "https://json-schema.org/draft/2020-12/schema", - "api_version": "2026-08-02", + "api_version": "2026-06-16", "description": "Apache Airflow SDK Supervisor Schema", "$defs": { "AssetAliasReferenceAssetEventDagRun": { diff --git a/task-sdk/src/airflow/sdk/execution_time/schema/versions/__init__.py b/task-sdk/src/airflow/sdk/execution_time/schema/versions/__init__.py index ad1c4df4f6be6..9491a8993fdc3 100644 --- a/task-sdk/src/airflow/sdk/execution_time/schema/versions/__init__.py +++ b/task-sdk/src/airflow/sdk/execution_time/schema/versions/__init__.py @@ -37,13 +37,8 @@ def get_bundle() -> VersionBundle: """ from cadwyn import HeadVersion, Version, VersionBundle - from airflow.sdk.execution_time.schema.versions.v2026_08_02 import ( - TrackDagProcessorSocketCleanupOwnership, - ) - return VersionBundle( HeadVersion(), - Version("2026-08-02", TrackDagProcessorSocketCleanupOwnership), Version("2026-06-16"), ) diff --git a/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_08_02.py b/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_08_02.py deleted file mode 100644 index c54321b77f2a3..0000000000000 --- a/task-sdk/src/airflow/sdk/execution_time/schema/versions/v2026_08_02.py +++ /dev/null @@ -1,26 +0,0 @@ -# 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 - -from cadwyn import VersionChange - - -class TrackDagProcessorSocketCleanupOwnership(VersionChange): - """Track the parser supervisor contract snapshot used by orphan-kill cleanup.""" - - description = __doc__ - instructions_to_migrate_to_previous_version = ()