1use super::history::JobHistory;
68use super::{
69 ApiResponse, FilterParams, PaginatedResponse, PaginationMeta, PaginationParams, error_reply,
70 internal_error, json_reply,
71};
72use hammerwork::{
73 JobId, JobStatus,
74 archive::{ArchivalConfig, ArchivalPolicy, ArchivalReason, ArchivalStats, ArchivedJob},
75};
76use serde::{Deserialize, Serialize};
77use std::sync::Arc;
78use warp::http::StatusCode;
79use warp::{Filter, Reply};
80
81#[derive(Debug, Deserialize)]
83pub struct ArchiveRequest {
84 pub queue_name: Option<String>,
86 pub reason: ArchivalReason,
88 pub archived_by: Option<String>,
90 pub dry_run: bool,
92 pub policy: Option<ArchivalPolicy>,
94 pub config: Option<ArchivalConfig>,
96}
97
98impl Default for ArchiveRequest {
99 fn default() -> Self {
100 Self {
101 queue_name: None,
102 reason: ArchivalReason::Manual,
103 archived_by: None,
104 dry_run: true,
105 policy: None,
106 config: None,
107 }
108 }
109}
110
111#[derive(Debug, Serialize)]
113pub struct ArchiveResponse {
114 pub stats: ArchivalStats,
116 pub dry_run: bool,
118 pub policy_used: ArchivalPolicy,
120 pub config_used: ArchivalConfig,
122}
123
124#[derive(Debug, Deserialize)]
126pub struct RestoreRequest {
127 pub reason: Option<String>,
129 pub restored_by: Option<String>,
131}
132
133#[derive(Debug, Serialize)]
135pub struct RestoreResponse {
136 pub job: hammerwork::Job,
138 pub restored_at: chrono::DateTime<chrono::Utc>,
140 pub restored_by: Option<String>,
142}
143
144#[derive(Debug, Deserialize)]
146pub struct PurgeRequest {
147 pub older_than: chrono::DateTime<chrono::Utc>,
149 pub dry_run: bool,
151 pub purged_by: Option<String>,
153}
154
155#[derive(Debug, Serialize)]
157pub struct PurgeResponse {
158 pub jobs_purged: u64,
160 pub dry_run: bool,
162 pub executed_at: chrono::DateTime<chrono::Utc>,
164}
165
166#[derive(Debug, Serialize)]
168pub struct ArchivedJobInfo {
169 pub id: JobId,
171 pub queue_name: String,
173 pub status: JobStatus,
175 pub created_at: chrono::DateTime<chrono::Utc>,
177 pub archived_at: chrono::DateTime<chrono::Utc>,
179 pub archival_reason: ArchivalReason,
181 pub original_payload_size: Option<usize>,
183 pub payload_compressed: bool,
185 pub archived_by: Option<String>,
187}
188
189impl From<ArchivedJob> for ArchivedJobInfo {
190 fn from(archived_job: ArchivedJob) -> Self {
191 Self {
192 id: archived_job.id,
193 queue_name: archived_job.queue_name,
194 status: archived_job.status,
195 created_at: archived_job.created_at,
196 archived_at: archived_job.archived_at,
197 archival_reason: archived_job.archival_reason,
198 original_payload_size: archived_job.original_payload_size,
199 payload_compressed: archived_job.payload_compressed,
200 archived_by: archived_job.archived_by,
201 }
202 }
203}
204
205#[derive(Debug, Serialize)]
207pub struct StatsResponse {
208 pub stats: ArchivalStats,
210 pub by_queue: std::collections::HashMap<String, ArchivalStats>,
212 pub recent_operations: Vec<RecentOperation>,
214}
215
216#[derive(Debug, Serialize)]
218pub struct RecentOperation {
219 pub operation_type: String,
221 pub queue_name: Option<String>,
223 pub jobs_affected: u64,
225 pub executed_at: chrono::DateTime<chrono::Utc>,
227 pub executed_by: Option<String>,
229 pub reason: Option<String>,
231}
232
233#[derive(Debug, Deserialize, Default)]
235pub struct ArchiveFilterParams {
236 pub queue: Option<String>,
238 pub reason: Option<String>,
240 pub archived_after: Option<chrono::DateTime<chrono::Utc>>,
242 pub archived_before: Option<chrono::DateTime<chrono::Utc>>,
244 pub archived_by: Option<String>,
246 pub compressed: Option<bool>,
248 pub original_status: Option<String>,
250}
251
252pub fn archive_routes<Q>(
254 queue: Arc<Q>,
255) -> impl Filter<Extract = impl Reply, Error = warp::Rejection> + Clone
256where
257 Q: JobHistory + 'static,
258{
259 let archive_jobs = warp::path!("archive" / "jobs")
260 .and(warp::post())
261 .and(crate::security::json_body())
262 .and(with_queue(queue.clone()))
263 .and_then(handle_archive_jobs);
264
265 let list_archived = warp::path!("archive" / "jobs")
266 .and(warp::get())
267 .and(super::with_pagination())
268 .and(with_archive_filters())
269 .and(with_queue(queue.clone()))
270 .and_then(handle_list_archived_jobs);
271
272 let restore_job = warp::path!("archive" / "jobs" / String / "restore")
273 .and(warp::post())
274 .and(crate::security::json_body())
275 .and(with_queue(queue.clone()))
276 .and_then(handle_restore_job);
277
278 let purge_jobs = warp::path!("archive" / "purge")
279 .and(warp::delete())
280 .and(crate::security::json_body())
281 .and(with_queue(queue.clone()))
282 .and_then(handle_purge_jobs);
283
284 let archive_stats = warp::path!("archive" / "stats")
285 .and(warp::get())
286 .and(warp::query::<FilterParams>())
287 .and(with_queue(queue.clone()))
288 .and_then(handle_archive_stats);
289
290 archive_jobs
291 .or(list_archived)
292 .or(restore_job)
293 .or(purge_jobs)
294 .or(archive_stats)
295}
296
297fn with_queue<Q>(
299 queue: Arc<Q>,
300) -> impl Filter<Extract = (Arc<Q>,), Error = std::convert::Infallible> + Clone
301where
302 Q: JobHistory + 'static,
303{
304 warp::any().map(move || queue.clone())
305}
306
307fn with_archive_filters()
309-> impl Filter<Extract = (ArchiveFilterParams,), Error = warp::Rejection> + Clone {
310 warp::query::<ArchiveFilterParams>()
311}
312
313async fn handle_archive_jobs<Q>(
315 request: ArchiveRequest,
316 queue: Arc<Q>,
317) -> Result<impl Reply, warp::Rejection>
318where
319 Q: JobHistory + 'static,
320{
321 let policy = request.policy.unwrap_or_default();
322 let config = request.config.unwrap_or_default();
323
324 if request.dry_run {
325 let response = ArchiveResponse {
327 stats: ArchivalStats::default(), dry_run: true,
329 policy_used: policy,
330 config_used: config,
331 };
332 return Ok(json_reply(&ApiResponse::success(response)));
333 }
334
335 match queue
336 .archive_jobs(
337 request.queue_name.as_deref(),
338 &policy,
339 &config,
340 request.reason,
341 request.archived_by.as_deref(),
342 )
343 .await
344 {
345 Ok(stats) => {
346 let response = ArchiveResponse {
347 stats,
348 dry_run: false,
349 policy_used: policy,
350 config_used: config,
351 };
352 Ok(json_reply(&ApiResponse::success(response)))
353 }
354 Err(e) => Ok(internal_error("Failed to archive jobs", &e)),
355 }
356}
357
358async fn handle_list_archived_jobs<Q>(
361 pagination: PaginationParams,
362 filters: ArchiveFilterParams,
363 queue: Arc<Q>,
364) -> Result<impl Reply, warp::Rejection>
365where
366 Q: JobHistory + 'static,
367{
368 let limit = pagination.get_limit();
369 let offset = pagination.get_offset();
370
371 let (page, total) = match queue.archived_jobs(&filters, limit, offset).await {
372 Ok(listing) => listing,
373 Err(e) => return Ok(internal_error("Failed to list archived jobs", &e)),
374 };
375
376 let response = PaginatedResponse {
377 items: page.into_iter().map(ArchivedJobInfo::from).collect(),
378 pagination: PaginationMeta::new(&pagination, total),
379 };
380 Ok(json_reply(&ApiResponse::success(response)))
381}
382
383async fn handle_restore_job<Q>(
385 job_id_str: String,
386 request: RestoreRequest,
387 queue: Arc<Q>,
388) -> Result<impl Reply, warp::Rejection>
389where
390 Q: JobHistory + 'static,
391{
392 let job_id = match uuid::Uuid::parse_str(&job_id_str) {
393 Ok(id) => id,
394 Err(_) => {
395 return Ok(error_reply(
396 StatusCode::BAD_REQUEST,
397 "Invalid job ID format",
398 ));
399 }
400 };
401
402 match queue.restore_archived_job(job_id).await {
403 Ok(job) => {
404 let response = RestoreResponse {
405 job,
406 restored_at: chrono::Utc::now(),
407 restored_by: request.restored_by,
408 };
409 Ok(json_reply(&ApiResponse::success(response)))
410 }
411 Err(e) => Ok(super::queue_error_reply("Failed to restore job", &e)),
412 }
413}
414
415async fn handle_purge_jobs<Q>(
417 request: PurgeRequest,
418 queue: Arc<Q>,
419) -> Result<impl Reply, warp::Rejection>
420where
421 Q: JobHistory + 'static,
422{
423 if request.dry_run {
424 let cutoff = ArchiveFilterParams {
427 archived_before: Some(request.older_than),
428 ..ArchiveFilterParams::default()
429 };
430 let count = match queue.archived_jobs(&cutoff, 0, 0).await {
431 Ok((_, count)) => count,
432 Err(e) => return Ok(internal_error("Failed to count archived jobs", &e)),
433 };
434
435 let response = PurgeResponse {
436 jobs_purged: count,
437 dry_run: true,
438 executed_at: chrono::Utc::now(),
439 };
440 return Ok(json_reply(&ApiResponse::success(response)));
441 }
442
443 match queue.purge_archived_jobs(request.older_than).await {
444 Ok(jobs_purged) => {
445 let response = PurgeResponse {
446 jobs_purged,
447 dry_run: false,
448 executed_at: chrono::Utc::now(),
449 };
450 Ok(json_reply(&ApiResponse::success(response)))
451 }
452 Err(e) => Ok(internal_error("Failed to purge archived jobs", &e)),
453 }
454}
455
456async fn handle_archive_stats<Q>(
458 filters: FilterParams,
459 queue: Arc<Q>,
460) -> Result<impl Reply, warp::Rejection>
461where
462 Q: JobHistory + 'static,
463{
464 match queue.get_archival_stats(filters.queue.as_deref()).await {
465 Ok(stats) => {
466 let mut by_queue = std::collections::HashMap::new();
468
469 if filters.queue.is_none() {
470 let queue_stats = match queue.get_all_queue_stats().await {
472 Ok(queue_stats) => queue_stats,
473 Err(e) => return Ok(internal_error("Failed to get queue stats", &e)),
474 };
475 for queue_stat in queue_stats {
476 match queue.get_archival_stats(Some(&queue_stat.queue_name)).await {
477 Ok(queue_archival_stats) => {
478 by_queue.insert(queue_stat.queue_name, queue_archival_stats);
479 }
480 Err(e) => {
481 return Ok(internal_error(
482 &format!(
483 "Failed to get archive stats for queue '{}'",
484 queue_stat.queue_name
485 ),
486 &e,
487 ));
488 }
489 }
490 }
491 }
492
493 let recent_operations: Vec<RecentOperation> = Vec::new();
496
497 let response = StatsResponse {
498 stats,
499 by_queue,
500 recent_operations,
501 };
502 Ok(json_reply(&ApiResponse::success(response)))
503 }
504 Err(e) => Ok(internal_error("Failed to get archive stats", &e)),
505 }
506}
507
508#[cfg(test)]
509mod tests {
510 use super::*;
511 use crate::api::test_support::{body_json, unreachable_queue};
512
513 #[tokio::test]
514 async fn test_dry_run_purge_returns_500_instead_of_zero_when_database_is_down() {
515 let request: PurgeRequest = serde_json::from_value(serde_json::json!({
516 "older_than": "2024-01-01T00:00:00Z",
517 "dry_run": true
518 }))
519 .unwrap();
520 let response = handle_purge_jobs(request, unreachable_queue())
521 .await
522 .unwrap()
523 .into_response();
524 let (status, body) = body_json(response).await;
525 assert_eq!(status, 500);
526 assert_eq!(body["success"], false);
527 assert!(body["data"].is_null());
528 }
529
530 #[tokio::test]
531 async fn test_archive_stats_returns_500_when_database_is_down() {
532 let filters: FilterParams = serde_json::from_value(serde_json::json!({})).unwrap();
533 let response = handle_archive_stats(filters, unreachable_queue())
534 .await
535 .unwrap()
536 .into_response();
537 let (status, _) = body_json(response).await;
538 assert_eq!(status, 500);
539 }
540 use chrono::Duration;
541
542 #[test]
543 fn test_archive_request_default() {
544 let request = ArchiveRequest::default();
545 assert!(request.dry_run);
546 assert_eq!(request.reason, ArchivalReason::Manual);
547 assert!(request.queue_name.is_none());
548 }
549
550 #[test]
551 fn test_archived_job_info_conversion() {
552 let archived_job = ArchivedJob {
553 id: uuid::Uuid::new_v4(),
554 queue_name: "test_queue".to_string(),
555 status: JobStatus::Completed,
556 created_at: chrono::Utc::now(),
557 archived_at: chrono::Utc::now(),
558 archival_reason: ArchivalReason::Automatic,
559 original_payload_size: Some(1024),
560 payload_compressed: true,
561 archived_by: Some("system".to_string()),
562 };
563
564 let info: ArchivedJobInfo = archived_job.into();
565 assert_eq!(info.queue_name, "test_queue");
566 assert!(info.payload_compressed);
567 assert_eq!(info.archival_reason, ArchivalReason::Automatic);
568 }
569
570 #[test]
571 fn test_purge_request_validation() {
572 let request = PurgeRequest {
573 older_than: chrono::Utc::now() - Duration::days(365),
574 dry_run: true,
575 purged_by: Some("admin".to_string()),
576 };
577
578 assert!(request.dry_run);
579 assert_eq!(request.purged_by, Some("admin".to_string()));
580 }
581}