From fe08ed631b5b2ff0e7c72367ead3858d768ca554 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 23:41:53 -0700 Subject: [PATCH] =?UTF-8?q?refactor(tui):=20server.py=20=E2=80=94=20=5Fsta?= =?UTF-8?q?rtup=5Fsystem=5Fprompt/=5Fhydrate=5Fsession=5Fcwd/=5Feffective?= =?UTF-8?q?=5Fapproval=5Fstate/=5Fturn=5Fstarted=5Fat=20helpers,=20=5Ftui?= =?UTF-8?q?=5Fnotice,=20compact=20session-info/toolset/runtime=20docstring?= =?UTF-8?q?s?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- tui_gateway/server.py | 619 ++++++++++++++++-------------------------- 1 file changed, 230 insertions(+), 389 deletions(-) diff --git a/tui_gateway/server.py b/tui_gateway/server.py index 08f8648a47..eb3fd170c7 100644 --- a/tui_gateway/server.py +++ b/tui_gateway/server.py @@ -1983,6 +1983,10 @@ def _enabled_mcp_server_names() -> tuple[set[str], set[str]]: return set(), set() +def _tui_notice(text: str) -> None: + print(text, file=sys.stderr, flush=True) + + def _resolve_explicit_toolsets(explicit: list[str], validate_toolset) -> list[str] | None | bool: """Resolve a HERMES_TUI_TOOLSETS pin: list, None for "all", False when nothing was valid.""" built_in = [name for name in explicit if validate_toolset(name)] @@ -1998,13 +2002,10 @@ def _resolve_explicit_toolsets(explicit: list[str], validate_toolset) -> list[st built_in.extend(plugin_valid) unresolved = [name for name in unresolved if name not in plugin_valid] if any(name in {"all", "*"} for name in built_in): - ignored = [name for name in explicit if name not in {"all", "*"}] - if ignored: - print( + if ignored := [name for name in explicit if name not in {"all", "*"}]: + _tui_notice( "[tui] HERMES_TUI_TOOLSETS=all enables every toolset; " - f"ignoring additional entries: {', '.join(ignored)}", - file=sys.stderr, - flush=True, + f"ignoring additional entries: {', '.join(ignored)}" ) return None if not unresolved: @@ -2014,38 +2015,26 @@ def _resolve_explicit_toolsets(explicit: list[str], validate_toolset) -> list[st disabled = [name for name in unresolved if name in mcp_disabled] unknown = [name for name in unresolved if name not in mcp_names and name not in mcp_disabled] if unknown: - print( - f"[tui] ignoring unknown HERMES_TUI_TOOLSETS entries: {', '.join(unknown)}", - file=sys.stderr, flush=True, - ) + _tui_notice(f"[tui] ignoring unknown HERMES_TUI_TOOLSETS entries: {', '.join(unknown)}") if disabled: - print( + _tui_notice( "[tui] ignoring disabled MCP servers in HERMES_TUI_TOOLSETS " "(set enabled: true in config.yaml to use): " - f"{', '.join(disabled)}", - file=sys.stderr, - flush=True, + f"{', '.join(disabled)}" ) return (built_in + mcp_valid) or False def _load_enabled_toolsets(platform: str | None = None) -> list[str] | None: - """Resolve the agent's toolsets for this desktop/TUI session (None = all). + """The agent's toolsets for this desktop/TUI session (None = all). - Order: an explicit HERMES_TUI_TOOLSETS pin; else the coding posture - (collapse to the coding toolset + enabled MCP servers when sitting in a code - workspace — agent/coding_context.py, config loaded lazily there); else the - configured CLI toolsets. The client-surface (pane/project) toolsets are off - _HERMES_CORE_TOOLS so no other platform carries their schema; this resolver - runs only in the desktop/TUI gateway, so folding them in here is the gate - that exposes them on exactly the surface that can answer them. + Order: an explicit HERMES_TUI_TOOLSETS pin; else the coding posture (agent/coding_context.py + collapses to the coding toolset + enabled MCP servers inside a code workspace); else the + configured CLI toolsets. The client-surface (pane/project) toolsets are folded in here because + this resolver runs only on the surface that can answer them. """ session_platform = platform or _resolve_session_platform() - explicit = [ - item.strip() - for item in os.environ.get("HERMES_TUI_TOOLSETS", "").split(",") - if item.strip() - ] + explicit = [item.strip() for item in os.environ.get("HERMES_TUI_TOOLSETS", "").split(",") if item.strip()] fallback_notice = None if not explicit: with contextlib.suppress(Exception): @@ -2061,28 +2050,21 @@ def _load_enabled_toolsets(platform: str | None = None) -> list[str] | None: resolved = _resolve_explicit_toolsets(explicit, validate_toolset) if resolved is not False: return resolved - fallback_notice = ( - "[tui] no valid HERMES_TUI_TOOLSETS entries; using configured CLI toolsets" - ) + fallback_notice = "[tui] no valid HERMES_TUI_TOOLSETS entries; using configured CLI toolsets" try: from hermes_cli.config import load_config from hermes_cli.tools_config import _get_platform_tools cfg = load_config() - # include_default_mcp_servers=True is the runtime variant (the agent - # must be able to call default MCP servers); False is the config-editing - # variant. Using the wrong one here silently drops MCP tools from the TUI. + # include_default_mcp_servers=True is the runtime variant (the agent must be able to call + # default MCP servers); the config-editing variant would silently drop MCP tools from the TUI. enabled = _get_platform_tools(cfg, "cli", include_default_mcp_servers=True) if fallback_notice is not None: - print(fallback_notice, file=sys.stderr, flush=True) - if not enabled: - return None - return sorted(enabled | _gui_surface_toolsets(session_platform)) + _tui_notice(fallback_notice) + return sorted(enabled | _gui_surface_toolsets(session_platform)) if enabled else None except Exception: if fallback_notice is not None: - print( - "[tui] no valid HERMES_TUI_TOOLSETS entries and configured CLI toolsets could not be loaded; enabling all toolsets", - file=sys.stderr, - flush=True, + _tui_notice( + "[tui] no valid HERMES_TUI_TOOLSETS entries and configured CLI toolsets could not be loaded; enabling all toolsets" ) return None @@ -2100,18 +2082,15 @@ def _tool_progress_enabled(sid: str) -> bool: def _tool_lifecycle_required_for_ui(name: str) -> bool: - """Return True for tool events that are interactive UI, not optional chrome.""" - # Desktop renders clarify / setup_mcp cards from the tool-call part; with - # tool progress off, suppressing them would leave only the sidebar dot. + """Tool events that are interactive UI, not optional chrome: Desktop renders clarify / setup_mcp + cards from the tool-call part; suppressing them with progress off would leave only the sidebar dot.""" return name in ("clarify", "setup_mcp") def _restart_slash_worker(sid: str, session: dict): worker = session.get("slash_worker") - # Nothing to replace for a session that never spawned a worker; spawning - # here would fork the per-worker MCP fleet for nothing. if worker is None: - return + return # never spawned one; spawning here would fork the per-worker MCP fleet for nothing with contextlib.suppress(Exception): worker.close() try: @@ -2122,8 +2101,8 @@ def _restart_slash_worker(sid: str, session: dict): except Exception: session["slash_worker"] = None return - # Store-iff-still-mapped: the post-turn restart races a close_on_disconnect - # reap, and a bare store would orphan the fresh worker. + # Store-iff-still-mapped: the post-turn restart races a close_on_disconnect reap, and a + # bare store would orphan the fresh worker. _attach_worker(sid, session, new_worker) @@ -2139,32 +2118,26 @@ def _get_usage(agent) -> dict: } comp = getattr(agent, "context_compressor", None) if comp: - # context_used is *current-window* occupancy. Never fall back to - # usage["total"] (cumulative lifetime) — an external engine without - # last_prompt_tokens then showed 1.9m/120k clamped to 100%. A falsy - # last_prompt_tokens emits NO gauge rather than a fabricated one; the -1 - # "compression just ran" sentinel is clamped to 0 for the same reason - # (matches cli.py _get_status_bar_snapshot). - last_prompt = getattr(comp, "last_prompt_tokens", 0) or 0 - if last_prompt < 0: - last_prompt = 0 + # context_used is *current-window* occupancy. Never fall back to usage["total"] (cumulative + # lifetime: an external engine without last_prompt_tokens showed 1.9m/120k clamped to 100%). + # A falsy last_prompt_tokens emits NO gauge rather than a fabricated one; the -1 "compression + # just ran" sentinel is clamped to 0 (matches cli.py _get_status_bar_snapshot). + last_prompt = max(0, getattr(comp, "last_prompt_tokens", 0) or 0) ctx_max = getattr(comp, "context_length", 0) or 0 if ctx_max and last_prompt: usage["context_used"] = last_prompt usage["context_max"] = ctx_max usage["context_percent"] = max(0, min(100, round(last_prompt / ctx_max * 100))) usage["compressions"] = getattr(comp, "compression_count", 0) or 0 - # Cache-hit ratio + rolling latency/tps (CLI status-bar parity): - # hit = cache_read / prompt_tokens (prompt = input + cache_read + cache_write); - # latency/tps read the per-call deque history from conversation_loop. - # Omitted, not fabricated, when there is no data (Codex reports no latency; - # zero cache reads shows no hit% rather than an alarming 0). + # Cache-hit ratio + rolling latency/tps (CLI status-bar parity): hit = cache_read / prompt_tokens; + # latency/tps read the per-call deque history. Omitted, not fabricated, when there is no data + # (Codex reports no latency; zero cache reads shows no hit% rather than an alarming 0). with contextlib.suppress(Exception): _prompt_total = int(getattr(agent, "session_prompt_tokens", 0) or 0) _cache_read = int(getattr(agent, "session_cache_read_tokens", 0) or 0) if _prompt_total > 0 and _cache_read > 0: usage["cache_hit_pct"] = max(0, min(100, round(_cache_read / _prompt_total * 100))) - try: + with contextlib.suppress(Exception): # a status-bar readout must never break usage reporting _lhist = list(getattr(agent, "_api_latency_history", []) or []) _ohist = list(getattr(agent, "_api_output_history", []) or []) _n = min(len(_lhist), len(_ohist)) @@ -2179,17 +2152,11 @@ def _get_usage(agent) -> dict: usage["avg_latency_s"] = round(float(_avg_lat), 1) if _avg_vel is not None and _avg_vel == _avg_vel and 0 < _avg_vel < 1e6: usage["avg_tps"] = round(float(_avg_vel), 1) - except Exception: - # A status-bar readout must never break usage reporting. - pass - # Live count of background/async subagents still running (delegate_task - # batches + background single delegations). Mirrors the classic CLI status - # bar's ⛓ indicator; sourced from the same async_delegation registry. + # Live count of background/async subagents (CLI status bar ⛓ parity, same async_delegation registry). with contextlib.suppress(Exception): from tools.async_delegation import active_count as _async_active_count usage["active_subagents"] = _async_active_count() - # Dev-only live credits-spent readout (L0 usage-aware-credits). Gated on - # HERMES_DEV_CREDITS so the payload stays clean when the flag is off. + # Dev-only live credits-spent readout, gated on HERMES_DEV_CREDITS so the payload stays clean otherwise. if is_truthy_value(os.environ.get("HERMES_DEV_CREDITS")): with contextlib.suppress(Exception): spent = agent.get_credits_spent_micros() @@ -2199,29 +2166,22 @@ def _get_usage(agent) -> dict: def _probe_credentials(agent) -> str: - """Light credential check at session creation — returns warning or ''. - - ``no-key-required`` is a valid sentinel for keyless custom providers; only - warn when the key is genuinely missing. - """ + """Light credential check at session creation — warning or ''. (``no-key-required`` is a valid + sentinel for keyless custom providers; only a genuinely missing key warns.)""" with contextlib.suppress(Exception): - key = getattr(agent, "api_key", "") or "" - provider = getattr(agent, "provider", "") or "" - if not key: + if not (getattr(agent, "api_key", "") or ""): + provider = getattr(agent, "provider", "") or "" return f"No API key configured for provider '{provider}'. First message will fail." return "" def _probe_config_health(cfg: dict) -> str: - """Flag bare YAML keys (`agent:` with no value → None) that silently - drop nested settings. Returns warning or ''.""" + """Flag bare YAML keys (`agent:` with no value → None) that silently drop nested settings, + and a ``display.personality`` naming no known overlay. Returns warning or ''.""" if not isinstance(cfg, dict): return "" warnings: list[str] = [] - null_keys = sorted(k for k, v in cfg.items() if v is None) - if not null_keys: - pass - else: + if null_keys := sorted(k for k, v in cfg.items() if v is None): keys = ", ".join(f"`{k}`" for k in null_keys) warnings.append( f"config.yaml has empty section(s): {keys}. " @@ -2229,7 +2189,6 @@ def _probe_config_health(cfg: dict) -> str: f"empty sections silently drop nested settings." ) display_cfg = cfg.get("display") - agent_cfg = cfg.get("agent") if isinstance(display_cfg, dict): personality = str(display_cfg.get("personality", "") or "").strip().lower() if personality and personality not in {"default", "none", "neutral"}: @@ -2241,7 +2200,6 @@ def _probe_config_health(cfg: dict) -> str: "built-in or `agent.personalities` entry; personality " "overlay will be skipped." ) - _ = agent_cfg # retained for shape parity; built-ins exist without config return " ".join(warnings).strip() @@ -2253,15 +2211,11 @@ def _current_profile_name() -> str: return "default" -# Monotonic GUI<->backend contract version: the desktop refuses a backend -# reporting less (or none) with a one-click "update to align" prompt. Bump -# whenever the desktop's backend contract changes. -# v2: adds the file.attach RPC (remote-gateway non-image file upload). -# v3: adds approvals.mode config RPCs and session.info reconciliation. -# v4: session.create fast=false is an explicit per-session normal-tier override. -# v5: uvicorn ws_max_size raised for one-shot base64 file.attach frames (>16 MiB). -# v6: plugins.manage list rows carry the canonical registry key; toggles are -# key-addressed (keyless rows render read-only in Desktop Settings). +# Monotonic GUI<->backend contract version: the desktop refuses a backend reporting less (or none) +# with a one-click "update to align" prompt. Bump whenever the desktop's backend contract changes. +# v2 file.attach RPC; v3 approvals.mode RPCs + session.info reconciliation; v4 session.create +# fast=false = explicit normal-tier override; v5 ws_max_size raised for >16 MiB file.attach frames; +# v6 plugins.manage rows carry the canonical registry key (toggles key-addressed). DESKTOP_BACKEND_CONTRACT = 6 @@ -2276,14 +2230,9 @@ def _session_usage_snapshot(session: dict | None) -> dict: def _project_info_for_cwd(cwd: str) -> dict | None: - """Return the first-class Project owning ``cwd`` for UI status surfaces. - - Backed by the per-profile projects.db (the same store the desktop's project - tree caches), so the TUI status label, the desktop status bar, and ``/status`` - all name the session's workspace identically. Only explicit, named projects - resolve here — an auto-discovered repo root has no projects.db row, so it - falls back to the cwd leaf on every surface. - """ + """The first-class Project owning ``cwd`` (per-profile projects.db, the store the desktop's project + tree caches) so TUI status, desktop status bar and ``/status`` name the workspace identically. + Only explicit named projects resolve; an auto-discovered repo root falls back to the cwd leaf.""" if not str(cwd or "").strip(): return None try: @@ -2292,68 +2241,56 @@ def _project_info_for_cwd(cwd: str) -> dict | None: project = pdb.project_for_path(conn, cwd) if project is None: return None - return { - "id": project.id, "slug": project.slug, "name": project.name, - "primary_path": project.primary_path, - } + return {"id": project.id, "slug": project.slug, "name": project.name, "primary_path": project.primary_path} except Exception: logger.debug("failed to resolve project for cwd", exc_info=True) return None +def _turn_started_at(session: dict | None) -> float | None: + """Epoch seconds the current turn started, or None when idle — lets the desktop keep the + turn-elapsed timer across session switches (cold resume) instead of resetting to 0:00.""" + inflight = (session or {}).get("inflight_turn") + return float(inflight["started_at"]) if isinstance(inflight, dict) and inflight.get("started_at") else None + + +def _effective_approval_state(session_key: str) -> tuple[bool, str]: + """(yolo, approval_mode): the same three sources check_all_command_guards() ORs — approvals.mode=off, + the process --yolo env, the per-session flag. The session flag alone would show YOLO "off" while + config silently auto-approves every dangerous command.""" + try: + from tools.approval import _YOLO_MODE_FROZEN, is_session_yolo_enabled + session_yolo = bool(is_session_yolo_enabled(session_key)) if session_key else False + approval_mode = _load_approval_mode() + return bool(_YOLO_MODE_FROZEN) or session_yolo or approval_mode == "off", approval_mode + except Exception: + return False, "manual" + + def _session_info(agent, session: dict | None = None) -> dict: if session is None: - for candidate in _sessions.values(): - if candidate.get("agent") is agent: - session = candidate - break + session = next((c for c in _sessions.values() if c.get("agent") is agent), None) + sess = session or {} mirror = _metadata_mirror(session) cwd = _display_session_cwd(session) - session_key = str((session or {}).get("session_key") or getattr(agent, "session_id", "") or "") - cfg_personality = _display_cfg().get("personality") or "" - personality = (session or {}).get("personality", cfg_personality) + session_key = str(sess.get("session_key") or getattr(agent, "session_id", "") or "") + personality = sess.get("personality", _display_cfg().get("personality") or "") reasoning_config = getattr(agent, "reasoning_config", None) reasoning_effort = "" if isinstance(reasoning_config, dict): - if reasoning_config.get("enabled") is False: - # Disabled must differ from unset ("" = provider default), or the - # desktop adopts "" after the first turn and loses "thinking off". - reasoning_effort = "none" - else: - reasoning_effort = str(reasoning_config.get("effort", "") or "") + # Disabled must differ from unset ("" = provider default), or the desktop adopts "" after + # the first turn and loses "thinking off". + reasoning_effort = "none" if reasoning_config.get("enabled") is False else str(reasoning_config.get("effort", "") or "") service_tier = getattr(agent, "service_tier", None) or mirror.get("service_tier") or "" - # Effective approval bypass = the same three sources check_all_command_guards() - # ORs: approvals.mode=off, the process --yolo env, the per-session flag. - # Reporting only the session flag would show YOLO "off" while config - # silently auto-approves every dangerous command. - yolo = False - approval_mode = "manual" - try: - from tools.approval import _YOLO_MODE_FROZEN, is_session_yolo_enabled - session_yolo = (bool(is_session_yolo_enabled(session_key)) if session_key else False) - approval_mode = _load_approval_mode() - yolo = bool(_YOLO_MODE_FROZEN) or session_yolo or approval_mode == "off" - except Exception: - yolo = False - # A switch queued mid-turn applies at the next turn start, so agent.model - # still reads the OLD model; report the pending pick so the end-of-turn - # settle doesn't blip the UI back before the switch lands. - pending_switch = (session or {}).get("pending_model_switch") or {} + yolo, approval_mode = _effective_approval_state(session_key) + # A switch queued mid-turn applies at the next turn start, so agent.model still reads the OLD + # model; report the pending pick so the end-of-turn settle doesn't blip the UI back first. + pending_switch = sess.get("pending_model_switch") or {} pending_model = str(pending_switch.get("display_model") or "").strip() pending_provider = str(pending_switch.get("display_provider") or "").strip() - # Epoch seconds the current turn started, or None when idle. Lets the - # desktop preserve the turn-elapsed timer across session switches (cold - # resume path) instead of resetting it to 0:00. - inflight = (session or {}).get("inflight_turn") - turn_started_at = ( - float(inflight["started_at"]) - if isinstance(inflight, dict) and inflight.get("started_at") - else None - ) info: dict = { "model": pending_model or mirror.get("model", getattr(agent, "model", "")), - "provider": pending_provider - or mirror.get("provider", getattr(agent, "provider", "")), + "provider": pending_provider or mirror.get("provider", getattr(agent, "provider", "")), "reasoning_effort": reasoning_effort, "service_tier": service_tier, "fast": service_tier == "priority", @@ -2366,9 +2303,9 @@ def _session_info(agent, session: dict | None = None) -> dict: "project": _project_info_for_cwd(cwd), "terminal_backend": _effective_terminal_backend(), "personality": str(personality or ""), - "running": bool((session or {}).get("running")), - "turn_started_at": turn_started_at, - "title": _session_live_title(session or {}, session_key) if session_key else "", + "running": bool(sess.get("running")), + "turn_started_at": _turn_started_at(session), + "title": _session_live_title(sess, session_key) if session_key else "", "stored_session_id": session_key or "", "desktop_contract": DESKTOP_BACKEND_CONTRACT, "version": "", @@ -2386,7 +2323,7 @@ def _session_info(agent, session: dict | None = None) -> dict: from hermes_cli import __version__, __release_date__ info["version"] = __version__ info["release_date"] = __release_date__ - live_agent = agent is not None and not (session or {}).get("_compute_host_active") + live_agent = agent is not None and not sess.get("_compute_host_active") if live_agent: with contextlib.suppress(Exception): from model_tools import get_toolset_for_tool @@ -2404,9 +2341,7 @@ def _session_info(agent, session: dict | None = None) -> dict: info["mcp_servers"] = [] with contextlib.suppress(Exception): info["system_prompt"] = ( - mirror.get("system_prompt") - if "system_prompt" in mirror - else getattr(agent, "_cached_system_prompt", "") or "" + mirror.get("system_prompt") if "system_prompt" in mirror else getattr(agent, "_cached_system_prompt", "") or "" ) with contextlib.suppress(Exception): from hermes_cli.banner import get_update_result @@ -2419,16 +2354,9 @@ def _session_info(agent, session: dict | None = None) -> dict: def _tool_ctx(name: str, args: dict) -> str: - """Argument preview for a tool row — never a phrased label. - - Clients own their own phrasing: the TUI wraps this as ``Terminal("...")`` - and the desktop prepends its own localized verb ("Running"/"Ran"). Sending - ``build_tool_label`` here instead of the raw preview stutters the verb on - both surfaces ("Running Running sleep 70 + 2 commands") and leaks a display - label into the desktop's ``args.context``, where it stands in for the real - command. The friendly labels belong on the CLI spinner, which builds them - from ``build_tool_label`` at its own call sites. - """ + """Argument preview for a tool row — never a phrased label: clients own their phrasing (the TUI + wraps it as ``Terminal("...")``, the desktop prepends a localized verb), so ``build_tool_label`` + here would stutter ("Running Running …") and leak a display label into the desktop's ``args.context``.""" try: from agent.display import build_tool_preview return build_tool_preview(name, args, max_len=80) or "" @@ -2445,35 +2373,23 @@ def _emit_session_info_for_session(sid: str, session: dict) -> None: def broadcast_session_info() -> None: - """Re-emit ``session.info`` to every live session. - - For approvals-config writers that bypass the ``config.set`` RPC (which - re-emits itself): the REST config saves and the ``/approvals`` slash - mirror. Only reaches sessions in THIS process; a spawned - ``tui_gateway.entry`` child gateway has its own ``_sessions``. - """ + """Re-emit ``session.info`` to every live session — for approvals-config writers that bypass the + self-re-emitting ``config.set`` RPC (REST config saves, the ``/approvals`` slash mirror). Only + reaches THIS process; a spawned ``tui_gateway.entry`` child gateway has its own ``_sessions``.""" with _sessions_lock: sessions = list(_sessions.items()) for sid, sess in sessions: _emit_session_info_for_session(sid, sess) -# Tool Args/Result text shipped to the TUI for the verbose trail line. The TUI -# renders only a small persisted preview (ui-tui VERBOSE_TRAIL_MAX_CHARS), kept -# all session and expanded by default — so shipping more than that is pure pipe def _schedule_mcp_late_refresh(sid: str, agent) -> None: """Refresh a session's tool snapshot when MCP discovery lands late. - The agent snapshots ``agent.tools`` once at build; ``_make_agent`` only - waits a bounded ``mcp_discovery_timeout`` (default 1.5s), so a slow server - (HTTP MCP on first connect) lands after the build and its tools are missing - for the whole session. A daemon waits for discovery to finish, then does - the same rebuild ``/reload-mcp`` performs and re-emits ``session.info``. - - Cache safety: the rebuild runs only while the session is pre-first-turn - (nothing cached to invalidate). Once a message was sent the snapshot stays - frozen — late tools then need an explicit, consent-gated ``/reload-mcp``. - No-op when discovery already finished before the build. + ``_make_agent`` waits only a bounded ``mcp_discovery_timeout``, so a slow server (HTTP MCP on + first connect) lands after the build and its tools are missing all session. A daemon joins + discovery, then does the same rebuild ``/reload-mcp`` performs and re-emits ``session.info``. + Cache safety: only while pre-first-turn (nothing cached to invalidate); afterwards the snapshot + stays frozen and late tools need an explicit, consent-gated ``/reload-mcp``. """ try: from tui_gateway.entry import mcp_discovery_in_flight, join_mcp_discovery @@ -2483,39 +2399,26 @@ def _schedule_mcp_late_refresh(sid: str, agent) -> None: return def _wait_then_refresh() -> None: - # Bounded but generous — a server still not connected after this is - # genuinely slow/dead; the user can /reload-mcp once it recovers. + # Bounded but generous — a server still not connected after this is genuinely slow/dead. if not join_mcp_discovery(timeout=30.0): return with _sessions_lock: session = _sessions.get(sid) - # Session may have been closed/reset while we waited. if session is None or session.get("agent") is not agent: - return - # Cache safety: never rebuild the tool list once the conversation - # has started — that would invalidate the cached prompt prefix. - if ( - int(getattr(agent, "_user_turn_count", 0) or 0) > 0 - or int(getattr(agent, "_api_call_count", 0) or 0) > 0 - ): - return + return # closed/reset while we waited + if int(getattr(agent, "_user_turn_count", 0) or 0) > 0 or int(getattr(agent, "_api_call_count", 0) or 0) > 0: + return # conversation started: a rebuild would invalidate the cached prompt prefix try: from tools.mcp_tool import refresh_agent_mcp_tools added = refresh_agent_mcp_tools(agent, quiet_mode=True) except Exception as exc: - logger.warning( - "Late MCP refresh: tool snapshot rebuild failed for %s: %s", sid, exc, - ) + logger.warning("Late MCP refresh: tool snapshot rebuild failed for %s: %s", sid, exc) return - # No new tools landed (discovery added nothing) → don't churn the client. if not added: - return + return # discovery added nothing → don't churn the client info = _session_info(agent, session) - # Emit outside the lock — write_json must not block under _sessions_lock. - _emit("session.info", sid, info) - threading.Thread( - target=_wait_then_refresh, name=f"tui-mcp-late-refresh-{sid}", daemon=True, - ).start() + _emit("session.info", sid, info) # outside the lock — write_json must not block under _sessions_lock + threading.Thread(target=_wait_then_refresh, name=f"tui-mcp-late-refresh-{sid}", daemon=True).start() class _RuntimeFallbackResolution(NamedTuple): @@ -2524,24 +2427,19 @@ class _RuntimeFallbackResolution(NamedTuple): used_fallback: bool -def _resolve_runtime_with_fallback( - resolve_kwargs: dict | None = None, -) -> _RuntimeFallbackResolution: +def _resolve_runtime_with_fallback(resolve_kwargs: dict | None = None) -> _RuntimeFallbackResolution: """Resolve the primary runtime or one complete provider/model fallback. - Setup-time auth fallback only accepts entries with both fields. Provider- - only entries are skipped so the unavailable primary model can never leak - into a different runtime. ``used_fallback`` remains explicit rather than - overloading a nullable model as control flow. + Setup-time auth fallback only accepts entries with both fields: provider-only entries are + skipped so the unavailable primary model can never leak into a different runtime. + ``used_fallback`` stays explicit rather than overloading a nullable model as control flow. """ from hermes_cli.auth import AuthError from hermes_cli.runtime_provider import resolve_runtime_provider - kwargs = resolve_kwargs or {} try: - return _RuntimeFallbackResolution(resolve_runtime_provider(**kwargs), None, False) + return _RuntimeFallbackResolution(resolve_runtime_provider(**(resolve_kwargs or {})), None, False) except AuthError as primary_exc: - fb_chain = _load_fallback_model() or [] - for entry in fb_chain: + for entry in _load_fallback_model() or []: if not isinstance(entry, dict): continue fb_provider = str(entry.get("provider") or "").strip() @@ -2553,14 +2451,11 @@ def _resolve_runtime_with_fallback( fb_kwargs: dict = {"requested": fb_provider, "target_model": fb_model} if entry.get("base_url"): fb_kwargs["explicit_base_url"] = entry["base_url"] - fb_api_key = resolve_entry_api_key(entry) - if fb_api_key: + if fb_api_key := resolve_entry_api_key(entry): fb_kwargs["explicit_api_key"] = fb_api_key runtime = resolve_runtime_provider(**fb_kwargs) - import logging logging.getLogger(__name__).warning( - "Primary auth failed (%s), falling back to %s model %s", primary_exc, - fb_provider, fb_model, + "Primary auth failed (%s), falling back to %s model %s", primary_exc, fb_provider, fb_model, ) return _RuntimeFallbackResolution(runtime, fb_model, True) except Exception: @@ -2569,16 +2464,13 @@ def _resolve_runtime_with_fallback( def _resolve_agent_model_runtime(model_override, provider_override) -> tuple[str, dict]: - """Resolve (model, runtime) for a new agent. + """(model, runtime) for a new agent; a per-session override (in-session /model switch or a + resumed row's persisted runtime) wins over global config/env. - A per-session override (prior in-session /model switch, or the persisted - runtime of a resumed row) wins over global config/env resolution. Rows - persisted before the custom-provider identity fix stored the resolved - provider "custom", which no named ``providers:`` entry matches — recover - the entry identity from the persisted base_url (falling back to the - configured provider) or the rebuild surfaces as "No LLM provider - configured". Persisted base_url/api_key/api_mode are honored only while - the original runtime is used; they must not leak into a fallback pair. + Rows persisted before the custom-provider identity fix stored the resolved provider "custom", + which no named ``providers:`` entry matches — recover the identity from the persisted base_url + or the rebuild surfaces as "No LLM provider configured". Persisted base_url/api_key/api_mode are + honored only while the original runtime is used; they must not leak into a fallback pair. """ if isinstance(model_override, dict) and model_override.get("model"): model = str(model_override.get("model") or "") @@ -2587,18 +2479,16 @@ def _resolve_agent_model_runtime(model_override, provider_override) -> tuple[str resolve_kwargs = {} if str(requested_provider or "").strip().lower() == "custom": from hermes_cli.runtime_provider import canonical_custom_identity - recovered = canonical_custom_identity(base_url=override_base_url or None, model=model or None) - if recovered: + if recovered := canonical_custom_identity(base_url=override_base_url or None, model=model or None): requested_provider = recovered if override_base_url: - # Failing identity recovery, still hand the base_url to the - # direct-alias branch so pool/env credentials resolve for it. + # Failing identity recovery, still hand the base_url to the direct-alias branch so + # pool/env credentials resolve for it. resolve_kwargs["explicit_base_url"] = override_base_url resolve_kwargs["requested"] = requested_provider resolve_kwargs["target_model"] = model or None overrides = { - "base_url": override_base_url, "api_key": model_override.get("api_key"), - "api_mode": model_override.get("api_mode"), + "base_url": override_base_url, "api_key": model_override.get("api_key"), "api_mode": model_override.get("api_mode"), } else: model, requested_provider = _resolve_startup_runtime() @@ -2614,12 +2504,35 @@ def _resolve_agent_model_runtime(model_override, provider_override) -> tuple[str if not resolution.selected_model: raise RuntimeError("Auth fallback resolved without a model") return resolution.selected_model, runtime - for k, v in overrides.items(): - if v: - runtime[k] = v + runtime.update({k: v for k, v in overrides.items() if v}) return model, runtime +def _startup_system_prompt(cfg: dict, task_id: str) -> str: + """Config ephemeral system prompt + the HERMES_TUI_SKILLS preload block. Hard-fails only when EVERY + requested skill is missing (cli.py parity): a typo'd name must not auto-block the Kanban task.""" + from hermes_cli.config import resolve_ephemeral_system_prompt_from_config + system_prompt = resolve_ephemeral_system_prompt_from_config(cfg) + startup_skills = _parse_tui_skills_env() + if not startup_skills: + return system_prompt + from agent.skill_commands import build_preloaded_skills_prompt + skills_prompt, loaded_skills, missing_skills = build_preloaded_skills_prompt(startup_skills, task_id=task_id) + if missing_skills: + missing_display = ", ".join(missing_skills) + if not loaded_skills: + raise ValueError(f"Unknown skill(s): {missing_display}") + logger.warning( + "Unknown skill(s) requested, skipping: %s. " + "Continuing with: %s. " + "List available skills with `hermes skills list`.", + missing_display, ", ".join(loaded_skills), + ) + if skills_prompt: + system_prompt = "\n\n".join(part for part in (system_prompt, skills_prompt) if part).strip() + return system_prompt + + def _make_agent( sid: str, key: str, session_id: str | None = None, session_db=None, model_override: dict | str | None = None, provider_override: str | None = None, @@ -2632,43 +2545,18 @@ def _make_agent( if synthetic is not None: return synthetic from run_agent import AIAgent - - # MCP discovery runs in a background daemon thread so a dead server can't - # freeze the shell; the agent snapshots its tool list once, so briefly - # (bounded) wait for in-flight discovery. Dashboard /api/ws uses - # hermes_cli.mcp_startup; TUI stdio keeps the tui_gateway.entry thread. + # MCP discovery runs in a background daemon thread so a dead server can't freeze the shell; the + # agent snapshots its tool list once, so briefly (bounded) wait for in-flight discovery. + # Dashboard /api/ws uses hermes_cli.mcp_startup; TUI stdio keeps the tui_gateway.entry thread. for _mod in ("hermes_cli.mcp_startup", "tui_gateway.entry"): with contextlib.suppress(Exception): importlib.import_module(_mod).wait_for_mcp_discovery() cfg = _load_cfg() - from hermes_cli.config import resolve_ephemeral_system_prompt_from_config - system_prompt = resolve_ephemeral_system_prompt_from_config(cfg) - startup_skills = _parse_tui_skills_env() - if startup_skills: - from agent.skill_commands import build_preloaded_skills_prompt - skills_prompt, loaded_skills, missing_skills = build_preloaded_skills_prompt( - startup_skills, task_id=session_id or key, - ) - if missing_skills: - missing_display = ", ".join(missing_skills) - # Hard-fail only when EVERY requested skill is missing (cli.py - # parity): a typo'd name must not auto-block the Kanban task. - if loaded_skills: - logger.warning( - "Unknown skill(s) requested, skipping: %s. " - "Continuing with: %s. " - "List available skills with `hermes skills list`.", - missing_display, - ", ".join(loaded_skills), - ) - else: - raise ValueError(f"Unknown skill(s): {missing_display}") - if skills_prompt: - system_prompt = "\n\n".join( - part for part in (system_prompt, skills_prompt) if part - ).strip() + system_prompt = _startup_system_prompt(cfg, session_id or key) model, runtime = _resolve_agent_model_runtime(model_override, provider_override) _pr = _load_provider_routing() + platform = _resolve_agent_platform(platform_override) + ignore_rules = is_truthy_value(os.environ.get("HERMES_IGNORE_RULES")) agent = AIAgent( model=model, max_iterations=_cfg_max_turns(cfg, 500), @@ -2682,16 +2570,10 @@ def _make_agent( quiet_mode=True, verbose_logging=False, # DEBUG agent logging; independent of tool_progress_mode reasoning_config=( - reasoning_config_override - if reasoning_config_override is not None - else _load_reasoning_config(str(model or "")) + reasoning_config_override if reasoning_config_override is not None else _load_reasoning_config(str(model or "")) ), - service_tier=( - service_tier_override - if service_tier_override is not None - else _load_service_tier() - ), - enabled_toolsets=_load_enabled_toolsets(_resolve_agent_platform(platform_override)), + service_tier=service_tier_override if service_tier_override is not None else _load_service_tier(), + enabled_toolsets=_load_enabled_toolsets(platform), # OpenRouter provider_routing prefs (gateway + CLI parity). providers_allowed=_pr.get("only"), providers_ignored=_pr.get("ignore"), @@ -2699,14 +2581,14 @@ def _make_agent( provider_sort=_pr.get("sort"), provider_require_parameters=_pr.get("require_parameters", False), provider_data_collection=_pr.get("data_collection"), - platform=_resolve_agent_platform(platform_override), + platform=platform, session_id=session_id or key, session_db=session_db if session_db is not None else _get_db(), ephemeral_system_prompt=system_prompt or None, checkpoints_enabled=is_truthy_value(os.environ.get("HERMES_TUI_CHECKPOINTS")), pass_session_id=is_truthy_value(os.environ.get("HERMES_TUI_PASS_SESSION_ID")), - skip_context_files=is_truthy_value(os.environ.get("HERMES_IGNORE_RULES")), - skip_memory=is_truthy_value(os.environ.get("HERMES_IGNORE_RULES")), + skip_context_files=ignore_rules, + skip_memory=ignore_rules, fallback_model=_load_fallback_model(), **_agent_cbs(sid), ) @@ -2718,61 +2600,23 @@ def _make_agent( return agent -def _init_session( - sid: str, key: str, agent, history: list, cols: int = 80, cwd: str | None = None, - session_db=None, source: str | None = None, profile_home: str | None = None, - explicit_cwd: bool = False, -): - now = time.time() - with _sessions_lock: - _sessions[sid] = { - "agent": agent, - "session_key": key, - "history": history, - "history_lock": threading.Lock(), - "history_version": 0, - "inflight_turn": None, - "created_at": now, - "last_active": now, - "running": False, - "attached_images": [], - "image_counter": 0, - "cwd": cwd or _completion_cwd(), - "explicit_cwd": bool(explicit_cwd), - "cols": cols, - "slash_worker": None, - "show_reasoning": _load_show_reasoning(), - "source": _resolve_session_source(source), - "tool_progress_mode": _load_tool_progress_mode(), - "edit_snapshots": {}, - "tool_started_at": {}, - # Profile-scoped HERMES_HOME (None = launch profile); SessionBranch - # copies the parent's so the child stays on the same state.db. - "profile_home": profile_home, - # In-session /model switch, honored on rebuild (/new, resume) so it - # never leaks into siblings via process-global env vars. - "model_override": None, - # Async events go to the transport that created the session - # (stdio for Ink, JSON-RPC WS for the dashboard sidebar). - "transport": current_transport() or _stdio_transport, - } - _session_todo_state(_sessions[sid]) - _init_owns_db = False +def _hydrate_session_cwd(sid: str, key: str, session_db, profile_home: str | None) -> None: + """Adopt the stored row's cwd, or persist the fresh session's cwd (+ schedule git meta) when the row has none.""" + owns_db = False if session_db is not None: db = session_db elif profile_home: try: db = _open_profile_session_db(profile_home) - _init_owns_db = True + owns_db = True except Exception: - # FAIL CLOSED (same class as the deferred-build bind): a named-profile - # session must never touch the launch state.db — skip cwd hydration - # (the row lands on the agent's own lazy-create once the store recovers). + # FAIL CLOSED (same class as the deferred-build bind): a named-profile session must never + # touch the launch state.db — skip cwd hydration (the row lands on the agent's own + # lazy-create once the store recovers). logger.warning( "profile session store unavailable for %s — skipping cwd " "hydration instead of touching the launch state.db", - profile_home, - exc_info=True, + profile_home, exc_info=True, ) db = None else: @@ -2792,12 +2636,38 @@ def _init_session( except Exception: logger.debug("failed to persist resumed session cwd", exc_info=True) finally: - if _init_owns_db and db is not None: + if owns_db and db is not None: with contextlib.suppress(Exception): db.close() + + +def _init_session( + sid: str, key: str, agent, history: list, cols: int = 80, cwd: str | None = None, + session_db=None, source: str | None = None, profile_home: str | None = None, + explicit_cwd: bool = False, +): + now = time.time() + with _sessions_lock: + _sessions[sid] = { + "agent": agent, "session_key": key, "history": history, "history_lock": threading.Lock(), + "history_version": 0, "inflight_turn": None, "created_at": now, "last_active": now, + "running": False, "attached_images": [], "image_counter": 0, "cwd": cwd or _completion_cwd(), + "explicit_cwd": bool(explicit_cwd), "cols": cols, "slash_worker": None, + "show_reasoning": _load_show_reasoning(), "source": _resolve_session_source(source), + "tool_progress_mode": _load_tool_progress_mode(), "edit_snapshots": {}, "tool_started_at": {}, + # Profile-scoped HERMES_HOME (None = launch profile); SessionBranch copies the parent's + # so the child stays on the same state.db. + "profile_home": profile_home, + # In-session /model switch, honored on rebuild (/new, resume) so it never leaks into + # siblings via process-global env vars. + "model_override": None, + # Async events go to the transport that created the session (stdio for Ink, WS for the dashboard). + "transport": current_transport() or _stdio_transport, + } + _session_todo_state(_sessions[sid]) + _hydrate_session_cwd(sid, key, session_db, profile_home) _register_session_cwd(_sessions[sid]) - # No eager slash-worker pre-warm (see _start_agent_build). - _wire_session_agent(sid, key, agent) + _wire_session_agent(sid, key, agent) # no eager slash-worker pre-warm (see _start_agent_build) _start_session_services(sid, key, _sessions.get(sid, {})) _emit("session.info", sid, _session_info(agent, _sessions.get(sid, {}))) _schedule_mcp_late_refresh(sid, agent) @@ -2825,11 +2695,8 @@ def _resolve_checkpoint_hash(mgr, cwd: str, ref: str) -> str: # ── Methods: session ───────────────────────────────────────────────── -def _lazy_resume_info( - cwd: str, *, model: str = "", provider: str = "", profile: str | None = None, -) -> dict: - """session.info for a not-yet-built session (the shape session.create - returns). tools/skills land later when the deferred build emits session.info.""" +def _lazy_resume_info(cwd: str, *, model: str = "", provider: str = "", profile: str | None = None) -> dict: + """session.info for a not-yet-built session (session.create's shape); tools/skills land with the deferred build.""" info = { "cwd": cwd, "branch": _git_branch_for_cwd(cwd), "project": _project_info_for_cwd(cwd), "model": model or _resolve_model(), "tools": {}, "skills": {}, "lazy": True, @@ -2848,8 +2715,7 @@ def _deferred_session_record( resume_runtime_overrides: dict | None = None, todo_state: dict | None = None, explicit_cwd: bool = False, ) -> dict: - """A live-session record whose AIAgent is built later (lazy watch / cold - resume) — _init_session's shape minus the agent.""" + """A live-session record whose AIAgent is built later (lazy watch / cold resume) — _init_session's shape minus the agent.""" now = time.time() return { "agent": None, "agent_error": None, "agent_ready": threading.Event(), "attached_images": [], @@ -2872,63 +2738,44 @@ _ANY_PROFILE = object() # default: match a live session regardless of profile def _live_profile_matches(session: dict, profile_home) -> bool: - """True when ``session`` belongs to ``profile_home`` (None = launch profile). - - Same string compare as session.resume's ``_find_live_unpersisted``: a - record with no ``profile_home`` is the launch profile's. ``_ANY_PROFILE`` - disables the check for callers that have no profile to scope by. - """ + """True when ``session`` belongs to ``profile_home`` (None = launch profile; a record with no + ``profile_home`` is the launch profile's). ``_ANY_PROFILE`` disables the check.""" if profile_home is _ANY_PROFILE: return True - want = str(profile_home) if profile_home else None - return (session.get("profile_home") or None) == want + return (session.get("profile_home") or None) == (str(profile_home) if profile_home else None) -def _claim_or_reuse_live( - sid: str, session_key: str, record: dict, lease -) -> tuple[str, dict] | None: - """Register ``record`` as the live session for ``session_key`` under the - resume lock, or — if a concurrent resume already won — release ``lease`` and - return the winner for the caller to reuse.""" - # The record carries the home this resume resolved; a live runtime of the - # same stored id under ANOTHER profile is not a winner to reuse (#100029). +def _claim_or_reuse_live(sid: str, session_key: str, record: dict, lease) -> tuple[str, dict] | None: + """Register ``record`` as the live session for ``session_key`` under the resume lock, or — if a + concurrent resume already won — release ``lease`` and return the winner for the caller to reuse.""" + # The record carries the home this resume resolved; a live runtime of the same stored id under + # ANOTHER profile is not a winner to reuse. profile_home = record.get("profile_home") with _session_resume_lock: live = _find_live_session_by_key(session_key, profile_home) if live is not None: if lease is not None: lease.release() - # The winner is being reattached by this resume: any pending - # ws-orphan reap for it must not fire against the reclaimed - # client (storm killer — see _cancel_ws_orphan_reap). + # The winner is being reattached by this resume: a pending ws-orphan reap must not + # fire against the reclaimed client (storm killer — see _cancel_ws_orphan_reap). _cancel_ws_orphan_reap(live[0]) return live with _sessions_lock: _sessions[sid] = record _register_session_cwd(_sessions[sid]) - # A PRIOR runtime for this stored id may still be sentinel-parked with - # a reap Timer armed; cancel + finalize it quietly so the reap doesn't - # broadcast session.reclaimed for a just-re-resumed session (storm). + # A PRIOR runtime for this stored id may still be sentinel-parked with a reap Timer armed; + # cancel + finalize it quietly so the reap doesn't broadcast session.reclaimed (storm). _cancel_ws_orphan_reap(sid) stale = _claim_parked_runtimes(session_key, keep_sid=sid, profile_home=profile_home) - # Slow finalization work stays OUTSIDE _session_resume_lock (see - # _pop_session_by_id) — the stale records are already claimed above. - _finalize_superseded_runtimes(stale) + _finalize_superseded_runtimes(stale) # slow finalization stays OUTSIDE _session_resume_lock return None -def _claim_parked_runtimes( - session_key: str, *, keep_sid: str, profile_home=_ANY_PROFILE -) -> list[tuple[str, dict]]: - """Claim sentinel-parked stale runtimes of ``session_key`` for supersession. - - When a resume mints a fresh runtime for stored session id ``session_key``, - any older runtime record for the same stored id that is still parked on - the detached-WS sentinel is superseded: its pending orphan-reap Timer is - cancelled and the record is atomically popped from ``_sessions`` here - (under the caller's _session_resume_lock), then finalized by - :func:`_finalize_superseded_runtimes` after the lock is released. - """ +def _claim_parked_runtimes(session_key: str, *, keep_sid: str, profile_home=_ANY_PROFILE) -> list[tuple[str, dict]]: + """Claim sentinel-parked stale runtimes of ``session_key`` for supersession: older records for the + same stored id still parked on the detached-WS sentinel get their orphan-reap Timer cancelled and + are popped from ``_sessions`` here (under the caller's _session_resume_lock), then finalized by + :func:`_finalize_superseded_runtimes` after the lock is released.""" stale: list[tuple[str, dict]] = [] with _sessions_lock: candidates = [ @@ -2949,15 +2796,10 @@ def _claim_parked_runtimes( def _finalize_superseded_runtimes(stale: list[tuple[str, dict]]) -> None: - """Quietly finalize runtimes claimed by :func:`_claim_parked_runtimes`. - - Ends them with end_reason ``superseded_by_resume`` — deliberately NOT in - _RECLAIM_END_REASONS, so no ``session.reclaimed`` broadcast fires (that - broadcast triggers client auto-re-resume and fed the - reap->broadcast->resume feedback loop). ``superseded_by_resume`` IS in - hermes_state_common._RECOVERABLE_END_REASONS so canonical Bot Chat - resurrection still applies to the stored session. - """ + """Quietly finalize runtimes claimed by :func:`_claim_parked_runtimes` with end_reason + ``superseded_by_resume`` — deliberately NOT in _RECLAIM_END_REASONS (no ``session.reclaimed`` + broadcast, which fed the reap->broadcast->resume loop) but IN _RECOVERABLE_END_REASONS so + canonical Bot Chat resurrection still applies.""" for old_sid, popped in stale: try: _teardown_popped_session(popped, end_reason="superseded_by_resume") @@ -2966,8 +2808,7 @@ def _finalize_superseded_runtimes(stale: list[tuple[str, dict]]) -> None: def _schedule_agent_build(sid: str, delay: float = 0.05) -> None: - """Pre-warm a deferred session's agent off the response path (session.create - and cold resume both build through here; _sess() also builds on demand).""" + """Pre-warm a deferred session's agent off the response path (session.create + cold resume; _sess() also builds on demand).""" def _run(): session = _sessions.get(sid)