Skip to main content

EventStore

Struct EventStore 

Source
pub struct EventStore { /* private fields */ }

Implementations§

Source§

impl EventStore

Source

pub fn new() -> Self

Create a new in-memory event store

Source

pub fn with_config(config: EventStoreConfig) -> Self

Create event store with custom configuration

Source

pub fn is_read_only(&self) -> bool

Whether this store was opened read-only (replica mode).

Source

pub fn ingest_with_expected_version( &self, event: &Event, expected_version: Option<u64>, ) -> Result<u64>

Ingest a new event with optional optimistic concurrency check.

If expected_version is Some(v), the write is rejected with VersionConflict unless the entity’s current version equals v. The version check and WAL append are atomic (locked together).

Returns the new entity version after the append.

Source

pub fn ingest(&self, event: &Event) -> Result<()>

Ingest a new event into the store

Source

pub fn ingest_batch(&self, batch: Vec<Event>) -> Result<()>

Ingest a batch of events with a single write lock acquisition.

All events are validated first. If any event fails validation, no events are stored (all-or-nothing validation). Events are then written to WAL, indexed, processed through projections, and pushed to the events vector under a single write lock.

Source

pub fn ingest_replicated(&self, event: &Event) -> Result<()>

Ingest a replicated event from the leader (follower mode).

Unlike ingest(), this method:

  • Skips WAL writing (the follower’s WalReceiver manages its own local WAL)
  • Skips schema validation (the leader already validated)
  • Still indexes, processes projections/pipelines, and broadcasts to WebSocket clients
Source

pub fn get_entity_version(&self, entity_id: &str) -> u64

Get the current version for an entity (number of events appended for it). Returns 0 if the entity has no events.

Source

pub fn consumer_registry(&self) -> &ConsumerRegistry

Get the consumer registry for durable subscriptions.

Source

pub fn subscribe_events(&self) -> Receiver<Arc<Event>>

Subscribe to every successfully-ingested event in this store.

Returns a tokio::sync::broadcast::Receiver that yields an Arc<Event> for each ingest. Always available — does not require the server feature. Lagging receivers surface RecvError::Lagged.

Source

pub fn set_consumer_registry(&mut self, registry: Arc<ConsumerRegistry>)

Replace the default in-memory consumer registry with a durable one.

Called during startup when system repositories are available, so that consumer cursors survive Core restarts via WAL persistence.

Source

pub fn total_events(&self) -> usize

Get the total number of events in the store (used as max offset for consumer ack).

Source

pub fn events_after_offset( &self, offset: u64, filters: &[String], limit: usize, ) -> Vec<(u64, Event)>

Get events after a given offset, optionally filtered by event type prefixes. Used by consumer polling to fetch unprocessed events.

Source

pub fn websocket_manager(&self) -> Arc<WebSocketManager> ⓘ

Get the WebSocket manager for this store

Source

pub fn snapshot_manager(&self) -> Arc<SnapshotManager> ⓘ

Get the snapshot manager for this store

Source

pub fn compaction_manager(&self) -> Option<Arc<CompactionManager>>

Get the compaction manager for this store

Source

pub fn schema_registry(&self) -> Arc<SchemaRegistry> ⓘ

Get the schema registry for this store (v0.5 feature)

Source

pub fn replay_manager(&self) -> Arc<ReplayManager> ⓘ

Get the replay manager for this store (v0.5 feature)

Source

pub fn pipeline_manager(&self) -> Arc<PipelineManager> ⓘ

Get the pipeline manager for this store (v0.5 feature)

Source

pub fn metrics(&self) -> Arc<MetricsRegistry> ⓘ

Get the metrics registry for this store (v0.6 feature)

Source

pub fn projection_manager(&self) -> RwLockReadGuard<'_, ProjectionManager>

Get the projection manager for this store (v0.7 feature)

Source

pub fn register_projection(&self, projection: Arc<dyn Projection>)

Register a custom projection at runtime.

The projection will receive all future events via process(). Historical events are not replayed — only events ingested after registration will be processed by this projection.

See register_projection_with_backfill to also process historical events.

Source

pub fn register_projection_with_backfill( &self, projection: &Arc<dyn Projection>, ) -> Result<()>

Register a custom projection and replay all existing events through it.

After registration, the projection will also receive all future events. Historical events are replayed under a read lock — the projection’s internal state (typically DashMap) handles concurrent access.

Replay is ordered by (timestamp, version). The in-memory pile can be physically out of order when Parquet is hydrated after WAL recovery — the WAL tail holds the newest events while Parquet holds the older history — so the backfill must sort before replaying. Projections with last-write-wins merge semantics produce wrong state otherwise.

