Skip to main content

PushStore

Trait PushStore 

Source
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§

Source

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.

Source

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,

One configuration.

§Errors

If the store cannot be reached.

Source

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,

Every configuration for a task.

§Errors

If the store cannot be reached.

Source

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.

Source

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.

Source

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.

Source

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.

Source

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,

Parked registrations, in the store’s stable order.

§Errors

If the store cannot be reached.

Source

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.

Source

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,

Forget one. Idempotent.

§Errors

If the store cannot be reached.

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.

Source

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".

Implementors§