| 1 | """Subprocess helpers: safe timeout + process-group cleanup. |
| 2 | |
| 3 | Used by bird_x.py (Node.js Bird search) and youtube_yt.py (yt-dlp search |
| 4 | and transcript download). Both need the same os.setsid/killpg cleanup |
| 5 | dance on timeout to avoid orphaning child processes. |
| 6 | """ |
| 7 | |
| 8 | from __future__ import annotations |
| 9 | |
| 10 | import os |
| 11 | import signal |
| 12 | import subprocess |
| 13 | import threading |
| 14 | import time |
| 15 | from dataclasses import dataclass |
| 16 | from typing import Optional, Sequence |
| 17 | |
| 18 | |
| 19 | class SubprocTimeout(Exception): |
| 20 | """Raised when a subprocess exceeds its timeout and is killed.""" |
| 21 | |
| 22 | |
| 23 | # Live run_with_timeout children, process-wide. Each child is a session |
| 24 | # leader (os.setsid), so a group kill aimed at the engine never reaches it; |
| 25 | # the engine's SIGTERM handler and atexit hook drain this set instead. It |
| 26 | # lives here, not in last30days.py, because the entrypoint runs as |
| 27 | # __main__: importing it by name from lib/ executes a second module copy |
| 28 | # with its own empty registry, and on a worker thread that copy cannot |
| 29 | # install its signal handler either. RLock because cleanup_children runs |
| 30 | # inside a signal handler on the main thread, which may already hold the |
| 31 | # lock in register_child_pid. |
| 32 | _child_pids: set[int] = set() |
| 33 | _child_pids_lock = threading.RLock() |
| 34 | # Set once cleanup_children starts. Worker threads keep running while the |
| 35 | # handler sleeps through the grace, so a source can spawn a child after the |
| 36 | # snapshot; register_child_pid kills such late children itself. |
| 37 | _shutting_down = False |
| 38 | |
| 39 | |
| 40 | def register_child_pid(pid: int) -> None: |
| 41 | with _child_pids_lock: |
| 42 | _child_pids.add(pid) |
| 43 | late = _shutting_down |
| 44 | if late: |
| 45 | _kill_child_group(pid, getattr(signal, "SIGKILL", signal.SIGTERM)) |
| 46 | |
| 47 | |
| 48 | def unregister_child_pid(pid: int) -> None: |
| 49 | with _child_pids_lock: |
| 50 | _child_pids.discard(pid) |
| 51 | |
| 52 | |
| 53 | # Upper bound on how long cleanup_children waits for SIGTERMed groups before |
| 54 | # SIGKILLing them. It runs inside the engine's SIGTERM handler, so it must |
| 55 | # finish well inside the MCP server's termGracePeriod |
| 56 | # (mcp/internal/engine/run.go), after which the engine group is SIGKILLed and |
| 57 | # the handler never reaches the escalation. Pinned by tests/test_subproc.py. |
| 58 | CLEANUP_TERM_GRACE_SECONDS = 0.8 |
| 59 | _CLEANUP_POLL_SECONDS = 0.05 |
| 60 | |
| 61 | |
| 62 | def _signal_group(pgid: int, sig: int) -> bool: |
| 63 | """Signal a child's process group; False once the group has no members.""" |
| 64 | try: |
| 65 | os.killpg(pgid, sig) |
| 66 | except (ProcessLookupError, PermissionError, OSError): |
| 67 | return False |
| 68 | return True |
| 69 | |
| 70 | |
| 71 | def _kill_child_group(pid: int, sig: int) -> None: |
| 72 | if hasattr(os, "setsid") and hasattr(os, "killpg"): |
| 73 | _signal_group(pid, sig) |
| 74 | return |
| 75 | try: |
| 76 | os.kill(pid, sig) |
| 77 | except (ProcessLookupError, PermissionError, OSError): |
| 78 | pass |
| 79 | |
| 80 | |
| 81 | def cleanup_children(grace: float = CLEANUP_TERM_GRACE_SECONDS) -> None: |
| 82 | """Terminate the process group of every registered child. |
| 83 | |
| 84 | SIGTERM every group, wait up to ``grace`` seconds for the groups to |
| 85 | empty, then SIGKILL the survivors. run_with_timeout starts each child |
| 86 | with os.setsid, so the child's pid is its pgid; signalling the pgid |
| 87 | directly still reaches grandchildren after the leader has been reaped, |
| 88 | and the kernel does not reuse a pid while a group of that id has |
| 89 | members. A group whose only member is an unreaped zombie still reads as |
| 90 | live, which at worst costs the full grace and a harmless SIGKILL. |
| 91 | Children registered after this call starts are SIGKILLed as they |
| 92 | register, since the caller is about to exit. |
| 93 | """ |
| 94 | global _shutting_down |
| 95 | with _child_pids_lock: |
| 96 | _shutting_down = True |
| 97 | pids = list(_child_pids) |
| 98 | if not pids: |
| 99 | return |
| 100 | if not (hasattr(os, "setsid") and hasattr(os, "killpg")): |
| 101 | for pid in pids: |
| 102 | _kill_child_group(pid, signal.SIGTERM) |
| 103 | return |
| 104 | live = [pgid for pgid in pids if _signal_group(pgid, signal.SIGTERM)] |
| 105 | deadline = time.monotonic() + grace |
| 106 | while live and time.monotonic() < deadline: |
| 107 | time.sleep(_CLEANUP_POLL_SECONDS) |
| 108 | live = [pgid for pgid in live if _signal_group(pgid, 0)] |
| 109 | for pgid in live: |
| 110 | _signal_group(pgid, signal.SIGKILL) |
| 111 | |
| 112 | |
| 113 | @dataclass |
| 114 | class SubprocResult: |
| 115 | """Result of a subprocess run that captured stdout and stderr.""" |
| 116 | |
| 117 | returncode: int |
| 118 | stdout: str |
| 119 | stderr: str |
| 120 | |
| 121 | |
| 122 | def run_with_timeout( |
| 123 | cmd: Sequence[str], |
| 124 | *, |
| 125 | timeout: int, |
| 126 | env: Optional[dict] = None, |
| 127 | on_pid: Optional[callable] = None, |
| 128 | ) -> SubprocResult: |
| 129 | """Run a subprocess with process-group cleanup on timeout. |
| 130 | |
| 131 | Spawns ``cmd`` inside its own process group via ``os.setsid`` where |
| 132 | available. If ``communicate(timeout=...)`` raises ``TimeoutExpired``, |
| 133 | signals ``SIGTERM`` to the entire group, falls back to ``proc.kill()`` |
| 134 | if the signal fails, then waits up to 5 seconds for cleanup, and |
| 135 | raises ``SubprocTimeout``. |
| 136 | |
| 137 | Args: |
| 138 | cmd: Command and arguments to spawn. |
| 139 | timeout: Timeout in seconds passed to ``communicate()``. |
| 140 | env: Optional environment dict. If None, inherits parent env. |
| 141 | on_pid: Optional callable invoked with the child PID right after |
| 142 | spawn. Exceptions raised by the callback are suppressed. The |
| 143 | child is registered for cleanup_children() regardless. |
| 144 | |
| 145 | Returns: |
| 146 | SubprocResult with returncode, stdout, and stderr as strings. |
| 147 | |
| 148 | Raises: |
| 149 | SubprocTimeout: If the process exceeded ``timeout``. |
| 150 | FileNotFoundError: If the executable is not found. |
| 151 | OSError: For other spawn failures. |
| 152 | """ |
| 153 | preexec = os.setsid if hasattr(os, "setsid") else None |
| 154 | |
| 155 | proc = subprocess.Popen( |
| 156 | list(cmd), |
| 157 | stdout=subprocess.PIPE, |
| 158 | stderr=subprocess.PIPE, |
| 159 | text=True, |
| 160 | encoding="utf-8", |
| 161 | errors="replace", |
| 162 | preexec_fn=preexec, |
| 163 | env=env, |
| 164 | ) |
| 165 | |
| 166 | if on_pid is not None: |
| 167 | try: |
| 168 | on_pid(proc.pid) |
| 169 | except Exception: |
| 170 | pass |
| 171 | |
| 172 | register_child_pid(proc.pid) |
| 173 | try: |
| 174 | try: |
| 175 | stdout, stderr = proc.communicate(timeout=timeout) |
| 176 | except subprocess.TimeoutExpired: |
| 177 | term_deadline = time.monotonic() + 5 |
| 178 | pgid = None |
| 179 | try: |
| 180 | if preexec is not None and hasattr(os, "killpg"): |
| 181 | pgid = proc.pid |
| 182 | os.killpg(pgid, signal.SIGTERM) |
| 183 | else: |
| 184 | proc.kill() |
| 185 | except (ProcessLookupError, PermissionError, OSError, AttributeError): |
| 186 | pgid = None |
| 187 | proc.kill() |
| 188 | try: |
| 189 | proc.wait(timeout=5) |
| 190 | except subprocess.TimeoutExpired: |
| 191 | # Child ignored SIGTERM (or our killpg lost the race); escalate. |
| 192 | try: |
| 193 | if pgid is not None: |
| 194 | os.killpg(pgid, signal.SIGKILL) |
| 195 | else: |
| 196 | proc.kill() |
| 197 | except (ProcessLookupError, PermissionError, OSError, AttributeError): |
| 198 | proc.kill() |
| 199 | try: |
| 200 | proc.wait(timeout=5) |
| 201 | except subprocess.TimeoutExpired: |
| 202 | pass # process unkillable (e.g. D-state); leave as zombie |
| 203 | else: |
| 204 | if pgid is not None: |
| 205 | # Reaping the leader does not end its process group. Keep |
| 206 | # the remaining members inside the original TERM grace. |
| 207 | while _signal_group(pgid, 0): |
| 208 | remaining = term_deadline - time.monotonic() |
| 209 | if remaining <= 0: |
| 210 | _signal_group(pgid, signal.SIGKILL) |
| 211 | break |
| 212 | time.sleep(min(_CLEANUP_POLL_SECONDS, remaining)) |
| 213 | raise SubprocTimeout(f"Command {cmd[0]} timed out after {timeout}s") |
| 214 | finally: |
| 215 | unregister_child_pid(proc.pid) |
| 216 | |
| 217 | return SubprocResult( |
| 218 | returncode=proc.returncode, |
| 219 | stdout=stdout or "", |
| 220 | stderr=stderr or "", |
| 221 | ) |
| 222 |