Skip to main content

allsource_core/
store.rs

1#[cfg(feature = "server")]
2use crate::application::services::webhook::WebhookRegistry;
3#[cfg(feature = "server")]
4use crate::infrastructure::observability::metrics::MetricsRegistry;
5#[cfg(feature = "server")]
6use crate::infrastructure::web::websocket::WebSocketManager;
7use crate::{
8    application::{
9        dto::QueryEventsRequest,
10        services::{
11            consumer::ConsumerRegistry,
12            exactly_once::{ExactlyOnceConfig, ExactlyOnceRegistry},
13            pipeline::PipelineManager,
14            projection::{EntitySnapshotProjection, EventCounterProjection, ProjectionManager},
15            replay::ReplayManager,
16            schema::{SchemaRegistry, SchemaRegistryConfig},
17            schema_evolution::SchemaEvolutionManager,
18        },
19    },
20    domain::entities::Event,
21    error::{AllSourceError, Result},
22    infrastructure::{
23        persistence::{
24            compaction::{CompactionConfig, CompactionManager},
25            index::{EventIndex, IndexEntry},
26            snapshot::{SnapshotConfig, SnapshotManager, SnapshotType},
27            storage::ParquetStorage,
28            tenant_loader::TenantLoader,
29            wal::{WALConfig, WriteAheadLog},
30        },
31        query::geospatial::GeoIndex,
32    },
33};
34use chrono::{DateTime, Utc};
35use dashmap::DashMap;
36use parking_lot::RwLock;
37use std::{path::PathBuf, sync::Arc};
38#[cfg(feature = "server")]
39use tokio::sync::mpsc;
40
41/// High-performance event store with columnar storage
42/// What a read credential is allowed to see, enforced inside the store.
43///
44/// A read key is otherwise all-or-nothing: "let this tool read the store" and
45/// "let this tool read every credential ever issued" are the same grant. That is
46/// not theoretical — auth events carry session tokens whose value IS the
47/// `entity_id`, so a routine query rendered live bearer tokens into an agent's
48/// context (#265). Redacting on write cannot fix it: the token is the lookup
49/// key.
50///
51/// Deliberately **not** `Deserialize`. This must never be settable from a
52/// request body, or a caller widens its own scope by omitting the field.
53///
54/// Allow-list, not deny-list, on purpose: a stream family added later is
55/// excluded until someone names it, so the failure mode of forgetting is a
56/// missing read rather than a leak.
57#[derive(Debug, Clone, Default)]
58pub struct ReadScope {
59    allow_entity_prefixes: Option<Vec<String>>,
60}
61
62impl ReadScope {
63    /// Everything in the tenant — the behaviour of every read before scopes.
64    pub fn unrestricted() -> Self {
65        Self {
66            allow_entity_prefixes: None,
67        }
68    }
69
70    /// Only entities whose id starts with one of `prefixes`.
71    ///
72    /// An empty list denies everything, which is the safe reading of "scoped to
73    /// nothing" and stops an accidentally-empty config granting full access.
74    pub fn allow_entity_prefixes<I, S>(prefixes: I) -> Self
75    where
76        I: IntoIterator<Item = S>,
77        S: Into<String>,
78    {
79        Self {
80            allow_entity_prefixes: Some(prefixes.into_iter().map(Into::into).collect()),
81        }
82    }
83
84    /// Whether this scope can see `entity_id`.
85    pub fn permits(&self, entity_id: &str) -> bool {
86        match &self.allow_entity_prefixes {
87            None => true,
88            Some(allowed) => allowed.iter().any(|p| entity_id.starts_with(p.as_str())),
89        }
90    }
91
92    /// True when this scope restricts nothing.
93    pub fn is_unrestricted(&self) -> bool {
94        self.allow_entity_prefixes.is_none()
95    }
96}
97
98pub struct EventStore {
99    /// In-memory event storage
100    events: Arc<RwLock<Vec<Event>>>,
101
102    /// High-performance concurrent index
103    index: Arc<EventIndex>,
104
105    /// Projection manager for real-time aggregations
106    pub(crate) projections: Arc<RwLock<ProjectionManager>>,
107
108    /// Optional persistent storage (v0.2 feature)
109    storage: Option<Arc<RwLock<ParquetStorage>>>,
110
111    /// WebSocket manager for real-time event streaming (v0.2 feature)
112    #[cfg(feature = "server")]
113    websocket_manager: Arc<WebSocketManager>,
114
115    /// Snapshot manager for fast state recovery (v0.2 feature)
116    snapshot_manager: Arc<SnapshotManager>,
117
118    /// Write-Ahead Log for durability (v0.2 feature)
119    wal: Option<Arc<WriteAheadLog>>,
120
121    /// Compaction manager for Parquet optimization (v0.2 feature)
122    compaction_manager: Option<Arc<CompactionManager>>,
123
124    /// Schema registry for event validation (v0.5 feature)
125    schema_registry: Arc<SchemaRegistry>,
126
127    /// Replay manager for event replay and projection rebuilding (v0.5 feature)
128    replay_manager: Arc<ReplayManager>,
129
130    /// Pipeline manager for stream processing (v0.5 feature)
131    pipeline_manager: Arc<PipelineManager>,
132
133    /// Prometheus metrics registry (v0.6 feature)
134    #[cfg(feature = "server")]
135    metrics: Arc<MetricsRegistry>,
136
137    /// Total events ingested (for metrics)
138    total_ingested: Arc<RwLock<u64>>,
139
140    /// Projection state cache for Query Service integration (v0.7 feature)
141    /// Key format: "{projection_name}:{entity_id}"
142    /// This DashMap provides O(1) access with ~11.9 μs latency
143    projection_state_cache: Arc<DashMap<String, serde_json::Value>>,
144
145    /// Projection status overrides (v0.13 feature)
146    /// Tracks pause/start state per projection name: "running" or "paused"
147    projection_status: Arc<DashMap<String, String>>,
148
149    /// Webhook registry for outbound event delivery (v0.11 feature)
150    #[cfg(feature = "server")]
151    webhook_registry: Arc<WebhookRegistry>,
152
153    /// Channel sender for async webhook delivery tasks
154    #[cfg(feature = "server")]
155    webhook_tx: Arc<RwLock<Option<mpsc::UnboundedSender<WebhookDeliveryTask>>>>,
156
157    /// Geospatial index for coordinate-based queries (v2.0 feature)
158    geo_index: Arc<GeoIndex>,
159
160    /// Exactly-once processing registry (v2.0 feature)
161    exactly_once: Arc<ExactlyOnceRegistry>,
162
163    /// Autonomous schema evolution manager (v2.0 feature)
164    schema_evolution: Arc<SchemaEvolutionManager>,
165
166    /// Per-entity version counters for optimistic concurrency control (v0.14 feature)
167    /// Key: entity_id string, Value: monotonic version (number of events for that entity)
168    entity_versions: Arc<DashMap<String, u64>>,
169
170    /// Durable consumer registry for subscription cursor tracking (v0.14 feature)
171    consumer_registry: Arc<ConsumerRegistry>,
172
173    /// In-process broadcast of every successfully-ingested event. Always enabled
174    /// so embedded consumers (TUI, web) can tail changes without the `server`
175    /// feature / HTTP stack. Lagging receivers see `RecvError::Lagged`.
176    event_broadcast_tx: tokio::sync::broadcast::Sender<Arc<Event>>,
177
178    /// Per-tenant lazy-load bookkeeping. Tracks which tenants have
179    /// been hydrated from Parquet into the in-memory pile and
180    /// serializes concurrent first-loads of the same tenant. See
181    /// `ensure_tenant_loaded` and Step 2 of the sustainable data
182    /// strategy.
183    tenant_loader: Arc<TenantLoader>,
184
185    /// Cadence of the runtime checkpoint loop (Step 6). `None` means
186    /// the loop is disabled — WAL grows until boot. Production reads
187    /// this from `ALLSOURCE_CHECKPOINT_INTERVAL_SECONDS`. Stored on
188    /// the store so background tasks can read it without re-parsing
189    /// the env.
190    checkpoint_interval_secs: Option<u64>,
191
192    /// Read-only (replica) mode. Set when a second process attaches to a
193    /// data-dir already owned by a live writer (see Prime's data-dir lock).
194    /// A read-only store replays the WAL + Parquet into memory at boot so it
195    /// can serve reads, but it MUST NOT truncate the WAL (that would unlink
196    /// the inode the owner is still appending to — issue #201) and rejects
197    /// all writes with `AllSourceError::ReadOnly`.
198    read_only: bool,
199
200    /// Held shared from an event's WAL append until it is in the Parquet
201    /// batch, and exclusively while `checkpoint` seals the WAL. Without it an
202    /// event can sit in a sealed segment yet miss the flush that retires it.
203    durability_gate: RwLock<()>,
204}
205
206/// A task queued for async webhook delivery
207#[cfg(feature = "server")]
208#[derive(Debug, Clone)]
209pub struct WebhookDeliveryTask {
210    pub webhook: crate::application::services::webhook::WebhookSubscription,
211    pub event: Event,
212}
213
214/// Reduce `items` to the `offset .. offset + limit` window under `order`,
215/// leaving that window sorted and everything past it dropped.
216///
217/// `limit == None` means "no window": the whole slice is sorted. Otherwise the
218/// window is partitioned off with `select_nth_unstable_by` — O(N) — and only
219/// the `offset + limit` retained items are sorted, O(k log k), instead of
220/// sorting all N (issue #251). `order` MUST be a total order (no two distinct
221/// items compare `Equal`); with one, this is observationally identical to
222/// sorting everything and then windowing, ties included.
223///
224/// The first `offset` items are kept, not skipped — the caller drops them,
225/// which is what makes `has_more`/`total` independent of the window.
226fn select_window<T>(
227    items: &mut Vec<T>,
228    offset: usize,
229    limit: Option<usize>,
230    order: impl Fn(&T, &T) -> std::cmp::Ordering + Copy,
231) {
232    match limit.map(|limit| offset.saturating_add(limit)) {
233        Some(0) => {
234            items.clear();
235            return;
236        }
237        Some(window_end) if window_end < items.len() => {
238            items.select_nth_unstable_by(window_end - 1, order);
239            items.truncate(window_end);
240        }
241        // No limit, or a window that already covers everything: nothing to
242        // partition off, the sort below orders the whole set.
243        _ => {}
244    }
245    items.sort_unstable_by(order);
246}
247
248impl EventStore {
249    /// Create a new in-memory event store
250    pub fn new() -> Self {
251        Self::with_config(EventStoreConfig::default())
252    }
253
254    /// Create event store with custom configuration
255    pub fn with_config(config: EventStoreConfig) -> Self {
256        let mut projections = ProjectionManager::new();
257
258        // Register built-in projections
259        projections.register(Arc::new(EntitySnapshotProjection::new("entity_snapshots")));
260        projections.register(Arc::new(EventCounterProjection::new("event_counters")));
261
262        // Initialize persistent storage if configured
263        let storage = config
264            .storage_dir
265            .as_ref()
266            .and_then(|dir| match ParquetStorage::new(dir) {
267                Ok(storage) => {
268                    tracing::info!("✅ Parquet persistence enabled at: {}", dir.display());
269                    Some(Arc::new(RwLock::new(storage)))
270                }
271                Err(e) => {
272                    tracing::error!("❌ Failed to initialize Parquet storage: {}", e);
273                    None
274                }
275            });
276
277        // Initialize WAL if configured (v0.2 feature)
278        let wal = config.wal_dir.as_ref().and_then(|dir| {
279            match WriteAheadLog::new(dir, config.wal_config.clone()) {
280                Ok(wal) => {
281                    tracing::info!("✅ WAL enabled at: {}", dir.display());
282                    Some(Arc::new(wal))
283                }
284                Err(e) => {
285                    tracing::error!("❌ Failed to initialize WAL: {}", e);
286                    None
287                }
288            }
289        });
290
291        // Initialize compaction manager if Parquet storage is enabled (v0.2 feature)
292        let compaction_manager = config.storage_dir.as_ref().map(|dir| {
293            let manager = CompactionManager::new(dir, config.compaction_config.clone());
294            Arc::new(manager)
295        });
296
297        // Initialize schema registry (v0.5 feature)
298        let schema_registry = Arc::new(SchemaRegistry::new(config.schema_registry_config.clone()));
299        tracing::info!("✅ Schema registry enabled");
300
301        // Initialize replay manager (v0.5 feature)
302        let replay_manager = Arc::new(ReplayManager::new());
303        tracing::info!("✅ Replay manager enabled");
304
305        // Initialize pipeline manager (v0.5 feature)
306        let pipeline_manager = Arc::new(PipelineManager::new());
307        tracing::info!("✅ Pipeline manager enabled");
308
309        // Initialize metrics registry (v0.6 feature)
310        #[cfg(feature = "server")]
311        let metrics = {
312            let m = MetricsRegistry::new();
313            tracing::info!("✅ Prometheus metrics registry initialized");
314            m
315        };
316
317        // Initialize projection state cache (v0.7 feature)
318        let projection_state_cache = Arc::new(DashMap::new());
319        tracing::info!("✅ Projection state cache initialized");
320
321        // Initialize webhook registry (v0.11 feature)
322        #[cfg(feature = "server")]
323        let webhook_registry = {
324            let w = Arc::new(WebhookRegistry::new());
325            tracing::info!("✅ Webhook registry initialized");
326            w
327        };
328
329        // Unconditional in-process event broadcaster so embedded consumers
330        // (TUI, web) can live-reload without the `server` feature.
331        let (event_broadcast_tx, _) = tokio::sync::broadcast::channel(1024);
332
333        let store = Self {
334            events: Arc::new(RwLock::new(Vec::new())),
335            index: Arc::new(EventIndex::new()),
336            projections: Arc::new(RwLock::new(projections)),
337            storage,
338            #[cfg(feature = "server")]
339            websocket_manager: Arc::new(WebSocketManager::new()),
340            snapshot_manager: Arc::new(SnapshotManager::new(config.snapshot_config)),
341            wal,
342            compaction_manager,
343            schema_registry,
344            replay_manager,
345            pipeline_manager,
346            #[cfg(feature = "server")]
347            metrics,
348            total_ingested: Arc::new(RwLock::new(0)),
349            projection_state_cache,
350            projection_status: Arc::new(DashMap::new()),
351            #[cfg(feature = "server")]
352            webhook_registry,
353            #[cfg(feature = "server")]
354            webhook_tx: Arc::new(RwLock::new(None)),
355            geo_index: Arc::new(GeoIndex::new()),
356            exactly_once: Arc::new(ExactlyOnceRegistry::new(ExactlyOnceConfig::default())),
357            schema_evolution: Arc::new(SchemaEvolutionManager::new()),
358            entity_versions: Arc::new(DashMap::new()),
359            consumer_registry: Arc::new(ConsumerRegistry::new()),
360            event_broadcast_tx,
361            tenant_loader: {
362                let loader = TenantLoader::new();
363                if let Some(budget) = config.cache_byte_budget {
364                    loader.set_byte_budget(budget);
365                    tracing::info!(
366                        "✅ Cache byte budget set to {} bytes ({:.2} GiB) — LRU eviction enabled",
367                        budget,
368                        budget as f64 / (1024.0 * 1024.0 * 1024.0)
369                    );
370                } else {
371                    tracing::info!(
372                        "✅ Cache budget unset — every loaded tenant stays resident \
373                         (set ALLSOURCE_CACHE_BYTES to enable eviction)"
374                    );
375                }
376                Arc::new(loader)
377            },
378            checkpoint_interval_secs: config.checkpoint_interval_secs,
379            read_only: config.read_only,
380            durability_gate: RwLock::new(()),
381        };
382
383        if config.read_only {
384            tracing::info!(
385                "📖 EventStore opened READ-ONLY (replica): WAL will be replayed for reads but \
386                 not truncated; writes are rejected"
387            );
388        }
389
390        // Boot is now O(1) regardless of dataset size (Step 2 of the
391        // sustainable data strategy). Pre-Step-2 we scanned every
392        // Parquet file at startup; on the production volume that
393        // grew past available memory and Core OOM'd during recovery
394        // (issue #160). Now:
395        //
396        //   - Parquet data stays on disk. Tenants are hydrated on
397        //     demand by `ensure_tenant_loaded`, called from the
398        //     query path on first access.
399        //   - WAL is still recovered eagerly. WAL is bounded
400        //     (rotates / truncates after each Parquet flush), so
401        //     replaying it is O(recent un-flushed writes) — small
402        //     by construction and required for correctness, since
403        //     those events aren't durably in Parquet yet.
404        //
405        // After WAL recovery we still checkpoint recovered events
406        // to Parquet and truncate the WAL. The lazy-load path's
407        // dedupe in `append_loaded_event` (index.get_by_id check)
408        // makes it safe for the same event to be reachable through
409        // both WAL recovery and a subsequent ensure_tenant_loaded
410        // pass on the same tenant.
411        if let Some(ref wal) = store.wal {
412            match wal.recover() {
413                Ok(recovered_events) if !recovered_events.is_empty() => {
414                    let mut wal_new = 0usize;
415                    for event in recovered_events {
416                        let offset = store.events.read().len();
417                        if let Err(e) = store.index.index_event(
418                            event.id,
419                            event.entity_id_str(),
420                            event.event_type_str(),
421                            event.timestamp,
422                            offset,
423                        ) {
424                            tracing::error!("Failed to re-index WAL event {}: {}", event.id, e);
425                        }
426
427                        if let Err(e) = store.projections.read().process_event(&event) {
428                            tracing::error!("Failed to re-process WAL event {}: {}", event.id, e);
429                        }
430
431                        *store
432                            .entity_versions
433                            .entry(event.entity_id_str().to_string())
434                            .or_insert(0) += 1;
435
436                        store.events.write().push(event);
437                        wal_new += 1;
438                    }
439
440                    // Step 6: expose replay size as a gauge so ops
441                    // can graph "how big was the last replay?" and
442                    // catch regressions where the checkpoint loop
443                    // stops draining the WAL.
444                    #[cfg(feature = "server")]
445                    store.metrics.wal_replay_events_total.set(wal_new as i64);
446
447                    if wal_new > 0 {
448                        let total = store.events.read().len();
449                        // total_ingested now reflects "events the
450                        // process knows about" rather than "events
451                        // ever written" — Parquet data isn't loaded
452                        // on boot, so the count grows as tenants
453                        // get hydrated. Consumers that need the
454                        // historical total should look at Parquet
455                        // file stats, not this counter.
456                        *store.total_ingested.write() = total as u64;
457                        tracing::info!(
458                            "✅ Recovered {} events from WAL (Parquet data stays cold until \
459                             first per-tenant query)",
460                            wal_new
461                        );
462
463                        // Checkpoint WAL events to Parquet — buffer
464                        // them into the per-tenant Parquet batches
465                        // first, then flush. Without this,
466                        // flush_storage() finds an empty current
467                        // batch and silently no-ops, the WAL gets
468                        // truncated, and the events exist only in
469                        // memory (lost on next restart).
470                        //
471                        // Skip entirely in read-only (replica) mode: a replica
472                        // does not own the WAL. `wal.truncate()` unlinks the WAL
473                        // file, which would delete the inode the owning writer is
474                        // still appending to and silently lose its in-flight
475                        // writes (issue #201). The replica keeps the recovered
476                        // events in memory for reads and leaves the file alone.
477                        if let Some(ref storage) = store.storage
478                            && !config.read_only
479                        {
480                            tracing::info!(
481                                "📸 Checkpointing {} WAL events to Parquet storage...",
482                                wal_new
483                            );
484                            let parquet = storage.read();
485                            let events = store.events.read();
486                            let mut buffered = 0usize;
487                            for event in events.iter().skip(events.len() - wal_new) {
488                                if let Err(e) = parquet.append_event(event.clone()) {
489                                    tracing::error!(
490                                        "Failed to buffer WAL event for Parquet: {}",
491                                        e
492                                    );
493                                } else {
494                                    buffered += 1;
495                                }
496                            }
497                            drop(events);
498                            drop(parquet);
499
500                            if buffered < wal_new {
501                                tracing::error!(
502                                    "Buffered {} of {} recovered WAL events; leaving the WAL \
503                                     in place so none are lost",
504                                    buffered,
505                                    wal_new
506                                );
507                            } else if buffered > 0 {
508                                if let Err(e) = store.flush_storage() {
509                                    tracing::error!("Failed to checkpoint to Parquet: {}", e);
510                                } else if let Err(e) = wal.truncate() {
511                                    tracing::error!(
512                                        "Failed to truncate WAL after checkpoint: {}",
513                                        e
514                                    );
515                                } else {
516                                    tracing::info!(
517                                        "✅ WAL checkpointed and truncated ({} events)",
518                                        buffered
519                                    );
520                                }
521                            }
522                        }
523                    }
524                }
525                Ok(_) => {
526                    tracing::debug!("No events to recover from WAL");
527                    #[cfg(feature = "server")]
528                    store.metrics.wal_replay_events_total.set(0);
529                }
530                Err(e) => {
531                    tracing::error!("❌ WAL recovery failed: {}", e);
532                }
533            }
534        } else if store.storage.is_some() {
535            tracing::info!(
536                "📂 Boot complete (lazy-load mode): Parquet data stays on disk until first \
537                 per-tenant query"
538            );
539        }
540
541        store
542    }
543
544    /// Whether this store was opened read-only (replica mode).
545    pub fn is_read_only(&self) -> bool {
546        self.read_only
547    }
548
549    /// Return `Err(AllSourceError::ReadOnly)` if the store is a read-only
550    /// replica. Called at the top of every write path so a replica never
551    /// appends to a WAL it does not own.
552    fn ensure_writable(&self) -> Result<()> {
553        if self.read_only {
554            return Err(crate::error::AllSourceError::ReadOnly(
555                "this AllSource instance is a read-only replica — the data directory is owned by \
556                 another running process. Stop the other process, or run a single shared writer \
557                 (e.g. Prime in --mode http) and point clients at it."
558                    .to_string(),
559            ));
560        }
561        Ok(())
562    }
563
564    /// Ingest a new event with optional optimistic concurrency check.
565    ///
566    /// If `expected_version` is `Some(v)`, the write is rejected with
567    /// `VersionConflict` unless the entity's current version equals `v`.
568    /// The version check and WAL append are atomic (locked together).
569    ///
570    /// Returns the new entity version after the append.
571    #[cfg_attr(feature = "hotpath", hotpath::measure)]
572    pub fn ingest_with_expected_version(
573        &self,
574        event: &Event,
575        expected_version: Option<u64>,
576    ) -> Result<u64> {
577        // Reject writes in read-only (replica) mode before touching the WAL.
578        self.ensure_writable()?;
579
580        // Validate event first (before any locking)
581        self.validate_event(event)?;
582
583        let entity_id = event.entity_id_str().to_string();
584        let _durable = self.durability_gate.read();
585
586        // Atomic version check + append: hold the DashMap entry lock
587        // to prevent TOCTOU races between check and write.
588        let new_version = {
589            let mut version_entry = self.entity_versions.entry(entity_id.clone()).or_insert(0);
590            let current = *version_entry;
591
592            if let Some(expected) = expected_version
593                && current != expected
594            {
595                return Err(crate::error::AllSourceError::VersionConflict { expected, current });
596            }
597
598            // Write to WAL FIRST for durability (under version lock to keep atomicity)
599            if let Some(ref wal) = self.wal {
600                wal.append(event.clone())?;
601            }
602
603            *version_entry += 1;
604            *version_entry
605        };
606
607        // From here on, the event is durable (WAL) and version is bumped.
608        // Continue with indexing, projections, storage, and broadcast.
609        self.ingest_post_wal(event)?;
610
611        Ok(new_version)
612    }
613
614    /// Post-WAL ingestion: index, projections, storage, broadcast.
615    /// Called after WAL append and version bump are complete.
616    #[cfg_attr(feature = "hotpath", hotpath::measure)]
617    fn ingest_post_wal(&self, event: &Event) -> Result<()> {
618        #[cfg(feature = "server")]
619        let timer = self.metrics.ingestion_duration_seconds.start_timer();
620
621        let mut events = self.events.write();
622        let offset = events.len();
623
624        // Index the event
625        self.index.index_event(
626            event.id,
627            event.entity_id_str(),
628            event.event_type_str(),
629            event.timestamp,
630            offset,
631        )?;
632
633        // Process through projections
634        let projections = self.projections.read();
635        projections.process_event(event)?;
636        drop(projections);
637
638        // Process through pipelines
639        let pipeline_results = self.pipeline_manager.process_event(event);
640        if !pipeline_results.is_empty() {
641            tracing::debug!(
642                "Event {} processed by {} pipeline(s)",
643                event.id,
644                pipeline_results.len()
645            );
646            for (pipeline_id, result) in pipeline_results {
647                tracing::trace!("Pipeline {} result: {:?}", pipeline_id, result);
648            }
649        }
650
651        // Persist to Parquet storage if enabled
652        if let Some(ref storage) = self.storage {
653            let storage = storage.read();
654            storage.append_event(event.clone())?;
655        }
656
657        // Store the event in memory
658        events.push(event.clone());
659        let total_events = events.len();
660        drop(events);
661
662        // Broadcast to in-process subscribers (always on) + optional WS.
663        let event_arc = Arc::new(event.clone());
664        let _ = self.event_broadcast_tx.send(Arc::clone(&event_arc));
665        #[cfg(feature = "server")]
666        self.websocket_manager.broadcast_event(event_arc);
667
668        // Dispatch to matching webhook subscriptions
669        #[cfg(feature = "server")]
670        self.dispatch_webhooks(event);
671
672        // Update geospatial index
673        self.geo_index.index_event(event);
674
675        // Autonomous schema evolution
676        self.schema_evolution
677            .analyze_event(event.event_type_str(), &event.payload);
678
679        // Check if automatic snapshot should be created
680        self.check_auto_snapshot(event.entity_id_str(), event);
681
682        // Update metrics
683        #[cfg(feature = "server")]
684        {
685            self.metrics.events_ingested_total.inc();
686            self.metrics
687                .events_ingested_by_type
688                .with_label_values(&[event.event_type_str()])
689                .inc();
690            self.metrics.storage_events_total.set(total_events as i64);
691        }
692
693        // Update legacy total counter
694        let mut total = self.total_ingested.write();
695        *total += 1;
696
697        #[cfg(feature = "server")]
698        timer.observe_duration();
699
700        tracing::debug!("Event ingested: {} (offset: {})", event.id, offset);
701
702        Ok(())
703    }
704
705    /// Ingest a new event into the store
706    #[cfg_attr(feature = "hotpath", hotpath::measure)]
707    pub fn ingest(&self, event: &Event) -> Result<()> {
708        // Start metrics timer (v0.6 feature)
709        #[cfg(feature = "server")]
710        let timer = self.metrics.ingestion_duration_seconds.start_timer();
711
712        // Reject writes in read-only (replica) mode before touching the WAL.
713        if let Err(e) = self.ensure_writable() {
714            #[cfg(feature = "server")]
715            {
716                self.metrics.ingestion_errors_total.inc();
717                timer.observe_duration();
718            }
719            return Err(e);
720        }
721
722        // Validate event
723        let validation_result = self.validate_event(event);
724        if let Err(e) = validation_result {
725            #[cfg(feature = "server")]
726            {
727                self.metrics.ingestion_errors_total.inc();
728                timer.observe_duration();
729            }
730            return Err(e);
731        }
732
733        let _durable = self.durability_gate.read();
734
735        // Write to WAL FIRST for durability (v0.2 feature)
736        // This ensures event is persisted before processing
737        if let Some(ref wal) = self.wal
738            && let Err(e) = wal.append(event.clone())
739        {
740            #[cfg(feature = "server")]
741            {
742                self.metrics.ingestion_errors_total.inc();
743                timer.observe_duration();
744            }
745            return Err(e);
746        }
747
748        // Track per-entity version (unconditional increment, no version check)
749        *self
750            .entity_versions
751            .entry(event.entity_id_str().to_string())
752            .or_insert(0) += 1;
753
754        let mut events = self.events.write();
755        let offset = events.len();
756
757        // Index the event
758        self.index.index_event(
759            event.id,
760            event.entity_id_str(),
761            event.event_type_str(),
762            event.timestamp,
763            offset,
764        )?;
765
766        // Process through projections
767        let projections = self.projections.read();
768        projections.process_event(event)?;
769        drop(projections); // Release lock
770
771        // Process through pipelines (v0.5 feature)
772        // Pipelines can transform, filter, and aggregate events in real-time
773        let pipeline_results = self.pipeline_manager.process_event(event);
774        if !pipeline_results.is_empty() {
775            tracing::debug!(
776                "Event {} processed by {} pipeline(s)",
777                event.id,
778                pipeline_results.len()
779            );
780            // Pipeline results could be stored, emitted, or forwarded elsewhere
781            // For now, we just log them for observability
782            for (pipeline_id, result) in pipeline_results {
783                tracing::trace!("Pipeline {} result: {:?}", pipeline_id, result);
784            }
785        }
786
787        // Persist to Parquet storage if enabled (v0.2)
788        if let Some(ref storage) = self.storage {
789            let storage = storage.read();
790            storage.append_event(event.clone())?;
791        }
792
793        // Store the event in memory
794        events.push(event.clone());
795        let total_events = events.len();
796        drop(events); // Release lock early
797
798        // Broadcast to in-process subscribers + optional WS.
799        let event_arc = Arc::new(event.clone());
800        let _ = self.event_broadcast_tx.send(Arc::clone(&event_arc));
801        #[cfg(feature = "server")]
802        self.websocket_manager.broadcast_event(event_arc);
803
804        // Dispatch to matching webhook subscriptions (v0.11 feature)
805        #[cfg(feature = "server")]
806        self.dispatch_webhooks(event);
807
808        // Update geospatial index (v2.0 feature)
809        self.geo_index.index_event(event);
810
811        // Autonomous schema evolution (v2.0 feature)
812        self.schema_evolution
813            .analyze_event(event.event_type_str(), &event.payload);
814
815        // Check if automatic snapshot should be created (v0.2 feature)
816        self.check_auto_snapshot(event.entity_id_str(), event);
817
818        // Update metrics (v0.6 feature)
819        #[cfg(feature = "server")]
820        {
821            self.metrics.events_ingested_total.inc();
822            self.metrics
823                .events_ingested_by_type
824                .with_label_values(&[event.event_type_str()])
825                .inc();
826            self.metrics.storage_events_total.set(total_events as i64);
827        }
828
829        // Update legacy total counter
830        let mut total = self.total_ingested.write();
831        *total += 1;
832
833        #[cfg(feature = "server")]
834        timer.observe_duration();
835
836        tracing::debug!("Event ingested: {} (offset: {})", event.id, offset);
837
838        Ok(())
839    }
840
841    /// Ingest a batch of events with a single write lock acquisition.
842    ///
843    /// All events are validated first. If any event fails validation, no
844    /// events are stored (all-or-nothing validation). Events are then written
845    /// to WAL, indexed, processed through projections, and pushed to the
846    /// events vector under a single write lock.
847    #[cfg_attr(feature = "hotpath", hotpath::measure)]
848    pub fn ingest_batch(&self, batch: Vec<Event>) -> Result<()> {
849        if batch.is_empty() {
850            return Ok(());
851        }
852
853        // Reject writes in read-only (replica) mode before touching the WAL.
854        self.ensure_writable()?;
855
856        // Phase 1: Validate all events before acquiring any locks
857        for event in &batch {
858            self.validate_event(event)?;
859        }
860
861        let _durable = self.durability_gate.read();
862
863        // Phase 2: Write all events to WAL (before write lock, for durability)
864        if let Some(ref wal) = self.wal {
865            for event in &batch {
866                wal.append(event.clone())?;
867            }
868        }
869
870        // Phase 3: Single write lock for index + projections + push
871        let mut events = self.events.write();
872        let projections = self.projections.read();
873
874        for event in batch {
875            let offset = events.len();
876
877            self.index.index_event(
878                event.id,
879                event.entity_id_str(),
880                event.event_type_str(),
881                event.timestamp,
882                offset,
883            )?;
884
885            projections.process_event(&event)?;
886            self.pipeline_manager.process_event(&event);
887
888            if let Some(ref storage) = self.storage {
889                let storage = storage.read();
890                storage.append_event(event.clone())?;
891            }
892
893            self.geo_index.index_event(&event);
894            self.schema_evolution
895                .analyze_event(event.event_type_str(), &event.payload);
896
897            // Track per-entity version
898            *self
899                .entity_versions
900                .entry(event.entity_id_str().to_string())
901                .or_insert(0) += 1;
902
903            // Broadcast to in-process subscribers
904            let _ = self.event_broadcast_tx.send(Arc::new(event.clone()));
905
906            events.push(event);
907        }
908
909        let total_events = events.len();
910        drop(projections);
911        drop(events);
912
913        let mut total = self.total_ingested.write();
914        *total += total_events as u64;
915
916        Ok(())
917    }
918
919    /// Ingest a replicated event from the leader (follower mode).
920    ///
921    /// Unlike `ingest()`, this method:
922    /// - Skips WAL writing (the follower's WalReceiver manages its own local WAL)
923    /// - Skips schema validation (the leader already validated)
924    /// - Still indexes, processes projections/pipelines, and broadcasts to WebSocket clients
925    #[cfg_attr(feature = "hotpath", hotpath::measure)]
926    pub fn ingest_replicated(&self, event: &Event) -> Result<()> {
927        #[cfg(feature = "server")]
928        let timer = self.metrics.ingestion_duration_seconds.start_timer();
929
930        let mut events = self.events.write();
931        let offset = events.len();
932
933        // Index the event
934        self.index.index_event(
935            event.id,
936            event.entity_id_str(),
937            event.event_type_str(),
938            event.timestamp,
939            offset,
940        )?;
941
942        // Process through projections
943        let projections = self.projections.read();
944        projections.process_event(event)?;
945        drop(projections);
946
947        // Process through pipelines
948        let pipeline_results = self.pipeline_manager.process_event(event);
949        if !pipeline_results.is_empty() {
950            tracing::debug!(
951                "Replicated event {} processed by {} pipeline(s)",
952                event.id,
953                pipeline_results.len()
954            );
955        }
956
957        // Track per-entity version
958        *self
959            .entity_versions
960            .entry(event.entity_id_str().to_string())
961            .or_insert(0) += 1;
962
963        // Store the event in memory
964        events.push(event.clone());
965        let total_events = events.len();
966        drop(events);
967
968        // Broadcast to in-process subscribers + optional WS.
969        let event_arc = Arc::new(event.clone());
970        let _ = self.event_broadcast_tx.send(Arc::clone(&event_arc));
971        #[cfg(feature = "server")]
972        self.websocket_manager.broadcast_event(event_arc);
973
974        // Update metrics
975        #[cfg(feature = "server")]
976        {
977            self.metrics.events_ingested_total.inc();
978            self.metrics
979                .events_ingested_by_type
980                .with_label_values(&[event.event_type_str()])
981                .inc();
982            self.metrics.storage_events_total.set(total_events as i64);
983        }
984
985        let mut total = self.total_ingested.write();
986        *total += 1;
987
988        #[cfg(feature = "server")]
989        timer.observe_duration();
990
991        tracing::debug!(
992            "Replicated event ingested: {} (offset: {})",
993            event.id,
994            offset
995        );
996
997        Ok(())
998    }
999
1000    /// Get the current version for an entity (number of events appended for it).
1001    /// Returns 0 if the entity has no events.
1002    #[cfg_attr(feature = "hotpath", hotpath::measure)]
1003    pub fn get_entity_version(&self, entity_id: &str) -> u64 {
1004        self.entity_versions.get(entity_id).map_or(0, |v| *v)
1005    }
1006
1007    /// Get the consumer registry for durable subscriptions.
1008    pub fn consumer_registry(&self) -> &ConsumerRegistry {
1009        &self.consumer_registry
1010    }
1011
1012    /// Subscribe to every successfully-ingested event in this store.
1013    ///
1014    /// Returns a `tokio::sync::broadcast::Receiver` that yields an `Arc<Event>`
1015    /// for each ingest. Always available — does not require the `server`
1016    /// feature. Lagging receivers surface `RecvError::Lagged`.
1017    pub fn subscribe_events(&self) -> tokio::sync::broadcast::Receiver<Arc<Event>> {
1018        self.event_broadcast_tx.subscribe()
1019    }
1020
1021    /// Replace the default in-memory consumer registry with a durable one.
1022    ///
1023    /// Called during startup when system repositories are available, so that
1024    /// consumer cursors survive Core restarts via WAL persistence.
1025    pub fn set_consumer_registry(&mut self, registry: Arc<ConsumerRegistry>) {
1026        self.consumer_registry = registry;
1027    }
1028
1029    /// Get the total number of events in the store (used as max offset for consumer ack).
1030    pub fn total_events(&self) -> usize {
1031        self.events.read().len()
1032    }
1033
1034    /// Get events after a given offset, optionally filtered by event type prefixes.
1035    /// Used by consumer polling to fetch unprocessed events.
1036    pub fn events_after_offset(
1037        &self,
1038        offset: u64,
1039        filters: &[String],
1040        limit: usize,
1041    ) -> Vec<(u64, Event)> {
1042        let events = self.events.read();
1043        let start = offset as usize;
1044        if start >= events.len() {
1045            return vec![];
1046        }
1047
1048        events[start..]
1049            .iter()
1050            .enumerate()
1051            .filter(|(_, event)| ConsumerRegistry::matches_filters(event.event_type_str(), filters))
1052            .take(limit)
1053            .map(|(i, event)| ((start + i + 1) as u64, event.clone()))
1054            .collect()
1055    }
1056
1057    /// Get the WebSocket manager for this store
1058    #[cfg(feature = "server")]
1059    pub fn websocket_manager(&self) -> Arc<WebSocketManager> {
1060        Arc::clone(&self.websocket_manager)
1061    }
1062
1063    /// Get the snapshot manager for this store
1064    pub fn snapshot_manager(&self) -> Arc<SnapshotManager> {
1065        Arc::clone(&self.snapshot_manager)
1066    }
1067
1068    /// Get the compaction manager for this store
1069    pub fn compaction_manager(&self) -> Option<Arc<CompactionManager>> {
1070        self.compaction_manager.as_ref().map(Arc::clone)
1071    }
1072
1073    /// Get the schema registry for this store (v0.5 feature)
1074    pub fn schema_registry(&self) -> Arc<SchemaRegistry> {
1075        Arc::clone(&self.schema_registry)
1076    }
1077
1078    /// Get the replay manager for this store (v0.5 feature)
1079    pub fn replay_manager(&self) -> Arc<ReplayManager> {
1080        Arc::clone(&self.replay_manager)
1081    }
1082
1083    /// Get the pipeline manager for this store (v0.5 feature)
1084    pub fn pipeline_manager(&self) -> Arc<PipelineManager> {
1085        Arc::clone(&self.pipeline_manager)
1086    }
1087
1088    /// Get the metrics registry for this store (v0.6 feature)
1089    #[cfg(feature = "server")]
1090    pub fn metrics(&self) -> Arc<MetricsRegistry> {
1091        Arc::clone(&self.metrics)
1092    }
1093
1094    /// Get the projection manager for this store (v0.7 feature)
1095    pub fn projection_manager(&self) -> parking_lot::RwLockReadGuard<'_, ProjectionManager> {
1096        self.projections.read()
1097    }
1098
1099    /// Register a custom projection at runtime.
1100    ///
1101    /// The projection will receive all future events via `process()`.
1102    /// Historical events are **not** replayed — only events ingested after
1103    /// registration will be processed by this projection.
1104    ///
1105    /// See [`register_projection_with_backfill`](Self::register_projection_with_backfill)
1106    /// to also process historical events.
1107    pub fn register_projection(
1108        &self,
1109        projection: Arc<dyn crate::application::services::projection::Projection>,
1110    ) {
1111        let mut pm = self.projections.write();
1112        pm.register(projection);
1113    }
1114
1115    /// Register a custom projection and replay all existing events through it.
1116    ///
1117    /// After registration, the projection will also receive all future events.
1118    /// Historical events are replayed under a read lock — the projection's
1119    /// internal state (typically DashMap) handles concurrent access.
1120    ///
1121    /// Replay is ordered by `(timestamp, version)`. The in-memory pile can be
1122    /// physically out of order when Parquet is hydrated after WAL recovery —
1123    /// the WAL tail holds the newest events while Parquet holds the older
1124    /// history — so the backfill must sort before replaying. Projections with
1125    /// last-write-wins merge semantics produce wrong state otherwise.
1126    pub fn register_projection_with_backfill(
1127        &self,
1128        projection: &Arc<dyn crate::application::services::projection::Projection>,
1129    ) -> Result<()> {
1130        // First register so future events are processed
1131        {
1132            let mut pm = self.projections.write();
1133            pm.register(Arc::clone(projection));
1134        }
1135
1136        // Then replay existing events in chronological order under read lock
1137        let events = self.events.read();
1138        let mut ordered: Vec<&Event> = events.iter().collect();
1139        ordered.sort_by(|a, b| {
1140            a.timestamp
1141                .cmp(&b.timestamp)
1142                .then_with(|| a.version.cmp(&b.version))
1143        });
1144        for event in ordered {
1145            projection.process(event)?;
1146        }
1147
1148        Ok(())
1149    }
1150
1151    /// Eagerly reconstruct the in-memory event pile from the full Parquet
1152    /// archive.
1153    ///
1154    /// The default boot path (Step 2, issue #160) keeps Parquet cold and
1155    /// hydrates tenants lazily on first query — the multi-tenant server
1156    /// cannot fit every tenant in memory. Embedded single-store consumers
1157    /// like Prime are the opposite case: their projections *are* the
1158    /// queryable surface and never trigger the lazy query path, so the
1159    /// projections must be backfilled from the complete history.
1160    ///
1161    /// Call this *before* registering projections. The dedupe in
1162    /// `append_loaded_event` makes it safe to run after WAL recovery —
1163    /// events already replayed from the WAL are not double-counted.
1164    /// No-op (and `Ok(0)`) when no Parquet storage is configured, e.g.
1165    /// in-memory test mode.
1166    ///
1167    /// Returns the number of events newly loaded from Parquet.
1168    pub fn hydrate_all_from_storage(&self) -> Result<usize> {
1169        let Some(storage) = self.storage.as_ref().map(Arc::clone) else {
1170            return Ok(0);
1171        };
1172
1173        let events = storage.read().load_all_events()?;
1174        let read_count = events.len();
1175        let before = self.events.read().len();
1176        for event in events {
1177            self.append_loaded_event(event);
1178        }
1179        let applied = self.events.read().len() - before;
1180
1181        // The pile is now the authoritative full history.
1182        *self.total_ingested.write() = self.events.read().len() as u64;
1183
1184        tracing::info!(
1185            read = read_count,
1186            applied = applied,
1187            "🔄 hydrate_all_from_storage: in-memory pile reconstructed from Parquet"
1188        );
1189        Ok(applied)
1190    }
1191
1192    /// Get the projection state cache for this store (v0.7 feature)
1193    /// Used by Elixir Query Service for state synchronization
1194    pub fn projection_state_cache(&self) -> Arc<DashMap<String, serde_json::Value>> {
1195        Arc::clone(&self.projection_state_cache)
1196    }
1197
1198    /// Get the projection status map (v0.13 feature)
1199    pub fn projection_status(&self) -> Arc<DashMap<String, String>> {
1200        Arc::clone(&self.projection_status)
1201    }
1202
1203    /// Get the webhook registry for this store (v0.11 feature)
1204    /// Geospatial index for coordinate-based queries (v2.0 feature)
1205    pub fn geo_index(&self) -> Arc<GeoIndex> {
1206        self.geo_index.clone()
1207    }
1208
1209    /// Exactly-once processing registry (v2.0 feature)
1210    pub fn exactly_once(&self) -> Arc<ExactlyOnceRegistry> {
1211        self.exactly_once.clone()
1212    }
1213
1214    /// Schema evolution manager (v2.0 feature)
1215    pub fn schema_evolution(&self) -> Arc<SchemaEvolutionManager> {
1216        self.schema_evolution.clone()
1217    }
1218
1219    /// Get a read-locked snapshot of all events (for EventQL/GraphQL queries).
1220    ///
1221    /// Returns an `Arc` reference to the internal events vec, avoiding a full
1222    /// clone. The caller holds a read lock for the duration of the `Arc`
1223    /// lifetime — prefer short-lived usage.
1224    pub fn snapshot_events(&self) -> Vec<Event> {
1225        self.events.read().clone()
1226    }
1227
1228    /// Compact token events for an entity by replacing matching events with a
1229    /// single merged event. Used by the embedded streaming feature.
1230    ///
1231    /// Returns `Ok(true)` if compaction was performed, `Ok(false)` if no
1232    /// matching events were found.
1233    ///
1234    /// **Note:** The merged event is processed through projections *without*
1235    /// clearing the removed events' projection state first. Projections that
1236    /// accumulate state (e.g., counters) should be designed to handle this
1237    /// (the merged event replaces individual tokens, not adds to them).
1238    ///
1239    /// **The merged event is a derived view and is never made durable.**
1240    /// Compaction reclaims memory; it does not rewrite history. The token
1241    /// events stay in the WAL and Parquet exactly as they were ingested, so a
1242    /// restart — or any reader that goes to the durable log — sees the original
1243    /// stream, not the merge. A caller that needs the merged result to survive
1244    /// must ingest it as an event of its own.
1245    ///
1246    /// The write lock is held for the swap + WAL write + index rebuild.
1247    /// The index rebuild is O(N) over all events, which is acceptable for
1248    /// embedded workloads but should not be called in hot paths for large stores.
1249    pub fn compact_entity_tokens(
1250        &self,
1251        entity_id: &str,
1252        token_event_type: &str,
1253        merged_event: Event,
1254    ) -> Result<bool> {
1255        // Reject writes in read-only (replica) mode before touching the WAL.
1256        self.ensure_writable()?;
1257
1258        // Phase 1: Read-only check — do we have anything to compact?
1259        {
1260            let events = self.events.read();
1261            let has_tokens = events
1262                .iter()
1263                .any(|e| e.entity_id_str() == entity_id && e.event_type_str() == token_event_type);
1264            if !has_tokens {
1265                return Ok(false);
1266            }
1267        }
1268
1269        // Phase 2: Process merged event through projections (no write lock held)
1270        let projections = self.projections.read();
1271        projections.process_event(&merged_event)?;
1272        drop(projections);
1273
1274        // Phase 3: Acquire write lock for the swap + WAL + index rebuild
1275        let mut events = self.events.write();
1276
1277        events.retain(|e| {
1278            !(e.entity_id_str() == entity_id && e.event_type_str() == token_event_type)
1279        });
1280
1281        events.push(merged_event);
1282
1283        // Rebuild entire index since retain() shifted event positions.
1284        // Errors here indicate a corrupt event (missing entity_id/event_type)
1285        // which should not happen for well-formed events. Log and continue
1286        // rather than failing the entire compaction.
1287        self.index.clear();
1288        for (offset, event) in events.iter().enumerate() {
1289            if let Err(e) = self.index.index_event(
1290                event.id,
1291                event.entity_id_str(),
1292                event.event_type_str(),
1293                event.timestamp,
1294                offset,
1295            ) {
1296                tracing::warn!(
1297                    event_id = %event.id,
1298                    offset,
1299                    "Failed to re-index event during compaction: {e}"
1300                );
1301            }
1302        }
1303
1304        Ok(true)
1305    }
1306
1307    #[cfg(feature = "server")]
1308    pub fn webhook_registry(&self) -> Arc<WebhookRegistry> {
1309        Arc::clone(&self.webhook_registry)
1310    }
1311
1312    /// Set the channel for async webhook delivery.
1313    /// Called during server startup to wire the delivery worker.
1314    #[cfg(feature = "server")]
1315    pub fn set_webhook_tx(&self, tx: mpsc::UnboundedSender<WebhookDeliveryTask>) {
1316        *self.webhook_tx.write() = Some(tx);
1317        tracing::info!("Webhook delivery channel connected");
1318    }
1319
1320    /// Dispatch matching webhooks for a given event (non-blocking).
1321    #[cfg(feature = "server")]
1322    fn dispatch_webhooks(&self, event: &Event) {
1323        let matching = self.webhook_registry.find_matching(event);
1324        if matching.is_empty() {
1325            return;
1326        }
1327
1328        let tx_guard = self.webhook_tx.read();
1329        if let Some(ref tx) = *tx_guard {
1330            for webhook in matching {
1331                let task = WebhookDeliveryTask {
1332                    webhook,
1333                    event: event.clone(),
1334                };
1335                if let Err(e) = tx.send(task) {
1336                    tracing::warn!("Failed to queue webhook delivery: {}", e);
1337                }
1338            }
1339        }
1340    }
1341
1342    /// Manually flush any pending events to persistent storage
1343    pub fn flush_storage(&self) -> Result<()> {
1344        if let Some(ref storage) = self.storage {
1345            let storage = storage.read();
1346            storage.flush()?;
1347            tracing::info!("✅ Flushed events to persistent storage");
1348        }
1349        Ok(())
1350    }
1351
1352    /// Run a checkpoint: flush pending Parquet batches, then truncate the
1353    /// WAL through the checkpoint point (Step 6 of the sustainable data
1354    /// strategy).
1355    ///
1356    /// Order matters. We flush Parquet first; only on success do we
1357    /// truncate the WAL. If the process crashes between the flush and the
1358    /// truncate, the WAL still contains the events that were just durably
1359    /// written, and recovery will replay them. The dedupe in
1360    /// `append_loaded_event` (index probe) makes that idempotent — the
1361    /// event is already in Parquet, so the lazy-load splice no-ops once
1362    /// the tenant is hydrated.
1363    ///
1364    /// The reverse order would be unsafe: a crash between truncate and
1365    /// flush would lose committed events.
1366    ///
1367    /// This bounds dirty-restart replay time to one checkpoint interval
1368    /// regardless of total dataset size — that's the load-bearing
1369    /// property for cold-start time as ingest rate grows.
1370    ///
1371    /// No-op when no WAL is configured (in-memory-only mode).
1372    pub fn checkpoint(&self) -> Result<()> {
1373        let Some(ref wal) = self.wal else {
1374            // No WAL (in-memory-only mode). Still refresh storage-size metrics in
1375            // case Parquet-only persistence is configured.
1376            #[cfg(feature = "server")]
1377            self.refresh_storage_metrics();
1378            return Ok(());
1379        };
1380
1381        if self.read_only {
1382            return Ok(());
1383        }
1384
1385        // Seal under the gate so every event in a sealed segment has already
1386        // reached the Parquet batch this flush drains. Events that arrive
1387        // after the seal land in the active segment, which is kept.
1388        let active = {
1389            let _sealing = self.durability_gate.write();
1390            wal.seal()?
1391        };
1392        self.flush_storage()?;
1393        wal.remove_sealed(&active)?;
1394        tracing::debug!("✅ Checkpoint complete: Parquet flushed, sealed WAL segments retired");
1395
1396        // Recompute on-disk storage size now that Parquet is flushed and the WAL
1397        // is truncated — the gauge reflects the post-checkpoint footprint.
1398        #[cfg(feature = "server")]
1399        self.refresh_storage_metrics();
1400
1401        Ok(())
1402    }
1403
1404    /// Public entrypoint to populate the on-disk storage gauges once, e.g. at boot
1405    /// so the dashboard's "storage" card is correct before the first checkpoint
1406    /// tick (and even when the checkpoint loop is disabled). Delegates to the
1407    /// internal refresh. No-op without persistent storage.
1408    #[cfg(feature = "server")]
1409    pub fn refresh_storage_metrics_now(&self) {
1410        self.refresh_storage_metrics();
1411    }
1412
1413    /// Recompute the on-disk storage gauges from the real Parquet + WAL files and
1414    /// publish them to Prometheus: `allsource_storage_size_bytes` (Parquet bytes +
1415    /// WAL segment bytes), `allsource_parquet_files_total`, and
1416    /// `allsource_wal_segments_total`.
1417    ///
1418    /// These gauges were registered but never set, so they read a constant 0 — the
1419    /// dashboard's "storage" card therefore showed `—`. This is the population.
1420    ///
1421    /// HONESTY: this is a **platform/process-wide** figure — the size of the whole
1422    /// data directory across all tenants, not any single tenant's storage. It is
1423    /// surfaced as a platform metric and must not be presented as a tenant number.
1424    ///
1425    /// Called from the checkpoint loop (default every 60s), not the ingest hot
1426    /// path: it does one `statx` per Parquet file + per WAL segment. Best-effort —
1427    /// a stat error logs and leaves the previous gauge value in place rather than
1428    /// resetting it to a misleading 0.
1429    #[cfg(feature = "server")]
1430    fn refresh_storage_metrics(&self) {
1431        let Some(ref storage) = self.storage else {
1432            return;
1433        };
1434
1435        let parquet_stats = match storage.read().stats() {
1436            Ok(stats) => stats,
1437            Err(e) => {
1438                tracing::warn!("storage-size metric refresh: failed to stat Parquet: {e}");
1439                return;
1440            }
1441        };
1442
1443        let (wal_bytes, wal_segments) = match self.wal.as_ref() {
1444            Some(wal) => match wal.on_disk_stats() {
1445                Ok(stats) => stats,
1446                Err(e) => {
1447                    tracing::warn!("storage-size metric refresh: failed to stat WAL: {e}");
1448                    (0, 0)
1449                }
1450            },
1451            None => (0, 0),
1452        };
1453
1454        let total_bytes = parquet_stats.total_size_bytes + wal_bytes;
1455
1456        self.metrics
1457            .storage_size_bytes
1458            .set(total_bytes.min(i64::MAX as u64) as i64);
1459        self.metrics
1460            .parquet_files_total
1461            .set(parquet_stats.total_files as i64);
1462        self.metrics.wal_segments_total.set(wal_segments as i64);
1463
1464        tracing::debug!(
1465            "storage-size metrics refreshed: {} bytes total ({} Parquet files, {} WAL segments)",
1466            total_bytes,
1467            parquet_stats.total_files,
1468            wal_segments
1469        );
1470    }
1471
1472    /// Get the configured checkpoint cadence (used by background tasks).
1473    pub fn checkpoint_interval(&self) -> Option<std::time::Duration> {
1474        self.checkpoint_interval_secs
1475            .map(std::time::Duration::from_secs)
1476    }
1477
1478    /// Hydrate `tenant_id`'s persisted Parquet data into the in-memory
1479    /// pile if it isn't already loaded. Cheap on the warm path
1480    /// (DashMap probe); on the cold path it walks just that tenant's
1481    /// subtree (`load_events_for_tenant`) and splices the events into
1482    /// `events`/`index`/`projections`/`entity_versions`.
1483    ///
1484    /// Concurrent first-callers for the same tenant serialize on a
1485    /// per-tenant Mutex (singleflight) so the disk read happens once.
1486    /// Other tenants are unaffected — distinct lock per tenant.
1487    ///
1488    /// Returns `Err` if the tenant_id fails the path-safety
1489    /// whitelist, the Parquet read fails, or another in-flight load
1490    /// holds the lock past the configured timeout. The caller (a
1491    /// query handler) is expected to surface that as a 5xx — see
1492    /// Step 2's "no infinite hangs" acceptance criterion.
1493    ///
1494    /// On failure, `loaded` is NOT marked, so a transient error is
1495    /// retried on the next request rather than poisoning the tenant
1496    /// permanently. A future commit may add a circuit breaker if
1497    /// thrash becomes an issue.
1498    ///
1499    /// No-op (and Ok) when no Parquet storage is configured — the
1500    /// in-memory-only mode used by tests has nothing to hydrate.
1501    pub fn ensure_tenant_loaded(&self, tenant_id: &str) -> Result<()> {
1502        // Fast path: warm tenant. Avoids the Mutex altogether.
1503        if self.tenant_loader.is_loaded(tenant_id) {
1504            return Ok(());
1505        }
1506
1507        let Some(storage) = self.storage.as_ref().map(Arc::clone) else {
1508            // No persistent storage to load from. Mark loaded so we
1509            // don't keep re-entering the slow path.
1510            self.tenant_loader.mark_loaded(tenant_id);
1511            return Ok(());
1512        };
1513
1514        // Singleflight: get-or-insert the per-tenant lock and try to
1515        // acquire it within the timeout budget.
1516        let lock = self.tenant_loader.lock_for(tenant_id);
1517        let timeout = self.tenant_loader.load_timeout();
1518        let _guard = lock.try_lock_for(timeout).ok_or_else(|| {
1519            AllSourceError::StorageError(format!(
1520                "ensure_tenant_loaded timed out after {timeout:?} waiting for in-flight load of \
1521                 tenant {tenant_id:?}"
1522            ))
1523        })?;
1524
1525        // Re-check inside the lock — another thread may have completed
1526        // the load while we were waiting.
1527        if self.tenant_loader.is_loaded(tenant_id) {
1528            return Ok(());
1529        }
1530
1531        let started = std::time::Instant::now();
1532        let events = storage.read().load_events_for_tenant(tenant_id)?;
1533        let read_count = events.len();
1534
1535        let before = self.events.read().len();
1536        for event in events {
1537            self.append_loaded_event(event);
1538        }
1539        let applied = self.events.read().len() - before;
1540
1541        // total_ingested only counts events newly added to memory.
1542        // Dedupe (e.g. WAL events re-checkpointed to Parquet) makes
1543        // applied < read_count possible.
1544        *self.total_ingested.write() += applied as u64;
1545        self.tenant_loader.mark_loaded(tenant_id);
1546
1547        tracing::info!(
1548            tenant_id = tenant_id,
1549            read = read_count,
1550            applied = applied,
1551            elapsed_ms = started.elapsed().as_millis() as u64,
1552            "ensure_tenant_loaded: tenant hydrated"
1553        );
1554
1555        // Budget check (Step 3 #3). After splicing the new
1556        // tenant in, evict LRU tenants until we're back under
1557        // budget. Excludes the just-loaded tenant from the
1558        // candidate set — otherwise a single oversized tenant
1559        // would evict itself in a tight loop. If no other tenant
1560        // is loaded, accept the over-budget state and log a
1561        // warning so ops can see it.
1562        self.enforce_cache_budget(tenant_id);
1563
1564        // Refresh the resident-bytes gauge — covers both the
1565        // load-with-no-eviction case and the post-eviction state.
1566        #[cfg(feature = "server")]
1567        self.metrics
1568            .cache_bytes
1569            .set(self.tenant_loader.total_bytes() as i64);
1570
1571        Ok(())
1572    }
1573
1574    /// Walks the LRU until total resident bytes are within the
1575    /// configured budget, calling `evict_tenant` on each victim.
1576    /// Excludes `recently_touched` from the candidate set so we
1577    /// don't evict the tenant that just triggered the call. Step 3
1578    /// #3 entry point.
1579    ///
1580    /// No-op when no budget is set or when already under budget —
1581    /// most queries take the early-return fast path. Eviction is
1582    /// the cold path; with a well-sized budget this only fires
1583    /// during warm-up of new tenants past the working-set size.
1584    fn enforce_cache_budget(&self, recently_touched: &str) {
1585        if !self.tenant_loader.over_budget() {
1586            return;
1587        }
1588        loop {
1589            let Some(victim) = self.tenant_loader.pick_lru_excluding(recently_touched) else {
1590                tracing::warn!(
1591                    cache_bytes = self.tenant_loader.total_bytes(),
1592                    budget = self.tenant_loader.byte_budget(),
1593                    recently_touched = recently_touched,
1594                    "cache over budget but no other tenant available to evict — \
1595                     a single tenant exceeds the budget; consider raising it"
1596                );
1597                return;
1598            };
1599            self.evict_tenant(&victim);
1600            if !self.tenant_loader.over_budget() {
1601                return;
1602            }
1603        }
1604    }
1605
1606    /// True iff `ensure_tenant_loaded` has previously succeeded for
1607    /// this tenant. Diagnostic / testing API.
1608    pub fn is_tenant_loaded(&self, tenant_id: &str) -> bool {
1609        self.tenant_loader.is_loaded(tenant_id)
1610    }
1611
1612    /// Drop `tenant_id` from the in-memory cache. Step 3 #2 of the
1613    /// sustainable data strategy.
1614    ///
1615    /// Removes every event for this tenant from the events Vec,
1616    /// rebuilds the index/entity_versions for the retained events
1617    /// (Vec offsets shift on remove, so the index has to be
1618    /// rebuilt), and resets the tenant_loader bookkeeping so a
1619    /// subsequent query triggers a fresh `ensure_tenant_loaded`.
1620    ///
1621    /// **Parquet is canonical, in-memory is just cache.** This
1622    /// only affects the in-memory side. Disk data is untouched —
1623    /// that's why eviction is safe even for tenants with
1624    /// recently-ingested data: a query after eviction transparently
1625    /// re-reads from Parquet.
1626    ///
1627    /// Projection state is NOT rolled back. Projections accumulate
1628    /// across boots and tenants (their durability story is
1629    /// separate); subtracting them would need replay support that
1630    /// doesn't exist in this commit. After eviction + re-load,
1631    /// projections may double-count the re-loaded events. Step 3
1632    /// #4's stress test only asserts the cache budget is held; a
1633    /// future commit will tackle projection-aware eviction.
1634    ///
1635    /// Locking: takes the events write lock for the full duration
1636    /// of the filter + re-index. Concurrent ingest blocks until
1637    /// done. Eviction is the cold path; the working set should
1638    /// stay in budget so this rarely fires.
1639    pub fn evict_tenant(&self, tenant_id: &str) {
1640        let mut events = self.events.write();
1641        let before = events.len();
1642        let evicted_bytes = self.tenant_loader.bytes_for(tenant_id);
1643
1644        events.retain(|e| e.tenant_id_str() != tenant_id);
1645        let after = events.len();
1646        let dropped = before - after;
1647
1648        if dropped == 0 {
1649            // Tenant had no events in memory. Still clear loader
1650            // state (e.g. a "loaded with zero events" marker) so
1651            // is_tenant_loaded reports the right thing.
1652            drop(events);
1653            self.tenant_loader.mark_unloaded(tenant_id);
1654            return;
1655        }
1656
1657        // Rebuild the index — Vec offsets shifted under retain().
1658        // Rebuild entity_versions from scratch too, since the
1659        // counter reflects "how many events of this entity remain".
1660        self.index.clear();
1661        self.entity_versions.clear();
1662        for (offset, event) in events.iter().enumerate() {
1663            if let Err(e) = self.index.index_event(
1664                event.id,
1665                event.entity_id_str(),
1666                event.event_type_str(),
1667                event.timestamp,
1668                offset,
1669            ) {
1670                tracing::error!(
1671                    "Failed to re-index event during eviction of {}: {}",
1672                    tenant_id,
1673                    e
1674                );
1675            }
1676            *self
1677                .entity_versions
1678                .entry(event.entity_id_str().to_string())
1679                .or_insert(0) += 1;
1680        }
1681        drop(events);
1682
1683        self.tenant_loader.mark_unloaded(tenant_id);
1684
1685        // total_ingested under Steps 2-3 means "events currently
1686        // resident in memory". Subtract what we just dropped.
1687        let mut t = self.total_ingested.write();
1688        *t = t.saturating_sub(dropped as u64);
1689        drop(t);
1690
1691        // Step 3 #4: cache observability. Increment the eviction
1692        // counter, refresh the resident-bytes gauge.
1693        #[cfg(feature = "server")]
1694        {
1695            self.metrics.cache_evictions_total.inc();
1696            self.metrics
1697                .cache_bytes
1698                .set(self.tenant_loader.total_bytes() as i64);
1699        }
1700
1701        tracing::info!(
1702            tenant_id = tenant_id,
1703            events_dropped = dropped,
1704            bytes_freed = evicted_bytes,
1705            "evicted tenant from memory cache"
1706        );
1707    }
1708
1709    /// Approximate resident bytes a single tenant occupies in the
1710    /// in-memory cache. Step 3 budget-tracking input. 0 for cold
1711    /// or evicted tenants.
1712    pub fn tenant_resident_bytes(&self, tenant_id: &str) -> u64 {
1713        self.tenant_loader.bytes_for(tenant_id)
1714    }
1715
1716    /// Sum of resident-byte estimates across every loaded tenant.
1717    /// What the budget check compares against.
1718    pub fn cache_resident_bytes(&self) -> u64 {
1719        self.tenant_loader.total_bytes()
1720    }
1721
1722    /// Splice a single loaded event into the in-memory structures
1723    /// (events vec, index, projections, entity_versions) atomically
1724    /// w.r.t. concurrent ingest. Used by `ensure_tenant_loaded`.
1725    ///
1726    /// The WAL recovery path on boot has its own (single-threaded)
1727    /// variant inline because boot can't race with ingest. This
1728    /// helper is the variant safe to call while traffic is flowing
1729    /// — it holds the events write lock across the index/offset
1730    /// assignment so (offset, push) stays atomic.
1731    ///
1732    /// Dedupes against events already in memory by event ID. Two
1733    /// paths can surface the same event:
1734    /// 1. WAL recovery on boot pushed it into memory.
1735    /// 2. The event was then checkpointed to Parquet and the
1736    ///    WAL truncated. A later ensure_tenant_loaded re-reads
1737    ///    the Parquet file, including this event.
1738    ///
1739    /// Without the dedupe, step 2 would double-count the event.
1740    /// The check is O(1) — DashMap probe by UUID — and the
1741    /// alternative (loading every tenant before truncating WAL)
1742    /// would defeat the lazy-load.
1743    fn append_loaded_event(&self, event: Event) {
1744        if self.index.get_by_id(&event.id).is_some() {
1745            return;
1746        }
1747
1748        let event_bytes = event.estimated_size_bytes();
1749        let tenant = event.tenant_id_str().to_string();
1750
1751        let mut events = self.events.write();
1752        let offset = events.len();
1753
1754        if let Err(e) = self.index.index_event(
1755            event.id,
1756            event.entity_id_str(),
1757            event.event_type_str(),
1758            event.timestamp,
1759            offset,
1760        ) {
1761            tracing::error!("Failed to index loaded event {}: {}", event.id, e);
1762        }
1763
1764        if let Err(e) = self.projections.read().process_event(&event) {
1765            tracing::error!("Failed to project loaded event {}: {}", event.id, e);
1766        }
1767
1768        *self
1769            .entity_versions
1770            .entry(event.entity_id_str().to_string())
1771            .or_insert(0) += 1;
1772
1773        events.push(event);
1774        // Account for the bytes AFTER the push so a panic in the
1775        // index/projection path doesn't leave the counter inflated.
1776        // The DashMap update is itself the last fallible step.
1777        self.tenant_loader.add_bytes(&tenant, event_bytes);
1778    }
1779
1780    /// Manually create a snapshot for an entity
1781    pub fn create_snapshot(&self, entity_id: &str) -> Result<()> {
1782        // Get all events for this entity
1783        let events = self.query(&QueryEventsRequest {
1784            entity_id: Some(entity_id.to_string()),
1785            event_type: None,
1786            tenant_id: None,
1787            as_of: None,
1788            since: None,
1789            until: None,
1790            limit: None,
1791            event_type_prefix: None,
1792            exclude_event_type_prefix: None,
1793            payload_filter: None,
1794        })?;
1795
1796        if events.is_empty() {
1797            return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
1798        }
1799
1800        // Build current state
1801        let mut state = serde_json::json!({});
1802        for event in &events {
1803            if let serde_json::Value::Object(ref mut state_map) = state
1804                && let serde_json::Value::Object(ref payload_map) = event.payload
1805            {
1806                for (key, value) in payload_map {
1807                    state_map.insert(key.clone(), value.clone());
1808                }
1809            }
1810        }
1811
1812        let last_event = events.last().unwrap();
1813        self.snapshot_manager.create_snapshot(
1814            entity_id,
1815            state,
1816            last_event.timestamp,
1817            events.len(),
1818            SnapshotType::Manual,
1819        )?;
1820
1821        Ok(())
1822    }
1823
1824    /// Check and create automatic snapshots if needed
1825    fn check_auto_snapshot(&self, entity_id: &str, event: &Event) {
1826        // Count events for this entity
1827        let entity_event_count = self
1828            .index
1829            .get_by_entity(entity_id)
1830            .map_or(0, |entries| entries.len());
1831
1832        if self.snapshot_manager.should_create_snapshot(
1833            entity_id,
1834            entity_event_count,
1835            event.timestamp,
1836        ) {
1837            // Create snapshot in background (don't block ingestion)
1838            if let Err(e) = self.create_snapshot(entity_id) {
1839                tracing::warn!(
1840                    "Failed to create automatic snapshot for {}: {}",
1841                    entity_id,
1842                    e
1843                );
1844            }
1845        }
1846    }
1847
1848    /// Validate an event before ingestion
1849    fn validate_event(&self, event: &Event) -> Result<()> {
1850        // EntityId and EventType value objects already validate non-empty in their constructors
1851        // So these checks are now redundant, but we keep them for explicit validation
1852        if event.entity_id_str().is_empty() {
1853            return Err(AllSourceError::ValidationError(
1854                "entity_id cannot be empty".to_string(),
1855            ));
1856        }
1857
1858        if event.event_type_str().is_empty() {
1859            return Err(AllSourceError::ValidationError(
1860                "event_type cannot be empty".to_string(),
1861            ));
1862        }
1863
1864        // Reject system namespace events from user-facing ingestion.
1865        // System events are written exclusively via SystemMetadataStore.
1866        if event.event_type().is_system() {
1867            return Err(AllSourceError::ValidationError(
1868                "Event types starting with '_system.' are reserved for internal use".to_string(),
1869            ));
1870        }
1871
1872        Ok(())
1873    }
1874
1875    /// Reset a projection by clearing its state and reprocessing all events
1876    pub fn reset_projection(&self, name: &str) -> Result<usize> {
1877        let projection_manager = self.projections.read();
1878        let projection = projection_manager.get_projection(name).ok_or_else(|| {
1879            AllSourceError::EntityNotFound(format!("Projection '{name}' not found"))
1880        })?;
1881
1882        // Clear existing state
1883        projection.clear();
1884
1885        // Clear cached state for this projection
1886        let prefix = format!("{name}:");
1887        let keys_to_remove: Vec<String> = self
1888            .projection_state_cache
1889            .iter()
1890            .filter(|entry| entry.key().starts_with(&prefix))
1891            .map(|entry| entry.key().clone())
1892            .collect();
1893        for key in keys_to_remove {
1894            self.projection_state_cache.remove(&key);
1895        }
1896
1897        // Reprocess all events through this projection
1898        let events = self.events.read();
1899        let mut reprocessed = 0usize;
1900        for event in events.iter() {
1901            if projection.process(event).is_ok() {
1902                reprocessed += 1;
1903            }
1904        }
1905
1906        Ok(reprocessed)
1907    }
1908
1909    /// Get a single event by its UUID
1910    pub fn get_event_by_id(&self, event_id: &uuid::Uuid) -> Result<Option<Event>> {
1911        if let Some(offset) = self.index.get_by_id(event_id) {
1912            let events = self.events.read();
1913            Ok(events.get(offset).cloned())
1914        } else {
1915            Ok(None)
1916        }
1917    }
1918
1919    /// Query events based on filters (optimized with indices).
1920    ///
1921    /// Ascending by `(timestamp, version)`, windowed by `request.limit`.
1922    #[cfg_attr(feature = "hotpath", hotpath::measure)]
1923    pub fn query(&self, request: &QueryEventsRequest) -> Result<Vec<Event>> {
1924        self.query_window(request, 0, false)
1925            .map(|(events, _)| events)
1926    }
1927
1928    /// [`Self::query`] under a read scope the request body cannot widen.
1929    pub fn query_scoped(
1930        &self,
1931        request: &QueryEventsRequest,
1932        scope: &ReadScope,
1933    ) -> Result<Vec<Event>> {
1934        self.query_window_scoped(request, 0, false, scope)
1935            .map(|(events, _)| events)
1936    }
1937
1938    /// Query events and return only the requested window, plus the total number
1939    /// of matches the window was taken from.
1940    ///
1941    /// `request.limit` bounds the window; `offset` skips that many matches
1942    /// first; `descending` returns newest first. The returned total is the
1943    /// pre-window match count, which is what `total_count`/`has_more` need.
1944    ///
1945    /// Cost (issue #251). Matching still costs an O(N) index scan over the
1946    /// entity's/type's entries — that is inherent to reporting `total`, and this
1947    /// method does not pretend otherwise. What `limit` *does* bound is the
1948    /// ordering and the materialization: matches are held as borrowed
1949    /// references, the `offset + limit` window is partitioned off with
1950    /// `select_nth_unstable_by` (O(N)) and only that window is sorted
1951    /// (O(k log k)) and cloned. So a `limit=1` read over a 100k-event entity
1952    /// scans 100k index entries but sorts and clones one event, where it
1953    /// previously sorted and cloned 100k.
1954    #[cfg_attr(feature = "hotpath", hotpath::measure)]
1955    pub fn query_window(
1956        &self,
1957        request: &QueryEventsRequest,
1958        offset: usize,
1959        descending: bool,
1960    ) -> Result<(Vec<Event>, usize)> {
1961        self.query_window_scoped(request, offset, descending, &ReadScope::unrestricted())
1962    }
1963
1964    /// [`Self::query_window`] under a [`ReadScope`].
1965    ///
1966    /// Every read path in this store funnels through here, which is the point:
1967    /// a scope enforced in one API surface leaves the others open, and the
1968    /// surfaces that matter (HTTP, the embedded API, MCP) all end up on this
1969    /// method. The scope is a separate argument rather than a field on
1970    /// `QueryEventsRequest` **because that type is deserialized from the
1971    /// request body** — a caller could then widen its own scope by omitting the
1972    /// field.
1973    pub fn query_window_scoped(
1974        &self,
1975        request: &QueryEventsRequest,
1976        offset: usize,
1977        descending: bool,
1978        scope: &ReadScope,
1979    ) -> Result<(Vec<Event>, usize)> {
1980        // Reject a payload filter that cannot be applied, BEFORE doing any
1981        // work. `apply_filters` parses it per event with `if let Ok(..)`, so an
1982        // unparseable filter degraded into NO filter: the query answered with
1983        // every event the caller had asked to exclude, `total_count` agreed,
1984        // and nothing in the response said the filter had been dropped. A
1985        // filter must fail closed — loudly — not open.
1986        if let Some(filter) = &request.payload_filter
1987            && serde_json::from_str::<serde_json::Map<String, serde_json::Value>>(filter).is_err()
1988        {
1989            return Err(AllSourceError::InvalidInput(format!(
1990                "invalid 'payload_filter': expected a JSON object of field/value \
1991                 pairs, got '{filter}'"
1992            )));
1993        }
1994
1995        // Lazy-load gate (Step 2): if the request scopes to a tenant,
1996        // make sure that tenant's persisted data is in memory before
1997        // running the in-memory index lookup. First call for a cold
1998        // tenant blocks here for the disk read (single-digit seconds
1999        // on ~100k events); warm tenants take the DashMap fast path
2000        // and add no measurable latency.
2001        //
2002        // Errors propagate as `Err`; the HTTP layer turns that into
2003        // a 5xx, which is the explicit "no infinite hangs" contract
2004        // from the Step 2 acceptance criteria.
2005        //
2006        // Unfiltered (cross-tenant) queries — `tenant_id = None` —
2007        // run against whatever is currently in memory. They cannot
2008        // pre-load every tenant without defeating the whole point
2009        // of the lazy-load model. In practice the gateway always
2010        // injects an auth-derived `tenant_id`; an unfiltered query
2011        // is admin-only and gets degraded results until a future
2012        // commit adds an explicit "load all tenants" admin path.
2013        if let Some(ref tenant_id) = request.tenant_id {
2014            self.ensure_tenant_loaded(tenant_id)?;
2015            // LRU touch — the most-recently-queried tenant moves
2016            // to the back of the eviction queue. Cheap (single
2017            // DashMap insert), called on every per-tenant query.
2018            self.tenant_loader.touch(tenant_id);
2019        }
2020
2021        // Determine query type for metrics (v0.6 feature)
2022        let query_type = if request.entity_id.is_some() {
2023            "entity"
2024        } else if request.event_type.is_some() {
2025            "type"
2026        } else if request.event_type_prefix.is_some() {
2027            "type_prefix"
2028        } else {
2029            "full_scan"
2030        };
2031
2032        // Start metrics timer (v0.6 feature)
2033        #[cfg(feature = "server")]
2034        let timer = self
2035            .metrics
2036            .query_duration_seconds
2037            .with_label_values(&[query_type])
2038            .start_timer();
2039
2040        // Increment query counter (v0.6 feature)
2041        #[cfg(feature = "server")]
2042        self.metrics
2043            .queries_total
2044            .with_label_values(&[query_type])
2045            .inc();
2046
2047        let events = self.events.read();
2048
2049        // Use index for fast lookups
2050        let offsets: Vec<usize> = if let Some(entity_id) = &request.entity_id {
2051            // Use entity index
2052            self.index
2053                .get_by_entity(entity_id)
2054                .map(|entries| self.filter_entries(entries, request))
2055                .unwrap_or_default()
2056        } else if let Some(event_type) = &request.event_type {
2057            // Use type index (exact match)
2058            self.index
2059                .get_by_type(event_type)
2060                .map(|entries| self.filter_entries(entries, request))
2061                .unwrap_or_default()
2062        } else if let Some(prefix) = &request.event_type_prefix {
2063            // Use type index (prefix match)
2064            let entries = self.index.get_by_type_prefix(prefix);
2065            self.filter_entries(entries, request)
2066        } else {
2067            // Full scan (less efficient but necessary for complex queries)
2068            (0..events.len()).collect()
2069        };
2070
2071        // Apply remaining filters against BORROWED events — cloning happens
2072        // only for the window that survives the selection below. Each match
2073        // carries its scan position, which is the final tie-breaker and makes
2074        // the ordering a *total* order; that is what lets the partial selection
2075        // below produce exactly what a stable full sort would have.
2076        let mut matches: Vec<(usize, &Event)> = offsets
2077            .iter()
2078            .filter_map(|&event_offset| events.get(event_offset))
2079            .filter(|event| scope.permits(event.entity_id().as_str()))
2080            .filter(|event| self.apply_filters(event, request))
2081            .enumerate()
2082            .collect();
2083
2084        // Total is the match count before windowing — callers surface it as
2085        // `total_count`/`has_more`.
2086        let total = matches.len();
2087
2088        // Order by timestamp ascending, with version as a deterministic
2089        // tie-breaker so events that share a timestamp keep a stable,
2090        // well-defined order — "the latest event" must be unambiguous
2091        // (issue #177) — then by scan position, so equal (timestamp, version)
2092        // pairs keep the order a stable sort gave them. `descending` is the
2093        // exact reverse of that total order.
2094        let order = |a: &(usize, &Event), b: &(usize, &Event)| {
2095            let ascending =
2096                a.1.timestamp
2097                    .cmp(&b.1.timestamp)
2098                    .then_with(|| a.1.version.cmp(&b.1.version))
2099                    .then_with(|| a.0.cmp(&b.0));
2100            if descending {
2101                ascending.reverse()
2102            } else {
2103                ascending
2104            }
2105        };
2106
2107        // Bounded selection (issue #251): sort only the requested window.
2108        select_window(&mut matches, offset, request.limit, order);
2109
2110        // Apply the offset. `limit` is already accounted for: either the window
2111        // was truncated by `select_window`, or it covered the whole match set.
2112        let results: Vec<Event> = matches
2113            .into_iter()
2114            .skip(offset)
2115            .map(|(_, event)| event.clone())
2116            .collect();
2117
2118        // Record query results count (v0.6 feature)
2119        #[cfg(feature = "server")]
2120        {
2121            self.metrics
2122                .query_results_total
2123                .with_label_values(&[query_type])
2124                .inc_by(results.len() as u64);
2125            timer.observe_duration();
2126        }
2127
2128        Ok((results, total))
2129    }
2130
2131    /// Filter index entries based on query parameters
2132    #[cfg_attr(feature = "hotpath", hotpath::measure)]
2133    fn filter_entries(&self, entries: Vec<IndexEntry>, request: &QueryEventsRequest) -> Vec<usize> {
2134        entries
2135            .into_iter()
2136            .filter(|entry| {
2137                // Time filters
2138                if let Some(as_of) = request.as_of
2139                    && entry.timestamp > as_of
2140                {
2141                    return false;
2142                }
2143                if let Some(since) = request.since
2144                    && entry.timestamp < since
2145                {
2146                    return false;
2147                }
2148                if let Some(until) = request.until
2149                    && entry.timestamp > until
2150                {
2151                    return false;
2152                }
2153                true
2154            })
2155            .map(|entry| entry.offset)
2156            .collect()
2157    }
2158
2159    /// Apply filters to an event
2160    #[cfg_attr(feature = "hotpath", hotpath::measure)]
2161    fn apply_filters(&self, event: &Event, request: &QueryEventsRequest) -> bool {
2162        // Tenant isolation: if a tenant_id is specified, only return events from that tenant
2163        if let Some(ref tid) = request.tenant_id
2164            && event.tenant_id_str() != tid
2165        {
2166            return false;
2167        }
2168
2169        // Time range. Also pre-applied to index entries in `filter_entries`
2170        // (cheaper — it skips the event fetch), but `filter_entries` runs ONLY
2171        // on the indexed branches. A query that no index narrows — scoped by
2172        // tenant alone, or by `payload_filter`/`exclude_event_type_prefix` —
2173        // takes the full-scan branch and never reaches it, so without this
2174        // `since`/`until`/`as_of` were silently dropped and the query answered
2175        // with the entire history (`total`/`has_more` included). Re-checking
2176        // here is idempotent for the indexed paths: same predicate, and an
2177        // index entry's timestamp is its event's.
2178        if let Some(as_of) = request.as_of
2179            && event.timestamp > as_of
2180        {
2181            return false;
2182        }
2183        if let Some(since) = request.since
2184            && event.timestamp < since
2185        {
2186            return false;
2187        }
2188        if let Some(until) = request.until
2189            && event.timestamp > until
2190        {
2191            return false;
2192        }
2193
2194        // Exclusion: drop events whose type starts with any excluded prefix
2195        // (comma-separated). Applied here, before sort+limit, so excluded events
2196        // never consume the result window.
2197        if let Some(ref excludes) = request.exclude_event_type_prefix {
2198            let et = event.event_type_str();
2199            if excludes
2200                .split(',')
2201                .map(str::trim)
2202                .filter(|p| !p.is_empty())
2203                .any(|p| et.starts_with(p))
2204            {
2205                return false;
2206            }
2207        }
2208
2209        // Additional type filter if entity was primary
2210        if request.entity_id.is_some()
2211            && let Some(ref event_type) = request.event_type
2212            && event.event_type_str() != event_type
2213        {
2214            return false;
2215        }
2216
2217        // Additional prefix filter if entity was primary
2218        if request.entity_id.is_some()
2219            && let Some(ref prefix) = request.event_type_prefix
2220            && !event.event_type_str().starts_with(prefix)
2221        {
2222            return false;
2223        }
2224
2225        // Payload field filtering
2226        if let Some(ref filter_str) = request.payload_filter
2227            && let Ok(filter_obj) =
2228                serde_json::from_str::<serde_json::Map<String, serde_json::Value>>(filter_str)
2229        {
2230            let payload = event.payload();
2231            for (key, expected_value) in &filter_obj {
2232                match payload.get(key) {
2233                    Some(actual_value) if actual_value == expected_value => {}
2234                    _ => return false,
2235                }
2236            }
2237        }
2238
2239        true
2240    }
2241
2242    /// Reconstruct entity state as of a specific timestamp
2243    /// v0.2: Now uses snapshots for fast reconstruction
2244    #[cfg_attr(feature = "hotpath", hotpath::measure)]
2245    pub fn reconstruct_state(
2246        &self,
2247        entity_id: &str,
2248        as_of: Option<DateTime<Utc>>,
2249    ) -> Result<serde_json::Value> {
2250        // Try to find a snapshot to use as a base (v0.2 optimization)
2251        let (merged_state, since_timestamp) = if let Some(as_of_time) = as_of {
2252            // Get snapshot closest to requested time
2253            if let Some(snapshot) = self
2254                .snapshot_manager
2255                .get_snapshot_as_of(entity_id, as_of_time)
2256            {
2257                tracing::debug!(
2258                    "Using snapshot from {} for entity {} (saved {} events)",
2259                    snapshot.as_of,
2260                    entity_id,
2261                    snapshot.event_count
2262                );
2263                (snapshot.state.clone(), Some(snapshot.as_of))
2264            } else {
2265                (serde_json::json!({}), None)
2266            }
2267        } else {
2268            // Get latest snapshot for current state
2269            if let Some(snapshot) = self.snapshot_manager.get_latest_snapshot(entity_id) {
2270                tracing::debug!(
2271                    "Using latest snapshot from {} for entity {}",
2272                    snapshot.as_of,
2273                    entity_id
2274                );
2275                (snapshot.state.clone(), Some(snapshot.as_of))
2276            } else {
2277                (serde_json::json!({}), None)
2278            }
2279        };
2280
2281        // Query events after the snapshot (or all if no snapshot)
2282        let events = self.query(&QueryEventsRequest {
2283            entity_id: Some(entity_id.to_string()),
2284            event_type: None,
2285            tenant_id: None,
2286            as_of,
2287            since: since_timestamp,
2288            until: None,
2289            limit: None,
2290            event_type_prefix: None,
2291            exclude_event_type_prefix: None,
2292            payload_filter: None,
2293        })?;
2294
2295        // If no events and no snapshot, entity not found
2296        if events.is_empty() && since_timestamp.is_none() {
2297            return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
2298        }
2299
2300        // Merge events on top of snapshot (or from scratch if no snapshot)
2301        let mut merged_state = merged_state;
2302        for event in &events {
2303            if let serde_json::Value::Object(ref mut state_map) = merged_state
2304                && let serde_json::Value::Object(ref payload_map) = event.payload
2305            {
2306                for (key, value) in payload_map {
2307                    state_map.insert(key.clone(), value.clone());
2308                }
2309            }
2310        }
2311
2312        // Wrap with metadata
2313        let state = serde_json::json!({
2314            "entity_id": entity_id,
2315            "last_updated": events.last().map(|e| e.timestamp),
2316            "event_count": events.len(),
2317            "as_of": as_of,
2318            "current_state": merged_state,
2319            "history": events.iter().map(|e| {
2320                serde_json::json!({
2321                    "event_id": e.id,
2322                    "type": e.event_type,
2323                    "timestamp": e.timestamp,
2324                    "payload": e.payload
2325                })
2326            }).collect::<Vec<_>>()
2327        });
2328
2329        Ok(state)
2330    }
2331
2332    /// Get snapshot from projection (faster than reconstructing)
2333    pub fn get_snapshot(&self, entity_id: &str) -> Result<serde_json::Value> {
2334        let projections = self.projections.read();
2335
2336        if let Some(snapshot_projection) = projections.get_projection("entity_snapshots")
2337            && let Some(state) = snapshot_projection.get_state(entity_id)
2338        {
2339            return Ok(serde_json::json!({
2340                "entity_id": entity_id,
2341                "snapshot": state,
2342                "from_projection": "entity_snapshots"
2343            }));
2344        }
2345
2346        Err(AllSourceError::EntityNotFound(entity_id.to_string()))
2347    }
2348
2349    /// Get statistics about the event store
2350    pub fn stats(&self) -> StoreStats {
2351        let events = self.events.read();
2352        let index_stats = self.index.stats();
2353
2354        StoreStats {
2355            total_events: events.len(),
2356            total_entities: index_stats.total_entities,
2357            total_event_types: index_stats.total_event_types,
2358            total_ingested: *self.total_ingested.read(),
2359        }
2360    }
2361
2362    /// Get all unique streams (entity_ids) in the store
2363    pub fn list_streams(&self) -> Vec<StreamInfo> {
2364        self.index
2365            .get_all_entities()
2366            .into_iter()
2367            .map(|entity_id| {
2368                let event_count = self
2369                    .index
2370                    .get_by_entity(&entity_id)
2371                    .map_or(0, |entries| entries.len());
2372                let last_event_at = self
2373                    .index
2374                    .get_by_entity(&entity_id)
2375                    .and_then(|entries| entries.last().map(|e| e.timestamp));
2376                StreamInfo {
2377                    stream_id: entity_id,
2378                    event_count,
2379                    last_event_at,
2380                }
2381            })
2382            .collect()
2383    }
2384
2385    /// Get all unique event types in the store
2386    pub fn list_event_types(&self) -> Vec<EventTypeInfo> {
2387        self.index
2388            .get_all_types()
2389            .into_iter()
2390            .map(|event_type| {
2391                let event_count = self
2392                    .index
2393                    .get_by_type(&event_type)
2394                    .map_or(0, |entries| entries.len());
2395                let last_event_at = self
2396                    .index
2397                    .get_by_type(&event_type)
2398                    .and_then(|entries| entries.last().map(|e| e.timestamp));
2399                EventTypeInfo {
2400                    event_type,
2401                    event_count,
2402                    last_event_at,
2403                }
2404            })
2405            .collect()
2406    }
2407
2408    // Tenant-scoped variants. The entity_index / type_index are GLOBAL (no tenant
2409    // dimension), so `list_streams` / `list_event_types` above span every tenant —
2410    // wrong + a cross-tenant spill on the per-tenant dashboard. These filter by the
2411    // event's tenant_id so the dashboard's "your" streams / event types / totals
2412    // reflect only the caller's tenant.
2413
2414    /// Distinct entities (+ per-entity event count) for ONE tenant.
2415    pub fn list_streams_for_tenant(&self, tenant_id: &str) -> Vec<StreamInfo> {
2416        let _ = self.ensure_tenant_loaded(tenant_id);
2417        let events = self.events.read();
2418        let mut by_entity: std::collections::HashMap<&str, (usize, chrono::DateTime<chrono::Utc>)> =
2419            std::collections::HashMap::new();
2420        for ev in events.iter() {
2421            if ev.tenant_id_str() != tenant_id {
2422                continue;
2423            }
2424            let e = by_entity
2425                .entry(ev.entity_id_str())
2426                .or_insert((0, ev.timestamp));
2427            e.0 += 1;
2428            if ev.timestamp > e.1 {
2429                e.1 = ev.timestamp;
2430            }
2431        }
2432        by_entity
2433            .into_iter()
2434            .map(|(entity_id, (count, last))| StreamInfo {
2435                stream_id: entity_id.to_string(),
2436                event_count: count,
2437                last_event_at: Some(last),
2438            })
2439            .collect()
2440    }
2441
2442    /// Distinct event types (+ per-type event count) for ONE tenant.
2443    pub fn list_event_types_for_tenant(&self, tenant_id: &str) -> Vec<EventTypeInfo> {
2444        let _ = self.ensure_tenant_loaded(tenant_id);
2445        let events = self.events.read();
2446        let mut by_type: std::collections::HashMap<&str, (usize, chrono::DateTime<chrono::Utc>)> =
2447            std::collections::HashMap::new();
2448        for ev in events.iter() {
2449            if ev.tenant_id_str() != tenant_id {
2450                continue;
2451            }
2452            let e = by_type
2453                .entry(ev.event_type_str())
2454                .or_insert((0, ev.timestamp));
2455            e.0 += 1;
2456            if ev.timestamp > e.1 {
2457                e.1 = ev.timestamp;
2458            }
2459        }
2460        by_type
2461            .into_iter()
2462            .map(|(event_type, (count, last))| EventTypeInfo {
2463                event_type: event_type.to_string(),
2464                event_count: count,
2465                last_event_at: Some(last),
2466            })
2467            .collect()
2468    }
2469
2470    /// Event-store statistics for ONE tenant.
2471    ///
2472    /// `stats()` above is global (`events.len()` + the global index), so it must
2473    /// never be served to a tenant — it would leak whole-fleet totals. This is
2474    /// the scoped equivalent the gateway exposes as `GET /api/v1/stats` with an
2475    /// auth-derived `tenant_id` (#230).
2476    ///
2477    /// Returns the same four counters as `StoreStats`, plus the per-type census
2478    /// and time range that callers otherwise reconstruct by paging events.
2479    pub fn stats_for_tenant(&self, tenant_id: &str) -> TenantStoreStats {
2480        let _ = self.ensure_tenant_loaded(tenant_id);
2481        let events = self.events.read();
2482
2483        let mut entities: std::collections::HashSet<&str> = std::collections::HashSet::new();
2484        let mut census: std::collections::HashMap<&str, usize> = std::collections::HashMap::new();
2485        let mut total_events = 0usize;
2486        let mut oldest: Option<chrono::DateTime<chrono::Utc>> = None;
2487        let mut newest: Option<chrono::DateTime<chrono::Utc>> = None;
2488
2489        for ev in events.iter() {
2490            if ev.tenant_id_str() != tenant_id {
2491                continue;
2492            }
2493
2494            total_events += 1;
2495            entities.insert(ev.entity_id_str());
2496            *census.entry(ev.event_type_str()).or_insert(0) += 1;
2497
2498            let ts = ev.timestamp;
2499            if oldest.is_none_or(|o| ts < o) {
2500                oldest = Some(ts);
2501            }
2502            if newest.is_none_or(|n| ts > n) {
2503                newest = Some(ts);
2504            }
2505        }
2506
2507        TenantStoreStats {
2508            total_events,
2509            total_entities: entities.len(),
2510            total_event_types: census.len(),
2511            // Scoped equivalent of `total_ingested`: that counter is a global
2512            // process-lifetime tally with no tenant dimension, so reporting it
2513            // here would leak. The tenant's own total is the honest answer.
2514            total_ingested: total_events as u64,
2515            event_types: census
2516                .into_iter()
2517                .map(|(k, v)| (k.to_string(), v))
2518                .collect(),
2519            oldest_event: oldest,
2520            newest_event: newest,
2521        }
2522    }
2523
2524    /// Reconstruct entity state for ONE tenant.
2525    ///
2526    /// Deliberately does NOT use the snapshot fast path that
2527    /// `reconstruct_state` takes: `Snapshot` carries no `tenant_id` and the
2528    /// snapshot manager is keyed by `entity_id` alone, so seeding from a
2529    /// snapshot could fold another tenant's state into this answer when two
2530    /// tenants share an entity_id. Folding tenant-filtered events is slower and
2531    /// fails closed (#230).
2532    pub fn reconstruct_state_for_tenant(
2533        &self,
2534        entity_id: &str,
2535        as_of: Option<DateTime<Utc>>,
2536        tenant_id: &str,
2537    ) -> Result<serde_json::Value> {
2538        let events = self.query(&QueryEventsRequest {
2539            entity_id: Some(entity_id.to_string()),
2540            event_type: None,
2541            tenant_id: Some(tenant_id.to_string()),
2542            as_of,
2543            since: None,
2544            until: None,
2545            limit: None,
2546            event_type_prefix: None,
2547            exclude_event_type_prefix: None,
2548            payload_filter: None,
2549        })?;
2550
2551        if events.is_empty() {
2552            return Err(AllSourceError::EntityNotFound(entity_id.to_string()));
2553        }
2554
2555        let mut merged_state = serde_json::json!({});
2556        for event in &events {
2557            if let serde_json::Value::Object(ref mut state_map) = merged_state
2558                && let serde_json::Value::Object(ref payload_map) = event.payload
2559            {
2560                for (key, value) in payload_map {
2561                    state_map.insert(key.clone(), value.clone());
2562                }
2563            }
2564        }
2565
2566        Ok(serde_json::json!({
2567            "entity_id": entity_id,
2568            "last_updated": events.last().map(|e| e.timestamp),
2569            "event_count": events.len(),
2570            "as_of": as_of,
2571            "current_state": merged_state,
2572            "history": events.iter().map(|e| {
2573                serde_json::json!({
2574                    "event_id": e.id,
2575                    "type": e.event_type,
2576                    "timestamp": e.timestamp,
2577                    "payload": e.payload
2578                })
2579            }).collect::<Vec<_>>()
2580        }))
2581    }
2582
2583    /// Attach a broadcast sender to the WAL for replication.
2584    ///
2585    /// Thread-safe: can be called through `Arc<EventStore>` at runtime.
2586    /// Used during initial setup and during follower → leader promotion.
2587    /// When set, every WAL append publishes the entry to the broadcast
2588    /// channel so the WAL shipper can stream it to followers.
2589    pub fn enable_wal_replication(
2590        &self,
2591        tx: tokio::sync::broadcast::Sender<crate::infrastructure::persistence::wal::WALEntry>,
2592    ) {
2593        if let Some(ref wal_arc) = self.wal {
2594            wal_arc.set_replication_tx(tx);
2595            tracing::info!("WAL replication broadcast enabled");
2596        } else {
2597            tracing::warn!("Cannot enable WAL replication: WAL is not configured");
2598        }
2599    }
2600
2601    /// Get a reference to the WAL (if configured).
2602    /// Used by the replication catch-up protocol to determine oldest available offset.
2603    pub fn wal(&self) -> Option<&Arc<WriteAheadLog>> {
2604        self.wal.as_ref()
2605    }
2606
2607    /// Get a reference to the Parquet storage (if configured).
2608    /// Used by the replication catch-up protocol to stream snapshot files to followers.
2609    pub fn parquet_storage(&self) -> Option<&Arc<RwLock<ParquetStorage>>> {
2610        self.storage.as_ref()
2611    }
2612}
2613
2614/// Configuration for EventStore
2615#[derive(Debug, Clone, Default)]
2616pub struct EventStoreConfig {
2617    /// Optional directory for persistent Parquet storage (v0.2 feature)
2618    pub storage_dir: Option<PathBuf>,
2619
2620    /// Snapshot configuration (v0.2 feature)
2621    pub snapshot_config: SnapshotConfig,
2622
2623    /// Optional directory for WAL (Write-Ahead Log) (v0.2 feature)
2624    pub wal_dir: Option<PathBuf>,
2625
2626    /// WAL configuration (v0.2 feature)
2627    pub wal_config: WALConfig,
2628
2629    /// Compaction configuration (v0.2 feature)
2630    pub compaction_config: CompactionConfig,
2631
2632    /// Schema registry configuration (v0.5 feature)
2633    pub schema_registry_config: SchemaRegistryConfig,
2634
2635    /// Optional directory for system metadata storage (dogfood feature).
2636    /// When set, operational metadata (tenants, config, audit) is stored
2637    /// using AllSource's own event store rather than an external database.
2638    /// Defaults to `{storage_dir}/__system/` when storage_dir is set.
2639    pub system_data_dir: Option<PathBuf>,
2640
2641    /// Name of the default tenant to auto-create on first boot.
2642    pub bootstrap_tenant: Option<String>,
2643
2644    /// In-memory cache budget in bytes (Step 3). When the resident
2645    /// total exceeds this after a load, the LRU tenant is evicted
2646    /// until the cache fits. `None` (the default in tests) disables
2647    /// the budget — every loaded tenant stays resident. Production
2648    /// reads this from the `ALLSOURCE_CACHE_BYTES` env var; see
2649    /// `from_env`.
2650    pub cache_byte_budget: Option<u64>,
2651
2652    /// Cadence of the runtime checkpoint loop, in seconds (Step 6).
2653    /// Each tick flushes pending Parquet batches and, on success,
2654    /// truncates the WAL up through the checkpoint. This bounds
2655    /// dirty-restart replay time to one interval of writes
2656    /// regardless of total dataset size.
2657    ///
2658    /// `None` disables the loop — the WAL still grows but is only
2659    /// truncated at boot, which is the pre-Step-6 behavior. Tests
2660    /// default to `None`; production reads
2661    /// `ALLSOURCE_CHECKPOINT_INTERVAL_SECONDS` (default 60s) via
2662    /// `from_env_vars`.
2663    pub checkpoint_interval_secs: Option<u64>,
2664
2665    /// Open the store read-only (replica mode). See `EventStore::read_only`.
2666    /// Defaults to `false` (read-write owner). Set by Prime when it fails to
2667    /// acquire the exclusive data-dir lock because another process owns it.
2668    pub read_only: bool,
2669}
2670
2671impl EventStoreConfig {
2672    /// Create config with persistent storage enabled
2673    pub fn with_persistence(storage_dir: impl Into<PathBuf>) -> Self {
2674        Self {
2675            storage_dir: Some(storage_dir.into()),
2676            ..Self::default()
2677        }
2678    }
2679
2680    /// Create config with custom snapshot settings
2681    pub fn with_snapshots(snapshot_config: SnapshotConfig) -> Self {
2682        Self {
2683            snapshot_config,
2684            ..Self::default()
2685        }
2686    }
2687
2688    /// Create config with WAL enabled
2689    pub fn with_wal(wal_dir: impl Into<PathBuf>, wal_config: WALConfig) -> Self {
2690        Self {
2691            wal_dir: Some(wal_dir.into()),
2692            wal_config,
2693            ..Self::default()
2694        }
2695    }
2696
2697    /// Create config with both persistence and snapshots
2698    pub fn with_all(storage_dir: impl Into<PathBuf>, snapshot_config: SnapshotConfig) -> Self {
2699        Self {
2700            storage_dir: Some(storage_dir.into()),
2701            snapshot_config,
2702            ..Self::default()
2703        }
2704    }
2705
2706    /// Create production config with all features enabled
2707    pub fn production(
2708        storage_dir: impl Into<PathBuf>,
2709        wal_dir: impl Into<PathBuf>,
2710        snapshot_config: SnapshotConfig,
2711        wal_config: WALConfig,
2712        compaction_config: CompactionConfig,
2713    ) -> Self {
2714        let storage_dir = storage_dir.into();
2715        let system_data_dir = storage_dir.join("__system");
2716        Self {
2717            storage_dir: Some(storage_dir),
2718            snapshot_config,
2719            wal_dir: Some(wal_dir.into()),
2720            wal_config,
2721            compaction_config,
2722            system_data_dir: Some(system_data_dir),
2723            ..Self::default()
2724        }
2725    }
2726
2727    /// Resolve the effective system data directory.
2728    ///
2729    /// If explicitly set, returns that. Otherwise, derives from storage_dir.
2730    /// Returns None if neither is configured (in-memory mode).
2731    pub fn effective_system_data_dir(&self) -> Option<PathBuf> {
2732        self.system_data_dir
2733            .clone()
2734            .or_else(|| self.storage_dir.as_ref().map(|d| d.join("__system")))
2735    }
2736
2737    /// Build config from environment variables.
2738    ///
2739    /// Reads `ALLSOURCE_DATA_DIR`, `ALLSOURCE_STORAGE_DIR`, `ALLSOURCE_WAL_DIR`,
2740    /// and `ALLSOURCE_WAL_ENABLED` to determine persistence mode.
2741    ///
2742    /// Returns `(config, description)` where description is a human-readable
2743    /// summary of the persistence mode for logging.
2744    pub fn from_env() -> (Self, &'static str) {
2745        Self::from_env_vars(
2746            std::env::var("ALLSOURCE_DATA_DIR")
2747                .ok()
2748                .filter(|s| !s.is_empty()),
2749            std::env::var("ALLSOURCE_STORAGE_DIR")
2750                .ok()
2751                .filter(|s| !s.is_empty()),
2752            std::env::var("ALLSOURCE_WAL_DIR")
2753                .ok()
2754                .filter(|s| !s.is_empty()),
2755            std::env::var("ALLSOURCE_WAL_ENABLED").ok(),
2756            std::env::var("ALLSOURCE_CACHE_BYTES").ok(),
2757            std::env::var("ALLSOURCE_SNAPSHOT_INTERVAL_SECONDS").ok(),
2758            std::env::var("ALLSOURCE_RETENTION_SYSTEM_DAYS").ok(),
2759            std::env::var("ALLSOURCE_CHECKPOINT_INTERVAL_SECONDS").ok(),
2760        )
2761    }
2762
2763    /// Build config from explicit env-var values (testable without mutating process env).
2764    pub fn from_env_vars(
2765        data_dir: Option<String>,
2766        explicit_storage_dir: Option<String>,
2767        explicit_wal_dir: Option<String>,
2768        wal_enabled_var: Option<String>,
2769        cache_bytes_var: Option<String>,
2770        snapshot_interval_var: Option<String>,
2771        retention_system_days_var: Option<String>,
2772        checkpoint_interval_var: Option<String>,
2773    ) -> (Self, &'static str) {
2774        let data_dir = data_dir.filter(|s| !s.is_empty());
2775        let storage_dir = explicit_storage_dir
2776            .filter(|s| !s.is_empty())
2777            .or_else(|| data_dir.as_ref().map(|d| format!("{d}/storage")));
2778        let wal_dir = explicit_wal_dir
2779            .filter(|s| !s.is_empty())
2780            .or_else(|| data_dir.as_ref().map(|d| format!("{d}/wal")));
2781        let wal_enabled = wal_enabled_var.is_none_or(|v| v == "true");
2782        // ALLSOURCE_CACHE_BYTES: parse decimal bytes. Unparseable
2783        // input is logged and ignored rather than failing boot —
2784        // the unbounded fallback is safe (worst case is the
2785        // original pre-Step-3 behavior).
2786        let cache_byte_budget =
2787            cache_bytes_var
2788                .filter(|s| !s.is_empty())
2789                .and_then(|s| match s.parse::<u64>() {
2790                    Ok(v) => Some(v),
2791                    Err(e) => {
2792                        tracing::warn!(
2793                            "ALLSOURCE_CACHE_BYTES={s:?} could not be parsed as u64: {e}; \
2794                         cache budget disabled"
2795                        );
2796                        None
2797                    }
2798                });
2799        let compaction_config =
2800            CompactionConfig::from_env_vars(snapshot_interval_var, retention_system_days_var);
2801
2802        // ALLSOURCE_CHECKPOINT_INTERVAL_SECONDS: parse decimal seconds. The
2803        // default (60s) only applies when WAL is enabled — there's no
2804        // checkpoint loop to run otherwise. Unparseable input is logged
2805        // and falls back to the default rather than failing boot.
2806        let checkpoint_interval_secs = if wal_enabled {
2807            checkpoint_interval_var
2808                .filter(|s| !s.is_empty())
2809                .map(|s| match s.parse::<u64>() {
2810                    Ok(v) => v,
2811                    Err(e) => {
2812                        tracing::warn!(
2813                            "ALLSOURCE_CHECKPOINT_INTERVAL_SECONDS={s:?} could not be parsed as \
2814                             u64: {e}; falling back to default 60s"
2815                        );
2816                        60
2817                    }
2818                })
2819                .or(Some(60))
2820        } else {
2821            None
2822        };
2823
2824        let mut config = match (&storage_dir, &wal_dir) {
2825            (Some(sd), Some(wd)) if wal_enabled => Self::production(
2826                sd,
2827                wd,
2828                SnapshotConfig::default(),
2829                WALConfig::default(),
2830                compaction_config,
2831            ),
2832            (Some(sd), _) => Self::with_persistence(sd),
2833            (_, Some(wd)) if wal_enabled => Self::with_wal(wd, WALConfig::default()),
2834            _ => Self::default(),
2835        };
2836        config.cache_byte_budget = cache_byte_budget;
2837        config.checkpoint_interval_secs = checkpoint_interval_secs;
2838
2839        let mode = match (&storage_dir, &wal_dir) {
2840            (Some(_), Some(_)) if wal_enabled => "wal+parquet",
2841            (Some(_), _) => "parquet-only",
2842            (_, Some(_)) if wal_enabled => "wal-only",
2843            _ => "in-memory",
2844        };
2845        (config, mode)
2846    }
2847}
2848
2849#[derive(Debug, serde::Serialize)]
2850pub struct StoreStats {
2851    pub total_events: usize,
2852    pub total_entities: usize,
2853    pub total_event_types: usize,
2854    pub total_ingested: u64,
2855}
2856
2857/// Tenant-scoped event-store statistics (#230).
2858///
2859/// Mirrors `StoreStats`' counters so existing consumers keep working, and adds
2860/// the per-type census and time range that clients previously derived by paging
2861/// the whole event stream.
2862#[derive(Debug, Clone, serde::Serialize)]
2863pub struct TenantStoreStats {
2864    pub total_events: usize,
2865    pub total_entities: usize,
2866    pub total_event_types: usize,
2867    pub total_ingested: u64,
2868    /// event_type -> count, for this tenant only.
2869    pub event_types: std::collections::HashMap<String, usize>,
2870    pub oldest_event: Option<chrono::DateTime<chrono::Utc>>,
2871    pub newest_event: Option<chrono::DateTime<chrono::Utc>>,
2872}
2873
2874/// Information about a stream (entity_id)
2875#[derive(Debug, Clone, serde::Serialize)]
2876pub struct StreamInfo {
2877    /// The stream identifier (entity_id)
2878    pub stream_id: String,
2879    /// Total number of events in this stream
2880    pub event_count: usize,
2881    /// Timestamp of the last event in this stream
2882    pub last_event_at: Option<chrono::DateTime<chrono::Utc>>,
2883}
2884
2885/// Information about an event type
2886#[derive(Debug, Clone, serde::Serialize)]
2887pub struct EventTypeInfo {
2888    /// The event type name
2889    pub event_type: String,
2890    /// Total number of events of this type
2891    pub event_count: usize,
2892    /// Timestamp of the last event of this type
2893    pub last_event_at: Option<chrono::DateTime<chrono::Utc>>,
2894}
2895
2896impl Default for EventStore {
2897    fn default() -> Self {
2898        Self::new()
2899    }
2900}
2901
2902#[cfg(test)]
2903mod tests {
2904    use super::*;
2905    use crate::domain::entities::Event;
2906    use tempfile::TempDir;
2907
2908    /// Recursively walk `dir` looking for `*.parquet` files.
2909    /// Tests that pre-date Step 1's tenant-partitioned layout used a
2910    /// flat `read_dir` here; after the move to <root>/<tenant>/<yyyy-mm>/
2911    /// they need to walk subdirectories.
2912    fn find_parquet_files(dir: &std::path::Path) -> Vec<std::path::PathBuf> {
2913        let mut out = Vec::new();
2914        let mut stack = vec![dir.to_path_buf()];
2915        while let Some(d) = stack.pop() {
2916            let Ok(entries) = std::fs::read_dir(&d) else {
2917                continue;
2918            };
2919            for e in entries.flatten() {
2920                let p = e.path();
2921                if p.is_dir() {
2922                    stack.push(p);
2923                } else if p.extension().and_then(|s| s.to_str()) == Some("parquet") {
2924                    out.push(p);
2925                }
2926            }
2927        }
2928        out
2929    }
2930
2931    fn create_test_event(entity_id: &str, event_type: &str) -> Event {
2932        Event::from_strings(
2933            event_type.to_string(),
2934            entity_id.to_string(),
2935            "default".to_string(),
2936            serde_json::json!({"name": "Test", "value": 42}),
2937            None,
2938        )
2939        .unwrap()
2940    }
2941
2942    fn create_test_event_with_payload(
2943        entity_id: &str,
2944        event_type: &str,
2945        payload: serde_json::Value,
2946    ) -> Event {
2947        Event::from_strings(
2948            event_type.to_string(),
2949            entity_id.to_string(),
2950            "default".to_string(),
2951            payload,
2952            None,
2953        )
2954        .unwrap()
2955    }
2956
2957    #[test]
2958    fn test_event_store_new() {
2959        let store = EventStore::new();
2960        assert_eq!(store.stats().total_events, 0);
2961        assert_eq!(store.stats().total_entities, 0);
2962    }
2963
2964    // -----------------------------------------------------------------
2965    // Step 2: ensure_tenant_loaded smoke tests. The full
2966    // cold-boot/lazy-hydrate paths land in commit #2 (skip boot
2967    // load) and commit #4 (integration test).
2968    // -----------------------------------------------------------------
2969
2970    #[test]
2971    fn test_ensure_tenant_loaded_no_storage_is_a_noop() {
2972        // An in-memory-only store (no ParquetStorage configured) has
2973        // nothing to hydrate. The method must succeed and mark the
2974        // tenant loaded so subsequent calls hit the fast path.
2975        let store = EventStore::new();
2976        assert!(!store.is_tenant_loaded("alice"));
2977        store.ensure_tenant_loaded("alice").unwrap();
2978        assert!(store.is_tenant_loaded("alice"));
2979        // Other tenants stay cold — the call is per-tenant.
2980        assert!(!store.is_tenant_loaded("bob"));
2981    }
2982
2983    #[test]
2984    fn test_ensure_tenant_loaded_warm_path_is_idempotent() {
2985        let store = EventStore::new();
2986        store.ensure_tenant_loaded("alice").unwrap();
2987        // Second call hits the DashMap fast path and returns Ok.
2988        store.ensure_tenant_loaded("alice").unwrap();
2989    }
2990
2991    #[test]
2992    fn test_ensure_tenant_loaded_rejects_unsafe_tenant_id() {
2993        // With persistence configured, the call has to walk a
2994        // tenant subtree, so the path-safety whitelist applies.
2995        // The error must propagate; the tenant must NOT be marked
2996        // loaded (otherwise an attacker probing path-traversal
2997        // strings could spam the loaded-set with junk).
2998        let temp_dir = TempDir::new().unwrap();
2999        let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3000        for unsafe_tid in ["..", "a/b", "a\\b", ""] {
3001            let result = store.ensure_tenant_loaded(unsafe_tid);
3002            assert!(
3003                result.is_err(),
3004                "tenant_id {unsafe_tid:?} should have been rejected"
3005            );
3006            assert!(
3007                !store.is_tenant_loaded(unsafe_tid),
3008                "rejected tenant {unsafe_tid:?} must not be marked loaded"
3009            );
3010        }
3011    }
3012
3013    #[test]
3014    fn test_ensure_tenant_loaded_no_subtree_marks_loaded_with_zero_events() {
3015        // A tenant that has no on-disk data (fresh tenant, never
3016        // persisted) must still succeed — load_events_for_tenant
3017        // returns empty, ensure_tenant_loaded marks it loaded so we
3018        // don't re-walk the empty subtree on every query.
3019        let temp_dir = TempDir::new().unwrap();
3020        let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3021        assert!(!store.is_tenant_loaded("never-existed"));
3022        store.ensure_tenant_loaded("never-existed").unwrap();
3023        assert!(store.is_tenant_loaded("never-existed"));
3024    }
3025
3026    #[test]
3027    fn test_evict_tenant_drops_events_and_resets_bytes() {
3028        // After eviction, the tenant's events are gone from memory,
3029        // its byte counter is reset, and is_tenant_loaded returns
3030        // false. Other tenants are untouched.
3031        let temp_dir = TempDir::new().unwrap();
3032        let storage_dir = temp_dir.path().to_path_buf();
3033
3034        {
3035            let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3036            for i in 0..3 {
3037                store
3038                    .ingest(
3039                        &Event::from_strings(
3040                            "test.event".to_string(),
3041                            format!("a-{i}"),
3042                            "alice".to_string(),
3043                            serde_json::json!({"i": i}),
3044                            None,
3045                        )
3046                        .unwrap(),
3047                    )
3048                    .unwrap();
3049            }
3050            for i in 0..2 {
3051                store
3052                    .ingest(
3053                        &Event::from_strings(
3054                            "test.event".to_string(),
3055                            format!("b-{i}"),
3056                            "bob".to_string(),
3057                            serde_json::json!({"i": i}),
3058                            None,
3059                        )
3060                        .unwrap(),
3061                    )
3062                    .unwrap();
3063            }
3064            store.flush_storage().unwrap();
3065        }
3066
3067        let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3068        store.ensure_tenant_loaded("alice").unwrap();
3069        store.ensure_tenant_loaded("bob").unwrap();
3070        assert_eq!(store.stats().total_events, 5);
3071        let alice_bytes = store.tenant_resident_bytes("alice");
3072        let bob_bytes = store.tenant_resident_bytes("bob");
3073        assert!(alice_bytes > 0 && bob_bytes > 0);
3074
3075        store.evict_tenant("alice");
3076
3077        assert!(!store.is_tenant_loaded("alice"));
3078        assert!(store.is_tenant_loaded("bob"));
3079        assert_eq!(store.tenant_resident_bytes("alice"), 0);
3080        assert_eq!(store.tenant_resident_bytes("bob"), bob_bytes);
3081        assert_eq!(store.stats().total_events, 2, "only bob's 2 events remain");
3082    }
3083
3084    #[test]
3085    fn test_evict_tenant_then_query_re_loads_from_disk() {
3086        // The transparent re-load behavior the bead's AC #5 calls
3087        // out: evict, then query the same tenant — its data comes
3088        // back via ensure_tenant_loaded, sourced from Parquet.
3089        let temp_dir = TempDir::new().unwrap();
3090        let storage_dir = temp_dir.path().to_path_buf();
3091
3092        {
3093            let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3094            for i in 0..4 {
3095                store
3096                    .ingest(
3097                        &Event::from_strings(
3098                            "test.event".to_string(),
3099                            format!("a-{i}"),
3100                            "alice".to_string(),
3101                            serde_json::json!({"i": i}),
3102                            None,
3103                        )
3104                        .unwrap(),
3105                    )
3106                    .unwrap();
3107            }
3108            store.flush_storage().unwrap();
3109        }
3110
3111        let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3112        store.ensure_tenant_loaded("alice").unwrap();
3113        store.evict_tenant("alice");
3114        assert_eq!(store.stats().total_events, 0);
3115
3116        // Query — re-load happens transparently.
3117        let results = store
3118            .query(&QueryEventsRequest {
3119                entity_id: None,
3120                event_type: None,
3121                tenant_id: Some("alice".to_string()),
3122                as_of: None,
3123                since: None,
3124                until: None,
3125                limit: None,
3126                event_type_prefix: None,
3127                exclude_event_type_prefix: None,
3128                payload_filter: None,
3129            })
3130            .unwrap();
3131        assert_eq!(results.len(), 4);
3132        assert!(store.is_tenant_loaded("alice"));
3133    }
3134
3135    #[test]
3136    fn test_evict_tenant_rebuilds_index_with_new_offsets() {
3137        // After eviction, the events Vec is compacted. The index
3138        // must be rebuilt against the new offsets — otherwise
3139        // queries return stale or wrong events. This test checks
3140        // index correctness end-to-end via a query for the
3141        // surviving tenant after the evicted tenant's events are
3142        // gone.
3143        let temp_dir = TempDir::new().unwrap();
3144        let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3145
3146        // Interleave: alice, bob, alice, bob, alice. After
3147        // evicting alice, the events Vec compacts to [bob, bob]
3148        // and the index must reflect the new layout.
3149        for i in 0..3 {
3150            store
3151                .ingest(
3152                    &Event::from_strings(
3153                        "test.event".to_string(),
3154                        format!("a-{i}"),
3155                        "alice".to_string(),
3156                        serde_json::json!({"i": i}),
3157                        None,
3158                    )
3159                    .unwrap(),
3160                )
3161                .unwrap();
3162            if i < 2 {
3163                store
3164                    .ingest(
3165                        &Event::from_strings(
3166                            "test.event".to_string(),
3167                            format!("b-{i}"),
3168                            "bob".to_string(),
3169                            serde_json::json!({"i": i}),
3170                            None,
3171                        )
3172                        .unwrap(),
3173                    )
3174                    .unwrap();
3175            }
3176        }
3177        // Mark both as loaded for accurate eviction bookkeeping.
3178        store.tenant_loader.mark_loaded("alice");
3179        store.tenant_loader.mark_loaded("bob");
3180
3181        store.evict_tenant("alice");
3182
3183        let bob_results = store
3184            .query(&QueryEventsRequest {
3185                entity_id: None,
3186                event_type: None,
3187                tenant_id: Some("bob".to_string()),
3188                as_of: None,
3189                since: None,
3190                until: None,
3191                limit: None,
3192                event_type_prefix: None,
3193                exclude_event_type_prefix: None,
3194                payload_filter: None,
3195            })
3196            .unwrap();
3197        assert_eq!(bob_results.len(), 2);
3198        for e in &bob_results {
3199            assert_eq!(e.tenant_id_str(), "bob");
3200        }
3201    }
3202
3203    #[test]
3204    fn test_budget_eviction_keeps_resident_set_bounded() {
3205        // Configure a tiny budget. Load three tenants in sequence;
3206        // the third load must evict the LRU tenant, keeping the
3207        // resident set under (or near) the budget.
3208        let temp_dir = TempDir::new().unwrap();
3209        let storage_dir = temp_dir.path().to_path_buf();
3210
3211        // Persist 5 events per tenant with ~1 KiB payloads. Each
3212        // tenant ends up at ~5 KiB + overhead.
3213        let big_payload = serde_json::json!({"data": "x".repeat(1000)});
3214        {
3215            let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3216            for tenant in ["alice", "bob", "carol"] {
3217                for i in 0..5 {
3218                    store
3219                        .ingest(
3220                            &Event::from_strings(
3221                                "test.event".to_string(),
3222                                format!("{tenant}-{i}"),
3223                                tenant.to_string(),
3224                                big_payload.clone(),
3225                                None,
3226                            )
3227                            .unwrap(),
3228                        )
3229                        .unwrap();
3230                }
3231            }
3232            store.flush_storage().unwrap();
3233        }
3234
3235        // Budget = 12 KiB. Two tenants (~6 KiB each = ~12 KiB) is
3236        // tight; loading a third must evict.
3237        let mut config = EventStoreConfig::with_persistence(&storage_dir);
3238        config.cache_byte_budget = Some(12_000);
3239        let store = EventStore::with_config(config);
3240
3241        // Load alice — under budget, no eviction.
3242        store.ensure_tenant_loaded("alice").unwrap();
3243        assert!(store.is_tenant_loaded("alice"));
3244
3245        // Touch alice and immediately load bob. Bob is the
3246        // freshly-loaded one, so bob is excluded from eviction.
3247        // Alice is the next-oldest. After the load, total may
3248        // exceed budget — if so, evict alice.
3249        store.tenant_loader.touch("alice");
3250        std::thread::sleep(std::time::Duration::from_millis(10));
3251        store.ensure_tenant_loaded("bob").unwrap();
3252        assert!(store.is_tenant_loaded("bob"));
3253
3254        // Touch bob, load carol. Carol is freshly-loaded; the LRU
3255        // candidate is the older of {alice, bob} — alice (since
3256        // bob was just touched).
3257        store.tenant_loader.touch("bob");
3258        std::thread::sleep(std::time::Duration::from_millis(10));
3259        store.ensure_tenant_loaded("carol").unwrap();
3260        assert!(store.is_tenant_loaded("carol"));
3261
3262        // After all loads, the cache must respect the budget OR
3263        // (if a single tenant alone exceeds it) we should at most
3264        // hold the just-loaded tenant. The test budget is small
3265        // enough that we expect at least one eviction.
3266        let resident = store.cache_resident_bytes();
3267        let budget = 12_000u64;
3268
3269        // Either we're within the budget, or only the freshly-loaded
3270        // tenant is left (the "single oversized tenant" fallback).
3271        if resident > budget {
3272            let loaded_count = ["alice", "bob", "carol"]
3273                .iter()
3274                .filter(|t| store.is_tenant_loaded(t))
3275                .count();
3276            assert_eq!(
3277                loaded_count, 1,
3278                "over budget but more than one tenant loaded — eviction policy didn't fire"
3279            );
3280        }
3281
3282        // Carol must still be loaded — it's the most recent and
3283        // never picked as a victim.
3284        assert!(store.is_tenant_loaded("carol"));
3285    }
3286
3287    #[test]
3288    fn test_query_after_eviction_re_loads_transparently() {
3289        // The end-to-end shape of AC #5: query → evict → query
3290        // again returns the right data.
3291        let temp_dir = TempDir::new().unwrap();
3292        let storage_dir = temp_dir.path().to_path_buf();
3293
3294        let big_payload = serde_json::json!({"data": "x".repeat(2000)});
3295        {
3296            let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3297            for tenant in ["alice", "bob"] {
3298                for i in 0..3 {
3299                    store
3300                        .ingest(
3301                            &Event::from_strings(
3302                                "test.event".to_string(),
3303                                format!("{tenant}-{i}"),
3304                                tenant.to_string(),
3305                                big_payload.clone(),
3306                                None,
3307                            )
3308                            .unwrap(),
3309                        )
3310                        .unwrap();
3311                }
3312            }
3313            store.flush_storage().unwrap();
3314        }
3315
3316        // Budget = 5 KiB — one tenant fits, two don't.
3317        let mut config = EventStoreConfig::with_persistence(&storage_dir);
3318        config.cache_byte_budget = Some(5_000);
3319        let store = EventStore::with_config(config);
3320
3321        // Query alice — sized at ~6 KiB, so over budget but no
3322        // peer to evict; alice stays as the single-oversized-tenant
3323        // case.
3324        let alice_first = store
3325            .query(&QueryEventsRequest {
3326                entity_id: None,
3327                event_type: None,
3328                tenant_id: Some("alice".to_string()),
3329                as_of: None,
3330                since: None,
3331                until: None,
3332                limit: None,
3333                event_type_prefix: None,
3334                exclude_event_type_prefix: None,
3335                payload_filter: None,
3336            })
3337            .unwrap();
3338        assert_eq!(alice_first.len(), 3);
3339
3340        // Sleep to make alice older than bob in the LRU ordering.
3341        std::thread::sleep(std::time::Duration::from_millis(15));
3342        // Query bob — alice will get evicted.
3343        let _bob = store
3344            .query(&QueryEventsRequest {
3345                entity_id: None,
3346                event_type: None,
3347                tenant_id: Some("bob".to_string()),
3348                as_of: None,
3349                since: None,
3350                until: None,
3351                limit: None,
3352                event_type_prefix: None,
3353                exclude_event_type_prefix: None,
3354                payload_filter: None,
3355            })
3356            .unwrap();
3357        assert!(
3358            !store.is_tenant_loaded("alice"),
3359            "alice should have been evicted"
3360        );
3361
3362        // Re-query alice — must transparently re-load.
3363        let alice_second = store
3364            .query(&QueryEventsRequest {
3365                entity_id: None,
3366                event_type: None,
3367                tenant_id: Some("alice".to_string()),
3368                as_of: None,
3369                since: None,
3370                until: None,
3371                limit: None,
3372                event_type_prefix: None,
3373                exclude_event_type_prefix: None,
3374                payload_filter: None,
3375            })
3376            .unwrap();
3377        assert_eq!(
3378            alice_second.len(),
3379            3,
3380            "alice's events come back via re-load"
3381        );
3382        assert!(store.is_tenant_loaded("alice"));
3383    }
3384
3385    #[test]
3386    #[cfg(feature = "server")]
3387    fn test_cache_metrics_track_evictions_and_bytes() {
3388        // Smoke test for the Step 3 #4 Prometheus metrics —
3389        // confirms the counter increments on eviction and the
3390        // gauge tracks the resident bytes.
3391        let temp_dir = TempDir::new().unwrap();
3392        let storage_dir = temp_dir.path().to_path_buf();
3393
3394        let big_payload = serde_json::json!({"data": "x".repeat(2000)});
3395        {
3396            let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3397            for tenant in ["alice", "bob"] {
3398                for i in 0..3 {
3399                    store
3400                        .ingest(
3401                            &Event::from_strings(
3402                                "test.event".to_string(),
3403                                format!("{tenant}-{i}"),
3404                                tenant.to_string(),
3405                                big_payload.clone(),
3406                                None,
3407                            )
3408                            .unwrap(),
3409                        )
3410                        .unwrap();
3411                }
3412            }
3413            store.flush_storage().unwrap();
3414        }
3415
3416        let mut config = EventStoreConfig::with_persistence(&storage_dir);
3417        config.cache_byte_budget = Some(5_000); // forces eviction
3418        let store = EventStore::with_config(config);
3419
3420        assert_eq!(store.metrics.cache_evictions_total.get(), 0);
3421        assert_eq!(store.metrics.cache_bytes.get(), 0);
3422
3423        store.ensure_tenant_loaded("alice").unwrap();
3424        // After loading alice, gauge reflects her bytes.
3425        let after_alice = store.metrics.cache_bytes.get();
3426        assert!(after_alice > 0, "gauge should reflect alice's bytes");
3427        // Single oversized tenant — no eviction yet.
3428        assert_eq!(store.metrics.cache_evictions_total.get(), 0);
3429
3430        std::thread::sleep(std::time::Duration::from_millis(10));
3431        store.ensure_tenant_loaded("bob").unwrap();
3432
3433        // Bob's load pushed total over budget; alice (older) was
3434        // evicted. Counter increments.
3435        assert_eq!(
3436            store.metrics.cache_evictions_total.get(),
3437            1,
3438            "exactly one tenant evicted after bob's load"
3439        );
3440        // Gauge now reflects only bob's bytes.
3441        let after_bob = store.metrics.cache_bytes.get();
3442        assert!(after_bob > 0);
3443        assert!(after_bob <= after_alice, "gauge dropped after eviction");
3444    }
3445
3446    #[test]
3447    #[cfg(feature = "server")]
3448    fn test_storage_size_gauge_populated_from_on_disk_bytes() {
3449        // The allsource_storage_size_bytes / _parquet_files_total / _wal_segments_total
3450        // gauges used to be registered but never set (constant 0), so the dashboard's
3451        // "storage" card showed "—". refresh_storage_metrics must populate them from
3452        // the real on-disk footprint (Parquet bytes + WAL segment bytes).
3453        let temp_dir = TempDir::new().unwrap();
3454        // Configure BOTH Parquet persistence and a WAL so the gauge exercises the
3455        // full `parquet_bytes + wal_bytes` summation (production runs with both).
3456        let config = EventStoreConfig {
3457            storage_dir: Some(temp_dir.path().join("parquet")),
3458            wal_dir: Some(temp_dir.path().join("wal")),
3459            ..EventStoreConfig::default()
3460        };
3461        let store = EventStore::with_config(config);
3462
3463        // Gauge starts at 0 before anything is written/refreshed.
3464        assert_eq!(store.metrics.storage_size_bytes.get(), 0);
3465
3466        let payload = serde_json::json!({ "data": "x".repeat(2000) });
3467        for i in 0..10 {
3468            store
3469                .ingest(
3470                    &Event::from_strings(
3471                        "test.event".to_string(),
3472                        format!("entity-{i}"),
3473                        "tenant-a".to_string(),
3474                        payload.clone(),
3475                        None,
3476                    )
3477                    .unwrap(),
3478                )
3479                .unwrap();
3480        }
3481        store.flush_storage().unwrap();
3482
3483        // Populate the gauges from disk (what the boot hook + checkpoint loop call).
3484        store.refresh_storage_metrics_now();
3485
3486        let size = store.metrics.storage_size_bytes.get();
3487        assert!(
3488            size > 0,
3489            "storage_size_bytes must reflect real on-disk bytes, got {size}"
3490        );
3491        assert!(
3492            store.metrics.parquet_files_total.get() >= 1,
3493            "at least one Parquet file should exist after a flush"
3494        );
3495
3496        // WAL segment count is surfaced too: with a WAL configured and writes done,
3497        // there is at least one segment on disk.
3498        assert!(
3499            store.metrics.wal_segments_total.get() >= 1,
3500            "at least one WAL segment should exist after writes"
3501        );
3502
3503        // Sanity: the reported size is the Parquet bytes plus the WAL bytes — i.e.
3504        // it's the real on-disk footprint, not a stale/placeholder value.
3505        let parquet_stats = store.storage.as_ref().unwrap().read().stats().unwrap();
3506        let (wal_bytes, _) = store.wal.as_ref().unwrap().on_disk_stats().unwrap();
3507        assert_eq!(
3508            size as u64,
3509            parquet_stats.total_size_bytes + wal_bytes,
3510            "gauge should equal Parquet bytes ({}) + WAL bytes ({wal_bytes})",
3511            parquet_stats.total_size_bytes
3512        );
3513    }
3514
3515    #[test]
3516    fn test_stress_resident_set_stays_near_budget_under_rolling_queries() {
3517        // Scaled-down version of the bead's stress test: the
3518        // bead's 10 × 50 MB / 100 MB ratio (10× tenants vs
3519        // budget-headroom) preserved at 500 KB / 1 MB to stay
3520        // unit-test-fast. The same correctness property: after
3521        // many rolling queries across more tenants than fit, the
3522        // resident set must stay at-or-near the budget.
3523        let temp_dir = TempDir::new().unwrap();
3524        let storage_dir = temp_dir.path().to_path_buf();
3525
3526        const TENANT_COUNT: usize = 10;
3527        const EVENTS_PER_TENANT: usize = 50;
3528        // Per-event payload ~10 KiB → tenant ~ 500 KiB.
3529        let big_payload = serde_json::json!({"data": "x".repeat(10_000)});
3530
3531        // Persist all tenants. Each ends up at ~500 KiB on disk
3532        // (and roughly the same in memory once loaded).
3533        {
3534            let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3535            for t in 0..TENANT_COUNT {
3536                let tenant = format!("tenant-{t}");
3537                for i in 0..EVENTS_PER_TENANT {
3538                    store
3539                        .ingest(
3540                            &Event::from_strings(
3541                                "test.event".to_string(),
3542                                format!("{tenant}-{i}"),
3543                                tenant.clone(),
3544                                big_payload.clone(),
3545                                None,
3546                            )
3547                            .unwrap(),
3548                        )
3549                        .unwrap();
3550                }
3551            }
3552            store.flush_storage().unwrap();
3553        }
3554
3555        // Budget = 1 MiB → fits ~2 tenants. We're going to query
3556        // all 10, so the LRU policy must hold the resident set
3557        // near 1 MiB across the rolling sequence.
3558        const BUDGET: u64 = 1_048_576;
3559        let mut config = EventStoreConfig::with_persistence(&storage_dir);
3560        config.cache_byte_budget = Some(BUDGET);
3561        let store = EventStore::with_config(config);
3562
3563        // Sweep through tenants in order. Each query loads its
3564        // tenant; if budget is exceeded after the load, an LRU
3565        // eviction fires.
3566        let mut peak_resident: u64 = 0;
3567        for t in 0..TENANT_COUNT {
3568            let tenant = format!("tenant-{t}");
3569            let results = store
3570                .query(&QueryEventsRequest {
3571                    entity_id: None,
3572                    event_type: None,
3573                    tenant_id: Some(tenant.clone()),
3574                    as_of: None,
3575                    since: None,
3576                    until: None,
3577                    limit: None,
3578                    event_type_prefix: None,
3579                    exclude_event_type_prefix: None,
3580                    payload_filter: None,
3581                })
3582                .unwrap();
3583            assert_eq!(
3584                results.len(),
3585                EVENTS_PER_TENANT,
3586                "every per-tenant query must return all of that tenant's events"
3587            );
3588            // Track peak resident bytes seen during the sweep.
3589            let resident = store.cache_resident_bytes();
3590            if resident > peak_resident {
3591                peak_resident = resident;
3592            }
3593        }
3594
3595        let final_resident = store.cache_resident_bytes();
3596
3597        // Tolerance: a tenant's bytes get added before eviction
3598        // fires, so peak transiently exceeds the budget by at
3599        // most one tenant's worth (~500 KiB). The final state
3600        // after the sweep should be well-bounded.
3601        let tolerance = BUDGET; // generous: 2× budget upper bound
3602        assert!(
3603            peak_resident <= BUDGET + tolerance,
3604            "peak resident {peak_resident} exceeds budget {BUDGET} by more than {tolerance} \
3605             — eviction policy not keeping up with the working-set churn"
3606        );
3607        assert!(
3608            final_resident <= BUDGET + tolerance,
3609            "final resident {final_resident} exceeds budget {BUDGET} by more than {tolerance}"
3610        );
3611
3612        // The most-recently-queried tenant must still be loaded
3613        // (it was just touched).
3614        let last_tenant = format!("tenant-{}", TENANT_COUNT - 1);
3615        assert!(
3616            store.is_tenant_loaded(&last_tenant),
3617            "the most-recent tenant must remain loaded after the sweep"
3618        );
3619
3620        // At least some tenants must have been evicted — otherwise
3621        // the budget didn't fire.
3622        let still_loaded = (0..TENANT_COUNT)
3623            .filter(|t| store.is_tenant_loaded(&format!("tenant-{t}")))
3624            .count();
3625        assert!(
3626            still_loaded < TENANT_COUNT,
3627            "no tenants evicted ({still_loaded}/{TENANT_COUNT} still loaded) — \
3628             budget enforcement didn't engage"
3629        );
3630    }
3631
3632    #[test]
3633    fn test_evict_tenant_when_not_loaded_is_a_noop() {
3634        // Eviction of a never-loaded tenant must not panic and
3635        // must not affect other tenants.
3636        let store = EventStore::new();
3637        store.evict_tenant("nobody"); // should not panic
3638        assert!(!store.is_tenant_loaded("nobody"));
3639    }
3640
3641    #[test]
3642    fn test_lazy_load_accounts_bytes_per_tenant() {
3643        // Step 3 #1: per-tenant byte tracking. Loading a tenant
3644        // should accumulate bytes proportional to its event
3645        // payload sizes; another tenant's counter must stay 0.
3646        let temp_dir = TempDir::new().unwrap();
3647        let storage_dir = temp_dir.path().to_path_buf();
3648
3649        // Persist 5 events for alice with measurable-size payloads.
3650        {
3651            let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3652            for i in 0..5 {
3653                store
3654                    .ingest(
3655                        &Event::from_strings(
3656                            "test.event".to_string(),
3657                            format!("a-{i}"),
3658                            "alice".to_string(),
3659                            serde_json::json!({"data": "x".repeat(1000)}),
3660                            None,
3661                        )
3662                        .unwrap(),
3663                    )
3664                    .unwrap();
3665            }
3666            store.flush_storage().unwrap();
3667        }
3668
3669        let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3670        // Cold: zero bytes accounted.
3671        assert_eq!(store.tenant_resident_bytes("alice"), 0);
3672        assert_eq!(store.cache_resident_bytes(), 0);
3673
3674        store.ensure_tenant_loaded("alice").unwrap();
3675
3676        // After load: alice's counter is non-trivial (5 events
3677        // each carrying ~1000 bytes of payload + overhead).
3678        let alice_bytes = store.tenant_resident_bytes("alice");
3679        assert!(
3680            alice_bytes >= 5 * 1000,
3681            "alice should have at least 5 KiB resident; got {alice_bytes}"
3682        );
3683        // Bob never loaded → 0.
3684        assert_eq!(store.tenant_resident_bytes("bob"), 0);
3685        // Total equals alice's portion (only loaded tenant).
3686        assert_eq!(store.cache_resident_bytes(), alice_bytes);
3687    }
3688
3689    #[test]
3690    fn test_query_lazy_loads_tenant_on_first_call() {
3691        // The end-to-end shape of Step 2: persist events for a
3692        // tenant in session 1, restart, and confirm session 2 boots
3693        // empty but a query for that tenant pulls them in.
3694        let temp_dir = TempDir::new().unwrap();
3695        let storage_dir = temp_dir.path().to_path_buf();
3696
3697        // Session 1: ingest 3 events for tenant "alice", flush, drop.
3698        {
3699            let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3700            for i in 0..3 {
3701                let event = Event::from_strings(
3702                    "test.event".to_string(),
3703                    format!("e-{i}"),
3704                    "alice".to_string(),
3705                    serde_json::json!({"i": i}),
3706                    None,
3707                )
3708                .unwrap();
3709                store.ingest(&event).unwrap();
3710            }
3711            store.flush_storage().unwrap();
3712        }
3713
3714        // Session 2: fresh boot. Events on disk, nothing in memory.
3715        let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3716        assert_eq!(
3717            store.stats().total_events,
3718            0,
3719            "boot must be O(1) — no Parquet pre-load"
3720        );
3721        assert!(!store.is_tenant_loaded("alice"));
3722        assert!(!store.is_tenant_loaded("bob"));
3723
3724        // First query for alice: triggers ensure_tenant_loaded.
3725        let results = store
3726            .query(&QueryEventsRequest {
3727                entity_id: None,
3728                event_type: None,
3729                tenant_id: Some("alice".to_string()),
3730                as_of: None,
3731                since: None,
3732                until: None,
3733                limit: None,
3734                event_type_prefix: None,
3735                exclude_event_type_prefix: None,
3736                payload_filter: None,
3737            })
3738            .unwrap();
3739        assert_eq!(results.len(), 3, "alice's 3 events are returned");
3740        assert!(store.is_tenant_loaded("alice"), "alice now warm");
3741        // bob untouched — load is per-tenant, so a query for alice
3742        // must not have hydrated bob.
3743        assert!(!store.is_tenant_loaded("bob"), "bob still cold");
3744    }
3745
3746    #[test]
3747    fn test_query_invalid_tenant_id_returns_error_no_hang() {
3748        // Step 2 acceptance criterion: in-flight load failures
3749        // surface as errors, not infinite hangs. Path-traversal
3750        // input fails fast at sanitization and propagates.
3751        let temp_dir = TempDir::new().unwrap();
3752        let store = EventStore::with_config(EventStoreConfig::with_persistence(temp_dir.path()));
3753
3754        let result = store.query(&QueryEventsRequest {
3755            entity_id: None,
3756            event_type: None,
3757            tenant_id: Some("../etc".to_string()),
3758            as_of: None,
3759            since: None,
3760            until: None,
3761            limit: None,
3762            event_type_prefix: None,
3763            exclude_event_type_prefix: None,
3764            payload_filter: None,
3765        });
3766        assert!(result.is_err(), "unsafe tenant_id must surface as error");
3767    }
3768
3769    #[test]
3770    fn test_query_concurrent_first_queries_for_same_tenant_all_succeed() {
3771        // Singleflight: N threads racing to query the same cold
3772        // tenant must all return the same correct result. The
3773        // tenant-load must happen exactly once (verified
3774        // structurally by the per-tenant Mutex in tenant_loader,
3775        // tested directly in test_singleflight_blocks_second_caller).
3776        // This integration test confirms the wiring at the query
3777        // level — no thread observes a half-loaded state.
3778        let temp_dir = TempDir::new().unwrap();
3779        let storage_dir = temp_dir.path().to_path_buf();
3780
3781        // Persist 25 events for tenant "alice".
3782        {
3783            let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3784            for i in 0..25 {
3785                let event = Event::from_strings(
3786                    "test.event".to_string(),
3787                    format!("e-{i}"),
3788                    "alice".to_string(),
3789                    serde_json::json!({"i": i}),
3790                    None,
3791                )
3792                .unwrap();
3793                store.ingest(&event).unwrap();
3794            }
3795            store.flush_storage().unwrap();
3796        }
3797
3798        // Fresh boot, then 8 threads simultaneously query alice.
3799        let store = Arc::new(EventStore::with_config(EventStoreConfig::with_persistence(
3800            &storage_dir,
3801        )));
3802        assert!(!store.is_tenant_loaded("alice"));
3803
3804        let mut handles = Vec::new();
3805        for _ in 0..8 {
3806            let s = store.clone();
3807            handles.push(std::thread::spawn(move || {
3808                s.query(&QueryEventsRequest {
3809                    entity_id: None,
3810                    event_type: None,
3811                    tenant_id: Some("alice".to_string()),
3812                    as_of: None,
3813                    since: None,
3814                    until: None,
3815                    limit: None,
3816                    event_type_prefix: None,
3817                    exclude_event_type_prefix: None,
3818                    payload_filter: None,
3819                })
3820            }));
3821        }
3822
3823        for h in handles {
3824            let result = h.join().unwrap().unwrap();
3825            assert_eq!(
3826                result.len(),
3827                25,
3828                "every concurrent caller must see all 25 events"
3829            );
3830        }
3831        assert!(store.is_tenant_loaded("alice"));
3832        // Memory has exactly 25 events — no double-load.
3833        assert_eq!(store.stats().total_events, 25);
3834    }
3835
3836    #[test]
3837    fn test_query_two_cold_tenants_load_independently() {
3838        // Querying tenant A loads only A; querying B then loads
3839        // only B. State after both queries: both tenants warm,
3840        // memory has exactly the expected event counts.
3841        let temp_dir = TempDir::new().unwrap();
3842        let storage_dir = temp_dir.path().to_path_buf();
3843
3844        {
3845            let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3846            for i in 0..3 {
3847                store
3848                    .ingest(
3849                        &Event::from_strings(
3850                            "test.event".to_string(),
3851                            format!("a-{i}"),
3852                            "alice".to_string(),
3853                            serde_json::json!({"i": i}),
3854                            None,
3855                        )
3856                        .unwrap(),
3857                    )
3858                    .unwrap();
3859            }
3860            for i in 0..5 {
3861                store
3862                    .ingest(
3863                        &Event::from_strings(
3864                            "test.event".to_string(),
3865                            format!("b-{i}"),
3866                            "bob".to_string(),
3867                            serde_json::json!({"i": i}),
3868                            None,
3869                        )
3870                        .unwrap(),
3871                    )
3872                    .unwrap();
3873            }
3874            store.flush_storage().unwrap();
3875        }
3876
3877        let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3878        assert_eq!(store.stats().total_events, 0);
3879
3880        // Query alice — bob stays cold.
3881        let alice = store
3882            .query(&QueryEventsRequest {
3883                entity_id: None,
3884                event_type: None,
3885                tenant_id: Some("alice".to_string()),
3886                as_of: None,
3887                since: None,
3888                until: None,
3889                limit: None,
3890                event_type_prefix: None,
3891                exclude_event_type_prefix: None,
3892                payload_filter: None,
3893            })
3894            .unwrap();
3895        assert_eq!(alice.len(), 3);
3896        assert!(store.is_tenant_loaded("alice"));
3897        assert!(!store.is_tenant_loaded("bob"));
3898        assert_eq!(store.stats().total_events, 3);
3899
3900        // Query bob — both warm now.
3901        let bob = store
3902            .query(&QueryEventsRequest {
3903                entity_id: None,
3904                event_type: None,
3905                tenant_id: Some("bob".to_string()),
3906                as_of: None,
3907                since: None,
3908                until: None,
3909                limit: None,
3910                event_type_prefix: None,
3911                exclude_event_type_prefix: None,
3912                payload_filter: None,
3913            })
3914            .unwrap();
3915        assert_eq!(bob.len(), 5);
3916        assert!(store.is_tenant_loaded("bob"));
3917        assert_eq!(store.stats().total_events, 8);
3918    }
3919
3920    #[test]
3921    fn test_boot_with_persisted_data_is_o1() {
3922        // Step 2's headline acceptance criterion: boot time does
3923        // not scale with persisted-data size. The 5M-events / <2s
3924        // target is too large for a unit test, so this asserts the
3925        // weaker but structural property: boot reads zero events
3926        // into memory regardless of how many are on disk.
3927        //
3928        // We persist 50 events across 3 tenants in session 1,
3929        // restart in session 2, and verify session 2's
3930        // total_events is 0. The actual boot wall-clock isn't
3931        // asserted here — it's machine-dependent — but the absence
3932        // of any in-memory data is the structural proxy that the
3933        // boot path no longer iterates Parquet.
3934        let temp_dir = TempDir::new().unwrap();
3935        let storage_dir = temp_dir.path().to_path_buf();
3936
3937        {
3938            let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3939            for tenant in ["alice", "bob", "carol"] {
3940                for i in 0..50 / 3 {
3941                    store
3942                        .ingest(
3943                            &Event::from_strings(
3944                                "test.event".to_string(),
3945                                format!("{tenant}-{i}"),
3946                                tenant.to_string(),
3947                                serde_json::json!({"i": i}),
3948                                None,
3949                            )
3950                            .unwrap(),
3951                        )
3952                        .unwrap();
3953                }
3954            }
3955            store.flush_storage().unwrap();
3956        }
3957
3958        // Confirm there is in fact data on disk to load.
3959        let on_disk = find_parquet_files(&storage_dir);
3960        assert!(
3961            !on_disk.is_empty(),
3962            "session 1 should have produced parquet files; pre-condition for the test"
3963        );
3964
3965        let started = std::time::Instant::now();
3966        let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3967        let boot_elapsed = started.elapsed();
3968
3969        assert_eq!(
3970            store.stats().total_events,
3971            0,
3972            "boot must not pre-load any Parquet events"
3973        );
3974
3975        // Sanity: even on a slow CI box, an O(1) boot finishes in
3976        // well under a second. If this trips it's a strong signal
3977        // the boot path regressed to scanning the Parquet tree.
3978        assert!(
3979            boot_elapsed < std::time::Duration::from_secs(2),
3980            "boot took {boot_elapsed:?} — Step 2 boot should be O(1)"
3981        );
3982    }
3983
3984    #[test]
3985    fn test_query_warm_tenant_does_not_re_read_disk() {
3986        // Performance contract: a warm tenant query goes through the
3987        // DashMap fast path. We can't easily assert "no disk read"
3988        // directly in a unit test, but we CAN assert the call
3989        // succeeds in O(in-memory-events) time even after the
3990        // on-disk file is removed — proving we didn't re-walk it.
3991        let temp_dir = TempDir::new().unwrap();
3992        let storage_dir = temp_dir.path().to_path_buf();
3993
3994        let store = EventStore::with_config(EventStoreConfig::with_persistence(&storage_dir));
3995        for i in 0..3 {
3996            let event = Event::from_strings(
3997                "test.event".to_string(),
3998                format!("e-{i}"),
3999                "alice".to_string(),
4000                serde_json::json!({"i": i}),
4001                None,
4002            )
4003            .unwrap();
4004            store.ingest(&event).unwrap();
4005        }
4006        store.flush_storage().unwrap();
4007
4008        // First query: cold, hits disk.
4009        let _ = store
4010            .query(&QueryEventsRequest {
4011                entity_id: None,
4012                event_type: None,
4013                tenant_id: Some("alice".to_string()),
4014                as_of: None,
4015                since: None,
4016                until: None,
4017                limit: None,
4018                event_type_prefix: None,
4019                exclude_event_type_prefix: None,
4020                payload_filter: None,
4021            })
4022            .unwrap();
4023        assert!(store.is_tenant_loaded("alice"));
4024
4025        // Now wipe the on-disk file. A warm-path query must still
4026        // succeed because it doesn't need disk.
4027        let parquet_files = find_parquet_files(&storage_dir);
4028        for f in parquet_files {
4029            std::fs::remove_file(&f).unwrap();
4030        }
4031
4032        let results = store
4033            .query(&QueryEventsRequest {
4034                entity_id: None,
4035                event_type: None,
4036                tenant_id: Some("alice".to_string()),
4037                as_of: None,
4038                since: None,
4039                until: None,
4040                limit: None,
4041                event_type_prefix: None,
4042                exclude_event_type_prefix: None,
4043                payload_filter: None,
4044            })
4045            .unwrap();
4046        assert_eq!(
4047            results.len(),
4048            3,
4049            "warm tenant query must not need disk; got {} events from a deleted parquet",
4050            results.len()
4051        );
4052    }
4053
4054    #[test]
4055    fn test_event_store_default() {
4056        let store = EventStore::default();
4057        assert_eq!(store.stats().total_events, 0);
4058    }
4059
4060    #[test]
4061    fn test_ingest_single_event() {
4062        let store = EventStore::new();
4063        let event = create_test_event("entity-1", "user.created");
4064
4065        store.ingest(&event).unwrap();
4066
4067        assert_eq!(store.stats().total_events, 1);
4068        assert_eq!(store.stats().total_ingested, 1);
4069    }
4070
4071    #[test]
4072    fn test_ingest_multiple_events() {
4073        let store = EventStore::new();
4074
4075        for i in 0..10 {
4076            let event = create_test_event(&format!("entity-{i}"), "user.created");
4077            store.ingest(&event).unwrap();
4078        }
4079
4080        assert_eq!(store.stats().total_events, 10);
4081        assert_eq!(store.stats().total_ingested, 10);
4082    }
4083
4084    #[test]
4085    fn test_query_by_entity_id() {
4086        let store = EventStore::new();
4087
4088        store
4089            .ingest(&create_test_event("entity-1", "user.created"))
4090            .unwrap();
4091        store
4092            .ingest(&create_test_event("entity-2", "user.created"))
4093            .unwrap();
4094        store
4095            .ingest(&create_test_event("entity-1", "user.updated"))
4096            .unwrap();
4097
4098        let results = store
4099            .query(&QueryEventsRequest {
4100                entity_id: Some("entity-1".to_string()),
4101                event_type: None,
4102                tenant_id: None,
4103                as_of: None,
4104                since: None,
4105                until: None,
4106                limit: None,
4107                event_type_prefix: None,
4108                exclude_event_type_prefix: None,
4109                payload_filter: None,
4110            })
4111            .unwrap();
4112
4113        assert_eq!(results.len(), 2);
4114    }
4115
4116    #[test]
4117    fn test_query_by_event_type() {
4118        let store = EventStore::new();
4119
4120        store
4121            .ingest(&create_test_event("entity-1", "user.created"))
4122            .unwrap();
4123        store
4124            .ingest(&create_test_event("entity-2", "user.updated"))
4125            .unwrap();
4126        store
4127            .ingest(&create_test_event("entity-3", "user.created"))
4128            .unwrap();
4129
4130        let results = store
4131            .query(&QueryEventsRequest {
4132                entity_id: None,
4133                event_type: Some("user.created".to_string()),
4134                tenant_id: None,
4135                as_of: None,
4136                since: None,
4137                until: None,
4138                limit: None,
4139                event_type_prefix: None,
4140                exclude_event_type_prefix: None,
4141                payload_filter: None,
4142            })
4143            .unwrap();
4144
4145        assert_eq!(results.len(), 2);
4146    }
4147
4148    #[test]
4149    fn test_query_with_limit() {
4150        let store = EventStore::new();
4151
4152        for i in 0..10 {
4153            let event = create_test_event(&format!("entity-{i}"), "user.created");
4154            store.ingest(&event).unwrap();
4155        }
4156
4157        let results = store
4158            .query(&QueryEventsRequest {
4159                entity_id: None,
4160                event_type: None,
4161                tenant_id: None,
4162                as_of: None,
4163                since: None,
4164                until: None,
4165                limit: Some(5),
4166                event_type_prefix: None,
4167                exclude_event_type_prefix: None,
4168                payload_filter: None,
4169            })
4170            .unwrap();
4171
4172        assert_eq!(results.len(), 5);
4173    }
4174
4175    #[test]
4176    fn test_query_empty_store() {
4177        let store = EventStore::new();
4178
4179        let results = store
4180            .query(&QueryEventsRequest {
4181                entity_id: Some("non-existent".to_string()),
4182                event_type: None,
4183                tenant_id: None,
4184                as_of: None,
4185                since: None,
4186                until: None,
4187                limit: None,
4188                event_type_prefix: None,
4189                exclude_event_type_prefix: None,
4190                payload_filter: None,
4191            })
4192            .unwrap();
4193
4194        assert!(results.is_empty());
4195    }
4196
4197    #[test]
4198    fn test_reconstruct_state() {
4199        let store = EventStore::new();
4200
4201        store
4202            .ingest(&create_test_event("entity-1", "user.created"))
4203            .unwrap();
4204
4205        let state = store.reconstruct_state("entity-1", None).unwrap();
4206        // The state is wrapped with metadata
4207        assert_eq!(state["current_state"]["name"], "Test");
4208        assert_eq!(state["current_state"]["value"], 42);
4209    }
4210
4211    #[test]
4212    fn test_reconstruct_state_not_found() {
4213        let store = EventStore::new();
4214
4215        let result = store.reconstruct_state("non-existent", None);
4216        assert!(result.is_err());
4217    }
4218
4219    #[test]
4220    fn test_get_snapshot_empty() {
4221        let store = EventStore::new();
4222
4223        let result = store.get_snapshot("non-existent");
4224        // Entity not found error is expected
4225        assert!(result.is_err());
4226    }
4227
4228    #[test]
4229    fn test_create_snapshot() {
4230        let store = EventStore::new();
4231
4232        store
4233            .ingest(&create_test_event("entity-1", "user.created"))
4234            .unwrap();
4235
4236        store.create_snapshot("entity-1").unwrap();
4237
4238        // Verify snapshot was created
4239        let snapshot = store.get_snapshot("entity-1").unwrap();
4240        assert_ne!(snapshot, serde_json::json!(null));
4241    }
4242
4243    #[test]
4244    fn test_create_snapshot_entity_not_found() {
4245        let store = EventStore::new();
4246
4247        let result = store.create_snapshot("non-existent");
4248        assert!(result.is_err());
4249    }
4250
4251    #[test]
4252    fn test_websocket_manager() {
4253        let store = EventStore::new();
4254        let manager = store.websocket_manager();
4255        // Manager should be accessible
4256        assert!(Arc::strong_count(&manager) >= 1);
4257    }
4258
4259    #[test]
4260    fn test_snapshot_manager() {
4261        let store = EventStore::new();
4262        let manager = store.snapshot_manager();
4263        assert!(Arc::strong_count(&manager) >= 1);
4264    }
4265
4266    #[test]
4267    fn test_compaction_manager_none() {
4268        let store = EventStore::new();
4269        // Without storage_dir, compaction manager should be None
4270        assert!(store.compaction_manager().is_none());
4271    }
4272
4273    #[test]
4274    fn test_schema_registry() {
4275        let store = EventStore::new();
4276        let registry = store.schema_registry();
4277        assert!(Arc::strong_count(&registry) >= 1);
4278    }
4279
4280    #[test]
4281    fn test_replay_manager() {
4282        let store = EventStore::new();
4283        let manager = store.replay_manager();
4284        assert!(Arc::strong_count(&manager) >= 1);
4285    }
4286
4287    #[test]
4288    fn test_pipeline_manager() {
4289        let store = EventStore::new();
4290        let manager = store.pipeline_manager();
4291        assert!(Arc::strong_count(&manager) >= 1);
4292    }
4293
4294    #[test]
4295    fn test_projection_manager() {
4296        let store = EventStore::new();
4297        let manager = store.projection_manager();
4298        // Built-in projections should be registered
4299        let projections = manager.list_projections();
4300        assert!(projections.len() >= 2); // entity_snapshots and event_counters
4301    }
4302
4303    #[test]
4304    fn test_projection_state_cache() {
4305        let store = EventStore::new();
4306        let cache = store.projection_state_cache();
4307
4308        cache.insert("test:key".to_string(), serde_json::json!({"value": 123}));
4309        assert_eq!(cache.len(), 1);
4310
4311        let value = cache.get("test:key").unwrap();
4312        assert_eq!(value["value"], 123);
4313    }
4314
4315    #[test]
4316    fn test_metrics() {
4317        let store = EventStore::new();
4318        let metrics = store.metrics();
4319        assert!(Arc::strong_count(&metrics) >= 1);
4320    }
4321
4322    #[test]
4323    fn test_store_stats() {
4324        let store = EventStore::new();
4325
4326        store
4327            .ingest(&create_test_event("entity-1", "user.created"))
4328            .unwrap();
4329        store
4330            .ingest(&create_test_event("entity-2", "order.placed"))
4331            .unwrap();
4332
4333        let stats = store.stats();
4334        assert_eq!(stats.total_events, 2);
4335        assert_eq!(stats.total_entities, 2);
4336        assert_eq!(stats.total_event_types, 2);
4337        assert_eq!(stats.total_ingested, 2);
4338    }
4339
4340    #[test]
4341    fn test_event_store_config_default() {
4342        let config = EventStoreConfig::default();
4343        assert!(config.storage_dir.is_none());
4344        assert!(config.wal_dir.is_none());
4345    }
4346
4347    #[test]
4348    fn test_event_store_config_with_persistence() {
4349        let temp_dir = TempDir::new().unwrap();
4350        let config = EventStoreConfig::with_persistence(temp_dir.path());
4351
4352        assert!(config.storage_dir.is_some());
4353        assert!(config.wal_dir.is_none());
4354    }
4355
4356    #[test]
4357    fn test_event_store_config_with_wal() {
4358        let temp_dir = TempDir::new().unwrap();
4359        let config = EventStoreConfig::with_wal(temp_dir.path(), WALConfig::default());
4360
4361        assert!(config.storage_dir.is_none());
4362        assert!(config.wal_dir.is_some());
4363    }
4364
4365    #[test]
4366    fn test_event_store_config_with_all() {
4367        let temp_dir = TempDir::new().unwrap();
4368        let config = EventStoreConfig::with_all(temp_dir.path(), SnapshotConfig::default());
4369
4370        assert!(config.storage_dir.is_some());
4371    }
4372
4373    #[test]
4374    fn test_event_store_config_production() {
4375        let storage_dir = TempDir::new().unwrap();
4376        let wal_dir = TempDir::new().unwrap();
4377        let config = EventStoreConfig::production(
4378            storage_dir.path(),
4379            wal_dir.path(),
4380            SnapshotConfig::default(),
4381            WALConfig::default(),
4382            CompactionConfig::default(),
4383        );
4384
4385        assert!(config.storage_dir.is_some());
4386        assert!(config.wal_dir.is_some());
4387    }
4388
4389    // -----------------------------------------------------------------------
4390    // from_env_vars tests — verifies the env-var-to-config wiring that
4391    // caused the durability bug (events lost on restart) in v0.10.3.
4392    // -----------------------------------------------------------------------
4393
4394    #[test]
4395    fn test_from_env_vars_data_dir_enables_full_persistence() {
4396        let (config, mode) = EventStoreConfig::from_env_vars(
4397            Some("/app/data".to_string()),
4398            None,
4399            None,
4400            None,
4401            None,
4402            None,
4403            None,
4404            None,
4405        );
4406        assert_eq!(mode, "wal+parquet");
4407        assert_eq!(
4408            config.storage_dir.unwrap().to_str().unwrap(),
4409            "/app/data/storage"
4410        );
4411        assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/app/data/wal");
4412    }
4413
4414    #[test]
4415    fn test_from_env_vars_explicit_dirs() {
4416        let (config, mode) = EventStoreConfig::from_env_vars(
4417            None,
4418            Some("/custom/storage".to_string()),
4419            Some("/custom/wal".to_string()),
4420            None,
4421            None,
4422            None,
4423            None,
4424            None,
4425        );
4426        assert_eq!(mode, "wal+parquet");
4427        assert_eq!(
4428            config.storage_dir.unwrap().to_str().unwrap(),
4429            "/custom/storage"
4430        );
4431        assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/custom/wal");
4432    }
4433
4434    #[test]
4435    fn test_from_env_vars_wal_disabled() {
4436        let (config, mode) = EventStoreConfig::from_env_vars(
4437            Some("/app/data".to_string()),
4438            None,
4439            None,
4440            Some("false".to_string()),
4441            None,
4442            None,
4443            None,
4444            None,
4445        );
4446        assert_eq!(mode, "parquet-only");
4447        assert!(config.storage_dir.is_some());
4448        assert!(config.wal_dir.is_none());
4449    }
4450
4451    #[test]
4452    fn test_from_env_vars_no_dirs_is_in_memory() {
4453        let (config, mode) =
4454            EventStoreConfig::from_env_vars(None, None, None, None, None, None, None, None);
4455        assert_eq!(mode, "in-memory");
4456        assert!(config.storage_dir.is_none());
4457        assert!(config.wal_dir.is_none());
4458    }
4459
4460    #[test]
4461    fn test_from_env_vars_empty_strings_treated_as_none() {
4462        let (_, mode) = EventStoreConfig::from_env_vars(
4463            Some(String::new()),
4464            Some(String::new()),
4465            Some(String::new()),
4466            None,
4467            None,
4468            None,
4469            None,
4470            None,
4471        );
4472        assert_eq!(mode, "in-memory");
4473    }
4474
4475    #[test]
4476    fn test_from_env_vars_explicit_overrides_data_dir() {
4477        let (config, mode) = EventStoreConfig::from_env_vars(
4478            Some("/app/data".to_string()),
4479            Some("/override/storage".to_string()),
4480            Some("/override/wal".to_string()),
4481            None,
4482            None,
4483            None,
4484            None,
4485            None,
4486        );
4487        assert_eq!(mode, "wal+parquet");
4488        assert_eq!(
4489            config.storage_dir.unwrap().to_str().unwrap(),
4490            "/override/storage"
4491        );
4492        assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/override/wal");
4493    }
4494
4495    #[test]
4496    fn test_from_env_vars_wal_only() {
4497        let (config, mode) = EventStoreConfig::from_env_vars(
4498            None,
4499            None,
4500            Some("/wal/only".to_string()),
4501            None,
4502            None,
4503            None,
4504            None,
4505            None,
4506        );
4507        assert_eq!(mode, "wal-only");
4508        assert!(config.storage_dir.is_none());
4509        assert_eq!(config.wal_dir.unwrap().to_str().unwrap(), "/wal/only");
4510    }
4511
4512    #[test]
4513    fn test_from_env_vars_cache_bytes_parses_decimal() {
4514        let (config, _) = EventStoreConfig::from_env_vars(
4515            Some("/app/data".to_string()),
4516            None,
4517            None,
4518            None,
4519            Some("536870912".to_string()),
4520            // 512 MiB
4521            None,
4522            None,
4523            None,
4524        );
4525        assert_eq!(config.cache_byte_budget, Some(536_870_912));
4526    }
4527
4528    #[test]
4529    fn test_from_env_vars_cache_bytes_unparseable_disables_budget() {
4530        // Garbage in CACHE_BYTES doesn't fail boot — we log and
4531        // fall back to no-budget. The unbounded fallback is safe
4532        // (just the pre-Step-3 behavior).
4533        let (config, _) = EventStoreConfig::from_env_vars(
4534            Some("/app/data".to_string()),
4535            None,
4536            None,
4537            None,
4538            Some("not-a-number".to_string()),
4539            None,
4540            None,
4541            None,
4542        );
4543        assert_eq!(config.cache_byte_budget, None);
4544    }
4545
4546    #[test]
4547    fn test_from_env_vars_cache_bytes_empty_disables_budget() {
4548        let (config, _) = EventStoreConfig::from_env_vars(
4549            Some("/app/data".to_string()),
4550            None,
4551            None,
4552            None,
4553            Some(String::new()),
4554            None,
4555            None,
4556            None,
4557        );
4558        assert_eq!(config.cache_byte_budget, None);
4559    }
4560
4561    #[test]
4562    fn test_from_env_vars_snapshot_interval_overrides_default() {
4563        // ALLSOURCE_SNAPSHOT_INTERVAL_SECONDS plumbs through to
4564        // CompactionConfig.compaction_interval_seconds. Default is
4565        // 3600s (hourly) per the bead.
4566        let (config, _) = EventStoreConfig::from_env_vars(
4567            Some("/app/data".to_string()),
4568            None,
4569            None,
4570            None,
4571            None,
4572            Some("60".to_string()),
4573            None,
4574            None,
4575        );
4576        assert_eq!(config.compaction_config.compaction_interval_seconds, 60);
4577    }
4578
4579    #[test]
4580    fn test_from_env_vars_snapshot_interval_default_is_hourly() {
4581        let (config, _) = EventStoreConfig::from_env_vars(
4582            Some("/app/data".to_string()),
4583            None,
4584            None,
4585            None,
4586            None,
4587            None,
4588            None,
4589            None,
4590        );
4591        assert_eq!(config.compaction_config.compaction_interval_seconds, 3600);
4592    }
4593
4594    #[test]
4595    fn test_from_env_vars_snapshot_interval_unparseable_falls_back() {
4596        let (config, _) = EventStoreConfig::from_env_vars(
4597            Some("/app/data".to_string()),
4598            None,
4599            None,
4600            None,
4601            None,
4602            Some("not-a-number".to_string()),
4603            None,
4604            None,
4605        );
4606        assert_eq!(config.compaction_config.compaction_interval_seconds, 3600);
4607    }
4608
4609    #[test]
4610    fn test_from_env_vars_retention_system_days_overrides_default() {
4611        // Step 5: ALLSOURCE_RETENTION_SYSTEM_DAYS overrides the
4612        // default 30-day TTL for the system tenant.
4613        let (config, _) = EventStoreConfig::from_env_vars(
4614            Some("/app/data".to_string()),
4615            None,
4616            None,
4617            None,
4618            None,
4619            None,
4620            Some("7".to_string()),
4621            None,
4622        );
4623        let ttl = config
4624            .compaction_config
4625            .retention
4626            .ttl_for("system")
4627            .unwrap();
4628        assert_eq!(ttl.as_secs(), 7 * 24 * 3600);
4629    }
4630
4631    #[test]
4632    fn test_from_env_vars_retention_default_is_30_days_for_system() {
4633        let (config, _) = EventStoreConfig::from_env_vars(
4634            Some("/app/data".to_string()),
4635            None,
4636            None,
4637            None,
4638            None,
4639            None,
4640            None,
4641            None,
4642        );
4643        let ttl = config
4644            .compaction_config
4645            .retention
4646            .ttl_for("system")
4647            .unwrap();
4648        assert_eq!(ttl.as_secs(), 30 * 24 * 3600);
4649        // Other tenants keep forever by default.
4650        assert!(config.compaction_config.retention.ttl_for("acme").is_none());
4651    }
4652
4653    #[test]
4654    fn test_store_stats_serde() {
4655        let stats = StoreStats {
4656            total_events: 100,
4657            total_entities: 50,
4658            total_event_types: 10,
4659            total_ingested: 100,
4660        };
4661
4662        let json = serde_json::to_string(&stats).unwrap();
4663        assert!(json.contains("\"total_events\":100"));
4664        assert!(json.contains("\"total_entities\":50"));
4665    }
4666
4667    #[test]
4668    fn test_query_with_entity_and_type() {
4669        let store = EventStore::new();
4670
4671        store
4672            .ingest(&create_test_event("entity-1", "user.created"))
4673            .unwrap();
4674        store
4675            .ingest(&create_test_event("entity-1", "user.updated"))
4676            .unwrap();
4677        store
4678            .ingest(&create_test_event("entity-2", "user.created"))
4679            .unwrap();
4680
4681        let results = store
4682            .query(&QueryEventsRequest {
4683                entity_id: Some("entity-1".to_string()),
4684                event_type: Some("user.created".to_string()),
4685                tenant_id: None,
4686                as_of: None,
4687                since: None,
4688                until: None,
4689                limit: None,
4690                event_type_prefix: None,
4691                exclude_event_type_prefix: None,
4692                payload_filter: None,
4693            })
4694            .unwrap();
4695
4696        assert_eq!(results.len(), 1);
4697        assert_eq!(results[0].event_type_str(), "user.created");
4698    }
4699
4700    #[test]
4701    fn test_query_by_event_type_prefix() {
4702        let store = EventStore::new();
4703
4704        // Ingest events with various types
4705        store
4706            .ingest(&create_test_event("entity-1", "index.created"))
4707            .unwrap();
4708        store
4709            .ingest(&create_test_event("entity-2", "index.updated"))
4710            .unwrap();
4711        store
4712            .ingest(&create_test_event("entity-3", "trade.created"))
4713            .unwrap();
4714        store
4715            .ingest(&create_test_event("entity-4", "trade.completed"))
4716            .unwrap();
4717        store
4718            .ingest(&create_test_event("entity-5", "balance.updated"))
4719            .unwrap();
4720
4721        // Query with prefix "index." should return exactly 2
4722        let results = store
4723            .query(&QueryEventsRequest {
4724                entity_id: None,
4725                event_type: None,
4726                tenant_id: None,
4727                as_of: None,
4728                since: None,
4729                until: None,
4730                limit: None,
4731                event_type_prefix: Some("index.".to_string()),
4732                exclude_event_type_prefix: None,
4733                payload_filter: None,
4734            })
4735            .unwrap();
4736
4737        assert_eq!(results.len(), 2);
4738        assert!(
4739            results
4740                .iter()
4741                .all(|e| e.event_type_str().starts_with("index."))
4742        );
4743    }
4744
4745    #[test]
4746    fn test_query_by_event_type_prefix_empty_returns_all() {
4747        let store = EventStore::new();
4748
4749        store
4750            .ingest(&create_test_event("entity-1", "index.created"))
4751            .unwrap();
4752        store
4753            .ingest(&create_test_event("entity-2", "trade.created"))
4754            .unwrap();
4755
4756        // Empty prefix matches all types
4757        let results = store
4758            .query(&QueryEventsRequest {
4759                entity_id: None,
4760                event_type: None,
4761                tenant_id: None,
4762                as_of: None,
4763                since: None,
4764                until: None,
4765                limit: None,
4766                event_type_prefix: Some(String::new()),
4767                exclude_event_type_prefix: None,
4768                payload_filter: None,
4769            })
4770            .unwrap();
4771
4772        assert_eq!(results.len(), 2);
4773    }
4774
4775    #[test]
4776    fn test_query_by_event_type_prefix_no_match() {
4777        let store = EventStore::new();
4778
4779        store
4780            .ingest(&create_test_event("entity-1", "index.created"))
4781            .unwrap();
4782
4783        let results = store
4784            .query(&QueryEventsRequest {
4785                entity_id: None,
4786                event_type: None,
4787                tenant_id: None,
4788                as_of: None,
4789                since: None,
4790                until: None,
4791                limit: None,
4792                event_type_prefix: Some("nonexistent.".to_string()),
4793                exclude_event_type_prefix: None,
4794                payload_filter: None,
4795            })
4796            .unwrap();
4797
4798        assert!(results.is_empty());
4799    }
4800
4801    #[test]
4802    fn test_query_by_entity_with_type_prefix() {
4803        let store = EventStore::new();
4804
4805        store
4806            .ingest(&create_test_event("entity-1", "index.created"))
4807            .unwrap();
4808        store
4809            .ingest(&create_test_event("entity-1", "trade.created"))
4810            .unwrap();
4811        store
4812            .ingest(&create_test_event("entity-2", "index.updated"))
4813            .unwrap();
4814
4815        // Query entity-1 with prefix "index." should return 1
4816        let results = store
4817            .query(&QueryEventsRequest {
4818                entity_id: Some("entity-1".to_string()),
4819                event_type: None,
4820                tenant_id: None,
4821                as_of: None,
4822                since: None,
4823                until: None,
4824                limit: None,
4825                event_type_prefix: Some("index.".to_string()),
4826                exclude_event_type_prefix: None,
4827                payload_filter: None,
4828            })
4829            .unwrap();
4830
4831        assert_eq!(results.len(), 1);
4832        assert_eq!(results[0].event_type_str(), "index.created");
4833    }
4834
4835    #[test]
4836    fn test_query_prefix_with_limit() {
4837        let store = EventStore::new();
4838
4839        for i in 0..5 {
4840            store
4841                .ingest(&create_test_event(&format!("entity-{i}"), "index.created"))
4842                .unwrap();
4843        }
4844
4845        let results = store
4846            .query(&QueryEventsRequest {
4847                entity_id: None,
4848                event_type: None,
4849                tenant_id: None,
4850                as_of: None,
4851                since: None,
4852                until: None,
4853                limit: Some(3),
4854                event_type_prefix: Some("index.".to_string()),
4855                exclude_event_type_prefix: None,
4856                payload_filter: None,
4857            })
4858            .unwrap();
4859
4860        assert_eq!(results.len(), 3);
4861    }
4862
4863    #[test]
4864    fn test_query_prefix_alongside_existing_filters() {
4865        let store = EventStore::new();
4866
4867        store
4868            .ingest(&create_test_event("entity-1", "index.created"))
4869            .unwrap();
4870        // Sleep briefly to ensure different timestamps
4871        std::thread::sleep(std::time::Duration::from_millis(10));
4872        store
4873            .ingest(&create_test_event("entity-2", "index.strategy.updated"))
4874            .unwrap();
4875        std::thread::sleep(std::time::Duration::from_millis(10));
4876        store
4877            .ingest(&create_test_event("entity-3", "index.deleted"))
4878            .unwrap();
4879
4880        // Prefix with limit
4881        let results = store
4882            .query(&QueryEventsRequest {
4883                entity_id: None,
4884                event_type: None,
4885                tenant_id: None,
4886                as_of: None,
4887                since: None,
4888                until: None,
4889                limit: Some(2),
4890                event_type_prefix: Some("index.".to_string()),
4891                exclude_event_type_prefix: None,
4892                payload_filter: None,
4893            })
4894            .unwrap();
4895
4896        assert_eq!(results.len(), 2);
4897    }
4898
4899    #[test]
4900    fn test_query_with_payload_filter() {
4901        let store = EventStore::new();
4902
4903        // Ingest 5 events with user_id=alice
4904        for i in 0..5 {
4905            store
4906                .ingest(&create_test_event_with_payload(
4907                    &format!("entity-{i}"),
4908                    "user.action",
4909                    serde_json::json!({"user_id": "alice", "action": "click"}),
4910                ))
4911                .unwrap();
4912        }
4913        // Ingest 5 events with user_id=bob
4914        for i in 5..10 {
4915            store
4916                .ingest(&create_test_event_with_payload(
4917                    &format!("entity-{i}"),
4918                    "user.action",
4919                    serde_json::json!({"user_id": "bob", "action": "view"}),
4920                ))
4921                .unwrap();
4922        }
4923
4924        // Filter for alice
4925        let results = store
4926            .query(&QueryEventsRequest {
4927                entity_id: None,
4928                event_type: Some("user.action".to_string()),
4929                tenant_id: None,
4930                as_of: None,
4931                since: None,
4932                until: None,
4933                limit: None,
4934                event_type_prefix: None,
4935                exclude_event_type_prefix: None,
4936                payload_filter: Some(r#"{"user_id":"alice"}"#.to_string()),
4937            })
4938            .unwrap();
4939
4940        assert_eq!(results.len(), 5);
4941    }
4942
4943    #[test]
4944    fn test_query_payload_filter_non_existent_field() {
4945        let store = EventStore::new();
4946
4947        store
4948            .ingest(&create_test_event_with_payload(
4949                "entity-1",
4950                "user.action",
4951                serde_json::json!({"user_id": "alice"}),
4952            ))
4953            .unwrap();
4954
4955        // Filter for a field that doesn't exist — returns 0, not error
4956        let results = store
4957            .query(&QueryEventsRequest {
4958                entity_id: None,
4959                event_type: None,
4960                tenant_id: None,
4961                as_of: None,
4962                since: None,
4963                until: None,
4964                limit: None,
4965                event_type_prefix: None,
4966                exclude_event_type_prefix: None,
4967                payload_filter: Some(r#"{"nonexistent":"value"}"#.to_string()),
4968            })
4969            .unwrap();
4970
4971        assert!(results.is_empty());
4972    }
4973
4974    #[test]
4975    fn test_query_payload_filter_with_prefix() {
4976        let store = EventStore::new();
4977
4978        store
4979            .ingest(&create_test_event_with_payload(
4980                "entity-1",
4981                "index.created",
4982                serde_json::json!({"status": "active"}),
4983            ))
4984            .unwrap();
4985        store
4986            .ingest(&create_test_event_with_payload(
4987                "entity-2",
4988                "index.created",
4989                serde_json::json!({"status": "inactive"}),
4990            ))
4991            .unwrap();
4992        store
4993            .ingest(&create_test_event_with_payload(
4994                "entity-3",
4995                "trade.created",
4996                serde_json::json!({"status": "active"}),
4997            ))
4998            .unwrap();
4999
5000        // Combine prefix + payload filter
5001        let results = store
5002            .query(&QueryEventsRequest {
5003                entity_id: None,
5004                event_type: None,
5005                tenant_id: None,
5006                as_of: None,
5007                since: None,
5008                until: None,
5009                limit: None,
5010                event_type_prefix: Some("index.".to_string()),
5011                exclude_event_type_prefix: None,
5012                payload_filter: Some(r#"{"status":"active"}"#.to_string()),
5013            })
5014            .unwrap();
5015
5016        assert_eq!(results.len(), 1);
5017        assert_eq!(results[0].entity_id().to_string(), "entity-1");
5018    }
5019
5020    #[test]
5021    fn test_flush_storage_no_storage() {
5022        let store = EventStore::new();
5023        // Without storage, flush should succeed (no-op)
5024        let result = store.flush_storage();
5025        assert!(result.is_ok());
5026    }
5027
5028    #[test]
5029    fn test_state_evolution() {
5030        let store = EventStore::new();
5031
5032        // Initial state
5033        store
5034            .ingest(
5035                &Event::from_strings(
5036                    "user.created".to_string(),
5037                    "user-1".to_string(),
5038                    "default".to_string(),
5039                    serde_json::json!({"name": "Alice", "age": 25}),
5040                    None,
5041                )
5042                .unwrap(),
5043            )
5044            .unwrap();
5045
5046        // Update state
5047        store
5048            .ingest(
5049                &Event::from_strings(
5050                    "user.updated".to_string(),
5051                    "user-1".to_string(),
5052                    "default".to_string(),
5053                    serde_json::json!({"age": 26}),
5054                    None,
5055                )
5056                .unwrap(),
5057            )
5058            .unwrap();
5059
5060        let state = store.reconstruct_state("user-1", None).unwrap();
5061        // The state is wrapped with metadata
5062        assert_eq!(state["current_state"]["name"], "Alice");
5063        assert_eq!(state["current_state"]["age"], 26);
5064    }
5065
5066    #[test]
5067    fn test_reject_system_event_types() {
5068        let store = EventStore::new();
5069
5070        // System event types should be rejected via user-facing ingestion
5071        let event = Event::reconstruct_from_strings(
5072            uuid::Uuid::new_v4(),
5073            "_system.tenant.created".to_string(),
5074            "_system:tenant:acme".to_string(),
5075            "_system".to_string(),
5076            serde_json::json!({"name": "ACME"}),
5077            chrono::Utc::now(),
5078            None,
5079            1,
5080        );
5081
5082        let result = store.ingest(&event);
5083        assert!(result.is_err());
5084        let err = result.unwrap_err();
5085        assert!(
5086            err.to_string().contains("reserved for internal use"),
5087            "Expected system namespace rejection, got: {err}"
5088        );
5089    }
5090
5091    // -----------------------------------------------------------------------
5092    // Crash recovery: WAL events survive restart via Parquet checkpoint.
5093    // Regression test for GitHub issue #84 — flush_storage() was a no-op
5094    // during recovery because events were never buffered into Parquet's
5095    // current_batch before flushing.
5096    // -----------------------------------------------------------------------
5097
5098    #[test]
5099    fn test_wal_recovery_checkpoints_to_parquet() {
5100        let data_dir = TempDir::new().unwrap();
5101        let storage_dir = data_dir.path().join("storage");
5102        let wal_dir = data_dir.path().join("wal");
5103
5104        // Session 1: ingest events with WAL + Parquet
5105        {
5106            let config = EventStoreConfig::production(
5107                &storage_dir,
5108                &wal_dir,
5109                SnapshotConfig::default(),
5110                WALConfig {
5111                    sync_on_write: true,
5112                    ..WALConfig::default()
5113                },
5114                CompactionConfig::default(),
5115            );
5116            let store = EventStore::with_config(config);
5117
5118            for i in 0..5 {
5119                let event = Event::from_strings(
5120                    "test.created".to_string(),
5121                    format!("entity-{i}"),
5122                    "default".to_string(),
5123                    serde_json::json!({"index": i}),
5124                    None,
5125                )
5126                .unwrap();
5127                store.ingest(&event).unwrap();
5128            }
5129
5130            assert_eq!(store.stats().total_events, 5);
5131
5132            // Do NOT call flush_storage or shutdown — simulate a crash.
5133            // Events are in WAL (sync_on_write: true) but NOT in Parquet.
5134        }
5135
5136        // Verify WAL file has data
5137        let wal_files: Vec<_> = std::fs::read_dir(&wal_dir)
5138            .unwrap()
5139            .filter_map(std::result::Result::ok)
5140            .filter(|e| e.path().extension().is_some_and(|ext| ext == "log"))
5141            .collect();
5142        assert!(!wal_files.is_empty(), "WAL file should exist");
5143        let wal_size = wal_files[0].metadata().unwrap().len();
5144        assert!(wal_size > 0, "WAL file should have data (got 0 bytes)");
5145
5146        // Session 2: reopen — recovery should checkpoint WAL to Parquet, then truncate
5147        {
5148            let config = EventStoreConfig::production(
5149                &storage_dir,
5150                &wal_dir,
5151                SnapshotConfig::default(),
5152                WALConfig {
5153                    sync_on_write: true,
5154                    ..WALConfig::default()
5155                },
5156                CompactionConfig::default(),
5157            );
5158            let store = EventStore::with_config(config);
5159
5160            // Events should be recovered
5161            assert_eq!(
5162                store.stats().total_events,
5163                5,
5164                "Session 2 should have all 5 events after WAL recovery"
5165            );
5166
5167            // Parquet should now have files (checkpoint happened).
5168            // After Step 1, files live under <root>/<tenant>/<yyyy-mm>/,
5169            // so walk recursively.
5170            let parquet_files = find_parquet_files(&storage_dir);
5171            assert!(
5172                !parquet_files.is_empty(),
5173                "Parquet file should exist after WAL checkpoint"
5174            );
5175        }
5176
5177        // Session 3: reopen again — events should be reachable via
5178        // lazy-load (Step 2: boot does not pre-load Parquet).
5179        {
5180            let config = EventStoreConfig::production(
5181                &storage_dir,
5182                &wal_dir,
5183                SnapshotConfig::default(),
5184                WALConfig {
5185                    sync_on_write: true,
5186                    ..WALConfig::default()
5187                },
5188                CompactionConfig::default(),
5189            );
5190            let store = EventStore::with_config(config);
5191
5192            // Boot is now O(1) — Parquet stays cold until first
5193            // per-tenant query. WAL was truncated in session 2,
5194            // so nothing is pre-loaded.
5195            assert_eq!(
5196                store.stats().total_events,
5197                0,
5198                "Session 3 boot should not pre-load Parquet (lazy-load mode)"
5199            );
5200
5201            // Trigger lazy load for the test tenant (events were
5202            // ingested with tenant_id=\"default\").
5203            store.ensure_tenant_loaded("default").unwrap();
5204            assert_eq!(
5205                store.stats().total_events,
5206                5,
5207                "Session 3 should have all 5 events after ensure_tenant_loaded"
5208            );
5209        }
5210    }
5211
5212    #[test]
5213    fn test_parquet_restore_surfaces_errors_not_silent() {
5214        // Write events with WAL+Parquet, flush to Parquet, then corrupt the
5215        // Parquet file. On reload, the error must be logged (not silently
5216        // swallowed as 0 events).
5217        let data_dir = TempDir::new().unwrap();
5218        let storage_dir = data_dir.path().join("storage");
5219        let wal_dir = data_dir.path().join("wal");
5220
5221        // Session 1: write events and flush to Parquet
5222        {
5223            let config = EventStoreConfig::production(
5224                &storage_dir,
5225                &wal_dir,
5226                SnapshotConfig::default(),
5227                WALConfig {
5228                    sync_on_write: true,
5229                    ..WALConfig::default()
5230                },
5231                CompactionConfig::default(),
5232            );
5233            let store = EventStore::with_config(config);
5234
5235            for i in 0..3 {
5236                let event = Event::from_strings(
5237                    "test.created".to_string(),
5238                    format!("entity-{i}"),
5239                    "default".to_string(),
5240                    serde_json::json!({"i": i}),
5241                    None,
5242                )
5243                .unwrap();
5244                store.ingest(&event).unwrap();
5245            }
5246
5247            store.flush_storage().unwrap();
5248            assert_eq!(store.stats().total_events, 3);
5249        }
5250
5251        // Verify parquet file exists. After Step 1 the file lives
5252        // under <root>/<tenant>/<yyyy-mm>/, so walk recursively.
5253        let parquet_files = find_parquet_files(&storage_dir);
5254        assert!(!parquet_files.is_empty(), "Parquet file must exist");
5255
5256        // Corrupt the parquet file
5257        std::fs::write(&parquet_files[0], b"corrupted data").unwrap();
5258
5259        // Truncate WAL so only Parquet matters
5260        for entry in std::fs::read_dir(&wal_dir).unwrap().flatten() {
5261            std::fs::write(entry.path(), b"").unwrap();
5262        }
5263
5264        // Session 2: reload — should NOT silently report 0 events.
5265        // The error is logged via tracing::error! which we can't capture in a
5266        // unit test, but we CAN verify the store has 0 events (previously this
5267        // looked identical to "no data on disk" — now there's an error log).
5268        // The key behavioral change is that with_config no longer uses a
5269        // let-chain that silently drops the Err variant.
5270        {
5271            let config = EventStoreConfig::production(
5272                &storage_dir,
5273                &wal_dir,
5274                SnapshotConfig::default(),
5275                WALConfig::default(),
5276                CompactionConfig::default(),
5277            );
5278            let store = EventStore::with_config(config);
5279
5280            // Store has 0 events because Parquet is corrupted — but the error
5281            // is now logged (not silently swallowed).
5282            assert_eq!(store.stats().total_events, 0);
5283        }
5284    }
5285
5286    // -----------------------------------------------------------------------
5287    // Step 6: Bounded WAL replay. Each successful checkpoint truncates the
5288    // WAL so cold-start replay is O(one checkpoint interval) regardless of
5289    // total dataset size.
5290    // -----------------------------------------------------------------------
5291
5292    /// Count entries in every WAL file under `wal_dir` (any line that
5293    /// parses as a valid JSON object — the line format is one
5294    /// JSON-serialized WALEntry per line, see WALFile::write_entry).
5295    fn count_wal_entries(wal_dir: &std::path::Path) -> usize {
5296        use std::io::{BufRead, BufReader};
5297        let mut total = 0usize;
5298        let Ok(entries) = std::fs::read_dir(wal_dir) else {
5299            return 0;
5300        };
5301        for entry in entries.flatten() {
5302            let path = entry.path();
5303            if path.extension().is_none_or(|e| e != "log") {
5304                continue;
5305            }
5306            let Ok(file) = std::fs::File::open(&path) else {
5307                continue;
5308            };
5309            for line in BufReader::new(file)
5310                .lines()
5311                .map_while(std::result::Result::ok)
5312            {
5313                if !line.trim().is_empty() {
5314                    total += 1;
5315                }
5316            }
5317        }
5318        total
5319    }
5320
5321    #[test]
5322    fn test_checkpoint_truncates_wal_after_flush() {
5323        // After a successful checkpoint, every previously-ingested event
5324        // should be in Parquet, and the WAL should be empty (truncated).
5325        // This is the load-bearing invariant for Step 6's bounded-replay
5326        // promise — without truncation, the WAL grows unboundedly.
5327        let data_dir = TempDir::new().unwrap();
5328        let storage_dir = data_dir.path().join("storage");
5329        let wal_dir = data_dir.path().join("wal");
5330
5331        let config = EventStoreConfig::production(
5332            &storage_dir,
5333            &wal_dir,
5334            SnapshotConfig::default(),
5335            WALConfig {
5336                sync_on_write: true,
5337                ..WALConfig::default()
5338            },
5339            CompactionConfig::default(),
5340        );
5341        let store = EventStore::with_config(config);
5342
5343        for i in 0..10 {
5344            let event = Event::from_strings(
5345                "test.created".to_string(),
5346                format!("entity-{i}"),
5347                "default".to_string(),
5348                serde_json::json!({"i": i}),
5349                None,
5350            )
5351            .unwrap();
5352            store.ingest(&event).unwrap();
5353        }
5354
5355        // Sanity: all 10 events are in the WAL pre-checkpoint.
5356        assert_eq!(
5357            count_wal_entries(&wal_dir),
5358            10,
5359            "WAL should have 10 events before checkpoint"
5360        );
5361
5362        store.checkpoint().unwrap();
5363
5364        assert_eq!(
5365            count_wal_entries(&wal_dir),
5366            0,
5367            "WAL should be empty after successful checkpoint"
5368        );
5369        let parquet_files = find_parquet_files(&storage_dir);
5370        assert!(!parquet_files.is_empty(), "Parquet should hold the events");
5371    }
5372
5373    #[test]
5374    fn test_replay_only_post_checkpoint_events_after_crash() {
5375        // Headline AC for the bead: write N events, checkpoint, write K
5376        // more, simulate a crash, restart, and verify only K events go
5377        // through replay (not N+K).
5378        //
5379        // Uses small N (50) and K (5) for test speed — the property
5380        // is the same as the spec's 1M+10k example, just scaled down.
5381        let data_dir = TempDir::new().unwrap();
5382        let storage_dir = data_dir.path().join("storage");
5383        let wal_dir = data_dir.path().join("wal");
5384
5385        let config_factory = || {
5386            EventStoreConfig::production(
5387                &storage_dir,
5388                &wal_dir,
5389                SnapshotConfig::default(),
5390                WALConfig {
5391                    sync_on_write: true,
5392                    ..WALConfig::default()
5393                },
5394                CompactionConfig::default(),
5395            )
5396        };
5397
5398        // Session 1: ingest N, checkpoint, ingest K, then drop without
5399        // a graceful shutdown — that's the crash.
5400        const N: usize = 50;
5401        const K: usize = 5;
5402        {
5403            let store = EventStore::with_config(config_factory());
5404            for i in 0..N {
5405                store
5406                    .ingest(
5407                        &Event::from_strings(
5408                            "pre.checkpoint".to_string(),
5409                            format!("e-{i}"),
5410                            "default".to_string(),
5411                            serde_json::json!({"i": i}),
5412                            None,
5413                        )
5414                        .unwrap(),
5415                    )
5416                    .unwrap();
5417            }
5418            store.checkpoint().unwrap();
5419            assert_eq!(
5420                count_wal_entries(&wal_dir),
5421                0,
5422                "WAL should be empty immediately after checkpoint"
5423            );
5424
5425            for i in 0..K {
5426                store
5427                    .ingest(
5428                        &Event::from_strings(
5429                            "post.checkpoint".to_string(),
5430                            format!("p-{i}"),
5431                            "default".to_string(),
5432                            serde_json::json!({"i": i}),
5433                            None,
5434                        )
5435                        .unwrap(),
5436                    )
5437                    .unwrap();
5438            }
5439            assert_eq!(
5440                count_wal_entries(&wal_dir),
5441                K,
5442                "WAL should hold only post-checkpoint events"
5443            );
5444            // Drop without flushing — simulates a crash mid-write.
5445        }
5446
5447        // Session 2: reopen. Recovery should replay only the K post-
5448        // checkpoint events from the WAL — the N pre-checkpoint events
5449        // are durable in Parquet and lazy-loaded on demand.
5450        {
5451            let store = EventStore::with_config(config_factory());
5452            // total_events reflects only WAL-recovered events at boot
5453            // (Step 2 — Parquet stays cold until first per-tenant
5454            // query). So the WAL replay size IS exactly K.
5455            assert_eq!(
5456                store.stats().total_events,
5457                K,
5458                "Boot should replay exactly K events from WAL (the post-checkpoint window), not N+K"
5459            );
5460
5461            // Lazy-load brings the rest in.
5462            store.ensure_tenant_loaded("default").unwrap();
5463            assert_eq!(
5464                store.stats().total_events,
5465                N + K,
5466                "After lazy-load, both pre- and post-checkpoint events should be reachable"
5467            );
5468        }
5469    }
5470
5471    #[test]
5472    fn test_checkpoint_is_idempotent() {
5473        // Calling checkpoint() twice in a row is safe: the second call
5474        // finds an empty WAL and an empty Parquet batch, and no-ops.
5475        let data_dir = TempDir::new().unwrap();
5476        let storage_dir = data_dir.path().join("storage");
5477        let wal_dir = data_dir.path().join("wal");
5478
5479        let store = EventStore::with_config(EventStoreConfig::production(
5480            &storage_dir,
5481            &wal_dir,
5482            SnapshotConfig::default(),
5483            WALConfig::default(),
5484            CompactionConfig::default(),
5485        ));
5486
5487        for i in 0..5 {
5488            store
5489                .ingest(
5490                    &Event::from_strings(
5491                        "x".to_string(),
5492                        format!("e-{i}"),
5493                        "default".to_string(),
5494                        serde_json::json!({}),
5495                        None,
5496                    )
5497                    .unwrap(),
5498                )
5499                .unwrap();
5500        }
5501
5502        store.checkpoint().unwrap();
5503        // Second call is a no-op and must not error.
5504        store.checkpoint().unwrap();
5505        assert_eq!(count_wal_entries(&wal_dir), 0);
5506    }
5507
5508    #[test]
5509    fn test_checkpoint_noop_in_memory_only_mode() {
5510        // Without WAL configured, checkpoint() is a no-op.
5511        let store = EventStore::new();
5512        store.checkpoint().unwrap();
5513    }
5514
5515    #[test]
5516    fn test_checkpoint_interval_from_env_defaults_to_60s_when_wal_enabled() {
5517        let (config, _) = EventStoreConfig::from_env_vars(
5518            Some("/app/data".to_string()),
5519            None,
5520            None,
5521            None,
5522            None,
5523            None,
5524            None,
5525            None,
5526        );
5527        assert_eq!(config.checkpoint_interval_secs, Some(60));
5528    }
5529
5530    #[test]
5531    fn test_checkpoint_interval_from_env_overrides_default() {
5532        let (config, _) = EventStoreConfig::from_env_vars(
5533            Some("/app/data".to_string()),
5534            None,
5535            None,
5536            None,
5537            None,
5538            None,
5539            None,
5540            Some("15".to_string()),
5541        );
5542        assert_eq!(config.checkpoint_interval_secs, Some(15));
5543    }
5544
5545    #[test]
5546    fn test_checkpoint_interval_disabled_when_wal_disabled() {
5547        // No WAL → no checkpoint loop, regardless of env var value.
5548        let (config, _) = EventStoreConfig::from_env_vars(
5549            Some("/app/data".to_string()),
5550            None,
5551            None,
5552            Some("false".to_string()),
5553            None,
5554            None,
5555            None,
5556            Some("15".to_string()),
5557        );
5558        assert_eq!(config.checkpoint_interval_secs, None);
5559    }
5560
5561    #[test]
5562    fn test_checkpoint_interval_unparseable_falls_back_to_default() {
5563        let (config, _) = EventStoreConfig::from_env_vars(
5564            Some("/app/data".to_string()),
5565            None,
5566            None,
5567            None,
5568            None,
5569            None,
5570            None,
5571            Some("not-a-number".to_string()),
5572        );
5573        assert_eq!(config.checkpoint_interval_secs, Some(60));
5574    }
5575
5576    // ── Tenant-scoped stats + state reconstruction (#230) ──────────────────
5577    //
5578    // `stats()` and `reconstruct_state()` are global. The gateway exposes the
5579    // scoped variants below to tenants, so cross-tenant isolation is the
5580    // property that matters most here.
5581
5582    fn seed_two_tenants() -> EventStore {
5583        let store = EventStore::new();
5584
5585        // alice: 3 events, 2 entities, 2 types
5586        for (entity, etype, payload) in [
5587            (
5588                "a-1",
5589                "created",
5590                serde_json::json!({"colour": "red", "size": 1}),
5591            ),
5592            ("a-1", "updated", serde_json::json!({"colour": "blue"})),
5593            ("a-2", "created", serde_json::json!({"colour": "green"})),
5594        ] {
5595            store
5596                .ingest(
5597                    &Event::from_strings(
5598                        etype.to_string(),
5599                        entity.to_string(),
5600                        "alice".to_string(),
5601                        payload,
5602                        None,
5603                    )
5604                    .unwrap(),
5605                )
5606                .unwrap();
5607        }
5608
5609        // bob: 1 event, and it deliberately reuses alice's entity_id "a-1"
5610        store
5611            .ingest(
5612                &Event::from_strings(
5613                    "created".to_string(),
5614                    "a-1".to_string(),
5615                    "bob".to_string(),
5616                    serde_json::json!({"colour": "BOB_SECRET", "bob_only": true}),
5617                    None,
5618                )
5619                .unwrap(),
5620            )
5621            .unwrap();
5622
5623        store
5624    }
5625
5626    #[test]
5627    fn test_stats_for_tenant_counts_only_that_tenant() {
5628        let store = seed_two_tenants();
5629
5630        let alice = store.stats_for_tenant("alice");
5631        assert_eq!(alice.total_events, 3);
5632        assert_eq!(alice.total_entities, 2);
5633        assert_eq!(alice.total_event_types, 2);
5634        assert_eq!(alice.event_types.get("created"), Some(&2));
5635        assert_eq!(alice.event_types.get("updated"), Some(&1));
5636
5637        let bob = store.stats_for_tenant("bob");
5638        assert_eq!(bob.total_events, 1);
5639        assert_eq!(bob.total_entities, 1);
5640        assert_eq!(bob.total_event_types, 1);
5641    }
5642
5643    #[test]
5644    fn test_stats_for_tenant_never_reports_global_totals() {
5645        let store = seed_two_tenants();
5646
5647        // The global view sees all 4; neither tenant may.
5648        assert_eq!(store.stats().total_events, 4);
5649        assert_eq!(store.stats_for_tenant("alice").total_events, 3);
5650        assert_eq!(store.stats_for_tenant("bob").total_events, 1);
5651
5652        // total_ingested is a global process counter, so the scoped form must
5653        // report the tenant's own total rather than leaking the fleet tally.
5654        assert_eq!(store.stats_for_tenant("bob").total_ingested, 1);
5655    }
5656
5657    #[test]
5658    fn test_stats_for_tenant_unknown_tenant_is_empty_not_global() {
5659        let store = seed_two_tenants();
5660        let nobody = store.stats_for_tenant("does-not-exist");
5661
5662        assert_eq!(nobody.total_events, 0);
5663        assert_eq!(nobody.total_entities, 0);
5664        assert!(nobody.event_types.is_empty());
5665        assert!(nobody.oldest_event.is_none());
5666        assert!(nobody.newest_event.is_none());
5667    }
5668
5669    #[test]
5670    fn test_stats_for_tenant_reports_time_range() {
5671        let store = seed_two_tenants();
5672        let alice = store.stats_for_tenant("alice");
5673
5674        let oldest = alice.oldest_event.expect("oldest");
5675        let newest = alice.newest_event.expect("newest");
5676        assert!(oldest <= newest);
5677    }
5678
5679    #[test]
5680    fn test_reconstruct_state_for_tenant_isolates_shared_entity_id() {
5681        let store = seed_two_tenants();
5682
5683        // Both tenants have an entity called "a-1". Each must see only its own.
5684        let alice = store
5685            .reconstruct_state_for_tenant("a-1", None, "alice")
5686            .unwrap();
5687        let alice_state = alice.get("current_state").unwrap();
5688        assert_eq!(alice_state.get("colour").unwrap(), "blue"); // last write wins
5689        assert_eq!(alice_state.get("size").unwrap(), 1);
5690        assert!(
5691            alice_state.get("bob_only").is_none(),
5692            "alice must not see bob's payload keys: {alice_state:?}"
5693        );
5694        assert_eq!(alice.get("event_count").unwrap(), 2);
5695
5696        let bob = store
5697            .reconstruct_state_for_tenant("a-1", None, "bob")
5698            .unwrap();
5699        let bob_state = bob.get("current_state").unwrap();
5700        assert_eq!(bob_state.get("colour").unwrap(), "BOB_SECRET");
5701        assert_eq!(bob.get("event_count").unwrap(), 1);
5702    }
5703
5704    #[test]
5705    fn test_reconstruct_state_for_tenant_rejects_another_tenants_entity() {
5706        let store = seed_two_tenants();
5707
5708        // "a-2" belongs to alice only.
5709        assert!(
5710            store
5711                .reconstruct_state_for_tenant("a-2", None, "alice")
5712                .is_ok()
5713        );
5714        assert!(
5715            store
5716                .reconstruct_state_for_tenant("a-2", None, "bob")
5717                .is_err(),
5718            "bob must not be able to read alice's entity"
5719        );
5720    }
5721
5722    #[test]
5723    fn test_global_reconstruct_state_still_spans_tenants() {
5724        // The unscoped form is unchanged — it is the internal/admin path and the
5725        // gateway never exposes it without a tenant_id.
5726        let store = seed_two_tenants();
5727        let all = store.reconstruct_state("a-1", None).unwrap();
5728        assert_eq!(all.get("event_count").unwrap(), 3);
5729    }
5730
5731    // -----------------------------------------------------------------
5732    // Issue #251: `limit` must bound the store's work, not just its
5733    // response. These guards count `Event` CLONES (crate::clone_probe),
5734    // not returned rows — `query_results_total` is incremented with
5735    // `results.len()`, so it reads 1 both when the store clones one event
5736    // and when it clones the whole history and throws the rest away.
5737    // Revert `query_window` to "clone every match, sort all N, truncate"
5738    // and these go red; row-shaped assertions alone would not.
5739    // -----------------------------------------------------------------
5740
5741    fn seed_hot_entity(history: usize) -> EventStore {
5742        let store = EventStore::new();
5743        for _ in 0..history {
5744            store
5745                .ingest(&create_test_event("entity-hot", "user.updated"))
5746                .unwrap();
5747        }
5748        store
5749    }
5750
5751    fn hot_request(limit: Option<usize>) -> QueryEventsRequest {
5752        QueryEventsRequest {
5753            entity_id: Some("entity-hot".to_string()),
5754            limit,
5755            ..QueryEventsRequest::default()
5756        }
5757    }
5758
5759    #[test]
5760    fn query_window_limit_bounds_materialization_not_just_the_response() {
5761        const HISTORY: usize = 400;
5762        let store = seed_hot_entity(HISTORY);
5763
5764        // The documented "latest event for an entity" read.
5765        let ((page, total), materialized) = crate::clone_probe::measure(|| {
5766            store.query_window(&hot_request(Some(1)), 0, true).unwrap()
5767        });
5768        assert_eq!(page.len(), 1, "limit=1 returns one event");
5769        assert_eq!(total, HISTORY, "total is still the full match count");
5770        assert_eq!(
5771            materialized, 1,
5772            "limit=1 cloned {materialized} of {HISTORY} events: `limit` must \
5773             bound what the store materializes, not just what it returns"
5774        );
5775
5776        // An offset page materializes the page, not everything up to it.
5777        let ((page, _), materialized) = crate::clone_probe::measure(|| {
5778            store
5779                .query_window(&hot_request(Some(5)), 300, false)
5780                .unwrap()
5781        });
5782        assert_eq!(page.len(), 5);
5783        assert_eq!(
5784            materialized, 5,
5785            "offset=300&limit=5 cloned {materialized} events, expected 5"
5786        );
5787
5788        // The unbounded query legitimately materializes everything — the
5789        // ceiling only applies where the caller asked for a window.
5790        let ((all, _), materialized) = crate::clone_probe::measure(|| {
5791            store.query_window(&hot_request(None), 0, true).unwrap()
5792        });
5793        assert_eq!(all.len(), HISTORY);
5794        assert_eq!(materialized, HISTORY as u64);
5795    }
5796
5797    #[test]
5798    fn query_window_cost_does_not_grow_with_entity_history() {
5799        // The same bounded page against a 10x longer history. Materialization
5800        // is bounded by the window, so history length must not show up in the
5801        // clone count — only the O(N) index scan that `total` inherently needs
5802        // still scales, and that clones no events.
5803        let short = seed_hot_entity(40);
5804        let long = seed_hot_entity(400);
5805
5806        let (_, short_cost) = crate::clone_probe::measure(|| {
5807            short.query_window(&hot_request(Some(1)), 0, true).unwrap()
5808        });
5809        let (_, long_cost) = crate::clone_probe::measure(|| {
5810            long.query_window(&hot_request(Some(1)), 0, true).unwrap()
5811        });
5812
5813        assert_eq!(
5814            (short_cost, long_cost),
5815            (1, 1),
5816            "a limit=1 page materialized {short_cost} events over a 40-event \
5817             history and {long_cost} over 400 — cost is tracking history length"
5818        );
5819    }
5820
5821    #[test]
5822    fn select_window_orders_only_the_window() {
5823        // The complexity claim in `select_window`'s doc comment, made
5824        // observable: counting comparisons. A full sort of N items costs
5825        // ~N·log2(N); partial selection of a small window costs ~N. Replace the
5826        // `select_nth_unstable_by` with a plain `sort_unstable_by` over all N
5827        // and the ratio collapses, failing this test.
5828        use std::cell::Cell;
5829        const N: usize = 4096;
5830
5831        let comparisons = Cell::new(0usize);
5832        let order = |a: &u64, b: &u64| {
5833            comparisons.set(comparisons.get() + 1);
5834            a.cmp(b)
5835        };
5836        let shuffled = || -> Vec<u64> {
5837            (0..N as u64)
5838                .map(|i| (i * 2_654_435_761) % 1_000_003)
5839                .collect()
5840        };
5841
5842        let mut bounded = shuffled();
5843        select_window(&mut bounded, 0, Some(1), order);
5844        let bounded_cost = comparisons.replace(0);
5845        assert_eq!(bounded.len(), 1, "window of 1 keeps 1 item");
5846
5847        let mut everything = shuffled();
5848        select_window(&mut everything, 0, None, order);
5849        let full_sort_cost = comparisons.get();
5850        assert_eq!(everything.len(), N);
5851
5852        assert!(
5853            bounded_cost < 4 * N,
5854            "selecting a 1-item window out of {N} took {bounded_cost} comparisons \
5855             (~{}·N) — that is sort-shaped, not selection-shaped",
5856            bounded_cost / N
5857        );
5858        assert!(
5859            bounded_cost * 3 < full_sort_cost,
5860            "a 1-item window cost {bounded_cost} comparisons against \
5861             {full_sort_cost} for sorting all {N}: the window is not bounding \
5862             the ordering work"
5863        );
5864    }
5865
5866    #[test]
5867    fn select_window_is_equivalent_to_sort_then_window() {
5868        // Exhaustive small-input check that partial selection cannot diverge
5869        // from sort-then-window, including duplicate keys (the total order is
5870        // supplied by the caller — here (key, position)).
5871        let order = |a: &(u32, usize), b: &(u32, usize)| a.cmp(b);
5872        let source: Vec<(u32, usize)> = [7, 3, 3, 9, 1, 3, 5, 9, 0, 2]
5873            .into_iter()
5874            .enumerate()
5875            .map(|(i, k)| (k, i))
5876            .collect();
5877
5878        let mut sorted = source.clone();
5879        sorted.sort_by(order);
5880
5881        for offset in 0..12 {
5882            for limit in [None, Some(0), Some(1), Some(3), Some(10), Some(50)] {
5883                let mut got = source.clone();
5884                select_window(&mut got, offset, limit, order);
5885                let got: Vec<_> = got.into_iter().skip(offset).collect();
5886                let expected: Vec<_> = sorted
5887                    .iter()
5888                    .copied()
5889                    .skip(offset)
5890                    .take(limit.unwrap_or(usize::MAX))
5891                    .collect();
5892                assert_eq!(got, expected, "offset={offset} limit={limit:?}");
5893            }
5894        }
5895    }
5896
5897    // `as_of`/`since`/`until` are applied in `filter_entries`, which only runs
5898    // on the three INDEXED branches (entity_id / event_type / event_type_prefix).
5899    // The full-scan branch — a query scoped only by tenant, or only by
5900    // `payload_filter`/`exclude_event_type_prefix`, which is exactly what the
5901    // gateway sends for "this tenant's activity since T" — builds its offsets
5902    // as `(0..events.len())` and never calls it, and `apply_filters` did not
5903    // check timestamps either. So the time window was silently dropped: the
5904    // query returned the whole history, and `total`/`has_more` described that
5905    // whole history too. Same class of "the parameter never reaches the code
5906    // that would honour it" as #250's offset.
5907    #[test]
5908    fn query_window_applies_time_filters_on_the_full_scan_path() {
5909        let store = EventStore::new();
5910        let base = Utc::now() - chrono::Duration::hours(24);
5911        let mut ids = Vec::new();
5912        for i in 0..5i64 {
5913            let mut event = create_test_event(&format!("e-{i}"), "user.created");
5914            event.timestamp = base + chrono::Duration::hours(i);
5915            event.version = i + 1;
5916            ids.push(event.id);
5917            store.ingest(&event).unwrap();
5918        }
5919        let at = |h: i64| base + chrono::Duration::hours(h);
5920
5921        // Tenant-only: no entity_id, no event_type, no event_type_prefix, so
5922        // nothing narrows the scan.
5923        let scoped = |mutate: &dyn Fn(&mut QueryEventsRequest)| {
5924            let mut req = QueryEventsRequest {
5925                tenant_id: Some("default".to_string()),
5926                ..QueryEventsRequest::default()
5927            };
5928            mutate(&mut req);
5929            req
5930        };
5931
5932        for (label, req, expected) in [
5933            (
5934                "since=T+2 keeps only events at or after T+2",
5935                scoped(&|r| r.since = Some(at(2))),
5936                vec![ids[2], ids[3], ids[4]],
5937            ),
5938            (
5939                "until=T+1 keeps only events at or before T+1",
5940                scoped(&|r| r.until = Some(at(1))),
5941                vec![ids[0], ids[1]],
5942            ),
5943            (
5944                "as_of=T+1 is time travel: nothing newer than T+1",
5945                scoped(&|r| r.as_of = Some(at(1))),
5946                vec![ids[0], ids[1]],
5947            ),
5948            (
5949                "since+until compose into a closed window",
5950                scoped(&|r| {
5951                    r.since = Some(at(1));
5952                    r.until = Some(at(3));
5953                }),
5954                vec![ids[1], ids[2], ids[3]],
5955            ),
5956        ] {
5957            let (events, total) = store.query_window(&req, 0, false).unwrap();
5958            let got: Vec<_> = events.iter().map(|e| e.id).collect();
5959            assert_eq!(got, expected, "{label}");
5960            assert_eq!(
5961                total,
5962                expected.len(),
5963                "{label}: total counts the events INSIDE the window — \
5964                 has_more is derived from it, so a full-history total makes a \
5965                 paginator walk events the caller filtered out"
5966            );
5967        }
5968
5969        // The indexed branches must keep agreeing with the full scan — the same
5970        // window asked for with an entity filter is the same window.
5971        let (indexed, total) = store
5972            .query_window(
5973                &scoped(&|r| {
5974                    r.entity_id = Some("e-3".to_string());
5975                    r.since = Some(at(2));
5976                }),
5977                0,
5978                false,
5979            )
5980            .unwrap();
5981        assert_eq!(
5982            indexed.iter().map(|e| e.id).collect::<Vec<_>>(),
5983            vec![ids[3]]
5984        );
5985        assert_eq!(total, 1);
5986
5987        // And the window composes with the #251 selection: a bounded page taken
5988        // out of a time window is a page of that window, not of the history.
5989        let (page, total) = store
5990            .query_window(&scoped(&|r| r.since = Some(at(2))), 1, true)
5991            .unwrap();
5992        assert_eq!(
5993            page.iter().map(|e| e.id).collect::<Vec<_>>(),
5994            vec![ids[3], ids[2]],
5995            "order=desc + offset=1 inside a since window"
5996        );
5997        assert_eq!(total, 3);
5998    }
5999
6000    #[test]
6001    fn query_window_bounded_selection_matches_a_full_sort() {
6002        // Bounded selection (partial select + sort of the window only) must be
6003        // observationally identical to sort-everything-then-window, including
6004        // ties and the descending path — otherwise the perf win silently
6005        // changes pagination.
6006        let store = EventStore::new();
6007        for i in 0..50 {
6008            store
6009                .ingest(&create_test_event(&format!("e-{i:02}"), "user.created"))
6010                .unwrap();
6011        }
6012        let all = |descending: bool| {
6013            let (events, _) = store
6014                .query_window(&QueryEventsRequest::default(), 0, descending)
6015                .unwrap();
6016            events
6017        };
6018
6019        for descending in [false, true] {
6020            let reference = all(descending);
6021            for offset in [0, 1, 7, 49, 50, 100] {
6022                for limit in [1, 3, 10, 50, 100] {
6023                    let (page, total) = store
6024                        .query_window(
6025                            &QueryEventsRequest {
6026                                limit: Some(limit),
6027                                ..QueryEventsRequest::default()
6028                            },
6029                            offset,
6030                            descending,
6031                        )
6032                        .unwrap();
6033                    let expected: Vec<_> = reference
6034                        .iter()
6035                        .skip(offset)
6036                        .take(limit)
6037                        .map(|e| e.id)
6038                        .collect();
6039                    let got: Vec<_> = page.iter().map(|e| e.id).collect();
6040                    assert_eq!(
6041                        got, expected,
6042                        "desc={descending} offset={offset} limit={limit}: windowed \
6043                         selection must match a full sort"
6044                    );
6045                    assert_eq!(total, 50, "total is always the full match count");
6046                }
6047            }
6048        }
6049    }
6050}