diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index de1a61421dd4f..20ac311c64f65 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -95,6 +95,7 @@ import threading import logging import time +import psutil from contextvars import ContextVar, Token from dataclasses import dataclass, field from pathlib import Path @@ -457,6 +458,14 @@ def _fire_dispatch_tick_hook( # signal lands, and the following tick reclaims cleanly. RECLAIM_DEFER_GRACE_SECONDS = 120 +# A pid-less claim whose local CLAIMER (dispatching gateway) is dead may still +# have a detached worker: Popen precedes the pid stamp. With no heartbeat on +# the current run yet, hold the claim this long after it was taken before +# treating the spawn as never-happened. Measured 2026-09-24 over 12,732 live +# runs: claim-to-first-heartbeat p50 6 s, p90 21 s, p99 69 s, max 708 s. +# The claim TTL (900 s) sits above that max. See _dead_claimer_release_at. +DEAD_CLAIMER_LAUNCH_BOUND_SECONDS = DEFAULT_CLAIM_TTL_SECONDS + def _resolve_claim_ttl_seconds(ttl_seconds: Optional[int] = None) -> int: """Return the effective claim TTL, honoring the kanban env override. @@ -7722,6 +7731,7 @@ def release_stale_claims( termination = _terminate_reclaimed_worker( row["worker_pid"], row["claim_lock"], signal_fn=signal_fn, + conn=conn, task_id=row["id"], ) # Never release a claim while our own worker is still alive: that would # spawn a duplicate beside it. Hold the claim and retry next tick. @@ -7809,7 +7819,8 @@ def reclaim_task( for the TTL to expire (e.g. after seeing a hallucination warning). If the worker was signalled but survived, or if this host holds the - claim but its worker pid cannot be resolved, reclamation FAILS CLOSED. + claim without a worker pid and its claimer may still be alive, + reclamation FAILS CLOSED. The card retains its owner and emits a ``reclaim_refused`` event marked ``needs_attention``. A human must resolve the worker outside this path; an operator request alone does not prove the worker is gone. @@ -7843,6 +7854,7 @@ def reclaim_task( prev_lock = row["claim_lock"] termination = _terminate_reclaimed_worker( row["worker_pid"], prev_lock, signal_fn=signal_fn, + conn=conn, task_id=task_id, ) # Never release a claim while our host-local worker is alive or its # liveness is unknown. This also covers NULL pid in the TTL and stale @@ -13464,11 +13476,63 @@ def _pid_alive(pid: Optional[int]) -> bool: return True +def _dead_claimer_release_at( + conn: sqlite3.Connection, task_id: str, +) -> tuple[Optional[int], str, Optional[str]]: + """Earliest time a pid-less claim with a DEAD claimer may be released. + + A dead claimer is not proof that no worker exists: ``_default_spawn`` + calls ``Popen(start_new_session=True)`` before ``_set_worker_pid`` + commits, so a gateway killed in that window leaves a live orphan with no + stamped pid. The row cannot tell such an orphan from a spawn that never + happened, so release is time-bounded on the CURRENT run only: + + * no worker evidence yet (no ``heartbeat``/``spawned`` event on this run): + release after ``DEAD_CLAIMER_LAUNCH_BOUND_SECONDS`` from the claim. An + orphan that has not heartbeated by then is outside every observed launch + (live ledger 2026-09-24, 12,732 runs: first heartbeat p99 69 s, max 708 s). + * worker evidence exists: release after the newest evidence is older than + ``DEFAULT_CLAIM_HEARTBEAT_MAX_STALE_SECONDS`` — the same staleness rule + ``release_stale_claims`` applies to a worker whose pid IS alive. Without + this bound an orphan that heartbeated once and then died held the card + forever (the stuck-running shape this path exists to end). + + Returns ``(release_at, basis, evidence_kind)``; ``release_at`` is None when + there is no current run to anchor the bound (held). + """ + run = conn.execute( + "SELECT r.id, r.started_at FROM tasks t " + "JOIN task_runs r ON r.id = t.current_run_id WHERE t.id = ?", + (task_id,), + ).fetchone() + if run is None or run["started_at"] is None: + return None, "no_current_run", None + evidence = conn.execute( + "SELECT kind, created_at FROM task_events WHERE task_id = ? " + "AND run_id = ? AND kind IN ('heartbeat', 'spawned') " + "ORDER BY created_at DESC, id DESC LIMIT 1", + (task_id, int(run["id"])), + ).fetchone() + if evidence is None: + return ( + int(run["started_at"]) + DEAD_CLAIMER_LAUNCH_BOUND_SECONDS, + "launch_bound", + None, + ) + return ( + int(evidence["created_at"]) + DEFAULT_CLAIM_HEARTBEAT_MAX_STALE_SECONDS, + "evidence_stale_bound", + evidence["kind"], + ) + + def _terminate_reclaimed_worker( pid: Optional[int], claim_lock: Optional[str], *, signal_fn=None, + conn: Optional[sqlite3.Connection] = None, + task_id: Optional[str] = None, ) -> dict[str, Any]: """Best-effort host-local worker termination for reclaim paths.""" import signal @@ -13489,16 +13553,43 @@ def _terminate_reclaimed_worker( info["host_local"] = True if not pid or pid <= 0: - # OUR host holds this claim but no worker pid was ever stamped, so - # there is nothing to signal — and, critically, no evidence of death - # either. Flag that explicitly instead of returning a payload that - # looks identical to a clean "nothing to do". - # - # The pre-fix code returned here BEFORE setting host_local, so the - # event claimed the worker was some other host's problem. Measured on - # t_09180e10 (2026-09-22): `prev_pid: null, host_local: false` on a - # card whose lock was `mac-studio-m3u:71817` — our own host. + # OUR host holds this claim but no worker pid was ever stamped. That + # is UNKNOWN liveness unless we can prove otherwise (t_09180e10). info["liveness_unprovable"] = True + if conn is None or task_id is None: + # Fail closed on omission: without the run row we cannot rule out + # a detached worker whose pid was never stamped, so a caller that + # forgets conn/task_id must hold the claim, never release it. + info["unstamped_worker_check"] = "skipped_no_run_context" + return info + claimer_pid = 0 + try: + claimer_pid = int(str(claim_lock)[len(host_prefix):]) + if claimer_pid <= 0: + return info + # Equivalent to kill(pid, 0), without Windows' destructive + # CTRL_C_EVENT behavior for signal 0. + psutil.Process(claimer_pid) + # Claimer alive: its spawn may still be in flight (launch grace). + return info + except psutil.NoSuchProcess: + pass + except (ValueError, OSError, psutil.AccessDenied): + return info + # The claimer is dead, but that alone does not prove no worker exists: + # Popen(start_new_session=True) precedes _set_worker_pid, so a claimer + # killed in between leaves a live, unstamped orphan. Release only once + # no worker can still be alive for THIS run (see _dead_claimer_release_at). + info["claimer_pid_dead"] = claimer_pid + release_at, basis, evidence_kind = _dead_claimer_release_at(conn, task_id) + info["dead_claimer_release_basis"] = basis + if evidence_kind: + info["unstamped_worker_evidence"] = evidence_kind + if release_at is None or int(time.time()) < release_at: + info["dead_claimer_hold_until"] = release_at + return info + info["liveness_unprovable"] = False + info["terminated"] = True return info kill = signal_fn if signal_fn is not None else ( @@ -13543,10 +13634,9 @@ def _worker_survived_termination(termination: dict) -> bool: """True when a host-local worker has NOT been proven gone. A signalled-but-still-alive worker is positive liveness evidence. A - missing pid is UNKNOWN liveness — not death evidence. Neither may release - the claim, because that would let the dispatcher spawn a second worker - beside the first. A proven-dead worker (including ProcessLookupError on - SIGTERM) and non-local claims use the normal release path. + missing worker pid is UNKNOWN liveness if the claimer may still launch or + the current run has evidence of an unstamped worker. Unknown liveness must + not release a claim: it could spawn a second worker beside the first. """ if not termination.get("host_local") or termination.get("terminated"): return False @@ -13968,7 +14058,9 @@ def detect_progress_stalls( _append_event(conn, row["id"], "stalled", evidence, run_id=rid) if age < reclaim_seconds: continue - termination = _terminate_reclaimed_worker(pid, lock) + termination = _terminate_reclaimed_worker( + pid, lock, conn=conn, task_id=row["id"], + ) if _worker_survived_termination(termination): _defer_reclaim_for_live_worker( conn, row["id"], lock, now, termination, reason="progress_stalled_worker_alive", @@ -14063,7 +14155,7 @@ def detect_stale_running( # Terminate the worker if it's still host-local. termination = _terminate_reclaimed_worker( - pid, lock, signal_fn=signal_fn, + pid, lock, signal_fn=signal_fn, conn=conn, task_id=tid, ) # Never release a claim while our own worker is still alive: that would @@ -15309,7 +15401,9 @@ def _abort_lost_claim_spawn( "run_id": task.current_run_id, } if pid: - termination = _terminate_reclaimed_worker(int(pid), task.claim_lock) + termination = _terminate_reclaimed_worker( + int(pid), task.claim_lock, conn=conn, task_id=task.id, + ) payload.update(termination) if not termination.get("terminated"): payload["needs_attention"] = True diff --git a/plugins/kanban/dashboard/plugin_api.py b/plugins/kanban/dashboard/plugin_api.py index aafa5107f41cd..558c0f6656bf0 100644 --- a/plugins/kanban/dashboard/plugin_api.py +++ b/plugins/kanban/dashboard/plugin_api.py @@ -1167,6 +1167,7 @@ def _set_status_direct( if held["status"] == "running" and new_status != "running": termination = kanban_db._terminate_reclaimed_worker( held["worker_pid"], held["claim_lock"], + conn=conn, task_id=task_id, ) if kanban_db._worker_survived_termination(termination): kanban_db._refuse_reclaim_unproven_death( diff --git a/tests/hermes_cli/test_kanban_batch_set_model.py b/tests/hermes_cli/test_kanban_batch_set_model.py index bd9a8bfbb287d..e7f1dc4d0ec5f 100644 --- a/tests/hermes_cli/test_kanban_batch_set_model.py +++ b/tests/hermes_cli/test_kanban_batch_set_model.py @@ -106,7 +106,7 @@ def test_set_model_reclaim_only_reclaims_selected_running_cards( "UPDATE tasks SET status='running', claim_lock=?, worker_pid=? WHERE id=?", (f"host:{pid}", pid, task_id), ) - monkeypatch.setattr(kb, "_terminate_reclaimed_worker", lambda pid, lock, signal_fn=None: signaled.append(pid) or {}) + monkeypatch.setattr(kb, "_terminate_reclaimed_worker", lambda pid, lock, **_kw: signaled.append(pid) or {}) out = kc.run_slash( "set-model model-a --provider batch-provider --reclaim " @@ -144,7 +144,7 @@ def test_set_model_reclaim_leaves_non_running_selected_cards_alone( ) monkeypatch.setattr( kb, "_terminate_reclaimed_worker", - lambda pid, lock, signal_fn=None: signaled.append(pid) or {}, + lambda pid, lock, **_kw: signaled.append(pid) or {}, ) out = kc.run_slash( diff --git a/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py b/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py index 1442e9ec15c30..78049b486510b 100644 --- a/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py +++ b/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py @@ -26,8 +26,12 @@ from __future__ import annotations import json +import os from pathlib import Path +import psutil import secrets +import subprocess +import sys import time import pytest @@ -110,6 +114,319 @@ def _row(conn, tid): ).fetchone() +def _age_claim(conn, tid, run_id, seconds): + """Backdate the claim (task + current run) by ``seconds``.""" + old = int(time.time()) - seconds + conn.execute("UPDATE tasks SET started_at=? WHERE id=?", (old, tid)) + conn.execute("UPDATE task_runs SET started_at=? WHERE id=?", (old, run_id)) + conn.commit() + + +_PAST_LAUNCH_BOUND = kb.DEAD_CLAIMER_LAUNCH_BOUND_SECONDS + 5 + + +def _dead_pid() -> int: + """A real pid that has exited (not a guessed number).""" + c = subprocess.Popen([sys.executable, "-c", "pass"], stdin=subprocess.DEVNULL) + c.wait(timeout=10) + assert not psutil.pid_exists(c.pid) + return c.pid + + +def _load_dashboard_api(): + import importlib.util + path = Path(__file__).resolve().parents[2] / "plugins/kanban/dashboard/plugin_api.py" + spec = importlib.util.spec_from_file_location("t9d7_plugin_api", path) + assert spec is not None and spec.loader is not None + mod = importlib.util.module_from_spec(spec) + sys.modules[spec.name] = mod + spec.loader.exec_module(mod) + return mod + + +def _release(conn, tid, run_id, path): + """Drive one release entry point; return True when it released the claim.""" + if path == "manual": + return kb.reclaim_task(conn, tid, reason="worker sweep: missing pid") + if path == "ttl": + conn.execute("UPDATE tasks SET claim_expires=? WHERE id=?", (int(time.time()) - 1, tid)) + conn.commit() + return kb.release_stale_claims(conn) == 1 + if path == "stale": + # Keep the claim age the caller set; only make the heartbeat column + # stale so detect_stale_running considers the row at all. + old = int(time.time()) - 7200 + conn.execute("UPDATE tasks SET last_heartbeat_at=? WHERE id=?", (old, tid)) + conn.commit() + return kb.detect_stale_running(conn, stale_timeout_seconds=30) == [tid] + if path == "dashboard": + return _load_dashboard_api()._set_status_direct(conn, tid, "ready") + raise AssertionError(path) + + +def test_dead_local_claimer_without_worker_pid_is_reclaimed(conn, monkeypatch): + """A dead claimer without any worker evidence can release its claim.""" + dead_pid = 999991 + lock = f"{kb._claimer_id().split(':', 1)[0]}:{dead_pid}" + tid, _, run_id = _running_card(conn, worker_pid=None, lock=lock) + _age_claim(conn, tid, run_id, _PAST_LAUNCH_BOUND) + + def dead_claimer(pid): + assert pid == dead_pid + raise psutil.NoSuchProcess(pid) + + monkeypatch.setattr(psutil, "Process", dead_claimer) + assert kb.reclaim_task(conn, tid, reason="worker sweep: missing pid") is True + assert _row(conn, tid)["status"] == "ready" + assert _row(conn, tid)["claim_lock"] is None + assert conn.execute("SELECT outcome FROM task_runs WHERE id=?", (run_id,)).fetchone()[0] == "reclaimed" + payload = _events(conn, tid, "reclaimed")[-1] + assert payload["terminated"] is True + assert payload["claimer_pid_dead"] == dead_pid + assert payload["dead_claimer_release_basis"] == "launch_bound" + assert not _events(conn, tid, "reclaim_refused") + + +def test_heartbeat_from_previous_run_does_not_hold_dead_claimer(conn, monkeypatch): + dead_pid = 999995 + lock = f"{kb._claimer_id().split(':', 1)[0]}:{dead_pid}" + tid, _, run_id = _running_card(conn, worker_pid=None, lock=lock) + _age_claim(conn, tid, run_id, _PAST_LAUNCH_BOUND) + with kb.write_txn(conn): + kb._append_event(conn, tid, "heartbeat", run_id=run_id - 1) + def dead(pid): + raise psutil.NoSuchProcess(pid) + monkeypatch.setattr(psutil, "Process", dead) + assert kb.reclaim_task(conn, tid) is True + + +@pytest.mark.parametrize("reclaim_path", ["manual", "ttl", "stale", "dashboard"]) +def test_dead_claimer_with_unstamped_orphan_before_first_heartbeat_keeps_claim( + conn, reclaim_path, +): + """Argus r2 B3: an orphan that has not heartbeated YET is still a worker. + + Real processes: the claimer has exited, a detached orphan is alive, no + heartbeat or pid landed, and the claim is 60 s old (past the sweep's 45 s + launch grace; live p99 first heartbeat is 69 s, max 708 s). No release + path may free the claim or let a second worker spawn. + """ + lock = f"{kb._claimer_id().split(':', 1)[0]}:{_dead_pid()}" + tid, _, run_id = _running_card(conn, worker_pid=None, lock=lock) + _age_claim(conn, tid, run_id, 60) + orphan = subprocess.Popen( + [sys.executable, "-c", "import time; time.sleep(60)"], + stdin=subprocess.DEVNULL, start_new_session=True, + ) + try: + assert _release(conn, tid, run_id, reclaim_path) is False + spawned = [] + kb.dispatch_once(conn, spawn_fn=lambda task, workspace, board=None: spawned.append(task.id) or 424242) + assert tid not in spawned + assert orphan.poll() is None + assert _row(conn, tid)["status"] == "running" + assert _row(conn, tid)["claim_lock"] == lock + assert not _events(conn, tid, "reclaimed") + finally: + orphan.kill() + orphan.wait(timeout=10) + + +@pytest.mark.parametrize("reclaim_path", ["manual", "ttl", "stale", "dashboard"]) +def test_dead_claimer_orphan_that_heartbeated_then_died_is_released_once_stale( + conn, reclaim_path, +): + """Argus r2 B5: heartbeat-then-die must not hold the card forever. + + Once the newest current-run worker evidence is older than + DEFAULT_CLAIM_HEARTBEAT_MAX_STALE_SECONDS (the bound release_stale_claims + already applies to a live-pid worker), the claim is released. + """ + lock = f"{kb._claimer_id().split(':', 1)[0]}:{_dead_pid()}" + tid, _, run_id = _running_card(conn, worker_pid=None, lock=lock) + orphan = subprocess.Popen( + [sys.executable, "-c", "import time; time.sleep(60)"], + stdin=subprocess.DEVNULL, start_new_session=True, + ) + assert kb.heartbeat_worker(conn, tid, note="orphan alive once") + orphan.kill() + orphan.wait(timeout=10) + aged = int(time.time()) - kb.DEFAULT_CLAIM_HEARTBEAT_MAX_STALE_SECONDS - 5 + conn.execute( + "UPDATE task_events SET created_at=? WHERE task_id=? AND run_id=? AND kind='heartbeat'", + (aged, tid, run_id), + ) + conn.commit() + _age_claim(conn, tid, run_id, kb.DEFAULT_CLAIM_HEARTBEAT_MAX_STALE_SECONDS + 10) + assert _release(conn, tid, run_id, reclaim_path) is True + assert _row(conn, tid)["claim_lock"] is None + assert _row(conn, tid)["status"] != "running" + + +def test_dead_claimer_fresh_heartbeat_holds_even_if_orphan_is_gone(conn): + """The bound, not a guess about the orphan, decides: fresh evidence holds.""" + lock = f"{kb._claimer_id().split(':', 1)[0]}:{_dead_pid()}" + tid, _, run_id = _running_card(conn, worker_pid=None, lock=lock) + _age_claim(conn, tid, run_id, _PAST_LAUNCH_BOUND) + assert kb.heartbeat_worker(conn, tid, note="recent") + assert kb.reclaim_task(conn, tid) is False + refused = _events(conn, tid, "reclaim_refused")[-1] + assert refused["dead_claimer_release_basis"] == "evidence_stale_bound" + assert refused["unstamped_worker_evidence"] == "heartbeat" + + +def test_dead_claimer_inside_launch_bound_is_held(conn): + lock = f"{kb._claimer_id().split(':', 1)[0]}:{_dead_pid()}" + tid, _, run_id = _running_card(conn, worker_pid=None, lock=lock) + _age_claim(conn, tid, run_id, kb.DEAD_CLAIMER_LAUNCH_BOUND_SECONDS - 60) + assert kb.reclaim_task(conn, tid) is False + assert _row(conn, tid)["claim_lock"] == lock + + +def test_missing_run_context_fails_closed_for_dead_claimer(conn): + """Omitting conn/task_id can never turn a dead claimer into a release.""" + lock = f"{kb._claimer_id().split(':', 1)[0]}:{_dead_pid()}" + info = kb._terminate_reclaimed_worker(None, lock) + assert kb._worker_survived_termination(info) is True + assert info["unstamped_worker_check"] == "skipped_no_run_context" + + +def test_every_gated_termination_call_site_passes_run_context(): + """Argus r2 B2 class guard: a release gate must hand the run to the probe. + + Any call whose result is bound to a name (i.e. used to decide a release) + must pass ``conn=`` and ``task_id=``. Bare post-commit kill calls are + fire-and-forget and do not gate a release. + """ + import ast + + root = Path(__file__).resolve().parents[2] + offenders = [] + for rel in ("hermes_cli/kanban_db.py", "plugins/kanban/dashboard/plugin_api.py"): + tree = ast.parse((root / rel).read_text(encoding="utf-8")) + for node in ast.walk(tree): + if not isinstance(node, ast.Assign) or not isinstance(node.value, ast.Call): + continue + fn = node.value.func + name = fn.attr if isinstance(fn, ast.Attribute) else getattr(fn, "id", None) + if name != "_terminate_reclaimed_worker": + continue + kws = {k.arg for k in node.value.keywords} + if not {"conn", "task_id"} <= kws: + offenders.append(f"{rel}:{node.lineno}") + assert not offenders, f"gated call sites missing conn/task_id: {offenders}" + + +@pytest.mark.parametrize("reclaim_path", ["manual", "ttl", "stale"]) +def test_dead_claimer_with_unstamped_heartbeating_worker_keeps_claim(conn, reclaim_path): + """Popen can succeed before the gateway commits _set_worker_pid. + + A detached worker may then outlive the gateway. Its heartbeat belongs to + the current run even if the PID was never stamped. Neither operator nor + automatic recovery may schedule a second worker beside it. + """ + host = kb._claimer_id().split(":", 1)[0] + claimer = subprocess.Popen([sys.executable, "-c", "pass"], stdin=subprocess.DEVNULL) + claimer.wait(timeout=10) + assert not psutil.pid_exists(claimer.pid) + lock = f"{host}:{claimer.pid}" + tid, _, run_id = _running_card(conn, worker_pid=None, lock=lock) + orphan = subprocess.Popen( + [sys.executable, "-c", "import time; time.sleep(60)"], + stdin=subprocess.DEVNULL, start_new_session=True, + ) + try: + assert kb.heartbeat_claim(conn, tid, claimer=lock) + assert kb.heartbeat_worker(conn, tid, note="unstamped worker alive") + assert _row(conn, tid)["worker_pid"] is None + if reclaim_path == "ttl": + conn.execute("UPDATE tasks SET claim_expires=? WHERE id=?", (int(time.time()) - 1, tid)) + conn.commit() + assert kb.release_stale_claims(conn) == 0 + elif reclaim_path == "stale": + old = int(time.time()) - 7200 + conn.execute("UPDATE task_runs SET started_at=? WHERE id=?", (old, run_id)) + conn.execute("UPDATE tasks SET last_heartbeat_at=? WHERE id=?", (old, tid)) + conn.commit() + assert kb.detect_stale_running(conn, stale_timeout_seconds=60) == [] + else: + assert kb.reclaim_task(conn, tid, reason="worker sweep: missing pid") is False + spawned = [] + kb.dispatch_once(conn, spawn_fn=lambda task, workspace, board=None: spawned.append(task.id) or 424242) + assert tid not in spawned + assert orphan.poll() is None + assert _row(conn, tid)["status"] == "running" + assert _row(conn, tid)["claim_lock"] == lock + assert not _events(conn, tid, "reclaimed") + finally: + orphan.kill() + orphan.wait(timeout=10) + + +def test_expired_dead_claimer_is_not_renewed(conn, monkeypatch): + dead_pid = 999992 + lock = f"{kb._claimer_id().split(':', 1)[0]}:{dead_pid}" + tid, _, run_id = _running_card(conn, worker_pid=None, lock=lock) + _age_claim(conn, tid, run_id, _PAST_LAUNCH_BOUND) + conn.execute("UPDATE tasks SET claim_expires=? WHERE id=?", (int(time.time()) - 1, tid)) + conn.commit() + + def dead_claimer(pid): + assert pid == dead_pid + raise psutil.NoSuchProcess(pid) + + monkeypatch.setattr(psutil, "Process", dead_claimer) + assert kb.release_stale_claims(conn) == 1 + assert _row(conn, tid)["status"] == "ready" + assert conn.execute("SELECT outcome FROM task_runs WHERE id=?", (run_id,)).fetchone()[0] == "reclaimed" + assert _events(conn, tid, "reclaimed")[-1]["claimer_pid_dead"] == dead_pid + assert not _events(conn, tid, "reclaim_deferred") + + +def test_live_local_claimer_without_worker_pid_still_holds_claim(conn, monkeypatch): + lock = f"{kb._claimer_id().split(':', 1)[0]}:{os.getpid()}" + tid, _, _ = _running_card(conn, worker_pid=None, lock=lock) + seen = [] + original_process = psutil.Process + + def live_claimer(pid): + seen.append(pid) + return original_process(pid) + + monkeypatch.setattr(psutil, "Process", live_claimer) + assert kb.reclaim_task(conn, tid) is False + assert seen == [os.getpid()] + assert _row(conn, tid)["claim_lock"] == lock + assert _events(conn, tid, "reclaim_refused")[-1]["reason"] == "liveness_unprovable" + + +@pytest.mark.parametrize("suffix, error", [ + ("not-a-pid", None), + ("999994", PermissionError), +]) +def test_unprovable_claimer_stays_held(conn, monkeypatch, suffix, error): + lock = f"{kb._claimer_id().split(':', 1)[0]}:{suffix}" + tid, _, _ = _running_card(conn, worker_pid=None, lock=lock) + + def inaccessible(pid): + raise psutil.AccessDenied(pid) + + if error: + monkeypatch.setattr(psutil, "Process", inaccessible) + assert kb.reclaim_task(conn, tid) is False + assert _row(conn, tid)["claim_lock"] == lock + assert not _events(conn, tid, "reclaimed") + + +def test_foreign_claimer_without_worker_pid_is_not_probed(conn, monkeypatch): + tid, _, _ = _running_card(conn, worker_pid=None, lock="other-host:999993") + monkeypatch.setattr(psutil, "Process", lambda *_: pytest.fail("foreign PID was probed")) + # The existing operator reclaim policy releases a foreign claim; do not + # assert local death for a process on a host we cannot inspect. + assert kb.reclaim_task(conn, tid) is True + assert _events(conn, tid, "reclaimed")[-1]["host_local"] is False + + # --------------------------------------------------------------------------- # 1. The incident: pid-less running card must NOT be re-queued ready. # ---------------------------------------------------------------------------