Files
hermes-agent/gateway/slash_commands.py
ethernet 97de4b4ef8 Merge origin/main into ethie/pm-clean
Conflict resolutions and semantic fixups:

- utils.py / hermes_yaml.py: main widened ruamel's round-trip emitter so a long
  double-quoted scalar is never folded after an escaped backslash. pm-clean builds
  every rt emitter through hermes_yaml.roundtrip_yaml(), so the width lives there
  (ROUNDTRIP_YAML_WIDTH moves with it); xai_retirement imports it from hermes_yaml.
- hermes_cli/banner.py: keep pm-clean's removal of the banner update check. Main's
  GIT_NO_LAZY_FETCH fix for it applies to its replacement, source_check: every
  read-only probe (source_git_env) now refuses promisor lazy fetches, and the
  partial-clone test targets that probe (red without the flag).
- .github/workflows/tests.yml: keep setup-pm; main's uv pin bump does not apply.
  Main's WAL-capable SQLite gates are kept, run against $HERMES_PYTHON (the
  PM-pinned interpreter, SQLite 3.53.1). The e2e step takes main's
  --include-integration invocation.
- apps/desktop: package.json has no build block here, so main's macOS locale-marker
  restore joins the darwin branch of the existing after-pack.mjs, and its test
  loads the hook from electron-builder.config.cjs and imports PlatformPackager
  from app-builder-lib's root (electron-builder 27 exports no ./out paths). The
  win32 row is dropped: this hook sanitizes and signs PE trees on win32 by design.
- reconciliation.ts: main's rowId hydration (#119326) was merged into the first of
  pm-clean's split helpers only; the resolver is now one helper both halves use.
- en.ts: both sides' keys kept. tests/tools/test_lazy_deps.py stays deleted.
- Tests main added with `import yaml` use hermes_yaml, like the rest of the tree.
2026-09-23 19:51:34 -04:00

1366 lines
78 KiB
Python

