1use crate::event_sourcing::{EventStoreTrait, EventStream};
7use crate::StreamEvent;
8use anyhow::{anyhow, Result};
9use chrono::{DateTime, Duration as ChronoDuration, Utc};
10use serde::{Deserialize, Serialize};
11use std::collections::{BTreeMap, HashMap, HashSet};
12use std::sync::Arc;
13use std::time::{Duration, Instant};
14use tokio::sync::RwLock;
15use tracing::{debug, error, info, warn};
16use uuid::Uuid;
17
18pub type CustomFilterFn = Box<dyn Fn(&StreamEvent) -> bool + Send + Sync>;
20
21#[derive(Debug, Clone, Serialize, Deserialize)]
23pub struct TimeTravelConfig {
24 pub max_time_window_days: u32,
26 pub enable_temporal_indexing: bool,
28 pub index_granularity_minutes: u32,
30 pub max_concurrent_queries: usize,
32 pub query_timeout_seconds: u64,
34 pub enable_result_caching: bool,
36 pub cache_ttl_minutes: u32,
38 pub max_cache_size_mb: usize,
40}
41
42impl Default for TimeTravelConfig {
43 fn default() -> Self {
44 Self {
45 max_time_window_days: 365,
46 enable_temporal_indexing: true,
47 index_granularity_minutes: 60,
48 max_concurrent_queries: 100,
49 query_timeout_seconds: 300,
50 enable_result_caching: true,
51 cache_ttl_minutes: 60,
52 max_cache_size_mb: 1024,
53 }
54 }
55}
56
57#[derive(Debug, Clone, Serialize, Deserialize)]
59pub enum TimePoint {
60 Timestamp(DateTime<Utc>),
62 RelativeTime(ChronoDuration),
64 Version(u64),
66 EventId(Uuid),
68 Snapshot(String),
70}
71
72#[derive(Debug, Clone, Serialize, Deserialize)]
74pub struct TimeRange {
75 pub start: TimePoint,
76 pub end: TimePoint,
77}
78
79#[derive(Debug, Clone)]
81pub struct TemporalQuery {
82 pub query_id: Uuid,
83 pub time_point: Option<TimePoint>,
84 pub time_range: Option<TimeRange>,
85 pub filter: TemporalFilter,
86 pub projection: TemporalProjection,
87 pub ordering: TemporalOrdering,
88 pub limit: Option<usize>,
89}
90
91impl Default for TemporalQuery {
92 fn default() -> Self {
93 Self::new()
94 }
95}
96
97impl TemporalQuery {
98 pub fn new() -> Self {
100 Self {
101 query_id: Uuid::new_v4(),
102 time_point: None,
103 time_range: None,
104 filter: TemporalFilter::default(),
105 projection: TemporalProjection::default(),
106 ordering: TemporalOrdering::default(),
107 limit: None,
108 }
109 }
110
111 pub fn at_time(mut self, time_point: TimePoint) -> Self {
113 self.time_point = Some(time_point);
114 self
115 }
116
117 pub fn in_range(mut self, time_range: TimeRange) -> Self {
119 self.time_range = Some(time_range);
120 self
121 }
122
123 pub fn filter(mut self, filter: TemporalFilter) -> Self {
125 self.filter = filter;
126 self
127 }
128
129 pub fn project(mut self, projection: TemporalProjection) -> Self {
131 self.projection = projection;
132 self
133 }
134
135 pub fn order_by(mut self, ordering: TemporalOrdering) -> Self {
137 self.ordering = ordering;
138 self
139 }
140
141 pub fn limit(mut self, limit: usize) -> Self {
143 self.limit = Some(limit);
144 self
145 }
146}
147
148#[derive(Default)]
150pub struct TemporalFilter {
151 pub event_types: Option<HashSet<String>>,
152 pub aggregate_ids: Option<HashSet<String>>,
153 pub user_ids: Option<HashSet<String>>,
154 pub sources: Option<HashSet<String>>,
155 pub custom_filters: Vec<CustomFilterFn>,
156}
157
158impl std::fmt::Debug for TemporalFilter {
159 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
160 f.debug_struct("TemporalFilter")
161 .field("event_types", &self.event_types)
162 .field("aggregate_ids", &self.aggregate_ids)
163 .field("user_ids", &self.user_ids)
164 .field("sources", &self.sources)
165 .field(
166 "custom_filters",
167 &format!("<{} filters>", self.custom_filters.len()),
168 )
169 .finish()
170 }
171}
172
173impl Clone for TemporalFilter {
174 fn clone(&self) -> Self {
175 Self {
176 event_types: self.event_types.clone(),
177 aggregate_ids: self.aggregate_ids.clone(),
178 user_ids: self.user_ids.clone(),
179 sources: self.sources.clone(),
180 custom_filters: Vec::new(), }
182 }
183}
184
185#[derive(Debug, Clone, Default)]
187pub enum TemporalProjection {
188 #[default]
190 FullEvents,
191 MetadataOnly,
193 Fields(Vec<String>),
195 Aggregation(AggregationType),
197}
198
199#[derive(Debug, Clone)]
201pub enum AggregationType {
202 Count,
203 CountBy(String),
204 Timeline(ChronoDuration),
205 Statistics,
206}
207
208#[derive(Debug, Clone, Default)]
210pub enum TemporalOrdering {
211 TimeAscending,
213 #[default]
215 TimeDescending,
216 VersionAscending,
218 VersionDescending,
220 Custom(String, bool), }
223
224#[derive(Debug, Clone)]
226pub struct TemporalQueryResult {
227 pub query_id: Uuid,
228 pub events: Vec<StreamEvent>,
229 pub metadata: TemporalResultMetadata,
230 pub aggregations: Option<TemporalAggregations>,
231 pub execution_time: Duration,
232 pub from_cache: bool,
233}
234
235#[derive(Debug, Clone)]
237pub struct TemporalResultMetadata {
238 pub total_events: usize,
239 pub time_range_covered: Option<(DateTime<Utc>, DateTime<Utc>)>,
240 pub version_range_covered: Option<(u64, u64)>,
241 pub aggregates_scanned: HashSet<String>,
242 pub index_hits: usize,
243 pub index_misses: usize,
244}
245
246#[derive(Debug, Clone)]
248pub struct TemporalAggregations {
249 pub count: usize,
250 pub count_by_type: HashMap<String, usize>,
251 pub timeline: Vec<TimelinePoint>,
252 pub statistics: TemporalStatistics,
253}
254
255#[derive(Debug, Clone)]
257pub struct TimelinePoint {
258 pub timestamp: DateTime<Utc>,
259 pub count: usize,
260 pub event_types: HashMap<String, usize>,
261}
262
263#[derive(Debug, Clone)]
265pub struct TemporalStatistics {
266 pub events_per_second: f64,
267 pub peak_throughput: f64,
268 pub average_event_size: f64,
269 pub unique_aggregates: usize,
270 pub unique_users: usize,
271 pub time_span: ChronoDuration,
272}
273
274#[derive(Debug)]
276struct TemporalIndex {
277 time_index: BTreeMap<DateTime<Utc>, Vec<Uuid>>,
279 version_index: BTreeMap<u64, EventIndexEntry>,
281 aggregate_index: HashMap<String, BTreeMap<DateTime<Utc>, Vec<Uuid>>>,
283 type_index: HashMap<String, BTreeMap<DateTime<Utc>, Vec<Uuid>>>,
285}
286
287#[derive(Debug, Clone)]
288struct EventIndexEntry {
289 pub event_id: Uuid,
290 pub timestamp: DateTime<Utc>,
291 pub aggregate_id: String,
292 pub event_type: String,
293 pub version: u64,
294}
295
296impl TemporalIndex {
297 fn new() -> Self {
298 Self {
299 time_index: BTreeMap::new(),
300 version_index: BTreeMap::new(),
301 aggregate_index: HashMap::new(),
302 type_index: HashMap::new(),
303 }
304 }
305
306 fn add_event(&mut self, event: &StreamEvent) {
307 let metadata = event.metadata();
308 let timestamp = metadata.timestamp;
309 let event_id = uuid::Uuid::parse_str(&metadata.event_id).unwrap_or(uuid::Uuid::new_v4());
310 let aggregate_id = metadata.context.clone().unwrap_or_default();
311 let event_type = format!("{event:?}");
312 let version = metadata.version.parse::<u64>().unwrap_or(0);
313
314 self.time_index.entry(timestamp).or_default().push(event_id);
316
317 self.version_index.insert(
319 version,
320 EventIndexEntry {
321 event_id,
322 timestamp,
323 aggregate_id: aggregate_id.clone(),
324 event_type: event_type.clone(),
325 version,
326 },
327 );
328
329 self.aggregate_index
331 .entry(aggregate_id)
332 .or_default()
333 .entry(timestamp)
334 .or_default()
335 .push(event_id);
336
337 self.type_index
339 .entry(event_type)
340 .or_default()
341 .entry(timestamp)
342 .or_default()
343 .push(event_id);
344 }
345
346 fn find_events_by_time_range(&self, start: DateTime<Utc>, end: DateTime<Utc>) -> Vec<Uuid> {
347 let mut event_ids = Vec::new();
348
349 for (_, ids) in self.time_index.range(start..=end) {
350 event_ids.extend_from_slice(ids);
351 }
352
353 event_ids
354 }
355
356 fn find_events_by_version_range(&self, start: u64, end: u64) -> Vec<Uuid> {
357 let mut event_ids = Vec::new();
358
359 for (_, entry) in self.version_index.range(start..=end) {
360 event_ids.push(entry.event_id);
361 }
362
363 event_ids
364 }
365
366 fn find_events_by_aggregate(
367 &self,
368 aggregate_id: &str,
369 start: DateTime<Utc>,
370 end: DateTime<Utc>,
371 ) -> Vec<Uuid> {
372 if let Some(time_map) = self.aggregate_index.get(aggregate_id) {
373 let mut event_ids = Vec::new();
374 for (_, ids) in time_map.range(start..=end) {
375 event_ids.extend_from_slice(ids);
376 }
377 event_ids
378 } else {
379 Vec::new()
380 }
381 }
382}
383
384pub struct TimeTravelEngine {
386 config: TimeTravelConfig,
387 event_store: Arc<dyn EventStoreTrait>,
388 event_stream: Arc<dyn EventStream>,
389 temporal_index: Arc<RwLock<TemporalIndex>>,
390 query_cache: Arc<RwLock<QueryCache>>,
391 query_semaphore: Arc<tokio::sync::Semaphore>,
392 metrics: Arc<RwLock<TimeTravelMetrics>>,
393}
394
395impl TimeTravelEngine {
396 pub fn new(
398 config: TimeTravelConfig,
399 event_store: Arc<dyn EventStoreTrait>,
400 event_stream: Arc<dyn EventStream>,
401 ) -> Self {
402 Self {
403 query_semaphore: Arc::new(tokio::sync::Semaphore::new(config.max_concurrent_queries)),
404 temporal_index: Arc::new(RwLock::new(TemporalIndex::new())),
405 query_cache: Arc::new(RwLock::new(QueryCache::new(config.clone()))),
406 config,
407 event_store,
408 event_stream,
409 metrics: Arc::new(RwLock::new(TimeTravelMetrics::default())),
410 }
411 }
412
413 pub async fn start(&self) -> Result<()> {
415 info!("Starting time-travel engine");
416
417 if self.config.enable_temporal_indexing {
419 self.build_temporal_index().await?;
420 }
421
422 let index = Arc::clone(&self.temporal_index);
424 let event_stream = Arc::clone(&self.event_stream);
425
426 tokio::spawn(async move {
427 let mut interval = tokio::time::interval(Duration::from_secs(60));
428 loop {
429 interval.tick().await;
430 if let Err(e) =
431 Self::update_index(Arc::clone(&index), Arc::clone(&event_stream)).await
432 {
433 error!("Failed to update temporal index: {}", e);
434 }
435 }
436 });
437
438 info!("Time-travel engine started successfully");
439 Ok(())
440 }
441
442 pub async fn execute_query(&self, query: TemporalQuery) -> Result<TemporalQueryResult> {
444 let start_time = Instant::now();
445 let query_id = query.query_id;
446
447 debug!("Executing temporal query {}", query_id);
448
449 let _permit = self.query_semaphore.acquire().await?;
451
452 {
454 let mut metrics = self.metrics.write().await;
455 metrics.queries_executed += 1;
456 metrics.active_queries += 1;
457 }
458
459 let cache_key = self.generate_cache_key(&query);
461 if self.config.enable_result_caching {
462 let cache = self.query_cache.read().await;
463 if let Some(cached_result) = cache.get(&cache_key) {
464 let mut metrics = self.metrics.write().await;
465 metrics.active_queries -= 1;
466 metrics.cache_hits += 1;
467
468 return Ok(TemporalQueryResult {
469 query_id,
470 events: cached_result.events,
471 metadata: cached_result.metadata,
472 aggregations: cached_result.aggregations,
473 execution_time: start_time.elapsed(),
474 from_cache: true,
475 });
476 }
477 }
478
479 let result = self.execute_query_internal(query).await;
480
481 {
483 let mut metrics = self.metrics.write().await;
484 metrics.active_queries -= 1;
485 match &result {
486 Ok(_) => {
487 metrics.queries_succeeded += 1;
488 if !self.config.enable_result_caching {
489 metrics.cache_misses += 1;
490 }
491 }
492 Err(_) => metrics.queries_failed += 1,
493 }
494 }
495
496 let execution_time = start_time.elapsed();
497 debug!(
498 "Temporal query {} executed in {:?}",
499 query_id, execution_time
500 );
501
502 if let Ok(ref res) = result {
503 if self.config.enable_result_caching {
505 let mut cache = self.query_cache.write().await;
506 cache.set(cache_key, res.clone());
507 }
508 }
509
510 result.map(|mut r| {
511 r.execution_time = execution_time;
512 r.from_cache = false;
513 r
514 })
515 }
516
517 async fn execute_query_internal(&self, query: TemporalQuery) -> Result<TemporalQueryResult> {
519 let query_id = query.query_id;
520
521 let (start_time, end_time) = self.resolve_time_range(&query).await?;
523
524 let candidate_event_ids = if self.config.enable_temporal_indexing {
526 self.find_events_with_index(&query, start_time, end_time)
527 .await?
528 } else {
529 self.find_events_without_index(&query, start_time, end_time)
530 .await?
531 };
532
533 let mut events = Vec::new();
535 for event_id in candidate_event_ids {
536 if let Some(event) = self.load_event(event_id).await? {
537 if self.matches_filter(&event, &query.filter) {
538 events.push(event);
539 }
540 }
541 }
542
543 self.apply_ordering(&mut events, &query.ordering);
545
546 if let Some(limit) = query.limit {
548 events.truncate(limit);
549 }
550
551 let metadata = self.generate_result_metadata(&events, start_time, end_time);
553
554 let aggregations = match query.projection {
556 TemporalProjection::Aggregation(ref agg_type) => {
557 Some(self.generate_aggregations(&events, agg_type, start_time, end_time)?)
558 }
559 _ => None,
560 };
561
562 let projected_events = self.apply_projection(events, &query.projection);
564
565 Ok(TemporalQueryResult {
566 query_id,
567 events: projected_events,
568 metadata,
569 aggregations,
570 execution_time: Duration::default(), from_cache: false,
572 })
573 }
574
575 pub async fn query_state_at_time(
577 &self,
578 aggregate_id: &str,
579 time_point: TimePoint,
580 ) -> Result<Vec<StreamEvent>> {
581 let query = TemporalQuery::new()
582 .at_time(time_point)
583 .filter(TemporalFilter {
584 aggregate_ids: Some(std::iter::once(aggregate_id.to_string()).collect()),
585 ..Default::default()
586 });
587
588 let result = self.execute_query(query).await?;
589 Ok(result.events)
590 }
591
592 pub async fn query_changes_between(
594 &self,
595 start: TimePoint,
596 end: TimePoint,
597 filter: Option<TemporalFilter>,
598 ) -> Result<Vec<StreamEvent>> {
599 let query = TemporalQuery::new()
600 .in_range(TimeRange { start, end })
601 .filter(filter.unwrap_or_default());
602
603 let result = self.execute_query(query).await?;
604 Ok(result.events)
605 }
606
607 pub async fn query_timeline(
609 &self,
610 time_range: TimeRange,
611 granularity: ChronoDuration,
612 filter: Option<TemporalFilter>,
613 ) -> Result<Vec<TimelinePoint>> {
614 let query = TemporalQuery::new()
615 .in_range(time_range)
616 .filter(filter.unwrap_or_default())
617 .project(TemporalProjection::Aggregation(AggregationType::Timeline(
618 granularity,
619 )));
620
621 let result = self.execute_query(query).await?;
622 Ok(result.aggregations.map(|a| a.timeline).unwrap_or_default())
623 }
624
625 async fn build_temporal_index(&self) -> Result<()> {
627 info!("Building temporal index");
628
629 let events = self
630 .event_stream
631 .read_events_from_position(0, usize::MAX)
632 .await?;
633 let mut index = self.temporal_index.write().await;
634
635 for stored_event in events {
636 index.add_event(&stored_event.event_data);
637 }
638
639 info!(
640 "Temporal index built with {} events",
641 index.time_index.len()
642 );
643 Ok(())
644 }
645
646 async fn update_index(
648 index: Arc<RwLock<TemporalIndex>>,
649 event_stream: Arc<dyn EventStream>,
650 ) -> Result<()> {
651 let events = event_stream.read_events_from_position(0, 10000).await?;
654 let mut idx = index.write().await;
655
656 for stored_event in events {
657 idx.add_event(&stored_event.event_data);
658 }
659
660 Ok(())
661 }
662
663 async fn resolve_time_range(
665 &self,
666 query: &TemporalQuery,
667 ) -> Result<(DateTime<Utc>, DateTime<Utc>)> {
668 let now = Utc::now();
669
670 match (&query.time_point, &query.time_range) {
671 (Some(time_point), None) => {
672 let timestamp = self.resolve_time_point(time_point).await?;
673 Ok((timestamp, timestamp))
674 }
675 (None, Some(time_range)) => {
676 let start = self.resolve_time_point(&time_range.start).await?;
677 let end = self.resolve_time_point(&time_range.end).await?;
678 Ok((start, end))
679 }
680 (None, None) => {
681 let start = now - ChronoDuration::hours(24);
683 Ok((start, now))
684 }
685 (Some(_), Some(_)) => Err(anyhow!("Cannot specify both time_point and time_range")),
686 }
687 }
688
689 async fn resolve_time_point(&self, time_point: &TimePoint) -> Result<DateTime<Utc>> {
691 match time_point {
692 TimePoint::Timestamp(timestamp) => Ok(*timestamp),
693 TimePoint::RelativeTime(duration) => Ok(Utc::now() + *duration),
694 TimePoint::Version(version) => {
695 let index = self.temporal_index.read().await;
697 if let Some(entry) = index.version_index.get(version) {
698 Ok(entry.timestamp)
699 } else {
700 Err(anyhow!("Version {} not found", version))
701 }
702 }
703 TimePoint::EventId(event_id) => {
704 if let Some(event) = self.load_event(*event_id).await? {
706 Ok(event.metadata().timestamp)
707 } else {
708 Err(anyhow!("Event {} not found", event_id))
709 }
710 }
711 TimePoint::Snapshot(name) => {
712 Err(anyhow!(
717 "Named snapshot '{}' cannot be resolved: the TimeTravelEngine requires an \
718 integrated snapshot store for named-snapshot resolution. \
719 Convert the snapshot to a Timestamp or Version TimePoint.",
720 name
721 ))
722 }
723 }
724 }
725
726 async fn find_events_with_index(
728 &self,
729 query: &TemporalQuery,
730 start_time: DateTime<Utc>,
731 end_time: DateTime<Utc>,
732 ) -> Result<Vec<Uuid>> {
733 let index = self.temporal_index.read().await;
734
735 if let Some(ref aggregate_ids) = query.filter.aggregate_ids {
737 if aggregate_ids.len() == 1 {
738 let aggregate_id = aggregate_ids
739 .iter()
740 .next()
741 .expect("aggregate_ids validated to have exactly 1 element");
742 return Ok(index.find_events_by_aggregate(aggregate_id, start_time, end_time));
743 }
744 }
745
746 Ok(index.find_events_by_time_range(start_time, end_time))
747 }
748
749 async fn find_events_without_index(
757 &self,
758 query: &TemporalQuery,
759 start_time: DateTime<Utc>,
760 end_time: DateTime<Utc>,
761 ) -> Result<Vec<Uuid>> {
762 let stored_events = self
763 .event_stream
764 .read_events_from_position(0, usize::MAX)
765 .await?;
766
767 let aggregate_filter = query.filter.aggregate_ids.as_ref();
768
769 let mut event_ids = Vec::new();
770 for stored in &stored_events {
771 let metadata = stored.event_data.metadata();
772
773 if metadata.timestamp < start_time || metadata.timestamp > end_time {
774 continue;
775 }
776
777 if let Some(aggregate_ids) = aggregate_filter {
778 let matches = metadata
779 .context
780 .as_ref()
781 .is_some_and(|ctx| aggregate_ids.contains(ctx));
782 if !matches {
783 continue;
784 }
785 }
786
787 event_ids.push(Self::event_uuid(stored));
788 }
789
790 debug!(
791 "Sequential scan found {} candidate event(s) in range {}..{}",
792 event_ids.len(),
793 start_time,
794 end_time
795 );
796
797 Ok(event_ids)
798 }
799
800 fn event_uuid(stored: &crate::event_sourcing::StoredEvent) -> Uuid {
808 Uuid::parse_str(&stored.event_data.metadata().event_id).unwrap_or(stored.event_id)
809 }
810
811 async fn load_event(&self, event_id: Uuid) -> Result<Option<StreamEvent>> {
818 let stored_events = self
819 .event_stream
820 .read_events_from_position(0, usize::MAX)
821 .await?;
822
823 Ok(stored_events
824 .into_iter()
825 .find(|stored| Self::event_uuid(stored) == event_id)
826 .map(|stored| stored.event_data))
827 }
828
829 fn matches_filter(&self, event: &StreamEvent, filter: &TemporalFilter) -> bool {
831 let metadata = event.metadata();
832 let event_type_str = format!("{event:?}");
833
834 if let Some(ref event_types) = filter.event_types {
835 if !event_types.contains(&event_type_str) {
836 return false;
837 }
838 }
839
840 if let Some(ref aggregate_ids) = filter.aggregate_ids {
841 if let Some(ref context) = metadata.context {
842 if !aggregate_ids.contains(context) {
843 return false;
844 }
845 } else {
846 return false;
847 }
848 }
849
850 if let Some(ref user_ids) = filter.user_ids {
851 if let Some(ref user) = metadata.user {
852 if !user_ids.contains(user) {
853 return false;
854 }
855 } else {
856 return false;
857 }
858 }
859
860 if let Some(ref sources) = filter.sources {
861 if !sources.contains(&metadata.source) {
862 return false;
863 }
864 }
865
866 for custom_filter in &filter.custom_filters {
868 if !custom_filter(event) {
869 return false;
870 }
871 }
872
873 true
874 }
875
876 fn apply_ordering(&self, events: &mut [StreamEvent], ordering: &TemporalOrdering) {
878 match ordering {
879 TemporalOrdering::TimeAscending => {
880 events.sort_by_key(|a| a.metadata().timestamp);
881 }
882 TemporalOrdering::TimeDescending => {
883 events.sort_by_key(|b| std::cmp::Reverse(b.metadata().timestamp));
884 }
885 TemporalOrdering::VersionAscending => {
886 events.sort_by(|a, b| a.metadata().version.cmp(&b.metadata().version));
887 }
888 TemporalOrdering::VersionDescending => {
889 events.sort_by(|a, b| b.metadata().version.cmp(&a.metadata().version));
890 }
891 TemporalOrdering::Custom(_field, _ascending) => {
892 warn!("Custom ordering not implemented");
894 }
895 }
896 }
897
898 fn apply_projection(
900 &self,
901 events: Vec<StreamEvent>,
902 projection: &TemporalProjection,
903 ) -> Vec<StreamEvent> {
904 match projection {
905 TemporalProjection::FullEvents => events,
906 TemporalProjection::MetadataOnly => {
907 events
910 }
911 TemporalProjection::Fields(_fields) => {
912 warn!("Field projection not implemented");
914 events
915 }
916 TemporalProjection::Aggregation(_) => {
917 Vec::new()
919 }
920 }
921 }
922
923 fn generate_result_metadata(
925 &self,
926 events: &[StreamEvent],
927 _start_time: DateTime<Utc>,
928 _end_time: DateTime<Utc>,
929 ) -> TemporalResultMetadata {
930 let total_events = events.len();
931
932 let time_range_covered = if !events.is_empty() {
933 let min_time = events
934 .iter()
935 .map(|e| e.metadata().timestamp)
936 .min()
937 .expect("events validated to be non-empty");
938 let max_time = events
939 .iter()
940 .map(|e| e.metadata().timestamp)
941 .max()
942 .expect("events validated to be non-empty");
943 Some((min_time, max_time))
944 } else {
945 None
946 };
947
948 let version_range_covered = if !events.is_empty() {
949 let min_version = events
950 .iter()
951 .filter_map(|e| e.metadata().version.parse::<u64>().ok())
952 .min();
953 let max_version = events
954 .iter()
955 .filter_map(|e| e.metadata().version.parse::<u64>().ok())
956 .max();
957 if let (Some(min), Some(max)) = (min_version, max_version) {
958 Some((min, max))
959 } else {
960 None
961 }
962 } else {
963 None
964 };
965
966 let aggregates_scanned: HashSet<String> = events
967 .iter()
968 .filter_map(|e| e.metadata().context.clone())
969 .collect();
970
971 TemporalResultMetadata {
972 total_events,
973 time_range_covered,
974 version_range_covered,
975 aggregates_scanned,
976 index_hits: 0, index_misses: 0,
978 }
979 }
980
981 fn generate_aggregations(
983 &self,
984 events: &[StreamEvent],
985 agg_type: &AggregationType,
986 start_time: DateTime<Utc>,
987 end_time: DateTime<Utc>,
988 ) -> Result<TemporalAggregations> {
989 match agg_type {
990 AggregationType::Count => Ok(TemporalAggregations {
991 count: events.len(),
992 count_by_type: HashMap::new(),
993 timeline: Vec::new(),
994 statistics: self.calculate_statistics(events, start_time, end_time),
995 }),
996 AggregationType::CountBy(field) => {
997 let mut count_by_type = HashMap::new();
998 for event in events {
999 if field == "event_type" {
1000 let event_type = format!("{event:?}");
1001 *count_by_type.entry(event_type).or_insert(0) += 1;
1002 }
1003 }
1005
1006 Ok(TemporalAggregations {
1007 count: events.len(),
1008 count_by_type,
1009 timeline: Vec::new(),
1010 statistics: self.calculate_statistics(events, start_time, end_time),
1011 })
1012 }
1013 AggregationType::Timeline(granularity) => {
1014 let timeline = self.generate_timeline(events, *granularity, start_time, end_time);
1015
1016 Ok(TemporalAggregations {
1017 count: events.len(),
1018 count_by_type: HashMap::new(),
1019 timeline,
1020 statistics: self.calculate_statistics(events, start_time, end_time),
1021 })
1022 }
1023 AggregationType::Statistics => Ok(TemporalAggregations {
1024 count: events.len(),
1025 count_by_type: HashMap::new(),
1026 timeline: Vec::new(),
1027 statistics: self.calculate_statistics(events, start_time, end_time),
1028 }),
1029 }
1030 }
1031
1032 fn generate_timeline(
1034 &self,
1035 events: &[StreamEvent],
1036 granularity: ChronoDuration,
1037 start_time: DateTime<Utc>,
1038 end_time: DateTime<Utc>,
1039 ) -> Vec<TimelinePoint> {
1040 let mut timeline = Vec::new();
1041 let mut current_time = start_time;
1042
1043 while current_time < end_time {
1044 let window_end = current_time + granularity;
1045
1046 let events_in_window: Vec<_> = events
1047 .iter()
1048 .filter(|e| {
1049 e.metadata().timestamp >= current_time && e.metadata().timestamp < window_end
1050 })
1051 .collect();
1052
1053 let mut event_types = HashMap::new();
1054 for event in &events_in_window {
1055 let event_type = format!("{event:?}");
1056 *event_types.entry(event_type).or_insert(0) += 1;
1057 }
1058
1059 timeline.push(TimelinePoint {
1060 timestamp: current_time,
1061 count: events_in_window.len(),
1062 event_types,
1063 });
1064
1065 current_time = window_end;
1066 }
1067
1068 timeline
1069 }
1070
1071 fn calculate_statistics(
1073 &self,
1074 events: &[StreamEvent],
1075 start_time: DateTime<Utc>,
1076 end_time: DateTime<Utc>,
1077 ) -> TemporalStatistics {
1078 let time_span = end_time.signed_duration_since(start_time);
1079 let events_per_second = if time_span.num_seconds() > 0 {
1080 events.len() as f64 / time_span.num_seconds() as f64
1081 } else {
1082 0.0
1083 };
1084
1085 let peak_throughput = if !events.is_empty() {
1087 let mut minute_counts = HashMap::new();
1088 for event in events {
1089 let minute = event
1090 .metadata()
1091 .timestamp
1092 .format("%Y-%m-%d %H:%M")
1093 .to_string();
1094 *minute_counts.entry(minute).or_insert(0) += 1;
1095 }
1096 minute_counts.values().max().copied().unwrap_or(0) as f64
1097 } else {
1098 0.0
1099 };
1100
1101 let total_size: usize = events.iter().map(|e| format!("{e:?}").len()).sum();
1103 let average_event_size = if !events.is_empty() {
1104 total_size as f64 / events.len() as f64
1105 } else {
1106 0.0
1107 };
1108
1109 let unique_aggregates = events
1110 .iter()
1111 .filter_map(|e| e.metadata().context.as_ref())
1112 .collect::<HashSet<_>>()
1113 .len();
1114
1115 let unique_users = events
1116 .iter()
1117 .filter_map(|e| e.metadata().user.as_ref())
1118 .collect::<HashSet<_>>()
1119 .len();
1120
1121 TemporalStatistics {
1122 events_per_second,
1123 peak_throughput,
1124 average_event_size,
1125 unique_aggregates,
1126 unique_users,
1127 time_span,
1128 }
1129 }
1130
1131 fn generate_cache_key(&self, query: &TemporalQuery) -> String {
1133 format!("temporal_query_{:?}", query.query_id)
1135 }
1136
1137 pub async fn get_metrics(&self) -> TimeTravelMetrics {
1139 self.metrics.read().await.clone()
1140 }
1141}
1142
1143#[derive(Debug)]
1145struct QueryCache {
1146 config: TimeTravelConfig,
1147 entries: HashMap<String, CachedResult>,
1148}
1149
1150#[derive(Debug, Clone)]
1151struct CachedResult {
1152 events: Vec<StreamEvent>,
1153 metadata: TemporalResultMetadata,
1154 aggregations: Option<TemporalAggregations>,
1155 cached_at: DateTime<Utc>,
1156}
1157
1158impl QueryCache {
1159 fn new(config: TimeTravelConfig) -> Self {
1160 Self {
1161 config,
1162 entries: HashMap::new(),
1163 }
1164 }
1165
1166 fn get(&self, key: &str) -> Option<CachedResult> {
1167 if let Some(entry) = self.entries.get(key) {
1168 let age = Utc::now().signed_duration_since(entry.cached_at);
1169 if age.num_minutes() < self.config.cache_ttl_minutes as i64 {
1170 return Some(entry.clone());
1171 }
1172 }
1173 None
1174 }
1175
1176 fn set(&mut self, key: String, result: TemporalQueryResult) {
1177 let entry = CachedResult {
1178 events: result.events,
1179 metadata: result.metadata,
1180 aggregations: result.aggregations,
1181 cached_at: Utc::now(),
1182 };
1183
1184 self.entries.insert(key, entry);
1185 self.evict_if_needed();
1186 }
1187
1188 fn evict_if_needed(&mut self) {
1189 let now = Utc::now();
1191 self.entries.retain(|_, entry| {
1192 let age = now.signed_duration_since(entry.cached_at);
1193 age.num_minutes() < self.config.cache_ttl_minutes as i64
1194 });
1195
1196 while self.entries.len() > 1000 {
1198 if let Some(oldest_key) = self
1199 .entries
1200 .iter()
1201 .min_by_key(|(_, entry)| entry.cached_at)
1202 .map(|(key, _)| key.clone())
1203 {
1204 self.entries.remove(&oldest_key);
1205 } else {
1206 break;
1207 }
1208 }
1209 }
1210}
1211
1212#[derive(Debug, Clone, Default)]
1214pub struct TimeTravelMetrics {
1215 pub queries_executed: u64,
1216 pub queries_succeeded: u64,
1217 pub queries_failed: u64,
1218 pub active_queries: u64,
1219 pub cache_hits: u64,
1220 pub cache_misses: u64,
1221 pub index_hits: u64,
1222 pub index_misses: u64,
1223 pub average_query_time_ms: f64,
1224}
1225
1226#[cfg(test)]
1227mod tests {
1228 use super::*;
1229 use crate::event_sourcing::{EventQuery, EventSnapshot, StoredEvent};
1230 use crate::EventMetadata;
1231
1232 struct MockEventStream {
1237 events: Vec<StoredEvent>,
1238 }
1239
1240 #[async_trait::async_trait]
1241 impl EventStream for MockEventStream {
1242 async fn next_event(&mut self) -> Option<StoredEvent> {
1243 None
1244 }
1245
1246 async fn has_events(&self) -> bool {
1247 !self.events.is_empty()
1248 }
1249
1250 async fn read_events_from_position(
1251 &self,
1252 position: u64,
1253 max_events: usize,
1254 ) -> Result<Vec<StoredEvent>> {
1255 Ok(self
1256 .events
1257 .iter()
1258 .skip(position as usize)
1259 .take(max_events)
1260 .cloned()
1261 .collect())
1262 }
1263 }
1264
1265 struct MockEventStore;
1269
1270 #[async_trait::async_trait]
1271 impl EventStoreTrait for MockEventStore {
1272 async fn store_event(&self, stream_id: String, event: StreamEvent) -> Result<StoredEvent> {
1273 Ok(test_stored_event(&stream_id, event))
1274 }
1275
1276 async fn query_events(&self, _query: EventQuery) -> Result<Vec<StoredEvent>> {
1277 Ok(Vec::new())
1278 }
1279
1280 async fn get_stream_events(
1281 &self,
1282 _stream_id: &str,
1283 _from_version: Option<u64>,
1284 ) -> Result<Vec<StoredEvent>> {
1285 Ok(Vec::new())
1286 }
1287
1288 async fn replay_from_timestamp(
1289 &self,
1290 _timestamp: DateTime<Utc>,
1291 ) -> Result<Vec<StoredEvent>> {
1292 Ok(Vec::new())
1293 }
1294
1295 async fn get_latest_snapshot(&self, _stream_id: &str) -> Result<Option<EventSnapshot>> {
1296 Ok(None)
1297 }
1298
1299 async fn rebuild_stream_state(&self, _stream_id: &str) -> Result<Vec<u8>> {
1300 Ok(Vec::new())
1301 }
1302
1303 async fn append_events(
1304 &self,
1305 _aggregate_id: &str,
1306 _events: &[StreamEvent],
1307 _expected_version: Option<u64>,
1308 ) -> Result<u64> {
1309 Ok(0)
1310 }
1311 }
1312
1313 fn test_event(source: &str, timestamp: DateTime<Utc>) -> StreamEvent {
1314 StreamEvent::TripleAdded {
1315 subject: "http://example.org/s".to_string(),
1316 predicate: "http://example.org/p".to_string(),
1317 object: "\"o\"".to_string(),
1318 graph: None,
1319 metadata: EventMetadata {
1320 event_id: Uuid::new_v4().to_string(),
1321 timestamp,
1322 source: source.to_string(),
1323 user: None,
1324 context: None,
1325 caused_by: None,
1326 version: "1".to_string(),
1327 properties: HashMap::new(),
1328 checksum: None,
1329 },
1330 }
1331 }
1332
1333 fn test_stored_event(stream_id: &str, event_data: StreamEvent) -> StoredEvent {
1334 let event_id =
1335 Uuid::parse_str(&event_data.metadata().event_id).unwrap_or_else(|_| Uuid::new_v4());
1336 StoredEvent {
1337 event_id,
1338 sequence_number: 0,
1339 stream_id: stream_id.to_string(),
1340 stream_version: 0,
1341 event_data,
1342 stored_at: Utc::now(),
1343 storage_metadata: crate::event_sourcing::StorageMetadata {
1344 checksum: String::new(),
1345 compressed_size: None,
1346 original_size: 0,
1347 storage_location: "test".to_string(),
1348 persistence_status: crate::event_sourcing::PersistenceStatus::Persisted,
1349 },
1350 }
1351 }
1352
1353 fn make_engine(events: Vec<StreamEvent>) -> TimeTravelEngine {
1354 let stored: Vec<StoredEvent> = events
1355 .into_iter()
1356 .map(|e| test_stored_event("test-stream", e))
1357 .collect();
1358 TimeTravelEngine::new(
1359 TimeTravelConfig::default(),
1360 Arc::new(MockEventStore),
1361 Arc::new(MockEventStream { events: stored }),
1362 )
1363 }
1364
1365 #[tokio::test]
1370 async fn test_find_events_without_index_scans_real_events() {
1371 let now = Utc::now();
1372 let in_range = test_event("svc-a", now);
1373 let out_of_range = test_event("svc-a", now - ChronoDuration::days(30));
1374 let engine = make_engine(vec![in_range.clone(), out_of_range]);
1375
1376 let query = TemporalQuery::new();
1377 let start = now - ChronoDuration::hours(1);
1378 let end = now + ChronoDuration::hours(1);
1379
1380 let found = engine
1381 .find_events_without_index(&query, start, end)
1382 .await
1383 .expect("sequential scan must not error");
1384
1385 assert_eq!(
1386 found.len(),
1387 1,
1388 "sequential scan must find exactly the in-range event, not silently return empty"
1389 );
1390 }
1391
1392 #[tokio::test]
1397 async fn test_load_event_finds_real_event() {
1398 let now = Utc::now();
1399 let event = test_event("svc-a", now);
1400 let event_id =
1401 Uuid::parse_str(&event.metadata().event_id).expect("test event_id is a valid UUID");
1402 let engine = make_engine(vec![event]);
1403
1404 let loaded = engine
1405 .load_event(event_id)
1406 .await
1407 .expect("load_event must not error");
1408
1409 assert!(
1410 loaded.is_some(),
1411 "load_event must find a previously stored event by ID, not always return None"
1412 );
1413 assert_eq!(
1414 Uuid::parse_str(&loaded.expect("checked above").metadata().event_id)
1415 .expect("stored event_id is a valid UUID"),
1416 event_id
1417 );
1418
1419 let missing = engine
1422 .load_event(Uuid::new_v4())
1423 .await
1424 .expect("load_event must not error for a missing ID");
1425 assert!(missing.is_none());
1426 }
1427
1428 #[tokio::test]
1429 async fn test_time_travel_config_defaults() {
1430 let config = TimeTravelConfig::default();
1431 assert_eq!(config.max_time_window_days, 365);
1432 assert!(config.enable_temporal_indexing);
1433 assert_eq!(config.index_granularity_minutes, 60);
1434 }
1435
1436 #[tokio::test]
1437 async fn test_temporal_query_builder() {
1438 let query = TemporalQuery::new()
1439 .at_time(TimePoint::Timestamp(Utc::now()))
1440 .filter(TemporalFilter::default())
1441 .order_by(TemporalOrdering::TimeDescending)
1442 .limit(100);
1443
1444 assert!(query.time_point.is_some());
1445 assert!(query.limit.is_some());
1446 assert_eq!(query.limit.unwrap(), 100);
1447 }
1448
1449 #[tokio::test]
1450 async fn test_time_point_resolution() {
1451 let now = Utc::now();
1452 let relative = TimePoint::RelativeTime(ChronoDuration::hours(-1));
1453
1454 match relative {
1455 TimePoint::RelativeTime(duration) => {
1456 let resolved = now + duration;
1457 assert!(resolved < now);
1458 }
1459 _ => panic!("Expected RelativeTime"),
1460 }
1461 }
1462
1463 #[tokio::test]
1464 async fn test_temporal_filter() {
1465 let filter = TemporalFilter {
1466 event_types: Some(std::iter::once("TestEvent".to_string()).collect()),
1467 ..Default::default()
1468 };
1469
1470 assert!(filter.event_types.is_some());
1471 assert!(filter.event_types.as_ref().unwrap().contains("TestEvent"));
1472 }
1473}