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