refactor(hermes_cli): fold short multi-line calls/signatures in g5 files (AST-identical)
This commit is contained in:
@@ -7,10 +7,7 @@ from hermes_cli.proxy.adapters.nous_portal import NousPortalAdapter
|
||||
from hermes_cli.proxy.adapters.xai import XAIGrokAdapter
|
||||
|
||||
# Keyed by the ``hermes proxy start --provider <name>`` value.
|
||||
ADAPTERS: Dict[str, Type[UpstreamAdapter]] = {
|
||||
"nous": NousPortalAdapter,
|
||||
"xai": XAIGrokAdapter,
|
||||
}
|
||||
ADAPTERS: Dict[str, Type[UpstreamAdapter]] = {"nous": NousPortalAdapter, "xai": XAIGrokAdapter}
|
||||
|
||||
|
||||
def get_adapter(name: str) -> UpstreamAdapter:
|
||||
|
||||
@@ -51,9 +51,7 @@ class NousPortalAdapter(UpstreamAdapter):
|
||||
def is_authenticated(self) -> bool:
|
||||
# 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"))
|
||||
)
|
||||
return bool(state.get("agent_key") or (state.get("refresh_token") and state.get("access_token")))
|
||||
|
||||
def get_credential(self) -> UpstreamCredential:
|
||||
return self._get_credential()
|
||||
@@ -67,9 +65,7 @@ class NousPortalAdapter(UpstreamAdapter):
|
||||
logger.info("proxy: Nous upstream rejected bearer; force-refreshing invoke JWT")
|
||||
return self._get_credential(force_refresh=True)
|
||||
|
||||
def _get_credential(
|
||||
self, *, force_refresh: bool = False
|
||||
) -> UpstreamCredential:
|
||||
def _get_credential(self, *, force_refresh: bool = False) -> UpstreamCredential:
|
||||
with self._lock:
|
||||
state = self._read_state()
|
||||
if state is None:
|
||||
@@ -80,9 +76,7 @@ class NousPortalAdapter(UpstreamAdapter):
|
||||
except Exception as exc:
|
||||
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"
|
||||
)
|
||||
self._save_state(state, quarantine_error=exc, quarantine_reason="proxy_refresh_failure")
|
||||
raise RuntimeError(f"Failed to refresh Nous Portal credentials: {exc}") from exc
|
||||
|
||||
runtime_key = refreshed.get("api_key")
|
||||
@@ -102,9 +96,7 @@ class NousPortalAdapter(UpstreamAdapter):
|
||||
or DEFAULT_NOUS_INFERENCE_URL
|
||||
).rstrip("/")
|
||||
|
||||
return UpstreamCredential(
|
||||
bearer=runtime_key, base_url=base_url, expires_at=refreshed.get("expires_at")
|
||||
)
|
||||
return UpstreamCredential(bearer=runtime_key, base_url=base_url, expires_at=refreshed.get("expires_at"))
|
||||
|
||||
# auth.json access — kept local so hermes_cli.auth's public surface does not grow.
|
||||
|
||||
|
||||
@@ -85,10 +85,7 @@ class XAIGrokAdapter(UpstreamAdapter):
|
||||
retry_cred = self._credential_from_entry(refreshed)
|
||||
if retry_cred.bearer == failed_credential.bearer:
|
||||
return None
|
||||
logger.info(
|
||||
"proxy: xAI upstream returned %s; retrying with rotated pool credential",
|
||||
status_code,
|
||||
)
|
||||
logger.info("proxy: xAI upstream returned %s; retrying with rotated pool credential", status_code)
|
||||
return retry_cred
|
||||
|
||||
def _load_pool(self) -> Optional[CredentialPool]:
|
||||
@@ -111,9 +108,7 @@ class XAIGrokAdapter(UpstreamAdapter):
|
||||
).strip().rstrip("/")
|
||||
|
||||
return UpstreamCredential(
|
||||
bearer=bearer,
|
||||
base_url=base_url or DEFAULT_XAI_OAUTH_BASE_URL,
|
||||
expires_at=entry.expires_at,
|
||||
bearer=bearer, base_url=base_url or DEFAULT_XAI_OAUTH_BASE_URL, expires_at=entry.expires_at
|
||||
)
|
||||
|
||||
|
||||
|
||||
@@ -115,6 +115,4 @@ def cmd_proxy(args: Any) -> int:
|
||||
return 0
|
||||
|
||||
|
||||
__all__ = [
|
||||
"cmd_proxy", "cmd_proxy_start", "cmd_proxy_status", "cmd_proxy_list_providers"
|
||||
]
|
||||
__all__ = ["cmd_proxy", "cmd_proxy_start", "cmd_proxy_status", "cmd_proxy_list_providers"]
|
||||
|
||||
+10
-36
@@ -44,9 +44,7 @@ MAX_REQUEST_BYTES = 10_000_000
|
||||
|
||||
def _require_aiohttp() -> None:
|
||||
if not AIOHTTP_AVAILABLE:
|
||||
raise RuntimeError(
|
||||
"aiohttp is required for `hermes proxy`. Run `hermes setup` to install it."
|
||||
)
|
||||
raise RuntimeError("aiohttp is required for `hermes proxy`. Run `hermes setup` to install it.")
|
||||
|
||||
|
||||
def _json_error(status: int, message: str, code: str = "proxy_error") -> "web.Response":
|
||||
@@ -68,23 +66,14 @@ async def _open_upstream(request: "web.Request", rel_path: str, body: bytes, cre
|
||||
upstream_url = f"{upstream_url}?{request.query_string}"
|
||||
fwd_headers = _filter_headers(request.headers)
|
||||
fwd_headers["Authorization"] = f"{cred.token_type} {cred.bearer}"
|
||||
logger.debug(
|
||||
"proxy: forwarding %s %s -> %s (body=%d bytes)",
|
||||
request.method, rel_path, upstream_url, len(body),
|
||||
)
|
||||
logger.debug("proxy: forwarding %s %s -> %s (body=%d bytes)", request.method, rel_path, upstream_url, len(body))
|
||||
try:
|
||||
session = aiohttp.ClientSession(
|
||||
timeout=aiohttp.ClientTimeout(total=None, sock_connect=15, sock_read=300)
|
||||
)
|
||||
session = aiohttp.ClientSession(timeout=aiohttp.ClientTimeout(total=None, sock_connect=15, sock_read=300))
|
||||
except Exception as exc: # pragma: no cover - aiohttp setup issue
|
||||
return _json_error(500, f"proxy session init failed: {exc}"), None
|
||||
try:
|
||||
upstream_resp = await session.request(
|
||||
request.method,
|
||||
upstream_url,
|
||||
data=body if body else None,
|
||||
headers=fwd_headers,
|
||||
allow_redirects=False,
|
||||
request.method, upstream_url, data=body if body else None, headers=fwd_headers, allow_redirects=False
|
||||
)
|
||||
except RuntimeError as exc:
|
||||
await session.close()
|
||||
@@ -92,9 +81,7 @@ async def _open_upstream(request: "web.Request", rel_path: str, body: bytes, cre
|
||||
except aiohttp.ClientError as exc:
|
||||
await session.close()
|
||||
logger.warning("proxy: upstream connection failed: %s", exc)
|
||||
return _json_error(
|
||||
502, f"upstream connection failed: {exc}", code="upstream_unreachable"
|
||||
), None
|
||||
return _json_error(502, f"upstream connection failed: {exc}", code="upstream_unreachable"), None
|
||||
except asyncio.TimeoutError:
|
||||
await session.close()
|
||||
return _json_error(504, "upstream request timed out", code="upstream_timeout"), None
|
||||
@@ -108,8 +95,7 @@ async def _stream_back(request: "web.Request", session, upstream_resp) -> "web.S
|
||||
"""Relay status + filtered headers, then the body chunk-by-chunk, appending a missing SSE
|
||||
``[DONE]`` only after a clean EOF."""
|
||||
resp = web.StreamResponse(
|
||||
status=upstream_resp.status,
|
||||
headers=_filter_headers(upstream_resp.headers, _RESPONSE_DROP_HEADERS),
|
||||
status=upstream_resp.status, headers=_filter_headers(upstream_resp.headers, _RESPONSE_DROP_HEADERS)
|
||||
)
|
||||
await resp.prepare(request)
|
||||
done_tracker: Optional[SseDoneTracker] = None
|
||||
@@ -153,18 +139,14 @@ def create_app(adapter: UpstreamAdapter) -> "web.Application":
|
||||
|
||||
async def handle_health(request: "web.Request") -> "web.Response":
|
||||
authenticated = await asyncio.to_thread(adapter.is_authenticated)
|
||||
return web.json_response(
|
||||
{"status": "ok", "upstream": adapter.display_name, "authenticated": authenticated}
|
||||
)
|
||||
return web.json_response({"status": "ok", "upstream": adapter.display_name, "authenticated": authenticated})
|
||||
|
||||
async def handle_proxy(request: "web.Request") -> "web.StreamResponse":
|
||||
rel_path = "/" + request.match_info.get("tail", "").lstrip("/")
|
||||
if rel_path not in adapter.allowed_paths:
|
||||
allowed = ", ".join(sorted(adapter.allowed_paths))
|
||||
return _json_error(
|
||||
404,
|
||||
f"Path /v1{rel_path} is not forwarded by this proxy. Allowed: {allowed}",
|
||||
code="path_not_allowed",
|
||||
404, f"Path /v1{rel_path} is not forwarded by this proxy. Allowed: {allowed}", code="path_not_allowed"
|
||||
)
|
||||
try:
|
||||
cred = await asyncio.to_thread(adapter.get_credential)
|
||||
@@ -184,9 +166,7 @@ def create_app(adapter: UpstreamAdapter) -> "web.Application":
|
||||
# POST under the auth lock; xAI: pool rotation).
|
||||
try:
|
||||
retry_cred = await asyncio.to_thread(
|
||||
adapter.get_retry_credential,
|
||||
failed_credential=cred,
|
||||
status_code=upstream_resp.status,
|
||||
adapter.get_retry_credential, failed_credential=cred, status_code=upstream_resp.status
|
||||
)
|
||||
except Exception as exc:
|
||||
logger.warning("proxy: retry credential resolution failed: %s", exc)
|
||||
@@ -238,10 +218,4 @@ async def run_server(
|
||||
await runner.cleanup()
|
||||
|
||||
|
||||
__all__ = [
|
||||
"create_app",
|
||||
"run_server",
|
||||
"DEFAULT_HOST",
|
||||
"DEFAULT_PORT",
|
||||
"AIOHTTP_AVAILABLE",
|
||||
]
|
||||
__all__ = ["create_app", "run_server", "DEFAULT_HOST", "DEFAULT_PORT", "AIOHTTP_AVAILABLE"]
|
||||
|
||||
@@ -117,6 +117,4 @@ def content_type_is_sse(headers) -> bool:
|
||||
return "text/event-stream" in str(value).lower()
|
||||
|
||||
|
||||
__all__ = [
|
||||
"DONE_SSE_FRAME", "SseDoneTracker", "content_type_is_sse"
|
||||
]
|
||||
__all__ = ["DONE_SSE_FRAME", "SseDoneTracker", "content_type_is_sse"]
|
||||
|
||||
+17
-51
@@ -26,8 +26,7 @@ def register_cli(parent_parser: argparse.ArgumentParser) -> None:
|
||||
# ``--help`` order, so keep it stable.
|
||||
commands = [
|
||||
("install", f"Download iron-proxy binary (v{ip._IRON_PROXY_VERSION})", cmd_install, [
|
||||
("--force", dict(action="store_true",
|
||||
help="Re-download even if a managed copy already exists")),
|
||||
("--force", dict(action="store_true", help="Re-download even if a managed copy already exists")),
|
||||
]),
|
||||
("setup", "Interactive wizard: install + CA + mint tokens + write config", cmd_setup, [
|
||||
("--tunnel-port", dict(
|
||||
@@ -119,9 +118,7 @@ def cmd_setup(args: argparse.Namespace) -> int:
|
||||
return 1
|
||||
_setup_restart_daemon(console, args, proxy_cfg)
|
||||
console.print()
|
||||
console.print(
|
||||
"[green]✓ iron-proxy is configured.[/green] Sandboxes will route outbound traffic through it."
|
||||
)
|
||||
console.print("[green]✓ iron-proxy is configured.[/green] Sandboxes will route outbound traffic through it.")
|
||||
console.print(
|
||||
" Start: [cyan]hermes egress start[/cyan]\n"
|
||||
" Restart: [cyan]hermes egress restart[/cyan] (after any re-setup)\n"
|
||||
@@ -174,10 +171,7 @@ def _setup_mint_tokens(console: Console, args: argparse.Namespace):
|
||||
# runs, NOT exported into an interactive shell); backfill so discovery sees them.
|
||||
loaded = _load_env_file_into_environ()
|
||||
if loaded:
|
||||
console.print(
|
||||
f" [dim]Loaded {loaded} provider key name(s) from "
|
||||
f"~/.hermes/.env for discovery.[/dim]"
|
||||
)
|
||||
console.print(f" [dim]Loaded {loaded} provider key name(s) from " f"~/.hermes/.env for discovery.[/dim]")
|
||||
|
||||
discovered = ip.discover_provider_mappings(available_env_names=available_env_names or None)
|
||||
|
||||
@@ -225,9 +219,7 @@ def _setup_mint_tokens(console: Console, args: argparse.Namespace):
|
||||
uncovered = ip.discover_uncovered_providers(available_env_names=available_env_names or None)
|
||||
if uncovered:
|
||||
console.print()
|
||||
console.print(
|
||||
" [yellow]⚠[/yellow] Detected provider env vars that the proxy does not yet cover:"
|
||||
)
|
||||
console.print(" [yellow]⚠[/yellow] Detected provider env vars that the proxy does not yet cover:")
|
||||
for name in uncovered:
|
||||
console.print(f" - {name}")
|
||||
console.print(
|
||||
@@ -266,9 +258,7 @@ def _setup_write_config(console: Console, args: argparse.Namespace, mappings, ca
|
||||
proxy_cfg["tunnel_port"] = tunnel_port
|
||||
|
||||
extra_hosts = list(proxy_cfg.get("extra_allowed_hosts") or [])
|
||||
allowed = list(ip._DEFAULT_ALLOWED_HOSTS) + [
|
||||
h for h in extra_hosts if h not in ip._DEFAULT_ALLOWED_HOSTS
|
||||
]
|
||||
allowed = list(ip._DEFAULT_ALLOWED_HOSTS) + [h for h in extra_hosts if h not in ip._DEFAULT_ALLOWED_HOSTS]
|
||||
|
||||
# Pre-create the audit log 0o600. The pinned v0.39 daemon never writes it (reserved for
|
||||
# v0.40+ per-request records), so a pre-create failure is a WARNING, not a setup abort.
|
||||
@@ -355,16 +345,10 @@ def _setup_restart_daemon(console: Console, args: argparse.Namespace, proxy_cfg:
|
||||
|
||||
if do_restart:
|
||||
try:
|
||||
new_status = ip.start_proxy(
|
||||
install_if_missing=bool(proxy_cfg.get("auto_install", True)),
|
||||
)
|
||||
new_status = ip.start_proxy(install_if_missing=bool(proxy_cfg.get("auto_install", True)))
|
||||
except Exception as exc: # noqa: BLE001 — user-facing funnel
|
||||
console.print(
|
||||
f" [yellow]⚠ could not start iron-proxy with the new config: {exc}[/yellow]"
|
||||
)
|
||||
console.print(
|
||||
" Run [cyan]hermes egress start[/cyan] manually before launching new Docker sandboxes."
|
||||
)
|
||||
console.print(f" [yellow]⚠ could not start iron-proxy with the new config: {exc}[/yellow]")
|
||||
console.print(" Run [cyan]hermes egress start[/cyan] manually before launching new Docker sandboxes.")
|
||||
else:
|
||||
listening = "listening" if new_status.listening else "not yet listening"
|
||||
verb = "restarted" if was_running else "started"
|
||||
@@ -392,9 +376,7 @@ def cmd_start(args: argparse.Namespace) -> int:
|
||||
# the rotation guarantee distinguishing it from ``env``.
|
||||
credential_source = proxy_cfg.get("credential_source", "env")
|
||||
bw_cfg = (cfg.get("secrets") or {}).get("bitwarden")
|
||||
refresh_bw = (
|
||||
credential_source == "bitwarden" and bw_cfg is not None and bool(bw_cfg.get("enabled"))
|
||||
)
|
||||
refresh_bw = (credential_source == "bitwarden" and bw_cfg is not None and bool(bw_cfg.get("enabled")))
|
||||
# Silent-degrade guard: bitwarden mode chosen but secrets.bitwarden disabled/removed. Refuse
|
||||
# (quietly starting on host env is the bug class BW mode exists to defeat) unless the
|
||||
# documented escape hatch is set.
|
||||
@@ -450,12 +432,9 @@ def cmd_start(args: argparse.Namespace) -> int:
|
||||
if not status.pid:
|
||||
console.print("[red]✗ iron-proxy did not come up cleanly[/red]")
|
||||
return 1
|
||||
listening = (
|
||||
"[green]listening[/green]" if status.listening else "[yellow]not yet listening[/yellow]"
|
||||
)
|
||||
listening = ("[green]listening[/green]" if status.listening else "[yellow]not yet listening[/yellow]")
|
||||
console.print(
|
||||
f"[green]✓[/green] iron-proxy running pid={status.pid} "
|
||||
f"port={status.tunnel_port} {listening}"
|
||||
f"[green]✓[/green] iron-proxy running pid={status.pid} " f"port={status.tunnel_port} {listening}"
|
||||
)
|
||||
return 0
|
||||
|
||||
@@ -489,9 +468,7 @@ def cmd_reload(args: argparse.Namespace) -> int:
|
||||
except Exception as exc: # noqa: BLE001 — top-level user-facing funnel
|
||||
console.print(f"[red]✗ reload failed:[/red] {exc}")
|
||||
return 1
|
||||
console.print(
|
||||
"[green]✓[/green] iron-proxy ruleset reloaded in-place (no restart, connections preserved)"
|
||||
)
|
||||
console.print("[green]✓[/green] iron-proxy ruleset reloaded in-place (no restart, connections preserved)")
|
||||
console.print(
|
||||
"[dim]Note: new upstream secrets (rotated keys, new providers) "
|
||||
"still need `hermes egress restart` — the daemon reads real "
|
||||
@@ -509,9 +486,7 @@ def format_status_text(*, show_tokens: bool = False) -> str:
|
||||
lines = ["Egress proxy status", ""]
|
||||
lines.extend(
|
||||
f"{label}: {value}"
|
||||
for label, value in _status_rows(
|
||||
proxy_cfg, status, yn=lambda v: "yes" if v else "no", dim=lambda t: t
|
||||
)
|
||||
for label, value in _status_rows(proxy_cfg, status, yn=lambda v: "yes" if v else "no", dim=lambda t: t)
|
||||
)
|
||||
lines.append("Scope: Docker backend only in this release")
|
||||
|
||||
@@ -524,9 +499,7 @@ def format_status_text(*, show_tokens: bool = False) -> str:
|
||||
|
||||
uncovered = ip.discover_uncovered_providers()
|
||||
if uncovered:
|
||||
lines.extend([
|
||||
"", "Uncovered providers (real credentials still visible inside the sandbox):"
|
||||
])
|
||||
lines.extend(["", "Uncovered providers (real credentials still visible inside the sandbox):"])
|
||||
for name in uncovered:
|
||||
lines.append(f" - {name}")
|
||||
|
||||
@@ -573,9 +546,7 @@ def cmd_status(args: argparse.Namespace) -> int:
|
||||
uncovered = ip.discover_uncovered_providers()
|
||||
if uncovered:
|
||||
console.print()
|
||||
console.print(
|
||||
"[yellow]Uncovered providers[/yellow] (real credentials still visible inside the sandbox):"
|
||||
)
|
||||
console.print("[yellow]Uncovered providers[/yellow] (real credentials still visible inside the sandbox):")
|
||||
for name in uncovered:
|
||||
console.print(f" - {name}")
|
||||
|
||||
@@ -623,9 +594,7 @@ def _bitwarden_env_names(console: Console) -> Optional[List[str]]:
|
||||
cfg = load_config()
|
||||
bw_cfg = (cfg.get("secrets") or {}).get("bitwarden") or {}
|
||||
if not bw_cfg.get("enabled"):
|
||||
console.print(
|
||||
" [red]✗ --from-bitwarden requested but secrets.bitwarden.enabled is false.[/red]"
|
||||
)
|
||||
console.print(" [red]✗ --from-bitwarden requested but secrets.bitwarden.enabled is false.[/red]")
|
||||
console.print(" Run `hermes secrets bitwarden setup` first, or omit --from-bitwarden.")
|
||||
return None
|
||||
try:
|
||||
@@ -639,10 +608,7 @@ def _bitwarden_env_names(console: Console) -> Optional[List[str]]:
|
||||
)
|
||||
return None
|
||||
secrets, _ = bw.fetch_bitwarden_secrets(
|
||||
access_token=access_token,
|
||||
project_id=bw_cfg.get("project_id", ""),
|
||||
cache_ttl_seconds=0,
|
||||
use_cache=False,
|
||||
access_token=access_token, project_id=bw_cfg.get("project_id", ""), cache_ttl_seconds=0, use_cache=False
|
||||
)
|
||||
names = list(secrets.keys())
|
||||
if not names:
|
||||
|
||||
@@ -60,9 +60,7 @@ def prepare_patched_psutil_sdist(archive: Path, destination: Path) -> Path:
|
||||
"""Safely extract the pinned psutil sdist and patch it for Android."""
|
||||
_safe_extract_tar_gz(archive, destination)
|
||||
|
||||
src_roots = [
|
||||
path for path in destination.iterdir() if path.is_dir() and path.name.startswith("psutil-")
|
||||
]
|
||||
src_roots = [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")
|
||||
|
||||
|
||||
+3
-11
@@ -87,10 +87,7 @@ def resolve_hermes_bin() -> Optional[str]:
|
||||
|
||||
|
||||
def build_relaunch_argv(
|
||||
extra_args: Sequence[str],
|
||||
*,
|
||||
preserve_inherited: bool = True,
|
||||
original_argv: Optional[Sequence[str]] = None,
|
||||
extra_args: Sequence[str], *, preserve_inherited: bool = True, original_argv: Optional[Sequence[str]] = None
|
||||
) -> list[str]:
|
||||
"""Construct an argv list for replacing the current process with hermes."""
|
||||
bin_path = resolve_hermes_bin()
|
||||
@@ -105,19 +102,14 @@ def build_relaunch_argv(
|
||||
|
||||
|
||||
def relaunch(
|
||||
extra_args: Sequence[str],
|
||||
*,
|
||||
preserve_inherited: bool = True,
|
||||
original_argv: Optional[Sequence[str]] = None,
|
||||
extra_args: Sequence[str], *, preserve_inherited: bool = True, original_argv: Optional[Sequence[str]] = None
|
||||
) -> None:
|
||||
"""Replace the current process with a fresh hermes invocation.
|
||||
|
||||
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
|
||||
)
|
||||
new_argv = build_relaunch_argv(extra_args, preserve_inherited=preserve_inherited, original_argv=original_argv)
|
||||
if sys.platform == "win32":
|
||||
import subprocess
|
||||
try:
|
||||
|
||||
@@ -33,13 +33,10 @@ def legacy_relay_plugin_keys(values: Any) -> tuple[str, ...]:
|
||||
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(
|
||||
env: Mapping[str, Any] | None,
|
||||
) -> tuple[str, ...]:
|
||||
def configured_legacy_relay_env_vars(env: Mapping[str, Any] | None) -> tuple[str, ...]:
|
||||
"""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()
|
||||
name for name in LEGACY_RELAY_EXPORT_ENV_VARS if env.get(name) is not None and str(env[name]).strip()
|
||||
))
|
||||
|
||||
@@ -19,9 +19,7 @@ DEFAULT_NOFILE_SOFT_LIMIT = int(DEFAULT_CONFIG["runtime"]["nofile_soft_limit"])
|
||||
_MISSING = object()
|
||||
|
||||
|
||||
def configured_nofile_soft_limit(
|
||||
config: Mapping[str, Any] | None = None,
|
||||
) -> int | None:
|
||||
def configured_nofile_soft_limit(config: Mapping[str, Any] | None = None) -> int | None:
|
||||
"""``runtime.nofile_soft_limit`` from a loaded config, or ``None`` when disabled/unresolvable.
|
||||
|
||||
Missing key → default. Explicit ``0``/``false``/``null`` disable; other non-int or negative
|
||||
@@ -55,9 +53,7 @@ def configured_nofile_soft_limit(
|
||||
return raw_value
|
||||
|
||||
|
||||
def apply_nofile_soft_limit(
|
||||
config: Mapping[str, Any] | None = None,
|
||||
) -> bool:
|
||||
def apply_nofile_soft_limit(config: Mapping[str, Any] | None = None) -> bool:
|
||||
"""Best-effort raise of this process's ``RLIMIT_NOFILE`` soft limit; ``True`` iff changed.
|
||||
|
||||
Target = ``runtime.nofile_soft_limit`` (default :data:`DEFAULT_NOFILE_SOFT_LIMIT`), clamped
|
||||
@@ -90,6 +86,4 @@ def apply_nofile_soft_limit(
|
||||
return False
|
||||
|
||||
|
||||
__all__ = [
|
||||
"DEFAULT_NOFILE_SOFT_LIMIT", "apply_nofile_soft_limit", "configured_nofile_soft_limit"
|
||||
]
|
||||
__all__ = ["DEFAULT_NOFILE_SOFT_LIMIT", "apply_nofile_soft_limit", "configured_nofile_soft_limit"]
|
||||
|
||||
@@ -46,11 +46,7 @@ class ServiceManager(Protocol):
|
||||
|
||||
def supports_runtime_registration(self) -> bool: ...
|
||||
def register_profile_gateway(
|
||||
self,
|
||||
profile: str,
|
||||
*,
|
||||
extra_env: dict[str, str] | None = None,
|
||||
start_now: bool = True,
|
||||
self, profile: str, *, extra_env: dict[str, str] | None = None, start_now: bool = True
|
||||
) -> None: ...
|
||||
def unregister_profile_gateway(self, profile: str) -> None: ...
|
||||
def list_profile_gateways(self) -> list[str]: ...
|
||||
@@ -139,16 +135,11 @@ class _HostServiceManager:
|
||||
|
||||
def _unsupported(self, verb: str) -> NotImplementedError:
|
||||
return NotImplementedError(
|
||||
f"{type(self).__name__} does not support runtime profile "
|
||||
f"gateway {verb} (container-only feature)"
|
||||
f"{type(self).__name__} does not support runtime profile gateway {verb} (container-only feature)"
|
||||
)
|
||||
|
||||
def register_profile_gateway(
|
||||
self,
|
||||
profile: str,
|
||||
*,
|
||||
extra_env: dict[str, str] | None = None,
|
||||
start_now: bool = True,
|
||||
self, profile: str, *, extra_env: dict[str, str] | None = None, start_now: bool = True
|
||||
) -> None:
|
||||
raise self._unsupported("registration")
|
||||
|
||||
@@ -202,10 +193,7 @@ class WindowsServiceManager(_HostServiceManager):
|
||||
elevated_handoff: bool = False,
|
||||
) -> None:
|
||||
self._backend_module().install(
|
||||
force=force,
|
||||
start_now=start_now,
|
||||
start_on_login=start_on_login,
|
||||
elevated_handoff=elevated_handoff,
|
||||
force=force, start_now=start_now, start_on_login=start_on_login, elevated_handoff=elevated_handoff
|
||||
)
|
||||
|
||||
def is_running(self, name: str) -> bool:
|
||||
@@ -375,9 +363,7 @@ class S6CommandError(S6Error):
|
||||
"""An s6 command failed for a reason other than a missing slot (EACCES on the control FIFO,
|
||||
unexpected non-zero exit); carries the command's stderr."""
|
||||
|
||||
def __init__(
|
||||
self, *, service: str, action: str, returncode: int, stderr: str,
|
||||
) -> None:
|
||||
def __init__(self, *, service: str, action: str, returncode: int, stderr: str) -> None:
|
||||
self.action = action
|
||||
self.returncode = returncode
|
||||
self.stderr = stderr
|
||||
@@ -401,10 +387,7 @@ class S6ServiceManager:
|
||||
return self.scandir / f"{S6_SERVICE_PREFIX}{profile}"
|
||||
|
||||
@staticmethod
|
||||
def _render_run_script(
|
||||
profile: str,
|
||||
extra_env: dict[str, str],
|
||||
) -> str:
|
||||
def _render_run_script(profile: str, extra_env: dict[str, str]) -> str:
|
||||
"""Run script for a profile-gateway s6 service.
|
||||
|
||||
Sources HERMES_HOME via with-contenv (run time, not baked in), resets ``HOME`` before the
|
||||
@@ -511,10 +494,7 @@ class S6ServiceManager:
|
||||
_s6_run("s6-svc", action_flag, str(service_dir), check=True)
|
||||
except subprocess.CalledProcessError as exc:
|
||||
raise S6CommandError(
|
||||
service=name,
|
||||
action=action_label,
|
||||
returncode=exc.returncode,
|
||||
stderr=exc.stderr or "",
|
||||
service=name, action=action_label, returncode=exc.returncode, stderr=exc.stderr or ""
|
||||
) from exc
|
||||
|
||||
def start(self, name: str) -> None:
|
||||
@@ -562,11 +542,7 @@ class S6ServiceManager:
|
||||
return True
|
||||
|
||||
def register_profile_gateway(
|
||||
self,
|
||||
profile: str,
|
||||
*,
|
||||
extra_env: dict[str, str] | None = None,
|
||||
start_now: bool = True,
|
||||
self, profile: str, *, extra_env: dict[str, str] | None = None, start_now: bool = True
|
||||
) -> None:
|
||||
"""Create the s6 service directory and ``s6-svscanctl -a`` so it is picked up immediately.
|
||||
|
||||
|
||||
Reference in New Issue
Block a user