use futures::{FutureExt, TryFutureExt}; use rand::Rng; use super::*; impl ReposControl { /// iterates through pending repositories, returning their state. #[allow(dead_code)] pub(crate) fn iter_pending( &self, cursor: Option, ) -> impl Iterator> { let start_bound = if let Some(cursor) = cursor { std::ops::Bound::Excluded(cursor.to_be_bytes().to_vec()) } else { std::ops::Bound::Unbounded }; let repos = self.0.db.repos.clone(); let state = self.0.clone(); self.0 .db .indexer .pending .range((start_bound, std::ops::Bound::Unbounded)) .map(move |g| { let (id_raw, did_key) = g.into_inner().into_diagnostic()?; let id = u64::from_be_bytes( id_raw .as_ref() .try_into() .into_diagnostic() .wrap_err("can't parse pending key")?, ); let Some(bytes) = repos.get(&did_key).into_diagnostic()? else { // stale pending that we forgot to delete? shouldn't happen though tracing::warn!(id, did = ?did_key, "stale pending???"); return Ok(None); }; let repo_state = crate::db::deser_repo_state(bytes.as_ref())?; let did = TrimmedDid::try_from(did_key.as_ref())?.to_did(); let metadata_key = keys::repo_metadata_key(&did); let metadata = state .db .repo_metadata .get(&metadata_key) .into_diagnostic()? .ok_or_else(|| miette::miette!("repo metadata not found for {}", did))?; let metadata = crate::db::deser_repo_meta(metadata.as_ref())?; Ok(Some(( id, repo_state_to_info(did, repo_state.into_static(), metadata.tracked), ))) }) .filter_map(|b| b.transpose()) } #[allow(dead_code)] pub(crate) fn iter_resync( &self, cursor: Option<&Did>, ) -> impl Iterator> { let start_bound = if let Some(cursor) = cursor { let did_key = keys::repo_key(cursor); std::ops::Bound::Excluded(did_key) } else { std::ops::Bound::Unbounded }; let repos = self.0.db.repos.clone(); let state = self.0.clone(); self.0 .db .indexer .resync .range((start_bound, std::ops::Bound::Unbounded)) .map(move |g| { let did_key = g.key().into_diagnostic()?; let Some(bytes) = repos.get(&did_key).into_diagnostic()? else { // stale pending that we forgot to delete? shouldn't happen though tracing::warn!(did = ?did_key, "stale resync???"); return Ok(None); }; let repo_state = crate::db::deser_repo_state(bytes.as_ref())?; let did = TrimmedDid::try_from(did_key.as_ref())?.to_did(); let metadata_key = keys::repo_metadata_key(&did); let metadata = state .db .repo_metadata .get(&metadata_key) .into_diagnostic()? .ok_or_else(|| miette::miette!("repo metadata not found for {}", did))?; let metadata = crate::db::deser_repo_meta(metadata.as_ref())?; Ok(Some(repo_state_to_info( did, repo_state.into_static(), metadata.tracked, ))) }) .filter_map(|b| b.transpose()) } pub(crate) fn _resync(db: &Db, did: &Did, txn: &mut crate::db::Txn<'_>) -> Result { if txn.lock_repo_and_is_excluded(did)? { return Ok(false); } let did_key = keys::repo_key(did); let metadata_key = keys::repo_metadata_key(did); let repo_bytes = db.repos.get(&did_key).into_diagnostic()?; let existing = repo_bytes .as_deref() .map(db::deser_repo_state) .transpose()?; if existing.is_some() { let metadata_bytes = db .repo_metadata .get(&metadata_key) .into_diagnostic()? .ok_or_else(|| miette::miette!("repo metadata not found for {}", did))?; let mut metadata = crate::db::deser_repo_meta(&metadata_bytes)?; // skip if already in pending queue let is_pending = db .indexer .pending .get(keys::pending_key(metadata.index_id)) .into_diagnostic()? .is_some(); if !is_pending { metadata.tracked = true; // insert into pending with new index_id let old_pending = keys::pending_key(metadata.index_id); txn.batch.remove(&db.indexer.pending, old_pending); metadata.index_id = rand::Rng::next_u64(&mut rand::rng()); txn.batch.insert( &db.indexer.pending, keys::pending_key(metadata.index_id), &did_key, ); txn.batch.remove(&db.indexer.resync, &did_key); txn.batch.insert( &db.repo_metadata, &metadata_key, crate::db::ser_repo_meta(&metadata)?, ); txn.transition_lifecycle(did, GaugeState::Pending)?; return Ok(true); } } Ok(false) } /// request one or more repositories to be resynced. /// /// note that they may not immediately start backfilling if: /// - other repos already filled the backfill concurrency limit, /// - or there are many repos pending already. /// /// this will also clear any error state the repo may have been in, /// allowing it to resync again. pub async fn resync(&self, dids: impl IntoIterator) -> Result> { let dids: Vec = dids.into_iter().collect(); let queued = self .0 .db .run(move |db| { let mut txn = crate::db::Txn::new(db); txn.hold_record_lock_indexes_sorted( dids.iter() .map(|did| crate::db::record_lock_index_for_did(did)), ); let mut queued: Vec = Vec::new(); for did in dids { if Self::_resync(db, &did, &mut txn)? { queued.push(did); } } txn.commit()?; db.persist()?; Ok(queued) }) .await?; if !queued.is_empty() { self.0.notify_backfill(); } Ok(queued) } /// explicitly track one or more repositories, enqueuing them for backfill if needed. /// /// - if a repo is new, a fresh [`RepoState`] is created and backfill is queued. /// - if a repo is already known but untracked, it is marked tracked and re-enqueued. /// - if a repo is already tracked, this is a no-op. pub async fn track(&self, dids: impl IntoIterator) -> Result> { let dids: Vec = dids.into_iter().collect(); let queued = self .0 .db .run(move |db| { let mut txn = crate::db::Txn::new(db); txn.hold_record_lock_indexes_sorted( dids.iter() .map(|did| crate::db::record_lock_index_for_did(did)), ); let mut queued: Vec = Vec::new(); for did in dids { if db .filter .contains_key(crate::db::filter::exclude_key(did.as_str())?) .into_diagnostic()? { continue; } let did_key = keys::repo_key(&did); let metadata_key = keys::repo_metadata_key(&did); let metadata_bytes = db.repo_metadata.get(&metadata_key).into_diagnostic()?; let existing_metadata = metadata_bytes .map(|b| crate::db::deser_repo_meta(&b)) .transpose()?; if let Some(metadata) = existing_metadata { if !metadata.tracked && Self::_resync(db, &did, &mut txn)? { queued.push(did); } } else { let repo_state = RepoState::backfilling(); let metadata = RepoMetadata::backfilling(rand::random()); txn.batch.insert( &db.repos, &did_key, crate::db::ser_repo_state(&repo_state)?, ); txn.batch.insert( &db.repo_metadata, &metadata_key, crate::db::ser_repo_meta(&metadata)?, ); txn.batch.insert( &db.indexer.pending, keys::pending_key(metadata.index_id), &did_key, ); txn.counts.add_repos(1); txn.transition_lifecycle(&did, GaugeState::Pending)?; queued.push(did); } } txn.commit()?; db.persist()?; Ok(queued) }) .await?; self.0.notify_backfill(); Ok(queued) } /// stop tracking one or more repositories. hydrant will stop processing new events /// for them and remove them from the pending/resync queues, but existing indexed /// records are **not** deleted. pub async fn untrack(&self, dids: impl IntoIterator) -> Result> { let dids: Vec = dids.into_iter().collect(); let untracked = self .0 .db .run(move |db| { let mut txn = crate::db::Txn::new(db); let mut untracked: Vec = Vec::new(); for did in dids { let did_key = keys::repo_key(&did); let metadata_key = keys::repo_metadata_key(&did); let repo_bytes = db.repos.get(&did_key).into_diagnostic()?; let existing = repo_bytes .as_deref() .map(db::deser_repo_state) .transpose()?; if existing.is_some() { let metadata_bytes = db.repo_metadata.get(&metadata_key).into_diagnostic()?; let existing_metadata = metadata_bytes .map(|b| crate::db::deser_repo_meta(&b)) .transpose()?; if let Some(mut metadata) = existing_metadata && metadata.tracked { metadata.tracked = false; txn.batch.insert( &db.repo_metadata, &metadata_key, crate::db::ser_repo_meta(&metadata)?, ); txn.batch .remove(&db.indexer.pending, keys::pending_key(metadata.index_id)); txn.batch.remove(&db.indexer.resync, &did_key); txn.transition_lifecycle(&did, GaugeState::Synced)?; untracked.push(did); } } } txn.commit()?; db.persist()?; Ok(untracked) }) .await?; 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, target: crate::types::DeleteBodyTarget, exclude: bool, ) -> Result { self.0.ensure_operator_body_deletion_supported()?; let did = did.clone(); 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); if exclude { 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, target) }) .await } } impl RepoHandle { /// 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(); 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() { miette::bail!("block storage is not available in this mode"); } let did = self.did.clone(); let db_key = keys::record_key(&did, collection, &DbRkey::new(rkey)); self.state .db .run(move |db| { use miette::WrapErr; let Some(body) = db.indexer.record(db_key).into_diagnostic()? else { return Ok(None); }; if crate::db::is_cid_record_value(&body) { return Ok(None); } let value = serde_ipld_dagcbor::from_slice::(&body) .into_diagnostic() .wrap_err("cant parse record body")?; // the cid is no longer stored; recompute it from the body let cid = jacquard_repo::mst::util::compute_cid(body.as_ref()) .into_diagnostic() .wrap_err("cant compute record cid")?; let cid = Cid::ipld(cid); Ok(Some(Record { did, cid, value })) }) .await } /// lists records from this repository. pub async fn list_records( &self, collection: &str, limit: usize, reverse: bool, cursor: Option<&str>, ) -> Result { if !self.state.stores_record_bodies() { miette::bail!("block storage is not available in this mode"); } let did = self.did.clone().into_static(); let prefix = keys::record_prefix_collection(&did, collection); let cursor = cursor.map(|c| c.to_smolstr()); self.state .db .run(move |db| { let mut results = Vec::new(); let mut next_cursor = None; let mut last_returned = None; let iter: Box> = if !reverse { // newest-first: walk backwards from the cursor, or from the // end of the collection. the exclusive upper bound is the // cursor key itself, or `prefix ff`: rkeys cannot contain // 0xff, so every key in the collection sorts below it let end_key = if let Some(cursor) = &cursor { let mut k = prefix.clone(); k.extend_from_slice(DbRkey::new(cursor).to_smolstr().as_bytes()); k } else { let mut k = prefix.clone(); k.push(0xff); k }; Box::new( db.indexer .record_range(prefix.as_slice()..end_key.as_slice()) .rev(), ) } else { let start_key = if let Some(cursor) = &cursor { // strictly after the cursor rkey: its key plus the next // possible byte let mut k = prefix.clone(); k.extend_from_slice(DbRkey::new(cursor).to_smolstr().as_bytes()); k.push(0); k } else { prefix.clone() }; Box::new(db.indexer.record_range(start_key.as_slice()..)) }; for item in iter { let (key, body) = item.into_inner().into_diagnostic()?; if !key.starts_with(prefix.as_slice()) { break; } if crate::db::is_cid_record_value(&body) { continue; } let rkey = keys::parse_rkey_text(&key[prefix.len()..])?; if results.len() >= limit { next_cursor = last_returned; break; } let value: Data = serde_ipld_dagcbor::from_slice(&body).unwrap_or(Data::Null); // the cid is no longer stored; recompute it from the body let cid = jacquard_repo::mst::util::compute_cid(body.as_ref()).into_diagnostic()?; let cid = Cid::ipld(cid); results.push(ListedRecord { rkey: Rkey::new(rkey.to_smolstr()).expect("that rkey is validated"), cid, value, }); last_returned = Some(rkey); } Ok((results, next_cursor)) }) .await .map(|(records, next_cursor)| RecordList { records, cursor: next_cursor .map(|rkey| Rkey::new(rkey.to_smolstr()).expect("that rkey is validated")), }) } /// generates a streaming CAR v1 response body for this repository. /// /// returns `None` if the repo has no commit yet (still backfilling) or is an /// unmigrated repo that does not have the necessary data to reconstruct the /// root commit from. /// /// ## notes /// - calling this if you are using collection allowlist will always result /// in an error since the commit root won't match the reconstructed CID. /// - calling this for big repositories will incur more resource cost due to /// hydrant's structure, the whole MST is always reconstructed. pub async fn generate_car( &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"); } use iroh_car::{CarHeader, CarWriter}; use jacquard_repo::{BlockStore, MemoryBlockStore, Mst}; use miette::WrapErr; use std::sync::Arc; let commit = match self.state().await? { Some(state) => match state.root { Some(c) => c, None => return Ok(None), }, None => return Ok(None), }; let atp_commit = match commit.into_atp_commit(self.did.clone().into_static()) { Some(c) => c, None => return Ok(None), }; let commit_cid = atp_commit.to_cid().into_diagnostic()?; let commit_cbor = atp_commit.to_cbor().into_diagnostic()?; let did = self.did.clone().into_static(); let app_state = self.state.clone(); // build mst and populate the block store in a single blocking pass let store = Arc::new(MemoryBlockStore::new()); let mst = Mst::new(store.clone()); let handle = tokio::runtime::Handle::current(); let mst = tokio::task::spawn_blocking(move || -> Result<_> { let mut mst = mst; let prefix = keys::record_prefix_did(&did); for guard in app_state.db.indexer.record_prefix(&prefix) { let (key, body) = guard.into_inner().into_diagnostic()?; 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() .wrap_err_with(|| format!("cant compute cid for record {mst_key}"))?; handle .block_on(mst.add_mut(&mst_key, ipld_cid)) .into_diagnostic()?; // we use put_many here to skip calculating the CID again handle .block_on( mst.storage() .put_many([(ipld_cid, bytes::Bytes::copy_from_slice(body.as_ref()))]), ) .into_diagnostic()?; } handle.block_on(mst.persist()).into_diagnostic()?; Result::<_>::Ok(mst) }) .await .into_diagnostic()??; // 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 { miette::bail!( "mst root mismatch for {}: computed {computed_root} != stored {}", self.did, atp_commit.data, ); } store .put_many([(commit_cid, bytes::Bytes::from(commit_cbor))]) .await .into_diagnostic()?; // stream the car directly to the response let (reader, writer) = tokio::io::duplex(64 * 1024); tokio::spawn( async move { let header = CarHeader::new_v1(vec![commit_cid]); let mut car_writer = CarWriter::new(header, writer); // write commit first, then mst nodes + leaf blocks let commit_data = store.get(&commit_cid).await?; if let Some(data) = commit_data { car_writer .write(commit_cid, &data) .await .into_diagnostic()?; } mst.write_blocks_to_car(&mut car_writer).await?; car_writer.finish().await.into_diagnostic()?; Result::<_, miette::Report>::Ok(()) } .inspect_err(|e| tracing::error!("can't generate car: {e}")), ); Ok(Some(tokio_util::io::ReaderStream::new(reader))) } /// gets how many records of a collection this repository has. pub async fn count_records(&self, collection: &str) -> Result { let did = self.did.clone().into_static(); let collection = collection.to_string(); self.state .db .run(move |db| db::get_record_count(db, &did, &collection)) .await } } #[cfg(test)] mod tests { use super::*; fn record_body(text: &str) -> Vec { serde_ipld_dagcbor::to_vec(&serde_json::json!({ "$type": "app.bsky.feed.post", "text": text, })) .unwrap() } #[tokio::test] async fn get_and_list_hide_redacted_heads() -> 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_static("did:plc:testredactedreads").into_diagnostic()?; let visible = record_body("visible"); let hidden_cid = jacquard_repo::mst::util::compute_cid(&record_body("hidden")).into_diagnostic()?; let mut batch = state.db.inner.batch(); for rkey in ["a", "c"] { state.db.indexer.stage_record( &mut batch, keys::record_key(&did, "app.bsky.feed.post", &DbRkey::new(rkey)), &visible, ); } state.db.indexer.stage_record( &mut batch, keys::record_key(&did, "app.bsky.feed.post", &DbRkey::new("b")), hidden_cid.to_bytes(), ); batch.commit().into_diagnostic()?; let handle = ReposControl(state).get(&did); let visible_record = handle .get_record("app.bsky.feed.post", "a") .await? .expect("visible record"); let visible_json = serde_json::to_value(&visible_record.value).into_diagnostic()?; assert_eq!(visible_json["$type"], "app.bsky.feed.post"); assert_eq!(visible_json["text"], "visible"); assert!( handle .get_record("app.bsky.feed.post", "b") .await? .is_none() ); let listed = handle .list_records("app.bsky.feed.post", 10, true, None) .await?; assert_eq!( listed .records .iter() .map(|record| record.rkey.as_str()) .collect::>(), ["a", "c"] ); let first = handle .list_records("app.bsky.feed.post", 1, true, None) .await?; assert_eq!(first.cursor.as_ref().map(|rkey| rkey.as_str()), Some("a")); let second = handle .list_records( "app.bsky.feed.post", 1, true, first.cursor.as_ref().map(|rkey| rkey.as_str()), ) .await?; assert_eq!(second.records[0].rkey.as_str(), "c"); 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_static("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, crate::types::DeleteBodyTarget::All, true) .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_static("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_static("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_static("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(()) } }