From 052582de2d7cb3573ff4ece883347eec04eeb5d7 Mon Sep 17 00:00:00 2001 From: Teknium <127238744+teknium1@users.noreply.github.com> Date: Wed, 2 Sep 2026 18:32:02 -0700 Subject: [PATCH] =?UTF-8?q?refactor(agent):=20compact=20monitoring/*=20and?= =?UTF-8?q?=20lsp/workspace=20=E2=80=94=20contextlib.suppress,=20shared=20?= =?UTF-8?q?platform=20error-code=20helper,=20nearest=5Froot=20marker=20pro?= =?UTF-8?q?be,=20docstring=20compaction?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit --- agent/lsp/workspace.py | 98 ++++++++------------ agent/monitoring/__init__.py | 7 +- agent/monitoring/cron_health.py | 43 ++------- agent/monitoring/emitter.py | 28 ++---- agent/monitoring/events.py | 6 +- agent/monitoring/gateway_health.py | 73 +++++---------- agent/monitoring/gateway_health_export.py | 106 ++++++--------------- agent/monitoring/otlp_exporter.py | 107 ++++++---------------- agent/monitoring/policy.py | 4 +- agent/monitoring/redaction.py | 29 ++---- 10 files changed, 144 insertions(+), 357 deletions(-) diff --git a/agent/lsp/workspace.py b/agent/lsp/workspace.py index a33e7908c6..5b220fcab7 100644 --- a/agent/lsp/workspace.py +++ b/agent/lsp/workspace.py @@ -1,11 +1,10 @@ """Workspace and project-root resolution for LSP. -1. **Workspace gate** — LSP only runs when the cwd (or the edited file) sits - inside a git worktree. Files outside any git root never trigger LSP, - which keeps gateway users on user-home cwd's from spawning daemons. -2. **nearest_root** — the per-server project-root walk: up from a start path - looking for marker files (``pyproject.toml``, ``Cargo.toml``, ...), - optionally bailing if an exclude marker shows up first. +1. **Workspace gate** — LSP only runs when the cwd (or the edited file) sits inside a git + worktree, so gateway users on user-home cwd's never spawn daemons. +2. **nearest_root** — the per-server project-root walk: up from a start path looking for marker + files (``pyproject.toml``, ``Cargo.toml``, ...), optionally bailing if an exclude marker + shows up first. """ from __future__ import annotations @@ -16,8 +15,7 @@ from typing import Iterable, Iterator, Optional, Tuple logger = logging.getLogger("agent.lsp.workspace") -# Cache: start dir → (worktree_root, is_git) so repeated calls don't re-stat. -# Cleared on shutdown. +# Cache: start dir → (worktree_root, is_git) so repeated calls don't re-stat. Cleared on shutdown. _workspace_cache: dict = {} # Walk cap: the deepest reasonable monorepo is well under 64 levels; bounds a @@ -26,11 +24,8 @@ _MAX_WALK = 64 def normalize_path(path: str) -> str: - """Expand ``~``, make absolute, collapse ``.``/``..``. - - Symlinks are deliberately NOT resolved — some servers (rust-analyzer's - Cargo workspace identity) care, and we want the path the user typed. - """ + """Expand ``~``, make absolute, collapse ``.``/``..``. Symlinks are deliberately NOT resolved — + some servers (rust-analyzer's Cargo workspace identity) care, and we want the path the user typed.""" return os.path.abspath(os.path.expanduser(path)) @@ -62,11 +57,9 @@ def find_git_worktree(start: str) -> Optional[str]: start_path = _start_dir(start) if start_path is None: return None - cached = _workspace_cache.get(str(start_path)) if cached is not None: return cached[0] - for cur in _walk_up(start_path): try: if (cur / ".git").exists(): @@ -74,30 +67,24 @@ def find_git_worktree(start: str) -> Optional[str]: _workspace_cache[str(start_path)] = (resolved, True) return resolved except OSError: - # Permission error on a parent dir — bail out cleanly. - break - + break # permission error on a parent dir — bail out cleanly _workspace_cache[str(start_path)] = (None, False) return None def is_inside_workspace(path: str, workspace_root: str) -> bool: - """True iff ``path`` is inside (or equal to) ``workspace_root``. - - Symlinks are not resolved: a symlink pointing outside still counts as - outside, matching servers that reject didOpen for unrelated files. - """ + """True iff ``path`` is inside (or equal to) ``workspace_root``. Symlinks are not resolved: a + symlink pointing outside still counts as outside, matching servers that reject didOpen for + unrelated files.""" p = normalize_path(path) root = normalize_path(workspace_root) if p == root: return True # commonpath handles case-insensitive filesystems on macOS/Windows. try: - common = os.path.commonpath([p, root]) + return os.path.commonpath([p, root]) == root except ValueError: - # Different drives on Windows. - return False - return common == root + return False # different drives on Windows def nearest_root( @@ -109,51 +96,42 @@ def nearest_root( ) -> Optional[str]: """Walk up from ``start`` for the directory containing the first matched marker. - Returns ``None`` past ``ceiling`` (or the filesystem root), or when an - exclude marker is found first — the server is gated off for that file - (e.g. typescript skips deno projects when ``deno.json`` precedes - ``package.json``). Marker names are exact filenames — no globs. + Returns ``None`` past ``ceiling`` (or the filesystem root), or when an exclude marker is found + first — the server is gated off for that file (e.g. typescript skips deno projects when + ``deno.json`` precedes ``package.json``). Marker names are exact filenames — no globs. """ start_path = _start_dir(start) if start_path is None: return None ceiling_path = Path(normalize_path(ceiling)) if ceiling else None - markers_list = list(markers) excludes_list = list(excludes) if excludes else [] + def present(cur: Path, names: list) -> bool: + for name in names: + try: + if (cur / name).exists(): + return True + except OSError: + continue + return False + for cur in _walk_up(start_path): # Excludes are checked before markers at each level. - for exc in excludes_list: - try: - if (cur / exc).exists(): - return None - except OSError: - continue - for marker in markers_list: - try: - if (cur / marker).exists(): - return str(cur) - except OSError: - continue + if present(cur, excludes_list): + return None + if present(cur, markers_list): + return str(cur) if ceiling_path is not None and cur == ceiling_path: return None return None -def resolve_workspace_for_file( - file_path: str, - *, - cwd: Optional[str] = None, -) -> Tuple[Optional[str], bool]: - """Return ``(workspace_root, gated_in)`` for a file. - - The cwd's worktree wins when the file is inside it; otherwise the file's - own worktree is the fallback anchor (monorepos / unrelated checkouts). - ``(None, False)`` when neither is in a git worktree. - """ - cwd = cwd or os.getcwd() - cwd_root = find_git_worktree(cwd) +def resolve_workspace_for_file(file_path: str, *, cwd: Optional[str] = None) -> Tuple[Optional[str], bool]: + """Return ``(workspace_root, gated_in)`` for a file. The cwd's worktree wins when the file is + inside it; otherwise the file's own worktree is the fallback anchor (monorepos / unrelated + checkouts). ``(None, False)`` when neither is in a git worktree.""" + cwd_root = find_git_worktree(cwd or os.getcwd()) if cwd_root is not None and is_inside_workspace(file_path, cwd_root): return cwd_root, True file_root = find_git_worktree(file_path) @@ -168,10 +146,6 @@ def clear_cache() -> None: __all__ = [ - "find_git_worktree", - "is_inside_workspace", - "nearest_root", - "normalize_path", - "resolve_workspace_for_file", + "find_git_worktree", "is_inside_workspace", "nearest_root", "normalize_path", "resolve_workspace_for_file", "clear_cache", ] diff --git a/agent/monitoring/__init__.py b/agent/monitoring/__init__.py index 20b4dc40f4..85f95932b9 100644 --- a/agent/monitoring/__init__.py +++ b/agent/monitoring/__init__.py @@ -14,9 +14,4 @@ from . import emitter, events emit = emitter.emit get_emitter = emitter.get_emitter -__all__ = [ - "emitter", - "events", - "emit", - "get_emitter", -] +__all__ = ["emitter", "events", "emit", "get_emitter"] diff --git a/agent/monitoring/cron_health.py b/agent/monitoring/cron_health.py index 04e5122f76..40a2344455 100644 --- a/agent/monitoring/cron_health.py +++ b/agent/monitoring/cron_health.py @@ -9,11 +9,7 @@ from datetime import datetime from typing import Any, Callable, Optional from agent.monitoring.events import CronExecutionEvent -from agent.monitoring.gateway_health import ( - GatewayMetric, - _contains_any, - _safe_instance_id, -) +from agent.monitoring.gateway_health import GatewayMetric, _contains_any, _safe_instance_id from cron.jobs import ( _compute_grace_seconds, get_catch_up_occurrence_count, @@ -27,9 +23,7 @@ from hermes_time import now as _now logger = logging.getLogger(__name__) _KNOWN_STATUSES = {"claimed", "running", "completed", "failed", "unknown"} _KNOWN_SOURCES = {"builtin", "direct", "external"} -_KNOWN_DELIVERY_OUTCOMES = { - "delivered", "failed", "suppressed", "suppressed_acked", "not_configured", -} +_KNOWN_DELIVERY_OUTCOMES = {"delivered", "failed", "suppressed", "suppressed_acked", "not_configured"} _TERMINAL_STATUSES = {"completed", "failed", "unknown"} @@ -47,7 +41,6 @@ _AUTH_RE = re.compile( r"|unauthorized|forbidden|bearer|401|403)\b|\b(?:access|api|refresh) token\b" ) - # Ordered (predicate, class) rules; first match wins. Auth uses word boundaries so # "oauth"/"tokenizer"/"HTTP 4015" do not false-positive. _CRON_ERROR_RULES: tuple[tuple[Callable[[str], bool], str], ...] = ( @@ -86,9 +79,7 @@ def _duration_ms(record: dict[str, Any]) -> Optional[int]: return max(0, duration) -def project_execution_event( - record: dict[str, Any], *, delivery_outcome: Optional[str] = None -) -> CronExecutionEvent: +def project_execution_event(record: dict[str, Any], *, delivery_outcome: Optional[str] = None) -> CronExecutionEvent: status = str(record.get("status") or "unknown").lower() source = str(record.get("source") or "unknown").lower() outcome = str(delivery_outcome).lower() if delivery_outcome is not None else None @@ -98,20 +89,12 @@ def project_execution_event( # Unknown non-empty sources are bucketed as "external" (not dropped to unknown). source=source if source in _KNOWN_SOURCES or source == "unknown" else "external", duration_ms=_duration_ms(record), - delivery_outcome=( - outcome if outcome in _KNOWN_DELIVERY_OUTCOMES else None - ), - error_class=( - classify_cron_error(record.get("error")) - if status in {"failed", "unknown"} - else None - ), + delivery_outcome=outcome if outcome in _KNOWN_DELIVERY_OUTCOMES else None, + error_class=classify_cron_error(record.get("error")) if status in {"failed", "unknown"} else None, ) -def emit_execution_state( - record: Optional[dict[str, Any]], *, delivery_outcome: Optional[str] = None -) -> None: +def emit_execution_state(record: Optional[dict[str, Any]], *, delivery_outcome: Optional[str] = None) -> None: """Best-effort lifecycle emit; terminal states synchronously cross the queue barrier.""" if not record: return @@ -137,8 +120,7 @@ def _is_overdue(job: dict[str, Any], now: datetime) -> bool: try: if next_run.tzinfo is None and now.tzinfo is not None: next_run = next_run.replace(tzinfo=now.tzinfo) - lateness = (now - next_run).total_seconds() - return lateness > _compute_grace_seconds(schedule) + return (now - next_run).total_seconds() > _compute_grace_seconds(schedule) except (TypeError, ValueError): return False @@ -146,11 +128,7 @@ def _is_overdue(job: dict[str, Any], now: datetime) -> bool: def _job_metrics(metrics: list[GatewayMetric]) -> None: enabled = [job for job in load_jobs() if job.get("enabled", True)] metrics.append(GatewayMetric("hermes.cron.jobs.enabled", len(enabled), {})) - metrics.append(GatewayMetric( - "hermes.cron.jobs.overdue", - sum(1 for job in enabled if _is_overdue(job, _now())), - {}, - )) + metrics.append(GatewayMetric("hermes.cron.jobs.overdue", sum(1 for job in enabled if _is_overdue(job, _now())), {})) def _freshness_metric(name: str, reader: Callable[[], Optional[float]]) -> Callable[[list[GatewayMetric]], None]: @@ -191,9 +169,6 @@ def build_cron_health_snapshot() -> CronHealthSnapshot: __all__ = [ - "CronHealthSnapshot", - "build_cron_health_snapshot", - "classify_cron_error", - "emit_execution_state", + "CronHealthSnapshot", "build_cron_health_snapshot", "classify_cron_error", "emit_execution_state", "project_execution_event", ] diff --git a/agent/monitoring/emitter.py b/agent/monitoring/emitter.py index d63ebb73ac..2858d257c6 100644 --- a/agent/monitoring/emitter.py +++ b/agent/monitoring/emitter.py @@ -14,6 +14,7 @@ import logging import queue import threading import time +from contextlib import suppress from typing import Any, Dict, Optional logger = logging.getLogger(__name__) @@ -67,9 +68,7 @@ class MonitoringEmitter: with self._lock: if self._started: return - self._thread = threading.Thread( - target=self._run, name="hermes-monitoring-dispatch", daemon=True - ) + self._thread = threading.Thread(target=self._run, name="hermes-monitoring-dispatch", daemon=True) self._thread.start() self._started = True @@ -106,10 +105,8 @@ class MonitoringEmitter: self._enabled = True def unsubscribe(self, callback) -> None: - try: + with suppress(ValueError): self._subscribers.remove(callback) - except ValueError: - pass if not self._subscribers: self._enabled = False @@ -118,19 +115,13 @@ class MonitoringEmitter: """Wait boundedly for queued and in-flight batches to finish dispatch.""" if timeout <= 0: return - finished = threading.Event() def _wait_for_completion() -> None: self._q.join() finished.set() - waiter = threading.Thread( - target=_wait_for_completion, - name="hermes-monitoring-flush", - daemon=True, - ) - waiter.start() + threading.Thread(target=_wait_for_completion, name="hermes-monitoring-flush", daemon=True).start() finished.wait(timeout=timeout) def stats(self) -> Dict[str, int]: @@ -175,16 +166,9 @@ def reset_emitter_for_tests(emitter: Optional[MonitoringEmitter] = None) -> None global _EMITTER with _EMITTER_LOCK: if _EMITTER is not None and emitter is not _EMITTER: - try: + with suppress(Exception): _EMITTER.close() - except Exception: - pass _EMITTER = emitter -__all__ = [ - "MonitoringEmitter", - "get_emitter", - "emit", - "reset_emitter_for_tests", -] +__all__ = ["MonitoringEmitter", "get_emitter", "emit", "reset_emitter_for_tests"] diff --git a/agent/monitoring/events.py b/agent/monitoring/events.py index ee51d7ba28..a9bb57dcb5 100644 --- a/agent/monitoring/events.py +++ b/agent/monitoring/events.py @@ -81,8 +81,4 @@ class CronExecutionEvent(_MonitoringEvent): ts_ns: int = field(default_factory=time.time_ns) -__all__ = [ - "GatewayHealthEvent", - "GatewayDiagnosticEvent", - "CronExecutionEvent", -] +__all__ = ["GatewayHealthEvent", "GatewayDiagnosticEvent", "CronExecutionEvent"] diff --git a/agent/monitoring/gateway_health.py b/agent/monitoring/gateway_health.py index f3ac45f8f3..ea7ca161ac 100644 --- a/agent/monitoring/gateway_health.py +++ b/agent/monitoring/gateway_health.py @@ -75,9 +75,7 @@ def classify_gateway_error(raw: Any) -> str: return next((label for match, label in _GATEWAY_ERROR_RULES if match(s)), "unknown") -def classify_exit_reason( - raw: Any, *, state: Any, restart_requested: bool -) -> Optional[str]: +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" @@ -142,9 +140,7 @@ def _parse_active_agents(raw: Any) -> int: 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 - ) + 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) @@ -179,6 +175,11 @@ def _platforms_of(runtime: Optional[dict[str, Any]]) -> dict[str, Any]: return raw if isinstance(raw, dict) else {} +def _platform_error_code(pdata: dict[str, Any]) -> str: + # classify_* is idempotent on its own labels, so error_class == error_code downstream. + return classify_gateway_error(pdata.get("error_code") or pdata.get("error_message")) + + def build_gateway_health_snapshot( runtime: Optional[dict[str, Any]], *, @@ -211,8 +212,7 @@ def build_gateway_health_snapshot( for platform, pdata in platforms.items(): pdata = pdata if isinstance(pdata, dict) else {} state = _bounded_state(pdata.get("state"), allowed=_KNOWN_PLATFORM_STATES) - # classify_* is idempotent on its own labels, so error_class == error_code here. - error_code = classify_gateway_error(pdata.get("error_code") or pdata.get("error_message")) + error_code = _platform_error_code(pdata) is_degraded = state in _FATAL_PLATFORM_STATES if is_degraded: fatal_count += 1 @@ -223,13 +223,8 @@ def build_gateway_health_snapshot( )) if is_degraded: events.append(GatewayDiagnosticEvent( - name="platform.fatal", - subsystem=f"platform.{platform}", - platform=str(platform), - error_code=error_code, - error_class=error_code, - profile=profile, - version=version, + name="platform.fatal", subsystem=f"platform.{platform}", platform=str(platform), + error_code=error_code, error_class=error_code, profile=profile, version=version, severity="error" if state == "fatal" else "warning", )) @@ -282,9 +277,7 @@ def _lifecycle_events( 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 - ), + 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, @@ -296,13 +289,8 @@ def _lifecycle_events( if new_state == "startup_failed": error_class = classify_gateway_error(current.get("exit_reason") or "startup_failed") out.append(GatewayDiagnosticEvent( - name="gateway.startup_failed", - subsystem="gateway", - error_class=error_class, - error_code=error_class, - profile=profile, - version=version, - severity="error", + name="gateway.startup_failed", subsystem="gateway", error_class=error_class, error_code=error_class, + profile=profile, version=version, severity="error", )) if new_state == "stopped": out.append(health("gateway.exit")) @@ -323,30 +311,20 @@ def _platform_events( new_state = _optional_state(pdata.get("state"), allowed=_KNOWN_PLATFORM_STATES) if old_state == new_state or not new_state: continue - error_code = classify_gateway_error(pdata.get("error_code") or pdata.get("error_message")) + error_code = _platform_error_code(pdata) common: dict[str, Any] = dict( - subsystem=f"platform.{platform}", - platform=str(platform), - error_code=error_code, - error_class=error_code, - profile=profile, - version=version, - severity="error" if new_state in {"fatal", "failed", "error"} else "warning", + subsystem=f"platform.{platform}", platform=str(platform), error_code=error_code, error_class=error_code, + profile=profile, version=version, severity="error" if new_state in {"fatal", "failed", "error"} else "warning", ) - out.append(GatewayDiagnosticEvent( - name="platform.state_change", old_state=old_state, new_state=new_state, **common - )) + out.append(GatewayDiagnosticEvent(name="platform.state_change", old_state=old_state, new_state=new_state, **common)) if new_state in _FATAL_PLATFORM_STATES: out.append(GatewayDiagnosticEvent(name="platform.fatal", **common)) return out def emit_runtime_status_transition(previous: Optional[dict[str, Any]], current: dict[str, Any]) -> None: - """Emit immediate content-free gateway events for runtime status changes. - - Called by gateway.status.write_runtime_status after persisting the new status. - Fully fail-open: failures never affect gateway status writes. - """ + """Emit immediate content-free gateway events for runtime status changes. Called by + gateway.status.write_runtime_status after persisting; fully fail-open.""" try: ctx = dict(profile=_safe_profile(), version=_safe_version()) for ev in _lifecycle_events(previous, current, **ctx) + _platform_events(previous, current, **ctx): @@ -379,7 +357,7 @@ class GatewayDiagnosticLogHandler(logging.Handler): return subsystem = subsystem_for_logger(record.name) error_class = classify_gateway_error(record.getMessage()) - event = GatewayDiagnosticEvent( + emitter.get_emitter().emit(GatewayDiagnosticEvent( name=f"gateway.log.{record.levelname.lower()}", subsystem=subsystem, source_logger=source_logger_for_export(record.name), @@ -389,17 +367,12 @@ class GatewayDiagnosticLogHandler(logging.Handler): profile=self.profile, version=self.version, severity=record.levelname.lower(), - ) - emitter.get_emitter().emit(event) + )) except Exception: logger.debug("gateway diagnostic emit failed", exc_info=True) __all__ = [ - "GatewayMetric", - "GatewayHealthSnapshot", - "GatewayDiagnosticLogHandler", - "build_gateway_health_snapshot", - "classify_gateway_error", - "source_logger_for_export", + "GatewayMetric", "GatewayHealthSnapshot", "GatewayDiagnosticLogHandler", + "build_gateway_health_snapshot", "classify_gateway_error", "source_logger_for_export", ] diff --git a/agent/monitoring/gateway_health_export.py b/agent/monitoring/gateway_health_export.py index fa7d1daee1..81f4bfd95d 100644 --- a/agent/monitoring/gateway_health_export.py +++ b/agent/monitoring/gateway_health_export.py @@ -1,9 +1,8 @@ """Gateway Health & Diagnostics OTLP export runtime. -Emits operator-owned gateway service-health metrics plus narrow redacted -diagnostic events. Deliberately in-process and fail-open so it works under -systemd, launchd, s6, containers, tmux, nohup, or a plain shell without a -sidecar/watchdog dependency. +Emits operator-owned gateway service-health metrics plus narrow redacted diagnostic events. +Deliberately in-process and fail-open so it works under systemd, launchd, s6, containers, +tmux, nohup, or a plain shell without a sidecar/watchdog dependency. """ from __future__ import annotations @@ -12,6 +11,7 @@ import importlib import logging import os import threading +from contextlib import suppress from dataclasses import dataclass from typing import Any, Callable, Dict, Optional @@ -47,21 +47,11 @@ _METRICS_SDK = ( ) # Every gauge the runtime snapshot can emit MUST be listed here or it is silently dropped. _OBSERVABLE_METRIC_NAMES = ( - "hermes.gateway.up", - "hermes.gateway.state", - "hermes.gateway.active_agents", - "hermes.gateway.busy", - "hermes.gateway.drainable", - "hermes.gateway.restart_requested", - "hermes.gateway.background_work", - "hermes.gateway.background_delegations", - "hermes.platform.up", - "hermes.platform.degraded", - "hermes.cron.scheduler.heartbeat_age_seconds", - "hermes.cron.scheduler.last_success_age_seconds", - "hermes.cron.scheduler.catch_up_occurrences", - "hermes.cron.jobs.enabled", - "hermes.cron.jobs.running", + "hermes.gateway.up", "hermes.gateway.state", "hermes.gateway.active_agents", "hermes.gateway.busy", + "hermes.gateway.drainable", "hermes.gateway.restart_requested", "hermes.gateway.background_work", + "hermes.gateway.background_delegations", "hermes.platform.up", "hermes.platform.degraded", + "hermes.cron.scheduler.heartbeat_age_seconds", "hermes.cron.scheduler.last_success_age_seconds", + "hermes.cron.scheduler.catch_up_occurrences", "hermes.cron.jobs.enabled", "hermes.cron.jobs.running", "hermes.cron.jobs.overdue", ) @@ -88,21 +78,17 @@ class GatewayHealthExportRuntime: if self.thread is not None: self.thread.join(timeout=0.25) if self.log_handler is not None: - try: + with suppress(Exception): logging.getLogger().removeHandler(self.log_handler) - except Exception: - pass # 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] - try: + with suppress(Exception): bus = emitter.get_emitter() bus.flush(timeout=1.0) for sub in subscribers: bus.unsubscribe(sub) - except Exception: - pass # Network flush/close runs under one bounded daemon-thread deadline so it can # never delay gateway teardown indefinitely. @@ -110,15 +96,11 @@ class GatewayHealthExportRuntime: def _close() -> None: for item in closeables: - try: + with suppress(Exception): item.shutdown() - except Exception: - pass if closeables: - worker = threading.Thread( - target=_close, name="hermes-gateway-health-export-shutdown", daemon=True - ) + worker = threading.Thread(target=_close, name="hermes-gateway-health-export-shutdown", daemon=True) worker.start() worker.join(timeout=2.0) @@ -172,11 +154,7 @@ def _read_gateway_snapshot(config: Dict[str, Any]): except Exception: runtime = {} return build_gateway_health_snapshot( - runtime, - gateway_running=True, - profile=_profile(), - install_id=_install_id(config), - version=_version(), + runtime, gateway_running=True, profile=_profile(), install_id=_install_id(config), version=_version(), supervision_mode=_supervision_mode(), ) @@ -197,27 +175,22 @@ def _count(failure_msg: str, module: str, read: Callable[[Any], Any]) -> int: def _read_background_work_count() -> int: - """Live background/subagent work that ``active_agents`` deliberately does NOT include. - - ``active_agents`` counts foreground turns + in-flight cron + API runs; backgrounded - ``delegate_task`` subagents, ``terminal(background=true)`` processes and kanban workers - are tracked only by the scale-to-zero guard, so without this a peer churning through - subagents shows ``active_agents=0``. TASK-granular: a fan-out batch of N contributes N - (real concurrent load), unlike the pool's one-slot-per-batch accounting. Content-free. - """ + """Live background/subagent work that ``active_agents`` deliberately does NOT include + (``active_agents`` = foreground turns + in-flight cron + API runs; backgrounded + ``delegate_task`` subagents, ``terminal(background=true)`` processes and kanban workers are + tracked only by the scale-to-zero guard). TASK-granular: a fan-out batch of N contributes N + (real concurrent load), unlike the pool's one-slot-per-batch accounting. Content-free.""" return ( - _count("background-work async-delegation count failed", "tools.async_delegation", - lambda m: m.active_task_count()) + _count("background-work async-delegation count failed", "tools.async_delegation", lambda m: m.active_task_count()) + _count("background-work process-registry count failed", "tools.process_registry", lambda m: m.process_registry.count_running()) ) def _read_background_delegations_count() -> int: - """Live async delegation UNITS (dispatch/pool slots): a batch counts ONE regardless of - fan-out width, matching the pool's capacity accounting — so operators can see slot - pressure (alert vs ``max_concurrent_children``) alongside ``background_work``'s real load. - Delegations only; terminal/kanban work is already folded into ``background_work``.""" + """Live async delegation UNITS (dispatch/pool slots): a batch counts ONE regardless of fan-out + width, matching the pool's capacity accounting — slot pressure (alert vs + ``max_concurrent_children``) alongside ``background_work``'s real load. Delegations only.""" return _count("background-delegations count failed", "tools.async_delegation", lambda m: m.active_count()) @@ -233,20 +206,14 @@ def _read_runtime_snapshot(config: Dict[str, Any]): ): gateway_snapshot.metrics.append(GatewayMetric(name=name, value=read(), attributes=base)) except Exception as exc: - logger.warning( - "background-work snapshot unavailable; metric not exported (error_type=%s)", - type(exc).__name__, - ) + logger.warning("background-work snapshot unavailable; metric not exported (error_type=%s)", type(exc).__name__) logger.debug("background-work snapshot traceback", exc_info=True) try: cron_snapshot = _read_cron_snapshot() except Exception as exc: # Cron telemetry silently dropping out is a release-relevant regression: WARN with only # the exception *type* (the message could carry paths); exc_info stays on DEBUG. - logger.warning( - "cron health snapshot unavailable; cron telemetry not exported (error_type=%s)", - type(exc).__name__, - ) + logger.warning("cron health snapshot unavailable; cron telemetry not exported (error_type=%s)", type(exc).__name__) logger.debug("cron health snapshot traceback", exc_info=True) return gateway_snapshot gateway_snapshot.metrics.extend(cron_snapshot.metrics) @@ -268,9 +235,7 @@ def _start_metric_provider(config: Dict[str, Any], sdk: Dict[str, Any]) -> Any: exporter = sdk["OTLPMetricExporter"](**_exporter_kwargs(config, "metrics")) interval_ms = max(5, int(gh.get("export_interval_seconds", 60))) * 1000 reader = sdk["PeriodicExportingMetricReader"](exporter, export_interval_millis=interval_ms) - provider = sdk["MeterProvider"]( - metric_readers=[reader], resource=_resource(config, sdk, "gateway_health") - ) + provider = sdk["MeterProvider"](metric_readers=[reader], resource=_resource(config, sdk, "gateway_health")) meter = provider.get_meter("hermes.gateway.health") Observation = sdk["Observation"] @@ -290,12 +255,7 @@ def _start_metric_provider(config: Dict[str, Any], sdk: Dict[str, Any]) -> Any: _SEVERITY_NAMES = { - "critical": "FATAL", - "fatal": "FATAL", - "error": "ERROR", - "info": "INFO", - "information": "INFO", - "debug": "DEBUG", + "critical": "FATAL", "fatal": "FATAL", "error": "ERROR", "info": "INFO", "information": "INFO", "debug": "DEBUG", } @@ -309,9 +269,7 @@ class GatewayDiagnosticLogStreamer(EmitterStreamer): def __init__(self, config: Dict[str, Any], sdk: Dict[str, Any]): self._provider = sdk["LoggerProvider"](resource=_resource(config, sdk, "gateway_diagnostics")) - self._processor = sdk["BatchLogRecordProcessor"]( - sdk["OTLPLogExporter"](**_exporter_kwargs(config, "logs")) - ) + self._processor = sdk["BatchLogRecordProcessor"](sdk["OTLPLogExporter"](**_exporter_kwargs(config, "logs"))) self._provider.add_log_record_processor(self._processor) self._logger = self._provider.get_logger(_DEFAULT_DIAGNOSTIC_SCOPE) self._sdk = sdk @@ -382,8 +340,7 @@ def start_gateway_health_export(config: Dict[str, Any]) -> GatewayHealthExportRu sdk = _require_metrics_sdk(prompt=False) except Exception: logger.warning( - "monitoring.gateway_health_export.enabled but OTLP SDK is unavailable; " - "install 'hermes-agent[otlp]'", + "monitoring.gateway_health_export.enabled but OTLP SDK is unavailable; install 'hermes-agent[otlp]'", exc_info=True, ) return GatewayHealthExportRuntime(enabled=False, reason="otlp_unavailable") @@ -423,7 +380,4 @@ def start_gateway_health_export(config: Dict[str, Any]) -> GatewayHealthExportRu return runtime -__all__ = [ - "GatewayHealthExportRuntime", - "start_gateway_health_export", -] +__all__ = ["GatewayHealthExportRuntime", "start_gateway_health_export"] diff --git a/agent/monitoring/otlp_exporter.py b/agent/monitoring/otlp_exporter.py index f1c12fa4dc..0c2d8e102c 100644 --- a/agent/monitoring/otlp_exporter.py +++ b/agent/monitoring/otlp_exporter.py @@ -1,18 +1,12 @@ """Export monitoring events to an OpenTelemetry Collector over OTLP/HTTP. -Maps gateway monitoring events to OTel spans and sends them to the operator -configured ``monitoring.export.otlp`` endpoint (no default destination ships). -Also hosts the OTLP plumbing shared with ``gateway_health_export`` (SDK loading, -header resolution, resource attributes, endpoint mapping). - -* The OTel SDK is an optional extra (``pip install hermes-agent[otlp]``), - imported lazily so it is only required when export is actually used. -* ``headers_env`` maps a header name to an environment variable name; values - are read from the environment at export time and never logged or stored. -* The continuous subscriber runs in the emitter's dispatcher thread and is - fail-isolated, so an export error cannot affect the gateway. The - ``event_filter`` seam keeps future planes sharing the emitter from silently - riding along on this exporter. +Maps gateway monitoring events to OTel spans for the operator-configured +``monitoring.export.otlp`` endpoint (no default destination ships) and hosts the OTLP +plumbing shared with ``gateway_health_export`` (SDK loading, header resolution, resource +attributes, endpoint mapping). The OTel SDK is an optional extra (``hermes-agent[otlp]``) +imported lazily; ``headers_env`` values are read from the environment at export time and +never logged or stored. The continuous subscriber runs on the emitter's dispatcher thread, +fail-isolated, and ``event_filter`` keeps other planes from riding along on this exporter. """ from __future__ import annotations @@ -21,6 +15,7 @@ import importlib import logging import os import re +from contextlib import suppress from typing import Any, Callable, Dict, Iterable, List, Optional from agent.monitoring.gateway_health import _safe_instance_id @@ -56,27 +51,19 @@ _SDK_SYMBOLS: Dict[str, str] = { _SPAN_SDK = ("TracerProvider", "BatchSpanProcessor", "Resource", "OTLPSpanExporter", "SpanKind") -def _require_sdk( - names: Iterable[str] = _SPAN_SDK, *, auto_install: bool = True, prompt: bool = True -) -> Dict[str, Any]: +def _require_sdk(names: Iterable[str] = _SPAN_SDK, *, auto_install: bool = True, prompt: bool = True) -> Dict[str, Any]: """Import the named OTel SDK symbols, lazily installing the extra on first use. - Routes through tools.lazy_deps (feature 'export.otlp') — gated by - security.allow_lazy_installs and TTY-prompted (``prompt=False`` from - non-interactive contexts). Any lazy-install failure falls through to the - import attempt, which raises OTLPUnavailable with a manual install hint. + Routes through tools.lazy_deps (feature 'export.otlp') — gated by security.allow_lazy_installs + and TTY-prompted (``prompt=False`` from non-interactive contexts). Any lazy-install failure + falls through to the import attempt, which raises OTLPUnavailable with a manual install hint. """ if auto_install: - try: + with suppress(Exception): from tools.lazy_deps import ensure as _lazy_ensure _lazy_ensure("export.otlp", prompt=prompt) - except Exception: - pass try: - return { - name: getattr(importlib.import_module(_SDK_SYMBOLS[name]), name) - for name in names - } + return {name: getattr(importlib.import_module(_SDK_SYMBOLS[name]), name) for name in names} except Exception as e: # ImportError or partial install raise OTLPUnavailable( "OTLP export requires the optional dependency. Install with:\n" @@ -122,15 +109,8 @@ def _signal_endpoint(endpoint: str, signal: str) -> str: _RESOURCE_ATTRIBUTE_KEYS = frozenset({ - "service.name", - "service.namespace", - "service.version", - "service.instance.id", - "deployment.environment.name", - "cloud.provider", - "cloud.platform", - "cloud.region", - "telemetry.scope", + "service.name", "service.namespace", "service.version", "service.instance.id", + "deployment.environment.name", "cloud.provider", "cloud.platform", "cloud.region", "telemetry.scope", }) _SAFE_RESOURCE_VALUE = re.compile(r"^[A-Za-z0-9._:/-]{1,128}$") @@ -158,12 +138,9 @@ def _safe_resource_attributes(raw: Any) -> Dict[str, str]: return attrs -def _runtime_resource_attributes( - config: Dict[str, Any], *, telemetry_scope: str -) -> Dict[str, str]: +def _runtime_resource_attributes(config: Dict[str, Any], *, telemetry_scope: str) -> Dict[str, str]: """Build the safe OTLP resource shared by spans, metrics and diagnostic logs.""" - gh = _monitoring_section(config, "gateway_health_export") - attrs = _safe_resource_attributes(gh.get("resource_attributes")) + attrs = _safe_resource_attributes(_monitoring_section(config, "gateway_health_export").get("resource_attributes")) attrs["service.name"] = "hermes-gateway" attrs["service.instance.id"] = _safe_instance_id(_install_id(config)) attrs["telemetry.scope"] = telemetry_scope @@ -177,8 +154,7 @@ def build_exporter(config: Dict[str, Any]): endpoint = otlp.get("endpoint") if not endpoint: raise ValueError("monitoring.export.otlp.endpoint is not set") - headers = _resolve_headers(otlp.get("headers_env")) - return sdk["OTLPSpanExporter"](endpoint=endpoint, headers=headers or None) + return sdk["OTLPSpanExporter"](endpoint=endpoint, headers=_resolve_headers(otlp.get("headers_env")) or None) def _resource_attributes(config: Dict[str, Any]) -> Dict[str, str]: @@ -187,8 +163,7 @@ def _resource_attributes(config: Dict[str, Any]) -> Dict[str, str]: def _make_provider(config: Dict[str, Any]): sdk = _require_sdk() - resource = sdk["Resource"].create(_resource_attributes(config)) - provider = sdk["TracerProvider"](resource=resource) + provider = sdk["TracerProvider"](resource=sdk["Resource"].create(_resource_attributes(config))) processor = sdk["BatchSpanProcessor"](build_exporter(config)) provider.add_span_processor(processor) return provider, processor @@ -234,9 +209,7 @@ def export_batch(provider, batch: List[Dict[str, Any]]) -> int: n = 0 for ev in batch: try: - name = f"hermes.{ev.get('event', 'event')}" - span = tracer.start_span(name, attributes=_span_attrs(ev)) - span.end() + tracer.start_span(f"hermes.{ev.get('event', 'event')}", attributes=_span_attrs(ev)).end() n += 1 except Exception: logger.debug("OTLP span map failed", exc_info=True) @@ -246,36 +219,25 @@ def export_batch(provider, batch: List[Dict[str, Any]]) -> int: # ── continuous streaming subscribers ───────────────────────────────────────── class EmitterStreamer: """Base for emitter subscribers owning an OTel provider + batch processor. - - Register with ``emitter.subscribe(streamer)``. Fail-isolated by the emitter. - """ + Register with ``emitter.subscribe(streamer)``. Fail-isolated by the emitter.""" _provider: Any _processor: Any exported: int = 0 def shutdown(self) -> None: - try: + with suppress(Exception): from agent.monitoring.emitter import get_emitter get_emitter().unsubscribe(self) - except Exception: - pass - try: + with suppress(Exception): self._processor.force_flush() self._provider.shutdown() - except Exception: - pass class OTLPStreamer(EmitterStreamer): """A live subscriber that pushes each emitter batch to OTLP as spans.""" - def __init__( - self, - config: Dict[str, Any], - *, - event_filter: Optional[Callable[[Dict[str, Any]], bool]] = None, - ): + def __init__(self, config: Dict[str, Any], *, event_filter: Optional[Callable[[Dict[str, Any]], bool]] = None): self._provider, self._processor = _make_provider(config) self._event_filter = event_filter self.exported = 0 @@ -303,16 +265,13 @@ def is_enabled(config: Dict[str, Any]) -> bool: def start_streaming( - config: Dict[str, Any], - *, - event_filter: Optional[Callable[[Dict[str, Any]], bool]] = None, + config: Dict[str, Any], *, event_filter: Optional[Callable[[Dict[str, Any]], bool]] = None, ) -> Optional[OTLPStreamer]: """If OTLP is enabled, attach a streamer to the singleton emitter. - ``event_filter`` scopes the exporter to its plane. Startup is non-interactive: - a configured-but-missing SDK is lazily installed once (prompt=False, gated by - security.allow_lazy_installs); if it still can't load, log and no-op — never - raise into startup. + ``event_filter`` scopes the exporter to its plane. Startup is non-interactive: a + configured-but-missing SDK is lazily installed once (prompt=False, gated by + security.allow_lazy_installs); if it still can't load, log and no-op — never raise into startup. """ if not is_enabled(config): return None @@ -329,11 +288,5 @@ def start_streaming( __all__ = [ - "OTLPUnavailable", - "OTLPStreamer", - "build_exporter", - "export_batch", - "is_available", - "is_enabled", - "start_streaming", + "OTLPUnavailable", "OTLPStreamer", "build_exporter", "export_batch", "is_available", "is_enabled", "start_streaming", ] diff --git a/agent/monitoring/policy.py b/agent/monitoring/policy.py index 0939fbbc67..a50f397c7a 100644 --- a/agent/monitoring/policy.py +++ b/agent/monitoring/policy.py @@ -46,6 +46,4 @@ def ensure_install_id(config: Dict[str, Any]) -> str: return minted -__all__ = [ - "ensure_install_id", -] +__all__ = ["ensure_install_id"] diff --git a/agent/monitoring/redaction.py b/agent/monitoring/redaction.py index 22ae1c0904..f716c00f08 100644 --- a/agent/monitoring/redaction.py +++ b/agent/monitoring/redaction.py @@ -1,11 +1,9 @@ """Redaction applied to monitoring data before egress. -One unconditional scrub, no modes, no knobs. Every string that leaves the -process passes through ``redact_for_export``: secrets first (wraps -``agent/redact.py::redact_sensitive_text(force=True)`` plus bearer/token -shapes, failing CLOSED so a broken redactor never emits the raw string), then +One unconditional scrub, no modes, no knobs. Every string that leaves the process passes +through ``redact_for_export``: secrets first (``agent/redact.py::redact_sensitive_text(force=True)`` +plus bearer/token shapes, failing CLOSED so a broken redactor never emits the raw string), then PII (e-mail, phone, UUID-shaped ids -> ``[email]`` / ``[phone]`` / ``[id]``). -There is deliberately no setting to weaken this. """ from __future__ import annotations @@ -15,9 +13,7 @@ from typing import Any, Optional # ── secret shapes (belt-and-suspenders on top of agent/redact.py) ─────────── _BEARER_RE = re.compile(r"\bBearer\s+[A-Za-z0-9._~+\-/]+=*", re.IGNORECASE) -_TOKEN_RE = re.compile( - r"\b(xox[baprs]-[A-Za-z0-9-]+|sk-[A-Za-z0-9_-]{8,}|gh[pousr]_[A-Za-z0-9_]{8,})\b" -) +_TOKEN_RE = re.compile(r"\b(xox[baprs]-[A-Za-z0-9-]+|sk-[A-Za-z0-9_-]{8,}|gh[pousr]_[A-Za-z0-9_]{8,})\b") _SECRET_LITERAL_RE = re.compile(r"\*{3,}") _BEARER_RESIDUE_RE = re.compile(r"\bBearer\s+\[[^\]]+\]", re.IGNORECASE) @@ -27,9 +23,7 @@ _EMAIL_RE = re.compile(r"[A-Za-z0-9._%+\-]+@[A-Za-z0-9.\-]+\.[A-Za-z]{2,}") _PHONE_RE = re.compile( r"(? Optional[str]: return out -def redact_bounded( - raw: Any, - *, - limit: int = 500, - empty: str = "[redacted]", - unavailable: str = UNAVAILABLE, -) -> str: +def redact_bounded(raw: Any, *, limit: int = 500, empty: str = "[redacted]", unavailable: str = UNAVAILABLE) -> str: """Redact ``str(raw or "")`` and length-bound it; ``empty`` replaces an empty result, ``unavailable`` is returned if redaction itself raises.""" try: @@ -73,7 +61,4 @@ def redact_bounded( return unavailable -__all__ = [ - "redact_for_export", - "redact_bounded", -] +__all__ = ["redact_for_export", "redact_bounded"]