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§
Sourcefn 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 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.
Sourcefn namespace(&self) -> PushNamespace
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§
Sourcefn terminal(&self, record: &Record) -> bool
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".