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