diff --git a/Cargo.lock b/Cargo.lock index a3710f58..d95ae5b5 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1785,6 +1785,20 @@ dependencies = [ "uuid", ] +[[package]] +name = "hive-jobq-metrics" +version = "0.1.0" +dependencies = [ + "anyhow", + "hive-jobq", + "hive-jobq-wire", + "opentelemetry", + "opentelemetry-otlp", + "opentelemetry_sdk", + "tokio", + "tracing", +] + [[package]] name = "hive-jobq-wire" version = "0.1.0" @@ -4558,6 +4572,7 @@ dependencies = [ "forgejo-api", "futures-util", "hive-jobq", + "hive-jobq-metrics", "hive-jobq-wire", "hive-types", "hmac 0.13.0", diff --git a/Cargo.toml b/Cargo.toml index f7445c9a..f2c58eb9 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -12,6 +12,7 @@ members = [ "hive-forge-notify", "hive-host-sock", "hive-jobq", + "hive-jobq-metrics", "hive-jobq-wire", "hive-matrix-mcp", "hive-metric", @@ -80,6 +81,7 @@ indicatif = "0.18" hive-sh4re = { path = "hive-sh4re" } hive-agent-sock = { path = "hive-agent-sock" } hive-jobq = { path = "hive-jobq" } +hive-jobq-metrics = { path = "hive-jobq-metrics" } hive-jobq-wire = { path = "hive-jobq-wire" } hive-core-agent-sock = { path = "hive-core-agent-sock" } hive-claude = "0.1" @@ -127,6 +129,21 @@ reqwest = { version = "0.13", default-features = false, features = [ ] } hyper = { version = "1", features = ["client", "http1"] } hyper-util = { version = "0.1", features = ["tokio"] } +# OTEL SDK, shared by every crate that pushes metrics (hive-c0re, hive-metric, +# swarm-controller). The blocking OTLP client is deliberate everywhere it's +# used: the metrics SDK's `PeriodicReader` drives its export from a +# background thread with no Tokio reactor, where the async client panics. +# `default-features = false` on the OTLP crate drops its own diagnostics +# ("internal-logs") unless a member opts back in — hive-c0re does, via +# `{ workspace = true, features = ["internal-logs"] }`. +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", +] } http-body-util = "0.1" # ⚠️ Keep at 0.11.1 or newer, and keep it on the SAME reqwest as everything # else. 0.11.0 links reqwest 0.12 while the workspace is on 0.13, and cargo diff --git a/hive-c0re/Cargo.toml b/hive-c0re/Cargo.toml index 8f23a07a..6ffc9549 100644 --- a/hive-c0re/Cargo.toml +++ b/hive-c0re/Cargo.toml @@ -23,20 +23,12 @@ clap.workspace = true clap_complete.workspace = true clap-markdown = "0.1" # OTEL SDK for the per-agent container-resource metrics exporter -# (stats/otel_metrics.rs). Same versions as hive-metric — the blocking OTLP -# client is deliberate: the metrics SDK's PeriodicReader runs on a background -# thread with no Tokio reactor, where the async client panics. -opentelemetry = "0.32" -opentelemetry_sdk = { version = "0.32", features = ["metrics"] } -opentelemetry-otlp = { version = "0.32", default-features = false, features = [ - # Not a transport: it is in `default`, so `default-features = false` drops the - # exporter's own diagnostics unless it is named here. - "internal-logs", - "metrics", - "http-json", - "reqwest-blocking-client", - "reqwest-rustls", -] } +# (stats/otel_metrics.rs). "internal-logs" on top of the workspace base is +# this crate's own opt-in: it's not a transport, it's in `default`, and +# `default-features = false` upstream drops it unless named here. +opentelemetry.workspace = true +opentelemetry_sdk.workspace = true +opentelemetry-otlp = { workspace = true, features = ["internal-logs"] } indicatif.workspace = true hive-core-agent-sock.workspace = true hive-sh4re.workspace = true diff --git a/hive-jobq-metrics/Cargo.toml b/hive-jobq-metrics/Cargo.toml new file mode 100644 index 00000000..e48094e3 --- /dev/null +++ b/hive-jobq-metrics/Cargo.toml @@ -0,0 +1,18 @@ +[package] +name = "hive-jobq-metrics" +edition.workspace = true +version.workspace = true +readme = "README.md" + +[lints] +workspace = true + +[dependencies] +anyhow.workspace = true +hive-jobq.workspace = true +hive-jobq-wire.workspace = true +opentelemetry.workspace = true +opentelemetry_sdk.workspace = true +opentelemetry-otlp.workspace = true +tokio.workspace = true +tracing.workspace = true diff --git a/hive-jobq-metrics/README.md b/hive-jobq-metrics/README.md new file mode 100644 index 00000000..4996f34f --- /dev/null +++ b/hive-jobq-metrics/README.md @@ -0,0 +1,20 @@ +# hive-jobq-metrics + +OTEL export of a [`hive-jobq`](../hive-jobq) graph's state rollup, using the +same [`hive-jobq-wire`](../hive-jobq-wire) `state_rollup()` counts a viewer's +`/rollup` endpoint would serve — ridden out to the collector on a timer +instead of only on request. + +**Why this is not part of `hive-jobq` or `hive-jobq-wire`.** Both of those +crates are dependency-light on purpose (no `tokio`, no HTTP client) — the +scheduler is logic, the wire crate is presentation, and neither wants to drag +the OTEL SDK, an async runtime, and an OTLP HTTP client into every consumer +that just wants to run a graph or serialize one to JSON. Metrics export is a +third concern with its own weight, so it gets its own crate rather than +bloating either of theirs. + +`spawn_exporter` is generic over `hive_jobq::scheduler::Scheduler` — any +host's jobq instance can call it, not just one hardcoded caller. It holds no +process-global state: call it once per graph you want exported, and keep the +returned `SdkMeterProvider` alive for as long as export should continue +(dropping it stops the `PeriodicReader`). diff --git a/hive-jobq-metrics/src/lib.rs b/hive-jobq-metrics/src/lib.rs new file mode 100644 index 00000000..b927ceef --- /dev/null +++ b/hive-jobq-metrics/src/lib.rs @@ -0,0 +1,263 @@ +//! OTEL export of a [`hive_jobq`] graph's state rollup — the same +//! [`hive_jobq_wire::state_rollup()`] counts a viewer's rollup endpoint +//! already serves, ridden out to the collector on a timer instead of only +//! on request. See the crate README for why this lives in its own crate +//! rather than inside `hive-jobq` or `hive-jobq-wire`. +//! +//! 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 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::hash::Hash; +use std::sync::{Arc, Mutex}; +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}; + +/// Default export cadence when `HYPERHIVE_OTEL_METRIC_INTERVAL_MS` is unset. +/// Matches `hive-c0re::stats::otel_metrics`'s default — no reason for the +/// exporters across this tree to disagree on how fresh "current" means by +/// default. +const DEFAULT_INTERVAL: Duration = Duration::from_mins(1); + +/// Spawn a jobq-rollup OTEL exporter for `jobq` if OTEL is configured +/// (`OTEL_EXPORTER_OTLP_ENDPOINT` non-empty) — `None` otherwise, same +/// graceful-absence shape optional OTEL wiring uses elsewhere in this tree. +/// `service_name` becomes the exported resource's `service.name` — the +/// caller's own binary name, since this crate has no opinion on who's +/// calling it. +/// +/// **No process-global state.** The snapshot is a plain local `Arc`, and +/// the provider is returned rather than stashed in a static — **the caller +/// owns it and must keep it alive** (bind it to a named variable, not `_`) +/// for as long as export should continue; dropping it stops the +/// `PeriodicReader`. Call this once per graph you want exported, from as +/// many call sites as you like — nothing here assumes there's only one. +/// +/// Generic over `N`/`R` (the same parameters [`hive_jobq::scheduler::Scheduler`] +/// itself takes) rather than any one host's node/resource vocabulary — the +/// whole point of this crate is a shape any `hive-jobq` host can plug in. +pub fn spawn_exporter( + jobq: Arc>>, + service_name: &str, +) -> Option +where + N: Send + 'static, + R: Clone + Eq + Hash + Send + 'static, +{ + let endpoint = endpoint()?; + let snapshot: Arc>> = Arc::new(Mutex::new(Vec::new())); + 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, service_name) { + Ok(provider) => { + tracing::info!(%endpoint, ?interval, %service_name, "otel jobq-metrics: exporter enabled"); + Some(provider) + } + Err(e) => { + tracing::warn!(error = ?e, "otel jobq-metrics: exporter init failed"); + None + } + } +} + +fn build_provider( + interval: Duration, + snapshot: Arc>>, + service_name: &str, +) -> 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(service_name)) + .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 = service_name` 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). +fn resource(service_name: &str) -> Resource { + let mut builder = Resource::builder().with_service_name(service_name.to_owned()); + 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 sharing it isn't worth a dependency edge, 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 crate 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/hive-metric/Cargo.toml b/hive-metric/Cargo.toml index e066970c..c6f6e045 100644 --- a/hive-metric/Cargo.toml +++ b/hive-metric/Cargo.toml @@ -11,9 +11,8 @@ path = "src/main.rs" [dependencies] anyhow.workspace = true clap.workspace = true -# OTEL SDK — only used by this crate, so kept local rather than in workspace.dependencies. -opentelemetry = "0.32" -opentelemetry_sdk = { version = "0.32", features = ["metrics"] } +# OTEL SDK, from the workspace base — see the root Cargo.toml's comment for +# the blocking-client rationale shared by every OTEL-pushing crate. # Blocking (not async) reqwest client on purpose: the metrics SDK drives the # OTLP push from a `PeriodicReader` background thread that has no Tokio runtime, # so the async client panics there with "no reactor running". The blocking @@ -21,9 +20,6 @@ opentelemetry_sdk = { version = "0.32", features = ["metrics"] } # records one point then `shutdown()`s, a synchronous send is exactly right. # No `internal-logs` here, unlike hive-c0re: this binary installs no tracing # subscriber, so the SDK's diagnostics would have nowhere to go. -opentelemetry-otlp = { version = "0.32", default-features = false, features = [ - "metrics", - "http-json", - "reqwest-blocking-client", - "reqwest-rustls", -] } +opentelemetry.workspace = true +opentelemetry_sdk.workspace = true +opentelemetry-otlp.workspace = true diff --git a/swarm-controller/Cargo.toml b/swarm-controller/Cargo.toml index cb7e0419..9716f4e8 100644 --- a/swarm-controller/Cargo.toml +++ b/swarm-controller/Cargo.toml @@ -38,6 +38,12 @@ 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 +# The jobq-rollup OTEL exporter, wired up in `main` via +# `hive_jobq_metrics::spawn_exporter` — moved to its own crate (rather than +# living here as `jobq_metrics.rs`) specifically so a future second caller +# (e.g. hive-c0re, for its own per-hive job graph) doesn't have to depend on +# this whole binary to reuse it. +hive-jobq-metrics.workspace = true # 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/main.rs b/swarm-controller/src/main.rs index d1423759..b3a84a69 100644 --- a/swarm-controller/src/main.rs +++ b/swarm-controller/src/main.rs @@ -916,6 +916,12 @@ async fn main() -> Result<()> { hive_jobq::resources::ResourceTable::new(), ))); spawn_jobq_worker(Arc::clone(&jobq), deps); + // Bound to a named variable, not `_` — dropping the provider stops its + // `PeriodicReader`, so it must live as long as `main` does (which it + // does here: this binding never goes out of scope before the process + // exits via `axum::serve(...).await` below). + let _jobq_metrics_provider = + hive_jobq_metrics::spawn_exporter(Arc::clone(&jobq), "swarm-controller"); // Same "log and carry on" shape as the queue/bridge/forge wiring above. // A controller that cannot hold a webhook secret still serves every