diff --git a/tests/tools/test_process_registry_list_exit.py b/tests/tools/test_process_registry_list_exit.py index 4ef71bc434..b46d8ff9b1 100644 --- a/tests/tools/test_process_registry_list_exit.py +++ b/tests/tools/test_process_registry_list_exit.py @@ -9,10 +9,48 @@ import signal import subprocess import sys import time +from types import SimpleNamespace +from typing import Any, cast import pytest +@pytest.mark.platforms("posix") +def test_list_leaves_live_reader_as_completion_owner(): + from tools.process_registry import ProcessRegistry, ProcessSession + + registry = ProcessRegistry() + session = ProcessSession( + id="proc_owned_reader", + command="owned-reader", + task_id="owner-task", + owner_task_id="owner-owner", + session_key="owner-session", + started_at=time.time(), + notify_on_complete=True, + ) + session.process = cast(Any, SimpleNamespace( + poll=lambda: 0, stdout=None, stderr=None, stdin=None, + )) + session._reader_selectable = True + session._reader_thread = cast(Any, SimpleNamespace(is_alive=lambda: True)) + registry._running[session.id] = session + + listed = registry.list_sessions(session_key="owner-session") + + assert listed[0]["status"] == "exited" + assert session._reader_finish_requested.is_set() + assert session.id in registry._running + assert registry.completion_queue.empty() + + session.append_output("owner-output") + registry._move_to_finished(session) + event = registry.completion_queue.get_nowait() + assert (event["owner_task_id"], event["output"]) == ( + "owner-owner", "owner-output", + ) + + @pytest.mark.platforms("linux") def test_list_reconciles_real_exit_without_consuming_owned_result(tmp_path): # A disposable subreaper owns even the orphaned writer; no global pytest diff --git a/tools/process_registry.py b/tools/process_registry.py index 300a4a0ec9..c30117f1df 100644 --- a/tools/process_registry.py +++ b/tools/process_registry.py @@ -574,6 +574,8 @@ class ProcessSession: _completion_event: threading.Event = field(default_factory=threading.Event, repr=False) _lock: threading.Lock = field(default_factory=threading.Lock) _reader_thread: Optional[threading.Thread] = field(default=None, repr=False) + _reader_finish_requested: threading.Event = field(default_factory=threading.Event, repr=False) + _reader_selectable: bool = field(default=False, repr=False) _pty: Any = field(default=None, repr=False) # ptyprocess handle (use_pty=True) def append_output(self, text: str) -> None: @@ -1382,6 +1384,7 @@ class ProcessRegistry(ProcessCheckpointMixin): fd = None if fd is not None: import select as _select + session._reader_selectable = True idle_after_exit = 0 while True: if fd is not None: @@ -1390,6 +1393,8 @@ class ProcessRegistry(ProcessCheckpointMixin): except (ValueError, OSError): break # fd already closed if not ready: + if session._reader_finish_requested.is_set(): + break # Direct child gone and pipe idle ~200ms: a few more cycles for a # buffered tail, then stop rather than wait forever on an orphaned # grandchild's pipe. @@ -1404,6 +1409,8 @@ class ProcessRegistry(ProcessCheckpointMixin): break # true EOF — all writers closed if chunk: _append_chunk(chunk) + if session._reader_finish_requested.is_set(): + break idle_after_exit = 0 except Exception as e: logger.debug("Process stdout reader ended: %s", e) @@ -1911,6 +1918,26 @@ class ProcessRegistry(ProcessCheckpointMixin): return if rc is None: return # Direct child still running — reader block is legitimate. + reader = session._reader_thread + if ( + not _IS_WINDOWS + and session._reader_selectable + and reader is not None + and reader.is_alive() + ): + # The reader owns the pipe and completion payload. Asking it to + # finish avoids a competing TextIOWrapper read here racing the + # reader, publishing an empty owner-stamped result, then closing + # the pipe before the buffered tail is ingested. It wakes within + # the reader's bounded select interval (or after one final chunk). + session._reader_finish_requested.set() + with session._lock: + session.mark_exited(rc) + logger.info( + "Reconciled session %s: direct child exited with code %s; " + "reader will publish the owned completion after its final drain.", + session.id, rc) + return # Best-effort non-blocking drain of whatever the reader hasn't consumed. stdout = getattr(proc, "stdout", None) if stdout is not None and not _IS_WINDOWS: