diff --git a/gateway/kanban_watchers_notifier.py b/gateway/kanban_watchers_notifier.py index 39330e401d..dba036b86a 100644 --- a/gateway/kanban_watchers_notifier.py +++ b/gateway/kanban_watchers_notifier.py @@ -9,6 +9,7 @@ from __future__ import annotations import re from pathlib import Path +import weakref from typing import Any, Callable, Optional from agent.i18n import t @@ -74,7 +75,7 @@ def _wake_scope_id(adapter: Any, sub: dict) -> Optional[str]: """ delivery_meta = sub.get("delivery_metadata") if isinstance(delivery_meta, dict): - for key in ("scope_id", "slack_team_id", "team_id"): + for key in ("scope_id", "guild_id", "slack_team_id", "team_id"): value = delivery_meta.get(key) if value: return str(value) @@ -95,6 +96,47 @@ def _platform_names(mapping: Any) -> set[str]: return {getattr(platform, "value", str(platform)).lower() for platform in mapping} +def _adapter_for_subscription(runner: Any, platform: Any, sub: dict, owner_profile: Optional[str]) -> Any: + """Resolve a durable route without turning a missing secondary bot into primary authority.""" + adapter = runner._authorization_adapter(platform, owner_profile) + config = getattr(runner, "config", None) + if not getattr(config, "multiplex_profiles", False): + return adapter + primary = runner.adapters.get(platform) + if adapter is not None and adapter is not primary: + return adapter + profile = owner_profile or getattr(runner, "_kanban_notifier_profile", None) + primary_profile = getattr(runner, "_primary_profile_name", None) or runner._active_profile_name() + profile = profile or primary_profile + # Empty maps are startup placeholders for route-only profiles; a connected + # secondary on ANY platform establishes an independent credential boundary. + if (getattr(runner, "_profile_adapters", {}) or {}).get(profile): + return None + metadata = sub.get("delivery_metadata") or {} + guild = metadata.get("scope_id") or metadata.get("guild_id") + parent = metadata.get("parent_chat_id") + chat, thread = sub.get("chat_id"), sub.get("thread_id") or None + thread_like = bool(thread) or (sub.get("chat_type") or metadata.get("chat_type")) in { + "thread", "forum", "forum_post", "forum-post", "topic", + } + # Preserve canonical route order, including equal-specificity ties. An older + # row missing an anchor must not skip a potentially winning route. Reuse the + # route's matcher (including platform-specific identity aliases), not a second + # hand-maintained equality implementation. + for route in getattr(config, "profile_routes", None) or []: + if route.matches(platform.value, guild_id=guild, chat_id=chat, + thread_id=thread, parent_chat_id=parent): + if route.profile != profile: + return None + from gateway.run import _multiplex_profile_homes + served = {name for name, _home in _multiplex_profile_homes(config)} + return primary if profile in served else None + if route.matches(platform.value, guild_id=guild or route.guild_id, chat_id=chat, + thread_id=thread, parent_chat_id=parent or (route.chat_id if thread_like else None)): + return None + return primary if profile == primary_profile else None + + # --- Collection (runs in a worker thread) --- @@ -102,6 +144,7 @@ class _Collector: """One tick's claim state: which profiles/platforms this gateway serves and the GC gate.""" def __init__(self, runner: Any, kb: Any, *, notifier_profile: Optional[str], gc_due: bool, gc_retention_days: int) -> None: + self.runner = runner self.kb = kb self.notifier_profile = notifier_profile self.gc_due = gc_due @@ -111,11 +154,15 @@ class _Collector: self.profile_adapters = getattr(runner, "_profile_adapters", {}) self.notifier_profiles = {notifier_profile} self.notifier_profiles.update(str(p).strip() for p in self.profile_adapters if str(p).strip()) + config = getattr(runner, "config", None) + if getattr(config, "multiplex_profiles", False): + self.notifier_profiles.update( + route.profile for route in config.profile_routes + if route.enabled and route.platform in _platform_names(runner.adapters) + ) # Include every platform any secondary profile has live. This is only a - # coarse pre-filter; the precise per-profile check (_authorization_adapter, - # no default fallback) runs at delivery and rewinds the claim if it - # resolves to None. An unclaimed event never retries, so dropping a - # secondary-profile sub here would lose it. + # coarse pre-filter; exact destination authorization runs before claim + # and again at delivery, rewinding if the route or adapter changed. self.active_platforms = _platform_names(runner.adapters).union( *(_platform_names(m) for m in self.profile_adapters.values())) @@ -169,15 +216,14 @@ class _Collector: def _claim_for_sub(self, conn: Any, slug: str, sub: dict) -> Optional[dict]: """Claim one subscription's unseen events; None when skipped or nothing new.""" owner_profile = sub.get("notifier_profile") or None - if owner_profile and owner_profile != self.notifier_profile and not self.profile_adapters.get(owner_profile): - logger.debug("kanban notifier: subscription for %s owned by profile %s; current profile %s has no adapter for it, skipping", - sub.get("task_id"), owner_profile, self.notifier_profile) - return None platform = (sub.get("platform") or "").lower() if platform not in self.active_platforms: logger.debug("kanban notifier: subscription for %s on %s skipped; adapter not connected", sub.get("task_id"), platform or "") return None + from gateway.config import Platform + if _adapter_for_subscription(self.runner, Platform(platform), sub, owner_profile or self.notifier_profile) is None: + return None old_cursor, cursor, events = _kbn().claim_unseen_events_for_sub( conn, task_id=sub["task_id"], platform=sub["platform"], chat_id=sub["chat_id"], thread_id=sub.get("thread_id") or "", kinds=TERMINAL_KINDS, @@ -460,17 +506,22 @@ class _KanbanNotification: # (#60600 rows) — fall back to that, then to "group" (the historical default that suits the # dashboard/group flows). handle_message() get_or_create_session's the target, so a mismatch only # ever degrades to a fresh session, never an exception. - _chat_type = str(sub.get("chat_type") or "").strip() - if not _chat_type: - _delivery_meta = sub.get("delivery_metadata") - if isinstance(_delivery_meta, dict): - _chat_type = str(_delivery_meta.get("chat_type") or "").strip() + _delivery_meta = sub.get("delivery_metadata") or {} + _chat_type = str(sub.get("chat_type") or _delivery_meta.get("chat_type") or "").strip() _source = SessionSource( platform=self.plat, chat_id=sub["chat_id"], chat_type=_chat_type or "group", thread_id=sub.get("thread_id") or None, user_id=sub.get("user_id"), user_id_alt=sub.get("user_id_alt"), profile=self.sub_profile or None, scope_id=_wake_scope_id(self.adapter, sub), + parent_chat_id=_delivery_meta.get("parent_chat_id"), ) - await deliver_wake(self.adapter, text=self.synth, session_id=self.session_key, source=_source) + _source._transport_adapter_ref = weakref.ref(self.adapter) + from gateway.run import _async_profile_runtime_scope + if self.sub_profile and getattr(getattr(self.runner, "config", None), "multiplex_profiles", False): + from hermes_cli.profiles import profile_exists + if not profile_exists(self.sub_profile): + raise RuntimeError(f"Kanban wake profile {self.sub_profile!r} no longer exists") + async with _async_profile_runtime_scope(self.runner._resolve_profile_home_for_source(_source)): + await deliver_wake(self.adapter, text=self.synth, session_id=self.session_key, source=_source) self._log_woke() async def _send_event(self, ev: Any, msg: str) -> None: @@ -535,11 +586,8 @@ class _KanbanNotification: except ValueError: await self.advance() return - # Same chokepoint as authorization: a stamped profile is served by ITS - # same-platform adapter and never falls back to the default profile's - # bot (cross-profile mis-delivery). None only when the profile (or - # default) has no adapter. - adapter = self.runner._authorization_adapter(self.plat, self.sub_profile or None) + # Recheck the exact route after claiming: config/adapters can change between ticks. + adapter = _adapter_for_subscription(self.runner, self.plat, self.sub, self.sub_profile or None) if adapter is None: logger.debug("kanban notifier: adapter %s disconnected before delivery for %s; rewinding claim", self.platform_str, self.task_id) diff --git a/tests/gateway/test_kanban_routed_transport.py b/tests/gateway/test_kanban_routed_transport.py new file mode 100644 index 0000000000..e00a9d8b5b --- /dev/null +++ b/tests/gateway/test_kanban_routed_transport.py @@ -0,0 +1,191 @@ +"""Persisted notification routes authorize exactly one transport, including route-only profiles.""" +import asyncio +from pathlib import Path + +from gateway.config import GatewayConfig, Platform +from gateway.kanban_watchers_notifier import _KanbanNotification, _notifier_collect +from gateway.profile_routing import parse_profile_routes +from gateway.run import GatewayRunner +from hermes_cli import kanban_db as kb, kanban_db_connect as kbc, kanban_db_notify as kbn + + +class RecordingAdapter: + supports_async_delivery = True + + def __init__(self): + self.sent = [] + self.handled = [] + + async def send(self, chat_id, text, **kwargs): + self.sent.append((chat_id, text, kwargs)) + + async def handle_message(self, event): + self.handled.append(event) + + +def setup_runner(tmp_path, monkeypatch): + monkeypatch.setattr(Path, "home", lambda: tmp_path) + home = tmp_path / ".hermes" + monkeypatch.setenv("HERMES_HOME", str(home)) + monkeypatch.setenv("HERMES_KANBAN_DB", str(home / "kanban.db")) + for name in ("yuki", "other"): + profile = home / "profiles" / name + profile.mkdir(parents=True) + (profile / "config.yaml").write_text("{}\n", encoding="utf-8") + runner = GatewayRunner.__new__(GatewayRunner) + runner.adapters = {Platform.DISCORD: RecordingAdapter()} + runner._profile_adapters = {"yuki": {}} + runner._primary_profile_name = "default" + runner._kanban_notifier_profile = "default" + runner._kanban_dispatcher_lock_handle = object() + runner.config = GatewayConfig(multiplex_profiles=True, profile_routes=parse_profile_routes([ + dict(platform="discord", guild_id="guild", chat_id="parent", profile="yuki"), + ])) + return runner + + +def completion(*, profile="yuki", metadata=None, chat="post", thread="post", mode="notify+wake"): + with kbc.connect() as conn: + task = kb.create_task(conn, title="route completion", assignee="worker") + kbn.add_notify_sub(conn, task_id=task, platform="discord", chat_id=chat, + thread_id=thread, chat_type="thread", user_id="creator", + notifier_profile=profile, delivery_mode=mode, + delivery_metadata=metadata if metadata is not None else + {"guild_id": "guild", "scope_id": "guild", "parent_chat_id": "parent"}) + kb.complete_task(conn, task, result="finished") + return task + + +def collect(runner): + return _notifier_collect(runner, kb, notifier_profile="default", gc_due=False, gc_retention_days=30) + + +async def deliver(runner, rows): + for row in rows: + await _KanbanNotification(runner, row, platform_cls=Platform, sub_fail_counts={}).deliver() + + +def unseen(task): + with kbc.connect() as conn: + return kbn.unseen_events_for_sub(conn, task_id=task, platform="discord", chat_id="post", + thread_id="post", kinds=["completed"])[1] + + +def test_exact_routed_profile_delivers_once_on_its_authorized_transport(tmp_path, monkeypatch): + runner = setup_runner(tmp_path, monkeypatch) + primary = runner.adapters[Platform.DISCORD] + task = completion(metadata={"scope_id": "guild", "guild_id": "stale-alias", "parent_chat_id": "parent"}) + rows = collect(runner) + assert [row["task"].id for row in rows] == [task] + asyncio.run(deliver(runner, rows)) + assert len(primary.sent) == len(primary.handled) == 1 + source = primary.handled[0].source + assert (source.profile, source.guild_id, source.scope_id, source.parent_chat_id) == ( + "yuki", "guild", "guild", "parent") + assert runner._adapter_for_source(source) is primary + assert not collect(runner) + + # A connected secondary owns its credential even where the primary route matches. + secondary = RecordingAdapter() + secondary.scope_id_for_chat = lambda chat: "stale-cache" + runner._profile_adapters["yuki"] = {Platform.DISCORD: secondary} + task = completion(metadata={"guild_id": "guild", "parent_chat_id": "parent"}) + asyncio.run(deliver(runner, collect(runner))) + assert len(primary.sent) == 1 + assert len(secondary.sent) == len(secondary.handled) == 1 + assert secondary.handled[0].source.scope_id == "guild" + assert runner._adapter_for_source(secondary.handled[0].source) is secondary + assert not unseen(task) + + +def test_route_denials_leave_events_retryable_at_claim_and_send(tmp_path, monkeypatch): + runner = setup_runner(tmp_path, monkeypatch) + primary = runner.adapters[Platform.DISCORD] + # Unknown owners, wrong/default owners, incomplete anchors, and partial credentials + # never become primary delivery authority. + tasks = [completion(profile=owner) for owner in ("other", "default", None)] + tasks += [completion(metadata=meta) for meta in ( + {"parent_chat_id": "parent"}, {"guild_id": "guild"}, + {"guild_id": "wrong", "parent_chat_id": "parent"}, + )] + assert not collect(runner) + assert all(unseen(task) for task in tasks) + + good = completion() + runner._profile_adapters["yuki"] = {Platform.TELEGRAM: RecordingAdapter()} + assert not collect(runner) + runner._profile_adapters["yuki"] = {} + runner.config.multiplex_profile_allowlist = ["other"] + assert not collect(runner) + runner.config.multiplex_profile_allowlist = ["yuki"] + rows = collect(runner) + assert [row["task"].id for row in rows] == [good] + # Reassignment after the claim must rewind, never send using stale authority. + runner.config.profile_routes = parse_profile_routes([ + dict(platform="discord", guild_id="guild", chat_id="parent", profile="other")]) + asyncio.run(deliver(runner, rows)) + assert primary.sent == primary.handled == [] + assert unseen(good) + + # Equal-specificity rules retain configuration order: an unknown parent + # cannot skip an earlier rule, but a known conflicting parent rules it out. + monkeypatch.setenv("HERMES_KANBAN_DB", str(tmp_path / "tied-routes.db")) + runner.config.profile_routes = parse_profile_routes([ + dict(platform="discord", guild_id="guild", chat_id="parent", profile="other"), + dict(platform="discord", guild_id="guild", chat_id="post", profile="yuki"), + ]) + ambiguous = completion(metadata={"scope_id": "guild"}) + exact = completion(metadata={"scope_id": "guild", "parent_chat_id": "different-parent"}) + rows = collect(runner) + assert [row["task"].id for row in rows] == [exact] + asyncio.run(deliver(runner, rows)) + assert len(primary.sent) == len(primary.handled) == 1 + assert unseen(ambiguous) + + +def test_kanban_wakes_install_the_destination_runtime_scope(tmp_path, monkeypatch): + from agent.secret_scope import get_secret + from gateway.run import _profile_runtime_scope + from hermes_constants import get_hermes_home + + runner = setup_runner(tmp_path, monkeypatch) + home = tmp_path / ".hermes" + (home / ".env").write_text("KANBAN_TEST_SECRET=primary\n", encoding="utf-8") + observed = [] + + class ScopedAdapter(RecordingAdapter): + async def handle_message(self, event): + # A real yield catches scopes that mutate process-global state. + await asyncio.sleep(0) + observed.append((event.source.profile, get_secret("KANBAN_TEST_SECRET"), get_hermes_home())) + await super().handle_message(event) + + for name in ("yuki", "other"): + (home / "profiles" / name / ".env").write_text(f"KANBAN_TEST_SECRET={name}\n", encoding="utf-8") + runner._profile_adapters[name] = {Platform.DISCORD: ScopedAdapter()} + completion(profile=name) + rows = collect(runner) + assert len(rows) == 2 + + async def concurrent_wakes(): + with _profile_runtime_scope(home): + await asyncio.gather(*(_KanbanNotification(runner, row, platform_cls=Platform, + sub_fail_counts={}).deliver() for row in rows)) + assert get_secret("KANBAN_TEST_SECRET") == "primary" + asyncio.run(concurrent_wakes()) + assert sorted(observed) == [(name, name, home / "profiles" / name) for name in ("other", "yuki")] + + +def test_removed_profile_never_wakes_under_the_primary_runtime(tmp_path, monkeypatch): + import shutil + + runner = setup_runner(tmp_path, monkeypatch) + secondary = RecordingAdapter() + runner._profile_adapters["yuki"] = {Platform.DISCORD: secondary} + task = completion(mode="wake") + rows = collect(runner) + assert len(rows) == 1 + shutil.rmtree(tmp_path / ".hermes" / "profiles" / "yuki") + asyncio.run(deliver(runner, rows)) + assert secondary.handled == [] + assert unseen(task) diff --git a/website/docs/user-guide/features/kanban.md b/website/docs/user-guide/features/kanban.md index 8374942f28..4eda0372d8 100644 --- a/website/docs/user-guide/features/kanban.md +++ b/website/docs/user-guide/features/kanban.md @@ -1118,6 +1118,14 @@ dispatch and delivery have separate owners: `writer` profile's Telegram gets its `completed`/`blocked` message delivered by the `writer` gateway, even though the `default` gateway did the dispatching. +- **Route-only multiplex profiles** can use the primary adapter when the + subscription's persisted platform, chat, thread, scope and parent-channel + anchors resolve to that exact served profile through `gateway.profile_routes`. + A connected secondary adapter remains authoritative; a partial secondary + adapter registry never falls back to the primary bot. Unmatched, reassigned, + disabled or ambiguous routes remain undelivered and retryable. Old rows + missing required routing anchors are not guessed into a profile. Wake turns keep + the destination profile's runtime scope and the authorized transport. - **Legacy subscriptions** created before profile stamping (no `notifier_profile` on the row) are delivered only by the gateway that holds the actual dispatcher singleton lock, so two gateways never race for them.