From 16e774f223372c4017fd0ce5abde21b5c2e9362c Mon Sep 17 00:00:00 2001 From: Victor Kyriazakos Date: Wed, 22 Jul 2026 01:59:55 +0000 Subject: [PATCH] fix(monitoring): harden gateway health OTLP egress --- agent/monitoring/emitter.py | 7 +- agent/monitoring/gateway_health.py | 79 +++++++-- agent/monitoring/gateway_health_export.py | 161 ++++++++++++++---- agent/monitoring/otlp_exporter.py | 10 +- docs/observability/monitoring.md | 5 +- hermes_cli/config.py | 2 +- tests/monitoring/test_emitter.py | 19 +++ .../monitoring/test_gateway_health_export.py | 154 ++++++++++++++++- 8 files changed, 379 insertions(+), 58 deletions(-) diff --git a/agent/monitoring/emitter.py b/agent/monitoring/emitter.py index f7dee1bdf4f..52bee5e05cf 100644 --- a/agent/monitoring/emitter.py +++ b/agent/monitoring/emitter.py @@ -114,12 +114,15 @@ class MonitoringEmitter: """Register a live batch subscriber (callable(batch: list[dict])).""" if callback not in self._subscribers: self._subscribers.append(callback) + self._enabled = True def unsubscribe(self, callback) -> None: try: self._subscribers.remove(callback) except ValueError: pass + if not self._subscribers: + self._enabled = False # ── introspection / shutdown (tests, CLI) ─────────────────────────────── def flush(self, timeout: float = 2.0) -> None: @@ -160,7 +163,9 @@ def get_emitter() -> MonitoringEmitter: return _EMITTER with _EMITTER_LOCK: if _EMITTER is None: - _EMITTER = MonitoringEmitter() + # Collection is opt-in. A plane exporter enables the singleton by + # attaching its first subscriber; until then producers are no-ops. + _EMITTER = MonitoringEmitter(enabled=False) return _EMITTER diff --git a/agent/monitoring/gateway_health.py b/agent/monitoring/gateway_health.py index 8ba31386e00..937b57ba442 100644 --- a/agent/monitoring/gateway_health.py +++ b/agent/monitoring/gateway_health.py @@ -8,6 +8,7 @@ session history, audit records, or product analytics belong here. from __future__ import annotations +import hashlib import logging from dataclasses import dataclass from typing import Any, Dict, List, Optional @@ -30,6 +31,10 @@ class GatewayHealthSnapshot: _RUNNING_PLATFORM_STATES = {"running", "connected", "ok", "ready"} _FATAL_PLATFORM_STATES = {"fatal", "degraded", "error", "failed"} +_KNOWN_STATES = _RUNNING_PLATFORM_STATES | _FATAL_PLATFORM_STATES | { + "starting", "draining", "stopping", "stopped", "startup_failed", "unknown" +} +_SUPERVISION_MODES = {"systemd", "s6", "container", "launchd", "manual", "unknown"} def _allowed_logger(name: str) -> bool: @@ -70,6 +75,46 @@ def classify_gateway_error(raw: Any) -> str: return "unknown" +def classify_exit_reason( + raw: Any, *, state: Any, restart_requested: bool +) -> Optional[str]: + """Reduce free-form shutdown text to a bounded operational class.""" + if restart_requested: + return "restart_requested" + state_name = str(state or "").lower() + if raw is None and state_name != "startup_failed": + return None + classified = classify_gateway_error(raw) + if state_name == "startup_failed": + return classified if classified != "unknown" else "startup_failed" + text = str(raw or "").lower() + if "signal" in text or "sigterm" in text or "sigint" in text: + return "signal" + if state_name == "stopped" and any(word in text for word in ("shutdown", "stop")): + return "planned_stop" + return classified + + +def _bounded_state(raw: Any) -> str: + state = str(raw or "unknown").lower() + return state if state in _KNOWN_STATES else "unknown" + + +def _safe_metric_value(raw: Any, *, limit: int = 128) -> str: + try: + from agent.monitoring.redaction import redact_for_export + value = redact_for_export(str(raw or "")) or "unknown" + except Exception: + return "unknown" + return value[:limit] + + +def _safe_instance_id(raw: Any) -> str: + """Return a stable opaque instance key without exporting the source ID.""" + value = str(raw or "unknown").encode("utf-8", errors="replace") + return f"sha256:{hashlib.sha256(value).hexdigest()[:24]}" + + def subsystem_for_logger(logger_name: str) -> str: if logger_name.startswith("gateway.platforms."): parts = logger_name.split(".") @@ -120,11 +165,11 @@ def _derive_drainable(gateway_running: bool, gateway_state: Any) -> bool: def _base_attrs(*, profile: str, install_id: str, version: str, supervision_mode: str) -> Dict[str, str]: + mode = str(supervision_mode or "unknown").lower() return { - "hermes.profile": str(profile), - "service.instance.id": str(install_id), - "service.version": str(version), - "hermes.supervision_mode": str(supervision_mode), + "service.instance.id": _safe_instance_id(install_id), + "service.version": _safe_metric_value(version, limit=64), + "hermes.supervision_mode": mode if mode in _SUPERVISION_MODES else "unknown", } @@ -132,7 +177,7 @@ def _metric(name: str, value: int | float, attrs: Dict[str, str], **extra: str) out = dict(attrs) for key, val in extra.items(): if val is not None: - out[key] = str(val) + out[key] = _safe_metric_value(val) return GatewayMetric(name=name, value=value, attributes=out) @@ -147,7 +192,7 @@ def build_gateway_health_snapshot( ) -> GatewayHealthSnapshot: """Convert gateway_state.json-compatible runtime state into P0 signals.""" runtime = runtime or {} - gateway_state = runtime.get("gateway_state") + gateway_state = _bounded_state(runtime.get("gateway_state")) active_agents = _parse_active_agents(runtime.get("active_agents", 0)) busy = _derive_busy(gateway_running, gateway_state, active_agents) drainable = _derive_drainable(gateway_running, gateway_state) @@ -168,7 +213,7 @@ def build_gateway_health_snapshot( events: list[GatewayHealthEvent | GatewayDiagnosticEvent] = [] for platform, pdata in platforms.items(): pdata = pdata if isinstance(pdata, dict) else {} - state = str(pdata.get("state") or "unknown").lower() + state = _bounded_state(pdata.get("state")) raw_error = pdata.get("error_code") or pdata.get("error_message") error_code = classify_gateway_error(raw_error) is_up = state in _RUNNING_PLATFORM_STATES @@ -244,15 +289,19 @@ def emit_runtime_status_transition(previous: Optional[dict[str, Any]], current: out: list[GatewayHealthEvent | GatewayDiagnosticEvent] = [] profile = _safe_profile() version = _safe_version() - old_gateway_state = str((previous or {}).get("gateway_state")) if (previous or {}).get("gateway_state") is not None else None - new_gateway_state = str(current.get("gateway_state")) if current.get("gateway_state") is not None else None + old_gateway_state = _bounded_state((previous or {}).get("gateway_state")) if (previous or {}).get("gateway_state") is not None else None + new_gateway_state = _bounded_state(current.get("gateway_state")) if current.get("gateway_state") is not None else None if old_gateway_state != new_gateway_state and new_gateway_state: out.append(GatewayHealthEvent( name="gateway.lifecycle", gateway_state=new_gateway_state, old_state=old_gateway_state, new_state=new_gateway_state, - exit_reason=current.get("exit_reason"), + exit_reason=classify_exit_reason( + current.get("exit_reason"), + state=new_gateway_state, + restart_requested=bool(current.get("restart_requested")), + ), restart_requested=bool(current.get("restart_requested")), active_agents=_parse_active_agents(current.get("active_agents", 0)), profile=profile, @@ -276,7 +325,11 @@ def emit_runtime_status_transition(previous: Optional[dict[str, Any]], current: gateway_state=new_gateway_state, old_state=old_gateway_state, new_state=new_gateway_state, - exit_reason=current.get("exit_reason"), + exit_reason=classify_exit_reason( + current.get("exit_reason"), + state=new_gateway_state, + restart_requested=bool(current.get("restart_requested")), + ), restart_requested=bool(current.get("restart_requested")), active_agents=_parse_active_agents(current.get("active_agents", 0)), profile=profile, @@ -292,8 +345,8 @@ def emit_runtime_status_transition(previous: Optional[dict[str, Any]], current: pdata = pdata if isinstance(pdata, dict) else {} prev_raw = old_platforms.get(platform, {}) prev = prev_raw if isinstance(prev_raw, dict) else {} - old_state = str(prev.get("state")) if prev.get("state") is not None else None - new_state = str(pdata.get("state")) if pdata.get("state") is not None else None + old_state = _bounded_state(prev.get("state")) if prev.get("state") is not None else None + new_state = _bounded_state(pdata.get("state")) if pdata.get("state") is not None else None if old_state == new_state or not new_state: continue error_code = classify_gateway_error(pdata.get("error_code") or pdata.get("error_message")) diff --git a/agent/monitoring/gateway_health_export.py b/agent/monitoring/gateway_health_export.py index 969fa970635..9295072da98 100644 --- a/agent/monitoring/gateway_health_export.py +++ b/agent/monitoring/gateway_health_export.py @@ -10,12 +10,77 @@ from __future__ import annotations import logging import os +import re import threading from dataclasses import dataclass from typing import Any, Dict, Optional logger = logging.getLogger(__name__) +_RESOURCE_ATTRIBUTE_KEYS = frozenset({ + "service.name", + "service.namespace", + "service.version", + "service.instance.id", + "deployment.environment.name", + "cloud.provider", + "cloud.platform", + "cloud.region", + "telemetry.scope", +}) +_DIAGNOSTIC_ATTRIBUTE_KEYS = frozenset({ + "name", + "subsystem", + "error_class", + "error_code", + "platform", + "old_state", + "new_state", + "version", + "severity", +}) +_SAFE_RESOURCE_VALUE = re.compile(r"^[A-Za-z0-9._:/-]{1,128}$") + + +def _redact_string(raw: Any, *, limit: int = 500) -> str: + try: + from agent.monitoring.redaction import redact_for_export + return (redact_for_export(str(raw or "")) or "[redacted]")[:limit] + except Exception: + return "[redaction-unavailable]" + + +def _safe_resource_attributes(raw: Any) -> Dict[str, str]: + """Allowlist bounded resource labels and reject values changed by redaction.""" + attrs: Dict[str, str] = {} + if not isinstance(raw, dict): + return attrs + for key, value in raw.items(): + key = str(key) + if key not in _RESOURCE_ATTRIBUTE_KEYS or value is None: + continue + if key == "service.instance.id": + from agent.monitoring.gateway_health import _safe_instance_id + attrs[key] = _safe_instance_id(value) + continue + text = str(value) + if not _SAFE_RESOURCE_VALUE.fullmatch(text): + continue + if _redact_string(text, limit=128) != text: + continue + attrs[key] = text + return attrs + + +def _diagnostic_log_attributes(event: Dict[str, Any]) -> Dict[str, Any]: + attrs: Dict[str, Any] = {} + for key in _DIAGNOSTIC_ATTRIBUTE_KEYS: + value = event.get(key) + if value is None: + continue + attrs[f"hermes.{key}"] = _redact_string(value) if isinstance(value, str) else value + return attrs + @dataclass(slots=True) class GatewayHealthExportRuntime: @@ -32,27 +97,52 @@ class GatewayHealthExportRuntime: if self.stop_event is not None: self.stop_event.set() if self.thread is not None: - self.thread.join(timeout=2.0) + self.thread.join(timeout=0.25) if self.log_handler is not None: try: logging.getLogger().removeHandler(self.log_handler) except Exception: pass - if self.streamer is not None: - try: - self.streamer.shutdown() - except Exception: - pass - if self.log_streamer is not None: - try: - self.log_streamer.shutdown() - except Exception: - pass - if self.metric_provider is not None: - try: - self.metric_provider.shutdown() - except Exception: - pass + + # Detach synchronously so no new records are collected during a slow + # exporter shutdown. Network flush/close then runs under one bounded + # daemon-thread deadline and can never delay gateway teardown. + try: + from agent.monitoring.emitter import get_emitter + emitter = get_emitter() + if self.streamer is not None: + emitter.unsubscribe(self.streamer) + if self.log_streamer is not None: + emitter.unsubscribe(self.log_streamer) + except Exception: + pass + + closeables = [ + item for item in (self.streamer, self.log_streamer, self.metric_provider) + if item is not None + ] + + def _close() -> None: + for item in closeables: + try: + item.shutdown() + except Exception: + pass + + if closeables: + worker = threading.Thread( + target=_close, + name="hermes-gateway-health-export-shutdown", + daemon=True, + ) + worker.start() + worker.join(timeout=2.0) + + self.streamer = None + self.log_streamer = None + self.metric_provider = None + self.thread = None + self.stop_event = None def _gateway_health_config(config: Dict[str, Any]) -> Dict[str, Any]: @@ -209,9 +299,9 @@ def _start_metric_provider(config: Dict[str, Any], sdk: Dict[str, Any]) -> Any: exporter = sdk["OTLPMetricExporter"](endpoint=endpoint, headers=headers or None) interval_ms = max(5, int(gh.get("export_interval_seconds", 60))) * 1000 reader = sdk["PeriodicExportingMetricReader"](exporter, export_interval_millis=interval_ms) - resource_attrs = dict((gh.get("resource_attributes") or {})) - resource_attrs.setdefault("service.name", "hermes-gateway") - resource_attrs.setdefault("telemetry.scope", "gateway_health") + resource_attrs = _safe_resource_attributes(gh.get("resource_attributes")) + resource_attrs["service.name"] = "hermes-gateway" + resource_attrs["telemetry.scope"] = "gateway_health" provider = sdk["MeterProvider"]( metric_readers=[reader], resource=sdk["Resource"].create(resource_attrs), @@ -267,9 +357,9 @@ class GatewayDiagnosticLogStreamer: gh = _gateway_health_config(config) headers = _resolve_headers(otlp.get("headers_env")) endpoint = _logs_endpoint(str(otlp.get("endpoint"))) - resource_attrs = dict((gh.get("resource_attributes") or {})) - resource_attrs.setdefault("service.name", "hermes-gateway") - resource_attrs.setdefault("telemetry.scope", "gateway_diagnostics") + resource_attrs = _safe_resource_attributes(gh.get("resource_attributes")) + resource_attrs["service.name"] = "hermes-gateway" + resource_attrs["telemetry.scope"] = "gateway_diagnostics" self._provider = sdk["LoggerProvider"](resource=sdk["Resource"].create(resource_attrs)) self._processor = sdk["BatchLogRecordProcessor"]( sdk["OTLPLogExporter"](endpoint=endpoint, headers=headers or None) @@ -284,11 +374,7 @@ class GatewayDiagnosticLogStreamer: for ev in batch: if ev.get("event") != "gateway_diagnostic": continue - attrs = { - f"hermes.{key}": val - for key, val in ev.items() - if key not in {"event", "redacted_message", "ts_ns"} and val is not None - } + attrs = _diagnostic_log_attributes(ev) body = ev.get("redacted_message") or ev.get("name") or "gateway diagnostic" record = self._LogRecord( timestamp=ev.get("ts_ns"), @@ -297,7 +383,7 @@ class GatewayDiagnosticLogStreamer: trace_flags=self._sdk["TraceFlags"].DEFAULT, severity_text=str(ev.get("severity") or "warning").upper(), severity_number=_severity_number(self._sdk, ev.get("severity")), - body=str(body), + body=_redact_string(body), attributes=attrs, ) self._logger.emit(record) @@ -337,7 +423,7 @@ def _start_snapshot_thread(config: Dict[str, Any], stop_event: threading.Event) def _attach_log_handler(config: Dict[str, Any]) -> Any: gh = _gateway_health_config(config) - if not gh.get("warning_error_events_enabled", True): + if not gh.get("diagnostic_events_enabled", True) or not gh.get("warning_error_events_enabled", True): return None from agent.monitoring.gateway_health import GatewayDiagnosticLogHandler handler = GatewayDiagnosticLogHandler(profile=_profile(), version=_version()) @@ -382,20 +468,25 @@ def start_gateway_health_export(config: Dict[str, Any]) -> GatewayHealthExportRu try: from agent.monitoring import otlp_exporter runtime.streamer = otlp_exporter.start_streaming(config, event_filter=_gateway_health_event) + if runtime.streamer is None: + raise RuntimeError("gateway health span streamer did not start") runtime.log_streamer = _start_diagnostic_log_streamer(config, sdk) except Exception: 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: logger.debug("gateway diagnostic log handler failed to attach", exc_info=True) - try: - _emit_snapshot_events(config) - runtime.stop_event = threading.Event() - runtime.thread = _start_snapshot_thread(config, runtime.stop_event) - except Exception: - logger.debug("gateway health snapshot thread failed to start", exc_info=True) + if gh.get("diagnostic_events_enabled", True): + try: + _emit_snapshot_events(config) + runtime.stop_event = threading.Event() + runtime.thread = _start_snapshot_thread(config, runtime.stop_event) + except Exception: + logger.debug("gateway health snapshot thread failed to start", exc_info=True) return runtime diff --git a/agent/monitoring/otlp_exporter.py b/agent/monitoring/otlp_exporter.py index 01eefa53506..27c600dcf8b 100644 --- a/agent/monitoring/otlp_exporter.py +++ b/agent/monitoring/otlp_exporter.py @@ -136,15 +136,21 @@ def _span_attrs(ev: Dict[str, Any]) -> Dict[str, Any]: "gateway_health": ("name", "gateway_state", "old_state", "new_state", "exit_reason", "restart_requested", "active_agents", "gateway_busy", "gateway_drainable", "platform_count", - "fatal_platform_count", "profile", "install_id", "version", + "fatal_platform_count", "version", "supervision_mode", "pid"), "gateway_diagnostic": ("name", "subsystem", "error_class", "error_code", "redacted_message", "platform", "old_state", "new_state", - "profile", "version", "severity"), + "version", "severity"), } for col in keep_by_kind.get(kind, ()): # type: ignore[arg-type] v = ev.get(col) if v is not None: + if isinstance(v, str): + try: + from agent.monitoring.redaction import redact_for_export + v = (redact_for_export(v) or "[redacted]")[:500] + except Exception: + v = "[redaction-unavailable]" attrs[f"hermes.{col}"] = v return attrs diff --git a/docs/observability/monitoring.md b/docs/observability/monitoring.md index 1ee9d4b5a8d..9ca8f2ce12a 100644 --- a/docs/observability/monitoring.md +++ b/docs/observability/monitoring.md @@ -19,8 +19,9 @@ trajectory capture is a separate plane served by the NeMo Relay integration | Health/lifecycle events | `/v1/traces` | `gateway.lifecycle` state transitions (`starting -> running -> draining -> stopped`, `startup_failed`, exit), `gateway.health_snapshot`, platform state changes | | Diagnostics | `/v1/logs` | Warning/error gateway log events with secrets AND PII scrubbed in-process before egress (`[redacted]` / `[email]`), bounded error classes | -Every signal carries resource attributes (`service.name`, profile, version, -install id, supervision mode) so an operator can tell instances apart. +Signals carry `service.name`, version, supervision mode, and a stable one-way +hash of the install id so an operator can distinguish instances without +exporting account/profile identity or the raw install identifier. ## Enabling diff --git a/hermes_cli/config.py b/hermes_cli/config.py index bd07d7d6a4e..1e2d1ba44a4 100644 --- a/hermes_cli/config.py +++ b/hermes_cli/config.py @@ -3051,7 +3051,7 @@ DEFAULT_CONFIG = { "logs_export_interval_seconds": 5, "resource_attributes": { "service.name": "hermes-gateway", - "deployment.environment": "production", + "deployment.environment.name": "production", }, }, # OTLP destination. headers_env maps header names to ENVIRONMENT diff --git a/tests/monitoring/test_emitter.py b/tests/monitoring/test_emitter.py index a0f5c2d889d..98ad7b0389b 100644 --- a/tests/monitoring/test_emitter.py +++ b/tests/monitoring/test_emitter.py @@ -15,6 +15,25 @@ def test_emit_never_raises_when_disabled(): em.close() +def test_process_singleton_stays_dormant_until_subscribed(): + from agent.monitoring import emitter + + emitter.reset_emitter_for_tests() + try: + emitter.emit({"event": "gateway_health", "name": "gateway.lifecycle"}) + singleton = emitter.get_emitter() + assert singleton.stats()["queued"] == 0 + assert singleton._started is False + + subscriber = lambda _batch: None # noqa: E731 + singleton.subscribe(subscriber) + emitter.emit({"event": "gateway_health", "name": "gateway.lifecycle"}) + assert singleton._started is True + singleton.unsubscribe(subscriber) + finally: + emitter.reset_emitter_for_tests() + + def test_emit_accepts_dataclass_and_dict(tmp_path): em = MonitoringEmitter() seen: list = [] diff --git a/tests/monitoring/test_gateway_health_export.py b/tests/monitoring/test_gateway_health_export.py index 3fe8ff190b1..715de681bf8 100644 --- a/tests/monitoring/test_gateway_health_export.py +++ b/tests/monitoring/test_gateway_health_export.py @@ -14,6 +14,8 @@ def test_default_config_keeps_gateway_health_export_disabled(): assert cfg["warning_error_events_enabled"] is True assert cfg["export_interval_seconds"] == 60 assert cfg["logs_export_interval_seconds"] == 5 + assert cfg["resource_attributes"]["deployment.environment.name"] == "production" + assert "deployment.environment" not in cfg["resource_attributes"] # Redaction is always-on and deliberately NOT configurable. assert "redaction" not in cfg @@ -59,11 +61,12 @@ def test_gateway_health_snapshot_maps_runtime_status_to_low_cardinality_metrics( active = next(m for m in snapshot.metrics if m.name == "hermes.gateway.active_agents") assert active.value == 2 assert active.attributes == { - "hermes.profile": "default", - "service.instance.id": "install-1", + "service.instance.id": active.attributes["service.instance.id"], "service.version": "2026.7.test", "hermes.supervision_mode": "manual", } + assert active.attributes["service.instance.id"].startswith("sha256:") + assert "install-1" not in active.attributes["service.instance.id"] busy = next(m for m in snapshot.metrics if m.name == "hermes.gateway.busy") drainable = next(m for m in snapshot.metrics if m.name == "hermes.gateway.drainable") @@ -183,6 +186,7 @@ def test_runtime_status_transition_emits_lifecycle_and_platform_events(monkeypat lifecycle = next(e for e in captured if e["name"] == "gateway.lifecycle") assert lifecycle["old_state"] == "starting" assert lifecycle["new_state"] == "running" + assert lifecycle["exit_reason"] is None platform = next(e for e in captured if e["name"] == "platform.state_change") assert platform["old_state"] == "running" assert platform["new_state"] == "fatal" @@ -198,8 +202,21 @@ def test_runtime_status_transition_emits_startup_failed_and_exit(): old = emitter.emit emitter.emit = lambda event: captured.append(event.to_dict()) # type: ignore[assignment] try: - emit_runtime_status_transition({"gateway_state": "starting"}, {"gateway_state": "startup_failed", "exit_reason": "startup token ***"}) - emit_runtime_status_transition({"gateway_state": "running"}, {"gateway_state": "stopped", "exit_reason": "shutdown", "restart_requested": True}) + emit_runtime_status_transition( + {"gateway_state": "starting"}, + { + "gateway_state": "startup_failed", + "exit_reason": "Bearer top-secret-token rejected for user@example.com", + }, + ) + emit_runtime_status_transition( + {"gateway_state": "running"}, + { + "gateway_state": "stopped", + "exit_reason": "shutdown requested by user@example.com", + "restart_requested": True, + }, + ) finally: emitter.emit = old # type: ignore[assignment] @@ -208,8 +225,11 @@ def test_runtime_status_transition_emits_startup_failed_and_exit(): assert "gateway.exit" in names failed = next(e for e in captured if e["name"] == "gateway.startup_failed") assert "***" not in failed["redacted_message"] + lifecycle = next(e for e in captured if e["name"] == "gateway.lifecycle") + assert lifecycle["exit_reason"] == "auth_failed" exit_event = next(e for e in captured if e["name"] == "gateway.exit") assert exit_event["restart_requested"] is True + assert exit_event["exit_reason"] == "restart_requested" def test_otlp_attrs_include_gateway_transition_fields(): @@ -230,6 +250,72 @@ def test_otlp_attrs_include_gateway_transition_fields(): assert attrs["hermes.restart_requested"] is True +def test_otlp_attrs_redact_strings_and_never_export_profile(): + from agent.monitoring.otlp_exporter import _span_attrs + + attrs = _span_attrs({ + "event": "gateway_health", + "name": "gateway.lifecycle", + "profile": "user@example.com", + "exit_reason": "Bearer top-secret-token for user@example.com", + }) + + assert "hermes.profile" not in attrs + assert "top-secret-token" not in str(attrs) + assert "user@example.com" not in str(attrs) + + +def test_resource_attributes_are_allowlisted_and_sanitized(): + from agent.monitoring.gateway_health_export import _safe_resource_attributes + + attrs = _safe_resource_attributes({ + "service.name": "hermes-gateway", + "service.instance.id": "install-1", + "deployment.environment.name": "staging", + "user.email": "user@example.com", + "authorization": "Bearer top-secret-token", + "custom.request.id": "unbounded", + }) + + assert attrs == { + "service.name": "hermes-gateway", + "service.instance.id": attrs["service.instance.id"], + "deployment.environment.name": "staging", + } + assert attrs["service.instance.id"].startswith("sha256:") + assert "install-1" not in attrs["service.instance.id"] + + +def test_instance_id_hash_is_stable_and_distinguishes_instances(): + from agent.monitoring.gateway_health import _safe_instance_id + + first = _safe_instance_id("install-1") + repeat = _safe_instance_id("install-1") + second = _safe_instance_id("install-2") + + assert first == repeat + assert first != second + assert first.startswith("sha256:") + assert "install-1" not in first + + +def test_diagnostic_log_attributes_are_allowlisted_redacted_and_profile_free(): + from agent.monitoring.gateway_health_export import _diagnostic_log_attributes + + attrs = _diagnostic_log_attributes({ + "event": "gateway_diagnostic", + "name": "platform.fatal", + "subsystem": "platform.slack", + "profile": "user@example.com", + "error_code": "Bearer top-secret-token", + "custom": "must-not-egress", + }) + + assert "hermes.profile" not in attrs + assert "hermes.custom" not in attrs + assert "top-secret-token" not in str(attrs) + + def test_gateway_health_export_start_is_fail_open_when_otlp_missing(monkeypatch): from agent.monitoring import gateway_health_export from agent.monitoring.gateway_health_export import GatewayHealthExportRuntime @@ -259,6 +345,7 @@ def test_gateway_health_export_streams_only_gateway_events(monkeypatch): monkeypatch.setattr(gateway_health_export, "_start_metric_provider", lambda *a, **k: None) monkeypatch.setattr(gateway_health_export, "_require_metrics_sdk", lambda *a, **k: {}) + monkeypatch.setattr(gateway_health_export, "_start_diagnostic_log_streamer", lambda *a, **k: object()) monkeypatch.setattr(gateway_health_export, "_attach_log_handler", lambda *a, **k: None) monkeypatch.setattr(gateway_health_export, "_emit_snapshot_events", lambda *a, **k: None) monkeypatch.setattr(gateway_health_export, "_start_snapshot_thread", lambda *a, **k: None) @@ -301,6 +388,65 @@ def test_gateway_health_export_metric_failure_does_not_start_streamer(monkeypatc assert started == [] +def test_gateway_health_export_diagnostic_partial_start_cleans_up(monkeypatch): + from agent.monitoring import emitter, gateway_health_export, otlp_exporter + + class Streamer: + def __call__(self, _batch): + pass + + def shutdown(self): + pass + + streamer = Streamer() + monkeypatch.setattr(gateway_health_export, "_require_metrics_sdk", lambda *a, **k: {}) + monkeypatch.setattr(gateway_health_export, "_start_metric_provider", lambda *a, **k: None) + monkeypatch.setattr(otlp_exporter, "start_streaming", lambda *a, **k: streamer) + monkeypatch.setattr( + gateway_health_export, + "_start_diagnostic_log_streamer", + lambda *a, **k: (_ for _ in ()).throw(RuntimeError("boom")), + ) + emitter.get_emitter().subscribe(streamer) + + runtime = gateway_health_export.start_gateway_health_export({ + "monitoring": { + "gateway_health_export": {"enabled": True, "metrics_enabled": False}, + "export": {"otlp": {"enabled": True, "endpoint": "http://collector:4318/v1/traces"}}, + } + }) + + assert runtime.enabled is False + assert runtime.reason == "diagnostics_start_failed" + assert streamer not in emitter.get_emitter()._subscribers + + +def test_gateway_health_export_shutdown_is_bounded(): + import threading + import time + + from agent.monitoring.gateway_health_export import GatewayHealthExportRuntime + + release = threading.Event() + + class Blocking: + def shutdown(self): + release.wait(10) + + runtime = GatewayHealthExportRuntime( + enabled=True, + streamer=Blocking(), + log_streamer=Blocking(), + metric_provider=Blocking(), + ) + started = time.monotonic() + runtime.shutdown() + elapsed = time.monotonic() - started + release.set() + + assert elapsed < 2.5 + + def test_otlp_streamer_shutdown_unsubscribes(monkeypatch): from agent.monitoring import emitter from agent.monitoring.otlp_exporter import OTLPStreamer