@@ -552,6 +552,26 @@ def _unblocked_sigttou() -> Iterator[None]:
552552 signal .pthread_sigmask (signal .SIG_SETMASK , previous_mask )
553553
554554
555+ @contextlib .contextmanager
556+ def _session_leader_job_stops () -> Iterator [None ]:
557+ """Spawn a process for a terminal pipeline's job with the same Ctrl-Z behavior as cmd2.
558+
559+ A session leader's job has no outer shell to resume it, so Ctrl-Z must not stop any part
560+ of it. A new process group would otherwise make SIGTSTP actionable again. An ignored
561+ signal stays ignored across exec, so ignore SIGTSTP while spawning.
562+ """
563+ import signal
564+
565+ if os .getpgrp () != os .getsid (0 ):
566+ yield
567+ return
568+ previous = signal .signal (signal .SIGTSTP , signal .SIG_IGN )
569+ try :
570+ yield
571+ finally :
572+ signal .signal (signal .SIGTSTP , previous )
573+
574+
555575class ProcReader :
556576 """Used to capture stdout and stderr from a Popen process if any of those were set to subprocess.PIPE.
557577
@@ -591,10 +611,13 @@ def __init__(
591611 self ._terminal_lock = threading .RLock ()
592612 self ._job_resumed = threading .Event ()
593613 # Set by a thread that relays a pipeline's stop to cmd2's own job. Otherwise Ctrl-Z
594- # reached cmd2 directly, and cmd2 stops the pipeline itself. The watcher skips the
595- # SIGSTOPs cmd2 sends that way, counted under _terminal_lock.
614+ # reached cmd2 directly, and cmd2 stops the pipeline itself.
596615 self ._relaying_stop = False
597- self ._own_stops = 0
616+ # One suspension of the whole job at a time, and how many have finished. A suspension
617+ # continues the whole pipeline as it ends, so a stop reported before it finished is
618+ # already dealt with: see _suspend_with_cmd2().
619+ self ._suspension_lock = threading .Lock ()
620+ self ._suspensions = 0
598621 if terminal_fd is not None :
599622 self ._original_group = os .tcgetpgrp (terminal_fd )
600623
@@ -618,21 +641,22 @@ def send_sigint(self) -> None:
618641 self ._proc .send_signal (signal .CTRL_BREAK_EVENT )
619642 else :
620643 # Since cmd2 uses shell=True in its Popen calls, we need to send the SIGINT to
621- # the whole process group to make sure it propagates further than the shell
622- try :
623- group_id = os .getpgid (self ._proc .pid )
624- except ProcessLookupError :
625- # Pipelines lead their own group. A shell command that joined it, such as
626- # `shell sleep 100 | head -1`, can outlive the reaped consumer. Find the group
627- # through that command, which cmd2 has not reaped yet. Never signal the
628- # consumer's own ID: once its group is gone, the system may reuse it.
629- for producer in self ._joined :
630- if producer .returncode is None :
631- with contextlib .suppress (ProcessLookupError ):
632- group_id = os .getpgid (producer .pid )
633- break
634- else :
635- return
644+ # the whole process group to make sure it propagates further than the shell.
645+ # Once reaped, the process's ID may already belong to another process.
646+ group_id = None
647+ if self ._proc .returncode is None :
648+ with contextlib .suppress (ProcessLookupError ):
649+ group_id = os .getpgid (self ._proc .pid )
650+ # Pipelines lead their own group. A shell command that joined it, such as
651+ # `shell sleep 100 | head -1`, can outlive the reaped consumer. Find the group
652+ # through that command, which cmd2 has not reaped yet. Never signal the
653+ # consumer's own ID: once its group is gone, the system may reuse it.
654+ for producer in self ._joined :
655+ if group_id is None and producer .returncode is None :
656+ with contextlib .suppress (ProcessLookupError ):
657+ group_id = os .getpgid (producer .pid )
658+ if group_id is None :
659+ return
636660 # Never re-signal our own group: other ProcReader callers may share it
637661 # and already received Ctrl-C.
638662 if group_id != os .getpgrp ():
@@ -647,8 +671,10 @@ def terminate(self) -> None:
647671 import signal
648672
649673 # Popen.terminate() polls first, which would compete with our waitpid thread.
650- with contextlib .suppress (ProcessLookupError ):
651- os .kill (self ._proc .pid , signal .SIGTERM )
674+ # Once the watcher has reaped the process, its ID may belong to another one.
675+ if self ._proc .returncode is None :
676+ with contextlib .suppress (ProcessLookupError ):
677+ os .kill (self ._proc .pid , signal .SIGTERM )
652678
653679 @property
654680 def _terminal_group (self ) -> int | None :
@@ -694,23 +720,20 @@ def _manage_terminal(self) -> Iterator[None]:
694720 def suspend_job (signum : int , frame : Any ) -> None :
695721 relayed = self ._relaying_stop
696722 self ._relaying_stop = False
697- stopped_pipeline = False
723+ own_suspension = False
698724 try :
699725 if previous_handler != signal .SIG_DFL :
700726 if callable (previous_handler ):
701727 previous_handler (signum , frame )
702728 return
703729 if os .tcgetpgrp (terminal_fd ) == self ._proc .pid :
704730 self ._set_foreground_group (terminal_fd , self ._original_group )
705- if not relayed and self ._proc .returncode is None :
706- # Ctrl-Z reached only cmd2's group, which owns the terminal between pipe
707- # writes. Stop the pipeline too, as a shell stops its whole job.
708- with self ._terminal_lock :
709- self ._own_stops += 1
710- stopped_pipeline = self ._signal_pipeline (signal .SIGSTOP )
711- if not stopped_pipeline :
712- with self ._terminal_lock :
713- self ._own_stops -= 1
731+ # Ctrl-Z reached only cmd2's group, which owns the terminal between pipe writes.
732+ # Stop the pipeline too, as a shell stops its whole job. Should another thread
733+ # be relaying a stop already, it has stopped the pipeline.
734+ if not relayed and self ._suspension_lock .acquire (blocking = False ):
735+ own_suspension = True
736+ self ._signal_pipeline (signal .SIGSTOP )
714737 # Ignore our group-directed copy, then stop this thread synchronously.
715738 # Wrappers in our job must stop too. Unlike SIGSTOP, SIGTSTP is
716739 # discarded for orphaned groups, which have no shell to resume them.
@@ -726,8 +749,10 @@ def suspend_job(signum: int, frame: Any) -> None:
726749 if self ._terminal_available .is_set () and os .tcgetpgrp (terminal_fd ) == self ._original_group :
727750 self ._set_foreground_group (terminal_fd , self ._proc .pid )
728751 finally :
729- if stopped_pipeline :
752+ if own_suspension :
730753 self ._signal_pipeline (signal .SIGCONT )
754+ self ._suspensions += 1
755+ self ._suspension_lock .release ()
731756 self ._job_resumed .set ()
732757
733758 signal .signal (signal .SIGTSTP , suspend_job )
@@ -820,16 +845,11 @@ def _watch_job(self, terminal_fd: int) -> None:
820845 import signal
821846
822847 while True :
848+ seen = self ._suspensions
823849 _ , status = os .waitpid (self ._proc .pid , os .WUNTRACED )
824850 if not os .WIFSTOPPED (status ):
825851 self ._proc .returncode = os .waitstatus_to_exitcode (status )
826852 return
827- if os .WSTOPSIG (status ) == signal .SIGSTOP :
828- with self ._terminal_lock :
829- if self ._own_stops :
830- # cmd2 stopped the pipeline along with itself, and continues it too.
831- self ._own_stops -= 1
832- continue
833853 if os .WSTOPSIG (status ) in (signal .SIGTTIN , signal .SIGTTOU ):
834854 # Command code owns the terminal between pipe writes. Defer
835855 # consumer terminal access until the next write or final wait.
@@ -838,10 +858,14 @@ def _watch_job(self, terminal_fd: int) -> None:
838858 with self ._terminal_lock :
839859 # A stopped consumer can be killed before another write.
840860 # Keep reaping even while command code owns the terminal.
861+ recent = self ._suspensions
841862 pid , pending_status = os .waitpid (self ._proc .pid , os .WNOHANG | os .WUNTRACED )
842- if pid and not os .WIFSTOPPED (pending_status ):
843- self ._proc .returncode = os .waitstatus_to_exitcode (pending_status )
844- return
863+ if pid :
864+ if not os .WIFSTOPPED (pending_status ):
865+ self ._proc .returncode = os .waitstatus_to_exitcode (pending_status )
866+ return
867+ # A newer stop: the one a suspension would deal with.
868+ seen = recent
845869 # A short write may already have returned the terminal.
846870 # Do not turn that ordinary handoff into a job suspension.
847871 if not self ._terminal_available .is_set ():
@@ -854,37 +878,55 @@ def _watch_job(self, terminal_fd: int) -> None:
854878 if foreground == self ._proc .pid :
855879 continue
856880
857- self ._suspend_with_cmd2 (terminal_fd )
881+ self ._suspend_with_cmd2 (terminal_fd , seen )
882+
883+ def _suspend_with_cmd2 (self , terminal_fd : int , seen : int ) -> None :
884+ """Relay a stop in the pipeline's job to cmd2's own, and continue the pipeline once cmd2 resumes.
885+
886+ Ctrl-Z stops every process of the job that does not ignore it, and each stop may be
887+ reported: the consumer's to its watcher, a producer's to the thread that waits for it.
888+ Only the first to arrive suspends the job. A suspension ends by continuing the whole
889+ pipeline, so a stop reported before it finished needs nothing more.
858890
859- def _suspend_with_cmd2 (self , terminal_fd : int ) -> None :
860- """Relay a pipeline's stop to cmd2's own job, and continue the pipeline once cmd2 resumes."""
891+ :param terminal_fd: the controlling terminal
892+ :param seen: the number of finished suspensions when the stop was reported
893+ """
861894 import signal
862895
863- with self ._terminal_lock :
864- if os .tcgetpgrp (terminal_fd ) == self ._proc .pid :
865- self ._set_foreground_group (terminal_fd , self ._original_group )
866- # Stop every terminal reader before returning control to the outer shell.
867- self ._signal_pipeline (signal .SIGSTOP )
868- self ._job_resumed .clear ()
869- self ._relaying_stop = True
870- # Signal the main thread itself. Only it runs Python signal handlers, and a
871- # process-directed signal may be taken by another thread while the main
872- # thread sleeps in a system call, which then never returns to run the handler.
873- signal .pthread_kill (threading .main_thread ().ident or 0 , signal .SIGTSTP )
874- self ._job_resumed .wait ()
875- self ._signal_pipeline (signal .SIGCONT )
876-
877- def _relay_producer_stop (self ) -> None :
878- """Suspend the shell's whole job for a stopped producer that outlived the consumer.
879-
880- Ctrl-Z reaches only the foreground group, and the producer may be all that is
881- left of it. The watcher ended with the consumer, so nothing else relays the stop.
896+ # Wait in short polls: on the main thread, the suspension being waited for may need
897+ # this thread to run the SIGTSTP handler.
898+ while not self ._suspension_lock .acquire (timeout = 0.1 ):
899+ pass
900+ try :
901+ if self ._suspensions != seen :
902+ return
903+ with self ._terminal_lock :
904+ if os .tcgetpgrp (terminal_fd ) == self ._proc .pid :
905+ self ._set_foreground_group (terminal_fd , self ._original_group )
906+ # Stop every terminal reader before returning control to the outer shell.
907+ self ._signal_pipeline (signal .SIGSTOP )
908+ self ._job_resumed .clear ()
909+ self ._relaying_stop = True
910+ # Signal the main thread itself. Only it runs Python signal handlers, and a
911+ # process-directed signal may be taken by another thread while the main
912+ # thread sleeps in a system call, which then never returns to run the handler.
913+ signal .pthread_kill (threading .main_thread ().ident or 0 , signal .SIGTSTP )
914+ self ._job_resumed .wait ()
915+ self ._signal_pipeline (signal .SIGCONT )
916+ self ._suspensions += 1
917+ finally :
918+ self ._suspension_lock .release ()
919+
920+ def _relay_producer_stop (self , seen : int ) -> None :
921+ """Suspend the shell's whole job for a stopped producer that joined this pipeline.
922+
923+ The consumer need not stop with it: it may ignore Ctrl-Z, or be gone already, and
924+ then nothing else would relay the stop.
925+
926+ :param seen: the number of finished suspensions when the stop was reported
882927 """
883- terminal_fd = self ._terminal_fd
884- if terminal_fd is None or not self ._process_done .is_set ():
885- # A live watcher relays the consumer's stop and continues the whole group.
886- return
887- self ._suspend_with_cmd2 (terminal_fd )
928+ if self ._terminal_fd is not None :
929+ self ._suspend_with_cmd2 (self ._terminal_fd , seen )
888930
889931 def _wait_for_producer (self , pipeline : "ProcReader" , timeout : float | None ) -> None :
890932 """Wait for a producer in a terminal pipeline's job, relaying its job-control stops.
@@ -896,6 +938,7 @@ def _wait_for_producer(self, pipeline: "ProcReader", timeout: float | None) -> N
896938
897939 deadline = None if timeout is None else time .monotonic () + timeout
898940 while self ._proc .returncode is None :
941+ seen = pipeline ._suspensions
899942 try :
900943 pid , status = os .waitpid (self ._proc .pid , os .WNOHANG | os .WUNTRACED )
901944 except ChildProcessError :
@@ -907,7 +950,7 @@ def _wait_for_producer(self, pipeline: "ProcReader", timeout: float | None) -> N
907950 raise subprocess .TimeoutExpired (self ._proc .args , timeout or 0 )
908951 time .sleep (0.05 )
909952 elif os .WIFSTOPPED (status ):
910- pipeline ._relay_producer_stop ()
953+ pipeline ._relay_producer_stop (seen )
911954 else :
912955 self ._proc .returncode = os .waitstatus_to_exitcode (status )
913956
@@ -1069,6 +1112,8 @@ def _relay(self) -> None:
10691112
10701113 poller = select .poll ()
10711114 poller .register (self ._out_fd , select .POLLOUT )
1115+ incoming = select .poll ()
1116+ incoming .register (self ._in_fd , select .POLLIN )
10721117 lend = contextlib .ExitStack ()
10731118 lending = False
10741119 try :
@@ -1079,11 +1124,15 @@ def _relay(self) -> None:
10791124 # No producer is waiting. Let command code have the terminal back.
10801125 lend .close ()
10811126 lending = False
1082- data = os .read (self ._in_fd , 65536 )
1083- if not data :
1084- return
1127+ # Wait for output outside the lock, then take and count it under the lock. Output
1128+ # taken out of the pipe but not yet counted would look passed on to idle() and
1129+ # flush(), and cmd2's next write could overtake it.
1130+ incoming .poll ()
10851131 with self ._lock :
1132+ data = os .read (self ._in_fd , 65536 )
10861133 self ._received += len (data )
1134+ if not data :
1135+ return
10871136 view = memoryview (data )
10881137 written = 0
10891138 while written < len (view ):
0 commit comments