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§
Sourcefn 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 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.
Sourcefn 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 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.
Sourcefn 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 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.
Sourcefn 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 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.
Sourcefn 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 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.
Sourcefn 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<'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.
Sourcefn 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 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.
Sourcefn 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 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.
Sourcefn 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 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.
Sourcefn 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 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.
Sourcefn 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 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.
Sourcefn 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 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.
Sourcefn 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 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.
Provided Methods§
Sourcefn tenant(&self) -> &str
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".