refactor(web): absorb single-router web_server helpers into their routers (75 defs); drop dead imports
This commit is contained in:
@@ -15,6 +15,9 @@ from fastapi import HTTPException, Request
|
||||
from hermes_cli import __version__
|
||||
from hermes_cli.config import format_docker_update_message, recommended_update_command_for_method
|
||||
from typing import Any, Dict, List, Optional
|
||||
from pathlib import Path
|
||||
import re
|
||||
import time
|
||||
|
||||
_log = logging.getLogger("hermes_cli.web_server")
|
||||
router = APIRouter()
|
||||
@@ -22,15 +25,127 @@ status_router = APIRouter()
|
||||
|
||||
# web_server helpers, late-bound so monkeypatch.setattr(web_server, ...) stays authoritative.
|
||||
_dashboard_local_update_managed_externally = late("_dashboard_local_update_managed_externally")
|
||||
_durable_completed_update_action_id = late("_durable_completed_update_action_id")
|
||||
_record_completed_action = late("_record_completed_action")
|
||||
_spawn_gateway_restart = late("_spawn_gateway_restart")
|
||||
_spawn_hermes_action = late("_spawn_hermes_action")
|
||||
_tail_lines = late("_tail_lines")
|
||||
detect_install_method = late("detect_install_method")
|
||||
get_hermes_home = late("get_hermes_home")
|
||||
|
||||
|
||||
_ACTION_LOG_TAIL_MAX_BYTES = 256 * 1024
|
||||
|
||||
|
||||
_ACTION_LOG_TAIL_INITIAL_CHUNK_BYTES = 8 * 1024
|
||||
|
||||
|
||||
_ACTION_LOG_TAIL_MAX_CHUNK_BYTES = 64 * 1024
|
||||
|
||||
|
||||
_UPDATE_ACTION_COMPLETED_RE = re.compile(
|
||||
r"^=== hermes-update completed ([0-9a-f]{32}) ===$"
|
||||
)
|
||||
|
||||
|
||||
def _record_completed_action(name: str, message: str, exit_code: int = 1) -> None:
|
||||
"""Record a non-spawned action result and write it to the action log."""
|
||||
from hermes_cli.web_server import (
|
||||
_ACTION_COMMANDS,
|
||||
_ACTION_IDS,
|
||||
_ACTION_LOG_DIR,
|
||||
_ACTION_LOG_FILES,
|
||||
_ACTION_PROCS,
|
||||
_ACTION_RESULTS,
|
||||
)
|
||||
log_file_name = _ACTION_LOG_FILES[name]
|
||||
_ACTION_LOG_DIR.mkdir(parents=True, exist_ok=True)
|
||||
log_path = _ACTION_LOG_DIR / log_file_name
|
||||
with open(log_path, "ab", buffering=0) as log_file:
|
||||
log_file.write(
|
||||
f"\n=== {name} completed {time.strftime('%Y-%m-%d %H:%M:%S')} ===\n".encode()
|
||||
)
|
||||
log_file.write(message.encode("utf-8", errors="replace"))
|
||||
if not message.endswith("\n"):
|
||||
log_file.write(b"\n")
|
||||
_ACTION_PROCS.pop(name, None)
|
||||
_ACTION_COMMANDS.pop(name, None)
|
||||
_ACTION_IDS.pop(name, None)
|
||||
_ACTION_RESULTS[name] = {"exit_code": exit_code, "pid": None}
|
||||
|
||||
|
||||
def _tail_lines(path: Path, n: int) -> List[str]:
|
||||
"""Return the last ``n`` lines of ``path`` without loading huge logs."""
|
||||
if n <= 0 or not path.exists():
|
||||
return []
|
||||
try:
|
||||
size = path.stat().st_size
|
||||
except OSError:
|
||||
return []
|
||||
if size <= 0:
|
||||
return []
|
||||
|
||||
min_offset = max(0, size - _ACTION_LOG_TAIL_MAX_BYTES)
|
||||
offset = size
|
||||
chunk_size = _ACTION_LOG_TAIL_INITIAL_CHUNK_BYTES
|
||||
newline_count = 0
|
||||
chunks: List[bytes] = []
|
||||
drop_partial_first_line = False
|
||||
|
||||
try:
|
||||
with path.open("rb") as handle:
|
||||
while offset > min_offset and newline_count <= n:
|
||||
read_size = min(chunk_size, offset - min_offset)
|
||||
offset -= read_size
|
||||
handle.seek(offset)
|
||||
chunk = handle.read(read_size)
|
||||
chunks.append(chunk)
|
||||
newline_count += chunk.count(b"\n")
|
||||
chunk_size = min(
|
||||
chunk_size * 2,
|
||||
_ACTION_LOG_TAIL_MAX_CHUNK_BYTES,
|
||||
)
|
||||
if offset > 0:
|
||||
handle.seek(offset - 1)
|
||||
drop_partial_first_line = handle.read(1) != b"\n"
|
||||
except OSError:
|
||||
return []
|
||||
|
||||
lines = (
|
||||
b"".join(reversed(chunks))
|
||||
.decode("utf-8", errors="replace")
|
||||
.splitlines()
|
||||
)
|
||||
if drop_partial_first_line and lines:
|
||||
lines = lines[1:]
|
||||
return lines[-n:]
|
||||
|
||||
|
||||
def _durable_completed_update_action_id(lines: List[str]) -> Optional[str]:
|
||||
"""Recover the latest successful update identity from ``update.log``.
|
||||
|
||||
The dashboard action process can restart the dashboard that spawned it.
|
||||
That loses the in-memory ``Popen``/result registries while the durable
|
||||
update log survives. Only accept a completion marker that occurs after
|
||||
the latest update-start marker, so a stale success cannot mask a newer
|
||||
failed attempt.
|
||||
"""
|
||||
last_start = -1
|
||||
last_completed = -1
|
||||
completed_action_id: Optional[str] = None
|
||||
|
||||
for index, line in enumerate(lines):
|
||||
if line.startswith("=== hermes update started "):
|
||||
last_start = index
|
||||
|
||||
match = _UPDATE_ACTION_COMPLETED_RE.fullmatch(line.strip())
|
||||
if match:
|
||||
last_completed = index
|
||||
completed_action_id = match.group(1)
|
||||
|
||||
if completed_action_id and last_completed > last_start:
|
||||
return completed_action_id
|
||||
|
||||
return None
|
||||
|
||||
|
||||
@router.post("/api/gateway/restart")
|
||||
async def restart_gateway(profile: Optional[str] = None):
|
||||
"""Kick off a ``hermes gateway restart`` in the background."""
|
||||
|
||||
@@ -27,7 +27,6 @@ _log = logging.getLogger("hermes_cli.web_server")
|
||||
router = APIRouter()
|
||||
|
||||
# web_server helpers, late-bound so monkeypatch.setattr(web_server, ...) stays authoritative.
|
||||
_audio_extension_for_mime = late("_audio_extension_for_mime")
|
||||
_config_profile_scope = late("_config_profile_scope")
|
||||
_split_text_for_speak_stream = late("_split_text_for_speak_stream")
|
||||
_voice_list_error_logged_once = late("_voice_list_error_logged_once")
|
||||
@@ -36,11 +35,35 @@ _ws_request_is_allowed = late("_ws_request_is_allowed")
|
||||
load_env = late("load_env")
|
||||
|
||||
|
||||
_AUDIO_MIME_EXTENSIONS: Dict[str, str] = {
|
||||
"audio/aac": ".aac",
|
||||
"audio/flac": ".flac",
|
||||
"audio/m4a": ".m4a",
|
||||
"audio/mp3": ".mp3",
|
||||
"audio/mp4": ".mp4",
|
||||
"audio/mpeg": ".mp3",
|
||||
"audio/ogg": ".ogg",
|
||||
"audio/wav": ".wav",
|
||||
"audio/wave": ".wav",
|
||||
"audio/webm": ".webm",
|
||||
"audio/x-m4a": ".m4a",
|
||||
"audio/x-wav": ".wav",
|
||||
"video/webm": ".webm",
|
||||
}
|
||||
|
||||
|
||||
_MAX_TRANSCRIPTION_UPLOAD_BYTES = 25 * 1024 * 1024
|
||||
|
||||
|
||||
def _audio_extension_for_mime(mime_type: str) -> str:
|
||||
normalized = (mime_type or "").split(";", 1)[0].strip().lower()
|
||||
return _AUDIO_MIME_EXTENSIONS.get(normalized, ".webm")
|
||||
|
||||
|
||||
@router.post("/api/audio/transcribe")
|
||||
async def transcribe_audio_upload(
|
||||
payload: AudioTranscriptionRequest, profile: Optional[str] = None
|
||||
):
|
||||
from hermes_cli.web_server import _MAX_TRANSCRIPTION_UPLOAD_BYTES
|
||||
data_url = (payload.data_url or "").strip()
|
||||
if not data_url.startswith("data:") or "," not in data_url:
|
||||
raise HTTPException(status_code=400, detail="Invalid audio payload")
|
||||
|
||||
@@ -14,32 +14,110 @@ from fastapi import HTTPException, WebSocket, WebSocketDisconnect
|
||||
from hermes_cli.pty_session import RegistryFull
|
||||
from pathlib import Path
|
||||
from typing import Any, Dict, Optional
|
||||
import re
|
||||
from fastapi import FastAPI
|
||||
|
||||
_log = logging.getLogger("hermes_cli.web_server")
|
||||
router = APIRouter()
|
||||
|
||||
# web_server helpers, late-bound so monkeypatch.setattr(web_server, ...) stays authoritative.
|
||||
_active_session_file_for_channel = late("_active_session_file_for_channel")
|
||||
_broadcast_event = late("_broadcast_event")
|
||||
_build_sidecar_url = late("_build_sidecar_url")
|
||||
_channel_or_close_code = late("_channel_or_close_code")
|
||||
_forget_active_session_file = late("_forget_active_session_file")
|
||||
_get_console_executor = late("_get_console_executor")
|
||||
_get_event_state = late("_get_event_state")
|
||||
_legacy_pump = late("_legacy_pump")
|
||||
_profile_scope = late("_profile_scope")
|
||||
_read_active_session_file = late("_read_active_session_file")
|
||||
_resolve_chat_argv_async = late("_resolve_chat_argv_async")
|
||||
_resolve_profile_dir = late("_resolve_profile_dir")
|
||||
_ws_auth_mode = late("_ws_auth_mode")
|
||||
_ws_auth_ok = late("_ws_auth_ok")
|
||||
_ws_auth_reason = late("_ws_auth_reason")
|
||||
_ws_client_reason = late("_ws_client_reason")
|
||||
_ws_close_reason = late("_ws_close_reason")
|
||||
_ws_host_origin_reason = late("_ws_host_origin_reason")
|
||||
_ws_request_is_allowed = late("_ws_request_is_allowed")
|
||||
|
||||
|
||||
def _get_event_state(app: "FastAPI"):
|
||||
"""Return (event_channels, event_lock) from app.state.
|
||||
|
||||
Lazily initialises the state if the lifespan hasn't run (e.g. when
|
||||
TestClient is constructed without a ``with`` block). The lifespan
|
||||
path is preferred because it guarantees the Lock is created on the
|
||||
correct event loop, but the lazy path lets existing non-``with``
|
||||
TestClient usages keep working.
|
||||
"""
|
||||
try:
|
||||
return app.state.event_channels, app.state.event_lock
|
||||
except AttributeError:
|
||||
app.state.event_channels = {}
|
||||
app.state.event_lock = asyncio.Lock()
|
||||
return app.state.event_channels, app.state.event_lock
|
||||
|
||||
|
||||
_VALID_CHANNEL_RE = re.compile(r"^[A-Za-z0-9._-]{1,128}$")
|
||||
|
||||
|
||||
def _ws_auth_mode() -> str:
|
||||
"""Short label for the active WS auth mode — logged on every connection."""
|
||||
from hermes_cli.web_server import _LOOPBACK_HOSTS, app
|
||||
if getattr(app.state, "auth_required", False):
|
||||
return "gated"
|
||||
bound_host = (getattr(app.state, "bound_host", "") or "").strip().lower()
|
||||
if bound_host and bound_host not in _LOOPBACK_HOSTS:
|
||||
return "insecure"
|
||||
return "loopback"
|
||||
|
||||
|
||||
async def _broadcast_event(app: Any, channel: str, payload: str) -> None:
|
||||
"""Fan out one publisher frame to every subscriber on `channel`."""
|
||||
event_channels, event_lock = _get_event_state(app)
|
||||
async with event_lock:
|
||||
subs = list(event_channels.get(channel, ()))
|
||||
|
||||
for sub in subs:
|
||||
try:
|
||||
await sub.send_text(payload)
|
||||
except Exception:
|
||||
# Subscriber went away mid-send; the /api/events finally clause
|
||||
# will remove it from the registry on its next iteration.
|
||||
_log.warning("broadcast send failed for subscriber on %s", channel, exc_info=True)
|
||||
|
||||
|
||||
def _channel_or_close_code(ws: WebSocket) -> Optional[str]:
|
||||
"""Return the channel id from the query string or None if invalid."""
|
||||
channel = ws.query_params.get("channel", "")
|
||||
|
||||
return channel if _VALID_CHANNEL_RE.match(channel) else None
|
||||
|
||||
|
||||
def _read_active_session_file(path: Path) -> Optional[str]:
|
||||
try:
|
||||
data = json.loads(path.read_text(encoding="utf-8"))
|
||||
except (OSError, json.JSONDecodeError):
|
||||
return None
|
||||
|
||||
session_id = str(data.get("session_id") or "").strip()
|
||||
return session_id or None
|
||||
|
||||
|
||||
def _forget_active_session_file(path: Path) -> None:
|
||||
try:
|
||||
path.unlink(missing_ok=True)
|
||||
except OSError:
|
||||
pass
|
||||
|
||||
|
||||
def _ws_close_reason(text: str) -> str:
|
||||
"""Clamp a WS close reason to the protocol's 123-byte UTF-8 limit.
|
||||
|
||||
RFC 6455 caps the close-frame reason at 123 bytes; uvicorn raises if a
|
||||
longer string is passed. Our reasons embed an attacker-controlled origin,
|
||||
so truncate defensively rather than crash the close handler.
|
||||
"""
|
||||
encoded = text.encode("utf-8", "replace")
|
||||
if len(encoded) <= 123:
|
||||
return text
|
||||
return encoded[:120].decode("utf-8", "ignore") + "..."
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# /api/console — safe Hermes Console command WebSocket.
|
||||
#
|
||||
@@ -231,7 +309,7 @@ def _console_json_payload(msg: Any) -> tuple[Optional[dict[str, Any]], Optional[
|
||||
|
||||
@router.websocket("/api/console")
|
||||
async def console_ws(ws: WebSocket) -> None:
|
||||
from hermes_cli.web_server import _DASHBOARD_EMBEDDED_CHAT_ENABLED, asyncio
|
||||
from hermes_cli.web_server import _DASHBOARD_EMBEDDED_CHAT_ENABLED
|
||||
peer = ws.client.host if ws.client else "?"
|
||||
|
||||
if not _DASHBOARD_EMBEDDED_CHAT_ENABLED:
|
||||
|
||||
@@ -41,6 +41,24 @@ save_config = late("save_config")
|
||||
save_env_value = late("save_env_value")
|
||||
|
||||
|
||||
# Simple rate limiter for the reveal endpoint
|
||||
_reveal_timestamps: List[float] = []
|
||||
|
||||
|
||||
_REVEAL_MAX_PER_WINDOW = 5
|
||||
|
||||
|
||||
_REVEAL_WINDOW_SECONDS = 30
|
||||
|
||||
|
||||
# Display order for tabs — unlisted categories sort alphabetically after these.
|
||||
_CATEGORY_ORDER = [
|
||||
"general", "agent", "terminal", "display", "delegation",
|
||||
"memory", "compression", "security", "browser", "voice",
|
||||
"tts", "stt", "logging", "discord", "auxiliary",
|
||||
]
|
||||
|
||||
|
||||
@config_router.get("/api/config")
|
||||
async def get_config(profile: Optional[str] = None):
|
||||
# _profile_scope blocks on the process-wide _SKILLS_PROFILE_LOCK and
|
||||
@@ -67,7 +85,6 @@ async def get_schema(profile: Optional[str] = None):
|
||||
# Discovery-driven provider options (voice command providers + memory
|
||||
# provider plugins) are merged per-request so providers added after server
|
||||
# start still show up, scoped to the requested profile's config.
|
||||
from hermes_cli.web_server import _CATEGORY_ORDER
|
||||
with _config_profile_scope(profile):
|
||||
fields = _schema_with_dynamic_provider_options()
|
||||
return {"fields": fields, "category_order": _CATEGORY_ORDER}
|
||||
@@ -800,11 +817,6 @@ async def reveal_env_var(
|
||||
- Rate limiting (max 5 reveals per 30s window)
|
||||
- Audit logging
|
||||
"""
|
||||
from hermes_cli.web_server import (
|
||||
_REVEAL_MAX_PER_WINDOW,
|
||||
_REVEAL_WINDOW_SECONDS,
|
||||
_reveal_timestamps,
|
||||
)
|
||||
# --- Token check ---
|
||||
_require_token(request)
|
||||
|
||||
|
||||
@@ -7,6 +7,7 @@ late-binding seam so ``monkeypatch.setattr(web_server, ...)`` keeps working.
|
||||
|
||||
import asyncio
|
||||
import functools
|
||||
import time
|
||||
from typing import Optional
|
||||
|
||||
from fastapi import APIRouter, HTTPException, Request
|
||||
@@ -19,19 +20,13 @@ from hermes_cli.web_models import (
|
||||
AutomationBlueprintInstantiate,
|
||||
)
|
||||
from hermes_cli.web_routers._common import log as _log
|
||||
from typing import Any, Dict, List
|
||||
from pathlib import Path
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
_run_cron_dashboard_io = late("_run_cron_dashboard_io")
|
||||
_list_cron_jobs_sync = late("_list_cron_jobs_sync")
|
||||
_get_cron_job_sync = late("_get_cron_job_sync")
|
||||
_list_cron_job_runs_sync = late("_list_cron_job_runs_sync")
|
||||
_create_cron_job_sync = late("_create_cron_job_sync")
|
||||
_update_cron_job_sync = late("_update_cron_job_sync")
|
||||
_pause_cron_job_sync = late("_pause_cron_job_sync")
|
||||
_resume_cron_job_sync = late("_resume_cron_job_sync")
|
||||
_trigger_cron_job_sync = late("_trigger_cron_job_sync")
|
||||
_delete_cron_job_sync = late("_delete_cron_job_sync")
|
||||
_find_cron_job_profile = late("_find_cron_job_profile")
|
||||
_fire_cron_job_for_profile = late("_fire_cron_job_for_profile")
|
||||
_forward_cron_fire_to_gateway = late("_forward_cron_fire_to_gateway")
|
||||
@@ -41,6 +36,228 @@ _call_cron_for_profile = late("_call_cron_for_profile")
|
||||
_raise_if_cron_registration_error = late("_raise_if_cron_registration_error")
|
||||
load_config = late("load_config")
|
||||
cfg_get = late("cfg_get")
|
||||
_cron_profile_dicts = late("_cron_profile_dicts")
|
||||
_cron_profile_home = late("_cron_profile_home")
|
||||
_mutate_cron_for_profile = late("_mutate_cron_for_profile")
|
||||
_open_session_db_for_profile = late("_open_session_db_for_profile")
|
||||
_validate_dashboard_cron_context_from = late("_validate_dashboard_cron_context_from")
|
||||
_validate_dashboard_cron_effective_job = late("_validate_dashboard_cron_effective_job")
|
||||
_cron_optional_text = late("_cron_optional_text")
|
||||
_cron_string_list = late("_cron_string_list")
|
||||
_normalize_dashboard_cron_script = late("_normalize_dashboard_cron_script")
|
||||
|
||||
|
||||
def _normalize_dashboard_cron_updates(
|
||||
updates: Dict[str, Any],
|
||||
profile_home: Path,
|
||||
) -> Dict[str, Any]:
|
||||
"""Normalize dashboard JSON into cron.jobs.update_job's storage shape.
|
||||
|
||||
This intentionally stays in the dashboard adapter layer: cron/jobs.py is the
|
||||
source of truth for scheduling behaviour; the dashboard only translates form
|
||||
payloads into the shapes that existing core functions already accept.
|
||||
"""
|
||||
normalized = dict(updates or {})
|
||||
|
||||
for key in ("model", "provider", "workdir"):
|
||||
if key in normalized:
|
||||
normalized[key] = _cron_optional_text(normalized[key])
|
||||
if "script" in normalized:
|
||||
normalized["script"] = _normalize_dashboard_cron_script(
|
||||
normalized["script"],
|
||||
profile_home,
|
||||
)
|
||||
if "base_url" in normalized:
|
||||
normalized["base_url"] = _cron_optional_text(
|
||||
normalized["base_url"], strip_trailing_slash=True
|
||||
)
|
||||
if "deliver" in normalized:
|
||||
normalized["deliver"] = _cron_optional_text(normalized["deliver"]) or "local"
|
||||
if "failure_deliver" in normalized:
|
||||
# Same text normalization as deliver, but empty CLEARS the override
|
||||
# (failures fall back to deliver) rather than coalescing to a target
|
||||
# — the field is optional by design (NS-788).
|
||||
normalized["failure_deliver"] = _cron_optional_text(
|
||||
normalized["failure_deliver"]
|
||||
)
|
||||
if "context_from" in normalized:
|
||||
normalized["context_from"] = _cron_string_list(normalized["context_from"])
|
||||
if "enabled_toolsets" in normalized:
|
||||
normalized["enabled_toolsets"] = _cron_string_list(normalized["enabled_toolsets"])
|
||||
return normalized
|
||||
|
||||
|
||||
def _list_cron_jobs_sync(profile: str = "all"):
|
||||
requested = (profile or "all").strip()
|
||||
if requested.lower() != "all":
|
||||
return _call_cron_for_profile(requested, "list_jobs", True)
|
||||
|
||||
jobs: List[Dict[str, Any]] = []
|
||||
for item in _cron_profile_dicts():
|
||||
name = str(item.get("name") or "")
|
||||
if not name:
|
||||
continue
|
||||
try:
|
||||
jobs.extend(_call_cron_for_profile(name, "list_jobs", True))
|
||||
except Exception:
|
||||
_log.exception("Failed to list cron jobs for profile %s", name)
|
||||
return jobs
|
||||
|
||||
|
||||
def _get_cron_job_sync(job_id: str, profile: Optional[str] = None):
|
||||
selected = profile or _find_cron_job_profile(job_id)
|
||||
if not selected:
|
||||
raise HTTPException(status_code=404, detail="Job not found")
|
||||
job = _call_cron_for_profile(selected, "get_job", job_id)
|
||||
if not job:
|
||||
raise HTTPException(status_code=404, detail="Job not found")
|
||||
return job
|
||||
|
||||
|
||||
def _list_cron_job_runs_sync(job_id: str, profile: Optional[str] = None, limit: int = 20):
|
||||
"""Run sessions produced by a cron job, newest first.
|
||||
|
||||
Cron runs are stored as ordinary sessions whose id is
|
||||
``cron_{job_id}_{timestamp}`` (see cron/scheduler.run_job). A job's history
|
||||
is therefore every session whose id carries that prefix; ``source='cron'``
|
||||
narrows it and the id prefix binds it to this job. Powers the run-history
|
||||
list under each job in the desktop cron detail. Same row shape as
|
||||
``/api/sessions`` so the frontend can reuse SessionInfo.
|
||||
|
||||
Backed by ``SessionDB.list_cron_job_runs`` — a bounded ``[prefix, hi)``
|
||||
id-range scan, not the compression-chain CTE used for the recents list,
|
||||
so the cost scales with the requested window and not the (unbounded) total
|
||||
cron history.
|
||||
"""
|
||||
selected = profile or _find_cron_job_profile(job_id)
|
||||
# job_id may be a human name; resolve to the canonical id used in run-session ids.
|
||||
canonical = job_id
|
||||
if selected:
|
||||
job = _call_cron_for_profile(selected, "get_job", job_id)
|
||||
if job and job.get("id"):
|
||||
canonical = str(job["id"])
|
||||
|
||||
try:
|
||||
limit_n = max(1, min(int(limit), 100))
|
||||
except (TypeError, ValueError):
|
||||
limit_n = 20
|
||||
|
||||
db = _open_session_db_for_profile(selected, read_only=True)
|
||||
try:
|
||||
runs = db.list_cron_job_runs(canonical, limit=limit_n, offset=0)
|
||||
now = time.time()
|
||||
for s in runs:
|
||||
s["is_active"] = (
|
||||
s.get("ended_at") is None
|
||||
and (now - s.get("last_active", s.get("started_at", 0))) < 300
|
||||
)
|
||||
s["archived"] = bool(s.get("archived"))
|
||||
if selected:
|
||||
s["profile"] = selected
|
||||
return {"runs": runs, "limit": limit_n}
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
def _update_cron_job_sync(job_id: str, body: CronJobUpdate, profile: Optional[str] = None):
|
||||
selected = profile or _find_cron_job_profile(job_id)
|
||||
if not selected:
|
||||
raise HTTPException(status_code=404, detail="Job not found")
|
||||
try:
|
||||
profile_name, profile_home = _cron_profile_home(selected)
|
||||
existing = _call_cron_for_profile(profile_name, "get_job", job_id)
|
||||
if not existing:
|
||||
raise HTTPException(status_code=404, detail="Job not found")
|
||||
updates = _normalize_dashboard_cron_updates(
|
||||
body.updates,
|
||||
profile_home,
|
||||
)
|
||||
if "context_from" in updates:
|
||||
_validate_dashboard_cron_context_from(
|
||||
updates.get("context_from"),
|
||||
profile_name,
|
||||
)
|
||||
execution_fields = {"prompt", "skill", "skills", "script", "no_agent"}
|
||||
if execution_fields.intersection(updates):
|
||||
effective = {**existing, **updates}
|
||||
if "skills" in updates and "skill" not in updates:
|
||||
effective["skill"] = None
|
||||
_validate_dashboard_cron_effective_job(effective)
|
||||
job = _mutate_cron_for_profile(profile_name, "update_job", job_id, updates)
|
||||
except HTTPException:
|
||||
raise
|
||||
except ValueError as exc:
|
||||
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
||||
if not job:
|
||||
raise HTTPException(status_code=404, detail="Job not found")
|
||||
return job
|
||||
|
||||
|
||||
def _pause_cron_job_sync(job_id: str, profile: Optional[str] = None):
|
||||
selected = profile or _find_cron_job_profile(job_id)
|
||||
if not selected:
|
||||
raise HTTPException(status_code=404, detail="Job not found")
|
||||
job = _mutate_cron_for_profile(selected, "pause_job", job_id)
|
||||
if not job:
|
||||
raise HTTPException(status_code=404, detail="Job not found")
|
||||
return job
|
||||
|
||||
|
||||
def _resume_cron_job_sync(job_id: str, profile: Optional[str] = None):
|
||||
selected = profile or _find_cron_job_profile(job_id)
|
||||
if not selected:
|
||||
raise HTTPException(status_code=404, detail="Job not found")
|
||||
job = _mutate_cron_for_profile(selected, "resume_job", job_id)
|
||||
if not job:
|
||||
raise HTTPException(status_code=404, detail="Job not found")
|
||||
return job
|
||||
|
||||
|
||||
def _trigger_cron_job_sync(job_id: str, profile: Optional[str] = None):
|
||||
selected = profile or _find_cron_job_profile(job_id)
|
||||
if not selected:
|
||||
raise HTTPException(status_code=404, detail="Job not found")
|
||||
job = _call_cron_for_profile(selected, "resolve_job_ref", job_id)
|
||||
if not job:
|
||||
raise HTTPException(status_code=404, detail="Job not found")
|
||||
# Do not expose the job as due before claiming it: the built-in ticker and
|
||||
# external/manual fire paths share the same durable claim, so only one can
|
||||
# execute this selected run even if they race across processes. Active jobs
|
||||
# keep the legacy provider call shape; paused jobs need the explicit force
|
||||
# flag to resume and claim atomically.
|
||||
force = not job.get("enabled", True) or job.get("state") == "paused"
|
||||
ran = _fire_cron_job_for_profile(selected, job["id"], force=force)
|
||||
refreshed = _call_cron_for_profile(selected, "get_job", job["id"])
|
||||
if refreshed and refreshed.get("last_run_at") != job.get("last_run_at"):
|
||||
return refreshed
|
||||
if not ran:
|
||||
raise HTTPException(
|
||||
status_code=409,
|
||||
detail="Job is already running or was claimed by another scheduler",
|
||||
)
|
||||
if refreshed:
|
||||
return refreshed
|
||||
# A one-shot may remove itself after exhausting repeat=1. Keep the response
|
||||
# shape compatible without inventing an outcome that is no longer present
|
||||
# in the job store; authoritative list refresh removes the completed row.
|
||||
return {
|
||||
**job,
|
||||
"enabled": False,
|
||||
"state": "completed",
|
||||
}
|
||||
|
||||
|
||||
def _delete_cron_job_sync(job_id: str, profile: Optional[str] = None):
|
||||
selected = profile or _find_cron_job_profile(job_id)
|
||||
if not selected:
|
||||
raise HTTPException(status_code=404, detail="Job not found")
|
||||
try:
|
||||
removed = _mutate_cron_for_profile(selected, "remove_job", job_id)
|
||||
except ValueError as exc:
|
||||
raise HTTPException(status_code=400, detail=str(exc)) from exc
|
||||
if not removed:
|
||||
raise HTTPException(status_code=404, detail="Job not found")
|
||||
return {"ok": True}
|
||||
|
||||
# Retry-After hint (seconds) on retryable cron-fire 503s: sized to clear a
|
||||
# scale-to-zero wake or gateway restart so a scheduler that honors it spaces
|
||||
|
||||
@@ -38,6 +38,84 @@ get_hermes_home = late("get_hermes_home")
|
||||
load_config = late("load_config")
|
||||
|
||||
|
||||
# Image MIME types this endpoint will serve. Extension-allowlisted so an
|
||||
# authenticated caller can't pull non-image files through it.
|
||||
_MEDIA_CONTENT_TYPES = {
|
||||
".png": "image/png",
|
||||
".jpg": "image/jpeg",
|
||||
".jpeg": "image/jpeg",
|
||||
".gif": "image/gif",
|
||||
".webp": "image/webp",
|
||||
".svg": "image/svg+xml",
|
||||
".bmp": "image/bmp",
|
||||
".ico": "image/x-icon",
|
||||
}
|
||||
|
||||
|
||||
_MEDIA_MAX_BYTES = 25 * 1024 * 1024
|
||||
|
||||
|
||||
_STREAMABLE_MEDIA_EXTENSIONS = frozenset(
|
||||
{
|
||||
".avi",
|
||||
".flac",
|
||||
".m4a",
|
||||
".mkv",
|
||||
".mov",
|
||||
".mp3",
|
||||
".mp4",
|
||||
".ogg",
|
||||
".opus",
|
||||
".wav",
|
||||
".webm",
|
||||
}
|
||||
)
|
||||
|
||||
|
||||
_FS_READDIR_HIDDEN = {
|
||||
".git",
|
||||
".hg",
|
||||
".svn",
|
||||
".cache",
|
||||
".next",
|
||||
".turbo",
|
||||
".venv",
|
||||
"__pycache__",
|
||||
"build",
|
||||
"dist",
|
||||
"node_modules",
|
||||
"target",
|
||||
"venv",
|
||||
}
|
||||
|
||||
|
||||
# Filenames that must never be listed, read, or downloaded through the
|
||||
# managed-files API. These typically contain credentials (API keys, tokens)
|
||||
# and exposing them through the dashboard file browser is a security leak —
|
||||
# see issue #57505. The set mirrors the credential-file basenames of the two
|
||||
# canonical credential guards elsewhere in the codebase
|
||||
# (agent.file_safety.get_read_block_error and
|
||||
# gateway.platforms.base._ROOT_CREDENTIAL_FILES) so the dashboard Files tab
|
||||
# doesn't lag behind them — an operator can point the managed root at
|
||||
# HERMES_HOME itself, at which point every one of these basenames is a live
|
||||
# secret store sitting in the browsable tree.
|
||||
_SENSITIVE_MANAGED_FILE_BASENAMES = frozenset({
|
||||
"auth.json",
|
||||
"auth.lock",
|
||||
"credentials",
|
||||
"config.yaml",
|
||||
".anthropic_oauth.json",
|
||||
"google_token.json",
|
||||
"google_oauth_pending.json",
|
||||
"google_oauth.json",
|
||||
"webhook_subscriptions.json",
|
||||
"bws_cache.json",
|
||||
"bws_cache.enc.json",
|
||||
# git's credential-store helper cache (agent.file_safety blocks this too).
|
||||
".git-credentials",
|
||||
})
|
||||
|
||||
|
||||
# Directory names whose entire subtree is credential material. Both canonical
|
||||
# guards deny these as directory trees, not basenames:
|
||||
# * gateway.platforms.base._ROOT_CREDENTIAL_DIRS = {"pairing", "mcp-tokens"}
|
||||
@@ -68,7 +146,6 @@ def _is_sensitive_filename(name: str) -> bool:
|
||||
(``mcp-tokens/``, ``pairing/``) that the canonical guards also deny,
|
||||
use :func:`_is_sensitive_path`, which the API call sites route through.
|
||||
"""
|
||||
from hermes_cli.web_server import _SENSITIVE_MANAGED_FILE_BASENAMES
|
||||
lowered = name.lower()
|
||||
if lowered == ".env" or lowered.startswith(".env.") or lowered == ".envrc":
|
||||
return True
|
||||
@@ -282,7 +359,6 @@ async def get_media(path: str):
|
||||
an image-extension allowlist, a size cap, AND the gateway's own media roots
|
||||
(resolved, symlink-safe) so it can't be used to read arbitrary files.
|
||||
"""
|
||||
from hermes_cli.web_server import _MEDIA_CONTENT_TYPES, _MEDIA_MAX_BYTES
|
||||
try:
|
||||
target = Path(path).expanduser().resolve()
|
||||
except (OSError, RuntimeError):
|
||||
@@ -503,7 +579,7 @@ def _managed_file_response(
|
||||
media_only: bool = False,
|
||||
) -> FileResponse:
|
||||
"""Build a range-aware response after applying managed-file policy."""
|
||||
from hermes_cli.web_server import _MANAGED_FILE_MAX_BYTES, _STREAMABLE_MEDIA_EXTENSIONS
|
||||
from hermes_cli.web_server import _MANAGED_FILE_MAX_BYTES
|
||||
policy, target, _display_path = _resolve_managed_path(path, request)
|
||||
if not target.exists():
|
||||
raise HTTPException(status_code=404, detail="File not found")
|
||||
@@ -714,7 +790,6 @@ async def delete_managed_file(payload: ManagedFileDelete, request: Request):
|
||||
|
||||
@router.get("/api/fs/list")
|
||||
async def fs_list(path: str):
|
||||
from hermes_cli.web_server import _FS_READDIR_HIDDEN
|
||||
target = _fs_path(path)
|
||||
try:
|
||||
entries = []
|
||||
|
||||
@@ -1,6 +1,7 @@
|
||||
"""Git dashboard routes — thin executor-offloaded wrappers over ``web_git``
|
||||
(``_git_op`` / ``_git_path`` live in web_server, reached via the late seam)."""
|
||||
|
||||
import asyncio
|
||||
from typing import Optional
|
||||
|
||||
from fastapi import APIRouter
|
||||
@@ -16,11 +17,36 @@ from hermes_cli.web_models import (
|
||||
GitWorktreeRemoveBody,
|
||||
GitBranchSwitchBody,
|
||||
)
|
||||
from fastapi import HTTPException
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
_git_op = late("_git_op")
|
||||
_git_path = late("_git_path")
|
||||
# web_server helpers, late-bound so monkeypatch.setattr(web_server, ...) stays authoritative.
|
||||
_fs_path = late("_fs_path")
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Git ops — the remote half of the desktop coding rail + review pane.
|
||||
#
|
||||
# The desktop runs these as Electron-local git on the user's machine; over a
|
||||
# remote gateway that's the wrong filesystem, so we mirror them here (same auth
|
||||
# gate + path hardening as /api/fs). Logic lives in ``hermes_cli.web_git``;
|
||||
# these are thin, executor-offloaded wrappers (git/gh can block).
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
async def _git_op(fn, *args):
|
||||
"""Run a (blocking) git op off the event loop; map a failed mutation to 400."""
|
||||
loop = asyncio.get_running_loop()
|
||||
try:
|
||||
return await loop.run_in_executor(None, fn, *args)
|
||||
except RuntimeError as exc:
|
||||
raise HTTPException(status_code=400, detail=str(exc) or "git operation failed")
|
||||
|
||||
|
||||
def _git_path(path: str) -> str:
|
||||
return str(_fs_path(path))
|
||||
|
||||
|
||||
|
||||
@router.get("/api/git/status")
|
||||
|
||||
@@ -29,13 +29,13 @@ from hermes_cli.web_routers._common import (
|
||||
log as _log,
|
||||
scoped_to_thread,
|
||||
)
|
||||
import hashlib
|
||||
import re
|
||||
import time
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
_config_profile_scope = late("_config_profile_scope")
|
||||
_gc_mcp_oauth_flows = late("_gc_mcp_oauth_flows")
|
||||
_mcp_install_action_name = late("_mcp_install_action_name")
|
||||
_mcp_oauth_callback_url = late("_mcp_oauth_callback_url")
|
||||
_mcp_server_summary = late("_mcp_server_summary")
|
||||
_normalize_mcp_server_create = late("_normalize_mcp_server_create")
|
||||
_require_token = late("_require_token")
|
||||
@@ -45,8 +45,56 @@ save_config = late("save_config")
|
||||
save_env_value = late("save_env_value")
|
||||
|
||||
_mcp_oauth_flows = LateState("_mcp_oauth_flows")
|
||||
_mcp_oauth_flows_lock = LateState("_mcp_oauth_flows_lock")
|
||||
_MAX_PENDING_MCP_OAUTH_FLOWS = LateState("_MAX_PENDING_MCP_OAUTH_FLOWS")
|
||||
|
||||
|
||||
_MCP_DASHBOARD_OAUTH_TTL = 15 * 60
|
||||
|
||||
|
||||
_mcp_oauth_flows_lock = threading.Lock()
|
||||
|
||||
|
||||
_MAX_PENDING_MCP_OAUTH_FLOWS = 8
|
||||
|
||||
|
||||
def _gc_mcp_oauth_flows() -> None:
|
||||
cutoff = time.time() - _MCP_DASHBOARD_OAUTH_TTL
|
||||
with _mcp_oauth_flows_lock:
|
||||
stale = [
|
||||
flow_id
|
||||
for flow_id, flow in _mcp_oauth_flows.items()
|
||||
if getattr(flow, "created_at", 0) < cutoff
|
||||
]
|
||||
for flow_id in stale:
|
||||
_mcp_oauth_flows.pop(flow_id, None)
|
||||
|
||||
|
||||
def _mcp_oauth_callback_url(request: Request, server_name: str) -> str:
|
||||
"""Build the externally reachable callback URL for a dashboard flow."""
|
||||
from urllib.parse import urlparse, urlunparse
|
||||
|
||||
from hermes_cli.dashboard_auth.prefix import prefix_from_request, resolve_public_url
|
||||
|
||||
from urllib.parse import quote
|
||||
|
||||
suffix = f"/api/mcp/oauth/callback/{quote(server_name, safe='')}"
|
||||
public_url = resolve_public_url()
|
||||
if public_url:
|
||||
return f"{public_url}{suffix}"
|
||||
base = urlparse(str(request.base_url))
|
||||
prefix = prefix_from_request(request)
|
||||
return urlunparse(base._replace(path=f"{prefix}{suffix}", params="", query="", fragment=""))
|
||||
|
||||
|
||||
def _mcp_install_action_name(name: str) -> str:
|
||||
"""Unique per-entry mcp-install action name (+ registered log file), so a
|
||||
re-click or a second catalog install doesn't overwrite the first's tracked
|
||||
process/log while its git clone is still running."""
|
||||
from hermes_cli.web_server import _ACTION_LOG_FILES
|
||||
slug = re.sub(r"[^a-z0-9]+", "-", name.lower()).strip("-")[:48] or "server"
|
||||
digest = hashlib.sha1(name.encode()).hexdigest()[:8]
|
||||
action = f"mcp-install-{slug}-{digest}"
|
||||
_ACTION_LOG_FILES.setdefault(action, f"action-{action}.log")
|
||||
return action
|
||||
|
||||
|
||||
@router.get("/api/mcp/servers")
|
||||
|
||||
@@ -15,23 +15,23 @@ from fastapi import HTTPException
|
||||
from hermes_cli.web_models import MemoryProviderConfigUpdate, MemoryProviderSetupRequest
|
||||
from plugins.memory.config_schema import get_provider_config_schema
|
||||
from typing import Any, Dict, List, Optional
|
||||
from plugins.memory.config_schema import ProviderConfigSchema, ProviderField, STORAGE_HONCHO_HOST_BLOCK
|
||||
import subprocess
|
||||
import json
|
||||
from pathlib import Path
|
||||
|
||||
_log = logging.getLogger("hermes_cli.web_server")
|
||||
router = APIRouter()
|
||||
|
||||
# web_server helpers, late-bound so monkeypatch.setattr(web_server, ...) stays authoritative.
|
||||
_coerce_bool = late("_coerce_bool")
|
||||
_command_result = late("_command_result")
|
||||
_declared_provider_payload = late("_declared_provider_payload")
|
||||
_discover_memory_provider_statuses = late("_discover_memory_provider_statuses")
|
||||
_field_default = late("_field_default")
|
||||
_field_is_set = late("_field_is_set")
|
||||
_field_value = late("_field_value")
|
||||
_field_visible = late("_field_visible")
|
||||
_install_memory_provider_pip_dependencies = late("_install_memory_provider_pip_dependencies")
|
||||
_invalidate_plugins_hub_cache = late("_invalidate_plugins_hub_cache")
|
||||
_load_memory_provider = late("_load_memory_provider")
|
||||
_memory_provider_label = late("_memory_provider_label")
|
||||
_memory_provider_manifest = late("_memory_provider_manifest")
|
||||
_memory_provider_setup_info = late("_memory_provider_setup_info")
|
||||
_memory_provider_setup_manifest = late("_memory_provider_setup_manifest")
|
||||
@@ -40,12 +40,439 @@ _profile_scope = late("_profile_scope")
|
||||
_read_memory_provider_existing_values = late("_read_memory_provider_existing_values")
|
||||
_require_memory_provider_ready = late("_require_memory_provider_ready")
|
||||
_run_setup_command = late("_run_setup_command")
|
||||
_stringify_submitted_values = late("_stringify_submitted_values")
|
||||
_update_memory_provider_config = late("_update_memory_provider_config")
|
||||
get_hermes_home = late("get_hermes_home")
|
||||
load_config = late("load_config")
|
||||
save_config = late("save_config")
|
||||
save_env_value = late("save_env_value")
|
||||
_dependency_importable = late("_dependency_importable")
|
||||
|
||||
|
||||
# Sentinel: remove this key so it falls back to the host or built-in default.
|
||||
_UNSET: Any = object()
|
||||
|
||||
|
||||
def _coerce_field_value(field: ProviderField, raw: str) -> Any:
|
||||
"""Coerce a submitted non-secret value to its native JSON type.
|
||||
|
||||
Values arrive as strings over the API; this converts them to the type the
|
||||
Honcho resolver expects (bool/number/list/dict), so e.g. a boolean is stored
|
||||
as a JSON ``false`` rather than the string ``"false"`` (which would read as
|
||||
truthy). Returns ``_UNSET`` when the field should be removed. Raises
|
||||
``ValueError`` on malformed input.
|
||||
"""
|
||||
|
||||
value = (raw or "").strip()
|
||||
kind = field.kind
|
||||
|
||||
if kind == "select":
|
||||
if not value:
|
||||
value = field.default
|
||||
if value not in field.allowed_values():
|
||||
raise ValueError(f"Invalid value for '{field.key}'")
|
||||
return value
|
||||
|
||||
if kind == "bool":
|
||||
from utils import is_truthy_value
|
||||
|
||||
return is_truthy_value(value)
|
||||
|
||||
if kind == "number":
|
||||
if not value:
|
||||
return _UNSET
|
||||
try:
|
||||
number = float(value)
|
||||
except ValueError as exc:
|
||||
raise ValueError(f"Invalid number for '{field.key}'") from exc
|
||||
return int(number) if number.is_integer() else number
|
||||
|
||||
if kind == "json":
|
||||
if not value:
|
||||
return _UNSET
|
||||
try:
|
||||
parsed = json.loads(value)
|
||||
except (ValueError, TypeError) as exc:
|
||||
raise ValueError(f"Invalid JSON for '{field.key}'") from exc
|
||||
if not isinstance(parsed, (dict, list)):
|
||||
raise ValueError(f"'{field.key}' must be a JSON object or array")
|
||||
return parsed
|
||||
|
||||
# text / secret — blank clears the key so it falls back to host/default.
|
||||
return value if value else _UNSET
|
||||
|
||||
|
||||
# — flat-json backend (default; reusable for simple providers) —
|
||||
|
||||
|
||||
def _flat_json_path(provider: ProviderConfigSchema) -> Path:
|
||||
return get_hermes_home() / provider.name / "config.json"
|
||||
|
||||
|
||||
def _read_flat_json(provider: ProviderConfigSchema) -> Dict[str, Any]:
|
||||
path = _flat_json_path(provider)
|
||||
if not path.exists():
|
||||
return {}
|
||||
try:
|
||||
data = json.loads(path.read_text(encoding="utf-8"))
|
||||
except Exception:
|
||||
_log.warning("Failed to read memory provider config from %s", path, exc_info=True)
|
||||
return {}
|
||||
return data if isinstance(data, dict) else {}
|
||||
|
||||
|
||||
# — honcho host-block backend —
|
||||
|
||||
|
||||
def _honcho_resolvers():
|
||||
"""Lazily import the Honcho plugin's resolvers (optional plugin)."""
|
||||
|
||||
from plugins.memory.honcho.client import _host_block, resolve_active_host, resolve_config_path
|
||||
|
||||
return resolve_active_host, resolve_config_path, _host_block
|
||||
|
||||
|
||||
def _apply_field_values(provider: ProviderConfigSchema, values: Dict[str, str], target_for) -> None:
|
||||
"""Apply submitted non-secret fields to their backend dict, in place.
|
||||
|
||||
Only keys present in ``values`` are touched, so a partial save never
|
||||
clobbers fields owned by another surface. ``_UNSET`` clears the key (and
|
||||
its aliases) so it falls back to the host/default mapping.
|
||||
"""
|
||||
|
||||
for field in provider.fields:
|
||||
if field.is_secret or field.key not in values:
|
||||
continue
|
||||
target = target_for(field)
|
||||
coerced = _coerce_field_value(field, values[field.key])
|
||||
if coerced is _UNSET:
|
||||
target.pop(field.key, None)
|
||||
for alias in field.aliases:
|
||||
target.pop(alias, None)
|
||||
else:
|
||||
target[field.key] = coerced
|
||||
|
||||
|
||||
def _trim_setup_output(value: Optional[str], limit: int = 4000) -> str:
|
||||
text = str(value or "")
|
||||
if len(text) <= limit:
|
||||
return text
|
||||
return f"{text[:limit]}\n... truncated ..."
|
||||
|
||||
|
||||
# ── Memory provider config: one generic GET/PUT pair, dispatching on storage ──
|
||||
|
||||
|
||||
def _provider_field_entry(field: ProviderField) -> Dict[str, Any]:
|
||||
"""Static, storage-independent shape of one field for the UI payload."""
|
||||
|
||||
return {
|
||||
"key": field.key,
|
||||
"label": field.label,
|
||||
"kind": field.kind,
|
||||
"description": field.description,
|
||||
"info": field.info,
|
||||
"placeholder": field.placeholder,
|
||||
"inline": field.inline,
|
||||
"group": field.group,
|
||||
"options": [
|
||||
{"value": opt.value, "label": opt.label, "description": opt.description}
|
||||
for opt in field.options
|
||||
],
|
||||
}
|
||||
|
||||
|
||||
def _serialize_field_value(field: ProviderField, value: Any) -> str:
|
||||
"""Render a stored native value as the string the generic UI edits.
|
||||
|
||||
``None`` (key absent) yields the field's declared default. Bools become
|
||||
``"true"``/``"false"``, JSON objects/arrays are re-encoded, numbers are
|
||||
stringified — so the renderer's per-kind controls always get the shape they
|
||||
expect regardless of how the value sits on disk.
|
||||
"""
|
||||
|
||||
if value is None:
|
||||
return field.default
|
||||
if field.kind == "bool":
|
||||
from utils import is_truthy_value
|
||||
|
||||
return "true" if is_truthy_value(value) else "false"
|
||||
if field.kind == "json":
|
||||
if isinstance(value, (dict, list)):
|
||||
return json.dumps(value)
|
||||
return str(value)
|
||||
return str(value)
|
||||
|
||||
|
||||
def _read_field(field: ProviderField, sources: tuple, env: Dict[str, str]) -> Any:
|
||||
"""Return the stored native value from the first source holding it, or ``None``.
|
||||
|
||||
Presence (``key in source``) decides, not truthiness, so a stored ``False``
|
||||
or ``0`` survives instead of being mistaken for "unset".
|
||||
"""
|
||||
|
||||
for source in sources:
|
||||
for source_key in (field.key, *field.aliases):
|
||||
if source_key in source and source[source_key] is not None:
|
||||
return source[source_key]
|
||||
for env_key in field.env_fallbacks:
|
||||
value = env.get(env_key)
|
||||
if value:
|
||||
return value
|
||||
return None
|
||||
|
||||
|
||||
def _declared_field_is_set(field: ProviderField, sources: tuple, env: Dict[str, str]) -> bool:
|
||||
for env_key in (field.env_key, *field.env_fallbacks):
|
||||
if env_key and env.get(env_key):
|
||||
return True
|
||||
return any(source.get(k) for source in sources for k in (field.key, *field.aliases))
|
||||
|
||||
|
||||
def _honcho_read_sources() -> tuple[Dict[str, Any], str, Dict[str, Any]]:
|
||||
"""Return (root config, active host key, host block) for the current profile."""
|
||||
|
||||
resolve_active_host, resolve_config_path, host_block_of = _honcho_resolvers()
|
||||
host = resolve_active_host()
|
||||
path = resolve_config_path()
|
||||
raw: Dict[str, Any] = {}
|
||||
if path.exists():
|
||||
try:
|
||||
loaded = json.loads(path.read_text(encoding="utf-8"))
|
||||
raw = loaded if isinstance(loaded, dict) else {}
|
||||
except Exception:
|
||||
_log.warning("Failed to read Honcho config from %s", path, exc_info=True)
|
||||
return raw, host, host_block_of(raw, host)
|
||||
|
||||
|
||||
def _write_provider_flat(provider: ProviderConfigSchema, values: Dict[str, str]) -> None:
|
||||
from utils import atomic_json_write
|
||||
|
||||
existing = _read_flat_json(provider)
|
||||
|
||||
for field in provider.fields:
|
||||
if field.is_secret:
|
||||
submitted = (values.get(field.key) or "").strip()
|
||||
if submitted and field.env_key:
|
||||
save_env_value(field.env_key, submitted)
|
||||
|
||||
_apply_field_values(provider, values, lambda field: existing)
|
||||
|
||||
path = _flat_json_path(provider)
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
atomic_json_write(path, existing, mode=0o600)
|
||||
|
||||
|
||||
def _write_provider_honcho(provider: ProviderConfigSchema, values: Dict[str, str]) -> None:
|
||||
"""Persist submitted fields to Honcho's real config for the active host.
|
||||
|
||||
Only keys present in ``values`` are touched, so a partial save (e.g. the
|
||||
inline panel) never clobbers fields owned by the full-config editor. Blank
|
||||
text clears a key so it falls back to the host/default mapping.
|
||||
"""
|
||||
|
||||
from plugins.memory.honcho.oauth import ACCESS_TOKEN_PREFIX, _config_refresh_lock
|
||||
from utils import atomic_json_write
|
||||
|
||||
resolve_active_host, resolve_config_path, host_block_of = _honcho_resolvers()
|
||||
host = resolve_active_host()
|
||||
# Write the file reads resolve, or a save shadows it with a sparse copy.
|
||||
path = resolve_config_path()
|
||||
|
||||
# OAuth rotation is single-use; an unlocked RMW here can revoke the grant.
|
||||
with _config_refresh_lock(path):
|
||||
cfg: Dict[str, Any] = {}
|
||||
if path.exists():
|
||||
try:
|
||||
loaded = json.loads(path.read_text(encoding="utf-8"))
|
||||
cfg = loaded if isinstance(loaded, dict) else {}
|
||||
except Exception:
|
||||
_log.warning("Failed to read Honcho config from %s", path, exc_info=True)
|
||||
|
||||
hosts = cfg.get("hosts")
|
||||
cfg["hosts"] = hosts = hosts if isinstance(hosts, dict) else {}
|
||||
# Update the block reads resolve (legacy dot-form included), never shadow it.
|
||||
existing = host_block_of(cfg, host)
|
||||
host_key = next((k for k, v in hosts.items() if v is existing), host) if existing else host
|
||||
host_block = hosts.setdefault(host_key, existing)
|
||||
|
||||
for field in provider.fields:
|
||||
if not field.is_secret:
|
||||
continue
|
||||
submitted = (values.get(field.key) or "").strip()
|
||||
if not submitted:
|
||||
continue
|
||||
if field.env_key:
|
||||
save_env_value(field.env_key, submitted)
|
||||
# Persist where the client reads first; an OAuth token owns that slot.
|
||||
stored = host_block.get(field.key)
|
||||
if not (isinstance(stored, str) and stored.startswith(ACCESS_TOKEN_PREFIX)):
|
||||
host_block[field.key] = submitted
|
||||
|
||||
_apply_field_values(provider, values, lambda field: host_block if field.scope == "host" else cfg)
|
||||
|
||||
path.parent.mkdir(parents=True, exist_ok=True)
|
||||
atomic_json_write(path, cfg, mode=0o600)
|
||||
|
||||
|
||||
def _command_result(
|
||||
*,
|
||||
kind: str,
|
||||
name: str,
|
||||
status: str,
|
||||
command: str = "",
|
||||
completed: Optional[subprocess.CompletedProcess] = None,
|
||||
error: Optional[str] = None,
|
||||
) -> Dict[str, Any]:
|
||||
return {
|
||||
"kind": kind,
|
||||
"name": name,
|
||||
"status": status,
|
||||
"command": command,
|
||||
"returncode": None if completed is None else completed.returncode,
|
||||
"stdout": "" if completed is None else _trim_setup_output(completed.stdout),
|
||||
"stderr": _trim_setup_output(error or ("" if completed is None else completed.stderr)),
|
||||
}
|
||||
|
||||
|
||||
def _declared_provider_payload(provider: ProviderConfigSchema) -> Dict[str, Any]:
|
||||
from hermes_cli.web_server import load_env
|
||||
fields: List[Dict[str, Any]] = []
|
||||
env = load_env()
|
||||
is_honcho = provider.storage == STORAGE_HONCHO_HOST_BLOCK
|
||||
|
||||
if is_honcho:
|
||||
raw, host, host_block = _honcho_read_sources()
|
||||
|
||||
def sources_for(field: ProviderField) -> tuple:
|
||||
return (host_block, raw) if field.scope == "host" else (raw,)
|
||||
else:
|
||||
host = ""
|
||||
data = _read_flat_json(provider)
|
||||
|
||||
def sources_for(field: ProviderField) -> tuple:
|
||||
return (data,)
|
||||
|
||||
for field in provider.fields:
|
||||
entry = _provider_field_entry(field)
|
||||
sources = sources_for(field)
|
||||
|
||||
if field.is_secret:
|
||||
entry["value"] = "" # secrets are write-only over the API
|
||||
entry["is_set"] = _declared_field_is_set(field, sources, env)
|
||||
fields.append(entry)
|
||||
continue
|
||||
|
||||
native = _read_field(field, sources, env)
|
||||
if is_honcho and not field.placeholder and field.key in {"workspace", "aiPeer"}:
|
||||
# Blank fields surface the resolved host Honcho will actually use.
|
||||
entry["placeholder"] = host
|
||||
|
||||
value = _serialize_field_value(field, native)
|
||||
if field.kind == "select" and value not in field.allowed_values():
|
||||
value = field.default
|
||||
entry["value"] = value
|
||||
# Presence, not truthiness — a stored False/0 is still "set".
|
||||
entry["is_set"] = native is not None if is_honcho else bool(value)
|
||||
fields.append(entry)
|
||||
|
||||
return {"name": provider.name, "label": provider.label, "docs_url": provider.docs_url, "fields": fields}
|
||||
|
||||
|
||||
def _stringify_submitted_values(values: Dict[str, Any]) -> Dict[str, str]:
|
||||
"""The declared-schema path edits strings; the dashboard may send natives."""
|
||||
|
||||
out: Dict[str, str] = {}
|
||||
for key, value in values.items():
|
||||
if value is None:
|
||||
out[key] = ""
|
||||
elif isinstance(value, str):
|
||||
out[key] = value
|
||||
elif isinstance(value, bool):
|
||||
out[key] = "true" if value else "false"
|
||||
elif isinstance(value, (dict, list)):
|
||||
out[key] = json.dumps(value)
|
||||
else:
|
||||
out[key] = str(value)
|
||||
return out
|
||||
|
||||
|
||||
def _update_memory_provider_config(provider: ProviderConfigSchema, values: Dict[str, str]) -> None:
|
||||
if provider.storage == STORAGE_HONCHO_HOST_BLOCK:
|
||||
_write_provider_honcho(provider, values)
|
||||
else:
|
||||
_write_provider_flat(provider, values)
|
||||
|
||||
config = load_config()
|
||||
memory_config = config.get("memory")
|
||||
if not isinstance(memory_config, dict):
|
||||
memory_config = {}
|
||||
config["memory"] = memory_config
|
||||
if memory_config.get("provider") != provider.name:
|
||||
memory_config["provider"] = provider.name
|
||||
save_config(config)
|
||||
|
||||
|
||||
def _memory_provider_label(name: str) -> str:
|
||||
return name.replace("_", " ").replace("-", " ").title()
|
||||
|
||||
|
||||
def _install_memory_provider_pip_dependencies(dependencies: List[str]) -> List[Dict[str, Any]]:
|
||||
missing = [dep for dep in dependencies if not _dependency_importable(dep)]
|
||||
if not dependencies:
|
||||
return []
|
||||
if not missing:
|
||||
return [
|
||||
_command_result(kind="pip", name=", ".join(dependencies), status="already_installed")
|
||||
]
|
||||
|
||||
# Route through the lazy-install pipeline (tools.lazy_deps.install_specs)
|
||||
# instead of shelling out to pip against sys.executable directly. That
|
||||
# pipeline is environment-aware: on hosted/immutable images the agent venv
|
||||
# under /opt/hermes is sealed read-only, and installs must be redirected
|
||||
# to the writable durable target on the data volume
|
||||
# (HERMES_LAZY_INSTALL_TARGET, e.g. /opt/data/lazy-packages) — the same
|
||||
# path every lazy backend already uses. A direct `pip install --python
|
||||
# sys.executable` on those images fails with a permission error (NS-605).
|
||||
# install_specs also activates the target on sys.path post-install so the
|
||||
# availability recheck below sees the new packages without a restart.
|
||||
try:
|
||||
from tools.lazy_deps import install_specs
|
||||
|
||||
outcome = install_specs(missing, timeout=240)
|
||||
except Exception as exc:
|
||||
return [
|
||||
_command_result(
|
||||
kind="pip",
|
||||
name=", ".join(missing),
|
||||
status="failed",
|
||||
error=str(exc),
|
||||
)
|
||||
]
|
||||
|
||||
if outcome.blocked:
|
||||
return [
|
||||
_command_result(
|
||||
kind="pip",
|
||||
name=", ".join(missing),
|
||||
status="failed",
|
||||
command=outcome.command,
|
||||
error=outcome.reason,
|
||||
)
|
||||
]
|
||||
|
||||
return [
|
||||
_command_result(
|
||||
kind="pip",
|
||||
name=", ".join(missing),
|
||||
status="installed" if outcome.ok else "failed",
|
||||
command=outcome.command,
|
||||
completed=subprocess.CompletedProcess(
|
||||
args=outcome.command,
|
||||
returncode=0 if outcome.ok else 1,
|
||||
stdout=outcome.stdout,
|
||||
stderr=outcome.stderr,
|
||||
),
|
||||
)
|
||||
]
|
||||
|
||||
|
||||
def _install_memory_provider_external_dependencies(
|
||||
|
||||
@@ -13,6 +13,7 @@ import json
|
||||
import secrets
|
||||
import time
|
||||
import urllib.parse
|
||||
import os
|
||||
from datetime import datetime, timezone
|
||||
from fastapi import APIRouter
|
||||
from hermes_cli.web_deps import late
|
||||
@@ -22,23 +23,22 @@ from hermes_cli.config import get_env_path
|
||||
from hermes_cli.web_models import MessagingPlatformUpdate, TelegramOnboardingStart, TelegramOnboardingApply, WhatsAppOnboardingStart, WhatsAppOnboardingApply
|
||||
from pathlib import Path
|
||||
from typing import Any, Optional
|
||||
from hermes_cli.config import redact_key
|
||||
from gateway.status import resolve_gateway_liveness
|
||||
from hermes_cli.config import OPTIONAL_ENV_VARS
|
||||
|
||||
_log = logging.getLogger("hermes_cli.web_server")
|
||||
router = APIRouter()
|
||||
|
||||
# web_server helpers, late-bound so monkeypatch.setattr(web_server, ...) stays authoritative.
|
||||
_catalog_lookup = late("_catalog_lookup")
|
||||
_config_profile_scope = late("_config_profile_scope")
|
||||
_gateway_display_command = late("_gateway_display_command")
|
||||
_messaging_platform_catalog = late("_messaging_platform_catalog")
|
||||
_messaging_platform_payload = late("_messaging_platform_payload")
|
||||
_profile_scope = late("_profile_scope")
|
||||
_resolve_profile_dir = late("_resolve_profile_dir")
|
||||
_restart_gateway_after_whatsapp_onboarding = late("_restart_gateway_after_whatsapp_onboarding")
|
||||
_restart_gateway_after = late("_restart_gateway_after")
|
||||
_telegram_onboarding_error_message = late("_telegram_onboarding_error_message")
|
||||
_telegram_onboarding_request_sync = late("_telegram_onboarding_request_sync")
|
||||
_validate_messaging_env_value = late("_validate_messaging_env_value")
|
||||
_whatsapp_onboarding_payload = late("_whatsapp_onboarding_payload")
|
||||
_whatsapp_session_path = late("_whatsapp_session_path")
|
||||
_write_platform_enabled = late("_write_platform_enabled")
|
||||
@@ -46,6 +46,417 @@ load_env = late("load_env")
|
||||
read_runtime_status = late("read_runtime_status")
|
||||
remove_env_value = late("remove_env_value")
|
||||
save_env_value = late("save_env_value")
|
||||
_gateway_subcommand = late("_gateway_subcommand")
|
||||
_probe_gateway_health = late("_probe_gateway_health")
|
||||
|
||||
|
||||
# Display labels for env vars not in OPTIONAL_ENV_VARS (HOME_CHANNEL_*, bridge
|
||||
# toggles, Twilio, HASS, Email, etc.). Anything missing from OPTIONAL_ENV_VARS
|
||||
# falls back here so the UI can still render a friendly label.
|
||||
_MESSAGING_ENV_FALLBACKS: dict[str, dict[str, Any]] = {
|
||||
"SIGNAL_HTTP_URL": {
|
||||
"description": "signal-cli REST API base URL, e.g. http://127.0.0.1:8080",
|
||||
"prompt": "Signal bridge URL",
|
||||
"url": "https://github.com/bbernhard/signal-cli-rest-api",
|
||||
},
|
||||
"SIGNAL_ACCOUNT": {
|
||||
"description": "Signal account phone number registered with the bridge",
|
||||
"prompt": "Signal account",
|
||||
},
|
||||
"SIGNAL_ALLOWED_USERS": {
|
||||
"description": "Comma-separated Signal users allowed to use the bot",
|
||||
"prompt": "Allowed Signal users",
|
||||
},
|
||||
"WHATSAPP_ENABLED": {
|
||||
"description": "Enable the WhatsApp gateway adapter",
|
||||
"prompt": "Enable WhatsApp",
|
||||
"advanced": True,
|
||||
},
|
||||
"WHATSAPP_MODE": {
|
||||
"description": "WhatsApp bridge mode",
|
||||
"prompt": "WhatsApp mode",
|
||||
"advanced": True,
|
||||
},
|
||||
"WHATSAPP_DM_POLICY": {
|
||||
"description": "How WhatsApp direct messages are authorized",
|
||||
"prompt": "WhatsApp DM policy",
|
||||
"advanced": True,
|
||||
},
|
||||
"WHATSAPP_ALLOWED_USERS": {
|
||||
"description": "Comma-separated WhatsApp users allowed to use the bot",
|
||||
"prompt": "Allowed WhatsApp users",
|
||||
},
|
||||
"HASS_URL": {
|
||||
"description": "Home Assistant base URL, e.g. https://homeassistant.local:8123",
|
||||
"prompt": "Home Assistant URL",
|
||||
},
|
||||
"HASS_TOKEN": {
|
||||
"description": "Long-lived access token from Home Assistant (Profile → Security)",
|
||||
"prompt": "Home Assistant access token",
|
||||
"password": True,
|
||||
},
|
||||
"EMAIL_ADDRESS": {
|
||||
"description": "Email address to send and receive from",
|
||||
"prompt": "Email address",
|
||||
},
|
||||
"EMAIL_PASSWORD": {
|
||||
"description": "Email account password or app password",
|
||||
"prompt": "Email password",
|
||||
"password": True,
|
||||
},
|
||||
"EMAIL_IMAP_HOST": {
|
||||
"description": "IMAP server host (e.g. imap.gmail.com)",
|
||||
"prompt": "IMAP host",
|
||||
},
|
||||
"EMAIL_SMTP_HOST": {
|
||||
"description": "SMTP server host (e.g. smtp.gmail.com)",
|
||||
"prompt": "SMTP host",
|
||||
},
|
||||
"TWILIO_ACCOUNT_SID": {
|
||||
"description": "Twilio Account SID",
|
||||
"prompt": "Twilio Account SID",
|
||||
"url": "https://www.twilio.com/console",
|
||||
},
|
||||
"TWILIO_AUTH_TOKEN": {
|
||||
"description": "Twilio Auth Token",
|
||||
"prompt": "Twilio Auth Token",
|
||||
"password": True,
|
||||
},
|
||||
"WECOM_BOT_ID": {"description": "WeCom group bot ID", "prompt": "WeCom Bot ID"},
|
||||
"WECOM_SECRET": {
|
||||
"description": "WeCom group bot secret",
|
||||
"prompt": "WeCom Secret",
|
||||
"password": True,
|
||||
},
|
||||
"WECOM_CALLBACK_CORP_ID": {
|
||||
"description": "WeCom corp ID",
|
||||
"prompt": "WeCom Corp ID",
|
||||
},
|
||||
"WECOM_CALLBACK_CORP_SECRET": {
|
||||
"description": "WeCom app corp secret",
|
||||
"prompt": "WeCom Corp Secret",
|
||||
"password": True,
|
||||
},
|
||||
"WECOM_CALLBACK_AGENT_ID": {
|
||||
"description": "WeCom app agent ID",
|
||||
"prompt": "WeCom Agent ID",
|
||||
},
|
||||
"WECOM_CALLBACK_TOKEN": {
|
||||
"description": "WeCom callback verification token",
|
||||
"prompt": "WeCom Token",
|
||||
},
|
||||
"WECOM_CALLBACK_ENCODING_AES_KEY": {
|
||||
"description": "WeCom callback AES encoding key",
|
||||
"prompt": "WeCom AES Key",
|
||||
"password": True,
|
||||
},
|
||||
"WEIXIN_ACCOUNT_ID": {
|
||||
"description": "iLink Bot account ID obtained through QR login in hermes gateway setup",
|
||||
"prompt": "iLink Bot account ID",
|
||||
},
|
||||
"WEIXIN_TOKEN": {
|
||||
"description": "iLink Bot token obtained through QR login in hermes gateway setup",
|
||||
"prompt": "iLink Bot token",
|
||||
"password": True,
|
||||
},
|
||||
"WEIXIN_BASE_URL": {
|
||||
"description": "iLink API base URL saved by QR login (default: https://ilinkai.weixin.qq.com)",
|
||||
"prompt": "iLink API base URL",
|
||||
},
|
||||
"FEISHU_APP_ID": {"description": "Feishu / Lark app ID", "prompt": "App ID"},
|
||||
"FEISHU_APP_SECRET": {
|
||||
"description": "Feishu / Lark app secret",
|
||||
"prompt": "App secret",
|
||||
"password": True,
|
||||
},
|
||||
"FEISHU_ENCRYPT_KEY": {
|
||||
"description": "Feishu / Lark encrypt key",
|
||||
"prompt": "Encrypt key",
|
||||
"password": True,
|
||||
},
|
||||
"FEISHU_VERIFICATION_TOKEN": {
|
||||
"description": "Feishu / Lark verification token",
|
||||
"prompt": "Verification token",
|
||||
"password": True,
|
||||
},
|
||||
"DINGTALK_CLIENT_ID": {
|
||||
"description": "DingTalk client ID (App key)",
|
||||
"prompt": "Client ID",
|
||||
},
|
||||
"DINGTALK_CLIENT_SECRET": {
|
||||
"description": "DingTalk client secret (App secret)",
|
||||
"prompt": "Client secret",
|
||||
"password": True,
|
||||
},
|
||||
}
|
||||
|
||||
|
||||
# Kept in sync with the corresponding frontend validation in ChannelsPage.tsx.
|
||||
_TELEGRAM_BOT_TOKEN_RE = re.compile(r"\d+:[A-Za-z0-9_-]{30,}")
|
||||
|
||||
|
||||
_TELEGRAM_USER_ID_RE = re.compile(r"\d+")
|
||||
|
||||
|
||||
_SLACK_MEMBER_ID_RE = re.compile(r"[UW][A-Z0-9]{2,}")
|
||||
|
||||
|
||||
def _messaging_env_info(key: str) -> dict[str, Any]:
|
||||
info = OPTIONAL_ENV_VARS.get(key) or _MESSAGING_ENV_FALLBACKS.get(key) or {}
|
||||
return {
|
||||
"description": info.get("description", ""),
|
||||
"prompt": info.get("prompt", key),
|
||||
"help": info.get("help", ""),
|
||||
"url": info.get("url"),
|
||||
"is_password": info.get("password", False),
|
||||
"advanced": info.get("advanced", False),
|
||||
}
|
||||
|
||||
|
||||
def _gateway_platform_config(platform_id: str):
|
||||
from gateway.config import Platform, load_gateway_config
|
||||
|
||||
config = load_gateway_config()
|
||||
platform = Platform(platform_id)
|
||||
platform_config = config.platforms.get(platform)
|
||||
return config, platform, platform_config
|
||||
|
||||
|
||||
def _gateway_display_command(profile: Optional[str], verb: str) -> str:
|
||||
return " ".join(["hermes", *_gateway_subcommand(profile, verb)])
|
||||
|
||||
|
||||
def _validate_messaging_env_value(platform_id: str, key: str, value: str) -> None:
|
||||
"""Reject platform credentials that are clearly in the wrong field."""
|
||||
if not value:
|
||||
return
|
||||
|
||||
if platform_id == "telegram":
|
||||
if key == "TELEGRAM_BOT_TOKEN" and not _TELEGRAM_BOT_TOKEN_RE.fullmatch(value):
|
||||
raise HTTPException(
|
||||
status_code=400,
|
||||
detail="Telegram bot token must be the complete token from @BotFather, such as 123456789:ABC…",
|
||||
)
|
||||
if key == "TELEGRAM_ALLOWED_USERS":
|
||||
user_ids = [part.strip() for part in value.split(",") if part.strip()]
|
||||
if any(not _TELEGRAM_USER_ID_RE.fullmatch(user_id) for user_id in user_ids):
|
||||
raise HTTPException(
|
||||
status_code=400,
|
||||
detail="Telegram allowed users must be comma-separated numeric user IDs.",
|
||||
)
|
||||
return
|
||||
|
||||
if platform_id != "slack":
|
||||
return
|
||||
|
||||
if key == "SLACK_BOT_TOKEN" and not value.startswith("xoxb-"):
|
||||
raise HTTPException(
|
||||
status_code=400,
|
||||
detail="Slack Bot Token must start with xoxb-. Paste the bot token from OAuth & Permissions.",
|
||||
)
|
||||
if key == "SLACK_APP_TOKEN" and not value.startswith("xapp-"):
|
||||
raise HTTPException(
|
||||
status_code=400,
|
||||
detail="Slack App Token must start with xapp-. Paste the app-level token from Basic Information > App-Level Tokens.",
|
||||
)
|
||||
if key == "SLACK_ALLOWED_USERS":
|
||||
# Mirror the gateway's parse (gateway/platforms/slack.py): split on comma,
|
||||
# strip, and drop empty entries so a trailing/interior comma isn't rejected
|
||||
# here when the runtime would accept it. "*" is the allow-all wildcard.
|
||||
user_ids = [part.strip() for part in value.split(",") if part.strip()]
|
||||
invalid = [
|
||||
user_id
|
||||
for user_id in user_ids
|
||||
if user_id != "*" and not _SLACK_MEMBER_ID_RE.fullmatch(user_id)
|
||||
]
|
||||
if invalid:
|
||||
raise HTTPException(
|
||||
status_code=400,
|
||||
detail="Slack allowed user IDs must be comma-separated member IDs like U01ABC2DEF3.",
|
||||
)
|
||||
|
||||
|
||||
def _catalog_lookup(platform_id: str) -> dict[str, Any] | None:
|
||||
for entry in _messaging_platform_catalog():
|
||||
if entry["id"] == platform_id:
|
||||
return entry
|
||||
return None
|
||||
|
||||
|
||||
def _messaging_platform_payload(
|
||||
entry: dict[str, Any],
|
||||
env_on_disk: dict[str, str],
|
||||
runtime: dict | None,
|
||||
scoped: bool = False,
|
||||
profile_home: Optional[Path] = None,
|
||||
) -> dict[str, Any]:
|
||||
from hermes_cli.web_server import (
|
||||
_GATEWAY_HEALTH_URL,
|
||||
get_running_pid_cached,
|
||||
get_runtime_status_running_pid,
|
||||
load_config,
|
||||
)
|
||||
platform_id = entry["id"]
|
||||
runtime_platforms = runtime.get("platforms") if runtime else {}
|
||||
runtime_platform = (
|
||||
runtime_platforms.get(platform_id, {})
|
||||
if isinstance(runtime_platforms, dict)
|
||||
else {}
|
||||
)
|
||||
# Same shared ladder /api/status uses. Before this was unified, the two
|
||||
# endpoints disagreed on the same page load — the sidebar strip read
|
||||
# "running" (it probed GATEWAY_HEALTH_URL and scoped to the requested
|
||||
# profile) while the Channels page rendered "The gateway is not running"
|
||||
# (it did neither). Cross-container, profile-scoped, and
|
||||
# launch-service-managed deployments each hit that split.
|
||||
#
|
||||
# profile_home is passed when the request was scoped to a named profile:
|
||||
# gateway/status readers resolve process-level paths and do NOT follow the
|
||||
# HERMES_HOME contextvar override (#56986 / #69143), so the profile's
|
||||
# directory has to be handed over explicitly or messaging silently reports
|
||||
# another profile's gateway (#71211).
|
||||
liveness = resolve_gateway_liveness(
|
||||
profile_dir=profile_home,
|
||||
runtime=runtime,
|
||||
health_probe=(
|
||||
_probe_gateway_health if _GATEWAY_HEALTH_URL else None
|
||||
),
|
||||
pid_probe=get_running_pid_cached,
|
||||
runtime_reader=read_runtime_status,
|
||||
runtime_pid_probe=get_runtime_status_running_pid,
|
||||
)
|
||||
gateway_running = liveness.running
|
||||
env_vars = []
|
||||
|
||||
for key in entry["env_vars"]:
|
||||
# When profile-scoped, judge only the profile's own .env — the
|
||||
# dashboard process's os.environ carries the ROOT install's .env
|
||||
# (loaded at startup) and would falsely report the root credentials
|
||||
# as the profile's.
|
||||
value = env_on_disk.get(key) or ("" if scoped else os.getenv(key, ""))
|
||||
env_vars.append(
|
||||
{
|
||||
"key": key,
|
||||
"required": key in entry["required_env"],
|
||||
"is_set": bool(value),
|
||||
"redacted_value": redact_key(value) if value else None,
|
||||
**_messaging_env_info(key),
|
||||
}
|
||||
)
|
||||
|
||||
if scoped:
|
||||
# Profile-scoped view: derive enablement/configuration from the
|
||||
# profile's config.yaml + .env only. load_gateway_config()'s
|
||||
# env-override layer reads os.environ and would leak the root
|
||||
# install's tokens into the profile's reported state.
|
||||
try:
|
||||
cfg = load_config()
|
||||
platforms_cfg = cfg.get("platforms") or {}
|
||||
plat_cfg = platforms_cfg.get(platform_id)
|
||||
if not isinstance(plat_cfg, dict):
|
||||
plat_cfg = {}
|
||||
enabled = bool(plat_cfg.get("enabled"))
|
||||
hc = plat_cfg.get("home_channel")
|
||||
home_channel = hc if isinstance(hc, dict) else None
|
||||
except Exception:
|
||||
enabled = False
|
||||
home_channel = None
|
||||
configured = all(env_on_disk.get(key) for key in entry["required_env"])
|
||||
else:
|
||||
try:
|
||||
gateway_config, platform, platform_config = _gateway_platform_config(
|
||||
platform_id
|
||||
)
|
||||
enabled = bool(platform_config and platform_config.enabled)
|
||||
configured = bool(
|
||||
platform_config
|
||||
and gateway_config._is_platform_connected(platform, platform_config)
|
||||
)
|
||||
home_channel = (
|
||||
platform_config.home_channel.to_dict()
|
||||
if platform_config and platform_config.home_channel
|
||||
else None
|
||||
)
|
||||
except Exception:
|
||||
enabled = False
|
||||
configured = all(
|
||||
env_on_disk.get(key) or os.getenv(key, "")
|
||||
for key in entry["required_env"]
|
||||
)
|
||||
home_channel = None
|
||||
|
||||
state = (
|
||||
runtime_platform.get("state") if isinstance(runtime_platform, dict) else None
|
||||
)
|
||||
runtime_gateway_state = runtime.get("gateway_state") if isinstance(runtime, dict) else None
|
||||
runtime_gateway_error = runtime.get("exit_reason") if isinstance(runtime, dict) else None
|
||||
if not enabled:
|
||||
state = "disabled"
|
||||
elif not configured:
|
||||
state = "not_configured"
|
||||
elif gateway_running and not state:
|
||||
state = "pending_restart"
|
||||
elif (
|
||||
not gateway_running
|
||||
and not state
|
||||
and runtime_gateway_state == "startup_failed"
|
||||
):
|
||||
state = "startup_failed"
|
||||
elif not gateway_running and not state:
|
||||
state = "gateway_stopped"
|
||||
|
||||
error_code = (
|
||||
runtime_platform.get("error_code")
|
||||
if isinstance(runtime_platform, dict)
|
||||
else None
|
||||
)
|
||||
error_message = (
|
||||
runtime_platform.get("error_message")
|
||||
if isinstance(runtime_platform, dict)
|
||||
else None
|
||||
)
|
||||
if state == "startup_failed":
|
||||
error_code = error_code or "startup_failed"
|
||||
error_message = error_message or runtime_gateway_error
|
||||
|
||||
whatsapp_setup = None
|
||||
if platform_id == "whatsapp":
|
||||
whatsapp_mode = (
|
||||
env_on_disk.get("WHATSAPP_MODE")
|
||||
or ("" if scoped else os.getenv("WHATSAPP_MODE", ""))
|
||||
).strip()
|
||||
allowed_users_value = (
|
||||
env_on_disk.get("WHATSAPP_ALLOWED_USERS")
|
||||
or ("" if scoped else os.getenv("WHATSAPP_ALLOWED_USERS", ""))
|
||||
).strip()
|
||||
whatsapp_setup = {
|
||||
"mode": whatsapp_mode if whatsapp_mode in {"bot", "self-chat"} else "",
|
||||
"allowed_users_set": bool(allowed_users_value),
|
||||
"home_channel_set": bool(home_channel),
|
||||
}
|
||||
|
||||
payload = {
|
||||
"id": platform_id,
|
||||
"name": entry["name"],
|
||||
"description": entry["description"],
|
||||
"docs_url": entry["docs_url"],
|
||||
"enabled": enabled,
|
||||
"configured": configured,
|
||||
"gateway_running": gateway_running,
|
||||
"state": state,
|
||||
"error_code": error_code,
|
||||
"error_message": error_message,
|
||||
"updated_at": (
|
||||
runtime_platform.get("updated_at")
|
||||
if isinstance(runtime_platform, dict)
|
||||
else None
|
||||
),
|
||||
"home_channel": home_channel,
|
||||
"env_vars": env_vars,
|
||||
}
|
||||
if whatsapp_setup is not None:
|
||||
payload["whatsapp_setup"] = whatsapp_setup
|
||||
return payload
|
||||
|
||||
|
||||
_WHATSAPP_ONBOARDING_TTL_SECONDS = 600
|
||||
@@ -516,7 +927,6 @@ def _prune_telegram_onboarding_pairings() -> None:
|
||||
|
||||
|
||||
def _normalize_telegram_user_id(value: Any) -> str | None:
|
||||
from hermes_cli.web_server import _TELEGRAM_USER_ID_RE
|
||||
normalized = str(value or "").strip()
|
||||
if _TELEGRAM_USER_ID_RE.fullmatch(normalized):
|
||||
return normalized
|
||||
|
||||
@@ -26,6 +26,16 @@ run_in_threadpool = late("run_in_threadpool")
|
||||
save_config = late("save_config")
|
||||
|
||||
|
||||
_EMPTY_MODEL_INFO: dict = {
|
||||
"model": "",
|
||||
"provider": "",
|
||||
"auto_context_length": 0,
|
||||
"config_context_length": 0,
|
||||
"effective_context_length": 0,
|
||||
"capabilities": {},
|
||||
}
|
||||
|
||||
|
||||
@router.get("/api/model/info")
|
||||
def get_model_info(profile: Optional[str] = None):
|
||||
"""Return resolved model metadata for the currently configured model.
|
||||
@@ -34,7 +44,6 @@ def get_model_info(profile: Optional[str] = None):
|
||||
frontend can display "Auto-detected: 200K" alongside the override field.
|
||||
Also returns model capabilities (vision, reasoning, tools) when available.
|
||||
"""
|
||||
from hermes_cli.web_server import _EMPTY_MODEL_INFO
|
||||
try:
|
||||
with _profile_scope(profile):
|
||||
cfg = load_config()
|
||||
|
||||
@@ -13,6 +13,9 @@ from hermes_cli.web_deps import late
|
||||
from fastapi import HTTPException, Request
|
||||
from hermes_cli.web_models import OAuthSubmitBody
|
||||
from typing import Any, Dict, Optional
|
||||
import threading
|
||||
import os
|
||||
import secrets
|
||||
|
||||
_log = logging.getLogger("hermes_cli.web_server")
|
||||
router = APIRouter()
|
||||
@@ -23,8 +26,534 @@ _oauth_profile_name = late("_oauth_profile_name")
|
||||
_profile_scope = late("_profile_scope")
|
||||
_require_token = late("_require_token")
|
||||
_resolve_profile_dir = late("_resolve_profile_dir")
|
||||
_resolve_provider_status = late("_resolve_provider_status")
|
||||
_start_device_code_flow = late("_start_device_code_flow")
|
||||
_minimax_poller = late("_minimax_poller")
|
||||
_nous_poller = late("_nous_poller")
|
||||
_truncate_token = late("_truncate_token")
|
||||
_xai_device_poller = late("_xai_device_poller")
|
||||
|
||||
|
||||
def _http_response_error_detail(resp: Any) -> str:
|
||||
"""Best-effort extraction of a short provider error detail."""
|
||||
try:
|
||||
payload = resp.json()
|
||||
except Exception:
|
||||
payload = None
|
||||
if isinstance(payload, dict):
|
||||
error = payload.get("error")
|
||||
if isinstance(error, dict):
|
||||
parts = [
|
||||
str(error.get(key, "")).strip()
|
||||
for key in ("message", "error_description", "code", "type")
|
||||
if str(error.get(key, "")).strip()
|
||||
]
|
||||
if parts:
|
||||
return ": ".join(parts)
|
||||
if isinstance(error, str) and error.strip():
|
||||
return error.strip()
|
||||
for key in ("detail", "message", "error_description"):
|
||||
value = payload.get(key)
|
||||
if isinstance(value, str) and value.strip():
|
||||
return value.strip()
|
||||
text = str(getattr(resp, "text", "") or "").strip()
|
||||
return text[:500]
|
||||
|
||||
|
||||
def _codex_device_code_start_error(resp: Any) -> str:
|
||||
"""Dashboard-facing OpenAI Codex device-code start failure."""
|
||||
status = getattr(resp, "status_code", "unknown")
|
||||
detail = _http_response_error_detail(resp)
|
||||
lower = detail.lower()
|
||||
if "device" in lower and ("authori" in lower or "enable" in lower):
|
||||
message = (
|
||||
"OpenAI rejected the device-code login request. Your OpenAI "
|
||||
"account may need device-code authorization enabled before Hermes "
|
||||
"can start this dashboard login. Enable device-code authorization "
|
||||
"in OpenAI, then return here and click Login again."
|
||||
)
|
||||
else:
|
||||
message = (
|
||||
"OpenAI rejected the device-code login request. Please try Login "
|
||||
"again from the dashboard after checking your OpenAI account settings."
|
||||
)
|
||||
if detail:
|
||||
return f"{message} (HTTP {status}: {detail})"
|
||||
return f"{message} (HTTP {status})"
|
||||
|
||||
|
||||
def _new_oauth_session(
|
||||
provider_id: str,
|
||||
flow: str,
|
||||
profile: Optional[str] = None,
|
||||
) -> tuple[str, Dict[str, Any]]:
|
||||
"""Create + register a new OAuth session, return (session_id, session_dict)."""
|
||||
from hermes_cli.web_server import _oauth_sessions, _oauth_sessions_lock
|
||||
sid = secrets.token_urlsafe(16)
|
||||
profile_name = _oauth_profile_name(profile)
|
||||
sess = {
|
||||
"session_id": sid,
|
||||
"provider": provider_id,
|
||||
"flow": flow,
|
||||
"profile": profile_name,
|
||||
"created_at": time.time(),
|
||||
"status": "pending", # pending | approved | denied | expired | error
|
||||
"error_message": None,
|
||||
}
|
||||
with _oauth_sessions_lock:
|
||||
_oauth_sessions[sid] = sess
|
||||
return sid, sess
|
||||
|
||||
|
||||
def _codex_full_login_worker(session_id: str) -> None:
|
||||
"""Run the complete OpenAI Codex device-code flow.
|
||||
|
||||
Codex doesn't use the standard OAuth device-code endpoints; it has its
|
||||
own ``/api/accounts/deviceauth/usercode`` (JSON body, returns
|
||||
``device_auth_id``) and ``/api/accounts/deviceauth/token`` (JSON body
|
||||
polled until 200). On success the response carries an
|
||||
``authorization_code`` + ``code_verifier`` that get exchanged at
|
||||
CODEX_OAUTH_TOKEN_URL with grant_type=authorization_code.
|
||||
|
||||
The flow is replicated inline (rather than calling
|
||||
_codex_device_code_login) because that helper prints/blocks/polls in a
|
||||
single function — we need to surface the user_code to the dashboard the
|
||||
moment we receive it, well before polling completes.
|
||||
"""
|
||||
from hermes_cli.web_server import _oauth_sessions, _oauth_sessions_lock
|
||||
try:
|
||||
import httpx
|
||||
from hermes_cli.auth import (
|
||||
CODEX_OAUTH_CLIENT_ID,
|
||||
CODEX_OAUTH_TOKEN_URL,
|
||||
)
|
||||
issuer = "https://auth.openai.com"
|
||||
|
||||
# Step 1: request device code
|
||||
with httpx.Client(timeout=httpx.Timeout(15.0)) as client:
|
||||
resp = client.post(
|
||||
f"{issuer}/api/accounts/deviceauth/usercode",
|
||||
json={"client_id": CODEX_OAUTH_CLIENT_ID},
|
||||
headers={"Content-Type": "application/json"},
|
||||
)
|
||||
if resp.status_code != 200:
|
||||
raise RuntimeError(_codex_device_code_start_error(resp))
|
||||
device_data = resp.json()
|
||||
user_code = device_data.get("user_code", "")
|
||||
device_auth_id = device_data.get("device_auth_id", "")
|
||||
poll_interval = max(3, int(device_data.get("interval", "5")))
|
||||
if not user_code or not device_auth_id:
|
||||
raise RuntimeError("device-code response missing user_code or device_auth_id")
|
||||
verification_url = f"{issuer}/codex/device"
|
||||
with _oauth_sessions_lock:
|
||||
sess = _oauth_sessions.get(session_id)
|
||||
if not sess:
|
||||
return
|
||||
sess["user_code"] = user_code
|
||||
sess["verification_url"] = verification_url
|
||||
sess["device_auth_id"] = device_auth_id
|
||||
sess["interval"] = poll_interval
|
||||
sess["expires_in"] = 15 * 60 # OpenAI's effective limit
|
||||
sess["expires_at"] = time.time() + sess["expires_in"]
|
||||
# Captured now (not re-derived after cancel pops the session) so a
|
||||
# cancelled session can never fall back to the caller's current
|
||||
# profile scope at save time.
|
||||
session_profile = sess.get("profile")
|
||||
|
||||
# Step 2: poll until authorized
|
||||
deadline = time.monotonic() + sess["expires_in"]
|
||||
code_resp = None
|
||||
with httpx.Client(timeout=httpx.Timeout(15.0)) as client:
|
||||
while time.monotonic() < deadline:
|
||||
if sess.get("cancelled"):
|
||||
_log.info("oauth/device: openai-codex login cancelled (session=%s)", session_id)
|
||||
return
|
||||
time.sleep(poll_interval)
|
||||
if sess.get("cancelled"):
|
||||
_log.info("oauth/device: openai-codex login cancelled (session=%s)", session_id)
|
||||
return
|
||||
poll = client.post(
|
||||
f"{issuer}/api/accounts/deviceauth/token",
|
||||
json={"device_auth_id": device_auth_id, "user_code": user_code},
|
||||
headers={"Content-Type": "application/json"},
|
||||
)
|
||||
if poll.status_code == 200:
|
||||
code_resp = poll.json()
|
||||
break
|
||||
if poll.status_code in {403, 404}:
|
||||
continue # user hasn't authorized yet
|
||||
raise RuntimeError(f"deviceauth/token poll returned {poll.status_code}")
|
||||
|
||||
if code_resp is None:
|
||||
with _oauth_sessions_lock:
|
||||
sess["status"] = "expired"
|
||||
sess["error_message"] = "Device code expired before approval"
|
||||
return
|
||||
|
||||
if sess.get("cancelled"):
|
||||
_log.info("oauth/device: openai-codex login cancelled before token exchange (session=%s)", session_id)
|
||||
return
|
||||
|
||||
# Step 3: exchange authorization_code for tokens
|
||||
authorization_code = code_resp.get("authorization_code", "")
|
||||
code_verifier = code_resp.get("code_verifier", "")
|
||||
if not authorization_code or not code_verifier:
|
||||
raise RuntimeError("device-auth response missing authorization_code/code_verifier")
|
||||
with httpx.Client(timeout=httpx.Timeout(15.0)) as client:
|
||||
token_resp = client.post(
|
||||
CODEX_OAUTH_TOKEN_URL,
|
||||
data={
|
||||
"grant_type": "authorization_code",
|
||||
"code": authorization_code,
|
||||
"redirect_uri": f"{issuer}/deviceauth/callback",
|
||||
"client_id": CODEX_OAUTH_CLIENT_ID,
|
||||
"code_verifier": code_verifier,
|
||||
},
|
||||
headers={"Content-Type": "application/x-www-form-urlencoded"},
|
||||
)
|
||||
if token_resp.status_code != 200:
|
||||
raise RuntimeError(f"token exchange returned {token_resp.status_code}")
|
||||
tokens = token_resp.json()
|
||||
access_token = tokens.get("access_token", "")
|
||||
refresh_token = tokens.get("refresh_token", "")
|
||||
if not access_token:
|
||||
raise RuntimeError("token exchange did not return access_token")
|
||||
|
||||
from hermes_cli.auth import _save_codex_tokens
|
||||
|
||||
# The cancellation check and the save must be one atomic critical
|
||||
# section under the same lock cancel_oauth_session() uses. Checking
|
||||
# "cancelled" and then saving as two separate steps left a window
|
||||
# where DELETE could flip the flag between them and the worker would
|
||||
# still persist tokens after the user believed the login was
|
||||
# aborted. Holding the lock across both closes that window: DELETE
|
||||
# either lands before this section (worker observes cancelled and
|
||||
# returns) or blocks until this section (and the save) is done.
|
||||
with _oauth_sessions_lock:
|
||||
if sess.get("cancelled"):
|
||||
_log.info("oauth/device: openai-codex login cancelled before token save (session=%s)", session_id)
|
||||
return
|
||||
with _profile_scope(session_profile):
|
||||
_save_codex_tokens({
|
||||
"access_token": access_token,
|
||||
"refresh_token": refresh_token,
|
||||
})
|
||||
sess["status"] = "approved"
|
||||
_log.info("oauth/device: openai-codex login completed (session=%s)", session_id)
|
||||
except Exception as e:
|
||||
_log.warning("codex device-code worker failed (session=%s): %s", session_id, e)
|
||||
with _oauth_sessions_lock:
|
||||
s = _oauth_sessions.get(session_id)
|
||||
if s:
|
||||
s["status"] = "error"
|
||||
s["error_message"] = str(e)
|
||||
|
||||
|
||||
def _resolve_provider_status(provider_id: str, status_fn) -> Dict[str, Any]:
|
||||
"""Dispatch to the right status helper for an OAuth provider entry."""
|
||||
if status_fn is not None:
|
||||
try:
|
||||
return status_fn()
|
||||
except Exception as e:
|
||||
return {"logged_in": False, "error": str(e)}
|
||||
try:
|
||||
from hermes_cli import auth as hauth
|
||||
if provider_id == "nous":
|
||||
# Read-only accounts-tab card: refresh-free snapshot so listing
|
||||
# providers never performs an OAuth refresh.
|
||||
raw = hauth.get_nous_auth_status_local()
|
||||
return {
|
||||
"logged_in": bool(raw.get("logged_in")),
|
||||
"source": "nous_portal",
|
||||
"source_label": raw.get("portal_base_url") or "Nous Portal",
|
||||
"token_preview": _truncate_token(raw.get("access_token")),
|
||||
"expires_at": raw.get("access_expires_at"),
|
||||
"has_refresh_token": bool(raw.get("has_refresh_token")),
|
||||
}
|
||||
if provider_id == "openai-codex":
|
||||
raw = hauth.get_codex_auth_status()
|
||||
return {
|
||||
"logged_in": bool(raw.get("logged_in")),
|
||||
"source": raw.get("source") or "openai_codex",
|
||||
"source_label": raw.get("auth_mode") or "OpenAI Codex",
|
||||
"token_preview": _truncate_token(raw.get("api_key")),
|
||||
"expires_at": None,
|
||||
"has_refresh_token": False,
|
||||
"last_refresh": raw.get("last_refresh"),
|
||||
}
|
||||
if provider_id == "qwen-oauth":
|
||||
raw = hauth.get_qwen_auth_status()
|
||||
return {
|
||||
"logged_in": bool(raw.get("logged_in")),
|
||||
"source": "qwen_cli",
|
||||
"source_label": raw.get("auth_store_path") or "Qwen CLI",
|
||||
"token_preview": _truncate_token(raw.get("access_token")),
|
||||
"expires_at": raw.get("expires_at"),
|
||||
"has_refresh_token": bool(raw.get("has_refresh_token")),
|
||||
}
|
||||
if provider_id == "minimax-oauth":
|
||||
raw = hauth.get_minimax_oauth_auth_status()
|
||||
return {
|
||||
"logged_in": bool(raw.get("logged_in")),
|
||||
"source": "minimax_oauth",
|
||||
"source_label": f"MiniMax ({raw.get('region', 'global')})",
|
||||
"token_preview": None,
|
||||
"expires_at": raw.get("expires_at"),
|
||||
"has_refresh_token": True,
|
||||
}
|
||||
if provider_id == "xai-oauth":
|
||||
raw = hauth.get_xai_oauth_auth_status()
|
||||
# source_label is meant to be a human-readable origin (auth-store
|
||||
# path / credential source), not the internal auth_mode string
|
||||
# ("oauth_pkce"). Prefer the store path, then the source slug.
|
||||
return {
|
||||
"logged_in": bool(raw.get("logged_in")),
|
||||
"source": raw.get("source") or "xai_oauth",
|
||||
"source_label": raw.get("auth_store") or raw.get("source") or "xAI Grok OAuth",
|
||||
"token_preview": _truncate_token(raw.get("api_key")),
|
||||
"expires_at": None,
|
||||
"has_refresh_token": True,
|
||||
"last_refresh": raw.get("last_refresh"),
|
||||
}
|
||||
# No hand-written branch for this provider id: fall through to the
|
||||
# canonical slug-driven dispatcher so accounts-tab providers derived
|
||||
# from the unified catalog (which carry status_fn=None) still reflect
|
||||
# real login state instead of rendering permanently logged-out. This
|
||||
# closes the membership-auto-extends-but-status-doesn't gap: add an
|
||||
# OAuth/account provider plugin and its card shows the right state.
|
||||
raw = hauth.get_auth_status(provider_id)
|
||||
if isinstance(raw, dict) and "logged_in" in raw:
|
||||
return {
|
||||
"logged_in": bool(raw.get("logged_in")),
|
||||
"source": raw.get("source") or raw.get("provider") or provider_id,
|
||||
"source_label": (
|
||||
raw.get("source_label")
|
||||
or raw.get("auth_store")
|
||||
or raw.get("auth_store_path")
|
||||
or raw.get("base_url")
|
||||
or raw.get("name")
|
||||
or ""
|
||||
),
|
||||
"token_preview": _truncate_token(
|
||||
raw.get("access_token") or raw.get("api_key")
|
||||
),
|
||||
"expires_at": raw.get("expires_at") or raw.get("access_expires_at"),
|
||||
"has_refresh_token": bool(raw.get("has_refresh_token")),
|
||||
}
|
||||
except Exception as e:
|
||||
return {"logged_in": False, "error": str(e)}
|
||||
return {"logged_in": False}
|
||||
|
||||
|
||||
async def _start_device_code_flow(
|
||||
provider_id: str,
|
||||
profile: Optional[str] = None,
|
||||
) -> Dict[str, Any]:
|
||||
"""Initiate a device-code flow (Nous, OpenAI Codex, MiniMax, or xAI).
|
||||
|
||||
Calls the provider's device-auth endpoint via the existing CLI helpers,
|
||||
then spawns a background poller. Returns the user-facing display fields
|
||||
so the UI can render the verification page link + user code.
|
||||
"""
|
||||
from hermes_cli.web_server import _oauth_sessions, _oauth_sessions_lock
|
||||
if provider_id == "nous":
|
||||
from hermes_cli.auth import (
|
||||
_request_device_code,
|
||||
PROVIDER_REGISTRY,
|
||||
)
|
||||
import httpx
|
||||
pconfig = PROVIDER_REGISTRY["nous"]
|
||||
portal_base_url = (
|
||||
os.getenv("HERMES_PORTAL_BASE_URL")
|
||||
or os.getenv("NOUS_PORTAL_BASE_URL")
|
||||
or pconfig.portal_base_url
|
||||
).rstrip("/")
|
||||
client_id = pconfig.client_id
|
||||
scope = pconfig.scope
|
||||
|
||||
def _do_nous_device_request():
|
||||
with httpx.Client(
|
||||
timeout=httpx.Timeout(15.0),
|
||||
headers={"Accept": "application/json"},
|
||||
) as client:
|
||||
return (
|
||||
_request_device_code(
|
||||
client=client,
|
||||
portal_base_url=portal_base_url,
|
||||
client_id=client_id,
|
||||
scope=scope,
|
||||
),
|
||||
scope,
|
||||
)
|
||||
|
||||
device_data, effective_scope = await asyncio.get_running_loop().run_in_executor(
|
||||
None, _do_nous_device_request
|
||||
)
|
||||
sid, sess = _new_oauth_session("nous", "device_code", profile=profile)
|
||||
sess["device_code"] = str(device_data["device_code"])
|
||||
sess["interval"] = int(device_data["interval"])
|
||||
sess["expires_at"] = time.time() + int(device_data["expires_in"])
|
||||
sess["portal_base_url"] = portal_base_url
|
||||
sess["client_id"] = client_id
|
||||
sess["scope"] = effective_scope
|
||||
threading.Thread(
|
||||
target=_nous_poller, args=(sid,), daemon=True, name=f"oauth-poll-{sid[:6]}"
|
||||
).start()
|
||||
return {
|
||||
"session_id": sid,
|
||||
"flow": "device_code",
|
||||
"user_code": str(device_data["user_code"]),
|
||||
"verification_url": str(device_data["verification_uri_complete"]),
|
||||
"expires_in": int(device_data["expires_in"]),
|
||||
"poll_interval": int(device_data["interval"]),
|
||||
}
|
||||
|
||||
if provider_id == "openai-codex":
|
||||
# Codex uses fixed OpenAI device-auth endpoints; reuse the helper.
|
||||
sid, _ = _new_oauth_session("openai-codex", "device_code", profile=profile)
|
||||
# Use the helper but in a thread because it polls inline.
|
||||
# We can't extract just the start step without refactoring auth.py,
|
||||
# so we run the full helper in a worker and proxy the user_code +
|
||||
# verification_url back via the session dict. The helper prints
|
||||
# to stdout — we capture nothing here, just status.
|
||||
threading.Thread(
|
||||
target=_codex_full_login_worker, args=(sid,), daemon=True,
|
||||
name=f"oauth-codex-{sid[:6]}",
|
||||
).start()
|
||||
# Block briefly until the worker has populated the user_code, OR error.
|
||||
deadline = time.monotonic() + 10
|
||||
while time.monotonic() < deadline:
|
||||
with _oauth_sessions_lock:
|
||||
s = _oauth_sessions.get(sid)
|
||||
if s and (s.get("user_code") or s["status"] != "pending"):
|
||||
break
|
||||
await asyncio.sleep(0.1)
|
||||
with _oauth_sessions_lock:
|
||||
s = _oauth_sessions.get(sid, {})
|
||||
if s.get("status") == "error":
|
||||
raise HTTPException(status_code=500, detail=s.get("error_message") or "device-auth failed")
|
||||
if not s.get("user_code"):
|
||||
raise HTTPException(status_code=504, detail="device-auth timed out before returning a user code")
|
||||
return {
|
||||
"session_id": sid,
|
||||
"flow": "device_code",
|
||||
"user_code": s["user_code"],
|
||||
"verification_url": s["verification_url"],
|
||||
"expires_in": int(s.get("expires_in") or 900),
|
||||
"poll_interval": int(s.get("interval") or 5),
|
||||
}
|
||||
|
||||
if provider_id == "minimax-oauth":
|
||||
# MiniMax uses a device-code-style flow (verification URI + user
|
||||
# code + background poll) with a PKCE extension on top. From the
|
||||
# operator's perspective it's identical to Nous's device-code
|
||||
# flow; the PKCE bit (verifier + challenge from
|
||||
# _minimax_pkce_pair) is a security extension that binds the
|
||||
# token exchange to the original session.
|
||||
from hermes_cli.auth import (
|
||||
_minimax_pkce_pair,
|
||||
_minimax_request_user_code,
|
||||
MINIMAX_OAUTH_CLIENT_ID,
|
||||
MINIMAX_OAUTH_GLOBAL_BASE,
|
||||
)
|
||||
import httpx
|
||||
verifier, challenge, state = _minimax_pkce_pair()
|
||||
portal_base_url = (
|
||||
os.getenv("MINIMAX_PORTAL_BASE_URL") or MINIMAX_OAUTH_GLOBAL_BASE
|
||||
).rstrip("/")
|
||||
def _do_minimax_request():
|
||||
with httpx.Client(
|
||||
timeout=httpx.Timeout(15.0),
|
||||
headers={"Accept": "application/json"},
|
||||
follow_redirects=True,
|
||||
) as client:
|
||||
return _minimax_request_user_code(
|
||||
client=client,
|
||||
portal_base_url=portal_base_url,
|
||||
client_id=MINIMAX_OAUTH_CLIENT_ID,
|
||||
code_challenge=challenge,
|
||||
state=state,
|
||||
)
|
||||
device_data = await asyncio.get_event_loop().run_in_executor(
|
||||
None, _do_minimax_request
|
||||
)
|
||||
sid, sess = _new_oauth_session("minimax-oauth", "device_code", profile=profile)
|
||||
# The CLI flow names this `interval_ms` because MiniMax's
|
||||
# `interval` field is in milliseconds (defensive default 2000ms
|
||||
# in _minimax_poll_token).
|
||||
interval_raw = device_data.get("interval")
|
||||
sess["interval_ms"] = (
|
||||
int(interval_raw) if interval_raw is not None else None
|
||||
)
|
||||
sess["user_code"] = str(device_data["user_code"])
|
||||
sess["code_verifier"] = verifier
|
||||
sess["state"] = state
|
||||
sess["portal_base_url"] = portal_base_url
|
||||
sess["client_id"] = MINIMAX_OAUTH_CLIENT_ID
|
||||
sess["region"] = "global"
|
||||
# `expired_in` from MiniMax is overloaded — could be a unix-ms
|
||||
# timestamp OR a seconds-from-now duration. Mirror the heuristic
|
||||
# in _minimax_poll_token. Stash the raw value for the poller;
|
||||
# compute a derived expires_at + UI-friendly expires_in seconds.
|
||||
expired_in_raw = int(device_data["expired_in"])
|
||||
sess["expired_in_raw"] = expired_in_raw
|
||||
if expired_in_raw > 1_000_000_000_000: # likely unix-ms
|
||||
expires_at_ts = expired_in_raw / 1000.0
|
||||
expires_in_seconds = max(0, int(expires_at_ts - time.time()))
|
||||
else:
|
||||
expires_at_ts = time.time() + expired_in_raw
|
||||
expires_in_seconds = expired_in_raw
|
||||
sess["expires_at"] = expires_at_ts
|
||||
threading.Thread(
|
||||
target=_minimax_poller,
|
||||
args=(sid,),
|
||||
daemon=True,
|
||||
name=f"oauth-poll-{sid[:6]}",
|
||||
).start()
|
||||
return {
|
||||
"session_id": sid,
|
||||
"flow": "device_code",
|
||||
"user_code": str(device_data["user_code"]),
|
||||
"verification_url": str(device_data["verification_uri"]),
|
||||
"expires_in": expires_in_seconds,
|
||||
"poll_interval": max(2, (sess["interval_ms"] or 2000) // 1000),
|
||||
}
|
||||
|
||||
if provider_id == "xai-oauth":
|
||||
from hermes_cli.auth import _xai_oauth_request_device_code
|
||||
import httpx
|
||||
|
||||
def _do_xai_device_request():
|
||||
with httpx.Client(
|
||||
timeout=httpx.Timeout(20.0),
|
||||
headers={"Accept": "application/json"},
|
||||
) as client:
|
||||
return _xai_oauth_request_device_code(client)
|
||||
|
||||
device_data = await asyncio.get_running_loop().run_in_executor(
|
||||
None, _do_xai_device_request
|
||||
)
|
||||
sid, sess = _new_oauth_session("xai-oauth", "device_code", profile=profile)
|
||||
sess["device_code"] = str(device_data["device_code"])
|
||||
sess["interval"] = int(device_data["interval"])
|
||||
sess["expires_at"] = time.time() + int(device_data["expires_in"])
|
||||
threading.Thread(
|
||||
target=_xai_device_poller,
|
||||
args=(sid,),
|
||||
daemon=True,
|
||||
name=f"oauth-poll-{sid[:6]}",
|
||||
).start()
|
||||
return {
|
||||
"session_id": sid,
|
||||
"flow": "device_code",
|
||||
"user_code": str(device_data["user_code"]),
|
||||
"verification_url": str(
|
||||
device_data.get("verification_uri_complete")
|
||||
or device_data["verification_uri"]
|
||||
),
|
||||
"expires_in": int(device_data["expires_in"]),
|
||||
"poll_interval": int(device_data["interval"]),
|
||||
}
|
||||
|
||||
raise HTTPException(status_code=400, detail=f"Provider {provider_id} does not support device-code flow")
|
||||
|
||||
|
||||
def _oauth_provider_disconnect_command(provider: Dict[str, Any]) -> Optional[str]:
|
||||
|
||||
@@ -32,12 +32,16 @@ _normalize_memory_provider_name = late("_normalize_memory_provider_name")
|
||||
_path_is_under = late("_path_is_under")
|
||||
_require_memory_provider_ready = late("_require_memory_provider_ready")
|
||||
_resolve_profile_dir = late("_resolve_profile_dir")
|
||||
_restart_gateway_after_webhook_enable = late("_restart_gateway_after_webhook_enable")
|
||||
_spawn_hermes_action = late("_spawn_hermes_action")
|
||||
_write_platform_enabled = late("_write_platform_enabled")
|
||||
get_hermes_home = late("get_hermes_home")
|
||||
load_config = late("load_config")
|
||||
save_config = late("save_config")
|
||||
_restart_gateway_after = late("_restart_gateway_after")
|
||||
|
||||
|
||||
def _restart_gateway_after_webhook_enable(profile: Optional[str] = None) -> dict[str, Any]:
|
||||
return _restart_gateway_after(profile, what="enabling webhooks", label="Webhook enable")
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
@@ -70,18 +70,113 @@ router = APIRouter()
|
||||
# Late-bound web_server helpers (resolved at call time; cycle-safe,
|
||||
# monkeypatch-transparent).
|
||||
_cron_profile_home = late("_cron_profile_home")
|
||||
_disable_unselected_skills = late("_disable_unselected_skills")
|
||||
_fallback_profile_dicts = late("_fallback_profile_dicts")
|
||||
_hub_action_name = late("_hub_action_name")
|
||||
_open_session_db_at_path = late("_open_session_db_at_path")
|
||||
_profile_setup_command = late("_profile_setup_command")
|
||||
_profile_to_dict = late("_profile_to_dict")
|
||||
_resolve_profile_dir = late("_resolve_profile_dir")
|
||||
_spawn_hermes_action = late("_spawn_hermes_action")
|
||||
run_in_threadpool = late("run_in_threadpool")
|
||||
_strip_session_list_rows = late("_strip_session_list_rows")
|
||||
_write_profile_mcp_servers = late("_write_profile_mcp_servers")
|
||||
_write_profile_model = late("_write_profile_model")
|
||||
_apply_main_model_assignment = late("_apply_main_model_assignment")
|
||||
_normalize_main_model_assignment = late("_normalize_main_model_assignment")
|
||||
|
||||
|
||||
# ---------------------------------------------------------------------------
|
||||
# Profile management endpoints (minimal — list/create/rename/delete + SOUL.md)
|
||||
# ---------------------------------------------------------------------------
|
||||
|
||||
|
||||
def _profile_attr(info, name: str, default: Any = None) -> Any:
|
||||
try:
|
||||
return getattr(info, name)
|
||||
except Exception:
|
||||
return default
|
||||
|
||||
|
||||
def _profile_to_dict(info) -> Dict[str, Any]:
|
||||
return {
|
||||
"name": _profile_attr(info, "name", ""),
|
||||
"path": str(_profile_attr(info, "path", "")),
|
||||
"is_default": bool(_profile_attr(info, "is_default", False)),
|
||||
"model": _profile_attr(info, "model"),
|
||||
"provider": _profile_attr(info, "provider"),
|
||||
"has_env": bool(_profile_attr(info, "has_env", False)),
|
||||
"skill_count": int(_profile_attr(info, "skill_count", 0) or 0),
|
||||
"gateway_running": bool(_profile_attr(info, "gateway_running", False)),
|
||||
"description": _profile_attr(info, "description", "") or "",
|
||||
"description_auto": bool(_profile_attr(info, "description_auto", False)),
|
||||
"display_name": _profile_attr(info, "display_name", "") or "",
|
||||
"distribution_name": _profile_attr(info, "distribution_name"),
|
||||
"distribution_version": _profile_attr(info, "distribution_version"),
|
||||
"distribution_source": _profile_attr(info, "distribution_source"),
|
||||
"has_alias": _profile_attr(info, "alias_path") is not None,
|
||||
}
|
||||
|
||||
|
||||
def _profile_setup_command(name: str) -> str:
|
||||
"""Return the shell command used to configure a profile in the CLI."""
|
||||
_resolve_profile_dir(name)
|
||||
return "hermes setup" if name == "default" else f"{name} setup"
|
||||
|
||||
|
||||
def _write_profile_model(profile_dir: Path, provider: str, model: str) -> None:
|
||||
"""Write the main model assignment into a specific profile's config.yaml.
|
||||
|
||||
Scopes ``load_config``/``save_config`` to ``profile_dir`` via the
|
||||
context-local HERMES_HOME override so the write lands in the target
|
||||
profile's config rather than the dashboard process's active profile.
|
||||
Clears any stale ``base_url`` / ``context_length`` the same way
|
||||
``POST /api/model/set`` does, since the new model may differ.
|
||||
"""
|
||||
from hermes_cli.web_server import load_config, save_config
|
||||
from hermes_constants import set_hermes_home_override, reset_hermes_home_override
|
||||
|
||||
token = set_hermes_home_override(str(profile_dir))
|
||||
try:
|
||||
provider, model = _normalize_main_model_assignment(provider, model)
|
||||
cfg = load_config()
|
||||
cfg["model"] = _apply_main_model_assignment(cfg.get("model", {}), provider, model)
|
||||
save_config(cfg)
|
||||
finally:
|
||||
reset_hermes_home_override(token)
|
||||
|
||||
|
||||
def _disable_unselected_skills(profile_dir: Path, keep: List[str]) -> int:
|
||||
"""Disable every installed skill in ``profile_dir`` not in ``keep``.
|
||||
|
||||
Profiles manage skill activation via a *disabled* list — all installed
|
||||
skills are active by default and users opt out. The builder's skill step
|
||||
uses "replace" semantics: the user picks exactly which seeded built-in /
|
||||
optional skills stay active, and everything else gets added to the disabled
|
||||
list. (Hub skills are installed separately via subprocess and are active on
|
||||
install.) Scoped to the profile via the HERMES_HOME override. Returns the
|
||||
number of skills newly disabled.
|
||||
"""
|
||||
from hermes_cli.web_server import load_config
|
||||
from hermes_constants import set_hermes_home_override, reset_hermes_home_override
|
||||
from hermes_cli.skills_config import get_disabled_skills, save_disabled_skills
|
||||
|
||||
keep_set = {s.strip() for s in keep if s and s.strip()}
|
||||
disabled_count = 0
|
||||
token = set_hermes_home_override(str(profile_dir))
|
||||
try:
|
||||
installed: List[str] = []
|
||||
skills_root = profile_dir / "skills"
|
||||
if skills_root.is_dir():
|
||||
for md in skills_root.rglob("SKILL.md"):
|
||||
installed.append(md.parent.name)
|
||||
cfg = load_config()
|
||||
disabled = get_disabled_skills(cfg)
|
||||
for name in installed:
|
||||
if name not in keep_set and name not in disabled:
|
||||
disabled.add(name)
|
||||
disabled_count += 1
|
||||
if disabled_count:
|
||||
save_disabled_skills(cfg, disabled)
|
||||
finally:
|
||||
reset_hermes_home_override(token)
|
||||
return disabled_count
|
||||
|
||||
# Returned by the offloaded file readers below to mean "the file is not there",
|
||||
# which a plain ``None`` cannot express: ``desktop.json`` may legitimately hold
|
||||
|
||||
@@ -29,6 +29,7 @@ from hermes_cli.web_models import (
|
||||
)
|
||||
from hermes_cli.web_routers._common import log as _log, http_failure
|
||||
from hermes_state import is_malformed_db_error, is_transient_sqlite_error
|
||||
from typing import Any, Dict
|
||||
|
||||
list_router = APIRouter()
|
||||
search_router = APIRouter()
|
||||
@@ -36,14 +37,136 @@ manage_router = APIRouter()
|
||||
|
||||
_cron_default_profile = late("_cron_default_profile")
|
||||
_cron_profile_home = late("_cron_profile_home")
|
||||
_import_sessions_for_profile = late("_import_sessions_for_profile")
|
||||
_maybe_auto_archive_for_profile = late("_maybe_auto_archive_for_profile")
|
||||
_open_session_db_for_profile = late("_open_session_db_for_profile")
|
||||
_prune_sessions = late("_prune_sessions")
|
||||
_read_session_import_body = late("_read_session_import_body")
|
||||
_session_latest_descendant = late("_session_latest_descendant")
|
||||
_strip_session_list_rows = late("_strip_session_list_rows")
|
||||
|
||||
|
||||
# CRITICAL — every literal-path route below MUST be declared BEFORE the
|
||||
# templated ``/api/sessions/{session_id}`` family that follows. FastAPI/
|
||||
# Starlette match routes in registration order, and the ``{session_id}``
|
||||
# pattern is unconstrained — it would otherwise swallow e.g.
|
||||
# ``DELETE /api/sessions/empty``, ``POST /api/sessions/bulk-delete``, or
|
||||
# ``GET /api/sessions/stats`` as "operate on the session with id
|
||||
# 'empty'" / "'bulk-delete'" / "'stats'", which would 404 (or worse,
|
||||
# succeed and delete the wrong row). Same story as the older
|
||||
# ``/api/sessions/search`` endpoint up at line ~1191. If you split or
|
||||
# reorder this block, move every route in it together.
|
||||
# Keep the dashboard import endpoint stream-safe: FastAPI otherwise parses and
|
||||
# buffers an arbitrarily large JSON body before SessionDB can enforce its own
|
||||
# per-session and transaction-work limits.
|
||||
_SESSION_IMPORT_MAX_BYTES = 25 * 1024 * 1024
|
||||
|
||||
|
||||
async def _read_session_import_body(request: Request) -> bytes:
|
||||
body = bytearray()
|
||||
async for chunk in request.stream():
|
||||
if len(body) + len(chunk) > _SESSION_IMPORT_MAX_BYTES:
|
||||
raise HTTPException(status_code=413, detail="Session import payload is too large")
|
||||
body.extend(chunk)
|
||||
return bytes(body)
|
||||
|
||||
|
||||
def _import_sessions_for_profile(profile: Optional[str], sessions: List[Dict[str, Any]]) -> Dict[str, Any]:
|
||||
db = _open_session_db_for_profile(profile, read_only=False)
|
||||
try:
|
||||
return db.import_sessions(sessions)
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
def _prune_sessions(body: SessionPrune):
|
||||
"""Delete ended sessions matching filters (mirrors `hermes sessions prune`)."""
|
||||
from hermes_cli.web_server import get_hermes_home
|
||||
has_window = (
|
||||
body.started_before is not None or body.started_after is not None
|
||||
)
|
||||
if body.older_than_days is not None and body.older_than_days < 1 and not has_window:
|
||||
raise HTTPException(status_code=400, detail="older_than_days must be >= 1")
|
||||
# Mirror the CLI: the implicit 90-day cutoff only applies to a truly bare
|
||||
# prune. Any attribute filter (source, title, model, ...) suppresses it
|
||||
# unless the caller explicitly sent older_than_days.
|
||||
_attr_filters_set = any(
|
||||
getattr(body, f) is not None
|
||||
for f in (
|
||||
"source", "title_like", "end_reason", "cwd_prefix",
|
||||
"min_messages", "max_messages", "model_like", "provider",
|
||||
"user_id", "chat_id", "chat_type", "branch_like",
|
||||
"min_tokens", "max_tokens", "min_cost", "max_cost",
|
||||
"min_tool_calls", "max_tool_calls",
|
||||
)
|
||||
)
|
||||
_older_than_explicit = "older_than_days" in body.model_fields_set
|
||||
_effective_older_than = body.older_than_days
|
||||
if has_window or (_attr_filters_set and not _older_than_explicit):
|
||||
_effective_older_than = None
|
||||
profile_home = _cron_profile_home(body.profile)[1] if body.profile else get_hermes_home()
|
||||
db = _open_session_db_for_profile(body.profile, read_only=False)
|
||||
try:
|
||||
filters = dict(
|
||||
older_than_days=_effective_older_than,
|
||||
source=(body.source or None),
|
||||
started_before=body.started_before,
|
||||
started_after=body.started_after,
|
||||
title_like=(body.title_like or None),
|
||||
end_reason=(body.end_reason or None),
|
||||
cwd_prefix=(body.cwd_prefix or None),
|
||||
min_messages=body.min_messages,
|
||||
max_messages=body.max_messages,
|
||||
model_like=(body.model_like or None),
|
||||
provider=(body.provider or None),
|
||||
user_id=(body.user_id or None),
|
||||
chat_id=(body.chat_id or None),
|
||||
chat_type=(body.chat_type or None),
|
||||
branch_like=(body.branch_like or None),
|
||||
min_tokens=body.min_tokens,
|
||||
max_tokens=body.max_tokens,
|
||||
min_cost=body.min_cost,
|
||||
max_cost=body.max_cost,
|
||||
min_tool_calls=body.min_tool_calls,
|
||||
max_tool_calls=body.max_tool_calls,
|
||||
archived=None if body.include_archived else False,
|
||||
)
|
||||
skipped_open = db.count_open_prune_matches(**filters)
|
||||
if body.dry_run:
|
||||
rows = db.list_prune_candidates(**filters)
|
||||
return {
|
||||
"ok": True,
|
||||
"removed": 0,
|
||||
"matched": len(rows),
|
||||
"skipped_open": skipped_open,
|
||||
# Rows are ordered by last activity, not creation time.
|
||||
"oldest_last_active": rows[0]["last_active"] if rows else None,
|
||||
"newest_last_active": rows[-1]["last_active"] if rows else None,
|
||||
"oldest_started_at": (
|
||||
min(r["started_at"] for r in rows) if rows else None
|
||||
),
|
||||
"newest_started_at": (
|
||||
max(r["started_at"] for r in rows) if rows else None
|
||||
),
|
||||
"sessions": [
|
||||
{
|
||||
"id": r["id"],
|
||||
"source": r["source"],
|
||||
"title": r.get("title"),
|
||||
"model": r.get("model"),
|
||||
"started_at": r["started_at"],
|
||||
"last_active": r["last_active"],
|
||||
"message_count": r["message_count"],
|
||||
}
|
||||
for r in rows
|
||||
],
|
||||
}
|
||||
sessions_dir = profile_home / "sessions"
|
||||
removed = db.prune_sessions(
|
||||
sessions_dir=sessions_dir if sessions_dir.exists() else None,
|
||||
**filters,
|
||||
)
|
||||
return {"ok": True, "removed": removed, "skipped_open": skipped_open}
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
_ACTIVE_WINDOW_S = 300
|
||||
|
||||
|
||||
|
||||
@@ -12,7 +12,7 @@ from typing import Optional
|
||||
|
||||
from fastapi import APIRouter, HTTPException
|
||||
|
||||
from hermes_cli.web_deps import late, LateState
|
||||
from hermes_cli.web_deps import late
|
||||
from hermes_cli.web_models import (
|
||||
SkillContentUpdate,
|
||||
SkillCreate,
|
||||
@@ -34,13 +34,50 @@ from hermes_cli.web_routers._common import (
|
||||
hub_router = APIRouter()
|
||||
router = APIRouter()
|
||||
|
||||
_clear_skills_prompt_cache = late("_clear_skills_prompt_cache")
|
||||
_config_profile_scope = late("_config_profile_scope")
|
||||
_hub_action_name = late("_hub_action_name")
|
||||
_installed_hub_identifiers = late("_installed_hub_identifiers")
|
||||
_skill_meta_to_payload = late("_skill_meta_to_payload")
|
||||
load_config = late("load_config")
|
||||
_SKILL_HUB_SOURCE_LABELS = LateState("_SKILL_HUB_SOURCE_LABELS")
|
||||
|
||||
|
||||
# Human-readable labels for each hub source id (matches `hermes skills search`
|
||||
# provenance). Keep in sync with create_source_router()'s source list.
|
||||
_SKILL_HUB_SOURCE_LABELS = {
|
||||
"official": "Official (Nous)",
|
||||
"hermes-index": "Hermes Index",
|
||||
"skills-sh": "skills.sh",
|
||||
"well-known": "Well-Known",
|
||||
"url": "Direct URL",
|
||||
"github": "GitHub",
|
||||
"clawhub": "ClawHub",
|
||||
"lobehub": "LobeHub",
|
||||
"browse-sh": "browse.sh",
|
||||
}
|
||||
|
||||
|
||||
def _skill_meta_to_payload(m) -> dict:
|
||||
return {
|
||||
"name": m.name,
|
||||
"description": m.description,
|
||||
"source": m.source,
|
||||
"identifier": m.identifier,
|
||||
"trust_level": m.trust_level,
|
||||
"repo": m.repo,
|
||||
"tags": list(m.tags or []),
|
||||
}
|
||||
|
||||
|
||||
def _clear_skills_prompt_cache() -> None:
|
||||
"""Best-effort: invalidate the skills system-prompt snapshot after a write.
|
||||
|
||||
Mirrors what ``skill_manage`` does so a dashboard-authored skill is picked
|
||||
up by the next session without a manual cache reset.
|
||||
"""
|
||||
try:
|
||||
from agent.prompt_builder import clear_skills_system_prompt_cache
|
||||
clear_skills_system_prompt_cache(clear_snapshot=True)
|
||||
except Exception:
|
||||
pass
|
||||
|
||||
|
||||
@hub_router.post("/api/skills/hub/install")
|
||||
|
||||
@@ -39,7 +39,6 @@ _resolve_profile_dir = late("_resolve_profile_dir")
|
||||
_resolve_restart_drain_timeout = late("_resolve_restart_drain_timeout")
|
||||
_spawn_hermes_action = late("_spawn_hermes_action")
|
||||
_ssh_runtime_intact = late("_ssh_runtime_intact")
|
||||
_status_active_sessions = late("_status_active_sessions")
|
||||
app = LateState("app") # the FastAPI instance (app.state.*)
|
||||
check_config_version = late("check_config_version")
|
||||
get_hermes_home = late("get_hermes_home")
|
||||
@@ -49,6 +48,58 @@ get_runtime_status_running_pid = late("get_runtime_status_running_pid")
|
||||
load_config = late("load_config")
|
||||
read_runtime_status = late("read_runtime_status")
|
||||
run_in_threadpool = late("run_in_threadpool")
|
||||
_open_session_db_for_profile = late("_open_session_db_for_profile")
|
||||
|
||||
|
||||
_STATUS_ACTIVE_SESSIONS_TIMEOUT = 0.75
|
||||
|
||||
|
||||
def _count_status_active_sessions() -> int:
|
||||
"""Return the dashboard status active-session count.
|
||||
|
||||
This is best-effort status garnish, not a critical path. Opens read-only
|
||||
(via the shared stale-schema heal, same as every other dashboard read
|
||||
path) so /api/status never routinely writes to state.db while another
|
||||
Hermes process is using it.
|
||||
"""
|
||||
from hermes_state import _default_db_path
|
||||
|
||||
# The heal helper bootstraps a missing store; this garnish must not — on
|
||||
# a fresh install /api/status polls would otherwise create state.db
|
||||
# before the user's first session.
|
||||
if not Path(_default_db_path()).exists():
|
||||
return 0
|
||||
|
||||
db = _open_session_db_for_profile(None, read_only=True)
|
||||
try:
|
||||
sessions = db.list_sessions_rich(limit=50, compact_rows=True)
|
||||
now = time.time()
|
||||
return sum(
|
||||
1 for s in sessions
|
||||
if s.get("ended_at") is None
|
||||
and (now - s.get("last_active", s.get("started_at", 0))) < 300
|
||||
)
|
||||
finally:
|
||||
db.close()
|
||||
|
||||
|
||||
_GATEWAY_HEALTH_ROUTE_TIMEOUT = 1.0
|
||||
|
||||
|
||||
async def _status_active_sessions() -> int:
|
||||
try:
|
||||
return await asyncio.wait_for(
|
||||
run_in_threadpool(_count_status_active_sessions),
|
||||
timeout=_STATUS_ACTIVE_SESSIONS_TIMEOUT,
|
||||
)
|
||||
except asyncio.TimeoutError:
|
||||
_log.debug(
|
||||
"/api/status active session count exceeded %.2fs; returning 0",
|
||||
_STATUS_ACTIVE_SESSIONS_TIMEOUT,
|
||||
)
|
||||
except Exception as exc:
|
||||
_log.debug("/api/status active session count unavailable: %s", exc)
|
||||
return 0
|
||||
|
||||
|
||||
@router.get("/api/ssh/ownership")
|
||||
@@ -158,11 +209,7 @@ def _merge_profile_gateway_platforms(
|
||||
|
||||
@router.get("/api/status")
|
||||
async def get_status(profile: Optional[str] = None):
|
||||
from hermes_cli.web_server import (
|
||||
DASHBOARD_HEALTH,
|
||||
_GATEWAY_HEALTH_ROUTE_TIMEOUT,
|
||||
_GATEWAY_HEALTH_URL,
|
||||
)
|
||||
from hermes_cli.web_server import DASHBOARD_HEALTH, _GATEWAY_HEALTH_URL
|
||||
status_scope = None
|
||||
requested_profile = (profile or "").strip()
|
||||
# Plain /api/status stays the machine-level public liveness probe. The
|
||||
|
||||
@@ -6,12 +6,14 @@ monkeypatching on web_server stays authoritative.
|
||||
"""
|
||||
|
||||
import asyncio
|
||||
import shutil
|
||||
import subprocess
|
||||
import sys
|
||||
from typing import Any, Dict, List, Optional
|
||||
|
||||
from fastapi import APIRouter, HTTPException
|
||||
|
||||
from hermes_cli.web_deps import late, LateState
|
||||
from hermes_cli.web_deps import late
|
||||
from hermes_cli.web_models import (
|
||||
TerminalBackendSelect,
|
||||
ToolsetEnvUpdate,
|
||||
@@ -33,18 +35,209 @@ from hermes_cli.web_routers._common import (
|
||||
|
||||
router = APIRouter()
|
||||
|
||||
_find_toolset_provider_row = late("_find_toolset_provider_row")
|
||||
_probe_terminal_backend = late("_probe_terminal_backend")
|
||||
_resolve_toolset_model_plugin = late("_resolve_toolset_model_plugin")
|
||||
_toolset_model_catalog = late("_toolset_model_catalog")
|
||||
load_config = late("load_config")
|
||||
save_config = late("save_config")
|
||||
run_in_threadpool = late("run_in_threadpool")
|
||||
_MODEL_CATALOG_TOOLSETS = LateState("_MODEL_CATALOG_TOOLSETS")
|
||||
_plugin_terminal_backend_rows = late("_plugin_terminal_backend_rows")
|
||||
|
||||
|
||||
def _terminal_cfg_value(terminal_cfg: dict, key: str, env_var: str) -> str:
|
||||
"""Read a terminal.* setting from config.yaml, falling back to its env var."""
|
||||
value = terminal_cfg.get(key)
|
||||
if value is not None and str(value).strip():
|
||||
return str(value).strip()
|
||||
try:
|
||||
from hermes_cli.config import get_env_value
|
||||
|
||||
return (get_env_value(env_var) or "").strip()
|
||||
except Exception:
|
||||
return ""
|
||||
|
||||
|
||||
def _terminal_backend_rows() -> List[Dict[str, str]]:
|
||||
"""Built-in picker rows plus plugin-registered backends (request time).
|
||||
|
||||
Computed per request (mirrors ``_schema_with_dynamic_provider_options``)
|
||||
so a plugin installed after server start still shows up.
|
||||
"""
|
||||
from hermes_cli.web_server import _TERMINAL_BACKENDS
|
||||
return [*_TERMINAL_BACKENDS, *_plugin_terminal_backend_rows()]
|
||||
|
||||
|
||||
def _probe_docker_backend() -> tuple:
|
||||
if not shutil.which("docker"):
|
||||
return (
|
||||
"needs_setup",
|
||||
"Docker CLI not found — install Docker Desktop or docker-ce.",
|
||||
)
|
||||
try:
|
||||
proc = subprocess.run(
|
||||
["docker", "info", "--format", "{{.ServerVersion}}"],
|
||||
capture_output=True,
|
||||
text=True,
|
||||
encoding="utf-8",
|
||||
errors="replace",
|
||||
timeout=2,
|
||||
)
|
||||
if proc.returncode == 0:
|
||||
return ("ready", "")
|
||||
return (
|
||||
"needs_setup",
|
||||
"Docker daemon not reachable — start Docker and retry.",
|
||||
)
|
||||
except subprocess.TimeoutExpired:
|
||||
return ("needs_setup", "Docker daemon not responding (timed out).")
|
||||
except Exception as exc:
|
||||
return ("unavailable", f"Docker probe failed: {exc}")
|
||||
|
||||
|
||||
def _probe_singularity_backend() -> tuple:
|
||||
if shutil.which("singularity") or shutil.which("apptainer"):
|
||||
return ("ready", "")
|
||||
return (
|
||||
"needs_setup",
|
||||
"Neither singularity nor apptainer found on PATH.",
|
||||
)
|
||||
|
||||
|
||||
def _probe_ssh_backend(terminal_cfg: dict) -> tuple:
|
||||
host = _terminal_cfg_value(terminal_cfg, "ssh_host", "TERMINAL_SSH_HOST")
|
||||
user = _terminal_cfg_value(terminal_cfg, "ssh_user", "TERMINAL_SSH_USER")
|
||||
missing = []
|
||||
if not host:
|
||||
missing.append("terminal.ssh_host")
|
||||
if not user:
|
||||
missing.append("terminal.ssh_user")
|
||||
if missing:
|
||||
return (
|
||||
"needs_setup",
|
||||
f"Set {' and '.join(missing)} in config.yaml (or the matching TERMINAL_SSH_* env vars).",
|
||||
)
|
||||
return ("ready", f"{user}@{host}")
|
||||
|
||||
|
||||
def _probe_modal_backend() -> tuple:
|
||||
try:
|
||||
from tools.tool_backend_helpers import has_direct_modal_credentials
|
||||
|
||||
if has_direct_modal_credentials():
|
||||
return ("ready", "")
|
||||
except Exception:
|
||||
pass
|
||||
try:
|
||||
from hermes_cli.config import get_env_value
|
||||
|
||||
if get_env_value("MODAL_TOKEN_ID") and get_env_value("MODAL_TOKEN_SECRET"):
|
||||
return ("ready", "")
|
||||
except Exception:
|
||||
pass
|
||||
return (
|
||||
"needs_setup",
|
||||
"Modal credentials not found — set MODAL_TOKEN_ID and MODAL_TOKEN_SECRET (or run `modal setup`).",
|
||||
)
|
||||
|
||||
|
||||
def _probe_daytona_backend() -> tuple:
|
||||
try:
|
||||
from hermes_cli.config import get_env_value
|
||||
|
||||
if get_env_value("DAYTONA_API_KEY"):
|
||||
return ("ready", "")
|
||||
except Exception:
|
||||
pass
|
||||
return ("needs_setup", "Set DAYTONA_API_KEY to use the Daytona backend.")
|
||||
# Built-ins + plugin-registered backends, computed per request so a plugin
|
||||
# installed after server start still shows up.
|
||||
_terminal_backend_rows = late("_terminal_backend_rows")
|
||||
_terminal_backend_names = late("_terminal_backend_names")
|
||||
|
||||
|
||||
# Toolsets whose backends carry a selectable model catalog, mapped to the
|
||||
# config.yaml section their `model` key lives in. Mirrors the CLI's
|
||||
# post-selection model pickers (`_configure_imagegen_model_for_plugin` /
|
||||
# `_configure_videogen_model_for_plugin` in tools_config.py).
|
||||
_MODEL_CATALOG_TOOLSETS = {
|
||||
"image_gen": "image_gen",
|
||||
"video_gen": "video_gen",
|
||||
}
|
||||
|
||||
|
||||
def _resolve_toolset_model_plugin(ts_key: str, provider_row: dict) -> Optional[str]:
|
||||
"""Map a provider picker row to its model-catalog plugin name.
|
||||
|
||||
Plugin-backed rows carry ``image_gen_plugin_name`` / ``video_gen_plugin_name``;
|
||||
the managed "Nous Subscription" image row instead carries the legacy
|
||||
``imagegen_backend: "fal"`` marker (same underlying FAL catalog).
|
||||
"""
|
||||
if ts_key == "image_gen":
|
||||
return provider_row.get("image_gen_plugin_name") or (
|
||||
"fal" if provider_row.get("imagegen_backend") else None
|
||||
)
|
||||
if ts_key == "video_gen":
|
||||
return provider_row.get("video_gen_plugin_name")
|
||||
return None
|
||||
|
||||
|
||||
def _toolset_model_catalog(ts_key: str, plugin_name: str):
|
||||
"""Return ``(catalog_dict, default_model)`` for a toolset's plugin backend."""
|
||||
from hermes_cli.tools_config import (
|
||||
_plugin_image_gen_catalog,
|
||||
_plugin_video_gen_catalog,
|
||||
)
|
||||
|
||||
if ts_key == "image_gen":
|
||||
return _plugin_image_gen_catalog(plugin_name)
|
||||
return _plugin_video_gen_catalog(plugin_name)
|
||||
|
||||
|
||||
def _find_toolset_provider_row(ts_key: str, config: dict, provider: Optional[str]) -> Optional[dict]:
|
||||
"""Resolve a provider picker row by name, or the active row when omitted."""
|
||||
from hermes_cli.tools_config import (
|
||||
TOOL_CATEGORIES,
|
||||
_is_provider_active,
|
||||
_visible_providers,
|
||||
)
|
||||
|
||||
cat = TOOL_CATEGORIES.get(ts_key)
|
||||
if cat is None:
|
||||
return None
|
||||
rows = _visible_providers(cat, config, force_fresh=True)
|
||||
if provider:
|
||||
return next((p for p in rows if p.get("name") == provider), None)
|
||||
return next(
|
||||
(p for p in rows if _is_provider_active(p, config, force_fresh=True)), None
|
||||
)
|
||||
|
||||
|
||||
def _terminal_backend_names() -> set:
|
||||
"""Valid ``terminal.backend`` values, including plugin backends."""
|
||||
return {row["name"] for row in _terminal_backend_rows()}
|
||||
|
||||
|
||||
def _probe_terminal_backend(name: str, terminal_cfg: dict) -> tuple:
|
||||
"""Return ``(status, detail)`` for one backend. Never raises."""
|
||||
try:
|
||||
if name == "local":
|
||||
return ("ready", "")
|
||||
if name == "docker":
|
||||
return _probe_docker_backend()
|
||||
if name == "singularity":
|
||||
return _probe_singularity_backend()
|
||||
if name == "ssh":
|
||||
return _probe_ssh_backend(terminal_cfg)
|
||||
if name == "modal":
|
||||
return _probe_modal_backend()
|
||||
if name == "daytona":
|
||||
return _probe_daytona_backend()
|
||||
try:
|
||||
from agent.terminal_env_registry import get_provider
|
||||
|
||||
provider = get_provider(name)
|
||||
if provider is not None:
|
||||
return provider.probe()
|
||||
except Exception:
|
||||
pass
|
||||
return ("unavailable", f"Unknown backend: {name}")
|
||||
except Exception as exc: # pragma: no cover — belt-and-braces guard
|
||||
return ("unavailable", f"Probe failed: {exc}")
|
||||
|
||||
|
||||
def _require_known_toolset(name: str) -> None:
|
||||
|
||||
File diff suppressed because it is too large
Load Diff
Reference in New Issue
Block a user