diff --git a/Cargo.lock b/Cargo.lock index a3710f58..bb9d6699 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -4561,6 +4561,9 @@ dependencies = [ "hive-jobq-wire", "hive-types", "hmac 0.13.0", + "opentelemetry", + "opentelemetry-otlp", + "opentelemetry_sdk", "problem_details", "reqwest", "serde", diff --git a/swarm-controller/Cargo.toml b/swarm-controller/Cargo.toml index cb7e0419..72a6f361 100644 --- a/swarm-controller/Cargo.toml +++ b/swarm-controller/Cargo.toml @@ -38,6 +38,20 @@ problem_details = { version = "0.9.0", features = ["axum"] } # same shape `hive-c0re/src/job_queue/scheduler.rs` uses over its own graph. hive-jobq.workspace = true hive-jobq-wire.workspace = true +# OTEL SDK for the jobq-rollup metrics exporter (`jobq_metrics.rs`). Same +# versions + blocking-client rationale as `hive-c0re`/`hive-metric`: the +# metrics SDK's `PeriodicReader` runs on a background thread with no Tokio +# reactor, where the async client panics. Kept local rather than in +# `workspace.dependencies` — same call `hive-metric` already made, this is +# still the only other crate that needs it. +opentelemetry = "0.32" +opentelemetry_sdk = { version = "0.32", features = ["metrics"] } +opentelemetry-otlp = { version = "0.32", default-features = false, features = [ + "metrics", + "http-json", + "reqwest-blocking-client", + "reqwest-rustls", +] } # The forge webhook HMAC (`webhook.rs`). Kept in this crate rather than # shared with hive-c0re's equivalent: c0re's copy is scheduled to be deleted # with its webhook routes once registration moves here, so the second holder diff --git a/swarm-controller/src/jobq_metrics.rs b/swarm-controller/src/jobq_metrics.rs new file mode 100644 index 00000000..a734bba5 --- /dev/null +++ b/swarm-controller/src/jobq_metrics.rs @@ -0,0 +1,259 @@ +//! OTEL export of the swarm-level job graph's state rollup — the same +//! `hive_jobq_wire::state_rollup()` counts `GET /api/jobq/rollup` already +//! serves, ridden out to the collector on a timer instead of only on +//! request. First consumer of a generic "add jobq metrics" ask: the same +//! shape works for any `hive-jobq` instance, this crate's is just the +//! first to wire it. +//! +//! Same OTEL SDK setup as `hive-c0re::stats::otel_metrics` (container +//! stats) and `hive-metric` (the one-shot CLI): the metrics SDK's +//! `PeriodicReader` drives its export from a background thread with no +//! Tokio reactor, so the blocking OTLP client is deliberate here too, not +//! an oversight. +//! +//! Bridging async→sync: reading `state_rollup()` needs the scheduler's +//! `std::sync::Mutex`, which is fine to lock briefly from a sync callback — +//! but holding it for the whole OTLP export would block every request +//! handler that also locks it for however long the export takes. So, same +//! pattern as the container-stats exporter: an async task refreshes a +//! shared snapshot on an interval, and the (sync) SDK callbacks only ever +//! read that snapshot, never the scheduler directly. + +use std::sync::{Arc, Mutex, OnceLock}; +use std::time::Duration; + +use anyhow::{Context, Result}; +use hive_jobq_wire::StateCount; +use opentelemetry::KeyValue; +use opentelemetry::metrics::MeterProvider as _; +use opentelemetry_otlp::{MetricExporter, Protocol, WithExportConfig}; +use opentelemetry_sdk::Resource; +use opentelemetry_sdk::metrics::{PeriodicReader, SdkMeterProvider}; + +use crate::{SwarmNodeKind, SwarmResourceKind}; + +/// Default export cadence when `HYPERHIVE_OTEL_METRIC_INTERVAL_MS` is unset. +/// Matches `hive-c0re::stats::otel_metrics`'s default — no reason for the +/// two exporters to disagree on how fresh "current" means by default. +const DEFAULT_INTERVAL: Duration = Duration::from_mins(1); + +/// Latest rollup snapshot: written by the async refresher, read by the sync +/// observable-instrument callbacks. +static SNAPSHOT: OnceLock>>> = OnceLock::new(); + +/// Keep the provider alive for the process lifetime — the `PeriodicReader` +/// exports only while the provider lives. +static PROVIDER: OnceLock = OnceLock::new(); + +/// Spawn the jobq-rollup OTEL exporter if OTEL is configured +/// (`OTEL_EXPORTER_OTLP_ENDPOINT` non-empty). No-op otherwise — same +/// graceful-absence shape every other optional wiring in this daemon uses +/// (queue, bridge, forge, webhook secret). Call once at startup. +pub fn spawn_exporter( + jobq: Arc>>, +) { + let Some(endpoint) = endpoint() else { + tracing::debug!("otel jobq-metrics: no endpoint configured, exporter disabled"); + return; + }; + let snapshot = SNAPSHOT + .get_or_init(|| Arc::new(Mutex::new(Vec::new()))) + .clone(); + let interval = interval(); + + let refresh = snapshot.clone(); + tokio::spawn(async move { + loop { + let fresh = { + let sched = jobq + .lock() + .unwrap_or_else(std::sync::PoisonError::into_inner); + let graph = sched.graph(); + let roots: Vec = graph.roots().map(|n| n.id).collect(); + hive_jobq_wire::state_rollup(graph, roots) + }; + if let Ok(mut g) = refresh.lock() { + *g = fresh; + } + tokio::time::sleep(interval).await; + } + }); + + match build_provider(interval, snapshot) { + Ok(provider) => { + let _ = PROVIDER.set(provider); + tracing::info!(%endpoint, ?interval, "otel jobq-metrics: exporter enabled"); + } + Err(e) => { + tracing::warn!(error = ?e, "otel jobq-metrics: exporter init failed"); + } + } +} + +fn build_provider( + interval: Duration, + snapshot: Arc>>, +) -> Result { + // Same http/json, endpoint-from-env-only construction as + // `hive-c0re::stats::otel_metrics` — see that module's doc comment for + // why `with_endpoint` is deliberately never called here. + let exporter = MetricExporter::builder() + .with_http() + .with_protocol(Protocol::HttpJson) + .build() + .context("build OTLP metric exporter")?; + let reader = PeriodicReader::builder(exporter) + .with_interval(interval) + .build(); + let provider = SdkMeterProvider::builder() + .with_reader(reader) + .with_resource(resource()) + .build(); + register_instruments(&provider, snapshot); + Ok(provider) +} + +/// Register the two observable gauges — `jobq.nodes` (every node at any +/// depth) and `jobq.roots` (only the root groups), each with a `state` +/// attribute — mirroring `StateCount`'s own two-count shape (see its doc +/// comment for why both numbers are kept rather than collapsed to one). +fn register_instruments(provider: &SdkMeterProvider, snapshot: Arc>>) { + let meter = provider.meter("hyperhive.jobq"); + + let snap = snapshot.clone(); + meter + .u64_observable_gauge("jobq.nodes") + .with_description("hive-jobq nodes by lifecycle state, at any depth") + .with_callback(move |obs| { + if let Ok(g) = snap.lock() { + for c in g.iter() { + obs.observe(c.nodes, &[state_attr(c)]); + } + } + }) + .build(); + + let snap = snapshot; + meter + .u64_observable_gauge("jobq.roots") + .with_description("hive-jobq root groups by lifecycle state") + .with_callback(move |obs| { + if let Ok(g) = snap.lock() { + for c in g.iter() { + obs.observe(c.roots, &[state_attr(c)]); + } + } + }) + .build(); +} + +/// The `state` attribute for one data point, spelled from `Debug` — the +/// `State` enum has no `Display`, and `Debug` on a fieldless variant is +/// exactly the variant name (`Pending`, `Running`, ...), so this costs no +/// new formatting code to keep in sync with the enum. +fn state_attr(c: &StateCount) -> KeyValue { + KeyValue::new("state", format!("{:?}", c.state)) +} + +/// Resource: `service.name = swarm-controller` plus whatever the operator +/// set in `HYPERHIVE_OTEL_EXTRA_RESOURCE_ATTRIBUTES` (same channel every +/// other OTEL exporter in this tree reads its `hive`/`swarm` labels from — +/// this daemon has no single "hive" of its own, so nothing is assumed here +/// beyond what the operator supplies). +fn resource() -> Resource { + let mut builder = Resource::builder().with_service_name("swarm-controller"); + for (k, v) in resource_attributes() { + builder = builder.with_attribute(KeyValue::new(k, v)); + } + builder.build() +} + +/// `OTEL_EXPORTER_OTLP_ENDPOINT`, non-empty. The enable signal — and the +/// very variable the SDK reads to build the exporter's URL, so "configured" +/// and "where it goes" cannot disagree. See +/// `hive-c0re::stats::otel_metrics::endpoint`'s doc comment for why this is +/// the only variable allowed to gate this exporter. +fn endpoint() -> Option { + std::env::var("OTEL_EXPORTER_OTLP_ENDPOINT") + .ok() + .map(|s| s.trim().to_owned()) + .filter(|s| !s.is_empty()) +} + +/// Export cadence from `HYPERHIVE_OTEL_METRIC_INTERVAL_MS`, else the +/// default. Same variable `hive-c0re::stats::otel_metrics` reads — one +/// cadence knob for every OTEL exporter in the tree, not a new one per +/// module. +fn interval() -> Duration { + std::env::var("HYPERHIVE_OTEL_METRIC_INTERVAL_MS") + .ok() + .and_then(|s| s.trim().parse::().ok()) + .filter(|ms| *ms > 0) + .map_or(DEFAULT_INTERVAL, Duration::from_millis) +} + +/// Extra resource attributes from `HYPERHIVE_OTEL_EXTRA_RESOURCE_ATTRIBUTES` +/// (`key=value,key=value`), same env `hive-c0re::stats::otel_metrics` reads. +fn resource_attributes() -> Vec<(String, String)> { + std::env::var("HYPERHIVE_OTEL_EXTRA_RESOURCE_ATTRIBUTES") + .ok() + .map(|s| parse_kv(&s)) + .unwrap_or_default() +} + +/// Parse `key=value` pairs separated by commas and/or newlines. The value +/// keeps any `=` after the first (so `Authorization=Bearer x=y` → `Bearer +/// x=y`). Byte-identical logic to `hive-c0re::stats::otel_metrics::parse_kv` +/// — small enough that a shared crate for it would cost more than it saves, +/// but kept in lockstep on purpose. +fn parse_kv(s: &str) -> Vec<(String, String)> { + s.split([',', '\n']) + .filter_map(|pair| { + let pair = pair.trim(); + if pair.is_empty() { + return None; + } + let (k, v) = pair.split_once('=')?; + let k = k.trim(); + if k.is_empty() { + None + } else { + Some((k.to_owned(), v.trim().to_owned())) + } + }) + .collect() +} + +#[cfg(test)] +mod tests { + use super::*; + use hive_jobq::State; + + #[test] + fn parse_kv_keeps_bearer_value() { + assert_eq!( + parse_kv("Authorization=Bearer abc=def"), + vec![("Authorization".to_owned(), "Bearer abc=def".to_owned())] + ); + assert_eq!( + parse_kv("a=1,\n b=2 ,=bad,"), + vec![ + ("a".to_owned(), "1".to_owned()), + ("b".to_owned(), "2".to_owned()) + ] + ); + } + + /// `state_attr` must spell the exact variant name — a Grafana query + /// filtering `state="Running"` should match what this daemon actually + /// emits, not a re-cased or re-worded version of it. + #[test] + fn state_attr_spells_the_variant_name() { + let c = StateCount { + state: State::Running, + nodes: 3, + roots: 1, + }; + assert_eq!(state_attr(&c), KeyValue::new("state", "Running")); + } +} diff --git a/swarm-controller/src/main.rs b/swarm-controller/src/main.rs index d1423759..bc86bc5a 100644 --- a/swarm-controller/src/main.rs +++ b/swarm-controller/src/main.rs @@ -37,6 +37,7 @@ use utoipa_axum::{router::OpenApiRouter, routes}; mod auth; mod forge; +mod jobq_metrics; mod status; mod webhook; @@ -916,6 +917,7 @@ async fn main() -> Result<()> { hive_jobq::resources::ResourceTable::new(), ))); spawn_jobq_worker(Arc::clone(&jobq), deps); + jobq_metrics::spawn_exporter(Arc::clone(&jobq)); // Same "log and carry on" shape as the queue/bridge/forge wiring above. // A controller that cannot hold a webhook secret still serves every