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.
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889//! `/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<JetstreamEphemeral>)>,}
impl JetstreamOutbox { pub(super) fn is_empty(&self) -> bool { self.events.is_empty() }
pub(crate) fn push(&mut self, event: PendingJetstreamEvent, json: Option<JetstreamEphemeral>) { self.events.push((event, json)); }
pub(super) fn stage<'a>( self, db: &'a JetstreamDb, assign: &mut Assign<'a>, batch: &mut OwnedWriteBatch, upstream: &StagedUpstream, ) -> Result<StagedJetstream> { 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::<Result<_>>()?; Ok(StagedJetstream { events }) }}
pub(super) struct StagedJetstream { events: Vec<JetstreamLive>,}
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); }}