Compare commits

...
9 changed files with 356 additions and 23 deletions

15
Cargo.lock generated
View file

@ -1785,6 +1785,20 @@ 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"
@ -4558,6 +4572,7 @@ dependencies = [
"forgejo-api",
"futures-util",
"hive-jobq",
"hive-jobq-metrics",
"hive-jobq-wire",
"hive-types",
"hmac 0.13.0",

View file

@ -12,6 +12,7 @@ members = [
"hive-forge-notify",
"hive-host-sock",
"hive-jobq",
"hive-jobq-metrics",
"hive-jobq-wire",
"hive-matrix-mcp",
"hive-metric",
@ -80,6 +81,7 @@ 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"
@ -127,6 +129,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

View file

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

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"));
}
}

View file

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

View file

@ -38,6 +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
# 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

View file

@ -916,6 +916,12 @@ 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