refactor(gateway): fail-open helpers in status mixin; tighten local import blocks
This commit is contained in:
+18
-63
@@ -10,6 +10,7 @@ run.py helpers are imported lazily inside handler bodies to avoid the import cyc
|
||||
from __future__ import annotations
|
||||
|
||||
import asyncio
|
||||
import contextlib
|
||||
import dataclasses
|
||||
import inspect
|
||||
import logging
|
||||
@@ -123,7 +124,6 @@ def _preview(text: str, limit: int = 60) -> str:
|
||||
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))
|
||||
|
||||
|
||||
@@ -161,10 +161,8 @@ def _spawn_detached_update(hermes_cmd, output_path, exit_code_path) -> None:
|
||||
"""
|
||||
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,
|
||||
@@ -231,12 +229,8 @@ class GatewaySlashCommandsMixin(
|
||||
cache = getattr(self, "_agent_cache", None)
|
||||
if cache is None:
|
||||
return None
|
||||
lock = getattr(self, "_agent_cache_lock", None)
|
||||
try:
|
||||
if lock is not None:
|
||||
with lock:
|
||||
entry = cache.get(session_key)
|
||||
else:
|
||||
with getattr(self, "_agent_cache_lock", None) or contextlib.nullcontext():
|
||||
entry = cache.get(session_key)
|
||||
except Exception:
|
||||
return None
|
||||
@@ -250,7 +244,6 @@ class GatewaySlashCommandsMixin(
|
||||
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
|
||||
@@ -279,15 +272,12 @@ class GatewaySlashCommandsMixin(
|
||||
"""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_kwargs = _checkpoint_agent_kwargs(_load_gateway_config())
|
||||
if not cp_kwargs["checkpoints_enabled"]:
|
||||
cp = _checkpoint_agent_kwargs(_load_gateway_config())
|
||||
if not cp["checkpoints_enabled"]:
|
||||
return None
|
||||
return CheckpointManager(
|
||||
enabled=True,
|
||||
max_snapshots=cp_kwargs["checkpoint_max_snapshots"],
|
||||
max_total_size_mb=cp_kwargs["checkpoint_max_total_size_mb"],
|
||||
max_file_size_mb=cp_kwargs["checkpoint_max_file_size_mb"],
|
||||
enabled=True, max_snapshots=cp["checkpoint_max_snapshots"],
|
||||
max_total_size_mb=cp["checkpoint_max_total_size_mb"], max_file_size_mb=cp["checkpoint_max_file_size_mb"],
|
||||
)
|
||||
|
||||
def _write_approval_setter(self, section: str, event: MessageEvent):
|
||||
@@ -326,15 +316,11 @@ class GatewaySlashCommandsMixin(
|
||||
if adapter:
|
||||
try:
|
||||
await adapter.send(
|
||||
source.chat_id,
|
||||
confirmation_text,
|
||||
reply_to=event.message_id,
|
||||
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,
|
||||
)
|
||||
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:
|
||||
@@ -346,7 +332,6 @@ class GatewaySlashCommandsMixin(
|
||||
|
||||
def _terminal_cwd(self) -> str:
|
||||
from tools.terminal_scope import terminal_env
|
||||
|
||||
return terminal_env("TERMINAL_CWD", str(Path.home()))
|
||||
|
||||
async def _handle_profile_command(self, event: MessageEvent) -> str:
|
||||
@@ -358,14 +343,12 @@ class GatewaySlashCommandsMixin(
|
||||
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:
|
||||
@@ -383,7 +366,6 @@ class GatewaySlashCommandsMixin(
|
||||
"""Handle /whoami — the user's slash command access on this scope (always allowed: slash_access
|
||||
floor). Reports platform, DM-vs-group scope, tier, and the commands the user can run here."""
|
||||
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 "?"
|
||||
@@ -658,7 +640,6 @@ class GatewaySlashCommandsMixin(
|
||||
# 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)
|
||||
if active_agents:
|
||||
@@ -700,10 +681,7 @@ class GatewaySlashCommandsMixin(
|
||||
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",
|
||||
)
|
||||
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(
|
||||
@@ -766,12 +744,8 @@ class GatewaySlashCommandsMixin(
|
||||
return await self._handle_voice_channel_leave(event)
|
||||
if args == "status":
|
||||
mode = self._voice_mode.get(voice_key, "off")
|
||||
labels = {
|
||||
"off": t("gateway.voice.label_off"),
|
||||
"voice_only": t("gateway.voice.label_voice_only"),
|
||||
"all": t("gateway.voice.label_all"),
|
||||
}
|
||||
mode_line = t("gateway.voice.status_mode", label=labels.get(mode, mode))
|
||||
label = t(f"gateway.voice.label_{mode}") if mode in ("off", "voice_only", "all") else mode
|
||||
mode_line = t("gateway.voice.status_mode", label=label)
|
||||
# Append voice channel info if connected
|
||||
guild_id = self._get_guild_id(event)
|
||||
info = adapter.get_voice_channel_info(guild_id) if guild_id and hasattr(adapter, "get_voice_channel_info") else None
|
||||
@@ -803,7 +777,6 @@ class GatewaySlashCommandsMixin(
|
||||
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")
|
||||
@@ -864,7 +837,6 @@ class GatewaySlashCommandsMixin(
|
||||
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"))
|
||||
@@ -992,7 +964,6 @@ class GatewaySlashCommandsMixin(
|
||||
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
|
||||
|
||||
args = event.get_command_args().strip().split()
|
||||
# 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).
|
||||
@@ -1013,7 +984,6 @@ class GatewaySlashCommandsMixin(
|
||||
"""
|
||||
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 ""
|
||||
if not wa.write_approval_enabled(wa.SKILLS) and sub not in {"approval", "mode"} and wa.pending_count(wa.SKILLS) == 0:
|
||||
@@ -1040,7 +1010,6 @@ class GatewaySlashCommandsMixin(
|
||||
"""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.
|
||||
@@ -1055,7 +1024,6 @@ class GatewaySlashCommandsMixin(
|
||||
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)
|
||||
@@ -1070,7 +1038,6 @@ class GatewaySlashCommandsMixin(
|
||||
→ log per *current platform*, saved to ``display.platforms.<platform>.tool_progress``.
|
||||
"""
|
||||
from gateway.run import _gateway_config_home, _load_gateway_config, _platform_config_key
|
||||
|
||||
config_path = _gateway_config_home() / "config.yaml"
|
||||
platform_key = _platform_config_key(event.source.platform)
|
||||
|
||||
@@ -1084,7 +1051,6 @@ class GatewaySlashCommandsMixin(
|
||||
|
||||
# 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")
|
||||
if current not in cycle:
|
||||
@@ -1123,7 +1089,6 @@ class GatewaySlashCommandsMixin(
|
||||
profile_name = self._busy_profile_name_for_source(event.source)
|
||||
if profile_name:
|
||||
from gateway.run import _load_gateway_runtime_config
|
||||
|
||||
self._snapshot_profile_busy_modes(profile_name, _load_gateway_runtime_config())
|
||||
else:
|
||||
self._busy_input_mode = arg
|
||||
@@ -1141,7 +1106,6 @@ class GatewaySlashCommandsMixin(
|
||||
"""Handle /footer command — toggle the runtime-metadata footer."""
|
||||
from gateway.run import _gateway_config_home, _load_gateway_config, _platform_config_key, _resolve_gateway_model
|
||||
from gateway.runtime_footer import resolve_footer_config
|
||||
|
||||
config_path = _gateway_config_home() / "config.yaml"
|
||||
platform_key = _platform_config_key(event.source.platform)
|
||||
|
||||
@@ -1150,8 +1114,7 @@ class GatewaySlashCommandsMixin(
|
||||
text = (getattr(event, "message", None) or "").strip()
|
||||
if text.startswith("/"):
|
||||
parts = text.split(None, 1)
|
||||
if len(parts) > 1:
|
||||
arg = parts[1].strip().lower()
|
||||
arg = parts[1].strip().lower() if len(parts) > 1 else ""
|
||||
except Exception:
|
||||
arg = ""
|
||||
|
||||
@@ -1264,10 +1227,8 @@ class GatewaySlashCommandsMixin(
|
||||
# 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)
|
||||
if not callable(refresh):
|
||||
continue
|
||||
try:
|
||||
maybe = refresh()
|
||||
maybe = refresh() if callable(refresh) else None
|
||||
if inspect.isawaitable(maybe):
|
||||
await maybe
|
||||
except Exception as exc:
|
||||
@@ -1305,7 +1266,6 @@ class GatewaySlashCommandsMixin(
|
||||
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:
|
||||
@@ -1346,7 +1306,6 @@ class GatewaySlashCommandsMixin(
|
||||
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
|
||||
@@ -1362,7 +1321,6 @@ class GatewaySlashCommandsMixin(
|
||||
command executes inline — same flow as the CLI's synchronous approval.
|
||||
"""
|
||||
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
|
||||
@@ -1389,7 +1347,6 @@ class GatewaySlashCommandsMixin(
|
||||
as in the CLI. ``/deny`` denies the oldest; ``/deny all`` denies everything.
|
||||
"""
|
||||
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
|
||||
@@ -1485,16 +1442,14 @@ class GatewaySlashCommandsMixin(
|
||||
pending_path = _hermes_home / ".update_pending.json"
|
||||
output_path = _hermes_home / ".update_output.txt"
|
||||
exit_code_path = _hermes_home / ".update_exit_code"
|
||||
src = event.source
|
||||
pending = {
|
||||
"platform": event.source.platform.value,
|
||||
"chat_id": event.source.chat_id,
|
||||
"chat_type": event.source.chat_type,
|
||||
"user_id": event.source.user_id,
|
||||
"session_key": self._session_key_for_source(event.source),
|
||||
"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(),
|
||||
}
|
||||
if event.source.thread_id:
|
||||
pending["thread_id"] = event.source.thread_id
|
||||
if src.thread_id:
|
||||
pending["thread_id"] = src.thread_id
|
||||
if event.message_id:
|
||||
pending["message_id"] = event.message_id
|
||||
_tmp_pending = pending_path.with_suffix(".tmp")
|
||||
|
||||
@@ -184,7 +184,6 @@ class GatewayGoalCommandsMixin:
|
||||
# Inline `field: value` lines parse into a completion contract; the remaining prose is
|
||||
# the goal headline. Plain free-form goals (no such lines) behave exactly as before.
|
||||
from hermes_cli.goals import parse_contract
|
||||
|
||||
headline, parsed = parse_contract(args)
|
||||
args = headline or args
|
||||
contract = parsed if not parsed.is_empty() else None
|
||||
@@ -211,7 +210,6 @@ class GatewayGoalCommandsMixin:
|
||||
prompt. The gateway-wide poller injects due heartbeats through the adapter FIFO as
|
||||
ordinary user turns, so alternation and caching hold."""
|
||||
from hermes_cli.heartbeat import parse_interval, format_interval, MIN_INTERVAL_SECONDS
|
||||
|
||||
args = (event.get_command_args() or "").strip()
|
||||
lower = args.lower()
|
||||
|
||||
@@ -338,7 +336,6 @@ class GatewayGoalCommandsMixin:
|
||||
token = set_current_session_key(quick_key)
|
||||
try:
|
||||
from agent.review_engine import start_review
|
||||
|
||||
return start_review(agent, snapshot, args)
|
||||
finally:
|
||||
reset_current_session_key(token)
|
||||
@@ -353,7 +350,6 @@ class GatewayGoalCommandsMixin:
|
||||
return f"/review failed to start: {exc}"
|
||||
|
||||
from agent.review_engine import format_dispatch_note
|
||||
|
||||
return format_dispatch_note(result, args)
|
||||
|
||||
async def _handle_subgoal_command(self, event: MessageEvent) -> str:
|
||||
@@ -429,14 +425,9 @@ class GatewayGoalCommandsMixin:
|
||||
src = event.source
|
||||
if src is not None:
|
||||
platform = getattr(src, "platform", "")
|
||||
route = {
|
||||
"platform": platform.value if hasattr(platform, "value") else str(platform or ""),
|
||||
"chat_id": str(getattr(src, "chat_id", "") or ""),
|
||||
"chat_type": str(getattr(src, "chat_type", "") or ""),
|
||||
"thread_id": str(getattr(src, "thread_id", "") or ""),
|
||||
"user_id": str(getattr(src, "user_id", "") or ""),
|
||||
"user_name": str(getattr(src, "user_name", "") or ""),
|
||||
}
|
||||
route = {"platform": platform.value if hasattr(platform, "value") else str(platform or "")}
|
||||
for key in ("chat_id", "chat_type", "thread_id", "user_id", "user_name"):
|
||||
route[key] = str(getattr(src, key, "") or "")
|
||||
route = {k: v for k, v in route.items() if v}
|
||||
except Exception:
|
||||
route = {}
|
||||
|
||||
@@ -58,6 +58,21 @@ def _chat_msgs(history) -> list[dict]:
|
||||
return [m for m in history if m.get("role") in {"user", "assistant"} and m.get("content")]
|
||||
|
||||
|
||||
async def _quiet(call, default=None):
|
||||
"""Await ``call()`` fail-open: any exception (sync or in the awaitable) yields *default*."""
|
||||
try:
|
||||
return await call()
|
||||
except Exception:
|
||||
return default
|
||||
|
||||
|
||||
def _quiet_sync(call, default=None):
|
||||
try:
|
||||
return call()
|
||||
except Exception:
|
||||
return default
|
||||
|
||||
|
||||
def _status_model_route(status_agent, persisted_route: dict, session_row: dict, session_entry):
|
||||
"""``(model, provider, context_used, context_total)`` for /status.
|
||||
|
||||
@@ -65,7 +80,6 @@ def _status_model_route(status_agent, persisted_route: dict, session_row: dict,
|
||||
(only loaded when something is still missing).
|
||||
"""
|
||||
from gateway.run import _AGENT_PENDING_SENTINEL, _load_gateway_config, _resolve_gateway_model
|
||||
|
||||
model_name = provider_name = ""
|
||||
context_used = context_total = 0
|
||||
if status_agent is not None and status_agent is not _AGENT_PENDING_SENTINEL:
|
||||
@@ -204,7 +218,6 @@ class GatewayStatusCommandsMixin:
|
||||
async def _handle_status_command(self, event: MessageEvent) -> str:
|
||||
"""Handle /status command."""
|
||||
from gateway.run import _AGENT_PENDING_SENTINEL
|
||||
|
||||
source = event.source
|
||||
session_entry = await self.async_session_store.get_or_create_session(source)
|
||||
session_key = session_entry.session_key
|
||||
@@ -272,33 +285,18 @@ class GatewayStatusCommandsMixin:
|
||||
agent's per-turn token deltas are persisted into sessions_db (run_agent.py), not into
|
||||
SessionEntry, so session_entry.total_tokens is always 0.
|
||||
"""
|
||||
title = None
|
||||
session_row: dict[str, Any] = {}
|
||||
db_total_tokens = 0
|
||||
persisted_route: dict[str, Any] = {}
|
||||
if not self._session_db:
|
||||
return title, session_row, db_total_tokens, persisted_route
|
||||
try:
|
||||
title = await self._session_db.get_session_title(session_id)
|
||||
except Exception:
|
||||
title = None
|
||||
try:
|
||||
row = await self._session_db.get_session(session_id)
|
||||
if isinstance(row, dict):
|
||||
session_row = row
|
||||
db_total_tokens = sum(
|
||||
_int_value(row.get(k))
|
||||
for k in ("input_tokens", "output_tokens", "cache_read_tokens", "cache_write_tokens", "reasoning_tokens")
|
||||
)
|
||||
except Exception:
|
||||
db_total_tokens = 0
|
||||
try:
|
||||
route = await self._session_db.get_dominant_session_model_route(session_id)
|
||||
if isinstance(route, dict):
|
||||
persisted_route = route
|
||||
except Exception:
|
||||
persisted_route = {}
|
||||
return title, session_row, db_total_tokens, persisted_route
|
||||
db = self._session_db
|
||||
if not db:
|
||||
return None, {}, 0, {}
|
||||
title = await _quiet(lambda: db.get_session_title(session_id))
|
||||
row = await _quiet(lambda: db.get_session(session_id))
|
||||
session_row = row if isinstance(row, dict) else {}
|
||||
db_total_tokens = sum(
|
||||
_int_value(session_row.get(k))
|
||||
for k in ("input_tokens", "output_tokens", "cache_read_tokens", "cache_write_tokens", "reasoning_tokens")
|
||||
)
|
||||
route = await _quiet(lambda: db.get_dominant_session_model_route(session_id))
|
||||
return title, session_row, db_total_tokens, route if isinstance(route, dict) else {}
|
||||
|
||||
@staticmethod
|
||||
def _redact_matrix_session_key(session_key: str) -> str:
|
||||
@@ -357,7 +355,6 @@ class GatewayStatusCommandsMixin:
|
||||
history = await self.async_session_store.load_transcript(session_entry.session_id)
|
||||
if history:
|
||||
from agent.model_metadata import estimate_messages_tokens_rough
|
||||
|
||||
msgs = _chat_msgs(history)
|
||||
return "\n".join([
|
||||
t("gateway.context.header"),
|
||||
@@ -380,41 +377,30 @@ class GatewayStatusCommandsMixin:
|
||||
if not used:
|
||||
used = _int_value(getattr(session_entry, "last_prompt_tokens", 0))
|
||||
if not model_name and self._session_db:
|
||||
try:
|
||||
row = await self._session_db.get_session(session_entry.session_id) or {}
|
||||
if isinstance(row, dict):
|
||||
model_name = _clean_str(row.get("model", ""))
|
||||
except Exception:
|
||||
model_name = ""
|
||||
row = await _quiet(lambda: self._session_db.get_session(session_entry.session_id))
|
||||
model_name = _clean_str(row.get("model", "")) if isinstance(row, dict) else ""
|
||||
if not context_length:
|
||||
try:
|
||||
from gateway.run import _profile_runtime_scope, _resolve_gateway_model_context
|
||||
from gateway.run import _profile_runtime_scope, _resolve_gateway_model_context
|
||||
|
||||
def _resolve_nonresident_context():
|
||||
if getattr(getattr(self, "config", None), "multiplex_profiles", False):
|
||||
with _profile_runtime_scope(self._resolve_profile_home_for_source(source)):
|
||||
return _resolve_gateway_model_context(model_name or None)
|
||||
return _resolve_gateway_model_context(model_name or None)
|
||||
def _resolve_nonresident_context():
|
||||
if getattr(getattr(self, "config", None), "multiplex_profiles", False):
|
||||
with _profile_runtime_scope(self._resolve_profile_home_for_source(source)):
|
||||
return _resolve_gateway_model_context(model_name or None)
|
||||
return _resolve_gateway_model_context(model_name or None)
|
||||
|
||||
resolved = await asyncio.to_thread(_resolve_nonresident_context)
|
||||
resolved = await _quiet(lambda: asyncio.to_thread(_resolve_nonresident_context))
|
||||
if resolved is not None:
|
||||
model_name = model_name or resolved.model
|
||||
context_length = _int_value(resolved.context_length)
|
||||
except Exception:
|
||||
context_length = 0
|
||||
if not context_length and model_name:
|
||||
try:
|
||||
from agent.model_metadata import get_model_context_length
|
||||
|
||||
context_length = _int_value(await asyncio.to_thread(get_model_context_length, model_name))
|
||||
except Exception:
|
||||
context_length = 0
|
||||
from agent.model_metadata import get_model_context_length
|
||||
context_length = _int_value(await _quiet(lambda: asyncio.to_thread(get_model_context_length, model_name)))
|
||||
return used, context_length, model_name
|
||||
|
||||
async def _handle_agents_command(self, event: MessageEvent) -> str:
|
||||
"""Handle /agents command - list active agents and running tasks."""
|
||||
from gateway.run import _AGENT_PENDING_SENTINEL
|
||||
from tools.process_registry import format_uptime_short, process_registry
|
||||
|
||||
now = time.time()
|
||||
current_session_key = self._session_key_for_source(event.source)
|
||||
running_started: dict = getattr(self, "_running_agents_ts", {}) or {}
|
||||
@@ -431,10 +417,8 @@ class GatewayStatusCommandsMixin:
|
||||
})
|
||||
agent_rows.sort(key=lambda row: row["elapsed"], reverse=True)
|
||||
|
||||
try:
|
||||
running_processes = [p for p in process_registry.list_sessions() if p.get("status") == "running"]
|
||||
except Exception:
|
||||
running_processes = []
|
||||
procs = _quiet_sync(process_registry.list_sessions, [])
|
||||
running_processes = [p for p in procs if p.get("status") == "running"]
|
||||
|
||||
background_tasks = [
|
||||
task for task in (getattr(self, "_background_tasks", set()) or set())
|
||||
@@ -442,14 +426,11 @@ class GatewayStatusCommandsMixin:
|
||||
]
|
||||
|
||||
# Background (async) delegations — delegate_task(background=true).
|
||||
try:
|
||||
from tools.async_delegation import list_async_delegations
|
||||
delegations = [
|
||||
d for d in list_async_delegations()
|
||||
if d.get("status") in ("running", "stalling", "finalizing")
|
||||
]
|
||||
except Exception:
|
||||
delegations = []
|
||||
from tools.async_delegation import list_async_delegations
|
||||
delegations = [
|
||||
d for d in _quiet_sync(list_async_delegations, [])
|
||||
if d.get("status") in ("running", "stalling", "finalizing")
|
||||
]
|
||||
|
||||
def _agent_row(idx_row):
|
||||
idx, row = idx_row
|
||||
@@ -483,12 +464,7 @@ class GatewayStatusCommandsMixin:
|
||||
shows the new balance. Fetched off the event loop; fail-open.
|
||||
"""
|
||||
from agent.account_usage import build_credits_view
|
||||
|
||||
try:
|
||||
view = await asyncio.to_thread(build_credits_view, markdown=True)
|
||||
except Exception:
|
||||
view = None
|
||||
|
||||
view = await _quiet(lambda: asyncio.to_thread(build_credits_view, markdown=True))
|
||||
if view is None or not view.logged_in:
|
||||
return t("gateway.credits.not_logged_in")
|
||||
|
||||
@@ -511,16 +487,12 @@ class GatewayStatusCommandsMixin:
|
||||
"""
|
||||
try:
|
||||
from agent.context_breakdown import compute_context_details, render_context_breakdown_lines
|
||||
|
||||
payload = self._session_context_breakdown(agent, source)
|
||||
if not (payload.get("categories") or []):
|
||||
return []
|
||||
details = None
|
||||
if expanded:
|
||||
try:
|
||||
details = compute_context_details(agent)
|
||||
except Exception:
|
||||
details = {"skills": [], "toolsets": []}
|
||||
details = _quiet_sync(lambda: compute_context_details(agent), {"skills": [], "toolsets": []})
|
||||
return render_context_breakdown_lines(payload, details=details, grid=False)
|
||||
except Exception:
|
||||
return []
|
||||
@@ -529,12 +501,11 @@ class GatewayStatusCommandsMixin:
|
||||
"""Per-category context estimate (chars/4) for *agent* over the session transcript (sync)."""
|
||||
from agent.context_breakdown import compute_session_context_breakdown
|
||||
|
||||
try:
|
||||
def _history():
|
||||
entry = self.session_store.get_or_create_session(source)
|
||||
history = self.session_store.load_transcript(entry.session_id) or []
|
||||
except Exception:
|
||||
history = []
|
||||
return compute_session_context_breakdown(agent, history)
|
||||
return self.session_store.load_transcript(entry.session_id) or []
|
||||
|
||||
return compute_session_context_breakdown(agent, _quiet_sync(_history, []))
|
||||
|
||||
def _context_breakdown_lines(self, agent, source) -> list[str]:
|
||||
"""Render the per-category context breakdown for /usage.
|
||||
@@ -590,7 +561,6 @@ class GatewayStatusCommandsMixin:
|
||||
if str(provider or "").strip().lower() != "openai-codex":
|
||||
return t("gateway.usage.reset_wrong_provider")
|
||||
from agent.account_usage import redeem_codex_reset_credit
|
||||
|
||||
result = await asyncio.to_thread(
|
||||
redeem_codex_reset_credit, base_url=base_url, api_key=api_key, force="--force" in args[1:],
|
||||
)
|
||||
@@ -600,24 +570,17 @@ class GatewayStatusCommandsMixin:
|
||||
# failures are non-fatal (account_lines stays []).
|
||||
account_lines: list[str] = []
|
||||
if provider:
|
||||
try:
|
||||
account_snapshot = await asyncio.to_thread(
|
||||
fetch_account_usage, provider, base_url=base_url, api_key=api_key,
|
||||
)
|
||||
except Exception:
|
||||
account_snapshot = None
|
||||
account_snapshot = await _quiet(
|
||||
lambda: asyncio.to_thread(fetch_account_usage, provider, base_url=base_url, api_key=api_key)
|
||||
)
|
||||
if account_snapshot:
|
||||
account_lines = render_account_usage_lines(account_snapshot, markdown=True)
|
||||
|
||||
# Nous credits + monthly-grant gauge (shared with CLI/TUI). Gates on "a Nous account is
|
||||
# logged in" — NOT the inference provider — so a Nous user inferring elsewhere still sees
|
||||
# a balance. Fail-open: never break /usage.
|
||||
try:
|
||||
from agent.account_usage import nous_credits_lines
|
||||
|
||||
credits_lines = await asyncio.to_thread(nous_credits_lines, markdown=True)
|
||||
except Exception:
|
||||
credits_lines = []
|
||||
from agent.account_usage import nous_credits_lines
|
||||
credits_lines = await _quiet(lambda: asyncio.to_thread(nous_credits_lines, markdown=True), [])
|
||||
|
||||
def _with_account_blocks(lines: list[str]) -> str:
|
||||
# Each block is preceded by a blank divider only when something precedes it.
|
||||
|
||||
Reference in New Issue
Block a user