diff --git a/AGENTS.md b/AGENTS.md index 63247b2bf4744..453bfcc464785 100644 --- a/AGENTS.md +++ b/AGENTS.md @@ -1122,6 +1122,10 @@ Isolation model: - After `kanban.failure_limit` consecutive non-success attempts on the same task (default: 2), the dispatcher auto-blocks it to prevent spin loops. +- Opt-in exception: with `kanban.retriage_on_timeout: true`, a breaker + trip caused by consecutive **timeouts** sends the task back to Triage + for decomposition (once per task) instead of blocking it — timeouts + are deterministic, so splitting beats blind retries. Full user-facing docs: `website/docs/user-guide/features/kanban.md`. diff --git a/gateway/kanban_watchers.py b/gateway/kanban_watchers.py index eb1c68ffd66b9..726ac1bc0839b 100644 --- a/gateway/kanban_watchers.py +++ b/gateway/kanban_watchers.py @@ -164,7 +164,11 @@ async def _kanban_notifier_watcher(self, interval: float = 5.0) -> None: # "status" covers dashboard drag-drop and `_set_status_direct()` # writes — surface those transitions to subscribers too. - TERMINAL_KINDS = ("completed", "blocked", "gave_up", "crashed", "timed_out", "status", "archived", "unblocked") + # "retriaged" (kanban.retriage_on_timeout) is delivered so a + # subscriber who just saw the task's `timed_out` event also sees + # that the dispatcher recovered it into Triage for decomposition + # rather than giving up. + TERMINAL_KINDS = ("completed", "blocked", "gave_up", "crashed", "timed_out", "retriaged", "status", "archived", "unblocked") # Subscriptions are removed only when the task reaches a truly final # status (done / archived). We used to also unsub on any terminal # event kind (gave_up / crashed / timed_out / blocked), but that @@ -874,6 +878,23 @@ async def _kanban_dispatcher_watcher(self) -> None: ) failure_limit = _kb.DEFAULT_FAILURE_LIMIT + # Retriage-on-timeout: opt-in — a task whose breaker trips on + # consecutive timeouts goes back to Triage for decomposition + # instead of blocking (see kanban_db._record_task_failure). + retriage_on_timeout = bool(kanban_cfg.get("retriage_on_timeout", False)) + if retriage_on_timeout: + logger.info("kanban dispatcher: retriage_on_timeout enabled") + _ad_enabled, _ = _resolve_auto_decompose_settings(_load_config) + if not _ad_enabled: + # Not fatal — manual `hermes kanban decompose` still works — + # but without auto-decompose a retriaged task sits in Triage + # until someone acts, which is easy to miss. + logger.warning( + "kanban dispatcher: retriage_on_timeout is enabled but " + "auto_decompose is disabled — retriaged tasks will wait " + "in Triage for a manual decompose" + ) + # Read stale_timeout_seconds — 0 disables stale detection. raw_stale = kanban_cfg.get("dispatch_stale_timeout_seconds", 0) try: @@ -1022,6 +1043,7 @@ def _tick_once_for_board(slug: str) -> "Optional[object]": stale_timeout_seconds=stale_timeout_seconds, default_assignee=default_assignee, max_in_progress_per_profile=max_in_progress_per_profile, + retriage_on_timeout=retriage_on_timeout, ) except sqlite3.DatabaseError as exc: if _is_corrupt_board_db_error(exc): diff --git a/hermes_cli/config.py b/hermes_cli/config.py index 060373555f0b9..006943ae403ef 100644 --- a/hermes_cli/config.py +++ b/hermes_cli/config.py @@ -2786,6 +2786,17 @@ def _ensure_hermes_home_managed(home: Path): # same task/profile (spawn_failed, timed_out, or crashed). Reassignment # resets the streak for the new profile. "failure_limit": 2, + # When true, a task whose circuit breaker trips on consecutive + # TIMEOUTS is sent back to Triage for decomposition (with a + # failure-context block appended to its body) instead of being + # auto-blocked. A timeout is deterministic — blind retries of a + # task that needs more than max_runtime_seconds fail identically + # — so subdivision is the productive recovery. At most one + # retriage per task; the second trip blocks normally. Works best + # with auto_decompose: true (otherwise the task waits in Triage + # for a manual `hermes kanban decompose`). Crashes and spawn + # failures never retriage. Default off — opt-in behavior change. + "retriage_on_timeout": False, # Worker stdout/stderr logs rotate at spawn time. Defaults preserve # the historical 2 MiB + one-backup behavior; long-running workers can # raise these to keep more early failure evidence. diff --git a/hermes_cli/kanban.py b/hermes_cli/kanban.py index 6f0766a738d1f..99cdd60430a7e 100644 --- a/hermes_cli/kanban.py +++ b/hermes_cli/kanban.py @@ -320,7 +320,13 @@ def build_parser(parent_subparsers: argparse._SubParsersAction) -> argparse.Argu "worktree under the project's primary repo with a " "deterministic branch. See `hermes project list`.") p_create.add_argument("--tenant", default=None, help="Tenant namespace") - p_create.add_argument("--priority", type=int, default=0, help="Priority tiebreaker") + p_create.add_argument( + "--priority", type=int, default=None, + help=( + "Priority tiebreaker (higher = picked sooner). Omitted: inherit " + "the highest parent priority when --parent is given, else 0." + ), + ) p_create.add_argument("--triage", action="store_true", help="Park in triage — a specifier will flesh out the spec and promote to todo") p_create.add_argument("--idempotency-key", default=None, @@ -2246,11 +2252,13 @@ def _coerce_positive_int(value): max_spawn = cli_max if cli_max is not None else _coerce_positive_int( _kanban_cfg.get("max_spawn") ) + retriage_on_timeout = bool(_kanban_cfg.get("retriage_on_timeout", False)) except Exception: default_assignee = None max_in_progress_per_profile = None max_in_progress = None max_spawn = getattr(args, "max", None) + retriage_on_timeout = False with kb.connect_closing() as conn: res = kb.dispatch_once( conn, @@ -2260,12 +2268,14 @@ def _coerce_positive_int(value): 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, + retriage_on_timeout=retriage_on_timeout, ) if getattr(args, "json", False): print(json.dumps({ "reclaimed": res.reclaimed, "crashed": res.crashed, "timed_out": res.timed_out, + "retriaged": res.retriaged, "stale": res.stale, "auto_blocked": res.auto_blocked, "promoted": res.promoted, @@ -2289,6 +2299,9 @@ def _coerce_positive_int(value): print(f"Timed out: {len(res.timed_out)}") if res.timed_out: print(f" {', '.join(res.timed_out)}") + if res.retriaged: + print(f"Retriaged: {len(res.retriaged)}") + print(f" {', '.join(res.retriaged)}") print(f"Stale: {len(res.stale)}") if res.stale: print(f" {', '.join(res.stale)}") diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index de332f36ee44e..bb2ad6ccd11ca 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -2395,7 +2395,7 @@ def create_task( workspace_path: Optional[str] = None, branch_name: Optional[str] = None, tenant: Optional[str] = None, - priority: int = 0, + priority: Optional[int] = None, parents: Iterable[str] = (), triage: bool = False, idempotency_key: Optional[str] = None, @@ -2431,6 +2431,13 @@ def create_task( each name to ``hermes --skills ...``. Use this to pin a task to a specialist skill (e.g. ``skills=["translation"]`` so the worker loads the translation skill regardless of the profile's default config). + + ``priority`` defaults to ``None``: when the caller does not specify + a priority and the task has parents, the child inherits the highest + parent priority (a child must be processed at least as urgently as + any parent it gates). With no parents — or parents whose priority is + NULL — the schema default 0 applies. An explicit priority always + wins; parent priorities are never modified. """ assignee = _canonical_assignee(assignee) if not title or not title.strip(): @@ -2629,6 +2636,33 @@ def create_task( except Exception: branch_name = None + # Priority inheritance: when the caller left priority + # unspecified and the task has parents, inherit the + # HIGHEST parent priority — a child that gates a + # high-priority parent must itself be picked sooner. + # Parents with NULL priority contribute nothing; if all + # are NULL (or there are no parents) the schema default + # 0 applies. The parents' own priorities are never + # modified here. + effective_priority = priority + if effective_priority is None and parents: + _p_rows = conn.execute( + "SELECT priority FROM tasks WHERE id IN (" + + ",".join("?" * len(parents)) + + ")", + parents, + ).fetchall() + effective_priority = max( + ( + r["priority"] + for r in _p_rows + if r["priority"] is not None + ), + default=0, + ) + if effective_priority is None: + effective_priority = 0 + conn.execute( """ INSERT INTO tasks ( @@ -2645,7 +2679,7 @@ def create_task( body, assignee, task_status, - priority, + effective_priority, created_by, now, workspace_kind, @@ -4584,6 +4618,47 @@ def _is_managed_scratch_path(p: Path) -> bool: return is_managed +def _cleanup_auto_commit_worktree(path: str, task_id: str) -> None: + """Resolve the actual worktree checkout dir and auto-commit changes. + + The DB stores the anchor *repo* path for worktree tasks, but the + actual checkout lives at ``/.worktrees/`` (or is + itself a linked checkout). Find the real dir and delegate to + ``_git_auto_commit_worktree``. + """ + anchor = Path(path).expanduser().resolve(strict=False) + # If the anchor itself is already a linked worktree checkout, use it. + try: + if anchor.exists() and _is_linked_worktree_checkout(anchor): + _git_auto_commit_worktree( + str(anchor), task_id, + reason="auto-commit on completion", + ) + return + except Exception: + pass + # Standard convention: worktree at /.worktrees/ + candidate = anchor / ".worktrees" / task_id + try: + if candidate.is_dir(): + _git_auto_commit_worktree( + str(candidate), task_id, + reason="auto-commit on completion", + ) + return + except Exception: + pass + # Fallback: try the anchor path itself (may be a worktree-like checkout) + try: + if anchor.is_dir(): + _git_auto_commit_worktree( + str(anchor), task_id, + reason="auto-commit on completion", + ) + except Exception: + pass + + def _cleanup_workspace(conn: sqlite3.Connection, task_id: str) -> None: """Remove a task's scratch workspace dir and kill its stale tmux session. @@ -4602,6 +4677,13 @@ def _cleanup_workspace(conn: sqlite3.Connection, task_id: str) -> None: kind: Optional[str] = row["workspace_kind"] path: Optional[str] = row["workspace_path"] if kind != "scratch" or not path: + # Worktree: auto-commit uncommitted changes so the worker's + # output survives into git history even if the worker forgot + # to commit (or timed out before doing so). The DB stores + # the anchor repo path, not the actual checkout, so we must + # resolve the worktree dir. + if kind == "worktree" and path: + _cleanup_auto_commit_worktree(path, task_id) # This task's own workspace isn't a removable scratch dir, but its # completion may still unblock a deferred parent scratch cleanup # (e.g. a 'dir' child whose scratch parent was waiting on it). #33774 @@ -5347,6 +5429,12 @@ def decompose_triage_task( - The root task is not in ``triage`` - A cycle would result (caller built a bad graph) + Children inherit the root task's priority exactly: a root with + ``priority=NULL`` yields the schema default 0, a root with a + concrete priority propagates that value to every child. The root's + own priority is never modified here (only ``status``/``assignee`` + are touched on the flip to ``todo``). + Validation of titles/assignees happens inside the same write_txn as the inserts so a malformed entry aborts the whole decomposition cleanly (no orphan children). @@ -5409,7 +5497,7 @@ def decompose_triage_task( child_ids: list[str] = [] with write_txn(conn): root_row = conn.execute( - "SELECT id, status, tenant, workspace_kind, workspace_path " + "SELECT id, status, tenant, workspace_kind, workspace_path, priority " "FROM tasks WHERE id = ?", (task_id,), ).fetchone() @@ -5448,14 +5536,18 @@ def decompose_triage_task( child_ws_path = None conn.execute( "INSERT INTO tasks " - "(id, title, body, assignee, status, workspace_kind, " + "(id, title, body, assignee, status, priority, workspace_kind, " " workspace_path, tenant, created_at, created_by) " - "VALUES (?, ?, ?, ?, 'todo', ?, ?, ?, ?, ?)", + "VALUES (?, ?, ?, ?, 'todo', ?, ?, ?, ?, ?, ?)", ( new_id, title, body if isinstance(body, str) else None, assignee, + # Children inherit the root's priority exactly. A + # NULL root priority maps to the schema default 0 — + # NOT a NULL insert, which would bypass the DEFAULT. + root_row["priority"] if root_row["priority"] is not None else 0, child_ws_kind, child_ws_path, tenant, @@ -5736,6 +5828,67 @@ def _repo_root_for_worktree_target(path: Path) -> Optional[Path]: current = current.parent +def _git_auto_commit_worktree( + workspace_path: Optional[str], + task_id: str, + *, + reason: str = "auto-commit on completion", +) -> None: + """Best-effort auto-commit of uncommitted changes in a git worktree. + + Checks if ``workspace_path`` is inside a git repo and has uncommitted + changes. If so, runs ``git add -A && git commit`` with a descriptive + message. Designed as a safety net for worktree kanban tasks whose + worker forgot (or timed out before) committing. + + Always best-effort — errors are logged and swallowed so the caller's + flow (completion, reclaim) is never blocked. + """ + if not workspace_path: + return + wdir = Path(workspace_path).expanduser() + if not wdir.is_dir(): + return + # Is the workspace inside a git repo? + toplevel = _git_toplevel(wdir) + if toplevel is None: + return + # Check for uncommitted changes (tracked or untracked). + try: + result = subprocess.run( + ["git", "-C", str(wdir), "status", "--porcelain"], + capture_output=True, text=True, timeout=30, check=False, + ) + except Exception: + return + if result.returncode != 0 or not (result.stdout or "").strip(): + return # clean — nothing to commit + lines = (result.stdout or "").strip().splitlines() + if not lines: + return + # We have uncommitted changes — commit them. + branch = _git_current_branch(wdir) or "unknown" + msg = f"wt/{task_id}: {reason} ({len(lines)} file{'s' if len(lines) != 1 else ''} dirty)" + try: + subprocess.run( + ["git", "-C", str(wdir), "add", "-A"], + capture_output=True, text=True, timeout=60, check=True, + ) + subprocess.run( + ["git", "-C", str(wdir), "commit", "-m", msg], + capture_output=True, text=True, timeout=60, check=False, + ) + _log.info( + "auto-committed %d dirty file(s) in worktree %s for task %s: %s", + len(lines), wdir, task_id, msg, + ) + except Exception as exc: + _log.warning( + "auto-commit failed for worktree %s task %s: %s", + wdir, task_id, exc, + ) + + def _ensure_git_worktree(repo_root: Path, target: Path, branch_name: str) -> None: """Materialize ``target`` as a linked git worktree under ``repo_root``.""" target = target.expanduser() @@ -6059,6 +6212,16 @@ class DispatchResult: """Task ids auto-blocked by the spawn-failure circuit breaker.""" timed_out: list[str] = field(default_factory=list) """Task ids whose workers exceeded ``max_runtime_seconds``.""" + retriaged: list[str] = field(default_factory=list) + """Task ids sent back to ``triage`` for decomposition instead of + being auto-blocked, because they kept timing out and + ``kanban.retriage_on_timeout`` is enabled. The auto-decomposer picks + them up on its next tick and splits them into smaller children. + + NOTE: these ids also appear in ``timed_out`` (the timeout really + happened; retriage is what the dispatcher did about it). Consumers + counting hard failures should treat ``timed_out`` minus ``retriaged`` + as the unrecovered set.""" stale: list[str] = field(default_factory=list) """Task ids reclaimed because no progress (heartbeat) was seen within ``dispatch_stale_timeout_seconds``.""" @@ -6421,6 +6584,7 @@ def enforce_max_runtime( conn: sqlite3.Connection, *, signal_fn=None, + retriage_on_timeout: bool = False, ) -> list[str]: """Terminate workers whose per-task ``max_runtime_seconds`` has elapsed. @@ -6430,19 +6594,29 @@ def enforce_max_runtime( breaker has already given up, in which case the task stays blocked where ``_record_spawn_failure`` parked it. + ``retriage_on_timeout`` is forwarded to ``_record_task_failure``: + when enabled, a task whose breaker trips on a timeout goes back to + ``triage`` for decomposition instead of ``blocked`` (see the + ``_record_task_failure`` docstring). Task ids retriaged this call + are stashed on ``enforce_max_runtime._last_retriaged`` (same + side-channel pattern as ``detect_crashed_workers``) so + ``dispatch_once`` can surface them in ``DispatchResult.retriaged``. + Runs host-local: only tasks claimed by this host are candidates (same reasoning as ``detect_crashed_workers``). ``signal_fn`` is a test hook; defaults to ``os.kill`` on POSIX. """ import signal timed_out: list[str] = [] + retriaged: list[str] = [] now = int(time.time()) host_prefix = f"{_claimer_id().split(':', 1)[0]}:" rows = conn.execute( "SELECT t.id, t.worker_pid, " " COALESCE(r.started_at, t.started_at) AS active_started_at, " - " t.max_runtime_seconds, t.claim_lock " + " t.max_runtime_seconds, t.claim_lock, " + " t.workspace_kind, t.workspace_path " "FROM tasks t " "LEFT JOIN task_runs r ON r.id = t.current_run_id " "WHERE t.status = 'running' AND t.max_runtime_seconds IS NOT NULL " @@ -6462,6 +6636,22 @@ def enforce_max_runtime( pid = int(row["worker_pid"]) tid = row["id"] + # Best-effort auto-commit worktree changes before killing the worker. + # The worker's changes are on disk but uncommitted — this ensures + # they survive into git history even on timeout. + _ws_kind: Optional[str] = row["workspace_kind"] + _ws_path: Optional[str] = row["workspace_path"] + if _ws_kind == "worktree" and _ws_path: + try: + _git_auto_commit_worktree( + _ws_path, tid, + reason="auto-commit on timeout", + ) + except Exception: + _log.warning( + "auto-commit failed on timeout for worktree task %s", tid, + exc_info=True, + ) # SIGTERM then SIGKILL. Keep it simple: 5 s grace. Workers that # want a cleaner shutdown can install their own SIGTERM handler # before the grace expires. @@ -6527,7 +6717,19 @@ def enforce_max_runtime( release_claim=False, end_run=False, event_payload_extra={"pid": pid, "sigkill": killed}, + retriage_on_timeout=retriage_on_timeout, ) + if retriage_on_timeout: + status_row = conn.execute( + "SELECT status FROM tasks WHERE id = ?", (tid,), + ).fetchone() + if status_row is not None and status_row["status"] == "triage": + retriaged.append(tid) + # Side-channel for dispatch_once — mirrors the + # ``detect_crashed_workers._last_auto_blocked`` pattern so the + # public list-of-timed-out return type stays stable for direct + # callers and existing tests. + enforce_max_runtime._last_retriaged = retriaged # type: ignore[attr-defined] return timed_out @@ -7019,6 +7221,105 @@ def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]: return crashed +def _was_retriaged(conn: sqlite3.Connection, task_id: str) -> bool: + """True when the task already went through retriage-on-timeout once. + + The ``retriaged`` event is the durable marker — no schema change + needed, and the audit trail doubles as the loop guard. Must be + called inside the caller's ``write_txn`` (plain SELECT, no writes). + """ + row = conn.execute( + "SELECT COUNT(*) FROM task_events " + "WHERE task_id = ? AND kind = 'retriaged'", + (task_id,), + ).fetchone() + return int(row[0]) > 0 + + +# Appended to a task body when retriage-on-timeout sends it back to the +# Triage column. Written for the decomposer LLM (kanban_decompose reads +# title + body): states what happened and biases the decomposition +# toward children that fit the runtime budget. +_RETRIAGE_BODY_TEMPLATE = """ + +--- +## Automated retriage after timeout + +{failures} consecutive attempt(s) at this task timed out ({error}). +Instead of giving up, the dispatcher moved it back to Triage for +decomposition. + +Decomposer: split the ORIGINAL task above into smaller child tasks that +are each independently completable{budget_clause}. Preserve the original +intent and acceptance criteria. Do NOT restate the task as a single +child — the whole point is that one worker could not finish it within +its runtime budget. +""" + + +def _retriage_task( + conn: sqlite3.Connection, + task_id: str, + *, + body: Optional[str], + max_runtime_seconds: Optional[int], + failures: int, + effective_limit: int, + limit_source: str, + error: str, + event_payload_extra: Optional[dict] = None, +) -> bool: + """Send a repeatedly-timed-out task back to ``triage`` for decomposition. + + Appends a failure-context block to the body (the decomposer LLM sees + it), resets the consecutive-failures counter so the decomposed root + gets a fresh breaker budget for its judge run, and emits a + ``retriaged`` event — which is also the once-only guard + ``_was_retriaged`` checks. Must be called inside the caller's + ``write_txn`` (this helper does NOT open its own). + + Returns True when the task was actually retriaged. The UPDATE matches + ``status = 'ready'`` ONLY: the timeout path parked the task there in + its own txn, and if another dispatcher claimed it back to ``running`` + in the window between that txn and this one, retriaging would yank a + freshly-spawned live worker off the task. Losing the race returns + False (no event, no body change) and the caller falls through to the + plain block semantics, which is exactly the pre-existing behavior for + that window. + """ + if max_runtime_seconds: + budget_clause = ( + f" well within {int(max_runtime_seconds)} seconds of work" + ) + else: + budget_clause = "" + new_body = (body or "") + _RETRIAGE_BODY_TEMPLATE.format( + failures=failures, + error=error, + budget_clause=budget_clause, + ) + cur = conn.execute( + "UPDATE tasks SET status = 'triage', body = ?, " + "consecutive_failures = 0, last_failure_error = ?, " + "claim_lock = NULL, claim_expires = NULL, worker_pid = NULL " + "WHERE id = ? AND status = 'ready'", + (new_body, error, task_id), + ) + if cur.rowcount != 1: + return False + payload = { + "failures": failures, + "effective_limit": effective_limit, + "limit_source": limit_source, + "error": error, + "trigger_outcome": "timed_out", + } + if event_payload_extra: + payload.update(event_payload_extra) + _append_event(conn, task_id, "retriaged", payload) + return True + + def _record_task_failure( conn: sqlite3.Connection, task_id: str, @@ -7030,6 +7331,7 @@ def _record_task_failure( release_claim: bool = False, end_run: bool = False, event_payload_extra: Optional[dict] = None, + retriage_on_timeout: bool = False, ) -> bool: """Record a non-success outcome (spawn_failed / crashed / timed_out) and maybe trip the circuit breaker. @@ -7073,13 +7375,29 @@ def _record_task_failure( ``detect_crashed_workers``, which resolves the per-task ``max_retries`` override against the violation streak itself. The failure is still counted into ``consecutive_failures``. + + ``retriage_on_timeout=True`` (config ``kanban.retriage_on_timeout``, + default off) changes what happens when the breaker trips with + ``outcome="timed_out"``: instead of parking the task in ``blocked`` + with a ``gave_up`` event, the task is sent back to ``triage`` with a + failure-context block appended to its body and its failure counter + reset. The auto-decomposer (``kanban.auto_decompose``) then splits it + into smaller child tasks on its next tick. Rationale: a timeout is a + *deterministic* failure — a task that needs more than + ``max_runtime_seconds`` fails identically on every blind retry, so + the productive recovery is subdivision, not repetition. Each task is + retriaged at most ONCE (guarded by the ``retriaged`` event); a second + trip blocks normally, so subdivision cannot recurse unbounded. + Non-timeout outcomes (crash / spawn_failed — plausibly transient) + are never retriaged. """ if failure_limit is None: failure_limit = DEFAULT_FAILURE_LIMIT blocked = False with write_txn(conn): row = conn.execute( - "SELECT consecutive_failures, status, max_retries " + "SELECT consecutive_failures, status, max_retries, body, " + "max_runtime_seconds " "FROM tasks WHERE id = ?", (task_id,), ).fetchone() if row is None: @@ -7100,6 +7418,32 @@ def _record_task_failure( limit_source = "dispatcher" if force_trip or failures >= effective_limit: + # Retriage-on-timeout: opt-in alternative to giving up when + # the trip was caused by a runtime timeout. ``force_trip`` + # callers already applied their own retry policy, so they + # keep the plain block semantics. + if ( + retriage_on_timeout + and not force_trip + and outcome == "timed_out" + and not _was_retriaged(conn, task_id) + ): + if _retriage_task( + conn, task_id, + body=row["body"], + max_runtime_seconds=( + row["max_runtime_seconds"] + if "max_runtime_seconds" in row.keys() else None + ), + failures=failures, + effective_limit=effective_limit, + limit_source=limit_source, + error=error[:500], + event_payload_extra=event_payload_extra, + ): + return False + # Lost the ready→running race to a concurrent claimer — + # fall through to the plain block semantics below. # Trip the breaker. if release_claim: # Spawn path: still running, also clear claim state. @@ -7449,6 +7793,7 @@ def dispatch_once( board: Optional[str] = None, default_assignee: Optional[str] = None, max_in_progress_per_profile: Optional[int] = None, + retriage_on_timeout: bool = False, ) -> DispatchResult: """Run one dispatcher tick under the board's single-writer lock. @@ -7483,6 +7828,7 @@ def dispatch_once( board=board, default_assignee=default_assignee, max_in_progress_per_profile=max_in_progress_per_profile, + retriage_on_timeout=retriage_on_timeout, ) with _dispatch_tick_lock(db_path) as held: if not held: @@ -7499,6 +7845,7 @@ def dispatch_once( board=board, default_assignee=default_assignee, max_in_progress_per_profile=max_in_progress_per_profile, + retriage_on_timeout=retriage_on_timeout, ) @@ -7515,6 +7862,7 @@ def _dispatch_once_locked( board: Optional[str] = None, default_assignee: Optional[str] = None, max_in_progress_per_profile: Optional[int] = None, + retriage_on_timeout: bool = False, ) -> DispatchResult: """Run one dispatcher tick. @@ -7570,7 +7918,14 @@ def _dispatch_once_locked( ) if _crash_rate_limited: result.rate_limited.extend(_crash_rate_limited) - result.timed_out = enforce_max_runtime(conn) + result.timed_out = enforce_max_runtime( + conn, retriage_on_timeout=retriage_on_timeout, + ) + # Tasks sent back to triage instead of blocked (retriage-on-timeout). + # Same side-channel pattern as detect_crashed_workers above. + result.retriaged = list( + getattr(enforce_max_runtime, "_last_retriaged", []) or [] + ) result.promoted = recompute_ready(conn, failure_limit=failure_limit) # Count tasks already running so max_spawn enforces concurrency rather diff --git a/tests/hermes_cli/test_kanban_db.py b/tests/hermes_cli/test_kanban_db.py index 25ed7223129f5..64943832fd6dd 100644 --- a/tests/hermes_cli/test_kanban_db.py +++ b/tests/hermes_cli/test_kanban_db.py @@ -236,6 +236,61 @@ def test_create_task_unknown_parent_errors(kanban_home): kb.create_task(conn, title="orphan", parents=["t_ghost"]) +def test_create_task_child_inherits_parent_priority(kanban_home): + with kb.connect() as conn: + p = kb.create_task(conn, title="parent", priority=7) + c = kb.create_task(conn, title="child", parents=[p]) + assert kb.get_task(conn, c).priority == 7 + + +def test_create_task_child_inherits_highest_parent_priority(kanban_home): + with kb.connect() as conn: + p1 = kb.create_task(conn, title="low parent", priority=3) + p2 = kb.create_task(conn, title="high parent", priority=9) + c = kb.create_task(conn, title="child", parents=[p1, p2]) + assert kb.get_task(conn, c).priority == 9 + + +def test_create_task_no_priority_without_parents_defaults_zero(kanban_home): + with kb.connect() as conn: + t = kb.create_task(conn, title="standalone") + assert kb.get_task(conn, t).priority == 0 + + +def test_create_task_null_priority_normalizes_to_zero(kanban_home): + # Passing priority=None explicitly is normalized to the schema + # default 0 (the old code stored NULL, which the dispatcher treats + # like a missing priority anyway). + with kb.connect() as conn: + t = kb.create_task(conn, title="standalone", priority=None) + assert kb.get_task(conn, t).priority == 0 + + +def test_create_task_child_of_null_priority_parent_yields_zero(kanban_home): + # A parent whose DB row carries NULL priority (e.g. legacy rows) + # must not leak NULL into the child — the default 0 applies. + with kb.connect() as conn: + p = kb.create_task(conn, title="parent", priority=3) + conn.execute("UPDATE tasks SET priority = NULL WHERE id = ?", (p,)) + c = kb.create_task(conn, title="child", parents=[p]) + assert kb.get_task(conn, p).priority is None + assert kb.get_task(conn, c).priority == 0 + + +def test_create_task_explicit_priority_wins_over_inheritance(kanban_home): + with kb.connect() as conn: + p = kb.create_task(conn, title="parent", priority=7) + c = kb.create_task(conn, title="child", parents=[p], priority=2) + assert kb.get_task(conn, c).priority == 2 + + +def test_create_task_child_creation_never_modifies_parent_priority(kanban_home): + with kb.connect() as conn: + p = kb.create_task(conn, title="parent", priority=5) + kb.create_task(conn, title="child", parents=[p]) + assert kb.get_task(conn, p).priority == 5 + + def test_workspace_kind_validation(kanban_home): with kb.connect() as conn, pytest.raises(ValueError, match="workspace_kind"): kb.create_task(conn, title="bad ws", workspace_kind="cloud") @@ -4959,3 +5014,61 @@ def test_bare_connect_does_not_close_on_context_exit(tmp_path): # Still usable after with-block exit (the leak). conn.execute("SELECT 1").fetchone() conn.close() # explicit close to avoid leaking THIS test + + +def test_complete_task_auto_commits_worktree(kanban_home, tmp_path): + """complete_task auto-commits the worktree when the worker left files dirty. + + Regression guard: workers that finish their work (build, test) but forget + to commit should have their work auto-committed on completion, so the + rescue pattern (Claude: git add -A && commit) becomes unnecessary. + """ + # ── 1. Create a git repo and a worktree task ── + repo = tmp_path / "repo" + _init_git_repo(repo) + with kb.connect() as conn: + tid = kb.create_task( + conn, + title="test auto-commit", + workspace_kind="worktree", + workspace_path=str(repo), + ) + task = kb.get_task(conn, tid) + assert task is not None + ws = kb.resolve_workspace(task) + + # ── 2. Write an uncommitted file in the worktree (worker forgot to commit) ── + dirty_file = ws / "feature.py" + dirty_file.write_text("print('hello')\n") + # Verify: git sees it as untracked + status = subprocess.run( + ["git", "-C", str(ws), "status", "--porcelain"], + capture_output=True, text=True, check=True, + ) + assert "feature.py" in status.stdout, ( + "prerequisite failed: file should be dirty before complete_task" + ) + + # ── 3. Complete the task ── + with kb.connect() as conn: + ok = kb.complete_task(conn, tid, result="done") + assert ok, "complete_task should succeed" + + # ── 4. Verify: worktree is now clean (auto-commit happened) ── + status = subprocess.run( + ["git", "-C", str(ws), "status", "--porcelain"], + capture_output=True, text=True, check=True, + ) + assert status.stdout.strip() == "", ( + "worktree should be clean after complete_task, " + f"but got: {status.stdout.strip()!r}" + ) + # Verify the commit exists with our expected message + log = subprocess.run( + ["git", "-C", str(ws), "log", "--oneline", "-1"], + capture_output=True, text=True, check=True, + ) + assert f"wt/{tid}" in log.stdout, ( + f"expected commit message to contain wt/{tid}, " + f"got: {log.stdout.strip()!r}" + ) diff --git a/tests/hermes_cli/test_kanban_decompose_db.py b/tests/hermes_cli/test_kanban_decompose_db.py index f0f11fb82fcc9..ce375786d3f8a 100644 --- a/tests/hermes_cli/test_kanban_decompose_db.py +++ b/tests/hermes_cli/test_kanban_decompose_db.py @@ -21,7 +21,8 @@ def kanban_home(tmp_path, monkeypatch): return home -def _create_triage(conn, title="rough idea", body=None, assignee=None, tenant=None): +def _create_triage(conn, title="rough idea", body=None, assignee=None, tenant=None, + priority=None): return kb.create_task( conn, title=title, @@ -29,6 +30,7 @@ def _create_triage(conn, title="rough idea", body=None, assignee=None, tenant=No assignee=assignee, tenant=tenant, triage=True, + priority=priority, ) @@ -68,6 +70,53 @@ def test_decompose_creates_children_and_promotes_root(kanban_home): assert c1.assignee == "engineer" +def test_decompose_children_inherit_root_priority(kanban_home): + with kb.connect() as conn: + tid = _create_triage(conn, title="prioritized idea", priority=7) + assert kb.get_task(conn, tid).priority == 7 + + children = [ + {"title": "research", "assignee": "researcher", "parents": []}, + {"title": "build it", "assignee": "engineer", "parents": [0]}, + ] + with kb.connect() as conn: + child_ids = kb.decompose_triage_task( + conn, tid, root_assignee="orchestrator", children=children, + ) + assert child_ids is not None + with kb.connect() as conn: + root = kb.get_task(conn, tid) + # Every child carries the root's exact priority… + for cid in child_ids: + assert kb.get_task(conn, cid).priority == 7 + # …and the root's own priority is untouched by the fan-out. + assert root.priority == 7 + + +def test_decompose_children_default_zero_when_root_priority_null(kanban_home): + # Legacy/edge rows can carry priority=NULL in the DB (only reachable + # via raw SQL — create_task normalizes None to 0). The fan-out must + # not leak NULL into children: it maps to the schema default 0. + with kb.connect() as conn: + tid = _create_triage(conn, title="unprioritized idea", priority=5) + conn.execute("UPDATE tasks SET priority = NULL WHERE id = ?", (tid,)) + assert kb.get_task(conn, tid).priority is None + + children = [ + {"title": "research", "assignee": "researcher", "parents": []}, + ] + with kb.connect() as conn: + child_ids = kb.decompose_triage_task( + conn, tid, root_assignee="orchestrator", children=children, + ) + assert child_ids is not None + with kb.connect() as conn: + # NULL root priority must NOT insert NULL children (which would + # bypass the schema DEFAULT) — it maps to the default 0. + for cid in child_ids: + assert kb.get_task(conn, cid).priority == 0 + + def test_decompose_returns_none_when_task_missing(kanban_home): with kb.connect() as conn: result = kb.decompose_triage_task( diff --git a/tests/hermes_cli/test_kanban_retriage_on_timeout.py b/tests/hermes_cli/test_kanban_retriage_on_timeout.py new file mode 100644 index 0000000000000..9af41f527df00 --- /dev/null +++ b/tests/hermes_cli/test_kanban_retriage_on_timeout.py @@ -0,0 +1,391 @@ +"""Tests for retriage-on-timeout (``kanban.retriage_on_timeout``). + +When enabled, a task whose circuit breaker trips on consecutive TIMEOUTS +is sent back to ``triage`` for decomposition (failure-context block +appended to its body, counter reset, ``retriaged`` event) instead of +being parked in ``blocked`` with a ``gave_up`` event. Rationale: a +timeout is deterministic — a task that needs more than +``max_runtime_seconds`` fails identically on every blind retry — so +subdivision via the existing auto-decompose pipeline is the productive +recovery, not repetition. + +Covers: + * default-off keeps the historical gave_up behavior byte-for-byte + * the retriage transition itself (status, body block, counter reset, + claim cleanup, event payload) + * the once-only guard (second trip blocks normally) + * non-timeout outcomes (crash / spawn_failed) never retriage + * ``force_trip`` callers keep plain block semantics + * per-task ``max_retries`` override interaction + * compatibility with ``decompose_triage_task`` (the retriaged task is + a valid decomposition root) + * ``dispatch_once`` plumbing end-to-end + ``DispatchResult.retriaged`` +""" + +from __future__ import annotations + +import os +import time +from pathlib import Path + +import pytest + +from hermes_cli import kanban_db as kb + + +# --------------------------------------------------------------------------- +# Fixtures / helpers +# --------------------------------------------------------------------------- + +@pytest.fixture +def kanban_home(tmp_path, monkeypatch): + home = tmp_path / ".hermes" + home.mkdir() + monkeypatch.setenv("HERMES_HOME", str(home)) + monkeypatch.setenv("HERMES_KANBAN_CRASH_GRACE_SECONDS", "0") + monkeypatch.setattr(Path, "home", lambda: tmp_path) + kb.init_db() + return home + + +def _age_active_run(conn, tid, *, seconds=30): + """Backdate the active run so elapsed > max_runtime_seconds.""" + old_started = int(time.time()) - seconds + with kb.write_txn(conn): + conn.execute( + "UPDATE tasks SET started_at = ? WHERE id = ?", + (old_started, tid), + ) + conn.execute( + "UPDATE task_runs SET started_at = ? " + "WHERE id = (SELECT current_run_id FROM tasks WHERE id = ?)", + (old_started, tid), + ) + + +def _run_one_timeout_cycle(conn, tid, *, retriage_on_timeout): + """Claim the task, backdate its run, and enforce max runtime once.""" + kb.claim_task(conn, tid) + kb._set_worker_pid(conn, tid, os.getpid()) + _age_active_run(conn, tid) + return kb.enforce_max_runtime( + conn, + signal_fn=lambda pid, sig: None, + retriage_on_timeout=retriage_on_timeout, + ) + + +@pytest.fixture +def fast_pid_dead(monkeypatch): + """Pretend SIGTERM works instantly so the grace poll exits fast.""" + monkeypatch.setattr(kb, "_pid_alive", lambda pid: False) + + +# --------------------------------------------------------------------------- +# Default off — historical behavior unchanged +# --------------------------------------------------------------------------- + +def test_default_off_keeps_gave_up_behavior(kanban_home, fast_pid_dead): + """Without the flag, repeated timeouts still block with gave_up.""" + conn = kb.connect() + try: + tid = kb.create_task( + conn, title="long job", assignee="worker", + max_runtime_seconds=1, + ) + for _ in range(2): + _run_one_timeout_cycle(conn, tid, retriage_on_timeout=False) + + task = kb.get_task(conn, tid) + assert task.status == "blocked" + events = kb.list_events(conn, tid) + assert any(e.kind == "gave_up" for e in events) + assert not any(e.kind == "retriaged" for e in events) + finally: + conn.close() + + +def test_config_default_is_off(): + """The config ships with retriage_on_timeout disabled (opt-in).""" + from hermes_cli.config import DEFAULT_CONFIG + assert DEFAULT_CONFIG["kanban"]["retriage_on_timeout"] is False + + +# --------------------------------------------------------------------------- +# The retriage transition +# --------------------------------------------------------------------------- + +def test_retriage_on_breaker_trip(kanban_home, fast_pid_dead): + """Second consecutive timeout retriages instead of blocking.""" + conn = kb.connect() + try: + tid = kb.create_task( + conn, title="long job", assignee="worker", + body="Original spec: do the big thing.", + max_runtime_seconds=1, + ) + # First timeout: below the default limit (2) — plain retry. + _run_one_timeout_cycle(conn, tid, retriage_on_timeout=True) + task = kb.get_task(conn, tid) + assert task.status == "ready" + assert task.consecutive_failures == 1 + + # Second timeout: breaker would trip — retriage instead. + timed_out = _run_one_timeout_cycle(conn, tid, retriage_on_timeout=True) + assert tid in timed_out, "retriaged task still counts as timed out" + assert tid in kb.enforce_max_runtime._last_retriaged + + task = kb.get_task(conn, tid) + assert task.status == "triage" + assert task.consecutive_failures == 0, "fresh breaker budget" + assert task.claim_lock is None + assert task.worker_pid is None + + # Body: original content preserved + failure-context block. + assert task.body.startswith("Original spec: do the big thing.") + assert "Automated retriage after timeout" in task.body + assert "within 1 seconds of work" in task.body + + events = kb.list_events(conn, tid) + assert not any(e.kind == "gave_up" for e in events) + retriaged = [e for e in events if e.kind == "retriaged"] + assert len(retriaged) == 1 + payload = retriaged[0].payload + assert payload["failures"] == 2 + assert payload["effective_limit"] == 2 + assert payload["trigger_outcome"] == "timed_out" + assert "error" in payload + finally: + conn.close() + + +def test_retriage_only_once_then_blocks(kanban_home, fast_pid_dead): + """A task is retriaged at most once; the second trip blocks normally. + + This is the guard against unbounded subdivision: if the decomposed + task somehow lands back in the run cycle and keeps timing out, the + breaker wins. + """ + conn = kb.connect() + try: + tid = kb.create_task( + conn, title="long job", assignee="worker", + max_runtime_seconds=1, + ) + for _ in range(2): + _run_one_timeout_cycle(conn, tid, retriage_on_timeout=True) + assert kb.get_task(conn, tid).status == "triage" + + # Simulate specify/promote putting it back into rotation. + with kb.write_txn(conn): + conn.execute( + "UPDATE tasks SET status = 'ready' WHERE id = ?", (tid,), + ) + + # Counter was reset — two more timeouts to reach the limit again. + for _ in range(2): + _run_one_timeout_cycle(conn, tid, retriage_on_timeout=True) + + task = kb.get_task(conn, tid) + assert task.status == "blocked", "second trip must give up" + events = kb.list_events(conn, tid) + assert sum(1 for e in events if e.kind == "retriaged") == 1 + assert any(e.kind == "gave_up" for e in events) + finally: + conn.close() + + +# --------------------------------------------------------------------------- +# Only deterministic timeouts retriage +# --------------------------------------------------------------------------- + +def test_crash_outcome_never_retriages(kanban_home): + """Crashes are plausibly transient — they keep the block semantics.""" + conn = kb.connect() + try: + tid = kb.create_task(conn, title="crashy", assignee="worker") + with kb.write_txn(conn): + conn.execute( + "UPDATE tasks SET status = 'ready', consecutive_failures = 1 " + "WHERE id = ?", (tid,), + ) + tripped = kb._record_task_failure( + conn, tid, "worker died", + outcome="crashed", + retriage_on_timeout=True, + ) + assert tripped is True + task = kb.get_task(conn, tid) + assert task.status == "blocked" + assert not any( + e.kind == "retriaged" for e in kb.list_events(conn, tid) + ) + finally: + conn.close() + + +def test_force_trip_never_retriages(kanban_home): + """force_trip callers applied their own retry policy — plain block.""" + conn = kb.connect() + try: + tid = kb.create_task(conn, title="violator", assignee="worker") + with kb.write_txn(conn): + conn.execute( + "UPDATE tasks SET status = 'ready' WHERE id = ?", (tid,), + ) + tripped = kb._record_task_failure( + conn, tid, "elapsed 99s > limit 1s", + outcome="timed_out", + force_trip=True, + retriage_on_timeout=True, + ) + assert tripped is True + assert kb.get_task(conn, tid).status == "blocked" + finally: + conn.close() + + +def test_retriage_lost_claim_race_falls_back_to_block(kanban_home): + """If a concurrent dispatcher re-claimed the task (ready → running) + between the timeout txn and the failure-recording txn, retriage must + NOT yank the live worker — it loses the race and the plain block + semantics apply (the pre-existing behavior for that window).""" + conn = kb.connect() + try: + tid = kb.create_task( + conn, title="long job", assignee="worker", + body="original body", + max_runtime_seconds=1, + ) + # Simulate the race: task is back in 'running' hands by the time + # the failure is recorded, with the counter one below the limit. + with kb.write_txn(conn): + conn.execute( + "UPDATE tasks SET status = 'running', " + "consecutive_failures = 1 WHERE id = ?", (tid,), + ) + tripped = kb._record_task_failure( + conn, tid, "elapsed 30s > limit 1s", + outcome="timed_out", + retriage_on_timeout=True, + ) + assert tripped is True, "race loser must fall back to the breaker" + task = kb.get_task(conn, tid) + assert task.status == "blocked" + assert task.body == "original body", "no failure block on race loss" + events = kb.list_events(conn, tid) + assert not any(e.kind == "retriaged" for e in events) + assert any(e.kind == "gave_up" for e in events) + finally: + conn.close() + + +def test_per_task_max_retries_retriages_on_first_timeout( + kanban_home, fast_pid_dead, +): + """max_retries=1 means the very first timeout is the trip → retriage.""" + conn = kb.connect() + try: + tid = kb.create_task( + conn, title="one shot", assignee="worker", + max_runtime_seconds=1, max_retries=1, + ) + _run_one_timeout_cycle(conn, tid, retriage_on_timeout=True) + task = kb.get_task(conn, tid) + assert task.status == "triage" + retriaged = [ + e for e in kb.list_events(conn, tid) if e.kind == "retriaged" + ] + assert retriaged and retriaged[0].payload["effective_limit"] == 1 + assert retriaged[0].payload["limit_source"] == "task" + finally: + conn.close() + + +# --------------------------------------------------------------------------- +# Decompose pipeline compatibility +# --------------------------------------------------------------------------- + +def test_retriaged_task_is_a_valid_decompose_root(kanban_home, fast_pid_dead): + """After retriage the task decomposes exactly like a created-in-triage + one: children get created, the root flips to todo, and the failure + block stays in the root body for the judge run.""" + conn = kb.connect() + try: + tid = kb.create_task( + conn, title="big job", assignee="worker", + body="Do A then B then C.", + max_runtime_seconds=1, + ) + for _ in range(2): + _run_one_timeout_cycle(conn, tid, retriage_on_timeout=True) + assert kb.get_task(conn, tid).status == "triage" + + child_ids = kb.decompose_triage_task( + conn, tid, + root_assignee="orchestrator", + children=[ + {"title": "Do A", "assignee": "worker", "parents": []}, + {"title": "Do B", "assignee": "worker", "parents": [0]}, + {"title": "Do C", "assignee": "worker", "parents": [1]}, + ], + author="auto-decomposer", + ) + assert child_ids is not None and len(child_ids) == 3 + + root = kb.get_task(conn, tid) + assert root.status == "todo" + assert "Automated retriage after timeout" in root.body + for cid in child_ids: + assert kb.get_task(conn, cid) is not None + finally: + conn.close() + + +# --------------------------------------------------------------------------- +# dispatch_once plumbing +# --------------------------------------------------------------------------- + +def test_dispatch_once_surfaces_retriaged(kanban_home, monkeypatch): + """The full chain dispatch_once → _dispatch_once_locked → + enforce_max_runtime → _record_task_failure works, and the retriaged + task id lands in DispatchResult.retriaged.""" + # Keep our own process alive in the eyes of detect_crashed_workers, + # neutralize the real signals + the grace poll sleep. + monkeypatch.setattr(kb, "_pid_alive", lambda pid: True) + monkeypatch.setattr(kb.os, "kill", lambda pid, sig: None) + monkeypatch.setattr(kb.time, "sleep", lambda s: None) + + conn = kb.connect() + try: + tid = kb.create_task( + conn, title="long job", assignee="worker", + max_runtime_seconds=1, + ) + kb.claim_task(conn, tid) + kb._set_worker_pid(conn, tid, os.getpid()) + _age_active_run(conn, tid) + # One failure already on the counter → this tick's timeout trips. + with kb.write_txn(conn): + conn.execute( + "UPDATE tasks SET consecutive_failures = 1 WHERE id = ?", + (tid,), + ) + + res = kb.dispatch_once( + conn, + spawn_fn=lambda task, ws, board=None: None, + retriage_on_timeout=True, + ) + assert tid in res.timed_out + assert tid in res.retriaged + assert kb.get_task(conn, tid).status == "triage" + + # Default-off dispatch keeps the field empty (and, with the task + # now in triage, nothing new times out). + res2 = kb.dispatch_once( + conn, spawn_fn=lambda task, ws, board=None: None, + ) + assert res2.retriaged == [] + finally: + conn.close() diff --git a/tests/tools/test_kanban_tools.py b/tests/tools/test_kanban_tools.py index 7964d2fe5bff7..c4770c084afab 100644 --- a/tests/tools/test_kanban_tools.py +++ b/tests/tools/test_kanban_tools.py @@ -1006,6 +1006,56 @@ def test_create_happy_path(worker_env): conn.close() +def test_create_transport_inherits_parent_priority(worker_env): + """Priority inheritance must survive the transport layer: a tool + call that omits ``priority`` must NOT be coerced to 0 before reaching + ``create_task`` (regression guard for the old ``else 0`` in + _handle_create). The child inherits the parent's priority; the + parent's own priority is untouched.""" + from tools import kanban_tools as kt + from hermes_cli import kanban_db as kb + + conn = kb.connect() + try: + # Give the spawning worker task a concrete priority (raw UPDATE + # — no helper mutates priority directly). + conn.execute("UPDATE tasks SET priority = 6 WHERE id = ?", (worker_env,)) + assert kb.get_task(conn, worker_env).priority == 6 + finally: + conn.close() + + # Omitted priority -> child inherits the parent's 6. + d = json.loads(kt._handle_create({ + "title": "inherit via transport", + "assignee": "peer", + "parents": [worker_env], + })) + assert d["ok"] is True + conn = kb.connect() + try: + child = kb.get_task(conn, d["task_id"]) + assert child.priority == 6 + # Parent never modified by the child's creation. + assert kb.get_task(conn, worker_env).priority == 6 + finally: + conn.close() + + # Explicit priority still wins over inheritance. + d2 = json.loads(kt._handle_create({ + "title": "explicit wins via transport", + "assignee": "peer", + "parents": [worker_env], + "priority": 1, + })) + assert d2["ok"] is True + conn = kb.connect() + try: + assert kb.get_task(conn, d2["task_id"]).priority == 1 + assert kb.get_task(conn, worker_env).priority == 6 + finally: + conn.close() + + def test_create_inherits_worker_dir_workspace(monkeypatch, worker_env): """A worker scoped to a dir: task that spawns a child without a workspace arg inherits the dir, not scratch (so follow-up code-gen diff --git a/tools/kanban_tools.py b/tools/kanban_tools.py index 8983a8b1f5094..37373172e613d 100644 --- a/tools/kanban_tools.py +++ b/tools/kanban_tools.py @@ -623,6 +623,24 @@ def _handle_complete(args: dict, **kw) -> str: f"and keep this task alive." ) + # Auto-commit any uncommitted worktree changes before marking done. + # This is a safety net for worktree kanban tasks whose worker + # forgot to commit (or timed out before committing) — the changes + # are already in the worktree on disk but would be invisible to + # git history without this guard. Best-effort: if the workspace + # is not a git worktree or commit fails, we complete anyway. + if task and task.workspace_kind == "worktree" and task.workspace_path: + try: + kb._git_auto_commit_worktree( + task.workspace_path, tid, + reason="auto-commit on kanban_complete", + ) + except Exception: + logger.warning( + "auto-commit failed for worktree task %s", tid, + exc_info=True, + ) + try: ok = kb.complete_task( conn, tid, @@ -1141,7 +1159,7 @@ def _handle_create(args: dict, **kw) -> str: assignee=str(assignee), parents=tuple(parents), tenant=tenant, - priority=int(priority) if priority is not None else 0, + priority=priority, workspace_kind=str(workspace_kind), workspace_path=workspace_path, project_id=project_id, @@ -1767,7 +1785,9 @@ def _board_schema_prop() -> dict[str, str]: "type": "integer", "description": ( "Dispatcher tiebreaker. Higher = picked sooner " - "when multiple ready tasks share an assignee." + "when multiple ready tasks share an assignee. " + "Omitted = inherit the highest parent priority " + "when parents are given; otherwise 0." ), }, "workspace_kind": { diff --git a/website/docs/user-guide/features/kanban.md b/website/docs/user-guide/features/kanban.md index c8841e9172ddd..5509b661fd60b 100644 --- a/website/docs/user-guide/features/kanban.md +++ b/website/docs/user-guide/features/kanban.md @@ -522,6 +522,7 @@ Config knobs (all under `kanban:` in `~/.hermes/config.yaml`): |---|---|---| | `auto_decompose` | `true` | Dispatcher auto-runs the decomposer every tick. | | `auto_decompose_per_tick` | `3` | Cap on decompositions per dispatcher tick. Excess defers to the next tick. | +| `retriage_on_timeout` | `false` | When the circuit breaker trips on consecutive **timeouts**, send the task back to Triage for decomposition (failure context appended to its body, counter reset) instead of auto-blocking it. A timeout is deterministic — a task that needs more than `max_runtime_seconds` fails identically on every blind retry — so subdivision is the productive recovery. At most one retriage per task; the second trip blocks normally. Crashes and spawn failures never retriage. Works best with `auto_decompose: true`. | | `orchestrator_profile` | `""` | Profile assigned to the root/orchestration task after decomposition. Empty = fall back to active default profile. | | `default_assignee` | `""` | Where a child task lands when the LLM picks an unknown profile. Empty = fall back to active default. | | `auto_subscribe_on_create` | `true` | When a worker calls `kanban_create` from inside a session with a persistent delivery channel (messaging gateway or TUI), the originating session is auto-subscribed to the new task's completion/block events. The dispatcher still drives the delivery — this only changes whether the caller's chat/key shows up in the notify-sub table. Set to `false` to require explicit `kanban_notify-subscribe` calls per task. | @@ -944,6 +945,7 @@ Every transition appends a row to `task_events`. Each row carries an optional `r | `spawn_failed` | `{error, failures}` | One spawn attempt failed (missing PATH, workspace unmountable, …). Counter increments; task returns to `ready` for retry. | | `protocol_violation` | `{pid, claimer, exit_code, protocol_violation}` | Worker exited successfully while the task was still `running`, usually because it answered without calling `kanban_complete` or `kanban_block`. Emitted on every violation (the payload's `protocol_violation: true` marker is copied into the run metadata and feeds the violation-only retry budget). Below the budget — up to `_PROTOCOL_VIOLATION_FAILURE_LIMIT` (default 3) *consecutive* violations, per-task `max_retries` overriding — the task simply returns to `ready` for another attempt; when the streak reaches the bound the dispatcher also emits `gave_up` and auto-blocks. | | `gave_up` | `{failures, effective_limit, limit_source, error}` | Circuit breaker fired after N consecutive non-successful attempts. Task auto-blocks with the last error. The effective limit resolves as task `max_retries`, then dispatcher `failure_limit` / `kanban.failure_limit`, then the built-in default. | +| `retriaged` | `{failures, effective_limit, limit_source, error, trigger_outcome, pid, sigkill}` | `kanban.retriage_on_timeout` fired instead of `gave_up`: the breaker tripped on consecutive timeouts, so the task went back to `triage` for decomposition with a failure-context block appended to its body and a fresh failure counter. Emitted at most once per task — a second trip emits `gave_up` and blocks normally. | `hermes kanban tail ` shows these for a single task. `hermes kanban watch` streams them board-wide.