//! 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`: 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); #[async_trait] impl HttpClient for BoxedHttpClient { async fn send_bytes(&self, request: Request) -> Result, 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( jobq: Arc>>, service_name: &str, http_client: Option>, ) -> 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, 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>>, service_name: &str, http_client: Option>, ) -> 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 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>>) { 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")); } }