From 294272c18ca0564d70b0e1e69c89d370c7a450e1 Mon Sep 17 00:00:00 2001 From: Dmytro Afanasiev <58473490+DmytroVolodymyrson@users.noreply.github.com> Date: Tue, 11 Aug 2026 22:54:30 -0400 Subject: [PATCH] fix(kanban): preserve notifier subscription after done --- gateway/kanban_watchers.py | 31 ++++---- tests/gateway/test_kanban_notifier.py | 104 ++++++++++++++++++++++++++ 2 files changed, 120 insertions(+), 15 deletions(-) diff --git a/gateway/kanban_watchers.py b/gateway/kanban_watchers.py index ba273a368a..a60caaf4bd 100644 --- a/gateway/kanban_watchers.py +++ b/gateway/kanban_watchers.py @@ -185,9 +185,10 @@ class GatewayKanbanWatchersMixin: ``review_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 reaches a truly final *status* (``done`` / ``archived``), not on - any terminal event kind — so review cycles and re-block loops keep - notifying. + task is ``archived``. A ``done`` task can be reopened for review or + continuation, so its subscription and origin-session ownership must + survive completion. Cursor advancement prevents old events replaying + when that happens. Runs in the gateway event loop; all SQLite work is pushed to a thread via ``asyncio.to_thread`` so the loop never blocks on the @@ -214,11 +215,13 @@ class GatewayKanbanWatchersMixin: # writes — surface those transitions to subscribers too. # ``review_requested`` wakes the origin subscriber like a block does, # but is not a block (see kanban_db.request_review); the task is not - # done/archived, so the subscription stays alive and later review + # 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") - # Subscriptions are removed only when the task reaches a truly final - # status (done / archived). We used to also unsub on any terminal + # 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 + # to also unsub on any terminal # event kind (gave_up / crashed / timed_out / blocked), but that # silently dropped the user out of the loop whenever the dispatcher # respawned the task: a worker that crashes, gets reclaimed, runs @@ -226,7 +229,7 @@ class GatewayKanbanWatchersMixin: # crash because the subscription was deleted after the first event. # Same shape as the reblock-after-unblock cycle that PR #22941 # fixed for `blocked`. Keeping the subscription alive until the - # task is genuinely done lets the cursor (advanced atomically by + # task is archived lets the cursor (advanced atomically by # claim_unseen_events_for_sub) handle dedup, and any retry-loop # event reaches the user. # Per-subscription send-failure counter. Adapter.send raising @@ -706,7 +709,7 @@ class GatewayKanbanWatchersMixin: # advances after it succeeds — a failure rewinds the # claim exactly like a failed send() above, so the # next tick retries. - task_terminal = task and task.status in {"done", "archived"} + task_terminal = task and task.status == "archived" _WAKE_KINDS = ("completed", "gave_up", "crashed", "timed_out", "blocked") _wake_kinds = ( {ev.kind for ev in d["events"] if ev.kind in _WAKE_KINDS} @@ -925,13 +928,11 @@ class GatewayKanbanWatchersMixin: # Nothing left to deliver on this path (the wake, # if any, already succeeded above). sub_fail_counts.pop(sub_key, None) - # Unsubscribe only when the task has reached a truly - # final status (done / archived). For blocked / - # gave_up / crashed / timed_out the subscription is - # kept alive so the user gets notified again if the - # dispatcher respawns the task and it cycles into the - # same state. See the longer comment on TERMINAL_KINDS - # above for the failure mode this prevents. + # Unsubscribe only on archive. Completion (``done``) + # remains reversible: controllers reopen completed + # work for review corrections and continuation. The + # retained cursor prevents replay while preserving the + # original delivery and wake ownership for that cycle. if _is_push_adapter and send_passive and _wake_kinds: # notify+wake: the text ping above was the # delivery and the cursor has advanced; the wake diff --git a/tests/gateway/test_kanban_notifier.py b/tests/gateway/test_kanban_notifier.py index 33273bbfbd..a698c13cc7 100644 --- a/tests/gateway/test_kanban_notifier.py +++ b/tests/gateway/test_kanban_notifier.py @@ -354,6 +354,110 @@ def test_notifier_redelivers_same_kind_on_dispatch_cycle(tmp_path, monkeypatch): assert "crashed" in adapter.sent[1]["text"].lower() +def test_notifier_subscription_survives_done_reopen_until_archive( + tmp_path, monkeypatch, +): + """Done is reversible; archive alone ends notification ownership.""" + db_path = tmp_path / "done-reopen-archive.db" + monkeypatch.setenv("HERMES_KANBAN_DB", str(db_path)) + kb.init_db() + + conn = kb.connect() + try: + tid = kb.create_task( + conn, + title="review continuation", + assignee="worker", + session_id="origin-session", + ) + kb.add_notify_sub( + conn, + task_id=tid, + platform="telegram", + chat_id="origin-chat", + thread_id="origin-thread", + user_id="origin-user", + chat_type="group", + notifier_profile="reviewer", + delivery_mode="notify+wake", + ) + assert kb.complete_task(conn, tid, summary="first completion") + finally: + conn.close() + + adapter = RecordingAdapter() + runner = _make_runner(adapter) + runner._active_profile_name = lambda: "reviewer" + asyncio.run(_run_one_notifier_tick(monkeypatch, runner)) + + assert len(adapter.sent) == 1 + assert len(adapter.handled) == 1 + assert adapter.sent[0]["chat_id"] == "origin-chat" + assert adapter.sent[0]["metadata"]["thread_id"] == "origin-thread" + assert adapter.handled[0].source.thread_id == "origin-thread" + assert adapter.handled[0].source.profile == "reviewer" + + conn = kb.connect() + try: + subs = kb.list_notify_subs(conn, tid) + assert len(subs) == 1, "completion must retain the origin subscription" + first_cursor = subs[0]["last_event_id"] + finally: + conn.close() + + # A quiet tick proves the completed event cannot replay after its cursor + # was advanced, even though the subscription now remains present. + runner = _make_runner(adapter) + runner._active_profile_name = lambda: "reviewer" + asyncio.run(_run_one_notifier_tick(monkeypatch, runner)) + assert len(adapter.sent) == 1 + assert len(adapter.handled) == 1 + + conn = kb.connect() + try: + with kb.write_txn(conn): + conn.execute("UPDATE tasks SET status = 'ready' WHERE id = ?", (tid,)) + kb._append_event(conn, tid, "status", {"status": "ready"}) + assert kb.complete_task(conn, tid, summary="corrected completion") + finally: + conn.close() + + runner = _make_runner(adapter) + runner._active_profile_name = lambda: "reviewer" + asyncio.run(_run_one_notifier_tick(monkeypatch, runner)) + + # The reopen status and second completion each deliver once, while only + # completion wakes the exact original session/thread. + assert len(adapter.sent) == 3 + assert len(adapter.handled) == 2 + assert all(item["chat_id"] == "origin-chat" for item in adapter.sent) + assert adapter.handled[-1].source.thread_id == "origin-thread" + assert adapter.handled[-1].source.profile == "reviewer" + + conn = kb.connect() + try: + subs = kb.list_notify_subs(conn, tid) + assert len(subs) == 1 + assert subs[0]["last_event_id"] > first_cursor + assert kb.archive_task(conn, tid) + finally: + conn.close() + + runner = _make_runner(adapter) + runner._active_profile_name = lambda: "reviewer" + asyncio.run(_run_one_notifier_tick(monkeypatch, runner)) + + # Archive itself is intentionally silent, but consumes its event and + # removes the subscription so no later historical event can replay. + assert len(adapter.sent) == 3 + assert len(adapter.handled) == 2 + conn = kb.connect() + try: + assert kb.list_notify_subs(conn, tid) == [] + finally: + conn.close() + + def test_notifier_wakeup_uses_subscription_chat_type(tmp_path, monkeypatch): db_path = tmp_path / "chat-type-wakeup.db" monkeypatch.setenv("HERMES_KANBAN_DB", str(db_path))