Compare commits
| Author | SHA1 | Date | |
|---|---|---|---|
|
|
f5ff8698f3 | ||
|
|
621245306f | ||
|
|
1f5a7b71ed |
9 changed files with 356 additions and 23 deletions
15
Cargo.lock
generated
15
Cargo.lock
generated
|
|
@ -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",
|
||||
|
|
|
|||
17
Cargo.toml
17
Cargo.toml
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
18
hive-jobq-metrics/Cargo.toml
Normal file
18
hive-jobq-metrics/Cargo.toml
Normal 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
|
||||
20
hive-jobq-metrics/README.md
Normal file
20
hive-jobq-metrics/README.md
Normal 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`).
|
||||
263
hive-jobq-metrics/src/lib.rs
Normal file
263
hive-jobq-metrics/src/lib.rs
Normal 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"));
|
||||
}
|
||||
}
|
||||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
|
|
@ -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
|
||||
|
|
|
|||
Loading…
Reference in a new issue