From db54f5448d7cd13dda34630564153e2edadcfe63 Mon Sep 17 00:00:00 2001 From: teknium1 <127238744+teknium1@users.noreply.github.com> Date: Wed, 16 Sep 2026 00:08:46 -0700 Subject: [PATCH] fix(kanban): worker fingerprint carries a boot witness; an uncaptured fingerprint never authorizes a signal #111617 review (andrexibiza P1 #3/#4, kvnloo nit): - worker_started_at persisted only gateway.status.get_process_start_time(): on Linux that is /proc//stat field 22, clock ticks since THIS boot. The threat is a row surviving a reboot, and that counter does not, so an unrelated process on a later boot with the same PID and the same tick value passed _start_times_agree(). The fingerprint is now "|" (boot_id + PID-1 start, the witness the drain marker already uses); both halves must match. Integer values on rows written before this change keep the start-time-only comparison. - A failed capture persisted NULL, which _pid_recycled treats as the legacy pre-fingerprint row and falls back to bare PID existence - a new spawn silently recreated the #89614/ #99558 kill authority. A failed capture now persists UNVERIFIED_WORKER_FINGERPRINT: the claim is held while the PID is live (never released beside it, never SIGTERM/SIGKILLed by timeout, stale-claim, manual reclaim, archive or the terminal reaper) and reclaimed once it is gone. NULL stays legacy-only. - Every tasks UPDATE that nulls worker_pid nulls worker_started_at too (archive_task and the reclaim/timeout/reopen paths): the fingerprint is part of the kill-authority tuple and must not outlive its pid. Live (real sleeper child): reboot-shaped row (same pid, same tick, other boot id) -> reclaimed to ready, child untouched; matching fingerprint -> SIGTERM delivered, exit -15. tests/hermes_cli/test_kanban_worker_pid_fingerprint.py: +2 hostile tests, both red on base. Not changed: the check-then-act window between _pid_recycled and kill (kvnloo P2) is real but needs pidfd_open/pidfd_send_signal (Linux 5.3+) to close atomically; left as the documented residual of "never kills a DETECTED recycled PID". --- hermes_cli/kanban_db.py | 21 ++-- hermes_cli/kanban_db_dispatch.py | 99 ++++++++++++++----- .../test_kanban_worker_pid_fingerprint.py | 71 +++++++++++++ 3 files changed, 157 insertions(+), 34 deletions(-) diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index 04198a301b..95e77e44de 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -890,9 +890,12 @@ CREATE TABLE IF NOT EXISTS tasks ( -- exceeds DEFAULT_FAILURE_LIMIT consecutive non-successes. consecutive_failures INTEGER NOT NULL DEFAULT 0, worker_pid INTEGER, - -- Start-time fingerprint of worker_pid (gateway.status.get_process_start_time) recorded at - -- spawn: liveness and kills require pid AND fingerprint to agree, so a PID recycled after a - -- reboot is never read as our worker or signalled. NULL = legacy row (pre-fingerprint spawn). + -- Restart-stable fingerprint of worker_pid ("|", + -- kanban_db_dispatch._process_fingerprint) recorded at spawn: liveness and kills require pid + -- AND fingerprint to agree, so a PID recycled after a reboot is never read as our worker or + -- signalled. NULL = legacy row (pre-fingerprint spawn); 'unverified' = capture failed at + -- spawn (held while live, never signalled). Column keeps its INTEGER affinity for the + -- start-time-only integer values older rows carry. worker_started_at INTEGER, -- Short excerpt of the most recent failure's error text. last_failure_error TEXT, @@ -2412,7 +2415,7 @@ def release_stale_claims( retry_status = _retry_status_for_run(conn, row["id"]) cur = conn.execute( "UPDATE tasks SET status = ?, claim_lock = NULL, " - "claim_expires = NULL, worker_pid = NULL " + "claim_expires = NULL, worker_pid = NULL, worker_started_at = NULL " "WHERE id = ? AND status = 'running' AND claim_lock IS ? " "AND claim_expires IS NOT NULL AND claim_expires < ?", (retry_status, row["id"], row["claim_lock"], now), @@ -2516,7 +2519,7 @@ def reclaim_task( retry_status = _retry_status_for_run(conn, task_id) cur = conn.execute( "UPDATE tasks SET status = ?, claim_lock = NULL, " - "claim_expires = NULL, worker_pid = NULL " + "claim_expires = NULL, worker_pid = NULL, worker_started_at = NULL " "WHERE id = ? AND status IN ('running', 'ready', 'blocked') " "AND claim_lock IS ?", (retry_status, task_id, prev_lock), ) @@ -3327,7 +3330,7 @@ def request_changes( assignee = COALESCE(?, assignee), claim_lock = NULL, claim_expires = NULL, - worker_pid = NULL + worker_pid = NULL, worker_started_at = NULL WHERE id = ? AND status = 'running' AND current_run_id = ? """, (new_status, implementer, task_id, int(current_run_id)), @@ -3495,7 +3498,7 @@ def reopen_review_task(conn: sqlite3.Connection, task_id: str) -> bool: # consecutive_failures deliberately PRESERVED: review reopen is not # a success signal; only complete_task resets the breaker (#35072). "UPDATE tasks SET status = ?, current_run_id = NULL, " - "claim_lock = NULL, claim_expires = NULL, worker_pid = NULL " + "claim_lock = NULL, claim_expires = NULL, worker_pid = NULL, worker_started_at = NULL " + (", assignee = ?" if implementer else "") + " WHERE id = ? AND status = 'review'", params, @@ -3570,7 +3573,7 @@ def invalidate_descendants_for_parent_reopen( # docstring for why this diverges from reopen_review_task. conn.execute( "UPDATE tasks SET status = 'todo', completed_at = NULL, " - "claim_lock = NULL, claim_expires = NULL, worker_pid = NULL, " + "claim_lock = NULL, claim_expires = NULL, worker_pid = NULL, worker_started_at = NULL, " "current_run_id = NULL, consecutive_failures = 0 WHERE id = ?", (row["id"],), ) entry = { @@ -3688,7 +3691,7 @@ def archive_task(conn: sqlite3.Connection, task_id: str, *, signal_fn=None) -> b prev_pid, prev_lock, prev_started = row["worker_pid"], row["claim_lock"], row["worker_started_at"] cur = conn.execute( "UPDATE tasks SET status = 'archived', " - " claim_lock = NULL, claim_expires = NULL, worker_pid = NULL " + " claim_lock = NULL, claim_expires = NULL, worker_pid = NULL, worker_started_at = NULL " "WHERE id = ? AND status != 'archived'", (task_id,), ) if cur.rowcount != 1: diff --git a/hermes_cli/kanban_db_dispatch.py b/hermes_cli/kanban_db_dispatch.py index a99382e79c..f8e32c1c61 100644 --- a/hermes_cli/kanban_db_dispatch.py +++ b/hermes_cli/kanban_db_dispatch.py @@ -295,20 +295,52 @@ def _pid_alive(pid: Optional[int]) -> bool: return True -def _worker_alive(pid: Optional[int], started_at: Optional[int]) -> bool: - """True when ``pid`` is live AND is still the worker we spawned. ``started_at`` is the start-time - fingerprint recorded by ``_set_worker_pid``; after a reboot (or any PID recycle) an unrelated - process can own the number, so bare existence is never enough to extend a claim or to signal. - A legacy row without a fingerprint keeps the existence answer: killing it is the pre-fingerprint - behaviour and the row is rewritten with a fingerprint on its next spawn.""" - return _kb._pid_alive(pid) and not _pid_recycled(pid, started_at) +# ``worker_started_at`` value for a spawn whose fingerprint could not be captured. Distinct from the +# NULL legacy row (pre-fingerprint spawn): such a worker is held (its claim is never released beside +# the live PID) but NEVER signalled — missing process identity is refusal, not permission (#99558). +UNVERIFIED_WORKER_FINGERPRINT = "unverified" -def _pid_recycled(pid: Optional[int], started_at: Optional[int]) -> bool: +def _process_fingerprint(pid: int) -> Optional[str]: + """Restart-stable identity of a live process: ``"|"``. The start + time alone (``/proc//stat`` field 22 on Linux) is clock ticks since THIS boot, so a row that + survives a reboot could match an unrelated process with the same PID and the same tick value; + ``gateway.drain_control.current_instantiation_epoch`` (``boot_id`` + PID-1 start) changes on every + reboot / container recreate, so the composed value never survives one. ``None`` when unreadable.""" + from gateway.drain_control import current_instantiation_epoch + from gateway.status import get_process_start_time + start = get_process_start_time(int(pid)) + if start is None: + return None + return f"{current_instantiation_epoch()}|{start}" + + +def _worker_alive(pid: Optional[int], started_at) -> bool: + """True when ``pid`` is live AND is still the worker we spawned. ``started_at`` is the fingerprint + recorded by ``_set_worker_pid``; after a reboot (or any PID recycle) an unrelated process can own + the number, so bare existence is never enough to extend a claim or to signal. A legacy row without + a fingerprint keeps the existence answer: killing it is the pre-fingerprint behaviour and the row is + rewritten with a fingerprint on its next spawn. An UNVERIFIED spawn also keeps the existence answer + (a claim is never released beside a possibly-live worker) but ``_terminate_reclaimed_worker`` + refuses to signal it.""" + if not _kb._pid_alive(pid): + return False + if started_at == UNVERIFIED_WORKER_FINGERPRINT: + return True + return not _pid_recycled(pid, started_at) + + +def _pid_recycled(pid: Optional[int], started_at) -> bool: """True when a live ``pid`` is NOT the process fingerprinted at spawn (or the fingerprint can no - longer be read). Signalling it would hit a stranger. ``None`` fingerprint = legacy row, never recycled.""" + longer be read). Signalling it would hit a stranger. ``None`` fingerprint = legacy row, never + recycled; the UNVERIFIED marker is always foreign. An integer fingerprint (rows written before the + boot witness was added) compares the start time only.""" if started_at is None or not pid: return False + if started_at == UNVERIFIED_WORKER_FINGERPRINT: + return True + if isinstance(started_at, str) and "|" in started_at: + return _process_fingerprint(int(pid)) != started_at from gateway.status import _start_times_agree, get_process_start_time current = get_process_start_time(int(pid)) if current is None: @@ -350,11 +382,14 @@ def _terminate_reclaimed_worker( claim_lock: Optional[str], *, signal_fn=None, - started_at: Optional[int] = None, + started_at=None, ) -> dict[str, Any]: """Best-effort host-local worker termination for reclaim paths. ``started_at`` is the spawn-time fingerprint: when the live process no longer matches it, the PID was recycled and nothing is - signalled — the worker is gone, which is what the reclaim wanted (``terminated`` = True).""" + signalled — the worker is gone, which is what the reclaim wanted (``terminated`` = True). An + UNVERIFIED spawn (fingerprint capture failed) that is still live is never signalled either, but + it is reported as surviving (``signal_refused``) so the reclaim holds the claim instead of + spawning a duplicate beside it.""" info: dict[str, Any] = { "prev_pid": int(pid) if pid else None, "host_local": False, @@ -371,6 +406,11 @@ def _terminate_reclaimed_worker( kill = _kill_fn(signal_fn) if kill is None: return info + if started_at == UNVERIFIED_WORKER_FINGERPRINT: + # Never signal by bare number: a dead PID is "gone" (reclaim proceeds), a live one is held. + info["signal_refused"] = True + info["terminated"] = not _kb._pid_alive(pid) + return info if _kb._pid_alive(pid) and _pid_recycled(pid, started_at): info["terminated"] = True info["pid_recycled"] = True @@ -429,9 +469,11 @@ def reap_terminal_workers(conn: sqlite3.Connection, *, signal_fn=None) -> list[s def _reap_terminal_worker_row(conn, row, host_prefix: str, signal_fn, reaped: list[str]) -> None: - pid, fingerprint = int(row["worker_pid"]), int(row["worker_started_at"]) + pid, fingerprint = int(row["worker_pid"]), row["worker_started_at"] if pid == os.getpid() or not str(row["claim_lock"] or "").startswith(host_prefix): return + if fingerprint == UNVERIFIED_WORKER_FINGERPRINT and _kb._pid_alive(pid): + return # unproven identity: never signalled; its evidence is cleared once the pid is gone alive = _worker_alive(pid, fingerprint) termination = None if alive: @@ -463,8 +505,8 @@ def _worker_survived_termination(termination: dict) -> bool: through to the normal release path since we cannot manage that worker anyway. """ return bool( - termination.get("termination_attempted") - and termination.get("host_local") + termination.get("host_local") + and (termination.get("termination_attempted") or termination.get("signal_refused")) and not termination.get("terminated") ) @@ -576,6 +618,12 @@ def enforce_max_runtime(conn: sqlite3.Connection, *, signal_fn=None) -> list[str pid = int(row["worker_pid"]) tid = row["id"] started_at = _kb._row_get(row, "worker_started_at") + if started_at == UNVERIFIED_WORKER_FINGERPRINT and _kb._pid_alive(pid): + # Fingerprint capture failed at spawn: we cannot prove this live PID is our worker, so + # it is neither signalled nor released beside (duplicate). It is reclaimed once it exits. + _kb._log.warning("kanban: task %s worker pid %s exceeded max runtime but has no verified " + "identity; not signalled", tid, pid) + continue # SIGTERM then SIGKILL after 5 s grace; workers wanting a cleaner # shutdown install their own SIGTERM handler. A recycled PID (fingerprint # mismatch) is never signalled: the worker is already gone. @@ -594,7 +642,7 @@ def enforce_max_runtime(conn: sqlite3.Connection, *, signal_fn=None) -> list[str retry_status = _kb._retry_status_for_run(conn, tid) cur = conn.execute( "UPDATE tasks SET status = ?, claim_lock = NULL, " - "claim_expires = NULL, worker_pid = NULL, " + "claim_expires = NULL, worker_pid = NULL, worker_started_at = NULL, " "last_heartbeat_at = NULL " "WHERE id = ? AND status = 'running' " " AND worker_pid = ? AND claim_lock IS ?", @@ -696,7 +744,7 @@ def detect_stale_running( retry_status = _kb._retry_status_for_run(conn, tid) cur = conn.execute( "UPDATE tasks SET status = ?, claim_lock = NULL, " - "claim_expires = NULL, worker_pid = NULL, " + "claim_expires = NULL, worker_pid = NULL, worker_started_at = NULL, " "last_heartbeat_at = NULL " "WHERE id = ? AND status = 'running' " " AND claim_lock IS ?", @@ -761,7 +809,7 @@ def reconcile_orphaned_running(conn: sqlite3.Connection) -> list[str]: with _kb.write_txn(conn): cur = conn.execute( "UPDATE tasks SET status = 'ready', claim_lock = NULL, " - "claim_expires = NULL, worker_pid = NULL, " + "claim_expires = NULL, worker_pid = NULL, worker_started_at = NULL, " "last_heartbeat_at = NULL " "WHERE id = ? AND status = 'running' " " AND claim_lock IS ? AND claim_expires IS ?", @@ -1015,7 +1063,7 @@ def _reclaim_dead_workers(conn: sqlite3.Connection, board: Optional[str] = None) dead.event_payload["retry_status"] = retry_status cur = conn.execute( "UPDATE tasks SET status = ?, claim_lock = NULL, " - "claim_expires = NULL, worker_pid = NULL " + "claim_expires = NULL, worker_pid = NULL, worker_started_at = NULL " "WHERE id = ? AND status = 'running' " " AND worker_pid = ? AND claim_lock IS ?", (retry_status, row["id"], pid, row["claim_lock"]), @@ -1221,7 +1269,7 @@ def _record_task_failure( # Spawn path: restore the claimed source phase + clear claim. conn.execute( "UPDATE tasks SET status = ?, claim_lock = NULL, " - "claim_expires = NULL, worker_pid = NULL, " + "claim_expires = NULL, worker_pid = NULL, worker_started_at = NULL, " "consecutive_failures = ?, last_failure_error = ? " "WHERE id = ? AND status = 'running'", (retry_status, failures, error, task_id), @@ -1249,7 +1297,7 @@ def _record_task_failure( # state; the timeout/crash path already did. conn.execute( "UPDATE tasks SET status = 'blocked', " - + ("claim_lock = NULL, claim_expires = NULL, worker_pid = NULL, " + + ("claim_lock = NULL, claim_expires = NULL, worker_pid = NULL, worker_started_at = NULL, " if release_claim else "") + "consecutive_failures = ?, last_failure_error = ? " "WHERE id = ? AND status IN ('running', 'ready', 'review')", @@ -1283,11 +1331,12 @@ def _record_task_failure( def _set_worker_pid(conn: sqlite3.Connection, task_id: str, pid: int) -> None: - """Record the spawned child's pid + its start-time fingerprint, and emit a ``spawned`` event - carrying them. The fingerprint is what lets every later liveness/kill decision tell OUR worker - from a process that recycled the PID after a reboot.""" - from gateway.status import get_process_start_time - started_at = get_process_start_time(int(pid)) + """Record the spawned child's pid + its restart-stable fingerprint (``_process_fingerprint``), and + emit a ``spawned`` event carrying them. The fingerprint is what lets every later liveness/kill + decision tell OUR worker from a process that recycled the PID after a reboot. A failed capture is + persisted as ``UNVERIFIED_WORKER_FINGERPRINT``, never NULL: NULL is the legacy pre-fingerprint row + whose bare-PID kill authority a new spawn must not inherit.""" + started_at = _process_fingerprint(int(pid)) or UNVERIFIED_WORKER_FINGERPRINT with _kb.write_txn(conn): conn.execute("UPDATE tasks SET worker_pid = ?, worker_started_at = ? WHERE id = ?", (int(pid), started_at, task_id)) diff --git a/tests/hermes_cli/test_kanban_worker_pid_fingerprint.py b/tests/hermes_cli/test_kanban_worker_pid_fingerprint.py index a7fc571987..d49d54400c 100644 --- a/tests/hermes_cli/test_kanban_worker_pid_fingerprint.py +++ b/tests/hermes_cli/test_kanban_worker_pid_fingerprint.py @@ -78,3 +78,74 @@ def test_matching_fingerprint_keeps_the_live_worker(board): conn.execute("UPDATE tasks SET max_runtime_seconds = 1 WHERE id = ?", (tid,)) kbd.enforce_max_runtime(conn, signal_fn=lambda pid, sig: killed.append((pid, sig))) assert killed and killed[0] == (os.getpid(), signal.SIGTERM) + + +def test_same_pid_and_start_tick_on_another_boot_is_foreign(board, monkeypatch): + """A row that survived a reboot: the PID AND the boot-relative start tick both match a process on + this boot (the Linux start time is clock ticks since boot, so that recurs), but the persisted + instantiation epoch does not. The worker is foreign: claim released, zero signals.""" + from gateway import drain_control + + conn = board + killed = [] + live_fingerprint = kbd._process_fingerprint(os.getpid()) + assert live_fingerprint is not None and live_fingerprint.split("|", 1)[1] == str( + __import__("gateway.status", fromlist=["x"]).get_process_start_time(os.getpid())) + tid = _claimed_running(conn, pid=os.getpid(), started_at=live_fingerprint, max_runtime=1) + assert kbd._worker_alive(os.getpid(), live_fingerprint) is True + + # Same PID, same start tick, different boot identity. + other_boot = "deadbeef-boot:1|" + live_fingerprint.split("|", 1)[1] + with kb.write_txn(conn): + conn.execute("UPDATE tasks SET worker_started_at = ? WHERE id = ?", (other_boot, tid)) + assert kbd._worker_alive(os.getpid(), other_boot) is False + assert tid in kbd.enforce_max_runtime(conn, signal_fn=lambda pid, sig: killed.append((pid, sig))) + assert killed == [] + task = kb.get_task(conn, tid) + assert task.status == "ready" and task.worker_pid is None + + # The same value re-derived on THIS boot still identifies our worker (the witness is stable + # within a boot, unlike the recorded epoch of a previous one). + drain_control.current_instantiation_epoch.cache_clear() + assert kbd._process_fingerprint(os.getpid()) == live_fingerprint + + +def test_unverified_fingerprint_capture_never_authorizes_a_signal(board, monkeypatch): + """Fingerprint capture fails for a new spawn: the row is NOT a legacy NULL row. A live PID under + it is never SIGTERM/SIGKILLed by any reclaim/timeout path, and the claim is held (not released + beside the live process); once the PID is gone the claim is reclaimed normally.""" + import gateway.status as status + + conn = board + killed = [] + monkeypatch.setattr(status, "_get_process_start_time", lambda pid: None) + tid = kb.create_task(conn, title="job", assignee="worker", max_runtime_seconds=1) + kb.claim_task(conn, tid) + kbd._set_worker_pid(conn, tid, os.getpid()) + row = conn.execute("SELECT worker_started_at FROM tasks WHERE id = ?", (tid,)).fetchone() + assert row["worker_started_at"] == kbd.UNVERIFIED_WORKER_FINGERPRINT + monkeypatch.undo() + old = int(time.time()) - 3600 + with kb.write_txn(conn): + conn.execute("UPDATE tasks SET started_at = ?, claim_expires = ? WHERE id = ?", (old, old, tid)) + conn.execute("UPDATE task_runs SET started_at = ? WHERE id = (SELECT current_run_id FROM tasks WHERE id = ?)", + (old, tid)) + + sig = lambda pid, s: killed.append((pid, s)) # noqa: E731 + assert kbd.enforce_max_runtime(conn, signal_fn=sig) == [] + assert kb.release_stale_claims(conn, signal_fn=sig) == 0 + assert killed == [] + assert kb.get_task(conn, tid).status == "running" + # An explicit operator reclaim releases the claim (human override) but still sends nothing. + assert kb.reclaim_task(conn, tid, reason="operator", signal_fn=sig) is True + assert killed == [] + + # The process is gone (a dead PID): the row is reclaimed like any dead worker, still no signal. + tid2 = kb.create_task(conn, title="job2", assignee="worker") + kb.claim_task(conn, tid2) + with kb.write_txn(conn): + conn.execute("UPDATE tasks SET worker_pid = ?, worker_started_at = ?, claim_expires = ? WHERE id = ?", + (os.getpid(), kbd.UNVERIFIED_WORKER_FINGERPRINT, old, tid2)) + monkeypatch.setattr(kb, "_pid_alive", lambda pid: False) + assert kb.release_stale_claims(conn, signal_fn=sig) == 1 + assert killed == [] and kb.get_task(conn, tid2).status == "ready"