diff --git a/crates/tranquil-db-traits/src/repo.rs b/crates/tranquil-db-traits/src/repo.rs index 3520753..62e47a1 100644 --- a/crates/tranquil-db-traits/src/repo.rs +++ b/crates/tranquil-db-traits/src/repo.rs @@ -551,6 +551,8 @@ pub trait RepoRepository: Send + Sync { record_uris: &[AtUri], blob_cids: &[CidLink], ) -> Result<(), DbError>; + + async fn mark_record_blobs_backfill_complete(&self, repo_id: Uuid) -> Result<(), DbError>; } #[async_trait] diff --git a/crates/tranquil-db/src/postgres/repo.rs b/crates/tranquil-db/src/postgres/repo.rs index ed04a06..6bc71c5 100644 --- a/crates/tranquil-db/src/postgres/repo.rs +++ b/crates/tranquil-db/src/postgres/repo.rs @@ -1630,25 +1630,25 @@ impl RepoRepository for PostgresRepoRepository { &self, limit: i64, ) -> Result, DbError> { - let rows = sqlx::query!( + let rows: Vec<(Uuid, String)> = sqlx::query_as( r#" - SELECT DISTINCT u.id as user_id, u.did + SELECT u.id, u.did FROM users u - JOIN records r ON r.repo_id = u.id - WHERE NOT EXISTS (SELECT 1 FROM record_blobs rb WHERE rb.repo_id = u.id) + JOIN repos r ON r.user_id = u.id + WHERE r.record_blobs_backfilled_at IS NULL LIMIT $1 "#, - limit ) + .bind(limit) .fetch_all(&self.pool) .await .map_err(map_sqlx_error)?; rows.into_iter() - .map(|r| { + .map(|(user_id, did)| { Ok(UserNeedingRecordBlobsBackfill { - user_id: r.user_id, - did: column(r.did, col::USERS_DID)?, + user_id, + did: column(did, col::USERS_DID)?, }) }) .collect() @@ -1680,4 +1680,14 @@ impl RepoRepository for PostgresRepoRepository { Ok(()) } + + async fn mark_record_blobs_backfill_complete(&self, repo_id: Uuid) -> Result<(), DbError> { + sqlx::query("UPDATE repos SET record_blobs_backfilled_at = NOW() WHERE user_id = $1") + .bind(repo_id) + .execute(&self.pool) + .await + .map_err(map_sqlx_error)?; + + Ok(()) + } } diff --git a/crates/tranquil-pds/src/scheduled.rs b/crates/tranquil-pds/src/scheduled.rs index 438e3c6..a71c5e9 100644 --- a/crates/tranquil-pds/src/scheduled.rs +++ b/crates/tranquil-pds/src/scheduled.rs @@ -1,5 +1,6 @@ use anyhow::Context; use cid::Cid; +use futures::{Future, StreamExt}; use ipld_core::ipld::Ipld; use jacquard_repo::commit::Commit; use jacquard_repo::storage::BlockStore; @@ -18,6 +19,19 @@ use crate::repo::AnyBlockStore; use crate::storage::BlobStorage; use crate::sync::car::encode_car_header; +const BACKGROUND_CONCURRENCY_LIMIT: usize = 8; + +async fn collect_background_tasks(tasks: I) -> Vec +where + I: IntoIterator, + F: Future, +{ + futures::stream::iter(tasks) + .buffer_unordered(BACKGROUND_CONCURRENCY_LIMIT) + .collect() + .await +} + async fn process_repo_rev( repo_repo: &dyn RepoRepository, block_store: &AnyBlockStore, @@ -64,7 +78,7 @@ pub async fn backfill_repo_rev(repo_repo: Arc, block_store: "Backfilling repo_rev for existing repos" ); - let results = futures::future::join_all(repos_missing_rev.into_iter().map(|repo| { + let results = collect_background_tasks(repos_missing_rev.into_iter().map(|repo| { let repo_repo = repo_repo.clone(); let block_store = block_store.clone(); async move { @@ -140,7 +154,7 @@ pub async fn backfill_user_blocks(repo_repo: Arc, block_stor "Backfilling user_blocks for existing repos" ); - let results = futures::future::join_all(users_without_blocks.into_iter().map(|user| { + let results = collect_background_tasks(users_without_blocks.into_iter().map(|user| { let repo_repo = repo_repo.clone(); let block_store = block_store.clone(); async move { @@ -270,7 +284,7 @@ async fn process_record_blobs( let mut batch_record_uris: Vec = Vec::new(); let mut batch_blob_cids: Vec = Vec::new(); - futures::future::join_all(records.into_iter().map(|record| { + collect_background_tasks(records.into_iter().map(|record| { let did = did.clone(); async move { let cid = Cid::from_str(&record.record_cid).ok()?; @@ -304,6 +318,10 @@ async fn process_record_blobs( .await .map_err(|_| (user_id, "failed to insert"))?; } + repo_repo + .mark_record_blobs_backfill_complete(user_id) + .await + .map_err(|_| (user_id, "failed to mark complete"))?; Ok((user_id, did, blob_refs_found)) } @@ -327,7 +345,7 @@ pub async fn backfill_record_blobs(repo_repo: Arc, block_sto "Backfilling record_blobs for existing repos" ); - let results = futures::future::join_all(users_needing_backfill.into_iter().map(|user| { + let results = collect_background_tasks(users_needing_backfill.into_iter().map(|user| { let repo_repo = repo_repo.clone(); let block_store = block_store.clone(); async move { @@ -707,7 +725,7 @@ async fn process_scheduled_deletions( "Processing scheduled account deletions" ); - futures::future::join_all(accounts_to_delete.into_iter().map(|account| async move { + collect_background_tasks(accounts_to_delete.into_iter().map(|account| async move { let result = delete_account_data(user_repo, blob_repo, blob_store, account.id, &account.did).await; (account.did, account.handle, result) @@ -886,6 +904,36 @@ pub async fn generate_repo_car_from_user_blocks( generate_repo_car(block_store, &actual_head_cid).await } +#[cfg(test)] +mod tests { + use super::*; + use std::sync::Arc; + use std::sync::atomic::{AtomicUsize, Ordering}; + + #[tokio::test] + async fn background_tasks_are_concurrency_limited() { + let active = Arc::new(AtomicUsize::new(0)); + let peak = Arc::new(AtomicUsize::new(0)); + + collect_background_tasks((0..BACKGROUND_CONCURRENCY_LIMIT * 4).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(10)).await; + active.fetch_sub(1, Ordering::SeqCst); + } + })) + .await; + + assert!( + peak.load(Ordering::SeqCst) <= BACKGROUND_CONCURRENCY_LIMIT, + "background work exceeded its concurrency limit" + ); + } +} + pub struct ReachabilityResult { pub repos_walked: u64, pub blocks_visited: u64, diff --git a/crates/tranquil-store/src/metastore/client.rs b/crates/tranquil-store/src/metastore/client.rs index dfbac2a..d0fc787 100644 --- a/crates/tranquil-store/src/metastore/client.rs +++ b/crates/tranquil-store/src/metastore/client.rs @@ -797,6 +797,14 @@ impl tranquil_db_traits::RepoRepository for MetastoreCli )))?; recv(rx).await } + + async fn mark_record_blobs_backfill_complete(&self, repo_id: Uuid) -> Result<(), DbError> { + let (tx, rx) = oneshot::channel(); + self.pool.send(MetastoreRequest::Commit(Box::new( + CommitRequest::MarkRecordBlobsBackfillComplete { repo_id, tx }, + )))?; + recv(rx).await + } } #[async_trait] diff --git a/crates/tranquil-store/src/metastore/commit_ops.rs b/crates/tranquil-store/src/metastore/commit_ops.rs index 0e9003a..9348418 100644 --- a/crates/tranquil-store/src/metastore/commit_ops.rs +++ b/crates/tranquil-store/src/metastore/commit_ops.rs @@ -73,6 +73,13 @@ pub(crate) fn record_blobs_user_prefix(user_hash: UserHash) -> SmallVec<[u8; 128 .build() } +pub(crate) fn record_blobs_backfilled_key(user_hash: UserHash) -> SmallVec<[u8; 128]> { + KeyBuilder::new() + .tag(KeyTag::RECORD_BLOBS_BACKFILLED) + .u64(user_hash.raw()) + .build() +} + pub struct CommitOps { db: Database, repo_data: Keyspace, @@ -366,7 +373,7 @@ impl CommitOps { let limit_usize = usize::try_from(limit).unwrap_or(0); self.scan_users_missing_prefix( - record_blobs_user_prefix, + record_blobs_backfilled_key, |meta, user_id| { let did = match meta.did { None => Err(MetastoreError::CorruptData("repo_meta missing DID field")), @@ -379,6 +386,17 @@ impl CommitOps { ) } + pub fn mark_record_blobs_backfill_complete(&self, repo_id: Uuid) -> Result<(), MetastoreError> { + let user_hash = self + .user_hashes + .get(&repo_id) + .ok_or(MetastoreError::InvalidInput("unknown user_id"))?; + let key = record_blobs_backfilled_key(user_hash); + self.repo_data + .insert(key.as_slice(), &[]) + .map_err(MetastoreError::Fjall) + } + pub fn get_users_without_blocks(&self) -> Result, MetastoreError> { const MAX_RESULTS: usize = 10_000; @@ -909,6 +927,16 @@ mod tests { let (user_id_a, did_a, _) = create_test_repo(&h, "teq", 60); let (user_id_b, _did_b, _) = create_test_repo(&h, "nel", 61); + let mut batch = h.metastore.database().batch(); + [user_id_a, user_id_b].iter().for_each(|user_id| { + let user_hash = ops.user_hashes.get(user_id).unwrap(); + batch.remove( + &ops.repo_data, + record_blobs_backfilled_key(user_hash).as_slice(), + ); + }); + batch.commit().unwrap(); + let needing = ops.get_users_needing_record_blobs_backfill(100).unwrap(); assert_eq!(needing.len(), 2); @@ -920,10 +948,18 @@ mod tests { let blob_cid = test_cid_link(62); ops.insert_record_blobs(user_id_a, &[uri], &[blob_cid]) .unwrap(); + ops.mark_record_blobs_backfill_complete(user_id_a).unwrap(); let needing_after = ops.get_users_needing_record_blobs_backfill(100).unwrap(); assert_eq!(needing_after.len(), 1); assert_eq!(needing_after[0].user_id, user_id_b); + + ops.mark_record_blobs_backfill_complete(user_id_b).unwrap(); + assert!( + ops.get_users_needing_record_blobs_backfill(100) + .unwrap() + .is_empty() + ); } #[test] diff --git a/crates/tranquil-store/src/metastore/handler.rs b/crates/tranquil-store/src/metastore/handler.rs index 64c2275..38f34f9 100644 --- a/crates/tranquil-store/src/metastore/handler.rs +++ b/crates/tranquil-store/src/metastore/handler.rs @@ -505,6 +505,10 @@ pub enum CommitRequest { blob_cids: Vec, tx: Tx<()>, }, + MarkRecordBlobsBackfillComplete { + repo_id: Uuid, + tx: Tx<()>, + }, } impl CommitRequest { @@ -514,6 +518,9 @@ impl CommitRequest { Self::ImportRepoData { user_id, .. } | Self::InsertRecordBlobs { repo_id: user_id, .. + } + | Self::MarkRecordBlobsBackfillComplete { + repo_id: user_id, .. } => uuid_to_routing(user_hashes, user_id), Self::GetUsersWithoutBlocks { .. } | Self::GetUsersNeedingRecordBlobsBackfill { .. } => Routing::Global, @@ -3146,6 +3153,14 @@ fn dispatch_commit(state: &HandlerState, req: CommitR .map_err(metastore_to_db), ); } + CommitRequest::MarkRecordBlobsBackfillComplete { repo_id, tx } => { + let _ = tx.send( + state + .commit_ops + .mark_record_blobs_backfill_complete(repo_id) + .map_err(metastore_to_db), + ); + } } } diff --git a/crates/tranquil-store/src/metastore/keys.rs b/crates/tranquil-store/src/metastore/keys.rs index 2b64151..6008d52 100644 --- a/crates/tranquil-store/src/metastore/keys.rs +++ b/crates/tranquil-store/src/metastore/keys.rs @@ -55,6 +55,7 @@ impl KeyTag { pub const BACKLINK_BY_USER: Self = Self(0x31); pub const RECORD_BY_CID: Self = Self(0x32); pub const RECORD_BY_CID_BUILT: Self = Self(0x33); + pub const RECORD_BLOBS_BACKFILLED: Self = Self(0x34); pub const USER_PRIMARY: Self = Self(0x40); pub const USER_BY_HANDLE: Self = Self(0x41); @@ -189,6 +190,7 @@ mod tests { KeyTag::BACKLINK_BY_USER, KeyTag::RECORD_BY_CID, KeyTag::RECORD_BY_CID_BUILT, + KeyTag::RECORD_BLOBS_BACKFILLED, KeyTag::USER_PRIMARY, KeyTag::USER_BY_HANDLE, KeyTag::USER_BY_EMAIL, diff --git a/crates/tranquil-store/src/metastore/repo_ops.rs b/crates/tranquil-store/src/metastore/repo_ops.rs index 6d18feb..3d172e3 100644 --- a/crates/tranquil-store/src/metastore/repo_ops.rs +++ b/crates/tranquil-store/src/metastore/repo_ops.rs @@ -4,6 +4,7 @@ use std::sync::Arc; use uuid::Uuid; use super::MetastoreError; +use super::commit_ops::record_blobs_backfilled_key; use super::encoding::KeyReader; use super::keys::{KeyTag, UserHash}; use super::records::{record_by_cid_user_prefix, record_user_prefix}; @@ -68,6 +69,11 @@ impl RepoOps { handle_key(&handle_lower).as_slice(), user_hash.raw().to_be_bytes(), ); + batch.insert( + &self.repo_data, + record_blobs_backfilled_key(user_hash).as_slice(), + [], + ); match batch.commit() { Ok(()) => Ok(()), @@ -246,6 +252,10 @@ impl RepoOps { let mut batch = db.batch(); stage_repo_meta_removal(&mut batch, &self.repo_data, user_hash, &meta.handle); + batch.remove( + &self.repo_data, + record_blobs_backfilled_key(user_hash).as_slice(), + ); self.user_hashes.stage_remove(&mut batch, &user_id); match batch.commit() { @@ -630,6 +640,7 @@ pub(super) fn stage_full_repo_data_removal( batch, user_block_user_prefix(user_hash).as_slice(), )?; + batch.remove(repo_data, record_blobs_backfilled_key(user_hash).as_slice()); Ok(()) } @@ -747,15 +758,29 @@ mod tests { let did = test_did("lyna"); let handle = test_handle("lyna"); let cid = test_cid_link(4); + let user_hash = UserHash::from_did(did.as_str()); + let backfill_key = record_blobs_backfilled_key(user_hash); ops.create_repo(ms.database(), user_id, &did, &handle, &cid, &test_rev(1)) .unwrap(); assert!(ops.get_repo(user_id).unwrap().is_some()); assert!(ops.lookup_handle(&handle).unwrap().is_some()); + assert!( + ops.repo_data + .get(backfill_key.as_slice()) + .unwrap() + .is_some() + ); ops.delete_repo(ms.database(), user_id).unwrap(); assert!(ops.get_repo(user_id).unwrap().is_none()); assert!(ops.lookup_handle(&handle).unwrap().is_none()); + assert!( + ops.repo_data + .get(backfill_key.as_slice()) + .unwrap() + .is_none() + ); } #[test] diff --git a/migrations/20260813_record_blobs_backfill_marker.sql b/migrations/20260813_record_blobs_backfill_marker.sql new file mode 100644 index 0000000..6c5f1e8 --- /dev/null +++ b/migrations/20260813_record_blobs_backfill_marker.sql @@ -0,0 +1,13 @@ +ALTER TABLE repos + ADD COLUMN record_blobs_backfilled_at TIMESTAMPTZ; + +ALTER TABLE repos + ALTER COLUMN record_blobs_backfilled_at SET DEFAULT NOW(); + +UPDATE repos r +SET record_blobs_backfilled_at = NOW() +WHERE EXISTS ( + SELECT 1 + FROM record_blobs rb + WHERE rb.repo_id = r.user_id +); \ No newline at end of file