From 621245306f3a86f1da96335be37942cc90fde371 Mon Sep 17 00:00:00 2001 From: damocles Date: Wed, 19 Aug 2026 17:32:38 +0200 Subject: [PATCH] swarm-controller: move OTEL deps to workspace level, drop jobq_metrics singleton --- Cargo.toml | 15 ++++++ hive-c0re/Cargo.toml | 20 +++----- hive-metric/Cargo.toml | 14 ++---- swarm-controller/Cargo.toml | 20 +++----- swarm-controller/src/jobq_metrics.rs | 73 +++++++++++++++++----------- swarm-controller/src/main.rs | 6 ++- 6 files changed, 81 insertions(+), 67 deletions(-) diff --git a/Cargo.toml b/Cargo.toml index f7445c9a..571b3455 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -127,6 +127,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-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 72a6f361..b1dba490 100644 --- a/swarm-controller/Cargo.toml +++ b/swarm-controller/Cargo.toml @@ -38,20 +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 -# OTEL SDK for the jobq-rollup metrics exporter (`jobq_metrics.rs`). Same -# versions + blocking-client rationale as `hive-c0re`/`hive-metric`: the -# metrics SDK's `PeriodicReader` runs on a background thread with no Tokio -# reactor, where the async client panics. Kept local rather than in -# `workspace.dependencies` — same call `hive-metric` already made, this is -# still the only other crate that needs it. -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", -] } +# OTEL SDK for the jobq-rollup metrics exporter (`jobq_metrics.rs`), from +# the workspace base — see the root Cargo.toml's comment for the +# blocking-client rationale shared by every OTEL-pushing crate. +opentelemetry.workspace = true +opentelemetry_sdk.workspace = true +opentelemetry-otlp.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/jobq_metrics.rs b/swarm-controller/src/jobq_metrics.rs index a734bba5..8dcc185f 100644 --- a/swarm-controller/src/jobq_metrics.rs +++ b/swarm-controller/src/jobq_metrics.rs @@ -1,9 +1,12 @@ -//! OTEL export of the swarm-level job graph's state rollup — the same +//! OTEL export of a job graph's state rollup — the same //! `hive_jobq_wire::state_rollup()` counts `GET /api/jobq/rollup` already //! serves, ridden out to the collector on a timer instead of only on -//! request. First consumer of a generic "add jobq metrics" ask: the same -//! shape works for any `hive-jobq` instance, this crate's is just the -//! first to wire it. +//! request. First consumer of a generic "add jobq metrics" ask: [`spawn_exporter`] +//! is generic over any `hive_jobq::scheduler::Scheduler`, not tied to +//! this crate's own `SwarmNodeKind`/`SwarmResourceKind` — swarm-controller's +//! is just the first instance to wire it. No process-global state either +//! (see [`spawn_exporter`]'s doc comment): call it once per graph you want +//! exported, from as many call sites as you like. //! //! Same OTEL SDK setup as `hive-c0re::stats::otel_metrics` (container //! stats) and `hive-metric` (the one-shot CLI): the metrics SDK's @@ -19,7 +22,8 @@ //! shared snapshot on an interval, and the (sync) SDK callbacks only ever //! read that snapshot, never the scheduler directly. -use std::sync::{Arc, Mutex, OnceLock}; +use std::hash::Hash; +use std::sync::{Arc, Mutex}; use std::time::Duration; use anyhow::{Context, Result}; @@ -30,35 +34,39 @@ use opentelemetry_otlp::{MetricExporter, Protocol, WithExportConfig}; use opentelemetry_sdk::Resource; use opentelemetry_sdk::metrics::{PeriodicReader, SdkMeterProvider}; -use crate::{SwarmNodeKind, SwarmResourceKind}; - /// Default export cadence when `HYPERHIVE_OTEL_METRIC_INTERVAL_MS` is unset. /// Matches `hive-c0re::stats::otel_metrics`'s default — no reason for the /// two exporters to disagree on how fresh "current" means by default. const DEFAULT_INTERVAL: Duration = Duration::from_mins(1); -/// Latest rollup snapshot: written by the async refresher, read by the sync -/// observable-instrument callbacks. -static SNAPSHOT: OnceLock>>> = OnceLock::new(); - -/// Keep the provider alive for the process lifetime — the `PeriodicReader` -/// exports only while the provider lives. -static PROVIDER: OnceLock = OnceLock::new(); - -/// Spawn the jobq-rollup OTEL exporter if OTEL is configured -/// (`OTEL_EXPORTER_OTLP_ENDPOINT` non-empty). No-op otherwise — same +/// Spawn a jobq-rollup OTEL exporter for `jobq` if OTEL is configured +/// (`OTEL_EXPORTER_OTLP_ENDPOINT` non-empty) — `None` otherwise, same /// graceful-absence shape every other optional wiring in this daemon uses -/// (queue, bridge, forge, webhook secret). Call once at startup. -pub fn spawn_exporter( - jobq: Arc>>, -) { - let Some(endpoint) = endpoint() else { - tracing::debug!("otel jobq-metrics: no endpoint configured, exporter disabled"); - return; - }; - let snapshot = SNAPSHOT - .get_or_init(|| Arc::new(Mutex::new(Vec::new()))) - .clone(); +/// (queue, bridge, forge, webhook secret). +/// +/// **No process-global state.** Earlier drafts of this held the snapshot and +/// the `SdkMeterProvider` behind `static OnceLock`s, which quietly made this +/// a swarm-controller singleton — a second call (a second graph, a test, a +/// future host with more than one jobq instance) would have silently reused +/// the first call's state instead of exporting its own. The snapshot is now +/// a plain local `Arc`, and the provider is returned rather than stashed — +/// **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`. +/// +/// Generic over `N`/`R` (the same parameters `Scheduler` itself takes) +/// rather than this crate's own `SwarmNodeKind`/`SwarmResourceKind` — the +/// whole point of the ask was a shape any `hive-jobq` host can plug in, not +/// one hardcoded to this crate's node/resource vocabulary. +pub fn spawn_exporter( + jobq: Arc>>, +) -> 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(); @@ -81,11 +89,12 @@ pub fn spawn_exporter( match build_provider(interval, snapshot) { Ok(provider) => { - let _ = PROVIDER.set(provider); tracing::info!(%endpoint, ?interval, "otel jobq-metrics: exporter enabled"); + Some(provider) } Err(e) => { tracing::warn!(error = ?e, "otel jobq-metrics: exporter init failed"); + None } } } @@ -160,6 +169,12 @@ fn state_attr(c: &StateCount) -> KeyValue { /// other OTEL exporter in this tree reads its `hive`/`swarm` labels from — /// this daemon has no single "hive" of its own, so nothing is assumed here /// beyond what the operator supplies). +/// +/// Hardcodes `swarm-controller` even though [`spawn_exporter`] is generic — +/// a future non-swarm-controller caller of this same function would want a +/// different `service.name`, but that's a real parameter to add when a +/// second caller actually exists, not a guess to make now for one that +/// doesn't. fn resource() -> Resource { let mut builder = Resource::builder().with_service_name("swarm-controller"); for (k, v) in resource_attributes() { diff --git a/swarm-controller/src/main.rs b/swarm-controller/src/main.rs index bc86bc5a..247f85a2 100644 --- a/swarm-controller/src/main.rs +++ b/swarm-controller/src/main.rs @@ -917,7 +917,11 @@ async fn main() -> Result<()> { hive_jobq::resources::ResourceTable::new(), ))); spawn_jobq_worker(Arc::clone(&jobq), deps); - jobq_metrics::spawn_exporter(Arc::clone(&jobq)); + // 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 = jobq_metrics::spawn_exporter(Arc::clone(&jobq)); // Same "log and carry on" shape as the queue/bridge/forge wiring above. // A controller that cannot hold a webhook secret still serves every