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 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 let mut conditions = Vec::new();
227 if let Some(queue_name) = &queue {
230 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 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 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}