From fcce28b5a407e360ae5a356a4728a915fa330c1a Mon Sep 17 00:00:00 2001 From: ethernet Date: Tue, 22 Sep 2026 01:59:23 -0400 Subject: [PATCH] fix(process): preserve reader-owned exit output during list List reconciliation could race the live stdout reader, publish a completion before buffered descendant output was ingested, then close the pipe underneath that reader. For selectable POSIX pipes, mark the direct child exited but ask the reader to perform the final drain and remain the sole completion publisher. Keep the existing fallback for readers that cannot be coordinated. --- .../tools/test_process_registry_list_exit.py | 38 +++++++++++++++++++ tools/process_registry.py | 27 +++++++++++++ 2 files changed, 65 insertions(+) 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: