Skip to content

Commit 46a1c1a

Browse files
committed
Fix pager hangs and a startup race in terminal pipelines
Three bugs in the terminal-pipeline job control, found in review and reproduced on a pseudo-terminal with real less: 1. Output written straight to the pipe's descriptor hung the pager. cmd2 lent the pager the terminal only during its own writes through PipelineWriter. A subprocess given self.stdout (for example subprocess.run(..., stdout=self.stdout) in a custom command, or a `!` command run by `run_script s.txt | less`) writes to the descriptor itself, so no lend happened. Once the pipe filled, less stopped with SIGTTIN reading the keyboard, the watcher waited for a lend that never came, and the producer blocked on the full pipe forever. On main these cases worked, because the pager ran in its own session. PipelineWriter.fileno() now hands out the write end of a relay pipe. A thread passes that output on to the consumer and lends it the terminal while a producer is blocked on the relay's full pipe, which is the only time the consumer needs the terminal to make progress. Once no producer is waiting, command code gets the terminal back, as between cmd2's own writes, so a later input() still works. cmd2's own writes first wait for pending relay output, which keeps output in order. When the consumer exits, the relay closes its pipe, so producers still get EPIPE or SIGPIPE. 2. A pipe nested in a piped command took the terminal from the outer pager. A pipeline counted as the terminal's job when either its stdout or its stderr was the foreground terminal. In `run_script s.txt | less` with `big | cat` in the script, the inner pipeline's stdout is the outer pipe, but its stderr is the terminal. It became a terminal job and its lends went to cat's group, so less was never lent the terminal, and the command hung as in (1). Only a pipeline whose stdout is the foreground terminal is now the terminal's job. Others, including nested pipelines and pipelines whose stdout cmd2 captures, run in their own session, as on main. 3. The pager could start before it owned the terminal. The consumer could run between Popen() and cmd2's first terminal lend. A pager such as less sets its terminal modes as it starts, and a background tcsetattr() stops it with SIGTTOU. On macOS the call then fails with EINTR once the process continues, and less carries on with a cooked terminal, so q needs Enter and keys are echoed. This was reproduced by delaying cmd2 after it starts the pipeline. A terminal pipeline now starts as `/bin/sh -c 'read -r _ || exit 1; exec "$SHELL" -c <command>'`. cmd2 sends a single newline down the consumer's stdin pipe only after making the pipeline's group the foreground one. From a pipe, the read builtin takes no more than that line, so no extra descriptor is needed. That matters because dash, /bin/sh on Debian and Ubuntu, cannot redirect descriptors above 9. Checked with sh, bash, dash and zsh. If cmd2 fails before sending the newline, it closes the pipe, and the held pipeline reads EOF and exits rather than waiting forever. Tests: - A PTY test covering the three hang cases in (1) and (2): a subprocess writer, a script's shell command, and a nested pipe. It fails without this change. - The pager-startup PTY test gains a variant that delays cmd2 after it starts the pipeline, which fails without this change. - Unit tests for the relay: lending only while a producer waits, giving the terminal back whether or not output is pending, ordering with cmd2's own writes, and passing the consumer's exit on to producers.
1 parent f4c1d67 commit 46a1c1a

5 files changed

Lines changed: 421 additions & 20 deletions

File tree

‎CHANGELOG.md‎

Lines changed: 2 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -6,8 +6,8 @@
66
does. It had run in a separate session that never received the terminal, so Ctrl-Z and `fg`
77
did not suspend and resume it together with cmd2. The program now owns the terminal as it
88
starts, so a pager can set its terminal modes, and Ctrl-C and Ctrl-Z reach the whole pipeline.
9-
Pipes started from a worker thread, or whose output cmd2 captures, still run in their own
10-
session
9+
Pipes started from a worker thread, or whose output does not go to the terminal, such as one
10+
nested in a command whose own output is piped, still run in their own session
1111
- A `shell` command piped to an interactive program, such as `shell git log | less`, now joins
1212
the pipeline's job, so both processes receive Ctrl-C and Ctrl-Z
1313

