refactor(hermes_cli): proxy adapters/cli/sse_done + relaunch/resource_limits/psutil_android/relay_cutover — dict dispatch, unified stderr printer, compact prose
This commit is contained in:
@@ -1,9 +1,4 @@
|
||||
"""Upstream adapter registry for the local proxy server.
|
||||
|
||||
Each adapter wraps a provider's OAuth state and exposes a uniform interface the proxy server can use
|
||||
to forward requests with a freshly-minted bearer token. See :class:`UpstreamAdapter` for the
|
||||
contract.
|
||||
"""
|
||||
"""Upstream adapter registry for the local proxy server (see :class:`UpstreamAdapter`)."""
|
||||
|
||||
from typing import Dict, Type
|
||||
|
||||
@@ -11,8 +6,7 @@ from hermes_cli.proxy.adapters.base import UpstreamAdapter
|
||||
from hermes_cli.proxy.adapters.nous_portal import NousPortalAdapter
|
||||
from hermes_cli.proxy.adapters.xai import XAIGrokAdapter
|
||||
|
||||
# Registry of available adapter classes keyed by provider name as used on
|
||||
# the ``hermes proxy start --provider <name>`` CLI flag.
|
||||
# Keyed by the ``hermes proxy start --provider <name>`` value.
|
||||
ADAPTERS: Dict[str, Type[UpstreamAdapter]] = {
|
||||
"nous": NousPortalAdapter,
|
||||
"xai": XAIGrokAdapter,
|
||||
@@ -24,9 +18,7 @@ def get_adapter(name: str) -> UpstreamAdapter:
|
||||
key = (name or "").strip().lower()
|
||||
if key not in ADAPTERS:
|
||||
available = ", ".join(sorted(ADAPTERS)) or "(none)"
|
||||
raise ValueError(
|
||||
f"Unknown proxy upstream provider: {name!r}. Available: {available}"
|
||||
)
|
||||
raise ValueError(f"Unknown proxy upstream provider: {name!r}. Available: {available}")
|
||||
return ADAPTERS[key]()
|
||||
|
||||
|
||||
|
||||
@@ -1,7 +1,4 @@
|
||||
"""Abstract base for proxy upstream adapters.
|
||||
|
||||
The proxy server is otherwise provider-agnostic.
|
||||
"""
|
||||
"""Abstract base for proxy upstream adapters; the proxy server is otherwise provider-agnostic."""
|
||||
|
||||
from __future__ import annotations
|
||||
|
||||
@@ -14,17 +11,10 @@ from typing import FrozenSet, Optional
|
||||
class UpstreamCredential:
|
||||
"""A resolved bearer + base URL ready to forward to."""
|
||||
|
||||
bearer: str
|
||||
"""Authorization header value to send upstream (token only, no ``Bearer`` prefix)."""
|
||||
|
||||
base_url: str
|
||||
"""Upstream base URL, e.g. ``https://inference-api.nousresearch.com/v1``."""
|
||||
|
||||
bearer: str # token only, no ``Bearer`` prefix
|
||||
base_url: str # e.g. ``https://inference-api.nousresearch.com/v1``
|
||||
token_type: str = "Bearer"
|
||||
"""Auth scheme — currently always ``Bearer`` for supported providers."""
|
||||
|
||||
expires_at: Optional[str] = None
|
||||
"""ISO-8601 expiry timestamp for the bearer, when known. Informational."""
|
||||
expires_at: Optional[str] = None # ISO-8601, informational
|
||||
|
||||
|
||||
class UpstreamAdapter(ABC):
|
||||
@@ -43,40 +33,24 @@ class UpstreamAdapter(ABC):
|
||||
@property
|
||||
@abstractmethod
|
||||
def allowed_paths(self) -> FrozenSet[str]:
|
||||
"""Set of relative request paths the upstream accepts.
|
||||
|
||||
Paths are relative to the proxy's ``/v1`` mount (``"/chat/completions"`` ⇒
|
||||
``/v1/chat/completions``). Requests outside this set get a 404 with a helpful body.
|
||||
"""
|
||||
"""Paths relative to the proxy's ``/v1`` mount (``"/chat/completions"`` ⇒
|
||||
``/v1/chat/completions``); anything else gets a 404 with a helpful body."""
|
||||
|
||||
@abstractmethod
|
||||
def is_authenticated(self) -> bool:
|
||||
"""Return True if the user has usable credentials for this upstream.
|
||||
|
||||
Should be cheap — no network calls. Used by ``proxy start`` for a clear up-front error
|
||||
before binding a port.
|
||||
"""
|
||||
"""Cheap (no network) usable-credentials check; ``proxy start`` uses it for a clear
|
||||
up-front error before binding a port."""
|
||||
|
||||
@abstractmethod
|
||||
def get_credential(self) -> UpstreamCredential:
|
||||
"""Return a fresh credential, refreshing or rotating if necessary.
|
||||
|
||||
Implementations refresh a near-expiry access token, rotate a near-expiry upstream bearer key
|
||||
and persist refreshed state to disk. Raises RuntimeError when unauthenticated or refresh
|
||||
fails; the proxy then returns 401 to the client.
|
||||
"""
|
||||
"""Fresh credential (refreshing/rotating + persisting as needed). Raises RuntimeError when
|
||||
unauthenticated or refresh fails; the proxy then returns 401 to the client."""
|
||||
|
||||
def get_retry_credential(
|
||||
self,
|
||||
*,
|
||||
failed_credential: UpstreamCredential,
|
||||
status_code: int,
|
||||
self, *, failed_credential: UpstreamCredential, status_code: int
|
||||
) -> Optional[UpstreamCredential]:
|
||||
"""Return an alternate credential after an upstream auth failure.
|
||||
|
||||
The default is no retry. Providers can override this for one-shot fallback paths after the
|
||||
upstream rejects the first request.
|
||||
"""
|
||||
"""Alternate credential for a one-shot retry after the upstream rejects the first request;
|
||||
default is no retry."""
|
||||
_ = failed_credential, status_code
|
||||
return None
|
||||
|
||||
|
||||
@@ -24,25 +24,16 @@ from hermes_cli.proxy.adapters.base import UpstreamAdapter, UpstreamCredential
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
# Endpoints inference-api.nousresearch.com actually serves. Anything else
|
||||
# the proxy will reject with 404 — keeps stray clients from leaking weird
|
||||
# requests to the upstream.
|
||||
_ALLOWED_PATHS: FrozenSet[str] = frozenset(
|
||||
{
|
||||
"/chat/completions",
|
||||
"/completions",
|
||||
"/embeddings",
|
||||
"/models",
|
||||
}
|
||||
)
|
||||
# Endpoints inference-api.nousresearch.com actually serves; anything else is a 404 so stray
|
||||
# clients cannot leak odd requests upstream.
|
||||
_ALLOWED_PATHS: FrozenSet[str] = frozenset({"/chat/completions", "/completions", "/embeddings", "/models"})
|
||||
|
||||
|
||||
class NousPortalAdapter(UpstreamAdapter):
|
||||
"""Proxy upstream for the Nous Portal inference API."""
|
||||
|
||||
def __init__(self) -> None:
|
||||
# Serialize proxy requests in this process; cross-process token refresh
|
||||
# and persistence are handled by resolve_nous_runtime_credentials().
|
||||
# In-process serialization; cross-process refresh/persistence is resolve_nous_runtime_credentials().
|
||||
self._lock = threading.Lock()
|
||||
|
||||
@property
|
||||
@@ -58,22 +49,17 @@ class NousPortalAdapter(UpstreamAdapter):
|
||||
return _ALLOWED_PATHS
|
||||
|
||||
def is_authenticated(self) -> bool:
|
||||
# We need either a usable inference JWT OR (refresh_token + access_token)
|
||||
# to recover. The refresh helper validates and refreshes as needed.
|
||||
# Usable inference JWT, OR refresh_token + access_token to recover via the refresh helper.
|
||||
state = self._read_state() or {}
|
||||
return bool(
|
||||
state.get("agent_key")
|
||||
or (state.get("refresh_token") and state.get("access_token"))
|
||||
state.get("agent_key") or (state.get("refresh_token") and state.get("access_token"))
|
||||
)
|
||||
|
||||
def get_credential(self) -> UpstreamCredential:
|
||||
return self._get_credential()
|
||||
|
||||
def get_retry_credential(
|
||||
self,
|
||||
*,
|
||||
failed_credential: UpstreamCredential,
|
||||
status_code: int,
|
||||
self, *, failed_credential: UpstreamCredential, status_code: int
|
||||
) -> Optional[UpstreamCredential]:
|
||||
_ = failed_credential
|
||||
if status_code != 401:
|
||||
@@ -82,16 +68,12 @@ class NousPortalAdapter(UpstreamAdapter):
|
||||
return self._get_credential(force_refresh=True)
|
||||
|
||||
def _get_credential(
|
||||
self,
|
||||
*,
|
||||
force_refresh: bool = False,
|
||||
self, *, force_refresh: bool = False
|
||||
) -> UpstreamCredential:
|
||||
with self._lock:
|
||||
state = self._read_state()
|
||||
if state is None:
|
||||
raise RuntimeError(
|
||||
"Not logged into Nous Portal. Run `hermes auth add nous` first."
|
||||
)
|
||||
raise RuntimeError("Not logged into Nous Portal. Run `hermes auth add nous` first.")
|
||||
|
||||
try:
|
||||
refreshed = resolve_nous_runtime_credentials(force_refresh=force_refresh)
|
||||
@@ -99,13 +81,9 @@ class NousPortalAdapter(UpstreamAdapter):
|
||||
if isinstance(exc, AuthError) and _is_terminal_nous_refresh_error(exc):
|
||||
_quarantine_nous_oauth_state(state, exc, reason="proxy_refresh_failure")
|
||||
self._save_state(
|
||||
state,
|
||||
quarantine_error=exc,
|
||||
quarantine_reason="proxy_refresh_failure",
|
||||
state, quarantine_error=exc, quarantine_reason="proxy_refresh_failure"
|
||||
)
|
||||
raise RuntimeError(
|
||||
f"Failed to refresh Nous Portal credentials: {exc}"
|
||||
) from exc
|
||||
raise RuntimeError(f"Failed to refresh Nous Portal credentials: {exc}") from exc
|
||||
|
||||
runtime_key = refreshed.get("api_key")
|
||||
if not runtime_key:
|
||||
@@ -114,31 +92,21 @@ class NousPortalAdapter(UpstreamAdapter):
|
||||
"Try `hermes auth add nous` to re-authenticate."
|
||||
)
|
||||
|
||||
# base_url returned by resolve_nous_runtime_credentials() already
|
||||
# honors the NOUS_INFERENCE_BASE_URL env override (the documented
|
||||
# dev/staging escape hatch). Re-validating it here against the prod
|
||||
# host allowlist would wrongly reject a legitimate staging override,
|
||||
# so layer the same env-first overlay on top of the network-validated
|
||||
# value: env override wins, else validate the returned URL, else
|
||||
# fall back to the production default (defense-in-depth for a future
|
||||
# source-layer bypass).
|
||||
# The returned base_url already honors the NOUS_INFERENCE_BASE_URL override (documented
|
||||
# dev/staging hatch); validating it against the prod allowlist would reject a legit
|
||||
# staging URL. So: env override wins, else network-validate the returned URL, else the
|
||||
# production default (defense-in-depth against a future source-layer bypass).
|
||||
base_url = (
|
||||
_nous_inference_env_override()
|
||||
or _validate_nous_inference_url_from_network(refreshed.get("base_url"))
|
||||
or DEFAULT_NOUS_INFERENCE_URL
|
||||
)
|
||||
base_url = base_url.rstrip("/")
|
||||
).rstrip("/")
|
||||
|
||||
return UpstreamCredential(
|
||||
bearer=runtime_key,
|
||||
base_url=base_url,
|
||||
expires_at=refreshed.get("expires_at"),
|
||||
bearer=runtime_key, base_url=base_url, expires_at=refreshed.get("expires_at")
|
||||
)
|
||||
|
||||
# ------------------------------------------------------------------
|
||||
# Internal helpers — auth.json access. Kept local rather than added
|
||||
# to hermes_cli.auth to avoid expanding that module's public surface.
|
||||
# ------------------------------------------------------------------
|
||||
# auth.json access — kept local so hermes_cli.auth's public surface does not grow.
|
||||
|
||||
def _read_state(self) -> Optional[Dict[str, Any]]:
|
||||
try:
|
||||
@@ -148,7 +116,6 @@ class NousPortalAdapter(UpstreamAdapter):
|
||||
logger.warning("proxy: failed to load auth store: %s", exc)
|
||||
return None
|
||||
state = (store.get("providers") or {}).get("nous")
|
||||
# copy so the refresh helper can mutate freely
|
||||
return dict(state) if isinstance(state, dict) else None
|
||||
|
||||
def _save_state(
|
||||
@@ -162,11 +129,7 @@ class NousPortalAdapter(UpstreamAdapter):
|
||||
with _auth_store_lock():
|
||||
store = _load_auth_store()
|
||||
if quarantine_error is not None and quarantine_reason:
|
||||
_quarantine_nous_pool_entries(
|
||||
store,
|
||||
quarantine_error,
|
||||
reason=quarantine_reason,
|
||||
)
|
||||
_quarantine_nous_pool_entries(store, quarantine_error, reason=quarantine_reason)
|
||||
providers = store.setdefault("providers", {})
|
||||
providers["nous"] = state
|
||||
_save_auth_store(store)
|
||||
|
||||
@@ -14,17 +14,9 @@ logger = logging.getLogger(__name__)
|
||||
|
||||
_POOL_PROVIDER = "xai-oauth"
|
||||
|
||||
# xAI's public API is OpenAI-compatible for the endpoints Hermes commonly
|
||||
# uses. The Responses endpoint is included because Hermes' native xAI runtime
|
||||
# uses codex_responses mode.
|
||||
# OpenAI-compatible endpoints; ``/responses`` because the native xAI runtime uses codex_responses.
|
||||
_ALLOWED_PATHS: FrozenSet[str] = frozenset(
|
||||
{
|
||||
"/responses",
|
||||
"/chat/completions",
|
||||
"/completions",
|
||||
"/embeddings",
|
||||
"/models",
|
||||
}
|
||||
{"/responses", "/chat/completions", "/completions", "/embeddings", "/models"}
|
||||
)
|
||||
|
||||
|
||||
@@ -74,10 +66,7 @@ class XAIGrokAdapter(UpstreamAdapter):
|
||||
return self._credential_from_entry(entry)
|
||||
|
||||
def get_retry_credential(
|
||||
self,
|
||||
*,
|
||||
failed_credential: UpstreamCredential,
|
||||
status_code: int,
|
||||
self, *, failed_credential: UpstreamCredential, status_code: int
|
||||
) -> Optional[UpstreamCredential]:
|
||||
if status_code not in {401, 429}:
|
||||
return None
|
||||
@@ -87,10 +76,8 @@ class XAIGrokAdapter(UpstreamAdapter):
|
||||
if pool is None:
|
||||
return None
|
||||
|
||||
# 401: try refreshing the current key first. 429: never refresh — mark
|
||||
# the rate-limited key with its 1-hour cooldown and rotate to the next
|
||||
# available credential. Returns None when the pool has no other key to
|
||||
# offer — the 429 will flow back to the client.
|
||||
# 401: refresh the current key first. 429: never refresh — cooldown the rate-limited
|
||||
# key and rotate. None when the pool has nothing else → the status flows to the client.
|
||||
refreshed = pool.try_refresh_current() if status_code == 401 else None
|
||||
if refreshed is None:
|
||||
refreshed = pool.mark_exhausted_and_rotate(status_code=status_code)
|
||||
|
||||
+27
-44
@@ -9,63 +9,52 @@ from typing import Any
|
||||
|
||||
from hermes_cli.proxy.adapters import ADAPTERS, get_adapter
|
||||
from hermes_cli.proxy.server import (
|
||||
AIOHTTP_AVAILABLE,
|
||||
DEFAULT_HOST,
|
||||
DEFAULT_PORT,
|
||||
run_server,
|
||||
AIOHTTP_AVAILABLE, DEFAULT_HOST, DEFAULT_PORT, run_server
|
||||
)
|
||||
|
||||
logger = logging.getLogger(__name__)
|
||||
|
||||
|
||||
def _print_aiohttp_missing() -> None:
|
||||
print(
|
||||
"hermes proxy requires aiohttp. Run `hermes setup` to install it.",
|
||||
file=sys.stderr,
|
||||
)
|
||||
def _err(msg: str) -> None:
|
||||
print(msg, file=sys.stderr)
|
||||
|
||||
|
||||
def cmd_proxy_start(args: Any) -> int:
|
||||
"""Run the proxy server in the foreground."""
|
||||
if not AIOHTTP_AVAILABLE:
|
||||
_print_aiohttp_missing()
|
||||
_err("hermes proxy requires aiohttp. Run `hermes setup` to install it.")
|
||||
return 1
|
||||
|
||||
provider = getattr(args, "provider", None) or "nous"
|
||||
try:
|
||||
adapter = get_adapter(provider)
|
||||
except ValueError as exc:
|
||||
print(f"Error: {exc}", file=sys.stderr)
|
||||
_err(f"Error: {exc}")
|
||||
return 2
|
||||
|
||||
if not adapter.is_authenticated():
|
||||
auth_hint = getattr(adapter, "auth_hint", f"hermes auth add {adapter.name}")
|
||||
print(
|
||||
f"Not logged into {adapter.display_name}. "
|
||||
f"Run `{auth_hint}` first.",
|
||||
file=sys.stderr,
|
||||
)
|
||||
_err(f"Not logged into {adapter.display_name}. Run `{auth_hint}` first.")
|
||||
return 2
|
||||
|
||||
host = getattr(args, "host", None) or DEFAULT_HOST
|
||||
port = getattr(args, "port", None) or DEFAULT_PORT
|
||||
|
||||
print(
|
||||
_err(
|
||||
f"Starting Hermes proxy for {adapter.display_name}\n"
|
||||
f" Listening on: http://{host}:{port}/v1\n"
|
||||
f" Forwarding to: (resolved per-request from your subscription)\n"
|
||||
f" Use any bearer token in the client — the proxy attaches your real credential.\n"
|
||||
f"\n"
|
||||
f"Press Ctrl+C to stop.",
|
||||
file=sys.stderr,
|
||||
f"Press Ctrl+C to stop."
|
||||
)
|
||||
|
||||
try:
|
||||
asyncio.run(run_server(adapter, host=host, port=port))
|
||||
except KeyboardInterrupt:
|
||||
print("\nproxy: stopped", file=sys.stderr)
|
||||
_err("\nproxy: stopped")
|
||||
except OSError as exc:
|
||||
print(f"proxy: failed to bind {host}:{port}: {exc}", file=sys.stderr)
|
||||
_err(f"proxy: failed to bind {host}:{port}: {exc}")
|
||||
return 1
|
||||
return 0
|
||||
|
||||
@@ -81,16 +70,11 @@ def cmd_proxy_status(args: Any) -> int:
|
||||
try:
|
||||
cred = adapter.get_credential()
|
||||
except Exception as exc:
|
||||
print(
|
||||
f" [{name:8s}] {adapter.display_name} — credentials need attention "
|
||||
f"({exc})"
|
||||
)
|
||||
print(f" [{name:8s}] {adapter.display_name} — credentials need attention ({exc})")
|
||||
continue
|
||||
expires = f" (bearer expires {cred.expires_at})" if cred.expires_at else ""
|
||||
print(f" [{name:8s}] {adapter.display_name} — ready{expires}")
|
||||
print(
|
||||
"\nStart the proxy with: hermes proxy start [--provider <name>]"
|
||||
)
|
||||
print("\nStart the proxy with: hermes proxy start [--provider <name>]")
|
||||
return 0
|
||||
|
||||
|
||||
@@ -103,17 +87,20 @@ def cmd_proxy_list_providers(args: Any) -> int:
|
||||
return 0
|
||||
|
||||
|
||||
_SUBCOMMANDS = {
|
||||
"start": cmd_proxy_start,
|
||||
"status": cmd_proxy_status,
|
||||
"providers": cmd_proxy_list_providers,
|
||||
"list": cmd_proxy_list_providers,
|
||||
}
|
||||
|
||||
|
||||
def cmd_proxy(args: Any) -> int:
|
||||
"""Dispatch ``hermes proxy <subcommand>``."""
|
||||
sub = getattr(args, "proxy_command", None)
|
||||
if sub == "start":
|
||||
return cmd_proxy_start(args)
|
||||
if sub == "status":
|
||||
return cmd_proxy_status(args)
|
||||
if sub in {"providers", "list"}:
|
||||
return cmd_proxy_list_providers(args)
|
||||
# No subcommand → print short help.
|
||||
print(
|
||||
"""Dispatch ``hermes proxy <subcommand>``; no/unknown subcommand prints the short help."""
|
||||
handler = _SUBCOMMANDS.get(getattr(args, "proxy_command", None))
|
||||
if handler is not None:
|
||||
return handler(args)
|
||||
_err(
|
||||
"hermes proxy — local OpenAI-compatible proxy that attaches your\n"
|
||||
"OAuth-authenticated provider credentials to outbound requests.\n"
|
||||
"\n"
|
||||
@@ -123,15 +110,11 @@ def cmd_proxy(args: Any) -> int:
|
||||
" hermes proxy status\n"
|
||||
" Show which upstream adapters are ready.\n"
|
||||
" hermes proxy providers\n"
|
||||
" List available upstream providers.\n",
|
||||
file=sys.stderr,
|
||||
" List available upstream providers.\n"
|
||||
)
|
||||
return 0
|
||||
|
||||
|
||||
__all__ = [
|
||||
"cmd_proxy",
|
||||
"cmd_proxy_start",
|
||||
"cmd_proxy_status",
|
||||
"cmd_proxy_list_providers",
|
||||
"cmd_proxy", "cmd_proxy_start", "cmd_proxy_status", "cmd_proxy_list_providers"
|
||||
]
|
||||
|
||||
@@ -1,13 +1,9 @@
|
||||
"""SSE ``[DONE]`` sentinel normalization for OpenAI-compatible proxies.
|
||||
|
||||
Strict OpenAI-compatible clients treat that shape as a truncated stream. This module watches the
|
||||
forwarded SSE byte stream and reports whether the proxy should append a single ``data: [DONE]``
|
||||
frame after a *clean* upstream EOF.
|
||||
|
||||
- Append ``[DONE]`` only after a complete terminal choice (``finish_reason`` non-null) **or** an
|
||||
upstream ``lastOne: true`` marker. - Never synthesize ``[DONE]`` after an error event, or when the
|
||||
stream was interrupted before clean EOF. - Never emit a second ``[DONE]`` when the upstream already
|
||||
sent one.
|
||||
Strict OpenAI clients treat a stream without ``data: [DONE]`` as truncated. The tracker watches the
|
||||
forwarded SSE bytes and says whether to append ONE ``[DONE]`` after a *clean* upstream EOF: only after
|
||||
a terminal choice (``finish_reason`` non-null) or ``lastOne: true``; never after an error event, an
|
||||
interrupted stream, or when the upstream already sent ``[DONE]``.
|
||||
"""
|
||||
|
||||
from __future__ import annotations
|
||||
@@ -49,60 +45,48 @@ class SseDoneTracker:
|
||||
"""Upstream stream ended via error/cancel — do not synthesize DONE."""
|
||||
self.interrupted = True
|
||||
|
||||
def _blocked(self) -> bool:
|
||||
return self.saw_done or self.saw_error_event or self.saw_malformed_event
|
||||
|
||||
def should_append_done(self) -> bool:
|
||||
"""True when a single terminal ``[DONE]`` should be appended."""
|
||||
if (
|
||||
self.interrupted
|
||||
or self.saw_done
|
||||
or self.saw_error_event
|
||||
or self.saw_malformed_event
|
||||
):
|
||||
if self.interrupted or self._blocked():
|
||||
return False
|
||||
# Flush any trailing line without a final newline (rare but valid),
|
||||
# then dispatch a final event that never saw its blank-line boundary.
|
||||
# Flush a trailing line without a final newline (rare but valid), then dispatch a final
|
||||
# event that never saw its blank-line boundary.
|
||||
if self._buf:
|
||||
self._consume_line(bytes(self._buf))
|
||||
self._buf.clear()
|
||||
self._dispatch_event()
|
||||
if self.saw_done or self.saw_error_event or self.saw_malformed_event:
|
||||
if self._blocked():
|
||||
return False
|
||||
return self.saw_terminal_finish or self.saw_last_one
|
||||
|
||||
def _consume_line(self, line: bytes) -> None:
|
||||
# Strip CR from CRLF-delimited SSE.
|
||||
if line.endswith(b"\r"):
|
||||
if line.endswith(b"\r"): # CRLF-delimited SSE
|
||||
line = line[:-1]
|
||||
if not line:
|
||||
# Blank line = SSE event boundary: dispatch accumulated data.
|
||||
if not line: # blank line = event boundary
|
||||
self._dispatch_event()
|
||||
return
|
||||
if not line.startswith(b"data:"):
|
||||
return
|
||||
# Per the SSE spec one event may span several consecutive ``data:``
|
||||
# lines whose payloads are joined with "\n" at dispatch time.
|
||||
# Parsing each line independently would misread a split JSON event
|
||||
# as two malformed fragments.
|
||||
# One event may span several ``data:`` lines joined with "\n" at dispatch; parsing each
|
||||
# line alone would misread a split JSON event as two malformed fragments.
|
||||
self._data_lines.append(line[5:].strip())
|
||||
|
||||
def _dispatch_event(self) -> None:
|
||||
if not self._data_lines:
|
||||
return
|
||||
payload = b"\n".join(self._data_lines)
|
||||
payload = b"\n".join(self._data_lines).strip()
|
||||
self._data_lines = []
|
||||
payload = payload.strip()
|
||||
if payload == b"[DONE]":
|
||||
self.saw_done = True
|
||||
return
|
||||
if not payload:
|
||||
return
|
||||
try:
|
||||
text = payload.decode("utf-8")
|
||||
except UnicodeDecodeError:
|
||||
self.saw_malformed_event = True
|
||||
return
|
||||
try:
|
||||
event = json.loads(text)
|
||||
except json.JSONDecodeError:
|
||||
event = json.loads(payload.decode("utf-8"))
|
||||
except (UnicodeDecodeError, json.JSONDecodeError):
|
||||
self.saw_malformed_event = True
|
||||
return
|
||||
if not isinstance(event, dict):
|
||||
@@ -110,23 +94,22 @@ class SseDoneTracker:
|
||||
if event.get("error") is not None:
|
||||
self.saw_error_event = True
|
||||
return
|
||||
# Accept integer-truthy sentinels too — relabelled upstreams have
|
||||
# been observed sending ``"lastOne": 1`` / ``"true"``.
|
||||
# Relabelled upstreams have been observed sending ``"lastOne": 1`` / ``"true"``.
|
||||
if event.get("lastOne") in (True, 1, "true"):
|
||||
self.saw_last_one = True
|
||||
for choice in event.get("choices") or []:
|
||||
if not isinstance(choice, dict):
|
||||
continue
|
||||
if choice.get("finish_reason") is not None:
|
||||
fr = choice.get("finish_reason")
|
||||
if fr is not None:
|
||||
self.saw_terminal_finish = True
|
||||
# OpenAI error-shaped finish reasons should not unlock DONE.
|
||||
fr = choice.get("finish_reason")
|
||||
if isinstance(fr, str) and fr.lower() in {"error", "provider_error"}:
|
||||
self.saw_error_event = True
|
||||
|
||||
|
||||
def content_type_is_sse(headers) -> bool:
|
||||
"""Return True when response headers advertise an SSE body."""
|
||||
"""True when response headers advertise an SSE body."""
|
||||
try:
|
||||
value = headers.get("Content-Type") or headers.get("content-type") or ""
|
||||
except Exception:
|
||||
@@ -135,7 +118,5 @@ def content_type_is_sse(headers) -> bool:
|
||||
|
||||
|
||||
__all__ = [
|
||||
"DONE_SSE_FRAME",
|
||||
"SseDoneTracker",
|
||||
"content_type_is_sse",
|
||||
"DONE_SSE_FRAME", "SseDoneTracker", "content_type_is_sse"
|
||||
]
|
||||
|
||||
@@ -6,12 +6,10 @@ import shutil
|
||||
import tarfile
|
||||
from pathlib import Path, PurePosixPath
|
||||
|
||||
# Pin a version we know patches cleanly. Update when a newer psutil
|
||||
# changes the marker line shape and we need to follow upstream.
|
||||
# Pinned to a version whose marker line patches cleanly; bump when upstream changes its shape.
|
||||
PSUTIL_URL = (
|
||||
"https://files.pythonhosted.org/packages/aa/c6/"
|
||||
"d1ddf4abb55e93cebc4f2ed8b5d6dbad109ecb8d63748dd2b20ab5e57ebe/"
|
||||
"psutil-7.2.2.tar.gz"
|
||||
"d1ddf4abb55e93cebc4f2ed8b5d6dbad109ecb8d63748dd2b20ab5e57ebe/psutil-7.2.2.tar.gz"
|
||||
)
|
||||
|
||||
MARKER = 'LINUX = sys.platform.startswith("linux")'
|
||||
@@ -26,9 +24,7 @@ def _normalize_member_parts(member_name: str) -> tuple[str, ...]:
|
||||
path = PurePosixPath(member_name)
|
||||
parts = tuple(part for part in path.parts if part not in ("", "."))
|
||||
if path.is_absolute() or ".." in parts or not parts:
|
||||
raise PsutilAndroidInstallError(
|
||||
f"Unsafe archive member path: {member_name!r}"
|
||||
)
|
||||
raise PsutilAndroidInstallError(f"Unsafe archive member path: {member_name!r}")
|
||||
return parts
|
||||
|
||||
|
||||
@@ -44,16 +40,12 @@ def _safe_extract_tar_gz(archive: Path, destination: Path) -> None:
|
||||
continue
|
||||
|
||||
if not member.isfile():
|
||||
raise PsutilAndroidInstallError(
|
||||
f"Unsupported archive member type: {member.name}"
|
||||
)
|
||||
raise PsutilAndroidInstallError(f"Unsupported archive member type: {member.name}")
|
||||
|
||||
target.parent.mkdir(parents=True, exist_ok=True)
|
||||
extracted = tf.extractfile(member)
|
||||
if extracted is None:
|
||||
raise PsutilAndroidInstallError(
|
||||
f"Cannot read archive member: {member.name}"
|
||||
)
|
||||
raise PsutilAndroidInstallError(f"Cannot read archive member: {member.name}")
|
||||
|
||||
with extracted, open(target, "wb") as dst:
|
||||
shutil.copyfileobj(extracted, dst)
|
||||
@@ -69,37 +61,24 @@ def prepare_patched_psutil_sdist(archive: Path, destination: Path) -> Path:
|
||||
_safe_extract_tar_gz(archive, destination)
|
||||
|
||||
src_roots = [
|
||||
path for path in destination.iterdir()
|
||||
if path.is_dir() and path.name.startswith("psutil-")
|
||||
path for path in destination.iterdir() if path.is_dir() and path.name.startswith("psutil-")
|
||||
]
|
||||
if not src_roots:
|
||||
raise PsutilAndroidInstallError(
|
||||
"psutil sdist did not contain a psutil-* directory"
|
||||
)
|
||||
raise PsutilAndroidInstallError("psutil sdist did not contain a psutil-* directory")
|
||||
|
||||
src_root = min(src_roots, key=lambda path: path.name)
|
||||
common_py = src_root / "psutil" / "_common.py"
|
||||
rel = common_py.relative_to(src_root)
|
||||
if not common_py.is_file():
|
||||
raise PsutilAndroidInstallError(
|
||||
f"psutil sdist did not contain {common_py.relative_to(src_root)!s}"
|
||||
)
|
||||
raise PsutilAndroidInstallError(f"psutil sdist did not contain {rel!s}")
|
||||
try:
|
||||
content = common_py.read_text(encoding="utf-8")
|
||||
except OSError as exc:
|
||||
raise PsutilAndroidInstallError(
|
||||
f"Failed to read {common_py.relative_to(src_root)!s}"
|
||||
) from exc
|
||||
raise PsutilAndroidInstallError(f"Failed to read {rel!s}") from exc
|
||||
if MARKER not in content:
|
||||
raise PsutilAndroidInstallError(
|
||||
"psutil Android compatibility patch marker not found"
|
||||
)
|
||||
raise PsutilAndroidInstallError("psutil Android compatibility patch marker not found")
|
||||
try:
|
||||
common_py.write_text(
|
||||
content.replace(MARKER, REPLACEMENT),
|
||||
encoding="utf-8",
|
||||
)
|
||||
common_py.write_text(content.replace(MARKER, REPLACEMENT), encoding="utf-8")
|
||||
except OSError as exc:
|
||||
raise PsutilAndroidInstallError(
|
||||
f"Failed to write {common_py.relative_to(src_root)!s}"
|
||||
) from exc
|
||||
raise PsutilAndroidInstallError(f"Failed to write {rel!s}") from exc
|
||||
return src_root
|
||||
|
||||
+21
-51
@@ -1,9 +1,5 @@
|
||||
"""Unified self-relaunch for Hermes CLI.
|
||||
|
||||
Preserves critical flags (--tui, --dev, --profile, --model, etc.) across process replacement so that
|
||||
``hermes sessions browse`` or post-setup relaunch doesn't silently drop the user's UI mode or other
|
||||
preferences.
|
||||
"""
|
||||
"""Unified self-relaunch for Hermes CLI: preserves inherited flags (--tui, --dev, --profile, --model…)
|
||||
across process replacement so ``hermes sessions browse`` / post-setup relaunch keep the user's mode."""
|
||||
|
||||
import os
|
||||
import shutil
|
||||
@@ -11,17 +7,13 @@ import sys
|
||||
from typing import Optional, Sequence
|
||||
|
||||
from hermes_cli._parser import (
|
||||
PRE_ARGPARSE_INHERITED_FLAGS,
|
||||
build_top_level_parser,
|
||||
PRE_ARGPARSE_INHERITED_FLAGS, build_top_level_parser
|
||||
)
|
||||
|
||||
|
||||
def _build_inherited_flag_table() -> list[tuple[str, bool]]:
|
||||
"""Build the ``(option_string, takes_value)`` table of flags that must survive a self-relaunch.
|
||||
|
||||
Introspects the real ``hermes`` parser: a flag participates iff its argparse Action carries
|
||||
``inherit_on_relaunch = True`` (set by ``_parser._inherited_flag``).
|
||||
"""
|
||||
"""``(option_string, takes_value)`` for every parser Action carrying ``inherit_on_relaunch``
|
||||
(set by ``_parser._inherited_flag``), plus the pre-argparse flags."""
|
||||
parser, _subparsers, chat_parser = build_top_level_parser()
|
||||
|
||||
table: list[tuple[str, bool]] = []
|
||||
@@ -73,35 +65,25 @@ def _extract_inherited_flags(argv: Sequence[str]) -> list[str]:
|
||||
|
||||
|
||||
def resolve_hermes_bin() -> Optional[str]:
|
||||
"""Find the hermes entry point.
|
||||
|
||||
Priority: 1. ``sys.argv[0]`` if it resolves to a real executable. 2. ``shutil.which("hermes")``
|
||||
on PATH. 3. ``None`` → caller should fall back to ``python -m hermes_cli.main``.
|
||||
"""
|
||||
"""Hermes entry point: ``sys.argv[0]`` if a real executable, else ``which hermes``, else ``None``
|
||||
(caller falls back to ``python -m hermes_cli.main``)."""
|
||||
argv0 = sys.argv[0]
|
||||
_is_windows = sys.platform == "win32"
|
||||
|
||||
def _is_python_script(p: str) -> bool:
|
||||
return p.lower().endswith((".py", ".pyc"))
|
||||
|
||||
# Absolute path to an executable (covers nix store, venv wrappers, etc.)
|
||||
if os.path.isabs(argv0) and os.path.isfile(argv0) and os.access(argv0, os.X_OK):
|
||||
if not (_is_windows and _is_python_script(argv0)):
|
||||
return argv0
|
||||
|
||||
# Relative path — resolve against CWD
|
||||
# Absolute executable (nix store, venv wrappers, …), then relative-to-CWD, then PATH.
|
||||
if (
|
||||
os.path.isabs(argv0) and os.path.isfile(argv0) and os.access(argv0, os.X_OK)
|
||||
and not (_is_windows and _is_python_script(argv0))
|
||||
):
|
||||
return argv0
|
||||
if not argv0.startswith("-") and os.path.isfile(argv0):
|
||||
abs_path = os.path.abspath(argv0)
|
||||
if os.access(abs_path, os.X_OK):
|
||||
if not (_is_windows and _is_python_script(abs_path)):
|
||||
return abs_path
|
||||
|
||||
# PATH lookup
|
||||
path_bin = shutil.which("hermes")
|
||||
if path_bin:
|
||||
return path_bin
|
||||
|
||||
return None
|
||||
if os.access(abs_path, os.X_OK) and not (_is_windows and _is_python_script(abs_path)):
|
||||
return abs_path
|
||||
return shutil.which("hermes") or None
|
||||
|
||||
|
||||
def build_relaunch_argv(
|
||||
@@ -112,12 +94,7 @@ def build_relaunch_argv(
|
||||
) -> list[str]:
|
||||
"""Construct an argv list for replacing the current process with hermes."""
|
||||
bin_path = resolve_hermes_bin()
|
||||
|
||||
if bin_path:
|
||||
argv = [bin_path]
|
||||
else:
|
||||
argv = [sys.executable, "-m", "hermes_cli.main"]
|
||||
|
||||
argv = [bin_path] if bin_path else [sys.executable, "-m", "hermes_cli.main"]
|
||||
src = list(original_argv) if original_argv is not None else list(sys.argv[1:])
|
||||
|
||||
if preserve_inherited:
|
||||
@@ -135,18 +112,13 @@ def relaunch(
|
||||
) -> None:
|
||||
"""Replace the current process with a fresh hermes invocation.
|
||||
|
||||
On POSIX we use ``os.execvp`` which replaces the running process with the new one in place —
|
||||
same PID, no double-fork. That's what the relaunch contract wants: "run hermes again as if the
|
||||
user had typed the new argv".
|
||||
|
||||
Windows has no native exec semantics — ``os.execvp`` on Windows *emulates* exec by spawning the
|
||||
child and exiting the parent, but only works when the target is a real Win32 executable.
|
||||
POSIX: ``os.execvp`` in place (same PID, no double-fork). Windows has no real exec — its
|
||||
``execvp`` emulation only works for a real Win32 executable, so spawn + exit instead.
|
||||
"""
|
||||
new_argv = build_relaunch_argv(
|
||||
extra_args, preserve_inherited=preserve_inherited, original_argv=original_argv
|
||||
)
|
||||
if sys.platform == "win32":
|
||||
# Windows: subprocess + exit, because execvp can't swap to .cmd/.exe shims.
|
||||
import subprocess
|
||||
try:
|
||||
result = subprocess.run(new_argv)
|
||||
@@ -154,10 +126,8 @@ def relaunch(
|
||||
except KeyboardInterrupt:
|
||||
sys.exit(130)
|
||||
except OSError as exc:
|
||||
# Surface a helpful error rather than the raw OSError — the
|
||||
# caller used to see ``[Errno 8] Exec format error`` which is
|
||||
# cryptic. Common causes: ``hermes`` not on PATH yet (install
|
||||
# hasn't propagated User PATH into this shell) or a stale shim.
|
||||
# Raw ``[Errno 8] Exec format error`` is cryptic; usual causes are ``hermes`` not on
|
||||
# PATH yet (install hasn't propagated User PATH into this shell) or a stale shim.
|
||||
print(
|
||||
f"\nHermes relaunch failed: {exc}\n"
|
||||
f"Command: {' '.join(new_argv)}\n"
|
||||
|
||||
@@ -8,44 +8,29 @@ from typing import Any
|
||||
|
||||
RELAY_PLUGINS_CONFIG_ENV = "HERMES_NEMO_RELAY_PLUGINS_TOML"
|
||||
|
||||
LEGACY_RELAY_PLUGIN_KEYS = frozenset(
|
||||
{
|
||||
"nemo_relay",
|
||||
"observability/nemo_relay",
|
||||
}
|
||||
)
|
||||
LEGACY_RELAY_PLUGIN_KEYS = frozenset({"nemo_relay", "observability/nemo_relay"})
|
||||
|
||||
LEGACY_RELAY_EXPORT_ENV_VARS = frozenset(
|
||||
{
|
||||
"HERMES_NEMO_RELAY_ATOF_ENABLED",
|
||||
"HERMES_NEMO_RELAY_ATOF_OUTPUT_DIRECTORY",
|
||||
"HERMES_NEMO_RELAY_ATOF_FILENAME",
|
||||
"HERMES_NEMO_RELAY_ATOF_MODE",
|
||||
"HERMES_NEMO_RELAY_ATIF_ENABLED",
|
||||
"HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY",
|
||||
"HERMES_NEMO_RELAY_ATIF_FILENAME_TEMPLATE",
|
||||
"HERMES_NEMO_RELAY_ATIF_AGENT_NAME",
|
||||
"HERMES_NEMO_RELAY_ATIF_AGENT_VERSION",
|
||||
"HERMES_NEMO_RELAY_ATIF_EXPORT_TIMEOUT_S",
|
||||
"HERMES_NEMO_RELAY_ATIF_MODEL_NAME",
|
||||
"HERMES_NEMO_RELAY_ATIF_SUBAGENT_EXPORT_MODE",
|
||||
}
|
||||
)
|
||||
LEGACY_RELAY_EXPORT_ENV_VARS = frozenset({
|
||||
"HERMES_NEMO_RELAY_ATOF_ENABLED",
|
||||
"HERMES_NEMO_RELAY_ATOF_OUTPUT_DIRECTORY",
|
||||
"HERMES_NEMO_RELAY_ATOF_FILENAME",
|
||||
"HERMES_NEMO_RELAY_ATOF_MODE",
|
||||
"HERMES_NEMO_RELAY_ATIF_ENABLED",
|
||||
"HERMES_NEMO_RELAY_ATIF_OUTPUT_DIRECTORY",
|
||||
"HERMES_NEMO_RELAY_ATIF_FILENAME_TEMPLATE",
|
||||
"HERMES_NEMO_RELAY_ATIF_AGENT_NAME",
|
||||
"HERMES_NEMO_RELAY_ATIF_AGENT_VERSION",
|
||||
"HERMES_NEMO_RELAY_ATIF_EXPORT_TIMEOUT_S",
|
||||
"HERMES_NEMO_RELAY_ATIF_MODEL_NAME",
|
||||
"HERMES_NEMO_RELAY_ATIF_SUBAGENT_EXPORT_MODE",
|
||||
})
|
||||
|
||||
|
||||
def legacy_relay_plugin_keys(values: Any) -> tuple[str, ...]:
|
||||
"""Return removed Relay plugin identities present in a config value."""
|
||||
if not isinstance(values, (list, tuple, set, frozenset)):
|
||||
return ()
|
||||
return tuple(
|
||||
sorted(
|
||||
{
|
||||
value
|
||||
for value in values
|
||||
if isinstance(value, str) and value in LEGACY_RELAY_PLUGIN_KEYS
|
||||
}
|
||||
)
|
||||
)
|
||||
return tuple(sorted({v for v in values if isinstance(v, str) and v in LEGACY_RELAY_PLUGIN_KEYS}))
|
||||
|
||||
|
||||
def configured_legacy_relay_env_vars(
|
||||
@@ -54,10 +39,7 @@ def configured_legacy_relay_env_vars(
|
||||
"""Return non-empty legacy Relay exporter variables in *env*."""
|
||||
if env is None:
|
||||
return ()
|
||||
return tuple(
|
||||
sorted(
|
||||
name
|
||||
for name in LEGACY_RELAY_EXPORT_ENV_VARS
|
||||
if env.get(name) is not None and str(env[name]).strip()
|
||||
)
|
||||
)
|
||||
return tuple(sorted(
|
||||
name for name in LEGACY_RELAY_EXPORT_ENV_VARS
|
||||
if env.get(name) is not None and str(env[name]).strip()
|
||||
))
|
||||
|
||||
@@ -22,17 +22,15 @@ _MISSING = object()
|
||||
def configured_nofile_soft_limit(
|
||||
config: Mapping[str, Any] | None = None,
|
||||
) -> int | None:
|
||||
"""Resolve ``runtime.nofile_soft_limit`` from a loaded config.
|
||||
"""``runtime.nofile_soft_limit`` from a loaded config, or ``None`` when disabled/unresolvable.
|
||||
|
||||
A missing key uses the default. Explicit ``0``, ``false``, and ``null`` disable the
|
||||
adjustment; other non-integer or negative values are ignored (caller fails open).
|
||||
Used by service-definition generators (e.g. the launchd plist) so persisted service limits
|
||||
and the in-process floor share one knob. ``None`` when disabled or unresolvable.
|
||||
Missing key → default. Explicit ``0``/``false``/``null`` disable; other non-int or negative
|
||||
values are ignored (caller fails open). Shared by the in-process floor and service-definition
|
||||
generators (launchd plist) so both use one knob.
|
||||
"""
|
||||
if config is None:
|
||||
try:
|
||||
# Use Hermes's real, profile-aware loader rather than reading YAML
|
||||
# here. This also applies managed-scope overlays and defaults.
|
||||
# Profile-aware loader (applies managed-scope overlays and defaults).
|
||||
from hermes_cli.config import load_config_readonly
|
||||
|
||||
config = load_config_readonly()
|
||||
@@ -60,15 +58,11 @@ def configured_nofile_soft_limit(
|
||||
def apply_nofile_soft_limit(
|
||||
config: Mapping[str, Any] | None = None,
|
||||
) -> bool:
|
||||
"""Raise this process's ``RLIMIT_NOFILE`` soft limit when possible.
|
||||
"""Best-effort raise of this process's ``RLIMIT_NOFILE`` soft limit; ``True`` iff changed.
|
||||
|
||||
The target defaults to :data:`DEFAULT_NOFILE_SOFT_LIMIT` and can be set with
|
||||
``runtime.nofile_soft_limit``. The target is clamped to a finite hard limit, never lowers an
|
||||
existing higher soft limit, and returns ``False`` for an explicit opt-out or when the
|
||||
platform/sandbox refuses the operation.
|
||||
|
||||
This is intentionally best-effort. Unsupported platforms, malformed settings, and denied
|
||||
``setrlimit`` calls must never prevent a server from starting.
|
||||
Target = ``runtime.nofile_soft_limit`` (default :data:`DEFAULT_NOFILE_SOFT_LIMIT`), clamped
|
||||
to a finite hard limit; never lowers a higher soft limit. Unsupported platforms, malformed
|
||||
settings, and denied ``setrlimit`` must never prevent a server from starting.
|
||||
"""
|
||||
if _resource is None:
|
||||
return False
|
||||
@@ -80,9 +74,8 @@ def apply_nofile_soft_limit(
|
||||
try:
|
||||
nofile = _resource.RLIMIT_NOFILE
|
||||
current_soft, current_hard = _resource.getrlimit(nofile)
|
||||
# On platforms where RLIM_INFINITY is represented as -1, ordinary
|
||||
# integer ordering would make an unlimited soft limit look lower than
|
||||
# every positive target. Never replace infinity with a finite limit.
|
||||
# RLIM_INFINITY may be -1, which ordinary ordering would treat as "lower than any
|
||||
# target"; never replace infinity with a finite limit.
|
||||
infinity = getattr(_resource, "RLIM_INFINITY", object())
|
||||
if current_soft == infinity or current_soft >= target:
|
||||
return False
|
||||
@@ -93,14 +86,10 @@ def apply_nofile_soft_limit(
|
||||
_resource.setrlimit(nofile, (new_soft, current_hard))
|
||||
return True
|
||||
except Exception:
|
||||
# This helper runs before server startup and must fail open for
|
||||
# unsupported/sandboxed environments and denied resource changes.
|
||||
logger.debug("Could not raise RLIMIT_NOFILE soft limit", exc_info=True)
|
||||
return False
|
||||
|
||||
|
||||
__all__ = [
|
||||
"DEFAULT_NOFILE_SOFT_LIMIT",
|
||||
"apply_nofile_soft_limit",
|
||||
"configured_nofile_soft_limit",
|
||||
"DEFAULT_NOFILE_SOFT_LIMIT", "apply_nofile_soft_limit", "configured_nofile_soft_limit"
|
||||
]
|
||||
|
||||
Reference in New Issue
Block a user