diff --git a/src/backfill/mod.rs b/src/backfill/mod.rs index d98197a..79022d7 100644 --- a/src/backfill/mod.rs +++ b/src/backfill/mod.rs @@ -711,7 +711,7 @@ async fn process_did( event_id, live: false, }; - crate::jetstream::stage_event(&mut batch, &app_state.db, jetstream)?; + crate::jetstream::stage_event(&mut batch, &app_state.db, jetstream, None)?; } } @@ -765,7 +765,7 @@ async fn process_did( event_id, live: false, }; - crate::jetstream::stage_event(&mut batch, &app_state.db, jetstream)?; + crate::jetstream::stage_event(&mut batch, &app_state.db, jetstream, None)?; } } diff --git a/src/control/jetstream.rs b/src/control/jetstream.rs index 01ec04d..720fbcc 100644 --- a/src/control/jetstream.rs +++ b/src/control/jetstream.rs @@ -248,7 +248,9 @@ mod tests { let wanted = opts.wanted_collections.unwrap(); assert!(wanted.prefixes.iter().any(|prefix| prefix == "app.bsky.")); - assert!(JetstreamSubscriberOptions::parse(&["app.bsky.feed.po*".into()], &[], 0, &[]).is_err()); + assert!( + JetstreamSubscriberOptions::parse(&["app.bsky.feed.po*".into()], &[], 0, &[]).is_err() + ); } #[test] diff --git a/src/control/stream.rs b/src/control/stream.rs index 7af7c10..8bb4d0a 100644 --- a/src/control/stream.rs +++ b/src/control/stream.rs @@ -893,6 +893,7 @@ fn read_jetstream_replay_chunk( id, time_us: time_us as i64, event: event.into_static(), + ephemeral: None, }); } @@ -905,17 +906,17 @@ fn read_jetstream_replay_chunk( #[cfg(feature = "jetstream")] #[derive(serde::Serialize)] -struct JetstreamEvent<'a> { - did: &'a str, - time_us: i64, +pub(crate) struct JetstreamEvent<'a> { + pub(crate) did: &'a str, + pub(crate) time_us: i64, #[serde(flatten)] - payload: JetstreamPayload<'a>, + pub(crate) payload: JetstreamPayload<'a>, } #[cfg(feature = "jetstream")] #[derive(serde::Serialize)] #[serde(tag = "kind", rename_all = "lowercase")] -enum JetstreamPayload<'a> { +pub(crate) enum JetstreamPayload<'a> { Commit { commit: JetstreamCommit<'a>, }, @@ -931,37 +932,37 @@ enum JetstreamPayload<'a> { #[cfg(feature = "jetstream")] #[derive(serde::Serialize)] -struct JetstreamCommit<'a> { - rev: &'a str, - operation: &'a str, - collection: &'a str, - rkey: &'a str, +pub(crate) struct JetstreamCommit<'a> { + pub(crate) rev: &'a str, + pub(crate) operation: &'a str, + pub(crate) collection: &'a str, + pub(crate) rkey: &'a str, #[serde(skip_serializing_if = "Option::is_none")] - record: Option<&'a serde_json::Value>, + pub(crate) record: Option<&'a serde_json::Value>, #[serde(skip_serializing_if = "Option::is_none")] - cid: Option, - live: bool, + pub(crate) cid: Option, + pub(crate) live: bool, } #[cfg(feature = "jetstream")] #[derive(serde::Serialize)] -struct JetstreamIdentity<'a> { - did: String, - seq: i64, - time: &'a crate::ingest::stream::Datetime, +pub(crate) struct JetstreamIdentity<'a> { + pub(crate) did: String, + pub(crate) seq: i64, + pub(crate) time: &'a crate::ingest::stream::Datetime, #[serde(skip_serializing_if = "Option::is_none")] - handle: Option, + pub(crate) handle: Option, } #[cfg(feature = "jetstream")] #[derive(serde::Serialize)] -struct JetstreamAccount<'a> { - active: bool, - did: String, - seq: i64, - time: &'a crate::ingest::stream::Datetime, +pub(crate) struct JetstreamAccount<'a> { + pub(crate) active: bool, + pub(crate) did: String, + pub(crate) seq: i64, + pub(crate) time: &'a crate::ingest::stream::Datetime, #[serde(skip_serializing_if = "Option::is_none")] - status: Option, + pub(crate) status: Option, } #[cfg(feature = "jetstream")] @@ -974,6 +975,11 @@ fn jetstream_event_to_bytes( return None; } + // live tailing: use pre-serialized ephemeral bytes to avoid db reads. + if let Some(bytes) = event.ephemeral { + return Some(bytes); + } + match &event.event { #[cfg(feature = "indexer_stream")] StoredJetstreamEvent::Commit { event_id, live, .. } => { diff --git a/src/ingest/indexer.rs b/src/ingest/indexer.rs index c060b2e..b0e0d8d 100644 --- a/src/ingest/indexer.rs +++ b/src/ingest/indexer.rs @@ -559,6 +559,7 @@ impl FirehoseWorker { seq: identity.seq, time: identity.time.clone(), }, + None, )?); if changed { @@ -600,6 +601,7 @@ impl FirehoseWorker { seq: account.seq, time: account.time.clone(), }, + None, )?); if is_inactive { diff --git a/src/ingest/relay.rs b/src/ingest/relay.rs index cce1c29..9a6e685 100644 --- a/src/ingest/relay.rs +++ b/src/ingest/relay.rs @@ -48,7 +48,10 @@ struct WorkerContext<'a> { #[cfg(feature = "relay")] pending_broadcasts: Vec, #[cfg(all(feature = "relay", feature = "jetstream"))] - pending_jetstream_events: Vec>, + pending_jetstream_events: Vec<( + crate::types::StoredJetstreamEvent<'static>, + Option, + )>, #[cfg(feature = "indexer")] pending_hook_messages: Vec, #[cfg(feature = "indexer")] @@ -226,8 +229,9 @@ impl RelayWorker { let res = { let _lock = ctx.state.db.jetstream_lock.lock(); let mut stage_res = Ok(()); - for event in ctx.pending_jetstream_events.drain(..) { - match crate::jetstream::stage_event(&mut batch, &ctx.state.db, event) { + for (event, ephemeral) in ctx.pending_jetstream_events.drain(..) { + match crate::jetstream::stage_event(&mut batch, &ctx.state.db, event, ephemeral) + { Ok(broadcast) => jetstream_broadcasts.push(broadcast), Err(e) => { stage_res = Err(e); @@ -383,15 +387,32 @@ impl RelayWorker { encode_frame("#commit", &commit) })?; #[cfg(feature = "jetstream")] - for (op_index, collection) in jetstream_ops { - // todo: build live jetstream frames from this decoded commit and use stored frames only for replay. - ctx.pending_jetstream_events - .push(StoredJetstreamEvent::RelayCommit { - did: jetstream_did.clone(), - collection, - relay_seq, - op_index, + { + let parsed_car = tokio::runtime::Handle::current() + .block_on(jacquard_repo::car::reader::parse_car_bytes( + commit.blocks.as_ref(), + )) + .ok(); + for (op_index, collection) in jetstream_ops { + let ephemeral = parsed_car.as_ref().and_then(|car| { + crate::jetstream::build_ephemeral_from_relay( + &commit, + op_index, + collection.as_str(), + true, + car, + ) }); + ctx.pending_jetstream_events.push(( + StoredJetstreamEvent::RelayCommit { + did: jetstream_did.clone(), + collection, + relay_seq: _relay_seq, + op_index, + }, + ephemeral, + )); + } } } @@ -530,11 +551,13 @@ impl RelayWorker { encode_frame("#identity", &identity) })?; #[cfg(feature = "jetstream")] - ctx.pending_jetstream_events - .push(StoredJetstreamEvent::RelayIdentity { + ctx.pending_jetstream_events.push(( + StoredJetstreamEvent::RelayIdentity { did: TrimmedDid::from(&identity.did).into_static(), - relay_seq, - }); + relay_seq: _relay_seq, + }, + None, + )); } ctx.batch.insert( @@ -628,11 +651,13 @@ impl RelayWorker { encode_frame("#account", &account) })?; #[cfg(feature = "jetstream")] - ctx.pending_jetstream_events - .push(StoredJetstreamEvent::RelayAccount { + ctx.pending_jetstream_events.push(( + StoredJetstreamEvent::RelayAccount { did: TrimmedDid::from(&account.did).into_static(), - relay_seq, - }); + relay_seq: _relay_seq, + }, + None, + )); } repo_state.touch(); diff --git a/src/jetstream.rs b/src/jetstream.rs index cda5678..2b0043f 100644 --- a/src/jetstream.rs +++ b/src/jetstream.rs @@ -1,18 +1,52 @@ use std::sync::atomic::Ordering; +use bytes::Bytes; use fjall::OwnedWriteBatch; use miette::{IntoDiagnostic, Result}; use crate::db::{Db, keys}; use crate::types::{JetstreamBroadcast, StoredJetstreamEvent}; +/// pre-built commit data for live jetstream tailing. +/// the stream thread serializes this into json with the assigned time_us. +#[derive(Clone)] +pub(crate) struct JetstreamEphemeral { + pub did: String, + pub rev: String, + pub operation: String, + pub collection: String, + pub rkey: String, + pub record: Option, + pub cid: Option, + pub live: bool, +} + pub(crate) fn stage_event( batch: &mut OwnedWriteBatch, db: &Db, event: StoredJetstreamEvent<'_>, + ephemeral: Option, ) -> Result { let id = db.next_jetstream_id.fetch_add(1, Ordering::SeqCst); let time_us = next_time_us(db); + let ephemeral = ephemeral.and_then(|data| { + let json_event = crate::control::stream::JetstreamEvent { + did: &data.did, + time_us, + payload: crate::control::stream::JetstreamPayload::Commit { + commit: crate::control::stream::JetstreamCommit { + rev: &data.rev, + operation: &data.operation, + collection: &data.collection, + rkey: &data.rkey, + record: data.record.as_ref(), + cid: data.cid, + live: data.live, + }, + }, + }; + serde_json::to_vec(&json_event).ok().map(Bytes::from) + }); let event = event.into_static(); let bytes = rmp_serde::to_vec(&event).into_diagnostic()?; batch.insert( @@ -20,7 +54,12 @@ pub(crate) fn stage_event( keys::jetstream_event_key(time_us as u64, id), bytes, ); - Ok(JetstreamBroadcast { id, time_us, event }) + Ok(JetstreamBroadcast { + id, + time_us, + event, + ephemeral, + }) } fn next_time_us(db: &Db) -> i64 { @@ -37,3 +76,89 @@ fn next_time_us(db: &Db) -> i64 { } } } + +#[cfg(all(feature = "indexer_stream", feature = "jetstream"))] +pub(crate) fn build_ephemeral_from_stored( + did: &str, + rev: &str, + operation: &str, + collection: &str, + rkey: &str, + data: &crate::types::StoredData, + inline_block: Option<&Bytes>, + live: bool, +) -> Option { + use crate::types::StoredData; + use jacquard_common::types::cid::{ATP_CID_HASH, IpldCid}; + use jacquard_repo::DAG_CBOR_CID_CODEC; + use sha2::{Digest, Sha256}; + + let (cid, record) = match data { + StoredData::Ptr(cid) => { + if let Some(bytes) = inline_block { + match serde_ipld_dagcbor::from_slice::(bytes) { + Ok(val) => (Some(cid.to_string()), Some(serde_json::to_value(val).ok()?)), + Err(_) => return None, + } + } else { + return None; + } + } + StoredData::Block(block) => { + let digest = Sha256::digest(block); + let hash = cid::multihash::Multihash::wrap(ATP_CID_HASH, &digest).ok()?; + let cid = IpldCid::new_v1(DAG_CBOR_CID_CODEC, hash); + match serde_ipld_dagcbor::from_slice::(block) { + Ok(val) => (Some(cid.to_string()), Some(serde_json::to_value(val).ok()?)), + Err(_) => return None, + } + } + StoredData::Nothing => (None, None), + }; + + Some(JetstreamEphemeral { + did: did.to_string(), + rev: rev.to_string(), + operation: operation.to_string(), + collection: collection.to_string(), + rkey: rkey.to_string(), + record, + cid, + live, + }) +} + +#[cfg(all(feature = "relay", feature = "jetstream"))] +pub(crate) fn build_ephemeral_from_relay( + commit: &crate::ingest::stream::Commit, + op_index: u32, + collection: &str, + live: bool, + parsed_car: &jacquard_repo::car::reader::ParsedCar, +) -> Option { + let op = commit.ops.get(op_index as usize)?; + let (_, rkey) = op.path.split_once('/')?; + let action = op.action.as_str(); + + let (record, cid) = if matches!(action, "create" | "update") { + let cid = op.cid.as_ref()?; + let cid_ipld = cid.to_ipld().ok()?; + let block = parsed_car.blocks.get(&cid_ipld)?; + let val = serde_ipld_dagcbor::from_slice::(block).ok()?; + let record = serde_json::to_value(val).ok()?; + (Some(record), Some(cid.to_string())) + } else { + (None, None) + }; + + Some(JetstreamEphemeral { + did: commit.repo.as_str().to_string(), + rev: commit.rev.as_str().to_string(), + operation: action.to_string(), + collection: collection.to_string(), + rkey: rkey.to_string(), + record, + cid, + live, + }) +} diff --git a/src/ops.rs b/src/ops.rs index b6e76ca..de11ca4 100644 --- a/src/ops.rs +++ b/src/ops.rs @@ -386,7 +386,19 @@ pub fn apply_commit<'s>( event_id, live: true, }; - jetstream_events.push(crate::jetstream::stage_event(batch, db, jetstream)?); + let ephemeral = crate::jetstream::build_ephemeral_from_stored( + did_trimmed.to_did().as_str(), + rev.to_tid().as_str(), + action.as_str(), + collection.as_str(), + rkey.to_smolstr().as_str(), + &data, + inline_block.as_ref(), + true, + ); + jetstream_events.push(crate::jetstream::stage_event( + batch, db, jetstream, ephemeral, + )?); } if should_broadcast_live { diff --git a/src/types.rs b/src/types.rs index 1e63d70..b2e045d 100644 --- a/src/types.rs +++ b/src/types.rs @@ -502,6 +502,9 @@ pub(crate) struct JetstreamBroadcast { pub id: u64, pub time_us: i64, pub event: StoredJetstreamEvent<'static>, + /// pre-serialized json bytes for live tailing. + /// replay events read from the database instead. + pub ephemeral: Option, } #[cfg(feature = "jetstream")]