//! `/stream` events staged by a transaction. use std::sync::Arc; use bytes::Bytes; use fjall::OwnedWriteBatch; use miette::Result; use crate::db::keys; use crate::db::keyspaces::StreamDb; use crate::db::sequencer::{Assign, publish_in_order}; use crate::types::{ AccountEvt, BroadcastEvent, EventType, IdentityEvt, LiveRecordEvent, MarshallableEvt, StoredEvent, }; /// a staged `/stream` event, known by its place in the outbox until the /// transaction commits and gives it an id. #[derive(Clone, Copy)] pub(crate) struct StreamEventRef(usize); impl StreamEventRef { pub(super) fn index(self) -> usize { self.0 } } #[derive(Default)] pub(crate) struct StreamOutbox { events: Vec, } enum Pending { Record { row: Vec, live: Option<(StoredEvent<'static>, Option)>, }, Identity(IdentityEvt), Account(AccountEvt<'static>), } impl StreamOutbox { pub(super) fn is_empty(&self) -> bool { self.events.is_empty() } /// a persisted record event. `live` is the copy to broadcast, with the /// record body when live subscribers may use it instead of a db read; /// without it, subscribers read the row once a persisted marker covers it. pub(crate) fn push_record( &mut self, row: Vec, live: Option<(StoredEvent<'static>, Option)>, ) -> StreamEventRef { self.events.push(Pending::Record { row, live }); StreamEventRef(self.events.len() - 1) } /// an identity event. broadcast only, never persisted or replayed. pub(crate) fn push_identity(&mut self, identity: IdentityEvt) { self.events.push(Pending::Identity(identity)); } /// an account event. broadcast only, never persisted or replayed. pub(crate) fn push_account(&mut self, account: AccountEvt<'static>) { self.events.push(Pending::Account(account)); } pub(super) fn stage<'a>( self, db: &'a StreamDb, assign: &mut Assign<'a>, batch: &mut OwnedWriteBatch, ) -> Result { let events = self .events .into_iter() .map(|event| { let id = db.ids.take(assign)?; Ok((id, event.stage(id, db, batch))) }) .collect::>()?; Ok(StagedStream { events }) } } impl Pending { /// write this event's row, if it has one, and build its broadcast. fn stage(self, id: u64, db: &StreamDb, batch: &mut OwnedWriteBatch) -> Option { let (kind, identity, account) = match self { Self::Record { row, live } => { batch.insert(&db.events, keys::event_key(id), row); return live.map(|(stored, inline_block)| { BroadcastEvent::LiveRecord(Arc::new(LiveRecordEvent::new( id, stored, inline_block, ))) }); } Self::Identity(identity) => (EventType::Identity, Some(identity), None), Self::Account(account) => (EventType::Account, None, Some(account)), }; Some(BroadcastEvent::Ephemeral(Box::new(MarshallableEvt { id, kind, record: None, identity, account, }))) } } /// staged events with their ids, in outbox order. pub(super) struct StagedStream { events: Vec<(u64, Option)>, } impl StagedStream { #[cfg(feature = "jetstream")] pub(super) fn id(&self, event: StreamEventRef) -> u64 { self.events[event.0].0 } /// broadcast the staged events, returning the first one's id. the rest /// follow it without gaps. pub(super) fn publish(self, db: &StreamDb) -> Option { let first = self.events.first().map(|(id, _)| *id); publish_in_order(&db.event_tx, self.events, BroadcastEvent::Persisted); first } }