Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -8,6 +8,7 @@
import hashlib
import json
import os
import selectors
import subprocess
import time
import weakref
Expand Down Expand Up @@ -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)
Expand Down
44 changes: 44 additions & 0 deletions tests/test_delegation_preview_reuse.py
Original file line number Diff line number Diff line change
Expand Up @@ -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):
Expand Down
Loading