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