very fast at protocol indexer with flexible filtering, xrpc queries, cursor-backed event stream, and more, built on fjall
rust fjall at-protocol atproto indexer
Something went wrong. Try again.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133134//! 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<Pending>,}
enum Pending { Commit(Box<Commit<'static>>), 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<StagedRelay> { 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::<Result<_>>()?; Ok(StagedRelay { frames }) }}
impl Pending { /// the frame to persist and broadcast at `seq`. fn encode(self, seq: u64) -> Result<Bytes> { 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<RelayBroadcast>)>,}
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); }}