From 6d07a9fa040832b65219064d58d49328a8662799 Mon Sep 17 00:00:00 2001 From: SoLo Date: Sat, 1 Aug 2026 06:48:19 -0400 Subject: [PATCH 1/4] feat(kanban): route implementation handoffs through review --- agent/prompt_builder.py | 13 +- hermes_cli/kanban.py | 57 ++++++++ hermes_cli/kanban_db.py | 126 +++++++++++++++++- .../test_kanban_review_lifecycle.py | 101 ++++++++++++++ tools/kanban_tools.py | 102 ++++++++++++++ .../features/kanban-worker-lanes.md | 14 +- 6 files changed, 395 insertions(+), 18 deletions(-) create mode 100644 tests/hermes_cli/test_kanban_review_lifecycle.py diff --git a/agent/prompt_builder.py b/agent/prompt_builder.py index 0f748f8b2bd24..c39eb4970ae17 100644 --- a/agent/prompt_builder.py +++ b/agent/prompt_builder.py @@ -243,14 +243,11 @@ def _strip_yaml_frontmatter(content: str) -> str: "(`{changed_files: [...], tests_run: N, decisions: [...]}`). Downstream " "workers read both via their own `kanban_show`. Never put secrets / " "tokens / raw PII in either field — run rows are durable forever. " - "Exception: if your output is a code change that needs human review " - "before counting as merged/done (most coding tasks), drop the " - "structured metadata (changed_files / tests_run / diff_path) into a " - "`kanban_comment` first, then end with " - "`kanban_block(reason=\"review-required: \")` so a " - "reviewer can approve+unblock or request changes. Reviewing-then-" - "completing is more honest than auto-completing work that still needs " - "eyes on it.\n" + "Exception: if your output is a code change that needs independent review, " + "call `kanban_submit_review(reviewer=..., summary=..., metadata=...)`. " + "It preserves implementation evidence and routes the card to the Review " + "lane; `kanban_block` remains for genuine human input, credentials, " + "capability, dependency, or transient failures.\n" "6. **If follow-up work appears, create it; don't do it.** Use " "`kanban_create(title=..., assignee=, parents=[your-task-id])` " "to spawn a child task for the appropriate specialist profile instead of " diff --git a/hermes_cli/kanban.py b/hermes_cli/kanban.py index a08cb8f9b4076..b04fb9c44881f 100644 --- a/hermes_cli/kanban.py +++ b/hermes_cli/kanban.py @@ -601,6 +601,21 @@ def build_parser(parent_subparsers: argparse._SubParsersAction) -> argparse.Argu help='JSON dict of structured facts (e.g. \'{"changed_files": [...], ' '"tests_run": 12}\'). Stored on the closing run.') + p_submit_review = sub.add_parser( + "submit-review", help="Submit a running implementation to the Review lane" + ) + p_submit_review.add_argument("task_id") + p_submit_review.add_argument("reviewer") + p_submit_review.add_argument("summary", nargs="+", help="Review handoff summary") + p_submit_review.add_argument("--metadata", default=None, help="JSON evidence object") + + p_review_changes = sub.add_parser( + "review-changes", help="Complete a review and create implementer remediation" + ) + p_review_changes.add_argument("task_id") + p_review_changes.add_argument("summary", nargs="+", help="Requested changes") + p_review_changes.add_argument("--metadata", default=None, help="JSON findings object") + p_edit = sub.add_parser( "edit", help="Edit recovery fields on an already-completed task", @@ -1062,6 +1077,8 @@ def kanban_command(args: argparse.Namespace) -> int: "attachments": _cmd_attachments, "attach-rm": _cmd_attach_rm, "complete": _cmd_complete, + "submit-review": _cmd_submit_review, + "review-changes": _cmd_review_changes, "edit": _cmd_edit, "block": _cmd_block, "schedule": _cmd_schedule, @@ -2138,6 +2155,46 @@ def _worker_run_id_for(task_id: str) -> Optional[int]: return None +def _cmd_submit_review(args: argparse.Namespace) -> int: + metadata = None + if args.metadata: + metadata = json.loads(args.metadata) + if not isinstance(metadata, dict): + raise ValueError("--metadata must be a JSON object") + with kb.connect_closing() as conn: + task = kb.get_task(conn, args.task_id) + run_id = task.current_run_id if task else None + if not kb.submit_for_review( + conn, args.task_id, reviewer=args.reviewer, + summary=" ".join(args.summary), metadata=metadata, + expected_run_id=run_id, + ): + print(f"cannot submit {args.task_id} for review", file=sys.stderr) + return 1 + print(f"Submitted {args.task_id} for review") + return 0 + + +def _cmd_review_changes(args: argparse.Namespace) -> int: + metadata = None + if args.metadata: + metadata = json.loads(args.metadata) + if not isinstance(metadata, dict): + raise ValueError("--metadata must be a JSON object") + with kb.connect_closing() as conn: + task = kb.get_task(conn, args.task_id) + run_id = task.current_run_id if task else None + remediation = kb.request_review_changes( + conn, args.task_id, summary=" ".join(args.summary), metadata=metadata, + expected_run_id=run_id, + ) + if not remediation: + print(f"cannot request changes for {args.task_id}", file=sys.stderr) + return 1 + print(f"Review changes recorded; remediation task: {remediation}") + return 0 + + def _cmd_complete(args: argparse.Namespace) -> int: """Mark one or more tasks done. Supports a single id or a list.""" ids = list(args.task_ids or []) diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index b64d54e53ab7c..104334864280b 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -100,7 +100,7 @@ # --------------------------------------------------------------------------- VALID_STATUSES = {"triage", "todo", "scheduled", "ready", "running", "blocked", "review", "done", "archived"} -VALID_INITIAL_STATUSES = {"running", "blocked"} +VALID_INITIAL_STATUSES = {"running", "blocked", "review"} # Typed block reasons. Distinguishes the two fundamentally different things a # worker (or human) means by "blocked", so each can be routed differently @@ -3157,12 +3157,15 @@ def create_task( for attempt in range(2): task_id = _new_task_id() try: - with write_txn(conn): + # A review changes-requested handoff may create a remediation + # while already holding the lifecycle transaction. SQLite has no + # nested BEGIN support, so reuse that transaction when present. + with (contextlib.nullcontext() if conn.in_transaction else write_txn(conn)): # Determine task status from parent status, unless the caller # parks it directly in blocked for human-ops review or in # triage for a specifier. - if initial_status == "blocked": - task_status = "blocked" + if initial_status in {"blocked", "review"}: + task_status = initial_status if parents: missing = _find_missing_parents(conn, parents) if missing: @@ -4420,6 +4423,121 @@ def claim_review_task( return get_task(conn, task_id) +def submit_for_review( + conn: sqlite3.Connection, + task_id: str, + *, + reviewer: str, + summary: str, + metadata: Optional[dict] = None, + expected_run_id: Optional[int] = None, +) -> bool: + """Move a running implementation to ``review`` with audit evidence. + + The implementation run is closed, but the task remains the canonical + review card. A later reviewer claim creates a new run, preserving both + sides of the handoff and preventing the implementation worker from being + respawned. + """ + reviewer = _canonical_assignee(reviewer) + if not reviewer: + raise ValueError("reviewer is required") + if not summary or not summary.strip(): + raise ValueError("review summary is required") + with write_txn(conn): + row = conn.execute( + "SELECT assignee, status FROM tasks WHERE id = ?", (task_id,) + ).fetchone() + if row is None or row["status"] != "running": + return False + original_assignee = str(row["assignee"] or "") + where = "id = ? AND status = 'running'" + params: tuple[Any, ...] = (task_id,) + if expected_run_id is not None: + where += " AND current_run_id = ?" + params += (int(expected_run_id),) + cur = conn.execute( + "UPDATE tasks SET status='review', assignee=?, claim_lock=NULL, " + "claim_expires=NULL, worker_pid=NULL WHERE " + where, + (reviewer, *params), + ) + if cur.rowcount != 1: + return False + handoff = dict(metadata or {}) + handoff.update({"reviewer": reviewer, "original_assignee": original_assignee}) + run_id = _end_run( + conn, task_id, outcome="submitted_for_review", status="review", + summary=summary.strip(), metadata=handoff, + ) + _append_event( + conn, task_id, "review_submitted", + {"reviewer": reviewer, "original_assignee": original_assignee, + "summary": summary.strip().splitlines()[0][:400], "metadata": handoff}, + run_id=run_id, + ) + return True + + +def request_review_changes( + conn: sqlite3.Connection, + task_id: str, + *, + summary: str, + metadata: Optional[dict] = None, + expected_run_id: Optional[int] = None, +) -> Optional[str]: + """Complete a review with findings and create one remediation card.""" + if not summary or not summary.strip(): + raise ValueError("changes-requested summary is required") + with write_txn(conn): + row = conn.execute("SELECT * FROM tasks WHERE id = ?", (task_id,)).fetchone() + if row is None or row["status"] != "running": + return None + if expected_run_id is not None and row["current_run_id"] != int(expected_run_id): + return None + event = conn.execute( + "SELECT payload FROM task_events WHERE task_id=? AND kind='review_submitted' " + "ORDER BY id DESC LIMIT 1", (task_id,) + ).fetchone() + handoff = json.loads(event["payload"]) if event and event["payload"] else {} + implementer = _canonical_assignee(handoff.get("original_assignee")) or "" + if not implementer: + return None + remediation_key = f"review-remediation:{task_id}:{row['current_run_id']}" + remediation_id = create_task( + conn, title=f"Address review feedback: {row['title']}", + body=f"Review task: {task_id}\n\nChanges requested:\n{summary.strip()}", + assignee=implementer, created_by=row["assignee"] or "reviewer", + tenant=row["tenant"], priority=row["priority"], + workspace_kind=row["workspace_kind"], workspace_path=row["workspace_path"], + branch_name=row["branch_name"], project_id=row["project_id"], + skills=json.loads(row["skills"]) if row["skills"] else None, + idempotency_key=remediation_key, + ) + review_metadata = dict(metadata or {}) + review_metadata.update({"approved": False, "remediation_task_id": remediation_id, + "original_assignee": implementer}) + where = "id=? AND status='running'" + params: tuple[Any, ...] = (summary.strip(), int(time.time()), task_id) + if expected_run_id is not None: + where += " AND current_run_id=?" + params += (int(expected_run_id),) + cur = conn.execute( + "UPDATE tasks SET status='done', result=?, completed_at=?, claim_lock=NULL, " + "claim_expires=NULL, worker_pid=NULL WHERE " + where, + params, + ) + if cur.rowcount != 1: + return None + run_id = _end_run( + conn, task_id, outcome="changes_requested", status="done", + summary=summary.strip(), metadata=review_metadata, + ) + _append_event(conn, task_id, "review_changes_requested", review_metadata, run_id=run_id) + recompute_ready(conn) + return remediation_id + + def heartbeat_claim( conn: sqlite3.Connection, task_id: str, diff --git a/tests/hermes_cli/test_kanban_review_lifecycle.py b/tests/hermes_cli/test_kanban_review_lifecycle.py new file mode 100644 index 0000000000000..1f44f8bb5a218 --- /dev/null +++ b/tests/hermes_cli/test_kanban_review_lifecycle.py @@ -0,0 +1,101 @@ +"""Behavioral tests for the upstream-aligned native review lifecycle.""" + +from pathlib import Path + +import pytest + +from hermes_cli import kanban_db as kb + + +@pytest.fixture +def board(tmp_path, monkeypatch): + home = tmp_path / ".hermes" + home.mkdir() + monkeypatch.setenv("HERMES_HOME", str(home)) + monkeypatch.setattr(Path, "home", lambda: tmp_path) + kb.init_db() + return kb.connect() + + +def test_implementation_handoff_is_claimable_by_reviewer(board): + with board as conn: + task_id = kb.create_task(conn, title="implement", assignee="dev") + implementation = kb.claim_task(conn, task_id, claimer="worker:dev") + assert implementation is not None + assert kb.submit_for_review( + conn, + task_id, + reviewer="reviewer", + summary="PR opened; focused tests pass", + metadata={"pr_url": "https://github.com/acme/repo/pull/1", "tests_run": 3}, + expected_run_id=implementation.current_run_id, + ) + task = kb.get_task(conn, task_id) + assert task.status == "review" + assert task.assignee == "reviewer" + assert task.claim_lock is None + review = kb.claim_review_task(conn, task_id, claimer="worker:reviewer") + assert review is not None + assert review.status == "running" + assert review.assignee == "reviewer" + + +def test_review_approval_completes_and_changes_create_one_remediation(board): + with board as conn: + task_id = kb.create_task(conn, title="implement", assignee="dev") + implementation = kb.claim_task(conn, task_id, claimer="worker:dev") + assert implementation is not None + assert kb.submit_for_review( + conn, task_id, reviewer="reviewer", summary="ready", expected_run_id=implementation.current_run_id + ) + review = kb.claim_review_task(conn, task_id, claimer="worker:reviewer") + assert review is not None + remediation_id = kb.request_review_changes( + conn, task_id, summary="Fix the regression test", expected_run_id=review.current_run_id + ) + assert remediation_id + remediation = kb.get_task(conn, remediation_id) + assert remediation is not None + assert remediation.assignee == "dev" + assert remediation.status == "ready" + assert kb.get_task(conn, task_id).status == "done" + # The closed review card is terminal; replaying the same reviewer run + # cannot create a second remediation. + assert kb.request_review_changes(conn, task_id, summary="Fix the regression test") is None + rows = conn.execute( + "SELECT COUNT(*) AS n FROM tasks WHERE idempotency_key LIKE ?", + (f"review-remediation:{task_id}:%",), + ).fetchone() + assert rows["n"] == 1 + + +def test_review_approval_preserves_proof_and_scheduled_is_not_dispatchable(board): + with board as conn: + task_id = kb.create_task(conn, title="implement", assignee="dev") + implementation = kb.claim_task(conn, task_id, claimer="worker:dev") + assert implementation is not None + assert kb.submit_for_review( + conn, + task_id, + reviewer="reviewer", + summary="Evidence attached", + metadata={"commit": "abc123", "changed_files": ["src/example.py"]}, + expected_run_id=implementation.current_run_id, + ) + review = kb.claim_review_task(conn, task_id, claimer="worker:reviewer") + assert review is not None + assert kb.complete_task( + conn, + task_id, + summary="Approved after independent review", + metadata={"approved": True, "commit": "abc123"}, + expected_run_id=review.current_run_id, + ) + run = kb.latest_run(conn, task_id) + assert run is not None + assert run.metadata["approved"] is True + + scheduled_id = kb.create_task(conn, title="later", assignee="dev") + assert kb.schedule_task(conn, scheduled_id, reason="wait for release") + assert kb.claim_task(conn, scheduled_id) is None + assert kb.get_task(conn, scheduled_id).status == "scheduled" diff --git a/tools/kanban_tools.py b/tools/kanban_tools.py index 76f8902790367..303c416ba9b34 100644 --- a/tools/kanban_tools.py +++ b/tools/kanban_tools.py @@ -875,6 +875,52 @@ def _handle_block(args: dict, **kw) -> str: return tool_error(f"kanban_block: {e}") +def _handle_submit_review(args: dict, **kw) -> str: + """Route a completed implementation to the canonical review lane.""" + tid = _default_task_id(args.get("task_id")) + reviewer = str(args.get("reviewer") or "").strip() + summary = str(args.get("summary") or "").strip() + if not tid or not reviewer or not summary: + return tool_error("task_id, reviewer, and summary are required") + try: + kb, conn = _connect(board=args.get("board")) + try: + ok = kb.submit_for_review( + conn, tid, reviewer=reviewer, summary=summary, + metadata=args.get("metadata"), expected_run_id=_worker_run_id(tid), + ) + return _ok(task_id=tid, status="review") if ok else tool_error( + f"could not submit {tid} for review (not the active implementation run)" + ) + finally: + conn.close() + except Exception as e: + return tool_error(f"kanban_submit_review: {e}") + + +def _handle_review_changes(args: dict, **kw) -> str: + """Close a review and create an implementer remediation card.""" + tid = _default_task_id(args.get("task_id")) + summary = str(args.get("summary") or "").strip() + if not tid or not summary: + return tool_error("task_id and summary are required") + try: + kb, conn = _connect(board=args.get("board")) + try: + remediation = kb.request_review_changes( + conn, tid, summary=summary, metadata=args.get("metadata"), + expected_run_id=_worker_run_id(tid), + ) + return _ok(task_id=tid, status="done", remediation_task_id=remediation) \ + if remediation else tool_error( + f"could not request changes for {tid} (not the active review run)" + ) + finally: + conn.close() + except Exception as e: + return tool_error(f"kanban_review_changes: {e}") + + def _handle_heartbeat(args: dict, **kw) -> str: """Signal that the worker is still alive during a long operation. @@ -1750,6 +1796,44 @@ def _board_schema_prop() -> dict[str, str]: }, } +KANBAN_SUBMIT_REVIEW_SCHEMA = { + "name": "kanban_submit_review", + "description": ( + "Submit the active implementation run to the Review lane. Preserve " + "evidence in metadata and name the independent reviewer. Use this " + "instead of kanban_block for normal code-review handoff." + ), + "parameters": { + "type": "object", + "properties": { + "task_id": {"type": "string", "description": _DESC_TASK_ID_DEFAULT}, + "reviewer": {"type": "string", "description": "Reviewer profile."}, + "summary": {"type": "string", "description": "Review handoff summary."}, + "metadata": {"type": "object", "description": "Evidence: PR URL, commit, tests, changed files."}, + "board": _board_schema_prop(), + }, + "required": ["reviewer", "summary"], + }, +} + +KANBAN_REVIEW_CHANGES_SCHEMA = { + "name": "kanban_review_changes", + "description": ( + "Record review findings, complete the active Review card, and create " + "one idempotent remediation task assigned to the original implementer." + ), + "parameters": { + "type": "object", + "properties": { + "task_id": {"type": "string", "description": _DESC_TASK_ID_DEFAULT}, + "summary": {"type": "string", "description": "Requested changes and evidence."}, + "metadata": {"type": "object", "description": "Structured review findings."}, + "board": _board_schema_prop(), + }, + "required": ["summary"], + }, +} + KANBAN_HEARTBEAT_SCHEMA = { "name": "kanban_heartbeat", "description": ( @@ -2162,6 +2246,24 @@ def _board_schema_prop() -> dict[str, str]: emoji="⏸", ) +registry.register( + name="kanban_submit_review", + toolset="kanban", + schema=KANBAN_SUBMIT_REVIEW_SCHEMA, + handler=_handle_submit_review, + check_fn=_check_kanban_mode, + emoji="🔎", +) + +registry.register( + name="kanban_review_changes", + toolset="kanban", + schema=KANBAN_REVIEW_CHANGES_SCHEMA, + handler=_handle_review_changes, + check_fn=_check_kanban_mode, + emoji="🛠", +) + registry.register( name="kanban_heartbeat", toolset="kanban", diff --git a/website/docs/user-guide/features/kanban-worker-lanes.md b/website/docs/user-guide/features/kanban-worker-lanes.md index 69f879c6b1132..1aafbc652c051 100644 --- a/website/docs/user-guide/features/kanban-worker-lanes.md +++ b/website/docs/user-guide/features/kanban-worker-lanes.md @@ -56,15 +56,17 @@ Every claim must end in exactly one of: The kanban kernel enforces that exactly one of these terminates each run. A worker that calls neither and exits normally is treated as crashed. -## Outputs and the review-required convention +## Outputs and the Review lane -For most code-changing tasks, the work isn't truly *done* the moment the worker finishes — it needs a human reviewer. The kanban kernel doesn't enforce this distinction (a "code-changing task" is fuzzy and forcing block-instead-of-complete on every code worker would break flows where no review is wanted). It's a convention layered on top: +For code-changing tasks, implementation is handed to an independent reviewer rather than masquerading as a human blocker: -- **Block instead of complete**, with `reason` prefixed `review-required: ` so the dashboard / `hermes kanban show` surfaces the row as awaiting review. -- **Drop structured metadata into a `kanban_comment` first** since `kanban_block` only carries the human-readable `reason`. Comments are the durable annotation channel — every audit-relevant field (changed_files, tests_run, diff_path or PR url, decisions) belongs there. -- **Reviewer either approves and unblocks**, which respawns the worker with the comment thread for follow-ups; or asks for changes via another comment, which the next worker run sees as part of `kanban_show`'s context. +- Call `kanban_submit_review(reviewer=..., summary=..., metadata=...)` with the PR/commit, changed files, tests, and other evidence. +- The task moves from `running` to `review`, preserving the implementation run and assigning the reviewer. The dispatcher claims review cards separately, so the implementer is not respawned. +- A reviewer approves with `kanban_complete(summary=..., metadata={"approved": true, ...})`. +- A reviewer requesting changes calls `kanban_review_changes(summary=..., metadata=...)`; the review card completes with findings and one idempotent remediation task is created for the original implementer. +- Use `kanban_block(reason=...)` only for genuine human input, credentials, capability, dependency, or transient failures. Scheduled tasks remain time-gated and distinct from blocked work. -The injected `KANBAN_GUIDANCE` covers both `kanban_complete` (truly terminal tasks — typo fixes, docs changes, research writeups) and the `review-required` block pattern. +The injected `KANBAN_GUIDANCE` covers both `kanban_complete` (truly terminal tasks) and the explicit Review-lane handoff. ## Logs and audit trail From 0acad10fdb02e2094c06d24b5a7269ee16ea74ab Mon Sep 17 00:00:00 2001 From: SoLoVision Personal <126707742+solovision24@users.noreply.github.com> Date: Sat, 1 Aug 2026 13:53:39 -0400 Subject: [PATCH 2/4] fix(kanban): allow requeued review workers past PR guard (#9) * fix(kanban): allow requeued review workers past PR guard * fix(kanban): preserve review routing after crash requeue * fix(kanban): preserve native review lane on crash * fix(kanban): apply retry guards to native reviews * fix(kanban): guard native review respawns during cooldown --------- Co-authored-by: SoLo --- hermes_cli/kanban_db.py | 105 ++++++++++++--- tests/hermes_cli/test_kanban_db.py | 199 +++++++++++++++++++++++++++++ 2 files changed, 289 insertions(+), 15 deletions(-) diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index 104334864280b..a394022061cdf 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -4232,6 +4232,7 @@ def claim_task( *, ttl_seconds: Optional[int] = None, claimer: Optional[str] = None, + source_status: Optional[str] = None, ) -> Optional[Task]: """Atomically transition ``ready -> running``. @@ -4332,11 +4333,11 @@ def claim_task( "UPDATE tasks SET current_run_id = ? WHERE id = ?", (run_id, task_id), ) - _append_event( - conn, task_id, "claimed", - {"lock": lock, "expires": expires, "run_id": run_id}, - run_id=run_id, - ) + claim_payload = {"lock": lock, "expires": expires, "run_id": run_id} + if source_status is not None: + claim_payload["source_status"] = source_status + claim_payload["assignee"] = trow["assignee"] if trow else None + _append_event(conn, task_id, "claimed", claim_payload, run_id=run_id) claimed = get_task(conn, task_id) _fire_kanban_lifecycle_hook( "kanban_task_claimed", @@ -4417,7 +4418,7 @@ def claim_review_task( _append_event( conn, task_id, "claimed", {"lock": lock, "expires": expires, "run_id": run_id, - "source_status": "review"}, + "source_status": "review", "assignee": trow["assignee"] if trow else None}, run_id=run_id, ) return get_task(conn, task_id) @@ -7757,12 +7758,32 @@ def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]: event_payload["exit_kind"] = kind event_payload["exit_code"] = code + # A reviewer crash must return to the native review column. Do + # not flatten it into an implementation-style ``ready`` card: + # the review dispatcher owns the claim/spawn semantics and will + # create the next run with the sdlc-review skill. + latest_claim = conn.execute( + """ + SELECT json_extract(payload, '$.source_status') AS source_status + FROM task_events + WHERE task_id = ? AND kind = 'claimed' + ORDER BY id DESC + LIMIT 1 + """, + (row["id"],), + ).fetchone() + requeued_review = bool( + latest_claim and latest_claim["source_status"] == "review" + ) + requeue_status = "review" if requeued_review else "ready" + if requeued_review: + event_payload["source_status"] = "review" 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"]), + (requeue_status, row["id"], pid, row["claim_lock"]), ) if cur.rowcount == 1: # Rate-limited requeues are a clean release, not a crash — @@ -7849,8 +7870,9 @@ def detect_crashed_workers(conn: sqlite3.Connection) -> list[str]: else _PROTOCOL_VIOLATION_FAILURE_LIMIT ) if streak < violation_limit: - # Below budget: the task is already back at ``ready`` - # (respawn allowed) with ``last_failure_error`` stamped. + # Below-budget: the task is already back at ``ready`` or + # its native ``review`` column (respawn allowed) with + # ``last_failure_error`` stamped. # Deliberately no ``_record_task_failure`` call — a # below-budget violation must not consume the unified # failure budget, just as other failure kinds don't @@ -7990,7 +8012,7 @@ 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: @@ -8000,7 +8022,7 @@ def _record_task_failure( 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', 'running', 'review')", (failures, error[:500], task_id), ) run_id = None @@ -8127,6 +8149,21 @@ def _clear_failure_counter(conn: sqlite3.Connection, task_id: str) -> None: _clear_spawn_failures = _clear_failure_counter +def _latest_claim_was_review(conn: sqlite3.Connection, task_id: str) -> bool: + """Return whether the most recent claim came from the review lane.""" + row = conn.execute( + """ + SELECT json_extract(payload, '$.source_status') AS source_status + FROM task_events + WHERE task_id = ? AND kind = 'claimed' + ORDER BY id DESC + LIMIT 1 + """, + (task_id,), + ).fetchone() + return bool(row and row["source_status"] == "review") + + 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. @@ -8176,12 +8213,21 @@ def check_respawn_guard(conn: sqlite3.Connection, task_id: str) -> Optional[str] genuinely dead (no live PID on this host). """ row = conn.execute( - "SELECT last_failure_error FROM tasks WHERE id = ?", + "SELECT status, last_failure_error FROM tasks WHERE id = ?", (task_id,), ).fetchone() if row is None: return None + # Native review cards are already routed to the reviewer lane and must not + # be treated as duplicate implementation work merely because the + # canonical PR is present in the card's comments. A reviewer crash is + # requeued as ``review``; retain the claimed event's ``source_status`` so + # a ready retry can pass the same PR guard after a status transition. + review_claim = ( + row["status"] == "review" or _latest_claim_was_review(conn, task_id) + ) + now = int(time.time()) # 1. Rate-limit cooldown. The most recent run ended ``rate_limited`` @@ -8220,6 +8266,12 @@ def check_respawn_guard(conn: sqlite3.Connection, task_id: str) -> Optional[str] # crash/completion supersedes it. return None + # A newly submitted review has no prior review claim to identify it, but + # it is still not an implementation retry. Apply the rate-limit check + # above, then leave the native review lane alone. + if row["status"] == "review": + return None + # 2. Quota / auth blocker: retrying immediately will not help. err = row["last_failure_error"] if err and _RESPAWN_BLOCKER_RE.search(err): @@ -8253,10 +8305,16 @@ def check_respawn_guard(conn: sqlite3.Connection, task_id: str) -> Optional[str] # 4. GitHub PR URL in a recent comment — prior worker already opened a PR. pr_cutoff = now - _RESPAWN_GUARD_PR_WINDOW for c in conn.execute( - "SELECT body FROM task_comments WHERE task_id = ? AND created_at >= ?", + "SELECT body, created_at FROM task_comments WHERE task_id = ? AND created_at >= ?", (task_id, pr_cutoff), ).fetchall(): if c["body"] and _RESPAWN_GUARD_PR_URL_RE.search(c["body"]): + # A PR opened for a native review belongs to the reviewer, not a + # duplicate implementation retry. Only bypass when the review + # claim is newer than the PR comment, so an ordinary implementation + # task with an unrelated historical review event remains guarded. + if review_claim: + return None return "active_pr" return None @@ -8646,7 +8704,11 @@ def _dispatch_once_locked( _per_profile_running.get(row_assignee, 0) + 1 ) continue - claimed = claim_task(conn, row["id"], ttl_seconds=ttl_seconds) + claimed = claim_task( + conn, + row["id"], + ttl_seconds=ttl_seconds, + ) if claimed is None: continue try: @@ -8735,6 +8797,19 @@ def _dispatch_once_locked( if profile_exists is not None and not profile_exists(row["assignee"]): result.skipped_nonspawnable.append(row["id"]) continue + # Review cards bypass the ready-task loop, so apply the respawn guard + # here as well. Otherwise a rate-limited reviewer is claimed again + # on every dispatcher tick during its cooldown. + 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, "lane": "review"}, + ) + continue if dry_run: result.spawned.append((row["id"], row["assignee"], "")) continue diff --git a/tests/hermes_cli/test_kanban_db.py b/tests/hermes_cli/test_kanban_db.py index 7de9b4c0b7c72..ff78b94d14f15 100644 --- a/tests/hermes_cli/test_kanban_db.py +++ b/tests/hermes_cli/test_kanban_db.py @@ -359,10 +359,209 @@ def test_respawn_guard_defers_rate_limited_within_cooldown( assert kb.check_respawn_guard(conn, tid) is None +def test_respawn_guard_keeps_ordinary_pr_retry_protected(kanban_home): + """A normal implementation retry still cannot duplicate its open PR.""" + with kb.connect() as conn: + tid = kb.create_task(conn, title="implementation", assignee="dev") + kb.add_comment(conn, tid, "dev", "PR: https://github.com/acme/repo/pull/42") + assert kb.check_respawn_guard(conn, tid) == "active_pr" + + +def test_respawn_guard_does_not_trap_native_review_card(kanban_home): + """The canonical PR must not prevent a native review card from running.""" + with kb.connect() as conn: + tid = kb.create_task(conn, title="review", assignee="reviewer", initial_status="review") + kb.add_comment(conn, tid, "dev", "PR: https://github.com/acme/repo/pull/42") + assert kb.check_respawn_guard(conn, tid) is None + + +def test_dispatch_requeues_review_worker_into_review_lane(kanban_home, monkeypatch): + """A crashed reviewer keeps review routing when ready is dispatched.""" + import json + import hermes_cli.kanban_db as _kb + import hermes_cli.profiles as profiles + + monkeypatch.setattr(profiles, "profile_exists", lambda _name: True) + monkeypatch.setenv("HERMES_KANBAN_CRASH_GRACE_SECONDS", "0") + monkeypatch.setattr(_kb, "_pid_alive", lambda _pid: False) + spawned: list[tuple[str, list[str]]] = [] + + def spawn(task, _workspace): + spawned.append((task.id, list(task.skills or []))) + return None + + with kb.connect() as conn: + review_id = kb.create_task( + conn, title="review", assignee="reviewer", initial_status="review", + ) + kb.add_comment(conn, review_id, "dev", "PR: https://github.com/acme/repo/pull/42") + host = _kb._claimer_id().split(":", 1)[0] + assert kb.claim_review_task(conn, review_id, claimer=f"{host}:reviewer") + _kb._set_worker_pid(conn, review_id, 98765) + implementation_id = kb.create_task( + conn, title="implementation", assignee="dev", + ) + kb.add_comment( + conn, implementation_id, "dev", + "PR: https://github.com/acme/repo/pull/43", + ) + + assert kb.detect_crashed_workers(conn) == [review_id] + result = kb.dispatch_once(conn, spawn_fn=spawn) + + assert len(result.spawned) == 1 + assert result.spawned[0][0] == review_id + assert spawned == [(review_id, ["sdlc-review"])] + assert (implementation_id, "active_pr") in result.respawn_guarded + claim = conn.execute( + "SELECT payload FROM task_events WHERE task_id=? AND kind='claimed' " + "ORDER BY id DESC LIMIT 1", + (review_id,), + ).fetchone() + assert claim is not None + assert json.loads(claim["payload"])["source_status"] == "review" + + +def test_respawn_guard_allows_requeued_review_worker_after_pr(kanban_home, monkeypatch): + """A crashed reviewer requeued to ready retains reviewer execution intent.""" + import hermes_cli.kanban_db as _kb + + with kb.connect() as conn: + tid = kb.create_task(conn, title="review", assignee="reviewer", initial_status="review") + kb.add_comment(conn, tid, "dev", "PR: https://github.com/acme/repo/pull/42") + host = _kb._claimer_id().split(":", 1)[0] + claimed = kb.claim_review_task(conn, tid, claimer=f"{host}:reviewer") + assert claimed is not None + _kb._set_worker_pid(conn, tid, 98765) + monkeypatch.setenv("HERMES_KANBAN_CRASH_GRACE_SECONDS", "0") + monkeypatch.setattr(_kb, "_pid_alive", lambda _pid: False) + assert kb.detect_crashed_workers(conn) == [tid] + requeued = kb.get_task(conn, tid) + assert requeued is not None + assert requeued.status == "review" + assert kb.check_respawn_guard(conn, tid) is None + + +def test_reviewer_crash_breaker_blocks_native_review_lane(kanban_home): + """Repeated reviewer crashes must trip the breaker instead of respawning.""" + with kb.connect() as conn: + tid = kb.create_task( + conn, title="review", assignee="reviewer", initial_status="review", + ) + + assert kb._record_task_failure( + conn, tid, "reviewer crashed once", outcome="crashed", + failure_limit=2, + ) is False + task = kb.get_task(conn, tid) + assert task is not None + assert task.status == "review" + assert task.consecutive_failures == 1 + + assert kb._record_task_failure( + conn, tid, "reviewer crashed twice", outcome="crashed", + failure_limit=2, + ) is True + task = kb.get_task(conn, tid) + assert task is not None + assert task.status == "blocked" + assert task.consecutive_failures == 2 + + +def test_review_respawn_guard_honors_rate_limit_cooldown(kanban_home, monkeypatch): + """A rate-limited reviewer must be deferred while its quota cools down.""" + import hermes_cli.kanban_db as _kb + + monkeypatch.setenv("HERMES_KANBAN_RATE_LIMIT_COOLDOWN_SECONDS", "300") + now = 5_000_000 + with kb.connect() as conn: + tid = kb.create_task( + conn, title="review", assignee="reviewer", initial_status="review", + ) + claimed = kb.claim_review_task(conn, tid) + assert claimed is not None + run_id = claimed.current_run_id + conn.execute( + "UPDATE task_runs SET outcome='rate_limited', status='rate_limited', " + "ended_at=? WHERE id=?", + (now, run_id), + ) + conn.execute( + "UPDATE tasks SET status='review', current_run_id=NULL, " + "claim_lock=NULL, claim_expires=NULL, worker_pid=NULL, " + "last_failure_error=? WHERE id=?", + ("pid 1 exited rate-limited (quota wall) — requeued", tid), + ) + conn.commit() + + monkeypatch.setattr(_kb.time, "time", lambda: now + 100) + assert kb.check_respawn_guard(conn, tid) == "rate_limit_cooldown" + + monkeypatch.setattr(_kb.time, "time", lambda: now + 400) + assert kb.check_respawn_guard(conn, tid) is None + + + + + + + + +def test_dispatch_review_lane_honors_rate_limit_cooldown(kanban_home, monkeypatch): + """The native review dispatcher must defer quota-wall reviewers.""" + import json + import hermes_cli.kanban_db as _kb + import hermes_cli.profiles as profiles + monkeypatch.setenv("HERMES_KANBAN_RATE_LIMIT_COOLDOWN_SECONDS", "300") + monkeypatch.setattr(profiles, "profile_exists", lambda _name: True) + now = 5_000_000 + spawned: list[tuple[str, list[str]]] = [] + def spawn(task, _workspace): + spawned.append((task.id, list(task.skills or []))) + return None + with kb.connect() as conn: + tid = kb.create_task( + conn, title="review", assignee="reviewer", initial_status="review", + ) + claimed = kb.claim_review_task(conn, tid) + assert claimed is not None + run_id = claimed.current_run_id + conn.execute( + "UPDATE task_runs SET outcome='rate_limited', status='rate_limited', " + "ended_at=? WHERE id=?", + (now, run_id), + ) + conn.execute( + "UPDATE tasks SET status='review', current_run_id=NULL, " + "claim_lock=NULL, claim_expires=NULL, worker_pid=NULL, " + "last_failure_error=? WHERE id=?", + ("pid 1 exited rate-limited (quota wall) — requeued", tid), + ) + conn.commit() + monkeypatch.setattr(_kb.time, "time", lambda: now + 100) + result = kb.dispatch_once(conn, spawn_fn=spawn) + assert result.spawned == [] + assert (tid, "rate_limit_cooldown") in result.respawn_guarded + assert spawned == [] + event = conn.execute( + "SELECT payload FROM task_events WHERE task_id=? " + "AND kind='respawn_guarded' ORDER BY id DESC LIMIT 1", + (tid,), + ).fetchone() + assert event is not None + assert json.loads(event["payload"]) == { + "reason": "rate_limit_cooldown", "lane": "review", + } + + monkeypatch.setattr(_kb.time, "time", lambda: now + 400) + result = kb.dispatch_once(conn, spawn_fn=spawn) + assert len(result.spawned) == 1 + assert result.spawned[0][0] == tid + assert spawned == [(tid, ["sdlc-review"])] # --------------------------------------------------------------------------- From 232d6b8994962d4490f0e7890321ba132ecd75ba Mon Sep 17 00:00:00 2001 From: SoLo Date: Sun, 2 Aug 2026 08:50:28 -0400 Subject: [PATCH 3/4] feat(kanban): enforce native review handoff evidence --- agent/prompt_builder.py | 10 +- hermes_cli/kanban_db.py | 64 +++++++- .../test_kanban_review_lifecycle.py | 72 ++++++--- tools/kanban_tools.py | 153 ++++++++++++++++++ toolsets.py | 2 + .../features/kanban-worker-lanes.md | 6 +- 6 files changed, 270 insertions(+), 37 deletions(-) diff --git a/agent/prompt_builder.py b/agent/prompt_builder.py index c39eb4970ae17..2f00d5d449a84 100644 --- a/agent/prompt_builder.py +++ b/agent/prompt_builder.py @@ -243,11 +243,11 @@ def _strip_yaml_frontmatter(content: str) -> str: "(`{changed_files: [...], tests_run: N, decisions: [...]}`). Downstream " "workers read both via their own `kanban_show`. Never put secrets / " "tokens / raw PII in either field — run rows are durable forever. " - "Exception: if your output is a code change that needs independent review, " - "call `kanban_submit_review(reviewer=..., summary=..., metadata=...)`. " - "It preserves implementation evidence and routes the card to the Review " - "lane; `kanban_block` remains for genuine human input, credentials, " - "capability, dependency, or transient failures.\n" + "Exception: code changes needing review should use " + "`kanban_submit_review` when a real reviewer and immutable PR evidence " + "(pr_url, repo, number, head_sha, verification_evidence) exist. " + "Otherwise put changed_files / tests_run / diff_path in a " + "`kanban_comment`, then `kanban_block(reason=\"review-required: ...\")`.\n" "6. **If follow-up work appears, create it; don't do it.** Use " "`kanban_create(title=..., assignee=, parents=[your-task-id])` " "to spawn a child task for the appropriate specialist profile instead of " diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index a394022061cdf..3d00e9fae2212 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -4424,6 +4424,44 @@ def claim_review_task( return get_task(conn, task_id) +_GITHUB_PR_URL_RE = re.compile( + r"^https://(?:www\.)?github\.com/([^/]+)/([^/]+)/pull/([1-9][0-9]*)(?:[/?#].*)?$", + re.IGNORECASE, +) +_GITHUB_HEAD_SHA_RE = re.compile(r"^[0-9a-f]{7,64}$", re.IGNORECASE) + + +def _canonical_review_metadata(metadata: Optional[dict]) -> tuple[dict, str]: + """Validate and canonicalize immutable PR identity and proof.""" + if not isinstance(metadata, dict): + raise ValueError("review metadata must include pr_url, repo, number, head_sha, and verification_evidence") + raw_url = metadata.get("pr_url") + match = _GITHUB_PR_URL_RE.fullmatch(raw_url.strip()) if isinstance(raw_url, str) else None + if not match: + raise ValueError("pr_url must be an HTTPS GitHub pull-request URL") + repo = f"{match.group(1)}/{match.group(2)}".casefold() + if not isinstance(metadata.get("repo"), str) or metadata["repo"].strip().casefold() != repo: + raise ValueError("review metadata repo must match pr_url") + number = int(match.group(3)) + try: + if int(metadata.get("number")) != number: + raise ValueError + except (TypeError, ValueError): + raise ValueError("review metadata number must match pr_url") from None + head_sha = metadata.get("head_sha") + if not isinstance(head_sha, str) or not _GITHUB_HEAD_SHA_RE.fullmatch(head_sha.strip()): + raise ValueError("review metadata requires a valid immutable head_sha") + evidence = metadata.get("verification_evidence") or metadata.get("verification") + if evidence in (None, {}, [], ""): + raise ValueError("review metadata requires verification_evidence") + canonical = dict(metadata) + canonical.update({"pr_url": f"https://github.com/{repo}/pull/{number}", "repo": repo, + "number": number, "head_sha": head_sha.strip().casefold(), + "verification_evidence": evidence}) + canonical.pop("verification", None) + return canonical, f"github-pr:{repo}:{number}:{canonical['head_sha']}" + + def submit_for_review( conn: sqlite3.Connection, task_id: str, @@ -4445,12 +4483,27 @@ def submit_for_review( raise ValueError("reviewer is required") if not summary or not summary.strip(): raise ValueError("review summary is required") + review_metadata, review_identity = _canonical_review_metadata(metadata) + try: + from hermes_cli.profiles import profile_exists + except Exception as exc: + raise ValueError(f"reviewer profile {reviewer!r} cannot be resolved") from exc + if not profile_exists(reviewer): + raise ValueError(f"reviewer profile {reviewer!r} does not exist") with write_txn(conn): row = conn.execute( - "SELECT assignee, status FROM tasks WHERE id = ?", (task_id,) + "SELECT assignee, status, current_run_id FROM tasks WHERE id = ?", (task_id,) ).fetchone() if row is None or row["status"] != "running": return False + if expected_run_id is not None and row["current_run_id"] != int(expected_run_id): + return False + duplicate = conn.execute( + "SELECT id FROM tasks WHERE idempotency_key = ? AND id != ? " + "AND status IN ('review', 'running') LIMIT 1", (review_identity, task_id) + ).fetchone() + if duplicate: + return False original_assignee = str(row["assignee"] or "") where = "id = ? AND status = 'running'" params: tuple[Any, ...] = (task_id,) @@ -4458,14 +4511,15 @@ def submit_for_review( where += " AND current_run_id = ?" params += (int(expected_run_id),) cur = conn.execute( - "UPDATE tasks SET status='review', assignee=?, claim_lock=NULL, " + "UPDATE tasks SET status='review', assignee=?, idempotency_key=?, claim_lock=NULL, " "claim_expires=NULL, worker_pid=NULL WHERE " + where, - (reviewer, *params), + (reviewer, review_identity, *params), ) if cur.rowcount != 1: return False - handoff = dict(metadata or {}) - handoff.update({"reviewer": reviewer, "original_assignee": original_assignee}) + handoff = dict(review_metadata) + handoff.update({"reviewer": reviewer, "original_assignee": original_assignee, + "review_identity": review_identity}) run_id = _end_run( conn, task_id, outcome="submitted_for_review", status="review", summary=summary.strip(), metadata=handoff, diff --git a/tests/hermes_cli/test_kanban_review_lifecycle.py b/tests/hermes_cli/test_kanban_review_lifecycle.py index 1f44f8bb5a218..991b8debdda07 100644 --- a/tests/hermes_cli/test_kanban_review_lifecycle.py +++ b/tests/hermes_cli/test_kanban_review_lifecycle.py @@ -7,6 +7,15 @@ from hermes_cli import kanban_db as kb +REVIEW_METADATA = { + "pr_url": "https://github.com/acme/repo/pull/1", + "repo": "acme/repo", + "number": 1, + "head_sha": "a" * 40, + "verification_evidence": {"tests_passed": 3}, +} + + @pytest.fixture def board(tmp_path, monkeypatch): home = tmp_path / ".hermes" @@ -23,21 +32,17 @@ def test_implementation_handoff_is_claimable_by_reviewer(board): implementation = kb.claim_task(conn, task_id, claimer="worker:dev") assert implementation is not None assert kb.submit_for_review( - conn, - task_id, - reviewer="reviewer", - summary="PR opened; focused tests pass", - metadata={"pr_url": "https://github.com/acme/repo/pull/1", "tests_run": 3}, - expected_run_id=implementation.current_run_id, + conn, task_id, reviewer="default", summary="PR opened; focused tests pass", + metadata=REVIEW_METADATA, expected_run_id=implementation.current_run_id, ) task = kb.get_task(conn, task_id) assert task.status == "review" - assert task.assignee == "reviewer" + assert task.assignee == "default" assert task.claim_lock is None - review = kb.claim_review_task(conn, task_id, claimer="worker:reviewer") + review = kb.claim_review_task(conn, task_id, claimer="worker:default") assert review is not None assert review.status == "running" - assert review.assignee == "reviewer" + assert review.assignee == "default" def test_review_approval_completes_and_changes_create_one_remediation(board): @@ -46,9 +51,10 @@ def test_review_approval_completes_and_changes_create_one_remediation(board): implementation = kb.claim_task(conn, task_id, claimer="worker:dev") assert implementation is not None assert kb.submit_for_review( - conn, task_id, reviewer="reviewer", summary="ready", expected_run_id=implementation.current_run_id + conn, task_id, reviewer="default", summary="ready", metadata=REVIEW_METADATA, + expected_run_id=implementation.current_run_id, ) - review = kb.claim_review_task(conn, task_id, claimer="worker:reviewer") + review = kb.claim_review_task(conn, task_id, claimer="worker:default") assert review is not None remediation_id = kb.request_review_changes( conn, task_id, summary="Fix the regression test", expected_run_id=review.current_run_id @@ -59,8 +65,6 @@ def test_review_approval_completes_and_changes_create_one_remediation(board): assert remediation.assignee == "dev" assert remediation.status == "ready" assert kb.get_task(conn, task_id).status == "done" - # The closed review card is terminal; replaying the same reviewer run - # cannot create a second remediation. assert kb.request_review_changes(conn, task_id, summary="Fix the regression test") is None rows = conn.execute( "SELECT COUNT(*) AS n FROM tasks WHERE idempotency_key LIKE ?", @@ -75,21 +79,16 @@ def test_review_approval_preserves_proof_and_scheduled_is_not_dispatchable(board implementation = kb.claim_task(conn, task_id, claimer="worker:dev") assert implementation is not None assert kb.submit_for_review( - conn, - task_id, - reviewer="reviewer", - summary="Evidence attached", - metadata={"commit": "abc123", "changed_files": ["src/example.py"]}, + conn, task_id, reviewer="default", summary="Evidence attached", + metadata={**REVIEW_METADATA, "head_sha": "b" * 40, + "commit": "abc123", "changed_files": ["src/example.py"]}, expected_run_id=implementation.current_run_id, ) - review = kb.claim_review_task(conn, task_id, claimer="worker:reviewer") + review = kb.claim_review_task(conn, task_id, claimer="worker:default") assert review is not None assert kb.complete_task( - conn, - task_id, - summary="Approved after independent review", - metadata={"approved": True, "commit": "abc123"}, - expected_run_id=review.current_run_id, + conn, task_id, summary="Approved after independent review", + metadata={"approved": True, "commit": "abc123"}, expected_run_id=review.current_run_id, ) run = kb.latest_run(conn, task_id) assert run is not None @@ -99,3 +98,28 @@ def test_review_approval_preserves_proof_and_scheduled_is_not_dispatchable(board assert kb.schedule_task(conn, scheduled_id, reason="wait for release") assert kb.claim_task(conn, scheduled_id) is None assert kb.get_task(conn, scheduled_id).status == "scheduled" + + +def test_unknown_reviewer_and_duplicate_head_do_not_mutate(board, monkeypatch): + from hermes_cli import profiles + + monkeypatch.setattr(profiles, "profile_exists", lambda name: name == "default") + with board as conn: + first_id = kb.create_task(conn, title="first", assignee="dev") + first_run = kb.claim_task(conn, first_id) + assert first_run is not None + with pytest.raises(ValueError, match="does not exist"): + kb.submit_for_review(conn, first_id, reviewer="missing", summary="ready", + metadata=REVIEW_METADATA, expected_run_id=first_run.current_run_id) + assert kb.get_task(conn, first_id).status == "running" + assert kb.submit_for_review(conn, first_id, reviewer="default", summary="ready", + metadata=REVIEW_METADATA, + expected_run_id=first_run.current_run_id) + + second_id = kb.create_task(conn, title="duplicate", assignee="dev") + second_run = kb.claim_task(conn, second_id) + assert second_run is not None + assert not kb.submit_for_review(conn, second_id, reviewer="default", summary="duplicate", + metadata=REVIEW_METADATA, + expected_run_id=second_run.current_run_id) + assert kb.get_task(conn, second_id).status == "running" diff --git a/tools/kanban_tools.py b/tools/kanban_tools.py index 303c416ba9b34..eb0cd16d4754c 100644 --- a/tools/kanban_tools.py +++ b/tools/kanban_tools.py @@ -794,6 +794,96 @@ def _handle_complete(args: dict, **kw) -> str: return tool_error(f"kanban_complete: {e}") +def _handle_submit_review(args: dict, **kw) -> str: + """Submit the current implementation run to the native review lane.""" + tid = _default_task_id(args.get("task_id")) + if not tid: + return tool_error("task_id is required (or set HERMES_KANBAN_TASK)") + ownership_err = _enforce_worker_task_ownership(tid) + if ownership_err: + return ownership_err + reviewer = args.get("reviewer") + summary = args.get("summary") + metadata = args.get("metadata") + if not reviewer or not str(reviewer).strip(): + return tool_error("reviewer is required") + if not summary or not str(summary).strip(): + return tool_error("summary is required") + if not isinstance(metadata, dict): + return tool_error( + "metadata is required and must include pr_url, repo, number, " + "head_sha, and verification_evidence" + ) + try: + kb, conn = _connect(board=args.get("board")) + try: + ok = kb.submit_for_review( + conn, + tid, + reviewer=str(reviewer), + summary=str(summary), + metadata=_stamp_worker_session_metadata(tid, metadata), + expected_run_id=_worker_run_id(tid), + ) + if not ok: + return tool_error( + f"could not submit {tid} for review; it may already have " + "an active review for this PR head" + ) + task = kb.get_task(conn, tid) + return _ok( + task_id=tid, + status=task.status if task else "review", + reviewer=task.assignee if task else str(reviewer), + ) + finally: + conn.close() + except ValueError as exc: + return tool_error(f"kanban_submit_review: {exc}") + except Exception as exc: + logger.exception("kanban_submit_review failed") + return tool_error(f"kanban_submit_review: {exc}") + + +def _handle_review_changes(args: dict, **kw) -> str: + """Finish a review with changes requested and create remediation work.""" + tid = _default_task_id(args.get("task_id")) + if not tid: + return tool_error("task_id is required (or set HERMES_KANBAN_TASK)") + ownership_err = _enforce_worker_task_ownership(tid) + if ownership_err: + return ownership_err + summary = args.get("summary") + metadata = args.get("metadata") + if not summary or not str(summary).strip(): + return tool_error("summary is required") + if metadata is not None and not isinstance(metadata, dict): + return tool_error("metadata must be an object/dict") + try: + kb, conn = _connect(board=args.get("board")) + try: + remediation = kb.request_review_changes( + conn, + tid, + summary=str(summary), + metadata=metadata, + expected_run_id=_worker_run_id(tid), + ) + if not remediation: + return tool_error( + "could not request review changes; task is not this review run " + "or lacks implementer provenance" + ) + return _ok(task_id=tid, status="done", remediation_task_id=remediation) + finally: + conn.close() + except ValueError as exc: + return tool_error(f"kanban_review_changes: {exc}") + except Exception as exc: + logger.exception("kanban_review_changes failed") + return tool_error(f"kanban_review_changes: {exc}") + + def _handle_block(args: dict, **kw) -> str: """Transition the task to blocked with a reason a human will read.""" delegated_err = _reject_delegated_child_mutation("kanban_block") @@ -1752,6 +1842,51 @@ def _board_schema_prop() -> dict[str, str]: }, } +KANBAN_SUBMIT_REVIEW_SCHEMA = { + "name": "kanban_submit_review", + "description": ( + "Submit the current implementation task for independent review. " + "The kernel validates the reviewer profile, preserves the original " + "implementer, and requires immutable PR identity plus verification " + "evidence before moving the card to review." + ), + "parameters": { + "type": "object", + "properties": { + "task_id": {"type": "string", "description": _DESC_TASK_ID_DEFAULT}, + "reviewer": {"type": "string", "description": "Existing profile that will perform the review."}, + "summary": {"type": "string", "description": "Implementation handoff summary and evidence."}, + "metadata": { + "type": "object", + "description": ( + "Required immutable PR evidence: pr_url, repo, number, " + "head_sha, and verification_evidence." + ), + }, + "board": _board_schema_prop(), + }, + "required": ["reviewer", "summary", "metadata"], + }, +} + +KANBAN_REVIEW_CHANGES_SCHEMA = { + "name": "kanban_review_changes", + "description": ( + "Finish the current review with changes requested. Records the review " + "evidence and creates a ready remediation task for the original implementer." + ), + "parameters": { + "type": "object", + "properties": { + "task_id": {"type": "string", "description": _DESC_TASK_ID_DEFAULT}, + "summary": {"type": "string", "description": "Concrete findings and required changes."}, + "metadata": {"type": "object", "description": "Structured review findings."}, + "board": _board_schema_prop(), + }, + "required": ["summary"], + }, +} + KANBAN_BLOCK_SCHEMA = { "name": "kanban_block", "description": ( @@ -2237,6 +2372,24 @@ def _board_schema_prop() -> dict[str, str]: emoji="✔", ) +registry.register( + name="kanban_submit_review", + toolset="kanban", + schema=KANBAN_SUBMIT_REVIEW_SCHEMA, + handler=_handle_submit_review, + check_fn=_check_kanban_mode, + emoji="🔎", +) + +registry.register( + name="kanban_review_changes", + toolset="kanban", + schema=KANBAN_REVIEW_CHANGES_SCHEMA, + handler=_handle_review_changes, + check_fn=_check_kanban_mode, + emoji="↩", +) + registry.register( name="kanban_block", toolset="kanban", diff --git a/toolsets.py b/toolsets.py index f4bb3c3343a99..0d7a1fd5d6fb3 100644 --- a/toolsets.py +++ b/toolsets.py @@ -78,6 +78,7 @@ # tools/kanban_tools.py. "kanban_show", "kanban_list", "kanban_complete", "kanban_block", "kanban_heartbeat", + "kanban_submit_review", "kanban_review_changes", "kanban_comment", "kanban_create", "kanban_link", "kanban_unblock", "kanban_attach", "kanban_attach_url", "kanban_attachments", @@ -297,6 +298,7 @@ "tools": [ "kanban_show", "kanban_list", "kanban_complete", "kanban_block", "kanban_heartbeat", "kanban_comment", + "kanban_submit_review", "kanban_review_changes", "kanban_create", "kanban_link", "kanban_unblock", "kanban_attach", "kanban_attach_url", "kanban_attachments", diff --git a/website/docs/user-guide/features/kanban-worker-lanes.md b/website/docs/user-guide/features/kanban-worker-lanes.md index 1aafbc652c051..70f60b58f6b3c 100644 --- a/website/docs/user-guide/features/kanban-worker-lanes.md +++ b/website/docs/user-guide/features/kanban-worker-lanes.md @@ -58,12 +58,12 @@ The kanban kernel enforces that exactly one of these terminates each run. A work ## Outputs and the Review lane -For code-changing tasks, implementation is handed to an independent reviewer rather than masquerading as a human blocker: +For code-changing tasks, implementation is handed to an independent reviewer rather than masquerading as a human blocker. Native submission is the SoLo compatibility boundary: `kanban_submit_review` validates the reviewer profile before mutation, requires canonical `pr_url`, `repo`, `number`, `head_sha`, and `verification_evidence`, and stores the kernel-owned implementer provenance. -- Call `kanban_submit_review(reviewer=..., summary=..., metadata=...)` with the PR/commit, changed files, tests, and other evidence. - The task moves from `running` to `review`, preserving the implementation run and assigning the reviewer. The dispatcher claims review cards separately, so the implementer is not respawned. +- Active submissions with the same repository, PR number, and head SHA are rejected without creating a second review card. A new head is a new immutable review identity. - A reviewer approves with `kanban_complete(summary=..., metadata={"approved": true, ...})`. -- A reviewer requesting changes calls `kanban_review_changes(summary=..., metadata=...)`; the review card completes with findings and one idempotent remediation task is created for the original implementer. +- A reviewer requesting changes calls `kanban_review_changes(summary=..., metadata=...)`; the review card completes with findings and one remediation task is created for the original implementer. - Use `kanban_block(reason=...)` only for genuine human input, credentials, capability, dependency, or transient failures. Scheduled tasks remain time-gated and distinct from blocked work. The injected `KANBAN_GUIDANCE` covers both `kanban_complete` (truly terminal tasks) and the explicit Review-lane handoff. From e0c5b8200554797fe94282761bf7f4625400decb Mon Sep 17 00:00:00 2001 From: SoLo Date: Sun, 2 Aug 2026 09:03:45 -0400 Subject: [PATCH 4/4] fix(kanban): make native review handoff canonical --- agent/prompt_builder.py | 12 ++- hermes_cli/kanban_db.py | 2 +- .../test_kanban_review_lifecycle.py | 22 ++++ tools/kanban_tools.py | 102 ------------------ .../features/kanban-worker-lanes.md | 19 ++++ 5 files changed, 49 insertions(+), 108 deletions(-) diff --git a/agent/prompt_builder.py b/agent/prompt_builder.py index 2f00d5d449a84..91c58c49d51fd 100644 --- a/agent/prompt_builder.py +++ b/agent/prompt_builder.py @@ -243,11 +243,13 @@ def _strip_yaml_frontmatter(content: str) -> str: "(`{changed_files: [...], tests_run: N, decisions: [...]}`). Downstream " "workers read both via their own `kanban_show`. Never put secrets / " "tokens / raw PII in either field — run rows are durable forever. " - "Exception: code changes needing review should use " - "`kanban_submit_review` when a real reviewer and immutable PR evidence " - "(pr_url, repo, number, head_sha, verification_evidence) exist. " - "Otherwise put changed_files / tests_run / diff_path in a " - "`kanban_comment`, then `kanban_block(reason=\"review-required: ...\")`.\n" + "Exception: code changes needing independent review must call " + "`kanban_submit_review(reviewer=..., summary=..., metadata=...)` with " + "the open PR URL, repository, PR number, exact 40-character head SHA, " + "verification evidence, deployment implications, and original implementer. " + "It preserves implementation evidence and routes the card to the Review " + "lane; `kanban_block` remains only for genuine human input, credentials, " + "capability, dependency, or transient failures.\n" "6. **If follow-up work appears, create it; don't do it.** Use " "`kanban_create(title=..., assignee=, parents=[your-task-id])` " "to spawn a child task for the appropriate specialist profile instead of " diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index 3d00e9fae2212..811b66f62e3ba 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -4428,7 +4428,7 @@ def claim_review_task( r"^https://(?:www\.)?github\.com/([^/]+)/([^/]+)/pull/([1-9][0-9]*)(?:[/?#].*)?$", re.IGNORECASE, ) -_GITHUB_HEAD_SHA_RE = re.compile(r"^[0-9a-f]{7,64}$", re.IGNORECASE) +_GITHUB_HEAD_SHA_RE = re.compile(r"^[0-9a-f]{40}$", re.IGNORECASE) def _canonical_review_metadata(metadata: Optional[dict]) -> tuple[dict, str]: diff --git a/tests/hermes_cli/test_kanban_review_lifecycle.py b/tests/hermes_cli/test_kanban_review_lifecycle.py index 991b8debdda07..2ff20f4db3137 100644 --- a/tests/hermes_cli/test_kanban_review_lifecycle.py +++ b/tests/hermes_cli/test_kanban_review_lifecycle.py @@ -123,3 +123,25 @@ def test_unknown_reviewer_and_duplicate_head_do_not_mutate(board, monkeypatch): metadata=REVIEW_METADATA, expected_run_id=second_run.current_run_id) assert kb.get_task(conn, second_id).status == "running" + + +def test_review_handoff_rejects_abbreviated_head_sha_without_mutation(board, monkeypatch): + from hermes_cli import profiles + + monkeypatch.setattr(profiles, "profile_exists", lambda name: name == "default") + with board as conn: + task_id = kb.create_task(conn, title="implement", assignee="dev") + implementation = kb.claim_task(conn, task_id) + assert implementation is not None + with pytest.raises(ValueError, match="immutable head_sha"): + kb.submit_for_review( + conn, + task_id, + reviewer="default", + summary="ready", + metadata={**REVIEW_METADATA, "head_sha": "a" * 12}, + expected_run_id=implementation.current_run_id, + ) + task = kb.get_task(conn, task_id) + assert task is not None + assert task.status == "running" diff --git a/tools/kanban_tools.py b/tools/kanban_tools.py index eb0cd16d4754c..27aa78bf610e5 100644 --- a/tools/kanban_tools.py +++ b/tools/kanban_tools.py @@ -965,52 +965,6 @@ def _handle_block(args: dict, **kw) -> str: return tool_error(f"kanban_block: {e}") -def _handle_submit_review(args: dict, **kw) -> str: - """Route a completed implementation to the canonical review lane.""" - tid = _default_task_id(args.get("task_id")) - reviewer = str(args.get("reviewer") or "").strip() - summary = str(args.get("summary") or "").strip() - if not tid or not reviewer or not summary: - return tool_error("task_id, reviewer, and summary are required") - try: - kb, conn = _connect(board=args.get("board")) - try: - ok = kb.submit_for_review( - conn, tid, reviewer=reviewer, summary=summary, - metadata=args.get("metadata"), expected_run_id=_worker_run_id(tid), - ) - return _ok(task_id=tid, status="review") if ok else tool_error( - f"could not submit {tid} for review (not the active implementation run)" - ) - finally: - conn.close() - except Exception as e: - return tool_error(f"kanban_submit_review: {e}") - - -def _handle_review_changes(args: dict, **kw) -> str: - """Close a review and create an implementer remediation card.""" - tid = _default_task_id(args.get("task_id")) - summary = str(args.get("summary") or "").strip() - if not tid or not summary: - return tool_error("task_id and summary are required") - try: - kb, conn = _connect(board=args.get("board")) - try: - remediation = kb.request_review_changes( - conn, tid, summary=summary, metadata=args.get("metadata"), - expected_run_id=_worker_run_id(tid), - ) - return _ok(task_id=tid, status="done", remediation_task_id=remediation) \ - if remediation else tool_error( - f"could not request changes for {tid} (not the active review run)" - ) - finally: - conn.close() - except Exception as e: - return tool_error(f"kanban_review_changes: {e}") - - def _handle_heartbeat(args: dict, **kw) -> str: """Signal that the worker is still alive during a long operation. @@ -1931,44 +1885,6 @@ def _board_schema_prop() -> dict[str, str]: }, } -KANBAN_SUBMIT_REVIEW_SCHEMA = { - "name": "kanban_submit_review", - "description": ( - "Submit the active implementation run to the Review lane. Preserve " - "evidence in metadata and name the independent reviewer. Use this " - "instead of kanban_block for normal code-review handoff." - ), - "parameters": { - "type": "object", - "properties": { - "task_id": {"type": "string", "description": _DESC_TASK_ID_DEFAULT}, - "reviewer": {"type": "string", "description": "Reviewer profile."}, - "summary": {"type": "string", "description": "Review handoff summary."}, - "metadata": {"type": "object", "description": "Evidence: PR URL, commit, tests, changed files."}, - "board": _board_schema_prop(), - }, - "required": ["reviewer", "summary"], - }, -} - -KANBAN_REVIEW_CHANGES_SCHEMA = { - "name": "kanban_review_changes", - "description": ( - "Record review findings, complete the active Review card, and create " - "one idempotent remediation task assigned to the original implementer." - ), - "parameters": { - "type": "object", - "properties": { - "task_id": {"type": "string", "description": _DESC_TASK_ID_DEFAULT}, - "summary": {"type": "string", "description": "Requested changes and evidence."}, - "metadata": {"type": "object", "description": "Structured review findings."}, - "board": _board_schema_prop(), - }, - "required": ["summary"], - }, -} - KANBAN_HEARTBEAT_SCHEMA = { "name": "kanban_heartbeat", "description": ( @@ -2399,24 +2315,6 @@ def _board_schema_prop() -> dict[str, str]: emoji="⏸", ) -registry.register( - name="kanban_submit_review", - toolset="kanban", - schema=KANBAN_SUBMIT_REVIEW_SCHEMA, - handler=_handle_submit_review, - check_fn=_check_kanban_mode, - emoji="🔎", -) - -registry.register( - name="kanban_review_changes", - toolset="kanban", - schema=KANBAN_REVIEW_CHANGES_SCHEMA, - handler=_handle_review_changes, - check_fn=_check_kanban_mode, - emoji="🛠", -) - registry.register( name="kanban_heartbeat", toolset="kanban", diff --git a/website/docs/user-guide/features/kanban-worker-lanes.md b/website/docs/user-guide/features/kanban-worker-lanes.md index 70f60b58f6b3c..c4f7e83e8cc41 100644 --- a/website/docs/user-guide/features/kanban-worker-lanes.md +++ b/website/docs/user-guide/features/kanban-worker-lanes.md @@ -68,6 +68,25 @@ For code-changing tasks, implementation is handed to an independent reviewer rat The injected `KANBAN_GUIDANCE` covers both `kanban_complete` (truly terminal tasks) and the explicit Review-lane handoff. +### Upstream boundary and upgrade checklist + +The native Review lane is intentionally a thin SoLoVision boundary over the upstream Kanban lifecycle: + +| Contract | Upstream behavior | Boundary behavior | Upgrade check | +| --- | --- | --- | --- | +| Completion | `kanban_complete` ends a run as `done` | unchanged | Run completion and goal-mode tests | +| Genuine blockers | `kanban_block` routes dependency/input/capability/transient failures | unchanged; never use it as a review queue | Run blocked-task and requeue tests | +| Review handoff | not an upstream terminal status | `kanban_submit_review` moves the same card to `review` with immutable PR proof | Run lifecycle, duplicate-head, and SHA validation tests | +| Review outcome | normal completion semantics | approval completes the review; changes requested create one remediation card | Run approval/remediation tests | + +Before upgrading the upstream Kanban implementation: + +1. Rebase this boundary onto the current upstream default branch, not a stale fork branch. +2. Confirm one definition, schema, and registry entry exists for each native Review tool. +3. Run `tests/hermes_cli/test_kanban_review_lifecycle.py` and the Kanban database/tool suites. +4. Verify the handoff still requires the open PR URL, matching repository/number, exact 40-character head SHA, verification evidence, deployment implications, and original implementer provenance. +5. Inspect the diff for restored `review-required` guidance or abbreviated SHA acceptance before opening the next PR. + ## Logs and audit trail The dispatcher writes per-task worker stdout/stderr to `/logs/.log`. Logs are auditable from kanban metadata: