mirror of
https://github.com/pewdiepie-archdaemon/odysseus.git
synced 2026-10-06 06:52:20 +02:00
fix(runtime): verify process identity before any teardown signal
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/<pid>/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.
This commit is contained in:
@@ -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():
|
||||
|
||||
+62
-1
@@ -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."
|
||||
|
||||
+142
-1
@@ -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",
|
||||
|
||||
@@ -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/<pid>/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/<pid>/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/<pid>/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
|
||||
@@ -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)
|
||||
+203
-30
@@ -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}
|
||||
|
||||
|
||||
@@ -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. ``"<session> <pane_pid>"`` 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"]
|
||||
|
||||
|
||||
|
||||
@@ -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
|
||||
@@ -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/<pid>/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()]
|
||||
Reference in New Issue
Block a user