From cbce4c071a797a21837357ac3f0bfe500e383a73 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Thu, 8 Oct 2026 21:21:06 +0300 Subject: [PATCH] [jetstream] send the live json as built, without a redaction read per reader jetstream already builds a live commit's json once at the commit and shares it, but every reader checked the record's redaction marker before sending it, a point read per reader per live event, and rendered it from the db when the marker was there. that check is gone, so the json goes out as built. a redaction that completes after the commit doesn't reach events already queued then, the same window /stream has, and docs/api/jetstream.md says so. the test that pinned the old read now pins this: a commit queued before the redaction still has its record. the redaction keys JetstreamLive and JetstreamEphemeral carried for it are gone too. --- docs/api/jetstream.md | 2 ++ src/api/debug.rs | 1 - src/control/stream/interleaving.rs | 10 ++++------ src/control/stream/jetstream.rs | 28 +++------------------------- src/db/keyspaces.rs | 6 ------ src/db/outbox/jetstream.rs | 5 +---- src/ingest/relay/sink/relay.rs | 1 - src/jetstream.rs | 16 ++-------------- src/types/event/jetstream.rs | 2 -- 9 files changed, 12 insertions(+), 59 deletions(-) diff --git a/docs/api/jetstream.md b/docs/api/jetstream.md index d3b9a1d..290f89e 100644 --- a/docs/api/jetstream.md +++ b/docs/api/jetstream.md @@ -69,6 +69,8 @@ fired when a record is created, updated, or deleted: `record` keeps the keys in the order the record's CBOR has them, shorter keys first, rather than sorting them. +a redaction stops a body from being served once it completes, except in events already queued for a reader at that moment. + if hydrant is running in indexer mode with `HYDRANT_ONLY_INDEX_LINKS=true`, record content blocks are not persisted. consequently, `record` and `cid` may be omitted from `commit` payloads. #### identity diff --git a/src/api/debug.rs b/src/api/debug.rs index f95314d..da5eeb3 100644 --- a/src/api/debug.rs +++ b/src/api/debug.rs @@ -616,7 +616,6 @@ 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 a7e79ae..35ebe68 100644 --- a/src/control/stream/interleaving.rs +++ b/src/control/stream/interleaving.rs @@ -281,7 +281,7 @@ mod jetstream { } #[test] - fn a_redaction_after_the_broadcast_keeps_the_body_out() { + fn a_redaction_doesnt_reach_events_already_queued() { let (_tmp, state, opts) = state(); let (tx, mut rx) = mpsc::channel(2); let filter = JetstreamFilter::new(JetstreamSubscriberOptions::default()); @@ -318,13 +318,11 @@ mod jetstream { serde_json::from_slice::(&item).unwrap()["commit"].take() }) .collect(); + // the third was queued before the redaction completed, so it goes out as committed assert!( - commits[..2] - .iter() - .all(|commit| commit.get("record").is_some()) + commits.iter().all(|commit| commit.get("record").is_some()), + "{commits:?}" ); - assert!(commits[2].get("record").is_none(), "{}", commits[2]); - assert!(commits[2].get("cid").is_some()); } #[test] diff --git a/src/control/stream/jetstream.rs b/src/control/stream/jetstream.rs index d4d64cb..193423f 100644 --- a/src/control/stream/jetstream.rs +++ b/src/control/stream/jetstream.rs @@ -97,13 +97,9 @@ impl StreamSource for JetstreamSource { if !self.filter.wants(&live.event) { return None; } - // 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(|| { + // json built at the commit goes out as it is, even past a redaction that completed + // since, see the redaction note in docs/api/jetstream.md + live.json.or_else(|| { stored_event_to_bytes(&self.state, live.position.time_us as i64, &live.event) }) } @@ -115,24 +111,6 @@ 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 5636bb9..4010eaa 100644 --- a/src/db/keyspaces.rs +++ b/src/db/keyspaces.rs @@ -103,12 +103,6 @@ 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 839ff9c..918edf0 100644 --- a/src/db/outbox/jetstream.rs +++ b/src/db/outbox/jetstream.rs @@ -57,14 +57,11 @@ impl JetstreamOutbox { return Ok(None); }; batch.insert(&db.events, position.key(), row); - let (json, redaction_key) = json.map_or((None, None), |json| { - (json.to_json(position.time_us as i64), json.redaction_key) - }); + let json = json.and_then(|json| json.to_json(position.time_us as i64)); Ok(Some(JetstreamLive { position, event, json, - redaction_key, })) }) .filter_map(Result::transpose) diff --git a/src/ingest/relay/sink/relay.rs b/src/ingest/relay/sink/relay.rs index ee706d4..0674d57 100644 --- a/src/ingest/relay/sink/relay.rs +++ b/src/ingest/relay/sink/relay.rs @@ -231,7 +231,6 @@ fn build_relay_commit_ephemeral( record, cid, live: true, - redaction_key: None, }) } diff --git a/src/jetstream.rs b/src/jetstream.rs index d1eb0e8..3e72224 100644 --- a/src/jetstream.rs +++ b/src/jetstream.rs @@ -17,9 +17,6 @@ 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 { @@ -55,23 +52,15 @@ pub(crate) fn build_ephemeral_from_stored( 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 { - let record = crate::types::record_json(bytes).ok()?; - let record_key = keys::record_key_trimmed(did, collection, rkey); - redaction_key = Some(keys::redaction_key(&record_key, cid)); - (Some(cid.to_string()), Some(record)) - } else { - return None; - } + let record = crate::types::record_json(inline_block?).ok()?; + (Some(cid.to_string()), Some(record)) } StoredData::Block(block) => { let digest = Sha256::digest(block); @@ -92,6 +81,5 @@ pub(crate) fn build_ephemeral_from_stored( record, cid, live, - redaction_key, }) } diff --git a/src/types/event/jetstream.rs b/src/types/event/jetstream.rs index d0f722c..4ff4109 100644 --- a/src/types/event/jetstream.rs +++ b/src/types/event/jetstream.rs @@ -114,8 +114,6 @@ 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> { -- 2.51.2