diff --git a/gateway/platforms/api_server.py b/gateway/platforms/api_server.py index 6287132243ad..e44c32b11884 100644 --- a/gateway/platforms/api_server.py +++ b/gateway/platforms/api_server.py @@ -27,6 +27,14 @@ AnythingLLM, NextChat, ChatBox, etc.) can connect to hermes-agent through this adapter by pointing at http://localhost:8642/v1 and authenticating with API_SERVER_KEY. +When ``gateway.multiplex_profiles`` is on, the default profile owns this +listener and secondary profiles are reached via a URL prefix — same contract +as the webhook adapter: + + GET /p//v1/models + POST /p//v1/chat/completions + ... + Requires: - aiohttp (already available in the gateway) """ @@ -35,7 +43,7 @@ import asyncio import hashlib import hmac import json -from contextlib import contextmanager +from contextlib import contextmanager, nullcontext from contextvars import ContextVar from functools import wraps import logging @@ -48,6 +56,17 @@ import uuid from pathlib import Path from typing import Any, Dict, List, Optional +# Sentinel returned by _resolve_request_profile when a /p// prefix +# names a profile this gateway does not serve (→ 404). Distinct from None +# (no prefix / multiplexing off → handle as the default profile). +_PROFILE_REJECTED = object() + +# Profile selected by the /p// URL prefix for the current request. +# Set by the profile-prefix middleware; read by handlers / _run_agent. +_api_request_profile: ContextVar[Optional[str]] = ContextVar( + "api_server_request_profile", default=None +) + def _approval_event_choices(*, smart_denied: bool, allow_permanent: bool) -> list[str]: if smart_denied: return ["once", "deny"] @@ -1228,6 +1247,117 @@ class APIServerAdapter(BasePlatformAdapter): status=401, ) + # ------------------------------------------------------------------ + # Multi-profile multiplexing (/p//…) + # ------------------------------------------------------------------ + + def _resolve_request_profile(self, request: "web.Request"): + """Resolve + validate the /p// URL prefix on an API request. + + Returns: + - ``None`` when no profile prefix is present, or multiplexing is off + (the prefix is ignored; request handled as the default profile). + - the profile name (str) when present, multiplexing is on, and the + profile is one this gateway serves. + - ``_PROFILE_REJECTED`` when a prefix is present but the profile is + unknown/unconfigured (handler/middleware returns 404). + """ + profile = (request.match_info.get("profile") or "").strip() + if not profile: + return None + runner = getattr(self, "gateway_runner", None) + cfg = getattr(runner, "config", None) + if not getattr(cfg, "multiplex_profiles", False): + # Prefix supplied but multiplexing is off — ignore it, behave as + # the single-profile gateway (don't 404 a would-be valid route). + return None + try: + from hermes_cli.profiles import profiles_to_serve + + served = {name for name, _ in profiles_to_serve(multiplex=True)} + except Exception: + return _PROFILE_REJECTED + if profile not in served: + return _PROFILE_REJECTED + return profile + + @staticmethod + def _profile_scope(profile: Optional[str]): + """Enter the multiplex profile runtime scope, or a no-op when unset.""" + if not profile: + return nullcontext() + from gateway.run import _profile_runtime_scope + from hermes_cli.profiles import get_profile_dir + + return _profile_runtime_scope(get_profile_dir(profile)) + + def _make_profile_prefix_middleware(self): + """Reject unknown /p// prefixes and scope the request home.""" + + @web.middleware + async def profile_prefix_middleware(request: "web.Request", handler): + profile = self._resolve_request_profile(request) + if profile is _PROFILE_REJECTED: + return web.json_response( + {"error": "Unknown or unconfigured profile"}, + status=404, + ) + token = _api_request_profile.set(profile) + try: + with self._profile_scope(profile): + return await handler(request) + finally: + _api_request_profile.reset(token) + + return profile_prefix_middleware + + def _http_route_table(self) -> List[tuple]: + """Return (method, path, handler) rows registered by ``connect()``. + + Kept as a method so multiplex tests can assert the /p// + mirrors without starting a real aiohttp listener. + """ + routes: List[tuple] = [ + ("GET", "/health", self._handle_health), + ("GET", "/health/detailed", self._handle_health_detailed), + ("GET", "/v1/health", self._handle_health), + ("GET", "/v1/models", self._handle_models), + ("GET", "/v1/capabilities", self._handle_capabilities), + ("GET", "/v1/skills", self._handle_skills), + ("GET", "/v1/toolsets", self._handle_toolsets), + ("GET", "/api/sessions", self._handle_list_sessions), + ("POST", "/api/sessions", self._handle_create_session), + ("GET", "/api/sessions/{session_id}", self._handle_get_session), + ("PATCH", "/api/sessions/{session_id}", self._handle_patch_session), + ("DELETE", "/api/sessions/{session_id}", self._handle_delete_session), + ("GET", "/api/sessions/{session_id}/messages", self._handle_session_messages), + ("POST", "/api/sessions/{session_id}/fork", self._handle_fork_session), + ("POST", "/api/sessions/{session_id}/chat", self._handle_session_chat), + ("POST", "/api/sessions/{session_id}/chat/stream", self._handle_session_chat_stream), + ("POST", "/v1/chat/completions", self._handle_chat_completions), + ("POST", "/v1/responses", self._handle_responses), + ("GET", "/v1/responses/{response_id}", self._handle_get_response), + ("DELETE", "/v1/responses/{response_id}", self._handle_delete_response), + ("GET", "/api/jobs", self._handle_list_jobs), + ("POST", "/api/jobs", self._handle_create_job), + ("GET", "/api/jobs/{job_id}", self._handle_get_job), + ("PATCH", "/api/jobs/{job_id}", self._handle_update_job), + ("DELETE", "/api/jobs/{job_id}", self._handle_delete_job), + ("POST", "/api/jobs/{job_id}/pause", self._handle_pause_job), + ("POST", "/api/jobs/{job_id}/resume", self._handle_resume_job), + ("POST", "/api/jobs/{job_id}/run", self._handle_run_job), + ("POST", "/v1/runs", self._handle_runs), + ("GET", "/v1/runs/{run_id}", self._handle_get_run), + ("GET", "/v1/runs/{run_id}/events", self._handle_run_events), + ("POST", "/v1/runs/{run_id}/approval", self._handle_run_approval), + ("POST", "/v1/runs/{run_id}/stop", self._handle_stop_run), + ] + if _CRON_AVAILABLE: + # Chronos managed-cron fire webhook (NAS → agent). Authenticated + # by a NAS-minted JWT (NOT API_SERVER_KEY). + routes.append(("POST", "/api/cron/fire", self._handle_cron_fire)) + return routes + # ------------------------------------------------------------------ # Session header helpers # ------------------------------------------------------------------ @@ -1298,18 +1428,39 @@ class APIServerAdapter(BasePlatformAdapter): # ------------------------------------------------------------------ def _ensure_session_db(self): - """Lazily initialise and return the shared SessionDB instance. + """Lazily initialise and return the SessionDB for the active profile home. Sessions are persisted to ``state.db`` so that ``hermes sessions list`` shows API-server conversations alongside CLI and gateway ones. + + Under multiplex ``/p//`` requests the profile runtime scope + redirects ``get_hermes_home()``, so each profile gets its own DB — + never the default profile's file. """ - if self._session_db is None: - try: - from hermes_state import SessionDB - self._session_db = SessionDB() - except Exception as e: - logger.debug("SessionDB unavailable for API server: %s", e) - return self._session_db + # Explicit override (tests / manual wiring) wins. Production never sets + # this externally, so the per-home cache below is the live path — and + # we deliberately do NOT write back into ``self._session_db`` there, or + # the first profile served would pin every later request to its DB. + if self._session_db is not None: + return self._session_db + try: + from hermes_constants import get_hermes_home + from hermes_state import SessionDB + + home = get_hermes_home() + cache = getattr(self, "_session_dbs", None) + if cache is None: + cache = {} + self._session_dbs = cache + key = str(home) + db = cache.get(key) + if db is None: + db = SessionDB(db_path=home / "state.db") + cache[key] = db + return db + except Exception as e: + logger.debug("SessionDB unavailable for API server: %s", e) + return None # ------------------------------------------------------------------ # Agent creation helper @@ -1586,20 +1737,32 @@ class APIServerAdapter(BasePlatformAdapter): }) async def _handle_models(self, request: "web.Request") -> "web.Response": - """GET /v1/models — list hermes-agent and any configured model_routes aliases.""" + """GET /v1/models — list hermes-agent and any configured model_routes aliases. + + Under ``/p//v1/models`` (multiplex on) the advertised primary + model id follows that profile's name/config, not the default adapter's + cached ``_model_name``. + """ auth_err = self._check_auth(request) if auth_err: return auth_err now = int(time.time()) + # Middleware already entered the profile runtime scope when a /p/ + # prefix was present, so get_active_profile_name() resolves correctly. + model_name = ( + self._resolve_model_name("") + if _api_request_profile.get() + else self._model_name + ) models = [ { - "id": self._model_name, + "id": model_name, "object": "model", "created": now, "owned_by": "hermes", "permission": [], - "root": self._model_name, + "root": model_name, "parent": None, } ] @@ -1607,7 +1770,7 @@ class APIServerAdapter(BasePlatformAdapter): # Only the alias and resolved model name are exposed — never provider # credentials. for alias, route_cfg in self._model_routes.items(): - if alias == self._model_name: + if alias == model_name: continue # already listed above models.append({ "id": alias, @@ -1616,7 +1779,7 @@ class APIServerAdapter(BasePlatformAdapter): "owned_by": "hermes", "permission": [], "root": route_cfg.get("model", alias), - "parent": self._model_name, + "parent": model_name, }) return web.json_response({"object": "list", "data": models}) @@ -4223,48 +4386,53 @@ class APIServerAdapter(BasePlatformAdapter): another thread to stop in-progress LLM calls. """ loop = asyncio.get_running_loop() + # Capture before hopping to the executor — ContextVars do not follow + # run_in_executor threads, so the profile scope must be re-entered + # inside _run() from this explicit value. + request_profile = _api_request_profile.get() def _run(): from gateway.session_context import clear_session_vars - tokens = self._bind_api_server_session( - chat_id=session_id or "", - session_key=gateway_session_key or session_id or "", - session_id=session_id or "", - ) - try: - agent = self._create_agent( - ephemeral_system_prompt=ephemeral_system_prompt, - session_id=session_id, - stream_delta_callback=stream_delta_callback, - tool_progress_callback=tool_progress_callback, - tool_start_callback=tool_start_callback, - tool_complete_callback=tool_complete_callback, - gateway_session_key=gateway_session_key, - route=route, + with self._profile_scope(request_profile): + tokens = self._bind_api_server_session( + chat_id=session_id or "", + session_key=gateway_session_key or session_id or "", + session_id=session_id or "", ) - if agent_ref is not None: - agent_ref[0] = agent - effective_task_id = session_id or str(uuid.uuid4()) - result = agent.run_conversation( - user_message=user_message, - conversation_history=conversation_history, - task_id=effective_task_id, - ) - usage = { - "input_tokens": getattr(agent, "session_prompt_tokens", 0) or 0, - "output_tokens": getattr(agent, "session_completion_tokens", 0) or 0, - "total_tokens": getattr(agent, "session_total_tokens", 0) or 0, - } - # Include the effective session ID in the result so callers - # (e.g. X-Hermes-Session-Id header) can track compression- - # triggered session rotations. (#16938) - _eff_sid = getattr(agent, "session_id", session_id) - if isinstance(_eff_sid, str) and _eff_sid: - result["session_id"] = _eff_sid - return result, usage - finally: - clear_session_vars(tokens) + try: + agent = self._create_agent( + ephemeral_system_prompt=ephemeral_system_prompt, + session_id=session_id, + stream_delta_callback=stream_delta_callback, + tool_progress_callback=tool_progress_callback, + tool_start_callback=tool_start_callback, + tool_complete_callback=tool_complete_callback, + gateway_session_key=gateway_session_key, + route=route, + ) + if agent_ref is not None: + agent_ref[0] = agent + effective_task_id = session_id or str(uuid.uuid4()) + result = agent.run_conversation( + user_message=user_message, + conversation_history=conversation_history, + task_id=effective_task_id, + ) + usage = { + "input_tokens": getattr(agent, "session_prompt_tokens", 0) or 0, + "output_tokens": getattr(agent, "session_completion_tokens", 0) or 0, + "total_tokens": getattr(agent, "session_total_tokens", 0) or 0, + } + # Include the effective session ID in the result so callers + # (e.g. X-Hermes-Session-Id header) can track compression- + # triggered session rotations. (#16938) + _eff_sid = getattr(agent, "session_id", session_id) + if isinstance(_eff_sid, str) and _eff_sid: + result["session_id"] = _eff_sid + return result, usage + finally: + clear_session_vars(tokens) self._activate_admitted_request() self._inflight_agent_runs += 1 @@ -4464,6 +4632,9 @@ class APIServerAdapter(BasePlatformAdapter): # Per-client model routing for /v1/runs (see model_routes). route = self._resolve_route(body.get("model")) + # Background task outlives the HTTP response (and thus the middleware + # profile scope). Capture now and re-enter inside the task/executor. + request_profile = _api_request_profile.get() async def _run_and_close(): try: @@ -4480,14 +4651,15 @@ class APIServerAdapter(BasePlatformAdapter): last_event="run.cancelled", ) return - agent = self._create_agent( - ephemeral_system_prompt=ephemeral_system_prompt, - session_id=session_id, - stream_delta_callback=_text_cb, - tool_progress_callback=event_cb, - gateway_session_key=gateway_session_key, - route=route, - ) + with self._profile_scope(request_profile): + agent = self._create_agent( + ephemeral_system_prompt=ephemeral_system_prompt, + session_id=session_id, + stream_delta_callback=_text_cb, + tool_progress_callback=event_cb, + gateway_session_key=gateway_session_key, + route=route, + ) self._active_run_agents[run_id] = agent def _approval_notify(approval_data: Dict[str, Any]) -> None: @@ -4531,40 +4703,41 @@ class APIServerAdapter(BasePlatformAdapter): effective_task_id = session_id or run_id approval_token = None session_tokens = [] - try: - # Bind approval/session identity for this API run via - # contextvars so concurrent runs do not share process - # environment state. - approval_token = set_current_session_key(approval_session_key) - session_tokens = self._bind_api_server_session( - session_key=approval_session_key, - ) - register_gateway_notify(approval_session_key, _approval_notify) - r = agent.run_conversation( - user_message=user_message, - conversation_history=conversation_history, - task_id=effective_task_id, - ) - finally: + with self._profile_scope(request_profile): try: - unregister_gateway_notify(approval_session_key) + # Bind approval/session identity for this API run via + # contextvars so concurrent runs do not share process + # environment state. + approval_token = set_current_session_key(approval_session_key) + session_tokens = self._bind_api_server_session( + session_key=approval_session_key, + ) + register_gateway_notify(approval_session_key, _approval_notify) + r = agent.run_conversation( + user_message=user_message, + conversation_history=conversation_history, + task_id=effective_task_id, + ) finally: - if approval_token is not None: - try: - reset_current_session_key(approval_token) - except Exception: - pass - if session_tokens: - try: - clear_session_vars(session_tokens) - except Exception: - pass - u = { - "input_tokens": getattr(agent, "session_prompt_tokens", 0) or 0, - "output_tokens": getattr(agent, "session_completion_tokens", 0) or 0, - "total_tokens": getattr(agent, "session_total_tokens", 0) or 0, - } - return r, u + try: + unregister_gateway_notify(approval_session_key) + finally: + if approval_token is not None: + try: + reset_current_session_key(approval_token) + except Exception: + pass + if session_tokens: + try: + clear_session_vars(session_tokens) + except Exception: + pass + u = { + "input_tokens": getattr(agent, "session_prompt_tokens", 0) or 0, + "output_tokens": getattr(agent, "session_completion_tokens", 0) or 0, + "total_tokens": getattr(agent, "session_total_tokens", 0) or 0, + } + return r, u result, usage = await asyncio.get_running_loop().run_in_executor(None, _run_sync) if run_id in self._stopping_run_ids: @@ -4970,50 +5143,24 @@ class APIServerAdapter(BasePlatformAdapter): return False try: - mws = [mw for mw in (cors_middleware, body_limit_middleware, security_headers_middleware) if mw is not None] + mws = [ + mw + for mw in ( + self._make_profile_prefix_middleware(), + cors_middleware, + body_limit_middleware, + security_headers_middleware, + ) + if mw is not None + ] self._app = web.Application(middlewares=mws, client_max_size=MAX_REQUEST_BYTES) assert self._app is not None - self._app.router.add_get("/health", self._handle_health) - self._app.router.add_get("/health/detailed", self._handle_health_detailed) - self._app.router.add_get("/v1/health", self._handle_health) - self._app.router.add_get("/v1/models", self._handle_models) - self._app.router.add_get("/v1/capabilities", self._handle_capabilities) - self._app.router.add_get("/v1/skills", self._handle_skills) - self._app.router.add_get("/v1/toolsets", self._handle_toolsets) - # Session/client control surface (thin wrappers over SessionDB + _run_agent) - self._app.router.add_get("/api/sessions", self._handle_list_sessions) - self._app.router.add_post("/api/sessions", self._handle_create_session) - self._app.router.add_get("/api/sessions/{session_id}", self._handle_get_session) - self._app.router.add_patch("/api/sessions/{session_id}", self._handle_patch_session) - self._app.router.add_delete("/api/sessions/{session_id}", self._handle_delete_session) - self._app.router.add_get("/api/sessions/{session_id}/messages", self._handle_session_messages) - self._app.router.add_post("/api/sessions/{session_id}/fork", self._handle_fork_session) - self._app.router.add_post("/api/sessions/{session_id}/chat", self._handle_session_chat) - self._app.router.add_post("/api/sessions/{session_id}/chat/stream", self._handle_session_chat_stream) - self._app.router.add_post("/v1/chat/completions", self._handle_chat_completions) - self._app.router.add_post("/v1/responses", self._handle_responses) - self._app.router.add_get("/v1/responses/{response_id}", self._handle_get_response) - self._app.router.add_delete("/v1/responses/{response_id}", self._handle_delete_response) - # Cron jobs management API - self._app.router.add_get("/api/jobs", self._handle_list_jobs) - self._app.router.add_post("/api/jobs", self._handle_create_job) - self._app.router.add_get("/api/jobs/{job_id}", self._handle_get_job) - self._app.router.add_patch("/api/jobs/{job_id}", self._handle_update_job) - self._app.router.add_delete("/api/jobs/{job_id}", self._handle_delete_job) - self._app.router.add_post("/api/jobs/{job_id}/pause", self._handle_pause_job) - self._app.router.add_post("/api/jobs/{job_id}/resume", self._handle_resume_job) - self._app.router.add_post("/api/jobs/{job_id}/run", self._handle_run_job) - - # Chronos managed-cron fire webhook (NAS → agent). Authenticated by a - # NAS-minted JWT (NOT API_SERVER_KEY), so it has its own auth path. - if _CRON_AVAILABLE: - self._app.router.add_post("/api/cron/fire", self._handle_cron_fire) - # Structured event streaming - self._app.router.add_post("/v1/runs", self._handle_runs) - self._app.router.add_get("/v1/runs/{run_id}", self._handle_get_run) - self._app.router.add_get("/v1/runs/{run_id}/events", self._handle_run_events) - self._app.router.add_post("/v1/runs/{run_id}/approval", self._handle_run_approval) - self._app.router.add_post("/v1/runs/{run_id}/stop", self._handle_stop_run) + # Native routes + multiplex /p//… mirrors. Same handlers; + # the profile-prefix middleware validates the prefix and scopes + # config/credentials to that profile when multiplexing is on. + for method, path, handler in self._http_route_table(): + self._app.router.add_route(method, path, handler) + self._app.router.add_route(method, f"/p/{{profile}}{path}", handler) # Store the adapter after native routes are registered. Local Hermes-Relay # bootstrap shims use this key as a feature-detection hook; registering # native routes first lets those shims no-op instead of shadowing the