1use super::history::{ErrorGroup, JobHistory, MAX_BUCKETS, hour_floor};
113use super::{ApiResponse, error_reply, json_reply};
114use hammerwork::queue::DatabaseQueue;
115use serde::{Deserialize, Serialize};
116use std::collections::HashMap;
117use std::sync::Arc;
118use warp::http::StatusCode;
119use warp::{Filter, Reply};
120
121#[derive(Debug, Serialize)]
123pub struct SystemOverview {
124 pub total_queues: u32,
125 pub total_jobs: u64,
126 pub pending_jobs: u64,
127 pub running_jobs: u64,
128 pub completed_jobs: u64,
129 pub failed_jobs: u64,
130 pub dead_jobs: u64,
131 pub overall_throughput: f64,
132 pub overall_error_rate: f64,
133 pub avg_processing_time_ms: f64,
134 pub system_health: SystemHealth,
135 pub uptime_seconds: u64,
136 pub last_updated: chrono::DateTime<chrono::Utc>,
137}
138
139#[derive(Debug, Serialize)]
141pub struct SystemHealth {
142 pub status: String, pub database_healthy: bool,
144 pub high_error_rate: bool,
145 pub queue_backlog: bool,
146 pub slow_processing: bool,
147 pub alerts: Vec<SystemAlert>,
148}
149
150#[derive(Debug, Serialize)]
152pub struct SystemAlert {
153 pub severity: String, pub message: String,
155 pub queue: Option<String>,
156 pub metric: Option<String>,
157 pub value: Option<f64>,
158 pub threshold: Option<f64>,
159 pub timestamp: chrono::DateTime<chrono::Utc>,
160}
161
162#[derive(Debug, Serialize)]
164pub struct DetailedStats {
165 pub overview: SystemOverview,
166 pub queue_stats: Vec<QueueStats>,
167 pub hourly_trends: Vec<HourlyTrend>,
168 pub error_patterns: Vec<ErrorPattern>,
169 pub performance_metrics: PerformanceMetrics,
170}
171
172#[derive(Debug, Serialize)]
174pub struct QueueStats {
175 pub name: String,
176 pub pending: u64,
177 pub running: u64,
178 pub completed_total: u64,
179 pub failed_total: u64,
180 pub dead_total: u64,
181 pub throughput_per_minute: f64,
182 pub avg_processing_time_ms: f64,
183 pub error_rate: f64,
184 pub oldest_pending_age_seconds: Option<u64>,
185 pub priority_distribution: HashMap<String, f32>,
186}
187
188#[derive(Debug, Serialize)]
190pub struct HourlyTrend {
191 pub hour: chrono::DateTime<chrono::Utc>,
192 pub completed: u64,
193 pub failed: u64,
194 pub throughput: f64,
195 pub avg_processing_time_ms: Option<f64>,
197 pub error_rate: f64,
198}
199
200#[derive(Debug, Serialize)]
202pub struct ErrorPattern {
203 pub error_type: String,
204 pub count: u64,
205 pub percentage: f64,
206 pub sample_message: String,
207 pub first_seen: chrono::DateTime<chrono::Utc>,
208 pub last_seen: chrono::DateTime<chrono::Utc>,
209 pub affected_queues: Vec<String>,
210}
211
212#[derive(Debug, Serialize)]
214pub struct PerformanceMetrics {
215 pub database_response_time_ms: Option<f64>,
217 pub average_queue_depth: f64,
218 pub jobs_per_second: f64,
219 pub memory_usage_mb: Option<f64>,
220 pub cpu_usage_percent: Option<f64>,
221 pub active_workers: Option<u32>,
223 pub worker_utilization: Option<f64>,
225}
226
227#[derive(Debug, Deserialize)]
229pub struct TimeRange {
230 pub start: chrono::DateTime<chrono::Utc>,
231 pub end: chrono::DateTime<chrono::Utc>,
232}
233
234#[derive(Debug, Deserialize)]
236pub struct StatsQuery {
237 pub time_range: Option<TimeRange>,
238 pub queues: Option<Vec<String>>,
239 pub granularity: Option<String>, pub hours: Option<u32>,
242}
243
244pub fn routes<T>(
246 queue: Arc<T>,
247 system_state: Arc<tokio::sync::RwLock<crate::api::system::SystemState>>,
248) -> impl Filter<Extract = impl Reply, Error = warp::Rejection> + Clone
249where
250 T: JobHistory + 'static,
251{
252 let queue_filter = warp::any().map(move || queue.clone());
253 let state_filter = warp::any().map(move || system_state.clone());
254
255 let overview = warp::path("stats")
256 .and(warp::path("overview"))
257 .and(warp::path::end())
258 .and(warp::get())
259 .and(queue_filter.clone())
260 .and(state_filter.clone())
261 .and_then(overview_handler);
262
263 let detailed = warp::path("stats")
264 .and(warp::path("detailed"))
265 .and(warp::path::end())
266 .and(warp::get())
267 .and(queue_filter.clone())
268 .and(state_filter.clone())
269 .and(warp::query::<StatsQuery>())
270 .and_then(detailed_stats_handler);
271
272 let trends = warp::path("stats")
273 .and(warp::path("trends"))
274 .and(warp::path::end())
275 .and(warp::get())
276 .and(queue_filter.clone())
277 .and(warp::query::<StatsQuery>())
278 .and_then(trends_handler);
279
280 let health = warp::path("stats")
281 .and(warp::path("health"))
282 .and(warp::path::end())
283 .and(warp::get())
284 .and(queue_filter)
285 .and_then(health_handler);
286
287 overview.or(detailed).or(trends).or(health)
288}
289
290async fn overview_handler<T>(
292 queue: Arc<T>,
293 system_state: Arc<tokio::sync::RwLock<crate::api::system::SystemState>>,
294) -> Result<impl Reply, warp::Rejection>
295where
296 T: DatabaseQueue + Send + Sync,
297{
298 match queue.get_all_queue_stats().await {
299 Ok(all_stats) => {
300 let mut total_pending = 0;
301 let mut total_running = 0;
302 let mut total_completed = 0;
303 let mut total_failed = 0;
304 let mut total_dead = 0;
305 let mut total_throughput = 0.0;
306 let mut total_processing_time = 0.0;
307 let mut queue_count = 0;
308
309 for stats in &all_stats {
310 total_pending += stats.pending_count;
311 total_running += stats.running_count;
312 total_completed += stats.completed_count;
313 total_failed += stats.dead_count + stats.timed_out_count;
314 total_dead += stats.dead_count;
315 total_throughput += stats.statistics.throughput_per_minute;
316 total_processing_time += stats.statistics.avg_processing_time_ms;
317 queue_count += 1;
318 }
319
320 let avg_processing_time = if queue_count > 0 {
321 total_processing_time / queue_count as f64
322 } else {
323 0.0
324 };
325
326 let total_jobs = total_pending + total_running + total_completed + total_failed;
327 let overall_error_rate = if total_jobs > 0 {
328 total_failed as f64 / total_jobs as f64
329 } else {
330 0.0
331 };
332
333 let health = assess_system_health(&all_stats);
335
336 let overview = SystemOverview {
337 total_queues: queue_count,
338 total_jobs,
339 pending_jobs: total_pending,
340 running_jobs: total_running,
341 completed_jobs: total_completed,
342 failed_jobs: total_failed,
343 dead_jobs: total_dead,
344 overall_throughput: total_throughput,
345 overall_error_rate,
346 avg_processing_time_ms: avg_processing_time,
347 system_health: health,
348 uptime_seconds: {
349 let state = system_state.read().await;
350 state.uptime_seconds() as u64
351 },
352 last_updated: chrono::Utc::now(),
353 };
354
355 Ok(json_reply(&ApiResponse::success(overview)))
356 }
357 Err(e) => Ok(error_reply(
358 StatusCode::INTERNAL_SERVER_ERROR,
359 format!("Failed to get statistics: {}", e),
360 )),
361 }
362}
363
364async fn detailed_stats_handler<T>(
366 queue: Arc<T>,
367 system_state: Arc<tokio::sync::RwLock<crate::api::system::SystemState>>,
368 query: StatsQuery,
369) -> Result<impl Reply, warp::Rejection>
370where
371 T: JobHistory,
372{
373 let now = chrono::Utc::now();
374 let (start, end) = match resolve_range(&query, now) {
375 Ok(range) => range,
376 Err(message) => return Ok(error_reply(StatusCode::BAD_REQUEST, message)),
377 };
378
379 let db_started = std::time::Instant::now();
380 match queue.get_all_queue_stats().await {
381 Ok(all_stats) => {
382 let db_ms = db_started.elapsed().as_secs_f64() * 1000.0;
383 let mut queue_stats: Vec<QueueStats> = Vec::new();
385 for stats in all_stats.iter() {
386 let oldest_pending_age_seconds = try_api!(
388 calculate_oldest_pending_age(&queue, &stats.queue_name).await,
389 "Failed to get oldest pending job age"
390 );
391
392 let priority_distribution = try_api!(
394 get_priority_distribution(&queue, &stats.queue_name).await,
395 "Failed to get priority distribution"
396 );
397
398 queue_stats.push(QueueStats {
399 name: stats.queue_name.clone(),
400 pending: stats.pending_count,
401 running: stats.running_count,
402 completed_total: stats.completed_count,
403 failed_total: stats.dead_count + stats.timed_out_count,
404 dead_total: stats.dead_count,
405 throughput_per_minute: stats.statistics.throughput_per_minute,
406 avg_processing_time_ms: stats.statistics.avg_processing_time_ms,
407 error_rate: stats.statistics.error_rate,
408 oldest_pending_age_seconds,
409 priority_distribution,
410 });
411 }
412
413 let hourly_trends = try_api!(
414 compute_trends(&*queue, start, end).await,
415 "Failed to compute hourly trends"
416 );
417 let error_patterns = try_api!(
418 generate_error_patterns(&*queue, start).await,
419 "Failed to compute error patterns"
420 );
421 let performance_metrics = calculate_performance_metrics(&all_stats, db_ms);
422
423 let uptime_seconds = system_state.read().await.uptime_seconds().max(0) as u64;
424 let overview = generate_overview_from_stats(&all_stats, uptime_seconds);
425
426 let detailed = DetailedStats {
427 overview,
428 queue_stats,
429 hourly_trends,
430 error_patterns,
431 performance_metrics,
432 };
433
434 Ok(json_reply(&ApiResponse::success(detailed)))
435 }
436 Err(e) => Ok(error_reply(
437 StatusCode::INTERNAL_SERVER_ERROR,
438 format!("Failed to get detailed statistics: {}", e),
439 )),
440 }
441}
442
443fn resolve_range(
446 query: &StatsQuery,
447 now: chrono::DateTime<chrono::Utc>,
448) -> Result<(chrono::DateTime<chrono::Utc>, chrono::DateTime<chrono::Utc>), String> {
449 if let Some(hours) = query.hours {
450 if hours == 0 || i64::from(hours) > MAX_BUCKETS {
451 return Err(format!("hours must be between 1 and {}", MAX_BUCKETS));
452 }
453 return Ok((
454 hour_floor(now) - chrono::Duration::hours(i64::from(hours) - 1),
455 now,
456 ));
457 }
458 match &query.time_range {
459 Some(range) => {
460 if range.start >= range.end {
461 Err("time_range.start must be before time_range.end".to_string())
462 } else if range.end - range.start > chrono::Duration::hours(MAX_BUCKETS) {
463 Err(format!("time_range may span at most {} hours", MAX_BUCKETS))
464 } else {
465 Ok((range.start, range.end))
466 }
467 }
468 None => Ok((hour_floor(now) - chrono::Duration::hours(23), now)),
469 }
470}
471
472async fn trends_handler<T>(queue: Arc<T>, query: StatsQuery) -> Result<impl Reply, warp::Rejection>
474where
475 T: JobHistory,
476{
477 let (start, end) = match resolve_range(&query, chrono::Utc::now()) {
478 Ok(range) => range,
479 Err(message) => return Ok(error_reply(StatusCode::BAD_REQUEST, message)),
480 };
481 let trends = try_api!(
482 compute_trends(&*queue, start, end).await,
483 "Failed to compute hourly trends"
484 );
485 Ok(json_reply(&ApiResponse::success(trends)))
486}
487
488async fn health_handler<T>(queue: Arc<T>) -> Result<impl Reply, warp::Rejection>
490where
491 T: DatabaseQueue + Send + Sync,
492{
493 match queue.get_all_queue_stats().await {
494 Ok(all_stats) => {
495 let health = assess_system_health(&all_stats);
496 Ok(json_reply(&ApiResponse::success(health)))
497 }
498 Err(e) => {
499 let health = SystemHealth {
500 status: "critical".to_string(),
501 database_healthy: false,
502 high_error_rate: false,
503 queue_backlog: false,
504 slow_processing: false,
505 alerts: vec![SystemAlert {
506 severity: "critical".to_string(),
507 message: format!("Database connection failed: {}", e),
508 queue: None,
509 metric: Some("database_connectivity".to_string()),
510 value: None,
511 threshold: None,
512 timestamp: chrono::Utc::now(),
513 }],
514 };
515 Ok(json_reply(&ApiResponse::success(health)))
516 }
517 }
518}
519
520fn assess_system_health(stats: &[hammerwork::stats::QueueStats]) -> SystemHealth {
522 let mut alerts = Vec::new();
523 let mut high_error_rate = false;
524 let mut queue_backlog = false;
525 let mut slow_processing = false;
526
527 for stat in stats {
528 if stat.statistics.error_rate > 0.1 {
530 high_error_rate = true;
532 alerts.push(SystemAlert {
533 severity: "warning".to_string(),
534 message: format!("High error rate in queue '{}'", stat.queue_name),
535 queue: Some(stat.queue_name.clone()),
536 metric: Some("error_rate".to_string()),
537 value: Some(stat.statistics.error_rate),
538 threshold: Some(0.1),
539 timestamp: chrono::Utc::now(),
540 });
541 }
542
543 if stat.pending_count > 1000 {
545 queue_backlog = true;
546 alerts.push(SystemAlert {
547 severity: "warning".to_string(),
548 message: format!("Large backlog in queue '{}'", stat.queue_name),
549 queue: Some(stat.queue_name.clone()),
550 metric: Some("pending_count".to_string()),
551 value: Some(stat.pending_count as f64),
552 threshold: Some(1000.0),
553 timestamp: chrono::Utc::now(),
554 });
555 }
556
557 if stat.statistics.avg_processing_time_ms > 30000.0 {
559 slow_processing = true;
561 alerts.push(SystemAlert {
562 severity: "info".to_string(),
563 message: format!("Slow processing in queue '{}'", stat.queue_name),
564 queue: Some(stat.queue_name.clone()),
565 metric: Some("avg_processing_time_ms".to_string()),
566 value: Some(stat.statistics.avg_processing_time_ms),
567 threshold: Some(30000.0),
568 timestamp: chrono::Utc::now(),
569 });
570 }
571 }
572
573 let status = if alerts.iter().any(|a| a.severity == "critical") {
574 "critical"
575 } else if alerts.iter().any(|a| a.severity == "warning") {
576 "degraded"
577 } else {
578 "healthy"
579 };
580
581 SystemHealth {
582 status: status.to_string(),
583 database_healthy: true, high_error_rate,
585 queue_backlog,
586 slow_processing,
587 alerts,
588 }
589}
590
591fn generate_overview_from_stats(
593 stats: &[hammerwork::stats::QueueStats],
594 uptime_seconds: u64,
595) -> SystemOverview {
596 let mut total_pending = 0;
597 let mut total_running = 0;
598 let mut total_completed = 0;
599 let mut total_failed = 0;
600 let mut total_dead = 0;
601 let mut total_throughput = 0.0;
602 let mut total_processing_time = 0.0;
603 let queue_count = stats.len();
604
605 for stat in stats {
606 total_pending += stat.pending_count;
607 total_running += stat.running_count;
608 total_completed += stat.completed_count;
609 total_failed += stat.dead_count + stat.timed_out_count;
610 total_dead += stat.dead_count;
611 total_throughput += stat.statistics.throughput_per_minute;
612 total_processing_time += stat.statistics.avg_processing_time_ms;
613 }
614
615 let avg_processing_time = if queue_count > 0 {
616 total_processing_time / queue_count as f64
617 } else {
618 0.0
619 };
620
621 let total_jobs = total_pending + total_running + total_completed + total_failed;
622 let overall_error_rate = if total_jobs > 0 {
623 total_failed as f64 / total_jobs as f64
624 } else {
625 0.0
626 };
627
628 let health = assess_system_health(stats);
629
630 SystemOverview {
631 total_queues: queue_count as u32,
632 total_jobs,
633 pending_jobs: total_pending,
634 running_jobs: total_running,
635 completed_jobs: total_completed,
636 failed_jobs: total_failed,
637 dead_jobs: total_dead,
638 overall_throughput: total_throughput,
639 overall_error_rate,
640 avg_processing_time_ms: avg_processing_time,
641 system_health: health,
642 uptime_seconds,
643 last_updated: chrono::Utc::now(),
644 }
645}
646
647async fn calculate_oldest_pending_age<T>(
649 queue: &Arc<T>,
650 queue_name: &str,
651) -> hammerwork::Result<Option<u64>>
652where
653 T: DatabaseQueue + Send + Sync,
654{
655 let jobs = queue.get_ready_jobs(queue_name, 100).await?;
657 let now = chrono::Utc::now();
658 Ok(jobs
659 .iter()
660 .filter(|job| matches!(job.status, hammerwork::job::JobStatus::Pending))
661 .map(|job| age_seconds(now, job.created_at))
662 .max())
663}
664
665fn age_seconds(
668 now: chrono::DateTime<chrono::Utc>,
669 created_at: chrono::DateTime<chrono::Utc>,
670) -> u64 {
671 u64::try_from((now - created_at).num_seconds()).unwrap_or(0)
672}
673
674async fn get_priority_distribution<T>(
676 queue: &Arc<T>,
677 queue_name: &str,
678) -> hammerwork::Result<HashMap<String, f32>>
679where
680 T: DatabaseQueue + Send + Sync,
681{
682 let priority_stats = queue.get_priority_stats(queue_name).await?;
683 Ok(priority_stats
684 .priority_distribution
685 .into_iter()
686 .map(|(priority, percentage)| {
687 let priority_name = match priority {
688 hammerwork::priority::JobPriority::Background => "background",
689 hammerwork::priority::JobPriority::Low => "low",
690 hammerwork::priority::JobPriority::Normal => "normal",
691 hammerwork::priority::JobPriority::High => "high",
692 hammerwork::priority::JobPriority::Critical => "critical",
693 };
694 (priority_name.to_string(), percentage)
695 })
696 .collect())
697}
698
699async fn compute_trends<T>(
701 queue: &T,
702 start: chrono::DateTime<chrono::Utc>,
703 end: chrono::DateTime<chrono::Utc>,
704) -> hammerwork::Result<Vec<HourlyTrend>>
705where
706 T: JobHistory,
707{
708 let buckets = queue.hourly_activity(None, start, end).await?;
709 Ok(buckets.into_iter().map(trend_from_bucket).collect())
710}
711
712fn trend_from_bucket(bucket: super::history::HourBucket) -> HourlyTrend {
713 let total = bucket.completed + bucket.failed;
714 HourlyTrend {
715 hour: bucket.hour,
716 completed: bucket.completed,
717 failed: bucket.failed,
718 throughput: total as f64 / 3600.0,
719 avg_processing_time_ms: bucket.avg_processing_time_ms,
720 error_rate: if total > 0 {
721 bucket.failed as f64 / total as f64
722 } else {
723 0.0
724 },
725 }
726}
727
728async fn generate_error_patterns<T>(
733 queue: &T,
734 since: chrono::DateTime<chrono::Utc>,
735) -> hammerwork::Result<Vec<ErrorPattern>>
736where
737 T: JobHistory,
738{
739 Ok(group_error_patterns(queue.error_groups(since, 500).await?))
740}
741
742fn group_error_patterns(groups: Vec<ErrorGroup>) -> Vec<ErrorPattern> {
745 use std::collections::{BTreeMap, BTreeSet};
746
747 struct Agg {
748 count: u64,
749 sample: (u64, String),
750 first_seen: chrono::DateTime<chrono::Utc>,
751 last_seen: chrono::DateTime<chrono::Utc>,
752 queues: BTreeSet<String>,
753 }
754
755 let total: u64 = groups.iter().map(|g| g.count).sum();
756 let mut by_type: BTreeMap<String, Agg> = BTreeMap::new();
757 for group in groups {
758 let error_type = extract_error_type(&group.message);
759 let agg = by_type.entry(error_type).or_insert_with(|| Agg {
760 count: 0,
761 sample: (0, group.message.clone()),
762 first_seen: group.first_seen,
763 last_seen: group.last_seen,
764 queues: BTreeSet::new(),
765 });
766 agg.count += group.count;
767 if group.count > agg.sample.0 {
768 agg.sample = (group.count, group.message.clone());
769 }
770 agg.first_seen = agg.first_seen.min(group.first_seen);
771 agg.last_seen = agg.last_seen.max(group.last_seen);
772 agg.queues.insert(group.queue_name);
773 }
774
775 let mut patterns: Vec<ErrorPattern> = by_type
776 .into_iter()
777 .map(|(error_type, agg)| ErrorPattern {
778 error_type,
779 count: agg.count,
780 percentage: if total > 0 {
781 agg.count as f64 / total as f64 * 100.0
782 } else {
783 0.0
784 },
785 sample_message: agg.sample.1,
786 first_seen: agg.first_seen,
787 last_seen: agg.last_seen,
788 affected_queues: agg.queues.into_iter().collect(),
789 })
790 .collect();
791 patterns.sort_by(|a, b| {
792 b.count
793 .cmp(&a.count)
794 .then_with(|| a.error_type.cmp(&b.error_type))
795 });
796 patterns
797}
798
799fn calculate_performance_metrics(
805 all_stats: &[hammerwork::stats::QueueStats],
806 database_response_time_ms: f64,
807) -> PerformanceMetrics {
808 let total_throughput = all_stats
809 .iter()
810 .map(|s| s.statistics.throughput_per_minute)
811 .sum::<f64>();
812
813 let average_queue_depth = if !all_stats.is_empty() {
814 all_stats
815 .iter()
816 .map(|s| s.pending_count as f64)
817 .sum::<f64>()
818 / all_stats.len() as f64
819 } else {
820 0.0
821 };
822
823 PerformanceMetrics {
824 database_response_time_ms: Some(database_response_time_ms),
825 average_queue_depth,
826 jobs_per_second: total_throughput / 60.0, memory_usage_mb: None,
828 cpu_usage_percent: None,
829 active_workers: None,
830 worker_utilization: None,
831 }
832}
833
834fn extract_error_type(error_msg: &str) -> String {
836 if error_msg.contains("timeout") || error_msg.contains("Timeout") {
838 "Timeout Error".to_string()
839 } else if error_msg.contains("connection") || error_msg.contains("Connection") {
840 "Connection Error".to_string()
841 } else if error_msg.contains("parse")
842 || error_msg.contains("Parse")
843 || error_msg.contains("invalid")
844 {
845 "Parse Error".to_string()
846 } else if error_msg.contains("permission")
847 || error_msg.contains("Permission")
848 || error_msg.contains("forbidden")
849 {
850 "Permission Error".to_string()
851 } else if error_msg.contains("not found") || error_msg.contains("Not Found") {
852 "Not Found Error".to_string()
853 } else {
854 error_msg
856 .split_whitespace()
857 .next()
858 .map(|s| format!("{} Error", s))
859 .unwrap_or_else(|| "Unknown Error".to_string())
860 }
861}
862
863#[cfg(test)]
864mod tests {
865 use super::*;
866 use crate::api::test_support::{body_json, unreachable_queue};
867
868 fn state() -> Arc<tokio::sync::RwLock<crate::api::system::SystemState>> {
869 Arc::new(tokio::sync::RwLock::new(
870 crate::api::system::SystemState::new(
871 crate::DashboardConfig::default(),
872 "PostgreSQL".to_string(),
873 1,
874 ),
875 ))
876 }
877
878 fn empty_query() -> StatsQuery {
879 serde_json::from_value(serde_json::json!({})).unwrap()
880 }
881
882 #[tokio::test]
883 async fn test_detailed_stats_returns_500_when_database_is_down() {
884 let response = detailed_stats_handler(unreachable_queue(), state(), empty_query())
885 .await
886 .unwrap()
887 .into_response();
888 let (status, body) = body_json(response).await;
889 assert_eq!(status, 500);
890 assert_eq!(body["success"], false);
891 }
892
893 #[tokio::test]
894 async fn test_trends_returns_500_when_database_is_down() {
895 let response = trends_handler(unreachable_queue(), empty_query())
896 .await
897 .unwrap()
898 .into_response();
899 let (status, body) = body_json(response).await;
900 assert_eq!(status, 500);
901 assert_eq!(body["success"], false);
902 }
903
904 #[tokio::test]
905 async fn test_trends_rejects_invalid_time_range() {
906 let inverted: StatsQuery = serde_json::from_value(serde_json::json!({
907 "time_range": {"start": "2024-01-02T00:00:00Z", "end": "2024-01-01T00:00:00Z"}
908 }))
909 .unwrap();
910 let response = trends_handler(unreachable_queue(), inverted)
911 .await
912 .unwrap()
913 .into_response();
914 assert_eq!(body_json(response).await.0, 400);
915
916 let too_long: StatsQuery = serde_json::from_value(serde_json::json!({
917 "time_range": {"start": "2020-01-01T00:00:00Z", "end": "2024-01-01T00:00:00Z"}
918 }))
919 .unwrap();
920 let response = trends_handler(unreachable_queue(), too_long)
921 .await
922 .unwrap()
923 .into_response();
924 assert_eq!(body_json(response).await.0, 400);
925 }
926
927 #[test]
928 fn test_default_range_is_24_hourly_buckets() {
929 let now = chrono::Utc::now();
930 let (start, end) = resolve_range(&empty_query(), now).unwrap();
931 assert_eq!(end, now);
932 assert_eq!(crate::api::history::hour_range(start, end).len(), 24);
933 }
934
935 #[test]
936 fn test_trend_from_bucket_computes_rates_without_inventing_values() {
937 let hour = hour_floor(chrono::Utc::now());
938 let busy = trend_from_bucket(crate::api::history::HourBucket {
939 hour,
940 completed: 3,
941 failed: 1,
942 avg_processing_time_ms: Some(40.0),
943 });
944 assert_eq!(busy.error_rate, 0.25);
945 assert_eq!(busy.throughput, 4.0 / 3600.0);
946 let idle = trend_from_bucket(crate::api::history::HourBucket {
947 hour,
948 completed: 0,
949 failed: 0,
950 avg_processing_time_ms: None,
951 });
952 assert_eq!(idle.error_rate, 0.0);
953 assert_eq!(idle.avg_processing_time_ms, None);
954 }
955
956 #[test]
957 fn test_group_error_patterns_uses_real_times_and_queues() {
958 let now = chrono::Utc::now();
959 let t = |h| now - chrono::Duration::hours(h);
960 let group = |queue: &str, msg: &str, count, first, last| ErrorGroup {
961 queue_name: queue.to_string(),
962 message: msg.to_string(),
963 count,
964 first_seen: t(first),
965 last_seen: t(last),
966 };
967 let patterns = group_error_patterns(vec![
968 group("emails", "connection refused by host a", 6, 10, 2),
969 group("reports", "connection reset", 2, 8, 1),
970 group("emails", "request timeout", 2, 5, 3),
971 ]);
972 assert_eq!(patterns.len(), 2);
973 let conn = &patterns[0];
974 assert_eq!(conn.error_type, "Connection Error");
975 assert_eq!(conn.count, 8);
976 assert_eq!(conn.percentage, 80.0);
977 assert_eq!(conn.sample_message, "connection refused by host a");
978 assert_eq!(conn.affected_queues, vec!["emails", "reports"]);
979 assert_eq!(conn.first_seen, t(10));
980 assert_eq!(conn.last_seen, t(1));
981 assert_eq!(patterns[1].affected_queues, vec!["emails"]);
982 assert!(group_error_patterns(Vec::new()).is_empty());
983 }
984
985 #[test]
986 fn test_performance_metrics_report_unknowns_as_none() {
987 let metrics = calculate_performance_metrics(&[], 3.5);
988 assert_eq!(metrics.database_response_time_ms, Some(3.5));
989 assert_eq!(metrics.active_workers, None);
990 assert_eq!(metrics.worker_utilization, None);
991 assert_eq!(metrics.memory_usage_mb, None);
992 }
993
994 #[test]
995 fn test_age_seconds_clamps_future_timestamps_to_zero() {
996 let now = chrono::Utc::now();
997 assert_eq!(age_seconds(now, now + chrono::Duration::seconds(30)), 0);
998 assert_eq!(age_seconds(now, now - chrono::Duration::seconds(30)), 30);
999 }
1000
1001 #[test]
1002 fn test_stats_query_deserialization() {
1003 let json = r#"{
1004 "time_range": {
1005 "start": "2024-01-01T00:00:00Z",
1006 "end": "2024-01-02T00:00:00Z"
1007 },
1008 "queues": ["email", "data-processing"],
1009 "granularity": "hour"
1010 }"#;
1011
1012 let query: StatsQuery = serde_json::from_str(json).unwrap();
1013 assert!(query.time_range.is_some());
1014 assert_eq!(query.queues.as_ref().unwrap().len(), 2);
1015 assert_eq!(query.granularity, Some("hour".to_string()));
1016 }
1017
1018 #[test]
1019 fn test_system_alert_serialization() {
1020 let alert = SystemAlert {
1021 severity: "warning".to_string(),
1022 message: "High error rate detected".to_string(),
1023 queue: Some("email".to_string()),
1024 metric: Some("error_rate".to_string()),
1025 value: Some(0.15),
1026 threshold: Some(0.1),
1027 timestamp: chrono::Utc::now(),
1028 };
1029
1030 let json = serde_json::to_string(&alert).unwrap();
1031 assert!(json.contains("warning"));
1032 assert!(json.contains("High error rate"));
1033 }
1034}