From a20a9cb28c77b8ec201aa9f74aa2b9d8efe6ad7c Mon Sep 17 00:00:00 2001 From: huangruiteng Date: Sun, 4 Oct 2026 12:20:15 +0800 Subject: [PATCH] fix: stop POSIX preview input at its parent deadline Signed-off-by: huangruiteng --- .../delegation_preview_transport.py | 29 +++++++++++- tests/test_delegation_preview_reuse.py | 44 +++++++++++++++++++ 2 files changed, 72 insertions(+), 1 deletion(-) diff --git a/loopx/control_plane/collaboration/delegation_preview_transport.py b/loopx/control_plane/collaboration/delegation_preview_transport.py index bfb7da4475..85c41ce47e 100644 --- a/loopx/control_plane/collaboration/delegation_preview_transport.py +++ b/loopx/control_plane/collaboration/delegation_preview_transport.py @@ -8,6 +8,7 @@ import hashlib import json import os +import selectors import subprocess import time import weakref @@ -163,9 +164,35 @@ def preview(self, *, command: list[str], workspace: Path, release: Path, def _send(self, value: dict, deadline: float, timeout: float) -> None: assert self._process is not None and self._process.stdin is not None stream = self._process.stdin - result: Queue = Queue(maxsize=1) data = (json.dumps(value) + "\n").encode("utf-8") + if os.name == "posix": + # A blocking write in another thread can retain the pipe after + # close(), or finish a frame after the caller's deadline. This + # transport exclusively owns stdin; keep its writes nonblocking + # so the original owner can observe EOF during timeout cleanup. + descriptor = stream.fileno() + os.set_blocking(descriptor, False) + remaining = memoryview(data) + with selectors.DefaultSelector() as writable: + writable.register(descriptor, selectors.EVENT_WRITE) + while remaining: + budget = deadline - time.monotonic() + if budget <= 0 or not writable.select(budget): + raise subprocess.TimeoutExpired(["delegation-preview"], timeout) + try: + written = os.write(descriptor, remaining) + except BlockingIOError: + continue + if written <= 0: + raise OSError("preview input closed") + remaining = remaining[written:] + return + + # Keep the existing Windows pipe writer path. Host cleanup remains + # owned by the same supervisor on every platform. + result: Queue = Queue(maxsize=1) + def write() -> None: try: remaining = memoryview(data) diff --git a/tests/test_delegation_preview_reuse.py b/tests/test_delegation_preview_reuse.py index 6543ea67e8..e29a94fe13 100644 --- a/tests/test_delegation_preview_reuse.py +++ b/tests/test_delegation_preview_reuse.py @@ -421,6 +421,50 @@ def measured_send(*args, **kwargs): transport.close() +@pytest.mark.skipif(os.name != "posix", reason="POSIX anonymous-pipe deadline") +@pytest.mark.parametrize("size", [65536, 524288]) +def test_backpressured_send_cannot_finish_a_frame_after_deadline(tmp_path, size): + from loopx.control_plane.collaboration.delegation_preview_transport import DelegationPreviewTransport + + trigger = tmp_path / "resume" + # Resume the actual pipe reader only after _send has reported a timeout. + # A surviving writer must not finish the request when capacity returns. + reader = f"""import json,os,select,sys,time +from pathlib import Path +trigger=Path({str(trigger)!r}) +print('ready',flush=True) +while not trigger.exists():time.sleep(.005) +fd=sys.stdin.fileno();os.set_blocking(fd,False);chunks=[] +while select.select([fd],[],[],.2)[0]: + data=os.read(fd,65536) + if not data:break + chunks.append(data) +data=b''.join(chunks) +print(json.dumps({{'bytes':len(data),'complete_frame':data.endswith(b'\\n')}}),flush=True) +""" + process = subprocess.Popen([sys.executable, "-c", reader], stdin=subprocess.PIPE, + stdout=subprocess.PIPE, text=True) + transport = DelegationPreviewTransport() + transport._process = process + value = {"kind": "request", "id": 1, "argv": ["x" * size], "timeout_ms": 100} + try: + assert process.stdout.readline().strip() == "ready" + started = time.monotonic() + with pytest.raises(subprocess.TimeoutExpired): + transport._send(value, started + 0.1, 0.1) + assert time.monotonic() - started < 0.4 + trigger.touch() + process.wait(timeout=3) + result = json.loads(process.stdout.readline()) + assert result["bytes"] < len((json.dumps(value) + "\n").encode()) + assert result["complete_frame"] is False + finally: + if process.poll() is None: + process.kill() + process.wait(timeout=3) + transport.close() + + @pytest.mark.skipif(sys.platform == "win32", reason="POSIX retirement cleanup fence") @pytest.mark.parametrize("retirement", ["idle", "lifetime", "broken_pipe"]) def test_unaccepted_request_recovers_only_after_owned_retirement(tmp_path, monkeypatch, retirement):