Skip to main content

allsource_core/application/services/
analytics.rs

1use crate::{
2    domain::entities::Event,
3    error::{AllSourceError, Result},
4    store::EventStore,
5};
6use chrono::{DateTime, Datelike, Duration, Timelike, Utc};
7use serde::{Deserialize, Serialize};
8use std::collections::HashMap;
9
10/// Time window granularity for analytics
11#[derive(Debug, Clone, Copy, Deserialize, Serialize)]
12#[serde(rename_all = "lowercase")]
13pub enum TimeWindow {
14    Minute,
15    Hour,
16    Day,
17    Week,
18    Month,
19}
20
21impl TimeWindow {
22    pub fn duration(&self) -> Duration {
23        match self {
24            TimeWindow::Minute => Duration::minutes(1),
25            TimeWindow::Hour => Duration::hours(1),
26            TimeWindow::Day => Duration::days(1),
27            TimeWindow::Week => Duration::weeks(1),
28            TimeWindow::Month => Duration::days(30),
29        }
30    }
31
32    pub fn truncate(&self, timestamp: DateTime<Utc>) -> DateTime<Utc> {
33        match self {
34            TimeWindow::Minute => timestamp
35                .with_second(0)
36                .unwrap()
37                .with_nanosecond(0)
38                .unwrap(),
39            TimeWindow::Hour => timestamp
40                .with_minute(0)
41                .unwrap()
42                .with_second(0)
43                .unwrap()
44                .with_nanosecond(0)
45                .unwrap(),
46            TimeWindow::Day => timestamp
47                .with_hour(0)
48                .unwrap()
49                .with_minute(0)
50                .unwrap()
51                .with_second(0)
52                .unwrap()
53                .with_nanosecond(0)
54                .unwrap(),
55            TimeWindow::Week => {
56                let days_from_monday = timestamp.weekday().num_days_from_monday();
57                (timestamp - Duration::days(i64::from(days_from_monday)))
58                    .with_hour(0)
59                    .unwrap()
60                    .with_minute(0)
61                    .unwrap()
62                    .with_second(0)
63                    .unwrap()
64                    .with_nanosecond(0)
65                    .unwrap()
66            }
67            TimeWindow::Month => timestamp
68                .with_day(1)
69                .unwrap()
70                .with_hour(0)
71                .unwrap()
72                .with_minute(0)
73                .unwrap()
74                .with_second(0)
75                .unwrap()
76                .with_nanosecond(0)
77                .unwrap(),
78        }
79    }
80}
81
82/// Request for event frequency analysis
83#[derive(Debug, Deserialize)]
84pub struct EventFrequencyRequest {
85    /// Filter by entity ID
86    pub entity_id: Option<String>,
87
88    /// Filter by event type
89    pub event_type: Option<String>,
90
91    /// Start time for analysis
92    pub since: DateTime<Utc>,
93
94    /// End time for analysis (defaults to now)
95    pub until: Option<DateTime<Utc>>,
96
97    /// Time window granularity
98    pub window: TimeWindow,
99}
100
101/// Time bucket with event count
102#[derive(Debug, Clone, Serialize)]
103pub struct TimeBucket {
104    pub timestamp: DateTime<Utc>,
105    pub count: usize,
106    pub event_types: HashMap<String, usize>,
107}
108
109/// Response containing time-series frequency data
110#[derive(Debug, Serialize)]
111pub struct EventFrequencyResponse {
112    pub buckets: Vec<TimeBucket>,
113    pub total_events: usize,
114    pub window: TimeWindow,
115    pub time_range: TimeRange,
116}
117
118#[derive(Debug, Serialize)]
119pub struct TimeRange {
120    pub from: DateTime<Utc>,
121    pub to: DateTime<Utc>,
122}
123
124/// Request for statistical summary
125#[derive(Debug, Deserialize)]
126pub struct StatsSummaryRequest {
127    /// Filter by entity ID
128    pub entity_id: Option<String>,
129
130    /// Filter by event type
131    pub event_type: Option<String>,
132
133    /// Start time for analysis
134    pub since: Option<DateTime<Utc>>,
135
136    /// End time for analysis
137    pub until: Option<DateTime<Utc>>,
138}
139
140/// Statistical summary response
141#[derive(Debug, Serialize)]
142pub struct StatsSummaryResponse {
143    pub total_events: usize,
144    pub unique_entities: usize,
145    pub unique_event_types: usize,
146    pub time_range: TimeRange,
147    pub events_per_day: f64,
148    pub top_event_types: Vec<EventTypeCount>,
149    pub top_entities: Vec<EntityCount>,
150    pub first_event: Option<DateTime<Utc>>,
151    pub last_event: Option<DateTime<Utc>>,
152}
153
154#[derive(Debug, Serialize)]
155pub struct EventTypeCount {
156    pub event_type: String,
157    pub count: usize,
158    pub percentage: f64,
159}
160
161#[derive(Debug, Serialize)]
162pub struct EntityCount {
163    pub entity_id: String,
164    pub count: usize,
165    pub percentage: f64,
166}
167
168/// Request for event correlation analysis
169#[derive(Debug, Deserialize)]
170pub struct CorrelationRequest {
171    /// First event type
172    pub event_type_a: String,
173
174    /// Second event type
175    pub event_type_b: String,
176
177    /// Maximum time window to consider events correlated
178    pub time_window_seconds: i64,
179
180    /// Start time for analysis
181    pub since: Option<DateTime<Utc>>,
182
183    /// End time for analysis
184    pub until: Option<DateTime<Utc>>,
185}
186
187/// Correlation analysis response
188#[derive(Debug, Serialize)]
189pub struct CorrelationResponse {
190    pub event_type_a: String,
191    pub event_type_b: String,
192    pub total_a: usize,
193    pub total_b: usize,
194    pub correlated_pairs: usize,
195    pub correlation_percentage: f64,
196    pub avg_time_between_seconds: f64,
197    pub examples: Vec<CorrelationExample>,
198}
199
200#[derive(Debug, Serialize)]
201pub struct CorrelationExample {
202    pub entity_id: String,
203    pub event_a_timestamp: DateTime<Utc>,
204    pub event_b_timestamp: DateTime<Utc>,
205    pub time_between_seconds: i64,
206}
207
208/// Analytics engine for time-series and statistical analysis
209pub struct AnalyticsEngine;
210
211impl AnalyticsEngine {
212    /// Analyze event frequency over time windows
213    pub fn event_frequency(
214        store: &EventStore,
215        request: &EventFrequencyRequest,
216    ) -> Result<EventFrequencyResponse> {
217        let until = request.until.unwrap_or_else(Utc::now);
218
219        // Query events in the time range
220        let events = store.query(&crate::application::dto::QueryEventsRequest {
221            entity_id: request.entity_id.clone(),
222            event_type: request.event_type.clone(),
223            tenant_id: None,
224            as_of: None,
225            since: Some(request.since),
226            until: Some(until),
227            limit: None,
228            event_type_prefix: None,
229            exclude_event_type_prefix: None,
230            payload_filter: None,
231        })?;
232
233        if events.is_empty() {
234            return Ok(EventFrequencyResponse {
235                buckets: Vec::new(),
236                total_events: 0,
237                window: request.window,
238                time_range: TimeRange {
239                    from: request.since,
240                    to: until,
241                },
242            });
243        }
244
245        // Create time buckets
246        let mut buckets_map: HashMap<DateTime<Utc>, HashMap<String, usize>> = HashMap::new();
247
248        for event in &events {
249            let bucket_time = request.window.truncate(event.timestamp);
250            let bucket = buckets_map.entry(bucket_time).or_default();
251            *bucket
252                .entry(event.event_type_str().to_string())
253                .or_insert(0) += 1;
254        }
255
256        // Convert to sorted vector
257        let mut buckets: Vec<TimeBucket> = buckets_map
258            .into_iter()
259            .map(|(timestamp, event_types)| {
260                let count = event_types.values().sum();
261                TimeBucket {
262                    timestamp,
263                    count,
264                    event_types,
265                }
266            })
267            .collect();
268
269        buckets.sort_by_key(|b| b.timestamp);
270
271        // Fill gaps in the timeline
272        let filled_buckets = Self::fill_time_gaps(&buckets, request.since, until, request.window);
273
274        Ok(EventFrequencyResponse {
275            total_events: events.len(),
276            buckets: filled_buckets,
277            window: request.window,
278            time_range: TimeRange {
279                from: request.since,
280                to: until,
281            },
282        })
283    }
284
285    /// Fill gaps in time buckets for continuous timeline
286    fn fill_time_gaps(
287        buckets: &[TimeBucket],
288        start: DateTime<Utc>,
289        end: DateTime<Utc>,
290        window: TimeWindow,
291    ) -> Vec<TimeBucket> {
292        if buckets.is_empty() {
293            return Vec::new();
294        }
295
296        let mut filled = Vec::new();
297        let mut current = window.truncate(start);
298        let end = window.truncate(end);
299
300        let bucket_map: HashMap<DateTime<Utc>, &TimeBucket> =
301            buckets.iter().map(|b| (b.timestamp, b)).collect();
302
303        while current <= end {
304            if let Some(bucket) = bucket_map.get(&current) {
305                filled.push((**bucket).clone());
306            } else {
307                filled.push(TimeBucket {
308                    timestamp: current,
309                    count: 0,
310                    event_types: HashMap::new(),
311                });
312            }
313            current += window.duration();
314        }
315
316        filled
317    }
318
319    /// Generate comprehensive statistical summary
320    pub fn stats_summary(
321        store: &EventStore,
322        request: &StatsSummaryRequest,
323    ) -> Result<StatsSummaryResponse> {
324        // Query events based on filters
325        let events = store.query(&crate::application::dto::QueryEventsRequest {
326            entity_id: request.entity_id.clone(),
327            event_type: request.event_type.clone(),
328            tenant_id: None,
329            as_of: None,
330            since: request.since,
331            until: request.until,
332            limit: None,
333            event_type_prefix: None,
334            exclude_event_type_prefix: None,
335            payload_filter: None,
336        })?;
337
338        if events.is_empty() {
339            return Err(AllSourceError::ValidationError(
340                "No events found for the specified criteria".to_string(),
341            ));
342        }
343
344        // Calculate statistics
345        let first_event = events.first().map(|e| e.timestamp);
346        let last_event = events.last().map(|e| e.timestamp);
347
348        let mut entity_counts: HashMap<String, usize> = HashMap::new();
349        let mut event_type_counts: HashMap<String, usize> = HashMap::new();
350
351        for event in &events {
352            *entity_counts
353                .entry(event.entity_id_str().to_string())
354                .or_insert(0) += 1;
355            *event_type_counts
356                .entry(event.event_type_str().to_string())
357                .or_insert(0) += 1;
358        }
359
360        // Calculate events per day
361        let time_span = if let (Some(first), Some(last)) = (first_event, last_event) {
362            (last - first).num_days().max(1) as f64
363        } else {
364            1.0
365        };
366
367        let events_per_day = events.len() as f64 / time_span;
368
369        // Top event types
370        let mut top_event_types: Vec<EventTypeCount> = event_type_counts
371            .into_iter()
372            .map(|(event_type, count)| EventTypeCount {
373                event_type,
374                count,
375                percentage: (count as f64 / events.len() as f64) * 100.0,
376            })
377            .collect();
378        top_event_types.sort_by_key(|x| std::cmp::Reverse(x.count));
379        top_event_types.truncate(10);
380
381        // Top entities
382        let mut top_entities: Vec<EntityCount> = entity_counts
383            .into_iter()
384            .map(|(entity_id, count)| EntityCount {
385                entity_id,
386                count,
387                percentage: (count as f64 / events.len() as f64) * 100.0,
388            })
389            .collect();
390        top_entities.sort_by_key(|x| std::cmp::Reverse(x.count));
391        top_entities.truncate(10);
392
393        let time_range = TimeRange {
394            from: first_event.unwrap_or_else(Utc::now),
395            to: last_event.unwrap_or_else(Utc::now),
396        };
397
398        Ok(StatsSummaryResponse {
399            total_events: events.len(),
400            unique_entities: top_entities.len(),
401            unique_event_types: top_event_types.len(),
402            time_range,
403            events_per_day,
404            top_event_types,
405            top_entities,
406            first_event,
407            last_event,
408        })
409    }
410
411    /// Analyze correlation between two event types
412    pub fn analyze_correlation(
413        store: &EventStore,
414        request: CorrelationRequest,
415    ) -> Result<CorrelationResponse> {
416        // Query both event types
417        let events_a = store.query(&crate::application::dto::QueryEventsRequest {
418            entity_id: None,
419            event_type: Some(request.event_type_a.clone()),
420            tenant_id: None,
421            as_of: None,
422            since: request.since,
423            until: request.until,
424            limit: None,
425            event_type_prefix: None,
426            exclude_event_type_prefix: None,
427            payload_filter: None,
428        })?;
429
430        let events_b = store.query(&crate::application::dto::QueryEventsRequest {
431            entity_id: None,
432            event_type: Some(request.event_type_b.clone()),
433            tenant_id: None,
434            as_of: None,
435            since: request.since,
436            until: request.until,
437            limit: None,
438            event_type_prefix: None,
439            exclude_event_type_prefix: None,
440            payload_filter: None,
441        })?;
442
443        // Group events by entity
444        let mut entity_events_a: HashMap<String, Vec<&Event>> = HashMap::new();
445        let mut entity_events_b: HashMap<String, Vec<&Event>> = HashMap::new();
446
447        for event in &events_a {
448            entity_events_a
449                .entry(event.entity_id_str().to_string())
450                .or_default()
451                .push(event);
452        }
453
454        for event in &events_b {
455            entity_events_b
456                .entry(event.entity_id_str().to_string())
457                .or_default()
458                .push(event);
459        }
460
461        // Find correlated pairs
462        let mut correlated_pairs = 0;
463        let mut total_time_between = 0i64;
464        let mut examples = Vec::new();
465
466        for (entity_id, a_events) in &entity_events_a {
467            if let Some(b_events) = entity_events_b.get(entity_id) {
468                for a_event in a_events {
469                    for b_event in b_events {
470                        let time_diff = (b_event.timestamp - a_event.timestamp).num_seconds().abs();
471
472                        if time_diff <= request.time_window_seconds {
473                            correlated_pairs += 1;
474                            total_time_between += time_diff;
475
476                            if examples.len() < 5 {
477                                examples.push(CorrelationExample {
478                                    entity_id: entity_id.clone(),
479                                    event_a_timestamp: a_event.timestamp,
480                                    event_b_timestamp: b_event.timestamp,
481                                    time_between_seconds: time_diff,
482                                });
483                            }
484                        }
485                    }
486                }
487            }
488        }
489
490        let correlation_percentage = if events_a.is_empty() {
491            0.0
492        } else {
493            (correlated_pairs as f64 / events_a.len() as f64) * 100.0
494        };
495
496        let avg_time_between = if correlated_pairs > 0 {
497            total_time_between as f64 / correlated_pairs as f64
498        } else {
499            0.0
500        };
501
502        Ok(CorrelationResponse {
503            event_type_a: request.event_type_a,
504            event_type_b: request.event_type_b,
505            total_a: events_a.len(),
506            total_b: events_b.len(),
507            correlated_pairs,
508            correlation_percentage,
509            avg_time_between_seconds: avg_time_between,
510            examples,
511        })
512    }
513}
514
515#[cfg(test)]
516mod tests {
517    use super::*;
518
519    // `GET /api/v1/analytics/frequency` takes `entity_id` and `event_type` as
520    // OPTIONAL filters, so the headline "how busy has this deployment been
521    // since T" call arrives with neither — and that is exactly the shape the
522    // store answers with a full scan, the one branch that never evaluated
523    // `since`/`until`. The response then reported the whole history's events
524    // under a `time_range` claiming the requested window: not an empty result a
525    // caller would notice, a wrong number that looks right.
526    #[test]
527    fn event_frequency_without_a_filter_still_respects_its_time_range() {
528        use crate::store::EventStore;
529
530        let store = EventStore::new();
531        let base = Utc::now() - chrono::Duration::hours(24);
532        for i in 0..6i64 {
533            let mut event = crate::domain::entities::Event::from_strings(
534                "user.created".to_string(),
535                format!("e-{i}"),
536                "default".to_string(),
537                serde_json::json!({}),
538                None,
539            )
540            .unwrap();
541            event.timestamp = base + chrono::Duration::hours(i);
542            event.version = i + 1;
543            store.ingest(&event).unwrap();
544        }
545
546        // No entity_id, no event_type: nothing for an index to narrow.
547        let response = AnalyticsEngine::event_frequency(
548            &store,
549            &EventFrequencyRequest {
550                entity_id: None,
551                event_type: None,
552                since: base + chrono::Duration::hours(2),
553                until: Some(base + chrono::Duration::hours(4)),
554                window: TimeWindow::Hour,
555            },
556        )
557        .unwrap();
558
559        assert_eq!(
560            response.total_events, 3,
561            "events at T+2, T+3 and T+4 are inside the range; the 6-event \
562             history is not"
563        );
564        let counted: usize = response.buckets.iter().map(|b| b.count).sum();
565        assert_eq!(
566            counted, 3,
567            "the buckets must add up to the events in the range"
568        );
569        assert!(
570            response.buckets.iter().all(|b| b.timestamp
571                >= TimeWindow::Hour.truncate(response.time_range.from)
572                && b.timestamp <= response.time_range.to),
573            "no bucket may fall outside the reported time_range"
574        );
575    }
576
577    #[test]
578    fn test_time_window_truncation() {
579        let timestamp = chrono::Utc::now();
580
581        let minute_truncated = TimeWindow::Minute.truncate(timestamp);
582        assert_eq!(minute_truncated.second(), 0);
583
584        let hour_truncated = TimeWindow::Hour.truncate(timestamp);
585        assert_eq!(hour_truncated.minute(), 0);
586        assert_eq!(hour_truncated.second(), 0);
587
588        let day_truncated = TimeWindow::Day.truncate(timestamp);
589        assert_eq!(day_truncated.hour(), 0);
590        assert_eq!(day_truncated.minute(), 0);
591    }
592}