Files
hermes-agent/tui_gateway/session_notifications.py
T
shannonsands a6ee31f55a feat(wisdom): add Hermes Collective Wisdom Agent V1 (#94266)
* feat(wisdom): add trusted publish and install foundation

* feat(wisdom): add private contribution loop

* feat(wisdom): add managed consumption workflows

* fix(wisdom): close cross-repository safety gaps

* fix(wisdom): align local package and lifecycle policy

* fix(wisdom): require explicit profile setup

* docs(wisdom): repin reconciled gateway head

* fix(wisdom): fence content downloads and approval receipts

* docs(wisdom): record generation-fenced downloads

* docs(wisdom): record unified delivery PR

* fix(ci): stop passing invalid classifier inputs

* docs(wisdom): remove internal requirements ledger

* feat(wisdom): localize dashboard and desktop copy

* feat(wisdom): complete local contribution and consumption UX

* style(wisdom): satisfy desktop lint

* chore(wisdom): refresh requirements pin

* test(dashboard): allow formatted profile copy

* test(wisdom): stabilize desktop interaction coverage

* fix(wisdom): surface dashboard action failures

* fix(wisdom): add repeatable Portal demo login

* feat(wisdom): add actionable skill notifications

* feat(wisdom): add notification install and update actions

* fix(wisdom): make Telegram skill alerts actionable

* fix(wisdom): always refresh demo Agent login

* feat(wisdom): embed Telegram notification actions

* fix(wisdom): preserve Telegram notifications after actions

* fix(wisdom): keep Telegram notification cards readable

* feat(wisdom): add Telegram candidate approval flow

* feat(wisdom): explain Telegram qualification reasons

* fix(wisdom): reconcile cross-surface candidate actions

* feat(telegram): add Collective Wisdom management command

* chore(wisdom): refresh Gateway contract pin

* chore(wisdom): advance Gateway contract pin

* feat(wisdom): align command UX across clients

* feat(slack): add Collective Wisdom management parity

* feat(wisdom): add security and professionalism reviews

* feat(wisdom): add first-time qualification guidance

* feat(wisdom): simplify qualification sharing choices

* feat(skills): add optional editorial metadata

* feat(wisdom): enrich legacy skill presentation

* fix(wisdom): harden review and update boundaries

* fix(wisdom): emit canonical review timestamps

* fix(wisdom): align with merged gateway and main

* wisdom: add agent-led sharing core (policy, evidence, schemas, templates, delivery, weekly job, share/install flows)

- hermes_wisdom/agent_led/: policy resolution (server > local > defaults),
  7-day evidence builder that excludes bundled/hub/managed skills and
  dismissed/handled/recently-suggested content hashes, strict pydantic
  schemas for agent output with repair-or-reject, fixed copy templates
  (Share / Teammate / Published / Update / Mute), idempotent retried
  delivery ledger with stale-action resolution, weekly review job,
  resumable Share and Install flows.
- prompts/: candidate review, recipient recommendation, share packaging.
- tests/wisdom/test_agent_led.py: 30 tests.

* wisdom: agent-led renderers and button action dispatcher

- render.py: Telegram HTML, Slack blocks, Desktop payload; editorial name
  is the emphasized line, product label stays separate.
- actions.py: resolve opaque wa:<action>:<dedup> targets via the delivery
  ledger; Not now -> dismissal, Mute -> fixed options, Share -> resumable
  packaging flow, Install/Update -> plan command. Never publishes/installs.

* wisdom: CLI verbs, agent_led config default, conversational catalog skill

- hermes wisdom browse/review-week/act/share/dismiss/mute (all --json).
- wisdom.agent_led config block, default enabled.
- SKILL.md rewritten so natural-language catalog questions map to the CLI
  verbs, share/install flows and fixed notification templates.

* wisdom: wire agent-led weekly review into gateway tick and Telegram buttons

- gateway housekeeping tick calls maybe_run_weekly_review with a home
  channel sender when a Telegram adapter is available.
- Telegram: wa: callbacks resolved through the ledger (stale-safe), mute
  duration keyboard, send_wisdom_agent_recommendation rich card + fallback.

* fix(wisdom): integrate local mediation and harden model and setup boundaries

* fix(wisdom): honor authoritative recommendation policy and defer on failure

* fix(wisdom): synchronize opaque suppression and recheck delivery preferences

* feat(wisdom): route weekly selection through the session-owned assessment queue

* fix(wisdom): prepare and submit the reviewed generated share package

* feat(wisdom): separate native Share preparation from publication consent

* feat(wisdom): sync native mute choices through a leased preference outbox

* feat(wisdom): bind native mute controls to durable preference choices

* feat(wisdom): add scoped desktop and dashboard notification settings

* fix(wisdom): revalidate feed recommendations before assessment and delivery

* fix(wisdom): persist validated delivery receipts before completing notices

* feat(wisdom): add private notification claim and receipt client

* Persist Wisdom send reservations and recover delivery acknowledgements

* Route legacy Wisdom controls through current native review

* Add typed private Wisdom operation outcome client

* fix(wisdom): make agent-led advice usable in the local demo

* fix(wisdom): keep requested consent outside proactive limits

* fix(wisdom): distinguish unavailable assessments and preserve digest text

* fix(wisdom): assess ongoing usefulness beyond the current task

* fix(wisdom): restore immediate qualification sharing controls

* fix(wisdom): separate qualification review from installation advice

* fix(wisdom): collapse review checklists and simplify sharing copy

* fix(wisdom): show compact sharing progress and publication receipts

* fix(wisdom): require credential prefixes rather than matching skill names

* fix(wisdom): finish package checks before presenting sharing consent

* fix(wisdom): scan local skills before qualification cards

* fix(wisdom): update moderation results on existing sharing cards

* fix(wisdom): keep sharing review accessible from receipt cards

* fix(wisdom): align mediated review cards and collapsible checks

* fix(wisdom): clarify clean security summary wording

* fix(wisdom): normalize consent plans and add explicit recheck

* fix(wisdom): keep install and update receipts concise

* fix(wisdom): collapse assessments and deduplicate operation cards

* fix(wisdom): restore private Portal review from native cards

* fix(wisdom): sync Portal publication to original consent card

* fix(wisdom): show local skill version on sharing cards

* fix(wisdom): skip agent recommendations for self-published versions

* fix(wisdom): simplify candidate notices and local-edit recovery copy

* feat(wisdom): submit locally reviewed packages with one confirmation

* feat(wisdom): expose safe receipt and outcome sync recovery

* wisdom: onboarding notice says detect and share, names the user's own skill

Copy review from the product owner on the first and returning
qualification notices (fixed delivery mode):
- the feature blurb now says the org enabled detection *and sharing*
- both notices say the detected skill is one the user created
- both close with an exclamation mark

Applied identically to hermes_wisdom.notice, the desktop and web i18n
strings, and the tests that assert the sentences.

* wisdom: one opener, no approval line, ask to share after the skill is shown

Product owner review of the candidate card.

- The Hermes written card now opens with the same sentence as the fixed card
  ("Your organisation has enabled Collective Wisdom, a feature designed to
  automatically detect and share useful skills across all team members.")
  instead of its own blurb, so there is one first time message.
- "Nothing is shared without your approval." removed from Telegram, Slack
  and Desktop. The buttons already make the permission explicit.
- "Would you like to share?" no longer appears before the skill is named.
  It is now the last line, after the skill name, description, why suggested
  and the checks, and reads "Would you like to share it?" (matching the
  agent led template wording).

Tests updated for the new order; proposalNotice removed from all desktop locales.

* wisdom: American spelling, organization

Product owner decision: user facing copy uses American spelling.
Changes "Your organisation" to "Your organization" in the chat notice,
the Hermes written card opener, the desktop and web strings, and the
tests that assert them. Identifiers such as nas_organisation:* and the
German and French locales are untouched.

* wisdom: candidate card copy round 4 (owner review)

Apply the product owner's round 4 copy decisions to the Hermes Collective
Wisdom candidate card on Telegram, Slack, Desktop and the shared views:

1. Hermes-written cards are titled "Hermes Collective Wisdom" instead of
   the bare "Collective Wisdom".
2. The "Reusable skill ready to review" line is gone from the candidate
   card (Telegram rich card and plain fallback, legacy agent-led share
   template).
3. The skill name and description are labelled: "Skill name: <name>" and
   "What it does: <description>" (Telegram, Slack, Desktop).
4. "Why suggested:" is now "Why others might benefit:".
5. A passing professionalism review reads "Safe to share at work ✓ (no
   inappropriate content found)" with no per-check bullets and no "Pass";
   a failed review reads "Needs a look before sharing at work (possible
   inappropriate content)" and lists only the checks that flagged
   something. Pending/unavailable wording is unchanged.
