1use 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#[derive(Debug, Clone, PartialEq, Serialize)]
28pub struct HourBucket {
29 pub hour: DateTime<Utc>,
31 pub completed: u64,
32 pub failed: u64,
33 pub avg_processing_time_ms: Option<f64>,
35}
36
37#[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#[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 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#[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 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 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 pub fn mysql_statement(self, table: &'static str) -> String {
112 format!("OPTIMIZE TABLE {table}")
113 }
114}
115
116pub 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
130pub 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#[derive(Debug, Clone, PartialEq, Serialize)]
145pub struct MaintenanceStatement {
146 pub table: &'static str,
147 pub statement: String,
148 pub messages: Vec<String>,
150}
151
152fn 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
162pub const MAX_BUCKETS: i64 = 24 * 31;
164
165pub 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
171pub 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
182pub 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
206pub trait JobHistory: DatabaseQueue + Send + Sync {
208 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 fn error_groups(
218 &self,
219 since: DateTime<Utc>,
220 limit: u32,
221 ) -> impl Future<Output = Result<Vec<ErrorGroup>>> + Send;
222
223 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 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 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 fn database_now(&self) -> impl Future<Output = Result<DateTime<Utc>>> + Send;
250
251 fn recent_job_changes(
255 &self,
256 since: DateTime<Utc>,
257 limit: u32,
258 ) -> impl Future<Output = Result<Vec<JobUpdate>>> + Send;
259
260 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
271const CHANGE_COLUMNS: [&str; 5] = [
274 "created_at",
275 "started_at",
276 "completed_at",
277 "failed_at",
278 "timed_out_at",
279];
280
281fn 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
292fn 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#[derive(Debug, Clone, PartialEq)]
315enum ArchiveBind {
316 Text(String),
317 Time(DateTime<Utc>),
318 Bool(bool),
319}
320
321fn 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
370macro_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
385const ARCHIVED_JOB_COLUMNS: &str = "id, queue_name, status, created_at, archived_at, \
387 archival_reason, original_payload_size, payload_compressed, archived_by";
388
389fn 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
408pub(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
423fn 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
452fn 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 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 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 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 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
945fn 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 #[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 #[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 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 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}