diff --git a/core/platform_compat.py b/core/platform_compat.py index 0db577f3f..cb135f51a 100644 --- a/core/platform_compat.py +++ b/core/platform_compat.py @@ -133,22 +133,14 @@ def pid_alive(pid: Optional[int]) -> bool: def kill_process_tree(pid: Optional[int], *, start_token=None, pgid=None, require_identity=False): """Use the runtime's shared escalating teardown and return verified death. - Callers retaining durable PIDs must validate their recorded identity before - calling this compatibility entry point. Native grants retain identity at - spawn and use containment.release directly. + Callers retaining durable PIDs must pass their recorded ``start_token`` + with ``require_identity=True``. Native grants retain identity at spawn and + use containment.release directly; this entry point owns no grant record. """ - from src import containment - if not pid or int(pid) <= 0: - return containment.ReleaseOutcome(dead=True, escalated=False) - spec = containment.ContainmentSpec(workspace=os.getcwd(), env={}, wall_clock_s=1, - required=frozenset()) - grant = containment.ContainmentGrant( - id="", mechanism="windows_tree" if IS_WINDOWS else "process_group", - workspace=spec.workspace, enforced=frozenset(), degraded=(), - unenforced_required=(), owner="compatibility", mode=containment.CONTAINMENT_MODE, - spec=spec, pid=int(pid), pgid=pgid or containment._pgid_of(int(pid)), + from src import process_lifecycle + return process_lifecycle.terminate_tree( + pid, pgid=pgid, start_token=start_token, require_identity=require_identity, ) - return containment.release(grant, start_token=start_token, require_identity=require_identity) # ── Shell / executable resolution ─────────────────────────────────────────── diff --git a/routes/shell_routes.py b/routes/shell_routes.py index 73e055c18..6a1c0f583 100644 --- a/routes/shell_routes.py +++ b/routes/shell_routes.py @@ -26,6 +26,7 @@ from src.host_docker_access import ( ) from src.optional_deps import prepare_optional_dependency_import from src.auth_helpers import _auth_disabled +from src import process_lifecycle # POSIX-only: `pty`/`fcntl` transitively import `termios`, which does NOT exist # on Windows, so importing them unconditionally crashed app startup there @@ -683,18 +684,18 @@ def _session_pgid(pid: int) -> int | None: """Process-group id of the session ``pid`` leads, or None if unavailable. Read this *before* the leader is reaped: once it is, ``getpgid`` fails and - the group id can no longer be recovered from the pid. + the group id can no longer be recovered from the pid. If the group is the + server's own — ``setsid`` did not take effect — there is no session group + to signal, and None makes teardown reach the child alone instead of the + whole server. """ - getpgid = getattr(os, "getpgid", None) - if getpgid is None: # no process groups (native Windows) - return None - try: - return getpgid(pid) - except OSError: + pgid = process_lifecycle.pgid_of(pid) + if pgid is None or pgid == process_lifecycle.own_pgid(): return None + return pgid -def _signal_session(pgid: int | None, pid: int, sig: int) -> bool: +def _signal_session(pgid: int | None, pid: int | None, sig: int) -> bool: """Send ``sig`` to the whole process group, or to the lone process. Returns whether anything was signalled, so a caller can tell "the session @@ -702,20 +703,7 @@ def _signal_session(pgid: int | None, pid: int, sig: int) -> bool: matters: if ``setsid`` did not take effect, or the platform has no process groups, teardown must still reach the child rather than do nothing. """ - killpg = getattr(os, "killpg", None) - if pgid is not None and killpg is not None: - try: - killpg(pgid, sig) - return True - except ProcessLookupError: - return False - except OSError: - pass # group signalling refused — fall through to the single pid - try: - os.kill(pid, sig) - return True - except OSError: - return False + return process_lifecycle.signal_group(pid, pgid, sig) def _session_alive(pgid: int | None, pid: int) -> bool: @@ -723,41 +711,35 @@ def _session_alive(pgid: int | None, pid: int) -> bool: An unreaped zombie is still signallable, so a True here can also mean the leader has exited but not yet been collected. Without a group id this can - only speak for the child itself, not for anything it spawned. + only speak for the child itself, not for anything it spawned. Only ESRCH + proves a group is gone; EPERM is a live group we may not signal, and + reporting a surviving session as contained is the one outcome teardown + must never produce. """ - killpg = getattr(os, "killpg", None) - if pgid is not None and killpg is not None: - try: - killpg(pgid, 0) - except ProcessLookupError: - return False # ESRCH — no member of the group is left - except OSError: - # Anything else (EPERM when the group holds a process we may not - # signal, EINVAL) answers the probe without proving the group is - # gone. Only ESRCH does that, so treat the rest as still running: - # reporting a surviving session as contained is the one outcome - # teardown must never produce. - return True - return True + if pgid is not None: + return process_lifecycle.group_present(pgid) return pid_alive(pid) -async def _await_session_exit(proc, pgid: int | None, pid: int) -> bool: - """Wait up to the grace period for the leader and its group to go away.""" - loop = asyncio.get_running_loop() - deadline = loop.time() + PTY_KILL_GRACE - while True: - remaining = deadline - loop.time() - if proc.returncode is None and remaining > 0: - # Reap the leader, otherwise its own zombie keeps the group alive - # and the liveness probe below can never come back clean. - with contextlib.suppress(asyncio.TimeoutError): - await asyncio.wait_for(proc.wait(), remaining) - if not _session_alive(pgid, pid): - return True - if loop.time() >= deadline: - return False - await asyncio.sleep(PTY_KILL_POLL_INTERVAL) +def _bind_pty_spawn_identity(proc) -> None: + """Freeze the PTY leader's identity and its session group at spawn. + + Called immediately after the spawn, while the pid is known to be the child + just created: we hold it unreaped, so the slot cannot have been reissued. + The group is recorded only when it is the leader's own (``setsid`` + applied: pgid == pid) and the identity still verifies after reading it. + Teardown works from this record alone and never re-derives ownership + from ``proc.pid``, which outlives the process it named. + """ + pid = getattr(proc, "pid", None) + if not pid: + return + identity = process_lifecycle.ProcessIdentity.capture(pid) + pgid = _session_pgid(pid) + if pgid != pid or identity.verdict() != process_lifecycle.OWNED: + pgid = None # No safe session group: teardown reaches the child alone. + proc._ody_pty_identity = process_lifecycle.ProcessIdentity( + pid=identity.pid, start_token=identity.start_token, pgid=pgid) async def _terminate_pty_session(proc) -> bool: @@ -770,18 +752,67 @@ async def _terminate_pty_session(proc) -> bool: command as terminated. Signal the group instead, escalate to SIGKILL if it outlives the grace period, and return whether the session is actually gone so the caller can say so rather than assume it. + + Ownership is the identity frozen at spawn (:func:`_bind_pty_spawn_identity`), + re-verified before every signal: + + * leader OWNED and still leading the recorded group → signal the group; + * leader GONE (exited and reaped) → the recorded group only, never the + pid: a group id is not reissued while the group lives, so a present + group with no process in its leader's slot is still ours; + * leader FOREIGN → the pid was reissued, which proves our group's + lifetime had already ended; nothing of ours is left to signal; + * leader UNVERIFIABLE, or no spawn identity at all → nothing is + signalled and the session is not reported gone. + + The ladder itself is :func:`src.process_lifecycle.escalate_async`; the + leader is reaped through ``proc.wait()`` inside each window, otherwise its + own zombie keeps the group alive and the probe can never come back clean. """ pid = getattr(proc, "pid", None) if pid is None: return True - pgid = _session_pgid(pid) + frozen = getattr(proc, "_ody_pty_identity", None) + if frozen is None: + logger.warning("PTY teardown for pid %s has no spawn identity; not signalling", pid) + return False + pgid = frozen.pgid - for sig in PTY_KILL_ESCALATION: - if not _signal_session(pgid, pid, sig): - break # nothing left to signal - if await _await_session_exit(proc, pgid, pid): + def _gone() -> bool: + verdict = frozen.verdict() + if verdict == process_lifecycle.FOREIGN: return True - return not _session_alive(pgid, pid) + if verdict == process_lifecycle.UNVERIFIABLE: + return False + if pgid is not None: + return not _session_alive(pgid, frozen.pid) + return verdict == process_lifecycle.GONE or process_lifecycle.is_zombie(frozen.pid) + + def _send(sig) -> bool: + verdict = frozen.verdict() + if verdict == process_lifecycle.OWNED: + if pgid is not None and process_lifecycle.pgid_of(frozen.pid) == pgid: + return _signal_session(pgid, frozen.pid, sig) + # Child-only: no safe group, or the leader no longer leads it. + return _signal_session(None, frozen.pid, sig) + if verdict == process_lifecycle.GONE and pgid is not None: + # Never fall back to the pid: it names no process of ours now. + return _signal_session(pgid, None, sig) + return False + + async def _reap_leader(): + if proc.returncode is None: + await proc.wait() + + result = await process_lifecycle.escalate_async( + _gone, + _send, + steps=tuple((sig, PTY_KILL_GRACE) for sig in PTY_KILL_ESCALATION), + wait=_reap_leader, + poll_s=PTY_KILL_POLL_INTERVAL, + wait_floor_s=0.0, + ) + return result.dead async def _terminate_pty_session_quietly(proc) -> None: @@ -829,6 +860,7 @@ async def _generate_pty(cmd: str, timeout: int, request: Request): cwd=str(Path.home()), preexec_fn=os.setsid, ) + _bind_pty_spawn_identity(proc) os.close(slave_fd) # parent doesn't need the slave side deadline = (loop.time() + timeout) if timeout else None diff --git a/src/agent_tools/web_tools.py b/src/agent_tools/web_tools.py index bff68ef66..700ab4a85 100644 --- a/src/agent_tools/web_tools.py +++ b/src/agent_tools/web_tools.py @@ -22,7 +22,7 @@ from pathlib import Path from typing import Dict, Any from core import platform_compat -from src import browser_lifecycle +from src import browser_lifecycle, process_lifecycle from src.constants import MAX_OUTPUT_CHARS PDF_EXTRACT_MAX_BYTES = 80_000_000 @@ -110,6 +110,13 @@ _BROWSER_CALL_PROCS: contextvars.ContextVar[list | None] = contextvars.ContextVa async def _spawn_browser_cli(*command, **kwargs): proc = await asyncio.create_subprocess_exec(*command, **kwargs) + # Identity taken while we hold the unreaped child: teardown later signals + # its group only while this identity still verifies, never on the pid + # alone (the event loop may reap it before returncode is observed). + pid = getattr(proc, "pid", None) + if isinstance(pid, int) and pid > 0: + with contextlib.suppress(Exception): + proc._ody_identity = process_lifecycle.ProcessIdentity.capture(pid, pgid=pid) tracked = _BROWSER_CALL_PROCS.get() if tracked is not None: tracked.append(proc) @@ -2403,12 +2410,22 @@ class PrivateBrowserTool: @staticmethod def _terminate_subprocess(proc) -> None: - """Terminate a browser CLI and descendants spawned for its session.""" + """Terminate a browser CLI and descendants spawned for its session. - pid = getattr(proc, "pid", None) - if pid: - with contextlib.suppress(ProcessLookupError, PermissionError, OSError): - os.killpg(os.getpgid(pid), signal.SIGKILL) + Every browser CLI is spawned with ``start_new_session``, so it leads a + group of its own. That group is signalled only while the identity + captured at spawn still verifies and still leads it; otherwise only + the held handle is killed. A pid alone — or a group derived from a pid + that may since have been reaped and reissued — is never signalled, + and neither is the server's own group. + """ + + identity = getattr(proc, "_ody_identity", None) + if identity is not None and getattr(proc, "returncode", None) is None: + verdict = process_lifecycle.group_ownership_verdict( + identity.pid, identity.pid, identity.start_token) + if verdict == process_lifecycle.OWNED: + process_lifecycle.signal_group(identity.pid, identity.pid, signal.SIGKILL) with contextlib.suppress(Exception): proc.kill() @@ -2440,21 +2457,20 @@ class PrivateBrowserTool: # correctness requirement. Leave those trees to the daemon's own # lifecycle instead of failing the whole shutdown path. return - pids: list[int] = [] + owned: list[process_lifecycle.ProcessIdentity] = [] for entry in platform_compat.PROC_ROOT.iterdir(): if not entry.name.isdigit(): continue - try: - command_line = (entry / "cmdline").read_bytes().replace(b"\0", b" ").decode( - "utf-8", errors="replace" - ) - except (OSError, UnicodeError): - continue - if "--user-data-dir=" + profile_prefix in command_line: - pids.append(int(entry.name)) - for pid in sorted(pids, reverse=True): - with contextlib.suppress(ProcessLookupError, PermissionError, OSError): - os.kill(pid, signal.SIGKILL) + # The command line that matches the profile and the identity that + # will be signalled are read from the same process. + seen = process_lifecycle.observe(int(entry.name), _process_command_line) + if seen is not None and "--user-data-dir=" + profile_prefix in seen.facts: + owned.append(seen.identity) + if owned: + process_lifecycle.terminate_identities( + sorted(owned, key=lambda identity: identity.pid, reverse=True), + steps=((signal.SIGKILL, 1.0),), poll_s=0.02, + ) @staticmethod def _terminate_owned_daemon( @@ -2488,7 +2504,10 @@ class PrivateBrowserTool: pid = int(pid_file.read_text().strip()) except (OSError, ValueError): continue - command_line = _process_command_line(pid) + # A pid file names a slot, not a process: the command-line match + # and the identity that authorises the signal come from one read. + seen = process_lifecycle.observe(pid, _process_command_line) if pid > 0 else None + command_line = seen.facts if seen is not None else None if command_line is None: # Either the daemon exited between writing its pid file and # this pass, or this host has no procfs to ask. Only the first @@ -2500,10 +2519,12 @@ class PrivateBrowserTool: pid_file.unlink() continue if "agent-browser" in command_line: - with contextlib.suppress(ProcessLookupError, PermissionError, OSError): - os.kill(pid, signal.SIGKILL) - with contextlib.suppress(FileNotFoundError, PermissionError, OSError): - pid_file.unlink() + sweep = process_lifecycle.terminate_identities( + [seen.identity], steps=((signal.SIGKILL, 1.0),), poll_s=0.02, + ) + if not sweep.survivors and not sweep.unverified: + with contextlib.suppress(FileNotFoundError, PermissionError, OSError): + pid_file.unlink() return receipt @staticmethod @@ -3819,7 +3840,7 @@ async def shutdown_private_browser_sessions() -> None: ): proc = None try: - proc = await asyncio.create_subprocess_exec( + proc = await _spawn_browser_cli( *command_prefix, "--session", session, "close", stdout=asyncio.subprocess.DEVNULL, stderr=asyncio.subprocess.DEVNULL, diff --git a/src/browser_lifecycle.py b/src/browser_lifecycle.py index 7d266ed63..58cf98622 100644 --- a/src/browser_lifecycle.py +++ b/src/browser_lifecycle.py @@ -6,9 +6,10 @@ POSIX session, so the daemon pid recorded in the session's own pid file identifies the complete browser tree. Cleanup here is limited to that tree, the session's runtime files and its ``agent-browser-chrome-*`` profile. -This is browser-specific ownership only. Generic process containment and -lifecycle primitives belong to the shared process layer; when those exist, -``kill_browser_tree`` is the single seam to replace. +This is browser-specific ownership only: session membership, the profile +prefix, runtime files and navigation state. Process identity, verified +signalling and death observation come from :mod:`src.process_lifecycle`, +reached through the single seam ``kill_browser_tree``. """ from __future__ import annotations @@ -24,6 +25,7 @@ from pathlib import Path from typing import Any, Callable from core import platform_compat +from src import process_lifecycle PROFILE_PREFIX = "agent-browser-chrome-" RUNTIME_SUFFIXES = (".pid", ".sock", ".stream", ".version", ".engine") @@ -121,34 +123,82 @@ def is_verified_daemon(pid: int | None) -> bool: return bool(command_line) and "agent-browser" in command_line -def browser_tree(leader: int) -> list[int]: +@dataclass(frozen=True) +class _Member: + """One process as the membership scan saw it, bound to its identity.""" + + identity: process_lifecycle.ProcessIdentity + state: str + pgid: int + sid: int + cmdline: str + + +def _read_member_facts(pid: int) -> tuple[tuple[str, int, int], str] | None: + stat = _read_stat(pid) + if stat is None: + return None + return stat, _read_cmdline(pid) or "" + + +def _snapshot() -> dict[int, _Member]: + """Every visible process, with the facts membership is decided from. + + Each pid's stat and command line are read between two start-token reads + (:func:`process_lifecycle.observe`), so the identity teardown later + verifies is the identity of the very process membership was decided for — + never one captured afterwards from a pid that may have changed hands. + """ + + snapshot: dict[int, _Member] = {} + for pid in _live_pids(): + seen = process_lifecycle.observe(pid, _read_member_facts) + if seen is None: + continue + (state, pgid, sid), cmdline = seen.facts + snapshot[pid] = _Member( + identity=process_lifecycle.ProcessIdentity( + pid=pid, start_token=seen.identity.start_token, pgid=pgid), + state=state, pgid=pgid, sid=sid, cmdline=cmdline, + ) + return snapshot + + +def browser_members(leader: int) -> list[_Member]: """Processes owned by the browser session whose daemon pid is ``leader``. While the daemon is verified alive, every member of its POSIX session is owned. Once the daemon is gone the pid may be reused, so only Chrome process groups whose root carries an agent-browser profile are claimed. + Decided from one identity-bound snapshot. """ - members: list[tuple[int, int]] = [] - for pid in _live_pids(): - stat = _read_stat(pid) - if stat is None or stat[0] == "Z" or stat[2] != leader: - continue - members.append((pid, stat[1])) + snapshot = _snapshot() + members = [ + member for member in snapshot.values() + if member.state != "Z" and member.sid == leader + ] if not members: return [] - if is_verified_daemon(leader): - return sorted(pid for pid, _ in members) - owned_groups = { - pgid for pid, pgid in members if _profile_dirs([pid]) - } - return sorted(pid for pid, pgid in members if pgid in owned_groups) + daemon = snapshot.get(leader) + if daemon is not None and daemon.state != "Z" and "agent-browser" in daemon.cmdline: + owned = members + else: + owned_groups = {member.pgid for member in members if _profiles_of([member.cmdline])} + owned = [member for member in members if member.pgid in owned_groups] + return sorted(owned, key=lambda member: member.identity.pid) -def _profile_dirs(pids: list[int]) -> set[Path]: +def browser_tree(leader: int) -> list[int]: + """Pids owned by the browser session whose daemon pid is ``leader``.""" + + return [member.identity.pid for member in browser_members(leader)] + + +def _profiles_of(cmdlines: list[str]) -> set[Path]: profiles: set[Path] = set() - for pid in pids: - for token in (_read_cmdline(pid) or "").split(): + for cmdline in cmdlines: + for token in cmdline.split(): if not token.startswith("--user-data-dir="): continue path = Path(token.split("=", 1)[1]) @@ -162,34 +212,30 @@ def kill_browser_tree(leader: int, *, settle_s: float = 1.0) -> tuple[list[int], Returns ``(killed, survivors, profile_dirs)``. Synchronous so it can run from cancellation and shutdown paths without awaiting. + + Which processes form the session is decided here (:func:`browser_members`); + how they are signalled is the generic lifecycle's. Each member's identity + is the one bound to the facts membership was decided from — never + recaptured afterwards — and it is re-verified before the signal, so a pid + freed and reissued at any point after the scan is never hit. A member + whose identity cannot be established is not signalled and is reported as + a survivor: the session still owns it, and its profile must not be + deleted from under it. """ - members = browser_tree(leader) - profiles = _profile_dirs(members) - ordered = [pid for pid in members if pid != leader] - if leader in members: - ordered.append(leader) - killed = [] - for pid in ordered: - try: - os.kill(pid, signal.SIGKILL) - killed.append(pid) - except (ProcessLookupError, PermissionError, OSError): - continue - deadline = time.monotonic() + settle_s - survivors = list(killed) - while survivors and time.monotonic() < deadline: - survivors = [pid for pid in survivors if _alive(pid)] - if survivors: - time.sleep(0.02) - return killed, survivors, profiles - - -def _alive(pid: int) -> bool: - stat = _read_stat(pid) - if stat is None: - return False - return stat[0] != "Z" + members = browser_members(leader) + profiles = _profiles_of([member.cmdline for member in members]) + ordered = [member for member in members if member.identity.pid != leader] + ordered += [member for member in members if member.identity.pid == leader] + # Browser semantics: Chrome is not asked to shut down here — the polite + # path is the agent-browser ``close`` command. This is the forced path. + sweep = process_lifecycle.terminate_identities( + [member.identity for member in ordered], + steps=((signal.SIGKILL, settle_s),), poll_s=0.02, + ) + survivors = [member.identity.pid for member in ordered + if member.identity.pid in sweep.survivors or member.identity.pid in sweep.unverified] + return list(sweep.killed), survivors, profiles @dataclass diff --git a/src/containment.py b/src/containment.py index 64a6505d3..5b2d45824 100644 --- a/src/containment.py +++ b/src/containment.py @@ -53,7 +53,6 @@ import json import logging import os import shutil -import select import signal import subprocess import sys @@ -67,7 +66,7 @@ from typing import Any, Awaitable, Callable, Mapping, Optional from core.atomic_io import atomic_write_json, store_transaction from core.platform_compat import IS_WINDOWS, find_bash, pid_alive -from src import process_ownership +from src import process_lifecycle, process_ownership from src.constants import ( CONTAINMENT_STATE_FILE, MAX_OUTPUT_CHARS, @@ -115,7 +114,7 @@ _RETENTION_S = 3600 # Teardown reads the group liveness probe this often while waiting out the # grace period. Short enough that a cooperative child is not waited on for the # full grace, long enough not to spin. -_DEATH_POLL_S = 0.05 +_DEATH_POLL_S = process_lifecycle.POLL_S # Destinations a bind must never overlay: replacing the private root, the # private /tmp or the workspace itself with a host directory would undo the @@ -285,29 +284,9 @@ class ContainmentResult: release: Optional["ReleaseOutcome"] = None -@dataclass(frozen=True) -class ReleaseOutcome: - """Whether the tree is actually gone, not whether a signal was sent.""" - - dead: bool - escalated: bool - survivors: tuple[int, ...] = () - mechanism: str = "" - #: The ownership verdict, when teardown had to establish one — a grant - #: recovered from the durable store after a restart. Empty for an - #: in-process teardown, where the caller holds the child and the question - #: does not arise. A non-empty value other than - #: :data:`process_ownership.OWNED` means **no signal was sent**. - ownership: str = "" - - def to_dict(self) -> dict[str, Any]: - return { - "dead": self.dead, - "escalated": self.escalated, - "survivors": list(self.survivors), - "mechanism": self.mechanism, - "ownership": self.ownership, - } +#: The termination receipt is generic lifecycle evidence, not a containment +#: concept; the name is kept because tool results and records already carry it. +ReleaseOutcome = process_lifecycle.TerminationOutcome # ── Mechanisms ────────────────────────────────────────────────────────────── @@ -1285,86 +1264,42 @@ async def _capture_namespace_identity(proc) -> None: # ── release ───────────────────────────────────────────────────────────────── -# core.platform_compat.kill_process_tree delegates here as well. Native tools, -# detached jobs and compatibility callers share escalation and death probes. +# Process mechanics — group probes, signalling, escalation, verified death — +# live in src.process_lifecycle, shared with the PTY shell, the Cookbook sweep, +# the browser lifecycle and core.platform_compat.kill_process_tree. What stays +# here is what a grant means: its record, its namespace init and its gate. +# The thin wrappers below are this module's seams; teardown resolves them at +# call time so a test can substitute one probe without replacing the engine. def _own_pgid() -> int: - try: - return os.getpgid(0) - except OSError: # pragma: no cover - getpgid(0) does not fail in practice - return -1 + return process_lifecycle.own_pgid() def _pgid_of(pid: Optional[int]) -> Optional[int]: - if not pid or IS_WINDOWS: - return None - try: - return os.getpgid(int(pid)) - except (OSError, ProcessLookupError, ValueError): + if IS_WINDOWS: return None + return process_lifecycle.pgid_of(pid) def _group_present(pgid: Optional[int]) -> bool: - """True while any process remains in ``pgid``. + """True while any process remains in ``pgid``; never true for our own group. - ``killpg(pgid, 0)`` is the authoritative probe: it raises - ``ProcessLookupError`` once the group is empty, which a per-pid check cannot - tell you — the leader can be gone while its children keep running. The - group id outlives the leader's pid, which is why teardown captures it at - spawn rather than deriving it afterwards. - - Our own group is never reported as present: if ``setsid`` had not applied, - probing it would describe the server, not the child. + EPERM is a live group we cannot signal, not verified death. """ - if not pgid or pgid <= 0 or IS_WINDOWS: + if IS_WINDOWS: return False - if pgid == _own_pgid(): - return False - try: - os.killpg(pgid, 0) - return True - except ProcessLookupError: - return False - except OSError: - return True # EPERM is a live group we cannot signal, not verified death. + return process_lifecycle.group_present(pgid, own=_own_pgid()) def _signal_tree(pid: Optional[int], pgid: Optional[int], sig: int) -> None: - """Signal the whole group, falling back to the leader alone. - - A group that is also *our* group is never signalled: if setsid failed, - killpg would take the server down with the child. - """ - if pgid and pgid > 0 and pgid != _own_pgid(): - try: - os.killpg(pgid, sig) - return - except (OSError, ProcessLookupError): - pass - if pid: - try: - os.kill(int(pid), sig) - except (OSError, ProcessLookupError, ValueError): - pass + """Signal the whole group, falling back to the leader; never our own group.""" + process_lifecycle.signal_group(pid, pgid, sig, own=_own_pgid()) def _reap_if_child(pid: Optional[int]) -> None: - """Clear a zombie we parented, so "alive" means running. - - ``os.kill(pid, 0)`` succeeds for a zombie and a zombie is still a member of - its process group, so without this a process we just killed is reported as a - survivor indefinitely — nothing else is going to reap it. A pid that is not - our child raises ``ChildProcessError`` and there is nothing to do. - - Only the synchronous :func:`release` reaps. A child being awaited is reaped - through ``proc.wait()`` instead, so this never races the event loop's own - child watcher. - """ - if not pid or IS_WINDOWS: + """Clear a zombie we parented, so "alive" means running (sync teardown only).""" + if IS_WINDOWS: return - try: - os.waitpid(int(pid), os.WNOHANG) - except (ChildProcessError, OSError, ValueError): - pass + process_lifecycle.reap_if_child(pid) def _tree_gone(pid: Optional[int], pgid: Optional[int], *, reap: bool = False) -> bool: @@ -1413,13 +1348,11 @@ def _ownership_gate( refusal is that an unreapable orphan has to remain visible instead of being closed out as handled. """ - verdict = process_ownership.verify(pid, token) + # A valid leader identity does not establish ownership of an arbitrary + # recorded process group: a stale or inconsistent PGID is UNVERIFIABLE. + verdict = process_lifecycle.group_ownership_verdict(pid, pgid, token, pgid_of=_pgid_of) if verdict == process_ownership.OWNED: - if IS_WINDOWS or not pgid or _pgid_of(pid) == pgid: - return None - # A valid leader identity does not establish ownership of an arbitrary - # recorded process group. Refuse a stale or inconsistent PGID. - verdict = process_ownership.UNVERIFIABLE + return None if verdict == process_ownership.GONE: # The leader is gone. Its group may still hold processes it @@ -1500,19 +1433,15 @@ def release(grant: ContainmentGrant, *, grace_s: float = 2.0, target = replace(grant, id=grant.id + ":namespace", mechanism="process_group", pid=namespace_pid, pgid=None, namespace_pid=None, namespace_start_token=None) - namespace_fd = None + # Opened before the identity gate inside _release_owner runs: + # a pidfd that still verifies afterwards names that process. + namespace_fd = process_lifecycle.open_pidfd(namespace_pid) try: - if hasattr(os, "pidfd_open") and hasattr(signal, "pidfd_send_signal"): - try: - namespace_fd = os.pidfd_open(namespace_pid) - except OSError: - pass namespace = _release_owner(target, grace_s=grace_s, start_token=namespace_token, require_identity=True, _record_release=False, _pidfd=namespace_fd) finally: - if namespace_fd is not None: - os.close(namespace_fd) + process_lifecycle.close_fd(namespace_fd) outcome = replace(owner, dead=owner.dead and namespace.dead, escalated=owner.escalated or namespace.escalated, survivors=tuple(dict.fromkeys((*owner.survivors, *namespace.survivors)))) @@ -1570,14 +1499,11 @@ def _release_owner(grant: ContainmentGrant, *, grace_s: float = 2.0, grant = replace(grant, pid=pid or None, pgid=pgid) def gone(): if _pidfd is not None: - return bool(select.select([_pidfd], [], [], 0)[0]) + return process_lifecycle.pidfd_exited(_pidfd) return _tree_gone(pid, pgid, reap=True) def send(sig): if _pidfd is not None: - try: - signal.pidfd_send_signal(_pidfd, sig) - except OSError: - pass + process_lifecycle.pidfd_signal(_pidfd, sig) else: _signal_tree(pid, pgid, sig) @@ -1593,15 +1519,7 @@ def _release_owner(grant: ContainmentGrant, *, grace_s: float = 2.0, return refusal if IS_WINDOWS: - try: - subprocess.run( - ["taskkill", "/F", "/T", "/PID", str(pid)], - stdout=subprocess.DEVNULL, - stderr=subprocess.DEVNULL, - creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0), - ) - except Exception: - logger.warning("containment: taskkill failed for pid %s", pid, exc_info=True) + process_lifecycle.taskkill_tree(pid) deadline = time.monotonic() + max(grace_s, 0.0) while time.monotonic() < deadline and pid_alive(pid): time.sleep(_DEATH_POLL_S) @@ -1609,32 +1527,21 @@ def _release_owner(grant: ContainmentGrant, *, grace_s: float = 2.0, finish(grant, outcome) return outcome - if gone(): - outcome = _outcome_for(grant, dead=True, escalated=False) - finish(grant, outcome) - return outcome - - send(signal.SIGTERM) - escalated = False - deadline = time.monotonic() + max(grace_s, 0.0) - while time.monotonic() < deadline and not gone(): - time.sleep(_DEATH_POLL_S) - if not gone(): - escalated = True + def regate(_sig): + # The grace period is long enough for the pid to be freed and reissued; + # a recovered claim must be re-proven before SIGKILL. if recovered or require_identity: - refusal = _ownership_gate(grant, pid, pgid, token) - if refusal is not None: - finish(grant, refusal) - return refusal - send(signal.SIGKILL) - # SIGKILL cannot be caught, so a short verification window is enough. - # Anything still here is out of our reach — a zombie whose parent is - # not us, or a pid we never owned. - deadline = time.monotonic() + 1.0 - while time.monotonic() < deadline and not gone(): - time.sleep(_DEATH_POLL_S) + return _ownership_gate(grant, pid, pgid, token) + return None - outcome = _outcome_for(grant, dead=gone(), escalated=escalated) + result = process_lifecycle.escalate( + gone, send, steps=process_lifecycle.term_kill_steps(grace_s), + poll_s=_DEATH_POLL_S, before_step=regate, + ) + if result.refusal is not None: + finish(grant, result.refusal) + return result.refusal + outcome = _outcome_for(grant, dead=result.dead, escalated=result.escalated) finish(grant, outcome) return outcome @@ -1739,7 +1646,7 @@ async def _release_awaited_impl( _update_record(grant.id, namespace_pid=namespace_pid, namespace_start_token=namespace_token) def namespace_gone(): if namespace_fd is not None: - return bool(select.select([namespace_fd], [], [], 0)[0]) + return process_lifecycle.pidfd_exited(namespace_fd) if namespace_pid: if not pid_alive(namespace_pid): return True @@ -1749,7 +1656,7 @@ async def _release_awaited_impl( return True def gone(): if pidfd is not None: - owner_gone = bool(select.select([pidfd], [], [], 0)[0]) + owner_gone = process_lifecycle.pidfd_exited(pidfd) elif grant.mechanism == "bubblewrap": owner_gone = proc.returncode is not None else: @@ -1757,41 +1664,20 @@ async def _release_awaited_impl( return owner_gone and namespace_gone() def send(sig): if pidfd is not None: - try: - signal.pidfd_send_signal(pidfd, sig) - except OSError: - pass + process_lifecycle.pidfd_signal(pidfd, sig) elif proc.returncode is None or grant.mechanism != "bubblewrap": _signal_tree(pid, pgid, sig) if namespace_fd is not None: - try: - signal.pidfd_send_signal(namespace_fd, sig) - except OSError: - pass + process_lifecycle.pidfd_signal(namespace_fd, sig) elif namespace_pid and process_ownership.verify(namespace_pid, namespace_token) == process_ownership.OWNED: _signal_tree(namespace_pid, None, sig) - send(signal.SIGTERM) - try: - await asyncio.wait_for(proc.wait(), timeout=max(grace_s, 0.05)) - except (asyncio.TimeoutError, ProcessLookupError): - pass - deadline = time.monotonic() + max(grace_s, 0.0) - while time.monotonic() < deadline and not gone(): - await asyncio.sleep(_DEATH_POLL_S) - - escalated = False - if not gone(): - escalated = True - send(signal.SIGKILL) - try: - await asyncio.wait_for(proc.wait(), timeout=1.0) - except (asyncio.TimeoutError, ProcessLookupError): - pass - deadline = time.monotonic() + 1.0 - while time.monotonic() < deadline and not gone(): - await asyncio.sleep(_DEATH_POLL_S) - - outcome = _outcome_for(grant, dead=gone(), escalated=escalated) + # No precheck: SIGTERM goes out first and the leader is reaped through + # proc.wait() before any group probe, or its zombie reads as a survivor. + result = await process_lifecycle.escalate_async( + gone, send, steps=process_lifecycle.term_kill_steps(grace_s), + wait=proc.wait, poll_s=_DEATH_POLL_S, precheck=False, + ) + outcome = _outcome_for(grant, dead=result.dead, escalated=result.escalated) if not namespace_gone(): outcome = replace(outcome, survivors=tuple(dict.fromkeys((*outcome.survivors, namespace_pid)))) _finish_release(grant, outcome) diff --git a/src/process_lifecycle.py b/src/process_lifecycle.py new file mode 100644 index 000000000..afa54ae6b --- /dev/null +++ b/src/process_lifecycle.py @@ -0,0 +1,670 @@ +"""Generic process lifecycle: identity, liveness, signalling, verified death. + +Every runtime-owned subprocess in this tree ends the same way — something has +to decide whether a process is still the one it started, signal it without +hitting a bystander, escalate when it ignores the polite signal, and report +death only when death was observed. Before this module that sequence was +written out four times (containment's sync and async release, the PTY shell, +the Cookbook survivor sweep) and a fifth time without identity at all (the +browser tree kill), and the copies disagreed on what "dead" means and on +whether the server's own process group is fair game. + +What lives here, and what deliberately does not +----------------------------------------------- +This module owns **mechanics**: process identity (pid + start token, never a +pid alone), group and pidfd probes, signal delivery, the TERM → verify → KILL +→ verify escalation, and the termination receipt. It owns no policy about +*which* processes belong to whom: + +* :mod:`src.containment` decides what a grant contains — dimensions, the + bubblewrap boundary, the namespace init, the durable grant store. +* :mod:`src.browser_lifecycle` decides which processes form a browser session + and which files and profiles that session owns. +* Request authority and resource identity decide whether anything runs at all. +* Effects/provenance consume :class:`TerminationOutcome` as evidence; they do + not kill. + +So the engines here take the caller's ``gone()`` and ``send(sig)`` rather than +a pid: the caller knows whether "the tree" is a process group, a pidfd, a +bubblewrap namespace init or a set of snapshot identities, and this module +only guarantees the ordering, the waits, the re-verification point before +escalation, and that the outcome is the observed one. + +Fail-closed rules, shared by every consumer +------------------------------------------- +* A signal requires :data:`process_ownership.OWNED`. GONE, FOREIGN and + UNVERIFIABLE are reasons not to signal, and UNVERIFIABLE is never death. +* A refused liveness probe (``EPERM``) is a live process, never a dead one. +* The server's own process group is never probed as a child's and never + signalled: if ``setsid`` did not apply, ``killpg`` would take the server down. +* An outcome reports ``dead=True`` only when death was observed after the last + signal, not when a signal was sent. +""" + +from __future__ import annotations + +import asyncio +import logging +import os +import select +import signal +import subprocess +import time +from dataclasses import dataclass +from typing import Any, Awaitable, Callable, Iterable, Mapping, Optional, Sequence + +from core import platform_compat +from core.platform_compat import IS_WINDOWS, pid_alive + +from src import process_ownership + +logger = logging.getLogger(__name__) + +OWNED = process_ownership.OWNED +GONE = process_ownership.GONE +FOREIGN = process_ownership.FOREIGN +UNVERIFIABLE = process_ownership.UNVERIFIABLE + +#: How often a waiting teardown re-reads its liveness probe. Short enough that a +#: cooperative process is not waited on for the full grace, long enough not to +#: spin. +POLL_S = 0.05 +#: SIGKILL cannot be caught, so a short window is enough to observe its effect. +#: Anything still present afterwards is out of reach — a zombie whose parent is +#: not us, or a process we were never entitled to signal. +KILL_WAIT_S = 1.0 + + +# ── Receipt ───────────────────────────────────────────────────────────────── +@dataclass(frozen=True) +class TerminationOutcome: + """Whether the target is actually gone, not whether a signal was sent. + + Re-exported as ``containment.ReleaseOutcome``; the ``to_dict`` shape is the + ``teardown`` block tool results and durable records already carry. + """ + + dead: bool + escalated: bool + survivors: tuple[int, ...] = () + mechanism: str = "" + #: The ownership verdict, when teardown had to establish one. A non-empty + #: value other than :data:`OWNED` means **no signal was sent**. + ownership: str = "" + + def to_dict(self) -> dict[str, Any]: + return { + "dead": self.dead, + "escalated": self.escalated, + "survivors": list(self.survivors), + "mechanism": self.mechanism, + "ownership": self.ownership, + } + + +# ── Identity ──────────────────────────────────────────────────────────────── +@dataclass(frozen=True) +class ProcessIdentity: + """A process as a durable claim: the pid slot *and* who occupied it. + + ``start_token`` binds the pid to one process on one boot (see + :mod:`src.process_ownership`). An identity without a token can only ever + verify as UNVERIFIABLE, so it can never authorise a signal. + """ + + pid: int + start_token: Optional[str] + pgid: Optional[int] = None + + @classmethod + def capture(cls, pid: int, *, pgid: Optional[int] = None) -> "ProcessIdentity": + """Identity of whatever holds ``pid`` now. Take it at launch or snapshot.""" + return cls(pid=int(pid), start_token=process_ownership.capture(pid)["start_token"], + pgid=pgid) + + @classmethod + def from_record( + cls, record: Mapping[str, Any], *, pid_key: str = "pid", + token_key: str = "start_token", pgid_key: str = "pgid", + ) -> Optional["ProcessIdentity"]: + """The identity a durable record claims, or None when it names no pid.""" + record = record or {} + try: + pid = int(record.get(pid_key) or 0) + except (TypeError, ValueError): + pid = 0 + if pid <= 0: + return None + try: + pgid = int(record.get(pgid_key) or 0) or None + except (TypeError, ValueError): + pgid = None + return cls(pid=pid, start_token=record.get(token_key) or None, pgid=pgid) + + def verdict(self) -> str: + return process_ownership.verify(self.pid, self.start_token) + + def owned(self) -> bool: + return self.verdict() == OWNED + + def exited(self) -> bool: + """True once the process this identity names has observably ended. + + GONE and FOREIGN both prove the original process is over — a reissued + pid cannot coexist with the process it was taken from. A zombie has + ended too: it runs no code and holds no resources but its exit status. + UNVERIFIABLE is **not** an exit. + """ + verdict = self.verdict() + if verdict in (GONE, FOREIGN): + return True + return verdict == OWNED and is_zombie(self.pid) + + def to_record(self) -> dict[str, Any]: + return {"pid": self.pid, "start_token": self.start_token, "pgid": self.pgid} + + +@dataclass(frozen=True) +class Observation: + """Facts read about a process, bound to the identity they were read from.""" + + identity: ProcessIdentity + facts: Any + + +def observe(pid: int, read: Callable[[int], Any]) -> Optional[Observation]: + """Read facts about ``pid`` and bind them to the process they describe. + + Membership is decided from facts — a session id, a parent, a command line — + and a signal is authorised by identity. If the two are read separately, + the pid can change hands in between and the identity of a stranger gets + attached to a decision made about our process. So the start token is read + *before* and *after* ``read(pid)``: equal tokens prove the facts belong to + that one process, because a reissued pid always carries a later start + time. A pid that exits or is reissued mid-read yields None — its facts + describe no one we can name. + + When this host cannot produce a token, the facts are kept with an + identity whose token is None: it can be reported but never signalled. + ``read`` returning None means the pid had nothing to read (gone). + """ + try: + before = process_ownership.start_token(pid) + except process_ownership.InspectionUnavailable: + before = None + unverifiable = True + else: + unverifiable = False + if before is None: + return None + facts = read(pid) + if facts is None: + return None + if not unverifiable: + try: + after = process_ownership.start_token(pid) + except process_ownership.InspectionUnavailable: + after, before = None, None + else: + if after != before: + return None + return Observation(identity=ProcessIdentity(pid=int(pid), start_token=before), facts=facts) + + +def bind_descendants( + roots: Iterable[int], *, exclude: Iterable[int] = (), +) -> list[Observation]: + """The processes under ``roots``, each bound to its identity. + + :func:`process_ownership.descendants` answers from one table snapshot, and + a token captured afterwards may belong to a process that reused a pid + after the snapshot. Here the tokens are taken between two snapshots, and a + pid is kept only if the second snapshot still places it under ``roots`` + with the same parent and its token has not changed since. Order is the + breadth-first order of the first snapshot; ``facts`` is the + :class:`process_ownership.ProcessInfo` row from the confirming snapshot. + A pid this host cannot identify is kept with a None token: reportable, + never signallable. + + :raises process_ownership.InspectionUnavailable: no process table. + """ + roots = [int(root) for root in roots if root] + excluded = {int(pid) for pid in exclude} + first = process_ownership.process_table() + candidates = [pid for pid in process_ownership.descendants(roots, table=first) + if pid not in excluded] + tokens: dict[int, Optional[str]] = {} + for pid in candidates: + try: + token = process_ownership.start_token(pid) + except process_ownership.InspectionUnavailable: + tokens[pid] = None # Reportable, never signallable. + continue + if token is not None: # None: already gone, nothing to bind. + tokens[pid] = token + second = process_ownership.process_table() + confirmed = set(process_ownership.descendants(roots, table=second)) + bound: list[Observation] = [] + for pid in candidates: + if pid not in tokens or pid not in confirmed or pid not in second or pid not in first: + continue + if second[pid].ppid != first[pid].ppid: + continue + token = tokens[pid] + if token is not None: + verdict = process_ownership.verify(pid, token) + if verdict in (GONE, FOREIGN): + continue # Exited or reissued since the token was taken. + if verdict != OWNED: + token = None # The binding cannot be confirmed. + bound.append(Observation(identity=ProcessIdentity(pid=pid, start_token=token), + facts=second[pid])) + return bound + + +def is_zombie(pid: Optional[int]) -> bool: + """True when procfs reports ``pid`` in state ``Z``. False when it cannot tell.""" + if not pid or not platform_compat.has_procfs(): + return False + try: + raw = (platform_compat.PROC_ROOT / str(int(pid)) / "stat").read_text( + encoding="utf-8", errors="replace") + except (OSError, ValueError): + return False + fields = raw.rpartition(")")[2].split() + return bool(fields) and fields[0] == "Z" + + +# ── Groups ────────────────────────────────────────────────────────────────── +def own_pgid() -> int: + try: + return os.getpgid(0) + except (OSError, AttributeError): # pragma: no cover - no process groups + return -1 + + +def pgid_of(pid: Optional[int]) -> Optional[int]: + """Process group of ``pid``, or None. Read it before the leader is reaped.""" + if not pid or IS_WINDOWS: + return None + try: + return os.getpgid(int(pid)) + except (OSError, ProcessLookupError, ValueError, AttributeError): + return None + + +def group_present(pgid: Optional[int], *, own: Optional[int] = None) -> bool: + """True while any process remains in ``pgid``. + + ``killpg(pgid, 0)`` raising ``ProcessLookupError`` is the only proof the + group is empty; any other refusal (EPERM) is a live group we may not + signal. Our own group is never reported: if ``setsid`` had not applied, + probing it would describe the server, not the child. + """ + if not pgid or pgid <= 0 or IS_WINDOWS: + return False + if pgid == (own_pgid() if own is None else own): + return False + try: + os.killpg(pgid, 0) + return True + except ProcessLookupError: + return False + except OSError: + return True + + +def signal_group(pid: Optional[int], pgid: Optional[int], sig: int, *, + own: Optional[int] = None) -> bool: + """Signal the whole group, falling back to the leader alone. + + Returns whether a signal was delivered. The server's own group is never + signalled; a pgid equal to it falls through to the single pid. + """ + if pgid and pgid > 0 and pgid != (own_pgid() if own is None else own): + try: + os.killpg(pgid, sig) + return True + except ProcessLookupError: + pass + except OSError: + pass # Group signalling refused; the lone pid may still be reachable. + if pid: + try: + os.kill(int(pid), sig) + return True + except (OSError, ValueError): + pass + return False + + +def reap_if_child(pid: Optional[int]) -> None: + """Collect a zombie we parented, so "alive" means running. + + A zombie still answers ``kill(pid, 0)`` and still belongs to its group, so a + process we just killed reads as a survivor until someone waits on it. Only + synchronous teardown calls this; an awaited child is reaped by its waiter. + """ + if not pid or IS_WINDOWS: + return + try: + os.waitpid(int(pid), os.WNOHANG) + except (ChildProcessError, OSError, ValueError): + pass + + +def tree_gone(pid: Optional[int], pgid: Optional[int], *, reap: bool = False, + own: Optional[int] = None) -> bool: + if reap: + reap_if_child(pid) + return not group_present(pgid, own=own) and not pid_alive(pid) + + +def group_ownership_verdict( + pid: Optional[int], pgid: Optional[int], token: Optional[str], *, + pgid_of: Callable[[Optional[int]], Optional[int]] = pgid_of, +) -> str: + """May a recorded (pid, pgid, token) be signalled as a group? + + The leader's identity must verify, and the recorded group must still be the + leader's group: a valid leader does not establish ownership of an + arbitrary recorded pgid. Anything short of that is UNVERIFIABLE. + """ + verdict = process_ownership.verify(pid, token) + if verdict == OWNED and not IS_WINDOWS and pgid and pgid_of(pid) != pgid: + return UNVERIFIABLE + return verdict + + +# ── pidfd ─────────────────────────────────────────────────────────────────── +def pidfd_supported() -> bool: + return hasattr(os, "pidfd_open") and hasattr(signal, "pidfd_send_signal") + + +def open_pidfd(pid: Optional[int]) -> Optional[int]: + """A pidfd for ``pid``, or None when unsupported or the pid is gone. + + A pidfd names the process it was opened on, not the slot. Callers that + open one for a recorded identity must verify the identity *after* opening: + if it still verifies, the handle refers to that process. + """ + if not pid or not pidfd_supported(): + return None + try: + return os.pidfd_open(int(pid)) + except (OSError, ValueError): + return None + + +def pidfd_exited(fd: int) -> bool: + return bool(select.select([fd], [], [], 0)[0]) + + +def pidfd_signal(fd: int, sig: int) -> bool: + try: + signal.pidfd_send_signal(fd, sig) + return True + except OSError: + return False + + +def close_fd(fd: Optional[int]) -> None: + if fd is not None: + try: + os.close(fd) + except OSError: + pass + + +# ── Windows ───────────────────────────────────────────────────────────────── +def taskkill_tree(pid: int) -> None: + """``taskkill /F /T``: Windows has no group escalation, only a forced tree kill.""" + try: + subprocess.run( + ["taskkill", "/F", "/T", "/PID", str(pid)], + stdout=subprocess.DEVNULL, + stderr=subprocess.DEVNULL, + creationflags=getattr(subprocess, "CREATE_NO_WINDOW", 0), + ) + except Exception: + logger.warning("process_lifecycle: taskkill failed for pid %s", pid, exc_info=True) + + +# ── Escalation ────────────────────────────────────────────────────────────── +@dataclass(frozen=True) +class Escalation: + """What an escalation observed. ``refusal`` is the caller's own re-gate result.""" + + dead: bool + escalated: bool + refusal: Any = None + + +def term_kill_steps(grace_s: float, kill_wait_s: float = KILL_WAIT_S) -> tuple[tuple[int, float], ...]: + """The standard POSIX ladder: SIGTERM, ``grace_s``, SIGKILL, ``kill_wait_s``.""" + return ((signal.SIGTERM, max(float(grace_s), 0.0)), + (signal.SIGKILL, max(float(kill_wait_s), 0.0))) + + +def escalate( + gone: Callable[[], bool], + send: Callable[[int], Optional[bool]], + *, + steps: Sequence[tuple[int, float]], + poll_s: float = POLL_S, + precheck: bool = True, + before_step: Optional[Callable[[int], Any]] = None, +) -> Escalation: + """Signal through ``steps`` until ``gone()``; report the observed outcome. + + Synchronous, so it can run from shutdown, cancellation and reaper paths + without an event loop. ``send`` returning ``False`` means nothing was left + to signal; the ladder stops and the final probe decides. ``before_step`` is + called before every step after the first — the point where the target's + identity must be re-established, because the grace period is exactly long + enough for a pid to be freed and reissued. A non-None return aborts with + that value as :attr:`Escalation.refusal`. + """ + escalated = False + for index, (sig, wait_s) in enumerate(steps): + if (precheck or index) and gone(): + return Escalation(dead=True, escalated=escalated) + if index: + escalated = True + if before_step is not None: + refusal = before_step(sig) + if refusal is not None: + return Escalation(dead=False, escalated=escalated, refusal=refusal) + if send(sig) is False: + break + deadline = time.monotonic() + max(wait_s, 0.0) + while time.monotonic() < deadline and not gone(): + time.sleep(poll_s) + return Escalation(dead=gone(), escalated=escalated) + + +async def escalate_async( + gone: Callable[[], bool], + send: Callable[[int], Optional[bool]], + *, + steps: Sequence[tuple[int, float]], + wait: Optional[Callable[[], Awaitable[Any]]] = None, + poll_s: float = POLL_S, + precheck: bool = False, + before_step: Optional[Callable[[int], Any]] = None, + wait_floor_s: float = 0.05, +) -> Escalation: + """:func:`escalate` for a child this coroutine owns. + + After each signal ``wait()`` (normally ``proc.wait``) is awaited within the + step's window before the probe is polled: an unreaped leader is a zombie, a + zombie is still a member of its group, and a group probe would otherwise + report survivors for a tree that has entirely exited. ``wait_floor_s`` + gives that reap a minimum window even when the grace is zero. + """ + escalated = False + loop = asyncio.get_running_loop() + for index, (sig, wait_s) in enumerate(steps): + if (precheck or index) and gone(): + return Escalation(dead=True, escalated=escalated) + if index: + escalated = True + if before_step is not None: + refusal = before_step(sig) + if refusal is not None: + return Escalation(dead=False, escalated=escalated, refusal=refusal) + if send(sig) is False: + break + window = max(wait_s, wait_floor_s if wait is not None else 0.0) + deadline = loop.time() + window + if wait is not None: + try: + await asyncio.wait_for(wait(), timeout=max(deadline - loop.time(), 0.0)) + except (asyncio.TimeoutError, ProcessLookupError, ChildProcessError): + pass + while loop.time() < deadline and not gone(): + await asyncio.sleep(poll_s) + return Escalation(dead=gone(), escalated=escalated) + + +# ── Snapshot identities ───────────────────────────────────────────────────── +@dataclass(frozen=True) +class IdentitySweep: + """Per-identity result of :func:`terminate_identities`. + + ``killed``: signalled by us and observed exited. ``survivors``: still the + same live process after the last step. ``unverified``: could not be + identified, so never signalled — reported, not silently dropped. Pids that + had already exited before any signal appear in none of the three. + """ + + killed: tuple[int, ...] = () + survivors: tuple[int, ...] = () + unverified: tuple[int, ...] = () + + @property + def dead(self) -> bool: + return not self.survivors and not self.unverified + + +def signal_identity(identity: ProcessIdentity, sig: int) -> bool: + """Signal ``identity`` only while it still verifies as OWNED. + + Re-verified immediately before the signal. A pid without a start token, or + one this host cannot inspect, is never signalled. + """ + if identity.verdict() != OWNED: + return False + try: + os.kill(identity.pid, sig) + return True + except (OSError, ValueError): + return False + + +def terminate_identities( + identities: Iterable[ProcessIdentity], + *, + steps: Sequence[tuple[int, float]], + poll_s: float = POLL_S, +) -> IdentitySweep: + """Escalate across a snapshot of identified processes, in the given order. + + For processes this run did not spawn and cannot hold a handle to — the + members of a tmux pane or a browser session, enumerated from the process + table. Every signal is preceded by a fresh verification, so a pid reissued + during the grace period is never hit; the residual window is the gap + between that verification and ``kill(2)`` itself. + """ + targets = list(dict.fromkeys(identities)) + unverified = [ident for ident in targets if ident.verdict() == UNVERIFIABLE] + pending = [ident for ident in targets if ident not in unverified and not ident.exited()] + signalled: list[ProcessIdentity] = [] + for sig, wait_s in steps: + pending = [ident for ident in pending if not ident.exited()] + if not pending: + break + for ident in pending: + if signal_identity(ident, sig) and ident not in signalled: + signalled.append(ident) + deadline = time.monotonic() + max(wait_s, 0.0) + while time.monotonic() < deadline and any(not ident.exited() for ident in pending): + time.sleep(poll_s) + survivors: list[int] = [] + for ident in pending: + if ident.exited(): + continue + if ident.verdict() == UNVERIFIABLE: + unverified.append(ident) + else: + survivors.append(ident.pid) + return IdentitySweep( + killed=tuple(ident.pid for ident in signalled if ident.pid not in survivors + and ident not in unverified), + survivors=tuple(survivors), + unverified=tuple(dict.fromkeys(ident.pid for ident in unverified)), + ) + + +# ── Compatibility teardown ────────────────────────────────────────────────── +def terminate_tree( + pid: Optional[int], + *, + pgid: Optional[int] = None, + start_token: Optional[str] = None, + require_identity: bool = False, + grace_s: float = 2.0, + mechanism: Optional[str] = None, +) -> TerminationOutcome: + """Escalating group teardown for a pid that is not a containment grant. + + With ``require_identity`` the recorded (pid, pgid, token) must pass + :func:`group_ownership_verdict` before the first signal and again before + SIGKILL; otherwise nothing is signalled and the verdict is reported. + """ + name = mechanism or ("windows_tree" if IS_WINDOWS else "process_group") + try: + pid = int(pid) if pid else 0 + except (TypeError, ValueError): + pid = 0 + if pid <= 0: + return TerminationOutcome(dead=True, escalated=False, mechanism=name) + if not pgid: + pgid = pgid_of(pid) + + def refusal() -> Optional[TerminationOutcome]: + if not require_identity: + return None + verdict = group_ownership_verdict(pid, pgid, start_token) + if verdict == OWNED: + return None + if verdict == GONE: + # Leader death does not prove group death, and without a leader + # nothing proves a surviving group is still ours to signal. + alive = group_present(pgid) + return TerminationOutcome(dead=not alive, escalated=False, mechanism=name, + survivors=(pgid,) if alive and pgid else (), + ownership=verdict) + # Never list a foreign or unidentified pid as *our* survivor. + return TerminationOutcome(dead=False, escalated=False, mechanism=name, ownership=verdict) + + refused = refusal() + if refused is not None: + return refused + if IS_WINDOWS: + taskkill_tree(pid) + deadline = time.monotonic() + max(grace_s, 0.0) + while time.monotonic() < deadline and pid_alive(pid): + time.sleep(POLL_S) + return TerminationOutcome(dead=not pid_alive(pid), escalated=True, mechanism=name) + result = escalate( + lambda: tree_gone(pid, pgid, reap=True), + lambda sig: signal_group(pid, pgid, sig), + steps=term_kill_steps(grace_s), + before_step=lambda _sig: refusal(), + ) + if result.refusal is not None: + return result.refusal + survivors = () if result.dead else tuple(dict.fromkeys(v for v in (pid, pgid) if v)) + return TerminationOutcome(dead=result.dead, escalated=result.escalated, + survivors=survivors, mechanism=name) diff --git a/src/process_reaper.py b/src/process_reaper.py index 1c5c7dab7..875c24ad7 100644 --- a/src/process_reaper.py +++ b/src/process_reaper.py @@ -35,7 +35,7 @@ from __future__ import annotations import logging from typing import Any, Dict -from src import process_ownership +from src import process_lifecycle, process_ownership logger = logging.getLogger(__name__) @@ -71,16 +71,19 @@ def reap_containment_grants() -> Dict[str, Any]: containment.forget(grant_id) report["already_gone"] += 1 continue - if record.get("lifetime") == "background" and process_ownership.verify( - record.get("supervisor_pid"), record.get("supervisor_token"), - ) == process_ownership.OWNED: + # A grant's holder — the detached supervisor of a background job, or + # the server process that acquired it — is an identity like any + # other: a live pid in its slot proves nothing without its token. + supervisor = process_lifecycle.ProcessIdentity.from_record( + record, pid_key="supervisor_pid", token_key="supervisor_token") + if record.get("lifetime") == "background" and supervisor and supervisor.owned(): # Detached jobs deliberately survive a server restart. Their # supervisor owns the wall clock and teardown, independently. report["background_kept"] = report.get("background_kept", 0) + 1 continue - if record.get("lifetime") != "cleanup" and record.get("manager_pid") and process_ownership.verify( - record["manager_pid"], record.get("manager_token"), - ) == process_ownership.OWNED: + manager = process_lifecycle.ProcessIdentity.from_record( + record, pid_key="manager_pid", token_key="manager_token") + if record.get("lifetime") != "cleanup" and manager and manager.owned(): report["manager_kept"] = report.get("manager_kept", 0) + 1 continue verdict = process_ownership.verify_record(record) @@ -218,10 +221,18 @@ def reap_legacy_agent_tmux() -> Dict[str, Any]: # current launcher must still match the observed tmux server. report["unverifiable"] += 1 continue - targets = process_ownership.descendants(roots, table=table) - identities = {pid: process_ownership.start_token(pid) for pid in targets} + # Membership and identity bound together: a descendant that changed + # hands after the table was read is dropped, not recorded under a + # stranger's token. One this host cannot identify keeps the whole + # session visible and unsignalled. + bound = process_lifecycle.bind_descendants(roots) + targets = [seen.identity.pid for seen in bound] + identities = {seen.identity.pid: seen.identity.start_token for seen in bound} + if any(token is None for token in identities.values()): + report["unverifiable"] += 1 + continue if snapshot().get(session_id) != panes or process_ownership.verify(server_pid, server_token) != process_ownership.OWNED or any( - process_ownership.verify(pid, identities[pid]) != process_ownership.OWNED for pid in roots + process_ownership.verify(pid, identities.get(pid)) != process_ownership.OWNED for pid in roots ): report["unverifiable"] += 1 continue diff --git a/src/tools/cookbook.py b/src/tools/cookbook.py index 48a70a53c..9318de02a 100644 --- a/src/tools/cookbook.py +++ b/src/tools/cookbook.py @@ -1080,45 +1080,25 @@ async def _capture_session_processes(session_id: str) -> tuple[List[Dict[str, An return [], "; the session had no live pane, so its processes could not be identified" def _snapshot() -> tuple[List[Dict[str, Any]], str]: + from src import process_lifecycle + + # Membership (descends from the pane) and identity (start token) are + # bound together: a pid that changed hands after the table was read + # is dropped instead of being recorded under a stranger's token. try: - table = process_ownership.process_table() + bound = process_lifecycle.bind_descendants( + pane_pids, exclude={os.getpid(), os.getppid()}) except process_ownership.InspectionUnavailable as exc: return [], f"; could not identify the session's processes ({exc})" - own = {os.getpid(), os.getppid()} - records = [] - for pid in process_ownership.descendants(pane_pids, table=table): - if pid in own: - continue - info = table.get(pid) - records.append({ - "pid": pid, - "start_token": process_ownership.capture(pid)["start_token"], - "command": info.command if info else "", - }) - return records, "" + return [{ + "pid": seen.identity.pid, + "start_token": seen.identity.start_token, + "command": seen.facts.command, + } for seen in bound], "" return await asyncio.to_thread(_snapshot) -def _signal_owned(pid: int, token: Optional[str], sig: int) -> bool: - """Signal ``pid`` only while it still verifies as the process we captured. - - Re-verified immediately before every signal, including the escalation: the - gap between SIGTERM and SIGKILL is exactly long enough for the pid to be - freed and reissued, and a SIGKILL aimed at whatever landed in the slot is - the bug this sweep exists to stop committing. - """ - from src import process_ownership - - if process_ownership.verify(pid, token) != process_ownership.OWNED: - return False - try: - os.kill(pid, sig) - return True - except (ProcessLookupError, PermissionError, OSError): - return False - - def _sweep_session_survivors( owned: List[Dict[str, Any]], tracked_cmd: str, capture_note: str, ) -> str: @@ -1133,47 +1113,36 @@ def _sweep_session_survivors( command line, so an identical one may well be a server the user started by hand, and killing it because it resembles ours is indistinguishable from killing ours. Naming it lets whoever is reading decide. - """ - import signal as _signal - import time as _time - from src import process_ownership + The TERM → verify → KILL → verify ladder, and the re-verification before + every signal, are :func:`src.process_lifecycle.terminate_identities`. + """ + from src import process_lifecycle if capture_note: return capture_note - live = [rec for rec in owned - if process_ownership.verify(rec["pid"], rec["start_token"]) == process_ownership.OWNED] - killed: List[int] = [] - survivors: List[int] = [] - for rec in live: - pid, token = rec["pid"], rec["start_token"] - if not _signal_owned(pid, token, _signal.SIGTERM): - continue - deadline = _time.monotonic() + _SWEEP_GRACE_S - while _time.monotonic() < deadline: - if process_ownership.verify(pid, token) != process_ownership.OWNED: - break - _time.sleep(_SWEEP_POLL_S) - if process_ownership.verify(pid, token) == process_ownership.OWNED: - _signal_owned(pid, token, _signal.SIGKILL) - deadline = _time.monotonic() + 1.0 - while _time.monotonic() < deadline: - if process_ownership.verify(pid, token) != process_ownership.OWNED: - break - _time.sleep(_SWEEP_POLL_S) - if process_ownership.verify(pid, token) == process_ownership.OWNED: - survivors.append(pid) - else: - killed.append(pid) + sweep = process_lifecycle.terminate_identities( + (process_lifecycle.ProcessIdentity(pid=int(rec["pid"]), start_token=rec["start_token"]) + for rec in owned), + steps=process_lifecycle.term_kill_steps(_SWEEP_GRACE_S), + poll_s=_SWEEP_POLL_S, + ) note = "" - if killed: - note += f"; killed {len(killed)} surviving process(es) owned by the session" - if survivors: + if sweep.killed: + note += f"; killed {len(sweep.killed)} surviving process(es) owned by the session" + if sweep.survivors: note += ( - f"; {len(survivors)} process(es) survived SIGKILL and are still " - f"running (pid {', '.join(str(pid) for pid in survivors)})" + f"; {len(sweep.survivors)} process(es) survived SIGKILL and are still " + f"running (pid {', '.join(str(pid) for pid in sweep.survivors)})" + ) + if sweep.unverified: + # Captured as the session's, but this host can no longer say whether + # the pid still holds that process — so it was not signalled. + note += ( + f"; {len(sweep.unverified)} process(es) could not be re-identified and " + f"were not signalled (pid {', '.join(str(pid) for pid in sweep.unverified)})" ) note += _unowned_match_note(tracked_cmd, {rec["pid"] for rec in owned}) return note diff --git a/tests/test_agent_tmux_retirement.py b/tests/test_agent_tmux_retirement.py index add2f72e8..4815a929a 100644 --- a/tests/test_agent_tmux_retirement.py +++ b/tests/test_agent_tmux_retirement.py @@ -132,3 +132,34 @@ def test_legacy_cleanup_against_a_private_real_tmux_server(tmp_path, monkeypatch assert containment.active_grants() == [] finally: subprocess.run([real_tmux, "-S", socket, "kill-server"], capture_output=True) + + +def test_a_descendant_reissued_after_the_table_read_is_never_released(legacy, monkeypatch): + reads = {"n": 0} + + def table(): + reads["n"] += 1 + rows = {4200: process_ownership.ProcessInfo(4200, 4100, "/bin/bash --noprofile --norc")} + if reads["n"] <= 2: + rows[4201] = process_ownership.ProcessInfo(4201, 4200, "sleep 60") + else: + # The child exited and its pid now names an unrelated process. + rows[4201] = process_ownership.ProcessInfo(4201, 1, "sshd: stranger") + return rows + + monkeypatch.setattr(process_ownership, "process_table", table) + process_reaper.reap_legacy_agent_tmux() + assert 4201 not in [pid for pid, _ in legacy["released"]] + + +def test_an_unidentifiable_descendant_keeps_the_session_unsignalled(legacy, monkeypatch): + def start_token(pid): + if pid == 4201: + raise process_ownership.InspectionUnavailable("/proc/4201/stat") + return f"token:{pid}" + + monkeypatch.setattr(process_ownership, "start_token", start_token) + report = process_reaper.reap_legacy_agent_tmux() + assert report["unverifiable"] == 1 + assert legacy["released"] == [] + assert all(call[1] != "kill-session" for call in legacy["calls"]) diff --git a/tests/test_browser_lifecycle.py b/tests/test_browser_lifecycle.py index 55450de38..529832e79 100644 --- a/tests/test_browser_lifecycle.py +++ b/tests/test_browser_lifecycle.py @@ -13,14 +13,18 @@ import pytest from core import platform_compat import src.agent_tools.web_tools as web_tools -from src import browser_lifecycle +from src import browser_lifecycle, process_ownership from src.agent_tools.web_tools import PrivateBrowserTool -def _fake_proc(root: Path, pid: int, *, ppid: int, pgid: int, sid: int, cmdline: str, state: str = "S") -> None: +def _fake_proc(root: Path, pid: int, *, ppid: int, pgid: int, sid: int, cmdline: str, + state: str = "S", starttime: int | None = None) -> None: entry = root / str(pid) entry.mkdir(parents=True) - (entry / "stat").write_text(f"{pid} (x y) {state} {ppid} {pgid} {sid} 0 0 0") + # A full stat line: fields after the comm up to starttime (field 22), which + # is what process identity is read from. + tail = " ".join(["0"] * 15 + [str(starttime if starttime is not None else 1000 + pid)]) + (entry / "stat").write_text(f"{pid} (x y) {state} {ppid} {pgid} {sid} {tail}") (entry / "cmdline").write_bytes(cmdline.replace(" ", "\0").encode()) @@ -28,7 +32,12 @@ def _fake_proc(root: Path, pid: int, *, ppid: int, pgid: int, sid: int, cmdline: def fake_procfs(monkeypatch, tmp_path): proc = tmp_path / "proc" proc.mkdir() + boot = proc / "sys/kernel/random/boot_id" + boot.parent.mkdir(parents=True) + boot.write_text("fake-boot\n") monkeypatch.setattr(platform_compat, "PROC_ROOT", proc) + # Identity reads its own binding of the proc root. + monkeypatch.setattr(process_ownership, "PROC_ROOT", proc) killed: list[tuple[int, int]] = [] def _kill(pid, sig): @@ -101,6 +110,69 @@ def test_forced_cleanup_kills_tree_and_removes_owned_resources(fake_procfs, tmp_ assert (proc / "900").exists() +def test_member_recycled_between_membership_and_identity_is_never_signalled( + fake_procfs, tmp_path, monkeypatch, +) -> None: + """Membership sees the old renderer; its pid is reused before any capture or signal. + + The replacement is an unrelated process occupying the same pid slot. Its + identity must never be the one teardown verifies, so it is never hit. + """ + proc, killed = fake_procfs + _browser_tree(proc, tmp_path / "agent-browser-chrome-a") + lifecycle = browser_lifecycle.process_lifecycle + recycled = [] + + def recycle_once(): + if not recycled: + recycled.append(True) + shutil.rmtree(proc / "502") + _fake_proc(proc, 502, ppid=1, pgid=502, sid=502, cmdline="sshd", starttime=99999) + + # Whichever comes first after membership is decided — an identity capture + # or the signalling sweep — the old renderer is gone and its pid reissued. + real_capture = lifecycle.ProcessIdentity.capture + real_terminate = lifecycle.terminate_identities + + def capture(pid, **kwargs): + recycle_once() + return real_capture(pid, **kwargs) + + def terminate(identities, **kwargs): + identities = list(identities) + recycle_once() + return real_terminate(identities, **kwargs) + + monkeypatch.setattr(lifecycle.ProcessIdentity, "capture", staticmethod(capture)) + monkeypatch.setattr(lifecycle, "terminate_identities", terminate) + + killed_pids, survivors, _ = browser_lifecycle.kill_browser_tree(500, settle_s=0.1) + + assert recycled, "the race was never staged" + assert 502 not in [pid for pid, _ in killed], "the replacement process was signalled" + assert (proc / "502").exists() + assert killed_pids == [501, 500] and survivors == [] + + +def test_member_without_identity_is_a_survivor_and_keeps_its_profile(fake_procfs, tmp_path) -> None: + proc, killed = fake_procfs + root = tmp_path / "rt" + root.mkdir() + (root / "ody-k.pid").write_text("500") + profile = tmp_path / "agent-browser-chrome-a" + profile.mkdir() + _browser_tree(proc, profile) + # The renderer's stat is unreadable as identity (no start time): this host + # cannot say which process holds the pid, so it must not be signalled. + (proc / "502" / "stat").write_text("502 (x y) S 501 501 500") + + receipt = browser_lifecycle.force_cleanup(root, "ody-k") + + assert 502 not in [pid for pid, _ in killed] + assert receipt.survivors == [502] and not receipt.verified + assert profile.exists() and (root / "ody-k.pid").exists() + + def test_forced_cleanup_without_procfs_never_kills_unverified_processes(monkeypatch, tmp_path) -> None: monkeypatch.setattr(platform_compat, "PROC_ROOT", tmp_path / "missing") monkeypatch.setattr(browser_lifecycle.os, "kill", lambda *a: pytest.fail("killed")) diff --git a/tests/test_cookbook_stop_without_procfs.py b/tests/test_cookbook_stop_without_procfs.py index 44f7392ce..2aab3b613 100644 --- a/tests/test_cookbook_stop_without_procfs.py +++ b/tests/test_cookbook_stop_without_procfs.py @@ -320,3 +320,69 @@ def test_model_process_scan_returns_empty_without_procfs(monkeypatch, tmp_path): monkeypatch.setattr(os, "listdir", _unexpected_listdir) assert tools._scan_running_model_processes() == [] + + +@pytest.mark.asyncio +async def test_stop_reports_a_survivor_it_can_no_longer_identify(monkeypatch, tmp_path): + """Captured as ours, unverifiable at sweep time: not signalled, and said so.""" + from src import process_ownership + + tracked_cmd = "python -m vllm.entrypoints.openai.api_server --model org/model" + state = _tracked_state(cmd=tracked_cmd) + posts = _install_httpx_client(monkeypatch, state) + _install_successful_tmux_kill(monkeypatch, panes="serve-abc123 900\n") + table = _fake_table(monkeypatch, {900: (1, "bash"), 101: (900, tracked_cmd)}) + asked = {"n": 0} + + def _token(pid): + if int(pid) == 101: + asked["n"] += 1 + if asked["n"] > 1: # the capture succeeded; every later look fails + raise process_ownership.InspectionUnavailable("/proc/101/stat") + return "token:101" + return f"token:{pid}" if int(pid or 0) in table else None + + monkeypatch.setattr(process_ownership, "start_token", _token) + signalled = _install_effective_kill(monkeypatch, table) + + result = await tools.do_stop_served_model(json.dumps({"session_id": "serve-abc123"})) + + assert result["exit_code"] == 0 + assert not any(pid == 101 for pid, _sig in signalled) + assert "could not be re-identified and were not signalled (pid 101)" in result["output"] + assert _stopped_statuses(posts, "serve-abc123") == ["stopped"] + + +@pytest.mark.asyncio +async def test_stop_never_signals_a_pid_reissued_between_the_table_and_its_capture( + monkeypatch, tmp_path +): + """The table places 101 under the pane; 101 is then reissued to a stranger. + + The stranger's token must never be the one recorded for the session. + """ + from src import process_ownership + + tracked_cmd = "python -m vllm.entrypoints.openai.api_server --model org/model" + state = _tracked_state(cmd=tracked_cmd) + posts = _install_httpx_client(monkeypatch, state) + _install_successful_tmux_kill(monkeypatch, panes="serve-abc123 900\n") + table = _fake_table(monkeypatch, {900: (1, "bash"), 101: (900, tracked_cmd)}) + reads = {"n": 0} + + def _table(): + reads["n"] += 1 + if reads["n"] > 1: + # After the first read: the server exited and its pid now belongs + # to an unrelated process with a fresh identity. + table[101] = process_ownership.ProcessInfo(101, 1, "sshd: stranger") + return dict(table) + + monkeypatch.setattr(process_ownership, "process_table", _table) + signalled = _install_effective_kill(monkeypatch, table) + + result = await tools.do_stop_served_model(json.dumps({"session_id": "serve-abc123"})) + + assert result["exit_code"] == 0 + assert not any(pid == 101 for pid, _sig in signalled) + assert _stopped_statuses(posts, "serve-abc123") == ["stopped"] diff --git a/tests/test_private_browser_tool.py b/tests/test_private_browser_tool.py index 8163c43df..5ac659126 100644 --- a/tests/test_private_browser_tool.py +++ b/tests/test_private_browser_tool.py @@ -2,6 +2,7 @@ import asyncio import pytest import base64 import json +import shutil from pathlib import Path from core import platform_compat @@ -1436,27 +1437,77 @@ def test_private_browser_screenshot_without_path_returns_image_payload(monkeypat }] -def test_private_browser_timeout_terminates_the_process_group(monkeypatch) -> None: - calls = [] - +def _cli_proc(calls, identity=None): class _Proc: pid = 1234 + returncode = None + _ody_identity = identity def kill(self): calls.append("fallback-kill") + return _Proc() + + +def test_private_browser_timeout_terminates_the_process_group(monkeypatch) -> None: + from src import process_lifecycle, process_ownership + + calls = [] monkeypatch.setattr(web_tools.os, "getpgid", lambda pid: pid) + monkeypatch.setattr(process_ownership, "verify", lambda pid, token: process_ownership.OWNED) monkeypatch.setattr( web_tools.os, "killpg", lambda pgid, signum: calls.append((pgid, signum)), ) - PrivateBrowserTool._terminate_subprocess(_Proc()) + PrivateBrowserTool._terminate_subprocess( + _cli_proc(calls, process_lifecycle.ProcessIdentity(1234, "spawned", pgid=1234))) assert calls == [(1234, web_tools.signal.SIGKILL), "fallback-kill"] +@pytest.mark.parametrize("verdict", ["foreign", "unverifiable", "gone"]) +def test_cli_group_is_not_signalled_once_its_identity_is_lost(monkeypatch, verdict) -> None: + """A reaped CLI's pid may be reissued; its group is then not ours to kill.""" + from src import process_lifecycle, process_ownership + + calls = [] + monkeypatch.setattr(web_tools.os, "getpgid", lambda pid: pid) + monkeypatch.setattr(process_ownership, "verify", lambda pid, token: verdict) + monkeypatch.setattr(web_tools.os, "killpg", lambda *a: pytest.fail("signalled an unowned group")) + + PrivateBrowserTool._terminate_subprocess( + _cli_proc(calls, process_lifecycle.ProcessIdentity(1234, "spawned", pgid=1234))) + + assert calls == ["fallback-kill"] + + +def test_cli_without_a_spawn_identity_is_never_group_signalled(monkeypatch) -> None: + calls = [] + monkeypatch.setattr(web_tools.os, "getpgid", lambda pid: pid) + monkeypatch.setattr(web_tools.os, "killpg", lambda *a: pytest.fail("signalled a bare pid's group")) + + PrivateBrowserTool._terminate_subprocess(_cli_proc(calls)) + + assert calls == ["fallback-kill"] + + +def test_cli_group_that_moved_is_not_signalled(monkeypatch) -> None: + """The identity verifies but no longer leads the recorded group.""" + from src import process_lifecycle, process_ownership + + calls = [] + monkeypatch.setattr(web_tools.os, "getpgid", lambda pid: 999) + monkeypatch.setattr(process_ownership, "verify", lambda pid, token: process_ownership.OWNED) + monkeypatch.setattr(web_tools.os, "killpg", lambda *a: pytest.fail("signalled a group it does not lead")) + + PrivateBrowserTool._terminate_subprocess( + _cli_proc(calls, process_lifecycle.ProcessIdentity(1234, "spawned", pgid=1234))) + + assert calls == ["fallback-kill"] + + def test_private_browser_retries_one_timed_out_local_open(monkeypatch, tmp_path) -> None: page = tmp_path / "output.html" page.write_text("