Skip to main content

renox_core/queue/
dashboard.rs

1//! The queue dashboard.
2
3use axum::extract::{Path, State};
4use axum::response::Redirect;
5use serde::Serialize;
6
7use super::{BatchStatus, FailedJob, unix_now};
8use crate::db::Transaction;
9use crate::{AppState, Module, Result, Routes, View, context, view};
10
11/// The gate that decides who sees the dashboard.
12pub const GATE: &str = "view-queue-dashboard";
13const DONE: &str = "renox:queue:done:";
14const FAILED: &str = "renox:queue:failed:";
15
16/// A page showing the queue (`/_renox/queue`): jobs waiting per queue,
17/// how long the oldest has waited, throughput, failed jobs (retry or
18/// forget them) and recent batches.
19///
20/// ```
21/// # use renox::prelude::*;
22/// # let _ =
23/// App::new()
24///     .module(renox::queue::Dashboard)
25///     // Who may see it; nobody else, not even in development.
26///     .gate("view-queue-dashboard", |user| user.email.ends_with("@example.com"))
27/// # ;
28/// ```
29pub struct Dashboard;
30
31impl Module for Dashboard {
32    fn name(&self) -> &'static str {
33        "queue-dashboard"
34    }
35
36    fn routes(&self) -> Routes {
37        Routes::new()
38            .get("/_renox/queue", show)
39            .name("queue.dashboard")
40            .post("/_renox/queue/failed/{id}/retry", retry)
41            .name("queue.retry")
42            .post("/_renox/queue/failed/{id}/forget", forget)
43            .name("queue.forget")
44            .post("/_renox/queue/failed/retry-all", retry_all)
45            .name("queue.retry_all")
46            .require_gate(GATE)
47    }
48}
49
50/// Jobs of one queue.
51#[derive(Debug, Clone, Serialize)]
52#[non_exhaustive]
53pub struct QueueCounts {
54    /// The queue's name.
55    pub queue: String,
56    /// Available now, waiting for a worker.
57    pub ready: i64,
58    /// Waiting for their delay or backoff.
59    pub delayed: i64,
60    /// Being run.
61    pub running: i64,
62}
63
64/// The queue at a glance; from [`Queue::stats`](super::Queue::stats).
65#[derive(Debug, Clone, Serialize)]
66#[non_exhaustive]
67pub struct QueueStats {
68    /// One entry per queue that has jobs, by name.
69    pub queues: Vec<QueueCounts>,
70    /// Seconds the oldest ready job has waited.
71    pub oldest_wait: Option<i64>,
72    /// Jobs finished in the last hour.
73    pub done_last_hour: i64,
74    /// Jobs that failed for good in the last hour.
75    pub failed_last_hour: i64,
76    /// Rows in `failed_jobs`, of any age.
77    pub failed_total: i64,
78}
79
80impl super::Queue {
81    /// Counts for monitoring (the dashboard shows them).
82    pub async fn stats(&self) -> Result<QueueStats> {
83        let now = unix_now();
84        let rows = crate::db::sql(
85            "SELECT queue, \
86             SUM(CASE WHEN reserved_at IS NULL AND available_at <= ? THEN 1 ELSE 0 END) AS ready, \
87             SUM(CASE WHEN reserved_at IS NULL AND available_at > ? THEN 1 ELSE 0 END) AS delayed, \
88             SUM(CASE WHEN reserved_at IS NOT NULL THEN 1 ELSE 0 END) AS running, \
89             MIN(CASE WHEN reserved_at IS NULL AND available_at <= ? THEN available_at END) AS oldest \
90             FROM jobs GROUP BY queue ORDER BY queue",
91        )
92        .bind(now)
93        .bind(now)
94        .bind(now)
95        .fetch_all(&self.db)
96        .await?;
97        let mut queues = Vec::new();
98        let mut oldest: Option<i64> = None;
99        for row in &rows {
100            queues.push(QueueCounts {
101                queue: row.try_get("queue")?,
102                ready: row.try_get("ready")?,
103                delayed: row.try_get("delayed")?,
104                running: row.try_get("running")?,
105            });
106            if let Some(at) = row.try_get::<Option<i64>>("oldest")? {
107                oldest = Some(oldest.map_or(at, |o| o.min(at)));
108            }
109        }
110        let counters: Vec<(String, String)> =
111            crate::db::sql("SELECT key, value FROM cache WHERE key LIKE ? OR key LIKE ?")
112                .bind(format!("{DONE}%"))
113                .bind(format!("{FAILED}%"))
114                .fetch_as(&self.db)
115                .await?;
116        let since = now / 60 - 60;
117        let (mut done, mut failed) = (0, 0);
118        for (key, value) in counters {
119            let (total, minute) = match key.strip_prefix(DONE) {
120                Some(minute) => (&mut done, minute),
121                None => (&mut failed, key.strip_prefix(FAILED).unwrap_or_default()),
122            };
123            if minute.parse::<i64>().is_ok_and(|m| m > since) {
124                *total += value.parse::<i64>().unwrap_or(0);
125            }
126        }
127        let failed_total: i64 = crate::db::sql("SELECT COUNT(*) FROM failed_jobs")
128            .scalar(&self.db)
129            .await?;
130        Ok(QueueStats {
131            queues,
132            oldest_wait: oldest.map(|at| now - at),
133            done_last_hour: done,
134            failed_last_hour: failed,
135            failed_total,
136        })
137    }
138
139    /// The most recent batches, newest first.
140    pub async fn recent_batches(&self, limit: u32) -> Result<Vec<BatchStatus>> {
141        let ids: Vec<i64> = crate::db::sql("SELECT id FROM job_batches ORDER BY id DESC LIMIT ?")
142            .bind(i64::from(limit))
143            .scalars(&self.db)
144            .await?;
145        let mut batches = Vec::with_capacity(ids.len());
146        for id in ids {
147            if let Some(batch) = self.batch_status(id).await? {
148                batches.push(batch);
149            }
150        }
151        Ok(batches)
152    }
153}
154
155/// Counts a job that finished (`failed`: for good) in this minute's
156/// throughput counter, kept two hours.
157pub(crate) async fn count_finished(tx: &mut Transaction, failed: bool) -> Result {
158    let now = unix_now();
159    let prefix = if failed { FAILED } else { DONE };
160    crate::db::sql(
161        "INSERT INTO cache (key, value, expires_at) VALUES (?, '1', ?) \
162         ON CONFLICT (key) DO UPDATE SET value = CAST(CAST(cache.value AS BIGINT) + 1 AS TEXT)",
163    )
164    .bind(format!("{prefix}{}", now / 60))
165    .bind(now + 2 * 60 * 60)
166    .execute(&mut *tx)
167    .await?;
168    Ok(())
169}
170
171#[derive(Serialize)]
172struct FailedRow {
173    id: i64,
174    job: String,
175    queue: String,
176    /// The error's first line.
177    error: String,
178    ago: String,
179}
180
181fn ago(seconds: i64) -> String {
182    match seconds.max(0) {
183        s if s < 60 => format!("{s}s"),
184        s if s < 3600 => format!("{}m", s / 60),
185        s if s < 86_400 => format!("{}h", s / 3600),
186        s => format!("{}d", s / 86_400),
187    }
188}
189
190async fn show(State(state): State<AppState>) -> Result<View> {
191    let queue = &state.queue;
192    let stats = queue.stats().await?;
193    let now = unix_now();
194    let mut failed: Vec<FailedJob> = queue.failed().await?;
195    failed.reverse();
196    failed.truncate(25);
197    let failed: Vec<FailedRow> = failed
198        .into_iter()
199        .map(|f| FailedRow {
200            id: f.id,
201            error: f
202                .error
203                .lines()
204                .next()
205                .unwrap_or_default()
206                .chars()
207                .take(200)
208                .collect(),
209            job: f.job,
210            queue: f.queue,
211            ago: ago(now - f.failed_at.timestamp()),
212        })
213        .collect();
214    let batches = queue.recent_batches(10).await?;
215    let batches: Vec<_> = batches
216        .into_iter()
217        .map(|b| {
218            let progress = b.progress();
219            context! { batch => b, progress => progress }
220        })
221        .collect();
222    Ok(view(
223        "renox/queue/dashboard.html",
224        context! {
225            stats => stats,
226            oldest_wait => stats.oldest_wait.map(ago),
227            failed => failed,
228            batches => batches,
229        },
230    ))
231}
232
233async fn retry(State(state): State<AppState>, Path(id): Path<i64>) -> Result<Redirect> {
234    state.queue.retry(id).await?;
235    Ok(Redirect::to("/_renox/queue"))
236}
237
238async fn forget(State(state): State<AppState>, Path(id): Path<i64>) -> Result<Redirect> {
239    state.queue.forget_failed(id).await?;
240    Ok(Redirect::to("/_renox/queue"))
241}
242
243async fn retry_all(State(state): State<AppState>) -> Result<Redirect> {
244    state.queue.retry_all().await?;
245    Ok(Redirect::to("/_renox/queue"))
246}
247
248#[cfg(test)]
249mod tests {
250    #[test]
251    fn ages_read_in_the_largest_whole_unit() {
252        let ago = super::ago;
253        assert_eq!(
254            [
255                ago(-5),
256                ago(59),
257                ago(60),
258                ago(3599),
259                ago(3600),
260                ago(86_399),
261                ago(86_400 * 3)
262            ],
263            ["0s", "59s", "1m", "59m", "1h", "23h", "3d"]
264        );
265    }
266}