Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
Show all changes
31 commits
Select commit Hold shift + click to select a range
3e28eac
feat(telemetry): local-first telemetry & observability
jquesnelle Jun 24, 2026
e829d03
refactor(telemetry): drop policy.resolve(); read config directly
jquesnelle Jun 26, 2026
72b7af7
refactor(telemetry): drop "plane" terminology
jquesnelle Jun 26, 2026
0d0549d
test(telemetry): end-to-end plugin dispatch coverage
jquesnelle Jun 26, 2026
24c86e1
fix(telemetry): aggregate requires local telemetry to be on
jquesnelle Jun 26, 2026
81a3415
docs(telemetry): clarify reserved subagent-lineage hooks
jquesnelle Jun 27, 2026
7c80b79
feat(telemetry): write tel_spans — reconstructable run -> calls trace
jquesnelle Jun 27, 2026
edac3a7
refactor(telemetry): cut dead schema; tests assert what's actually wr…
jquesnelle Jun 27, 2026
fdb5ae0
docs(telemetry): align observability docs with the trimmed schema
jquesnelle Jun 27, 2026
505d12f
refactor(monitoring): scope telemetry substrate to gateway health/dia…
victor-kyriazakos Jul 14, 2026
87a1573
fix(monitoring): address review — persist install_id, drop dead confi…
victor-kyriazakos Jul 14, 2026
16e774f
fix(monitoring): harden gateway health OTLP egress
Jul 22, 2026
98a0260
fix(otel): identify Hermes resources safely
Jul 22, 2026
73190cd
fix(monitoring): complete production OTLP runtime
Jul 22, 2026
4dc4f43
chore(monitoring): align smoke and attribution metadata
Jul 22, 2026
18b4664
fix(monitoring): drain terminal lifecycle events
Jul 22, 2026
42bc6c8
fix(monitoring): keep diagnostics content-free
Jul 22, 2026
9b098e7
fix(monitoring): enrich gateway diagnostic scope
Jul 23, 2026
efbf2fd
feat(monitoring): add cron operational telemetry
Jul 24, 2026
a65a647
fix(monitoring): correct cron operational signals
Jul 24, 2026
2a46acc
docs(monitoring): clarify bounded cron flush
Jul 24, 2026
2e2c77c
fix(monitoring): normalize cron delivery outcomes
Jul 24, 2026
58a7491
docs(monitoring): add fleet operations guide
Jul 24, 2026
6de6e36
fix(monitoring): surface cron snapshot failures at WARNING (content-f…
Jul 24, 2026
9edcb7b
feat(monitoring): emit hermes.gateway.background_work (subagent/bg jobs)
Jul 25, 2026
b49b78b
docs(monitoring): background_work signal + plane maintenance/extensio…
Jul 25, 2026
3c567f7
test(monitoring): assert background_work in snapshot + metric_names r…
Jul 25, 2026
0ef0816
feat(monitoring): count background_work task-granular (expand delegat…
Jul 25, 2026
cc80dba
feat(monitoring): add hermes.gateway.background_delegations (unit/slo…
Jul 25, 2026
6a174e9
Merge origin/main into feat/gateway-health-diagnostics
victor-kyriazakos Jul 28, 2026
1773752
Merge remote-tracking branch 'origin/main' into feat/gateway-health-d…
victor-kyriazakos Jul 29, 2026
File filter

Filter by extension

Filter by extension


Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
8 changes: 6 additions & 2 deletions Dockerfile
Original file line number Diff line number Diff line change
Expand Up @@ -227,7 +227,7 @@ RUN cd plugins/platforms/photon/sidecar && \
# frontend stats the readme path during dep resolution, so we `touch` an
# empty placeholder — the real README is restored by `COPY . .` below.
#
# `uv sync --frozen --no-install-project --extra all --extra messaging`
# `uv sync --frozen --no-install-project --extra all --extra messaging --extra otlp`
# installs the deps reachable through the composite `[all]` extra
# (handpicked set intended for the production image — excludes `[dev]`),
# plus gateway messaging adapters that should work in the published image
Expand All @@ -240,6 +240,10 @@ RUN cd plugins/platforms/photon/sidecar && \
# so Docker users can use these providers without requiring runtime
# lazy-install access to PyPI (often blocked in containerized envs).
#
# The [otlp] extra contains the SDK/exporter imported by Hermes when Gateway
# Health export is enabled. Collector and observability-backend dependencies
# remain external and are not part of the Hermes production image.
#
# The hindsight memory provider's client (hindsight-client) is baked in
# for the same reason: it lazy-installs into /opt/hermes/.venv at first
# use, which lives inside the (immutable) image layer rather than the
Expand All @@ -257,7 +261,7 @@ RUN cd plugins/platforms/photon/sidecar && \
# The editable link is created after the source copy below.
COPY pyproject.toml uv.lock ./
RUN touch ./README.md
RUN uv sync --frozen --no-install-project --extra all --extra messaging --extra anthropic --extra bedrock --extra azure-identity --extra hindsight --extra matrix
RUN uv sync --frozen --no-install-project --extra all --extra messaging --extra otlp --extra anthropic --extra bedrock --extra azure-identity --extra hindsight --extra matrix

# ---------- Frontend build (cached independently from Python source) ----------
# Copy only the frontend source trees first so that Python-only changes don't
Expand Down
29 changes: 29 additions & 0 deletions agent/monitoring/__init__.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,29 @@
"""Hermes gateway monitoring.

Service health monitoring plus redacted operational diagnostics for the
gateway daemon, exported over OTLP to an operator-configured endpoint.

``emitter`` is the in-process event bus: producers (gateway status hooks,
the diagnostic log handler) hand typed events to a fire-and-forget queue,
and subscribers (the OTLP streamers) consume them off the hot path. The
emitter never blocks or raises into gateway code (the hot-path invariant),
and nothing is persisted locally — monitoring is an egress path, not a store.

Deliberately out of scope here: run/model/tool trajectory capture, usage
analytics, and any content-bearing signal. Those planes are served by the
NeMo Relay integration and its Hermes-owned subscribers.
"""

from __future__ import annotations

from . import emitter, events

emit = emitter.emit
get_emitter = emitter.get_emitter

__all__ = [
"emitter",
"events",
"emit",
"get_emitter",
]
201 changes: 201 additions & 0 deletions agent/monitoring/cron_health.py
Original file line number Diff line number Diff line change
@@ -0,0 +1,201 @@
"""Content-free cron service-health and execution telemetry projection."""

from __future__ import annotations

import hashlib
import logging
import re
from dataclasses import dataclass
from datetime import datetime
from typing import Any, Optional

from agent.monitoring.events import CronExecutionEvent
from agent.monitoring.gateway_health import GatewayHealthSnapshot, GatewayMetric
from cron.jobs import (
_compute_grace_seconds,
get_catch_up_occurrence_count,
get_ticker_heartbeat_age,
get_ticker_success_age,
load_jobs,
)
from cron.scheduler import get_running_job_ids
from hermes_time import now as _hermes_now

logger = logging.getLogger(__name__)
_KNOWN_STATUSES = {"claimed", "running", "completed", "failed", "unknown"}
_KNOWN_SOURCES = {"builtin", "direct", "external"}
_KNOWN_DELIVERY_OUTCOMES = {"delivered", "failed", "suppressed", "not_configured"}


@dataclass(frozen=True, slots=True)
class CronHealthSnapshot:
metrics: list[GatewayMetric]
events: list[CronExecutionEvent]


def _now() -> datetime:
return _hermes_now()


def _job_key(raw: Any) -> str:
value = str(raw or "unknown").encode("utf-8", errors="replace")
return f"sha256:{hashlib.sha256(value).hexdigest()[:24]}"


def classify_cron_error(raw: Any) -> str:
text = str(raw or "").lower()
if (
re.search(r"\b(?:authentication|authenticated|authenticate|authorization|authorized|authorize|unauthorized|forbidden)\b", text)
or re.search(r"\bbearer\b", text)
or re.search(r"\b(?:access|api|refresh) token\b", text)
or re.search(r"\b(?:401|403)\b", text)
):
return "auth_failed"
if "rate limit" in text or "429" in text or "quota" in text:
return "rate_limited"
if "timeout" in text or "timed out" in text:
return "timeout"
if any(value in text for value in ("network", "connection", "dns", "socket", "unreachable")):
return "network_error"
if "dispatch" in text or "executor" in text:
return "dispatch_failed"
if "interrupt" in text or "owner exited" in text or "restarted" in text:
return "interrupted"
if "empty response" in text:
return "empty_response"
if any(value in text for value in ("config", "missing", "invalid")):
return "invalid_config"
return "unknown"


def _parse_time(raw: Any) -> Optional[datetime]:
try:
return datetime.fromisoformat(str(raw)) if raw else None
except (TypeError, ValueError):
return None


def _duration_ms(record: dict[str, Any]) -> Optional[int]:
start = _parse_time(record.get("started_at")) or _parse_time(record.get("claimed_at"))
finish = _parse_time(record.get("finished_at"))
if start is None or finish is None:
return None
try:
duration = int((finish - start).total_seconds() * 1000)
except (TypeError, ValueError):
return None
return max(0, duration)


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()
if source not in _KNOWN_SOURCES and source != "unknown":
source = "external"
outcome = str(delivery_outcome).lower() if delivery_outcome is not None else None
return CronExecutionEvent(
status=status if status in _KNOWN_STATUSES else "unknown",
job_key=_job_key(record.get("job_id")),
source=source if source in _KNOWN_SOURCES else "unknown",
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
),
)


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
try:
from agent.monitoring import emitter

event = project_execution_event(record, delivery_outcome=delivery_outcome)
target = emitter.get_emitter()
target.emit(event)
if event.status in {"completed", "failed", "unknown"}:
target.flush(timeout=1.0)
except Exception:
logger.debug("cron execution telemetry emit failed", exc_info=True)


def _is_overdue(job: dict[str, Any], now: datetime) -> bool:
if not job.get("enabled", True):
return False
next_run = _parse_time(job.get("next_run_at"))
schedule = job.get("schedule")
if next_run is None or not isinstance(schedule, dict):
return False
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)
except (TypeError, ValueError):
return False


def build_cron_health_snapshot() -> CronHealthSnapshot:
metrics: list[GatewayMetric] = []
for name, reader in (
("hermes.cron.scheduler.heartbeat_age_seconds", get_ticker_heartbeat_age),
("hermes.cron.scheduler.last_success_age_seconds", get_ticker_success_age),
):
try:
value = reader()
if value is not None:
metrics.append(GatewayMetric(name, max(0.0, float(value)), {}))
except Exception:
logger.debug("cron freshness metric unavailable", exc_info=True)

try:
metrics.append(
GatewayMetric(
"hermes.cron.scheduler.catch_up_occurrences",
get_catch_up_occurrence_count(),
{},
)
)
except Exception:
logger.debug("cron catch-up metric unavailable", exc_info=True)

try:
jobs = load_jobs()
enabled = [job for job in 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())),
{},
)
)
except Exception:
logger.debug("cron job metrics unavailable", exc_info=True)

try:
metrics.append(
GatewayMetric("hermes.cron.jobs.running", len(get_running_job_ids()), {})
)
except Exception:
logger.debug("cron running-job metric unavailable", exc_info=True)
return CronHealthSnapshot(metrics=metrics, events=[])


__all__ = [
"CronHealthSnapshot",
"build_cron_health_snapshot",
"classify_cron_error",
"emit_execution_state",
"project_execution_event",
]
Loading
Loading