Skip to main content

EventStore

Trait EventStore 

Source
pub trait EventStore:
    Send
    + Sync
    + Debug {
Show 15 methods // Required methods fn buffer<'life0, 'life1, 'async_trait>( &'life0 self, event: &'life1 InboundEvent, at: Timestamp, ) -> Pin<Box<dyn Future<Output = Result<bool, StoreError>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn subscribe<'life0, 'life1, 'async_trait>( &'life0 self, sub: &'life1 Subscription, at: Timestamp, ) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn claim_for<'life0, 'life1, 'async_trait>( &'life0 self, sub: &'life1 Subscription, at: Timestamp, ) -> Pin<Box<dyn Future<Output = Result<Option<BufferedEvent>, StoreError>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn match_waiter<'life0, 'life1, 'async_trait>( &'life0 self, event: &'life1 InboundEvent, at: Timestamp, ) -> Pin<Box<dyn Future<Output = Result<Option<Subscription>, StoreError>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn deliver_to<'life0, 'life1, 'async_trait>( &'life0 self, run: RunId, event: &'life1 InboundEvent, at: Timestamp, ) -> Pin<Box<dyn Future<Output = Result<TargetedDelivery, StoreError>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn unsubscribe<'life0, 'async_trait>( &'life0 self, run: RunId, effect: EffectKey, ) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn unsubscribe_run<'life0, 'life1, 'async_trait>( &'life0 self, run: RunId, unanswered: &'life1 [EffectKey], ) -> Pin<Box<dyn Future<Output = Result<Retired, StoreError>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn park_wait<'life0, 'life1, 'async_trait>( &'life0 self, sub: &'life1 Subscription, at: Timestamp, ) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn parked_waits<'life0, 'async_trait>( &'life0 self, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<Subscription>, StoreError>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn erase_payload<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, source: &'life1 str, id: &'life2 str, ) -> Pin<Box<dyn Future<Output = Result<bool, StoreError>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait; fn minter<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, source: &'life1 str, id: &'life2 str, ) -> Pin<Box<dyn Future<Output = Result<Option<Minter>, StoreError>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait; fn sweep_unclaimed<'life0, 'life1, 'async_trait>( &'life0 self, older_than: Timestamp, reason: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<usize, StoreError>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait; fn dead_letters<'life0, 'async_trait>( &'life0 self, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<DeadLetter>, StoreError>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; fn waiting<'life0, 'async_trait>( &'life0 self, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<Subscription>, StoreError>> + Send + 'async_trait>> where Self: 'async_trait, 'life0: 'async_trait; // Provided method fn tenant(&self) -> &str { ... }
}

Required Methods§

Source

fn buffer<'life0, 'life1, 'async_trait>( &'life0 self, event: &'life1 InboundEvent, at: Timestamp, ) -> Pin<Box<dyn Future<Output = Result<bool, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Record an inbound event, returning false if this (source, id) was already seen.

Deduplication is by the pair, never the bare id — an id is unique only within one producer, and keying on it alone lets two counterparties swallow each other’s messages as apparent retries. The pair is InboundEvent::dedup_key, the one implementation of that identity.

Source

fn subscribe<'life0, 'life1, 'async_trait>( &'life0 self, sub: &'life1 Subscription, at: Timestamp, ) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Register a run’s interest in a future event.

Source

fn claim_for<'life0, 'life1, 'async_trait>( &'life0 self, sub: &'life1 Subscription, at: Timestamp, ) -> Pin<Box<dyn Future<Output = Result<Option<BufferedEvent>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Atomically find and claim a buffered event matching a subscription.

Claiming is what makes an event single-delivery: two runs waiting on the same key cannot both consume one message.

An event already claimed by this subscription’s run is returned again rather than filtered out. That is crash recovery, not a second delivery: match_waiter claims durably and the run resumes in a separate step, so a crash between the two leaves an event claimed for a run that never saw it — and a claim_for that hid the run’s own claim from it would strand the wait until its deadline breached, losing a message that arrived in time. Single delivery is untouched, because only the claiming run can re-claim.