6. Telegram button toasts: "Will ask later...", "Preparing more
   details...", "Sharing...".
7. Qualification reasons: "You used this skill consistently across many
   days." and "You've really refined this skill."
8. prompts/wisdom_candidate_review.md asks for a compelling
   editorial_name, a simple one_line_description and a compelling
   why_coworkers_benefit under 300 characters; "Be concise and
   convincing." becomes "Be concise and compelling: the goal is that the
   user wants to share it."

Tests updated for the new strings; review_text() gains direct coverage.

* wisdom: re-apply owner copy after rebase

- Native share cards (advice_view/interaction_view): drop the approval line, ask "Would you like to share it?" as the last line after the checks
- Hermes-written completion card titled "Hermes Collective Wisdom"
- Qualification reasons use the owner wording (consistently across many days / really refined)
- American spelling (organization) in remaining English copy
- Desktop test asserts the current Share button; web test matches the returning notice

* fix(wisdom): pin reconciled Gateway and verify Unicode hash vectors

Pin Gateway 60cd2d6b613ae3cd4a6e65155d1142006d907e78 and byte-identical producer artifacts. Verify every content-order case and package-manifest binding. Validation: 186 focused Python tests, Ruff and contract verifier.

* fix(wisdom): reconcile optional SDK tests and frontend lint

* fix(wisdom): default to agent-written notification summaries

* fix(wisdom): restore deferred install review and browse controls

* feat(wisdom): inspect installed setup with exact package provenance

* feat(wisdom): run native-approved installed setup steps with durable evidence

* fix(wisdom): recover interrupted setup with explicit native consent

* feat(wisdom): hand native installs into guided setup review

* fix(wisdom): continue requested setup with fixed notification copy

* fix(wisdom): preserve setup while waiting for a session model

* fix(wisdom): expose canonical setup review controls on desktop

* fix(wisdom): resume setup after recorded automatic updates

* fix(wisdom): make missing setup prerequisites recheckable

* chore(wisdom): align Agent with verified Gateway contract

* fix(wisdom): stop guessing team slugs in portal links

* fix(wisdom): retire pending advice on account sign-out

* fix(wisdom): cancel advice after terminal account revocation

* fix(wisdom): fence feed responses across account sign-out

* fix(wisdom): checkpoint signed-out feed before reactivation

