From fd2bfa1893ebc7b085f93af093eb011978ba2f38 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 15:57:11 -0700 Subject: [PATCH] refactor(gateway/kanban): lift notifier delivery into _KanbanNotification + event-formatter table, dispatcher into _DispatcherSettings/_KanbanDispatcher; unify board enumeration/tick sleep --- gateway/kanban_watchers.py | 1085 ++----------------------- gateway/kanban_watchers_common.py | 35 + gateway/kanban_watchers_dispatcher.py | 375 +++++++++ gateway/kanban_watchers_notifier.py | 722 ++++++++++++++++ 4 files changed, 1190 insertions(+), 1027 deletions(-) create mode 100644 gateway/kanban_watchers_common.py create mode 100644 gateway/kanban_watchers_dispatcher.py create mode 100644 gateway/kanban_watchers_notifier.py diff --git a/gateway/kanban_watchers.py b/gateway/kanban_watchers.py index 2df15b25e6..5929061f89 100644 --- a/gateway/kanban_watchers.py +++ b/gateway/kanban_watchers.py @@ -2,48 +2,30 @@ Background loops that subscribe to kanban boards, deliver notifications and artifacts, and drive the multi-agent dispatcher. They use only ``self`` state, -so they live on a mixin ``GatewayRunner`` inherits. +so they live on a mixin ``GatewayRunner`` inherits. Per-tick work lives in +``kanban_watchers_notifier`` / ``kanban_watchers_dispatcher``. """ from __future__ import annotations import asyncio -import logging +import contextlib import os -import re -import sqlite3 import time -from contextvars import Context from pathlib import Path from typing import Any, Callable, Optional -from agent.i18n import t -import contextlib - -# Keep the logger name run.py used so extracted log records are unchanged. -logger = logging.getLogger("gateway.run") - - -_LOCAL_PATH_RE = re.compile( - r"(? str: - """Return a mobile-friendly review reason safe for external delivery.""" - from agent.redact import redact_sensitive_text - - reason = redact_sensitive_text( - "" if value is None else str(value), - force=True, - redact_url_credentials=True, - ) - reason = _LOCAL_PATH_RE.sub("[local path]", reason) - reason = " ".join(reason.split()) - if len(reason) > limit: - reason = reason[: limit - 1].rstrip() + "…" - return reason def _resolve_auto_decompose_settings( @@ -66,9 +48,7 @@ def _resolve_auto_decompose_settings( per_tick = int(kcfg.get("auto_decompose_per_tick", 3) or 3) except (TypeError, ValueError): per_tick = 3 - if per_tick < 1: - per_tick = 1 - return enabled, per_tick + return enabled, max(per_tick, 1) def _kanban_dispatch_allowed() -> bool: @@ -84,22 +64,6 @@ def _kanban_dispatch_allowed() -> bool: return not check_paused("kanban", logger) -def _run_in_fresh_context(func: Callable[..., Any], /, *args: Any) -> Any: - """Run *func* in an empty ``Context`` so request-local ContextVars stay behind. - - ``asyncio.to_thread`` copies the caller's context; a lingering - ``delegate_task`` child marker would make ``write_txn`` false-trip for - these process-owned writers. An empty Context keeps the DB guard intact - for real children without exempting dispatcher writes. - """ - return Context().run(func, *args) - - -async def _to_thread_process_service(func: Callable[..., Any], /, *args: Any) -> Any: - """Offload blocking process-service work without inheriting request ContextVars.""" - return await asyncio.to_thread(_run_in_fresh_context, func, *args) - - def _acquire_singleton_lock(lock_path) -> "tuple[Optional[object], str]": """Take the exclusive, non-blocking advisory lock for the sole dispatcher. @@ -141,37 +105,15 @@ def _release_singleton_lock(handle) -> None: handle.close() -def _wake_scope_id(adapter: Any, sub: dict) -> Optional[str]: - """Return the tenant scope (Slack workspace) a subscription's wake keys to. +def _gc_retention_days() -> int: + """``kanban.done_sub_retention_days`` (default 30; 0 disables), re-read per sweep; fails safe to 30.""" + try: + from hermes_cli.config import load_config as _load_cfg - ``build_session_key()`` includes ``scope_id`` on multi-tenant platforms, - so the wake must carry the same scope as inbound messages. Persisted - ``delivery_metadata`` wins (it records the creating scope); the adapter's - live chat → scope map only covers rows without metadata. ``None`` means - unscoped, matching an unscoped platform's key. - """ - delivery_meta = sub.get("delivery_metadata") - if isinstance(delivery_meta, dict): - for key in ("scope_id", "slack_team_id", "team_id"): - value = delivery_meta.get(key) - if value: - return str(value) - resolver = getattr(adapter, "scope_id_for_chat", None) - if callable(resolver): - try: - resolved = resolver(str(sub.get("chat_id") or "")) - except Exception as exc: - # An adapter-side lookup failure yields no scope, never an error. - logger.debug( - "kanban notifier: scope lookup failed for chat %s: %s", - sub.get("chat_id"), - exc, - exc_info=True, - ) - return None - if resolved: - return str(resolved) - return None + _kanban_cfg = (_load_cfg() or {}).get("kanban") or {} + return int(_kanban_cfg.get("done_sub_retention_days", 30)) + except Exception: + return 30 class GatewayKanbanWatchersMixin: @@ -187,6 +129,14 @@ class GatewayKanbanWatchersMixin: self._kanban_dispatcher_lock_handle = None _release_singleton_lock(handle) + async def _sleep_between_ticks(self, interval: float) -> None: + """Sleep *interval* (floored to 1s) in 1s slices so stop() never waits a full interval.""" + interval = max(interval, 1.0) + slept = 0.0 + while slept < interval and self._running: + await asyncio.sleep(min(1.0, interval - slept)) + slept += 1.0 + async def _kanban_notifier_watcher(self, interval: float = 5.0) -> None: """Poll ``kanban_notify_subs`` and deliver terminal events to users. @@ -199,10 +149,7 @@ class GatewayKanbanWatchersMixin: users when the dispatcher respawned a crashed task. All SQLite work runs in a thread; one tick's failure never stops the - next. Iterates every board on disk per tick; each gateway polls only - subscriptions owned by profiles whose adapters it hosts, and legacy - rows without a profile stamp are visible only to the process holding - the singleton dispatcher lock. + next. Iterates every board on disk per tick. """ from gateway.config import Platform as _Platform try: @@ -211,15 +158,6 @@ class GatewayKanbanWatchersMixin: logger.warning("kanban notifier: kanban_db not importable; notifier disabled") return - # "status" covers dashboard drag-drop and `_set_status_direct()`. - # ``review_requested`` wakes the origin like a block but is not one; - # the task is not archived so later review cycles keep notifying. - TERMINAL_KINDS = ("completed", "blocked", "gave_up", "crashed", "timed_out", "status", "archived", "unblocked", "block_loop_detected", "review_requested", "changes_requested") - # Consecutive send failures (adapter raised OR reported - # SendResult(success=False)) before a sub is dropped as a dead chat. - # 12 ≈ 60s at the 5s cadence: a transient API outage must not - # permanently unsubscribe a live review-gate channel. - MAX_SEND_FAILURES = 12 sub_fail_counts: dict[tuple, int] = getattr( self, "_kanban_sub_fail_counts", {} ) @@ -234,608 +172,31 @@ class GatewayKanbanWatchersMixin: # Stale done-sub GC: subs survive ``done``, so boards that never # archive would accumulate rows scanned every tick. One DELETE per - # board, at startup and at most hourly; retention is - # kanban.done_sub_retention_days (default 30; 0 disables), re-read - # at each sweep. + # board, at startup and at most hourly. _GC_INTERVAL_SECONDS = 3600.0 _gc_next_at = 0.0 # 0 → sweep on the first tick after startup while self._running: try: _gc_due = time.monotonic() >= _gc_next_at - _gc_retention_days = 30 + _retention = 30 if _gc_due: _gc_next_at = time.monotonic() + _GC_INTERVAL_SECONDS - try: - from hermes_cli.config import load_config as _load_cfg + _retention = _gc_retention_days() - _kanban_cfg = (_load_cfg() or {}).get("kanban") or {} - _gc_retention_days = int( - _kanban_cfg.get("done_sub_retention_days", 30) - ) - except Exception: - _gc_retention_days = 30 # fail safe on the shipped default - - def _collect(): - deliveries: list[dict] = [] - include_unowned = self._owns_kanban_dispatcher_lock() - notifier_profiles = {notifier_profile} - notifier_profiles.update( - str(profile).strip() - for profile in getattr(self, "_profile_adapters", {}) - if str(profile).strip() - ) - active_platforms = { - getattr(platform, "value", str(platform)).lower() - for platform in self.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 getattr(self, "_profile_adapters", {}).values(): - active_platforms.update( - getattr(platform, "value", str(platform)).lower() - for platform in _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. - try: - boards = _kb.list_boards(include_archived=False) - except Exception: - boards = [_kb.read_board_metadata(_kb.DEFAULT_BOARD)] - seen_db_paths: set[str] = set() - for board_meta in boards: - 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) - # Cheap read-only probe before the writable connect() - # (schema init, WAL sidecars, checkpoints) — a board - # with no subscriptions has nothing to notify. - 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), - ) - continue - except Exception as exc: - logger.debug( - "kanban notifier: read-only subscription probe failed " - "for board %s (%s); falling back to writable open", - slug, exc, - ) - try: - conn = _kb.connect(board=slug) - except Exception as exc: - logger.debug("kanban notifier: cannot open board %s: %s", slug, exc) - continue - try: - if _gc_due: - # Best-effort: 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, - ) - # 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: - owner_profile = sub.get("notifier_profile") or None - if owner_profile and owner_profile != notifier_profile: - _owner_adapters = getattr(self, "_profile_adapters", {}).get(owner_profile) - if not _owner_adapters: - 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, - ) - continue - 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 "", - ) - continue - 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: - continue - 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, - ) - deliveries.append({ - "sub": sub, - "old_cursor": old_cursor, - "cursor": cursor, - "events": events, - "task": task, - "board": slug, - }) - 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 deliveries - - deliveries = await asyncio.to_thread(_collect) + deliveries = await asyncio.to_thread( + _notifier_collect, self, _kb, + notifier_profile=notifier_profile, + gc_due=_gc_due, + gc_retention_days=_retention, + ) for d in deliveries: - sub = d["sub"] - task = d["task"] - board_slug = d.get("board") - platform_str = (sub["platform"] or "").lower() - - async def _rewind() -> None: - await _to_thread_process_service( - self._kanban_rewind, - sub, - d["cursor"], - d.get("old_cursor", 0), - board_slug, - ) - - try: - plat = _Platform(platform_str) - except ValueError: - # Unknown platform: advance the cursor so it can't replay forever. - await _to_thread_process_service( - self._kanban_advance, sub, d["cursor"], board_slug, - ) - continue - sub_profile = sub.get("notifier_profile") or "" - # 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._authorization_adapter(plat, sub_profile or None) - if adapter is None: - logger.debug( - "kanban notifier: adapter %s disconnected before delivery for %s; rewinding claim", - platform_str, sub["task_id"], - ) - await _rewind() - continue - title = (task.title if task else sub["task_id"])[:120] - board_tag = f"[{board_slug}] " if board_slug else "" - # Hoisted: the wake self-post path (the loop's ``else``) - # needs the key even when every event was skipped. - sub_key = ( - sub["task_id"], sub["platform"], - sub["chat_id"], sub.get("thread_id") or "", - ) - - async def _delivery_failed(fmt: str, prefix: tuple, drop_fmt: str, exc: Exception, exc_info: bool) -> None: - """Bump the failure counter; drop the sub past the limit, else rewind the claim so the next tick retries.""" - fails = sub_fail_counts.get(sub_key, 0) + 1 - sub_fail_counts[sub_key] = fails - logger.warning(fmt, *prefix, fails, MAX_SEND_FAILURES, exc, exc_info=exc_info) - if fails >= MAX_SEND_FAILURES: - logger.warning(drop_fmt, sub["task_id"], platform_str, fails) - await _to_thread_process_service(self._kanban_unsub, sub, board_slug) - sub_fail_counts.pop(sub_key, None) - else: - await _rewind() - - mode = sub.get("delivery_mode") or "notify" - wake_agent = mode in ("notify+wake", "wake") - send_passive = mode != "wake" - # Worker handoff carried into the synthetic wake turn so the - # woken creator doesn't re-decompose work already on the board. - wake_handoff = "" - wake_review_detail = "" - from gateway.wake import adapter_supports_push - - for ev in d["events"]: - kind = ev.kind - # Attribute the ping to the worker that did the work. - who = (task.assignee if task and task.assignee else None) - tag = f"@{who} " if who else "" - if kind == "completed": - # Prefer the run summary from the event payload; - # fall back to task.result for legacy rows. - handoff = "" - payload_summary = None - if ev.payload and ev.payload.get("summary"): - payload_summary = str(ev.payload["summary"]) - if payload_summary: - lines = payload_summary.strip().splitlines() - h = lines[0][:200] if lines else payload_summary[:200] - handoff = f"\n{h}" - wake_handoff = h - elif task and task.result: - lines = task.result.strip().splitlines() - r = lines[0][:160] if lines else task.result[:160] - handoff = f"\n{r}" - wake_handoff = r - msg = ( - f"✔ {board_tag}{tag}Kanban {sub['task_id']} done" - f" — {title}{handoff}" - ) - elif kind == "blocked": - reason = "" - if ev.payload and ev.payload.get("reason"): - reason = f": {str(ev.payload['reason'])[:160]}" - msg = f"⏸ {board_tag}{tag}Kanban {sub['task_id']} blocked{reason}" - elif kind == "gave_up": - err = "" - if ev.payload and ev.payload.get("error"): - err = f"\n{str(ev.payload['error'])[:200]}" - msg = ( - f"✖ {board_tag}{tag}Kanban {sub['task_id']} gave up " - f"after repeated spawn failures{err}" - ) - elif kind == "crashed": - msg = ( - f"✖ {board_tag}{tag}Kanban {sub['task_id']} worker crashed " - f"(pid gone); dispatcher will retry" - ) - elif kind == "timed_out": - limit = 0 - if ev.payload and ev.payload.get("limit_seconds"): - limit = int(ev.payload["limit_seconds"]) - msg = ( - f"⏱ {board_tag}{tag}Kanban {sub['task_id']} timed out " - f"(max_runtime={limit}s); will retry" - ) - elif kind == "status": - new_status = "" - if ev.payload and ev.payload.get("status"): - new_status = str(ev.payload["status"]) - msg = f"🔄 {board_tag}{tag}Kanban {sub['task_id']} → {new_status}" - elif kind == "review_requested": - # Implementation done; task moved to the review lane. - handoff = "" - if ev.payload and ev.payload.get("summary"): - summary = str(ev.payload["summary"]) - handoff = f"\n{summary[:200]}" - # Carry the handoff into the wake turn like - # ``completed`` so the reviewer needn't re-read the board. - lines = summary.strip().splitlines() - wake_handoff = ( - lines[0][:200] if lines else summary[:200] - ) - msg = ( - f"👀 {board_tag}{tag}Kanban {sub['task_id']} ready for review" - f" — {title}{handoff}" - ) - elif kind == "changes_requested": - payload = ev.payload or {} - reason = _safe_review_reason(payload.get("reason")) - reviewer = _safe_review_reason(payload.get("reviewer"), 48) - implementer = _safe_review_reason(payload.get("implementer"), 48) - reason_text = reason or "reviewer feedback requires changes" - provenance = "" - if reviewer: - provenance += f" — reviewer @{reviewer}" - if implementer: - provenance += f" → implementer @{implementer}" - msg = ( - f"🛑 {board_tag}Kanban {sub['task_id']} review requested " - f"changes/BLOCK: {reason_text}{provenance}" - ) - wake_review_detail = reason_text - elif kind == "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. - reason = "" - recurrences = None - if ev.payload: - if ev.payload.get("reason"): - reason = f": {str(ev.payload['reason'])[:160]}" - recurrences = ev.payload.get("recurrences") - rc = f" (blocked {recurrences}x for the same cause)" if recurrences else "" - msg = ( - f"🛑 {board_tag}{tag}Kanban {sub['task_id']} routed to TRIAGE" - f" — needs a human decision{rc}{reason}" - ) - else: - # archived / unblocked are claimed (so the cursor - # advances past them) but intentionally silent, and - # excluded from _WAKE_KINDS so they never wake the creator. - continue - delivery_metadata = sub.get("delivery_metadata") - metadata: dict[str, Any] = ( - dict(delivery_metadata) - if isinstance(delivery_metadata, dict) - else {} - ) - - if sub.get("thread_id") and not metadata.get("thread_id"): - metadata["thread_id"] = sub["thread_id"] - # Non-push adapters (api_server) always report - # SendResult(success=False) from send(); treating that - # as failure would drop the sub forever and make the - # wake path in this loop's ``else`` unreachable. Skip - # the doomed send; the self-post below IS the delivery. - if not adapter_supports_push(adapter) and wake_agent: - logger.debug( - "kanban notifier: adapter %s has no push " - "channel; skipping text ping for %s, relying " - "on wake self-post instead", - platform_str, sub["task_id"], - ) - # Counter is resolved by the self-post outcome, not here. - continue - if not send_passive: - # Wake-only: the wake path below is the sole delivery - # and resolves the failure counter. - continue - try: - _send_res = await adapter.send( - sub["chat_id"], msg, metadata=metadata, - ) - # SendResult(success=False) without an exception is - # a FAILED delivery (else the event is lost); None / - # non-SendResult keeps the "no exception == delivered" contract. - if getattr(_send_res, "success", True) is False: - raise RuntimeError( - "adapter send() reported failure: " - f"{getattr(_send_res, 'error', None) or 'unknown error'}" - ) - logger.debug( - "kanban notifier: delivered %s event for %s to %s/%s on board %s", - kind, sub["task_id"], platform_str, sub["chat_id"], board_slug, - ) - # Upload artifact paths from the completion payload - # / legacy result as native files. Only on - # ``completed`` so retries never spam attachments. - if kind == "completed": - try: - await self._deliver_kanban_artifacts( - adapter=adapter, - chat_id=sub["chat_id"], - metadata=metadata, - event_payload=getattr(ev, "payload", None), - task=task, - ) - except Exception as art_exc: - logger.debug( - "kanban notifier: artifact delivery for %s failed: %s", - sub["task_id"], art_exc, - ) - sub_fail_counts.pop(sub_key, None) - except Exception as exc: - await _delivery_failed( - "kanban notifier: send failed for %s on %s " - "(attempt %d/%d): %s", - (sub["task_id"], platform_str), - "kanban notifier: dropping subscription " - "%s on %s after %d consecutive send failures", - exc, False, - ) - break - else: - # All text pings delivered (or skipped for non-push / - # wake-only). Cursor advance ordering by adapter class: - # * push + notify: the text send WAS the delivery → - # advance now; wake injection stays best-effort. - # * non-push or wake-only: the wake IS the delivery → - # it runs FIRST and the cursor advances only after it - # succeeds; failure rewinds like a failed send(). - task_terminal = task and task.status == "archived" - # Kinds that hand a decision back to the origin, which - # must take a turn. status/archived/unblocked are bookkeeping. - _WAKE_KINDS = ( - "completed", "gave_up", "crashed", "timed_out", - "blocked", "review_requested", "changes_requested", - "block_loop_detected", - ) - _wake_kinds = ( - {ev.kind for ev in d["events"] if ev.kind in _WAKE_KINDS} - if wake_agent - else set() - ) - _is_push_adapter = adapter_supports_push(adapter) - _session_key = "" - _synth = "" - if _wake_kinds: - if _is_push_adapter: - _session_key = getattr(task, "session_id", None) or "" - else: - # Non-push wakes target sub["chat_id"] (the raw - # session id the subscriber registered). - # task.session_id may be a WORKER session for - # child tasks; use it only for legacy rows. - _session_key = ( - sub["chat_id"] - or getattr(task, "session_id", None) - or "" - ) - _title = (task.title if task else sub["task_id"])[:120] - _assignee = task.assignee if task else "" - # i18n keys: gateway.kanban.wake. for each _WAKE_KINDS entry. - _parts = [t(f"gateway.kanban.wake.{k}") for k in _WAKE_KINDS if k in _wake_kinds] - _status = t("gateway.kanban.wake.status_joiner").join(_parts) or t("gateway.kanban.wake.status_default") - _synth = t( - "gateway.kanban.wake.message", - task_id=sub["task_id"], - status=_status, - title=_title, - assignee=_assignee, - board=board_slug, - ) - # Label as an automatic notification and carry the - # handoff so the creator inspects the board instead - # of re-decomposing. - if wake_handoff: - _synth += "\n" + t( - "gateway.kanban.wake.handoff", - summary=wake_handoff, - ) - if wake_review_detail: - _synth += "\n" + t( - "gateway.kanban.wake.review_detail", - reason=wake_review_detail, - ) - _synth += "\n\n" + t( - "gateway.kanban.wake.guidance" - ) - - if not _is_push_adapter and _wake_kinds and _session_key: - # Self-post IS the delivery: must succeed BEFORE the cursor advances. - from gateway.wake import deliver_wake - - try: - await deliver_wake( - adapter, - text=_synth, - session_id=_session_key, - ) - logger.info( - "kanban notifier: woke agent for %s on %s/%s profile=%s events=%s", - sub["task_id"], platform_str, sub["chat_id"], sub_profile or "default", _wake_kinds, - ) - sub_fail_counts.pop(sub_key, None) - except Exception as _wk_err: - await _delivery_failed( - "kanban notifier: wake self-post failed " - "for %s (attempt %d/%d): %s", - (sub["task_id"],), - "kanban notifier: dropping subscription " - "%s on %s after %d consecutive wake failures", - _wk_err, True, - ) - continue - - async def _push_wake() -> None: - """Wake the creator session behind a push adapter; raises on failure.""" - from gateway.session import SessionSource - from gateway.wake import deliver_wake - # 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. Legacy rows may carry chat_type in - # delivery_metadata; last resort is "group". A - # mismatch only degrades to a fresh session. - _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() - _chat_type = _chat_type or "group" - _source = SessionSource( - platform=plat, - chat_id=sub["chat_id"], - chat_type=_chat_type, - thread_id=sub.get("thread_id") or None, - user_id=sub.get("user_id"), - user_id_alt=sub.get("user_id_alt"), - profile=sub_profile or None, - scope_id=_wake_scope_id(adapter, sub), - ) - await deliver_wake( - adapter, - text=_synth, - session_id=_session_key, - source=_source, - ) - logger.info( - "kanban notifier: woke agent for %s on %s/%s profile=%s events=%s", - sub["task_id"], platform_str, sub["chat_id"], sub_profile or "default", _wake_kinds, - ) - - if _is_push_adapter and not send_passive and _wake_kinds: - # Wake-only push sub: the wake is the sole delivery - # and must succeed BEFORE the cursor advances. - try: - await _push_wake() - sub_fail_counts.pop(sub_key, None) - except Exception as _wk_err: - await _delivery_failed( - "kanban notifier: wake-only delivery failed " - "for %s (attempt %d/%d): %s", - (sub["task_id"],), - "kanban notifier: dropping subscription " - "%s on %s after %d consecutive wake failures", - _wk_err, True, - ) - continue - - # Delivery complete: advance the cursor (the dedup mechanism). - await _to_thread_process_service( - self._kanban_advance, sub, d["cursor"], board_slug, - ) - if not _is_push_adapter: - sub_fail_counts.pop(sub_key, None) - if _is_push_adapter and send_passive and _wake_kinds: - # notify+wake: text ping was the delivery and the - # cursor has advanced; the wake stays best-effort, - # but log at WARNING so a persistently failing wake is visible. - try: - await _push_wake() - except Exception as _wk_err: - logger.warning( - "kanban notifier: wakeup injection failed for %s: %s", - sub["task_id"], _wk_err, exc_info=True, - ) - # Unsubscribe only on archive; ``done`` is reversible. - if task_terminal: - await _to_thread_process_service( - self._kanban_unsub, sub, board_slug, - ) + await _KanbanNotification( + self, d, platform_cls=_Platform, sub_fail_counts=sub_fail_counts, + ).deliver() except Exception as exc: logger.warning("kanban notifier tick failed: %s", exc) - # Sleep with cancellation checks. - for _ in range(int(max(1, interval))): - if not self._running: - return - await asyncio.sleep(1) + await self._sleep_between_ticks(interval) def _kanban_sub_op(self, board: Optional[str], op: str, sub: dict, **extra: Any) -> None: """Sync helper (runs in to_thread): call ``kanban_db.`` for one subscription on its board.""" @@ -891,8 +252,6 @@ class GatewayKanbanWatchersMixin: deduplicated, missing files are skipped (may be mentioned for reference only), and upload errors are logged, never raised. """ - from pathlib import Path as _Path - candidates: list[str] = [] seen: set[str] = set() @@ -900,9 +259,7 @@ class GatewayKanbanWatchersMixin: if not path: return expanded = os.path.expanduser(path) - if expanded in seen: - return - if not os.path.isfile(expanded): + if expanded in seen or not os.path.isfile(expanded): return seen.add(expanded) candidates.append(expanded) @@ -921,8 +278,7 @@ class GatewayKanbanWatchersMixin: _add(p) if task is not None and getattr(task, "result", None): - result_text = str(task.result) - paths, _ = adapter.extract_local_files(result_text) + paths, _ = adapter.extract_local_files(str(task.result)) for p in paths: _add(p) @@ -940,8 +296,8 @@ class GatewayKanbanWatchersMixin: from urllib.parse import quote as _quote # Images ride one send_multiple_images call (batch uploads on Signal/Slack). - image_paths = [p for p in candidates if _Path(p).suffix.lower() in _IMAGE_EXTS] - other_paths = [p for p in candidates if _Path(p).suffix.lower() not in _IMAGE_EXTS] + image_paths = [p for p in candidates if Path(p).suffix.lower() in _IMAGE_EXTS] + other_paths = [p for p in candidates if Path(p).suffix.lower() not in _IMAGE_EXTS] if image_paths: try: @@ -955,9 +311,8 @@ class GatewayKanbanWatchersMixin: ) for path in other_paths: - ext = _Path(path).suffix.lower() try: - if ext in _VIDEO_EXTS: + if Path(path).suffix.lower() in _VIDEO_EXTS: await adapter.send_video( chat_id=chat_id, video_path=path, metadata=metadata, ) @@ -1031,98 +386,8 @@ class GatewayKanbanWatchersMixin: "on config control alone.", _lock_path, ) - try: - interval = float(kanban_cfg.get("dispatch_interval_seconds", 60) or 60) - except (ValueError, TypeError): - logger.warning( - "kanban dispatcher: invalid dispatch_interval_seconds=%r, using default 60", - kanban_cfg.get("dispatch_interval_seconds"), - ) - interval = 60.0 - interval = max(interval, 1.0) # sanity floor — tighter than this is a footgun - - max_spawn = kanban_cfg.get("max_spawn", None) - if max_spawn is not None: - logger.info("kanban dispatcher: max_spawn=%s", max_spawn) - - def _positive_int_setting(key: str) -> Optional[int]: - """Parse an optional ``kanban.`` int cap; None when unset or invalid (< 1 is invalid).""" - raw = kanban_cfg.get(key, None) - if raw is None: - return None - try: - value = int(raw) - except (TypeError, ValueError): - logger.warning("kanban dispatcher: invalid kanban.%s=%r; ignoring", key, raw) - return None - if value < 1: - logger.warning("kanban dispatcher: kanban.%s=%r is below 1; ignoring", key, raw) - return None - logger.info("kanban dispatcher: %s=%d", key, value) - return value - - # Cap simultaneously running tasks so slow workers don't pile up and - # time out. Explicit config wins; otherwise a memory-derived default - # (unbounded fan-out swap-thrashes small hosts), or None where total - # memory can't be read. - max_in_progress = _positive_int_setting("max_in_progress") - effective_max_in_progress = _kb.resolve_max_in_progress(max_in_progress) - if max_in_progress is None and effective_max_in_progress is not None: - logger.info( - "kanban dispatcher: kanban.max_in_progress unset; using " - "memory-derived default max_in_progress=%d " - "(set kanban.max_in_progress in config.yaml to override)", - effective_max_in_progress, - ) - max_in_progress = effective_max_in_progress - - raw_failure_limit = kanban_cfg.get("failure_limit", _kb.DEFAULT_FAILURE_LIMIT) - try: - failure_limit = int(raw_failure_limit) - except (TypeError, ValueError): - logger.warning( - "kanban dispatcher: invalid kanban.failure_limit=%r; using default %d", - raw_failure_limit, - _kb.DEFAULT_FAILURE_LIMIT, - ) - failure_limit = _kb.DEFAULT_FAILURE_LIMIT - if failure_limit < 1: - logger.warning( - "kanban dispatcher: kanban.failure_limit=%r is below 1; using default %d", - raw_failure_limit, - _kb.DEFAULT_FAILURE_LIMIT, - ) - failure_limit = _kb.DEFAULT_FAILURE_LIMIT - - # 0 disables stale detection. - raw_stale = kanban_cfg.get("dispatch_stale_timeout_seconds", 0) - try: - stale_timeout_seconds = int(raw_stale or 0) - except (TypeError, ValueError): - logger.warning( - "kanban dispatcher: invalid kanban.dispatch_stale_timeout_seconds=%r; " - "disabling stale detection", - raw_stale, - ) - stale_timeout_seconds = 0 - - # Requeue 'running' cards with broken claim bookkeeping (zombie-card - # reconciliation); false keeps orphans frozen for manual forensics. - reconcile_orphans = bool(kanban_cfg.get("reconcile_orphans", True)) - - # Fallback profile for tasks created without an assignee (e.g. via the - # dashboard). Empty (the schema default) keeps skipping them. - default_assignee = (kanban_cfg.get("default_assignee") or "").strip() or None - if default_assignee: - logger.info( - "kanban dispatcher: default_assignee=%r (unassigned ready tasks " - "will route to this profile)", - default_assignee, - ) - - # Per-profile concurrency cap: no single profile's local model / API - # quota / browser pool gets overwhelmed by a fan-out. - max_in_progress_per_profile = _positive_int_setting("max_in_progress_per_profile") + settings = _resolve_dispatcher_settings(kanban_cfg, _kb) + interval = settings.interval # Initial delay so adapters are wired before workers spawn (matches the notifier). await asyncio.sleep(5) @@ -1133,218 +398,7 @@ class GatewayKanbanWatchersMixin: HEALTH_WINDOW = 6 bad_ticks = 0 last_warn_at = 0 - # Quarantine corrupt-looking board DBs, but retry after a while: - # transient WAL/open races can look like "malformed" for one tick. - CORRUPT_BOARD_RETRY_AFTER_SECONDS = 300 - disabled_corrupt_boards: dict[ - str, tuple[tuple[str, int | None, int | None], float] - ] = {} - - def _board_db_fingerprint(slug: str) -> tuple[str, int | None, int | None]: - path = _kb.kanban_db_path(slug) - try: - resolved = str(path.expanduser().resolve()) - except Exception: - resolved = str(path) - try: - stat = path.stat() - except OSError: - return (resolved, None, None) - return (resolved, stat.st_mtime_ns, stat.st_size) - - def _is_corrupt_board_db_error(exc: Exception) -> bool: - corrupt_guard_error = getattr(_kb, "KanbanDbCorruptError", None) - if corrupt_guard_error is not None and isinstance(exc, corrupt_guard_error): - return True - if not isinstance(exc, sqlite3.DatabaseError): - return False - msg = str(exc).lower() - return ( - "file is not a database" in msg - or "database disk image is malformed" in msg - ) - - def _tick_once_for_board(slug: str) -> "Optional[object]": - """Run one dispatch_once for a specific board (in a worker thread). - - The per-board DB is opened explicitly so boards never share a - connection or claim across each other. - """ - conn = None - fingerprint = _board_db_fingerprint(slug) - disabled_entry = disabled_corrupt_boards.get(slug) - if disabled_entry is not None: - disabled_fingerprint, disabled_at = disabled_entry - age = time.monotonic() - disabled_at - if ( - disabled_fingerprint == fingerprint - and age < CORRUPT_BOARD_RETRY_AFTER_SECONDS - ): - return None - if disabled_fingerprint == fingerprint: - logger.info( - "kanban dispatcher: board %s database fingerprint unchanged " - "after %.0fs quarantine; retrying dispatch", - slug, - age, - ) - else: - logger.info( - "kanban dispatcher: board %s database changed; retrying dispatch", - slug, - ) - disabled_corrupt_boards.pop(slug, None) - try: - # No explicit init_db(): connect() runs the migration once per - # process (see the matching note in _kanban_notifier_watcher). - conn = _kb.connect(board=slug) - return _kb.dispatch_once( - conn, - board=slug, - max_spawn=max_spawn, - max_in_progress=max_in_progress, - failure_limit=failure_limit, - stale_timeout_seconds=stale_timeout_seconds, - default_assignee=default_assignee, - max_in_progress_per_profile=max_in_progress_per_profile, - reconcile_orphans=reconcile_orphans, - ) - except Exception as exc: - if _is_corrupt_board_db_error(exc): - disabled_corrupt_boards[slug] = (fingerprint, time.monotonic()) - logger.error( - "kanban dispatcher: board %s database %s is not a valid " - "SQLite database; pausing dispatch for this board until " - "the file changes, the gateway restarts, or the " - "quarantine timer expires. Move or restore the file, " - "then run `hermes kanban init` if you need a fresh board.", - slug, - fingerprint[0], - ) - return None - logger.exception("kanban dispatcher: tick failed on board %s", slug) - return None - finally: - if conn is not None: - with contextlib.suppress(Exception): - conn.close() - - def _list_boards() -> list: - try: - return _kb.list_boards(include_archived=False) - except Exception: - return [_kb.read_board_metadata(_kb.DEFAULT_BOARD)] - - def _tick_once() -> "list[tuple[str, Optional[object]]]": - """Run one dispatch_once per board. Returns (slug, result) pairs. - - Boards are enumerated every tick so a board created mid-run is - picked up without a restart. - """ - out: list[tuple[str, "Optional[object]"]] = [] - for b in _list_boards(): - slug = b.get("slug") or _kb.DEFAULT_BOARD - out.append((slug, _tick_once_for_board(slug))) - return out - - def _ready_nonempty() -> bool: - """Is there a ready+assigned+unclaimed task on ANY board that the - dispatcher would actually spawn for? - - Control-plane lanes (e.g. ``orion-cc``) are pulled by terminals - via ``claim_task`` and never spawnable — a queue full of those is - "correctly idle", not "stuck". The review column is probed only - when review dispatch is on (same gate as the dispatcher): a task - waiting for a human reviewer is idle, not stuck. - """ - _review_probe = _kb.review_dispatch_enabled() - for b in _list_boards(): - slug = b.get("slug") or _kb.DEFAULT_BOARD - conn = None - try: - conn = _kb.connect(board=slug) - if _kb.has_spawnable_ready(conn): - return True - if _review_probe and _kb.has_spawnable_review(conn): - return True - except Exception: - continue - finally: - if conn is not None: - with contextlib.suppress(Exception): - conn.close() - return False - - def _auto_decompose_tick(auto_decompose_per_tick: int) -> int: - """Auto-decompose up to N triage tasks across all boards into - ready workgraphs before dispatch fans out; the per-tick cap keeps - a bulk triage load from burst-spending the aux LLM. Returns the - number decomposed/specified. - """ - try: - from hermes_cli import kanban_decompose as _decomp - except Exception as exc: # pragma: no cover - logger.warning( - "kanban auto-decompose: import failed (%s); skipping", exc, - ) - return 0 - attempted = 0 - successes = 0 - for b in _list_boards(): - slug = b.get("slug") or _kb.DEFAULT_BOARD - if attempted >= auto_decompose_per_tick: - break - # Pin the board via env for the call: the decomposer connects - # with no board kwarg (same pattern as the dashboard specify endpoint). - prev_env = os.environ.get("HERMES_KANBAN_BOARD") - try: - os.environ["HERMES_KANBAN_BOARD"] = slug - try: - triage_ids = _decomp.list_triage_ids() - except Exception as exc: - logger.debug( - "kanban auto-decompose: list_triage_ids failed on board %s (%s)", - slug, exc, - ) - triage_ids = [] - for tid in triage_ids: - if attempted >= auto_decompose_per_tick: - break - attempted += 1 - try: - outcome = _decomp.decompose_task( - tid, author="auto-decomposer", - ) - except Exception: - logger.exception( - "kanban auto-decompose: decompose_task crashed on %s", - tid, - ) - continue - if outcome.ok: - successes += 1 - if outcome.fanout and outcome.child_ids: - logger.info( - "kanban auto-decompose [%s]: %s → %d children", - slug, tid, len(outcome.child_ids), - ) - else: - logger.info( - "kanban auto-decompose [%s]: %s → single task (no fanout)", - slug, tid, - ) - else: - # Common no-op reasons (no aux client) must not spam logs every tick. - logger.debug( - "kanban auto-decompose [%s]: %s skipped: %s", - slug, tid, outcome.reason, - ) - finally: - if prev_env is None: - os.environ.pop("HERMES_KANBAN_BOARD", None) - else: - os.environ["HERMES_KANBAN_BOARD"] = prev_env - return successes + dispatcher = _KanbanDispatcher(_kb, settings) logger.info( "kanban dispatcher: embedded in gateway (interval=%.1fs)", interval @@ -1367,37 +421,18 @@ class GatewayKanbanWatchersMixin: # Emergency stop (`hermes pause`): no auto-decompose or # dispatch while paused; running workers finish naturally. if not _kanban_dispatch_allowed(): - ready_pending = False bad_ticks = 0 else: # Re-read the auto-decompose toggle live so disabling it # takes effect on the next tick, not on restart. _ad_enabled, _ad_per_tick = _resolve_auto_decompose_settings(_load_config) if _ad_enabled: - await _to_thread_process_service(_auto_decompose_tick, _ad_per_tick) - results = await _to_thread_process_service(_tick_once) - any_spawned = False - for slug, res in (results or []): - if res is not None and getattr(res, "spawned", None): - any_spawned = True - # Quiet by default: an idle gateway stays silent. - logger.info( - "kanban dispatcher [%s]: spawned=%d reclaimed=%d " - "crashed=%d timed_out=%d promoted=%d auto_blocked=%d", - slug, - len(res.spawned), - res.reclaimed, - len(res.crashed) if hasattr(res.crashed, "__len__") else 0, - len(res.timed_out) if hasattr(res.timed_out, "__len__") else 0, - res.promoted, - len(res.auto_blocked) if hasattr(res.auto_blocked, "__len__") else 0, - ) + await _to_thread_process_service(dispatcher.auto_decompose_tick, _ad_per_tick) + results = await _to_thread_process_service(dispatcher.tick_once) + any_spawned = _log_spawn_results(results) # Health telemetry (aggregate across boards) - ready_pending = await _to_thread_process_service(_ready_nonempty) - if ready_pending and not any_spawned: - bad_ticks += 1 - else: - bad_ticks = 0 + ready_pending = await _to_thread_process_service(dispatcher.ready_nonempty) + bad_ticks = bad_ticks + 1 if ready_pending and not any_spawned else 0 if bad_ticks >= HEALTH_WINDOW: now = int(time.time()) if now - last_warn_at >= 300: @@ -1416,10 +451,6 @@ class GatewayKanbanWatchersMixin: except Exception: logger.exception("kanban dispatcher: unexpected watcher error") - # Sleep in 1s slices so stop() never waits a full interval. - slept = 0.0 - while slept < interval and self._running: - await asyncio.sleep(min(1.0, interval - slept)) - slept += 1.0 + await self._sleep_between_ticks(interval) self._release_kanban_dispatcher_lock() diff --git a/gateway/kanban_watchers_common.py b/gateway/kanban_watchers_common.py new file mode 100644 index 0000000000..22833575ee --- /dev/null +++ b/gateway/kanban_watchers_common.py @@ -0,0 +1,35 @@ +"""Plumbing shared by the kanban notifier and dispatcher loops.""" + +from __future__ import annotations + +import asyncio +import logging +from contextvars import Context +from typing import Any, Callable + +# Keep the logger name run.py used so extracted log records are unchanged. +logger = logging.getLogger("gateway.run") + + +def _run_in_fresh_context(func: Callable[..., Any], /, *args: Any) -> Any: + """Run *func* in an empty ``Context`` so request-local ContextVars stay behind. + + ``asyncio.to_thread`` copies the caller's context; a lingering + ``delegate_task`` child marker would make ``write_txn`` false-trip for + these process-owned writers. An empty Context keeps the DB guard intact + for real children without exempting dispatcher writes. + """ + return Context().run(func, *args) + + +async def _to_thread_process_service(func: Callable[..., Any], /, *args: Any) -> Any: + """Offload blocking process-service work without inheriting request ContextVars.""" + return await asyncio.to_thread(_run_in_fresh_context, func, *args) + + +def _list_boards(kb: Any) -> list: + """Enumerate live boards; fall back to the default board when listing fails.""" + try: + return kb.list_boards(include_archived=False) + except Exception: + return [kb.read_board_metadata(kb.DEFAULT_BOARD)] diff --git a/gateway/kanban_watchers_dispatcher.py b/gateway/kanban_watchers_dispatcher.py new file mode 100644 index 0000000000..b768906245 --- /dev/null +++ b/gateway/kanban_watchers_dispatcher.py @@ -0,0 +1,375 @@ +"""Embedded kanban dispatcher: settings resolution and per-tick board work. + +``GatewayKanbanWatchersMixin._kanban_dispatcher_watcher`` owns the loop, +the singleton lock and the health telemetry; everything that only needs the +``kanban_db`` module and the resolved settings lives here. +""" + +from __future__ import annotations + +import contextlib +import os +import sqlite3 +import time +from dataclasses import dataclass +from typing import Any, Optional + +from gateway.kanban_watchers_common import _list_boards, logger + +_CORRUPT_DB_MARKERS = ("file is not a database", "database disk image is malformed") + + +@dataclass +class _DispatcherSettings: + """``kanban.*`` dispatch settings, read once at boot (restart to apply).""" + + interval: float + max_spawn: Any + max_in_progress: Optional[int] + failure_limit: int + stale_timeout_seconds: int + reconcile_orphans: bool + default_assignee: Optional[str] + max_in_progress_per_profile: Optional[int] + + +def _positive_int_setting(kanban_cfg: dict, key: str) -> Optional[int]: + """Parse an optional ``kanban.`` int cap; None when unset or invalid (< 1 is invalid).""" + raw = kanban_cfg.get(key) + if raw is None: + return None + try: + value = int(raw) + except (TypeError, ValueError): + logger.warning("kanban dispatcher: invalid kanban.%s=%r; ignoring", key, raw) + return None + if value < 1: + logger.warning("kanban dispatcher: kanban.%s=%r is below 1; ignoring", key, raw) + return None + logger.info("kanban dispatcher: %s=%d", key, value) + return value + + +def _resolve_dispatcher_settings(kanban_cfg: dict, kb: Any) -> _DispatcherSettings: + """Parse and log the dispatcher settings in their established order.""" + try: + interval = float(kanban_cfg.get("dispatch_interval_seconds", 60) or 60) + except (ValueError, TypeError): + logger.warning( + "kanban dispatcher: invalid dispatch_interval_seconds=%r, using default 60", + kanban_cfg.get("dispatch_interval_seconds"), + ) + interval = 60.0 + interval = max(interval, 1.0) # sanity floor — tighter than this is a footgun + + max_spawn = kanban_cfg.get("max_spawn") + if max_spawn is not None: + logger.info("kanban dispatcher: max_spawn=%s", max_spawn) + + # Cap simultaneously running tasks so slow workers don't pile up and + # time out. Explicit config wins; otherwise a memory-derived default + # (unbounded fan-out swap-thrashes small hosts), or None where total + # memory can't be read. + max_in_progress = _positive_int_setting(kanban_cfg, "max_in_progress") + effective_max_in_progress = kb.resolve_max_in_progress(max_in_progress) + if max_in_progress is None and effective_max_in_progress is not None: + logger.info( + "kanban dispatcher: kanban.max_in_progress unset; using " + "memory-derived default max_in_progress=%d " + "(set kanban.max_in_progress in config.yaml to override)", + effective_max_in_progress, + ) + + raw_failure_limit = kanban_cfg.get("failure_limit", kb.DEFAULT_FAILURE_LIMIT) + try: + failure_limit = int(raw_failure_limit) + except (TypeError, ValueError): + logger.warning( + "kanban dispatcher: invalid kanban.failure_limit=%r; using default %d", + raw_failure_limit, + kb.DEFAULT_FAILURE_LIMIT, + ) + failure_limit = kb.DEFAULT_FAILURE_LIMIT + if failure_limit < 1: + logger.warning( + "kanban dispatcher: kanban.failure_limit=%r is below 1; using default %d", + raw_failure_limit, + kb.DEFAULT_FAILURE_LIMIT, + ) + failure_limit = kb.DEFAULT_FAILURE_LIMIT + + # 0 disables stale detection. + raw_stale = kanban_cfg.get("dispatch_stale_timeout_seconds", 0) + try: + stale_timeout_seconds = int(raw_stale or 0) + except (TypeError, ValueError): + logger.warning( + "kanban dispatcher: invalid kanban.dispatch_stale_timeout_seconds=%r; " + "disabling stale detection", + raw_stale, + ) + stale_timeout_seconds = 0 + + # Fallback profile for tasks created without an assignee (e.g. via the + # dashboard). Empty (the schema default) keeps skipping them. + default_assignee = (kanban_cfg.get("default_assignee") or "").strip() or None + if default_assignee: + logger.info( + "kanban dispatcher: default_assignee=%r (unassigned ready tasks " + "will route to this profile)", + default_assignee, + ) + + return _DispatcherSettings( + interval=interval, + max_spawn=max_spawn, + max_in_progress=effective_max_in_progress, + failure_limit=failure_limit, + stale_timeout_seconds=stale_timeout_seconds, + # Requeue 'running' cards with broken claim bookkeeping (zombie-card + # reconciliation); false keeps orphans frozen for manual forensics. + reconcile_orphans=bool(kanban_cfg.get("reconcile_orphans", True)), + default_assignee=default_assignee, + # Per-profile concurrency cap: no single profile's local model / API + # quota / browser pool gets overwhelmed by a fan-out. + max_in_progress_per_profile=_positive_int_setting(kanban_cfg, "max_in_progress_per_profile"), + ) + + +class _KanbanDispatcher: + """Per-tick board work for the embedded dispatcher (runs in worker threads). + + Boards are enumerated every tick so a board created mid-run is picked up + without a restart. Corrupt-looking board DBs are quarantined per + fingerprint and retried after ``CORRUPT_BOARD_RETRY_AFTER_SECONDS``: + transient WAL/open races can look like "malformed" for one tick. + """ + + CORRUPT_BOARD_RETRY_AFTER_SECONDS = 300 + + def __init__(self, kb: Any, settings: _DispatcherSettings) -> None: + self.kb = kb + self.settings = settings + self.disabled_corrupt_boards: dict[ + str, tuple[tuple[str, int | None, int | None], float] + ] = {} + + def _board_slugs(self) -> list: + return [b.get("slug") or self.kb.DEFAULT_BOARD for b in _list_boards(self.kb)] + + def board_db_fingerprint(self, slug: str) -> tuple[str, int | None, int | None]: + path = self.kb.kanban_db_path(slug) + try: + resolved = str(path.expanduser().resolve()) + except Exception: + resolved = str(path) + try: + stat = path.stat() + except OSError: + return (resolved, None, None) + return (resolved, stat.st_mtime_ns, stat.st_size) + + def is_corrupt_board_db_error(self, exc: Exception) -> bool: + corrupt_guard_error = getattr(self.kb, "KanbanDbCorruptError", None) + if corrupt_guard_error is not None and isinstance(exc, corrupt_guard_error): + return True + if not isinstance(exc, sqlite3.DatabaseError): + return False + msg = str(exc).lower() + return any(marker in msg for marker in _CORRUPT_DB_MARKERS) + + def _quarantine_lifted(self, slug: str, fingerprint: tuple) -> bool: + """Return False while *slug* stays quarantined; lift (and log) otherwise.""" + disabled_entry = self.disabled_corrupt_boards.get(slug) + if disabled_entry is None: + return True + disabled_fingerprint, disabled_at = disabled_entry + age = time.monotonic() - disabled_at + if disabled_fingerprint == fingerprint and age < self.CORRUPT_BOARD_RETRY_AFTER_SECONDS: + return False + if disabled_fingerprint == fingerprint: + logger.info( + "kanban dispatcher: board %s database fingerprint unchanged " + "after %.0fs quarantine; retrying dispatch", + slug, + age, + ) + else: + logger.info( + "kanban dispatcher: board %s database changed; retrying dispatch", + slug, + ) + self.disabled_corrupt_boards.pop(slug, None) + return True + + def tick_once_for_board(self, slug: str) -> Optional[object]: + """Run one dispatch_once for a specific board. + + The per-board DB is opened explicitly so boards never share a + connection or claim across each other. + """ + conn = None + fingerprint = self.board_db_fingerprint(slug) + if not self._quarantine_lifted(slug, fingerprint): + return None + s = self.settings + try: + # No explicit init_db(): connect() runs the migration once per + # process (see the matching note in the notifier collector). + conn = self.kb.connect(board=slug) + return self.kb.dispatch_once( + conn, + board=slug, + max_spawn=s.max_spawn, + max_in_progress=s.max_in_progress, + failure_limit=s.failure_limit, + stale_timeout_seconds=s.stale_timeout_seconds, + default_assignee=s.default_assignee, + max_in_progress_per_profile=s.max_in_progress_per_profile, + reconcile_orphans=s.reconcile_orphans, + ) + except Exception as exc: + if self.is_corrupt_board_db_error(exc): + self.disabled_corrupt_boards[slug] = (fingerprint, time.monotonic()) + logger.error( + "kanban dispatcher: board %s database %s is not a valid " + "SQLite database; pausing dispatch for this board until " + "the file changes, the gateway restarts, or the " + "quarantine timer expires. Move or restore the file, " + "then run `hermes kanban init` if you need a fresh board.", + slug, + fingerprint[0], + ) + return None + logger.exception("kanban dispatcher: tick failed on board %s", slug) + return None + finally: + if conn is not None: + with contextlib.suppress(Exception): + conn.close() + + def tick_once(self) -> list[tuple[str, Optional[object]]]: + """Run one dispatch_once per board. Returns (slug, result) pairs.""" + return [(slug, self.tick_once_for_board(slug)) for slug in self._board_slugs()] + + def ready_nonempty(self) -> bool: + """Is there a ready+assigned+unclaimed task on ANY board that the + dispatcher would actually spawn for? + + Control-plane lanes (e.g. ``orion-cc``) are pulled by terminals + via ``claim_task`` and never spawnable — a queue full of those is + "correctly idle", not "stuck". The review column is probed only + when review dispatch is on (same gate as the dispatcher): a task + waiting for a human reviewer is idle, not stuck. + """ + kb = self.kb + _review_probe = kb.review_dispatch_enabled() + for slug in self._board_slugs(): + conn = None + try: + conn = kb.connect(board=slug) + if kb.has_spawnable_ready(conn): + return True + if _review_probe and kb.has_spawnable_review(conn): + return True + except Exception: + continue + finally: + if conn is not None: + with contextlib.suppress(Exception): + conn.close() + return False + + def auto_decompose_tick(self, auto_decompose_per_tick: int) -> int: + """Auto-decompose up to N triage tasks across all boards into + ready workgraphs before dispatch fans out; the per-tick cap keeps + a bulk triage load from burst-spending the aux LLM. Returns the + number decomposed/specified. + """ + try: + from hermes_cli import kanban_decompose as _decomp + except Exception as exc: # pragma: no cover + logger.warning( + "kanban auto-decompose: import failed (%s); skipping", exc, + ) + return 0 + attempted = 0 + successes = 0 + for slug in self._board_slugs(): + if attempted >= auto_decompose_per_tick: + break + # Pin the board via env for the call: the decomposer connects + # with no board kwarg (same pattern as the dashboard specify endpoint). + prev_env = os.environ.get("HERMES_KANBAN_BOARD") + try: + os.environ["HERMES_KANBAN_BOARD"] = slug + try: + triage_ids = _decomp.list_triage_ids() + except Exception as exc: + logger.debug( + "kanban auto-decompose: list_triage_ids failed on board %s (%s)", + slug, exc, + ) + triage_ids = [] + for tid in triage_ids: + if attempted >= auto_decompose_per_tick: + break + attempted += 1 + successes += self._decompose_one(_decomp, slug, tid) + finally: + if prev_env is None: + os.environ.pop("HERMES_KANBAN_BOARD", None) + else: + os.environ["HERMES_KANBAN_BOARD"] = prev_env + return successes + + @staticmethod + def _decompose_one(_decomp: Any, slug: str, tid: str) -> int: + """Decompose one triage task; returns 1 on success, 0 otherwise.""" + try: + outcome = _decomp.decompose_task(tid, author="auto-decomposer") + except Exception: + logger.exception( + "kanban auto-decompose: decompose_task crashed on %s", + tid, + ) + return 0 + if not outcome.ok: + # Common no-op reasons (no aux client) must not spam logs every tick. + logger.debug( + "kanban auto-decompose [%s]: %s skipped: %s", + slug, tid, outcome.reason, + ) + return 0 + if outcome.fanout and outcome.child_ids: + logger.info( + "kanban auto-decompose [%s]: %s → %d children", + slug, tid, len(outcome.child_ids), + ) + else: + logger.info( + "kanban auto-decompose [%s]: %s → single task (no fanout)", + slug, tid, + ) + return 1 + + +def _log_spawn_results(results: Optional[list]) -> bool: + """Log per-board spawn summaries; returns whether any board spawned.""" + any_spawned = False + for slug, res in (results or []): + if res is not None and getattr(res, "spawned", None): + any_spawned = True + # Quiet by default: an idle gateway stays silent. + logger.info( + "kanban dispatcher [%s]: spawned=%d reclaimed=%d " + "crashed=%d timed_out=%d promoted=%d auto_blocked=%d", + slug, + len(res.spawned), + res.reclaimed, + len(res.crashed) if hasattr(res.crashed, "__len__") else 0, + len(res.timed_out) if hasattr(res.timed_out, "__len__") else 0, + res.promoted, + len(res.auto_blocked) if hasattr(res.auto_blocked, "__len__") else 0, + ) + return any_spawned diff --git a/gateway/kanban_watchers_notifier.py b/gateway/kanban_watchers_notifier.py new file mode 100644 index 0000000000..2abb00890e --- /dev/null +++ b/gateway/kanban_watchers_notifier.py @@ -0,0 +1,722 @@ +"""Kanban notifier: claim terminal task events per subscription and deliver them. + +``GatewayKanbanWatchersMixin._kanban_notifier_watcher`` owns the loop and +the GC cadence; the per-tick claim (``_notifier_collect``) and the +per-subscription delivery (``_KanbanNotification``) live here. +""" + +from __future__ import annotations + +import re +from pathlib import Path +from typing import Any, Callable, Optional + +from agent.i18n import t + +from gateway.kanban_watchers_common import _list_boards, _to_thread_process_service, logger + +# "status" covers dashboard drag-drop and `_set_status_direct()`. +# ``review_requested`` wakes the origin like a block but is not one; +# the task is not archived so later review cycles keep notifying. +TERMINAL_KINDS = ("completed", "blocked", "gave_up", "crashed", "timed_out", "status", "archived", "unblocked", "block_loop_detected", "review_requested", "changes_requested") +# Kinds that hand a decision back to the origin, which must take a turn. +# status/archived/unblocked are bookkeeping. +_WAKE_KINDS = ( + "completed", "gave_up", "crashed", "timed_out", + "blocked", "review_requested", "changes_requested", + "block_loop_detected", +) +# Consecutive send failures (adapter raised OR reported +# SendResult(success=False)) before a sub is dropped as a dead chat. +# 12 ≈ 60s at the 5s cadence: a transient API outage must not +# permanently unsubscribe a live review-gate channel. +MAX_SEND_FAILURES = 12 + +_LOCAL_PATH_RE = re.compile( + r"(? str: + """Return a mobile-friendly review reason safe for external delivery.""" + from agent.redact import redact_sensitive_text + + reason = redact_sensitive_text( + "" if value is None else str(value), + force=True, + redact_url_credentials=True, + ) + reason = _LOCAL_PATH_RE.sub("[local path]", reason) + reason = " ".join(reason.split()) + if len(reason) > limit: + reason = reason[: limit - 1].rstrip() + "…" + return reason + + +def _wake_scope_id(adapter: Any, sub: dict) -> Optional[str]: + """Return the tenant scope (Slack workspace) a subscription's wake keys to. + + ``build_session_key()`` includes ``scope_id`` on multi-tenant platforms, + so the wake must carry the same scope as inbound messages. Persisted + ``delivery_metadata`` wins (it records the creating scope); the adapter's + live chat → scope map only covers rows without metadata. ``None`` means + unscoped, matching an unscoped platform's key. + """ + delivery_meta = sub.get("delivery_metadata") + if isinstance(delivery_meta, dict): + for key in ("scope_id", "slack_team_id", "team_id"): + value = delivery_meta.get(key) + if value: + return str(value) + resolver = getattr(adapter, "scope_id_for_chat", None) + if callable(resolver): + try: + resolved = resolver(str(sub.get("chat_id") or "")) + except Exception as exc: + # An adapter-side lookup failure yields no scope, never an error. + logger.debug( + "kanban notifier: scope lookup failed for chat %s: %s", + sub.get("chat_id"), + exc, + exc_info=True, + ) + return None + if resolved: + return str(resolved) + return None + + +def _platform_names(mapping: Any) -> set[str]: + """Lower-cased platform names of an adapters mapping (Platform enums or strings).""" + return {getattr(platform, "value", str(platform)).lower() for platform in mapping} + + +# --------------------------------------------------------------------------- +# Collection (runs in a worker thread) +# --------------------------------------------------------------------------- + + +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. + + Each gateway polls only subscriptions owned by profiles whose adapters it + 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 _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*.""" + # Cheap read-only probe before the writable connect() (schema init, WAL + # sidecars, checkpoints) — a board with no subscriptions has nothing to notify. + 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 + except Exception as exc: + logger.debug( + "kanban notifier: read-only subscription probe failed " + "for board %s (%s); falling back to writable open", + slug, exc, + ) + 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: + # Best-effort: 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, + ) + # 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: + 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, + ) + continue + 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 "", + ) + continue + 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: + continue + 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, + ) + deliveries.append({ + "sub": sub, + "old_cursor": old_cursor, + "cursor": cursor, + "events": events, + "task": task, + "board": slug, + }) + 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() + + +# --------------------------------------------------------------------------- +# Per-event message formatting: kind -> (msg, wake_handoff, wake_review_detail) +# --------------------------------------------------------------------------- +# ``None`` for handoff / review_detail leaves the accumulated wake value +# untouched. ``_payload(ev, key)`` is the shared "payload present and truthy" read. + + +def _payload(ev: Any, key: str) -> Any: + return ev.payload.get(key) if ev.payload and ev.payload.get(key) else None + + +def _first_line(text: str, limit: int) -> str: + lines = text.strip().splitlines() + return lines[0][:limit] if lines else text[:limit] + + +def _fmt_completed(ev, n) -> tuple: + # Prefer the run summary from the event payload; fall back to + # task.result for legacy rows. + handoff = "" + wake_handoff = None + payload_summary = _payload(ev, "summary") + if payload_summary: + wake_handoff = _first_line(str(payload_summary), 200) + handoff = f"\n{wake_handoff}" + elif n.task and n.task.result: + wake_handoff = _first_line(n.task.result, 160) + handoff = f"\n{wake_handoff}" + msg = ( + f"✔ {n.board_tag}{n.tag}Kanban {n.task_id} done" + f" — {n.title}{handoff}" + ) + return msg, wake_handoff, None + + +def _fmt_blocked(ev, n) -> tuple: + reason = _payload(ev, "reason") + reason = f": {str(reason)[:160]}" if reason else "" + return f"⏸ {n.board_tag}{n.tag}Kanban {n.task_id} blocked{reason}", None, None + + +def _fmt_gave_up(ev, n) -> tuple: + err = _payload(ev, "error") + err = f"\n{str(err)[:200]}" if err else "" + msg = ( + f"✖ {n.board_tag}{n.tag}Kanban {n.task_id} gave up " + f"after repeated spawn failures{err}" + ) + return msg, None, None + + +def _fmt_crashed(ev, n) -> tuple: + msg = ( + f"✖ {n.board_tag}{n.tag}Kanban {n.task_id} worker crashed " + f"(pid gone); dispatcher will retry" + ) + return msg, None, None + + +def _fmt_timed_out(ev, n) -> tuple: + limit = _payload(ev, "limit_seconds") + limit = int(limit) if limit else 0 + msg = ( + f"⏱ {n.board_tag}{n.tag}Kanban {n.task_id} timed out " + f"(max_runtime={limit}s); will retry" + ) + return msg, None, None + + +def _fmt_status(ev, n) -> tuple: + new_status = _payload(ev, "status") + new_status = str(new_status) if new_status else "" + return f"🔄 {n.board_tag}{n.tag}Kanban {n.task_id} → {new_status}", 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. + handoff = "" + wake_handoff = None + summary = _payload(ev, "summary") + if summary: + summary = str(summary) + handoff = f"\n{summary[:200]}" + wake_handoff = _first_line(summary, 200) + msg = ( + f"👀 {n.board_tag}{n.tag}Kanban {n.task_id} ready for review" + f" — {n.title}{handoff}" + ) + return msg, wake_handoff, None + + +def _fmt_changes_requested(ev, n) -> tuple: + payload = ev.payload or {} + reason = _safe_review_reason(payload.get("reason")) + reviewer = _safe_review_reason(payload.get("reviewer"), 48) + implementer = _safe_review_reason(payload.get("implementer"), 48) + reason_text = reason or "reviewer feedback requires changes" + provenance = "" + if reviewer: + provenance += f" — reviewer @{reviewer}" + if implementer: + provenance += f" → implementer @{implementer}" + msg = ( + f"🛑 {n.board_tag}Kanban {n.task_id} review requested " + f"changes/BLOCK: {reason_text}{provenance}" + ) + 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. + reason = _payload(ev, "reason") + reason = f": {str(reason)[:160]}" if reason else "" + recurrences = ev.payload.get("recurrences") if ev.payload else None + rc = f" (blocked {recurrences}x for the same cause)" if recurrences else "" + msg = ( + f"🛑 {n.board_tag}{n.tag}Kanban {n.task_id} routed to TRIAGE" + f" — needs a human decision{rc}{reason}" + ) + return msg, 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, + "review_requested": _fmt_review_requested, + "changes_requested": _fmt_changes_requested, + "block_loop_detected": _fmt_block_loop_detected, +} + + +# --------------------------------------------------------------------------- +# Delivery of one claimed batch (one subscription, N events) +# --------------------------------------------------------------------------- + + +class _KanbanNotification: + """Deliver one subscription's claimed events, then settle the cursor. + + Cursor advance ordering by adapter class: + * push + notify: the text send WAS the delivery → advance now; wake + injection stays best-effort. + * non-push or wake-only: the wake IS the delivery → it runs FIRST and the + cursor advances only after it succeeds; failure rewinds like a failed + send(). An unknown platform advances the cursor so it can't replay forever. + """ + + def __init__(self, runner: Any, d: dict, *, platform_cls: Any, sub_fail_counts: dict) -> None: + self.runner = runner + self.d = d + self.platform_cls = platform_cls + self.sub_fail_counts = sub_fail_counts + self.sub = sub = d["sub"] + self.task = task = d["task"] + self.board_slug = d.get("board") + self.platform_str = (sub["platform"] or "").lower() + self.task_id = sub["task_id"] + self.sub_profile = sub.get("notifier_profile") or "" + self.title = (task.title if task else sub["task_id"])[:120] + self.board_tag = f"[{self.board_slug}] " if self.board_slug else "" + # Attribute the ping to the worker that did the work. + who = task.assignee if task and task.assignee else None + self.tag = f"@{who} " if who else "" + # The wake self-post path needs the key even when every event was skipped. + self.sub_key = (sub["task_id"], sub["platform"], sub["chat_id"], sub.get("thread_id") or "") + mode = sub.get("delivery_mode") or "notify" + self.wake_agent = mode in ("notify+wake", "wake") + self.send_passive = mode != "wake" + # Worker handoff carried into the synthetic wake turn so the woken + # creator doesn't re-decompose work already on the board. + self.wake_handoff = "" + self.wake_review_detail = "" + self.plat: Any = None + self.adapter: Any = None + self.is_push_adapter = True + self.wake_kinds: set = set() + self.session_key = "" + self.synth = "" + + # -- cursor / subscription ops (blocking, run in a fresh-context thread) -- + + async def rewind(self) -> None: + await _to_thread_process_service( + self.runner._kanban_rewind, self.sub, self.d["cursor"], self.d.get("old_cursor", 0), self.board_slug, + ) + + async def advance(self) -> None: + await _to_thread_process_service(self.runner._kanban_advance, self.sub, self.d["cursor"], self.board_slug) + + async def unsub(self) -> None: + await _to_thread_process_service(self.runner._kanban_unsub, self.sub, self.board_slug) + + def clear_failures(self) -> None: + self.sub_fail_counts.pop(self.sub_key, None) + + async def delivery_failed(self, fmt: str, prefix: tuple, drop_fmt: str, exc: Exception, exc_info: bool) -> None: + """Bump the failure counter; drop the sub past the limit, else rewind the claim so the next tick retries.""" + fails = self.sub_fail_counts.get(self.sub_key, 0) + 1 + self.sub_fail_counts[self.sub_key] = fails + logger.warning(fmt, *prefix, fails, MAX_SEND_FAILURES, exc, exc_info=exc_info) + if fails >= MAX_SEND_FAILURES: + logger.warning(drop_fmt, self.task_id, self.platform_str, fails) + await self.unsub() + self.clear_failures() + else: + await self.rewind() + + async def _wake_failed(self, fmt: str, exc: Exception) -> None: + await self.delivery_failed( + fmt, (self.task_id,), + "kanban notifier: dropping subscription %s on %s after %d consecutive wake failures", + exc, True, + ) + + # -- formatting -- + + def format_event(self, ev: Any) -> Optional[str]: + """Render one event; accumulates wake handoff/review detail. None → silent kind.""" + formatter = _EVENT_FORMATTERS.get(ev.kind) + if formatter is None: + return None + msg, handoff, review_detail = formatter(ev, self) + if handoff is not None: + self.wake_handoff = handoff + if review_detail is not None: + self.wake_review_detail = review_detail + return msg + + def build_wake_text(self) -> None: + """Set ``wake_kinds`` / ``session_key`` / ``synth`` for the wake paths.""" + task, sub = self.task, self.sub + self.wake_kinds = ( + {ev.kind for ev in self.d["events"] if ev.kind in _WAKE_KINDS} + if self.wake_agent + else set() + ) + if not self.wake_kinds: + return + if self.is_push_adapter: + self.session_key = getattr(task, "session_id", None) or "" + else: + # Non-push wakes target sub["chat_id"] (the raw session id the + # subscriber registered). task.session_id may be a WORKER session + # for child tasks; use it only for legacy rows. + self.session_key = sub["chat_id"] or getattr(task, "session_id", None) or "" + # i18n keys: gateway.kanban.wake. for each _WAKE_KINDS entry. + _parts = [t(f"gateway.kanban.wake.{k}") for k in _WAKE_KINDS if k in self.wake_kinds] + _status = t("gateway.kanban.wake.status_joiner").join(_parts) or t("gateway.kanban.wake.status_default") + synth = t( + "gateway.kanban.wake.message", + task_id=sub["task_id"], + status=_status, + title=self.title, + assignee=task.assignee if task else "", + board=self.board_slug, + ) + # Label as an automatic notification and carry the handoff so the + # creator inspects the board instead of re-decomposing. + if self.wake_handoff: + synth += "\n" + t("gateway.kanban.wake.handoff", summary=self.wake_handoff) + if self.wake_review_detail: + synth += "\n" + t("gateway.kanban.wake.review_detail", reason=self.wake_review_detail) + self.synth = synth + "\n\n" + t("gateway.kanban.wake.guidance") + + def _log_woke(self) -> None: + logger.info( + "kanban notifier: woke agent for %s on %s/%s profile=%s events=%s", + 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 + from gateway.wake import deliver_wake + sub = self.sub + # 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. Legacy rows may carry chat_type in + # delivery_metadata; last resort is "group". A mismatch only degrades + # to a fresh session. + _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() + _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), + ) + 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: + """Send one text ping; raises on adapter exception or SendResult(success=False).""" + sub, adapter = self.sub, self.adapter + delivery_metadata = sub.get("delivery_metadata") + metadata: dict[str, Any] = dict(delivery_metadata) if isinstance(delivery_metadata, dict) else {} + if sub.get("thread_id") and not metadata.get("thread_id"): + metadata["thread_id"] = sub["thread_id"] + _send_res = await adapter.send(sub["chat_id"], msg, metadata=metadata) + # SendResult(success=False) without an exception is a FAILED delivery + # (else the event is lost); None / non-SendResult keeps the + # "no exception == delivered" contract. + if getattr(_send_res, "success", True) is False: + raise RuntimeError( + "adapter send() reported failure: " + f"{getattr(_send_res, 'error', None) or 'unknown error'}" + ) + logger.debug( + "kanban notifier: delivered %s event for %s to %s/%s on board %s", + ev.kind, self.task_id, self.platform_str, sub["chat_id"], self.board_slug, + ) + # Upload artifact paths from the completion payload / legacy result as + # native files. Only on ``completed`` so retries never spam attachments. + if ev.kind == "completed": + try: + await self.runner._deliver_kanban_artifacts( + adapter=adapter, + chat_id=sub["chat_id"], + metadata=metadata, + event_payload=getattr(ev, "payload", None), + task=self.task, + ) + except Exception as art_exc: + logger.debug( + "kanban notifier: artifact delivery for %s failed: %s", + self.task_id, art_exc, + ) + + async def _send_pings(self) -> bool: + """Send every text ping; False when a send failed (claim already rewound/dropped).""" + for ev in self.d["events"]: + msg = self.format_event(ev) + if msg is None: + continue + # Non-push adapters (api_server) always report SendResult(success=False) + # from send(); treating that as failure would drop the sub forever and + # make the wake path unreachable. Skip the doomed send; the self-post + # IS the delivery and resolves the failure counter. + if not self.is_push_adapter and self.wake_agent: + logger.debug( + "kanban notifier: adapter %s has no push " + "channel; skipping text ping for %s, relying " + "on wake self-post instead", + self.platform_str, self.task_id, + ) + continue + if not self.send_passive: + # Wake-only: the wake path is the sole delivery and resolves the counter. + continue + try: + await self._send_event(ev, msg) + self.clear_failures() + except Exception as exc: + await self.delivery_failed( + "kanban notifier: send failed for %s on %s (attempt %d/%d): %s", + (self.task_id, self.platform_str), + "kanban notifier: dropping subscription %s on %s after %d consecutive send failures", + exc, False, + ) + return False + return True + + async def deliver(self) -> None: + try: + self.plat = self.platform_cls(self.platform_str) + 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) + if adapter is None: + logger.debug( + "kanban notifier: adapter %s disconnected before delivery for %s; rewinding claim", + self.platform_str, self.task_id, + ) + await self.rewind() + return + self.adapter = adapter + from gateway.wake import adapter_supports_push + self.is_push_adapter = adapter_supports_push(adapter) + + if not await self._send_pings(): + return + # All text pings delivered (or skipped for non-push / wake-only). + task_terminal = self.task and self.task.status == "archived" + 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. + from gateway.wake import deliver_wake + try: + await deliver_wake(adapter, text=self.synth, session_id=self.session_key) + self._log_woke() + 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) + return + + # Delivery complete: advance the cursor (the dedup mechanism). + await self.advance() + if not is_push: + self.clear_failures() + if is_push and self.send_passive and wake_kinds: + # notify+wake: text ping was the delivery and the cursor has + # advanced; the wake stays best-effort, but log at WARNING so a + # persistently failing wake is visible. + try: + await self.push_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. + if task_terminal: + await self.unsub()