Skip to main content

hammerwork_web/api/
history.rs

1//! Historical job queries backing the dashboard's trend, error-pattern and clear
2//! endpoints.
3//!
4//! [`DatabaseQueue`] has no way to ask "how many jobs completed or failed in each
5//! hour", "which errors occurred where" or "delete this queue's completed jobs". The
6//! dashboard talks to a concrete [`JobQueue`], so [`JobHistory`] answers those with
7//! parameterized SQL against `hammerwork_jobs` for both backends.
8//!
9//! "Failed" here means jobs *currently* in `Failed`, `Dead` or `TimedOut`, bucketed by
10//! when they failed (`failed_at`, falling back to `timed_out_at`). A job that failed and
11//! was then retried to success is counted only as completed, because retrying clears its
12//! failure timestamps.
13
14use super::archive::ArchiveFilterParams;
15use chrono::{DateTime, Duration, NaiveDateTime, TimeZone, Utc};
16use hammerwork::archive::{ArchivalReason, ArchivedJob};
17use hammerwork::queue::DatabaseQueue;
18use hammerwork::{JobQueue, JobStatus, Result};
19use serde::Serialize;
20use sqlx::Row;
21use std::collections::BTreeMap;
22use std::future::Future;
23
24use crate::websocket::JobUpdate;
25
26/// Completed/failed activity in one UTC hour.
27#[derive(Debug, Clone, PartialEq, Serialize)]
28pub struct HourBucket {
29    /// Start of the hour (UTC).
30    pub hour: DateTime<Utc>,
31    pub completed: u64,
32    pub failed: u64,
33    /// Mean run time of the jobs completed in the hour, if any recorded one.
34    pub avg_processing_time_ms: Option<f64>,
35}
36
37/// One distinct error message with how often and where it happened.
38#[derive(Debug, Clone, PartialEq)]
39pub struct ErrorGroup {
40    pub queue_name: String,
41    pub message: String,
42    pub count: u64,
43    pub first_seen: DateTime<Utc>,
44    pub last_seen: DateTime<Utc>,
45}
46
47/// Which finished jobs [`JobHistory::count_jobs`] / [`JobHistory::delete_jobs`] target.
48#[derive(Debug, Clone, Copy, PartialEq, Eq)]
49pub enum FinishedKind {
50    Completed,
51    Dead,
52}
53
54impl FinishedKind {
55    fn status(self) -> &'static str {
56        match self {
57            FinishedKind::Completed => "Completed",
58            FinishedKind::Dead => "Dead",
59        }
60    }
61
62    /// Column the `older_than` cutoff applies to.
63    fn age_column(self) -> &'static str {
64        match self {
65            FinishedKind::Completed => "completed_at",
66            FinishedKind::Dead => "COALESCE(failed_at, timed_out_at, created_at)",
67        }
68    }
69}
70
71/// A table maintenance operation run from the dashboard.
72///
73/// | Operation  | PostgreSQL               | MySQL            |
74/// |------------|--------------------------|------------------|
75/// | `vacuum`   | `VACUUM (ANALYZE) <t>`   | `OPTIMIZE TABLE` |
76/// | `reindex`  | `REINDEX TABLE <t>`      | `OPTIMIZE TABLE` |
77/// | `optimize` | `ANALYZE <t>`            | `OPTIMIZE TABLE` |
78///
79/// MySQL has neither `VACUUM` nor `REINDEX`. For InnoDB, `OPTIMIZE TABLE` rebuilds the
80/// table and its indexes, reclaims free space and refreshes the statistics, which covers
81/// all three operations.
82#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize)]
83#[serde(rename_all = "lowercase")]
84pub enum TableMaintenance {
85    Vacuum,
86    Reindex,
87    Optimize,
88}
89
90impl TableMaintenance {
91    /// Parses an operation name (`vacuum`, `reindex`, `optimize`).
92    pub fn parse(name: &str) -> Option<Self> {
93        match name {
94            "vacuum" => Some(Self::Vacuum),
95            "reindex" => Some(Self::Reindex),
96            "optimize" => Some(Self::Optimize),
97            _ => None,
98        }
99    }
100
101    /// The PostgreSQL statement for `table`, one of [`HAMMERWORK_TABLES`].
102    pub fn postgres_statement(self, table: &'static str) -> String {
103        match self {
104            Self::Vacuum => format!("VACUUM (ANALYZE) {table}"),
105            Self::Reindex => format!("REINDEX TABLE {table}"),
106            Self::Optimize => format!("ANALYZE {table}"),
107        }
108    }
109
110    /// The MySQL statement for `table`, one of [`HAMMERWORK_TABLES`].
111    pub fn mysql_statement(self, table: &'static str) -> String {
112        format!("OPTIMIZE TABLE {table}")
113    }
114}
115
116/// Every table Hammerwork owns. Table maintenance only ever names these fixed identifiers;
117/// the ones absent from the database (e.g. before a migration created them) are skipped.
118pub const HAMMERWORK_TABLES: &[&str] = &[
119    "hammerwork_jobs",
120    "hammerwork_jobs_archive",
121    "hammerwork_batches",
122    "hammerwork_workflows",
123    "hammerwork_queue_pause",
124    "hammerwork_encryption_keys",
125    "hammerwork_kms_data_keys",
126    "hammerwork_key_audit_log",
127    "hammerwork_migrations",
128];
129
130/// Resolves an optional table name to the matching entries of [`HAMMERWORK_TABLES`].
131/// `None` selects every table; a name that is not a Hammerwork table is an error.
132pub fn maintenance_tables(target: Option<&str>) -> std::result::Result<Vec<&'static str>, String> {
133    match target {
134        None => Ok(HAMMERWORK_TABLES.to_vec()),
135        Some(name) => HAMMERWORK_TABLES
136            .iter()
137            .find(|t| **t == name)
138            .map(|t| vec![*t])
139            .ok_or_else(|| format!("'{name}' is not a Hammerwork table")),
140    }
141}
142
143/// One statement a table maintenance operation ran (or, on a dry run, would run).
144#[derive(Debug, Clone, PartialEq, Serialize)]
145pub struct MaintenanceStatement {
146    pub table: &'static str,
147    pub statement: String,
148    /// What the database reported (MySQL `OPTIMIZE TABLE` result rows); empty otherwise.
149    pub messages: Vec<String>,
150}
151
152/// The entries of `candidates` whose names appear in `existing` (names from the catalog).
153/// Statements are built from the fixed `candidates`, never from catalog strings.
154fn present_tables(candidates: &[&'static str], existing: &[String]) -> Vec<&'static str> {
155    candidates
156        .iter()
157        .copied()
158        .filter(|t| existing.iter().any(|e| e == t))
159        .collect()
160}
161
162/// Maximum number of hourly buckets a single request may span.
163pub const MAX_BUCKETS: i64 = 24 * 31;
164
165/// Truncates to the start of the UTC hour.
166pub fn hour_floor(t: DateTime<Utc>) -> DateTime<Utc> {
167    let secs = t.timestamp() - t.timestamp().rem_euclid(3600);
168    Utc.timestamp_opt(secs, 0).single().unwrap_or(t)
169}
170
171/// Dense list of hour starts covering `[start, end)`, oldest first.
172pub fn hour_range(start: DateTime<Utc>, end: DateTime<Utc>) -> Vec<DateTime<Utc>> {
173    let mut hours = Vec::new();
174    let mut cursor = hour_floor(start);
175    while cursor < end && (hours.len() as i64) < MAX_BUCKETS {
176        hours.push(cursor);
177        cursor += Duration::hours(1);
178    }
179    hours
180}
181
182/// Merges sparse per-hour counts into a zero-filled series over `hours`.
183pub fn fill_buckets(
184    hours: &[DateTime<Utc>],
185    completed: &BTreeMap<DateTime<Utc>, (u64, Option<f64>)>,
186    failed: &BTreeMap<DateTime<Utc>, u64>,
187) -> Vec<HourBucket> {
188    hours
189        .iter()
190        .map(|hour| {
191            let (completed, avg) = completed.get(hour).copied().unwrap_or((0, None));
192            HourBucket {
193                hour: *hour,
194                completed,
195                failed: failed.get(hour).copied().unwrap_or(0),
196                avg_processing_time_ms: avg,
197            }
198        })
199        .collect()
200}
201
202fn to_u64(n: i64) -> u64 {
203    u64::try_from(n).unwrap_or(0)
204}
205
206/// SQL-backed history queries. Implemented for PostgreSQL and MySQL job queues.
207pub trait JobHistory: DatabaseQueue + Send + Sync {
208    /// Zero-filled hourly activity over `[start, end)`, optionally for one queue.
209    fn hourly_activity(
210        &self,
211        queue: Option<&str>,
212        start: DateTime<Utc>,
213        end: DateTime<Utc>,
214    ) -> impl Future<Output = Result<Vec<HourBucket>>> + Send;
215
216    /// Distinct error messages of failed jobs since `since`, most frequent first.
217    fn error_groups(
218        &self,
219        since: DateTime<Utc>,
220        limit: u32,
221    ) -> impl Future<Output = Result<Vec<ErrorGroup>>> + Send;
222
223    /// Number of finished jobs of `kind`, optionally per queue and older than a cutoff.
224    fn count_jobs(
225        &self,
226        queue: Option<&str>,
227        kind: FinishedKind,
228        older_than: Option<DateTime<Utc>>,
229    ) -> impl Future<Output = Result<u64>> + Send;
230
231    /// Deletes the jobs [`count_jobs`](Self::count_jobs) counts; returns how many.
232    fn delete_jobs(
233        &self,
234        queue: Option<&str>,
235        kind: FinishedKind,
236        older_than: Option<DateTime<Utc>>,
237    ) -> impl Future<Output = Result<u64>> + Send;
238
239    /// One page of the archived jobs matching `filter`, newest first, and the number of
240    /// matching archived jobs. Filtering, counting and paging all happen in the database.
241    fn archived_jobs(
242        &self,
243        filter: &ArchiveFilterParams,
244        limit: u32,
245        offset: u32,
246    ) -> impl Future<Output = Result<(Vec<ArchivedJob>, u64)>> + Send;
247
248    /// The database's current time, the clock the job state timestamps are written with.
249    fn database_now(&self) -> impl Future<Output = Result<DateTime<Utc>>> + Send;
250
251    /// Jobs whose state changed after `since`, newest first, at most `limit`: jobs that
252    /// were created, started, completed, failed or timed out since then. Each is reported
253    /// with the time of its latest such change. Rows that cannot be decoded are skipped.
254    fn recent_job_changes(
255        &self,
256        since: DateTime<Utc>,
257        limit: u32,
258    ) -> impl Future<Output = Result<Vec<JobUpdate>>> + Send;
259
260    /// Runs `operation` on each of `tables` (entries of [`HAMMERWORK_TABLES`]) that
261    /// exists, one statement per table, and returns what ran. With `dry_run` nothing is
262    /// executed and the statements that would run are returned.
263    fn table_maintenance(
264        &self,
265        operation: TableMaintenance,
266        tables: &[&'static str],
267        dry_run: bool,
268    ) -> impl Future<Output = Result<Vec<MaintenanceStatement>>> + Send;
269}
270
271/// The job timestamps that mark a change of state, newest first in
272/// [`JobHistory::recent_job_changes`].
273const CHANGE_COLUMNS: [&str; 5] = [
274    "created_at",
275    "started_at",
276    "completed_at",
277    "failed_at",
278    "timed_out_at",
279];
280
281/// `WHERE` clause of [`JobHistory::recent_job_changes`]: any change column after the
282/// cutoff. Each comparison can use an index on its column.
283fn changed_since_filter(postgres: bool) -> String {
284    let placeholder = if postgres { "$1" } else { "?" };
285    CHANGE_COLUMNS
286        .iter()
287        .map(|column| format!("{column} > {placeholder}"))
288        .collect::<Vec<_>>()
289        .join(" OR ")
290}
291
292/// A changed job as sent to dashboard clients.
293fn job_change(
294    id: String,
295    queue_name: String,
296    status: &str,
297    priority: i32,
298    attempts: i32,
299    changed_at: DateTime<Utc>,
300) -> JobUpdate {
301    JobUpdate {
302        id,
303        queue_name,
304        status: status.trim_matches('"').to_string(),
305        priority: hammerwork::JobPriority::from_i32(priority)
306            .unwrap_or_default()
307            .to_string(),
308        attempts,
309        updated_at: changed_at,
310    }
311}
312
313/// A value bound to an archive filter placeholder.
314#[derive(Debug, Clone, PartialEq)]
315enum ArchiveBind {
316    Text(String),
317    Time(DateTime<Utc>),
318    Bool(bool),
319}
320
321/// The `WHERE` clause (empty without filters) for `filter`, with `$n`/`?` placeholders, and
322/// the values to bind to them in order. Statuses and reasons compare case-insensitively and
323/// ignore the JSON quotes older rows may carry.
324fn archive_conditions(filter: &ArchiveFilterParams, postgres: bool) -> (String, Vec<ArchiveBind>) {
325    let mut clauses = Vec::new();
326    let mut binds = Vec::new();
327    let mut add = |clause: &str, bind: ArchiveBind| {
328        binds.push(bind);
329        let placeholder = if postgres {
330            format!("${}", binds.len())
331        } else {
332            "?".to_string()
333        };
334        clauses.push(clause.replace('?', &placeholder));
335    };
336    if let Some(queue) = &filter.queue {
337        add("queue_name = ?", ArchiveBind::Text(queue.clone()));
338    }
339    if let Some(reason) = &filter.reason {
340        add(
341            "LOWER(TRIM(BOTH '\"' FROM archival_reason)) = ?",
342            ArchiveBind::Text(reason.to_lowercase()),
343        );
344    }
345    if let Some(after) = filter.archived_after {
346        add("archived_at >= ?", ArchiveBind::Time(after));
347    }
348    if let Some(before) = filter.archived_before {
349        add("archived_at <= ?", ArchiveBind::Time(before));
350    }
351    if let Some(by) = &filter.archived_by {
352        add("archived_by = ?", ArchiveBind::Text(by.clone()));
353    }
354    if let Some(compressed) = filter.compressed {
355        add("payload_compressed = ?", ArchiveBind::Bool(compressed));
356    }
357    if let Some(status) = &filter.original_status {
358        add(
359            "LOWER(TRIM(BOTH '\"' FROM status)) = ?",
360            ArchiveBind::Text(status.to_lowercase()),
361        );
362    }
363    if clauses.is_empty() {
364        (String::new(), binds)
365    } else {
366        (format!("WHERE {}", clauses.join(" AND ")), binds)
367    }
368}
369
370/// Binds every [`ArchiveBind`] to a query, in order.
371macro_rules! bind_archive_filter {
372    ($query:expr, $binds:expr) => {{
373        let mut query = $query;
374        for bind in $binds {
375            query = match bind {
376                ArchiveBind::Text(value) => query.bind(value.clone()),
377                ArchiveBind::Time(value) => query.bind(*value),
378                ArchiveBind::Bool(value) => query.bind(*value),
379            };
380        }
381        query
382    }};
383}
384
385/// Columns of an archived job listing.
386const ARCHIVED_JOB_COLUMNS: &str = "id, queue_name, status, created_at, archived_at, \
387     archival_reason, original_payload_size, payload_compressed, archived_by";
388
389/// The status stored in the archive (possibly JSON-quoted); unknown values read as `Dead`,
390/// as the core library does.
391fn archived_status(value: &str) -> JobStatus {
392    let value = value.trim_matches('"');
393    [
394        JobStatus::Pending,
395        JobStatus::Running,
396        JobStatus::Completed,
397        JobStatus::Failed,
398        JobStatus::Dead,
399        JobStatus::TimedOut,
400        JobStatus::Retrying,
401        JobStatus::Archived,
402    ]
403    .into_iter()
404    .find(|status| status.as_str() == value)
405    .unwrap_or(JobStatus::Dead)
406}
407
408/// The rows of a dashboard listing that could be decoded. A row that cannot be decoded
409/// is skipped with a warning naming its id, so one corrupt row does not turn the whole
410/// page into an error.
411pub(crate) fn decodable<T>(rows: impl IntoIterator<Item = (String, Result<T>)>) -> Vec<T> {
412    rows.into_iter()
413        .filter_map(|(id, decoded)| match decoded {
414            Ok(value) => Some(value),
415            Err(error) => {
416                tracing::warn!(job_id = %id, error = %error, "Skipping an undecodable row");
417                None
418            }
419        })
420        .collect()
421}
422
423/// An archived job from a row of [`ARCHIVED_JOB_COLUMNS`], with the id already decoded.
424fn archived_job<R>(row: &R, id: uuid::Uuid) -> Result<ArchivedJob>
425where
426    R: Row,
427    for<'a> &'a str: sqlx::ColumnIndex<R>,
428    for<'a> String: sqlx::Decode<'a, R::Database> + sqlx::Type<R::Database>,
429    for<'a> Option<String>: sqlx::Decode<'a, R::Database> + sqlx::Type<R::Database>,
430    for<'a> Option<i32>: sqlx::Decode<'a, R::Database> + sqlx::Type<R::Database>,
431    for<'a> bool: sqlx::Decode<'a, R::Database> + sqlx::Type<R::Database>,
432    for<'a> DateTime<Utc>: sqlx::Decode<'a, R::Database> + sqlx::Type<R::Database>,
433{
434    Ok(ArchivedJob {
435        id,
436        queue_name: row.try_get("queue_name")?,
437        status: archived_status(&row.try_get::<String, _>("status")?),
438        created_at: row.try_get("created_at")?,
439        archived_at: row.try_get("archived_at")?,
440        archival_reason: ArchivalReason::parse_from_db(
441            &row.try_get::<String, _>("archival_reason")?,
442        )
443        .unwrap_or_default(),
444        original_payload_size: row
445            .try_get::<Option<i32>, _>("original_payload_size")?
446            .and_then(|size| usize::try_from(size).ok()),
447        payload_compressed: row.try_get("payload_compressed")?,
448        archived_by: row.try_get("archived_by")?,
449    })
450}
451
452/// Filters shared by count and delete: `$n`/`?` placeholders, queue then cutoff.
453fn finished_filter(kind: FinishedKind, queue: bool, cutoff: bool, postgres: bool) -> String {
454    let mut n = 0;
455    let mut ph = || {
456        n += 1;
457        if postgres {
458            format!("${n}")
459        } else {
460            "?".to_string()
461        }
462    };
463    let mut sql = format!("WHERE status = '{}'", kind.status());
464    if queue {
465        sql.push_str(&format!(" AND queue_name = {}", ph()));
466    }
467    if cutoff {
468        sql.push_str(&format!(" AND {} < {}", kind.age_column(), ph()));
469    }
470    sql
471}
472
473impl JobHistory for JobQueue<sqlx::Postgres> {
474    async fn hourly_activity(
475        &self,
476        queue: Option<&str>,
477        start: DateTime<Utc>,
478        end: DateTime<Utc>,
479    ) -> Result<Vec<HourBucket>> {
480        let hours = hour_range(start, end);
481        let (Some(first), Some(last)) = (hours.first().copied(), hours.last().copied()) else {
482            return Ok(Vec::new());
483        };
484        let upper = last + Duration::hours(1);
485        let queue_clause = if queue.is_some() {
486            " AND queue_name = $3"
487        } else {
488            ""
489        };
490
491        let completed_sql = format!(
492            "SELECT date_trunc('hour', completed_at AT TIME ZONE 'UTC') AS bucket, \
493             COUNT(*)::bigint AS n, \
494             AVG(EXTRACT(EPOCH FROM (completed_at - started_at)) * 1000)::float8 AS avg_ms \
495             FROM hammerwork_jobs WHERE status = 'Completed' \
496             AND completed_at >= $1 AND completed_at < $2{queue_clause} GROUP BY 1"
497        );
498        let mut q = sqlx::query(&completed_sql).bind(first).bind(upper);
499        if let Some(name) = queue {
500            q = q.bind(name.to_string());
501        }
502        let mut completed = BTreeMap::new();
503        for row in q.fetch_all(&self.pool).await? {
504            let bucket: NaiveDateTime = row.try_get("bucket")?;
505            let n: i64 = row.try_get("n")?;
506            let avg: Option<f64> = row.try_get("avg_ms")?;
507            completed.insert(bucket.and_utc(), (to_u64(n), avg));
508        }
509
510        let failed_sql = format!(
511            "SELECT date_trunc('hour', COALESCE(failed_at, timed_out_at) AT TIME ZONE 'UTC') AS bucket, \
512             COUNT(*)::bigint AS n FROM hammerwork_jobs \
513             WHERE status IN ('Failed', 'Dead', 'TimedOut') \
514             AND COALESCE(failed_at, timed_out_at) >= $1 AND COALESCE(failed_at, timed_out_at) < $2{queue_clause} \
515             GROUP BY 1"
516        );
517        let mut q = sqlx::query(&failed_sql).bind(first).bind(upper);
518        if let Some(name) = queue {
519            q = q.bind(name.to_string());
520        }
521        let mut failed = BTreeMap::new();
522        for row in q.fetch_all(&self.pool).await? {
523            let bucket: NaiveDateTime = row.try_get("bucket")?;
524            let n: i64 = row.try_get("n")?;
525            failed.insert(bucket.and_utc(), to_u64(n));
526        }
527
528        Ok(fill_buckets(&hours, &completed, &failed))
529    }
530
531    async fn error_groups(&self, since: DateTime<Utc>, limit: u32) -> Result<Vec<ErrorGroup>> {
532        let rows = sqlx::query(
533            "SELECT queue_name, error_message, COUNT(*)::bigint AS n, \
534             MIN(COALESCE(failed_at, timed_out_at)) AS first_seen, \
535             MAX(COALESCE(failed_at, timed_out_at)) AS last_seen \
536             FROM hammerwork_jobs \
537             WHERE status IN ('Failed', 'Dead', 'TimedOut') AND error_message IS NOT NULL \
538             AND COALESCE(failed_at, timed_out_at) >= $1 \
539             GROUP BY queue_name, error_message ORDER BY n DESC LIMIT $2",
540        )
541        .bind(since)
542        .bind(i64::from(limit))
543        .fetch_all(&self.pool)
544        .await?;
545        let mut groups = Vec::with_capacity(rows.len());
546        for row in rows {
547            groups.push(ErrorGroup {
548                queue_name: row.try_get("queue_name")?,
549                message: row.try_get("error_message")?,
550                count: to_u64(row.try_get("n")?),
551                first_seen: row.try_get("first_seen")?,
552                last_seen: row.try_get("last_seen")?,
553            });
554        }
555        Ok(groups)
556    }
557
558    async fn count_jobs(
559        &self,
560        queue: Option<&str>,
561        kind: FinishedKind,
562        older_than: Option<DateTime<Utc>>,
563    ) -> Result<u64> {
564        let sql = format!(
565            "SELECT COUNT(*)::bigint FROM hammerwork_jobs {}",
566            finished_filter(kind, queue.is_some(), older_than.is_some(), true)
567        );
568        let mut q = sqlx::query_scalar::<_, i64>(&sql);
569        if let Some(name) = queue {
570            q = q.bind(name.to_string());
571        }
572        if let Some(cutoff) = older_than {
573            q = q.bind(cutoff);
574        }
575        Ok(to_u64(q.fetch_one(&self.pool).await?))
576    }
577
578    async fn delete_jobs(
579        &self,
580        queue: Option<&str>,
581        kind: FinishedKind,
582        older_than: Option<DateTime<Utc>>,
583    ) -> Result<u64> {
584        let sql = format!(
585            "DELETE FROM hammerwork_jobs {}",
586            finished_filter(kind, queue.is_some(), older_than.is_some(), true)
587        );
588        let mut q = sqlx::query(&sql);
589        if let Some(name) = queue {
590            q = q.bind(name.to_string());
591        }
592        if let Some(cutoff) = older_than {
593            q = q.bind(cutoff);
594        }
595        Ok(q.execute(&self.pool).await?.rows_affected())
596    }
597
598    async fn archived_jobs(
599        &self,
600        filter: &ArchiveFilterParams,
601        limit: u32,
602        offset: u32,
603    ) -> Result<(Vec<ArchivedJob>, u64)> {
604        let (conditions, binds) = archive_conditions(filter, true);
605        let count_sql = format!("SELECT COUNT(*) FROM hammerwork_jobs_archive {conditions}");
606        let total = bind_archive_filter!(sqlx::query_scalar::<_, i64>(&count_sql), &binds)
607            .fetch_one(&self.pool)
608            .await?;
609        let (limit_ph, offset_ph) = (
610            format!("${}", binds.len() + 1),
611            format!("${}", binds.len() + 2),
612        );
613        let page_sql = format!(
614            "SELECT {ARCHIVED_JOB_COLUMNS} FROM hammerwork_jobs_archive {conditions} \
615             ORDER BY archived_at DESC, id DESC LIMIT {limit_ph} OFFSET {offset_ph}"
616        );
617        let rows = bind_archive_filter!(sqlx::query(&page_sql), &binds)
618            .bind(i64::from(limit))
619            .bind(i64::from(offset))
620            .fetch_all(&self.pool)
621            .await?;
622        let jobs = decodable(rows.iter().map(|row| {
623            let id = row.try_get::<uuid::Uuid, _>("id");
624            let label = id.as_ref().map(|id| id.to_string()).unwrap_or_default();
625            (
626                label,
627                id.map_err(Into::into).and_then(|id| archived_job(row, id)),
628            )
629        }));
630        Ok((jobs, to_u64(total)))
631    }
632
633    async fn database_now(&self) -> Result<DateTime<Utc>> {
634        Ok(sqlx::query_scalar("SELECT NOW()")
635            .fetch_one(&self.pool)
636            .await?)
637    }
638
639    async fn recent_job_changes(&self, since: DateTime<Utc>, limit: u32) -> Result<Vec<JobUpdate>> {
640        // GREATEST ignores NULLs on PostgreSQL.
641        let sql = format!(
642            "SELECT id, queue_name, status, priority, attempts, GREATEST({}) AS changed_at \
643             FROM hammerwork_jobs WHERE {} ORDER BY changed_at DESC LIMIT $2",
644            CHANGE_COLUMNS.join(", "),
645            changed_since_filter(true)
646        );
647        let rows = sqlx::query(&sql)
648            .bind(since)
649            .bind(i64::from(limit))
650            .fetch_all(&self.pool)
651            .await?;
652        Ok(decodable(rows.iter().map(|row| {
653            let id = row.try_get::<uuid::Uuid, _>("id");
654            let label = id.as_ref().map(|id| id.to_string()).unwrap_or_default();
655            let change = id.map_err(Into::into).and_then(|id| {
656                Ok(job_change(
657                    id.to_string(),
658                    row.try_get("queue_name")?,
659                    &row.try_get::<String, _>("status")?,
660                    row.try_get("priority")?,
661                    row.try_get("attempts")?,
662                    row.try_get("changed_at")?,
663                ))
664            });
665            (label, change)
666        })))
667    }
668
669    async fn table_maintenance(
670        &self,
671        operation: TableMaintenance,
672        tables: &[&'static str],
673        dry_run: bool,
674    ) -> Result<Vec<MaintenanceStatement>> {
675        let names: Vec<String> = tables.iter().map(|t| t.to_string()).collect();
676        let existing: Vec<String> = sqlx::query_scalar(
677            "SELECT tablename::text FROM pg_tables \
678             WHERE schemaname = current_schema() AND tablename = ANY($1)",
679        )
680        .bind(&names)
681        .fetch_all(&self.pool)
682        .await?;
683        let mut ran = Vec::new();
684        for table in present_tables(tables, &existing) {
685            let statement = operation.postgres_statement(table);
686            if !dry_run {
687                // Simple-query protocol: VACUUM may not run inside a transaction block.
688                sqlx::raw_sql(&statement).execute(&self.pool).await?;
689            }
690            ran.push(MaintenanceStatement {
691                table,
692                statement,
693                messages: Vec::new(),
694            });
695        }
696        Ok(ran)
697    }
698}
699
700impl JobHistory for JobQueue<sqlx::MySql> {
701    async fn hourly_activity(
702        &self,
703        queue: Option<&str>,
704        start: DateTime<Utc>,
705        end: DateTime<Utc>,
706    ) -> Result<Vec<HourBucket>> {
707        let hours = hour_range(start, end);
708        let (Some(first), Some(last)) = (hours.first().copied(), hours.last().copied()) else {
709            return Ok(Vec::new());
710        };
711        let upper = last + Duration::hours(1);
712        let queue_clause = if queue.is_some() {
713            " AND queue_name = ?"
714        } else {
715            ""
716        };
717
718        let completed_sql = format!(
719            "SELECT DATE_FORMAT(completed_at, '%Y-%m-%d %H:00:00') AS bucket, \
720             CAST(COUNT(*) AS SIGNED) AS n, \
721             CAST(AVG(TIMESTAMPDIFF(MICROSECOND, started_at, completed_at)) / 1000 AS DOUBLE) AS avg_ms \
722             FROM hammerwork_jobs WHERE status = 'Completed' \
723             AND completed_at >= ? AND completed_at < ?{queue_clause} GROUP BY bucket"
724        );
725        let mut q = sqlx::query(&completed_sql).bind(first).bind(upper);
726        if let Some(name) = queue {
727            q = q.bind(name.to_string());
728        }
729        let mut completed = BTreeMap::new();
730        for row in q.fetch_all(&self.pool).await? {
731            let bucket: String = row.try_get("bucket")?;
732            let n: i64 = row.try_get("n")?;
733            let avg: Option<f64> = row.try_get("avg_ms")?;
734            if let Some(hour) = parse_bucket(&bucket) {
735                completed.insert(hour, (to_u64(n), avg));
736            }
737        }
738
739        let failed_sql = format!(
740            "SELECT DATE_FORMAT(COALESCE(failed_at, timed_out_at), '%Y-%m-%d %H:00:00') AS bucket, \
741             CAST(COUNT(*) AS SIGNED) AS n FROM hammerwork_jobs \
742             WHERE status IN ('Failed', 'Dead', 'TimedOut') \
743             AND COALESCE(failed_at, timed_out_at) >= ? AND COALESCE(failed_at, timed_out_at) < ?{queue_clause} \
744             GROUP BY bucket"
745        );
746        let mut q = sqlx::query(&failed_sql).bind(first).bind(upper);
747        if let Some(name) = queue {
748            q = q.bind(name.to_string());
749        }
750        let mut failed = BTreeMap::new();
751        for row in q.fetch_all(&self.pool).await? {
752            let bucket: String = row.try_get("bucket")?;
753            let n: i64 = row.try_get("n")?;
754            if let Some(hour) = parse_bucket(&bucket) {
755                failed.insert(hour, to_u64(n));
756            }
757        }
758
759        Ok(fill_buckets(&hours, &completed, &failed))
760    }
761
762    async fn error_groups(&self, since: DateTime<Utc>, limit: u32) -> Result<Vec<ErrorGroup>> {
763        let rows = sqlx::query(
764            "SELECT queue_name, error_message, CAST(COUNT(*) AS SIGNED) AS n, \
765             MIN(COALESCE(failed_at, timed_out_at)) AS first_seen, \
766             MAX(COALESCE(failed_at, timed_out_at)) AS last_seen \
767             FROM hammerwork_jobs \
768             WHERE status IN ('Failed', 'Dead', 'TimedOut') AND error_message IS NOT NULL \
769             AND COALESCE(failed_at, timed_out_at) >= ? \
770             GROUP BY queue_name, error_message ORDER BY n DESC LIMIT ?",
771        )
772        .bind(since)
773        .bind(i64::from(limit))
774        .fetch_all(&self.pool)
775        .await?;
776        let mut groups = Vec::with_capacity(rows.len());
777        for row in rows {
778            groups.push(ErrorGroup {
779                queue_name: row.try_get("queue_name")?,
780                message: row.try_get("error_message")?,
781                count: to_u64(row.try_get("n")?),
782                first_seen: row.try_get("first_seen")?,
783                last_seen: row.try_get("last_seen")?,
784            });
785        }
786        Ok(groups)
787    }
788
789    async fn count_jobs(
790        &self,
791        queue: Option<&str>,
792        kind: FinishedKind,
793        older_than: Option<DateTime<Utc>>,
794    ) -> Result<u64> {
795        let sql = format!(
796            "SELECT CAST(COUNT(*) AS SIGNED) FROM hammerwork_jobs {}",
797            finished_filter(kind, queue.is_some(), older_than.is_some(), false)
798        );
799        let mut q = sqlx::query_scalar::<_, i64>(&sql);
800        if let Some(name) = queue {
801            q = q.bind(name.to_string());
802        }
803        if let Some(cutoff) = older_than {
804            q = q.bind(cutoff);
805        }
806        Ok(to_u64(q.fetch_one(&self.pool).await?))
807    }
808
809    async fn delete_jobs(
810        &self,
811        queue: Option<&str>,
812        kind: FinishedKind,
813        older_than: Option<DateTime<Utc>>,
814    ) -> Result<u64> {
815        let sql = format!(
816            "DELETE FROM hammerwork_jobs {}",
817            finished_filter(kind, queue.is_some(), older_than.is_some(), false)
818        );
819        let mut q = sqlx::query(&sql);
820        if let Some(name) = queue {
821            q = q.bind(name.to_string());
822        }
823        if let Some(cutoff) = older_than {
824            q = q.bind(cutoff);
825        }
826        Ok(q.execute(&self.pool).await?.rows_affected())
827    }
828
829    async fn archived_jobs(
830        &self,
831        filter: &ArchiveFilterParams,
832        limit: u32,
833        offset: u32,
834    ) -> Result<(Vec<ArchivedJob>, u64)> {
835        let (conditions, binds) = archive_conditions(filter, false);
836        let count_sql = format!("SELECT COUNT(*) FROM hammerwork_jobs_archive {conditions}");
837        let total = bind_archive_filter!(sqlx::query_scalar::<_, i64>(&count_sql), &binds)
838            .fetch_one(&self.pool)
839            .await?;
840        let page_sql = format!(
841            "SELECT {ARCHIVED_JOB_COLUMNS} FROM hammerwork_jobs_archive {conditions} \
842             ORDER BY archived_at DESC, id DESC LIMIT ? OFFSET ?"
843        );
844        let rows = bind_archive_filter!(sqlx::query(&page_sql), &binds)
845            .bind(i64::from(limit))
846            .bind(i64::from(offset))
847            .fetch_all(&self.pool)
848            .await?;
849        let jobs = decodable(rows.iter().map(|row| {
850            let label = row.try_get::<String, _>("id").unwrap_or_default();
851            let job = uuid::Uuid::parse_str(&label)
852                .map_err(Into::into)
853                .and_then(|id| archived_job(row, id));
854            (label, job)
855        }));
856        Ok((jobs, to_u64(total)))
857    }
858
859    async fn database_now(&self) -> Result<DateTime<Utc>> {
860        Ok(sqlx::query_scalar("SELECT UTC_TIMESTAMP(6)")
861            .fetch_one(&self.pool)
862            .await?)
863    }
864
865    async fn recent_job_changes(&self, since: DateTime<Utc>, limit: u32) -> Result<Vec<JobUpdate>> {
866        // GREATEST is NULL if any argument is on MySQL, so missing timestamps fall back to
867        // `created_at`.
868        let latest = CHANGE_COLUMNS
869            .iter()
870            .map(|column| format!("COALESCE({column}, created_at)"))
871            .collect::<Vec<_>>()
872            .join(", ");
873        let sql = format!(
874            "SELECT id, queue_name, status, priority, attempts, \
875             CAST(GREATEST({latest}) AS DATETIME(6)) AS changed_at \
876             FROM hammerwork_jobs WHERE {} ORDER BY changed_at DESC LIMIT ?",
877            changed_since_filter(false)
878        );
879        let mut query = sqlx::query(&sql);
880        for _ in CHANGE_COLUMNS {
881            query = query.bind(since);
882        }
883        let rows = query.bind(i64::from(limit)).fetch_all(&self.pool).await?;
884        Ok(decodable(rows.iter().map(|row| {
885            let label = row.try_get::<String, _>("id").unwrap_or_default();
886            let change = (|| {
887                Ok(job_change(
888                    row.try_get("id")?,
889                    row.try_get("queue_name")?,
890                    &row.try_get::<String, _>("status")?,
891                    row.try_get("priority")?,
892                    row.try_get("attempts")?,
893                    row.try_get("changed_at")?,
894                ))
895            })();
896            (label, change)
897        })))
898    }
899    async fn table_maintenance(
900        &self,
901        operation: TableMaintenance,
902        tables: &[&'static str],
903        dry_run: bool,
904    ) -> Result<Vec<MaintenanceStatement>> {
905        if tables.is_empty() {
906            return Ok(Vec::new());
907        }
908        let placeholders = vec!["?"; tables.len()].join(", ");
909        let sql = format!(
910            "SELECT CAST(TABLE_NAME AS CHAR) FROM information_schema.TABLES \
911             WHERE TABLE_SCHEMA = DATABASE() AND TABLE_NAME IN ({placeholders})"
912        );
913        let mut q = sqlx::query_scalar::<_, String>(&sql);
914        for table in tables {
915            q = q.bind(*table);
916        }
917        let existing = q.fetch_all(&self.pool).await?;
918        let mut ran = Vec::new();
919        for table in present_tables(tables, &existing) {
920            let statement = operation.mysql_statement(table);
921            let mut messages = Vec::new();
922            if !dry_run {
923                // OPTIMIZE TABLE reports failures as result rows, not as SQL errors.
924                for row in sqlx::raw_sql(&statement).fetch_all(&self.pool).await? {
925                    let kind = text_column(&row, "Msg_type");
926                    let text = text_column(&row, "Msg_text");
927                    if kind.eq_ignore_ascii_case("error") {
928                        return Err(hammerwork::HammerworkError::Queue {
929                            message: format!("{statement} failed: {text}"),
930                        });
931                    }
932                    messages.push(format!("{kind}: {text}"));
933                }
934            }
935            ran.push(MaintenanceStatement {
936                table,
937                statement,
938                messages,
939            });
940        }
941        Ok(ran)
942    }
943}
944
945/// A text column of a MySQL admin statement's result, which may arrive as binary.
946fn text_column(row: &sqlx::mysql::MySqlRow, name: &str) -> String {
947    row.try_get::<String, _>(name)
948        .or_else(|_| {
949            row.try_get::<Vec<u8>, _>(name)
950                .map(|bytes| String::from_utf8_lossy(&bytes).into_owned())
951        })
952        .unwrap_or_default()
953}
954
955fn parse_bucket(bucket: &str) -> Option<DateTime<Utc>> {
956    NaiveDateTime::parse_from_str(bucket, "%Y-%m-%d %H:%M:%S")
957        .ok()
958        .map(|t| t.and_utc())
959}
960
961#[cfg(test)]
962mod tests {
963    use super::*;
964
965    fn at(h: u32, m: u32) -> DateTime<Utc> {
966        Utc.with_ymd_and_hms(2026, 1, 2, h, m, 0).unwrap()
967    }
968
969    /// #71: a dashboard listing skips the rows it cannot decode.
970    #[test]
971    fn undecodable_rows_are_skipped() {
972        let bad: Result<u8> = Err(hammerwork::HammerworkError::Processing("bad".into()));
973        let rows = vec![("a".to_string(), Ok(1u8)), ("b".to_string(), bad)];
974        assert_eq!(decodable(rows), [1]);
975    }
976
977    #[test]
978    fn job_changes_normalise_status_and_priority() {
979        let at = Utc::now();
980        let change = job_change("id".into(), "q".into(), "\"Running\"", 3, 2, at);
981        assert_eq!(change.status, "Running");
982        assert_eq!(change.priority, "high");
983        assert_eq!(change.attempts, 2);
984        assert_eq!(change.updated_at, at);
985        assert_eq!(
986            job_change("id".into(), "q".into(), "Pending", 99, 0, at).priority,
987            "normal"
988        );
989        assert_eq!(
990            changed_since_filter(true),
991            "created_at > $1 OR started_at > $1 OR completed_at > $1 OR failed_at > $1 \
992             OR timed_out_at > $1"
993        );
994        assert!(changed_since_filter(false).starts_with("created_at > ? OR started_at > ?"));
995    }
996
997    #[test]
998    fn test_hour_floor_and_range() {
999        assert_eq!(
1000            hour_floor(at(5, 59)),
1001            Utc.with_ymd_and_hms(2026, 1, 2, 5, 0, 0).unwrap()
1002        );
1003        let hours = hour_range(at(5, 10), at(8, 0));
1004        assert_eq!(hours.len(), 3);
1005        assert_eq!(hours[2], Utc.with_ymd_and_hms(2026, 1, 2, 7, 0, 0).unwrap());
1006        assert!(hour_range(at(8, 0), at(5, 0)).is_empty());
1007    }
1008
1009    #[test]
1010    fn test_hour_range_is_capped() {
1011        let end = at(0, 0);
1012        let hours = hour_range(end - Duration::days(400), end);
1013        assert_eq!(hours.len() as i64, MAX_BUCKETS);
1014    }
1015
1016    #[test]
1017    fn test_fill_buckets_zero_fills() {
1018        let hours = hour_range(at(5, 0), at(8, 0));
1019        let mut completed = BTreeMap::new();
1020        completed.insert(hours[0], (4, Some(12.5)));
1021        let mut failed = BTreeMap::new();
1022        failed.insert(hours[2], 2);
1023        let buckets = fill_buckets(&hours, &completed, &failed);
1024        assert_eq!(buckets.len(), 3);
1025        assert_eq!((buckets[0].completed, buckets[0].failed), (4, 0));
1026        assert_eq!(buckets[0].avg_processing_time_ms, Some(12.5));
1027        assert_eq!((buckets[1].completed, buckets[1].failed), (0, 0));
1028        assert_eq!(buckets[1].avg_processing_time_ms, None);
1029        assert_eq!((buckets[2].completed, buckets[2].failed), (0, 2));
1030    }
1031
1032    /// M11: archive filters become SQL, so the listing never reads the whole archive.
1033    #[test]
1034    fn archive_filters_become_sql_conditions() {
1035        let (sql, binds) = archive_conditions(&ArchiveFilterParams::default(), true);
1036        assert_eq!(sql, "");
1037        assert!(binds.is_empty());
1038
1039        let at = Utc.with_ymd_and_hms(2024, 5, 1, 0, 0, 0).unwrap();
1040        let filter = ArchiveFilterParams {
1041            queue: Some("emails".into()),
1042            reason: Some("Manual".into()),
1043            archived_after: Some(at),
1044            archived_before: Some(at),
1045            archived_by: Some("ops".into()),
1046            compressed: Some(true),
1047            original_status: Some("COMPLETED".into()),
1048        };
1049        let (pg, binds) = archive_conditions(&filter, true);
1050        assert_eq!(
1051            pg,
1052            "WHERE queue_name = $1 AND LOWER(TRIM(BOTH '\"' FROM archival_reason)) = $2 \
1053             AND archived_at >= $3 AND archived_at <= $4 AND archived_by = $5 \
1054             AND payload_compressed = $6 AND LOWER(TRIM(BOTH '\"' FROM status)) = $7"
1055        );
1056        assert_eq!(
1057            binds,
1058            vec![
1059                ArchiveBind::Text("emails".into()),
1060                ArchiveBind::Text("manual".into()),
1061                ArchiveBind::Time(at),
1062                ArchiveBind::Time(at),
1063                ArchiveBind::Text("ops".into()),
1064                ArchiveBind::Bool(true),
1065                ArchiveBind::Text("completed".into()),
1066            ]
1067        );
1068        let (mysql, _) = archive_conditions(&filter, false);
1069        assert_eq!(mysql.matches('?').count(), 7);
1070        assert!(!mysql.contains('$'));
1071    }
1072
1073    #[test]
1074    fn archived_statuses_parse_like_the_core_library() {
1075        assert_eq!(archived_status("Completed"), JobStatus::Completed);
1076        assert_eq!(archived_status("\"TimedOut\""), JobStatus::TimedOut);
1077        assert_eq!(archived_status("Archived"), JobStatus::Archived);
1078        assert_eq!(archived_status("whatever"), JobStatus::Dead);
1079    }
1080
1081    #[test]
1082    fn test_finished_filter_placeholders() {
1083        let pg = finished_filter(FinishedKind::Completed, true, true, true);
1084        assert_eq!(
1085            pg,
1086            "WHERE status = 'Completed' AND queue_name = $1 AND completed_at < $2"
1087        );
1088        let my = finished_filter(FinishedKind::Dead, false, true, false);
1089        assert!(my.contains("status = 'Dead'"));
1090        assert!(my.ends_with("< ?"));
1091        assert!(!my.contains("queue_name"));
1092    }
1093
1094    #[test]
1095    fn test_parse_bucket() {
1096        assert_eq!(
1097            parse_bucket("2026-01-02 05:00:00"),
1098            Some(Utc.with_ymd_and_hms(2026, 1, 2, 5, 0, 0).unwrap())
1099        );
1100        assert_eq!(parse_bucket("garbage"), None);
1101    }
1102
1103    #[test]
1104    fn test_table_maintenance_statements() {
1105        assert_eq!(
1106            TableMaintenance::parse("vacuum"),
1107            Some(TableMaintenance::Vacuum)
1108        );
1109        assert_eq!(
1110            TableMaintenance::parse("reindex"),
1111            Some(TableMaintenance::Reindex)
1112        );
1113        assert_eq!(
1114            TableMaintenance::parse("optimize"),
1115            Some(TableMaintenance::Optimize)
1116        );
1117        assert_eq!(TableMaintenance::parse("VACUUM"), None);
1118        assert_eq!(TableMaintenance::parse("cleanup"), None);
1119
1120        let t = "hammerwork_jobs";
1121        assert_eq!(
1122            TableMaintenance::Vacuum.postgres_statement(t),
1123            "VACUUM (ANALYZE) hammerwork_jobs"
1124        );
1125        assert_eq!(
1126            TableMaintenance::Reindex.postgres_statement(t),
1127            "REINDEX TABLE hammerwork_jobs"
1128        );
1129        assert_eq!(
1130            TableMaintenance::Optimize.postgres_statement(t),
1131            "ANALYZE hammerwork_jobs"
1132        );
1133        for op in [
1134            TableMaintenance::Vacuum,
1135            TableMaintenance::Reindex,
1136            TableMaintenance::Optimize,
1137        ] {
1138            assert_eq!(op.mysql_statement(t), "OPTIMIZE TABLE hammerwork_jobs");
1139        }
1140        assert_eq!(
1141            serde_json::to_value(TableMaintenance::Reindex).unwrap(),
1142            "reindex"
1143        );
1144    }
1145
1146    #[test]
1147    fn test_maintenance_tables_only_hammerwork_tables() {
1148        assert_eq!(maintenance_tables(None).unwrap(), HAMMERWORK_TABLES);
1149        assert_eq!(
1150            maintenance_tables(Some("hammerwork_jobs_archive")).unwrap(),
1151            vec!["hammerwork_jobs_archive"]
1152        );
1153        for bad in [
1154            "users",
1155            "hammerwork_jobs; DROP TABLE x",
1156            "HAMMERWORK_JOBS",
1157            "",
1158        ] {
1159            assert!(maintenance_tables(Some(bad)).is_err(), "{bad}");
1160        }
1161        assert!(
1162            HAMMERWORK_TABLES
1163                .iter()
1164                .all(|t| t.starts_with("hammerwork_")
1165                    && t.chars().all(|c| c.is_ascii_lowercase() || c == '_'))
1166        );
1167    }
1168
1169    #[test]
1170    fn test_present_tables_keeps_fixed_names_in_order() {
1171        let existing = vec!["hammerwork_batches".to_string(), "other".to_string()];
1172        assert_eq!(
1173            present_tables(&["hammerwork_jobs", "hammerwork_batches"], &existing),
1174            vec!["hammerwork_batches"]
1175        );
1176        assert!(present_tables(&["hammerwork_jobs"], &[]).is_empty());
1177    }
1178}
1179
1180#[cfg(test)]
1181mod db_tests {
1182    use super::*;
1183    use hammerwork::JobQueue;
1184
1185    struct Seed {
1186        queue: String,
1187        status: &'static str,
1188        started: Option<DateTime<Utc>>,
1189        completed: Option<DateTime<Utc>>,
1190        failed: Option<DateTime<Utc>>,
1191        error: Option<&'static str>,
1192    }
1193
1194    fn seeds(queue: &str, other: &str, now: DateTime<Utc>) -> Vec<Seed> {
1195        let base = hour_floor(now);
1196        let done = |q: &str, at: DateTime<Utc>| Seed {
1197            queue: q.to_string(),
1198            status: "Completed",
1199            started: Some(at - Duration::milliseconds(100)),
1200            completed: Some(at),
1201            failed: None,
1202            error: None,
1203        };
1204        vec![
1205            done(queue, base - Duration::hours(2) + Duration::minutes(5)),
1206            done(queue, base - Duration::hours(2) + Duration::minutes(30)),
1207            done(other, base - Duration::hours(2) + Duration::minutes(10)),
1208            Seed {
1209                queue: queue.to_string(),
1210                status: "Dead",
1211                started: None,
1212                completed: None,
1213                failed: Some(base - Duration::hours(1) + Duration::minutes(15)),
1214                error: Some("connection refused"),
1215            },
1216        ]
1217    }
1218
1219    async fn assert_history<Q: JobHistory>(q: &Q, queue: &str, other: &str, now: DateTime<Utc>) {
1220        let base = hour_floor(now);
1221        let buckets = q
1222            .hourly_activity(Some(queue), base - Duration::hours(3), now)
1223            .await
1224            .unwrap();
1225        assert_eq!(buckets.len(), 4, "3h ago .. current hour, zero filled");
1226        assert_eq!((buckets[0].completed, buckets[0].failed), (0, 0));
1227        assert_eq!((buckets[1].completed, buckets[1].failed), (2, 0));
1228        let avg = buckets[1].avg_processing_time_ms.expect("avg");
1229        assert!((avg - 100.0).abs() < 5.0, "avg was {avg}");
1230        assert_eq!((buckets[2].completed, buckets[2].failed), (0, 1));
1231        assert_eq!(buckets[2].avg_processing_time_ms, None);
1232        assert_eq!(buckets[1].hour, base - Duration::hours(2));
1233
1234        // Without a queue filter the other queue's completion is included too.
1235        let all = q
1236            .hourly_activity(None, base - Duration::hours(3), now)
1237            .await
1238            .unwrap();
1239        assert!(all[1].completed >= 3);
1240
1241        let groups = q.error_groups(now - Duration::hours(5), 500).await.unwrap();
1242        let ours: Vec<_> = groups.iter().filter(|g| g.queue_name == queue).collect();
1243        assert_eq!(ours.len(), 1);
1244        assert_eq!(ours[0].message, "connection refused");
1245        assert_eq!(ours[0].count, 1);
1246        assert_eq!(
1247            ours[0].first_seen.timestamp(),
1248            (base - Duration::hours(1) + Duration::minutes(15)).timestamp()
1249        );
1250        assert_eq!(ours[0].last_seen, ours[0].first_seen);
1251
1252        // Counting and deleting are scoped to the queue and kind.
1253        let completed = FinishedKind::Completed;
1254        let dead = FinishedKind::Dead;
1255        assert_eq!(q.count_jobs(Some(queue), completed, None).await.unwrap(), 2);
1256        assert_eq!(q.count_jobs(Some(queue), dead, None).await.unwrap(), 1);
1257        let week_ago = Some(now - Duration::days(7));
1258        assert_eq!(q.count_jobs(Some(queue), dead, week_ago).await.unwrap(), 0);
1259        assert_eq!(
1260            q.delete_jobs(Some(queue), completed, None).await.unwrap(),
1261            2
1262        );
1263        assert_eq!(q.count_jobs(Some(queue), completed, None).await.unwrap(), 0);
1264        assert_eq!(q.count_jobs(Some(queue), dead, None).await.unwrap(), 1);
1265        assert_eq!(q.count_jobs(Some(other), completed, None).await.unwrap(), 1);
1266    }
1267
1268    #[tokio::test]
1269    #[ignore = "requires DATABASE_URL (PostgreSQL)"]
1270    async fn test_job_history_queries_postgres() {
1271        let url = std::env::var("DATABASE_URL").expect("DATABASE_URL");
1272        let pool = sqlx::PgPool::connect(&url).await.unwrap();
1273        let tag = uuid::Uuid::new_v4().simple().to_string();
1274        let (queue, other) = (format!("hist_a_{tag}"), format!("hist_b_{tag}"));
1275        let now = Utc::now();
1276        for s in seeds(&queue, &other, now) {
1277            sqlx::query(
1278                "INSERT INTO hammerwork_jobs (id, queue_name, payload, status, priority, attempts, \
1279                 max_attempts, created_at, scheduled_at, started_at, completed_at, failed_at, error_message) \
1280                 VALUES ($1, $2, '{}'::jsonb, $3, 2, 1, 3, NOW(), NOW(), $4, $5, $6, $7)",
1281            )
1282            .bind(uuid::Uuid::new_v4())
1283            .bind(&s.queue)
1284            .bind(s.status)
1285            .bind(s.started)
1286            .bind(s.completed)
1287            .bind(s.failed)
1288            .bind(s.error)
1289            .execute(&pool)
1290            .await
1291            .unwrap();
1292        }
1293        let jq = JobQueue::<sqlx::Postgres>::new(pool.clone());
1294        assert_history(&jq, &queue, &other, now).await;
1295        for name in [&queue, &other] {
1296            sqlx::query("DELETE FROM hammerwork_jobs WHERE queue_name = $1")
1297                .bind(name)
1298                .execute(&pool)
1299                .await
1300                .unwrap();
1301        }
1302    }
1303
1304    #[tokio::test]
1305    #[ignore = "requires MYSQL_DATABASE_URL"]
1306    async fn test_job_history_queries_mysql() {
1307        let url = std::env::var("MYSQL_DATABASE_URL").expect("MYSQL_DATABASE_URL");
1308        let pool = sqlx::MySqlPool::connect(&url).await.unwrap();
1309        let tag = uuid::Uuid::new_v4().simple().to_string();
1310        let (queue, other) = (format!("hist_a_{tag}"), format!("hist_b_{tag}"));
1311        let now = Utc::now();
1312        for s in seeds(&queue, &other, now) {
1313            sqlx::query(
1314                "INSERT INTO hammerwork_jobs (id, queue_name, payload, status, priority, attempts, \
1315                 max_attempts, created_at, scheduled_at, started_at, completed_at, failed_at, error_message) \
1316                 VALUES (?, ?, '{}', ?, 2, 1, 3, UTC_TIMESTAMP(6), UTC_TIMESTAMP(6), ?, ?, ?, ?)",
1317            )
1318            .bind(uuid::Uuid::new_v4().to_string())
1319            .bind(&s.queue)
1320            .bind(s.status)
1321            .bind(s.started)
1322            .bind(s.completed)
1323            .bind(s.failed)
1324            .bind(s.error)
1325            .execute(&pool)
1326            .await
1327            .unwrap();
1328        }
1329        let jq = JobQueue::<sqlx::MySql>::new(pool.clone());
1330        assert_history(&jq, &queue, &other, now).await;
1331        for name in [&queue, &other] {
1332            sqlx::query("DELETE FROM hammerwork_jobs WHERE queue_name = ?")
1333                .bind(name)
1334                .execute(&pool)
1335                .await
1336                .unwrap();
1337        }
1338    }
1339}