atlas-iac/services/hermes/scripts/kanban_supervisor.py
2026-08-21 05:55:46 +00:00

456 lines
17 KiB
Python

#!/usr/bin/env python3
"""In-pod autonomous supervisor that sustains multi-hour unattended work.
The deployed Hermes agent already decomposes objectives and iterates a within-
card goal loop, but nothing in-pod spawns the *cross-card* review->repair->
re-review follow-ups after a card finishes; today that chain only exists as an
external Claude Code "codex-shepherd" session on an operator's workstation, so
unattended runs stall once the first implementation card completes. This
supervisor closes that gap from inside the pod.
It runs a bounded poll loop that reads board state through ``hermes_cli.kanban_db``
and, for each terminal card, applies exactly one follow-up:
* a done implementation with a produced PR/commit and no existing review ->
create a review card (routed to a review assignee via the normal cli lane);
* a review card that returned SHIP -> mark the implementation ready for a human
to merge (a comment/flag; the supervisor never merges, approves, or clears WIP);
* a review card that returned BLOCK -> create a bounded repair card;
* a repair card that produced a new commit -> create the re-review card.
All decision logic lives in :mod:`supervisor_policy` (pure, unit tested). This
module is only the I/O shell: config, board iteration, and turning a
:class:`supervisor_policy.Decision` into ``create_task``/``block_task``/
``add_comment`` calls. It imports no provider client and holds no metered path;
the sole spawn is a Kanban card that routes through the existing subscription
lanes. It is inert until ``kanban.auto_supervise`` is set true in the deployed
config, re-read every tick like ``kanban.auto_decompose``.
"""
from __future__ import annotations
import json
import os
import sys
import time
from dataclasses import dataclass
from pathlib import Path
from typing import Any, Callable
import yaml
import supervisor_policy as policy
from cli_lane_config import (
DEFAULT_DIRECT_HEARTBEAT_TIMEOUT_SECONDS,
DIRECT_CLAIM_LOCK,
EXTERNAL_PREFIX,
)
from cli_lane_recovery import _has_pending_finalization
DATA_ROOT = Path(os.environ.get("HERMES_HOME", "/opt/data"))
CONFIG_PATH = DATA_ROOT / "config.yaml"
# Persistent record of terminal (ship / escalate) emissions so each fires at
# most once across the stateless poll loop, and survives a pod restart.
LEDGER_PATH = DATA_ROOT / "supervisor" / "emitted.json"
DEFAULT_INTERVAL_SECONDS = 30
DEFAULT_MAX_CYCLES = 5
DEFAULT_MAX_CHAINS = 20
DEFAULT_REVIEW_ASSIGNEE = "cli-claude-xhigh"
DEFAULT_REPAIR_ASSIGNEE = "cli-auto"
@dataclass(frozen=True)
class Settings:
"""Per-tick configuration, re-read from the deployed config each pass."""
enabled: bool
interval: int
limits: policy.Limits
direct_lane_heartbeat_timeout_seconds: int = (
DEFAULT_DIRECT_HEARTBEAT_TIMEOUT_SECONDS
)
def _kanban_config() -> dict[str, Any]:
try:
document = yaml.safe_load(CONFIG_PATH.read_text(encoding="utf-8"))
except (OSError, yaml.YAMLError):
return {}
kanban = document.get("kanban") if isinstance(document, dict) else None
return kanban if isinstance(kanban, dict) else {}
def _bool(value: Any, default: bool) -> bool:
return value if isinstance(value, bool) else default
def _positive_int(value: Any, default: int) -> int:
if isinstance(value, bool) or not isinstance(value, (int, float)):
return default
parsed = int(value)
return parsed if parsed > 0 else default
def _text(value: Any, default: str) -> str:
return value if isinstance(value, str) and value.strip() else default
def load_settings() -> Settings:
"""Read the ``kanban`` config block; default to inert and conservative."""
cfg = _kanban_config()
return Settings(
enabled=_bool(cfg.get("auto_supervise"), False),
interval=_positive_int(
cfg.get("supervise_interval_seconds"), DEFAULT_INTERVAL_SECONDS
),
limits=policy.Limits(
max_cycles=_positive_int(cfg.get("supervise_max_cycles"), DEFAULT_MAX_CYCLES),
max_chains=_positive_int(cfg.get("supervise_max_chains"), DEFAULT_MAX_CHAINS),
review_assignee=_text(
cfg.get("supervise_review_assignee"), DEFAULT_REVIEW_ASSIGNEE
),
repair_assignee=_text(
cfg.get("supervise_repair_assignee"), DEFAULT_REPAIR_ASSIGNEE
),
),
direct_lane_heartbeat_timeout_seconds=_positive_int(
cfg.get("direct_lane_heartbeat_timeout_seconds"),
DEFAULT_DIRECT_HEARTBEAT_TIMEOUT_SECONDS,
),
)
def _board_slug(board: Any) -> str:
if isinstance(board, dict):
return str(board.get("slug") or board.get("id") or "")
return str(getattr(board, "slug", None) or getattr(board, "id", None) or board or "")
def _log(message: str) -> None:
print(f"kanban-supervisor: {message}", file=sys.stderr, flush=True)
class Ledger:
"""Durable set of emission keys the supervisor itself has already fired.
The poll loop is stateless and re-plans every terminal card each tick, so a
SHIP or a cycle-limit escalation would otherwise re-comment its target on
every pass. This ledger is the authoritative read-before-write guard: an
emission is performed only when its key is absent, then recorded, so each
fires at most once regardless of whether the underlying ``kanban_db`` moves
the source card out of its done state.
"""
def __init__(self, path: Path):
self._path = path
self._keys = self._load(path)
@staticmethod
def _load(path: Path) -> set[str]:
try:
data = json.loads(path.read_text(encoding="utf-8"))
except (OSError, ValueError):
return set()
return {str(key) for key in data} if isinstance(data, list) else set()
def has(self, key: str) -> bool:
return key in self._keys
def record(self, key: str) -> None:
self._keys.add(key)
try:
self._path.parent.mkdir(parents=True, exist_ok=True)
tmp = self._path.with_suffix(".tmp")
tmp.write_text(json.dumps(sorted(self._keys)), encoding="utf-8")
os.replace(tmp, self._path)
except OSError as error:
_log(f"could not persist emission ledger: {error}")
def _iter_boards(kanban_db: Any) -> list[str]:
try:
boards = kanban_db.list_boards(include_archived=False)
except Exception as error: # noqa: BLE001 - one bad registry never stops the tick
_log(f"could not list boards: {error}")
return []
slugs = [_board_slug(raw) for raw in boards]
return [slug for slug in slugs if slug]
def _comment(kanban_db: Any, conn: Any, task_id: str, body: str) -> None:
try:
kanban_db.add_comment(conn, task_id, policy.SUPERVISOR_AUTHOR, body)
except Exception as error: # noqa: BLE001 - a comment is best-effort signalling
_log(f"could not comment on {task_id}: {error}")
def _value(record: Any, name: str, default: Any = None) -> Any:
if isinstance(record, dict):
return record.get(name, default)
return getattr(record, name, default)
def _watchdog_board(
kanban_db: Any,
board: str,
heartbeat_timeout_seconds: int,
now: int,
) -> int:
"""Atomically reclaim heartbeat-dead direct-lane runs on one board."""
reclaimed = 0
cutoff = now - heartbeat_timeout_seconds
with kanban_db.scoped_current_board(board):
conn = kanban_db.connect(board=board)
try:
tasks = list(kanban_db.list_tasks(conn))
for task in tasks:
task_id = str(_value(task, "id", "") or "")
run_id = _value(task, "current_run_id")
if (
not task_id
or _value(task, "status") != "running"
or _value(task, "claim_lock") != DIRECT_CLAIM_LOCK
or not str(_value(task, "assignee", "") or "").startswith(
EXTERNAL_PREFIX
)
or type(run_id) is not int
):
continue
try:
run = kanban_db.get_run(conn, run_id)
if (
run is None
or _value(run, "id") != run_id
or _value(run, "task_id") != task_id
or _value(run, "status") != "running"
or _value(run, "ended_at") is not None
):
continue
heartbeat_at = _value(run, "last_heartbeat_at")
liveness_at = (
int(heartbeat_at)
if heartbeat_at is not None
else int(_value(run, "started_at"))
)
if liveness_at >= cutoff:
continue
# Accepted terminal evidence outranks liveness recovery.
# The lane finalizer will either commit it or classify it
# before this watchdog is allowed to release the run.
if _has_pending_finalization(board, task_id, run_id):
continue
age = max(0, now - liveness_at)
reason = (
f"direct CLI heartbeat absent for {age}s "
f"(limit {heartbeat_timeout_seconds}s)"
)
won = kanban_db.reclaim_task(
conn,
task_id,
reason=reason,
expected_run_id=run_id,
expected_run_liveness_before=cutoff,
)
if not won:
continue
_comment(
kanban_db,
conn,
task_id,
(
f"Watchdog reclaimed heartbeat-dead direct CLI run "
f"{run_id}: {reason}. The claim was released and the "
"task requeued without counting a worker failure."
),
)
reclaimed += 1
except Exception as error: # noqa: BLE001 - isolate each run
_log(f"watchdog could not inspect {board}/{task_id}: {error}")
finally:
conn.close()
return reclaimed
def watchdog_once(
kanban_db: Any,
heartbeat_timeout_seconds: int = DEFAULT_DIRECT_HEARTBEAT_TIMEOUT_SECONDS,
) -> int:
"""One independent heartbeat-death pass over every active board."""
total = 0
now = int(time.time())
for board in _iter_boards(kanban_db):
try:
total += _watchdog_board(
kanban_db, board, heartbeat_timeout_seconds, now
)
except Exception as error: # noqa: BLE001 - isolate a faulty board
_log(f"watchdog temporarily skipping board {board!r}: {error}")
return total
def _mark_ready_for_human(kanban_db: Any, conn: Any, task_id: str, body: str) -> None:
"""Flag the implementation as human-mergeable. Never merges or clears WIP."""
setter = getattr(kanban_db, "set_task_metadata", None) or getattr(
kanban_db, "update_task_metadata", None
)
if callable(setter):
try:
setter(conn, task_id, {"supervisor_ready_for_human_merge": True})
except Exception as error: # noqa: BLE001 - the comment is the durable flag
_log(f"could not set ready flag on {task_id}: {error}")
_comment(kanban_db, conn, task_id, body)
def _escalate(kanban_db: Any, conn: Any, task_id: str, reason: str) -> None:
body = (
f"supervisor fail-closed escalation: {reason}. Human attention required; "
"no follow-up card was auto-created."
)
try:
kanban_db.block_task(conn, task_id, reason=body, kind="supervisor")
except Exception as error: # noqa: BLE001 - still record the human-visible comment
_log(f"could not block {task_id}: {error}")
_comment(kanban_db, conn, task_id, body)
def _already_created(kanban_db: Any, conn: Any, payload: dict[str, Any]) -> bool:
"""True if a card matching this spawn payload's stamp already exists."""
supervisor_stamp = (payload.get("metadata") or {}).get("supervisor") or {}
kind = supervisor_stamp.get("kind")
root = str(supervisor_stamp.get("root") or "")
commit = str(supervisor_stamp.get("head_commit") or "")
if not kind or not root:
return False
try:
tasks = list(kanban_db.list_tasks(conn))
except Exception as error: # noqa: BLE001 - if we cannot confirm, do not skip
_log(f"could not re-scan before spawn retry: {error}")
return False
return policy.existing_followup(tasks, kind, root, commit)
def _spawn(kanban_db: Any, conn: Any, decision: policy.Decision) -> None:
payload = dict(decision.payload or {})
try:
kanban_db.create_task(conn, **payload)
except TypeError:
# create_task may reject idempotency_key on older runtimes. That is
# raised at call binding, before any insert, but a post-insert TypeError
# is also possible, so re-run the dedup scan first: only retry the create
# when no matching card exists, so a partial insert is never doubled.
if not _already_created(kanban_db, conn, payload):
payload.pop("idempotency_key", None)
kanban_db.create_task(conn, **payload)
_comment(kanban_db, conn, decision.target_id, f"supervisor: {decision.reason}")
def apply_decision(
kanban_db: Any, conn: Any, decision: policy.Decision, ledger: Ledger
) -> bool:
"""Execute one decision. Returns True when an action was taken.
Spawns are self-limiting across ticks (the created card fails the next dedup
scan). SHIP and escalate targets stay in their done state, so they are gated
on the persistent ledger and fire at most once per source card.
"""
if decision.action == "none":
return False
policy.assert_safe(decision)
if decision.action == "spawn":
_spawn(kanban_db, conn, decision)
return True
key = f"{decision.action}:{decision.target_id}"
if ledger.has(key):
return False # already emitted on a prior tick; never re-emit or re-block
if decision.action == "ship":
info = decision.payload or {}
commit = info.get("commit", "")
pr = info.get("pr", "") or "branch on record"
body = (
f"supervisor: review returned SHIP for commit {commit} ({pr}). "
"READY FOR HUMAN MERGE - a human must merge; the supervisor never "
"merges, approves, closes, or clears WIP."
)
_mark_ready_for_human(kanban_db, conn, decision.target_id, body)
else: # escalate
_escalate(kanban_db, conn, decision.target_id, decision.reason)
ledger.record(key)
return True
def supervise_board(
kanban_db: Any, board: str, limits: policy.Limits, ledger: Ledger
) -> int:
actions = 0
with kanban_db.scoped_current_board(board):
conn = kanban_db.connect(board=board)
try:
tasks = list(kanban_db.list_tasks(conn))
for task in tasks:
decision = policy.plan(task, tasks, limits)
try:
if apply_decision(kanban_db, conn, decision, ledger):
actions += 1
except Exception as error: # noqa: BLE001 - one card never stops others
_log(f"action failed on board {board!r}: {error}")
finally:
conn.close()
return actions
def supervise_once(
kanban_db: Any, limits: policy.Limits, ledger: Ledger | None = None
) -> int:
"""One bounded pass over every non-archived board."""
if ledger is None:
ledger = Ledger(LEDGER_PATH)
total = 0
for board in _iter_boards(kanban_db):
try:
total += supervise_board(kanban_db, board, limits, ledger)
except Exception as error: # noqa: BLE001 - isolate a faulty board
_log(f"temporarily skipping board {board!r}: {error}")
return total
def run_forever(
kanban_db: Any,
*,
sleep: Callable[[float], None] = time.sleep,
load: Callable[[], Settings] = load_settings,
max_ticks: int | None = None,
) -> int:
"""Poll loop. Inert while ``auto_supervise`` is false; re-reads config a tick."""
ticks = 0
while max_ticks is None or ticks < max_ticks:
settings = load()
try:
watchdog_once(
kanban_db,
heartbeat_timeout_seconds=(
settings.direct_lane_heartbeat_timeout_seconds
),
)
except Exception as error: # noqa: BLE001 - never stop the safety loop
_log(f"watchdog tick failed: {error}")
if settings.enabled:
try:
supervise_once(kanban_db, settings.limits)
except Exception as error: # noqa: BLE001 - never let one tick kill the loop
_log(f"tick failed: {error}")
sleep(settings.interval)
ticks += 1
return ticks
def main(argv: list[str] | None = None) -> int:
from hermes_cli import kanban_db
run_forever(kanban_db)
return 0
if __name__ == "__main__": # pragma: no cover - process entry point
sys.exit(main(sys.argv[1:]))