Skip to main content

oxirs_stream/
time_travel.rs

1//! # Time-Travel Query System
2//!
3//! Advanced temporal query capabilities for OxiRS Stream, enabling querying
4//! data at any point in time, temporal analytics, and historical state reconstruction.
5
6use 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
18/// Type alias for custom filter functions to reduce complexity
19pub type CustomFilterFn = Box<dyn Fn(&StreamEvent) -> bool + Send + Sync>;
20
21/// Time-travel query configuration
22#[derive(Debug, Clone, Serialize, Deserialize)]
23pub struct TimeTravelConfig {
24    /// Maximum time window for time-travel queries
25    pub max_time_window_days: u32,
26    /// Enable temporal indexing for faster queries
27    pub enable_temporal_indexing: bool,
28    /// Temporal index granularity (minutes)
29    pub index_granularity_minutes: u32,
30    /// Maximum concurrent time-travel queries
31    pub max_concurrent_queries: usize,
32    /// Query timeout in seconds
33    pub query_timeout_seconds: u64,
34    /// Enable result caching
35    pub enable_result_caching: bool,
36    /// Cache TTL in minutes
37    pub cache_ttl_minutes: u32,
38    /// Maximum cache size in MB
39    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/// Time point specification for queries
58#[derive(Debug, Clone, Serialize, Deserialize)]
59pub enum TimePoint {
60    /// Specific timestamp
61    Timestamp(DateTime<Utc>),
62    /// Relative time from now
63    RelativeTime(ChronoDuration),
64    /// Event version number
65    Version(u64),
66    /// Event ID
67    EventId(Uuid),
68    /// Named snapshot
69    Snapshot(String),
70}
71
72/// Time range specification for queries
73#[derive(Debug, Clone, Serialize, Deserialize)]
74pub struct TimeRange {
75    pub start: TimePoint,
76    pub end: TimePoint,
77}
78
79/// Temporal query specification
80#[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    /// Create a new temporal query
99    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    /// Query at specific time point
112    pub fn at_time(mut self, time_point: TimePoint) -> Self {
113        self.time_point = Some(time_point);
114        self
115    }
116
117    /// Query within time range
118    pub fn in_range(mut self, time_range: TimeRange) -> Self {
119        self.time_range = Some(time_range);
120        self
121    }
122
123    /// Add filter
124    pub fn filter(mut self, filter: TemporalFilter) -> Self {
125        self.filter = filter;
126        self
127    }
128
129    /// Set projection
130    pub fn project(mut self, projection: TemporalProjection) -> Self {
131        self.projection = projection;
132        self
133    }
134
135    /// Set ordering
136    pub fn order_by(mut self, ordering: TemporalOrdering) -> Self {
137        self.ordering = ordering;
138        self
139    }
140
141    /// Set limit
142    pub fn limit(mut self, limit: usize) -> Self {
143        self.limit = Some(limit);
144        self
145    }
146}
147
148/// Temporal filter for events
149#[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(), // Cannot clone function pointers
181        }
182    }
183}
184
185/// Temporal projection specification
186#[derive(Debug, Clone, Default)]
187pub enum TemporalProjection {
188    /// Return full events
189    #[default]
190    FullEvents,
191    /// Return only metadata
192    MetadataOnly,
193    /// Return specific fields
194    Fields(Vec<String>),
195    /// Return aggregated data
196    Aggregation(AggregationType),
197}
198
199/// Aggregation type for temporal queries
200#[derive(Debug, Clone)]
201pub enum AggregationType {
202    Count,
203    CountBy(String),
204    Timeline(ChronoDuration),
205    Statistics,
206}
207
208/// Temporal ordering specification
209#[derive(Debug, Clone, Default)]
210pub enum TemporalOrdering {
211    /// Order by timestamp ascending
212    TimeAscending,
213    /// Order by timestamp descending
214    #[default]
215    TimeDescending,
216    /// Order by version ascending
217    VersionAscending,
218    /// Order by version descending
219    VersionDescending,
220    /// Order by custom field
221    Custom(String, bool), // field, ascending
222}
223
224/// Result of a temporal query
225#[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/// Metadata about temporal query results
236#[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/// Aggregated data from temporal queries
247#[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/// Point in timeline aggregation
256#[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/// Statistical data from temporal queries
264#[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/// Temporal index for efficient time-travel queries
275#[derive(Debug)]
276struct TemporalIndex {
277    /// Time-based index: timestamp -> event IDs
278    time_index: BTreeMap<DateTime<Utc>, Vec<Uuid>>,
279    /// Version-based index: version -> event metadata
280    version_index: BTreeMap<u64, EventIndexEntry>,
281    /// Aggregate-based index: aggregate_id -> time-ordered events
282    aggregate_index: HashMap<String, BTreeMap<DateTime<Utc>, Vec<Uuid>>>,
283    /// Type-based index: event_type -> time-ordered events
284    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        // Time index
315        self.time_index.entry(timestamp).or_default().push(event_id);
316
317        // Version index
318        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        // Aggregate index
330        self.aggregate_index
331            .entry(aggregate_id)
332            .or_default()
333            .entry(timestamp)
334            .or_default()
335            .push(event_id);
336
337        // Type index
338        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
384/// Time-travel query engine
385pub 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    /// Create a new time-travel engine
397    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    /// Start the time-travel engine
414    pub async fn start(&self) -> Result<()> {
415        info!("Starting time-travel engine");
416
417        // Build initial index if enabled
418        if self.config.enable_temporal_indexing {
419            self.build_temporal_index().await?;
420        }
421
422        // Start index maintenance task
423        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    /// Execute a temporal query
443    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        // Acquire semaphore for concurrency control
450        let _permit = self.query_semaphore.acquire().await?;
451
452        // Update metrics
453        {
454            let mut metrics = self.metrics.write().await;
455            metrics.queries_executed += 1;
456            metrics.active_queries += 1;
457        }
458
459        // Check cache first
460        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        // Update metrics
482        {
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            // Cache result if applicable
504            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    /// Execute query internally
518    async fn execute_query_internal(&self, query: TemporalQuery) -> Result<TemporalQueryResult> {
519        let query_id = query.query_id;
520
521        // Resolve time points to actual timestamps
522        let (start_time, end_time) = self.resolve_time_range(&query).await?;
523
524        // Find candidate events
525        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        // Load full events
534        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        // Apply ordering
544        self.apply_ordering(&mut events, &query.ordering);
545
546        // Apply limit
547        if let Some(limit) = query.limit {
548            events.truncate(limit);
549        }
550
551        // Generate metadata
552        let metadata = self.generate_result_metadata(&events, start_time, end_time);
553
554        // Generate aggregations if requested
555        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        // Apply projection
563        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(), // Will be set by caller
571            from_cache: false,
572        })
573    }
574
575    /// Query state at specific time point
576    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    /// Query changes between two time points
593    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    /// Query timeline aggregation
608    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    /// Build temporal index from existing events
626    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    /// Update index with new events
647    async fn update_index(
648        index: Arc<RwLock<TemporalIndex>>,
649        event_stream: Arc<dyn EventStream>,
650    ) -> Result<()> {
651        // This would typically track the last processed position
652        // For simplicity, we'll just rebuild periodically
653        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    /// Resolve time range from query specification
664    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                // Default to last 24 hours
682                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    /// Resolve a time point to an actual timestamp
690    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                // Find timestamp for this version
696                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                // Find timestamp for this event ID
705                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                // Named snapshots require an external snapshot store integration.
713                // The TimeTravelEngine does not embed a snapshot store; callers that need
714                // named-snapshot resolution should resolve the snapshot to a concrete
715                // Timestamp or Version and use that TimePoint variant instead.
716                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    /// Find events using temporal index
727    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        // Use most specific index available
736        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    /// Find events without using index (sequential scan).
750    ///
751    /// Used whenever `enable_temporal_indexing` is off (or as a correctness
752    /// baseline): scans the full event stream and applies the same time-range
753    /// (and, if present, single-aggregate) filtering that
754    /// [`Self::find_events_with_index`] gets "for free" from `TemporalIndex`,
755    /// so results are consistent between the indexed and non-indexed paths.
756    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    /// Resolve the stable UUID identity used to reference a stored event
801    /// throughout `TimeTravelEngine` (matches the derivation in
802    /// [`TemporalIndex::add_event`]): parsed from `metadata.event_id` when
803    /// that's a valid UUID, falling back to the event store's own
804    /// `StoredEvent::event_id` otherwise (unlike the index's random fallback,
805    /// this is deterministic, which is required for `load_event` to be able
806    /// to look events back up by the ID `find_events_without_index` returned).
807    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    /// Load a specific event by ID.
812    ///
813    /// The event store trait has no direct get-by-ID lookup, so this scans
814    /// the event stream (same source [`Self::find_events_without_index`] and
815    /// [`Self::build_temporal_index`] use) and returns the first event whose
816    /// resolved UUID identity (see [`Self::event_uuid`]) matches.
817    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    /// Check if event matches filter
830    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        // Apply custom filters
867        for custom_filter in &filter.custom_filters {
868            if !custom_filter(event) {
869                return false;
870            }
871        }
872
873        true
874    }
875
876    /// Apply ordering to events
877    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                // Custom field ordering would be implemented here
893                warn!("Custom ordering not implemented");
894            }
895        }
896    }
897
898    /// Apply projection to events
899    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                // Return events with only metadata (simplified data)
908                // For metadata-only projection, we keep the event but could filter data in a real implementation
909                events
910            }
911            TemporalProjection::Fields(_fields) => {
912                // Field projection would be implemented here
913                warn!("Field projection not implemented");
914                events
915            }
916            TemporalProjection::Aggregation(_) => {
917                // Aggregation results are handled separately
918                Vec::new()
919            }
920        }
921    }
922
923    /// Generate result metadata
924    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, // Would be tracked during execution
977            index_misses: 0,
978        }
979    }
980
981    /// Generate aggregations
982    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                    // Other fields would be handled here
1004                }
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    /// Generate timeline aggregation
1033    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    /// Calculate temporal statistics
1072    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        // Calculate peak throughput (events per second in busiest minute)
1086        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        // Calculate average event size
1102        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    /// Generate cache key for query
1132    fn generate_cache_key(&self, query: &TemporalQuery) -> String {
1133        // Simple cache key based on query structure
1134        format!("temporal_query_{:?}", query.query_id)
1135    }
1136
1137    /// Get time-travel metrics
1138    pub async fn get_metrics(&self) -> TimeTravelMetrics {
1139        self.metrics.read().await.clone()
1140    }
1141}
1142
1143/// Query cache for temporal queries
1144#[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        // Remove expired entries
1190        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        // Simple memory management (could be more sophisticated)
1197        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/// Time-travel engine metrics
1213#[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    /// Minimal in-memory `EventStream` used only to exercise
1233    /// `find_events_without_index` / `load_event` against real (if
1234    /// synthetic) stored events, since no production `EventStream`
1235    /// implementor exists to construct a `TimeTravelEngine` against.
1236    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    /// `EventStoreTrait` is required to construct a `TimeTravelEngine` but
1266    /// is untouched by the methods under test here, so every method is a
1267    /// harmless stub.
1268    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    /// Regression test for time_travel.rs:750 — `find_events_without_index`
1366    /// used to unconditionally `warn!` and return an empty `Vec`, regardless
1367    /// of what was actually in the event store. It must now genuinely scan
1368    /// the event stream and return events whose timestamp falls in range.
1369    #[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    /// Regression test for time_travel.rs:763 — `load_event` used to always
1393    /// return `Ok(None)` regardless of the requested ID, breaking
1394    /// `TimePoint::EventId` resolution. It must now actually locate and
1395    /// return the matching event from the event store.
1396    #[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        // A random, never-stored ID must still resolve to `None` (correct
1420        // "not found", not a fabricated hit).
1421        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}