Skip to main content

cargo_hammerwork/commands/
job.rs

1use anyhow::Result;
2use clap::Subcommand;
3use hammerwork::queue::DatabaseQueue;
4use hammerwork::{Job, JobPriority};
5use sqlx::Row;
6use tracing::info;
7
8use crate::config::Config;
9use crate::utils::database::DatabasePool;
10use crate::utils::display::JobTable;
11use crate::utils::validation::{validate_json_payload, validate_priority, validate_status};
12
13#[derive(Subcommand)]
14pub enum JobCommand {
15    #[command(about = "List jobs in the queue")]
16    List {
17        #[arg(short = 'u', long, help = "Database connection URL")]
18        database_url: Option<String>,
19        #[arg(short = 'n', long, help = "Queue name to filter by")]
20        queue: Option<String>,
21        #[arg(short = 't', long, help = "Job status to filter by")]
22        status: Option<String>,
23        #[arg(short = 'r', long, help = "Job priority to filter by")]
24        priority: Option<String>,
25        #[arg(short, long, help = "Maximum number of jobs to display")]
26        limit: Option<u32>,
27        #[arg(long, help = "Show only failed jobs")]
28        failed: bool,
29        #[arg(long, help = "Show only completed jobs")]
30        completed: bool,
31        #[arg(long, help = "Show jobs from last N hours")]
32        last_hours: Option<u32>,
33    },
34    #[command(about = "Show details of a specific job")]
35    Show {
36        #[arg(short = 'u', long, help = "Database connection URL")]
37        database_url: Option<String>,
38        #[arg(help = "Job ID")]
39        job_id: String,
40    },
41    #[command(about = "Enqueue a new job")]
42    Enqueue {
43        #[arg(short = 'u', long, help = "Database connection URL")]
44        database_url: Option<String>,
45        #[arg(short = 'n', long, help = "Queue name")]
46        queue: String,
47        #[arg(short = 'j', long, help = "Job payload as JSON")]
48        payload: String,
49        #[arg(short = 'r', long, help = "Job priority")]
50        priority: Option<String>,
51        #[arg(long, help = "Delay in seconds before job becomes available")]
52        delay: Option<u64>,
53        #[arg(long, help = "Maximum number of retry attempts")]
54        max_attempts: Option<u32>,
55        #[arg(long, help = "Timeout in seconds")]
56        timeout: Option<u32>,
57    },
58    #[command(about = "Retry failed jobs")]
59    Retry {
60        #[arg(short = 'u', long, help = "Database connection URL")]
61        database_url: Option<String>,
62        #[arg(long, help = "Specific job ID to retry")]
63        job_id: Option<String>,
64        #[arg(short = 'n', long, help = "Queue name to retry all failed jobs")]
65        queue: Option<String>,
66        #[arg(long, help = "Retry all failed jobs")]
67        all: bool,
68    },
69    #[command(about = "Cancel/delete jobs")]
70    Cancel {
71        #[arg(short = 'u', long, help = "Database connection URL")]
72        database_url: Option<String>,
73        #[arg(long, help = "Specific job ID to cancel")]
74        job_id: Option<String>,
75        #[arg(short = 'n', long, help = "Queue name to cancel pending jobs")]
76        queue: Option<String>,
77        #[arg(long, help = "Cancel all pending jobs")]
78        all_pending: bool,
79    },
80    #[command(about = "Purge completed or dead jobs")]
81    Purge {
82        #[arg(short = 'u', long, help = "Database connection URL")]
83        database_url: Option<String>,
84        #[arg(short, long, help = "Queue name to filter by")]
85        queue: Option<String>,
86        #[arg(long, help = "Only purge completed jobs")]
87        completed: bool,
88        #[arg(long, help = "Only purge dead jobs")]
89        dead: bool,
90        #[arg(long, help = "Only purge failed jobs")]
91        failed: bool,
92        #[arg(long, help = "Purge jobs older than N days")]
93        older_than_days: Option<u32>,
94        #[arg(long, help = "Confirm the purge operation")]
95        confirm: bool,
96    },
97}
98
99impl JobCommand {
100    pub async fn execute(&self, config: &Config) -> Result<()> {
101        let db_url = self.get_database_url(config)?;
102        let pool = DatabasePool::connect(&db_url, config.get_connection_pool_size()).await?;
103
104        match self {
105            JobCommand::List {
106                queue,
107                status,
108                priority,
109                limit,
110                failed,
111                completed,
112                last_hours,
113                ..
114            } => {
115                list_jobs(
116                    pool,
117                    queue.clone(),
118                    status.clone(),
119                    priority.clone(),
120                    limit.unwrap_or(config.get_default_limit()),
121                    *failed,
122                    *completed,
123                    *last_hours,
124                )
125                .await?;
126            }
127            JobCommand::Show { job_id, .. } => {
128                show_job_details(pool, job_id).await?;
129            }
130            JobCommand::Enqueue {
131                queue,
132                payload,
133                priority,
134                delay,
135                max_attempts,
136                timeout,
137                ..
138            } => {
139                enqueue_job(
140                    pool,
141                    queue,
142                    payload,
143                    priority,
144                    *delay,
145                    *max_attempts,
146                    *timeout,
147                )
148                .await?;
149            }
150            JobCommand::Retry {
151                job_id, queue, all, ..
152            } => {
153                retry_jobs(pool, job_id.clone(), queue.clone(), *all).await?;
154            }
155            JobCommand::Cancel {
156                job_id,
157                queue,
158                all_pending,
159                ..
160            } => {
161                cancel_jobs(pool, job_id.clone(), queue.clone(), *all_pending).await?;
162            }
163            JobCommand::Purge {
164                queue,
165                completed,
166                dead,
167                failed,
168                older_than_days,
169                confirm,
170                ..
171            } => {
172                purge_jobs(
173                    pool,
174                    queue.clone(),
175                    *completed,
176                    *dead,
177                    *failed,
178                    *older_than_days,
179                    *confirm,
180                )
181                .await?;
182            }
183        }
184        Ok(())
185    }
186
187    fn get_database_url(&self, config: &Config) -> Result<String> {
188        let url_option = match self {
189            JobCommand::List { database_url, .. } => database_url,
190            JobCommand::Show { database_url, .. } => database_url,
191            JobCommand::Enqueue { database_url, .. } => database_url,
192            JobCommand::Retry { database_url, .. } => database_url,
193            JobCommand::Cancel { database_url, .. } => database_url,
194            JobCommand::Purge { database_url, .. } => database_url,
195        };
196
197        url_option
198            .as_ref()
199            .map(|s| s.as_str())
200            .or(config.get_database_url())
201            .map(|s| s.to_string())
202            .ok_or_else(|| anyhow::anyhow!("Database URL is required"))
203    }
204}
205
206#[allow(clippy::too_many_arguments)]
207async fn list_jobs(
208    pool: DatabasePool,
209    queue: Option<String>,
210    status: Option<String>,
211    priority: Option<String>,
212    limit: u32,
213    failed: bool,
214    completed: bool,
215    last_hours: Option<u32>,
216) -> Result<()> {
217    // Validate inputs
218    if let Some(ref s) = status {
219        validate_status(s)?;
220    }
221    if let Some(ref p) = priority {
222        validate_priority(p)?;
223    }
224
225    // Build dynamic query conditions
226    let mut conditions = Vec::new();
227    // Dynamic query building for complex filtering
228
229    if let Some(queue_name) = &queue {
230        // Escape single quotes to prevent SQL injection
231        let escaped_queue = queue_name.replace("'", "''");
232        conditions.push(format!("queue_name = '{}'", escaped_queue));
233    }
234
235    if let Some(status_str) = &status {
236        let escaped_status = status_str.replace("'", "''");
237        conditions.push(format!("status = '{}'", escaped_status));
238    }
239
240    if let Some(priority_str) = &priority {
241        let escaped_priority = priority_str.replace("'", "''");
242        conditions.push(format!("priority = '{}'", escaped_priority));
243    }
244
245    if failed {
246        conditions.push("status IN ('failed', 'dead')".to_string());
247    }
248
249    if completed {
250        conditions.push("status = 'completed'".to_string());
251    }
252
253    let mut job_table = JobTable::new();
254
255    match pool {
256        DatabasePool::Postgres(pg_pool) => {
257            // Build PostgreSQL query
258            let mut query = "SELECT id, queue_name, status, priority, attempts, created_at, scheduled_at FROM hammerwork_jobs".to_string();
259
260            if !conditions.is_empty() {
261                query.push_str(&format!(" WHERE {}", conditions.join(" AND ")));
262            }
263
264            if let Some(hours) = last_hours {
265                if conditions.is_empty() {
266                    query.push_str(&format!(
267                        " WHERE created_at > NOW() - INTERVAL '{} hours'",
268                        hours
269                    ));
270                } else {
271                    query.push_str(&format!(
272                        " AND created_at > NOW() - INTERVAL '{} hours'",
273                        hours
274                    ));
275                }
276            }
277
278            query.push_str(" ORDER BY created_at DESC");
279            query.push_str(&format!(" LIMIT {}", limit));
280
281            let rows = sqlx::query(&query).fetch_all(&pg_pool).await?;
282            for row in rows {
283                let id: uuid::Uuid = row.try_get("id")?;
284                let queue_name: String = row.try_get("queue_name")?;
285                let status: String = row.try_get("status")?;
286                let priority: String = row.try_get("priority")?;
287                let attempts: i32 = row.try_get("attempts")?;
288                let created_at: chrono::DateTime<chrono::Utc> = row.try_get("created_at")?;
289                let scheduled_at: chrono::DateTime<chrono::Utc> = row.try_get("scheduled_at")?;
290
291                job_table.add_job_row(
292                    &id.to_string(),
293                    &queue_name,
294                    &status,
295                    &priority,
296                    attempts,
297                    &created_at.format("%Y-%m-%d %H:%M:%S").to_string(),
298                    &scheduled_at.format("%Y-%m-%d %H:%M:%S").to_string(),
299                );
300            }
301        }
302        DatabasePool::MySQL(mysql_pool) => {
303            // Build MySQL query with different interval syntax
304            let mut query = "SELECT id, queue_name, status, priority, attempts, created_at, scheduled_at FROM hammerwork_jobs".to_string();
305
306            if !conditions.is_empty() {
307                query.push_str(&format!(" WHERE {}", conditions.join(" AND ")));
308            }
309
310            if let Some(hours) = last_hours {
311                if conditions.is_empty() {
312                    query.push_str(&format!(
313                        " WHERE created_at > DATE_SUB(NOW(), INTERVAL {} HOUR)",
314                        hours
315                    ));
316                } else {
317                    query.push_str(&format!(
318                        " AND created_at > DATE_SUB(NOW(), INTERVAL {} HOUR)",
319                        hours
320                    ));
321                }
322            }
323
324            query.push_str(" ORDER BY created_at DESC");
325            query.push_str(&format!(" LIMIT {}", limit));
326
327            let rows = sqlx::query(&query).fetch_all(&mysql_pool).await?;
328            for row in rows {
329                let id: String = row.try_get("id")?;
330                let queue_name: String = row.try_get("queue_name")?;
331                let status: String = row.try_get("status")?;
332                let priority: String = row.try_get("priority")?;
333                let attempts: i32 = row.try_get("attempts")?;
334                let created_at: chrono::DateTime<chrono::Utc> = row.try_get("created_at")?;
335                let scheduled_at: chrono::DateTime<chrono::Utc> = row.try_get("scheduled_at")?;
336
337                job_table.add_job_row(
338                    &id,
339                    &queue_name,
340                    &status,
341                    &priority,
342                    attempts,
343                    &created_at.format("%Y-%m-%d %H:%M:%S").to_string(),
344                    &scheduled_at.format("%Y-%m-%d %H:%M:%S").to_string(),
345                );
346            }
347        }
348    }
349
350    println!("{}", job_table);
351    Ok(())
352}
353
354async fn show_job_details(pool: DatabasePool, job_id: &str) -> Result<()> {
355    match pool {
356        DatabasePool::Postgres(pg_pool) => {
357            let job_uuid = uuid::Uuid::parse_str(job_id)?;
358            let row = sqlx::query("SELECT * FROM hammerwork_jobs WHERE id = $1")
359                .bind(job_uuid)
360                .fetch_optional(&pg_pool)
361                .await?;
362
363            if let Some(row) = row {
364                print_job_details_postgres(&row)?;
365            } else {
366                println!("❌ Job not found: {}", job_id);
367            }
368        }
369        DatabasePool::MySQL(mysql_pool) => {
370            let row = sqlx::query("SELECT * FROM hammerwork_jobs WHERE id = ?")
371                .bind(job_id)
372                .fetch_optional(&mysql_pool)
373                .await?;
374
375            if let Some(row) = row {
376                print_job_details_mysql(&row)?;
377            } else {
378                println!("❌ Job not found: {}", job_id);
379            }
380        }
381    }
382
383    Ok(())
384}
385
386fn print_job_details_postgres(row: &sqlx::postgres::PgRow) -> Result<()> {
387    println!("📋 Job Details");
388    println!("═══════════════");
389    println!("ID: {}", row.try_get::<String, _>("id")?);
390    println!("Queue: {}", row.try_get::<String, _>("queue_name")?);
391    println!("Status: {}", row.try_get::<String, _>("status")?);
392    println!("Priority: {}", row.try_get::<String, _>("priority")?);
393    println!(
394        "Attempts: {}/{}",
395        row.try_get::<i32, _>("attempts")?,
396        row.try_get::<i32, _>("max_attempts")?
397    );
398
399    let created_at: chrono::DateTime<chrono::Utc> = row.try_get("created_at")?;
400    let scheduled_at: chrono::DateTime<chrono::Utc> = row.try_get("scheduled_at")?;
401
402    println!("Created: {}", created_at.format("%Y-%m-%d %H:%M:%S UTC"));
403    println!(
404        "Scheduled: {}",
405        scheduled_at.format("%Y-%m-%d %H:%M:%S UTC")
406    );
407
408    if let Ok(Some(started)) = row.try_get::<Option<chrono::DateTime<chrono::Utc>>, _>("started_at")
409    {
410        println!("Started: {}", started.format("%Y-%m-%d %H:%M:%S UTC"));
411    }
412
413    if let Ok(Some(completed)) =
414        row.try_get::<Option<chrono::DateTime<chrono::Utc>>, _>("completed_at")
415    {
416        println!("Completed: {}", completed.format("%Y-%m-%d %H:%M:%S UTC"));
417    }
418
419    if let Ok(Some(failed)) = row.try_get::<Option<chrono::DateTime<chrono::Utc>>, _>("failed_at") {
420        println!("Failed: {}", failed.format("%Y-%m-%d %H:%M:%S UTC"));
421    }
422
423    if let Ok(Some(error)) = row.try_get::<Option<String>, _>("error_message") {
424        println!("Error: {}", error);
425    }
426
427    let payload: serde_json::Value = row.try_get("payload")?;
428    println!("Payload: {}", serde_json::to_string_pretty(&payload)?);
429
430    Ok(())
431}
432
433fn print_job_details_mysql(row: &sqlx::mysql::MySqlRow) -> Result<()> {
434    println!("📋 Job Details");
435    println!("═══════════════");
436    println!("ID: {}", row.try_get::<String, _>("id")?);
437    println!("Queue: {}", row.try_get::<String, _>("queue_name")?);
438    println!("Status: {}", row.try_get::<String, _>("status")?);
439    println!("Priority: {}", row.try_get::<String, _>("priority")?);
440    println!(
441        "Attempts: {}/{}",
442        row.try_get::<i32, _>("attempts")?,
443        row.try_get::<i32, _>("max_attempts")?
444    );
445
446    let created_at: chrono::DateTime<chrono::Utc> = row.try_get("created_at")?;
447    let scheduled_at: chrono::DateTime<chrono::Utc> = row.try_get("scheduled_at")?;
448
449    println!("Created: {}", created_at.format("%Y-%m-%d %H:%M:%S UTC"));
450    println!(
451        "Scheduled: {}",
452        scheduled_at.format("%Y-%m-%d %H:%M:%S UTC")
453    );
454
455    if let Ok(Some(started)) = row.try_get::<Option<chrono::DateTime<chrono::Utc>>, _>("started_at")
456    {
457        println!("Started: {}", started.format("%Y-%m-%d %H:%M:%S UTC"));
458    }
459
460    if let Ok(Some(completed)) =
461        row.try_get::<Option<chrono::DateTime<chrono::Utc>>, _>("completed_at")
462    {
463        println!("Completed: {}", completed.format("%Y-%m-%d %H:%M:%S UTC"));
464    }
465
466    if let Ok(Some(failed)) = row.try_get::<Option<chrono::DateTime<chrono::Utc>>, _>("failed_at") {
467        println!("Failed: {}", failed.format("%Y-%m-%d %H:%M:%S UTC"));
468    }
469
470    if let Ok(Some(error)) = row.try_get::<Option<String>, _>("error_message") {
471        println!("Error: {}", error);
472    }
473
474    let payload: serde_json::Value = row.try_get("payload")?;
475    println!("Payload: {}", serde_json::to_string_pretty(&payload)?);
476
477    Ok(())
478}
479
480async fn enqueue_job(
481    pool: DatabasePool,
482    queue: &str,
483    payload: &str,
484    priority: &Option<String>,
485    delay: Option<u64>,
486    max_attempts: Option<u32>,
487    timeout: Option<u32>,
488) -> Result<()> {
489    let payload_value = validate_json_payload(payload)?;
490    let job_priority = if let Some(p) = priority {
491        validate_priority(p)?
492    } else {
493        JobPriority::Normal
494    };
495
496    let job_queue = pool.create_job_queue();
497    let mut job = Job::new(queue.to_string(), payload_value);
498    job.priority = job_priority;
499
500    if let Some(max_att) = max_attempts {
501        job.max_attempts = max_att as i32;
502    }
503
504    if let Some(timeout_secs) = timeout {
505        job.timeout = Some(std::time::Duration::from_secs(timeout_secs as u64));
506    }
507
508    if let Some(delay_secs) = delay {
509        let scheduled_at = chrono::Utc::now() + chrono::Duration::seconds(delay_secs as i64);
510        job.scheduled_at = scheduled_at;
511    }
512
513    match job_queue {
514        crate::utils::database::JobQueueWrapper::Postgres(queue) => {
515            let job_id = queue.enqueue(job).await?;
516            info!("✅ Job enqueued successfully: {}", job_id);
517        }
518        crate::utils::database::JobQueueWrapper::MySQL(queue) => {
519            let job_id = queue.enqueue(job).await?;
520            info!("✅ Job enqueued successfully: {}", job_id);
521        }
522    }
523
524    Ok(())
525}
526
527async fn retry_jobs(
528    pool: DatabasePool,
529    job_id: Option<String>,
530    queue: Option<String>,
531    all: bool,
532) -> Result<()> {
533    if !all && job_id.is_none() && queue.is_none() {
534        return Err(anyhow::anyhow!("Must specify --job-id, --queue, or --all"));
535    }
536
537    let (query, affected) = match pool {
538        DatabasePool::Postgres(ref pg_pool) => {
539            if let Some(id) = job_id {
540                let job_uuid = uuid::Uuid::parse_str(&id)?;
541                let result = sqlx::query(
542                    "UPDATE hammerwork_jobs SET status = 'pending', attempts = 0, scheduled_at = NOW() 
543                     WHERE id = $1 AND status IN ('failed', 'dead')"
544                )
545                .bind(job_uuid)
546                .execute(pg_pool).await?;
547                ("single job".to_string(), result.rows_affected())
548            } else if let Some(queue_name) = queue {
549                let result = sqlx::query(
550                    "UPDATE hammerwork_jobs SET status = 'pending', attempts = 0, scheduled_at = NOW() 
551                     WHERE queue_name = $1 AND status IN ('failed', 'dead')"
552                )
553                .bind(&queue_name)
554                .execute(pg_pool).await?;
555                (format!("queue '{}'", queue_name), result.rows_affected())
556            } else {
557                let result = sqlx::query(
558                    "UPDATE hammerwork_jobs SET status = 'pending', attempts = 0, scheduled_at = NOW() 
559                     WHERE status IN ('failed', 'dead')"
560                )
561                .execute(pg_pool).await?;
562                ("all failed jobs".to_string(), result.rows_affected())
563            }
564        }
565        DatabasePool::MySQL(ref mysql_pool) => {
566            if let Some(id) = job_id {
567                let result = sqlx::query(
568                    "UPDATE hammerwork_jobs SET status = 'pending', attempts = 0, scheduled_at = NOW() 
569                     WHERE id = ? AND status IN ('failed', 'dead')"
570                )
571                .bind(id)
572                .execute(mysql_pool).await?;
573                ("single job".to_string(), result.rows_affected())
574            } else if let Some(queue_name) = queue {
575                let result = sqlx::query(
576                    "UPDATE hammerwork_jobs SET status = 'pending', attempts = 0, scheduled_at = NOW() 
577                     WHERE queue_name = ? AND status IN ('failed', 'dead')"
578                )
579                .bind(&queue_name)
580                .execute(mysql_pool).await?;
581                (format!("queue '{}'", queue_name), result.rows_affected())
582            } else {
583                let result = sqlx::query(
584                    "UPDATE hammerwork_jobs SET status = 'pending', attempts = 0, scheduled_at = NOW() 
585                     WHERE status IN ('failed', 'dead')"
586                )
587                .execute(mysql_pool).await?;
588                ("all failed jobs".to_string(), result.rows_affected())
589            }
590        }
591    };
592
593    info!("✅ Retried {} jobs for {}", affected, query);
594    Ok(())
595}
596
597async fn cancel_jobs(
598    pool: DatabasePool,
599    job_id: Option<String>,
600    queue: Option<String>,
601    all_pending: bool,
602) -> Result<()> {
603    if !all_pending && job_id.is_none() && queue.is_none() {
604        return Err(anyhow::anyhow!(
605            "Must specify --job-id, --queue, or --all-pending"
606        ));
607    }
608
609    let (query, affected) = match pool {
610        DatabasePool::Postgres(ref pg_pool) => {
611            if let Some(id) = job_id {
612                let job_uuid = uuid::Uuid::parse_str(&id)?;
613                let result =
614                    sqlx::query("DELETE FROM hammerwork_jobs WHERE id = $1 AND status = 'pending'")
615                        .bind(job_uuid)
616                        .execute(pg_pool)
617                        .await?;
618                ("single job".to_string(), result.rows_affected())
619            } else if let Some(queue_name) = queue {
620                let result = sqlx::query(
621                    "DELETE FROM hammerwork_jobs WHERE queue_name = $1 AND status = 'pending'",
622                )
623                .bind(&queue_name)
624                .execute(pg_pool)
625                .await?;
626                (format!("queue '{}'", queue_name), result.rows_affected())
627            } else {
628                let result = sqlx::query("DELETE FROM hammerwork_jobs WHERE status = 'pending'")
629                    .execute(pg_pool)
630                    .await?;
631                ("all pending jobs".to_string(), result.rows_affected())
632            }
633        }
634        DatabasePool::MySQL(ref mysql_pool) => {
635            if let Some(id) = job_id {
636                let result =
637                    sqlx::query("DELETE FROM hammerwork_jobs WHERE id = ? AND status = 'pending'")
638                        .bind(id)
639                        .execute(mysql_pool)
640                        .await?;
641                ("single job".to_string(), result.rows_affected())
642            } else if let Some(queue_name) = queue {
643                let result = sqlx::query(
644                    "DELETE FROM hammerwork_jobs WHERE queue_name = ? AND status = 'pending'",
645                )
646                .bind(&queue_name)
647                .execute(mysql_pool)
648                .await?;
649                (format!("queue '{}'", queue_name), result.rows_affected())
650            } else {
651                let result = sqlx::query("DELETE FROM hammerwork_jobs WHERE status = 'pending'")
652                    .execute(mysql_pool)
653                    .await?;
654                ("all pending jobs".to_string(), result.rows_affected())
655            }
656        }
657    };
658
659    info!("✅ Cancelled {} jobs for {}", affected, query);
660    Ok(())
661}
662
663async fn purge_jobs(
664    pool: DatabasePool,
665    queue: Option<String>,
666    completed: bool,
667    dead: bool,
668    failed: bool,
669    older_than_days: Option<u32>,
670    confirm: bool,
671) -> Result<()> {
672    if !completed && !dead && !failed {
673        return Err(anyhow::anyhow!(
674            "Must specify at least one of: --completed, --dead, --failed"
675        ));
676    }
677
678    if !confirm {
679        println!("⚠️  This will permanently delete jobs. Use --confirm to proceed.");
680        return Ok(());
681    }
682
683    let mut conditions = Vec::new();
684
685    if completed {
686        conditions.push("status = 'completed'");
687    }
688    if dead {
689        conditions.push("status = 'dead'");
690    }
691    if failed {
692        conditions.push("status = 'failed'");
693    }
694
695    let status_condition = format!("({})", conditions.join(" OR "));
696
697    let affected = match pool {
698        DatabasePool::Postgres(ref pg_pool) => {
699            let mut query = format!("DELETE FROM hammerwork_jobs WHERE {}", status_condition);
700
701            if let Some(queue_name) = queue {
702                query.push_str(&format!(" AND queue_name = '{}'", queue_name));
703            }
704
705            if let Some(days) = older_than_days {
706                query.push_str(&format!(
707                    " AND created_at < NOW() - INTERVAL '{} days'",
708                    days
709                ));
710            }
711
712            let result = sqlx::query(&query).execute(pg_pool).await?;
713            result.rows_affected()
714        }
715        DatabasePool::MySQL(ref mysql_pool) => {
716            let mut query = format!("DELETE FROM hammerwork_jobs WHERE {}", status_condition);
717
718            if let Some(queue_name) = queue {
719                query.push_str(&format!(" AND queue_name = '{}'", queue_name));
720            }
721
722            if let Some(days) = older_than_days {
723                query.push_str(&format!(
724                    " AND created_at < DATE_SUB(NOW(), INTERVAL {} DAY)",
725                    days
726                ));
727            }
728
729            let result = sqlx::query(&query).execute(mysql_pool).await?;
730            result.rows_affected()
731        }
732    };
733
734    info!("✅ Purged {} jobs", affected);
735    Ok(())
736}