From 79c704f34f65ec1e71fcc9fa31fe2484f0989b5e Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 18:48:19 -0700 Subject: [PATCH] =?UTF-8?q?refactor(agent):=20monitoring/iron=5Fproxy=20th?= =?UTF-8?q?ird=20pass=20=E2=80=94=20gateway.status=20dispatch=20helper,=20?= =?UTF-8?q?SDK=20symbol=20table=20by=20module,=20compact=20comments=20and?= =?UTF-8?q?=20inner=20blanks?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- agent/monitoring/cron_health.py | 2 +- agent/monitoring/events.py | 6 - agent/monitoring/gateway_health.py | 97 ++++------ agent/monitoring/gateway_health_export.py | 27 +-- agent/monitoring/otlp_exporter.py | 51 +++-- agent/monitoring/policy.py | 2 - agent/outbound_webhooks.py | 12 -- agent/proxy_sources/iron_proxy.py | 223 ++++++++-------------- 8 files changed, 150 insertions(+), 270 deletions(-) diff --git a/agent/monitoring/cron_health.py b/agent/monitoring/cron_health.py index 40a2344455..de4e9ea73d 100644 --- a/agent/monitoring/cron_health.py +++ b/agent/monitoring/cron_health.py @@ -100,7 +100,6 @@ def emit_execution_state(record: Optional[dict[str, Any]], *, delivery_outcome: return try: from agent.monitoring import emitter - event = project_execution_event(record, delivery_outcome=delivery_outcome) target = emitter.get_emitter() target.emit(event) @@ -136,6 +135,7 @@ def _freshness_metric(name: str, reader: Callable[[], Optional[float]]) -> Calla value = reader() if value is not None: metrics.append(GatewayMetric(name, max(0.0, float(value)), {})) + return build diff --git a/agent/monitoring/events.py b/agent/monitoring/events.py index a9bb57dcb5..877a537fd3 100644 --- a/agent/monitoring/events.py +++ b/agent/monitoring/events.py @@ -24,9 +24,7 @@ class _MonitoringEvent: @dataclass(slots=True) class GatewayHealthEvent(_MonitoringEvent): """Content-free gateway health snapshot or lifecycle event.""" - EVENT: ClassVar[str] = "gateway_health" - name: str gateway_state: Optional[str] = None old_state: Optional[str] = None @@ -49,9 +47,7 @@ class GatewayHealthEvent(_MonitoringEvent): @dataclass(slots=True) class GatewayDiagnosticEvent(_MonitoringEvent): """Redacted gateway diagnostic event for operator-owned observability.""" - EVENT: ClassVar[str] = "gateway_diagnostic" - name: str subsystem: str error_class: str = "unknown" @@ -69,9 +65,7 @@ class GatewayDiagnosticEvent(_MonitoringEvent): @dataclass(slots=True) class CronExecutionEvent(_MonitoringEvent): """Content-free durable cron execution lifecycle projection.""" - EVENT: ClassVar[str] = "cron_execution" - status: str job_key: str source: str = "unknown" diff --git a/agent/monitoring/gateway_health.py b/agent/monitoring/gateway_health.py index ea7ca161ac..c305110896 100644 --- a/agent/monitoring/gateway_health.py +++ b/agent/monitoring/gateway_health.py @@ -126,31 +126,42 @@ def platform_for_subsystem(subsystem: str) -> Optional[str]: return (subsystem.split(".", 1)[1] or None) if subsystem.startswith("platform.") else None -def _parse_active_agents(raw: Any) -> int: +def _gateway_status(name: str, fallback: Callable[..., Any], /, **kwargs: Any) -> Any: + """Prefer ``gateway.status.`` (the runtime-status contract); fall back to the local approximation.""" try: - from gateway.status import parse_active_agents - return parse_active_agents(raw) + import gateway.status as status + return getattr(status, name)(**kwargs) except Exception: - try: - return max(0, int(raw)) - except (TypeError, ValueError): - return 0 + return fallback(**kwargs) + + +def _int_or_zero(raw: Any) -> int: + try: + return max(0, int(raw)) + except (TypeError, ValueError): + return 0 + + +def _parse_active_agents(raw: Any) -> int: + return _gateway_status("parse_active_agents", _int_or_zero, raw=raw) def _derive_busy(gateway_running: bool, gateway_state: Any, active_agents: Any) -> bool: - try: - from gateway.status import derive_gateway_busy - return derive_gateway_busy(gateway_running=gateway_running, gateway_state=gateway_state, active_agents=active_agents) - except Exception: - return bool(gateway_running and gateway_state == "running" and _parse_active_agents(active_agents) > 0) + return _gateway_status( + "derive_gateway_busy", + lambda gateway_running, gateway_state, active_agents: bool( + gateway_running and gateway_state == "running" and _parse_active_agents(active_agents) > 0 + ), + gateway_running=gateway_running, gateway_state=gateway_state, active_agents=active_agents, + ) def _derive_drainable(gateway_running: bool, gateway_state: Any) -> bool: - try: - from gateway.status import derive_gateway_drainable - return derive_gateway_drainable(gateway_running=gateway_running, gateway_state=gateway_state) - except Exception: - return bool(gateway_running and gateway_state == "running") + return _gateway_status( + "derive_gateway_drainable", + lambda gateway_running, gateway_state: bool(gateway_running and gateway_state == "running"), + gateway_running=gateway_running, gateway_state=gateway_state, + ) def _base_attrs(*, install_id: str, version: str, supervision_mode: str) -> Dict[str, str]: @@ -181,12 +192,7 @@ def _platform_error_code(pdata: dict[str, Any]) -> str: def build_gateway_health_snapshot( - runtime: Optional[dict[str, Any]], - *, - gateway_running: bool, - profile: str, - install_id: str, - version: str, + runtime: Optional[dict[str, Any]], *, gateway_running: bool, profile: str, install_id: str, version: str, supervision_mode: str = "unknown", ) -> GatewayHealthSnapshot: """Convert gateway_state.json-compatible runtime state into P0 signals.""" @@ -197,7 +203,6 @@ def build_gateway_health_snapshot( drainable = _derive_drainable(gateway_running, gateway_state) platforms = _platforms_of(runtime) base = _base_attrs(install_id=install_id, version=version, supervision_mode=supervision_mode) - metrics: list[GatewayMetric] = [ _metric("hermes.gateway.up", 1 if gateway_running else 0, base), _metric("hermes.gateway.active_agents", active_agents, base), @@ -206,7 +211,6 @@ def build_gateway_health_snapshot( _metric("hermes.gateway.restart_requested", 1 if runtime.get("restart_requested") else 0, base), _metric("hermes.gateway.state", 1, base, **{"hermes.gateway.state": gateway_state}), ] - fatal_count = 0 events: list[GatewayHealthEvent | GatewayDiagnosticEvent] = [] for platform, pdata in platforms.items(): @@ -227,20 +231,10 @@ def build_gateway_health_snapshot( error_code=error_code, error_class=error_code, profile=profile, version=version, severity="error" if state == "fatal" else "warning", )) - events.insert(0, GatewayHealthEvent( - name="gateway.health_snapshot", - gateway_state=gateway_state, - active_agents=active_agents, - gateway_busy=busy, - gateway_drainable=drainable, - platform_count=len(platforms), - fatal_platform_count=fatal_count, - profile=profile, - install_id=install_id, - version=version, - supervision_mode=supervision_mode, - pid=_coerce_pid(runtime.get("pid")), + name="gateway.health_snapshot", gateway_state=gateway_state, active_agents=active_agents, gateway_busy=busy, + gateway_drainable=drainable, platform_count=len(platforms), fatal_platform_count=fatal_count, profile=profile, + install_id=install_id, version=version, supervision_mode=supervision_mode, pid=_coerce_pid(runtime.get("pid")), )) return GatewayHealthSnapshot(metrics=metrics, events=events) @@ -273,16 +267,10 @@ def _lifecycle_events( def health(name: str) -> GatewayHealthEvent: return GatewayHealthEvent( - name=name, - gateway_state=new_state, - old_state=old_state, - new_state=new_state, + name=name, gateway_state=new_state, old_state=old_state, new_state=new_state, exit_reason=classify_exit_reason(current.get("exit_reason"), state=new_state, restart_requested=restart_requested), - restart_requested=restart_requested, - active_agents=_parse_active_agents(current.get("active_agents", 0)), - profile=profile, - version=version, - pid=_coerce_pid(current.get("pid")), + restart_requested=restart_requested, active_agents=_parse_active_agents(current.get("active_agents", 0)), + profile=profile, version=version, pid=_coerce_pid(current.get("pid")), ) out: list[GatewayHealthEvent | GatewayDiagnosticEvent] = [health("gateway.lifecycle")] @@ -351,21 +339,14 @@ class GatewayDiagnosticLogHandler(logging.Handler): def emit(self, record: logging.LogRecord) -> None: try: - if record.levelno < logging.WARNING: - return - if not (record.name == "gateway" or record.name.startswith("gateway.")): + if record.levelno < logging.WARNING or not (record.name == "gateway" or record.name.startswith("gateway.")): return subsystem = subsystem_for_logger(record.name) error_class = classify_gateway_error(record.getMessage()) emitter.get_emitter().emit(GatewayDiagnosticEvent( - name=f"gateway.log.{record.levelname.lower()}", - subsystem=subsystem, - source_logger=source_logger_for_export(record.name), - platform=platform_for_subsystem(subsystem), - error_class=error_class, - error_code=error_class, - profile=self.profile, - version=self.version, + name=f"gateway.log.{record.levelname.lower()}", subsystem=subsystem, + source_logger=source_logger_for_export(record.name), platform=platform_for_subsystem(subsystem), + error_class=error_class, error_code=error_class, profile=self.profile, version=self.version, severity=record.levelname.lower(), )) except Exception: diff --git a/agent/monitoring/gateway_health_export.py b/agent/monitoring/gateway_health_export.py index 81f4bfd95d..643298be61 100644 --- a/agent/monitoring/gateway_health_export.py +++ b/agent/monitoring/gateway_health_export.py @@ -80,7 +80,6 @@ class GatewayHealthExportRuntime: if self.log_handler is not None: with suppress(Exception): logging.getLogger().removeHandler(self.log_handler) - # Producers are stopped; drain queued/in-flight events BEFORE detaching subscribers # so the terminal lifecycle event cannot race exporter shutdown. Bounded, fail-open. subscribers = [item for item in (self.streamer, self.log_streamer) if item is not None] @@ -89,7 +88,6 @@ class GatewayHealthExportRuntime: bus.flush(timeout=1.0) for sub in subscribers: bus.unsubscribe(sub) - # Network flush/close runs under one bounded daemon-thread deadline so it can # never delay gateway teardown indefinitely. closeables = subscribers + ([self.metric_provider] if self.metric_provider is not None else []) @@ -103,7 +101,6 @@ class GatewayHealthExportRuntime: worker = threading.Thread(target=_close, name="hermes-gateway-health-export-shutdown", daemon=True) worker.start() worker.join(timeout=2.0) - self.streamer = self.log_streamer = self.metric_provider = self.thread = self.stop_event = None @@ -161,7 +158,6 @@ def _read_gateway_snapshot(config: Dict[str, Any]): def _read_cron_snapshot(): from agent.monitoring.cron_health import build_cron_health_snapshot - return build_cron_health_snapshot() @@ -242,11 +238,11 @@ def _start_metric_provider(config: Dict[str, Any], sdk: Dict[str, Any]) -> Any: def callback(name: str): def _cb(_options=None): try: - snapshot = _read_runtime_snapshot(config) - return [Observation(m.value, m.attributes) for m in snapshot.metrics if m.name == name] + return [Observation(m.value, m.attributes) for m in _read_runtime_snapshot(config).metrics if m.name == name] except Exception: logger.debug("gateway metric callback failed", exc_info=True) return [] + return _cb for metric_name in _OBSERVABLE_METRIC_NAMES: @@ -286,13 +282,9 @@ class GatewayDiagnosticLogStreamer(EmitterStreamer): source_logger = source_logger_for_export(ev.get("source_logger")) otel_logger = self._provider.get_logger(source_logger) if source_logger is not None else self._logger otel_logger.emit(sdk["LogRecord"]( - timestamp=ev.get("ts_ns"), - trace_id=sdk["INVALID_TRACE_ID"], - span_id=sdk["INVALID_SPAN_ID"], - trace_flags=sdk["TraceFlags"].DEFAULT, - severity_text=str(ev.get("severity") or "warning").upper(), - severity_number=_severity_number(sdk, ev.get("severity")), - body=redact_bounded("gateway diagnostic"), + timestamp=ev.get("ts_ns"), trace_id=sdk["INVALID_TRACE_ID"], span_id=sdk["INVALID_SPAN_ID"], + trace_flags=sdk["TraceFlags"].DEFAULT, severity_text=str(ev.get("severity") or "warning").upper(), + severity_number=_severity_number(sdk, ev.get("severity")), body=redact_bounded("gateway diagnostic"), attributes=_diagnostic_log_attributes(ev), )) self.exported += 1 @@ -334,17 +326,12 @@ def start_gateway_health_export(config: Dict[str, Any]) -> GatewayHealthExportRu diagnostics_on = gh.get("diagnostic_events_enabled", True) runtime = GatewayHealthExportRuntime(enabled=True, reason="enabled") sdk: Optional[Dict[str, Any]] = None - if metrics_on or diagnostics_on: try: sdk = _require_metrics_sdk(prompt=False) except Exception: - logger.warning( - "monitoring.gateway_health_export.enabled but OTLP SDK is unavailable; install 'hermes-agent[otlp]'", - exc_info=True, - ) + logger.warning("monitoring.gateway_health_export.enabled but OTLP SDK is unavailable; install 'hermes-agent[otlp]'", exc_info=True) return GatewayHealthExportRuntime(enabled=False, reason="otlp_unavailable") - if metrics_on and sdk is not None: try: runtime.metric_provider = _start_metric_provider(config, sdk) @@ -352,7 +339,6 @@ def start_gateway_health_export(config: Dict[str, Any]) -> GatewayHealthExportRu logger.warning("gateway health OTLP metrics failed to start", exc_info=True) runtime.shutdown() return GatewayHealthExportRuntime(enabled=False, reason="metrics_start_failed") - if diagnostics_on and sdk is not None: try: runtime.streamer = otlp_exporter.start_streaming(config, event_filter=_gateway_health_event) @@ -365,7 +351,6 @@ def start_gateway_health_export(config: Dict[str, Any]) -> GatewayHealthExportRu logger.debug("gateway diagnostic OTLP export failed to start", exc_info=True) runtime.shutdown() return GatewayHealthExportRuntime(enabled=False, reason="diagnostics_start_failed") - try: runtime.log_handler = _attach_log_handler(config) except Exception: diff --git a/agent/monitoring/otlp_exporter.py b/agent/monitoring/otlp_exporter.py index 0c2d8e102c..2f845569ff 100644 --- a/agent/monitoring/otlp_exporter.py +++ b/agent/monitoring/otlp_exporter.py @@ -29,25 +29,23 @@ class OTLPUnavailable(RuntimeError): # ── SDK loading ────────────────────────────────────────────────────────────── -_SDK_SYMBOLS: Dict[str, str] = { - "TracerProvider": "opentelemetry.sdk.trace", - "BatchSpanProcessor": "opentelemetry.sdk.trace.export", - "Resource": "opentelemetry.sdk.resources", - "OTLPSpanExporter": "opentelemetry.exporter.otlp.proto.http.trace_exporter", - "SpanKind": "opentelemetry.trace", - "OTLPLogExporter": "opentelemetry.exporter.otlp.proto.http._log_exporter", - "OTLPMetricExporter": "opentelemetry.exporter.otlp.proto.http.metric_exporter", - "Observation": "opentelemetry.metrics", - "INVALID_SPAN_ID": "opentelemetry.trace", - "INVALID_TRACE_ID": "opentelemetry.trace", - "TraceFlags": "opentelemetry.trace", - "LogRecord": "opentelemetry._logs", - "SeverityNumber": "opentelemetry._logs.severity", - "LoggerProvider": "opentelemetry.sdk._logs", - "BatchLogRecordProcessor": "opentelemetry.sdk._logs.export", - "MeterProvider": "opentelemetry.sdk.metrics", - "PeriodicExportingMetricReader": "opentelemetry.sdk.metrics.export", +_SDK_MODULES: Dict[str, tuple[str, ...]] = { + "opentelemetry.sdk.trace": ("TracerProvider",), + "opentelemetry.sdk.trace.export": ("BatchSpanProcessor",), + "opentelemetry.sdk.resources": ("Resource",), + "opentelemetry.exporter.otlp.proto.http.trace_exporter": ("OTLPSpanExporter",), + "opentelemetry.exporter.otlp.proto.http._log_exporter": ("OTLPLogExporter",), + "opentelemetry.exporter.otlp.proto.http.metric_exporter": ("OTLPMetricExporter",), + "opentelemetry.trace": ("SpanKind", "INVALID_SPAN_ID", "INVALID_TRACE_ID", "TraceFlags"), + "opentelemetry.metrics": ("Observation",), + "opentelemetry._logs": ("LogRecord",), + "opentelemetry._logs.severity": ("SeverityNumber",), + "opentelemetry.sdk._logs": ("LoggerProvider",), + "opentelemetry.sdk._logs.export": ("BatchLogRecordProcessor",), + "opentelemetry.sdk.metrics": ("MeterProvider",), + "opentelemetry.sdk.metrics.export": ("PeriodicExportingMetricReader",), } +_SDK_SYMBOLS: Dict[str, str] = {name: module for module, names in _SDK_MODULES.items() for name in names} _SPAN_SDK = ("TracerProvider", "BatchSpanProcessor", "Resource", "OTLPSpanExporter", "SpanKind") @@ -172,16 +170,12 @@ def _make_provider(config: Dict[str, Any]): # ── event -> span attribute mapping ────────────────────────────────────────── # Per-kind attribute allowlists: everything else (profile, install_id, ...) never egresses. _KEEP_BY_KIND: Dict[str, tuple[str, ...]] = { - "gateway_health": ("name", "gateway_state", "old_state", "new_state", - "exit_reason", "restart_requested", "active_agents", - "gateway_busy", "gateway_drainable", "platform_count", - "fatal_platform_count", "version", - "supervision_mode", "pid"), - "gateway_diagnostic": ("name", "subsystem", "error_class", "error_code", - "platform", "old_state", "new_state", - "version", "severity"), - "cron_execution": ("status", "job_key", "source", "duration_ms", - "delivery_outcome", "error_class"), + "gateway_health": ( + "name", "gateway_state", "old_state", "new_state", "exit_reason", "restart_requested", "active_agents", + "gateway_busy", "gateway_drainable", "platform_count", "fatal_platform_count", "version", "supervision_mode", "pid", + ), + "gateway_diagnostic": ("name", "subsystem", "error_class", "error_code", "platform", "old_state", "new_state", "version", "severity"), + "cron_execution": ("status", "job_key", "source", "duration_ms", "delivery_outcome", "error_class"), } @@ -220,7 +214,6 @@ def export_batch(provider, batch: List[Dict[str, Any]]) -> int: class EmitterStreamer: """Base for emitter subscribers owning an OTel provider + batch processor. Register with ``emitter.subscribe(streamer)``. Fail-isolated by the emitter.""" - _provider: Any _processor: Any exported: int = 0 diff --git a/agent/monitoring/policy.py b/agent/monitoring/policy.py index a50f397c7a..1c1620a04f 100644 --- a/agent/monitoring/policy.py +++ b/agent/monitoring/policy.py @@ -25,11 +25,9 @@ def ensure_install_id(config: Dict[str, Any]) -> str: existing = (mon or {}).get("install_id") if isinstance(mon, dict) else None if isinstance(existing, str) and existing.strip(): return existing - minted = str(uuid.uuid4()) try: from hermes_cli.config import load_config, save_config - fresh = load_config() if isinstance(fresh, dict): slot = fresh.setdefault("monitoring", {}) diff --git a/agent/outbound_webhooks.py b/agent/outbound_webhooks.py index 4604504147..4504838888 100644 --- a/agent/outbound_webhooks.py +++ b/agent/outbound_webhooks.py @@ -59,9 +59,7 @@ _worker: Optional[threading.Thread] = None @dataclass class WebhookTarget(_ToolMatcherMixin): """Parsed and validated representation of one ``hooks.outbound`` entry.""" - _MATCHER_KIND = "outbound webhook" - url: str events: List[str] name: str = "" @@ -86,7 +84,6 @@ def register_from_config(cfg: Optional[Dict[str, Any]]) -> List[WebhookTarget]: if not isinstance(cfg, dict): return [] from utils import env_var_enabled - if env_var_enabled("HERMES_SAFE_MODE"): logger.info("HERMES_SAFE_MODE=1 — outbound webhook registration skipped") return [] @@ -94,7 +91,6 @@ def register_from_config(cfg: Optional[Dict[str, Any]]) -> List[WebhookTarget]: if not targets: return [] from hermes_cli.plugins import get_plugin_manager - manager = get_plugin_manager() home_key = _home_key() registered: List[WebhookTarget] = [] @@ -149,7 +145,6 @@ def re_register_config_hooks() -> None: current home's idempotence keys are cleared so a force-reload in one profile cannot invalidate another profile's still-live registration.""" from hermes_cli.config import load_config - _forget_home_registrations(_registered, _registered_lock) register_from_config(load_config()) @@ -177,7 +172,6 @@ def _parse_single_target(index: int, raw: Any) -> Optional[WebhookTarget]: if not isinstance(raw, dict): warn(" must be a mapping with 'url' and 'events' keys; got %s", type(raw).__name__) return None - url = raw.get("url") if not isinstance(url, str) or not url.strip(): warn(" is missing a non-empty 'url'") @@ -188,7 +182,6 @@ def _parse_single_target(index: int, raw: Any) -> Optional[WebhookTarget]: return None if url.lower().startswith("http://"): warn(".url uses plain http:// — payloads (including tool inputs) travel unencrypted. Prefer https.") - events_raw = raw.get("events") valid_list = ", ".join(sorted(VALID_HOOKS)) if not isinstance(events_raw, list) or not events_raw: @@ -203,7 +196,6 @@ def _parse_single_target(index: int, raw: Any) -> Optional[WebhookTarget]: if not events: warn(" has no valid events — skipped") return None - matcher = raw.get("matcher") if matcher is not None and not isinstance(matcher, str): warn(".matcher must be a string regex; ignoring") @@ -211,14 +203,12 @@ def _parse_single_target(index: int, raw: Any) -> Optional[WebhookTarget]: if matcher is not None and not any(e in _TOOL_SCOPED_EVENTS for e in events): warn(".matcher=%r will be ignored — matcher is only honored for pre_tool_call / post_tool_call.", matcher) matcher = None - timeout_raw = raw.get("timeout", DEFAULT_TIMEOUT_SECONDS) try: timeout = int(timeout_raw) except (TypeError, ValueError): warn(".timeout must be an int (got %r); using default %ds", timeout_raw, DEFAULT_TIMEOUT_SECONDS) timeout = DEFAULT_TIMEOUT_SECONDS - name = raw.get("name") return WebhookTarget( url=url, @@ -276,7 +266,6 @@ def _serialize_payload(event: str, kwargs: Dict[str, Any], delivery_id: str) -> body, so they double as replay protection.""" # Profile resolved at fire time so a multiplexed gateway's receivers can tell which profile emitted. from hermes_cli.profiles import get_active_profile_name - payload = { "hook_event_name": event, "profile": get_active_profile_name(), @@ -383,7 +372,6 @@ def _deliver(delivery: Dict[str, Any]) -> None: last_error = str(exc) or type(exc).__name__ if attempt < MAX_DELIVERY_ATTEMPTS: time.sleep(RETRY_BACKOFF_SECONDS * attempt) - logger.warning( "outbound webhook delivery failed after %d attempt(s) (event=%s target=%s): %s", MAX_DELIVERY_ATTEMPTS, event, label, last_error, diff --git a/agent/proxy_sources/iron_proxy.py b/agent/proxy_sources/iron_proxy.py index 6ab6c18aba..d5cd271165 100644 --- a/agent/proxy_sources/iron_proxy.py +++ b/agent/proxy_sources/iron_proxy.py @@ -31,13 +31,11 @@ from typing import Dict, List, Optional, Tuple logger = logging.getLogger(__name__) -# Pinned upstream version. Never auto-resolve "latest": the YAML schema may change between -# releases and updates must be deliberate. +# Pinned: never auto-resolve "latest" — the YAML schema may change between releases. _IRON_PROXY_VERSION = "0.39.0" _IRON_PROXY_RELEASE_BASE = f"https://github.com/ironsh/iron-proxy/releases/download/v{_IRON_PROXY_VERSION}" _IRON_PROXY_CHECKSUM_NAME = "checksums.txt" -# Detached signature + signing key for optional GPG verification of checksums.txt (SHA-256 -# alone only protects the archive if checksums.txt itself came from an untampered channel). +# Optional GPG verification of checksums.txt (SHA-256 alone trusts the release channel). _IRON_PROXY_CHECKSUM_SIG_NAME = "checksums.txt.asc" _IRON_PROXY_PUBKEY_NAME = "public-key.asc" @@ -45,9 +43,8 @@ _DOWNLOAD_TIMEOUT = 120 # binary is ~16MB _RUN_TIMEOUT = 30 _STARTUP_GRACE_SECONDS = 5 -# Management API (v0.39): authenticated loopback listener whose POST /v1/reload hot-swaps the -# ruleset. Bearer key minted at setup, stored 0600 at /management.token, injected into -# the daemon env under this name; v0.39 refuses to start if the named var is empty. +# Management API (v0.39): loopback POST /v1/reload hot-swaps the ruleset. Bearer key minted at +# setup (0600 at /management.token), injected under this env name; empty => daemon refuses to start. _MGMT_API_KEY_ENV = "HERMES_IRON_PROXY_MGMT_KEY" _MGMT_PORT_OFFSET = 2 # tunnel_port is CONNECT/MITM, +1 is plain-HTTP forward, +2 is management _MGMT_RELOAD_TIMEOUT = 15 @@ -74,12 +71,11 @@ _BEARER_PROVIDERS: Dict[str, Tuple[str, ...]] = { "NOUS_API_KEY": ("inference.nousresearch.com",), } -# Providers authenticating with a non-Authorization header (v0.39 ``match_headers`` is -# case-insensitive). ``aliases`` are interchangeable env names for the SAME credential and -# MUST collapse into one mapping: every rule is ``require: true`` and two require-rules on one -# host would reject each other's requests. The sandbox gets the token under every name. -# Authorization is also matched for Anthropic/Azure so an SDK sending the token as Bearer still -# swaps; Gemini's ``?key=`` SDK style is covered by match_query. +# Non-Authorization-header providers (v0.39 ``match_headers`` is case-insensitive). ``aliases`` +# name the SAME credential and MUST collapse into one mapping: every rule is ``require: true`` and +# two require-rules on one host would reject each other's requests; the sandbox gets the token +# under every name. Authorization is also matched for Anthropic/Azure (SDKs may send Bearer); +# Gemini's ``?key=`` style is covered by match_query. _HEADER_AUTH_PROVIDERS: Dict[str, Dict[str, Tuple[str, ...]]] = { "ANTHROPIC_API_KEY": {"hosts": ("api.anthropic.com",), "match_headers": ("x-api-key", "Authorization"), "aliases": ()}, "AZURE_OPENAI_API_KEY": { @@ -94,14 +90,13 @@ _HEADER_AUTH_PROVIDERS: Dict[str, Dict[str, Tuple[str, ...]]] = { }, } -# Recognized creds that cannot be swapped by static header replacement (SigV4 signing, -# SDK-minted OAuth). Surfaced as a warning only — they're generic cloud creds, never blocking. +# Creds that static header replacement can't swap (SigV4, SDK-minted OAuth): warning only. _NON_BEARER_PROVIDERS: Tuple[str, ...] = ( "AWS_ACCESS_KEY_ID", "AWS_SECRET_ACCESS_KEY", "GOOGLE_APPLICATION_CREDENTIALS", ) -# Default SSRF deny list for outbound traffic (docs promise: cloud metadata IPs refused -# regardless of allowlist). Callers may pass [] to disable (hermetic tests only). +# Default SSRF deny list (docs promise: cloud metadata IPs refused regardless of allowlist); +# callers pass [] to disable (hermetic tests only). _DEFAULT_UPSTREAM_DENY_CIDRS: Tuple[str, ...] = ( "127.0.0.0/8", "::1/128", # loopback v4 / v6 "169.254.0.0/16", "fe80::/10", # link-local incl. AWS/GCP/Azure IMDS @@ -111,8 +106,7 @@ _DEFAULT_UPSTREAM_DENY_CIDRS: Tuple[str, ...] = ( "198.18.0.0/15", # RFC2544 benchmark range ) -# Minimal env the iron-proxy subprocess needs; everything else is stripped so -# /proc//environ never exposes the operator's shell secrets. +# Minimal daemon env; everything else is stripped so /proc//environ never exposes operator secrets. _PROXY_SUBPROCESS_ENV_ALLOWLIST: Tuple[str, ...] = ( "PATH", "HOME", "TMPDIR", "TZ", "LANG", "LC_ALL", "LC_CTYPE", "NO_COLOR", "SSL_CERT_DIR", "SSL_CERT_FILE", @@ -129,11 +123,11 @@ _KILL_SIGNAL = getattr(signal, "SIGKILL", signal.SIGTERM) # O_NOFOLLOW is POSIX-only; 0 is a no-op flag elsewhere. _O_NOFOLLOW = getattr(os, "O_NOFOLLOW", 0) -# ``iron-proxy --version`` output keyed by binary path (get_status runs per container create). +# ``--version`` output keyed by binary path (get_status runs per container create). _VERSION_CACHE: Dict[str, str] = {} -# Nonce planted in the daemon env at start so ``_pid_alive`` can prove a PID is still *our* -# binary across PID recycling (a fresh process can't inherit our arbitrary env value). +# Nonce planted in the daemon env so ``_pid_alive`` can prove a PID is still *our* binary across +# PID recycling (a fresh process can't inherit our arbitrary env value). _HERMES_IRON_PROXY_NONCE_ENV = "HERMES_IRON_PROXY_NONCE" _proxy_nonce: Optional[str] = None @@ -162,9 +156,8 @@ class ProxyStatus: @dataclass class TokenMapping: - """Map a sandbox-visible proxy token to a real upstream credential lookup. ``real_env_name`` - is read from iron-proxy's OWN env at egress time; ``alias_env_names`` are extra names the - SANDBOX receives the same token under (not emitted in the proxy config).""" + """Sandbox-visible proxy token -> upstream credential lookup. ``real_env_name`` is read from + iron-proxy's OWN env at egress; ``alias_env_names`` are extra SANDBOX names for the same token.""" proxy_token: str real_env_name: str upstream_hosts: Tuple[str, ...] @@ -184,8 +177,8 @@ def _proxy_state_dir_ro() -> Path: def _proxy_state_dir() -> Path: - """Proxy state dir, created 0o700 (holds the CA key, pidfile, logs); chmod is unconditional - so a pre-existing dir with a slack umask gets tightened.""" + """Proxy state dir (CA key, pidfile, logs), created 0o700; chmod is unconditional so a + pre-existing slack-umask dir gets tightened.""" d = _proxy_state_dir_ro() d.mkdir(parents=True, exist_ok=True) with suppress(OSError): # Windows no-op / shared fs we don't own; files still get explicit perms @@ -236,7 +229,6 @@ def install_iron_proxy(*, force: bool = False) -> Path: target = bin_dir / _platform_binary_name() if target.exists() and not force: return target - asset_name = _platform_asset_name() with tempfile.TemporaryDirectory(prefix="hermes-iron-proxy-") as tmpdir: tmp = Path(tmpdir) @@ -251,24 +243,20 @@ def install_iron_proxy(*, force: bool = False) -> Path: actual = _sha256_file(archive_path) if expected.lower() != actual.lower(): raise RuntimeError(f"Checksum mismatch for {asset_name}: expected {expected}, got {actual}") - with tarfile.open(archive_path, "r:gz") as tf: member = _pick_tar_member(tf, _platform_binary_name()) - # PEP 706 data filter rejects escaping links; Python < 3.12 lacks the kwarg and - # relies on _pick_tar_member's path sanitization. + # PEP 706 data filter rejects escaping links; < 3.12 relies on _pick_tar_member's sanitization. try: tf.extract(member, tmp, filter="data") # noqa: S202 except TypeError: tf.extract(member, tmp) # noqa: S202 extracted = tmp / member.name - # Stage then atomically rename so the binary is never visible half-written. fd, staged = tempfile.mkstemp(dir=str(bin_dir), prefix=".iron-proxy_") os.close(fd) shutil.copy2(extracted, staged) os.chmod(staged, 0o755) os.replace(staged, target) - # A freshly-installed binary must re-probe --version on the next get_status(). _VERSION_CACHE.pop(str(target), None) logger.info("Installed iron-proxy %s at %s", _IRON_PROXY_VERSION, target) @@ -287,14 +275,13 @@ def _release_asset(name: str, dest: Path) -> None: def _verify_checksums_signature(tmp: Path, checksum_path: Path) -> bool: - """Best-effort GPG verification of checksums.txt in an ephemeral keyring. Returns False (with - a warning) when gpg or the signature assets are unavailable — SHA-256 stays enforced and gpg - is never a hard dependency. Raises ONLY on a present-but-bad signature (tamper signal).""" + """Best-effort GPG check of checksums.txt in an ephemeral keyring. False (with a warning) when + gpg or the signature assets are unavailable — SHA-256 stays enforced, gpg is never a hard + dependency. Raises ONLY on a present-but-bad signature (tamper signal).""" gpg = shutil.which("gpg") if not gpg: logger.warning("gpg not found on PATH — skipping iron-proxy release-signature verification (SHA-256 checksum check still enforced).") return False - sig_path, pubkey_path = tmp / _IRON_PROXY_CHECKSUM_SIG_NAME, tmp / _IRON_PROXY_PUBKEY_NAME try: _release_asset(_IRON_PROXY_CHECKSUM_SIG_NAME, sig_path) @@ -302,7 +289,6 @@ def _verify_checksums_signature(tmp: Path, checksum_path: Path) -> bool: except RuntimeError as exc: logger.warning("iron-proxy release signature assets unavailable (%s) — skipping GPG verification (SHA-256 checksum check still enforced).", exc) return False - gnupg_home = tmp / "gnupg" gnupg_home.mkdir(mode=0o700, exist_ok=True) gpg_base = [gpg, "--homedir", str(gnupg_home), "--batch", "--no-tty"] @@ -341,8 +327,7 @@ def _pick_tar_member(tf: tarfile.TarFile, binary_name: str) -> tarfile.TarInfo: """Find the binary in the archive (flat or one dir deep); reject abs paths and ``..``.""" candidates = [ m for m in tf.getmembers() - if m.isfile() and not m.name.startswith("/") and ".." not in Path(m.name).parts - and Path(m.name).name == binary_name + if m.isfile() and not m.name.startswith("/") and ".." not in Path(m.name).parts and Path(m.name).name == binary_name ] if not candidates: raise RuntimeError(f"Could not find {binary_name} inside downloaded archive (members: {[m.name for m in tf.getmembers()[:5]]}...)") @@ -378,8 +363,8 @@ def iron_proxy_version(binary: Path) -> str: def _write_private_file(path: Path, data: bytes) -> None: - """Create/truncate ``path`` 0o600 from the first byte (no chmod-after TOCTOU); O_NOFOLLOW - defends against a planted symlink; fchmod tightens a pre-existing file.""" + """Create/truncate ``path`` 0o600 from the first byte (no chmod-after TOCTOU), O_NOFOLLOW + against a planted symlink, fchmod to tighten a pre-existing file.""" fd = os.open(str(path), os.O_WRONLY | os.O_CREAT | os.O_TRUNC | _O_NOFOLLOW, 0o600) try: with suppress(OSError, AttributeError): @@ -406,17 +391,14 @@ def ensure_ca_cert(*, force: bool = False) -> Tuple[Path, Path]: return ca_crt, ca_key if shutil.which("openssl") is None: raise RuntimeError("openssl not found on PATH. Install OpenSSL (apt: `openssl`, brew: `openssl`) to generate the iron-proxy CA cert.") - with tempfile.TemporaryDirectory(prefix="hermes-proxy-ca-") as tmpdir: tmp_key, tmp_crt = Path(tmpdir) / "ca.key", Path(tmpdir) / "ca.crt" _run(["openssl", "genrsa", "-out", str(tmp_key), "4096"], timeout=60, check=True) _run([ - "openssl", "req", "-x509", "-new", "-nodes", "-key", str(tmp_key), - "-sha256", "-days", "3650", "-subj", "/CN=hermes iron-proxy CA", - "-addext", "basicConstraints=critical,CA:TRUE", "-addext", "keyUsage=critical,keyCertSign", - "-out", str(tmp_crt), + "openssl", "req", "-x509", "-new", "-nodes", "-key", str(tmp_key), "-sha256", "-days", "3650", + "-subj", "/CN=hermes iron-proxy CA", "-addext", "basicConstraints=critical,CA:TRUE", + "-addext", "keyUsage=critical,keyCertSign", "-out", str(tmp_crt), ], timeout=60, check=True) - # Key: stage 0o600 against a fresh inode, then atomically rename into place. key_staged = ca_key.with_suffix(ca_key.suffix + ".staged") key_staged.unlink(missing_ok=True) @@ -425,7 +407,6 @@ def ensure_ca_cert(*, force: bool = False) -> Tuple[Path, Path]: # Cert is public — 0o644 matches typical PEM layout. ca_crt.write_bytes(tmp_crt.read_bytes()) os.chmod(ca_crt, 0o644) - logger.info("Generated iron-proxy CA at %s", ca_crt) return ca_crt, ca_key @@ -505,8 +486,8 @@ _RELOAD_HTTP_ERRORS = { def reload_proxy() -> bool: - """Hot-reload the ruleset via ``POST /v1/reload`` (validation failures leave the running - config untouched). Raises RuntimeError with an actionable message on any failure.""" + """``POST /v1/reload`` (validation failures leave the running config untouched); RuntimeError + with an actionable message on any failure.""" pid = _read_pid() if not pid or not _pid_alive(pid): raise RuntimeError("iron-proxy is not running — nothing to reload. Run `hermes egress start`.") @@ -519,11 +500,8 @@ def reload_proxy() -> bool: token = _read_text_or_none(_proxy_state_dir_ro() / "management.token") if not token: raise RuntimeError("management.token is missing — re-run `hermes egress setup`, then `hermes egress restart`.") - host, port = mgmt - req = urllib.request.Request( - f"http://{host}:{port}/v1/reload", method="POST", headers={"Authorization": f"Bearer {token}"}, data=b"", - ) + req = urllib.request.Request(f"http://{host}:{port}/v1/reload", method="POST", headers={"Authorization": f"Bearer {token}"}, data=b"") try: with urllib.request.urlopen(req, timeout=_MGMT_RELOAD_TIMEOUT) as resp: if resp.status == 200: @@ -544,10 +522,9 @@ def reload_proxy() -> bool: def _default_http_listen(tunnel_port: int) -> List[str]: - """Single bind (v0.39 allows one per daemon): docker bridge gateway on Linux — that's what - ``host.docker.internal`` resolves to there, loopback is unreachable from containers — and - loopback on macOS/Windows Docker Desktop (VPNkit). NEVER 0.0.0.0: a LAN peer with a leaked - sandbox token could spend the operator's API quota.""" + """Single bind (v0.39 allows one): docker bridge on Linux (what ``host.docker.internal`` + resolves to; loopback is unreachable from containers), loopback on Docker Desktop (VPNkit). + NEVER 0.0.0.0: a LAN peer with a leaked sandbox token could spend the operator's API quota.""" if platform.system() == "Linux": bridge_ip = _detect_docker_bridge_ip() if bridge_ip and bridge_ip != "127.0.0.1": @@ -561,8 +538,7 @@ def _default_http_listen(tunnel_port: int) -> List[str]: def _detect_docker_bridge_ip() -> Optional[str]: """docker0 IPv4 via ``ip -4 addr show docker0``, or None. SECURITY: validated through - :mod:`ipaddress` so a hostile ``ip`` shim on PATH can't inject ``0.0.0.0`` (or - loopback/multicast/link-local/public) and reopen INADDR_ANY binding.""" + :mod:`ipaddress` so a hostile ``ip`` shim can't inject 0.0.0.0/loopback/multicast/link-local/public.""" try: res = _run(["ip", "-4", "-o", "addr", "show", "docker0"], timeout=2, text=True) except (OSError, subprocess.TimeoutExpired): @@ -592,20 +568,18 @@ def build_proxy_config( audit_log: Optional[Path] = None, allowed_hosts: Optional[List[str]] = None, upstream_deny_cidrs: Optional[List[str]] = None, http_listen: Optional[List[str]] = None, ) -> Dict: - """Build the iron-proxy YAML config dict (schema as of v0.39.0). Real secrets are read from - iron-proxy's OWN env (``source: {type: env}``); the sandbox never sees them. - ``upstream_deny_cidrs=None`` means the default SSRF deny list; ``[]`` opts out. ``audit_log`` - is accepted for forward-compat but unused: v0.39's ``config.Log`` rejects ``audit_path``.""" + """iron-proxy YAML config dict (v0.39.0 schema). Real secrets come from iron-proxy's OWN env + (``source: {type: env}``); the sandbox never sees them. ``upstream_deny_cidrs=None`` = default + SSRF deny list, ``[]`` opts out. ``audit_log`` is forward-compat only (v0.39 rejects ``audit_path``).""" hosts: List[str] = list(allowed_hosts or _DEFAULT_ALLOWED_HOSTS) for m in mappings: for h in m.upstream_hosts: if h not in hosts: hosts.append(h) deny_cidrs = list(_DEFAULT_UPSTREAM_DENY_CIDRS if upstream_deny_cidrs is None else upstream_deny_cidrs) - - # Per-provider headers (case-insensitive in v0.39); query scan covers ``?key=`` SDKs; - # body inspection deliberately off. ``require`` fails closed: an allowlisted-host request - # WITHOUT the proxy token is rejected, so a real key sent directly can't cross the boundary. + # Query scan covers ``?key=`` SDKs; body inspection deliberately off. ``require`` fails + # closed: an allowlisted-host request WITHOUT the proxy token is rejected, so a real key sent + # directly can't cross the boundary. secrets_rules = [ { "source": {"type": "env", "var": m.real_env_name}, @@ -617,29 +591,23 @@ def build_proxy_config( } for m in mappings ] - - # v0.39 takes ONE string per listener field (no plural form). tunnel_listen is the - # CONNECT+MITM listener sandboxes must reach via HTTPS_PROXY (a CONNECT to http_listen is - # forwarded upstream and 400s); http_listen is plain-HTTP forward on tunnel_port+1. + # ONE string per listener field. tunnel_listen is the CONNECT+MITM listener sandboxes reach via + # HTTPS_PROXY (a CONNECT to http_listen is forwarded upstream and 400s); http_listen is plain-HTTP + # forward on tunnel_port+1. listens = list(http_listen) if http_listen else _default_http_listen(tunnel_port) primary_listen = listens[0] if listens else f"127.0.0.1:{tunnel_port}" bind_host = primary_listen.rsplit(":", 1)[0] or "127.0.0.1" - return { # Required by the parser; tunnel-only mode never binds an exposed DNS port. "dns": {"listen": "127.0.0.1:0", "proxy_ip": "127.0.0.1"}, "proxy": { # Both bind the docker bridge on Linux / loopback on Docker Desktop — NEVER 0.0.0.0. - "tunnel_listen": primary_listen, - "http_listen": f"{bind_host}:{tunnel_port + 1}", + "tunnel_listen": primary_listen, "http_listen": f"{bind_host}:{tunnel_port + 1}", "https_listen": "127.0.0.1:0", # direct-TLS listener is not exposed - "max_request_body_bytes": 16 * 1024 * 1024, - "max_response_body_bytes": 0, - "upstream_response_header_timeout": "120s", - "upstream_deny_cidrs": deny_cidrs, + "max_request_body_bytes": 16 * 1024 * 1024, "max_response_body_bytes": 0, + "upstream_response_header_timeout": "120s", "upstream_deny_cidrs": deny_cidrs, }, - # v0.39 defaults its metrics server to :9090 — the same as our default tunnel_port — - # so pin it to an ephemeral loopback port (metrics effectively undiscoverable). + # v0.39 defaults metrics to :9090 (our default tunnel_port) — pin to an ephemeral loopback port. "metrics": {"listen": "127.0.0.1:0"}, # Loopback only: sandboxes must never reach the management surface. "management": {"listen": f"127.0.0.1:{tunnel_port + _MGMT_PORT_OFFSET}", "api_key_env": _MGMT_API_KEY_ENV}, @@ -653,9 +621,8 @@ def build_proxy_config( def _open_private_append(path: Path, *, strict_chmod: bool) -> int: - """Open ``path`` O_APPEND|O_CREAT 0o600 with O_NOFOLLOW (planted symlinks refused) and tighten - a pre-existing file via fchmod (its failure is fatal iff ``strict_chmod``). Raises OSError; - caller owns the fd.""" + """Open ``path`` O_APPEND|O_CREAT 0o600 with O_NOFOLLOW (planted symlinks refused); fchmod + tightens a pre-existing file (failure fatal iff ``strict_chmod``). Raises OSError; caller owns the fd.""" fd = os.open(str(path), os.O_WRONLY | os.O_CREAT | os.O_APPEND | _O_NOFOLLOW, 0o600) try: os.fchmod(fd, 0o600) @@ -667,8 +634,8 @@ def _open_private_append(path: Path, *, strict_chmod: bool) -> int: def ensure_audit_log(audit_path: Path) -> None: - """Pre-create the audit log 0o600 (forward-compat: v0.39 never writes it, no ``log.audit_path``). - Raises RuntimeError on any OSError (planted symlink, immutable dir, full disk).""" + """Pre-create the audit log 0o600 (forward-compat: v0.39 never writes it). RuntimeError on any + OSError (planted symlink, immutable dir, full disk).""" try: os.close(_open_private_append(audit_path, strict_chmod=True)) except OSError as exc: @@ -679,8 +646,8 @@ def ensure_audit_log(audit_path: Path) -> None: def _write_state_file_atomic(state: Path, name: str, dump) -> Path: - """Write ``/`` via a 0600 temp file + atomic replace: the file holds proxy - token values, and chmod-after-replace would leave a world-readable TOCTOU window.""" + """0600 temp file + atomic replace: the file holds proxy tokens, and chmod-after-replace would + leave a world-readable TOCTOU window.""" tmp_path = state / f".{name}.tmp" with open(tmp_path, "w", encoding="utf-8") as f: dump(f) @@ -744,9 +711,8 @@ def _env_names(available_env_names: Optional[List[str]]) -> set: def discover_provider_mappings(*, available_env_names: Optional[List[str]] = None) -> List[TokenMapping]: - """Mint a TokenMapping for every known provider whose env var is set (bearer providers first). - Canonical OR any alias present -> ONE mapping keyed on the canonical name (the subprocess-env - builder mirrors an alias value into it).""" + """One TokenMapping per known provider whose env var is set (bearer providers first). Canonical + OR any alias present -> ONE mapping on the canonical name (the subprocess-env builder mirrors aliases).""" names = _env_names(available_env_names) specs = [(n, h, ("Authorization",), ()) for n, h in _BEARER_PROVIDERS.items()] + [ (n, tuple(s["hosts"]), tuple(s["match_headers"]), tuple(s.get("aliases") or ())) @@ -766,9 +732,8 @@ def discover_uncovered_providers(*, available_env_names: Optional[List[str]] = N def merge_mappings(*, existing: List[TokenMapping], discovered: List[TokenMapping], rotate: bool = False) -> List[TokenMapping]: - """Combine existing and freshly discovered mappings. Tokens already in ``existing`` are - preserved (containers baked with them keep working) while hosts/headers/aliases are refreshed - from ``discovered``; ``rotate=True`` re-mints everything; undiscovered providers are dropped.""" + """Tokens already in ``existing`` are preserved (containers baked with them keep working) while + hosts/headers/aliases refresh from ``discovered``; ``rotate=True`` re-mints; undiscovered providers drop.""" by_name = {} if rotate else {m.real_env_name: m for m in existing} return [ replace(d, proxy_token=by_name[d.real_env_name].proxy_token) if d.real_env_name in by_name else d @@ -807,8 +772,8 @@ def _persisted_nonce_path() -> Path: def _read_persisted_nonce() -> Optional[str]: - """Nonce from disk, or None if missing/unreadable/empty/not owned by us (callers fall back - to argv0 matching). O_NOFOLLOW: this read decides whether stop_proxy SIGKILLs a PID.""" + """Nonce from disk, or None if missing/unreadable/empty/not owned by us (callers fall back to + argv0 matching). O_NOFOLLOW: this read decides whether stop_proxy SIGKILLs a PID.""" try: fd = os.open(str(_persisted_nonce_path()), os.O_RDONLY | _O_NOFOLLOW) except OSError: @@ -822,9 +787,9 @@ def _read_persisted_nonce() -> Optional[str]: def _pid_alive(pid: int) -> bool: - """True iff ``pid`` is alive AND is an iron-proxy process. PID-reuse defense, in priority - order: nonce in /proc//environ, argv[0] basename in /proc//cmdline, ``ps -o comm=`` - basename. A loose ``"iron-proxy" in cmdline`` match would hit ``tail iron-proxy.log``.""" + """True iff ``pid`` is alive AND an iron-proxy process. PID-reuse defense, in priority order: + nonce in /proc//environ, argv[0] basename in /proc//cmdline, ``ps -o comm=`` basename + (a loose ``"iron-proxy" in cmdline`` would hit ``tail iron-proxy.log``).""" if pid <= 0: return False try: @@ -838,7 +803,6 @@ def _pid_alive(pid: int) -> bool: os.kill(pid, 0) # windows-footgun: ok — POSIX-only branch except (ProcessLookupError, PermissionError, OSError): return False - # Strong proof: nonce from this process's start_proxy and/or the on-disk sibling file. nonce_candidates = [n for n in dict.fromkeys((_proxy_nonce, _read_persisted_nonce())) if n] if nonce_candidates: @@ -864,10 +828,8 @@ def start_proxy( install_if_missing: bool = True, refresh_secrets_from_bitwarden: bool = False, bitwarden_config: Optional[Dict] = None, ) -> ProxyStatus: """Spawn iron-proxy as a managed background subprocess (idempotent if already running). - ``refresh_secrets_from_bitwarden=True`` re-fetches upstream secrets from BWS at startup — - the rotation promise of ``credential_source: bitwarden``.""" + ``refresh_secrets_from_bitwarden`` re-fetches secrets from BWS — the ``credential_source: bitwarden`` rotation promise.""" global _proxy_nonce - existing = _read_pid() if existing and _pid_alive(existing): return get_status() @@ -877,7 +839,6 @@ def start_proxy( cfg = config_path or (_proxy_state_dir() / "proxy.yaml") if not cfg.exists(): raise RuntimeError(f"iron-proxy config not found at {cfg}. Run `hermes egress setup` first.") - # Minimal env: os.environ.copy() would expose every operator secret via /proc//environ. env = _build_proxy_subprocess_env( extra_env=extra_env, refresh_from_bitwarden=refresh_secrets_from_bitwarden, bitwarden_config=bitwarden_config, @@ -888,12 +849,9 @@ def start_proxy( # Per-start nonce for PID-recycling defense; module-global is fine (one proxy per process). _proxy_nonce = hashlib.sha256(os.urandom(16)).hexdigest() env[_HERMES_IRON_PROXY_NONCE_ENV] = _proxy_nonce - log_path = _proxy_state_dir() / "iron-proxy.log" proc = _spawn_daemon(bin_path, cfg, env, log_path) - - # Pidfile BEFORE the listening poll: if the parent dies mid-poll, `hermes egress stop` - # can still clean up the orphan. + # Pidfile BEFORE the listening poll so `hermes egress stop` can clean an orphan if the parent dies mid-poll. pidfile = _pidfile() try: _write_pidfile_safely(pidfile, proc.pid) @@ -916,8 +874,7 @@ def start_proxy( pidfile.unlink(missing_ok=True) raise KeyboardInterrupt() - # Poll with timeout (the Go binary is up in <200ms). Probe the CONFIGURED bind host — on - # Linux that's the docker bridge, where loopback never connects. + # Probe the CONFIGURED bind host (on Linux the docker bridge, where loopback never connects). probe_host, tunnel_port = _probe_target() install_handlers = platform.system() != "Windows" and threading.current_thread() is threading.main_thread() prev_sigint = prev_sigterm = None @@ -930,22 +887,20 @@ def start_proxy( if install_handlers: signal.signal(signal.SIGINT, prev_sigint) signal.signal(signal.SIGTERM, prev_sigterm) - # Process may have died right at deadline. if proc.poll() is not None: raise _exited_error() # Alive-but-not-listening is a failure: an orphan holding the port breaks every restart. if not listening: raise _abort(f"iron-proxy did not bind {probe_host}:{tunnel_port} within {_STARTUP_GRACE_SECONDS}s. Process was killed. ", kill=True) - logger.info("Started iron-proxy pid=%s config=%s", proc.pid, cfg) return get_status() def _spawn_daemon(bin_path: Path, cfg: Path, env: Dict[str, str], log_path: Path) -> "subprocess.Popen": - """Popen the daemon with stdout/stderr appended to ``log_path`` (0o600 from the first byte, - O_NOFOLLOW so a planted symlink e.g. to authorized_keys can't receive output, owner-checked). - Our log fd is closed right after Popen — the child has its dup.""" + """Popen with stdout/stderr appended to ``log_path`` (0o600 from the first byte, O_NOFOLLOW so a + planted symlink e.g. to authorized_keys can't receive output, owner-checked). Our log fd is + closed right after Popen — the child has its dup.""" try: log_fd = _open_private_append(log_path, strict_chmod=False) except OSError as exc: @@ -954,7 +909,6 @@ def _spawn_daemon(bin_path: Path, cfg: Path, env: Dict[str, str], log_path: Path st = os.fstat(log_fd) os.close(log_fd) raise RuntimeError(f"iron-proxy log {log_path} has unexpected owner uid={st.st_uid}; refusing to write.") - try: # start_new_session is POSIX-only (Windows isn't supported anyway — no upstream binary). popen_kwargs: Dict = dict(env=env, stdin=subprocess.DEVNULL, stdout=log_fd, stderr=subprocess.STDOUT) @@ -969,8 +923,8 @@ def _spawn_daemon(bin_path: Path, cfg: Path, env: Dict[str, str], log_path: Path def _await_listening(proc: "subprocess.Popen", host: str, port: int, *, on_exit) -> bool: - """Poll until ``host:port`` accepts or the grace window lapses (do-while: checks at least - once even with a 0s window). Raises ``on_exit()`` if the child dies meanwhile.""" + """Poll until ``host:port`` accepts or the grace window lapses (do-while: at least one check + even with a 0s window). Raises ``on_exit()`` if the child dies meanwhile.""" deadline = time.time() + _STARTUP_GRACE_SECONDS while True: if proc.poll() is not None: @@ -983,9 +937,8 @@ def _await_listening(proc: "subprocess.Popen", host: str, port: int, *, on_exit) def _write_pidfile_safely(pidfile: Path, pid: int) -> None: - """Write the pidfile with O_EXCL + O_NOFOLLOW + ownership check, then persist the nonce. - O_EXCL on an existing pidfile means either a concurrent start (fail cleanly) or a stale - crash leftover (unlink and retry once).""" + """O_EXCL + O_NOFOLLOW + ownership check, then persist the nonce. An existing pidfile is either + a concurrent start (fail cleanly) or a stale crash leftover (unlink and retry once).""" open_flags = os.O_WRONLY | os.O_CREAT | os.O_EXCL | _O_NOFOLLOW try: fd = os.open(str(pidfile), open_flags, 0o600) @@ -998,14 +951,12 @@ def _write_pidfile_safely(pidfile: Path, pid: int) -> None: except OSError as exc: # ELOOP from a planted symlink at the pidfile path. raise RuntimeError(f"Refusing to write pidfile {pidfile}: {exc}. Remove that path manually and retry.") from exc - try: if not _fd_owned_by_us(fd): raise RuntimeError(f"pidfile {pidfile} has unexpected owner uid={os.fstat(fd).st_uid}") os.write(fd, str(pid).encode("utf-8")) finally: os.close(fd) - # Best-effort nonce sibling (0o600); without it stop falls back to argv0 matching. if _proxy_nonce: with suppress(OSError): @@ -1030,15 +981,12 @@ def _kill_and_wait(proc: "subprocess.Popen", *, grace_seconds: int = 2) -> None: def _build_proxy_subprocess_env( *, extra_env: Optional[Dict[str, str]] = None, refresh_from_bitwarden: bool = False, bitwarden_config: Optional[Dict] = None, ) -> Dict[str, str]: - """Minimal env for the daemon: allowlisted infra vars + the secrets named in mappings. With - ``refresh_from_bitwarden`` and a populated ``bitwarden_config``, secrets are fetched from BWS - at startup (the rotation guarantee); without ``allow_env_fallback`` any BWS shortfall fails - closed rather than silently keeping stale host-env values.""" + """Allowlisted infra vars + the secrets named in mappings. With ``refresh_from_bitwarden`` and a + populated ``bitwarden_config`` secrets come from BWS (the rotation guarantee); without + ``allow_env_fallback`` any BWS shortfall fails closed instead of keeping stale host-env values.""" env = _allowlisted_env() parent = os.environ - - # Forward ONLY the secrets named in mappings. For alias providers the rule is keyed on the - # canonical name; mirror the alias value into it when only the alias is set. + # Forward ONLY mapped secrets; the rule is keyed on the canonical name, so mirror an alias value into it. mappings = load_mappings() needed = {m.real_env_name for m in mappings} alias_sources = {m.real_env_name: m.alias_env_names for m in mappings if m.alias_env_names} @@ -1046,10 +994,8 @@ def _build_proxy_subprocess_env( source = name if name in parent else next((a for a in alias_sources.get(name, ()) if parent.get(a)), None) if source is not None: env[name] = parent[source] - if refresh_from_bitwarden and bitwarden_config: _refresh_secrets_from_bitwarden(env, needed, bitwarden_config, bool(bitwarden_config.get("allow_env_fallback"))) - # Caller overrides win (wizard test secrets), then strip proxy-recursion vars regardless. if extra_env: env.update(extra_env) @@ -1060,16 +1006,15 @@ def _build_proxy_subprocess_env( def _bitwarden_shortfall(allow_env_fallback: bool, error: str, warning: str, *args) -> None: - """Fail closed with ``error`` unless the operator opted into the legacy host-env fallback, - in which case only ``warning`` (``%``-formatted with ``args``) is logged.""" + """Raise ``error`` unless the operator opted into the legacy host-env fallback (then log ``warning``).""" if not allow_env_fallback: raise RuntimeError(error) logger.warning(warning, *args) def _refresh_secrets_from_bitwarden(env: Dict[str, str], needed: set, bitwarden_config: Dict, allow_env_fallback: bool) -> None: - """Overwrite ``env[needed]`` with fresh BWS values (uncached fetch). Only mapped names are - injected so unrelated BWS secrets never leak into the daemon env.""" + """Overwrite ``env[needed]`` with fresh (uncached) BWS values; only mapped names are injected so + unrelated BWS secrets never leak into the daemon env.""" try: # Lazy: the bitwarden module isn't importable in every install. from agent.secret_sources import bitwarden as bw @@ -1095,7 +1040,6 @@ def _refresh_secrets_from_bitwarden(env: Dict[str, str], needed: set, bitwarden_ ) from exc logger.warning("Bitwarden refresh module unavailable at proxy start, falling back to parent env (allow_env_fallback=true): %s", exc) return - missing = sorted(needed - set(secrets)) for n in needed: if n in secrets: @@ -1128,7 +1072,6 @@ def stop_proxy() -> bool: if not pid or not _pid_alive(pid): _forget_daemon() return False - # Capture starttime BEFORE signalling: if the pid is recycled mid-wait, abort the SIGKILL. starttime_before = _pid_proc_starttime(pid) try: @@ -1136,7 +1079,6 @@ def stop_proxy() -> bool: except ProcessLookupError: _forget_daemon() return False - # Up to 5s for graceful exit, then SIGKILL — unless the pid was recycled meanwhile. deadline = time.time() + 5.0 while time.time() < deadline: @@ -1151,7 +1093,6 @@ def stop_proxy() -> bool: else: with suppress(ProcessLookupError): os.kill(pid, _KILL_SIGNAL) - _forget_daemon() logger.info("Stopped iron-proxy pid=%s", pid) return True