* fix(wisdom): link proactive advice to scoped notification settings

* fix(wisdom): coalesce queued publication recommendations by version

* fix(wisdom): keep package review navigation local and deferable

* fix(wisdom): reflect installed state in discovery controls

* fix(wisdom): show exact checks before command confirmation

* chore(wisdom): pin bounded analytics privacy contract

* chore(wisdom): pin retired legacy notification contract

* feat(wisdom): review publisher usage with exact sharing copy

* fix(wisdom): align discovery and review check summaries

* fix(wisdom): show expired consent before confirmation

* fix(wisdom): require fresh review for legacy install controls

* fix(wisdom): preserve review expiry across check toggles

* fix(wisdom): retain update policy in native install reviews

* fix(wisdom): surface failed native card edits

* fix(wisdom): persist local command approval reviews

* fix(wisdom): use saved approvals for messaging commands

* test(wisdom): provide scan result in setup handoff fixture

* test(wisdom): exercise Telegram approvals with saved review state

* fix(wisdom): retain suppression policy for offline deferral

* fix(wisdom): reconsider candidates after deferred suppression expires

* fix(wisdom): bind review checks and report verified readiness separately

* fix(wisdom): persist accepted publication intent and recover exact outcomes

* fix(sync): pin UTF-8 tree ordering across writers

* chore(wisdom): pin organisation-scoped Gateway authorization

* fix(wisdom): restrict consent delivery to user-facing sessions

* chore(wisdom): refresh reviewed Gateway contract pin

* fix(wisdom): preserve kept tools in Blank Slate exclusions

* test(auth): reset anonymous fixture with a profile-scoped cache

* fix(wisdom): gate local surfaces and work on current profile entitlement

* fix(wisdom): invalidate quiet tool cache on entitlement changes

* test(wisdom): authorize local consent gateway fixtures

* fix(wisdom): keep entitlement decoding free of native crypto imports

* test(wisdom): provide local entitlement to demo CLI subprocess

* ci: leave upstream workflow unchanged in Wisdom PR

* fix(wisdom): ship package and contracts in Nix wheels

---------

Co-authored-by: hbizi <36184542+hbizi@users.noreply.github.com>
2026-09-11 19:04:06 +10:00

688 lines
34 KiB
Python

