fix(kanban): give a finished worker a grace window before the terminal reaper signals it
reap_terminal_workers signalled any retained worker on the first tick after its run closed, but a healthy worker is still alive for a moment after kanban_complete / kanban_request_review returns (final assistant turn, session persistence), so slow-but-healthy workers were killed mid-finalisation and logged as terminal_worker_reaped. Reap only runs whose ended_at is at least TERMINAL_WORKER_REAP_GRACE_SECONDS (120 s, two default ticks) old; the fingerprint check is unchanged. Each row is now handled on its own so a signal or /proc failure on one run is logged and skips only that run. Tests: a just-closed run is not signalled and keeps its evidence, then is reaped once the grace has passed (red before); one raising row no longer aborts the sweep for the others (red before).
This commit is contained in:
@@ -41,6 +41,12 @@ DEFAULT_LOG_BACKUP_COUNT = 1
|
||||
# and call kanban_block/kanban_complete before max_runtime_seconds kills it.
|
||||
KANBAN_TERMINAL_TIMEOUT_GRACE_SECONDS = 30
|
||||
|
||||
# A healthy worker is still alive for a while after kanban_complete /
|
||||
# kanban_request_review returns (final assistant turn, session persistence), so
|
||||
# a run's retained worker is only reaped once ended_at is at least this old
|
||||
# (two default dispatch ticks).
|
||||
TERMINAL_WORKER_REAP_GRACE_SECONDS = 120
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Respawn guard constants
|
||||
# ---------------------------------------------------------------------------
|
||||
@@ -398,41 +404,56 @@ def reap_terminal_workers(conn: sqlite3.Connection, *, signal_fn=None) -> list[s
|
||||
fds open and no ``running``-only sweep can see it once ``tasks.worker_pid`` is
|
||||
cleared. Keys on the closed ``task_runs`` row's retained pid + spawn
|
||||
fingerprint: a legacy row (NULL fingerprint) or a recycled PID is never
|
||||
signalled; a pid that is simply gone just has its evidence cleared. Returns
|
||||
the task ids whose worker was terminated."""
|
||||
signalled; a pid that is simply gone just has its evidence cleared. A run
|
||||
that ended less than ``TERMINAL_WORKER_REAP_GRACE_SECONDS`` ago is left
|
||||
alone so a worker still finalising after its own transition is not killed.
|
||||
One row's failure (signal, /proc probe) is logged and skips only that row.
|
||||
Returns the task ids whose worker was terminated."""
|
||||
rows = conn.execute(
|
||||
"SELECT id, task_id, worker_pid, worker_started_at, claim_lock FROM task_runs "
|
||||
"WHERE ended_at IS NOT NULL AND worker_pid IS NOT NULL AND worker_started_at IS NOT NULL"
|
||||
"WHERE ended_at IS NOT NULL AND ended_at <= ? "
|
||||
"AND worker_pid IS NOT NULL AND worker_started_at IS NOT NULL",
|
||||
(int(time.time()) - TERMINAL_WORKER_REAP_GRACE_SECONDS,),
|
||||
).fetchall()
|
||||
host_prefix = _kb._host_prefix()
|
||||
reaped: list[str] = []
|
||||
for row in rows:
|
||||
pid, fingerprint = int(row["worker_pid"]), int(row["worker_started_at"])
|
||||
if pid == os.getpid() or not str(row["claim_lock"] or "").startswith(host_prefix):
|
||||
continue
|
||||
alive = _worker_alive(pid, fingerprint)
|
||||
termination = None
|
||||
if alive:
|
||||
termination = _terminate_reclaimed_worker(
|
||||
pid, row["claim_lock"], signal_fn=signal_fn, started_at=fingerprint)
|
||||
if not termination["terminated"]:
|
||||
continue # still alive: try again next tick
|
||||
with _kb.write_txn(conn):
|
||||
conn.execute(
|
||||
"UPDATE task_runs SET worker_pid = NULL, worker_started_at = NULL "
|
||||
"WHERE id = ? AND worker_pid = ? AND worker_started_at = ?",
|
||||
(row["id"], pid, fingerprint),
|
||||
try:
|
||||
_reap_terminal_worker_row(conn, row, host_prefix, signal_fn, reaped)
|
||||
except Exception:
|
||||
_kb._log.debug(
|
||||
"kanban dispatch: terminal worker reap failed for run %s (task %s)",
|
||||
row["id"], row["task_id"], exc_info=True,
|
||||
)
|
||||
if alive:
|
||||
_kb._append_event(
|
||||
conn, row["task_id"], "terminal_worker_reaped",
|
||||
{"pid": pid, "worker_started_at": fingerprint, **termination}, run_id=row["id"],
|
||||
)
|
||||
if alive:
|
||||
reaped.append(row["task_id"])
|
||||
return reaped
|
||||
|
||||
|
||||
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"])
|
||||
if pid == os.getpid() or not str(row["claim_lock"] or "").startswith(host_prefix):
|
||||
return
|
||||
alive = _worker_alive(pid, fingerprint)
|
||||
termination = None
|
||||
if alive:
|
||||
termination = _terminate_reclaimed_worker(
|
||||
pid, row["claim_lock"], signal_fn=signal_fn, started_at=fingerprint)
|
||||
if not termination["terminated"]:
|
||||
return # still alive: try again next tick
|
||||
with _kb.write_txn(conn):
|
||||
conn.execute(
|
||||
"UPDATE task_runs SET worker_pid = NULL, worker_started_at = NULL "
|
||||
"WHERE id = ? AND worker_pid = ? AND worker_started_at = ?",
|
||||
(row["id"], pid, fingerprint),
|
||||
)
|
||||
if alive:
|
||||
_kb._append_event(
|
||||
conn, row["task_id"], "terminal_worker_reaped",
|
||||
{"pid": pid, "worker_started_at": fingerprint, **termination}, run_id=row["id"],
|
||||
)
|
||||
if alive:
|
||||
reaped.append(row["task_id"])
|
||||
|
||||
|
||||
def _worker_survived_termination(termination: dict) -> bool:
|
||||
"""True when we tried to kill our own host-local worker and it is still alive.
|
||||
|
||||
|
||||
@@ -11,6 +11,7 @@ the next tick — never a recycled PID, never a legacy row without a fingerprint
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
import os
|
||||
import subprocess
|
||||
import sys
|
||||
import time
|
||||
@@ -45,12 +46,14 @@ def _sleeper():
|
||||
return proc
|
||||
|
||||
|
||||
def _completed_card_with_worker(conn, proc) -> tuple[str, int]:
|
||||
def _completed_card_with_worker(conn, proc, *, ended_ago: int = 600) -> tuple[str, int]:
|
||||
tid = kb.create_task(conn, title="finished", assignee="coder")
|
||||
kb.claim_task(conn, tid, claimer=kb._claimer_id())
|
||||
run_id = kb._current_run_id(conn, tid)
|
||||
kbd._set_worker_pid(conn, tid, proc.pid)
|
||||
assert kb.complete_task(conn, tid, result="done", expected_run_id=run_id) is True
|
||||
# Default: the run closed long enough ago that the reaper's grace window has passed.
|
||||
conn.execute("UPDATE task_runs SET ended_at = ended_at - ? WHERE id=?", (ended_ago, run_id))
|
||||
return tid, run_id
|
||||
|
||||
|
||||
@@ -96,3 +99,51 @@ def test_recycled_pid_and_legacy_row_are_never_signalled(conn):
|
||||
for p in (stranger, legacy):
|
||||
p.kill()
|
||||
p.wait()
|
||||
|
||||
|
||||
def test_fresh_terminal_run_is_left_alone_until_grace_passes(conn):
|
||||
"""A worker is still alive for a moment after its own kanban_complete returns
|
||||
(final turn, session persistence): a just-closed run is not signalled."""
|
||||
proc = _sleeper()
|
||||
signals = []
|
||||
|
||||
def signal_fn(pid, sig):
|
||||
signals.append((pid, sig))
|
||||
os.kill(pid, sig)
|
||||
|
||||
try:
|
||||
tid, run_id = _completed_card_with_worker(conn, proc, ended_ago=0)
|
||||
|
||||
assert kbd.reap_terminal_workers(conn, signal_fn=signal_fn) == []
|
||||
|
||||
assert signals == [] and proc.poll() is None
|
||||
run = conn.execute("SELECT worker_pid FROM task_runs WHERE id=?", (run_id,)).fetchone()
|
||||
assert run["worker_pid"] == proc.pid # evidence kept for a later tick
|
||||
conn.execute(
|
||||
"UPDATE task_runs SET ended_at = ended_at - ? WHERE id=?",
|
||||
(kbd.TERMINAL_WORKER_REAP_GRACE_SECONDS, run_id),
|
||||
)
|
||||
assert kbd.reap_terminal_workers(conn, signal_fn=signal_fn) == [tid]
|
||||
assert [pid for pid, _ in signals] == [proc.pid]
|
||||
finally:
|
||||
proc.kill()
|
||||
proc.wait()
|
||||
|
||||
|
||||
def test_one_failing_row_does_not_abort_the_sweep(conn):
|
||||
"""A signal failure on one run is logged and skipped; the other rows are still reaped."""
|
||||
broken, healthy = _sleeper(), _sleeper()
|
||||
try:
|
||||
_completed_card_with_worker(conn, broken)
|
||||
tid, _ = _completed_card_with_worker(conn, healthy)
|
||||
|
||||
def signal_fn(pid, sig):
|
||||
if pid == broken.pid:
|
||||
raise RuntimeError("boom")
|
||||
healthy.kill()
|
||||
|
||||
assert kbd.reap_terminal_workers(conn, signal_fn=signal_fn) == [tid]
|
||||
finally:
|
||||
for p in (broken, healthy):
|
||||
p.kill()
|
||||
p.wait()
|
||||
|
||||
@@ -123,7 +123,7 @@ They coexist: a kanban worker may call `delegate_task` internally during its run
|
||||
- `scratch` (default) — fresh tmp dir under `~/.hermes/kanban/workspaces/<id>/` (or `~/.hermes/kanban/boards/<slug>/workspaces/<id>/` on non-default boards). **Deleted when the task completes** — scratch is ephemeral by design. Files explicitly declared through `kanban_complete(artifacts=[...])` or `kanban_request_review(artifacts=[...])` are copied into durable per-task attachment storage before cleanup (a review handoff stages them at handoff time, since the reviewer's later completion is what deletes the scratch workspace); existing deliverable paths in legacy completion summaries receive the same treatment. Other scratch files are removed. A missing declared scratch artifact keeps the task in-flight so the worker can correct the path and retry. Use `worktree:` or `dir:<path>` when the whole workspace should remain available. The first time a scratch workspace is created on an install, the dispatcher logs a warning and emits a `tip_scratch_workspace` event on the task (visible via `hermes kanban show <id>`).
|
||||
- `dir:<path>` — an existing shared directory (Obsidian vault, mail ops dir, per-account folder). **Must be an absolute path.** Relative paths like `dir:../tenants/foo/` are rejected at dispatch because they'd resolve against whatever CWD the dispatcher happens to be in, which is ambiguous and a confused-deputy escape vector. The path is otherwise trusted — it's your box, your filesystem, the worker runs with your uid. This is the trusted-local-user threat model; kanban is single-host by design. **Preserved on completion.**
|
||||
- `worktree` — a git worktree under `.worktrees/<id>/` for coding tasks. Use `worktree:<path>` to pin the exact target path. Worker-side `git worktree add` creates it, using `--branch` when provided. **Preserved on completion.**
|
||||
- **Dispatcher** — a long-lived loop that, every N seconds (default 60): reclaims stale claims, reclaims crashed workers (PID gone but TTL not yet expired), reaps workers that outlived their finished run (a worker still alive after its own `kanban_complete`/`kanban_block` is terminated on the next tick — matched by PID *and* spawn-time fingerprint, so a recycled PID is never signalled; recorded as a `terminal_worker_reaped` event), promotes ready tasks, atomically claims, spawns assigned profiles. Runs **inside the gateway** by default (`kanban.dispatch_in_gateway: true`). One dispatcher sweeps all boards per tick; workers are spawned with `HERMES_KANBAN_BOARD` pinned so they can't see other boards. After `kanban.failure_limit` consecutive spawn failures on the same task (default: 2) the dispatcher auto-blocks it with the last error as the reason — prevents thrashing on tasks whose profile doesn't exist, workspace can't mount, etc.
|
||||
- **Dispatcher** — a long-lived loop that, every N seconds (default 60): reclaims stale claims, reclaims crashed workers (PID gone but TTL not yet expired), reaps workers that outlived their finished run (a worker still alive after its own `kanban_complete`/`kanban_block` is terminated once its run has been closed for two minutes, leaving it time to finish its final turn — matched by PID *and* spawn-time fingerprint, so a recycled PID is never signalled; recorded as a `terminal_worker_reaped` event), promotes ready tasks, atomically claims, spawns assigned profiles. Runs **inside the gateway** by default (`kanban.dispatch_in_gateway: true`). One dispatcher sweeps all boards per tick; workers are spawned with `HERMES_KANBAN_BOARD` pinned so they can't see other boards. After `kanban.failure_limit` consecutive spawn failures on the same task (default: 2) the dispatcher auto-blocks it with the last error as the reason — prevents thrashing on tasks whose profile doesn't exist, workspace can't mount, etc.
|
||||
- **Tenant** — optional string namespace *within* a board. One specialist fleet can serve multiple businesses (`--tenant business-a`) with data isolation by workspace path and memory key prefix. Tenants are a soft filter; boards are the hard isolation boundary.
|
||||
|
||||
## Boards (multi-project)
|
||||
|
||||
Reference in New Issue
Block a user