diff --git a/hermes_cli/kanban.py b/hermes_cli/kanban.py index 70b136b2f0dd..b478a1db1614 100644 --- a/hermes_cli/kanban.py +++ b/hermes_cli/kanban.py @@ -207,7 +207,7 @@ def _profile_author() -> str: _DELEGATED_CHILD_DENIED_ACTIONS: frozenset[str] = frozenset({ "init", "create", "swarm", "assign", "reclaim", "reassign", "link", "unlink", "claim", "comment", "attach", "attach-rm", "complete", "edit", "block", - "schedule", "unblock", "promote", "archive", "dispatch", "daemon", "repair", + "schedule", "unblock", "requeue", "promote", "archive", "dispatch", "daemon", "repair", "heartbeat", "notify-subscribe", "notify-unsubscribe", "specify", "decompose", "request-review", "request-changes", "reopen-review", "gc", @@ -1031,6 +1031,18 @@ def _cmd_unblock(args: argparse.Namespace) -> int: lambda tid: f"cannot unblock {tid} (not blocked/scheduled?)") +def _cmd_requeue(args: argparse.Namespace) -> int: + if os.environ.get("HERMES_KANBAN_TASK"): + return _err("kanban requeue is orchestrator-only") + reason = " ".join(args.reason).strip() + with kbc.connect_closing() as conn: + ok, error = kb.requeue_task(conn, args.task_id, actor=_profile_author(), reason=reason) + if not ok: + return _err(f"cannot requeue {args.task_id}: {error}") + print(f"Requeued {args.task_id}: {reason}") + return 0 + + def _cmd_request_review(args: argparse.Namespace) -> int: tid = args.task_id summary = _stripped_or_none(getattr(args, "summary", None)) @@ -1326,7 +1338,7 @@ def _cmd_decompose(args: argparse.Namespace) -> int: "comment": _cmd_comment, "attach": _cmd_attach, "attachments": _cmd_attachments, "attach-rm": _cmd_attach_rm, "complete": _cmd_complete, "edit": _cmd_edit, "block": _cmd_block, - "schedule": _cmd_schedule, "unblock": _cmd_unblock, + "schedule": _cmd_schedule, "unblock": _cmd_unblock, "requeue": _cmd_requeue, "request-review": _cmd_request_review, "request-changes": _cmd_request_changes, "reopen-review": _cmd_reopen_review, "promote": _cmd_promote, "archive": _cmd_archive, "tail": _cmd_tail, "dispatch": _cmd_dispatch, diff --git a/hermes_cli/kanban_db.py b/hermes_cli/kanban_db.py index b034a68054d9..9005cf9dff94 100644 --- a/hermes_cli/kanban_db.py +++ b/hermes_cli/kanban_db.py @@ -1576,6 +1576,22 @@ def assign_task(conn: sqlite3.Connection, task_id: str, profile: Optional[str]) return True +def requeue_task(conn: sqlite3.Connection, task_id: str, *, actor: str, reason: str) -> tuple[bool, Optional[str]]: + """Explicitly retry a READY card, recording operator intent without changing its status.""" + if not reason.strip(): + return False, "a reason is required" + with write_txn(conn): + row = conn.execute("SELECT status, claim_lock FROM tasks WHERE id = ?", (task_id,)).fetchone() + if row is None: + return False, "task not found" + if row["status"] != "ready" or row["claim_lock"] is not None: + return False, "requeue requires an unclaimed READY task" + _append_event(conn, task_id, "requeued", { + "actor": actor, "reason": reason.strip(), + }) + return True, None + + def set_model_override( conn: sqlite3.Connection, task_id: str, model: Optional[str], provider: Optional[str] = None, ) -> bool: @@ -1940,6 +1956,13 @@ def _append_event( run_id: Optional[int] = None, ) -> None: """Insert an event row inside the caller's txn; ``run_id`` groups it by attempt (NULL = task-scoped).""" + if kind in {"assigned", "changes_requested", "review_reopened", "requeued"} or ( + kind == "dependency_wait" and (payload or {}).get("kind") == "dependency" + ): + payload = dict(payload or {}) + payload["after_comment_id"] = conn.execute( + "SELECT COALESCE(MAX(id), 0) FROM task_comments WHERE task_id = ?", (task_id,), + ).fetchone()[0] conn.execute( "INSERT INTO task_events (task_id, run_id, kind, payload, created_at) " "VALUES (?, ?, ?, ?, ?)", (task_id, run_id, kind, _json_or_null(payload), int(time.time())), diff --git a/hermes_cli/kanban_db_dispatch.py b/hermes_cli/kanban_db_dispatch.py index 7b71be5569ef..1c8dc08bf8b7 100644 --- a/hermes_cli/kanban_db_dispatch.py +++ b/hermes_cli/kanban_db_dispatch.py @@ -1578,7 +1578,7 @@ def check_respawn_guard( requeued_after = conn.execute( "SELECT 1 FROM task_events " "WHERE task_id = ? AND created_at >= ? " - "AND kind IN ('status', 'promoted', 'unblocked', 'reclaimed') " + "AND kind IN ('status', 'promoted', 'unblocked', 'reclaimed', 'requeued') " "LIMIT 1", (task_id, completed_at), ).fetchone() @@ -1593,22 +1593,45 @@ def check_respawn_guard( # so the worker that opened the PR is still not re-spawned against it. pr_cutoff = now - _RESPAWN_GUARD_PR_WINDOW for c in conn.execute( - "SELECT body, created_at FROM task_comments " - "WHERE task_id = ? AND created_at >= ? ORDER BY created_at DESC", + "SELECT id, author, body, created_at FROM task_comments " + "WHERE task_id = ? AND created_at >= ? ORDER BY id DESC", (task_id, pr_cutoff), ).fetchall(): body = _kb._lossy_text(c["body"]) if not (body and _RESPAWN_GUARD_PR_URL_RE.search(body)): continue - events = conn.execute( - # Strictly after: a same-second tie stays guarded (fail closed). - "SELECT kind, payload FROM task_events " - "WHERE task_id = ? AND created_at > ? " - "AND kind IN ('assigned', 'changes_requested', 'review_reopened')", - (task_id, int(c["created_at"] or 0)), + # Intent events snapshot the latest comment id within their write txn. + # Historical events without that marker resume on an equal-second tie: + # a duplicate worker is recoverable; a stranded READY card is not. + intent_rows = conn.execute( + "SELECT id, kind, payload, created_at FROM task_events " + "WHERE task_id = ? AND kind IN " + "('assigned', 'changes_requested', 'review_reopened', 'requeued', 'dependency_wait') " + "ORDER BY id DESC", (task_id,), ).fetchall() - if any(_is_handoff_event(e["kind"], e["payload"]) for e in events): - return None + for event in intent_rows: + kind = event["kind"] + if kind == "dependency_wait": + if _kb._json_or(event["payload"], {}).get("kind") != "dependency": + continue + if not conn.execute( + "SELECT 1 FROM task_events WHERE task_id = ? " + "AND kind = 'promoted' AND id > ? LIMIT 1", + (task_id, event["id"]), + ).fetchone(): + continue + elif kind not in {"requeued"} and not _is_handoff_event(kind, event["payload"]): + continue + marker = _kb._json_or(event["payload"], {}).get("after_comment_id") + if not (marker >= c["id"] if marker is not None + else event["created_at"] >= c["created_at"]): + continue + if not conn.execute( + "SELECT 1 FROM task_events WHERE task_id = ? " + "AND kind = 'spawned' AND id > ? LIMIT 1", + (task_id, event["id"]), + ).fetchone(): + return None return "active_pr" return None diff --git a/hermes_cli/kanban_parser.py b/hermes_cli/kanban_parser.py index 8a65d4134671..c4b169122322 100644 --- a/hermes_cli/kanban_parser.py +++ b/hermes_cli/kanban_parser.py @@ -332,6 +332,10 @@ def _bulk_ids(verb: str): _reason("Optional reason/note — recorded as a comment before unblocking. Quote multi-word reasons."), _TASK_IDS, ], help="Return blocked/scheduled tasks to ready, or todo while parents remain open"), + _cmd("requeue", [ + _TASK_ID, + _arg("reason", nargs="+", help="Required operator reason for retrying a READY card"), + ], help="Explicitly retry a READY card held by the respawn guard"), _cmd("request-review", [ _TASK_ID, _arg("--summary", help="What was implemented and how it was verified — shown to the reviewer."), diff --git a/tests/hermes_cli/test_kanban_guard_resume.py b/tests/hermes_cli/test_kanban_guard_resume.py new file mode 100644 index 000000000000..89cdf17267b1 --- /dev/null +++ b/tests/hermes_cli/test_kanban_guard_resume.py @@ -0,0 +1,210 @@ +"""An open PR is not a duplicate when a worker explicitly waits on parents.""" + +from pathlib import Path + +import pytest + +from hermes_cli import kanban_db as kb +from hermes_cli import kanban_db_connect as kbc +from hermes_cli import kanban_db_dispatch as kbd + + +@pytest.fixture +def board(tmp_path, monkeypatch): + monkeypatch.setenv("HERMES_HOME", str(tmp_path)) + monkeypatch.setattr(Path, "home", lambda: tmp_path) + kb.init_db() + + +def test_worker_dependency_wait_promoted_continues_pr_once(board): + with kbc.connect() as conn: + child = kb.create_task(conn, title="implement", assignee="worker") + claim = kb.claim_task(conn, child) + assert claim is not None + kb.add_comment(conn, child, "worker", "https://github.com/o/r/pull/9") + parent = kb.create_task(conn, title="prerequisite", assignee="worker") + kb.link_tasks(conn, parent, child, expected_child_run_id=claim.current_run_id) + assert kb.block_task(conn, child, kind="dependency", reason="resume then complete") + assert kb.complete_task(conn, parent, summary="prerequisite complete") + kb.recompute_ready(conn) + assert kb.get_task(conn, child).status == "ready" + assert kbd.check_respawn_guard(conn, child) is None + assert kb.claim_task(conn, child) + with kb.write_txn(conn): + kb._append_event(conn, child, "spawned", {"pid": 42}) + conn.execute("UPDATE tasks SET status='ready', claim_lock=NULL, claim_expires=NULL, current_run_id=NULL WHERE id=?", (child,)) + assert kbd.check_respawn_guard(conn, child) == "active_pr" + + +def test_operator_requeue_ready_card(board): + from hermes_cli import kanban as kc + from hermes_cli.kanban_parser import build_parser + import argparse + + with kbc.connect() as conn: + task_id = kb.create_task(conn, title="retry", assignee="worker") + kb.add_comment(conn, task_id, "worker", "https://github.com/o/r/pull/9") + assert kbd.check_respawn_guard(conn, task_id) == "active_pr" + assert kb.requeue_task(conn, task_id, actor="operator", reason=" ") == (False, "a reason is required") + parser = argparse.ArgumentParser() + build_parser(parser.add_subparsers(dest="cmd")) + assert kc.kanban_command(parser.parse_args(["kanban", "requeue", task_id, "continue", "PR"])) == 0 + with kbc.connect() as conn: + assert kbd.check_respawn_guard(conn, task_id) is None + kb.add_comment(conn, task_id, "worker", "new https://github.com/o/r/pull/10") + assert kbd.check_respawn_guard(conn, task_id) == "active_pr" + assert kb.requeue_task(conn, task_id, actor="operator", reason="again") == (True, None) + assert kbd.check_respawn_guard(conn, task_id) is None + kb.add_comment(conn, task_id, "worker", "status update without a PR") + assert kbd.check_respawn_guard(conn, task_id) is None + with kb.write_txn(conn): + kb._append_event(conn, task_id, "spawned", {"pid": 42}) + assert kbd.check_respawn_guard(conn, task_id) == "active_pr" + assert kb.block_task(conn, task_id, kind="needs_input", reason="wait") + assert kb.requeue_task(conn, task_id, actor="operator", reason="again")[0] is False + + +def test_ordinary_comment_after_wait_does_not_cancel_pr_resume(board): + with kbc.connect() as conn: + child = kb.create_task(conn, title="implement", assignee="worker") + claim = kb.claim_task(conn, child) + assert claim is not None + kb.add_comment(conn, child, "worker", "https://github.com/o/r/pull/9") + parent = kb.create_task(conn, title="prerequisite", assignee="worker") + kb.link_tasks(conn, parent, child, expected_child_run_id=claim.current_run_id) + assert kb.block_task(conn, child, kind="dependency", reason="resume") + kb.add_comment(conn, child, "worker", "parent still pending") + assert kb.complete_task(conn, parent, summary="prerequisite complete") + kb.recompute_ready(conn) + assert kbd.check_respawn_guard(conn, child) is None + + +def test_inline_audit_comment_does_not_shift_ready_requeue(board, monkeypatch): + import time + from hermes_cli import kanban_db as kb + + now = int(time.time()) + monkeypatch.setattr(kbd.time, "time", lambda: now) + with kbc.connect() as conn: + task_id = kb.create_task(conn, title="inline audit", assignee="worker") + with kb.write_txn(conn): + kb._insert_comment(conn, task_id, "worker", "x" * len("https://github.com/o/r/pull/9"), now) + kb.add_comment(conn, task_id, "worker", "https://github.com/o/r/pull/9") + assert kbd.check_respawn_guard(conn, task_id) == "active_pr" + assert kb.requeue_task(conn, task_id, actor="operator", reason="continue PR") == (True, None) + assert kbd.check_respawn_guard(conn, task_id) is None + kb.add_comment(conn, task_id, "worker", "https://github.com/o/r/pull/10") + assert kbd.check_respawn_guard(conn, task_id) == "active_pr" + + +def test_legacy_equal_length_inline_comment_requeues(board, monkeypatch): + import time + + now = int(time.time()) + monkeypatch.setattr(kbd.time, "time", lambda: now) + with kbc.connect() as conn: + task_id = kb.create_task(conn, title="legacy inline audit", assignee="worker") + pr = "https://github.com/o/r/pull/9" + with kb.write_txn(conn): + kb._insert_comment(conn, task_id, "worker", "x" * len(pr), now) + kb.add_comment(conn, task_id, "worker", pr) + with kb.write_txn(conn): + conn.execute("UPDATE task_events SET payload=json_remove(payload, '$.comment_id') " + "WHERE task_id=? AND kind='commented'", (task_id,)) + assert kbd.check_respawn_guard(conn, task_id) == "active_pr" + assert kb.requeue_task(conn, task_id, actor="operator", reason="continue PR") == (True, None) + assert kbd.check_respawn_guard(conn, task_id) is None + kb.add_comment(conn, task_id, "worker", "https://github.com/o/r/pull/10") + assert kbd.check_respawn_guard(conn, task_id) == "active_pr" + + +def test_dependency_intent_cannot_attach_to_second_promotion(board): + with kbc.connect() as conn: + child = kb.create_task(conn, title="implement", assignee="worker") + claim = kb.claim_task(conn, child) + assert claim is not None + kb.add_comment(conn, child, "worker", "https://github.com/o/r/pull/9") + parent = kb.create_task(conn, title="prerequisite", assignee="worker") + kb.link_tasks(conn, parent, child, expected_child_run_id=claim.current_run_id) + assert kb.block_task(conn, child, kind="dependency", reason="resume") + assert kb.complete_task(conn, parent, summary="done") + kb.recompute_ready(conn) + assert kbd.check_respawn_guard(conn, child) is None + with kb.write_txn(conn): + kb._append_event(conn, child, "spawned", {"pid": 42}) + kb._append_event(conn, child, "promoted", {"status": "ready"}) + assert kbd.check_respawn_guard(conn, child) == "active_pr" + + +@pytest.mark.parametrize("inline_before", [False, True]) +@pytest.mark.parametrize("padded", [False, True]) +def test_historical_pr_wait_resumes_with_trimmed_or_inline_comment(board, monkeypatch, inline_before, padded): + import time + + now = int(time.time()) + clock = {"now": now} + monkeypatch.setattr(kbd.time, "time", lambda: clock["now"]) + with kbc.connect() as conn: + child = kb.create_task(conn, title="implement", assignee="worker") + claim = kb.claim_task(conn, child) + assert claim is not None + pr = "https://github.com/o/r/pull/9" + if inline_before: + with kb.write_txn(conn): + kb._insert_comment(conn, child, "worker", "x" * len(pr), now) + kb.add_comment(conn, child, "worker", f" {pr} " if padded else pr) + with kb.write_txn(conn): + conn.execute("UPDATE task_events SET payload=json_remove(payload, '$.comment_id') " + "WHERE task_id=? AND kind='commented'", (child,)) + parent = kb.create_task(conn, title="prerequisite", assignee="worker") + kb.link_tasks(conn, parent, child, expected_child_run_id=claim.current_run_id) + clock["now"] += 2 + assert kb.block_task(conn, child, kind="dependency", reason="resume") + assert kb.complete_task(conn, parent, summary="prerequisite complete") + kb.recompute_ready(conn) + assert kbd.check_respawn_guard(conn, child) is None + + +def test_legacy_intent_equal_second_resumes_without_comment_event_mapping(board, monkeypatch): + import time + + now = int(time.time()) + monkeypatch.setattr(kbd.time, "time", lambda: now) + with kbc.connect() as conn: + task_id = kb.create_task(conn, title="legacy tie", assignee="worker") + kb.add_comment(conn, task_id, "worker", " https://github.com/o/r/pull/9 ") + assert kb.requeue_task(conn, task_id, actor="operator", reason="continue") == (True, None) + with kb.write_txn(conn): + conn.execute( + "UPDATE task_events SET payload=json_remove(payload, '$.after_comment_id') " + "WHERE task_id=? AND kind='requeued'", (task_id,), + ) + assert kbd.check_respawn_guard(conn, task_id) is None + + +def test_pr_guard_does_not_correlate_comment_rows_with_events(): + import ast + import inspect + + source = inspect.getsource(kbd.check_respawn_guard) + sql_literals = [node.value.lower() for node in ast.walk(ast.parse(source)) + if isinstance(node, ast.Constant) and isinstance(node.value, str)] + assert "after_comment_id" in source + assert not any("kind = 'commented'" in value for value in sql_literals) + assert not any("join task_comments" in value or "join task_events" in value + for value in sql_literals) + + +def test_new_pr_comment_after_wait_does_not_resume(board): + with kbc.connect() as conn: + child = kb.create_task(conn, title="implement", assignee="worker") + claim = kb.claim_task(conn, child) + assert claim is not None + kb.add_comment(conn, child, "worker", "https://github.com/o/r/pull/9") + parent = kb.create_task(conn, title="prerequisite", assignee="worker") + kb.link_tasks(conn, parent, child, expected_child_run_id=claim.current_run_id) + assert kb.block_task(conn, child, kind="dependency", reason="resume") + kb.add_comment(conn, child, "worker", "https://github.com/o/r/pull/10") + assert kb.complete_task(conn, parent, summary="prerequisite complete") + kb.recompute_ready(conn) + assert kbd.check_respawn_guard(conn, child) == "active_pr"