diff --git a/Cargo.lock b/Cargo.lock index d95ae5b5..a3710f58 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1785,20 +1785,6 @@ 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" @@ -4572,7 +4558,6 @@ 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 f2c58eb9..f7445c9a 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -12,7 +12,6 @@ members = [ "hive-forge-notify", "hive-host-sock", "hive-jobq", - "hive-jobq-metrics", "hive-jobq-wire", "hive-matrix-mcp", "hive-metric", @@ -81,7 +80,6 @@ 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" @@ -129,21 +127,6 @@ 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 6ffc9549..8f23a07a 100644 --- a/hive-c0re/Cargo.toml +++ b/hive-c0re/Cargo.toml @@ -23,12 +23,20 @@ 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). "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"] } +# (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", +] } 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 deleted file mode 100644 index e48094e3..00000000 --- a/hive-jobq-metrics/Cargo.toml +++ /dev/null @@ -1,18 +0,0 @@ -[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 deleted file mode 100644 index 4996f34f..00000000 --- a/hive-jobq-metrics/README.md +++ /dev/null @@ -1,20 +0,0 @@ -# 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 deleted file mode 100644 index b927ceef..00000000 --- a/hive-jobq-metrics/src/lib.rs +++ /dev/null @@ -1,263 +0,0 @@ -//! 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 c6f6e045..e066970c 100644 --- a/hive-metric/Cargo.toml +++ b/hive-metric/Cargo.toml @@ -11,8 +11,9 @@ path = "src/main.rs" [dependencies] anyhow.workspace = true clap.workspace = true -# OTEL SDK, from the workspace base — see the root Cargo.toml's comment for -# the blocking-client rationale shared by every OTEL-pushing crate. +# 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"] } # 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 @@ -20,6 +21,9 @@ clap.workspace = true # 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.workspace = true -opentelemetry_sdk.workspace = true -opentelemetry-otlp.workspace = true +opentelemetry-otlp = { version = "0.32", default-features = false, features = [ + "metrics", + "http-json", + "reqwest-blocking-client", + "reqwest-rustls", +] } diff --git a/swarm-controller/Cargo.toml b/swarm-controller/Cargo.toml index 9716f4e8..cb7e0419 100644 --- a/swarm-controller/Cargo.toml +++ b/swarm-controller/Cargo.toml @@ -38,12 +38,6 @@ 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 b3a84a69..d1423759 100644 --- a/swarm-controller/src/main.rs +++ b/swarm-controller/src/main.rs @@ -916,12 +916,6 @@ 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