Source

pub fn hydrate_all_from_storage(&self) -> Result<usize>

Eagerly reconstruct the in-memory event pile from the full Parquet archive.

The default boot path (Step 2, issue #160) keeps Parquet cold and hydrates tenants lazily on first query — the multi-tenant server cannot fit every tenant in memory. Embedded single-store consumers like Prime are the opposite case: their projections are the queryable surface and never trigger the lazy query path, so the projections must be backfilled from the complete history.

Call this before registering projections. The dedupe in append_loaded_event makes it safe to run after WAL recovery — events already replayed from the WAL are not double-counted. No-op (and Ok(0)) when no Parquet storage is configured, e.g. in-memory test mode.

Marks every tenant present in the archive as loaded, so a later ensure_tenant_loaded for it takes the warm path.

Returns the number of events newly loaded from Parquet.

Source

pub fn projection_state_cache(&self) -> Arc<DashMap<String, Value>> ⓘ

Get the projection state cache for this store (v0.7 feature) Used by Elixir Query Service for state synchronization

Source

pub fn projection_status(&self) -> Arc<DashMap<String, String>> ⓘ

Get the projection status map (v0.13 feature)

Source

pub fn geo_index(&self) -> Arc<GeoIndex> ⓘ

Get the webhook registry for this store (v0.11 feature) Geospatial index for coordinate-based queries (v2.0 feature)

Source

pub fn exactly_once(&self) -> Arc<ExactlyOnceRegistry> ⓘ

Exactly-once processing registry (v2.0 feature)

Source

pub fn schema_evolution(&self) -> Arc<SchemaEvolutionManager> ⓘ

Schema evolution manager (v2.0 feature)

Source

pub fn snapshot_events(&self) -> Vec<Event>

Get a read-locked snapshot of all events (for EventQL/GraphQL queries).

Returns an Arc reference to the internal events vec, avoiding a full clone. The caller holds a read lock for the duration of the Arc lifetime — prefer short-lived usage.

Source

pub fn compact_entity_tokens( &self, entity_id: &str, token_event_type: &str, merged_event: Event, ) -> Result<bool>

Compact token events for an entity by replacing matching events with a single merged event. Used by the embedded streaming feature.

Returns Ok(true) if compaction was performed, Ok(false) if no matching events were found.

Note: The merged event is processed through projections without clearing the removed events’ projection state first. Projections that accumulate state (e.g., counters) should be designed to handle this (the merged event replaces individual tokens, not adds to them).

The merged event is a derived view and is never made durable. Compaction reclaims memory; it does not rewrite history. The token events stay in the WAL and Parquet exactly as they were ingested, so a restart — or any reader that goes to the durable log — sees the original stream, not the merge. A caller that needs the merged result to survive must ingest it as an event of its own.

The write lock is held for the swap + WAL write + index rebuild. The index rebuild is O(N) over all events, which is acceptable for embedded workloads but should not be called in hot paths for large stores.

Source

pub fn webhook_registry(&self) -> Arc<WebhookRegistry> ⓘ

Source

pub fn set_webhook_tx(&self, tx: UnboundedSender<WebhookDeliveryTask>)

Set the channel for async webhook delivery. Called during server startup to wire the delivery worker.

Source

pub fn flush_storage(&self) -> Result<()>

Manually flush any pending events to persistent storage

Source

pub fn checkpoint(&self) -> Result<()>

Run a checkpoint: flush pending Parquet batches, then truncate the WAL through the checkpoint point (Step 6 of the sustainable data strategy).

Order matters. We flush Parquet first; only on success do we truncate the WAL. If the process crashes between the flush and the truncate, the WAL still contains the events that were just durably written, and recovery will replay them. The dedupe in append_loaded_event (index probe) makes that idempotent — the event is already in Parquet, so the lazy-load splice no-ops once the tenant is hydrated.

The reverse order would be unsafe: a crash between truncate and flush would lose committed events.

This bounds dirty-restart replay time to one checkpoint interval regardless of total dataset size — that’s the load-bearing property for cold-start time as ingest rate grows.

No-op when no WAL is configured (in-memory-only mode).

Source

pub fn refresh_storage_metrics_now(&self)

Public entrypoint to populate the on-disk storage gauges once, e.g. at boot so the dashboard’s “storage” card is correct before the first checkpoint tick (and even when the checkpoint loop is disabled). Delegates to the internal refresh. No-op without persistent storage.

Source

pub fn checkpoint_interval(&self) -> Option<Duration>

Get the configured checkpoint cadence (used by background tasks).

Source

pub fn ensure_tenant_loaded(&self, tenant_id: &str) -> Result<()>

Hydrate tenant_id’s persisted Parquet data into the in-memory pile if it isn’t already loaded. Cheap on the warm path (DashMap probe); on the cold path it walks just that tenant’s subtree (load_events_for_tenant) and splices the events into events/index/projections/entity_versions.

Concurrent first-callers for the same tenant serialize on a per-tenant Mutex (singleflight) so the disk read happens once. Other tenants are unaffected — distinct lock per tenant.

Returns Err if the tenant_id fails the path-safety whitelist, the Parquet read fails, or another in-flight load holds the lock past the configured timeout. The caller (a query handler) is expected to surface that as a 5xx — see Step 2’s “no infinite hangs” acceptance criterion.

On failure, loaded is NOT marked, so a transient error is retried on the next request rather than poisoning the tenant permanently. A future commit may add a circuit breaker if thrash becomes an issue.

No-op (and Ok) when no Parquet storage is configured — the in-memory-only mode used by tests has nothing to hydrate.

Source

pub fn is_tenant_loaded(&self, tenant_id: &str) -> bool

True iff this tenant is resident: ensure_tenant_loaded succeeded for it, or hydrate_all_from_storage found it in the archive, and it has not been evicted since. Diagnostic / testing API.

Source

pub fn evict_tenant(&self, tenant_id: &str)

Drop tenant_id from the in-memory cache. Step 3 #2 of the sustainable data strategy.

Removes every event for this tenant from the events Vec, rebuilds the index/entity_versions for the retained events (Vec offsets shift on remove, so the index has to be rebuilt), and resets the tenant_loader bookkeeping so a subsequent query triggers a fresh ensure_tenant_loaded.

Parquet is canonical, in-memory is just cache. This only affects the in-memory side. Disk data is untouched — that’s why eviction is safe even for tenants with recently-ingested data: a query after eviction transparently re-reads from Parquet.

Projection state is NOT rolled back. Projections accumulate across boots and tenants (their durability story is separate); subtracting them would need replay support that doesn’t exist in this commit. After eviction + re-load, projections may double-count the re-loaded events. Step 3 #4’s stress test only asserts the cache budget is held; a future commit will tackle projection-aware eviction.

Locking: takes the events write lock for the full duration of the filter + re-index. Concurrent ingest blocks until done. Eviction is the cold path; the working set should stay in budget so this rarely fires.

Source

pub fn tenant_resident_bytes(&self, tenant_id: &str) -> u64

Approximate resident bytes a single tenant occupies in the in-memory cache. Step 3 budget-tracking input. 0 for cold or evicted tenants.

Source

pub fn cache_resident_bytes(&self) -> u64

Sum of resident-byte estimates across every loaded tenant. What the budget check compares against.

Source

pub fn create_snapshot(&self, entity_id: &str) -> Result<()>

Manually create a snapshot for an entity

Source

pub fn reset_projection(&self, name: &str) -> Result<usize>

Reset a projection by clearing its state and reprocessing all events

Source

pub fn get_event_by_id(&self, event_id: &Uuid) -> Result<Option<Event>>

Get a single event by its UUID

Source

pub fn query(&self, request: &QueryEventsRequest) -> Result<Vec<Event>>

Query events based on filters (optimized with indices).

Ascending by (timestamp, version), windowed by request.limit.

Source

pub fn query_scoped( &self, request: &QueryEventsRequest, scope: &ReadScope, ) -> Result<Vec<Event>>

Self::query under a read scope the request body cannot widen.

Source

pub fn query_window( &self, request: &QueryEventsRequest, offset: usize, descending: bool, ) -> Result<(Vec<Event>, usize)>

Query events and return only the requested window, plus the total number of matches the window was taken from.

request.limit bounds the window; offset skips that many matches first; descending returns newest first. The returned total is the pre-window match count, which is what total_count/has_more need.

Cost (issue #251). Matching still costs an O(N) index scan over the entity’s/type’s entries — that is inherent to reporting total, and this method does not pretend otherwise. What limit does bound is the ordering and the materialization: matches are held as borrowed references, the offset + limit window is partitioned off with select_nth_unstable_by (O(N)) and only that window is sorted (O(k log k)) and cloned. So a limit=1 read over a 100k-event entity scans 100k index entries but sorts and clones one event, where it previously sorted and cloned 100k.

Source

pub fn query_window_scoped( &self, request: &QueryEventsRequest, offset: usize, descending: bool, scope: &ReadScope, ) -> Result<(Vec<Event>, usize)>

Self::query_window under a ReadScope.

Every read path in this store funnels through here, which is the point: a scope enforced in one API surface leaves the others open, and the surfaces that matter (HTTP, the embedded API, MCP) all end up on this method. The scope is a separate argument rather than a field on QueryEventsRequest because that type is deserialized from the request body — a caller could then widen its own scope by omitting the field.

Source

pub fn reconstruct_state( &self, entity_id: &str, as_of: Option<DateTime<Utc>>, ) -> Result<Value>

Reconstruct entity state as of a specific timestamp v0.2: Now uses snapshots for fast reconstruction

Source

pub fn get_snapshot(&self, entity_id: &str) -> Result<Value>

Get snapshot from projection (faster than reconstructing)

Source

pub fn stats(&self) -> StoreStats

Get statistics about the event store

Source

pub fn list_streams(&self) -> Vec<StreamInfo>

Get all unique streams (entity_ids) in the store

Source

pub fn list_event_types(&self) -> Vec<EventTypeInfo>

Get all unique event types in the store

Source

pub fn list_streams_for_tenant(&self, tenant_id: &str) -> Vec<StreamInfo>

Distinct entities (+ per-entity event count) for ONE tenant.

Source

pub fn list_event_types_for_tenant(&self, tenant_id: &str) -> Vec<EventTypeInfo>

Distinct event types (+ per-type event count) for ONE tenant.

Source

pub fn stats_for_tenant(&self, tenant_id: &str) -> TenantStoreStats

Event-store statistics for ONE tenant.

stats() above is global (events.len() + the global index), so it must never be served to a tenant — it would leak whole-fleet totals. This is the scoped equivalent the gateway exposes as GET /api/v1/stats with an auth-derived tenant_id (#230).

Returns the same four counters as StoreStats, plus the per-type census and time range that callers otherwise reconstruct by paging events.

Source

pub fn reconstruct_state_for_tenant( &self, entity_id: &str, as_of: Option<DateTime<Utc>>, tenant_id: &str, ) -> Result<Value>

Reconstruct entity state for ONE tenant.

Deliberately does NOT use the snapshot fast path that reconstruct_state takes: Snapshot carries no tenant_id and the snapshot manager is keyed by entity_id alone, so seeding from a snapshot could fold another tenant’s state into this answer when two tenants share an entity_id. Folding tenant-filtered events is slower and fails closed (#230).

Source

pub fn enable_wal_replication(&self, tx: Sender<WALEntry>)

Attach a broadcast sender to the WAL for replication.

Thread-safe: can be called through Arc<EventStore> at runtime. Used during initial setup and during follower → leader promotion. When set, every WAL append publishes the entry to the broadcast channel so the WAL shipper can stream it to followers.

Source

pub fn wal(&self) -> Option<&Arc<WriteAheadLog>>

Get a reference to the WAL (if configured). Used by the replication catch-up protocol to determine oldest available offset.

Source

pub fn parquet_storage(&self) -> Option<&Arc<RwLock<ParquetStorage>>>

Get a reference to the Parquet storage (if configured). Used by the replication catch-up protocol to stream snapshot files to followers.

Trait Implementations§

Source§

impl Default for EventStore

Source§

fn default() -> Self

Returns the “default value” for a type. Read more

Auto Trait Implementations§

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<A, B, T> HttpServerConnExec<A, B> for T
where B: Body,

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T> Pointable for T

Source§

const ALIGN: usize

The alignment of pointer.
Source§

type Init = T

The type for initializers.
Source§

unsafe fn init(init: <T as Pointable>::Init) -> usize

Initializes a with the given initializer. Read more
Source§

unsafe fn deref<'a>(ptr: usize) -> &'a T

Dereferences the given pointer. Read more
Source§

unsafe fn deref_mut<'a>(ptr: usize) -> &'a mut T

Mutably dereferences the given pointer. Read more
Source§

unsafe fn drop(ptr: usize)

Drops the object pointed to by the given pointer. Read more
Source§

impl<T> PolicyExt for T
where T: ?Sized,

Source§

fn and<P, B, E>(self, other: P) -> And<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow only if self and other return Action::Follow. Read more
Source§

fn or<P, B, E>(self, other: P) -> Or<T, P>
where T: Sized + Policy<B, E>, P: Policy<B, E>,

Create a new Policy that returns Action::Follow if either self or other returns Action::Follow. Read more
Source§

impl<T> Same for T

Source§

type Output = T

Should always be Self
Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more