Skip to content

Commit 589b468

Browse files
committed
Keep the pipeline descriptor blocking for shell producers
Making the pipe non-blocking regressed shell commands piped to a terminal consumer: do_shell() hands the same open file description to the child, which writes to it directly and failed with EAGAIN once the pipe was full. Leave the descriptor blocking. PipelineWriter still returns to Python regularly, as the job-control relay requires, by polling for room and writing at most PIPE_BUF bytes at a time, which cannot block once the pipe reports it is writable. The job-control test's timeout diagnostics now cope with a sandbox that cannot run ps.
1 parent d3b2630 commit 589b468

3 files changed

Lines changed: 60 additions & 13 deletions

File tree

‎cmd2/utils.py‎

Lines changed: 11 additions & 10 deletions
Original file line numberDiff line numberDiff line change
@@ -836,7 +836,6 @@ class PipelineWriter(io.FileIO):
836836
def __init__(self, fd: int, reader: ProcReader) -> None:
837837
"""Take ownership of a pipe descriptor managed by reader."""
838838
super().__init__(fd, "w")
839-
os.set_blocking(fd, False)
840839
self._reader = reader
841840

842841
def write(self, b: Any) -> int:
@@ -846,25 +845,27 @@ def write(self, b: Any) -> int:
846845
writes, even for an instant, would stop a consumer that had just resumed a
847846
terminal read with SIGTTIN.
848847
849-
A full pipe is awaited in short polls rather than a blocking write. Only the
850-
main thread runs Python signal handlers, and the job-control stop ProcReader
851-
relays may wake another thread, so the main thread has to return to Python
852-
code on its own for the handler to run.
848+
A full pipe is awaited in short polls rather than in one blocking write. Only
849+
the main thread runs Python signal handlers, and the job-control stop ProcReader
850+
relays may wake another thread, so the main thread has to return to Python code
851+
on its own for the handler to run. The descriptor itself stays blocking: a shell
852+
command inherits it, and a producer that found it non-blocking would fail with
853+
EAGAIN once the pipe filled.
853854
"""
854855
import select
855856
import signal
856857

857858
view = memoryview(b).cast("B")
859+
fd = self.fileno()
858860
poller = select.poll()
859-
poller.register(self.fileno(), select.POLLOUT)
861+
poller.register(fd, select.POLLOUT)
860862
try:
861863
with self._reader.lend_terminal():
862864
written = 0
863865
while written < len(view):
864-
try:
865-
written += os.write(self.fileno(), view[written:])
866-
except BlockingIOError:
867-
poller.poll(100)
866+
# Once there is room, a write of at most PIPE_BUF bytes does not block.
867+
if poller.poll(100):
868+
written += os.write(fd, view[written : written + select.PIPE_BUF])
868869
return written
869870
except BrokenPipeError:
870871
# Ctrl-C during a blocking write must cancel the command, even if it

‎tests/test_pipeline_job_control.py‎

Lines changed: 6 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -29,9 +29,12 @@ def describe_processes(root: int, master: int) -> str:
2929
foreground: object = os.tcgetpgrp(master)
3030
except OSError as error:
3131
foreground = error
32-
listing = subprocess.run(
33-
["ps", "-e", "-o", "pid,ppid,pgid,stat,wchan,command"], capture_output=True, text=True, check=False
34-
)
32+
try:
33+
listing = subprocess.run(
34+
["ps", "-e", "-o", "pid,ppid,pgid,stat,wchan,command"], capture_output=True, text=True, check=False
35+
)
36+
except OSError as error:
37+
return f"foreground process group: {foreground}\nno process listing: {error}"
3538
rows = listing.stdout.splitlines()
3639
parents = {}
3740
for row in rows[1:]:

‎tests/test_utils.py‎

Lines changed: 43 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -1,5 +1,6 @@
11
"""Unit testing for cmd2/utils.py module."""
22

3+
import contextlib
34
import errno
45
import math
56
import os
@@ -280,6 +281,48 @@ def test_proc_reader_terminal_group() -> None:
280281
assert reader.terminal_group is None
281282

282283

284+
@pytest.mark.skipif(sys.platform == "win32", reason="POSIX pipes")
285+
@pytest.mark.parametrize("producer", ["writer", "child"])
286+
def test_pipeline_writer_delivers_more_than_the_pipe_holds(producer) -> None:
287+
"""Both cmd2's writes and a child inheriting the descriptor must wait for a slow consumer.
288+
289+
A shell producer gets the descriptor itself, so it must stay blocking: a child that
290+
inherits O_NONBLOCK fails with EAGAIN once the pipe is full.
291+
"""
292+
import subprocess
293+
import threading
294+
295+
payload = b"x" * 4 * 1024 * 1024
296+
read_fd, write_fd = os.pipe()
297+
received = bytearray()
298+
299+
def drain() -> None:
300+
while chunk := os.read(read_fd, 65536):
301+
received.extend(chunk)
302+
time.sleep(0.001)
303+
304+
reader = mock.Mock(lend_terminal=contextlib.nullcontext)
305+
writer = cu.PipelineWriter(write_fd, reader)
306+
consumer = threading.Thread(target=drain)
307+
consumer.start()
308+
try:
309+
if producer == "writer":
310+
assert writer.write(payload) == len(payload)
311+
else:
312+
child = subprocess.run(
313+
[sys.executable, "-c", f"import sys; sys.stdout.buffer.write(b'x' * {len(payload)})"],
314+
stdout=writer.fileno(),
315+
stderr=subprocess.PIPE,
316+
check=False,
317+
)
318+
assert child.returncode == 0, child.stderr.decode()
319+
finally:
320+
writer.close()
321+
consumer.join()
322+
os.close(read_fd)
323+
assert bytes(received) == payload
324+
325+
283326
def test_proc_reader_terminate(pr_none) -> None:
284327
assert pr_none._proc.poll() is None
285328
pr_none.terminate()

0 commit comments

Comments
 (0)