//! `/subscribe` events staged by a transaction. use fjall::OwnedWriteBatch; use miette::Result; use tracing::error; use crate::db::keyspaces::JetstreamDb; use crate::db::sequencer::{Assign, publish_in_order}; use crate::jetstream::JetstreamEphemeral; use crate::types::{JetstreamBroadcast, JetstreamLive, JetstreamPosition, StoredJetstreamEvent}; #[cfg(feature = "relay")] use super::relay::{RelayFrameRef as UpstreamRef, StagedRelay as StagedUpstream}; #[cfg(feature = "indexer_stream")] use super::stream::{StagedStream as StagedUpstream, StreamEventRef as UpstreamRef}; /// a jetstream event staged next to the upstream row its body resolves /// through, which is in the same outbox. pub(crate) type PendingJetstreamEvent = StoredJetstreamEvent<'static, UpstreamRef>; #[derive(Default)] pub(crate) struct JetstreamOutbox { /// each event with the parts of its live json, when a subscriber was /// listening as it was staged events: Vec<(PendingJetstreamEvent, Option)>, } impl JetstreamOutbox { pub(super) fn is_empty(&self) -> bool { self.events.is_empty() } pub(crate) fn push(&mut self, event: PendingJetstreamEvent, json: Option) { self.events.push((event, json)); } pub(super) fn stage<'a>( self, db: &'a JetstreamDb, assign: &mut Assign<'a>, batch: &mut OwnedWriteBatch, upstream: &StagedUpstream, ) -> Result { let events = self .events .into_iter() .map(|(event, json)| { let event = event.resolve(|upstream_ref| upstream.id(upstream_ref)); let position = JetstreamPosition { time_us: db.clock.tick(assign)?, id: db.ids.take(assign)?, }; // like a relay frame that fails to encode, only this event is lost let Ok(row) = rmp_serde::to_vec(&event).inspect_err(|e| { error!(err = %e, ?position, "dropping a jetstream event that failed to serialize") }) else { return Ok(None); }; batch.insert(&db.events, position.key(), row); let json = json.and_then(|json| json.to_json(position.time_us as i64)); Ok(Some(JetstreamLive { position, event, json, })) }) .filter_map(Result::transpose) .collect::>()?; Ok(StagedJetstream { events }) } } pub(super) struct StagedJetstream { events: Vec, } impl StagedJetstream { /// broadcast the live events. backfilled commits are only announced by a /// marker, for the subscribers that want them to read. pub(super) fn publish(self, db: &JetstreamDb) { let events = self.events.into_iter().map(|live| { let position = live.position; let broadcast = live.event.is_live().then(|| JetstreamBroadcast::Live(live)); (position, broadcast) }); publish_in_order(&db.tx, events, JetstreamBroadcast::Historical); } }