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 archive_budget::{ArchiveReadBudget, ArchiveReadLimits, check_cancelled},
25 compaction::{CompactionConfig, CompactionManager},
26 index::{EventIndex, IndexEntry},
27 snapshot::{SnapshotConfig, SnapshotManager, SnapshotType},
28 storage::ParquetStorage,
29 tenant_loader::TenantLoader,
30 wal::{WALConfig, WriteAheadLog},
31 },
32 query::geospatial::GeoIndex,
33 },
34};
35use chrono::{DateTime, Utc};
36use dashmap::DashMap;
37use parking_lot::RwLock;
38use std::{path::PathBuf, sync::Arc};
39#[cfg(feature = "server")]
40use tokio::sync::mpsc;
41
42#[derive(Debug, Clone, Default)]
59pub struct ReadScope {
60 allow_entity_prefixes: Option<Vec<String>>,
61}
62
63impl ReadScope {
64 pub fn unrestricted() -> Self {
66 Self {
67 allow_entity_prefixes: None,
68 }
69 }
70
71 pub fn allow_entity_prefixes<I, S>(prefixes: I) -> Self
76 where
77 I: IntoIterator<Item = S>,
78 S: Into<String>,
79 {
80 Self {
81 allow_entity_prefixes: Some(prefixes.into_iter().map(Into::into).collect()),
82 }
83 }
84
85 pub fn permits(&self, entity_id: &str) -> bool {
87 match &self.allow_entity_prefixes {
88 None => true,
89 Some(allowed) => allowed.iter().any(|p| entity_id.starts_with(p.as_str())),
90 }
91 }
92
93 pub fn is_unrestricted(&self) -> bool {
95 self.allow_entity_prefixes.is_none()
96 }
97}
98
99pub struct EventStore {
100 events: Arc<RwLock<Vec<Event>>>,
102
103 index: Arc<EventIndex>,
105
106 pub(crate) projections: Arc<RwLock<ProjectionManager>>,
108
109 storage: Option<Arc<RwLock<ParquetStorage>>>,
111
112 #[cfg(feature = "server")]
114 websocket_manager: Arc<WebSocketManager>,
115
116 snapshot_manager: Arc<SnapshotManager>,
118
119 wal: Option<Arc<WriteAheadLog>>,
121
122 compaction_manager: Option<Arc<CompactionManager>>,
124
125 schema_registry: Arc<SchemaRegistry>,
127
128 replay_manager: Arc<ReplayManager>,
130
131 pipeline_manager: Arc<PipelineManager>,
133
134 #[cfg(feature = "server")]
136 metrics: Arc<MetricsRegistry>,
137
138 total_ingested: Arc<RwLock<u64>>,
140
141 projection_state_cache: Arc<DashMap<String, serde_json::Value>>,
145
146 projection_status: Arc<DashMap<String, String>>,
149
150 #[cfg(feature = "server")]
152 webhook_registry: Arc<WebhookRegistry>,
153
154 #[cfg(feature = "server")]
156 webhook_tx: Arc<RwLock<Option<mpsc::UnboundedSender<WebhookDeliveryTask>>>>,
157
158 geo_index: Arc<GeoIndex>,
160
161 exactly_once: Arc<ExactlyOnceRegistry>,
163
164 schema_evolution: Arc<SchemaEvolutionManager>,
166
167 entity_versions: Arc<DashMap<String, u64>>,
170
171 consumer_registry: Arc<ConsumerRegistry>,
173
174 event_broadcast_tx: tokio::sync::broadcast::Sender<Arc<Event>>,
178
179 tenant_loader: Arc<TenantLoader>,
185 strict_archive_limits: ArchiveReadLimits,
186 #[cfg(feature = "server")]
187 http_archive_warmup_timeout: Option<std::time::Duration>,
188 cache_residency_gate: RwLock<()>,
191 cache_generations: DashMap<String, u64>,
192
193 checkpoint_interval_secs: Option<u64>,
199
200 read_only: bool,
207
208 durability_gate: RwLock<()>,
212
213 refresh_state: parking_lot::Mutex<refresh::RefreshState>,
214}
215
216#[cfg(feature = "server")]
218#[derive(Debug, Clone)]
219pub struct WebhookDeliveryTask {
220 pub webhook: crate::application::services::webhook::WebhookSubscription,
221 pub event: Event,
222}
223
224fn select_window<T>(
237 items: &mut Vec<T>,
238 offset: usize,
239 limit: Option<usize>,
240 order: impl Fn(&T, &T) -> std::cmp::Ordering + Copy,
241) {
242 match limit.map(|limit| offset.saturating_add(limit)) {
243 Some(0) => {
244 items.clear();
245 return;
246 }
247 Some(window_end) if window_end < items.len() => {
248 items.select_nth_unstable_by(window_end - 1, order);
249 items.truncate(window_end);
250 }
251 _ => {}
254 }
255 items.sort_unstable_by(order);
256}
257
258impl EventStore {
259 pub fn new() -> Self {
261 Self::with_config(EventStoreConfig::default())
262 }
263
264 pub fn with_config(config: EventStoreConfig) -> Self {
266 let mut projections = ProjectionManager::new();
267
268 projections.register(Arc::new(EntitySnapshotProjection::new("entity_snapshots")));
270 projections.register(Arc::new(EventCounterProjection::new("event_counters")));
271
272 let storage = config
274 .storage_dir
275 .as_ref()
276 .and_then(|dir| match ParquetStorage::new(dir) {
277 Ok(storage) => {
278 tracing::info!("✅ Parquet persistence enabled at: {}", dir.display());
279 Some(Arc::new(RwLock::new(storage)))
280 }
281 Err(e) => {
282 tracing::error!("❌ Failed to initialize Parquet storage: {}", e);
283 None
284 }
285 });
286
287 let wal = config.wal_dir.as_ref().and_then(|dir| {
289 match WriteAheadLog::new(dir, config.wal_config.clone()) {
290 Ok(wal) => {
291 tracing::info!("✅ WAL enabled at: {}", dir.display());
292 Some(Arc::new(wal))
293 }
294 Err(e) => {
295 tracing::error!("❌ Failed to initialize WAL: {}", e);
296 None
297 }
298 }
299 });
300
301 let compaction_manager = config.storage_dir.as_ref().map(|dir| {
303 let manager = CompactionManager::new(dir, config.compaction_config.clone());
304 Arc::new(manager)
305 });
306
307 let schema_registry = Arc::new(SchemaRegistry::new(config.schema_registry_config.clone()));
309 tracing::info!("✅ Schema registry enabled");
310
311 let replay_manager = Arc::new(ReplayManager::new());
313 tracing::info!("✅ Replay manager enabled");
314
315 let pipeline_manager = Arc::new(PipelineManager::new());
317 tracing::info!("✅ Pipeline manager enabled");
318
319 #[cfg(feature = "server")]
321 let metrics = {
322 let m = MetricsRegistry::new();
323 tracing::info!("✅ Prometheus metrics registry initialized");
324 m
325 };
326
327 let projection_state_cache = Arc::new(DashMap::new());
329 tracing::info!("✅ Projection state cache initialized");
330
331 #[cfg(feature = "server")]
333 let webhook_registry = {
334 let w = Arc::new(WebhookRegistry::new());
335 tracing::info!("✅ Webhook registry initialized");
336 w
337 };
338
339 let (event_broadcast_tx, _) = tokio::sync::broadcast::channel(1024);
342
343 let store = Self {
344 strict_archive_limits: config.strict_archive_limits,
345 #[cfg(feature = "server")]
346 http_archive_warmup_timeout: config.http_archive_warmup_timeout,
347 cache_residency_gate: RwLock::new(()),
348 cache_generations: DashMap::new(),
349 events: Arc::new(RwLock::new(Vec::new())),
350 index: Arc::new(EventIndex::new()),
351 projections: Arc::new(RwLock::new(projections)),
352 storage,
353 #[cfg(feature = "server")]
354 websocket_manager: Arc::new(WebSocketManager::new()),
355 snapshot_manager: Arc::new(SnapshotManager::new(config.snapshot_config)),
356 wal,
357 compaction_manager,
358 schema_registry,
359 replay_manager,
360 pipeline_manager,
361 #[cfg(feature = "server")]
362 metrics,
363 total_ingested: Arc::new(RwLock::new(0)),
364 projection_state_cache,
365 projection_status: Arc::new(DashMap::new()),
366 #[cfg(feature = "server")]
367 webhook_registry,
368 #[cfg(feature = "server")]
369 webhook_tx: Arc::new(RwLock::new(None)),
370 geo_index: Arc::new(GeoIndex::new()),
371 exactly_once: Arc::new(ExactlyOnceRegistry::new(ExactlyOnceConfig::default())),
372 schema_evolution: Arc::new(SchemaEvolutionManager::new()),
373 entity_versions: Arc::new(DashMap::new()),
374 consumer_registry: Arc::new(ConsumerRegistry::new()),
375 event_broadcast_tx,
376 tenant_loader: {
377 let loader = TenantLoader::new();
378 if let Some(budget) = config.cache_byte_budget {
379 loader.set_byte_budget(budget);
380 tracing::info!(
381 "✅ Cache byte budget set to {} bytes ({:.2} GiB) — LRU eviction enabled",
382 budget,
383 budget as f64 / (1024.0 * 1024.0 * 1024.0)
384 );
385 } else {
386 tracing::info!(
387 "✅ Cache budget unset — every loaded tenant stays resident \
388 (set ALLSOURCE_CACHE_BYTES to enable eviction)"
389 );
390 }
391 Arc::new(loader)
392 },
393 checkpoint_interval_secs: config.checkpoint_interval_secs,
394 read_only: config.read_only,
395 durability_gate: RwLock::new(()),
396 refresh_state: parking_lot::Mutex::new(refresh::RefreshState::default()),
397 };
398
399 if config.read_only {
400 tracing::info!(
401 "📖 EventStore opened READ-ONLY (replica): WAL will be replayed for reads but \
402 not truncated; writes are rejected"
403 );
404 }
405
406 if let Some(ref wal) = store.wal {
428 if config.read_only
429 && let Ok(stamps) = wal.segment_stamps()
430 {
431 store.refresh_state.lock().set_wal_stamps(stamps);
432 }
433 match wal.recover() {
434 Ok(recovered_events) if !recovered_events.is_empty() => {
435 let mut wal_new = 0usize;
436 for event in recovered_events {
437 let offset = store.events.read().len();
438 if let Err(e) = store.index.index_event(
439 event.id,
440 event.entity_id_str(),
441 event.event_type_str(),
442 event.timestamp,
443 offset,
444 ) {
445 tracing::error!("Failed to re-index WAL event {}: {}", event.id, e);
446 }
447
448 if let Err(e) = store.projections.read().process_event(&event) {
449 tracing::error!("Failed to re-process WAL event {}: {}", event.id, e);
450 }
451
452 *store
453 .entity_versions
454 .entry(event.entity_id_str().to_string())
455 .or_insert(0) += 1;
456
457 store.events.write().push(event);
458 wal_new += 1;
459 }
460
461 #[cfg(feature = "server")]
466 store.metrics.wal_replay_events_total.set(wal_new as i64);
467
468 if wal_new > 0 {
469 let total = store.events.read().len();
470 *store.total_ingested.write() = total as u64;
478 tracing::info!(
479 "✅ Recovered {} events from WAL (Parquet data stays cold until \
480 first per-tenant query)",
481 wal_new
482 );
483
484 if let Some(ref storage) = store.storage
499 && !config.read_only
500 {
501 tracing::info!(
502 "📸 Checkpointing {} WAL events to Parquet storage...",
503 wal_new
504 );
505 let parquet = storage.read();
506 let events = store.events.read();
507 let mut buffered = 0usize;
508 for event in events.iter().skip(events.len() - wal_new) {
509 if let Err(e) = parquet.append_event(event.clone()) {
510 tracing::error!(
511 "Failed to buffer WAL event for Parquet: {}",
512 e
513 );
514 } else {
515 buffered += 1;
516 }
517 }
518 drop(events);
519 drop(parquet);
520
521 if buffered < wal_new {
522 tracing::error!(
523 "Buffered {} of {} recovered WAL events; leaving the WAL \
524 in place so none are lost",
525 buffered,
526 wal_new
527 );
528 } else if buffered > 0 {
529 if let Err(e) = store.flush_storage() {
530 tracing::error!("Failed to checkpoint to Parquet: {}", e);
531 } else if let Err(e) = wal.truncate() {
532 tracing::error!(
533 "Failed to truncate WAL after checkpoint: {}",
534 e
535 );
536 } else {
537 tracing::info!(
538 "✅ WAL checkpointed and truncated ({} events)",
539 buffered
540 );
541 }
542 }
543 }
544 }
545 }
546 Ok(_) => {
547 tracing::debug!("No events to recover from WAL");
548 #[cfg(feature = "server")]
549 store.metrics.wal_replay_events_total.set(0);
550 }
551 Err(e) => {
552 tracing::error!("❌ WAL recovery failed: {}", e);
553 }
554 }
555 } else if store.storage.is_some() {
556 tracing::info!(
557 "📂 Boot complete (lazy-load mode): Parquet data stays on disk until first \
558 per-tenant query"
559 );
560 }
561
562 store
563 }
564
565 pub fn is_read_only(&self) -> bool {
567 self.read_only
568 }
569
570 fn ensure_writable(&self) -> Result<()> {
574 if self.read_only {
575 return Err(crate::error::AllSourceError::ReadOnly(
576 "this AllSource instance is a read-only replica — the data directory is owned by \
577 another running process. Stop the other process, or run a single shared writer \
578 (e.g. Prime in --mode http) and point clients at it."
579 .to_string(),
580 ));
581 }
582 Ok(())
583 }
584
585 #[cfg_attr(feature = "hotpath", hotpath::measure)]
593 pub fn ingest_with_expected_version(
594 &self,
595 event: &Event,
596 expected_version: Option<u64>,
597 ) -> Result<u64> {
598 self.ingest_with_expected_version_cancellable(event, expected_version, None)
599 }
600
601 pub(crate) fn ingest_with_expected_version_cancellable(
602 &self,
603 event: &Event,
604 expected_version: Option<u64>,
605 cancellation: Option<&Arc<std::sync::atomic::AtomicBool>>,
606 ) -> Result<u64> {
607 check_cancelled(cancellation.map(Arc::as_ref))?;
608 self.ensure_writable()?;
610
611 self.validate_event(event)?;
613 if expected_version.is_some() {
616 self.ensure_tenant_loaded_budgeted(event.tenant_id_str(), true, cancellation.cloned())?;
617 }
618
619 let _resident = self.cache_residency_gate.read();
620 if expected_version.is_some() && !self.tenant_loader.is_complete(event.tenant_id_str()) {
621 return Err(AllSourceError::StorageError(
622 "Verified archive was evicted before the conditional write".into(),
623 ));
624 }
625 let entity_id = event.entity_id_str().to_string();
626 let _durable = self.durability_gate.read();
627 check_cancelled(cancellation.map(Arc::as_ref))?;
628 let mut stored_event = event.clone();
629
630 let new_version = {
633 let mut version_entry = self.entity_versions.entry(entity_id.clone()).or_insert(0);
634 check_cancelled(cancellation.map(Arc::as_ref))?;
635 let current = *version_entry;
636
637 if let Some(expected) = expected_version
638 && current != expected
639 {
640 return Err(crate::error::AllSourceError::VersionConflict { expected, current });
641 }
642
643 let next = current.checked_add(1).ok_or_else(|| {
644 crate::error::AllSourceError::InvalidInput("Entity version exhausted".into())
645 })?;
646 stored_event.version = i64::try_from(next).map_err(|_| {
647 crate::error::AllSourceError::InvalidInput("Entity version exhausted".into())
648 })?;
649
650 if let Some(ref wal) = self.wal {
652 wal.append(stored_event.clone())?;
653 }
654
655 *version_entry = next;
656 next
657 };
658
659 self.ingest_post_wal(&stored_event)?;
662
663 Ok(new_version)
664 }
665
666 #[cfg_attr(feature = "hotpath", hotpath::measure)]
669 fn ingest_post_wal(&self, event: &Event) -> Result<()> {
670 #[cfg(feature = "server")]
671 let timer = self.metrics.ingestion_duration_seconds.start_timer();
672
673 let mut events = self.events.write();
674 let offset = events.len();
675
676 self.index.index_event(
678 event.id,
679 event.entity_id_str(),
680 event.event_type_str(),
681 event.timestamp,
682 offset,
683 )?;
684
685 let projections = self.projections.read();
687 projections.process_event(event)?;
688 drop(projections);
689
690 let pipeline_results = self.pipeline_manager.process_event(event);
692 if !pipeline_results.is_empty() {
693 tracing::debug!(
694 "Event {} processed by {} pipeline(s)",
695 event.id,
696 pipeline_results.len()
697 );
698 for (pipeline_id, result) in pipeline_results {
699 tracing::trace!("Pipeline {} result: {:?}", pipeline_id, result);
700 }
701 }
702
703 if let Some(ref storage) = self.storage {
705 let storage = storage.read();
706 storage.append_event(event.clone())?;
707 }
708
709 events.push(event.clone());
711 let total_events = events.len();
712 drop(events);
713
714 let event_arc = Arc::new(event.clone());
716 let _ = self.event_broadcast_tx.send(Arc::clone(&event_arc));
717 #[cfg(feature = "server")]
718 self.websocket_manager.broadcast_event(event_arc);
719
720 #[cfg(feature = "server")]
722 self.dispatch_webhooks(event);
723
724 self.geo_index.index_event(event);
726
727 self.schema_evolution
729 .analyze_event(event.event_type_str(), &event.payload);
730
731 self.check_auto_snapshot(event.entity_id_str(), event);
733
734 #[cfg(feature = "server")]
736 {
737 self.metrics.events_ingested_total.inc();
738 self.metrics
739 .events_ingested_by_type
740 .with_label_values(&[event.event_type_str()])
741 .inc();
742 self.metrics.storage_events_total.set(total_events as i64);
743 }
744
745 let mut total = self.total_ingested.write();
747 *total += 1;
748
749 #[cfg(feature = "server")]
750 timer.observe_duration();
751
752 tracing::debug!("Event ingested: {} (offset: {})", event.id, offset);
753
754 Ok(())
755 }
756
757 #[cfg_attr(feature = "hotpath", hotpath::measure)]
759 pub fn ingest(&self, event: &Event) -> Result<()> {
760 #[cfg(feature = "server")]
762 let timer = self.metrics.ingestion_duration_seconds.start_timer();
763
764 if let Err(e) = self.ensure_writable() {
766 #[cfg(feature = "server")]
767 {
768 self.metrics.ingestion_errors_total.inc();
769 timer.observe_duration();
770 }
771 return Err(e);
772 }
773
774 let validation_result = self.validate_event(event);
776 if let Err(e) = validation_result {
777 #[cfg(feature = "server")]
778 {
779 self.metrics.ingestion_errors_total.inc();
780 timer.observe_duration();
781 }
782 return Err(e);
783 }
784
785 let _resident = self.cache_residency_gate.read();
786 let _durable = self.durability_gate.read();
787
788 if let Some(ref wal) = self.wal
791 && let Err(e) = wal.append(event.clone())
792 {
793 #[cfg(feature = "server")]
794 {
795 self.metrics.ingestion_errors_total.inc();
796 timer.observe_duration();
797 }
798 return Err(e);
799 }
800
801 *self
803 .entity_versions
804 .entry(event.entity_id_str().to_string())
805 .or_insert(0) += 1;
806
807 let mut events = self.events.write();
808 let offset = events.len();
809
810 self.index.index_event(
812 event.id,
813 event.entity_id_str(),
814 event.event_type_str(),
815 event.timestamp,
816 offset,
817 )?;
818
819 let projections = self.projections.read();
821 projections.process_event(event)?;
822 drop(projections); let pipeline_results = self.pipeline_manager.process_event(event);
827 if !pipeline_results.is_empty() {
828 tracing::debug!(
829 "Event {} processed by {} pipeline(s)",
830 event.id,
831 pipeline_results.len()
832 );
833 for (pipeline_id, result) in pipeline_results {
836 tracing::trace!("Pipeline {} result: {:?}", pipeline_id, result);
837 }
838 }
839
840 if let Some(ref storage) = self.storage {
842 let storage = storage.read();
843 storage.append_event(event.clone())?;
844 }
845
846 events.push(event.clone());
848 let total_events = events.len();
849 drop(events); let event_arc = Arc::new(event.clone());
853 let _ = self.event_broadcast_tx.send(Arc::clone(&event_arc));
854 #[cfg(feature = "server")]
855 self.websocket_manager.broadcast_event(event_arc);
856
857 #[cfg(feature = "server")]
859 self.dispatch_webhooks(event);
860
861 self.geo_index.index_event(event);
863
864 self.schema_evolution
866 .analyze_event(event.event_type_str(), &event.payload);
867
868 self.check_auto_snapshot(event.entity_id_str(), event);
870
871 #[cfg(feature = "server")]
873 {
874 self.metrics.events_ingested_total.inc();
875 self.metrics
876 .events_ingested_by_type
877 .with_label_values(&[event.event_type_str()])
878 .inc();
879 self.metrics.storage_events_total.set(total_events as i64);
880 }
881
882 let mut total = self.total_ingested.write();
884 *total += 1;
885
886 #[cfg(feature = "server")]
887 timer.observe_duration();
888
889 tracing::debug!("Event ingested: {} (offset: {})", event.id, offset);
890
891 Ok(())
892 }
893
894 #[cfg_attr(feature = "hotpath", hotpath::measure)]
901 pub fn ingest_batch(&self, batch: Vec<Event>) -> Result<()> {
902 if batch.is_empty() {
903 return Ok(());
904 }
905
906 self.ensure_writable()?;
908
909 for event in &batch {
911 self.validate_event(event)?;
912 }
913
914 let _resident = self.cache_residency_gate.read();
915 let _durable = self.durability_gate.read();
916 let batch_count = batch.len();
917
918 if let Some(ref wal) = self.wal {
920 for event in &batch {
921 wal.append(event.clone())?;
922 }
923 }
924
925 let mut events = self.events.write();
927 let projections = self.projections.read();
928
929 for event in batch {
930 let offset = events.len();
931
932 self.index.index_event(
933 event.id,
934 event.entity_id_str(),
935 event.event_type_str(),
936 event.timestamp,
937 offset,
938 )?;
939
940 projections.process_event(&event)?;
941 self.pipeline_manager.process_event(&event);
942
943 if let Some(ref storage) = self.storage {
944 let storage = storage.read();
945 storage.append_event(event.clone())?;
946 }
947
948 self.geo_index.index_event(&event);
949 self.schema_evolution
950 .analyze_event(event.event_type_str(), &event.payload);
951
952 *self
954 .entity_versions
955 .entry(event.entity_id_str().to_string())
956 .or_insert(0) += 1;
957
958 let _ = self.event_broadcast_tx.send(Arc::new(event.clone()));
960
961 events.push(event);
962 }
963
964 drop(projections);
965 drop(events);
966
967 let mut total = self.total_ingested.write();
968 *total += batch_count as u64;
969
970 Ok(())
971 }
972
973 #[cfg_attr(feature = "hotpath", hotpath::measure)]
980 pub fn ingest_replicated(&self, event: &Event) -> Result<()> {
981 #[cfg(feature = "server")]
982 let timer = self.metrics.ingestion_duration_seconds.start_timer();
983
984 let _resident = self.cache_residency_gate.read();
985 let mut events = self.events.write();
986 let offset = events.len();
987
988 self.index.index_event(
990 event.id,
991 event.entity_id_str(),
992 event.event_type_str(),
993 event.timestamp,
994 offset,
995 )?;
996
997 let projections = self.projections.read();
999 projections.process_event(event)?;
1000 drop(projections);
1001
1002 let pipeline_results = self.pipeline_manager.process_event(event);
1004 if !pipeline_results.is_empty() {
1005 tracing::debug!(
1006 "Replicated event {} processed by {} pipeline(s)",
1007 event.id,
1008 pipeline_results.len()
1009 );
1010 }
1011
1012 *self
1014 .entity_versions
1015 .entry(event.entity_id_str().to_string())
1016 .or_insert(0) += 1;
1017
1018 if !self.read_only
1019 && let Some(storage) = &self.storage
1020 {
1021 storage.read().append_event(event.clone())?;
1022 }
1023
1024 events.push(event.clone());
1026 let total_events = events.len();
1027 drop(events);
1028
1029 let event_arc = Arc::new(event.clone());
1031 let _ = self.event_broadcast_tx.send(Arc::clone(&event_arc));
1032 #[cfg(feature = "server")]
1033 self.websocket_manager.broadcast_event(event_arc);
1034
1035 #[cfg(feature = "server")]
1037 {
1038 self.metrics.events_ingested_total.inc();
1039 self.metrics
1040 .events_ingested_by_type
1041 .with_label_values(&[event.event_type_str()])
1042 .inc();
1043 self.metrics.storage_events_total.set(total_events as i64);
1044 }
1045
1046 let mut total = self.total_ingested.write();
1047 *total += 1;
1048
1049 #[cfg(feature = "server")]
1050 timer.observe_duration();
1051
1052 tracing::debug!(
1053 "Replicated event ingested: {} (offset: {})",
1054 event.id,
1055 offset
1056 );
1057
1058 Ok(())
1059 }
1060
1061 #[cfg_attr(feature = "hotpath", hotpath::measure)]
1064 pub fn get_entity_version(&self, entity_id: &str) -> u64 {
1065 self.entity_versions.get(entity_id).map_or(0, |v| *v)
1066 }
1067
1068 pub fn consumer_registry(&self) -> &ConsumerRegistry {
1070 &self.consumer_registry
1071 }
1072
1073 pub fn subscribe_events(&self) -> tokio::sync::broadcast::Receiver<Arc<Event>> {
1079 self.event_broadcast_tx.subscribe()
1080 }
1081
1082 pub fn set_consumer_registry(&mut self, registry: Arc<ConsumerRegistry>) {
1087 self.consumer_registry = registry;
1088 }
1089
1090 pub fn total_events(&self) -> usize {
1092 self.events.read().len()
1093 }
1094
1095 pub fn events_after_offset(
1098 &self,
1099 offset: u64,
1100 filters: &[String],
1101 limit: usize,
1102 ) -> Vec<(u64, Event)> {
1103 let events = self.events.read();
1104 let start = offset as usize;
1105 if start >= events.len() {
1106 return vec![];
1107 }
1108
1109 events[start..]
1110 .iter()
1111 .enumerate()
1112 .filter(|(_, event)| ConsumerRegistry::matches_filters(event.event_type_str(), filters))
1113 .take(limit)
1114 .map(|(i, event)| ((start + i + 1) as u64, event.clone()))
1115 .collect()
1116 }
1117
1118 #[cfg(feature = "server")]
1120 pub fn websocket_manager(&self) -> Arc<WebSocketManager> {
1121 Arc::clone(&self.websocket_manager)
1122 }
1123
1124 pub fn snapshot_manager(&self) -> Arc<SnapshotManager> {
1126 Arc::clone(&self.snapshot_manager)
1127 }
1128
1129 pub fn compaction_manager(&self) -> Option<Arc<CompactionManager>> {
1131 self.compaction_manager.as_ref().map(Arc::clone)
1132 }
1133
1134 pub fn schema_registry(&self) -> Arc<SchemaRegistry> {
1136 Arc::clone(&self.schema_registry)
1137 }
1138
1139 pub fn replay_manager(&self) -> Arc<ReplayManager> {
1141 Arc::clone(&self.replay_manager)
1142 }
1143
1144 pub fn pipeline_manager(&self) -> Arc<PipelineManager> {
1146 Arc::clone(&self.pipeline_manager)
1147 }
1148
1149 #[cfg(feature = "server")]
1151 pub fn metrics(&self) -> Arc<MetricsRegistry> {
1152 Arc::clone(&self.metrics)
1153 }
1154
1155 pub fn projection_manager(&self) -> parking_lot::RwLockReadGuard<'_, ProjectionManager> {
1157 self.projections.read()
1158 }
1159
1160 pub fn register_projection(
1169 &self,
1170 projection: Arc<dyn crate::application::services::projection::Projection>,
1171 ) {
1172 let mut pm = self.projections.write();
1173 pm.register(projection);
1174 }
1175
1176 pub fn register_projection_with_backfill(
1188 &self,
1189 projection: &Arc<dyn crate::application::services::projection::Projection>,
1190 ) -> Result<()> {
1191 {
1193 let mut pm = self.projections.write();
1194 pm.register(Arc::clone(projection));
1195 }
1196
1197 let events = self.events.read();
1199 let mut ordered: Vec<&Event> = events.iter().collect();
1200 ordered.sort_by(|a, b| {
1201 a.timestamp
1202 .cmp(&b.timestamp)
1203 .then_with(|| a.version.cmp(&b.version))
1204 });
1205 for event in ordered {
1206 projection.process(event)?;
1207 }
1208
1209 Ok(())
1210 }
1211
1212 pub fn hydrate_all_from_storage(&self) -> Result<usize> {
1233 let Some(storage) = self.storage.as_ref().map(Arc::clone) else {
1234 return Ok(0);
1235 };
1236
1237 let _resident = self.cache_residency_gate.read();
1238 let files_before_load = storage.read().list_parquet_files()?;
1241 let events = storage.read().load_all_events()?;
1242 self.refresh_state
1243 .lock()
1244 .mark_parquet_seen(files_before_load);
1245 let read_count = events.len();
1246 let tenants: Vec<String> = events
1247 .iter()
1248 .map(Event::tenant_id_str)
1249 .collect::<std::collections::HashSet<_>>()
1250 .into_iter()
1251 .map(str::to_owned)
1252 .collect();
1253 let mut applied = 0;
1254 for event in events {
1255 applied += usize::from(self.append_loaded_event(event));
1256 }
1257 for tenant in &tenants {
1261 self.tenant_loader.mark_loaded(tenant);
1262 }
1263
1264 tracing::info!(
1265 read = read_count,
1266 applied = applied,
1267 "🔄 hydrate_all_from_storage: in-memory pile reconstructed from Parquet"
1268 );
1269 Ok(applied)
1270 }
1271
1272 pub fn projection_state_cache(&self) -> Arc<DashMap<String, serde_json::Value>> {
1275 Arc::clone(&self.projection_state_cache)
1276 }
1277
1278 pub fn projection_status(&self) -> Arc<DashMap<String, String>> {
1280 Arc::clone(&self.projection_status)
1281 }
1282
1283 pub fn geo_index(&self) -> Arc<GeoIndex> {
1286 self.geo_index.clone()
1287 }
1288
1289 pub fn exactly_once(&self) -> Arc<ExactlyOnceRegistry> {
1291 self.exactly_once.clone()
1292 }
1293
1294 pub fn schema_evolution(&self) -> Arc<SchemaEvolutionManager> {
1296 self.schema_evolution.clone()
1297 }
1298
1299 pub fn snapshot_events(&self) -> Vec<Event> {
1305 self.events.read().clone()
1306 }
1307
1308 pub fn compact_entity_tokens(
1330 &self,
1331 entity_id: &str,
1332 token_event_type: &str,
1333 merged_event: Event,
1334 ) -> Result<bool> {
1335 self.ensure_writable()?;
1337
1338 {
1340 let events = self.events.read();
1341 let has_tokens = events
1342 .iter()
1343 .any(|e| e.entity_id_str() == entity_id && e.event_type_str() == token_event_type);
1344 if !has_tokens {
1345 return Ok(false);
1346 }
1347 }
1348
1349 let projections = self.projections.read();
1351 projections.process_event(&merged_event)?;
1352 drop(projections);
1353
1354 let mut events = self.events.write();
1356
1357 events.retain(|e| {
1358 !(e.entity_id_str() == entity_id && e.event_type_str() == token_event_type)
1359 });
1360
1361 events.push(merged_event);
1362
1363 self.index.clear();
1368 for (offset, event) in events.iter().enumerate() {
1369 if let Err(e) = self.index.index_event(
1370 event.id,
1371 event.entity_id_str(),
1372 event.event_type_str(),
1373 event.timestamp,
1374 offset,
1375 ) {
1376 tracing::warn!(
1377 event_id = %event.id,
1378 offset,
1379 "Failed to re-index event during compaction: {e}"
1380 );
1381 }
1382 }
1383
1384 Ok(true)
1385 }
1386
1387 #[cfg(feature = "server")]
1388 pub fn webhook_registry(&self) -> Arc<WebhookRegistry> {
1389 Arc::clone(&self.webhook_registry)
1390 }
1391
1392 #[cfg(feature = "server")]
1395 pub fn set_webhook_tx(&self, tx: mpsc::UnboundedSender<WebhookDeliveryTask>) {
1396 *self.webhook_tx.write() = Some(tx);
1397 tracing::info!("Webhook delivery channel connected");
1398 }
1399
1400 #[cfg(feature = "server")]
1402 fn dispatch_webhooks(&self, event: &Event) {
1403 let matching = self.webhook_registry.find_matching(event);
1404 if matching.is_empty() {
1405 return;
1406 }
1407
1408 let tx_guard = self.webhook_tx.read();
1409 if let Some(ref tx) = *tx_guard {
1410 for webhook in matching {
1411 let task = WebhookDeliveryTask {
1412 webhook,
1413 event: event.clone(),
1414 };
1415 if let Err(e) = tx.send(task) {
1416 tracing::warn!("Failed to queue webhook delivery: {}", e);
1417 }
1418 }
1419 }
1420 }
1421
1422 pub fn flush_storage(&self) -> Result<()> {
1424 if let Some(ref storage) = self.storage {
1425 let storage = storage.read();
1426 storage.flush()?;
1427 tracing::info!("✅ Flushed events to persistent storage");
1428 }
1429 Ok(())
1430 }
1431
1432 pub fn checkpoint(&self) -> Result<()> {
1453 let Some(ref wal) = self.wal else {
1454 #[cfg(feature = "server")]
1457 self.refresh_storage_metrics();
1458 return Ok(());
1459 };
1460
1461 if self.read_only {
1462 return Ok(());
1463 }
1464
1465 let active = {
1469 let _sealing = self.durability_gate.write();
1470 wal.seal()?
1471 };
1472 self.flush_storage()?;
1473 wal.remove_sealed(&active)?;
1474 tracing::debug!("✅ Checkpoint complete: Parquet flushed, sealed WAL segments retired");
1475
1476 #[cfg(feature = "server")]
1479 self.refresh_storage_metrics();
1480
1481 Ok(())
1482 }
1483
1484 #[cfg(feature = "server")]
1489 pub fn refresh_storage_metrics_now(&self) {
1490 self.refresh_storage_metrics();
1491 }
1492
1493 #[cfg(feature = "server")]
1510 fn refresh_storage_metrics(&self) {
1511 let Some(ref storage) = self.storage else {
1512 return;
1513 };
1514
1515 let parquet_stats = match storage.read().stats() {
1516 Ok(stats) => stats,
1517 Err(e) => {
1518 tracing::warn!("storage-size metric refresh: failed to stat Parquet: {e}");
1519 return;
1520 }
1521 };
1522
1523 let (wal_bytes, wal_segments) = match self.wal.as_ref() {
1524 Some(wal) => match wal.on_disk_stats() {
1525 Ok(stats) => stats,
1526 Err(e) => {
1527 tracing::warn!("storage-size metric refresh: failed to stat WAL: {e}");
1528 (0, 0)
1529 }
1530 },
1531 None => (0, 0),
1532 };
1533
1534 let total_bytes = parquet_stats.total_size_bytes + wal_bytes;
1535
1536 self.metrics
1537 .storage_size_bytes
1538 .set(total_bytes.min(i64::MAX as u64) as i64);
1539 self.metrics
1540 .parquet_files_total
1541 .set(parquet_stats.total_files as i64);
1542 self.metrics.wal_segments_total.set(wal_segments as i64);
1543
1544 tracing::debug!(
1545 "storage-size metrics refreshed: {} bytes total ({} Parquet files, {} WAL segments)",
1546 total_bytes,
1547 parquet_stats.total_files,
1548 wal_segments
1549 );
1550 }
1551
1552 pub fn checkpoint_interval(&self) -> Option<std::time::Duration> {
1554 self.checkpoint_interval_secs
1555 .map(std::time::Duration::from_secs)
1556 }
1557
1558 pub fn ensure_tenant_loaded(&self, tenant_id: &str) -> Result<()> {
1584 self.ensure_tenant_loaded_with_integrity(tenant_id, false)
1585 }
1586
1587 fn ensure_tenant_loaded_with_integrity(
1588 &self,
1589 tenant_id: &str,
1590 require_complete: bool,
1591 ) -> Result<()> {
1592 self.ensure_tenant_loaded_budgeted(tenant_id, require_complete, None)
1593 }
1594
1595 fn ensure_tenant_loaded_budgeted(
1596 &self,
1597 tenant_id: &str,
1598 require_complete: bool,
1599 cancellation: Option<Arc<std::sync::atomic::AtomicBool>>,
1600 ) -> Result<()> {
1601 self.ensure_tenant_loaded_with_limits(
1602 tenant_id,
1603 require_complete,
1604 cancellation,
1605 &self.strict_archive_limits,
1606 )
1607 }
1608
1609 #[cfg(feature = "server")]
1610 pub(crate) fn prepare_http_append(
1611 &self,
1612 event: &Event,
1613 cancellation: &Arc<std::sync::atomic::AtomicBool>,
1614 ) -> Result<()> {
1615 self.ensure_writable()?;
1616 self.validate_event(event)?;
1617 self.prepare_http_archive(event.tenant_id_str(), cancellation)
1618 }
1619
1620 #[cfg(feature = "server")]
1623 pub(crate) fn prepare_http_archive(
1624 &self,
1625 tenant_id: &str,
1626 cancellation: &Arc<std::sync::atomic::AtomicBool>,
1627 ) -> Result<()> {
1628 crate::domain::value_objects::TenantId::new(tenant_id.to_string())?;
1629 check_cancelled(Some(cancellation.as_ref()))?;
1630 if let Some(timeout) = self.http_archive_warmup_timeout {
1631 let limits = ArchiveReadLimits {
1632 timeout,
1633 ..self.strict_archive_limits.clone()
1634 };
1635 self.ensure_tenant_loaded_with_limits(tenant_id, true, None, &limits)?;
1636 }
1637 check_cancelled(Some(cancellation.as_ref()))
1638 }
1639
1640 fn ensure_tenant_loaded_with_limits(
1641 &self,
1642 tenant_id: &str,
1643 require_complete: bool,
1644 cancellation: Option<Arc<std::sync::atomic::AtomicBool>>,
1645 limits: &ArchiveReadLimits,
1646 ) -> Result<()> {
1647 let is_loaded = || {
1648 self.tenant_loader.is_loaded(tenant_id)
1649 && (!require_complete || self.tenant_loader.is_complete(tenant_id))
1650 };
1651 if is_loaded() {
1653 return Ok(());
1654 }
1655
1656 let mut budget = require_complete
1657 .then(|| ArchiveReadBudget::new(limits.clone()).with_cancellation(cancellation));
1658
1659 let Some(storage) = self.storage.as_ref().map(Arc::clone) else {
1660 self.tenant_loader
1663 .mark_loaded_with_integrity(tenant_id, true);
1664 return Ok(());
1665 };
1666
1667 let lock = self.tenant_loader.lock_for(tenant_id);
1670 let timeout = if let Some(budget) = &budget {
1671 self.tenant_loader.load_timeout().min(budget.remaining()?)
1672 } else {
1673 self.tenant_loader.load_timeout()
1674 };
1675 let _guard = lock.try_lock_for(timeout).ok_or_else(|| {
1676 AllSourceError::StorageError(format!(
1677 "ensure_tenant_loaded timed out after {timeout:?} waiting for in-flight load of \
1678 tenant {tenant_id:?}"
1679 ))
1680 })?;
1681
1682 if is_loaded() {
1685 return Ok(());
1686 }
1687
1688 let generation = {
1689 let _resident = if let Some(budget) = &budget {
1690 self.cache_residency_gate
1691 .try_read_for(budget.remaining()?)
1692 .ok_or_else(|| {
1693 AllSourceError::StorageError(
1694 "Strict archive residency lock timed out".into(),
1695 )
1696 })?
1697 } else {
1698 self.cache_residency_gate.read()
1699 };
1700 self.cache_generations
1701 .get(tenant_id)
1702 .map_or(0, |value| *value)
1703 };
1704 let started = std::time::Instant::now();
1705 let (events, complete) = {
1706 let storage = if let Some(budget) = &budget {
1707 storage.try_read_for(budget.remaining()?).ok_or_else(|| {
1708 AllSourceError::StorageError("Strict archive storage lock timed out".into())
1709 })?
1710 } else {
1711 storage.read()
1712 };
1713 storage.load_events_for_tenant_with_integrity(tenant_id, budget.as_mut())?
1714 };
1715 let read_count = events.len();
1716
1717 let resident = if let Some(budget) = &budget {
1718 self.cache_residency_gate
1719 .try_read_for(budget.remaining()?)
1720 .ok_or_else(|| {
1721 AllSourceError::StorageError("Strict archive residency lock timed out".into())
1722 })?
1723 } else {
1724 self.cache_residency_gate.read()
1725 };
1726 if self
1727 .cache_generations
1728 .get(tenant_id)
1729 .map_or(0, |value| *value)
1730 != generation
1731 {
1732 return Err(AllSourceError::StorageError(
1733 "Archive cache changed while loading retained history".into(),
1734 ));
1735 }
1736 let mut applied = 0;
1737 let load_result = (|| -> Result<()> {
1738 for event in events {
1739 if let Some(budget) = &budget {
1740 budget.check()?;
1741 }
1742 applied += usize::from(self.append_loaded_event(event));
1743 }
1744 if let Some(budget) = &budget {
1745 budget.check()?;
1746 }
1747 Ok(())
1748 })();
1749 load_result?;
1750 self.tenant_loader
1751 .mark_loaded_with_integrity(tenant_id, complete);
1752 drop(resident);
1753
1754 tracing::info!(
1755 tenant_id = tenant_id,
1756 read = read_count,
1757 applied = applied,
1758 elapsed_ms = started.elapsed().as_millis() as u64,
1759 "ensure_tenant_loaded: tenant hydrated"
1760 );
1761
1762 self.enforce_cache_budget(tenant_id);
1770
1771 #[cfg(feature = "server")]
1774 self.metrics
1775 .cache_bytes
1776 .set(self.tenant_loader.total_bytes() as i64);
1777
1778 Ok(())
1779 }
1780
1781 fn enforce_cache_budget(&self, recently_touched: &str) {
1792 if !self.tenant_loader.over_budget() {
1793 return;
1794 }
1795 loop {
1796 let Some(victim) = self.tenant_loader.pick_lru_excluding(recently_touched) else {
1797 tracing::warn!(
1798 cache_bytes = self.tenant_loader.total_bytes(),
1799 budget = self.tenant_loader.byte_budget(),
1800 recently_touched = recently_touched,
1801 "cache over budget but no other tenant available to evict — \
1802 a single tenant exceeds the budget; consider raising it"
1803 );
1804 return;
1805 };
1806 if !self.try_evict_tenant(&victim) {
1807 tracing::warn!(
1808 tenant_id = victim,
1809 "cache eviction refused; retaining resident history"
1810 );
1811 return;
1812 }
1813 if !self.tenant_loader.over_budget() {
1814 return;
1815 }
1816 }
1817 }
1818
1819 pub fn is_tenant_loaded(&self, tenant_id: &str) -> bool {
1823 self.tenant_loader.is_loaded(tenant_id)
1824 }
1825
1826 pub fn evict_tenant(&self, tenant_id: &str) {
1853 self.try_evict_tenant(tenant_id);
1854 }
1855
1856 fn try_evict_tenant(&self, tenant_id: &str) -> bool {
1857 let _resident = self.cache_residency_gate.write();
1858 let Some(storage) = &self.storage else {
1859 return false;
1861 };
1862 if self.read_only {
1863 return false;
1864 }
1865 let Some(storage) = storage.try_write() else {
1868 return false;
1869 };
1870 if storage.has_pending_tenant_events(tenant_id) {
1871 return false;
1872 }
1873 {
1874 let mut generation = self
1875 .cache_generations
1876 .entry(tenant_id.to_string())
1877 .or_insert(0);
1878 let Some(next) = generation.checked_add(1) else {
1879 return false;
1880 };
1881 *generation = next;
1882 }
1883 let mut events = self.events.write();
1884 let before = events.len();
1885 let evicted_bytes = self.tenant_loader.bytes_for(tenant_id);
1886
1887 events.retain(|e| e.tenant_id_str() != tenant_id);
1888 let after = events.len();
1889 let dropped = before - after;
1890
1891 if dropped == 0 {
1892 self.tenant_loader.mark_unloaded(tenant_id);
1896 return true;
1897 }
1898
1899 self.index.clear();
1903 self.entity_versions.clear();
1904 for (offset, event) in events.iter().enumerate() {
1905 if let Err(e) = self.index.index_event(
1906 event.id,
1907 event.entity_id_str(),
1908 event.event_type_str(),
1909 event.timestamp,
1910 offset,
1911 ) {
1912 tracing::error!(
1913 "Failed to re-index event during eviction of {}: {}",
1914 tenant_id,
1915 e
1916 );
1917 }
1918 *self
1919 .entity_versions
1920 .entry(event.entity_id_str().to_string())
1921 .or_insert(0) += 1;
1922 }
1923 self.tenant_loader.mark_unloaded(tenant_id);
1924
1925 let mut t = self.total_ingested.write();
1928 *t = t.saturating_sub(dropped as u64);
1929 drop(t);
1930 drop(events);
1931
1932 #[cfg(feature = "server")]
1935 {
1936 self.metrics.cache_evictions_total.inc();
1937 self.metrics
1938 .cache_bytes
1939 .set(self.tenant_loader.total_bytes() as i64);
1940 }
1941
1942 tracing::info!(
1943 tenant_id = tenant_id,
1944 events_dropped = dropped,
1945 bytes_freed = evicted_bytes,
1946 "evicted tenant from memory cache"
1947 );
1948 true
1949 }
1950
1951 pub fn tenant_resident_bytes(&self, tenant_id: &str) -> u64 {
1955 self.tenant_loader.bytes_for(tenant_id)
1956 }
1957
1958 pub fn cache_resident_bytes(&self) -> u64 {
1961 self.tenant_loader.total_bytes()
1962 }
1963
1964 fn append_loaded_event(&self, event: Event) -> bool {
1986 let mut events = self.events.write();
1987 if self.index.get_by_id(&event.id).is_some() {
1988 return false;
1989 }
1990
1991 let event_bytes = event.estimated_size_bytes();
1992 let tenant = event.tenant_id_str().to_string();
1993
1994 let offset = events.len();
1995
1996 if let Err(e) = self.index.index_event(
1997 event.id,
1998 event.entity_id_str(),
1999 event.event_type_str(),
2000 event.timestamp,
2001 offset,
2002 ) {
2003 tracing::error!("Failed to index loaded event {}: {}", event.id, e);
2004 }
2005
2006 if let Err(e) = self.projections.read().process_event(&event) {
2007 tracing::error!("Failed to project loaded event {}: {}", event.id, e);
2008 }
2009
2010 *self
2011 .entity_versions
2012 .entry(event.entity_id_str().to_string())
2013 .or_insert(0) += 1;
2014
2015 events.push(event);
2016 self.tenant_loader.add_bytes(&tenant, event_bytes);
2020 *self.total_ingested.write() += 1;
2021 true
2022 }
2023
2024 pub fn create_snapshot(&self, entity_id: &str) -> Result<()> {
2026 let events = self.query(&QueryEventsRequest {
2028 entity_id: Some(entity_id.to_string()),
2029 event_type: None,
2030 tenant_id: None,
2031 as_of: None,
2032 since: None,
2033 until: None,
2034 limit: None,
2035 event_type_prefix: None,
2036 exclude_event_type_prefix: None,
2037 payload_filter: None,
2038 })?;
2039
2040 if events.is_empty() {
2041 return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
2042 }
2043
2044 let mut state = serde_json::json!({});
2046 for event in &events {
2047 if let serde_json::Value::Object(ref mut state_map) = state
2048 && let serde_json::Value::Object(ref payload_map) = event.payload
2049 {
2050 for (key, value) in payload_map {
2051 state_map.insert(key.clone(), value.clone());
2052 }
2053 }
2054 }
2055
2056 let last_event = events.last().unwrap();
2057 self.snapshot_manager.create_snapshot(
2058 entity_id,
2059 state,
2060 last_event.timestamp,
2061 events.len(),
2062 SnapshotType::Manual,
2063 )?;
2064
2065 Ok(())
2066 }
2067
2068 fn check_auto_snapshot(&self, entity_id: &str, event: &Event) {
2070 let entity_event_count = self
2072 .index
2073 .get_by_entity(entity_id)
2074 .map_or(0, |entries| entries.len());
2075
2076 if self.snapshot_manager.should_create_snapshot(
2077 entity_id,
2078 entity_event_count,
2079 event.timestamp,
2080 ) {
2081 if let Err(e) = self.create_snapshot(entity_id) {
2083 tracing::warn!(
2084 "Failed to create automatic snapshot for {}: {}",
2085 entity_id,
2086 e
2087 );
2088 }
2089 }
2090 }
2091
2092 fn validate_event(&self, event: &Event) -> Result<()> {
2094 if event.entity_id_str().is_empty() {
2097 return Err(AllSourceError::ValidationError(
2098 "entity_id cannot be empty".to_string(),
2099 ));
2100 }
2101
2102 if event.event_type_str().is_empty() {
2103 return Err(AllSourceError::ValidationError(
2104 "event_type cannot be empty".to_string(),
2105 ));
2106 }
2107
2108 if event.event_type().is_system() {
2111 return Err(AllSourceError::ValidationError(
2112 "Event types starting with '_system.' are reserved for internal use".to_string(),
2113 ));
2114 }
2115
2116 Ok(())
2117 }
2118
2119 pub fn reset_projection(&self, name: &str) -> Result<usize> {
2121 let projection_manager = self.projections.read();
2122 let projection = projection_manager.get_projection(name).ok_or_else(|| {
2123 AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
2124 })?;
2125
2126 projection.clear();
2128
2129 let prefix = format!("{name}:");
2131 let keys_to_remove: Vec<String> = self
2132 .projection_state_cache
2133 .iter()
2134 .filter(|entry| entry.key().starts_with(&prefix))
2135 .map(|entry| entry.key().clone())
2136 .collect();
2137 for key in keys_to_remove {
2138 self.projection_state_cache.remove(&key);
2139 }
2140
2141 let events = self.events.read();
2143 let mut reprocessed = 0usize;
2144 for event in events.iter() {
2145 if projection.process(event).is_ok() {
2146 reprocessed += 1;
2147 }
2148 }
2149
2150 Ok(reprocessed)
2151 }
2152
2153 pub fn get_event_by_id(&self, event_id: &uuid::Uuid) -> Result<Option<Event>> {
2155 if let Some(offset) = self.index.get_by_id(event_id) {
2156 let events = self.events.read();
2157 Ok(events.get(offset).cloned())
2158 } else {
2159 Ok(None)
2160 }
2161 }
2162
2163 #[cfg_attr(feature = "hotpath", hotpath::measure)]
2167 pub fn query(&self, request: &QueryEventsRequest) -> Result<Vec<Event>> {
2168 self.query_window(request, 0, false)
2169 .map(|(events, _)| events)
2170 }
2171
2172 pub fn query_scoped(
2174 &self,
2175 request: &QueryEventsRequest,
2176 scope: &ReadScope,
2177 ) -> Result<Vec<Event>> {
2178 self.query_window_scoped(request, 0, false, scope)
2179 .map(|(events, _)| events)
2180 }
2181
2182 #[cfg_attr(feature = "hotpath", hotpath::measure)]
2199 pub fn query_window(
2200 &self,
2201 request: &QueryEventsRequest,
2202 offset: usize,
2203 descending: bool,
2204 ) -> Result<(Vec<Event>, usize)> {
2205 self.query_window_scoped(request, offset, descending, &ReadScope::unrestricted())
2206 }
2207
2208 pub fn query_window_scoped(
2218 &self,
2219 request: &QueryEventsRequest,
2220 offset: usize,
2221 descending: bool,
2222 scope: &ReadScope,
2223 ) -> Result<(Vec<Event>, usize)> {
2224 if let Some(filter) = &request.payload_filter
2231 && serde_json::from_str::<serde_json::Map<String, serde_json::Value>>(filter).is_err()
2232 {
2233 return Err(AllSourceError::InvalidInput(format!(
2234 "invalid 'payload_filter': expected a JSON object of field/value \
2235 pairs, got '{filter}'"
2236 )));
2237 }
2238
2239 if let Some(ref tenant_id) = request.tenant_id {
2258 self.ensure_tenant_loaded(tenant_id)?;
2259 self.tenant_loader.touch(tenant_id);
2263 }
2264
2265 let query_type = if request.entity_id.is_some() {
2267 "entity"
2268 } else if request.event_type.is_some() {
2269 "type"
2270 } else if request.event_type_prefix.is_some() {
2271 "type_prefix"
2272 } else {
2273 "full_scan"
2274 };
2275
2276 #[cfg(feature = "server")]
2278 let timer = self
2279 .metrics
2280 .query_duration_seconds
2281 .with_label_values(&[query_type])
2282 .start_timer();
2283
2284 #[cfg(feature = "server")]
2286 self.metrics
2287 .queries_total
2288 .with_label_values(&[query_type])
2289 .inc();
2290
2291 let events = self.events.read();
2292
2293 let offsets: Vec<usize> = if let Some(entity_id) = &request.entity_id {
2295 self.index
2297 .get_by_entity(entity_id)
2298 .map(|entries| self.filter_entries(entries, request))
2299 .unwrap_or_default()
2300 } else if let Some(event_type) = &request.event_type {
2301 self.index
2303 .get_by_type(event_type)
2304 .map(|entries| self.filter_entries(entries, request))
2305 .unwrap_or_default()
2306 } else if let Some(prefix) = &request.event_type_prefix {
2307 let entries = self.index.get_by_type_prefix(prefix);
2309 self.filter_entries(entries, request)
2310 } else {
2311 (0..events.len()).collect()
2313 };
2314
2315 let mut matches: Vec<(usize, &Event)> = offsets
2321 .iter()
2322 .filter_map(|&event_offset| events.get(event_offset))
2323 .filter(|event| scope.permits(event.entity_id().as_str()))
2324 .filter(|event| self.apply_filters(event, request))
2325 .enumerate()
2326 .collect();
2327
2328 let total = matches.len();
2331
2332 let order = |a: &(usize, &Event), b: &(usize, &Event)| {
2339 let ascending =
2340 a.1.timestamp
2341 .cmp(&b.1.timestamp)
2342 .then_with(|| a.1.version.cmp(&b.1.version))
2343 .then_with(|| a.0.cmp(&b.0));
2344 if descending {
2345 ascending.reverse()
2346 } else {
2347 ascending
2348 }
2349 };
2350
2351 select_window(&mut matches, offset, request.limit, order);
2353
2354 let results: Vec<Event> = matches
2357 .into_iter()
2358 .skip(offset)
2359 .map(|(_, event)| event.clone())
2360 .collect();
2361
2362 #[cfg(feature = "server")]
2364 {
2365 self.metrics
2366 .query_results_total
2367 .with_label_values(&[query_type])
2368 .inc_by(results.len() as u64);
2369 timer.observe_duration();
2370 }
2371
2372 Ok((results, total))
2373 }
2374
2375 #[cfg_attr(feature = "hotpath", hotpath::measure)]
2377 fn filter_entries(&self, entries: Vec<IndexEntry>, request: &QueryEventsRequest) -> Vec<usize> {
2378 entries
2379 .into_iter()
2380 .filter(|entry| {
2381 if let Some(as_of) = request.as_of
2383 && entry.timestamp > as_of
2384 {
2385 return false;
2386 }
2387 if let Some(since) = request.since
2388 && entry.timestamp < since
2389 {
2390 return false;
2391 }
2392 if let Some(until) = request.until
2393 && entry.timestamp > until
2394 {
2395 return false;
2396 }
2397 true
2398 })
2399 .map(|entry| entry.offset)
2400 .collect()
2401 }
2402
2403 #[cfg_attr(feature = "hotpath", hotpath::measure)]
2405 fn apply_filters(&self, event: &Event, request: &QueryEventsRequest) -> bool {
2406 if let Some(ref tid) = request.tenant_id
2408 && event.tenant_id_str() != tid
2409 {
2410 return false;
2411 }
2412
2413 if let Some(as_of) = request.as_of
2423 && event.timestamp > as_of
2424 {
2425 return false;
2426 }
2427 if let Some(since) = request.since
2428 && event.timestamp < since
2429 {
2430 return false;
2431 }
2432 if let Some(until) = request.until
2433 && event.timestamp > until
2434 {
2435 return false;
2436 }
2437
2438 if let Some(ref excludes) = request.exclude_event_type_prefix {
2442 let et = event.event_type_str();
2443 if excludes
2444 .split(',')
2445 .map(str::trim)
2446 .filter(|p| !p.is_empty())
2447 .any(|p| et.starts_with(p))
2448 {
2449 return false;
2450 }
2451 }
2452
2453 if request.entity_id.is_some()
2455 && let Some(ref event_type) = request.event_type
2456 && event.event_type_str() != event_type
2457 {
2458 return false;
2459 }
2460
2461 if request.entity_id.is_some()
2463 && let Some(ref prefix) = request.event_type_prefix
2464 && !event.event_type_str().starts_with(prefix)
2465 {
2466 return false;
2467 }
2468
2469 if let Some(ref filter_str) = request.payload_filter
2471 && let Ok(filter_obj) =
2472 serde_json::from_str::<serde_json::Map<String, serde_json::Value>>(filter_str)
2473 {
2474 let payload = event.payload();
2475 for (key, expected_value) in &filter_obj {
2476 match payload.get(key) {
2477 Some(actual_value) if actual_value == expected_value => {}
2478 _ => return false,
2479 }
2480 }
2481 }
2482
2483 true
2484 }
2485
2486 #[cfg_attr(feature = "hotpath", hotpath::measure)]
2489 pub fn reconstruct_state(
2490 &self,
2491 entity_id: &str,
2492 as_of: Option<DateTime<Utc>>,
2493 ) -> Result<serde_json::Value> {
2494 let (merged_state, since_timestamp) = if let Some(as_of_time) = as_of {
2496 if let Some(snapshot) = self
2498 .snapshot_manager
2499 .get_snapshot_as_of(entity_id, as_of_time)
2500 {
2501 tracing::debug!(
2502 "Using snapshot from {} for entity {} (saved {} events)",
2503 snapshot.as_of,
2504 entity_id,
2505 snapshot.event_count
2506 );
2507 (snapshot.state.clone(), Some(snapshot.as_of))
2508 } else {
2509 (serde_json::json!({}), None)
2510 }
2511 } else {
2512 if let Some(snapshot) = self.snapshot_manager.get_latest_snapshot(entity_id) {
2514 tracing::debug!(
2515 "Using latest snapshot from {} for entity {}",
2516 snapshot.as_of,
2517 entity_id
2518 );
2519 (snapshot.state.clone(), Some(snapshot.as_of))
2520 } else {
2521 (serde_json::json!({}), None)
2522 }
2523 };
2524
2525 let events = self.query(&QueryEventsRequest {
2527 entity_id: Some(entity_id.to_string()),
2528 event_type: None,
2529 tenant_id: None,
2530 as_of,
2531 since: since_timestamp,
2532 until: None,
2533 limit: None,
2534 event_type_prefix: None,
2535 exclude_event_type_prefix: None,
2536 payload_filter: None,
2537 })?;
2538
2539 if events.is_empty() && since_timestamp.is_none() {
2541 return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
2542 }
2543
2544 let mut merged_state = merged_state;
2546 for event in &events {
2547 if let serde_json::Value::Object(ref mut state_map) = merged_state
2548 && let serde_json::Value::Object(ref payload_map) = event.payload
2549 {
2550 for (key, value) in payload_map {
2551 state_map.insert(key.clone(), value.clone());
2552 }
2553 }
2554 }
2555
2556 let state = serde_json::json!({
2558 "entity_id": entity_id,
2559 "last_updated": events.last().map(|e| e.timestamp),
2560 "event_count": events.len(),
2561 "as_of": as_of,
2562 "current_state": merged_state,
2563 "history": events.iter().map(|e| {
2564 serde_json::json!({
2565 "event_id": e.id,
2566 "type": e.event_type,
2567 "timestamp": e.timestamp,
2568 "payload": e.payload
2569 })
2570 }).collect::<Vec<_>>()
2571 });
2572
2573 Ok(state)
2574 }
2575
2576 pub fn get_snapshot(&self, entity_id: &str) -> Result<serde_json::Value> {
2578 let projections = self.projections.read();
2579
2580 if let Some(snapshot_projection) = projections.get_projection("entity_snapshots")
2581 && let Some(state) = snapshot_projection.get_state(entity_id)
2582 {
2583 return Ok(serde_json::json!({
2584 "entity_id": entity_id,
2585 "snapshot": state,
2586 "from_projection": "entity_snapshots"
2587 }));
2588 }
2589
2590 Err(AllSourceError::EntityNotFound(entity_id.to_string()))
2591 }
2592
2593 pub fn stats(&self) -> StoreStats {
2595 let events = self.events.read();
2596 let index_stats = self.index.stats();
2597
2598 StoreStats {
2599 total_events: events.len(),
2600 total_entities: index_stats.total_entities,
2601 total_event_types: index_stats.total_event_types,
2602 total_ingested: *self.total_ingested.read(),
2603 }
2604 }
2605
2606 pub fn list_streams(&self) -> Vec<StreamInfo> {
2608 self.index
2609 .get_all_entities()
2610 .into_iter()
2611 .map(|entity_id| {
2612 let event_count = self
2613 .index
2614 .get_by_entity(&entity_id)
2615 .map_or(0, |entries| entries.len());
2616 let last_event_at = self
2617 .index
2618 .get_by_entity(&entity_id)
2619 .and_then(|entries| entries.last().map(|e| e.timestamp));
2620 StreamInfo {
2621 stream_id: entity_id,
2622 event_count,
2623 last_event_at,
2624 }
2625 })
2626 .collect()
2627 }
2628
2629 pub fn list_event_types(&self) -> Vec<EventTypeInfo> {
2631 self.index
2632 .get_all_types()
2633 .into_iter()
2634 .map(|event_type| {
2635 let event_count = self
2636 .index
2637 .get_by_type(&event_type)
2638 .map_or(0, |entries| entries.len());
2639 let last_event_at = self
2640 .index
2641 .get_by_type(&event_type)
2642 .and_then(|entries| entries.last().map(|e| e.timestamp));
2643 EventTypeInfo {
2644 event_type,
2645 event_count,
2646 last_event_at,
2647 }
2648 })
2649 .collect()
2650 }
2651
2652 pub fn list_streams_for_tenant(&self, tenant_id: &str) -> Vec<StreamInfo> {
2660 let _ = self.ensure_tenant_loaded(tenant_id);
2661 let events = self.events.read();
2662 let mut by_entity: std::collections::HashMap<&str, (usize, chrono::DateTime<chrono::Utc>)> =
2663 std::collections::HashMap::new();
2664 for ev in events.iter() {
2665 if ev.tenant_id_str() != tenant_id {
2666 continue;
2667 }
2668 let e = by_entity
2669 .entry(ev.entity_id_str())
2670 .or_insert((0, ev.timestamp));
2671 e.0 += 1;
2672 if ev.timestamp > e.1 {
2673 e.1 = ev.timestamp;
2674 }
2675 }
2676 by_entity
2677 .into_iter()
2678 .map(|(entity_id, (count, last))| StreamInfo {
2679 stream_id: entity_id.to_string(),
2680 event_count: count,
2681 last_event_at: Some(last),
2682 })
2683 .collect()
2684 }
2685
2686 pub fn list_event_types_for_tenant(&self, tenant_id: &str) -> Vec<EventTypeInfo> {
2688 let _ = self.ensure_tenant_loaded(tenant_id);
2689 let events = self.events.read();
2690 let mut by_type: std::collections::HashMap<&str, (usize, chrono::DateTime<chrono::Utc>)> =
2691 std::collections::HashMap::new();
2692 for ev in events.iter() {
2693 if ev.tenant_id_str() != tenant_id {
2694 continue;
2695 }
2696 let e = by_type
2697 .entry(ev.event_type_str())
2698 .or_insert((0, ev.timestamp));
2699 e.0 += 1;
2700 if ev.timestamp > e.1 {
2701 e.1 = ev.timestamp;
2702 }
2703 }
2704 by_type
2705 .into_iter()
2706 .map(|(event_type, (count, last))| EventTypeInfo {
2707 event_type: event_type.to_string(),
2708 event_count: count,
2709 last_event_at: Some(last),
2710 })
2711 .collect()
2712 }
2713
2714 pub fn stats_for_tenant(&self, tenant_id: &str) -> TenantStoreStats {
2724 let _ = self.ensure_tenant_loaded(tenant_id);
2725 let events = self.events.read();
2726
2727 let mut entities: std::collections::HashSet<&str> = std::collections::HashSet::new();
2728 let mut census: std::collections::HashMap<&str, usize> = std::collections::HashMap::new();
2729 let mut total_events = 0usize;
2730 let mut oldest: Option<chrono::DateTime<chrono::Utc>> = None;
2731 let mut newest: Option<chrono::DateTime<chrono::Utc>> = None;
2732
2733 for ev in events.iter() {
2734 if ev.tenant_id_str() != tenant_id {
2735 continue;
2736 }
2737
2738 total_events += 1;
2739 entities.insert(ev.entity_id_str());
2740 *census.entry(ev.event_type_str()).or_insert(0) += 1;
2741
2742 let ts = ev.timestamp;
2743 if oldest.is_none_or(|o| ts < o) {
2744 oldest = Some(ts);
2745 }
2746 if newest.is_none_or(|n| ts > n) {
2747 newest = Some(ts);
2748 }
2749 }
2750
2751 TenantStoreStats {
2752 total_events,
2753 total_entities: entities.len(),
2754 total_event_types: census.len(),
2755 total_ingested: total_events as u64,
2759 event_types: census
2760 .into_iter()
2761 .map(|(k, v)| (k.to_string(), v))
2762 .collect(),
2763 oldest_event: oldest,
2764 newest_event: newest,
2765 }
2766 }
2767
2768 pub fn reconstruct_state_for_tenant(
2777 &self,
2778 entity_id: &str,
2779 as_of: Option<DateTime<Utc>>,
2780 tenant_id: &str,
2781 ) -> Result<serde_json::Value> {
2782 let events = self.query(&QueryEventsRequest {
2783 entity_id: Some(entity_id.to_string()),
2784 event_type: None,
2785 tenant_id: Some(tenant_id.to_string()),
2786 as_of,
2787 since: None,
2788 until: None,
2789 limit: None,
2790 event_type_prefix: None,
2791 exclude_event_type_prefix: None,
2792 payload_filter: None,
2793 })?;
2794
2795 if events.is_empty() {
2796 return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
2797 }
2798
2799 let mut merged_state = serde_json::json!({});
2800 for event in &events {
2801 if let serde_json::Value::Object(ref mut state_map) = merged_state
2802 && let serde_json::Value::Object(ref payload_map) = event.payload
2803 {
2804 for (key, value) in payload_map {
2805 state_map.insert(key.clone(), value.clone());
2806 }
2807 }
2808 }
2809
2810 Ok(serde_json::json!({
2811 "entity_id": entity_id,
2812 "last_updated": events.last().map(|e| e.timestamp),
2813 "event_count": events.len(),
2814 "as_of": as_of,
2815 "current_state": merged_state,
2816 "history": events.iter().map(|e| {
2817 serde_json::json!({
2818 "event_id": e.id,
2819 "type": e.event_type,
2820 "timestamp": e.timestamp,
2821 "payload": e.payload
2822 })
2823 }).collect::<Vec<_>>()
2824 }))
2825 }
2826
2827 pub fn enable_wal_replication(
2834 &self,
2835 tx: tokio::sync::broadcast::Sender<crate::infrastructure::persistence::wal::WALEntry>,
2836 ) {
2837 if let Some(ref wal_arc) = self.wal {
2838 wal_arc.set_replication_tx(tx);
2839 tracing::info!("WAL replication broadcast enabled");
2840 } else {
2841 tracing::warn!("Cannot enable WAL replication: WAL is not configured");
2842 }
2843 }
2844
2845 pub fn wal(&self) -> Option<&Arc<WriteAheadLog>> {
2848 self.wal.as_ref()
2849 }
2850
2851 pub fn parquet_storage(&self) -> Option<&Arc<RwLock<ParquetStorage>>> {
2854 self.storage.as_ref()
2855 }
2856}
2857
2858#[derive(Debug, Clone, Default)]
2860pub struct EventStoreConfig {
2861 pub strict_archive_limits: ArchiveReadLimits,
2864 pub http_archive_warmup_timeout: Option<std::time::Duration>,
2868 pub storage_dir: Option<PathBuf>,
2870
2871 pub snapshot_config: SnapshotConfig,
2873
2874 pub wal_dir: Option<PathBuf>,
2876
2877 pub wal_config: WALConfig,
2879
2880 pub compaction_config: CompactionConfig,
2882
2883 pub schema_registry_config: SchemaRegistryConfig,
2885
2886 pub system_data_dir: Option<PathBuf>,
2891
2892 pub bootstrap_tenant: Option<String>,
2894
2895 pub cache_byte_budget: Option<u64>,
2902
2903 pub checkpoint_interval_secs: Option<u64>,
2915
2916 pub read_only: bool,
2920}
2921
2922impl EventStoreConfig {
2923 pub fn with_persistence(storage_dir: impl Into<PathBuf>) -> Self {
2925 Self {
2926 storage_dir: Some(storage_dir.into()),
2927 ..Self::default()
2928 }
2929 }
2930
2931 pub fn with_snapshots(snapshot_config: SnapshotConfig) -> Self {
2933 Self {
2934 snapshot_config,
2935 ..Self::default()
2936 }
2937 }
2938
2939 pub fn with_wal(wal_dir: impl Into<PathBuf>, wal_config: WALConfig) -> Self {
2941 Self {
2942 wal_dir: Some(wal_dir.into()),
2943 wal_config,
2944 ..Self::default()
2945 }
2946 }
2947
2948 pub fn with_all(storage_dir: impl Into<PathBuf>, snapshot_config: SnapshotConfig) -> Self {
2950 Self {
2951 storage_dir: Some(storage_dir.into()),
2952 snapshot_config,
2953 ..Self::default()
2954 }
2955 }
2956
2957 pub fn production(
2959 storage_dir: impl Into<PathBuf>,
2960 wal_dir: impl Into<PathBuf>,
2961 snapshot_config: SnapshotConfig,
2962 wal_config: WALConfig,
2963 compaction_config: CompactionConfig,
2964 ) -> Self {
2965 let storage_dir = storage_dir.into();
2966 let system_data_dir = storage_dir.join("__system");
2967 Self {
2968 storage_dir: Some(storage_dir),
2969 snapshot_config,
2970 wal_dir: Some(wal_dir.into()),
2971 wal_config,
2972 compaction_config,
2973 system_data_dir: Some(system_data_dir),
2974 ..Self::default()
2975 }
2976 }
2977
2978 pub fn effective_system_data_dir(&self) -> Option<PathBuf> {
2983 self.system_data_dir
2984 .clone()
2985 .or_else(|| self.storage_dir.as_ref().map(|d| d.join("__system")))
2986 }
2987
2988 pub fn from_env() -> (Self, &'static str) {
2996 Self::from_env_vars(
2997 std::env::var("ALLSOURCE_DATA_DIR")
2998 .ok()
2999 .filter(|s| !s.is_empty()),
3000 std::env::var("ALLSOURCE_STORAGE_DIR")
3001 .ok()
3002 .filter(|s| !s.is_empty()),
3003 std::env::var("ALLSOURCE_WAL_DIR")
3004 .ok()
3005 .filter(|s| !s.is_empty()),
3006 std::env::var("ALLSOURCE_WAL_ENABLED").ok(),
3007 std::env::var("ALLSOURCE_CACHE_BYTES").ok(),
3008 std::env::var("ALLSOURCE_SNAPSHOT_INTERVAL_SECONDS").ok(),
3009 std::env::var("ALLSOURCE_RETENTION_SYSTEM_DAYS").ok(),
3010 std::env::var("ALLSOURCE_CHECKPOINT_INTERVAL_SECONDS").ok(),
3011 )
3012 }
3013
3014 pub fn from_env_vars(
3016 data_dir: Option<String>,
3017 explicit_storage_dir: Option<String>,
3018 explicit_wal_dir: Option<String>,
3019 wal_enabled_var: Option<String>,
3020 cache_bytes_var: Option<String>,
3021 snapshot_interval_var: Option<String>,
3022 retention_system_days_var: Option<String>,
3023 checkpoint_interval_var: Option<String>,
3024 ) -> (Self, &'static str) {
3025 let data_dir = data_dir.filter(|s| !s.is_empty());
3026 let storage_dir = explicit_storage_dir
3027 .filter(|s| !s.is_empty())
3028 .or_else(|| data_dir.as_ref().map(|d| format!("{d}/storage")));
3029 let wal_dir = explicit_wal_dir
3030 .filter(|s| !s.is_empty())
3031 .or_else(|| data_dir.as_ref().map(|d| format!("{d}/wal")));
3032 let wal_enabled = wal_enabled_var.is_none_or(|v| v == "true");
3033 let cache_byte_budget =
3038 cache_bytes_var
3039 .filter(|s| !s.is_empty())
3040 .and_then(|s| match s.parse::<u64>() {
3041 Ok(v) => Some(v),
3042 Err(e) => {
3043 tracing::warn!(
3044 "ALLSOURCE_CACHE_BYTES={s:?} could not be parsed as u64: {e}; \
3045 cache budget disabled"
3046 );
3047 None
3048 }
3049 });
3050 let compaction_config =
3051 CompactionConfig::from_env_vars(snapshot_interval_var, retention_system_days_var);
3052
3053 let checkpoint_interval_secs = if wal_enabled {
3058 checkpoint_interval_var
3059 .filter(|s| !s.is_empty())
3060 .map(|s| match s.parse::<u64>() {
3061 Ok(v) => v,
3062 Err(e) => {
3063 tracing::warn!(
3064 "ALLSOURCE_CHECKPOINT_INTERVAL_SECONDS={s:?} could not be parsed as \
3065 u64: {e}; falling back to default 60s"
3066 );
3067 60
3068 }
3069 })
3070 .or(Some(60))
3071 } else {
3072 None
3073 };
3074
3075 let mut config = match (&storage_dir, &wal_dir) {
3076 (Some(sd), Some(wd)) if wal_enabled => Self::production(
3077 sd,
3078 wd,
3079 SnapshotConfig::default(),
3080 WALConfig::default(),
3081 compaction_config,
3082 ),
3083 (Some(sd), _) => Self::with_persistence(sd),
3084 (_, Some(wd)) if wal_enabled => Self::with_wal(wd, WALConfig::default()),
3085 _ => Self::default(),
3086 };
3087 config.cache_byte_budget = cache_byte_budget;
3088 config.checkpoint_interval_secs = checkpoint_interval_secs;
3089 config.http_archive_warmup_timeout = storage_dir
3090 .as_ref()
3091 .map(|_| std::time::Duration::from_secs(30));
3092
3093 let mode = match (&storage_dir, &wal_dir) {
3094 (Some(_), Some(_)) if wal_enabled => "wal+parquet",
3095 (Some(_), _) => "parquet-only",
3096 (_, Some(_)) if wal_enabled => "wal-only",
3097 _ => "in-memory",
3098 };
3099 (config, mode)
3100 }
3101}
3102
3103#[derive(Debug, serde::Serialize)]
3104pub struct StoreStats {
3105 pub total_events: usize,
3106 pub total_entities: usize,
3107 pub total_event_types: usize,
3108 pub total_ingested: u64,
3109}
3110
3111#[derive(Debug, Clone, serde::Serialize)]
3117pub struct TenantStoreStats {
3118 pub total_events: usize,
3119 pub total_entities: usize,
3120 pub total_event_types: usize,
3121 pub total_ingested: u64,
3122 pub event_types: std::collections::HashMap<String, usize>,
3124 pub oldest_event: Option<chrono::DateTime<chrono::Utc>>,
3125 pub newest_event: Option<chrono::DateTime<chrono::Utc>>,
3126}
3127
3128#[derive(Debug, Clone, serde::Serialize)]
3130pub struct StreamInfo {
3131 pub stream_id: String,
3133 pub event_count: usize,
3135 pub last_event_at: Option<chrono::DateTime<chrono::Utc>>,
3137}
3138
3139#[derive(Debug, Clone, serde::Serialize)]
3141pub struct EventTypeInfo {
3142 pub event_type: String,
3144 pub event_count: usize,
3146 pub last_event_at: Option<chrono::DateTime<chrono::Utc>>,
3148}
3149
3150impl Default for EventStore {
3151 fn default() -> Self {
3152 Self::new()
3153 }
3154}
3155
3156#[path = "store_strict_read.rs"]
3157mod strict_read;
3158
3159#[path = "store_refresh.rs"]
3160mod refresh;
3161pub use refresh::RefreshReport;
3162
3163#[cfg(all(test, feature = "server"))]
3164#[path = "store_archive_work_tests.rs"]
3165mod archive_work_tests;
3166
3167#[cfg(test)]
3168#[path = "store_archive_consistency_tests.rs"]
3169mod archive_consistency_tests;
3170
3171#[cfg(test)]
3172mod tests {
3173 use super::*;
3174 use crate::domain::entities::Event;
3175 use tempfile::TempDir;
3176
3177 fn find_parquet_files(dir: &std::path::Path) -> Vec<std::path::PathBuf> {
3182 let mut out = Vec::new();
3183 let mut stack = vec![dir.to_path_buf()];
3184 while let Some(d) = stack.pop() {
3185 let Ok(entries) = std::fs::read_dir(&d) else {
3186 continue;
3187 };
3188 for e in entries.flatten() {
3189 let p = e.path();
3190 if p.is_dir() {
3191 stack.push(p);
3192 } else if p.extension().and_then(|s| s.to_str()) == Some("parquet") {
3193 out.push(p);
3194 }
3195 }
3196 }
3197 out
3198 }
3199
3200 fn create_test_event(entity_id: &str, event_type: &str) -> Event {
3201 Event::from_strings(
3202 event_type.to_string(),
3203 entity_id.to_string(),
3204 "default".to_string(),
3205 serde_json::json!({"name": "Test", "value": 42}),
3206 None,
3207 )
3208 .unwrap()
3209 }
3210
3211 fn create_test_event_with_payload(
3212 entity_id: &str,
3213 event_type: &str,
3214 payload: serde_json::Value,
3215 ) -> Event {
3216 Event::from_strings(
3217 event_type.to_string(),
3218 entity_id.to_string(),
3219 "default".to_string(),
3220 payload,
3221 None,
3222 )
3223 .unwrap()
3224 }
3225
3226 #[test]
3227 fn test_event_store_new() {
3228 let store = EventStore::new();
3229 assert_eq!(store.stats().total_events, 0);
3230 assert_eq!(store.stats().total_entities, 0);
3231 }
3232
3233 #[test]
3240 fn test_ensure_tenant_loaded_no_storage_is_a_noop() {
3241 let store = EventStore::new();
3245 assert!(!store.is_tenant_loaded("alice"));
3246 store.ensure_tenant_loaded("alice").unwrap();
3247 assert!(store.is_tenant_loaded("alice"));
3248 assert!(!store.is_tenant_loaded("bob"));
3250 }
3251
3252 #[test]
3253 fn test_ensure_tenant_loaded_warm_path_is_idempotent() {
3254 let store = EventStore::new();
3255 store.ensure_tenant_loaded("alice").unwrap();
3256 store.ensure_tenant_loaded("alice").unwrap();
3258 }
3259
3260 #[test]
3261 fn test_ensure_tenant_loaded_rejects_unsafe_tenant_id() {
3262 let temp_dir = TempDir::new().unwrap();
3268 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3269 for unsafe_tid in ["..", "a/b", "a\\b", ""] {
3270 let result = store.ensure_tenant_loaded(unsafe_tid);
3271 assert!(
3272 result.is_err(),
3273 "tenant_id {unsafe_tid:?} should have been rejected"
3274 );
3275 assert!(
3276 !store.is_tenant_loaded(unsafe_tid),
3277 "rejected tenant {unsafe_tid:?} must not be marked loaded"
3278 );
3279 }
3280 }
3281
3282 #[test]
3283 fn test_ensure_tenant_loaded_no_subtree_marks_loaded_with_zero_events() {
3284 let temp_dir = TempDir::new().unwrap();
3289 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3290 assert!(!store.is_tenant_loaded("never-existed"));
3291 store.ensure_tenant_loaded("never-existed").unwrap();
3292 assert!(store.is_tenant_loaded("never-existed"));
3293 }
3294
3295 #[test]
3296 fn test_evict_tenant_drops_events_and_resets_bytes() {
3297 let temp_dir = TempDir::new().unwrap();
3301 let storage_dir = temp_dir.path().to_path_buf();
3302
3303 {
3304 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3305 for i in 0..3 {
3306 store
3307 .ingest(
3308 &Event::from_strings(
3309 "test.event".to_string(),
3310 format!("a-{i}"),
3311 "alice".to_string(),
3312 serde_json::json!({"i": i}),
3313 None,
3314 )
3315 .unwrap(),
3316 )
3317 .unwrap();
3318 }
3319 for i in 0..2 {
3320 store
3321 .ingest(
3322 &Event::from_strings(
3323 "test.event".to_string(),
3324 format!("b-{i}"),
3325 "bob".to_string(),
3326 serde_json::json!({"i": i}),
3327 None,
3328 )
3329 .unwrap(),
3330 )
3331 .unwrap();
3332 }
3333 store.flush_storage().unwrap();
3334 }
3335
3336 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3337 store.ensure_tenant_loaded("alice").unwrap();
3338 store.ensure_tenant_loaded("bob").unwrap();
3339 assert_eq!(store.stats().total_events, 5);
3340 let alice_bytes = store.tenant_resident_bytes("alice");
3341 let bob_bytes = store.tenant_resident_bytes("bob");
3342 assert!(alice_bytes > 0 && bob_bytes > 0);
3343
3344 store.evict_tenant("alice");
3345
3346 assert!(!store.is_tenant_loaded("alice"));
3347 assert!(store.is_tenant_loaded("bob"));
3348 assert_eq!(store.tenant_resident_bytes("alice"), 0);
3349 assert_eq!(store.tenant_resident_bytes("bob"), bob_bytes);
3350 assert_eq!(store.stats().total_events, 2, "only bob's 2 events remain");
3351 }
3352
3353 #[test]
3354 fn test_evict_tenant_then_query_re_loads_from_disk() {
3355 let temp_dir = TempDir::new().unwrap();
3359 let storage_dir = temp_dir.path().to_path_buf();
3360
3361 {
3362 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3363 for i in 0..4 {
3364 store
3365 .ingest(
3366 &Event::from_strings(
3367 "test.event".to_string(),
3368 format!("a-{i}"),
3369 "alice".to_string(),
3370 serde_json::json!({"i": i}),
3371 None,
3372 )
3373 .unwrap(),
3374 )
3375 .unwrap();
3376 }
3377 store.flush_storage().unwrap();
3378 }
3379
3380 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3381 store.ensure_tenant_loaded("alice").unwrap();
3382 store.evict_tenant("alice");
3383 assert_eq!(store.stats().total_events, 0);
3384
3385 let results = store
3387 .query(&QueryEventsRequest {
3388 entity_id: None,
3389 event_type: None,
3390 tenant_id: Some("alice".to_string()),
3391 as_of: None,
3392 since: None,
3393 until: None,
3394 limit: None,
3395 event_type_prefix: None,
3396 exclude_event_type_prefix: None,
3397 payload_filter: None,
3398 })
3399 .unwrap();
3400 assert_eq!(results.len(), 4);
3401 assert!(store.is_tenant_loaded("alice"));
3402 }
3403
3404 #[test]
3405 fn test_evict_tenant_rebuilds_index_with_new_offsets() {
3406 let temp_dir = TempDir::new().unwrap();
3413 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3414
3415 for i in 0..3 {
3419 store
3420 .ingest(
3421 &Event::from_strings(
3422 "test.event".to_string(),
3423 format!("a-{i}"),
3424 "alice".to_string(),
3425 serde_json::json!({"i": i}),
3426 None,
3427 )
3428 .unwrap(),
3429 )
3430 .unwrap();
3431 if i < 2 {
3432 store
3433 .ingest(
3434 &Event::from_strings(
3435 "test.event".to_string(),
3436 format!("b-{i}"),
3437 "bob".to_string(),
3438 serde_json::json!({"i": i}),
3439 None,
3440 )
3441 .unwrap(),
3442 )
3443 .unwrap();
3444 }
3445 }
3446 store.tenant_loader.mark_loaded("alice");
3448 store.tenant_loader.mark_loaded("bob");
3449
3450 store.evict_tenant("alice");
3451
3452 let bob_results = store
3453 .query(&QueryEventsRequest {
3454 entity_id: None,
3455 event_type: None,
3456 tenant_id: Some("bob".to_string()),
3457 as_of: None,
3458 since: None,
3459 until: None,
3460 limit: None,
3461 event_type_prefix: None,
3462 exclude_event_type_prefix: None,
3463 payload_filter: None,
3464 })
3465 .unwrap();
3466 assert_eq!(bob_results.len(), 2);
3467 for e in &bob_results {
3468 assert_eq!(e.tenant_id_str(), "bob");
3469 }
3470 }
3471
3472 #[test]
3473 fn test_budget_eviction_keeps_resident_set_bounded() {
3474 let temp_dir = TempDir::new().unwrap();
3478 let storage_dir = temp_dir.path().to_path_buf();
3479
3480 let big_payload = serde_json::json!({"data": "x".repeat(1000)});
3483 {
3484 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3485 for tenant in ["alice", "bob", "carol"] {
3486 for i in 0..5 {
3487 store
3488 .ingest(
3489 &Event::from_strings(
3490 "test.event".to_string(),
3491 format!("{tenant}-{i}"),
3492 tenant.to_string(),
3493 big_payload.clone(),
3494 None,
3495 )
3496 .unwrap(),
3497 )
3498 .unwrap();
3499 }
3500 }
3501 store.flush_storage().unwrap();
3502 }
3503
3504 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3507 config.cache_byte_budget = Some(12_000);
3508 let store = EventStore::with_config(config);
3509
3510 store.ensure_tenant_loaded("alice").unwrap();
3512 assert!(store.is_tenant_loaded("alice"));
3513
3514 store.tenant_loader.touch("alice");
3519 std::thread::sleep(std::time::Duration::from_millis(10));
3520 store.ensure_tenant_loaded("bob").unwrap();
3521 assert!(store.is_tenant_loaded("bob"));
3522
3523 store.tenant_loader.touch("bob");
3527 std::thread::sleep(std::time::Duration::from_millis(10));
3528 store.ensure_tenant_loaded("carol").unwrap();
3529 assert!(store.is_tenant_loaded("carol"));
3530
3531 let resident = store.cache_resident_bytes();
3536 let budget = 12_000u64;
3537
3538 if resident > budget {
3541 let loaded_count = ["alice", "bob", "carol"]
3542 .iter()
3543 .filter(|t| store.is_tenant_loaded(t))
3544 .count();
3545 assert_eq!(
3546 loaded_count, 1,
3547 "over budget but more than one tenant loaded — eviction policy didn't fire"
3548 );
3549 }
3550
3551 assert!(store.is_tenant_loaded("carol"));
3554 }
3555
3556 #[test]
3557 fn test_query_after_eviction_re_loads_transparently() {
3558 let temp_dir = TempDir::new().unwrap();
3561 let storage_dir = temp_dir.path().to_path_buf();
3562
3563 let big_payload = serde_json::json!({"data": "x".repeat(2000)});
3564 {
3565 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3566 for tenant in ["alice", "bob"] {
3567 for i in 0..3 {
3568 store
3569 .ingest(
3570 &Event::from_strings(
3571 "test.event".to_string(),
3572 format!("{tenant}-{i}"),
3573 tenant.to_string(),
3574 big_payload.clone(),
3575 None,
3576 )
3577 .unwrap(),
3578 )
3579 .unwrap();
3580 }
3581 }
3582 store.flush_storage().unwrap();
3583 }
3584
3585 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3587 config.cache_byte_budget = Some(5_000);
3588 let store = EventStore::with_config(config);
3589
3590 let alice_first = store
3594 .query(&QueryEventsRequest {
3595 entity_id: None,
3596 event_type: None,
3597 tenant_id: Some("alice".to_string()),
3598 as_of: None,
3599 since: None,
3600 until: None,
3601 limit: None,
3602 event_type_prefix: None,
3603 exclude_event_type_prefix: None,
3604 payload_filter: None,
3605 })
3606 .unwrap();
3607 assert_eq!(alice_first.len(), 3);
3608
3609 std::thread::sleep(std::time::Duration::from_millis(15));
3611 let _bob = store
3613 .query(&QueryEventsRequest {
3614 entity_id: None,
3615 event_type: None,
3616 tenant_id: Some("bob".to_string()),
3617 as_of: None,
3618 since: None,
3619 until: None,
3620 limit: None,
3621 event_type_prefix: None,
3622 exclude_event_type_prefix: None,
3623 payload_filter: None,
3624 })
3625 .unwrap();
3626 assert!(
3627 !store.is_tenant_loaded("alice"),
3628 "alice should have been evicted"
3629 );
3630
3631 let alice_second = store
3633 .query(&QueryEventsRequest {
3634 entity_id: None,
3635 event_type: None,
3636 tenant_id: Some("alice".to_string()),
3637 as_of: None,
3638 since: None,
3639 until: None,
3640 limit: None,
3641 event_type_prefix: None,
3642 exclude_event_type_prefix: None,
3643 payload_filter: None,
3644 })
3645 .unwrap();
3646 assert_eq!(
3647 alice_second.len(),
3648 3,
3649 "alice's events come back via re-load"
3650 );
3651 assert!(store.is_tenant_loaded("alice"));
3652 }
3653
3654 #[test]
3655 #[cfg(feature = "server")]
3656 fn test_cache_metrics_track_evictions_and_bytes() {
3657 let temp_dir = TempDir::new().unwrap();
3661 let storage_dir = temp_dir.path().to_path_buf();
3662
3663 let big_payload = serde_json::json!({"data": "x".repeat(2000)});
3664 {
3665 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3666 for tenant in ["alice", "bob"] {
3667 for i in 0..3 {
3668 store
3669 .ingest(
3670 &Event::from_strings(
3671 "test.event".to_string(),
3672 format!("{tenant}-{i}"),
3673 tenant.to_string(),
3674 big_payload.clone(),
3675 None,
3676 )
3677 .unwrap(),
3678 )
3679 .unwrap();
3680 }
3681 }
3682 store.flush_storage().unwrap();
3683 }
3684
3685 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3686 config.cache_byte_budget = Some(5_000); let store = EventStore::with_config(config);
3688
3689 assert_eq!(store.metrics.cache_evictions_total.get(), 0);
3690 assert_eq!(store.metrics.cache_bytes.get(), 0);
3691
3692 store.ensure_tenant_loaded("alice").unwrap();
3693 let after_alice = store.metrics.cache_bytes.get();
3695 assert!(after_alice > 0, "gauge should reflect alice's bytes");
3696 assert_eq!(store.metrics.cache_evictions_total.get(), 0);
3698
3699 std::thread::sleep(std::time::Duration::from_millis(10));
3700 store.ensure_tenant_loaded("bob").unwrap();
3701
3702 assert_eq!(
3705 store.metrics.cache_evictions_total.get(),
3706 1,
3707 "exactly one tenant evicted after bob's load"
3708 );
3709 let after_bob = store.metrics.cache_bytes.get();
3711 assert!(after_bob > 0);
3712 assert!(after_bob <= after_alice, "gauge dropped after eviction");
3713 }
3714
3715 #[test]
3716 #[cfg(feature = "server")]
3717 fn test_storage_size_gauge_populated_from_on_disk_bytes() {
3718 let temp_dir = TempDir::new().unwrap();
3723 let config = EventStoreConfig {
3726 storage_dir: Some(temp_dir.path().join("parquet")),
3727 wal_dir: Some(temp_dir.path().join("wal")),
3728 ..EventStoreConfig::default()
3729 };
3730 let store = EventStore::with_config(config);
3731
3732 assert_eq!(store.metrics.storage_size_bytes.get(), 0);
3734
3735 let payload = serde_json::json!({ "data": "x".repeat(2000) });
3736 for i in 0..10 {
3737 store
3738 .ingest(
3739 &Event::from_strings(
3740 "test.event".to_string(),
3741 format!("entity-{i}"),
3742 "tenant-a".to_string(),
3743 payload.clone(),
3744 None,
3745 )
3746 .unwrap(),
3747 )
3748 .unwrap();
3749 }
3750 store.flush_storage().unwrap();
3751
3752 store.refresh_storage_metrics_now();
3754
3755 let size = store.metrics.storage_size_bytes.get();
3756 assert!(
3757 size > 0,
3758 "storage_size_bytes must reflect real on-disk bytes, got {size}"
3759 );
3760 assert!(
3761 store.metrics.parquet_files_total.get() >= 1,
3762 "at least one Parquet file should exist after a flush"
3763 );
3764
3765 assert!(
3768 store.metrics.wal_segments_total.get() >= 1,
3769 "at least one WAL segment should exist after writes"
3770 );
3771
3772 let parquet_stats = store.storage.as_ref().unwrap().read().stats().unwrap();
3775 let (wal_bytes, _) = store.wal.as_ref().unwrap().on_disk_stats().unwrap();
3776 assert_eq!(
3777 size as u64,
3778 parquet_stats.total_size_bytes + wal_bytes,
3779 "gauge should equal Parquet bytes ({}) + WAL bytes ({wal_bytes})",
3780 parquet_stats.total_size_bytes
3781 );
3782 }
3783
3784 #[test]
3785 fn test_stress_resident_set_stays_near_budget_under_rolling_queries() {
3786 let temp_dir = TempDir::new().unwrap();
3793 let storage_dir = temp_dir.path().to_path_buf();
3794
3795 const TENANT_COUNT: usize = 10;
3796 const EVENTS_PER_TENANT: usize = 50;
3797 let big_payload = serde_json::json!({"data": "x".repeat(10_000)});
3799
3800 {
3803 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3804 for t in 0..TENANT_COUNT {
3805 let tenant = format!("tenant-{t}");
3806 for i in 0..EVENTS_PER_TENANT {
3807 store
3808 .ingest(
3809 &Event::from_strings(
3810 "test.event".to_string(),
3811 format!("{tenant}-{i}"),
3812 tenant.clone(),
3813 big_payload.clone(),
3814 None,
3815 )
3816 .unwrap(),
3817 )
3818 .unwrap();
3819 }
3820 }
3821 store.flush_storage().unwrap();
3822 }
3823
3824 const BUDGET: u64 = 1_048_576;
3828 let mut config = EventStoreConfig::with_persistence(&storage_dir);
3829 config.cache_byte_budget = Some(BUDGET);
3830 let store = EventStore::with_config(config);
3831
3832 let mut peak_resident: u64 = 0;
3836 for t in 0..TENANT_COUNT {
3837 let tenant = format!("tenant-{t}");
3838 let results = store
3839 .query(&QueryEventsRequest {
3840 entity_id: None,
3841 event_type: None,
3842 tenant_id: Some(tenant.clone()),
3843 as_of: None,
3844 since: None,
3845 until: None,
3846 limit: None,
3847 event_type_prefix: None,
3848 exclude_event_type_prefix: None,
3849 payload_filter: None,
3850 })
3851 .unwrap();
3852 assert_eq!(
3853 results.len(),
3854 EVENTS_PER_TENANT,
3855 "every per-tenant query must return all of that tenant's events"
3856 );
3857 let resident = store.cache_resident_bytes();
3859 if resident > peak_resident {
3860 peak_resident = resident;
3861 }
3862 }
3863
3864 let final_resident = store.cache_resident_bytes();
3865
3866 let tolerance = BUDGET; assert!(
3872 peak_resident <= BUDGET + tolerance,
3873 "peak resident {peak_resident} exceeds budget {BUDGET} by more than {tolerance} \
3874 — eviction policy not keeping up with the working-set churn"
3875 );
3876 assert!(
3877 final_resident <= BUDGET + tolerance,
3878 "final resident {final_resident} exceeds budget {BUDGET} by more than {tolerance}"
3879 );
3880
3881 let last_tenant = format!("tenant-{}", TENANT_COUNT - 1);
3884 assert!(
3885 store.is_tenant_loaded(&last_tenant),
3886 "the most-recent tenant must remain loaded after the sweep"
3887 );
3888
3889 let still_loaded = (0..TENANT_COUNT)
3892 .filter(|t| store.is_tenant_loaded(&format!("tenant-{t}")))
3893 .count();
3894 assert!(
3895 still_loaded < TENANT_COUNT,
3896 "no tenants evicted ({still_loaded}/{TENANT_COUNT} still loaded) — \
3897 budget enforcement didn't engage"
3898 );
3899 }
3900
3901 #[test]
3902 fn test_evict_tenant_when_not_loaded_is_a_noop() {
3903 let store = EventStore::new();
3906 store.evict_tenant("nobody"); assert!(!store.is_tenant_loaded("nobody"));
3908 }
3909
3910 #[test]
3911 fn test_lazy_load_accounts_bytes_per_tenant() {
3912 let temp_dir = TempDir::new().unwrap();
3916 let storage_dir = temp_dir.path().to_path_buf();
3917
3918 {
3920 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3921 for i in 0..5 {
3922 store
3923 .ingest(
3924 &Event::from_strings(
3925 "test.event".to_string(),
3926 format!("a-{i}"),
3927 "alice".to_string(),
3928 serde_json::json!({"data": "x".repeat(1000)}),
3929 None,
3930 )
3931 .unwrap(),
3932 )
3933 .unwrap();
3934 }
3935 store.flush_storage().unwrap();
3936 }
3937
3938 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3939 assert_eq!(store.tenant_resident_bytes("alice"), 0);
3941 assert_eq!(store.cache_resident_bytes(), 0);
3942
3943 store.ensure_tenant_loaded("alice").unwrap();
3944
3945 let alice_bytes = store.tenant_resident_bytes("alice");
3948 assert!(
3949 alice_bytes >= 5 * 1000,
3950 "alice should have at least 5 KiB resident; got {alice_bytes}"
3951 );
3952 assert_eq!(store.tenant_resident_bytes("bob"), 0);
3954 assert_eq!(store.cache_resident_bytes(), alice_bytes);
3956 }
3957
3958 #[test]
3959 fn test_query_lazy_loads_tenant_on_first_call() {
3960 let temp_dir = TempDir::new().unwrap();
3964 let storage_dir = temp_dir.path().to_path_buf();
3965
3966 {
3968 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3969 for i in 0..3 {
3970 let event = Event::from_strings(
3971 "test.event".to_string(),
3972 format!("e-{i}"),
3973 "alice".to_string(),
3974 serde_json::json!({"i": i}),
3975 None,
3976 )
3977 .unwrap();
3978 store.ingest(&event).unwrap();
3979 }
3980 store.flush_storage().unwrap();
3981 }
3982
3983 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3985 assert_eq!(
3986 store.stats().total_events,
3987 0,
3988 "boot must be O(1) — no Parquet pre-load"
3989 );
3990 assert!(!store.is_tenant_loaded("alice"));
3991 assert!(!store.is_tenant_loaded("bob"));
3992
3993 let results = store
3995 .query(&QueryEventsRequest {
3996 entity_id: None,
3997 event_type: None,
3998 tenant_id: Some("alice".to_string()),
3999 as_of: None,
4000 since: None,
4001 until: None,
4002 limit: None,
4003 event_type_prefix: None,
4004 exclude_event_type_prefix: None,
4005 payload_filter: None,
4006 })
4007 .unwrap();
4008 assert_eq!(results.len(), 3, "alice's 3 events are returned");
4009 assert!(store.is_tenant_loaded("alice"), "alice now warm");
4010 assert!(!store.is_tenant_loaded("bob"), "bob still cold");
4013 }
4014
4015 #[test]
4016 fn test_query_invalid_tenant_id_returns_error_no_hang() {
4017 let temp_dir = TempDir::new().unwrap();
4021 let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
4022
4023 let result = store.query(&QueryEventsRequest {
4024 entity_id: None,
4025 event_type: None,
4026 tenant_id: Some("../etc".to_string()),
4027 as_of: None,
4028 since: None,
4029 until: None,
4030 limit: None,
4031 event_type_prefix: None,
4032 exclude_event_type_prefix: None,
4033 payload_filter: None,
4034 });
4035 assert!(result.is_err(), "unsafe tenant_id must surface as error");
4036 }
4037
4038 #[test]
4039 fn test_query_concurrent_first_queries_for_same_tenant_all_succeed() {
4040 let temp_dir = TempDir::new().unwrap();
4048 let storage_dir = temp_dir.path().to_path_buf();
4049
4050 {
4052 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
4053 for i in 0..25 {
4054 let event = Event::from_strings(
4055 "test.event".to_string(),
4056 format!("e-{i}"),
4057 "alice".to_string(),
4058 serde_json::json!({"i": i}),
4059 None,
4060 )
4061 .unwrap();
4062 store.ingest(&event).unwrap();
4063 }
4064 store.flush_storage().unwrap();
4065 }
4066
4067 let store = Arc::new(EventStore::with_config(EventStoreConfig::with_persistence(
4069 &storage_dir,
4070 )));
4071 assert!(!store.is_tenant_loaded("alice"));
4072
4073 let mut handles = Vec::new();
4074 for _ in 0..8 {
4075 let s = store.clone();
4076 handles.push(std::thread::spawn(move || {
4077 s.query(&QueryEventsRequest {
4078 entity_id: None,
4079 event_type: None,
4080 tenant_id: Some("alice".to_string()),
4081 as_of: None,
4082 since: None,
4083 until: None,
4084 limit: None,
4085 event_type_prefix: None,
4086 exclude_event_type_prefix: None,
4087 payload_filter: None,
4088 })
4089 }));
4090 }
4091
4092 for h in handles {
4093 let result = h.join().unwrap().unwrap();
4094 assert_eq!(
4095 result.len(),
4096 25,
4097 "every concurrent caller must see all 25 events"
4098 );
4099 }
4100 assert!(store.is_tenant_loaded("alice"));
4101 assert_eq!(store.stats().total_events, 25);
4103 }
4104
4105 #[test]
4106 fn test_query_two_cold_tenants_load_independently() {
4107 let temp_dir = TempDir::new().unwrap();
4111 let storage_dir = temp_dir.path().to_path_buf();
4112
4113 {
4114 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
4115 for i in 0..3 {
4116 store
4117 .ingest(
4118 &Event::from_strings(
4119 "test.event".to_string(),
4120 format!("a-{i}"),
4121 "alice".to_string(),
4122 serde_json::json!({"i": i}),
4123 None,
4124 )
4125 .unwrap(),
4126 )
4127 .unwrap();
4128 }
4129 for i in 0..5 {
4130 store
4131 .ingest(
4132 &Event::from_strings(
4133 "test.event".to_string(),
4134 format!("b-{i}"),
4135 "bob".to_string(),
4136 serde_json::json!({"i": i}),
4137 None,
4138 )
4139 .unwrap(),
4140 )
4141 .unwrap();
4142 }
4143 store.flush_storage().unwrap();
4144 }
4145
4146 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
4147 assert_eq!(store.stats().total_events, 0);
4148
4149 let alice = store
4151 .query(&QueryEventsRequest {
4152 entity_id: None,
4153 event_type: None,
4154 tenant_id: Some("alice".to_string()),
4155 as_of: None,
4156 since: None,
4157 until: None,
4158 limit: None,
4159 event_type_prefix: None,
4160 exclude_event_type_prefix: None,
4161 payload_filter: None,
4162 })
4163 .unwrap();
4164 assert_eq!(alice.len(), 3);
4165 assert!(store.is_tenant_loaded("alice"));
4166 assert!(!store.is_tenant_loaded("bob"));
4167 assert_eq!(store.stats().total_events, 3);
4168
4169 let bob = store
4171 .query(&QueryEventsRequest {
4172 entity_id: None,
4173 event_type: None,
4174 tenant_id: Some("bob".to_string()),
4175 as_of: None,
4176 since: None,
4177 until: None,
4178 limit: None,
4179 event_type_prefix: None,
4180 exclude_event_type_prefix: None,
4181 payload_filter: None,
4182 })
4183 .unwrap();
4184 assert_eq!(bob.len(), 5);
4185 assert!(store.is_tenant_loaded("bob"));
4186 assert_eq!(store.stats().total_events, 8);
4187 }
4188
4189 #[test]
4190 fn test_boot_with_persisted_data_is_o1() {
4191 let temp_dir = TempDir::new().unwrap();
4204 let storage_dir = temp_dir.path().to_path_buf();
4205
4206 {
4207 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
4208 for tenant in ["alice", "bob", "carol"] {
4209 for i in 0..50 / 3 {
4210 store
4211 .ingest(
4212 &Event::from_strings(
4213 "test.event".to_string(),
4214 format!("{tenant}-{i}"),
4215 tenant.to_string(),
4216 serde_json::json!({"i": i}),
4217 None,
4218 )
4219 .unwrap(),
4220 )
4221 .unwrap();
4222 }
4223 }
4224 store.flush_storage().unwrap();
4225 }
4226
4227 let on_disk = find_parquet_files(&storage_dir);
4229 assert!(
4230 !on_disk.is_empty(),
4231 "session 1 should have produced parquet files; pre-condition for the test"
4232 );
4233
4234 let started = std::time::Instant::now();
4235 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
4236 let boot_elapsed = started.elapsed();
4237
4238 assert_eq!(
4239 store.stats().total_events,
4240 0,
4241 "boot must not pre-load any Parquet events"
4242 );
4243
4244 assert!(
4248 boot_elapsed < std::time::Duration::from_secs(2),
4249 "boot took {boot_elapsed:?} — Step 2 boot should be O(1)"
4250 );
4251 }
4252
4253 #[test]
4254 fn test_query_warm_tenant_does_not_re_read_disk() {
4255 let temp_dir = TempDir::new().unwrap();
4261 let storage_dir = temp_dir.path().to_path_buf();
4262
4263 let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
4264 for i in 0..3 {
4265 let event = Event::from_strings(
4266 "test.event".to_string(),
4267 format!("e-{i}"),
4268 "alice".to_string(),
4269 serde_json::json!({"i": i}),
4270 None,
4271 )
4272 .unwrap();
4273 store.ingest(&event).unwrap();
4274 }
4275 store.flush_storage().unwrap();
4276
4277 let _ = store
4279 .query(&QueryEventsRequest {
4280 entity_id: None,
4281 event_type: None,
4282 tenant_id: Some("alice".to_string()),
4283 as_of: None,
4284 since: None,
4285 until: None,
4286 limit: None,
4287 event_type_prefix: None,
4288 exclude_event_type_prefix: None,
4289 payload_filter: None,
4290 })
4291 .unwrap();
4292 assert!(store.is_tenant_loaded("alice"));
4293
4294 let parquet_files = find_parquet_files(&storage_dir);
4297 for f in parquet_files {
4298 std::fs::remove_file(&f).unwrap();
4299 }
4300
4301 let results = store
4302 .query(&QueryEventsRequest {
4303 entity_id: None,
4304 event_type: None,
4305 tenant_id: Some("alice".to_string()),
4306 as_of: None,
4307 since: None,
4308 until: None,
4309 limit: None,
4310 event_type_prefix: None,
4311 exclude_event_type_prefix: None,
4312 payload_filter: None,
4313 })
4314 .unwrap();
4315 assert_eq!(
4316 results.len(),
4317 3,
4318 "warm tenant query must not need disk; got {} events from a deleted parquet",
4319 results.len()
4320 );
4321 }
4322
4323 #[test]
4324 fn test_event_store_default() {
4325 let store = EventStore::default();
4326 assert_eq!(store.stats().total_events, 0);
4327 }
4328
4329 #[test]
4330 fn test_ingest_single_event() {
4331 let store = EventStore::new();
4332 let event = create_test_event("entity-1", "user.created");
4333
4334 store.ingest(&event).unwrap();
4335
4336 assert_eq!(store.stats().total_events, 1);
4337 assert_eq!(store.stats().total_ingested, 1);
4338 }
4339
4340 #[test]
4341 fn test_ingest_multiple_events() {
4342 let store = EventStore::new();
4343
4344 for i in 0..10 {
4345 let event = create_test_event(&format!("entity-{i}"), "user.created");
4346 store.ingest(&event).unwrap();
4347 }
4348
4349 assert_eq!(store.stats().total_events, 10);
4350 assert_eq!(store.stats().total_ingested, 10);
4351 }
4352
4353 #[test]
4354 fn test_query_by_entity_id() {
4355 let store = EventStore::new();
4356
4357 store
4358 .ingest(&create_test_event("entity-1", "user.created"))
4359 .unwrap();
4360 store
4361 .ingest(&create_test_event("entity-2", "user.created"))
4362 .unwrap();
4363 store
4364 .ingest(&create_test_event("entity-1", "user.updated"))
4365 .unwrap();
4366
4367 let results = store
4368 .query(&QueryEventsRequest {
4369 entity_id: Some("entity-1".to_string()),
4370 event_type: None,
4371 tenant_id: None,
4372 as_of: None,
4373 since: None,
4374 until: None,
4375 limit: None,
4376 event_type_prefix: None,
4377 exclude_event_type_prefix: None,
4378 payload_filter: None,
4379 })
4380 .unwrap();
4381
4382 assert_eq!(results.len(), 2);
4383 }
4384
4385 #[test]
4386 fn test_query_by_event_type() {
4387 let store = EventStore::new();
4388
4389 store
4390 .ingest(&create_test_event("entity-1", "user.created"))
4391 .unwrap();
4392 store
4393 .ingest(&create_test_event("entity-2", "user.updated"))
4394 .unwrap();
4395 store
4396 .ingest(&create_test_event("entity-3", "user.created"))
4397 .unwrap();
4398
4399 let results = store
4400 .query(&QueryEventsRequest {
4401 entity_id: None,
4402 event_type: Some("user.created".to_string()),
4403 tenant_id: None,
4404 as_of: None,
4405 since: None,
4406 until: None,
4407 limit: None,
4408 event_type_prefix: None,
4409 exclude_event_type_prefix: None,
4410 payload_filter: None,
4411 })
4412 .unwrap();
4413
4414 assert_eq!(results.len(), 2);
4415 }
4416
4417 #[test]
4418 fn test_query_with_limit() {
4419 let store = EventStore::new();
4420
4421 for i in 0..10 {
4422 let event = create_test_event(&format!("entity-{i}"), "user.created");
4423 store.ingest(&event).unwrap();
4424 }
4425
4426 let results = store
4427 .query(&QueryEventsRequest {
4428 entity_id: None,
4429 event_type: None,
4430 tenant_id: None,
4431 as_of: None,
4432 since: None,
4433 until: None,
4434 limit: Some(5),
4435 event_type_prefix: None,
4436 exclude_event_type_prefix: None,
4437 payload_filter: None,
4438 })
4439 .unwrap();
4440
4441 assert_eq!(results.len(), 5);
4442 }
4443
4444 #[test]
4445 fn test_query_empty_store() {
4446 let store = EventStore::new();
4447
4448 let results = store
4449 .query(&QueryEventsRequest {
4450 entity_id: Some("non-existent".to_string()),
4451 event_type: None,
4452 tenant_id: None,
4453 as_of: None,
4454 since: None,
4455 until: None,
4456 limit: None,
4457 event_type_prefix: None,
4458 exclude_event_type_prefix: None,
4459 payload_filter: None,
4460 })
4461 .unwrap();
4462
4463 assert!(results.is_empty());
4464 }
4465
4466 #[test]
4467 fn test_reconstruct_state() {
4468 let store = EventStore::new();
4469
4470 store
4471 .ingest(&create_test_event("entity-1", "user.created"))
4472 .unwrap();
4473
4474 let state = store.reconstruct_state("entity-1", None).unwrap();
4475 assert_eq!(state["current_state"]["name"], "Test");
4477 assert_eq!(state["current_state"]["value"], 42);
4478 }
4479
4480 #[test]
4481 fn test_reconstruct_state_not_found() {
4482 let store = EventStore::new();
4483
4484 let result = store.reconstruct_state("non-existent", None);
4485 assert!(result.is_err());
4486 }
4487
4488 #[test]
4489 fn test_get_snapshot_empty() {
4490 let store = EventStore::new();
4491
4492 let result = store.get_snapshot("non-existent");
4493 assert!(result.is_err());
4495 }
4496
4497 #[test]
4498 fn test_create_snapshot() {
4499 let store = EventStore::new();
4500
4501 store
4502 .ingest(&create_test_event("entity-1", "user.created"))
4503 .unwrap();
4504
4505 store.create_snapshot("entity-1").unwrap();
4506
4507 let snapshot = store.get_snapshot("entity-1").unwrap();
4509 assert_ne!(snapshot, serde_json::json!(null));
4510 }
4511
4512 #[test]
4513 fn test_create_snapshot_entity_not_found() {
4514 let store = EventStore::new();
4515
4516 let result = store.create_snapshot("non-existent");
4517 assert!(result.is_err());
4518 }
4519
4520 #[test]
4521 fn test_websocket_manager() {
4522 let store = EventStore::new();
4523 let manager = store.websocket_manager();
4524 assert!(Arc::strong_count(&manager) >= 1);
4526 }
4527
4528 #[test]
4529 fn test_snapshot_manager() {
4530 let store = EventStore::new();
4531 let manager = store.snapshot_manager();
4532 assert!(Arc::strong_count(&manager) >= 1);
4533 }
4534
4535 #[test]
4536 fn test_compaction_manager_none() {
4537 let store = EventStore::new();
4538 assert!(store.compaction_manager().is_none());
4540 }
4541
4542 #[test]
4543 fn test_schema_registry() {
4544 let store = EventStore::new();
4545 let registry = store.schema_registry();
4546 assert!(Arc::strong_count(®istry) >= 1);
4547 }
4548
4549 #[test]
4550 fn test_replay_manager() {
4551 let store = EventStore::new();
4552 let manager = store.replay_manager();
4553 assert!(Arc::strong_count(&manager) >= 1);
4554 }
4555
4556 #[test]
4557 fn test_pipeline_manager() {
4558 let store = EventStore::new();
4559 let manager = store.pipeline_manager();
4560 assert!(Arc::strong_count(&manager) >= 1);
4561 }
4562
4563 #[test]
4564 fn test_projection_manager() {
4565 let store = EventStore::new();
4566 let manager = store.projection_manager();
4567 let projections = manager.list_projections();
4569 assert!(projections.len() >= 2); }
4571
4572 #[test]
4573 fn test_projection_state_cache() {
4574 let store = EventStore::new();
4575 let cache = store.projection_state_cache();
4576
4577 cache.insert("test:key".to_string(), serde_json::json!({"value": 123}));
4578 assert_eq!(cache.len(), 1);
4579
4580 let value = cache.get("test:key").unwrap();
4581 assert_eq!(value["value"], 123);
4582 }
4583
4584 #[test]
4585 fn test_metrics() {
4586 let store = EventStore::new();
4587 let metrics = store.metrics();
4588 assert!(Arc::strong_count(&metrics) >= 1);
4589 }
4590
4591 #[test]
4592 fn test_store_stats() {
4593 let store = EventStore::new();
4594
4595 store
4596 .ingest(&create_test_event("entity-1", "user.created"))
4597 .unwrap();
4598 store
4599 .ingest(&create_test_event("entity-2", "order.placed"))
4600 .unwrap();
4601
4602 let stats = store.stats();
4603 assert_eq!(stats.total_events, 2);
4604 assert_eq!(stats.total_entities, 2);
4605 assert_eq!(stats.total_event_types, 2);
4606 assert_eq!(stats.total_ingested, 2);
4607 }
4608
4609 #[test]
4610 fn test_event_store_config_default() {
4611 let config = EventStoreConfig::default();
4612 assert!(config.storage_dir.is_none());
4613 assert!(config.wal_dir.is_none());
4614 }
4615
4616 #[test]
4617 fn test_event_store_config_with_persistence() {
4618 let temp_dir = TempDir::new().unwrap();
4619 let config = EventStoreConfig::with_persistence(temp_dir.path());
4620
4621 assert!(config.storage_dir.is_some());
4622 assert!(config.wal_dir.is_none());
4623 }
4624
4625 #[test]
4626 fn test_event_store_config_with_wal() {
4627 let temp_dir = TempDir::new().unwrap();
4628 let config = EventStoreConfig::with_wal(temp_dir.path(), WALConfig::default());
4629
4630 assert!(config.storage_dir.is_none());
4631 assert!(config.wal_dir.is_some());
4632 }
4633
4634 #[test]
4635 fn test_event_store_config_with_all() {
4636 let temp_dir = TempDir::new().unwrap();
4637 let config = EventStoreConfig::with_all(temp_dir.path(), SnapshotConfig::default());
4638
4639 assert!(config.storage_dir.is_some());
4640 }
4641
4642 #[test]
4643 fn test_event_store_config_production() {
4644 let storage_dir = TempDir::new().unwrap();
4645 let wal_dir = TempDir::new().unwrap();
4646 let config = EventStoreConfig::production(
4647 storage_dir.path(),
4648 wal_dir.path(),
4649 SnapshotConfig::default(),
4650 WALConfig::default(),
4651 CompactionConfig::default(),
4652 );
4653
4654 assert!(config.storage_dir.is_some());
4655 assert!(config.wal_dir.is_some());
4656 }
4657
4658 #[test]
4664 fn test_from_env_vars_data_dir_enables_full_persistence() {
4665 let (config, mode) = EventStoreConfig::from_env_vars(
4666 Some("/app/data".to_string()),
4667 None,
4668 None,
4669 None,
4670 None,
4671 None,
4672 None,
4673 None,
4674 );
4675 assert_eq!(mode, "wal+parquet");
4676 assert_eq!(
4677 config.storage_dir.unwrap().to_str().unwrap(),
4678 "/app/data/storage"
4679 );
4680 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/app/data/wal");
4681 }
4682
4683 #[test]
4684 fn test_from_env_vars_explicit_dirs() {
4685 let (config, mode) = EventStoreConfig::from_env_vars(
4686 None,
4687 Some("/custom/storage".to_string()),
4688 Some("/custom/wal".to_string()),
4689 None,
4690 None,
4691 None,
4692 None,
4693 None,
4694 );
4695 assert_eq!(mode, "wal+parquet");
4696 assert_eq!(
4697 config.storage_dir.unwrap().to_str().unwrap(),
4698 "/custom/storage"
4699 );
4700 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/custom/wal");
4701 }
4702
4703 #[test]
4704 fn test_from_env_vars_wal_disabled() {
4705 let (config, mode) = EventStoreConfig::from_env_vars(
4706 Some("/app/data".to_string()),
4707 None,
4708 None,
4709 Some("false".to_string()),
4710 None,
4711 None,
4712 None,
4713 None,
4714 );
4715 assert_eq!(mode, "parquet-only");
4716 assert!(config.storage_dir.is_some());
4717 assert!(config.wal_dir.is_none());
4718 }
4719
4720 #[test]
4721 fn test_from_env_vars_no_dirs_is_in_memory() {
4722 let (config, mode) =
4723 EventStoreConfig::from_env_vars(None, None, None, None, None, None, None, None);
4724 assert_eq!(mode, "in-memory");
4725 assert!(config.storage_dir.is_none());
4726 assert!(config.wal_dir.is_none());
4727 }
4728
4729 #[test]
4730 fn test_from_env_vars_empty_strings_treated_as_none() {
4731 let (_, mode) = EventStoreConfig::from_env_vars(
4732 Some(String::new()),
4733 Some(String::new()),
4734 Some(String::new()),
4735 None,
4736 None,
4737 None,
4738 None,
4739 None,
4740 );
4741 assert_eq!(mode, "in-memory");
4742 }
4743
4744 #[test]
4745 fn test_from_env_vars_explicit_overrides_data_dir() {
4746 let (config, mode) = EventStoreConfig::from_env_vars(
4747 Some("/app/data".to_string()),
4748 Some("/override/storage".to_string()),
4749 Some("/override/wal".to_string()),
4750 None,
4751 None,
4752 None,
4753 None,
4754 None,
4755 );
4756 assert_eq!(mode, "wal+parquet");
4757 assert_eq!(
4758 config.storage_dir.unwrap().to_str().unwrap(),
4759 "/override/storage"
4760 );
4761 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/override/wal");
4762 }
4763
4764 #[test]
4765 fn test_from_env_vars_wal_only() {
4766 let (config, mode) = EventStoreConfig::from_env_vars(
4767 None,
4768 None,
4769 Some("/wal/only".to_string()),
4770 None,
4771 None,
4772 None,
4773 None,
4774 None,
4775 );
4776 assert_eq!(mode, "wal-only");
4777 assert!(config.storage_dir.is_none());
4778 assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/wal/only");
4779 }
4780
4781 #[test]
4782 fn test_from_env_vars_cache_bytes_parses_decimal() {
4783 let (config, _) = EventStoreConfig::from_env_vars(
4784 Some("/app/data".to_string()),
4785 None,
4786 None,
4787 None,
4788 Some("536870912".to_string()),
4789 None,
4791 None,
4792 None,
4793 );
4794 assert_eq!(config.cache_byte_budget, Some(536_870_912));
4795 }
4796
4797 #[test]
4798 fn test_from_env_vars_cache_bytes_unparseable_disables_budget() {
4799 let (config, _) = EventStoreConfig::from_env_vars(
4803 Some("/app/data".to_string()),
4804 None,
4805 None,
4806 None,
4807 Some("not-a-number".to_string()),
4808 None,
4809 None,
4810 None,
4811 );
4812 assert_eq!(config.cache_byte_budget, None);
4813 }
4814
4815 #[test]
4816 fn test_from_env_vars_cache_bytes_empty_disables_budget() {
4817 let (config, _) = EventStoreConfig::from_env_vars(
4818 Some("/app/data".to_string()),
4819 None,
4820 None,
4821 None,
4822 Some(String::new()),
4823 None,
4824 None,
4825 None,
4826 );
4827 assert_eq!(config.cache_byte_budget, None);
4828 }
4829
4830 #[test]
4831 fn test_from_env_vars_snapshot_interval_overrides_default() {
4832 let (config, _) = EventStoreConfig::from_env_vars(
4836 Some("/app/data".to_string()),
4837 None,
4838 None,
4839 None,
4840 None,
4841 Some("60".to_string()),
4842 None,
4843 None,
4844 );
4845 assert_eq!(config.compaction_config.compaction_interval_seconds, 60);
4846 }
4847
4848 #[test]
4849 fn test_from_env_vars_snapshot_interval_default_is_hourly() {
4850 let (config, _) = EventStoreConfig::from_env_vars(
4851 Some("/app/data".to_string()),
4852 None,
4853 None,
4854 None,
4855 None,
4856 None,
4857 None,
4858 None,
4859 );
4860 assert_eq!(config.compaction_config.compaction_interval_seconds, 3600);
4861 }
4862
4863 #[test]
4864 fn test_from_env_vars_snapshot_interval_unparseable_falls_back() {
4865 let (config, _) = EventStoreConfig::from_env_vars(
4866 Some("/app/data".to_string()),
4867 None,
4868 None,
4869 None,
4870 None,
4871 Some("not-a-number".to_string()),
4872 None,
4873 None,
4874 );
4875 assert_eq!(config.compaction_config.compaction_interval_seconds, 3600);
4876 }
4877
4878 #[test]
4879 fn test_from_env_vars_retention_system_days_overrides_default() {
4880 let (config, _) = EventStoreConfig::from_env_vars(
4883 Some("/app/data".to_string()),
4884 None,
4885 None,
4886 None,
4887 None,
4888 None,
4889 Some("7".to_string()),
4890 None,
4891 );
4892 let ttl = config
4893 .compaction_config
4894 .retention
4895 .ttl_for("system")
4896 .unwrap();
4897 assert_eq!(ttl.as_secs(), 7 * 24 * 3600);
4898 }
4899
4900 #[test]
4901 fn test_from_env_vars_retention_default_is_30_days_for_system() {
4902 let (config, _) = EventStoreConfig::from_env_vars(
4903 Some("/app/data".to_string()),
4904 None,
4905 None,
4906 None,
4907 None,
4908 None,
4909 None,
4910 None,
4911 );
4912 let ttl = config
4913 .compaction_config
4914 .retention
4915 .ttl_for("system")
4916 .unwrap();
4917 assert_eq!(ttl.as_secs(), 30 * 24 * 3600);
4918 assert!(config.compaction_config.retention.ttl_for("acme").is_none());
4920 }
4921
4922 #[test]
4923 fn test_store_stats_serde() {
4924 let stats = StoreStats {
4925 total_events: 100,
4926 total_entities: 50,
4927 total_event_types: 10,
4928 total_ingested: 100,
4929 };
4930
4931 let json = serde_json::to_string(&stats).unwrap();
4932 assert!(json.contains("\"total_events\":100"));
4933 assert!(json.contains("\"total_entities\":50"));
4934 }
4935
4936 #[test]
4937 fn test_query_with_entity_and_type() {
4938 let store = EventStore::new();
4939
4940 store
4941 .ingest(&create_test_event("entity-1", "user.created"))
4942 .unwrap();
4943 store
4944 .ingest(&create_test_event("entity-1", "user.updated"))
4945 .unwrap();
4946 store
4947 .ingest(&create_test_event("entity-2", "user.created"))
4948 .unwrap();
4949
4950 let results = store
4951 .query(&QueryEventsRequest {
4952 entity_id: Some("entity-1".to_string()),
4953 event_type: Some("user.created".to_string()),
4954 tenant_id: None,
4955 as_of: None,
4956 since: None,
4957 until: None,
4958 limit: None,
4959 event_type_prefix: None,
4960 exclude_event_type_prefix: None,
4961 payload_filter: None,
4962 })
4963 .unwrap();
4964
4965 assert_eq!(results.len(), 1);
4966 assert_eq!(results[0].event_type_str(), "user.created");
4967 }
4968
4969 #[test]
4970 fn test_query_by_event_type_prefix() {
4971 let store = EventStore::new();
4972
4973 store
4975 .ingest(&create_test_event("entity-1", "index.created"))
4976 .unwrap();
4977 store
4978 .ingest(&create_test_event("entity-2", "index.updated"))
4979 .unwrap();
4980 store
4981 .ingest(&create_test_event("entity-3", "trade.created"))
4982 .unwrap();
4983 store
4984 .ingest(&create_test_event("entity-4", "trade.completed"))
4985 .unwrap();
4986 store
4987 .ingest(&create_test_event("entity-5", "balance.updated"))
4988 .unwrap();
4989
4990 let results = store
4992 .query(&QueryEventsRequest {
4993 entity_id: None,
4994 event_type: None,
4995 tenant_id: None,
4996 as_of: None,
4997 since: None,
4998 until: None,
4999 limit: None,
5000 event_type_prefix: Some("index.".to_string()),
5001 exclude_event_type_prefix: None,
5002 payload_filter: None,
5003 })
5004 .unwrap();
5005
5006 assert_eq!(results.len(), 2);
5007 assert!(
5008 results
5009 .iter()
5010 .all(|e| e.event_type_str().starts_with("index."))
5011 );
5012 }
5013
5014 #[test]
5015 fn test_query_by_event_type_prefix_empty_returns_all() {
5016 let store = EventStore::new();
5017
5018 store
5019 .ingest(&create_test_event("entity-1", "index.created"))
5020 .unwrap();
5021 store
5022 .ingest(&create_test_event("entity-2", "trade.created"))
5023 .unwrap();
5024
5025 let results = store
5027 .query(&QueryEventsRequest {
5028 entity_id: None,
5029 event_type: None,
5030 tenant_id: None,
5031 as_of: None,
5032 since: None,
5033 until: None,
5034 limit: None,
5035 event_type_prefix: Some(String::new()),
5036 exclude_event_type_prefix: None,
5037 payload_filter: None,
5038 })
5039 .unwrap();
5040
5041 assert_eq!(results.len(), 2);
5042 }
5043
5044 #[test]
5045 fn test_query_by_event_type_prefix_no_match() {
5046 let store = EventStore::new();
5047
5048 store
5049 .ingest(&create_test_event("entity-1", "index.created"))
5050 .unwrap();
5051
5052 let results = store
5053 .query(&QueryEventsRequest {
5054 entity_id: None,
5055 event_type: None,
5056 tenant_id: None,
5057 as_of: None,
5058 since: None,
5059 until: None,
5060 limit: None,
5061 event_type_prefix: Some("nonexistent.".to_string()),
5062 exclude_event_type_prefix: None,
5063 payload_filter: None,
5064 })
5065 .unwrap();
5066
5067 assert!(results.is_empty());
5068 }
5069
5070 #[test]
5071 fn test_query_by_entity_with_type_prefix() {
5072 let store = EventStore::new();
5073
5074 store
5075 .ingest(&create_test_event("entity-1", "index.created"))
5076 .unwrap();
5077 store
5078 .ingest(&create_test_event("entity-1", "trade.created"))
5079 .unwrap();
5080 store
5081 .ingest(&create_test_event("entity-2", "index.updated"))
5082 .unwrap();
5083
5084 let results = store
5086 .query(&QueryEventsRequest {
5087 entity_id: Some("entity-1".to_string()),
5088 event_type: None,
5089 tenant_id: None,
5090 as_of: None,
5091 since: None,
5092 until: None,
5093 limit: None,
5094 event_type_prefix: Some("index.".to_string()),
5095 exclude_event_type_prefix: None,
5096 payload_filter: None,
5097 })
5098 .unwrap();
5099
5100 assert_eq!(results.len(), 1);
5101 assert_eq!(results[0].event_type_str(), "index.created");
5102 }
5103
5104 #[test]
5105 fn test_query_prefix_with_limit() {
5106 let store = EventStore::new();
5107
5108 for i in 0..5 {
5109 store
5110 .ingest(&create_test_event(&format!("entity-{i}"), "index.created"))
5111 .unwrap();
5112 }
5113
5114 let results = store
5115 .query(&QueryEventsRequest {
5116 entity_id: None,
5117 event_type: None,
5118 tenant_id: None,
5119 as_of: None,
5120 since: None,
5121 until: None,
5122 limit: Some(3),
5123 event_type_prefix: Some("index.".to_string()),
5124 exclude_event_type_prefix: None,
5125 payload_filter: None,
5126 })
5127 .unwrap();
5128
5129 assert_eq!(results.len(), 3);
5130 }
5131
5132 #[test]
5133 fn test_query_prefix_alongside_existing_filters() {
5134 let store = EventStore::new();
5135
5136 store
5137 .ingest(&create_test_event("entity-1", "index.created"))
5138 .unwrap();
5139 std::thread::sleep(std::time::Duration::from_millis(10));
5141 store
5142 .ingest(&create_test_event("entity-2", "index.strategy.updated"))
5143 .unwrap();
5144 std::thread::sleep(std::time::Duration::from_millis(10));
5145 store
5146 .ingest(&create_test_event("entity-3", "index.deleted"))
5147 .unwrap();
5148
5149 let results = store
5151 .query(&QueryEventsRequest {
5152 entity_id: None,
5153 event_type: None,
5154 tenant_id: None,
5155 as_of: None,
5156 since: None,
5157 until: None,
5158 limit: Some(2),
5159 event_type_prefix: Some("index.".to_string()),
5160 exclude_event_type_prefix: None,
5161 payload_filter: None,
5162 })
5163 .unwrap();
5164
5165 assert_eq!(results.len(), 2);
5166 }
5167
5168 #[test]
5169 fn test_query_with_payload_filter() {
5170 let store = EventStore::new();
5171
5172 for i in 0..5 {
5174 store
5175 .ingest(&create_test_event_with_payload(
5176 &format!("entity-{i}"),
5177 "user.action",
5178 serde_json::json!({"user_id": "alice", "action": "click"}),
5179 ))
5180 .unwrap();
5181 }
5182 for i in 5..10 {
5184 store
5185 .ingest(&create_test_event_with_payload(
5186 &format!("entity-{i}"),
5187 "user.action",
5188 serde_json::json!({"user_id": "bob", "action": "view"}),
5189 ))
5190 .unwrap();
5191 }
5192
5193 let results = store
5195 .query(&QueryEventsRequest {
5196 entity_id: None,
5197 event_type: Some("user.action".to_string()),
5198 tenant_id: None,
5199 as_of: None,
5200 since: None,
5201 until: None,
5202 limit: None,
5203 event_type_prefix: None,
5204 exclude_event_type_prefix: None,
5205 payload_filter: Some(r#"{"user_id":"alice"}"#.to_string()),
5206 })
5207 .unwrap();
5208
5209 assert_eq!(results.len(), 5);
5210 }
5211
5212 #[test]
5213 fn test_query_payload_filter_non_existent_field() {
5214 let store = EventStore::new();
5215
5216 store
5217 .ingest(&create_test_event_with_payload(
5218 "entity-1",
5219 "user.action",
5220 serde_json::json!({"user_id": "alice"}),
5221 ))
5222 .unwrap();
5223
5224 let results = store
5226 .query(&QueryEventsRequest {
5227 entity_id: None,
5228 event_type: None,
5229 tenant_id: None,
5230 as_of: None,
5231 since: None,
5232 until: None,
5233 limit: None,
5234 event_type_prefix: None,
5235 exclude_event_type_prefix: None,
5236 payload_filter: Some(r#"{"nonexistent":"value"}"#.to_string()),
5237 })
5238 .unwrap();
5239
5240 assert!(results.is_empty());
5241 }
5242
5243 #[test]
5244 fn test_query_payload_filter_with_prefix() {
5245 let store = EventStore::new();
5246
5247 store
5248 .ingest(&create_test_event_with_payload(
5249 "entity-1",
5250 "index.created",
5251 serde_json::json!({"status": "active"}),
5252 ))
5253 .unwrap();
5254 store
5255 .ingest(&create_test_event_with_payload(
5256 "entity-2",
5257 "index.created",
5258 serde_json::json!({"status": "inactive"}),
5259 ))
5260 .unwrap();
5261 store
5262 .ingest(&create_test_event_with_payload(
5263 "entity-3",
5264 "trade.created",
5265 serde_json::json!({"status": "active"}),
5266 ))
5267 .unwrap();
5268
5269 let results = store
5271 .query(&QueryEventsRequest {
5272 entity_id: None,
5273 event_type: None,
5274 tenant_id: None,
5275 as_of: None,
5276 since: None,
5277 until: None,
5278 limit: None,
5279 event_type_prefix: Some("index.".to_string()),
5280 exclude_event_type_prefix: None,
5281 payload_filter: Some(r#"{"status":"active"}"#.to_string()),
5282 })
5283 .unwrap();
5284
5285 assert_eq!(results.len(), 1);
5286 assert_eq!(results[0].entity_id().to_string(), "entity-1");
5287 }
5288
5289 #[test]
5290 fn test_flush_storage_no_storage() {
5291 let store = EventStore::new();
5292 let result = store.flush_storage();
5294 assert!(result.is_ok());
5295 }
5296
5297 #[test]
5298 fn test_state_evolution() {
5299 let store = EventStore::new();
5300
5301 store
5303 .ingest(
5304 &Event::from_strings(
5305 "user.created".to_string(),
5306 "user-1".to_string(),
5307 "default".to_string(),
5308 serde_json::json!({"name": "Alice", "age": 25}),
5309 None,
5310 )
5311 .unwrap(),
5312 )
5313 .unwrap();
5314
5315 store
5317 .ingest(
5318 &Event::from_strings(
5319 "user.updated".to_string(),
5320 "user-1".to_string(),
5321 "default".to_string(),
5322 serde_json::json!({"age": 26}),
5323 None,
5324 )
5325 .unwrap(),
5326 )
5327 .unwrap();
5328
5329 let state = store.reconstruct_state("user-1", None).unwrap();
5330 assert_eq!(state["current_state"]["name"], "Alice");
5332 assert_eq!(state["current_state"]["age"], 26);
5333 }
5334
5335 #[test]
5336 fn test_reject_system_event_types() {
5337 let store = EventStore::new();
5338
5339 let event = Event::reconstruct_from_strings(
5341 uuid::Uuid::new_v4(),
5342 "_system.tenant.created".to_string(),
5343 "_system:tenant:acme".to_string(),
5344 "_system".to_string(),
5345 serde_json::json!({"name": "ACME"}),
5346 chrono::Utc::now(),
5347 None,
5348 1,
5349 );
5350
5351 let result = store.ingest(&event);
5352 assert!(result.is_err());
5353 let err = result.unwrap_err();
5354 assert!(
5355 err.to_string().contains("reserved for internal use"),
5356 "Expected system namespace rejection, got: {err}"
5357 );
5358 }
5359
5360 #[test]
5368 fn test_wal_recovery_checkpoints_to_parquet() {
5369 let data_dir = TempDir::new().unwrap();
5370 let storage_dir = data_dir.path().join("storage");
5371 let wal_dir = data_dir.path().join("wal");
5372
5373 {
5375 let config = EventStoreConfig::production(
5376 &storage_dir,
5377 &wal_dir,
5378 SnapshotConfig::default(),
5379 WALConfig {
5380 sync_on_write: true,
5381 ..WALConfig::default()
5382 },
5383 CompactionConfig::default(),
5384 );
5385 let store = EventStore::with_config(config);
5386
5387 for i in 0..5 {
5388 let event = Event::from_strings(
5389 "test.created".to_string(),
5390 format!("entity-{i}"),
5391 "default".to_string(),
5392 serde_json::json!({"index": i}),
5393 None,
5394 )
5395 .unwrap();
5396 store.ingest(&event).unwrap();
5397 }
5398
5399 assert_eq!(store.stats().total_events, 5);
5400
5401 }
5404
5405 let wal_files: Vec<_> = std::fs::read_dir(&wal_dir)
5407 .unwrap()
5408 .filter_map(std::result::Result::ok)
5409 .filter(|e| e.path().extension().is_some_and(|ext| ext == "log"))
5410 .collect();
5411 assert!(!wal_files.is_empty(), "WAL file should exist");
5412 let wal_size = wal_files[0].metadata().unwrap().len();
5413 assert!(wal_size > 0, "WAL file should have data (got 0 bytes)");
5414
5415 {
5417 let config = EventStoreConfig::production(
5418 &storage_dir,
5419 &wal_dir,
5420 SnapshotConfig::default(),
5421 WALConfig {
5422 sync_on_write: true,
5423 ..WALConfig::default()
5424 },
5425 CompactionConfig::default(),
5426 );
5427 let store = EventStore::with_config(config);
5428
5429 assert_eq!(
5431 store.stats().total_events,
5432 5,
5433 "Session 2 should have all 5 events after WAL recovery"
5434 );
5435
5436 let parquet_files = find_parquet_files(&storage_dir);
5440 assert!(
5441 !parquet_files.is_empty(),
5442 "Parquet file should exist after WAL checkpoint"
5443 );
5444 }
5445
5446 {
5449 let config = EventStoreConfig::production(
5450 &storage_dir,
5451 &wal_dir,
5452 SnapshotConfig::default(),
5453 WALConfig {
5454 sync_on_write: true,
5455 ..WALConfig::default()
5456 },
5457 CompactionConfig::default(),
5458 );
5459 let store = EventStore::with_config(config);
5460
5461 assert_eq!(
5465 store.stats().total_events,
5466 0,
5467 "Session 3 boot should not pre-load Parquet (lazy-load mode)"
5468 );
5469
5470 store.ensure_tenant_loaded("default").unwrap();
5473 assert_eq!(
5474 store.stats().total_events,
5475 5,
5476 "Session 3 should have all 5 events after ensure_tenant_loaded"
5477 );
5478 }
5479 }
5480
5481 #[test]
5482 fn test_parquet_restore_surfaces_errors_not_silent() {
5483 let data_dir = TempDir::new().unwrap();
5487 let storage_dir = data_dir.path().join("storage");
5488 let wal_dir = data_dir.path().join("wal");
5489
5490 {
5492 let config = EventStoreConfig::production(
5493 &storage_dir,
5494 &wal_dir,
5495 SnapshotConfig::default(),
5496 WALConfig {
5497 sync_on_write: true,
5498 ..WALConfig::default()
5499 },
5500 CompactionConfig::default(),
5501 );
5502 let store = EventStore::with_config(config);
5503
5504 for i in 0..3 {
5505 let event = Event::from_strings(
5506 "test.created".to_string(),
5507 format!("entity-{i}"),
5508 "default".to_string(),
5509 serde_json::json!({"i": i}),
5510 None,
5511 )
5512 .unwrap();
5513 store.ingest(&event).unwrap();
5514 }
5515
5516 store.flush_storage().unwrap();
5517 assert_eq!(store.stats().total_events, 3);
5518 }
5519
5520 let parquet_files = find_parquet_files(&storage_dir);
5523 assert!(!parquet_files.is_empty(), "Parquet file must exist");
5524
5525 std::fs::write(&parquet_files[0], b"corrupted data").unwrap();
5527
5528 for entry in std::fs::read_dir(&wal_dir).unwrap().flatten() {
5530 std::fs::write(entry.path(), b"").unwrap();
5531 }
5532
5533 {
5540 let config = EventStoreConfig::production(
5541 &storage_dir,
5542 &wal_dir,
5543 SnapshotConfig::default(),
5544 WALConfig::default(),
5545 CompactionConfig::default(),
5546 );
5547 let store = EventStore::with_config(config);
5548
5549 assert_eq!(store.stats().total_events, 0);
5552 }
5553 }
5554
5555 fn count_wal_entries(wal_dir: &std::path::Path) -> usize {
5565 use std::io::{BufRead, BufReader};
5566 let mut total = 0usize;
5567 let Ok(entries) = std::fs::read_dir(wal_dir) else {
5568 return 0;
5569 };
5570 for entry in entries.flatten() {
5571 let path = entry.path();
5572 if path.extension().is_none_or(|e| e != "log") {
5573 continue;
5574 }
5575 let Ok(file) = std::fs::File::open(&path) else {
5576 continue;
5577 };
5578 for line in BufReader::new(file)
5579 .lines()
5580 .map_while(std::result::Result::ok)
5581 {
5582 if !line.trim().is_empty() {
5583 total += 1;
5584 }
5585 }
5586 }
5587 total
5588 }
5589
5590 #[test]
5591 fn test_checkpoint_truncates_wal_after_flush() {
5592 let data_dir = TempDir::new().unwrap();
5597 let storage_dir = data_dir.path().join("storage");
5598 let wal_dir = data_dir.path().join("wal");
5599
5600 let config = EventStoreConfig::production(
5601 &storage_dir,
5602 &wal_dir,
5603 SnapshotConfig::default(),
5604 WALConfig {
5605 sync_on_write: true,
5606 ..WALConfig::default()
5607 },
5608 CompactionConfig::default(),
5609 );
5610 let store = EventStore::with_config(config);
5611
5612 for i in 0..10 {
5613 let event = Event::from_strings(
5614 "test.created".to_string(),
5615 format!("entity-{i}"),
5616 "default".to_string(),
5617 serde_json::json!({"i": i}),
5618 None,
5619 )
5620 .unwrap();
5621 store.ingest(&event).unwrap();
5622 }
5623
5624 assert_eq!(
5626 count_wal_entries(&wal_dir),
5627 10,
5628 "WAL should have 10 events before checkpoint"
5629 );
5630
5631 store.checkpoint().unwrap();
5632
5633 assert_eq!(
5634 count_wal_entries(&wal_dir),
5635 0,
5636 "WAL should be empty after successful checkpoint"
5637 );
5638 let parquet_files = find_parquet_files(&storage_dir);
5639 assert!(!parquet_files.is_empty(), "Parquet should hold the events");
5640 }
5641
5642 #[test]
5643 fn test_replay_only_post_checkpoint_events_after_crash() {
5644 let data_dir = TempDir::new().unwrap();
5651 let storage_dir = data_dir.path().join("storage");
5652 let wal_dir = data_dir.path().join("wal");
5653
5654 let config_factory = || {
5655 EventStoreConfig::production(
5656 &storage_dir,
5657 &wal_dir,
5658 SnapshotConfig::default(),
5659 WALConfig {
5660 sync_on_write: true,
5661 ..WALConfig::default()
5662 },
5663 CompactionConfig::default(),
5664 )
5665 };
5666
5667 const N: usize = 50;
5670 const K: usize = 5;
5671 {
5672 let store = EventStore::with_config(config_factory());
5673 for i in 0..N {
5674 store
5675 .ingest(
5676 &Event::from_strings(
5677 "pre.checkpoint".to_string(),
5678 format!("e-{i}"),
5679 "default".to_string(),
5680 serde_json::json!({"i": i}),
5681 None,
5682 )
5683 .unwrap(),
5684 )
5685 .unwrap();
5686 }
5687 store.checkpoint().unwrap();
5688 assert_eq!(
5689 count_wal_entries(&wal_dir),
5690 0,
5691 "WAL should be empty immediately after checkpoint"
5692 );
5693
5694 for i in 0..K {
5695 store
5696 .ingest(
5697 &Event::from_strings(
5698 "post.checkpoint".to_string(),
5699 format!("p-{i}"),
5700 "default".to_string(),
5701 serde_json::json!({"i": i}),
5702 None,
5703 )
5704 .unwrap(),
5705 )
5706 .unwrap();
5707 }
5708 assert_eq!(
5709 count_wal_entries(&wal_dir),
5710 K,
5711 "WAL should hold only post-checkpoint events"
5712 );
5713 }
5715
5716 {
5720 let store = EventStore::with_config(config_factory());
5721 assert_eq!(
5725 store.stats().total_events,
5726 K,
5727 "Boot should replay exactly K events from WAL (the post-checkpoint window), not N+K"
5728 );
5729
5730 store.ensure_tenant_loaded("default").unwrap();
5732 assert_eq!(
5733 store.stats().total_events,
5734 N + K,
5735 "After lazy-load, both pre- and post-checkpoint events should be reachable"
5736 );
5737 }
5738 }
5739
5740 #[test]
5741 fn test_checkpoint_is_idempotent() {
5742 let data_dir = TempDir::new().unwrap();
5745 let storage_dir = data_dir.path().join("storage");
5746 let wal_dir = data_dir.path().join("wal");
5747
5748 let store = EventStore::with_config(EventStoreConfig::production(
5749 &storage_dir,
5750 &wal_dir,
5751 SnapshotConfig::default(),
5752 WALConfig::default(),
5753 CompactionConfig::default(),
5754 ));
5755
5756 for i in 0..5 {
5757 store
5758 .ingest(
5759 &Event::from_strings(
5760 "x".to_string(),
5761 format!("e-{i}"),
5762 "default".to_string(),
5763 serde_json::json!({}),
5764 None,
5765 )
5766 .unwrap(),
5767 )
5768 .unwrap();
5769 }
5770
5771 store.checkpoint().unwrap();
5772 store.checkpoint().unwrap();
5774 assert_eq!(count_wal_entries(&wal_dir), 0);
5775 }
5776
5777 #[test]
5778 fn test_checkpoint_noop_in_memory_only_mode() {
5779 let store = EventStore::new();
5781 store.checkpoint().unwrap();
5782 }
5783
5784 #[test]
5785 fn test_checkpoint_interval_from_env_defaults_to_60s_when_wal_enabled() {
5786 let (config, _) = EventStoreConfig::from_env_vars(
5787 Some("/app/data".to_string()),
5788 None,
5789 None,
5790 None,
5791 None,
5792 None,
5793 None,
5794 None,
5795 );
5796 assert_eq!(config.checkpoint_interval_secs, Some(60));
5797 }
5798
5799 #[test]
5800 fn test_checkpoint_interval_from_env_overrides_default() {
5801 let (config, _) = EventStoreConfig::from_env_vars(
5802 Some("/app/data".to_string()),
5803 None,
5804 None,
5805 None,
5806 None,
5807 None,
5808 None,
5809 Some("15".to_string()),
5810 );
5811 assert_eq!(config.checkpoint_interval_secs, Some(15));
5812 }
5813
5814 #[test]
5815 fn test_checkpoint_interval_disabled_when_wal_disabled() {
5816 let (config, _) = EventStoreConfig::from_env_vars(
5818 Some("/app/data".to_string()),
5819 None,
5820 None,
5821 Some("false".to_string()),
5822 None,
5823 None,
5824 None,
5825 Some("15".to_string()),
5826 );
5827 assert_eq!(config.checkpoint_interval_secs, None);
5828 }
5829
5830 #[test]
5831 fn test_checkpoint_interval_unparseable_falls_back_to_default() {
5832 let (config, _) = EventStoreConfig::from_env_vars(
5833 Some("/app/data".to_string()),
5834 None,
5835 None,
5836 None,
5837 None,
5838 None,
5839 None,
5840 Some("not-a-number".to_string()),
5841 );
5842 assert_eq!(config.checkpoint_interval_secs, Some(60));
5843 }
5844
5845 fn seed_two_tenants() -> EventStore {
5852 let store = EventStore::new();
5853
5854 for (entity, etype, payload) in [
5856 (
5857 "a-1",
5858 "created",
5859 serde_json::json!({"colour": "red", "size": 1}),
5860 ),
5861 ("a-1", "updated", serde_json::json!({"colour": "blue"})),
5862 ("a-2", "created", serde_json::json!({"colour": "green"})),
5863 ] {
5864 store
5865 .ingest(
5866 &Event::from_strings(
5867 etype.to_string(),
5868 entity.to_string(),
5869 "alice".to_string(),
5870 payload,
5871 None,
5872 )
5873 .unwrap(),
5874 )
5875 .unwrap();
5876 }
5877
5878 store
5880 .ingest(
5881 &Event::from_strings(
5882 "created".to_string(),
5883 "a-1".to_string(),
5884 "bob".to_string(),
5885 serde_json::json!({"colour": "BOB_SECRET", "bob_only": true}),
5886 None,
5887 )
5888 .unwrap(),
5889 )
5890 .unwrap();
5891
5892 store
5893 }
5894
5895 #[test]
5896 fn test_stats_for_tenant_counts_only_that_tenant() {
5897 let store = seed_two_tenants();
5898
5899 let alice = store.stats_for_tenant("alice");
5900 assert_eq!(alice.total_events, 3);
5901 assert_eq!(alice.total_entities, 2);
5902 assert_eq!(alice.total_event_types, 2);
5903 assert_eq!(alice.event_types.get("created"), Some(&2));
5904 assert_eq!(alice.event_types.get("updated"), Some(&1));
5905
5906 let bob = store.stats_for_tenant("bob");
5907 assert_eq!(bob.total_events, 1);
5908 assert_eq!(bob.total_entities, 1);
5909 assert_eq!(bob.total_event_types, 1);
5910 }
5911
5912 #[test]
5913 fn test_stats_for_tenant_never_reports_global_totals() {
5914 let store = seed_two_tenants();
5915
5916 assert_eq!(store.stats().total_events, 4);
5918 assert_eq!(store.stats_for_tenant("alice").total_events, 3);
5919 assert_eq!(store.stats_for_tenant("bob").total_events, 1);
5920
5921 assert_eq!(store.stats_for_tenant("bob").total_ingested, 1);
5924 }
5925
5926 #[test]
5927 fn test_stats_for_tenant_unknown_tenant_is_empty_not_global() {
5928 let store = seed_two_tenants();
5929 let nobody = store.stats_for_tenant("does-not-exist");
5930
5931 assert_eq!(nobody.total_events, 0);
5932 assert_eq!(nobody.total_entities, 0);
5933 assert!(nobody.event_types.is_empty());
5934 assert!(nobody.oldest_event.is_none());
5935 assert!(nobody.newest_event.is_none());
5936 }
5937
5938 #[test]
5939 fn test_stats_for_tenant_reports_time_range() {
5940 let store = seed_two_tenants();
5941 let alice = store.stats_for_tenant("alice");
5942
5943 let oldest = alice.oldest_event.expect("oldest");
5944 let newest = alice.newest_event.expect("newest");
5945 assert!(oldest <= newest);
5946 }
5947
5948 #[test]
5949 fn test_reconstruct_state_for_tenant_isolates_shared_entity_id() {
5950 let store = seed_two_tenants();
5951
5952 let alice = store
5954 .reconstruct_state_for_tenant("a-1", None, "alice")
5955 .unwrap();
5956 let alice_state = alice.get("current_state").unwrap();
5957 assert_eq!(alice_state.get("colour").unwrap(), "blue"); assert_eq!(alice_state.get("size").unwrap(), 1);
5959 assert!(
5960 alice_state.get("bob_only").is_none(),
5961 "alice must not see bob's payload keys: {alice_state:?}"
5962 );
5963 assert_eq!(alice.get("event_count").unwrap(), 2);
5964
5965 let bob = store
5966 .reconstruct_state_for_tenant("a-1", None, "bob")
5967 .unwrap();
5968 let bob_state = bob.get("current_state").unwrap();
5969 assert_eq!(bob_state.get("colour").unwrap(), "BOB_SECRET");
5970 assert_eq!(bob.get("event_count").unwrap(), 1);
5971 }
5972
5973 #[test]
5974 fn test_reconstruct_state_for_tenant_rejects_another_tenants_entity() {
5975 let store = seed_two_tenants();
5976
5977 assert!(
5979 store
5980 .reconstruct_state_for_tenant("a-2", None, "alice")
5981 .is_ok()
5982 );
5983 assert!(
5984 store
5985 .reconstruct_state_for_tenant("a-2", None, "bob")
5986 .is_err(),
5987 "bob must not be able to read alice's entity"
5988 );
5989 }
5990
5991 #[test]
5992 fn test_global_reconstruct_state_still_spans_tenants() {
5993 let store = seed_two_tenants();
5996 let all = store.reconstruct_state("a-1", None).unwrap();
5997 assert_eq!(all.get("event_count").unwrap(), 3);
5998 }
5999
6000 fn seed_hot_entity(history: usize) -> EventStore {
6011 let store = EventStore::new();
6012 for _ in 0..history {
6013 store
6014 .ingest(&create_test_event("entity-hot", "user.updated"))
6015 .unwrap();
6016 }
6017 store
6018 }
6019
6020 fn hot_request(limit: Option<usize>) -> QueryEventsRequest {
6021 QueryEventsRequest {
6022 entity_id: Some("entity-hot".to_string()),
6023 limit,
6024 ..QueryEventsRequest::default()
6025 }
6026 }
6027
6028 #[test]
6029 fn query_window_limit_bounds_materialization_not_just_the_response() {
6030 const HISTORY: usize = 400;
6031 let store = seed_hot_entity(HISTORY);
6032
6033 let ((page, total), materialized) = crate::clone_probe::measure(|| {
6035 store.query_window(&hot_request(Some(1)), 0, true).unwrap()
6036 });
6037 assert_eq!(page.len(), 1, "limit=1 returns one event");
6038 assert_eq!(total, HISTORY, "total is still the full match count");
6039 assert_eq!(
6040 materialized, 1,
6041 "limit=1 cloned {materialized} of {HISTORY} events: `limit` must \
6042 bound what the store materializes, not just what it returns"
6043 );
6044
6045 let ((page, _), materialized) = crate::clone_probe::measure(|| {
6047 store
6048 .query_window(&hot_request(Some(5)), 300, false)
6049 .unwrap()
6050 });
6051 assert_eq!(page.len(), 5);
6052 assert_eq!(
6053 materialized, 5,
6054 "offset=300&limit=5 cloned {materialized} events, expected 5"
6055 );
6056
6057 let ((all, _), materialized) = crate::clone_probe::measure(|| {
6060 store.query_window(&hot_request(None), 0, true).unwrap()
6061 });
6062 assert_eq!(all.len(), HISTORY);
6063 assert_eq!(materialized, HISTORY as u64);
6064 }
6065
6066 #[test]
6067 fn query_window_cost_does_not_grow_with_entity_history() {
6068 let short = seed_hot_entity(40);
6073 let long = seed_hot_entity(400);
6074
6075 let (_, short_cost) = crate::clone_probe::measure(|| {
6076 short.query_window(&hot_request(Some(1)), 0, true).unwrap()
6077 });
6078 let (_, long_cost) = crate::clone_probe::measure(|| {
6079 long.query_window(&hot_request(Some(1)), 0, true).unwrap()
6080 });
6081
6082 assert_eq!(
6083 (short_cost, long_cost),
6084 (1, 1),
6085 "a limit=1 page materialized {short_cost} events over a 40-event \
6086 history and {long_cost} over 400 — cost is tracking history length"
6087 );
6088 }
6089
6090 #[test]
6091 fn select_window_orders_only_the_window() {
6092 use std::cell::Cell;
6098 const N: usize = 4096;
6099
6100 let comparisons = Cell::new(0usize);
6101 let order = |a: &u64, b: &u64| {
6102 comparisons.set(comparisons.get() + 1);
6103 a.cmp(b)
6104 };
6105 let shuffled = || -> Vec<u64> {
6106 (0..N as u64)
6107 .map(|i| (i * 2_654_435_761) % 1_000_003)
6108 .collect()
6109 };
6110
6111 let mut bounded = shuffled();
6112 select_window(&mut bounded, 0, Some(1), order);
6113 let bounded_cost = comparisons.replace(0);
6114 assert_eq!(bounded.len(), 1, "window of 1 keeps 1 item");
6115
6116 let mut everything = shuffled();
6117 select_window(&mut everything, 0, None, order);
6118 let full_sort_cost = comparisons.get();
6119 assert_eq!(everything.len(), N);
6120
6121 assert!(
6122 bounded_cost < 4 * N,
6123 "selecting a 1-item window out of {N} took {bounded_cost} comparisons \
6124 (~{}·N) — that is sort-shaped, not selection-shaped",
6125 bounded_cost / N
6126 );
6127 assert!(
6128 bounded_cost * 3 < full_sort_cost,
6129 "a 1-item window cost {bounded_cost} comparisons against \
6130 {full_sort_cost} for sorting all {N}: the window is not bounding \
6131 the ordering work"
6132 );
6133 }
6134
6135 #[test]
6136 fn select_window_is_equivalent_to_sort_then_window() {
6137 let order = |a: &(u32, usize), b: &(u32, usize)| a.cmp(b);
6141 let source: Vec<(u32, usize)> = [7, 3, 3, 9, 1, 3, 5, 9, 0, 2]
6142 .into_iter()
6143 .enumerate()
6144 .map(|(i, k)| (k, i))
6145 .collect();
6146
6147 let mut sorted = source.clone();
6148 sorted.sort_by(order);
6149
6150 for offset in 0..12 {
6151 for limit in [None, Some(0), Some(1), Some(3), Some(10), Some(50)] {
6152 let mut got = source.clone();
6153 select_window(&mut got, offset, limit, order);
6154 let got: Vec<_> = got.into_iter().skip(offset).collect();
6155 let expected: Vec<_> = sorted
6156 .iter()
6157 .copied()
6158 .skip(offset)
6159 .take(limit.unwrap_or(usize::MAX))
6160 .collect();
6161 assert_eq!(got, expected, "offset={offset} limit={limit:?}");
6162 }
6163 }
6164 }
6165
6166 #[test]
6177 fn query_window_applies_time_filters_on_the_full_scan_path() {
6178 let store = EventStore::new();
6179 let base = Utc::now() - chrono::Duration::hours(24);
6180 let mut ids = Vec::new();
6181 for i in 0..5i64 {
6182 let mut event = create_test_event(&format!("e-{i}"), "user.created");
6183 event.timestamp = base + chrono::Duration::hours(i);
6184 event.version = i + 1;
6185 ids.push(event.id);
6186 store.ingest(&event).unwrap();
6187 }
6188 let at = |h: i64| base + chrono::Duration::hours(h);
6189
6190 let scoped = |mutate: &dyn Fn(&mut QueryEventsRequest)| {
6193 let mut req = QueryEventsRequest {
6194 tenant_id: Some("default".to_string()),
6195 ..QueryEventsRequest::default()
6196 };
6197 mutate(&mut req);
6198 req
6199 };
6200
6201 for (label, req, expected) in [
6202 (
6203 "since=T+2 keeps only events at or after T+2",
6204 scoped(&|r| r.since = Some(at(2))),
6205 vec![ids[2], ids[3], ids[4]],
6206 ),
6207 (
6208 "until=T+1 keeps only events at or before T+1",
6209 scoped(&|r| r.until = Some(at(1))),
6210 vec![ids[0], ids[1]],
6211 ),
6212 (
6213 "as_of=T+1 is time travel: nothing newer than T+1",
6214 scoped(&|r| r.as_of = Some(at(1))),
6215 vec![ids[0], ids[1]],
6216 ),
6217 (
6218 "since+until compose into a closed window",
6219 scoped(&|r| {
6220 r.since = Some(at(1));
6221 r.until = Some(at(3));
6222 }),
6223 vec![ids[1], ids[2], ids[3]],
6224 ),
6225 ] {
6226 let (events, total) = store.query_window(&req, 0, false).unwrap();
6227 let got: Vec<_> = events.iter().map(|e| e.id).collect();
6228 assert_eq!(got, expected, "{label}");
6229 assert_eq!(
6230 total,
6231 expected.len(),
6232 "{label}: total counts the events INSIDE the window — \
6233 has_more is derived from it, so a full-history total makes a \
6234 paginator walk events the caller filtered out"
6235 );
6236 }
6237
6238 let (indexed, total) = store
6241 .query_window(
6242 &scoped(&|r| {
6243 r.entity_id = Some("e-3".to_string());
6244 r.since = Some(at(2));
6245 }),
6246 0,
6247 false,
6248 )
6249 .unwrap();
6250 assert_eq!(
6251 indexed.iter().map(|e| e.id).collect::<Vec<_>>(),
6252 vec![ids[3]]
6253 );
6254 assert_eq!(total, 1);
6255
6256 let (page, total) = store
6259 .query_window(&scoped(&|r| r.since = Some(at(2))), 1, true)
6260 .unwrap();
6261 assert_eq!(
6262 page.iter().map(|e| e.id).collect::<Vec<_>>(),
6263 vec![ids[3], ids[2]],
6264 "order=desc + offset=1 inside a since window"
6265 );
6266 assert_eq!(total, 3);
6267 }
6268
6269 #[test]
6270 fn query_window_bounded_selection_matches_a_full_sort() {
6271 let store = EventStore::new();
6276 for i in 0..50 {
6277 store
6278 .ingest(&create_test_event(&format!("e-{i:02}"), "user.created"))
6279 .unwrap();
6280 }
6281 let all = |descending: bool| {
6282 let (events, _) = store
6283 .query_window(&QueryEventsRequest::default(), 0, descending)
6284 .unwrap();
6285 events
6286 };
6287
6288 for descending in [false, true] {
6289 let reference = all(descending);
6290 for offset in [0, 1, 7, 49, 50, 100] {
6291 for limit in [1, 3, 10, 50, 100] {
6292 let (page, total) = store
6293 .query_window(
6294 &QueryEventsRequest {
6295 limit: Some(limit),
6296 ..QueryEventsRequest::default()
6297 },
6298 offset,
6299 descending,
6300 )
6301 .unwrap();
6302 let expected: Vec<_> = reference
6303 .iter()
6304 .skip(offset)
6305 .take(limit)
6306 .map(|e| e.id)
6307 .collect();
6308 let got: Vec<_> = page.iter().map(|e| e.id).collect();
6309 assert_eq!(
6310 got, expected,
6311 "desc={descending} offset={offset} limit={limit}: windowed \
6312 selection must match a full sort"
6313 );
6314 assert_eq!(total, 50, "total is always the full match count");
6315 }
6316 }
6317 }
6318 }
6319}