fix(kanban): preserve notifier subscription after done

This commit is contained in:
Dmytro Afanasiev
2026-08-11 22:54:30 -04:00
committed by Teknium
parent dcbfc79f8d
commit 294272c18c
2 changed files with 120 additions and 15 deletions
+16 -15
View File
@@ -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
+104
View File
@@ -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))