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.
123456789101112131415161718192021222324252627282930313233343536373839404142434445464748495051525354555657585960616263646566676869707172737475767778798081828384858687888990919293949596979899100101102103104105106107108109110111112113114115116117118119120121122123124125126127128129130131132133//! `/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<Pending>,}
enum Pending { Record { row: Vec<u8>, live: Option<(StoredEvent<'static>, Option<Bytes>)>, }, 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<u8>, live: Option<(StoredEvent<'static>, Option<Bytes>)>, ) -> 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<StagedStream> { let events = self .events .into_iter() .map(|event| { let id = db.ids.take(assign)?; Ok((id, event.stage(id, db, batch))) }) .collect::<Result<_>>()?; 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<BroadcastEvent> { 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<BroadcastEvent>)>,}
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<u64> { let first = self.events.first().map(|(id, _)| *id); publish_in_order(&db.event_tx, self.events, BroadcastEvent::Persisted); first }}