swarm-controller: extract jobq_metrics into its own hive-jobq-metrics crate

This commit is contained in:
damocles 2026-08-19 17:52:12 +02:00 committed by mara
commit f5ff8698f3
7 changed files with 100 additions and 59 deletions

View file

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

View file

@ -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<N, R>` — 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`).

View file

@ -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<N, R>(
jobq: Arc<Mutex<hive_jobq::scheduler::Scheduler<N, R>>>,
service_name: &str,
) -> Option<SdkMeterProvider>
where
N: Send + 'static,
R: Clone + Eq + Hash + Send + 'static,
{
let endpoint = endpoint()?;
let snapshot: Arc<Mutex<Vec<StateCount>>> = 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<hive_jobq::NodeId> = 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<Mutex<Vec<StateCount>>>,
service_name: &str,
) -> Result<SdkMeterProvider> {
// 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<Mutex<Vec<StateCount>>>) {
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<String> {
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::<u64>().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"));
}
}