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/<pid>/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
  "<gateway.drain_control.current_instantiation_epoch()>|<start>" (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".
This commit is contained in:
teknium1
2026-09-16 00:08:46 -07:00
committed by Teknium
parent 8609743389
commit db54f5448d
3 changed files with 157 additions and 34 deletions
+12 -9
View File
@@ -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 ("<boot/instantiation epoch>|<start time>",
-- 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:
+74 -25
View File
@@ -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: ``"<instantiation epoch>|<start time>"``. The start
time alone (``/proc/<pid>/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))
@@ -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"