hyperhive/hive-jobq-metrics/src/lib.rs

304 lines
12 KiB
Rust

//! 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 async_trait::async_trait;
use hive_jobq_wire::StateCount;
use opentelemetry::KeyValue;
use opentelemetry::metrics::MeterProvider as _;
use opentelemetry_http::{Bytes, HttpClient, HttpError, Request, Response};
use opentelemetry_otlp::{MetricExporter, Protocol, WithExportConfig, WithHttpConfig};
use opentelemetry_sdk::Resource;
use opentelemetry_sdk::metrics::{PeriodicReader, SdkMeterProvider};
/// Type-erases a caller's concrete [`HttpClient`] impl so [`spawn_exporter`]
/// can accept "any client, or none" without this crate depending on any one
/// caller's concrete type (`swarm-controller`'s own
/// `otel_http_client::AuthenticatedHttpClient` today, conceivably something
/// else from a future caller) — the same genericity
/// this crate already has over `N`/`R`, applied to the one other caller-
/// supplied thing it touches.
///
/// A thin delegating newtype rather than a blanket `impl HttpClient for
/// Box<dyn HttpClient>`: the orphan rule refuses that blanket impl outright
/// (`Box` and `dyn HttpClient` are both foreign to this crate), so a local
/// wrapper is the standard way around it, not a workaround chosen over a
/// cleaner alternative.
#[derive(Debug)]
struct BoxedHttpClient(Box<dyn HttpClient>);
#[async_trait]
impl HttpClient for BoxedHttpClient {
async fn send_bytes(&self, request: Request<Bytes>) -> Result<Response<Bytes>, HttpError> {
self.0.send_bytes(request).await
}
}
/// 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.
///
/// `http_client`: `None` builds the exporter with its default (endpoint-
/// only, no per-request auth) transport — the shape every caller used
/// before this parameter existed. `Some(client)` routes every export
/// through that [`HttpClient`] instead, for a caller whose destination
/// checks a bearer token per request rather than a static header
/// (`swarm-controller`'s case — see its own `otel_http_client` module doc
/// for why a static header doesn't fit an expiring token).
pub fn spawn_exporter<N, R>(
jobq: Arc<Mutex<hive_jobq::scheduler::Scheduler<N, R>>>,
service_name: &str,
http_client: Option<Box<dyn HttpClient>>,
) -> 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, http_client) {
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,
http_client: Option<Box<dyn HttpClient>>,
) -> 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 mut builder = MetricExporter::builder()
.with_http()
.with_protocol(Protocol::HttpJson);
// Only when the caller supplied one — see `spawn_exporter`'s doc for
// who does and why. Wrapped in `BoxedHttpClient` because
// `with_http_client` wants an owned `T: HttpClient`, not the trait
// object this parameter is typed as (see that newtype's own doc).
if let Some(client) = http_client {
builder = builder.with_http_client(BoxedHttpClient(client));
}
let exporter = builder.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"));
}
}