Source

fn match_waiter<'life0, 'life1, 'async_trait>( &'life0 self, event: &'life1 InboundEvent, at: Timestamp, ) -> Pin<Box<dyn Future<Output = Result<Option<Subscription>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Atomically find and claim a subscription matching an arrived event.

The mirror of claim_for: one looks for a waiter given an event, the other for an event given a waiter. Both directions are needed precisely because either can arrive first.

A subscription whose run is sealed is passed over: the event goes to the next live waiter, or stays buffered, rather than being claimed for a run that will never consume it. So is a parked one (park_wait): it already holds its claimed event, and a second event elected for it would stay claimed for a satisfied wait for ever.

Source

fn deliver_to<'life0, 'life1, 'async_trait>( &'life0 self, run: RunId, event: &'life1 InboundEvent, at: Timestamp, ) -> Pin<Box<dyn Future<Output = Result<TargetedDelivery, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Atomically buffer event and claim it for the matching subscription of exactly run.

Unlike buffer followed by match_waiter, a failed targeted delivery leaves no unclaimed event behind for another run to consume. Implementations must decide duplicate/not-waiting/matched in one transaction.

Source

fn unsubscribe<'life0, 'async_trait>( &'life0 self, run: RunId, effect: EffectKey, ) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Drop a subscription once it has been satisfied.

Also the point where the buffer’s copy of the run’s delivered payloads is shed. Stripping earlier — at the claim — would lose the payload for a run that crashed between claim and resume, whose recovery re-reads it from the buffer; unsubscribe is the run’s own signal that the wait is over, so the buffer row keeps only its (source, id) identity for dedup from here on.

That places an ordering obligation on the caller: journal the delivered payload before unsubscribing. The delivery worker does — it appends EffectDone and only then retires the subscription. A caller that unsubscribes first has a crash window in which the buffer copy is gone and the journal copy never landed, and recovery then resumes the wait on a stripped row. What stripping here does not cover: unclaimed and dead-lettered rows, which never reach an unsubscribe and are erased through erase_payload instead.

Source

fn unsubscribe_run<'life0, 'life1, 'async_trait>( &'life0 self, run: RunId, unanswered: &'life1 [EffectKey], ) -> Pin<Box<dyn Future<Output = Result<Retired, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Drop every subscription a run holds, and hand back what it held claimed for a wait it never answered.

Called when the run concludes closed. Left registered, a closed run’s wait is the oldest waiter on its key, and the next matching event is claimed for a run that will never consume it — while a live run waiting on the same key starves.

unanswered names the run’s waits whose answer its journal does not hold. A message claimed for one of them reached nobody: it goes back to the buffer unclaimed, with its payload, and is returned so the caller can offer it to the next waiter. Shedding it would lose a message that arrived in time for a run that happened to conclude first. What the run holds claimed for an answered wait is shed, as unsubscribe sheds it — releasing that would let a second run consume a message the first already journaled.

Source

fn park_wait<'life0, 'life1, 'async_trait>( &'life0 self, sub: &'life1 Subscription, at: Timestamp, ) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Register a wait that already holds a claimed event, and mark it for redelivery.

The delivery that claimed an event and could not resume the run — its owner held the lease past the retry window — parks the pair here, and parked_waits is what the sweep’s redelivery pass walks. A targeted delivery’s claim is parked by deliver_to itself. The mark goes with the subscription: unsubscribe and the claim that retires a wait both clear it.

A parked wait is listed for redelivery and recovers its own event through claim_for, but it is not matchable: no other event is claimed for it.

Source

fn parked_waits<'life0, 'async_trait>( &'life0 self, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<Subscription>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Waits holding a claimed event nothing has delivered yet.

A separate listing from waiting because that one is every registered wait in registration order, and a plane is expected to hold far more legitimately long waits than a redelivery page: the parked pair registered after them would never be reached.

A sealed run’s parked waits are retired rather than listed — it can record no delivery, so listing them has the pass fail on them every tick — and the payloads it held claimed are shed with them.

Source

fn erase_payload<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, source: &'life1 str, id: &'life2 str, ) -> Pin<Box<dyn Future<Output = Result<bool, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Remove one buffered event’s payload while keeping its identity.

The erasure verb the buffer was missing. The buffer keeps its copy of an inbound payload indefinitely — claimed rows stay for dedup, and dead-lettered rows stay for the operator — so a message that becomes the object of an erasure request had no path to erasure at all: the journal’s copy is key-erasable, the buffer’s was immortal.

The row survives; only the payload goes. Dedup needs exactly (source, id), so a replay of the erased message is still refused rather than accepted as new — erasure must not reopen the door it closed. Dead-letter entries likewise keep their identity, correlation keys and reason, because “what went unclaimed and why” is operational truth about the deployment, not the counterparty’s content.

A row nobody had claimed also leaves the claimable set: it becomes a dead letter whose reason is ERASED_REASON, in the same write. Left live, the next matching waiter would claim the emptied row and be handed null as if the counterparty had sent it.

So does a row claimed and not yet delivered — claimed, with its payload not yet shed by the claimant’s unsubscribe. The claim is released and the claimant’s wait that was parked holding it is unparked, so the wait stays open for its deadline to bound; the recovery that would have re-read the row finds nothing to hand over. Every erased row is marked erased, and claim_for never returns one, whatever its claim says.

What this does not cover: the journaled copy of a delivered payload, which lives under the run’s case and is erased by that case’s key; and the correlation keys, which are business identifiers the row is filed under, not content. Returns whether a row existed.

Source

fn minter<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, source: &'life1 str, id: &'life2 str, ) -> Pin<Box<dyn Future<Output = Result<Option<Minter>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Who minted the buffered event (source, id), or None when no such event was ever buffered.