"""Per-session notification poller: kanban/loop/delegation events routed to the owning session,
desktop UI wiring, HUD surface note. Bodies are rebound onto server.py's globals at install time
(method_ctx.bind_module), so they reference server.py globals bare."""
from __future__ import annotations
import contextlib
from .method_ctx import bind_module
def _notif_locked_sessions(fn, default):
"""Run ``fn(_sessions)`` under ``_sessions_lock``; ``default`` on failure (poller must never crash)."""
try:
with _sessions_lock:
return fn(_sessions)
except Exception:
return default
def _notif_current_keys(sid: str, session: dict) -> set:
return {str(session.get("session_key") or ""), _session_lookup_key(session, fallback=sid)}
def _notif_session_matches(s: dict, keys) -> bool:
return str(s.get("session_key") or "") in keys or _session_lookup_key(s, fallback="") in keys
def _notif_live_session_matches(keys, exclude: dict | None = None) -> bool:
"""Any non-finalized live session (other than ``exclude``) matches ``keys``; False if the registry can't be read
(fail open rather than drop the event)."""
return _notif_locked_sessions(
lambda ss: any(s is not exclude and not s.get("_finalized") and _notif_session_matches(s, keys)
for s in ss.values()),
False)
def _notif_resolve_event_key(evt_key: str) -> str:
"""Resolve a compression-rotated session key to its continuation tip (or itself)."""
try:
db = _get_db()
return (db.resolve_resume_session_id(evt_key) if db is not None else evt_key) or evt_key
except Exception:
return evt_key
def _notification_event_belongs_elsewhere(sid: str, session: dict, evt: dict) -> bool:
"""True if ``evt`` is owned by a *different* live session. Background completions carry the ``session_key`` of the
session that started the work; async delegation completions also carry ``origin_ui_session_id`` (the live TUI tab)."""
evt_ui_sid = str(evt.get("origin_ui_session_id") or "")
if evt_ui_sid:
if evt_ui_sid == str(sid or "") and not session.get("_finalized"):
return False
if _notif_locked_sessions(lambda ss: evt_ui_sid in ss and not ss[evt_ui_sid].get("_finalized"), False):
return True
# Exact UI tab gone: fall through to durable session_key routing so a resumed continuation with the same
# key/lineage can still claim it.
evt_key = str(evt.get("session_key") or "")
if not evt_key:
return False
current_keys = _notif_current_keys(sid, session)
# Compression can rotate AIAgent.session_id while the detached child is still running: map the event's original
# key to its continuation tip so it reaches the live session instead of becoming an orphan any poller may consume.
# A live continuation wins over the compressed parent, else a stale parent tab could consume the event first.
resolved_key = _notif_resolve_event_key(evt_key)
if resolved_key != evt_key:
if resolved_key in current_keys:
return False
if _notif_live_session_matches({resolved_key}):
return True
if evt_key in current_keys:
return False
return _notif_live_session_matches({evt_key, resolved_key}, exclude=session)
def _session_owns_notification_event(sid: str, session: dict, evt: dict) -> bool:
"""True iff *this* session PROVABLY owns ``evt`` (UI origin is this live session, or ``session_key`` raw/compression-
resolved matches) — the fail-closed gate for addressed notifications, without the orphan-adoption fallback."""
if session.get("_finalized"):
return False
if str(evt.get("origin_ui_session_id") or "") == str(sid or ""):
return True
evt_key = str(evt.get("session_key") or "")
current_keys = _notif_current_keys(sid, session)
return bool(evt_key) and (evt_key in current_keys or _notif_resolve_event_key(evt_key) in current_keys)
def _notification_event_requires_owner(evt: dict) -> bool:
"""Whether ``evt`` must be positively claimed before TUI delivery."""
return evt.get("type") == "async_delegation" or bool(evt.get("origin_ui_session_id") or evt.get("session_key"))
# Extra dedup fields per event type. Completions are terminal (one-shot per process session); watch events are not —
# one process can match patterns many times, so their content is part of the key.
_DEDUP_EXTRA_FIELDS = {
"watch_match": ("command", "pattern", "output", "suppressed", "message_id"),
"watch_disabled": ("command", "message", "suppressed"),
"watch_overflow_": ("command", "message", "suppressed"), # prefix match
}
def _notification_event_dedup_key(evt: dict) -> tuple:
"""UI-emission identity for a process notification event."""
evt_type = evt.get("type", "completion")
if evt_type == "async_delegation":
# No process session_id: else every completion keys as ("", "async_delegation") and the second is suppressed forever.
# An early per-task failure notice must not collapse with the batch's final result (nor with a sibling's notice).
if evt.get("task_failure_notice"):
task_idx = ((evt.get("results") or [{}])[0] or {}).get("task_index", "")
return (evt.get("delegation_id", ""), evt_type, "task_failure", task_idx)
return (evt.get("delegation_id", ""), evt_type)
extra = _DEDUP_EXTRA_FIELDS.get("watch_overflow_" if evt_type.startswith("watch_overflow_") else evt_type, ())
return (evt.get("session_id", ""), evt_type, *(evt.get(f, 0 if f == "suppressed" else "") for f in extra))
# Mirror gateway/kanban_watchers.py TERMINAL_KINDS: claim silent kinds (archived/unblocked) too so the cursor advances
# past them and they can't wedge a later completed/blocked event behind an unclaimed row.
_KANBAN_NOTIFY_KINDS = ("completed", "blocked", "gave_up", "crashed", "timed_out", "status", "archived", "unblocked")
_KANBAN_POLL_SECONDS = _LOOP_POLL_SECONDS = 5.0 # /loop and /heartbeat share one idle-poll cadence
def _notif_release_turn(session: dict) -> None:
with session["history_lock"]:
session["running"] = False
def _notif_claim_turn(session: dict) -> bool:
"""Claim the idle session (running=True) under history_lock; False if a turn is live."""
with session["history_lock"]:
claimed = not session.get("running")
session["running"] = True
return claimed
def _notif_log_failure(what: str, exc: BaseException) -> None:
print(f"[tui_gateway] {what}: {type(exc).__name__}: {exc}", file=sys.stderr)
def _notif_submit(rid: str, sid: str, session: dict, text: str, what: str, **kwargs) -> None:
"""message.start + _run_prompt_submit for a claimed (running=True) turn; releases on failure."""
try:
_emit("message.start", sid)
_run_prompt_submit(rid, sid, session, text, **kwargs)
except Exception as exc:
_notif_log_failure(what, exc)
_notif_release_turn(session)
raise
def _notif_loop_status(sid: str, text: str) -> None:
_emit("status.update", sid, {"kind": "loop", "text": text})
def _notif_slash_loop_tick(rid: str, sid: str, session: dict, mgr, wakeup: str) -> None:
"""Slash-command /loop wakeup: route through the slash pipeline, not the model. No model reply to evaluate, so the
tick completes immediately — unless the command resolves to a prompt (skill command etc.), which runs as a normal
turn whose post-turn hook completes the tick."""
_notif_release_turn(session)
try:
parts = wakeup.lstrip()[1:].split(None, 1)
resp = _methods["command.dispatch"](
rid, {"name": parts[0] if parts else "", "arg": parts[1] if len(parts) > 1 else "", "session_id": sid})
payload = (resp or {}).get("result") or {}
if out := str(payload.get("output") or "").strip():
_notif_loop_status(sid, out)
if payload.get("type") == "send" and payload.get("message"):
if not _notif_claim_turn(session):
mgr.abandon_tick()
return
_emit("message.start", sid)
_run_prompt_submit(rid, sid, session, payload["message"])
return
except Exception:
pass
decision = mgr.complete_tick("")
if decision.get("message"):
_notif_loop_status(sid, decision["message"])
def _maybe_fire_tui_heartbeat_tick(sid: str, session: dict) -> None:
"""Fire a due /heartbeat prompt for an idle TUI/Desktop/dashboard session (#102056, #103044).
``/heartbeat`` runs in the slash worker, whose CLI watchdog queues the due prompt into a
``_pending_input`` no turn loop ever drains — armed-but-dead. State is durable in SessionDB, so
the session-owner process drives firing exactly like ``_maybe_fire_tui_loop_tick``: claim the
idle session first (a racing user prompt wins), then re-enter through ``_run_prompt_submit`` as a
plain user turn. A dispatch that never starts a turn rewinds the persisted fire so the tick stays
due instead of being silently consumed.
"""
try:
from hermes_cli.heartbeat import HeartbeatManager
except Exception:
return
if not (sid_key := session.get("session_key") or ""):
return
mgr = HeartbeatManager(session_id=sid_key)
if not mgr.is_active() or not mgr.state.is_due() or not _notif_claim_turn(session):
return # not due, or busy — the tick coalesces to the next idle poll
if not (prompt := mgr.due_prompt()):
_notif_release_turn(session)
return
started = False
try:
_emit("status.update", sid, {"kind": "heartbeat", "text": f"♥ heartbeat #{mgr.state.fire_count} firing…"})
started = bool(_run_prompt_submit(f"__heartbeat__{int(time.time() * 1000)}", sid, session, prompt))
except Exception as exc:
_notif_log_failure("heartbeat dispatch failed", exc)
if not started:
# _run_prompt_submit releases ``running`` itself when it refuses the turn; make it unconditional.
_notif_release_turn(session)
with contextlib.suppress(Exception):
mgr.abandon_fire()
def _maybe_fire_tui_loop_tick(sid: str, session: dict) -> None:
"""Fire a due /loop wakeup for an idle TUI/Desktop/dashboard session (per-session poller, coarse cadence). Claims
the session (running=True) before dispatching so a racing user prompt wins; the post-turn hook completes the tick."""
try:
from hermes_cli.loops import LoopManager, goal_blocks_loop_tick
except Exception:
return
if not (sid_key := session.get("session_key") or ""):
return
mgr = LoopManager(session_id=sid_key)
if not mgr.is_due() or goal_blocks_loop_tick(sid_key) or not _notif_claim_turn(session):
return # busy — stays due, next poll retries
if not (wakeup := mgr.fire_tick()):
_notif_release_turn(session)
return
rid = f"__loop__{int(time.time() * 1000)}"
try:
_notif_loop_status(sid, f"↻ /loop wakeup #{mgr.state.ticks_fired if mgr.state else '?'} firing…")
if wakeup.lstrip().startswith("/"):
_notif_slash_loop_tick(rid, sid, session, mgr, wakeup)
else:
_emit("message.start", sid)
_run_prompt_submit(rid, sid, session, wakeup)
except Exception as exc:
_notif_log_failure("loop wakeup dispatch failed", exc)
_notif_release_turn(session)
with contextlib.suppress(Exception):
mgr.abandon_tick()
def _kb_first_line(value: Any, limit: int) -> str:
lines = str(value).strip().splitlines()
return f"\n{lines[0][:limit]}" if lines else ""
def _kb_completed(task, payload: dict, title: str) -> str:
handoff = (_kb_first_line(payload["summary"], 200) if payload.get("summary")
else _kb_first_line(task.result, 160) if getattr(task, "result", None) else "")
return f" done — {title}{handoff}"
def _kb_timed_out(task, payload: dict, title: str) -> str:
with contextlib.suppress(TypeError, ValueError):
return f" timed out (max_runtime={int(payload.get('limit_seconds') or 0)}s); will retry"
return " timed out (max_runtime=0s); will retry"
# kind -> (glyph, suffix after "Kanban <id>"); silent kinds (archived/unblocked) are absent → None.
_KANBAN_EVENT_FORMATTERS = {
"completed": ("✔", _kb_completed),
"blocked": ("⏸", lambda t, p, title: " blocked" + (f": {str(p.get('reason'))[:160]}" if p.get("reason") else "")),
"gave_up": ("✖", lambda t, p, title: " gave up after repeated spawn failures"
+ (f"\n{str(p.get('error'))[:200]}" if p.get("error") else "")),
"crashed": ("✖", lambda t, p, title: " worker crashed (pid gone); dispatcher will retry"),
"timed_out": ("⏱", _kb_timed_out),
"status": ("🔄", lambda t, p, title: f" → {p.get('status') or ''}"),
}
def _format_kanban_event_text(sub: dict, task, ev, board_slug: str) -> Optional[str]:
"""Single-line notification text for one kanban event; wording mirrors gateway/kanban_watchers.py (reads the same
as on Telegram). None for silent kinds."""
if (entry := _KANBAN_EVENT_FORMATTERS.get(getattr(ev, "kind", ""))) is None:
return None
glyph, fmt = entry
task_id = sub.get("task_id", "")
title = (getattr(task, "title", None) or task_id)[:120]
who = getattr(task, "assignee", None) or ""
prefix = f"{glyph} " + (f"[{board_slug}] " if board_slug else "") + (f"@{who} " if who else "")
return f"{prefix}Kanban {task_id}{fmt(task, getattr(ev, 'payload', None) or {}, title)}"
def _kb_board_key(_kb, board_meta) -> tuple[str, str]:
"""(slug, resolved DB identity) — multiple slugs can point at one DB when HERMES_KANBAN_DB pins it."""
slug = (board_meta or {}).get("slug") or _kb.DEFAULT_BOARD
db_path = (board_meta or {}).get("db_path")
try:
return slug, str(Path(db_path).expanduser().resolve() if db_path else _kb.kanban_db_path(slug).resolve())
except Exception:
return slug, f"slug:{slug}"
def _kb_poll_board(_kb, slug: str, session_key: str) -> list:
"""Claim + format this session's unseen events on one board. One poller per live session: the board is not opened
writable unless it has a subscription owned by this exact session (a failed read-only probe — locked/corrupt DB —
falls through so delivery is preserved)."""
from hermes_cli import kanban_db_connect as _kbc
from hermes_cli import kanban_db_notify as _kbn
with contextlib.suppress(Exception):
if _kbn.count_notify_subs(board=slug, platform="tui", chat_id=session_key) == 0:
return []
try:
conn = _kbc.connect(board=slug)
except Exception:
return []
texts: list = []
with contextlib.closing(conn):
try:
subs = _kbn.list_notify_subs(conn)
except Exception:
return []
for sub in subs:
if (sub.get("platform") or "").lower() != "tui" or sub.get("chat_id") != session_key:
continue
sub_ident = dict(task_id=sub["task_id"], platform=sub["platform"], chat_id=sub["chat_id"],
thread_id=sub.get("thread_id") or "")
_old, _new, events = _kbn.claim_unseen_events_for_sub(conn, kinds=_KANBAN_NOTIFY_KINDS, **sub_ident)
if not events:
continue
task = _kb.get_task(conn, sub["task_id"])
texts.extend(t for t in (_format_kanban_event_text(sub, task, ev, slug) for ev in events) if t)
# Unsubscribe only on archive: ``done`` is reversible in review/controller flows, so keeping the sub lets a
# later reopen notify the same session. The claimed cursor prevents replay.
if task and getattr(task, "status", "") == "archived":
with contextlib.suppress(Exception):
_kbn.remove_notify_sub(conn, **sub_ident)
return texts
def _collect_kanban_notifications(session: dict) -> list:
"""Claim unseen terminal kanban events for this session's ``platform="tui"`` subscriptions (``kanban_create``
auto-subscribes with ``chat_id=HERMES_SESSION_KEY``; no "tui" messaging adapter exists, so this poller is the
delivery path). Same atomic cursor-claim as the gateway notifier: exactly-once even if a gateway polls the same DB.
See #59890.
"""
session_key = str(session.get("session_key") or "")
if not session_key or session.get("_finalized"):
return []
try:
from hermes_cli import kanban_db as _kb
except Exception:
return []
try:
boards = _kb.list_boards(include_archived=False)
except Exception:
try:
boards = [_kb.read_board_metadata(_kb.DEFAULT_BOARD)]
except Exception:
return []
# dict keyed by resolved DB identity: first slug per DB wins (a pinned HERMES_KANBAN_DB aliases slugs).
unique = {}
for slug, resolved in (_kb_board_key(_kb, board_meta) for board_meta in boards):
unique.setdefault(resolved, slug)
return [t for slug in unique.values() for t in _kb_poll_board(_kb, slug, session_key)]
def _notif_poll_kanban(sid: str, session: dict) -> None:
"""One kanban poll: emit new texts, buffer them, and run the buffered batch as a turn if idle. Events are
cursor-claimed (never re-queued), so they wait in the buffer instead of dropping the agent turn."""
try:
texts = _collect_kanban_notifications(session)
except Exception as exc:
_notif_log_failure("kanban notification poll failed", exc)
texts = []
for text in texts:
_emit("status.update", sid, {"kind": "process", "text": text})
if texts:
session.setdefault("_kanban_pending", []).extend(texts)
if not session.get("_kanban_pending") or not _notif_claim_turn(session):
return
with session["history_lock"]:
batch, session["_kanban_pending"] = list(session.get("_kanban_pending") or []), []
with contextlib.suppress(Exception):
_notif_submit(f"__notif__{int(time.time() * 1000)}", sid, session, "\n".join(batch), "kanban notification dispatch failed")
def _notif_dispatch_event(sid: str, session: dict, evt: dict, text: str) -> None:
"""Run the claimed (running=True) agent turn for one notification event."""
from tools.async_delegation import claim_event_delivery, complete_event_delivery, release_event_delivery
if (claim := claim_event_delivery(evt, "tui-poller")) is None:
return
kwargs = ({"display_kind": "async_delegation_complete", "display_metadata": _async_delegation_display_metadata(evt)}
if evt.get("type") == "async_delegation" else {})
try:
_notif_submit(f"__notif__{int(time.time() * 1000)}", sid, session, text, "notification poller dispatch failed", **kwargs)
except Exception:
release_event_delivery(evt, claim)
return
complete_event_delivery(evt, claim)
def _notif_handle_event(sid, session, evt, emitted, registry, fmt, deferred, completions=None, *, owned=False) -> bool:
"""Route one dequeued event: foreign (another live session owns it) → requeued, or onto ``deferred`` during the
shutdown drain; unowned (addressed but unprovable — never adopt an orphan) → dropped, except delegation payloads
deferred for a resume; ours (or ownerless legacy, kept process-global) → status.update once, then an agent turn if
idle. False = the drain must stop (session busy). ``owned`` skips the ownership gates for events a caller already
drained through ``_session_owns_notification_event`` (the post-turn safety net), so lineage resolves once."""
queue = registry.completion_queue
evt_type, is_delegation = evt.get("type", "completion"), evt.get("type") == "async_delegation"
if not owned and _notification_event_belongs_elsewhere(sid, session, evt):
if deferred is not None:
deferred.append(evt)
else: # otherwise a process started in session A surfaces in whichever poller wakes first
queue.put(evt)
time.sleep(0.1)
return True
if not owned and _notification_event_requires_owner(evt) and not _session_owns_notification_event(sid, session, evt):
origin, key = str(evt.get("origin_ui_session_id") or ""), str(evt.get("session_key") or "")
if deferred is None:
(logger.warning if is_delegation else logger.debug)(
"Dropping unowned %s notification (origin=%r key=%r) instead of delivering to session %s",
evt_type, origin, key, sid)
elif is_delegation:
deferred.append(evt)
else:
logger.debug("Dropping unowned %s notification during shutdown drain (origin=%r key=%r)", evt_type, origin, key)
return True
if evt_type == "completion" and registry.is_completion_consumed(evt.get("session_id", "")):
return True
text = fmt(evt)
if not text:
return True
# Emit once per dedup key: a re-queued completion would otherwise re-emit every 0.5s while the session is busy,
# while distinct watch_match events from one process must stay visible.
dedup_key = _notification_event_dedup_key(evt)
if dedup_key not in emitted:
from tools.process_registry_notifications import async_delegation_display_text
display_text = async_delegation_display_text(evt) if is_delegation else text
_emit("status.update", sid, {"kind": "process", "text": display_text})
emitted.add(dedup_key)
if evt_type == "completion" and completions is not None:
completions.append((evt, text))
return True
if not _notif_claim_turn(session):
queue.put(evt)
if deferred is not None:
return False
time.sleep(0.25) # back off: the re-queued event keeps the queue non-empty, else this loop spins at 100% CPU
return True
_notif_dispatch_event(sid, session, evt, text)
return True
def _notif_dispatch_completions(sid, session, notifications, registry, deferred):
from tools.process_registry_notifications import ProcessNotificationBatch
from tools.async_delegation import claim_event_delivery, complete_event_delivery, release_event_delivery
if not notifications:
return
if not _notif_claim_turn(session):
for event, _text in notifications:
(deferred.append if deferred is not None else registry.completion_queue.put)(event)
if deferred is None:
time.sleep(0.25)
return
claimed = [(event, text, claim) for event, text in notifications
if (claim := claim_event_delivery(event, "tui-completion-batch")) is not None]
text = ProcessNotificationBatch(tuple((event, text) for event, text, _claim in claimed)).render(registry)
if text is None:
_notif_release_turn(session)
try:
if text is not None:
_notif_submit(f"__notif__{int(time.time() * 1000)}", sid, session, text,
"completion batch dispatch failed")
except Exception:
for event, _text, claim in claimed:
release_event_delivery(event, claim)
return
for event, _text, claim in claimed:
complete_event_delivery(event, claim)
def _notif_handle_ready(sid, session, events, emitted, registry, fmt, deferred, *, owned=False):
"""One ready snapshot: ownership and UI emission per event, one turn per completion run."""
completions = []
for index, event in enumerate(events):
if event.get("type", "completion") != "completion":
_notif_dispatch_completions(sid, session, completions, registry, deferred)
completions = []
if not _notif_handle_event(sid, session, event, emitted, registry, fmt, deferred, completions, owned=owned):
for remaining in events[index + 1:]:
(deferred.append if deferred is not None else registry.completion_queue.put)(remaining)
break
_notif_dispatch_completions(sid, session, completions, registry, deferred)
def _poll_bot_live_delivery_once(sid: str, session: dict) -> bool:
"""Run one durable envelope only after local FIFO/continuations yield the idle boundary."""
from tools.bot_live_delivery import claim_pending_delivery, complete_delivery, find_canonical_live_owner
home = _session_home(session)
with session["history_lock"]:
if any(session.get(key) for key in (
"running", "_closing", "_finalized", "queued_prompt", "queued_prompts",
"_auto_continue_scheduled")) or session.get("agent") is None:
return False
lease = session.get("active_session_lease")
if lease is None or getattr(lease, "released", False):
return False
owner = find_canonical_live_owner(home)
if (not owner or owner.get("lease_id") != lease.lease_id
or owner.get("live_session_id") != sid
or owner.get("session_id") != session.get("session_key")):
return False
# The mailbox matches each envelope to this pinned lease/live id and compression lineage.
claimed = claim_pending_delivery(home, owner)
if claimed is None:
return False
session["running"] = True
delivery_id = str(claimed["id"])
def terminal_receipt(terminal: dict) -> None:
status = str(terminal.get("status") or "failed")
error = str(terminal.get("error") or "")
reason = "cancelled" if status == "cancelled" else ""
if status not in {"settled", "cancelled"}:
from tools.bot_failure_reasons import classify_agent_error
reason = classify_agent_error(error)
# Let a failed write propagate: the turn must not retire its crash marker without its receipt.
complete_delivery(home, delivery_id, status=status,
reply=str(terminal.get("text") or "") if status == "settled" else "",
error=error, reason=reason)
try:
started = _run_prompt_submit(f"__bot_dm__{delivery_id}", sid, session, claimed["message"],
image_paths=[], terminal_callback=terminal_receipt,
turn_author=claimed.get("author") or None)
except Exception as exc:
_notif_release_turn(session)
terminal_receipt({"status": "failed", "error": str(exc)})
raise
if not started:
_notif_release_turn(session)
terminal_receipt({"status": "failed", "error": "live session owner could not start the delivery turn"})
return started
def _notification_poller_loop(stop_event: threading.Event, sid: str, session: dict) -> None:
"""Daemon thread (started by _init_session()) that drains the process-global completion_queue for this session
(ownership routing: _notif_handle_event) and polls ``kanban_notify_subs`` every ``_KANBAN_POLL_SECONDS`` — the
delivery path for platform="tui" rows.
Also polls ``kanban_notify_subs`` every ``_KANBAN_POLL_SECONDS`` for this session's TUI kanban
subscriptions and delivers terminal task events the same way (status.update + agent turn) — the delivery
path tools/kanban_tools.py documents for platform="tui" rows (issue #59890).
"""
from tools.process_registry import process_registry
from tools.process_registry_notifications import format_process_notification
queue = process_registry.completion_queue
emitted = session.setdefault("_notification_emitted", set())
handle = lambda events, deferred: _notif_handle_ready( # noqa: E731
sid, session, events, emitted, process_registry, format_process_notification, deferred)
last_kanban_poll = last_loop_poll = 0.0
last_wisdom_poll = 0.0
while not stop_event.is_set() and not session.get("_finalized"):
now = time.monotonic()
if not session.get("running") and now - last_wisdom_poll >= _WISDOM_POLL_SECONDS:
last_wisdom_poll = now
_sync_wisdom_activity_notice(sid, session)
try:
_poll_bot_live_delivery_once(sid, session)
except Exception:
logger.warning("Bot live-owner delivery poll failed", exc_info=True)
# /loop and /heartbeat wakeup drivers: fire a due tick for THIS session while idle (same claim-under-lock
# as kanban dispatch). An active non-parked /goal owns the idle boundary and defers the loop tick.
if now - last_loop_poll >= _LOOP_POLL_SECONDS:
last_loop_poll = now
for what, fire in (("loop wakeup", _maybe_fire_tui_loop_tick), ("heartbeat", _maybe_fire_tui_heartbeat_tick)):
try:
fire(sid, session)
except Exception as tick_exc:
_notif_log_failure(f"{what} poll failed", tick_exc)
if now - last_kanban_poll >= _KANBAN_POLL_SECONDS:
last_kanban_poll = now
_notif_poll_kanban(sid, session)
try:
evt = queue.get(timeout=0.5)
except Exception:
continue
ready = [evt]
for _ in range(queue.qsize()):
try:
ready.append(queue.get_nowait())
except Exception:
break
handle(ready, None)
# Drain remaining events after the stop signal so nothing is lost on shutdown; foreign and orphaned-delegation
# events are handed back to the shared queue afterwards.
deferred: list = []
ready = []
for _ in range(queue.qsize()):
try:
ready.append(queue.get_nowait())
except Exception:
break
handle(ready, deferred)
for evt in deferred:
queue.put(evt)
def _async_delegation_display_metadata(evt: dict) -> dict:
"""Build display-only metadata before the completion event is formatted."""
from tools.process_registry_notifications import async_delegation_display_text
raw_results = evt.get("results")
results: list[dict] = [r for r in raw_results if isinstance(r, dict)] if isinstance(raw_results, list) else []
task_count = len(results) or 1
completed_count = sum(1 for r in results if r.get("status") in {"completed", "success"})
failed_count = sum(1 for r in results if r.get("status") in {"failed", "error"})
duration = evt.get("total_duration_seconds") or evt.get("duration_seconds")
return {"display_text": async_delegation_display_text(evt),
"delegation_id": str(evt.get("delegation_id") or ""), "task_count": task_count,
"completed_count": completed_count or task_count - failed_count, "failed_count": failed_count,
**({"duration_seconds": duration} if isinstance(duration, (int, float)) else {})}
_desktop_ui_wired = False
def _wire_desktop_sinks() -> None:
"""Idempotently wire process-registry and desktop-tool sinks to renderer events: `agent.terminal.output` and
`terminal.close` (drops a tab without killing the process) route to the window owning the process; desktop-only
tools pass the turn's ``HERMES_UI_SESSION_ID`` as ``sid``. `_emit` is thread-safe."""
global _desktop_ui_wired
from tools.process_registry import process_registry
def _owner_sid(session) -> str:
# session may be None (process already finished/pruned) — the tab can still linger and be closed.
session_key = str(getattr(session, "session_key", "") or "") if session is not None else ""
if not session_key:
return ""
with _sessions_lock:
return next((sid for sid, s in _sessions.items() if str(s.get("session_key") or "") == session_key), "")
if getattr(process_registry, "on_output", None) is None:
process_registry.on_output = lambda session, chunk: _emit(
"agent.terminal.output", _owner_sid(session), {"process_id": session.id, "chunk": chunk})
if getattr(process_registry, "on_close", None) is None:
process_registry.on_close = lambda session, pid: _emit("terminal.close", _owner_sid(session), {"process_id": pid})
if not _desktop_ui_wired:
with contextlib.suppress(Exception):
from tools import desktop_ui
desktop_ui.set_emitter(lambda sid, event, payload: _emit(event, sid, payload))
_desktop_ui_wired = True
# (stop_event, thread) for every poller started in this process, pruned of dead threads on each spawn. Test teardowns
# reap leaked pollers through it: an unjoined poller steals events off the process-global queue in a LATER test.
_notification_pollers: list = []
def _start_notification_poller(sid: str, session: dict) -> threading.Event:
"""Start the background notification poller for a TUI session (thread name is greppable)."""
_wire_desktop_sinks()
stop = threading.Event()
t = threading.Thread(target=_notification_poller_loop, args=(stop, sid, session), daemon=True, name=f"tui-notif-poller-{sid}")
_notification_pollers[:] = [(s, th) for (s, th) in _notification_pollers if th.is_alive()] + [(stop, t)]
t.start()
return stop
def _hud_surface_note(session: dict) -> str:
"""The HUD-mode note for this turn, or "" when it was not typed there."""
if session.get("client_surface") != "hud":
return ""
from agent.prompt_builder import hud_surface_note
return hud_surface_note(getattr(session.get("agent"), "valid_tool_names", None))
def _prepend_note(run_message: Any, note: str) -> Any:
"""Prefix a per-turn note onto the MODEL INPUT, leaving the prompt alone: everything the model must know that the
user did not type (interrupted reply, reactions, surface) arrives this way, so no scaffolding reaches the
transcript and no sent message is rewritten — the cached prefix survives."""
if note and isinstance(run_message, str):
return f"{note}\n\n{run_message}"
if note and isinstance(run_message, list):
return [{"type": "text", "text": note}, *run_message]
return run_message
def register(server) -> None:
"""Publish this module's helpers + handlers onto ``server``, rebound to its globals."""
bind_module(globals(), server, skip=("_",))