diff --git a/crates/tranquil-pds/src/scheduled.rs b/crates/tranquil-pds/src/scheduled.rs index a71c5e9..d808958 100644 --- a/crates/tranquil-pds/src/scheduled.rs +++ b/crates/tranquil-pds/src/scheduled.rs @@ -32,6 +32,17 @@ where .await } +async fn collect_serial_tasks(tasks: I) -> Vec +where + I: IntoIterator, + F: Future, +{ + futures::stream::iter(tasks) + .then(|task| task) + .collect() + .await +} + async fn process_repo_rev( repo_repo: &dyn RepoRepository, block_store: &AnyBlockStore, @@ -284,7 +295,7 @@ async fn process_record_blobs( let mut batch_record_uris: Vec = Vec::new(); let mut batch_blob_cids: Vec = Vec::new(); - collect_background_tasks(records.into_iter().map(|record| { + collect_serial_tasks(records.into_iter().map(|record| { let did = did.clone(); async move { let cid = Cid::from_str(&record.record_cid).ok()?; @@ -932,6 +943,36 @@ mod tests { "background work exceeded its concurrency limit" ); } + + #[tokio::test] + async fn nested_record_tasks_do_not_multiply_background_concurrency() { + let active = Arc::new(AtomicUsize::new(0)); + let peak = Arc::new(AtomicUsize::new(0)); + + collect_background_tasks((0..BACKGROUND_CONCURRENCY_LIMIT * 2).map(|_| { + let active = Arc::clone(&active); + let peak = Arc::clone(&peak); + async move { + collect_serial_tasks((0..BACKGROUND_CONCURRENCY_LIMIT).map(|_| { + let active = Arc::clone(&active); + let peak = Arc::clone(&peak); + async move { + let now_active = active.fetch_add(1, Ordering::SeqCst) + 1; + peak.fetch_max(now_active, Ordering::SeqCst); + tokio::time::sleep(Duration::from_millis(5)).await; + active.fetch_sub(1, Ordering::SeqCst); + } + })) + .await; + } + })) + .await; + + assert!( + peak.load(Ordering::SeqCst) <= BACKGROUND_CONCURRENCY_LIMIT, + "nested record work multiplied background concurrency" + ); + } } pub struct ReachabilityResult { diff --git a/crates/tranquil-server/src/main.rs b/crates/tranquil-server/src/main.rs index 5c22a9a..ae38ebb 100644 --- a/crates/tranquil-server/src/main.rs +++ b/crates/tranquil-server/src/main.rs @@ -193,11 +193,9 @@ async fn run() -> Result<(), Box> { let backfill_repo_repo = state.repos.repo.clone(); let backfill_block_store = state.block_store.clone(); tokio::spawn(async move { - tokio::join!( - backfill_repo_rev(backfill_repo_repo.clone(), backfill_block_store.clone()), - backfill_user_blocks(backfill_repo_repo.clone(), backfill_block_store.clone()), - backfill_record_blobs(backfill_repo_repo, backfill_block_store), - ); + backfill_repo_rev(backfill_repo_repo.clone(), backfill_block_store.clone()).await; + backfill_user_blocks(backfill_repo_repo.clone(), backfill_block_store.clone()).await; + backfill_record_blobs(backfill_repo_repo, backfill_block_store).await; }); let mut comms_service = CommsService::new(state.repos.infra.clone());