mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-10-08 07:52:20 +02:00
refactor(runtime): centralize verified process lifecycle
Extract the generic process lifecycle layer (src/process_lifecycle.py) shared by runtime-owned subprocesses: process identity (pid + boot-bound start token), identity-bound observation, group and pidfd probes, the TERM -> verify -> KILL -> verify escalation with re-gating before escalation, identity-scoped sweeps, and the termination receipt. Containment, the PTY shell, the Cookbook survivor sweep, the browser lifecycle, web_tools browser cleanup, kill_process_tree and the startup reaper consume it while keeping their own ownership semantics. Safety corrections: - browser membership and identity are bound in one snapshot; no identity is recaptured after membership is decided - web_tools legacy pid-file and profile-match kills signal only verified identities; browser CLI groups only while their spawn identity verifies - Cookbook and legacy-tmux descendant capture bind membership to identity - PTY teardown never signals the server's own process group - unverifiable processes are reported, never signalled
This commit is contained in:
@@ -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
@@ -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
@@ -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)
|
||||
|
||||
@@ -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
@@ -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
@@ -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
|
||||
|
||||
Reference in New Issue
Block a user