From feabfd3b2b0ba62d0a0b49cee00c667ca48b547d Mon Sep 17 00:00:00 2001 From: dawn <90008@klbr.net> Date: Wed, 30 Sep 2026 08:02:12 +0300 Subject: [PATCH] [backfill] walk full repo msts sequentially off the async runtime --- src/backfill/worker/process.rs | 23 ++++++++++++++++------- src/car.rs | 2 ++ 2 files changed, 18 insertions(+), 7 deletions(-) diff --git a/src/backfill/worker/process.rs b/src/backfill/worker/process.rs index abd331a..f36e661 100644 --- a/src/backfill/worker/process.rs +++ b/src/backfill/worker/process.rs @@ -316,14 +316,23 @@ pub(super) async fn process_did( let root_commit = Commit::from(root_commit); - // 5. walk mst and fetch every record block in one batch + // 5. walk mst and fetch every record block in one batch, off the runtime. the sequential + // walk decodes every node on this blocking thread; `leaves()` would spawn a task per + // subtree onto the runtime workers and spend 3-4x the cpu doing it. let start = Instant::now(); - let mst: Mst = Mst::load(store, root_commit.data, None); - let leaves = mst.leaves().await.into_diagnostic()?; - let leaf_cids = leaves.iter().map(|(_, cid)| *cid).collect::>(); - let leaf_blocks = mst.storage().get_many(&leaf_cids).await.into_diagnostic()?; - let records = leaves.into_iter().zip(leaf_blocks); - drop(mst); + let mst_root = root_commit.data; + let handle = tokio::runtime::Handle::current(); + let records = tokio::task::spawn_blocking(move || { + let mst: Mst = Mst::load(store, mst_root, None); + let leaves = handle.block_on(mst.leaves_sequential()).into_diagnostic()?; + let leaf_cids = leaves.iter().map(|(_, cid)| *cid).collect::>(); + let leaf_blocks = handle + .block_on(mst.storage().get_many(&leaf_cids)) + .into_diagnostic()?; + Ok::<_, miette::Report>(leaves.into_iter().zip(leaf_blocks)) + }) + .await + .into_diagnostic()??; trace!(elapsed = %start.elapsed().as_secs_f32(), "walked mst"); // 6. insert records into db diff --git a/src/car.rs b/src/car.rs index 4f8ec61..276c1cc 100644 --- a/src/car.rs +++ b/src/car.rs @@ -350,6 +350,8 @@ mod tests { let loaded = Mst::load(Arc::new(CarBlockStore::new(parsed.blocks)), root, None); let expected = mst.leaves().await.unwrap(); assert_eq!(loaded.leaves().await.unwrap(), expected); + // backfill walks with leaves_sequential; pin that it matches the parallel walk + assert_eq!(loaded.leaves_sequential().await.unwrap(), expected); let leaf_cids: Vec<_> = expected.iter().map(|(_, cid)| *cid).collect(); let bodies = loaded.storage().get_many(&leaf_cids).await.unwrap(); assert!(bodies.iter().all(Option::is_some)); -- 2.51.2