use std::sync::Arc; use std::sync::atomic::Ordering; use tokio::sync::mpsc; use tracing::{debug, error}; use crate::db::types::{DbRkey, DbTid, TrimmedDid}; use crate::db::{self, keys}; use crate::state::AppState; use crate::types::{BroadcastEvent, MarshallableEvt, RecordEvt, StoredData, StoredEvent}; use jacquard_common::types::cid::{ATP_CID_HASH, IpldCid}; use jacquard_common::types::nsid::Nsid; use jacquard_common::types::string::Rkey; use jacquard_common::{CowStr, IntoStatic, RawData}; use jacquard_repo::DAG_CBOR_CID_CODEC; use sha2::{Digest, Sha256}; use super::{ReplayChunk, StreamBroadcast, StreamOptions, run_ordered_stream, stream_seq_after}; use crate::control::{Event, StreamError}; pub(crate) fn event_stream_thread( state: Arc, tx: mpsc::Sender>, cursor: Option, opts: StreamOptions, ) { let db = &state.db; let event_rx = db.stream.event_tx.subscribe(); let current_id = match cursor { Some(c) => c.checked_sub(1), None => db .stream .next_event_id .load(Ordering::SeqCst) .checked_sub(1), }; let catch_up_target = cursor .and_then(|_| { db.stream .next_event_id .load(Ordering::SeqCst) .checked_sub(1) }) .filter(|target| stream_seq_after(*target, current_id)); let replay_state = state.clone(); run_ordered_stream( tx, event_rx, current_id, catch_up_target, opts, move |current_id, target, chunk_size| { read_event_replay_chunk(&replay_state, current_id, target, chunk_size) }, move |event| broadcast_to_event(&state, event), ); } fn read_event_replay_chunk( state: &AppState, current_id: Option, target: u64, chunk_size: usize, ) -> ReplayChunk { let start = current_id.map(|id| id.saturating_add(1)).unwrap_or(0); if start > target { return ReplayChunk { events: Vec::new(), last_seen_seq: current_id, exhausted: true, }; } let mut events = Vec::with_capacity(chunk_size); let mut last_seen_seq = current_id; let mut exhausted = false; let max_scanned = chunk_size.saturating_mul(4).max(chunk_size); let mut scanned = 0usize; let mut iter = state .db .stream .event_range(keys::event_key(start)..=keys::event_key(target)); while events.len() < chunk_size && scanned < max_scanned { let Some(item) = iter.next() else { exhausted = true; break; }; scanned += 1; let (k, v) = match item.into_inner() { Ok(kv) => kv, Err(e) => { error!(err = %e, "failed to read event from db"); exhausted = true; break; } }; let id = match k.as_ref().try_into().map(u64::from_be_bytes) { Ok(id) => id, Err(_) => { error!("failed to parse event id"); continue; } }; last_seen_seq = Some(id); let stored: StoredEvent = match rmp_serde::from_slice(&v) { Ok(e) => e, Err(e) => { error!(err = %e, "failed to deserialize stored event"); continue; } }; let Some(out_evt) = stored_to_event(state, id, stored, None) else { continue; }; events.push(out_evt); } ReplayChunk { events, last_seen_seq, exhausted, } } fn broadcast_to_event(state: &AppState, event: BroadcastEvent) -> Option { match event { BroadcastEvent::Persisted(_) => None, BroadcastEvent::LiveRecord(evt) => { let stored = evt.stored.clone(); stored_to_event(state, evt.id, stored, evt.inline_block.clone()) } BroadcastEvent::Ephemeral(evt) => Some(*evt), } } pub(crate) fn stored_to_event( state: &AppState, id: u64, stored: StoredEvent<'_>, inline_block: Option, ) -> Option { let StoredEvent { live, did, rev, collection, rkey, action, data, } = stored; let (cid, record) = match data { StoredData::Ptr(cid) => { // an inline live-tail body can outlive its committing transaction. // re-resolve its CID before exposing it so an operator redaction // that completed between commit and broadcast cannot leak the // already-queued bytes. let persisted = resolve_event_body(state, id, &did, collection.as_str(), &rkey, &rev, &cid); let bytes = inline_block .filter(|_| persisted.is_some()) .or_else(|| persisted.map(|body| bytes::Bytes::copy_from_slice(&body))); let record = match bytes { Some(bytes) => match serde_ipld_dagcbor::from_slice::(&bytes) { Ok(val) => serde_json::value::to_raw_value(&val).ok().map(Arc::from), Err(e) => { error!(err = %e, id, "cant parse record body"); None } }, None => { if state.history_ttl.is_some() { // bounded history: the body was most likely trimmed by // retention, which is expected, not a storage bug state.history_trim_misses.fetch_add(1, Ordering::Relaxed); debug!(id, "record body not found for event (bounded history)"); } else { debug!(id, "record body unavailable for event"); } None } }; (Some(cid), record) } StoredData::Block(block) => { let digest = Sha256::digest(&block); let hash = cid::multihash::Multihash::wrap(ATP_CID_HASH, &digest).expect("valid sha256 hash"); let cid = IpldCid::new_v1(DAG_CBOR_CID_CODEC, hash); let record = match serde_ipld_dagcbor::from_slice::(&block) { Ok(val) => serde_json::value::to_raw_value(&val).ok().map(Arc::from), Err(e) => { error!(err = %e, id, "cant parse inline record body"); None } }; (Some(cid), record) } StoredData::Nothing => (None, None), }; Some(MarshallableEvt { id, kind: crate::types::EventType::Record, record: Some(RecordEvt { live, did: did.to_did(), rev: rev.to_tid(), collection: match Nsid::new_cow(collection.clone().into_static()) { Ok(nsid) => nsid, Err(e) => { error!(err = %e, ?collection, "stored event has invalid collection NSID"); return None; } }, rkey: match Rkey::new_cow(CowStr::Owned(rkey.to_smolstr())) { Ok(rk) => rk, Err(e) => { error!(err = %e, ?rkey, "stored event has invalid rkey"); return None; } }, action: CowStr::Borrowed(action.as_str()), record, cid, }), identity: None, account: None, }) } /// resolve the body an event's record pointer refers to. /// /// current heads and the exact pre-v10 compatibility archive are point-read /// first; ordinary deaths after this event's rev are then scanned for the body /// whose cid the event claims. every candidate is content-verified by the db /// resolver before it can be returned. pub(crate) fn resolve_event_body( state: &AppState, id: u64, did: &TrimmedDid, collection: &str, rkey: &DbRkey, rev: &DbTid, cid: &IpldCid, ) -> Option { let key = keys::record_key_trimmed(did, collection, rkey); match state.db.resolve_event_record_body(&key, rev, cid) { Ok(body) => body, Err(e) => { error!(err = %e, id, "cant resolve event record body"); db::check_poisoned_report(&e); None } } } #[cfg(test)] mod tests { use super::*; use crate::config::Config; use jacquard_common::types::string::Tid; use tempfile::TempDir; const DID: &str = "did:plc:ewvi7nxzyoun6zhxrhs64oiz"; const COL: &str = "app.bsky.feed.post"; const RKEY: &str = "3kzbif5moe22m"; fn test_state() -> (TempDir, AppState) { let tmp = tempfile::tempdir().unwrap(); let mut config = Config::default(); config.database_path = tmp.path().to_path_buf(); let state = AppState::new(&config).unwrap(); (tmp, state) } fn rev(s: &str) -> DbTid { DbTid::from(&Tid::new(s).unwrap()) } fn resolve(state: &AppState, id: u64, event_rev: &DbTid, cid: &IpldCid) -> Option> { let did_full = jacquard_common::types::string::Did::new(DID).unwrap(); let did = TrimmedDid::from(&did_full); resolve_event_body(state, id, &did, COL, &DbRkey::new(RKEY), event_rev, cid) .map(|b| b.to_vec()) } fn record_key() -> Vec { let did = jacquard_common::types::string::Did::new(DID).unwrap(); keys::record_key(&did, COL, &DbRkey::new(RKEY)) } #[test] fn legacy_body_resolves_from_compatibility_archive() { let (_tmp, state) = test_state(); let cid = jacquard_repo::mst::util::compute_cid(b"legacy body").unwrap(); let mut batch = state.db.inner.batch(); batch.insert( &state.db.stream.event_bodies, keys::event_body_key(&record_key(), &cid), b"legacy body", ); batch.commit().unwrap(); assert_eq!( resolve(&state, 41, &rev("3kzbif5moe22m"), &cid), Some(b"legacy body".to_vec()) ); } #[test] fn falls_to_head_when_record_never_died() { let (_tmp, state) = test_state(); let mut batch = state.db.inner.batch(); state .db .indexer .stage_record(&mut batch, record_key(), b"head body".as_slice()); batch.commit().unwrap(); let cid = jacquard_repo::mst::util::compute_cid(b"head body").unwrap(); assert_eq!( resolve(&state, 42, &rev("3kzbif5moe22m"), &cid), Some(b"head body".to_vec()) ); } #[test] fn serves_the_body_killed_after_the_event() { let (_tmp, state) = test_state(); // the event wrote "old body" at rev m; the record died at rev n > m let mut batch = state.db.inner.batch(); state.db.indexer.stage_history( &mut batch, keys::history_key(&record_key(), &rev("3kzbif5mof33m")), b"old body".as_slice(), ); state .db .indexer .stage_record(&mut batch, record_key(), b"new body".as_slice()); batch.commit().unwrap(); let cid = jacquard_repo::mst::util::compute_cid(b"old body").unwrap(); assert_eq!( resolve(&state, 42, &rev("3kzbif5moe22m"), &cid), Some(b"old body".to_vec()) ); } #[test] fn deaths_before_the_event_are_not_served() { let (_tmp, state) = test_state(); // the only death predates the event: the event's body must be the head let mut batch = state.db.inner.batch(); state.db.indexer.stage_history( &mut batch, keys::history_key(&record_key(), &rev("3kzbif5mod22m")), b"ancient body".as_slice(), ); state .db .indexer .stage_record(&mut batch, record_key(), b"event body".as_slice()); batch.commit().unwrap(); let cid = jacquard_repo::mst::util::compute_cid(b"event body").unwrap(); assert_eq!( resolve(&state, 42, &rev("3kzbif5moe22m"), &cid), Some(b"event body".to_vec()) ); } #[test] fn delete_tombstone_resolves_to_no_body() { let (_tmp, state) = test_state(); // a delete of a never-seen record leaves an empty tombstone: the // record was already deleted at that point, so there is no body let mut batch = state.db.inner.batch(); state.db.indexer.stage_history( &mut batch, keys::history_key(&record_key(), &rev("3kzbif5mof33m")), &[], ); batch.commit().unwrap(); let cid = jacquard_repo::mst::util::compute_cid(b"anything").unwrap(); assert_eq!(resolve(&state, 42, &rev("3kzbif5moe22m"), &cid), None); } #[test] fn block_event_inflates_with_inline_body() { let (_tmp, state) = test_state(); let record_val = serde_json::json!({ "$type": "app.bsky.feed.post", "text": "hello" }); let record_cbor = serde_ipld_dagcbor::to_vec(&record_val).unwrap(); let expected_cid = jacquard_repo::mst::util::compute_cid(&record_cbor).unwrap(); let did_full = jacquard_common::types::string::Did::new(DID).unwrap(); let stored = StoredEvent { live: false, did: TrimmedDid::from(&did_full), rev: rev("3kzbif5moe22m"), collection: CowStr::Borrowed(COL), rkey: DbRkey::new(RKEY), action: crate::db::types::DbAction::Create, data: StoredData::Block(bytes::Bytes::from(record_cbor)), }; let serialized = rmp_serde::to_vec(&stored).unwrap(); let deserialized: StoredEvent = rmp_serde::from_slice(&serialized).unwrap(); let event = stored_to_event(&state, 100, deserialized, None).unwrap(); let rec = event.record.unwrap(); assert_eq!(rec.cid, Some(expected_cid)); assert!(rec.record.is_some()); } #[test] fn redacted_head_suppresses_queued_inline_body_but_keeps_cid_metadata() { let (_tmp, state) = test_state(); let body = serde_ipld_dagcbor::to_vec(&serde_json::json!({"text": "private"})).unwrap(); let cid = jacquard_repo::mst::util::compute_cid(&body).unwrap(); let mut batch = state.db.inner.batch(); state .db .indexer .stage_record(&mut batch, record_key(), cid.to_bytes()); batch.commit().unwrap(); let did_full = jacquard_common::types::string::Did::new(DID).unwrap(); let stored = StoredEvent { live: true, did: TrimmedDid::from(&did_full), rev: rev("3kzbif5moe22m"), collection: CowStr::Borrowed(COL), rkey: DbRkey::new(RKEY), action: crate::db::types::DbAction::Create, data: StoredData::Ptr(cid), }; let event = stored_to_event(&state, 42, stored, Some(bytes::Bytes::from(body))).unwrap(); let record = event.record.unwrap(); assert_eq!(record.cid, Some(cid)); assert!(record.record.is_none()); } #[test] fn history_body_cid_mismatch_returns_none() { let (_tmp, mut state) = test_state(); state.history_ttl = Some(std::time::Duration::from_secs(3600)); let body1 = serde_ipld_dagcbor::to_vec(&serde_json::json!({"text": "body 1"})).unwrap(); let body2 = serde_ipld_dagcbor::to_vec(&serde_json::json!({"text": "body 2"})).unwrap(); let cid1 = jacquard_repo::mst::util::compute_cid(&body1).unwrap(); let cid2 = jacquard_repo::mst::util::compute_cid(&body2).unwrap(); assert_ne!(cid1, cid2); let mut batch = state.db.inner.batch(); state.db.indexer.stage_history( &mut batch, keys::history_key(&record_key(), &rev("3kzbif5mof33m")), &body2, ); batch.commit().unwrap(); let resolved = resolve(&state, 42, &rev("3kzbif5moe22m"), &cid1); assert_eq!(resolved, None); let did_full = jacquard_common::types::string::Did::new(DID).unwrap(); let stored = StoredEvent { live: false, did: TrimmedDid::from(&did_full), rev: rev("3kzbif5moe22m"), collection: CowStr::Borrowed(COL), rkey: DbRkey::new(RKEY), action: crate::db::types::DbAction::Create, data: StoredData::Ptr(cid1), }; let initial_misses = state.history_trim_misses.load(Ordering::Relaxed); let event = stored_to_event(&state, 42, stored, None).unwrap(); let record = event.record.unwrap(); assert_eq!(record.cid, Some(cid1)); assert!(record.record.is_none()); assert_eq!( state.history_trim_misses.load(Ordering::Relaxed), initial_misses + 1 ); } #[test] fn head_body_cid_mismatch_returns_none() { let (_tmp, mut state) = test_state(); state.history_ttl = Some(std::time::Duration::from_secs(3600)); let body1 = serde_ipld_dagcbor::to_vec(&serde_json::json!({"text": "body 1"})).unwrap(); let body2 = serde_ipld_dagcbor::to_vec(&serde_json::json!({"text": "body 2"})).unwrap(); let cid1 = jacquard_repo::mst::util::compute_cid(&body1).unwrap(); let cid2 = jacquard_repo::mst::util::compute_cid(&body2).unwrap(); assert_ne!(cid1, cid2); let mut batch = state.db.inner.batch(); state .db .indexer .stage_record(&mut batch, record_key(), &body2); batch.commit().unwrap(); let resolved = resolve(&state, 42, &rev("3kzbif5moe22m"), &cid1); assert_eq!(resolved, None); let did_full = jacquard_common::types::string::Did::new(DID).unwrap(); let stored = StoredEvent { live: false, did: TrimmedDid::from(&did_full), rev: rev("3kzbif5moe22m"), collection: CowStr::Borrowed(COL), rkey: DbRkey::new(RKEY), action: crate::db::types::DbAction::Create, data: StoredData::Ptr(cid1), }; let initial_misses = state.history_trim_misses.load(Ordering::Relaxed); let event = stored_to_event(&state, 42, stored, None).unwrap(); let record = event.record.unwrap(); assert_eq!(record.cid, Some(cid1)); assert!(record.record.is_none()); assert_eq!( state.history_trim_misses.load(Ordering::Relaxed), initial_misses + 1 ); } }