use chrono::{DateTime, Utc};
use serde::Serialize;
use sqlx::PgPool;
#[derive(Debug, Serialize)]
pub struct Metrics {
pub at: DateTime<Utc>,
pub queue: String,
pub runnable_queue_depth: i64,
pub jobs_per_sec: f64,
pub success_rate: f64,
pub retry_rate: f64,
pub mean_latency_ms: f64,
}
#[derive(Clone)]
pub struct MetricsRepo {
pool: PgPool,
}
impl MetricsRepo {
pub fn new(pool: PgPool) -> Self {
Self { pool }
}
pub async fn snapshot_all(&self) -> anyhow::Result<Vec<Metrics>> {
let queues: Vec<String> = sqlx::query_scalar(
r#"
SELECT DISTINCT queue
FROM jobs
ORDER BY queue
"#,
)
.fetch_all(&self.pool)
.await?;
let mut out = Vec::with_capacity(queues.len());
for queue in queues {
out.push(self.snapshot_for_queue(&queue).await?);
}
Ok(out)
}
pub async fn snapshot_for_queue(&self, queue: &str) -> anyhow::Result<Metrics> {
let depth: i64 = sqlx::query_scalar(
r#"
SELECT COUNT(*)
FROM jobs
WHERE queue = $1
AND status = 'queued'
AND run_at <= now()
"#,
)
.bind(queue)
.fetch_one(&self.pool)
.await?;
let row = sqlx::query_as::<
_,
(
Option<f64>,
Option<f64>,
Option<f64>,
Option<f64>,
Option<f64>,
),
>(
r#"
WITH a AS (
SELECT a.*
FROM job_attempts a
JOIN jobs j ON j.id = a.job_id AND j.dataset_id = a.dataset_id
WHERE j.queue = $1
AND a.started_at >= now() - interval '60 seconds'
),
finished AS (
SELECT *
FROM a
WHERE finished_at IS NOT NULL
)
SELECT
(SELECT COUNT(*) FROM finished)::float8 AS finished_count,
(SELECT COUNT(*) FROM finished WHERE status = 'succeeded')::float8 AS succeeded_count,
(SELECT COUNT(*) FROM a WHERE attempt_no >= 2)::float8 AS retry_count,
(SELECT COUNT(*) FROM a)::float8 AS started_count,
COALESCE((SELECT AVG(latency_ms)::float8 FROM finished), 0.0) AS mean_latency_ms
"#,
)
.bind(queue)
.fetch_one(&self.pool)
.await?;
let finished_count = row.0.unwrap_or(0.0);
let succeeded_count = row.1.unwrap_or(0.0);
let retry_count = row.2.unwrap_or(0.0);
let started_count = row.3.unwrap_or(0.0);
let mean_latency_ms = row.4.unwrap_or(0.0);
let jobs_per_sec = finished_count / 60.0;
let success_rate = if finished_count > 0.0 {
succeeded_count / finished_count
} else {
0.0
};
let retry_rate = if started_count > 0.0 {
retry_count / started_count
} else {
0.0
};
Ok(Metrics {
at: Utc::now(),
queue: queue.to_string(),
runnable_queue_depth: depth,
jobs_per_sec,
success_rate,
retry_rate,
mean_latency_ms,
})
}
}