"""Gateway slash-command handlers for GatewayRunner: lifted out of ``gateway/run.py`` into a mixin
so ``self._handle_*_command`` keeps resolving via the MRO. Cohesive clusters live in the sibling
mixins (``slash_commands_model/_session/_status/_goals``); this module keeps the shared helpers plus
the one-off commands. run.py helpers are imported lazily."""
from __future__ import annotations
import asyncio
import contextlib
import dataclasses
import inspect
import logging
import os
import re
import shlex
import sys
import time
from datetime import datetime
from pathlib import Path
from typing import Optional, Union
from agent.i18n import t
from gateway.config import HomeChannel, Platform, PlatformConfig, persist_home_channel
from gateway.platforms.base import EphemeralReply
from gateway.platforms.event import MessageEvent
from gateway.session import AsyncSessionStore
from gateway.session_transcript import TranscriptReadError
from gateway.slash_commands_goals import GatewayGoalCommandsMixin
from gateway.slash_commands_model import GatewayModelCommandsMixin
from gateway.slash_commands_session import GatewaySessionCommandsMixin
from gateway.slash_commands_login import GatewayLoginCommandsMixin
from gateway.slash_commands_status import HISTORY_UNREADABLE, GatewayStatusCommandsMixin
from hermes_cli.config import atomic_config_write, cfg_get
from utils import atomic_json_write, is_truthy_value
logger = logging.getLogger("gateway.run")
# /rollback result keys -> i18n line for files the safe restore left alone.
_ROLLBACK_SKIP_LINES = (("skipped_user_edits", "gateway.rollback.kept_user_edits"),
("skipped_oversize", "gateway.rollback.kept_oversize"),
("failed_deletes", "gateway.rollback.failed_deletes"))
# /busy input modes -> (status-card behavior, set-confirmation behavior).
_BUSY_MODE_BEHAVIOR = {
"queue": ("queues for next turn", "Messages will be queued for the next turn while Hermes is busy."),
"steer": ("steers into current run (after next tool call)",
"Messages will be steered into the current run (after the next tool call)."),
"interrupt": ("interrupts current run", "Messages will interrupt the current run while Hermes is busy."),
}
# /diff argument -> diff mode (unknown args leave the mode unchanged).
_DIFF_MODE_BY_ARG = {**dict.fromkeys(("staged", "--staged", "cached", "--cached"), "staged"),
**dict.fromkeys(("all", "--all", "head"), "all"), "session": "session"}
# /voice subcommand -> stored mode (None = auto-TTS disabled), confirmation i18n key.
_VOICE_MODE_BY_ARG = {
**dict.fromkeys(("on", "enable"), ("voice_only", "gateway.voice.enabled_voice_only")),
**dict.fromkeys(("off", "disable"), ("off", "gateway.voice.disabled_text")),
"tts": ("all", "gateway.voice.tts_enabled")}
# /footer argument -> new enabled state ("" toggles; anything else is a usage error).
_FOOTER_STATE_BY_ARG = {**dict.fromkeys(("on", "enable", "true", "1"), True),
**dict.fromkeys(("off", "disable", "false", "0"), False)}
# /approve modifier tokens -> approval choice (default "once").
_APPROVE_CHOICE_BY_ARG = {**dict.fromkeys(("always", "permanent", "permanently"), "always"),
**dict.fromkeys(("session", "ses"), "session")}
_PLATFORM_USAGE = ("Usage: /platform <list|pause|resume> [name]\n"
" /platform list — show platform status\n"
" /platform pause <name> — stop retrying a failing platform\n"
" /platform resume <name> — re-queue a paused platform")
_WINDOWS_UPDATE_HELPER = """
import os, subprocess, sys
output_path, exit_code_path, cmd = sys.argv[1], sys.argv[2], sys.argv[3:]
env = dict(os.environ, PYTHONUNBUFFERED="1")
with open(output_path, "wb") as f:
rc = subprocess.Popen(cmd, stdout=f, stderr=subprocess.STDOUT, env=env).wait(timeout=3600)
with open(exit_code_path, "w", encoding="utf-8") as f:
f.write(str(rc))
""".strip()
def _nested_dict(root: dict, *keys: str) -> dict:
"""Walk/create ``root[k1][k2]...`` as dicts, replacing any non-dict value on the path."""
for k in keys:
if not isinstance(root.get(k), dict):
root[k] = {}
root = root[k]
return root
def _write_raw_config_leaf(config_path: Path, keys: tuple, value) -> None:
"""Set one leaf through a strict raw round-trip. The behavioral read is fail-open (``{}``) and
expanded, so writing it back wipes the file after a read error and persists ``${VAR}`` values."""
from hermes_cli.config import read_user_config_raw
raw = read_user_config_raw(config_path)
*parents, leaf = keys
_nested_dict(raw, *parents)[leaf] = value
atomic_config_write(config_path, raw)
def _preview(text: str, limit: int = 60) -> str:
return text[:limit] + ("..." if len(text) > limit else "")
def _execute(command: str, **ctx_kwargs):
"""Run *command* through the shared slash executor on the gateway surface."""
from hermes_cli.slash_exec import CommandContext, execute_command
return execute_command(command, CommandContext(surface="gateway", **ctx_kwargs))
def _restart_notify_payload(event: MessageEvent) -> dict:
"""Requester routing info so the new gateway process can notify them once back online.
``profile`` is persisted so the notice leaves through the requester's own profile bot after the
restart (a bare platform lookup would resolve the default profile's adapter)."""
source = event.source
data = {"platform": source.platform.value if source.platform else None,
"chat_id": source.chat_id, "chat_type": source.chat_type}
if source.delivered_via_upstream_relay is True:
data["delivered_via_upstream_relay"] = True
data.update({k: getattr(source, k) for k in ("user_id", "scope_id") if getattr(source, k)})
optional = (("thread_id", source.thread_id), ("message_id", event.message_id),
("profile", getattr(source, "profile", None)))
data.update({k: v for k, v in optional if v})
return data
def _spawn_detached_update(hermes_cmd, output_path, exit_code_path) -> None:
"""Spawn ``hermes update --gateway`` detached so it survives the gateway restart it may trigger.
setsid is portable (works where ``systemd-run --user`` lacks a D-Bus session); ``--gateway``
enables file-based IPC so interactive prompts are forwarded; PYTHONUNBUFFERED lets the gateway
stream output live. Windows has no setsid: an inline helper runs the updater as a module under
this interpreter (not venv\\Scripts\\hermes.exe — that shim holds its own file open, and the
update must replace it), redirects both outputs to one file and writes the exit code."""
import shutil
import subprocess
if sys.platform == "win32":
from hermes_cli._subprocess_compat import windows_detach_popen_kwargs
subprocess.Popen(
[sys.executable, "-c", _WINDOWS_UPDATE_HELPER, str(output_path), str(exit_code_path),
sys.executable, "-m", "hermes_cli.main", "update", "--gateway"],
stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, **windows_detach_popen_kwargs())
return
hermes_cmd_str = " ".join(shlex.quote(part) for part in hermes_cmd)
update_cmd = (
f"PYTHONUNBUFFERED=1 {hermes_cmd_str} update --gateway"
f" > {shlex.quote(str(output_path))} 2>&1; "
# Avoid `status=$?`: `status` is read-only in zsh and this template is reused in
# macOS/zsh operator wrappers, so keep it zsh-safe even though bash runs it here.
f"rc=$?; printf '%s' \"$rc\" > {shlex.quote(str(exit_code_path))}")
# Preferred: setsid creates a new session, fully detached; fallback start_new_session=True
# calls os.setsid() in the child.
setsid_bin = shutil.which("setsid")
argv = [setsid_bin, "bash", "-c", update_cmd] if setsid_bin else ["bash", "-c", update_cmd]
subprocess.Popen(argv, stdout=subprocess.DEVNULL, stderr=subprocess.DEVNULL, start_new_session=True)
def _home_thread_from_source(source) -> Optional[str]:
"""The thread id /sethome should persist on the home target, or None. Slack thread-per-message
keying stamps a top-level message's own id as ``source.thread_id`` (a session key, not a
location); persisting it would pin HOME to that ephemeral thread. A thread id equal to the
message's own id is synthetic and dropped; a real thread (id = parent's) is kept."""
thread_id = getattr(source, "thread_id", None)
if not thread_id:
return None
synthetic = (getattr(source, "platform", None) == Platform.SLACK and getattr(source, "message_id", None)
and str(thread_id) == str(source.message_id))
return None if synthetic else str(thread_id)
class GatewaySlashCommandsMixin(
GatewayLoginCommandsMixin,
GatewayModelCommandsMixin,
GatewaySessionCommandsMixin,
GatewayStatusCommandsMixin,
GatewayGoalCommandsMixin):
"""In-session slash-command handlers for GatewayRunner (plus the helpers the sibling mixins share)."""
async_session_store: AsyncSessionStore
# ------------------------------------------------------------------ shared helpers
def _cached_agent_for(self, session_key: str, *, lockless_fallback: bool = False):
"""Peek the cached AIAgent for *session_key* without evicting it, or None. Entries are
``(agent, signature, ...)`` tuples (bare agents from test doubles accepted). Historical callers
read the cache ONLY under ``_agent_cache_lock`` and got None when a fixture that skipped
``__init__`` had no lock; the manual codex ``/compress`` path was the one exception that read
lock-free (``lockless_fallback=True``)."""
cache = getattr(self, "_agent_cache", None)
lock = getattr(self, "_agent_cache_lock", None)
if cache is None or (lock is None and not lockless_fallback):
return None
try:
if lock:
with lock:
entry = cache.get(session_key)
else:
entry = cache.get(session_key)
except Exception:
return None
return (entry[0] if entry else None) if isinstance(entry, (tuple, list)) else entry or None
def _resident_agent_for(self, session_key: str):
"""The live running agent for *session_key*, else the cached one, else None. The pending
sentinel (a run that is starting) never counts as a usable agent."""
from gateway.run import _AGENT_PENDING_SENTINEL
agent = self._running_agents.get(session_key)
if agent is not None and agent is not _AGENT_PENDING_SENTINEL:
return agent
return self._cached_agent_for(session_key)
@staticmethod
def _session_db_unavailable_reply() -> str:
from hermes_state import format_session_db_unavailable
return format_session_db_unavailable(prefix=t("gateway.shared.session_db_unavailable_prefix"))
def _reply_metadata(self, event: MessageEvent):
"""Thread/reply metadata for an outbound send anchored on *event*."""
return self._thread_metadata_for_source(event.source, self._reply_anchor_for_event(event))
def _adapter_and_key_for(self, event: MessageEvent):
"""``(adapter, session_key)`` for the event's source, either None when no source. The source's
OWN transport (profile-aware, fail-closed) — ``self.adapters`` is the default profile's map."""
if not event.source:
return None, None
return self._delivery_adapter_for(event.source), self._session_key_for_source(event.source)
def _telegramized_command_reply(self, event: MessageEvent, text: str) -> str:
from gateway.run import _telegramize_command_mentions
return _telegramize_command_mentions(text, getattr(getattr(event, "source", None), "platform", None))
def _checkpoint_manager(self):
"""A CheckpointManager from gateway config, or None when checkpoints are disabled."""
from gateway.run import _checkpoint_agent_kwargs, _load_gateway_config
from tools.checkpoint_manager import CheckpointManager
cp = _checkpoint_agent_kwargs(_load_gateway_config())
if not cp["checkpoints_enabled"]:
return None
# AIAgent kwargs are ``checkpoint_<field>``; CheckpointManager takes the bare field names.
fields = {k[len("checkpoint_"):]: v for k, v in cp.items() if k.startswith("checkpoint_")}
return CheckpointManager(enabled=True, **fields)
def _write_approval_setter(self, section: str, event: MessageEvent):
"""``set_mode_fn`` for /memory and /skills: persist ``<section>.write_approval``. Raw read is
correct for the write-back round-trip (merged defaults must not be persisted back to the
user's file); the cached agent is dropped so the setting takes effect next message."""
from gateway.run import _gateway_config_home
# Persist to config (default) unless --session opted out, mirroring the text /model command path
# above so a picked model survives across sessions like a typed one (#49066).
from hermes_cli.config import read_user_config_raw
config_path = _gateway_config_home() / "config.yaml"
session_key = self._session_key_for_source(event.source)
def _set_approval(enabled: bool):
user_config = read_user_config_raw(config_path)
user_config.setdefault(section, {})["write_approval"] = bool(enabled)
atomic_config_write(config_path, user_config)
# Evict any cached agent for this session so the next message rebuilds with the correct
# session_id end-to-end — mirrors /branch and /reset. Without this, the cached AIAgent (and its
# memory provider, which cached `_session_id` during initialize()) keeps writing into the wrong
# session's record. See #6672.
self._evict_cached_agent(session_key)
return _set_approval
async def _deliver_approval_confirmation(self, event: MessageEvent, confirmation_text: str, verb: str):
"""Return *confirmation_text* for normal delivery, or push it on native-streaming adapters
(WeCom msgtype:"stream"), which need it sent directly with control-lane metadata (reliable
proactive send, not the finalized reply stream). ``is not True``: mocks auto-create attrs."""
source = event.source
adapter = self._delivery_adapter_for(source) # the receiving bot, not the default profile's
if adapter:
adapter.resume_typing_for_chat(source.chat_id) # agent is about to continue
if getattr(adapter, "SUPPORTS_NATIVE_STREAMING", False) is not True:
return confirmation_text
if adapter:
try:
await adapter.send(
source.chat_id, confirmation_text, reply_to=event.message_id,
metadata={"is_approval_prompt": True, "force_proactive_send": True})
except Exception as exc:
logger.warning("Failed to send /%s confirmation to %s: %s", verb, source.chat_id,
exc, exc_info=True)
return None
def _typed_command_prefix_for(self, platform) -> str:
"""The prefix users can always type to reach Hermes commands (adapter ``typed_command_prefix``,
default "/"). Slack and Matrix use "!" because typed "/" is blocked/reserved there; their
adapters rewrite "!command" to "/command"."""
adapter = self.adapters.get(platform) if getattr(self, "adapters", None) else None
return getattr(adapter, "typed_command_prefix", "/") if adapter is not None else "/"
def _terminal_cwd(self) -> str:
from tools.terminal_scope import terminal_env
return terminal_env("TERMINAL_CWD", str(Path.home()))
@staticmethod
def _display_config_target(event: MessageEvent):
"""``(config.yaml path, platform config key)`` for the per-platform display settings."""
from gateway.run import _gateway_config_home, _platform_config_key
return _gateway_config_home() / "config.yaml", _platform_config_key(event.source.platform)
async def _handle_profile_command(self, event: MessageEvent) -> str:
"""Handle /profile — show the profile serving this source and its home. On a multiplexed
gateway the process-level profile is the multiplexer's own ("default" in every chat), so
with ``multiplex_profiles`` on report ``source.profile`` and resolve home under that
profile's runtime scope; when off the stamp is ignored, mirroring ``_run_agent``."""
from hermes_constants import display_hermes_home
source = getattr(event, "source", None)
profile_name = display = ""
if getattr(getattr(self, "config", None), "multiplex_profiles", False):
profile_name = (getattr(source, "profile", "") or "").strip()
try:
from gateway.run import _profile_runtime_scope
with _profile_runtime_scope(self._resolve_profile_home_for_source(source)):
display = display_hermes_home()
except Exception:
display = display_hermes_home()
# Shared executor resolves process-level fallbacks; the multiplexed per-source overrides
# (when any) ride in via options.
reply = _execute("profile", options={"profile_name": profile_name, "home_display": display})
return "\n".join([t("gateway.profile.header", profile=reply.data["profile"]),
t("gateway.profile.home", home=reply.data["home"])])
async def _handle_whoami_command(self, event: MessageEvent) -> str:
"""Handle /whoami — platform, DM-vs-group scope, tier and runnable commands (always allowed)."""
from gateway.slash_access import policy_for_source
source = event.source
policy = policy_for_source(self.config, source)
platform = source.platform.value if source and source.platform else "?"
chat_type = ((source.chat_type if source else "") or "dm").lower()
scope = "DM" if chat_type in {"dm", "direct", "private", ""} else "group/channel"
user_id = (source.user_id if source else None) or "?"
head = f"**You** — {platform} ({scope})\nUser ID: `{user_id}`\n"
if not policy.enabled:
return head + "Tier: unrestricted (no admin list configured for this scope)\nSlash commands: all available"
if policy.is_admin(user_id):
return head + "Tier: **admin**\nSlash commands: all available"
# Non-admin: floor first (mirrors slash_access._ALWAYS_ALLOWED_FOR_USERS), then operator
# additions, deduped in order.
runnable = list(dict.fromkeys(["help", "whoami"] + sorted(policy.user_allowed_commands)))
runnable_str = ", ".join(f"/{c}" for c in runnable) if runnable else "(none)"
return head + f"Tier: user\nSlash commands you can run: {runnable_str}"
async def _handle_kanban_command(self, event: MessageEvent) -> str:
"""Handle /kanban — delegate to the shared kanban CLI (DB work in a thread pool). Allowed
while an agent runs: the board is profile-agnostic and never touches agent state."""
from hermes_cli.kanban import run_slash
# Strip the leading "/kanban" (with or without slash), leaving args.
text = (event.text or "").strip().lstrip("/")
if text.startswith("kanban"):
text = text[len("kanban"):].lstrip()
requested_board = action = None
tokens = iter(shlex.split(text) if text else [])
for tok in tokens: # leading --board/--board=<b> options, then the action verb
if tok == "--board":
requested_board = next(tokens, requested_board)
elif tok.startswith("--board="):
requested_board = tok.split("=", 1)[1]
else:
action = tok
break
try:
output = await asyncio.to_thread(run_slash, text)
except Exception as exc: # pragma: no cover - defensive
return t("gateway.kanban.error_prefix", error=exc)
# Auto-subscribe on create, parsing the task id from the CLI's standard success line
# ("Created t_abcd (ready, ...)"). With --json there is no such line, so a scripting user
# gets no subscription and can call /kanban notify-subscribe explicitly.
m = re.search(r"Created\s+(t_[0-9a-f]+)\b", output) if action == "create" and output else None
if m:
task_id = m.group(1)
try:
if await self._kanban_auto_subscribe(event, task_id, requested_board):
output = output.rstrip() + "\n" + t("gateway.kanban.subscribed_suffix", task_id=task_id)
except Exception as exc:
logger.warning("kanban create auto-subscribe failed: %s", exc)
# Gateway messages have practical length caps; truncate long listings.
if len(output) > 3800:
output = output[:3800] + "\n" + t("gateway.kanban.truncated_suffix")
return output or t("gateway.kanban.no_output")
async def _kanban_auto_subscribe(self, event: MessageEvent, task_id: str, requested_board) -> bool:
"""Subscribe the event's chat to *task_id* notifications (notify+wake). False when the
source has no platform/chat to route back to."""
source = event.source
def _field(name: str) -> Optional[str]:
return str(getattr(source, name, "") or "") or None
platform = getattr(source, "platform", None)
platform_str = (platform.value if hasattr(platform, "value") else str(platform or "")).lower()
chat_id, chat_type = _field("chat_id"), _field("chat_type")
delivery_metadata = self._reply_metadata(event) or None
if isinstance(delivery_metadata, dict) and chat_type:
delivery_metadata.setdefault("chat_type", chat_type)
if not (platform_str and chat_id):
return False
def _sub():
from hermes_cli import kanban_db as _kb
from hermes_cli import kanban_db_connect as _kbc
from hermes_cli import kanban_db_notify as _kbn
conn = _kbc.connect(board=requested_board)
try:
_kbn.add_notify_sub(
conn, task_id=task_id, platform=platform_str, chat_id=chat_id, chat_type=chat_type,
thread_id=_field("thread_id"), user_id=_field("user_id"),
# Also persist the stable alt id (Signal UUID, Feishu union_id): build_session_key
# keys the participant on ``user_id_alt or user_id``, so a replayed wake rebuilds
# the same session key only when the alt id survives the round-trip.
user_id_alt=_field("user_id_alt"),
notifier_profile=_field("profile") or getattr(self, "_kanban_notifier_profile", None) or self._active_profile_name(),
# Subscribing from chat: deliver the passive message and wake the destination agent.
delivery_mode="notify+wake", delivery_metadata=delivery_metadata)
finally:
conn.close()
await asyncio.to_thread(_sub)
return True
async def _handle_stop_command(self, event: MessageEvent) -> Union[str, EphemeralReply]:
"""Handle /stop command - interrupt a running agent. A truly hung agent (blocked thread
never checking _interrupt_requested) is caught by the early intercept in _handle_message();
this handler runs via normal dispatch or as a fallback, and force-cleans the session lock in
all cases. The session is preserved so the user can continue."""
from gateway.run import _AGENT_PENDING_SENTINEL, _INTERRUPT_REASON_STOP
source = event.source
session_entry = await self.async_session_store.get_or_create_session(source)
session_key = session_entry.session_key
async def _stop(key: str, invalidation_reason: str) -> None:
await self._interrupt_and_clear_session(
key, source, interrupt_reason=_INTERRUPT_REASON_STOP,
invalidation_reason=invalidation_reason)
agent = self._running_agents.get(session_key)
if agent is _AGENT_PENDING_SENTINEL: # force-clean the sentinel so the session is unlocked
await _stop(session_key, "stop_command_pending")
logger.info("STOP (pending) for session %s — sentinel cleared", session_key)
return EphemeralReply(t("gateway.stop.stopped_pending"))
if agent: # force-clean the session lock so a truly hung agent doesn't keep it forever
await _stop(session_key, "stop_command_handler")
return EphemeralReply(t("gateway.stop.stopped"))
# No run under the caller's own key: a live turn in THIS chat may still carry a differently
# shaped key. One scan feeds both tiers; the chat tier is a superset of the thread-sibling
# tier (a sibling needs the caller's own thread slot, which satisfies the chat predicate), so
# it is the set to act on — acting on the sibling subset alone would reply "Stopped" while a
# same-thread run under a differently shaped key kept going. See `_chat_scoped_run_keys` for
# the shapes and isolation bounds; both tiers are authorization-gated.
runs = self._same_chat_runs(source, session_key)
sibling_keys = self._sibling_thread_run_keys(source, runs)
fallback_keys = self._chat_scoped_run_keys(source, runs)
# Reason is per-stop, not per-key: a stop that only ever had thread siblings keeps its own
# label for hook consumers, anything wider is a chat-scope stop.
reason = (
"stop_command_thread_sibling"
if fallback_keys == sibling_keys
else "stop_command_chat_scope"
)
if fallback_keys and self._is_user_authorized_for_source(source):
for fallback_key in fallback_keys:
await _stop(fallback_key, reason)
logger.info("STOP (%s) by %s — interrupted %d run(s): %s",
reason, session_key, len(fallback_keys), ", ".join(fallback_keys))
return EphemeralReply(t("gateway.stop.stopped"))
# No running agent anywhere for this scope. Background delegations the session dispatched in an
# earlier turn still count as "active": stop them; each returns as an interrupted completion.
from tools.async_delegation import interrupt_for_session
if interrupt_for_session(session_key=session_key, reason="stop_command",
parent_session_id=str(getattr(session_entry, "session_id", "") or "")):
return EphemeralReply(t("gateway.stop.stopped"))
# A platform status indicator can still be stuck —
# e.g. Slack's persistent assistant.threads.setStatus survives a gateway restart or a turn
# that died without a final send.
# Best-effort clear so /stop always dismisses a phantom "is thinking...". See #32295.
adapter = getattr(self, "adapters", {}).get(source.platform)
try:
if adapter and hasattr(adapter, "_stop_typing_with_metadata"):
await adapter._stop_typing_with_metadata(source.chat_id, self._reply_metadata(event))
except Exception:
logger.debug("Failed to clear typing on /stop with no active agent", exc_info=True)
return t("gateway.stop.no_active")
async def _handle_platform_command(self, event: MessageEvent) -> str:
"""Handle ``/platform list|pause|resume [name]`` — inspect and manually control failed/paused
adapters (pause stops the reconnect watcher; resume re-queues for retry)."""
# Strip the leading "/platform" (or "/PLATFORM") token if present
parts = (getattr(event, "content", "") or "").strip().split(maxsplit=2)
if parts and parts[0].lower().lstrip("/").startswith("platform"):
parts = parts[1:]
action = (parts[0] if parts else "list").lower()
target = parts[1].lower() if len(parts) > 1 else ""
failed = getattr(self, "_failed_platforms", {}) or {}
if action == "list":
connected = ", ".join(sorted(p.value for p in self.adapters)) or "(none)"
lines = ["**Gateway platforms**", f"Connected: {connected}"]
for p, info in failed.items():
if info.get("paused"):
reason = info.get("pause_reason") or "paused"
lines.append(f" · {p.value} — PAUSED ({reason}). Resume with `/platform resume {p.value}`.")
else:
lines.append(f" · {p.value} — retrying (attempt {info.get('attempts', 0)})")
return "\n".join(lines + ([] if failed else ["Failed/paused: (none)"]))
if action not in {"pause", "resume"}:
return _PLATFORM_USAGE
if not target:
return f"Usage: /platform {action} <name>"
# Resolve platform name (case-insensitive, value match)
platform = next((p for p in Platform.__members__.values() if p.value.lower() == target), None)
if platform is None:
return f"Unknown platform: {target}"
name = platform.value
queued = platform in failed
paused = queued and bool(failed[platform].get("paused"))
if action == "pause":
if not queued:
return f"{name} is not in the retry queue (it's either connected or not enabled)."
if paused:
return f"{name} is already paused."
self._pause_failed_platform(platform, reason="paused via /platform pause")
return f"✓ {name} paused. Resume with `/platform resume {name}` or `hermes gateway restart` to reset."
if not queued:
return f"{name} is not in the retry queue — nothing to resume."
if not paused:
return f"{name} is already retrying — no resume needed."
self._resume_paused_platform(platform)
return f"✓ {name} resumed — retrying on next watcher tick."
async def _handle_restart_command(self, event: MessageEvent) -> Union[str, EphemeralReply]:
"""Handle /restart command - drain active work, then restart the gateway."""
from gateway.run import _hermes_home
# Idempotency check: if the previous gateway process recorded this same /restart (platform +
# update_id) and we see it *again*, it's a redelivery from PTB's graceful-shutdown get_updates
# ACK failing on the way out. Ignoring it prevents a loop where every fresh gateway re-restarts.
if self._is_stale_restart_redelivery(event):
src = event.source
logger.info("Ignoring redelivered /restart (platform=%s, update_id=%s) — "
"already processed by a previous gateway instance.",
src.platform.value if src and src.platform else "?",
event.platform_update_id)
return ""
if self._restart_requested or self._draining:
count = self._running_agent_count()
return t("gateway.draining", count=count) if count else EphemeralReply(t("gateway.restart.in_progress"))
async def _write_marker(name: str, build, label: str) -> None:
try:
await asyncio.to_thread(atomic_json_write, _hermes_home / name, build(), indent=None)
except Exception as e:
logger.debug("Failed to write restart %s: %s", label, e)
def _notify_payload() -> dict:
data = _restart_notify_payload(event)
mid = str(event.message_id) if event.message_id is not None else event.source.message_id
try:
self._restart_command_source = dataclasses.replace(event.source, message_id=mid)
except Exception:
self._restart_command_source = event.source
return data
def _dedup_payload() -> dict:
# Platform + update_id of the triggering /restart, for redelivery detection.
data = {"platform": event.source.platform.value if event.source.platform else None,
"requested_at": time.time()}
if event.platform_update_id is not None:
data["update_id"] = event.platform_update_id
return data
# Save the requester's routing info so the new gateway process can notify them once back.
await _write_marker(".restart_notify.json", _notify_payload, "notify file")
# Record the triggering platform + update_id in a dedicated dedup marker. Unlike
# .restart_notify.json (unlinked once the new gateway sends its notification) this persists
# so a delayed Telegram redelivery is still detectable. Overwritten on every /restart.
await _write_marker(".restart_last_processed.json", _dedup_payload, "dedup marker")
active_agents = self._running_agent_count()
# Under a service manager (systemd/launchd) or Docker/Podman, exit 75 so the supervisor /
# restart policy restarts us — detached setsid+bash fails there (systemd KillMode=mixed kills
# the cgroup; tini exits with the gateway). The explicit marker covers ``sudo env -i`` wrappers.
from gateway.restart import is_container_restart_context, is_gateway_supervisor_process
via_service = is_gateway_supervisor_process() or is_container_restart_context()
self.request_restart(detached=not via_service, via_service=via_service)
# Track sessions that were active at shutdown for stuck-loop detection (#7536). On each restart, the
# counter increments for sessions that were running. If a session hits the threshold (3 consecutive
# restarts while active), the next startup auto-suspends it — breaking the loop.
if active_agents:
return t("gateway.draining", count=active_agents)
return EphemeralReply(t("gateway.restart.restarting"))
async def _handle_version_command(self, event: MessageEvent) -> str:
"""Handle /version — show the running Hermes Agent version."""
return _execute("version").text
def _catalog_options(self, event: MessageEvent) -> dict:
"""``allowed_commands`` for /help and /commands when the caller is a gated non-admin:
the slash-access floor + ``user_allowed_commands`` (mirrors /whoami), so the catalog
never advertises commands ``_check_slash_access`` would refuse. Admins / ungated -> {}."""
from gateway.slash_access import policy_for_source
source = event.source
# ``getattr``: partially-constructed runners (``GatewayRunner.__new__`` in tests) have
# no ``config``; policy_for_source treats None as ungated.
policy = policy_for_source(getattr(self, "config", None), source)
if policy.enabled and not policy.is_admin(source.user_id if source else None):
return {"allowed_commands": {"help", "whoami", *policy.user_allowed_commands}}
return {}
async def _handle_help_command(self, event: MessageEvent) -> str:
"""Handle /help command - list available commands."""
return self._telegramized_command_reply(
event, _execute("help", options=self._catalog_options(event)).text)
async def _handle_commands_command(self, event: MessageEvent) -> str:
# Page size is a surface parameter (Telegram messages are shorter).
page_size = 15 if event.source.platform == Platform.TELEGRAM else 20
options = {"page_size": page_size, **self._catalog_options(event)}
reply = _execute("commands", args=event.get_command_args(), options=options)
return self._telegramized_command_reply(event, reply.text)
async def _handle_set_home_command(self, event: MessageEvent) -> str:
"""Handle /sethome command -- set the current chat as the platform's home channel."""
from gateway.run import _home_target_env_var, _home_thread_env_var
source = event.source
platform_name = source.platform.value if source.platform else "unknown"
chat_id = source.chat_id
chat_name = source.chat_name or chat_id
if source.platform is None:
return t("gateway.set_home.save_failed", error="Missing logical platform")
via_relay = getattr(source, "delivered_via_upstream_relay", False) is True
if via_relay:
adapter_for_source = getattr(self, "_intake_adapter_for", None)
relay_adapter = adapter_for_source(source) if callable(adapter_for_source) else None
fronts_platform = getattr(relay_adapter, "fronts_platform", None)
if (source.platform in {None, Platform.LOCAL, Platform.RELAY}
or not getattr(source, "user_id", None)
or not callable(fronts_platform) or not fronts_platform(source.platform)):
return t("gateway.set_home.save_failed",
error="Relay does not authenticate this logical home target")
thread_id = _home_thread_from_source(source)
home = HomeChannel(
platform=source.platform, chat_id=str(chat_id), name=chat_name, thread_id=thread_id,
user_id=str(source.user_id) if getattr(source, "user_id", None) else None,
scope_id=str(source.scope_id) if getattr(source, "scope_id", None) else None)
# config.yaml is canonical because it can persist the authenticated logical-target
# provenance required by Relay after a restart.
try:
persist_home_channel(home, enabled_if_new=not via_relay)
except Exception as e:
return t("gateway.set_home.save_failed", error=e)
# Preserve legacy home env vars for existing cron/setup consumers.
try:
from hermes_cli.config import save_env_value
save_env_value(_home_target_env_var(platform_name), str(chat_id))
save_env_value(_home_thread_env_var(platform_name), str(thread_id or ""))
except Exception as e:
logger.warning("Home config saved but legacy env persistence failed: %s", e)
# Keep the running gateway config in sync too. The pre-restart notification path reads
# self.config before the process reloads config.
platform_config = self.config.platforms.setdefault(source.platform, PlatformConfig(enabled=not via_relay))
platform_config.home_channel = home
return t("gateway.set_home.success", name=chat_name, chat_id=chat_id)
async def _handle_voice_command(self, event: MessageEvent) -> str:
"""Handle /voice [on|off|tts|channel|leave|status] command."""
args = event.get_command_args().strip().lower()
chat_id = event.source.chat_id
# Voice state belongs to the (bot, chat) pair: resolve the adapter that received the
# command and key the mode by its owning profile so two multiplexed bots in one chat keep
# independent /voice state.
# See #75198.
voice_key = self._voice_key_for_source(event.source)
adapter = self._delivery_adapter_for(event.source)
def _set_mode(mode: str) -> None:
self._voice_mode[voice_key] = mode
self._save_voice_modes()
if not adapter:
return
if mode == "off":
self._set_adapter_auto_tts_disabled(adapter, chat_id, disabled=True)
else:
self._set_adapter_auto_tts_enabled(adapter, chat_id, enabled=True)
if args in _VOICE_MODE_BY_ARG:
mode, reply_key = _VOICE_MODE_BY_ARG[args]
_set_mode(mode)
return t(reply_key)
if args in {"channel", "join"}:
return await self._handle_voice_channel_join(event)
if args == "leave":
return await self._handle_voice_channel_leave(event)
if args == "status":
mode = self._voice_mode.get(voice_key, "off")
label = t(f"gateway.voice.label_{mode}") if mode in ("off", "voice_only", "all") else mode
lines = [t("gateway.voice.status_mode", label=label)]
guild_id = self._get_guild_id(event) # append voice channel info if connected
info = adapter.get_voice_channel_info(guild_id) if guild_id and hasattr(adapter, "get_voice_channel_info") else None
if info:
lines += [t("gateway.voice.status_channel", channel=info['channel_name']),
t("gateway.voice.status_participants", count=info['member_count'])]
for m in info["members"]:
status = t("gateway.voice.speaking") if m.get("is_speaking") else ""
lines.append(t("gateway.voice.status_member", name=m['display_name'], status=status))
return "\n".join(lines)
# Toggle: off → on, on/all → off
turning_on = self._voice_mode.get(voice_key, "off") == "off"
_set_mode("voice_only" if turning_on else "off")
toggle_line = t("gateway.voice.enabled_short" if turning_on else "gateway.voice.disabled_short")
# Bare /voice still toggles, but append an explainer so users discover the on/off/tts/status
# subcommands (and, on Discord, live voice-channel join/leave). Toggle result shows first.
supports_voice_channels = adapter is not None and hasattr(adapter, "join_voice_channel")
channels = t("gateway.voice.help_channels") if supports_voice_channels else ""
return t("gateway.voice.help", toggle=toggle_line, channels=channels)
async def _handle_rollback_command(self, event: MessageEvent) -> str:
"""Handle /rollback command — list or restore filesystem checkpoints."""
from tools.checkpoint_manager import format_checkpoint_list
mgr = self._checkpoint_manager()
if mgr is None:
return t("gateway.rollback.not_enabled")
cwd = self._terminal_cwd()
# --all / --force: classic full restore, overwriting user edits too.
tokens = event.get_command_args().strip().split()
restore_all = any(tok.lower() in ("--all", "--force") for tok in tokens)
arg = " ".join(tok for tok in tokens if tok.lower() not in ("--all", "--force"))
# Container-backed session: host checkpoints belong to another tree, so a restore is
# refused; the bare listing stays visible, prefixed with the reason (same as the CLI).
reason = mgr.unsupported_backend_reason()
if reason and arg:
return reason
checkpoints = mgr.list_checkpoints(cwd)
if not arg:
listing = format_checkpoint_list(checkpoints, cwd)
return f"{reason}\n{listing}" if reason else listing
if not checkpoints:
return t("gateway.rollback.none_found", cwd=cwd)
# Restore by number or hash
try:
idx = int(arg) - 1
except ValueError:
target_hash = arg
else:
if not 0 <= idx < len(checkpoints):
return t("gateway.rollback.invalid_number", max=len(checkpoints))
target_hash = checkpoints[idx]["hash"]
result = mgr.restore(cwd, target_hash, safe=not restore_all)
if not result["success"]:
return t("gateway.rollback.restore_failed", error=result["error"])
msg = t("gateway.rollback.restored", hash=result["restored_to"], reason=result["reason"])
for result_key, i18n_key in _ROLLBACK_SKIP_LINES:
files = result.get(result_key) or []
if files:
more = f" (+{len(files) - 5})" if len(files) > 5 else ""
msg += "\n" + t(i18n_key, files=", ".join(files[:5]) + more)
return msg
async def _handle_diff_command(self, event: MessageEvent) -> str:
"""Handle /diff — show git changes in the working directory. Diff body is truncated hard
here (chat is not a pager); platform senders clamp further."""
args = [a.lower() for a in event.get_command_args().strip().split()]
stat_only = bool({"--stat", "stat"} & set(args))
mode = "working"
for low in args:
mode = _DIFF_MODE_BY_ARG.get(low, mode)
cwd = self._terminal_cwd()
if mode == "session":
# Cumulative checkpoint-baseline diff.
mgr = self._checkpoint_manager()
if mgr is None:
return t("gateway.diff.not_enabled")
if reason := mgr.unsupported_backend_reason(): # host baseline is not this session's tree
return reason
result = await asyncio.to_thread(mgr.session_diff, cwd)
else:
from tools.working_diff import collect_working_diff
result = await asyncio.to_thread(collect_working_diff, cwd, mode)
if not result.get("success"):
return t("gateway.diff.failed", error=result.get("error", "Could not generate diff"))
return self._render_diff_result(result, stat_only)
def _render_diff_result(self, result: dict, stat_only: bool) -> str:
"""Render a working/session diff result: stat block, untracked list, fenced (truncated) diff."""
stat = result.get("stat", "")
diff = result.get("diff", "")
untracked = result.get("untracked", [])
if result.get("empty") or (not stat and not diff and not untracked):
return t("gateway.diff.no_changes")
out: list[str] = []
if stat:
out.append(f"```\n{stat}\n```")
if untracked:
shown = "\n".join(f"+ {rel}" for rel in untracked[:15])
more = f"\n... and {len(untracked) - 15} more" if len(untracked) > 15 else ""
out.append(f"**Untracked:**\n```\n{shown}{more}\n```")
if not stat_only and diff:
out.append(self._fenced_truncated_diff(diff))
return "\n\n".join(out)
@staticmethod
def _fenced_truncated_diff(diff: str, max_lines: int = 60, max_chars: int = 3000) -> str:
"""Fence a diff body, truncating to messaging-friendly size."""
diff_lines = diff.splitlines()
truncated = len(diff_lines) > max_lines
if truncated:
diff = "\n".join(diff_lines[:max_lines])
if len(diff) > max_chars:
diff = diff[:max_chars]
truncated = True
note = ""
if truncated:
note = f"\n... (truncated — {len(diff_lines)} lines total; use /diff --stat for a summary)"
return f"```diff\n{diff}{note}\n```"
def _track_background_task(self, coro) -> None:
"""Fire-and-forget *coro*, keeping a strong ref in ``_background_tasks`` until it finishes."""
task = asyncio.create_task(coro)
self._background_tasks.add(task)
task.add_done_callback(self._background_tasks.discard)
async def _handle_background_command(self, event: MessageEvent) -> str:
"""Handle /bg <prompt> — run a prompt in a background thread with its own session; the
result is sent to the same chat without touching the active session's history."""
prompt = event.get_command_args().strip()
if not prompt:
return t("gateway.background.usage")
task_id = f"bg_{datetime.now().strftime('%H%M%S')}_{os.urandom(3).hex()}"
self._track_background_task(self._run_background_task(
prompt, event.source, task_id, event_message_id=self._reply_anchor_for_event(event),
# Forward image/audio attachments so the background agent can see them.
media_urls=list(event.media_urls or []), media_types=list(event.media_types or [])))
return t("gateway.background.started", preview=_preview(prompt), task_id=task_id)
async def _handle_btw_command(self, event: MessageEvent) -> str:
"""Handle /btw <question> — one-shot auxiliary LLM call on a transcript snapshot; live history
is never touched (alternation + prompt cache intact, current turn keeps running). Unlike /bg,
which spawns a fresh contextless session."""
question = event.get_command_args().strip()
if not question:
return t("gateway.btw.usage")
source = event.source
session_entry = await self.async_session_store.get_or_create_session(source)
try:
history = await self.async_session_store.load_transcript(session_entry.session_id)
except TranscriptReadError:
return HISTORY_UNREADABLE
if not history:
return t("gateway.btw.no_history")
try:
model, rt = self._resolve_session_agent_runtime(source=source)
except Exception:
model, rt = None, {}
if not rt.get("api_key"):
return t("gateway.btw.no_provider")
main_runtime = {
"model": model,
**{k: rt.get(k) for k in ("provider", "base_url", "api_key", "api_mode")},
"session_id": session_entry.session_id,
}
history_snapshot = list(history)
# Prefer the cache-parity fork when a live cached AIAgent exists: it replays the snapshot
# against the warm provider prefix cache, giving FULL context at cache-read prices. With no
# cached agent the cache is cold anyway — answer_side_question's digest fallback handles it.
try:
parent_agent = self._cached_agent_for(self._session_key_for_source(source))
except Exception:
parent_agent = None
_thread_metadata = self._reply_metadata(event)
adapter = self._delivery_adapter_for(source)
preview = _preview(question)
async def _run_side_question() -> None:
from agent.side_question import answer_side_question
try:
answer = await asyncio.to_thread(
answer_side_question, question, history_snapshot,
parent_agent=parent_agent, main_runtime=main_runtime)
reply = t("gateway.btw.answer", preview=preview, answer=answer or "")
except Exception as e:
logger.warning("/btw side question failed: %s", e)
reply = t("gateway.btw.failed", preview=preview, error=str(e))
if adapter is not None:
await adapter.send(source.chat_id, reply, metadata=_thread_metadata)
self._track_background_task(_run_side_question())
return t("gateway.btw.started", preview=preview)
async def _handle_memory_command(self, event: MessageEvent) -> str:
"""Handle /memory — review pending memory writes + toggle the approval gate. Entries are small
enough to review inline, so the full flow works on every platform."""
from hermes_cli.write_approval_commands import handle_pending_subcommand
from tools import write_approval as wa
from tools.memory_tool import load_on_disk_store
# Apply approved writes against a fresh on-disk store (the gateway has no long-lived agent;
# the store persists to the same MEMORY/USER.md and honors the configured char limits).
out = handle_pending_subcommand(
wa.MEMORY, event.get_command_args().strip().split(), memory_store=load_on_disk_store(),
set_mode_fn=self._write_approval_setter("memory", event))
return out if out is not None else (
"Unknown /memory subcommand. Use: pending, approve <id>, reject <id>, approval <on|off>."
)
async def _handle_skills_command(self, event: MessageEvent) -> str:
"""Handle /skills on the gateway — pending skill-write review only (hub stays CLI-only). Gated
by ``skills.write_approval`` but still answers when staged writes exist after the gate is off
(never stranded). ``diff`` is truncated for chat."""
from hermes_cli.write_approval_commands import handle_pending_subcommand
from tools import write_approval as wa
args = event.get_command_args().strip().split()
sub = args[0].lower() if args else ""
gate_off = not wa.write_approval_enabled(wa.SKILLS) and sub not in {"approval", "mode"}
if gate_off and wa.pending_count(wa.SKILLS) == 0:
return ("Skill write approval is off (skills.write_approval). "
"Enable it with /skills approval on, then review staged "
"writes here with /skills pending.")
out = handle_pending_subcommand(
wa.SKILLS, args, set_mode_fn=self._write_approval_setter("skills", event))
if out is None:
return ("Unknown /skills subcommand on this platform. Use: pending, "
"approve <id>, reject <id>, diff <id>, approval <on|off>. "
"(Search/install are CLI-only.)")
# Chat bubbles can't hold a full skill diff — truncate and point at the pending JSON file
# (NOT `hermes skills diff <name>`, which diffs a bundled skill against its stock version).
if sub == "diff" and len(out) > 3000:
pending_id = args[1] if len(args) > 1 else "<id>"
out = (out[:3000]
+ "\n… (truncated — full diff in "
f"~/.hermes/pending/skills/{pending_id}.json)")
return out
async def _handle_approvals_command(self, event: MessageEvent) -> str:
"""Show or persist the profile-wide dangerous-command approval mode."""
from gateway.slash_access import policy_for_source
from hermes_cli.approval_mode import run_approval_mode_command
requested = event.get_command_args().strip() or None
# This mutates profile-wide security policy. The central slash gate can allow selected
# commands to non-admin users, so enforce admin again at this side-effect boundary.
# Unconfigured policies remain unrestricted.
policy = policy_for_source(self.config, event.source)
if requested and not policy.is_admin(event.source.user_id):
return "Only gateway admins can change the persistent approval mode."
# Approval checks load config dynamically; do not evict the cached agent or alter its
# system prompt/tool schema (prompt-cache prefix is sacred).
return run_approval_mode_command(requested).message
async def _handle_yolo_command(self, event: MessageEvent) -> Union[str, EphemeralReply]:
"""Handle /yolo — toggle dangerous command approval bypass for this session only."""
from tools.approval import disable_session_yolo, enable_session_yolo, is_session_yolo_enabled
session_key = self._session_key_for_source(event.source)
if is_session_yolo_enabled(session_key):
disable_session_yolo(session_key)
return EphemeralReply(t("gateway.yolo.disabled"))
enable_session_yolo(session_key)
return EphemeralReply(t("gateway.yolo.enabled"))
async def _handle_verbose_command(self, event: MessageEvent) -> str:
"""Handle /verbose — cycle tool progress display mode (off → new → all → verbose → log) per
*current platform*, saved to ``display.platforms.<platform>.tool_progress``. Gated by
``display.tool_progress_command`` (default off)."""
from gateway.run import _load_gateway_config
config_path, platform_key = self._display_config_target(event)
try:
user_config = _load_gateway_config()
gate_enabled = is_truthy_value(cfg_get(user_config, "display", "tool_progress_command"),
default=False)
except Exception:
gate_enabled = False
if not gate_enabled:
return t("gateway.verbose.not_enabled")
# Cycle mode (per-platform), reading the current effective mode via the resolver.
from gateway.display_config import resolve_display_setting
cycle = ["off", "new", "all", "verbose", "log"]
current = resolve_display_setting(user_config, platform_key, "tool_progress", "all")
new_mode = cycle[(cycle.index(current if current in cycle else "all") + 1) % len(cycle)]
description = t(f"gateway.verbose.mode_{new_mode}")
try:
_write_raw_config_leaf(config_path, ("display", "platforms", platform_key, "tool_progress"), new_mode)
return f"{description}\n" + t("gateway.verbose.saved_suffix", platform=platform_key)
except Exception as e:
logger.warning("Failed to save tool_progress mode: %s", e)
return f"{description}\n" + t("gateway.verbose.save_failed", error=e)
async def _handle_busy_command(self, event: MessageEvent) -> Union[str, EphemeralReply]:
"""Handle /busy — control what happens when messaging while Hermes is working."""
arg = event.get_command_args().strip().lower()
if not arg or arg == "status":
mode = self._effective_busy_input_mode(event.source)
behavior = _BUSY_MODE_BEHAVIOR.get(mode, _BUSY_MODE_BEHAVIOR["interrupt"])[0]
return EphemeralReply(
f"**Busy input mode: `{mode}`\nMessages while busy: _{behavior}_\n"
f"Change with `/busy queue`, `/busy steer`, or `/busy interrupt`.")
if arg not in _BUSY_MODE_BEHAVIOR:
return EphemeralReply(
f"Unknown mode `{arg}`. Use `/busy queue`, `/busy steer`, or `/busy interrupt`.")
# Persist before mutate
from cli import save_config_value
if not save_config_value("display.busy_input_mode", arg):
return EphemeralReply("Busy input mode could not be saved to config. Mode unchanged.")
profile_name = self._busy_profile_name_for_source(event.source)
if profile_name:
from gateway.run import _load_gateway_config
self._snapshot_profile_busy_modes(profile_name, _load_gateway_config())
else:
self._busy_input_mode = arg
# busy_input_mode is also the source of truth for the text mode — re-derive it so the
# adapter refresh below doesn't keep a stale value and keep interrupting.
self._busy_text_mode = self._load_busy_text_mode()
adapter = self._delivery_adapter_for(event.source)
if adapter is not None:
adapter._busy_text_mode = self._effective_busy_text_mode(event.source)
return EphemeralReply(
f"Busy input mode set to **`{arg}`** (saved).\n_{_BUSY_MODE_BEHAVIOR[arg][1]}_")
async def _handle_footer_command(self, event: MessageEvent) -> str:
"""Handle /footer command — toggle the runtime-metadata footer."""
from gateway.run import _load_gateway_config, _resolve_gateway_model
from gateway.runtime_footer import format_runtime_footer, resolve_footer_config
config_path, platform_key = self._display_config_target(event)
arg = ""
try:
text = (getattr(event, "message", None) or "").strip()
if text.startswith("/"):
parts = text.split(None, 1)
arg = parts[1].strip().lower() if len(parts) > 1 else ""
except Exception:
arg = ""
try:
user_config: dict = _load_gateway_config()
except Exception as e:
return t("gateway.config_read_failed", error=e)
effective = resolve_footer_config(user_config, platform_key)
def _state(enabled: bool) -> str:
return t("gateway.footer.state_on") if enabled else t("gateway.footer.state_off")
if arg in {"status", "?"}:
return t("gateway.footer.status", state=_state(effective["enabled"]),
fields=", ".join(effective.get("fields") or []), platform=platform_key)
if arg and arg not in _FOOTER_STATE_BY_ARG:
return t("gateway.footer.usage")
new_state = _FOOTER_STATE_BY_ARG[arg] if arg else not effective["enabled"]
try:
_write_raw_config_leaf(config_path, ("display", "runtime_footer", "enabled"), new_state)
except Exception as e:
logger.warning("Failed to save runtime_footer.enabled: %s", e)
return t("gateway.config_save_failed", error=e)
example = ""
if new_state:
# Show a preview using current agent state if available.
preview = format_runtime_footer(
model=_resolve_gateway_model(user_config) or None, context_tokens=0, context_length=None,
fields=effective.get("fields") or ["model", "context_pct", "cwd"])
if preview:
example = t("gateway.footer.example_line", preview=preview)
return t("gateway.footer.saved", state=_state(new_state), example=example)
async def _handle_reload_mcp_command(self, event: MessageEvent) -> Optional[str]:
"""Handle /reload-mcp — reconnect MCP servers and rebuild the cached agent. Reloading
invalidates the provider prompt cache (tool schemas live in the system prompt), so it routes
through slash-confirm; "Always Approve" persists ``approvals.mcp_reload_confirm: false``."""
session_key = self._session_key_for_source(event.source)
# Read the gate fresh from disk so a prior "always" click takes effect on the next
# invocation without restarting the gateway.
user_config = self._read_user_config()
approvals = user_config.get("approvals") if isinstance(user_config, dict) else None
if isinstance(approvals, dict) and not approvals.get("mcp_reload_confirm", True):
return await self._execute_mcp_reload(event)
# Route through slash-confirm. The primitive sends the prompt and stores the resume handler;
# the button/text response triggers ``_resolve_slash_confirm`` which invokes the handler
# with the chosen outcome.
async def _on_confirm(choice: str) -> Optional[str]:
if choice == "cancel":
return t("gateway.reload_mcp.cancelled")
if choice == "always":
# Persist the opt-out and run the reload.
try:
from cli import save_config_value
save_config_value("approvals.mcp_reload_confirm", False)
logger.info("User opted out of /reload-mcp confirmation (session=%s)", session_key)
except Exception as exc:
logger.warning("Failed to persist mcp_reload_confirm=false: %s", exc)
# once / always → run the reload
result = await self._execute_mcp_reload(event)
if choice == "always":
return f"{result}\n\n" + t("gateway.reload_mcp.always_followup")
return result
return await self._request_slash_confirm(
event=event, command="reload-mcp", title="/reload-mcp",
message=t("gateway.reload_mcp.confirm_prompt"), handler=_on_confirm)
async def _handle_reload_skills_command(self, event: MessageEvent) -> str:
"""Handle /reload-skills — rescan skills dir, queue a note for next turn. Skills are invoked at
runtime, not baked into the system prompt, so this does NOT clear the prompt cache. The diff
goes into ``_pending_skills_reload_notes[session_key]``, prepended to the NEXT user message —
nothing out-of-band, so alternation is preserved."""
try:
from agent.skill_commands import reload_skills
# _run_in_executor_with_context, not a bare hop: the rescan walks
# get_hermes_home()/skills, a contextvar override under multiplex.
result = await self._run_in_executor_with_context(reload_skills)
added, removed = result.get("added", []), result.get("removed", []) # [{"name", "description"}]
total = result.get("total", 0)
# Let adapters refresh platform-side state that cached the skill list at startup (today:
# Discord /skill autocomplete — otherwise new skills stay invisible and deleted ones
# error). Adapters without refresh_skill_group are skipped; the in-process reload suffices.
for adapter in list(self.adapters.values()):
refresh = getattr(adapter, "refresh_skill_group", None)
try:
maybe = refresh() if callable(refresh) else None
if inspect.isawaitable(maybe):
await maybe
except Exception as exc:
logger.warning("Adapter %s refresh_skill_group raised: %s",
getattr(adapter, "name", adapter), exc)
lines = [t("gateway.reload_skills.header")]
if not added and not removed:
lines += [t("gateway.reload_skills.no_new"), t("gateway.reload_skills.total", count=total)]
return "\n".join(lines)
def _fmt_line(item: dict) -> str:
nm, desc = item.get("name", ""), item.get("description", "")
return (t("gateway.reload_skills.item_with_desc", name=nm, desc=desc) if desc
else t("gateway.reload_skills.item_no_desc", name=nm))
# Queue a one-shot note for the next user turn in this session too. Format matches how
# the system prompt renders pre-existing skills (`` - name: description``) so the
# model reads the diff in the same shape as its original skill catalog.
sections = ["[USER INITIATED SKILLS RELOAD:"]
for i18n_key, note_header, items in (
("gateway.reload_skills.added_header", "Added Skills:", added),
("gateway.reload_skills.removed_header", "Removed Skills:", removed)):
if items:
formatted = [_fmt_line(item) for item in items]
lines += [t(i18n_key)] + formatted
sections += ["", note_header] + formatted
lines.append(t("gateway.reload_skills.total", count=total))
sections += ["", "Use skills_list to see the updated catalog.]"]
session_key = self._session_key_for_source(event.source)
if not hasattr(self, "_pending_skills_reload_notes"):
self._pending_skills_reload_notes = {}
if session_key:
self._pending_skills_reload_notes[session_key] = "\n".join(sections)
return "\n".join(lines)
except Exception as e:
logger.warning("Skills reload failed: %s", e)
return t("gateway.reload_skills.failed", error=e)
async def _handle_bundles_command(self, event: MessageEvent) -> str:
"""Handle /bundles — list installed skill bundles (mirrors the CLI handler). Bundles are
loaded by invoking their own ``/<slug>`` command, not by this one."""
reply = _execute("bundles")
if "error" in reply.data:
logger.warning("Bundles command unavailable: %s", reply.data["error"])
return reply.text
bundles = reply.data["bundles"]
if not bundles:
return ("No skill bundles installed.\nCreate one on the host with:\n"
" `hermes bundles create <name> --skill <s1> --skill <s2>`\n"
f"Directory: `{reply.data['dir']}`")
lines = [f"**Skill Bundles** ({len(bundles)} installed):", ""]
for info in bundles:
skills = info.get("skills", [])
desc = info.get("description") or f"Load {len(skills)} skills"
lines += [f"• `/{info['slug']}` — {desc} _({len(skills)} skills)_"] + [f" · {s}" for s in skills]
return "\n".join(lines + ["", "Invoke a bundle with `/<slug>` to load all its skills."])
def _blocking_approval_or_stale(self, event: MessageEvent, stale_key: str, none_key: str):
"""``(session_key, None)`` when an agent thread is blocked on approval, else the reply to send.
A pending-approvals entry with no blocked thread is a stale prompt: drop it and say so."""
from tools.approval import has_blocking_approval
session_key = self._session_key_for_source(event.source)
if has_blocking_approval(session_key):
return session_key, None
if session_key in self._pending_approvals:
self._pending_approvals.pop(session_key)
return session_key, t(stale_key)
return session_key, t(none_key)
async def _handle_approve_command(self, event: MessageEvent) -> Optional[str]:
"""Handle /approve — unblock waiting agent thread(s). They block inside tools/approval.py;
signalling the event resumes them so the command executes inline (same flow as the CLI)."""
from tools.approval import resolve_gateway_approval
session_key, stale = self._blocking_approval_or_stale(event, "gateway.approval_expired",
"gateway.approve.no_pending")
if stale:
return stale
# Args: "all", "all session", "all always", "session", "always" ("always" beats "session").
args = event.get_command_args().strip().lower().split()
choices = {_APPROVE_CHOICE_BY_ARG[a] for a in args if a in _APPROVE_CHOICE_BY_ARG}
choice = "always" if "always" in choices else "session" if "session" in choices else "once"
count = resolve_gateway_approval(session_key, choice, resolve_all="all" in args)
if not count:
return t("gateway.approve.no_pending")
confirmation_text = t(f"gateway.approve.{choice}_{'plural' if count > 1 else 'singular'}", count=count)
logger.info("User approved %d dangerous command(s) via /approve (%s)", count, choice)
return await self._deliver_approval_confirmation(event, confirmation_text, "approve")
async def _handle_deny_command(self, event: MessageEvent) -> str:
"""Handle /deny — reject pending dangerous command(s) with a definitive BLOCKED result, as in
the CLI. ``/deny`` denies the oldest; ``/deny all`` denies everything.
``/deny <reason>`` (or ``/deny all <reason>``) attaches a one-line reason that is relayed back to
the agent so it can adapt instead of only hearing "denied". Ported from qwibitai/nanoclaw#2832.
"""
from tools.approval import resolve_gateway_approval
session_key, stale = self._blocking_approval_or_stale(event, "gateway.deny.stale",
"gateway.deny.no_pending")
if stale:
return stale
# A leading "all" denies every pending command; the rest (or the whole arg string without
# "all") is the optional deny reason relayed to the agent, capped to a sane one-liner.
raw_args = event.get_command_args().strip()
tokens = raw_args.split()
resolve_all = bool(tokens) and tokens[0].lower() == "all"
reason = (raw_args[len(tokens[0]):].strip() if resolve_all else raw_args)[:280].strip()
count = resolve_gateway_approval(session_key, "deny", resolve_all=resolve_all, reason=reason or None)
if not count:
return t("gateway.deny.no_pending")
logger.info("User denied %d dangerous command(s) via /deny%s", count,
" (with reason)" if reason else "")
key = "gateway.deny.denied" + ("_reason" if reason else "") + ("_plural" if count > 1 else "_singular")
confirmation_text = t(key, count=count, reason=reason)
return await self._deliver_approval_confirmation(event, confirmation_text, "deny")
async def _handle_debug_command(self, event: MessageEvent) -> str:
"""Handle /debug — upload ONLY the summary (system info + log tails), never full logs, to
protect privacy; ``hermes debug share`` from the CLI does full uploads."""
from hermes_cli.debug import (_GATEWAY_PRIVACY_NOTICE, _best_effort_sweep_expired_pastes,
_capture_dump, _is_dpaste_url, _schedule_auto_delete,
collect_debug_report, upload_to_pastebin)
def _collect_and_upload(): # blocking I/O (dump capture, log reads, uploads) -> thread
_best_effort_sweep_expired_pastes()
report = collect_debug_report(log_lines=200, dump_text=_capture_dump())
try:
urls = {"Report": upload_to_pastebin(report)}
except Exception as exc:
return t("gateway.debug.upload_failed", error=exc)
_schedule_auto_delete(list(urls.values())) # paste.rs only; dpaste.com has no delete
label_width = max(len(k) for k in urls)
# The 6-hour line is only true for paste.rs; the privacy notice above already states
# the dpaste.com fallback retention, so drop the line rather than contradict it.
auto_delete = [] if any(map(_is_dpaste_url, urls.values())) else [
t("gateway.debug.auto_delete")]
return "\n".join([_GATEWAY_PRIVACY_NOTICE, "", t("gateway.debug.header"), "",
*(f"`{label:<{label_width}}` {url}" for label, url in urls.items()),
"", *auto_delete, t("gateway.debug.full_logs_hint"),
t("gateway.debug.share_hint")])
# _run_in_executor_with_context, not a bare hop: this collects the profile's logs/config off
# ``get_hermes_home()`` and uploads them to a public paste. Losing the contextvar override
# would publish the DEFAULT profile's diagnostics from another profile's chat.
return await self._run_in_executor_with_context(_collect_and_upload)
async def _handle_update_command(self, event: MessageEvent) -> str:
"""Handle /update — spawn ``hermes update`` detached (``setsid``) so it survives the gateway
restart it may trigger; marker files let this or the next gateway process notify the user."""
import json
from gateway.run import _hermes_home, _resolve_hermes_bin
from hermes_cli.config import is_managed, format_managed_message
# Block non-messaging platforms (API server, webhooks, ACP); plugin platforms with
# allow_update_command=True are also allowed.
src = event.source
if src.platform not in self._UPDATE_ALLOWED_PLATFORMS:
try:
from gateway.platform_registry import platform_registry
entry = platform_registry.get(src.platform.value)
if not entry or not entry.allow_update_command:
return t("gateway.update.platform_not_messaging")
except Exception:
return t("gateway.update.platform_not_messaging")
if is_managed():
return f"✗ {format_managed_message('update Hermes Agent')}"
project_root = Path(__file__).parent.parent.resolve()
# Not a git-managed install (docker/nix/desktop-app/source): refuse
# with the steward's own update mechanism instead of git-pulling a
# tree `hermes update` does not own.
try:
from hermes_cli.config import (
detect_install_method,
recommended_update_command_for_method,
)
method = detect_install_method(project_root)
if method not in {"git", "unknown"}:
return (
f"✗ `hermes update` does not apply to this install ({method}).\n"
f"Update with: {recommended_update_command_for_method(method)}"
)
except Exception:
pass # config unreadable — fall through to the .git check below
git_dir = project_root / '.git'
if not git_dir.exists():
return t("gateway.update.not_git_repo")
hermes_cmd = _resolve_hermes_bin()
if not hermes_cmd:
return t("gateway.update.hermes_cmd_not_found")
pending_path = _hermes_home / ".update_pending.json"
output_path = _hermes_home / ".update_output.txt"
exit_code_path = _hermes_home / ".update_exit_code"
pending = {
"platform": src.platform.value, "chat_id": src.chat_id, "chat_type": src.chat_type,
"user_id": src.user_id, "session_key": self._session_key_for_source(src),
"timestamp": datetime.now().isoformat()}
# ``profile``: the update watcher (possibly the NEXT gateway process) must answer through the
# requester's own profile bot, not the default profile's adapter for the same platform.
pending.update({k: v for k, v in (("thread_id", src.thread_id), ("message_id", event.message_id),
("profile", getattr(src, "profile", None))) if v})
_tmp_pending = pending_path.with_suffix(".tmp")
_tmp_pending.write_text(json.dumps(pending), encoding="utf-8")
_tmp_pending.replace(pending_path)
exit_code_path.unlink(missing_ok=True)
try:
_spawn_detached_update(hermes_cmd, output_path, exit_code_path)
except Exception as e:
pending_path.unlink(missing_ok=True)
exit_code_path.unlink(missing_ok=True)
return t("gateway.update.start_failed", error=e)
self._schedule_update_notification_watch()
return t("gateway.update.starting")
# ---- BEGIN PLUGIN-COMPAT (revert-scheduled; see COMPAT_MANIFEST.md) ----
# Names external plugins imported from this module before the Sep 2026 decomposition.
# Internal code MUST NOT use these (scripts/check_compat_pointers.py fails CI if it does).
# The whole block is removed by reverting the commit that added it.
from typing import Any # noqa: F401,E402
import hashlib # noqa: F401,E402
_PLUGIN_COMPAT_LAZY = {
'HISTORY_UNREADABLE': ('gateway.slash_commands_status', 'HISTORY_UNREADABLE'),
'MessageType': ('gateway.platforms.event', 'MessageType'),
'SessionSource': ('gateway.session', 'SessionSource'),
'base_url_host_matches': ('utils', 'base_url_host_matches'),
'build_session_key': ('gateway.session', 'build_session_key'),
'clear_model_endpoint_credentials': ('hermes_cli.config', 'clear_model_endpoint_credentials'),
'extract_api_content_sidecar': ('agent.turn_context', 'extract_api_content_sidecar'),
'fetch_account_usage': ('agent.account_usage', 'fetch_account_usage'),
'is_shared_multi_user_session': ('gateway.session', 'is_shared_multi_user_session'),
'render_account_usage_lines': ('agent.account_usage', 'render_account_usage_lines'),
}
def __getattr__(name): # PEP 562 — lazy so no import cycles
target = _PLUGIN_COMPAT_LAZY.get(name)
if target is None:
raise AttributeError(f"module {__name__!r} has no attribute {name!r}")
import importlib
from hermes_cli.plugin_compat import warn_once
warn_once(__name__, name, *target)
return getattr(importlib.import_module(target[0]), target[1])
# ---- END PLUGIN-COMPAT ----