Skip to main content

hammerwork_web/api/
queues.rs

1//! Queue management API endpoints.
2//!
3//! This module provides REST API endpoints for managing and monitoring Hammerwork job queues,
4//! including queue statistics, actions, and job management within specific queues.
5//!
6//! # API Endpoints
7//!
8//! - `GET /api/queues` - List all queues with statistics
9//! - `GET /api/queues/{name}` - Get detailed statistics for a specific queue
10//! - `POST /api/queues/{name}/actions` - Perform actions on a queue (pause, resume, clear)
11//! - `GET /api/queues/{name}/jobs` - List jobs in a specific queue
12//!
13//! # Examples
14//!
15//! ## Queue Information Structure
16//!
17//! ```rust
18//! use hammerwork_web::api::queues::QueueInfo;
19//! use chrono::Utc;
20//!
21//! let queue_info = QueueInfo {
22//!     name: "email_queue".to_string(),
23//!     pending_count: 25,
24//!     running_count: 3,
25//!     completed_count: 1500,
26//!     failed_count: 12,
27//!     dead_count: 2,
28//!     avg_processing_time_ms: 250.5,
29//!     throughput_per_minute: 45.0,
30//!     error_rate: 0.008,
31//!     last_job_at: Some(Utc::now()),
32//!     oldest_pending_job: Some(Utc::now()),
33//!     is_paused: false,
34//!     paused_at: None,
35//!     paused_by: None,
36//! };
37//!
38//! assert_eq!(queue_info.name, "email_queue");
39//! assert_eq!(queue_info.pending_count, 25);
40//! assert_eq!(queue_info.running_count, 3);
41//! ```
42//!
43//! ## Queue Actions
44//!
45//! ```rust
46//! use hammerwork_web::api::queues::QueueActionRequest;
47//!
48//! let clear_dead_request = QueueActionRequest {
49//!     action: "clear_dead".to_string(),
50//!     confirm: Some(true),
51//! };
52//!
53//! let pause_request = QueueActionRequest {
54//!     action: "pause".to_string(),
55//!     confirm: None,
56//! };
57//!
58//! assert_eq!(clear_dead_request.action, "clear_dead");
59//! assert_eq!(pause_request.action, "pause");
60//! ```
61//!
62//! ## Detailed Queue Statistics
63//!
64//! ```rust
65//! use hammerwork_web::api::queues::{DetailedQueueStats, QueueInfo, HourlyThroughput, RecentError};
66//! use std::collections::HashMap;
67//! use chrono::Utc;
68//!
69//! let queue_info = QueueInfo {
70//!     name: "default".to_string(),
71//!     pending_count: 10,
72//!     running_count: 2,
73//!     completed_count: 500,
74//!     failed_count: 5,
75//!     dead_count: 1,
76//!     avg_processing_time_ms: 180.0,
77//!     throughput_per_minute: 30.0,
78//!     error_rate: 0.01,
79//!     last_job_at: None,
80//!     oldest_pending_job: None,
81//!     is_paused: false,
82//!     paused_at: None,
83//!     paused_by: None,
84//! };
85//!
86//! let mut priority_breakdown = HashMap::new();
87//! priority_breakdown.insert("high".to_string(), 5);
88//! priority_breakdown.insert("normal".to_string(), 15);
89//!
90//! let detailed_stats = DetailedQueueStats {
91//!     queue_info,
92//!     priority_breakdown,
93//!     status_breakdown: HashMap::new(),
94//!     hourly_throughput: vec![],
95//!     recent_errors: vec![],
96//! };
97//!
98//! assert_eq!(detailed_stats.queue_info.name, "default");
99//! assert_eq!(detailed_stats.priority_breakdown.get("high"), Some(&5));
100//! ```
101
102use super::history::{FinishedKind, JobHistory};
103use super::{
104    ApiResponse, FilterParams, PaginatedResponse, PaginationMeta, PaginationParams, SortParams,
105    with_filters, with_pagination, with_sort,
106};
107use super::{error_reply, json_reply};
108use hammerwork::{JobPriority, queue::DatabaseQueue};
109use serde::{Deserialize, Serialize};
110use std::sync::Arc;
111use warp::http::StatusCode;
112use warp::{Filter, Reply};
113
114/// Queue information for API responses
115#[derive(Debug, Serialize, Clone)]
116pub struct QueueInfo {
117    pub name: String,
118    pub pending_count: u64,
119    pub running_count: u64,
120    pub completed_count: u64,
121    pub failed_count: u64,
122    pub dead_count: u64,
123    pub avg_processing_time_ms: f64,
124    pub throughput_per_minute: f64,
125    pub error_rate: f64,
126    pub last_job_at: Option<chrono::DateTime<chrono::Utc>>,
127    pub oldest_pending_job: Option<chrono::DateTime<chrono::Utc>>,
128    pub is_paused: bool,
129    pub paused_at: Option<chrono::DateTime<chrono::Utc>>,
130    pub paused_by: Option<String>,
131}
132
133/// Detailed queue statistics
134#[derive(Debug, Serialize)]
135pub struct DetailedQueueStats {
136    pub queue_info: QueueInfo,
137    pub priority_breakdown: std::collections::HashMap<String, u64>,
138    pub status_breakdown: std::collections::HashMap<String, u64>,
139    pub hourly_throughput: Vec<HourlyThroughput>,
140    pub recent_errors: Vec<RecentError>,
141}
142
143/// Hourly throughput data point
144#[derive(Debug, Serialize)]
145pub struct HourlyThroughput {
146    pub hour: chrono::DateTime<chrono::Utc>,
147    pub completed: u64,
148    pub failed: u64,
149}
150
151/// Recent error information
152#[derive(Debug, Serialize)]
153pub struct RecentError {
154    pub job_id: String,
155    pub error_message: String,
156    pub occurred_at: chrono::DateTime<chrono::Utc>,
157    pub attempts: i32,
158}
159
160/// Queue action request
161#[derive(Debug, Deserialize)]
162pub struct QueueActionRequest {
163    pub action: String, // "pause", "resume", "clear_dead", "clear_completed"
164    pub confirm: Option<bool>,
165}
166
167/// Create queue routes
168pub fn routes<T>(
169    queue: Arc<T>,
170) -> impl Filter<Extract = impl Reply, Error = warp::Rejection> + Clone
171where
172    T: JobHistory + 'static,
173{
174    let queue_filter = warp::any().map(move || queue.clone());
175
176    let list_queues = warp::path("queues")
177        .and(warp::path::end())
178        .and(warp::get())
179        .and(queue_filter.clone())
180        .and(with_pagination())
181        .and(with_filters())
182        .and(with_sort())
183        .and_then(list_queues_handler);
184
185    let get_queue = warp::path("queues")
186        .and(warp::path::param::<String>())
187        .and(warp::path::end())
188        .and(warp::get())
189        .and(queue_filter.clone())
190        .and_then(get_queue_handler);
191
192    let queue_action = warp::path("queues")
193        .and(warp::path::param::<String>())
194        .and(warp::path("actions"))
195        .and(warp::path::end())
196        .and(warp::post())
197        .and(queue_filter.clone())
198        .and(crate::security::json_body())
199        .and_then(queue_action_handler);
200
201    let queue_jobs = warp::path("queues")
202        .and(warp::path::param::<String>())
203        .and(warp::path("jobs"))
204        .and(warp::path::end())
205        .and(warp::get())
206        .and(queue_filter)
207        .and(with_pagination())
208        .and(with_filters())
209        .and(with_sort())
210        .and_then(queue_jobs_handler);
211
212    list_queues.or(get_queue).or(queue_action).or(queue_jobs)
213}
214
215/// Handler for listing all queues
216async fn list_queues_handler<T>(
217    queue: Arc<T>,
218    pagination: PaginationParams,
219    _filters: FilterParams,
220    _sort: SortParams,
221) -> Result<impl Reply, warp::Rejection>
222where
223    T: DatabaseQueue + Send + Sync,
224{
225    // Get all queue statistics
226    match queue.get_all_queue_stats().await {
227        Ok(all_stats) => {
228            let mut queue_infos: Vec<QueueInfo> = Vec::new();
229
230            for stats in all_stats {
231                // Get pause information for this queue
232                let pause_info = try_api!(
233                    queue.get_queue_pause_info(&stats.queue_name).await,
234                    "Failed to get queue pause info"
235                );
236
237                let queue_info = QueueInfo {
238                    name: stats.queue_name.clone(),
239                    pending_count: stats.pending_count,
240                    running_count: stats.running_count,
241                    completed_count: stats.completed_count,
242                    failed_count: stats.dead_count + stats.timed_out_count,
243                    dead_count: stats.dead_count,
244                    avg_processing_time_ms: stats.statistics.avg_processing_time_ms,
245                    throughput_per_minute: stats.statistics.throughput_per_minute,
246                    error_rate: stats.statistics.error_rate,
247                    last_job_at: try_api!(
248                        get_last_job_time(&queue, &stats.queue_name).await,
249                        "Failed to get last job time"
250                    ),
251                    oldest_pending_job: try_api!(
252                        get_oldest_pending_job(&queue, &stats.queue_name).await,
253                        "Failed to get oldest pending job"
254                    ),
255                    is_paused: pause_info.is_some(),
256                    paused_at: pause_info.as_ref().map(|p| p.paused_at),
257                    paused_by: pause_info.as_ref().and_then(|p| p.paused_by.clone()),
258                };
259                queue_infos.push(queue_info);
260            }
261
262            // Apply pagination
263            let total = queue_infos.len() as u64;
264            let offset = pagination.get_offset() as usize;
265            let limit = pagination.get_limit() as usize;
266
267            let items = if offset < queue_infos.len() {
268                let end = (offset + limit).min(queue_infos.len());
269                queue_infos[offset..end].to_vec()
270            } else {
271                Vec::new()
272            };
273
274            let response = PaginatedResponse {
275                items,
276                pagination: PaginationMeta::new(&pagination, total),
277            };
278
279            Ok(json_reply(&ApiResponse::success(response)))
280        }
281        Err(e) => Ok(error_reply(
282            StatusCode::INTERNAL_SERVER_ERROR,
283            format!("Failed to get queue statistics: {}", e),
284        )),
285    }
286}
287
288/// Handler for getting a specific queue
289async fn get_queue_handler<T>(
290    queue_name: String,
291    queue: Arc<T>,
292) -> Result<impl Reply, warp::Rejection>
293where
294    T: JobHistory,
295{
296    match queue.get_all_queue_stats().await {
297        Ok(all_stats) => {
298            if let Some(stats) = all_stats.into_iter().find(|s| s.queue_name == queue_name) {
299                // Get additional details for this specific queue
300                let priority_breakdown = try_api!(
301                    get_priority_breakdown(&queue, &queue_name).await,
302                    "Failed to get priority breakdown"
303                );
304                let status_breakdown = try_api!(
305                    get_status_breakdown(&queue, &queue_name).await,
306                    "Failed to get status breakdown"
307                );
308                let hourly_throughput = try_api!(
309                    get_hourly_throughput(&*queue, &queue_name).await,
310                    "Failed to get hourly throughput"
311                );
312                let recent_errors = try_api!(
313                    get_recent_errors(&queue, &queue_name).await,
314                    "Failed to get recent errors"
315                );
316
317                // Get pause information for this queue
318                let pause_info = try_api!(
319                    queue.get_queue_pause_info(&queue_name).await,
320                    "Failed to get queue pause info"
321                );
322
323                let queue_info = QueueInfo {
324                    name: stats.queue_name.clone(),
325                    pending_count: stats.pending_count,
326                    running_count: stats.running_count,
327                    completed_count: stats.completed_count,
328                    failed_count: stats.dead_count + stats.timed_out_count,
329                    dead_count: stats.dead_count,
330                    avg_processing_time_ms: stats.statistics.avg_processing_time_ms,
331                    throughput_per_minute: stats.statistics.throughput_per_minute,
332                    error_rate: stats.statistics.error_rate,
333                    last_job_at: try_api!(
334                        get_last_job_time(&queue, &stats.queue_name).await,
335                        "Failed to get last job time"
336                    ),
337                    oldest_pending_job: try_api!(
338                        get_oldest_pending_job(&queue, &stats.queue_name).await,
339                        "Failed to get oldest pending job"
340                    ),
341                    is_paused: pause_info.is_some(),
342                    paused_at: pause_info.as_ref().map(|p| p.paused_at),
343                    paused_by: pause_info.as_ref().and_then(|p| p.paused_by.clone()),
344                };
345
346                let detailed_stats = DetailedQueueStats {
347                    queue_info,
348                    priority_breakdown,
349                    status_breakdown,
350                    hourly_throughput,
351                    recent_errors,
352                };
353
354                Ok(json_reply(&ApiResponse::success(detailed_stats)))
355            } else {
356                Ok(error_reply(
357                    StatusCode::NOT_FOUND,
358                    format!("Queue '{}' not found", queue_name),
359                ))
360            }
361        }
362        Err(e) => Ok(error_reply(
363            StatusCode::INTERNAL_SERVER_ERROR,
364            format!("Failed to get queue statistics: {}", e),
365        )),
366    }
367}
368
369/// Handler for queue actions (pause, resume, clear, etc.)
370async fn queue_action_handler<T>(
371    queue_name: String,
372    queue: Arc<T>,
373    action_request: QueueActionRequest,
374) -> Result<impl Reply, warp::Rejection>
375where
376    T: JobHistory,
377{
378    match action_request.action.as_str() {
379        "clear_dead" => {
380            let older_than = chrono::Utc::now() - chrono::Duration::days(7); // Remove jobs older than 7 days
381            match queue
382                .delete_jobs(Some(&queue_name), FinishedKind::Dead, Some(older_than))
383                .await
384            {
385                Ok(count) => {
386                    let response = ApiResponse::success(serde_json::json!({
387                        "message": format!(
388                            "Cleared {} dead jobs older than 7 days from queue '{}'",
389                            count, queue_name
390                        ),
391                        "queue": queue_name,
392                        "count": count
393                    }));
394                    Ok(json_reply(&response))
395                }
396                Err(e) => Ok(error_reply(
397                    StatusCode::INTERNAL_SERVER_ERROR,
398                    format!("Failed to clear dead jobs: {}", e),
399                )),
400            }
401        }
402        "clear_completed" => match queue
403            .delete_jobs(Some(&queue_name), FinishedKind::Completed, None)
404            .await
405        {
406            Ok(count) => {
407                let response = ApiResponse::success(serde_json::json!({
408                    "message": format!("Cleared {} completed jobs from queue '{}'", count, queue_name),
409                    "queue": queue_name,
410                    "cleared_count": count
411                }));
412                Ok(json_reply(&response))
413            }
414            Err(e) => Ok(error_reply(
415                StatusCode::INTERNAL_SERVER_ERROR,
416                format!("Failed to clear completed jobs: {}", e),
417            )),
418        },
419        "pause" => match queue.pause_queue(&queue_name, Some("web-ui")).await {
420            Ok(()) => {
421                let response = ApiResponse::success(serde_json::json!({
422                    "message": format!("Queue '{}' has been paused", queue_name),
423                    "queue": queue_name,
424                    "action": "pause"
425                }));
426                Ok(json_reply(&response))
427            }
428            Err(e) => Ok(error_reply(
429                StatusCode::INTERNAL_SERVER_ERROR,
430                format!("Failed to pause queue: {}", e),
431            )),
432        },
433        "resume" => match queue.resume_queue(&queue_name, Some("web-ui")).await {
434            Ok(()) => {
435                let response = ApiResponse::success(serde_json::json!({
436                    "message": format!("Queue '{}' has been resumed", queue_name),
437                    "queue": queue_name,
438                    "action": "resume"
439                }));
440                Ok(json_reply(&response))
441            }
442            Err(e) => Ok(error_reply(
443                StatusCode::INTERNAL_SERVER_ERROR,
444                format!("Failed to resume queue: {}", e),
445            )),
446        },
447        _ => Ok(error_reply(
448            StatusCode::BAD_REQUEST,
449            format!("Unknown action: {}", action_request.action),
450        )),
451    }
452}
453
454/// Handler for getting jobs in a specific queue: the jobs listing scoped to the queue.
455async fn queue_jobs_handler<T>(
456    queue_name: String,
457    queue: Arc<T>,
458    pagination: PaginationParams,
459    mut filters: FilterParams,
460    sort: SortParams,
461) -> Result<impl Reply, warp::Rejection>
462where
463    T: DatabaseQueue + Send + Sync,
464{
465    filters.queue = Some(queue_name);
466    super::jobs::list_jobs_handler(queue, pagination, filters, sort).await
467}
468
469/// Helper function to get the last job time for a queue
470async fn get_last_job_time<T>(
471    queue: &Arc<T>,
472    queue_name: &str,
473) -> hammerwork::Result<Option<chrono::DateTime<chrono::Utc>>>
474where
475    T: DatabaseQueue + Send + Sync,
476{
477    // Get recent jobs from multiple sources and find the most recent timestamp
478    let mut latest_time: Option<chrono::DateTime<chrono::Utc>> = None;
479
480    // Check ready jobs
481    let ready_jobs = queue.get_ready_jobs(queue_name, 10).await?;
482    for job in ready_jobs {
483        if let Some(time) = job.completed_at.or(job.started_at).or(Some(job.created_at)) {
484            latest_time = match latest_time {
485                Some(current) if time > current => Some(time),
486                None => Some(time),
487                _ => latest_time,
488            };
489        }
490    }
491
492    // Check dead jobs
493    let dead_jobs = queue
494        .get_dead_jobs_by_queue(queue_name, Some(10), Some(0))
495        .await?;
496    for job in dead_jobs {
497        if let Some(time) = job
498            .failed_at
499            .or(job.completed_at)
500            .or(job.started_at)
501            .or(Some(job.created_at))
502        {
503            latest_time = match latest_time {
504                Some(current) if time > current => Some(time),
505                None => Some(time),
506                _ => latest_time,
507            };
508        }
509    }
510
511    Ok(latest_time)
512}
513
514/// Helper function to get the oldest pending job time for a queue
515async fn get_oldest_pending_job<T>(
516    queue: &Arc<T>,
517    queue_name: &str,
518) -> hammerwork::Result<Option<chrono::DateTime<chrono::Utc>>>
519where
520    T: DatabaseQueue + Send + Sync,
521{
522    // Get ready jobs (these are pending jobs) and find the oldest
523    let ready_jobs = queue.get_ready_jobs(queue_name, 100).await?;
524    Ok(ready_jobs
525        .iter()
526        .filter(|job| matches!(job.status, hammerwork::job::JobStatus::Pending))
527        .map(|job| job.created_at)
528        .min())
529}
530
531/// Helper function to get priority breakdown for a queue
532async fn get_priority_breakdown<T>(
533    queue: &Arc<T>,
534    queue_name: &str,
535) -> hammerwork::Result<std::collections::HashMap<String, u64>>
536where
537    T: DatabaseQueue + Send + Sync,
538{
539    // Use the new get_priority_stats method
540    let priority_stats = queue.get_priority_stats(queue_name).await?;
541    let mut breakdown = std::collections::HashMap::new();
542    for (priority, count) in priority_stats.job_counts {
543        let priority_name = match priority {
544            JobPriority::Background => "background",
545            JobPriority::Low => "low",
546            JobPriority::Normal => "normal",
547            JobPriority::High => "high",
548            JobPriority::Critical => "critical",
549        };
550        breakdown.insert(priority_name.to_string(), count);
551    }
552    Ok(breakdown)
553}
554
555/// Helper function to get status breakdown for a queue
556async fn get_status_breakdown<T>(
557    queue: &Arc<T>,
558    queue_name: &str,
559) -> hammerwork::Result<std::collections::HashMap<String, u64>>
560where
561    T: DatabaseQueue + Send + Sync,
562{
563    // Use existing job counts method
564    let counts = queue.get_job_counts_by_status(queue_name).await?;
565    Ok(counts.into_iter().collect())
566}
567
568/// Completed/failed counts per hour over the last 24 hours (23 whole hours plus the
569/// current one) for a queue.
570async fn get_hourly_throughput<T>(
571    queue: &T,
572    queue_name: &str,
573) -> hammerwork::Result<Vec<HourlyThroughput>>
574where
575    T: JobHistory,
576{
577    let now = chrono::Utc::now();
578    let start = super::history::hour_floor(now) - chrono::Duration::hours(23);
579    let buckets = queue.hourly_activity(Some(queue_name), start, now).await?;
580    Ok(buckets
581        .into_iter()
582        .map(|b| HourlyThroughput {
583            hour: b.hour,
584            completed: b.completed,
585            failed: b.failed,
586        })
587        .collect())
588}
589
590/// Helper function to get recent errors for a queue
591async fn get_recent_errors<T>(
592    queue: &Arc<T>,
593    queue_name: &str,
594) -> hammerwork::Result<Vec<RecentError>>
595where
596    T: DatabaseQueue + Send + Sync,
597{
598    // Get dead jobs which contain failed jobs with error messages
599    let dead_jobs = queue
600        .get_dead_jobs_by_queue(queue_name, Some(20), Some(0))
601        .await?;
602    Ok(dead_jobs
603        .into_iter()
604        .filter_map(|job| {
605            job.error_message.map(|error_msg| RecentError {
606                job_id: job.id.to_string(),
607                error_message: error_msg,
608                occurred_at: job.failed_at.unwrap_or(job.created_at),
609                attempts: job.attempts,
610            })
611        })
612        .collect())
613}
614
615#[cfg(test)]
616mod tests {
617    use super::*;
618    use crate::api::test_support::{body_json, unreachable_queue};
619
620    #[tokio::test]
621    async fn test_list_queues_returns_500_when_database_is_down() {
622        let response = list_queues_handler(
623            unreachable_queue(),
624            PaginationParams::default(),
625            serde_json::from_value(serde_json::json!({})).unwrap(),
626            SortParams {
627                sort_by: None,
628                sort_order: None,
629            },
630        )
631        .await
632        .unwrap()
633        .into_response();
634        let (status, body) = body_json(response).await;
635        assert_eq!(status, 500);
636        assert_eq!(body["success"], false);
637    }
638
639    #[tokio::test]
640    async fn test_clear_actions_report_database_failures_not_success() {
641        for action in ["clear_completed", "clear_dead"] {
642            let response = queue_action_handler(
643                "q".to_string(),
644                unreachable_queue(),
645                QueueActionRequest {
646                    action: action.to_string(),
647                    confirm: Some(true),
648                },
649            )
650            .await
651            .unwrap()
652            .into_response();
653            let (status, body) = body_json(response).await;
654            assert_eq!(status, 500, "{action}");
655            assert_eq!(body["success"], false);
656        }
657    }
658
659    #[tokio::test]
660    async fn test_queue_jobs_delegates_to_job_listing() {
661        let response = queue_jobs_handler(
662            "q".to_string(),
663            unreachable_queue(),
664            PaginationParams::default(),
665            serde_json::from_value(serde_json::json!({})).unwrap(),
666            SortParams {
667                sort_by: None,
668                sort_order: None,
669            },
670        )
671        .await
672        .unwrap()
673        .into_response();
674        let (status, body) = body_json(response).await;
675        // The real listing ran (and failed on the dead database), not a stub message.
676        assert_eq!(status, 500);
677        assert!(body.get("data").is_none_or(|d| d.is_null()));
678    }
679
680    #[test]
681    fn test_queue_action_request_deserialization() {
682        let json = r#"{"action": "clear_dead", "confirm": true}"#;
683        let request: QueueActionRequest = serde_json::from_str(json).unwrap();
684        assert_eq!(request.action, "clear_dead");
685        assert_eq!(request.confirm, Some(true));
686    }
687
688    #[test]
689    fn test_queue_info_serialization() {
690        let queue_info = QueueInfo {
691            name: "test_queue".to_string(),
692            pending_count: 42,
693            running_count: 3,
694            completed_count: 1000,
695            failed_count: 5,
696            dead_count: 2,
697            avg_processing_time_ms: 150.5,
698            throughput_per_minute: 25.0,
699            error_rate: 0.05,
700            last_job_at: None,
701            oldest_pending_job: None,
702            is_paused: false,
703            paused_at: None,
704            paused_by: None,
705        };
706
707        let json = serde_json::to_string(&queue_info).unwrap();
708        assert!(json.contains("test_queue"));
709        assert!(json.contains("42"));
710        assert!(json.contains("is_paused"));
711    }
712}