diff --git a/hermes_cli/kanban.py b/hermes_cli/kanban.py index 0e0f3615e1..7d2cd8010c 100644 --- a/hermes_cli/kanban.py +++ b/hermes_cli/kanban.py @@ -1,11 +1,9 @@ """CLI for the Hermes Kanban board — ``hermes kanban …`` subcommand. -All DB work is delegated to ``kanban_db``. This module holds dispatch -(``kanban_command``), the task-verb handlers, and ``run_slash`` for -``/kanban …`` from the CLI and gateway. Siblings: ``kanban_parser`` -(argparse tree, re-exported here as ``build_parser``), ``kanban_output`` -(text/``--json`` helpers), ``kanban_boards`` (``boards …``), ``kanban_ops`` -(``dispatch``/``daemon``/``tail``/``watch``/``gc``/``repair``). +All DB work is delegated to ``kanban_db``. This module holds dispatch (``kanban_command``), the +task-verb handlers, and ``run_slash`` for ``/kanban …`` from the CLI and gateway. Siblings: +``kanban_parser`` (argparse tree, re-exported as ``build_parser``), ``kanban_output`` (text/--json +helpers), ``kanban_boards`` (``boards …``), ``kanban_ops`` (dispatch/daemon/tail/watch/gc/repair). """ from __future__ import annotations @@ -66,10 +64,7 @@ def _run_state_kwargs(args: argparse.Namespace, cmd: str) -> tuple[Optional[dict def _parse_workspace_flag(value: str) -> tuple[str, Optional[str]]: - """Parse ``--workspace`` into ``(kind, path|None)``. - - Accepts: ``scratch``, ``worktree``, ``worktree:``, ``dir:``. - """ + """``--workspace`` -> ``(kind, path|None)``: ``scratch``, ``worktree``, ``worktree:

``, ``dir:

``.""" if not value: return ("scratch", None) v = value.strip() @@ -80,9 +75,7 @@ def _parse_workspace_flag(value: str) -> tuple[str, Optional[str]]: continue path = v[len(prefix):].strip() if not path: - raise argparse.ArgumentTypeError( - f"--workspace {prefix} requires a path after the colon" - ) + raise argparse.ArgumentTypeError(f"--workspace {prefix} requires a path after the colon") return (kind, os.path.expanduser(path)) raise argparse.ArgumentTypeError( f"unknown --workspace value {value!r}: use scratch, worktree, " @@ -104,24 +97,19 @@ def _parse_branch_flag(value: Optional[str]) -> Optional[str]: return branch -def _check_dispatcher_presence( - hermes_home: Optional[Path] = None, -) -> tuple[bool, str]: +def _check_dispatcher_presence(hermes_home: Optional[Path] = None) -> tuple[bool, str]: """``(running, message)`` for the "will anything dispatch this?" warning. - ``running=True`` when a gateway is alive for this HERMES_HOME with - ``kanban.dispatch_in_gateway`` on; otherwise ``False`` plus human guidance. - Fails OPEN — import/probe/config errors return ``(True, "")`` — since a - missed warning beats crying wolf. ``hermes_home`` scopes the probe to a - profile dir (the dashboard backend may run under a different HERMES_HOME - than the profile it serves); CLI callers pass ``None``. + ``running=True`` when a gateway is alive for this HERMES_HOME with ``kanban.dispatch_in_gateway`` + on; otherwise ``False`` plus human guidance. Fails OPEN — import/probe/config errors return + ``(True, "")`` — since a missed warning beats crying wolf. ``hermes_home`` scopes the probe to a + profile dir (the dashboard backend may run under a different HERMES_HOME); CLI callers pass None. """ try: from gateway.status import resolve_gateway_liveness # type: ignore - # Same ladder as the dashboard status endpoints so PID-file-less or - # cross-container gateways aren't misreported; use_cache=False because - # this one-shot probe must see the gateway's state right now. + # Same ladder as the dashboard status endpoints so PID-file-less or cross-container + # gateways aren't misreported; use_cache=False because this one-shot probe must see now. liveness = resolve_gateway_liveness(profile_dir=hermes_home, use_cache=False) except Exception: return (True, "") # can't probe — silent @@ -129,19 +117,17 @@ def _check_dispatcher_presence( # The resolver swallows per-rung failures; "can't tell" != "no gateway". return (True, "") pid = liveness.pid - # Even if the gateway is up, dispatch_in_gateway may be off (can't tell -> assume default). dispatch_on = bool(_kanban_config().get("dispatch_in_gateway", True)) - if pid and dispatch_on: return (True, f"gateway pid={pid}, dispatch enabled") - if pid and not dispatch_on: + if pid: return ( False, "Gateway is running but kanban.dispatch_in_gateway=false in " "config.yaml — the task will sit in 'ready' until you flip it " "back on and restart the gateway, OR run the legacy " - "standalone daemon (`hermes kanban daemon --force`)." + "standalone daemon (`hermes kanban daemon --force`).", ) return ( False, @@ -150,7 +136,7 @@ def _check_dispatcher_presence( " hermes gateway start\n" "The gateway hosts an embedded dispatcher (tick interval 60s by " "default); your task will be picked up on the next tick after " - "the gateway comes up." + "the gateway comes up.", ) @@ -173,19 +159,18 @@ def kanban_command(args: argparse.Namespace) -> int: ) return 0 - # Fast-fail for UX only; the durable trust boundary is in kanban_db, since - # children can import DB mutators directly. + # Fast-fail for UX only; the durable trust boundary is in kanban_db, since children can + # import DB mutators directly. if _is_delegated_child_cli_mutation(args): return _err("kanban: delegate_task child contexts cannot mutate Kanban tasks via the CLI") - # `boards …` manages board metadata and the current-board pointer itself, so - # it must ignore the `--board` routing override (else `--board beta boards - # show` reports beta). + # `boards …` manages board metadata and the current-board pointer itself, so it must ignore + # the `--board` routing override (else `--board beta boards show` reports beta). if action == "boards": return _dispatch_boards(args) - # `--board ` pins HERMES_KANBAN_BOARD for the duration of this call so - # it inherits the exact resolution the dispatcher uses for workers. + # `--board ` pins HERMES_KANBAN_BOARD for the duration of this call so it inherits the + # exact resolution the dispatcher uses for workers. board_override = getattr(args, "board", None) board_scope = contextlib.nullcontext() if board_override: @@ -195,8 +180,8 @@ def kanban_command(args: argparse.Namespace) -> int: return _err(f"kanban: {exc}", 2) if not normed: return _err("kanban: --board requires a slug", 2) - # Boards other than 'default' must already exist — typoed slugs - # would otherwise silently create an empty board. + # Boards other than 'default' must already exist — typoed slugs would otherwise silently + # create an empty board. if normed != kb.DEFAULT_BOARD and not kb.board_exists(normed): return _err( f"kanban: board {normed!r} does not exist. " @@ -205,13 +190,12 @@ def kanban_command(args: argparse.Namespace) -> int: board_scope = kb.scoped_current_board(normed) with board_scope: - # `repair` dispatches BEFORE auto-init: on a corrupt DB init_db() itself - # raises KanbanDbCorruptError, which would turn every repair into - # "could not initialize database". + # `repair` dispatches BEFORE auto-init: on a corrupt DB init_db() itself raises + # KanbanDbCorruptError, which would turn every repair into "could not initialize database". if action == "repair": return _cmd_repair(args) - # init_db is idempotent (one sqlite_master SELECT when tables exist) and - # prevents "no such table: tasks" on first use from a fresh HERMES_HOME. + # init_db is idempotent (one sqlite_master SELECT when tables exist) and prevents + # "no such table: tasks" on first use from a fresh HERMES_HOME. try: kb.init_db() except Exception as exc: @@ -260,8 +244,7 @@ _DELEGATED_CHILD_DENIED_BOARD_ACTIONS: frozenset[str] = frozenset({ def _is_delegated_child_cli_mutation(args: argparse.Namespace) -> bool: action = getattr(args, "kanban_action", None) if action == "boards": - boards_action = getattr(args, "boards_action", None) or "list" - if boards_action not in _DELEGATED_CHILD_DENIED_BOARD_ACTIONS: + if (getattr(args, "boards_action", None) or "list") not in _DELEGATED_CHILD_DENIED_BOARD_ACTIONS: return False elif action not in _DELEGATED_CHILD_DENIED_ACTIONS: return False @@ -305,17 +288,15 @@ def _require_ids(args: argparse.Namespace) -> tuple[list[str], int]: def _parse_duration(val) -> Optional[int]: - """``30s`` / ``5m`` / ``2h`` / ``1d`` or a raw integer → seconds; None for - empty input; ValueError on malformed input.""" + """``30s`` / ``5m`` / ``2h`` / ``1d`` or a raw integer → seconds; None for empty input; + ValueError on malformed input.""" if val is None or val == "": return None s = str(val).strip().lower() - # Bare integer → seconds. try: - return int(s) + return int(s) # bare integer → seconds except ValueError: pass - # Suffixed form. units = {"s": 1, "m": 60, "h": 3600, "d": 86400} if s and s[-1] in units: try: @@ -329,7 +310,6 @@ def _parse_duration(val) -> Optional[int]: def _cmd_init(args: argparse.Namespace) -> int: path = kb.init_db() print(f"Kanban DB initialized at {path}") - print() # Profiles on disk == assignees already addressable. try: @@ -337,18 +317,15 @@ def _cmd_init(args: argparse.Namespace) -> int: except Exception: profiles = [] if profiles: - print(f"Discovered {len(profiles)} profile(s) on disk; any of these can " - f"be an --assignee:") + print(f"Discovered {len(profiles)} profile(s) on disk; any of these can be an --assignee:") for name in profiles: print(f" {name}") else: - print("No profiles found under ~/.hermes/profiles/.") - print("Create one with `hermes -p setup` before assigning tasks.") - print() - print("Next step: start the gateway so ready tasks actually get picked up.") - print(" hermes gateway start") - print() + print("No profiles found under ~/.hermes/profiles/.\n" + "Create one with `hermes -p setup` before assigning tasks.") print( + "\nNext step: start the gateway so ready tasks actually get picked up.\n" + " hermes gateway start\n\n" "The gateway hosts an embedded dispatcher that ticks every 60 seconds\n" "by default (config: kanban.dispatch_interval_seconds). Without a\n" "running gateway, tasks stay in 'ready' forever." @@ -359,9 +336,7 @@ def _cmd_init(args: argparse.Namespace) -> int: def _cmd_heartbeat(args: argparse.Namespace) -> int: with kb.connect_closing() as conn: ok = kb.heartbeat_worker( - conn, - args.task_id, - note=getattr(args, "note", None), + conn, args.task_id, note=getattr(args, "note", None), expected_run_id=_worker_run_id_for(args.task_id), ) return _ok_or_err(ok, f"cannot heartbeat {args.task_id} (not running?)", @@ -404,24 +379,14 @@ def _cmd_create(args: argparse.Namespace) -> int: ) with kb.connect_closing() as conn: task_id = kb.create_task( - conn, - title=args.title, - body=args.body, - assignee=args.assignee, + conn, title=args.title, body=args.body, assignee=args.assignee, created_by=args.created_by or _profile_author(), - workspace_kind=ws_kind, - workspace_path=ws_path, - branch_name=branch_name, - project_id=getattr(args, "project", None), - tenant=args.tenant, - priority=args.priority, - parents=tuple(args.parent or ()), - triage=bool(getattr(args, "triage", False)), + workspace_kind=ws_kind, workspace_path=ws_path, branch_name=branch_name, + project_id=getattr(args, "project", None), tenant=args.tenant, priority=args.priority, + parents=tuple(args.parent or ()), triage=bool(getattr(args, "triage", False)), idempotency_key=getattr(args, "idempotency_key", None), - max_runtime_seconds=max_runtime, - skills=getattr(args, "skills", None) or None, - max_retries=max_retries, - model_override=getattr(args, "model_override", None), + max_runtime_seconds=max_runtime, skills=getattr(args, "skills", None) or None, + max_retries=max_retries, model_override=getattr(args, "model_override", None), provider_override=getattr(args, "provider_override", None), goal_mode=bool(getattr(args, "goal_mode", False)), goal_max_turns=getattr(args, "goal_max_turns", None), @@ -432,9 +397,8 @@ def _cmd_create(args: argparse.Namespace) -> int: _print_json(_task_to_dict(task)) else: print(f"Created {task_id} ({task.status}, assignee={task.assignee or '-'})") - # Warn only for ready+assigned tasks that would sit without a dispatcher - # (triage/todo idle by design, unassigned can't dispatch); skipped under - # --json so stdout stays machine-parseable. + # Warn only for ready+assigned tasks that would sit without a dispatcher (triage/todo idle + # by design, unassigned can't dispatch); skipped under --json so stdout stays parseable. if task.status == "ready" and task.assignee: running, message = _check_dispatcher_presence() if not running and message: @@ -451,23 +415,18 @@ def _cmd_swarm(args: argparse.Namespace) -> int: return _err("kanban swarm: at least one --worker is required", 2) with kb.connect_closing() as conn: created = ks.create_swarm( - conn, - goal=args.goal, - workers=workers, - verifier_assignee=args.verifier, - synthesizer_assignee=args.synthesizer, - tenant=args.tenant, - created_by=args.created_by or _profile_author(), - priority=args.priority, + conn, goal=args.goal, workers=workers, verifier_assignee=args.verifier, + synthesizer_assignee=args.synthesizer, tenant=args.tenant, + created_by=args.created_by or _profile_author(), priority=args.priority, idempotency_key=getattr(args, "idempotency_key", None), ) if getattr(args, "json", False): _print_json(created.as_dict()) else: - print(f"Swarm root: {created.root_id}") - print("Workers: " + ", ".join(created.worker_ids)) - print(f"Verifier: {created.verifier_id}") - print(f"Synthesizer: {created.synthesizer_id}") + print(f"Swarm root: {created.root_id}\n" + "Workers: " + ", ".join(created.worker_ids) + "\n" + f"Verifier: {created.verifier_id}\n" + f"Synthesizer: {created.synthesizer_id}") return 0 @@ -479,15 +438,9 @@ def _cmd_list(args: argparse.Namespace) -> int: # Cheap mini-dispatch so list reflects dependencies cleared since the last tick. kb.recompute_ready(conn) tasks = kb.list_tasks( - conn, - assignee=assignee, - status=args.status, - tenant=args.tenant, - session_id=args.session, - include_archived=args.archived, - order_by=getattr(args, "sort", None), - workflow_template_id=args.workflow_template_id, - current_step_key=args.current_step_key, + conn, assignee=assignee, status=args.status, tenant=args.tenant, session_id=args.session, + include_archived=args.archived, order_by=getattr(args, "sort", None), + workflow_template_id=args.workflow_template_id, current_step_key=args.current_step_key, ) if _json_out(args, [_task_to_dict(t) for t in tasks]): return 0 @@ -497,10 +450,9 @@ def _cmd_list(args: argparse.Namespace) -> int: except Exception: all_boards = [] if len(all_boards) > 1: - current = kb.get_current_board() other_count = len(all_boards) - 1 print( - f"Board: {current} " + f"Board: {kb.get_current_board()} " f"({other_count} other board{'s' if other_count != 1 else ''} — " f"`hermes kanban boards list`)\n" ) @@ -530,11 +482,20 @@ def _print_diagnostics(diags, indent: str, *, with_kind: bool) -> None: print(f"{indent} → {a.label}") +def _print_section(title: str, lines) -> None: + """Blank line, ``title``, then each line (``show`` body sections).""" + print() + print(title) + for line in lines: + print(line) + + def _cmd_show(args: argparse.Namespace) -> int: rsk, rc = _run_state_kwargs(args, "show") if rc: return rc graph = None + want_json = getattr(args, "json", False) with kb.connect_closing() as conn: task = kb.get_task(conn, args.task_id) if not task: @@ -546,10 +507,10 @@ def _cmd_show(args: argparse.Namespace) -> int: runs = kb.list_runs(conn, args.task_id, **rsk) # Workers hand off via task_runs.summary; tasks.result stays NULL unless set. latest_summary = kb.latest_summary(conn, args.task_id) - if not getattr(args, "json", False): + if not want_json: graph = kb.task_graph_context(conn, task.id) - if getattr(args, "json", False): + if want_json: _print_json({ "task": _task_to_dict(task), "latest_summary": latest_summary, @@ -561,20 +522,22 @@ def _cmd_show(args: argparse.Namespace) -> int: }) return 0 + def field(label: str, value) -> None: + print(f" {label + ':':<11}{value}") + print(f"Task {task.id}: {task.title}") - print(f" status: {task.status}") - print(f" assignee: {task.assignee or '-'}") + field("status", task.status) + field("assignee", task.assignee or "-") if task.tenant: - print(f" tenant: {task.tenant}") - print(f" workspace: {task.workspace_kind}" + - (f" @ {task.workspace_path}" if task.workspace_path else "")) + field("tenant", task.tenant) + field("workspace", f"{task.workspace_kind}" + (f" @ {task.workspace_path}" if task.workspace_path else "")) if task.branch_name: - print(f" branch: {task.branch_name}") + field("branch", task.branch_name) if task.skills: - print(f" skills: {', '.join(task.skills)}") + field("skills", ", ".join(task.skills)) if task.model_override: _prov = f" (provider: {task.provider_override})" if task.provider_override else "" - print(f" model: {task.model_override}{_prov}") + field("model", f"{task.model_override}{_prov}") # Effective retry threshold (task > config > default) explains auto-blocks. if task.max_retries is not None: print(f" max-retries: {task.max_retries} (task)") @@ -584,7 +547,7 @@ def _cmd_show(args: argparse.Namespace) -> int: print(f" max-retries: {int(cfg_val)} (config kanban.failure_limit)") else: print(f" max-retries: {kb.DEFAULT_FAILURE_LIMIT} (default)") - print(f" created: {_fmt_ts(task.created_at)} by {task.created_by or '-'}") + field("created", f"{_fmt_ts(task.created_at)} by {task.created_by or '-'}") # Diagnostics up top so CLI users see distress signals before scrolling. from hermes_cli import kanban_diagnostics as kd @@ -593,48 +556,37 @@ def _cmd_show(args: argparse.Namespace) -> int: print(f"\n Diagnostics ({len(diags)}):") _print_diagnostics(diags, " ", with_kind=False) if task.started_at: - print(f" started: {_fmt_ts(task.started_at)}") + field("started", _fmt_ts(task.started_at)) if task.completed_at: - print(f" completed: {_fmt_ts(task.completed_at)}") + field("completed", _fmt_ts(task.completed_at)) if parents: - print(f" parents: {', '.join(parents)}") + field("parents", ", ".join(parents)) if children: - print(f" children: {', '.join(children)}") + field("children", ", ".join(children)) if task.body: - print() - print("Body:") - print(task.body) + _print_section("Body:", [task.body]) if task.result: - print() - print("Result:") - print(task.result) + _print_section("Result:", [task.result]) elif latest_summary: - print() - print("Latest summary:") - print(latest_summary) + _print_section("Latest summary:", [latest_summary]) if comments: - print() - print(f"Comments ({len(comments)}):") - for c in comments: - print(f" [{_fmt_ts(c.created_at)}] {c.author}: {c.body}") + _print_section(f"Comments ({len(comments)}):", + (f" [{_fmt_ts(c.created_at)}] {c.author}: {c.body}" for c in comments)) if events: - print() - print(f"Events ({len(events)}):") - for e in events[-20:]: - pl = f" {e.payload}" if e.payload else "" - run_tag = f" [run {e.run_id}]" if e.run_id else "" - print(f" [{_fmt_ts(e.created_at)}]{run_tag} {e.kind}{pl}") + _print_section(f"Events ({len(events)}):", ( + f" [{_fmt_ts(e.created_at)}]{f' [run {e.run_id}]' if e.run_id else ''} {e.kind}" + f"{f' {e.payload}' if e.payload else ''}" + for e in events[-20:] + )) if runs: print() print(f"Runs ({len(runs)}):") for r in runs: # Clamp to 0 so NTP backward-jumps don't print negative seconds. - elapsed = (max(0, r.ended_at - r.started_at) - if r.ended_at else None) + elapsed = max(0, r.ended_at - r.started_at) if r.ended_at else None el = f"{elapsed}s" if elapsed is not None else "active" outcome = r.outcome or r.status or "active" - print(f" #{r.id:<3} {outcome:<12} @{r.profile or '-'} {el} " - f"{_fmt_ts(r.started_at)}") + print(f" #{r.id:<3} {outcome:<12} @{r.profile or '-'} {el} {_fmt_ts(r.started_at)}") if r.summary: print(f" → {r.summary.splitlines()[0][:160]}") if r.error: @@ -664,37 +616,30 @@ def _cmd_set_model(args: argparse.Namespace) -> int: return _err(f"no such task: {args.task_id}") if model: label = f"{provider}:{model}" if provider else model - print(f"Set model override on {args.task_id}: {label} " - "(applies on next dispatch)") + print(f"Set model override on {args.task_id}: {label} (applies on next dispatch)") else: - print(f"Cleared model override on {args.task_id} " - "(worker uses its profile default)") + print(f"Cleared model override on {args.task_id} (worker uses its profile default)") return 0 def _cmd_reclaim(args: argparse.Namespace) -> int: with kb.connect_closing() as conn: - ok = kb.reclaim_task( - conn, args.task_id, - reason=getattr(args, "reason", None), - ) + ok = kb.reclaim_task(conn, args.task_id, reason=getattr(args, "reason", None)) 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: profile = _none_profile(args.profile) + reclaim = bool(getattr(args, "reclaim", False)) with kb.connect_closing() as conn: ok = kb.reassign_task( - conn, args.task_id, profile, - reclaim_first=bool(getattr(args, "reclaim", False)), - reason=getattr(args, "reason", None), + conn, args.task_id, profile, reclaim_first=reclaim, reason=getattr(args, "reason", None), ) 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 ""), + f"Reassigned {args.task_id} to {profile or '(unassigned)'}" + (" (claim reclaimed)" if reclaim else ""), ) @@ -724,18 +669,13 @@ def _cmd_diagnostics(args: argparse.Namespace) -> int: return _err(f"no such task: {args.task}") diags_by_task = { args.task: kd.compute_task_diagnostics( - task, - kb.list_events(conn, args.task), - kb.list_runs(conn, args.task), - graph=kb.task_graph_context(conn, args.task), - config=diag_config, + task, kb.list_events(conn, args.task), kb.list_runs(conn, args.task), + graph=kb.task_graph_context(conn, args.task), config=diag_config, ) } else: # Fleet mode: pull all non-archived tasks + their events/runs. - rows = list(conn.execute( - "SELECT * FROM tasks WHERE status != 'archived'" - ).fetchall()) + rows = list(conn.execute("SELECT * FROM tasks WHERE status != 'archived'").fetchall()) ids = [r["id"] for r in rows] diags_by_task = {} if ids: @@ -745,16 +685,12 @@ def _cmd_diagnostics(args: argparse.Namespace) -> int: for r in rows: tid = r["id"] dl = kd.compute_task_diagnostics( - r, - ev_by.get(tid, []), - run_by.get(tid, []), - graph=graph_by.get(tid), - config=diag_config, + r, ev_by.get(tid, []), run_by.get(tid, []), + graph=graph_by.get(tid), config=diag_config, ) if dl: diags_by_task[tid] = dl - # Severity filter. sev = getattr(args, "severity", None) if sev: floor = kd.SEVERITY_ORDER.index(sev) @@ -776,11 +712,7 @@ def _cmd_diagnostics(args: argparse.Namespace) -> int: if getattr(args, "json", False): _print_json([ - { - "task_id": tid, - **meta.get(tid, {}), - "diagnostics": [d.to_dict() for d in dl], - } + {"task_id": tid, **meta.get(tid, {}), "diagnostics": [d.to_dict() for d in dl]} for tid, dl in diags_by_task.items() ]) return 0 @@ -790,10 +722,7 @@ def _cmd_diagnostics(args: argparse.Namespace) -> int: return 0 total = sum(len(dl) for dl in diags_by_task.values()) - print( - f"{total} active diagnostic(s) across " - f"{len(diags_by_task)} task(s):\n" - ) + print(f"{total} active diagnostic(s) across {len(diags_by_task)} task(s):\n") for tid, dl in diags_by_task.items(): m = meta.get(tid, {}) title = m.get("title") or "(untitled)" @@ -832,8 +761,7 @@ def _cmd_claim(args: argparse.Namespace) -> int: ) workspace = kb.resolve_workspace(task) kb.set_workspace_path(conn, task.id, str(workspace)) - print(f"Claimed {task.id}") - print(f"Workspace: {workspace}") + print(f"Claimed {task.id}\nWorkspace: {workspace}") return 0 @@ -853,8 +781,8 @@ def _cmd_comment(args: argparse.Namespace) -> int: def _cmd_attach(args: argparse.Namespace) -> int: - """Attach a local file via the shared ``store_attachment_bytes`` path (same - 25 MB cap and name sanitisation as the dashboard upload and agent tool).""" + """Attach a local file via the shared ``store_attachment_bytes`` path (same 25 MB cap and name + sanitisation as the dashboard upload and agent tool).""" import mimetypes src = Path(args.path).expanduser() @@ -867,12 +795,7 @@ def _cmd_attach(args: argparse.Namespace) -> int: try: with kb.connect_closing() as conn: att_id = kb.store_attachment_bytes( - conn, - args.task_id, - name, - data, - content_type=content_type, - uploaded_by=uploaded_by, + conn, args.task_id, name, data, content_type=content_type, uploaded_by=uploaded_by, ) except kb.AttachmentTooLarge as exc: return _err(f"kanban: {exc}") @@ -922,9 +845,9 @@ def _worker_run_id_for(task_id: str) -> Optional[int]: def _goal_mode_handoff_rejection(task: Optional[kb.Task], evidence: str): """Goal judge for every terminal worker handoff (including review). - Returns ``(verdict, reason_or_None)``: ``"done"`` allows; ``"blocked"`` = - judge ruled the goal unachievable; ``"continue"``/``"wait"`` reject with - the judge's reason. Judge failures allow the handoff (logged). + Returns ``(verdict, reason_or_None)``: ``"done"`` allows; ``"blocked"`` = judge ruled the goal + unachievable; ``"continue"``/``"wait"`` reject with the judge's reason. Judge failures allow + the handoff (logged). """ if task is None or not task.goal_mode: return ("done", None) @@ -943,25 +866,22 @@ def _goal_mode_handoff_rejection(task: Optional[kb.Task], evidence: str): reason = "" try: verdict, reason, _, _, _ = judge_goal( - goal=f"{task.title}\n\n{task.body or ''}".strip(), - last_response=evidence.strip(), + goal=f"{task.title}\n\n{task.body or ''}".strip(), last_response=evidence.strip(), ) except Exception as judge_exc: import logging as _logging _logging.getLogger(__name__).warning( - "goal judge check failed, allowing lifecycle handoff: %s", - judge_exc, - exc_info=True, + "goal judge check failed, allowing lifecycle handoff: %s", judge_exc, exc_info=True, ) return (verdict, None if verdict == "done" else reason) def _goal_gate_error(conn, tid: str, evidence: str, handoff: str, blocked_hint: str, continue_hint: str) -> Optional[str]: - """Goal-mode judge gate shared by ``complete`` / ``request-review`` - (mirrors tools/kanban_tools.py); applied to every terminal handoff so - request-review can't bypass it. Returns the error line, or None to allow.""" + """Goal-mode judge gate shared by ``complete`` / ``request-review`` (mirrors tools/kanban_tools.py); + applied to every terminal handoff so request-review can't bypass it. Returns the error line, or + None to allow.""" verdict, rejection = _goal_mode_handoff_rejection(kb.get_task(conn, tid), evidence) if verdict == "blocked": return (f"kanban: goal {handoff} of {tid} rejected: judge ruled " @@ -1003,10 +923,7 @@ def _cmd_complete(args: argparse.Namespace) -> int: return False fail_msg[tid] = f"cannot complete {tid} (unknown id or terminal state)" return kb.complete_task( - conn, tid, - result=args.result, - summary=summary, - metadata=metadata, + conn, tid, result=args.result, summary=summary, metadata=metadata, expected_run_id=_worker_run_id_for(tid), ) @@ -1019,10 +936,7 @@ def _cmd_edit(args: argparse.Namespace) -> int: return rc with kb.connect_closing() as conn: if not kb.edit_completed_task_result( - conn, - args.task_id, - result=args.result, - summary=getattr(args, "summary", None), + conn, args.task_id, result=args.result, summary=getattr(args, "summary", None), metadata=metadata, ): return _err(f"cannot edit {args.task_id} (unknown id or task is not done)") @@ -1030,6 +944,15 @@ def _cmd_edit(args: argparse.Namespace) -> int: return 0 +def _commented(conn, reason: Optional[str], author, prefix: str, op): + """Wrap a per-task ``op`` so a ``reason`` is first recorded as a ``PREFIX: reason`` comment.""" + def run(tid): + if reason: + kb.add_comment(conn, tid, author, f"{prefix}: {reason}") + return op(tid) + return run + + def _cmd_block(args: argparse.Namespace) -> int: reason = _joined_words(args.reason) kind = getattr(args, "kind", None) @@ -1037,14 +960,6 @@ def _cmd_block(args: argparse.Namespace) -> int: ids = _bulk_ids(args) suffix = f": {reason}" if reason else "" with kb.connect_closing() as conn: - def op(tid): - if reason: - kb.add_comment(conn, tid, author, f"BLOCKED: {reason}") - return kb.block_task( - conn, tid, reason=reason, kind=kind, - expected_run_id=_worker_run_id_for(tid), - ) - def ok_msg(tid): # Report where it landed: dependency blocks -> todo, tripped unblock-loop breaker -> triage. landed = kb.get_task(conn, tid) @@ -1052,10 +967,12 @@ def _cmd_block(args: argparse.Namespace) -> int: if where == "todo": return f"{tid} → todo (dependency wait){suffix}" if where == "triage": - return (f"{tid} → triage (unblock loop detected — needs a " - f"human decision){suffix}") + return f"{tid} → triage (unblock loop detected — needs a human decision){suffix}" return f"Blocked {tid}{suffix}" + op = _commented(conn, reason, author, "BLOCKED", lambda tid: kb.block_task( + conn, tid, reason=reason, kind=kind, expected_run_id=_worker_run_id_for(tid), + )) return _bulk_apply(ids, op, ok_msg, lambda tid: f"cannot block {tid}") @@ -1065,13 +982,9 @@ def _cmd_schedule(args: argparse.Namespace) -> int: ids = _bulk_ids(args) suffix = f": {reason}" if reason else "" with kb.connect_closing() as conn: - def op(tid): - if reason: - kb.add_comment(conn, tid, author, f"SCHEDULED: {reason}") - return kb.schedule_task( - conn, tid, reason=reason, expected_run_id=_worker_run_id_for(tid), - ) - + op = _commented(conn, reason, author, "SCHEDULED", lambda tid: kb.schedule_task( + conn, tid, reason=reason, expected_run_id=_worker_run_id_for(tid), + )) return _bulk_apply( ids, op, lambda tid: f"Scheduled {tid}{suffix}", lambda tid: f"cannot schedule {tid}", ) @@ -1085,11 +998,7 @@ def _cmd_unblock(args: argparse.Namespace) -> int: author = _profile_author() if reason else None suffix = f": {reason}" if reason else "" with kb.connect_closing() as conn: - def op(tid): - if reason: - kb.add_comment(conn, tid, author, f"UNBLOCK: {reason}") - return kb.unblock_task(conn, tid) - + op = _commented(conn, reason, author, "UNBLOCK", lambda tid: kb.unblock_task(conn, tid)) return _bulk_apply( ids, op, lambda tid: f"Unblocked {tid}{suffix}", lambda tid: f"cannot unblock {tid} (not blocked/scheduled?)", @@ -1102,7 +1011,6 @@ def _cmd_request_review(args: argparse.Namespace) -> int: metadata, rc = _parse_metadata_flag(getattr(args, "metadata", None)) if rc: return rc - reviewer = getattr(args, "reviewer", None) with kb.connect_closing() as conn: gate_err = _goal_gate_error( conn, tid, summary or "", "review handoff", @@ -1112,23 +1020,15 @@ def _cmd_request_review(args: argparse.Namespace) -> int: if gate_err: return _err(gate_err) ok, reason = kb.request_review( - conn, - tid, - summary=summary, - metadata=metadata, - reviewer=reviewer, - expected_run_id=_worker_run_id_for(tid), - force=bool(getattr(args, "force", False)), + conn, tid, summary=summary, metadata=metadata, reviewer=getattr(args, "reviewer", None), + expected_run_id=_worker_run_id_for(tid), force=bool(getattr(args, "force", False)), with_reason=True, ) if not ok: return _err(f"cannot request review for {tid}: {reason or 'not running/ready?'}") persisted_run = kb.latest_run(conn, tid) display_summary = persisted_run.summary if persisted_run else None - print( - f"Requested review for {tid}" - + (f": {display_summary}" if display_summary else "") - ) + print(f"Requested review for {tid}" + (f": {display_summary}" if display_summary else "")) return 0 @@ -1136,18 +1036,10 @@ def _cmd_request_changes(args: argparse.Namespace) -> int: tid = args.task_id reason = " ".join(args.reason).strip() with kb.connect_closing() as conn: - ok, detail = kb.request_changes( - conn, - tid, - reason=reason, - expected_run_id=_worker_run_id_for(tid), - ) + ok, detail = kb.request_changes(conn, tid, reason=reason, expected_run_id=_worker_run_id_for(tid)) if not ok: return _err(f"cannot request changes for {tid}: {detail or 'invalid review state'}") - print( - f"Requested changes for {tid}" - + (f"; routed to {detail}" if detail else "") - ) + print(f"Requested changes for {tid}" + (f"; routed to {detail}" if detail else "")) return 0 @@ -1169,8 +1061,7 @@ def _cmd_reopen_review(args: argparse.Namespace) -> int: return True return _bulk_apply( - ids, op, lambda tid: f"Reopened {tid}{suffix}", - lambda tid: f"cannot reopen {tid} (not in review?)", + ids, op, lambda tid: f"Reopened {tid}{suffix}", lambda tid: f"cannot reopen {tid} (not in review?)", ) @@ -1179,25 +1070,15 @@ def _cmd_promote(args: argparse.Namespace) -> int: author = _profile_author() # Dedupe while preserving order; positional task_id always first. ids = list(dict.fromkeys(_bulk_ids(args))) + dry_run, force = bool(args.dry_run), bool(args.force) results: list[dict[str, object]] = [] with kb.connect_closing() as conn: for tid in ids: - ok, err = kb.promote_task( - conn, - tid, - actor=author, - reason=reason, - force=bool(args.force), - dry_run=bool(args.dry_run), - ) + ok, err = kb.promote_task(conn, tid, actor=author, reason=reason, force=force, dry_run=dry_run) results.append({ - "task_id": tid, - "promoted": ok, - "dry_run": bool(args.dry_run), - "forced": bool(args.force), - "reason": reason, - "error": err, + "task_id": tid, "promoted": ok, "dry_run": dry_run, "forced": force, + "reason": reason, "error": err, }) failed = [r for r in results if not r["promoted"]] @@ -1206,11 +1087,11 @@ def _cmd_promote(args: argparse.Namespace) -> int: _print_json(results[0] if len(results) == 1 else results) return 0 if not failed else 1 - tag = " (dry)" if args.dry_run else "" - label = "Would promote" if args.dry_run else "Promoted" + tag = " (dry)" if dry_run else "" + label = "Would promote" if dry_run else "Promoted" + suffix = f": {reason}" if reason else "" for r in results: if r["promoted"]: - suffix = f": {reason}" if reason else "" print(f"{label} {r['task_id']} -> ready{tag}{suffix}") else: print(f"cannot promote {r['task_id']}: {r['error']}", file=sys.stderr) @@ -1228,8 +1109,7 @@ def _cmd_archive(args: argparse.Namespace) -> int: if purge_ids: return _bulk_apply( purge_ids, lambda tid: kb.delete_archived_task(conn, tid), - lambda tid: f"Deleted {tid}", - lambda tid: f"cannot delete {tid} (must already be archived)", + lambda tid: f"Deleted {tid}", lambda tid: f"cannot delete {tid} (must already be archived)", ) return _bulk_apply( ids, lambda tid: kb.archive_task(conn, tid), @@ -1260,10 +1140,8 @@ def _cmd_notify_subscribe(args: argparse.Namespace) -> int: if kb.get_task(conn, args.task_id) is None: return _err(f"no such task: {args.task_id}") kb.add_notify_sub( - conn, task_id=args.task_id, - platform=args.platform, chat_id=args.chat_id, - chat_type=args.chat_type, - thread_id=args.thread_id, user_id=args.user_id, + conn, task_id=args.task_id, platform=args.platform, chat_id=args.chat_id, + chat_type=args.chat_type, thread_id=args.thread_id, user_id=args.user_id, user_id_alt=getattr(args, "user_id_alt", None), notifier_profile=args.notifier_profile or _profile_author(), delivery_mode=getattr(args, "delivery_mode", None), @@ -1298,8 +1176,7 @@ def _cmd_notify_list(args: argparse.Namespace) -> int: def _cmd_notify_unsubscribe(args: argparse.Namespace) -> int: with kb.connect_closing() as conn: ok = kb.remove_notify_sub( - conn, task_id=args.task_id, - platform=args.platform, chat_id=args.chat_id, + conn, task_id=args.task_id, platform=args.platform, chat_id=args.chat_id, thread_id=args.thread_id, ) return _ok_or_err(ok, "(no such subscription)", f"Unsubscribed from {args.task_id}") @@ -1332,12 +1209,8 @@ def _cmd_runs(args: argparse.Namespace) -> int: end = r.ended_at or int(time.time()) # Clamp to 0 so NTP backward-jumps don't print negative durations. elapsed = max(0, end - r.started_at) - if elapsed < 60: - el = f"{elapsed}s" - elif elapsed < 3600: - el = f"{elapsed // 60}m" - else: - el = f"{elapsed / 3600:.1f}h" + el = (f"{elapsed}s" if elapsed < 60 else f"{elapsed // 60}m" if elapsed < 3600 + else f"{elapsed / 3600:.1f}h") outcome = r.outcome or ("(running)" if not r.ended_at else r.status) print(f"{i:3d} {outcome:12s} {(r.profile or '-'):16s} {el:>8s} {_fmt_ts(r.started_at)}") if r.summary: @@ -1379,8 +1252,8 @@ def _triage_sweep_ids(args: argparse.Namespace, verb: str, list_triage_ids, json def _run_triage_sweep(args: argparse.Namespace, verb: str, mod, run_one, json_key: str, json_fields: tuple[str, ...], human_ok) -> int: - """Shared driver for ``specify`` / ``decompose``: validate ids, run - ``run_one(tid, author=...)`` per id, print JSON or human lines, exit code.""" + """Shared driver for ``specify`` / ``decompose``: validate ids, run ``run_one(tid, author=...)`` + per id, print JSON or human lines, exit code.""" all_flag = bool(getattr(args, "all_triage", False)) author = getattr(args, "author", None) or _profile_author() want_json = bool(getattr(args, "json", False)) @@ -1414,8 +1287,7 @@ def _cmd_specify(args: argparse.Namespace) -> int: from hermes_cli import kanban_specify as spec return _run_triage_sweep( - args, "specify", spec, spec.specify_task, "specified", - ("task_id", "ok", "reason", "new_title"), + args, "specify", spec, spec.specify_task, "specified", ("task_id", "ok", "reason", "new_title"), lambda o: f"Specified {o.task_id} → todo{_retitled_suffix(o)}", ) @@ -1433,8 +1305,7 @@ def _cmd_decompose(args: argparse.Namespace) -> int: return _run_triage_sweep( args, "decompose", decomp, decomp.decompose_task, "decomposed", - ("task_id", "ok", "reason", "fanout", "child_ids", "new_title"), - _decompose_ok_line, + ("task_id", "ok", "reason", "fanout", "child_ids", "new_title"), _decompose_ok_line, ) @@ -1491,39 +1362,36 @@ Read-only commands are safe while an agent is running.\ def run_slash(rest: str) -> str: - """Execute a ``/kanban …`` string (``rest`` = everything after ``/kanban``) - and return captured stdout/stderr. Shared by the interactive CLI and the - gateway so formatting is identical.""" + """Execute a ``/kanban …`` string (``rest`` = everything after ``/kanban``) and return captured + stdout/stderr. Shared by the interactive CLI and the gateway so formatting is identical.""" import io tokens = shlex.split(rest) if rest and rest.strip() else [] - # Bare ``/kanban`` / ``help`` / ``-h``: curated short block, not argparse's - # full tree (garbage in a chat bubble). ``/kanban foo -h`` still works. + # Bare ``/kanban`` / ``help`` / ``-h``: curated short block, not argparse's full tree (garbage + # in a chat bubble). ``/kanban foo -h`` still works. if not tokens or tokens[0] in {"help", "--help", "-h", "?"}: return _SLASH_KANBAN_HELP - # build_parser() needs a subparsers action to attach to: build a throwaway - # one and drive kanban_parser directly so usage/error text reads ``/kanban``. + # build_parser() needs a subparsers action to attach to: build a throwaway one and drive + # kanban_parser directly so usage/error text reads ``/kanban``. _wrap = argparse.ArgumentParser(prog="/kanban-wrap", add_help=False) _wrap.exit_on_error = False # type: ignore[attr-defined] - _top_sub = _wrap.add_subparsers(dest="_top") - kanban_parser = build_parser(_top_sub) + kanban_parser = build_parser(_wrap.add_subparsers(dest="_top")) kanban_parser.prog = "/kanban" kanban_parser.exit_on_error = False # type: ignore[attr-defined] - for _action in kanban_parser._actions: - if isinstance(_action, argparse._SubParsersAction): - for _name, _choice in _action.choices.items(): - _choice.prog = f"/kanban {_name}" - _choice.exit_on_error = False # type: ignore[attr-defined] + subparsers = [a for a in kanban_parser._actions if isinstance(a, argparse._SubParsersAction)] + for _action in subparsers: + for _name, _choice in _action.choices.items(): + _choice.prog = f"/kanban {_name}" + _choice.exit_on_error = False # type: ignore[attr-defined] def _usage_for_error() -> str: if tokens: - for _action in kanban_parser._actions: - if isinstance(_action, argparse._SubParsersAction): - subparser = _action.choices.get(tokens[0]) - if subparser is not None: - return subparser.format_usage().rstrip() + for _action in subparsers: + subparser = _action.choices.get(tokens[0]) + if subparser is not None: + return subparser.format_usage().rstrip() return kanban_parser.format_usage().rstrip() buf_out = io.StringIO()