refactor(tui): server.py — _startup_system_prompt/_hydrate_session_cwd/_effective_approval_state/_turn_started_at helpers, _tui_notice, compact session-info/toolset/runtime docstrings
This commit is contained in:
+230
-389
@@ -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)
|
||||
|
||||
Reference in New Issue
Block a user