fix(process): reap children after early stdout EOF

This commit is contained in:
fangliquan
2026-09-14 22:13:04 +08:00
committed by Teknium
parent 8bc5894a3a
commit 1bfbeff4e5
2 changed files with 60 additions and 3 deletions

View File

@@ -314,6 +314,57 @@ def test_reader_loop_streams_incremental_chunks_from_read1(registry, monkeypatch
assert moved == ["proc_reader_live"]
def test_reader_waits_past_early_stdout_eof_before_publishing_completion(registry, monkeypatch):
"""Closing stdout is not process completion; the reader must still reap the child."""
class _EarlyEofStdout:
def read(self, _n):
return ""
class _StillRunningProcess:
stdout = _EarlyEofStdout()
def __init__(self):
self.returncode = None
self.wait_timeouts = []
def wait(self, timeout=None):
self.wait_timeouts.append(timeout)
if timeout is not None:
raise subprocess.TimeoutExpired("render", timeout)
self.returncode = 0
return 0
session = _make_session(sid="proc_early_eof")
session.process = _StillRunningProcess()
monkeypatch.setattr(registry, "_move_to_finished", lambda _s: None)
registry._reader_loop(session)
assert session.process.wait_timeouts == [None]
assert session.exited is True
assert session.exit_code == 0
def test_failed_reader_wait_does_not_publish_false_completion(registry, monkeypatch):
"""A failed reap must leave the session running for later reconciliation."""
session = _make_session(sid="proc_wait_failed")
moved = []
monkeypatch.setattr(registry, "_move_to_finished", lambda _s: moved.append(_s.id))
registry._finish_reader(
session,
MagicMock(decode=MagicMock(return_value="")),
lambda _text: None,
"Process",
MagicMock(side_effect=OSError("wait failed")),
lambda: None,
)
assert session.exited is False
assert moved == []
# =========================================================================
# Incremental UTF-8 decoding across chunk boundaries
# (ported from openclaw/openclaw#112325)

View File

@@ -1161,11 +1161,16 @@ class ProcessRegistry(ProcessCheckpointMixin):
finally:
self._finish_reader(
session, decoder, _append_chunk, "Process",
lambda: session.process.wait(timeout=5), lambda: session.process.returncode)
session.process.wait, lambda: session.process.returncode)
def _finish_reader(self, session, decoder, append, label, wait, exit_code) -> None:
"""Reader-thread teardown: flush the decoder (a truncated multibyte tail becomes
one U+FFFD instead of vanishing), reap the child (no zombies), record the exit."""
one U+FFFD instead of vanishing), reap the child (no zombies), record the exit.
A process may close stdout long before it exits. The reader owns a dedicated
daemon thread, so it must keep waiting rather than publish a false completion
and discard the only ``Popen`` handle that can reap the child.
"""
with suppress(Exception):
tail = decoder.decode(b"", final=True)
if tail:
@@ -1173,7 +1178,8 @@ class ProcessRegistry(ProcessCheckpointMixin):
try:
wait()
except Exception as e:
logger.debug("%s wait timed out or failed: %s", label, e)
logger.warning("%s wait failed; leaving process tracked: %s", label, e)
return
self._finish_exited(session, exit_code())
@staticmethod