refactor(gateway): kanban notifier collector class, lambda formatters, single wake path

This commit is contained in:
Teknium
2026-09-02 18:52:31 -07:00
parent 3b6529d5ec
commit 320a7caa6a
+179 -199
View File
@@ -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 "<missing>",
)
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 "<missing>",
)
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.