From e66429eb0a067bbcb6a14fc037430e5bb74b8d62 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 21:33:44 -0700 Subject: [PATCH] refactor(kanban_db_dispatch,connect): split god functions into phase helpers, unify probe/kill/profile helpers, AST-neutral line joins - detect_crashed_workers (239) -> _reclaim_dead_workers/_classify_dead_worker(_DeadWorker)/_account_crashes(_CrashSweep) - _dispatch_once_locked (193) -> _run_reclaim_phase/_tick_spawn_budget/_lane_rows/_any_spawnable_review/_resolve_default_assignee - _default_spawn (164) -> _worker_argv/_open_worker_log; TERMINAL_* overrides looped - connect: _probe_for_corruption/_missing_or_empty shared by guard+repair_db, _backfill_legacy_inflight_runs, _column_names/_table_exists, unified rebuild INSERT - dispatch: _kill_fn/_poll_worker_exit/_sigkill shared by enforce_max_runtime and _terminate_reclaimed_worker; dead _log loggers dropped - SQL trace parity identical; 934/934 kanban tests green --- hermes_cli/kanban_db_connect.py | 391 ++++---- hermes_cli/kanban_db_dispatch.py | 1387 +++++++++++++---------------- hermes_cli/kanban_db_notify.py | 10 +- hermes_cli/kanban_db_workspace.py | 29 +- 4 files changed, 803 insertions(+), 1014 deletions(-) diff --git a/hermes_cli/kanban_db_connect.py b/hermes_cli/kanban_db_connect.py index f2cacb6d8f..44afe3fc95 100644 --- a/hermes_cli/kanban_db_connect.py +++ b/hermes_cli/kanban_db_connect.py @@ -9,7 +9,6 @@ from __future__ import annotations import contextlib import hashlib -import logging import random import re import secrets @@ -24,39 +23,26 @@ from pathlib import Path from typing import Any from typing import Optional -# Log-record parity with the origin module. -_log = logging.getLogger("hermes_cli.kanban_db") - # --------------------------------------------------------------------------- # Connection helpers # --------------------------------------------------------------------------- _INITIALIZED_PATHS: set[str] = set() - - _INIT_LOCK = threading.RLock() - - _SQLITE_HEADER = b"SQLite format 3\x00" - - DEFAULT_BUSY_TIMEOUT_MS = 120_000 - # Cap on ``.corrupt..bak`` quarantines per board: content-addressing # dedupes identical bytes, but mutating corruption mints a new fingerprint each # time (one user hit 124). Oldest-by-mtime beyond the cap are pruned after each # new backup. _CORRUPT_BACKUP_RETENTION = 10 - # Bounded init-lock acquire: a bare blocking flock let a wedged holder block the # dispatcher's next-tick connect forever. Poll non-blocking until the deadline, # then proceed without the lock (in-process _INIT_LOCK + idempotent init backstop). _INIT_LOCK_TIMEOUT_SECONDS = 10.0 - - _INIT_LOCK_POLL_SECONDS = 0.05 @@ -227,11 +213,7 @@ def _dispatch_tick_lock(db_path: Path): # ``journal_size_limit`` (set at connection init) on the writer's natural # post-checkpoint reset. Best-effort, keyed per resolved DB path. _WAL_CHECKPOINT_INTERVAL_SECONDS = 300.0 - - _LAST_WAL_CHECKPOINT: dict[str, float] = {} - - _WAL_CHECKPOINT_LOCK = threading.Lock() @@ -266,9 +248,7 @@ def _looks_like_tls_record_at(data: bytes, offset: int) -> bool: """Return True for a TLS record header at ``data[offset:]``.""" if len(data) < offset + 5: return False - content_type = data[offset] - major = data[offset + 1] - minor = data[offset + 2] + content_type, major, minor = data[offset], data[offset + 1], data[offset + 2] length = int.from_bytes(data[offset + 3:offset + 5], "big") return ( content_type in {0x14, 0x15, 0x16, 0x17} @@ -285,22 +265,17 @@ def _validate_sqlite_header(path: Path) -> None: collapsed into a generic PRAGMA error and the gateway's corrupt-board handling can identify the board by fingerprint.""" try: - stat = path.stat() - except FileNotFoundError: - return + if path.stat().st_size == 0: + return except OSError: return - if stat.st_size == 0: - return # Byte-level probe: must run BEFORE any connection to this path exists # (read_header_bytes_preopen refuses once one is live, because the close() # would cancel this process's POSIX locks). from hermes_cli.sqlite_safe_read import read_header_bytes_preopen head = read_header_bytes_preopen(path, length=64) - if head is None: - return - if head.startswith(_SQLITE_HEADER): + if head is None or head.startswith(_SQLITE_HEADER): return signature = "" if head.startswith(b"SQLit") and _looks_like_tls_record_at(head, 5): @@ -322,16 +297,17 @@ class KanbanDbCorruptError(RuntimeError): self.db_path = db_path self.backup_path = backup_path self.reason = reason - backup_str = str(backup_path) if backup_path is not None else "" super().__init__( f"Refusing to open corrupt kanban DB at {db_path}: {reason}. " - f"Original preserved; backup at {backup_str}." + f"Original preserved; backup at {_backup_label(backup_path)}." ) -def _prune_corrupt_backups( - parent: Path, base_name: str, keep: Optional[Path] = None, -) -> None: +def _backup_label(backup_path: Optional[Path]) -> str: + return str(backup_path) if backup_path is not None else "" + + +def _prune_corrupt_backups(parent: Path, base_name: str, keep: Optional[Path] = None) -> None: """Keep only the ``_CORRUPT_BACKUP_RETENTION`` newest (by mtime) ``.corrupt..bak`` files plus their ``-wal``/``-shm`` copies. ``keep`` (the just-created backup) is never pruned regardless of mtime — @@ -346,8 +322,7 @@ def _prune_corrupt_backups( ] except OSError: return - budget = _kb._CORRUPT_BACKUP_RETENTION - (1 if keep is not None else 0) - budget = max(budget, 0) + budget = max(_kb._CORRUPT_BACKUP_RETENTION - (1 if keep is not None else 0), 0) if len(backups) <= budget: return @@ -359,11 +334,7 @@ def _prune_corrupt_backups( backups.sort(key=_mtime, reverse=True) for stale in backups[budget:]: - for victim in ( - stale, - stale.with_name(stale.name + "-wal"), - stale.with_name(stale.name + "-shm"), - ): + for victim in (stale, stale.with_name(stale.name + "-wal"), stale.with_name(stale.name + "-shm")): with contextlib.suppress(OSError): victim.unlink(missing_ok=True) @@ -402,8 +373,7 @@ def _backup_corrupt_db(path: Path) -> Optional[Path]: digest.update(chunk) except OSError: return None - token = digest.hexdigest()[:16] - candidate = parent / f"{base_name}.corrupt.{token}.bak" + candidate = parent / f"{base_name}.corrupt.{digest.hexdigest()[:16]}.bak" # Defensive: candidate must still be inside parent after construction. if candidate.parent != parent: return None @@ -417,10 +387,11 @@ def _backup_corrupt_db(path: Path) -> Optional[Path]: _prune_corrupt_backups(parent, base_name, keep=candidate) for suffix in ("-wal", "-shm"): sidecar = parent / (base_name + suffix) - if sidecar.parent != parent or not sidecar.exists(): - continue sidecar_backup = parent / (candidate.name + suffix) - if sidecar_backup.parent != parent or sidecar_backup.exists(): + if ( + sidecar.parent != parent or not sidecar.exists() + or sidecar_backup.parent != parent or sidecar_backup.exists() + ): continue with contextlib.suppress(OSError): shutil.copy2(sidecar, sidecar_backup) @@ -464,7 +435,6 @@ def _repairable_index_names(messages: list[str]) -> Optional[list[str]]: (caller fails closed; also ``None`` for no messages). First-appearance order is preserved so the REINDEX pass is deterministic.""" names: list[str] = [] - saw_any = False for raw in messages: message = (raw or "").strip() if not message: @@ -475,18 +445,13 @@ def _repairable_index_names(messages: list[str]) -> Optional[list[str]]: break else: return None - saw_any = True name = match.group("index").strip() if name and name not in names: names.append(name) - if not saw_any or not names: - return None - return names + return names or None -def _attempt_index_reindex_repair( - path: Path, index_names: list[str], -) -> tuple[bool, list[str]]: +def _attempt_index_reindex_repair(path: Path, index_names: list[str]) -> tuple[bool, list[str]]: """REINDEX the named indexes (per-index first; bare ``REINDEX`` fallback if a parsed name is an internal/auto index), then re-run integrity_check. Returns ``(clean, post_repair_messages)``; never raises. Callers must hold @@ -511,6 +476,30 @@ def _attempt_index_reindex_repair( return _integrity_messages_ok(messages), messages +def _missing_or_empty(resolved: Path) -> bool: + """True for a missing / zero-byte / unstat-able DB file — nothing to probe.""" + try: + return not resolved.exists() or resolved.stat().st_size == 0 + except OSError: + return True + + +def _probe_for_corruption(resolved: Path) -> tuple[Optional[list[str]], Optional[str]]: + """``(messages, reason)`` from an integrity probe; ``reason`` is ``None`` + when healthy and ``messages`` is ``None`` when sqlite refused to open the + file at all. ``OperationalError`` (lock/busy) is NOT corruption and + propagates raw so a locked healthy DB is never quarantined.""" + try: + messages = _probe_integrity(resolved) + except sqlite3.OperationalError: + raise + except sqlite3.DatabaseError as exc: + return None, f"sqlite refused to open file: {exc}" + if _integrity_messages_ok(messages): + return messages, None + return messages, f"integrity_check returned {messages[0] if messages else ''!r}" + + def _guard_existing_db_is_healthy(path: Path) -> None: """Run ``PRAGMA integrity_check`` on an existing non-empty DB file. @@ -528,47 +517,27 @@ def _guard_existing_db_is_healthy(path: Path) -> None: resolved = path.resolve() except OSError: return - try: - if not resolved.exists() or resolved.stat().st_size == 0: - return - except OSError: + if _missing_or_empty(resolved) or str(resolved) in _INITIALIZED_PATHS: return - if str(resolved) in _INITIALIZED_PATHS: - return - reason: Optional[str] = None - messages: list[str] = [] - try: - messages = _probe_integrity(resolved) - if not _integrity_messages_ok(messages): - reason = ( - f"integrity_check returned " - f"{messages[0] if messages else ''!r}" - ) - except sqlite3.OperationalError: - # Lock contention / busy — not corruption. Propagate. - raise - except sqlite3.DatabaseError as exc: - reason = f"sqlite refused to open file: {exc}" + messages, reason = _probe_for_corruption(resolved) if reason is None: return # Quarantine FIRST — both the repair and fail-closed paths preserve the # pre-touch bytes before anything mutates the file. backup = _backup_corrupt_db(resolved) - index_names = _repairable_index_names(messages) + index_names = _repairable_index_names(messages or []) if index_names: _kb._log.warning( "kanban DB %s failed integrity_check with index-only errors " "(%s); pre-repair backup at %s — attempting REINDEX auto-repair.", - resolved, ", ".join(index_names), - backup if backup is not None else "", + resolved, ", ".join(index_names), _backup_label(backup), ) repaired, post = _attempt_index_reindex_repair(resolved, index_names) if repaired: _kb._log.warning( "kanban DB %s auto-repaired via REINDEX (%s); " "integrity_check now clean. Pre-repair copy kept at %s.", - resolved, ", ".join(index_names), - backup if backup is not None else "", + resolved, ", ".join(index_names), _backup_label(backup), ) return reason = ( @@ -593,11 +562,7 @@ class RepairResult: reindexed: list[str] = field(default_factory=list) -def repair_db( - db_path: Optional[Path] = None, - *, - board: Optional[str] = None, -) -> RepairResult: +def repair_db(db_path: Optional[Path] = None, *, board: Optional[str] = None) -> RepairResult: """Probe a kanban DB and apply the narrow index-REINDEX repair if needed. Same policy as :func:`_guard_existing_db_is_healthy` (quarantine BEFORE any mutation; REINDEX under the init flock; anything non-index stays @@ -610,41 +575,26 @@ def repair_db( resolved = path.resolve() except OSError: resolved = path - try: - if not resolved.exists() or resolved.stat().st_size == 0: - return RepairResult(status="missing", db_path=resolved) - except OSError: + if _missing_or_empty(resolved): return RepairResult(status="missing", db_path=resolved) with _kb._cross_process_init_lock(resolved): - messages: list[str] = [] - try: - messages = _probe_integrity(resolved) - except sqlite3.OperationalError: - # Locked/busy — not corruption; let the caller report it raw. - raise - except sqlite3.DatabaseError as exc: + messages, reason = _probe_for_corruption(resolved) + if messages is None: # Same quarantine the connect-time guard takes when sqlite # refuses to open the file at all. return RepairResult( - status="corrupt", - db_path=resolved, - messages=[f"sqlite refused to open file: {exc}"], + status="corrupt", db_path=resolved, messages=[str(reason)], backup_path=_backup_corrupt_db(resolved), ) - if _integrity_messages_ok(messages): + if reason is None: return RepairResult(status="ok", db_path=resolved, messages=messages) # Quarantine FIRST — identical policy to the connect-time guard. backup = _backup_corrupt_db(resolved) index_names = _repairable_index_names(messages) if not index_names: - return RepairResult( - status="corrupt", - db_path=resolved, - messages=messages, - backup_path=backup, - ) + return RepairResult(status="corrupt", db_path=resolved, messages=messages, backup_path=backup) repaired, post = _attempt_index_reindex_repair(resolved, index_names) # The file changed on disk; force the next connect() in this process # to re-probe instead of trusting the stale healthy-path cache. @@ -710,11 +660,7 @@ def _open_configured(path: Path, under_lock) -> tuple[sqlite3.Connection, Any]: return conn, out -def connect( - db_path: Optional[Path] = None, - *, - board: Optional[str] = None, -) -> sqlite3.Connection: +def connect(db_path: Optional[Path] = None, *, board: Optional[str] = None) -> sqlite3.Connection: """Open (and initialize if needed) the kanban DB. WAL is (re)enabled on every connection so a re-created file stays robust; the first connection per path auto-runs :func:`init_db`, later ones skip via @@ -771,11 +717,7 @@ def connect( @contextlib.contextmanager -def connect_closing( - db_path: Optional[Path] = None, - *, - board: Optional[str] = None, -): +def connect_closing(db_path: Optional[Path] = None, *, board: Optional[str] = None): """Open a kanban DB connection and guarantee it is closed on exit. Use instead of ``with kb.connect() as conn:`` — sqlite3's context manager only commits/rolls back, it does NOT close the fd, so long-lived processes @@ -788,21 +730,16 @@ def connect_closing( conn.close() -def init_db( - db_path: Optional[Path] = None, - *, - board: Optional[str] = None, -) -> Path: +def init_db(db_path: Optional[Path] = None, *, board: Optional[str] = None) -> Path: """Create the schema if it doesn't exist; return the path used. Unlike :func:`connect`'s cached first-time auto-init, this always re-runs the migration pass — callers that know the on-disk schema may have drifted (tests writing legacy event kinds, external upgrades) use it to force it.""" path = db_path if db_path is not None else _kb.kanban_db_path(board=board) path.parent.mkdir(parents=True, exist_ok=True) - resolved = str(path.resolve()) # Clear the cache entry so connect() re-runs schema + migrations. with _INIT_LOCK: - _INITIALIZED_PATHS.discard(resolved) + _INITIALIZED_PATHS.discard(str(path.resolve())) with contextlib.closing(_kb.connect(path)): pass return path @@ -818,7 +755,6 @@ _EARLY_TASK_COLUMNS = ( ("idempotency_key", "idempotency_key TEXT"), ) - # (new column, ddl, legacy source column, copy statement) — see the # RENAME-avoidance note in ``_migrate_add_optional_columns``. _RENAMED_TASK_COLUMNS = ( @@ -833,7 +769,6 @@ _RENAMED_TASK_COLUMNS = ( ), ) - # NULL / 0 defaults below reproduce the behaviour existing rows had before the # column existed. _LATER_TASK_COLUMNS = ( @@ -859,10 +794,31 @@ _LATER_TASK_COLUMNS = ( ("block_recurrences", "block_recurrences INTEGER NOT NULL DEFAULT 0"), ) +_NOTIFY_SUB_COLUMNS = ( + ("notifier_profile", "notifier_profile TEXT"), + ("delivery_mode", "delivery_mode TEXT NOT NULL DEFAULT 'notify'"), + ("chat_type", "chat_type TEXT"), + # Platform-specific stable alt ID (Signal UUID, Feishu union_id, ...) + # so an active-wake replay reconstructs the SAME ``build_session_key`` + # (which prefers ``user_id_alt``). NULL is inert. + ("user_id_alt", "user_id_alt TEXT"), + ("delivery_metadata", "delivery_metadata TEXT"), +) + + +def _column_names(conn: sqlite3.Connection, table: str) -> set[str]: + return {row["name"] for row in conn.execute(f"PRAGMA table_info({table})")} + + +def _table_exists(conn: sqlite3.Connection, table: str) -> bool: + return conn.execute( + f"SELECT name FROM sqlite_master WHERE type='table' AND name='{table}'" + ).fetchone() is not None + def _migrate_add_optional_columns(conn: sqlite3.Connection) -> None: """Add columns introduced after v1 to legacy DBs (called via ``init_db``).""" - cols = {row["name"] for row in conn.execute("PRAGMA table_info(tasks)")} + cols = _column_names(conn, "tasks") for name, ddl in _EARLY_TASK_COLUMNS: if name not in cols: _add_column_if_missing(conn, "tasks", name, ddl) @@ -870,7 +826,7 @@ def _migrate_add_optional_columns(conn: sqlite3.Connection) -> None: # Re-snapshot: DBs partially migrated by older releases may already carry # later columns (e.g. ``consecutive_failures``), keeping the legacy-column # migration idempotent. - cols = {row["name"] for row in conn.execute("PRAGMA table_info(tasks)")} + cols = _column_names(conn, "tasks") # Legacy renames via ADD-then-copy rather than ``RENAME COLUMN``: very old # DBs may lack the legacy column entirely (RENAME raises "no such column"), @@ -894,42 +850,20 @@ def _migrate_add_optional_columns(conn: sqlite3.Connection) -> None: # init on legacy boards before the ALTER TABLE pass runs. ``IF NOT EXISTS`` # keeps re-running here cheap and correct on fresh DBs. conn.execute("CREATE INDEX IF NOT EXISTS idx_tasks_tenant ON tasks(tenant)") - conn.execute( - "CREATE INDEX IF NOT EXISTS idx_tasks_idempotency ON tasks(idempotency_key)" - ) - conn.execute( - "CREATE INDEX IF NOT EXISTS idx_tasks_session_id ON tasks(session_id)" - ) + conn.execute("CREATE INDEX IF NOT EXISTS idx_tasks_idempotency ON tasks(idempotency_key)") + conn.execute("CREATE INDEX IF NOT EXISTS idx_tasks_session_id ON tasks(session_id)") # task_events.run_id back-fills as NULL for historical events (they predate # runs and can't be attributed). - ev_cols = {row["name"] for row in conn.execute("PRAGMA table_info(task_events)")} - if "run_id" not in ev_cols: + if "run_id" not in _column_names(conn, "task_events"): _add_column_if_missing(conn, "task_events", "run_id", "run_id INTEGER") # Same ordering rule as the ``tasks`` indexes above: index after column. - conn.execute( - "CREATE INDEX IF NOT EXISTS idx_events_run " - "ON task_events(run_id, id)" - ) + conn.execute("CREATE INDEX IF NOT EXISTS idx_events_run ON task_events(run_id, id)") - notify_table_exists = conn.execute( - "SELECT name FROM sqlite_master WHERE type='table' AND name='kanban_notify_subs'" - ).fetchone() is not None - if notify_table_exists: - notify_cols = { - row["name"] for row in conn.execute("PRAGMA table_info(kanban_notify_subs)") - } - for name, ddl in ( - ("notifier_profile", "notifier_profile TEXT"), - ("delivery_mode", "delivery_mode TEXT NOT NULL DEFAULT 'notify'"), - ("chat_type", "chat_type TEXT"), - # Platform-specific stable alt ID (Signal UUID, Feishu union_id, ...) - # so an active-wake replay reconstructs the SAME ``build_session_key`` - # (which prefers ``user_id_alt``). NULL is inert. - ("user_id_alt", "user_id_alt TEXT"), - ("delivery_metadata", "delivery_metadata TEXT"), - ): + if _table_exists(conn, "kanban_notify_subs"): + notify_cols = _column_names(conn, "kanban_notify_subs") + for name, ddl in _NOTIFY_SUB_COLUMNS: if name in notify_cols: continue _add_column_if_missing(conn, "kanban_notify_subs", name, ddl) @@ -944,57 +878,8 @@ def _migrate_add_optional_columns(conn: sqlite3.Connection) -> None: "WHERE platform != 'tui'" ) - # One-shot backfill: tasks 'running' before runs existed carried - # claim_lock / claim_expires / worker_pid on the task row; synthesize a - # matching task_runs row so end-run / heartbeat have something to write. - # write_txn serializes against concurrent dispatchers, and the per-row - # UPDATE uses ``current_run_id IS NULL`` as a CAS guard so a racing claim - # can't produce an orphaned row. - runs_exist = conn.execute( - "SELECT name FROM sqlite_master WHERE type='table' AND name='task_runs'" - ).fetchone() is not None - if runs_exist: - with write_txn(conn): - inflight = conn.execute( - "SELECT id, assignee, claim_lock, claim_expires, worker_pid, " - " max_runtime_seconds, last_heartbeat_at, started_at " - "FROM tasks " - "WHERE status = 'running' AND current_run_id IS NULL" - ).fetchall() - for row in inflight: - started = row["started_at"] or int(time.time()) - cur = conn.execute( - """ - INSERT INTO task_runs ( - task_id, profile, status, - claim_lock, claim_expires, worker_pid, - max_runtime_seconds, last_heartbeat_at, - started_at - ) VALUES (?, ?, 'running', ?, ?, ?, ?, ?, ?) - """, - ( - row["id"], row["assignee"], row["claim_lock"], - row["claim_expires"], row["worker_pid"], - row["max_runtime_seconds"], row["last_heartbeat_at"], - started, - ), - ) - # CAS: only install the pointer if nothing claimed the task - # since our SELECT (belt-and-suspenders under write_txn). On - # failure mark the orphan run row reclaimed so it doesn't - # look in-flight. - upd = conn.execute( - "UPDATE tasks SET current_run_id = ? " - "WHERE id = ? AND current_run_id IS NULL", - (cur.lastrowid, row["id"]), - ) - if upd.rowcount != 1: - conn.execute( - "UPDATE task_runs SET status = 'reclaimed', " - " outcome = 'reclaimed', ended_at = ? " - "WHERE id = ?", - (int(time.time()), cur.lastrowid), - ) + if _table_exists(conn, "task_runs"): + _backfill_legacy_inflight_runs(conn) # One-shot event-kind rename: old names still worked but were awkward on # the wire. Fires once per DB — after the UPDATE no rows match. @@ -1008,6 +893,56 @@ def _migrate_add_optional_columns(conn: sqlite3.Connection) -> None: _rebuild_drifted_tables(conn) +def _backfill_legacy_inflight_runs(conn: sqlite3.Connection) -> None: + """One-shot backfill: tasks 'running' before runs existed carried + claim_lock / claim_expires / worker_pid on the task row; synthesize a + matching task_runs row so end-run / heartbeat have something to write. + write_txn serializes against concurrent dispatchers, and the per-row + UPDATE uses ``current_run_id IS NULL`` as a CAS guard so a racing claim + can't produce an orphaned row.""" + with write_txn(conn): + inflight = conn.execute( + "SELECT id, assignee, claim_lock, claim_expires, worker_pid, " + " max_runtime_seconds, last_heartbeat_at, started_at " + "FROM tasks " + "WHERE status = 'running' AND current_run_id IS NULL" + ).fetchall() + for row in inflight: + started = row["started_at"] or int(time.time()) + cur = conn.execute( + """ + INSERT INTO task_runs ( + task_id, profile, status, + claim_lock, claim_expires, worker_pid, + max_runtime_seconds, last_heartbeat_at, + started_at + ) VALUES (?, ?, 'running', ?, ?, ?, ?, ?, ?) + """, + ( + row["id"], row["assignee"], row["claim_lock"], + row["claim_expires"], row["worker_pid"], + row["max_runtime_seconds"], row["last_heartbeat_at"], + started, + ), + ) + # CAS: only install the pointer if nothing claimed the task + # since our SELECT (belt-and-suspenders under write_txn). On + # failure mark the orphan run row reclaimed so it doesn't + # look in-flight. + upd = conn.execute( + "UPDATE tasks SET current_run_id = ? " + "WHERE id = ? AND current_run_id IS NULL", + (cur.lastrowid, row["id"]), + ) + if upd.rowcount != 1: + conn.execute( + "UPDATE task_runs SET status = 'reclaimed', " + " outcome = 'reclaimed', ended_at = ? " + "WHERE id = ?", + (int(time.time()), cur.lastrowid), + ) + + # Legacy DBs defined these tables with a ``TEXT PRIMARY KEY`` id (or a nullable # ``TEXT last_event_id`` for ``kanban_notify_subs``); the current schema uses # ``INTEGER PRIMARY KEY AUTOINCREMENT`` / ``INTEGER NOT NULL DEFAULT 0``. @@ -1105,24 +1040,19 @@ def _rebuild_drifted_tables(conn: sqlite3.Connection) -> None: _kb._log.info("kanban migration: rebuilding %s to match current schema", table) conn.execute(f"ALTER TABLE {table} RENAME TO {table}_legacy") conn.execute(create_sql) - new_cols = {c["name"] for c in conn.execute(f"PRAGMA table_info({table})")} + new_cols = _column_names(conn, table) if table == "kanban_notify_subs": # Cast the legacy TEXT cursor to INTEGER; NULL / non-numeric → 0. - shared = [c for c in old_cols if c in new_cols and c != "last_event_id"] - cols_csv = ", ".join(shared) - conn.execute( - f"INSERT INTO {table} ({cols_csv}, last_event_id) " - f"SELECT {cols_csv}, COALESCE(CAST(last_event_id AS INTEGER), 0) " - f"FROM {table}_legacy" - ) + drop, extra_cols = "last_event_id", ", last_event_id" + extra_select = ", COALESCE(CAST(last_event_id AS INTEGER), 0)" else: # Drop the legacy TEXT id; AUTOINCREMENT reassigns it. - shared = [c for c in old_cols if c in new_cols and c != "id"] - cols_csv = ", ".join(shared) - conn.execute( - f"INSERT INTO {table} ({cols_csv}) " - f"SELECT {cols_csv} FROM {table}_legacy" - ) + drop, extra_cols, extra_select = "id", "", "" + cols_csv = ", ".join(c for c in old_cols if c in new_cols and c != drop) + conn.execute( + f"INSERT INTO {table} ({cols_csv}{extra_cols}) " + f"SELECT {cols_csv}{extra_select} FROM {table}_legacy" + ) conn.execute(f"DROP TABLE {table}_legacy") for index_sql in index_sqls: conn.execute(index_sql) @@ -1155,8 +1085,7 @@ def _check_file_length_invariant(conn: sqlite3.Connection) -> None: if journal_mode == "wal": return - ok = file_length_matches_header(conn) - if ok is False: + if file_length_matches_header(conn) is False: raise sqlite3.DatabaseError( "torn-extend detected: the database file is shorter than its " "header page count claims" @@ -1172,11 +1101,7 @@ def _check_file_length_invariant(conn: sqlite3.Connection) -> None: # 15): the 120s busy_timeout absorbs most waits; this is the backstop for the # tail where SQLite returns BUSY immediately. _BUSY_MAX_RETRIES = 5 - - _BUSY_RETRY_MIN_S = 0.020 # 20ms - - _BUSY_RETRY_MAX_S = 0.150 # 150ms @@ -1227,11 +1152,9 @@ def write_txn(conn: sqlite3.Connection, *, allow_nested: bool = False): try: yield conn except Exception: - try: + with contextlib.suppress(sqlite3.OperationalError): conn.execute(f"ROLLBACK TO {savepoint}") conn.execute(f"RELEASE {savepoint}") - except sqlite3.OperationalError: - pass raise else: conn.execute(f"RELEASE {savepoint}") @@ -1241,12 +1164,10 @@ def write_txn(conn: sqlite3.Connection, *, allow_nested: bool = False): try: yield conn except Exception: - try: + # SQLite may already have auto-rolled-back (EIO, contention, corruption); + # don't let this secondary failure shadow the real one. + with contextlib.suppress(sqlite3.OperationalError): conn.execute("ROLLBACK") - except sqlite3.OperationalError: - # SQLite already auto-rolled-back (EIO, contention, corruption); - # don't let this secondary failure shadow the real one. - pass raise else: try: @@ -1261,6 +1182,6 @@ def write_txn(conn: sqlite3.Connection, *, allow_nested: bool = False): _kb._check_file_length_invariant(conn) -# Late-bound origin namespace (see module docstring). Imported LAST so this +# Late-bound origin namespace (see module docstring); imported LAST so this # module is fully populated before ``kanban_db`` re-exports from it. from hermes_cli import kanban_db as _kb # noqa: E402 diff --git a/hermes_cli/kanban_db_dispatch.py b/hermes_cli/kanban_db_dispatch.py index 0381dfd2c0..2f5d7b75c1 100644 --- a/hermes_cli/kanban_db_dispatch.py +++ b/hermes_cli/kanban_db_dispatch.py @@ -8,9 +8,9 @@ origin-resident helpers are reached late-bound via ``_kb`` so monkeypatching from __future__ import annotations import contextlib -import logging import os import re +import signal import sqlite3 import subprocess import sys @@ -19,6 +19,7 @@ from dataclasses import dataclass from dataclasses import field from pathlib import Path from typing import Any +from typing import Callable from typing import Mapping from typing import Optional from typing import TYPE_CHECKING @@ -26,31 +27,21 @@ from typing import TYPE_CHECKING if TYPE_CHECKING: from hermes_cli.kanban_db import Task -# Log-record parity with the origin module. -_log = logging.getLogger("hermes_cli.kanban_db") - # After this many consecutive non-success attempts on a task/profile the # dispatcher parks the task in ``blocked`` with a reason — prevents retry storms. DEFAULT_FAILURE_LIMIT = 2 - - # Legacy alias — callers / tests still reference the old name. DEFAULT_SPAWN_FAILURE_LIMIT = DEFAULT_FAILURE_LIMIT - # Worker log files larger than this at spawn time are rotated. DEFAULT_LOG_ROTATE_BYTES = 2 * 1024 * 1024 # 2 MiB - - DEFAULT_LOG_BACKUP_COUNT = 1 - # Keep a little wall-clock budget for the worker to observe a terminal timeout # and call kanban_block/kanban_complete before max_runtime_seconds kills it. KANBAN_TERMINAL_TIMEOUT_GRACE_SECONDS = 30 - # --------------------------------------------------------------------------- # Respawn guard constants # --------------------------------------------------------------------------- @@ -65,23 +56,18 @@ _RESPAWN_BLOCKER_RE = re.compile( re.IGNORECASE, ) - # Within this window a completed run counts as "recent proof"; don't re-spawn. _RESPAWN_GUARD_SUCCESS_WINDOW = 3600 # 1 hour - # Cooldown after a rate-limited (quota-wall) requeue before re-spawning. Without # it the task would re-spawn on the very next tick and bounce off the same quota # wall, burning a worker slot every tick for hours. Overridable via # ``HERMES_KANBAN_RATE_LIMIT_COOLDOWN_SECONDS``. DEFAULT_RATE_LIMIT_COOLDOWN_SECONDS = 300 # 5 minutes - # Within this window a GitHub PR URL in a comment blocks re-spawn. _RESPAWN_GUARD_PR_WINDOW = 86400 # 24 hours - -# Pattern matching a GitHub PR URL in task comments. _RESPAWN_GUARD_PR_URL_RE = re.compile( r"https?://github\.com/[^/\s]+/[^/\s]+/pull/\d+", re.IGNORECASE, @@ -95,35 +81,26 @@ class DispatchResult: reclaimed: int = 0 promoted: int = 0 reconciled_orphans: list[str] = field(default_factory=list) - """Task ids requeued by :func:`reconcile_orphaned_running` this tick — - ``running`` cards whose claim bookkeeping was broken (no valid claim, - dead/gone worker). See the reconciliation pass for details.""" + """``running`` cards requeued by :func:`reconcile_orphaned_running` (broken + claim bookkeeping, dead/gone worker).""" spawned: list[tuple[str, str, str]] = field(default_factory=list) - """List of ``(task_id, assignee, workspace_path)`` triples.""" + """``(task_id, assignee, workspace_path)`` triples.""" skipped_unassigned: list[str] = field(default_factory=list) - """Ready task ids skipped because they have no assignee at all. - Operator-actionable — usually a misfiled task waiting for routing.""" + """Ready task ids with no assignee at all — operator-actionable (usually a + misfiled task waiting for routing).""" auto_assigned_default: list[str] = field(default_factory=list) - """Task ids that were unassigned in the DB and had - ``kanban.default_assignee`` applied this tick before spawning (#27145). - Surfaces the auto-assignment to telemetry / CLI / dashboard so the - operator can see when the dispatcher is acting on the fallback rule - rather than on explicit per-task assignments.""" + """Unassigned task ids that had ``kanban.default_assignee`` applied this + tick before spawning, so telemetry/CLI/dashboard can show the dispatcher + acting on the fallback rule rather than explicit assignments.""" skipped_nonspawnable: list[str] = field(default_factory=list) - """Ready task ids skipped because their assignee names a control-plane - lane (a Claude Code terminal like ``orion-cc``) rather than a Hermes - profile. Expected steady-state on multi-lane setups; NOT an - operator-actionable failure. Tracked separately so health telemetry - can distinguish "real stuck" (nothing spawned but spawnable work - available) from "correctly idle" (nothing spawnable in the queue).""" + """Ready task ids whose assignee names a control-plane lane (e.g. a Claude + Code terminal like ``orion-cc``), not a Hermes profile. Expected steady-state + on multi-lane setups, NOT operator-actionable; tracked apart so health + telemetry can tell "stuck" from "correctly idle".""" skipped_per_profile_capped: list[tuple[str, str, int]] = field(default_factory=list) - """Tasks deferred this tick because their assignee is already at - ``kanban.max_in_progress_per_profile`` (#21582). Each entry is - ``(task_id, assignee, current_running_count)``. NOT an - operator-actionable failure — the task will be picked up on a - subsequent tick when the assignee has capacity. Separate bucket so - telemetry / dashboards can show "this profile is busy" vs - "task is genuinely stuck".""" + """``(task_id, assignee, current_running_count)`` deferred because the + assignee is at ``kanban.max_in_progress_per_profile``. Picked up on a later + tick; separate bucket so dashboards show "profile busy" vs "stuck".""" crashed: list[str] = field(default_factory=list) """Task ids reclaimed because their worker PID disappeared.""" auto_blocked: list[str] = field(default_factory=list) @@ -131,32 +108,22 @@ class DispatchResult: timed_out: list[str] = field(default_factory=list) """Task ids whose workers exceeded ``max_runtime_seconds``.""" stale: list[str] = field(default_factory=list) - """Task ids reclaimed because no progress (heartbeat) was seen - within ``dispatch_stale_timeout_seconds``.""" + """Task ids reclaimed for no heartbeat within ``dispatch_stale_timeout_seconds``.""" respawn_guarded: list[tuple[str, str]] = field(default_factory=list) - """Tasks skipped by the respawn guard, as ``(task_id, reason)`` pairs. - - Reasons: ``"blocker_auth"`` (quota/auth error — also auto-blocked), - ``"recent_success"`` (completed run within guard window), - ``"active_pr"`` (GitHub PR URL in a recent comment).""" + """``(task_id, reason)`` skipped by the respawn guard: ``"blocker_auth"`` + (quota/auth error — also auto-blocked), ``"recent_success"`` (completed run + within guard window), ``"active_pr"`` (GitHub PR URL in a recent comment).""" rate_limited: list[str] = field(default_factory=list) """Task ids whose workers bailed on a provider rate-limit / quota wall - (EX_TEMPFAIL sentinel exit) and were released back to ``ready`` WITHOUT - counting a failure. These never trip the circuit breaker — a long quota - window just makes the task bounce cheaply until the window clears.""" + (EX_TEMPFAIL sentinel exit) and were released to ``ready`` WITHOUT counting + a failure — a long quota window must never trip the circuit breaker.""" skipped_locked: bool = False - """True when this tick was skipped because another process already held - the board's dispatch lock (issue #35240). A losing dispatcher does no - DB writes this tick — the lock holder is making progress on the same - board. This is the steady-state signal that a single-writer guard is - actively preventing two dispatchers from racing on ``kanban.db``.""" + """True when another process held the board's dispatch lock: this tick did + no DB writes; the lock holder is making progress on the same board.""" memory_pressure: Optional[str] = None - """System memory pressure observed at spawn time when the memory guard - restricted this tick (OOF-30/OOF-77): ``"critical"`` — no new workers - were spawned this tick; ``"elevated"`` — at most one new worker was - spawned. ``None`` when memory was fine/unknown and the guard imposed - no restriction. Reclaim/promotion bookkeeping still ran either way; - deferred tasks stay queued for the next tick.""" + """Memory pressure that restricted this tick: ``"critical"`` (no new + workers), ``"elevated"`` (at most one), ``None`` (no restriction). + Reclaim/promotion bookkeeping still ran; deferred tasks stay queued.""" # Bounded registry of recently-reaped worker exits, filled by the reap loop in @@ -165,11 +132,7 @@ class DispatchResult: # both WIFEXITED/WEXITSTATUS and WIFSIGNALED can be consulted. Trimmed by age # plus a total size cap. _RECENT_WORKER_EXIT_TTL_SECONDS = 600 - - _RECENT_WORKER_EXITS_MAX = 4096 - - _recent_worker_exits: "dict[int, tuple[int, float]]" = {} @@ -289,6 +252,32 @@ def _pid_alive(pid: Optional[int]) -> bool: return True +def _kill_fn(signal_fn) -> Optional[Callable[[int, int], None]]: + """``signal_fn`` test hook, else ``os.kill`` when the platform has one.""" + if signal_fn is not None: + return signal_fn + return os.kill if hasattr(os, "kill") else None + + +def _poll_worker_exit(pid: int) -> bool: + """Poll ~5 s (10 x 0.5 s) for ``pid`` to die; True once it is gone.""" + for _ in range(10): + if not _kb._pid_alive(pid): + return True + time.sleep(0.5) + return False + + +def _sigkill(kill, pid: int) -> bool: + """Best-effort SIGKILL; True when the signal was delivered.""" + try: + # signal.SIGKILL doesn't exist on Windows; SIGTERM maps to TerminateProcess. + kill(int(pid), getattr(signal, "SIGKILL", signal.SIGTERM)) + return True + except (ProcessLookupError, OSError): + return False + + def _terminate_reclaimed_worker( pid: Optional[int], claim_lock: Optional[str], @@ -296,8 +285,6 @@ def _terminate_reclaimed_worker( signal_fn=None, ) -> dict[str, Any]: """Best-effort host-local worker termination for reclaim paths.""" - import signal - info: dict[str, Any] = { "prev_pid": int(pid) if pid else None, "host_local": False, @@ -307,15 +294,11 @@ def _terminate_reclaimed_worker( } if not pid or pid <= 0 or not claim_lock: return info - - host_prefix = _kb._host_prefix() - if not str(claim_lock).startswith(host_prefix): + if not str(claim_lock).startswith(_kb._host_prefix()): return info info["host_local"] = True - kill = signal_fn if signal_fn is not None else ( - os.kill if hasattr(os, "kill") else None - ) + kill = _kill_fn(signal_fn) if kill is None: return info @@ -330,21 +313,13 @@ def _terminate_reclaimed_worker( except OSError: return info - for _ in range(10): - if not _kb._pid_alive(pid): - info["terminated"] = True - return info - time.sleep(0.5) - + if _poll_worker_exit(pid): + info["terminated"] = True + return info if _kb._pid_alive(pid): - try: - # signal.SIGKILL doesn't exist on Windows; SIGTERM maps to TerminateProcess. - _sigkill = getattr(signal, "SIGKILL", signal.SIGTERM) - kill(int(pid), _sigkill) - info["sigkill"] = True - except (ProcessLookupError, OSError): + if not _sigkill(kill, pid): return info - + info["sigkill"] = True info["terminated"] = not _kb._pid_alive(pid) return info @@ -391,15 +366,8 @@ def _defer_reclaim_for_live_worker( return run_id = _kb._current_run_id(conn, task_id) if run_id is not None: - conn.execute( - "UPDATE task_runs SET claim_expires = ? WHERE id = ?", - (grace, run_id), - ) - payload = { - "reason": reason, - "claim_lock": claim_lock, - "claim_expires_now": grace, - } + conn.execute("UPDATE task_runs SET claim_expires = ? WHERE id = ?", (grace, run_id)) + payload = {"reason": reason, "claim_lock": claim_lock, "claim_expires_now": grace} payload.update(termination) _kb._append_event(conn, task_id, "reclaim_deferred", payload, run_id=run_id) @@ -439,10 +407,7 @@ def heartbeat_worker( else _kb._current_run_id(conn, task_id) ) if run_id is not None: - conn.execute( - "UPDATE task_runs SET last_heartbeat_at = ? WHERE id = ?", - (now, run_id), - ) + conn.execute("UPDATE task_runs SET last_heartbeat_at = ? WHERE id = ?", (now, run_id)) _kb._append_event( conn, task_id, "heartbeat", {"note": note} if note else None, @@ -451,11 +416,7 @@ def heartbeat_worker( return True -def enforce_max_runtime( - conn: sqlite3.Connection, - *, - signal_fn=None, -) -> list[str]: +def enforce_max_runtime(conn: sqlite3.Connection, *, signal_fn=None) -> list[str]: """Terminate workers whose per-task ``max_runtime_seconds`` has elapsed. SIGTERM, short grace, then SIGKILL. Emits ``timed_out`` and restores the @@ -463,7 +424,6 @@ def enforce_max_runtime( unless the circuit breaker already gave up, leaving it blocked. Host-local only (same reasoning as ``detect_crashed_workers``). ``signal_fn`` is a test hook. """ - import signal timed_out: list[str] = [] now = int(time.time()) host_prefix = _kb._host_prefix() @@ -485,7 +445,8 @@ def enforce_max_runtime( # Runtime is per attempt: ``tasks.started_at`` records the FIRST start, # so retries must be measured from the active task_runs row. elapsed = now - int(row["active_started_at"]) - if elapsed < int(row["max_runtime_seconds"]): + limit = int(row["max_runtime_seconds"]) + if elapsed < limit: continue pid = int(row["worker_pid"]) @@ -493,26 +454,16 @@ def enforce_max_runtime( # SIGTERM then SIGKILL after 5 s grace; workers wanting a cleaner # shutdown install their own SIGTERM handler. killed = False - kill = signal_fn if signal_fn is not None else ( - os.kill if hasattr(os, "kill") else None - ) + kill = _kill_fn(signal_fn) if kill is not None: with contextlib.suppress(ProcessLookupError, OSError): kill(pid, signal.SIGTERM) # Short polling wait — no time.sleep on the write txn. - for _ in range(10): - if not _kb._pid_alive(pid): - break - time.sleep(0.5) + _poll_worker_exit(pid) if _kb._pid_alive(pid): - try: - # signal.SIGKILL doesn't exist on Windows. - _sigkill = getattr(signal, "SIGKILL", signal.SIGTERM) - kill(pid, _sigkill) - killed = True - except (ProcessLookupError, OSError): - pass + killed = _sigkill(kill, pid) + error = f"elapsed {int(elapsed)}s > limit {limit}s" with _kb.write_txn(conn): retry_status = _kb._retry_status_for_run(conn, tid) cur = conn.execute( @@ -527,19 +478,15 @@ def enforce_max_runtime( payload = { "pid": pid, "elapsed_seconds": int(elapsed), - "limit_seconds": int(row["max_runtime_seconds"]), + "limit_seconds": limit, "sigkill": killed, "retry_status": retry_status, } run_id = _kb._end_run( - conn, tid, - outcome="timed_out", status="timed_out", - error=f"elapsed {int(elapsed)}s > limit {int(row['max_runtime_seconds'])}s", - metadata=payload, - ) - _kb._append_event( - conn, tid, "timed_out", payload, run_id=run_id, + conn, tid, outcome="timed_out", status="timed_out", + error=error, metadata=payload, ) + _kb._append_event(conn, tid, "timed_out", payload, run_id=run_id) timed_out.append(tid) # Outside the write_txn above because ``_record_task_failure`` opens its # own. If the breaker trips this flips the task to ``blocked`` and emits @@ -547,15 +494,11 @@ def enforce_max_runtime( if cur.rowcount == 1: _kb._record_task_failure( conn, tid, - error=f"elapsed {int(elapsed)}s > limit {int(row['max_runtime_seconds'])}s", + error=error, outcome="timed_out", release_claim=False, end_run=False, - event_payload_extra={ - "pid": pid, - "sigkill": killed, - "retry_status": retry_status, - }, + event_payload_extra={"pid": pid, "sigkill": killed, "retry_status": retry_status}, ) return timed_out @@ -578,11 +521,15 @@ def detect_stale_running( than ``_STALE_HEARTBEAT_GAP_SECONDS`` or NULL. Task returns to its source phase, run closes ``outcome='stale'``, live host-local worker is terminated. ``stale_timeout_seconds=0`` disables the check. ``signal_fn`` is a test hook. + + Deliberately does NOT call ``_record_task_failure``: stale reclaim is + dispatcher-side detection of an absent heartbeat, not a worker failure; + counting it would let legitimately long-running tasks (>4h, no explicit + heartbeat) trip the breaker. The ``stale`` event is the audit surface. """ if stale_timeout_seconds <= 0: return [] - now = int(time.time()) reclaimed: list[str] = [] @@ -597,7 +544,6 @@ def detect_stale_running( for row in rows: if row["active_started_at"] is None: continue - elapsed = now - int(row["active_started_at"]) if elapsed < stale_timeout_seconds: continue @@ -611,9 +557,7 @@ def detect_stale_running( tid = row["id"] lock = row["claim_lock"] or "" - termination = _kb._terminate_reclaimed_worker( - pid, lock, signal_fn=signal_fn, - ) + termination = _kb._terminate_reclaimed_worker(pid, lock, signal_fn=signal_fn) # 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. @@ -657,22 +601,13 @@ def detect_stale_running( ) + f" after {int(elapsed)}s running", metadata=payload, ) - _kb._append_event( - conn, tid, "stale", payload, run_id=run_id, - ) + _kb._append_event(conn, tid, "stale", payload, run_id=run_id) reclaimed.append(tid) - # Intentionally NOT calling _record_task_failure: stale reclaim is - # dispatcher-side detection of an absent heartbeat, not a worker failure. - # Counting it would let legitimately long-running tasks (>4h, no explicit - # heartbeat) trip the breaker. The 'stale' event is the audit surface. - return reclaimed -def reconcile_orphaned_running( - conn: sqlite3.Connection, -) -> list[str]: +def reconcile_orphaned_running(conn: sqlite3.Connection) -> list[str]: """Requeue ``running`` cards with broken claim bookkeeping; returns their ids. A task ``running`` with NULL ``claim_lock``/``claim_expires`` (crash @@ -752,7 +687,6 @@ def _error_fingerprint(error_text: str) -> str: # nor extend it. Per-task ``max_retries`` overrides it. _PROTOCOL_VIOLATION_FAILURE_LIMIT = 3 - # Closed runs to walk when counting the streak; it trips at a handful anyway. _PROTOCOL_VIOLATION_SCAN_LIMIT = 50 @@ -778,17 +712,229 @@ def _protocol_violation_streak(conn: sqlite3.Connection, task_id: str) -> int: outcome = row["outcome"] or "" if outcome == "rate_limited": continue - if outcome == "crashed": - is_violation = bool(_kb._json_dict(row["metadata"]).get("protocol_violation")) - if not is_violation: - is_violation = "protocol violation" in (row["error"] or "") - if is_violation: - streak += 1 - continue + if outcome == "crashed" and ( + _kb._json_dict(row["metadata"]).get("protocol_violation") + or "protocol violation" in (row["error"] or "") + ): + streak += 1 + continue break return streak +_PROTOCOL_VIOLATION_ERROR = ( + "worker exited cleanly (rc=0) without calling " + "kanban_complete or kanban_block — protocol violation. " + "If the prior run already did the work, verify it and " + "report the result via kanban_complete; a run that ends " + "without a terminal kanban call counts as failed no " + "matter what it did." +) + + +@dataclass +class _DeadWorker: + """How ``detect_crashed_workers`` should book one dead worker.""" + + kind: str + code: Optional[int] + error_text: str + event_kind: str + event_payload: dict + protocol_violation: bool = False + rate_limited: bool = False + + @property + def run_outcome(self) -> str: + # A rate-limited requeue is recorded as ``rate_limited`` so board history + # doesn't show a phantom crash for a quota wall. + return "rate_limited" if self.rate_limited else "crashed" + + +def _classify_dead_worker(pid: int, claimer: Optional[str]) -> _DeadWorker: + """Map a dead worker's reaped exit status to its reclaim bookkeeping.""" + kind, code = _kb._classify_worker_exit(pid) + if kind == "clean_exit": + # rc=0 while still ``running``: usually the work succeeded and only the + # paperwork was skipped; the corrective sentence reaches the retry + # worker via ``build_worker_context``. + return _DeadWorker( + kind, code, _PROTOCOL_VIOLATION_ERROR, "protocol_violation", + # ``protocol_violation`` is the durable marker for + # _protocol_violation_streak: _end_run copies this payload into the + # run metadata. + {"pid": pid, "claimer": claimer, "exit_code": code, "protocol_violation": True}, + protocol_violation=True, + ) + if kind == "rate_limited": + # Quota wall — NOT a task failure. Release to the source phase and do + # NOT count a failure so a long quota window can't trip the breaker. + return _DeadWorker( + kind, code, + f"pid {pid} exited rate-limited (quota wall) — requeued without counting a failure", + "rate_limited", + {"pid": pid, "claimer": claimer, "exit_code": code}, + rate_limited=True, + ) + if kind == "nonzero_exit": + error_text = f"pid {pid} exited with code {code}" + elif kind == "signaled": + error_text = f"pid {pid} killed by signal {code}" + else: + error_text = f"pid {pid} not alive" + event_payload = {"pid": pid, "claimer": claimer} + if code is not None and kind != "unknown": + event_payload["exit_kind"] = kind + event_payload["exit_code"] = code + return _DeadWorker(kind, code, error_text, "crashed", event_payload) + + +@dataclass +class _CrashSweep: + """Everything ``detect_crashed_workers`` collects inside its reclaim txn.""" + + crashed: list[str] = field(default_factory=list) + rate_limited: list[str] = field(default_factory=list) + # ``(task_id, pid, claimer, protocol_violation, error_text)``: accounted + # after the txn via ``_record_task_failure`` (needs its own write_txn). + crash_details: list[tuple[str, int, str, bool, str]] = field(default_factory=list) + # Worker-exit observer payloads, fired only after every reclaim/accounting + # txn has committed. + exited_hook_payloads: list[dict] = field(default_factory=list) + + +def _reclaim_dead_workers(conn: sqlite3.Connection) -> _CrashSweep: + """Release every host-local ``running`` task whose worker PID is dead.""" + sweep = _CrashSweep() + with _kb.write_txn(conn): + rows = conn.execute( + "SELECT id, worker_pid, claim_lock, started_at, assignee " + "FROM tasks " + "WHERE status = 'running' AND worker_pid IS NOT NULL" + ).fetchall() + host_prefix = _kb._host_prefix() + for row in rows: + lock = row["claim_lock"] or "" + if not lock.startswith(host_prefix): + continue + # Launch-window grace so a freshly-spawned worker isn't reclaimed + # before its PID is visible on /proc. + started_at = _kb._row_get(row, "started_at") + if started_at is not None and time.time() - started_at < _kb._resolve_crash_grace_seconds(): + continue + if _kb._pid_alive(row["worker_pid"]): + continue + + pid = int(row["worker_pid"]) + dead = _classify_dead_worker(pid, row["claim_lock"]) + retry_status = _kb._retry_status_for_run(conn, row["id"]) + dead.event_payload["retry_status"] = retry_status + cur = conn.execute( + "UPDATE tasks SET status = ?, claim_lock = NULL, " + "claim_expires = NULL, worker_pid = NULL " + "WHERE id = ? AND status = 'running' " + " AND worker_pid = ? AND claim_lock IS ?", + (retry_status, row["id"], pid, row["claim_lock"]), + ) + if cur.rowcount != 1: + continue + run_id = _kb._end_run( + conn, row["id"], + outcome=dead.run_outcome, status=dead.run_outcome, + error=dead.error_text, + metadata=dict(dead.event_payload), + ) + _kb._append_event(conn, row["id"], dead.event_kind, dead.event_payload, run_id=run_id) + sweep.exited_hook_payloads.append({ + "task_id": row["id"], + "assignee": row["assignee"], + "run_id": run_id, + "worker_pid": pid, + "exit_kind": dead.kind, + "exit_code": dead.code, + "outcome": dead.run_outcome, + "retry_status": retry_status, + }) + if dead.rate_limited or dead.protocol_violation: + # Stamp last_failure_error WITHOUT touching ``consecutive_failures``: + # a rate-limited requeue must show ``check_respawn_guard`` a quota + # blocker; a below-budget protocol violation never reaches + # ``_record_task_failure`` (which stamps this column), yet the + # board UI and retry worker need the corrective message. + conn.execute( + "UPDATE tasks SET last_failure_error = ? WHERE id = ?", + (dead.error_text[:500], row["id"]), + ) + if dead.rate_limited: + sweep.rate_limited.append(row["id"]) + else: + sweep.crashed.append(row["id"]) + sweep.crash_details.append( + (row["id"], pid, row["claim_lock"], dead.protocol_violation, dead.error_text) + ) + return sweep + + +def _account_crashes(conn: sqlite3.Connection, crash_details: list) -> list[str]: + """Count each crash against the breaker; returns the task ids it tripped. + + Protocol violations get a BOUNDED violation-only budget independent of + ``consecutive_failures`` (per-task ``max_retries`` takes precedence); + systemic same-error crashes (>= 3 identical fingerprints this tick) trip + immediately. + """ + auto_blocked: list[str] = [] + fp_counts: dict[str, int] = {} + for _, _, _, _, err_text in crash_details: + fp = _error_fingerprint(err_text) + fp_counts[fp] = fp_counts.get(fp, 0) + 1 + for tid, pid, claimer, protocol_violation, error_text in crash_details: + if protocol_violation: + streak = _protocol_violation_streak(conn, tid) + trow = conn.execute("SELECT max_retries FROM tasks WHERE id = ?", (tid,)).fetchone() + if trow is None: + continue # task deleted mid-loop + task_override = _kb._row_get(trow, "max_retries") + violation_limit = ( + int(task_override) if task_override is not None else _PROTOCOL_VIOLATION_FAILURE_LIMIT + ) + if streak < violation_limit: + # Below budget: already back at ``ready`` with the error stamped. + # No ``_record_task_failure`` — must not consume the unified budget. + continue + # ``force_trip``: the decision (incl. per-task ``max_retries``) was + # already made against the violation streak above. + tripped = _kb._record_task_failure( + conn, tid, + error=error_text, + outcome="crashed", + failure_limit=violation_limit, + force_trip=True, + release_claim=False, + end_run=False, + event_payload_extra={ + "pid": pid, + "claimer": claimer, + "protocol_violations": streak, + "protocol_violation_limit": violation_limit, + }, + ) + else: + is_systemic = fp_counts.get(_error_fingerprint(error_text), 0) >= 3 + tripped = _kb._record_task_failure( + conn, tid, + error=error_text, + outcome="crashed", + failure_limit=1 if is_systemic else None, + release_claim=False, + end_run=False, + event_payload_extra={"pid": pid, "claimer": claimer}, + ) + if tripped: + auto_blocked.append(tid) + return auto_blocked + + def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]: """Reclaim ``running`` tasks whose worker PID is no longer alive. @@ -805,221 +951,19 @@ def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]: ids surface via the ``_last_rate_limited`` function attribute; the return stays crashed-only. """ - crashed: list[str] = [] - rate_limited: list[str] = [] - # Collected inside the main txn, used after it closes to run - # ``_record_task_failure`` (needs its own write_txn, can't nest). - # ``protocol_violation`` marks the clean-exit case, accounted against its - # own bounded streak instead of the unified failure counter. - crash_details: list[tuple[str, int, str, bool, str]] = [] - # (task_id, pid, claimer, protocol_violation, error_text) - # Worker-exit observer payloads, fired only after every reclaim/accounting - # txn has committed. - exited_hook_payloads: list[dict] = [] - with _kb.write_txn(conn): - rows = conn.execute( - "SELECT id, worker_pid, claim_lock, started_at, assignee " - "FROM tasks " - "WHERE status = 'running' AND worker_pid IS NOT NULL" - ).fetchall() - host_prefix = _kb._host_prefix() - for row in rows: - lock = row["claim_lock"] or "" - if not lock.startswith(host_prefix): - continue - # Launch-window grace so a freshly-spawned worker isn't reclaimed - # before its PID is visible on /proc. - started_at = _kb._row_get(row, "started_at") - if started_at is not None: - grace = _kb._resolve_crash_grace_seconds() - if time.time() - started_at < grace: - continue - if _kb._pid_alive(row["worker_pid"]): - continue - - pid = int(row["worker_pid"]) - kind, code = _kb._classify_worker_exit(pid) - rate_limited_exit = False - if kind == "clean_exit": - # rc=0 while still ``running``: usually the work succeeded and - # only the paperwork was skipped; the corrective sentence reaches - # the retry worker via ``build_worker_context``. - protocol_violation = True - error_text = ( - "worker exited cleanly (rc=0) without calling " - "kanban_complete or kanban_block — protocol violation. " - "If the prior run already did the work, verify it and " - "report the result via kanban_complete; a run that ends " - "without a terminal kanban call counts as failed no " - "matter what it did." - ) - event_kind = "protocol_violation" - event_payload = { - "pid": pid, - "claimer": row["claim_lock"], - "exit_code": code, - # Durable marker for _protocol_violation_streak: _end_run - # copies this payload into the run metadata. - "protocol_violation": True, - } - elif kind == "rate_limited": - # Quota wall — NOT a task failure. Release to the source phase - # and do NOT count a failure so a long quota window can't trip the breaker. - protocol_violation = False - rate_limited_exit = True - error_text = ( - f"pid {pid} exited rate-limited (quota wall) — " - f"requeued without counting a failure" - ) - event_kind = "rate_limited" - event_payload = { - "pid": pid, - "claimer": row["claim_lock"], - "exit_code": code, - } - else: - protocol_violation = False - if kind == "nonzero_exit": - error_text = f"pid {pid} exited with code {code}" - elif kind == "signaled": - error_text = f"pid {pid} killed by signal {code}" - else: - error_text = f"pid {pid} not alive" - event_kind = "crashed" - event_payload = {"pid": pid, "claimer": row["claim_lock"]} - if code is not None and kind != "unknown": - event_payload["exit_kind"] = kind - event_payload["exit_code"] = code - - retry_status = _kb._retry_status_for_run(conn, row["id"]) - event_payload["retry_status"] = retry_status - cur = conn.execute( - "UPDATE tasks SET status = ?, claim_lock = NULL, " - "claim_expires = NULL, worker_pid = NULL " - "WHERE id = ? AND status = 'running' " - " AND worker_pid = ? AND claim_lock IS ?", - (retry_status, row["id"], pid, row["claim_lock"]), - ) - if cur.rowcount == 1: - # Record a rate-limited requeue as ``rate_limited`` so board - # history doesn't show a phantom crash for a quota wall. - _run_outcome = "rate_limited" if rate_limited_exit else "crashed" - run_id = _kb._end_run( - conn, row["id"], - outcome=_run_outcome, status=_run_outcome, - error=error_text, - metadata=dict(event_payload), - ) - _kb._append_event( - conn, row["id"], event_kind, - event_payload, - run_id=run_id, - ) - exited_hook_payloads.append({ - "task_id": row["id"], - "assignee": row["assignee"], - "run_id": run_id, - "worker_pid": pid, - "exit_kind": kind, - "exit_code": code, - "outcome": _run_outcome, - "retry_status": retry_status, - }) - if rate_limited_exit: - # Stamp last_failure_error so ``check_respawn_guard`` sees a - # quota blocker — WITHOUT touching ``consecutive_failures``. - conn.execute( - "UPDATE tasks SET last_failure_error = ? WHERE id = ?", - (error_text[:500], row["id"]), - ) - rate_limited.append(row["id"]) - else: - if protocol_violation: - # A below-budget violation never reaches - # ``_record_task_failure`` (which stamps this column), yet - # the board UI and retry worker need the corrective message. - conn.execute( - "UPDATE tasks SET last_failure_error = ? " - "WHERE id = ?", - (error_text[:500], row["id"]), - ) - crashed.append(row["id"]) - crash_details.append( - (row["id"], pid, row["claim_lock"], - protocol_violation, error_text) - ) + sweep = _reclaim_dead_workers(conn) # Outside the main txn: account each crash and maybe trip the breaker. - # Protocol violations get a BOUNDED violation-only budget independent of - # ``consecutive_failures`` (per-task ``max_retries`` takes precedence); - # systemic same-error crashes still trip immediately. - auto_blocked: list[str] = [] - if crash_details: - _fp_counts: dict[str, int] = {} - for _, _, _, _, err_text in crash_details: - fp = _error_fingerprint(err_text) - _fp_counts[fp] = _fp_counts.get(fp, 0) + 1 - for tid, pid, claimer, protocol_violation, error_text in crash_details: - if protocol_violation: - streak = _protocol_violation_streak(conn, tid) - trow = conn.execute( - "SELECT max_retries FROM tasks WHERE id = ?", (tid,), - ).fetchone() - if trow is None: - continue # task deleted mid-loop - task_override = _kb._row_get(trow, "max_retries") - violation_limit = ( - int(task_override) - if task_override is not None - else _PROTOCOL_VIOLATION_FAILURE_LIMIT - ) - if streak < violation_limit: - # Below budget: already back at ``ready`` with the error - # stamped. No ``_record_task_failure`` — must not consume - # the unified failure budget. - continue - # ``force_trip``: the decision (incl. per-task ``max_retries``) - # was already made against the violation streak above. - tripped = _kb._record_task_failure( - conn, tid, - error=error_text, - outcome="crashed", - failure_limit=violation_limit, - force_trip=True, - release_claim=False, - end_run=False, - event_payload_extra={ - "pid": pid, - "claimer": claimer, - "protocol_violations": streak, - "protocol_violation_limit": violation_limit, - }, - ) - if tripped: - auto_blocked.append(tid) - continue - fp = _error_fingerprint(error_text) - is_systemic = _fp_counts.get(fp, 0) >= 3 - tripped = _kb._record_task_failure( - conn, tid, - error=error_text, - outcome="crashed", - failure_limit=1 if is_systemic else None, - release_claim=False, - end_run=False, - event_payload_extra={"pid": pid, "claimer": claimer}, - ) - if tripped: - auto_blocked.append(tid) + auto_blocked = _account_crashes(conn, sweep.crash_details) if sweep.crash_details else [] # Side-channel attributes keep the public ``list[str]`` return stable; # ``dispatch_once`` reads them to populate ``DispatchResult``. Rate-limited # requeues did NOT count a failure and are NOT crashes. detect_crashed_workers._last_auto_blocked = auto_blocked # type: ignore[attr-defined] - detect_crashed_workers._last_rate_limited = rate_limited # type: ignore[attr-defined] + detect_crashed_workers._last_rate_limited = sweep.rate_limited # type: ignore[attr-defined] # Fired only now, after the reclaim txn AND breaker accounting have # committed, so subscribers always observe fully durable board state. - if exited_hook_payloads and _kb._kanban_observer_consumed("on_kanban_worker_exited"): + if sweep.exited_hook_payloads and _kb._kanban_observer_consumed("on_kanban_worker_exited"): _board = _kb.get_current_board() - for hook_fields in exited_hook_payloads: + for hook_fields in sweep.exited_hook_payloads: hook_fields = dict(hook_fields) _kb._fire_kanban_lifecycle_hook( "on_kanban_worker_exited", @@ -1027,7 +971,7 @@ def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]: board=_board, **hook_fields, ) - return crashed + return sweep.crashed def _record_task_failure( @@ -1059,7 +1003,7 @@ def _record_task_failure( """ if failure_limit is None: failure_limit = DEFAULT_FAILURE_LIMIT - blocked = False + error = error[:500] with _kb.write_txn(conn): row = conn.execute( "SELECT consecutive_failures, status, max_retries, current_run_id " @@ -1077,53 +1021,11 @@ def _record_task_failure( # Per-task override wins over caller-supplied and default thresholds. task_override = _kb._row_get(row, "max_retries") if task_override is not None: - effective_limit = int(task_override) - limit_source = "task" + effective_limit, limit_source = int(task_override), "task" else: - effective_limit = int(failure_limit) - limit_source = "dispatcher" + effective_limit, limit_source = int(failure_limit), "dispatcher" - if force_trip or failures >= effective_limit: - # Spawn path (release_claim) is still running and also clears claim - # state; the timeout/crash path already did. - conn.execute( - "UPDATE tasks SET status = 'blocked', " - + ("claim_lock = NULL, claim_expires = NULL, worker_pid = NULL, " - if release_claim else "") - + "consecutive_failures = ?, last_failure_error = ? " - "WHERE id = ? AND status IN ('running', 'ready', 'review')", - (failures, error[:500], task_id), - ) - run_id = None - if end_run: - # Only the spawn path has an open run to close. - run_id = _kb._end_run( - conn, task_id, - outcome="gave_up", status="gave_up", - error=error[:500], - metadata={ - "failures": failures, - "trigger_outcome": outcome, - "effective_limit": effective_limit, - "limit_source": limit_source, - "retry_status": retry_status, - }, - ) - payload = { - "failures": failures, - "effective_limit": effective_limit, - "limit_source": limit_source, - "error": error[:500], - "trigger_outcome": outcome, - "retry_status": retry_status, - } - if event_payload_extra: - payload.update(event_payload_extra) - _kb._append_event( - conn, task_id, "gave_up", payload, run_id=run_id, - ) - blocked = True - else: + if not (force_trip or failures >= effective_limit): if release_claim: # Spawn path: restore the claimed source phase + clear claim. conn.execute( @@ -1131,35 +1033,62 @@ def _record_task_failure( "claim_expires = NULL, worker_pid = NULL, " "consecutive_failures = ?, last_failure_error = ? " "WHERE id = ? AND status = 'running'", - (retry_status, failures, error[:500], task_id), + (retry_status, failures, error, task_id), ) else: conn.execute( "UPDATE tasks SET consecutive_failures = ?, " "last_failure_error = ? WHERE id = ?", - (failures, error[:500], task_id), + (failures, error, task_id), ) + # Timeout/crash path's caller already emitted its own event. if end_run: run_id = _kb._end_run( - conn, task_id, - outcome=outcome, status=outcome, - error=error[:500], - metadata={ - "failures": failures, - "retry_status": retry_status, - }, + conn, task_id, outcome=outcome, status=outcome, error=error, + metadata={"failures": failures, "retry_status": retry_status}, ) _kb._append_event( conn, task_id, outcome, - { - "error": error[:500], - "failures": failures, - "retry_status": retry_status, - }, + {"error": error, "failures": failures, "retry_status": retry_status}, run_id=run_id, ) - # Timeout/crash path's caller already emitted its own event. - return blocked + return False + + # Spawn path (release_claim) is still running and also clears claim + # state; the timeout/crash path already did. + conn.execute( + "UPDATE tasks SET status = 'blocked', " + + ("claim_lock = NULL, claim_expires = NULL, worker_pid = NULL, " + if release_claim else "") + + "consecutive_failures = ?, last_failure_error = ? " + "WHERE id = ? AND status IN ('running', 'ready', 'review')", + (failures, error, task_id), + ) + payload = { + "failures": failures, + "effective_limit": effective_limit, + "limit_source": limit_source, + "error": error, + "trigger_outcome": outcome, + "retry_status": retry_status, + } + run_id = None + if end_run: + # Only the spawn path has an open run to close. + run_id = _kb._end_run( + conn, task_id, outcome="gave_up", status="gave_up", error=error, + metadata={ + "failures": failures, + "trigger_outcome": outcome, + "effective_limit": effective_limit, + "limit_source": limit_source, + "retry_status": retry_status, + }, + ) + if event_payload_extra: + payload.update(event_payload_extra) + _kb._append_event(conn, task_id, "gave_up", payload, run_id=run_id) + return True # Backward-compat alias; new code should call ``_record_task_failure``. @@ -1182,16 +1111,10 @@ def _record_spawn_failure( def _set_worker_pid(conn: sqlite3.Connection, task_id: str, pid: int) -> None: """Record the spawned child's pid + emit a ``spawned`` event carrying it.""" with _kb.write_txn(conn): - conn.execute( - "UPDATE tasks SET worker_pid = ? WHERE id = ?", - (int(pid), task_id), - ) + conn.execute("UPDATE tasks SET worker_pid = ? WHERE id = ?", (int(pid), task_id)) run_id = _kb._current_run_id(conn, task_id) if run_id is not None: - conn.execute( - "UPDATE task_runs SET worker_pid = ? WHERE id = ?", - (int(pid), run_id), - ) + conn.execute("UPDATE task_runs SET worker_pid = ? WHERE id = ?", (int(pid), run_id)) _kb._append_event(conn, task_id, "spawned", {"pid": int(pid)}, run_id=run_id) @@ -1210,10 +1133,6 @@ def _clear_failure_counter(conn: sqlite3.Connection, task_id: str) -> None: ) -# Legacy alias for test-code and anything else that still imports it. -_clear_spawn_failures = _clear_failure_counter - - def check_respawn_guard( conn: sqlite3.Connection, task_id: str, *, lane: str = "ready", ) -> Optional[str]: @@ -1263,10 +1182,7 @@ def check_respawn_guard( "ORDER BY ended_at DESC LIMIT 1", (task_id,), ).fetchone() - if ( - latest_run is not None - and latest_run["outcome"] == "rate_limited" - ): + if latest_run is not None and latest_run["outcome"] == "rate_limited": if rl_cooldown <= 0: # Cooldown disabled — respawn immediately, skipping blocker_auth so # the stamped rate-limit text doesn't re-trap the task. @@ -1324,6 +1240,17 @@ def check_respawn_guard( return None +def _profile_exists_fn() -> Optional[Callable[[str], bool]]: + """``hermes_cli.profiles.profile_exists``, or ``None`` when it cannot be + imported (local import avoids a cycle; callers fall back to trusting the + assignee).""" + try: + from hermes_cli.profiles import profile_exists + except Exception: + return None + return profile_exists + + def _has_spawnable(conn: sqlite3.Connection, status: str) -> bool: rows = conn.execute( "SELECT DISTINCT assignee FROM tasks " @@ -1332,9 +1259,8 @@ def _has_spawnable(conn: sqlite3.Connection, status: str) -> bool: ).fetchall() if not rows: return False - try: - from hermes_cli.profiles import profile_exists # local import: avoids cycle - except Exception: + profile_exists = _profile_exists_fn() + if profile_exists is None: # Can't introspect — assume spawnable, preserve legacy behavior. return True return any(profile_exists(row["assignee"]) for row in rows) @@ -1361,9 +1287,7 @@ def review_dispatch_enabled() -> bool: """ try: from hermes_cli.config import load_config - return bool( - (load_config() or {}).get("kanban", {}).get("review_dispatch", True) - ) + return bool((load_config() or {}).get("kanban", {}).get("review_dispatch", True)) except Exception: return True @@ -1383,12 +1307,9 @@ def review_dispatch_enabled() -> bool: # so the cap errs toward fewer workers on small VMs. MEMORY_GUARD_MB_PER_WORKER = 512 - # Derived default bounds: never below 2 (smallest VM must still progress), # never above 8 (more fan-out must be explicit in config). DERIVED_MAX_IN_PROGRESS_FLOOR = 2 - - DERIVED_MAX_IN_PROGRESS_CEILING = 8 @@ -1418,10 +1339,7 @@ def derive_default_max_in_progress(sample: Optional[Mapping[str, Any]] = None) - if isinstance(total_kib, bool) or not isinstance(total_kib, int) or total_kib <= 0: return None workers = (total_kib // 1024) // MEMORY_GUARD_MB_PER_WORKER - return max( - DERIVED_MAX_IN_PROGRESS_FLOOR, - min(workers, DERIVED_MAX_IN_PROGRESS_CEILING), - ) + return max(DERIVED_MAX_IN_PROGRESS_FLOOR, min(workers, DERIVED_MAX_IN_PROGRESS_CEILING)) def resolve_max_in_progress(configured: Optional[int]) -> Optional[int]: @@ -1442,9 +1360,7 @@ def configured_max_in_progress() -> Optional[int]: """ try: from hermes_cli.config import load_config_readonly - raw = (load_config_readonly() or {}).get("kanban", {}).get( - "max_in_progress" - ) + raw = (load_config_readonly() or {}).get("kanban", {}).get("max_in_progress") except Exception: return None if raw is None: @@ -1524,9 +1440,7 @@ def _memory_pressure_level(sample: Optional[Mapping[str, Any]] = None) -> str: return "unknown" try: from gateway.memory_status import classify_pressure - return classify_pressure( - sample.get("mem_available_kib"), sample.get("mem_total_kib") - ) + return classify_pressure(sample.get("mem_available_kib"), sample.get("mem_total_kib")) except Exception: return "unknown" @@ -1554,11 +1468,8 @@ def dispatch_once( ``skipped_locked=True`` and writes nothing; the lock is keyed on the resolved DB path so unrelated boards tick in parallel. """ - try: - db_path = _kb.kanban_db_path(board=board) - except Exception: - # Must not lose the tick — fall through to an unguarded dispatch. - result = _dispatch_once_locked( + def _locked_tick() -> DispatchResult: + return _dispatch_once_locked( conn, spawn_fn=spawn_fn, ttl_seconds=ttl_seconds, @@ -1572,26 +1483,19 @@ def dispatch_once( max_in_progress_per_profile=max_in_progress_per_profile, reconcile_orphans=reconcile_orphans, ) + + try: + db_path = _kb.kanban_db_path(board=board) + except Exception: + # Must not lose the tick — fall through to an unguarded dispatch. + result = _locked_tick() _kb._fire_dispatch_tick_hook(result, board=board, dry_run=dry_run) return result with _kb._dispatch_tick_lock(db_path) as held: if not held: result = DispatchResult(skipped_locked=True) else: - result = _dispatch_once_locked( - conn, - spawn_fn=spawn_fn, - ttl_seconds=ttl_seconds, - dry_run=dry_run, - max_spawn=max_spawn, - max_in_progress=max_in_progress, - failure_limit=failure_limit, - stale_timeout_seconds=stale_timeout_seconds, - board=board, - default_assignee=default_assignee, - max_in_progress_per_profile=max_in_progress_per_profile, - reconcile_orphans=reconcile_orphans, - ) + result = _locked_tick() # Still under the dispatch lock: periodic PASSIVE WAL checkpoint. _kb._maybe_checkpoint_wal(conn, db_path) # Lock released. Fire the tick observer strictly OUTSIDE the critical @@ -1600,6 +1504,19 @@ def dispatch_once( return result +def _call_spawn_fn(spawn_fn, task: Task, workspace: str, board: Optional[str]) -> Optional[int]: + """Back-compat: older spawn_fn signatures (and test stubs) accept only + ``(task, workspace)``; pass ``board`` only when the callable supports it.""" + import inspect + try: + sig = inspect.signature(spawn_fn) + if "board" in sig.parameters: + return spawn_fn(task, workspace, board=board) + return spawn_fn(task, workspace) + except (TypeError, ValueError): + return spawn_fn(task, workspace) + + def _dispatch_lane_task( conn: sqlite3.Connection, row: sqlite3.Row, @@ -1624,10 +1541,7 @@ def _dispatch_lane_task( # would fail ``hermes -p `` at startup and loop ready→crash→ready # forever. Bucketed apart from skipped_unassigned: the operator cannot fix # it by assigning a profile, and health telemetry suppresses "stuck" for it. - try: - from hermes_cli.profiles import profile_exists # local import: avoids cycle - except Exception: - profile_exists = None # type: ignore[assignment] + profile_exists = _profile_exists_fn() if profile_exists is not None and not profile_exists(assignee): result.skipped_nonspawnable.append(task_id) return False @@ -1646,11 +1560,16 @@ def _dispatch_lane_task( with _kb.write_txn(conn): _kb._append_event(conn, task_id, "respawn_guarded", {"reason": guard_reason}) return False + + def _count_spawn(name: str) -> None: + # Later rows in this tick respect the per-profile cap; subsequent + # ticks re-query from the DB. + if per_profile_cap is not None and name: + per_profile_running[name] = per_profile_running.get(name, 0) + 1 + if dry_run: result.spawned.append((task_id, assignee, "")) - # Count the would-be spawn so the cap check sees it on later rows. - if per_profile_cap is not None and assignee: - per_profile_running[assignee] = per_profile_running.get(assignee, 0) + 1 + _count_spawn(assignee) return True claim = _kb.claim_review_task if lane == "review" else _kb.claim_task claimed = claim(conn, task_id, ttl_seconds=ttl_seconds) @@ -1674,19 +1593,8 @@ def _dispatch_lane_task( # Force-load sdlc-review; the kanban lifecycle is already in every # worker's system prompt via KANBAN_GUIDANCE. claimed.skills = list(dict.fromkeys([*(claimed.skills or []), "sdlc-review"])) - _spawn = spawn_fn if spawn_fn is not None else _default_spawn try: - # Back-compat: older spawn_fn signatures (and test stubs) accept only - # (task, workspace); pass ``board`` only when the callable supports it. - import inspect - try: - sig = inspect.signature(_spawn) - if "board" in sig.parameters: - pid = _spawn(claimed, str(workspace), board=board) - else: - pid = _spawn(claimed, str(workspace)) - except (TypeError, ValueError): - pid = _spawn(claimed, str(workspace)) + pid = _call_spawn_fn(spawn_fn if spawn_fn is not None else _default_spawn, claimed, str(workspace), board) if pid: _set_worker_pid(conn, claimed.id, int(pid)) # Fires AFTER the PID (when reported) is durably persisted. Best-effort. @@ -1695,10 +1603,7 @@ def _dispatch_lane_task( # spawn would let a task that keeps timing out loop forever. Cleared # only on successful completion (complete_task). result.spawned.append((claimed.id, claimed.assignee or "", str(workspace))) - # Later rows in this tick respect the per-profile cap; subsequent - # ticks re-query from the DB. - if per_profile_cap is not None and claimed.assignee: - per_profile_running[claimed.assignee] = per_profile_running.get(claimed.assignee, 0) + 1 + _count_spawn(claimed.assignee) return True except Exception as exc: if _record_spawn_failure(conn, claimed.id, str(exc), failure_limit=failure_limit): @@ -1737,6 +1642,126 @@ def _apply_default_assignee( return True +def _run_reclaim_phase( + conn: sqlite3.Connection, + result: DispatchResult, + *, + stale_timeout_seconds: int, + failure_limit: int, + reconcile_orphans: bool, +) -> None: + """Reclaim stale/orphaned/crashed/timed-out running tasks, then promote.""" + reap_worker_zombies() + result.reclaimed = _kb.release_stale_claims(conn) + if reconcile_orphans: + result.reconciled_orphans = reconcile_orphaned_running(conn) + result.stale = detect_stale_running(conn, stale_timeout_seconds=stale_timeout_seconds) + result.crashed = detect_crashed_workers(conn) + # Side-channel attributes (see detect_crashed_workers); rate-limited tasks + # went back to ``ready`` and the respawn guard defers them until quota clears. + result.auto_blocked.extend(getattr(detect_crashed_workers, "_last_auto_blocked", [])) + result.rate_limited.extend(getattr(detect_crashed_workers, "_last_rate_limited", [])) + result.timed_out = enforce_max_runtime(conn) + result.promoted = _kb.recompute_ready(conn, failure_limit=failure_limit) + + +def _tick_spawn_budget( + conn: sqlite3.Connection, + result: DispatchResult, + *, + max_spawn: Optional[int], + max_in_progress: Optional[int], + board: Optional[str], +) -> tuple[bool, Optional[int]]: + """``(may_spawn, spawn_budget)`` for this tick; ``budget None`` = uncapped. + + ``max_spawn`` is a live per-board concurrency cap (running + this tick's + spawns), not a per-tick budget — a per-tick reading would grow concurrency + by N every tick. ``max_in_progress`` is a HOST-level cap: running workers on + every other board count against the same budget, else N boards multiply the + cap by N — exactly the fan-out the memory-derived default exists to prevent. + """ + # Count already-running tasks so max_spawn enforces concurrency, not a + # per-tick budget: "running" tasks stay running until the worker calls + # kanban_complete/kanban_block or the TTL reclaims them. + running_count = 0 + spawn_budget: Optional[int] = None + if max_spawn is not None or max_in_progress is not None: + running_count = count_running_tasks(conn) + + # Both ready and review loops consume from the same budget. + if max_spawn is not None: + if running_count >= max_spawn: + return False, None + spawn_budget = max_spawn - running_count + + if max_in_progress is not None: + total_running = running_count + count_running_tasks_other_boards(board) + if total_running >= max_in_progress: + return False, None + remaining = max_in_progress - total_running + if spawn_budget is None or spawn_budget > remaining: + spawn_budget = remaining + + # Memory-pressure guard: a static cap can't see the host's actual state. + # critical -> spawn nothing this tick; elevated -> at most one new worker. + # Reclaim/promotion already ran, so bookkeeping stays live; deferred tasks + # wait for a later tick. "unknown" imposes no restriction. + pressure = _memory_pressure_level() + if pressure == "critical": + result.memory_pressure = pressure + _kb._log.warning( + "kanban dispatch: system memory pressure is critical; " + "spawning no new workers this tick (deferred, not dropped)" + ) + return False, None + if pressure == "elevated": + result.memory_pressure = pressure + if spawn_budget is None or spawn_budget > 1: + _kb._log.warning( + "kanban dispatch: system memory pressure is elevated; " + "limiting to at most 1 new worker this tick" + ) + spawn_budget = 1 + return True, spawn_budget + + +def _lane_rows(conn: sqlite3.Connection, status: str) -> list[sqlite3.Row]: + """Unclaimed rows of one lane in dispatch order.""" + return conn.execute( + "SELECT id, assignee FROM tasks " + f"WHERE status = '{status}' AND claim_lock IS NULL " + "ORDER BY priority DESC, created_at ASC" + ).fetchall() + + +def _any_spawnable_review(review_rows: list[sqlite3.Row]) -> bool: + """Mirrors the review loop's own gate so human-pulled control-plane lanes + don't tax ready throughput; assumes spawnable when profiles are unimportable.""" + if not review_rows: + return False + profile_exists = _profile_exists_fn() + if profile_exists is None: + return any(row["assignee"] for row in review_rows) + return any(row["assignee"] and profile_exists(row["assignee"]) for row in review_rows) + + +def _resolve_default_assignee(default_assignee: Optional[str]) -> Optional[str]: + """``kanban.default_assignee`` when it names a real profile. When the + profiles module isn't importable trust the operator's config: the + downstream profile_exists check still buckets a missing profile as + nonspawnable.""" + name = (default_assignee or "").strip() or None + if name: + try: + from hermes_cli.profiles import profile_exists + if not profile_exists(name): + return None + except Exception: + pass + return name + + def _dispatch_once_locked( conn: sqlite3.Connection, *, @@ -1757,138 +1782,52 @@ def _dispatch_once_locked( parents are done, then for each ready task with an assignee atomically claim and call ``spawn_fn(task, workspace_path, board) -> Optional[int]``, recording the PID as ``worker_pid`` so later ticks catch crashes before - the TTL expires. - - After ``failure_limit`` consecutive failures the task is auto-blocked. - ``max_spawn`` is a live per-board concurrency cap (running + this tick's - spawns), not a per-tick budget — a per-tick reading would grow concurrency - by N every tick. ``max_in_progress`` is a host-level cap over every active - board: workers share one machine's memory, so a per-board reading would - multiply it by the board count. + the TTL expires. After ``failure_limit`` consecutive failures the task is + auto-blocked. Cap semantics: :func:`_tick_spawn_budget`. """ - reap_worker_zombies() - result = DispatchResult() - result.reclaimed = _kb.release_stale_claims(conn) - if reconcile_orphans: - result.reconciled_orphans = reconcile_orphaned_running(conn) - result.stale = detect_stale_running( - conn, stale_timeout_seconds=stale_timeout_seconds, + _run_reclaim_phase( + conn, result, stale_timeout_seconds=stale_timeout_seconds, + failure_limit=failure_limit, reconcile_orphans=reconcile_orphans, ) - result.crashed = detect_crashed_workers(conn) - # Side-channel attributes (see detect_crashed_workers); rate-limited tasks - # went back to ``ready`` and the respawn guard defers them until quota clears. - result.auto_blocked.extend(getattr(detect_crashed_workers, "_last_auto_blocked", [])) - result.rate_limited.extend(getattr(detect_crashed_workers, "_last_rate_limited", [])) - result.timed_out = enforce_max_runtime(conn) - result.promoted = _kb.recompute_ready(conn, failure_limit=failure_limit) - - # Count already-running tasks so max_spawn enforces concurrency, not a - # per-tick budget: "running" tasks stay running until the worker calls - # kanban_complete/kanban_block or the TTL reclaims them. - running_count = 0 - spawn_budget: Optional[int] = None - if max_spawn is not None or max_in_progress is not None: - running_count = count_running_tasks(conn) - - # Both ready and review loops consume from the same budget. - if max_spawn is not None: - if running_count >= max_spawn: - return result - spawn_budget = max_spawn - running_count - - # max_in_progress is a HOST-level cap: running workers on every other board - # count against the same budget, else N boards multiply the cap by N — - # exactly the fan-out the memory-derived default exists to prevent. - if max_in_progress is not None: - total_running = running_count + count_running_tasks_other_boards(board) - if total_running >= max_in_progress: - return result - remaining = max_in_progress - total_running - if spawn_budget is None or spawn_budget > remaining: - spawn_budget = remaining - - # Memory-pressure guard: a static cap can't see the host's actual state. - # critical -> spawn nothing this tick; elevated -> at most one new worker. - # Reclaim/promotion already ran, so bookkeeping stays live; deferred tasks - # wait for a later tick. "unknown" imposes no restriction. - pressure = _memory_pressure_level() - if pressure == "critical": - result.memory_pressure = pressure - _kb._log.warning( - "kanban dispatch: system memory pressure is critical; " - "spawning no new workers this tick (deferred, not dropped)" - ) + may_spawn, spawn_budget = _tick_spawn_budget( + conn, result, max_spawn=max_spawn, max_in_progress=max_in_progress, board=board, + ) + if not may_spawn: return result - if pressure == "elevated": - result.memory_pressure = pressure - if spawn_budget is None or spawn_budget > 1: - _kb._log.warning( - "kanban dispatch: system memory pressure is elevated; " - "limiting to at most 1 new worker this tick" - ) - spawn_budget = 1 - ready_rows = conn.execute( - "SELECT id, assignee FROM tasks " - "WHERE status = 'ready' AND claim_lock IS NULL " - "ORDER BY priority DESC, created_at ASC" - ).fetchall() + ready_rows = _lane_rows(conn, "ready") # Review rows are enumerated up front so the budget split can see whether # review work exists at all. - review_rows = [] - if review_dispatch_enabled(): - review_rows = conn.execute( - "SELECT id, assignee FROM tasks " - "WHERE status = 'review' AND claim_lock IS NULL " - "ORDER BY priority DESC, created_at ASC" - ).fetchall() + review_rows = _lane_rows(conn, "review") if review_dispatch_enabled() else [] # Review-lane reservation: the ready loop runs first and would otherwise # consume the ENTIRE shared budget, starving reviews under a sustained ready # backlog. When spawnable review work exists and there is any budget, hold - # one slot back. "Spawnable" mirrors the review loop's own gate so - # human-pulled control-plane lanes don't tax ready throughput. - def _any_spawnable_review() -> bool: - if not review_rows: - return False - try: - from hermes_cli.profiles import profile_exists as _rpe - except Exception: - # Assume spawnable, matching the review loop's own fallback. - return any(row["assignee"] for row in review_rows) - return any( - row["assignee"] and _rpe(row["assignee"]) for row in review_rows - ) - + # one slot back. ready_budget = spawn_budget - if spawn_budget is not None and spawn_budget > 0 and _any_spawnable_review(): + if spawn_budget is not None and spawn_budget > 0 and _any_spawnable_review(review_rows): ready_budget = max(spawn_budget - 1, 0) - spawned = 0 # Per-profile cap. Deferred tasks go to skipped_per_profile_capped, not # skipped_unassigned — "busy, retry later" differs from "needs routing". - _per_profile_cap = max_in_progress_per_profile if ( + per_profile_cap = max_in_progress_per_profile if ( isinstance(max_in_progress_per_profile, int) and max_in_progress_per_profile > 0 ) else None - _per_profile_running: dict[str, int] = {} - if _per_profile_cap is not None: + per_profile_running: dict[str, int] = {} + if per_profile_cap is not None: for prow in conn.execute( "SELECT assignee, COUNT(*) AS n FROM tasks " "WHERE status = 'running' AND assignee IS NOT NULL " "GROUP BY assignee" ): - _per_profile_running[prow["assignee"]] = int(prow["n"]) - # When the profiles module isn't importable trust the operator's config: - # the downstream profile_exists check still buckets a missing profile as - # nonspawnable. - _default_assignee = (default_assignee or "").strip() or None - if _default_assignee: - try: - from hermes_cli.profiles import profile_exists as _pe - if not _pe(_default_assignee): - _default_assignee = None - except Exception: - pass + per_profile_running[prow["assignee"]] = int(prow["n"]) + lane_kwargs: dict[str, Any] = dict( + dry_run=dry_run, ttl_seconds=ttl_seconds, board=board, + failure_limit=failure_limit, spawn_fn=spawn_fn, + per_profile_cap=per_profile_cap, per_profile_running=per_profile_running, + ) + default_assignee = _resolve_default_assignee(default_assignee) + spawned = 0 for row in ready_rows: if ready_budget is not None and spawned >= ready_budget: break @@ -1896,22 +1835,16 @@ def _dispatch_once_locked( if not row_assignee: # Honour kanban.default_assignee so an unassigned task doesn't # park in 'ready' forever. - if not _default_assignee or not _apply_default_assignee( - conn, row["id"], _default_assignee, dry_run=dry_run, + if not default_assignee or not _apply_default_assignee( + conn, row["id"], default_assignee, dry_run=dry_run, ): result.skipped_unassigned.append(row["id"]) continue - row_assignee = _default_assignee + row_assignee = default_assignee result.auto_assigned_default.append(row["id"]) - if _dispatch_lane_task( - conn, row, row_assignee, result, lane="ready", dry_run=dry_run, - ttl_seconds=ttl_seconds, board=board, failure_limit=failure_limit, - spawn_fn=spawn_fn, per_profile_cap=_per_profile_cap, - per_profile_running=_per_profile_running, - ): + if _dispatch_lane_task(conn, row, row_assignee, result, lane="ready", **lane_kwargs): spawned += 1 - # ---- review column dispatch ---- # A review agent (sdlc-review) approves (→ done) or requests changes # (→ ready/todo). Review spawns share max_spawn with ready tasks. The loop # checks the FULL shared ``spawn_budget`` — the reservation above caps the @@ -1922,12 +1855,7 @@ def _dispatch_once_locked( if not row["assignee"]: result.skipped_unassigned.append(row["id"]) continue - if _dispatch_lane_task( - conn, row, row["assignee"], result, lane="review", dry_run=dry_run, - ttl_seconds=ttl_seconds, board=board, failure_limit=failure_limit, - spawn_fn=spawn_fn, per_profile_cap=_per_profile_cap, - per_profile_running=_per_profile_running, - ): + if _dispatch_lane_task(conn, row, row["assignee"], result, lane="review", **lane_kwargs): spawned += 1 return result @@ -1952,16 +1880,9 @@ def worker_log_rotation_config(kanban_cfg: Optional[dict] = None) -> tuple[int, kanban_cfg = (load_config().get("kanban") or {}) except Exception: kanban_cfg = {} - max_bytes = _positive_int( - (kanban_cfg or {}).get("worker_log_rotate_bytes"), - DEFAULT_LOG_ROTATE_BYTES, - minimum=1, - ) - backup_count = _positive_int( - (kanban_cfg or {}).get("worker_log_backup_count"), - DEFAULT_LOG_BACKUP_COUNT, - minimum=0, - ) + kanban_cfg = kanban_cfg or {} + max_bytes = _positive_int(kanban_cfg.get("worker_log_rotate_bytes"), DEFAULT_LOG_ROTATE_BYTES, minimum=1) + backup_count = _positive_int(kanban_cfg.get("worker_log_backup_count"), DEFAULT_LOG_BACKUP_COUNT, minimum=0) return max_bytes, backup_count @@ -1978,24 +1899,16 @@ def _rotate_worker_log( older generations shift up to ``backup_count``. """ try: - if not log_path.exists(): + if not log_path.exists() or log_path.stat().st_size <= max_bytes: return - if log_path.stat().st_size <= max_bytes: - return - backup_count = _positive_int( - backup_count, - DEFAULT_LOG_BACKUP_COUNT, - minimum=0, - ) + backup_count = _positive_int(backup_count, DEFAULT_LOG_BACKUP_COUNT, minimum=0) if backup_count == 0: log_path.unlink() return oldest = _rotated_log_path(log_path, backup_count) - try: + with contextlib.suppress(OSError): if oldest.exists(): oldest.unlink() - except OSError: - pass for generation in range(backup_count - 1, 0, -1): src = _rotated_log_path(log_path, generation) if not src.exists(): @@ -2008,9 +1921,8 @@ def _rotate_worker_log( def _module_hermes_argv() -> list[str]: - """Return the interpreter-bound Hermes CLI invocation.""" - # ``hermes_cli.main`` is the console-script target — there is no top-level - # ``hermes`` package to import. + """Interpreter-bound Hermes CLI invocation (``hermes_cli.main`` is the + console-script target — there is no top-level ``hermes`` package).""" return [sys.executable, "-m", "hermes_cli.main"] @@ -2042,8 +1954,7 @@ def _path_search_names(command: str) -> list[str]: if not _kb._IS_WINDOWS or os.path.splitext(command)[1]: return [command] raw = os.environ.get("PATHEXT") or ".COM;.EXE;.BAT;.CMD" - exts = [ext for ext in raw.split(";") if ext] - return [command + ext for ext in exts] + return [command + ext for ext in raw.split(";") if ext] def _safe_which_no_cwd(command: str) -> Optional[str]: @@ -2053,26 +1964,21 @@ def _safe_which_no_cwd(command: str) -> Optional[str]: for bare names — unsafe for a dispatcher. Only explicit PATH entries are considered; empty / ``.`` entries are skipped. """ - path_env = os.environ.get("PATH", "") - for raw_dir in path_env.split(os.pathsep): + for raw_dir in os.environ.get("PATH", "").split(os.pathsep): if not raw_dir or raw_dir == ".": continue directory = os.path.expanduser(raw_dir) for name in _path_search_names(command): candidate = os.path.join(directory, name) - if not os.path.isfile(candidate): - continue - if _kb._IS_WINDOWS or os.access(candidate, os.X_OK): + if os.path.isfile(candidate) and (_kb._IS_WINDOWS or os.access(candidate, os.X_OK)): return candidate return None def _hermes_path_argv(path: str) -> list[str]: - """Return argv for a resolved Hermes executable path. - - Windows batch shims (`.cmd` / `.bat`) are unsafe as argv[0] because the - argument vector includes task-derived values; prefer the module form. - """ + """argv for a resolved Hermes executable path. Windows batch shims + (``.cmd``/``.bat``) are unsafe as argv[0] because the argument vector + includes task-derived values; prefer the module form.""" if _kb._IS_WINDOWS and _is_windows_batch_shim(path): return _module_hermes_argv() return [_absolute_hermes_path(path)] @@ -2195,12 +2101,60 @@ def _retag_legacy_worker_sessions(workspaces_root_path: str) -> None: _kb._log.debug("kanban worker: legacy session retag skipped (%s)", exc) -def _default_spawn( - task: Task, - workspace: str, - *, - board: Optional[str] = None, -) -> Optional[int]: +def _worker_argv(task: Task, profile_arg: str, hermes_home: Optional[str]) -> list[str]: + """Build the ``hermes -p --cli ... chat -q ...`` worker command.""" + cmd = [ + *_kb._resolve_hermes_argv(), + "-p", profile_arg, + # A worker must NEVER boot the interactive TUI: its no-TTY bail-out + # exits 0 without doing the task → "protocol violation" every attempt. + "--cli", + # Workers run under a profile-scoped HERMES_HOME and so see that + # profile's shell-hook allowlist; pass --accept-hooks explicitly so + # configured hooks still register. + "--accept-hooks", + ] + # One `--skills X` pair per name: easier to read in `ps` and avoids quoting + # ambiguity if a skill name contains unusual chars. + for sk in task.skills or (): + if sk: + cmd.extend(["--skills", sk]) + if task.model_override: + cmd.extend(["-m", task.model_override]) + # Pin the provider too so the worker resolves the model against the + # intended backend (model X with provider Y is the classic board-stall). + if task.provider_override: + cmd.extend(["--provider", task.provider_override]) + # Independent of the model override — a task can run the profile's own + # model at a different depth. + if task.reasoning_effort: + cmd.extend(["--reasoning", task.reasoning_effort]) + worker_toolsets = _resolve_worker_cli_toolsets(hermes_home) + if worker_toolsets: + cmd.extend(["--toolsets", ",".join(worker_toolsets)]) + cmd.extend(["chat", "-q", f"work kanban task {task.id}"]) + if task.goal_mode: + # The kanban goal-loop hook only runs in cli.py's fully-quiet branch. + # Without -Q the worker gets one turn, prints text, exits rc=0, and the + # dispatcher records a protocol violation. + cmd.append("-Q") + return cmd + + +def _open_worker_log(task: Task, board: Optional[str]): + """Append-mode per-task log (a re-run on unblock appends, never overwrites), + rotated first. Anchored at the board root (not the shared kanban root) so + `hermes kanban log` reads its own file and boards sharing task ids don't + collide.""" + log_dir = _kb.worker_logs_dir(board=board) + log_dir.mkdir(parents=True, exist_ok=True) + log_path = log_dir / f"{task.id}.log" + rotate_bytes, backup_count = worker_log_rotation_config() + _rotate_worker_log(log_path, rotate_bytes, backup_count) + return open(log_path, "ab") + + +def _default_spawn(task: Task, workspace: str, *, board: Optional[str] = None) -> Optional[int]: """Fire-and-forget ``hermes -p chat -q ...`` subprocess. Returns the child's PID so the dispatcher can detect crashes before the @@ -2209,15 +2163,13 @@ def _default_spawn( ``HERMES_KANBAN_DB`` / ``HERMES_KANBAN_BOARD`` / workspaces_root to the board the task was claimed from, so workers cannot see other boards. """ - import subprocess if not task.assignee: raise ValueError(f"task {task.id} has no assignee") - from hermes_cli.profiles import normalize_profile_name + from hermes_cli.profiles import normalize_profile_name, resolve_profile_env profile_arg = normalize_profile_name(task.assignee) - prompt = f"work kanban task {task.id}" env = dict(os.environ) # The dispatcher is detached from every conversation; its worker must never # inherit routing mirrored by a previous gateway turn. @@ -2229,7 +2181,6 @@ def _default_spawn( # without it the child's get_hermes_home() falls back to the DEFAULT # profile root because `hermes -p` applies its override before # hermes_constants is imported. - from hermes_cli.profiles import resolve_profile_env try: env["HERMES_HOME"] = resolve_profile_env(profile_arg) except FileNotFoundError: @@ -2261,18 +2212,10 @@ def _default_spawn( env["HERMES_KANBAN_GOAL_MODE"] = "1" if task.goal_max_turns is not None: env["HERMES_KANBAN_GOAL_MAX_TURNS"] = str(int(task.goal_max_turns)) - terminal_timeout = _worker_terminal_timeout_env( - task.max_runtime_seconds, - env.get("TERMINAL_TIMEOUT"), - ) - if terminal_timeout is not None: - env["TERMINAL_TIMEOUT"] = terminal_timeout - foreground_timeout = _worker_terminal_timeout_env( - task.max_runtime_seconds, - env.get("TERMINAL_MAX_FOREGROUND_TIMEOUT"), - ) - if foreground_timeout is not None: - env["TERMINAL_MAX_FOREGROUND_TIMEOUT"] = foreground_timeout + for var in ("TERMINAL_TIMEOUT", "TERMINAL_MAX_FOREGROUND_TIMEOUT"): + override = _worker_terminal_timeout_env(task.max_runtime_seconds, env.get(var)) + if override is not None: + env[var] = override # Pin the board DB + workspaces root so the worker's kanban paths still # match after `hermes -p` rewrites HERMES_HOME (symlink / Docker layouts). env["HERMES_KANBAN_DB"] = str(_kb.kanban_db_path(board=board)) @@ -2280,65 +2223,16 @@ def _default_spawn( _kb._retag_legacy_worker_sessions(env["HERMES_KANBAN_WORKSPACES_ROOT"]) # Board slug — defense-in-depth pin if a path is resolved without the # DB / workspaces env vars. - resolved_board = _kb._normalize_board_slug(board) or _kb.get_current_board() - env["HERMES_KANBAN_BOARD"] = resolved_board + env["HERMES_KANBAN_BOARD"] = _kb._normalize_board_slug(board) or _kb.get_current_board() # kanban_comment reads HERMES_PROFILE for its default author; `-p` alone # doesn't set the env var. env["HERMES_PROFILE"] = profile_arg - - # A worker must NEVER boot the interactive TUI: its no-TTY bail-out exits - # 0 without doing the task → "protocol violation" every attempt. `--cli` - # is the highest-precedence override; dropping HERMES_TUI covers older - # hermes builds on PATH that predate the flag's precedence. + # `--cli` is the highest-precedence TUI override; dropping HERMES_TUI covers + # older hermes builds on PATH that predate the flag's precedence. env.pop("HERMES_TUI", None) - cmd = [ - *_kb._resolve_hermes_argv(), - "-p", profile_arg, - "--cli", - # Workers run under a profile-scoped HERMES_HOME and so see that - # profile's shell-hook allowlist; pass --accept-hooks explicitly so - # configured hooks still register. - "--accept-hooks", - ] - # One `--skills X` pair per name: easier to read in `ps` and avoids quoting - # ambiguity if a skill name contains unusual chars. - if task.skills: - for sk in task.skills: - if sk: - cmd.extend(["--skills", sk]) - if task.model_override: - cmd.extend(["-m", task.model_override]) - # Pin the provider too so the worker resolves the model against the - # intended backend (model X with provider Y is the classic board-stall). - if task.provider_override: - cmd.extend(["--provider", task.provider_override]) - # Independent of the model override — a task can run the profile's own - # model at a different depth. - if task.reasoning_effort: - cmd.extend(["--reasoning", task.reasoning_effort]) - worker_toolsets = _resolve_worker_cli_toolsets(env.get("HERMES_HOME")) - if worker_toolsets: - cmd.extend(["--toolsets", ",".join(worker_toolsets)]) - cmd.extend([ - "chat", - "-q", prompt, - ]) - if task.goal_mode: - # The kanban goal-loop hook only runs in cli.py's fully-quiet branch. - # Without -Q the worker gets one turn, prints text, exits rc=0, and the - # dispatcher records a protocol violation. - cmd.append("-Q") - # Per-task log anchored at the board root (not the shared kanban root) so - # `hermes kanban log` reads its own file and boards sharing task ids don't collide. - log_dir = _kb.worker_logs_dir(board=board) - log_dir.mkdir(parents=True, exist_ok=True) - log_path = log_dir / f"{task.id}.log" - rotate_bytes, backup_count = worker_log_rotation_config() - _rotate_worker_log(log_path, rotate_bytes, backup_count) - - # Use 'a' so a re-run on unblock appends rather than overwrites. - log_f = open(log_path, "ab") + cmd = _worker_argv(task, profile_arg, env.get("HERMES_HOME")) + log_f = _open_worker_log(task, board) try: proc = subprocess.Popen( # noqa: S603 -- argv is a fixed list built above cmd, @@ -2381,7 +2275,6 @@ def run_daemon( the gateway dispatcher and ``hermes kanban dispatch`` — the standalone daemon must not be the one uncapped entry point. """ - import signal import threading if stop_event is None: @@ -2403,9 +2296,7 @@ def run_daemon( try: # Re-resolved every tick (config load is mtime-cached) so operator # edits apply without a restart. - max_in_progress = resolve_max_in_progress( - _kb.configured_max_in_progress() - ) + max_in_progress = resolve_max_in_progress(_kb.configured_max_in_progress()) with contextlib.closing(_kb.connect()) as conn: res = _kb.dispatch_once( conn, @@ -2423,6 +2314,6 @@ def run_daemon( stop_event.wait(timeout=interval) -# Late-bound origin namespace (see module docstring). Imported LAST so this +# Late-bound origin namespace (see module docstring); imported LAST so this # module is fully populated before ``kanban_db`` re-exports from it. from hermes_cli import kanban_db as _kb # noqa: E402 diff --git a/hermes_cli/kanban_db_notify.py b/hermes_cli/kanban_db_notify.py index 077fad8e0e..5365ebb5d0 100644 --- a/hermes_cli/kanban_db_notify.py +++ b/hermes_cli/kanban_db_notify.py @@ -36,9 +36,7 @@ def _sub_key(task_id: str, platform: str, chat_id: str, thread_id: Optional[str] return (task_id, platform, chat_id, thread_id or "") -def _encode_notify_delivery_metadata( - metadata: Optional[Mapping[str, Any]], -) -> Optional[str]: +def _encode_notify_delivery_metadata(metadata: Optional[Mapping[str, Any]]) -> Optional[str]: """Serialize platform send metadata stored on notification subscriptions.""" if not isinstance(metadata, Mapping): return None @@ -267,11 +265,7 @@ def remove_notify_sub( return cur.rowcount > 0 -def purge_stale_done_notify_subs( - conn: sqlite3.Connection, - *, - max_age_days: int = 30, -) -> int: +def purge_stale_done_notify_subs(conn: sqlite3.Connection, *, max_age_days: int = 30) -> int: """Delete notify subs whose task has sat in ``done``/``blocked`` untouched for longer than ``max_age_days``. diff --git a/hermes_cli/kanban_db_workspace.py b/hermes_cli/kanban_db_workspace.py index 29b80eb5f9..2511c2d1e2 100644 --- a/hermes_cli/kanban_db_workspace.py +++ b/hermes_cli/kanban_db_workspace.py @@ -7,7 +7,6 @@ origin-resident helpers are reached late-bound via ``_kb`` so monkeypatching from __future__ import annotations -import logging import os import shutil import sqlite3 @@ -20,9 +19,6 @@ import contextlib if TYPE_CHECKING: from hermes_cli.kanban_db import Task -# Log-record parity with the origin module. -_log = logging.getLogger("hermes_cli.kanban_db") - _REMOVABLE_KINDS = ("scratch", "worktree") # Statuses after which a child no longer needs its parent's workspace artifacts. @@ -237,9 +233,7 @@ def _try_cleanup_parent_workspaces(conn: sqlite3.Connection, task_id: str) -> No ): continue if row["workspace_kind"] == "worktree": - _cleanup_worktree_workspace( - parent_id, row["workspace_path"], row["branch_name"] - ) + _cleanup_worktree_workspace(parent_id, row["workspace_path"], row["branch_name"]) continue wp = Path(row["workspace_path"]) if wp.is_dir() and _is_managed_scratch_path(wp): @@ -264,10 +258,7 @@ def _cleanup_worker_tmux(conn: sqlite3.Connection, task_id: str) -> None: capture_output=True, text=True, encoding='utf-8', errors='replace', timeout=5, ) if out.stdout.strip() == "1": - subprocess.run( - ["tmux", "kill-session", "-t", session], - capture_output=True, timeout=5, - ) + subprocess.run(["tmux", "kill-session", "-t", session], capture_output=True, timeout=5) _kb._log.debug("Killed stale tmux session: %s", session) except Exception: pass # best-effort — never block completion @@ -422,9 +413,7 @@ def _anchored_worktree(repo_root: Path, task_id: str, branch_name: str) -> tuple return target, branch_name -def _resolve_worktree_workspace( - task: Task, *, board: Optional[str] = None -) -> tuple[Path, str]: +def _resolve_worktree_workspace(task: Task, *, board: Optional[str] = None) -> tuple[Path, str]: """Resolve + materialize a linked git worktree for ``task``. With no ``task.workspace_path`` the anchor is the board's ``default_workdir`` so every worktree lands under a board-owned repo (``/.worktrees/``) @@ -529,9 +518,7 @@ def resolve_workspace(task: Task, *, board: Optional[str] = None) -> Path: ) elif kind == "dir": if not task.workspace_path: - raise ValueError( - f"task {task.id} has workspace_kind=dir but no workspace_path" - ) + raise ValueError(f"task {task.id} has workspace_kind=dir but no workspace_path") p = Path(task.workspace_path).expanduser() if not p.is_absolute(): raise ValueError( @@ -550,15 +537,11 @@ def _set_task_column(conn: sqlite3.Connection, task_id: str, column: str, value: conn.execute(f"UPDATE tasks SET {column} = ? WHERE id = ?", (value, task_id)) -def set_workspace_path( - conn: sqlite3.Connection, task_id: str, path: Path | str -) -> None: +def set_workspace_path(conn: sqlite3.Connection, task_id: str, path: Path | str) -> None: _set_task_column(conn, task_id, "workspace_path", str(path)) -def set_branch_name( - conn: sqlite3.Connection, task_id: str, branch_name: str -) -> None: +def set_branch_name(conn: sqlite3.Connection, task_id: str, branch_name: str) -> None: _set_task_column(conn, task_id, "branch_name", str(branch_name))