1use 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#[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#[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#[derive(Debug, Serialize)]
145pub struct HourlyThroughput {
146 pub hour: chrono::DateTime<chrono::Utc>,
147 pub completed: u64,
148 pub failed: u64,
149}
150
151#[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#[derive(Debug, Deserialize)]
162pub struct QueueActionRequest {
163 pub action: String, pub confirm: Option<bool>,
165}
166
167pub 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
215async 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 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 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 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
288async 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 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 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
369async 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); 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
454async 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
469async 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 let mut latest_time: Option<chrono::DateTime<chrono::Utc>> = None;
479
480 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 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
514async 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 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
531async 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 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
555async 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 let counts = queue.get_job_counts_by_status(queue_name).await?;
565 Ok(counts.into_iter().collect())
566}
567
568async 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
590async 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 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 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}