From 37b160d04eb85bd401606e6edcb68efeb0bb25ba Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 16:32:06 -0700 Subject: [PATCH] refactor(kanban-cli): extract dispatcher/maintenance verbs to kanban_ops; unify single-mutation ok/err reporting --- hermes_cli/kanban.py | 465 +++------------------------------------ hermes_cli/kanban_ops.py | 426 +++++++++++++++++++++++++++++++++++ 2 files changed, 451 insertions(+), 440 deletions(-) create mode 100644 hermes_cli/kanban_ops.py diff --git a/hermes_cli/kanban.py b/hermes_cli/kanban.py index 75ed34ded3..aa25b3a790 100644 --- a/hermes_cli/kanban.py +++ b/hermes_cli/kanban.py @@ -25,6 +25,9 @@ from hermes_cli.kanban_output import ( _task_to_dict, ) from hermes_cli.kanban_boards import _dispatch_boards +from hermes_cli.kanban_ops import ( + _cmd_daemon, _kanban_config, _cmd_dispatch, _cmd_gc, _cmd_repair, _cmd_tail, _cmd_watch, +) from hermes_cli.kanban_parser import build_parser # noqa: F401 (re-exported: hermes_cli.main, run_slash) @@ -234,16 +237,6 @@ def kanban_command(args: argparse.Namespace) -> int: # Handlers # --------------------------------------------------------------------------- -def _kanban_config() -> dict: - """``config.yaml`` ``kanban:`` section, or ``{}`` when config can't be loaded.""" - try: - from hermes_cli.config import load_config - cfg = load_config() - return (cfg.get("kanban", {}) if isinstance(cfg, dict) else {}) or {} - except Exception: - return {} - - def _profile_author() -> str: """Best-effort author name for an interactive CLI call.""" for env in ("HERMES_PROFILE_NAME", "HERMES_PROFILE"): @@ -300,6 +293,14 @@ def _stripped_or_none(value: Optional[str]) -> Optional[str]: return None if value is None else (value.strip() or None) +def _ok_or_err(ok, fail: str, done: str) -> int: + """Single-mutation handlers: print ``done`` (rc 0) or ``fail`` to stderr (rc 1).""" + if not ok: + return _err(fail) + print(done) + return 0 + + def _bulk_ids(args: argparse.Namespace) -> list[str]: """Positional ``task_id`` plus ``--ids`` extras (bulk verbs).""" return [args.task_id] + list(getattr(args, "ids", None) or []) @@ -374,10 +375,8 @@ def _cmd_heartbeat(args: argparse.Namespace) -> int: note=getattr(args, "note", None), expected_run_id=_worker_run_id_for(args.task_id), ) - if not ok: - return _err(f"cannot heartbeat {args.task_id} (not running?)") - print(f"Heartbeat recorded for {args.task_id}") - return 0 + return _ok_or_err(ok, f"cannot heartbeat {args.task_id} (not running?)", + f"Heartbeat recorded for {args.task_id}") def _cmd_assignees(args: argparse.Namespace) -> int: @@ -663,10 +662,8 @@ def _cmd_assign(args: argparse.Namespace) -> int: profile = _none_profile(args.profile) with kb.connect_closing() as conn: ok = kb.assign_task(conn, args.task_id, profile) - if not ok: - return _err(f"no such task: {args.task_id}") - print(f"Assigned {args.task_id} to {profile or '(unassigned)'}") - return 0 + return _ok_or_err(ok, f"no such task: {args.task_id}", + f"Assigned {args.task_id} to {profile or '(unassigned)'}") def _cmd_set_model(args: argparse.Namespace) -> int: @@ -697,10 +694,8 @@ def _cmd_reclaim(args: argparse.Namespace) -> int: conn, args.task_id, reason=getattr(args, "reason", None), ) - if not ok: - return _err(f"cannot reclaim {args.task_id} (not running or unknown id)") - print(f"Reclaimed {args.task_id}") - return 0 + return _ok_or_err(ok, f"cannot reclaim {args.task_id} (not running or unknown id)", + f"Reclaimed {args.task_id}") def _cmd_reassign(args: argparse.Namespace) -> int: @@ -711,17 +706,12 @@ def _cmd_reassign(args: argparse.Namespace) -> int: reclaim_first=bool(getattr(args, "reclaim", False)), reason=getattr(args, "reason", None), ) - if not ok: - return _err( - f"cannot reassign {args.task_id} " - f"(unknown id, or still running — pass --reclaim to release first)" - ) - print( - f"Reassigned {args.task_id} to " - f"{profile or '(unassigned)'}" - + (" (claim reclaimed)" if getattr(args, "reclaim", False) else "") + return _ok_or_err( + ok, + f"cannot reassign {args.task_id} (unknown id, or still running — pass --reclaim to release first)", + f"Reassigned {args.task_id} to {profile or '(unassigned)'}" + + (" (claim reclaimed)" if getattr(args, "reclaim", False) else ""), ) - return 0 def _rows_by_task(conn, table: str, ids: list[str]) -> dict[str, list]: @@ -841,10 +831,8 @@ def _cmd_link(args: argparse.Namespace) -> int: def _cmd_unlink(args: argparse.Namespace) -> int: with kb.connect_closing() as conn: ok = kb.unlink_tasks(conn, args.parent_id, args.child_id) - if not ok: - return _err(f"No such link: {args.parent_id} -> {args.child_id}") - print(f"Unlinked {args.parent_id} -> {args.child_id}") - return 0 + return _ok_or_err(ok, f"No such link: {args.parent_id} -> {args.child_id}", + f"Unlinked {args.parent_id} -> {args.child_id}") def _cmd_claim(args: argparse.Namespace) -> int: @@ -1267,292 +1255,6 @@ def _cmd_archive(args: argparse.Namespace) -> int: ) -def _cmd_tail(args: argparse.Namespace) -> int: - last_id = 0 - print(f"Tailing events for {args.task_id}. Ctrl-C to stop.") - try: - while True: - with kb.connect_closing() as conn: - events = kb.list_events(conn, args.task_id) - for e in events: - if e.id > last_id: - pl = f" {e.payload}" if e.payload else "" - print(f"[{_fmt_ts(e.created_at)}] {e.kind}{pl}", flush=True) - last_id = e.id - time.sleep(max(0.1, args.interval)) - except KeyboardInterrupt: - print("\n(stopped)") - return 0 - - -def _coerce_positive_int(value): - if value is None: - return None - try: - ival = int(value) - except (TypeError, ValueError): - return None - return ival if ival >= 1 else None - - -def _cmd_dispatch(args: argparse.Namespace) -> int: - # Honour kanban.default_assignee, kanban.max_in_progress, - # kanban.max_in_progress_per_profile and kanban.max_spawn with the same - # semantics as the gateway dispatch path. - try: - from hermes_cli.config import load_config - _cfg = load_config() - _kanban_cfg = _cfg.get("kanban", {}) if isinstance(_cfg, dict) else {} - default_assignee = (_kanban_cfg.get("default_assignee") or "").strip() or None - max_in_progress_per_profile = _coerce_positive_int( - _kanban_cfg.get("max_in_progress_per_profile") - ) - # Memory-derived default when unset — same fallback the gateway applies. - max_in_progress = kb.resolve_max_in_progress( - _coerce_positive_int(_kanban_cfg.get("max_in_progress")) - ) - # CLI --max is the more explicit signal, so it wins over kanban.max_spawn. - cli_max = getattr(args, "max", None) - max_spawn = cli_max if cli_max is not None else _coerce_positive_int( - _kanban_cfg.get("max_spawn") - ) - except Exception: - default_assignee = None - max_in_progress_per_profile = None - max_in_progress = None - max_spawn = getattr(args, "max", None) - with kb.connect_closing() as conn: - res = kb.dispatch_once( - conn, - dry_run=args.dry_run, - max_spawn=max_spawn, - max_in_progress=max_in_progress, - failure_limit=getattr(args, "failure_limit", kb.DEFAULT_SPAWN_FAILURE_LIMIT), - default_assignee=default_assignee, - max_in_progress_per_profile=max_in_progress_per_profile, - ) - if getattr(args, "json", False): - _print_json({ - "reclaimed": res.reclaimed, - "crashed": res.crashed, - "timed_out": res.timed_out, - "stale": res.stale, - "auto_blocked": res.auto_blocked, - "promoted": res.promoted, - "spawned": [ - {"task_id": tid, "assignee": who, "workspace": ws} - for (tid, who, ws) in res.spawned - ], - "skipped_unassigned": res.skipped_unassigned, - "skipped_nonspawnable": res.skipped_nonspawnable, - "skipped_per_profile_capped": [ - {"task_id": tid, "assignee": who, "current": current} - for (tid, who, current) in res.skipped_per_profile_capped - ], - "auto_assigned_default": res.auto_assigned_default, - }, ascii=True) - return 0 - print(f"Reclaimed: {res.reclaimed}") - for label, items in ( - ("Crashed: ", res.crashed), - ("Timed out: ", res.timed_out), - ("Stale: ", res.stale), - ("Auto-blocked:", res.auto_blocked), - ): - print(f"{label} {len(items)}") - if items: - print(f" {', '.join(items)}") - print(f"Promoted: {res.promoted}") - print(f"Spawned: {len(res.spawned)}") - tag = " (dry)" if args.dry_run else "" - for tid, who, ws in res.spawned: - print(f" - {tid} -> {who} @ {ws or '-'}{tag}") - if res.auto_assigned_default: - print( - f"Auto-assigned to kanban.default_assignee={default_assignee!r}: " - f"{', '.join(res.auto_assigned_default)}" - ) - if res.skipped_unassigned: - print(f"Skipped (unassigned): {', '.join(res.skipped_unassigned)}") - for tid, who, current in res.skipped_per_profile_capped: - print(f"Deferred ({who} at per-profile cap, {current} running): {tid}") - if res.skipped_nonspawnable: - print( - f"Skipped (non-spawnable assignee — terminal lane, OK): " - f"{', '.join(res.skipped_nonspawnable)}" - ) - return 0 - - -def _cmd_daemon(args: argparse.Namespace) -> int: - """Deprecated — the dispatcher now runs inside the gateway. - - Kept as a stub so old scripts/systemd units get a clear migration message. - ``--force`` (hidden from --help) keeps the standalone loop for hosts that - truly cannot run the gateway; the default path exits 2 so nobody - accidentally runs two dispatchers against the same kanban.db. - """ - if not getattr(args, "force", False): - return _err( - "hermes kanban daemon: DEPRECATED — the dispatcher now runs\n" - "inside the gateway. To use kanban:\n" - "\n" - " hermes gateway start # starts the gateway + embedded dispatcher\n" - "\n" - "Ready tasks will be picked up on the next dispatcher tick\n" - "(default: every 60 seconds). Configure via config.yaml:\n" - "\n" - " kanban:\n" - " dispatch_in_gateway: true # default\n" - " dispatch_interval_seconds: 60\n" - " failure_limit: 2 # consecutive non-success attempts before auto-block\n" - "\n" - "Running both the gateway AND this standalone daemon will\n" - "race for claims. If you truly need the old standalone\n" - "daemon (no gateway available), rerun with --force.", - 2, - ) - - # Init before printing "started" so the DB path is right and init errors - # surface immediately. - kb.init_db() - - pidfile = getattr(args, "pidfile", None) - if pidfile: - try: - Path(pidfile).parent.mkdir(parents=True, exist_ok=True) - Path(pidfile).write_text(str(os.getpid()), encoding="utf-8") - except OSError as exc: - print(f"warning: could not write pidfile {pidfile}: {exc}", file=sys.stderr) - - verbose = bool(getattr(args, "verbose", False)) - print( - f"Kanban dispatcher running STANDALONE via --force " - f"(interval={args.interval}s, pid={os.getpid()}). " - f"Ctrl-C to stop. NOTE: if a gateway is also running with " - f"dispatch_in_gateway=true (default), you have two dispatchers " - f"racing for claims.", - file=sys.stderr, - ) - - # Health telemetry: warn when every tick finds ready work but spawns - # nothing (broken profile, PATH drift, missing venv, credential loss) — - # the per-task breaker auto-blocks quietly, so the operator needs a signal. - HEALTH_WINDOW = 6 # ticks (default 30s at interval=5) - health_state = {"bad_ticks": 0, "last_warn_at": 0} - - def _on_tick(res): - ready_pending = bool(res.skipped_unassigned) or _ready_queue_nonempty() - spawned_any = bool(res.spawned) - if ready_pending and not spawned_any: - health_state["bad_ticks"] += 1 - else: - health_state["bad_ticks"] = 0 - # Warn once per HEALTH_WINDOW bad ticks, at most every 5 minutes. - if health_state["bad_ticks"] >= HEALTH_WINDOW: - now = int(time.time()) - if now - health_state["last_warn_at"] >= 300: - print( - f"[{_fmt_ts(now)}] WARN dispatcher stuck: " - f"ready queue non-empty for {health_state['bad_ticks']} " - f"consecutive ticks but 0 workers spawned successfully. " - f"Check profile health (venv, PATH, credentials) and " - f"`hermes kanban list --status ready` / " - f"`hermes kanban list --status blocked` for recent " - f"spawn_failed tasks.", - file=sys.stderr, flush=True, - ) - health_state["last_warn_at"] = now - if not verbose: - return - did_work = ( - res.reclaimed or res.crashed or res.timed_out or res.promoted - or res.spawned or res.auto_blocked or res.stale - ) - if did_work: - print( - f"[{_fmt_ts(int(time.time()))}] " - f"reclaimed={res.reclaimed} crashed={len(res.crashed)} " - f"timed_out={len(res.timed_out)} stale={len(res.stale)} " - f"promoted={res.promoted} spawned={len(res.spawned)} " - f"auto_blocked={len(res.auto_blocked)}", - flush=True, - ) - - def _ready_queue_nonempty() -> bool: - """Is there a ready+assigned+unclaimed task the dispatcher would spawn for? - Control-plane lanes pulled via ``claim_task`` are correctly idle, not stuck.""" - try: - with kb.connect_closing() as conn: - return kb.has_spawnable_ready(conn) - except Exception: - return False - - try: - kb.run_daemon( - interval=args.interval, - max_spawn=args.max, - failure_limit=getattr(args, "failure_limit", kb.DEFAULT_SPAWN_FAILURE_LIMIT), - on_tick=_on_tick, - ) - finally: - if pidfile: - try: - Path(pidfile).unlink() - except OSError: - pass - print("(dispatcher stopped)") - return 0 - - -def _cmd_watch(args: argparse.Namespace) -> int: - """Live-stream task_events to the terminal.""" - kinds = ( - {k.strip() for k in args.kinds.split(",") if k.strip()} - if args.kinds else None - ) - print("Watching kanban events. Ctrl-C to stop.", flush=True) - # Seed cursor at the latest id so we don't replay history. - with kb.connect_closing() as conn: - row = conn.execute( - "SELECT COALESCE(MAX(id), 0) AS m FROM task_events" - ).fetchone() - cursor = int(row["m"]) - - try: - while True: - with kb.connect_closing() as conn: - rows = conn.execute( - "SELECT e.id, e.task_id, e.kind, e.payload, e.created_at, " - " t.assignee, t.tenant " - "FROM task_events e LEFT JOIN tasks t ON t.id = e.task_id " - "WHERE e.id > ? ORDER BY e.id ASC LIMIT 200", - (cursor,), - ).fetchall() - for r in rows: - cursor = max(cursor, int(r["id"])) - if kinds and r["kind"] not in kinds: - continue - if args.assignee and r["assignee"] != args.assignee: - continue - if args.tenant and r["tenant"] != args.tenant: - continue - try: - payload = json.loads(r["payload"]) if r["payload"] else None - except Exception: - payload = None - pl = f" {payload}" if payload else "" - print( - f"[{_fmt_ts(r['created_at'])}] {r['task_id']:10s} " - f"{r['kind']:18s} (@{r['assignee'] or '-'}){pl}", - flush=True, - ) - time.sleep(max(0.1, args.interval)) - except KeyboardInterrupt: - print("\n(stopped)") - return 0 - - def _cmd_stats(args: argparse.Namespace) -> int: with kb.connect_closing() as conn: stats = kb.board_stats(conn) @@ -1618,10 +1320,7 @@ def _cmd_notify_unsubscribe(args: argparse.Namespace) -> int: platform=args.platform, chat_id=args.chat_id, thread_id=args.thread_id, ) - if not ok: - return _err("(no such subscription)") - print(f"Unsubscribed from {args.task_id}") - return 0 + return _ok_or_err(ok, "(no such subscription)", f"Unsubscribed from {args.task_id}") def _cmd_log(args: argparse.Namespace) -> int: @@ -1760,120 +1459,6 @@ def _cmd_decompose(args: argparse.Namespace) -> int: ) -def _cmd_gc(args: argparse.Namespace) -> int: - """Remove archived tasks' scratch workspaces, old events, and old worker logs.""" - import shutil - scratch_root = kb.workspaces_root() - removed_ws = 0 - with kb.connect_closing() as conn: - rows = conn.execute( - "SELECT id, workspace_kind, workspace_path, branch_name FROM tasks " - "WHERE status = 'archived'" - ).fetchall() - for row in rows: - if row["workspace_kind"] == "worktree": - # Backstop for worktrees that escaped the completion/archive hook. - # Same safety predicate: only clean, fully-pushed worktrees go. - wt_path = row["workspace_path"] - if wt_path and Path(wt_path).is_dir(): - kb._cleanup_worktree_workspace(row["id"], wt_path, row["branch_name"]) - if not Path(wt_path).is_dir(): - removed_ws += 1 - continue - if row["workspace_kind"] != "scratch": - continue - path = Path(row["workspace_path"] or (scratch_root / row["id"])) - try: - path = path.resolve() - except OSError: - continue - try: - path.relative_to(scratch_root.resolve()) - except ValueError: - # Safety: never delete outside the scratch root. - continue - if path.exists() and path.is_dir(): - shutil.rmtree(path, ignore_errors=True) - removed_ws += 1 - - event_days = getattr(args, "event_retention_days", 30) - log_days = getattr(args, "log_retention_days", 30) - with kb.connect_closing() as conn: - removed_events = kb.gc_events( - conn, older_than_seconds=event_days * 24 * 3600, - ) - removed_logs = kb.gc_worker_logs( - older_than_seconds=log_days * 24 * 3600, - ) - print(f"GC complete: {removed_ws} workspace(s), " - f"{removed_events} event row(s), {removed_logs} log file(s) removed") - return 0 - - -def _cmd_repair(args: argparse.Namespace) -> int: - """Integrity check + narrow index-REINDEX auto-repair. Dispatched BEFORE - the auto ``kb.init_db()`` (init refuses corrupt DBs). Exit 0 = healthy / - repaired / no DB file, 1 = still corrupt.""" - try: - report = kb.repair_db() - except Exception as exc: # locked/busy probe, unexpected I/O - return _err(f"kanban repair: {exc}") - - if getattr(args, "json", False): - _print_json({ - "status": report.status, - "db_path": str(report.db_path), - "messages": report.messages, - "post_repair_messages": report.post_repair_messages, - "backup_path": ( - str(report.backup_path) if report.backup_path else None - ), - "reindexed": report.reindexed, - }, ascii=True) - return 0 if report.status in {"ok", "repaired", "missing"} else 1 - - if report.status == "missing": - print(f"No kanban DB at {report.db_path} — nothing to repair.") - return 0 - if report.status == "ok": - print(f"{report.db_path}: integrity_check ok — no repair needed.") - return 0 - if report.status == "repaired": - print(f"{report.db_path}: repaired.") - print(f" reindexed: {', '.join(report.reindexed)}") - if report.backup_path: - print(f" pre-repair backup: {report.backup_path}") - print(" integrity_check now ok.") - return 0 - # still corrupt - print(f"{report.db_path}: CORRUPT.", file=sys.stderr) - for line in (report.messages or [])[:10]: - print(f" {line}", file=sys.stderr) - if report.reindexed: - print( - f" REINDEX ({', '.join(report.reindexed)}) attempted but " - f"integrity_check is still failing:", - file=sys.stderr, - ) - for line in (report.post_repair_messages or [])[:10]: - print(f" {line}", file=sys.stderr) - else: - print( - " Not an index-only failure — automatic REINDEX repair does " - "not apply (fail-closed).", - file=sys.stderr, - ) - if report.backup_path: - print(f" corrupt copy quarantined at: {report.backup_path}", - file=sys.stderr) - print( - " Recover manually (e.g. `sqlite3 kanban.db \".recover\"` into a " - "fresh file) or move the file aside to start a new board.", - file=sys.stderr, - ) - return 1 - - _HANDLERS = { "init": _cmd_init, "create": _cmd_create, "swarm": _cmd_swarm, "list": _cmd_list, "ls": _cmd_list, "show": _cmd_show, diff --git a/hermes_cli/kanban_ops.py b/hermes_cli/kanban_ops.py new file mode 100644 index 0000000000..0b1540b98c --- /dev/null +++ b/hermes_cli/kanban_ops.py @@ -0,0 +1,426 @@ +"""Dispatcher and maintenance verbs for ``hermes kanban``: ``dispatch``, +``daemon`` (deprecated standalone loop), ``tail``/``watch`` event streaming, +``gc`` and ``repair``. +""" + +from __future__ import annotations + +import argparse +import json +import os +import sys +import time +from pathlib import Path + +from hermes_cli import kanban_db as kb +from hermes_cli.kanban_output import _err, _fmt_ts, _print_json + + +def _kanban_config() -> dict: + """``config.yaml`` ``kanban:`` section, or ``{}`` when config can't be loaded.""" + try: + from hermes_cli.config import load_config + cfg = load_config() + return (cfg.get("kanban", {}) if isinstance(cfg, dict) else {}) or {} + except Exception: + return {} + + +def _cmd_tail(args: argparse.Namespace) -> int: + last_id = 0 + print(f"Tailing events for {args.task_id}. Ctrl-C to stop.") + try: + while True: + with kb.connect_closing() as conn: + events = kb.list_events(conn, args.task_id) + for e in events: + if e.id > last_id: + pl = f" {e.payload}" if e.payload else "" + print(f"[{_fmt_ts(e.created_at)}] {e.kind}{pl}", flush=True) + last_id = e.id + time.sleep(max(0.1, args.interval)) + except KeyboardInterrupt: + print("\n(stopped)") + return 0 + + +def _coerce_positive_int(value): + if value is None: + return None + try: + ival = int(value) + except (TypeError, ValueError): + return None + return ival if ival >= 1 else None + + +def _cmd_dispatch(args: argparse.Namespace) -> int: + # Honour kanban.default_assignee, kanban.max_in_progress, + # kanban.max_in_progress_per_profile and kanban.max_spawn with the same + # semantics as the gateway dispatch path. + try: + from hermes_cli.config import load_config + _cfg = load_config() + _kanban_cfg = _cfg.get("kanban", {}) if isinstance(_cfg, dict) else {} + default_assignee = (_kanban_cfg.get("default_assignee") or "").strip() or None + max_in_progress_per_profile = _coerce_positive_int( + _kanban_cfg.get("max_in_progress_per_profile") + ) + # Memory-derived default when unset — same fallback the gateway applies. + max_in_progress = kb.resolve_max_in_progress( + _coerce_positive_int(_kanban_cfg.get("max_in_progress")) + ) + # CLI --max is the more explicit signal, so it wins over kanban.max_spawn. + cli_max = getattr(args, "max", None) + max_spawn = cli_max if cli_max is not None else _coerce_positive_int( + _kanban_cfg.get("max_spawn") + ) + except Exception: + default_assignee = None + max_in_progress_per_profile = None + max_in_progress = None + max_spawn = getattr(args, "max", None) + with kb.connect_closing() as conn: + res = kb.dispatch_once( + conn, + dry_run=args.dry_run, + max_spawn=max_spawn, + max_in_progress=max_in_progress, + failure_limit=getattr(args, "failure_limit", kb.DEFAULT_SPAWN_FAILURE_LIMIT), + default_assignee=default_assignee, + max_in_progress_per_profile=max_in_progress_per_profile, + ) + if getattr(args, "json", False): + _print_json({ + "reclaimed": res.reclaimed, + "crashed": res.crashed, + "timed_out": res.timed_out, + "stale": res.stale, + "auto_blocked": res.auto_blocked, + "promoted": res.promoted, + "spawned": [ + {"task_id": tid, "assignee": who, "workspace": ws} + for (tid, who, ws) in res.spawned + ], + "skipped_unassigned": res.skipped_unassigned, + "skipped_nonspawnable": res.skipped_nonspawnable, + "skipped_per_profile_capped": [ + {"task_id": tid, "assignee": who, "current": current} + for (tid, who, current) in res.skipped_per_profile_capped + ], + "auto_assigned_default": res.auto_assigned_default, + }, ascii=True) + return 0 + print(f"Reclaimed: {res.reclaimed}") + for label, items in ( + ("Crashed: ", res.crashed), + ("Timed out: ", res.timed_out), + ("Stale: ", res.stale), + ("Auto-blocked:", res.auto_blocked), + ): + print(f"{label} {len(items)}") + if items: + print(f" {', '.join(items)}") + print(f"Promoted: {res.promoted}") + print(f"Spawned: {len(res.spawned)}") + tag = " (dry)" if args.dry_run else "" + for tid, who, ws in res.spawned: + print(f" - {tid} -> {who} @ {ws or '-'}{tag}") + if res.auto_assigned_default: + print( + f"Auto-assigned to kanban.default_assignee={default_assignee!r}: " + f"{', '.join(res.auto_assigned_default)}" + ) + if res.skipped_unassigned: + print(f"Skipped (unassigned): {', '.join(res.skipped_unassigned)}") + for tid, who, current in res.skipped_per_profile_capped: + print(f"Deferred ({who} at per-profile cap, {current} running): {tid}") + if res.skipped_nonspawnable: + print( + f"Skipped (non-spawnable assignee — terminal lane, OK): " + f"{', '.join(res.skipped_nonspawnable)}" + ) + return 0 + + +def _cmd_daemon(args: argparse.Namespace) -> int: + """Deprecated — the dispatcher now runs inside the gateway. + + Kept as a stub so old scripts/systemd units get a clear migration message. + ``--force`` (hidden from --help) keeps the standalone loop for hosts that + truly cannot run the gateway; the default path exits 2 so nobody + accidentally runs two dispatchers against the same kanban.db. + """ + if not getattr(args, "force", False): + return _err( + "hermes kanban daemon: DEPRECATED — the dispatcher now runs\n" + "inside the gateway. To use kanban:\n" + "\n" + " hermes gateway start # starts the gateway + embedded dispatcher\n" + "\n" + "Ready tasks will be picked up on the next dispatcher tick\n" + "(default: every 60 seconds). Configure via config.yaml:\n" + "\n" + " kanban:\n" + " dispatch_in_gateway: true # default\n" + " dispatch_interval_seconds: 60\n" + " failure_limit: 2 # consecutive non-success attempts before auto-block\n" + "\n" + "Running both the gateway AND this standalone daemon will\n" + "race for claims. If you truly need the old standalone\n" + "daemon (no gateway available), rerun with --force.", + 2, + ) + + # Init before printing "started" so the DB path is right and init errors + # surface immediately. + kb.init_db() + + pidfile = getattr(args, "pidfile", None) + if pidfile: + try: + Path(pidfile).parent.mkdir(parents=True, exist_ok=True) + Path(pidfile).write_text(str(os.getpid()), encoding="utf-8") + except OSError as exc: + print(f"warning: could not write pidfile {pidfile}: {exc}", file=sys.stderr) + + verbose = bool(getattr(args, "verbose", False)) + print( + f"Kanban dispatcher running STANDALONE via --force " + f"(interval={args.interval}s, pid={os.getpid()}). " + f"Ctrl-C to stop. NOTE: if a gateway is also running with " + f"dispatch_in_gateway=true (default), you have two dispatchers " + f"racing for claims.", + file=sys.stderr, + ) + + # Health telemetry: warn when every tick finds ready work but spawns + # nothing (broken profile, PATH drift, missing venv, credential loss) — + # the per-task breaker auto-blocks quietly, so the operator needs a signal. + HEALTH_WINDOW = 6 # ticks (default 30s at interval=5) + health_state = {"bad_ticks": 0, "last_warn_at": 0} + + def _on_tick(res): + ready_pending = bool(res.skipped_unassigned) or _ready_queue_nonempty() + spawned_any = bool(res.spawned) + if ready_pending and not spawned_any: + health_state["bad_ticks"] += 1 + else: + health_state["bad_ticks"] = 0 + # Warn once per HEALTH_WINDOW bad ticks, at most every 5 minutes. + if health_state["bad_ticks"] >= HEALTH_WINDOW: + now = int(time.time()) + if now - health_state["last_warn_at"] >= 300: + print( + f"[{_fmt_ts(now)}] WARN dispatcher stuck: " + f"ready queue non-empty for {health_state['bad_ticks']} " + f"consecutive ticks but 0 workers spawned successfully. " + f"Check profile health (venv, PATH, credentials) and " + f"`hermes kanban list --status ready` / " + f"`hermes kanban list --status blocked` for recent " + f"spawn_failed tasks.", + file=sys.stderr, flush=True, + ) + health_state["last_warn_at"] = now + if not verbose: + return + did_work = ( + res.reclaimed or res.crashed or res.timed_out or res.promoted + or res.spawned or res.auto_blocked or res.stale + ) + if did_work: + print( + f"[{_fmt_ts(int(time.time()))}] " + f"reclaimed={res.reclaimed} crashed={len(res.crashed)} " + f"timed_out={len(res.timed_out)} stale={len(res.stale)} " + f"promoted={res.promoted} spawned={len(res.spawned)} " + f"auto_blocked={len(res.auto_blocked)}", + flush=True, + ) + + def _ready_queue_nonempty() -> bool: + """Is there a ready+assigned+unclaimed task the dispatcher would spawn for? + Control-plane lanes pulled via ``claim_task`` are correctly idle, not stuck.""" + try: + with kb.connect_closing() as conn: + return kb.has_spawnable_ready(conn) + except Exception: + return False + + try: + kb.run_daemon( + interval=args.interval, + max_spawn=args.max, + failure_limit=getattr(args, "failure_limit", kb.DEFAULT_SPAWN_FAILURE_LIMIT), + on_tick=_on_tick, + ) + finally: + if pidfile: + try: + Path(pidfile).unlink() + except OSError: + pass + print("(dispatcher stopped)") + return 0 + + +def _cmd_watch(args: argparse.Namespace) -> int: + """Live-stream task_events to the terminal.""" + kinds = ( + {k.strip() for k in args.kinds.split(",") if k.strip()} + if args.kinds else None + ) + print("Watching kanban events. Ctrl-C to stop.", flush=True) + # Seed cursor at the latest id so we don't replay history. + with kb.connect_closing() as conn: + row = conn.execute( + "SELECT COALESCE(MAX(id), 0) AS m FROM task_events" + ).fetchone() + cursor = int(row["m"]) + + try: + while True: + with kb.connect_closing() as conn: + rows = conn.execute( + "SELECT e.id, e.task_id, e.kind, e.payload, e.created_at, " + " t.assignee, t.tenant " + "FROM task_events e LEFT JOIN tasks t ON t.id = e.task_id " + "WHERE e.id > ? ORDER BY e.id ASC LIMIT 200", + (cursor,), + ).fetchall() + for r in rows: + cursor = max(cursor, int(r["id"])) + if kinds and r["kind"] not in kinds: + continue + if args.assignee and r["assignee"] != args.assignee: + continue + if args.tenant and r["tenant"] != args.tenant: + continue + try: + payload = json.loads(r["payload"]) if r["payload"] else None + except Exception: + payload = None + pl = f" {payload}" if payload else "" + print( + f"[{_fmt_ts(r['created_at'])}] {r['task_id']:10s} " + f"{r['kind']:18s} (@{r['assignee'] or '-'}){pl}", + flush=True, + ) + time.sleep(max(0.1, args.interval)) + except KeyboardInterrupt: + print("\n(stopped)") + return 0 + + +def _cmd_gc(args: argparse.Namespace) -> int: + """Remove archived tasks' scratch workspaces, old events, and old worker logs.""" + import shutil + scratch_root = kb.workspaces_root() + removed_ws = 0 + with kb.connect_closing() as conn: + rows = conn.execute( + "SELECT id, workspace_kind, workspace_path, branch_name FROM tasks " + "WHERE status = 'archived'" + ).fetchall() + for row in rows: + if row["workspace_kind"] == "worktree": + # Backstop for worktrees that escaped the completion/archive hook. + # Same safety predicate: only clean, fully-pushed worktrees go. + wt_path = row["workspace_path"] + if wt_path and Path(wt_path).is_dir(): + kb._cleanup_worktree_workspace(row["id"], wt_path, row["branch_name"]) + if not Path(wt_path).is_dir(): + removed_ws += 1 + continue + if row["workspace_kind"] != "scratch": + continue + path = Path(row["workspace_path"] or (scratch_root / row["id"])) + try: + path = path.resolve() + except OSError: + continue + try: + path.relative_to(scratch_root.resolve()) + except ValueError: + # Safety: never delete outside the scratch root. + continue + if path.exists() and path.is_dir(): + shutil.rmtree(path, ignore_errors=True) + removed_ws += 1 + + event_days = getattr(args, "event_retention_days", 30) + log_days = getattr(args, "log_retention_days", 30) + with kb.connect_closing() as conn: + removed_events = kb.gc_events( + conn, older_than_seconds=event_days * 24 * 3600, + ) + removed_logs = kb.gc_worker_logs( + older_than_seconds=log_days * 24 * 3600, + ) + print(f"GC complete: {removed_ws} workspace(s), " + f"{removed_events} event row(s), {removed_logs} log file(s) removed") + return 0 + + +def _cmd_repair(args: argparse.Namespace) -> int: + """Integrity check + narrow index-REINDEX auto-repair. Dispatched BEFORE + the auto ``kb.init_db()`` (init refuses corrupt DBs). Exit 0 = healthy / + repaired / no DB file, 1 = still corrupt.""" + try: + report = kb.repair_db() + except Exception as exc: # locked/busy probe, unexpected I/O + return _err(f"kanban repair: {exc}") + + if getattr(args, "json", False): + _print_json({ + "status": report.status, + "db_path": str(report.db_path), + "messages": report.messages, + "post_repair_messages": report.post_repair_messages, + "backup_path": ( + str(report.backup_path) if report.backup_path else None + ), + "reindexed": report.reindexed, + }, ascii=True) + return 0 if report.status in {"ok", "repaired", "missing"} else 1 + + if report.status == "missing": + print(f"No kanban DB at {report.db_path} — nothing to repair.") + return 0 + if report.status == "ok": + print(f"{report.db_path}: integrity_check ok — no repair needed.") + return 0 + if report.status == "repaired": + print(f"{report.db_path}: repaired.") + print(f" reindexed: {', '.join(report.reindexed)}") + if report.backup_path: + print(f" pre-repair backup: {report.backup_path}") + print(" integrity_check now ok.") + return 0 + # still corrupt + print(f"{report.db_path}: CORRUPT.", file=sys.stderr) + for line in (report.messages or [])[:10]: + print(f" {line}", file=sys.stderr) + if report.reindexed: + print( + f" REINDEX ({', '.join(report.reindexed)}) attempted but " + f"integrity_check is still failing:", + file=sys.stderr, + ) + for line in (report.post_repair_messages or [])[:10]: + print(f" {line}", file=sys.stderr) + else: + print( + " Not an index-only failure — automatic REINDEX repair does " + "not apply (fail-closed).", + file=sys.stderr, + ) + if report.backup_path: + print(f" corrupt copy quarantined at: {report.backup_path}", + file=sys.stderr) + print( + " Recover manually (e.g. `sqlite3 kanban.db \".recover\"` into a " + "fresh file) or move the file aside to start a new board.", + file=sys.stderr, + ) + return 1