Survives both payload erasure and the unsubscribe that sheds a delivered payload, as the minter does on the row: who on this plane minted a message is not the counterparty’s content. It is how a settler that meets a duplicate answer to a human task learns whose answer is already on record.

Source

fn sweep_unclaimed<'life0, 'life1, 'async_trait>( &'life0 self, older_than: Timestamp, reason: &'life1 str, ) -> Pin<Box<dyn Future<Output = Result<usize, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Move events nobody claimed within the window to the dead-letter list.

Dead-lettering deliberately happens here and not on arrival: “nobody is waiting yet” and “nobody will ever want this” are different claims, and only the second is safe to act on. Returns how many were retired.

Source

fn dead_letters<'life0, 'async_trait>( &'life0 self, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<DeadLetter>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Events that aged out unclaimed, newest first.

A non-empty list means a correlation key is wrong somewhere. That is the failure which otherwise presents as a process silently never completing.

Source

fn waiting<'life0, 'async_trait>( &'life0 self, limit: usize, ) -> Pin<Box<dyn Future<Output = Result<Vec<Subscription>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Runs currently waiting, for operational visibility.

Provided Methods§

Source

fn tenant(&self) -> &str

Whose rows this handle can reach.

Defaults to TenantId::DEFAULT, the tenant a store serves until told otherwise. Override it with the tenant the handle is actually scoped to.

This exists so a mismatch with the plane’s tenant is a startup refusal. When a key ring is wired, build() seals this state under the plane’s tenant while the store writes rows under its own; the two disagreeing is not a leak — the scopes simply differ — but it seals the state under a scope erasure will never destroy. That is an erasure that reports success and misses, which is the one failure a deletion guarantee cannot have.

Dyn Compatibility§

This trait is dyn compatible.

In older versions of Rust, dyn compatibility was called "object safety".

Implementors§