From 750747985dbcd616817deaf22b09dced4893a83e Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Mon, 5 Oct 2026 13:45:54 +0300 Subject: [PATCH] [backfill] start a sparse scan from the latest commit when the probe has no proof pdses that can't prove a record is absent answer the sparse probe with RecordNotFound or a bare 404, which sent every repo on them to a full getRepo. those now start from getLatestCommit and getBlocks(commit): the commit is trusted through its signature like a proof's, so the scan finds the same records with a couple more requests. the mst root is fetched up front so auto mode can still pick full for tiny trees. both the sparse and full paths now reject a commit that names another did. --- src/backfill/sparse.rs | 487 ++++++++++++++++++++++++++++++--- src/backfill/worker/process.rs | 3 +- src/backfill/worker/task.rs | 26 +- 3 files changed, 471 insertions(+), 45 deletions(-) diff --git a/src/backfill/sparse.rs b/src/backfill/sparse.rs index 08e4f53..01f8eb5 100644 --- a/src/backfill/sparse.rs +++ b/src/backfill/sparse.rs @@ -1,6 +1,6 @@ use crate::backfill::client::ThrottledHttpClient; use crate::backfill::error::BackfillError; -use crate::car::{Blocks, CarBlock, CommitBlocks}; +use crate::car::{Blocks, CarBlock, CommitBlocks, CommitCar}; use crate::config::{BackfillStrategy, RateTier}; use crate::db::types::{DbAction, DbRkey}; use crate::db::{self, Txn as DbTxn, keys}; @@ -13,12 +13,14 @@ use crate::types::{Commit, RepoState}; use futures::{StreamExt, stream}; use jacquard_api::com_atproto::sync::get_blocks::GetBlocksError; +use jacquard_api::com_atproto::sync::get_latest_commit::{GetLatestCommit, GetLatestCommitError}; use jacquard_api::com_atproto::sync::get_record::{GetRecord, GetRecordError}; +use jacquard_common::error::{ClientError, ClientErrorKind}; use jacquard_common::http_client::HttpClient; use jacquard_common::types::cid::{Cid as AtCid, IpldCid}; use jacquard_common::types::did::Did; use jacquard_common::types::string::{Nsid, RecordKey}; -use jacquard_common::xrpc::{XrpcError, XrpcExt}; +use jacquard_common::xrpc::{Response, XrpcError, XrpcExt, XrpcRequest}; use miette::IntoDiagnostic; use reqwest::StatusCode; @@ -103,51 +105,48 @@ pub(crate) async fn process_did_sparse( let tier = app_state.resolve_pds_tier(pds.host_str().unwrap_or("")); - let resp = { - let _permit = throttle.acquire().await; - if throttle.is_throttled() { - return Err(BackfillError::PreemptivelyThrottled); - } - throttle.wait_for_allow(1, &tier).await; - if throttle.is_throttled() { - return Err(BackfillError::PreemptivelyThrottled); - } - match http - .xrpc(url_to_fluent_uri_ref(pds_endpoint.url())) - .send(&req) - .await - { - Ok(resp) => { - if !throttle.is_throttled() { - throttle.record_success(); - } - resp - } - Err(e) => return Err(BackfillError::from_sparse_client(e, &throttle)), - } + let start = match send_throttled(http, &pds_endpoint, &throttle, &tier, &req).await? { + Ok(resp) => match resp.into_output() { + Ok(proof) => Some(crate::car::parse_commit_car( + proof.body, + app_state.verify_cids, + )?), + Err(XrpcError::Xrpc(GetRecordError::RecordNotFound(_))) => None, + Err(XrpcError::Xrpc( + GetRecordError::RepoNotFound(_) + | GetRecordError::RepoTakendown(_) + | GetRecordError::RepoSuspended(_) + | GetRecordError::RepoDeactivated(_), + )) => return Ok(SparseBackfillResult::Skipped), + Err(e) => Err(e).into_diagnostic()?, + }, + // the client only reads an xrpc error body on a 400, so a pds answering + // RecordNotFound with a 404 lands here too + Err(e) if is_http_status(&e, StatusCode::NOT_FOUND) => None, + Err(e) => return Err(BackfillError::from_sparse_client(e, &throttle)), }; - let proof = match resp.into_output() { - Ok(o) => o, - Err(XrpcError::Xrpc(GetRecordError::RecordNotFound(_))) => { - return Ok(SparseBackfillResult::Skipped); + let start = match start { + Some(proof) => proof, + None => { + debug!("sparse probe has no proof, starting from the latest commit"); + let latest_commit = + latest_commit_start(app_state, http, did, pds, &pds_endpoint, &throttle, &tier) + .await?; + let Some(start) = latest_commit else { + return Ok(SparseBackfillResult::Skipped); + }; + start } - Err(XrpcError::Xrpc( - GetRecordError::RepoNotFound(_) - | GetRecordError::RepoTakendown(_) - | GetRecordError::RepoSuspended(_) - | GetRecordError::RepoDeactivated(_), - )) => return Ok(SparseBackfillResult::Skipped), - Err(e) => Err(e).into_diagnostic()?, }; - let parsed = crate::car::parse_commit_car(proof.body, app_state.verify_cids)?; - let root_bytes = parsed + let root_bytes = start .blocks - .get(&parsed.root) + .get(&start.root) .ok_or_else(|| miette::miette!("root block missing from sparse proof CAR"))?; let root_commit = jacquard_repo::commit::Commit::::from_cbor(root_bytes) .into_diagnostic()?; + check_commit_did(root_commit.did.as_str(), did)?; if verify_signatures { let pubkey = app_state.resolver.resolve_signing_key(did).await?; @@ -157,9 +156,30 @@ pub(crate) async fn process_did_sparse( } let root_cid = root_commit.data; + let mut blocks = start.blocks; + // a start from the latest commit has no mst nodes yet. fetching the root now is the + // scan's first step anyway, and lets auto mode see how deep the tree is + if blocks.get(&root_cid).is_none() { + blocks.extend( + fetch_blocks( + http, + pds, + did, + &[root_cid], + &throttle, + &tier, + app_state.verify_cids, + app_state.max_car_body_bytes, + ) + .await?, + ); + if blocks.get(&root_cid).is_none() { + debug!(%root_cid, "pds did not serve the mst root"); + return Ok(SparseBackfillResult::Skipped); + } + } let auto_full_layer = if strategy == BackfillStrategy::Auto { - parsed - .blocks + blocks .get(&root_cid) .map(|root_bytes| mst_node_layer(root_bytes)) .transpose()? @@ -170,7 +190,7 @@ pub(crate) async fn process_did_sparse( }; let root_commit = Commit::from(root_commit); - let mut scanner = SparseScanner::new(ranges, parsed.blocks); + let mut scanner = SparseScanner::new(ranges, blocks); let mut scan_rounds = 0; let scan = loop { match scanner.scan(root_cid)? { @@ -282,6 +302,97 @@ pub(crate) async fn process_did_sparse( })) } +/// sends one sparse request once the pds's throttle allows it. the client's own result comes +/// back as is, so a caller can read a status that only means "no" before calling it a failure +async fn send_throttled( + http: &ThrottledHttpClient, + pds: &PublicEndpoint, + throttle: &ThrottleHandle, + tier: &RateTier, + req: &R, +) -> Result, ClientError>, BackfillError> +where + R: XrpcRequest + serde::Serialize, + R::Response: Send + Sync, +{ + let _permit = throttle.acquire().await; + if throttle.is_throttled() { + return Err(BackfillError::PreemptivelyThrottled); + } + throttle.wait_for_allow(1, tier).await; + if throttle.is_throttled() { + return Err(BackfillError::PreemptivelyThrottled); + } + let resp = http.xrpc(url_to_fluent_uri_ref(pds.url())).send(req).await; + if resp.is_ok() && !throttle.is_throttled() { + throttle.record_success(); + } + Ok(resp) +} + +fn is_http_status(err: &ClientError, wanted: StatusCode) -> bool { + matches!(err.kind(), ClientErrorKind::Http { status } if *status == wanted) +} + +/// the repo's latest commit block on its own, for a pds that can't prove a record is absent. +/// the commit is trusted through its signature exactly like a proof's, so starting from it +/// costs a couple more requests but finds the same records. `None` sends the repo to a full +/// getRepo, which also records an inactive repo's status. +async fn latest_commit_start( + app_state: &AppState, + http: &ThrottledHttpClient, + did: &Did, + pds: &url::Url, + pds_endpoint: &PublicEndpoint, + throttle: &ThrottleHandle, + tier: &RateTier, +) -> Result, BackfillError> { + let req = GetLatestCommit::new().did(did.clone()).build(); + let latest = match send_throttled(http, pds_endpoint, throttle, tier, &req).await? { + Ok(resp) => match resp.into_output() { + Ok(latest) => latest, + Err(XrpcError::Xrpc( + GetLatestCommitError::RepoNotFound(_) + | GetLatestCommitError::RepoTakendown(_) + | GetLatestCommitError::RepoSuspended(_) + | GetLatestCommitError::RepoDeactivated(_), + )) => return Ok(None), + Err(e) => Err(e).into_diagnostic()?, + }, + Err(e) if is_http_status(&e, StatusCode::NOT_FOUND) => return Ok(None), + Err(e) => return Err(BackfillError::from_sparse_client(e, throttle)), + }; + let commit = latest.cid.to_ipld().into_diagnostic()?; + let blocks = fetch_blocks( + http, + pds, + did, + &[commit], + throttle, + tier, + app_state.verify_cids, + app_state.max_car_body_bytes, + ) + .await?; + if blocks.get(&commit).is_none() { + debug!(%commit, "pds did not serve its latest commit"); + return Ok(None); + } + Ok(Some(CommitCar { + root: commit, + blocks, + })) +} + +/// a commit signed by this did's key could still name another did if two dids share a key, +/// and its records would then be stored under the wrong repo +pub(crate) fn check_commit_did(commit_did: &str, did: &Did) -> Result<(), BackfillError> { + if commit_did != did.as_str() { + return Err(miette::miette!("commit is for {commit_did}, not {did}").into()); + } + Ok(()) +} + async fn fetch_blocks( http: &ThrottledHttpClient, pds: &url::Url, @@ -571,6 +682,7 @@ mod tests { use cid::multihash::Multihash; use jacquard_common::types::crypto::{DAG_CBOR, SHA2_256}; use jacquard_repo::mst::util::compute_cid; + use std::collections::BTreeMap; use std::sync::atomic::{AtomicUsize, Ordering}; use std::time::Duration; use tokio::sync::Barrier; @@ -723,6 +835,10 @@ mod tests { }), ) }); + serve(app).await + } + + async fn serve(app: Router) -> std::net::SocketAddr { let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await.unwrap(); let address = listener.local_addr().unwrap(); tokio::spawn(async move { axum::serve(listener, app).await.unwrap() }); @@ -832,6 +948,18 @@ mod tests { ) -> ( Result, Option, + ) { + let address = spawn_car_server(routes).await; + backfill_at(address, verify_cids, BackfillStrategy::SparseFilter).await + } + + async fn backfill_at( + address: std::net::SocketAddr, + verify_cids: bool, + strategy: BackfillStrategy, + ) -> ( + Result, + Option, ) { let tmp = tempfile::tempdir().unwrap(); let config = crate::config::Config { @@ -850,7 +978,6 @@ mod tests { ); batch.commit().unwrap(); - let address = spawn_car_server(routes).await; let (pds, http) = public_client(address, &state.throttler); let result = process_did_sparse( &state, @@ -859,7 +986,7 @@ mod tests { &pds, RepoState::backfilling(), false, - BackfillStrategy::SparseFilter, + strategy, ) .await; let stored = state @@ -874,6 +1001,280 @@ mod tests { (result, stored) } + const LATEST_REV: &str = "3l2abcdefgh22"; + const OTHER_DID: &str = "did:plc:bbbbbbbbbbbbbbbbbbbbbbbb"; + + /// one record under `PROOF_COLLECTION/PROOF_RKEY`, its mst root and a commit naming + /// `commit_did`, by cid + struct TestRepo { + commit: IpldCid, + blocks: BTreeMap>, + } + + fn test_repo(commit_did: &'static str, record: &[u8]) -> TestRepo { + let key = format!("{PROOF_COLLECTION}/{PROOF_RKEY}"); + let record_cid = compute_cid(record).unwrap(); + let root = serde_ipld_dagcbor::to_vec(&node(vec![(key.as_str(), record_cid, None)], None)) + .unwrap(); + let root_cid = compute_cid(&root).unwrap(); + let commit = jacquard_repo::commit::Commit:: { + did: Did::new_static(commit_did).unwrap(), + version: 3, + data: root_cid, + rev: LATEST_REV.parse().unwrap(), + prev: None, + sig: bytes::Bytes::new(), + } + .to_cbor() + .unwrap(); + let commit_cid = compute_cid(&commit).unwrap(); + TestRepo { + commit: commit_cid, + blocks: BTreeMap::from([ + (commit_cid, commit.to_vec()), + (root_cid, root), + (record_cid, record.to_vec()), + ]), + } + } + + type Reply = (StatusCode, &'static str, String); + + fn json_error(status: StatusCode, error: &str) -> Reply { + let body = serde_json::json!({ "error": error, "message": "nope" }).to_string(); + (status, "application/json", body) + } + + fn reply_route(reply: Reply) -> axum::routing::MethodRouter { + get(move || { + let (status, content_type, body) = reply.clone(); + async move { + ( + status, + [(reqwest::header::CONTENT_TYPE.as_str(), content_type)], + body, + ) + } + }) + } + + /// a pds without exclusion proofs: the getRecord probe gets `probe`, getLatestCommit gets + /// `latest` (pointing at `repo`'s commit when `None`), and getBlocks serves only the blocks + /// asked for, from `blocks` + fn no_proof_pds( + repo: &TestRepo, + probe: Reply, + latest: Option, + blocks: Option>>, + ) -> Router { + let latest = latest.unwrap_or_else(|| { + let body = serde_json::json!({ "cid": repo.commit.to_string(), "rev": LATEST_REV }); + (StatusCode::OK, "application/json", body.to_string()) + }); + let app = Router::new() + .route(&format!("/xrpc/{GET_RECORD}"), reply_route(probe)) + .route( + "/xrpc/com.atproto.sync.getLatestCommit", + reply_route(latest), + ); + let Some(blocks) = blocks else { + return app; + }; + let blocks = Arc::new(blocks); + app.route( + &format!("/xrpc/{GET_BLOCKS}"), + get( + move |axum::extract::RawQuery(query): axum::extract::RawQuery| { + let blocks = blocks.clone(); + async move { + let found = + url::form_urlencoded::parse(query.unwrap_or_default().as_bytes()) + .filter(|(name, _)| name == "cids") + .filter_map(|(_, cid)| cid.parse::().ok()) + .filter_map(|cid| { + blocks.get(&cid).map(|bytes| (cid, bytes.clone())) + }) + .collect::>(); + if found.is_empty() { + let (status, content_type, body) = + json_error(StatusCode::BAD_REQUEST, "BlockNotFound"); + let headers = [(reqwest::header::CONTENT_TYPE.as_str(), content_type)]; + return (status, headers, body.into_bytes()); + } + let mut car = Vec::new(); + let header = iroh_car::CarHeader::new_v1(Vec::new()); + let mut writer = iroh_car::CarWriter::new(header, &mut car); + for (cid, bytes) in found { + writer.write(cid, bytes).await.unwrap(); + } + writer.finish().await.unwrap(); + let headers = [( + reqwest::header::CONTENT_TYPE.as_str(), + "application/vnd.ipld.car", + )]; + (StatusCode::OK, headers, car) + } + }, + ), + ) + } + + async fn backfill_no_proof( + app: Router, + strategy: BackfillStrategy, + ) -> ( + Result, + Option, + ) { + backfill_at(serve(app).await, true, strategy).await + } + + #[tokio::test] + async fn probe_without_proof_starts_from_the_latest_commit() { + let honest = tangled_repo("honest"); + let repo = test_repo(PROOF_DID, &honest); + // a pds that can't prove absence (an atcr hold's text 404), tranquil's 404, and the + // 400 an xrpc client expects + let probes = [ + ( + StatusCode::NOT_FOUND, + "text/plain", + "mst: not found".to_string(), + ), + json_error(StatusCode::NOT_FOUND, "RecordNotFound"), + json_error(StatusCode::BAD_REQUEST, "RecordNotFound"), + ]; + for probe in probes { + let app = no_proof_pds(&repo, probe.clone(), None, Some(repo.blocks.clone())); + + let (result, stored) = backfill_no_proof(app, BackfillStrategy::SparseFilter).await; + + let Ok(SparseBackfillResult::Imported(success)) = result else { + panic!("probe {probe:?}: expected an import, got {result:?}"); + }; + assert_eq!(success.records, 1, "probe {probe:?}"); + assert_eq!( + stored.as_deref(), + Some(honest.as_slice()), + "probe {probe:?}" + ); + } + } + + #[tokio::test] + async fn latest_commit_for_another_did_is_rejected() { + let repo = test_repo(OTHER_DID, &tangled_repo("honest")); + let probe = json_error(StatusCode::NOT_FOUND, "RecordNotFound"); + let app = no_proof_pds(&repo, probe, None, Some(repo.blocks.clone())); + + let (result, stored) = backfill_no_proof(app, BackfillStrategy::SparseFilter).await; + + let err = result.unwrap_err(); + assert!(err.to_string().contains("commit is for"), "{err}"); + assert_eq!(stored, None); + } + + #[tokio::test] + async fn proof_commit_for_another_did_is_rejected() { + let repo = test_repo(OTHER_DID, &tangled_repo("honest")); + let mut car = Vec::new(); + let mut writer = + iroh_car::CarWriter::new(iroh_car::CarHeader::new_v1(vec![repo.commit]), &mut car); + for (cid, bytes) in &repo.blocks { + writer.write(*cid, bytes.clone()).await.unwrap(); + } + writer.finish().await.unwrap(); + + let (result, stored) = backfill_from(vec![(GET_RECORD, car.into())], true).await; + + let err = result.unwrap_err(); + assert!(err.to_string().contains("commit is for"), "{err}"); + assert_eq!(stored, None); + } + + #[tokio::test] + async fn forged_latest_commit_block_is_rejected() { + // getLatestCommit names one commit, and getBlocks hands back another under its cid + let honest = tangled_repo("honest"); + let repo = test_repo(PROOF_DID, &honest); + let forged = test_repo(PROOF_DID, &tangled_repo("forged")); + let mut blocks = forged.blocks.clone(); + blocks.insert(repo.commit, forged.blocks[&forged.commit].clone()); + let probe = json_error(StatusCode::NOT_FOUND, "RecordNotFound"); + let app = no_proof_pds(&repo, probe, None, Some(blocks)); + + let (result, stored) = backfill_no_proof(app, BackfillStrategy::SparseFilter).await; + + let err = result.unwrap_err(); + assert!(err.to_string().contains("CAR block CID mismatch"), "{err}"); + assert_eq!(stored, None); + } + + #[tokio::test] + async fn latest_commit_that_ends_the_sparse_attempt_skips_to_full() { + let repo = test_repo(PROOF_DID, &tangled_repo("honest")); + let probe = json_error(StatusCode::NOT_FOUND, "RecordNotFound"); + let cases = [ + ( + "inactive repo", + Some(json_error(StatusCode::BAD_REQUEST, "RepoDeactivated")), + Some(repo.blocks.clone()), + ), + ( + "no getLatestCommit", + Some((StatusCode::NOT_FOUND, "text/plain", String::new())), + Some(repo.blocks.clone()), + ), + ("commit not served", None, Some(BTreeMap::new())), + ]; + for (case, latest, blocks) in cases { + let app = no_proof_pds(&repo, probe.clone(), latest, blocks); + + let (result, stored) = backfill_no_proof(app, BackfillStrategy::SparseFilter).await; + + assert!( + matches!(result, Ok(SparseBackfillResult::Skipped)), + "{case}: expected a skip, got {result:?}" + ); + assert_eq!(stored, None, "{case}"); + } + } + + #[tokio::test] + async fn pds_without_get_blocks_fails_the_sparse_attempt() { + let repo = test_repo(PROOF_DID, &tangled_repo("honest")); + let probe = ( + StatusCode::NOT_FOUND, + "text/plain", + "mst: not found".to_string(), + ); + let app = no_proof_pds(&repo, probe, None, None); + + let (result, stored) = backfill_no_proof(app, BackfillStrategy::SparseFilter).await; + + let err = result.unwrap_err(); + assert!( + err.to_string().contains("getBlocks failed with HTTP 404"), + "{err}" + ); + assert_eq!(stored, None); + } + + #[tokio::test] + async fn auto_still_takes_full_for_a_tiny_tree_reached_from_the_latest_commit() { + let repo = test_repo(PROOF_DID, &tangled_repo("honest")); + let probe = json_error(StatusCode::NOT_FOUND, "RecordNotFound"); + let app = no_proof_pds(&repo, probe, None, Some(repo.blocks.clone())); + + let (result, stored) = backfill_no_proof(app, BackfillStrategy::Auto).await; + + assert!( + matches!(result, Ok(SparseBackfillResult::Skipped)), + "expected auto to pick full, got {result:?}" + ); + assert_eq!(stored, None); + } + fn assert_imported_one(result: Result) { let Ok(SparseBackfillResult::Imported(success)) = result else { panic!("expected an import, got {result:?}"); diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index 3d05ba9..8d4e54d 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -17,7 +17,7 @@ use jacquard_common::types::did::Did; use crate::backfill::admission::BackfillAdmission; use crate::backfill::client::{ThrottledHttpClient, collect_body_bounded}; use crate::backfill::error::BackfillError; -use crate::backfill::sparse::{SparseBackfillResult, process_did_sparse}; +use crate::backfill::sparse::{SparseBackfillResult, check_commit_did, process_did_sparse}; use crate::config::{BackfillStrategy, RateTier}; use crate::db::types::{DbAction, DbRkey}; use crate::db::{self, Txn as DbTxn, keys}; @@ -292,6 +292,7 @@ pub(super) async fn process_did( let root_commit = jacquard_repo::commit::Commit::::from_cbor(&root_bytes) .into_diagnostic()?; + check_commit_did(root_commit.did.as_str(), did)?; debug!( rev = %root_commit.rev, cid = %root_commit.data, diff --git a/src/backfill/worker/task.rs b/src/backfill/worker/task.rs index ec9bebc..de60ff2 100644 --- a/src/backfill/worker/task.rs +++ b/src/backfill/worker/task.rs @@ -511,6 +511,14 @@ mod tests { } async fn build_car(rev: &str, records: &[(&str, &[u8])]) -> miette::Result { + build_car_for(DID, rev, records).await + } + + async fn build_car_for( + commit_did: &'static str, + 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 { @@ -525,7 +533,7 @@ mod tests { } 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 did: Did = Did::new_static(commit_did).into_diagnostic()?; let commit = AtpCommit { did, version: 3, @@ -779,6 +787,22 @@ mod tests { Ok(()) } + #[tokio::test] + async fn car_whose_commit_names_another_did_is_rejected() -> miette::Result<()> { + let fixture = Fixture::new().await?; + fixture.seed(&RepoState::backfilling().into_static(), 9, &[])?; + let body = post("a"); + let other = "did:plc:bbbbbbbbbbbbbbbbbbbbbbbb"; + let car = build_car_for(other, "3jzfcijpj2z2a", &[("a", body.as_slice())]).await?; + fixture.pds.serve_car(car.bytes).await; + + let err = fixture.run(keys::pending_key(9)).await.unwrap_err(); + + assert!(err.to_string().contains("commit is for"), "{err}"); + assert!(fixture.record_rkeys()?.is_empty()); + Ok(()) + } + #[tokio::test] async fn repo_not_found_keeps_a_rooted_snapshot_and_parks_a_retry() -> miette::Result<()> { let fixture = Fixture::new().await?; -- 2.51.2