‎cmd2/cmd2.py‎

Lines changed: 33 additions & 16 deletions
Original file line numberDiff line numberDiff line change
@@ -3346,22 +3346,40 @@ def _redirect_output(self, statement: Statement) -> utils.RedirectionSavedState:
33463346
pipe_stderr = None if isinstance(sys.stderr, utils.StdSim) else sys.stderr
33473347

33483348
terminal_fd = None
3349+
popen_command = statement.redirect_to
33493350
if sys.platform != "win32":
3350-
# Job control installs signal handlers, which only the main thread may do.
3351-
# Elsewhere, keep the pipeline in its own session as before.
3352-
if threading.current_thread() is threading.main_thread():
3353-
for stream in (pipe_stdout, pipe_stderr):
3354-
if stream is not None and stream.isatty():
3355-
with contextlib.suppress(OSError, ValueError):
3356-
if os.tcgetpgrp(stream.fileno()) == os.getpgrp():
3357-
terminal_fd = stream.fileno()
3358-
break
3351+
# Only a pipeline whose output goes to the terminal is the terminal's job. One
3352+
# nested in a command whose output is piped, for instance, feeds that outer
3353+
# pipeline, whose consumer needs the terminal instead. Job control installs
3354+
# signal handlers, which only the main thread may do. Otherwise, keep the
3355+
# pipeline in its own session as before.
3356+
if threading.current_thread() is threading.main_thread() and pipe_stdout is not None and pipe_stdout.isatty():
3357+
with contextlib.suppress(OSError, ValueError):
3358+
if os.tcgetpgrp(pipe_stdout.fileno()) == os.getpgrp():
3359+
terminal_fd = pipe_stdout.fileno()
33593360
if terminal_fd is None:
33603361
kwargs["start_new_session"] = True
33613362
else:
33623363
kwargs["process_group"] = 0
33633364

