Merge pull request #54 from pewdiepie-archdaemon/feature/runtime-process-lifecycle

refactor(runtime): centralize verified process lifecycle
This commit is contained in:
Alexandre Teixeira
2026-10-02 10:21:22 +01:00
committed by GitHub
17 changed files with 1837 additions and 417 deletions
+6 -14
View File
@@ -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 ───────────────────────────────────────────
+91 -59
View File
@@ -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
+45 -24
View File
@@ -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,
+91 -45
View File
@@ -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
+58 -172
View File
@@ -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)
+670
View File
@@ -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)
+21 -10
View File
@@ -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
+34 -65
View File
@@ -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
+31
View File
@@ -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"])
+75 -3
View File
@@ -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"))
@@ -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"]
+179 -15
View File
@@ -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("<html><title>retry</title></html>")
@@ -1900,29 +1951,142 @@ def test_terminate_owned_chrome_kills_only_this_runtimes_profile(
) -> None:
"""With procfs present, match on the runtime-owned profile prefix alone."""
proc = tmp_path / "proc"
proc = _fake_procfs(monkeypatch, tmp_path)
tmpdir = tmp_path / "runtime-tmp"
tmpdir.mkdir()
profile_prefix = str(tmpdir.resolve() / "agent-browser-chrome-")
def _write_pid(pid: str, cmdline: str) -> None:
entry = proc / pid
entry.mkdir(parents=True)
(entry / "cmdline").write_bytes(cmdline.replace(" ", "\0").encode())
_write_pid("101", f"chrome --user-data-dir={profile_prefix}abc")
_write_pid("202", "chrome --user-data-dir=/Users/someone/Library/Chrome")
_fake_process(proc, 101, f"chrome --user-data-dir={profile_prefix}abc")
_fake_process(proc, 202, "chrome --user-data-dir=/Users/someone/Library/Chrome")
(proc / "self").mkdir()
monkeypatch.setattr(platform_compat, "PROC_ROOT", proc)
killed: list[int] = []
monkeypatch.setattr(web_tools.os, "kill", lambda pid, sig: killed.append(pid))
killed = _install_lethal_kill(monkeypatch, proc)
PrivateBrowserTool._terminate_owned_chrome({"TMPDIR": str(tmpdir)})
assert killed == [101]
def _fake_procfs(monkeypatch, tmp_path):
from src import process_ownership
proc = tmp_path / "proc"
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)
monkeypatch.setattr(process_ownership, "PROC_ROOT", proc)
return proc
def _fake_process(proc, pid: int, cmdline: str, *, starttime: int | None = None) -> None:
entry = proc / str(pid)
entry.mkdir(parents=True)
tail = " ".join(["0"] * 15 + [str(starttime if starttime is not None else 1000 + pid)])
(entry / "stat").write_text(f"{pid} (x) S 1 {pid} {pid} {tail}")
(entry / "cmdline").write_bytes(cmdline.replace(" ", "\0").encode())
def _install_lethal_kill(monkeypatch, proc, *, before_kill=None):
killed: list[int] = []
def _kill(pid, sig):
if before_kill is not None:
before_kill(pid)
entry = proc / str(pid)
if not entry.exists():
raise ProcessLookupError(pid)
killed.append(pid)
shutil.rmtree(entry)
monkeypatch.setattr(web_tools.os, "kill", _kill)
return killed
def test_owned_chrome_sweep_never_signals_a_reused_pid(monkeypatch, tmp_path) -> None:
"""Matched by profile, then reissued to a stranger before the signal."""
from src import process_lifecycle
proc = _fake_procfs(monkeypatch, tmp_path)
tmpdir = tmp_path / "runtime-tmp"
tmpdir.mkdir()
_fake_process(proc, 101, f"chrome --user-data-dir={tmpdir.resolve()}/agent-browser-chrome-x")
killed = _install_lethal_kill(monkeypatch, proc)
real_terminate = process_lifecycle.terminate_identities
def recycle_then_terminate(identities, **kwargs):
identities = list(identities)
shutil.rmtree(proc / "101")
_fake_process(proc, 101, "postgres", starttime=99999)
return real_terminate(identities, **kwargs)
monkeypatch.setattr(process_lifecycle, "terminate_identities", recycle_then_terminate)
PrivateBrowserTool._terminate_owned_chrome({"TMPDIR": str(tmpdir)})
assert killed == [] and (proc / "101").exists()
def _legacy_pid_file_for(tmp_path, monkeypatch, namespace, session, pid):
"""A pid file in the legacy namespace layout, which only the fallback loop reads."""
monkeypatch.setenv("XDG_RUNTIME_DIR", str(tmp_path))
monkeypatch.setenv("ODYSSEUS_BROWSER_NAMESPACE", namespace)
legacy_run = (tmp_path / "agent-browser" / "namespaces"
/ web_tools._bounded_browser_identity(namespace) / "run")
legacy_run.mkdir(parents=True, exist_ok=True)
target = legacy_run / f"ody-{web_tools._bounded_browser_identity(session)}.pid"
assert target in web_tools._browser_pid_file_candidates(tmp_path, namespace, session)
target.write_text(str(pid))
return target
def test_legacy_daemon_pid_file_kills_only_the_verified_daemon(monkeypatch, tmp_path) -> None:
proc = _fake_procfs(monkeypatch, tmp_path)
_fake_process(proc, 4401, "node agent-browser --serve")
killed = _install_lethal_kill(monkeypatch, proc)
pid_file = _legacy_pid_file_for(tmp_path, monkeypatch, "clawmm-test", "session-7", 4401)
PrivateBrowserTool._terminate_owned_daemon({}, "session-7")
assert killed == [4401] and not pid_file.exists()
def test_legacy_daemon_pid_reused_before_the_signal_is_spared(monkeypatch, tmp_path) -> None:
"""The pid file still names the slot; the daemon in it was replaced."""
from src import process_lifecycle
proc = _fake_procfs(monkeypatch, tmp_path)
_fake_process(proc, 4402, "node agent-browser --serve")
killed = _install_lethal_kill(monkeypatch, proc)
pid_file = _legacy_pid_file_for(tmp_path, monkeypatch, "clawmm-test", "session-8", 4402)
real_terminate = process_lifecycle.terminate_identities
def recycle_then_terminate(identities, **kwargs):
identities = list(identities)
shutil.rmtree(proc / "4402")
_fake_process(proc, 4402, "node agent-browser --serve", starttime=99999)
return real_terminate(identities, **kwargs)
monkeypatch.setattr(process_lifecycle, "terminate_identities", recycle_then_terminate)
PrivateBrowserTool._terminate_owned_daemon({}, "session-8")
# Even a lookalike command line is not the process the match was made on.
assert killed == [] and (proc / "4402").exists()
def test_legacy_daemon_without_identity_keeps_its_pid_file(monkeypatch, tmp_path) -> None:
proc = _fake_procfs(monkeypatch, tmp_path)
_fake_process(proc, 4403, "node agent-browser --serve")
(proc / "4403" / "stat").write_text("4403 (x) S 1") # no start time: unidentifiable
monkeypatch.setattr(web_tools.os, "kill", lambda *a: pytest.fail("signalled an unidentified pid"))
pid_file = _legacy_pid_file_for(tmp_path, monkeypatch, "clawmm-test", "session-9", 4403)
PrivateBrowserTool._terminate_owned_daemon({}, "session-9")
assert pid_file.exists()
def _pid_file_for(tmp_path, monkeypatch, namespace, session, pid):
"""Write a pid file where the daemon helpers will look for it."""
monkeypatch.setenv("XDG_RUNTIME_DIR", str(tmp_path))
+365
View File
@@ -0,0 +1,365 @@
"""The generic process lifecycle: identity, probes, escalation, verified death.
These pin the rules every consumer — containment, the PTY shell, the Cookbook
sweep, the browser lifecycle and kill_process_tree — inherits from one place.
"""
import asyncio
import errno
import os
import signal
import subprocess
import sys
import time
import pytest
from core import platform_compat
from src import process_lifecycle, process_ownership
posix_only = pytest.mark.skipif(os.name == "nt", reason="POSIX process groups and signals")
_IGNORE_TERM = (
"import signal, sys, time\n"
"signal.signal(signal.SIGTERM, signal.SIG_IGN)\n"
"print('ready', flush=True)\n"
"time.sleep(60)\n"
)
def _spawn(code: str = _IGNORE_TERM) -> subprocess.Popen:
proc = subprocess.Popen([sys.executable, "-c", code], stdout=subprocess.PIPE,
stderr=subprocess.DEVNULL, start_new_session=True)
assert proc.stdout.readline().strip() == b"ready"
return proc
def _cleanup(proc: subprocess.Popen) -> None:
if proc.poll() is None:
proc.kill()
proc.wait(timeout=5)
proc.stdout.close()
# ── Groups ──────────────────────────────────────────────────────────────────
@posix_only
def test_our_own_group_is_never_present_and_never_signalled(monkeypatch):
own = os.getpgid(0)
sent = []
monkeypatch.setattr(os, "killpg", lambda pgid, sig: sent.append(("group", pgid, sig)))
monkeypatch.setattr(os, "kill", lambda pid, sig: sent.append(("pid", pid, sig)))
assert process_lifecycle.group_present(own) is False
assert process_lifecycle.signal_group(4242, own, signal.SIGTERM) is True
assert sent == [("pid", 4242, signal.SIGTERM)]
@posix_only
def test_a_refused_group_probe_is_a_live_group(monkeypatch):
def refuse(*_args):
raise PermissionError(errno.EPERM, "Operation not permitted")
monkeypatch.setattr(os, "killpg", refuse)
assert process_lifecycle.group_present(987654, own=1) is True
def gone(*_args):
raise ProcessLookupError(errno.ESRCH, "No such process")
monkeypatch.setattr(os, "killpg", gone)
assert process_lifecycle.group_present(987654, own=1) is False
@posix_only
def test_signal_group_reports_when_nothing_was_left_to_signal(monkeypatch):
def gone(*_args):
raise ProcessLookupError(errno.ESRCH, "No such process")
monkeypatch.setattr(os, "killpg", gone)
monkeypatch.setattr(os, "kill", gone)
assert process_lifecycle.signal_group(4242, 4242, signal.SIGTERM, own=1) is False
# ── Escalation ──────────────────────────────────────────────────────────────
def _ladder():
return ((signal.SIGTERM, 0.01), (getattr(signal, "SIGKILL", signal.SIGTERM), 0.01))
def test_escalate_does_not_signal_a_target_already_gone():
sent = []
result = process_lifecycle.escalate(lambda: True, sent.append, steps=_ladder())
assert result.dead is True and result.escalated is False and sent == []
def test_escalate_terms_then_kills_and_reports_the_observed_outcome():
sent = []
result = process_lifecycle.escalate(lambda: False, sent.append, steps=_ladder(), poll_s=0.001)
assert sent == [step[0] for step in _ladder()]
assert result.dead is False and result.escalated is True
def test_escalate_stops_when_term_is_enough():
sent = []
result = process_lifecycle.escalate(lambda: bool(sent), sent.append, steps=_ladder())
assert sent == [signal.SIGTERM]
assert result.dead is True and result.escalated is False
def test_escalate_reverifies_before_kill_and_refuses_on_a_lost_identity():
sent = []
result = process_lifecycle.escalate(
lambda: False, sent.append, steps=_ladder(), poll_s=0.001,
before_step=lambda sig: "foreign",
)
assert sent == [signal.SIGTERM]
assert result.refusal == "foreign" and result.dead is False
def test_escalate_stops_the_ladder_when_nothing_is_left_to_signal():
sent = []
def send(sig):
sent.append(sig)
return False
result = process_lifecycle.escalate(lambda: False, send, steps=_ladder())
assert sent == [signal.SIGTERM] and result.dead is False
async def test_escalate_async_signals_first_and_reaps_inside_the_window():
sent, waited = [], []
state = {"gone": False}
async def wait():
waited.append(True)
state["gone"] = True
result = await process_lifecycle.escalate_async(
lambda: state["gone"], sent.append, steps=_ladder(), wait=wait,
)
# No precheck: an awaited leader is reaped only after the first signal.
assert sent == [signal.SIGTERM] and waited == [True]
assert result.dead is True and result.escalated is False
async def test_escalate_async_gives_the_reap_a_floor_even_at_zero_grace():
seen = {}
async def wait():
seen["ran"] = True
await process_lifecycle.escalate_async(
lambda: False, lambda sig: None, steps=((signal.SIGTERM, 0.0),), wait=wait,
)
assert seen == {"ran": True}
# ── Identity ────────────────────────────────────────────────────────────────
def test_identity_from_record_needs_a_pid():
assert process_lifecycle.ProcessIdentity.from_record({}) is None
assert process_lifecycle.ProcessIdentity.from_record({"pid": "x"}) is None
ident = process_lifecycle.ProcessIdentity.from_record(
{"supervisor_pid": 42, "supervisor_token": "t", "pgid": "7"},
pid_key="supervisor_pid", token_key="supervisor_token",
)
assert ident == process_lifecycle.ProcessIdentity(pid=42, start_token="t", pgid=7)
def test_a_pid_without_a_token_is_never_owned_and_never_exited(monkeypatch):
monkeypatch.setattr(process_ownership, "start_token", lambda pid: "live")
ident = process_lifecycle.ProcessIdentity(pid=4242, start_token=None)
assert ident.verdict() == process_ownership.UNVERIFIABLE
assert ident.owned() is False and ident.exited() is False
@pytest.mark.parametrize("current,exited", [(None, True), ("other", True), ("mine", False)])
def test_exited_means_gone_or_recycled(monkeypatch, current, exited):
monkeypatch.setattr(process_ownership, "start_token", lambda pid: current)
monkeypatch.setattr(process_lifecycle, "is_zombie", lambda pid: False)
ident = process_lifecycle.ProcessIdentity(pid=4242, start_token="mine")
assert ident.exited() is exited
@pytest.mark.skipif(not platform_compat.has_procfs(), reason="zombie state needs procfs")
def test_an_unreaped_zombie_has_exited():
proc = subprocess.Popen([sys.executable, "-c", "pass"])
ident = process_lifecycle.ProcessIdentity.capture(proc.pid)
try:
for _ in range(250):
if process_lifecycle.is_zombie(proc.pid):
break
time.sleep(0.02)
assert ident.verdict() == process_ownership.OWNED # still in its slot
assert ident.exited() is True
finally:
proc.wait(timeout=5)
# ── Snapshot identities ─────────────────────────────────────────────────────
def test_terminate_identities_never_signals_foreign_or_unverifiable(monkeypatch):
tokens = {1: "now-someone-else", 2: None}
def start_token(pid):
if pid == 3:
raise process_ownership.InspectionUnavailable("stat")
return tokens.get(pid)
monkeypatch.setattr(process_ownership, "start_token", start_token)
monkeypatch.setattr(os, "kill", lambda *a: pytest.fail("signalled an unowned pid"))
sweep = process_lifecycle.terminate_identities(
[process_lifecycle.ProcessIdentity(1, "mine"),
process_lifecycle.ProcessIdentity(2, "mine"),
process_lifecycle.ProcessIdentity(3, "mine")],
steps=_ladder(),
)
assert sweep.killed == () and sweep.survivors == ()
assert sweep.unverified == (3,) and sweep.dead is False
def test_terminate_identities_kills_in_order_and_reverifies_each_signal(monkeypatch):
alive = {10: "a", 11: "b", 12: "c"}
sent = []
monkeypatch.setattr(process_ownership, "start_token", lambda pid: alive.get(pid))
monkeypatch.setattr(process_lifecycle, "is_zombie", lambda pid: False)
def kill(pid, sig):
sent.append((pid, sig))
alive.pop(pid, None)
monkeypatch.setattr(os, "kill", kill)
sweep = process_lifecycle.terminate_identities(
[process_lifecycle.ProcessIdentity(pid, token) for pid, token in list(alive.items())],
steps=((signal.SIGTERM, 0.05),),
)
assert [pid for pid, _ in sent] == [10, 11, 12]
assert sweep.killed == (10, 11, 12) and sweep.dead
@posix_only
@pytest.mark.skipif(not platform_compat.has_procfs(), reason="real identity needs procfs")
def test_terminate_identities_escalates_past_an_ignored_sigterm():
proc = _spawn()
try:
ident = process_lifecycle.ProcessIdentity.capture(proc.pid)
sweep = process_lifecycle.terminate_identities(
[ident], steps=process_lifecycle.term_kill_steps(0.2))
assert sweep.killed == (proc.pid,) and sweep.survivors == ()
assert proc.wait(timeout=5) == -signal.SIGKILL
finally:
_cleanup(proc)
# ── terminate_tree / kill_process_tree ──────────────────────────────────────
@posix_only
@pytest.mark.skipif(not platform_compat.has_procfs(), reason="real identity needs procfs")
def test_terminate_tree_requires_a_matching_identity():
proc = _spawn()
try:
refused = process_lifecycle.terminate_tree(
proc.pid, pgid=proc.pid, start_token="procfs:not-this-process",
require_identity=True, grace_s=0.1)
assert refused.dead is False and refused.ownership == process_ownership.FOREIGN
assert refused.survivors == ()
assert proc.poll() is None
token = process_ownership.start_token(proc.pid)
outcome = process_lifecycle.terminate_tree(
proc.pid, pgid=proc.pid, start_token=token, require_identity=True, grace_s=0.2)
assert outcome.dead is True and outcome.escalated is True
finally:
_cleanup(proc)
@posix_only
def test_terminate_tree_keeps_a_leaderless_group_visible(monkeypatch):
monkeypatch.setattr(process_ownership, "start_token", lambda pid: None)
monkeypatch.setattr(process_lifecycle, "group_present", lambda pgid, **kw: True)
monkeypatch.setattr(process_lifecycle, "signal_group",
lambda *a, **kw: pytest.fail("signalled an unproven group"))
outcome = process_lifecycle.terminate_tree(4242, pgid=4242, start_token="t",
require_identity=True)
assert outcome.dead is False and outcome.survivors == (4242,)
assert outcome.ownership == process_ownership.GONE
@posix_only
def test_kill_process_tree_writes_no_containment_record(monkeypatch):
from src import containment
monkeypatch.setattr(containment, "_update_record", lambda *a, **k: pytest.fail("grant store touched"))
proc = _spawn("import time\nprint('ready', flush=True)\ntime.sleep(60)\n")
try:
outcome = platform_compat.kill_process_tree(proc.pid)
assert outcome.dead is True
assert isinstance(outcome, containment.ReleaseOutcome)
finally:
_cleanup(proc)
def test_termination_outcome_shape_is_the_teardown_block():
outcome = process_lifecycle.TerminationOutcome(dead=False, escalated=True, survivors=(1,),
mechanism="process_group", ownership="owned")
assert outcome.to_dict() == {"dead": False, "escalated": True, "survivors": [1],
"mechanism": "process_group", "ownership": "owned"}
# ── Identity-bound observation ──────────────────────────────────────────────
def _tokens(monkeypatch, sequence):
calls = iter(sequence)
def start_token(pid):
value = next(calls)
if isinstance(value, Exception):
raise value
return value
monkeypatch.setattr(process_ownership, "start_token", start_token)
def test_observe_binds_facts_to_the_token_read_around_them(monkeypatch):
_tokens(monkeypatch, ["t1", "t1"])
seen = process_lifecycle.observe(42, lambda pid: "facts")
assert seen.identity == process_lifecycle.ProcessIdentity(42, "t1") and seen.facts == "facts"
def test_observe_drops_facts_read_across_a_pid_reuse(monkeypatch):
_tokens(monkeypatch, ["old", "new"])
assert process_lifecycle.observe(42, lambda pid: "whose facts?") is None
def test_observe_of_a_gone_pid_reads_nothing(monkeypatch):
_tokens(monkeypatch, [None])
assert process_lifecycle.observe(42, lambda pid: pytest.fail("read a gone pid")) is None
@pytest.mark.parametrize("sequence", [
[process_ownership.InspectionUnavailable("stat")],
["t1", process_ownership.InspectionUnavailable("stat")],
])
def test_observe_keeps_unidentifiable_facts_without_a_token(monkeypatch, sequence):
_tokens(monkeypatch, sequence)
seen = process_lifecycle.observe(42, lambda pid: "facts")
assert seen.identity.start_token is None and seen.facts == "facts"
assert seen.identity.verdict() == process_ownership.UNVERIFIABLE
def test_bind_descendants_drops_a_pid_reparented_between_snapshots(monkeypatch):
Info = process_ownership.ProcessInfo
tables = iter([
{10: Info(10, 1, "shell"), 11: Info(11, 10, "server")},
{10: Info(10, 1, "shell"), 11: Info(11, 1, "stranger")},
])
monkeypatch.setattr(process_ownership, "process_table", lambda: next(tables))
monkeypatch.setattr(process_ownership, "start_token", lambda pid: f"t{pid}")
bound = process_lifecycle.bind_descendants([10])
assert [seen.identity.pid for seen in bound] == [10]
def test_bind_descendants_drops_a_pid_whose_token_changed(monkeypatch):
Info = process_ownership.ProcessInfo
table = {10: Info(10, 1, "shell"), 11: Info(11, 10, "server")}
monkeypatch.setattr(process_ownership, "process_table", lambda: dict(table))
monkeypatch.setattr(process_ownership, "start_token", lambda pid: f"t{pid}")
monkeypatch.setattr(process_ownership, "verify",
lambda pid, token: process_ownership.FOREIGN if pid == 11 else process_ownership.OWNED)
bound = process_lifecycle.bind_descendants([10], exclude={99})
assert [(seen.identity.pid, seen.facts.command) for seen in bound] == [(10, "shell")]
+20 -1
View File
@@ -243,6 +243,9 @@ def test_native_terminal_runtime_adds_offered_tools_deliberately():
# belonging to the user, or to another worktree, must survive it.
def test_chrome_sweep_kills_only_this_runtimes_profile(monkeypatch, tmp_path):
import shutil
from src import process_ownership
from src.agent_tools.web_tools import PrivateBrowserTool
proc = tmp_path / "proc"
@@ -250,9 +253,19 @@ def test_chrome_sweep_kills_only_this_runtimes_profile(monkeypatch, tmp_path):
tmpdir.mkdir()
ours = str(tmpdir.resolve() / "agent-browser-chrome-")
# A fake process needs an identity, not only a command line: the sweep
# signals a process only after verifying the start token it matched on.
# Without a stat here the token would be read from the host's real /proc,
# so the outcome would depend on whether this pid happens to exist.
boot = proc / "sys/kernel/random/boot_id"
boot.parent.mkdir(parents=True)
boot.write_text("fake-boot\n")
def _pid(pid, cmdline):
entry = proc / pid
entry.mkdir(parents=True)
starttime = " ".join(["0"] * 15 + [str(1000 + int(pid))])
(entry / "stat").write_text(f"{pid} (chrome) S 1 {pid} {pid} {starttime}")
(entry / "cmdline").write_bytes(cmdline.replace(" ", "\0").encode())
_pid("101", f"chrome --user-data-dir={ours}session-a")
@@ -261,8 +274,14 @@ def test_chrome_sweep_kills_only_this_runtimes_profile(monkeypatch, tmp_path):
(proc / "self").mkdir()
monkeypatch.setattr(platform_compat, "PROC_ROOT", proc)
monkeypatch.setattr(process_ownership, "PROC_ROOT", proc)
killed = []
monkeypatch.setattr(al_web.os, "kill", lambda pid, sig: killed.append(pid))
def _kill(pid, sig):
killed.append(pid)
shutil.rmtree(proc / str(pid), ignore_errors=True)
monkeypatch.setattr(al_web.os, "kill", _kill)
PrivateBrowserTool._terminate_owned_chrome({"TMPDIR": str(tmpdir)})
+82 -8
View File
@@ -119,13 +119,31 @@ pty_session = pytest.mark.skipif(
async def _spawn_pty_style_session(script: str):
"""Spawn `script` the way _generate_pty does: its own session via setsid."""
return await asyncio.create_subprocess_shell(
"""Spawn `script` the way _generate_pty does: its own session via setsid,
with the leader's identity and group bound at spawn."""
import routes.shell_routes as shell_routes
proc = await asyncio.create_subprocess_shell(
script,
stdout=asyncio.subprocess.DEVNULL,
stderr=asyncio.subprocess.DEVNULL,
preexec_fn=os.setsid,
)
shell_routes._bind_pty_spawn_identity(proc)
return proc
def _bound_fake_leader(monkeypatch, pid=4242):
"""A fake PTY leader whose spawn identity verifies and still leads its group."""
from src import process_lifecycle, process_ownership
real_getpgid = os.getpgid
monkeypatch.setattr(process_ownership, "verify", lambda p, token: process_ownership.OWNED)
monkeypatch.setattr(os, "getpgid", lambda p: pid if p == pid else real_getpgid(p))
return SimpleNamespace(
pid=pid, returncode=0, wait=None,
_ody_pty_identity=process_lifecycle.ProcessIdentity(pid, "spawn-token", pgid=pid),
)
def _stubborn_child(pid_file: Path, ignore: tuple[str, ...]) -> str:
@@ -281,9 +299,8 @@ async def test_terminate_pty_session_reports_a_session_it_could_not_kill(
monkeypatch.setattr(shell_routes, "PTY_KILL_GRACE", 0.01)
monkeypatch.setattr(shell_routes, "_signal_session", lambda *_: True)
monkeypatch.setattr(shell_routes, "_session_alive", lambda *_: True)
monkeypatch.setattr(shell_routes, "_session_pgid", lambda _: 4242)
proc = SimpleNamespace(pid=4242, returncode=0, wait=None)
proc = _bound_fake_leader(monkeypatch)
assert await shell_routes._terminate_pty_session(proc) is False
@@ -293,7 +310,6 @@ async def test_terminate_pty_session_escalates_before_giving_up(monkeypatch):
sent = []
monkeypatch.setattr(shell_routes, "PTY_KILL_GRACE", 0.01)
monkeypatch.setattr(shell_routes, "_session_pgid", lambda _: 4242)
monkeypatch.setattr(shell_routes, "_session_alive", lambda *_: True)
monkeypatch.setattr(
shell_routes,
@@ -301,7 +317,7 @@ async def test_terminate_pty_session_escalates_before_giving_up(monkeypatch):
lambda pgid, pid, sig: sent.append(sig) or True,
)
proc = SimpleNamespace(pid=4242, returncode=0, wait=None)
proc = _bound_fake_leader(monkeypatch)
await shell_routes._terminate_pty_session(proc)
assert sent == [signal.SIGTERM, signal.SIGKILL]
@@ -338,6 +354,65 @@ async def test_generate_pty_timeout_says_so_when_the_session_survives(
)
@pytest.mark.skipif(os.name == "nt", reason="POSIX process groups")
async def test_terminate_pty_session_never_signals_the_servers_own_group(monkeypatch):
"""If setsid did not apply, the child's group is ours: reach the child alone."""
import routes.shell_routes as shell_routes
from src import process_ownership
own = os.getpgid(0)
sent = []
monkeypatch.setattr(shell_routes, "PTY_KILL_GRACE", 0.01)
monkeypatch.setattr(shell_routes.process_lifecycle, "pgid_of", lambda _pid: own)
monkeypatch.setattr(process_ownership, "start_token", lambda pid: "the-leader")
proc = SimpleNamespace(pid=987654, returncode=0, wait=None)
assert shell_routes._session_pgid(proc.pid) is None
shell_routes._bind_pty_spawn_identity(proc)
assert proc._ody_pty_identity.pgid is None # no safe session group recorded
monkeypatch.setattr(os, "killpg", lambda pgid, sig: sent.append(("group", pgid, sig)))
monkeypatch.setattr(os, "kill", lambda pid, sig: sent.append(("pid", pid, sig)))
await shell_routes._terminate_pty_session(proc)
assert sent and all(kind == "pid" and target == 987654 for kind, target, _ in sent), sent
@pytest.mark.skipif(os.name == "nt", reason="POSIX process groups")
async def test_terminate_pty_session_never_signals_a_reused_leader_pid(monkeypatch):
"""Leader spawned and bound → reaped → pid reissued → teardown signals nothing.
The replacement is the worst case: an unrelated session leader, so both
its pid and its process group carry the number our leader had.
"""
import routes.shell_routes as shell_routes
from src import process_ownership
pid = 987650
real_getpgid = os.getpgid
occupant = {"token": "leader-token"}
monkeypatch.setattr(process_ownership, "start_token",
lambda p: occupant["token"] if int(p) == pid else None)
monkeypatch.setattr(os, "getpgid", lambda p: pid if p == pid else real_getpgid(p))
proc = SimpleNamespace(pid=pid, returncode=None, wait=None)
shell_routes._bind_pty_spawn_identity(proc) # spawn time: the leader we just created
assert proc._ody_pty_identity.pgid == pid
# The leader exits and is reaped; the kernel reissues its pid to a stranger.
proc.returncode = 0
occupant["token"] = "replacement-token"
signalled = []
monkeypatch.setattr(os, "killpg", lambda g, sig: sig and signalled.append(("group", g, sig)))
monkeypatch.setattr(os, "kill", lambda p, sig: sig and signalled.append(("pid", p, sig)))
monkeypatch.setattr(shell_routes, "PTY_KILL_GRACE", 0.01)
await shell_routes._terminate_pty_session(proc)
assert signalled == [], f"teardown signalled the replacement: {signalled}"
def test_session_alive_treats_a_refused_probe_as_alive(monkeypatch):
"""EPERM says the group exists but we may not signal it, not that it died.
@@ -372,11 +447,10 @@ async def test_terminate_pty_session_reports_a_group_it_may_not_signal(monkeypat
raise PermissionError(errno.EPERM, "Operation not permitted")
monkeypatch.setattr(shell_routes, "PTY_KILL_GRACE", 0.01)
monkeypatch.setattr(shell_routes, "_session_pgid", lambda _: 4242)
proc = _bound_fake_leader(monkeypatch)
monkeypatch.setattr(shell_routes.os, "killpg", refuse)
monkeypatch.setattr(shell_routes.os, "kill", refuse)
proc = SimpleNamespace(pid=4242, returncode=0, wait=None)
assert await shell_routes._terminate_pty_session(proc) is False
@@ -1,5 +1,7 @@
from pathlib import Path
import pytest
def test_unoffered_artifact_recovery_is_bounded():
from src.agent_loop import _artifact_unoffered_recovery_exhausted
+1 -1
View File
@@ -90,7 +90,7 @@ The source tree reads **109** `ODYSSEUS_*` variables: 79 an operator may want to
| `ODYSSEUS_BROWSER_MCP_REQUIRE_CACHE` | `''` | `src/builtin_mcp.py:90` | Truthy refuses to start the browser MCP server unless its npm package is already in the npx cache, instead of installing it at startup. |
| `ODYSSEUS_BROWSER_NAMESPACE` | `'odysseus-ui'` | `src/agent_tools/web_tools.py:100` (+3 more) | Namespace for the detached agent-browser daemon's pid files, so two runtimes on one machine do not terminate each other's browsers. |
| `ODYSSEUS_BROWSER_NO_SANDBOX` | `'1'` | `src/builtin_mcp.py:142` | Security-relevant. On by default, adding `--no-sandbox` because the Docker image cannot use the Chromium sandbox. Set 0, false or no to keep it. |
| `ODYSSEUS_BROWSER_SCREENSHOT_DIR` | *unset* | `src/agent_tools/web_tools.py:3458` | Where private-browser screenshots are written. Falls back to the container path, then the system temp directory. |
| `ODYSSEUS_BROWSER_SCREENSHOT_DIR` | *unset* | `src/agent_tools/web_tools.py:3479` | Where private-browser screenshots are written. Falls back to the container path, then the system temp directory. |
### Container and workspace mounts