prop/rust/prop_grid_rs/src/metrics.rs
Graham McIntire b8eea1bd45
prom: wire rust workers into prometheus
prop-grid-rs already exposed prop_grid_rs_chain_step_duration_seconds
+ tasks_in_flight on :9100 but the comment still pointed at the old
prom server (10.0.15.25). hrrr-point-rs had no metrics listener at all.

- metrics.rs: new POINT_BATCH_DURATION / POINT_BATCHES_TOTAL /
  POINTS_PROCESSED_TOTAL series + record_point_batch helper, plus unit
  tests serialized via a static Mutex (the prometheus crate's global
  registry races otherwise)
- hrrr_point_worker.rs: spawn metrics::serve on METRICS_ADDR (default
  0.0.0.0:9100), wrap process_batch in InFlightGuard, record outcome +
  per-point counts on every batch
- deployment-hrrr-point-rs.yaml: prometheus.io/{scrape,port,path,job}
  annotations, METRICS_ADDR env, ports.metrics, /health readiness probe
- deployment-grid-rs.yaml: refresh stale 10.0.15.25 comment, add
  prometheus.io/job for clean dashboard up{} selectors

Also fix Microwaveprop.PromExTest: `assert plugins != []` triggered
Elixir 1.19 type-checker warning ("non_empty_list != empty_list always
true"); replaced with membership assertions for the InstrumentPlugin +
Plugins.Beam.
2026-05-08 10:36:49 -05:00

279 lines
9.2 KiB
Rust

//! Prometheus metrics + `/metrics` HTTP endpoint.
//!
//! Scraped by the Prometheus server at 10.0.15.31. Exposed on
//! `METRICS_ADDR` (default `0.0.0.0:9100`). Both binaries
//! (`worker` and `hrrr_point_worker`) share this module.
use std::sync::LazyLock;
use std::time::Duration;
use axum::{extract::State, http::StatusCode, response::IntoResponse, routing::get, Router};
use prometheus::{
register_histogram, register_histogram_vec, register_int_counter_vec, register_int_gauge,
Encoder, Histogram, HistogramVec, IntCounterVec, IntGauge, TextEncoder,
};
use tracing::{info, warn};
/// End-to-end duration of one forecast-hour chain step (fetch → decode
/// → score → write). Labelled by outcome so p99 of successful steps can
/// be tracked independently from failure tails.
pub static CHAIN_STEP_DURATION: LazyLock<HistogramVec> = LazyLock::new(|| {
register_histogram_vec!(
"prop_grid_rs_chain_step_duration_seconds",
"Wall time per grid_tasks chain step",
&["outcome"],
vec![5.0, 10.0, 20.0, 30.0, 45.0, 60.0, 90.0, 120.0, 180.0, 300.0, 600.0]
)
.expect("register histogram")
});
/// Count of chain steps by outcome. Makes the success/failure ratio
/// visible without having to derive it from the histogram.
pub static CHAIN_STEPS_TOTAL: LazyLock<IntCounterVec> = LazyLock::new(|| {
register_int_counter_vec!(
"prop_grid_rs_chain_steps_total",
"Total grid_tasks chain steps processed",
&["outcome"]
)
.expect("register counter")
});
/// Tasks currently in flight. A gauge so both "idle pod" (0) and
/// "fully saturated" (== PROP_GRID_RS_PARALLELISM) are distinguishable.
pub static TASKS_IN_FLIGHT: LazyLock<IntGauge> = LazyLock::new(|| {
register_int_gauge!(
"prop_grid_rs_tasks_in_flight",
"Number of grid_tasks rows currently being processed"
)
.expect("register gauge")
});
/// Decode-only duration (the `spawn_blocking(wgrib2)` portion). Helps
/// answer "is wgrib2 fork overhead dominant?" without instrumenting the
/// full step.
pub static DECODE_DURATION: LazyLock<Histogram> = LazyLock::new(|| {
register_histogram!(
"prop_grid_rs_decode_duration_seconds",
"wgrib2 subprocess duration per product (surface or pressure)",
vec![0.5, 1.0, 2.0, 5.0, 10.0, 20.0, 30.0, 60.0]
)
.expect("register histogram")
});
/// Per-batch wall time of `hrrr_point_worker::process_batch`. Labelled
/// by outcome so retry loops can be tracked separately from the happy
/// path.
pub static POINT_BATCH_DURATION: LazyLock<HistogramVec> = LazyLock::new(|| {
register_histogram_vec!(
"prop_grid_rs_point_batch_duration_seconds",
"Wall time per hrrr_fetch_tasks batch (fetch + decode + upsert)",
&["outcome"],
vec![0.5, 1.0, 2.5, 5.0, 10.0, 20.0, 30.0, 60.0, 120.0]
)
.expect("register histogram")
});
/// Counter of point batches by outcome (success/failure).
pub static POINT_BATCHES_TOTAL: LazyLock<IntCounterVec> = LazyLock::new(|| {
register_int_counter_vec!(
"prop_grid_rs_point_batches_total",
"Total hrrr_fetch_tasks batches processed",
&["outcome"]
)
.expect("register counter")
});
/// Counter of points by lifecycle stage (`requested` from the task,
/// `inserted` into `hrrr_profiles`). The ratio gives the cache-fill
/// rate without needing a separate gauge.
pub static POINTS_PROCESSED_TOTAL: LazyLock<IntCounterVec> = LazyLock::new(|| {
register_int_counter_vec!(
"prop_grid_rs_points_processed_total",
"Per-point counts across hrrr_fetch_tasks batches",
&["kind"]
)
.expect("register counter")
});
/// Install the lazy statics so the `/metrics` endpoint lists every
/// series from the first scrape, even before any chain step runs.
pub fn init() {
LazyLock::force(&CHAIN_STEP_DURATION);
LazyLock::force(&CHAIN_STEPS_TOTAL);
LazyLock::force(&TASKS_IN_FLIGHT);
LazyLock::force(&DECODE_DURATION);
LazyLock::force(&POINT_BATCH_DURATION);
LazyLock::force(&POINT_BATCHES_TOTAL);
LazyLock::force(&POINTS_PROCESSED_TOTAL);
}
/// RAII helper: increment on construction, decrement on drop. Guarantees
/// the in-flight gauge returns to 0 even if the worker panics.
pub struct InFlightGuard;
impl InFlightGuard {
pub fn new() -> Self {
TASKS_IN_FLIGHT.inc();
Self
}
}
impl Drop for InFlightGuard {
fn drop(&mut self) {
TASKS_IN_FLIGHT.dec();
}
}
impl Default for InFlightGuard {
fn default() -> Self {
Self::new()
}
}
pub fn record_chain_step(duration: Duration, success: bool) {
let outcome = if success { "success" } else { "failure" };
CHAIN_STEP_DURATION
.with_label_values(&[outcome])
.observe(duration.as_secs_f64());
CHAIN_STEPS_TOTAL.with_label_values(&[outcome]).inc();
}
pub fn record_point_batch(duration: Duration, success: bool, requested: u64, inserted: u64) {
let outcome = if success { "success" } else { "failure" };
POINT_BATCH_DURATION
.with_label_values(&[outcome])
.observe(duration.as_secs_f64());
POINT_BATCHES_TOTAL.with_label_values(&[outcome]).inc();
POINTS_PROCESSED_TOTAL
.with_label_values(&["requested"])
.inc_by(requested);
POINTS_PROCESSED_TOTAL
.with_label_values(&["inserted"])
.inc_by(inserted);
}
async fn metrics_handler(State(()): State<()>) -> impl IntoResponse {
let mut buf = Vec::new();
let encoder = TextEncoder::new();
let metric_families = prometheus::gather();
match encoder.encode(&metric_families, &mut buf) {
Ok(()) => (
StatusCode::OK,
[("content-type", TextEncoder::new().format_type().to_string())],
buf,
),
Err(e) => {
warn!(error = %e, "metrics encode failed");
(
StatusCode::INTERNAL_SERVER_ERROR,
[("content-type", "text/plain".to_string())],
Vec::new(),
)
}
}
}
pub async fn serve(addr: std::net::SocketAddr) {
init();
let app = Router::new()
.route("/metrics", get(metrics_handler))
.route("/health", get(|| async { "ok" }))
.with_state(());
match tokio::net::TcpListener::bind(addr).await {
Ok(listener) => {
info!(%addr, "metrics server listening");
if let Err(e) = axum::serve(listener, app).await {
warn!(error = %e, "metrics server exited");
}
}
Err(e) => {
warn!(error = %e, %addr, "failed to bind metrics listener; continuing without /metrics");
}
}
}
#[cfg(test)]
mod tests {
use std::sync::Mutex;
use super::*;
// The prometheus crate keeps a process-wide registry so parallel
// tests would race on counter deltas. Serialize via a static Mutex.
static SERIAL: Mutex<()> = Mutex::new(());
#[test]
fn record_point_batch_success_increments_counters_and_observes_duration() {
let _g = SERIAL.lock().unwrap_or_else(|e| e.into_inner());
init();
let baseline_success = POINT_BATCHES_TOTAL.with_label_values(&["success"]).get();
let baseline_requested = POINTS_PROCESSED_TOTAL.with_label_values(&["requested"]).get();
let baseline_inserted = POINTS_PROCESSED_TOTAL.with_label_values(&["inserted"]).get();
let baseline_hist_count = POINT_BATCH_DURATION
.with_label_values(&["success"])
.get_sample_count();
record_point_batch(Duration::from_millis(250), true, 12, 9);
assert_eq!(
POINT_BATCHES_TOTAL.with_label_values(&["success"]).get(),
baseline_success + 1
);
assert_eq!(
POINTS_PROCESSED_TOTAL.with_label_values(&["requested"]).get(),
baseline_requested + 12
);
assert_eq!(
POINTS_PROCESSED_TOTAL.with_label_values(&["inserted"]).get(),
baseline_inserted + 9
);
assert_eq!(
POINT_BATCH_DURATION
.with_label_values(&["success"])
.get_sample_count(),
baseline_hist_count + 1
);
}
#[test]
fn record_point_batch_failure_uses_failure_label() {
let _g = SERIAL.lock().unwrap_or_else(|e| e.into_inner());
init();
let baseline = POINT_BATCHES_TOTAL.with_label_values(&["failure"]).get();
record_point_batch(Duration::from_millis(50), false, 5, 0);
assert_eq!(
POINT_BATCHES_TOTAL.with_label_values(&["failure"]).get(),
baseline + 1
);
}
#[test]
fn record_chain_step_increments_counter() {
let _g = SERIAL.lock().unwrap_or_else(|e| e.into_inner());
init();
let baseline = CHAIN_STEPS_TOTAL.with_label_values(&["success"]).get();
record_chain_step(Duration::from_secs(15), true);
assert_eq!(
CHAIN_STEPS_TOTAL.with_label_values(&["success"]).get(),
baseline + 1
);
}
#[test]
fn in_flight_guard_inc_dec_pairs() {
let _g = SERIAL.lock().unwrap_or_else(|e| e.into_inner());
init();
let baseline = TASKS_IN_FLIGHT.get();
{
let _g = InFlightGuard::new();
assert_eq!(TASKS_IN_FLIGHT.get(), baseline + 1);
}
assert_eq!(TASKS_IN_FLIGHT.get(), baseline);
}
}