From 320a7caa6afb780cab0a520a54f44bfc442858af Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 18:52:31 -0700 Subject: [PATCH] refactor(gateway): kanban notifier collector class, lambda formatters, single wake path --- gateway/kanban_watchers_notifier.py | 378 +++++++++++++--------------- 1 file changed, 179 insertions(+), 199 deletions(-) diff --git a/gateway/kanban_watchers_notifier.py b/gateway/kanban_watchers_notifier.py index 1c0b0c283f..ac50b1ea34 100644 --- a/gateway/kanban_watchers_notifier.py +++ b/gateway/kanban_watchers_notifier.py @@ -86,6 +86,145 @@ def _platform_names(mapping: Any) -> set[str]: # --- Collection (runs in a worker thread) --- +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.kb = kb + self.notifier_profile = notifier_profile + self.gc_due = gc_due + self.gc_retention_days = gc_retention_days + self.deliveries: list[dict] = [] + self.include_unowned = runner._owns_kanban_dispatcher_lock() + 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()) + self.active_platforms = _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. + for _profile_adapter_map in self.profile_adapters.values(): + self.active_platforms.update(_platform_names(_profile_adapter_map)) + + def collect(self) -> list[dict]: + if not self.active_platforms: + logger.debug("kanban notifier: no connected adapters; skipping tick") + return self.deliveries + # Poll each resolved DB path once: several slugs can map to one DB when + # HERMES_KANBAN_DB pins the board path. + kb = self.kb + seen_db_paths: set[str] = set() + for board_meta in _list_boards(kb): + slug = board_meta.get("slug") or kb.DEFAULT_BOARD + db_path = board_meta.get("db_path") + try: + resolved_db_path = str(Path(db_path).expanduser().resolve()) if db_path else str(kb.kanban_db_path(slug).resolve()) + except Exception: + resolved_db_path = f"slug:{slug}" + if resolved_db_path in seen_db_paths: + logger.debug("kanban notifier: skipping duplicate board slug %s for DB %s", slug, resolved_db_path) + continue + seen_db_paths.add(resolved_db_path) + self.collect_board(slug) + return self.deliveries + + def _board_has_subs(self, slug: str) -> bool: + """Cheap read-only probe before the writable connect() (schema init, WAL + sidecars, checkpoints); a probe failure falls back to the writable open.""" + try: + if self.kb.count_notify_subs( + board=slug, notifier_profiles=self.notifier_profiles, include_unowned=self.include_unowned, + ) == 0: + logger.debug( + "kanban notifier: board %s has no subscriptions owned by %s; skipping open", + slug, sorted(self.notifier_profiles), + ) + return False + except Exception as exc: + logger.debug( + "kanban notifier: read-only subscription probe failed " + "for board %s (%s); falling back to writable open", + slug, exc, + ) + return True + + def _gc_stale_subs(self, conn: Any, slug: str) -> None: + """Best-effort stale-sub sweep: a failed sweep never blocks delivery; the next hourly gate retries.""" + try: + _purged = self.kb.purge_stale_done_notify_subs(conn, max_age_days=self.gc_retention_days) + if _purged: + logger.info( + "kanban notifier: purged %d stale done/blocked-task subscription(s) on board %s (retention %dd)", + _purged, slug, self.gc_retention_days, + ) + except Exception as _gc_exc: + logger.debug("kanban notifier: stale-sub GC failed for board %s: %s", slug, _gc_exc) + + 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 + old_cursor, cursor, events = self.kb.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, + ) + if not events: + return None + task = self.kb.get_task(conn, sub["task_id"]) + logger.debug( + "kanban notifier: claimed %d event(s) for %s on board %s cursor %s→%s", + len(events), sub["task_id"], slug, old_cursor, cursor, + ) + return {"sub": sub, "old_cursor": old_cursor, "cursor": cursor, "events": events, "task": task, "board": slug} + + def collect_board(self, slug: str) -> None: + """Claim events on one board, appending delivery dicts to ``deliveries``.""" + if not self._board_has_subs(slug): + return + kb = self.kb + try: + conn = kb.connect(board=slug) + except Exception as exc: + logger.debug("kanban notifier: cannot open board %s: %s", slug, exc) + return + try: + if self.gc_due: + self._gc_stale_subs(conn, slug) + # No explicit init_db(): connect() already runs the migration once per + # process, and init_db() would re-run it on a second connection racing + # the first. + subs = kb.list_notify_subs(conn, notifier_profiles=self.notifier_profiles, include_unowned=self.include_unowned) + if not subs: + logger.debug("kanban notifier: board %s has no subscriptions", slug) + for sub in subs: + try: + claimed = self._claim_for_sub(conn, slug, sub) + if claimed is not None: + self.deliveries.append(claimed) + except Exception as sub_exc: + # One bad subscription must not block the rest of the tick. + logger.warning( + "kanban notifier: subscription for %s on board %s failed: %s", + sub.get("task_id"), slug, sub_exc, + ) + finally: + conn.close() + + def _notifier_collect(runner: Any, kb: Any, *, notifier_profile: Optional[str], gc_due: bool, gc_retention_days: int) -> list[dict]: """Claim unseen terminal events for every owned subscription on every board. @@ -93,146 +232,9 @@ def _notifier_collect(runner: Any, kb: Any, *, notifier_profile: Optional[str], hosts; legacy rows without a profile stamp are visible only to the process holding the singleton dispatcher lock. """ - deliveries: list[dict] = [] - include_unowned = runner._owns_kanban_dispatcher_lock() - profile_adapters = getattr(runner, "_profile_adapters", {}) - notifier_profiles = {notifier_profile} - notifier_profiles.update(str(profile).strip() for profile in profile_adapters if str(profile).strip()) - active_platforms = _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. - for _profile_adapter_map in profile_adapters.values(): - active_platforms.update(_platform_names(_profile_adapter_map)) - if not active_platforms: - logger.debug("kanban notifier: no connected adapters; skipping tick") - return deliveries - - # Poll each resolved DB path once: several slugs can map to one DB when - # HERMES_KANBAN_DB pins the board path. - seen_db_paths: set[str] = set() - for board_meta in _list_boards(kb): - slug = board_meta.get("slug") or kb.DEFAULT_BOARD - db_path = board_meta.get("db_path") - try: - resolved_db_path = str(Path(db_path).expanduser().resolve()) if db_path else str(kb.kanban_db_path(slug).resolve()) - except Exception: - resolved_db_path = f"slug:{slug}" - if resolved_db_path in seen_db_paths: - logger.debug("kanban notifier: skipping duplicate board slug %s for DB %s", slug, resolved_db_path) - continue - seen_db_paths.add(resolved_db_path) - _notifier_collect_board( - kb, slug, deliveries, - notifier_profile=notifier_profile, notifier_profiles=notifier_profiles, - include_unowned=include_unowned, profile_adapters=profile_adapters, - active_platforms=active_platforms, gc_due=gc_due, gc_retention_days=gc_retention_days, - ) - return deliveries - - -def _board_has_subs(kb: Any, slug: str, notifier_profiles: set, include_unowned: bool) -> bool: - """Cheap read-only probe before the writable connect() (schema init, WAL - sidecars, checkpoints); a probe failure falls back to the writable open.""" - try: - if kb.count_notify_subs(board=slug, notifier_profiles=notifier_profiles, include_unowned=include_unowned) == 0: - logger.debug( - "kanban notifier: board %s has no subscriptions owned by %s; skipping open", - slug, sorted(notifier_profiles), - ) - return False - except Exception as exc: - logger.debug( - "kanban notifier: read-only subscription probe failed " - "for board %s (%s); falling back to writable open", - slug, exc, - ) - return True - - -def _gc_stale_subs(kb: Any, conn: Any, slug: str, gc_retention_days: int) -> None: - """Best-effort stale-sub sweep: a failed sweep never blocks delivery; the next hourly gate retries.""" - try: - _purged = kb.purge_stale_done_notify_subs(conn, max_age_days=gc_retention_days) - if _purged: - logger.info( - "kanban notifier: purged %d stale done/blocked-task subscription(s) on board %s (retention %dd)", - _purged, slug, gc_retention_days, - ) - except Exception as _gc_exc: - logger.debug("kanban notifier: stale-sub GC failed for board %s: %s", slug, _gc_exc) - - -def _claim_for_sub(kb: Any, conn: Any, slug: str, sub: dict, *, notifier_profile: Optional[str], profile_adapters: dict, active_platforms: set[str]) -> 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 != notifier_profile and not 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, notifier_profile, - ) - return None - platform = (sub.get("platform") or "").lower() - if platform not in active_platforms: - logger.debug( - "kanban notifier: subscription for %s on %s skipped; adapter not connected", - sub.get("task_id"), platform or "", - ) - return None - old_cursor, cursor, events = kb.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, - ) - if not events: - return None - task = kb.get_task(conn, sub["task_id"]) - logger.debug( - "kanban notifier: claimed %d event(s) for %s on board %s cursor %s→%s", - len(events), sub["task_id"], slug, old_cursor, cursor, - ) - return {"sub": sub, "old_cursor": old_cursor, "cursor": cursor, "events": events, "task": task, "board": slug} - - -def _notifier_collect_board( - kb: Any, slug: str, deliveries: list[dict], *, - notifier_profile: Optional[str], notifier_profiles: set, include_unowned: bool, - profile_adapters: dict, active_platforms: set[str], gc_due: bool, gc_retention_days: int, -) -> None: - """Claim events on one board, appending delivery dicts to *deliveries*.""" - if not _board_has_subs(kb, slug, notifier_profiles, include_unowned): - return - try: - conn = kb.connect(board=slug) - except Exception as exc: - logger.debug("kanban notifier: cannot open board %s: %s", slug, exc) - return - try: - if gc_due: - _gc_stale_subs(kb, conn, slug, gc_retention_days) - # No explicit init_db(): connect() already runs the migration once per - # process, and init_db() would re-run it on a second connection racing - # the first. - subs = kb.list_notify_subs(conn, notifier_profiles=notifier_profiles, include_unowned=include_unowned) - if not subs: - logger.debug("kanban notifier: board %s has no subscriptions", slug) - for sub in subs: - try: - claimed = _claim_for_sub( - kb, conn, slug, sub, notifier_profile=notifier_profile, - profile_adapters=profile_adapters, active_platforms=active_platforms, - ) - if claimed is not None: - deliveries.append(claimed) - except Exception as sub_exc: - # One bad subscription must not block the rest of the tick. - logger.warning( - "kanban notifier: subscription for %s on board %s failed: %s", - sub.get("task_id"), slug, sub_exc, - ) - finally: - conn.close() + return _Collector( + runner, kb, notifier_profile=notifier_profile, gc_due=gc_due, gc_retention_days=gc_retention_days, + ).collect() # --- Per-event message formatting: kind -> (msg, wake_handoff, wake_review_detail) --- @@ -250,6 +252,9 @@ def _clip(ev: Any, key: str, fmt: str, limit: int) -> str: return fmt.format(str(value)[:limit]) if value else "" +_NL = "\n{}" + + def _first_line(text: str, limit: int) -> str: lines = text.strip().splitlines() return lines[0][:limit] if lines else text[:limit] @@ -267,30 +272,6 @@ def _fmt_completed(ev, n) -> tuple: return f"✔ {n.head} done — {n.title}{handoff}", wake_handoff, None -def _fmt_blocked(ev, n) -> tuple: - reason = _clip(ev, "reason", ": {}", 160) - return f"⏸ {n.head} blocked{reason}", None, None - - -def _fmt_gave_up(ev, n) -> tuple: - err = _clip(ev, "error", "\n{}", 200) - return f"✖ {n.head} gave up after repeated spawn failures{err}", None, None - - -def _fmt_crashed(ev, n) -> tuple: - return f"✖ {n.head} worker crashed (pid gone); dispatcher will retry", None, None - - -def _fmt_timed_out(ev, n) -> tuple: - limit = _payload(ev, "limit_seconds") - return f"⏱ {n.head} timed out (max_runtime={int(limit) if limit else 0}s); will retry", None, None - - -def _fmt_status(ev, n) -> tuple: - new_status = _payload(ev, "status") - return f"🔄 {n.head} → {str(new_status) if new_status else ''}", None, None - - def _fmt_review_requested(ev, n) -> tuple: # Implementation done; task moved to the review lane. Carry the handoff # into the wake turn like ``completed`` so the reviewer needn't re-read the board. @@ -317,28 +298,29 @@ def _fmt_changes_requested(ev, n) -> tuple: return msg, None, reason_text -def _fmt_block_loop_detected(ev, n) -> tuple: - # Re-blocked for the same cause past the limit and routed to `triage` - # for a human. It emits no blocked/status event, so ping loudly here. - recurrences = _payload(ev, "recurrences") - rc = f" (blocked {recurrences}x for the same cause)" if recurrences else "" - reason = _clip(ev, "reason", ": {}", 160) - return f"🛑 {n.head} routed to TRIAGE — needs a human decision{rc}{reason}", None, None - - # archived / unblocked are claimed (so the cursor advances past them) but # intentionally silent (no formatter), and excluded from _WAKE_KINDS so they # never wake the creator. _EVENT_FORMATTERS: dict[str, Callable[[Any, "_KanbanNotification"], tuple]] = { "completed": _fmt_completed, - "blocked": _fmt_blocked, - "gave_up": _fmt_gave_up, - "crashed": _fmt_crashed, - "timed_out": _fmt_timed_out, - "status": _fmt_status, + "blocked": lambda ev, n: (f"⏸ {n.head} blocked{_clip(ev, 'reason', ': {}', 160)}", None, None), + "gave_up": lambda ev, n: ( + f"✖ {n.head} gave up after repeated spawn failures{_clip(ev, 'error', _NL, 200)}", None, None, + ), + "crashed": lambda ev, n: (f"✖ {n.head} worker crashed (pid gone); dispatcher will retry", None, None), + "timed_out": lambda ev, n: ( + f"⏱ {n.head} timed out (max_runtime={int(_payload(ev, 'limit_seconds') or 0)}s); will retry", None, None, + ), + "status": lambda ev, n: (f"🔄 {n.head} → {_payload(ev, 'status') or ''}", None, None), "review_requested": _fmt_review_requested, "changes_requested": _fmt_changes_requested, - "block_loop_detected": _fmt_block_loop_detected, + # Re-blocked for the same cause past the limit and routed to `triage` for a + # human. It emits no blocked/status event, so ping loudly here. + "block_loop_detected": lambda ev, n: ( + f"🛑 {n.head} routed to TRIAGE — needs a human decision" + f"{_clip(ev, 'recurrences', ' (blocked {}x for the same cause)', 200)}{_clip(ev, 'reason', ': {}', 160)}", + None, None, + ), } @@ -473,11 +455,15 @@ class _KanbanNotification: self.task_id, self.platform_str, self.sub["chat_id"], self.sub_profile or "default", self.wake_kinds, ) - async def push_wake(self) -> None: - """Wake the creator session behind a push adapter; raises on failure.""" - from gateway.session import SessionSource + async def wake(self) -> None: + """Wake the creator session (raises on failure): push adapters get a full SessionSource, non-push a raw self-post.""" from gateway.wake import deliver_wake sub = self.sub + if not self.is_push_adapter: + await deliver_wake(self.adapter, text=self.synth, session_id=self.session_key) + self._log_woke() + return + from gateway.session import SessionSource # Rebuild the creator's real session scope from the persisted chat_type: # build_session_key() keys DMs differently from group/thread, so a # hardcoded "group" mis-routed DM/thread creators into a fresh session. @@ -580,7 +566,7 @@ class _KanbanNotification: await self.rewind() return self.adapter = adapter - from gateway.wake import adapter_supports_push, deliver_wake + from gateway.wake import adapter_supports_push self.is_push_adapter = adapter_supports_push(adapter) if not await self._send_pings(): @@ -590,24 +576,18 @@ class _KanbanNotification: self.build_wake_text() wake_kinds, is_push = self.wake_kinds, self.is_push_adapter - if not is_push and wake_kinds and self.session_key: - # Self-post IS the delivery: must succeed BEFORE the cursor advances. + # Non-push self-post, or wake-only push sub: the wake IS the delivery + # and must succeed BEFORE the cursor advances. + if wake_kinds and (not self.send_passive if is_push else bool(self.session_key)): try: - await deliver_wake(adapter, text=self.synth, session_id=self.session_key) - self._log_woke() + await self.wake() self.clear_failures() except Exception as _wk_err: - await self._wake_failed("kanban notifier: wake self-post failed for %s (attempt %d/%d): %s", _wk_err) - return - - if is_push and not self.send_passive and wake_kinds: - # Wake-only push sub: the wake is the sole delivery and must - # succeed BEFORE the cursor advances. - try: - await self.push_wake() - self.clear_failures() - except Exception as _wk_err: - await self._wake_failed("kanban notifier: wake-only delivery failed for %s (attempt %d/%d): %s", _wk_err) + await self._wake_failed( + "kanban notifier: wake-only delivery failed for %s (attempt %d/%d): %s" if is_push + else "kanban notifier: wake self-post failed for %s (attempt %d/%d): %s", + _wk_err, + ) return # Delivery complete: advance the cursor (the dedup mechanism). @@ -619,7 +599,7 @@ class _KanbanNotification: # advanced; the wake stays best-effort, but log at WARNING so a # persistently failing wake is visible. try: - await self.push_wake() + await self.wake() except Exception as _wk_err: logger.warning("kanban notifier: wakeup injection failed for %s: %s", self.task_id, _wk_err, exc_info=True) # Unsubscribe only on archive; ``done`` is reversible.