refactor(tui_gateway): compact voice/config/profiles docstrings, hug payloads, squeeze body blanks
This commit is contained in:
@@ -1,12 +1,7 @@
|
||||
"""Config / projects / setup JSON-RPC handlers.
|
||||
|
||||
Handlers and module-level helpers are rebound onto server.py's globals at
|
||||
install time (see method_ctx.bind_module), so bodies reference server.py
|
||||
globals bare (``_ok``, ``_err``, ``_load_cfg``, ``_sessions``, ...).
|
||||
``config.set`` still lives in server.py.
|
||||
"""Config / projects / setup JSON-RPC handlers. Bodies are rebound onto server.py's globals
|
||||
(method_ctx.bind_module) and reference them bare. ``config.set`` lives in methods_config_set.
|
||||
"""
|
||||
|
||||
|
||||
from .method_ctx import HandlerRegistry, bind_module
|
||||
|
||||
from hermes_constants import DEFAULT_INDICATOR_STYLE, INDICATOR_STYLES
|
||||
@@ -18,8 +13,7 @@ _profile_scoped = _registry.profile_scoped
|
||||
|
||||
def _reconcile_repo_discovery(pdb, conn, policy, policy_key):
|
||||
pdb.reconcile_discovered_repos_policy(
|
||||
conn, policy_key, preserve_unversioned=_repo_discovery_policy_is_default(policy)
|
||||
)
|
||||
conn, policy_key, preserve_unversioned=_repo_discovery_policy_is_default(policy))
|
||||
|
||||
|
||||
@method("projects.discover_repos")
|
||||
@@ -34,9 +28,8 @@ def _(rid, params: dict) -> dict:
|
||||
policy = _repo_discovery_policy()
|
||||
with pdb.connect_closing() as conn:
|
||||
_reconcile_repo_discovery(pdb, conn, policy, _repo_discovery_policy_key(policy))
|
||||
# `scan=true` (desktop in remote-gateway mode): the desktop's
|
||||
# native scan only sees its local filesystem, so ask the host
|
||||
# to scan the policy roots itself so zero-session repos surface.
|
||||
# `scan=true` (remote-gateway desktop): its native scan only sees its own
|
||||
# filesystem, so the host scans the policy roots so zero-session repos surface.
|
||||
if params.get("scan") and policy["enabled"]:
|
||||
_scan_discovered_repos_remote(conn, policy)
|
||||
repos = _discover_repos_payload(db, conn=conn, include_cached=policy["enabled"])
|
||||
@@ -48,26 +41,24 @@ def _(rid, params: dict) -> dict:
|
||||
@method("projects.record_repos")
|
||||
@_profile_scoped
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""Persist git repo roots found by the client's filesystem scan (the native
|
||||
crawl runs on the desktop), then return the merged repo list."""
|
||||
"""Persist repo roots found by the client's (desktop-side) scan; return the merged 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_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 = (
|
||||
_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))
|
||||
pairs: list[tuple[str, str | None]] = []
|
||||
for item in params.get("repos") or []:
|
||||
if isinstance(item, str):
|
||||
pairs.append((item, None))
|
||||
elif isinstance(item, dict) and item.get("root"):
|
||||
pairs.append((str(item["root"]), item.get("label")))
|
||||
|
||||
with pdb.connect_closing() as conn:
|
||||
_reconcile_repo_discovery(pdb, conn, policy, policy_key)
|
||||
accepted = bool(policy["enabled"] and (incoming_matches or accept_legacy_default))
|
||||
@@ -75,9 +66,9 @@ def _(rid, params: dict) -> dict:
|
||||
pdb.record_discovered_repos(conn, pairs, replace=True, policy_key=policy_key)
|
||||
elif not policy["enabled"]:
|
||||
pdb.clear_discovered_repos(conn, policy_key=policy_key)
|
||||
|
||||
with _profile_db(params) as db:
|
||||
repos = _discover_repos_payload(db, include_cached=policy["enabled"]) if db is not None else []
|
||||
repos = ([] if db is None
|
||||
else _discover_repos_payload(db, include_cached=policy["enabled"]))
|
||||
return _ok(rid, {"repos": repos, "accepted": accepted, "discovery_policy": policy})
|
||||
except Exception as e:
|
||||
return _err(rid, 5061, str(e))
|
||||
@@ -94,23 +85,19 @@ def _stamped_project_tree(db, params, **kwargs):
|
||||
@method("projects.tree")
|
||||
@_profile_scoped
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""Authoritative project overview: project -> repo -> lane structure with
|
||||
counts + a few preview sessions per project, plus the flat set of session
|
||||
ids claimed by any project (so the desktop excludes them from flat Recents).
|
||||
Lanes carry no session rows here; drill-in uses ``projects.project_sessions``.
|
||||
"""
|
||||
"""Project -> repo -> lane overview with counts + a few preview sessions per project, plus
|
||||
the flat set of session ids claimed by any project (excluded from flat Recents). Lanes carry
|
||||
no session rows here; drill-in uses ``projects.project_sessions``."""
|
||||
try:
|
||||
with _profile_db(params) as db:
|
||||
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,
|
||||
)
|
||||
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"],
|
||||
})
|
||||
"scoped_session_ids": tree["scoped_session_ids"]})
|
||||
except Exception as e:
|
||||
return _err(rid, 5061, str(e))
|
||||
|
||||
@@ -118,33 +105,26 @@ def _(rid, params: dict) -> dict:
|
||||
@method("projects.project_sessions")
|
||||
@_profile_scoped
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""Fully hydrated lanes (repo -> lane -> session rows) for one project,
|
||||
built from the same authoritative grouping as ``projects.tree`` so ids and
|
||||
membership match exactly."""
|
||||
"""Fully hydrated lanes for one project, from the same grouping as ``projects.tree``."""
|
||||
try:
|
||||
project_id = str(params.get("project_id") or "")
|
||||
if not project_id:
|
||||
return _err(rid, 5063, "project_id required")
|
||||
|
||||
with _profile_db(params) as db:
|
||||
if db is None:
|
||||
return _ok(rid, {"project": None})
|
||||
# Drill-in only needs the entered project (which has sessions):
|
||||
# skip the zero-session discovery tier.
|
||||
# Drill-in only needs the entered project: 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,
|
||||
)
|
||||
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})
|
||||
except Exception as e:
|
||||
return _err(rid, 5061, str(e))
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# config.get — one getter per key. Each returns the result payload (a dict) or
|
||||
# a full ``_err`` response (dicts containing "error" pass through untouched).
|
||||
# ---------------------------------------------------------------------------
|
||||
# ── config.get — one getter per key; returns the result payload or a full ``_err`` response
|
||||
# (dicts containing "error" pass through untouched).
|
||||
|
||||
|
||||
def _display_mode(cfg: dict, key: str, allowed: frozenset, default: str) -> str:
|
||||
@@ -163,8 +143,7 @@ def _cfg_get_provider(rid, params):
|
||||
return {
|
||||
"model": model,
|
||||
"provider": normalize_provider(parts[0]) if len(parts) > 1 else "unknown",
|
||||
"providers": list_available_providers(),
|
||||
}
|
||||
"providers": list_available_providers()}
|
||||
except Exception as e:
|
||||
return _err(rid, 5013, str(e))
|
||||
|
||||
@@ -182,16 +161,14 @@ def _cfg_get_project(rid, params):
|
||||
|
||||
|
||||
def _cfg_get_indicator(rid, params):
|
||||
# Normalize so a hand-edited config.yaml (stray casing / unknown value)
|
||||
# reads back the SAME value the TUI rendered (frontend falls back to
|
||||
# DEFAULT_INDICATOR_STYLE for the same inputs).
|
||||
# Normalize so a hand-edited config.yaml (stray casing / unknown value) reads back the SAME
|
||||
# value the TUI rendered (frontend falls back to DEFAULT_INDICATOR_STYLE for the same inputs).
|
||||
norm = str((_load_cfg().get("display") or {}).get("tui_status_indicator", "")).strip().lower()
|
||||
return {"value": norm if norm in INDICATOR_STYLES else DEFAULT_INDICATOR_STYLE}
|
||||
|
||||
|
||||
def _cfg_get_personality(rid, params):
|
||||
# EFFECTIVE personality via the single owner — a stale/unknown name in
|
||||
# config must not display as active.
|
||||
# EFFECTIVE personality via the single owner — a stale/unknown name must not show as active.
|
||||
from hermes_cli.personality import active_personality_name
|
||||
return {"value": active_personality_name(_load_cfg()) or "none"}
|
||||
|
||||
@@ -205,10 +182,8 @@ def _cfg_get_reasoning(rid, params):
|
||||
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"
|
||||
else:
|
||||
effort = str(reasoning_config.get("effort") or "medium")
|
||||
enabled = reasoning_config.get("enabled") is not False
|
||||
effort = str(reasoning_config.get("effort") or "medium") if enabled else "none"
|
||||
else:
|
||||
raw_effort = (cfg.get("agent") or {}).get("reasoning_effort", "")
|
||||
# YAML `reasoning_effort: false` means thinking disabled, not "unset".
|
||||
@@ -218,9 +193,8 @@ def _cfg_get_reasoning(rid, params):
|
||||
|
||||
|
||||
def _cfg_get_fast(rid, params):
|
||||
# `config.set fast` is session-scoped, so prefer the session's live/pinned
|
||||
# value over the global key; a pre-build session keeps its pin in
|
||||
# create_service_tier_override.
|
||||
# `config.set fast` is session-scoped: prefer the session's live/pinned value over the
|
||||
# global key (a pre-build session keeps its pin in create_service_tier_override).
|
||||
session = _sessions.get(params.get("session_id", ""))
|
||||
tier = None
|
||||
if session is not None:
|
||||
@@ -255,27 +229,20 @@ def _cfg_get_theme(rid, params):
|
||||
return {"value": raw if raw in {"auto", "light", "dark"} else "auto"}
|
||||
|
||||
|
||||
def _cfg_get_focus(rid, params):
|
||||
on = bool(_display_cfg().get("focus_view", False))
|
||||
return {"value": "on" if on else "off", "tool_progress": _load_tool_progress_mode()}
|
||||
|
||||
|
||||
def _cfg_get_mtime(rid, params):
|
||||
cfg_path = _hermes_home / "config.yaml"
|
||||
try:
|
||||
mtime = cfg_path.stat().st_mtime if cfg_path.exists() else 0
|
||||
except Exception:
|
||||
return {"mtime": 0}
|
||||
# mcp_rev: hash of the MCP-relevant config sections so the TUI's poller
|
||||
# reloads MCP servers only when their config changed — a /skin write bumps
|
||||
# mtime but must not cost a multi-second MCP reconnect.
|
||||
# mcp_rev: hash of the MCP-relevant sections so the poller reloads MCP servers only when
|
||||
# their config changed — a /skin write bumps mtime but must not cost an MCP reconnect.
|
||||
return {"mtime": mtime, "mcp_rev": _compute_mcp_rev()}
|
||||
|
||||
|
||||
def _config_getters() -> dict:
|
||||
"""key -> getter(rid, params). Built inside a function (not a module-level
|
||||
dict) so, once rebound onto server.py, every entry resolves to the rebound
|
||||
helper copies rather than this module's un-rebound originals."""
|
||||
"""key -> getter(rid, params). Built per call so, once rebound onto server.py, every entry
|
||||
resolves to the rebound helper copies rather than this module's originals."""
|
||||
return {
|
||||
"provider": _cfg_get_provider,
|
||||
"profile": _cfg_get_profile,
|
||||
@@ -291,20 +258,19 @@ 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", _DETAIL_MODES, "collapsed")
|
||||
},
|
||||
"value": _display_mode(_load_cfg(), "details_mode", _DETAIL_MODES, "collapsed")},
|
||||
"thinking_mode": _cfg_get_thinking_mode,
|
||||
"density": lambda rid, params: {
|
||||
"value": "on" if bool((_load_cfg().get("display") or {}).get("tui_compact", False)) else "off"
|
||||
},
|
||||
"theme": _cfg_get_theme,
|
||||
"statusbar": lambda rid, params: {
|
||||
"value": _coerce_statusbar(_display_cfg().get("tui_statusbar", "top"))
|
||||
},
|
||||
"focus": _cfg_get_focus,
|
||||
"value": _coerce_statusbar(_display_cfg().get("tui_statusbar", "top"))},
|
||||
"focus": lambda rid, params: {
|
||||
"value": "on" if bool(_display_cfg().get("focus_view", False)) else "off",
|
||||
"tool_progress": _load_tool_progress_mode()},
|
||||
"mouse": lambda rid, params: {"value": _display_mouse_tracking(_load_cfg().get("display"))},
|
||||
"mtime": _cfg_get_mtime,
|
||||
}
|
||||
"mtime": _cfg_get_mtime}
|
||||
|
||||
|
||||
@method("config.get")
|
||||
@@ -320,20 +286,14 @@ def _(rid, params: dict) -> dict:
|
||||
return _ok(rid, payload)
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# setup readiness
|
||||
# ---------------------------------------------------------------------------
|
||||
# ── setup readiness
|
||||
|
||||
|
||||
def _readiness_profile_scope(params: dict):
|
||||
"""Resolve the optional ``profile`` param of the setup readiness RPCs.
|
||||
|
||||
Returns ``(profile, scope)``: ``scope`` binds that profile's HERMES_HOME and
|
||||
``.env`` secret scope (ContextVars, so concurrent checks stay isolated); the
|
||||
launch profile / no param yields ``("", nullcontext())``. A profile unknown
|
||||
to this host raises ``FileNotFoundError`` — a readiness check must never
|
||||
quietly answer for the launch profile instead.
|
||||
"""
|
||||
"""``(profile, scope)`` for the readiness RPCs' optional ``profile`` param: ``scope`` binds
|
||||
that profile's HERMES_HOME + ``.env`` secret scope (ContextVars, so concurrent checks stay
|
||||
isolated); no param yields ``("", nullcontext())``. An unknown profile raises
|
||||
``FileNotFoundError`` — never quietly answer for the launch profile instead."""
|
||||
import contextlib
|
||||
profile = str(params.get("profile") or "").strip() if isinstance(params, dict) else ""
|
||||
if not profile:
|
||||
@@ -348,11 +308,8 @@ def _readiness_profile_scope(params: dict):
|
||||
|
||||
|
||||
def _readiness_check(rid, params, probe):
|
||||
"""Shared shell of setup.status / setup.runtime_check.
|
||||
|
||||
``probe(profile)`` runs inside the profile scope and returns the payload;
|
||||
an unknown profile answers ``ok=False`` (never a JSON-RPC error).
|
||||
"""
|
||||
"""Shared shell of setup.status / setup.runtime_check: ``probe(profile)`` runs inside the
|
||||
profile scope; an unknown profile answers ``ok=False`` (never a JSON-RPC error)."""
|
||||
try:
|
||||
profile, scope = _readiness_profile_scope(params)
|
||||
except FileNotFoundError as e:
|
||||
@@ -367,10 +324,10 @@ 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)))
|
||||
return {"provider_configured": configured, **({"profile": profile} if profile else {})}
|
||||
|
||||
return _readiness_check(rid, params, probe)
|
||||
except Exception as e:
|
||||
return _err(rid, 5016, str(e))
|
||||
@@ -379,15 +336,10 @@ def _(rid, params: dict) -> dict:
|
||||
@method("setup.runtime_check")
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""Strict provider check: does the configured/default model resolve to a usable runtime?
|
||||
|
||||
Unlike setup.status (True if ANY provider auth state is discoverable, incl.
|
||||
indirect fallbacks like ``gh auth token``), this runs the same
|
||||
resolve_runtime_provider() the agent uses on session creation and returns
|
||||
ok=False with the auth error when the model cannot actually be served, so
|
||||
UIs can surface onboarding before a doomed prompt. ``profile`` (optional)
|
||||
answers for THAT profile's config.yaml pin and ``.env``; an unknown profile
|
||||
answers ``ok=False`` rather than the launch profile's readiness.
|
||||
"""
|
||||
Unlike setup.status (True if ANY provider auth state is discoverable), this runs the same
|
||||
resolve_runtime_provider() the agent uses on session creation and returns ok=False with the
|
||||
auth error when the model can't be served, so UIs surface onboarding before a doomed prompt.
|
||||
``profile`` answers for THAT profile's pin and ``.env``; unknown -> ``ok=False``."""
|
||||
try:
|
||||
from hermes_cli.runtime_provider import resolve_runtime_provider
|
||||
from hermes_cli.auth import has_usable_secret
|
||||
@@ -397,8 +349,7 @@ def _(rid, params: dict) -> dict:
|
||||
def probe(profile):
|
||||
runtime = resolve_runtime_provider(requested=requested)
|
||||
provider_configured = bool(
|
||||
_has_any_provider_configured(strict_profile_scope=bool(profile))
|
||||
)
|
||||
_has_any_provider_configured(strict_profile_scope=bool(profile)))
|
||||
scoped = {"profile": profile} if profile else {}
|
||||
provider = runtime.get("provider") or "provider"
|
||||
source = str(runtime.get("source") or "")
|
||||
@@ -406,7 +357,6 @@ def _(rid, params: dict) -> dict:
|
||||
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)
|
||||
@@ -414,13 +364,11 @@ def _(rid, params: dict) -> dict:
|
||||
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"))
|
||||
)
|
||||
or has_usable_secret(api_key_text) or bool(runtime.get("command")))
|
||||
if not credential_ok:
|
||||
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:
|
||||
return _ok(rid, {"ok": False, "error": str(e)})
|
||||
@@ -428,66 +376,51 @@ def _(rid, params: dict) -> dict:
|
||||
|
||||
@method("diagnostics.share_nous")
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""Upload a redacted debug bundle to Nous-internal diagnostics storage.
|
||||
|
||||
Same collection + force-redaction pipeline as ``hermes debug share --nous``;
|
||||
redaction is NOT client-controllable. Consent lives with the CALLER (the
|
||||
desktop shows the privacy notice + Upload button first). Structured
|
||||
``ok``/``error`` envelope rather than JSON-RPC errors so the client can
|
||||
render upload failures inline.
|
||||
|
||||
Params (optional): ``error_context`` (client text about the failure,
|
||||
redacted, attached as ``error-context.txt``), ``extra_files`` ({label →
|
||||
text} client-side artifacts such as a remote desktop.log; force-redacted,
|
||||
labels sanitized and size-capped), ``log_lines`` (default 200).
|
||||
"""
|
||||
"""Upload a redacted debug bundle to Nous-internal diagnostics storage — same collection +
|
||||
force-redaction pipeline as ``hermes debug share --nous``; redaction is NOT
|
||||
client-controllable and consent lives with the CALLER (privacy notice first). Structured
|
||||
``ok``/``error`` envelope so the client renders upload failures inline. Optional params:
|
||||
``error_context`` (redacted, attached as ``error-context.txt``), ``extra_files`` ({label ->
|
||||
text}, force-redacted, labels sanitized and size-capped), ``log_lines`` (default 200)."""
|
||||
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
|
||||
|
||||
bundle = collect_share_bundle(log_lines=log_lines, redact=True)
|
||||
|
||||
# Client text goes through the SAME upload-safe redactor as backend
|
||||
# logs (force secret redaction + email masking), never the weaker bare
|
||||
# secret pass.
|
||||
# Client text goes through the SAME upload-safe redactor as backend logs (force secret
|
||||
# redaction + email masking), never the weaker bare secret pass.
|
||||
error_context = params.get("error_context")
|
||||
if isinstance(error_context, str) and error_context.strip():
|
||||
bundle["error-context.txt"] = _redact_log_text(error_context.strip()[:8_000])
|
||||
|
||||
# Bounded: at most 4 files, 512KB each, sanitized labels — a
|
||||
# diagnostics channel, not an arbitrary upload surface.
|
||||
# Bounded: at most 4 files, 512KB each, sanitized labels — not an arbitrary upload surface.
|
||||
extra_files = params.get("extra_files")
|
||||
if isinstance(extra_files, dict):
|
||||
for label, text in list(extra_files.items())[:4]:
|
||||
if not isinstance(label, str) or not isinstance(text, str):
|
||||
continue
|
||||
safe_label = "".join(ch for ch in label if ch.isalnum() or ch in "._- ()").strip()[:64]
|
||||
# Collapse dot-runs / leading dots so traversal-shaped labels
|
||||
# can't survive even cosmetically.
|
||||
# Collapse dot-runs / leading dots so traversal-shaped labels can't survive.
|
||||
while ".." in safe_label:
|
||||
safe_label = safe_label.replace("..", ".")
|
||||
safe_label = safe_label.lstrip(".").strip()
|
||||
if not safe_label or not text.strip():
|
||||
continue
|
||||
bundle[f"client/{safe_label}"] = _redact_log_text(text[:524_288])
|
||||
|
||||
res = share_to_nous(build_nous_bundle(bundle, redact=True))
|
||||
view_url = res.get("viewUrl") or res.get("view_url")
|
||||
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": 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"),
|
||||
})
|
||||
"expires_at": res.get("expiresAt") or res.get("expires_at")})
|
||||
except Exception as e:
|
||||
return _ok(rid, {"ok": False, "error": str(e)})
|
||||
|
||||
|
||||
def register(server) -> None:
|
||||
"""Publish helpers + handlers onto ``server``, rebound to its globals."""
|
||||
bind_module(globals(), server, skip=("_",))
|
||||
|
||||
@@ -1,8 +1,7 @@
|
||||
"""``config.set`` — one JSON-RPC method, dispatched on ``key`` through a table.
|
||||
|
||||
Bodies are rebound onto server.py's globals (method_ctx.bind_module) and reference them bare.
|
||||
Each ``_set_*`` handler takes ``(rid, params, key, value, session)`` and returns the JSON-RPC
|
||||
envelope. Keys match exactly except ``details_mode.<section>`` (prefix) and ``_DISPLAY_TOGGLE_KEYS``.
|
||||
"""``config.set`` — one JSON-RPC method, dispatched on ``key`` through a table. Bodies are
|
||||
rebound onto server.py's globals (method_ctx.bind_module) and reference them bare. Each
|
||||
``_set_*`` takes ``(rid, params, key, value, session)`` and returns the JSON-RPC envelope.
|
||||
Keys match exactly except ``details_mode.<section>`` (prefix) and ``_DISPLAY_TOGGLE_KEYS``.
|
||||
"""
|
||||
|
||||
import os
|
||||
@@ -16,18 +15,13 @@ method = _registry.method
|
||||
_profile_scoped = _registry.profile_scoped
|
||||
|
||||
|
||||
# ── shared helpers ────────────────────────────────────────────────────
|
||||
|
||||
|
||||
def _display_section(cfg: dict) -> dict:
|
||||
display = cfg.get("display")
|
||||
return display if isinstance(display, dict) else {}
|
||||
|
||||
# ── shared helpers
|
||||
|
||||
def _write_display_sections(*, sections=None, drop_sections=(), **display_fields) -> None:
|
||||
"""Persist ``display.<field>`` + ``display.sections`` edits via the raw (uncached) config write-back."""
|
||||
"""Persist ``display.<field>`` + ``display.sections`` edits via the raw (uncached) write-back."""
|
||||
cfg = _load_cfg_raw()
|
||||
display = _display_section(cfg)
|
||||
display = cfg.get("display")
|
||||
display = display if isinstance(display, dict) else {}
|
||||
cur = display.get("sections")
|
||||
cur = cur if isinstance(cur, dict) else {}
|
||||
display.update(display_fields)
|
||||
@@ -81,8 +75,7 @@ def _cfgset_model_ok(rid, key, value, warning, confirm_required, confirm_message
|
||||
"confirm_message": confirm_message, "scope": scope, **extra})
|
||||
|
||||
|
||||
# ── per-key handlers ──────────────────────────────────────────────────
|
||||
|
||||
# ── per-key handlers
|
||||
|
||||
def _set_model(rid, params, key, value, session):
|
||||
"""Live/deferred model switch; see _apply_model_switch and _apply_pending_model_switch."""
|
||||
@@ -93,9 +86,8 @@ def _set_model(rid, params, key, value, session):
|
||||
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
|
||||
# apply it at the NEXT turn start (_apply_pending_model_switch).
|
||||
# No live swap while a turn streams (agent.switch_model() mutates fields the worker
|
||||
# thread reads every iteration): stash the pick for the NEXT turn start.
|
||||
if session.get("running"):
|
||||
parsed = parse_model_switch_args(value)
|
||||
try:
|
||||
@@ -104,12 +96,12 @@ def _set_model(rid, params, key, value, session):
|
||||
pending_model = str(value)
|
||||
pending_provider = (getattr(parsed, "explicit_provider", "") or "").strip()
|
||||
# Selection guards run HERE (the only moment a confirm round-trip is possible);
|
||||
# otherwise an unconfirmed stashed pick is dropped at turn start, never confirmed.
|
||||
# otherwise an unconfirmed stashed pick is dropped at turn start.
|
||||
if not confirmed:
|
||||
pending_warning = _pending_switch_selection_warning(pending_model, pending_provider)
|
||||
if pending_warning is not None:
|
||||
# Nothing stashed; the client re-sends with confirm_expensive_model.
|
||||
# `confirm_message` is canonical, `warning` its legacy alias — identical.
|
||||
# `confirm_message` is canonical, `warning` its legacy alias.
|
||||
return _cfgset_model_ok(
|
||||
rid, key, pending_model, pending_warning, True, pending_warning, "session", deferred=False
|
||||
)
|
||||
@@ -117,10 +109,9 @@ def _set_model(rid, params, key, value, session):
|
||||
"raw": value,
|
||||
"confirm_expensive_model": confirmed,
|
||||
# _session_info reports these while pending so the end-of-turn settle keeps
|
||||
# showing the user's pick instead of the still-live old model.
|
||||
# showing the user's pick, not the still-live old model.
|
||||
"display_model": pending_model,
|
||||
"display_provider": pending_provider,
|
||||
}
|
||||
"display_provider": pending_provider}
|
||||
return _cfgset_model_ok(rid, key, pending_model, "", False, "", "session", deferred=True)
|
||||
parsed_flags = parse_model_switch_args(value)
|
||||
explicit_provider = parsed_flags.explicit_provider
|
||||
@@ -133,8 +124,7 @@ def _set_model(rid, params, key, value, session):
|
||||
return _err(rid, 5032, "agent initialization timed out")
|
||||
failed_agent_init = (
|
||||
failed_agent_init and session.get("agent") is None and session.get("agent_error") is not None
|
||||
and session.get("agent_ready") is failed_ready and failed_ready.is_set()
|
||||
)
|
||||
and session.get("agent_ready") is failed_ready and failed_ready.is_set())
|
||||
if session.get("agent") is None and not explicit_provider.strip() and not failed_agent_init:
|
||||
_start_agent_build(sid, session)
|
||||
if init_err := _cfgset_await_agent(session, rid):
|
||||
@@ -153,13 +143,13 @@ def _set_model(rid, params, key, value, session):
|
||||
result = _apply_model_switch("", {"agent": None}, value, confirm_expensive_model=confirmed)
|
||||
return _cfgset_model_ok(
|
||||
rid, key, result["value"], result["warning"], result.get("confirm_required", False),
|
||||
result.get("confirm_message", ""), result.get("scope", "session"),
|
||||
)
|
||||
result.get("confirm_message", ""), result.get("scope", "session"))
|
||||
except Exception as e:
|
||||
return _err(rid, 5001, str(e))
|
||||
|
||||
|
||||
_FAST_WORDS = {"fast": "fast", "on": "fast", "normal": "normal", "off": "normal", "auto": "auto", "cold": "cold"}
|
||||
_FAST_WORDS = {"fast": "fast", "on": "fast", "normal": "normal", "off": "normal",
|
||||
"auto": "auto", "cold": "cold"}
|
||||
|
||||
|
||||
def _set_fast(rid, params, key, value, session):
|
||||
@@ -173,13 +163,12 @@ def _set_fast(rid, params, key, value, session):
|
||||
else:
|
||||
current_tier = _load_service_tier()
|
||||
current_fast = current_tier == "priority"
|
||||
|
||||
if raw == "status":
|
||||
return _ok(rid, {"key": key, "value": {"priority": "fast", None: "normal"}.get(current_tier, current_tier)})
|
||||
nv = _FAST_WORDS.get(raw, ("normal" if current_fast else "fast") if raw in {"", "toggle"} else None)
|
||||
toggled = ("normal" if current_fast else "fast") if raw in {"", "toggle"} else None
|
||||
nv = _FAST_WORDS.get(raw, toggled)
|
||||
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
|
||||
@@ -191,15 +180,14 @@ def _set_fast(rid, params, key, value, session):
|
||||
target_model = (isinstance(session_override, dict) and session_override.get("model")) or _resolve_model()
|
||||
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))
|
||||
overrides = resolve_fast_mode_overrides(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")
|
||||
|
||||
if session is not None:
|
||||
# Session-scoped like `reasoning` (global persistence is `--global` / Settings → Model):
|
||||
# writing config.yaml here flipped fast mode for every other session/profile/CLI/gateway.
|
||||
# The create override keeps the choice across lazy builds and rebuilds; "" pins normal.
|
||||
# writing config.yaml here flipped fast mode for every other surface. The create
|
||||
# override keeps the choice across lazy builds and rebuilds; "" pins normal.
|
||||
session["create_service_tier_override"] = {"fast": "priority", "normal": ""}.get(nv, nv)
|
||||
else:
|
||||
_write_config_key("agent.service_tier", nv)
|
||||
@@ -245,8 +233,8 @@ def _set_verbose(rid, params, key, value, session):
|
||||
|
||||
|
||||
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.
|
||||
# Focus view (/focus): enabling stashes the configured tool_progress mode and pins it
|
||||
# "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_cfg()
|
||||
cur_focus = bool(d_f.get("focus_view", False))
|
||||
@@ -255,7 +243,6 @@ def _set_focus(rid, params, key, value, session):
|
||||
return _err(rid, 4002, f"unknown focus value: {value} (use on|off|status)")
|
||||
if action == "status" or target is None:
|
||||
return _ok(rid, {"key": key, "value": "on" if cur_focus else "off", "tool_progress": _load_tool_progress_mode()})
|
||||
|
||||
if target:
|
||||
saved = (cur_focus and d_f.get("focus_saved_tool_progress")) or _load_tool_progress_mode()
|
||||
_write_config_key("display.focus_saved_tool_progress", normalize_tool_progress_mode(saved))
|
||||
@@ -265,7 +252,6 @@ def _set_focus(rid, params, key, value, session):
|
||||
effective = normalize_tool_progress_mode(d_f.get("focus_saved_tool_progress") or "all")
|
||||
_write_config_key("display.tool_progress", effective)
|
||||
_write_config_key("display.focus_view", bool(target))
|
||||
|
||||
if session:
|
||||
session["focus_view"] = bool(target)
|
||||
session["tool_progress_mode"] = effective
|
||||
@@ -285,9 +271,8 @@ def _set_approval_mode(rid, params, key, value, session):
|
||||
|
||||
|
||||
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="session" (default; Shift+Tab) toggles ONLY this session's flag. scope="global"
|
||||
# (Shift+click the zap) flips persistent approvals.mode between "off" and "manual".
|
||||
scope = _word(params.get("scope") or "session")
|
||||
try:
|
||||
from tools.approval import disable_session_yolo, enable_session_yolo, is_session_yolo_enabled
|
||||
@@ -295,7 +280,6 @@ def _set_yolo(rid, params, key, value, session):
|
||||
|
||||
def _resolve_toggle(current: bool) -> bool:
|
||||
return _BOOL_WORDS.get(raw, not current)
|
||||
|
||||
if scope == "global":
|
||||
from tools.approval import _normalize_approval_mode
|
||||
appr = _load_cfg().get("approvals")
|
||||
@@ -305,7 +289,6 @@ def _set_yolo(rid, params, key, value, session):
|
||||
_write_config_key("approvals.mode", "off" if enable else "manual")
|
||||
_emit_all_session_info() # reflect the flip in every live indicator
|
||||
return _ok(rid, {"key": key, "value": "1" if enable else "0", "scope": "global"})
|
||||
|
||||
if session:
|
||||
skey = session["session_key"]
|
||||
enable = _resolve_toggle(is_session_yolo_enabled(skey))
|
||||
@@ -323,14 +306,12 @@ def _set_yolo(rid, params, key, value, session):
|
||||
|
||||
|
||||
# /reasoning display words: (accepted inputs, reported value, display field, sections.thinking,
|
||||
# session show_reasoning or None). full/clamp mirror the CLI's reasoning_full toggle; the TUI
|
||||
# renders thinking as an expand/collapse section and display.reasoning_full is persisted too.
|
||||
# session show_reasoning or None). full/clamp mirror the CLI's reasoning_full toggle.
|
||||
_REASONING_DISPLAY_WORDS = (
|
||||
({"show", "on"}, "show", {"show_reasoning": True}, "expanded", True),
|
||||
({"hide", "off"}, "hide", {"show_reasoning": False}, "hidden", False),
|
||||
({"full", "all"}, "full", {"reasoning_full": True}, "expanded", None),
|
||||
({"clamp", "collapse", "short"}, "clamp", {"reasoning_full": False}, "collapsed", None),
|
||||
)
|
||||
({"clamp", "collapse", "short"}, "clamp", {"reasoning_full": False}, "collapsed", None))
|
||||
|
||||
|
||||
def _set_reasoning(rid, params, key, value, session):
|
||||
@@ -344,7 +325,6 @@ def _set_reasoning(rid, params, key, value, session):
|
||||
if show is not None and session:
|
||||
session["show_reasoning"] = show
|
||||
return _ok(rid, {"key": key, "value": reported})
|
||||
|
||||
parsed = parse_reasoning_effort(arg)
|
||||
if parsed is None:
|
||||
return _err(rid, 4002, f"unknown reasoning value: {value}")
|
||||
@@ -353,8 +333,8 @@ def _set_reasoning(rid, params, key, value, session):
|
||||
if session is not None:
|
||||
session.pop("create_reasoning_override", None)
|
||||
else:
|
||||
# Session-scoped like the messaging gateway's `/reasoning <level>`; otherwise every
|
||||
# desktop model-menu pick rewrote the global default.
|
||||
# Session-scoped like the gateway's `/reasoning <level>`; otherwise every desktop
|
||||
# model-menu pick rewrote the global default.
|
||||
session["create_reasoning_override"] = parsed
|
||||
if session and session.get("agent") is not None:
|
||||
session["agent"].reasoning_config = parsed
|
||||
@@ -374,8 +354,8 @@ def _set_details_mode(rid, params, key, value, session):
|
||||
|
||||
|
||||
def _set_details_section(rid, params, key, value, session):
|
||||
# `details_mode.<section>` -> `display.sections.<section>`. Empty value clears the explicit
|
||||
# override so the frontend applies built-in section defaults before the global details_mode.
|
||||
# `details_mode.<section>` -> `display.sections.<section>`; empty clears the override so the
|
||||
# frontend applies built-in section defaults before the global details_mode.
|
||||
section = key.split(".", 1)[1]
|
||||
if section not in _DETAIL_SECTION_NAMES:
|
||||
return _err(rid, 4002, f"unknown section: {section}")
|
||||
@@ -432,8 +412,8 @@ def _set_statusbar(rid, params, key, value, session):
|
||||
|
||||
|
||||
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.
|
||||
# Explicit None check so falsy non-string inputs (0, False) 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_cfg())
|
||||
if raw in {"", "toggle"}:
|
||||
@@ -480,8 +460,8 @@ def _set_prompt_like(rid, params, key, value, session):
|
||||
_save_cfg(cfg)
|
||||
elif key == "personality":
|
||||
pname, new_prompt = _validate_personality(str(value or ""), cfg)
|
||||
# Personality text is an in-session overlay; persistence goes through
|
||||
# hermes_cli.personality (single owner), never the user-owned global system prompt.
|
||||
# Personality persists 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")
|
||||
@@ -492,8 +472,8 @@ def _set_prompt_like(rid, params, key, value, session):
|
||||
else:
|
||||
_write_config_key(f"display.{key}", value)
|
||||
if key == "skin":
|
||||
# Every connected surface repaints; then sync the watcher baseline so the poll
|
||||
# loop doesn't re-broadcast the skin this RPC just applied.
|
||||
# Every surface repaints; sync the watcher baseline so the poll loop doesn't
|
||||
# re-broadcast the skin this RPC just applied.
|
||||
_broadcast_global_event("skin.changed", resolve_skin())
|
||||
_note_skin_broadcast()
|
||||
return _ok(rid, resp)
|
||||
@@ -509,7 +489,7 @@ def _set_display_toggle(rid, params, key, value, session):
|
||||
return _ok(rid, {"key": key, "value": on})
|
||||
|
||||
|
||||
# ── dispatch ──────────────────────────────────────────────────────────
|
||||
# ── dispatch
|
||||
|
||||
_CONFIG_SETTERS = {
|
||||
"model": _set_model, "fast": _set_fast, "busy": _set_busy, "verbose": _set_verbose, "focus": _set_focus,
|
||||
@@ -518,19 +498,7 @@ _CONFIG_SETTERS = {
|
||||
"density": _set_density, "battery": _set_battery, "theme": _set_theme, "statusbar": _set_statusbar,
|
||||
"mouse": _set_mouse, "indicator": _set_indicator,
|
||||
"cwd": _set_cwd, "terminal.cwd": _set_cwd, "workdir": _set_cwd,
|
||||
"prompt": _set_prompt_like, "personality": _set_prompt_like, "skin": _set_prompt_like,
|
||||
}
|
||||
|
||||
|
||||
def _config_setter(key: str):
|
||||
handler = _CONFIG_SETTERS.get(key)
|
||||
if handler is not None:
|
||||
return handler
|
||||
if key.startswith("details_mode."):
|
||||
return _set_details_section
|
||||
if key in _DISPLAY_TOGGLE_KEYS:
|
||||
return _set_display_toggle
|
||||
return None
|
||||
"prompt": _set_prompt_like, "personality": _set_prompt_like, "skin": _set_prompt_like}
|
||||
|
||||
|
||||
@method("config.set")
|
||||
@@ -538,12 +506,15 @@ def _config_setter(key: str):
|
||||
def _(rid, params: dict) -> dict:
|
||||
key, value = params.get("key", ""), params.get("value", "")
|
||||
session = _sessions.get(params.get("session_id", ""))
|
||||
handler = _config_setter(key)
|
||||
handler = _CONFIG_SETTERS.get(key)
|
||||
if handler is None and key.startswith("details_mode."):
|
||||
handler = _set_details_section
|
||||
elif handler is None and key in _DISPLAY_TOGGLE_KEYS:
|
||||
handler = _set_display_toggle
|
||||
if handler is None:
|
||||
return _err(rid, 4002, f"unknown config key: {key}")
|
||||
return handler(rid, params, key, value, session)
|
||||
|
||||
|
||||
def register(server) -> None:
|
||||
"""Publish helpers + the config.set handler onto ``server``, rebound to its globals."""
|
||||
bind_module(globals(), server, skip=("_",))
|
||||
|
||||
@@ -1,10 +1,8 @@
|
||||
"""Image-generation JSON-RPC handler (ws twin of the image_generate tool), so UI
|
||||
surfaces (avatar pickers, artifact panes) can generate directly. The result is
|
||||
returned as a data URL: a remote desktop can't read a gateway file path and hosted
|
||||
URLs are often CORS-opaque to a renderer canvas.
|
||||
"""Image-generation JSON-RPC handler (ws twin of the image_generate tool) for UI surfaces
|
||||
(avatar pickers, artifact panes). The result is a data URL: a remote desktop can't read a
|
||||
gateway file path and hosted URLs are often CORS-opaque to a renderer canvas.
|
||||
|
||||
Bodies are rebound onto server.py's globals at install time (see
|
||||
method_ctx.bind_module), so they reference server.py globals bare.
|
||||
Bodies are rebound onto server.py's globals (method_ctx.bind_module) and reference them bare.
|
||||
"""
|
||||
|
||||
from .method_ctx import HandlerRegistry, bind_module
|
||||
@@ -16,7 +14,6 @@ method = _registry.method
|
||||
def _image_gen_available() -> bool:
|
||||
try:
|
||||
from tools.image_generation_tool import check_image_generation_requirements
|
||||
|
||||
return bool(check_image_generation_requirements())
|
||||
except Exception:
|
||||
return False
|
||||
@@ -27,11 +24,9 @@ def _image_to_data_url(ref: str, cap: int):
|
||||
import base64
|
||||
import mimetypes
|
||||
import os
|
||||
|
||||
try:
|
||||
if ref.startswith(("http://", "https://")):
|
||||
import urllib.request
|
||||
|
||||
req = urllib.request.Request(ref, headers={"User-Agent": "hermes-agent"})
|
||||
with urllib.request.urlopen(req, timeout=60) as resp:
|
||||
if resp.length is not None and resp.length > cap:
|
||||
@@ -66,46 +61,28 @@ def _(rid, params: dict) -> dict:
|
||||
if is_truthy_value(params.get("probe", False)):
|
||||
return _ok(rid, {"available": available})
|
||||
if not available:
|
||||
return _ok(
|
||||
rid,
|
||||
{
|
||||
"available": False,
|
||||
"success": False,
|
||||
"error": "No image generation backend configured (run `hermes tools` to enable one).",
|
||||
},
|
||||
)
|
||||
|
||||
return _ok(rid, {
|
||||
"available": False, "success": False,
|
||||
"error": "No image generation backend configured (run `hermes tools` to enable one).",
|
||||
})
|
||||
prompt = str(params.get("prompt") or "").strip()
|
||||
if not prompt:
|
||||
return _err(rid, 4071, "prompt required")
|
||||
|
||||
aspect = str(params.get("aspect_ratio") or "square").strip().lower()
|
||||
try:
|
||||
cap = min(int(params.get("max_bytes", 8_000_000) or 8_000_000), 16_000_000)
|
||||
except (TypeError, ValueError):
|
||||
cap = 8_000_000
|
||||
|
||||
try:
|
||||
from tools.image_generation_tool import _handle_image_generate
|
||||
|
||||
# Full provider dispatcher — same path as the model tool (source-image
|
||||
# confinement, plugin providers, managed routing, FAL fallback); calling
|
||||
# the FAL leaf directly bypassed configured providers.
|
||||
raw = _handle_image_generate({"prompt": prompt, "aspect_ratio": aspect})
|
||||
result = json.loads(raw)
|
||||
# Full provider dispatcher — same path as the model tool (source-image confinement,
|
||||
# plugin providers, managed routing, FAL fallback); the FAL leaf bypassed providers.
|
||||
result = json.loads(_handle_image_generate({"prompt": prompt, "aspect_ratio": aspect}))
|
||||
except Exception as e:
|
||||
return _err(rid, 5071, str(e))
|
||||
|
||||
if not result.get("success"):
|
||||
return _ok(
|
||||
rid,
|
||||
{
|
||||
"available": True,
|
||||
"success": False,
|
||||
"error": str(result.get("error") or "generation failed"),
|
||||
},
|
||||
)
|
||||
|
||||
return _ok(rid, {"available": True, "success": False,
|
||||
"error": str(result.get("error") or "generation failed")})
|
||||
image_ref = str(result.get("image") or "")
|
||||
payload = {"available": True, "success": True, "image": image_ref}
|
||||
data_url = _image_to_data_url(image_ref, cap) if image_ref else None
|
||||
@@ -115,5 +92,4 @@ def _(rid, params: dict) -> dict:
|
||||
|
||||
|
||||
def register(server) -> None:
|
||||
"""Publish this module's helpers + handlers onto ``server``, rebound to its globals."""
|
||||
bind_module(globals(), server, skip=("_",))
|
||||
|
||||
@@ -1,9 +1,8 @@
|
||||
"""Profile JSON-RPC handlers — the ws twin of the dashboard's /api/profiles (desktop plugins
|
||||
only have the ws door), on the same `hermes_cli.profiles` primitives.
|
||||
|
||||
Bodies are rebound onto server.py's globals (method_ctx.bind_module) and use them bare
|
||||
(`_ok`, `_err`, `os`, `json`, `Path`, `is_truthy_value`, `get_hermes_home`, ...); module-level
|
||||
names are published onto server.py, so they must not collide with its globals.
|
||||
Bodies are rebound onto server.py's globals (method_ctx.bind_module) and use them bare;
|
||||
module-level names are published onto server.py, so they must not collide with its globals.
|
||||
"""
|
||||
|
||||
import contextlib
|
||||
@@ -31,8 +30,8 @@ def _profile_handler(name: str, code: int):
|
||||
|
||||
|
||||
def _lazy(module, name):
|
||||
"""Late-bound attribute lookup (heavy / cyclic modules). ``__import__`` builtin on purpose:
|
||||
rebound bodies see only server.py globals, not this module's imports."""
|
||||
"""Late-bound attribute lookup; ``__import__`` builtin on purpose (rebound bodies see only
|
||||
server.py globals, not this module's imports)."""
|
||||
return getattr(__import__(module, fromlist=[name]), name)
|
||||
|
||||
|
||||
@@ -58,10 +57,6 @@ def _best_effort(fn) -> bool:
|
||||
return _try(lambda: (fn(), True)[1], False)
|
||||
|
||||
|
||||
def _read_text_if_file(path) -> str:
|
||||
return _try(lambda: path.read_text(encoding="utf-8", errors="replace") if path.is_file() else "", "")
|
||||
|
||||
|
||||
@contextlib.contextmanager
|
||||
def _hermes_home_scope(path):
|
||||
"""Scope config/auth resolution to ``path`` for the block."""
|
||||
@@ -98,9 +93,8 @@ def _clean_revisions(raw: dict) -> dict:
|
||||
|
||||
|
||||
def _latest_message_preview(db, session_id):
|
||||
"""≤80-char excerpt of the NEWEST active user/assistant message, or "" (roster semantics:
|
||||
latest exchange, not the first-message preview). Same query shape as
|
||||
``SessionDB.latest_message_row_id`` — keep them in step."""
|
||||
"""≤80-char excerpt of the NEWEST active user/assistant message, or "" (roster semantics).
|
||||
Same query shape as ``SessionDB.latest_message_row_id`` — keep them in step."""
|
||||
try:
|
||||
with db._lock:
|
||||
row = db._conn.execute(
|
||||
@@ -108,8 +102,7 @@ def _latest_message_preview(db, session_id):
|
||||
" WHERE session_id = ? AND role IN ('user', 'assistant')"
|
||||
" AND active = 1 AND content IS NOT NULL AND TRIM(content) != ''"
|
||||
" ORDER BY id DESC LIMIT 1",
|
||||
(session_id,),
|
||||
).fetchone()
|
||||
(session_id,)).fetchone()
|
||||
except Exception:
|
||||
return ""
|
||||
if not row:
|
||||
@@ -119,8 +112,8 @@ def _latest_message_preview(db, session_id):
|
||||
|
||||
|
||||
def _open_profile_session_db_readonly(profile_path):
|
||||
"""Read-only attach for roster previews, or None. A writable ``SessionDB()`` waits up to
|
||||
20s for the write lock + runs DDL; the 5s roster poll stalled past the desktop timeout."""
|
||||
"""Read-only attach for roster previews, or None (a writable ``SessionDB()`` waits up to 20s
|
||||
for the write lock + runs DDL and stalled the 5s roster poll)."""
|
||||
db_path = Path(profile_path) / "state.db"
|
||||
if not _try(db_path.exists, False):
|
||||
return None
|
||||
@@ -128,8 +121,8 @@ def _open_profile_session_db_readonly(profile_path):
|
||||
|
||||
|
||||
def _resurrect_recoverable_canonical(db, profile_path, session_id):
|
||||
"""Un-archive an accidentally archived canonical row, or False. Recoverability is judged on
|
||||
the read-only handle; the write uses a short-lived writable handle."""
|
||||
"""Un-archive an accidentally archived canonical row (judged read-only, written via a
|
||||
short-lived writable handle), or False."""
|
||||
try:
|
||||
row = db.get_session(session_id)
|
||||
if not row or not row.get("archived"):
|
||||
@@ -149,10 +142,9 @@ def _resurrect_recoverable_canonical(db, profile_path, session_id):
|
||||
|
||||
|
||||
def _canonical_session_row(db, profile_path):
|
||||
"""Summary of the profile's canonical "Bot Chat" row, or None. Identity is the NAME, so
|
||||
preview and click target agree without a client pointer. Hidden rows resolve; lineages
|
||||
via ``get_compression_tip`` (NOT the resume walker's unmarked-child fallback); worker
|
||||
sources count as absent. ``id`` is the registry row, ``resolved_id`` the live tip."""
|
||||
"""Summary of the profile's canonical "Bot Chat" row (identity is the NAME), or None.
|
||||
Lineages via ``get_compression_tip`` (NOT the resume walker's unmarked-child fallback);
|
||||
worker sources count as absent. ``id`` is the registry row, ``resolved_id`` the live tip."""
|
||||
if db is None:
|
||||
return None
|
||||
try:
|
||||
@@ -173,8 +165,7 @@ def _canonical_session_row(db, profile_path):
|
||||
"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,
|
||||
}
|
||||
"message_count": tip_row.get("message_count") or 0}
|
||||
except Exception:
|
||||
return None
|
||||
|
||||
@@ -201,8 +192,7 @@ def _latest_profile_session_rows(db):
|
||||
"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,
|
||||
}
|
||||
"message_count": s.get("message_count") or 0}
|
||||
if worker is not None:
|
||||
break
|
||||
return human, worker
|
||||
@@ -224,11 +214,8 @@ def _profile_session_fields(row, profile_path):
|
||||
|
||||
def _profile_ui_meta_fields(row: dict, profile_dir) -> None:
|
||||
"""Attach ``ui_meta`` / ``ui_meta_revisions`` / ``has_avatar`` from profile.yaml + assets.
|
||||
|
||||
Client-agnostic UI metadata lives in profile.yaml so every client paints the
|
||||
same roster. ``ui_meta_revisions`` is always present: it feature-detects
|
||||
gateway-owned CAS even for a brand-new profile.
|
||||
"""
|
||||
``ui_meta_revisions`` is always present: it feature-detects gateway-owned CAS even for a
|
||||
brand-new profile."""
|
||||
row["ui_meta_revisions"] = {}
|
||||
raw_meta = _try(lambda: _read_profile_yaml(profile_dir), {})
|
||||
ui_meta = raw_meta.get("ui_meta")
|
||||
@@ -243,11 +230,8 @@ def _profile_ui_meta_fields(row: dict, profile_dir) -> None:
|
||||
|
||||
@_profile_handler("profiles.list", 5061)
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""List Hermes profiles (name, path, model, description, skill count).
|
||||
|
||||
``include_sessions`` (default true) adds ``last_session`` / ``worker_session``
|
||||
/ ``canonical_session`` so a roster paints per-agent previews without N calls.
|
||||
"""
|
||||
"""List Hermes profiles. ``include_sessions`` (default true) adds ``last_session`` /
|
||||
``worker_session`` / ``canonical_session`` so a roster paints previews without N calls."""
|
||||
from hermes_cli.profiles import list_profiles
|
||||
include_sessions = is_truthy_value(params.get("include_sessions", True))
|
||||
out = []
|
||||
@@ -256,8 +240,7 @@ def _(rid, params: dict) -> dict:
|
||||
"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,
|
||||
}
|
||||
"skill_count": p.skill_count or 0}
|
||||
if include_sessions:
|
||||
_profile_session_fields(row, p.path)
|
||||
_profile_ui_meta_fields(row, Path(str(p.path)))
|
||||
@@ -270,7 +253,7 @@ def _(rid, params: dict) -> dict:
|
||||
def _has_real_env_content(env_path) -> bool:
|
||||
"""True when .env has any non-comment, non-blank line."""
|
||||
lines = env_path.read_text(encoding="utf-8", errors="replace").splitlines()
|
||||
return any(s and not s.startswith("#") for s in (line.strip() for line in lines))
|
||||
return any(s and not s.startswith("#") for s in map(str.strip, lines))
|
||||
|
||||
|
||||
def _copy_secret_file(src, dst) -> None:
|
||||
@@ -312,8 +295,7 @@ def _mirror_voice_sections(path) -> bool:
|
||||
if not sections:
|
||||
return False
|
||||
with _hermes_home_scope(path):
|
||||
# RAW file: load_config() merges DEFAULT_CONFIG (every section looks present
|
||||
# and save_config would persist the whole default tree).
|
||||
# RAW file: load_config() merges DEFAULT_CONFIG (every section would look present).
|
||||
dst_cfg = read_user_config_raw() or {}
|
||||
missing = {k: v for k, v in sections.items() if k not in dst_cfg}
|
||||
if missing:
|
||||
@@ -342,10 +324,8 @@ def _inherit_launch_model(path) -> bool:
|
||||
|
||||
def _mirror_launch_credentials(path, params: dict) -> dict:
|
||||
"""Copy launch .env / auth.json / voice sections into a new profile (best-effort per item).
|
||||
|
||||
``share_auth`` reports ``auth: "shared"`` and skips the auth copy; ``mirror_credentials``
|
||||
false skips everything. ``model_inherited`` is filled in by the caller.
|
||||
"""
|
||||
false skips everything. ``model_inherited`` is filled in by the caller."""
|
||||
mirrored = {"env": False, "auth": False, "model_inherited": False, "voice": False}
|
||||
share_auth = is_truthy_value(params.get("share_auth", False))
|
||||
if share_auth:
|
||||
@@ -362,13 +342,11 @@ def _mirror_launch_credentials(path, params: dict) -> dict:
|
||||
|
||||
@method("profiles.create")
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""Create a profile (ws twin of POST /api/profiles).
|
||||
|
||||
Params: ``name``, ``description``, ``clone_from`` (omitted = fresh + bundled skills),
|
||||
``clone_all``, ``no_skills``, ``soul``, ``model`` + ``provider``, ``share_auth``,
|
||||
``mirror_credentials`` (default true). Mirroring exists because ``create_profile()``
|
||||
seeds a comment-only .env and no auth.json — a headless profile had NO provider.
|
||||
"""
|
||||
"""Create a profile (ws twin of POST /api/profiles). Params: ``name``, ``description``,
|
||||
``clone_from`` (omitted = fresh + bundled skills), ``clone_all``, ``no_skills``, ``soul``,
|
||||
``model`` + ``provider``, ``share_auth``, ``mirror_credentials`` (default true — a
|
||||
``create_profile()`` seeds a comment-only .env and no auth.json, so a headless profile had
|
||||
NO provider)."""
|
||||
name = str(params.get("name") or "").strip()
|
||||
if not name:
|
||||
return _err(rid, 4061, "name required")
|
||||
@@ -380,13 +358,11 @@ def _(rid, params: dict) -> dict:
|
||||
name=name, clone_from=clone_from, clone_all=clone_all,
|
||||
clone_config=bool(clone_from) and not clone_all,
|
||||
no_skills=is_truthy_value(params.get("no_skills", False)),
|
||||
description=str(params.get("description") or "").strip() or None,
|
||||
)
|
||||
description=str(params.get("description") or "").strip() or None)
|
||||
except (ValueError, FileExistsError, FileNotFoundError) as e:
|
||||
return _err(rid, 4062, str(e))
|
||||
except Exception as e:
|
||||
return _err(rid, 5062, str(e))
|
||||
|
||||
# CLI/REST create flow: bundled skills for fresh profiles, then the alias wrapper.
|
||||
if not clone_from:
|
||||
_best_effort(lambda: profiles_mod.seed_profile_skills(path, quiet=True))
|
||||
@@ -410,7 +386,8 @@ def _(rid, params: dict) -> dict:
|
||||
def _describe_toolsets(cfg):
|
||||
"""``(toolsets, pinned_set)`` as the `hermes tools` checklist presents them (the raw registry
|
||||
leaks platform composites and reports everything "enabled" without a pin)."""
|
||||
from hermes_cli.tools_config import _get_effective_configurable_toolsets, _get_platform_tools, _toolset_allowed_for_platform
|
||||
from hermes_cli.tools_config import (
|
||||
_get_effective_configurable_toolsets, _get_platform_tools, _toolset_allowed_for_platform)
|
||||
from toolsets import resolve_toolset
|
||||
pinned = (cfg.get("tools") if isinstance(cfg.get("tools"), dict) else {}).get("enabled_toolsets")
|
||||
pinned_set = _clean_names(pinned) if isinstance(pinned, list) else None
|
||||
@@ -438,14 +415,14 @@ def _describe_mcp_servers(cfg):
|
||||
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)
|
||||
for srv_name in sorted(mcp_cfg.keys()) for entry in (mcp_cfg[srv_name],)
|
||||
if isinstance(entry, dict)
|
||||
], [])
|
||||
|
||||
|
||||
@_profile_handler("profiles.describe", 5063)
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""Editor snapshot: ``{name, description, soul, model, skills: [{name, enabled}], toolsets,
|
||||
toolsets_pinned, mcp_servers}``; installed skills are enabled unless in ``skills.disabled``."""
|
||||
"""Editor snapshot; installed skills are enabled unless in ``skills.disabled``."""
|
||||
name, profile_dir, err = _resolve_profile(rid, params)
|
||||
if err is not None:
|
||||
return err
|
||||
@@ -457,10 +434,10 @@ def _(rid, params: dict) -> dict:
|
||||
skills_root = profile_dir / "skills"
|
||||
installed = [
|
||||
{"name": md.parent.name, "enabled": md.parent.name.lower() not in disabled}
|
||||
for md in (sorted(skills_root.rglob("SKILL.md")) if skills_root.is_dir() else ())
|
||||
]
|
||||
for md in (sorted(skills_root.rglob("SKILL.md")) if skills_root.is_dir() else ())]
|
||||
toolsets_out, pinned_set = _describe_toolsets(cfg)
|
||||
soul = _read_text_if_file(profile_dir / "SOUL.md")
|
||||
soul_path = profile_dir / "SOUL.md"
|
||||
soul = _try(lambda: soul_path.read_text(encoding="utf-8", errors="replace") if soul_path.is_file() else "", "")
|
||||
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), {})
|
||||
@@ -468,9 +445,8 @@ def _(rid, params: dict) -> dict:
|
||||
"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,
|
||||
})
|
||||
"skills": installed, "toolsets": toolsets_out,
|
||||
"toolsets_pinned": pinned_set is not None, "mcp_servers": mcp_out})
|
||||
|
||||
|
||||
def _configure_ui_meta(profile_dir, params, applied) -> None:
|
||||
@@ -521,18 +497,16 @@ def _configure_ui_meta(profile_dir, params, applied) -> None:
|
||||
|
||||
|
||||
def _configure_model(profile_dir, params, applied):
|
||||
"""Apply a ``model`` + ``provider`` pin, or return a confirm message and write NOTHING (the
|
||||
``config.set model`` handshake: client resends with ``confirm_expensive_model``). A failing
|
||||
guard counts as "no warning", matching ``_apply_model_switch``."""
|
||||
"""Apply a ``model`` + ``provider`` pin, or return a confirm message and write NOTHING (client
|
||||
resends with ``confirm_expensive_model``). A failing guard = "no warning" (as _apply_model_switch)."""
|
||||
model = str(params.get("model") or "").strip()
|
||||
provider = str(params.get("provider") or "").strip()
|
||||
confirm_message = None
|
||||
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)
|
||||
warn = _lazy("hermes_cli.model_selection_guards", "combined_selection_warning")
|
||||
confirm_message = _try(lambda: getattr(warn(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
|
||||
@@ -540,8 +514,8 @@ def _configure_model(profile_dir, params, applied):
|
||||
|
||||
def _configure_cfg_sections(profile_dir, params, applied) -> None:
|
||||
"""Apply ``disabled_skills`` / ``enabled_toolsets`` / ``enabled_mcp_servers`` (replace
|
||||
semantics; empty toolsets clears the pin). Enabling an undefined MCP server copies its
|
||||
definition from the LAUNCH catalog (unknown names skipped); credentials stay in .env/auth."""
|
||||
semantics; empty toolsets clears the pin). An undefined MCP server is copied from the LAUNCH
|
||||
catalog (unknown names skipped); credentials stay in .env/auth."""
|
||||
want_mcp = isinstance(params.get("enabled_mcp_servers"), list)
|
||||
# Launch catalog read BEFORE the home override flips config resolution.
|
||||
launch_mcp = _try(_launch_mcp_catalog, {}) if want_mcp else {}
|
||||
@@ -609,7 +583,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")(
|
||||
write_meta = _lazy("hermes_cli.profiles", "write_profile_meta")
|
||||
applied["description"] = _best_effort(lambda: write_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")):
|
||||
@@ -642,7 +617,7 @@ def _unlink_asset_files(assets_dir, asset) -> int:
|
||||
@_profile_handler("profiles.set_asset", 5065)
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""Store ``assets/<asset>.<ext>`` atomically. Params: ``name``, ``asset`` (``"avatar"`` only),
|
||||
``data`` (data URL or base64; PNG/JPEG/WebP ≤2MB) or ``clear: true``. Result ``{ok, asset, size}``."""
|
||||
``data`` (data URL or base64; PNG/JPEG/WebP ≤2MB) or ``clear: true``."""
|
||||
asset = str(params.get("asset") or "avatar").strip().lower()
|
||||
if not str(params.get("name") or "").strip():
|
||||
return _err(rid, 4063, "name required")
|
||||
@@ -681,7 +656,7 @@ def _(rid, params: dict) -> dict:
|
||||
|
||||
@_profile_handler("profiles.get_asset", 5066)
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""Profile asset as a data URL: ``{found, data?, mime?, size?}``; absent is ``found: false``, not an error."""
|
||||
"""Profile asset as a data URL; absent is ``found: false``, not an error."""
|
||||
asset = str(params.get("asset") or "avatar").strip().lower()
|
||||
import base64
|
||||
_name, profile_dir, err = _resolve_profile(rid, params)
|
||||
|
||||
@@ -1,7 +1,6 @@
|
||||
"""Voice / TTS / wake-word JSON-RPC handlers and their process-global state (one microphone, one speaker per process).
|
||||
|
||||
Bodies are rebound onto server.py's globals at install time (see
|
||||
method_ctx.bind_module), so they reference server.py globals bare.
|
||||
"""Voice / TTS / wake-word JSON-RPC handlers and their process-global state (one mic, one
|
||||
speaker per process). Bodies are rebound onto server.py's globals (method_ctx.bind_module)
|
||||
and reference them bare.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -27,8 +26,7 @@ def _caller_transport():
|
||||
|
||||
|
||||
def _voice_emit(event: str, payload: dict | None = None) -> None:
|
||||
"""Emit toward the session that most recently turned voice on (one mic → one
|
||||
target sid; the TUI treats an empty sid as "active session")."""
|
||||
"""Emit toward the session that most recently turned voice on (empty sid = active session)."""
|
||||
with _voice_sid_lock:
|
||||
sid = _voice_event_sid
|
||||
_emit(event, sid, payload)
|
||||
@@ -43,8 +41,7 @@ def _resume_voice_wake() -> None:
|
||||
|
||||
|
||||
def _voice_mode_enabled() -> bool:
|
||||
"""Runtime-only flag (CLI parity): env var only, never config.yaml, so the TUI
|
||||
can't auto-start in REC because voice was on in a prior session."""
|
||||
"""Runtime-only flag (env, never config.yaml) so a prior session can't auto-start REC."""
|
||||
return os.environ.get("HERMES_VOICE", "").strip() == "1"
|
||||
|
||||
|
||||
@@ -54,8 +51,7 @@ def _voice_tts_enabled() -> bool:
|
||||
|
||||
|
||||
def _end_voice_chat(*, stop_loop: bool, stop_tts: bool) -> None:
|
||||
"""Flip voice + TTS mode off; optionally halt the continuous loop / cut live TTS.
|
||||
Every step is best-effort."""
|
||||
"""Flip voice + TTS off; optionally halt the continuous loop / cut live TTS (best-effort)."""
|
||||
os.environ["HERMES_VOICE"] = "0"
|
||||
os.environ["HERMES_VOICE_TTS"] = "0"
|
||||
if stop_loop:
|
||||
@@ -68,23 +64,19 @@ def _end_voice_chat(*, stop_loop: bool, stop_tts: bool) -> None:
|
||||
|
||||
|
||||
def _tts_lease_async(lease: str, active: bool) -> None:
|
||||
"""Acquire/release a TTS engine lease off the RPC thread: acquiring warms the
|
||||
provider (local engines load a model, maybe download a voice) and must not
|
||||
block the toggle's reply. Best-effort — failure never affects the toggle."""
|
||||
|
||||
"""Acquire/release a TTS engine lease off the RPC thread: acquiring warms the provider
|
||||
(local engines load a model) and must not block the toggle's reply. Best-effort."""
|
||||
def _run():
|
||||
try:
|
||||
from tools.tts_tool import acquire_tts_lease, release_tts_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)
|
||||
|
||||
threading.Thread(target=_run, name=f"tts-lease-{lease}", daemon=True).start()
|
||||
|
||||
|
||||
def _any_session_running() -> bool:
|
||||
"""Voice busy-probe (``hermes_cli.voice.set_voice_busy_probe``): silent capture
|
||||
cycles during a long agent turn must not count toward the no-speech limit."""
|
||||
"""Voice busy-probe: silent captures during a long turn don't count toward the no-speech limit."""
|
||||
try:
|
||||
with _sessions_lock:
|
||||
return any(s.get("running") for s in _sessions.values())
|
||||
@@ -92,9 +84,8 @@ def _any_session_running() -> bool:
|
||||
return False
|
||||
|
||||
|
||||
# ── Streaming TTS (one active pipeline per process — one speaker) ──────────
|
||||
# Token deltas feed a sentence-buffering consumer (tools.tts_tool.stream_tts_to_speaker)
|
||||
# so speech starts on the first sentence; a new turn's pipeline barges in on the previous.
|
||||
# ── Streaming TTS: one pipeline per process (one speaker); a new turn's pipeline barges in
|
||||
# on the previous. Token deltas feed a sentence-buffering consumer (stream_tts_to_speaker).
|
||||
|
||||
_tts_stream_lock = threading.Lock()
|
||||
_tts_stream_state: Optional[dict] = None
|
||||
@@ -113,7 +104,8 @@ def _tts_stream_begin() -> Optional[queue.Queue]:
|
||||
_tts_stream_stop()
|
||||
text_queue: queue.Queue = queue.Queue()
|
||||
stop, done = threading.Event(), threading.Event()
|
||||
threading.Thread(target=stream_tts_to_speaker, args=(text_queue, stop, done), daemon=True).start()
|
||||
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}
|
||||
@@ -122,8 +114,8 @@ def _tts_stream_begin() -> Optional[queue.Queue]:
|
||||
|
||||
|
||||
def _tts_stream_stop(user_barge: bool = True) -> None:
|
||||
"""Cut any in-flight streaming TTS. *user_barge* latches the interruption for
|
||||
the next turn's model note — pass ``False`` for mode changes (/voice off)."""
|
||||
"""Cut in-flight streaming TTS. *user_barge* latches the interruption for the next turn's
|
||||
model note — pass ``False`` for mode changes (/voice off)."""
|
||||
global _tts_stream_state
|
||||
with _tts_stream_lock:
|
||||
state, _tts_stream_state = _tts_stream_state, None
|
||||
@@ -141,15 +133,14 @@ def _tts_stream_stop(user_barge: bool = True) -> None:
|
||||
stop_playback()
|
||||
|
||||
|
||||
# ── Full-duplex agent-turn listener (one mic, whole turn) ──────────────────
|
||||
# Arms at utterance-submit, spans generation AND playback (per-playback barge
|
||||
# monitors were deaf during generation and mis-calibrated against speaker bleed),
|
||||
# disarms when no session runs, no TTS is pending, and no audio flows.
|
||||
# ── Full-duplex agent-turn listener: arms at utterance-submit, spans generation AND playback
|
||||
# (per-playback monitors were deaf during generation and mis-calibrated against speaker
|
||||
# bleed), disarms when no session runs, no TTS is pending, and no audio flows.
|
||||
|
||||
_fd_listener_lock = threading.Lock()
|
||||
_fd_listener_active = False
|
||||
# (stop, done) pairs of fallback whole-reply speak paths: the listener must cut
|
||||
# their private stop events too, and keep listening while any is still speaking.
|
||||
# (stop, done) pairs of fallback whole-reply speak paths: the listener cuts their private stop
|
||||
# events too, and keeps listening while any is still speaking.
|
||||
_fd_speak_pipelines: "set[tuple[threading.Event, threading.Event]]" = set()
|
||||
|
||||
|
||||
@@ -184,17 +175,12 @@ def _fd_tts_pending() -> bool:
|
||||
|
||||
|
||||
def _full_duplex_listener() -> None:
|
||||
"""Mic live from utterance-submit to turn-complete; phase-aware trip.
|
||||
|
||||
Generation phase: user speech interrupts every running session's turn (the
|
||||
``agent.interrupt()`` seam ``session.interrupt`` uses) and cuts pending TTS so
|
||||
the stale reply never plays. Playback phase: cuts TTS (streaming + fallback
|
||||
speak paths + file player). Either way the utterance is transcribed and emitted
|
||||
as ``voice.transcript``; a bare stop phrase also ends the voice chat.
|
||||
"""
|
||||
"""Mic live from utterance-submit to turn-complete; a trip (see ``_fd_trip``) transcribes
|
||||
the utterance and emits it as ``voice.transcript``."""
|
||||
global _fd_listener_active
|
||||
try:
|
||||
from tools.voice_mode import full_duplex_listen, is_audio_output_active, transcribe_recording
|
||||
from tools.voice_mode import (full_duplex_listen, is_audio_output_active,
|
||||
transcribe_recording)
|
||||
|
||||
def _should_stop() -> bool:
|
||||
if not _voice_mode_enabled():
|
||||
@@ -202,16 +188,15 @@ def _full_duplex_listener() -> None:
|
||||
if _any_session_running() or _fd_tts_pending():
|
||||
return False
|
||||
return not is_audio_output_active()
|
||||
|
||||
tripped = threading.Event()
|
||||
|
||||
def _on_trigger(phase: str) -> None:
|
||||
tripped.set()
|
||||
_fd_trip(phase)
|
||||
|
||||
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)
|
||||
on_trigger=_on_trigger, multiplier=mult or None,
|
||||
grace_ms=grace_ms)
|
||||
if not (wav_path and tripped.is_set()):
|
||||
return
|
||||
try:
|
||||
@@ -245,7 +230,6 @@ def _fd_barge_params(cfg: dict) -> tuple[float, int]:
|
||||
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()
|
||||
@@ -256,13 +240,13 @@ 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)")
|
||||
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:
|
||||
@@ -296,11 +280,9 @@ def _deliver_fd_transcript(text: str) -> None:
|
||||
|
||||
|
||||
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."""
|
||||
"""Speak via hermes_cli.voice.speak_text, registered in ``_fd_speak_pipelines`` so the
|
||||
full-duplex listener can cut it and keeps listening while it is pending."""
|
||||
from hermes_cli.voice import speak_text
|
||||
|
||||
stop, done = threading.Event(), threading.Event()
|
||||
with _fd_listener_lock:
|
||||
_fd_speak_pipelines.add((stop, done))
|
||||
@@ -315,22 +297,19 @@ def _speak_text_with_barge(text: str) -> None:
|
||||
done.set()
|
||||
with _fd_listener_lock:
|
||||
_fd_speak_pipelines.discard((stop, done))
|
||||
|
||||
threading.Thread(target=_speak, daemon=True).start()
|
||||
_arm_barge_listener_if_enabled()
|
||||
|
||||
|
||||
def _voice_cfg_dict() -> dict:
|
||||
"""Shape-safe ``voice:`` block. ``_load_cfg()`` doesn't deep-merge defaults, so
|
||||
root and ``voice`` may be any YAML scalar/list/None; malformed → {}."""
|
||||
"""Shape-safe ``voice:`` block (no deep-merged defaults: any YAML shape possible; bad → {})."""
|
||||
cfg = _load_cfg()
|
||||
voice_cfg = cfg.get("voice") if isinstance(cfg, dict) else None
|
||||
return voice_cfg if isinstance(voice_cfg, dict) else {}
|
||||
|
||||
|
||||
def _voice_cfg_number(value, default):
|
||||
"""Numeric config value, else *default*. bool is excluded explicitly (int
|
||||
subclass): ``silence_threshold: true`` must not forward as ``1``."""
|
||||
"""Numeric config value, else *default*; bool excluded (``silence_threshold: true`` ≠ 1)."""
|
||||
return value if isinstance(value, (int, float)) and not isinstance(value, bool) else default
|
||||
|
||||
|
||||
@@ -340,12 +319,10 @@ def _voice_record_key() -> str:
|
||||
return str(record_key) if isinstance(record_key, str) and record_key else "ctrl+b"
|
||||
|
||||
|
||||
# ── Wake word ("Hey Hermes") ──────────────────────────────────────────────
|
||||
# Process-global detector (one mic). The first eligible transport to call
|
||||
# wake.start owns it until stop, disconnect, or stream failure. On detection we
|
||||
# emit wake.detected; the client opens a session and starts its own capture. The
|
||||
# detector yields the mic to voice.record (pause/resume) and to the desktop's
|
||||
# browser mic (wake.pause/resume RPCs).
|
||||
# ── Wake word ("Hey Hermes"): process-global detector (one mic). The first eligible transport
|
||||
# to call wake.start owns it until stop, disconnect, or stream failure; on detection we emit
|
||||
# wake.detected and the client opens a session + its own capture. The detector yields the mic
|
||||
# to voice.record (pause/resume) and to the desktop's browser mic (wake.pause/resume RPCs).
|
||||
_wake_lock = threading.Lock()
|
||||
_wake_owner_transport: "Optional[Transport]" = None
|
||||
_wake_owner_surface = ""
|
||||
@@ -390,20 +367,17 @@ def _wake_resume_if_owner(owner: "Transport", *, retry_seconds: float = 15.0,
|
||||
retry_interval: float = 1.0) -> bool:
|
||||
"""Resume the wake detector for ``owner``; self-heal a busy microphone.
|
||||
|
||||
Reopening the mic right after a voice turn can fail while the device is still
|
||||
being released (browser WebRTC tracks release async). On an exception we retry
|
||||
in a background thread until it sticks, the lease changes hands, or
|
||||
``retry_seconds`` elapses. ``False`` from ``resume_listening`` (lease gone /
|
||||
different owner) is final — never retried, so this can't steal another
|
||||
surface's mic.
|
||||
Reopening the mic right after a voice turn can fail while the device is still being
|
||||
released (browser WebRTC tracks release async): on an exception retry in a background
|
||||
thread until it sticks, the lease changes hands, or ``retry_seconds`` elapses. ``False``
|
||||
from ``resume_listening`` (lease gone / other owner) is final — never retried, so this
|
||||
can't steal another surface's mic.
|
||||
"""
|
||||
from tools.wake_word import resume_listening
|
||||
|
||||
try:
|
||||
return resume_listening(owner=owner)
|
||||
except Exception as e:
|
||||
logger.debug("wake resume failed (will retry): %s", e)
|
||||
|
||||
global _wake_resume_retry_active
|
||||
with _wake_resume_retry_lock:
|
||||
if _wake_resume_retry_active:
|
||||
@@ -429,14 +403,12 @@ def _wake_resume_if_owner(owner: "Transport", *, retry_seconds: float = 15.0,
|
||||
finally:
|
||||
with _wake_resume_retry_lock:
|
||||
_wake_resume_retry_active = False
|
||||
|
||||
threading.Thread(target=_retry, daemon=True, name="wake-resume-retry").start()
|
||||
return False
|
||||
|
||||
|
||||
def _persist_wake_enabled(enabled: bool) -> bool:
|
||||
"""Write ``wake_word.enabled`` to config.yaml. Only for explicit user gestures
|
||||
(ear toggle, /wake on|off) — never passive auto-arm paths."""
|
||||
"""Write ``wake_word.enabled``; only for explicit gestures (ear toggle, /wake on|off)."""
|
||||
try:
|
||||
from cli import save_config_value
|
||||
return bool(save_config_value("wake_word.enabled", enabled))
|
||||
@@ -446,34 +418,27 @@ def _persist_wake_enabled(enabled: bool) -> bool:
|
||||
|
||||
|
||||
def _wake_prefers_client(params: dict, surface: str) -> bool:
|
||||
"""Desktop remote (gui) prefers client capture (Mac mic → wake.feed PCM) while
|
||||
the engine runs on the backend; CLI/TUI stay local."""
|
||||
"""Desktop (gui) prefers client capture (Mac mic → wake.feed PCM); CLI/TUI stay local."""
|
||||
return surface in ("gui", "desktop") or bool(params.get("client_capture"))
|
||||
|
||||
|
||||
def _wake_probe(cfg: dict, prefer_client: bool) -> tuple[str, dict]:
|
||||
"""``(capture_mode, requirements)`` with capture stamped so the probe matches
|
||||
the mode that would actually arm."""
|
||||
"""``(capture_mode, requirements)``; capture stamped so the probe matches what would arm."""
|
||||
from tools.wake_word import check_wake_word_requirements, resolve_capture_mode
|
||||
|
||||
capture_mode = resolve_capture_mode(cfg, prefer_client=prefer_client)
|
||||
return capture_mode, check_wake_word_requirements({**cfg, "capture": capture_mode})
|
||||
|
||||
|
||||
def _wake_detect_handler(transport, sid: str, phrase: str, new_session: bool):
|
||||
"""Build the on-detect callback: pause, verify ownership, emit ``wake.detected``
|
||||
on the owner's transport."""
|
||||
|
||||
"""On-detect callback: pause, verify ownership, emit ``wake.detected`` on the owner's transport."""
|
||||
def _on_detect() -> None:
|
||||
from tools.wake_word import get_last_match, owns_listener, pause_listening
|
||||
|
||||
if not pause_listening(owner=transport) or not owns_listener(transport):
|
||||
return
|
||||
if _transport_is_dead(transport):
|
||||
_release_wake_for_transport(transport)
|
||||
return
|
||||
# Multi-phrase engines report WHICH phrase fired and its profile, so one
|
||||
# listener wakes any enrolled profile; single-phrase engines fall back.
|
||||
# Multi-phrase engines report WHICH phrase/profile fired; single-phrase engines fall back.
|
||||
matched_phrase, matched_profile = get_last_match() or (phrase, "")
|
||||
logger.info("wake.detected: emitting to sid=%r (transport=%s, profile=%r)",
|
||||
sid, type(transport).__name__, matched_profile)
|
||||
@@ -481,87 +446,79 @@ def _wake_detect_handler(transport, sid: str, phrase: str, new_session: bool):
|
||||
try:
|
||||
_emit("wake.detected", sid, {
|
||||
"phrase": matched_phrase or phrase, "profile": matched_profile or None,
|
||||
"start_new_session": new_session,
|
||||
})
|
||||
"start_new_session": new_session})
|
||||
finally:
|
||||
reset_transport(token)
|
||||
|
||||
return _on_detect
|
||||
|
||||
|
||||
@method("gateway.capabilities")
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""Advertise what THIS BUILD enforces. A client can't tell a gateway that fences
|
||||
concurrent writers from one that doesn't (both accept the same calls), so it
|
||||
withholds unless the guarantee is advertised. Sourced from the enforcing module,
|
||||
never config: a believed-but-absent capability is worse than none."""
|
||||
"""Advertise what THIS BUILD enforces: a client can't tell a gateway that fences concurrent
|
||||
writers from one that doesn't, so it withholds unless 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)})
|
||||
|
||||
|
||||
@method("ping")
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""Cheapest liveness probe, answered on the WS reader thread so it works while
|
||||
every agent is mid-turn: lets the desktop tell a half-open TCP socket after
|
||||
sleep/wake from a healthy one and force a reconnect."""
|
||||
"""Cheapest liveness probe, answered on the WS reader thread so it works while every agent
|
||||
is mid-turn: lets the desktop tell a half-open socket after sleep/wake and reconnect."""
|
||||
return _ok(rid, {"pong": True})
|
||||
|
||||
|
||||
@method("wake.start")
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""Arm the wake-word listener for the calling surface ("tui" | "gui").
|
||||
|
||||
Idempotent and gated: ``{started: False, reason}`` when disabled, scoped to
|
||||
another surface, or deps/mic aren't ready. ``persist: true`` marks an explicit
|
||||
user gesture: when disabled in config it flips ``wake_word.enabled`` on before
|
||||
arming; passive auto-arm callers omit it and keep the config-gated refusal.
|
||||
"""
|
||||
"""Arm the wake-word listener for the calling surface ("tui" | "gui"). Idempotent and gated:
|
||||
``{started: False, reason}`` when disabled, scoped to another surface, or deps/mic aren't
|
||||
ready. ``persist: true`` marks an explicit user gesture: when disabled in config it flips
|
||||
``wake_word.enabled`` on before arming; passive auto-arm callers omit it."""
|
||||
surface = str(params.get("surface") or "auto").strip().lower()
|
||||
persist = bool(params.get("persist"))
|
||||
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}")
|
||||
|
||||
cfg = load_wake_word_config()
|
||||
capture_mode, reqs = _wake_probe(cfg, _wake_prefers_client(params, surface))
|
||||
external_audio = capture_mode == "client"
|
||||
# Requirements first: a gesture on an un-armable setup must refuse WITHOUT
|
||||
# flipping wake_word.enabled — else config says on while nothing can arm.
|
||||
# Requirements first: a gesture on an un-armable setup must refuse WITHOUT flipping
|
||||
# wake_word.enabled — else config says on while nothing can arm.
|
||||
if not reqs["available"]:
|
||||
logger.warning("wake.start(%s): not available — %s", surface, reqs.get("hint"))
|
||||
return _ok(rid, {
|
||||
"started": False, "reason": "unavailable", "hint": reqs.get("hint") or "",
|
||||
"capture": capture_mode,
|
||||
})
|
||||
"capture": capture_mode})
|
||||
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).
|
||||
# "disabled" (persist:true can turn it on) vs "disabled_for_surface" (explicit
|
||||
# wake_word.surface choice, which persist does NOT override).
|
||||
reason = "disabled" if not cfg.get("enabled") else "disabled_for_surface"
|
||||
logger.info("wake.start(%s): %s (enabled=%s, surface=%s)",
|
||||
surface, reason, cfg.get("enabled"), cfg.get("surface"))
|
||||
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, 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:
|
||||
on_detect = _wake_detect_handler(transport, sid, wake_phrase(cfg), bool(cfg.get("start_new_session", True)))
|
||||
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))
|
||||
@@ -569,20 +526,17 @@ def _(rid, params: dict) -> dict:
|
||||
frame = detector_frame_info()
|
||||
logger.info(
|
||||
"wake.start(%s): listening for %r (%s) capture=%s frame=%s",
|
||||
surface, reqs["phrase"], reqs["provider"], capture_mode, frame.get("frame_length"),
|
||||
)
|
||||
surface, reqs["phrase"], reqs["provider"], capture_mode, frame.get("frame_length"))
|
||||
return _ok(rid, {
|
||||
"started": True, "phrase": reqs["phrase"], "provider": reqs["provider"],
|
||||
"owner_surface": surface, "enabled_persisted": enabled_persisted, "capture": capture_mode,
|
||||
"sample_rate": frame.get("sample_rate", 16000),
|
||||
"frame_length": frame.get("frame_length", 1280),
|
||||
})
|
||||
"frame_length": frame.get("frame_length", 1280)})
|
||||
|
||||
|
||||
@method("wake.stop")
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""Stop this surface's listener. ``persist: true`` also writes
|
||||
``wake_word.enabled: false`` so auto-arm stays off in future sessions."""
|
||||
"""Stop this surface's listener; ``persist: true`` also writes ``wake_word.enabled: false``."""
|
||||
transport = _caller_transport()
|
||||
stopped = _release_wake_for_transport(transport)
|
||||
disabled_persisted = False
|
||||
@@ -596,8 +550,7 @@ def _(rid, params: dict) -> dict:
|
||||
disabled_persisted = _persist_wake_enabled(False)
|
||||
return _ok(rid, {
|
||||
"stopped": stopped, "reason": None if stopped else "not_owner",
|
||||
"disabled_persisted": disabled_persisted,
|
||||
})
|
||||
"disabled_persisted": disabled_persisted})
|
||||
|
||||
|
||||
@method("wake.pause")
|
||||
@@ -627,8 +580,7 @@ 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,
|
||||
)
|
||||
load_wake_word_config, owns_listener, silent_audio_hint)
|
||||
cfg = load_wake_word_config()
|
||||
surface = str(params.get("surface") or "").strip().lower()
|
||||
probe_capture, reqs = _wake_probe(cfg, _wake_prefers_client(params, surface))
|
||||
@@ -643,9 +595,8 @@ def _(rid, params: dict) -> dict:
|
||||
hint = f"Wake-word input device could not be resolved: {input_device['error']}"
|
||||
if silent and not hint:
|
||||
hint = silent_audio_hint(input_device)
|
||||
# Effective capture: prefer the *armed* detector over config/auto, else
|
||||
# with capture:auto a bare status probe reports "local" and the desktop
|
||||
# never reattaches the PCM feeder after wake.detected.
|
||||
# Effective capture: prefer the *armed* detector over config/auto, else with capture:auto
|
||||
# a bare status probe reports "local" and the desktop never reattaches the PCM feeder.
|
||||
frame = detector_frame_info()
|
||||
if owned_by_caller and frame.get("external_audio"):
|
||||
capture = "client"
|
||||
@@ -664,7 +615,8 @@ def _(rid, params: dict) -> dict:
|
||||
# Armed but deaf despite an open stream; see platform-specific hint.
|
||||
"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))
|
||||
@@ -672,9 +624,8 @@ def _(rid, params: dict) -> dict:
|
||||
|
||||
@method("wake.feed")
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""Push client-captured PCM (``pcm``/``pcm_b64``: base64 int16 mono LE, 16 kHz
|
||||
only) into the armed detector — used when ``wake.start`` returned
|
||||
``capture: "client"`` so mic-less remote backends can run openWakeWord."""
|
||||
"""Push client-captured PCM (``pcm``/``pcm_b64``: base64 int16 mono LE, 16 kHz only) into
|
||||
the armed detector — for ``capture: "client"`` so mic-less remote backends run openWakeWord."""
|
||||
transport = _caller_transport()
|
||||
raw_b64 = params.get("pcm") or params.get("pcm_b64") or ""
|
||||
if not isinstance(raw_b64, str) or not raw_b64.strip():
|
||||
@@ -702,27 +653,27 @@ 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()}
|
||||
# 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()}
|
||||
try:
|
||||
from tools.voice_mode import check_voice_requirements
|
||||
reqs = check_voice_requirements()
|
||||
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 "")
|
||||
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)
|
||||
|
||||
return _ok(rid, payload)
|
||||
|
||||
|
||||
def _voice_toggle_mode(rid, params: dict) -> dict:
|
||||
enabled = params.get("action") == "on"
|
||||
# Runtime-only flag (CLI parity) — never persisted, so the next TUI launch
|
||||
# starts with voice OFF instead of auto-REC from a stale toggle.
|
||||
# Runtime-only flag — never persisted, so the next launch starts with voice OFF.
|
||||
os.environ["HERMES_VOICE"] = "1" if enabled else "0"
|
||||
|
||||
stop_hint = ""
|
||||
if enabled:
|
||||
# Spoken-stop hint for the client; sourced from voice.stop_phrases, empty when disabled.
|
||||
@@ -747,8 +698,8 @@ def _voice_toggle_mode(rid, params: dict) -> dict:
|
||||
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:
|
||||
@@ -758,8 +709,7 @@ def _voice_toggle_tts(rid, params: dict) -> dict:
|
||||
os.environ["HERMES_VOICE_TTS"] = "1" if new_value else "0"
|
||||
if not new_value:
|
||||
_tts_stream_stop(user_barge=False)
|
||||
# on → pre-load the engine so the first reply starts hot; off → release the
|
||||
# lease (last holder gone = resident local model freed).
|
||||
# on → pre-load the engine so the first reply starts hot; off → release the lease.
|
||||
_tts_lease_async("tui:voice-tts", new_value)
|
||||
# record_key on every branch so a tts toggle never resets a custom binding.
|
||||
return _ok(rid, {"enabled": True, "record_key": _voice_record_key(), "tts": new_value})
|
||||
@@ -767,15 +717,13 @@ def _voice_toggle_tts(rid, params: dict) -> dict:
|
||||
|
||||
_VOICE_TOGGLE_ACTIONS = {
|
||||
"status": _voice_toggle_status, "on": _voice_toggle_mode, "off": _voice_toggle_mode,
|
||||
"tts": _voice_toggle_tts,
|
||||
}
|
||||
"tts": _voice_toggle_tts}
|
||||
|
||||
|
||||
@method("voice.toggle")
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""CLI parity for ``/voice``: ``status``; ``on``/``off`` flip voice *mode* (off
|
||||
also tears down the continuous loop; recording itself is driven by
|
||||
``voice.record``/Ctrl+B); ``tts`` toggles speech output (requires mode on)."""
|
||||
"""CLI parity for ``/voice``: ``status``; ``on``/``off`` flip voice *mode* (off also tears
|
||||
down the continuous loop); ``tts`` toggles speech output (requires mode on)."""
|
||||
action = params.get("action", "status")
|
||||
handler = _VOICE_TOGGLE_ACTIONS.get(action) if isinstance(action, str) else None
|
||||
if handler is None:
|
||||
@@ -795,9 +743,8 @@ def _vr_on_silent():
|
||||
|
||||
|
||||
def _vr_on_stop_phrase(t):
|
||||
# The user SAID a bare stop phrase: end the chat like /voice off and emit a
|
||||
# distinct signal so clients end the conversation instead of treating it as
|
||||
# a no-speech timeout. The continuous loop has already halted.
|
||||
# A SPOKEN bare stop phrase: end the chat like /voice off and emit a distinct signal so
|
||||
# clients end the conversation instead of treating it as a no-speech timeout.
|
||||
_end_voice_chat(stop_loop=False, stop_tts=True)
|
||||
_voice_emit("voice.transcript", {"stop_phrase": True, "text": t})
|
||||
_resume_voice_wake()
|
||||
@@ -811,21 +758,17 @@ def _vr_on_status(state):
|
||||
|
||||
@method("voice.record")
|
||||
def _(rid, params: dict) -> dict:
|
||||
"""VAD-bounded push-to-talk capture, CLI-parity. ``start`` begins one capture
|
||||
and emits ``voice.transcript`` when silence stops it; ``stop`` forces
|
||||
transcription of the active buffer. The wrapper retains no-speech counts across
|
||||
starts, so three silent captures emit ``no_speech_limit=True``."""
|
||||
"""VAD-bounded push-to-talk capture. ``start`` begins one capture and emits
|
||||
``voice.transcript`` when silence stops it; ``stop`` forces transcription of the active
|
||||
buffer. No-speech counts persist across starts: three silent captures emit ``no_speech_limit``."""
|
||||
action = params.get("action", "start")
|
||||
wake_paused = False
|
||||
|
||||
if action not in {"start", "stop"}:
|
||||
return _err(rid, 4019, f"unknown voice action: {action}")
|
||||
|
||||
transport = _caller_transport()
|
||||
wake_owner, _surface = _wake_owner_snapshot()
|
||||
if wake_owner is not None and wake_owner is not transport:
|
||||
return _ok(rid, {"status": "busy", "reason": "wake_owned"})
|
||||
|
||||
try:
|
||||
global _voice_event_sid, _voice_wake_owner
|
||||
if action == "start" and not _voice_mode_enabled():
|
||||
@@ -837,7 +780,6 @@ def _(rid, params: dict) -> dict:
|
||||
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.
|
||||
@@ -861,8 +803,9 @@ def _(rid, params: dict) -> dict:
|
||||
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,
|
||||
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()
|
||||
@@ -882,8 +825,7 @@ def _(rid, params: dict) -> dict:
|
||||
if not text:
|
||||
return _err(rid, 4020, "text required")
|
||||
try:
|
||||
# Import check up front so a missing voice module returns 5026 instead
|
||||
# of failing silently in the thread.
|
||||
# Import check up front so a missing voice module returns 5026, not a silent thread death.
|
||||
import hermes_cli.voice # noqa: F401
|
||||
threading.Thread(target=_speak_text_with_barge, args=(text,), daemon=True).start()
|
||||
return _ok(rid, {"status": "speaking"})
|
||||
@@ -894,5 +836,4 @@ def _(rid, params: dict) -> dict:
|
||||
|
||||
|
||||
def register(server) -> None:
|
||||
"""Publish this module's helpers + handlers onto ``server``, rebound to its globals."""
|
||||
bind_module(globals(), server, skip=("_",))
|
||||
|
||||
Reference in New Issue
Block a user