From 92a2c243f8ac48b56a247c30ca1276efa23dafdf Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Sat, 12 Sep 2026 21:04:37 +0300 Subject: [PATCH] [backfill] preserve repos on non-authoritative not-found --- src/backfill/error.rs | 2 + src/backfill/worker/process.rs | 26 +- src/backfill/worker/task.rs | 557 +++++++++++++++++++++++++++++++++ 3 files changed, 568 insertions(+), 17 deletions(-) diff --git a/src/backfill/error.rs b/src/backfill/error.rs index 5a3fd5b..8261244 100644 --- a/src/backfill/error.rs +++ b/src/backfill/error.rs @@ -21,6 +21,8 @@ pub enum BackfillError { Transport(SmolStr), #[error("repo was concurrently deleted")] Deleted, + #[error("getRepo returned RepoNotFound")] + RepoNotFound, } impl From for BackfillError { diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index 5b80499..95e1071 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -7,7 +7,7 @@ use fjall::Slice; use miette::{IntoDiagnostic, Result}; use reqwest::StatusCode; use smol_str::{SmolStr, ToSmolStr}; -use tracing::{debug, error, trace, warn}; +use tracing::{debug, trace, warn}; use jacquard_api::com_atproto::sync::get_repo::GetRepoError; use jacquard_common::IntoStatic; @@ -190,21 +190,12 @@ pub(super) async fn process_did( { FullRepoOutcome::Car(body) => body, FullRepoOutcome::NotFound => { - warn!("repo not found, deleting"); - let mut txn = DbTxn::new(db); - if txn.lock_repo_and_is_excluded(did)? { - return Ok(None); - } - let applied = - txn.transition_pending_key(did, pending_key.as_ref(), GaugeState::Synced)?; - txn.batch.remove(&db.indexer.pending, pending_key.clone()); - let _ = applied - .then(|| crate::ops::delete_repo(&mut txn, db, did, &state)) - .and_then(Result::err) - .inspect(|e| error!(err = %e, "failed to wipe repo during backfill")); - txn.commit()?; - // return None so did_task skips sending BackfillFinished (nothing to drain for a deleted repo) - return Ok(None); + // `RepoNotFound` is not lifecycle evidence: it says the pds has no repo, not that the + // account was deleted, and the repo may only be missing on this target. erasure stays + // reserved for authoritative firehose and listing `Deleted` events, so fail retryably + // and let the task error path schedule the existing jittered retry. + warn!("repo not found on pds, retrying"); + return Err(BackfillError::RepoNotFound); } FullRepoOutcome::Inactive(status) => { warn!(?status, "repo is inactive, stopping backfill"); @@ -515,7 +506,8 @@ const ERROR_BODY_MAX_BYTES: usize = 64 * 1024; enum FullRepoOutcome { /// the streamed CAR body, within the configured size ceiling. Car(Bytes), - /// the PDS reported that the repository does not exist (`RepoNotFound`). + /// the PDS has no repository for this DID (`RepoNotFound`). this is a fetch result, not + /// lifecycle evidence: callers must not delete the stored snapshot for it. NotFound, /// the PDS reported the repository is inactive; carries the status to record. Inactive(RepoStatus), diff --git a/src/backfill/worker/task.rs b/src/backfill/worker/task.rs index 67b87cd..3ec42ec 100644 --- a/src/backfill/worker/task.rs +++ b/src/backfill/worker/task.rs @@ -199,6 +199,9 @@ pub(super) async fn did_task( BackfillError::Generic(e) => { error!(err = %e, "failed"); } + BackfillError::RepoNotFound => { + warn!("repo not found, retrying through existing backoff"); + } BackfillError::Deleted => unreachable!("already handled"), BackfillError::PdsBusy => unreachable!("already handled"), } @@ -209,6 +212,7 @@ pub(super) async fn did_task( } BackfillError::Transport(_) => ResyncErrorKind::Transport, BackfillError::Generic(_) => ResyncErrorKind::Generic, + BackfillError::RepoNotFound => ResyncErrorKind::Generic, BackfillError::Deleted => unreachable!("already handled"), BackfillError::PdsBusy => unreachable!("already handled"), }; @@ -315,7 +319,26 @@ pub(super) async fn did_task( #[cfg(test)] mod tests { use super::*; + use crate::backfill::admission::BackfillAdmission; use crate::config::Config; + use crate::db; + use crate::db::types::DbRkey; + use crate::ingest::indexer::FirehoseWorker; + use crate::types::RepoMetadata; + use axum::Router; + use axum::routing::get; + use bytes::Bytes; + use jacquard_common::CowStr; + use jacquard_common::types::cid::IpldCid; + use jacquard_common::types::string::Handle; + use jacquard_common::types::tid::Tid; + use jacquard_repo::commit::Commit as AtpCommit; + use jacquard_repo::mst::Mst; + use jacquard_repo::{BlockStore, MemoryBlockStore}; + use reqwest::StatusCode; + use std::sync::atomic::{AtomicUsize, Ordering}; + use tokio::sync::{RwLock, Semaphore}; + use url::Url; #[tokio::test] async fn excluded_orphan_pending_entry_is_removed_before_work() -> miette::Result<()> { @@ -351,4 +374,538 @@ mod tests { ); Ok(()) } + + // --- full-backfill `RepoNotFound` boundary --------------------------------------------- + // + // `getRepo` `RepoNotFound` is a fetch result, not lifecycle evidence. these tests pin the + // non-destructive policy at the task boundary: the snapshot survives, the queue parks a + // retryable error, and only a later valid CAR or an authoritative deletion moves state. + + const DID: &str = "did:plc:ewvi7nxzyoun6zhxrhs64oiz"; + const NOT_FOUND_BODY: &[u8] = br#"{"error":"RepoNotFound","message":"not here"}"#; + const COLLECTION: &str = "app.bsky.feed.post"; + + fn post(text: &str) -> Vec { + serde_ipld_dagcbor::to_vec(&serde_json::json!({ + "$type": COLLECTION, + "text": text, + })) + .unwrap() + } + + struct PdsStub { + base: Url, + response: Arc>, + requests: Arc, + } + + impl PdsStub { + async fn spawn() -> Self { + let response = Arc::new(RwLock::new(( + StatusCode::BAD_REQUEST, + Bytes::from_static(NOT_FOUND_BODY), + ))); + let requests = Arc::new(AtomicUsize::new(0)); + let app = Router::new().route( + "/xrpc/com.atproto.sync.getRepo", + get({ + let response = response.clone(); + let requests = requests.clone(); + move || { + let response = response.clone(); + let requests = requests.clone(); + async move { + requests.fetch_add(1, Ordering::SeqCst); + let (status, body) = response.read().await.clone(); + (status, body) + } + } + }), + ); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let base = Url::parse(&format!("http://{}/", listener.local_addr().unwrap())).unwrap(); + tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + Self { + base, + response, + requests, + } + } + + async fn serve_car(&self, car: Bytes) { + *self.response.write().await = (StatusCode::OK, car); + } + + fn requests(&self) -> usize { + self.requests.load(Ordering::SeqCst) + } + } + + async fn spawn_doc_stub(doc: String) -> Url { + let app = Router::new().fallback({ + move |_uri: axum::http::Uri| { + let doc = doc.clone(); + async move { + ( + [(axum::http::header::CONTENT_TYPE, "application/json")], + doc, + ) + } + } + }); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); + let base = Url::parse(&format!("http://{}/", listener.local_addr().unwrap())).unwrap(); + tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); + base + } + + fn did_doc(base: &Url) -> String { + serde_json::json!({ + "@context": ["https://www.w3.org/ns/did/v1"], + "id": DID, + "service": [{ + "id": "#atproto_pds", + "type": "AtprotoPersonalDataServer", + "serviceEndpoint": base.as_str(), + }], + }) + .to_string() + } + + fn rooted(rev: &str) -> miette::Result> { + let mut state = RepoState::synced(); + state.root = Some(crate::types::Commit { + version: 3, + rev: crate::db::types::DbTid::from(&Tid::new(rev).into_diagnostic()?), + data: jacquard_repo::mst::util::compute_cid(b"seeded-root").into_diagnostic()?, + prev: None, + sig: Bytes::from_static(b"signature"), + }); + state.pds = Some(CowStr::Borrowed("https://old-pds.example")); + state.handle = Some(Handle::new_static("alice.test").into_diagnostic()?); + state.last_message_time = Some(101); + state.last_updated_at = 5; + Ok(state.into_static()) + } + + struct BuiltCar { + bytes: Bytes, + mst_root: IpldCid, + rev: Tid, + } + + async fn build_car(rev: &str, records: &[(&str, &[u8])]) -> miette::Result { + let store = Arc::new(MemoryBlockStore::new()); + let mut mst = Mst::new(store.clone()); + for &(rkey, body) in records { + let cid = jacquard_repo::mst::util::compute_cid(body).into_diagnostic()?; + mst.add_mut(&format!("{COLLECTION}/{rkey}"), cid) + .await + .into_diagnostic()?; + store + .put_many([(cid, Bytes::copy_from_slice(body))]) + .await + .into_diagnostic()?; + } + let mst_root = mst.persist().await.into_diagnostic()?; + let rev = Tid::new(rev).into_diagnostic()?; + let did: Did = Did::new_static(DID).into_diagnostic()?; + let commit = AtpCommit { + did, + version: 3, + data: mst_root, + rev: rev.clone(), + prev: None, + sig: Bytes::from_static(b"signature"), + }; + let commit_cid = commit.to_cid().into_diagnostic()?; + let commit_cbor = commit.to_cbor().into_diagnostic()?; + store + .put_many([(commit_cid, Bytes::from(commit_cbor.clone()))]) + .await + .into_diagnostic()?; + + let mut bytes = Vec::new(); + let mut writer = + iroh_car::CarWriter::new(iroh_car::CarHeader::new_v1(vec![commit_cid]), &mut bytes); + writer + .write(commit_cid, &commit_cbor) + .await + .into_diagnostic()?; + mst.write_blocks_to_car(&mut writer) + .await + .into_diagnostic()?; + writer.finish().await.into_diagnostic()?; + Ok(BuiltCar { + bytes: Bytes::from(bytes), + mst_root, + rev, + }) + } + + struct Fixture { + state: Arc, + did: Did, + pds: PdsStub, + _tmp: tempfile::TempDir, + } + + impl Fixture { + async fn new() -> miette::Result { + let tmp = tempfile::tempdir().into_diagnostic()?; + let pds = PdsStub::spawn().await; + let doc_base = spawn_doc_stub(did_doc(&pds.base)).await; + let config = Config { + database_path: tmp.path().to_path_buf(), + plc_urls: vec![doc_base], + identity_cache_size: 64, + ..Default::default() + }; + Ok(Self { + state: Arc::new(AppState::new(&config)?), + did: Did::new_static(DID).into_diagnostic()?, + pds, + _tmp: tmp, + }) + } + + fn seed( + &self, + repo: &RepoState<'static>, + index_id: u64, + records: &[(&str, &[u8])], + ) -> miette::Result<()> { + let mut batch = self.state.db.inner.batch(); + batch.insert( + &self.state.db.repos, + keys::repo_key(&self.did), + db::ser_repo_state(repo)?, + ); + batch.insert( + &self.state.db.repo_metadata, + keys::repo_metadata_key(&self.did), + db::ser_repo_meta(&RepoMetadata::backfilling(index_id))?, + ); + batch.insert( + &self.state.db.indexer.pending, + keys::pending_key(index_id), + keys::repo_key(&self.did), + ); + for &(rkey, body) in records { + self.state.db.indexer.stage_record( + &mut batch, + keys::record_key(&self.did, COLLECTION, &DbRkey::new(rkey)), + body, + ); + } + batch.commit().into_diagnostic() + } + + /// mirrors the retry worker: a completed queue entry gives way to a fresh pending key. + fn reindex_pending(&self, from: u64, to: u64) -> miette::Result<()> { + let mut batch = self.state.db.inner.batch(); + batch.remove(&self.state.db.indexer.pending, keys::pending_key(from)); + batch.insert( + &self.state.db.indexer.pending, + keys::pending_key(to), + keys::repo_key(&self.did), + ); + batch.insert( + &self.state.db.repo_metadata, + keys::repo_metadata_key(&self.did), + db::ser_repo_meta(&RepoMetadata::backfilling(to))?, + ); + batch.commit().into_diagnostic() + } + + fn clear_resync(&self) -> miette::Result<()> { + self.state + .db + .indexer + .resync + .remove(keys::repo_key(&self.did)) + .into_diagnostic()?; + Ok(()) + } + + /// the authoritative deletion primitive plus its `Deleted` tombstone: the same pair the + /// firehose and list-repos paths apply under the repository lock. + fn delete_authoritatively(&self, repo: &RepoState<'static>) -> miette::Result<()> { + let mut txn = DbTxn::new(&self.state.db); + crate::ops::delete_repo(&mut txn, &self.state.db, &self.did, repo)?; + crate::ops::transition_repo( + &mut txn.batch, + &self.state.db, + &mut Vec::new(), + &self.did, + repo.clone(), + RepoStatus::Deleted, + )?; + txn.commit() + } + + async fn run(&self, pending_key: [u8; 8]) -> Result { + let (tx, _worker) = FirehoseWorker::new(self.state.clone(), 1); + let http = ThrottledHttpClient::new( + vec![reqwest::Client::new()], + self.state.throttler.clone(), + ); + let admission = BackfillAdmission::new(4, 4); + let permit = Arc::new(Semaphore::new(1)).acquire_owned().await.unwrap(); + did_task( + &self.state, + http, + tx, + &self.did, + Slice::from(pending_key), + admission, + permit, + false, + BackfillStrategy::Full, + ) + .await + } + + fn repo(&self) -> miette::Result> { + let bytes = self + .state + .db + .repos + .get(keys::repo_key(&self.did)) + .into_diagnostic()? + .expect("repository state row") + .to_vec(); + Ok(db::deser_repo_state(&bytes)?.into_static()) + } + + fn resync(&self) -> miette::Result> { + self.state + .db + .indexer + .resync + .get(keys::repo_key(&self.did)) + .into_diagnostic()? + .map(|bytes| rmp_serde::from_slice::(&bytes).into_diagnostic()) + .transpose() + } + + fn pending_keys(&self) -> miette::Result>> { + self.state + .db + .indexer + .pending + .iter() + .map(|guard| Ok(guard.into_inner().into_diagnostic()?.0.to_vec())) + .collect() + } + + fn record_rkeys(&self) -> miette::Result> { + let prefix = keys::record_prefix_did(&self.did); + self.state + .db + .indexer + .record_prefix(&prefix) + .map(|guard| { + let (key, _) = guard.into_inner().into_diagnostic()?; + let (_, rkey) = keys::split_record_suffix(&key[prefix.len()..])?; + Ok(rkey.to_string()) + }) + .collect() + } + } + + #[tokio::test] + async fn repo_not_found_keeps_a_rooted_snapshot_and_parks_a_retry() -> miette::Result<()> { + let fixture = Fixture::new().await?; + let seeded = rooted("3jzfcijpj2z2a")?; + let body = post("a"); + fixture.seed(&seeded, 7, &[("a", body.as_slice())])?; + + let err = fixture.run(keys::pending_key(7)).await.unwrap_err(); + + assert!( + matches!(err, BackfillError::RepoNotFound), + "unexpected error: {err}" + ); + let after = fixture.repo()?; + assert_eq!( + after.root.as_ref().map(|root| root.data), + seeded.root.as_ref().map(|root| root.data), + "a non-authoritative fetch result must not clear the root" + ); + assert!( + matches!(&after.status, RepoStatus::Error(_)), + "expected a retryable error status, got {:?}", + after.status + ); + assert!(after.active); + assert_eq!(after.last_message_time, Some(101)); + assert_eq!(fixture.record_rkeys()?, ["a"]); + assert!( + fixture + .state + .db + .repo_metadata + .get(keys::repo_metadata_key(&fixture.did)) + .into_diagnostic()? + .is_some(), + "metadata must survive" + ); + assert!( + fixture.pending_keys()?.is_empty(), + "parked work must leave pending, or the worker spins" + ); + + let Some(ResyncState::Error { + kind, + retry_count, + next_retry, + }) = fixture.resync()? + else { + panic!("expected a durable retryable error"); + }; + assert_eq!(kind, ResyncErrorKind::Generic); + assert_eq!(retry_count, 1); + assert!( + next_retry > chrono::Utc::now().timestamp(), + "the retry worker may only requeue a due error" + ); + Ok(()) + } + + #[tokio::test] + async fn repo_not_found_on_a_never_rooted_repo_uses_the_same_retry_policy() -> miette::Result<()> + { + let fixture = Fixture::new().await?; + fixture.seed(&RepoState::backfilling().into_static(), 9, &[])?; + + let err = fixture.run(keys::pending_key(9)).await.unwrap_err(); + + assert!( + matches!(err, BackfillError::RepoNotFound), + "unexpected error: {err}" + ); + let after = fixture.repo()?; + assert!(after.root.is_none(), "no root may be invented"); + assert!( + matches!(&after.status, RepoStatus::Error(_)), + "expected a retryable error status, got {:?}", + after.status + ); + assert!(fixture.record_rkeys()?.is_empty()); + assert!(fixture.pending_keys()?.is_empty()); + assert!(matches!( + fixture.resync()?, + Some(ResyncState::Error { + kind: ResyncErrorKind::Generic, + .. + }) + )); + Ok(()) + } + + #[tokio::test] + async fn repo_not_found_cannot_overwrite_a_replaced_pending_entry() -> miette::Result<()> { + let fixture = Fixture::new().await?; + let seeded = rooted("3jzfcijpj2z2a")?; + let body = post("a"); + fixture.seed(&seeded, 7, &[("a", body.as_slice())])?; + fixture.reindex_pending(7, 11)?; + + let err = fixture.run(keys::pending_key(7)).await.unwrap_err(); + + assert!( + matches!(err, BackfillError::RepoNotFound), + "unexpected error: {err}" + ); + assert_eq!(fixture.pending_keys()?, [keys::pending_key(11).to_vec()]); + assert!( + fixture.resync()?.is_none(), + "stale work must not park a retry against a newer queue entry" + ); + let after = fixture.repo()?; + assert_eq!(after.status, RepoStatus::Synced); + assert_eq!( + after.root.as_ref().map(|root| root.data), + seeded.root.as_ref().map(|root| root.data) + ); + assert_eq!(fixture.record_rkeys()?, ["a"]); + Ok(()) + } + + #[tokio::test] + async fn authoritative_deletion_is_not_resurrected_by_a_late_repo_not_found() + -> miette::Result<()> { + let fixture = Fixture::new().await?; + let seeded = rooted("3jzfcijpj2z2a")?; + let body = post("a"); + fixture.seed(&seeded, 7, &[("a", body.as_slice())])?; + fixture.delete_authoritatively(&seeded)?; + + let err = fixture.run(keys::pending_key(7)).await.unwrap_err(); + + assert!( + matches!(err, BackfillError::RepoNotFound), + "unexpected error: {err}" + ); + let after = fixture.repo()?; + assert_eq!(after.status, RepoStatus::Deleted); + assert!(!after.active, "the tombstone must stay inactive"); + assert!( + fixture.record_rkeys()?.is_empty(), + "erased records must stay erased" + ); + assert!(fixture.pending_keys()?.is_empty()); + assert!( + fixture.resync()?.is_none(), + "a deleted repo must not be re-queued" + ); + Ok(()) + } + + #[tokio::test] + async fn a_later_valid_car_replaces_the_error_after_requeue() -> miette::Result<()> { + let fixture = Fixture::new().await?; + let seeded = rooted("3jzfcijpj2z2a")?; + let first = post("a"); + fixture.seed(&seeded, 7, &[("a", first.as_slice())])?; + assert!(matches!( + fixture.run(keys::pending_key(7)).await.unwrap_err(), + BackfillError::RepoNotFound + )); + assert_eq!(fixture.pds.requests(), 1); + + fixture.clear_resync()?; + fixture.reindex_pending(7, 21)?; + let second = post("b"); + let car = build_car( + "3jzfcijpj2z2b", + &[("a", first.as_slice()), ("b", second.as_slice())], + ) + .await?; + fixture.pds.serve_car(car.bytes.clone()).await; + + assert_eq!( + fixture.run(keys::pending_key(21)).await.unwrap(), + TaskDisposition::Finished + ); + + let after = fixture.repo()?; + assert_eq!( + after.root.as_ref().map(|root| root.data), + Some(car.mst_root) + ); + assert_eq!( + after.root.as_ref().map(|root| root.rev.to_tid()), + Some(car.rev.clone()) + ); + assert_eq!(fixture.record_rkeys()?, ["a", "b"]); + assert!( + fixture.resync()?.is_none(), + "a successful import clears the error row" + ); + assert!(fixture.pending_keys()?.is_empty()); + assert_eq!(fixture.pds.requests(), 2); + Ok(()) + } } -- 2.51.2