refactor(kanban-cli): extract dispatcher/maintenance verbs to kanban_ops; unify single-mutation ok/err reporting
This commit is contained in:
@@ -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,
|
||||
|
||||
426
hermes_cli/kanban_ops.py
Normal file
426
hermes_cli/kanban_ops.py
Normal file
@@ -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
|
||||
Reference in New Issue
Block a user