fix(kanban): enforce review lifecycle invariants

This commit is contained in:
Jakub Wolniewicz
2026-07-31 15:17:00 +02:00
committed by Teknium
parent 4ab998a7de
commit 6d7e86c262
8 changed files with 384 additions and 69 deletions
+6 -1
View File
@@ -2428,7 +2428,12 @@ def _cmd_request_review(args: argparse.Namespace) -> int:
file=sys.stderr,
)
return 1
print(f"Requested review for {tid}" + (f": {summary}" if summary else ""))
persisted_run = kb.latest_run(conn, tid)
display_summary = persisted_run.summary if persisted_run else None
print(
f"Requested review for {tid}"
+ (f": {display_summary}" if display_summary else "")
)
return 0
+50 -5
View File
@@ -4217,6 +4217,8 @@ def _resume_status_from_events(conn: sqlite3.Connection, task_id: str) -> str:
payload = json.loads(row["payload"]) if row and row["payload"] else {}
except (json.JSONDecodeError, TypeError):
payload = {}
if not isinstance(payload, dict):
payload = {}
for key in ("resume_status", "retry_status", "source_status"):
if payload.get(key) == "review":
return "review"
@@ -4572,6 +4574,8 @@ def _retry_status_for_run(
payload = json.loads(event["payload"]) if event and event["payload"] else {}
except (json.JSONDecodeError, TypeError):
payload = {}
if not isinstance(payload, dict):
payload = {}
return "review" if payload.get("source_status") == "review" else "ready"
@@ -5077,6 +5081,11 @@ def complete_task(
# ``review`` or ``running``.
if not _parents_satisfied(conn, task_id):
return False
prior = conn.execute(
"SELECT status FROM tasks WHERE id = ?",
(task_id,),
).fetchone()
prior_status = prior["status"] if prior else None
if expected_run_id is None:
cur = conn.execute(
"""
@@ -5136,18 +5145,31 @@ def complete_task(
# blocked → done with no run in flight), synthesize a
# zero-duration run so the handoff fields are persisted in
# attempt history instead of silently lost.
if run_id is None and (summary or metadata or result):
if run_id is None and (
summary or metadata or result or prior_status == "review"
):
synth_summary = summary if summary is not None else result
synth_metadata = metadata
if prior_status == "review" and not synth_summary and not synth_metadata:
synth_summary = "Review approved without additional evidence."
synth_metadata = {
"source_status": "review",
"approval": "manual",
}
run_id = _synthesize_ended_run(
conn, task_id,
outcome="completed",
summary=summary if summary is not None else result,
metadata=metadata,
summary=synth_summary,
metadata=synth_metadata,
)
# Carry the handoff summary in the event payload so gateway
# notifiers and dashboard WS consumers can render it without a
# second SQL round-trip. First line only, 400 char cap — the
# full summary stays on the run row.
_ev_lines = ((summary if summary is not None else result) or "").strip().splitlines()
event_summary = summary if summary is not None else result
if prior_status == "review" and not event_summary:
event_summary = "Review approved without additional evidence."
_ev_lines = (event_summary or "").strip().splitlines()
ev_summary = _ev_lines[0][:400] if _ev_lines else ""
completed_payload: dict = {
"result_len": len(result) if result else 0,
@@ -6016,6 +6038,21 @@ def block_task(
def _redact_review_value(value: Any) -> Any:
"""Redact secrets at the domain boundary for durable review handoffs."""
if isinstance(value, str):
from agent.redact import redact_sensitive_text
return redact_sensitive_text(value, force=True)
if isinstance(value, dict):
return {key: _redact_review_value(item) for key, item in value.items()}
if isinstance(value, list):
return [_redact_review_value(item) for item in value]
if isinstance(value, tuple):
return tuple(_redact_review_value(item) for item in value)
return value
def request_review(
conn: sqlite3.Connection,
task_id: str,
@@ -6033,6 +6070,8 @@ def request_review(
right profile. Supplying ``reviewer`` also reassigns the task before it is
exposed to the review dispatcher.
"""
summary = _redact_review_value(summary)
metadata = _redact_review_value(metadata)
with write_txn(conn):
if not _parents_satisfied(conn, task_id):
return False
@@ -6117,7 +6156,7 @@ def request_changes(
``changes_requested`` event. The second tuple item is the implementer on
success or a diagnostic reason on failure.
"""
reason = (reason or "").strip()
reason = str(_redact_review_value(reason or "")).strip()
if not reason:
return False, "reason is required"
@@ -6148,6 +6187,8 @@ def request_changes(
)
except (json.JSONDecodeError, TypeError):
claimed_payload = {}
if not isinstance(claimed_payload, dict):
claimed_payload = {}
if claimed_payload.get("source_status") != "review":
return False, "active run was not claimed from review"
@@ -6167,6 +6208,8 @@ def request_changes(
)
except (json.JSONDecodeError, TypeError):
requested_payload = {}
if not isinstance(requested_payload, dict):
requested_payload = {}
implementer = requested_payload.get("implementer")
if not isinstance(implementer, str) or not implementer.strip():
return False, "review handoff has no valid implementer provenance"
@@ -9146,6 +9189,7 @@ def _dispatch_once_locked(
continue
if dry_run:
result.spawned.append((row["id"], row_assignee, ""))
spawned += 1
# Increment per-profile counter even in dry_run so the cap
# check sees the would-be spawn on subsequent iterations.
# Without this, dry_run reports every task as spawnable and
@@ -9269,6 +9313,7 @@ def _dispatch_once_locked(
continue
if dry_run:
result.spawned.append((row["id"], row["assignee"], ""))
spawned += 1
if _per_profile_cap is not None:
_per_profile_running[row["assignee"]] = (
_per_profile_running.get(row["assignee"], 0) + 1
+110 -45
View File
@@ -1070,6 +1070,69 @@ def _parents_blocking_ready(
]
def _invalidate_descendants_for_parent_reopen(
conn: sqlite3.Connection,
parent_id: str,
terminations: list[tuple[Optional[int], Optional[str]]],
) -> None:
"""Retract every dispatchable/completed descendant of a reopened parent."""
rows = conn.execute(
"""
WITH RECURSIVE descendants(id) AS (
SELECT child_id FROM task_links WHERE parent_id = ?
UNION
SELECT l.child_id
FROM task_links l
JOIN descendants d ON d.id = l.parent_id
)
SELECT t.id, t.status, t.current_run_id, t.worker_pid, t.claim_lock
FROM descendants d
JOIN tasks t ON t.id = d.id
ORDER BY t.id
""",
(parent_id,),
).fetchall()
for row in rows:
previous_status = row["status"]
if previous_status not in {"ready", "review", "running", "done"}:
continue
resume_status = "ready"
run_id = None
if previous_status == "review":
resume_status = "review"
elif previous_status == "running":
resume_status = kanban_db._retry_status_for_run(
conn, row["id"], row["current_run_id"]
)
terminations.append((row["worker_pid"], row["claim_lock"]))
run_id = kanban_db._end_run(
conn,
row["id"],
outcome="reclaimed",
status="todo",
summary=f"ancestor {parent_id} reopened",
)
conn.execute(
"UPDATE tasks SET status = 'todo', completed_at = NULL, "
"claim_lock = NULL, claim_expires = NULL, worker_pid = NULL, "
"current_run_id = NULL WHERE id = ?",
(row["id"],),
)
kanban_db._append_event(
conn,
row["id"],
"status",
{
"status": "todo",
"reason": "ancestor_reopened",
"parent": parent_id,
"previous_status": previous_status,
"resume_status": resume_status,
},
run_id=run_id,
)
def _set_status_direct(
conn: sqlite3.Connection, task_id: str, new_status: str,
) -> bool:
@@ -1083,19 +1146,33 @@ def _set_status_direct(
orphaned. ``running -> ready`` via drag-drop is the common case
(user yanking a stuck worker back to the queue).
"""
terminations: list[tuple[Optional[int], Optional[str]]] = []
effective_status = new_status
with kanban_db.write_txn(conn):
# Snapshot current state so we know whether to close a run.
prev = conn.execute(
"SELECT status, current_run_id FROM tasks WHERE id = ?",
"SELECT status, current_run_id, worker_pid, claim_lock "
"FROM tasks WHERE id = ?",
(task_id,),
).fetchone()
if prev is None:
return False
if prev["status"] == "running" and new_status == "ready":
resume_status = kanban_db._retry_status_for_run(
conn, task_id, prev["current_run_id"]
)
if resume_status == "review":
effective_status = (
"review"
if kanban_db._parents_satisfied(conn, task_id)
else "todo"
)
# Guard: don't allow promoting to 'ready' unless all parents are done.
# Prevents the dispatcher from spawning a child whose upstream work
# hasn't completed (e.g. T4 dispatched while T3 is still blocked).
if new_status == "ready":
if effective_status == "ready":
parent_statuses = conn.execute(
"SELECT t.status FROM tasks t "
"JOIN task_links l ON l.parent_id = t.id "
@@ -1103,14 +1180,14 @@ def _set_status_direct(
(task_id,),
).fetchall()
if parent_statuses and not all(
p["status"] == "done" for p in parent_statuses
p["status"] in {"done", "archived"} for p in parent_statuses
):
return False
was_running = prev["status"] == "running"
reopening_satisfied_parent = (
prev["status"] in {"done", "archived"}
and new_status not in {"done", "archived"}
and effective_status not in {"done", "archived"}
)
cur = conn.execute(
@@ -1119,61 +1196,49 @@ def _set_status_direct(
" claim_expires = CASE WHEN ? = 'running' THEN claim_expires ELSE NULL END, "
" worker_pid = CASE WHEN ? = 'running' THEN worker_pid ELSE NULL END "
"WHERE id = ?",
(new_status, new_status, new_status, new_status, task_id),
(
effective_status,
effective_status,
effective_status,
effective_status,
task_id,
),
)
if cur.rowcount != 1:
return False
run_id = None
if was_running and new_status != "running" and prev["current_run_id"]:
if was_running and effective_status != "running" and prev["current_run_id"]:
run_id = kanban_db._end_run(
conn, task_id,
outcome="reclaimed", status="reclaimed",
summary=f"status changed to {new_status} (dashboard/direct)",
summary=f"status changed to {effective_status} (dashboard/direct)",
)
terminations.append((prev["worker_pid"], prev["claim_lock"]))
conn.execute(
"INSERT INTO task_events (task_id, run_id, kind, payload, created_at) "
"VALUES (?, ?, 'status', ?, ?)",
(task_id, run_id, json.dumps({"status": new_status}), int(time.time())),
(
task_id,
run_id,
json.dumps(
{
"status": effective_status,
"requested_status": new_status,
}
),
int(time.time()),
),
)
if reopening_satisfied_parent:
# A parent leaving done/archived invalidates any direct child that
# was sitting in ready solely because that parent used to satisfy
# the dependency gate. Demote those children immediately so the
# dashboard does not keep advertising stale-ready work.
for row in conn.execute(
"SELECT c.id AS child_id, c.status AS child_status "
"FROM task_links l JOIN tasks c ON c.id = l.child_id "
"WHERE l.parent_id = ? ORDER BY c.id",
(task_id,),
).fetchall():
child_id = row["child_id"]
child_status = row["child_status"]
demoted = conn.execute(
"UPDATE tasks SET status = 'todo' "
"WHERE id = ? AND status IN ('ready', 'review')",
(child_id,),
)
if demoted.rowcount == 1 and child_status in {"ready", "review"}:
conn.execute(
"INSERT INTO task_events (task_id, kind, payload, created_at) "
"VALUES (?, 'status', ?, ?)",
(
child_id,
json.dumps(
{
"status": "todo",
"reason": "parent_reopened",
"parent": task_id,
"resume_status": (
"review" if child_status == "review" else "ready"
),
}
),
int(time.time()),
),
)
_invalidate_descendants_for_parent_reopen(
conn,
task_id,
terminations,
)
for pid, claim_lock in terminations:
kanban_db._terminate_reclaimed_worker(pid, claim_lock)
# If we re-opened something, children may have gone stale.
if new_status in {"done", "ready"}:
if effective_status in {"done", "ready", "review"}:
kanban_db.recompute_ready(conn)
return True
@@ -400,7 +400,8 @@ def test_review_dispatch_honors_global_and_per_profile_caps(
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
running = kb.claim_task(conn, running_id)
assert running is not None
review_ids: list[str] = []
for title in ("review one", "review two"):
@@ -424,6 +425,20 @@ def test_review_dispatch_honors_global_and_per_profile_caps(
task for task in globally_capped.spawned if task[0] in review_ids
]
assert kb.complete_task(
conn,
running_id,
expected_run_id=running.current_run_id,
)
global_dry_run = kb.dispatch_once(
conn,
dry_run=True,
max_in_progress=1,
)
assert len([
task for task in global_dry_run.spawned if task[0] in review_ids
]) == 1
per_profile_capped = kb.dispatch_once(
conn,
dry_run=True,
@@ -208,7 +208,11 @@ def test_parent_reopen_blocks_request_review_until_parent_is_done(conn) -> None:
)
def test_request_changes_fails_closed_on_malformed_review_provenance(conn):
@pytest.mark.parametrize("bad_payload", ["{not-json", "[]"])
def test_request_changes_fails_closed_on_malformed_review_provenance(
conn,
bad_payload: str,
):
task_id = kb.create_task(conn, title="Malformed handoff", assignee="builder")
implementation = kb.claim_task(conn, task_id, claimer="builder:1")
assert implementation is not None
@@ -222,7 +226,7 @@ def test_request_changes_fails_closed_on_malformed_review_provenance(conn):
conn.execute(
"UPDATE task_events SET payload = ? "
"WHERE task_id = ? AND kind = 'review_requested'",
("{malformed-json", task_id),
(bad_payload, task_id),
)
conn.commit()
review = kb.claim_review_task(conn, task_id, claimer="reviewer:1")
@@ -243,6 +247,21 @@ def test_request_changes_fails_closed_on_malformed_review_provenance(conn):
assert task.current_run_id == review.current_run_id
def test_reclaim_fails_safe_on_non_object_claim_provenance(conn) -> None:
task_id, _review = _claimed_review(conn, "Non-object claimed payload")
with kb.write_txn(conn):
conn.execute(
"UPDATE task_events SET payload = '[]' "
"WHERE task_id = ? AND kind = 'claimed' "
"AND run_id = (SELECT current_run_id FROM tasks WHERE id = ?)",
(task_id, task_id),
)
assert kb.reclaim_task(conn, task_id, signal_fn=lambda *_args: None)
task = kb.get_task(conn, task_id)
assert task is not None
assert task.status == "ready"
@pytest.mark.parametrize(
"reclaim_kind",
["spawn_failure", "expired_claim", "manual_reclaim", "stale_heartbeat"],
@@ -415,6 +434,24 @@ def test_crashed_and_timed_out_review_runs_retry_in_review_phase(
assert crashed.status == "review"
def test_parked_review_approval_without_evidence_still_creates_audit_run(conn) -> None:
task_id = kb.create_task(conn, title="Manual approval", assignee="reviewer")
assert kb.request_review(conn, task_id, summary="implementation handoff")
assert kb.complete_task(conn, task_id)
completed_event = _event(kb.list_events(conn, task_id), "completed")
assert completed_event.run_id is not None
run = kb.latest_run(conn, task_id)
assert run is not None
assert run.id == completed_event.run_id
assert run.outcome == "completed"
assert run.profile == "reviewer"
assert run.summary == "Review approved without additional evidence."
assert run.metadata == {
"source_status": "review",
"approval": "manual",
}
def test_legacy_review_child_deadlock_is_reported_immediately(conn):
implementation_id = kb.create_task(
conn,
@@ -159,6 +159,72 @@ def test_review_cli_round_trip_preserves_handoff(
assert task.assignee == "builder"
def test_domain_and_cli_review_handoffs_redact_before_persistence(
monkeypatch: pytest.MonkeyPatch,
tmp_path: Path,
) -> None:
home = tmp_path / ".hermes"
home.mkdir()
monkeypatch.setenv("HERMES_HOME", str(home))
secret = "ghp_" + "R" * 40
with kb.connect() as conn:
direct_id = kb.create_task(conn, title="direct redaction", assignee="builder")
direct_run = kb.claim_task(conn, direct_id)
assert direct_run is not None
assert kb.request_review(
conn,
direct_id,
summary=f"direct {secret}",
metadata={"nested": [secret]},
expected_run_id=direct_run.current_run_id,
)
run = kb.latest_run(conn, direct_id)
event = [
item for item in kb.list_events(conn, direct_id)
if item.kind == "review_requested"
][-1]
assert run is not None
assert secret not in str(run.summary)
assert secret not in json.dumps(run.metadata)
assert secret not in json.dumps(event.payload)
review = kb.claim_review_task(conn, direct_id)
assert review is not None
assert kb.request_changes(
conn,
direct_id,
reason=f"change {secret}",
expected_run_id=review.current_run_id,
) == (True, "builder")
run = kb.latest_run(conn, direct_id)
event = [
item for item in kb.list_events(conn, direct_id)
if item.kind == "changes_requested"
][-1]
assert run is not None
assert secret not in str(run.summary)
assert secret not in json.dumps(event.payload)
cli_id = kb.create_task(conn, title="CLI redaction", assignee="builder")
cli_output = kc.run_slash(
f'request-review {cli_id} --summary "cli {secret}" '
f"--metadata '{{\"token\":\"{secret}\"}}'"
)
assert "Requested review" in cli_output
assert secret not in cli_output
with kb.connect() as conn:
run = kb.latest_run(conn, cli_id)
event = [
item for item in kb.list_events(conn, cli_id)
if item.kind == "review_requested"
][-1]
assert run is not None
assert secret not in str(run.summary)
assert secret not in json.dumps(run.metadata)
assert secret not in json.dumps(event.payload)
def test_worker_guidance_distinguishes_same_card_and_downstream_review() -> None:
from agent.prompt_builder import KANBAN_GUIDANCE
from hermes_cli.config_defaults import DEFAULT_CONFIG
+96 -14
View File
@@ -8,6 +8,7 @@ REST surface without spinning up the whole dashboard.
from __future__ import annotations
import importlib.util
import json
import os
import subprocess
import sys
@@ -225,6 +226,7 @@ def test_task_detail_includes_links_and_events(client):
def test_patch_review_lifecycle_preserves_handoff_and_reopens(client):
secret = "ghp_" + "D" * 40
task = client.post(
"/api/plugins/kanban/tasks", json={"title": "review me", "assignee": "builder"},
).json()["task"]
@@ -233,8 +235,8 @@ def test_patch_review_lifecycle_preserves_handoff_and_reopens(client):
f"/api/plugins/kanban/tasks/{task['id']}",
json={
"status": "review",
"summary": "Implementation ready.",
"metadata": {"tests_run": 4},
"summary": f"Implementation ready. {secret}",
"metadata": {"tests_run": 4, "token": secret},
},
)
assert response.status_code == 200, response.text
@@ -243,7 +245,15 @@ def test_patch_review_lifecycle_preserves_handoff_and_reopens(client):
run = kb.latest_run(conn, task["id"])
assert run is not None
assert run.outcome == "review_requested"
assert run.metadata == {"tests_run": 4}
assert run.metadata is not None
assert run.metadata["tests_run"] == 4
assert secret not in str(run.summary)
assert secret not in json.dumps(run.metadata)
review_event = [
event for event in kb.list_events(conn, task["id"])
if event.kind == "review_requested"
][-1]
assert secret not in json.dumps(review_event.payload)
assert kb.assign_task(conn, task["id"], "reviewer")
response = client.patch(
@@ -331,20 +341,14 @@ def test_reopening_parent_retracts_review_and_blocks_approval(client):
assert response.status_code == 200, response.text
with kb.connect() as conn:
child = kb.get_task(conn, child_id)
assert child is not None
assert child.status == "running"
assert not kb.complete_task(
conn,
child_id,
summary="must not approve",
expected_run_id=active_review.current_run_id,
)
assert kb.reclaim_task(conn, child_id, signal_fn=lambda *_args: None)
assert kb.claim_review_task(conn, child_id) is None
child = kb.get_task(conn, child_id)
assert child is not None
assert child.status == "todo"
reclaimed = kb.latest_run(conn, child_id)
assert reclaimed is not None
assert reclaimed.outcome == "reclaimed"
assert kb.claim_review_task(conn, child_id) is None
assert not kb.complete_task(conn, child_id, summary="must not approve")
grandchild = kb.get_task(conn, grandchild_id)
assert grandchild is not None
assert grandchild.status == "todo"
@@ -372,6 +376,84 @@ def test_reopening_parent_retracts_review_and_blocks_approval(client):
assert grandchild.status == "ready"
def test_reopening_parent_recursively_retracts_done_and_running_descendants(client):
with kb.connect() as conn:
parent_id = kb.create_task(conn, title="root", assignee="planner")
assert kb.complete_task(conn, parent_id)
child_id = kb.create_task(
conn,
title="accepted child",
assignee="builder",
parents=[parent_id],
)
assert kb.complete_task(conn, child_id)
grandchild_id = kb.create_task(
conn,
title="running grandchild",
assignee="writer",
parents=[child_id],
)
grandchild_run = kb.claim_task(conn, grandchild_id)
assert grandchild_run is not None
response = client.patch(
f"/api/plugins/kanban/tasks/{parent_id}",
json={"status": "ready"},
)
assert response.status_code == 200, response.text
with kb.connect() as conn:
child = kb.get_task(conn, child_id)
grandchild = kb.get_task(conn, grandchild_id)
assert child is not None and child.status == "todo"
assert grandchild is not None and grandchild.status == "todo"
assert grandchild.current_run_id is None
assert kb.claim_task(conn, grandchild_id) is None
reclaimed = kb.latest_run(conn, grandchild_id)
assert reclaimed is not None
assert reclaimed.outcome == "reclaimed"
response = client.patch(
f"/api/plugins/kanban/tasks/{parent_id}",
json={"status": "done"},
)
assert response.status_code == 200, response.text
with kb.connect() as conn:
child = kb.get_task(conn, child_id)
grandchild = kb.get_task(conn, grandchild_id)
assert child is not None and child.status == "ready"
assert grandchild is not None and grandchild.status == "todo"
def test_dashboard_reclaim_of_active_review_preserves_review_phase(client):
with kb.connect() as conn:
task_id = kb.create_task(conn, title="active review", 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 = kb.claim_review_task(conn, task_id)
assert review is not None
response = client.patch(
f"/api/plugins/kanban/tasks/{task_id}",
json={"status": "ready"},
)
assert response.status_code == 200, response.text
assert response.json()["task"]["status"] == "review"
assert response.json()["task"]["assignee"] == "reviewer"
with kb.connect() as conn:
run = kb.latest_run(conn, task_id)
assert run is not None
assert run.outcome == "reclaimed"
next_review = kb.claim_review_task(conn, task_id)
assert next_review is not None
# ---------------------------------------------------------------------------
# DELETE /tasks/:id
# ---------------------------------------------------------------------------
+1 -1
View File
@@ -621,7 +621,7 @@ Multi-profile, multi-project collaboration board. Each install can host many boa
| `request-changes <id> <reason>` | Reviewer verdict for an active review run: close the review attempt and route the task back to its original implementer. |
| `reopen-review <id>...` | Send review task(s) back for changes (`review` → ready/todo). Flag: `--reason` (appended as a comment). |
| `schedule <id> "<reason>"` | Park time-delay/follow-up work in `scheduled` so it is not shown as a human blocker. |
| `unblock <id>` | Return a blocked or scheduled task to ready (or `todo` if dependencies are still open). |
| `unblock <id>` | Restore a blocked task to its source phase (`review` or `ready`), or `todo` while dependencies remain open. |
| `archive <id>` | Hide from default list. `gc` will remove scratch workspaces. |
| `tail <id>` | Follow a task's event stream. |
| `dispatch` | One dispatcher pass on the active board. Flags: `--dry-run`, `--max N`, `--failure-limit N`, `--json`. |