Skip to main content

Projection

Trait Projection 

Source
pub trait Projection:
    Send
    + Sync
    + Debug {
    // Required methods
    fn messages<'life0, 'life1, 'async_trait>(
        &'life0 self,
        record: &'life1 Record,
    ) -> Pin<Box<dyn Future<Output = Result<Vec<PushMessage>, StoreError>> + Send + 'async_trait>>
       where Self: 'async_trait,
             'life0: 'async_trait,
             'life1: 'async_trait;
    fn namespace(&self) -> PushNamespace;

    // Provided method
    fn terminal(&self, record: &Record) -> bool { ... }
}
Expand description

What one journal record should be delivered as.

Returning several messages for one record is deliberate: A2A’s status and artifact events are two messages derived from one append, and a projection that could only answer with one would force the caller to invent a second cursor. An empty vector means nothing to send for this record, and the cursor still advances past it — which is how a projection filters.

Required Methods§

Source

fn messages<'life0, 'life1, 'async_trait>( &'life0 self, record: &'life1 Record, ) -> Pin<Box<dyn Future<Output = Result<Vec<PushMessage>, StoreError>> + Send + 'async_trait>>
where Self: 'async_trait, 'life0: 'async_trait, 'life1: 'async_trait,

The messages this record becomes, in delivery order.

Each carries its own media type and its own identity, because those are facts about the payload and not about the loop: two projections share this worker and speak different wires, and a receiver that cannot recognise a repeat has no defence against at-least-once delivery. See PushMessage.

§Errors

StoreError when the record cannot be projected — a materialisation that needs a blob store, say. The worker treats this as transient and this plane’s own fault: the cursor does not move, so the same record is tried again, and the retry ceiling still applies because a record that cannot be projected now will not project on the next tick either.

Source

fn namespace(&self) -> PushNamespace

Which id namespace this worker serves.

Two workers share one store: the A2A worker serves caller-registered webhooks, the outbox worker serves operator destinations. Without this each would deliver the other’s rows with its own projection — an operator’s CloudEvents message to a peer’s A2A webhook, and a StreamResponse to the deployment’s bus.

A declaration rather than a per-row predicate, because the split is the store’s discriminator: the worker hands it to PushStore::due_in so the query itself filters, instead of reading a bounded window and dropping the rows it does not own — which starves this worker’s rows the moment the other namespace fills the window.

Provided Methods§

Source

fn terminal(&self, record: &Record) -> bool

Whether this record ends the stream for its run.

The registration is deleted once its payloads are acknowledged, because a cursor sitting forever at the end of a sealed run’s journal is a row that is scanned on every tick and can never move.

The default is RecordKind::RunConcluded, which is the honest answer for anything keyed on a run: nothing is appended after a seal.

Dyn Compatibility§

This trait is dyn compatible.

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

Implementors§