refactor(agent): compact monitoring/* and lsp/workspace — contextlib.suppress, shared platform error-code helper, nearest_root marker probe, docstring compaction

This commit is contained in:
Teknium
2026-09-02 18:32:02 -07:00
parent 1a1716ee64
commit 052582de2d
10 changed files with 144 additions and 357 deletions

View File

@@ -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",
]

View File

@@ -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"]

View File

@@ -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",
]

View File

@@ -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"]

View File

@@ -81,8 +81,4 @@ class CronExecutionEvent(_MonitoringEvent):
ts_ns: int = field(default_factory=time.time_ns)
__all__ = [
"GatewayHealthEvent",
"GatewayDiagnosticEvent",
"CronExecutionEvent",
]
__all__ = ["GatewayHealthEvent", "GatewayDiagnosticEvent", "CronExecutionEvent"]

View File

@@ -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",
]

View File

@@ -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"]

View File

@@ -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",
]

View File

@@ -46,6 +46,4 @@ def ensure_install_id(config: Dict[str, Any]) -> str:
return minted
__all__ = [
"ensure_install_id",
]
__all__ = ["ensure_install_id"]

View File

@@ -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"(?<!\w)(?:\+?\d{1,3}[\s.\-]?)?(?:\(\d{2,4}\)[\s.\-]?)?\d{3}[\s.\-]?\d{3,4}(?:[\s.\-]?\d{2,4})?(?!\w)"
)
_UUID_RE = re.compile(
r"\b[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}\b"
)
_UUID_RE = re.compile(r"\b[0-9a-fA-F]{8}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{4}-[0-9a-fA-F]{12}\b")
UNAVAILABLE = "[redaction-unavailable]"
@@ -58,13 +52,7 @@ def redact_for_export(text: Optional[str]) -> 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"]