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
41#[derive(Debug, Clone, Default)]
58pub struct ReadScope {
59 allow_entity_prefixes: Option<Vec<String>>,
60}
61
62impl ReadScope {
63 pub fn unrestricted() -> Self {
65 Self {
66 allow_entity_prefixes: None,
67 }
68 }
69
70 pub fn allow_entity_prefixes<I, S>(prefixes: I) -> Self
75 where
76 I: IntoIterator<Item = S>,
77 S: Into<String>,
78 {
79 Self {
80 allow_entity_prefixes: Some(prefixes.into_iter().map(Into::into).collect()),
81 }
82 }
83
84 pub fn permits(&self, entity_id: &str) -> bool {
86 match &self.allow_entity_prefixes {
87 None => true,
88 Some(allowed) => allowed.iter().any(|p| entity_id.starts_with(p.as_str())),
89 }
90 }
91
92 pub fn is_unrestricted(&self) -> bool {
94 self.allow_entity_prefixes.is_none()
95 }
96}
97
98pub struct EventStore {
99 events: Arc<RwLock<Vec<Event>>>,
101
102 index: Arc<EventIndex>,
104
105 pub(crate) projections: Arc<RwLock<ProjectionManager>>,
107
108 storage: Option<Arc<RwLock<ParquetStorage>>>,
110
111 #[cfg(feature = "server")]
113 websocket_manager: Arc<WebSocketManager>,
114
115 snapshot_manager: Arc<SnapshotManager>,
117
118 wal: Option<Arc<WriteAheadLog>>,
120
121 compaction_manager: Option<Arc<CompactionManager>>,
123
124 schema_registry: Arc<SchemaRegistry>,
126
127 replay_manager: Arc<ReplayManager>,
129
130 pipeline_manager: Arc<PipelineManager>,
132
133 #[cfg(feature = "server")]
135 metrics: Arc<MetricsRegistry>,
136
137 total_ingested: Arc<RwLock<u64>>,
139
140 projection_state_cache: Arc<DashMap<String, serde_json::Value>>,
144
145 projection_status: Arc<DashMap<String, String>>,
148
149 #[cfg(feature = "server")]
151 webhook_registry: Arc<WebhookRegistry>,
152
153 #[cfg(feature = "server")]
155 webhook_tx: Arc<RwLock<Option<mpsc::UnboundedSender<WebhookDeliveryTask>>>>,
156
157 geo_index: Arc<GeoIndex>,
159
160 exactly_once: Arc<ExactlyOnceRegistry>,
162
163 schema_evolution: Arc<SchemaEvolutionManager>,
165
166 entity_versions: Arc<DashMap<String, u64>>,
169
170 consumer_registry: Arc<ConsumerRegistry>,
172
173 event_broadcast_tx: tokio::sync::broadcast::Sender<Arc<Event>>,
177
178 tenant_loader: Arc<TenantLoader>,
184
185 checkpoint_interval_secs: Option<u64>,
191
192 read_only: bool,
199
200 durability_gate: RwLock<()>,
204}
205
206#[cfg(feature = "server")]
208#[derive(Debug, Clone)]
209pub struct WebhookDeliveryTask {
210 pub webhook: crate::application::services::webhook::WebhookSubscription,
211 pub event: Event,
212}
213
214fn select_window<T>(
227 items: &mut Vec<T>,
228 offset: usize,
229 limit: Option<usize>,
230 order: impl Fn(&T, &T) -> std::cmp::Ordering + Copy,
231) {
232 match limit.map(|limit| offset.saturating_add(limit)) {
233 Some(0) => {
234 items.clear();
235 return;
236 }
237 Some(window_end) if window_end < items.len() => {
238 items.select_nth_unstable_by(window_end - 1, order);
239 items.truncate(window_end);
240 }
241 _ => {}
244 }
245 items.sort_unstable_by(order);
246}
247
248impl EventStore {
249 pub fn new() -> Self {
251 Self::with_config(EventStoreConfig::default())
252 }
253
254 pub fn with_config(config: EventStoreConfig) -> Self {
256 let mut projections = ProjectionManager::new();
257
258 projections.register(Arc::new(EntitySnapshotProjection::new("entity_snapshots")));
260 projections.register(Arc::new(EventCounterProjection::new("event_counters")));
261
262 let storage = config
264 .storage_dir
265 .as_ref()
266 .and_then(|dir| match ParquetStorage::new(dir) {
267 Ok(storage) => {
268 tracing::info!("✅ Parquet persistence enabled at: {}", dir.display());
269 Some(Arc::new(RwLock::new(storage)))
270 }
271 Err(e) => {
272 tracing::error!("❌ Failed to initialize Parquet storage: {}", e);
273 None
274 }
275 });
276
277 let wal = config.wal_dir.as_ref().and_then(|dir| {
279 match WriteAheadLog::new(dir, config.wal_config.clone()) {
280 Ok(wal) => {
281 tracing::info!("✅ WAL enabled at: {}", dir.display());
282 Some(Arc::new(wal))
283 }
284 Err(e) => {
285 tracing::error!("❌ Failed to initialize WAL: {}", e);
286 None
287 }
288 }
289 });
290
291 let compaction_manager = config.storage_dir.as_ref().map(|dir| {
293 let manager = CompactionManager::new(dir, config.compaction_config.clone());
294 Arc::new(manager)
295 });
296
297 let schema_registry = Arc::new(SchemaRegistry::new(config.schema_registry_config.clone()));
299 tracing::info!("✅ Schema registry enabled");
300
301 let replay_manager = Arc::new(ReplayManager::new());
303 tracing::info!("✅ Replay manager enabled");
304
305 let pipeline_manager = Arc::new(PipelineManager::new());
307 tracing::info!("✅ Pipeline manager enabled");
308
309 #[cfg(feature = "server")]
311 let metrics = {
312 let m = MetricsRegistry::new();
313 tracing::info!("✅ Prometheus metrics registry initialized");
314 m
315 };
316
317 let projection_state_cache = Arc::new(DashMap::new());
319 tracing::info!("✅ Projection state cache initialized");
320
321 #[cfg(feature = "server")]
323 let webhook_registry = {
324 let w = Arc::new(WebhookRegistry::new());
325 tracing::info!("✅ Webhook registry initialized");
326 w
327 };
328
329 let (event_broadcast_tx, _) = tokio::sync::broadcast::channel(1024);
332
333 let store = Self {
334 events: Arc::new(RwLock::new(Vec::new())),
335 index: Arc::new(EventIndex::new()),
336 projections: Arc::new(RwLock::new(projections)),
337 storage,
338 #[cfg(feature = "server")]
339 websocket_manager: Arc::new(WebSocketManager::new()),
340 snapshot_manager: Arc::new(SnapshotManager::new(config.snapshot_config)),
341 wal,
342 compaction_manager,
343 schema_registry,
344 replay_manager,
345 pipeline_manager,
346 #[cfg(feature = "server")]
347 metrics,
348 total_ingested: Arc::new(RwLock::new(0)),
349 projection_state_cache,
350 projection_status: Arc::new(DashMap::new()),
351 #[cfg(feature = "server")]
352 webhook_registry,
353 #[cfg(feature = "server")]
354 webhook_tx: Arc::new(RwLock::new(None)),
355 geo_index: Arc::new(GeoIndex::new()),
356 exactly_once: Arc::new(ExactlyOnceRegistry::new(ExactlyOnceConfig::default())),
357 schema_evolution: Arc::new(SchemaEvolutionManager::new()),
358 entity_versions: Arc::new(DashMap::new()),
359 consumer_registry: Arc::new(ConsumerRegistry::new()),
360 event_broadcast_tx,
361 tenant_loader: {
362 let loader = TenantLoader::new();
363 if let Some(budget) = config.cache_byte_budget {
364 loader.set_byte_budget(budget);
365 tracing::info!(
366 "✅ Cache byte budget set to {} bytes ({:.2} GiB) — LRU eviction enabled",
367 budget,
368 budget as f64 / (1024.0 * 1024.0 * 1024.0)
369 );
370 } else {
371 tracing::info!(
372 "✅ Cache budget unset — every loaded tenant stays resident \
373 (set ALLSOURCE_CACHE_BYTES to enable eviction)"
374 );
375 }
376 Arc::new(loader)
377 },
378 checkpoint_interval_secs: config.checkpoint_interval_secs,
379 read_only: config.read_only,
380 durability_gate: RwLock::new(()),
381 };
382
383 if config.read_only {
384 tracing::info!(
385 "📖 EventStore opened READ-ONLY (replica): WAL will be replayed for reads but \
386 not truncated; writes are rejected"
387 );
388 }
389
390 if let Some(ref wal) = store.wal {
412 match wal.recover() {
413 Ok(recovered_events) if !recovered_events.is_empty() => {
414 let mut wal_new = 0usize;
415 for event in recovered_events {
416 let offset = store.events.read().len();
417 if let Err(e) = store.index.index_event(
418 event.id,
419 event.entity_id_str(),
420 event.event_type_str(),
421 event.timestamp,
422 offset,
423 ) {
424 tracing::error!("Failed to re-index WAL event {}: {}", event.id, e);
425 }
426
427 if let Err(e) = store.projections.read().process_event(&event) {
428 tracing::error!("Failed to re-process WAL event {}: {}", event.id, e);
429 }
430
431 *store
432 .entity_versions
433 .entry(event.entity_id_str().to_string())
434 .or_insert(0) += 1;
435
436 store.events.write().push(event);
437 wal_new += 1;
438 }
439
440 #[cfg(feature = "server")]
445 store.metrics.wal_replay_events_total.set(wal_new as i64);
446
447 if wal_new > 0 {
448 let total = store.events.read().len();
449 *store.total_ingested.write() = total as u64;
457 tracing::info!(
458 "✅ Recovered {} events from WAL (Parquet data stays cold until \
459 first per-tenant query)",
460 wal_new
461 );
462
463 if let Some(ref storage) = store.storage
478 && !config.read_only
479 {
480 tracing::info!(
481 "📸 Checkpointing {} WAL events to Parquet storage...",
482 wal_new
483 );
484 let parquet = storage.read();
485 let events = store.events.read();
486 let mut buffered = 0usize;
487 for event in events.iter().skip(events.len() - wal_new) {
488 if let Err(e) = parquet.append_event(event.clone()) {
489 tracing::error!(
490 "Failed to buffer WAL event for Parquet: {}",
491 e
492 );
493 } else {
494 buffered += 1;
495 }
496 }
497 drop(events);
498 drop(parquet);
499
500 if buffered < wal_new {
501 tracing::error!(
502 "Buffered {} of {} recovered WAL events; leaving the WAL \
503 in place so none are lost",
504 buffered,
505 wal_new
506 );
507 } else if buffered > 0 {
508 if let Err(e) = store.flush_storage() {
509 tracing::error!("Failed to checkpoint to Parquet: {}", e);
510 } else if let Err(e) = wal.truncate() {
511 tracing::error!(
512 "Failed to truncate WAL after checkpoint: {}",
513 e
514 );
515 } else {
516 tracing::info!(
517 "✅ WAL checkpointed and truncated ({} events)",
518 buffered
519 );
520 }
521 }
522 }
523 }
524 }
525 Ok(_) => {
526 tracing::debug!("No events to recover from WAL");
527 #[cfg(feature = "server")]
528 store.metrics.wal_replay_events_total.set(0);
529 }
530 Err(e) => {
531 tracing::error!("❌ WAL recovery failed: {}", e);
532 }
533 }
534 } else if store.storage.is_some() {
535 tracing::info!(
536 "📂 Boot complete (lazy-load mode): Parquet data stays on disk until first \
537 per-tenant query"
538 );
539 }
540
541 store
542 }
543
544 pub fn is_read_only(&self) -> bool {
546 self.read_only
547 }
548
549 fn ensure_writable(&self) -> Result<()> {
553 if self.read_only {
554 return Err(crate::error::AllSourceError::ReadOnly(
555 "this AllSource instance is a read-only replica — the data directory is owned by \
556 another running process. Stop the other process, or run a single shared writer \
557 (e.g. Prime in --mode http) and point clients at it."
558 .to_string(),
559 ));
560 }
561 Ok(())
562 }
563
564 #[cfg_attr(feature = "hotpath", hotpath::measure)]
572 pub fn ingest_with_expected_version(
573 &self,
574 event: &Event,
575 expected_version: Option<u64>,
576 ) -> Result<u64> {
577 self.ensure_writable()?;
579
580 self.validate_event(event)?;
582
583 let entity_id = event.entity_id_str().to_string();
584 let _durable = self.durability_gate.read();
585
586 let new_version = {
589 let mut version_entry = self.entity_versions.entry(entity_id.clone()).or_insert(0);
590 let current = *version_entry;
591
592 if let Some(expected) = expected_version
593 && current != expected
594 {
595 return Err(crate::error::AllSourceError::VersionConflict { expected, current });
596 }
597
598 if let Some(ref wal) = self.wal {
600 wal.append(event.clone())?;
601 }
602
603 *version_entry += 1;
604 *version_entry
605 };
606
607 self.ingest_post_wal(event)?;
610
611 Ok(new_version)
612 }
613
614 #[cfg_attr(feature = "hotpath", hotpath::measure)]
617 fn ingest_post_wal(&self, event: &Event) -> Result<()> {
618 #[cfg(feature = "server")]
619 let timer = self.metrics.ingestion_duration_seconds.start_timer();
620
621 let mut events = self.events.write();
622 let offset = events.len();
623
624 self.index.index_event(
626 event.id,
627 event.entity_id_str(),
628 event.event_type_str(),
629 event.timestamp,
630 offset,
631 )?;
632
633 let projections = self.projections.read();
635 projections.process_event(event)?;
636 drop(projections);
637
638 let pipeline_results = self.pipeline_manager.process_event(event);
640 if !pipeline_results.is_empty() {
641 tracing::debug!(
642 "Event {} processed by {} pipeline(s)",
643 event.id,
644 pipeline_results.len()
645 );
646 for (pipeline_id, result) in pipeline_results {
647 tracing::trace!("Pipeline {} result: {:?}", pipeline_id, result);
648 }
649 }
650
651 if let Some(ref storage) = self.storage {
653 let storage = storage.read();
654 storage.append_event(event.clone())?;
655 }
656
657 events.push(event.clone());
659 let total_events = events.len();
660 drop(events);
661
662 let event_arc = Arc::new(event.clone());
664 let _ = self.event_broadcast_tx.send(Arc::clone(&event_arc));
665 #[cfg(feature = "server")]
666 self.websocket_manager.broadcast_event(event_arc);
667
668 #[cfg(feature = "server")]
670 self.dispatch_webhooks(event);
671
672 self.geo_index.index_event(event);
674
675 self.schema_evolution
677 .analyze_event(event.event_type_str(), &event.payload);
678
679 self.check_auto_snapshot(event.entity_id_str(), event);
681
682 #[cfg(feature = "server")]
684 {
685 self.metrics.events_ingested_total.inc();
686 self.metrics
687 .events_ingested_by_type
688 .with_label_values(&[event.event_type_str()])
689 .inc();
690 self.metrics.storage_events_total.set(total_events as i64);
691 }
692
693 let mut total = self.total_ingested.write();
695 *total += 1;
696
697 #[cfg(feature = "server")]
698 timer.observe_duration();
699
700 tracing::debug!("Event ingested: {} (offset: {})", event.id, offset);
701
702 Ok(())
703 }
704
705 #[cfg_attr(feature = "hotpath", hotpath::measure)]
707 pub fn ingest(&self, event: &Event) -> Result<()> {
708 #[cfg(feature = "server")]
710 let timer = self.metrics.ingestion_duration_seconds.start_timer();
711
712 if let Err(e) = self.ensure_writable() {
714 #[cfg(feature = "server")]
715 {
716 self.metrics.ingestion_errors_total.inc();
717 timer.observe_duration();
718 }
719 return Err(e);
720 }
721
722 let validation_result = self.validate_event(event);
724 if let Err(e) = validation_result {
725 #[cfg(feature = "server")]
726 {
727 self.metrics.ingestion_errors_total.inc();
728 timer.observe_duration();
729 }
730 return Err(e);
731 }
732
733 let _durable = self.durability_gate.read();
734
735 if let Some(ref wal) = self.wal
738 && let Err(e) = wal.append(event.clone())
739 {
740 #[cfg(feature = "server")]
741 {
742 self.metrics.ingestion_errors_total.inc();
743 timer.observe_duration();
744 }
745 return Err(e);
746 }
747
748 *self
750 .entity_versions
751 .entry(event.entity_id_str().to_string())
752 .or_insert(0) += 1;
753
754 let mut events = self.events.write();
755 let offset = events.len();
756
757 self.index.index_event(
759 event.id,
760 event.entity_id_str(),
761 event.event_type_str(),
762 event.timestamp,
763 offset,
764 )?;
765
766 let projections = self.projections.read();
768 projections.process_event(event)?;
769 drop(projections); let pipeline_results = self.pipeline_manager.process_event(event);
774 if !pipeline_results.is_empty() {
775 tracing::debug!(
776 "Event {} processed by {} pipeline(s)",
777 event.id,
778 pipeline_results.len()
779 );
780 for (pipeline_id, result) in pipeline_results {
783 tracing::trace!("Pipeline {} result: {:?}", pipeline_id, result);
784 }
785 }
786
787 if let Some(ref storage) = self.storage {
789 let storage = storage.read();
790 storage.append_event(event.clone())?;
791 }
792
793 events.push(event.clone());
795 let total_events = events.len();
796 drop(events); let event_arc = Arc::new(event.clone());
800 let _ = self.event_broadcast_tx.send(Arc::clone(&event_arc));
801 #[cfg(feature = "server")]
802 self.websocket_manager.broadcast_event(event_arc);
803
804 #[cfg(feature = "server")]
806 self.dispatch_webhooks(event);
807
808 self.geo_index.index_event(event);
810
811 self.schema_evolution
813 .analyze_event(event.event_type_str(), &event.payload);
814
815 self.check_auto_snapshot(event.entity_id_str(), event);
817
818 #[cfg(feature = "server")]
820 {
821 self.metrics.events_ingested_total.inc();
822 self.metrics
823 .events_ingested_by_type
824 .with_label_values(&[event.event_type_str()])
825 .inc();
826 self.metrics.storage_events_total.set(total_events as i64);
827 }
828
829 let mut total = self.total_ingested.write();
831 *total += 1;
832
833 #[cfg(feature = "server")]
834 timer.observe_duration();
835
836 tracing::debug!("Event ingested: {} (offset: {})", event.id, offset);
837
838 Ok(())
839 }
840
841 #[cfg_attr(feature = "hotpath", hotpath::measure)]
848 pub fn ingest_batch(&self, batch: Vec<Event>) -> Result<()> {
849 if batch.is_empty() {
850 return Ok(());
851 }
852
853 self.ensure_writable()?;
855
856 for event in &batch {
858 self.validate_event(event)?;
859 }
860
861 let _durable = self.durability_gate.read();
862
863 if let Some(ref wal) = self.wal {
865 for event in &batch {
866 wal.append(event.clone())?;
867 }
868 }
869
870 let mut events = self.events.write();
872 let projections = self.projections.read();
873
874 for event in batch {
875 let offset = events.len();
876
877 self.index.index_event(
878 event.id,
879 event.entity_id_str(),
880 event.event_type_str(),
881 event.timestamp,
882 offset,
883 )?;
884
885 projections.process_event(&event)?;
886 self.pipeline_manager.process_event(&event);
887
888 if let Some(ref storage) = self.storage {
889 let storage = storage.read();
890 storage.append_event(event.clone())?;
891 }
892
893 self.geo_index.index_event(&event);
894 self.schema_evolution
895 .analyze_event(event.event_type_str(), &event.payload);
896
897 *self
899 .entity_versions
900 .entry(event.entity_id_str().to_string())
901 .or_insert(0) += 1;
902
903 let _ = self.event_broadcast_tx.send(Arc::new(event.clone()));
905
906 events.push(event);
907 }
908
909 let total_events = events.len();
910 drop(projections);
911 drop(events);
912
913 let mut total = self.total_ingested.write();
914 *total += total_events as u64;
915
916 Ok(())
917 }
918
919 #[cfg_attr(feature = "hotpath", hotpath::measure)]
926 pub fn ingest_replicated(&self, event: &Event) -> Result<()> {
927 #[cfg(feature = "server")]
928 let timer = self.metrics.ingestion_duration_seconds.start_timer();
929
930 let mut events = self.events.write();
931 let offset = events.len();
932
933 self.index.index_event(
935 event.id,
936 event.entity_id_str(),
937 event.event_type_str(),
938 event.timestamp,
939 offset,
940 )?;
941
942 let projections = self.projections.read();
944 projections.process_event(event)?;
945 drop(projections);
946
947 let pipeline_results = self.pipeline_manager.process_event(event);
949 if !pipeline_results.is_empty() {
950 tracing::debug!(
951 "Replicated event {} processed by {} pipeline(s)",
952 event.id,
953 pipeline_results.len()
954 );
955 }
956
957 *self
959 .entity_versions
960 .entry(event.entity_id_str().to_string())
961 .or_insert(0) += 1;
962
963 events.push(event.clone());
965 let total_events = events.len();
966 drop(events);
967
968 let event_arc = Arc::new(event.clone());
970 let _ = self.event_broadcast_tx.send(Arc::clone(&event_arc));
971 #[cfg(feature = "server")]
972 self.websocket_manager.broadcast_event(event_arc);
973
974 #[cfg(feature = "server")]
976 {
977 self.metrics.events_ingested_total.inc();
978 self.metrics
979 .events_ingested_by_type
980 .with_label_values(&[event.event_type_str()])
981 .inc();
982 self.metrics.storage_events_total.set(total_events as i64);
983 }
984
985 let mut total = self.total_ingested.write();
986 *total += 1;
987
988 #[cfg(feature = "server")]
989 timer.observe_duration();
990
991 tracing::debug!(
992 "Replicated event ingested: {} (offset: {})",
993 event.id,
994 offset
995 );
996
997 Ok(())
998 }
999
1000 #[cfg_attr(feature = "hotpath", hotpath::measure)]
1003 pub fn get_entity_version(&self, entity_id: &str) -> u64 {
1004 self.entity_versions.get(entity_id).map_or(0, |v| *v)
1005 }
1006
1007 pub fn consumer_registry(&self) -> &ConsumerRegistry {
1009 &self.consumer_registry
1010 }
1011
1012 pub fn subscribe_events(&self) -> tokio::sync::broadcast::Receiver<Arc<Event>> {
1018 self.event_broadcast_tx.subscribe()
1019 }
1020
1021 pub fn set_consumer_registry(&mut self, registry: Arc<ConsumerRegistry>) {
1026 self.consumer_registry = registry;
1027 }
1028
1029 pub fn total_events(&self) -> usize {
1031 self.events.read().len()
1032 }
1033
1034 pub fn events_after_offset(
1037 &self,
1038 offset: u64,
1039 filters: &[String],
1040 limit: usize,
1041 ) -> Vec<(u64, Event)> {
1042 let events = self.events.read();
1043 let start = offset as usize;
1044 if start >= events.len() {
1045 return vec![];
1046 }
1047
1048 events[start..]
1049 .iter()
1050 .enumerate()
1051 .filter(|(_, event)| ConsumerRegistry::matches_filters(event.event_type_str(), filters))
1052 .take(limit)
1053 .map(|(i, event)| ((start + i + 1) as u64, event.clone()))
1054 .collect()
1055 }
1056
1057 #[cfg(feature = "server")]
1059 pub fn websocket_manager(&self) -> Arc<WebSocketManager> {
1060 Arc::clone(&self.websocket_manager)
1061 }
1062
1063 pub fn snapshot_manager(&self) -> Arc<SnapshotManager> {
1065 Arc::clone(&self.snapshot_manager)
1066 }
1067
1068 pub fn compaction_manager(&self) -> Option<Arc<CompactionManager>> {
1070 self.compaction_manager.as_ref().map(Arc::clone)
1071 }
1072
1073 pub fn schema_registry(&self) -> Arc<SchemaRegistry> {
1075 Arc::clone(&self.schema_registry)
1076 }
1077
1078 pub fn replay_manager(&self) -> Arc<ReplayManager> {
1080 Arc::clone(&self.replay_manager)
1081 }
1082
1083 pub fn pipeline_manager(&self) -> Arc<PipelineManager> {
1085 Arc::clone(&self.pipeline_manager)
1086 }
1087
1088 #[cfg(feature = "server")]
1090 pub fn metrics(&self) -> Arc<MetricsRegistry> {
1091 Arc::clone(&self.metrics)
1092 }
1093
1094 pub fn projection_manager(&self) -> parking_lot::RwLockReadGuard<'_, ProjectionManager> {
1096 self.projections.read()
1097 }
1098
1099 pub fn register_projection(
1108 &self,
1109 projection: Arc<dyn crate::application::services::projection::Projection>,
1110 ) {
1111 let mut pm = self.projections.write();
1112 pm.register(projection);
1113 }
1114
1115 pub fn register_projection_with_backfill(
1127 &self,
1128 projection: &Arc<dyn crate::application::services::projection::Projection>,
1129 ) -> Result<()> {
1130 {
1132 let mut pm = self.projections.write();
1133 pm.register(Arc::clone(projection));
1134 }
1135
1136 let events = self.events.read();
1138 let mut ordered: Vec<&Event> = events.iter().collect();
1139 ordered.sort_by(|a, b| {
1140 a.timestamp
1141 .cmp(&b.timestamp)
1142 .then_with(|| a.version.cmp(&b.version))
1143 });
1144 for event in ordered {
1145 projection.process(event)?;
1146 }
1147
1148 Ok(())
1149 }
1150
1151 pub fn hydrate_all_from_storage(&self) -> Result<usize> {
1169 let Some(storage) = self.storage.as_ref().map(Arc::clone) else {
1170 return Ok(0);
1171 };
1172
1173 let events = storage.read().load_all_events()?;
1174 let read_count = events.len();
1175 let before = self.events.read().len();
1176 for event in events {
1177 self.append_loaded_event(event);
1178 }
1179 let applied = self.events.read().len() - before;
1180
1181 *self.total_ingested.write() = self.events.read().len() as u64;
1183
1184 tracing::info!(
1185 read = read_count,
1186 applied = applied,
1187 "🔄 hydrate_all_from_storage: in-memory pile reconstructed from Parquet"
1188 );
1189 Ok(applied)
1190 }
1191
1192 pub fn projection_state_cache(&self) -> Arc<DashMap<String, serde_json::Value>> {
1195 Arc::clone(&self.projection_state_cache)
1196 }
1197
1198 pub fn projection_status(&self) -> Arc<DashMap<String, String>> {
1200 Arc::clone(&self.projection_status)
1201 }
1202
1203 pub fn geo_index(&self) -> Arc<GeoIndex> {
1206 self.geo_index.clone()
1207 }
1208
1209 pub fn exactly_once(&self) -> Arc<ExactlyOnceRegistry> {
1211 self.exactly_once.clone()
1212 }
1213
1214 pub fn schema_evolution(&self) -> Arc<SchemaEvolutionManager> {
1216 self.schema_evolution.clone()
1217 }
1218
1219 pub fn snapshot_events(&self) -> Vec<Event> {
1225 self.events.read().clone()
1226 }
1227
1228 pub fn compact_entity_tokens(
1250 &self,
1251 entity_id: &str,
1252 token_event_type: &str,
1253 merged_event: Event,
1254 ) -> Result<bool> {
1255 self.ensure_writable()?;
1257
1258 {
1260 let events = self.events.read();
1261 let has_tokens = events
1262 .iter()
1263 .any(|e| e.entity_id_str() == entity_id && e.event_type_str() == token_event_type);
1264 if !has_tokens {
1265 return Ok(false);
1266 }
1267 }
1268
1269 let projections = self.projections.read();
1271 projections.process_event(&merged_event)?;
1272 drop(projections);
1273
1274 let mut events = self.events.write();
1276
1277 events.retain(|e| {
1278 !(e.entity_id_str() == entity_id && e.event_type_str() == token_event_type)
1279 });
1280
1281 events.push(merged_event);
1282
1283 self.index.clear();
1288 for (offset, event) in events.iter().enumerate() {
1289 if let Err(e) = self.index.index_event(
1290 event.id,
1291 event.entity_id_str(),
1292 event.event_type_str(),
1293 event.timestamp,
1294 offset,
1295 ) {
1296 tracing::warn!(
1297 event_id = %event.id,
1298 offset,
1299 "Failed to re-index event during compaction: {e}"
1300 );
1301 }
1302 }
1303
1304 Ok(true)
1305 }
1306
1307 #[cfg(feature = "server")]
1308 pub fn webhook_registry(&self) -> Arc<WebhookRegistry> {
1309 Arc::clone(&self.webhook_registry)
1310 }
1311
1312 #[cfg(feature = "server")]
1315 pub fn set_webhook_tx(&self, tx: mpsc::UnboundedSender<WebhookDeliveryTask>) {
1316 *self.webhook_tx.write() = Some(tx);
1317 tracing::info!("Webhook delivery channel connected");
1318 }
1319
1320 #[cfg(feature = "server")]
1322 fn dispatch_webhooks(&self, event: &Event) {
1323 let matching = self.webhook_registry.find_matching(event);
1324 if matching.is_empty() {
1325 return;
1326 }
1327
1328 let tx_guard = self.webhook_tx.read();
1329 if let Some(ref tx) = *tx_guard {
1330 for webhook in matching {
1331 let task = WebhookDeliveryTask {
1332 webhook,
1333 event: event.clone(),
1334 };
1335 if let Err(e) = tx.send(task) {
1336 tracing::warn!("Failed to queue webhook delivery: {}", e);
1337 }
1338 }
1339 }
1340 }
1341
1342 pub fn flush_storage(&self) -> Result<()> {
1344 if let Some(ref storage) = self.storage {
1345 let storage = storage.read();
1346 storage.flush()?;
1347 tracing::info!("✅ Flushed events to persistent storage");
1348 }
1349 Ok(())
1350 }
1351
1352 pub fn checkpoint(&self) -> Result<()> {
1373 let Some(ref wal) = self.wal else {
1374 #[cfg(feature = "server")]
1377 self.refresh_storage_metrics();
1378 return Ok(());
1379 };
1380
1381 if self.read_only {
1382 return Ok(());
1383 }
1384
1385 let active = {
1389 let _sealing = self.durability_gate.write();
1390 wal.seal()?
1391 };
1392 self.flush_storage()?;
1393 wal.remove_sealed(&active)?;
1394 tracing::debug!("✅ Checkpoint complete: Parquet flushed, sealed WAL segments retired");
1395
1396 #[cfg(feature = "server")]
1399 self.refresh_storage_metrics();
1400
1401 Ok(())
1402 }
1403
1404 #[cfg(feature = "server")]
1409 pub fn refresh_storage_metrics_now(&self) {
1410 self.refresh_storage_metrics();
1411 }
1412
1413 #[cfg(feature = "server")]
1430 fn refresh_storage_metrics(&self) {
1431 let Some(ref storage) = self.storage else {
1432 return;
1433 };
1434
1435 let parquet_stats = match storage.read().stats() {
1436 Ok(stats) => stats,
1437 Err(e) => {
1438 tracing::warn!("storage-size metric refresh: failed to stat Parquet: {e}");
1439 return;
1440 }
1441 };
1442
1443 let (wal_bytes, wal_segments) = match self.wal.as_ref() {
1444 Some(wal) => match wal.on_disk_stats() {
1445 Ok(stats) => stats,
1446 Err(e) => {
1447 tracing::warn!("storage-size metric refresh: failed to stat WAL: {e}");
1448 (0, 0)
1449 }
1450 },
1451 None => (0, 0),
1452 };
1453
1454 let total_bytes = parquet_stats.total_size_bytes + wal_bytes;
1455
1456 self.metrics
1457 .storage_size_bytes
1458 .set(total_bytes.min(i64::MAX as u64) as i64);
1459 self.metrics
1460 .parquet_files_total
1461 .set(parquet_stats.total_files as i64);
1462 self.metrics.wal_segments_total.set(wal_segments as i64);
1463
1464 tracing::debug!(
1465 "storage-size metrics refreshed: {} bytes total ({} Parquet files, {} WAL segments)",
1466 total_bytes,
1467 parquet_stats.total_files,
1468 wal_segments
1469 );
1470 }
1471
1472 pub fn checkpoint_interval(&self) -> Option<std::time::Duration> {
1474 self.checkpoint_interval_secs
1475 .map(std::time::Duration::from_secs)
1476 }
1477
1478 pub fn ensure_tenant_loaded(&self, tenant_id: &str) -> Result<()> {
1502 if self.tenant_loader.is_loaded(tenant_id) {
1504 return Ok(());
1505 }
1506
1507 let Some(storage) = self.storage.as_ref().map(Arc::clone) else {
1508 self.tenant_loader.mark_loaded(tenant_id);
1511 return Ok(());
1512 };
1513
1514 let lock = self.tenant_loader.lock_for(tenant_id);
1517 let timeout = self.tenant_loader.load_timeout();
1518 let _guard = lock.try_lock_for(timeout).ok_or_else(|| {
1519 AllSourceError::StorageError(format!(
1520 "ensure_tenant_loaded timed out after {timeout:?} waiting for in-flight load of \
1521 tenant {tenant_id:?}"
1522 ))
1523 })?;
1524
1525 if self.tenant_loader.is_loaded(tenant_id) {
1528 return Ok(());
1529 }
1530
1531 let started = std::time::Instant::now();
1532 let events = storage.read().load_events_for_tenant(tenant_id)?;
1533 let read_count = events.len();
1534
1535 let before = self.events.read().len();
1536 for event in events {
1537 self.append_loaded_event(event);
1538 }
1539 let applied = self.events.read().len() - before;
1540
1541 *self.total_ingested.write() += applied as u64;
1545 self.tenant_loader.mark_loaded(tenant_id);
1546
1547 tracing::info!(
1548 tenant_id = tenant_id,
1549 read = read_count,
1550 applied = applied,
1551 elapsed_ms = started.elapsed().as_millis() as u64,
1552 "ensure_tenant_loaded: tenant hydrated"
1553 );
1554
1555 self.enforce_cache_budget(tenant_id);
1563
1564 #[cfg(feature = "server")]
1567 self.metrics
1568 .cache_bytes
1569 .set(self.tenant_loader.total_bytes() as i64);
1570
1571 Ok(())
1572 }
1573
1574 fn enforce_cache_budget(&self, recently_touched: &str) {
1585 if !self.tenant_loader.over_budget() {
1586 return;
1587 }
1588 loop {
1589 let Some(victim) = self.tenant_loader.pick_lru_excluding(recently_touched) else {
1590 tracing::warn!(
1591 cache_bytes = self.tenant_loader.total_bytes(),
1592 budget = self.tenant_loader.byte_budget(),
1593 recently_touched = recently_touched,
1594 "cache over budget but no other tenant available to evict — \
1595 a single tenant exceeds the budget; consider raising it"
1596 );
1597 return;
1598 };
1599 self.evict_tenant(&victim);
1600 if !self.tenant_loader.over_budget() {
1601 return;
1602 }
1603 }
1604 }
1605
1606 pub fn is_tenant_loaded(&self, tenant_id: &str) -> bool {
1609 self.tenant_loader.is_loaded(tenant_id)
1610 }
1611
1612 pub fn evict_tenant(&self, tenant_id: &str) {
1640 let mut events = self.events.write();
1641 let before = events.len();
1642 let evicted_bytes = self.tenant_loader.bytes_for(tenant_id);
1643
1644 events.retain(|e| e.tenant_id_str() != tenant_id);
1645 let after = events.len();
1646 let dropped = before - after;
1647
1648 if dropped == 0 {
1649 drop(events);
1653 self.tenant_loader.mark_unloaded(tenant_id);
1654 return;
1655 }
1656
1657 self.index.clear();
1661 self.entity_versions.clear();
1662 for (offset, event) in events.iter().enumerate() {
1663 if let Err(e) = self.index.index_event(
1664 event.id,
1665 event.entity_id_str(),
1666 event.event_type_str(),
1667 event.timestamp,
1668 offset,
1669 ) {
1670 tracing::error!(
1671 "Failed to re-index event during eviction of {}: {}",
1672 tenant_id,
1673 e
1674 );
1675 }
1676 *self
1677 .entity_versions
1678 .entry(event.entity_id_str().to_string())
1679 .or_insert(0) += 1;
1680 }
1681 drop(events);
1682
1683 self.tenant_loader.mark_unloaded(tenant_id);
1684
1685 let mut t = self.total_ingested.write();
1688 *t = t.saturating_sub(dropped as u64);
1689 drop(t);
1690
1691 #[cfg(feature = "server")]
1694 {
1695 self.metrics.cache_evictions_total.inc();
1696 self.metrics
1697 .cache_bytes
1698 .set(self.tenant_loader.total_bytes() as i64);
1699 }
1700
1701 tracing::info!(
1702 tenant_id = tenant_id,
1703 events_dropped = dropped,
1704 bytes_freed = evicted_bytes,
1705 "evicted tenant from memory cache"
1706 );
1707 }
1708
1709 pub fn tenant_resident_bytes(&self, tenant_id: &str) -> u64 {
1713 self.tenant_loader.bytes_for(tenant_id)
1714 }
1715
1716 pub fn cache_resident_bytes(&self) -> u64 {
1719 self.tenant_loader.total_bytes()
1720 }
1721
1722 fn append_loaded_event(&self, event: Event) {
1744 if self.index.get_by_id(&event.id).is_some() {
1745 return;
1746 }
1747
1748 let event_bytes = event.estimated_size_bytes();
1749 let tenant = event.tenant_id_str().to_string();
1750
1751 let mut events = self.events.write();
1752 let offset = events.len();
1753
1754 if let Err(e) = self.index.index_event(
1755 event.id,
1756 event.entity_id_str(),
1757 event.event_type_str(),
1758 event.timestamp,
1759 offset,
1760 ) {
1761 tracing::error!("Failed to index loaded event {}: {}", event.id, e);
1762 }
1763
1764 if let Err(e) = self.projections.read().process_event(&event) {
1765 tracing::error!("Failed to project loaded event {}: {}", event.id, e);
1766 }
1767
1768 *self
1769 .entity_versions
1770 .entry(event.entity_id_str().to_string())
1771 .or_insert(0) += 1;
1772
1773 events.push(event);
1774 self.tenant_loader.add_bytes(&tenant, event_bytes);
1778 }
1779
1780 pub fn create_snapshot(&self, entity_id: &str) -> Result<()> {
1782 let events = self.query(&QueryEventsRequest {
1784 entity_id: Some(entity_id.to_string()),
1785 event_type: None,
1786 tenant_id: None,
1787 as_of: None,
1788 since: None,
1789 until: None,
1790 limit: None,
1791 event_type_prefix: None,
1792 exclude_event_type_prefix: None,
1793 payload_filter: None,
1794 })?;
1795
1796 if events.is_empty() {
1797 return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
1798 }
1799
1800 let mut state = serde_json::json!({});
1802 for event in &events {
1803 if let serde_json::Value::Object(ref mut state_map) = state
1804 && let serde_json::Value::Object(ref payload_map) = event.payload
1805 {
1806 for (key, value) in payload_map {
1807 state_map.insert(key.clone(), value.clone());
1808 }
1809 }
1810 }
1811
1812 let last_event = events.last().unwrap();
1813 self.snapshot_manager.create_snapshot(
1814 entity_id,
1815 state,
1816 last_event.timestamp,
1817 events.len(),
1818 SnapshotType::Manual,
1819 )?;
1820
1821 Ok(())
1822 }
1823
1824 fn check_auto_snapshot(&self, entity_id: &str, event: &Event) {
1826 let entity_event_count = self
1828 .index
1829 .get_by_entity(entity_id)
1830 .map_or(0, |entries| entries.len());
1831
1832 if self.snapshot_manager.should_create_snapshot(
1833 entity_id,
1834 entity_event_count,
1835 event.timestamp,
1836 ) {
1837 if let Err(e) = self.create_snapshot(entity_id) {
1839 tracing::warn!(
1840 "Failed to create automatic snapshot for {}: {}",
1841 entity_id,
1842 e
1843 );
1844 }
1845 }
1846 }
1847
1848 fn validate_event(&self, event: &Event) -> Result<()> {
1850 if event.entity_id_str().is_empty() {
1853 return Err(AllSourceError::ValidationError(
1854 "entity_id cannot be empty".to_string(),
1855 ));
1856 }
1857
1858 if event.event_type_str().is_empty() {
1859 return Err(AllSourceError::ValidationError(
1860 "event_type cannot be empty".to_string(),
1861 ));
1862 }
1863
1864 if event.event_type().is_system() {
1867 return Err(AllSourceError::ValidationError(
1868 "Event types starting with '_system.' are reserved for internal use".to_string(),
1869 ));
1870 }
1871
1872 Ok(())
1873 }
1874
1875 pub fn reset_projection(&self, name: &str) -> Result<usize> {
1877 let projection_manager = self.projections.read();
1878 let projection = projection_manager.get_projection(name).ok_or_else(|| {
1879 AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1880 })?;
1881
1882 projection.clear();
1884
1885 let prefix = format!("{name}:");
1887 let keys_to_remove: Vec<String> = self
1888 .projection_state_cache
1889 .iter()
1890 .filter(|entry| entry.key().starts_with(&prefix))
1891 .map(|entry| entry.key().clone())
1892 .collect();
1893 for key in keys_to_remove {
1894 self.projection_state_cache.remove(&key);
1895 }
1896
1897 let events = self.events.read();
1899 let mut reprocessed = 0usize;
1900 for event in events.iter() {
1901 if projection.process(event).is_ok() {
1902 reprocessed += 1;
1903 }
1904 }
1905
1906 Ok(reprocessed)
1907 }
1908
1909 pub fn get_event_by_id(&self, event_id: &uuid::Uuid) -> Result<Option<Event>> {
1911 if let Some(offset) = self.index.get_by_id(event_id) {
1912 let events = self.events.read();
1913 Ok(events.get(offset).cloned())
1914 } else {
1915 Ok(None)
1916 }
1917 }
1918
1919 #[cfg_attr(feature = "hotpath", hotpath::measure)]
1923 pub fn query(&self, request: &QueryEventsRequest) -> Result<Vec<Event>> {
1924 self.query_window(request, 0, false)
1925 .map(|(events, _)| events)
1926 }
1927
1928 pub fn query_scoped(
1930 &self,
1931 request: &QueryEventsRequest,
1932 scope: &ReadScope,
1933 ) -> Result<Vec<Event>> {
1934 self.query_window_scoped(request, 0, false, scope)
1935 .map(|(events, _)| events)
1936 }
1937
1938 #[cfg_attr(feature = "hotpath", hotpath::measure)]
1955 pub fn query_window(
1956 &self,
1957 request: &QueryEventsRequest,
1958 offset: usize,
1959 descending: bool,
1960 ) -> Result<(Vec<Event>, usize)> {
1961 self.query_window_scoped(request, offset, descending, &ReadScope::unrestricted())
1962 }
1963
1964 pub fn query_window_scoped(
1974 &self,
1975 request: &QueryEventsRequest,
1976 offset: usize,
1977 descending: bool,
1978 scope: &ReadScope,
1979 ) -> Result<(Vec<Event>, usize)> {
1980 if let Some(filter) = &request.payload_filter
1987 && serde_json::from_str::<serde_json::Map<String, serde_json::Value>>(filter).is_err()
1988 {
1989 return Err(AllSourceError::InvalidInput(format!(
1990 "invalid 'payload_filter': expected a JSON object of field/value \
1991 pairs, got '{filter}'"
1992 )));
1993 }
1994
1995 if let Some(ref tenant_id) = request.tenant_id {
2014 self.ensure_tenant_loaded(tenant_id)?;
2015 self.tenant_loader.touch(tenant_id);
2019 }
2020
2021 let query_type = if request.entity_id.is_some() {
2023 "entity"
2024 } else if request.event_type.is_some() {
2025 "type"
2026 } else if request.event_type_prefix.is_some() {
2027 "type_prefix"
2028 } else {
2029 "full_scan"
2030 };
2031
2032 #[cfg(feature = "server")]
2034 let timer = self
2035 .metrics
2036 .query_duration_seconds
2037 .with_label_values(&[query_type])
2038 .start_timer();
2039
2040 #[cfg(feature = "server")]
2042 self.metrics
2043 .queries_total
2044 .with_label_values(&[query_type])
2045 .inc();
2046
2047 let events = self.events.read();
2048
2049 let offsets: Vec<usize> = if let Some(entity_id) = &request.entity_id {
2051 self.index
2053 .get_by_entity(entity_id)
2054 .map(|entries| self.filter_entries(entries, request))
2055 .unwrap_or_default()
2056 } else if let Some(event_type) = &request.event_type {
2057 self.index
2059 .get_by_type(event_type)
2060 .map(|entries| self.filter_entries(entries, request))
2061 .unwrap_or_default()
2062 } else if let Some(prefix) = &request.event_type_prefix {
2063 let entries = self.index.get_by_type_prefix(prefix);
2065 self.filter_entries(entries, request)
2066 } else {
2067 (0..events.len()).collect()
2069 };
2070
2071 let mut matches: Vec<(usize, &Event)> = offsets
2077 .iter()
2078 .filter_map(|&event_offset| events.get(event_offset))
2079 .filter(|event| scope.permits(event.entity_id().as_str()))
2080 .filter(|event| self.apply_filters(event, request))
2081 .enumerate()
2082 .collect();
2083
2084 let total = matches.len();
2087
2088 let order = |a: &(usize, &Event), b: &(usize, &Event)| {
2095 let ascending =
2096 a.1.timestamp
2097 .cmp(&b.1.timestamp)
2098 .then_with(|| a.1.version.cmp(&b.1.version))
2099 .then_with(|| a.0.cmp(&b.0));
2100 if descending {
2101 ascending.reverse()
2102 } else {
2103 ascending
2104 }
2105 };
2106
2107 select_window(&mut matches, offset, request.limit, order);
2109
2110 let results: Vec<Event> = matches
2113 .into_iter()
2114 .skip(offset)
2115 .map(|(_, event)| event.clone())
2116 .collect();
2117
2118 #[cfg(feature = "server")]
2120 {
2121 self.metrics
2122 .query_results_total
2123 .with_label_values(&[query_type])
2124 .inc_by(results.len() as u64);
2125 timer.observe_duration();
2126 }
2127
2128 Ok((results, total))
2129 }
2130
2131 #[cfg_attr(feature = "hotpath", hotpath::measure)]
2133 fn filter_entries(&self, entries: Vec<IndexEntry>, request: &QueryEventsRequest) -> Vec<usize> {
2134 entries
2135 .into_iter()
2136 .filter(|entry| {
2137 if let Some(as_of) = request.as_of
2139 && entry.timestamp > as_of
2140 {
2141 return false;
2142 }
2143 if let Some(since) = request.since
2144 && entry.timestamp < since
2145 {
2146 return false;
2147 }
2148 if let Some(until) = request.until
2149 && entry.timestamp > until
2150 {
2151 return false;
2152 }
2153 true
2154 })
2155 .map(|entry| entry.offset)
2156 .collect()
2157 }
2158
2159 #[cfg_attr(feature = "hotpath", hotpath::measure)]
2161 fn apply_filters(&self, event: &Event, request: &QueryEventsRequest) -> bool {
2162 if let Some(ref tid) = request.tenant_id
2164 && event.tenant_id_str() != tid
2165 {
2166 return false;
2167 }
2168
2169 if let Some(as_of) = request.as_of
2179 && event.timestamp > as_of
2180 {
2181 return false;
2182 }
2183 if let Some(since) = request.since
2184 && event.timestamp < since
2185 {
2186 return false;
2187 }
2188 if let Some(until) = request.until
2189 && event.timestamp > until
2190 {
2191 return false;
2192 }
2193
2194 if let Some(ref excludes) = request.exclude_event_type_prefix {
2198 let et = event.event_type_str();
2199 if excludes
2200 .split(',')
2201 .map(str::trim)
2202 .filter(|p| !p.is_empty())
2203 .any(|p| et.starts_with(p))
2204 {
2205 return false;
2206 }
2207 }
2208
2209 if request.entity_id.is_some()
2211 && let Some(ref event_type) = request.event_type
2212 && event.event_type_str() != event_type
2213 {
2214 return false;
2215 }
2216
2217 if request.entity_id.is_some()
2219 && let Some(ref prefix) = request.event_type_prefix
2220 && !event.event_type_str().starts_with(prefix)
2221 {
2222 return false;
2223 }
2224
2225 if let Some(ref filter_str) = request.payload_filter
2227 && let Ok(filter_obj) =
2228 serde_json::from_str::<serde_json::Map<String, serde_json::Value>>(filter_str)
2229 {
2230 let payload = event.payload();
2231 for (key, expected_value) in &filter_obj {
2232 match payload.get(key) {
2233 Some(actual_value) if actual_value == expected_value => {}
2234 _ => return false,
2235 }
2236 }
2237 }
2238
2239 true
2240 }
2241
2242 #[cfg_attr(feature = "hotpath", hotpath::measure)]
2245 pub fn reconstruct_state(
2246 &self,
2247 entity_id: &str,
2248 as_of: Option<DateTime<Utc>>,
2249 ) -> Result<serde_json::Value> {
2250 let (merged_state, since_timestamp) = if let Some(as_of_time) = as_of {
2252 if let Some(snapshot) = self
2254 .snapshot_manager
2255 .get_snapshot_as_of(entity_id, as_of_time)
2256 {
2257 tracing::debug!(
2258 "Using snapshot from {} for entity {} (saved {} events)",
2259 snapshot.as_of,
2260 entity_id,
2261 snapshot.event_count
2262 );
2263 (snapshot.state.clone(), Some(snapshot.as_of))
2264 } else {
2265 (serde_json::json!({}), None)
2266 }
2267 } else {
2268 if let Some(snapshot) = self.snapshot_manager.get_latest_snapshot(entity_id) {
2270 tracing::debug!(
2271 "Using latest snapshot from {} for entity {}",
2272 snapshot.as_of,
2273 entity_id
2274 );
2275 (snapshot.state.clone(), Some(snapshot.as_of))
2276 } else {
2277 (serde_json::json!({}), None)
2278 }
2279 };
2280
2281 let events = self.query(&QueryEventsRequest {
2283 entity_id: Some(entity_id.to_string()),
2284 event_type: None,
2285 tenant_id: None,
2286 as_of,
2287 since: since_timestamp,
2288 until: None,
2289 limit: None,
2290 event_type_prefix: None,
2291 exclude_event_type_prefix: None,
2292 payload_filter: None,
2293 })?;
2294
2295 if events.is_empty() && since_timestamp.is_none() {
2297 return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
2298 }
2299
2300 let mut merged_state = merged_state;
2302 for event in &events {
2303 if let serde_json::Value::Object(ref mut state_map) = merged_state
2304 && let serde_json::Value::Object(ref payload_map) = event.payload
2305 {
2306 for (key, value) in payload_map {
2307 state_map.insert(key.clone(), value.clone());
2308 }
2309 }
2310 }
2311
2312 let state = serde_json::json!({
2314 "entity_id": entity_id,
2315 "last_updated": events.last().map(|e| e.timestamp),
2316 "event_count": events.len(),
2317 "as_of": as_of,
2318 "current_state": merged_state,
2319 "history": events.iter().map(|e| {
2320 serde_json::json!({
2321 "event_id": e.id,
2322 "type": e.event_type,
2323 "timestamp": e.timestamp,
2324 "payload": e.payload
2325 })
2326 }).collect::<Vec<_>>()
2327 });
2328
2329 Ok(state)
2330 }
2331
2332 pub fn get_snapshot(&self, entity_id: &str) -> Result<serde_json::Value> {
2334 let projections = self.projections.read();
2335
2336 if let Some(snapshot_projection) = projections.get_projection("entity_snapshots")
2337 && let Some(state) = snapshot_projection.get_state(entity_id)
2338 {
2339 return Ok(serde_json::json!({
2340 "entity_id": entity_id,
2341 "snapshot": state,
2342 "from_projection": "entity_snapshots"
2343 }));
2344 }
2345
2346 Err(AllSourceError::EntityNotFound(entity_id.to_string()))
2347 }
2348
2349 pub fn stats(&self) -> StoreStats {
2351 let events = self.events.read();
2352 let index_stats = self.index.stats();
2353
2354 StoreStats {
2355 total_events: events.len(),
2356 total_entities: index_stats.total_entities,
2357 total_event_types: index_stats.total_event_types,
2358 total_ingested: *self.total_ingested.read(),
2359 }
2360 }
2361
2362 pub fn list_streams(&self) -> Vec<StreamInfo> {
2364 self.index
2365 .get_all_entities()
2366 .into_iter()
2367 .map(|entity_id| {
2368 let event_count = self
2369 .index
2370 .get_by_entity(&entity_id)
2371 .map_or(0, |entries| entries.len());
2372 let last_event_at = self
2373 .index
2374 .get_by_entity(&entity_id)
2375 .and_then(|entries| entries.last().map(|e| e.timestamp));
2376 StreamInfo {
2377 stream_id: entity_id,
2378 event_count,
2379 last_event_at,
2380 }
2381 })
2382 .collect()
2383 }
2384
2385 pub fn list_event_types(&self) -> Vec<EventTypeInfo> {
2387 self.index
2388 .get_all_types()
2389 .into_iter()
2390 .map(|event_type| {
2391 let event_count = self
2392 .index
2393 .get_by_type(&event_type)
2394 .map_or(0, |entries| entries.len());
2395 let last_event_at = self
2396 .index
2397 .get_by_type(&event_type)
2398 .and_then(|entries| entries.last().map(|e| e.timestamp));
2399 EventTypeInfo {
2400 event_type,
2401 event_count,
2402 last_event_at,
2403 }
2404 })
2405 .collect()
2406 }
2407
2408 pub fn list_streams_for_tenant(&self, tenant_id: &str) -> Vec<StreamInfo> {
2416 let _ = self.ensure_tenant_loaded(tenant_id);
2417 let events = self.events.read();
2418 let mut by_entity: std::collections::HashMap<&str, (usize, chrono::DateTime<chrono::Utc>)> =
2419 std::collections::HashMap::new();
2420 for ev in events.iter() {
2421 if ev.tenant_id_str() != tenant_id {
2422 continue;
2423 }
2424 let e = by_entity
2425 .entry(ev.entity_id_str())
2426 .or_insert((0, ev.timestamp));
2427 e.0 += 1;
2428 if ev.timestamp > e.1 {
2429 e.1 = ev.timestamp;
2430 }
2431 }
2432 by_entity
2433 .into_iter()
2434 .map(|(entity_id, (count, last))| StreamInfo {
2435 stream_id: entity_id.to_string(),
2436 event_count: count,
2437 last_event_at: Some(last),
2438 })
2439 .collect()
2440 }
2441
2442 pub fn list_event_types_for_tenant(&self, tenant_id: &str) -> Vec<EventTypeInfo> {
2444 let _ = self.ensure_tenant_loaded(tenant_id);
2445 let events = self.events.read();
2446 let mut by_type: std::collections::HashMap<&str, (usize, chrono::DateTime<chrono::Utc>)> =
2447 std::collections::HashMap::new();
2448 for ev in events.iter() {
2449 if ev.tenant_id_str() != tenant_id {
2450 continue;
2451 }
2452 let e = by_type
2453 .entry(ev.event_type_str())
2454 .or_insert((0, ev.timestamp));
2455 e.0 += 1;
2456 if ev.timestamp > e.1 {
2457 e.1 = ev.timestamp;
2458 }
2459 }
2460 by_type
2461 .into_iter()
2462 .map(|(event_type, (count, last))| EventTypeInfo {
2463 event_type: event_type.to_string(),
2464 event_count: count,
2465 last_event_at: Some(last),
2466 })
2467 .collect()
2468 }
2469
2470 pub fn stats_for_tenant(&self, tenant_id: &str) -> TenantStoreStats {
2480 let _ = self.ensure_tenant_loaded(tenant_id);
2481 let events = self.events.read();
2482
2483 let mut entities: std::collections::HashSet<&str> = std::collections::HashSet::new();
2484 let mut census: std::collections::HashMap<&str, usize> = std::collections::HashMap::new();
2485 let mut total_events = 0usize;
2486 let mut oldest: Option<chrono::DateTime<chrono::Utc>> = None;
2487 let mut newest: Option<chrono::DateTime<chrono::Utc>> = None;
2488
2489 for ev in events.iter() {
2490 if ev.tenant_id_str() != tenant_id {
2491 continue;
2492 }
2493
2494 total_events += 1;
2495 entities.insert(ev.entity_id_str());
2496 *census.entry(ev.event_type_str()).or_insert(0) += 1;
2497
2498 let ts = ev.timestamp;
2499 if oldest.is_none_or(|o| ts < o) {
2500 oldest = Some(ts);
2501 }
2502 if newest.is_none_or(|n| ts > n) {
2503 newest = Some(ts);
2504 }
2505 }
2506
2507 TenantStoreStats {
2508 total_events,
2509 total_entities: entities.len(),
2510 total_event_types: census.len(),
2511 total_ingested: total_events as u64,
2515 event_types: census
2516 .into_iter()
2517 .map(|(k, v)| (k.to_string(), v))
2518 .collect(),
2519 oldest_event: oldest,
2520 newest_event: newest,
2521 }
2522 }
2523
2524 pub fn reconstruct_state_for_tenant(
2533 &self,
2534 entity_id: &str,
2535 as_of: Option<DateTime<Utc>>,
2536 tenant_id: &str,
2537 ) -> Result<serde_json::Value> {
2538 let events = self.query(&QueryEventsRequest {
2539 entity_id: Some(entity_id.to_string()),
2540 event_type: None,
2541 tenant_id: Some(tenant_id.to_string()),
2542 as_of,
2543 since: None,
2544 until: None,
2545 limit: None,
2546 event_type_prefix: None,
2547 exclude_event_type_prefix: None,
2548 payload_filter: None,
2549 })?;
2550
2551 if events.is_empty() {
2552 return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
2553 }
2554
2555 let mut merged_state = serde_json::json!({});
2556 for event in &events {
2557 if let serde_json::Value::Object(ref mut state_map) = merged_state
2558 && let serde_json::Value::Object(ref payload_map) = event.payload
2559 {
2560 for (key, value) in payload_map {
2561 state_map.insert(key.clone(), value.clone());
2562 }
2563 }
2564 }
2565
2566 Ok(serde_json::json!({
2567 "entity_id": entity_id,
2568 "last_updated": events.last().map(|e| e.timestamp),
2569 "event_count": events.len(),
2570 "as_of": as_of,
2571 "current_state": merged_state,
2572 "history": events.iter().map(|e| {
2573 serde_json::json!({
2574 "event_id": e.id,
2575 "type": e.event_type,
2576 "timestamp": e.timestamp,
2577 "payload": e.payload
2578 })
2579 }).collect::<Vec<_>>()
2580 }))
2581 }
2582
2583 pub fn enable_wal_replication(
2590 &self,
2591 tx: tokio::sync::broadcast::Sender<crate::infrastructure::persistence::wal::WALEntry>,
2592 ) {
2593 if let Some(ref wal_arc) = self.wal {
2594 wal_arc.set_replication_tx(tx);
2595 tracing::info!("WAL replication broadcast enabled");
2596 } else {
2597 tracing::warn!("Cannot enable WAL replication: WAL is not configured");
2598 }
2599 }
2600
2601 pub fn wal(&self) -> Option<&Arc<WriteAheadLog>> {
2604 self.wal.as_ref()
2605 }
2606
2607 pub fn parquet_storage(&self) -> Option<&Arc<RwLock<ParquetStorage>>> {
2610 self.storage.as_ref()
2611 }
2612}
2613
2614#[derive(Debug, Clone, Default)]
2616pub struct EventStoreConfig {
2617 pub storage_dir: Option<PathBuf>,
2619
2620 pub snapshot_config: SnapshotConfig,
2622
2623 pub wal_dir: Option<PathBuf>,
2625
2626 pub wal_config: WALConfig,
2628
2629 pub compaction_config: CompactionConfig,
2631
2632 pub schema_registry_config: SchemaRegistryConfig,
2634
2635 pub system_data_dir: Option<PathBuf>,
2640
2641 pub bootstrap_tenant: Option<String>,
2643
2644 pub cache_byte_budget: Option<u64>,
2651
2652 pub checkpoint_interval_secs: Option<u64>,
2664
2665 pub read_only: bool,
2669}
2670
2671impl EventStoreConfig {
2672 pub fn with_persistence(storage_dir: impl Into<PathBuf>) -> Self {
2674 Self {
2675 storage_dir: Some(storage_dir.into()),
2676 ..Self::default()
2677 }
2678 }
2679
2680 pub fn with_snapshots(snapshot_config: SnapshotConfig) -> Self {
2682 Self {
2683 snapshot_config,
2684 ..Self::default()
2685 }
2686 }
2687
2688 pub fn with_wal(wal_dir: impl Into<PathBuf>, wal_config: WALConfig) -> Self {
2690 Self {
2691 wal_dir: Some(wal_dir.into()),
2692 wal_config,
2693 ..Self::default()
2694 }
2695 }
2696
2697 pub fn with_all(storage_dir: impl Into<PathBuf>, snapshot_config: SnapshotConfig) -> Self {
2699 Self {
2700 storage_dir: Some(storage_dir.into()),
2701 snapshot_config,
2702 ..Self::default()
2703 }
2704 }
2705
2706 pub fn production(
2708 storage_dir: impl Into<PathBuf>,
2709 wal_dir: impl Into<PathBuf>,
2710 snapshot_config: SnapshotConfig,
2711 wal_config: WALConfig,
2712 compaction_config: CompactionConfig,
2713 ) -> Self {
2714 let storage_dir = storage_dir.into();
2715 let system_data_dir = storage_dir.join("__system");
2716 Self {
2717 storage_dir: Some(storage_dir),
2718 snapshot_config,
2719 wal_dir: Some(wal_dir.into()),
2720 wal_config,
2721 compaction_config,
2722 system_data_dir: Some(system_data_dir),
2723 ..Self::default()
2724 }
2725 }
2726
2727 pub fn effective_system_data_dir(&self) -> Option<PathBuf> {
2732 self.system_data_dir
2733 .clone()
2734 .or_else(|| self.storage_dir.as_ref().map(|d| d.join("__system")))
2735 }
2736
2737 pub fn from_env() -> (Self, &'static str) {
2745 Self::from_env_vars(
2746 std::env::var("ALLSOURCE_DATA_DIR")
2747 .ok()
2748 .filter(|s| !s.is_empty()),
2749 std::env::var("ALLSOURCE_STORAGE_DIR")
2750 .ok()
2751 .filter(|s| !s.is_empty()),
2752 std::env::var("ALLSOURCE_WAL_DIR")
2753 .ok()
2754 .filter(|s| !s.is_empty()),
2755 std::env::var("ALLSOURCE_WAL_ENABLED").ok(),
2756 std::env::var("ALLSOURCE_CACHE_BYTES").ok(),
2757 std::env::var("ALLSOURCE_SNAPSHOT_INTERVAL_SECONDS").ok(),
2758 std::env::var("ALLSOURCE_RETENTION_SYSTEM_DAYS").ok(),
2759 std::env::var("ALLSOURCE_CHECKPOINT_INTERVAL_SECONDS").ok(),
2760 )
2761 }
2762
2763 pub fn from_env_vars(
2765 data_dir: Option<String>,
2766 explicit_storage_dir: Option<String>,
2767 explicit_wal_dir: Option<String>,
2768 wal_enabled_var: Option<String>,
2769 cache_bytes_var: Option<String>,
2770 snapshot_interval_var: Option<String>,
2771 retention_system_days_var: Option<String>,
2772 checkpoint_interval_var: Option<String>,
2773 ) -> (Self, &'static str) {
2774 let data_dir = data_dir.filter(|s| !s.is_empty());
2775 let storage_dir = explicit_storage_dir
2776 .filter(|s| !s.is_empty())
2777 .or_else(|| data_dir.as_ref().map(|d| format!("{d}/storage")));
2778 let wal_dir = explicit_wal_dir
2779 .filter(|s| !s.is_empty())
2780 .or_else(|| data_dir.as_ref().map(|d| format!("{d}/wal")));
2781 let wal_enabled = wal_enabled_var.is_none_or(|v| v == "true");
2782 let cache_byte_budget =
2787 cache_bytes_var
2788 .filter(|s| !s.is_empty())
2789 .and_then(|s| match s.parse::<u64>() {
2790 Ok(v) => Some(v),
2791 Err(e) => {
2792 tracing::warn!(
2793 "ALLSOURCE_CACHE_BYTES={s:?} could not be parsed as u64: {e}; \
2794 cache budget disabled"
2795 );
2796 None
2797 }
2798 });
2799 let compaction_config =
2800 CompactionConfig::from_env_vars(snapshot_interval_var, retention_system_days_var);
2801
2802 let checkpoint_interval_secs = if wal_enabled {
2807 checkpoint_interval_var
2808 .filter(|s| !s.is_empty())
2809 .map(|s| match s.parse::<u64>() {
2810 Ok(v) => v,
2811 Err(e) => {
2812 tracing::warn!(
2813 "ALLSOURCE_CHECKPOINT_INTERVAL_SECONDS={s:?} could not be parsed as \
2814 u64: {e}; falling back to default 60s"
2815 );
2816 60
2817 }
2818 })
2819 .or(Some(60))
2820 } else {
2821 None
2822 };
2823
2824 let mut config = match (&storage_dir, &wal_dir) {
2825 (Some(sd), Some(wd)) if wal_enabled => Self::production(
2826 sd,
2827 wd,
2828 SnapshotConfig::default(),
2829 WALConfig::default(),
2830 compaction_config,
2831 ),
2832 (Some(sd), _) => Self::with_persistence(sd),
2833 (_, Some(wd)) if wal_enabled => Self::with_wal(wd, WALConfig::default()),
2834 _ => Self::default(),
2835 };
2836 config.cache_byte_budget = cache_byte_budget;
2837 config.checkpoint_interval_secs = checkpoint_interval_secs;
2838
2839 let mode = match (&storage_dir, &wal_dir) {
2840 (Some(_), Some(_)) if wal_enabled => "wal+parquet",
2841 (Some(_), _) => "parquet-only",
2842 (_, Some(_)) if wal_enabled => "wal-only",
2843 _ => "in-memory",
2844 };
2845 (config, mode)
2846 }
2847}
2848
2849#[derive(Debug, serde::Serialize)]
2850pub struct StoreStats {
2851 pub total_events: usize,
2852 pub total_entities: usize,
2853 pub total_event_types: usize,
2854 pub total_ingested: u64,
2855}
2856
2857#[derive(Debug, Clone, serde::Serialize)]
2863pub struct TenantStoreStats {
2864 pub total_events: usize,
2865 pub total_entities: usize,
2866 pub total_event_types: usize,
2867 pub total_ingested: u64,
2868 pub event_types: std::collections::HashMap<String, usize>,
2870 pub oldest_event: Option<chrono::DateTime<chrono::Utc>>,
2871 pub newest_event: Option<chrono::DateTime<chrono::Utc>>,
2872}
2873
2874#[derive(Debug, Clone, serde::Serialize)]
2876pub struct StreamInfo {
2877 pub stream_id: String,
2879 pub event_count: usize,
2881 pub last_event_at: Option<chrono::DateTime<chrono::Utc>>,
2883}
2884
2885#[derive(Debug, Clone, serde::Serialize)]
2887pub struct EventTypeInfo {
2888 pub event_type: String,
2890 pub event_count: usize,
2892 pub last_event_at: Option<chrono::DateTime<chrono::Utc>>,
2894}
2895
2896impl Default for EventStore {
2897 fn default() -> Self {
2898 Self::new()
2899 }
2900}
2901
2902#[cfg(test)]
2903mod tests {
2904 use super::*;
2905 use crate::domain::entities::Event;
2906 use tempfile::TempDir;
2907
2908 fn find_parquet_files(dir: &std::path::Path) -> Vec<std::path::PathBuf> {
2913 let mut out = Vec::new();
2914 let mut stack = vec![dir.to_path_buf()];
2915 while let Some(d) = stack.pop() {
2916 let Ok(entries) = std::fs::read_dir(&d) else {
2917 continue;
2918 };
2919 for e in entries.flatten() {
2920 let p = e.path();
2921 if p.is_dir() {
2922 stack.push(p);
2923 } else if p.extension().and_then(|s| s.to_str()) == Some("parquet") {
2924 out.push(p);
2925 }
2926 }
2927 }
2928 out
2929 }
2930
2931 fn create_test_event(entity_id: &str, event_type: &str) -> Event {
2932 Event::from_strings(
2933 event_type.to_string(),
2934 entity_id.to_string(),
2935 "default".to_string(),
2936 serde_json::json!({"name": "Test", "value": 42}),
2937 None,
2938 )
2939 .unwrap()
2940 }
2941
2942 fn create_test_event_with_payload(
2943 entity_id: &str,
2944 event_type: &str,
2945 payload: serde_json::Value,
2946 ) -> Event {
2947 Event::from_strings(
2948 event_type.to_string(),
2949 entity_id.to_string(),
2950 "default".to_string(),
2951 payload,
2952 None,
2953 )
2954 .unwrap()
2955 }
2956
2957 #[test]
2958 fn test_event_store_new() {
2959 let store = EventStore::new();
2960 assert_eq!(store.stats().total_events, 0);
2961 assert_eq!(store.stats().total_entities, 0);
2962 }
2963
2964 #[test]
2971 fn test_ensure_tenant_loaded_no_storage_is_a_noop() {
2972 let store = EventStore::new();
2976 assert!(!store.is_tenant_loaded("alice"));
2977 store.ensure_tenant_loaded("alice").unwrap();
2978 assert!(store.is_tenant_loaded("alice"));
2979 assert!(!store.is_tenant_loaded("bob"));
2981 }
2982
2983 #[test]
2984 fn test_ensure_tenant_loaded_warm_path_is_idempotent() {
2985 let store = EventStore::new();
2986 store.ensure_tenant_loaded("alice").unwrap();
2987 store.ensure_tenant_loaded("alice").unwrap();
2989 }
2990
2991 #[test]
2992 fn test_ensure_tenant_loaded_rejects_unsafe_tenant_id() {
2993 let temp_dir = TempDir::new().unwrap();
2999 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3000 for unsafe_tid in ["..", "a/b", "a\\b", ""] {
3001 let result = store.ensure_tenant_loaded(unsafe_tid);
3002 assert!(
3003 result.is_err(),
3004 "tenant_id {unsafe_tid:?} should have been rejected"
3005 );
3006 assert!(
3007 !store.is_tenant_loaded(unsafe_tid),
3008 "rejected tenant {unsafe_tid:?} must not be marked loaded"
3009 );
3010 }
3011 }
3012
3013 #[test]
3014 fn test_ensure_tenant_loaded_no_subtree_marks_loaded_with_zero_events() {
3015 let temp_dir = TempDir::new().unwrap();
3020 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3021 assert!(!store.is_tenant_loaded("never-existed"));
3022 store.ensure_tenant_loaded("never-existed").unwrap();
3023 assert!(store.is_tenant_loaded("never-existed"));
3024 }
3025
3026 #[test]
3027 fn test_evict_tenant_drops_events_and_resets_bytes() {
3028 let temp_dir = TempDir::new().unwrap();
3032 let storage_dir = temp_dir.path().to_path_buf();
3033
3034 {
3035 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3036 for i in 0..3 {
3037 store
3038 .ingest(
3039 &Event::from_strings(
3040 "test.event".to_string(),
3041 format!("a-{i}"),
3042 "alice".to_string(),
3043 serde_json::json!({"i": i}),
3044 None,
3045 )
3046 .unwrap(),
3047 )
3048 .unwrap();
3049 }
3050 for i in 0..2 {
3051 store
3052 .ingest(
3053 &Event::from_strings(
3054 "test.event".to_string(),
3055 format!("b-{i}"),
3056 "bob".to_string(),
3057 serde_json::json!({"i": i}),
3058 None,
3059 )
3060 .unwrap(),
3061 )
3062 .unwrap();
3063 }
3064 store.flush_storage().unwrap();
3065 }
3066
3067 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3068 store.ensure_tenant_loaded("alice").unwrap();
3069 store.ensure_tenant_loaded("bob").unwrap();
3070 assert_eq!(store.stats().total_events, 5);
3071 let alice_bytes = store.tenant_resident_bytes("alice");
3072 let bob_bytes = store.tenant_resident_bytes("bob");
3073 assert!(alice_bytes > 0 && bob_bytes > 0);
3074
3075 store.evict_tenant("alice");
3076
3077 assert!(!store.is_tenant_loaded("alice"));
3078 assert!(store.is_tenant_loaded("bob"));
3079 assert_eq!(store.tenant_resident_bytes("alice"), 0);
3080 assert_eq!(store.tenant_resident_bytes("bob"), bob_bytes);
3081 assert_eq!(store.stats().total_events, 2, "only bob's 2 events remain");
3082 }
3083
3084 #[test]
3085 fn test_evict_tenant_then_query_re_loads_from_disk() {
3086 let temp_dir = TempDir::new().unwrap();
3090 let storage_dir = temp_dir.path().to_path_buf();
3091
3092 {
3093 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3094 for i in 0..4 {
3095 store
3096 .ingest(
3097 &Event::from_strings(
3098 "test.event".to_string(),
3099 format!("a-{i}"),
3100 "alice".to_string(),
3101 serde_json::json!({"i": i}),
3102 None,
3103 )
3104 .unwrap(),
3105 )
3106 .unwrap();
3107 }
3108 store.flush_storage().unwrap();
3109 }
3110
3111 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3112 store.ensure_tenant_loaded("alice").unwrap();
3113 store.evict_tenant("alice");
3114 assert_eq!(store.stats().total_events, 0);
3115
3116 let results = store
3118 .query(&QueryEventsRequest {
3119 entity_id: None,
3120 event_type: None,
3121 tenant_id: Some("alice".to_string()),
3122 as_of: None,
3123 since: None,
3124 until: None,
3125 limit: None,
3126 event_type_prefix: None,
3127 exclude_event_type_prefix: None,
3128 payload_filter: None,
3129 })
3130 .unwrap();
3131 assert_eq!(results.len(), 4);
3132 assert!(store.is_tenant_loaded("alice"));
3133 }
3134
3135 #[test]
3136 fn test_evict_tenant_rebuilds_index_with_new_offsets() {
3137 let temp_dir = TempDir::new().unwrap();
3144 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3145
3146 for i in 0..3 {
3150 store
3151 .ingest(
3152 &Event::from_strings(
3153 "test.event".to_string(),
3154 format!("a-{i}"),
3155 "alice".to_string(),
3156 serde_json::json!({"i": i}),
3157 None,
3158 )
3159 .unwrap(),
3160 )
3161 .unwrap();
3162 if i < 2 {
3163 store
3164 .ingest(
3165 &Event::from_strings(
3166 "test.event".to_string(),
3167 format!("b-{i}"),
3168 "bob".to_string(),
3169 serde_json::json!({"i": i}),
3170 None,
3171 )
3172 .unwrap(),
3173 )
3174 .unwrap();
3175 }
3176 }
3177 store.tenant_loader.mark_loaded("alice");
3179 store.tenant_loader.mark_loaded("bob");
3180
3181 store.evict_tenant("alice");
3182
3183 let bob_results = store
3184 .query(&QueryEventsRequest {
3185 entity_id: None,
3186 event_type: None,
3187 tenant_id: Some("bob".to_string()),
3188 as_of: None,
3189 since: None,
3190 until: None,
3191 limit: None,
3192 event_type_prefix: None,
3193 exclude_event_type_prefix: None,
3194 payload_filter: None,
3195 })
3196 .unwrap();
3197 assert_eq!(bob_results.len(), 2);
3198 for e in &bob_results {
3199 assert_eq!(e.tenant_id_str(), "bob");
3200 }
3201 }
3202
3203 #[test]
3204 fn test_budget_eviction_keeps_resident_set_bounded() {
3205 let temp_dir = TempDir::new().unwrap();
3209 let storage_dir = temp_dir.path().to_path_buf();
3210
3211 let big_payload = serde_json::json!({"data": "x".repeat(1000)});
3214 {
3215 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3216 for tenant in ["alice", "bob", "carol"] {
3217 for i in 0..5 {
3218 store
3219 .ingest(
3220 &Event::from_strings(
3221 "test.event".to_string(),
3222 format!("{tenant}-{i}"),
3223 tenant.to_string(),
3224 big_payload.clone(),
3225 None,
3226 )
3227 .unwrap(),
3228 )
3229 .unwrap();
3230 }
3231 }
3232 store.flush_storage().unwrap();
3233 }
3234
3235 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3238 config.cache_byte_budget = Some(12_000);
3239 let store = EventStore::with_config(config);
3240
3241 store.ensure_tenant_loaded("alice").unwrap();
3243 assert!(store.is_tenant_loaded("alice"));
3244
3245 store.tenant_loader.touch("alice");
3250 std::thread::sleep(std::time::Duration::from_millis(10));
3251 store.ensure_tenant_loaded("bob").unwrap();
3252 assert!(store.is_tenant_loaded("bob"));
3253
3254 store.tenant_loader.touch("bob");
3258 std::thread::sleep(std::time::Duration::from_millis(10));
3259 store.ensure_tenant_loaded("carol").unwrap();
3260 assert!(store.is_tenant_loaded("carol"));
3261
3262 let resident = store.cache_resident_bytes();
3267 let budget = 12_000u64;
3268
3269 if resident > budget {
3272 let loaded_count = ["alice", "bob", "carol"]
3273 .iter()
3274 .filter(|t| store.is_tenant_loaded(t))
3275 .count();
3276 assert_eq!(
3277 loaded_count, 1,
3278 "over budget but more than one tenant loaded — eviction policy didn't fire"
3279 );
3280 }
3281
3282 assert!(store.is_tenant_loaded("carol"));
3285 }
3286
3287 #[test]
3288 fn test_query_after_eviction_re_loads_transparently() {
3289 let temp_dir = TempDir::new().unwrap();
3292 let storage_dir = temp_dir.path().to_path_buf();
3293
3294 let big_payload = serde_json::json!({"data": "x".repeat(2000)});
3295 {
3296 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3297 for tenant in ["alice", "bob"] {
3298 for i in 0..3 {
3299 store
3300 .ingest(
3301 &Event::from_strings(
3302 "test.event".to_string(),
3303 format!("{tenant}-{i}"),
3304 tenant.to_string(),
3305 big_payload.clone(),
3306 None,
3307 )
3308 .unwrap(),
3309 )
3310 .unwrap();
3311 }
3312 }
3313 store.flush_storage().unwrap();
3314 }
3315
3316 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3318 config.cache_byte_budget = Some(5_000);
3319 let store = EventStore::with_config(config);
3320
3321 let alice_first = store
3325 .query(&QueryEventsRequest {
3326 entity_id: None,
3327 event_type: None,
3328 tenant_id: Some("alice".to_string()),
3329 as_of: None,
3330 since: None,
3331 until: None,
3332 limit: None,
3333 event_type_prefix: None,
3334 exclude_event_type_prefix: None,
3335 payload_filter: None,
3336 })
3337 .unwrap();
3338 assert_eq!(alice_first.len(), 3);
3339
3340 std::thread::sleep(std::time::Duration::from_millis(15));
3342 let _bob = store
3344 .query(&QueryEventsRequest {
3345 entity_id: None,
3346 event_type: None,
3347 tenant_id: Some("bob".to_string()),
3348 as_of: None,
3349 since: None,
3350 until: None,
3351 limit: None,
3352 event_type_prefix: None,
3353 exclude_event_type_prefix: None,
3354 payload_filter: None,
3355 })
3356 .unwrap();
3357 assert!(
3358 !store.is_tenant_loaded("alice"),
3359 "alice should have been evicted"
3360 );
3361
3362 let alice_second = store
3364 .query(&QueryEventsRequest {
3365 entity_id: None,
3366 event_type: None,
3367 tenant_id: Some("alice".to_string()),
3368 as_of: None,
3369 since: None,
3370 until: None,
3371 limit: None,
3372 event_type_prefix: None,
3373 exclude_event_type_prefix: None,
3374 payload_filter: None,
3375 })
3376 .unwrap();
3377 assert_eq!(
3378 alice_second.len(),
3379 3,
3380 "alice's events come back via re-load"
3381 );
3382 assert!(store.is_tenant_loaded("alice"));
3383 }
3384
3385 #[test]
3386 #[cfg(feature = "server")]
3387 fn test_cache_metrics_track_evictions_and_bytes() {
3388 let temp_dir = TempDir::new().unwrap();
3392 let storage_dir = temp_dir.path().to_path_buf();
3393
3394 let big_payload = serde_json::json!({"data": "x".repeat(2000)});
3395 {
3396 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3397 for tenant in ["alice", "bob"] {
3398 for i in 0..3 {
3399 store
3400 .ingest(
3401 &Event::from_strings(
3402 "test.event".to_string(),
3403 format!("{tenant}-{i}"),
3404 tenant.to_string(),
3405 big_payload.clone(),
3406 None,
3407 )
3408 .unwrap(),
3409 )
3410 .unwrap();
3411 }
3412 }
3413 store.flush_storage().unwrap();
3414 }
3415
3416 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3417 config.cache_byte_budget = Some(5_000); let store = EventStore::with_config(config);
3419
3420 assert_eq!(store.metrics.cache_evictions_total.get(), 0);
3421 assert_eq!(store.metrics.cache_bytes.get(), 0);
3422
3423 store.ensure_tenant_loaded("alice").unwrap();
3424 let after_alice = store.metrics.cache_bytes.get();
3426 assert!(after_alice > 0, "gauge should reflect alice's bytes");
3427 assert_eq!(store.metrics.cache_evictions_total.get(), 0);
3429
3430 std::thread::sleep(std::time::Duration::from_millis(10));
3431 store.ensure_tenant_loaded("bob").unwrap();
3432
3433 assert_eq!(
3436 store.metrics.cache_evictions_total.get(),
3437 1,
3438 "exactly one tenant evicted after bob's load"
3439 );
3440 let after_bob = store.metrics.cache_bytes.get();
3442 assert!(after_bob > 0);
3443 assert!(after_bob <= after_alice, "gauge dropped after eviction");
3444 }
3445
3446 #[test]
3447 #[cfg(feature = "server")]
3448 fn test_storage_size_gauge_populated_from_on_disk_bytes() {
3449 let temp_dir = TempDir::new().unwrap();
3454 let config = EventStoreConfig {
3457 storage_dir: Some(temp_dir.path().join("parquet")),
3458 wal_dir: Some(temp_dir.path().join("wal")),
3459 ..EventStoreConfig::default()
3460 };
3461 let store = EventStore::with_config(config);
3462
3463 assert_eq!(store.metrics.storage_size_bytes.get(), 0);
3465
3466 let payload = serde_json::json!({ "data": "x".repeat(2000) });
3467 for i in 0..10 {
3468 store
3469 .ingest(
3470 &Event::from_strings(
3471 "test.event".to_string(),
3472 format!("entity-{i}"),
3473 "tenant-a".to_string(),
3474 payload.clone(),
3475 None,
3476 )
3477 .unwrap(),
3478 )
3479 .unwrap();
3480 }
3481 store.flush_storage().unwrap();
3482
3483 store.refresh_storage_metrics_now();
3485
3486 let size = store.metrics.storage_size_bytes.get();
3487 assert!(
3488 size > 0,
3489 "storage_size_bytes must reflect real on-disk bytes, got {size}"
3490 );
3491 assert!(
3492 store.metrics.parquet_files_total.get() >= 1,
3493 "at least one Parquet file should exist after a flush"
3494 );
3495
3496 assert!(
3499 store.metrics.wal_segments_total.get() >= 1,
3500 "at least one WAL segment should exist after writes"
3501 );
3502
3503 let parquet_stats = store.storage.as_ref().unwrap().read().stats().unwrap();
3506 let (wal_bytes, _) = store.wal.as_ref().unwrap().on_disk_stats().unwrap();
3507 assert_eq!(
3508 size as u64,
3509 parquet_stats.total_size_bytes + wal_bytes,
3510 "gauge should equal Parquet bytes ({}) + WAL bytes ({wal_bytes})",
3511 parquet_stats.total_size_bytes
3512 );
3513 }
3514
3515 #[test]
3516 fn test_stress_resident_set_stays_near_budget_under_rolling_queries() {
3517 let temp_dir = TempDir::new().unwrap();
3524 let storage_dir = temp_dir.path().to_path_buf();
3525
3526 const TENANT_COUNT: usize = 10;
3527 const EVENTS_PER_TENANT: usize = 50;
3528 let big_payload = serde_json::json!({"data": "x".repeat(10_000)});
3530
3531 {
3534 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3535 for t in 0..TENANT_COUNT {
3536 let tenant = format!("tenant-{t}");
3537 for i in 0..EVENTS_PER_TENANT {
3538 store
3539 .ingest(
3540 &Event::from_strings(
3541 "test.event".to_string(),
3542 format!("{tenant}-{i}"),
3543 tenant.clone(),
3544 big_payload.clone(),
3545 None,
3546 )
3547 .unwrap(),
3548 )
3549 .unwrap();
3550 }
3551 }
3552 store.flush_storage().unwrap();
3553 }
3554
3555 const BUDGET: u64 = 1_048_576;
3559 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3560 config.cache_byte_budget = Some(BUDGET);
3561 let store = EventStore::with_config(config);
3562
3563 let mut peak_resident: u64 = 0;
3567 for t in 0..TENANT_COUNT {
3568 let tenant = format!("tenant-{t}");
3569 let results = store
3570 .query(&QueryEventsRequest {
3571 entity_id: None,
3572 event_type: None,
3573 tenant_id: Some(tenant.clone()),
3574 as_of: None,
3575 since: None,
3576 until: None,
3577 limit: None,
3578 event_type_prefix: None,
3579 exclude_event_type_prefix: None,
3580 payload_filter: None,
3581 })
3582 .unwrap();
3583 assert_eq!(
3584 results.len(),
3585 EVENTS_PER_TENANT,
3586 "every per-tenant query must return all of that tenant's events"
3587 );
3588 let resident = store.cache_resident_bytes();
3590 if resident > peak_resident {
3591 peak_resident = resident;
3592 }
3593 }
3594
3595 let final_resident = store.cache_resident_bytes();
3596
3597 let tolerance = BUDGET; assert!(
3603 peak_resident <= BUDGET + tolerance,
3604 "peak resident {peak_resident} exceeds budget {BUDGET} by more than {tolerance} \
3605 — eviction policy not keeping up with the working-set churn"
3606 );
3607 assert!(
3608 final_resident <= BUDGET + tolerance,
3609 "final resident {final_resident} exceeds budget {BUDGET} by more than {tolerance}"
3610 );
3611
3612 let last_tenant = format!("tenant-{}", TENANT_COUNT - 1);
3615 assert!(
3616 store.is_tenant_loaded(&last_tenant),
3617 "the most-recent tenant must remain loaded after the sweep"
3618 );
3619
3620 let still_loaded = (0..TENANT_COUNT)
3623 .filter(|t| store.is_tenant_loaded(&format!("tenant-{t}")))
3624 .count();
3625 assert!(
3626 still_loaded < TENANT_COUNT,
3627 "no tenants evicted ({still_loaded}/{TENANT_COUNT} still loaded) — \
3628 budget enforcement didn't engage"
3629 );
3630 }
3631
3632 #[test]
3633 fn test_evict_tenant_when_not_loaded_is_a_noop() {
3634 let store = EventStore::new();
3637 store.evict_tenant("nobody"); assert!(!store.is_tenant_loaded("nobody"));
3639 }
3640
3641 #[test]
3642 fn test_lazy_load_accounts_bytes_per_tenant() {
3643 let temp_dir = TempDir::new().unwrap();
3647 let storage_dir = temp_dir.path().to_path_buf();
3648
3649 {
3651 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3652 for i in 0..5 {
3653 store
3654 .ingest(
3655 &Event::from_strings(
3656 "test.event".to_string(),
3657 format!("a-{i}"),
3658 "alice".to_string(),
3659 serde_json::json!({"data": "x".repeat(1000)}),
3660 None,
3661 )
3662 .unwrap(),
3663 )
3664 .unwrap();
3665 }
3666 store.flush_storage().unwrap();
3667 }
3668
3669 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3670 assert_eq!(store.tenant_resident_bytes("alice"), 0);
3672 assert_eq!(store.cache_resident_bytes(), 0);
3673
3674 store.ensure_tenant_loaded("alice").unwrap();
3675
3676 let alice_bytes = store.tenant_resident_bytes("alice");
3679 assert!(
3680 alice_bytes >= 5 * 1000,
3681 "alice should have at least 5 KiB resident; got {alice_bytes}"
3682 );
3683 assert_eq!(store.tenant_resident_bytes("bob"), 0);
3685 assert_eq!(store.cache_resident_bytes(), alice_bytes);
3687 }
3688
3689 #[test]
3690 fn test_query_lazy_loads_tenant_on_first_call() {
3691 let temp_dir = TempDir::new().unwrap();
3695 let storage_dir = temp_dir.path().to_path_buf();
3696
3697 {
3699 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3700 for i in 0..3 {
3701 let event = Event::from_strings(
3702 "test.event".to_string(),
3703 format!("e-{i}"),
3704 "alice".to_string(),
3705 serde_json::json!({"i": i}),
3706 None,
3707 )
3708 .unwrap();
3709 store.ingest(&event).unwrap();
3710 }
3711 store.flush_storage().unwrap();
3712 }
3713
3714 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3716 assert_eq!(
3717 store.stats().total_events,
3718 0,
3719 "boot must be O(1) — no Parquet pre-load"
3720 );
3721 assert!(!store.is_tenant_loaded("alice"));
3722 assert!(!store.is_tenant_loaded("bob"));
3723
3724 let results = store
3726 .query(&QueryEventsRequest {
3727 entity_id: None,
3728 event_type: None,
3729 tenant_id: Some("alice".to_string()),
3730 as_of: None,
3731 since: None,
3732 until: None,
3733 limit: None,
3734 event_type_prefix: None,
3735 exclude_event_type_prefix: None,
3736 payload_filter: None,
3737 })
3738 .unwrap();
3739 assert_eq!(results.len(), 3, "alice's 3 events are returned");
3740 assert!(store.is_tenant_loaded("alice"), "alice now warm");
3741 assert!(!store.is_tenant_loaded("bob"), "bob still cold");
3744 }
3745
3746 #[test]
3747 fn test_query_invalid_tenant_id_returns_error_no_hang() {
3748 let temp_dir = TempDir::new().unwrap();
3752 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3753
3754 let result = store.query(&QueryEventsRequest {
3755 entity_id: None,
3756 event_type: None,
3757 tenant_id: Some("../etc".to_string()),
3758 as_of: None,
3759 since: None,
3760 until: None,
3761 limit: None,
3762 event_type_prefix: None,
3763 exclude_event_type_prefix: None,
3764 payload_filter: None,
3765 });
3766 assert!(result.is_err(), "unsafe tenant_id must surface as error");
3767 }
3768
3769 #[test]
3770 fn test_query_concurrent_first_queries_for_same_tenant_all_succeed() {
3771 let temp_dir = TempDir::new().unwrap();
3779 let storage_dir = temp_dir.path().to_path_buf();
3780
3781 {
3783 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3784 for i in 0..25 {
3785 let event = Event::from_strings(
3786 "test.event".to_string(),
3787 format!("e-{i}"),
3788 "alice".to_string(),
3789 serde_json::json!({"i": i}),
3790 None,
3791 )
3792 .unwrap();
3793 store.ingest(&event).unwrap();
3794 }
3795 store.flush_storage().unwrap();
3796 }
3797
3798 let store = Arc::new(EventStore::with_config(EventStoreConfig::with_persistence(
3800 &storage_dir,
3801 )));
3802 assert!(!store.is_tenant_loaded("alice"));
3803
3804 let mut handles = Vec::new();
3805 for _ in 0..8 {
3806 let s = store.clone();
3807 handles.push(std::thread::spawn(move || {
3808 s.query(&QueryEventsRequest {
3809 entity_id: None,
3810 event_type: None,
3811 tenant_id: Some("alice".to_string()),
3812 as_of: None,
3813 since: None,
3814 until: None,
3815 limit: None,
3816 event_type_prefix: None,
3817 exclude_event_type_prefix: None,
3818 payload_filter: None,
3819 })
3820 }));
3821 }
3822
3823 for h in handles {
3824 let result = h.join().unwrap().unwrap();
3825 assert_eq!(
3826 result.len(),
3827 25,
3828 "every concurrent caller must see all 25 events"
3829 );
3830 }
3831 assert!(store.is_tenant_loaded("alice"));
3832 assert_eq!(store.stats().total_events, 25);
3834 }
3835
3836 #[test]
3837 fn test_query_two_cold_tenants_load_independently() {
3838 let temp_dir = TempDir::new().unwrap();
3842 let storage_dir = temp_dir.path().to_path_buf();
3843
3844 {
3845 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3846 for i in 0..3 {
3847 store
3848 .ingest(
3849 &Event::from_strings(
3850 "test.event".to_string(),
3851 format!("a-{i}"),
3852 "alice".to_string(),
3853 serde_json::json!({"i": i}),
3854 None,
3855 )
3856 .unwrap(),
3857 )
3858 .unwrap();
3859 }
3860 for i in 0..5 {
3861 store
3862 .ingest(
3863 &Event::from_strings(
3864 "test.event".to_string(),
3865 format!("b-{i}"),
3866 "bob".to_string(),
3867 serde_json::json!({"i": i}),
3868 None,
3869 )
3870 .unwrap(),
3871 )
3872 .unwrap();
3873 }
3874 store.flush_storage().unwrap();
3875 }
3876
3877 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3878 assert_eq!(store.stats().total_events, 0);
3879
3880 let alice = store
3882 .query(&QueryEventsRequest {
3883 entity_id: None,
3884 event_type: None,
3885 tenant_id: Some("alice".to_string()),
3886 as_of: None,
3887 since: None,
3888 until: None,
3889 limit: None,
3890 event_type_prefix: None,
3891 exclude_event_type_prefix: None,
3892 payload_filter: None,
3893 })
3894 .unwrap();
3895 assert_eq!(alice.len(), 3);
3896 assert!(store.is_tenant_loaded("alice"));
3897 assert!(!store.is_tenant_loaded("bob"));
3898 assert_eq!(store.stats().total_events, 3);
3899
3900 let bob = store
3902 .query(&QueryEventsRequest {
3903 entity_id: None,
3904 event_type: None,
3905 tenant_id: Some("bob".to_string()),
3906 as_of: None,
3907 since: None,
3908 until: None,
3909 limit: None,
3910 event_type_prefix: None,
3911 exclude_event_type_prefix: None,
3912 payload_filter: None,
3913 })
3914 .unwrap();
3915 assert_eq!(bob.len(), 5);
3916 assert!(store.is_tenant_loaded("bob"));
3917 assert_eq!(store.stats().total_events, 8);
3918 }
3919
3920 #[test]
3921 fn test_boot_with_persisted_data_is_o1() {
3922 let temp_dir = TempDir::new().unwrap();
3935 let storage_dir = temp_dir.path().to_path_buf();
3936
3937 {
3938 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3939 for tenant in ["alice", "bob", "carol"] {
3940 for i in 0..50 / 3 {
3941 store
3942 .ingest(
3943 &Event::from_strings(
3944 "test.event".to_string(),
3945 format!("{tenant}-{i}"),
3946 tenant.to_string(),
3947 serde_json::json!({"i": i}),
3948 None,
3949 )
3950 .unwrap(),
3951 )
3952 .unwrap();
3953 }
3954 }
3955 store.flush_storage().unwrap();
3956 }
3957
3958 let on_disk = find_parquet_files(&storage_dir);
3960 assert!(
3961 !on_disk.is_empty(),
3962 "session 1 should have produced parquet files; pre-condition for the test"
3963 );
3964
3965 let started = std::time::Instant::now();
3966 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3967 let boot_elapsed = started.elapsed();
3968
3969 assert_eq!(
3970 store.stats().total_events,
3971 0,
3972 "boot must not pre-load any Parquet events"
3973 );
3974
3975 assert!(
3979 boot_elapsed < std::time::Duration::from_secs(2),
3980 "boot took {boot_elapsed:?} — Step 2 boot should be O(1)"
3981 );
3982 }
3983
3984 #[test]
3985 fn test_query_warm_tenant_does_not_re_read_disk() {
3986 let temp_dir = TempDir::new().unwrap();
3992 let storage_dir = temp_dir.path().to_path_buf();
3993
3994 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3995 for i in 0..3 {
3996 let event = Event::from_strings(
3997 "test.event".to_string(),
3998 format!("e-{i}"),
3999 "alice".to_string(),
4000 serde_json::json!({"i": i}),
4001 None,
4002 )
4003 .unwrap();
4004 store.ingest(&event).unwrap();
4005 }
4006 store.flush_storage().unwrap();
4007
4008 let _ = store
4010 .query(&QueryEventsRequest {
4011 entity_id: None,
4012 event_type: None,
4013 tenant_id: Some("alice".to_string()),
4014 as_of: None,
4015 since: None,
4016 until: None,
4017 limit: None,
4018 event_type_prefix: None,
4019 exclude_event_type_prefix: None,
4020 payload_filter: None,
4021 })
4022 .unwrap();
4023 assert!(store.is_tenant_loaded("alice"));
4024
4025 let parquet_files = find_parquet_files(&storage_dir);
4028 for f in parquet_files {
4029 std::fs::remove_file(&f).unwrap();
4030 }
4031
4032 let results = store
4033 .query(&QueryEventsRequest {
4034 entity_id: None,
4035 event_type: None,
4036 tenant_id: Some("alice".to_string()),
4037 as_of: None,
4038 since: None,
4039 until: None,
4040 limit: None,
4041 event_type_prefix: None,
4042 exclude_event_type_prefix: None,
4043 payload_filter: None,
4044 })
4045 .unwrap();
4046 assert_eq!(
4047 results.len(),
4048 3,
4049 "warm tenant query must not need disk; got {} events from a deleted parquet",
4050 results.len()
4051 );
4052 }
4053
4054 #[test]
4055 fn test_event_store_default() {
4056 let store = EventStore::default();
4057 assert_eq!(store.stats().total_events, 0);
4058 }
4059
4060 #[test]
4061 fn test_ingest_single_event() {
4062 let store = EventStore::new();
4063 let event = create_test_event("entity-1", "user.created");
4064
4065 store.ingest(&event).unwrap();
4066
4067 assert_eq!(store.stats().total_events, 1);
4068 assert_eq!(store.stats().total_ingested, 1);
4069 }
4070
4071 #[test]
4072 fn test_ingest_multiple_events() {
4073 let store = EventStore::new();
4074
4075 for i in 0..10 {
4076 let event = create_test_event(&format!("entity-{i}"), "user.created");
4077 store.ingest(&event).unwrap();
4078 }
4079
4080 assert_eq!(store.stats().total_events, 10);
4081 assert_eq!(store.stats().total_ingested, 10);
4082 }
4083
4084 #[test]
4085 fn test_query_by_entity_id() {
4086 let store = EventStore::new();
4087
4088 store
4089 .ingest(&create_test_event("entity-1", "user.created"))
4090 .unwrap();
4091 store
4092 .ingest(&create_test_event("entity-2", "user.created"))
4093 .unwrap();
4094 store
4095 .ingest(&create_test_event("entity-1", "user.updated"))
4096 .unwrap();
4097
4098 let results = store
4099 .query(&QueryEventsRequest {
4100 entity_id: Some("entity-1".to_string()),
4101 event_type: None,
4102 tenant_id: None,
4103 as_of: None,
4104 since: None,
4105 until: None,
4106 limit: None,
4107 event_type_prefix: None,
4108 exclude_event_type_prefix: None,
4109 payload_filter: None,
4110 })
4111 .unwrap();
4112
4113 assert_eq!(results.len(), 2);
4114 }
4115
4116 #[test]
4117 fn test_query_by_event_type() {
4118 let store = EventStore::new();
4119
4120 store
4121 .ingest(&create_test_event("entity-1", "user.created"))
4122 .unwrap();
4123 store
4124 .ingest(&create_test_event("entity-2", "user.updated"))
4125 .unwrap();
4126 store
4127 .ingest(&create_test_event("entity-3", "user.created"))
4128 .unwrap();
4129
4130 let results = store
4131 .query(&QueryEventsRequest {
4132 entity_id: None,
4133 event_type: Some("user.created".to_string()),
4134 tenant_id: None,
4135 as_of: None,
4136 since: None,
4137 until: None,
4138 limit: None,
4139 event_type_prefix: None,
4140 exclude_event_type_prefix: None,
4141 payload_filter: None,
4142 })
4143 .unwrap();
4144
4145 assert_eq!(results.len(), 2);
4146 }
4147
4148 #[test]
4149 fn test_query_with_limit() {
4150 let store = EventStore::new();
4151
4152 for i in 0..10 {
4153 let event = create_test_event(&format!("entity-{i}"), "user.created");
4154 store.ingest(&event).unwrap();
4155 }
4156
4157 let results = store
4158 .query(&QueryEventsRequest {
4159 entity_id: None,
4160 event_type: None,
4161 tenant_id: None,
4162 as_of: None,
4163 since: None,
4164 until: None,
4165 limit: Some(5),
4166 event_type_prefix: None,
4167 exclude_event_type_prefix: None,
4168 payload_filter: None,
4169 })
4170 .unwrap();
4171
4172 assert_eq!(results.len(), 5);
4173 }
4174
4175 #[test]
4176 fn test_query_empty_store() {
4177 let store = EventStore::new();
4178
4179 let results = store
4180 .query(&QueryEventsRequest {
4181 entity_id: Some("non-existent".to_string()),
4182 event_type: None,
4183 tenant_id: None,
4184 as_of: None,
4185 since: None,
4186 until: None,
4187 limit: None,
4188 event_type_prefix: None,
4189 exclude_event_type_prefix: None,
4190 payload_filter: None,
4191 })
4192 .unwrap();
4193
4194 assert!(results.is_empty());
4195 }
4196
4197 #[test]
4198 fn test_reconstruct_state() {
4199 let store = EventStore::new();
4200
4201 store
4202 .ingest(&create_test_event("entity-1", "user.created"))
4203 .unwrap();
4204
4205 let state = store.reconstruct_state("entity-1", None).unwrap();
4206 assert_eq!(state["current_state"]["name"], "Test");
4208 assert_eq!(state["current_state"]["value"], 42);
4209 }
4210
4211 #[test]
4212 fn test_reconstruct_state_not_found() {
4213 let store = EventStore::new();
4214
4215 let result = store.reconstruct_state("non-existent", None);
4216 assert!(result.is_err());
4217 }
4218
4219 #[test]
4220 fn test_get_snapshot_empty() {
4221 let store = EventStore::new();
4222
4223 let result = store.get_snapshot("non-existent");
4224 assert!(result.is_err());
4226 }
4227
4228 #[test]
4229 fn test_create_snapshot() {
4230 let store = EventStore::new();
4231
4232 store
4233 .ingest(&create_test_event("entity-1", "user.created"))
4234 .unwrap();
4235
4236 store.create_snapshot("entity-1").unwrap();
4237
4238 let snapshot = store.get_snapshot("entity-1").unwrap();
4240 assert_ne!(snapshot, serde_json::json!(null));
4241 }
4242
4243 #[test]
4244 fn test_create_snapshot_entity_not_found() {
4245 let store = EventStore::new();
4246
4247 let result = store.create_snapshot("non-existent");
4248 assert!(result.is_err());
4249 }
4250
4251 #[test]
4252 fn test_websocket_manager() {
4253 let store = EventStore::new();
4254 let manager = store.websocket_manager();
4255 assert!(Arc::strong_count(&manager) >= 1);
4257 }
4258
4259 #[test]
4260 fn test_snapshot_manager() {
4261 let store = EventStore::new();
4262 let manager = store.snapshot_manager();
4263 assert!(Arc::strong_count(&manager) >= 1);
4264 }
4265
4266 #[test]
4267 fn test_compaction_manager_none() {
4268 let store = EventStore::new();
4269 assert!(store.compaction_manager().is_none());
4271 }
4272
4273 #[test]
4274 fn test_schema_registry() {
4275 let store = EventStore::new();
4276 let registry = store.schema_registry();
4277 assert!(Arc::strong_count(®istry) >= 1);
4278 }
4279
4280 #[test]
4281 fn test_replay_manager() {
4282 let store = EventStore::new();
4283 let manager = store.replay_manager();
4284 assert!(Arc::strong_count(&manager) >= 1);
4285 }
4286
4287 #[test]
4288 fn test_pipeline_manager() {
4289 let store = EventStore::new();
4290 let manager = store.pipeline_manager();
4291 assert!(Arc::strong_count(&manager) >= 1);
4292 }
4293
4294 #[test]
4295 fn test_projection_manager() {
4296 let store = EventStore::new();
4297 let manager = store.projection_manager();
4298 let projections = manager.list_projections();
4300 assert!(projections.len() >= 2); }
4302
4303 #[test]
4304 fn test_projection_state_cache() {
4305 let store = EventStore::new();
4306 let cache = store.projection_state_cache();
4307
4308 cache.insert("test:key".to_string(), serde_json::json!({"value": 123}));
4309 assert_eq!(cache.len(), 1);
4310
4311 let value = cache.get("test:key").unwrap();
4312 assert_eq!(value["value"], 123);
4313 }
4314
4315 #[test]
4316 fn test_metrics() {
4317 let store = EventStore::new();
4318 let metrics = store.metrics();
4319 assert!(Arc::strong_count(&metrics) >= 1);
4320 }
4321
4322 #[test]
4323 fn test_store_stats() {
4324 let store = EventStore::new();
4325
4326 store
4327 .ingest(&create_test_event("entity-1", "user.created"))
4328 .unwrap();
4329 store
4330 .ingest(&create_test_event("entity-2", "order.placed"))
4331 .unwrap();
4332
4333 let stats = store.stats();
4334 assert_eq!(stats.total_events, 2);
4335 assert_eq!(stats.total_entities, 2);
4336 assert_eq!(stats.total_event_types, 2);
4337 assert_eq!(stats.total_ingested, 2);
4338 }
4339
4340 #[test]
4341 fn test_event_store_config_default() {
4342 let config = EventStoreConfig::default();
4343 assert!(config.storage_dir.is_none());
4344 assert!(config.wal_dir.is_none());
4345 }
4346
4347 #[test]
4348 fn test_event_store_config_with_persistence() {
4349 let temp_dir = TempDir::new().unwrap();
4350 let config = EventStoreConfig::with_persistence(temp_dir.path());
4351
4352 assert!(config.storage_dir.is_some());
4353 assert!(config.wal_dir.is_none());
4354 }
4355
4356 #[test]
4357 fn test_event_store_config_with_wal() {
4358 let temp_dir = TempDir::new().unwrap();
4359 let config = EventStoreConfig::with_wal(temp_dir.path(), WALConfig::default());
4360
4361 assert!(config.storage_dir.is_none());
4362 assert!(config.wal_dir.is_some());
4363 }
4364
4365 #[test]
4366 fn test_event_store_config_with_all() {
4367 let temp_dir = TempDir::new().unwrap();
4368 let config = EventStoreConfig::with_all(temp_dir.path(), SnapshotConfig::default());
4369
4370 assert!(config.storage_dir.is_some());
4371 }
4372
4373 #[test]
4374 fn test_event_store_config_production() {
4375 let storage_dir = TempDir::new().unwrap();
4376 let wal_dir = TempDir::new().unwrap();
4377 let config = EventStoreConfig::production(
4378 storage_dir.path(),
4379 wal_dir.path(),
4380 SnapshotConfig::default(),
4381 WALConfig::default(),
4382 CompactionConfig::default(),
4383 );
4384
4385 assert!(config.storage_dir.is_some());
4386 assert!(config.wal_dir.is_some());
4387 }
4388
4389 #[test]
4395 fn test_from_env_vars_data_dir_enables_full_persistence() {
4396 let (config, mode) = EventStoreConfig::from_env_vars(
4397 Some("/app/data".to_string()),
4398 None,
4399 None,
4400 None,
4401 None,
4402 None,
4403 None,
4404 None,
4405 );
4406 assert_eq!(mode, "wal+parquet");
4407 assert_eq!(
4408 config.storage_dir.unwrap().to_str().unwrap(),
4409 "/app/data/storage"
4410 );
4411 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/app/data/wal");
4412 }
4413
4414 #[test]
4415 fn test_from_env_vars_explicit_dirs() {
4416 let (config, mode) = EventStoreConfig::from_env_vars(
4417 None,
4418 Some("/custom/storage".to_string()),
4419 Some("/custom/wal".to_string()),
4420 None,
4421 None,
4422 None,
4423 None,
4424 None,
4425 );
4426 assert_eq!(mode, "wal+parquet");
4427 assert_eq!(
4428 config.storage_dir.unwrap().to_str().unwrap(),
4429 "/custom/storage"
4430 );
4431 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/custom/wal");
4432 }
4433
4434 #[test]
4435 fn test_from_env_vars_wal_disabled() {
4436 let (config, mode) = EventStoreConfig::from_env_vars(
4437 Some("/app/data".to_string()),
4438 None,
4439 None,
4440 Some("false".to_string()),
4441 None,
4442 None,
4443 None,
4444 None,
4445 );
4446 assert_eq!(mode, "parquet-only");
4447 assert!(config.storage_dir.is_some());
4448 assert!(config.wal_dir.is_none());
4449 }
4450
4451 #[test]
4452 fn test_from_env_vars_no_dirs_is_in_memory() {
4453 let (config, mode) =
4454 EventStoreConfig::from_env_vars(None, None, None, None, None, None, None, None);
4455 assert_eq!(mode, "in-memory");
4456 assert!(config.storage_dir.is_none());
4457 assert!(config.wal_dir.is_none());
4458 }
4459
4460 #[test]
4461 fn test_from_env_vars_empty_strings_treated_as_none() {
4462 let (_, mode) = EventStoreConfig::from_env_vars(
4463 Some(String::new()),
4464 Some(String::new()),
4465 Some(String::new()),
4466 None,
4467 None,
4468 None,
4469 None,
4470 None,
4471 );
4472 assert_eq!(mode, "in-memory");
4473 }
4474
4475 #[test]
4476 fn test_from_env_vars_explicit_overrides_data_dir() {
4477 let (config, mode) = EventStoreConfig::from_env_vars(
4478 Some("/app/data".to_string()),
4479 Some("/override/storage".to_string()),
4480 Some("/override/wal".to_string()),
4481 None,
4482 None,
4483 None,
4484 None,
4485 None,
4486 );
4487 assert_eq!(mode, "wal+parquet");
4488 assert_eq!(
4489 config.storage_dir.unwrap().to_str().unwrap(),
4490 "/override/storage"
4491 );
4492 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/override/wal");
4493 }
4494
4495 #[test]
4496 fn test_from_env_vars_wal_only() {
4497 let (config, mode) = EventStoreConfig::from_env_vars(
4498 None,
4499 None,
4500 Some("/wal/only".to_string()),
4501 None,
4502 None,
4503 None,
4504 None,
4505 None,
4506 );
4507 assert_eq!(mode, "wal-only");
4508 assert!(config.storage_dir.is_none());
4509 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/wal/only");
4510 }
4511
4512 #[test]
4513 fn test_from_env_vars_cache_bytes_parses_decimal() {
4514 let (config, _) = EventStoreConfig::from_env_vars(
4515 Some("/app/data".to_string()),
4516 None,
4517 None,
4518 None,
4519 Some("536870912".to_string()),
4520 None,
4522 None,
4523 None,
4524 );
4525 assert_eq!(config.cache_byte_budget, Some(536_870_912));
4526 }
4527
4528 #[test]
4529 fn test_from_env_vars_cache_bytes_unparseable_disables_budget() {
4530 let (config, _) = EventStoreConfig::from_env_vars(
4534 Some("/app/data".to_string()),
4535 None,
4536 None,
4537 None,
4538 Some("not-a-number".to_string()),
4539 None,
4540 None,
4541 None,
4542 );
4543 assert_eq!(config.cache_byte_budget, None);
4544 }
4545
4546 #[test]
4547 fn test_from_env_vars_cache_bytes_empty_disables_budget() {
4548 let (config, _) = EventStoreConfig::from_env_vars(
4549 Some("/app/data".to_string()),
4550 None,
4551 None,
4552 None,
4553 Some(String::new()),
4554 None,
4555 None,
4556 None,
4557 );
4558 assert_eq!(config.cache_byte_budget, None);
4559 }
4560
4561 #[test]
4562 fn test_from_env_vars_snapshot_interval_overrides_default() {
4563 let (config, _) = EventStoreConfig::from_env_vars(
4567 Some("/app/data".to_string()),
4568 None,
4569 None,
4570 None,
4571 None,
4572 Some("60".to_string()),
4573 None,
4574 None,
4575 );
4576 assert_eq!(config.compaction_config.compaction_interval_seconds, 60);
4577 }
4578
4579 #[test]
4580 fn test_from_env_vars_snapshot_interval_default_is_hourly() {
4581 let (config, _) = EventStoreConfig::from_env_vars(
4582 Some("/app/data".to_string()),
4583 None,
4584 None,
4585 None,
4586 None,
4587 None,
4588 None,
4589 None,
4590 );
4591 assert_eq!(config.compaction_config.compaction_interval_seconds, 3600);
4592 }
4593
4594 #[test]
4595 fn test_from_env_vars_snapshot_interval_unparseable_falls_back() {
4596 let (config, _) = EventStoreConfig::from_env_vars(
4597 Some("/app/data".to_string()),
4598 None,
4599 None,
4600 None,
4601 None,
4602 Some("not-a-number".to_string()),
4603 None,
4604 None,
4605 );
4606 assert_eq!(config.compaction_config.compaction_interval_seconds, 3600);
4607 }
4608
4609 #[test]
4610 fn test_from_env_vars_retention_system_days_overrides_default() {
4611 let (config, _) = EventStoreConfig::from_env_vars(
4614 Some("/app/data".to_string()),
4615 None,
4616 None,
4617 None,
4618 None,
4619 None,
4620 Some("7".to_string()),
4621 None,
4622 );
4623 let ttl = config
4624 .compaction_config
4625 .retention
4626 .ttl_for("system")
4627 .unwrap();
4628 assert_eq!(ttl.as_secs(), 7 * 24 * 3600);
4629 }
4630
4631 #[test]
4632 fn test_from_env_vars_retention_default_is_30_days_for_system() {
4633 let (config, _) = EventStoreConfig::from_env_vars(
4634 Some("/app/data".to_string()),
4635 None,
4636 None,
4637 None,
4638 None,
4639 None,
4640 None,
4641 None,
4642 );
4643 let ttl = config
4644 .compaction_config
4645 .retention
4646 .ttl_for("system")
4647 .unwrap();
4648 assert_eq!(ttl.as_secs(), 30 * 24 * 3600);
4649 assert!(config.compaction_config.retention.ttl_for("acme").is_none());
4651 }
4652
4653 #[test]
4654 fn test_store_stats_serde() {
4655 let stats = StoreStats {
4656 total_events: 100,
4657 total_entities: 50,
4658 total_event_types: 10,
4659 total_ingested: 100,
4660 };
4661
4662 let json = serde_json::to_string(&stats).unwrap();
4663 assert!(json.contains("\"total_events\":100"));
4664 assert!(json.contains("\"total_entities\":50"));
4665 }
4666
4667 #[test]
4668 fn test_query_with_entity_and_type() {
4669 let store = EventStore::new();
4670
4671 store
4672 .ingest(&create_test_event("entity-1", "user.created"))
4673 .unwrap();
4674 store
4675 .ingest(&create_test_event("entity-1", "user.updated"))
4676 .unwrap();
4677 store
4678 .ingest(&create_test_event("entity-2", "user.created"))
4679 .unwrap();
4680
4681 let results = store
4682 .query(&QueryEventsRequest {
4683 entity_id: Some("entity-1".to_string()),
4684 event_type: Some("user.created".to_string()),
4685 tenant_id: None,
4686 as_of: None,
4687 since: None,
4688 until: None,
4689 limit: None,
4690 event_type_prefix: None,
4691 exclude_event_type_prefix: None,
4692 payload_filter: None,
4693 })
4694 .unwrap();
4695
4696 assert_eq!(results.len(), 1);
4697 assert_eq!(results[0].event_type_str(), "user.created");
4698 }
4699
4700 #[test]
4701 fn test_query_by_event_type_prefix() {
4702 let store = EventStore::new();
4703
4704 store
4706 .ingest(&create_test_event("entity-1", "index.created"))
4707 .unwrap();
4708 store
4709 .ingest(&create_test_event("entity-2", "index.updated"))
4710 .unwrap();
4711 store
4712 .ingest(&create_test_event("entity-3", "trade.created"))
4713 .unwrap();
4714 store
4715 .ingest(&create_test_event("entity-4", "trade.completed"))
4716 .unwrap();
4717 store
4718 .ingest(&create_test_event("entity-5", "balance.updated"))
4719 .unwrap();
4720
4721 let results = store
4723 .query(&QueryEventsRequest {
4724 entity_id: None,
4725 event_type: None,
4726 tenant_id: None,
4727 as_of: None,
4728 since: None,
4729 until: None,
4730 limit: None,
4731 event_type_prefix: Some("index.".to_string()),
4732 exclude_event_type_prefix: None,
4733 payload_filter: None,
4734 })
4735 .unwrap();
4736
4737 assert_eq!(results.len(), 2);
4738 assert!(
4739 results
4740 .iter()
4741 .all(|e| e.event_type_str().starts_with("index."))
4742 );
4743 }
4744
4745 #[test]
4746 fn test_query_by_event_type_prefix_empty_returns_all() {
4747 let store = EventStore::new();
4748
4749 store
4750 .ingest(&create_test_event("entity-1", "index.created"))
4751 .unwrap();
4752 store
4753 .ingest(&create_test_event("entity-2", "trade.created"))
4754 .unwrap();
4755
4756 let results = store
4758 .query(&QueryEventsRequest {
4759 entity_id: None,
4760 event_type: None,
4761 tenant_id: None,
4762 as_of: None,
4763 since: None,
4764 until: None,
4765 limit: None,
4766 event_type_prefix: Some(String::new()),
4767 exclude_event_type_prefix: None,
4768 payload_filter: None,
4769 })
4770 .unwrap();
4771
4772 assert_eq!(results.len(), 2);
4773 }
4774
4775 #[test]
4776 fn test_query_by_event_type_prefix_no_match() {
4777 let store = EventStore::new();
4778
4779 store
4780 .ingest(&create_test_event("entity-1", "index.created"))
4781 .unwrap();
4782
4783 let results = store
4784 .query(&QueryEventsRequest {
4785 entity_id: None,
4786 event_type: None,
4787 tenant_id: None,
4788 as_of: None,
4789 since: None,
4790 until: None,
4791 limit: None,
4792 event_type_prefix: Some("nonexistent.".to_string()),
4793 exclude_event_type_prefix: None,
4794 payload_filter: None,
4795 })
4796 .unwrap();
4797
4798 assert!(results.is_empty());
4799 }
4800
4801 #[test]
4802 fn test_query_by_entity_with_type_prefix() {
4803 let store = EventStore::new();
4804
4805 store
4806 .ingest(&create_test_event("entity-1", "index.created"))
4807 .unwrap();
4808 store
4809 .ingest(&create_test_event("entity-1", "trade.created"))
4810 .unwrap();
4811 store
4812 .ingest(&create_test_event("entity-2", "index.updated"))
4813 .unwrap();
4814
4815 let results = store
4817 .query(&QueryEventsRequest {
4818 entity_id: Some("entity-1".to_string()),
4819 event_type: None,
4820 tenant_id: None,
4821 as_of: None,
4822 since: None,
4823 until: None,
4824 limit: None,
4825 event_type_prefix: Some("index.".to_string()),
4826 exclude_event_type_prefix: None,
4827 payload_filter: None,
4828 })
4829 .unwrap();
4830
4831 assert_eq!(results.len(), 1);
4832 assert_eq!(results[0].event_type_str(), "index.created");
4833 }
4834
4835 #[test]
4836 fn test_query_prefix_with_limit() {
4837 let store = EventStore::new();
4838
4839 for i in 0..5 {
4840 store
4841 .ingest(&create_test_event(&format!("entity-{i}"), "index.created"))
4842 .unwrap();
4843 }
4844
4845 let results = store
4846 .query(&QueryEventsRequest {
4847 entity_id: None,
4848 event_type: None,
4849 tenant_id: None,
4850 as_of: None,
4851 since: None,
4852 until: None,
4853 limit: Some(3),
4854 event_type_prefix: Some("index.".to_string()),
4855 exclude_event_type_prefix: None,
4856 payload_filter: None,
4857 })
4858 .unwrap();
4859
4860 assert_eq!(results.len(), 3);
4861 }
4862
4863 #[test]
4864 fn test_query_prefix_alongside_existing_filters() {
4865 let store = EventStore::new();
4866
4867 store
4868 .ingest(&create_test_event("entity-1", "index.created"))
4869 .unwrap();
4870 std::thread::sleep(std::time::Duration::from_millis(10));
4872 store
4873 .ingest(&create_test_event("entity-2", "index.strategy.updated"))
4874 .unwrap();
4875 std::thread::sleep(std::time::Duration::from_millis(10));
4876 store
4877 .ingest(&create_test_event("entity-3", "index.deleted"))
4878 .unwrap();
4879
4880 let results = store
4882 .query(&QueryEventsRequest {
4883 entity_id: None,
4884 event_type: None,
4885 tenant_id: None,
4886 as_of: None,
4887 since: None,
4888 until: None,
4889 limit: Some(2),
4890 event_type_prefix: Some("index.".to_string()),
4891 exclude_event_type_prefix: None,
4892 payload_filter: None,
4893 })
4894 .unwrap();
4895
4896 assert_eq!(results.len(), 2);
4897 }
4898
4899 #[test]
4900 fn test_query_with_payload_filter() {
4901 let store = EventStore::new();
4902
4903 for i in 0..5 {
4905 store
4906 .ingest(&create_test_event_with_payload(
4907 &format!("entity-{i}"),
4908 "user.action",
4909 serde_json::json!({"user_id": "alice", "action": "click"}),
4910 ))
4911 .unwrap();
4912 }
4913 for i in 5..10 {
4915 store
4916 .ingest(&create_test_event_with_payload(
4917 &format!("entity-{i}"),
4918 "user.action",
4919 serde_json::json!({"user_id": "bob", "action": "view"}),
4920 ))
4921 .unwrap();
4922 }
4923
4924 let results = store
4926 .query(&QueryEventsRequest {
4927 entity_id: None,
4928 event_type: Some("user.action".to_string()),
4929 tenant_id: None,
4930 as_of: None,
4931 since: None,
4932 until: None,
4933 limit: None,
4934 event_type_prefix: None,
4935 exclude_event_type_prefix: None,
4936 payload_filter: Some(r#"{"user_id":"alice"}"#.to_string()),
4937 })
4938 .unwrap();
4939
4940 assert_eq!(results.len(), 5);
4941 }
4942
4943 #[test]
4944 fn test_query_payload_filter_non_existent_field() {
4945 let store = EventStore::new();
4946
4947 store
4948 .ingest(&create_test_event_with_payload(
4949 "entity-1",
4950 "user.action",
4951 serde_json::json!({"user_id": "alice"}),
4952 ))
4953 .unwrap();
4954
4955 let results = store
4957 .query(&QueryEventsRequest {
4958 entity_id: None,
4959 event_type: None,
4960 tenant_id: None,
4961 as_of: None,
4962 since: None,
4963 until: None,
4964 limit: None,
4965 event_type_prefix: None,
4966 exclude_event_type_prefix: None,
4967 payload_filter: Some(r#"{"nonexistent":"value"}"#.to_string()),
4968 })
4969 .unwrap();
4970
4971 assert!(results.is_empty());
4972 }
4973
4974 #[test]
4975 fn test_query_payload_filter_with_prefix() {
4976 let store = EventStore::new();
4977
4978 store
4979 .ingest(&create_test_event_with_payload(
4980 "entity-1",
4981 "index.created",
4982 serde_json::json!({"status": "active"}),
4983 ))
4984 .unwrap();
4985 store
4986 .ingest(&create_test_event_with_payload(
4987 "entity-2",
4988 "index.created",
4989 serde_json::json!({"status": "inactive"}),
4990 ))
4991 .unwrap();
4992 store
4993 .ingest(&create_test_event_with_payload(
4994 "entity-3",
4995 "trade.created",
4996 serde_json::json!({"status": "active"}),
4997 ))
4998 .unwrap();
4999
5000 let results = store
5002 .query(&QueryEventsRequest {
5003 entity_id: None,
5004 event_type: None,
5005 tenant_id: None,
5006 as_of: None,
5007 since: None,
5008 until: None,
5009 limit: None,
5010 event_type_prefix: Some("index.".to_string()),
5011 exclude_event_type_prefix: None,
5012 payload_filter: Some(r#"{"status":"active"}"#.to_string()),
5013 })
5014 .unwrap();
5015
5016 assert_eq!(results.len(), 1);
5017 assert_eq!(results[0].entity_id().to_string(), "entity-1");
5018 }
5019
5020 #[test]
5021 fn test_flush_storage_no_storage() {
5022 let store = EventStore::new();
5023 let result = store.flush_storage();
5025 assert!(result.is_ok());
5026 }
5027
5028 #[test]
5029 fn test_state_evolution() {
5030 let store = EventStore::new();
5031
5032 store
5034 .ingest(
5035 &Event::from_strings(
5036 "user.created".to_string(),
5037 "user-1".to_string(),
5038 "default".to_string(),
5039 serde_json::json!({"name": "Alice", "age": 25}),
5040 None,
5041 )
5042 .unwrap(),
5043 )
5044 .unwrap();
5045
5046 store
5048 .ingest(
5049 &Event::from_strings(
5050 "user.updated".to_string(),
5051 "user-1".to_string(),
5052 "default".to_string(),
5053 serde_json::json!({"age": 26}),
5054 None,
5055 )
5056 .unwrap(),
5057 )
5058 .unwrap();
5059
5060 let state = store.reconstruct_state("user-1", None).unwrap();
5061 assert_eq!(state["current_state"]["name"], "Alice");
5063 assert_eq!(state["current_state"]["age"], 26);
5064 }
5065
5066 #[test]
5067 fn test_reject_system_event_types() {
5068 let store = EventStore::new();
5069
5070 let event = Event::reconstruct_from_strings(
5072 uuid::Uuid::new_v4(),
5073 "_system.tenant.created".to_string(),
5074 "_system:tenant:acme".to_string(),
5075 "_system".to_string(),
5076 serde_json::json!({"name": "ACME"}),
5077 chrono::Utc::now(),
5078 None,
5079 1,
5080 );
5081
5082 let result = store.ingest(&event);
5083 assert!(result.is_err());
5084 let err = result.unwrap_err();
5085 assert!(
5086 err.to_string().contains("reserved for internal use"),
5087 "Expected system namespace rejection, got: {err}"
5088 );
5089 }
5090
5091 #[test]
5099 fn test_wal_recovery_checkpoints_to_parquet() {
5100 let data_dir = TempDir::new().unwrap();
5101 let storage_dir = data_dir.path().join("storage");
5102 let wal_dir = data_dir.path().join("wal");
5103
5104 {
5106 let config = EventStoreConfig::production(
5107 &storage_dir,
5108 &wal_dir,
5109 SnapshotConfig::default(),
5110 WALConfig {
5111 sync_on_write: true,
5112 ..WALConfig::default()
5113 },
5114 CompactionConfig::default(),
5115 );
5116 let store = EventStore::with_config(config);
5117
5118 for i in 0..5 {
5119 let event = Event::from_strings(
5120 "test.created".to_string(),
5121 format!("entity-{i}"),
5122 "default".to_string(),
5123 serde_json::json!({"index": i}),
5124 None,
5125 )
5126 .unwrap();
5127 store.ingest(&event).unwrap();
5128 }
5129
5130 assert_eq!(store.stats().total_events, 5);
5131
5132 }
5135
5136 let wal_files: Vec<_> = std::fs::read_dir(&wal_dir)
5138 .unwrap()
5139 .filter_map(std::result::Result::ok)
5140 .filter(|e| e.path().extension().is_some_and(|ext| ext == "log"))
5141 .collect();
5142 assert!(!wal_files.is_empty(), "WAL file should exist");
5143 let wal_size = wal_files[0].metadata().unwrap().len();
5144 assert!(wal_size > 0, "WAL file should have data (got 0 bytes)");
5145
5146 {
5148 let config = EventStoreConfig::production(
5149 &storage_dir,
5150 &wal_dir,
5151 SnapshotConfig::default(),
5152 WALConfig {
5153 sync_on_write: true,
5154 ..WALConfig::default()
5155 },
5156 CompactionConfig::default(),
5157 );
5158 let store = EventStore::with_config(config);
5159
5160 assert_eq!(
5162 store.stats().total_events,
5163 5,
5164 "Session 2 should have all 5 events after WAL recovery"
5165 );
5166
5167 let parquet_files = find_parquet_files(&storage_dir);
5171 assert!(
5172 !parquet_files.is_empty(),
5173 "Parquet file should exist after WAL checkpoint"
5174 );
5175 }
5176
5177 {
5180 let config = EventStoreConfig::production(
5181 &storage_dir,
5182 &wal_dir,
5183 SnapshotConfig::default(),
5184 WALConfig {
5185 sync_on_write: true,
5186 ..WALConfig::default()
5187 },
5188 CompactionConfig::default(),
5189 );
5190 let store = EventStore::with_config(config);
5191
5192 assert_eq!(
5196 store.stats().total_events,
5197 0,
5198 "Session 3 boot should not pre-load Parquet (lazy-load mode)"
5199 );
5200
5201 store.ensure_tenant_loaded("default").unwrap();
5204 assert_eq!(
5205 store.stats().total_events,
5206 5,
5207 "Session 3 should have all 5 events after ensure_tenant_loaded"
5208 );
5209 }
5210 }
5211
5212 #[test]
5213 fn test_parquet_restore_surfaces_errors_not_silent() {
5214 let data_dir = TempDir::new().unwrap();
5218 let storage_dir = data_dir.path().join("storage");
5219 let wal_dir = data_dir.path().join("wal");
5220
5221 {
5223 let config = EventStoreConfig::production(
5224 &storage_dir,
5225 &wal_dir,
5226 SnapshotConfig::default(),
5227 WALConfig {
5228 sync_on_write: true,
5229 ..WALConfig::default()
5230 },
5231 CompactionConfig::default(),
5232 );
5233 let store = EventStore::with_config(config);
5234
5235 for i in 0..3 {
5236 let event = Event::from_strings(
5237 "test.created".to_string(),
5238 format!("entity-{i}"),
5239 "default".to_string(),
5240 serde_json::json!({"i": i}),
5241 None,
5242 )
5243 .unwrap();
5244 store.ingest(&event).unwrap();
5245 }
5246
5247 store.flush_storage().unwrap();
5248 assert_eq!(store.stats().total_events, 3);
5249 }
5250
5251 let parquet_files = find_parquet_files(&storage_dir);
5254 assert!(!parquet_files.is_empty(), "Parquet file must exist");
5255
5256 std::fs::write(&parquet_files[0], b"corrupted data").unwrap();
5258
5259 for entry in std::fs::read_dir(&wal_dir).unwrap().flatten() {
5261 std::fs::write(entry.path(), b"").unwrap();
5262 }
5263
5264 {
5271 let config = EventStoreConfig::production(
5272 &storage_dir,
5273 &wal_dir,
5274 SnapshotConfig::default(),
5275 WALConfig::default(),
5276 CompactionConfig::default(),
5277 );
5278 let store = EventStore::with_config(config);
5279
5280 assert_eq!(store.stats().total_events, 0);
5283 }
5284 }
5285
5286 fn count_wal_entries(wal_dir: &std::path::Path) -> usize {
5296 use std::io::{BufRead, BufReader};
5297 let mut total = 0usize;
5298 let Ok(entries) = std::fs::read_dir(wal_dir) else {
5299 return 0;
5300 };
5301 for entry in entries.flatten() {
5302 let path = entry.path();
5303 if path.extension().is_none_or(|e| e != "log") {
5304 continue;
5305 }
5306 let Ok(file) = std::fs::File::open(&path) else {
5307 continue;
5308 };
5309 for line in BufReader::new(file)
5310 .lines()
5311 .map_while(std::result::Result::ok)
5312 {
5313 if !line.trim().is_empty() {
5314 total += 1;
5315 }
5316 }
5317 }
5318 total
5319 }
5320
5321 #[test]
5322 fn test_checkpoint_truncates_wal_after_flush() {
5323 let data_dir = TempDir::new().unwrap();
5328 let storage_dir = data_dir.path().join("storage");
5329 let wal_dir = data_dir.path().join("wal");
5330
5331 let config = EventStoreConfig::production(
5332 &storage_dir,
5333 &wal_dir,
5334 SnapshotConfig::default(),
5335 WALConfig {
5336 sync_on_write: true,
5337 ..WALConfig::default()
5338 },
5339 CompactionConfig::default(),
5340 );
5341 let store = EventStore::with_config(config);
5342
5343 for i in 0..10 {
5344 let event = Event::from_strings(
5345 "test.created".to_string(),
5346 format!("entity-{i}"),
5347 "default".to_string(),
5348 serde_json::json!({"i": i}),
5349 None,
5350 )
5351 .unwrap();
5352 store.ingest(&event).unwrap();
5353 }
5354
5355 assert_eq!(
5357 count_wal_entries(&wal_dir),
5358 10,
5359 "WAL should have 10 events before checkpoint"
5360 );
5361
5362 store.checkpoint().unwrap();
5363
5364 assert_eq!(
5365 count_wal_entries(&wal_dir),
5366 0,
5367 "WAL should be empty after successful checkpoint"
5368 );
5369 let parquet_files = find_parquet_files(&storage_dir);
5370 assert!(!parquet_files.is_empty(), "Parquet should hold the events");
5371 }
5372
5373 #[test]
5374 fn test_replay_only_post_checkpoint_events_after_crash() {
5375 let data_dir = TempDir::new().unwrap();
5382 let storage_dir = data_dir.path().join("storage");
5383 let wal_dir = data_dir.path().join("wal");
5384
5385 let config_factory = || {
5386 EventStoreConfig::production(
5387 &storage_dir,
5388 &wal_dir,
5389 SnapshotConfig::default(),
5390 WALConfig {
5391 sync_on_write: true,
5392 ..WALConfig::default()
5393 },
5394 CompactionConfig::default(),
5395 )
5396 };
5397
5398 const N: usize = 50;
5401 const K: usize = 5;
5402 {
5403 let store = EventStore::with_config(config_factory());
5404 for i in 0..N {
5405 store
5406 .ingest(
5407 &Event::from_strings(
5408 "pre.checkpoint".to_string(),
5409 format!("e-{i}"),
5410 "default".to_string(),
5411 serde_json::json!({"i": i}),
5412 None,
5413 )
5414 .unwrap(),
5415 )
5416 .unwrap();
5417 }
5418 store.checkpoint().unwrap();
5419 assert_eq!(
5420 count_wal_entries(&wal_dir),
5421 0,
5422 "WAL should be empty immediately after checkpoint"
5423 );
5424
5425 for i in 0..K {
5426 store
5427 .ingest(
5428 &Event::from_strings(
5429 "post.checkpoint".to_string(),
5430 format!("p-{i}"),
5431 "default".to_string(),
5432 serde_json::json!({"i": i}),
5433 None,
5434 )
5435 .unwrap(),
5436 )
5437 .unwrap();
5438 }
5439 assert_eq!(
5440 count_wal_entries(&wal_dir),
5441 K,
5442 "WAL should hold only post-checkpoint events"
5443 );
5444 }
5446
5447 {
5451 let store = EventStore::with_config(config_factory());
5452 assert_eq!(
5456 store.stats().total_events,
5457 K,
5458 "Boot should replay exactly K events from WAL (the post-checkpoint window), not N+K"
5459 );
5460
5461 store.ensure_tenant_loaded("default").unwrap();
5463 assert_eq!(
5464 store.stats().total_events,
5465 N + K,
5466 "After lazy-load, both pre- and post-checkpoint events should be reachable"
5467 );
5468 }
5469 }
5470
5471 #[test]
5472 fn test_checkpoint_is_idempotent() {
5473 let data_dir = TempDir::new().unwrap();
5476 let storage_dir = data_dir.path().join("storage");
5477 let wal_dir = data_dir.path().join("wal");
5478
5479 let store = EventStore::with_config(EventStoreConfig::production(
5480 &storage_dir,
5481 &wal_dir,
5482 SnapshotConfig::default(),
5483 WALConfig::default(),
5484 CompactionConfig::default(),
5485 ));
5486
5487 for i in 0..5 {
5488 store
5489 .ingest(
5490 &Event::from_strings(
5491 "x".to_string(),
5492 format!("e-{i}"),
5493 "default".to_string(),
5494 serde_json::json!({}),
5495 None,
5496 )
5497 .unwrap(),
5498 )
5499 .unwrap();
5500 }
5501
5502 store.checkpoint().unwrap();
5503 store.checkpoint().unwrap();
5505 assert_eq!(count_wal_entries(&wal_dir), 0);
5506 }
5507
5508 #[test]
5509 fn test_checkpoint_noop_in_memory_only_mode() {
5510 let store = EventStore::new();
5512 store.checkpoint().unwrap();
5513 }
5514
5515 #[test]
5516 fn test_checkpoint_interval_from_env_defaults_to_60s_when_wal_enabled() {
5517 let (config, _) = EventStoreConfig::from_env_vars(
5518 Some("/app/data".to_string()),
5519 None,
5520 None,
5521 None,
5522 None,
5523 None,
5524 None,
5525 None,
5526 );
5527 assert_eq!(config.checkpoint_interval_secs, Some(60));
5528 }
5529
5530 #[test]
5531 fn test_checkpoint_interval_from_env_overrides_default() {
5532 let (config, _) = EventStoreConfig::from_env_vars(
5533 Some("/app/data".to_string()),
5534 None,
5535 None,
5536 None,
5537 None,
5538 None,
5539 None,
5540 Some("15".to_string()),
5541 );
5542 assert_eq!(config.checkpoint_interval_secs, Some(15));
5543 }
5544
5545 #[test]
5546 fn test_checkpoint_interval_disabled_when_wal_disabled() {
5547 let (config, _) = EventStoreConfig::from_env_vars(
5549 Some("/app/data".to_string()),
5550 None,
5551 None,
5552 Some("false".to_string()),
5553 None,
5554 None,
5555 None,
5556 Some("15".to_string()),
5557 );
5558 assert_eq!(config.checkpoint_interval_secs, None);
5559 }
5560
5561 #[test]
5562 fn test_checkpoint_interval_unparseable_falls_back_to_default() {
5563 let (config, _) = EventStoreConfig::from_env_vars(
5564 Some("/app/data".to_string()),
5565 None,
5566 None,
5567 None,
5568 None,
5569 None,
5570 None,
5571 Some("not-a-number".to_string()),
5572 );
5573 assert_eq!(config.checkpoint_interval_secs, Some(60));
5574 }
5575
5576 fn seed_two_tenants() -> EventStore {
5583 let store = EventStore::new();
5584
5585 for (entity, etype, payload) in [
5587 (
5588 "a-1",
5589 "created",
5590 serde_json::json!({"colour": "red", "size": 1}),
5591 ),
5592 ("a-1", "updated", serde_json::json!({"colour": "blue"})),
5593 ("a-2", "created", serde_json::json!({"colour": "green"})),
5594 ] {
5595 store
5596 .ingest(
5597 &Event::from_strings(
5598 etype.to_string(),
5599 entity.to_string(),
5600 "alice".to_string(),
5601 payload,
5602 None,
5603 )
5604 .unwrap(),
5605 )
5606 .unwrap();
5607 }
5608
5609 store
5611 .ingest(
5612 &Event::from_strings(
5613 "created".to_string(),
5614 "a-1".to_string(),
5615 "bob".to_string(),
5616 serde_json::json!({"colour": "BOB_SECRET", "bob_only": true}),
5617 None,
5618 )
5619 .unwrap(),
5620 )
5621 .unwrap();
5622
5623 store
5624 }
5625
5626 #[test]
5627 fn test_stats_for_tenant_counts_only_that_tenant() {
5628 let store = seed_two_tenants();
5629
5630 let alice = store.stats_for_tenant("alice");
5631 assert_eq!(alice.total_events, 3);
5632 assert_eq!(alice.total_entities, 2);
5633 assert_eq!(alice.total_event_types, 2);
5634 assert_eq!(alice.event_types.get("created"), Some(&2));
5635 assert_eq!(alice.event_types.get("updated"), Some(&1));
5636
5637 let bob = store.stats_for_tenant("bob");
5638 assert_eq!(bob.total_events, 1);
5639 assert_eq!(bob.total_entities, 1);
5640 assert_eq!(bob.total_event_types, 1);
5641 }
5642
5643 #[test]
5644 fn test_stats_for_tenant_never_reports_global_totals() {
5645 let store = seed_two_tenants();
5646
5647 assert_eq!(store.stats().total_events, 4);
5649 assert_eq!(store.stats_for_tenant("alice").total_events, 3);
5650 assert_eq!(store.stats_for_tenant("bob").total_events, 1);
5651
5652 assert_eq!(store.stats_for_tenant("bob").total_ingested, 1);
5655 }
5656
5657 #[test]
5658 fn test_stats_for_tenant_unknown_tenant_is_empty_not_global() {
5659 let store = seed_two_tenants();
5660 let nobody = store.stats_for_tenant("does-not-exist");
5661
5662 assert_eq!(nobody.total_events, 0);
5663 assert_eq!(nobody.total_entities, 0);
5664 assert!(nobody.event_types.is_empty());
5665 assert!(nobody.oldest_event.is_none());
5666 assert!(nobody.newest_event.is_none());
5667 }
5668
5669 #[test]
5670 fn test_stats_for_tenant_reports_time_range() {
5671 let store = seed_two_tenants();
5672 let alice = store.stats_for_tenant("alice");
5673
5674 let oldest = alice.oldest_event.expect("oldest");
5675 let newest = alice.newest_event.expect("newest");
5676 assert!(oldest <= newest);
5677 }
5678
5679 #[test]
5680 fn test_reconstruct_state_for_tenant_isolates_shared_entity_id() {
5681 let store = seed_two_tenants();
5682
5683 let alice = store
5685 .reconstruct_state_for_tenant("a-1", None, "alice")
5686 .unwrap();
5687 let alice_state = alice.get("current_state").unwrap();
5688 assert_eq!(alice_state.get("colour").unwrap(), "blue"); assert_eq!(alice_state.get("size").unwrap(), 1);
5690 assert!(
5691 alice_state.get("bob_only").is_none(),
5692 "alice must not see bob's payload keys: {alice_state:?}"
5693 );
5694 assert_eq!(alice.get("event_count").unwrap(), 2);
5695
5696 let bob = store
5697 .reconstruct_state_for_tenant("a-1", None, "bob")
5698 .unwrap();
5699 let bob_state = bob.get("current_state").unwrap();
5700 assert_eq!(bob_state.get("colour").unwrap(), "BOB_SECRET");
5701 assert_eq!(bob.get("event_count").unwrap(), 1);
5702 }
5703
5704 #[test]
5705 fn test_reconstruct_state_for_tenant_rejects_another_tenants_entity() {
5706 let store = seed_two_tenants();
5707
5708 assert!(
5710 store
5711 .reconstruct_state_for_tenant("a-2", None, "alice")
5712 .is_ok()
5713 );
5714 assert!(
5715 store
5716 .reconstruct_state_for_tenant("a-2", None, "bob")
5717 .is_err(),
5718 "bob must not be able to read alice's entity"
5719 );
5720 }
5721
5722 #[test]
5723 fn test_global_reconstruct_state_still_spans_tenants() {
5724 let store = seed_two_tenants();
5727 let all = store.reconstruct_state("a-1", None).unwrap();
5728 assert_eq!(all.get("event_count").unwrap(), 3);
5729 }
5730
5731 fn seed_hot_entity(history: usize) -> EventStore {
5742 let store = EventStore::new();
5743 for _ in 0..history {
5744 store
5745 .ingest(&create_test_event("entity-hot", "user.updated"))
5746 .unwrap();
5747 }
5748 store
5749 }
5750
5751 fn hot_request(limit: Option<usize>) -> QueryEventsRequest {
5752 QueryEventsRequest {
5753 entity_id: Some("entity-hot".to_string()),
5754 limit,
5755 ..QueryEventsRequest::default()
5756 }
5757 }
5758
5759 #[test]
5760 fn query_window_limit_bounds_materialization_not_just_the_response() {
5761 const HISTORY: usize = 400;
5762 let store = seed_hot_entity(HISTORY);
5763
5764 let ((page, total), materialized) = crate::clone_probe::measure(|| {
5766 store.query_window(&hot_request(Some(1)), 0, true).unwrap()
5767 });
5768 assert_eq!(page.len(), 1, "limit=1 returns one event");
5769 assert_eq!(total, HISTORY, "total is still the full match count");
5770 assert_eq!(
5771 materialized, 1,
5772 "limit=1 cloned {materialized} of {HISTORY} events: `limit` must \
5773 bound what the store materializes, not just what it returns"
5774 );
5775
5776 let ((page, _), materialized) = crate::clone_probe::measure(|| {
5778 store
5779 .query_window(&hot_request(Some(5)), 300, false)
5780 .unwrap()
5781 });
5782 assert_eq!(page.len(), 5);
5783 assert_eq!(
5784 materialized, 5,
5785 "offset=300&limit=5 cloned {materialized} events, expected 5"
5786 );
5787
5788 let ((all, _), materialized) = crate::clone_probe::measure(|| {
5791 store.query_window(&hot_request(None), 0, true).unwrap()
5792 });
5793 assert_eq!(all.len(), HISTORY);
5794 assert_eq!(materialized, HISTORY as u64);
5795 }
5796
5797 #[test]
5798 fn query_window_cost_does_not_grow_with_entity_history() {
5799 let short = seed_hot_entity(40);
5804 let long = seed_hot_entity(400);
5805
5806 let (_, short_cost) = crate::clone_probe::measure(|| {
5807 short.query_window(&hot_request(Some(1)), 0, true).unwrap()
5808 });
5809 let (_, long_cost) = crate::clone_probe::measure(|| {
5810 long.query_window(&hot_request(Some(1)), 0, true).unwrap()
5811 });
5812
5813 assert_eq!(
5814 (short_cost, long_cost),
5815 (1, 1),
5816 "a limit=1 page materialized {short_cost} events over a 40-event \
5817 history and {long_cost} over 400 — cost is tracking history length"
5818 );
5819 }
5820
5821 #[test]
5822 fn select_window_orders_only_the_window() {
5823 use std::cell::Cell;
5829 const N: usize = 4096;
5830
5831 let comparisons = Cell::new(0usize);
5832 let order = |a: &u64, b: &u64| {
5833 comparisons.set(comparisons.get() + 1);
5834 a.cmp(b)
5835 };
5836 let shuffled = || -> Vec<u64> {
5837 (0..N as u64)
5838 .map(|i| (i * 2_654_435_761) % 1_000_003)
5839 .collect()
5840 };
5841
5842 let mut bounded = shuffled();
5843 select_window(&mut bounded, 0, Some(1), order);
5844 let bounded_cost = comparisons.replace(0);
5845 assert_eq!(bounded.len(), 1, "window of 1 keeps 1 item");
5846
5847 let mut everything = shuffled();
5848 select_window(&mut everything, 0, None, order);
5849 let full_sort_cost = comparisons.get();
5850 assert_eq!(everything.len(), N);
5851
5852 assert!(
5853 bounded_cost < 4 * N,
5854 "selecting a 1-item window out of {N} took {bounded_cost} comparisons \
5855 (~{}·N) — that is sort-shaped, not selection-shaped",
5856 bounded_cost / N
5857 );
5858 assert!(
5859 bounded_cost * 3 < full_sort_cost,
5860 "a 1-item window cost {bounded_cost} comparisons against \
5861 {full_sort_cost} for sorting all {N}: the window is not bounding \
5862 the ordering work"
5863 );
5864 }
5865
5866 #[test]
5867 fn select_window_is_equivalent_to_sort_then_window() {
5868 let order = |a: &(u32, usize), b: &(u32, usize)| a.cmp(b);
5872 let source: Vec<(u32, usize)> = [7, 3, 3, 9, 1, 3, 5, 9, 0, 2]
5873 .into_iter()
5874 .enumerate()
5875 .map(|(i, k)| (k, i))
5876 .collect();
5877
5878 let mut sorted = source.clone();
5879 sorted.sort_by(order);
5880
5881 for offset in 0..12 {
5882 for limit in [None, Some(0), Some(1), Some(3), Some(10), Some(50)] {
5883 let mut got = source.clone();
5884 select_window(&mut got, offset, limit, order);
5885 let got: Vec<_> = got.into_iter().skip(offset).collect();
5886 let expected: Vec<_> = sorted
5887 .iter()
5888 .copied()
5889 .skip(offset)
5890 .take(limit.unwrap_or(usize::MAX))
5891 .collect();
5892 assert_eq!(got, expected, "offset={offset} limit={limit:?}");
5893 }
5894 }
5895 }
5896
5897 #[test]
5908 fn query_window_applies_time_filters_on_the_full_scan_path() {
5909 let store = EventStore::new();
5910 let base = Utc::now() - chrono::Duration::hours(24);
5911 let mut ids = Vec::new();
5912 for i in 0..5i64 {
5913 let mut event = create_test_event(&format!("e-{i}"), "user.created");
5914 event.timestamp = base + chrono::Duration::hours(i);
5915 event.version = i + 1;
5916 ids.push(event.id);
5917 store.ingest(&event).unwrap();
5918 }
5919 let at = |h: i64| base + chrono::Duration::hours(h);
5920
5921 let scoped = |mutate: &dyn Fn(&mut QueryEventsRequest)| {
5924 let mut req = QueryEventsRequest {
5925 tenant_id: Some("default".to_string()),
5926 ..QueryEventsRequest::default()
5927 };
5928 mutate(&mut req);
5929 req
5930 };
5931
5932 for (label, req, expected) in [
5933 (
5934 "since=T+2 keeps only events at or after T+2",
5935 scoped(&|r| r.since = Some(at(2))),
5936 vec![ids[2], ids[3], ids[4]],
5937 ),
5938 (
5939 "until=T+1 keeps only events at or before T+1",
5940 scoped(&|r| r.until = Some(at(1))),
5941 vec![ids[0], ids[1]],
5942 ),
5943 (
5944 "as_of=T+1 is time travel: nothing newer than T+1",
5945 scoped(&|r| r.as_of = Some(at(1))),
5946 vec![ids[0], ids[1]],
5947 ),
5948 (
5949 "since+until compose into a closed window",
5950 scoped(&|r| {
5951 r.since = Some(at(1));
5952 r.until = Some(at(3));
5953 }),
5954 vec![ids[1], ids[2], ids[3]],
5955 ),
5956 ] {
5957 let (events, total) = store.query_window(&req, 0, false).unwrap();
5958 let got: Vec<_> = events.iter().map(|e| e.id).collect();
5959 assert_eq!(got, expected, "{label}");
5960 assert_eq!(
5961 total,
5962 expected.len(),
5963 "{label}: total counts the events INSIDE the window — \
5964 has_more is derived from it, so a full-history total makes a \
5965 paginator walk events the caller filtered out"
5966 );
5967 }
5968
5969 let (indexed, total) = store
5972 .query_window(
5973 &scoped(&|r| {
5974 r.entity_id = Some("e-3".to_string());
5975 r.since = Some(at(2));
5976 }),
5977 0,
5978 false,
5979 )
5980 .unwrap();
5981 assert_eq!(
5982 indexed.iter().map(|e| e.id).collect::<Vec<_>>(),
5983 vec![ids[3]]
5984 );
5985 assert_eq!(total, 1);
5986
5987 let (page, total) = store
5990 .query_window(&scoped(&|r| r.since = Some(at(2))), 1, true)
5991 .unwrap();
5992 assert_eq!(
5993 page.iter().map(|e| e.id).collect::<Vec<_>>(),
5994 vec![ids[3], ids[2]],
5995 "order=desc + offset=1 inside a since window"
5996 );
5997 assert_eq!(total, 3);
5998 }
5999
6000 #[test]
6001 fn query_window_bounded_selection_matches_a_full_sort() {
6002 let store = EventStore::new();
6007 for i in 0..50 {
6008 store
6009 .ingest(&create_test_event(&format!("e-{i:02}"), "user.created"))
6010 .unwrap();
6011 }
6012 let all = |descending: bool| {
6013 let (events, _) = store
6014 .query_window(&QueryEventsRequest::default(), 0, descending)
6015 .unwrap();
6016 events
6017 };
6018
6019 for descending in [false, true] {
6020 let reference = all(descending);
6021 for offset in [0, 1, 7, 49, 50, 100] {
6022 for limit in [1, 3, 10, 50, 100] {
6023 let (page, total) = store
6024 .query_window(
6025 &QueryEventsRequest {
6026 limit: Some(limit),
6027 ..QueryEventsRequest::default()
6028 },
6029 offset,
6030 descending,
6031 )
6032 .unwrap();
6033 let expected: Vec<_> = reference
6034 .iter()
6035 .skip(offset)
6036 .take(limit)
6037 .map(|e| e.id)
6038 .collect();
6039 let got: Vec<_> = page.iter().map(|e| e.id).collect();
6040 assert_eq!(
6041 got, expected,
6042 "desc={descending} offset={offset} limit={limit}: windowed \
6043 selection must match a full sort"
6044 );
6045 assert_eq!(total, 50, "total is always the full match count");
6046 }
6047 }
6048 }
6049 }
6050}