From c004a26d46ea0b0b1422fd48f51a9dc20e73cd0c Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?L=C3=A9o?= Date: Thu, 1 Oct 2026 18:55:50 +0200 Subject: [PATCH] fix(runtime): verify process identity before any teardown signal MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit A recorded pid is a claim, not a handle. The containment grant store, the background-job store and the Cookbook task list all outlive the process that wrote them — deliberately, so a restart keeps a job and its result — and the kernel reuses pids. Any teardown driven off one of those records can therefore land on a process we never started. ODY-86 was exactly this, and the Cookbook survivor sweep still terminated any process whose full command line matched a tracked one, which is the same mistake spelled differently. Identity is (pid, start token). The token comes from /proc//stat on Linux, ps -o lstart= on macOS and the BSDs, and GetProcessTimes on Windows; the kernel will not hand a pid to a process that started earlier, so comparing the token recorded at launch against the token read now answers "is this still ours" without a handle or a supervisor. verify() returns owned, gone, foreign or unverifiable, and only owned permits a signal. Keeping "unverifiable" out of the other two is the point. Process inspection has broken off Linux four times here — ODY-70, -86, -94, -99 — every time because an absent mechanism read as a successful answer. Folding it into "ours" signals strangers; folding it into "gone" abandons live processes. It is a containment failure and every caller treats it as one. Wired into the three places that signal: - containment.release() gates a grant recovered from the durable store, and leaves an in-process teardown alone, where the caller holds the child and no identity question arises. The verdict lands on the record, so "why is this grant still here" is answerable afterwards. - A startup reaper. Nothing read either store before, so a crashed run left every grant permanently active and every job permanently running, and the first thing to touch such a record was a teardown aimed at a reassigned pid. The two stores get opposite treatment: an orphaned grant has no caller left and is torn down, while a detached job is documented to survive a restart and is only corrected, never killed. - The Cookbook sweep takes its ownership from the tmux pane's process tree, captured before the kill destroys the only link between a surviving model server and the session that started it. A process that merely matches the tracked command line is now reported rather than killed: the Cookbook composed that command line, so an identical one is just as likely to be a server the user started by hand. The sweep also runs on hosts with no procfs instead of silently skipping, and says so when it could not look at all. --- app.py | 8 + src/bg_jobs.py | 63 +++- src/containment.py | 143 ++++++- src/process_ownership.py | 405 ++++++++++++++++++++ src/process_reaper.py | 158 ++++++++ src/tools/cookbook.py | 233 ++++++++++-- tests/test_cookbook_stop_without_procfs.py | 226 +++++++++-- tests/test_orphan_reaping.py | 412 +++++++++++++++++++++ tests/test_process_ownership.py | 327 ++++++++++++++++ 9 files changed, 1908 insertions(+), 67 deletions(-) create mode 100644 src/process_ownership.py create mode 100644 src/process_reaper.py create mode 100644 tests/test_orphan_reaping.py create mode 100644 tests/test_process_ownership.py diff --git a/app.py b/app.py index bd1e1580f..5215966fb 100644 --- a/app.py +++ b/app.py @@ -1345,6 +1345,14 @@ async def _startup_event(): from src.cookbook_serve_lifecycle import cookbook_serve_lifecycle_loop _startup_tasks.append(asyncio.create_task(cookbook_serve_lifecycle_loop())) + # Reconcile the processes a previous run left behind: tear down orphaned + # containment grants, and stop trusting background-job records whose pid the + # kernel has since reassigned. Runs once, and deliberately runs *here* — + # every record it sees predates this run, which is what makes "I cannot + # identify this process" a safe thing to act on. See src/process_reaper.py. + from src.process_reaper import reap_orphans_at_startup + _startup_tasks.append(asyncio.create_task(reap_orphans_at_startup())) + logger.info("Application startup complete") async def _shutdown_event(): diff --git a/src/bg_jobs.py b/src/bg_jobs.py index f864f8ef1..7d1d44cf5 100644 --- a/src/bg_jobs.py +++ b/src/bg_jobs.py @@ -38,6 +38,7 @@ from core.platform_compat import ( pid_alive, ) +from src import process_ownership from src.constants import BG_JOBS_DIR, BG_JOBS_FILE _JOBS_DIR = Path(BG_JOBS_DIR) @@ -152,6 +153,11 @@ def launch(command: str, session_id: str, cwd: Optional[str] = None, "followed_up": False, # has the agent been re-invoked with the result? "log_path": str(log_path), "exit_path": str(exit_path), + # Identity, not just a slot. The pid above is reused by the kernel, and + # this record outlives the process and the server; the token is what a + # later run compares before it signals anything. See + # src/process_ownership.py. + "start_token": process_ownership.capture(proc.pid)["start_token"], } jobs = _load() jobs[job_id] = rec @@ -283,10 +289,65 @@ def kill(job_id: str) -> Optional[Dict[str, Any]]: return rec +def disown_unverified() -> Dict[str, Any]: + """Stop tracking running jobs whose process can no longer be proven ours. + + Called once at startup by :mod:`src.process_reaper`, never from the poll + loop — every record it sees was written by an earlier run, which is what + makes "unidentifiable" a statement about a previous run's child rather than + about a job this run just launched. + + Signals nothing. A detached job is meant to survive a restart, so a job that + verifies as ours is left alone and its result is still collected. What is + corrected is the record that would otherwise be signalled later on a pid the + kernel has reassigned: the max-runtime branch of :func:`refresh` sends + SIGTERM then SIGKILL to ``rec["pid"]`` an hour in, and on a reused pid that + lands on a bystander. + + Fail closed: a job that cannot be verified is retired too, not kept. + Retiring loses a result, which is visible; keeping it leaves a pid this + server will eventually signal without knowing what it is pointing at, which + is not. + """ + jobs = _load() + report = {"seen": 0, "retired": 0, "kept": 0} + changed = False + now = time.time() + for rec in jobs.values(): + if rec.get("status") != "running": + continue + report["seen"] += 1 + verdict = process_ownership.verify(rec.get("pid"), rec.get("start_token")) + if verdict in (process_ownership.OWNED, process_ownership.GONE): + # OWNED: still ours, still running, still watched. GONE: refresh() + # already turns an absent process into a "died" record, and it may + # yet find an exit-code file the job wrote before it went. + report["kept"] += 1 + continue + rec["status"] = "failed" + rec["exit_code"] = -1 + rec["ended_at"] = now + rec["ownership_lost"] = verdict + # followed_up stays False: the agent asked for this job and is owed an + # answer, even when the answer is that we lost track of it. + report["retired"] += 1 + changed = True + if changed: + _save(jobs) + return report + + def result_text(rec: Dict[str, Any]) -> str: """Human/agent-readable summary of a finished job, for the follow-up.""" out = _read_output(rec) - if rec.get("killed"): + if rec.get("ownership_lost"): + head = ( + "Background job was abandoned across a server restart: its process " + f"could not be identified as ours ({rec.get('ownership_lost')}), so it was " + "neither waited on nor signalled. Any output below is what it had " + "written by then; if the work matters, re-run it." + ) + elif rec.get("killed"): head = "Background job was killed." elif rec.get("timed_out"): head = f"Background job timed out after {rec.get('max_runtime_s')}s." diff --git a/src/containment.py b/src/containment.py index 1f46f07d5..bd84cc7ff 100644 --- a/src/containment.py +++ b/src/containment.py @@ -66,6 +66,7 @@ from typing import Any, Awaitable, Callable, Mapping, Optional from core.atomic_io import atomic_write_json from core.platform_compat import IS_WINDOWS, find_bash, pid_alive +from src import process_ownership from src.constants import CONTAINMENT_STATE_FILE, MAX_OUTPUT_CHARS logger = logging.getLogger(__name__) @@ -256,6 +257,12 @@ class ReleaseOutcome: 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 { @@ -263,6 +270,7 @@ class ReleaseOutcome: "escalated": self.escalated, "survivors": list(self.survivors), "mechanism": self.mechanism, + "ownership": self.ownership, } @@ -906,7 +914,17 @@ async def run( # exits, getpgid can no longer tell us which group its children are in. pgid = None if IS_WINDOWS else (_pgid_of(proc.pid) or proc.pid) live = replace(grant, pid=proc.pid, pgid=pgid) - _update_record(grant.id, pid=proc.pid, pgid=pgid, started_at=time.time()) + # The start token is what makes this record signallable by a *later* + # process. Without it a restart reaper holds a pid and no way to tell + # whether the pid is still this child or something the kernel has since + # handed to a stranger; see src/process_ownership.py. + _update_record( + grant.id, + pid=proc.pid, + pgid=pgid, + started_at=time.time(), + start_token=process_ownership.capture(proc.pid)["start_token"], + ) out_buf: list[str] = [] err_buf: list[str] = [] @@ -1092,6 +1110,75 @@ def _outcome_for( ) +def _ownership_gate( + grant: ContainmentGrant, + pid: int, + pgid: Optional[int], + token: Optional[str], +) -> Optional[ReleaseOutcome]: + """Decide whether a recovered grant may be signalled at all. + + Returns None to let teardown proceed, or the outcome to report instead. + Reached only for a grant recovered from the durable store — the restart and + reaper path, where the recorded pid is a claim rather than a child this + process is holding. + + The rule is fail-closed: **a signal requires a positive identity.** Anything + else is reported as an undead tree rather than silently killed, because the + alternative is sending SIGKILL to whatever the kernel has since given that + pid to. ODY-86 was this defect; the reason the record stays active on a + refusal is that an unreapable orphan has to remain visible instead of being + closed out as handled. + """ + verdict = process_ownership.verify(pid, token) + if verdict == process_ownership.OWNED: + return None + + if verdict == process_ownership.GONE: + # The leader is gone. Its group may still hold processes it + # backgrounded, but with the leader unverifiable there is nothing left + # to prove the group is still ours, and a recycled group id would mean + # killpg hits strangers. An empty group is the clean case. + if not _group_present(pgid): + return replace( + _outcome_for(grant, dead=True, escalated=False), ownership=verdict, + ) + logger.warning( + "containment: grant %s leader pid %s is gone but group %s still has " + "members; not signalling a group whose ownership cannot be proven", + grant.id, pid, pgid, + ) + return ReleaseOutcome( + dead=False, + escalated=False, + survivors=(pgid,) if pgid else (), + mechanism=grant.mechanism, + ownership=verdict, + ) + + if verdict == process_ownership.FOREIGN: + logger.warning( + "containment: grant %s records pid %s, which now belongs to a " + "different process; refusing to signal it", + grant.id, pid, + ) + else: + logger.warning( + "containment: grant %s pid %s cannot be verified on this host (%s); " + "refusing to signal an unidentified process", + grant.id, pid, process_ownership.inspection_mechanism(), + ) + return ReleaseOutcome( + dead=False, + escalated=False, + # Not ours to enumerate, and listing a foreign pid as a survivor of + # *our* grant would invite the next reaper to kill it. + survivors=(), + mechanism=grant.mechanism, + ownership=verdict, + ) + + def release(grant: ContainmentGrant, *, grace_s: float = 2.0) -> ReleaseOutcome: """Authoritative teardown: signal the group, escalate, then verify. @@ -1106,10 +1193,17 @@ def release(grant: ContainmentGrant, *, grace_s: float = 2.0) -> ReleaseOutcome: otherwise report a tree that is already gone. """ pid, pgid = grant.pid, grant.pgid + # A grant that carries its own pid belongs to the process holding it: this + # caller launched the child and no identity question arises. A grant whose + # pid had to be recovered from the durable store is the restart case, and + # there the pid is a *claim* about a process this run never started. + recovered = pid is None + token: Optional[str] = None if pid is None or (pgid is None and not IS_WINDOWS): record = _load_records().get(grant.id) or {} pid = pid if pid is not None else record.get("pid") pgid = pgid if pgid is not None else record.get("pgid") + token = record.get("start_token") try: pid = int(pid) if pid else 0 except (TypeError, ValueError): @@ -1125,6 +1219,12 @@ def release(grant: ContainmentGrant, *, grace_s: float = 2.0) -> ReleaseOutcome: _finish_release(grant, outcome) return outcome + if recovered: + refusal = _ownership_gate(grant, pid, pgid, token) + if refusal is not None: + _finish_release(grant, refusal) + return refusal + if IS_WINDOWS: try: subprocess.run( @@ -1167,6 +1267,47 @@ def release(grant: ContainmentGrant, *, grace_s: float = 2.0) -> ReleaseOutcome: return outcome +def reap_record(record: Mapping[str, Any], *, grace_s: float = 2.0) -> ReleaseOutcome: + """Tear down a grant known only by its durable record. + + The entry point for a reaper after a restart: the process that acquired the + grant is gone, so there is no :class:`ContainmentGrant` in memory, only the + row :func:`active_grants` returned. Reconstructs the minimum + :func:`release` needs and goes through the same ownership gate — a record is + a claim about a pid, and a reaper is exactly the caller that must not treat + it as more than that. + + ``env`` is not reconstructed because it is never persisted (it is where + credentials live) and teardown does not use it. + """ + record = dict(record or {}) + spec = ContainmentSpec( + workspace=record.get("workspace") or os.getcwd(), + env={}, + wall_clock_s=int(record.get("wall_clock_s") or 1), + required=frozenset(record.get("required") or ()), + ) + grant = ContainmentGrant( + id=str(record.get("id") or ""), + mechanism=str(record.get("mechanism") or "none"), + workspace=spec.workspace, + enforced=frozenset(record.get("enforced") or ()), + degraded=tuple(record.get("degraded") or ()), + unenforced_required=tuple(record.get("unenforced_required") or ()), + owner=str(record.get("owner") or "reaper"), + mode=str(record.get("mode") or CONTAINMENT_MODE), + spec=spec, + external=bool(record.get("external")), + # Left as None on purpose: release() then recovers pid, pgid and the + # start token from the store itself and routes through the ownership + # gate. Passing them here would mark the grant as held in-process and + # skip the very check this path exists to apply. + pid=None, + pgid=None, + ) + return release(grant, grace_s=grace_s) + + async def _release_awaited( grant: ContainmentGrant, proc: "asyncio.subprocess.Process", diff --git a/src/process_ownership.py b/src/process_ownership.py new file mode 100644 index 000000000..d0e814272 --- /dev/null +++ b/src/process_ownership.py @@ -0,0 +1,405 @@ +"""Process identity: is this pid still the process we started? + +A recorded pid is not an identity. The kernel reuses pids, and every store in +this tree that remembers a process — ``data/bg_jobs.json``, +``data/containment_grants.json``, the Cookbook's task list — outlives the +process that wrote it, by design: those records exist so a restart does not lose +a job. The combination is the defect this module closes. A record that says +``pid 4242`` and a live ``pid 4242`` are not the same claim, and signalling the +second because the first was written is how a teardown kills a stranger. + +That is not hypothetical here. ODY-86 was pid files unlinked while the daemons +they named were still live, with ownership never verified; the Cookbook survivor +sweep still terminates *any* process whose command line matches a tracked one, +which is a different spelling of the same mistake. + +**The identity is (pid, start token).** A pid identifies a slot; the start token +identifies which process is occupying it. The kernel will not reissue a pid to a +process that started earlier, so comparing the token recorded at launch with the +token read now answers "is this still ours" without a handle, a lock file or a +supervisor. + +Four verdicts, and the fourth is the point +------------------------------------------ +:data:`OWNED`, :data:`GONE` and :data:`FOREIGN` are the answers. The fourth, +:data:`UNVERIFIABLE`, is what this host could not determine — no procfs, no +``ps``, a probe that raised, or a record written before anything recorded a +token. It is deliberately **not** collapsed into either "ours" (which would +signal strangers) or "gone" (which would abandon live processes). + +Process inspection has broken off Linux four times in this tree — ODY-70, -86, +-94, -99 — every time because an inspection mechanism that was absent read as a +successful answer. So :data:`UNVERIFIABLE` is a containment failure and callers +must treat it as one: do not signal, and do not report a teardown that was not +performed. Refusing to act is the only honest option when you cannot tell what +you would be acting on. + +Token granularity, stated because it bounds the guarantee +--------------------------------------------------------- +======== ============================= =============== +Host Source Resolution +======== ============================= =============== +Linux ``/proc//stat`` field 22 ~10 ms (1 tick) +macOS ``ps -o lstart=`` 1 s +Windows ``GetProcessTimes`` 100 ns +======== ============================= =============== + +A pid recycled *within one token tick* is indistinguishable from the original. +On Linux and Windows that window is too small to hit in practice. On macOS it is +one second, which a pid wrap could theoretically land inside — so the token +narrows the risk by many orders of magnitude there without eliminating it. It is +a strictly better claim than the pid alone, which is the comparison that +matters; it is not a proof of identity and this module does not claim one. +""" + +from __future__ import annotations + +import logging +import os +import shutil +import subprocess +from typing import Any, Iterable, Mapping, NamedTuple, Optional + +from core.platform_compat import IS_WINDOWS, PROC_ROOT, has_procfs + +logger = logging.getLogger(__name__) + + +# ── Verdicts ──────────────────────────────────────────────────────────────── +#: The pid is running and is the same process the token was taken from. +OWNED = "owned" +#: No process holds the pid. Nothing to signal and nothing to reap. +GONE = "gone" +#: A process holds the pid, and it is **not** ours — the pid was recycled. +#: Never signal a foreign pid; that is the defect, not the fix. +FOREIGN = "foreign" +#: This host could not answer. A containment failure, not a default. +UNVERIFIABLE = "unverifiable" + +#: Verdicts that permit a signal. Exactly one. +SIGNALLABLE = frozenset({OWNED}) + + +# ── Inspection mechanisms ─────────────────────────────────────────────────── +MECHANISM_PROCFS = "procfs" +MECHANISM_PS = "ps" +MECHANISM_WIN32 = "win32" +#: No way to inspect processes on this host. Every verdict becomes +#: UNVERIFIABLE, which is the honest answer and not a permissive one. +MECHANISM_NONE = "none" + +# The failure path only: a wedged `ps` must never hold up a teardown decision. +_PS_TIMEOUT_S = 5 + +#: Field 22 of ``/proc//stat`` (1-indexed) is the process start time in +#: clock ticks since boot. Fields 1 and 2 are skipped by splitting on the last +#: ``)`` first, because a comm can itself contain spaces and parentheses. +_PROC_STAT_STARTTIME_INDEX = 19 + + +class InspectionUnavailable(RuntimeError): + """This host offers no way to inspect a process. + + Raised by the probes rather than returned, so a caller that forgets to + handle it fails loudly instead of silently reading an absent mechanism as + "the process is gone". :func:`verify` catches it and reports + :data:`UNVERIFIABLE`. + """ + + def __init__(self, what: str) -> None: + super().__init__(f"process inspection unavailable: cannot read {what}") + self.what = what + + +def inspection_mechanism() -> str: + """Which mechanism this host can answer identity questions with. + + Probed per call rather than cached at import: the tests substitute + ``PROC_ROOT`` to exercise both branches on either kind of host, and a cached + answer would pin whichever host happened to import the module first. + """ + if IS_WINDOWS: + return MECHANISM_WIN32 + if has_procfs(): + return MECHANISM_PROCFS + if shutil.which("ps"): + return MECHANISM_PS + return MECHANISM_NONE + + +def inspection_available() -> bool: + return inspection_mechanism() != MECHANISM_NONE + + +# ── Start tokens ──────────────────────────────────────────────────────────── +def _procfs_token(pid: int) -> Optional[str]: + try: + raw = (PROC_ROOT / str(pid) / "stat").read_text(encoding="utf-8", errors="replace") + except (FileNotFoundError, ProcessLookupError): + return None + except (OSError, PermissionError) as exc: + # The pid exists but is not readable. "I cannot tell" is not "it is + # gone", so this must not return None. + raise InspectionUnavailable(f"/proc/{pid}/stat ({exc})") from exc + # comm is parenthesised and may contain spaces and ')' — split past the last. + _, _, rest = raw.rpartition(")") + fields = rest.split() + try: + return f"procfs:{fields[_PROC_STAT_STARTTIME_INDEX]}" + except IndexError: + raise InspectionUnavailable(f"/proc/{pid}/stat (unexpected layout)") from None + + +def _ps_token(pid: int) -> Optional[str]: + try: + completed = subprocess.run( + ["ps", "-p", str(pid), "-o", "lstart="], + stdout=subprocess.PIPE, + stderr=subprocess.DEVNULL, + timeout=_PS_TIMEOUT_S, + text=True, + ) + except (OSError, subprocess.SubprocessError) as exc: + raise InspectionUnavailable(f"ps -p {pid} ({exc})") from exc + value = (completed.stdout or "").strip() + if completed.returncode != 0: + # ps exits non-zero for a pid that does not exist. With no output that + # is an absent process; with output it is a mechanism that misbehaved. + if not value: + return None + raise InspectionUnavailable(f"ps -p {pid} (exit {completed.returncode})") + if not value: + return None + return f"ps:{' '.join(value.split())}" + + +def _win32_token(pid: int) -> Optional[str]: + import ctypes + from ctypes import wintypes + + PROCESS_QUERY_LIMITED_INFORMATION = 0x1000 + kernel32 = ctypes.windll.kernel32 + handle = kernel32.OpenProcess(PROCESS_QUERY_LIMITED_INFORMATION, False, int(pid)) + if not handle: + return None + try: + creation = wintypes.FILETIME() + exit_time = wintypes.FILETIME() + kernel_time = wintypes.FILETIME() + user_time = wintypes.FILETIME() + ok = kernel32.GetProcessTimes( + handle, + ctypes.byref(creation), + ctypes.byref(exit_time), + ctypes.byref(kernel_time), + ctypes.byref(user_time), + ) + if not ok: + raise InspectionUnavailable(f"GetProcessTimes({pid})") + stamp = (int(creation.dwHighDateTime) << 32) | int(creation.dwLowDateTime) + return f"win32:{stamp}" + finally: + kernel32.CloseHandle(handle) + + +def start_token(pid: Optional[int]) -> Optional[str]: + """An opaque token identifying the process currently holding ``pid``. + + Returns None when no process holds the pid. Raises + :class:`InspectionUnavailable` when this host cannot answer — never a + token, and never None, for a question it could not ask. + + Record this at launch next to the pid. Compare it before signalling. + """ + if not pid: + return None + try: + pid = int(pid) + except (TypeError, ValueError): + return None + if pid <= 0: + return None + mechanism = inspection_mechanism() + if mechanism == MECHANISM_WIN32: + return _win32_token(pid) + if mechanism == MECHANISM_PROCFS: + return _procfs_token(pid) + if mechanism == MECHANISM_PS: + return _ps_token(pid) + raise InspectionUnavailable("process start time on this host") + + +def verify(pid: Optional[int], token: Optional[str]) -> str: + """Is the process now holding ``pid`` the one ``token`` was taken from? + + Returns :data:`OWNED`, :data:`GONE`, :data:`FOREIGN` or + :data:`UNVERIFIABLE`. Only :data:`OWNED` permits a signal. + + A missing or empty ``token`` is :data:`UNVERIFIABLE`, not :data:`OWNED`: + a record that never captured an identity cannot establish one afterwards, + and treating "we did not write it down" as "it is ours" is precisely the + assumption that makes a recycled pid lethal. + """ + if not pid: + return GONE + if not token: + return UNVERIFIABLE + try: + current = start_token(pid) + except InspectionUnavailable as exc: + logger.warning("process_ownership: cannot verify pid %s: %s", pid, exc) + return UNVERIFIABLE + if current is None: + return GONE + return OWNED if current == str(token) else FOREIGN + + +def verify_record( + record: Mapping[str, Any], *, pid_key: str = "pid", token_key: str = "start_token", +) -> str: + """:func:`verify` against a stored record. Convenience for the reaper.""" + return verify((record or {}).get(pid_key), (record or {}).get(token_key)) + + +def capture(pid: Optional[int]) -> dict[str, Any]: + """The identity fields to persist for a process at launch. + + Always returns both keys, with ``start_token`` None when the host could not + produce one, so a record's shape never depends on the host and a later + reader can tell "no token" from "no field". + """ + try: + token = start_token(pid) + except InspectionUnavailable as exc: + logger.warning("process_ownership: launched pid %s without an identity: %s", pid, exc) + token = None + return {"pid": int(pid) if pid else None, "start_token": token} + + +# ── The process table ─────────────────────────────────────────────────────── +class ProcessInfo(NamedTuple): + pid: int + ppid: int + command: str + + +#: Field 4 of ``/proc//stat`` (1-indexed) is the parent pid; it lands at +#: index 1 of the fields that follow the comm's closing paren. +_PROC_STAT_PPID_INDEX = 1 + + +def _procfs_process_table() -> dict[int, ProcessInfo]: + # Guarded here and not only in process_table(): a procfs scan whose + # existence check sits in a caller is one refactor away from being an + # unguarded scan, which is the defect tests/test_procfs_scan_guard.py pins. + if not has_procfs(): + raise InspectionUnavailable(f"the process table via {PROC_ROOT}") + table: dict[int, ProcessInfo] = {} + for entry in os.listdir(PROC_ROOT): + if not entry.isdigit(): + continue + pid = int(entry) + try: + raw = (PROC_ROOT / entry / "cmdline").read_bytes() + command = raw.replace(b"\x00", b" ").decode("utf-8", errors="replace").strip() + except (OSError, PermissionError): + continue + ppid = 0 + try: + stat = (PROC_ROOT / entry / "stat").read_text(encoding="utf-8", errors="replace") + _, _, rest = stat.rpartition(")") + ppid = int(rest.split()[_PROC_STAT_PPID_INDEX]) + except (OSError, PermissionError, IndexError, ValueError): + # A kernel thread or a pid that exited mid-walk. Keeping the row + # with ppid 0 is better than dropping it: a command-line match + # still works, only the descendant walk loses this link. + pass + if command: + table[pid] = ProcessInfo(pid=pid, ppid=ppid, command=command) + return table + + +def _ps_process_table() -> dict[int, ProcessInfo]: + try: + completed = subprocess.run( + # -ww defeats ps's default truncation to terminal width; without it + # a long serve command is clipped and no match can ever be exact. + ["ps", "-axww", "-o", "pid=,ppid=,command="], + stdout=subprocess.PIPE, + stderr=subprocess.DEVNULL, + timeout=_PS_TIMEOUT_S, + text=True, + ) + except (OSError, subprocess.SubprocessError) as exc: + raise InspectionUnavailable(f"ps -axww ({exc})") from exc + if completed.returncode != 0: + raise InspectionUnavailable(f"ps -axww (exit {completed.returncode})") + table: dict[int, ProcessInfo] = {} + for line in (completed.stdout or "").splitlines(): + parts = line.strip().split(None, 2) + if len(parts) < 3 or not parts[0].isdigit() or not parts[1].isdigit(): + continue + command = parts[2].strip() + if command: + pid = int(parts[0]) + table[pid] = ProcessInfo(pid=pid, ppid=int(parts[1]), command=command) + return table + + +def process_table() -> dict[int, ProcessInfo]: + """Every visible process, by pid, with its parent and full command line. + + Raises :class:`InspectionUnavailable` when the host cannot enumerate + processes, so a caller reports that it could not look rather than reporting + that it found nothing. Those are different answers and this tree has + conflated them before (ODY-94). + + ``ps`` covers macOS and the BSDs, which have no procfs to walk — the reason + this exists rather than another ``/proc`` scan. procfs is preferred where + present because it needs no subprocess. + """ + mechanism = inspection_mechanism() + if mechanism == MECHANISM_PROCFS: + return _procfs_process_table() + if mechanism == MECHANISM_PS: + return _ps_process_table() + # Windows: tasklist cannot report a full command line without WMI, and a + # truncated one cannot be matched exactly. Claiming an empty table would + # read as "no survivors". + raise InspectionUnavailable(f"the process table via {mechanism}") + + +def command_lines() -> dict[int, str]: + """Every visible pid mapped to its full command line.""" + return {pid: info.command for pid, info in process_table().items()} + + +def descendants( + roots: "Iterable[int]", *, table: Optional[Mapping[int, ProcessInfo]] = None, +) -> list[int]: + """Every process under ``roots``, roots included, breadth-first. + + The point of taking several roots and one table is that the answer is a + *snapshot*: walking the tree one subprocess call at a time lets a child be + reparented between calls and vanish from the result. Callers that need to + act on a tree should capture it once, before they start tearing it down. + + A pid that is its own parent, or a cycle the table reports, terminates the + walk rather than looping. + """ + rows = dict(table) if table is not None else process_table() + children: dict[int, list[int]] = {} + for info in rows.values(): + children.setdefault(info.ppid, []).append(info.pid) + + found: list[int] = [] + seen: set[int] = set() + queue = [int(root) for root in roots if root] + while queue: + pid = queue.pop(0) + if pid in seen: + continue + seen.add(pid) + found.append(pid) + queue.extend(child for child in children.get(pid, ()) if child not in seen) + return found diff --git a/src/process_reaper.py b/src/process_reaper.py new file mode 100644 index 000000000..63a2dbff7 --- /dev/null +++ b/src/process_reaper.py @@ -0,0 +1,158 @@ +"""Startup reconciliation for processes a previous run left behind. + +Two stores in this tree outlive the process that wrote them, on purpose: +``data/containment_grants.json`` so a restart can reap rather than orphan, and +``data/bg_jobs.json`` so a restart never loses a detached job or its result. +Until now nothing read either of them at startup. A crashed or restarted server +therefore left every grant permanently "active" and every background job +permanently "running", and the first thing to touch one of those records was a +teardown aimed at a pid that had been reassigned in the meantime. + +This module runs once, during startup, before anything of this run exists. That +timing is what makes its rules safe: every record it sees was written by an +earlier run, so "I cannot identify this process" is information about a previous +run's child and not about one of ours. + +The two stores get **opposite** treatment, which is the whole reason this is a +module and not a loop: + +* A **containment grant** is tied to a tool call that no longer has a caller. + A live process under an abandoned grant is by definition an orphan, so it is + torn down. +* A **background job** is detached deliberately and is documented to survive a + uvicorn restart. Killing one here would break the feature, so its record is + only corrected, never reaped. What gets fixed is identity: a job whose pid now + belongs to someone else is retired so that nothing later signals the stranger. + +Fail closed in both: a signal requires a positive identity from +:mod:`src.process_ownership`, and every other verdict is recorded rather than +acted on. Containment that cannot identify its target is not containment, and +the honest failure is a visible orphan rather than a dead bystander. +""" + +from __future__ import annotations + +import logging +from typing import Any, Dict + +from src import process_ownership + +logger = logging.getLogger(__name__) + + +def reap_containment_grants() -> Dict[str, Any]: + """Tear down or retire every grant a previous run left active. + + Per grant: a verified live process is torn down through + :func:`src.containment.reap_record`; a grant whose process is gone is + dropped; a grant naming a pid that is now someone else's is dropped + *without a signal*, because the only thing left to do with it is stop + believing it. A grant that cannot be verified at all is **kept**, so the + orphan stays visible in ``active_grants()`` instead of being quietly + written off as handled. + """ + from src import containment + + report: Dict[str, Any] = { + "seen": 0, "torn_down": 0, "already_gone": 0, + "foreign": 0, "unverifiable": 0, "failed": 0, + } + try: + records = containment.active_grants() + except Exception: + logger.warning("process_reaper: containment grant store unreadable", exc_info=True) + return report + + for record in records: + report["seen"] += 1 + grant_id = str(record.get("id") or "") + if record.get("external"): + # Nothing local ever ran, so there is nothing local to reap. + containment.forget(grant_id) + report["already_gone"] += 1 + continue + verdict = process_ownership.verify_record(record) + if verdict == process_ownership.GONE: + containment.forget(grant_id) + report["already_gone"] += 1 + continue + if verdict == process_ownership.FOREIGN: + logger.warning( + "process_reaper: grant %s named pid %s, which now belongs to a " + "different process; dropping the record unsignalled", + grant_id, record.get("pid"), + ) + containment.forget(grant_id) + report["foreign"] += 1 + continue + if verdict == process_ownership.UNVERIFIABLE: + logger.error( + "process_reaper: grant %s (pid %s, owner %s) cannot be verified " + "via %s; leaving it active and unsignalled — this is a " + "containment failure, not a clean start", + grant_id, record.get("pid"), record.get("owner"), + process_ownership.inspection_mechanism(), + ) + report["unverifiable"] += 1 + continue + try: + outcome = containment.reap_record(record) + except Exception: + logger.warning("process_reaper: tearing down grant %s failed", grant_id, exc_info=True) + report["failed"] += 1 + continue + if outcome.dead: + containment.forget(grant_id) + report["torn_down"] += 1 + else: + logger.error( + "process_reaper: grant %s survived teardown; survivors=%s", + grant_id, list(outcome.survivors), + ) + report["failed"] += 1 + return report + + +def reap_bg_jobs() -> Dict[str, Any]: + """Correct the identity of background jobs a previous run launched. + + Deliberately kills nothing: a ``#!bg`` job is detached so that it outlives + the request *and* the server, and the store exists so its result is still + collected afterwards. The defect being closed is narrower — a record whose + pid has been reassigned will be signalled by the max-runtime reaper an hour + later, and that signal lands on whatever now holds the pid. + """ + from src import bg_jobs + + try: + return bg_jobs.disown_unverified() + except Exception: + logger.warning("process_reaper: background job store unreadable", exc_info=True) + return {"seen": 0, "retired": 0, "kept": 0} + + +def reap_orphans() -> Dict[str, Any]: + """Run both reconciliations. Returns a report; raises nothing. + + Blocking: a teardown escalates SIGTERM → grace → SIGKILL and waits for the + process to actually go. Call it off the event loop. + """ + report = { + "mechanism": process_ownership.inspection_mechanism(), + "grants": reap_containment_grants(), + "bg_jobs": reap_bg_jobs(), + } + if report["mechanism"] == process_ownership.MECHANISM_NONE: + logger.error( + "process_reaper: this host offers no process inspection; no orphan " + "from a previous run can be identified or reaped" + ) + logger.info("process_reaper: startup reconciliation %s", report) + return report + + +async def reap_orphans_at_startup() -> Dict[str, Any]: + """:func:`reap_orphans` off the event loop, for an app startup task.""" + import asyncio + + return await asyncio.to_thread(reap_orphans) diff --git a/src/tools/cookbook.py b/src/tools/cookbook.py index c2def959b..48a70a53c 100644 --- a/src/tools/cookbook.py +++ b/src/tools/cookbook.py @@ -13,6 +13,7 @@ import asyncio import contextlib import json import logging +import os import re from typing import Any, Dict, List, Optional @@ -1024,6 +1025,189 @@ async def do_list_served_models(content: str, owner: Optional[str] = None) -> Di return {"output": "\n".join(lines), "tasks": merged, "exit_code": 0} +# How long a tmux query may take before the stop gives up on identifying the +# session's processes and says so. The kill itself does not depend on it. +_TMUX_QUERY_TIMEOUT_S = 5 +# Grace between SIGTERM and SIGKILL for a model server that ignored SIGHUP. +_SWEEP_GRACE_S = 2.0 +_SWEEP_POLL_S = 0.05 + + +async def _capture_session_processes(session_id: str) -> tuple[List[Dict[str, Any]], str]: + """Snapshot the processes belonging to a local tmux session. + + This is what gives the survivor sweep an *ownership* record rather than a + resemblance. tmux knows which pane hosts the session, the pane pid's + descendants are the processes that session started, and a start token taken + now is what lets the sweep prove, after the kill, that a pid it is about to + signal is still one of them. + + Returns ``(records, note)``. An empty list with a note is the honest + outcome when the session cannot be enumerated — the note reaches the tool + result, because "I found no survivors" and "I could not look" are different + answers and the sweep used to give the first for both (ODY-94). + """ + from src import process_ownership + + try: + proc = await asyncio.create_subprocess_exec( + "tmux", "list-panes", "-a", "-F", "#{session_name} #{pane_pid}", + stdout=asyncio.subprocess.PIPE, + stderr=asyncio.subprocess.PIPE, + ) + try: + stdout, _stderr = await asyncio.wait_for( + proc.communicate(), timeout=_TMUX_QUERY_TIMEOUT_S + ) + except asyncio.TimeoutError: + with contextlib.suppress(Exception): + proc.kill() + await proc.communicate() + return [], "; could not identify the session's processes (tmux timed out)" + except (OSError, FileNotFoundError) as exc: + return [], f"; could not identify the session's processes (tmux unavailable: {exc})" + if proc.returncode not in (0, None): + return [], "; could not identify the session's processes (tmux listed no panes)" + + pane_pids: List[int] = [] + for line in (stdout or b"").decode("utf-8", errors="replace").splitlines(): + name, _, pid_text = line.strip().rpartition(" ") + if name == session_id and pid_text.isdigit(): + pane_pids.append(int(pid_text)) + if not pane_pids: + # The session is already gone, so nothing links a survivor to it. Said + # out loud rather than reported as a clean sweep. + return [], "; the session had no live pane, so its processes could not be identified" + + def _snapshot() -> tuple[List[Dict[str, Any]], str]: + try: + table = process_ownership.process_table() + 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 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: + """Terminate the captured processes that outlived the tmux kill. + + Blocking; call it off the event loop. Returns the note to append to the + tool result — the sweep's outcome is part of whether the stop worked, and + silence here is what let a half-stopped server read as stopped. + + Only captured pids are signalled. A process that merely matches + ``tracked_cmd`` is reported and left alone: the Cookbook composed that + 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 + + 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) + + note = "" + if killed: + note += f"; killed {len(killed)} surviving process(es) owned by the session" + if survivors: + note += ( + f"; {len(survivors)} process(es) survived SIGKILL and are still " + f"running (pid {', '.join(str(pid) for pid in survivors)})" + ) + note += _unowned_match_note(tracked_cmd, {rec["pid"] for rec in owned}) + return note + + +def _unowned_match_note(tracked_cmd: str, owned_pids: set) -> str: + """Report, without signalling, processes that look like the tracked command. + + The old sweep killed these. It could not tell them apart from the server it + started, and neither can this — so it names them instead. Reporting keeps + the information the old behaviour acted on while giving up the one thing it + was never entitled to do. + """ + from src import process_ownership + + if not tracked_cmd: + return "" + try: + table = process_ownership.process_table() + except process_ownership.InspectionUnavailable: + return "; could not check for unowned processes matching the command" + strangers = sorted( + pid for pid, info in table.items() + if info.command == tracked_cmd and pid not in owned_pids and pid != os.getpid() + ) + if not strangers: + return "" + return ( + f"; note: {len(strangers)} other process(es) match this server's command " + f"line (pid {', '.join(str(pid) for pid in strangers)}) — not signalled, " + f"because nothing identifies them as started by this session" + ) + + async def _cookbook_kill_session(session_id: str, *, remote_host: str = "", ssh_port: str = "", verb: str = "Stopped") -> Dict: """Kill a cookbook tmux session — remote-aware — AND mark the task @@ -1077,6 +1261,16 @@ async def _cookbook_kill_session(session_id: str, *, remote_host: str = "", cmd = f"tmux kill-session -t {shlex.quote(session_id)}" target_label = session_id + # Capture what this session owns BEFORE the kill. Once tmux tears the + # session down the pane is gone, and with it the only evidence linking a + # surviving model server to the session that started it. A sweep that looks + # afterwards has nothing left but the command line, which identifies a + # *kind* of process and not one we started. + owned: List[Dict[str, Any]] = [] + owned_note = "" + if not remote and isinstance(matched, dict): + owned, owned_note = await _capture_session_processes(session_id) + try: if remote: async with httpx.AsyncClient(timeout=15) as client: @@ -1117,37 +1311,16 @@ async def _cookbook_kill_session(session_id: str, *, remote_host: str = "", if kill_failed and not already_gone: return {"error": f"Failed to {verb.lower()} {target_label}: {kill_err or 'kill-session returned non-zero'}", "exit_code": 1} - # Some model servers survive the tmux session's SIGHUP. For local - # tracked tasks only, terminate processes whose full command line - # exactly matches the command saved by the Cookbook launcher. + # Some model servers survive the tmux session's SIGHUP. Terminate the + # ones this session actually owns — captured above, each verified by + # identity at signal time — and report, without signalling, anything + # that merely looks like the tracked command. + sweep_note = "" if not remote and isinstance(matched, dict): - import os - import signal tracked_cmd = str((matched.get("payload") or {}).get("_cmd") or "").strip() - matched_pids: list[int] = [] - # No procfs means no way to match a survivor by its command line. - # The tmux kill above already stopped the session, so skip the - # sweep instead of failing a stop that worked. - if tracked_cmd and platform_compat.has_procfs(): - proc_root = platform_compat.PROC_ROOT - for pid_name in os.listdir(proc_root): - if not pid_name.isdigit() or int(pid_name) == os.getpid(): - continue - try: - raw = (proc_root / pid_name / "cmdline").read_bytes() - process_cmd = raw.replace(b"\x00", b" ").decode("utf-8", errors="replace").strip() - except (OSError, PermissionError): - continue - if process_cmd == tracked_cmd: - matched_pids.append(int(pid_name)) - with contextlib.suppress(ProcessLookupError, PermissionError): - os.kill(int(pid_name), signal.SIGTERM) - if matched_pids: - await asyncio.sleep(0.5) - for pid in matched_pids: - with contextlib.suppress(ProcessLookupError, PermissionError): - os.kill(pid, 0) - os.kill(pid, signal.SIGKILL) + sweep_note = await asyncio.to_thread( + _sweep_session_survivors, owned, tracked_cmd, owned_note, + ) # Update state: mark stopped (so the UI + list reflect reality). if matched is not None: @@ -1160,7 +1333,7 @@ async def _cookbook_kill_session(session_id: str, *, remote_host: str = "", logger.debug(f"failed to mark {session_id} stopped in state: {e}") suffix = " (was already gone)" if already_gone else "" - return {"output": f"{verb} {target_label}{suffix}", "exit_code": 0} + return {"output": f"{verb} {target_label}{suffix}{sweep_note}", "exit_code": 0} except Exception as e: return {"error": str(e), "exit_code": 1} diff --git a/tests/test_cookbook_stop_without_procfs.py b/tests/test_cookbook_stop_without_procfs.py index ad6e225a8..44f7392ce 100644 --- a/tests/test_cookbook_stop_without_procfs.py +++ b/tests/test_cookbook_stop_without_procfs.py @@ -1,13 +1,26 @@ -"""Stopping a Cookbook server must succeed on a host with no procfs. +"""Stopping a Cookbook server, on a host with procfs and on one without. -The tmux kill is what actually stops the server; the pid sweep that follows -it only catches model servers that survive the session's SIGHUP. On macOS and -Windows there is no ``/proc`` to sweep, and letting that raise turned a -successful stop into a reported failure *and* skipped the state write that -marks the session stopped for the Cookbook UI. +The tmux kill is what actually stops the server; the pid sweep that follows it +only catches model servers that survive the session's SIGHUP. Two invariants +live here. + +**The stop must not fail because the host cannot be inspected.** Letting a +procfs scan raise on macOS turned a successful stop into a reported failure and +skipped the state write that marks the session stopped for the Cookbook UI +(ODY-94). Skipping the sweep silently fixed the crash and left the other half: +the stop then claimed success without having looked at all. So the sweep now +runs through ``ps`` where there is no procfs, and says so when it cannot look. + +**The sweep signals only processes the session owns.** It used to kill anything +whose full command line matched the tracked one. The Cookbook composed that +command line, so an identical one is just as likely to be a server the user +started by hand — killing it is indistinguishable from killing ours, which is +the "stop only what we started" failure. Ownership now comes from the tmux +pane's process tree, captured before the kill; a lookalike is reported instead. """ import asyncio import json +import os import signal import pytest @@ -67,20 +80,34 @@ def _install_httpx_client(monkeypatch, state): return posts -def _install_successful_tmux_kill(monkeypatch): - """Replace the real ``tmux kill-session`` with a process that succeeds.""" +def _install_successful_tmux_kill(monkeypatch, panes=""): + """Fake the two tmux calls a stop makes: list-panes, then kill-session. + + ``panes`` is the ``list-panes`` stdout, i.e. ``" "`` per + line — the stop reads it to learn which processes the session owns before + the kill destroys that link. + """ + calls = [] class FakeProc: returncode = 0 + def __init__(self, stdout=b""): + self._stdout = stdout + async def communicate(self): - return b"", b"" + return self._stdout, b"" async def fake_exec(*argv, **kwargs): - assert argv[:2] == ("tmux", "kill-session") + calls.append(argv) + assert argv[0] == "tmux" + if argv[1] == "list-panes": + return FakeProc(panes.encode()) + assert argv[1] == "kill-session" return FakeProc() monkeypatch.setattr(asyncio, "create_subprocess_exec", fake_exec) + return calls def _stopped_statuses(posts, session_id): @@ -92,63 +119,192 @@ def _stopped_statuses(posts, session_id): return out +def _fake_table(monkeypatch, rows): + """Substitute the process table. ``rows`` is {pid: (ppid, command)}. + + Returns the live dict, so a test can model a process actually dying by + removing it: ``start_token`` reads from the same dict, and a pid that is no + longer in it has no token, which :func:`process_ownership.verify` reports as + ``GONE``. + """ + from src import process_ownership + + table = { + pid: process_ownership.ProcessInfo(pid=pid, ppid=ppid, command=command) + for pid, (ppid, command) in rows.items() + } + monkeypatch.setattr(process_ownership, "process_table", lambda: dict(table)) + # Identity is what authorises a signal, so every pid in the fake table has + # one. A pid absent from the table has no token and cannot be signalled. + monkeypatch.setattr( + process_ownership, "start_token", + lambda pid: f"token:{pid}" if int(pid or 0) in table else None, + ) + return table + + +def _install_effective_kill(monkeypatch, table): + """Record signals, and let SIGTERM actually remove the process. + + Keeps the sweep off its escalation path, which would otherwise spend the + full SIGTERM grace plus the SIGKILL confirmation window on every pid. + """ + signalled = [] + + def _kill(pid, sig): + signalled.append((pid, sig)) + table.pop(int(pid), None) + + monkeypatch.setattr(os, "kill", _kill) + return signalled + + @pytest.mark.asyncio async def test_stop_marks_session_stopped_when_the_host_has_no_procfs( monkeypatch, tmp_path ): + """The ODY-94 regression: no procfs must not turn a working stop into a failure.""" state = _tracked_state() posts = _install_httpx_client(monkeypatch, state) - _install_successful_tmux_kill(monkeypatch) + _install_successful_tmux_kill(monkeypatch, panes="serve-abc123 900\n") monkeypatch.setattr(platform_compat, "PROC_ROOT", tmp_path / "no-procfs") - - import os - - def _unexpected_listdir(*args, **kwargs): - raise AssertionError("the pid sweep must not run without procfs") - - monkeypatch.setattr(os, "listdir", _unexpected_listdir) + # ps is the mechanism on a procfs-less host; the sweep goes through it + # instead of being skipped. + _fake_table(monkeypatch, {900: (1, "bash")}) result = await tools.do_stop_served_model( json.dumps({"session_id": "serve-abc123"}) ) - assert result == {"output": "Stopped server serve-abc123", "exit_code": 0} + assert result["exit_code"] == 0 + assert result["output"].startswith("Stopped server serve-abc123") assert _stopped_statuses(posts, "serve-abc123") == ["stopped"] @pytest.mark.asyncio -async def test_stop_sweeps_surviving_pids_when_procfs_is_present( +async def test_stop_says_so_when_the_session_cannot_be_inspected( monkeypatch, tmp_path ): - tracked_cmd = "python -m vllm.entrypoints.openai.api_server --model org/model" - state = _tracked_state(cmd=tracked_cmd) + """A sweep that could not look must not read as a sweep that found nothing. + + This is the half of ODY-94 that the procfs guard left behind: skipping the + sweep stopped the crash and still reported plain success. + """ + from src import process_ownership + + state = _tracked_state() posts = _install_httpx_client(monkeypatch, state) - _install_successful_tmux_kill(monkeypatch) + _install_successful_tmux_kill(monkeypatch, panes="serve-abc123 900\n") + monkeypatch.setattr(platform_compat, "PROC_ROOT", tmp_path / "no-procfs") - proc = tmp_path / "proc" + def _no_inspection(): + raise process_ownership.InspectionUnavailable("the process table") - def _write_pid(pid, cmdline): - entry = proc / pid - entry.mkdir(parents=True) - (entry / "cmdline").write_bytes(cmdline.replace(" ", "\0").encode()) - - _write_pid("101", tracked_cmd) - _write_pid("202", "python -m http.server") - (proc / "self").mkdir() - monkeypatch.setattr(platform_compat, "PROC_ROOT", proc) + monkeypatch.setattr(process_ownership, "process_table", _no_inspection) signalled = [] - import os - monkeypatch.setattr(os, "kill", lambda pid, sig: signalled.append((pid, sig))) result = await tools.do_stop_served_model( json.dumps({"session_id": "serve-abc123"}) ) + assert result["exit_code"] == 0 + assert "could not identify the session's processes" in result["output"] + assert signalled == [] + assert _stopped_statuses(posts, "serve-abc123") == ["stopped"] + + +@pytest.mark.asyncio +async def test_stop_kills_the_sessions_own_survivor(monkeypatch, tmp_path): + """A process under the session's pane is ours, so it gets signalled.""" + 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"), # the pane shell + 101: (900, tracked_cmd), # the model server it started — ours + }) + 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 (101, signal.SIGTERM) in signalled + assert "killed 2 surviving process(es)" in result["output"] + assert _stopped_statuses(posts, "serve-abc123") == ["stopped"] + + +@pytest.mark.asyncio +async def test_stop_reports_a_command_line_lookalike_without_signalling_it( + monkeypatch, tmp_path +): + """The headline change: matching the command line is not owning the process. + + pid 202 runs exactly the tracked command but descends from nothing this + session started — a server the user launched by hand looks precisely like + this. The old sweep killed it. + """ + 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"), + 202: (1, tracked_cmd), # same command, different lineage + }) + 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 == 202 for pid, _sig in signalled) + # Reported rather than silently dropped: the old behaviour acted on this + # information, so giving it up entirely would be a regression of its own. + assert "202" in result["output"] + assert "not signalled" in result["output"] + assert _stopped_statuses(posts, "serve-abc123") == ["stopped"] + + +@pytest.mark.asyncio +async def test_stop_does_not_signal_a_pid_whose_identity_changed( + monkeypatch, tmp_path +): + """Captured before the kill, recycled before the sweep: do not signal it.""" + 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)}) + + # pid 101's slot reads differently every time it is asked, so whatever the + # capture recorded, the sweep's re-check cannot match it: the pid was + # recycled in between. Every other pid keeps a stable identity. + drift = {"n": 0} + + def _drifting_token(pid): + if int(pid) == 101: + drift["n"] += 1 + return f"token:101:{drift['n']}" + return f"token:{pid}" if int(pid or 0) in table else None + + monkeypatch.setattr(process_ownership, "start_token", _drifting_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 + # The pane shell is genuinely ours and is signalled; 101 never is. + assert not any(pid == 101 for pid, _sig in signalled) assert _stopped_statuses(posts, "serve-abc123") == ["stopped"] diff --git a/tests/test_orphan_reaping.py b/tests/test_orphan_reaping.py new file mode 100644 index 000000000..d89b45810 --- /dev/null +++ b/tests/test_orphan_reaping.py @@ -0,0 +1,412 @@ +"""Teardown across a restart: the pid in a store is a claim, not a handle. + +Three stores here outlive the process that wrote them, deliberately — a restart +is supposed to keep a background job and its result. The consequence nobody had +closed is that the recorded pid is reassignable, so a teardown driven off an old +record can land on a process the kernel has since given to somebody else. That +was ODY-86's shape. + +Covered: + +* :func:`src.containment.release` gating a grant it recovered from the store, + and *not* gating one whose process the caller is holding. +* :func:`src.containment.reap_record`, the entry point for a reaper that has a + row and no grant object. +* :mod:`src.process_reaper`, which gives the two stores opposite treatment — + orphaned grants are torn down, detached jobs are only corrected. +* :func:`src.bg_jobs.disown_unverified`. + +The verdicts come from :mod:`src.process_ownership`, substituted here so each +case is driven exactly; that module's own tests pin it against real processes. +""" + +import json + +import pytest + +from src import bg_jobs, containment, process_ownership, process_reaper + + +@pytest.fixture +def grant_store(tmp_path, monkeypatch): + """Redirect the containment grant store. Returns a reader for it.""" + path = tmp_path / "containment_grants.json" + monkeypatch.setattr(containment, "_store_path", lambda: path) + + def _read(): + return json.loads(path.read_text(encoding="utf-8")) if path.exists() else {} + + return _read + + +@pytest.fixture +def job_store(tmp_path, monkeypatch): + """Redirect the background-job store and its spool directory.""" + monkeypatch.setattr(bg_jobs, "_STORE", tmp_path / "bg_jobs.json") + monkeypatch.setattr(bg_jobs, "_JOBS_DIR", tmp_path / "bg_jobs") + (tmp_path / "bg_jobs").mkdir() + + +def verdicts(monkeypatch, mapping, default=process_ownership.OWNED): + """Pin verify()'s answer per pid.""" + monkeypatch.setattr( + process_ownership, "verify", + lambda pid, _token: mapping.get(int(pid or 0), default), + ) + + +def seed_grant(pid=4242, pgid=4242, token="token:4242", **extra): + """Write a grant record the way a previous run would have left it.""" + record = { + "id": "grant-1", "owner": "session-7", "mechanism": "process_group", + "mode": containment.CONTAINMENT_MODE, "workspace": "/tmp", + "enforced": ["filesystem", "process_tree", "wall_clock"], + "degraded": [], "unenforced_required": [], + "required": ["filesystem", "process_tree", "wall_clock"], + "wall_clock_s": 60, "max_memory_bytes": None, "max_processes": None, + "network": "inherit", "external": False, "pid": pid, "pgid": pgid, + "start_token": token, "acquired_at": 0.0, "released_at": None, "release": None, + } + record.update(extra) + containment._save_records({record["id"]: record}) + return record + + +def seed_job(pid=4242, token="token:4242", status="running", **extra): + record = { + "id": "job-1", "session_id": "chat-1", "command": "sleep 300", + "status": status, "pid": pid, "start_token": token, "started_at": 0.0, + "ended_at": None, "exit_code": None, "max_runtime_s": 3600, + "followed_up": False, "log_path": "", "exit_path": "", + } + record.update(extra) + bg_jobs._save({record["id"]: record}) + return record + + +# ── The gate inside release() ─────────────────────────────────────────────── +def test_a_recovered_grant_naming_a_recycled_pid_is_not_signalled( + grant_store, monkeypatch +): + """The headline case. The pid is live, and it is not ours.""" + seed_grant() + verdicts(monkeypatch, {4242: process_ownership.FOREIGN}) + + def _no_signals(*_args, **_kwargs): + raise AssertionError("a foreign pid must never be signalled") + + monkeypatch.setattr(containment, "_signal_tree", _no_signals) + + outcome = containment.reap_record(grant_store()["grant-1"]) + + assert outcome.ownership == process_ownership.FOREIGN + assert outcome.dead is False + # Not listed as a survivor of *our* grant either: naming a stranger's pid + # there invites the next reaper to kill it. + assert outcome.survivors == () + + +def test_an_unverifiable_grant_is_not_signalled_and_stays_active( + grant_store, monkeypatch +): + """An inspection mechanism this host does not have is a containment failure. + + Reported as an undead tree and left in the store, so the orphan stays + visible in ``active_grants()`` rather than being written off as handled. + """ + seed_grant() + verdicts(monkeypatch, {4242: process_ownership.UNVERIFIABLE}) + monkeypatch.setattr( + containment, "_signal_tree", + lambda *_a, **_k: pytest.fail("an unidentified pid must never be signalled"), + ) + + outcome = containment.reap_record(grant_store()["grant-1"]) + + assert outcome.ownership == process_ownership.UNVERIFIABLE + assert outcome.dead is False + assert containment.active_grants(), "the orphan must remain visible" + + +def test_a_recovered_grant_whose_process_is_gone_is_released_clean( + grant_store, monkeypatch +): + seed_grant() + verdicts(monkeypatch, {4242: process_ownership.GONE}) + monkeypatch.setattr(containment, "_group_present", lambda _pgid: False) + + outcome = containment.reap_record(grant_store()["grant-1"]) + + assert outcome.dead is True + assert outcome.ownership == process_ownership.GONE + assert containment.active_grants() == [] + + +def test_a_gone_leader_with_a_live_group_is_reported_not_killed( + grant_store, monkeypatch +): + """Children outlive the leader, but with the leader gone nothing proves the + group is still ours — and a recycled group id would mean killpg hits + strangers. The orphan is reported instead of guessed at.""" + seed_grant() + verdicts(monkeypatch, {4242: process_ownership.GONE}) + monkeypatch.setattr(containment, "_group_present", lambda _pgid: True) + monkeypatch.setattr( + containment, "_signal_tree", + lambda *_a, **_k: pytest.fail("an unprovable group must not be signalled"), + ) + + outcome = containment.reap_record(grant_store()["grant-1"]) + + assert outcome.dead is False + assert outcome.survivors == (4242,) + + +def test_a_verified_grant_is_torn_down_normally(grant_store, monkeypatch): + seed_grant() + verdicts(monkeypatch, {4242: process_ownership.OWNED}) + signals = [] + monkeypatch.setattr( + containment, "_signal_tree", + lambda pid, pgid, sig: signals.append((pid, pgid, sig)), + ) + # Dead on the first probe, so the teardown does not wait out its grace. + monkeypatch.setattr(containment, "_tree_gone", lambda *_a, **_k: True) + + outcome = containment.reap_record(grant_store()["grant-1"]) + + assert outcome.dead is True + assert outcome.ownership == "" + + +def test_an_in_process_grant_is_not_subjected_to_the_gate(monkeypatch, tmp_path): + """A grant carrying its own pid belongs to the caller holding it. + + The caller launched the child, so there is no identity question — and + demanding a token here would refuse teardown of a perfectly ordinary tool + call on a host with no inspection mechanism. + """ + monkeypatch.setattr(containment, "_store_path", lambda: tmp_path / "grants.json") + monkeypatch.setattr( + process_ownership, "verify", + lambda *_a, **_k: pytest.fail("an in-process teardown must not consult ownership"), + ) + spec = containment.ContainmentSpec( + workspace=str(tmp_path), env={}, wall_clock_s=5, + ) + grant = containment.ContainmentGrant( + id="live-1", mechanism="process_group", workspace=str(tmp_path), + enforced=containment.DEFAULT_REQUIRED, degraded=(), unenforced_required=(), + owner="session-7", mode=containment.CONTAINMENT_MODE, spec=spec, + pid=4242, pgid=4242, + ) + monkeypatch.setattr(containment, "_tree_gone", lambda *_a, **_k: True) + monkeypatch.setattr(containment, "_group_present", lambda _pgid: False) + + outcome = containment.release(grant) + + assert outcome.dead is True + assert outcome.ownership == "" + + +def test_the_release_block_names_the_ownership_verdict(grant_store, monkeypatch): + """The verdict reaches the record, so "why is this still here" is answerable.""" + seed_grant() + verdicts(monkeypatch, {4242: process_ownership.FOREIGN}) + + outcome = containment.reap_record(grant_store()["grant-1"]) + + assert outcome.to_dict()["ownership"] == process_ownership.FOREIGN + assert grant_store()["grant-1"]["release"]["ownership"] == process_ownership.FOREIGN + + +# ── The reaper ────────────────────────────────────────────────────────────── +def test_the_reaper_drops_a_foreign_grant_without_signalling_it( + grant_store, job_store, monkeypatch +): + seed_grant() + verdicts(monkeypatch, {4242: process_ownership.FOREIGN}) + monkeypatch.setattr( + containment, "_signal_tree", + lambda *_a, **_k: pytest.fail("the reaper must not signal a foreign pid"), + ) + + report = process_reaper.reap_containment_grants() + + assert report["foreign"] == 1 + # Dropped rather than retried: the only thing left to do with a record + # about someone else's process is stop believing it. + assert containment.active_grants() == [] + + +def test_the_reaper_keeps_an_unverifiable_grant_visible( + grant_store, job_store, monkeypatch +): + seed_grant() + verdicts(monkeypatch, {4242: process_ownership.UNVERIFIABLE}) + + report = process_reaper.reap_containment_grants() + + assert report["unverifiable"] == 1 + assert len(containment.active_grants()) == 1 + + +def test_the_reaper_tears_down_a_verified_orphan(grant_store, job_store, monkeypatch): + """A live process under an abandoned grant has no caller left. It goes.""" + seed_grant() + verdicts(monkeypatch, {4242: process_ownership.OWNED}) + torn_down = [] + + def _reap(record, **_kwargs): + torn_down.append(record["id"]) + return containment.ReleaseOutcome(dead=True, escalated=True, mechanism="process_group") + + monkeypatch.setattr(containment, "reap_record", _reap) + + report = process_reaper.reap_containment_grants() + + assert torn_down == ["grant-1"] + assert report["torn_down"] == 1 + assert containment.active_grants() == [] + + +def test_the_reaper_keeps_a_grant_that_survived_its_teardown( + grant_store, job_store, monkeypatch +): + seed_grant() + verdicts(monkeypatch, {4242: process_ownership.OWNED}) + monkeypatch.setattr( + containment, "reap_record", + lambda record, **_k: containment.ReleaseOutcome( + dead=False, escalated=True, survivors=(4242,), mechanism="process_group", + ), + ) + + report = process_reaper.reap_containment_grants() + + assert report["failed"] == 1 + assert len(containment.active_grants()) == 1 + + +def test_the_reaper_forgets_an_external_grant_without_inspecting_anything( + grant_store, job_store, monkeypatch +): + """Nothing local ever ran, so there is nothing local to reap.""" + seed_grant(external=True) + monkeypatch.setattr( + process_ownership, "verify", + lambda *_a, **_k: pytest.fail("an external grant has no local pid to verify"), + ) + + report = process_reaper.reap_containment_grants() + + assert report["already_gone"] == 1 + assert containment.active_grants() == [] + + +def test_the_reaper_survives_an_unreadable_store(monkeypatch, job_store): + monkeypatch.setattr( + containment, "active_grants", + lambda: (_ for _ in ()).throw(RuntimeError("store on fire")), + ) + + assert process_reaper.reap_containment_grants()["seen"] == 0 + + +# ── Background jobs: corrected, never killed ──────────────────────────────── +def test_a_job_whose_pid_was_reassigned_is_retired_unsignalled( + job_store, monkeypatch +): + """Left alone, refresh() would SIGKILL this pid at max-runtime. + + An hour after a restart, aimed at whatever now holds it. + """ + seed_job() + verdicts(monkeypatch, {4242: process_ownership.FOREIGN}) + monkeypatch.setattr( + bg_jobs, "_kill", + lambda *_a, **_k: pytest.fail("disowning a job must not signal anything"), + ) + + report = bg_jobs.disown_unverified() + + assert report == {"seen": 1, "retired": 1, "kept": 0} + record = bg_jobs._load()["job-1"] + assert record["status"] == "failed" + assert record["ownership_lost"] == process_ownership.FOREIGN + # The agent asked for this job and is still owed an answer. + assert record["followed_up"] is False + + +def test_an_unverifiable_job_is_also_retired(job_store, monkeypatch): + """Fail closed. Retiring loses a result, which is visible; keeping it leaves + a pid this server will later signal without knowing what it points at.""" + seed_job(token=None) + verdicts(monkeypatch, {4242: process_ownership.UNVERIFIABLE}) + + assert bg_jobs.disown_unverified()["retired"] == 1 + + +def test_a_job_that_is_still_ours_keeps_running(job_store, monkeypatch): + """A detached job is documented to survive a restart. Killing it here would + break the feature the store exists for.""" + seed_job() + verdicts(monkeypatch, {4242: process_ownership.OWNED}) + + report = bg_jobs.disown_unverified() + + assert report == {"seen": 1, "retired": 0, "kept": 1} + assert bg_jobs._load()["job-1"]["status"] == "running" + + +def test_a_job_whose_process_is_gone_is_left_for_refresh(job_store, monkeypatch): + """refresh() may still find an exit-code file the job wrote before it went, + so retiring it here would discard a result that exists.""" + seed_job() + verdicts(monkeypatch, {4242: process_ownership.GONE}) + + assert bg_jobs.disown_unverified()["kept"] == 1 + assert bg_jobs._load()["job-1"]["status"] == "running" + + +def test_already_finished_jobs_are_not_reconsidered(job_store, monkeypatch): + seed_job(status="done") + verdicts(monkeypatch, {4242: process_ownership.FOREIGN}) + + assert bg_jobs.disown_unverified() == {"seen": 0, "retired": 0, "kept": 0} + + +def test_a_launched_job_records_an_identity_next_to_its_pid(job_store): + """Without this the record is unverifiable forever and the reaper can only + refuse — the token has to be captured at launch or not at all.""" + record = bg_jobs.launch("true", "chat-1") + + assert "start_token" in record + assert process_ownership.verify(record["pid"], record["start_token"]) in ( + process_ownership.OWNED, process_ownership.GONE, + ) + + +def test_an_abandoned_job_says_so_in_its_follow_up(job_store): + """The agent is told the job was lost, not that it failed for its own reasons.""" + record = seed_job(status="failed", ownership_lost=process_ownership.FOREIGN) + + text = bg_jobs.result_text(record) + + assert "abandoned across a server restart" in text + assert "neither waited on nor signalled" in text + + +def test_reap_orphans_reports_both_stores_and_the_mechanism( + grant_store, job_store, monkeypatch +): + seed_grant() + seed_job() + verdicts(monkeypatch, {4242: process_ownership.GONE}) + monkeypatch.setattr(containment, "_group_present", lambda _pgid: False) + + report = process_reaper.reap_orphans() + + assert report["mechanism"] == process_ownership.inspection_mechanism() + assert report["grants"]["already_gone"] == 1 + assert report["bg_jobs"]["kept"] == 1 diff --git a/tests/test_process_ownership.py b/tests/test_process_ownership.py new file mode 100644 index 000000000..f4d8244fa --- /dev/null +++ b/tests/test_process_ownership.py @@ -0,0 +1,327 @@ +"""Process identity: the four verdicts, and that the fourth is never permissive. + +Split in two. The verdict tests use **real processes**, because the claim under +test is about the kernel's behaviour — a pid that has been reaped, a pid that was +never issued, a pid whose start time differs from the one recorded — and a fake +process table cannot be wrong about that in the same ways. The mechanism tests +substitute the inspection layer, so both the procfs branch and the ``ps`` branch +are exercised on whichever kind of host happens to be running them. + +The invariant worth most here is negative: :data:`process_ownership.OWNED` is the +only verdict that permits a signal, and nothing — a missing token, an absent +mechanism, a probe that raised — may produce it by default. Process inspection +has broken off Linux four times in this tree (ODY-70, -86, -94, -99), every time +because an absent mechanism read as a successful answer. +""" + +import os +import subprocess + +import pytest + +from core import platform_compat +from src import process_ownership as po + + +@pytest.fixture +def sleeper(): + """A real, short-lived child in its own session. Always reaped.""" + procs = [] + + def _spawn(argv=("sleep", "30")): + proc = subprocess.Popen(list(argv), start_new_session=True) + procs.append(proc) + return proc + + yield _spawn + for proc in procs: + try: + proc.kill() + proc.wait(timeout=5) + except Exception: + pass + + +# ── Verdicts, against real processes ──────────────────────────────────────── +def test_a_live_process_with_its_own_token_is_owned(sleeper): + proc = sleeper() + token = po.start_token(proc.pid) + + assert token + assert po.verify(proc.pid, token) == po.OWNED + + +def test_the_same_pid_with_a_different_token_is_foreign(sleeper): + """The whole point: a pid is a slot, and the token says who is in it.""" + proc = sleeper() + token = po.start_token(proc.pid) + + assert po.verify(proc.pid, str(token) + "-not-this-one") == po.FOREIGN + + +def test_a_reaped_process_is_gone(sleeper): + proc = sleeper() + token = po.start_token(proc.pid) + proc.kill() + proc.wait(timeout=5) + + assert po.verify(proc.pid, token) == po.GONE + + +def test_a_pid_that_was_never_issued_is_gone(): + # Above any plausible pid_max, so this cannot collide with a real process. + assert po.verify(2 ** 30, "token:anything") == po.GONE + + +def test_a_missing_token_is_unverifiable_and_never_owned(sleeper): + """A record that captured no identity cannot acquire one afterwards. + + This is the pre-upgrade record, and the reason it must not be OWNED is that + treating "we did not write it down" as "it is ours" is what makes a recycled + pid lethal. + """ + proc = sleeper() + + assert po.verify(proc.pid, None) == po.UNVERIFIABLE + assert po.verify(proc.pid, "") == po.UNVERIFIABLE + assert po.UNVERIFIABLE not in po.SIGNALLABLE + + +def test_only_owned_permits_a_signal(): + assert po.SIGNALLABLE == frozenset({po.OWNED}) + + +def test_a_falsy_pid_is_gone_rather_than_unverifiable(): + """Nothing to identify and nothing to signal; the record is just empty.""" + assert po.verify(None, "token:x") == po.GONE + assert po.verify(0, "token:x") == po.GONE + + +def test_capture_always_returns_both_fields(sleeper): + proc = sleeper() + captured = po.capture(proc.pid) + + assert set(captured) == {"pid", "start_token"} + assert captured["pid"] == proc.pid + assert po.verify(captured["pid"], captured["start_token"]) == po.OWNED + + +def test_this_process_verifies_as_itself(): + assert po.verify(os.getpid(), po.start_token(os.getpid())) == po.OWNED + + +# ── An unavailable mechanism is a failure, not a default ──────────────────── +def test_no_inspection_mechanism_yields_unverifiable(monkeypatch, sleeper): + proc = sleeper() + token = po.start_token(proc.pid) + monkeypatch.setattr(po, "inspection_mechanism", lambda: po.MECHANISM_NONE) + + # Not GONE (which would abandon a live process) and not OWNED (which would + # license a signal at an unidentified one). + assert po.verify(proc.pid, token) == po.UNVERIFIABLE + + +def test_no_mechanism_makes_start_token_raise_rather_than_return_none(monkeypatch): + """None means "no such process". A question we could not ask is not that.""" + monkeypatch.setattr(po, "inspection_mechanism", lambda: po.MECHANISM_NONE) + + with pytest.raises(po.InspectionUnavailable): + po.start_token(os.getpid()) + + +def test_a_probe_that_raises_is_unverifiable_not_owned(monkeypatch, sleeper): + proc = sleeper() + + def _broken(_pid): + raise po.InspectionUnavailable("deliberately broken probe") + + monkeypatch.setattr(po, "start_token", _broken) + + assert po.verify(proc.pid, "token:whatever") == po.UNVERIFIABLE + + +def test_capture_records_no_token_rather_than_failing(monkeypatch): + """A host that cannot identify its children must still be able to launch. + + The record then reads UNVERIFIABLE forever, which is the honest outcome: + the launch is allowed, and the later teardown refuses. + """ + monkeypatch.setattr(po, "inspection_mechanism", lambda: po.MECHANISM_NONE) + + captured = po.capture(4242) + + assert captured == {"pid": 4242, "start_token": None} + assert po.verify(4242, captured["start_token"]) == po.UNVERIFIABLE + + +def test_process_table_raises_without_any_mechanism(monkeypatch): + monkeypatch.setattr(po, "inspection_mechanism", lambda: po.MECHANISM_NONE) + + with pytest.raises(po.InspectionUnavailable): + po.process_table() + + +def test_the_procfs_table_guards_its_own_scan(monkeypatch, tmp_path): + """Guarded in the function that scans, not only in its caller. + + tests/test_procfs_scan_guard.py pins this structurally; this pins the + behaviour, so calling the branch directly on a procfs-less host raises + instead of FileNotFoundError. + """ + monkeypatch.setattr(platform_compat, "PROC_ROOT", tmp_path / "absent") + + with pytest.raises(po.InspectionUnavailable): + po._procfs_process_table() + + +# ── Mechanism selection ───────────────────────────────────────────────────── +def test_procfs_is_preferred_where_it_exists(monkeypatch, tmp_path): + procfs = tmp_path / "proc" + procfs.mkdir() + monkeypatch.setattr(platform_compat, "PROC_ROOT", procfs) + monkeypatch.setattr(po, "PROC_ROOT", procfs) + monkeypatch.setattr(po, "IS_WINDOWS", False) + + assert po.inspection_mechanism() == po.MECHANISM_PROCFS + + +def test_ps_covers_hosts_with_no_procfs(monkeypatch, tmp_path): + """macOS and the BSDs. The reason this module is not another /proc scan.""" + monkeypatch.setattr(platform_compat, "PROC_ROOT", tmp_path / "absent") + monkeypatch.setattr(po, "IS_WINDOWS", False) + monkeypatch.setattr(po.shutil, "which", lambda name: "/bin/ps" if name == "ps" else None) + + assert po.inspection_mechanism() == po.MECHANISM_PS + assert po.inspection_available() + + +def test_a_host_with_neither_reports_none(monkeypatch, tmp_path): + monkeypatch.setattr(platform_compat, "PROC_ROOT", tmp_path / "absent") + monkeypatch.setattr(po, "IS_WINDOWS", False) + monkeypatch.setattr(po.shutil, "which", lambda _name: None) + + assert po.inspection_mechanism() == po.MECHANISM_NONE + assert not po.inspection_available() + + +def test_the_procfs_token_reads_a_comm_containing_spaces_and_parens(monkeypatch, tmp_path): + """``/proc//stat`` field 2 is attacker-adjacent: it is the executable name. + + A process called ``my (weird) prog`` would shift every field after it if the + parser split on whitespace, which would silently read the wrong number as the + start time and make every verdict wrong. + """ + procfs = tmp_path / "proc" + (procfs / "77").mkdir(parents=True) + # "77 (comm) S" are fields 1-3, so the filler starts numbering at 4 and + # each value equals its own field number. + fields = " ".join(str(index) for index in range(4, 54)) + (procfs / "77" / "stat").write_text(f"77 (my (weird) prog) S {fields}\n") + monkeypatch.setattr(platform_compat, "PROC_ROOT", procfs) + monkeypatch.setattr(po, "PROC_ROOT", procfs) + monkeypatch.setattr(po, "IS_WINDOWS", False) + + assert po.start_token(77) == "procfs:22" + + +# ── The process tree ──────────────────────────────────────────────────────── +def test_descendants_walks_a_real_tree(sleeper): + """The grandchild case: a shell that backgrounds work and the work itself.""" + proc = sleeper(("bash", "-c", "sleep 30 & sleep 30")) + # Wait for the shell to have actually forked, without sleeping on a clock: + # poll the table until the children appear or the attempts run out. + found = [] + for _attempt in range(100): + found = po.descendants([proc.pid]) + if len(found) >= 3: + break + + assert proc.pid in found + assert len(found) >= 3, f"expected the shell and its two children, got {found}" + + +def test_descendants_includes_the_root_even_with_no_children(sleeper): + proc = sleeper() + + assert po.descendants([proc.pid]) == [proc.pid] + + +def test_descendants_takes_a_single_table_snapshot(monkeypatch): + """One table in, one answer out — reparenting cannot hide a process. + + Walking the tree with a fresh query per level lets a child be reparented + between queries and drop out of the result, which for a teardown means a + process nobody signals. + """ + rows = { + 10: po.ProcessInfo(pid=10, ppid=1, command="root"), + 11: po.ProcessInfo(pid=11, ppid=10, command="child"), + 12: po.ProcessInfo(pid=12, ppid=11, command="grandchild"), + 13: po.ProcessInfo(pid=13, ppid=1, command="unrelated"), + } + + def _explode(): + raise AssertionError("descendants must use the table it was given") + + monkeypatch.setattr(po, "process_table", _explode) + + assert po.descendants([10], table=rows) == [10, 11, 12] + + +def test_descendants_terminates_on_a_parent_cycle(): + """A table can report a cycle; the walk must not spin on it.""" + rows = { + 20: po.ProcessInfo(pid=20, ppid=21, command="a"), + 21: po.ProcessInfo(pid=21, ppid=20, command="b"), + } + + assert sorted(po.descendants([20], table=rows)) == [20, 21] + + +def test_the_procfs_table_reads_the_parent_pid(monkeypatch, tmp_path): + """The procfs branch of the tree walk, exercised on a host without procfs. + + macOS runs this suite and takes the ``ps`` branch, so without a substituted + ``/proc`` the Linux parse — which is what the deployed image uses — would be + covered by nothing. + """ + procfs = tmp_path / "proc" + for pid, ppid in ((10, 1), (11, 10)): + (procfs / str(pid)).mkdir(parents=True) + (procfs / str(pid) / "cmdline").write_bytes(f"proc-{pid}\0--flag\0".encode()) + filler = " ".join(str(index) for index in range(5, 54)) + (procfs / str(pid) / "stat").write_text(f"{pid} (proc) S {ppid} {filler}\n") + monkeypatch.setattr(platform_compat, "PROC_ROOT", procfs) + monkeypatch.setattr(po, "PROC_ROOT", procfs) + monkeypatch.setattr(po, "IS_WINDOWS", False) + + table = po.process_table() + + assert table[11].ppid == 10 + assert table[10].command == "proc-10 --flag" + assert po.descendants([10], table=table) == [10, 11] + + +def test_a_procfs_row_with_an_unreadable_stat_keeps_its_command_line( + monkeypatch, tmp_path +): + """A kernel thread or a pid that exits mid-walk still matters to a + command-line match; dropping the row entirely would hide it.""" + procfs = tmp_path / "proc" + (procfs / "12").mkdir(parents=True) + (procfs / "12" / "cmdline").write_bytes(b"orphan-cmd\0") + monkeypatch.setattr(platform_compat, "PROC_ROOT", procfs) + monkeypatch.setattr(po, "PROC_ROOT", procfs) + monkeypatch.setattr(po, "IS_WINDOWS", False) + + table = po.process_table() + + assert table[12].command == "orphan-cmd" + assert table[12].ppid == 0 + + +def test_command_lines_sees_this_process(): + table = po.command_lines() + + assert os.getpid() in table + assert table[os.getpid()]