pub trait PushStore:
Send
+ Sync
+ Debug {
// Required methods
fn put<'life0, 'life1, 'async_trait>(
&'life0 self,
config: &'life1 PushConfig,
next_seq: Seq,
) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>
where Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait;
fn get<'life0, 'life1, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<PushConfig>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait;
fn list<'life0, 'async_trait>(
&'life0 self,
task: RunId,
) -> Pin<Box<dyn Future<Output = Result<Vec<PushConfig>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait,
'life0: 'async_trait;
fn due<'life0, 'async_trait>(
&'life0 self,
at: u64,
limit: usize,
) -> Pin<Box<dyn Future<Output = Result<Vec<PushRegistration>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait,
'life0: 'async_trait;
fn advance<'life0, 'life1, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
next_seq: Seq,
) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>
where Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait;
fn retry<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
next_attempt_at: u64,
error: &'life2 str,
) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>
where Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait;
fn park<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
error: &'life2 str,
) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>
where Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait;
fn parked<'life0, 'async_trait>(
&'life0 self,
limit: usize,
) -> Pin<Box<dyn Future<Output = Result<Vec<PushRegistration>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait,
'life0: 'async_trait;
fn unpark<'life0, 'life1, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
at: u64,
) -> Pin<Box<dyn Future<Output = Result<bool, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait;
fn delete<'life0, 'life1, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>
where Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait;
// Provided methods
fn tenant(&self) -> &str { ... }
fn due_in<'life0, 'async_trait>(
&'life0 self,
at: u64,
limit: usize,
namespace: PushNamespace,
) -> Pin<Box<dyn Future<Output = Result<DueBatch, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait,
'life0: 'async_trait { ... }
}Expand description
Durable storage for webhook registrations.
Durable because a task outlives the connection that created it, and a registration that vanished on restart would leave a peer waiting for a notification nobody remembers promising.
Required Methods§
Sourcefn put<'life0, 'life1, 'async_trait>(
&'life0 self,
config: &'life1 PushConfig,
next_seq: Seq,
) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn put<'life0, 'life1, 'async_trait>(
&'life0 self,
config: &'life1 PushConfig,
next_seq: Seq,
) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Register or replace a configuration.
Replacement preserves the existing acknowledgement cursor. Changing a URL or credentials is not permission to discard updates that receiver has not accepted.
§Errors
If the store cannot be reached.
Sourcefn get<'life0, 'life1, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<PushConfig>, StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn get<'life0, 'life1, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
) -> Pin<Box<dyn Future<Output = Result<Option<PushConfig>, StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Sourcefn list<'life0, 'async_trait>(
&'life0 self,
task: RunId,
) -> Pin<Box<dyn Future<Output = Result<Vec<PushConfig>, StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn list<'life0, 'async_trait>(
&'life0 self,
task: RunId,
) -> Pin<Box<dyn Future<Output = Result<Vec<PushConfig>, StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Sourcefn due<'life0, 'async_trait>(
&'life0 self,
at: u64,
limit: usize,
) -> Pin<Box<dyn Future<Output = Result<Vec<PushRegistration>, StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn due<'life0, 'async_trait>(
&'life0 self,
at: u64,
limit: usize,
) -> Pin<Box<dyn Future<Output = Result<Vec<PushRegistration>, StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Registrations whose retry instant has arrived, in stable order.
Parked rows are excluded — see park.
Sourcefn advance<'life0, 'life1, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
next_seq: Seq,
) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn advance<'life0, 'life1, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
next_seq: Seq,
) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Acknowledge every record before next_seq.
Sourcefn retry<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
next_attempt_at: u64,
error: &'life2 str,
) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
fn retry<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
next_attempt_at: u64,
error: &'life2 str,
) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
Record a failed attempt without advancing the cursor.
Clears park: an attempt was made, so the registration is
live again whatever it was before.
Sourcefn park<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
error: &'life2 str,
) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
fn park<'life0, 'life1, 'life2, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
error: &'life2 str,
) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
'life2: 'async_trait,
Stop delivering to this registration, and keep its cursor.
What a worker does when a receiver has answered permanently, or has
stayed silent past the retry ceiling. Deleting the row instead would be
the one place this design loses an event: the cursor is the only record
of how far a receiver got, and once it is gone the undelivered tail of
that run is unrecoverable without a scan nobody schedules. A parked row
is not returned by due or due_in, so it
costs no sweep; it is listed by parked and re-armed by
unpark, which is the difference between a backlog an
operator can act on and a warning line in yesterday’s logs.
§Errors
If the store cannot be reached.
Sourcefn parked<'life0, 'async_trait>(
&'life0 self,
limit: usize,
) -> Pin<Box<dyn Future<Output = Result<Vec<PushRegistration>, StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn parked<'life0, 'async_trait>(
&'life0 self,
limit: usize,
) -> Pin<Box<dyn Future<Output = Result<Vec<PushRegistration>, StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
Sourcefn unpark<'life0, 'life1, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
at: u64,
) -> Pin<Box<dyn Future<Output = Result<bool, StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
fn unpark<'life0, 'life1, 'async_trait>(
&'life0 self,
task: RunId,
id: &'life1 str,
at: u64,
) -> Pin<Box<dyn Future<Output = Result<bool, StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
'life1: 'async_trait,
Re-arm a parked registration: due at at, with its attempt count reset.
The cursor is untouched, so delivery resumes at the first record the receiver never acknowledged rather than at the head of the run. Returns whether a parked registration was found — an operator re-arming one that is already live, or that never existed, must be told so rather than left waiting for a sweep that has nothing to do.
§Errors
If the store cannot be reached.
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.
Sourcefn due_in<'life0, 'async_trait>(
&'life0 self,
at: u64,
limit: usize,
namespace: PushNamespace,
) -> Pin<Box<dyn Future<Output = Result<DueBatch, StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
fn due_in<'life0, 'async_trait>(
&'life0 self,
at: u64,
limit: usize,
namespace: PushNamespace,
) -> Pin<Box<dyn Future<Output = Result<DueBatch, StoreError>> + Send + 'async_trait>>where
Self: 'async_trait,
'life0: 'async_trait,
The due registrations of one namespace, however deep in the stable order they sit — plus a count of the other namespace’s due rows.
The filter belongs in the query because it cannot live after it: a
worker that reads a bounded window and then drops the rows it does not
own is starved the moment the other namespace fills the window, and its
own rows beyond it are never read at all. The id prefix is already the
discriminator — see PushNamespace.
The default implementation pages over due with a growing
window until it holds limit rows of the namespace or the store is
exhausted. Correct against any backend, and linear in the other
namespace’s backlog — a backend with an index should override it with an
in-query filter on the id prefix.
§Errors
If the store cannot be reached.
Dyn Compatibility§
This trait is dyn compatible.
In older versions of Rust, dyn compatibility was called "object safety".