Skip to main content

BatchStore

Trait BatchStore 

Source
pub trait BatchStore:
    Send
    + Sync
    + Debug {
    // Required methods
    fn tenant(&self) -> &str;
    fn open<'life0, 'life1, 'async_trait>(
        &'life0 self,
        id: BatchId,
        plan_digest: &'life1 str,
    ) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait,
             'life1: 'async_trait;
    fn plan_digest<'life0, 'async_trait>(
        &'life0 self,
        id: BatchId,
    ) -> Pin<Box<dyn Future<Output = Result<Option<String>, StoreError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait;
    fn mark_exhausted<'life0, 'async_trait>(
        &'life0 self,
        id: BatchId,
    ) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait;
    fn is_exhausted<'life0, 'async_trait>(
        &'life0 self,
        id: BatchId,
    ) -> Pin<Box<dyn Future<Output = Result<bool, StoreError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait;
    fn reserve<'life0, 'life1, 'async_trait>(
        &'life0 self,
        batch: BatchId,
        key: &'life1 str,
        run: RunId,
    ) -> Pin<Box<dyn Future<Output = Result<ItemRecord, StoreError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait,
             'life1: 'async_trait;
    fn record<'life0, 'life1, 'life2, 'async_trait>(
        &'life0 self,
        batch: BatchId,
        key: &'life1 str,
        outcome: &'life2 ItemOutcome,
        spend: Spend,
    ) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait,
             'life1: 'async_trait,
             'life2: 'async_trait;
    fn cursor<'life0, 'async_trait>(
        &'life0 self,
        batch: BatchId,
    ) -> Pin<Box<dyn Future<Output = Result<Option<String>, StoreError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait;
    fn census<'life0, 'async_trait>(
        &'life0 self,
        batch: BatchId,
    ) -> Pin<Box<dyn Future<Output = Result<BatchCensus, StoreError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait;
    fn items<'life0, 'async_trait>(
        &'life0 self,
        batch: BatchId,
        limit: usize,
    ) -> Pin<Box<dyn Future<Output = Result<Vec<ItemRecord>, StoreError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait;
    fn items_needing_attention<'life0, 'async_trait>(
        &'life0 self,
        batch: BatchId,
        limit: usize,
    ) -> Pin<Box<dyn Future<Output = Result<Vec<ItemRecord>, StoreError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait;
}
Expand description

Durable batch state.

Separate from the journal because a batch is not a run: it has no hash chain of its own, and its integrity comes from the per-item runs it points at. What it needs is a cursor, a reservation, and a census — three queries, not a log.

Required Methods§

Source

fn tenant(&self) -> &str

Which tenant this handle’s batches and item reservations belong to.

Source

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

Register a batch. Idempotent on id, so a retried submission does not fork one act into two.

Idempotent on the id and held to the digest: one batch runs one frozen plan, and that sentence is this method’s to enforce, because the runner cannot — by the time it executes an item, the store’s row is the only witness to what the batch was opened with.

§Errors

StoreError::BatchPlanChanged when the batch exists under a different plan_digest. A resume offering an edited plan must be refused here, in the store, or items settle under a plan the batch’s record does not name.

Source

fn plan_digest<'life0, 'async_trait>( &'life0 self, id: BatchId, ) -> Pin<Box<dyn Future<Output = Result<Option<String>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

The plan digest this batch was opened with, or None for no such batch.

The existence question, answered from the batch’s own row. A census cannot answer it — a batch with no items yet and a batch that does not exist both count zero rows — and the difference matters to an operator: one is work not started, the other is a mistyped id.

Source

fn mark_exhausted<'life0, 'async_trait>( &'life0 self, id: BatchId, ) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Record that the source produced its last item.

Without this a batch cannot tell “every item I have stored is terminal” from “I am finished” — and those differ every time processing stops early. A batch halted after 10,000 of 100,000 items has no unfinished item anywhere in its store, so a census alone would report it complete with 90,000 meters unsettled.

Durable rather than in-memory because the distinction has to survive the process: a resumed batch must know whether it ever reached the end.

§Errors

StoreError::NotFound when no such batch exists, for the reason record refuses an unreserved item: this mark is the one bit that lets a census read as finished, and reporting it written when nothing was is the quietest possible way to lose it.

Source

fn is_exhausted<'life0, 'async_trait>( &'life0 self, id: BatchId, ) -> Pin<Box<dyn Future<Output = Result<bool, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Whether the source has been read to the end.

Source

fn reserve<'life0, 'life1, 'async_trait>( &'life0 self, batch: BatchId, key: &'life1 str, run: RunId, ) -> Pin<Box<dyn Future<Output = Result<ItemRecord, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

Claim an item and bind it to a run id, before the run starts.

Returns the existing record if the item was already reserved — which is what makes a crashed batch resumable: the second attempt gets the first attempt’s run id back and replays it rather than starting fresh.

Source

fn record<'life0, 'life1, 'life2, 'async_trait>( &'life0 self, batch: BatchId, key: &'life1 str, outcome: &'life2 ItemOutcome, spend: Spend, ) -> Pin<Box<dyn Future<Output = Result<(), StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait, 'life2: 'async_trait,

Record how an item ended, and what it consumed.

§Errors

StoreError::NotFound when the item was never reserved — a refusal, never a silent no-op: Ok over a write that matched nothing tells the caller recorded about an outcome that vanished, the same lie a release that freed nothing tells. The row count the write already produces is what catches it.

Source

fn cursor<'life0, 'async_trait>( &'life0 self, batch: BatchId, ) -> Pin<Box<dyn Future<Output = Result<Option<String>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

The highest key whose item reached a terminal outcome, with no unfinished item before it.

The contiguous prefix, not the maximum: an item suspended at key 400 must hold the cursor at 399 even if 401 through 500 have finished, or a resume would step over it and the batch would report complete with an item still waiting.

Source

fn census<'life0, 'async_trait>( &'life0 self, batch: BatchId, ) -> Pin<Box<dyn Future<Output = Result<BatchCensus, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait,

Counts by outcome, plus reserved-but-unfinished.

Source

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

Every item record, oldest key first. For operators and for tests; the driver uses cursor and census.

Source

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

The items that are not settled, oldest key first.

census answers 43 failed; the question an operator has is which 43, and over the size a batch exists for the only other route is paging items through a hundred thousand rows that are almost all successes — the finding indexed and reaching nobody.

Unsettled, not un-terminal: a failed item and a suspended one are both here, because finished and finished with are different questions. Each record’s ItemOutcome says whether the next move is a re-run, a raised ceiling, or a look.

Oldest-first is legitimate because entries leave: resolving an item drops it. An ascending page over a listing nothing empties would have a permanent head and an unreachable tail.

§Errors

Backend failures. An unknown batch is an empty listing, not an error: plan_digest is the existence question.

Dyn Compatibility§

This trait is dyn compatible.

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

Implementors§