//! events a transaction stages, written and broadcast when it commits. //! //! stream rows are staged here without a position. [`Outbox::commit`] gives //! them one inside the [`Sequencer`](super::sequencer::Sequencer)'s critical //! section, which is the only place stream positions come from. use fjall::OwnedWriteBatch; use miette::{IntoDiagnostic, Result}; use super::Db; #[cfg(feature = "jetstream")] mod jetstream; #[cfg(feature = "relay")] mod relay; #[cfg(feature = "indexer_stream")] mod stream; #[cfg(feature = "jetstream")] pub(crate) use jetstream::JetstreamOutbox; #[cfg(feature = "relay")] pub(crate) use relay::{RelayFrameRef, RelayOutbox}; #[cfg(feature = "indexer_stream")] pub(crate) use stream::{StreamEventRef, StreamOutbox}; #[derive(Default)] pub(crate) struct Outbox { #[cfg(feature = "indexer_stream")] pub(crate) stream: StreamOutbox, #[cfg(feature = "relay")] pub(crate) relay: RelayOutbox, #[cfg(feature = "jetstream")] pub(crate) jetstream: JetstreamOutbox, } /// the positions a commit gave its events. #[derive(Default)] pub(crate) struct Assigned { #[cfg(feature = "indexer_stream")] first_stream_id: Option, } #[cfg(feature = "indexer_stream")] impl Assigned { pub(crate) fn stream_id(&self, event: StreamEventRef) -> Option { self.first_stream_id .map(|first| first + event.index() as u64) } } #[cfg(not(any(feature = "indexer_stream", feature = "relay")))] impl Outbox { pub(super) fn commit( self, _db: &Db, batch: OwnedWriteBatch, committed: impl FnOnce(), ) -> Result { batch.commit().into_diagnostic()?; committed(); Ok(Assigned::default()) } } #[cfg(any(feature = "indexer_stream", feature = "relay"))] impl Outbox { /// commit `batch` together with this outbox's rows, run `committed`, then /// broadcast them. pub(super) fn commit( self, db: &Db, batch: OwnedWriteBatch, committed: impl FnOnce(), ) -> Result { if self.is_empty() { batch.commit().into_diagnostic()?; committed(); return Ok(Assigned::default()); } db.sequencer.commit( batch, |assign, batch| self.stage(db, assign, batch), |staged| { committed(); staged.publish(db) }, ) } fn is_empty(&self) -> bool { #[cfg(feature = "indexer_stream")] if !self.stream.is_empty() { return false; } #[cfg(feature = "relay")] if !self.relay.is_empty() { return false; } #[cfg(feature = "jetstream")] if !self.jetstream.is_empty() { return false; } true } fn stage<'a>( self, db: &'a Db, assign: &mut super::sequencer::Assign<'a>, batch: &mut OwnedWriteBatch, ) -> Result { #[cfg(feature = "indexer_stream")] let stream = self.stream.stage(&db.stream, assign, batch)?; #[cfg(feature = "relay")] let relay = self.relay.stage(&db.relay, assign, batch)?; // jetstream rows refer to the upstream rows staged just above #[cfg(all(feature = "jetstream", feature = "indexer_stream"))] let jetstream = self .jetstream .stage(&db.jetstream, assign, batch, &stream)?; #[cfg(all(feature = "jetstream", feature = "relay"))] let jetstream = self.jetstream.stage(&db.jetstream, assign, batch, &relay)?; Ok(Staged { #[cfg(feature = "indexer_stream")] stream, #[cfg(feature = "relay")] relay, #[cfg(feature = "jetstream")] jetstream, }) } } #[cfg(any(feature = "indexer_stream", feature = "relay"))] struct Staged { #[cfg(feature = "indexer_stream")] stream: stream::StagedStream, #[cfg(feature = "relay")] relay: relay::StagedRelay, #[cfg(feature = "jetstream")] jetstream: jetstream::StagedJetstream, } #[cfg(any(feature = "indexer_stream", feature = "relay"))] impl Staged { fn publish(self, db: &Db) -> Assigned { #[cfg(feature = "relay")] self.relay.publish(&db.relay); #[cfg(feature = "jetstream")] self.jetstream.publish(&db.jetstream); Assigned { #[cfg(feature = "indexer_stream")] first_stream_id: self.stream.publish(&db.stream), } } } #[cfg(all(test, feature = "indexer_stream"))] mod tests { use std::sync::Arc; use jacquard_common::types::did::Did; use jacquard_common::types::string::Tid; use crate::config::Config; use crate::db::Txn; use crate::db::types::{DbAction, DbRkey, DbTid}; use crate::state::AppState; use crate::types::BroadcastEvent; const BACKFILLED: &str = "did:plc:aaaaaaaaaaaaaaaaaaaaaaaa"; const LIVE: &str = "did:plc:zzzzzzzzzzzzzzzzzzzzzzzz"; const COLLECTION: &str = "app.bsky.feed.post"; fn state() -> (tempfile::TempDir, Arc) { let tmp = tempfile::tempdir().unwrap(); let state = AppState::new(&Config { database_path: tmp.path().to_path_buf(), ..Default::default() }) .unwrap(); (tmp, Arc::new(state)) } fn put(records: &mut crate::db::txn::RecordTxn<'_, '_, '_>, n: usize) { let block = crate::car::CarBlock::from_body( serde_ipld_dagcbor::to_vec(&serde_json::json!({ "$type": COLLECTION, "n": n })) .unwrap(), ); let rkey = DbRkey::Str(smol_str::format_smolstr!("r{n}")); records .put_record(COLLECTION, &rkey, &block, DbAction::Create) .unwrap(); } /// a backfill batch is staged, a live commit lands, then the backfill /// commits: the order ingest produces when both run at once. fn interleave(state: &AppState, backfilled: usize) { let rev = DbTid::from(&Tid::now_0()); let did = Did::new_static(BACKFILLED).unwrap(); let mut backfill = Txn::new(&state.db); let mut records = backfill.backfill_records(state, &rev, &did).unwrap(); (0..backfilled).for_each(|n| put(&mut records, n)); records.finish().unwrap(); let did = Did::new_static(LIVE).unwrap(); let mut live = Txn::new(&state.db); let mut records = live.records(state, &rev, &did).unwrap(); put(&mut records, backfilled); records.finish().unwrap(); live.commit().unwrap(); backfill.commit().unwrap(); } fn event_did(state: &AppState, id: u64) -> String { let row = state .db .stream .event_range(crate::db::keys::event_key(id)..=crate::db::keys::event_key(id)) .next() .unwrap() .value() .unwrap(); let event: crate::types::StoredEvent<'_> = rmp_serde::from_slice(&row).unwrap(); event.did.to_did().as_str().to_owned() } #[test] fn stream_ids_follow_commit_order() { let (_tmp, state) = state(); let mut rx = state.db.stream.event_tx.subscribe(); interleave(&state, 3); // the live commit took the first id even though it staged last let first = state.db.stream.ids.head().unwrap() - 3; assert_eq!(event_did(&state, first), LIVE); for id in first + 1..=first + 3 { assert_eq!(event_did(&state, id), BACKFILLED); } let Ok(BroadcastEvent::LiveRecord(live)) = rx.try_recv() else { panic!("expected the live record first"); }; assert_eq!(live.id, first); let Ok(BroadcastEvent::Persisted(head)) = rx.try_recv() else { panic!("expected the backfill's persisted marker"); }; assert_eq!(head, first + 3); assert!(rx.try_recv().is_err()); } #[test] fn the_committed_hook_runs_before_the_broadcast() { let (_tmp, state) = state(); let rx = state.db.stream.event_tx.subscribe(); let mut outbox = super::Outbox::default(); outbox.stream.push_identity(crate::types::IdentityEvt { did: Did::new_static(LIVE).unwrap(), handle: None, }); let mut ran = false; outbox .commit(&state.db, state.db.inner.batch(), || { assert!(rx.is_empty(), "broadcast before the hook ran"); ran = true; }) .unwrap(); assert!(ran); assert!(!rx.is_empty()); } #[test] fn stream_ids_are_not_reused_after_a_restart() { let tmp = tempfile::tempdir().unwrap(); let config = Config { database_path: tmp.path().to_path_buf(), ..Default::default() }; let head = { let db = crate::db::Db::open(&config).unwrap(); // identity events take ids but leave no rows let mut txn = Txn::new(&db); txn.outbox.stream.push_identity(crate::types::IdentityEvt { did: Did::new_static(LIVE).unwrap(), handle: None, }); txn.commit().unwrap(); db.persist().unwrap(); db.stream.ids.head() }; let db = crate::db::Db::open(&config).unwrap(); let head = head.unwrap(); // past every id taken, and no further than one reserved block let next = db.stream.ids.next(); assert!(head < next && next <= head + 1 + crate::db::sequencer::SEQUENCE_BLOCK); } #[cfg(feature = "jetstream")] #[test] fn jetstream_keys_follow_commit_order() { use crate::types::StoredJetstreamEvent; let (_tmp, state) = state(); interleave(&state, 3); let rows: Vec<_> = state .db .jetstream .events .iter() .map(|guard| { let (key, value) = guard.into_inner().unwrap(); let position = crate::db::keys::parse_jetstream_event_key(&key).unwrap(); let event: StoredJetstreamEvent<'_> = rmp_serde::from_slice(&value).unwrap(); let StoredJetstreamEvent::Commit { event_id, live, .. } = event else { panic!("expected commit rows"); }; (position, event_id, live) }) .collect(); assert_eq!(rows.len(), 4); assert!(rows.windows(2).all(|pair| { let ((time_a, id_a), _, _) = pair[0]; let ((time_b, id_b), _, _) = pair[1]; time_a < time_b && id_a + 1 == id_b })); let (_, live_event, live) = rows[0]; assert!(live); assert_eq!(event_did(&state, live_event), LIVE); for &(_, event_id, live) in &rows[1..] { assert!(!live); assert_eq!(event_did(&state, event_id), BACKFILLED); } } } #[cfg(all(test, feature = "relay"))] mod relay_tests { use jacquard_common::types::string::Did; use crate::config::Config; use crate::db::{Db, Txn}; use crate::ingest::stream::{Datetime, Identity, SubscribeReposMessage, decode_frame}; use crate::types::RelayBroadcast; const FIRST: &str = "did:plc:aaaaaaaaaaaaaaaaaaaaaaaa"; const SECOND: &str = "did:plc:zzzzzzzzzzzzzzzzzzzzzzzz"; fn identity(did: &'static str) -> Identity { Identity { did: Did::new_static(did).unwrap(), handle: None, seq: 0, time: Datetime::try_from("2026-08-10T12:34:56Z".to_owned()).unwrap(), } } fn stage<'a>(db: &'a Db, did: &'static str) -> Txn<'a> { let mut txn = Txn::new(db); let _frame = txn.outbox.relay.push_identity(identity(did)); #[cfg(feature = "jetstream")] txn.outbox.jetstream.push( crate::types::StoredJetstreamEvent::RelayIdentity { did: crate::db::types::TrimmedDid::from(&Did::new_static(did).unwrap()) .into_static(), relay_seq: _frame, }, None, ); txn } #[test] fn relay_seqs_continue_after_their_rows_are_pruned() { let tmp = tempfile::tempdir().unwrap(); let config = Config { database_path: tmp.path().to_path_buf(), ..Default::default() }; let head = { let db = Db::open(&config).unwrap(); stage(&db, FIRST).commit().unwrap(); stage(&db, SECOND).commit().unwrap(); // ttl pruning can remove every row let mut batch = db.inner.batch(); for guard in db.relay.events.iter() { batch.remove(&db.relay.events, guard.key().unwrap()); } batch.commit().unwrap(); db.persist().unwrap(); db.relay.seqs.head() }; let db = Db::open(&config).unwrap(); let head = head.unwrap(); // past every seq taken, and no further than one reserved block let next = db.relay.seqs.next(); assert!(head < next && next <= head + 1 + crate::db::sequencer::SEQUENCE_BLOCK); } #[test] fn a_frame_that_fails_to_encode_is_dropped_alone() { let tmp = tempfile::tempdir().unwrap(); let db = Db::open(&Config { database_path: tmp.path().to_path_buf(), ..Default::default() }) .unwrap(); let mut rx = db.relay.broadcast_tx.subscribe(); let mut txn = stage(&db, FIRST); txn.outbox.relay.push_unencodable(); let _ = txn.outbox.relay.push_identity(identity(SECOND)); txn.batch.insert( &db.cursors, b"test_cursor".as_slice(), b"written".as_slice(), ); txn.commit().unwrap(); // the rest of the transaction lands, and the lost frame's seq stays a hole assert!(db.cursors.get(b"test_cursor").unwrap().is_some()); let seqs: Vec = db .relay .events .iter() .map(|guard| u64::from_be_bytes(guard.key().unwrap().as_ref().try_into().unwrap())) .collect(); let head = db.relay.seqs.head().unwrap(); assert_eq!(seqs, [head - 2, head]); let broadcast: Vec<_> = std::iter::from_fn(|| rx.try_recv().ok()).collect(); assert!(matches!( broadcast.as_slice(), [ RelayBroadcast::Ephemeral(first, _), RelayBroadcast::Persisted(hole), RelayBroadcast::Ephemeral(second, _), ] if *first == head - 2 && *hole == head - 1 && *second == head )); } #[test] fn relay_seqs_follow_commit_order() { let tmp = tempfile::tempdir().unwrap(); let db = Db::open(&Config { database_path: tmp.path().to_path_buf(), ..Default::default() }) .unwrap(); let mut rx = db.relay.broadcast_tx.subscribe(); let first = stage(&db, FIRST); stage(&db, SECOND).commit().unwrap(); first.commit().unwrap(); let frames: Vec<_> = db .relay .events .iter() .map(|guard| { let (key, frame) = guard.into_inner().unwrap(); let seq = u64::from_be_bytes(key.as_ref().try_into().unwrap()); let SubscribeReposMessage::Identity(identity) = decode_frame(&frame).unwrap() else { panic!("expected identity frames"); }; assert_eq!(identity.seq as u64, seq, "a frame carries its own seq"); (seq, identity.did.as_str().to_owned()) }) .collect(); let head = db.relay.seqs.head().unwrap(); assert_eq!( frames, [(head - 1, SECOND.to_owned()), (head, FIRST.to_owned())] ); let broadcast: Vec<_> = std::iter::from_fn(|| rx.try_recv().ok()) .map(|event| match event { RelayBroadcast::Ephemeral(seq, _) => seq, RelayBroadcast::Persisted(seq) => panic!("unexpected marker at {seq}"), }) .collect(); assert_eq!(broadcast, [head - 1, head]); #[cfg(feature = "jetstream")] for guard in db.jetstream.events.iter() { let (_, value) = guard.into_inner().unwrap(); let event: crate::types::StoredJetstreamEvent<'_> = rmp_serde::from_slice(&value).unwrap(); let crate::types::StoredJetstreamEvent::RelayIdentity { did, relay_seq } = event else { panic!("expected relay identity rows"); }; let frame = db .relay .events .get(crate::db::keys::relay_event_key(relay_seq)) .unwrap() .unwrap(); let SubscribeReposMessage::Identity(identity) = decode_frame(&frame).unwrap() else { panic!("jetstream row points at a non-identity frame"); }; assert_eq!(identity.did.as_str(), did.to_did().as_str()); } } }