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> {
1172 let Some(storage) = self.storage.as_ref().map(Arc::clone) else {
1173 return Ok(0);
1174 };
1175
1176 let events = storage.read().load_all_events()?;
1177 let read_count = events.len();
1178 let tenants: Vec<String> = events
1179 .iter()
1180 .map(Event::tenant_id_str)
1181 .collect::<std::collections::HashSet<_>>()
1182 .into_iter()
1183 .map(str::to_owned)
1184 .collect();
1185 let before = self.events.read().len();
1186 for event in events {
1187 self.append_loaded_event(event);
1188 }
1189 let applied = self.events.read().len() - before;
1190
1191 *self.total_ingested.write() = self.events.read().len() as u64;
1193 for tenant in &tenants {
1197 self.tenant_loader.mark_loaded(tenant);
1198 }
1199
1200 tracing::info!(
1201 read = read_count,
1202 applied = applied,
1203 "🔄 hydrate_all_from_storage: in-memory pile reconstructed from Parquet"
1204 );
1205 Ok(applied)
1206 }
1207
1208 pub fn projection_state_cache(&self) -> Arc<DashMap<String, serde_json::Value>> {
1211 Arc::clone(&self.projection_state_cache)
1212 }
1213
1214 pub fn projection_status(&self) -> Arc<DashMap<String, String>> {
1216 Arc::clone(&self.projection_status)
1217 }
1218
1219 pub fn geo_index(&self) -> Arc<GeoIndex> {
1222 self.geo_index.clone()
1223 }
1224
1225 pub fn exactly_once(&self) -> Arc<ExactlyOnceRegistry> {
1227 self.exactly_once.clone()
1228 }
1229
1230 pub fn schema_evolution(&self) -> Arc<SchemaEvolutionManager> {
1232 self.schema_evolution.clone()
1233 }
1234
1235 pub fn snapshot_events(&self) -> Vec<Event> {
1241 self.events.read().clone()
1242 }
1243
1244 pub fn compact_entity_tokens(
1266 &self,
1267 entity_id: &str,
1268 token_event_type: &str,
1269 merged_event: Event,
1270 ) -> Result<bool> {
1271 self.ensure_writable()?;
1273
1274 {
1276 let events = self.events.read();
1277 let has_tokens = events
1278 .iter()
1279 .any(|e| e.entity_id_str() == entity_id && e.event_type_str() == token_event_type);
1280 if !has_tokens {
1281 return Ok(false);
1282 }
1283 }
1284
1285 let projections = self.projections.read();
1287 projections.process_event(&merged_event)?;
1288 drop(projections);
1289
1290 let mut events = self.events.write();
1292
1293 events.retain(|e| {
1294 !(e.entity_id_str() == entity_id && e.event_type_str() == token_event_type)
1295 });
1296
1297 events.push(merged_event);
1298
1299 self.index.clear();
1304 for (offset, event) in events.iter().enumerate() {
1305 if let Err(e) = self.index.index_event(
1306 event.id,
1307 event.entity_id_str(),
1308 event.event_type_str(),
1309 event.timestamp,
1310 offset,
1311 ) {
1312 tracing::warn!(
1313 event_id = %event.id,
1314 offset,
1315 "Failed to re-index event during compaction: {e}"
1316 );
1317 }
1318 }
1319
1320 Ok(true)
1321 }
1322
1323 #[cfg(feature = "server")]
1324 pub fn webhook_registry(&self) -> Arc<WebhookRegistry> {
1325 Arc::clone(&self.webhook_registry)
1326 }
1327
1328 #[cfg(feature = "server")]
1331 pub fn set_webhook_tx(&self, tx: mpsc::UnboundedSender<WebhookDeliveryTask>) {
1332 *self.webhook_tx.write() = Some(tx);
1333 tracing::info!("Webhook delivery channel connected");
1334 }
1335
1336 #[cfg(feature = "server")]
1338 fn dispatch_webhooks(&self, event: &Event) {
1339 let matching = self.webhook_registry.find_matching(event);
1340 if matching.is_empty() {
1341 return;
1342 }
1343
1344 let tx_guard = self.webhook_tx.read();
1345 if let Some(ref tx) = *tx_guard {
1346 for webhook in matching {
1347 let task = WebhookDeliveryTask {
1348 webhook,
1349 event: event.clone(),
1350 };
1351 if let Err(e) = tx.send(task) {
1352 tracing::warn!("Failed to queue webhook delivery: {}", e);
1353 }
1354 }
1355 }
1356 }
1357
1358 pub fn flush_storage(&self) -> Result<()> {
1360 if let Some(ref storage) = self.storage {
1361 let storage = storage.read();
1362 storage.flush()?;
1363 tracing::info!("✅ Flushed events to persistent storage");
1364 }
1365 Ok(())
1366 }
1367
1368 pub fn checkpoint(&self) -> Result<()> {
1389 let Some(ref wal) = self.wal else {
1390 #[cfg(feature = "server")]
1393 self.refresh_storage_metrics();
1394 return Ok(());
1395 };
1396
1397 if self.read_only {
1398 return Ok(());
1399 }
1400
1401 let active = {
1405 let _sealing = self.durability_gate.write();
1406 wal.seal()?
1407 };
1408 self.flush_storage()?;
1409 wal.remove_sealed(&active)?;
1410 tracing::debug!("✅ Checkpoint complete: Parquet flushed, sealed WAL segments retired");
1411
1412 #[cfg(feature = "server")]
1415 self.refresh_storage_metrics();
1416
1417 Ok(())
1418 }
1419
1420 #[cfg(feature = "server")]
1425 pub fn refresh_storage_metrics_now(&self) {
1426 self.refresh_storage_metrics();
1427 }
1428
1429 #[cfg(feature = "server")]
1446 fn refresh_storage_metrics(&self) {
1447 let Some(ref storage) = self.storage else {
1448 return;
1449 };
1450
1451 let parquet_stats = match storage.read().stats() {
1452 Ok(stats) => stats,
1453 Err(e) => {
1454 tracing::warn!("storage-size metric refresh: failed to stat Parquet: {e}");
1455 return;
1456 }
1457 };
1458
1459 let (wal_bytes, wal_segments) = match self.wal.as_ref() {
1460 Some(wal) => match wal.on_disk_stats() {
1461 Ok(stats) => stats,
1462 Err(e) => {
1463 tracing::warn!("storage-size metric refresh: failed to stat WAL: {e}");
1464 (0, 0)
1465 }
1466 },
1467 None => (0, 0),
1468 };
1469
1470 let total_bytes = parquet_stats.total_size_bytes + wal_bytes;
1471
1472 self.metrics
1473 .storage_size_bytes
1474 .set(total_bytes.min(i64::MAX as u64) as i64);
1475 self.metrics
1476 .parquet_files_total
1477 .set(parquet_stats.total_files as i64);
1478 self.metrics.wal_segments_total.set(wal_segments as i64);
1479
1480 tracing::debug!(
1481 "storage-size metrics refreshed: {} bytes total ({} Parquet files, {} WAL segments)",
1482 total_bytes,
1483 parquet_stats.total_files,
1484 wal_segments
1485 );
1486 }
1487
1488 pub fn checkpoint_interval(&self) -> Option<std::time::Duration> {
1490 self.checkpoint_interval_secs
1491 .map(std::time::Duration::from_secs)
1492 }
1493
1494 pub fn ensure_tenant_loaded(&self, tenant_id: &str) -> Result<()> {
1518 if self.tenant_loader.is_loaded(tenant_id) {
1520 return Ok(());
1521 }
1522
1523 let Some(storage) = self.storage.as_ref().map(Arc::clone) else {
1524 self.tenant_loader.mark_loaded(tenant_id);
1527 return Ok(());
1528 };
1529
1530 let lock = self.tenant_loader.lock_for(tenant_id);
1533 let timeout = self.tenant_loader.load_timeout();
1534 let _guard = lock.try_lock_for(timeout).ok_or_else(|| {
1535 AllSourceError::StorageError(format!(
1536 "ensure_tenant_loaded timed out after {timeout:?} waiting for in-flight load of \
1537 tenant {tenant_id:?}"
1538 ))
1539 })?;
1540
1541 if self.tenant_loader.is_loaded(tenant_id) {
1544 return Ok(());
1545 }
1546
1547 let started = std::time::Instant::now();
1548 let events = storage.read().load_events_for_tenant(tenant_id)?;
1549 let read_count = events.len();
1550
1551 let before = self.events.read().len();
1552 for event in events {
1553 self.append_loaded_event(event);
1554 }
1555 let applied = self.events.read().len() - before;
1556
1557 *self.total_ingested.write() += applied as u64;
1561 self.tenant_loader.mark_loaded(tenant_id);
1562
1563 tracing::info!(
1564 tenant_id = tenant_id,
1565 read = read_count,
1566 applied = applied,
1567 elapsed_ms = started.elapsed().as_millis() as u64,
1568 "ensure_tenant_loaded: tenant hydrated"
1569 );
1570
1571 self.enforce_cache_budget(tenant_id);
1579
1580 #[cfg(feature = "server")]
1583 self.metrics
1584 .cache_bytes
1585 .set(self.tenant_loader.total_bytes() as i64);
1586
1587 Ok(())
1588 }
1589
1590 fn enforce_cache_budget(&self, recently_touched: &str) {
1601 if !self.tenant_loader.over_budget() {
1602 return;
1603 }
1604 loop {
1605 let Some(victim) = self.tenant_loader.pick_lru_excluding(recently_touched) else {
1606 tracing::warn!(
1607 cache_bytes = self.tenant_loader.total_bytes(),
1608 budget = self.tenant_loader.byte_budget(),
1609 recently_touched = recently_touched,
1610 "cache over budget but no other tenant available to evict — \
1611 a single tenant exceeds the budget; consider raising it"
1612 );
1613 return;
1614 };
1615 self.evict_tenant(&victim);
1616 if !self.tenant_loader.over_budget() {
1617 return;
1618 }
1619 }
1620 }
1621
1622 pub fn is_tenant_loaded(&self, tenant_id: &str) -> bool {
1626 self.tenant_loader.is_loaded(tenant_id)
1627 }
1628
1629 pub fn evict_tenant(&self, tenant_id: &str) {
1657 let mut events = self.events.write();
1658 let before = events.len();
1659 let evicted_bytes = self.tenant_loader.bytes_for(tenant_id);
1660
1661 events.retain(|e| e.tenant_id_str() != tenant_id);
1662 let after = events.len();
1663 let dropped = before - after;
1664
1665 if dropped == 0 {
1666 drop(events);
1670 self.tenant_loader.mark_unloaded(tenant_id);
1671 return;
1672 }
1673
1674 self.index.clear();
1678 self.entity_versions.clear();
1679 for (offset, event) in events.iter().enumerate() {
1680 if let Err(e) = self.index.index_event(
1681 event.id,
1682 event.entity_id_str(),
1683 event.event_type_str(),
1684 event.timestamp,
1685 offset,
1686 ) {
1687 tracing::error!(
1688 "Failed to re-index event during eviction of {}: {}",
1689 tenant_id,
1690 e
1691 );
1692 }
1693 *self
1694 .entity_versions
1695 .entry(event.entity_id_str().to_string())
1696 .or_insert(0) += 1;
1697 }
1698 drop(events);
1699
1700 self.tenant_loader.mark_unloaded(tenant_id);
1701
1702 let mut t = self.total_ingested.write();
1705 *t = t.saturating_sub(dropped as u64);
1706 drop(t);
1707
1708 #[cfg(feature = "server")]
1711 {
1712 self.metrics.cache_evictions_total.inc();
1713 self.metrics
1714 .cache_bytes
1715 .set(self.tenant_loader.total_bytes() as i64);
1716 }
1717
1718 tracing::info!(
1719 tenant_id = tenant_id,
1720 events_dropped = dropped,
1721 bytes_freed = evicted_bytes,
1722 "evicted tenant from memory cache"
1723 );
1724 }
1725
1726 pub fn tenant_resident_bytes(&self, tenant_id: &str) -> u64 {
1730 self.tenant_loader.bytes_for(tenant_id)
1731 }
1732
1733 pub fn cache_resident_bytes(&self) -> u64 {
1736 self.tenant_loader.total_bytes()
1737 }
1738
1739 fn append_loaded_event(&self, event: Event) {
1761 if self.index.get_by_id(&event.id).is_some() {
1762 return;
1763 }
1764
1765 let event_bytes = event.estimated_size_bytes();
1766 let tenant = event.tenant_id_str().to_string();
1767
1768 let mut events = self.events.write();
1769 let offset = events.len();
1770
1771 if let Err(e) = self.index.index_event(
1772 event.id,
1773 event.entity_id_str(),
1774 event.event_type_str(),
1775 event.timestamp,
1776 offset,
1777 ) {
1778 tracing::error!("Failed to index loaded event {}: {}", event.id, e);
1779 }
1780
1781 if let Err(e) = self.projections.read().process_event(&event) {
1782 tracing::error!("Failed to project loaded event {}: {}", event.id, e);
1783 }
1784
1785 *self
1786 .entity_versions
1787 .entry(event.entity_id_str().to_string())
1788 .or_insert(0) += 1;
1789
1790 events.push(event);
1791 self.tenant_loader.add_bytes(&tenant, event_bytes);
1795 }
1796
1797 pub fn create_snapshot(&self, entity_id: &str) -> Result<()> {
1799 let events = self.query(&QueryEventsRequest {
1801 entity_id: Some(entity_id.to_string()),
1802 event_type: None,
1803 tenant_id: None,
1804 as_of: None,
1805 since: None,
1806 until: None,
1807 limit: None,
1808 event_type_prefix: None,
1809 exclude_event_type_prefix: None,
1810 payload_filter: None,
1811 })?;
1812
1813 if events.is_empty() {
1814 return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
1815 }
1816
1817 let mut state = serde_json::json!({});
1819 for event in &events {
1820 if let serde_json::Value::Object(ref mut state_map) = state
1821 && let serde_json::Value::Object(ref payload_map) = event.payload
1822 {
1823 for (key, value) in payload_map {
1824 state_map.insert(key.clone(), value.clone());
1825 }
1826 }
1827 }
1828
1829 let last_event = events.last().unwrap();
1830 self.snapshot_manager.create_snapshot(
1831 entity_id,
1832 state,
1833 last_event.timestamp,
1834 events.len(),
1835 SnapshotType::Manual,
1836 )?;
1837
1838 Ok(())
1839 }
1840
1841 fn check_auto_snapshot(&self, entity_id: &str, event: &Event) {
1843 let entity_event_count = self
1845 .index
1846 .get_by_entity(entity_id)
1847 .map_or(0, |entries| entries.len());
1848
1849 if self.snapshot_manager.should_create_snapshot(
1850 entity_id,
1851 entity_event_count,
1852 event.timestamp,
1853 ) {
1854 if let Err(e) = self.create_snapshot(entity_id) {
1856 tracing::warn!(
1857 "Failed to create automatic snapshot for {}: {}",
1858 entity_id,
1859 e
1860 );
1861 }
1862 }
1863 }
1864
1865 fn validate_event(&self, event: &Event) -> Result<()> {
1867 if event.entity_id_str().is_empty() {
1870 return Err(AllSourceError::ValidationError(
1871 "entity_id cannot be empty".to_string(),
1872 ));
1873 }
1874
1875 if event.event_type_str().is_empty() {
1876 return Err(AllSourceError::ValidationError(
1877 "event_type cannot be empty".to_string(),
1878 ));
1879 }
1880
1881 if event.event_type().is_system() {
1884 return Err(AllSourceError::ValidationError(
1885 "Event types starting with '_system.' are reserved for internal use".to_string(),
1886 ));
1887 }
1888
1889 Ok(())
1890 }
1891
1892 pub fn reset_projection(&self, name: &str) -> Result<usize> {
1894 let projection_manager = self.projections.read();
1895 let projection = projection_manager.get_projection(name).ok_or_else(|| {
1896 AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1897 })?;
1898
1899 projection.clear();
1901
1902 let prefix = format!("{name}:");
1904 let keys_to_remove: Vec<String> = self
1905 .projection_state_cache
1906 .iter()
1907 .filter(|entry| entry.key().starts_with(&prefix))
1908 .map(|entry| entry.key().clone())
1909 .collect();
1910 for key in keys_to_remove {
1911 self.projection_state_cache.remove(&key);
1912 }
1913
1914 let events = self.events.read();
1916 let mut reprocessed = 0usize;
1917 for event in events.iter() {
1918 if projection.process(event).is_ok() {
1919 reprocessed += 1;
1920 }
1921 }
1922
1923 Ok(reprocessed)
1924 }
1925
1926 pub fn get_event_by_id(&self, event_id: &uuid::Uuid) -> Result<Option<Event>> {
1928 if let Some(offset) = self.index.get_by_id(event_id) {
1929 let events = self.events.read();
1930 Ok(events.get(offset).cloned())
1931 } else {
1932 Ok(None)
1933 }
1934 }
1935
1936 #[cfg_attr(feature = "hotpath", hotpath::measure)]
1940 pub fn query(&self, request: &QueryEventsRequest) -> Result<Vec<Event>> {
1941 self.query_window(request, 0, false)
1942 .map(|(events, _)| events)
1943 }
1944
1945 pub fn query_scoped(
1947 &self,
1948 request: &QueryEventsRequest,
1949 scope: &ReadScope,
1950 ) -> Result<Vec<Event>> {
1951 self.query_window_scoped(request, 0, false, scope)
1952 .map(|(events, _)| events)
1953 }
1954
1955 #[cfg_attr(feature = "hotpath", hotpath::measure)]
1972 pub fn query_window(
1973 &self,
1974 request: &QueryEventsRequest,
1975 offset: usize,
1976 descending: bool,
1977 ) -> Result<(Vec<Event>, usize)> {
1978 self.query_window_scoped(request, offset, descending, &ReadScope::unrestricted())
1979 }
1980
1981 pub fn query_window_scoped(
1991 &self,
1992 request: &QueryEventsRequest,
1993 offset: usize,
1994 descending: bool,
1995 scope: &ReadScope,
1996 ) -> Result<(Vec<Event>, usize)> {
1997 if let Some(filter) = &request.payload_filter
2004 && serde_json::from_str::<serde_json::Map<String, serde_json::Value>>(filter).is_err()
2005 {
2006 return Err(AllSourceError::InvalidInput(format!(
2007 "invalid 'payload_filter': expected a JSON object of field/value \
2008 pairs, got '{filter}'"
2009 )));
2010 }
2011
2012 if let Some(ref tenant_id) = request.tenant_id {
2031 self.ensure_tenant_loaded(tenant_id)?;
2032 self.tenant_loader.touch(tenant_id);
2036 }
2037
2038 let query_type = if request.entity_id.is_some() {
2040 "entity"
2041 } else if request.event_type.is_some() {
2042 "type"
2043 } else if request.event_type_prefix.is_some() {
2044 "type_prefix"
2045 } else {
2046 "full_scan"
2047 };
2048
2049 #[cfg(feature = "server")]
2051 let timer = self
2052 .metrics
2053 .query_duration_seconds
2054 .with_label_values(&[query_type])
2055 .start_timer();
2056
2057 #[cfg(feature = "server")]
2059 self.metrics
2060 .queries_total
2061 .with_label_values(&[query_type])
2062 .inc();
2063
2064 let events = self.events.read();
2065
2066 let offsets: Vec<usize> = if let Some(entity_id) = &request.entity_id {
2068 self.index
2070 .get_by_entity(entity_id)
2071 .map(|entries| self.filter_entries(entries, request))
2072 .unwrap_or_default()
2073 } else if let Some(event_type) = &request.event_type {
2074 self.index
2076 .get_by_type(event_type)
2077 .map(|entries| self.filter_entries(entries, request))
2078 .unwrap_or_default()
2079 } else if let Some(prefix) = &request.event_type_prefix {
2080 let entries = self.index.get_by_type_prefix(prefix);
2082 self.filter_entries(entries, request)
2083 } else {
2084 (0..events.len()).collect()
2086 };
2087
2088 let mut matches: Vec<(usize, &Event)> = offsets
2094 .iter()
2095 .filter_map(|&event_offset| events.get(event_offset))
2096 .filter(|event| scope.permits(event.entity_id().as_str()))
2097 .filter(|event| self.apply_filters(event, request))
2098 .enumerate()
2099 .collect();
2100
2101 let total = matches.len();
2104
2105 let order = |a: &(usize, &Event), b: &(usize, &Event)| {
2112 let ascending =
2113 a.1.timestamp
2114 .cmp(&b.1.timestamp)
2115 .then_with(|| a.1.version.cmp(&b.1.version))
2116 .then_with(|| a.0.cmp(&b.0));
2117 if descending {
2118 ascending.reverse()
2119 } else {
2120 ascending
2121 }
2122 };
2123
2124 select_window(&mut matches, offset, request.limit, order);
2126
2127 let results: Vec<Event> = matches
2130 .into_iter()
2131 .skip(offset)
2132 .map(|(_, event)| event.clone())
2133 .collect();
2134
2135 #[cfg(feature = "server")]
2137 {
2138 self.metrics
2139 .query_results_total
2140 .with_label_values(&[query_type])
2141 .inc_by(results.len() as u64);
2142 timer.observe_duration();
2143 }
2144
2145 Ok((results, total))
2146 }
2147
2148 #[cfg_attr(feature = "hotpath", hotpath::measure)]
2150 fn filter_entries(&self, entries: Vec<IndexEntry>, request: &QueryEventsRequest) -> Vec<usize> {
2151 entries
2152 .into_iter()
2153 .filter(|entry| {
2154 if let Some(as_of) = request.as_of
2156 && entry.timestamp > as_of
2157 {
2158 return false;
2159 }
2160 if let Some(since) = request.since
2161 && entry.timestamp < since
2162 {
2163 return false;
2164 }
2165 if let Some(until) = request.until
2166 && entry.timestamp > until
2167 {
2168 return false;
2169 }
2170 true
2171 })
2172 .map(|entry| entry.offset)
2173 .collect()
2174 }
2175
2176 #[cfg_attr(feature = "hotpath", hotpath::measure)]
2178 fn apply_filters(&self, event: &Event, request: &QueryEventsRequest) -> bool {
2179 if let Some(ref tid) = request.tenant_id
2181 && event.tenant_id_str() != tid
2182 {
2183 return false;
2184 }
2185
2186 if let Some(as_of) = request.as_of
2196 && event.timestamp > as_of
2197 {
2198 return false;
2199 }
2200 if let Some(since) = request.since
2201 && event.timestamp < since
2202 {
2203 return false;
2204 }
2205 if let Some(until) = request.until
2206 && event.timestamp > until
2207 {
2208 return false;
2209 }
2210
2211 if let Some(ref excludes) = request.exclude_event_type_prefix {
2215 let et = event.event_type_str();
2216 if excludes
2217 .split(',')
2218 .map(str::trim)
2219 .filter(|p| !p.is_empty())
2220 .any(|p| et.starts_with(p))
2221 {
2222 return false;
2223 }
2224 }
2225
2226 if request.entity_id.is_some()
2228 && let Some(ref event_type) = request.event_type
2229 && event.event_type_str() != event_type
2230 {
2231 return false;
2232 }
2233
2234 if request.entity_id.is_some()
2236 && let Some(ref prefix) = request.event_type_prefix
2237 && !event.event_type_str().starts_with(prefix)
2238 {
2239 return false;
2240 }
2241
2242 if let Some(ref filter_str) = request.payload_filter
2244 && let Ok(filter_obj) =
2245 serde_json::from_str::<serde_json::Map<String, serde_json::Value>>(filter_str)
2246 {
2247 let payload = event.payload();
2248 for (key, expected_value) in &filter_obj {
2249 match payload.get(key) {
2250 Some(actual_value) if actual_value == expected_value => {}
2251 _ => return false,
2252 }
2253 }
2254 }
2255
2256 true
2257 }
2258
2259 #[cfg_attr(feature = "hotpath", hotpath::measure)]
2262 pub fn reconstruct_state(
2263 &self,
2264 entity_id: &str,
2265 as_of: Option<DateTime<Utc>>,
2266 ) -> Result<serde_json::Value> {
2267 let (merged_state, since_timestamp) = if let Some(as_of_time) = as_of {
2269 if let Some(snapshot) = self
2271 .snapshot_manager
2272 .get_snapshot_as_of(entity_id, as_of_time)
2273 {
2274 tracing::debug!(
2275 "Using snapshot from {} for entity {} (saved {} events)",
2276 snapshot.as_of,
2277 entity_id,
2278 snapshot.event_count
2279 );
2280 (snapshot.state.clone(), Some(snapshot.as_of))
2281 } else {
2282 (serde_json::json!({}), None)
2283 }
2284 } else {
2285 if let Some(snapshot) = self.snapshot_manager.get_latest_snapshot(entity_id) {
2287 tracing::debug!(
2288 "Using latest snapshot from {} for entity {}",
2289 snapshot.as_of,
2290 entity_id
2291 );
2292 (snapshot.state.clone(), Some(snapshot.as_of))
2293 } else {
2294 (serde_json::json!({}), None)
2295 }
2296 };
2297
2298 let events = self.query(&QueryEventsRequest {
2300 entity_id: Some(entity_id.to_string()),
2301 event_type: None,
2302 tenant_id: None,
2303 as_of,
2304 since: since_timestamp,
2305 until: None,
2306 limit: None,
2307 event_type_prefix: None,
2308 exclude_event_type_prefix: None,
2309 payload_filter: None,
2310 })?;
2311
2312 if events.is_empty() && since_timestamp.is_none() {
2314 return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
2315 }
2316
2317 let mut merged_state = merged_state;
2319 for event in &events {
2320 if let serde_json::Value::Object(ref mut state_map) = merged_state
2321 && let serde_json::Value::Object(ref payload_map) = event.payload
2322 {
2323 for (key, value) in payload_map {
2324 state_map.insert(key.clone(), value.clone());
2325 }
2326 }
2327 }
2328
2329 let state = serde_json::json!({
2331 "entity_id": entity_id,
2332 "last_updated": events.last().map(|e| e.timestamp),
2333 "event_count": events.len(),
2334 "as_of": as_of,
2335 "current_state": merged_state,
2336 "history": events.iter().map(|e| {
2337 serde_json::json!({
2338 "event_id": e.id,
2339 "type": e.event_type,
2340 "timestamp": e.timestamp,
2341 "payload": e.payload
2342 })
2343 }).collect::<Vec<_>>()
2344 });
2345
2346 Ok(state)
2347 }
2348
2349 pub fn get_snapshot(&self, entity_id: &str) -> Result<serde_json::Value> {
2351 let projections = self.projections.read();
2352
2353 if let Some(snapshot_projection) = projections.get_projection("entity_snapshots")
2354 && let Some(state) = snapshot_projection.get_state(entity_id)
2355 {
2356 return Ok(serde_json::json!({
2357 "entity_id": entity_id,
2358 "snapshot": state,
2359 "from_projection": "entity_snapshots"
2360 }));
2361 }
2362
2363 Err(AllSourceError::EntityNotFound(entity_id.to_string()))
2364 }
2365
2366 pub fn stats(&self) -> StoreStats {
2368 let events = self.events.read();
2369 let index_stats = self.index.stats();
2370
2371 StoreStats {
2372 total_events: events.len(),
2373 total_entities: index_stats.total_entities,
2374 total_event_types: index_stats.total_event_types,
2375 total_ingested: *self.total_ingested.read(),
2376 }
2377 }
2378
2379 pub fn list_streams(&self) -> Vec<StreamInfo> {
2381 self.index
2382 .get_all_entities()
2383 .into_iter()
2384 .map(|entity_id| {
2385 let event_count = self
2386 .index
2387 .get_by_entity(&entity_id)
2388 .map_or(0, |entries| entries.len());
2389 let last_event_at = self
2390 .index
2391 .get_by_entity(&entity_id)
2392 .and_then(|entries| entries.last().map(|e| e.timestamp));
2393 StreamInfo {
2394 stream_id: entity_id,
2395 event_count,
2396 last_event_at,
2397 }
2398 })
2399 .collect()
2400 }
2401
2402 pub fn list_event_types(&self) -> Vec<EventTypeInfo> {
2404 self.index
2405 .get_all_types()
2406 .into_iter()
2407 .map(|event_type| {
2408 let event_count = self
2409 .index
2410 .get_by_type(&event_type)
2411 .map_or(0, |entries| entries.len());
2412 let last_event_at = self
2413 .index
2414 .get_by_type(&event_type)
2415 .and_then(|entries| entries.last().map(|e| e.timestamp));
2416 EventTypeInfo {
2417 event_type,
2418 event_count,
2419 last_event_at,
2420 }
2421 })
2422 .collect()
2423 }
2424
2425 pub fn list_streams_for_tenant(&self, tenant_id: &str) -> Vec<StreamInfo> {
2433 let _ = self.ensure_tenant_loaded(tenant_id);
2434 let events = self.events.read();
2435 let mut by_entity: std::collections::HashMap<&str, (usize, chrono::DateTime<chrono::Utc>)> =
2436 std::collections::HashMap::new();
2437 for ev in events.iter() {
2438 if ev.tenant_id_str() != tenant_id {
2439 continue;
2440 }
2441 let e = by_entity
2442 .entry(ev.entity_id_str())
2443 .or_insert((0, ev.timestamp));
2444 e.0 += 1;
2445 if ev.timestamp > e.1 {
2446 e.1 = ev.timestamp;
2447 }
2448 }
2449 by_entity
2450 .into_iter()
2451 .map(|(entity_id, (count, last))| StreamInfo {
2452 stream_id: entity_id.to_string(),
2453 event_count: count,
2454 last_event_at: Some(last),
2455 })
2456 .collect()
2457 }
2458
2459 pub fn list_event_types_for_tenant(&self, tenant_id: &str) -> Vec<EventTypeInfo> {
2461 let _ = self.ensure_tenant_loaded(tenant_id);
2462 let events = self.events.read();
2463 let mut by_type: std::collections::HashMap<&str, (usize, chrono::DateTime<chrono::Utc>)> =
2464 std::collections::HashMap::new();
2465 for ev in events.iter() {
2466 if ev.tenant_id_str() != tenant_id {
2467 continue;
2468 }
2469 let e = by_type
2470 .entry(ev.event_type_str())
2471 .or_insert((0, ev.timestamp));
2472 e.0 += 1;
2473 if ev.timestamp > e.1 {
2474 e.1 = ev.timestamp;
2475 }
2476 }
2477 by_type
2478 .into_iter()
2479 .map(|(event_type, (count, last))| EventTypeInfo {
2480 event_type: event_type.to_string(),
2481 event_count: count,
2482 last_event_at: Some(last),
2483 })
2484 .collect()
2485 }
2486
2487 pub fn stats_for_tenant(&self, tenant_id: &str) -> TenantStoreStats {
2497 let _ = self.ensure_tenant_loaded(tenant_id);
2498 let events = self.events.read();
2499
2500 let mut entities: std::collections::HashSet<&str> = std::collections::HashSet::new();
2501 let mut census: std::collections::HashMap<&str, usize> = std::collections::HashMap::new();
2502 let mut total_events = 0usize;
2503 let mut oldest: Option<chrono::DateTime<chrono::Utc>> = None;
2504 let mut newest: Option<chrono::DateTime<chrono::Utc>> = None;
2505
2506 for ev in events.iter() {
2507 if ev.tenant_id_str() != tenant_id {
2508 continue;
2509 }
2510
2511 total_events += 1;
2512 entities.insert(ev.entity_id_str());
2513 *census.entry(ev.event_type_str()).or_insert(0) += 1;
2514
2515 let ts = ev.timestamp;
2516 if oldest.is_none_or(|o| ts < o) {
2517 oldest = Some(ts);
2518 }
2519 if newest.is_none_or(|n| ts > n) {
2520 newest = Some(ts);
2521 }
2522 }
2523
2524 TenantStoreStats {
2525 total_events,
2526 total_entities: entities.len(),
2527 total_event_types: census.len(),
2528 total_ingested: total_events as u64,
2532 event_types: census
2533 .into_iter()
2534 .map(|(k, v)| (k.to_string(), v))
2535 .collect(),
2536 oldest_event: oldest,
2537 newest_event: newest,
2538 }
2539 }
2540
2541 pub fn reconstruct_state_for_tenant(
2550 &self,
2551 entity_id: &str,
2552 as_of: Option<DateTime<Utc>>,
2553 tenant_id: &str,
2554 ) -> Result<serde_json::Value> {
2555 let events = self.query(&QueryEventsRequest {
2556 entity_id: Some(entity_id.to_string()),
2557 event_type: None,
2558 tenant_id: Some(tenant_id.to_string()),
2559 as_of,
2560 since: None,
2561 until: None,
2562 limit: None,
2563 event_type_prefix: None,
2564 exclude_event_type_prefix: None,
2565 payload_filter: None,
2566 })?;
2567
2568 if events.is_empty() {
2569 return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
2570 }
2571
2572 let mut merged_state = serde_json::json!({});
2573 for event in &events {
2574 if let serde_json::Value::Object(ref mut state_map) = merged_state
2575 && let serde_json::Value::Object(ref payload_map) = event.payload
2576 {
2577 for (key, value) in payload_map {
2578 state_map.insert(key.clone(), value.clone());
2579 }
2580 }
2581 }
2582
2583 Ok(serde_json::json!({
2584 "entity_id": entity_id,
2585 "last_updated": events.last().map(|e| e.timestamp),
2586 "event_count": events.len(),
2587 "as_of": as_of,
2588 "current_state": merged_state,
2589 "history": events.iter().map(|e| {
2590 serde_json::json!({
2591 "event_id": e.id,
2592 "type": e.event_type,
2593 "timestamp": e.timestamp,
2594 "payload": e.payload
2595 })
2596 }).collect::<Vec<_>>()
2597 }))
2598 }
2599
2600 pub fn enable_wal_replication(
2607 &self,
2608 tx: tokio::sync::broadcast::Sender<crate::infrastructure::persistence::wal::WALEntry>,
2609 ) {
2610 if let Some(ref wal_arc) = self.wal {
2611 wal_arc.set_replication_tx(tx);
2612 tracing::info!("WAL replication broadcast enabled");
2613 } else {
2614 tracing::warn!("Cannot enable WAL replication: WAL is not configured");
2615 }
2616 }
2617
2618 pub fn wal(&self) -> Option<&Arc<WriteAheadLog>> {
2621 self.wal.as_ref()
2622 }
2623
2624 pub fn parquet_storage(&self) -> Option<&Arc<RwLock<ParquetStorage>>> {
2627 self.storage.as_ref()
2628 }
2629}
2630
2631#[derive(Debug, Clone, Default)]
2633pub struct EventStoreConfig {
2634 pub storage_dir: Option<PathBuf>,
2636
2637 pub snapshot_config: SnapshotConfig,
2639
2640 pub wal_dir: Option<PathBuf>,
2642
2643 pub wal_config: WALConfig,
2645
2646 pub compaction_config: CompactionConfig,
2648
2649 pub schema_registry_config: SchemaRegistryConfig,
2651
2652 pub system_data_dir: Option<PathBuf>,
2657
2658 pub bootstrap_tenant: Option<String>,
2660
2661 pub cache_byte_budget: Option<u64>,
2668
2669 pub checkpoint_interval_secs: Option<u64>,
2681
2682 pub read_only: bool,
2686}
2687
2688impl EventStoreConfig {
2689 pub fn with_persistence(storage_dir: impl Into<PathBuf>) -> Self {
2691 Self {
2692 storage_dir: Some(storage_dir.into()),
2693 ..Self::default()
2694 }
2695 }
2696
2697 pub fn with_snapshots(snapshot_config: SnapshotConfig) -> Self {
2699 Self {
2700 snapshot_config,
2701 ..Self::default()
2702 }
2703 }
2704
2705 pub fn with_wal(wal_dir: impl Into<PathBuf>, wal_config: WALConfig) -> Self {
2707 Self {
2708 wal_dir: Some(wal_dir.into()),
2709 wal_config,
2710 ..Self::default()
2711 }
2712 }
2713
2714 pub fn with_all(storage_dir: impl Into<PathBuf>, snapshot_config: SnapshotConfig) -> Self {
2716 Self {
2717 storage_dir: Some(storage_dir.into()),
2718 snapshot_config,
2719 ..Self::default()
2720 }
2721 }
2722
2723 pub fn production(
2725 storage_dir: impl Into<PathBuf>,
2726 wal_dir: impl Into<PathBuf>,
2727 snapshot_config: SnapshotConfig,
2728 wal_config: WALConfig,
2729 compaction_config: CompactionConfig,
2730 ) -> Self {
2731 let storage_dir = storage_dir.into();
2732 let system_data_dir = storage_dir.join("__system");
2733 Self {
2734 storage_dir: Some(storage_dir),
2735 snapshot_config,
2736 wal_dir: Some(wal_dir.into()),
2737 wal_config,
2738 compaction_config,
2739 system_data_dir: Some(system_data_dir),
2740 ..Self::default()
2741 }
2742 }
2743
2744 pub fn effective_system_data_dir(&self) -> Option<PathBuf> {
2749 self.system_data_dir
2750 .clone()
2751 .or_else(|| self.storage_dir.as_ref().map(|d| d.join("__system")))
2752 }
2753
2754 pub fn from_env() -> (Self, &'static str) {
2762 Self::from_env_vars(
2763 std::env::var("ALLSOURCE_DATA_DIR")
2764 .ok()
2765 .filter(|s| !s.is_empty()),
2766 std::env::var("ALLSOURCE_STORAGE_DIR")
2767 .ok()
2768 .filter(|s| !s.is_empty()),
2769 std::env::var("ALLSOURCE_WAL_DIR")
2770 .ok()
2771 .filter(|s| !s.is_empty()),
2772 std::env::var("ALLSOURCE_WAL_ENABLED").ok(),
2773 std::env::var("ALLSOURCE_CACHE_BYTES").ok(),
2774 std::env::var("ALLSOURCE_SNAPSHOT_INTERVAL_SECONDS").ok(),
2775 std::env::var("ALLSOURCE_RETENTION_SYSTEM_DAYS").ok(),
2776 std::env::var("ALLSOURCE_CHECKPOINT_INTERVAL_SECONDS").ok(),
2777 )
2778 }
2779
2780 pub fn from_env_vars(
2782 data_dir: Option<String>,
2783 explicit_storage_dir: Option<String>,
2784 explicit_wal_dir: Option<String>,
2785 wal_enabled_var: Option<String>,
2786 cache_bytes_var: Option<String>,
2787 snapshot_interval_var: Option<String>,
2788 retention_system_days_var: Option<String>,
2789 checkpoint_interval_var: Option<String>,
2790 ) -> (Self, &'static str) {
2791 let data_dir = data_dir.filter(|s| !s.is_empty());
2792 let storage_dir = explicit_storage_dir
2793 .filter(|s| !s.is_empty())
2794 .or_else(|| data_dir.as_ref().map(|d| format!("{d}/storage")));
2795 let wal_dir = explicit_wal_dir
2796 .filter(|s| !s.is_empty())
2797 .or_else(|| data_dir.as_ref().map(|d| format!("{d}/wal")));
2798 let wal_enabled = wal_enabled_var.is_none_or(|v| v == "true");
2799 let cache_byte_budget =
2804 cache_bytes_var
2805 .filter(|s| !s.is_empty())
2806 .and_then(|s| match s.parse::<u64>() {
2807 Ok(v) => Some(v),
2808 Err(e) => {
2809 tracing::warn!(
2810 "ALLSOURCE_CACHE_BYTES={s:?} could not be parsed as u64: {e}; \
2811 cache budget disabled"
2812 );
2813 None
2814 }
2815 });
2816 let compaction_config =
2817 CompactionConfig::from_env_vars(snapshot_interval_var, retention_system_days_var);
2818
2819 let checkpoint_interval_secs = if wal_enabled {
2824 checkpoint_interval_var
2825 .filter(|s| !s.is_empty())
2826 .map(|s| match s.parse::<u64>() {
2827 Ok(v) => v,
2828 Err(e) => {
2829 tracing::warn!(
2830 "ALLSOURCE_CHECKPOINT_INTERVAL_SECONDS={s:?} could not be parsed as \
2831 u64: {e}; falling back to default 60s"
2832 );
2833 60
2834 }
2835 })
2836 .or(Some(60))
2837 } else {
2838 None
2839 };
2840
2841 let mut config = match (&storage_dir, &wal_dir) {
2842 (Some(sd), Some(wd)) if wal_enabled => Self::production(
2843 sd,
2844 wd,
2845 SnapshotConfig::default(),
2846 WALConfig::default(),
2847 compaction_config,
2848 ),
2849 (Some(sd), _) => Self::with_persistence(sd),
2850 (_, Some(wd)) if wal_enabled => Self::with_wal(wd, WALConfig::default()),
2851 _ => Self::default(),
2852 };
2853 config.cache_byte_budget = cache_byte_budget;
2854 config.checkpoint_interval_secs = checkpoint_interval_secs;
2855
2856 let mode = match (&storage_dir, &wal_dir) {
2857 (Some(_), Some(_)) if wal_enabled => "wal+parquet",
2858 (Some(_), _) => "parquet-only",
2859 (_, Some(_)) if wal_enabled => "wal-only",
2860 _ => "in-memory",
2861 };
2862 (config, mode)
2863 }
2864}
2865
2866#[derive(Debug, serde::Serialize)]
2867pub struct StoreStats {
2868 pub total_events: usize,
2869 pub total_entities: usize,
2870 pub total_event_types: usize,
2871 pub total_ingested: u64,
2872}
2873
2874#[derive(Debug, Clone, serde::Serialize)]
2880pub struct TenantStoreStats {
2881 pub total_events: usize,
2882 pub total_entities: usize,
2883 pub total_event_types: usize,
2884 pub total_ingested: u64,
2885 pub event_types: std::collections::HashMap<String, usize>,
2887 pub oldest_event: Option<chrono::DateTime<chrono::Utc>>,
2888 pub newest_event: Option<chrono::DateTime<chrono::Utc>>,
2889}
2890
2891#[derive(Debug, Clone, serde::Serialize)]
2893pub struct StreamInfo {
2894 pub stream_id: String,
2896 pub event_count: usize,
2898 pub last_event_at: Option<chrono::DateTime<chrono::Utc>>,
2900}
2901
2902#[derive(Debug, Clone, serde::Serialize)]
2904pub struct EventTypeInfo {
2905 pub event_type: String,
2907 pub event_count: usize,
2909 pub last_event_at: Option<chrono::DateTime<chrono::Utc>>,
2911}
2912
2913impl Default for EventStore {
2914 fn default() -> Self {
2915 Self::new()
2916 }
2917}
2918
2919#[cfg(test)]
2920mod tests {
2921 use super::*;
2922 use crate::domain::entities::Event;
2923 use tempfile::TempDir;
2924
2925 fn find_parquet_files(dir: &std::path::Path) -> Vec<std::path::PathBuf> {
2930 let mut out = Vec::new();
2931 let mut stack = vec![dir.to_path_buf()];
2932 while let Some(d) = stack.pop() {
2933 let Ok(entries) = std::fs::read_dir(&d) else {
2934 continue;
2935 };
2936 for e in entries.flatten() {
2937 let p = e.path();
2938 if p.is_dir() {
2939 stack.push(p);
2940 } else if p.extension().and_then(|s| s.to_str()) == Some("parquet") {
2941 out.push(p);
2942 }
2943 }
2944 }
2945 out
2946 }
2947
2948 fn create_test_event(entity_id: &str, event_type: &str) -> Event {
2949 Event::from_strings(
2950 event_type.to_string(),
2951 entity_id.to_string(),
2952 "default".to_string(),
2953 serde_json::json!({"name": "Test", "value": 42}),
2954 None,
2955 )
2956 .unwrap()
2957 }
2958
2959 fn create_test_event_with_payload(
2960 entity_id: &str,
2961 event_type: &str,
2962 payload: serde_json::Value,
2963 ) -> Event {
2964 Event::from_strings(
2965 event_type.to_string(),
2966 entity_id.to_string(),
2967 "default".to_string(),
2968 payload,
2969 None,
2970 )
2971 .unwrap()
2972 }
2973
2974 #[test]
2975 fn test_event_store_new() {
2976 let store = EventStore::new();
2977 assert_eq!(store.stats().total_events, 0);
2978 assert_eq!(store.stats().total_entities, 0);
2979 }
2980
2981 #[test]
2988 fn test_ensure_tenant_loaded_no_storage_is_a_noop() {
2989 let store = EventStore::new();
2993 assert!(!store.is_tenant_loaded("alice"));
2994 store.ensure_tenant_loaded("alice").unwrap();
2995 assert!(store.is_tenant_loaded("alice"));
2996 assert!(!store.is_tenant_loaded("bob"));
2998 }
2999
3000 #[test]
3001 fn test_ensure_tenant_loaded_warm_path_is_idempotent() {
3002 let store = EventStore::new();
3003 store.ensure_tenant_loaded("alice").unwrap();
3004 store.ensure_tenant_loaded("alice").unwrap();
3006 }
3007
3008 #[test]
3009 fn test_ensure_tenant_loaded_rejects_unsafe_tenant_id() {
3010 let temp_dir = TempDir::new().unwrap();
3016 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3017 for unsafe_tid in ["..", "a/b", "a\\b", ""] {
3018 let result = store.ensure_tenant_loaded(unsafe_tid);
3019 assert!(
3020 result.is_err(),
3021 "tenant_id {unsafe_tid:?} should have been rejected"
3022 );
3023 assert!(
3024 !store.is_tenant_loaded(unsafe_tid),
3025 "rejected tenant {unsafe_tid:?} must not be marked loaded"
3026 );
3027 }
3028 }
3029
3030 #[test]
3031 fn test_ensure_tenant_loaded_no_subtree_marks_loaded_with_zero_events() {
3032 let temp_dir = TempDir::new().unwrap();
3037 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3038 assert!(!store.is_tenant_loaded("never-existed"));
3039 store.ensure_tenant_loaded("never-existed").unwrap();
3040 assert!(store.is_tenant_loaded("never-existed"));
3041 }
3042
3043 #[test]
3044 fn test_evict_tenant_drops_events_and_resets_bytes() {
3045 let temp_dir = TempDir::new().unwrap();
3049 let storage_dir = temp_dir.path().to_path_buf();
3050
3051 {
3052 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3053 for i in 0..3 {
3054 store
3055 .ingest(
3056 &Event::from_strings(
3057 "test.event".to_string(),
3058 format!("a-{i}"),
3059 "alice".to_string(),
3060 serde_json::json!({"i": i}),
3061 None,
3062 )
3063 .unwrap(),
3064 )
3065 .unwrap();
3066 }
3067 for i in 0..2 {
3068 store
3069 .ingest(
3070 &Event::from_strings(
3071 "test.event".to_string(),
3072 format!("b-{i}"),
3073 "bob".to_string(),
3074 serde_json::json!({"i": i}),
3075 None,
3076 )
3077 .unwrap(),
3078 )
3079 .unwrap();
3080 }
3081 store.flush_storage().unwrap();
3082 }
3083
3084 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3085 store.ensure_tenant_loaded("alice").unwrap();
3086 store.ensure_tenant_loaded("bob").unwrap();
3087 assert_eq!(store.stats().total_events, 5);
3088 let alice_bytes = store.tenant_resident_bytes("alice");
3089 let bob_bytes = store.tenant_resident_bytes("bob");
3090 assert!(alice_bytes > 0 && bob_bytes > 0);
3091
3092 store.evict_tenant("alice");
3093
3094 assert!(!store.is_tenant_loaded("alice"));
3095 assert!(store.is_tenant_loaded("bob"));
3096 assert_eq!(store.tenant_resident_bytes("alice"), 0);
3097 assert_eq!(store.tenant_resident_bytes("bob"), bob_bytes);
3098 assert_eq!(store.stats().total_events, 2, "only bob's 2 events remain");
3099 }
3100
3101 #[test]
3102 fn test_evict_tenant_then_query_re_loads_from_disk() {
3103 let temp_dir = TempDir::new().unwrap();
3107 let storage_dir = temp_dir.path().to_path_buf();
3108
3109 {
3110 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3111 for i in 0..4 {
3112 store
3113 .ingest(
3114 &Event::from_strings(
3115 "test.event".to_string(),
3116 format!("a-{i}"),
3117 "alice".to_string(),
3118 serde_json::json!({"i": i}),
3119 None,
3120 )
3121 .unwrap(),
3122 )
3123 .unwrap();
3124 }
3125 store.flush_storage().unwrap();
3126 }
3127
3128 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3129 store.ensure_tenant_loaded("alice").unwrap();
3130 store.evict_tenant("alice");
3131 assert_eq!(store.stats().total_events, 0);
3132
3133 let results = store
3135 .query(&QueryEventsRequest {
3136 entity_id: None,
3137 event_type: None,
3138 tenant_id: Some("alice".to_string()),
3139 as_of: None,
3140 since: None,
3141 until: None,
3142 limit: None,
3143 event_type_prefix: None,
3144 exclude_event_type_prefix: None,
3145 payload_filter: None,
3146 })
3147 .unwrap();
3148 assert_eq!(results.len(), 4);
3149 assert!(store.is_tenant_loaded("alice"));
3150 }
3151
3152 #[test]
3153 fn test_evict_tenant_rebuilds_index_with_new_offsets() {
3154 let temp_dir = TempDir::new().unwrap();
3161 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3162
3163 for i in 0..3 {
3167 store
3168 .ingest(
3169 &Event::from_strings(
3170 "test.event".to_string(),
3171 format!("a-{i}"),
3172 "alice".to_string(),
3173 serde_json::json!({"i": i}),
3174 None,
3175 )
3176 .unwrap(),
3177 )
3178 .unwrap();
3179 if i < 2 {
3180 store
3181 .ingest(
3182 &Event::from_strings(
3183 "test.event".to_string(),
3184 format!("b-{i}"),
3185 "bob".to_string(),
3186 serde_json::json!({"i": i}),
3187 None,
3188 )
3189 .unwrap(),
3190 )
3191 .unwrap();
3192 }
3193 }
3194 store.tenant_loader.mark_loaded("alice");
3196 store.tenant_loader.mark_loaded("bob");
3197
3198 store.evict_tenant("alice");
3199
3200 let bob_results = store
3201 .query(&QueryEventsRequest {
3202 entity_id: None,
3203 event_type: None,
3204 tenant_id: Some("bob".to_string()),
3205 as_of: None,
3206 since: None,
3207 until: None,
3208 limit: None,
3209 event_type_prefix: None,
3210 exclude_event_type_prefix: None,
3211 payload_filter: None,
3212 })
3213 .unwrap();
3214 assert_eq!(bob_results.len(), 2);
3215 for e in &bob_results {
3216 assert_eq!(e.tenant_id_str(), "bob");
3217 }
3218 }
3219
3220 #[test]
3221 fn test_budget_eviction_keeps_resident_set_bounded() {
3222 let temp_dir = TempDir::new().unwrap();
3226 let storage_dir = temp_dir.path().to_path_buf();
3227
3228 let big_payload = serde_json::json!({"data": "x".repeat(1000)});
3231 {
3232 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3233 for tenant in ["alice", "bob", "carol"] {
3234 for i in 0..5 {
3235 store
3236 .ingest(
3237 &Event::from_strings(
3238 "test.event".to_string(),
3239 format!("{tenant}-{i}"),
3240 tenant.to_string(),
3241 big_payload.clone(),
3242 None,
3243 )
3244 .unwrap(),
3245 )
3246 .unwrap();
3247 }
3248 }
3249 store.flush_storage().unwrap();
3250 }
3251
3252 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3255 config.cache_byte_budget = Some(12_000);
3256 let store = EventStore::with_config(config);
3257
3258 store.ensure_tenant_loaded("alice").unwrap();
3260 assert!(store.is_tenant_loaded("alice"));
3261
3262 store.tenant_loader.touch("alice");
3267 std::thread::sleep(std::time::Duration::from_millis(10));
3268 store.ensure_tenant_loaded("bob").unwrap();
3269 assert!(store.is_tenant_loaded("bob"));
3270
3271 store.tenant_loader.touch("bob");
3275 std::thread::sleep(std::time::Duration::from_millis(10));
3276 store.ensure_tenant_loaded("carol").unwrap();
3277 assert!(store.is_tenant_loaded("carol"));
3278
3279 let resident = store.cache_resident_bytes();
3284 let budget = 12_000u64;
3285
3286 if resident > budget {
3289 let loaded_count = ["alice", "bob", "carol"]
3290 .iter()
3291 .filter(|t| store.is_tenant_loaded(t))
3292 .count();
3293 assert_eq!(
3294 loaded_count, 1,
3295 "over budget but more than one tenant loaded — eviction policy didn't fire"
3296 );
3297 }
3298
3299 assert!(store.is_tenant_loaded("carol"));
3302 }
3303
3304 #[test]
3305 fn test_query_after_eviction_re_loads_transparently() {
3306 let temp_dir = TempDir::new().unwrap();
3309 let storage_dir = temp_dir.path().to_path_buf();
3310
3311 let big_payload = serde_json::json!({"data": "x".repeat(2000)});
3312 {
3313 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3314 for tenant in ["alice", "bob"] {
3315 for i in 0..3 {
3316 store
3317 .ingest(
3318 &Event::from_strings(
3319 "test.event".to_string(),
3320 format!("{tenant}-{i}"),
3321 tenant.to_string(),
3322 big_payload.clone(),
3323 None,
3324 )
3325 .unwrap(),
3326 )
3327 .unwrap();
3328 }
3329 }
3330 store.flush_storage().unwrap();
3331 }
3332
3333 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3335 config.cache_byte_budget = Some(5_000);
3336 let store = EventStore::with_config(config);
3337
3338 let alice_first = store
3342 .query(&QueryEventsRequest {
3343 entity_id: None,
3344 event_type: None,
3345 tenant_id: Some("alice".to_string()),
3346 as_of: None,
3347 since: None,
3348 until: None,
3349 limit: None,
3350 event_type_prefix: None,
3351 exclude_event_type_prefix: None,
3352 payload_filter: None,
3353 })
3354 .unwrap();
3355 assert_eq!(alice_first.len(), 3);
3356
3357 std::thread::sleep(std::time::Duration::from_millis(15));
3359 let _bob = store
3361 .query(&QueryEventsRequest {
3362 entity_id: None,
3363 event_type: None,
3364 tenant_id: Some("bob".to_string()),
3365 as_of: None,
3366 since: None,
3367 until: None,
3368 limit: None,
3369 event_type_prefix: None,
3370 exclude_event_type_prefix: None,
3371 payload_filter: None,
3372 })
3373 .unwrap();
3374 assert!(
3375 !store.is_tenant_loaded("alice"),
3376 "alice should have been evicted"
3377 );
3378
3379 let alice_second = store
3381 .query(&QueryEventsRequest {
3382 entity_id: None,
3383 event_type: None,
3384 tenant_id: Some("alice".to_string()),
3385 as_of: None,
3386 since: None,
3387 until: None,
3388 limit: None,
3389 event_type_prefix: None,
3390 exclude_event_type_prefix: None,
3391 payload_filter: None,
3392 })
3393 .unwrap();
3394 assert_eq!(
3395 alice_second.len(),
3396 3,
3397 "alice's events come back via re-load"
3398 );
3399 assert!(store.is_tenant_loaded("alice"));
3400 }
3401
3402 #[test]
3403 #[cfg(feature = "server")]
3404 fn test_cache_metrics_track_evictions_and_bytes() {
3405 let temp_dir = TempDir::new().unwrap();
3409 let storage_dir = temp_dir.path().to_path_buf();
3410
3411 let big_payload = serde_json::json!({"data": "x".repeat(2000)});
3412 {
3413 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3414 for tenant in ["alice", "bob"] {
3415 for i in 0..3 {
3416 store
3417 .ingest(
3418 &Event::from_strings(
3419 "test.event".to_string(),
3420 format!("{tenant}-{i}"),
3421 tenant.to_string(),
3422 big_payload.clone(),
3423 None,
3424 )
3425 .unwrap(),
3426 )
3427 .unwrap();
3428 }
3429 }
3430 store.flush_storage().unwrap();
3431 }
3432
3433 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3434 config.cache_byte_budget = Some(5_000); let store = EventStore::with_config(config);
3436
3437 assert_eq!(store.metrics.cache_evictions_total.get(), 0);
3438 assert_eq!(store.metrics.cache_bytes.get(), 0);
3439
3440 store.ensure_tenant_loaded("alice").unwrap();
3441 let after_alice = store.metrics.cache_bytes.get();
3443 assert!(after_alice > 0, "gauge should reflect alice's bytes");
3444 assert_eq!(store.metrics.cache_evictions_total.get(), 0);
3446
3447 std::thread::sleep(std::time::Duration::from_millis(10));
3448 store.ensure_tenant_loaded("bob").unwrap();
3449
3450 assert_eq!(
3453 store.metrics.cache_evictions_total.get(),
3454 1,
3455 "exactly one tenant evicted after bob's load"
3456 );
3457 let after_bob = store.metrics.cache_bytes.get();
3459 assert!(after_bob > 0);
3460 assert!(after_bob <= after_alice, "gauge dropped after eviction");
3461 }
3462
3463 #[test]
3464 #[cfg(feature = "server")]
3465 fn test_storage_size_gauge_populated_from_on_disk_bytes() {
3466 let temp_dir = TempDir::new().unwrap();
3471 let config = EventStoreConfig {
3474 storage_dir: Some(temp_dir.path().join("parquet")),
3475 wal_dir: Some(temp_dir.path().join("wal")),
3476 ..EventStoreConfig::default()
3477 };
3478 let store = EventStore::with_config(config);
3479
3480 assert_eq!(store.metrics.storage_size_bytes.get(), 0);
3482
3483 let payload = serde_json::json!({ "data": "x".repeat(2000) });
3484 for i in 0..10 {
3485 store
3486 .ingest(
3487 &Event::from_strings(
3488 "test.event".to_string(),
3489 format!("entity-{i}"),
3490 "tenant-a".to_string(),
3491 payload.clone(),
3492 None,
3493 )
3494 .unwrap(),
3495 )
3496 .unwrap();
3497 }
3498 store.flush_storage().unwrap();
3499
3500 store.refresh_storage_metrics_now();
3502
3503 let size = store.metrics.storage_size_bytes.get();
3504 assert!(
3505 size > 0,
3506 "storage_size_bytes must reflect real on-disk bytes, got {size}"
3507 );
3508 assert!(
3509 store.metrics.parquet_files_total.get() >= 1,
3510 "at least one Parquet file should exist after a flush"
3511 );
3512
3513 assert!(
3516 store.metrics.wal_segments_total.get() >= 1,
3517 "at least one WAL segment should exist after writes"
3518 );
3519
3520 let parquet_stats = store.storage.as_ref().unwrap().read().stats().unwrap();
3523 let (wal_bytes, _) = store.wal.as_ref().unwrap().on_disk_stats().unwrap();
3524 assert_eq!(
3525 size as u64,
3526 parquet_stats.total_size_bytes + wal_bytes,
3527 "gauge should equal Parquet bytes ({}) + WAL bytes ({wal_bytes})",
3528 parquet_stats.total_size_bytes
3529 );
3530 }
3531
3532 #[test]
3533 fn test_stress_resident_set_stays_near_budget_under_rolling_queries() {
3534 let temp_dir = TempDir::new().unwrap();
3541 let storage_dir = temp_dir.path().to_path_buf();
3542
3543 const TENANT_COUNT: usize = 10;
3544 const EVENTS_PER_TENANT: usize = 50;
3545 let big_payload = serde_json::json!({"data": "x".repeat(10_000)});
3547
3548 {
3551 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3552 for t in 0..TENANT_COUNT {
3553 let tenant = format!("tenant-{t}");
3554 for i in 0..EVENTS_PER_TENANT {
3555 store
3556 .ingest(
3557 &Event::from_strings(
3558 "test.event".to_string(),
3559 format!("{tenant}-{i}"),
3560 tenant.clone(),
3561 big_payload.clone(),
3562 None,
3563 )
3564 .unwrap(),
3565 )
3566 .unwrap();
3567 }
3568 }
3569 store.flush_storage().unwrap();
3570 }
3571
3572 const BUDGET: u64 = 1_048_576;
3576 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3577 config.cache_byte_budget = Some(BUDGET);
3578 let store = EventStore::with_config(config);
3579
3580 let mut peak_resident: u64 = 0;
3584 for t in 0..TENANT_COUNT {
3585 let tenant = format!("tenant-{t}");
3586 let results = store
3587 .query(&QueryEventsRequest {
3588 entity_id: None,
3589 event_type: None,
3590 tenant_id: Some(tenant.clone()),
3591 as_of: None,
3592 since: None,
3593 until: None,
3594 limit: None,
3595 event_type_prefix: None,
3596 exclude_event_type_prefix: None,
3597 payload_filter: None,
3598 })
3599 .unwrap();
3600 assert_eq!(
3601 results.len(),
3602 EVENTS_PER_TENANT,
3603 "every per-tenant query must return all of that tenant's events"
3604 );
3605 let resident = store.cache_resident_bytes();
3607 if resident > peak_resident {
3608 peak_resident = resident;
3609 }
3610 }
3611
3612 let final_resident = store.cache_resident_bytes();
3613
3614 let tolerance = BUDGET; assert!(
3620 peak_resident <= BUDGET + tolerance,
3621 "peak resident {peak_resident} exceeds budget {BUDGET} by more than {tolerance} \
3622 — eviction policy not keeping up with the working-set churn"
3623 );
3624 assert!(
3625 final_resident <= BUDGET + tolerance,
3626 "final resident {final_resident} exceeds budget {BUDGET} by more than {tolerance}"
3627 );
3628
3629 let last_tenant = format!("tenant-{}", TENANT_COUNT - 1);
3632 assert!(
3633 store.is_tenant_loaded(&last_tenant),
3634 "the most-recent tenant must remain loaded after the sweep"
3635 );
3636
3637 let still_loaded = (0..TENANT_COUNT)
3640 .filter(|t| store.is_tenant_loaded(&format!("tenant-{t}")))
3641 .count();
3642 assert!(
3643 still_loaded < TENANT_COUNT,
3644 "no tenants evicted ({still_loaded}/{TENANT_COUNT} still loaded) — \
3645 budget enforcement didn't engage"
3646 );
3647 }
3648
3649 #[test]
3650 fn test_evict_tenant_when_not_loaded_is_a_noop() {
3651 let store = EventStore::new();
3654 store.evict_tenant("nobody"); assert!(!store.is_tenant_loaded("nobody"));
3656 }
3657
3658 #[test]
3659 fn test_lazy_load_accounts_bytes_per_tenant() {
3660 let temp_dir = TempDir::new().unwrap();
3664 let storage_dir = temp_dir.path().to_path_buf();
3665
3666 {
3668 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3669 for i in 0..5 {
3670 store
3671 .ingest(
3672 &Event::from_strings(
3673 "test.event".to_string(),
3674 format!("a-{i}"),
3675 "alice".to_string(),
3676 serde_json::json!({"data": "x".repeat(1000)}),
3677 None,
3678 )
3679 .unwrap(),
3680 )
3681 .unwrap();
3682 }
3683 store.flush_storage().unwrap();
3684 }
3685
3686 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3687 assert_eq!(store.tenant_resident_bytes("alice"), 0);
3689 assert_eq!(store.cache_resident_bytes(), 0);
3690
3691 store.ensure_tenant_loaded("alice").unwrap();
3692
3693 let alice_bytes = store.tenant_resident_bytes("alice");
3696 assert!(
3697 alice_bytes >= 5 * 1000,
3698 "alice should have at least 5 KiB resident; got {alice_bytes}"
3699 );
3700 assert_eq!(store.tenant_resident_bytes("bob"), 0);
3702 assert_eq!(store.cache_resident_bytes(), alice_bytes);
3704 }
3705
3706 #[test]
3707 fn test_query_lazy_loads_tenant_on_first_call() {
3708 let temp_dir = TempDir::new().unwrap();
3712 let storage_dir = temp_dir.path().to_path_buf();
3713
3714 {
3716 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3717 for i in 0..3 {
3718 let event = Event::from_strings(
3719 "test.event".to_string(),
3720 format!("e-{i}"),
3721 "alice".to_string(),
3722 serde_json::json!({"i": i}),
3723 None,
3724 )
3725 .unwrap();
3726 store.ingest(&event).unwrap();
3727 }
3728 store.flush_storage().unwrap();
3729 }
3730
3731 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3733 assert_eq!(
3734 store.stats().total_events,
3735 0,
3736 "boot must be O(1) — no Parquet pre-load"
3737 );
3738 assert!(!store.is_tenant_loaded("alice"));
3739 assert!(!store.is_tenant_loaded("bob"));
3740
3741 let results = store
3743 .query(&QueryEventsRequest {
3744 entity_id: None,
3745 event_type: None,
3746 tenant_id: Some("alice".to_string()),
3747 as_of: None,
3748 since: None,
3749 until: None,
3750 limit: None,
3751 event_type_prefix: None,
3752 exclude_event_type_prefix: None,
3753 payload_filter: None,
3754 })
3755 .unwrap();
3756 assert_eq!(results.len(), 3, "alice's 3 events are returned");
3757 assert!(store.is_tenant_loaded("alice"), "alice now warm");
3758 assert!(!store.is_tenant_loaded("bob"), "bob still cold");
3761 }
3762
3763 #[test]
3764 fn test_query_invalid_tenant_id_returns_error_no_hang() {
3765 let temp_dir = TempDir::new().unwrap();
3769 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3770
3771 let result = store.query(&QueryEventsRequest {
3772 entity_id: None,
3773 event_type: None,
3774 tenant_id: Some("../etc".to_string()),
3775 as_of: None,
3776 since: None,
3777 until: None,
3778 limit: None,
3779 event_type_prefix: None,
3780 exclude_event_type_prefix: None,
3781 payload_filter: None,
3782 });
3783 assert!(result.is_err(), "unsafe tenant_id must surface as error");
3784 }
3785
3786 #[test]
3787 fn test_query_concurrent_first_queries_for_same_tenant_all_succeed() {
3788 let temp_dir = TempDir::new().unwrap();
3796 let storage_dir = temp_dir.path().to_path_buf();
3797
3798 {
3800 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3801 for i in 0..25 {
3802 let event = Event::from_strings(
3803 "test.event".to_string(),
3804 format!("e-{i}"),
3805 "alice".to_string(),
3806 serde_json::json!({"i": i}),
3807 None,
3808 )
3809 .unwrap();
3810 store.ingest(&event).unwrap();
3811 }
3812 store.flush_storage().unwrap();
3813 }
3814
3815 let store = Arc::new(EventStore::with_config(EventStoreConfig::with_persistence(
3817 &storage_dir,
3818 )));
3819 assert!(!store.is_tenant_loaded("alice"));
3820
3821 let mut handles = Vec::new();
3822 for _ in 0..8 {
3823 let s = store.clone();
3824 handles.push(std::thread::spawn(move || {
3825 s.query(&QueryEventsRequest {
3826 entity_id: None,
3827 event_type: None,
3828 tenant_id: Some("alice".to_string()),
3829 as_of: None,
3830 since: None,
3831 until: None,
3832 limit: None,
3833 event_type_prefix: None,
3834 exclude_event_type_prefix: None,
3835 payload_filter: None,
3836 })
3837 }));
3838 }
3839
3840 for h in handles {
3841 let result = h.join().unwrap().unwrap();
3842 assert_eq!(
3843 result.len(),
3844 25,
3845 "every concurrent caller must see all 25 events"
3846 );
3847 }
3848 assert!(store.is_tenant_loaded("alice"));
3849 assert_eq!(store.stats().total_events, 25);
3851 }
3852
3853 #[test]
3854 fn test_query_two_cold_tenants_load_independently() {
3855 let temp_dir = TempDir::new().unwrap();
3859 let storage_dir = temp_dir.path().to_path_buf();
3860
3861 {
3862 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3863 for i in 0..3 {
3864 store
3865 .ingest(
3866 &Event::from_strings(
3867 "test.event".to_string(),
3868 format!("a-{i}"),
3869 "alice".to_string(),
3870 serde_json::json!({"i": i}),
3871 None,
3872 )
3873 .unwrap(),
3874 )
3875 .unwrap();
3876 }
3877 for i in 0..5 {
3878 store
3879 .ingest(
3880 &Event::from_strings(
3881 "test.event".to_string(),
3882 format!("b-{i}"),
3883 "bob".to_string(),
3884 serde_json::json!({"i": i}),
3885 None,
3886 )
3887 .unwrap(),
3888 )
3889 .unwrap();
3890 }
3891 store.flush_storage().unwrap();
3892 }
3893
3894 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3895 assert_eq!(store.stats().total_events, 0);
3896
3897 let alice = store
3899 .query(&QueryEventsRequest {
3900 entity_id: None,
3901 event_type: None,
3902 tenant_id: Some("alice".to_string()),
3903 as_of: None,
3904 since: None,
3905 until: None,
3906 limit: None,
3907 event_type_prefix: None,
3908 exclude_event_type_prefix: None,
3909 payload_filter: None,
3910 })
3911 .unwrap();
3912 assert_eq!(alice.len(), 3);
3913 assert!(store.is_tenant_loaded("alice"));
3914 assert!(!store.is_tenant_loaded("bob"));
3915 assert_eq!(store.stats().total_events, 3);
3916
3917 let bob = store
3919 .query(&QueryEventsRequest {
3920 entity_id: None,
3921 event_type: None,
3922 tenant_id: Some("bob".to_string()),
3923 as_of: None,
3924 since: None,
3925 until: None,
3926 limit: None,
3927 event_type_prefix: None,
3928 exclude_event_type_prefix: None,
3929 payload_filter: None,
3930 })
3931 .unwrap();
3932 assert_eq!(bob.len(), 5);
3933 assert!(store.is_tenant_loaded("bob"));
3934 assert_eq!(store.stats().total_events, 8);
3935 }
3936
3937 #[test]
3938 fn test_boot_with_persisted_data_is_o1() {
3939 let temp_dir = TempDir::new().unwrap();
3952 let storage_dir = temp_dir.path().to_path_buf();
3953
3954 {
3955 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3956 for tenant in ["alice", "bob", "carol"] {
3957 for i in 0..50 / 3 {
3958 store
3959 .ingest(
3960 &Event::from_strings(
3961 "test.event".to_string(),
3962 format!("{tenant}-{i}"),
3963 tenant.to_string(),
3964 serde_json::json!({"i": i}),
3965 None,
3966 )
3967 .unwrap(),
3968 )
3969 .unwrap();
3970 }
3971 }
3972 store.flush_storage().unwrap();
3973 }
3974
3975 let on_disk = find_parquet_files(&storage_dir);
3977 assert!(
3978 !on_disk.is_empty(),
3979 "session 1 should have produced parquet files; pre-condition for the test"
3980 );
3981
3982 let started = std::time::Instant::now();
3983 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3984 let boot_elapsed = started.elapsed();
3985
3986 assert_eq!(
3987 store.stats().total_events,
3988 0,
3989 "boot must not pre-load any Parquet events"
3990 );
3991
3992 assert!(
3996 boot_elapsed < std::time::Duration::from_secs(2),
3997 "boot took {boot_elapsed:?} — Step 2 boot should be O(1)"
3998 );
3999 }
4000
4001 #[test]
4002 fn test_query_warm_tenant_does_not_re_read_disk() {
4003 let temp_dir = TempDir::new().unwrap();
4009 let storage_dir = temp_dir.path().to_path_buf();
4010
4011 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
4012 for i in 0..3 {
4013 let event = Event::from_strings(
4014 "test.event".to_string(),
4015 format!("e-{i}"),
4016 "alice".to_string(),
4017 serde_json::json!({"i": i}),
4018 None,
4019 )
4020 .unwrap();
4021 store.ingest(&event).unwrap();
4022 }
4023 store.flush_storage().unwrap();
4024
4025 let _ = store
4027 .query(&QueryEventsRequest {
4028 entity_id: None,
4029 event_type: None,
4030 tenant_id: Some("alice".to_string()),
4031 as_of: None,
4032 since: None,
4033 until: None,
4034 limit: None,
4035 event_type_prefix: None,
4036 exclude_event_type_prefix: None,
4037 payload_filter: None,
4038 })
4039 .unwrap();
4040 assert!(store.is_tenant_loaded("alice"));
4041
4042 let parquet_files = find_parquet_files(&storage_dir);
4045 for f in parquet_files {
4046 std::fs::remove_file(&f).unwrap();
4047 }
4048
4049 let results = store
4050 .query(&QueryEventsRequest {
4051 entity_id: None,
4052 event_type: None,
4053 tenant_id: Some("alice".to_string()),
4054 as_of: None,
4055 since: None,
4056 until: None,
4057 limit: None,
4058 event_type_prefix: None,
4059 exclude_event_type_prefix: None,
4060 payload_filter: None,
4061 })
4062 .unwrap();
4063 assert_eq!(
4064 results.len(),
4065 3,
4066 "warm tenant query must not need disk; got {} events from a deleted parquet",
4067 results.len()
4068 );
4069 }
4070
4071 #[test]
4072 fn test_event_store_default() {
4073 let store = EventStore::default();
4074 assert_eq!(store.stats().total_events, 0);
4075 }
4076
4077 #[test]
4078 fn test_ingest_single_event() {
4079 let store = EventStore::new();
4080 let event = create_test_event("entity-1", "user.created");
4081
4082 store.ingest(&event).unwrap();
4083
4084 assert_eq!(store.stats().total_events, 1);
4085 assert_eq!(store.stats().total_ingested, 1);
4086 }
4087
4088 #[test]
4089 fn test_ingest_multiple_events() {
4090 let store = EventStore::new();
4091
4092 for i in 0..10 {
4093 let event = create_test_event(&format!("entity-{i}"), "user.created");
4094 store.ingest(&event).unwrap();
4095 }
4096
4097 assert_eq!(store.stats().total_events, 10);
4098 assert_eq!(store.stats().total_ingested, 10);
4099 }
4100
4101 #[test]
4102 fn test_query_by_entity_id() {
4103 let store = EventStore::new();
4104
4105 store
4106 .ingest(&create_test_event("entity-1", "user.created"))
4107 .unwrap();
4108 store
4109 .ingest(&create_test_event("entity-2", "user.created"))
4110 .unwrap();
4111 store
4112 .ingest(&create_test_event("entity-1", "user.updated"))
4113 .unwrap();
4114
4115 let results = store
4116 .query(&QueryEventsRequest {
4117 entity_id: Some("entity-1".to_string()),
4118 event_type: None,
4119 tenant_id: None,
4120 as_of: None,
4121 since: None,
4122 until: None,
4123 limit: None,
4124 event_type_prefix: None,
4125 exclude_event_type_prefix: None,
4126 payload_filter: None,
4127 })
4128 .unwrap();
4129
4130 assert_eq!(results.len(), 2);
4131 }
4132
4133 #[test]
4134 fn test_query_by_event_type() {
4135 let store = EventStore::new();
4136
4137 store
4138 .ingest(&create_test_event("entity-1", "user.created"))
4139 .unwrap();
4140 store
4141 .ingest(&create_test_event("entity-2", "user.updated"))
4142 .unwrap();
4143 store
4144 .ingest(&create_test_event("entity-3", "user.created"))
4145 .unwrap();
4146
4147 let results = store
4148 .query(&QueryEventsRequest {
4149 entity_id: None,
4150 event_type: Some("user.created".to_string()),
4151 tenant_id: None,
4152 as_of: None,
4153 since: None,
4154 until: None,
4155 limit: None,
4156 event_type_prefix: None,
4157 exclude_event_type_prefix: None,
4158 payload_filter: None,
4159 })
4160 .unwrap();
4161
4162 assert_eq!(results.len(), 2);
4163 }
4164
4165 #[test]
4166 fn test_query_with_limit() {
4167 let store = EventStore::new();
4168
4169 for i in 0..10 {
4170 let event = create_test_event(&format!("entity-{i}"), "user.created");
4171 store.ingest(&event).unwrap();
4172 }
4173
4174 let results = store
4175 .query(&QueryEventsRequest {
4176 entity_id: None,
4177 event_type: None,
4178 tenant_id: None,
4179 as_of: None,
4180 since: None,
4181 until: None,
4182 limit: Some(5),
4183 event_type_prefix: None,
4184 exclude_event_type_prefix: None,
4185 payload_filter: None,
4186 })
4187 .unwrap();
4188
4189 assert_eq!(results.len(), 5);
4190 }
4191
4192 #[test]
4193 fn test_query_empty_store() {
4194 let store = EventStore::new();
4195
4196 let results = store
4197 .query(&QueryEventsRequest {
4198 entity_id: Some("non-existent".to_string()),
4199 event_type: None,
4200 tenant_id: None,
4201 as_of: None,
4202 since: None,
4203 until: None,
4204 limit: None,
4205 event_type_prefix: None,
4206 exclude_event_type_prefix: None,
4207 payload_filter: None,
4208 })
4209 .unwrap();
4210
4211 assert!(results.is_empty());
4212 }
4213
4214 #[test]
4215 fn test_reconstruct_state() {
4216 let store = EventStore::new();
4217
4218 store
4219 .ingest(&create_test_event("entity-1", "user.created"))
4220 .unwrap();
4221
4222 let state = store.reconstruct_state("entity-1", None).unwrap();
4223 assert_eq!(state["current_state"]["name"], "Test");
4225 assert_eq!(state["current_state"]["value"], 42);
4226 }
4227
4228 #[test]
4229 fn test_reconstruct_state_not_found() {
4230 let store = EventStore::new();
4231
4232 let result = store.reconstruct_state("non-existent", None);
4233 assert!(result.is_err());
4234 }
4235
4236 #[test]
4237 fn test_get_snapshot_empty() {
4238 let store = EventStore::new();
4239
4240 let result = store.get_snapshot("non-existent");
4241 assert!(result.is_err());
4243 }
4244
4245 #[test]
4246 fn test_create_snapshot() {
4247 let store = EventStore::new();
4248
4249 store
4250 .ingest(&create_test_event("entity-1", "user.created"))
4251 .unwrap();
4252
4253 store.create_snapshot("entity-1").unwrap();
4254
4255 let snapshot = store.get_snapshot("entity-1").unwrap();
4257 assert_ne!(snapshot, serde_json::json!(null));
4258 }
4259
4260 #[test]
4261 fn test_create_snapshot_entity_not_found() {
4262 let store = EventStore::new();
4263
4264 let result = store.create_snapshot("non-existent");
4265 assert!(result.is_err());
4266 }
4267
4268 #[test]
4269 fn test_websocket_manager() {
4270 let store = EventStore::new();
4271 let manager = store.websocket_manager();
4272 assert!(Arc::strong_count(&manager) >= 1);
4274 }
4275
4276 #[test]
4277 fn test_snapshot_manager() {
4278 let store = EventStore::new();
4279 let manager = store.snapshot_manager();
4280 assert!(Arc::strong_count(&manager) >= 1);
4281 }
4282
4283 #[test]
4284 fn test_compaction_manager_none() {
4285 let store = EventStore::new();
4286 assert!(store.compaction_manager().is_none());
4288 }
4289
4290 #[test]
4291 fn test_schema_registry() {
4292 let store = EventStore::new();
4293 let registry = store.schema_registry();
4294 assert!(Arc::strong_count(®istry) >= 1);
4295 }
4296
4297 #[test]
4298 fn test_replay_manager() {
4299 let store = EventStore::new();
4300 let manager = store.replay_manager();
4301 assert!(Arc::strong_count(&manager) >= 1);
4302 }
4303
4304 #[test]
4305 fn test_pipeline_manager() {
4306 let store = EventStore::new();
4307 let manager = store.pipeline_manager();
4308 assert!(Arc::strong_count(&manager) >= 1);
4309 }
4310
4311 #[test]
4312 fn test_projection_manager() {
4313 let store = EventStore::new();
4314 let manager = store.projection_manager();
4315 let projections = manager.list_projections();
4317 assert!(projections.len() >= 2); }
4319
4320 #[test]
4321 fn test_projection_state_cache() {
4322 let store = EventStore::new();
4323 let cache = store.projection_state_cache();
4324
4325 cache.insert("test:key".to_string(), serde_json::json!({"value": 123}));
4326 assert_eq!(cache.len(), 1);
4327
4328 let value = cache.get("test:key").unwrap();
4329 assert_eq!(value["value"], 123);
4330 }
4331
4332 #[test]
4333 fn test_metrics() {
4334 let store = EventStore::new();
4335 let metrics = store.metrics();
4336 assert!(Arc::strong_count(&metrics) >= 1);
4337 }
4338
4339 #[test]
4340 fn test_store_stats() {
4341 let store = EventStore::new();
4342
4343 store
4344 .ingest(&create_test_event("entity-1", "user.created"))
4345 .unwrap();
4346 store
4347 .ingest(&create_test_event("entity-2", "order.placed"))
4348 .unwrap();
4349
4350 let stats = store.stats();
4351 assert_eq!(stats.total_events, 2);
4352 assert_eq!(stats.total_entities, 2);
4353 assert_eq!(stats.total_event_types, 2);
4354 assert_eq!(stats.total_ingested, 2);
4355 }
4356
4357 #[test]
4358 fn test_event_store_config_default() {
4359 let config = EventStoreConfig::default();
4360 assert!(config.storage_dir.is_none());
4361 assert!(config.wal_dir.is_none());
4362 }
4363
4364 #[test]
4365 fn test_event_store_config_with_persistence() {
4366 let temp_dir = TempDir::new().unwrap();
4367 let config = EventStoreConfig::with_persistence(temp_dir.path());
4368
4369 assert!(config.storage_dir.is_some());
4370 assert!(config.wal_dir.is_none());
4371 }
4372
4373 #[test]
4374 fn test_event_store_config_with_wal() {
4375 let temp_dir = TempDir::new().unwrap();
4376 let config = EventStoreConfig::with_wal(temp_dir.path(), WALConfig::default());
4377
4378 assert!(config.storage_dir.is_none());
4379 assert!(config.wal_dir.is_some());
4380 }
4381
4382 #[test]
4383 fn test_event_store_config_with_all() {
4384 let temp_dir = TempDir::new().unwrap();
4385 let config = EventStoreConfig::with_all(temp_dir.path(), SnapshotConfig::default());
4386
4387 assert!(config.storage_dir.is_some());
4388 }
4389
4390 #[test]
4391 fn test_event_store_config_production() {
4392 let storage_dir = TempDir::new().unwrap();
4393 let wal_dir = TempDir::new().unwrap();
4394 let config = EventStoreConfig::production(
4395 storage_dir.path(),
4396 wal_dir.path(),
4397 SnapshotConfig::default(),
4398 WALConfig::default(),
4399 CompactionConfig::default(),
4400 );
4401
4402 assert!(config.storage_dir.is_some());
4403 assert!(config.wal_dir.is_some());
4404 }
4405
4406 #[test]
4412 fn test_from_env_vars_data_dir_enables_full_persistence() {
4413 let (config, mode) = EventStoreConfig::from_env_vars(
4414 Some("/app/data".to_string()),
4415 None,
4416 None,
4417 None,
4418 None,
4419 None,
4420 None,
4421 None,
4422 );
4423 assert_eq!(mode, "wal+parquet");
4424 assert_eq!(
4425 config.storage_dir.unwrap().to_str().unwrap(),
4426 "/app/data/storage"
4427 );
4428 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/app/data/wal");
4429 }
4430
4431 #[test]
4432 fn test_from_env_vars_explicit_dirs() {
4433 let (config, mode) = EventStoreConfig::from_env_vars(
4434 None,
4435 Some("/custom/storage".to_string()),
4436 Some("/custom/wal".to_string()),
4437 None,
4438 None,
4439 None,
4440 None,
4441 None,
4442 );
4443 assert_eq!(mode, "wal+parquet");
4444 assert_eq!(
4445 config.storage_dir.unwrap().to_str().unwrap(),
4446 "/custom/storage"
4447 );
4448 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/custom/wal");
4449 }
4450
4451 #[test]
4452 fn test_from_env_vars_wal_disabled() {
4453 let (config, mode) = EventStoreConfig::from_env_vars(
4454 Some("/app/data".to_string()),
4455 None,
4456 None,
4457 Some("false".to_string()),
4458 None,
4459 None,
4460 None,
4461 None,
4462 );
4463 assert_eq!(mode, "parquet-only");
4464 assert!(config.storage_dir.is_some());
4465 assert!(config.wal_dir.is_none());
4466 }
4467
4468 #[test]
4469 fn test_from_env_vars_no_dirs_is_in_memory() {
4470 let (config, mode) =
4471 EventStoreConfig::from_env_vars(None, None, None, None, None, None, None, None);
4472 assert_eq!(mode, "in-memory");
4473 assert!(config.storage_dir.is_none());
4474 assert!(config.wal_dir.is_none());
4475 }
4476
4477 #[test]
4478 fn test_from_env_vars_empty_strings_treated_as_none() {
4479 let (_, mode) = EventStoreConfig::from_env_vars(
4480 Some(String::new()),
4481 Some(String::new()),
4482 Some(String::new()),
4483 None,
4484 None,
4485 None,
4486 None,
4487 None,
4488 );
4489 assert_eq!(mode, "in-memory");
4490 }
4491
4492 #[test]
4493 fn test_from_env_vars_explicit_overrides_data_dir() {
4494 let (config, mode) = EventStoreConfig::from_env_vars(
4495 Some("/app/data".to_string()),
4496 Some("/override/storage".to_string()),
4497 Some("/override/wal".to_string()),
4498 None,
4499 None,
4500 None,
4501 None,
4502 None,
4503 );
4504 assert_eq!(mode, "wal+parquet");
4505 assert_eq!(
4506 config.storage_dir.unwrap().to_str().unwrap(),
4507 "/override/storage"
4508 );
4509 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/override/wal");
4510 }
4511
4512 #[test]
4513 fn test_from_env_vars_wal_only() {
4514 let (config, mode) = EventStoreConfig::from_env_vars(
4515 None,
4516 None,
4517 Some("/wal/only".to_string()),
4518 None,
4519 None,
4520 None,
4521 None,
4522 None,
4523 );
4524 assert_eq!(mode, "wal-only");
4525 assert!(config.storage_dir.is_none());
4526 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/wal/only");
4527 }
4528
4529 #[test]
4530 fn test_from_env_vars_cache_bytes_parses_decimal() {
4531 let (config, _) = EventStoreConfig::from_env_vars(
4532 Some("/app/data".to_string()),
4533 None,
4534 None,
4535 None,
4536 Some("536870912".to_string()),
4537 None,
4539 None,
4540 None,
4541 );
4542 assert_eq!(config.cache_byte_budget, Some(536_870_912));
4543 }
4544
4545 #[test]
4546 fn test_from_env_vars_cache_bytes_unparseable_disables_budget() {
4547 let (config, _) = EventStoreConfig::from_env_vars(
4551 Some("/app/data".to_string()),
4552 None,
4553 None,
4554 None,
4555 Some("not-a-number".to_string()),
4556 None,
4557 None,
4558 None,
4559 );
4560 assert_eq!(config.cache_byte_budget, None);
4561 }
4562
4563 #[test]
4564 fn test_from_env_vars_cache_bytes_empty_disables_budget() {
4565 let (config, _) = EventStoreConfig::from_env_vars(
4566 Some("/app/data".to_string()),
4567 None,
4568 None,
4569 None,
4570 Some(String::new()),
4571 None,
4572 None,
4573 None,
4574 );
4575 assert_eq!(config.cache_byte_budget, None);
4576 }
4577
4578 #[test]
4579 fn test_from_env_vars_snapshot_interval_overrides_default() {
4580 let (config, _) = EventStoreConfig::from_env_vars(
4584 Some("/app/data".to_string()),
4585 None,
4586 None,
4587 None,
4588 None,
4589 Some("60".to_string()),
4590 None,
4591 None,
4592 );
4593 assert_eq!(config.compaction_config.compaction_interval_seconds, 60);
4594 }
4595
4596 #[test]
4597 fn test_from_env_vars_snapshot_interval_default_is_hourly() {
4598 let (config, _) = EventStoreConfig::from_env_vars(
4599 Some("/app/data".to_string()),
4600 None,
4601 None,
4602 None,
4603 None,
4604 None,
4605 None,
4606 None,
4607 );
4608 assert_eq!(config.compaction_config.compaction_interval_seconds, 3600);
4609 }
4610
4611 #[test]
4612 fn test_from_env_vars_snapshot_interval_unparseable_falls_back() {
4613 let (config, _) = EventStoreConfig::from_env_vars(
4614 Some("/app/data".to_string()),
4615 None,
4616 None,
4617 None,
4618 None,
4619 Some("not-a-number".to_string()),
4620 None,
4621 None,
4622 );
4623 assert_eq!(config.compaction_config.compaction_interval_seconds, 3600);
4624 }
4625
4626 #[test]
4627 fn test_from_env_vars_retention_system_days_overrides_default() {
4628 let (config, _) = EventStoreConfig::from_env_vars(
4631 Some("/app/data".to_string()),
4632 None,
4633 None,
4634 None,
4635 None,
4636 None,
4637 Some("7".to_string()),
4638 None,
4639 );
4640 let ttl = config
4641 .compaction_config
4642 .retention
4643 .ttl_for("system")
4644 .unwrap();
4645 assert_eq!(ttl.as_secs(), 7 * 24 * 3600);
4646 }
4647
4648 #[test]
4649 fn test_from_env_vars_retention_default_is_30_days_for_system() {
4650 let (config, _) = EventStoreConfig::from_env_vars(
4651 Some("/app/data".to_string()),
4652 None,
4653 None,
4654 None,
4655 None,
4656 None,
4657 None,
4658 None,
4659 );
4660 let ttl = config
4661 .compaction_config
4662 .retention
4663 .ttl_for("system")
4664 .unwrap();
4665 assert_eq!(ttl.as_secs(), 30 * 24 * 3600);
4666 assert!(config.compaction_config.retention.ttl_for("acme").is_none());
4668 }
4669
4670 #[test]
4671 fn test_store_stats_serde() {
4672 let stats = StoreStats {
4673 total_events: 100,
4674 total_entities: 50,
4675 total_event_types: 10,
4676 total_ingested: 100,
4677 };
4678
4679 let json = serde_json::to_string(&stats).unwrap();
4680 assert!(json.contains("\"total_events\":100"));
4681 assert!(json.contains("\"total_entities\":50"));
4682 }
4683
4684 #[test]
4685 fn test_query_with_entity_and_type() {
4686 let store = EventStore::new();
4687
4688 store
4689 .ingest(&create_test_event("entity-1", "user.created"))
4690 .unwrap();
4691 store
4692 .ingest(&create_test_event("entity-1", "user.updated"))
4693 .unwrap();
4694 store
4695 .ingest(&create_test_event("entity-2", "user.created"))
4696 .unwrap();
4697
4698 let results = store
4699 .query(&QueryEventsRequest {
4700 entity_id: Some("entity-1".to_string()),
4701 event_type: Some("user.created".to_string()),
4702 tenant_id: None,
4703 as_of: None,
4704 since: None,
4705 until: None,
4706 limit: None,
4707 event_type_prefix: None,
4708 exclude_event_type_prefix: None,
4709 payload_filter: None,
4710 })
4711 .unwrap();
4712
4713 assert_eq!(results.len(), 1);
4714 assert_eq!(results[0].event_type_str(), "user.created");
4715 }
4716
4717 #[test]
4718 fn test_query_by_event_type_prefix() {
4719 let store = EventStore::new();
4720
4721 store
4723 .ingest(&create_test_event("entity-1", "index.created"))
4724 .unwrap();
4725 store
4726 .ingest(&create_test_event("entity-2", "index.updated"))
4727 .unwrap();
4728 store
4729 .ingest(&create_test_event("entity-3", "trade.created"))
4730 .unwrap();
4731 store
4732 .ingest(&create_test_event("entity-4", "trade.completed"))
4733 .unwrap();
4734 store
4735 .ingest(&create_test_event("entity-5", "balance.updated"))
4736 .unwrap();
4737
4738 let results = store
4740 .query(&QueryEventsRequest {
4741 entity_id: None,
4742 event_type: None,
4743 tenant_id: None,
4744 as_of: None,
4745 since: None,
4746 until: None,
4747 limit: None,
4748 event_type_prefix: Some("index.".to_string()),
4749 exclude_event_type_prefix: None,
4750 payload_filter: None,
4751 })
4752 .unwrap();
4753
4754 assert_eq!(results.len(), 2);
4755 assert!(
4756 results
4757 .iter()
4758 .all(|e| e.event_type_str().starts_with("index."))
4759 );
4760 }
4761
4762 #[test]
4763 fn test_query_by_event_type_prefix_empty_returns_all() {
4764 let store = EventStore::new();
4765
4766 store
4767 .ingest(&create_test_event("entity-1", "index.created"))
4768 .unwrap();
4769 store
4770 .ingest(&create_test_event("entity-2", "trade.created"))
4771 .unwrap();
4772
4773 let results = store
4775 .query(&QueryEventsRequest {
4776 entity_id: None,
4777 event_type: None,
4778 tenant_id: None,
4779 as_of: None,
4780 since: None,
4781 until: None,
4782 limit: None,
4783 event_type_prefix: Some(String::new()),
4784 exclude_event_type_prefix: None,
4785 payload_filter: None,
4786 })
4787 .unwrap();
4788
4789 assert_eq!(results.len(), 2);
4790 }
4791
4792 #[test]
4793 fn test_query_by_event_type_prefix_no_match() {
4794 let store = EventStore::new();
4795
4796 store
4797 .ingest(&create_test_event("entity-1", "index.created"))
4798 .unwrap();
4799
4800 let results = store
4801 .query(&QueryEventsRequest {
4802 entity_id: None,
4803 event_type: None,
4804 tenant_id: None,
4805 as_of: None,
4806 since: None,
4807 until: None,
4808 limit: None,
4809 event_type_prefix: Some("nonexistent.".to_string()),
4810 exclude_event_type_prefix: None,
4811 payload_filter: None,
4812 })
4813 .unwrap();
4814
4815 assert!(results.is_empty());
4816 }
4817
4818 #[test]
4819 fn test_query_by_entity_with_type_prefix() {
4820 let store = EventStore::new();
4821
4822 store
4823 .ingest(&create_test_event("entity-1", "index.created"))
4824 .unwrap();
4825 store
4826 .ingest(&create_test_event("entity-1", "trade.created"))
4827 .unwrap();
4828 store
4829 .ingest(&create_test_event("entity-2", "index.updated"))
4830 .unwrap();
4831
4832 let results = store
4834 .query(&QueryEventsRequest {
4835 entity_id: Some("entity-1".to_string()),
4836 event_type: None,
4837 tenant_id: None,
4838 as_of: None,
4839 since: None,
4840 until: None,
4841 limit: None,
4842 event_type_prefix: Some("index.".to_string()),
4843 exclude_event_type_prefix: None,
4844 payload_filter: None,
4845 })
4846 .unwrap();
4847
4848 assert_eq!(results.len(), 1);
4849 assert_eq!(results[0].event_type_str(), "index.created");
4850 }
4851
4852 #[test]
4853 fn test_query_prefix_with_limit() {
4854 let store = EventStore::new();
4855
4856 for i in 0..5 {
4857 store
4858 .ingest(&create_test_event(&format!("entity-{i}"), "index.created"))
4859 .unwrap();
4860 }
4861
4862 let results = store
4863 .query(&QueryEventsRequest {
4864 entity_id: None,
4865 event_type: None,
4866 tenant_id: None,
4867 as_of: None,
4868 since: None,
4869 until: None,
4870 limit: Some(3),
4871 event_type_prefix: Some("index.".to_string()),
4872 exclude_event_type_prefix: None,
4873 payload_filter: None,
4874 })
4875 .unwrap();
4876
4877 assert_eq!(results.len(), 3);
4878 }
4879
4880 #[test]
4881 fn test_query_prefix_alongside_existing_filters() {
4882 let store = EventStore::new();
4883
4884 store
4885 .ingest(&create_test_event("entity-1", "index.created"))
4886 .unwrap();
4887 std::thread::sleep(std::time::Duration::from_millis(10));
4889 store
4890 .ingest(&create_test_event("entity-2", "index.strategy.updated"))
4891 .unwrap();
4892 std::thread::sleep(std::time::Duration::from_millis(10));
4893 store
4894 .ingest(&create_test_event("entity-3", "index.deleted"))
4895 .unwrap();
4896
4897 let results = store
4899 .query(&QueryEventsRequest {
4900 entity_id: None,
4901 event_type: None,
4902 tenant_id: None,
4903 as_of: None,
4904 since: None,
4905 until: None,
4906 limit: Some(2),
4907 event_type_prefix: Some("index.".to_string()),
4908 exclude_event_type_prefix: None,
4909 payload_filter: None,
4910 })
4911 .unwrap();
4912
4913 assert_eq!(results.len(), 2);
4914 }
4915
4916 #[test]
4917 fn test_query_with_payload_filter() {
4918 let store = EventStore::new();
4919
4920 for i in 0..5 {
4922 store
4923 .ingest(&create_test_event_with_payload(
4924 &format!("entity-{i}"),
4925 "user.action",
4926 serde_json::json!({"user_id": "alice", "action": "click"}),
4927 ))
4928 .unwrap();
4929 }
4930 for i in 5..10 {
4932 store
4933 .ingest(&create_test_event_with_payload(
4934 &format!("entity-{i}"),
4935 "user.action",
4936 serde_json::json!({"user_id": "bob", "action": "view"}),
4937 ))
4938 .unwrap();
4939 }
4940
4941 let results = store
4943 .query(&QueryEventsRequest {
4944 entity_id: None,
4945 event_type: Some("user.action".to_string()),
4946 tenant_id: None,
4947 as_of: None,
4948 since: None,
4949 until: None,
4950 limit: None,
4951 event_type_prefix: None,
4952 exclude_event_type_prefix: None,
4953 payload_filter: Some(r#"{"user_id":"alice"}"#.to_string()),
4954 })
4955 .unwrap();
4956
4957 assert_eq!(results.len(), 5);
4958 }
4959
4960 #[test]
4961 fn test_query_payload_filter_non_existent_field() {
4962 let store = EventStore::new();
4963
4964 store
4965 .ingest(&create_test_event_with_payload(
4966 "entity-1",
4967 "user.action",
4968 serde_json::json!({"user_id": "alice"}),
4969 ))
4970 .unwrap();
4971
4972 let results = store
4974 .query(&QueryEventsRequest {
4975 entity_id: None,
4976 event_type: None,
4977 tenant_id: None,
4978 as_of: None,
4979 since: None,
4980 until: None,
4981 limit: None,
4982 event_type_prefix: None,
4983 exclude_event_type_prefix: None,
4984 payload_filter: Some(r#"{"nonexistent":"value"}"#.to_string()),
4985 })
4986 .unwrap();
4987
4988 assert!(results.is_empty());
4989 }
4990
4991 #[test]
4992 fn test_query_payload_filter_with_prefix() {
4993 let store = EventStore::new();
4994
4995 store
4996 .ingest(&create_test_event_with_payload(
4997 "entity-1",
4998 "index.created",
4999 serde_json::json!({"status": "active"}),
5000 ))
5001 .unwrap();
5002 store
5003 .ingest(&create_test_event_with_payload(
5004 "entity-2",
5005 "index.created",
5006 serde_json::json!({"status": "inactive"}),
5007 ))
5008 .unwrap();
5009 store
5010 .ingest(&create_test_event_with_payload(
5011 "entity-3",
5012 "trade.created",
5013 serde_json::json!({"status": "active"}),
5014 ))
5015 .unwrap();
5016
5017 let results = store
5019 .query(&QueryEventsRequest {
5020 entity_id: None,
5021 event_type: None,
5022 tenant_id: None,
5023 as_of: None,
5024 since: None,
5025 until: None,
5026 limit: None,
5027 event_type_prefix: Some("index.".to_string()),
5028 exclude_event_type_prefix: None,
5029 payload_filter: Some(r#"{"status":"active"}"#.to_string()),
5030 })
5031 .unwrap();
5032
5033 assert_eq!(results.len(), 1);
5034 assert_eq!(results[0].entity_id().to_string(), "entity-1");
5035 }
5036
5037 #[test]
5038 fn test_flush_storage_no_storage() {
5039 let store = EventStore::new();
5040 let result = store.flush_storage();
5042 assert!(result.is_ok());
5043 }
5044
5045 #[test]
5046 fn test_state_evolution() {
5047 let store = EventStore::new();
5048
5049 store
5051 .ingest(
5052 &Event::from_strings(
5053 "user.created".to_string(),
5054 "user-1".to_string(),
5055 "default".to_string(),
5056 serde_json::json!({"name": "Alice", "age": 25}),
5057 None,
5058 )
5059 .unwrap(),
5060 )
5061 .unwrap();
5062
5063 store
5065 .ingest(
5066 &Event::from_strings(
5067 "user.updated".to_string(),
5068 "user-1".to_string(),
5069 "default".to_string(),
5070 serde_json::json!({"age": 26}),
5071 None,
5072 )
5073 .unwrap(),
5074 )
5075 .unwrap();
5076
5077 let state = store.reconstruct_state("user-1", None).unwrap();
5078 assert_eq!(state["current_state"]["name"], "Alice");
5080 assert_eq!(state["current_state"]["age"], 26);
5081 }
5082
5083 #[test]
5084 fn test_reject_system_event_types() {
5085 let store = EventStore::new();
5086
5087 let event = Event::reconstruct_from_strings(
5089 uuid::Uuid::new_v4(),
5090 "_system.tenant.created".to_string(),
5091 "_system:tenant:acme".to_string(),
5092 "_system".to_string(),
5093 serde_json::json!({"name": "ACME"}),
5094 chrono::Utc::now(),
5095 None,
5096 1,
5097 );
5098
5099 let result = store.ingest(&event);
5100 assert!(result.is_err());
5101 let err = result.unwrap_err();
5102 assert!(
5103 err.to_string().contains("reserved for internal use"),
5104 "Expected system namespace rejection, got: {err}"
5105 );
5106 }
5107
5108 #[test]
5116 fn test_wal_recovery_checkpoints_to_parquet() {
5117 let data_dir = TempDir::new().unwrap();
5118 let storage_dir = data_dir.path().join("storage");
5119 let wal_dir = data_dir.path().join("wal");
5120
5121 {
5123 let config = EventStoreConfig::production(
5124 &storage_dir,
5125 &wal_dir,
5126 SnapshotConfig::default(),
5127 WALConfig {
5128 sync_on_write: true,
5129 ..WALConfig::default()
5130 },
5131 CompactionConfig::default(),
5132 );
5133 let store = EventStore::with_config(config);
5134
5135 for i in 0..5 {
5136 let event = Event::from_strings(
5137 "test.created".to_string(),
5138 format!("entity-{i}"),
5139 "default".to_string(),
5140 serde_json::json!({"index": i}),
5141 None,
5142 )
5143 .unwrap();
5144 store.ingest(&event).unwrap();
5145 }
5146
5147 assert_eq!(store.stats().total_events, 5);
5148
5149 }
5152
5153 let wal_files: Vec<_> = std::fs::read_dir(&wal_dir)
5155 .unwrap()
5156 .filter_map(std::result::Result::ok)
5157 .filter(|e| e.path().extension().is_some_and(|ext| ext == "log"))
5158 .collect();
5159 assert!(!wal_files.is_empty(), "WAL file should exist");
5160 let wal_size = wal_files[0].metadata().unwrap().len();
5161 assert!(wal_size > 0, "WAL file should have data (got 0 bytes)");
5162
5163 {
5165 let config = EventStoreConfig::production(
5166 &storage_dir,
5167 &wal_dir,
5168 SnapshotConfig::default(),
5169 WALConfig {
5170 sync_on_write: true,
5171 ..WALConfig::default()
5172 },
5173 CompactionConfig::default(),
5174 );
5175 let store = EventStore::with_config(config);
5176
5177 assert_eq!(
5179 store.stats().total_events,
5180 5,
5181 "Session 2 should have all 5 events after WAL recovery"
5182 );
5183
5184 let parquet_files = find_parquet_files(&storage_dir);
5188 assert!(
5189 !parquet_files.is_empty(),
5190 "Parquet file should exist after WAL checkpoint"
5191 );
5192 }
5193
5194 {
5197 let config = EventStoreConfig::production(
5198 &storage_dir,
5199 &wal_dir,
5200 SnapshotConfig::default(),
5201 WALConfig {
5202 sync_on_write: true,
5203 ..WALConfig::default()
5204 },
5205 CompactionConfig::default(),
5206 );
5207 let store = EventStore::with_config(config);
5208
5209 assert_eq!(
5213 store.stats().total_events,
5214 0,
5215 "Session 3 boot should not pre-load Parquet (lazy-load mode)"
5216 );
5217
5218 store.ensure_tenant_loaded("default").unwrap();
5221 assert_eq!(
5222 store.stats().total_events,
5223 5,
5224 "Session 3 should have all 5 events after ensure_tenant_loaded"
5225 );
5226 }
5227 }
5228
5229 #[test]
5230 fn test_parquet_restore_surfaces_errors_not_silent() {
5231 let data_dir = TempDir::new().unwrap();
5235 let storage_dir = data_dir.path().join("storage");
5236 let wal_dir = data_dir.path().join("wal");
5237
5238 {
5240 let config = EventStoreConfig::production(
5241 &storage_dir,
5242 &wal_dir,
5243 SnapshotConfig::default(),
5244 WALConfig {
5245 sync_on_write: true,
5246 ..WALConfig::default()
5247 },
5248 CompactionConfig::default(),
5249 );
5250 let store = EventStore::with_config(config);
5251
5252 for i in 0..3 {
5253 let event = Event::from_strings(
5254 "test.created".to_string(),
5255 format!("entity-{i}"),
5256 "default".to_string(),
5257 serde_json::json!({"i": i}),
5258 None,
5259 )
5260 .unwrap();
5261 store.ingest(&event).unwrap();
5262 }
5263
5264 store.flush_storage().unwrap();
5265 assert_eq!(store.stats().total_events, 3);
5266 }
5267
5268 let parquet_files = find_parquet_files(&storage_dir);
5271 assert!(!parquet_files.is_empty(), "Parquet file must exist");
5272
5273 std::fs::write(&parquet_files[0], b"corrupted data").unwrap();
5275
5276 for entry in std::fs::read_dir(&wal_dir).unwrap().flatten() {
5278 std::fs::write(entry.path(), b"").unwrap();
5279 }
5280
5281 {
5288 let config = EventStoreConfig::production(
5289 &storage_dir,
5290 &wal_dir,
5291 SnapshotConfig::default(),
5292 WALConfig::default(),
5293 CompactionConfig::default(),
5294 );
5295 let store = EventStore::with_config(config);
5296
5297 assert_eq!(store.stats().total_events, 0);
5300 }
5301 }
5302
5303 fn count_wal_entries(wal_dir: &std::path::Path) -> usize {
5313 use std::io::{BufRead, BufReader};
5314 let mut total = 0usize;
5315 let Ok(entries) = std::fs::read_dir(wal_dir) else {
5316 return 0;
5317 };
5318 for entry in entries.flatten() {
5319 let path = entry.path();
5320 if path.extension().is_none_or(|e| e != "log") {
5321 continue;
5322 }
5323 let Ok(file) = std::fs::File::open(&path) else {
5324 continue;
5325 };
5326 for line in BufReader::new(file)
5327 .lines()
5328 .map_while(std::result::Result::ok)
5329 {
5330 if !line.trim().is_empty() {
5331 total += 1;
5332 }
5333 }
5334 }
5335 total
5336 }
5337
5338 #[test]
5339 fn test_checkpoint_truncates_wal_after_flush() {
5340 let data_dir = TempDir::new().unwrap();
5345 let storage_dir = data_dir.path().join("storage");
5346 let wal_dir = data_dir.path().join("wal");
5347
5348 let config = EventStoreConfig::production(
5349 &storage_dir,
5350 &wal_dir,
5351 SnapshotConfig::default(),
5352 WALConfig {
5353 sync_on_write: true,
5354 ..WALConfig::default()
5355 },
5356 CompactionConfig::default(),
5357 );
5358 let store = EventStore::with_config(config);
5359
5360 for i in 0..10 {
5361 let event = Event::from_strings(
5362 "test.created".to_string(),
5363 format!("entity-{i}"),
5364 "default".to_string(),
5365 serde_json::json!({"i": i}),
5366 None,
5367 )
5368 .unwrap();
5369 store.ingest(&event).unwrap();
5370 }
5371
5372 assert_eq!(
5374 count_wal_entries(&wal_dir),
5375 10,
5376 "WAL should have 10 events before checkpoint"
5377 );
5378
5379 store.checkpoint().unwrap();
5380
5381 assert_eq!(
5382 count_wal_entries(&wal_dir),
5383 0,
5384 "WAL should be empty after successful checkpoint"
5385 );
5386 let parquet_files = find_parquet_files(&storage_dir);
5387 assert!(!parquet_files.is_empty(), "Parquet should hold the events");
5388 }
5389
5390 #[test]
5391 fn test_replay_only_post_checkpoint_events_after_crash() {
5392 let data_dir = TempDir::new().unwrap();
5399 let storage_dir = data_dir.path().join("storage");
5400 let wal_dir = data_dir.path().join("wal");
5401
5402 let config_factory = || {
5403 EventStoreConfig::production(
5404 &storage_dir,
5405 &wal_dir,
5406 SnapshotConfig::default(),
5407 WALConfig {
5408 sync_on_write: true,
5409 ..WALConfig::default()
5410 },
5411 CompactionConfig::default(),
5412 )
5413 };
5414
5415 const N: usize = 50;
5418 const K: usize = 5;
5419 {
5420 let store = EventStore::with_config(config_factory());
5421 for i in 0..N {
5422 store
5423 .ingest(
5424 &Event::from_strings(
5425 "pre.checkpoint".to_string(),
5426 format!("e-{i}"),
5427 "default".to_string(),
5428 serde_json::json!({"i": i}),
5429 None,
5430 )
5431 .unwrap(),
5432 )
5433 .unwrap();
5434 }
5435 store.checkpoint().unwrap();
5436 assert_eq!(
5437 count_wal_entries(&wal_dir),
5438 0,
5439 "WAL should be empty immediately after checkpoint"
5440 );
5441
5442 for i in 0..K {
5443 store
5444 .ingest(
5445 &Event::from_strings(
5446 "post.checkpoint".to_string(),
5447 format!("p-{i}"),
5448 "default".to_string(),
5449 serde_json::json!({"i": i}),
5450 None,
5451 )
5452 .unwrap(),
5453 )
5454 .unwrap();
5455 }
5456 assert_eq!(
5457 count_wal_entries(&wal_dir),
5458 K,
5459 "WAL should hold only post-checkpoint events"
5460 );
5461 }
5463
5464 {
5468 let store = EventStore::with_config(config_factory());
5469 assert_eq!(
5473 store.stats().total_events,
5474 K,
5475 "Boot should replay exactly K events from WAL (the post-checkpoint window), not N+K"
5476 );
5477
5478 store.ensure_tenant_loaded("default").unwrap();
5480 assert_eq!(
5481 store.stats().total_events,
5482 N + K,
5483 "After lazy-load, both pre- and post-checkpoint events should be reachable"
5484 );
5485 }
5486 }
5487
5488 #[test]
5489 fn test_checkpoint_is_idempotent() {
5490 let data_dir = TempDir::new().unwrap();
5493 let storage_dir = data_dir.path().join("storage");
5494 let wal_dir = data_dir.path().join("wal");
5495
5496 let store = EventStore::with_config(EventStoreConfig::production(
5497 &storage_dir,
5498 &wal_dir,
5499 SnapshotConfig::default(),
5500 WALConfig::default(),
5501 CompactionConfig::default(),
5502 ));
5503
5504 for i in 0..5 {
5505 store
5506 .ingest(
5507 &Event::from_strings(
5508 "x".to_string(),
5509 format!("e-{i}"),
5510 "default".to_string(),
5511 serde_json::json!({}),
5512 None,
5513 )
5514 .unwrap(),
5515 )
5516 .unwrap();
5517 }
5518
5519 store.checkpoint().unwrap();
5520 store.checkpoint().unwrap();
5522 assert_eq!(count_wal_entries(&wal_dir), 0);
5523 }
5524
5525 #[test]
5526 fn test_checkpoint_noop_in_memory_only_mode() {
5527 let store = EventStore::new();
5529 store.checkpoint().unwrap();
5530 }
5531
5532 #[test]
5533 fn test_checkpoint_interval_from_env_defaults_to_60s_when_wal_enabled() {
5534 let (config, _) = EventStoreConfig::from_env_vars(
5535 Some("/app/data".to_string()),
5536 None,
5537 None,
5538 None,
5539 None,
5540 None,
5541 None,
5542 None,
5543 );
5544 assert_eq!(config.checkpoint_interval_secs, Some(60));
5545 }
5546
5547 #[test]
5548 fn test_checkpoint_interval_from_env_overrides_default() {
5549 let (config, _) = EventStoreConfig::from_env_vars(
5550 Some("/app/data".to_string()),
5551 None,
5552 None,
5553 None,
5554 None,
5555 None,
5556 None,
5557 Some("15".to_string()),
5558 );
5559 assert_eq!(config.checkpoint_interval_secs, Some(15));
5560 }
5561
5562 #[test]
5563 fn test_checkpoint_interval_disabled_when_wal_disabled() {
5564 let (config, _) = EventStoreConfig::from_env_vars(
5566 Some("/app/data".to_string()),
5567 None,
5568 None,
5569 Some("false".to_string()),
5570 None,
5571 None,
5572 None,
5573 Some("15".to_string()),
5574 );
5575 assert_eq!(config.checkpoint_interval_secs, None);
5576 }
5577
5578 #[test]
5579 fn test_checkpoint_interval_unparseable_falls_back_to_default() {
5580 let (config, _) = EventStoreConfig::from_env_vars(
5581 Some("/app/data".to_string()),
5582 None,
5583 None,
5584 None,
5585 None,
5586 None,
5587 None,
5588 Some("not-a-number".to_string()),
5589 );
5590 assert_eq!(config.checkpoint_interval_secs, Some(60));
5591 }
5592
5593 fn seed_two_tenants() -> EventStore {
5600 let store = EventStore::new();
5601
5602 for (entity, etype, payload) in [
5604 (
5605 "a-1",
5606 "created",
5607 serde_json::json!({"colour": "red", "size": 1}),
5608 ),
5609 ("a-1", "updated", serde_json::json!({"colour": "blue"})),
5610 ("a-2", "created", serde_json::json!({"colour": "green"})),
5611 ] {
5612 store
5613 .ingest(
5614 &Event::from_strings(
5615 etype.to_string(),
5616 entity.to_string(),
5617 "alice".to_string(),
5618 payload,
5619 None,
5620 )
5621 .unwrap(),
5622 )
5623 .unwrap();
5624 }
5625
5626 store
5628 .ingest(
5629 &Event::from_strings(
5630 "created".to_string(),
5631 "a-1".to_string(),
5632 "bob".to_string(),
5633 serde_json::json!({"colour": "BOB_SECRET", "bob_only": true}),
5634 None,
5635 )
5636 .unwrap(),
5637 )
5638 .unwrap();
5639
5640 store
5641 }
5642
5643 #[test]
5644 fn test_stats_for_tenant_counts_only_that_tenant() {
5645 let store = seed_two_tenants();
5646
5647 let alice = store.stats_for_tenant("alice");
5648 assert_eq!(alice.total_events, 3);
5649 assert_eq!(alice.total_entities, 2);
5650 assert_eq!(alice.total_event_types, 2);
5651 assert_eq!(alice.event_types.get("created"), Some(&2));
5652 assert_eq!(alice.event_types.get("updated"), Some(&1));
5653
5654 let bob = store.stats_for_tenant("bob");
5655 assert_eq!(bob.total_events, 1);
5656 assert_eq!(bob.total_entities, 1);
5657 assert_eq!(bob.total_event_types, 1);
5658 }
5659
5660 #[test]
5661 fn test_stats_for_tenant_never_reports_global_totals() {
5662 let store = seed_two_tenants();
5663
5664 assert_eq!(store.stats().total_events, 4);
5666 assert_eq!(store.stats_for_tenant("alice").total_events, 3);
5667 assert_eq!(store.stats_for_tenant("bob").total_events, 1);
5668
5669 assert_eq!(store.stats_for_tenant("bob").total_ingested, 1);
5672 }
5673
5674 #[test]
5675 fn test_stats_for_tenant_unknown_tenant_is_empty_not_global() {
5676 let store = seed_two_tenants();
5677 let nobody = store.stats_for_tenant("does-not-exist");
5678
5679 assert_eq!(nobody.total_events, 0);
5680 assert_eq!(nobody.total_entities, 0);
5681 assert!(nobody.event_types.is_empty());
5682 assert!(nobody.oldest_event.is_none());
5683 assert!(nobody.newest_event.is_none());
5684 }
5685
5686 #[test]
5687 fn test_stats_for_tenant_reports_time_range() {
5688 let store = seed_two_tenants();
5689 let alice = store.stats_for_tenant("alice");
5690
5691 let oldest = alice.oldest_event.expect("oldest");
5692 let newest = alice.newest_event.expect("newest");
5693 assert!(oldest <= newest);
5694 }
5695
5696 #[test]
5697 fn test_reconstruct_state_for_tenant_isolates_shared_entity_id() {
5698 let store = seed_two_tenants();
5699
5700 let alice = store
5702 .reconstruct_state_for_tenant("a-1", None, "alice")
5703 .unwrap();
5704 let alice_state = alice.get("current_state").unwrap();
5705 assert_eq!(alice_state.get("colour").unwrap(), "blue"); assert_eq!(alice_state.get("size").unwrap(), 1);
5707 assert!(
5708 alice_state.get("bob_only").is_none(),
5709 "alice must not see bob's payload keys: {alice_state:?}"
5710 );
5711 assert_eq!(alice.get("event_count").unwrap(), 2);
5712
5713 let bob = store
5714 .reconstruct_state_for_tenant("a-1", None, "bob")
5715 .unwrap();
5716 let bob_state = bob.get("current_state").unwrap();
5717 assert_eq!(bob_state.get("colour").unwrap(), "BOB_SECRET");
5718 assert_eq!(bob.get("event_count").unwrap(), 1);
5719 }
5720
5721 #[test]
5722 fn test_reconstruct_state_for_tenant_rejects_another_tenants_entity() {
5723 let store = seed_two_tenants();
5724
5725 assert!(
5727 store
5728 .reconstruct_state_for_tenant("a-2", None, "alice")
5729 .is_ok()
5730 );
5731 assert!(
5732 store
5733 .reconstruct_state_for_tenant("a-2", None, "bob")
5734 .is_err(),
5735 "bob must not be able to read alice's entity"
5736 );
5737 }
5738
5739 #[test]
5740 fn test_global_reconstruct_state_still_spans_tenants() {
5741 let store = seed_two_tenants();
5744 let all = store.reconstruct_state("a-1", None).unwrap();
5745 assert_eq!(all.get("event_count").unwrap(), 3);
5746 }
5747
5748 fn seed_hot_entity(history: usize) -> EventStore {
5759 let store = EventStore::new();
5760 for _ in 0..history {
5761 store
5762 .ingest(&create_test_event("entity-hot", "user.updated"))
5763 .unwrap();
5764 }
5765 store
5766 }
5767
5768 fn hot_request(limit: Option<usize>) -> QueryEventsRequest {
5769 QueryEventsRequest {
5770 entity_id: Some("entity-hot".to_string()),
5771 limit,
5772 ..QueryEventsRequest::default()
5773 }
5774 }
5775
5776 #[test]
5777 fn query_window_limit_bounds_materialization_not_just_the_response() {
5778 const HISTORY: usize = 400;
5779 let store = seed_hot_entity(HISTORY);
5780
5781 let ((page, total), materialized) = crate::clone_probe::measure(|| {
5783 store.query_window(&hot_request(Some(1)), 0, true).unwrap()
5784 });
5785 assert_eq!(page.len(), 1, "limit=1 returns one event");
5786 assert_eq!(total, HISTORY, "total is still the full match count");
5787 assert_eq!(
5788 materialized, 1,
5789 "limit=1 cloned {materialized} of {HISTORY} events: `limit` must \
5790 bound what the store materializes, not just what it returns"
5791 );
5792
5793 let ((page, _), materialized) = crate::clone_probe::measure(|| {
5795 store
5796 .query_window(&hot_request(Some(5)), 300, false)
5797 .unwrap()
5798 });
5799 assert_eq!(page.len(), 5);
5800 assert_eq!(
5801 materialized, 5,
5802 "offset=300&limit=5 cloned {materialized} events, expected 5"
5803 );
5804
5805 let ((all, _), materialized) = crate::clone_probe::measure(|| {
5808 store.query_window(&hot_request(None), 0, true).unwrap()
5809 });
5810 assert_eq!(all.len(), HISTORY);
5811 assert_eq!(materialized, HISTORY as u64);
5812 }
5813
5814 #[test]
5815 fn query_window_cost_does_not_grow_with_entity_history() {
5816 let short = seed_hot_entity(40);
5821 let long = seed_hot_entity(400);
5822
5823 let (_, short_cost) = crate::clone_probe::measure(|| {
5824 short.query_window(&hot_request(Some(1)), 0, true).unwrap()
5825 });
5826 let (_, long_cost) = crate::clone_probe::measure(|| {
5827 long.query_window(&hot_request(Some(1)), 0, true).unwrap()
5828 });
5829
5830 assert_eq!(
5831 (short_cost, long_cost),
5832 (1, 1),
5833 "a limit=1 page materialized {short_cost} events over a 40-event \
5834 history and {long_cost} over 400 — cost is tracking history length"
5835 );
5836 }
5837
5838 #[test]
5839 fn select_window_orders_only_the_window() {
5840 use std::cell::Cell;
5846 const N: usize = 4096;
5847
5848 let comparisons = Cell::new(0usize);
5849 let order = |a: &u64, b: &u64| {
5850 comparisons.set(comparisons.get() + 1);
5851 a.cmp(b)
5852 };
5853 let shuffled = || -> Vec<u64> {
5854 (0..N as u64)
5855 .map(|i| (i * 2_654_435_761) % 1_000_003)
5856 .collect()
5857 };
5858
5859 let mut bounded = shuffled();
5860 select_window(&mut bounded, 0, Some(1), order);
5861 let bounded_cost = comparisons.replace(0);
5862 assert_eq!(bounded.len(), 1, "window of 1 keeps 1 item");
5863
5864 let mut everything = shuffled();
5865 select_window(&mut everything, 0, None, order);
5866 let full_sort_cost = comparisons.get();
5867 assert_eq!(everything.len(), N);
5868
5869 assert!(
5870 bounded_cost < 4 * N,
5871 "selecting a 1-item window out of {N} took {bounded_cost} comparisons \
5872 (~{}·N) — that is sort-shaped, not selection-shaped",
5873 bounded_cost / N
5874 );
5875 assert!(
5876 bounded_cost * 3 < full_sort_cost,
5877 "a 1-item window cost {bounded_cost} comparisons against \
5878 {full_sort_cost} for sorting all {N}: the window is not bounding \
5879 the ordering work"
5880 );
5881 }
5882
5883 #[test]
5884 fn select_window_is_equivalent_to_sort_then_window() {
5885 let order = |a: &(u32, usize), b: &(u32, usize)| a.cmp(b);
5889 let source: Vec<(u32, usize)> = [7, 3, 3, 9, 1, 3, 5, 9, 0, 2]
5890 .into_iter()
5891 .enumerate()
5892 .map(|(i, k)| (k, i))
5893 .collect();
5894
5895 let mut sorted = source.clone();
5896 sorted.sort_by(order);
5897
5898 for offset in 0..12 {
5899 for limit in [None, Some(0), Some(1), Some(3), Some(10), Some(50)] {
5900 let mut got = source.clone();
5901 select_window(&mut got, offset, limit, order);
5902 let got: Vec<_> = got.into_iter().skip(offset).collect();
5903 let expected: Vec<_> = sorted
5904 .iter()
5905 .copied()
5906 .skip(offset)
5907 .take(limit.unwrap_or(usize::MAX))
5908 .collect();
5909 assert_eq!(got, expected, "offset={offset} limit={limit:?}");
5910 }
5911 }
5912 }
5913
5914 #[test]
5925 fn query_window_applies_time_filters_on_the_full_scan_path() {
5926 let store = EventStore::new();
5927 let base = Utc::now() - chrono::Duration::hours(24);
5928 let mut ids = Vec::new();
5929 for i in 0..5i64 {
5930 let mut event = create_test_event(&format!("e-{i}"), "user.created");
5931 event.timestamp = base + chrono::Duration::hours(i);
5932 event.version = i + 1;
5933 ids.push(event.id);
5934 store.ingest(&event).unwrap();
5935 }
5936 let at = |h: i64| base + chrono::Duration::hours(h);
5937
5938 let scoped = |mutate: &dyn Fn(&mut QueryEventsRequest)| {
5941 let mut req = QueryEventsRequest {
5942 tenant_id: Some("default".to_string()),
5943 ..QueryEventsRequest::default()
5944 };
5945 mutate(&mut req);
5946 req
5947 };
5948
5949 for (label, req, expected) in [
5950 (
5951 "since=T+2 keeps only events at or after T+2",
5952 scoped(&|r| r.since = Some(at(2))),
5953 vec![ids[2], ids[3], ids[4]],
5954 ),
5955 (
5956 "until=T+1 keeps only events at or before T+1",
5957 scoped(&|r| r.until = Some(at(1))),
5958 vec![ids[0], ids[1]],
5959 ),
5960 (
5961 "as_of=T+1 is time travel: nothing newer than T+1",
5962 scoped(&|r| r.as_of = Some(at(1))),
5963 vec![ids[0], ids[1]],
5964 ),
5965 (
5966 "since+until compose into a closed window",
5967 scoped(&|r| {
5968 r.since = Some(at(1));
5969 r.until = Some(at(3));
5970 }),
5971 vec![ids[1], ids[2], ids[3]],
5972 ),
5973 ] {
5974 let (events, total) = store.query_window(&req, 0, false).unwrap();
5975 let got: Vec<_> = events.iter().map(|e| e.id).collect();
5976 assert_eq!(got, expected, "{label}");
5977 assert_eq!(
5978 total,
5979 expected.len(),
5980 "{label}: total counts the events INSIDE the window — \
5981 has_more is derived from it, so a full-history total makes a \
5982 paginator walk events the caller filtered out"
5983 );
5984 }
5985
5986 let (indexed, total) = store
5989 .query_window(
5990 &scoped(&|r| {
5991 r.entity_id = Some("e-3".to_string());
5992 r.since = Some(at(2));
5993 }),
5994 0,
5995 false,
5996 )
5997 .unwrap();
5998 assert_eq!(
5999 indexed.iter().map(|e| e.id).collect::<Vec<_>>(),
6000 vec![ids[3]]
6001 );
6002 assert_eq!(total, 1);
6003
6004 let (page, total) = store
6007 .query_window(&scoped(&|r| r.since = Some(at(2))), 1, true)
6008 .unwrap();
6009 assert_eq!(
6010 page.iter().map(|e| e.id).collect::<Vec<_>>(),
6011 vec![ids[3], ids[2]],
6012 "order=desc + offset=1 inside a since window"
6013 );
6014 assert_eq!(total, 3);
6015 }
6016
6017 #[test]
6018 fn query_window_bounded_selection_matches_a_full_sort() {
6019 let store = EventStore::new();
6024 for i in 0..50 {
6025 store
6026 .ingest(&create_test_event(&format!("e-{i:02}"), "user.created"))
6027 .unwrap();
6028 }
6029 let all = |descending: bool| {
6030 let (events, _) = store
6031 .query_window(&QueryEventsRequest::default(), 0, descending)
6032 .unwrap();
6033 events
6034 };
6035
6036 for descending in [false, true] {
6037 let reference = all(descending);
6038 for offset in [0, 1, 7, 49, 50, 100] {
6039 for limit in [1, 3, 10, 50, 100] {
6040 let (page, total) = store
6041 .query_window(
6042 &QueryEventsRequest {
6043 limit: Some(limit),
6044 ..QueryEventsRequest::default()
6045 },
6046 offset,
6047 descending,
6048 )
6049 .unwrap();
6050 let expected: Vec<_> = reference
6051 .iter()
6052 .skip(offset)
6053 .take(limit)
6054 .map(|e| e.id)
6055 .collect();
6056 let got: Vec<_> = page.iter().map(|e| e.id).collect();
6057 assert_eq!(
6058 got, expected,
6059 "desc={descending} offset={offset} limit={limit}: windowed \
6060 selection must match a full sort"
6061 );
6062 assert_eq!(total, 50, "total is always the full match count");
6063 }
6064 }
6065 }
6066 }
6067}