From b90da8243b07ae1c1ebafeb205a8ed54c51df32d Mon Sep 17 00:00:00 2001 From: Jakub Wolniewicz <4850809+frizikk@users.noreply.github.com> Date: Fri, 31 Jul 2026 14:14:11 +0200 Subject: [PATCH] fix(kanban): preserve review phase across retries --- hermes_cli/kanban_db.py | 277 +++++++++++++----- .../test_kanban_review_lifecycle.py | 67 +++++ .../test_kanban_review_lifecycle_complete.py | 204 +++++++++++++ 3 files changed, 481 insertions(+), 67 deletions(-) diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index e47e9972bc545..4438847b404a9 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -4196,6 +4196,33 @@ def _has_sticky_block(conn: sqlite3.Connection, task_id: str) -> bool: return bool(row) and row["kind"] == "blocked" +def _resume_status_from_events(conn: sqlite3.Connection, task_id: str) -> str: + """Return the durable phase a blocked/dependency-wait task should resume. + + Events written by review workers carry ``source_status``/``retry_status``; + an explicit unblock that must wait for parents carries ``resume_status``. + Legacy events omit these fields and therefore retain the historical + ``ready`` behavior. + """ + row = conn.execute( + "SELECT payload FROM task_events " + "WHERE task_id = ? AND kind IN (" + "'blocked', 'block_loop_detected', 'dependency_wait', 'gave_up', " + "'unblocked', 'changes_requested', 'review_reopened', 'reclaimed', " + "'stale', 'timed_out', 'crashed', 'spawn_failed', 'rate_limited'" + ") ORDER BY id DESC LIMIT 1", + (task_id,), + ).fetchone() + try: + payload = json.loads(row["payload"]) if row and row["payload"] else {} + except (json.JSONDecodeError, TypeError): + payload = {} + for key in ("resume_status", "retry_status", "source_status"): + if payload.get(key) == "review": + return "review" + return "ready" + + def recompute_ready( conn: sqlite3.Connection, failure_limit: int = None, ) -> int: @@ -4251,6 +4278,7 @@ def recompute_ready( (task_id,), ).fetchall() if all(p["status"] in ("done", "archived") for p in parents): + resume_status = _resume_status_from_events(conn, task_id) if cur_status == "blocked": # Don't auto-recover tasks that have hit the # circuit-breaker failure limit. Without this @@ -4269,16 +4297,19 @@ def recompute_ready( if failures >= effective_limit: continue conn.execute( - "UPDATE tasks SET status = 'ready' " + "UPDATE tasks SET status = ? " "WHERE id = ? AND status = 'blocked'", - (task_id,), + (resume_status, task_id), ) else: conn.execute( - "UPDATE tasks SET status = 'ready' WHERE id = ? AND status = 'todo'", - (task_id,), + "UPDATE tasks SET status = ? WHERE id = ? AND status = 'todo'", + (resume_status, task_id), ) - _append_event(conn, task_id, "promoted", None) + _append_event( + conn, task_id, "promoted", + {"status": resume_status} if resume_status != "ready" else None, + ) promoted += 1 return promoted @@ -4484,6 +4515,39 @@ def claim_review_task( return get_task(conn, task_id) +def _retry_status_for_run( + conn: sqlite3.Connection, + task_id: str, + run_id: Optional[int] = None, +) -> str: + """Return the non-running phase an interrupted run must resume from. + + Review claims record ``source_status=review`` on their claimed event. All + other and legacy runs retry from ``ready``. Keeping this decision in one + place prevents crash/timeout/reclaim paths from silently converting a + reviewer run into an implementation run. + """ + if run_id is None: + row = conn.execute( + "SELECT current_run_id FROM tasks WHERE id = ?", + (task_id,), + ).fetchone() + run_id = row["current_run_id"] if row else None + if run_id is None: + return "ready" + event = conn.execute( + "SELECT payload FROM task_events " + "WHERE task_id = ? AND run_id = ? AND kind = 'claimed' " + "ORDER BY id DESC LIMIT 1", + (task_id, int(run_id)), + ).fetchone() + try: + payload = json.loads(event["payload"]) if event and event["payload"] else {} + except (json.JSONDecodeError, TypeError): + payload = {} + return "review" if payload.get("source_status") == "review" else "ready" + + def heartbeat_claim( conn: sqlite3.Connection, task_id: str, @@ -4621,12 +4685,13 @@ def release_stale_claims( ) continue with write_txn(conn): + retry_status = _retry_status_for_run(conn, row["id"]) cur = conn.execute( - "UPDATE tasks SET status = 'ready', claim_lock = NULL, " + "UPDATE tasks SET status = ?, claim_lock = NULL, " "claim_expires = NULL, worker_pid = NULL " "WHERE id = ? AND status = 'running' AND claim_lock IS ? " "AND claim_expires IS NOT NULL AND claim_expires < ?", - (row["id"], row["claim_lock"], now), + (retry_status, row["id"], row["claim_lock"], now), ) if cur.rowcount != 1: continue @@ -4650,6 +4715,7 @@ def release_stale_claims( "now": now, "host_local": host_local, "heartbeat_stale": bool(heartbeat_stale), + "retry_status": retry_status, } payload.update(termination) _append_event( @@ -4668,7 +4734,7 @@ def reclaim_task( reason: Optional[str] = None, signal_fn=None, ) -> bool: - """Operator-driven reclaim: release the claim and reset to ``ready``. + """Operator-driven reclaim: release the claim and restore its source phase. Unlike :func:`release_stale_claims` which only acts on tasks whose ``claim_expires`` has passed, this function reclaims immediately @@ -4693,12 +4759,13 @@ def reclaim_task( row["worker_pid"], prev_lock, signal_fn=signal_fn, ) with write_txn(conn): + retry_status = _retry_status_for_run(conn, task_id) cur = conn.execute( - "UPDATE tasks SET status = 'ready', claim_lock = NULL, " + "UPDATE tasks SET status = ?, claim_lock = NULL, " "claim_expires = NULL, worker_pid = NULL " "WHERE id = ? AND status IN ('running', 'ready', 'blocked') " "AND claim_lock IS ?", - (task_id, prev_lock), + (retry_status, task_id, prev_lock), ) if cur.rowcount != 1: return False @@ -4715,6 +4782,7 @@ def reclaim_task( "manual": True, "reason": reason, "prev_lock": prev_lock, + "retry_status": retry_status, } payload.update(termination) _append_event( @@ -5728,6 +5796,11 @@ def block_task( ).fetchone() if cur_row is None: return False + source_status = ( + _retry_status_for_run(conn, task_id) + if cur_row["status"] == "running" + else "ready" + ) prev_kind = cur_row["block_kind"] if "block_kind" in cur_row.keys() else None prev_recurrences = ( int(cur_row["block_recurrences"]) @@ -5768,7 +5841,12 @@ def block_task( ) _append_event( conn, task_id, "dependency_wait", - {"reason": reason, "kind": kind}, run_id=run_id, + { + "reason": reason, + "kind": kind, + "source_status": source_status, + }, + run_id=run_id, ) _blocked_task = get_task(conn, task_id) _fire_kanban_lifecycle_hook( @@ -5826,6 +5904,7 @@ def block_task( "kind": kind, "recurrences": recurrences, "limit": BLOCK_RECURRENCE_LIMIT, + "source_status": source_status, }, run_id=run_id, ) @@ -5878,7 +5957,12 @@ def block_task( ) _append_event( conn, task_id, "blocked", - {"reason": reason, "kind": kind, "recurrences": recurrences}, + { + "reason": reason, + "kind": kind, + "recurrences": recurrences, + "source_status": source_status, + }, run_id=run_id, ) _blocked_task = get_task(conn, task_id) @@ -6205,7 +6289,7 @@ def _landing_status_after_parents(conn: sqlite3.Connection, task_id: str) -> str def unblock_task(conn: sqlite3.Connection, task_id: str) -> bool: - """Transition ``blocked``/``scheduled`` -> ready or todo. + """Transition ``blocked``/``scheduled`` to its safe resumable phase. Defensively closes any stale ``current_run_id`` pointer before flipping status. In the common path (``block_task`` closed the run already) this @@ -6216,13 +6300,26 @@ def unblock_task(conn: sqlite3.Connection, task_id: str) -> bool: """ now = int(time.time()) with write_txn(conn): + current = conn.execute( + "SELECT status FROM tasks WHERE id = ?", + (task_id,), + ).fetchone() + resume_status = ( + _resume_status_from_events(conn, task_id) + if current and current["status"] == "blocked" + else "ready" + ) _reclaim_dangling_run( conn, task_id, statuses=("blocked", "scheduled"), now=now, note="invariant recovery on unblock", ) - # Re-gate on parent completion before flipping 'blocked' back to - # 'ready' (see :func:`_landing_status_after_parents`). - new_status = _landing_status_after_parents(conn, task_id) + # Re-gate on parent completion before restoring the source phase. + landing_status = _landing_status_after_parents(conn, task_id) + new_status = ( + "review" + if landing_status == "ready" and resume_status == "review" + else landing_status + ) # NOTE: deliberately does NOT touch ``block_recurrences`` or # ``block_kind``. Resetting the recurrence counter on unblock is exactly # the amnesia that let a cron unblock → worker re-block loop run @@ -6243,7 +6340,11 @@ def unblock_task(conn: sqlite3.Connection, task_id: str) -> bool: return False _append_event( conn, task_id, "unblocked", - {"status": new_status} if new_status != "ready" else None, + ( + {"status": new_status, "resume_status": resume_status} + if new_status != "ready" or resume_status != "ready" + else None + ), ) return True @@ -7546,8 +7647,8 @@ def enforce_max_runtime( """Terminate workers whose per-task ``max_runtime_seconds`` has elapsed. Sends SIGTERM, waits a short grace window, then SIGKILL. Emits a - ``timed_out`` event and drops the task back to ``ready`` so the next - dispatcher tick re-spawns it — unless the spawn-failure circuit + ``timed_out`` event and restores the task's source phase so the next + dispatcher tick re-spawns the same kind of worker — unless the circuit breaker has already given up, in which case the task stays blocked where ``_record_spawn_failure`` parked it. @@ -7610,13 +7711,14 @@ def enforce_max_runtime( pass with write_txn(conn): + retry_status = _retry_status_for_run(conn, tid) cur = conn.execute( - "UPDATE tasks SET status = 'ready', claim_lock = NULL, " + "UPDATE tasks SET status = ?, claim_lock = NULL, " "claim_expires = NULL, worker_pid = NULL, " "last_heartbeat_at = NULL " "WHERE id = ? AND status = 'running' " " AND worker_pid = ? AND claim_lock IS ?", - (tid, pid, row["claim_lock"]), + (retry_status, tid, pid, row["claim_lock"]), ) if cur.rowcount == 1: payload = { @@ -7624,6 +7726,7 @@ def enforce_max_runtime( "elapsed_seconds": int(elapsed), "limit_seconds": int(row["max_runtime_seconds"]), "sigkill": killed, + "retry_status": retry_status, } run_id = _end_run( conn, tid, @@ -7637,7 +7740,7 @@ def enforce_max_runtime( timed_out.append(tid) # Increment the unified failure counter. Outside the write_txn # above because ``_record_task_failure`` opens its own. If the - # breaker trips, this flips the task ``ready → blocked`` and + # breaker trips, this flips the retried task to ``blocked`` and # emits a ``gave_up`` event on top of the ``timed_out`` we # already emitted. if cur.rowcount == 1: @@ -7647,7 +7750,11 @@ def enforce_max_runtime( outcome="timed_out", release_claim=False, end_run=False, - event_payload_extra={"pid": pid, "sigkill": killed}, + event_payload_extra={ + "pid": pid, + "sigkill": killed, + "retry_status": retry_status, + }, ) return timed_out @@ -7676,7 +7783,7 @@ def detect_stale_running( 2. Its ``last_heartbeat_at`` is older than ``_STALE_HEARTBEAT_GAP_SECONDS`` (or NULL — never sent a heartbeat). - On reclaim the task is reset to ``ready``, the run is closed with + On reclaim the task is restored to its source phase, the run is closed with ``outcome='stale'``, and the host-local worker (if still running) is terminated. @@ -7735,13 +7842,14 @@ def detect_stale_running( continue with write_txn(conn): + retry_status = _retry_status_for_run(conn, tid) cur = conn.execute( - "UPDATE tasks SET status = 'ready', claim_lock = NULL, " + "UPDATE tasks SET status = ?, claim_lock = NULL, " "claim_expires = NULL, worker_pid = NULL, " "last_heartbeat_at = NULL " "WHERE id = ? AND status = 'running' " " AND claim_lock IS ?", - (tid, row["claim_lock"]), + (retry_status, tid, row["claim_lock"]), ) if cur.rowcount != 1: continue @@ -7756,6 +7864,7 @@ def detect_stale_running( ), "timeout_seconds": stale_timeout_seconds, "pid": int(pid) if pid else None, + "retry_status": retry_status, } payload.update(termination) @@ -7776,7 +7885,7 @@ def detect_stale_running( # Intentionally NOT calling _record_task_failure here. Stale reclaim # is dispatcher-side detection of an absent heartbeat; the task is - # going straight back to ``ready`` for re-dispatch. Counting it as + # going straight back to its source phase for re-dispatch. Counting it as # a worker failure would let two legitimately-long-running tasks # (>4h without explicit heartbeat) trip the circuit breaker and # auto-block, even though no worker actually failed. The 'stale' @@ -7962,7 +8071,7 @@ def _protocol_violation_streak(conn: sqlite3.Connection, task_id: str) -> int: def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]: """Reclaim ``running`` tasks whose worker PID is no longer alive. - Appends a ``crashed`` event and drops the task back to ``ready``. + Appends a ``crashed`` event and restores the task's source phase. Different from ``release_stale_claims``: this checks liveness immediately rather than waiting for the claim TTL. @@ -7981,7 +8090,7 @@ def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]: When the reap registry shows the worker exited with the rate-limit sentinel (``KANBAN_RATE_LIMIT_EXIT_CODE``), the worker bailed on a provider quota wall, NOT a task failure. Such tasks are released back - to ``ready`` WITHOUT counting a failure (so a long quota window can't + to its source phase WITHOUT counting a failure (so a long quota window can't trip the breaker) and stamped with a quota-blocker error so ``check_respawn_guard`` defers their respawn until the window clears. The ids are returned via the ``_last_rate_limited`` function attribute @@ -8053,7 +8162,7 @@ def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]: # Worker bailed because the provider rate-limited / exhausted # quota (EX_TEMPFAIL sentinel). This is NOT a task failure — # the task is fine, the account just hit a wall. Release it - # back to ``ready`` so the respawn guard defers it until the + # back to its source phase so the respawn guard defers it until the # quota window clears, and crucially do NOT count a failure # (skip ``_record_task_failure``) so a long quota window can't # trip the circuit breaker and permanently block the card. @@ -8083,12 +8192,14 @@ def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]: event_payload["exit_kind"] = kind event_payload["exit_code"] = code + retry_status = _retry_status_for_run(conn, row["id"]) + event_payload["retry_status"] = retry_status cur = conn.execute( - "UPDATE tasks SET status = 'ready', claim_lock = NULL, " + "UPDATE tasks SET status = ?, claim_lock = NULL, " "claim_expires = NULL, worker_pid = NULL " "WHERE id = ? AND status = 'running' " " AND worker_pid = ? AND claim_lock IS ?", - (row["id"], pid, row["claim_lock"]), + (retry_status, row["id"], pid, row["claim_lock"]), ) if cur.rowcount == 1: # Rate-limited requeues are a clean release, not a crash — @@ -8136,7 +8247,7 @@ def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]: protocol_violation, error_text) ) # Outside the main txn: account each crashed task and maybe trip the - # breaker (the task transitions ready → blocked with a ``gave_up`` event + # breaker (the retried task transitions to blocked with a ``gave_up`` event # on top of the event we already emitted). # # Protocol-violation crashes (clean exit, no terminal tool call) get a @@ -8256,14 +8367,14 @@ def _record_task_failure( * ``release_claim=True, end_run=True`` — spawn-failure path. Caller has a running task with an open run; this transitions - it back to ``ready`` (or ``blocked`` when the breaker trips), + it back to its source phase (or ``blocked`` when the breaker trips), releases the claim, and closes the run with ``outcome=``. * ``release_claim=False, end_run=False`` — timeout/crash path. - Caller has ALREADY flipped the task to ``ready`` and closed the + Caller has ALREADY restored the task's source phase and closed the run with the appropriate outcome. This just increments the counter; if the breaker trips, the task is re-transitioned - ``ready → blocked`` and a ``gave_up`` event is emitted. + into ``blocked`` and a ``gave_up`` event is emitted. ``event_payload_extra`` merges into the ``gave_up`` event payload when the breaker trips, so callers can include outcome-specific @@ -8289,11 +8400,16 @@ def _record_task_failure( blocked = False with write_txn(conn): row = conn.execute( - "SELECT consecutive_failures, status, max_retries " + "SELECT consecutive_failures, status, max_retries, current_run_id " "FROM tasks WHERE id = ?", (task_id,), ).fetchone() if row is None: return False + retry_status = ( + _retry_status_for_run(conn, task_id, row["current_run_id"]) + if release_claim + else ("review" if row["status"] == "review" else "ready") + ) failures = int(row["consecutive_failures"]) + 1 # Per-task override wins over both caller-supplied and default @@ -8316,17 +8432,17 @@ def _record_task_failure( "UPDATE tasks SET status = 'blocked', claim_lock = NULL, " "claim_expires = NULL, worker_pid = NULL, " "consecutive_failures = ?, last_failure_error = ? " - "WHERE id = ? AND status IN ('running', 'ready')", + "WHERE id = ? AND status IN ('running', 'ready', 'review')", (failures, error[:500], task_id), ) else: - # Timeout/crash path: task is already at ``ready`` - # with claim cleared; just flip to blocked + update + # Timeout/crash path: source phase already restored with claim + # cleared; just flip to blocked + update # counter fields. conn.execute( "UPDATE tasks SET status = 'blocked', " "consecutive_failures = ?, last_failure_error = ? " - "WHERE id = ? AND status IN ('ready', 'running')", + "WHERE id = ? AND status IN ('ready', 'review', 'running')", (failures, error[:500], task_id), ) run_id = None @@ -8341,6 +8457,7 @@ def _record_task_failure( "trigger_outcome": outcome, "effective_limit": effective_limit, "limit_source": limit_source, + "retry_status": retry_status, }, ) payload = { @@ -8349,6 +8466,7 @@ def _record_task_failure( "limit_source": limit_source, "error": error[:500], "trigger_outcome": outcome, + "retry_status": retry_status, } if event_payload_extra: payload.update(event_payload_extra) @@ -8359,17 +8477,16 @@ def _record_task_failure( else: # Below threshold. if release_claim: - # Spawn path: transition running → ready + clear claim. + # Spawn path: restore the claimed source phase + clear claim. conn.execute( - "UPDATE tasks SET status = 'ready', claim_lock = NULL, " + "UPDATE tasks SET status = ?, claim_lock = NULL, " "claim_expires = NULL, worker_pid = NULL, " "consecutive_failures = ?, last_failure_error = ? " "WHERE id = ? AND status = 'running'", - (failures, error[:500], task_id), + (retry_status, failures, error[:500], task_id), ) else: - # Timeout/crash path: task is already at ``ready`` via - # its own UPDATE. Just bookkeep the counter + last error. + # Timeout/crash path: caller already restored the source phase. conn.execute( "UPDATE tasks SET consecutive_failures = ?, " "last_failure_error = ? WHERE id = ?", @@ -8381,11 +8498,18 @@ def _record_task_failure( conn, task_id, outcome=outcome, status=outcome, error=error[:500], - metadata={"failures": failures}, + metadata={ + "failures": failures, + "retry_status": retry_status, + }, ) _append_event( conn, task_id, outcome, - {"error": error[:500], "failures": failures}, + { + "error": error[:500], + "failures": failures, + "retry_status": retry_status, + }, run_id=run_id, ) # Timeout/crash path's caller already emitted its own event. @@ -8456,9 +8580,9 @@ _clear_spawn_failures = _clear_failure_counter def check_respawn_guard(conn: sqlite3.Connection, task_id: str) -> Optional[str]: """Return a guard reason if ``task_id`` should NOT be re-spawned, else None. - Called per ready task in ``dispatch_once`` before any claim attempt. - Returning a reason defers the spawn this tick; the task stays in - ``ready`` and gets another chance on the next dispatcher tick. + Called per ready/review task in ``dispatch_once`` before any claim attempt. + Returning a reason defers the spawn this tick; the task stays in its + source phase and gets another chance on the next dispatcher tick. Checks in priority order: @@ -8811,6 +8935,14 @@ def _dispatch_once_locked( result.timed_out = enforce_max_runtime(conn) result.promoted = recompute_ready(conn, failure_limit=failure_limit) + # Both knobs are total in-flight caps. Collapse them before either lane + # dispatches so ready and review workers consume the same budget without + # subtracting the already-running count twice. + if max_in_progress is not None and ( + max_spawn is None or max_in_progress < max_spawn + ): + max_spawn = max_in_progress + # Count tasks already running so max_spawn enforces concurrency rather # than a per-tick spawn budget. See the docstring above for the full # rationale; the short version is that a 60-second tick interval with a @@ -8831,20 +8963,6 @@ def _dispatch_once_locked( "WHERE status = 'ready' AND claim_lock IS NULL " "ORDER BY priority DESC, created_at ASC" ).fetchall() - # Honour kanban.max_in_progress: if the board already has enough running - # tasks, skip spawning this tick so slow workers (local LLMs, - # resource-constrained hosts) can finish what they have before more tasks - # pile up and time out. - if max_in_progress is not None and ready_rows: - in_progress = conn.execute( - "SELECT COUNT(*) FROM tasks WHERE status = 'running'" - ).fetchone()[0] - if in_progress >= max_in_progress: - return result - # Only spawn enough to reach the cap, respecting max_spawn too. - remaining = max_in_progress - in_progress - if max_spawn is None or max_spawn > remaining: - max_spawn = remaining spawned = 0 # Per-profile concurrency cap (#21582): when set, track how many # workers each assignee already has in flight, and refuse to spawn @@ -9062,8 +9180,8 @@ def _dispatch_once_locked( # ---- review column dispatch ---- # Review tasks are tasks that a worker moved to 'review' after # creating a PR. The dispatcher spawns a review agent (loading - # sdlc-review skill) that verifies the PR and either merges (→ done) - # or rejects (→ back to running for the worker to fix). + # sdlc-review skill) that verifies the candidate and either approves + # (→ done) or requests changes (→ ready/todo for the implementer). # # Same concurrency model as ready dispatch: review spawns count # against max_spawn alongside ready tasks, so the total number of @@ -9092,8 +9210,29 @@ def _dispatch_once_locked( if profile_exists is not None and not profile_exists(row["assignee"]): result.skipped_nonspawnable.append(row["id"]) continue + if _per_profile_cap is not None: + current = _per_profile_running.get(row["assignee"], 0) + if current >= _per_profile_cap: + result.skipped_per_profile_capped.append( + (row["id"], row["assignee"], current) + ) + continue + guard_reason = check_respawn_guard(conn, row["id"]) + if guard_reason is not None: + result.respawn_guarded.append((row["id"], guard_reason)) + if not dry_run: + with write_txn(conn): + _append_event( + conn, row["id"], "respawn_guarded", + {"reason": guard_reason}, + ) + continue if dry_run: result.spawned.append((row["id"], row["assignee"], "")) + if _per_profile_cap is not None: + _per_profile_running[row["assignee"]] = ( + _per_profile_running.get(row["assignee"], 0) + 1 + ) continue claimed = claim_review_task(conn, row["id"], ttl_seconds=ttl_seconds) if claimed is None: @@ -9140,6 +9279,10 @@ def _dispatch_once_locked( _set_worker_pid(conn, claimed.id, int(pid)) result.spawned.append((claimed.id, claimed.assignee or "", str(workspace))) spawned += 1 + if _per_profile_cap is not None and claimed.assignee: + _per_profile_running[claimed.assignee] = ( + _per_profile_running.get(claimed.assignee, 0) + 1 + ) except Exception as exc: auto = _record_spawn_failure( conn, claimed.id, str(exc), diff --git a/tests/hermes_cli/test_kanban_review_lifecycle.py b/tests/hermes_cli/test_kanban_review_lifecycle.py index c3337a980da4c..9fb8d54100a2a 100644 --- a/tests/hermes_cli/test_kanban_review_lifecycle.py +++ b/tests/hermes_cli/test_kanban_review_lifecycle.py @@ -365,12 +365,79 @@ def test_review_dispatch_preserves_task_skills_and_adds_reviewer_skill( summary="ready", expected_run_id=implementation.current_run_id, ) + monkeypatch.setattr( + kb, + "check_respawn_guard", + lambda _conn, _task_id: "rate_limit_cooldown", + ) + guarded = kb.dispatch_once(conn, spawn_fn=spawn) + assert guarded.respawn_guarded == [(task_id, "rate_limit_cooldown")] + assert not guarded.spawned + guarded_task = kb.get_task(conn, task_id) + assert guarded_task is not None + assert guarded_task.status == "review" + + monkeypatch.setattr(kb, "check_respawn_guard", lambda _conn, _task_id: None) result = kb.dispatch_once(conn, spawn_fn=spawn) assert task_id in [task[0] for task in result.spawned] assert captured == [["domain-specific-review", "sdlc-review"]] +def test_review_dispatch_honors_global_and_per_profile_caps( + kanban_home: Path, + monkeypatch: pytest.MonkeyPatch, +) -> None: + import hermes_cli.config as cfgmod + import hermes_cli.profiles as profmod + + monkeypatch.setattr(profmod, "profile_exists", lambda _name: True) + monkeypatch.setattr( + cfgmod, + "load_config", + lambda *args, **kwargs: {"kanban": {"review_dispatch": True}}, + ) + + with kb.connect() as conn: + running_id = kb.create_task(conn, title="already running", assignee="builder") + assert kb.claim_task(conn, running_id) is not None + + review_ids: list[str] = [] + for title in ("review one", "review two"): + task_id = kb.create_task(conn, title=title, assignee="reviewer") + implementation = kb.claim_task(conn, task_id) + assert implementation is not None + assert kb.request_review( + conn, + task_id, + summary="ready", + expected_run_id=implementation.current_run_id, + ) + review_ids.append(task_id) + + globally_capped = kb.dispatch_once( + conn, + dry_run=True, + max_in_progress=1, + ) + assert not [ + task for task in globally_capped.spawned if task[0] in review_ids + ] + + per_profile_capped = kb.dispatch_once( + conn, + dry_run=True, + max_in_progress=10, + max_in_progress_per_profile=1, + ) + spawned_reviews = [ + task for task in per_profile_capped.spawned if task[0] in review_ids + ] + assert len(spawned_reviews) == 1 + assert len(per_profile_capped.skipped_per_profile_capped) == 1 + assert per_profile_capped.skipped_per_profile_capped[0][0] in review_ids + + # --------------------------------------------------------------------------- # reopen: a follow-up sends a review task back out for another pass # --------------------------------------------------------------------------- diff --git a/tests/hermes_cli/test_kanban_review_lifecycle_complete.py b/tests/hermes_cli/test_kanban_review_lifecycle_complete.py index 97cf76b7c966f..165bbab2797f5 100644 --- a/tests/hermes_cli/test_kanban_review_lifecycle_complete.py +++ b/tests/hermes_cli/test_kanban_review_lifecycle_complete.py @@ -11,6 +11,7 @@ These tests cover the two review models that must coexist: from __future__ import annotations +import time from pathlib import Path import pytest @@ -36,6 +37,37 @@ def _run(runs, outcome: str): return [run for run in runs if run.outcome == outcome][-1] +def _claimed_review( + conn, + title: str, + *, + ttl_seconds: int | None = None, + max_runtime_seconds: int | None = None, +): + task_id = kb.create_task( + conn, + title=title, + assignee="builder", + max_runtime_seconds=max_runtime_seconds, + ) + implementation = kb.claim_task(conn, task_id, claimer="builder:test") + assert implementation is not None + assert kb.request_review( + conn, + task_id, + summary="ready for independent review", + reviewer="reviewer", + expected_run_id=implementation.current_run_id, + ) + review = kb.claim_review_task( + conn, + task_id, + ttl_seconds=ttl_seconds, + ) + assert review is not None + return task_id, review + + def test_same_card_review_supports_changes_and_approval_without_block_loop(conn): task_id = kb.create_task(conn, title="Implement guarded export", assignee="builder") implementation = kb.claim_task(conn, task_id, claimer="builder:1") @@ -180,6 +212,178 @@ def test_request_changes_fails_closed_on_malformed_review_provenance(conn): assert task.current_run_id == review.current_run_id +@pytest.mark.parametrize( + "reclaim_kind", + ["spawn_failure", "expired_claim", "manual_reclaim", "stale_heartbeat"], +) +def test_interrupted_review_runs_retry_in_review_phase( + conn, + reclaim_kind: str, +) -> None: + task_id, review = _claimed_review( + conn, + f"Retry review after {reclaim_kind}", + ttl_seconds=-1 if reclaim_kind == "expired_claim" else None, + ) + + if reclaim_kind == "spawn_failure": + assert not kb._record_spawn_failure( + conn, + task_id, + "reviewer process failed to spawn", + failure_limit=3, + ) + elif reclaim_kind == "expired_claim": + with kb.write_txn(conn): + conn.execute( + "UPDATE tasks SET claim_expires = ? WHERE id = ?", + (int(time.time()) - 1, task_id), + ) + assert kb.release_stale_claims(conn) == 1 + elif reclaim_kind == "manual_reclaim": + assert kb.reclaim_task(conn, task_id, reason="operator retry") + else: + old = int(time.time()) - 1_000 + with kb.write_txn(conn): + conn.execute( + "UPDATE tasks SET started_at = ?, last_heartbeat_at = NULL " + "WHERE id = ?", + (old, task_id), + ) + conn.execute( + "UPDATE task_runs SET started_at = ? WHERE id = ?", + (old, review.current_run_id), + ) + assert kb.detect_stale_running(conn, stale_timeout_seconds=1) == [task_id] + + retried = kb.get_task(conn, task_id) + assert retried is not None + assert retried.status == "review" + assert retried.current_run_id is None + event = kb.list_events(conn, task_id=task_id)[-1] + assert event.payload is not None + assert event.payload.get("retry_status") == "review" + + +def test_review_retry_still_trips_the_failure_breaker(conn) -> None: + task_id, _review = _claimed_review(conn, "Reviewer repeatedly fails") + assert kb._record_spawn_failure( + conn, + task_id, + "reviewer cannot start", + failure_limit=1, + ) + blocked = kb.get_task(conn, task_id) + assert blocked is not None + assert blocked.status == "blocked" + gave_up = _event(kb.list_events(conn, task_id), "gave_up") + assert gave_up.payload is not None + assert gave_up.payload["retry_status"] == "review" + assert kb.unblock_task(conn, task_id) + unblocked = kb.get_task(conn, task_id) + assert unblocked is not None + assert unblocked.status == "review" + + +def test_review_escalation_unblocks_back_to_review(conn) -> None: + task_id, review = _claimed_review(conn, "External review escalation") + assert kb.block_task( + conn, + task_id, + reason="needs_input: maintainer decision required", + kind="needs_input", + expected_run_id=review.current_run_id, + ) + blocked_event = _event(kb.list_events(conn, task_id), "blocked") + assert blocked_event.payload is not None + assert blocked_event.payload["source_status"] == "review" + assert kb.unblock_task(conn, task_id) + resumed = kb.get_task(conn, task_id) + assert resumed is not None + assert resumed.status == "review" + + +def test_review_dependency_wait_reenters_review_after_parent_finishes(conn) -> None: + parent_id = kb.create_task(conn, title="Parent", assignee="planner") + assert kb.complete_task(conn, parent_id) + task_id = kb.create_task( + conn, + title="Review after dependency refresh", + assignee="builder", + parents=[parent_id], + ) + implementation = kb.claim_task(conn, task_id) + assert implementation is not None + assert kb.request_review( + conn, + task_id, + summary="ready", + reviewer="reviewer", + expected_run_id=implementation.current_run_id, + ) + review = kb.claim_review_task(conn, task_id) + assert review is not None + with kb.write_txn(conn): + conn.execute("UPDATE tasks SET status = 'ready' WHERE id = ?", (parent_id,)) + assert kb.block_task( + conn, + task_id, + reason="dependency: parent contract is being refreshed", + kind="dependency", + expected_run_id=review.current_run_id, + ) + waiting = kb.get_task(conn, task_id) + assert waiting is not None + assert waiting.status == "todo" + assert kb.complete_task(conn, parent_id) + resumed = kb.get_task(conn, task_id) + assert resumed is not None + assert resumed.status == "review" + + +def test_crashed_and_timed_out_review_runs_retry_in_review_phase( + conn, + monkeypatch: pytest.MonkeyPatch, +) -> None: + monkeypatch.setattr(kb, "_pid_alive", lambda _pid: False) + monkeypatch.setattr(kb, "_classify_worker_exit", lambda _pid: ("nonzero_exit", 1)) + old = int(time.time()) - 1_000 + + timed_out_id, timed_out_run = _claimed_review( + conn, + "Timeout during review", + max_runtime_seconds=1, + ) + with kb.write_txn(conn): + conn.execute( + "UPDATE tasks SET worker_pid = ?, started_at = ? WHERE id = ?", + (999_998, old, timed_out_id), + ) + conn.execute( + "UPDATE task_runs SET worker_pid = ?, started_at = ? WHERE id = ?", + (999_998, old, timed_out_run.current_run_id), + ) + assert timed_out_id in kb.enforce_max_runtime(conn, signal_fn=lambda *_: None) + timed_out = kb.get_task(conn, timed_out_id) + assert timed_out is not None + assert timed_out.status == "review" + + crashed_id, crashed_run = _claimed_review(conn, "Crash during review") + with kb.write_txn(conn): + conn.execute( + "UPDATE tasks SET worker_pid = ?, started_at = ? WHERE id = ?", + (999_999, old, crashed_id), + ) + conn.execute( + "UPDATE task_runs SET worker_pid = ?, started_at = ? WHERE id = ?", + (999_999, old, crashed_run.current_run_id), + ) + assert crashed_id in kb.detect_crashed_workers(conn) + crashed = kb.get_task(conn, crashed_id) + assert crashed is not None + assert crashed.status == "review" + + def test_legacy_review_child_deadlock_is_reported_immediately(conn): implementation_id = kb.create_task( conn,