diff --git a/tui_gateway/methods_complete_helpers.py b/tui_gateway/methods_complete_helpers.py new file mode 100644 index 0000000000..f11e0827dd --- /dev/null +++ b/tui_gateway/methods_complete_helpers.py @@ -0,0 +1,345 @@ +"""Completion helpers (@-mention / path fuzzy ranking, repo file listing) for the complete.* RPCs. + +Bodies are rebound onto server.py's globals at install time (see +method_ctx.bind_module), so they reference server.py globals bare. +""" + +from __future__ import annotations + +import threading + +from .method_ctx import HandlerRegistry, bind_module + +_registry = HandlerRegistry() + + +# ── Methods: complete ───────────────────────────────────────────────── + +_FUZZY_CACHE_TTL_S = 5.0 +_FUZZY_CACHE_MAX_FILES = 20000 +_FUZZY_FALLBACK_EXCLUDES = frozenset( + { + ".git", + ".hg", + ".svn", + ".next", + ".cache", + ".venv", + "venv", + "node_modules", + "__pycache__", + "dist", + "build", + "target", + ".mypy_cache", + ".pytest_cache", + ".ruff_cache", + } +) +_fuzzy_cache_lock = threading.Lock() +_fuzzy_cache: dict[str, tuple[float, list[str]]] = {} + + +def _list_repo_files(root: str) -> list[str]: + """Return file paths relative to ``root``. + + Uses ``git ls-files`` from the repo top (resolved via + ``rev-parse --show-toplevel``) so the listing covers tracked + untracked + files anywhere in the repo, then converts each path back to be relative + to ``root``. Files outside ``root`` (parent directories of cwd, sibling + subtrees) are excluded so the picker stays scoped to what's reachable + from the gateway's cwd. Falls back to a bounded ``os.walk(root)`` when + ``root`` isn't inside a git repo. Result cached per-root for + ``_FUZZY_CACHE_TTL_S`` so rapid keystrokes don't respawn git processes. + """ + now = time.monotonic() + with _fuzzy_cache_lock: + cached = _fuzzy_cache.get(root) + if cached and now - cached[0] < _FUZZY_CACHE_TTL_S: + return cached[1] + + files: list[str] = [] + from hermes_cli._subprocess_compat import windows_hide_flags + + _creationflags = windows_hide_flags() + try: + top_result = subprocess.run( + ["git", "-C", root, "rev-parse", "--show-toplevel"], + capture_output=True, + timeout=2.0, + check=False, + stdin=subprocess.DEVNULL, + creationflags=_creationflags, + ) + if top_result.returncode == 0: + top = top_result.stdout.decode("utf-8", "replace").strip() + list_result = subprocess.run( + [ + "git", + "-C", + top, + "ls-files", + "-z", + "--cached", + "--others", + "--exclude-standard", + ], + capture_output=True, + timeout=2.0, + check=False, + stdin=subprocess.DEVNULL, + creationflags=_creationflags, + ) + if list_result.returncode == 0: + for p in list_result.stdout.decode("utf-8", "replace").split("\0"): + if not p: + continue + rel = os.path.relpath(os.path.join(top, p), root).replace( + os.sep, "/" + ) + # Skip parents/siblings of cwd — keep the picker scoped + # to root-and-below, matching Cmd-P workspace semantics. + if rel.startswith("../"): + continue + files.append(rel) + if len(files) >= _FUZZY_CACHE_MAX_FILES: + break + except (OSError, subprocess.TimeoutExpired): + pass + + if not files: + # Fallback walk: skip vendor/build dirs + dot-dirs so the walk stays + # tractable. Dotfiles themselves survive — the ranker decides based + # on whether the query starts with `.`. + try: + for dirpath, dirnames, filenames in os.walk(root, followlinks=False): + dirnames[:] = [ + d + for d in dirnames + if d not in _FUZZY_FALLBACK_EXCLUDES and not d.startswith(".") + ] + rel_dir = os.path.relpath(dirpath, root) + for f in filenames: + rel = f if rel_dir == "." else f"{rel_dir}/{f}" + files.append(rel.replace(os.sep, "/")) + if len(files) >= _FUZZY_CACHE_MAX_FILES: + break + if len(files) >= _FUZZY_CACHE_MAX_FILES: + break + except OSError: + pass + + with _fuzzy_cache_lock: + _fuzzy_cache[root] = (now, files) + + return files + + +def _fuzzy_basename_rank(name: str, query: str) -> tuple[int, int] | None: + """Rank ``name`` against ``query``; lower is better. Returns None to reject. + + Tiers (kind): + 0 — exact basename + 1 — basename prefix (e.g. `app` → `appChrome.tsx`) + 2 — word-boundary / camelCase hit (e.g. `chrome` → `appChrome.tsx`) + 3 — substring anywhere in basename + 4 — subsequence match (every query char appears in order) + + Secondary key is `len(name)` so shorter names win ties. + """ + if not query: + return (3, len(name)) + + nl = name.lower() + ql = query.lower() + + if nl == ql: + return (0, len(name)) + + if nl.startswith(ql): + return (1, len(name)) + + # Word-boundary split: `foo-bar_baz.qux` → ["foo","bar","baz","qux"]. + # camelCase split: `appChrome` → ["app","Chrome"]. Cheap approximation; + # falls through to substring/subsequence if it misses. + parts: list[str] = [] + buf = "" + for ch in name: + if ch in "-_." or (ch.isupper() and buf and not buf[-1].isupper()): + if buf: + parts.append(buf) + buf = ch if ch not in "-_." else "" + else: + buf += ch + if buf: + parts.append(buf) + for p in parts: + if p.lower().startswith(ql): + return (2, len(name)) + + if ql in nl: + return (3, len(name)) + + i = 0 + for ch in nl: + if ch == ql[i]: + i += 1 + if i == len(ql): + return (4, len(name)) + + return None + + +def _abs_completion_prefix_exists(path_part: str) -> bool: + """True when ``path_part`` reads sensibly as an absolute path. + + A leading `/` is only meant literally if something is actually there: + the parent directory has to exist, and a partially-typed final segment + has to match at least one of its entries. Used to decide whether + `@/foo` is the absolute `/foo` or shorthand for `foo` under the cwd. + """ + expanded = _normalize_completion_path(path_part) + parent = os.path.dirname(expanded.rstrip("/")) or "/" + tail = os.path.basename(expanded.rstrip("/")) + + if not os.path.isdir(parent): + return False + + if not tail or expanded.endswith("/"): + return os.path.isdir(expanded) or expanded == "/" + + try: + tail_lower = tail.lower() + return any(e.lower().startswith(tail_lower) for e in os.listdir(parent)) + except OSError: + return False + + +def _details_completion_item(value: str, meta: str = "") -> dict: + return {"text": value, "display": value, "meta": meta} + + +def _details_root_completion_item( + value: str, meta: str, needs_leading_space: bool +) -> dict: + return _details_completion_item( + f" {value}" if needs_leading_space else value, + meta, + ) + + +def _details_completions(text: str) -> list[dict] | None: + if not text.lower().startswith("/details"): + return None + + stripped = text.strip() + if stripped and not "/details".startswith(stripped.lower().split()[0]): + return None + + body = text[len("/details") :] + if body.startswith(" "): + body = body[1:] + parts = body.split() + has_trailing_space = text.endswith(" ") + sections = ("thinking", "tools", "subagents", "activity") + modes = ("hidden", "collapsed", "expanded") + + if not body or (len(parts) == 0 and has_trailing_space): + return [ + *[ + _details_root_completion_item( + mode, "global mode", not has_trailing_space + ) + for mode in modes + ], + _details_root_completion_item( + "cycle", "cycle global mode", not has_trailing_space + ), + *[ + _details_root_completion_item( + section, "section override", not has_trailing_space + ) + for section in sections + ], + ] + + if len(parts) == 1 and not has_trailing_space: + prefix = parts[0].lower() + candidates = [*modes, "cycle", *sections] + return [ + _details_completion_item( + candidate, + ( + "section override" + if candidate in sections + else "cycle global mode" if candidate == "cycle" else "global mode" + ), + ) + for candidate in candidates + if candidate.startswith(prefix) and candidate != prefix + ] + + if len(parts) == 1 and has_trailing_space and parts[0].lower() in sections: + return [ + *[ + _details_completion_item(mode, f"set {parts[0].lower()}") + for mode in modes + ], + _details_completion_item("reset", f"clear {parts[0].lower()} override"), + ] + + if len(parts) == 2 and not has_trailing_space and parts[0].lower() in sections: + prefix = parts[1].lower() + return [ + _details_completion_item( + candidate, + ( + f"clear {parts[0].lower()} override" + if candidate == "reset" + else f"set {parts[0].lower()}" + ), + ) + for candidate in (*modes, "reset") + if candidate.startswith(prefix) and candidate != prefix + ] + + return [] + + +def _model_picker_context(agent): + """Layer live session state onto config without losing custom identity.""" + from hermes_cli.inventory import load_picker_context + + ctx = load_picker_context() + provider = getattr(agent, "provider", "") if agent else "" + base_url = getattr(agent, "base_url", "") if agent else "" + if str(provider or "").strip().lower() == "custom": + try: + from hermes_cli.runtime_provider import canonical_custom_identity + + provider = ( + canonical_custom_identity( + base_url=base_url or None, + config_provider=ctx.current_provider, + model=(getattr(agent, "model", "") if agent else "") + or None, + ) + or provider + ) + except Exception: + logger.debug( + "custom provider identity recovery failed (model picker)", + exc_info=True, + ) + + return ctx.with_overrides( + current_provider=provider, + current_model=(getattr(agent, "model", "") if agent else "") + or _resolve_model(), + current_base_url=base_url, + ) + + +def register(server) -> None: + """Publish this module's helpers + handlers onto ``server``, rebound to its globals.""" + bind_module(globals(), server, skip=("_",)) diff --git a/tui_gateway/methods_session.py b/tui_gateway/methods_session.py index 497216d188..5f52c63905 100644 --- a/tui_gateway/methods_session.py +++ b/tui_gateway/methods_session.py @@ -32,6 +32,24 @@ def _with_session(fn): return handler +def _with_session_db(code: int): + """:func:`_with_session` plus the session's db as a 4th arg (``_db_unavailable_error(code)`` when None).""" + + def deco(fn): + def handler(rid, params: dict) -> dict: + session, err = _sess_nowait(params, rid) + if err: + return err + with _session_db(session) as db: + if db is None: + return _db_unavailable_error(rid, code=code) + return fn(rid, params, session, db) + + return handler + + return deco + + def _with_live_session(fn): """Like :func:`_with_session` but via ``_sess`` (waits for the agent build).""" @@ -46,10 +64,7 @@ def _with_live_session(fn): def _new_runtime_ids(params: dict) -> tuple[str, str]: """Fresh runtime sid + resolved DB ``source`` for a session minted from ``params``.""" - return ( - uuid.uuid4().hex[:8], - _resolve_session_source(str(params.get("source") or "").strip() or None), - ) + return (uuid.uuid4().hex[:8], _resolve_session_source(str(params.get("source") or "").strip() or None)) @contextlib.contextmanager @@ -125,7 +140,6 @@ def _pet_display_cfg() -> dict: """``display.pet`` config block, ``{}`` when config is unreadable.""" try: from hermes_cli.config import load_config - cfg = load_config() display = cfg.get("display", {}) if isinstance(cfg.get("display"), dict) else {} return display.get("pet", {}) if isinstance(display.get("pet"), dict) else {} @@ -162,7 +176,6 @@ def _billing_call(rid, fn, extra: dict | None = None) -> dict: TUI reuses on retry); the success payload is whatever ``fn`` returns. """ from hermes_cli.nous_billing import BillingError - try: return _ok(rid, fn()) except BillingError as exc: @@ -178,10 +191,8 @@ def _billing_invalid(rid, message: str, error: str = "invalid_request") -> dict: # ── session.create / list / most_recent / facts ────────────────────── -def _create_branch_rows( - db, new_key: str, parent_key: str, title: str, history: list, *, source, cwd, profile_name, copy_fields=() -) -> None: - """Create a branch child row + copy the parent transcript in bounded-chunk transactions. +def _create_branch_row(db, new_key: str, parent_key: str, *, source, cwd, profile_name) -> None: + """Create a branch child row. ``_branched_from`` is the stable marker that keeps the branch visible in list_sessions_rich(): the TUI branch leaves the parent live (no @@ -198,6 +209,10 @@ def _create_branch_rows( cwd=cwd, profile_name=profile_name, ) + + +def _copy_branch_transcript(db, new_key: str, title: str, history: list, copy_fields=()) -> None: + """Copy the parent transcript in bounded-chunk transactions, then title the child.""" db.append_messages_batch( new_key, [ @@ -227,23 +242,21 @@ def _seed_branch_row(sid: str, key: str, parent_session_id: str, history: list, if db is None: return branch_title = _branch_title(db, parent_session_id) + _create_branch_row( + db, + key, + parent_session_id, + source=source, + cwd=_sessions[sid]["cwd"], + profile_name=(Path(profile_home).name if profile_home else None), + ) try: - _create_branch_rows( - db, - key, - parent_session_id, - branch_title, - history, - source=source, - cwd=_sessions[sid]["cwd"], - profile_name=(Path(profile_home).name if profile_home else None), - ) + _copy_branch_transcript(db, key, branch_title, history) except Exception as exc: - # Compensation: if the transcript copy / title write failed AFTER the - # row committed, a durable-but-empty row would defeat the INSERT OR + # Compensation: the row committed but the transcript copy / title + # write failed; a durable-but-empty row would defeat the INSERT OR # IGNORE first-prompt seed. Roll back just this child so it can retry. from hermes_state import is_disk_full_error - if is_disk_full_error(exc): raise try: @@ -260,6 +273,33 @@ def _seed_branch_row(sid: str, key: str, parent_session_id: str, history: list, ) +def _create_overrides(params: dict) -> tuple: + """(model_override, reasoning_override, service_tier_override) from the composer's UI state. + + PER-SESSION overrides only — never a global config write, so picking a model + for a new chat can't mutate the profile default. provider is optional (resolved + at build). ``fast`` presence is the contract: omitted inherits the profile, true + pins priority, false pins normal ("" — _make_agent uses None for inheritance). + """ + create_model = str(params.get("model") or "").strip() + model_override = ( + {"model": create_model, "provider": str(params.get("provider") or "").strip() or None} + if create_model + else None + ) + reasoning_override = None + if effort := str(params.get("reasoning_effort") or "").strip(): + try: + from hermes_constants import parse_reasoning_effort + reasoning_override = parse_reasoning_effort(effort) + except Exception: + reasoning_override = None + service_tier_override = None + if "fast" in params: + service_tier_override = "priority" if is_truthy_value(params.get("fast")) else "" + return model_override, reasoning_override, service_tier_override + + @method("session.create") def _(rid, params: dict) -> dict: sid = uuid.uuid4().hex[:8] @@ -286,27 +326,7 @@ def _(rid, params: dict) -> dict: profile = (params.get("profile") or "").strip() or None profile_home = _profile_home(profile) - # Composer model/effort/fast are PER-SESSION overrides, never a global config - # write. provider is optional (resolved at build). - create_model = str(params.get("model") or "").strip() - session_model_override = ( - {"model": create_model, "provider": str(params.get("provider") or "").strip() or None} - if create_model - else None - ) - create_reasoning_override = None - if effort := str(params.get("reasoning_effort") or "").strip(): - try: - from hermes_constants import parse_reasoning_effort - - create_reasoning_override = parse_reasoning_effort(effort) - except Exception: - create_reasoning_override = None - # ``fast`` presence is the contract: omitted inherits the profile, true pins - # priority, false pins normal ("" — _make_agent uses None for inheritance). - create_service_tier_override = None - if "fast" in params: - create_service_tier_override = "priority" if is_truthy_value(params.get("fast")) else "" + session_model_override, create_reasoning_override, create_service_tier_override = _create_overrides(params) now = time.time() with _sessions_lock: @@ -401,7 +421,6 @@ def _session_list_by_title(rid, db, title_lookup: str) -> dict: row = db.get_session_by_title(title_lookup) if row and row.get("archived"): from tools.bot_mode_probe import BOT_CHAT_TITLE - # The canonical Bot Chat is identity-scoped: an archive stamped by the # ws-orphan reaper / agent_close is an accident, and hiding it makes the # desktop mint transient replacements forever. Resurrect recoverable @@ -471,9 +490,7 @@ def _(rid, params: dict) -> dict: try: # Generous over-fetch so heavy sub-agent users (many ``tool`` rows) # don't get a false "no eligible session". - rows = db.list_sessions_rich( - source=None, limit=200, order_by_last_active=True, compact_rows=True - ) + rows = db.list_sessions_rich(source=None, limit=200, order_by_last_active=True, compact_rows=True) for row in rows: if _denied_source(row): continue @@ -501,7 +518,6 @@ def _(rid, params: dict) -> dict: """ try: from agent.coding_context import project_facts_for - return _ok(rid, {"facts": project_facts_for(params.get("cwd"))}) except Exception: logger.exception("project.facts failed") @@ -518,7 +534,6 @@ def _(rid, params: dict) -> dict: """ try: from agent.verification_evidence import verification_status - return _ok( rid, { @@ -536,7 +551,9 @@ def _(rid, params: dict) -> dict: # ── session.resume ─────────────────────────────────────────────────── -@dataclass +# repr/eq off: bind_module rebinds every function on the class onto server.py's +# globals, and dataclasses' generated __repr__ wrapper reads its own module globals. +@dataclass(repr=False, eq=False) class _Resume: """Per-call state for ``session.resume`` shared by the path helpers below. @@ -560,13 +577,16 @@ class _Resume: found: dict | None = None profile_resume_cwd: str = "" - def record(self, source: str, history: list, **extra) -> dict: + def cwd(self) -> str: + return self.profile_resume_cwd or _default_session_cwd() + + def record(self, source: str, cwd: str, history: list, **extra) -> dict: """``_deferred_session_record`` with this resume's common fields; the active-session lease is always claimed lazily on the first turn (_ensure_active_session_slot).""" return _deferred_session_record( self.target, cols=self.cols, - cwd=self.profile_resume_cwd or _default_session_cwd(), + cwd=cwd, history=history, lease=None, source=source, @@ -605,7 +625,6 @@ def _resume_live_unpersisted(ctx: _Resume, live_sid: str, live: dict) -> dict: if ctx.owns_db: with contextlib.suppress(Exception): from hermes_state import release_or_close - release_or_close(ctx.db) live["last_active"] = time.time() transport = current_transport() @@ -705,7 +724,6 @@ def _resume_follow_tip(ctx: _Resume) -> None: return try: from tools.bot_mode_probe import BOT_CHAT_TITLE - if (ctx.found.get("title") or "").strip() == BOT_CHAT_TITLE: tip = ctx.db.get_compression_tip(ctx.target) or ctx.target else: @@ -728,7 +746,6 @@ def _resume_guard(ctx: _Resume) -> dict | None: — only a genuine over-limit blocks. """ from hermes_state import SessionResumeTooLargeError, resolved_max_resume_messages - guard_tip_only = ctx.lazy or ctx.omit_messages or (ctx.defer_history and not ctx.eager_build) safety_check = getattr(ctx.db, "assert_resume_safe", None) try: @@ -745,9 +762,7 @@ def _resume_guard(ctx: _Resume) -> dict | None: except SessionResumeTooLargeError as exc: return _err(ctx.rid, 4130, str(exc)) except Exception as exc: - logger.warning( - "resume safety check failed for %s (proceeding without guard): %s", ctx.target, exc - ) + logger.warning("resume safety check failed for %s (proceeding without guard): %s", ctx.target, exc) return None @@ -803,21 +818,25 @@ def _resume_response( sid: str, record: dict, *, - messages: list, - message_count: int, info: dict, - running: bool, - status: str, + display: list = (), + count_source: list | None = None, + messages: list | None = None, + message_count: int | None = None, + running: bool = False, + status: str = "idle", hydrating: bool | None = None, started_at=None, auto_continue=None, ) -> dict: - payload = { - "session_id": sid, - "resumed": ctx.target, - "message_count": message_count, - "messages": messages, - } + """Common resume payload. ``messages`` is the display projection (empty when + omit_messages); the count then falls back to ``count_source`` so the client still + learns the stored size. ``hydrating`` (deferred path) replaces ``messages_omitted``.""" + if messages is None: + messages = [] if ctx.omit_messages else _history_to_messages(display) + if message_count is None: + message_count = len(count_source) if ctx.omit_messages else len(messages) + payload = {"session_id": sid, "resumed": ctx.target, "message_count": message_count, "messages": messages} if hydrating is None: payload["messages_omitted"] = ctx.omit_messages else: @@ -866,7 +885,8 @@ def _resume_lazy(ctx: _Resume) -> dict: ) except Exception as e: return ctx.resume_failed(e) - record = ctx.record(source, history, lazy=True, todo_state=_todo_state_from_history(history)) + cwd = ctx.cwd() + record = ctx.record(source, cwd, history, lazy=True, todo_state=_todo_state_from_history(history)) if (live := _claim_or_reuse_live(sid, ctx.target, record, None)) is not None: return _resume_reuse_live(ctx, *live) # A delegated child mid-run emits no session events of its own — report @@ -882,14 +902,13 @@ def _resume_lazy(ctx: _Resume) -> dict: except Exception: logger.debug("child-watch display projection read failed", exc_info=True) display_history = history - messages = [] if ctx.omit_messages else _history_to_messages(display_history) return _resume_response( ctx, sid, record, - messages=messages, - message_count=len(display_history) if ctx.omit_messages else len(messages), - info=_lazy_resume_info(record["cwd"], profile=ctx.profile), + info=_lazy_resume_info(cwd, profile=ctx.profile), + display=display_history, + count_source=display_history, running=child_running, status="streaming" if child_running else "idle", ) @@ -903,8 +922,10 @@ def _resume_deferred(ctx: _Resume) -> dict: sid, source = _new_runtime_ids(ctx.params) _enable_gateway_prompts() overrides = _stored_session_runtime_overrides(ctx.found) or {} + cwd = ctx.cwd() record = ctx.record( source, + cwd, [], model_override=overrides.get("model_override"), resume_runtime_overrides=overrides or None, @@ -924,10 +945,9 @@ def _resume_deferred(ctx: _Resume) -> dict: ctx, sid, record, + info=_resume_info(ctx, cwd, overrides), messages=[], message_count=record["resume_message_count"], - info=_resume_info(ctx, record["cwd"], overrides), - running=False, status="resuming", hydrating=True, ) @@ -953,8 +973,10 @@ def _resume_cold(ctx: _Resume) -> dict: # Restore model/provider/reasoning/tier so the deferred build matches the # eager path — without them the build drops the provider. overrides = _stored_session_runtime_overrides(ctx.found) or {} + cwd = ctx.cwd() record = ctx.record( source, + cwd, history, display_history_prefix=prefix, model_override=overrides.get("model_override"), @@ -968,16 +990,13 @@ def _resume_cold(ctx: _Resume) -> dict: _schedule_session_cap_enforcement() # trim detached idle sessions over the cap auto_continue = _maybe_schedule_auto_continue(sid, record, ctx.target) - messages = [] if ctx.omit_messages else _history_to_messages(display_history) return _resume_response( ctx, sid, record, - messages=messages, - message_count=len(raw_history) if ctx.omit_messages else len(messages), - info=_resume_info(ctx, record["cwd"], overrides), - running=False, - status="idle", + info=_resume_info(ctx, cwd, overrides), + display=display_history, + count_source=raw_history, auto_continue=auto_continue, ) @@ -1077,11 +1096,9 @@ def _resume_eager(ctx: _Resume) -> dict: ctx, sid, session, - messages=messages, - message_count=len(raw_history) if ctx.omit_messages else len(messages), info=_session_info(agent, session), - running=False, - status="idle", + messages=messages, + count_source=raw_history, started_at=float(session.get("created_at") or time.time()), auto_continue=auto_continue, ) @@ -1116,7 +1133,6 @@ def _(rid, params: dict) -> dict: # otherwise the shared launch db, which outlives the RPC and is never closed here. if ctx.profile_home is not None: from hermes_state import get_shared_session_db - ctx.db = get_shared_session_db(ctx.profile_home / "state.db") ctx.owns_db = True else: @@ -1192,7 +1208,6 @@ def _(rid, params: dict) -> dict: if not raw: return _err(rid, 4016, "cwd required") from hermes_constants import translate_cwd_for_wsl_backend - resolved = os.path.abspath(os.path.expanduser(translate_cwd_for_wsl_backend(raw))) if not os.path.isdir(resolved): return _err(rid, 4017, f"working directory does not exist: {raw}") @@ -1339,47 +1354,44 @@ def _title_read(rid, params: dict, session: dict, db) -> dict: @method("session.title") -@_with_session -def _(rid, params: dict, session: dict) -> dict: - with _session_db(session) as db: - if db is None: - return _db_unavailable_error(rid, code=5007) - if "title" not in params: - return _title_read(rid, params, session, db) - key = session["session_key"] - title = (params.get("title", "") or "").strip() - if not title: - return _err(rid, 4021, "title required") - sid = params.get("session_id", "") +@_with_session_db(5007) +def _(rid, params: dict, session: dict, db) -> dict: + if "title" not in params: + return _title_read(rid, params, session, db) + key = session["session_key"] + title = (params.get("title", "") or "").strip() + if not title: + return _err(rid, 4021, "title required") + sid = params.get("session_id", "") - def _done(pending: bool, value: str) -> dict: - session["pending_title"] = value if pending else None - _emit_session_info_for_session(sid, session) - return _ok(rid, {"pending": pending, "title": value}) + def _done(pending: bool, value: str) -> dict: + session["pending_title"] = value if pending else None + _emit_session_info_for_session(sid, session) + return _ok(rid, {"pending": pending, "title": value}) - try: - if db.set_session_title(key, title): + try: + if db.set_session_title(key, title): + return _done(False, title) + # rowcount == 0 can mean "same value" as well as "missing row". + existing_row = db.get_session(key) + if existing_row: + return _done(False, existing_row.get("title") or title) + # No row yet (deferred to the first prompt). An explicit /title is clear + # intent, so persist the row NOW (mirrors the messaging gateway's + # _handle_title_command) instead of queuing pending_title and hoping the + # post-turn apply block lands under this key. The min-messages sidebar + # filter keeps a titled 0-message row hidden. + _ensure_session_db_row(session) + with _session_db(session) as scoped_db: + if scoped_db is not None and scoped_db.set_session_title(key, title): return _done(False, title) - # rowcount == 0 can mean "same value" as well as "missing row". - existing_row = db.get_session(key) - if existing_row: - return _done(False, existing_row.get("title") or title) - # No row yet (deferred to the first prompt). An explicit /title is - # clear intent, so persist the row NOW (mirrors the messaging - # gateway's _handle_title_command) instead of queuing pending_title - # and hoping the post-turn apply block lands under this key. The - # min-messages sidebar filter keeps a titled 0-message row hidden. - _ensure_session_db_row(session) - with _session_db(session) as scoped_db: - if scoped_db is not None and scoped_db.set_session_title(key, title): - return _done(False, title) - # Row creation didn't take (DB unavailable / concurrent writer) — - # queue so the post-turn apply block can still recover. - return _done(True, title) - except ValueError as e: - return _err(rid, 4022, str(e)) - except Exception as e: - return _err(rid, 5007, str(e)) + # Row creation didn't take (DB unavailable / concurrent writer) — queue + # so the post-turn apply block can still recover. + return _done(True, title) + except ValueError as e: + return _err(rid, 4022, str(e)) + except Exception as e: + return _err(rid, 5007, str(e)) @method("session.set_hidden") @@ -1497,7 +1509,6 @@ def _(rid, params: dict) -> dict: try: from agent.oneshot import run_oneshot - text = run_oneshot( instructions=instructions, user_input=user_input, @@ -1533,9 +1544,7 @@ def _(rid, params: dict, session: dict) -> dict: desktop then polls ``handoff.state``. """ if session.get("running"): - return _err( - rid, 4009, "session busy — wait for the current turn to finish, then retry the handoff" - ) + return _err(rid, 4009, "session busy — wait for the current turn to finish, then retry the handoff") platform_name = (params.get("platform", "") or "").strip().lower() if not platform_name: @@ -1584,24 +1593,15 @@ def _(rid, params: dict, session: dict) -> dict: return _err(rid, 5007, str(e)) if not ok: - return _err( - rid, 4027, "session is already in flight for handoff — wait for it to settle, then retry" - ) - return _ok( - rid, {"queued": True, "session_key": key, "platform": platform_name, "home_name": home.name} - ) + return _err(rid, 4027, "session is already in flight for handoff — wait for it to settle, then retry") + return _ok(rid, {"queued": True, "session_key": key, "platform": platform_name, "home_name": home.name}) @method("handoff.state") -@_with_session -def _(rid, params: dict, session: dict) -> dict: +@_with_session_db(5007) +def _(rid, params: dict, session: dict, db) -> dict: """Poll ``{state, platform, error}``; ``state`` is pending|running|completed|failed or empty.""" - with _session_db(session) as db: - if db is None: - return _db_unavailable_error(rid, code=5007) - record = db.get_handoff_state(session["session_key"]) - - record = record or {} + record = db.get_handoff_state(session["session_key"]) or {} return _ok( rid, { @@ -1660,7 +1660,6 @@ def _(rid, params: dict, session: dict) -> dict: # zero API calls. Fail-open: absent when not logged in / portal hiccup. try: from agent.account_usage import nous_credits_lines - credits = nous_credits_lines() if credits: usage["credits_lines"] = credits @@ -1690,7 +1689,6 @@ def _(rid, params: dict, session: dict) -> dict: history = list(session.get("history", [])) try: from agent.context_breakdown import compute_session_context_breakdown - payload = compute_session_context_breakdown(agent, history) except Exception as exc: return _err(rid, 5000, f"Could not compute context breakdown: {exc}") @@ -1755,7 +1753,6 @@ def _(rid, params: dict) -> dict: """ from agent.pet import constants, render, store from agent.pet.render import PetRenderer - pet_cfg = _pet_display_cfg() if not is_truthy_value(pet_cfg.get("enabled"), default=False): return _ok(rid, {"enabled": False}) @@ -1821,7 +1818,6 @@ def _(rid, params: dict) -> dict: """ local_only = bool(params.get("localOnly")) from agent.pet import store - pet_cfg = _pet_display_cfg() installed = {p.slug: p for p in store.installed_pets()} @@ -1829,7 +1825,6 @@ def _(rid, params: dict) -> dict: seen: set[str] = set() try: from agent.pet.manifest import fetch_manifest, prefetch - # Local-only still warms the manifest cache in the background. if local_only: prefetch() @@ -1894,7 +1889,6 @@ def _(rid, params: dict, slug: str) -> dict: from agent.pet import store from agent.pet.manifest import ManifestError from hermes_cli.pets import _set_active - try: pet = store.install_pet(slug) except (store.PetStoreError, ManifestError) as exc: @@ -1911,7 +1905,6 @@ def _(rid, params: dict, slug: str) -> dict: """Uninstall a pet (delete its directory); if it was active, turn the display off.""" from agent.pet import store from hermes_cli.pets import _clear_active_if - removed = store.remove_pet(slug) try: _clear_active_if(slug) @@ -1927,7 +1920,6 @@ def _(rid, params: dict, slug: str) -> dict: def _(rid, params: dict, slug: str) -> dict: """Export an installed pet as a re-importable ``.zip`` → ``{ok, filename, zipBase64}``.""" import base64 - from agent.pet import store filename, data = store.export_pet(slug) @@ -1947,14 +1939,12 @@ def _(rid, params: dict, slug: str) -> dict: if not name: return _err(rid, 4004, "missing name") from agent.pet import store - new_slug = store.rename_pet(slug, name) if not new_slug: return _err(rid, 5031, "pet.rename failed") if new_slug != slug: try: from hermes_cli.pets import _rename_active_if - _rename_active_if(slug, new_slug) except Exception as exc: # noqa: BLE001 - rename already succeeded logger.debug("pet.rename config update failed: %s", exc) @@ -1969,7 +1959,6 @@ def _(rid, params: dict, slug: str) -> dict: """Small idle-frame PNG data URI for the picker preview (same-origin; the desktop CSP / R2 hotlink rules break a CDN ````). ``url`` serves not-yet-installed pets.""" import base64 - from agent.pet import store data = store.thumbnail_png(slug, source_url=str(params.get("url") or "")) @@ -1991,7 +1980,6 @@ def _(rid, params: dict, slug: str) -> dict: def _(rid, params: dict) -> dict: """``display.pet.enabled=false`` from the desktop picker.""" from hermes_cli.pets import _set_enabled - _set_enabled(False) return _ok(rid, {"ok": True}) @@ -2002,7 +1990,6 @@ def _(rid, params: dict) -> dict: def _(rid, params: dict) -> dict: """Persist ``display.pet.scale`` (clamped to engine bounds) from the desktop slider.""" from hermes_cli.pets import set_pet_scale - scale, err = set_pet_scale(params.get("scale")) if err: return _err(rid, 4004, err) @@ -2026,7 +2013,6 @@ def _(rid, params: dict) -> dict: def _(rid, params: dict) -> dict: """Whether pet generation is possible: a reference-capable image backend is configured.""" from agent.pet.generate.imagegen import GenerationError, list_sprite_providers, resolve_provider - try: resolve_provider(require_references=True) available = True @@ -2060,10 +2046,8 @@ def _(rid, params: dict) -> dict: style = str(params.get("style") or "auto").strip() or "auto" import shutil - from agent.pet.generate import generate_base_drafts from agent.pet.generate.imagegen import GenerationError, resolve_provider - root = _pet_gen_root() _pet_gen_sweep(root) @@ -2170,7 +2154,6 @@ def _(rid, params: dict) -> dict: from agent.pet import store from agent.pet.generate import hatch_pet from agent.pet.generate.imagegen import GenerationError, resolve_provider - base = _pet_gen_root() / token / f"draft-{index}.png" if not base.is_file(): return _err(rid, 4004, "draft expired — generate again") @@ -2242,7 +2225,6 @@ def _(rid, params: dict) -> dict: """GET /api/billing/state → serialized BillingState. No scope required.""" try: from agent.billing_view import build_billing_state - return _ok(rid, _serialize_billing_state(build_billing_state())) except Exception: return _ok(rid, {"ok": True, "logged_in": False, "error": "could not load billing state"}) @@ -2253,7 +2235,6 @@ def _(rid, params: dict) -> dict: """Shared dollar usage model (two-bar view) for /usage + /subscription.""" try: from agent.billing_usage import build_usage_model - return _ok(rid, _serialize_usage_model(build_usage_model())) except Exception: return _ok(rid, {"ok": True, "available": False}) @@ -2264,7 +2245,6 @@ def _(rid, params: dict) -> dict: """GET /api/billing/subscription → serialized SubscriptionState (read-only).""" try: from agent.subscription_view import build_subscription_state - return _ok(rid, _serialize_subscription_state(build_subscription_state())) except Exception: return _ok(rid, {"ok": True, "logged_in": False, "error": "could not load subscription state"}) @@ -2275,7 +2255,6 @@ def _(rid, params: dict) -> dict: """POST /api/billing/subscription/preview → chargeless effect quote. billing:manage.""" from agent.subscription_view import subscription_change_preview_from_payload from hermes_cli.nous_billing import post_subscription_preview - tier_id = params.get("subscription_type_id") if not tier_id: return _billing_invalid(rid, "subscription_type_id is required") @@ -2292,7 +2271,6 @@ def _(rid, params: dict) -> dict: """PUT /api/billing/subscription/pending-change: schedule a downgrade / same-price change OR a period-end cancellation (chargeless). billing:manage.""" from hermes_cli.nous_billing import put_subscription_pending_change - cancel = bool(params.get("cancel")) tier_id = params.get("subscription_type_id") if not cancel and not tier_id: @@ -2310,7 +2288,6 @@ def _(rid, params: dict) -> dict: """DELETE /api/billing/subscription/pending-change: clear a scheduled downgrade / cancellation. Re-enables recurring spend → billing:manage + kill-switch.""" from hermes_cli.nous_billing import delete_subscription_pending_change - def call(): result = delete_subscription_pending_change() return {"ok": True, "message": result.get("message"), "payload": result} @@ -2328,7 +2305,6 @@ def _(rid, params: dict) -> dict: """ from agent.billing_view import new_idempotency_key from hermes_cli.nous_billing import post_subscription_upgrade - tier_id = params.get("subscription_type_id") if not tier_id: return _billing_invalid(rid, "subscription_type_id is required") @@ -2354,7 +2330,6 @@ def _(rid, params: dict) -> dict: and echoed (also on error) so the TUI reuses it on retry of the SAME purchase.""" from hermes_cli.nous_billing import post_charge from agent.billing_view import new_idempotency_key - amount = params.get("amount_usd") if amount is None: return _billing_invalid(rid, "amount_usd is required") @@ -2371,7 +2346,6 @@ def _(rid, params: dict) -> dict: def _(rid, params: dict) -> dict: """GET /api/billing/charge/{id} — a single status read; the caller drives the poll cadence.""" from hermes_cli.nous_billing import get_charge_status - charge_id = params.get("charge_id") if not charge_id: return _billing_invalid(rid, "charge_id is required", error="invalid_charge_id") @@ -2393,7 +2367,6 @@ def _(rid, params: dict) -> dict: def _(rid, params: dict) -> dict: """PATCH /api/billing/auto-top-up. params: {enabled, threshold, top_up_amount}.""" from hermes_cli.nous_billing import patch_auto_top_up - enabled = bool(params.get("enabled")) threshold = params.get("threshold") top_up_amount = params.get("top_up_amount") @@ -2422,7 +2395,6 @@ def _(rid, params: dict) -> dict: def call(): from hermes_cli.auth import step_up_nous_billing_scope - def _on_verification(url: str, code: str) -> None: _emit("billing.step_up.verification", sid, {"verification_url": url, "user_code": code}) @@ -2439,7 +2411,6 @@ def _(rid, params: dict) -> dict: @_with_session def _(rid, params: dict, session: dict) -> dict: from hermes_constants import display_hermes_home - key = session.get("session_key") or params.get("session_id") or "" agent = session.get("agent") @@ -2533,7 +2504,6 @@ def _(rid, params: dict, session: dict) -> dict: # Truncate from the last *real* user turn: popping trailing assistant/tool # then one user left timeline markers / compaction handoffs as the target. from agent.context_compressor import user_originated_turn_view - user_indices = [ index for index, message in enumerate(history) if user_originated_turn_view(message) is not None ] @@ -2623,13 +2593,11 @@ def _(rid, params: dict) -> dict: if session.get("running"): return _err(rid, 4009, "session busy — /interrupt the current turn before /compress") from agent.conversation_compression import finalize_context_engine_compression_notification - sid = params.get("session_id", "") focus_topic = str(params.get("focus_topic", "") or "").strip() try: from agent.manual_compression_feedback import summarize_manual_compression from agent.model_metadata import estimate_request_tokens_rough - with session["history_lock"]: before_messages = list(session.get("history", [])) history_version = int(session.get("history_version", 0)) @@ -2703,7 +2671,6 @@ def _(rid, params: dict) -> dict: except CompressionLockHeld as e: _status_update(sid, "ready") from agent.manual_compression_feedback import describe_compression_lock_skip - return _ok(rid, {"compressed": False, "lock_held": True, "message": describe_compression_lock_skip(e.holder)}) except Exception as e: finalize_context_engine_compression_notification(session["agent"], committed=False) @@ -2849,19 +2816,17 @@ def _(rid, params: dict, session: dict) -> dict: source = _session_source(session) try: title = params.get("name", "") or _branch_title(db, old_key) - _create_branch_rows( + _create_branch_row( db, new_key, old_key, - title, - history, source=source, cwd=_session_cwd(session), profile_name=( Path(session["profile_home"]).name if session.get("profile_home") else _current_profile_name() ), - copy_fields=_BRANCH_COPY_FIELDS, ) + _copy_branch_transcript(db, new_key, title, history, _BRANCH_COPY_FIELDS) except Exception as e: return _err(rid, 5008, f"branch failed: {e}") # Bound before the try so the ownership finally can never see them unbound. @@ -2876,7 +2841,6 @@ def _(rid, params: dict, session: dict) -> dict: # DEDICATED handle, same ownership rule as session.resume: ours until # the branched agent takes it below. from hermes_state import get_shared_session_db - branch_db = get_shared_session_db(Path(parent_home) / "state.db") branch_owns_db = True with _profile_build_scope(parent_home): @@ -2916,7 +2880,6 @@ def _(rid, params: dict, session: dict) -> dict: if branch_owns_db and branch_db is not None: with contextlib.suppress(Exception): from hermes_state import release_or_close - release_or_close(branch_db) branched_session = _sessions.get(new_sid) return _ok( @@ -3067,14 +3030,12 @@ def _(rid, params: dict) -> dict: @method("delegation.pause") def _(rid, params: dict) -> dict: from tools.delegate_tool import set_spawn_paused - return _ok(rid, {"paused": set_spawn_paused(bool(params.get("paused", True)))}) @method("subagent.interrupt") def _(rid, params: dict) -> dict: from tools.delegate_tool import interrupt_subagent - subagent_id = str(params.get("subagent_id") or "").strip() if not subagent_id: return _err(rid, 4000, "subagent_id required") @@ -3091,7 +3052,6 @@ def _(rid, params: dict) -> dict: ``missed_steer`` on the parent's completion entry. """ from tools.delegate_tool import steer_subagent - subagent_id = str(params.get("subagent_id") or "").strip() if not subagent_id: return _err(rid, 4000, "subagent_id required") @@ -3240,7 +3200,6 @@ def _(rid, params: dict) -> dict: except (TypeError, ValueError): return _err(rid, -32602, "invalid params: last_seen must be an integer") from tui_gateway import event_replay - frames = event_replay.events_since(sid, last_seen) return _ok( rid, @@ -3260,7 +3219,6 @@ def _(rid, params: dict) -> dict: def _(rid, params: dict) -> dict: """Replay-buffer telemetry (ops/debug).""" from tui_gateway import event_replay - return _ok(rid, event_replay.replay_stats()) diff --git a/tui_gateway/methods_slash.py b/tui_gateway/methods_slash.py new file mode 100644 index 0000000000..70e66be887 --- /dev/null +++ b/tui_gateway/methods_slash.py @@ -0,0 +1,497 @@ +"""slash.exec helpers: command resolution, side-effect mirroring after a slash command ran in the worker. + +Bodies are rebound onto server.py's globals at install time (see +method_ctx.bind_module), so they reference server.py globals bare. +""" + +from __future__ import annotations + + +from .method_ctx import HandlerRegistry, bind_module + +_registry = HandlerRegistry() + + +# ── Methods: slash.exec ────────────────────────────────────────────── + + +_LIVE_SESSION_DIRECT_COMMANDS = frozenset( + { + "clear", + "compress", + "effort", + "history", + "models", + "prompt", + "rename", + "review", + "status", + "usage", + } +) + +_ISOLATED_SESSION_READ_COMMANDS = frozenset({"context", "tools", "help"}) + + +def _format_live_review_output(session: Optional[dict], arg: str) -> str: + """Dispatch /review against the live TUI/desktop session's agent. + + Spawns the reviewer subagent on the async delegation rail; the TUI + notification poller already drains async-delegation completions for the + owning session, so the finished review re-enters this chat as a normal + completion turn. The dispatch stamps the parent agent's durable + session_id as the completion's session_key (the delegate_task CLI-path + fallback), which is exactly what ``_session_owns_notification_event`` + matches against. + """ + if session is None: + return "Nothing to review yet — send a message first." + if _session_uses_compute_host(session): + return ( + "/review runs on the local agent only for now — this session's " + "agent lives on a remote compute host." + ) + agent = session.get("agent") + if agent is None: + return "Nothing to review yet — send a message first." + if session.get("running"): + return "session busy — wait for the current turn to finish, then /review" + + history_lock = session.get("history_lock") + if history_lock is not None: + with history_lock: + snapshot = list(session.get("history", [])) + else: + snapshot = list(session.get("history", [])) + if not snapshot: + snapshot = list(getattr(agent, "_session_messages", None) or []) + + try: + from agent.review_engine import format_dispatch_note, start_review + + result = start_review(agent, snapshot, arg or "") + except ValueError as exc: + return str(exc) + except Exception as exc: + return f"/review failed to start: {exc}" + return format_dispatch_note(result, arg or "") + + +def _format_live_usage_output(session: dict) -> str: + agent = session.get("agent") + usage = _session_usage_snapshot(session) + if agent is None and not usage: + return "(._.) No active agent -- send a message first." + if session.get("_metadata_message_count") is not None: + message_count = int(session.get("_metadata_message_count") or 0) + else: + with session["history_lock"]: + message_count = len(session.get("history", [])) + lines = [ + "Session Token Usage", + "────────────────────────────────────────", + f"Model: {usage.get('model') or _metadata_mirror(session).get('model') or getattr(agent, 'model', '') or '(unknown)'}", + f"Input tokens: {int(usage.get('input') or 0):,}", + f"Output tokens: {int(usage.get('output') or 0):,}", + ] + reasoning = int(usage.get("reasoning") or 0) + if reasoning: + lines.append(f"Reasoning tokens: {reasoning:,}") + lines.extend( + [ + f"Prompt tokens: {int(usage.get('prompt') or 0):,}", + f"Completion tokens: {int(usage.get('completion') or 0):,}", + f"Total tokens: {int(usage.get('total') or 0):,}", + f"API calls: {int(usage.get('calls') or 0):,}", + ] + ) + if usage.get("context_max"): + lines.append( + "Current context: " + f"{int(usage.get('context_used') or 0):,} / " + f"{int(usage.get('context_max') or 0):,} " + f"({int(usage.get('context_percent') or 0)}%)" + ) + lines.extend( + [ + f"Messages: {message_count:,}", + f"Compressions: {int(usage.get('compressions') or 0):,}", + ] + ) + return "\n".join(lines) + + +def _format_live_history_output(session: dict) -> str: + with session["history_lock"]: + history = list(session.get("history", [])) + # _session_db, not _get_db(): a profile session's transcript lives in its + # own profile's state.db, and this read is scoped by session id — through + # the launch handle it comes back empty and /history renders nothing. + with _session_db(session) as db: + if db is not None and session.get("session_key"): + try: + history = db.get_messages_as_conversation( + session["session_key"], include_ancestors=True, include_row_ids=True + ) + except Exception: + pass + messages = _history_to_messages(history) + if not messages: + return "No conversation history yet." + lines = ["Conversation History", "────────────────────────────────────────"] + for idx, message in enumerate(messages, start=1): + role = str(message.get("role") or "unknown") + label = "You" if role == "user" else "Hermes" if role == "assistant" else role.title() + text = str(message.get("text") or message.get("context") or "").strip() + if len(text) > 400: + text = f"{text[:400]}..." + lines.append(f"[{label} #{idx}] {text or '(no text)'}") + return "\n".join(lines) + + +def _format_live_prompt_output(session: dict) -> str: + agent = session.get("agent") + mirror = _metadata_mirror(session) + if agent is None and "system_prompt" not in mirror: + return "No active agent -- send a message first." + prompt = ( + mirror.get("system_prompt") + or getattr(agent, "ephemeral_system_prompt", None) + or getattr(agent, "_cached_system_prompt", None) + or "" + ) + if not prompt: + return "Current system prompt is not built yet; send a message first." + return f"Current system prompt:\n{prompt}" + + +def _format_live_context_output(session: dict) -> str: + messages = [] + # Same session-scoped read as /history — resolve it against the db that + # owns this session's rows, not the launch profile's handle. + with _session_db(session) as db: + if db is not None and session.get("session_key"): + try: + messages = _history_to_messages( + db.get_messages_as_conversation( + session["session_key"], include_ancestors=True, include_row_ids=True + ) + ) + except Exception: + messages = [] + if not messages: + with session["history_lock"]: + messages = _history_to_messages(list(session.get("history", []))) + usage = _session_usage_snapshot(session) + mirror = _metadata_mirror(session) + lines = [ + f"Conversation: {len(messages)} messages" if messages else "Conversation is empty (no messages yet)." + ] + roles: dict[str, int] = {} + for msg in messages: + role = str(msg.get("role") or "unknown") + roles[role] = roles.get(role, 0) + 1 + lines.append( + f" user: {roles.get('user', 0)}, assistant: {roles.get('assistant', 0)}, " + f"tool: {roles.get('tool', 0)}, system: {roles.get('system', 0)}" + ) + model = mirror.get("model") or usage.get("model") or "" + provider = mirror.get("provider") or "auto" + if model: + lines.append(f"Model: {model}") + lines.append(f"Provider: {provider}") + context_used = int(usage.get("context_used") or usage.get("total") or 0) + context_max = int(usage.get("context_max") or 0) + if context_used: + if context_max: + usage_pct = (context_used / context_max) * 100 + lines.append( + f"Context usage: ~{context_used:,} / {context_max:,} tokens ({usage_pct:.1f}%)" + ) + else: + lines.append(f"Context usage: ~{context_used:,} tokens") + if usage.get("compressions"): + lines.append(f"Compressions: {int(usage.get('compressions') or 0):,}") + return "\n".join(lines) + + +def _format_live_tools_output(session: dict) -> str: + info = _session_info(session.get("agent"), session) + groups = info.get("tools") if isinstance(info, dict) else {} + if not isinstance(groups, dict) or not groups: + return "No tools available." + names: list[str] = [] + for group_names in groups.values(): + if isinstance(group_names, list): + names.extend(str(name) for name in group_names) + names = sorted(set(names)) + if not names: + return "No tools available." + return "Available tools ({}):\n{}".format( + len(names), "\n".join(f" {name}" for name in names) + ) + + +def _format_live_help_output() -> str: + try: + from hermes_cli.commands import COMMANDS_BY_CATEGORY + + lines = ["Available commands:", ""] + for category, commands in COMMANDS_BY_CATEGORY.items(): + lines.append(f"{category}:") + for cmd, desc in commands.items(): + lines.append(f" {cmd:<15} {desc}") + return "\n".join(lines) + except Exception as exc: + return f"help unavailable: {exc}" + + +def _format_live_model_output(session: dict) -> str: + agent = session.get("agent") + model = getattr(agent, "model", "") if agent is not None else "" + provider = getattr(agent, "provider", "") if agent is not None else "" + if model and provider: + return f"Current model: {model} ({provider})" + if model: + return f"Current model: {model}" + return "Current model: (unknown)" + + +def _live_slash_command_output(sid: str, session: Optional[dict], name: str, arg: str) -> Optional[str]: + name = (name or "").lstrip("/").lower() + arg = arg or "" + if name == "model" and not arg.strip(): + return _format_live_model_output(session or {}) + if name not in _LIVE_SESSION_DIRECT_COMMANDS: + if not ( + name in _ISOLATED_SESSION_READ_COMMANDS + and session is not None + and _session_uses_compute_host(session) + ): + return None + + if name in _ISOLATED_SESSION_READ_COMMANDS and not ( + session is not None and _session_uses_compute_host(session) + ): + return None + if name == "compress": + if session is None: + return "no active session for /compress" + return _mirror_slash_side_effects(sid, session, f"/compress {arg}".strip()) + if name == "usage": + if session is None: + return "(._.) No active agent -- send a message first." + return _format_live_usage_output(session) + if name == "review": + return _format_live_review_output(session, arg) + if name == "history": + if session is None: + return "No conversation history yet." + return _format_live_history_output(session) + if name == "prompt": + if session is None: + return "No active agent -- send a message first." + return _format_live_prompt_output(session) + if name == "status": + response = _methods["session.status"]("status", {"session_id": sid}) + if response.get("error"): + return str(response["error"].get("message") or "status unavailable") + return str(response.get("result", {}).get("output") or "") + if name == "context": + if session is None: + return "Conversation is empty (no messages yet)." + return _format_live_context_output(session) + if name == "tools": + if session is None: + return "No tools available." + return _format_live_tools_output(session) + if name == "help": + return _format_live_help_output() + if name == "clear": + return "Screen clear is terminal-only; desktop/TUI chat left unchanged." + if name == "models": + return "Use /model to view or switch the current model; desktop users can also open the model picker." + if name == "rename": + return "Use /title to rename this session." + if name == "effort": + return "Use /reasoning to change reasoning effort." + return None + + + +def _mirror_slash_side_effects(sid: str, session: dict, command: str) -> str: + """Apply side effects that must also hit the gateway's live agent.""" + parts = command.lstrip("/").split(None, 1) + if not parts: + return "" + name, arg, agent = ( + parts[0], + (parts[1].strip() if len(parts) > 1 else ""), + session.get("agent"), + ) + if name == "compact": + # /compact is an alias of /compress in every host. The compute-host + # slash.compress control forwards the user's raw alias verbatim, so + # without normalizing here the child mirror silently no-ops — the + # session never compresses and the deferred context-engine + # notification wiring below is never exercised for that route. + name = "compress" + + # Reject agent-mutating commands during an in-flight turn. These + # all do read-then-mutate on live agent/session state that the + # worker thread running agent.run_conversation is using. Parity + # with the session.compress / session.undo guards and the gateway + # runner's running-agent /model guard. + _MUTATES_WHILE_RUNNING = {"model", "personality", "prompt", "compress"} + if _session_uses_compute_host(session) and name in _MUTATES_WHILE_RUNNING: + route_name = f"slash.{name}" + is_compress = name == "compress" + _late_session = session + + def _on_late_ack(late: dict, _sid=sid) -> None: + _adopt_late_compute_host_compress_ack(_sid, _late_session, late, route_name=route_name) + + try: + ack = _send_compute_host_control( + sid, + route_name=route_name, + command=command, + wait=True, + **( + {"timeout": _compute_host_compress_wait_seconds(), "on_late_ack": _on_late_ack} + if is_compress + else {} + ), + ) + except queue.Empty: + if is_compress: + return "compression still running in the background; the transcript will refresh when it finishes" + return f"compute-host {route_name} failed: timed out" + except Exception as exc: + return f"compute-host {route_name} failed: {exc}" + if ack.get("type") in {"control.error", "error"}: + return str(ack.get("message") or f"compute-host {route_name} failed") + _apply_compute_host_metadata_mirror(session, ack) + return str(ack.get("output") or "") + if name in _MUTATES_WHILE_RUNNING and session.get("running"): + return f"session busy — /interrupt the current turn before running /{name}" + + try: + if name == "model" and arg and agent: + result = _apply_model_switch(sid, session, arg) + return result.get("warning", "") + elif name == "approvals" and arg: + # The slash worker already persisted the new approvals.mode; the + # bare (read-only) form has no arg and needs no repaint. + broadcast_session_info() + elif name == "personality" and arg and agent: + pname, new_prompt = _validate_personality(arg, _load_cfg()) + # Persist through the single owner so this surface can never + # drift from the others (the old TUI slash path applied the + # overlay in-session but skipped persistence entirely). + from hermes_cli.personality import persist_personality + + persist_personality(pname) + _apply_personality_to_session(sid, session, new_prompt, pname) + elif name == "prompt" and agent: + cfg = _load_cfg() + new_prompt = _prompt_text((cfg.get("agent") or {}).get("system_prompt", "")) + agent.ephemeral_system_prompt = new_prompt or None + agent._cached_system_prompt = None + elif name == "compress" and agent: + # Mirror the session.compress RPC: build a before/after summary so + # the user gets feedback (#46686). The slash path previously just + # compressed + emitted session.info and returned "", so the TUI + # showed no "compressed N → M messages / ~X → ~Y tokens" stats + # while CLI and gateway both did. + from agent.manual_compression_feedback import summarize_manual_compression + from agent.model_metadata import estimate_request_tokens_rough + from agent.conversation_compression import ( + finalize_context_engine_compression_notification, + ) + + with session["history_lock"]: + _before_messages = list(session.get("history", [])) + _before_count = len(_before_messages) + _sys_prompt = getattr(agent, "_cached_system_prompt", "") or "" + _tools = getattr(agent, "tools", None) or None + _before_tokens = ( + estimate_request_tokens_rough( + _before_messages, system_prompt=_sys_prompt, tools=_tools + ) + if _before_count + else 0 + ) + + # The raw argument goes through unparsed: _compress_session_history + # (the choke point shared by all three manual-compress routes) + # parses the boundary-aware forms (here [N], up to here, --keep N) + # and does the partial head/tail split there (#35533). + try: + _compress_session_history(session, arg) + except CompressionLockHeld as e: + from agent.manual_compression_feedback import ( + describe_compression_lock_skip, + ) + return describe_compression_lock_skip(e.holder) + _sync_session_key_after_compress(sid, session) + + with session["history_lock"]: + _after_messages = list(session.get("history", [])) + _sys_prompt_after = getattr(agent, "_cached_system_prompt", "") or _sys_prompt + _tools_after = getattr(agent, "tools", None) or _tools + _after_tokens = ( + estimate_request_tokens_rough( + _after_messages, system_prompt=_sys_prompt_after, tools=_tools_after + ) + if _after_messages + else 0 + ) + _emit("session.info", sid, _session_info(agent, session)) + _fb = summarize_manual_compression( + _before_messages, + _after_messages, + _before_tokens, + _after_tokens, + compression_state=getattr(agent, "context_compressor", None), + ) + _lines = [_fb["headline"], _fb["token_line"]] + if _fb.get("note"): + _lines.append(_fb["note"]) + finalize_context_engine_compression_notification( + agent, + committed=True, + ) + return "\n".join(_lines) + elif name == "fast" and agent: + mode = arg.lower() + if mode in {"fast", "on"}: + agent.service_tier = "priority" + elif mode in {"normal", "off"}: + agent.service_tier = None + elif mode in {"auto", "cold"}: + agent.service_tier = mode + _emit("session.info", sid, _session_info(agent, session)) + elif name == "reload-mcp" and agent and hasattr(agent, "reload_mcp_tools"): + agent.reload_mcp_tools() + elif name == "stop": + from tools.process_registry import process_registry + + process_registry.kill_all() + except Exception as e: + if name == "compress" and agent: + from agent.conversation_compression import ( + finalize_context_engine_compression_notification, + ) + + finalize_context_engine_compression_notification( + agent, + committed=False, + ) + return f"live session sync failed: {e}" + return "" + + +def register(server) -> None: + """Publish this module's helpers + handlers onto ``server``, rebound to its globals.""" + bind_module(globals(), server, skip=("_",)) diff --git a/tui_gateway/server.py b/tui_gateway/server.py index defc5598ce..476789b270 100644 --- a/tui_gateway/server.py +++ b/tui_gateway/server.py @@ -15403,813 +15403,6 @@ def _resolve_name(name: str) -> str: _paste_counter = 0 -# ── Methods: complete ───────────────────────────────────────────────── - -_FUZZY_CACHE_TTL_S = 5.0 -_FUZZY_CACHE_MAX_FILES = 20000 -_FUZZY_FALLBACK_EXCLUDES = frozenset( - { - ".git", - ".hg", - ".svn", - ".next", - ".cache", - ".venv", - "venv", - "node_modules", - "__pycache__", - "dist", - "build", - "target", - ".mypy_cache", - ".pytest_cache", - ".ruff_cache", - } -) -_fuzzy_cache_lock = threading.Lock() -_fuzzy_cache: dict[str, tuple[float, list[str]]] = {} - - -def _list_repo_files(root: str) -> list[str]: - """Return file paths relative to ``root``. - - Uses ``git ls-files`` from the repo top (resolved via - ``rev-parse --show-toplevel``) so the listing covers tracked + untracked - files anywhere in the repo, then converts each path back to be relative - to ``root``. Files outside ``root`` (parent directories of cwd, sibling - subtrees) are excluded so the picker stays scoped to what's reachable - from the gateway's cwd. Falls back to a bounded ``os.walk(root)`` when - ``root`` isn't inside a git repo. Result cached per-root for - ``_FUZZY_CACHE_TTL_S`` so rapid keystrokes don't respawn git processes. - """ - now = time.monotonic() - with _fuzzy_cache_lock: - cached = _fuzzy_cache.get(root) - if cached and now - cached[0] < _FUZZY_CACHE_TTL_S: - return cached[1] - - files: list[str] = [] - from hermes_cli._subprocess_compat import windows_hide_flags - - _creationflags = windows_hide_flags() - try: - top_result = subprocess.run( - ["git", "-C", root, "rev-parse", "--show-toplevel"], - capture_output=True, - timeout=2.0, - check=False, - stdin=subprocess.DEVNULL, - creationflags=_creationflags, - ) - if top_result.returncode == 0: - top = top_result.stdout.decode("utf-8", "replace").strip() - list_result = subprocess.run( - [ - "git", - "-C", - top, - "ls-files", - "-z", - "--cached", - "--others", - "--exclude-standard", - ], - capture_output=True, - timeout=2.0, - check=False, - stdin=subprocess.DEVNULL, - creationflags=_creationflags, - ) - if list_result.returncode == 0: - for p in list_result.stdout.decode("utf-8", "replace").split("\0"): - if not p: - continue - rel = os.path.relpath(os.path.join(top, p), root).replace( - os.sep, "/" - ) - # Skip parents/siblings of cwd — keep the picker scoped - # to root-and-below, matching Cmd-P workspace semantics. - if rel.startswith("../"): - continue - files.append(rel) - if len(files) >= _FUZZY_CACHE_MAX_FILES: - break - except (OSError, subprocess.TimeoutExpired): - pass - - if not files: - # Fallback walk: skip vendor/build dirs + dot-dirs so the walk stays - # tractable. Dotfiles themselves survive — the ranker decides based - # on whether the query starts with `.`. - try: - for dirpath, dirnames, filenames in os.walk(root, followlinks=False): - dirnames[:] = [ - d - for d in dirnames - if d not in _FUZZY_FALLBACK_EXCLUDES and not d.startswith(".") - ] - rel_dir = os.path.relpath(dirpath, root) - for f in filenames: - rel = f if rel_dir == "." else f"{rel_dir}/{f}" - files.append(rel.replace(os.sep, "/")) - if len(files) >= _FUZZY_CACHE_MAX_FILES: - break - if len(files) >= _FUZZY_CACHE_MAX_FILES: - break - except OSError: - pass - - with _fuzzy_cache_lock: - _fuzzy_cache[root] = (now, files) - - return files - - -def _fuzzy_basename_rank(name: str, query: str) -> tuple[int, int] | None: - """Rank ``name`` against ``query``; lower is better. Returns None to reject. - - Tiers (kind): - 0 — exact basename - 1 — basename prefix (e.g. `app` → `appChrome.tsx`) - 2 — word-boundary / camelCase hit (e.g. `chrome` → `appChrome.tsx`) - 3 — substring anywhere in basename - 4 — subsequence match (every query char appears in order) - - Secondary key is `len(name)` so shorter names win ties. - """ - if not query: - return (3, len(name)) - - nl = name.lower() - ql = query.lower() - - if nl == ql: - return (0, len(name)) - - if nl.startswith(ql): - return (1, len(name)) - - # Word-boundary split: `foo-bar_baz.qux` → ["foo","bar","baz","qux"]. - # camelCase split: `appChrome` → ["app","Chrome"]. Cheap approximation; - # falls through to substring/subsequence if it misses. - parts: list[str] = [] - buf = "" - for ch in name: - if ch in "-_." or (ch.isupper() and buf and not buf[-1].isupper()): - if buf: - parts.append(buf) - buf = ch if ch not in "-_." else "" - else: - buf += ch - if buf: - parts.append(buf) - for p in parts: - if p.lower().startswith(ql): - return (2, len(name)) - - if ql in nl: - return (3, len(name)) - - i = 0 - for ch in nl: - if ch == ql[i]: - i += 1 - if i == len(ql): - return (4, len(name)) - - return None - - -def _abs_completion_prefix_exists(path_part: str) -> bool: - """True when ``path_part`` reads sensibly as an absolute path. - - A leading `/` is only meant literally if something is actually there: - the parent directory has to exist, and a partially-typed final segment - has to match at least one of its entries. Used to decide whether - `@/foo` is the absolute `/foo` or shorthand for `foo` under the cwd. - """ - expanded = _normalize_completion_path(path_part) - parent = os.path.dirname(expanded.rstrip("/")) or "/" - tail = os.path.basename(expanded.rstrip("/")) - - if not os.path.isdir(parent): - return False - - if not tail or expanded.endswith("/"): - return os.path.isdir(expanded) or expanded == "/" - - try: - tail_lower = tail.lower() - return any(e.lower().startswith(tail_lower) for e in os.listdir(parent)) - except OSError: - return False - - -def _details_completion_item(value: str, meta: str = "") -> dict: - return {"text": value, "display": value, "meta": meta} - - -def _details_root_completion_item( - value: str, meta: str, needs_leading_space: bool -) -> dict: - return _details_completion_item( - f" {value}" if needs_leading_space else value, - meta, - ) - - -def _details_completions(text: str) -> list[dict] | None: - if not text.lower().startswith("/details"): - return None - - stripped = text.strip() - if stripped and not "/details".startswith(stripped.lower().split()[0]): - return None - - body = text[len("/details") :] - if body.startswith(" "): - body = body[1:] - parts = body.split() - has_trailing_space = text.endswith(" ") - sections = ("thinking", "tools", "subagents", "activity") - modes = ("hidden", "collapsed", "expanded") - - if not body or (len(parts) == 0 and has_trailing_space): - return [ - *[ - _details_root_completion_item( - mode, "global mode", not has_trailing_space - ) - for mode in modes - ], - _details_root_completion_item( - "cycle", "cycle global mode", not has_trailing_space - ), - *[ - _details_root_completion_item( - section, "section override", not has_trailing_space - ) - for section in sections - ], - ] - - if len(parts) == 1 and not has_trailing_space: - prefix = parts[0].lower() - candidates = [*modes, "cycle", *sections] - return [ - _details_completion_item( - candidate, - ( - "section override" - if candidate in sections - else "cycle global mode" if candidate == "cycle" else "global mode" - ), - ) - for candidate in candidates - if candidate.startswith(prefix) and candidate != prefix - ] - - if len(parts) == 1 and has_trailing_space and parts[0].lower() in sections: - return [ - *[ - _details_completion_item(mode, f"set {parts[0].lower()}") - for mode in modes - ], - _details_completion_item("reset", f"clear {parts[0].lower()} override"), - ] - - if len(parts) == 2 and not has_trailing_space and parts[0].lower() in sections: - prefix = parts[1].lower() - return [ - _details_completion_item( - candidate, - ( - f"clear {parts[0].lower()} override" - if candidate == "reset" - else f"set {parts[0].lower()}" - ), - ) - for candidate in (*modes, "reset") - if candidate.startswith(prefix) and candidate != prefix - ] - - return [] - - -def _model_picker_context(agent): - """Layer live session state onto config without losing custom identity.""" - from hermes_cli.inventory import load_picker_context - - ctx = load_picker_context() - provider = getattr(agent, "provider", "") if agent else "" - base_url = getattr(agent, "base_url", "") if agent else "" - if str(provider or "").strip().lower() == "custom": - try: - from hermes_cli.runtime_provider import canonical_custom_identity - - provider = ( - canonical_custom_identity( - base_url=base_url or None, - config_provider=ctx.current_provider, - model=(getattr(agent, "model", "") if agent else "") - or None, - ) - or provider - ) - except Exception: - logger.debug( - "custom provider identity recovery failed (model picker)", - exc_info=True, - ) - - return ctx.with_overrides( - current_provider=provider, - current_model=(getattr(agent, "model", "") if agent else "") - or _resolve_model(), - current_base_url=base_url, - ) - - -# ── Methods: slash.exec ────────────────────────────────────────────── - - -_LIVE_SESSION_DIRECT_COMMANDS = frozenset( - { - "clear", - "compress", - "effort", - "history", - "models", - "prompt", - "rename", - "review", - "status", - "usage", - } -) - -_ISOLATED_SESSION_READ_COMMANDS = frozenset({"context", "tools", "help"}) - - -def _format_live_review_output(session: Optional[dict], arg: str) -> str: - """Dispatch /review against the live TUI/desktop session's agent. - - Spawns the reviewer subagent on the async delegation rail; the TUI - notification poller already drains async-delegation completions for the - owning session, so the finished review re-enters this chat as a normal - completion turn. The dispatch stamps the parent agent's durable - session_id as the completion's session_key (the delegate_task CLI-path - fallback), which is exactly what ``_session_owns_notification_event`` - matches against. - """ - if session is None: - return "Nothing to review yet — send a message first." - if _session_uses_compute_host(session): - return ( - "/review runs on the local agent only for now — this session's " - "agent lives on a remote compute host." - ) - agent = session.get("agent") - if agent is None: - return "Nothing to review yet — send a message first." - if session.get("running"): - return "session busy — wait for the current turn to finish, then /review" - - history_lock = session.get("history_lock") - if history_lock is not None: - with history_lock: - snapshot = list(session.get("history", [])) - else: - snapshot = list(session.get("history", [])) - if not snapshot: - snapshot = list(getattr(agent, "_session_messages", None) or []) - - try: - from agent.review_engine import format_dispatch_note, start_review - - result = start_review(agent, snapshot, arg or "") - except ValueError as exc: - return str(exc) - except Exception as exc: - return f"/review failed to start: {exc}" - return format_dispatch_note(result, arg or "") - - -def _format_live_usage_output(session: dict) -> str: - agent = session.get("agent") - usage = _session_usage_snapshot(session) - if agent is None and not usage: - return "(._.) No active agent -- send a message first." - if session.get("_metadata_message_count") is not None: - message_count = int(session.get("_metadata_message_count") or 0) - else: - with session["history_lock"]: - message_count = len(session.get("history", [])) - lines = [ - "Session Token Usage", - "────────────────────────────────────────", - f"Model: {usage.get('model') or _metadata_mirror(session).get('model') or getattr(agent, 'model', '') or '(unknown)'}", - f"Input tokens: {int(usage.get('input') or 0):,}", - f"Output tokens: {int(usage.get('output') or 0):,}", - ] - reasoning = int(usage.get("reasoning") or 0) - if reasoning: - lines.append(f"Reasoning tokens: {reasoning:,}") - lines.extend( - [ - f"Prompt tokens: {int(usage.get('prompt') or 0):,}", - f"Completion tokens: {int(usage.get('completion') or 0):,}", - f"Total tokens: {int(usage.get('total') or 0):,}", - f"API calls: {int(usage.get('calls') or 0):,}", - ] - ) - if usage.get("context_max"): - lines.append( - "Current context: " - f"{int(usage.get('context_used') or 0):,} / " - f"{int(usage.get('context_max') or 0):,} " - f"({int(usage.get('context_percent') or 0)}%)" - ) - lines.extend( - [ - f"Messages: {message_count:,}", - f"Compressions: {int(usage.get('compressions') or 0):,}", - ] - ) - return "\n".join(lines) - - -def _format_live_history_output(session: dict) -> str: - with session["history_lock"]: - history = list(session.get("history", [])) - # _session_db, not _get_db(): a profile session's transcript lives in its - # own profile's state.db, and this read is scoped by session id — through - # the launch handle it comes back empty and /history renders nothing. - with _session_db(session) as db: - if db is not None and session.get("session_key"): - try: - history = db.get_messages_as_conversation( - session["session_key"], include_ancestors=True, include_row_ids=True - ) - except Exception: - pass - messages = _history_to_messages(history) - if not messages: - return "No conversation history yet." - lines = ["Conversation History", "────────────────────────────────────────"] - for idx, message in enumerate(messages, start=1): - role = str(message.get("role") or "unknown") - label = "You" if role == "user" else "Hermes" if role == "assistant" else role.title() - text = str(message.get("text") or message.get("context") or "").strip() - if len(text) > 400: - text = f"{text[:400]}..." - lines.append(f"[{label} #{idx}] {text or '(no text)'}") - return "\n".join(lines) - - -def _format_live_prompt_output(session: dict) -> str: - agent = session.get("agent") - mirror = _metadata_mirror(session) - if agent is None and "system_prompt" not in mirror: - return "No active agent -- send a message first." - prompt = ( - mirror.get("system_prompt") - or getattr(agent, "ephemeral_system_prompt", None) - or getattr(agent, "_cached_system_prompt", None) - or "" - ) - if not prompt: - return "Current system prompt is not built yet; send a message first." - return f"Current system prompt:\n{prompt}" - - -def _format_live_context_output(session: dict) -> str: - messages = [] - # Same session-scoped read as /history — resolve it against the db that - # owns this session's rows, not the launch profile's handle. - with _session_db(session) as db: - if db is not None and session.get("session_key"): - try: - messages = _history_to_messages( - db.get_messages_as_conversation( - session["session_key"], include_ancestors=True, include_row_ids=True - ) - ) - except Exception: - messages = [] - if not messages: - with session["history_lock"]: - messages = _history_to_messages(list(session.get("history", []))) - usage = _session_usage_snapshot(session) - mirror = _metadata_mirror(session) - lines = [ - f"Conversation: {len(messages)} messages" if messages else "Conversation is empty (no messages yet)." - ] - roles: dict[str, int] = {} - for msg in messages: - role = str(msg.get("role") or "unknown") - roles[role] = roles.get(role, 0) + 1 - lines.append( - f" user: {roles.get('user', 0)}, assistant: {roles.get('assistant', 0)}, " - f"tool: {roles.get('tool', 0)}, system: {roles.get('system', 0)}" - ) - model = mirror.get("model") or usage.get("model") or "" - provider = mirror.get("provider") or "auto" - if model: - lines.append(f"Model: {model}") - lines.append(f"Provider: {provider}") - context_used = int(usage.get("context_used") or usage.get("total") or 0) - context_max = int(usage.get("context_max") or 0) - if context_used: - if context_max: - usage_pct = (context_used / context_max) * 100 - lines.append( - f"Context usage: ~{context_used:,} / {context_max:,} tokens ({usage_pct:.1f}%)" - ) - else: - lines.append(f"Context usage: ~{context_used:,} tokens") - if usage.get("compressions"): - lines.append(f"Compressions: {int(usage.get('compressions') or 0):,}") - return "\n".join(lines) - - -def _format_live_tools_output(session: dict) -> str: - info = _session_info(session.get("agent"), session) - groups = info.get("tools") if isinstance(info, dict) else {} - if not isinstance(groups, dict) or not groups: - return "No tools available." - names: list[str] = [] - for group_names in groups.values(): - if isinstance(group_names, list): - names.extend(str(name) for name in group_names) - names = sorted(set(names)) - if not names: - return "No tools available." - return "Available tools ({}):\n{}".format( - len(names), "\n".join(f" {name}" for name in names) - ) - - -def _format_live_help_output() -> str: - try: - from hermes_cli.commands import COMMANDS_BY_CATEGORY - - lines = ["Available commands:", ""] - for category, commands in COMMANDS_BY_CATEGORY.items(): - lines.append(f"{category}:") - for cmd, desc in commands.items(): - lines.append(f" {cmd:<15} {desc}") - return "\n".join(lines) - except Exception as exc: - return f"help unavailable: {exc}" - - -def _format_live_model_output(session: dict) -> str: - agent = session.get("agent") - model = getattr(agent, "model", "") if agent is not None else "" - provider = getattr(agent, "provider", "") if agent is not None else "" - if model and provider: - return f"Current model: {model} ({provider})" - if model: - return f"Current model: {model}" - return "Current model: (unknown)" - - -def _live_slash_command_output(sid: str, session: Optional[dict], name: str, arg: str) -> Optional[str]: - name = (name or "").lstrip("/").lower() - arg = arg or "" - if name == "model" and not arg.strip(): - return _format_live_model_output(session or {}) - if name not in _LIVE_SESSION_DIRECT_COMMANDS: - if not ( - name in _ISOLATED_SESSION_READ_COMMANDS - and session is not None - and _session_uses_compute_host(session) - ): - return None - - if name in _ISOLATED_SESSION_READ_COMMANDS and not ( - session is not None and _session_uses_compute_host(session) - ): - return None - if name == "compress": - if session is None: - return "no active session for /compress" - return _mirror_slash_side_effects(sid, session, f"/compress {arg}".strip()) - if name == "usage": - if session is None: - return "(._.) No active agent -- send a message first." - return _format_live_usage_output(session) - if name == "review": - return _format_live_review_output(session, arg) - if name == "history": - if session is None: - return "No conversation history yet." - return _format_live_history_output(session) - if name == "prompt": - if session is None: - return "No active agent -- send a message first." - return _format_live_prompt_output(session) - if name == "status": - response = _methods["session.status"]("status", {"session_id": sid}) - if response.get("error"): - return str(response["error"].get("message") or "status unavailable") - return str(response.get("result", {}).get("output") or "") - if name == "context": - if session is None: - return "Conversation is empty (no messages yet)." - return _format_live_context_output(session) - if name == "tools": - if session is None: - return "No tools available." - return _format_live_tools_output(session) - if name == "help": - return _format_live_help_output() - if name == "clear": - return "Screen clear is terminal-only; desktop/TUI chat left unchanged." - if name == "models": - return "Use /model to view or switch the current model; desktop users can also open the model picker." - if name == "rename": - return "Use /title to rename this session." - if name == "effort": - return "Use /reasoning to change reasoning effort." - return None - - - -def _mirror_slash_side_effects(sid: str, session: dict, command: str) -> str: - """Apply side effects that must also hit the gateway's live agent.""" - parts = command.lstrip("/").split(None, 1) - if not parts: - return "" - name, arg, agent = ( - parts[0], - (parts[1].strip() if len(parts) > 1 else ""), - session.get("agent"), - ) - if name == "compact": - # /compact is an alias of /compress in every host. The compute-host - # slash.compress control forwards the user's raw alias verbatim, so - # without normalizing here the child mirror silently no-ops — the - # session never compresses and the deferred context-engine - # notification wiring below is never exercised for that route. - name = "compress" - - # Reject agent-mutating commands during an in-flight turn. These - # all do read-then-mutate on live agent/session state that the - # worker thread running agent.run_conversation is using. Parity - # with the session.compress / session.undo guards and the gateway - # runner's running-agent /model guard. - _MUTATES_WHILE_RUNNING = {"model", "personality", "prompt", "compress"} - if _session_uses_compute_host(session) and name in _MUTATES_WHILE_RUNNING: - route_name = f"slash.{name}" - is_compress = name == "compress" - _late_session = session - - def _on_late_ack(late: dict, _sid=sid) -> None: - _adopt_late_compute_host_compress_ack(_sid, _late_session, late, route_name=route_name) - - try: - ack = _send_compute_host_control( - sid, - route_name=route_name, - command=command, - wait=True, - **( - {"timeout": _compute_host_compress_wait_seconds(), "on_late_ack": _on_late_ack} - if is_compress - else {} - ), - ) - except queue.Empty: - if is_compress: - return "compression still running in the background; the transcript will refresh when it finishes" - return f"compute-host {route_name} failed: timed out" - except Exception as exc: - return f"compute-host {route_name} failed: {exc}" - if ack.get("type") in {"control.error", "error"}: - return str(ack.get("message") or f"compute-host {route_name} failed") - _apply_compute_host_metadata_mirror(session, ack) - return str(ack.get("output") or "") - if name in _MUTATES_WHILE_RUNNING and session.get("running"): - return f"session busy — /interrupt the current turn before running /{name}" - - try: - if name == "model" and arg and agent: - result = _apply_model_switch(sid, session, arg) - return result.get("warning", "") - elif name == "approvals" and arg: - # The slash worker already persisted the new approvals.mode; the - # bare (read-only) form has no arg and needs no repaint. - broadcast_session_info() - elif name == "personality" and arg and agent: - pname, new_prompt = _validate_personality(arg, _load_cfg()) - # Persist through the single owner so this surface can never - # drift from the others (the old TUI slash path applied the - # overlay in-session but skipped persistence entirely). - from hermes_cli.personality import persist_personality - - persist_personality(pname) - _apply_personality_to_session(sid, session, new_prompt, pname) - elif name == "prompt" and agent: - cfg = _load_cfg() - new_prompt = _prompt_text((cfg.get("agent") or {}).get("system_prompt", "")) - agent.ephemeral_system_prompt = new_prompt or None - agent._cached_system_prompt = None - elif name == "compress" and agent: - # Mirror the session.compress RPC: build a before/after summary so - # the user gets feedback (#46686). The slash path previously just - # compressed + emitted session.info and returned "", so the TUI - # showed no "compressed N → M messages / ~X → ~Y tokens" stats - # while CLI and gateway both did. - from agent.manual_compression_feedback import summarize_manual_compression - from agent.model_metadata import estimate_request_tokens_rough - from agent.conversation_compression import ( - finalize_context_engine_compression_notification, - ) - - with session["history_lock"]: - _before_messages = list(session.get("history", [])) - _before_count = len(_before_messages) - _sys_prompt = getattr(agent, "_cached_system_prompt", "") or "" - _tools = getattr(agent, "tools", None) or None - _before_tokens = ( - estimate_request_tokens_rough( - _before_messages, system_prompt=_sys_prompt, tools=_tools - ) - if _before_count - else 0 - ) - - # The raw argument goes through unparsed: _compress_session_history - # (the choke point shared by all three manual-compress routes) - # parses the boundary-aware forms (here [N], up to here, --keep N) - # and does the partial head/tail split there (#35533). - try: - _compress_session_history(session, arg) - except CompressionLockHeld as e: - from agent.manual_compression_feedback import ( - describe_compression_lock_skip, - ) - return describe_compression_lock_skip(e.holder) - _sync_session_key_after_compress(sid, session) - - with session["history_lock"]: - _after_messages = list(session.get("history", [])) - _sys_prompt_after = getattr(agent, "_cached_system_prompt", "") or _sys_prompt - _tools_after = getattr(agent, "tools", None) or _tools - _after_tokens = ( - estimate_request_tokens_rough( - _after_messages, system_prompt=_sys_prompt_after, tools=_tools_after - ) - if _after_messages - else 0 - ) - _emit("session.info", sid, _session_info(agent, session)) - _fb = summarize_manual_compression( - _before_messages, - _after_messages, - _before_tokens, - _after_tokens, - compression_state=getattr(agent, "context_compressor", None), - ) - _lines = [_fb["headline"], _fb["token_line"]] - if _fb.get("note"): - _lines.append(_fb["note"]) - finalize_context_engine_compression_notification( - agent, - committed=True, - ) - return "\n".join(_lines) - elif name == "fast" and agent: - mode = arg.lower() - if mode in {"fast", "on"}: - agent.service_tier = "priority" - elif mode in {"normal", "off"}: - agent.service_tier = None - elif mode in {"auto", "cold"}: - agent.service_tier = mode - _emit("session.info", sid, _session_info(agent, session)) - elif name == "reload-mcp" and agent and hasattr(agent, "reload_mcp_tools"): - agent.reload_mcp_tools() - elif name == "stop": - from tools.process_registry import process_registry - - process_registry.kill_all() - except Exception as e: - if name == "compress" and agent: - from agent.conversation_compression import ( - finalize_context_engine_compression_notification, - ) - - finalize_context_engine_compression_notification( - agent, - committed=False, - ) - return f"live session sync failed: {e}" - return "" - - # ── Methods: insights ──────────────────────────────────────────────── @@ -16244,6 +15437,8 @@ def _mcp_summarize_server(name, cfg): # noqa: E402 from . import ( # noqa: E402 methods_voice as _methods_voice, methods_browser as _methods_browser, + methods_slash as _methods_slash, + methods_complete_helpers as _methods_complete_helpers, methods_browser_control as _methods_browser_control, methods_bot_relay as _methods_bot_relay, methods_complete as _methods_complete, @@ -16257,6 +15452,8 @@ from . import ( # noqa: E402 ) for _m in ( + _methods_complete_helpers, + _methods_slash, _methods_voice, _methods_browser, _methods_browser_control,