refactor(tui_gateway): entry discovery spawn helper, row/kwargs table tightening
- entry: _spawn_discovery + logged _mcp_startup_call replace the two duplicated start_background_mcp_discovery try/except blocks (same log levels + messages); docstrings/comments compacted keeping every signal/thread WHY. - methods_projects/agent_callbacks: _project_tree_row and preview kwargs tightened. Goldens identical; WIRE-PARITY-OK.
This commit is contained in:
@@ -319,9 +319,8 @@ def _background_agent_kwargs(agent, task_id: str) -> dict:
|
||||
|
||||
|
||||
def _ephemeral_preview_agent_kwargs(agent, task_id: str) -> dict:
|
||||
kwargs = _background_agent_kwargs(agent, task_id)
|
||||
kwargs.update({"enabled_toolsets": ["terminal", "file"], "session_db": None, "skip_memory": True})
|
||||
return kwargs
|
||||
return {**_background_agent_kwargs(agent, task_id),
|
||||
"enabled_toolsets": ["terminal", "file"], "session_db": None, "skip_memory": True}
|
||||
|
||||
|
||||
_PREVIEW_HISTORY_ROLES = ("user", "assistant", "tool", "system")
|
||||
|
||||
@@ -36,11 +36,8 @@ _mcp_discovery_enabled = False
|
||||
|
||||
|
||||
def _install_sidecar_publisher() -> None:
|
||||
"""Mirror every dispatcher emit to the dashboard sidebar via WS.
|
||||
|
||||
Activated by `HERMES_TUI_SIDECAR_URL` (set by the dashboard's ``/api/pty``
|
||||
endpoint). Best-effort: connect failure or runtime drop falls back to stdio-only.
|
||||
"""
|
||||
"""Mirror every dispatcher emit to the dashboard sidebar via WS when
|
||||
`HERMES_TUI_SIDECAR_URL` is set (best-effort: a dropped WS falls back to stdio-only)."""
|
||||
url = os.environ.get("HERMES_TUI_SIDECAR_URL")
|
||||
if not url:
|
||||
return
|
||||
@@ -63,15 +60,23 @@ def _stamp() -> str:
|
||||
return time.strftime("%Y-%m-%d %H:%M:%S")
|
||||
|
||||
|
||||
def _mcp_startup_call(name: str, *args, default=None, **kwargs):
|
||||
"""Call ``hermes_cli.mcp_startup.<name>`` (lazy import); ``default`` on any failure."""
|
||||
def _mcp_startup_call(name: str, *args, default=None, log=None, **kwargs):
|
||||
"""Call ``hermes_cli.mcp_startup.<name>`` (lazy import); ``default`` on any failure,
|
||||
optionally logged as ``(level, message)``."""
|
||||
try:
|
||||
from hermes_cli import mcp_startup
|
||||
return getattr(mcp_startup, name)(*args, **kwargs)
|
||||
except Exception:
|
||||
if log:
|
||||
getattr(logger, log[0])(log[1], exc_info=True)
|
||||
return default
|
||||
|
||||
|
||||
def _spawn_discovery(log: tuple) -> None:
|
||||
_mcp_startup_call(
|
||||
"start_background_mcp_discovery", logger=logger, thread_name="tui-mcp-discovery", log=log)
|
||||
|
||||
|
||||
def _append_crash_log(header: str, dump=None) -> None:
|
||||
"""Best-effort ``=== header ===`` entry in the crash log; ``dump(f)`` adds detail."""
|
||||
with suppress(Exception):
|
||||
@@ -84,11 +89,9 @@ def _append_crash_log(header: str, dump=None) -> None:
|
||||
|
||||
def _log_signal(signum: int, frame) -> None:
|
||||
"""Capture WHICH thread and WHERE a termination signal hit us, then exit.
|
||||
|
||||
``sys.exit(0)`` alone raced the worker pool — a thread holding ``_stdout_lock``
|
||||
mid-flush blocks interpreter shutdown indefinitely — so log all thread stacks,
|
||||
give the configured grace to drain, then ``os._exit(0)``.
|
||||
"""
|
||||
``sys.exit(0)`` alone raced the worker pool (a thread holding ``_stdout_lock``
|
||||
mid-flush blocks interpreter shutdown), so: log all thread stacks, give the
|
||||
configured grace to drain, then ``os._exit(0)``."""
|
||||
# SIGPIPE/SIGHUP don't exist on Windows — only look up attributes present.
|
||||
names = {int(sig): attr for attr in ("SIGPIPE", "SIGTERM", "SIGHUP", "SIGINT", "SIGBREAK")
|
||||
if (sig := getattr(signal, attr, None)) is not None}
|
||||
@@ -111,26 +114,21 @@ def _log_signal(signum: int, frame) -> None:
|
||||
timer = threading.Timer(_shutdown_grace_seconds(), lambda: os._exit(0))
|
||||
timer.daemon = True
|
||||
timer.start()
|
||||
# The atexit handler (_shutdown_sessions) can be blocked past the grace window
|
||||
# by a worker holding the GIL/_stdout_lock; finalize explicitly so unpersisted
|
||||
# messages reach state.db before the hard-exit timer fires.
|
||||
# atexit (_shutdown_sessions) can be blocked past the grace window by a worker
|
||||
# holding the GIL/_stdout_lock; finalize explicitly so unpersisted messages reach
|
||||
# state.db before the hard-exit timer fires.
|
||||
with suppress(Exception):
|
||||
from tui_gateway.server import _shutdown_sessions
|
||||
_shutdown_sessions()
|
||||
|
||||
# Unwind the main thread so atexit + finalisers run inside the grace window;
|
||||
# the daemon timer is the safety net if that unwind hangs.
|
||||
sys.exit(0)
|
||||
|
||||
|
||||
def _install_signal(signame, handler):
|
||||
"""Install a signal handler if legal in this thread and platform.
|
||||
|
||||
signal.signal() raises ValueError outside the main thread; skip silently so a
|
||||
worker-thread first import (Desktop build path: server._build imports entry)
|
||||
doesn't abort. Handlers are process-global, so any main-thread import installs
|
||||
them for everyone. Missing signals (Windows: SIGPIPE/SIGHUP) are skipped too.
|
||||
"""
|
||||
"""Install a signal handler if legal here: signal.signal() raises off the main
|
||||
thread (Desktop build path: server._build imports entry from a worker), and
|
||||
Windows lacks SIGPIPE/SIGHUP — both are skipped. Handlers are process-global."""
|
||||
sig = getattr(signal, signame, None)
|
||||
if sig is None or threading.current_thread() is not threading.main_thread():
|
||||
return
|
||||
@@ -139,12 +137,10 @@ def _install_signal(signame, handler):
|
||||
signal.signal(sig, handler)
|
||||
|
||||
|
||||
# SIGPIPE: ignore, don't exit. SIG_DFL killed the process silently whenever a
|
||||
# *background* thread (TTS, beep, voice status) wrote to a pipe the TUI had gone
|
||||
# quiet on. Ignoring lets the write raise BrokenPipeError (write_json handles it with
|
||||
# a clean sys.exit(0) + _log_exit) so the gateway lives as long as the command pipe
|
||||
# is readable. Terminal signals route through _log_signal so kills/hangups are
|
||||
# diagnosable; SIGBREAK (Windows Ctrl+Break) is the weaker SIGHUP.
|
||||
# SIGPIPE: ignore, don't exit — SIG_DFL killed the process silently whenever a
|
||||
# *background* thread (TTS, beep) wrote to a pipe the TUI had gone quiet on; ignoring
|
||||
# lets write_json see BrokenPipeError and exit cleanly via _log_exit. Terminal signals
|
||||
# route through _log_signal so kills/hangups are diagnosable (SIGBREAK = Windows SIGHUP).
|
||||
_install_signal("SIGPIPE", signal.SIG_IGN)
|
||||
_install_signal("SIGTERM", _log_signal)
|
||||
if hasattr(signal, "SIGHUP"):
|
||||
@@ -155,39 +151,29 @@ _install_signal("SIGINT", signal.SIG_IGN)
|
||||
|
||||
|
||||
def _log_exit(reason: str) -> None:
|
||||
"""Record why the gateway is shutting down: every exit path collapses into a
|
||||
silent sys.exit(0), and without this trail the TUI shows "gateway exited" with
|
||||
no clue WHICH broken pipe or message triggered it."""
|
||||
"""Record why the gateway exits: every path collapses into a silent sys.exit(0),
|
||||
and without this trail the TUI can't tell WHICH broken pipe triggered it."""
|
||||
_append_crash_log(f"gateway exit · {_stamp()} · reason={reason}")
|
||||
print(f"[gateway-exit] {reason}", file=sys.stderr, flush=True)
|
||||
|
||||
|
||||
def wait_for_mcp_discovery(timeout: "float | None" = None) -> None:
|
||||
"""Block until background MCP discovery finishes, up to the resolved bound.
|
||||
|
||||
The agent snapshots its tool list ONCE at build time, so a bounded join before
|
||||
the first build lets already-spawning servers land (no-MCP startups pay ~0s)
|
||||
without re-introducing the startup hang. Bound: ``mcp_discovery_timeout`` from
|
||||
config; ``timeout`` overrides it.
|
||||
"""
|
||||
"""Block until background MCP discovery finishes, up to the resolved bound
|
||||
(``mcp_discovery_timeout`` from config; ``timeout`` overrides). The agent snapshots
|
||||
its tool list ONCE at build time, so this bounded join lets already-spawning
|
||||
servers land without re-introducing the startup hang."""
|
||||
thread = _mcp_discovery_thread
|
||||
if thread is not None and thread.is_alive():
|
||||
fallback = timeout if timeout is not None else 0.75
|
||||
bound = _mcp_startup_call("_resolve_discovery_timeout", timeout, default=fallback)
|
||||
thread.join(timeout=bound)
|
||||
return
|
||||
# Shared-owner path. Re-invoke the idempotent spawn first so a previous
|
||||
# zero-connected run gets its retry instead of latching the process MCP-less.
|
||||
# It runs under the CALLER's profile context (agent build binds the session
|
||||
# profile's HERMES_HOME first), so a launch profile without mcp_servers doesn't
|
||||
# starve selected profiles. Gated so non-MCP sessions skip the mcp_tool import.
|
||||
# Shared-owner path: re-invoke the idempotent spawn first so a zero-connected run
|
||||
# gets its retry instead of latching the process MCP-less. Runs under the CALLER's
|
||||
# profile context (agent build binds the session profile's HERMES_HOME first).
|
||||
if not _mcp_discovery_enabled:
|
||||
return
|
||||
try:
|
||||
from hermes_cli.mcp_startup import start_background_mcp_discovery
|
||||
start_background_mcp_discovery(logger=logger, thread_name="tui-mcp-discovery")
|
||||
except Exception:
|
||||
logger.debug("TUI MCP discovery retry-spawn failed", exc_info=True)
|
||||
_spawn_discovery(("debug", "TUI MCP discovery retry-spawn failed"))
|
||||
_mcp_startup_call("wait_for_mcp_discovery", timeout)
|
||||
|
||||
|
||||
@@ -226,24 +212,16 @@ def _has_configured_mcp_servers() -> bool:
|
||||
|
||||
def ensure_mcp_discovery_started() -> None:
|
||||
"""Start background MCP discovery for the current profile context, once.
|
||||
|
||||
``main()`` calls this for stdio; WS/Desktop skip ``main()``, so
|
||||
``server._start_agent_build`` also calls it AFTER binding the session profile's
|
||||
HERMES_HOME (the shared owner captures that override, so discovery reads the
|
||||
SELECTED profile's ``mcp_servers``). Delegating keeps the process-wide start lock,
|
||||
retry-after-zero-connected allowance and interactive-OAuth suppression. MCP
|
||||
registration is process-global: the FIRST profile to build an agent wins.
|
||||
"""
|
||||
HERMES_HOME so discovery reads the SELECTED profile's ``mcp_servers``. MCP
|
||||
registration is process-global: the FIRST profile to build an agent wins."""
|
||||
global _mcp_discovery_enabled
|
||||
|
||||
if not _has_configured_mcp_servers():
|
||||
return
|
||||
_mcp_discovery_enabled = True
|
||||
try:
|
||||
from hermes_cli.mcp_startup import start_background_mcp_discovery
|
||||
start_background_mcp_discovery(logger=logger, thread_name="tui-mcp-discovery")
|
||||
except Exception:
|
||||
logger.warning("Background MCP tool discovery failed to start", exc_info=True)
|
||||
_spawn_discovery(("warning", "Background MCP tool discovery failed to start"))
|
||||
|
||||
|
||||
def _write_or_exit(payload: dict, reason: str) -> None:
|
||||
@@ -255,9 +233,8 @@ def _write_or_exit(payload: dict, reason: str) -> None:
|
||||
def main():
|
||||
_install_sidecar_publisher()
|
||||
|
||||
# Heartbeat row lets the orphan sweep tell "live but idle backend" from "truly
|
||||
# orphaned"; must run BEFORE the sweep so it sees our row. The sweep itself is
|
||||
# once-per-process and config-gated (the handle_ws call site becomes a no-op).
|
||||
# Heartbeat row lets the orphan sweep tell "live but idle" from "truly orphaned";
|
||||
# it must run BEFORE the sweep. The sweep is once-per-process and config-gated.
|
||||
for start, what in (
|
||||
(server._start_backend_heartbeat_refresher, "backend heartbeat refresher start"),
|
||||
(server._schedule_startup_orphan_sweep, "startup orphan sweep scheduling"),
|
||||
@@ -267,9 +244,9 @@ def main():
|
||||
except Exception:
|
||||
logger.warning("%s failed", what, exc_info=True)
|
||||
|
||||
# Backgrounded so a dead MCP server (~7s of connect retries) can't freeze
|
||||
# startup; _make_agent briefly joins it (wait_for_mcp_discovery). The config
|
||||
# gate inside keeps the ~200ms MCP SDK import off the no-mcp_servers path.
|
||||
# Backgrounded so a dead MCP server (~7s of retries) can't freeze startup;
|
||||
# _make_agent briefly joins it. The config gate keeps the MCP SDK import off
|
||||
# the no-mcp_servers path.
|
||||
ensure_mcp_discovery_started()
|
||||
|
||||
_write_or_exit({
|
||||
@@ -288,9 +265,8 @@ def main():
|
||||
# Live-apply skins Hermes activates mid-conversation.
|
||||
server._ensure_skin_watcher()
|
||||
|
||||
# Warm the /model picker's provider-models cache during this idle window
|
||||
# (mirrors the classic CLI loop); otherwise the first /model open blocks on
|
||||
# serial /v1/models fetches. Fire-and-forget, once-per-process.
|
||||
# Warm the /model picker's provider-models cache in this idle window, else the
|
||||
# first /model open blocks on serial /v1/models fetches. Fire-and-forget.
|
||||
try:
|
||||
from hermes_cli.model_switch import prewarm_picker_cache_async
|
||||
prewarm_picker_cache_async()
|
||||
|
||||
@@ -348,10 +348,9 @@ def _project_tree_row(r: dict) -> dict:
|
||||
source=r.get("source"), archived=bool(r.get("archived")))
|
||||
row.update({k: r.get(k) or 0 for k in (
|
||||
"message_count", "tool_call_count", "input_tokens", "output_tokens")})
|
||||
row.update(
|
||||
actual_cost_usd=r.get("actual_cost_usd"), estimated_cost_usd=r.get("estimated_cost_usd"),
|
||||
model=r.get("model"), is_active=False, cwd=r.get("cwd"), git_branch=r.get("git_branch"),
|
||||
git_repo_root=r.get("git_repo_root"))
|
||||
row.update({k: r.get(k) for k in ("actual_cost_usd", "estimated_cost_usd", "model")})
|
||||
row["is_active"] = False
|
||||
row.update({k: r.get(k) for k in ("cwd", "git_branch", "git_repo_root")})
|
||||
return row
|
||||
|
||||
|
||||
@@ -361,17 +360,12 @@ def _project_tree_inputs(
|
||||
"""Gather (sessions, projects, discovered_repos, active_id) for build_tree.
|
||||
``include_discovered`` is the zero-session-repo overview tier; drill-in skips it,
|
||||
avoiding the distinct-cwd scan + git probes on that per-turn path."""
|
||||
# compact_rows: `_project_tree_row` drops the system-prompt blob; selecting it
|
||||
# only to discard it costs tens of MB of B-tree reads per build on a big DB.
|
||||
rows = db.list_sessions_rich(
|
||||
limit=session_limit,
|
||||
offset=0,
|
||||
order_by_last_active=True,
|
||||
min_message_count=1,
|
||||
include_children=False,
|
||||
exclude_sources=_PROJECT_TREE_EXCLUDED_SOURCES,
|
||||
include_archived=False,
|
||||
# `_project_tree_row` drops the system-prompt blob; selecting it only to
|
||||
# discard it costs tens of MB of B-tree reads per build on a big DB.
|
||||
compact_rows=True)
|
||||
limit=session_limit, offset=0, order_by_last_active=True, min_message_count=1,
|
||||
include_children=False, exclude_sources=_PROJECT_TREE_EXCLUDED_SOURCES,
|
||||
include_archived=False, compact_rows=True)
|
||||
sessions = [_project_tree_row(r) for r in rows]
|
||||
# Parallel-warm the git cache so build_tree's resolver reads it instead of
|
||||
# cold-probing each cwd in sequence (matters on the drill-in path).
|
||||
@@ -386,10 +380,10 @@ def _project_tree_inputs(
|
||||
projects = [p.to_dict() for p in pdb.list_projects(conn)]
|
||||
active_id = pdb.get_active_id(conn)
|
||||
# backfill stays off the hot tree path — grouping uses the live resolver.
|
||||
discovered = (
|
||||
_discover_repos_payload(db, conn=conn, backfill=False, include_cached=policy["enabled"])
|
||||
if include_discovered
|
||||
else [])
|
||||
discovered = []
|
||||
if include_discovered:
|
||||
discovered = _discover_repos_payload(
|
||||
db, conn=conn, backfill=False, include_cached=policy["enabled"])
|
||||
return sessions, projects, discovered, active_id
|
||||
|
||||
|
||||
|
||||
Reference in New Issue
Block a user