Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
69 changes: 60 additions & 9 deletions gateway/kanban_watchers.py
Original file line number Diff line number Diff line change
Expand Up @@ -202,6 +202,24 @@ def _collect():
getattr(platform, "value", str(platform)).lower()
for platform in self.adapters.keys()
}
# Widen to every platform any secondary profile has live,
# not just the default profile's. This is only a coarse
# pre-filter to skip claiming events for subs nobody can
# possibly deliver — the precise per-profile check (via
# gateway/authz_mixin.py::_authorization_adapter, which
# forbids default-profile fallback) still runs at delivery
# time below, rewinding the claim if it resolves to None.
# Without this, a subscription owned by a secondary
# profile on a platform the DEFAULT profile never
# connected (e.g. beta owns discord, default doesn't) was
# dropped here before ever being claimed — no rewind
# applies to an unclaimed event, so it silently never
# retries.
for _profile_adapter_map in getattr(self, "_profile_adapters", {}).values():
active_platforms.update(
getattr(platform, "value", str(platform)).lower()
for platform in _profile_adapter_map.keys()
)
if not active_platforms:
logger.debug("kanban notifier: no connected adapters; skipping tick")
return deliveries
Expand Down Expand Up @@ -311,13 +329,16 @@ def _collect():
)
continue
sub_profile = sub.get("notifier_profile") or ""
adapter = None
if sub_profile:
_profile_map = getattr(self, "_profile_adapters", {}).get(sub_profile)
if _profile_map:
adapter = _profile_map.get(plat)
if adapter is None:
adapter = self.adapters.get(plat)
# Route via the SAME chokepoint the authorization path uses
# (gateway/authz_mixin.py::_authorization_adapter): a stamped
# profile with its own adapter-registry entry must be served
# by THAT profile's same-platform adapter and must NOT silently
# fall back to the default profile's adapter — otherwise a
# secondary profile's task notification is delivered by the
# wrong bot (the cross-profile mis-delivery this whole change
# exists to fix). The helper returns None only when the profile
# (or default) genuinely has no adapter for the platform.
adapter = self._authorization_adapter(plat, sub_profile or None)
if adapter is None:
logger.debug(
"kanban notifier: adapter %s disconnected before delivery for %s; rewinding claim",
Expand Down Expand Up @@ -394,6 +415,13 @@ def _collect():
new_status = str(ev.payload["status"])
msg = f"🔄 {board_tag}{tag}Kanban {sub['task_id']} → {new_status}"
else:
# archived / unblocked are claimed by TERMINAL_KINDS
# (so the cursor advances past them and they can't
# wedge a later completed/blocked event behind an
# unclaimed row) but are intentionally SILENT: an
# archive needs no user ping, and unblocked is an
# internal transition. They are also excluded from
# _WAKE_KINDS below, so they never wake the creator.
continue
metadata: dict[str, Any] = {}
if sub.get("thread_id"):
Expand Down Expand Up @@ -504,6 +532,23 @@ def _collect():
)
from gateway.session import SessionSource
from gateway.platforms.base import MessageEvent, MessageType
# KNOWN LIMITATION (tracked follow-up): the
# subscription row does not persist the
# creator's chat_type, and it is not carried
# on the session-context bridge, so we cannot
# faithfully reconstruct the creator's real
# session key here. build_session_key() keys
# DMs (":dm:<chat_id>") on a wholly different
# shape from group/thread, so any hardcoded
# value mis-routes some creators. "group" is
# the least-surprising default for the
# dashboard/group flows this wake primarily
# serves; DM-originated creators are handled
# by the follow-up that stamps + persists
# chat_type end-to-end. handle_message()
# get_or_create_session's the target, so a
# mismatch degrades to "wake lands in a fresh
# group session" — never an exception.
_source = SessionSource(
platform=plat,
chat_id=sub["chat_id"],
Expand All @@ -524,9 +569,15 @@ def _collect():
sub["task_id"], platform_str, sub["chat_id"], sub_profile or "default", _wake_kinds,
)
except Exception as _wk_err:
logger.debug(
# Best-effort: the notification itself already
# delivered and the cursor has advanced, so a
# broken wake path must not wedge the tick — but
# log at WARNING with a traceback rather than
# DEBUG so a persistently-failing wake is visible
# in normal logs instead of silently no-op'ing.
logger.warning(
"kanban notifier: wakeup injection failed for %s: %s",
sub["task_id"], _wk_err,
sub["task_id"], _wk_err, exc_info=True,
)
if task_terminal:
await asyncio.to_thread(
Expand Down
123 changes: 123 additions & 0 deletions tests/gateway/test_kanban_notifier.py
Original file line number Diff line number Diff line change
Expand Up @@ -233,3 +233,126 @@ def test_notifier_redelivers_same_kind_on_dispatch_cycle(tmp_path, monkeypatch):
f"deliveries (texts: {[d['text'] for d in adapter.sent]})"
)
assert "crashed" in adapter.sent[1]["text"].lower()


def test_notifier_owning_profile_adapter_no_default_fallback(tmp_path, monkeypatch):
"""A subscription owned by a secondary profile whose profile-adapter
registry entry EXISTS but lacks this platform must NOT fall back to the
default profile's same-platform adapter — the notifier must route through
the shared ``_authorization_adapter`` chokepoint, which forbids that
fallback (gateway/authz_mixin.py). Delivering via the default profile's bot
is the exact cross-profile mis-delivery this whole change exists to fix
(`[230002] Bot can NOT be out of the chat`).

Mutation check: reverting kanban_watchers.py's adapter selection to the old
inline ``if adapter is None: adapter = self.adapters.get(plat)`` fallback
makes this test FAIL (the default adapter receives the delivery).
"""
db_path = tmp_path / "profile-no-fallback.db"
monkeypatch.setenv("HERMES_KANBAN_DB", str(db_path))
kb.init_db()

conn = kb.connect()
try:
tid = kb.create_task(conn, title="owned by beta", assignee="worker")
# Subscription is owned by profile "beta".
kb.add_notify_sub(
conn, task_id=tid, platform="telegram", chat_id="chat-beta",
notifier_profile="beta",
)
kb.complete_task(conn, tid, summary="done")
finally:
conn.close()

default_adapter = RecordingAdapter()
other_adapter = RecordingAdapter()
runner = GatewayRunner.__new__(GatewayRunner)
runner._running = True
# Default profile has a telegram adapter …
runner.adapters = {Platform.TELEGRAM: default_adapter}
# … and profile "beta" HAS a non-empty registry entry (so it passes the
# notifier's upstream skip-filter, which only skips owning profiles with NO
# adapter at all), but that entry does NOT contain a telegram adapter — beta
# connected a different platform (discord). The telegram sub owned by beta
# must therefore resolve to NO adapter, not silently borrow the default
# profile's telegram bot.
runner._profile_adapters = {"beta": {Platform.DISCORD: other_adapter}}
runner._kanban_sub_fail_counts = {}

asyncio.run(_run_one_notifier_tick(monkeypatch, runner))

# The default profile's adapter must never receive beta's notification.
assert default_adapter.sent == [], (
"Owning-profile subscription must not fall back to the default "
f"profile's adapter; got {default_adapter.sent!r}"
)
assert other_adapter.sent == [], (
f"beta's discord adapter must not receive a telegram sub; got {other_adapter.sent!r}"
)
# The claim is rewound (adapter resolved to None → treated as disconnected),
# so the event is still unseen and will deliver once beta's adapter connects.
assert [ev.kind for ev in _unseen_terminal_events_for(tid, "chat-beta")] == ["completed"]


def test_notifier_claims_platform_only_a_secondary_profile_owns(tmp_path, monkeypatch):
"""A subscription owned by a secondary profile on a platform the DEFAULT
profile never connected must still be claimed and delivered.

Regression: the ``_collect()`` pre-filter built ``active_platforms``
solely from ``self.adapters`` (the default profile). A sub owned by
profile "beta" on "discord", where beta genuinely has a live discord
adapter but the default profile has no discord adapter at all, was
dropped by that pre-filter (``platform not in active_platforms``)
before ``claim_unseen_events_for_sub`` ever ran — unlike the
disconnected-adapter path, an unclaimed event is never rewound, so this
was a permanent, silent notification loss, not a retryable one. This
directly contradicts the feature's own purpose (routing notifications
via the owning profile), and is the same cross-profile-adapter-lookup
class the delivery-side chokepoint in
``test_notifier_owning_profile_adapter_no_default_fallback`` already
guards — just one gate earlier.
"""
db_path = tmp_path / "secondary-only-platform.db"
monkeypatch.setenv("HERMES_KANBAN_DB", str(db_path))
kb.init_db()

conn = kb.connect()
try:
tid = kb.create_task(conn, title="owned by beta on discord", assignee="worker")
kb.add_notify_sub(
conn, task_id=tid, platform="discord", chat_id="chat-beta",
notifier_profile="beta",
)
kb.complete_task(conn, tid, summary="done")
finally:
conn.close()

beta_adapter = RecordingAdapter()
runner = GatewayRunner.__new__(GatewayRunner)
runner._running = True
# Default profile has NO discord adapter at all.
runner.adapters = {Platform.TELEGRAM: RecordingAdapter()}
# Secondary profile "beta" has a live discord adapter.
runner._profile_adapters = {"beta": {Platform.DISCORD: beta_adapter}}
runner._kanban_sub_fail_counts = {}

asyncio.run(_run_one_notifier_tick(monkeypatch, runner))

assert len(beta_adapter.sent) == 1, (
f"beta's discord adapter should have received the notification; got {beta_adapter.sent!r}"
)


def _unseen_terminal_events_for(tid, chat_id):
conn = kb.connect()
try:
_, events = kb.unseen_events_for_sub(
conn,
task_id=tid,
platform="telegram",
chat_id=chat_id,
kinds=["completed", "blocked", "gave_up", "crashed", "timed_out"],
)
return events
finally:
conn.close()
Loading