From f997235f1bcd2c0505efaec4227c028b2c19887d Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 23:02:07 -0700 Subject: [PATCH] refactor(tui_gateway): dedupe config/voice/profiles helpers, collapse ladders and hug payloads --- tui_gateway/methods_config.py | 148 +++------- tui_gateway/methods_config_set.py | 99 +++---- tui_gateway/methods_profiles.py | 98 +++---- tui_gateway/methods_voice.py | 447 ++++++++++++------------------ 4 files changed, 282 insertions(+), 510 deletions(-) diff --git a/tui_gateway/methods_config.py b/tui_gateway/methods_config.py index 7ba7a45feb..2d448a28ef 100644 --- a/tui_gateway/methods_config.py +++ b/tui_gateway/methods_config.py @@ -31,7 +31,6 @@ def _(rid, params: dict) -> dict: if db is None: return _ok(rid, {"repos": []}) from hermes_cli import projects_db as pdb - policy = _repo_discovery_policy() with pdb.connect_closing() as conn: _reconcile_repo_discovery(pdb, conn, policy, _repo_discovery_policy_key(policy)) @@ -53,20 +52,14 @@ def _(rid, params: dict) -> dict: crawl runs on the desktop), then return the merged repo list.""" try: from hermes_cli import projects_db as pdb - policy = _repo_discovery_policy() policy_key = _repo_discovery_policy_key(policy) incoming_raw = params.get("discovery_policy") - incoming_policy = ( - _repo_discovery_policy(incoming_raw) if isinstance(incoming_raw, dict) else None - ) + incoming_policy = _repo_discovery_policy(incoming_raw) if isinstance(incoming_raw, dict) else None incoming_matches = ( - incoming_policy is not None - and _repo_discovery_policy_key(incoming_policy) == policy_key - ) - accept_legacy_default = ( - incoming_policy is None and _repo_discovery_policy_is_default(policy) + incoming_policy is not None and _repo_discovery_policy_key(incoming_policy) == policy_key ) + accept_legacy_default = incoming_policy is None and _repo_discovery_policy_is_default(policy) pairs: list[tuple[str, str | None]] = [] for item in params.get("repos") or []: @@ -84,16 +77,8 @@ def _(rid, params: dict) -> dict: pdb.clear_discovered_repos(conn, policy_key=policy_key) with _profile_db(params) as db: - return _ok( - rid, - { - "repos": _discover_repos_payload(db, include_cached=policy["enabled"]) - if db is not None - else [], - "accepted": accepted, - "discovery_policy": policy, - }, - ) + repos = _discover_repos_payload(db, include_cached=policy["enabled"]) if db is not None else [] + return _ok(rid, {"repos": repos, "accepted": accepted, "discovery_policy": policy}) except Exception as e: return _err(rid, 5061, str(e)) @@ -101,7 +86,6 @@ def _(rid, params: dict) -> dict: def _stamped_project_tree(db, params, **kwargs): """``_build_project_tree`` + profile stamping shared by the two tree RPCs.""" from tui_gateway.project_tree import stamp_profile - tree, active_id = _build_project_tree(db, **kwargs) stamp_profile(tree["projects"], _response_profile_name(params.get("profile"))) return tree, active_id @@ -120,21 +104,13 @@ def _(rid, params: dict) -> dict: if db is None: return _ok(rid, {"projects": [], "active_id": None, "scoped_session_ids": []}) tree, active_id = _stamped_project_tree( - db, - params, - preview_limit=int(params.get("preview_limit") or 3), - hydrate=False, - session_limit=int(params.get("session_limit") or 2000), - include_discovered=True, - ) - return _ok( - rid, - { - "projects": tree["projects"], - "active_id": active_id, - "scoped_session_ids": tree["scoped_session_ids"], - }, + db, params, preview_limit=int(params.get("preview_limit") or 3), hydrate=False, + session_limit=int(params.get("session_limit") or 2000), include_discovered=True, ) + return _ok(rid, { + "projects": tree["projects"], "active_id": active_id, + "scoped_session_ids": tree["scoped_session_ids"], + }) except Exception as e: return _err(rid, 5061, str(e)) @@ -156,12 +132,8 @@ def _(rid, params: dict) -> dict: # Drill-in only needs the entered project (which has sessions): # skip the zero-session discovery tier. tree, _active = _stamped_project_tree( - db, - params, - preview_limit=0, - hydrate=True, - session_limit=int(params.get("session_limit") or 5000), - include_discovered=False, + db, params, preview_limit=0, hydrate=True, + session_limit=int(params.get("session_limit") or 5000), include_discovered=False, ) proj = next((p for p in tree["projects"] if p["id"] == project_id), None) return _ok(rid, {"project": proj}) @@ -175,24 +147,17 @@ def _(rid, params: dict) -> dict: # --------------------------------------------------------------------------- -def _display_cfg() -> dict: - display = _load_cfg().get("display") - return display if isinstance(display, dict) else {} - - def _display_mode(cfg: dict, key: str, allowed: frozenset, default: str) -> str: raw = str((cfg.get("display") or {}).get(key, default) or default).strip().lower() return raw if raw in allowed else default -_DETAILS_MODES = frozenset({"hidden", "collapsed", "expanded"}) _THINKING_MODES = frozenset({"collapsed", "truncated", "full"}) def _cfg_get_provider(rid, params): try: from hermes_cli.models import list_available_providers, normalize_provider - model = _resolve_model() parts = model.split("/", 1) return { @@ -206,7 +171,6 @@ def _cfg_get_provider(rid, params): def _cfg_get_profile(rid, params): from hermes_constants import display_hermes_home - return {"home": str(_hermes_home), "display": display_hermes_home()} @@ -229,7 +193,6 @@ def _cfg_get_personality(rid, params): # EFFECTIVE personality via the single owner — a stale/unknown name in # config must not display as active. from hermes_cli.personality import active_personality_name - return {"value": active_personality_name(_load_cfg()) or "none"} @@ -238,13 +201,9 @@ def _cfg_get_reasoning(rid, params): session = _sessions.get(params.get("session_id", "")) reasoning_config = None if session is not None: - if isinstance(session.get("create_reasoning_override"), dict): - reasoning_config = session.get("create_reasoning_override") - else: - agent_reasoning = getattr(session.get("agent"), "reasoning_config", None) - if isinstance(agent_reasoning, dict): - reasoning_config = agent_reasoning - + reasoning_config = session.get("create_reasoning_override") + if not isinstance(reasoning_config, dict): + reasoning_config = getattr(session.get("agent"), "reasoning_config", None) if isinstance(reasoning_config, dict): if reasoning_config.get("enabled") is False: effort = "none" @@ -287,7 +246,7 @@ def _cfg_get_thinking_mode(rid, params): raw = str((cfg.get("display") or {}).get("thinking_mode", "") or "").strip().lower() if raw in _THINKING_MODES: return {"value": raw} - dm = _display_mode(cfg, "details_mode", _DETAILS_MODES, "collapsed") + dm = _display_mode(cfg, "details_mode", _DETAIL_MODES, "collapsed") return {"value": "full" if dm == "expanded" else "collapsed"} @@ -332,7 +291,7 @@ def _config_getters() -> dict: "approval_mode": _cfg_get_approval_mode, "approvals.mode": _cfg_get_approval_mode, "details_mode": lambda rid, params: { - "value": _display_mode(_load_cfg(), "details_mode", _DETAILS_MODES, "collapsed") + "value": _display_mode(_load_cfg(), "details_mode", _DETAIL_MODES, "collapsed") }, "thinking_mode": _cfg_get_thinking_mode, "density": lambda rid, params: { @@ -376,12 +335,10 @@ def _readiness_profile_scope(params: dict): quietly answer for the launch profile instead. """ import contextlib - profile = str(params.get("profile") or "").strip() if isinstance(params, dict) else "" if not profile: return "", contextlib.nullcontext() from hermes_cli import profiles as profiles_mod - if not profiles_mod.profile_exists(profile): raise FileNotFoundError(f"Profile '{profile}' does not exist on this backend.") home = _profile_home(profile) @@ -410,13 +367,9 @@ def _(rid, params: dict) -> dict: """Loose provider check; ``profile`` (optional) scopes it to that profile's home.""" try: from hermes_cli.main import _has_any_provider_configured - def probe(profile): configured = bool(_has_any_provider_configured(strict_profile_scope=bool(profile))) - payload = {"provider_configured": configured} - if profile: - payload["profile"] = profile - return payload + return {"provider_configured": configured, **({"profile": profile} if profile else {})} return _readiness_check(rid, params, probe) except Exception as e: @@ -439,7 +392,6 @@ def _(rid, params: dict) -> dict: from hermes_cli.runtime_provider import resolve_runtime_provider from hermes_cli.auth import has_usable_secret from hermes_cli.main import _has_any_provider_configured - requested = str(params.get("provider") or "").strip() or None def probe(profile): @@ -450,44 +402,24 @@ def _(rid, params: dict) -> dict: scoped = {"profile": profile} if profile else {} provider = runtime.get("provider") or "provider" source = str(runtime.get("source") or "") - if ( - not provider_configured - and provider == "bedrock" - and source in {"iam-role", "aws-sdk-default-chain"} - ): - return { - "ok": False, - "provider": provider, - "model": runtime.get("model"), - "source": source, - "error": "No Hermes provider is configured.", - **scoped, - } + def fail(error, src): + return {"ok": False, "provider": provider, "model": runtime.get("model"), + "source": src, "error": error, **scoped} + + if (not provider_configured and provider == "bedrock" + and source in {"iam-role", "aws-sdk-default-chain"}): + return fail("No Hermes provider is configured.", source) api_key = runtime.get("api_key") api_key_text = "" if callable(api_key) else str(api_key or "").strip() credential_ok = ( - callable(api_key) - or api_key_text in {"aws-sdk", "no-key-required"} - or has_usable_secret(api_key_text) - or bool(runtime.get("command")) + callable(api_key) or api_key_text in {"aws-sdk", "no-key-required"} + or has_usable_secret(api_key_text) or bool(runtime.get("command")) ) if not credential_ok: - return { - "ok": False, - "provider": provider, - "model": runtime.get("model"), - "source": runtime.get("source"), - "error": f"No usable credentials found for {provider}.", - **scoped, - } - return { - "ok": True, - "provider": runtime.get("provider"), - "model": runtime.get("model"), - "source": runtime.get("source"), - **scoped, - } + return fail(f"No usable credentials found for {provider}.", runtime.get("source")) + return {"ok": True, "provider": runtime.get("provider"), "model": runtime.get("model"), + "source": runtime.get("source"), **scoped} return _readiness_check(rid, params, probe) except Exception as e: @@ -512,7 +444,6 @@ def _(rid, params: dict) -> dict: try: from hermes_cli.debug import _redact_log_text, build_nous_bundle, collect_share_bundle from hermes_cli.diagnostics_upload import share_to_nous - log_lines = params.get("log_lines") if not isinstance(log_lines, int) or not (10 <= log_lines <= 2000): log_lines = 200 @@ -548,18 +479,11 @@ def _(rid, params: dict) -> dict: upload_id = res.get("id") if not view_url and not upload_id: # An upload the user can't reference is useless to support. - return _ok( - rid, {"ok": False, "error": "upload succeeded but returned no view URL or id"} - ) - return _ok( - rid, - { - "ok": True, - "view_url": view_url, - "upload_id": upload_id, - "expires_at": res.get("expiresAt") or res.get("expires_at"), - }, - ) + return _ok(rid, {"ok": False, "error": "upload succeeded but returned no view URL or id"}) + return _ok(rid, { + "ok": True, "view_url": view_url, "upload_id": upload_id, + "expires_at": res.get("expiresAt") or res.get("expires_at"), + }) except Exception as e: return _ok(rid, {"ok": False, "error": str(e)}) diff --git a/tui_gateway/methods_config_set.py b/tui_gateway/methods_config_set.py index aa1b316c92..12cdfe3781 100644 --- a/tui_gateway/methods_config_set.py +++ b/tui_gateway/methods_config_set.py @@ -52,28 +52,28 @@ def _emit_all_session_info() -> None: def _toggle_display_bool(rid, key, value, *, cfg_key, on_words, off_words): """Shared body of the on/off/toggle display booleans (``density``, ``battery``).""" - raw = str(value or "").strip().lower() - cur_b = bool(_display_section(_load_cfg()).get(cfg_key, False)) + raw = _word(value) + cur_b = bool(_display_cfg().get(cfg_key, False)) if raw in {"", "toggle"}: nv_b = not cur_b - elif raw in on_words: - nv_b = True - elif raw in off_words: - nv_b = False + elif raw in on_words or raw in off_words: + nv_b = raw in on_words else: return _err(rid, 4002, f"unknown {key} value: {value}") _write_config_key(f"display.{cfg_key}", nv_b) return _ok(rid, {"key": key, "value": "on" if nv_b else "off"}) +def _word(value) -> str: + return str(value or "").strip().lower() + + def _cfgset_await_agent(session, rid): """Wait for an in-progress agent build; the error envelope if it failed, else None.""" init_err = _wait_agent(session, rid) if init_err: return init_err - if session.get("agent") is None: - return _err(rid, 5032, "agent initialization failed") - return None + return _err(rid, 5032, "agent initialization failed") if session.get("agent") is None else None def _cfgset_model_ok(rid, key, value, warning, confirm_required, confirm_message, scope, **extra): @@ -92,7 +92,6 @@ def _set_model(rid, params, key, value, session): confirmed = bool(params.get("confirm_expensive_model", False)) if session: from hermes_cli.model_switch import parse_model_switch_args - sid = params.get("session_id", "") # No live swap while a turn streams: agent.switch_model() mutates model/provider/ # base_url/client that the worker thread reads every iteration. Stash the pick and @@ -138,8 +137,7 @@ def _set_model(rid, params, key, value, session): ) if session.get("agent") is None and not explicit_provider.strip() and not failed_agent_init: _start_agent_build(sid, session) - init_err = _cfgset_await_agent(session, rid) - if init_err: + if init_err := _cfgset_await_agent(session, rid): return init_err with _session_profile_runtime_scope(session): result = _apply_model_switch( @@ -147,8 +145,7 @@ def _set_model(rid, params, key, value, session): ) if failed_agent_init and not result.get("confirm_required"): _restart_completed_failed_agent_build(sid, session, failed_ready) - init_err = _cfgset_await_agent(session, rid) - if init_err: + if init_err := _cfgset_await_agent(session, rid): return init_err with _session_profile_runtime_scope(session): _persist_live_session_runtime(session) @@ -162,8 +159,11 @@ def _set_model(rid, params, key, value, session): return _err(rid, 5001, str(e)) +_FAST_WORDS = {"fast": "fast", "on": "fast", "normal": "normal", "off": "normal", "auto": "auto", "cold": "cold"} + + def _set_fast(rid, params, key, value, session): - raw = str(value or "").strip().lower() + raw = _word(value) agent = session.get("agent") if session else None if agent is not None: current_tier = getattr(agent, "service_tier", None) @@ -174,23 +174,15 @@ def _set_fast(rid, params, key, value, session): current_tier = _load_service_tier() current_fast = current_tier == "priority" - if raw in {"status"}: + if raw == "status": return _ok(rid, {"key": key, "value": {"priority": "fast", None: "normal"}.get(current_tier, current_tier)}) - if raw in {"", "toggle"}: - nv = "normal" if current_fast else "fast" - elif raw in {"fast", "on"}: - nv = "fast" - elif raw in {"normal", "off"}: - nv = "normal" - elif raw in {"auto", "cold"}: - nv = raw - else: + nv = _FAST_WORDS.get(raw, ("normal" if current_fast else "fast") if raw in {"", "toggle"} else None) + if nv is None: return _err(rid, 4002, f"unknown fast mode: {value}") overrides = None if nv == "fast": from hermes_cli.models import resolve_fast_mode_overrides - if agent is not None: target_model = getattr(agent, "model", None) else: @@ -200,8 +192,7 @@ def _set_fast(rid, params, key, value, session): if not target_model: return _err(rid, 4002, "fast mode is not available without a selected model") overrides = resolve_fast_mode_overrides( - target_model, provider=getattr(agent, "provider", None), base_url=getattr(agent, "base_url", None) - ) + target_model, provider=getattr(agent, "provider", None), base_url=getattr(agent, "base_url", None)) if overrides is None: return _err(rid, 4002, "fast mode is not available for this model") @@ -226,7 +217,7 @@ def _set_fast(rid, params, key, value, session): def _set_busy(rid, params, key, value, session): - raw = str(value or "").strip().lower() + raw = _word(value) if raw in {"", "status"}: return _ok(rid, {"key": key, "value": _load_busy_input_mode()}) if raw not in {"queue", "steer", "interrupt"}: @@ -248,9 +239,8 @@ def _set_verbose(rid, params, key, value, session): _write_config_key("display.tool_progress", nv) if session: session["tool_progress_mode"] = nv - agent = session.get("agent") - if agent is not None: - agent.verbose_logging = nv == "verbose" + if session.get("agent") is not None: + session["agent"].verbose_logging = nv == "verbose" return _ok(rid, {"key": key, "value": nv}) @@ -258,8 +248,7 @@ def _set_focus(rid, params, key, value, session): # Focus view (/focus): display-only reduced output composed with tool_progress — enabling # stashes the configured mode and pins tool_progress "off"; disabling restores the stash. from hermes_cli.focus_view import FOCUS_TOOL_PROGRESS_MODE, normalize_tool_progress_mode, resolve_focus_arg - - d_f = _display_section(_load_cfg()) + d_f = _display_cfg() cur_focus = bool(d_f.get("focus_view", False)) action, target = resolve_focus_arg(str(value or ""), cur_focus) if action == "usage": @@ -280,15 +269,14 @@ def _set_focus(rid, params, key, value, session): if session: session["focus_view"] = bool(target) session["tool_progress_mode"] = effective - agent_f = session.get("agent") - if agent_f is not None: + if session.get("agent") is not None: with contextlib.suppress(Exception): - agent_f.tool_progress_mode = effective + session["agent"].tool_progress_mode = effective return _ok(rid, {"key": key, "value": "on" if target else "off", "tool_progress": effective}) def _set_approval_mode(rid, params, key, value, session): - raw = str(value or "").strip().lower() + raw = _word(value) if raw not in _APPROVAL_MODES: return _err(rid, 4002, f"unknown approval mode: {value}; pick one of manual|smart|off") _write_config_key("approvals.mode", raw) @@ -300,20 +288,17 @@ def _set_yolo(rid, params, key, value, session): # Approval bypass. scope="session" (default; TUI Shift+Tab) toggles ONLY this session's flag. # scope="global" (Shift+click the zap) flips persistent approvals.mode between "off" (bypass # on) and "manual" (bypass off) for every surface, surviving restarts. - scope = str(params.get("scope") or "session").strip().lower() + scope = _word(params.get("scope") or "session") try: from tools.approval import disable_session_yolo, enable_session_yolo, is_session_yolo_enabled - - raw = str(value or "").strip().lower() + raw = _word(value) def _resolve_toggle(current: bool) -> bool: return _BOOL_WORDS.get(raw, not current) if scope == "global": from tools.approval import _normalize_approval_mode - - cfg = _load_cfg() - appr = cfg.get("approvals") if isinstance(cfg, dict) else None + appr = _load_cfg().get("approvals") appr = appr if isinstance(appr, dict) else {} enable = _resolve_toggle(_normalize_approval_mode(appr.get("mode", "manual")) == "off") # Binary affordance: no restore of a prior "smart"/custom mode (those live in config.yaml). @@ -351,9 +336,8 @@ _REASONING_DISPLAY_WORDS = ( def _set_reasoning(rid, params, key, value, session): try: from hermes_constants import parse_reasoning_effort - - arg = str(value or "").strip().lower() - scope = str(params.get("scope") or "").strip().lower() + arg = _word(value) + scope = _word(params.get("scope")) for words, reported, fields, thinking, show in _REASONING_DISPLAY_WORDS: if arg in words: _write_display_sections(sections={"thinking": thinking}, **fields) @@ -382,7 +366,7 @@ def _set_reasoning(rid, params, key, value, session): def _set_details_mode(rid, params, key, value, session): - nv = str(value or "").strip().lower() + nv = _word(value) if nv not in _DETAIL_MODES: return _err(rid, 4002, f"unknown details_mode: {value}") _write_display_sections(sections={section: nv for section in _DETAIL_SECTION_NAMES}, details_mode=nv) @@ -395,7 +379,7 @@ def _set_details_section(rid, params, key, value, session): section = key.split(".", 1)[1] if section not in _DETAIL_SECTION_NAMES: return _err(rid, 4002, f"unknown section: {section}") - nv = str(value or "").strip().lower() + nv = _word(value) if not nv: _write_display_sections(drop_sections=(section,)) elif nv not in _DETAIL_MODES: @@ -406,7 +390,7 @@ def _set_details_section(rid, params, key, value, session): def _set_thinking_mode(rid, params, key, value, session): - nv = str(value or "").strip().lower() + nv = _word(value) if nv not in {"collapsed", "truncated", "full"}: return _err(rid, 4002, f"unknown thinking_mode: {value}") _write_config_key("display.thinking_mode", nv) @@ -427,7 +411,7 @@ def _set_battery(rid, params, key, value, session): def _set_theme(rid, params, key, value, session): # 'light'/'dark' pin beats background auto-detection (xterm.js hosts misreport OSC 11). - raw = str(value or "").strip().lower() + raw = _word(value) if raw not in {"auto", "light", "dark"}: return _err(rid, 4002, f"unknown theme value: {value} (use auto|light|dark)") _write_config_key("display.tui_theme", raw) @@ -435,14 +419,12 @@ def _set_theme(rid, params, key, value, session): def _set_statusbar(rid, params, key, value, session): - raw = str(value or "").strip().lower() - current = _coerce_statusbar(_display_section(_load_cfg()).get("tui_statusbar", "top")) + raw = _word(value) + current = _coerce_statusbar(_display_cfg().get("tui_statusbar", "top")) if raw in {"", "toggle"}: nv = "top" if current == "off" else "off" - elif raw == "on": - nv = "top" - elif raw in _STATUSBAR_MODES: - nv = raw + elif raw == "on" or raw in _STATUSBAR_MODES: + nv = "top" if raw == "on" else raw else: return _err(rid, 4002, f"unknown statusbar value: {value}") _write_config_key("display.tui_statusbar", nv) @@ -453,7 +435,7 @@ def _set_mouse(rid, params, key, value, session): # Explicit None check (not `value or ""`) so falsy non-string inputs (0, False from # programmatic callers) reach the alias map as themselves (-> 'off') instead of toggling. raw = ("" if value is None else str(value)).strip().lower() - current = _display_mouse_tracking(_display_section(_load_cfg())) + current = _display_mouse_tracking(_display_cfg()) if raw in {"", "toggle"}: nv = "all" if current == "off" else "off" elif raw in _MOUSE_TRACKING_ALIASES: @@ -501,7 +483,6 @@ def _set_prompt_like(rid, params, key, value, session): # Personality text is an in-session overlay; persistence goes through # hermes_cli.personality (single owner), never the user-owned global system prompt. from hermes_cli.personality import persist_personality - persist_personality(pname) resp["value"] = str(value or "none") history_reset, info = _apply_personality_to_session(params.get("session_id", ""), session, new_prompt, pname) diff --git a/tui_gateway/methods_profiles.py b/tui_gateway/methods_profiles.py index f1c88105c2..a623f4d90f 100644 --- a/tui_gateway/methods_profiles.py +++ b/tui_gateway/methods_profiles.py @@ -169,11 +169,8 @@ def _canonical_session_row(db, profile_path): tip_row = db.get_session(tip) or row started = row.get("started_at") or 0 return { - "id": session_id, - "resolved_id": tip, - "root_title": row.get("title") or "", - "title": tip_row.get("title") or "", - "preview": _latest_message_preview(db, tip), + "id": session_id, "resolved_id": tip, "root_title": row.get("title") or "", + "title": tip_row.get("title") or "", "preview": _latest_message_preview(db, tip), "started_at": tip_row.get("started_at") or started, "last_active": tip_row.get("last_activity_at") or tip_row.get("started_at") or started, "message_count": tip_row.get("message_count") or 0, @@ -199,16 +196,13 @@ def _latest_profile_session_rows(db): continue if human is not None: continue + # Rosters want "where the conversation IS": prefer the newest text. human = { - "id": s["id"], - "title": title, - "preview": s.get("preview") or "", - "started_at": s.get("started_at") or 0, - "last_active": last_active, + "id": s["id"], "title": title, + "preview": _latest_message_preview(db, s["id"]) or s.get("preview") or "", + "started_at": s.get("started_at") or 0, "last_active": last_active, "message_count": s.get("message_count") or 0, } - # Rosters want "where the conversation IS": prefer the newest text. - human["preview"] = _latest_message_preview(db, s["id"]) or human["preview"] if worker is not None: break return human, worker @@ -259,14 +253,10 @@ def _(rid, params: dict) -> dict: out = [] for p in list_profiles(): row = { - "name": p.name, - "path": str(p.path), - "is_default": bool(p.is_default), - "model": p.model, - "provider": p.provider, - "description": getattr(p, "description", "") or "", - "display_name": getattr(p, "display_name", "") or "", - "skill_count": getattr(p, "skill_count", 0) or 0, + "name": p.name, "path": str(p.path), "is_default": bool(p.is_default), + "model": p.model, "provider": p.provider, + "description": p.description or "", "display_name": p.display_name or "", + "skill_count": p.skill_count or 0, } if include_sessions: _profile_session_fields(row, p.path) @@ -413,10 +403,8 @@ def _(rid, params: dict) -> dict: model_set = _best_effort(lambda: _pin_profile_model(path, provider, model)) elif is_truthy_value(params.get("mirror_credentials", True)): mirrored["model_inherited"] = _try(lambda: _inherit_launch_model(path), False) - return _ok( - rid, - {"ok": True, "name": name, "path": str(path), "soul_written": soul_written, "model_set": model_set, "mirrored": mirrored}, - ) + return _ok(rid, {"ok": True, "name": name, "path": str(path), "soul_written": soul_written, + "model_set": model_set, "mirrored": mirrored}) def _describe_toolsets(cfg): @@ -437,9 +425,8 @@ def _describe_toolsets(cfg): if (ts_name in default_off or ts_name == "yuanbao") and not enabled: continue tool_count = _try(lambda: len(set(resolve_toolset(ts_name))), 0) - toolsets_out.append( - {"name": ts_name, "label": ts_label, "description": ts_desc or "", "tool_count": tool_count, "enabled": enabled} - ) + toolsets_out.append({"name": ts_name, "label": ts_label, "description": ts_desc or "", + "tool_count": tool_count, "enabled": enabled}) return toolsets_out, pinned_set @@ -448,19 +435,11 @@ def _describe_mcp_servers(cfg): mcp_cfg = cfg.get("mcp_servers") if not isinstance(mcp_cfg, dict): return [] - return _try( - lambda: [ - { - "name": str(srv_name), - "enabled": not is_truthy_value(entry.get("disabled", False)), - "transport": str(entry.get("transport") or "http") if entry.get("url") else "stdio", - } - for srv_name in sorted(mcp_cfg.keys()) - for entry in (mcp_cfg[srv_name],) - if isinstance(entry, dict) - ], - [], - ) + return _try(lambda: [ + {"name": str(srv_name), "enabled": not is_truthy_value(entry.get("disabled", False)), + "transport": str(entry.get("transport") or "http") if entry.get("url") else "stdio"} + for srv_name in sorted(mcp_cfg.keys()) for entry in (mcp_cfg[srv_name],) if isinstance(entry, dict) + ], []) @_profile_handler("profiles.describe", 5063) @@ -485,17 +464,13 @@ def _(rid, params: dict) -> dict: mcp_out = _describe_mcp_servers(cfg) model_cfg = cfg.get("model") if isinstance(cfg.get("model"), dict) else {} meta = _try(lambda: _lazy("hermes_cli.profiles", "read_profile_meta")(profile_dir), {}) - result = { - "name": name, - "description": str(meta.get("description") or ""), - "soul": soul, - "model": {"provider": str(model_cfg.get("provider") or ""), "default": str(model_cfg.get("default") or "")}, - "skills": installed, - "toolsets": toolsets_out, - "toolsets_pinned": pinned_set is not None, + return _ok(rid, { + "name": name, "description": str(meta.get("description") or ""), "soul": soul, + "model": {"provider": str(model_cfg.get("provider") or ""), + "default": str(model_cfg.get("default") or "")}, + "skills": installed, "toolsets": toolsets_out, "toolsets_pinned": pinned_set is not None, "mcp_servers": mcp_out, - } - return _ok(rid, result) + }) def _configure_ui_meta(profile_dir, params, applied) -> None: @@ -555,14 +530,9 @@ def _configure_model(profile_dir, params, applied): if not (model and provider): return None if not is_truthy_value(params.get("confirm_expensive_model", False)): - confirm_message = _try( - lambda: getattr( - _lazy("hermes_cli.model_selection_guards", "combined_selection_warning")(model, provider=provider or None), - "message", - None, - ), - None, - ) + confirm_message = _try(lambda: getattr( + _lazy("hermes_cli.model_selection_guards", "combined_selection_warning")(model, provider=provider or None), + "message", None), None) if confirm_message is None: applied["model"] = _best_effort(lambda: _pin_profile_model(profile_dir, provider, model)) return confirm_message @@ -589,9 +559,8 @@ def _configure_cfg_sections(profile_dir, params, applied) -> None: if isinstance(params.get("enabled_toolsets"), list): applied["toolsets"] = _best_effort(lambda: _save_toolset_pin(cfg, params["enabled_toolsets"], save_config)) if want_mcp: - applied["mcp_servers"] = _best_effort( - lambda: _save_mcp_toggles(load_config() or {}, params["enabled_mcp_servers"], launch_mcp, save_config) - ) + applied["mcp_servers"] = _best_effort(lambda: _save_mcp_toggles( + load_config() or {}, params["enabled_mcp_servers"], launch_mcp, save_config)) def _clean_names(values) -> set: @@ -640,11 +609,8 @@ def _(rid, params: dict) -> dict: if isinstance(params.get("soul"), str): applied["soul"] = _best_effort(lambda: (profile_dir / "SOUL.md").write_text(params["soul"], encoding="utf-8")) if isinstance(params.get("description"), str): - applied["description"] = _best_effort( - lambda: _lazy("hermes_cli.profiles", "write_profile_meta")( - profile_dir, description=params["description"].strip(), description_auto=False - ) - ) + applied["description"] = _best_effort(lambda: _lazy("hermes_cli.profiles", "write_profile_meta")( + profile_dir, description=params["description"].strip(), description_auto=False)) confirm_message = _configure_model(profile_dir, params, applied) if any(isinstance(params.get(k), list) for k in ("disabled_skills", "enabled_toolsets", "enabled_mcp_servers")): _configure_cfg_sections(profile_dir, params, applied) diff --git a/tui_gateway/methods_voice.py b/tui_gateway/methods_voice.py index ef5d28ead6..c420022b15 100644 --- a/tui_gateway/methods_voice.py +++ b/tui_gateway/methods_voice.py @@ -6,6 +6,7 @@ method_ctx.bind_module), so they reference server.py globals bare. from __future__ import annotations +import contextlib import threading from .method_ctx import HandlerRegistry, bind_module @@ -58,17 +59,12 @@ def _end_voice_chat(*, stop_loop: bool, stop_tts: bool) -> None: os.environ["HERMES_VOICE"] = "0" os.environ["HERMES_VOICE_TTS"] = "0" if stop_loop: - try: + with contextlib.suppress(Exception): from hermes_cli.voice import stop_continuous - stop_continuous() - except Exception: - pass if stop_tts: - try: + with contextlib.suppress(Exception): _tts_stream_stop(user_barge=False) - except Exception: - pass def _tts_lease_async(lease: str, active: bool) -> None: @@ -79,11 +75,7 @@ def _tts_lease_async(lease: str, active: bool) -> None: def _run(): try: from tools.tts_tool import acquire_tts_lease, release_tts_lease - - if active: - acquire_tts_lease(lease) - else: - release_tts_lease(lease) + (acquire_tts_lease if active else release_tts_lease)(lease) except Exception as e: logger.debug("voice: tts lease %s active=%s failed: %s", lease, active, e) @@ -114,27 +106,18 @@ def _tts_stream_begin() -> Optional[queue.Queue]: return None try: from tools.tts_tool import check_tts_requirements, stream_tts_to_speaker - if not check_tts_requirements(): return None except Exception: return None - _tts_stream_stop() text_queue: queue.Queue = queue.Queue() - stop = threading.Event() - done = threading.Event() - threading.Thread( - target=stream_tts_to_speaker, args=(text_queue, stop, done), daemon=True - ).start() - + stop, done = threading.Event(), threading.Event() + threading.Thread(target=stream_tts_to_speaker, args=(text_queue, stop, done), daemon=True).start() global _tts_stream_state with _tts_stream_lock: _tts_stream_state = {"stop": stop, "done": done} - - if _voice_mode_enabled() and _voice_cfg_dict().get("barge_in", True): - _arm_full_duplex_listener() - + _arm_barge_listener_if_enabled() return text_queue @@ -148,21 +131,14 @@ def _tts_stream_stop(user_barge: bool = True) -> None: return if user_barge and not state["done"].is_set(): import traceback as _tb - logger.debug( - "TTS CUT: _tts_stream_stop(user_barge=True) — new turn or " - "interrupt cutting in-flight TTS\n%s", - "".join(_tb.format_stack()), - ) + logger.debug("TTS CUT: _tts_stream_stop(user_barge=True) — new turn or " + "interrupt cutting in-flight TTS\n%s", "".join(_tb.format_stack())) from tools.tts_streaming import mark_speech_interrupted - mark_speech_interrupted() state["stop"].set() - try: + with contextlib.suppress(Exception): from tools.voice_mode import stop_playback - stop_playback() - except Exception: - pass # ── Full-duplex agent-turn listener (one mic, whole turn) ────────────────── @@ -187,15 +163,24 @@ def _arm_full_duplex_listener() -> None: threading.Thread(target=_full_duplex_listener, daemon=True, name="voice-full-duplex").start() +def _arm_barge_listener_if_enabled() -> None: + """Arm the listener when voice mode is on and ``voice.barge_in`` isn't disabled.""" + if _voice_mode_enabled() and _voice_cfg_dict().get("barge_in", True): + _arm_full_duplex_listener() + + +def _fd_speak_pipelines_snapshot() -> list: + with _fd_listener_lock: + return list(_fd_speak_pipelines) + + def _fd_tts_pending() -> bool: """True while any TTS (streaming pipeline or fallback speak) is unfinished.""" with _tts_stream_lock: state = _tts_stream_state if state is not None and not state["done"].is_set(): return True - with _fd_listener_lock: - pipelines = list(_fd_speak_pipelines) - return any(not done.is_set() for _stop, done in pipelines) + return any(not done.is_set() for _stop, done in _fd_speak_pipelines_snapshot()) def _full_duplex_listener() -> None: @@ -209,99 +194,34 @@ def _full_duplex_listener() -> None: """ global _fd_listener_active try: - from tools.tts_streaming import mark_speech_interrupted - from tools.voice_mode import ( - full_duplex_listen, - is_audio_output_active, - stop_playback, - transcribe_recording, - ) - - cfg = _voice_cfg_dict() - try: - _mult = float(cfg.get("barge_in_threshold_multiplier", 0) or 0) - except (TypeError, ValueError): - _mult = 0.0 - try: - _grace_ms = int(float(cfg.get("barge_in_grace_seconds", 0.5)) * 1000) - except (TypeError, ValueError): - _grace_ms = 500 + from tools.voice_mode import full_duplex_listen, is_audio_output_active, transcribe_recording def _should_stop() -> bool: if not _voice_mode_enabled(): return True - if _any_session_running(): - return False - if _fd_tts_pending(): + if _any_session_running() or _fd_tts_pending(): return False return not is_audio_output_active() tripped = threading.Event() - def _cut_all_tts() -> None: - _tts_stream_stop(user_barge=True) - with _fd_listener_lock: - pipelines = list(_fd_speak_pipelines) - for _stop, _done in pipelines: - _stop.set() - stop_playback() - def _on_trigger(phase: str) -> None: tripped.set() - mark_speech_interrupted() - if phase == "playback": - logger.debug("TTS CUT: full-duplex listener tripped during playback") - _cut_all_tts() - else: - logger.debug( - "full-duplex listener tripped during generation — " - "interrupting running turn(s)" - ) - # Cut pending TTS FIRST so the stale reply can never speak. - _cut_all_tts() - try: - with _sessions_lock: - running = [s for s in _sessions.values() if s.get("running")] - for s in running: - agent = s.get("agent") - if agent is not None and hasattr(agent, "interrupt"): - try: - agent.interrupt() - except Exception: - pass - except Exception as e: - logger.debug("voice interjection interrupt failed: %s", e) - _voice_emit("voice.interrupted") + _fd_trip(phase) - wav_path = full_duplex_listen( - _should_stop, is_playing=is_audio_output_active, on_trigger=_on_trigger, - multiplier=_mult or None, grace_ms=max(0, _grace_ms), - ) + mult, grace_ms = _fd_barge_params(_voice_cfg_dict()) + wav_path = full_duplex_listen(_should_stop, is_playing=is_audio_output_active, + on_trigger=_on_trigger, multiplier=mult or None, grace_ms=grace_ms) if not (wav_path and tripped.is_set()): return try: result = transcribe_recording(wav_path) text = (result.get("transcript") or "").strip() if result.get("success") else "" if text: - # Stop-check must never break transcript delivery (stubbed - # voice_mode in tests, partial installs) — treat as not-a-stop. - try: - from tools.voice_mode import is_voice_stop_phrase - _is_stop = is_voice_stop_phrase(text) - except Exception: - _is_stop = False - - if _is_stop: - # Turn already interrupted / TTS cut at trip time; now end the chat. - _end_voice_chat(stop_loop=True, stop_tts=False) - _voice_emit("voice.transcript", {"stop_phrase": True, "text": text}) - else: - _voice_emit("voice.transcript", {"text": text}) + _deliver_fd_transcript(text) finally: - try: + with contextlib.suppress(OSError): os.unlink(wav_path) - except OSError: - pass except Exception as e: logger.debug("full-duplex listener failed: %s", e) finally: @@ -309,14 +229,79 @@ def _full_duplex_listener() -> None: _fd_listener_active = False +def _fd_barge_params(cfg: dict) -> tuple[float, int]: + """``(threshold multiplier, grace ms)`` from the voice config; malformed -> defaults.""" + try: + mult = float(cfg.get("barge_in_threshold_multiplier", 0) or 0) + except (TypeError, ValueError): + mult = 0.0 + try: + grace_ms = int(float(cfg.get("barge_in_grace_seconds", 0.5)) * 1000) + except (TypeError, ValueError): + grace_ms = 500 + return mult, max(0, grace_ms) + + +def _cut_all_tts() -> None: + """Cut streaming TTS, every fallback speak pipeline, and the file player.""" + from tools.voice_mode import stop_playback + + _tts_stream_stop(user_barge=True) + for _stop, _done in _fd_speak_pipelines_snapshot(): + _stop.set() + stop_playback() + + +def _fd_trip(phase: str) -> None: + """Listener tripped: latch the interruption, cut TTS, and during generation also + interrupt every running turn (the ``agent.interrupt()`` seam ``session.interrupt`` uses).""" + from tools.tts_streaming import mark_speech_interrupted + + mark_speech_interrupted() + if phase == "playback": + logger.debug("TTS CUT: full-duplex listener tripped during playback") + _cut_all_tts() + else: + logger.debug("full-duplex listener tripped during generation — interrupting running turn(s)") + # Cut pending TTS FIRST so the stale reply can never speak. + _cut_all_tts() + try: + with _sessions_lock: + running = [s for s in _sessions.values() if s.get("running")] + for s in running: + agent = s.get("agent") + if agent is not None and hasattr(agent, "interrupt"): + with contextlib.suppress(Exception): + agent.interrupt() + except Exception as e: + logger.debug("voice interjection interrupt failed: %s", e) + _voice_emit("voice.interrupted") + + +def _deliver_fd_transcript(text: str) -> None: + """Emit the captured interjection; a bare stop phrase also ends the voice chat.""" + # Stop-check must never break transcript delivery (stubbed voice_mode in tests, + # partial installs) — treat as not-a-stop. + try: + from tools.voice_mode import is_voice_stop_phrase + is_stop = is_voice_stop_phrase(text) + except Exception: + is_stop = False + if is_stop: + # Turn already interrupted / TTS cut at trip time; now end the chat. + _end_voice_chat(stop_loop=True, stop_tts=False) + _voice_emit("voice.transcript", {"stop_phrase": True, "text": text}) + else: + _voice_emit("voice.transcript", {"text": text}) + + def _speak_text_with_barge(text: str) -> None: """Speak via hermes_cli.voice.speak_text with spoken barge-in: the (stop, done) pair is registered in ``_fd_speak_pipelines`` so the full-duplex listener can cut it on a playback trip and keeps listening while it is pending.""" from hermes_cli.voice import speak_text - stop = threading.Event() - done = threading.Event() + stop, done = threading.Event(), threading.Event() with _fd_listener_lock: _fd_speak_pipelines.add((stop, done)) @@ -332,8 +317,7 @@ def _speak_text_with_barge(text: str) -> None: _fd_speak_pipelines.discard((stop, done)) threading.Thread(target=_speak, daemon=True).start() - if _voice_mode_enabled() and _voice_cfg_dict().get("barge_in", True): - _arm_full_duplex_listener() + _arm_barge_listener_if_enabled() def _voice_cfg_dict() -> dict: @@ -341,7 +325,6 @@ def _voice_cfg_dict() -> dict: root and ``voice`` may be any YAML scalar/list/None; malformed → {}.""" cfg = _load_cfg() voice_cfg = cfg.get("voice") if isinstance(cfg, dict) else None - return voice_cfg if isinstance(voice_cfg, dict) else {} @@ -354,7 +337,6 @@ def _voice_cfg_number(value, default): def _voice_record_key() -> str: """Current ``voice.record_key`` value, documented default on error.""" record_key = _voice_cfg_dict().get("record_key") - return str(record_key) if isinstance(record_key, str) and record_key else "ctrl+b" @@ -374,17 +356,21 @@ def _wake_owner_snapshot(): return _wake_owner_transport, _wake_owner_surface +def _set_wake_owner(transport, surface: str) -> None: + global _wake_owner_transport, _wake_owner_surface + with _wake_lock: + _wake_owner_transport, _wake_owner_surface = transport, surface + + def _release_wake_for_transport(transport: "Transport") -> bool: """Release the wake lease iff ``transport`` is the current gateway owner.""" global _wake_owner_transport, _wake_owner_surface with _wake_lock: if _wake_owner_transport is not transport: return False - _wake_owner_transport = None - _wake_owner_surface = "" + _wake_owner_transport, _wake_owner_surface = None, "" try: from tools.wake_word import stop_listening - stop_listening(owner=transport) except Exception as e: logger.debug("wake stop failed: %s", e) @@ -438,10 +424,8 @@ def _wake_resume_if_owner(owner: "Transport", *, retry_seconds: float = 15.0, continue # False — detector gone or lease moved: stop, don't fight it. return - logger.warning( - "wake: could not resume detector after voice turn " - "(microphone still busy?) — toggle the wake word to re-arm" - ) + logger.warning("wake: could not resume detector after voice turn " + "(microphone still busy?) — toggle the wake word to re-arm") finally: with _wake_resume_retry_lock: _wake_resume_retry_active = False @@ -455,7 +439,6 @@ def _persist_wake_enabled(enabled: bool) -> bool: (ear toggle, /wake on|off) — never passive auto-arm paths.""" try: from cli import save_config_value - return bool(save_config_value("wake_word.enabled", enabled)) except Exception as e: logger.warning("wake: failed to persist wake_word.enabled=%s: %s", enabled, e) @@ -474,9 +457,7 @@ def _wake_probe(cfg: dict, prefer_client: bool) -> tuple[str, dict]: from tools.wake_word import check_wake_word_requirements, resolve_capture_mode capture_mode = resolve_capture_mode(cfg, prefer_client=prefer_client) - probe_cfg = dict(cfg) - probe_cfg["capture"] = capture_mode - return capture_mode, check_wake_word_requirements(probe_cfg) + return capture_mode, check_wake_word_requirements({**cfg, "capture": capture_mode}) def _wake_detect_handler(transport, sid: str, phrase: str, new_session: bool): @@ -486,9 +467,7 @@ def _wake_detect_handler(transport, sid: str, phrase: str, new_session: bool): def _on_detect() -> None: from tools.wake_word import get_last_match, owns_listener, pause_listening - if not pause_listening(owner=transport): - return - if not owns_listener(transport): + if not pause_listening(owner=transport) or not owns_listener(transport): return if _transport_is_dead(transport): _release_wake_for_transport(transport) @@ -517,7 +496,6 @@ def _(rid, params: dict) -> dict: withholds unless the guarantee is advertised. Sourced from the enforcing module, never config: a believed-but-absent capability is worse than none.""" from hermes_cli.active_sessions import PER_SESSION_EXCLUSIVE_SUBMIT - return _ok(rid, {"per_session_exclusive_submit": bool(PER_SESSION_EXCLUSIVE_SUBMIT)}) @@ -543,13 +521,8 @@ def _(rid, params: dict) -> dict: transport = _caller_transport() try: from tools.wake_word import ( - WakeWordInUse, - detector_frame_info, - load_wake_word_config, - owns_listener, - start_listening, - wake_phrase, - wake_surface_enabled, + WakeWordInUse, detector_frame_info, load_wake_word_config, owns_listener, start_listening, + wake_phrase, wake_surface_enabled, ) except Exception as e: return _err(rid, 5026, f"wake module unavailable: {e}") @@ -565,12 +538,9 @@ def _(rid, params: dict) -> dict: "started": False, "reason": "unavailable", "hint": reqs.get("hint") or "", "capture": capture_mode, }) - enabled_persisted = False - if persist and not cfg.get("enabled"): - enabled_persisted = _persist_wake_enabled(True) - if enabled_persisted: - cfg = dict(cfg) - cfg["enabled"] = True + enabled_persisted = bool(persist and not cfg.get("enabled") and _persist_wake_enabled(True)) + if enabled_persisted: + cfg = {**cfg, "enabled": True} if not wake_surface_enabled(surface, cfg): # "disabled" (a persist:true retry can turn it on) vs "disabled_for_surface" # (explicit wake_word.surface choice, which persist does NOT override). @@ -580,36 +550,22 @@ def _(rid, params: dict) -> dict: return _ok(rid, {"started": False, "reason": reason}) existing_owner, existing_surface = _wake_owner_snapshot() - if existing_owner is not None and ( - _transport_is_dead(existing_owner) or not owns_listener(existing_owner) - ): + if existing_owner is not None and (_transport_is_dead(existing_owner) or not owns_listener(existing_owner)): _release_wake_for_transport(existing_owner) - existing_owner = None - existing_surface = "" + existing_owner, existing_surface = None, "" if existing_owner is not None and existing_owner is not transport: return _ok(rid, {"started": False, "reason": "owned", "owner_surface": existing_surface}) sid = str(params.get("session_id") or "") try: - start_listening( - _wake_detect_handler( - transport, sid, wake_phrase(cfg), bool(cfg.get("start_new_session", True)) - ), - owner=transport, - config=cfg, - external_audio=external_audio, - ) + on_detect = _wake_detect_handler(transport, sid, wake_phrase(cfg), bool(cfg.get("start_new_session", True))) + start_listening(on_detect, owner=transport, config=cfg, external_audio=external_audio) except WakeWordInUse: - return _ok(rid, { - "started": False, "reason": "owned", "owner_surface": existing_surface or None, - }) + return _ok(rid, {"started": False, "reason": "owned", "owner_surface": existing_surface or None}) except Exception as e: logger.warning("wake.start(%s): failed to start listener: %s", surface, e) return _err(rid, 5026, str(e)) - global _wake_owner_transport, _wake_owner_surface - with _wake_lock: - _wake_owner_transport = transport - _wake_owner_surface = surface + _set_wake_owner(transport, surface) frame = detector_frame_info() logger.info( "wake.start(%s): listening for %r (%s) capture=%s frame=%s", @@ -633,7 +589,6 @@ def _(rid, params: dict) -> dict: if bool(params.get("persist")): try: from tools.wake_word import load_wake_word_config - currently_enabled = bool(load_wake_word_config().get("enabled")) except Exception: currently_enabled = True @@ -651,7 +606,6 @@ def _(rid, params: dict) -> dict: transport = _caller_transport() try: from tools.wake_word import pause_listening - paused = pause_listening(owner=transport) logger.info("wake.pause: detector paused=%s", paused) except Exception as e: @@ -672,18 +626,12 @@ def _(rid, params: dict) -> dict: def _(rid, params: dict) -> dict: try: from tools.wake_word import ( - audio_is_silent, - detector_frame_info, - get_input_device_status, - is_listening, - load_wake_word_config, - owns_listener, - silent_audio_hint, + audio_is_silent, detector_frame_info, get_input_device_status, is_listening, + load_wake_word_config, owns_listener, silent_audio_hint, ) cfg = load_wake_word_config() - probe_capture, reqs = _wake_probe( - cfg, _wake_prefers_client(params, str(params.get("surface") or "").strip().lower()) - ) + surface = str(params.get("surface") or "").strip().lower() + probe_capture, reqs = _wake_probe(cfg, _wake_prefers_client(params, surface)) transport = _caller_transport() owner, owner_surface = _wake_owner_snapshot() owned_by_caller = owns_listener(transport) @@ -706,23 +654,17 @@ def _(rid, params: dict) -> dict: else: capture = probe_capture or reqs.get("capture") or str(cfg.get("capture") or "auto") return _ok(rid, { - "listening": listening, - "owned_by_caller": owned_by_caller, + "listening": listening, "owned_by_caller": owned_by_caller, "owner_surface": owner_surface if owner is not None else None, - "phrase": reqs["phrase"], - "provider": reqs["provider"], + "phrase": reqs["phrase"], "provider": reqs["provider"], "configured_surface": str(cfg.get("surface") or "auto"), - "input_device": input_device, - "available": reqs["available"], - "hint": hint, + "input_device": input_device, "available": reqs["available"], "hint": hint, # Config truth: clients re-arm after a voice turn ("permanent on") from this. "enabled": bool(cfg.get("enabled")), # Armed but deaf despite an open stream; see platform-specific hint. - "audio_silent": silent, - "capture": capture, + "audio_silent": silent, "capture": capture, "local_input_available": bool(reqs.get("local_input_available")), - "sample_rate": frame.get("sample_rate", 16000), - "frame_length": frame.get("frame_length", 1280), + "sample_rate": frame.get("sample_rate", 16000), "frame_length": frame.get("frame_length", 1280), }) except Exception as e: return _err(rid, 5026, str(e)) @@ -762,19 +704,12 @@ def _(rid, params: dict) -> dict: def _voice_toggle_status(rid, params: dict) -> dict: # Mirrors CLI _show_voice_status: STT/TTS availability tells the user WHY # voice isn't working; record_key lets the TUI bind and display the shortcut. - payload: dict = { - "enabled": _voice_mode_enabled(), - "record_key": _voice_record_key(), - "tts": _voice_tts_enabled(), - } + payload: dict = {"enabled": _voice_mode_enabled(), "record_key": _voice_record_key(), "tts": _voice_tts_enabled()} try: from tools.voice_mode import check_voice_requirements - reqs = check_voice_requirements() - payload["available"] = bool(reqs.get("available")) - payload["audio_available"] = bool(reqs.get("audio_available")) - payload["stt_available"] = bool(reqs.get("stt_available")) - payload["details"] = reqs.get("details") or "" + payload.update(available=bool(reqs.get("available")), audio_available=bool(reqs.get("audio_available")), + stt_available=bool(reqs.get("stt_available")), details=reqs.get("details") or "") except Exception as e: # Optional transcription deps — /voice status must always answer. logger.warning("voice.toggle status: requirements probe failed: %s", e) @@ -790,44 +725,30 @@ def _voice_toggle_mode(rid, params: dict) -> dict: stop_hint = "" if enabled: - # Spoken-stop hint for the client; sourced from voice.stop_phrases, - # empty when the feature is disabled. + # Spoken-stop hint for the client; sourced from voice.stop_phrases, empty when disabled. try: from tools.voice_mode import voice_stop_hint - stop_hint = voice_stop_hint() except Exception: stop_hint = "" - # Speech output already on → warm the engine now, not on the first reply. if _voice_tts_enabled(): _tts_lease_async("tui:voice-tts", True) - - if not enabled: + else: # The continuous loop holds the microphone; tear it down with the mode. try: from hermes_cli.voice import stop_continuous - stop_continuous() except ImportError: pass except Exception as e: logger.warning("voice: stop_continuous failed during toggle off: %s", e) - # Clear TTS so it can be toggled independently later; silence live speech. os.environ["HERMES_VOICE_TTS"] = "0" _tts_stream_stop(user_barge=False) _tts_lease_async("tui:voice-tts", False) - - return _ok( - rid, - { - "enabled": enabled, - "record_key": _voice_record_key(), - "tts": _voice_tts_enabled(), - "stop_hint": stop_hint, - }, - ) + return _ok(rid, {"enabled": enabled, "record_key": _voice_record_key(), "tts": _voice_tts_enabled(), + "stop_hint": stop_hint}) def _voice_toggle_tts(rid, params: dict) -> dict: @@ -907,63 +828,46 @@ def _(rid, params: dict) -> dict: try: global _voice_event_sid, _voice_wake_owner - if action == "start": - if not _voice_mode_enabled(): - return _err(rid, 4015, "voice mode is off — enable with /voice on") - - with _voice_sid_lock: - _voice_event_sid = params.get("session_id") or _voice_event_sid - - from hermes_cli.voice import start_continuous - - # Busy probe holds the no-speech counter during long agent turns. - # Safe to re-register every start; older wrappers lack the setter. - try: - from hermes_cli.voice import set_voice_busy_probe - - set_voice_busy_probe(_any_session_running) - except Exception: - pass - - # Shape-safe: malformed voice YAML falls back to documented defaults. - voice_cfg = _voice_cfg_dict() - safe_threshold = _voice_cfg_number(voice_cfg.get("silence_threshold"), 200) - safe_duration = _voice_cfg_number(voice_cfg.get("silence_duration"), 3.0) - # max_recording_seconds: explicit numeric <= 0 disables the cap (0.0). - max_rec = _voice_cfg_number(voice_cfg.get("max_recording_seconds"), 120.0) - safe_max_rec = max_rec if max_rec > 0 else 0.0 - # Hand the mic to STT if the wake detector holds it; resume on a - # terminal capture event so wake-triggered and manual captures coexist. - try: - from tools.wake_word import pause_listening - - wake_paused = pause_listening(owner=transport) - except Exception: - wake_paused = False - if wake_paused: - with _voice_sid_lock: - _voice_wake_owner = transport - - started = start_continuous( - on_transcript=_vr_on_transcript, on_status=_vr_on_status, - on_silent_limit=_vr_on_silent, silence_threshold=safe_threshold, - silence_duration=safe_duration, auto_restart=False, - max_recording_seconds=safe_max_rec, on_stop_phrase=_vr_on_stop_phrase, - ) - if started is False: - _resume_voice_wake() - return _ok(rid, {"status": "busy"}) - return _ok(rid, {"status": "recording"}) - - # action == "stop" + if action == "start" and not _voice_mode_enabled(): + return _err(rid, 4015, "voice mode is off — enable with /voice on") with _voice_sid_lock: _voice_event_sid = params.get("session_id") or _voice_event_sid + if action == "stop": + from hermes_cli.voice import stop_continuous + stop_continuous(force_transcribe=True) + _resume_voice_wake() + return _ok(rid, {"status": "stopped"}) - from hermes_cli.voice import stop_continuous - - stop_continuous(force_transcribe=True) - _resume_voice_wake() - return _ok(rid, {"status": "stopped"}) + from hermes_cli.voice import start_continuous + # Busy probe holds the no-speech counter during long agent turns. + # Safe to re-register every start; older wrappers lack the setter. + with contextlib.suppress(Exception): + from hermes_cli.voice import set_voice_busy_probe + set_voice_busy_probe(_any_session_running) + # Shape-safe: malformed voice YAML falls back to documented defaults. + # max_recording_seconds: explicit numeric <= 0 disables the cap (0.0). + voice_cfg = _voice_cfg_dict() + max_rec = _voice_cfg_number(voice_cfg.get("max_recording_seconds"), 120.0) + # Hand the mic to STT if the wake detector holds it; resume on a terminal + # capture event so wake-triggered and manual captures coexist. + try: + from tools.wake_word import pause_listening + wake_paused = pause_listening(owner=transport) + except Exception: + wake_paused = False + if wake_paused: + with _voice_sid_lock: + _voice_wake_owner = transport + started = start_continuous( + on_transcript=_vr_on_transcript, on_status=_vr_on_status, on_silent_limit=_vr_on_silent, + silence_threshold=_voice_cfg_number(voice_cfg.get("silence_threshold"), 200), + silence_duration=_voice_cfg_number(voice_cfg.get("silence_duration"), 3.0), auto_restart=False, + max_recording_seconds=max_rec if max_rec > 0 else 0.0, on_stop_phrase=_vr_on_stop_phrase, + ) + if started is False: + _resume_voice_wake() + return _ok(rid, {"status": "busy"}) + return _ok(rid, {"status": "recording"}) except Exception as e: if wake_paused or action == "stop": _resume_voice_wake() @@ -981,10 +885,7 @@ def _(rid, params: dict) -> dict: # Import check up front so a missing voice module returns 5026 instead # of failing silently in the thread. import hermes_cli.voice # noqa: F401 - - threading.Thread( - target=_speak_text_with_barge, args=(text,), daemon=True - ).start() + threading.Thread(target=_speak_text_with_barge, args=(text,), daemon=True).start() return _ok(rid, {"status": "speaking"}) except ImportError: return _err(rid, 5026, "voice module not available")