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#[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#[derive(Debug, Deserialize)]
84pub struct EventFrequencyRequest {
85 pub entity_id: Option<String>,
87
88 pub event_type: Option<String>,
90
91 pub since: DateTime<Utc>,
93
94 pub until: Option<DateTime<Utc>>,
96
97 pub window: TimeWindow,
99}
100
101#[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#[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#[derive(Debug, Deserialize)]
126pub struct StatsSummaryRequest {
127 pub entity_id: Option<String>,
129
130 pub event_type: Option<String>,
132
133 pub since: Option<DateTime<Utc>>,
135
136 pub until: Option<DateTime<Utc>>,
138}
139
140#[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#[derive(Debug, Deserialize)]
170pub struct CorrelationRequest {
171 pub event_type_a: String,
173
174 pub event_type_b: String,
176
177 pub time_window_seconds: i64,
179
180 pub since: Option<DateTime<Utc>>,
182
183 pub until: Option<DateTime<Utc>>,
185}
186
187#[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
208pub struct AnalyticsEngine;
210
211impl AnalyticsEngine {
212 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 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 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 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 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 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(¤t) {
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 pub fn stats_summary(
321 store: &EventStore,
322 request: &StatsSummaryRequest,
323 ) -> Result<StatsSummaryResponse> {
324 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 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 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 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 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 pub fn analyze_correlation(
413 store: &EventStore,
414 request: CorrelationRequest,
415 ) -> Result<CorrelationResponse> {
416 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 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 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 #[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 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}