From a699234f812598f3e27071e1de12388ff437210f Mon Sep 17 00:00:00 2001 From: DmytroVolodymyrson <58473490+DmytroVolodymyrson@users.noreply.github.com> Date: Wed, 26 Aug 2026 15:22:09 -0700 Subject: [PATCH] fix(kanban): wake controllers on review changes Deliver changes_requested review outcomes through kanban subscriptions and wake the origin for notify+wake / wake modes. Review-specific only: no task mutation. Reasons are redacted, path-scrubbed and truncated before delivery. Salvaged from #88694; conflicts with the #87733 wake-kinds expansion resolved keep-both. --- docs/kanban/multi-gateway.md | 7 + gateway/kanban_watchers.py | 64 +++++- locales/en.yaml | 2 + .../test_kanban_changes_requested_notifier.py | 182 ++++++++++++++++++ 4 files changed, 248 insertions(+), 7 deletions(-) create mode 100644 tests/gateway/test_kanban_changes_requested_notifier.py diff --git a/docs/kanban/multi-gateway.md b/docs/kanban/multi-gateway.md index fb6eec981f..18d55799f7 100644 --- a/docs/kanban/multi-gateway.md +++ b/docs/kanban/multi-gateway.md @@ -4,6 +4,13 @@ Hermes supports multiple gateway processes running concurrently — one per prof (default, writer, admin, coder, researcher). Each gateway opens its own connection to platform APIs and delivers messages for its profile's subscribers. +Task subscriptions also cover review feedback. A `changes_requested` review +event is delivered as an actionable review-BLOCK notification. Subscriptions +using `notify+wake` additionally wake the exact originating chat/thread/session +so the controller inspects the existing card and current run; `notify` remains +passive-only and `wake` remains wake-only. Review feedback never creates, +unblocks, requeues, or otherwise mutates a task. + ## Single-dispatcher posture Only one gateway owns the kanban dispatcher. The owning gateway keeps diff --git a/gateway/kanban_watchers.py b/gateway/kanban_watchers.py index df6342adbd..ad8c64ae8b 100644 --- a/gateway/kanban_watchers.py +++ b/gateway/kanban_watchers.py @@ -13,6 +13,7 @@ from __future__ import annotations import asyncio import logging import os +import re import sqlite3 import time from contextvars import Context @@ -26,6 +27,28 @@ from agent.i18n import t logger = logging.getLogger("gateway.run") +_LOCAL_PATH_RE = re.compile( + r"(? str: + """Return a mobile-friendly review reason safe for external delivery.""" + from agent.redact import redact_sensitive_text + + reason = redact_sensitive_text( + "" if value is None else str(value), + force=True, + redact_url_credentials=True, + ) + reason = _LOCAL_PATH_RE.sub("[local path]", reason) + reason = " ".join(reason.split()) + if len(reason) > limit: + reason = reason[: limit - 1].rstrip() + "…" + return reason + + def _resolve_auto_decompose_settings( load_config: Callable[[], Any], ) -> "tuple[bool, int]": @@ -204,7 +227,8 @@ class GatewayKanbanWatchersMixin: For each subscription row, fetches ``task_events`` newer than the stored cursor with kind in the terminal set (``completed``, ``blocked``, ``gave_up``, ``crashed``, ``timed_out``, - ``review_requested``, ``block_loop_detected``). Sends one + ``review_requested``, ``changes_requested``, + ``block_loop_detected``). Sends one message per new event to ``(platform, chat_id, thread_id)``, then advances the cursor. The subscription is removed only when the task is ``archived``. A ``done`` task can be reopened for review or @@ -239,7 +263,7 @@ class GatewayKanbanWatchersMixin: # but is not a block (see kanban_db.request_review); the task is not # archived, so the subscription stays alive and later review # cycles keep notifying. - TERMINAL_KINDS = ("completed", "blocked", "gave_up", "crashed", "timed_out", "status", "archived", "unblocked", "block_loop_detected", "review_requested") + TERMINAL_KINDS = ("completed", "blocked", "gave_up", "crashed", "timed_out", "status", "archived", "unblocked", "block_loop_detected", "review_requested", "changes_requested") # Subscriptions are removed only when the task reaches the irreversible # archived status. ``done`` is reversible in review/controller flows, # so removing its subscription would silence a later reopen. We used @@ -546,6 +570,7 @@ class GatewayKanbanWatchersMixin: # "Task X completed" and re-decomposes work that already # exists on the board. wake_handoff = "" + wake_review_detail = "" for ev in d["events"]: kind = ev.kind # Identity prefix: attribute terminal pings to the @@ -627,6 +652,22 @@ class GatewayKanbanWatchersMixin: f"👀 {board_tag}{tag}Kanban {sub['task_id']} ready for review" f" — {title}{handoff}" ) + elif kind == "changes_requested": + payload = ev.payload or {} + reason = _safe_review_reason(payload.get("reason")) + reviewer = _safe_review_reason(payload.get("reviewer"), 48) + implementer = _safe_review_reason(payload.get("implementer"), 48) + reason_text = reason or "reviewer feedback requires changes" + provenance = "" + if reviewer: + provenance += f" — reviewer @{reviewer}" + if implementer: + provenance += f" → implementer @{implementer}" + msg = ( + f"🛑 {board_tag}Kanban {sub['task_id']} review requested " + f"changes/BLOCK: {reason_text}{provenance}" + ) + wake_review_detail = reason_text elif kind == "block_loop_detected": # A task re-blocked for the same cause past the # recurrence limit and was routed to `triage` for a @@ -792,13 +833,16 @@ class GatewayKanbanWatchersMixin: task_terminal = task and task.status == "archived" # Kinds that hand a decision back to the origin, so the # origin has to take a turn. ``review_requested`` (the - # implementation is done and waits for a reviewer) and - # ``block_loop_detected`` (routed to triage) belong here - # for the same reason ``blocked`` does. ``status`` / - # ``archived`` / ``unblocked`` stay out: bookkeeping. + # implementation is done and waits for a reviewer), + # ``changes_requested`` (a reviewer BLOCKed and work + # returns to the implementer) and ``block_loop_detected`` + # (routed to triage) belong here for the same reason + # ``blocked`` does. ``status`` / ``archived`` / + # ``unblocked`` stay out: bookkeeping. _WAKE_KINDS = ( "completed", "gave_up", "crashed", "timed_out", - "blocked", "review_requested", "block_loop_detected", + "blocked", "review_requested", "changes_requested", + "block_loop_detected", ) _wake_kinds = ( {ev.kind for ev in d["events"] if ev.kind in _WAKE_KINDS} @@ -837,6 +881,7 @@ class GatewayKanbanWatchersMixin: if "timed_out" in _wake_kinds: _parts.append(t("gateway.kanban.wake.timed_out")) if "blocked" in _wake_kinds: _parts.append(t("gateway.kanban.wake.blocked")) if "review_requested" in _wake_kinds: _parts.append(t("gateway.kanban.wake.review_requested")) + if "changes_requested" in _wake_kinds: _parts.append(t("gateway.kanban.wake.changes_requested")) if "block_loop_detected" in _wake_kinds: _parts.append(t("gateway.kanban.wake.block_loop_detected")) _status = t("gateway.kanban.wake.status_joiner").join(_parts) or t("gateway.kanban.wake.status_default") _synth = t( @@ -857,6 +902,11 @@ class GatewayKanbanWatchersMixin: "gateway.kanban.wake.handoff", summary=wake_handoff, ) + if wake_review_detail: + _synth += "\n" + t( + "gateway.kanban.wake.review_detail", + reason=wake_review_detail, + ) _synth += "\n\n" + t( "gateway.kanban.wake.guidance" ) diff --git a/locales/en.yaml b/locales/en.yaml index 7a7ced00fe..a636a23c27 100644 --- a/locales/en.yaml +++ b/locales/en.yaml @@ -182,11 +182,13 @@ gateway: timed_out: "timed out; dispatcher will retry" blocked: "blocked; needs attention" review_requested: "handed off for review; the implementation is done" + changes_requested: "review requested changes (BLOCK); implementation is not approved" block_loop_detected: "routed to triage after repeated blocks; needs a decision" status_default: "status changed" status_joiner: ", " message: "[kanban] Task {task_id} {status}.\nTitle: {title}\nAssignee: @{assignee}\nBoard: {board}\n\nCheck the result or decide the next step." handoff: "Result: {summary}" + review_detail: "Review feedback: {reason}\nInspect the existing card and its current review run; return work to the same implementation task. Do not treat this BLOCK as approval and do not create a duplicate task." guidance: "This is an automatic task-status notification, not a request to decompose the task again. Inspect the current board before creating follow-up tasks; do not recreate tasks or graphs that already exist." personality: diff --git a/tests/gateway/test_kanban_changes_requested_notifier.py b/tests/gateway/test_kanban_changes_requested_notifier.py new file mode 100644 index 0000000000..a8f657aff6 --- /dev/null +++ b/tests/gateway/test_kanban_changes_requested_notifier.py @@ -0,0 +1,182 @@ +import asyncio + +from gateway.config import Platform +from gateway.run import GatewayRunner +from hermes_cli import kanban_db as kb + + +class RecordingAdapter: + def __init__(self, *, fail_send=False): + self.sent = [] + self.handled = [] + self.fail_send = fail_send + + async def send(self, chat_id, text, metadata=None): + self.sent.append({"chat_id": chat_id, "text": text, "metadata": metadata or {}}) + if self.fail_send: + raise RuntimeError("transient send failure") + + async def handle_message(self, event): + self.handled.append(event) + + +async def _run_one_tick(monkeypatch, runner): + real_sleep = asyncio.sleep + + async def fake_sleep(delay): + if delay == 5: + return None + runner._running = False + await real_sleep(0) + + monkeypatch.setattr(asyncio, "sleep", fake_sleep) + await runner._kanban_notifier_watcher(interval=1) + + +def _runner(adapter): + runner = GatewayRunner.__new__(GatewayRunner) + runner._running = True + runner.adapters = {Platform.TELEGRAM: adapter} + runner._kanban_sub_fail_counts = {} + runner._kanban_dispatcher_lock_handle = object() + return runner + + +def _create_review_block(delivery_mode, *, reason="Tests need updates"): + conn = kb.connect() + try: + task_id = kb.create_task( + conn, + title="existing implementation card", + assignee="implementer", + session_id="agent:main:telegram:thread:chat-1:topic-7", + ) + kb.add_notify_sub( + conn, + task_id=task_id, + platform="telegram", + chat_id="chat-1", + thread_id="topic-7", + chat_type="thread", + delivery_mode=delivery_mode, + delivery_metadata={"thread_id": "topic-7", "chat_type": "thread"}, + ) + kb._append_event( + conn, + task_id, + kind="changes_requested", + payload={ + "reason": reason, + "reviewer": "claude-qa", + "implementer": "codex-cua", + "status": "ready", + }, + ) + return task_id + finally: + conn.close() + + +def _unseen(task_id): + conn = kb.connect() + try: + _, events = kb.unseen_events_for_sub( + conn, + task_id=task_id, + platform="telegram", + chat_id="chat-1", + thread_id="topic-7", + kinds=["changes_requested"], + ) + return events + finally: + conn.close() + + +def test_changes_requested_notify_wake_is_actionable_and_exactly_routed(tmp_path, monkeypatch): + monkeypatch.setenv("HERMES_KANBAN_DB", str(tmp_path / "review-block.db")) + kb.init_db() + task_id = _create_review_block("notify+wake") + adapter = RecordingAdapter() + + asyncio.run(_run_one_tick(monkeypatch, _runner(adapter))) + + assert len(adapter.sent) == 1 + text = adapter.sent[0]["text"] + assert text.startswith(f"🛑 [default] Kanban {task_id} review requested changes/BLOCK: Tests need updates") + assert "reviewer @claude-qa → implementer @codex-cua" in text + assert adapter.sent[0]["metadata"]["thread_id"] == "topic-7" + assert len(adapter.handled) == 1 + wake = adapter.handled[0] + assert wake.source.chat_id == "chat-1" + assert wake.source.chat_type == "thread" + assert wake.source.thread_id == "topic-7" + assert "implementation is not approved" in wake.text + assert "Inspect the existing card and its current review run" in wake.text + assert "do not create a duplicate task" in wake.text + assert wake.text.count("do not create a duplicate task") == 1 + assert _unseen(task_id) == [] + + # A fresh watcher after restart cannot replay an event whose cursor advanced. + asyncio.run(_run_one_tick(monkeypatch, _runner(adapter))) + assert len(adapter.sent) == 1 + assert len(adapter.handled) == 1 + + +def test_changes_requested_notify_is_passive_only(tmp_path, monkeypatch): + monkeypatch.setenv("HERMES_KANBAN_DB", str(tmp_path / "notify.db")) + kb.init_db() + task_id = _create_review_block("notify") + adapter = RecordingAdapter() + + asyncio.run(_run_one_tick(monkeypatch, _runner(adapter))) + + assert len(adapter.sent) == 1 + assert adapter.handled == [] + assert _unseen(task_id) == [] + + +def test_changes_requested_wake_only_has_no_passive_post(tmp_path, monkeypatch): + monkeypatch.setenv("HERMES_KANBAN_DB", str(tmp_path / "wake.db")) + kb.init_db() + task_id = _create_review_block("wake") + adapter = RecordingAdapter() + + asyncio.run(_run_one_tick(monkeypatch, _runner(adapter))) + + assert adapter.sent == [] + assert len(adapter.handled) == 1 + assert _unseen(task_id) == [] + + +def test_changes_requested_send_failure_retries_without_event_loss(tmp_path, monkeypatch): + monkeypatch.setenv("HERMES_KANBAN_DB", str(tmp_path / "retry.db")) + kb.init_db() + task_id = _create_review_block("notify") + failing = RecordingAdapter(fail_send=True) + + asyncio.run(_run_one_tick(monkeypatch, _runner(failing))) + assert len(failing.sent) == 1 + assert len(_unseen(task_id)) == 1 + + healthy = RecordingAdapter() + asyncio.run(_run_one_tick(monkeypatch, _runner(healthy))) + assert len(healthy.sent) == 1 + assert _unseen(task_id) == [] + + +def test_changes_requested_reason_is_redacted_path_safe_and_truncated(tmp_path, monkeypatch): + monkeypatch.setenv("HERMES_KANBAN_DB", str(tmp_path / "redact.db")) + kb.init_db() + secret = "sk-abcdefghijklmnopqrstuvwxyz123456" + reason = f"See /Users/alice/private/review.log token={secret} " + ("x" * 300) + _create_review_block("notify", reason=reason) + adapter = RecordingAdapter() + + asyncio.run(_run_one_tick(monkeypatch, _runner(adapter))) + + text = adapter.sent[0]["text"] + assert "/Users/alice" not in text + assert "abcdefghijklmnopqrstuvwxyz" not in text + assert "[local path]" in text + assert "… — reviewer @claude-qa" in text