Skip to main content

hammerwork_web/api/
stats.rs

1//! Statistics and monitoring API endpoints.
2//!
3//! This module provides comprehensive monitoring and analytics endpoints for tracking
4//! system health, performance metrics, and operational insights across all job queues.
5//!
6//! # API Endpoints
7//!
8//! - `GET /api/stats/overview` - System overview with key metrics
9//! - `GET /api/stats/detailed` - Detailed statistics with historical trends
10//! - `GET /api/stats/trends` - Hourly/daily trend analysis
11//! - `GET /api/stats/health` - System health check and alerts
12//!
13//! # Examples
14//!
15//! ## System Overview
16//!
17//! ```rust
18//! use hammerwork_web::api::stats::{SystemOverview, SystemHealth, SystemAlert};
19//! use chrono::Utc;
20//!
21//! let overview = SystemOverview {
22//!     total_queues: 5,
23//!     total_jobs: 10000,
24//!     pending_jobs: 50,
25//!     running_jobs: 10,
26//!     completed_jobs: 9800,
27//!     failed_jobs: 125,
28//!     dead_jobs: 15,
29//!     overall_throughput: 150.5,
30//!     overall_error_rate: 0.0125,
31//!     avg_processing_time_ms: 250.0,
32//!     system_health: SystemHealth {
33//!         status: "healthy".to_string(),
34//!         database_healthy: true,
35//!         high_error_rate: false,
36//!         queue_backlog: false,
37//!         slow_processing: false,
38//!         alerts: vec![],
39//!     },
40//!     uptime_seconds: 86400,
41//!     last_updated: Utc::now(),
42//! };
43//!
44//! assert_eq!(overview.total_queues, 5);
45//! assert_eq!(overview.overall_error_rate, 0.0125);
46//! assert_eq!(overview.system_health.status, "healthy");
47//! ```
48//!
49//! ## Statistics Queries
50//!
51//! ```rust
52//! use hammerwork_web::api::stats::{StatsQuery, TimeRange};
53//! use chrono::{Utc, Duration};
54//!
55//! let time_range = TimeRange {
56//!     start: Utc::now() - Duration::hours(24),
57//!     end: Utc::now(),
58//! };
59//!
60//! let query = StatsQuery {
61//!     time_range: Some(time_range),
62//!     queues: Some(vec!["email".to_string(), "notifications".to_string()]),
63//!     granularity: Some("hour".to_string()),
64//!     hours: None,
65//! };
66//!
67//! assert!(query.time_range.is_some());
68//! assert_eq!(query.queues.as_ref().unwrap().len(), 2);
69//! assert_eq!(query.granularity, Some("hour".to_string()));
70//! ```
71//!
72//! ## System Alerts
73//!
74//! ```rust
75//! use hammerwork_web::api::stats::SystemAlert;
76//! use chrono::Utc;
77//!
78//! let alert = SystemAlert {
79//!     severity: "warning".to_string(),
80//!     message: "Queue backlog detected".to_string(),
81//!     queue: Some("image_processing".to_string()),
82//!     metric: Some("pending_count".to_string()),
83//!     value: Some(1500.0),
84//!     threshold: Some(1000.0),
85//!     timestamp: Utc::now(),
86//! };
87//!
88//! assert_eq!(alert.severity, "warning");
89//! assert_eq!(alert.queue, Some("image_processing".to_string()));
90//! assert_eq!(alert.value, Some(1500.0));
91//! ```
92//!
93//! ## Performance Metrics
94//!
95//! ```rust
96//! use hammerwork_web::api::stats::PerformanceMetrics;
97//!
98//! let metrics = PerformanceMetrics {
99//!     database_response_time_ms: Some(5.2),
100//!     average_queue_depth: 15.5,
101//!     jobs_per_second: 8.3,
102//!     memory_usage_mb: Some(512.0),
103//!     cpu_usage_percent: Some(45.2),
104//!     active_workers: None, // not tracked: there is no worker registry
105//!     worker_utilization: None,
106//! };
107//!
108//! assert_eq!(metrics.database_response_time_ms, Some(5.2));
109//! assert_eq!(metrics.active_workers, None);
110//! ```
111
112use 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/// System overview statistics
122#[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/// System health status
140#[derive(Debug, Serialize)]
141pub struct SystemHealth {
142    pub status: String, // "healthy", "degraded", "critical"
143    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/// System alert
151#[derive(Debug, Serialize)]
152pub struct SystemAlert {
153    pub severity: String, // "info", "warning", "error", "critical"
154    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/// Detailed statistics for monitoring
163#[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/// Queue statistics
173#[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/// Hourly trend data
189#[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    /// Mean run time of jobs completed in the hour; `null` when none completed.
196    pub avg_processing_time_ms: Option<f64>,
197    pub error_rate: f64,
198}
199
200/// Error pattern analysis
201#[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/// Performance metrics
213#[derive(Debug, Serialize)]
214pub struct PerformanceMetrics {
215    /// Measured time to fetch the queue statistics for this response.
216    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    /// Always `null`: Hammerwork has no worker registry, so worker counts are unknown.
222    pub active_workers: Option<u32>,
223    /// Always `null`, see `active_workers`.
224    pub worker_utilization: Option<f64>,
225}
226
227/// Time range for statistics queries
228#[derive(Debug, Deserialize)]
229pub struct TimeRange {
230    pub start: chrono::DateTime<chrono::Utc>,
231    pub end: chrono::DateTime<chrono::Utc>,
232}
233
234/// Statistics query parameters
235#[derive(Debug, Deserialize)]
236pub struct StatsQuery {
237    pub time_range: Option<TimeRange>,
238    pub queues: Option<Vec<String>>,
239    pub granularity: Option<String>, // "hour", "day", "week"
240    /// The last N hours (alternative to `time_range`, which a query string cannot carry).
241    pub hours: Option<u32>,
242}
243
244/// Create statistics routes
245pub 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
290/// Handler for system overview statistics
291async 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            // Generate system health assessment
334            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
364/// Handler for detailed statistics
365async 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            // Convert hammerwork stats to our API format
384            let mut queue_stats: Vec<QueueStats> = Vec::new();
385            for stats in all_stats.iter() {
386                // Calculate oldest pending age seconds
387                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                // Get priority distribution from priority stats
393                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
443/// The `[start, end)` window of a stats request: `time_range` when given (validated),
444/// otherwise the last 24 hours (23 whole hours plus the current one).
445fn 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
472/// Handler for trend analysis: hourly completed/failed counts from the database.
473async 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
488/// Handler for system health check
489async 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
520/// Assess overall system health based on queue statistics
521fn 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        // Check error rate
529        if stat.statistics.error_rate > 0.1 {
530            // > 10% error rate
531            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        // Check queue backlog
544        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        // Check processing time
558        if stat.statistics.avg_processing_time_ms > 30000.0 {
559            // > 30 seconds
560            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, // If we got here, DB is accessible
584        high_error_rate,
585        queue_backlog,
586        slow_processing,
587        alerts,
588    }
589}
590
591/// Generate system overview from queue statistics
592fn 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
647/// Calculate the oldest pending job age in seconds for a queue
648async 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    // Get ready jobs (pending jobs) and find the oldest
656    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
665/// Age in whole seconds between `created_at` and `now`, clamped to zero when
666/// `created_at` is in the future (clock skew) instead of wrapping to ~1.8e19.
667fn 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
674/// Get priority distribution from priority stats for a queue
675async 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
699/// Hourly trends for `[start, end)` computed from the database (zero-filled).
700async 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
728/// Error patterns of failures since `since`, grouped by error type.
729///
730/// Built from the 500 most frequent distinct error messages, so counts and percentages
731/// are exact unless the failure set has more distinct messages than that.
732async 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
742/// Folds per-(queue, message) groups into one pattern per error type with real
743/// first/last-seen times and the set of affected queues.
744fn 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
799/// Calculate performance metrics from queue statistics.
800///
801/// `database_response_time_ms` is the measured time of the statistics query. Worker
802/// counts and utilization are `None`: running jobs are not workers, and there is no
803/// worker registry to count them from.
804fn 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, // Convert from per minute to per second
827        memory_usage_mb: None,
828        cpu_usage_percent: None,
829        active_workers: None,
830        worker_utilization: None,
831    }
832}
833
834/// Extract error type from error message for grouping
835fn extract_error_type(error_msg: &str) -> String {
836    // Simple error classification logic
837    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        // Use first word of error message as type
855        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}