Skip to main content

hammerwork_web/api/
jobs.rs

1//! Job management API endpoints.
2//!
3//! This module provides comprehensive REST API endpoints for managing Hammerwork jobs,
4//! including creating, listing, searching, and performing actions on jobs.
5//!
6//! # API Endpoints
7//!
8//! - `GET /api/jobs` - List jobs with filtering, pagination, and sorting
9//! - `POST /api/jobs` - Create a new job
10//! - `GET /api/jobs/{id}` - Get details of a specific job
11//! - `POST /api/jobs/{id}/actions` - Perform actions on a job (retry, cancel, delete)
12//! - `POST /api/jobs/bulk` - Perform bulk actions on multiple jobs
13//! - `POST /api/jobs/search` - Search jobs with full-text queries
14//!
15//! # Examples
16//!
17//! ## Creating a Job
18//!
19//! ```rust
20//! use hammerwork_web::api::jobs::CreateJobRequest;
21//! use serde_json::json;
22//!
23//! let create_request = CreateJobRequest {
24//!     queue_name: "email_queue".to_string(),
25//!     payload: json!({
26//!         "to": "user@example.com",
27//!         "subject": "Welcome!",
28//!         "template": "welcome_email"
29//!     }),
30//!     priority: Some("high".to_string()),
31//!     scheduled_at: None,
32//!     max_attempts: Some(3),
33//!     cron_schedule: None,
34//!     trace_id: Some("trace-123".to_string()),
35//!     correlation_id: Some("corr-456".to_string()),
36//! };
37//!
38//! // This would be sent as JSON in a POST request to /api/jobs
39//! let json_payload = serde_json::to_string(&create_request).unwrap();
40//! assert!(json_payload.contains("email_queue"));
41//! assert!(json_payload.contains("high"));
42//! ```
43//!
44//! ## Job Actions
45//!
46//! ```rust
47//! use hammerwork_web::api::jobs::JobActionRequest;
48//!
49//! let retry_request = JobActionRequest {
50//!     action: "retry".to_string(),
51//!     reason: Some("Network issue resolved".to_string()),
52//! };
53//!
54//! let cancel_request = JobActionRequest {
55//!     action: "cancel".to_string(),
56//!     reason: Some("No longer needed".to_string()),
57//! };
58//!
59//! assert_eq!(retry_request.action, "retry");
60//! assert_eq!(cancel_request.action, "cancel");
61//! ```
62//!
63//! ## Bulk Operations
64//!
65//! ```rust
66//! use hammerwork_web::api::jobs::BulkJobActionRequest;
67//!
68//! let bulk_delete = BulkJobActionRequest {
69//!     job_ids: vec![
70//!         "550e8400-e29b-41d4-a716-446655440000".to_string(),
71//!         "550e8400-e29b-41d4-a716-446655440001".to_string(),
72//!     ],
73//!     action: "delete".to_string(),
74//!     reason: Some("Cleanup old failed jobs".to_string()),
75//! };
76//!
77//! assert_eq!(bulk_delete.job_ids.len(), 2);
78//! assert_eq!(bulk_delete.action, "delete");
79//! ```
80
81use super::{
82    ApiResponse, FilterParams, PaginatedResponse, PaginationMeta, PaginationParams, SortParams,
83    with_filters, with_pagination, with_sort,
84};
85use super::{error_reply, json_reply};
86use hammerwork::queue::DatabaseQueue;
87use serde::{Deserialize, Serialize};
88use std::sync::Arc;
89use warp::http::StatusCode;
90use warp::{Filter, Reply};
91
92/// Job information for API responses
93#[derive(Debug, Serialize)]
94pub struct JobInfo {
95    pub id: String,
96    pub queue_name: String,
97    pub status: String,
98    pub priority: String,
99    pub attempts: i32,
100    pub max_attempts: i32,
101    pub payload: serde_json::Value,
102    pub created_at: chrono::DateTime<chrono::Utc>,
103    pub scheduled_at: chrono::DateTime<chrono::Utc>,
104    pub started_at: Option<chrono::DateTime<chrono::Utc>>,
105    pub completed_at: Option<chrono::DateTime<chrono::Utc>>,
106    pub failed_at: Option<chrono::DateTime<chrono::Utc>>,
107    pub error_message: Option<String>,
108    pub processing_time_ms: Option<i64>,
109    pub cron_schedule: Option<String>,
110    pub is_recurring: bool,
111    pub trace_id: Option<String>,
112    pub correlation_id: Option<String>,
113}
114
115/// Job creation request
116#[derive(Debug, Deserialize, Serialize)]
117pub struct CreateJobRequest {
118    pub queue_name: String,
119    pub payload: serde_json::Value,
120    pub priority: Option<String>,
121    pub scheduled_at: Option<chrono::DateTime<chrono::Utc>>,
122    pub max_attempts: Option<i32>,
123    pub cron_schedule: Option<String>,
124    pub trace_id: Option<String>,
125    pub correlation_id: Option<String>,
126}
127
128/// Job action request
129#[derive(Debug, Deserialize)]
130pub struct JobActionRequest {
131    pub action: String, // "retry", "cancel", "delete"
132    pub reason: Option<String>,
133}
134
135/// The most job IDs one bulk action may name.
136pub const MAX_BULK_JOB_IDS: usize = 1000;
137
138/// Bulk job action request (at most [`MAX_BULK_JOB_IDS`] IDs)
139#[derive(Debug, Deserialize)]
140pub struct BulkJobActionRequest {
141    pub job_ids: Vec<String>,
142    pub action: String,
143    pub reason: Option<String>,
144}
145
146/// Job search request
147#[derive(Debug, Deserialize)]
148pub struct JobSearchRequest {
149    pub query: String,
150    pub queues: Option<Vec<String>>,
151    pub statuses: Option<Vec<String>>,
152    pub priorities: Option<Vec<String>>,
153    pub created_after: Option<chrono::DateTime<chrono::Utc>>,
154    pub created_before: Option<chrono::DateTime<chrono::Utc>>,
155}
156
157/// Create job routes
158pub fn routes<T>(
159    queue: Arc<T>,
160) -> impl Filter<Extract = impl Reply, Error = warp::Rejection> + Clone
161where
162    T: DatabaseQueue + Send + Sync + 'static,
163{
164    let queue_filter = warp::any().map(move || queue.clone());
165
166    let list_jobs = warp::path("jobs")
167        .and(warp::path::end())
168        .and(warp::get())
169        .and(queue_filter.clone())
170        .and(with_pagination())
171        .and(with_filters())
172        .and(with_sort())
173        .and_then(list_jobs_handler);
174
175    let create_job = warp::path("jobs")
176        .and(warp::path::end())
177        .and(warp::post())
178        .and(queue_filter.clone())
179        .and(crate::security::json_body())
180        .and_then(create_job_handler);
181
182    let get_job = warp::path("jobs")
183        .and(warp::path::param::<String>())
184        .and(warp::path::end())
185        .and(warp::get())
186        .and(queue_filter.clone())
187        .and_then(get_job_handler);
188
189    let job_action = warp::path("jobs")
190        .and(warp::path::param::<String>())
191        .and(warp::path("actions"))
192        .and(warp::path::end())
193        .and(warp::post())
194        .and(queue_filter.clone())
195        .and(crate::security::json_body())
196        .and_then(job_action_handler);
197
198    let bulk_action = warp::path("jobs")
199        .and(warp::path("bulk"))
200        .and(warp::path::end())
201        .and(warp::post())
202        .and(queue_filter.clone())
203        .and(crate::security::json_body())
204        .and_then(bulk_job_action_handler);
205
206    let search_jobs = warp::path("jobs")
207        .and(warp::path("search"))
208        .and(warp::path::end())
209        .and(warp::post())
210        .and(queue_filter)
211        .and(crate::security::json_body())
212        .and(with_pagination())
213        .and_then(search_jobs_handler);
214
215    list_jobs
216        .or(create_job)
217        .or(get_job)
218        .or(job_action)
219        .or(bulk_action)
220        .or(search_jobs)
221}
222
223/// The API representation of a job.
224fn job_info(job: &hammerwork::Job) -> JobInfo {
225    let end = job.completed_at.or(job.failed_at).or(job.timed_out_at);
226    JobInfo {
227        id: job.id.to_string(),
228        queue_name: job.queue_name.clone(),
229        status: job.status.as_str().to_string(),
230        priority: job.priority.to_string(),
231        attempts: job.attempts,
232        max_attempts: job.max_attempts,
233        payload: job.payload.clone(),
234        created_at: job.created_at,
235        scheduled_at: job.scheduled_at,
236        started_at: job.started_at,
237        completed_at: job.completed_at,
238        failed_at: job.failed_at,
239        error_message: job.error_message.clone(),
240        processing_time_ms: job
241            .started_at
242            .zip(end)
243            .map(|(start, end)| (end - start).num_milliseconds()),
244        cron_schedule: job.cron_schedule.clone(),
245        is_recurring: job.is_recurring(),
246        trace_id: job.trace_id.clone(),
247        correlation_id: job.correlation_id.clone(),
248    }
249}
250
251/// Whether a job satisfies the `status` filter of the listing: `failed` covers every
252/// unsuccessful terminal status, `recurring` selects recurring jobs, anything else must equal
253/// the job's status (ignoring case).
254fn status_matches(filter: &str, job: &JobInfo) -> bool {
255    let status = job.status.to_lowercase();
256    match filter {
257        "failed" => matches!(status.as_str(), "failed" | "dead" | "timedout"),
258        "recurring" => job.is_recurring,
259        other => status == other,
260    }
261}
262
263/// Order of priorities for sorting (`background` lowest).
264fn priority_rank(priority: &str) -> i32 {
265    priority
266        .parse::<hammerwork::JobPriority>()
267        .map(|p| p.as_i32())
268        .unwrap_or(-1)
269}
270
271/// The jobs of `queue_name` the API can enumerate: pending jobs ready to run, dead jobs
272/// and recurring jobs, without duplicates (a recurring job can also be ready).
273async fn collect_jobs<T>(
274    queue: &T,
275    queue_name: &str,
276    include_ready: bool,
277    include_dead: bool,
278    include_recurring: bool,
279    limit: u32,
280) -> hammerwork::Result<Vec<hammerwork::Job>>
281where
282    T: DatabaseQueue + Send + Sync,
283{
284    let mut jobs = Vec::new();
285    if include_ready {
286        jobs.extend(queue.get_ready_jobs(queue_name, limit).await?);
287    }
288    if include_dead {
289        jobs.extend(
290            queue
291                .get_dead_jobs_by_queue(queue_name, Some(limit), Some(0))
292                .await?,
293        );
294    }
295    if include_recurring {
296        jobs.extend(queue.get_recurring_jobs(queue_name).await?);
297    }
298    let mut seen = std::collections::HashSet::new();
299    jobs.retain(|job| seen.insert(job.id));
300    Ok(jobs)
301}
302
303/// Page `items` according to `pagination` (default 20 per page, at most 100).
304fn paginate<I>(items: Vec<I>, pagination: &PaginationParams) -> PaginatedResponse<I> {
305    let limit = pagination.limit.unwrap_or(20).clamp(1, 100);
306    let page = pagination.page.unwrap_or(1).max(1);
307    let offset = pagination
308        .offset
309        .unwrap_or_else(|| (page - 1).saturating_mul(limit));
310    let total = items.len() as u64;
311    let items = items
312        .into_iter()
313        .skip(offset as usize)
314        .take(limit as usize)
315        .collect();
316    // Describe the page that was actually served, not the requested limit.
317    let served = PaginationParams {
318        page: Some((offset / limit).saturating_add(1)),
319        limit: Some(limit),
320        offset: Some(offset),
321    };
322    PaginatedResponse {
323        items,
324        pagination: PaginationMeta::new(&served, total),
325    }
326}
327
328/// Handler for listing jobs
329pub(crate) async fn list_jobs_handler<T>(
330    queue: Arc<T>,
331    pagination: PaginationParams,
332    filters: FilterParams,
333    sort: SortParams,
334) -> Result<impl Reply, warp::Rejection>
335where
336    T: DatabaseQueue + Send + Sync,
337{
338    // Since DatabaseQueue doesn't provide direct list methods with filters,
339    // we'll use the available methods to gather jobs
340    let mut all_jobs = Vec::new();
341
342    // Get queue stats to find available queues
343    let queue_stats = match queue.get_all_queue_stats().await {
344        Ok(stats) => stats,
345        Err(e) => {
346            return Ok(error_reply(
347                StatusCode::INTERNAL_SERVER_ERROR,
348                format!("Failed to get queue stats: {}", e),
349            ));
350        }
351    };
352
353    // Filter by queue if specified
354    let target_queues: Vec<String> = if let Some(ref queue_name) = filters.queue {
355        vec![queue_name.clone()]
356    } else {
357        queue_stats.iter().map(|s| s.queue_name.clone()).collect()
358    };
359
360    let status_filter = filters.status.as_deref().map(str::to_lowercase);
361    let wants = |names: &[&str]| {
362        status_filter
363            .as_deref()
364            .is_none_or(|status| names.contains(&status))
365    };
366
367    // For each queue, get jobs from different sources based on status filter
368    for queue_name in &target_queues {
369        let queue_jobs = try_api!(
370            collect_jobs(
371                queue.as_ref(),
372                queue_name,
373                wants(&["pending"]),
374                wants(&["failed", "dead"]),
375                wants(&["recurring"]),
376                100
377            )
378            .await,
379            "Failed to list jobs"
380        );
381
382        for job in &queue_jobs {
383            let job_info = job_info(job);
384
385            // Apply status filter
386            if let Some(ref status) = status_filter
387                && !status_matches(status, &job_info)
388            {
389                continue;
390            }
391
392            // Apply priority filter
393            if let Some(ref priority) = filters.priority
394                && !job_info.priority.eq_ignore_ascii_case(priority)
395            {
396                continue;
397            }
398
399            all_jobs.push(job_info);
400        }
401    }
402
403    // Sort jobs
404    let ascending = sort.sort_order.as_deref() == Some("asc");
405    match sort.sort_by.as_deref() {
406        Some("scheduled_at") => all_jobs.sort_by_key(|j| j.scheduled_at),
407        Some("priority") => all_jobs.sort_by_key(|j| priority_rank(&j.priority)),
408        Some("created_at") => all_jobs.sort_by_key(|j| j.created_at),
409        _ => {
410            // Default sort by created_at desc
411            all_jobs.sort_by_key(|j| std::cmp::Reverse(j.created_at));
412        }
413    }
414    let sorted_by_default = !matches!(
415        sort.sort_by.as_deref(),
416        Some("scheduled_at") | Some("priority") | Some("created_at")
417    );
418    if !ascending && !sorted_by_default {
419        all_jobs.reverse();
420    }
421
422    Ok(json_reply(&ApiResponse::success(paginate(
423        all_jobs,
424        &pagination,
425    ))))
426}
427
428/// Handler for creating a new job
429async fn create_job_handler<T>(
430    queue: Arc<T>,
431    request: CreateJobRequest,
432) -> Result<impl Reply, warp::Rejection>
433where
434    T: DatabaseQueue + Send + Sync,
435{
436    use hammerwork::{CronSchedule, Job, JobPriority};
437
438    if request.queue_name.trim().is_empty() {
439        return Ok(error_reply(
440            StatusCode::BAD_REQUEST,
441            "queue_name must not be empty",
442        ));
443    }
444
445    let priority = match request.priority.as_deref() {
446        None => JobPriority::Normal,
447        Some(name) => match name.parse::<JobPriority>() {
448            Ok(priority) => priority,
449            Err(_) => {
450                return Ok(error_reply(
451                    StatusCode::BAD_REQUEST,
452                    format!(
453                        "Invalid priority '{}'. Valid options: background, low, normal, high, critical",
454                        name
455                    ),
456                ));
457            }
458        },
459    };
460
461    if let Some(max_attempts) = request.max_attempts
462        && max_attempts < 1
463    {
464        return Ok(error_reply(
465            StatusCode::BAD_REQUEST,
466            "max_attempts must be at least 1",
467        ));
468    }
469
470    let mut job = Job::new(request.queue_name, request.payload).with_priority(priority);
471
472    if let Some(expression) = request.cron_schedule.as_deref() {
473        let schedule = match CronSchedule::new(expression) {
474            Ok(schedule) => schedule,
475            Err(e) => {
476                return Ok(error_reply(
477                    StatusCode::BAD_REQUEST,
478                    format!("Invalid cron schedule: {}", e),
479                ));
480            }
481        };
482        job = match job.with_cron(schedule) {
483            Ok(job) => job,
484            Err(e) => {
485                return Ok(error_reply(
486                    StatusCode::BAD_REQUEST,
487                    format!("Invalid cron schedule: {}", e),
488                ));
489            }
490        };
491    } else if let Some(scheduled_at) = request.scheduled_at {
492        job.scheduled_at = scheduled_at;
493    }
494
495    if let Some(max_attempts) = request.max_attempts {
496        job = job.with_max_attempts(max_attempts);
497    }
498
499    if let Some(trace_id) = request.trace_id {
500        job.trace_id = Some(trace_id);
501    }
502
503    if let Some(correlation_id) = request.correlation_id {
504        job.correlation_id = Some(correlation_id);
505    }
506
507    match queue.enqueue(job).await {
508        Ok(job_id) => {
509            let response = ApiResponse::success(serde_json::json!({
510                "message": "Job created successfully",
511                "job_id": job_id.to_string()
512            }));
513            Ok(json_reply(&response))
514        }
515        Err(e) => Ok(error_reply(
516            StatusCode::INTERNAL_SERVER_ERROR,
517            format!("Failed to create job: {}", e),
518        )),
519    }
520}
521
522/// Handler for getting a specific job
523async fn get_job_handler<T>(job_id: String, queue: Arc<T>) -> Result<impl Reply, warp::Rejection>
524where
525    T: DatabaseQueue + Send + Sync,
526{
527    let job_uuid = match uuid::Uuid::parse_str(&job_id) {
528        Ok(uuid) => uuid,
529        Err(_) => {
530            return Ok(error_reply(
531                StatusCode::BAD_REQUEST,
532                "Invalid job ID format".to_string(),
533            ));
534        }
535    };
536
537    match queue.get_job(job_uuid).await {
538        Ok(Some(job)) => Ok(json_reply(&ApiResponse::success(job_info(&job)))),
539        Ok(None) => Ok(error_reply(
540            StatusCode::NOT_FOUND,
541            format!("Job '{}' not found", job_id),
542        )),
543        Err(e) => Ok(error_reply(
544            StatusCode::INTERNAL_SERVER_ERROR,
545            format!("Failed to get job: {}", e),
546        )),
547    }
548}
549
550/// Handler for job actions
551async fn job_action_handler<T>(
552    job_id: String,
553    queue: Arc<T>,
554    action_request: JobActionRequest,
555) -> Result<impl Reply, warp::Rejection>
556where
557    T: DatabaseQueue + Send + Sync,
558{
559    let job_uuid = match uuid::Uuid::parse_str(&job_id) {
560        Ok(uuid) => uuid,
561        Err(_) => {
562            return Ok(error_reply(
563                StatusCode::BAD_REQUEST,
564                "Invalid job ID format".to_string(),
565            ));
566        }
567    };
568
569    match action_request.action.as_str() {
570        "retry" => match retry_job_action(queue.as_ref(), job_uuid).await {
571            Ok(()) => {
572                let response = ApiResponse::success(serde_json::json!({
573                    "message": format!("Job '{}' scheduled for retry", job_id)
574                }));
575                Ok(json_reply(&response))
576            }
577            Err(e) => Ok(super::queue_error_reply("Failed to retry job", &e)),
578        },
579        "cancel" | "delete" => match delete_job_action(queue.as_ref(), job_uuid).await {
580            Ok(()) => {
581                let response = ApiResponse::success(serde_json::json!({
582                    "message": format!("Job '{}' deleted", job_id)
583                }));
584                Ok(json_reply(&response))
585            }
586            Err(e) => Ok(super::queue_error_reply("Failed to delete job", &e)),
587        },
588        _ => Ok(error_reply(
589            StatusCode::BAD_REQUEST,
590            format!("Unknown action: {}", action_request.action),
591        )),
592    }
593}
594
595/// Delete a job, reporting a missing job as `JobNotFound` instead of silently succeeding.
596async fn delete_job_action<T>(queue: &T, job_id: uuid::Uuid) -> hammerwork::Result<()>
597where
598    T: DatabaseQueue + Send + Sync,
599{
600    if queue.get_job(job_id).await?.is_none() {
601        return Err(hammerwork::HammerworkError::JobNotFound {
602            id: job_id.to_string(),
603        });
604    }
605    queue.delete_job(job_id).await
606}
607
608/// Re-run a job now: `Dead` and `TimedOut` jobs go through `retry_dead_job` (which also
609/// resets their attempts); other retryable statuses through `retry_job`. Statuses that
610/// cannot be retried (e.g. `Completed`) return an `InvalidJobTransition` error, and a
611/// missing job `JobNotFound`.
612async fn retry_job_action<T>(queue: &T, job_id: uuid::Uuid) -> hammerwork::Result<()>
613where
614    T: DatabaseQueue + Send + Sync,
615{
616    let job =
617        queue
618            .get_job(job_id)
619            .await?
620            .ok_or_else(|| hammerwork::HammerworkError::JobNotFound {
621                id: job_id.to_string(),
622            })?;
623    match job.status {
624        hammerwork::JobStatus::Dead | hammerwork::JobStatus::TimedOut => {
625            queue.retry_dead_job(job_id).await
626        }
627        _ => queue.retry_job(job_id, chrono::Utc::now()).await,
628    }
629}
630
631/// Handler for bulk job actions
632async fn bulk_job_action_handler<T>(
633    queue: Arc<T>,
634    request: BulkJobActionRequest,
635) -> Result<impl Reply, warp::Rejection>
636where
637    T: DatabaseQueue + Send + Sync,
638{
639    if !matches!(request.action.as_str(), "retry" | "delete") {
640        return Ok(error_reply(
641            StatusCode::BAD_REQUEST,
642            format!("Unknown action: {}", request.action),
643        ));
644    }
645    if request.job_ids.len() > MAX_BULK_JOB_IDS {
646        return Ok(error_reply(
647            StatusCode::BAD_REQUEST,
648            format!(
649                "Too many job IDs: {} (at most {} per request)",
650                request.job_ids.len(),
651                MAX_BULK_JOB_IDS
652            ),
653        ));
654    }
655
656    let mut successful = 0;
657    let mut failed = 0;
658    let mut errors = Vec::new();
659
660    for job_id_str in &request.job_ids {
661        let job_uuid = match uuid::Uuid::parse_str(job_id_str) {
662            Ok(uuid) => uuid,
663            Err(_) => {
664                failed += 1;
665                errors.push(format!("Invalid job ID: {}", job_id_str));
666                continue;
667            }
668        };
669
670        let result = match request.action.as_str() {
671            "retry" => retry_job_action(queue.as_ref(), job_uuid).await,
672            _ => delete_job_action(queue.as_ref(), job_uuid).await,
673        };
674
675        match result {
676            Ok(()) => successful += 1,
677            Err(e) => {
678                failed += 1;
679                errors.push(format!("Job {}: {}", job_id_str, e));
680            }
681        }
682    }
683
684    let response = ApiResponse::success(serde_json::json!({
685        "successful": successful,
686        "failed": failed,
687        "errors": errors,
688        "message": format!("Bulk {} completed: {} successful, {} failed", request.action, successful, failed)
689    }));
690
691    Ok(json_reply(&response))
692}
693
694/// Whether `job` matches the lowercase search `term` (id, queue, payload, error, trace ids).
695fn job_matches_search(job: &hammerwork::Job, term: &str) -> bool {
696    job.id.to_string().contains(term)
697        || job.queue_name.to_lowercase().contains(term)
698        || payload_search_text(&job.payload).contains(term)
699        || [&job.error_message, &job.trace_id, &job.correlation_id]
700            .into_iter()
701            .flatten()
702            .any(|text| text.to_lowercase().contains(term))
703}
704
705/// Handler for searching jobs
706async fn search_jobs_handler<T>(
707    queue: Arc<T>,
708    search_request: JobSearchRequest,
709    pagination: PaginationParams,
710) -> Result<impl Reply, warp::Rejection>
711where
712    T: DatabaseQueue + Send + Sync,
713{
714    // Since we don't have direct search methods, we'll gather jobs and filter in memory
715    let mut matching_jobs = Vec::new();
716    let search_term = search_request.query.to_lowercase();
717
718    // Get queue stats to find available queues
719    let queue_stats = match queue.get_all_queue_stats().await {
720        Ok(stats) => stats,
721        Err(e) => {
722            return Ok(error_reply(
723                StatusCode::INTERNAL_SERVER_ERROR,
724                format!("Failed to get queue stats: {}", e),
725            ));
726        }
727    };
728
729    // Filter by specified queues or use all
730    let target_queues: Vec<String> = if let Some(ref queue_names) = search_request.queues {
731        queue_names.clone()
732    } else {
733        queue_stats.iter().map(|s| s.queue_name.clone()).collect()
734    };
735
736    for queue_name in &target_queues {
737        // Search every source: ready, dead and recurring jobs
738        let queue_jobs = try_api!(
739            collect_jobs(queue.as_ref(), queue_name, true, true, true, 200).await,
740            "Failed to list jobs"
741        );
742
743        for job in &queue_jobs {
744            if !job_matches_search(job, &search_term) {
745                continue;
746            }
747
748            // Apply status filter
749            if let Some(ref statuses) = search_request.statuses
750                && !statuses
751                    .iter()
752                    .any(|s| s.eq_ignore_ascii_case(job.status.as_str()))
753            {
754                continue;
755            }
756
757            // Apply priority filter
758            if let Some(ref priorities) = search_request.priorities {
759                let job_priority = job.priority.to_string();
760                if !priorities
761                    .iter()
762                    .any(|p| p.eq_ignore_ascii_case(&job_priority))
763                {
764                    continue;
765                }
766            }
767
768            // Apply date filters
769            if let Some(ref created_after) = search_request.created_after
770                && job.created_at < *created_after
771            {
772                continue;
773            }
774
775            if let Some(ref created_before) = search_request.created_before
776                && job.created_at > *created_before
777            {
778                continue;
779            }
780
781            matching_jobs.push(job_info(job));
782        }
783    }
784
785    // Sort by created_at desc by default
786    matching_jobs.sort_by_key(|j| std::cmp::Reverse(j.created_at));
787
788    Ok(json_reply(&ApiResponse::success(paginate(
789        matching_jobs,
790        &pagination,
791    ))))
792}
793
794/// Lowercased text of a job payload used for substring search.
795///
796/// `serde_json::Value` always serializes, so this does not lose data; it is
797/// kept infallible so a (theoretical) failure cannot make a job unsearchable
798/// silently: such a payload falls back to `Value`'s `Display` output.
799fn payload_search_text(payload: &serde_json::Value) -> String {
800    match serde_json::to_string(payload) {
801        Ok(text) => text.to_lowercase(),
802        Err(e) => {
803            tracing::warn!(error = %e, "failed to serialize job payload for search");
804            payload.to_string().to_lowercase()
805        }
806    }
807}
808
809#[cfg(test)]
810mod tests {
811    use super::*;
812    use crate::api::test_support::{body_json, unreachable_queue};
813
814    #[tokio::test]
815    async fn test_list_jobs_returns_500_when_database_is_down() {
816        let response = list_jobs_handler(
817            unreachable_queue(),
818            PaginationParams::default(),
819            serde_json::from_value(serde_json::json!({})).unwrap(),
820            SortParams {
821                sort_by: None,
822                sort_order: None,
823            },
824        )
825        .await
826        .unwrap()
827        .into_response();
828        let (status, body) = body_json(response).await;
829        assert_eq!(status, 500);
830        assert_eq!(body["success"], false);
831        assert!(
832            body["error"]
833                .as_str()
834                .unwrap()
835                .contains("Failed to get queue stats")
836        );
837    }
838
839    #[tokio::test]
840    async fn test_invalid_job_id_is_400() {
841        let response = get_job_handler("not-a-uuid".to_string(), unreachable_queue())
842            .await
843            .unwrap()
844            .into_response();
845        let (status, _) = body_json(response).await;
846        assert_eq!(status, 400);
847    }
848
849    #[test]
850    fn test_payload_search_text_lowercases_json() {
851        let payload = serde_json::json!({"To": "User@Example.com"});
852        let text = payload_search_text(&payload);
853        assert!(text.contains("user@example.com"));
854        assert!(text.contains("\"to\""));
855    }
856
857    #[test]
858    fn test_create_job_request_deserialization() {
859        let json = r#"{
860            "queue_name": "email",
861            "payload": {"to": "user@example.com", "subject": "Hello"},
862            "priority": "high",
863            "max_attempts": 5
864        }"#;
865
866        let request: CreateJobRequest = serde_json::from_str(json).unwrap();
867        assert_eq!(request.queue_name, "email");
868        assert_eq!(request.priority, Some("high".to_string()));
869        assert_eq!(request.max_attempts, Some(5));
870    }
871
872    #[test]
873    fn test_job_action_request_deserialization() {
874        let json = r#"{"action": "retry", "reason": "Network error resolved"}"#;
875        let request: JobActionRequest = serde_json::from_str(json).unwrap();
876        assert_eq!(request.action, "retry");
877        assert_eq!(request.reason, Some("Network error resolved".to_string()));
878    }
879
880    #[test]
881    fn test_bulk_job_action_request() {
882        let json = r#"{
883            "job_ids": ["job-1", "job-2", "job-3"],
884            "action": "delete",
885            "reason": "Cleanup old jobs"
886        }"#;
887
888        let request: BulkJobActionRequest = serde_json::from_str(json).unwrap();
889        assert_eq!(request.job_ids.len(), 3);
890        assert_eq!(request.action, "delete");
891    }
892}