From d9d8ec516f84048c2eab5961c3dfffd300285c80 Mon Sep 17 00:00:00 2001 From: Apollo Date: Thu, 24 Sep 2026 08:39:03 -0700 Subject: [PATCH 1/5] fix(kanban): reclaim pid-less claim when local claimer is dead Probe the local claim owner with kill(pid, 0) before treating a missing worker PID as unknown. An absent claimer cannot finish its launch; preserve fail-closed grace for live or inaccessible claimers. Cover manual and TTL reclaim plus foreign-host behavior. --- hermes_cli/kanban_db.py | 33 ++++--- ...test_kanban_reclaim_unprovable_liveness.py | 86 +++++++++++++++++++ 2 files changed, 105 insertions(+), 14 deletions(-) diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index de1a61421dd4f..85237fbbf60e8 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -7809,7 +7809,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. @@ -13489,16 +13490,20 @@ 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. + # The local claimer may still be launching the worker. Only an absent + # claimer proves a never-stamped worker cannot be launched later. + # An alive, inaccessible or malformed claimer remains unprovable. info["liveness_unprovable"] = True + try: + claimer_pid = int(str(claim_lock)[len(host_prefix):]) + if claimer_pid > 0 and hasattr(os, "kill"): + os.kill(claimer_pid, 0) + except ProcessLookupError: + info["liveness_unprovable"] = False + info["terminated"] = True + info["claimer_pid_dead"] = claimer_pid + except (ValueError, OSError): + pass return info kill = signal_fn if signal_fn is not None else ( @@ -13543,10 +13548,10 @@ 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 unless the host-local claimer is + proven dead. Unknown liveness must not release a claim: it could spawn a + second worker beside the first. Proven-dead and non-local claims use the + normal release path. """ if not termination.get("host_local") or termination.get("terminated"): return False diff --git a/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py b/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py index 1442e9ec15c30..471f18bad1e53 100644 --- a/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py +++ b/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py @@ -26,6 +26,7 @@ from __future__ import annotations import json +import os from pathlib import Path import secrets import time @@ -110,6 +111,91 @@ def _row(conn, tid): ).fetchone() +def test_dead_local_claimer_without_worker_pid_is_reclaimed(conn, monkeypatch): + """A gateway that died before stamping a worker cannot still spawn one.""" + dead_pid = 999991 + lock = f"{kb._claimer_id().split(':', 1)[0]}:{dead_pid}" + tid, _, run_id = _running_card(conn, worker_pid=None, lock=lock) + + def dead_claimer(pid, sig): + assert (pid, sig) == (dead_pid, 0) + raise ProcessLookupError(pid) + + monkeypatch.setattr(kb.os, "kill", 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 not _events(conn, tid, "reclaim_refused") + + +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) + conn.execute("UPDATE tasks SET claim_expires=? WHERE id=?", (int(time.time()) - 1, tid)) + conn.commit() + + def dead_claimer(pid, sig): + assert (pid, sig) == (dead_pid, 0) + raise ProcessLookupError(pid) + + monkeypatch.setattr(kb.os, "kill", 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_kill = os.kill + + def live_claimer(pid, sig): + seen.append((pid, sig)) + return original_kill(pid, sig) + + monkeypatch.setattr(kb.os, "kill", live_claimer) + assert kb.reclaim_task(conn, tid) is False + assert seen == [(os.getpid(), 0)] + 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, sig): + assert sig == 0 + raise error(pid) + + if error: + monkeypatch.setattr(kb.os, "kill", 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(kb.os, "kill", 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. # --------------------------------------------------------------------------- From cdae9f0a7fd0f92b3c91919cf9fd29488c3eb9c0 Mon Sep 17 00:00:00 2001 From: Apollo Date: Thu, 24 Sep 2026 08:50:01 -0700 Subject: [PATCH 2/5] fix(kanban): initialize claimer probe PID before exception path --- hermes_cli/kanban_db.py | 1 + 1 file changed, 1 insertion(+) diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index 85237fbbf60e8..45e7ead5a2976 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -13494,6 +13494,7 @@ def _terminate_reclaimed_worker( # claimer proves a never-stamped worker cannot be launched later. # An alive, inaccessible or malformed claimer remains unprovable. info["liveness_unprovable"] = True + claimer_pid = 0 try: claimer_pid = int(str(claim_lock)[len(host_prefix):]) if claimer_pid > 0 and hasattr(os, "kill"): From 5ad2916979c5e50537c2b7f4f5cb1ae9fee93390 Mon Sep 17 00:00:00 2001 From: Apollo Date: Thu, 24 Sep 2026 09:07:16 -0700 Subject: [PATCH 3/5] fix(kanban): probe claim owner portably on Windows --- hermes_cli/kanban_db.py | 11 ++++-- ...test_kanban_reclaim_unprovable_liveness.py | 38 +++++++++---------- 2 files changed, 26 insertions(+), 23 deletions(-) diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index 45e7ead5a2976..b9f2880612822 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 @@ -13497,13 +13498,15 @@ def _terminate_reclaimed_worker( claimer_pid = 0 try: claimer_pid = int(str(claim_lock)[len(host_prefix):]) - if claimer_pid > 0 and hasattr(os, "kill"): - os.kill(claimer_pid, 0) - except ProcessLookupError: + if claimer_pid > 0: + # Equivalent to kill(pid, 0), without Windows' destructive + # CTRL_C_EVENT behavior for signal 0. + psutil.Process(claimer_pid) + except psutil.NoSuchProcess: info["liveness_unprovable"] = False info["terminated"] = True info["claimer_pid_dead"] = claimer_pid - except (ValueError, OSError): + except (ValueError, OSError, psutil.AccessDenied): pass return info diff --git a/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py b/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py index 471f18bad1e53..4883cb39bb18c 100644 --- a/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py +++ b/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py @@ -28,6 +28,7 @@ import json import os from pathlib import Path +import psutil import secrets import time @@ -117,11 +118,11 @@ def test_dead_local_claimer_without_worker_pid_is_reclaimed(conn, monkeypatch): lock = f"{kb._claimer_id().split(':', 1)[0]}:{dead_pid}" tid, _, run_id = _running_card(conn, worker_pid=None, lock=lock) - def dead_claimer(pid, sig): - assert (pid, sig) == (dead_pid, 0) - raise ProcessLookupError(pid) + def dead_claimer(pid): + assert pid == dead_pid + raise psutil.NoSuchProcess(pid) - monkeypatch.setattr(kb.os, "kill", dead_claimer) + 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 @@ -139,11 +140,11 @@ def test_expired_dead_claimer_is_not_renewed(conn, monkeypatch): conn.execute("UPDATE tasks SET claim_expires=? WHERE id=?", (int(time.time()) - 1, tid)) conn.commit() - def dead_claimer(pid, sig): - assert (pid, sig) == (dead_pid, 0) - raise ProcessLookupError(pid) + def dead_claimer(pid): + assert pid == dead_pid + raise psutil.NoSuchProcess(pid) - monkeypatch.setattr(kb.os, "kill", dead_claimer) + 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" @@ -155,15 +156,15 @@ def test_live_local_claimer_without_worker_pid_still_holds_claim(conn, monkeypat lock = f"{kb._claimer_id().split(':', 1)[0]}:{os.getpid()}" tid, _, _ = _running_card(conn, worker_pid=None, lock=lock) seen = [] - original_kill = os.kill + original_process = psutil.Process - def live_claimer(pid, sig): - seen.append((pid, sig)) - return original_kill(pid, sig) + def live_claimer(pid): + seen.append(pid) + return original_process(pid) - monkeypatch.setattr(kb.os, "kill", live_claimer) + monkeypatch.setattr(psutil, "Process", live_claimer) assert kb.reclaim_task(conn, tid) is False - assert seen == [(os.getpid(), 0)] + assert seen == [os.getpid()] assert _row(conn, tid)["claim_lock"] == lock assert _events(conn, tid, "reclaim_refused")[-1]["reason"] == "liveness_unprovable" @@ -176,12 +177,11 @@ 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, sig): - assert sig == 0 - raise error(pid) + def inaccessible(pid): + raise psutil.AccessDenied(pid) if error: - monkeypatch.setattr(kb.os, "kill", inaccessible) + 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") @@ -189,7 +189,7 @@ def inaccessible(pid, sig): 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(kb.os, "kill", lambda *_: pytest.fail("foreign PID was probed")) + 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 From 5121a1ef88ae8fb310dd9a9bf00b77e1de7a5df9 Mon Sep 17 00:00:00 2001 From: Apollo Date: Thu, 24 Sep 2026 12:11:52 -0700 Subject: [PATCH 4/5] fix(kanban): hold pid-less claim when current run has worker evidence A dead claimer does not prove no worker was spawned: _default_spawn's Popen(start_new_session=True) precedes the _set_worker_pid commit, so a gateway killed in that window leaves a detached, heartbeating worker with no stamped pid. Releasing that claim hands the card to a second worker (t_09180e10 double-worker shape). The no-pid branch now looks for heartbeat/spawned events bound to the CURRENT run before probing the claimer; any such evidence keeps the claim liveness_unprovable. Prior-run evidence is ignored so a never-spawned retry still releases. Applies to reclaim_task, release_stale_claims and detect_stale_running (all three pass conn/task_id). Regression: real-subprocess orphan arm (manual/ttl/stale) is RED on 5ad2916979 and GREEN here. --- hermes_cli/kanban_db.py | 32 +++++++--- ...test_kanban_reclaim_unprovable_liveness.py | 62 ++++++++++++++++++- 2 files changed, 85 insertions(+), 9 deletions(-) diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index b9f2880612822..105e13250c289 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -7723,6 +7723,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. @@ -7845,6 +7846,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 @@ -13471,6 +13473,8 @@ def _terminate_reclaimed_worker( 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 @@ -13491,10 +13495,23 @@ def _terminate_reclaimed_worker( info["host_local"] = True if not pid or pid <= 0: - # The local claimer may still be launching the worker. Only an absent - # claimer proves a never-stamped worker cannot be launched later. - # An alive, inaccessible or malformed claimer remains unprovable. + # Popen precedes _set_worker_pid. A dead claimer can leave a detached, + # unstamped worker behind; any heartbeat/spawned event on THIS run is + # evidence that a worker existed, even when its heartbeat is stale. + # Without a run-bound check, a prior attempt's heartbeat would wedge + # a genuinely never-spawned current attempt. info["liveness_unprovable"] = True + if conn is not None and task_id is not None: + run_id = _current_run_id(conn, task_id) + if run_id is not None: + evidence = conn.execute( + "SELECT kind FROM task_events WHERE task_id=? AND run_id=? " + "AND kind IN ('heartbeat', 'spawned') LIMIT 1", + (task_id, run_id), + ).fetchone() + if evidence is not None: + info["unstamped_worker_evidence"] = evidence["kind"] + return info claimer_pid = 0 try: claimer_pid = int(str(claim_lock)[len(host_prefix):]) @@ -13552,10 +13569,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 worker pid is UNKNOWN liveness unless the host-local claimer is - proven dead. Unknown liveness must not release a claim: it could spawn a - second worker beside the first. Proven-dead 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 @@ -14072,7 +14088,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 diff --git a/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py b/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py index 4883cb39bb18c..c6dd5a2ea725d 100644 --- a/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py +++ b/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py @@ -30,6 +30,8 @@ from pathlib import Path import psutil import secrets +import subprocess +import sys import time import pytest @@ -113,7 +115,7 @@ def _row(conn, tid): def test_dead_local_claimer_without_worker_pid_is_reclaimed(conn, monkeypatch): - """A gateway that died before stamping a worker cannot still spawn one.""" + """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) @@ -133,6 +135,64 @@ def dead_claimer(pid): 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) + 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"]) +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}" From 28b437a6712c98c07667a726ec7371c89d94cd84 Mon Sep 17 00:00:00 2001 From: Apollo Date: Thu, 24 Sep 2026 14:21:56 -0700 Subject: [PATCH 5/5] fix(kanban): time-bound dead-claimer release; fail closed without run context Argus r2 (t_9d7e0142): - B2: dashboard _set_status_direct now passes conn/task_id; the no-pid branch fails closed when either is omitted, and an AST guard asserts every gated _terminate_reclaimed_worker call site passes them. - B3: a dead claimer with no current-run worker evidence is held until DEAD_CLAIMER_LAUNCH_BOUND_SECONDS (= claim TTL, 900 s) after the claim; live ledger first-heartbeat max is 708 s. - B5: with evidence, release once the newest heartbeat/spawned event on the current run is older than DEFAULT_CLAIM_HEARTBEAT_MAX_STALE_SECONDS. - B4: set-model test fakes accept the new kwargs. --- hermes_cli/kanban_db.py | 121 ++++++++++--- plugins/kanban/dashboard/plugin_api.py | 1 + .../hermes_cli/test_kanban_batch_set_model.py | 4 +- ...test_kanban_reclaim_unprovable_liveness.py | 171 ++++++++++++++++++ 4 files changed, 269 insertions(+), 28 deletions(-) diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index 105e13250c289..20ac311c64f65 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -458,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. @@ -13468,6 +13476,56 @@ 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], @@ -13495,36 +13553,43 @@ def _terminate_reclaimed_worker( info["host_local"] = True if not pid or pid <= 0: - # Popen precedes _set_worker_pid. A dead claimer can leave a detached, - # unstamped worker behind; any heartbeat/spawned event on THIS run is - # evidence that a worker existed, even when its heartbeat is stale. - # Without a run-bound check, a prior attempt's heartbeat would wedge - # a genuinely never-spawned current attempt. + # 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 not None and task_id is not None: - run_id = _current_run_id(conn, task_id) - if run_id is not None: - evidence = conn.execute( - "SELECT kind FROM task_events WHERE task_id=? AND run_id=? " - "AND kind IN ('heartbeat', 'spawned') LIMIT 1", - (task_id, run_id), - ).fetchone() - if evidence is not None: - info["unstamped_worker_evidence"] = evidence["kind"] - return info + 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: - # Equivalent to kill(pid, 0), without Windows' destructive - # CTRL_C_EVENT behavior for signal 0. - psutil.Process(claimer_pid) + 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: - info["liveness_unprovable"] = False - info["terminated"] = True - info["claimer_pid_dead"] = claimer_pid - except (ValueError, OSError, psutil.AccessDenied): 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 ( @@ -13993,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", @@ -15334,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 c6dd5a2ea725d..78049b486510b 100644 --- a/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py +++ b/tests/hermes_cli/test_kanban_reclaim_unprovable_liveness.py @@ -114,11 +114,62 @@ 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 @@ -132,6 +183,7 @@ def dead_claimer(pid): 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") @@ -139,6 +191,7 @@ def test_heartbeat_from_previous_run_does_not_hold_dead_claimer(conn, monkeypatc 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): @@ -147,6 +200,123 @@ def dead(pid): 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. @@ -197,6 +367,7 @@ 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()