1use 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#[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#[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#[derive(Debug, Deserialize)]
130pub struct JobActionRequest {
131 pub action: String, pub reason: Option<String>,
133}
134
135pub const MAX_BULK_JOB_IDS: usize = 1000;
137
138#[derive(Debug, Deserialize)]
140pub struct BulkJobActionRequest {
141 pub job_ids: Vec<String>,
142 pub action: String,
143 pub reason: Option<String>,
144}
145
146#[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
157pub 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
223fn 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
251fn 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
263fn priority_rank(priority: &str) -> i32 {
265 priority
266 .parse::<hammerwork::JobPriority>()
267 .map(|p| p.as_i32())
268 .unwrap_or(-1)
269}
270
271async 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
303fn 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 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
328pub(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 let mut all_jobs = Vec::new();
341
342 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 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 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 if let Some(ref status) = status_filter
387 && !status_matches(status, &job_info)
388 {
389 continue;
390 }
391
392 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 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 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
428async 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
522async 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
550async 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
595async 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
608async 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
631async 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
694fn 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
705async 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 let mut matching_jobs = Vec::new();
716 let search_term = search_request.query.to_lowercase();
717
718 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 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 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 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 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 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 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
794fn 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}