diff --git a/src/api/debug.rs b/src/api/debug.rs index 365e058..b51b4f7 100644 --- a/src/api/debug.rs +++ b/src/api/debug.rs @@ -591,6 +591,7 @@ pub async fn handle_debug_seed_relay_jetstream_commits( record: None, cid: None, live: true, + redaction_key: None, }), ); } diff --git a/src/control/stream/interleaving.rs b/src/control/stream/interleaving.rs index 5e07052..0847bcf 100644 --- a/src/control/stream/interleaving.rs +++ b/src/control/stream/interleaving.rs @@ -279,6 +279,53 @@ mod jetstream { assert_eq!(seen.counts(), (0, 1)); } + #[test] + fn a_redaction_after_the_broadcast_keeps_the_body_out() { + let (_tmp, state, opts) = state(); + let (tx, mut rx) = mpsc::channel(2); + let filter = JetstreamFilter::new(JetstreamSubscriberOptions::default()); + let thread_state = state.clone(); + std::thread::spawn(move || { + super::super::jetstream_stream_thread(thread_state, tx, None, filter, opts) + }); + // live bodies are only kept inline while /stream has a subscriber too + let _stream = stream(&state, opts, None); + wait_subscribed(1, || state.db.jetstream.tx.receiver_count()); + wait_subscribed(1, || state.db.stream.event_tx.receiver_count()); + let taken = || wait_subscribed(1, || 1 - state.db.jetstream.tx.len().min(1)); + + // fill the output and leave the subscriber waiting to send the next, + // so the one after sits unrendered in its backlog + commit_live_records(&state, REPO_SIZE..REPO_SIZE + 2); + taken(); + commit_live_records(&state, REPO_SIZE + 2..REPO_SIZE + 3); + taken(); + let did = Did::new_static(LIVE).unwrap(); + let rkey = DbRkey::Str(smol_str::format_smolstr!("r{}", REPO_SIZE + 2)); + crate::db::redact_record_bodies( + &state.db, + &did, + COLLECTION, + &rkey, + crate::types::DeleteBodyTarget::Head, + ) + .unwrap(); + + let commits: Vec = (0..3) + .map(|_| { + let item = rx.blocking_recv().unwrap().unwrap(); + serde_json::from_slice::(&item).unwrap()["commit"].take() + }) + .collect(); + assert!( + commits[..2] + .iter() + .all(|commit| commit.get("record").is_some()) + ); + assert!(commits[2].get("record").is_none(), "{}", commits[2]); + assert!(commits[2].get("cid").is_some()); + } + #[test] fn replay_sends_backfills_by_default() { let (_tmp, state, opts) = state(); diff --git a/src/control/stream/jetstream.rs b/src/control/stream/jetstream.rs index 8eff521..f3fecd6 100644 --- a/src/control/stream/jetstream.rs +++ b/src/control/stream/jetstream.rs @@ -89,8 +89,13 @@ impl StreamSource for JetstreamSource { if !self.filter.wants(&live.event) { return None; } - // pre-serialized json avoids a db read - live.json.or_else(|| { + // pre-serialized json avoids a db read, unless an operator redaction + // that committed after the broadcast erased the body it embeds + let redacted = live + .redaction_key + .as_deref() + .is_some_and(|key| is_redacted(&self.state, key)); + live.json.filter(|_| !redacted).or_else(|| { stored_event_to_bytes(&self.state, live.position.time_us as i64, &live.event) }) } @@ -102,6 +107,24 @@ impl StreamSource for JetstreamSource { } } +/// a marker that can't be read counts as a redaction, so the event is +/// rendered from the db, which resolves the body again. +#[cfg(feature = "indexer_stream")] +fn is_redacted(state: &AppState, redaction_key: &[u8]) -> bool { + state + .db + .indexer + .is_redacted(redaction_key) + .inspect_err(|e| error!(err = %e, "failed to check a live event for redaction")) + .unwrap_or(true) +} + +/// relay mode stores no records to redact +#[cfg(feature = "relay")] +fn is_redacted(_state: &AppState, _redaction_key: &[u8]) -> bool { + false +} + fn read_jetstream_replay_chunk( state: &AppState, after: Option, diff --git a/src/db/keyspaces.rs b/src/db/keyspaces.rs index d63496f..c3083fb 100644 --- a/src/db/keyspaces.rs +++ b/src/db/keyspaces.rs @@ -98,6 +98,12 @@ impl IndexerDb { self.records.range(range) } + /// whether an operator erased the version a redaction marker key names. + #[cfg(feature = "jetstream")] + pub(crate) fn is_redacted(&self, redaction_key: &[u8]) -> fjall::Result { + self.redactions.contains_key(redaction_key) + } + // read-side only: stream inflation is today's sole history reader #[cfg(feature = "indexer_stream")] pub(crate) fn history_range(&self, range: R) -> fjall::Iter diff --git a/src/db/outbox/jetstream.rs b/src/db/outbox/jetstream.rs index 579c73e..adf6f8c 100644 --- a/src/db/outbox/jetstream.rs +++ b/src/db/outbox/jetstream.rs @@ -57,11 +57,14 @@ impl JetstreamOutbox { }) .ok()?; batch.insert(&db.events, position.key(), row); - let json = json.and_then(|json| json.to_json(position.time_us as i64)); + let (json, redaction_key) = json.map_or((None, None), |json| { + (json.to_json(position.time_us as i64), json.redaction_key) + }); Some(JetstreamLive { position, event, json, + redaction_key, }) }) .collect(); diff --git a/src/ingest/relay/sink/relay.rs b/src/ingest/relay/sink/relay.rs index c9d8c5a..5fd5b09 100644 --- a/src/ingest/relay/sink/relay.rs +++ b/src/ingest/relay/sink/relay.rs @@ -229,6 +229,7 @@ fn build_relay_commit_ephemeral( record, cid, live: true, + redaction_key: None, }) } diff --git a/src/jetstream.rs b/src/jetstream.rs index 6166a87..69256b7 100644 --- a/src/jetstream.rs +++ b/src/jetstream.rs @@ -17,6 +17,9 @@ pub(crate) struct JetstreamEphemeral { pub record: Option>, pub cid: Option, pub live: bool, + /// the key of the redaction marker that would erase `record`, when it + /// embeds a stored version + pub redaction_key: Option>, } impl JetstreamEphemeral { @@ -43,28 +46,34 @@ impl JetstreamEphemeral { #[cfg(all(feature = "indexer_stream", feature = "jetstream"))] pub(crate) fn build_ephemeral_from_stored( - did: &str, + did: &crate::db::types::TrimmedDid, rev: &str, operation: &str, collection: &str, - rkey: &str, + rkey: &crate::db::types::DbRkey, data: &crate::types::StoredData, inline_block: Option<&Bytes>, live: bool, ) -> Option { + use crate::db::keys; 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 mut redaction_key = None; 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(Arc::from(serde_json::value::to_raw_value(&val).ok()?)), - ), + Ok(val) => { + let record_key = keys::record_key_trimmed(did, collection, rkey); + redaction_key = Some(keys::redaction_key(&record_key, cid)); + ( + Some(cid.to_string()), + Some(Arc::from(serde_json::value::to_raw_value(&val).ok()?)), + ) + } Err(_) => return None, } } else { @@ -87,13 +96,14 @@ pub(crate) fn build_ephemeral_from_stored( }; Some(JetstreamEphemeral { - did: did.to_string(), + did: did.to_did().to_string(), rev: rev.to_string(), operation: operation.to_string(), collection: collection.to_string(), - rkey: rkey.to_string(), + rkey: rkey.to_smolstr().to_string(), record, cid, live, + redaction_key, }) } diff --git a/src/ops/record_events.rs b/src/ops/record_events.rs index 52e9bb6..2093962 100644 --- a/src/ops/record_events.rs +++ b/src/ops/record_events.rs @@ -110,11 +110,11 @@ mod enabled { { let json = self.should_stage_jetstream.then_some(()).and_then(|()| { crate::jetstream::build_ephemeral_from_stored( - did_trimmed.to_did().as_str(), + &did_trimmed, self.rev.to_tid().as_str(), op.action.as_str(), collection.as_str(), - op.rkey.to_smolstr().as_str(), + op.rkey, &data, inline_block.as_ref(), self.live, diff --git a/src/types/event/jetstream.rs b/src/types/event/jetstream.rs index 4ff4109..d0f722c 100644 --- a/src/types/event/jetstream.rs +++ b/src/types/event/jetstream.rs @@ -114,6 +114,8 @@ pub(crate) struct JetstreamLive { /// pre-serialized json, when a subscriber was listening as the event was /// staged. otherwise it is rendered from the database. pub(crate) json: Option, + /// the redaction marker that erases the record `json` embeds + pub(crate) redaction_key: Option>, } impl<'i, Id> StoredJetstreamEvent<'i, Id> {