3364-
with contextlib.ExitStack() as terminal_stack:
3365+
# A pager such as less sets its terminal modes as it starts, before it reads
3366+
# the pipe. It must own the terminal by then: a background tcsetattr() stops
3367+
# it with SIGTTOU, and on macOS that call fails with EINTR when the process
3368+
# is continued instead of being restarted. less ignores the failure and runs
3369+
# on a cooked terminal. So hold the pipeline in a POSIX sh until cmd2 has
3370+
# made its group the foreground one and sent a newline down the pipe. The
3371+
# read builtin takes no more than that line from a pipe. Then exec the
3372+
# user's shell as before.
3373+
import shlex
3374+
3375+
user_shell = shlex.quote(kwargs.get("executable", "/bin/sh"))
3376+
popen_command = f"read -r _ || exit 1; exec {user_shell} -c {shlex.quote(statement.redirect_to)}"
3377+
kwargs["executable"] = "/bin/sh"
3378+
3379+
with contextlib.ExitStack() as terminal_stack, contextlib.ExitStack() as gate_stack:
3380+
if terminal_fd is not None:
3381+
# Should cmd2 fail before opening the gate, the held pipeline reads EOF and exits.
3382+
gate_stack.callback(new_stdout.close)
33653383
with contextlib.ExitStack() as spawn_stack:
33663384
if terminal_fd is not None and os.getpgrp() == os.getsid(0):
33673385
import signal
@@ -3372,7 +3390,7 @@ def _redirect_output(self, statement: Statement) -> utils.RedirectionSavedState:
33723390
previous_tstp = signal.signal(signal.SIGTSTP, signal.SIG_IGN)
33733391
spawn_stack.callback(signal.signal, signal.SIGTSTP, previous_tstp)
33743392
proc = subprocess.Popen( # noqa: S602
3375-
statement.redirect_to,
3393+
popen_command,
33763394
stdin=subproc_stdin,
33773395
stdout=subprocess.PIPE if pipe_stdout is None else pipe_stdout,
33783396
stderr=subprocess.PIPE if pipe_stderr is None else pipe_stderr,
@@ -3394,12 +3412,11 @@ def _redirect_output(self, statement: Statement) -> utils.RedirectionSavedState:
33943412
if cmd_pipe_proc_reader is None:
33953413
proc.wait(0.2)
33963414
else:
3397-
# A pager such as less sets its terminal modes as it starts, before it
3398-
# reads the pipe. It must own the terminal by then: a background
3399-
# tcsetattr() stops it with SIGTTOU, and on macOS that call fails with
3400-
# EINTR when the process is continued instead of being restarted. less
3401-
# ignores the failure and runs on a cooked terminal.
3415+
# Open the start gate only once the pipeline owns the terminal.
34023416
with cmd_pipe_proc_reader.lend_terminal():
3417+
with contextlib.suppress(OSError):
3418+
os.write(new_stdout.fileno(), b"\n")
3419+
gate_stack.pop_all()
34033420
cmd_pipe_proc_reader.wait_for_exit(0.2)
34043421

34053422
# Check if the pipe process already exited

‎cmd2/utils.py‎

Lines changed: 148 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -921,13 +921,158 @@ def _write_bytes(stream: StdSim | TextIO, to_write: bytes | str) -> None:
921921
stream.buffer.write(to_write)
922922

923923

924+
class _DescriptorRelay:
925+
"""Carry output that subprocesses write to a terminal pipeline's descriptor.
926+
927+
A subprocess given :meth:`PipelineWriter.fileno` writes on its own, so cmd2 cannot lend
928+
the terminal write by write. It writes into this relay's pipe instead, and a thread passes
929+
that output on to the consumer. The consumer is lent the terminal while a producer is
930+
blocked on the relay's full pipe, since only then does it need the terminal to make
931+
progress. Otherwise command code keeps it, as it does between cmd2's own writes.
932+
"""
933+
934+
def __init__(self, out_fd: int, reader: ProcReader) -> None:
935+
"""Start relaying to out_fd, a descriptor the relay takes ownership of.
936+
937+
:param out_fd: the consumer's pipe
938+
:param reader: the terminal pipeline that reads out_fd
939+
"""
940+
self._out_fd = out_fd
941+
self._reader = reader
942+
# write_fd is the descriptor handed to producers. PipelineWriter closes it; once every
943+
# producer has closed its copy too, the relay reads EOF and closes the consumer's pipe.
944+
self._in_fd, self.write_fd = os.pipe()
945+
self._write_fd_open = True
946+
self._lock = threading.Condition()
947+
# Bytes read from the relay's pipe, and bytes passed on to the consumer's pipe
948+
self._received = 0
949+
self._sent = 0
950+
self._done = False
951+
threading.Thread(name="pipe_relay", target=self._relay, daemon=True).start()
952+
953+
def close_write_fd(self) -> None:
954+
"""Close cmd2's copy of the producers' descriptor. The relay finishes once theirs close."""
955+
with self._lock:
956+
if self._write_fd_open:
957+
os.close(self.write_fd)
958+
self._write_fd_open = False
959+
960+
def _unread(self) -> int:
961+
"""Bytes producers have written that the relay has not read yet. Requires _lock."""
962+
import fcntl
963+
import struct
964+
import termios
965+
966+
if self._done:
967+
return 0
968+
return int(struct.unpack("i", fcntl.ioctl(self._in_fd, termios.FIONREAD, b"\0" * 4))[0])
969+
970+
def _producer_blocked(self) -> bool:
971+
"""Whether the relay's pipe is full, which means a producer is waiting on the consumer."""
972+
import select
973+
974+
with self._lock:
975+
if not self._write_fd_open:
976+
# The command is done. ProcReader.wait() lends the terminal from here on.
977+
return False
978+
poller = select.poll()
979+
poller.register(self.write_fd, select.POLLOUT)
980+
return not poller.poll(0)
981+
982+
def flush(self) -> None:
983+
"""Wait until output already written to the relay has reached the consumer's pipe.
984+
985+
cmd2's own writes go straight to the consumer's pipe, so they wait for this first to
986+
keep their order with the output of a producer that has finished. The wait is in short
987+
polls so that the main thread still runs Python signal handlers.
988+
"""
989+
with self._lock:
990+
target = self._received + self._unread()
991+
while not self._done and self._sent < target:
992+
self._lock.wait(0.1)
993+
994+
def _relay(self) -> None:
995+
"""Pass producer output on to the consumer until producers close or the consumer exits."""
996+
import select
997+
998+
poller = select.poll()
999+
poller.register(self._out_fd, select.POLLOUT)
1000+
lend = contextlib.ExitStack()
1001+
lending = False
1002+
try:
1003+
while True:
1004+
with self._lock:
1005+
idle = not self._unread()
1006+
if idle and lending:
1007+
# No producer is waiting. Let command code have the terminal back.
1008+
lend.close()
1009+
lending = False
1010+
data = os.read(self._in_fd, 65536)
1011+
if not data:
1012+
return
1013+
with self._lock:
1014+
self._received += len(data)
1015+
view = memoryview(data)
1016+
written = 0
1017+
while written < len(view):
1018+
# Once there is room, a write of at most PIPE_BUF bytes does not block.
1019+
if poller.poll(100):
1020+
count = os.write(self._out_fd, view[written : written + select.PIPE_BUF])
1021+
written += count
1022+
with self._lock:
1023+
self._sent += count
1024+
self._lock.notify_all()
1025+
elif self._producer_blocked():
1026+
if not lending:
1027+
lend.enter_context(self._reader.lend_terminal())
1028+
lending = True
1029+
elif lending:
1030+
# A producer that stopped writing does not need the consumer to go on.
1031+
lend.close()
1032+
lending = False
1033+
except OSError:
1034+
# The consumer exited. Closing the relay's pipe below passes that on to producers,
1035+
# which get EPIPE or SIGPIPE just as they would writing to the consumer directly.
1036+
return
1037+
finally:
1038+
with contextlib.suppress(OSError):
1039+
lend.close()
1040+
with self._lock:
1041+
self._done = True
1042+
os.close(self._in_fd)
1043+
os.close(self._out_fd)
1044+
self._lock.notify_all()
1045+
1046+
9241047
class PipelineWriter(io.FileIO):
9251048
"""A pipe whose blocking writes temporarily give the consumer terminal access."""
9261049

9271050
def __init__(self, fd: int, reader: ProcReader) -> None:
9281051
"""Take ownership of a pipe descriptor managed by reader."""
9291052
super().__init__(fd, "w")
9301053
self._reader = reader
1054+
self._relay: _DescriptorRelay | None = None
1055+
1056+
def fileno(self) -> int:
1057+
"""Return a descriptor for subprocesses, such as a shell command's stdout.
1058+
1059+
A subprocess writes to it directly, bypassing :meth:`write`. It is the write end of a
1060+
relay (see :class:`_DescriptorRelay`), which lends the consumer the terminal whenever
1061+
such a producer is waiting for the consumer to drain the pipe.
1062+
"""
1063+
if self.closed:
1064+
raise ValueError("I/O operation on closed file")
1065+
if self._relay is None:
1066+
self._relay = _DescriptorRelay(os.dup(super().fileno()), self._reader)
1067+
return self._relay.write_fd
1068+
1069+
def close(self) -> None:
1070+
"""Close the pipe. The consumer sees EOF once every producer has closed its descriptor too."""
1071+
try:
1072+
super().close()
1073+
finally:
1074+
if self._relay is not None:
1075+
self._relay.close_write_fd()
9311076

9321077
def write(self, b: Any) -> int:
9331078
"""Write all of b while the consumer can interact with the terminal.
@@ -947,11 +1092,13 @@ def write(self, b: Any) -> int:
9471092
import signal
9481093

9491094
view = memoryview(b).cast("B")
950-
fd = self.fileno()
1095+
fd = super().fileno()
9511096
poller = select.poll()
9521097
poller.register(fd, select.POLLOUT)
9531098
try:
9541099
with self._reader.lend_terminal():
1100+
if self._relay is not None:
1101+
self._relay.flush()
9551102
written = 0
9561103
while written < len(view):
9571104
# Once there is room, a write of at most PIPE_BUF bytes does not block.

‎tests/test_pipeline_job_control.py‎

Lines changed: 107 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -461,12 +461,16 @@ def wait_until(predicate):
461461
process.wait(timeout=5)
462462

463463

464-
def test_pipeline_pager_can_set_terminal_modes_at_startup(tmp_path) -> None:
464+
@pytest.mark.parametrize("parent_delay", [False, True])
465+
def test_pipeline_pager_can_set_terminal_modes_at_startup(tmp_path, parent_delay) -> None:
465466
"""A pager such as less puts the terminal in raw mode as it starts, before reading its pipe.
466467
467468
It has to own the terminal by then. A background tcsetattr() stops it with SIGTTOU, and
468469
on macOS the call then fails with EINTR once it is continued rather than being restarted.
469470
less ignores that failure, leaving a cooked terminal: q needs Enter and keys are echoed.
471+
472+
However late cmd2 gets to hand the terminal over after starting the pipeline, the pager
473+
must not start before it has.
470474
"""
471475
import pty
472476
import termios
@@ -511,6 +515,13 @@ def test_pipeline_pager_can_set_terminal_modes_at_startup(tmp_path) -> None:
511515
" time.sleep(0.01)\n"
512516
" return startup_wait(reader, timeout)\n"
513517
"utils.ProcReader.wait_for_exit = held_startup_wait\n"
518+
f"if {parent_delay!r}:\n"
519+
# Stall between starting the pipeline and handing it the terminal.
520+
" start_job_control = utils.ProcReader.manage_terminal\n"
521+
" def late_job_control(reader):\n"
522+
" time.sleep(0.5)\n"
523+
" return start_job_control(reader)\n"
524+
" utils.ProcReader.manage_terminal = late_job_control\n"
514525
"app = Cmd()\n"
515526
"app.prompt = 'TEST> '\n"
516527
"app.cmdloop()\n",
@@ -731,3 +742,98 @@ def wait_until(predicate):
731742
os.close(master)
732743
process.kill()
733744
process.wait(timeout=5)
745+
746+
747+
@pytest.mark.parametrize("producer", ["subprocess", "script_shell", "nested_pipe"])
748+
def test_pager_gets_the_terminal_for_output_written_to_the_descriptor(tmp_path, producer) -> None:
749+
"""A producer that writes to the pipe's descriptor itself, bypassing cmd2's writes, cannot trigger a lend per write.
750+
751+
Such as a subprocess given self.stdout, a shell command run by a script, or a pipeline nested
752+
in a piped command. Once the pipe is full, the pager has to be lent the terminal to read the
753+
keys that let it go on. Otherwise it stops with SIGTTIN while the producer waits on it forever.
754+
"""
755+
import pty
756+
757+
shell = shutil.which("bash")
758+
if shell is None:
759+
pytest.skip("requires an interactive bash shell")
760+
big = tmp_path / "big.py"
761+
big.write_text("import sys\nsys.stdout.write('x' * 1048576)\n", encoding="utf-8")
762+
pager = tmp_path / "pager.py"
763+
pager.write_text(
764+
"import os, signal, sys\n"
765+
# Interactive bash leaves TTIN ignored in what it execs, which turns a background read into EIO.
766+
"signal.signal(signal.SIGTTIN, signal.SIG_DFL)\n"
767+
"sys.stdin.buffer.read(4096)\n"
768+
"os.write(2, b'PAGER_READY\\n')\n"
769+
# Like less, read the keyboard from an inherited terminal descriptor. Quit without
770+
# draining the pipe, which the producer must then learn of.
771+
"os.write(2, b'PAGER_GOT ' + os.read(2, 1) + b'\\n')\n",
772+
encoding="utf-8",
773+
)
774+
python = shlex.quote(sys.executable)
775+
script = tmp_path / "script.txt"
776+
script.write_text(
777+
f"!{python} {shlex.quote(str(big))}\n" if producer == "script_shell" else "big | cat\n", encoding="utf-8"
778+
)
779+
application = tmp_path / "application.py"
780+
application.write_text(
781+
"import subprocess, sys\n"
782+
"from cmd2 import Cmd\n"
783+
"class App(Cmd):\n"
784+
" def do_sub(self, _):\n"
785+
f" subprocess.run([sys.executable, {str(big)!r}], stdout=self.stdout, check=False)\n"
786+
" def do_big(self, _):\n"
787+
" self.poutput('y' * 1048576)\n"
788+
"app = App()\n"
789+
"app.prompt = 'TEST> '\n"
790+
"app.cmdloop()\n",
791+
encoding="utf-8",
792+
)
793+
master, slave = pty.openpty()
794+
bootstrap = (
795+
"import os, fcntl, termios; os.setsid(); "
796+
"fcntl.ioctl(0, termios.TIOCSCTTY, 0); "
797+
"os.execv(os.environ['TEST_SHELL'], ['bash', '--noprofile', '--norc', '-i'])"
798+
)
799+
env = dict(os.environ, TERM="xterm-256color", PS1="OUTER> ", TEST_SHELL=shell, SHELL=shell)
800+
env["PYTHONPATH"] = str(Path(__file__).resolve().parents[1])
801+
process = subprocess.Popen([sys.executable, "-c", bootstrap], stdin=slave, stdout=slave, stderr=slave, env=env)
802+
os.close(slave)
803+
decoder = codecs.getincrementaldecoder("utf-8")("replace")
804+
transcript = ""
805+
806+
def wait_until(predicate):
807+
nonlocal transcript
808+
deadline = time.monotonic() + 10
809+
while time.monotonic() < deadline:
810+
if select.select([master], [], [], 0.05)[0]:
811+
data = decoder.decode(os.read(master, 65536))
812+
transcript += data
813+
if "\x1b[6n" in data:
814+
# Answer prompt-toolkit's cursor-position request as a terminal would.
815+
os.write(master, b"\x1b[1;1R")
816+
if predicate():
817+
return
818+
pytest.fail(f"terminal condition timed out:\n{transcript}\n{describe_processes(process.pid, master)}")
819+
820+
command = "sub" if producer == "subprocess" else f"run_script {shlex.quote(str(script))}"
821+
try:
822+
wait_until(lambda: "OUTER> " in transcript)
823+
os.write(master, f"{python} {shlex.quote(str(application))}\n".encode())
824+
wait_until(lambda: "TEST>" in transcript)
825+
os.write(master, f"{command} | {python} {shlex.quote(str(pager))}\n".encode())
826+
wait_until(lambda: "PAGER_READY\r\n" in transcript)
827+
# The pipe fills once the pager has read its first chunk.
828+
time.sleep(0.5)
829+
start = len(transcript)
830+
# The terminal is still canonical: this pager, unlike less, sets no modes.
831+
os.write(master, b"q\n")
832+
wait_until(lambda: "PAGER_GOT q" in transcript[start:])
833+
wait_until(lambda: "TEST>" in transcript[start:])
834+
os.write(master, b"quit\n")
835+
wait_until(lambda: os.tcgetpgrp(master) == process.pid)
836+
finally:
837+
os.close(master)
838+
process.kill()
839+
process.wait(timeout=5)

0 commit comments

Comments
 (0)