//! relay `subscribeRepos` frames staged by a transaction. use bytes::Bytes; use fjall::OwnedWriteBatch; use miette::Result; use tracing::error; use crate::db::keys; use crate::db::keyspaces::RelayDb; use crate::db::sequencer::{Assign, publish_in_order}; use crate::ingest::stream::{Account, Commit, Identity, Sync, encode_frame}; use crate::types::RelayBroadcast; /// a staged relay frame, known by its place in the outbox until the /// transaction commits and gives it a seq. #[derive(Clone, Copy)] pub(crate) struct RelayFrameRef(#[cfg_attr(not(feature = "jetstream"), allow(dead_code))] usize); #[derive(Default)] pub(crate) struct RelayOutbox { frames: Vec, } enum Pending { Commit(Box>), Sync(Sync<'static>), Identity(Identity), Account(Account<'static>), /// a frame whose encoding fails #[cfg(test)] Unencodable, } impl RelayOutbox { pub(super) fn is_empty(&self) -> bool { self.frames.is_empty() } fn push(&mut self, frame: Pending) -> RelayFrameRef { self.frames.push(frame); RelayFrameRef(self.frames.len() - 1) } pub(crate) fn push_commit(&mut self, commit: Commit<'static>) -> RelayFrameRef { self.push(Pending::Commit(Box::new(commit))) } pub(crate) fn push_sync(&mut self, sync: Sync<'static>) -> RelayFrameRef { self.push(Pending::Sync(sync)) } pub(crate) fn push_identity(&mut self, identity: Identity) -> RelayFrameRef { self.push(Pending::Identity(identity)) } pub(crate) fn push_account(&mut self, account: Account<'static>) -> RelayFrameRef { self.push(Pending::Account(account)) } #[cfg(test)] pub(crate) fn push_unencodable(&mut self) -> RelayFrameRef { self.push(Pending::Unencodable) } pub(super) fn stage<'a>( self, db: &'a RelayDb, assign: &mut Assign<'a>, batch: &mut OwnedWriteBatch, ) -> Result { let frames = self .frames .into_iter() .map(|frame| { let seq = db.seqs.take(assign)?; // encoding is deterministic, so failing the commit would lose // the rest of its writes for good. only this frame is lost. let encoded = frame.encode(seq).inspect_err( |e| error!(err = %e, seq, "dropping a relay frame that failed to encode"), ); let Ok(bytes) = encoded else { return Ok((seq, None)); }; batch.insert(&db.events, keys::relay_event_key(seq), bytes.as_ref()); Ok((seq, Some(RelayBroadcast::Ephemeral(seq, bytes)))) }) .collect::>()?; Ok(StagedRelay { frames }) } } impl Pending { /// the frame to persist and broadcast at `seq`. fn encode(self, seq: u64) -> Result { let at = seq as i64; match self { Self::Commit(mut commit) => { commit.seq = at; encode_frame("#commit", &commit) } Self::Sync(mut sync) => { sync.seq = at; encode_frame("#sync", &sync) } Self::Identity(mut identity) => { identity.seq = at; encode_frame("#identity", &identity) } Self::Account(mut account) => { account.seq = at; encode_frame("#account", &account) } #[cfg(test)] Self::Unencodable => Err(miette::miette!("unencodable test frame")), } } } /// staged frames with their seqs, in outbox order. pub(super) struct StagedRelay { frames: Vec<(u64, Option)>, } impl StagedRelay { #[cfg(feature = "jetstream")] pub(super) fn id(&self, frame: RelayFrameRef) -> u64 { self.frames[frame.0].0 } pub(super) fn publish(self, db: &RelayDb) { publish_in_order(&db.broadcast_tx, self.frames, RelayBroadcast::Persisted); } }