1#[cfg(feature = "server")]
2use crate::application::services::webhook::WebhookRegistry;
3#[cfg(feature = "server")]
4use crate::infrastructure::observability::metrics::MetricsRegistry;
5#[cfg(feature = "server")]
6use crate::infrastructure::web::websocket::WebSocketManager;
7use crate::{
8 application::{
9 dto::QueryEventsRequest,
10 services::{
11 consumer::ConsumerRegistry,
12 exactly_once::{ExactlyOnceConfig, ExactlyOnceRegistry},
13 pipeline::PipelineManager,
14 projection::{EntitySnapshotProjection, EventCounterProjection, ProjectionManager},
15 replay::ReplayManager,
16 schema::{SchemaRegistry, SchemaRegistryConfig},
17 schema_evolution::SchemaEvolutionManager,
18 },
19 },
20 domain::entities::Event,
21 error::{AllSourceError, Result},
22 infrastructure::{
23 persistence::{
24 compaction::{CompactionConfig, CompactionManager},
25 index::{EventIndex, IndexEntry},
26 snapshot::{SnapshotConfig, SnapshotManager, SnapshotType},
27 storage::ParquetStorage,
28 tenant_loader::TenantLoader,
29 wal::{WALConfig, WriteAheadLog},
30 },
31 query::geospatial::GeoIndex,
32 },
33};
34use chrono::{DateTime, Utc};
35use dashmap::DashMap;
36use parking_lot::RwLock;
37use std::{path::PathBuf, sync::Arc};
38#[cfg(feature = "server")]
39use tokio::sync::mpsc;
40
41pub struct EventStore {
43 events: Arc<RwLock<Vec<Event>>>,
45
46 index: Arc<EventIndex>,
48
49 pub(crate) projections: Arc<RwLock<ProjectionManager>>,
51
52 storage: Option<Arc<RwLock<ParquetStorage>>>,
54
55 #[cfg(feature = "server")]
57 websocket_manager: Arc<WebSocketManager>,
58
59 snapshot_manager: Arc<SnapshotManager>,
61
62 wal: Option<Arc<WriteAheadLog>>,
64
65 compaction_manager: Option<Arc<CompactionManager>>,
67
68 schema_registry: Arc<SchemaRegistry>,
70
71 replay_manager: Arc<ReplayManager>,
73
74 pipeline_manager: Arc<PipelineManager>,
76
77 #[cfg(feature = "server")]
79 metrics: Arc<MetricsRegistry>,
80
81 total_ingested: Arc<RwLock<u64>>,
83
84 projection_state_cache: Arc<DashMap<String, serde_json::Value>>,
88
89 projection_status: Arc<DashMap<String, String>>,
92
93 #[cfg(feature = "server")]
95 webhook_registry: Arc<WebhookRegistry>,
96
97 #[cfg(feature = "server")]
99 webhook_tx: Arc<RwLock<Option<mpsc::UnboundedSender<WebhookDeliveryTask>>>>,
100
101 geo_index: Arc<GeoIndex>,
103
104 exactly_once: Arc<ExactlyOnceRegistry>,
106
107 schema_evolution: Arc<SchemaEvolutionManager>,
109
110 entity_versions: Arc<DashMap<String, u64>>,
113
114 consumer_registry: Arc<ConsumerRegistry>,
116
117 event_broadcast_tx: tokio::sync::broadcast::Sender<Arc<Event>>,
121
122 tenant_loader: Arc<TenantLoader>,
128
129 checkpoint_interval_secs: Option<u64>,
135
136 read_only: bool,
143}
144
145#[cfg(feature = "server")]
147#[derive(Debug, Clone)]
148pub struct WebhookDeliveryTask {
149 pub webhook: crate::application::services::webhook::WebhookSubscription,
150 pub event: Event,
151}
152
153fn select_window<T>(
166 items: &mut Vec<T>,
167 offset: usize,
168 limit: Option<usize>,
169 order: impl Fn(&T, &T) -> std::cmp::Ordering + Copy,
170) {
171 match limit.map(|limit| offset.saturating_add(limit)) {
172 Some(0) => {
173 items.clear();
174 return;
175 }
176 Some(window_end) if window_end < items.len() => {
177 items.select_nth_unstable_by(window_end - 1, order);
178 items.truncate(window_end);
179 }
180 _ => {}
183 }
184 items.sort_unstable_by(order);
185}
186
187impl EventStore {
188 pub fn new() -> Self {
190 Self::with_config(EventStoreConfig::default())
191 }
192
193 pub fn with_config(config: EventStoreConfig) -> Self {
195 let mut projections = ProjectionManager::new();
196
197 projections.register(Arc::new(EntitySnapshotProjection::new("entity_snapshots")));
199 projections.register(Arc::new(EventCounterProjection::new("event_counters")));
200
201 let storage = config
203 .storage_dir
204 .as_ref()
205 .and_then(|dir| match ParquetStorage::new(dir) {
206 Ok(storage) => {
207 tracing::info!("✅ Parquet persistence enabled at: {}", dir.display());
208 Some(Arc::new(RwLock::new(storage)))
209 }
210 Err(e) => {
211 tracing::error!("❌ Failed to initialize Parquet storage: {}", e);
212 None
213 }
214 });
215
216 let wal = config.wal_dir.as_ref().and_then(|dir| {
218 match WriteAheadLog::new(dir, config.wal_config.clone()) {
219 Ok(wal) => {
220 tracing::info!("✅ WAL enabled at: {}", dir.display());
221 Some(Arc::new(wal))
222 }
223 Err(e) => {
224 tracing::error!("❌ Failed to initialize WAL: {}", e);
225 None
226 }
227 }
228 });
229
230 let compaction_manager = config.storage_dir.as_ref().map(|dir| {
232 let manager = CompactionManager::new(dir, config.compaction_config.clone());
233 Arc::new(manager)
234 });
235
236 let schema_registry = Arc::new(SchemaRegistry::new(config.schema_registry_config.clone()));
238 tracing::info!("✅ Schema registry enabled");
239
240 let replay_manager = Arc::new(ReplayManager::new());
242 tracing::info!("✅ Replay manager enabled");
243
244 let pipeline_manager = Arc::new(PipelineManager::new());
246 tracing::info!("✅ Pipeline manager enabled");
247
248 #[cfg(feature = "server")]
250 let metrics = {
251 let m = MetricsRegistry::new();
252 tracing::info!("✅ Prometheus metrics registry initialized");
253 m
254 };
255
256 let projection_state_cache = Arc::new(DashMap::new());
258 tracing::info!("✅ Projection state cache initialized");
259
260 #[cfg(feature = "server")]
262 let webhook_registry = {
263 let w = Arc::new(WebhookRegistry::new());
264 tracing::info!("✅ Webhook registry initialized");
265 w
266 };
267
268 let (event_broadcast_tx, _) = tokio::sync::broadcast::channel(1024);
271
272 let store = Self {
273 events: Arc::new(RwLock::new(Vec::new())),
274 index: Arc::new(EventIndex::new()),
275 projections: Arc::new(RwLock::new(projections)),
276 storage,
277 #[cfg(feature = "server")]
278 websocket_manager: Arc::new(WebSocketManager::new()),
279 snapshot_manager: Arc::new(SnapshotManager::new(config.snapshot_config)),
280 wal,
281 compaction_manager,
282 schema_registry,
283 replay_manager,
284 pipeline_manager,
285 #[cfg(feature = "server")]
286 metrics,
287 total_ingested: Arc::new(RwLock::new(0)),
288 projection_state_cache,
289 projection_status: Arc::new(DashMap::new()),
290 #[cfg(feature = "server")]
291 webhook_registry,
292 #[cfg(feature = "server")]
293 webhook_tx: Arc::new(RwLock::new(None)),
294 geo_index: Arc::new(GeoIndex::new()),
295 exactly_once: Arc::new(ExactlyOnceRegistry::new(ExactlyOnceConfig::default())),
296 schema_evolution: Arc::new(SchemaEvolutionManager::new()),
297 entity_versions: Arc::new(DashMap::new()),
298 consumer_registry: Arc::new(ConsumerRegistry::new()),
299 event_broadcast_tx,
300 tenant_loader: {
301 let loader = TenantLoader::new();
302 if let Some(budget) = config.cache_byte_budget {
303 loader.set_byte_budget(budget);
304 tracing::info!(
305 "✅ Cache byte budget set to {} bytes ({:.2} GiB) — LRU eviction enabled",
306 budget,
307 budget as f64 / (1024.0 * 1024.0 * 1024.0)
308 );
309 } else {
310 tracing::info!(
311 "✅ Cache budget unset — every loaded tenant stays resident \
312 (set ALLSOURCE_CACHE_BYTES to enable eviction)"
313 );
314 }
315 Arc::new(loader)
316 },
317 checkpoint_interval_secs: config.checkpoint_interval_secs,
318 read_only: config.read_only,
319 };
320
321 if config.read_only {
322 tracing::info!(
323 "📖 EventStore opened READ-ONLY (replica): WAL will be replayed for reads but \
324 not truncated; writes are rejected"
325 );
326 }
327
328 if let Some(ref wal) = store.wal {
350 match wal.recover() {
351 Ok(recovered_events) if !recovered_events.is_empty() => {
352 let mut wal_new = 0usize;
353 for event in recovered_events {
354 let offset = store.events.read().len();
355 if let Err(e) = store.index.index_event(
356 event.id,
357 event.entity_id_str(),
358 event.event_type_str(),
359 event.timestamp,
360 offset,
361 ) {
362 tracing::error!("Failed to re-index WAL event {}: {}", event.id, e);
363 }
364
365 if let Err(e) = store.projections.read().process_event(&event) {
366 tracing::error!("Failed to re-process WAL event {}: {}", event.id, e);
367 }
368
369 *store
370 .entity_versions
371 .entry(event.entity_id_str().to_string())
372 .or_insert(0) += 1;
373
374 store.events.write().push(event);
375 wal_new += 1;
376 }
377
378 #[cfg(feature = "server")]
383 store.metrics.wal_replay_events_total.set(wal_new as i64);
384
385 if wal_new > 0 {
386 let total = store.events.read().len();
387 *store.total_ingested.write() = total as u64;
395 tracing::info!(
396 "✅ Recovered {} events from WAL (Parquet data stays cold until \
397 first per-tenant query)",
398 wal_new
399 );
400
401 if let Some(ref storage) = store.storage
416 && !config.read_only
417 {
418 tracing::info!(
419 "📸 Checkpointing {} WAL events to Parquet storage...",
420 wal_new
421 );
422 let parquet = storage.read();
423 let events = store.events.read();
424 let mut buffered = 0usize;
425 for event in events.iter().skip(events.len() - wal_new) {
426 if let Err(e) = parquet.append_event(event.clone()) {
427 tracing::error!(
428 "Failed to buffer WAL event for Parquet: {}",
429 e
430 );
431 } else {
432 buffered += 1;
433 }
434 }
435 drop(events);
436 drop(parquet);
437
438 if buffered > 0 {
439 if let Err(e) = store.flush_storage() {
440 tracing::error!("Failed to checkpoint to Parquet: {}", e);
441 } else if let Err(e) = wal.truncate() {
442 tracing::error!(
443 "Failed to truncate WAL after checkpoint: {}",
444 e
445 );
446 } else {
447 tracing::info!(
448 "✅ WAL checkpointed and truncated ({} events)",
449 buffered
450 );
451 }
452 }
453 }
454 }
455 }
456 Ok(_) => {
457 tracing::debug!("No events to recover from WAL");
458 #[cfg(feature = "server")]
459 store.metrics.wal_replay_events_total.set(0);
460 }
461 Err(e) => {
462 tracing::error!("❌ WAL recovery failed: {}", e);
463 }
464 }
465 } else if store.storage.is_some() {
466 tracing::info!(
467 "📂 Boot complete (lazy-load mode): Parquet data stays on disk until first \
468 per-tenant query"
469 );
470 }
471
472 store
473 }
474
475 pub fn is_read_only(&self) -> bool {
477 self.read_only
478 }
479
480 fn ensure_writable(&self) -> Result<()> {
484 if self.read_only {
485 return Err(crate::error::AllSourceError::ReadOnly(
486 "this AllSource instance is a read-only replica — the data directory is owned by \
487 another running process. Stop the other process, or run a single shared writer \
488 (e.g. Prime in --mode http) and point clients at it."
489 .to_string(),
490 ));
491 }
492 Ok(())
493 }
494
495 #[cfg_attr(feature = "hotpath", hotpath::measure)]
503 pub fn ingest_with_expected_version(
504 &self,
505 event: &Event,
506 expected_version: Option<u64>,
507 ) -> Result<u64> {
508 self.ensure_writable()?;
510
511 self.validate_event(event)?;
513
514 let entity_id = event.entity_id_str().to_string();
515
516 let new_version = {
519 let mut version_entry = self.entity_versions.entry(entity_id.clone()).or_insert(0);
520 let current = *version_entry;
521
522 if let Some(expected) = expected_version
523 && current != expected
524 {
525 return Err(crate::error::AllSourceError::VersionConflict { expected, current });
526 }
527
528 if let Some(ref wal) = self.wal {
530 wal.append(event.clone())?;
531 }
532
533 *version_entry += 1;
534 *version_entry
535 };
536
537 self.ingest_post_wal(event)?;
540
541 Ok(new_version)
542 }
543
544 #[cfg_attr(feature = "hotpath", hotpath::measure)]
547 fn ingest_post_wal(&self, event: &Event) -> Result<()> {
548 #[cfg(feature = "server")]
549 let timer = self.metrics.ingestion_duration_seconds.start_timer();
550
551 let mut events = self.events.write();
552 let offset = events.len();
553
554 self.index.index_event(
556 event.id,
557 event.entity_id_str(),
558 event.event_type_str(),
559 event.timestamp,
560 offset,
561 )?;
562
563 let projections = self.projections.read();
565 projections.process_event(event)?;
566 drop(projections);
567
568 let pipeline_results = self.pipeline_manager.process_event(event);
570 if !pipeline_results.is_empty() {
571 tracing::debug!(
572 "Event {} processed by {} pipeline(s)",
573 event.id,
574 pipeline_results.len()
575 );
576 for (pipeline_id, result) in pipeline_results {
577 tracing::trace!("Pipeline {} result: {:?}", pipeline_id, result);
578 }
579 }
580
581 if let Some(ref storage) = self.storage {
583 let storage = storage.read();
584 storage.append_event(event.clone())?;
585 }
586
587 events.push(event.clone());
589 let total_events = events.len();
590 drop(events);
591
592 let event_arc = Arc::new(event.clone());
594 let _ = self.event_broadcast_tx.send(Arc::clone(&event_arc));
595 #[cfg(feature = "server")]
596 self.websocket_manager.broadcast_event(event_arc);
597
598 #[cfg(feature = "server")]
600 self.dispatch_webhooks(event);
601
602 self.geo_index.index_event(event);
604
605 self.schema_evolution
607 .analyze_event(event.event_type_str(), &event.payload);
608
609 self.check_auto_snapshot(event.entity_id_str(), event);
611
612 #[cfg(feature = "server")]
614 {
615 self.metrics.events_ingested_total.inc();
616 self.metrics
617 .events_ingested_by_type
618 .with_label_values(&[event.event_type_str()])
619 .inc();
620 self.metrics.storage_events_total.set(total_events as i64);
621 }
622
623 let mut total = self.total_ingested.write();
625 *total += 1;
626
627 #[cfg(feature = "server")]
628 timer.observe_duration();
629
630 tracing::debug!("Event ingested: {} (offset: {})", event.id, offset);
631
632 Ok(())
633 }
634
635 #[cfg_attr(feature = "hotpath", hotpath::measure)]
637 pub fn ingest(&self, event: &Event) -> Result<()> {
638 #[cfg(feature = "server")]
640 let timer = self.metrics.ingestion_duration_seconds.start_timer();
641
642 if let Err(e) = self.ensure_writable() {
644 #[cfg(feature = "server")]
645 {
646 self.metrics.ingestion_errors_total.inc();
647 timer.observe_duration();
648 }
649 return Err(e);
650 }
651
652 let validation_result = self.validate_event(event);
654 if let Err(e) = validation_result {
655 #[cfg(feature = "server")]
656 {
657 self.metrics.ingestion_errors_total.inc();
658 timer.observe_duration();
659 }
660 return Err(e);
661 }
662
663 if let Some(ref wal) = self.wal
666 && let Err(e) = wal.append(event.clone())
667 {
668 #[cfg(feature = "server")]
669 {
670 self.metrics.ingestion_errors_total.inc();
671 timer.observe_duration();
672 }
673 return Err(e);
674 }
675
676 *self
678 .entity_versions
679 .entry(event.entity_id_str().to_string())
680 .or_insert(0) += 1;
681
682 let mut events = self.events.write();
683 let offset = events.len();
684
685 self.index.index_event(
687 event.id,
688 event.entity_id_str(),
689 event.event_type_str(),
690 event.timestamp,
691 offset,
692 )?;
693
694 let projections = self.projections.read();
696 projections.process_event(event)?;
697 drop(projections); let pipeline_results = self.pipeline_manager.process_event(event);
702 if !pipeline_results.is_empty() {
703 tracing::debug!(
704 "Event {} processed by {} pipeline(s)",
705 event.id,
706 pipeline_results.len()
707 );
708 for (pipeline_id, result) in pipeline_results {
711 tracing::trace!("Pipeline {} result: {:?}", pipeline_id, result);
712 }
713 }
714
715 if let Some(ref storage) = self.storage {
717 let storage = storage.read();
718 storage.append_event(event.clone())?;
719 }
720
721 events.push(event.clone());
723 let total_events = events.len();
724 drop(events); let event_arc = Arc::new(event.clone());
728 let _ = self.event_broadcast_tx.send(Arc::clone(&event_arc));
729 #[cfg(feature = "server")]
730 self.websocket_manager.broadcast_event(event_arc);
731
732 #[cfg(feature = "server")]
734 self.dispatch_webhooks(event);
735
736 self.geo_index.index_event(event);
738
739 self.schema_evolution
741 .analyze_event(event.event_type_str(), &event.payload);
742
743 self.check_auto_snapshot(event.entity_id_str(), event);
745
746 #[cfg(feature = "server")]
748 {
749 self.metrics.events_ingested_total.inc();
750 self.metrics
751 .events_ingested_by_type
752 .with_label_values(&[event.event_type_str()])
753 .inc();
754 self.metrics.storage_events_total.set(total_events as i64);
755 }
756
757 let mut total = self.total_ingested.write();
759 *total += 1;
760
761 #[cfg(feature = "server")]
762 timer.observe_duration();
763
764 tracing::debug!("Event ingested: {} (offset: {})", event.id, offset);
765
766 Ok(())
767 }
768
769 #[cfg_attr(feature = "hotpath", hotpath::measure)]
776 pub fn ingest_batch(&self, batch: Vec<Event>) -> Result<()> {
777 if batch.is_empty() {
778 return Ok(());
779 }
780
781 self.ensure_writable()?;
783
784 for event in &batch {
786 self.validate_event(event)?;
787 }
788
789 if let Some(ref wal) = self.wal {
791 for event in &batch {
792 wal.append(event.clone())?;
793 }
794 }
795
796 let mut events = self.events.write();
798 let projections = self.projections.read();
799
800 for event in batch {
801 let offset = events.len();
802
803 self.index.index_event(
804 event.id,
805 event.entity_id_str(),
806 event.event_type_str(),
807 event.timestamp,
808 offset,
809 )?;
810
811 projections.process_event(&event)?;
812 self.pipeline_manager.process_event(&event);
813
814 if let Some(ref storage) = self.storage {
815 let storage = storage.read();
816 storage.append_event(event.clone())?;
817 }
818
819 self.geo_index.index_event(&event);
820 self.schema_evolution
821 .analyze_event(event.event_type_str(), &event.payload);
822
823 *self
825 .entity_versions
826 .entry(event.entity_id_str().to_string())
827 .or_insert(0) += 1;
828
829 let _ = self.event_broadcast_tx.send(Arc::new(event.clone()));
831
832 events.push(event);
833 }
834
835 let total_events = events.len();
836 drop(projections);
837 drop(events);
838
839 let mut total = self.total_ingested.write();
840 *total += total_events as u64;
841
842 Ok(())
843 }
844
845 #[cfg_attr(feature = "hotpath", hotpath::measure)]
852 pub fn ingest_replicated(&self, event: &Event) -> Result<()> {
853 #[cfg(feature = "server")]
854 let timer = self.metrics.ingestion_duration_seconds.start_timer();
855
856 let mut events = self.events.write();
857 let offset = events.len();
858
859 self.index.index_event(
861 event.id,
862 event.entity_id_str(),
863 event.event_type_str(),
864 event.timestamp,
865 offset,
866 )?;
867
868 let projections = self.projections.read();
870 projections.process_event(event)?;
871 drop(projections);
872
873 let pipeline_results = self.pipeline_manager.process_event(event);
875 if !pipeline_results.is_empty() {
876 tracing::debug!(
877 "Replicated event {} processed by {} pipeline(s)",
878 event.id,
879 pipeline_results.len()
880 );
881 }
882
883 *self
885 .entity_versions
886 .entry(event.entity_id_str().to_string())
887 .or_insert(0) += 1;
888
889 events.push(event.clone());
891 let total_events = events.len();
892 drop(events);
893
894 let event_arc = Arc::new(event.clone());
896 let _ = self.event_broadcast_tx.send(Arc::clone(&event_arc));
897 #[cfg(feature = "server")]
898 self.websocket_manager.broadcast_event(event_arc);
899
900 #[cfg(feature = "server")]
902 {
903 self.metrics.events_ingested_total.inc();
904 self.metrics
905 .events_ingested_by_type
906 .with_label_values(&[event.event_type_str()])
907 .inc();
908 self.metrics.storage_events_total.set(total_events as i64);
909 }
910
911 let mut total = self.total_ingested.write();
912 *total += 1;
913
914 #[cfg(feature = "server")]
915 timer.observe_duration();
916
917 tracing::debug!(
918 "Replicated event ingested: {} (offset: {})",
919 event.id,
920 offset
921 );
922
923 Ok(())
924 }
925
926 #[cfg_attr(feature = "hotpath", hotpath::measure)]
929 pub fn get_entity_version(&self, entity_id: &str) -> u64 {
930 self.entity_versions.get(entity_id).map_or(0, |v| *v)
931 }
932
933 pub fn consumer_registry(&self) -> &ConsumerRegistry {
935 &self.consumer_registry
936 }
937
938 pub fn subscribe_events(&self) -> tokio::sync::broadcast::Receiver<Arc<Event>> {
944 self.event_broadcast_tx.subscribe()
945 }
946
947 pub fn set_consumer_registry(&mut self, registry: Arc<ConsumerRegistry>) {
952 self.consumer_registry = registry;
953 }
954
955 pub fn total_events(&self) -> usize {
957 self.events.read().len()
958 }
959
960 pub fn events_after_offset(
963 &self,
964 offset: u64,
965 filters: &[String],
966 limit: usize,
967 ) -> Vec<(u64, Event)> {
968 let events = self.events.read();
969 let start = offset as usize;
970 if start >= events.len() {
971 return vec![];
972 }
973
974 events[start..]
975 .iter()
976 .enumerate()
977 .filter(|(_, event)| ConsumerRegistry::matches_filters(event.event_type_str(), filters))
978 .take(limit)
979 .map(|(i, event)| ((start + i + 1) as u64, event.clone()))
980 .collect()
981 }
982
983 #[cfg(feature = "server")]
985 pub fn websocket_manager(&self) -> Arc<WebSocketManager> {
986 Arc::clone(&self.websocket_manager)
987 }
988
989 pub fn snapshot_manager(&self) -> Arc<SnapshotManager> {
991 Arc::clone(&self.snapshot_manager)
992 }
993
994 pub fn compaction_manager(&self) -> Option<Arc<CompactionManager>> {
996 self.compaction_manager.as_ref().map(Arc::clone)
997 }
998
999 pub fn schema_registry(&self) -> Arc<SchemaRegistry> {
1001 Arc::clone(&self.schema_registry)
1002 }
1003
1004 pub fn replay_manager(&self) -> Arc<ReplayManager> {
1006 Arc::clone(&self.replay_manager)
1007 }
1008
1009 pub fn pipeline_manager(&self) -> Arc<PipelineManager> {
1011 Arc::clone(&self.pipeline_manager)
1012 }
1013
1014 #[cfg(feature = "server")]
1016 pub fn metrics(&self) -> Arc<MetricsRegistry> {
1017 Arc::clone(&self.metrics)
1018 }
1019
1020 pub fn projection_manager(&self) -> parking_lot::RwLockReadGuard<'_, ProjectionManager> {
1022 self.projections.read()
1023 }
1024
1025 pub fn register_projection(
1034 &self,
1035 projection: Arc<dyn crate::application::services::projection::Projection>,
1036 ) {
1037 let mut pm = self.projections.write();
1038 pm.register(projection);
1039 }
1040
1041 pub fn register_projection_with_backfill(
1053 &self,
1054 projection: &Arc<dyn crate::application::services::projection::Projection>,
1055 ) -> Result<()> {
1056 {
1058 let mut pm = self.projections.write();
1059 pm.register(Arc::clone(projection));
1060 }
1061
1062 let events = self.events.read();
1064 let mut ordered: Vec<&Event> = events.iter().collect();
1065 ordered.sort_by(|a, b| {
1066 a.timestamp
1067 .cmp(&b.timestamp)
1068 .then_with(|| a.version.cmp(&b.version))
1069 });
1070 for event in ordered {
1071 projection.process(event)?;
1072 }
1073
1074 Ok(())
1075 }
1076
1077 pub fn hydrate_all_from_storage(&self) -> Result<usize> {
1095 let Some(storage) = self.storage.as_ref().map(Arc::clone) else {
1096 return Ok(0);
1097 };
1098
1099 let events = storage.read().load_all_events()?;
1100 let read_count = events.len();
1101 let before = self.events.read().len();
1102 for event in events {
1103 self.append_loaded_event(event);
1104 }
1105 let applied = self.events.read().len() - before;
1106
1107 *self.total_ingested.write() = self.events.read().len() as u64;
1109
1110 tracing::info!(
1111 read = read_count,
1112 applied = applied,
1113 "🔄 hydrate_all_from_storage: in-memory pile reconstructed from Parquet"
1114 );
1115 Ok(applied)
1116 }
1117
1118 pub fn projection_state_cache(&self) -> Arc<DashMap<String, serde_json::Value>> {
1121 Arc::clone(&self.projection_state_cache)
1122 }
1123
1124 pub fn projection_status(&self) -> Arc<DashMap<String, String>> {
1126 Arc::clone(&self.projection_status)
1127 }
1128
1129 pub fn geo_index(&self) -> Arc<GeoIndex> {
1132 self.geo_index.clone()
1133 }
1134
1135 pub fn exactly_once(&self) -> Arc<ExactlyOnceRegistry> {
1137 self.exactly_once.clone()
1138 }
1139
1140 pub fn schema_evolution(&self) -> Arc<SchemaEvolutionManager> {
1142 self.schema_evolution.clone()
1143 }
1144
1145 pub fn snapshot_events(&self) -> Vec<Event> {
1151 self.events.read().clone()
1152 }
1153
1154 pub fn compact_entity_tokens(
1176 &self,
1177 entity_id: &str,
1178 token_event_type: &str,
1179 merged_event: Event,
1180 ) -> Result<bool> {
1181 self.ensure_writable()?;
1183
1184 {
1186 let events = self.events.read();
1187 let has_tokens = events
1188 .iter()
1189 .any(|e| e.entity_id_str() == entity_id && e.event_type_str() == token_event_type);
1190 if !has_tokens {
1191 return Ok(false);
1192 }
1193 }
1194
1195 let projections = self.projections.read();
1197 projections.process_event(&merged_event)?;
1198 drop(projections);
1199
1200 let mut events = self.events.write();
1202
1203 events.retain(|e| {
1204 !(e.entity_id_str() == entity_id && e.event_type_str() == token_event_type)
1205 });
1206
1207 events.push(merged_event.clone());
1208
1209 if let Some(ref wal) = self.wal {
1213 wal.append(merged_event)?;
1214 }
1215
1216 self.index.clear();
1221 for (offset, event) in events.iter().enumerate() {
1222 if let Err(e) = self.index.index_event(
1223 event.id,
1224 event.entity_id_str(),
1225 event.event_type_str(),
1226 event.timestamp,
1227 offset,
1228 ) {
1229 tracing::warn!(
1230 event_id = %event.id,
1231 offset,
1232 "Failed to re-index event during compaction: {e}"
1233 );
1234 }
1235 }
1236
1237 Ok(true)
1238 }
1239
1240 #[cfg(feature = "server")]
1241 pub fn webhook_registry(&self) -> Arc<WebhookRegistry> {
1242 Arc::clone(&self.webhook_registry)
1243 }
1244
1245 #[cfg(feature = "server")]
1248 pub fn set_webhook_tx(&self, tx: mpsc::UnboundedSender<WebhookDeliveryTask>) {
1249 *self.webhook_tx.write() = Some(tx);
1250 tracing::info!("Webhook delivery channel connected");
1251 }
1252
1253 #[cfg(feature = "server")]
1255 fn dispatch_webhooks(&self, event: &Event) {
1256 let matching = self.webhook_registry.find_matching(event);
1257 if matching.is_empty() {
1258 return;
1259 }
1260
1261 let tx_guard = self.webhook_tx.read();
1262 if let Some(ref tx) = *tx_guard {
1263 for webhook in matching {
1264 let task = WebhookDeliveryTask {
1265 webhook,
1266 event: event.clone(),
1267 };
1268 if let Err(e) = tx.send(task) {
1269 tracing::warn!("Failed to queue webhook delivery: {}", e);
1270 }
1271 }
1272 }
1273 }
1274
1275 pub fn flush_storage(&self) -> Result<()> {
1277 if let Some(ref storage) = self.storage {
1278 let storage = storage.read();
1279 storage.flush()?;
1280 tracing::info!("✅ Flushed events to persistent storage");
1281 }
1282 Ok(())
1283 }
1284
1285 pub fn checkpoint(&self) -> Result<()> {
1306 let Some(ref wal) = self.wal else {
1307 #[cfg(feature = "server")]
1310 self.refresh_storage_metrics();
1311 return Ok(());
1312 };
1313
1314 self.flush_storage()?;
1319 wal.truncate()?;
1320 tracing::debug!("✅ Checkpoint complete: Parquet flushed, WAL truncated");
1321
1322 #[cfg(feature = "server")]
1325 self.refresh_storage_metrics();
1326
1327 Ok(())
1328 }
1329
1330 #[cfg(feature = "server")]
1335 pub fn refresh_storage_metrics_now(&self) {
1336 self.refresh_storage_metrics();
1337 }
1338
1339 #[cfg(feature = "server")]
1356 fn refresh_storage_metrics(&self) {
1357 let Some(ref storage) = self.storage else {
1358 return;
1359 };
1360
1361 let parquet_stats = match storage.read().stats() {
1362 Ok(stats) => stats,
1363 Err(e) => {
1364 tracing::warn!("storage-size metric refresh: failed to stat Parquet: {e}");
1365 return;
1366 }
1367 };
1368
1369 let (wal_bytes, wal_segments) = match self.wal.as_ref() {
1370 Some(wal) => match wal.on_disk_stats() {
1371 Ok(stats) => stats,
1372 Err(e) => {
1373 tracing::warn!("storage-size metric refresh: failed to stat WAL: {e}");
1374 (0, 0)
1375 }
1376 },
1377 None => (0, 0),
1378 };
1379
1380 let total_bytes = parquet_stats.total_size_bytes + wal_bytes;
1381
1382 self.metrics
1383 .storage_size_bytes
1384 .set(total_bytes.min(i64::MAX as u64) as i64);
1385 self.metrics
1386 .parquet_files_total
1387 .set(parquet_stats.total_files as i64);
1388 self.metrics.wal_segments_total.set(wal_segments as i64);
1389
1390 tracing::debug!(
1391 "storage-size metrics refreshed: {} bytes total ({} Parquet files, {} WAL segments)",
1392 total_bytes,
1393 parquet_stats.total_files,
1394 wal_segments
1395 );
1396 }
1397
1398 pub fn checkpoint_interval(&self) -> Option<std::time::Duration> {
1400 self.checkpoint_interval_secs
1401 .map(std::time::Duration::from_secs)
1402 }
1403
1404 pub fn ensure_tenant_loaded(&self, tenant_id: &str) -> Result<()> {
1428 if self.tenant_loader.is_loaded(tenant_id) {
1430 return Ok(());
1431 }
1432
1433 let Some(storage) = self.storage.as_ref().map(Arc::clone) else {
1434 self.tenant_loader.mark_loaded(tenant_id);
1437 return Ok(());
1438 };
1439
1440 let lock = self.tenant_loader.lock_for(tenant_id);
1443 let timeout = self.tenant_loader.load_timeout();
1444 let _guard = lock.try_lock_for(timeout).ok_or_else(|| {
1445 AllSourceError::StorageError(format!(
1446 "ensure_tenant_loaded timed out after {timeout:?} waiting for in-flight load of \
1447 tenant {tenant_id:?}"
1448 ))
1449 })?;
1450
1451 if self.tenant_loader.is_loaded(tenant_id) {
1454 return Ok(());
1455 }
1456
1457 let started = std::time::Instant::now();
1458 let events = storage.read().load_events_for_tenant(tenant_id)?;
1459 let read_count = events.len();
1460
1461 let before = self.events.read().len();
1462 for event in events {
1463 self.append_loaded_event(event);
1464 }
1465 let applied = self.events.read().len() - before;
1466
1467 *self.total_ingested.write() += applied as u64;
1471 self.tenant_loader.mark_loaded(tenant_id);
1472
1473 tracing::info!(
1474 tenant_id = tenant_id,
1475 read = read_count,
1476 applied = applied,
1477 elapsed_ms = started.elapsed().as_millis() as u64,
1478 "ensure_tenant_loaded: tenant hydrated"
1479 );
1480
1481 self.enforce_cache_budget(tenant_id);
1489
1490 #[cfg(feature = "server")]
1493 self.metrics
1494 .cache_bytes
1495 .set(self.tenant_loader.total_bytes() as i64);
1496
1497 Ok(())
1498 }
1499
1500 fn enforce_cache_budget(&self, recently_touched: &str) {
1511 if !self.tenant_loader.over_budget() {
1512 return;
1513 }
1514 loop {
1515 let Some(victim) = self.tenant_loader.pick_lru_excluding(recently_touched) else {
1516 tracing::warn!(
1517 cache_bytes = self.tenant_loader.total_bytes(),
1518 budget = self.tenant_loader.byte_budget(),
1519 recently_touched = recently_touched,
1520 "cache over budget but no other tenant available to evict — \
1521 a single tenant exceeds the budget; consider raising it"
1522 );
1523 return;
1524 };
1525 self.evict_tenant(&victim);
1526 if !self.tenant_loader.over_budget() {
1527 return;
1528 }
1529 }
1530 }
1531
1532 pub fn is_tenant_loaded(&self, tenant_id: &str) -> bool {
1535 self.tenant_loader.is_loaded(tenant_id)
1536 }
1537
1538 pub fn evict_tenant(&self, tenant_id: &str) {
1566 let mut events = self.events.write();
1567 let before = events.len();
1568 let evicted_bytes = self.tenant_loader.bytes_for(tenant_id);
1569
1570 events.retain(|e| e.tenant_id_str() != tenant_id);
1571 let after = events.len();
1572 let dropped = before - after;
1573
1574 if dropped == 0 {
1575 drop(events);
1579 self.tenant_loader.mark_unloaded(tenant_id);
1580 return;
1581 }
1582
1583 self.index.clear();
1587 self.entity_versions.clear();
1588 for (offset, event) in events.iter().enumerate() {
1589 if let Err(e) = self.index.index_event(
1590 event.id,
1591 event.entity_id_str(),
1592 event.event_type_str(),
1593 event.timestamp,
1594 offset,
1595 ) {
1596 tracing::error!(
1597 "Failed to re-index event during eviction of {}: {}",
1598 tenant_id,
1599 e
1600 );
1601 }
1602 *self
1603 .entity_versions
1604 .entry(event.entity_id_str().to_string())
1605 .or_insert(0) += 1;
1606 }
1607 drop(events);
1608
1609 self.tenant_loader.mark_unloaded(tenant_id);
1610
1611 let mut t = self.total_ingested.write();
1614 *t = t.saturating_sub(dropped as u64);
1615 drop(t);
1616
1617 #[cfg(feature = "server")]
1620 {
1621 self.metrics.cache_evictions_total.inc();
1622 self.metrics
1623 .cache_bytes
1624 .set(self.tenant_loader.total_bytes() as i64);
1625 }
1626
1627 tracing::info!(
1628 tenant_id = tenant_id,
1629 events_dropped = dropped,
1630 bytes_freed = evicted_bytes,
1631 "evicted tenant from memory cache"
1632 );
1633 }
1634
1635 pub fn tenant_resident_bytes(&self, tenant_id: &str) -> u64 {
1639 self.tenant_loader.bytes_for(tenant_id)
1640 }
1641
1642 pub fn cache_resident_bytes(&self) -> u64 {
1645 self.tenant_loader.total_bytes()
1646 }
1647
1648 fn append_loaded_event(&self, event: Event) {
1670 if self.index.get_by_id(&event.id).is_some() {
1671 return;
1672 }
1673
1674 let event_bytes = event.estimated_size_bytes();
1675 let tenant = event.tenant_id_str().to_string();
1676
1677 let mut events = self.events.write();
1678 let offset = events.len();
1679
1680 if let Err(e) = self.index.index_event(
1681 event.id,
1682 event.entity_id_str(),
1683 event.event_type_str(),
1684 event.timestamp,
1685 offset,
1686 ) {
1687 tracing::error!("Failed to index loaded event {}: {}", event.id, e);
1688 }
1689
1690 if let Err(e) = self.projections.read().process_event(&event) {
1691 tracing::error!("Failed to project loaded event {}: {}", event.id, e);
1692 }
1693
1694 *self
1695 .entity_versions
1696 .entry(event.entity_id_str().to_string())
1697 .or_insert(0) += 1;
1698
1699 events.push(event);
1700 self.tenant_loader.add_bytes(&tenant, event_bytes);
1704 }
1705
1706 pub fn create_snapshot(&self, entity_id: &str) -> Result<()> {
1708 let events = self.query(&QueryEventsRequest {
1710 entity_id: Some(entity_id.to_string()),
1711 event_type: None,
1712 tenant_id: None,
1713 as_of: None,
1714 since: None,
1715 until: None,
1716 limit: None,
1717 event_type_prefix: None,
1718 exclude_event_type_prefix: None,
1719 payload_filter: None,
1720 })?;
1721
1722 if events.is_empty() {
1723 return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
1724 }
1725
1726 let mut state = serde_json::json!({});
1728 for event in &events {
1729 if let serde_json::Value::Object(ref mut state_map) = state
1730 && let serde_json::Value::Object(ref payload_map) = event.payload
1731 {
1732 for (key, value) in payload_map {
1733 state_map.insert(key.clone(), value.clone());
1734 }
1735 }
1736 }
1737
1738 let last_event = events.last().unwrap();
1739 self.snapshot_manager.create_snapshot(
1740 entity_id,
1741 state,
1742 last_event.timestamp,
1743 events.len(),
1744 SnapshotType::Manual,
1745 )?;
1746
1747 Ok(())
1748 }
1749
1750 fn check_auto_snapshot(&self, entity_id: &str, event: &Event) {
1752 let entity_event_count = self
1754 .index
1755 .get_by_entity(entity_id)
1756 .map_or(0, |entries| entries.len());
1757
1758 if self.snapshot_manager.should_create_snapshot(
1759 entity_id,
1760 entity_event_count,
1761 event.timestamp,
1762 ) {
1763 if let Err(e) = self.create_snapshot(entity_id) {
1765 tracing::warn!(
1766 "Failed to create automatic snapshot for {}: {}",
1767 entity_id,
1768 e
1769 );
1770 }
1771 }
1772 }
1773
1774 fn validate_event(&self, event: &Event) -> Result<()> {
1776 if event.entity_id_str().is_empty() {
1779 return Err(AllSourceError::ValidationError(
1780 "entity_id cannot be empty".to_string(),
1781 ));
1782 }
1783
1784 if event.event_type_str().is_empty() {
1785 return Err(AllSourceError::ValidationError(
1786 "event_type cannot be empty".to_string(),
1787 ));
1788 }
1789
1790 if event.event_type().is_system() {
1793 return Err(AllSourceError::ValidationError(
1794 "Event types starting with '_system.' are reserved for internal use".to_string(),
1795 ));
1796 }
1797
1798 Ok(())
1799 }
1800
1801 pub fn reset_projection(&self, name: &str) -> Result<usize> {
1803 let projection_manager = self.projections.read();
1804 let projection = projection_manager.get_projection(name).ok_or_else(|| {
1805 AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1806 })?;
1807
1808 projection.clear();
1810
1811 let prefix = format!("{name}:");
1813 let keys_to_remove: Vec<String> = self
1814 .projection_state_cache
1815 .iter()
1816 .filter(|entry| entry.key().starts_with(&prefix))
1817 .map(|entry| entry.key().clone())
1818 .collect();
1819 for key in keys_to_remove {
1820 self.projection_state_cache.remove(&key);
1821 }
1822
1823 let events = self.events.read();
1825 let mut reprocessed = 0usize;
1826 for event in events.iter() {
1827 if projection.process(event).is_ok() {
1828 reprocessed += 1;
1829 }
1830 }
1831
1832 Ok(reprocessed)
1833 }
1834
1835 pub fn get_event_by_id(&self, event_id: &uuid::Uuid) -> Result<Option<Event>> {
1837 if let Some(offset) = self.index.get_by_id(event_id) {
1838 let events = self.events.read();
1839 Ok(events.get(offset).cloned())
1840 } else {
1841 Ok(None)
1842 }
1843 }
1844
1845 #[cfg_attr(feature = "hotpath", hotpath::measure)]
1849 pub fn query(&self, request: &QueryEventsRequest) -> Result<Vec<Event>> {
1850 self.query_window(request, 0, false)
1851 .map(|(events, _)| events)
1852 }
1853
1854 #[cfg_attr(feature = "hotpath", hotpath::measure)]
1871 pub fn query_window(
1872 &self,
1873 request: &QueryEventsRequest,
1874 offset: usize,
1875 descending: bool,
1876 ) -> Result<(Vec<Event>, usize)> {
1877 if let Some(filter) = &request.payload_filter
1884 && serde_json::from_str::<serde_json::Map<String, serde_json::Value>>(filter).is_err()
1885 {
1886 return Err(AllSourceError::InvalidInput(format!(
1887 "invalid 'payload_filter': expected a JSON object of field/value \
1888 pairs, got '{filter}'"
1889 )));
1890 }
1891
1892 if let Some(ref tenant_id) = request.tenant_id {
1911 self.ensure_tenant_loaded(tenant_id)?;
1912 self.tenant_loader.touch(tenant_id);
1916 }
1917
1918 let query_type = if request.entity_id.is_some() {
1920 "entity"
1921 } else if request.event_type.is_some() {
1922 "type"
1923 } else if request.event_type_prefix.is_some() {
1924 "type_prefix"
1925 } else {
1926 "full_scan"
1927 };
1928
1929 #[cfg(feature = "server")]
1931 let timer = self
1932 .metrics
1933 .query_duration_seconds
1934 .with_label_values(&[query_type])
1935 .start_timer();
1936
1937 #[cfg(feature = "server")]
1939 self.metrics
1940 .queries_total
1941 .with_label_values(&[query_type])
1942 .inc();
1943
1944 let events = self.events.read();
1945
1946 let offsets: Vec<usize> = if let Some(entity_id) = &request.entity_id {
1948 self.index
1950 .get_by_entity(entity_id)
1951 .map(|entries| self.filter_entries(entries, request))
1952 .unwrap_or_default()
1953 } else if let Some(event_type) = &request.event_type {
1954 self.index
1956 .get_by_type(event_type)
1957 .map(|entries| self.filter_entries(entries, request))
1958 .unwrap_or_default()
1959 } else if let Some(prefix) = &request.event_type_prefix {
1960 let entries = self.index.get_by_type_prefix(prefix);
1962 self.filter_entries(entries, request)
1963 } else {
1964 (0..events.len()).collect()
1966 };
1967
1968 let mut matches: Vec<(usize, &Event)> = offsets
1974 .iter()
1975 .filter_map(|&event_offset| events.get(event_offset))
1976 .filter(|event| self.apply_filters(event, request))
1977 .enumerate()
1978 .collect();
1979
1980 let total = matches.len();
1983
1984 let order = |a: &(usize, &Event), b: &(usize, &Event)| {
1991 let ascending =
1992 a.1.timestamp
1993 .cmp(&b.1.timestamp)
1994 .then_with(|| a.1.version.cmp(&b.1.version))
1995 .then_with(|| a.0.cmp(&b.0));
1996 if descending {
1997 ascending.reverse()
1998 } else {
1999 ascending
2000 }
2001 };
2002
2003 select_window(&mut matches, offset, request.limit, order);
2005
2006 let results: Vec<Event> = matches
2009 .into_iter()
2010 .skip(offset)
2011 .map(|(_, event)| event.clone())
2012 .collect();
2013
2014 #[cfg(feature = "server")]
2016 {
2017 self.metrics
2018 .query_results_total
2019 .with_label_values(&[query_type])
2020 .inc_by(results.len() as u64);
2021 timer.observe_duration();
2022 }
2023
2024 Ok((results, total))
2025 }
2026
2027 #[cfg_attr(feature = "hotpath", hotpath::measure)]
2029 fn filter_entries(&self, entries: Vec<IndexEntry>, request: &QueryEventsRequest) -> Vec<usize> {
2030 entries
2031 .into_iter()
2032 .filter(|entry| {
2033 if let Some(as_of) = request.as_of
2035 && entry.timestamp > as_of
2036 {
2037 return false;
2038 }
2039 if let Some(since) = request.since
2040 && entry.timestamp < since
2041 {
2042 return false;
2043 }
2044 if let Some(until) = request.until
2045 && entry.timestamp > until
2046 {
2047 return false;
2048 }
2049 true
2050 })
2051 .map(|entry| entry.offset)
2052 .collect()
2053 }
2054
2055 #[cfg_attr(feature = "hotpath", hotpath::measure)]
2057 fn apply_filters(&self, event: &Event, request: &QueryEventsRequest) -> bool {
2058 if let Some(ref tid) = request.tenant_id
2060 && event.tenant_id_str() != tid
2061 {
2062 return false;
2063 }
2064
2065 if let Some(as_of) = request.as_of
2075 && event.timestamp > as_of
2076 {
2077 return false;
2078 }
2079 if let Some(since) = request.since
2080 && event.timestamp < since
2081 {
2082 return false;
2083 }
2084 if let Some(until) = request.until
2085 && event.timestamp > until
2086 {
2087 return false;
2088 }
2089
2090 if let Some(ref excludes) = request.exclude_event_type_prefix {
2094 let et = event.event_type_str();
2095 if excludes
2096 .split(',')
2097 .map(str::trim)
2098 .filter(|p| !p.is_empty())
2099 .any(|p| et.starts_with(p))
2100 {
2101 return false;
2102 }
2103 }
2104
2105 if request.entity_id.is_some()
2107 && let Some(ref event_type) = request.event_type
2108 && event.event_type_str() != event_type
2109 {
2110 return false;
2111 }
2112
2113 if request.entity_id.is_some()
2115 && let Some(ref prefix) = request.event_type_prefix
2116 && !event.event_type_str().starts_with(prefix)
2117 {
2118 return false;
2119 }
2120
2121 if let Some(ref filter_str) = request.payload_filter
2123 && let Ok(filter_obj) =
2124 serde_json::from_str::<serde_json::Map<String, serde_json::Value>>(filter_str)
2125 {
2126 let payload = event.payload();
2127 for (key, expected_value) in &filter_obj {
2128 match payload.get(key) {
2129 Some(actual_value) if actual_value == expected_value => {}
2130 _ => return false,
2131 }
2132 }
2133 }
2134
2135 true
2136 }
2137
2138 #[cfg_attr(feature = "hotpath", hotpath::measure)]
2141 pub fn reconstruct_state(
2142 &self,
2143 entity_id: &str,
2144 as_of: Option<DateTime<Utc>>,
2145 ) -> Result<serde_json::Value> {
2146 let (merged_state, since_timestamp) = if let Some(as_of_time) = as_of {
2148 if let Some(snapshot) = self
2150 .snapshot_manager
2151 .get_snapshot_as_of(entity_id, as_of_time)
2152 {
2153 tracing::debug!(
2154 "Using snapshot from {} for entity {} (saved {} events)",
2155 snapshot.as_of,
2156 entity_id,
2157 snapshot.event_count
2158 );
2159 (snapshot.state.clone(), Some(snapshot.as_of))
2160 } else {
2161 (serde_json::json!({}), None)
2162 }
2163 } else {
2164 if let Some(snapshot) = self.snapshot_manager.get_latest_snapshot(entity_id) {
2166 tracing::debug!(
2167 "Using latest snapshot from {} for entity {}",
2168 snapshot.as_of,
2169 entity_id
2170 );
2171 (snapshot.state.clone(), Some(snapshot.as_of))
2172 } else {
2173 (serde_json::json!({}), None)
2174 }
2175 };
2176
2177 let events = self.query(&QueryEventsRequest {
2179 entity_id: Some(entity_id.to_string()),
2180 event_type: None,
2181 tenant_id: None,
2182 as_of,
2183 since: since_timestamp,
2184 until: None,
2185 limit: None,
2186 event_type_prefix: None,
2187 exclude_event_type_prefix: None,
2188 payload_filter: None,
2189 })?;
2190
2191 if events.is_empty() && since_timestamp.is_none() {
2193 return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
2194 }
2195
2196 let mut merged_state = merged_state;
2198 for event in &events {
2199 if let serde_json::Value::Object(ref mut state_map) = merged_state
2200 && let serde_json::Value::Object(ref payload_map) = event.payload
2201 {
2202 for (key, value) in payload_map {
2203 state_map.insert(key.clone(), value.clone());
2204 }
2205 }
2206 }
2207
2208 let state = serde_json::json!({
2210 "entity_id": entity_id,
2211 "last_updated": events.last().map(|e| e.timestamp),
2212 "event_count": events.len(),
2213 "as_of": as_of,
2214 "current_state": merged_state,
2215 "history": events.iter().map(|e| {
2216 serde_json::json!({
2217 "event_id": e.id,
2218 "type": e.event_type,
2219 "timestamp": e.timestamp,
2220 "payload": e.payload
2221 })
2222 }).collect::<Vec<_>>()
2223 });
2224
2225 Ok(state)
2226 }
2227
2228 pub fn get_snapshot(&self, entity_id: &str) -> Result<serde_json::Value> {
2230 let projections = self.projections.read();
2231
2232 if let Some(snapshot_projection) = projections.get_projection("entity_snapshots")
2233 && let Some(state) = snapshot_projection.get_state(entity_id)
2234 {
2235 return Ok(serde_json::json!({
2236 "entity_id": entity_id,
2237 "snapshot": state,
2238 "from_projection": "entity_snapshots"
2239 }));
2240 }
2241
2242 Err(AllSourceError::EntityNotFound(entity_id.to_string()))
2243 }
2244
2245 pub fn stats(&self) -> StoreStats {
2247 let events = self.events.read();
2248 let index_stats = self.index.stats();
2249
2250 StoreStats {
2251 total_events: events.len(),
2252 total_entities: index_stats.total_entities,
2253 total_event_types: index_stats.total_event_types,
2254 total_ingested: *self.total_ingested.read(),
2255 }
2256 }
2257
2258 pub fn list_streams(&self) -> Vec<StreamInfo> {
2260 self.index
2261 .get_all_entities()
2262 .into_iter()
2263 .map(|entity_id| {
2264 let event_count = self
2265 .index
2266 .get_by_entity(&entity_id)
2267 .map_or(0, |entries| entries.len());
2268 let last_event_at = self
2269 .index
2270 .get_by_entity(&entity_id)
2271 .and_then(|entries| entries.last().map(|e| e.timestamp));
2272 StreamInfo {
2273 stream_id: entity_id,
2274 event_count,
2275 last_event_at,
2276 }
2277 })
2278 .collect()
2279 }
2280
2281 pub fn list_event_types(&self) -> Vec<EventTypeInfo> {
2283 self.index
2284 .get_all_types()
2285 .into_iter()
2286 .map(|event_type| {
2287 let event_count = self
2288 .index
2289 .get_by_type(&event_type)
2290 .map_or(0, |entries| entries.len());
2291 let last_event_at = self
2292 .index
2293 .get_by_type(&event_type)
2294 .and_then(|entries| entries.last().map(|e| e.timestamp));
2295 EventTypeInfo {
2296 event_type,
2297 event_count,
2298 last_event_at,
2299 }
2300 })
2301 .collect()
2302 }
2303
2304 pub fn list_streams_for_tenant(&self, tenant_id: &str) -> Vec<StreamInfo> {
2312 let _ = self.ensure_tenant_loaded(tenant_id);
2313 let events = self.events.read();
2314 let mut by_entity: std::collections::HashMap<&str, (usize, chrono::DateTime<chrono::Utc>)> =
2315 std::collections::HashMap::new();
2316 for ev in events.iter() {
2317 if ev.tenant_id_str() != tenant_id {
2318 continue;
2319 }
2320 let e = by_entity
2321 .entry(ev.entity_id_str())
2322 .or_insert((0, ev.timestamp));
2323 e.0 += 1;
2324 if ev.timestamp > e.1 {
2325 e.1 = ev.timestamp;
2326 }
2327 }
2328 by_entity
2329 .into_iter()
2330 .map(|(entity_id, (count, last))| StreamInfo {
2331 stream_id: entity_id.to_string(),
2332 event_count: count,
2333 last_event_at: Some(last),
2334 })
2335 .collect()
2336 }
2337
2338 pub fn list_event_types_for_tenant(&self, tenant_id: &str) -> Vec<EventTypeInfo> {
2340 let _ = self.ensure_tenant_loaded(tenant_id);
2341 let events = self.events.read();
2342 let mut by_type: std::collections::HashMap<&str, (usize, chrono::DateTime<chrono::Utc>)> =
2343 std::collections::HashMap::new();
2344 for ev in events.iter() {
2345 if ev.tenant_id_str() != tenant_id {
2346 continue;
2347 }
2348 let e = by_type
2349 .entry(ev.event_type_str())
2350 .or_insert((0, ev.timestamp));
2351 e.0 += 1;
2352 if ev.timestamp > e.1 {
2353 e.1 = ev.timestamp;
2354 }
2355 }
2356 by_type
2357 .into_iter()
2358 .map(|(event_type, (count, last))| EventTypeInfo {
2359 event_type: event_type.to_string(),
2360 event_count: count,
2361 last_event_at: Some(last),
2362 })
2363 .collect()
2364 }
2365
2366 pub fn stats_for_tenant(&self, tenant_id: &str) -> TenantStoreStats {
2376 let _ = self.ensure_tenant_loaded(tenant_id);
2377 let events = self.events.read();
2378
2379 let mut entities: std::collections::HashSet<&str> = std::collections::HashSet::new();
2380 let mut census: std::collections::HashMap<&str, usize> = std::collections::HashMap::new();
2381 let mut total_events = 0usize;
2382 let mut oldest: Option<chrono::DateTime<chrono::Utc>> = None;
2383 let mut newest: Option<chrono::DateTime<chrono::Utc>> = None;
2384
2385 for ev in events.iter() {
2386 if ev.tenant_id_str() != tenant_id {
2387 continue;
2388 }
2389
2390 total_events += 1;
2391 entities.insert(ev.entity_id_str());
2392 *census.entry(ev.event_type_str()).or_insert(0) += 1;
2393
2394 let ts = ev.timestamp;
2395 if oldest.is_none_or(|o| ts < o) {
2396 oldest = Some(ts);
2397 }
2398 if newest.is_none_or(|n| ts > n) {
2399 newest = Some(ts);
2400 }
2401 }
2402
2403 TenantStoreStats {
2404 total_events,
2405 total_entities: entities.len(),
2406 total_event_types: census.len(),
2407 total_ingested: total_events as u64,
2411 event_types: census
2412 .into_iter()
2413 .map(|(k, v)| (k.to_string(), v))
2414 .collect(),
2415 oldest_event: oldest,
2416 newest_event: newest,
2417 }
2418 }
2419
2420 pub fn reconstruct_state_for_tenant(
2429 &self,
2430 entity_id: &str,
2431 as_of: Option<DateTime<Utc>>,
2432 tenant_id: &str,
2433 ) -> Result<serde_json::Value> {
2434 let events = self.query(&QueryEventsRequest {
2435 entity_id: Some(entity_id.to_string()),
2436 event_type: None,
2437 tenant_id: Some(tenant_id.to_string()),
2438 as_of,
2439 since: None,
2440 until: None,
2441 limit: None,
2442 event_type_prefix: None,
2443 exclude_event_type_prefix: None,
2444 payload_filter: None,
2445 })?;
2446
2447 if events.is_empty() {
2448 return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
2449 }
2450
2451 let mut merged_state = serde_json::json!({});
2452 for event in &events {
2453 if let serde_json::Value::Object(ref mut state_map) = merged_state
2454 && let serde_json::Value::Object(ref payload_map) = event.payload
2455 {
2456 for (key, value) in payload_map {
2457 state_map.insert(key.clone(), value.clone());
2458 }
2459 }
2460 }
2461
2462 Ok(serde_json::json!({
2463 "entity_id": entity_id,
2464 "last_updated": events.last().map(|e| e.timestamp),
2465 "event_count": events.len(),
2466 "as_of": as_of,
2467 "current_state": merged_state,
2468 "history": events.iter().map(|e| {
2469 serde_json::json!({
2470 "event_id": e.id,
2471 "type": e.event_type,
2472 "timestamp": e.timestamp,
2473 "payload": e.payload
2474 })
2475 }).collect::<Vec<_>>()
2476 }))
2477 }
2478
2479 pub fn enable_wal_replication(
2486 &self,
2487 tx: tokio::sync::broadcast::Sender<crate::infrastructure::persistence::wal::WALEntry>,
2488 ) {
2489 if let Some(ref wal_arc) = self.wal {
2490 wal_arc.set_replication_tx(tx);
2491 tracing::info!("WAL replication broadcast enabled");
2492 } else {
2493 tracing::warn!("Cannot enable WAL replication: WAL is not configured");
2494 }
2495 }
2496
2497 pub fn wal(&self) -> Option<&Arc<WriteAheadLog>> {
2500 self.wal.as_ref()
2501 }
2502
2503 pub fn parquet_storage(&self) -> Option<&Arc<RwLock<ParquetStorage>>> {
2506 self.storage.as_ref()
2507 }
2508}
2509
2510#[derive(Debug, Clone, Default)]
2512pub struct EventStoreConfig {
2513 pub storage_dir: Option<PathBuf>,
2515
2516 pub snapshot_config: SnapshotConfig,
2518
2519 pub wal_dir: Option<PathBuf>,
2521
2522 pub wal_config: WALConfig,
2524
2525 pub compaction_config: CompactionConfig,
2527
2528 pub schema_registry_config: SchemaRegistryConfig,
2530
2531 pub system_data_dir: Option<PathBuf>,
2536
2537 pub bootstrap_tenant: Option<String>,
2539
2540 pub cache_byte_budget: Option<u64>,
2547
2548 pub checkpoint_interval_secs: Option<u64>,
2560
2561 pub read_only: bool,
2565}
2566
2567impl EventStoreConfig {
2568 pub fn with_persistence(storage_dir: impl Into<PathBuf>) -> Self {
2570 Self {
2571 storage_dir: Some(storage_dir.into()),
2572 ..Self::default()
2573 }
2574 }
2575
2576 pub fn with_snapshots(snapshot_config: SnapshotConfig) -> Self {
2578 Self {
2579 snapshot_config,
2580 ..Self::default()
2581 }
2582 }
2583
2584 pub fn with_wal(wal_dir: impl Into<PathBuf>, wal_config: WALConfig) -> Self {
2586 Self {
2587 wal_dir: Some(wal_dir.into()),
2588 wal_config,
2589 ..Self::default()
2590 }
2591 }
2592
2593 pub fn with_all(storage_dir: impl Into<PathBuf>, snapshot_config: SnapshotConfig) -> Self {
2595 Self {
2596 storage_dir: Some(storage_dir.into()),
2597 snapshot_config,
2598 ..Self::default()
2599 }
2600 }
2601
2602 pub fn production(
2604 storage_dir: impl Into<PathBuf>,
2605 wal_dir: impl Into<PathBuf>,
2606 snapshot_config: SnapshotConfig,
2607 wal_config: WALConfig,
2608 compaction_config: CompactionConfig,
2609 ) -> Self {
2610 let storage_dir = storage_dir.into();
2611 let system_data_dir = storage_dir.join("__system");
2612 Self {
2613 storage_dir: Some(storage_dir),
2614 snapshot_config,
2615 wal_dir: Some(wal_dir.into()),
2616 wal_config,
2617 compaction_config,
2618 system_data_dir: Some(system_data_dir),
2619 ..Self::default()
2620 }
2621 }
2622
2623 pub fn effective_system_data_dir(&self) -> Option<PathBuf> {
2628 self.system_data_dir
2629 .clone()
2630 .or_else(|| self.storage_dir.as_ref().map(|d| d.join("__system")))
2631 }
2632
2633 pub fn from_env() -> (Self, &'static str) {
2641 Self::from_env_vars(
2642 std::env::var("ALLSOURCE_DATA_DIR")
2643 .ok()
2644 .filter(|s| !s.is_empty()),
2645 std::env::var("ALLSOURCE_STORAGE_DIR")
2646 .ok()
2647 .filter(|s| !s.is_empty()),
2648 std::env::var("ALLSOURCE_WAL_DIR")
2649 .ok()
2650 .filter(|s| !s.is_empty()),
2651 std::env::var("ALLSOURCE_WAL_ENABLED").ok(),
2652 std::env::var("ALLSOURCE_CACHE_BYTES").ok(),
2653 std::env::var("ALLSOURCE_SNAPSHOT_INTERVAL_SECONDS").ok(),
2654 std::env::var("ALLSOURCE_RETENTION_SYSTEM_DAYS").ok(),
2655 std::env::var("ALLSOURCE_CHECKPOINT_INTERVAL_SECONDS").ok(),
2656 )
2657 }
2658
2659 pub fn from_env_vars(
2661 data_dir: Option<String>,
2662 explicit_storage_dir: Option<String>,
2663 explicit_wal_dir: Option<String>,
2664 wal_enabled_var: Option<String>,
2665 cache_bytes_var: Option<String>,
2666 snapshot_interval_var: Option<String>,
2667 retention_system_days_var: Option<String>,
2668 checkpoint_interval_var: Option<String>,
2669 ) -> (Self, &'static str) {
2670 let data_dir = data_dir.filter(|s| !s.is_empty());
2671 let storage_dir = explicit_storage_dir
2672 .filter(|s| !s.is_empty())
2673 .or_else(|| data_dir.as_ref().map(|d| format!("{d}/storage")));
2674 let wal_dir = explicit_wal_dir
2675 .filter(|s| !s.is_empty())
2676 .or_else(|| data_dir.as_ref().map(|d| format!("{d}/wal")));
2677 let wal_enabled = wal_enabled_var.is_none_or(|v| v == "true");
2678 let cache_byte_budget =
2683 cache_bytes_var
2684 .filter(|s| !s.is_empty())
2685 .and_then(|s| match s.parse::<u64>() {
2686 Ok(v) => Some(v),
2687 Err(e) => {
2688 tracing::warn!(
2689 "ALLSOURCE_CACHE_BYTES={s:?} could not be parsed as u64: {e}; \
2690 cache budget disabled"
2691 );
2692 None
2693 }
2694 });
2695 let compaction_config =
2696 CompactionConfig::from_env_vars(snapshot_interval_var, retention_system_days_var);
2697
2698 let checkpoint_interval_secs = if wal_enabled {
2703 checkpoint_interval_var
2704 .filter(|s| !s.is_empty())
2705 .map(|s| match s.parse::<u64>() {
2706 Ok(v) => v,
2707 Err(e) => {
2708 tracing::warn!(
2709 "ALLSOURCE_CHECKPOINT_INTERVAL_SECONDS={s:?} could not be parsed as \
2710 u64: {e}; falling back to default 60s"
2711 );
2712 60
2713 }
2714 })
2715 .or(Some(60))
2716 } else {
2717 None
2718 };
2719
2720 let mut config = match (&storage_dir, &wal_dir) {
2721 (Some(sd), Some(wd)) if wal_enabled => Self::production(
2722 sd,
2723 wd,
2724 SnapshotConfig::default(),
2725 WALConfig::default(),
2726 compaction_config,
2727 ),
2728 (Some(sd), _) => Self::with_persistence(sd),
2729 (_, Some(wd)) if wal_enabled => Self::with_wal(wd, WALConfig::default()),
2730 _ => Self::default(),
2731 };
2732 config.cache_byte_budget = cache_byte_budget;
2733 config.checkpoint_interval_secs = checkpoint_interval_secs;
2734
2735 let mode = match (&storage_dir, &wal_dir) {
2736 (Some(_), Some(_)) if wal_enabled => "wal+parquet",
2737 (Some(_), _) => "parquet-only",
2738 (_, Some(_)) if wal_enabled => "wal-only",
2739 _ => "in-memory",
2740 };
2741 (config, mode)
2742 }
2743}
2744
2745#[derive(Debug, serde::Serialize)]
2746pub struct StoreStats {
2747 pub total_events: usize,
2748 pub total_entities: usize,
2749 pub total_event_types: usize,
2750 pub total_ingested: u64,
2751}
2752
2753#[derive(Debug, Clone, serde::Serialize)]
2759pub struct TenantStoreStats {
2760 pub total_events: usize,
2761 pub total_entities: usize,
2762 pub total_event_types: usize,
2763 pub total_ingested: u64,
2764 pub event_types: std::collections::HashMap<String, usize>,
2766 pub oldest_event: Option<chrono::DateTime<chrono::Utc>>,
2767 pub newest_event: Option<chrono::DateTime<chrono::Utc>>,
2768}
2769
2770#[derive(Debug, Clone, serde::Serialize)]
2772pub struct StreamInfo {
2773 pub stream_id: String,
2775 pub event_count: usize,
2777 pub last_event_at: Option<chrono::DateTime<chrono::Utc>>,
2779}
2780
2781#[derive(Debug, Clone, serde::Serialize)]
2783pub struct EventTypeInfo {
2784 pub event_type: String,
2786 pub event_count: usize,
2788 pub last_event_at: Option<chrono::DateTime<chrono::Utc>>,
2790}
2791
2792impl Default for EventStore {
2793 fn default() -> Self {
2794 Self::new()
2795 }
2796}
2797
2798#[cfg(test)]
2799mod tests {
2800 use super::*;
2801 use crate::domain::entities::Event;
2802 use tempfile::TempDir;
2803
2804 fn find_parquet_files(dir: &std::path::Path) -> Vec<std::path::PathBuf> {
2809 let mut out = Vec::new();
2810 let mut stack = vec![dir.to_path_buf()];
2811 while let Some(d) = stack.pop() {
2812 let Ok(entries) = std::fs::read_dir(&d) else {
2813 continue;
2814 };
2815 for e in entries.flatten() {
2816 let p = e.path();
2817 if p.is_dir() {
2818 stack.push(p);
2819 } else if p.extension().and_then(|s| s.to_str()) == Some("parquet") {
2820 out.push(p);
2821 }
2822 }
2823 }
2824 out
2825 }
2826
2827 fn create_test_event(entity_id: &str, event_type: &str) -> Event {
2828 Event::from_strings(
2829 event_type.to_string(),
2830 entity_id.to_string(),
2831 "default".to_string(),
2832 serde_json::json!({"name": "Test", "value": 42}),
2833 None,
2834 )
2835 .unwrap()
2836 }
2837
2838 fn create_test_event_with_payload(
2839 entity_id: &str,
2840 event_type: &str,
2841 payload: serde_json::Value,
2842 ) -> Event {
2843 Event::from_strings(
2844 event_type.to_string(),
2845 entity_id.to_string(),
2846 "default".to_string(),
2847 payload,
2848 None,
2849 )
2850 .unwrap()
2851 }
2852
2853 #[test]
2854 fn test_event_store_new() {
2855 let store = EventStore::new();
2856 assert_eq!(store.stats().total_events, 0);
2857 assert_eq!(store.stats().total_entities, 0);
2858 }
2859
2860 #[test]
2867 fn test_ensure_tenant_loaded_no_storage_is_a_noop() {
2868 let store = EventStore::new();
2872 assert!(!store.is_tenant_loaded("alice"));
2873 store.ensure_tenant_loaded("alice").unwrap();
2874 assert!(store.is_tenant_loaded("alice"));
2875 assert!(!store.is_tenant_loaded("bob"));
2877 }
2878
2879 #[test]
2880 fn test_ensure_tenant_loaded_warm_path_is_idempotent() {
2881 let store = EventStore::new();
2882 store.ensure_tenant_loaded("alice").unwrap();
2883 store.ensure_tenant_loaded("alice").unwrap();
2885 }
2886
2887 #[test]
2888 fn test_ensure_tenant_loaded_rejects_unsafe_tenant_id() {
2889 let temp_dir = TempDir::new().unwrap();
2895 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
2896 for unsafe_tid in ["..", "a/b", "a\\b", ""] {
2897 let result = store.ensure_tenant_loaded(unsafe_tid);
2898 assert!(
2899 result.is_err(),
2900 "tenant_id {unsafe_tid:?} should have been rejected"
2901 );
2902 assert!(
2903 !store.is_tenant_loaded(unsafe_tid),
2904 "rejected tenant {unsafe_tid:?} must not be marked loaded"
2905 );
2906 }
2907 }
2908
2909 #[test]
2910 fn test_ensure_tenant_loaded_no_subtree_marks_loaded_with_zero_events() {
2911 let temp_dir = TempDir::new().unwrap();
2916 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
2917 assert!(!store.is_tenant_loaded("never-existed"));
2918 store.ensure_tenant_loaded("never-existed").unwrap();
2919 assert!(store.is_tenant_loaded("never-existed"));
2920 }
2921
2922 #[test]
2923 fn test_evict_tenant_drops_events_and_resets_bytes() {
2924 let temp_dir = TempDir::new().unwrap();
2928 let storage_dir = temp_dir.path().to_path_buf();
2929
2930 {
2931 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
2932 for i in 0..3 {
2933 store
2934 .ingest(
2935 &Event::from_strings(
2936 "test.event".to_string(),
2937 format!("a-{i}"),
2938 "alice".to_string(),
2939 serde_json::json!({"i": i}),
2940 None,
2941 )
2942 .unwrap(),
2943 )
2944 .unwrap();
2945 }
2946 for i in 0..2 {
2947 store
2948 .ingest(
2949 &Event::from_strings(
2950 "test.event".to_string(),
2951 format!("b-{i}"),
2952 "bob".to_string(),
2953 serde_json::json!({"i": i}),
2954 None,
2955 )
2956 .unwrap(),
2957 )
2958 .unwrap();
2959 }
2960 store.flush_storage().unwrap();
2961 }
2962
2963 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
2964 store.ensure_tenant_loaded("alice").unwrap();
2965 store.ensure_tenant_loaded("bob").unwrap();
2966 assert_eq!(store.stats().total_events, 5);
2967 let alice_bytes = store.tenant_resident_bytes("alice");
2968 let bob_bytes = store.tenant_resident_bytes("bob");
2969 assert!(alice_bytes > 0 && bob_bytes > 0);
2970
2971 store.evict_tenant("alice");
2972
2973 assert!(!store.is_tenant_loaded("alice"));
2974 assert!(store.is_tenant_loaded("bob"));
2975 assert_eq!(store.tenant_resident_bytes("alice"), 0);
2976 assert_eq!(store.tenant_resident_bytes("bob"), bob_bytes);
2977 assert_eq!(store.stats().total_events, 2, "only bob's 2 events remain");
2978 }
2979
2980 #[test]
2981 fn test_evict_tenant_then_query_re_loads_from_disk() {
2982 let temp_dir = TempDir::new().unwrap();
2986 let storage_dir = temp_dir.path().to_path_buf();
2987
2988 {
2989 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
2990 for i in 0..4 {
2991 store
2992 .ingest(
2993 &Event::from_strings(
2994 "test.event".to_string(),
2995 format!("a-{i}"),
2996 "alice".to_string(),
2997 serde_json::json!({"i": i}),
2998 None,
2999 )
3000 .unwrap(),
3001 )
3002 .unwrap();
3003 }
3004 store.flush_storage().unwrap();
3005 }
3006
3007 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3008 store.ensure_tenant_loaded("alice").unwrap();
3009 store.evict_tenant("alice");
3010 assert_eq!(store.stats().total_events, 0);
3011
3012 let results = store
3014 .query(&QueryEventsRequest {
3015 entity_id: None,
3016 event_type: None,
3017 tenant_id: Some("alice".to_string()),
3018 as_of: None,
3019 since: None,
3020 until: None,
3021 limit: None,
3022 event_type_prefix: None,
3023 exclude_event_type_prefix: None,
3024 payload_filter: None,
3025 })
3026 .unwrap();
3027 assert_eq!(results.len(), 4);
3028 assert!(store.is_tenant_loaded("alice"));
3029 }
3030
3031 #[test]
3032 fn test_evict_tenant_rebuilds_index_with_new_offsets() {
3033 let temp_dir = TempDir::new().unwrap();
3040 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3041
3042 for i in 0..3 {
3046 store
3047 .ingest(
3048 &Event::from_strings(
3049 "test.event".to_string(),
3050 format!("a-{i}"),
3051 "alice".to_string(),
3052 serde_json::json!({"i": i}),
3053 None,
3054 )
3055 .unwrap(),
3056 )
3057 .unwrap();
3058 if i < 2 {
3059 store
3060 .ingest(
3061 &Event::from_strings(
3062 "test.event".to_string(),
3063 format!("b-{i}"),
3064 "bob".to_string(),
3065 serde_json::json!({"i": i}),
3066 None,
3067 )
3068 .unwrap(),
3069 )
3070 .unwrap();
3071 }
3072 }
3073 store.tenant_loader.mark_loaded("alice");
3075 store.tenant_loader.mark_loaded("bob");
3076
3077 store.evict_tenant("alice");
3078
3079 let bob_results = store
3080 .query(&QueryEventsRequest {
3081 entity_id: None,
3082 event_type: None,
3083 tenant_id: Some("bob".to_string()),
3084 as_of: None,
3085 since: None,
3086 until: None,
3087 limit: None,
3088 event_type_prefix: None,
3089 exclude_event_type_prefix: None,
3090 payload_filter: None,
3091 })
3092 .unwrap();
3093 assert_eq!(bob_results.len(), 2);
3094 for e in &bob_results {
3095 assert_eq!(e.tenant_id_str(), "bob");
3096 }
3097 }
3098
3099 #[test]
3100 fn test_budget_eviction_keeps_resident_set_bounded() {
3101 let temp_dir = TempDir::new().unwrap();
3105 let storage_dir = temp_dir.path().to_path_buf();
3106
3107 let big_payload = serde_json::json!({"data": "x".repeat(1000)});
3110 {
3111 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3112 for tenant in ["alice", "bob", "carol"] {
3113 for i in 0..5 {
3114 store
3115 .ingest(
3116 &Event::from_strings(
3117 "test.event".to_string(),
3118 format!("{tenant}-{i}"),
3119 tenant.to_string(),
3120 big_payload.clone(),
3121 None,
3122 )
3123 .unwrap(),
3124 )
3125 .unwrap();
3126 }
3127 }
3128 store.flush_storage().unwrap();
3129 }
3130
3131 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3134 config.cache_byte_budget = Some(12_000);
3135 let store = EventStore::with_config(config);
3136
3137 store.ensure_tenant_loaded("alice").unwrap();
3139 assert!(store.is_tenant_loaded("alice"));
3140
3141 store.tenant_loader.touch("alice");
3146 std::thread::sleep(std::time::Duration::from_millis(10));
3147 store.ensure_tenant_loaded("bob").unwrap();
3148 assert!(store.is_tenant_loaded("bob"));
3149
3150 store.tenant_loader.touch("bob");
3154 std::thread::sleep(std::time::Duration::from_millis(10));
3155 store.ensure_tenant_loaded("carol").unwrap();
3156 assert!(store.is_tenant_loaded("carol"));
3157
3158 let resident = store.cache_resident_bytes();
3163 let budget = 12_000u64;
3164
3165 if resident > budget {
3168 let loaded_count = ["alice", "bob", "carol"]
3169 .iter()
3170 .filter(|t| store.is_tenant_loaded(t))
3171 .count();
3172 assert_eq!(
3173 loaded_count, 1,
3174 "over budget but more than one tenant loaded — eviction policy didn't fire"
3175 );
3176 }
3177
3178 assert!(store.is_tenant_loaded("carol"));
3181 }
3182
3183 #[test]
3184 fn test_query_after_eviction_re_loads_transparently() {
3185 let temp_dir = TempDir::new().unwrap();
3188 let storage_dir = temp_dir.path().to_path_buf();
3189
3190 let big_payload = serde_json::json!({"data": "x".repeat(2000)});
3191 {
3192 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3193 for tenant in ["alice", "bob"] {
3194 for i in 0..3 {
3195 store
3196 .ingest(
3197 &Event::from_strings(
3198 "test.event".to_string(),
3199 format!("{tenant}-{i}"),
3200 tenant.to_string(),
3201 big_payload.clone(),
3202 None,
3203 )
3204 .unwrap(),
3205 )
3206 .unwrap();
3207 }
3208 }
3209 store.flush_storage().unwrap();
3210 }
3211
3212 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3214 config.cache_byte_budget = Some(5_000);
3215 let store = EventStore::with_config(config);
3216
3217 let alice_first = store
3221 .query(&QueryEventsRequest {
3222 entity_id: None,
3223 event_type: None,
3224 tenant_id: Some("alice".to_string()),
3225 as_of: None,
3226 since: None,
3227 until: None,
3228 limit: None,
3229 event_type_prefix: None,
3230 exclude_event_type_prefix: None,
3231 payload_filter: None,
3232 })
3233 .unwrap();
3234 assert_eq!(alice_first.len(), 3);
3235
3236 std::thread::sleep(std::time::Duration::from_millis(15));
3238 let _bob = store
3240 .query(&QueryEventsRequest {
3241 entity_id: None,
3242 event_type: None,
3243 tenant_id: Some("bob".to_string()),
3244 as_of: None,
3245 since: None,
3246 until: None,
3247 limit: None,
3248 event_type_prefix: None,
3249 exclude_event_type_prefix: None,
3250 payload_filter: None,
3251 })
3252 .unwrap();
3253 assert!(
3254 !store.is_tenant_loaded("alice"),
3255 "alice should have been evicted"
3256 );
3257
3258 let alice_second = store
3260 .query(&QueryEventsRequest {
3261 entity_id: None,
3262 event_type: None,
3263 tenant_id: Some("alice".to_string()),
3264 as_of: None,
3265 since: None,
3266 until: None,
3267 limit: None,
3268 event_type_prefix: None,
3269 exclude_event_type_prefix: None,
3270 payload_filter: None,
3271 })
3272 .unwrap();
3273 assert_eq!(
3274 alice_second.len(),
3275 3,
3276 "alice's events come back via re-load"
3277 );
3278 assert!(store.is_tenant_loaded("alice"));
3279 }
3280
3281 #[test]
3282 #[cfg(feature = "server")]
3283 fn test_cache_metrics_track_evictions_and_bytes() {
3284 let temp_dir = TempDir::new().unwrap();
3288 let storage_dir = temp_dir.path().to_path_buf();
3289
3290 let big_payload = serde_json::json!({"data": "x".repeat(2000)});
3291 {
3292 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3293 for tenant in ["alice", "bob"] {
3294 for i in 0..3 {
3295 store
3296 .ingest(
3297 &Event::from_strings(
3298 "test.event".to_string(),
3299 format!("{tenant}-{i}"),
3300 tenant.to_string(),
3301 big_payload.clone(),
3302 None,
3303 )
3304 .unwrap(),
3305 )
3306 .unwrap();
3307 }
3308 }
3309 store.flush_storage().unwrap();
3310 }
3311
3312 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3313 config.cache_byte_budget = Some(5_000); let store = EventStore::with_config(config);
3315
3316 assert_eq!(store.metrics.cache_evictions_total.get(), 0);
3317 assert_eq!(store.metrics.cache_bytes.get(), 0);
3318
3319 store.ensure_tenant_loaded("alice").unwrap();
3320 let after_alice = store.metrics.cache_bytes.get();
3322 assert!(after_alice > 0, "gauge should reflect alice's bytes");
3323 assert_eq!(store.metrics.cache_evictions_total.get(), 0);
3325
3326 std::thread::sleep(std::time::Duration::from_millis(10));
3327 store.ensure_tenant_loaded("bob").unwrap();
3328
3329 assert_eq!(
3332 store.metrics.cache_evictions_total.get(),
3333 1,
3334 "exactly one tenant evicted after bob's load"
3335 );
3336 let after_bob = store.metrics.cache_bytes.get();
3338 assert!(after_bob > 0);
3339 assert!(after_bob <= after_alice, "gauge dropped after eviction");
3340 }
3341
3342 #[test]
3343 #[cfg(feature = "server")]
3344 fn test_storage_size_gauge_populated_from_on_disk_bytes() {
3345 let temp_dir = TempDir::new().unwrap();
3350 let config = EventStoreConfig {
3353 storage_dir: Some(temp_dir.path().join("parquet")),
3354 wal_dir: Some(temp_dir.path().join("wal")),
3355 ..EventStoreConfig::default()
3356 };
3357 let store = EventStore::with_config(config);
3358
3359 assert_eq!(store.metrics.storage_size_bytes.get(), 0);
3361
3362 let payload = serde_json::json!({ "data": "x".repeat(2000) });
3363 for i in 0..10 {
3364 store
3365 .ingest(
3366 &Event::from_strings(
3367 "test.event".to_string(),
3368 format!("entity-{i}"),
3369 "tenant-a".to_string(),
3370 payload.clone(),
3371 None,
3372 )
3373 .unwrap(),
3374 )
3375 .unwrap();
3376 }
3377 store.flush_storage().unwrap();
3378
3379 store.refresh_storage_metrics_now();
3381
3382 let size = store.metrics.storage_size_bytes.get();
3383 assert!(
3384 size > 0,
3385 "storage_size_bytes must reflect real on-disk bytes, got {size}"
3386 );
3387 assert!(
3388 store.metrics.parquet_files_total.get() >= 1,
3389 "at least one Parquet file should exist after a flush"
3390 );
3391
3392 assert!(
3395 store.metrics.wal_segments_total.get() >= 1,
3396 "at least one WAL segment should exist after writes"
3397 );
3398
3399 let parquet_stats = store.storage.as_ref().unwrap().read().stats().unwrap();
3402 let (wal_bytes, _) = store.wal.as_ref().unwrap().on_disk_stats().unwrap();
3403 assert_eq!(
3404 size as u64,
3405 parquet_stats.total_size_bytes + wal_bytes,
3406 "gauge should equal Parquet bytes ({}) + WAL bytes ({wal_bytes})",
3407 parquet_stats.total_size_bytes
3408 );
3409 }
3410
3411 #[test]
3412 fn test_stress_resident_set_stays_near_budget_under_rolling_queries() {
3413 let temp_dir = TempDir::new().unwrap();
3420 let storage_dir = temp_dir.path().to_path_buf();
3421
3422 const TENANT_COUNT: usize = 10;
3423 const EVENTS_PER_TENANT: usize = 50;
3424 let big_payload = serde_json::json!({"data": "x".repeat(10_000)});
3426
3427 {
3430 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3431 for t in 0..TENANT_COUNT {
3432 let tenant = format!("tenant-{t}");
3433 for i in 0..EVENTS_PER_TENANT {
3434 store
3435 .ingest(
3436 &Event::from_strings(
3437 "test.event".to_string(),
3438 format!("{tenant}-{i}"),
3439 tenant.clone(),
3440 big_payload.clone(),
3441 None,
3442 )
3443 .unwrap(),
3444 )
3445 .unwrap();
3446 }
3447 }
3448 store.flush_storage().unwrap();
3449 }
3450
3451 const BUDGET: u64 = 1_048_576;
3455 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3456 config.cache_byte_budget = Some(BUDGET);
3457 let store = EventStore::with_config(config);
3458
3459 let mut peak_resident: u64 = 0;
3463 for t in 0..TENANT_COUNT {
3464 let tenant = format!("tenant-{t}");
3465 let results = store
3466 .query(&QueryEventsRequest {
3467 entity_id: None,
3468 event_type: None,
3469 tenant_id: Some(tenant.clone()),
3470 as_of: None,
3471 since: None,
3472 until: None,
3473 limit: None,
3474 event_type_prefix: None,
3475 exclude_event_type_prefix: None,
3476 payload_filter: None,
3477 })
3478 .unwrap();
3479 assert_eq!(
3480 results.len(),
3481 EVENTS_PER_TENANT,
3482 "every per-tenant query must return all of that tenant's events"
3483 );
3484 let resident = store.cache_resident_bytes();
3486 if resident > peak_resident {
3487 peak_resident = resident;
3488 }
3489 }
3490
3491 let final_resident = store.cache_resident_bytes();
3492
3493 let tolerance = BUDGET; assert!(
3499 peak_resident <= BUDGET + tolerance,
3500 "peak resident {peak_resident} exceeds budget {BUDGET} by more than {tolerance} \
3501 — eviction policy not keeping up with the working-set churn"
3502 );
3503 assert!(
3504 final_resident <= BUDGET + tolerance,
3505 "final resident {final_resident} exceeds budget {BUDGET} by more than {tolerance}"
3506 );
3507
3508 let last_tenant = format!("tenant-{}", TENANT_COUNT - 1);
3511 assert!(
3512 store.is_tenant_loaded(&last_tenant),
3513 "the most-recent tenant must remain loaded after the sweep"
3514 );
3515
3516 let still_loaded = (0..TENANT_COUNT)
3519 .filter(|t| store.is_tenant_loaded(&format!("tenant-{t}")))
3520 .count();
3521 assert!(
3522 still_loaded < TENANT_COUNT,
3523 "no tenants evicted ({still_loaded}/{TENANT_COUNT} still loaded) — \
3524 budget enforcement didn't engage"
3525 );
3526 }
3527
3528 #[test]
3529 fn test_evict_tenant_when_not_loaded_is_a_noop() {
3530 let store = EventStore::new();
3533 store.evict_tenant("nobody"); assert!(!store.is_tenant_loaded("nobody"));
3535 }
3536
3537 #[test]
3538 fn test_lazy_load_accounts_bytes_per_tenant() {
3539 let temp_dir = TempDir::new().unwrap();
3543 let storage_dir = temp_dir.path().to_path_buf();
3544
3545 {
3547 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3548 for i in 0..5 {
3549 store
3550 .ingest(
3551 &Event::from_strings(
3552 "test.event".to_string(),
3553 format!("a-{i}"),
3554 "alice".to_string(),
3555 serde_json::json!({"data": "x".repeat(1000)}),
3556 None,
3557 )
3558 .unwrap(),
3559 )
3560 .unwrap();
3561 }
3562 store.flush_storage().unwrap();
3563 }
3564
3565 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3566 assert_eq!(store.tenant_resident_bytes("alice"), 0);
3568 assert_eq!(store.cache_resident_bytes(), 0);
3569
3570 store.ensure_tenant_loaded("alice").unwrap();
3571
3572 let alice_bytes = store.tenant_resident_bytes("alice");
3575 assert!(
3576 alice_bytes >= 5 * 1000,
3577 "alice should have at least 5 KiB resident; got {alice_bytes}"
3578 );
3579 assert_eq!(store.tenant_resident_bytes("bob"), 0);
3581 assert_eq!(store.cache_resident_bytes(), alice_bytes);
3583 }
3584
3585 #[test]
3586 fn test_query_lazy_loads_tenant_on_first_call() {
3587 let temp_dir = TempDir::new().unwrap();
3591 let storage_dir = temp_dir.path().to_path_buf();
3592
3593 {
3595 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3596 for i in 0..3 {
3597 let event = Event::from_strings(
3598 "test.event".to_string(),
3599 format!("e-{i}"),
3600 "alice".to_string(),
3601 serde_json::json!({"i": i}),
3602 None,
3603 )
3604 .unwrap();
3605 store.ingest(&event).unwrap();
3606 }
3607 store.flush_storage().unwrap();
3608 }
3609
3610 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3612 assert_eq!(
3613 store.stats().total_events,
3614 0,
3615 "boot must be O(1) — no Parquet pre-load"
3616 );
3617 assert!(!store.is_tenant_loaded("alice"));
3618 assert!(!store.is_tenant_loaded("bob"));
3619
3620 let results = store
3622 .query(&QueryEventsRequest {
3623 entity_id: None,
3624 event_type: None,
3625 tenant_id: Some("alice".to_string()),
3626 as_of: None,
3627 since: None,
3628 until: None,
3629 limit: None,
3630 event_type_prefix: None,
3631 exclude_event_type_prefix: None,
3632 payload_filter: None,
3633 })
3634 .unwrap();
3635 assert_eq!(results.len(), 3, "alice's 3 events are returned");
3636 assert!(store.is_tenant_loaded("alice"), "alice now warm");
3637 assert!(!store.is_tenant_loaded("bob"), "bob still cold");
3640 }
3641
3642 #[test]
3643 fn test_query_invalid_tenant_id_returns_error_no_hang() {
3644 let temp_dir = TempDir::new().unwrap();
3648 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3649
3650 let result = store.query(&QueryEventsRequest {
3651 entity_id: None,
3652 event_type: None,
3653 tenant_id: Some("../etc".to_string()),
3654 as_of: None,
3655 since: None,
3656 until: None,
3657 limit: None,
3658 event_type_prefix: None,
3659 exclude_event_type_prefix: None,
3660 payload_filter: None,
3661 });
3662 assert!(result.is_err(), "unsafe tenant_id must surface as error");
3663 }
3664
3665 #[test]
3666 fn test_query_concurrent_first_queries_for_same_tenant_all_succeed() {
3667 let temp_dir = TempDir::new().unwrap();
3675 let storage_dir = temp_dir.path().to_path_buf();
3676
3677 {
3679 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3680 for i in 0..25 {
3681 let event = Event::from_strings(
3682 "test.event".to_string(),
3683 format!("e-{i}"),
3684 "alice".to_string(),
3685 serde_json::json!({"i": i}),
3686 None,
3687 )
3688 .unwrap();
3689 store.ingest(&event).unwrap();
3690 }
3691 store.flush_storage().unwrap();
3692 }
3693
3694 let store = Arc::new(EventStore::with_config(EventStoreConfig::with_persistence(
3696 &storage_dir,
3697 )));
3698 assert!(!store.is_tenant_loaded("alice"));
3699
3700 let mut handles = Vec::new();
3701 for _ in 0..8 {
3702 let s = store.clone();
3703 handles.push(std::thread::spawn(move || {
3704 s.query(&QueryEventsRequest {
3705 entity_id: None,
3706 event_type: None,
3707 tenant_id: Some("alice".to_string()),
3708 as_of: None,
3709 since: None,
3710 until: None,
3711 limit: None,
3712 event_type_prefix: None,
3713 exclude_event_type_prefix: None,
3714 payload_filter: None,
3715 })
3716 }));
3717 }
3718
3719 for h in handles {
3720 let result = h.join().unwrap().unwrap();
3721 assert_eq!(
3722 result.len(),
3723 25,
3724 "every concurrent caller must see all 25 events"
3725 );
3726 }
3727 assert!(store.is_tenant_loaded("alice"));
3728 assert_eq!(store.stats().total_events, 25);
3730 }
3731
3732 #[test]
3733 fn test_query_two_cold_tenants_load_independently() {
3734 let temp_dir = TempDir::new().unwrap();
3738 let storage_dir = temp_dir.path().to_path_buf();
3739
3740 {
3741 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3742 for i in 0..3 {
3743 store
3744 .ingest(
3745 &Event::from_strings(
3746 "test.event".to_string(),
3747 format!("a-{i}"),
3748 "alice".to_string(),
3749 serde_json::json!({"i": i}),
3750 None,
3751 )
3752 .unwrap(),
3753 )
3754 .unwrap();
3755 }
3756 for i in 0..5 {
3757 store
3758 .ingest(
3759 &Event::from_strings(
3760 "test.event".to_string(),
3761 format!("b-{i}"),
3762 "bob".to_string(),
3763 serde_json::json!({"i": i}),
3764 None,
3765 )
3766 .unwrap(),
3767 )
3768 .unwrap();
3769 }
3770 store.flush_storage().unwrap();
3771 }
3772
3773 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3774 assert_eq!(store.stats().total_events, 0);
3775
3776 let alice = store
3778 .query(&QueryEventsRequest {
3779 entity_id: None,
3780 event_type: None,
3781 tenant_id: Some("alice".to_string()),
3782 as_of: None,
3783 since: None,
3784 until: None,
3785 limit: None,
3786 event_type_prefix: None,
3787 exclude_event_type_prefix: None,
3788 payload_filter: None,
3789 })
3790 .unwrap();
3791 assert_eq!(alice.len(), 3);
3792 assert!(store.is_tenant_loaded("alice"));
3793 assert!(!store.is_tenant_loaded("bob"));
3794 assert_eq!(store.stats().total_events, 3);
3795
3796 let bob = store
3798 .query(&QueryEventsRequest {
3799 entity_id: None,
3800 event_type: None,
3801 tenant_id: Some("bob".to_string()),
3802 as_of: None,
3803 since: None,
3804 until: None,
3805 limit: None,
3806 event_type_prefix: None,
3807 exclude_event_type_prefix: None,
3808 payload_filter: None,
3809 })
3810 .unwrap();
3811 assert_eq!(bob.len(), 5);
3812 assert!(store.is_tenant_loaded("bob"));
3813 assert_eq!(store.stats().total_events, 8);
3814 }
3815
3816 #[test]
3817 fn test_boot_with_persisted_data_is_o1() {
3818 let temp_dir = TempDir::new().unwrap();
3831 let storage_dir = temp_dir.path().to_path_buf();
3832
3833 {
3834 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3835 for tenant in ["alice", "bob", "carol"] {
3836 for i in 0..50 / 3 {
3837 store
3838 .ingest(
3839 &Event::from_strings(
3840 "test.event".to_string(),
3841 format!("{tenant}-{i}"),
3842 tenant.to_string(),
3843 serde_json::json!({"i": i}),
3844 None,
3845 )
3846 .unwrap(),
3847 )
3848 .unwrap();
3849 }
3850 }
3851 store.flush_storage().unwrap();
3852 }
3853
3854 let on_disk = find_parquet_files(&storage_dir);
3856 assert!(
3857 !on_disk.is_empty(),
3858 "session 1 should have produced parquet files; pre-condition for the test"
3859 );
3860
3861 let started = std::time::Instant::now();
3862 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3863 let boot_elapsed = started.elapsed();
3864
3865 assert_eq!(
3866 store.stats().total_events,
3867 0,
3868 "boot must not pre-load any Parquet events"
3869 );
3870
3871 assert!(
3875 boot_elapsed < std::time::Duration::from_secs(2),
3876 "boot took {boot_elapsed:?} — Step 2 boot should be O(1)"
3877 );
3878 }
3879
3880 #[test]
3881 fn test_query_warm_tenant_does_not_re_read_disk() {
3882 let temp_dir = TempDir::new().unwrap();
3888 let storage_dir = temp_dir.path().to_path_buf();
3889
3890 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3891 for i in 0..3 {
3892 let event = Event::from_strings(
3893 "test.event".to_string(),
3894 format!("e-{i}"),
3895 "alice".to_string(),
3896 serde_json::json!({"i": i}),
3897 None,
3898 )
3899 .unwrap();
3900 store.ingest(&event).unwrap();
3901 }
3902 store.flush_storage().unwrap();
3903
3904 let _ = store
3906 .query(&QueryEventsRequest {
3907 entity_id: None,
3908 event_type: None,
3909 tenant_id: Some("alice".to_string()),
3910 as_of: None,
3911 since: None,
3912 until: None,
3913 limit: None,
3914 event_type_prefix: None,
3915 exclude_event_type_prefix: None,
3916 payload_filter: None,
3917 })
3918 .unwrap();
3919 assert!(store.is_tenant_loaded("alice"));
3920
3921 let parquet_files = find_parquet_files(&storage_dir);
3924 for f in parquet_files {
3925 std::fs::remove_file(&f).unwrap();
3926 }
3927
3928 let results = store
3929 .query(&QueryEventsRequest {
3930 entity_id: None,
3931 event_type: None,
3932 tenant_id: Some("alice".to_string()),
3933 as_of: None,
3934 since: None,
3935 until: None,
3936 limit: None,
3937 event_type_prefix: None,
3938 exclude_event_type_prefix: None,
3939 payload_filter: None,
3940 })
3941 .unwrap();
3942 assert_eq!(
3943 results.len(),
3944 3,
3945 "warm tenant query must not need disk; got {} events from a deleted parquet",
3946 results.len()
3947 );
3948 }
3949
3950 #[test]
3951 fn test_event_store_default() {
3952 let store = EventStore::default();
3953 assert_eq!(store.stats().total_events, 0);
3954 }
3955
3956 #[test]
3957 fn test_ingest_single_event() {
3958 let store = EventStore::new();
3959 let event = create_test_event("entity-1", "user.created");
3960
3961 store.ingest(&event).unwrap();
3962
3963 assert_eq!(store.stats().total_events, 1);
3964 assert_eq!(store.stats().total_ingested, 1);
3965 }
3966
3967 #[test]
3968 fn test_ingest_multiple_events() {
3969 let store = EventStore::new();
3970
3971 for i in 0..10 {
3972 let event = create_test_event(&format!("entity-{i}"), "user.created");
3973 store.ingest(&event).unwrap();
3974 }
3975
3976 assert_eq!(store.stats().total_events, 10);
3977 assert_eq!(store.stats().total_ingested, 10);
3978 }
3979
3980 #[test]
3981 fn test_query_by_entity_id() {
3982 let store = EventStore::new();
3983
3984 store
3985 .ingest(&create_test_event("entity-1", "user.created"))
3986 .unwrap();
3987 store
3988 .ingest(&create_test_event("entity-2", "user.created"))
3989 .unwrap();
3990 store
3991 .ingest(&create_test_event("entity-1", "user.updated"))
3992 .unwrap();
3993
3994 let results = store
3995 .query(&QueryEventsRequest {
3996 entity_id: Some("entity-1".to_string()),
3997 event_type: None,
3998 tenant_id: None,
3999 as_of: None,
4000 since: None,
4001 until: None,
4002 limit: None,
4003 event_type_prefix: None,
4004 exclude_event_type_prefix: None,
4005 payload_filter: None,
4006 })
4007 .unwrap();
4008
4009 assert_eq!(results.len(), 2);
4010 }
4011
4012 #[test]
4013 fn test_query_by_event_type() {
4014 let store = EventStore::new();
4015
4016 store
4017 .ingest(&create_test_event("entity-1", "user.created"))
4018 .unwrap();
4019 store
4020 .ingest(&create_test_event("entity-2", "user.updated"))
4021 .unwrap();
4022 store
4023 .ingest(&create_test_event("entity-3", "user.created"))
4024 .unwrap();
4025
4026 let results = store
4027 .query(&QueryEventsRequest {
4028 entity_id: None,
4029 event_type: Some("user.created".to_string()),
4030 tenant_id: None,
4031 as_of: None,
4032 since: None,
4033 until: None,
4034 limit: None,
4035 event_type_prefix: None,
4036 exclude_event_type_prefix: None,
4037 payload_filter: None,
4038 })
4039 .unwrap();
4040
4041 assert_eq!(results.len(), 2);
4042 }
4043
4044 #[test]
4045 fn test_query_with_limit() {
4046 let store = EventStore::new();
4047
4048 for i in 0..10 {
4049 let event = create_test_event(&format!("entity-{i}"), "user.created");
4050 store.ingest(&event).unwrap();
4051 }
4052
4053 let results = store
4054 .query(&QueryEventsRequest {
4055 entity_id: None,
4056 event_type: None,
4057 tenant_id: None,
4058 as_of: None,
4059 since: None,
4060 until: None,
4061 limit: Some(5),
4062 event_type_prefix: None,
4063 exclude_event_type_prefix: None,
4064 payload_filter: None,
4065 })
4066 .unwrap();
4067
4068 assert_eq!(results.len(), 5);
4069 }
4070
4071 #[test]
4072 fn test_query_empty_store() {
4073 let store = EventStore::new();
4074
4075 let results = store
4076 .query(&QueryEventsRequest {
4077 entity_id: Some("non-existent".to_string()),
4078 event_type: None,
4079 tenant_id: None,
4080 as_of: None,
4081 since: None,
4082 until: None,
4083 limit: None,
4084 event_type_prefix: None,
4085 exclude_event_type_prefix: None,
4086 payload_filter: None,
4087 })
4088 .unwrap();
4089
4090 assert!(results.is_empty());
4091 }
4092
4093 #[test]
4094 fn test_reconstruct_state() {
4095 let store = EventStore::new();
4096
4097 store
4098 .ingest(&create_test_event("entity-1", "user.created"))
4099 .unwrap();
4100
4101 let state = store.reconstruct_state("entity-1", None).unwrap();
4102 assert_eq!(state["current_state"]["name"], "Test");
4104 assert_eq!(state["current_state"]["value"], 42);
4105 }
4106
4107 #[test]
4108 fn test_reconstruct_state_not_found() {
4109 let store = EventStore::new();
4110
4111 let result = store.reconstruct_state("non-existent", None);
4112 assert!(result.is_err());
4113 }
4114
4115 #[test]
4116 fn test_get_snapshot_empty() {
4117 let store = EventStore::new();
4118
4119 let result = store.get_snapshot("non-existent");
4120 assert!(result.is_err());
4122 }
4123
4124 #[test]
4125 fn test_create_snapshot() {
4126 let store = EventStore::new();
4127
4128 store
4129 .ingest(&create_test_event("entity-1", "user.created"))
4130 .unwrap();
4131
4132 store.create_snapshot("entity-1").unwrap();
4133
4134 let snapshot = store.get_snapshot("entity-1").unwrap();
4136 assert_ne!(snapshot, serde_json::json!(null));
4137 }
4138
4139 #[test]
4140 fn test_create_snapshot_entity_not_found() {
4141 let store = EventStore::new();
4142
4143 let result = store.create_snapshot("non-existent");
4144 assert!(result.is_err());
4145 }
4146
4147 #[test]
4148 fn test_websocket_manager() {
4149 let store = EventStore::new();
4150 let manager = store.websocket_manager();
4151 assert!(Arc::strong_count(&manager) >= 1);
4153 }
4154
4155 #[test]
4156 fn test_snapshot_manager() {
4157 let store = EventStore::new();
4158 let manager = store.snapshot_manager();
4159 assert!(Arc::strong_count(&manager) >= 1);
4160 }
4161
4162 #[test]
4163 fn test_compaction_manager_none() {
4164 let store = EventStore::new();
4165 assert!(store.compaction_manager().is_none());
4167 }
4168
4169 #[test]
4170 fn test_schema_registry() {
4171 let store = EventStore::new();
4172 let registry = store.schema_registry();
4173 assert!(Arc::strong_count(®istry) >= 1);
4174 }
4175
4176 #[test]
4177 fn test_replay_manager() {
4178 let store = EventStore::new();
4179 let manager = store.replay_manager();
4180 assert!(Arc::strong_count(&manager) >= 1);
4181 }
4182
4183 #[test]
4184 fn test_pipeline_manager() {
4185 let store = EventStore::new();
4186 let manager = store.pipeline_manager();
4187 assert!(Arc::strong_count(&manager) >= 1);
4188 }
4189
4190 #[test]
4191 fn test_projection_manager() {
4192 let store = EventStore::new();
4193 let manager = store.projection_manager();
4194 let projections = manager.list_projections();
4196 assert!(projections.len() >= 2); }
4198
4199 #[test]
4200 fn test_projection_state_cache() {
4201 let store = EventStore::new();
4202 let cache = store.projection_state_cache();
4203
4204 cache.insert("test:key".to_string(), serde_json::json!({"value": 123}));
4205 assert_eq!(cache.len(), 1);
4206
4207 let value = cache.get("test:key").unwrap();
4208 assert_eq!(value["value"], 123);
4209 }
4210
4211 #[test]
4212 fn test_metrics() {
4213 let store = EventStore::new();
4214 let metrics = store.metrics();
4215 assert!(Arc::strong_count(&metrics) >= 1);
4216 }
4217
4218 #[test]
4219 fn test_store_stats() {
4220 let store = EventStore::new();
4221
4222 store
4223 .ingest(&create_test_event("entity-1", "user.created"))
4224 .unwrap();
4225 store
4226 .ingest(&create_test_event("entity-2", "order.placed"))
4227 .unwrap();
4228
4229 let stats = store.stats();
4230 assert_eq!(stats.total_events, 2);
4231 assert_eq!(stats.total_entities, 2);
4232 assert_eq!(stats.total_event_types, 2);
4233 assert_eq!(stats.total_ingested, 2);
4234 }
4235
4236 #[test]
4237 fn test_event_store_config_default() {
4238 let config = EventStoreConfig::default();
4239 assert!(config.storage_dir.is_none());
4240 assert!(config.wal_dir.is_none());
4241 }
4242
4243 #[test]
4244 fn test_event_store_config_with_persistence() {
4245 let temp_dir = TempDir::new().unwrap();
4246 let config = EventStoreConfig::with_persistence(temp_dir.path());
4247
4248 assert!(config.storage_dir.is_some());
4249 assert!(config.wal_dir.is_none());
4250 }
4251
4252 #[test]
4253 fn test_event_store_config_with_wal() {
4254 let temp_dir = TempDir::new().unwrap();
4255 let config = EventStoreConfig::with_wal(temp_dir.path(), WALConfig::default());
4256
4257 assert!(config.storage_dir.is_none());
4258 assert!(config.wal_dir.is_some());
4259 }
4260
4261 #[test]
4262 fn test_event_store_config_with_all() {
4263 let temp_dir = TempDir::new().unwrap();
4264 let config = EventStoreConfig::with_all(temp_dir.path(), SnapshotConfig::default());
4265
4266 assert!(config.storage_dir.is_some());
4267 }
4268
4269 #[test]
4270 fn test_event_store_config_production() {
4271 let storage_dir = TempDir::new().unwrap();
4272 let wal_dir = TempDir::new().unwrap();
4273 let config = EventStoreConfig::production(
4274 storage_dir.path(),
4275 wal_dir.path(),
4276 SnapshotConfig::default(),
4277 WALConfig::default(),
4278 CompactionConfig::default(),
4279 );
4280
4281 assert!(config.storage_dir.is_some());
4282 assert!(config.wal_dir.is_some());
4283 }
4284
4285 #[test]
4291 fn test_from_env_vars_data_dir_enables_full_persistence() {
4292 let (config, mode) = EventStoreConfig::from_env_vars(
4293 Some("/app/data".to_string()),
4294 None,
4295 None,
4296 None,
4297 None,
4298 None,
4299 None,
4300 None,
4301 );
4302 assert_eq!(mode, "wal+parquet");
4303 assert_eq!(
4304 config.storage_dir.unwrap().to_str().unwrap(),
4305 "/app/data/storage"
4306 );
4307 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/app/data/wal");
4308 }
4309
4310 #[test]
4311 fn test_from_env_vars_explicit_dirs() {
4312 let (config, mode) = EventStoreConfig::from_env_vars(
4313 None,
4314 Some("/custom/storage".to_string()),
4315 Some("/custom/wal".to_string()),
4316 None,
4317 None,
4318 None,
4319 None,
4320 None,
4321 );
4322 assert_eq!(mode, "wal+parquet");
4323 assert_eq!(
4324 config.storage_dir.unwrap().to_str().unwrap(),
4325 "/custom/storage"
4326 );
4327 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/custom/wal");
4328 }
4329
4330 #[test]
4331 fn test_from_env_vars_wal_disabled() {
4332 let (config, mode) = EventStoreConfig::from_env_vars(
4333 Some("/app/data".to_string()),
4334 None,
4335 None,
4336 Some("false".to_string()),
4337 None,
4338 None,
4339 None,
4340 None,
4341 );
4342 assert_eq!(mode, "parquet-only");
4343 assert!(config.storage_dir.is_some());
4344 assert!(config.wal_dir.is_none());
4345 }
4346
4347 #[test]
4348 fn test_from_env_vars_no_dirs_is_in_memory() {
4349 let (config, mode) =
4350 EventStoreConfig::from_env_vars(None, None, None, None, None, None, None, None);
4351 assert_eq!(mode, "in-memory");
4352 assert!(config.storage_dir.is_none());
4353 assert!(config.wal_dir.is_none());
4354 }
4355
4356 #[test]
4357 fn test_from_env_vars_empty_strings_treated_as_none() {
4358 let (_, mode) = EventStoreConfig::from_env_vars(
4359 Some(String::new()),
4360 Some(String::new()),
4361 Some(String::new()),
4362 None,
4363 None,
4364 None,
4365 None,
4366 None,
4367 );
4368 assert_eq!(mode, "in-memory");
4369 }
4370
4371 #[test]
4372 fn test_from_env_vars_explicit_overrides_data_dir() {
4373 let (config, mode) = EventStoreConfig::from_env_vars(
4374 Some("/app/data".to_string()),
4375 Some("/override/storage".to_string()),
4376 Some("/override/wal".to_string()),
4377 None,
4378 None,
4379 None,
4380 None,
4381 None,
4382 );
4383 assert_eq!(mode, "wal+parquet");
4384 assert_eq!(
4385 config.storage_dir.unwrap().to_str().unwrap(),
4386 "/override/storage"
4387 );
4388 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/override/wal");
4389 }
4390
4391 #[test]
4392 fn test_from_env_vars_wal_only() {
4393 let (config, mode) = EventStoreConfig::from_env_vars(
4394 None,
4395 None,
4396 Some("/wal/only".to_string()),
4397 None,
4398 None,
4399 None,
4400 None,
4401 None,
4402 );
4403 assert_eq!(mode, "wal-only");
4404 assert!(config.storage_dir.is_none());
4405 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/wal/only");
4406 }
4407
4408 #[test]
4409 fn test_from_env_vars_cache_bytes_parses_decimal() {
4410 let (config, _) = EventStoreConfig::from_env_vars(
4411 Some("/app/data".to_string()),
4412 None,
4413 None,
4414 None,
4415 Some("536870912".to_string()),
4416 None,
4418 None,
4419 None,
4420 );
4421 assert_eq!(config.cache_byte_budget, Some(536_870_912));
4422 }
4423
4424 #[test]
4425 fn test_from_env_vars_cache_bytes_unparseable_disables_budget() {
4426 let (config, _) = EventStoreConfig::from_env_vars(
4430 Some("/app/data".to_string()),
4431 None,
4432 None,
4433 None,
4434 Some("not-a-number".to_string()),
4435 None,
4436 None,
4437 None,
4438 );
4439 assert_eq!(config.cache_byte_budget, None);
4440 }
4441
4442 #[test]
4443 fn test_from_env_vars_cache_bytes_empty_disables_budget() {
4444 let (config, _) = EventStoreConfig::from_env_vars(
4445 Some("/app/data".to_string()),
4446 None,
4447 None,
4448 None,
4449 Some(String::new()),
4450 None,
4451 None,
4452 None,
4453 );
4454 assert_eq!(config.cache_byte_budget, None);
4455 }
4456
4457 #[test]
4458 fn test_from_env_vars_snapshot_interval_overrides_default() {
4459 let (config, _) = EventStoreConfig::from_env_vars(
4463 Some("/app/data".to_string()),
4464 None,
4465 None,
4466 None,
4467 None,
4468 Some("60".to_string()),
4469 None,
4470 None,
4471 );
4472 assert_eq!(config.compaction_config.compaction_interval_seconds, 60);
4473 }
4474
4475 #[test]
4476 fn test_from_env_vars_snapshot_interval_default_is_hourly() {
4477 let (config, _) = EventStoreConfig::from_env_vars(
4478 Some("/app/data".to_string()),
4479 None,
4480 None,
4481 None,
4482 None,
4483 None,
4484 None,
4485 None,
4486 );
4487 assert_eq!(config.compaction_config.compaction_interval_seconds, 3600);
4488 }
4489
4490 #[test]
4491 fn test_from_env_vars_snapshot_interval_unparseable_falls_back() {
4492 let (config, _) = EventStoreConfig::from_env_vars(
4493 Some("/app/data".to_string()),
4494 None,
4495 None,
4496 None,
4497 None,
4498 Some("not-a-number".to_string()),
4499 None,
4500 None,
4501 );
4502 assert_eq!(config.compaction_config.compaction_interval_seconds, 3600);
4503 }
4504
4505 #[test]
4506 fn test_from_env_vars_retention_system_days_overrides_default() {
4507 let (config, _) = EventStoreConfig::from_env_vars(
4510 Some("/app/data".to_string()),
4511 None,
4512 None,
4513 None,
4514 None,
4515 None,
4516 Some("7".to_string()),
4517 None,
4518 );
4519 let ttl = config
4520 .compaction_config
4521 .retention
4522 .ttl_for("system")
4523 .unwrap();
4524 assert_eq!(ttl.as_secs(), 7 * 24 * 3600);
4525 }
4526
4527 #[test]
4528 fn test_from_env_vars_retention_default_is_30_days_for_system() {
4529 let (config, _) = EventStoreConfig::from_env_vars(
4530 Some("/app/data".to_string()),
4531 None,
4532 None,
4533 None,
4534 None,
4535 None,
4536 None,
4537 None,
4538 );
4539 let ttl = config
4540 .compaction_config
4541 .retention
4542 .ttl_for("system")
4543 .unwrap();
4544 assert_eq!(ttl.as_secs(), 30 * 24 * 3600);
4545 assert!(config.compaction_config.retention.ttl_for("acme").is_none());
4547 }
4548
4549 #[test]
4550 fn test_store_stats_serde() {
4551 let stats = StoreStats {
4552 total_events: 100,
4553 total_entities: 50,
4554 total_event_types: 10,
4555 total_ingested: 100,
4556 };
4557
4558 let json = serde_json::to_string(&stats).unwrap();
4559 assert!(json.contains("\"total_events\":100"));
4560 assert!(json.contains("\"total_entities\":50"));
4561 }
4562
4563 #[test]
4564 fn test_query_with_entity_and_type() {
4565 let store = EventStore::new();
4566
4567 store
4568 .ingest(&create_test_event("entity-1", "user.created"))
4569 .unwrap();
4570 store
4571 .ingest(&create_test_event("entity-1", "user.updated"))
4572 .unwrap();
4573 store
4574 .ingest(&create_test_event("entity-2", "user.created"))
4575 .unwrap();
4576
4577 let results = store
4578 .query(&QueryEventsRequest {
4579 entity_id: Some("entity-1".to_string()),
4580 event_type: Some("user.created".to_string()),
4581 tenant_id: None,
4582 as_of: None,
4583 since: None,
4584 until: None,
4585 limit: None,
4586 event_type_prefix: None,
4587 exclude_event_type_prefix: None,
4588 payload_filter: None,
4589 })
4590 .unwrap();
4591
4592 assert_eq!(results.len(), 1);
4593 assert_eq!(results[0].event_type_str(), "user.created");
4594 }
4595
4596 #[test]
4597 fn test_query_by_event_type_prefix() {
4598 let store = EventStore::new();
4599
4600 store
4602 .ingest(&create_test_event("entity-1", "index.created"))
4603 .unwrap();
4604 store
4605 .ingest(&create_test_event("entity-2", "index.updated"))
4606 .unwrap();
4607 store
4608 .ingest(&create_test_event("entity-3", "trade.created"))
4609 .unwrap();
4610 store
4611 .ingest(&create_test_event("entity-4", "trade.completed"))
4612 .unwrap();
4613 store
4614 .ingest(&create_test_event("entity-5", "balance.updated"))
4615 .unwrap();
4616
4617 let results = store
4619 .query(&QueryEventsRequest {
4620 entity_id: None,
4621 event_type: None,
4622 tenant_id: None,
4623 as_of: None,
4624 since: None,
4625 until: None,
4626 limit: None,
4627 event_type_prefix: Some("index.".to_string()),
4628 exclude_event_type_prefix: None,
4629 payload_filter: None,
4630 })
4631 .unwrap();
4632
4633 assert_eq!(results.len(), 2);
4634 assert!(
4635 results
4636 .iter()
4637 .all(|e| e.event_type_str().starts_with("index."))
4638 );
4639 }
4640
4641 #[test]
4642 fn test_query_by_event_type_prefix_empty_returns_all() {
4643 let store = EventStore::new();
4644
4645 store
4646 .ingest(&create_test_event("entity-1", "index.created"))
4647 .unwrap();
4648 store
4649 .ingest(&create_test_event("entity-2", "trade.created"))
4650 .unwrap();
4651
4652 let results = store
4654 .query(&QueryEventsRequest {
4655 entity_id: None,
4656 event_type: None,
4657 tenant_id: None,
4658 as_of: None,
4659 since: None,
4660 until: None,
4661 limit: None,
4662 event_type_prefix: Some(String::new()),
4663 exclude_event_type_prefix: None,
4664 payload_filter: None,
4665 })
4666 .unwrap();
4667
4668 assert_eq!(results.len(), 2);
4669 }
4670
4671 #[test]
4672 fn test_query_by_event_type_prefix_no_match() {
4673 let store = EventStore::new();
4674
4675 store
4676 .ingest(&create_test_event("entity-1", "index.created"))
4677 .unwrap();
4678
4679 let results = store
4680 .query(&QueryEventsRequest {
4681 entity_id: None,
4682 event_type: None,
4683 tenant_id: None,
4684 as_of: None,
4685 since: None,
4686 until: None,
4687 limit: None,
4688 event_type_prefix: Some("nonexistent.".to_string()),
4689 exclude_event_type_prefix: None,
4690 payload_filter: None,
4691 })
4692 .unwrap();
4693
4694 assert!(results.is_empty());
4695 }
4696
4697 #[test]
4698 fn test_query_by_entity_with_type_prefix() {
4699 let store = EventStore::new();
4700
4701 store
4702 .ingest(&create_test_event("entity-1", "index.created"))
4703 .unwrap();
4704 store
4705 .ingest(&create_test_event("entity-1", "trade.created"))
4706 .unwrap();
4707 store
4708 .ingest(&create_test_event("entity-2", "index.updated"))
4709 .unwrap();
4710
4711 let results = store
4713 .query(&QueryEventsRequest {
4714 entity_id: Some("entity-1".to_string()),
4715 event_type: None,
4716 tenant_id: None,
4717 as_of: None,
4718 since: None,
4719 until: None,
4720 limit: None,
4721 event_type_prefix: Some("index.".to_string()),
4722 exclude_event_type_prefix: None,
4723 payload_filter: None,
4724 })
4725 .unwrap();
4726
4727 assert_eq!(results.len(), 1);
4728 assert_eq!(results[0].event_type_str(), "index.created");
4729 }
4730
4731 #[test]
4732 fn test_query_prefix_with_limit() {
4733 let store = EventStore::new();
4734
4735 for i in 0..5 {
4736 store
4737 .ingest(&create_test_event(&format!("entity-{i}"), "index.created"))
4738 .unwrap();
4739 }
4740
4741 let results = store
4742 .query(&QueryEventsRequest {
4743 entity_id: None,
4744 event_type: None,
4745 tenant_id: None,
4746 as_of: None,
4747 since: None,
4748 until: None,
4749 limit: Some(3),
4750 event_type_prefix: Some("index.".to_string()),
4751 exclude_event_type_prefix: None,
4752 payload_filter: None,
4753 })
4754 .unwrap();
4755
4756 assert_eq!(results.len(), 3);
4757 }
4758
4759 #[test]
4760 fn test_query_prefix_alongside_existing_filters() {
4761 let store = EventStore::new();
4762
4763 store
4764 .ingest(&create_test_event("entity-1", "index.created"))
4765 .unwrap();
4766 std::thread::sleep(std::time::Duration::from_millis(10));
4768 store
4769 .ingest(&create_test_event("entity-2", "index.strategy.updated"))
4770 .unwrap();
4771 std::thread::sleep(std::time::Duration::from_millis(10));
4772 store
4773 .ingest(&create_test_event("entity-3", "index.deleted"))
4774 .unwrap();
4775
4776 let results = store
4778 .query(&QueryEventsRequest {
4779 entity_id: None,
4780 event_type: None,
4781 tenant_id: None,
4782 as_of: None,
4783 since: None,
4784 until: None,
4785 limit: Some(2),
4786 event_type_prefix: Some("index.".to_string()),
4787 exclude_event_type_prefix: None,
4788 payload_filter: None,
4789 })
4790 .unwrap();
4791
4792 assert_eq!(results.len(), 2);
4793 }
4794
4795 #[test]
4796 fn test_query_with_payload_filter() {
4797 let store = EventStore::new();
4798
4799 for i in 0..5 {
4801 store
4802 .ingest(&create_test_event_with_payload(
4803 &format!("entity-{i}"),
4804 "user.action",
4805 serde_json::json!({"user_id": "alice", "action": "click"}),
4806 ))
4807 .unwrap();
4808 }
4809 for i in 5..10 {
4811 store
4812 .ingest(&create_test_event_with_payload(
4813 &format!("entity-{i}"),
4814 "user.action",
4815 serde_json::json!({"user_id": "bob", "action": "view"}),
4816 ))
4817 .unwrap();
4818 }
4819
4820 let results = store
4822 .query(&QueryEventsRequest {
4823 entity_id: None,
4824 event_type: Some("user.action".to_string()),
4825 tenant_id: None,
4826 as_of: None,
4827 since: None,
4828 until: None,
4829 limit: None,
4830 event_type_prefix: None,
4831 exclude_event_type_prefix: None,
4832 payload_filter: Some(r#"{"user_id":"alice"}"#.to_string()),
4833 })
4834 .unwrap();
4835
4836 assert_eq!(results.len(), 5);
4837 }
4838
4839 #[test]
4840 fn test_query_payload_filter_non_existent_field() {
4841 let store = EventStore::new();
4842
4843 store
4844 .ingest(&create_test_event_with_payload(
4845 "entity-1",
4846 "user.action",
4847 serde_json::json!({"user_id": "alice"}),
4848 ))
4849 .unwrap();
4850
4851 let results = store
4853 .query(&QueryEventsRequest {
4854 entity_id: None,
4855 event_type: None,
4856 tenant_id: None,
4857 as_of: None,
4858 since: None,
4859 until: None,
4860 limit: None,
4861 event_type_prefix: None,
4862 exclude_event_type_prefix: None,
4863 payload_filter: Some(r#"{"nonexistent":"value"}"#.to_string()),
4864 })
4865 .unwrap();
4866
4867 assert!(results.is_empty());
4868 }
4869
4870 #[test]
4871 fn test_query_payload_filter_with_prefix() {
4872 let store = EventStore::new();
4873
4874 store
4875 .ingest(&create_test_event_with_payload(
4876 "entity-1",
4877 "index.created",
4878 serde_json::json!({"status": "active"}),
4879 ))
4880 .unwrap();
4881 store
4882 .ingest(&create_test_event_with_payload(
4883 "entity-2",
4884 "index.created",
4885 serde_json::json!({"status": "inactive"}),
4886 ))
4887 .unwrap();
4888 store
4889 .ingest(&create_test_event_with_payload(
4890 "entity-3",
4891 "trade.created",
4892 serde_json::json!({"status": "active"}),
4893 ))
4894 .unwrap();
4895
4896 let results = store
4898 .query(&QueryEventsRequest {
4899 entity_id: None,
4900 event_type: None,
4901 tenant_id: None,
4902 as_of: None,
4903 since: None,
4904 until: None,
4905 limit: None,
4906 event_type_prefix: Some("index.".to_string()),
4907 exclude_event_type_prefix: None,
4908 payload_filter: Some(r#"{"status":"active"}"#.to_string()),
4909 })
4910 .unwrap();
4911
4912 assert_eq!(results.len(), 1);
4913 assert_eq!(results[0].entity_id().to_string(), "entity-1");
4914 }
4915
4916 #[test]
4917 fn test_flush_storage_no_storage() {
4918 let store = EventStore::new();
4919 let result = store.flush_storage();
4921 assert!(result.is_ok());
4922 }
4923
4924 #[test]
4925 fn test_state_evolution() {
4926 let store = EventStore::new();
4927
4928 store
4930 .ingest(
4931 &Event::from_strings(
4932 "user.created".to_string(),
4933 "user-1".to_string(),
4934 "default".to_string(),
4935 serde_json::json!({"name": "Alice", "age": 25}),
4936 None,
4937 )
4938 .unwrap(),
4939 )
4940 .unwrap();
4941
4942 store
4944 .ingest(
4945 &Event::from_strings(
4946 "user.updated".to_string(),
4947 "user-1".to_string(),
4948 "default".to_string(),
4949 serde_json::json!({"age": 26}),
4950 None,
4951 )
4952 .unwrap(),
4953 )
4954 .unwrap();
4955
4956 let state = store.reconstruct_state("user-1", None).unwrap();
4957 assert_eq!(state["current_state"]["name"], "Alice");
4959 assert_eq!(state["current_state"]["age"], 26);
4960 }
4961
4962 #[test]
4963 fn test_reject_system_event_types() {
4964 let store = EventStore::new();
4965
4966 let event = Event::reconstruct_from_strings(
4968 uuid::Uuid::new_v4(),
4969 "_system.tenant.created".to_string(),
4970 "_system:tenant:acme".to_string(),
4971 "_system".to_string(),
4972 serde_json::json!({"name": "ACME"}),
4973 chrono::Utc::now(),
4974 None,
4975 1,
4976 );
4977
4978 let result = store.ingest(&event);
4979 assert!(result.is_err());
4980 let err = result.unwrap_err();
4981 assert!(
4982 err.to_string().contains("reserved for internal use"),
4983 "Expected system namespace rejection, got: {err}"
4984 );
4985 }
4986
4987 #[test]
4995 fn test_wal_recovery_checkpoints_to_parquet() {
4996 let data_dir = TempDir::new().unwrap();
4997 let storage_dir = data_dir.path().join("storage");
4998 let wal_dir = data_dir.path().join("wal");
4999
5000 {
5002 let config = EventStoreConfig::production(
5003 &storage_dir,
5004 &wal_dir,
5005 SnapshotConfig::default(),
5006 WALConfig {
5007 sync_on_write: true,
5008 ..WALConfig::default()
5009 },
5010 CompactionConfig::default(),
5011 );
5012 let store = EventStore::with_config(config);
5013
5014 for i in 0..5 {
5015 let event = Event::from_strings(
5016 "test.created".to_string(),
5017 format!("entity-{i}"),
5018 "default".to_string(),
5019 serde_json::json!({"index": i}),
5020 None,
5021 )
5022 .unwrap();
5023 store.ingest(&event).unwrap();
5024 }
5025
5026 assert_eq!(store.stats().total_events, 5);
5027
5028 }
5031
5032 let wal_files: Vec<_> = std::fs::read_dir(&wal_dir)
5034 .unwrap()
5035 .filter_map(std::result::Result::ok)
5036 .filter(|e| e.path().extension().is_some_and(|ext| ext == "log"))
5037 .collect();
5038 assert!(!wal_files.is_empty(), "WAL file should exist");
5039 let wal_size = wal_files[0].metadata().unwrap().len();
5040 assert!(wal_size > 0, "WAL file should have data (got 0 bytes)");
5041
5042 {
5044 let config = EventStoreConfig::production(
5045 &storage_dir,
5046 &wal_dir,
5047 SnapshotConfig::default(),
5048 WALConfig {
5049 sync_on_write: true,
5050 ..WALConfig::default()
5051 },
5052 CompactionConfig::default(),
5053 );
5054 let store = EventStore::with_config(config);
5055
5056 assert_eq!(
5058 store.stats().total_events,
5059 5,
5060 "Session 2 should have all 5 events after WAL recovery"
5061 );
5062
5063 let parquet_files = find_parquet_files(&storage_dir);
5067 assert!(
5068 !parquet_files.is_empty(),
5069 "Parquet file should exist after WAL checkpoint"
5070 );
5071 }
5072
5073 {
5076 let config = EventStoreConfig::production(
5077 &storage_dir,
5078 &wal_dir,
5079 SnapshotConfig::default(),
5080 WALConfig {
5081 sync_on_write: true,
5082 ..WALConfig::default()
5083 },
5084 CompactionConfig::default(),
5085 );
5086 let store = EventStore::with_config(config);
5087
5088 assert_eq!(
5092 store.stats().total_events,
5093 0,
5094 "Session 3 boot should not pre-load Parquet (lazy-load mode)"
5095 );
5096
5097 store.ensure_tenant_loaded("default").unwrap();
5100 assert_eq!(
5101 store.stats().total_events,
5102 5,
5103 "Session 3 should have all 5 events after ensure_tenant_loaded"
5104 );
5105 }
5106 }
5107
5108 #[test]
5109 fn test_parquet_restore_surfaces_errors_not_silent() {
5110 let data_dir = TempDir::new().unwrap();
5114 let storage_dir = data_dir.path().join("storage");
5115 let wal_dir = data_dir.path().join("wal");
5116
5117 {
5119 let config = EventStoreConfig::production(
5120 &storage_dir,
5121 &wal_dir,
5122 SnapshotConfig::default(),
5123 WALConfig {
5124 sync_on_write: true,
5125 ..WALConfig::default()
5126 },
5127 CompactionConfig::default(),
5128 );
5129 let store = EventStore::with_config(config);
5130
5131 for i in 0..3 {
5132 let event = Event::from_strings(
5133 "test.created".to_string(),
5134 format!("entity-{i}"),
5135 "default".to_string(),
5136 serde_json::json!({"i": i}),
5137 None,
5138 )
5139 .unwrap();
5140 store.ingest(&event).unwrap();
5141 }
5142
5143 store.flush_storage().unwrap();
5144 assert_eq!(store.stats().total_events, 3);
5145 }
5146
5147 let parquet_files = find_parquet_files(&storage_dir);
5150 assert!(!parquet_files.is_empty(), "Parquet file must exist");
5151
5152 std::fs::write(&parquet_files[0], b"corrupted data").unwrap();
5154
5155 for entry in std::fs::read_dir(&wal_dir).unwrap().flatten() {
5157 std::fs::write(entry.path(), b"").unwrap();
5158 }
5159
5160 {
5167 let config = EventStoreConfig::production(
5168 &storage_dir,
5169 &wal_dir,
5170 SnapshotConfig::default(),
5171 WALConfig::default(),
5172 CompactionConfig::default(),
5173 );
5174 let store = EventStore::with_config(config);
5175
5176 assert_eq!(store.stats().total_events, 0);
5179 }
5180 }
5181
5182 fn count_wal_entries(wal_dir: &std::path::Path) -> usize {
5192 use std::io::{BufRead, BufReader};
5193 let mut total = 0usize;
5194 let Ok(entries) = std::fs::read_dir(wal_dir) else {
5195 return 0;
5196 };
5197 for entry in entries.flatten() {
5198 let path = entry.path();
5199 if path.extension().is_none_or(|e| e != "log") {
5200 continue;
5201 }
5202 let Ok(file) = std::fs::File::open(&path) else {
5203 continue;
5204 };
5205 for line in BufReader::new(file)
5206 .lines()
5207 .map_while(std::result::Result::ok)
5208 {
5209 if !line.trim().is_empty() {
5210 total += 1;
5211 }
5212 }
5213 }
5214 total
5215 }
5216
5217 #[test]
5218 fn test_checkpoint_truncates_wal_after_flush() {
5219 let data_dir = TempDir::new().unwrap();
5224 let storage_dir = data_dir.path().join("storage");
5225 let wal_dir = data_dir.path().join("wal");
5226
5227 let config = EventStoreConfig::production(
5228 &storage_dir,
5229 &wal_dir,
5230 SnapshotConfig::default(),
5231 WALConfig {
5232 sync_on_write: true,
5233 ..WALConfig::default()
5234 },
5235 CompactionConfig::default(),
5236 );
5237 let store = EventStore::with_config(config);
5238
5239 for i in 0..10 {
5240 let event = Event::from_strings(
5241 "test.created".to_string(),
5242 format!("entity-{i}"),
5243 "default".to_string(),
5244 serde_json::json!({"i": i}),
5245 None,
5246 )
5247 .unwrap();
5248 store.ingest(&event).unwrap();
5249 }
5250
5251 assert_eq!(
5253 count_wal_entries(&wal_dir),
5254 10,
5255 "WAL should have 10 events before checkpoint"
5256 );
5257
5258 store.checkpoint().unwrap();
5259
5260 assert_eq!(
5261 count_wal_entries(&wal_dir),
5262 0,
5263 "WAL should be empty after successful checkpoint"
5264 );
5265 let parquet_files = find_parquet_files(&storage_dir);
5266 assert!(!parquet_files.is_empty(), "Parquet should hold the events");
5267 }
5268
5269 #[test]
5270 fn test_replay_only_post_checkpoint_events_after_crash() {
5271 let data_dir = TempDir::new().unwrap();
5278 let storage_dir = data_dir.path().join("storage");
5279 let wal_dir = data_dir.path().join("wal");
5280
5281 let config_factory = || {
5282 EventStoreConfig::production(
5283 &storage_dir,
5284 &wal_dir,
5285 SnapshotConfig::default(),
5286 WALConfig {
5287 sync_on_write: true,
5288 ..WALConfig::default()
5289 },
5290 CompactionConfig::default(),
5291 )
5292 };
5293
5294 const N: usize = 50;
5297 const K: usize = 5;
5298 {
5299 let store = EventStore::with_config(config_factory());
5300 for i in 0..N {
5301 store
5302 .ingest(
5303 &Event::from_strings(
5304 "pre.checkpoint".to_string(),
5305 format!("e-{i}"),
5306 "default".to_string(),
5307 serde_json::json!({"i": i}),
5308 None,
5309 )
5310 .unwrap(),
5311 )
5312 .unwrap();
5313 }
5314 store.checkpoint().unwrap();
5315 assert_eq!(
5316 count_wal_entries(&wal_dir),
5317 0,
5318 "WAL should be empty immediately after checkpoint"
5319 );
5320
5321 for i in 0..K {
5322 store
5323 .ingest(
5324 &Event::from_strings(
5325 "post.checkpoint".to_string(),
5326 format!("p-{i}"),
5327 "default".to_string(),
5328 serde_json::json!({"i": i}),
5329 None,
5330 )
5331 .unwrap(),
5332 )
5333 .unwrap();
5334 }
5335 assert_eq!(
5336 count_wal_entries(&wal_dir),
5337 K,
5338 "WAL should hold only post-checkpoint events"
5339 );
5340 }
5342
5343 {
5347 let store = EventStore::with_config(config_factory());
5348 assert_eq!(
5352 store.stats().total_events,
5353 K,
5354 "Boot should replay exactly K events from WAL (the post-checkpoint window), not N+K"
5355 );
5356
5357 store.ensure_tenant_loaded("default").unwrap();
5359 assert_eq!(
5360 store.stats().total_events,
5361 N + K,
5362 "After lazy-load, both pre- and post-checkpoint events should be reachable"
5363 );
5364 }
5365 }
5366
5367 #[test]
5368 fn test_checkpoint_is_idempotent() {
5369 let data_dir = TempDir::new().unwrap();
5372 let storage_dir = data_dir.path().join("storage");
5373 let wal_dir = data_dir.path().join("wal");
5374
5375 let store = EventStore::with_config(EventStoreConfig::production(
5376 &storage_dir,
5377 &wal_dir,
5378 SnapshotConfig::default(),
5379 WALConfig::default(),
5380 CompactionConfig::default(),
5381 ));
5382
5383 for i in 0..5 {
5384 store
5385 .ingest(
5386 &Event::from_strings(
5387 "x".to_string(),
5388 format!("e-{i}"),
5389 "default".to_string(),
5390 serde_json::json!({}),
5391 None,
5392 )
5393 .unwrap(),
5394 )
5395 .unwrap();
5396 }
5397
5398 store.checkpoint().unwrap();
5399 store.checkpoint().unwrap();
5401 assert_eq!(count_wal_entries(&wal_dir), 0);
5402 }
5403
5404 #[test]
5405 fn test_checkpoint_noop_in_memory_only_mode() {
5406 let store = EventStore::new();
5408 store.checkpoint().unwrap();
5409 }
5410
5411 #[test]
5412 fn test_checkpoint_interval_from_env_defaults_to_60s_when_wal_enabled() {
5413 let (config, _) = EventStoreConfig::from_env_vars(
5414 Some("/app/data".to_string()),
5415 None,
5416 None,
5417 None,
5418 None,
5419 None,
5420 None,
5421 None,
5422 );
5423 assert_eq!(config.checkpoint_interval_secs, Some(60));
5424 }
5425
5426 #[test]
5427 fn test_checkpoint_interval_from_env_overrides_default() {
5428 let (config, _) = EventStoreConfig::from_env_vars(
5429 Some("/app/data".to_string()),
5430 None,
5431 None,
5432 None,
5433 None,
5434 None,
5435 None,
5436 Some("15".to_string()),
5437 );
5438 assert_eq!(config.checkpoint_interval_secs, Some(15));
5439 }
5440
5441 #[test]
5442 fn test_checkpoint_interval_disabled_when_wal_disabled() {
5443 let (config, _) = EventStoreConfig::from_env_vars(
5445 Some("/app/data".to_string()),
5446 None,
5447 None,
5448 Some("false".to_string()),
5449 None,
5450 None,
5451 None,
5452 Some("15".to_string()),
5453 );
5454 assert_eq!(config.checkpoint_interval_secs, None);
5455 }
5456
5457 #[test]
5458 fn test_checkpoint_interval_unparseable_falls_back_to_default() {
5459 let (config, _) = EventStoreConfig::from_env_vars(
5460 Some("/app/data".to_string()),
5461 None,
5462 None,
5463 None,
5464 None,
5465 None,
5466 None,
5467 Some("not-a-number".to_string()),
5468 );
5469 assert_eq!(config.checkpoint_interval_secs, Some(60));
5470 }
5471
5472 fn seed_two_tenants() -> EventStore {
5479 let store = EventStore::new();
5480
5481 for (entity, etype, payload) in [
5483 (
5484 "a-1",
5485 "created",
5486 serde_json::json!({"colour": "red", "size": 1}),
5487 ),
5488 ("a-1", "updated", serde_json::json!({"colour": "blue"})),
5489 ("a-2", "created", serde_json::json!({"colour": "green"})),
5490 ] {
5491 store
5492 .ingest(
5493 &Event::from_strings(
5494 etype.to_string(),
5495 entity.to_string(),
5496 "alice".to_string(),
5497 payload,
5498 None,
5499 )
5500 .unwrap(),
5501 )
5502 .unwrap();
5503 }
5504
5505 store
5507 .ingest(
5508 &Event::from_strings(
5509 "created".to_string(),
5510 "a-1".to_string(),
5511 "bob".to_string(),
5512 serde_json::json!({"colour": "BOB_SECRET", "bob_only": true}),
5513 None,
5514 )
5515 .unwrap(),
5516 )
5517 .unwrap();
5518
5519 store
5520 }
5521
5522 #[test]
5523 fn test_stats_for_tenant_counts_only_that_tenant() {
5524 let store = seed_two_tenants();
5525
5526 let alice = store.stats_for_tenant("alice");
5527 assert_eq!(alice.total_events, 3);
5528 assert_eq!(alice.total_entities, 2);
5529 assert_eq!(alice.total_event_types, 2);
5530 assert_eq!(alice.event_types.get("created"), Some(&2));
5531 assert_eq!(alice.event_types.get("updated"), Some(&1));
5532
5533 let bob = store.stats_for_tenant("bob");
5534 assert_eq!(bob.total_events, 1);
5535 assert_eq!(bob.total_entities, 1);
5536 assert_eq!(bob.total_event_types, 1);
5537 }
5538
5539 #[test]
5540 fn test_stats_for_tenant_never_reports_global_totals() {
5541 let store = seed_two_tenants();
5542
5543 assert_eq!(store.stats().total_events, 4);
5545 assert_eq!(store.stats_for_tenant("alice").total_events, 3);
5546 assert_eq!(store.stats_for_tenant("bob").total_events, 1);
5547
5548 assert_eq!(store.stats_for_tenant("bob").total_ingested, 1);
5551 }
5552
5553 #[test]
5554 fn test_stats_for_tenant_unknown_tenant_is_empty_not_global() {
5555 let store = seed_two_tenants();
5556 let nobody = store.stats_for_tenant("does-not-exist");
5557
5558 assert_eq!(nobody.total_events, 0);
5559 assert_eq!(nobody.total_entities, 0);
5560 assert!(nobody.event_types.is_empty());
5561 assert!(nobody.oldest_event.is_none());
5562 assert!(nobody.newest_event.is_none());
5563 }
5564
5565 #[test]
5566 fn test_stats_for_tenant_reports_time_range() {
5567 let store = seed_two_tenants();
5568 let alice = store.stats_for_tenant("alice");
5569
5570 let oldest = alice.oldest_event.expect("oldest");
5571 let newest = alice.newest_event.expect("newest");
5572 assert!(oldest <= newest);
5573 }
5574
5575 #[test]
5576 fn test_reconstruct_state_for_tenant_isolates_shared_entity_id() {
5577 let store = seed_two_tenants();
5578
5579 let alice = store
5581 .reconstruct_state_for_tenant("a-1", None, "alice")
5582 .unwrap();
5583 let alice_state = alice.get("current_state").unwrap();
5584 assert_eq!(alice_state.get("colour").unwrap(), "blue"); assert_eq!(alice_state.get("size").unwrap(), 1);
5586 assert!(
5587 alice_state.get("bob_only").is_none(),
5588 "alice must not see bob's payload keys: {alice_state:?}"
5589 );
5590 assert_eq!(alice.get("event_count").unwrap(), 2);
5591
5592 let bob = store
5593 .reconstruct_state_for_tenant("a-1", None, "bob")
5594 .unwrap();
5595 let bob_state = bob.get("current_state").unwrap();
5596 assert_eq!(bob_state.get("colour").unwrap(), "BOB_SECRET");
5597 assert_eq!(bob.get("event_count").unwrap(), 1);
5598 }
5599
5600 #[test]
5601 fn test_reconstruct_state_for_tenant_rejects_another_tenants_entity() {
5602 let store = seed_two_tenants();
5603
5604 assert!(
5606 store
5607 .reconstruct_state_for_tenant("a-2", None, "alice")
5608 .is_ok()
5609 );
5610 assert!(
5611 store
5612 .reconstruct_state_for_tenant("a-2", None, "bob")
5613 .is_err(),
5614 "bob must not be able to read alice's entity"
5615 );
5616 }
5617
5618 #[test]
5619 fn test_global_reconstruct_state_still_spans_tenants() {
5620 let store = seed_two_tenants();
5623 let all = store.reconstruct_state("a-1", None).unwrap();
5624 assert_eq!(all.get("event_count").unwrap(), 3);
5625 }
5626
5627 fn seed_hot_entity(history: usize) -> EventStore {
5638 let store = EventStore::new();
5639 for _ in 0..history {
5640 store
5641 .ingest(&create_test_event("entity-hot", "user.updated"))
5642 .unwrap();
5643 }
5644 store
5645 }
5646
5647 fn hot_request(limit: Option<usize>) -> QueryEventsRequest {
5648 QueryEventsRequest {
5649 entity_id: Some("entity-hot".to_string()),
5650 limit,
5651 ..QueryEventsRequest::default()
5652 }
5653 }
5654
5655 #[test]
5656 fn query_window_limit_bounds_materialization_not_just_the_response() {
5657 const HISTORY: usize = 400;
5658 let store = seed_hot_entity(HISTORY);
5659
5660 let ((page, total), materialized) = crate::clone_probe::measure(|| {
5662 store.query_window(&hot_request(Some(1)), 0, true).unwrap()
5663 });
5664 assert_eq!(page.len(), 1, "limit=1 returns one event");
5665 assert_eq!(total, HISTORY, "total is still the full match count");
5666 assert_eq!(
5667 materialized, 1,
5668 "limit=1 cloned {materialized} of {HISTORY} events: `limit` must \
5669 bound what the store materializes, not just what it returns"
5670 );
5671
5672 let ((page, _), materialized) = crate::clone_probe::measure(|| {
5674 store
5675 .query_window(&hot_request(Some(5)), 300, false)
5676 .unwrap()
5677 });
5678 assert_eq!(page.len(), 5);
5679 assert_eq!(
5680 materialized, 5,
5681 "offset=300&limit=5 cloned {materialized} events, expected 5"
5682 );
5683
5684 let ((all, _), materialized) = crate::clone_probe::measure(|| {
5687 store.query_window(&hot_request(None), 0, true).unwrap()
5688 });
5689 assert_eq!(all.len(), HISTORY);
5690 assert_eq!(materialized, HISTORY as u64);
5691 }
5692
5693 #[test]
5694 fn query_window_cost_does_not_grow_with_entity_history() {
5695 let short = seed_hot_entity(40);
5700 let long = seed_hot_entity(400);
5701
5702 let (_, short_cost) = crate::clone_probe::measure(|| {
5703 short.query_window(&hot_request(Some(1)), 0, true).unwrap()
5704 });
5705 let (_, long_cost) = crate::clone_probe::measure(|| {
5706 long.query_window(&hot_request(Some(1)), 0, true).unwrap()
5707 });
5708
5709 assert_eq!(
5710 (short_cost, long_cost),
5711 (1, 1),
5712 "a limit=1 page materialized {short_cost} events over a 40-event \
5713 history and {long_cost} over 400 — cost is tracking history length"
5714 );
5715 }
5716
5717 #[test]
5718 fn select_window_orders_only_the_window() {
5719 use std::cell::Cell;
5725 const N: usize = 4096;
5726
5727 let comparisons = Cell::new(0usize);
5728 let order = |a: &u64, b: &u64| {
5729 comparisons.set(comparisons.get() + 1);
5730 a.cmp(b)
5731 };
5732 let shuffled = || -> Vec<u64> {
5733 (0..N as u64)
5734 .map(|i| (i * 2_654_435_761) % 1_000_003)
5735 .collect()
5736 };
5737
5738 let mut bounded = shuffled();
5739 select_window(&mut bounded, 0, Some(1), order);
5740 let bounded_cost = comparisons.replace(0);
5741 assert_eq!(bounded.len(), 1, "window of 1 keeps 1 item");
5742
5743 let mut everything = shuffled();
5744 select_window(&mut everything, 0, None, order);
5745 let full_sort_cost = comparisons.get();
5746 assert_eq!(everything.len(), N);
5747
5748 assert!(
5749 bounded_cost < 4 * N,
5750 "selecting a 1-item window out of {N} took {bounded_cost} comparisons \
5751 (~{}·N) — that is sort-shaped, not selection-shaped",
5752 bounded_cost / N
5753 );
5754 assert!(
5755 bounded_cost * 3 < full_sort_cost,
5756 "a 1-item window cost {bounded_cost} comparisons against \
5757 {full_sort_cost} for sorting all {N}: the window is not bounding \
5758 the ordering work"
5759 );
5760 }
5761
5762 #[test]
5763 fn select_window_is_equivalent_to_sort_then_window() {
5764 let order = |a: &(u32, usize), b: &(u32, usize)| a.cmp(b);
5768 let source: Vec<(u32, usize)> = [7, 3, 3, 9, 1, 3, 5, 9, 0, 2]
5769 .into_iter()
5770 .enumerate()
5771 .map(|(i, k)| (k, i))
5772 .collect();
5773
5774 let mut sorted = source.clone();
5775 sorted.sort_by(order);
5776
5777 for offset in 0..12 {
5778 for limit in [None, Some(0), Some(1), Some(3), Some(10), Some(50)] {
5779 let mut got = source.clone();
5780 select_window(&mut got, offset, limit, order);
5781 let got: Vec<_> = got.into_iter().skip(offset).collect();
5782 let expected: Vec<_> = sorted
5783 .iter()
5784 .copied()
5785 .skip(offset)
5786 .take(limit.unwrap_or(usize::MAX))
5787 .collect();
5788 assert_eq!(got, expected, "offset={offset} limit={limit:?}");
5789 }
5790 }
5791 }
5792
5793 #[test]
5804 fn query_window_applies_time_filters_on_the_full_scan_path() {
5805 let store = EventStore::new();
5806 let base = Utc::now() - chrono::Duration::hours(24);
5807 let mut ids = Vec::new();
5808 for i in 0..5i64 {
5809 let mut event = create_test_event(&format!("e-{i}"), "user.created");
5810 event.timestamp = base + chrono::Duration::hours(i);
5811 event.version = i + 1;
5812 ids.push(event.id);
5813 store.ingest(&event).unwrap();
5814 }
5815 let at = |h: i64| base + chrono::Duration::hours(h);
5816
5817 let scoped = |mutate: &dyn Fn(&mut QueryEventsRequest)| {
5820 let mut req = QueryEventsRequest {
5821 tenant_id: Some("default".to_string()),
5822 ..QueryEventsRequest::default()
5823 };
5824 mutate(&mut req);
5825 req
5826 };
5827
5828 for (label, req, expected) in [
5829 (
5830 "since=T+2 keeps only events at or after T+2",
5831 scoped(&|r| r.since = Some(at(2))),
5832 vec![ids[2], ids[3], ids[4]],
5833 ),
5834 (
5835 "until=T+1 keeps only events at or before T+1",
5836 scoped(&|r| r.until = Some(at(1))),
5837 vec![ids[0], ids[1]],
5838 ),
5839 (
5840 "as_of=T+1 is time travel: nothing newer than T+1",
5841 scoped(&|r| r.as_of = Some(at(1))),
5842 vec![ids[0], ids[1]],
5843 ),
5844 (
5845 "since+until compose into a closed window",
5846 scoped(&|r| {
5847 r.since = Some(at(1));
5848 r.until = Some(at(3));
5849 }),
5850 vec![ids[1], ids[2], ids[3]],
5851 ),
5852 ] {
5853 let (events, total) = store.query_window(&req, 0, false).unwrap();
5854 let got: Vec<_> = events.iter().map(|e| e.id).collect();
5855 assert_eq!(got, expected, "{label}");
5856 assert_eq!(
5857 total,
5858 expected.len(),
5859 "{label}: total counts the events INSIDE the window — \
5860 has_more is derived from it, so a full-history total makes a \
5861 paginator walk events the caller filtered out"
5862 );
5863 }
5864
5865 let (indexed, total) = store
5868 .query_window(
5869 &scoped(&|r| {
5870 r.entity_id = Some("e-3".to_string());
5871 r.since = Some(at(2));
5872 }),
5873 0,
5874 false,
5875 )
5876 .unwrap();
5877 assert_eq!(
5878 indexed.iter().map(|e| e.id).collect::<Vec<_>>(),
5879 vec![ids[3]]
5880 );
5881 assert_eq!(total, 1);
5882
5883 let (page, total) = store
5886 .query_window(&scoped(&|r| r.since = Some(at(2))), 1, true)
5887 .unwrap();
5888 assert_eq!(
5889 page.iter().map(|e| e.id).collect::<Vec<_>>(),
5890 vec![ids[3], ids[2]],
5891 "order=desc + offset=1 inside a since window"
5892 );
5893 assert_eq!(total, 3);
5894 }
5895
5896 #[test]
5897 fn query_window_bounded_selection_matches_a_full_sort() {
5898 let store = EventStore::new();
5903 for i in 0..50 {
5904 store
5905 .ingest(&create_test_event(&format!("e-{i:02}"), "user.created"))
5906 .unwrap();
5907 }
5908 let all = |descending: bool| {
5909 let (events, _) = store
5910 .query_window(&QueryEventsRequest::default(), 0, descending)
5911 .unwrap();
5912 events
5913 };
5914
5915 for descending in [false, true] {
5916 let reference = all(descending);
5917 for offset in [0, 1, 7, 49, 50, 100] {
5918 for limit in [1, 3, 10, 50, 100] {
5919 let (page, total) = store
5920 .query_window(
5921 &QueryEventsRequest {
5922 limit: Some(limit),
5923 ..QueryEventsRequest::default()
5924 },
5925 offset,
5926 descending,
5927 )
5928 .unwrap();
5929 let expected: Vec<_> = reference
5930 .iter()
5931 .skip(offset)
5932 .take(limit)
5933 .map(|e| e.id)
5934 .collect();
5935 let got: Vec<_> = page.iter().map(|e| e.id).collect();
5936 assert_eq!(
5937 got, expected,
5938 "desc={descending} offset={offset} limit={limit}: windowed \
5939 selection must match a full sort"
5940 );
5941 assert_eq!(total, 50, "total is always the full match count");
5942 }
5943 }
5944 }
5945 }
5946}