From f15c6a83a455ace25ff6ee2d677222e48077c536 Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Wed, 05 Aug 2026 15:20:40 +0000 Subject: [PATCH] [api] add operator record and repository deletion private management routes to redact stored record bodies without touching logical identity: DELETE /repos/{did} excludes, untracks, and erases every body a repo owns (heads, history, the compatibility archive, buffered resync commits), and DELETE /repos/{did}/records/{collection}/{rkey} scopes the same to one record with a head/history/all target. redaction replaces bodies with their raw CID bytes and writes permanent markers so a replayed version of the same record+CID stays suppressed even after an intervening version or history compaction. both routes are idempotent, chunk their rewrites in bounded key-ordered batches, and are rejected in ephemeral and links-only modes. issue: hydrant-0gg --- src/state.rs | 11 +++++++++++ src/types.rs | 34 ++++++++++++++++++++++++++++++++++ tests/api_repos.nu | 36 ++++++++++++++++++++++++++++++++++-- docs/api/repos.md | 8 ++++++++ src/api/repos.rs | 69 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++------- src/db/indexer.rs | 616 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++- src/db/txn.rs | 259 +++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++ src/control/repos/indexer.rs | 301 ++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++++----- src/db/keys/indexer.rs | 6 ++++++ 9 file(s) changed, 1325 insertion(s)(+), 15 deletion(s)(-) diff --git a/src/state.rs b/src/state.rs --- a/src/state.rs +++ b/src/state.rs @@ -154,6 +154,17 @@ } #[cfg(feature = "indexer")] + pub(crate) fn ensure_operator_body_deletion_supported(&self) -> Result<()> { + if self.ephemeral { + miette::bail!("operator body deletion is not supported in ephemeral mode"); + } + if !self.stores_record_bodies() { + miette::bail!("record bodies are not stored in links-only mode"); + } + Ok(()) + } + + #[cfg(feature = "indexer")] pub fn notify_backfill(&self) { self.backfill_notify.notify_one(); } diff --git a/src/types.rs b/src/types.rs --- a/src/types.rs +++ b/src/types.rs @@ -20,6 +20,40 @@ pub(crate) use event::*; pub(crate) use v7::*; +/// which record-scoped body copies an operator redaction removes. +#[cfg(feature = "indexer")] +#[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] +#[serde(rename_all = "snake_case")] +pub enum DeleteBodyTarget { + Head, + History, + All, +} + +#[cfg(feature = "indexer")] +impl DeleteBodyTarget { + pub(crate) fn includes_head(self) -> bool { + matches!(self, Self::Head | Self::All) + } + + pub(crate) fn includes_history(self) -> bool { + matches!(self, Self::History | Self::All) + } +} + +/// idempotent result of deleting indexed body bytes while preserving logical +/// record identity and repository topology. +#[cfg(feature = "indexer")] +#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)] +pub struct DeleteBodiesReport { + pub heads_redacted: u64, + pub history_redacted: u64, + pub event_bodies_redacted: u64, + pub body_bytes_redacted: u64, + pub buffered_commits_deleted: u64, + pub buffered_commit_bytes_deleted: u64, +} + impl<'c> From> for Commit { fn from(value: AtpCommit<'c>) -> Self { Self { diff --git a/tests/api_repos.nu b/tests/api_repos.nu --- a/tests/api_repos.nu +++ b/tests/api_repos.nu @@ -34,11 +34,11 @@ print "verifying /repos error handling..." print " testing PUT /repos with invalid DID..." - http put -f -e -t application/json $"($url)/repos" { did: "invalid" } + http put -f -e -t application/json $"($url)/repos" [{ did: "invalid" }] | assert-status 400 "PUT /repos invalid DID" print " testing DELETE /repos with invalid DID..." - http delete -f -e -t application/json $"($url)/repos" --data { did: "invalid" } + http delete -f -e -t application/json $"($url)/repos" --data [{ did: "invalid" }] | assert-status 400 "DELETE /repos invalid DID" print " testing GET /repos with invalid cursor..." @@ -46,6 +46,37 @@ | assert-status 400 "GET /repos invalid cursor" print "/repos validation errors passed!" +} + +def test-body-deletion-routes [url: string, did: string] { + print "verifying operator body deletion routes..." + + let record_url = $"($url)/repos/($did)/records/app.bsky.feed.post/3kzbif5moe22m" + http delete -f -e $record_url + | assert-status 400 "record deletion requires target" + + http delete -f -e $"($record_url)?target=nope" + | assert-status 400 "record deletion rejects invalid target" + + let record_report = (http delete $"($record_url)?target=all") + if $record_report.heads_redacted != 0 or $record_report.history_redacted != 0 { + fail "empty record deletion should be idempotent" + } + + http delete -f -e $"($url)/repos/not-a-did" + | assert-status 400 "repo drop rejects invalid DID" + + let repo_report = (http delete $"($url)/repos/($did)") + if $repo_report.heads_redacted != 0 or $repo_report.history_redacted != 0 { + fail "empty repo drop should be idempotent" + } + + let info = (http get $"($url)/repos/($did)") + if $info.tracked { + fail "repo drop should untrack the DID" + } + + print "operator body deletion routes passed!" } def main [] { @@ -71,6 +102,7 @@ test-repos-pagination $url test-repos-validation-errors $url + test-body-deletion-routes $url ($dids | first) kill $instance.pid } diff --git a/docs/api/repos.md b/docs/api/repos.md --- a/docs/api/repos.md +++ b/docs/api/repos.md @@ -25,6 +25,14 @@ untrack repositories. accepts an NDJSON body of `{"did": "..."}` (or JSON array of the same). only affects repositories that are currently tracked. returns a list of the DIDs that were untracked. +## DELETE /repos/{did} + +operator-only destructive drop. hydrant persistently excludes and untracks the DID, replaces every current/history body with its CID marker, removes legacy event-body copies, clears buffered commits and retry state, and removes outgoing backlinks. logical CIDs, counts, and repository topology remain as redacted metadata. remove the DID from `/filter` excludes before intentionally indexing future versions again. + +## DELETE /repos/{did}/records/{collection}/{rkey} + +delete body copies for one record. the required `target` query parameter is one of `head`, `history`, or `all`. a durable record+CID marker keeps the same version redacted across replay, protocol deletion, intervening versions, and history compaction; a genuinely new CID may store a new body. the response reports redacted copies and bytes. operator deletion is rejected in ephemeral and links-only modes. + ## POST /repos/resync force a new backfill for one or more repositories. accepts an NDJSON body of `{"did": "..."}` (or JSON array of the same). only affects repositories hydrant already knows about. returns a list of the DIDs that were queued. diff --git a/src/api/repos.rs b/src/api/repos.rs --- a/src/api/repos.rs +++ b/src/api/repos.rs @@ -12,6 +12,8 @@ routing::get, }; use jacquard_common::types::did::Did; +#[cfg(feature = "indexer")] +use jacquard_common::types::{nsid::Nsid, string::Rkey}; use serde::Deserialize; pub fn router() -> Router { @@ -24,7 +26,12 @@ let r = r .route("/repos/resync", post(handle_post_resync)) .route("/repos", put(handle_put_repos)) - .route("/repos", delete(handle_delete_repos)); + .route("/repos", delete(handle_delete_repos)) + .route("/repos/{did}", delete(handle_drop_repo)) + .route( + "/repos/{did}/records/{collection}/{rkey}", + delete(handle_delete_record_bodies), + ); r } @@ -33,6 +40,12 @@ #[derive(Deserialize, Debug)] pub struct RepoRequest { pub did: String, +} + +#[cfg(feature = "indexer")] +#[derive(Deserialize, Debug)] +pub struct DeleteRecordParams { + pub target: crate::types::DeleteBodyTarget, } #[derive(Deserialize)] @@ -109,8 +122,8 @@ let dids: Vec> = items .into_iter() - .filter_map(|item| Did::new_owned(&item.did).ok()) - .collect(); + .map(|item| Did::new_owned(&item.did).map_err(bad_request)) + .collect::>()?; let queued = hydrant .repos @@ -131,8 +144,8 @@ let dids: Vec> = items .into_iter() - .filter_map(|item| Did::new_owned(&item.did).ok()) - .collect(); + .map(|item| Did::new_owned(&item.did).map_err(bad_request)) + .collect::>()?; let untracked = hydrant .repos @@ -141,6 +154,48 @@ .map_err(|e| (StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; Ok(did_list_response(untracked, &headers)) +} + +#[cfg(feature = "indexer")] +pub async fn handle_drop_repo( + State(hydrant): State, + Path(did_str): Path, +) -> Result, (StatusCode, String)> { + ensure_body_deletion_supported(&hydrant)?; + let did = Did::new(&did_str).map_err(bad_request)?; + hydrant + .repos + .drop_repo(&did) + .await + .map(Json) + .map_err(internal) +} + +#[cfg(feature = "indexer")] +pub async fn handle_delete_record_bodies( + State(hydrant): State, + Path((did_str, collection, rkey)): Path<(String, String, String)>, + Query(params): Query, +) -> Result, (StatusCode, String)> { + ensure_body_deletion_supported(&hydrant)?; + let did = Did::new(&did_str).map_err(bad_request)?; + Nsid::new(&collection).map_err(bad_request)?; + Rkey::new(&rkey).map_err(bad_request)?; + hydrant + .repos + .get(&did) + .delete_record_bodies(&collection, &rkey, params.target) + .await + .map(Json) + .map_err(internal) +} + +#[cfg(feature = "indexer")] +fn ensure_body_deletion_supported(hydrant: &Hydrant) -> Result<(), (StatusCode, String)> { + hydrant + .state + .ensure_operator_body_deletion_supported() + .map_err(|e| (StatusCode::CONFLICT, e.to_string())) } #[cfg(feature = "indexer")] @@ -153,8 +208,8 @@ let dids: Vec> = items .into_iter() - .filter_map(|item| Did::new_owned(&item.did).ok()) - .collect(); + .map(|item| Did::new_owned(&item.did).map_err(bad_request)) + .collect::>()?; let queued = hydrant.repos.resync(dids).await.map_err(internal)?; diff --git a/src/db/indexer.rs b/src/db/indexer.rs --- a/src/db/indexer.rs +++ b/src/db/indexer.rs @@ -8,7 +8,7 @@ use crate::db::types::DbTid; use crate::db::{Db, Txn, deser_repo_state, keys, ser_repo_state}; -use crate::types::RepoState; +use crate::types::{DeleteBodiesReport, DeleteBodyTarget, RepoState}; use std::collections::BTreeMap; impl Db { @@ -107,6 +107,291 @@ pub(crate) fn is_cid_record_value(value: &[u8]) -> bool { value.len() == 36 && value.starts_with(&[0x01, 0x71, 0x12, 0x20]) } + +struct RedactedValue { + cid: jacquard_common::types::cid::IpldCid, + body_bytes: Option, +} + +fn redact_value( + batch: &mut OwnedWriteBatch, + keyspace: &Keyspace, + key: &[u8], + value: &[u8], +) -> Result> { + if value.is_empty() { + return Ok(None); + } + if is_cid_record_value(value) { + return Ok(Some(RedactedValue { + cid: ipld_cid_from_record_value(value)?, + body_bytes: None, + })); + } + let cid = jacquard_repo::mst::util::compute_cid(value).into_diagnostic()?; + batch.insert(keyspace, key, cid.to_bytes()); + Ok(Some(RedactedValue { + cid, + body_bytes: Some(value.len() as u64), + })) +} + +fn stage_redaction( + batch: &mut OwnedWriteBatch, + db: &Db, + record_key: &[u8], + cid: &jacquard_common::types::cid::IpldCid, +) { + batch.insert( + &db.indexer.redactions, + keys::redaction_key(record_key, cid), + [], + ); +} + +#[cfg(feature = "indexer_stream")] +fn delete_event_body( + batch: &mut OwnedWriteBatch, + db: &Db, + key: &[u8], + value: &[u8], +) -> Result> { + let Some((record_key, cid)) = keys::split_record_cid_key(key) else { + miette::bail!("invalid compatibility event-body key"); + }; + let computed = jacquard_repo::mst::util::compute_cid(value).into_diagnostic()?; + if computed != cid { + miette::bail!("compatibility event body does not match its cid"); + } + stage_redaction(batch, db, record_key, &cid); + batch.remove(db.stream.event_bodies.raw(), key); + Ok(Some(value.len() as u64)) +} + +const REDACTION_CHUNK_ENTRIES: usize = 4_096; +const REDACTION_CHUNK_BYTES: u64 = 32 * 1024 * 1024; + +/// rewrite one prefix in bounded, key-ordered batches. the caller holds the +/// repository write lock for the whole operation, so the cursor cannot skip a +/// concurrent insert. committed chunks are independently valid and retrying +/// after a crash is idempotent. +fn rewrite_prefix( + db: &Db, + keyspace: &Keyspace, + prefix: &[u8], + mut rewrite: F, +) -> Result<(u64, u64)> +where + F: FnMut(&mut OwnedWriteBatch, &[u8], &[u8]) -> Result>, +{ + use std::ops::Bound; + + let mut redacted = 0_u64; + let mut bytes = 0_u64; + let mut cursor: Option> = None; + + loop { + let start = cursor + .as_ref() + .map_or(Bound::Included(prefix.to_vec()), |key| { + Bound::Excluded(key.clone()) + }); + let mut batch = db.inner.batch(); + let mut scanned = 0_usize; + let mut staged_bytes = 0_u64; + let mut last = None; + + for guard in keyspace.range((start, Bound::>::Unbounded)) { + let (key, value) = guard.into_inner().into_diagnostic()?; + if !key.starts_with(prefix) { + break; + } + if scanned > 0 + && (scanned >= REDACTION_CHUNK_ENTRIES || staged_bytes >= REDACTION_CHUNK_BYTES) + { + break; + } + if let Some(len) = rewrite(&mut batch, &key, &value)? { + redacted += 1; + bytes += len; + staged_bytes += len; + } + scanned += 1; + last = Some(key.to_vec()); + } + + let Some(last) = last else { break }; + batch.commit().into_diagnostic()?; + cursor = Some(last); + } + Ok((redacted, bytes)) +} + +/// erase one record's body bytes without changing its CID identity, counts, +/// or the repository MST root. history tombstones are metadata and remain. +pub(crate) fn redact_record_bodies( + db: &Db, + did: &Did<'_>, + collection: &str, + rkey: &crate::db::types::DbRkey, + target: DeleteBodyTarget, +) -> Result { + let lock_index = crate::db::record_lock_index_for_did(did); + let _guard = db.record_write_locks[lock_index as usize] + .lock() + .unwrap_or_else(|e| e.into_inner()); + let record_key = keys::record_key(did, collection, rkey); + let mut report = DeleteBodiesReport::default(); + + if target.includes_head() { + let mut batch = db.inner.batch(); + if let Some(value) = db.indexer.record(&record_key).into_diagnostic()? + && let Some(redacted) = + redact_value(&mut batch, db.indexer.records.raw(), &record_key, &value)? + { + stage_redaction(&mut batch, db, &record_key, &redacted.cid); + if let Some(len) = redacted.body_bytes { + report.heads_redacted = 1; + report.body_bytes_redacted += len; + } + crate::ops::backlink_ops::delete_record( + &mut batch, + db, + did, + collection, + &rkey.to_smolstr(), + )?; + } + batch.commit().into_diagnostic()?; + } + + if target.includes_history() { + let mut prefix = record_key.clone(); + prefix.push(keys::indexer::REC_SEP); + let (count, bytes) = rewrite_prefix( + db, + db.indexer.history.raw(), + &prefix, + |batch, key, value| { + let Some(redacted) = redact_value(batch, db.indexer.history.raw(), key, value)? + else { + return Ok(None); + }; + stage_redaction(batch, db, &record_key, &redacted.cid); + Ok(redacted.body_bytes) + }, + )?; + report.history_redacted = count; + report.body_bytes_redacted += bytes; + + #[cfg(feature = "indexer_stream")] + { + let (count, bytes) = rewrite_prefix( + db, + db.stream.event_bodies.raw(), + &prefix, + |batch, key, value| delete_event_body(batch, db, key, value), + )?; + report.event_bodies_redacted = count; + report.body_bytes_redacted += bytes; + } + } + + db.persist()?; + Ok(report) +} + +/// erase every body owned by a repository while retaining CIDs, counts, and +/// repo state. this makes the operation safe and idempotent without forging a +/// new authoritative commit. +pub(crate) fn redact_repo_bodies( + db: &Db, + did: &Did<'_>, + target: DeleteBodyTarget, +) -> Result { + let lock_index = crate::db::record_lock_index_for_did(did); + let _guard = db.record_write_locks[lock_index as usize] + .lock() + .unwrap_or_else(|e| e.into_inner()); + let prefix = keys::record_prefix_did(did); + let mut report = DeleteBodiesReport::default(); + + if target.includes_head() { + let (count, bytes) = rewrite_prefix( + db, + db.indexer.records.raw(), + &prefix, + |batch, key, value| { + let Some(redacted) = redact_value(batch, db.indexer.records.raw(), key, value)? + else { + return Ok(None); + }; + stage_redaction(batch, db, key, &redacted.cid); + let suffix = &key[prefix.len()..]; + let (collection, rkey) = keys::split_record_suffix(suffix)?; + crate::ops::backlink_ops::delete_record( + batch, + db, + did, + collection, + &rkey.to_smolstr(), + )?; + Ok(redacted.body_bytes) + }, + )?; + report.heads_redacted = count; + report.body_bytes_redacted += bytes; + } + + if target.includes_history() { + let (count, bytes) = rewrite_prefix( + db, + db.indexer.history.raw(), + &prefix, + |batch, key, value| { + let Some(redacted) = redact_value(batch, db.indexer.history.raw(), key, value)? + else { + return Ok(None); + }; + let record_key = keys::history_record_key(key) + .ok_or_else(|| miette::miette!("invalid history key during redaction"))?; + stage_redaction(batch, db, record_key, &redacted.cid); + Ok(redacted.body_bytes) + }, + )?; + report.history_redacted = count; + report.body_bytes_redacted += bytes; + + #[cfg(feature = "indexer_stream")] + { + let (count, bytes) = rewrite_prefix( + db, + db.stream.event_bodies.raw(), + &prefix, + |batch, key, value| delete_event_body(batch, db, key, value), + )?; + report.event_bodies_redacted = count; + report.body_bytes_redacted += bytes; + } + } + + let buffer_prefix = keys::resync_buffer_prefix(did); + let (count, bytes) = rewrite_prefix( + db, + db.indexer.resync_buffer.raw(), + &buffer_prefix, + |batch, key, value| { + batch.remove(db.indexer.resync_buffer.raw(), key); + Ok(Some(value.len() as u64)) + }, + )?; + report.buffered_commits_deleted = count; + report.buffered_commit_bytes_deleted = bytes; + + db.persist()?; + Ok(report) +} + pub(crate) fn delete_repo_records( txn: &mut Txn<'_>, db: &Db, @@ -288,6 +573,335 @@ #[cfg(test)] mod tests { use super::*; + + const DID: &str = "did:plc:yk4q3id7id6p5z3bypvshc64"; + const OTHER_DID: &str = "did:plc:ewvi7nxzyoun6zhxrhs64oiz"; + const COLLECTION: &str = "app.bsky.feed.post"; + const RKEY: &str = "3kzbif5moe22m"; + + fn test_db() -> Result<(tempfile::TempDir, Db)> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let cfg = crate::config::Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }; + let db = Db::open(&cfg)?; + Ok((tmp, db)) + } + + fn body(text: &str) -> Vec { + serde_ipld_dagcbor::to_vec(&serde_json::json!({ + "$type": COLLECTION, + "text": text, + })) + .unwrap() + } + + fn did(value: &'static str) -> Did<'static> { + Did::new(value).unwrap() + } + + fn record_key(did: &Did<'_>) -> Vec { + keys::record_key(did, COLLECTION, &crate::db::types::DbRkey::new(RKEY)) + } + + fn stage_body_copies(db: &Db, did: &Did<'_>) -> Result<(Vec, Vec, Vec)> { + let head = body("head"); + let old_a = body("old a"); + let old_b = body("old b"); + #[cfg(feature = "indexer_stream")] + let archive = body("archive"); + #[cfg(feature = "indexer_stream")] + let archive_cid = jacquard_repo::mst::util::compute_cid(&archive).into_diagnostic()?; + let key = record_key(did); + let mut batch = db.inner.batch(); + batch.insert(&db.indexer.records, &key, &head); + batch.insert( + &db.indexer.history, + keys::history_key(&key, &DbTid::new_from_bytes(1_u64.to_be_bytes())), + &old_a, + ); + batch.insert( + &db.indexer.history, + keys::history_key(&key, &DbTid::new_from_bytes(2_u64.to_be_bytes())), + &old_b, + ); + batch.insert( + &db.indexer.history, + keys::history_key(&key, &DbTid::new_from_bytes(3_u64.to_be_bytes())), + [], + ); + #[cfg(feature = "indexer_stream")] + batch.insert( + &db.stream.event_bodies, + keys::event_body_key(&key, &archive_cid), + &archive, + ); + batch.commit().into_diagnostic()?; + Ok((head, old_a, old_b)) + } + + #[test] + fn record_head_redaction_is_idempotent_and_keeps_history() -> Result<()> { + let (_tmp, db) = test_db()?; + let did = did(DID); + let (head, old_a, _) = stage_body_copies(&db, &did)?; + + let report = redact_record_bodies( + &db, + &did, + COLLECTION, + &crate::db::types::DbRkey::new(RKEY), + DeleteBodyTarget::Head, + )?; + assert_eq!(report.heads_redacted, 1); + assert_eq!(report.history_redacted, 0); + assert_eq!(report.body_bytes_redacted, head.len() as u64); + + let key = record_key(&did); + let marker = db.indexer.record(&key).into_diagnostic()?.unwrap(); + assert!(is_cid_record_value(&marker)); + assert_eq!( + db.indexer + .history + .get(keys::history_key( + &key, + &DbTid::new_from_bytes(1_u64.to_be_bytes()) + )) + .into_diagnostic()? + .as_deref(), + Some(old_a.as_slice()) + ); + + let again = redact_record_bodies( + &db, + &did, + COLLECTION, + &crate::db::types::DbRkey::new(RKEY), + DeleteBodyTarget::Head, + )?; + assert_eq!(again, DeleteBodiesReport::default()); + Ok(()) + } + + #[test] + fn record_history_redaction_keeps_tombstones_and_removes_archive() -> Result<()> { + let (_tmp, db) = test_db()?; + let did = did(DID); + let (head, old_a, old_b) = stage_body_copies(&db, &did)?; + + let report = redact_record_bodies( + &db, + &did, + COLLECTION, + &crate::db::types::DbRkey::new(RKEY), + DeleteBodyTarget::History, + )?; + assert_eq!(report.heads_redacted, 0); + assert_eq!(report.history_redacted, 2); + #[cfg(feature = "indexer_stream")] + assert_eq!(report.event_bodies_redacted, 1); + assert_eq!( + report.body_bytes_redacted, + (old_a.len() + + old_b.len() + + if cfg!(feature = "indexer_stream") { + body("archive").len() + } else { + 0 + }) as u64 + ); + + let key = record_key(&did); + assert_eq!( + db.indexer.record(&key).into_diagnostic()?.as_deref(), + Some(head.as_slice()) + ); + assert!( + db.indexer + .history + .get(keys::history_key( + &key, + &DbTid::new_from_bytes(1_u64.to_be_bytes()) + )) + .into_diagnostic()? + .is_some_and(|value| is_cid_record_value(&value)) + ); + assert!( + db.indexer + .history + .get(keys::history_key( + &key, + &DbTid::new_from_bytes(3_u64.to_be_bytes()) + )) + .into_diagnostic()? + .is_some_and(|value| value.is_empty()) + ); + #[cfg(feature = "indexer_stream")] + assert!( + db.stream + .event_bodies + .prefix([key.as_slice(), &[keys::indexer::REC_SEP]].concat()) + .next() + .is_none() + ); + Ok(()) + } + + #[test] + fn record_redaction_does_not_discard_repo_resync_buffer() -> Result<()> { + let (_tmp, db) = test_db()?; + let did = did(DID); + stage_body_copies(&db, &did)?; + let buffer_key = keys::resync_buffer_key(&did, DbTid::new_from_bytes(4_u64.to_be_bytes())); + db.indexer + .resync_buffer + .insert(&buffer_key, b"unrelated buffered commit") + .into_diagnostic()?; + + redact_record_bodies( + &db, + &did, + COLLECTION, + &crate::db::types::DbRkey::new(RKEY), + DeleteBodyTarget::All, + )?; + + assert_eq!( + db.indexer + .resync_buffer + .get(buffer_key) + .into_diagnostic()? + .as_deref(), + Some(b"unrelated buffered commit".as_slice()) + ); + Ok(()) + } + + #[test] + fn repo_redaction_is_did_prefix_isolated() -> Result<()> { + let (_tmp, db) = test_db()?; + let target = did(DID); + let other = did(OTHER_DID); + stage_body_copies(&db, &target)?; + let (other_head, other_old, _) = stage_body_copies(&db, &other)?; + + let report = redact_repo_bodies(&db, &target, DeleteBodyTarget::All)?; + assert_eq!(report.heads_redacted, 1); + assert_eq!(report.history_redacted, 2); + + assert!( + db.indexer + .record(record_key(&target)) + .into_diagnostic()? + .is_some_and(|value| is_cid_record_value(&value)) + ); + let other_key = record_key(&other); + assert_eq!( + db.indexer.record(&other_key).into_diagnostic()?.as_deref(), + Some(other_head.as_slice()) + ); + assert_eq!( + db.indexer + .history + .get(keys::history_key( + &other_key, + &DbTid::new_from_bytes(1_u64.to_be_bytes()) + )) + .into_diagnostic()? + .as_deref(), + Some(other_old.as_slice()) + ); + Ok(()) + } + + #[test] + fn redaction_chunks_large_history_prefixes() -> Result<()> { + let (_tmp, db) = test_db()?; + let did = did(DID); + let key = record_key(&did); + let body = body("version"); + let mut batch = db.inner.batch(); + for i in 0..(REDACTION_CHUNK_ENTRIES as u64 + 1) { + batch.insert( + &db.indexer.history, + keys::history_key(&key, &DbTid::new_from_bytes(i.to_be_bytes())), + &body, + ); + } + batch.commit().into_diagnostic()?; + + let report = redact_record_bodies( + &db, + &did, + COLLECTION, + &crate::db::types::DbRkey::new(RKEY), + DeleteBodyTarget::History, + )?; + assert_eq!(report.history_redacted, REDACTION_CHUNK_ENTRIES as u64 + 1); + assert!( + db.indexer + .history + .prefix([key.as_slice(), &[keys::indexer::REC_SEP]].concat()) + .all(|item| item.value().is_ok_and(|value| is_cid_record_value(&value))) + ); + Ok(()) + } + + #[cfg(feature = "backlinks")] + #[test] + fn only_head_redaction_removes_outgoing_backlinks() -> Result<()> { + let (_tmp, db) = test_db()?; + let did = did(DID); + let rkey = crate::db::types::DbRkey::new(RKEY); + let value = serde_json::json!({ + "$type": "app.bsky.feed.like", + "subject": { + "uri": "at://did:plc:target/app.bsky.feed.post/3jz", + "cid": "bafyreig2toyab2fv7kqv3wcldlgl63p2oqxv54htv23s3kbsnshs7j2yli" + }, + "createdAt": "2026-08-04T00:00:00Z" + }); + let body = serde_ipld_dagcbor::to_vec(&value).into_diagnostic()?; + let data = serde_ipld_dagcbor::from_slice(&body).into_diagnostic()?; + let mut batch = db.inner.batch(); + db.indexer.stage_record( + &mut batch, + keys::record_key(&did, "app.bsky.feed.like", &rkey), + &body, + ); + crate::backlinks::store::index_record( + &mut batch, + &db.backlinks, + did.as_str(), + "app.bsky.feed.like", + RKEY, + &data, + )?; + batch.commit().into_diagnostic()?; + let forward = + crate::backlinks::store::forward_key(did.as_str(), "app.bsky.feed.like", RKEY); + assert!(db.backlinks.contains_key(&forward).into_diagnostic()?); + + redact_record_bodies( + &db, + &did, + "app.bsky.feed.like", + &rkey, + DeleteBodyTarget::History, + )?; + assert!(db.backlinks.contains_key(&forward).into_diagnostic()?); + + redact_record_bodies( + &db, + &did, + "app.bsky.feed.like", + &rkey, + DeleteBodyTarget::Head, + )?; + assert!(!db.backlinks.contains_key(&forward).into_diagnostic()?); + Ok(()) + } #[test] fn replace_record_counts_clears_stale_collections() -> Result<()> { diff --git a/src/db/txn.rs b/src/db/txn.rs --- a/src/db/txn.rs +++ b/src/db/txn.rs @@ -888,6 +888,265 @@ } #[test] + fn same_cid_replay_does_not_restore_a_redacted_head() { + let (_tmp, state) = test_state(); + let collection = "app.bsky.feed.post"; + let rkey = "3kzbif5moe22m"; + let body = body("v1"); + write( + &state, + &rev("3kzbif5moe22m"), + collection, + rkey, + &body, + DbAction::Create, + ); + + crate::db::redact_record_bodies( + &state.db, + &did(), + collection, + &DbRkey::new(rkey), + crate::types::DeleteBodyTarget::Head, + ) + .unwrap(); + write( + &state, + &rev("3kzbif5mof33m"), + collection, + rkey, + &body, + DbAction::Update, + ); + + let head = head(&state, collection, rkey).unwrap(); + assert!(crate::db::is_cid_record_value(&head)); + assert!(history_bodies(&state, collection, rkey).is_empty()); + } + + #[test] + fn genuinely_new_cid_replaces_a_redacted_head_without_archiving_marker() { + let (_tmp, state) = test_state(); + let collection = "app.bsky.feed.post"; + let rkey = "3kzbif5moe22m"; + write( + &state, + &rev("3kzbif5moe22m"), + collection, + rkey, + &body("v1"), + DbAction::Create, + ); + crate::db::redact_record_bodies( + &state.db, + &did(), + collection, + &DbRkey::new(rkey), + crate::types::DeleteBodyTarget::Head, + ) + .unwrap(); + + write( + &state, + &rev("3kzbif5mof33m"), + collection, + rkey, + &body("v2"), + DbAction::Update, + ); + + assert_eq!(head(&state, collection, rkey), Some(b"v2".to_vec())); + assert!(history_bodies(&state, collection, rkey).is_empty()); + } + + #[test] + fn redacted_cid_stays_redacted_after_an_intervening_version() { + let (_tmp, state) = test_state(); + let collection = "app.bsky.feed.post"; + let rkey = "3kzbif5moe22m"; + let v1 = body("v1"); + write( + &state, + &rev("3kzbif5moe22m"), + collection, + rkey, + &v1, + DbAction::Create, + ); + crate::db::redact_record_bodies( + &state.db, + &did(), + collection, + &DbRkey::new(rkey), + crate::types::DeleteBodyTarget::Head, + ) + .unwrap(); + + write( + &state, + &rev("3kzbif5mof33m"), + collection, + rkey, + &body("v2"), + DbAction::Update, + ); + write( + &state, + &rev("3kzbif5mog44m"), + collection, + rkey, + &v1, + DbAction::Update, + ); + + let head = head(&state, collection, rkey).unwrap(); + assert!(crate::db::is_cid_record_value(&head)); + } + + #[test] + fn redacted_cid_stays_redacted_after_reopen() { + let tmp = tempfile::tempdir().unwrap(); + let config = crate::config::Config { + database_path: tmp.path().to_path_buf(), + ..Default::default() + }; + let collection = "app.bsky.feed.post"; + let rkey = "3kzbif5moe22m"; + let v1 = body("v1"); + + { + let state = AppState::new(&config).unwrap(); + write( + &state, + &rev("3kzbif5moe22m"), + collection, + rkey, + &v1, + DbAction::Create, + ); + crate::db::redact_record_bodies( + &state.db, + &did(), + collection, + &DbRkey::new(rkey), + crate::types::DeleteBodyTarget::Head, + ) + .unwrap(); + state.db.persist().unwrap(); + } + + let state = AppState::new(&config).unwrap(); + write( + &state, + &rev("3kzbif5mof33m"), + collection, + rkey, + &body("v2"), + DbAction::Update, + ); + write( + &state, + &rev("3kzbif5mog44m"), + collection, + rkey, + &v1, + DbAction::Update, + ); + + assert!(crate::db::is_cid_record_value( + &head(&state, collection, rkey).unwrap() + )); + } + + #[test] + fn protocol_delete_of_redacted_head_writes_only_a_tombstone() { + let (_tmp, state) = test_state(); + let collection = "app.bsky.feed.post"; + let rkey = "3kzbif5moe22m"; + write( + &state, + &rev("3kzbif5moe22m"), + collection, + rkey, + &body("v1"), + DbAction::Create, + ); + crate::db::redact_record_bodies( + &state.db, + &did(), + collection, + &DbRkey::new(rkey), + crate::types::DeleteBodyTarget::Head, + ) + .unwrap(); + + delete(&state, &rev("3kzbif5mof33m"), collection, rkey); + + assert_eq!(head(&state, collection, rkey), None); + assert_eq!( + history_bodies(&state, collection, rkey), + vec![Vec::::new()] + ); + + write( + &state, + &rev("3kzbif5mog44m"), + collection, + rkey, + &body("v1"), + DbAction::Create, + ); + assert!(crate::db::is_cid_record_value( + &head(&state, collection, rkey).unwrap() + )); + } + + #[test] + fn history_redaction_suppresses_that_version_when_it_returns_to_head() { + let (_tmp, state) = test_state(); + let collection = "app.bsky.feed.post"; + let rkey = "3kzbif5moe22m"; + let v1 = body("v1"); + write( + &state, + &rev("3kzbif5moe22m"), + collection, + rkey, + &v1, + DbAction::Create, + ); + write( + &state, + &rev("3kzbif5mof33m"), + collection, + rkey, + &body("v2"), + DbAction::Update, + ); + crate::db::redact_record_bodies( + &state.db, + &did(), + collection, + &DbRkey::new(rkey), + crate::types::DeleteBodyTarget::History, + ) + .unwrap(); + + write( + &state, + &rev("3kzbif5mog44m"), + collection, + rkey, + &v1, + DbAction::Update, + ); + + assert!(crate::db::is_cid_record_value( + &head(&state, collection, rkey).unwrap() + )); + } + + #[test] fn create_over_differing_head_preserves_it() { let (_tmp, state) = test_state(); write( diff --git a/src/control/repos/indexer.rs b/src/control/repos/indexer.rs --- a/src/control/repos/indexer.rs +++ b/src/control/repos/indexer.rs @@ -321,9 +321,76 @@ Ok(untracked) } + /// permanently stop ingesting a repository and erase all of its stored + /// record bodies. logical CIDs and repo state remain as redacted metadata; + /// removing the DID from the filter excludes set permits future indexing. + pub async fn drop_repo(&self, did: &Did<'_>) -> Result { + self.0.ensure_operator_body_deletion_supported()?; + + let did = did.clone().into_static(); + self.0 + .db + .run(move |db| { + // make exclusion and untracking durable before redaction. + // writers re-check exclusion under this same record lock. + let mut txn = crate::db::Txn::new(db); + txn.hold_repo_write_lock(&did); + txn.batch.insert( + &db.filter, + crate::db::filter::exclude_key(did.as_str())?, + [], + ); + + let metadata_key = keys::repo_metadata_key(&did); + if let Some(bytes) = db.repo_metadata.get(&metadata_key).into_diagnostic()? { + let mut metadata = crate::db::deser_repo_meta(&bytes)?; + txn.batch + .remove(&db.indexer.pending, keys::pending_key(metadata.index_id)); + metadata.tracked = false; + txn.batch.insert( + &db.repo_metadata, + metadata_key, + crate::db::ser_repo_meta(&metadata)?, + ); + } + + let repo_key = keys::repo_key(&did); + txn.batch.remove(&db.indexer.resync, &repo_key); + txn.batch.remove(&db.crawler, keys::crawler_retry_key(&did)); + if db.repos.contains_key(&repo_key).into_diagnostic()? { + txn.transition_lifecycle(&did, GaugeState::Synced)?; + } + txn.commit()?; + db.persist()?; + + crate::db::redact_repo_bodies(db, &did, crate::types::DeleteBodyTarget::All) + }) + .await + } } impl<'i> RepoHandle<'i> { + /// erase selected body copies for one record while preserving its CID and + /// the repository's authoritative topology. + pub async fn delete_record_bodies( + &self, + collection: &str, + rkey: &str, + target: crate::types::DeleteBodyTarget, + ) -> Result { + self.state.ensure_operator_body_deletion_supported()?; + + Nsid::new(collection).into_diagnostic()?; + Rkey::new(rkey).into_diagnostic()?; + let did = self.did.clone().into_static(); + let collection = collection.to_owned(); + let rkey = DbRkey::new(rkey); + self.state + .db + .run(move |db| crate::db::redact_record_bodies(db, &did, &collection, &rkey, target)) + .await + } + /// gets a record from this repository. pub async fn get_record(&self, collection: &str, rkey: &str) -> Result> { if !self.state.stores_record_bodies() { @@ -479,6 +546,9 @@ &self, ) -> Result> + Send + 'static>> { + if self.state.ephemeral { + miette::bail!("CAR generation is not supported in ephemeral mode"); + } if !self.state.stores_record_bodies() { miette::bail!("block storage is not available in this mode"); } @@ -520,6 +590,10 @@ let (collection, rkey) = keys::split_record_suffix(&key[prefix.len()..])?; let mst_key = format!("{collection}/{rkey}"); + if crate::db::is_cid_record_value(&body) { + miette::bail!("record body was deleted for {}/{mst_key}", did); + } + // the cid is no longer stored; recompute it from the body let ipld_cid = jacquard_repo::mst::util::compute_cid(body.as_ref()) .into_diagnostic() @@ -547,11 +621,10 @@ // sanity check: rebuilt root should match stored commit data in full-index mode let computed_root = mst.get_pointer().await.into_diagnostic()?; if computed_root != atp_commit.data { - tracing::warn!( - computed = %computed_root, - stored = %atp_commit.data, - did = %self.did, - "mst root mismatch (expected in filter mode)", + miette::bail!( + "mst root mismatch for {}: computed {computed_root} != stored {}", + self.did, + atp_commit.data, ); } store @@ -667,4 +740,222 @@ Ok(()) } + + #[tokio::test] + async fn drop_repo_excludes_untracks_and_redacts_all_body_storage() -> miette::Result<()> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let mut config = crate::config::Config::default(); + config.database_path = tmp.path().to_path_buf(); + let state = Arc::new(AppState::new(&config)?); + let did = Did::new("did:plc:testoperatorrepo").into_diagnostic()?; + let metadata = crate::types::RepoMetadata::backfilling(7); + let head = record_body("head"); + let old = record_body("old"); + #[cfg(feature = "indexer_stream")] + let archived = record_body("archived"); + #[cfg(feature = "indexer_stream")] + let archived_cid = jacquard_repo::mst::util::compute_cid(&archived).into_diagnostic()?; + let record_key = + keys::record_key(&did, "app.bsky.feed.post", &DbRkey::new("3kzbif5moe22m")); + let mut batch = state.db.inner.batch(); + batch.insert( + &state.db.repos, + keys::repo_key(&did), + crate::db::ser_repo_state(&crate::types::RepoState::synced())?, + ); + batch.insert( + &state.db.repo_metadata, + keys::repo_metadata_key(&did), + crate::db::ser_repo_meta(&metadata)?, + ); + batch.insert( + &state.db.indexer.pending, + keys::pending_key(7), + keys::repo_key(&did), + ); + batch.insert(&state.db.crawler, keys::crawler_retry_key(&did), []); + state + .db + .indexer + .stage_record(&mut batch, &record_key, &head); + state.db.indexer.stage_history( + &mut batch, + keys::history_key( + &record_key, + &crate::db::types::DbTid::new_from_bytes(1_u64.to_be_bytes()), + ), + &old, + ); + #[cfg(feature = "indexer_stream")] + batch.insert( + &state.db.stream.event_bodies, + keys::event_body_key(&record_key, &archived_cid), + &archived, + ); + batch.insert( + &state.db.indexer.resync_buffer, + keys::resync_buffer_key( + &did, + crate::db::types::DbTid::new_from_bytes(2_u64.to_be_bytes()), + ), + b"buffered body bytes", + ); + batch.commit().into_diagnostic()?; + + let repos = ReposControl(state.clone()); + let report = repos.drop_repo(&did).await?; + assert_eq!(report.heads_redacted, 1); + assert_eq!(report.history_redacted, 1); + assert_eq!(report.buffered_commits_deleted, 1); + assert!( + state + .db + .filter + .contains_key(crate::db::filter::exclude_key(did.as_str())?) + .into_diagnostic()? + ); + let stored_meta = state + .db + .repo_metadata + .get(keys::repo_metadata_key(&did)) + .into_diagnostic()? + .expect("repo metadata remains"); + assert!(!crate::db::deser_repo_meta(&stored_meta)?.tracked); + assert!(repos.track([did.clone()]).await?.is_empty()); + assert!( + !state + .db + .crawler + .contains_key(keys::crawler_retry_key(&did)) + .into_diagnostic()? + ); + assert!( + state + .db + .indexer + .record(&record_key) + .into_diagnostic()? + .is_some_and(|value| crate::db::is_cid_record_value(&value)) + ); + assert!( + state + .db + .indexer + .resync_buffer + .prefix(keys::resync_buffer_prefix(&did)) + .next() + .is_none() + ); + #[cfg(feature = "indexer_stream")] + assert!( + state + .db + .stream + .event_bodies + .prefix([record_key.as_slice(), &[keys::indexer::REC_SEP]].concat()) + .next() + .is_none() + ); + Ok(()) + } + + #[tokio::test] + #[cfg(not(feature = "backlinks"))] + async fn test_generate_car_ephemeral_rejected() -> miette::Result<()> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let mut config = crate::config::Config::default(); + config.database_path = tmp.path().to_path_buf(); + config.ephemeral = true; + + let state = Arc::new(AppState::new(&config)?); + let repos = ReposControl(state); + let did = Did::new("did:plc:testephemeral").into_diagnostic()?; + let handle = repos.get(&did); + let res = handle.generate_car().await; + let Err(err) = res else { + panic!("generate_car in ephemeral mode should fail"); + }; + assert!(err.to_string().contains("ephemeral mode")); + Ok(()) + } + #[tokio::test] + async fn test_generate_car_root_mismatch_fails() -> miette::Result<()> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let mut config = crate::config::Config::default(); + config.database_path = tmp.path().to_path_buf(); + config.ephemeral = false; + + let state = Arc::new(AppState::new(&config)?); + let did = Did::new("did:plc:testmismatch").into_diagnostic()?; + + let dummy_cid = jacquard_repo::mst::util::compute_cid(b"dummy_root").into_diagnostic()?; + let repo_state = crate::types::RepoState { + root: Some(crate::types::v2::Commit { + version: 3, + rev: crate::db::types::DbTid::from(&jacquard_common::types::string::Tid::now_0()), + data: dummy_cid, + prev: None, + sig: bytes::Bytes::new(), + }), + ..crate::types::RepoState::synced() + }; + + let did_key = keys::repo_key(&did); + state + .db + .repos + .insert(&did_key, crate::db::ser_repo_state(&repo_state)?) + .into_diagnostic()?; + + let repos = ReposControl(state); + let handle = repos.get(&did); + + let res = handle.generate_car().await; + let Err(err) = res else { + panic!("generate_car with mismatched root should fail"); + }; + assert!(err.to_string().contains("mst root mismatch")); + Ok(()) + } + + #[tokio::test] + async fn generate_car_reports_a_deleted_head_body() -> miette::Result<()> { + let tmp = tempfile::tempdir().into_diagnostic()?; + let mut config = crate::config::Config::default(); + config.database_path = tmp.path().to_path_buf(); + let state = Arc::new(AppState::new(&config)?); + let did = Did::new("did:plc:testredactedcar").into_diagnostic()?; + let body = record_body("deleted"); + let cid = jacquard_repo::mst::util::compute_cid(&body).into_diagnostic()?; + let repo_state = crate::types::RepoState { + root: Some(crate::types::v2::Commit { + version: 3, + rev: crate::db::types::DbTid::from(&jacquard_common::types::string::Tid::now_0()), + data: cid, + prev: None, + sig: bytes::Bytes::new(), + }), + ..crate::types::RepoState::synced() + }; + let mut batch = state.db.inner.batch(); + batch.insert( + &state.db.repos, + keys::repo_key(&did), + crate::db::ser_repo_state(&repo_state)?, + ); + state.db.indexer.stage_record( + &mut batch, + keys::record_key(&did, "app.bsky.feed.post", &DbRkey::new("self")), + cid.to_bytes(), + ); + batch.commit().into_diagnostic()?; + + let handle = ReposControl(state).get(&did); + let err = match handle.generate_car().await { + Ok(_) => panic!("redacted head must make getRepo fail"), + Err(err) => err, + }; + assert!(err.to_string().contains("record body was deleted")); + Ok(()) + } } diff --git a/src/db/keys/indexer.rs b/src/db/keys/indexer.rs --- a/src/db/keys/indexer.rs +++ b/src/db/keys/indexer.rs @@ -126,6 +126,12 @@ key } +/// recover the record-key prefix from a canonical death-keyed history entry. +pub fn history_record_key(key: &[u8]) -> Option<&[u8]> { + let record_end = key.len().checked_sub(9)?; + (key.get(record_end) == Some(&REC_SEP)).then_some(&key[..record_end]) +} + /// durable operator-redaction identity: `{record key} 00 {cid}`. /// /// this lives in its own keyspace. keeping the cid in the key makes replay -- tangled.sh