Something went wrong. Try again.
lightweight `com.atproto.sync.listReposByCollection`
Something went wrong. Try again.
Rust
1234567891011121314151617181920212223242526272829303132333435363738394041424344454647484950515253545556575859606162636465666768697071727374757677787980818283848586878889909192939495969798991001011021031041051061071081091101111121131141151161171181191201211221231241251261271281291301311321331341351361371381391401411421431441451461471481491501511521531541551561571581591601611621631641651661671681691701711721731741751761771781791801811821831841851861871881891901911921931941951961971981992002012022032042052062072082092102112122132142152162172182192202212222232242252262272282292302312322332342352362372382392402412422432442452462472482492502512522532542552562572582592602612622632642652662672682692702712722732742752762772782792802812822832842852862872882892902912922932942952962972982993003013023033043053063073083093103113123133143153163173183193203213223233243253263273283293303313323333343353363373383393403413423433443453463473483493503513523533543553563573583593603613623633643653663673683693703713723733743753763773783793803813823833843853863873883893903913923933943953963973983994004014024034044054064074084094104114124134144154164174184194204214224234244254264274284294304314324334344354364374384394404414424434444454464474484494504514524534544554564574584594604614624634644654664674684694704714724734744754764774784794804814824834844854864874884894904914924934944954964974984995005015025035045055065075085095105115125135145155165175185195205215225235245255265275285295305315325335345355365375385395405415425435445455465475485495505515525535545555565575585595605615625635645655665675685695705715725735745755765775785795805815825835845855865875885895905915925935945955965975985996006016026036046056066076086096106116126136146156166176186196206216226236246256266276286296306316326336346356366376386396406416426436446456466476486496506516526536546556566576586596606616626636646656666676686696706716726736746756766776786796806816826836846856866876886896906916926936946956966976986997007017027037047057067077087097107117127137147157167177187197207217227237247257267277287297307317327337347357367377387397407417427437447457467477487497507517527537547557567577587597607617627637647657667677687697707717727737747757767777787797807817827837847857867877887897907917927937947957967977987998008018028038048058068078088098108118128138148158168178188198208218228238248258268278288298308318328338348358368378388398408418428438448458468478488498508518528538548558568578588598608618628638648658668678688698708718728738748758768778788798808818828838848858868878888898908918928938948958968978988999009019029039049059069079089099109119129139149159169179189199209219229239249259269279289299309319329339349359369379389399409419429439449459469479489499509519529539549559569579589599609619629639649659669679689699709719729739749759769779789799809819829839849859869879889899909919929939949959969979989991000100110021003//! Two-level host-partitioned resync queue.//!//! Keys: `"rsH" <host_len:u8> <host_bytes> \0 <ts_be:u64> \0 <did>`//! Values: `[u16 BE retry_count][u16 BE reason_len][reason_bytes][commit_cbor_bytes]`//!//! Items are partitioned by PDS hostname so that the dispatcher can round-robin//! across hosts for fair scheduling. An in-memory [`HostIndex`] tracks the set//! of non-empty hosts and per-host counts so that `next_host` is O(hosts) rather//! than O(queue depth).
use std::collections::{BTreeSet, HashMap, HashSet};use std::sync::atomic::Ordering;
use jacquard_common::types::string::Did;use tracing::{debug, trace, warn};
use crate::storage::{ DbRef, PREFIX_RESYNC_HOST_QUEUE, PREFIX_RESYNC_QUEUE, error::{StorageError, StorageResult}, repo,};
const NUL: u8 = b'\0';
// ---------------------------------------------------------------------------// In-memory host index// ---------------------------------------------------------------------------
/// Tracks the set of hosts that have queued items for round-robin scheduling.pub struct HostIndex { hosts: BTreeSet<String>, host_counts: HashMap<String, u64>, rr_cursor: usize,}
impl HostIndex { pub fn new() -> Self { Self { hosts: BTreeSet::new(), host_counts: HashMap::new(), rr_cursor: 0, } }
fn increment(&mut self, host: &str) { let count = self.host_counts.entry(host.to_owned()).or_insert(0); *count += 1; self.hosts.insert(host.to_owned()); }
fn decrement(&mut self, host: &str) { if let Some(count) = self.host_counts.get_mut(host) { *count = count.saturating_sub(1); if *count == 0 { self.host_counts.remove(host); self.hosts.remove(host); // If cursor was past the removed host, adjust if self.rr_cursor > self.hosts.len() { self.rr_cursor = 0; } } } }}
// ---------------------------------------------------------------------------// Key encoding (new host-partitioned format)// ---------------------------------------------------------------------------
/// `"rsH" <host_len:u8> <host_bytes> \0 <ts_be:u64> \0 <did>`fn host_key(host: &str, ts: u64, did: &Did<'_>) -> Vec<u8> { let d = did.as_str(); let hb = host.as_bytes(); let mut k = Vec::with_capacity(PREFIX_RESYNC_HOST_QUEUE.len() + 1 + hb.len() + 1 + 8 + 1 + d.len()); k.extend_from_slice(&PREFIX_RESYNC_HOST_QUEUE); k.push(hb.len() as u8); k.extend_from_slice(hb); k.push(NUL); k.extend_from_slice(&ts.to_be_bytes()); k.push(NUL); k.extend_from_slice(d.as_bytes()); k}
/// Prefix for scanning all items for one host: `"rsH" <host_len:u8> <host_bytes> \0`fn host_prefix(host: &str) -> Vec<u8> { let hb = host.as_bytes(); let mut k = Vec::with_capacity(PREFIX_RESYNC_HOST_QUEUE.len() + 1 + hb.len() + 1); k.extend_from_slice(&PREFIX_RESYNC_HOST_QUEUE); k.push(hb.len() as u8); k.extend_from_slice(hb); k.push(NUL); k}
/// Prefix for scanning the entire host-partitioned queue.fn key_prefix_all() -> Vec<u8> { PREFIX_RESYNC_HOST_QUEUE.to_vec()}
/// Upper bound for timestamp within a host prefix (exclusive):/// `"rsH" <host_len:u8> <host_bytes> \0 <ts_be:u64> \0`////// Used as the exclusive end of a range scan so that we only see items with/// timestamps strictly less than `ts`.fn host_ts_upper(host: &str, ts: u64) -> Vec<u8> { let hb = host.as_bytes(); let mut k = Vec::with_capacity(PREFIX_RESYNC_HOST_QUEUE.len() + 1 + hb.len() + 1 + 8 + 1); k.extend_from_slice(&PREFIX_RESYNC_HOST_QUEUE); k.push(hb.len() as u8); k.extend_from_slice(hb); k.push(NUL); k.extend_from_slice(&ts.to_be_bytes()); k.push(NUL); k}
/// Parse host, timestamp, and DID from a full host-partitioned key.fn host_key_parse(raw: &[u8]) -> StorageResult<(String, u64, Did<'static>)> { let key_str = String::from_utf8_lossy(raw); let rest = raw .strip_prefix(&PREFIX_RESYNC_HOST_QUEUE) .ok_or(StorageError::Corrupt { key: key_str.to_string(), reason: "wrong prefix for host resync queue", })?; if rest.is_empty() { return Err(StorageError::Corrupt { key: key_str.to_string(), reason: "missing host_len byte", }); } let host_len = rest[0] as usize; let rest = &rest[1..]; if rest.len() < host_len { return Err(StorageError::Corrupt { key: key_str.to_string(), reason: "host bytes truncated", }); } let host = std::str::from_utf8(&rest[..host_len]) .map_err(|_| StorageError::Corrupt { key: key_str.to_string(), reason: "host not valid UTF-8", })? .to_owned(); let rest = rest[host_len..] .strip_prefix(&[NUL]) .ok_or(StorageError::Corrupt { key: key_str.to_string(), reason: "missing NUL after host", })?; if rest.len() < 9 { return Err(StorageError::Corrupt { key: key_str.to_string(), reason: "not enough bytes for timestamp + NUL", }); } let ts_bytes: [u8; 8] = rest[..8].try_into().map_err(|_| StorageError::Corrupt { key: key_str.to_string(), reason: "timestamp conversion failed", })?; let ts = u64::from_be_bytes(ts_bytes); let rest = rest[8..] .strip_prefix(&[NUL]) .ok_or(StorageError::Corrupt { key: key_str.to_string(), reason: "missing NUL after timestamp", })?; let did_str = std::str::from_utf8(rest).map_err(|_| StorageError::Corrupt { key: key_str.to_string(), reason: "invalid UTF-8 for DID", })?; let did = Did::new_owned(did_str).map_err(|_| StorageError::Corrupt { key: key_str.to_string(), reason: "invalid DID", })?; Ok((host, ts, did))}
// ---------------------------------------------------------------------------// Legacy key helpers (for migration)// ---------------------------------------------------------------------------
/// Parse timestamp and DID from old `"rsq"` format key.fn legacy_key_parse(raw: &[u8]) -> StorageResult<(u64, Did<'static>)> { let key_str = String::from_utf8_lossy(raw); let rest = raw .strip_prefix(&PREFIX_RESYNC_QUEUE) .ok_or(StorageError::Corrupt { key: key_str.to_string(), reason: "wrong prefix for legacy resync queue", })?; if rest.len() < 9 { return Err(StorageError::Corrupt { key: key_str.to_string(), reason: "not enough suffix bytes for legacy resync queue", }); } let ts_bytes: [u8; 8] = rest[..8].try_into().map_err(|_| StorageError::Corrupt { key: key_str.to_string(), reason: "not enough bytes for timestamp", })?; let ts = u64::from_be_bytes(ts_bytes); let rest = rest[8..] .strip_prefix(&[NUL]) .ok_or(StorageError::Corrupt { key: key_str.to_string(), reason: "missing NUL separator in legacy key", })?; let did_str = std::str::from_utf8(rest).map_err(|_| StorageError::Corrupt { key: key_str.to_string(), reason: "invalid UTF-8 for DID in legacy key", })?; let did = Did::new_owned(did_str).map_err(|_| StorageError::Corrupt { key: key_str.to_string(), reason: "invalid DID in legacy key", })?; Ok((ts, did))}
// ---------------------------------------------------------------------------// Value encoding// ---------------------------------------------------------------------------
/// An item waiting in the resync queue.#[derive(Debug, Clone)]pub struct ResyncItem { pub did: Did<'static>, pub retry_count: u16, pub retry_reason: String, /// Raw CBOR of the triggering firehose commit. pub commit_cbor: Vec<u8>,}
fn encode(item: &ResyncItem) -> Vec<u8> { let reason = item.retry_reason.as_bytes(); let mut v = Vec::with_capacity(2 + 2 + reason.len() + item.commit_cbor.len()); v.extend_from_slice(&item.retry_count.to_be_bytes()); v.extend_from_slice(&(reason.len() as u16).to_be_bytes()); v.extend_from_slice(reason); v.extend_from_slice(&item.commit_cbor); v}
fn decode(bytes: &[u8], key_str: &str, did: Did<'static>) -> StorageResult<ResyncItem> { if bytes.len() < 4 { return Err(StorageError::Corrupt { key: key_str.to_owned(), reason: "value too short", }); } let retry_count = u16::from_be_bytes([bytes[0], bytes[1]]); let reason_len = u16::from_be_bytes([bytes[2], bytes[3]]) as usize; let rest = &bytes[4..]; if rest.len() < reason_len { return Err(StorageError::Corrupt { key: key_str.to_owned(), reason: "reason truncated", }); } let retry_reason = std::str::from_utf8(&rest[..reason_len]) .map_err(|_| StorageError::Corrupt { key: key_str.to_owned(), reason: "reason not UTF-8", })? .to_owned(); let commit_cbor = rest[reason_len..].to_vec(); Ok(ResyncItem { did, retry_count, retry_reason, commit_cbor, })}
// ---------------------------------------------------------------------------// Public item type for claimed items// ---------------------------------------------------------------------------
/// A claimed item with its host.pub struct ClaimedItem { pub item: ResyncItem, pub host: String,}
// ---------------------------------------------------------------------------// CRUD// ---------------------------------------------------------------------------
/// Queue a repo for resync into a batch.////// `host` is the PDS hostname; `None` maps to the empty-string (unknown) bucket.pub fn enqueue_into( batch: &mut fjall::OwnedWriteBatch, db: &DbRef, ts: u64, item: &ResyncItem, host: Option<&str>,) { let h = host.unwrap_or(""); if item.retry_reason == "backfill" { trace!( did = item.did.as_str(), ts, host = h, reason = %item.retry_reason, retry = item.retry_count, "enqueue resync to batch" ); } else { debug!( did = item.did.as_str(), ts, host = h, reason = %item.retry_reason, retry = item.retry_count, "enqueue resync to batch" ); } batch.insert(&db.ks, host_key(h, ts, &item.did), encode(item)); db.stats.resync_queue_depth.fetch_add(1, Ordering::Relaxed);
let mut idx = db.host_index.lock().expect("host_index poisoned"); idx.increment(h);}
/// Enqueue a repo for resync at the given Unix timestamp (seconds).////// `host` is the PDS hostname; `None` maps to the empty-string (unknown) bucket.pub fn enqueue(db: &DbRef, ts: u64, item: &ResyncItem, host: Option<&str>) -> StorageResult<()> { let mut batch = db.database.batch(); enqueue_into(&mut batch, db, ts, item, host); batch.commit()?; Ok(())}
/// Count the total number of entries currently in the resync queue.////// Performs a full prefix scan; use only for admin/diagnostic views.pub fn count_queued(db: &DbRef) -> usize { db.ks.prefix(key_prefix_all()).count()}
/// Advance the round-robin cursor and return the next host that has queued items/// and is not in `skip_hosts`. Returns `None` if all hosts are skipped or the/// queue is empty.pub fn next_host(db: &DbRef, skip_hosts: &HashSet<String>) -> Option<String> { let idx = db.host_index.lock().expect("host_index poisoned"); let n = idx.hosts.len(); if n == 0 { return None; }
// Start from current cursor position and wrap around once. let start = idx.rr_cursor % n; let hosts_vec: Vec<&String> = idx.hosts.iter().collect();
for i in 0..n { let pos = (start + i) % n; let host = hosts_vec[pos]; if !skip_hosts.contains(host) { // We need to drop the immutable borrow to take a mutable one. // Instead, capture the result and the new cursor position. let result = host.clone(); let new_cursor = (pos + 1) % n; drop(idx); let mut idx = db.host_index.lock().expect("host_index poisoned"); idx.rr_cursor = new_cursor; return Some(result); } } None}
/// Claim the oldest ready item from a specific host's sub-queue, skipping busy/// DIDs. Atomically removes from queue and transitions repo state to Resyncing.////// Returns `None` if no claimable item exists for this host.pub fn claim_from_host( db: &DbRef, host: &str, now: u64, busy: &HashSet<Did<'_>>,) -> StorageResult<Option<ClaimedItem>> { let prefix = host_prefix(host); let upper = host_ts_upper(host, now);
for guard in db.ks.range(prefix..upper) { let (key_slice, val_slice) = guard.into_inner()?; let key_bytes = key_slice.as_ref(); let (_, _, did) = host_key_parse(key_bytes)?;
if busy.contains(&did) { debug!(did = did.as_str(), host, "skip busy did in host queue"); continue; }
let key_str = String::from_utf8_lossy(key_bytes).into_owned(); let item = decode(val_slice.as_ref(), &key_str, did.clone())?;
// Read current repo info to preserve account status across the transition. let repo_key = repo::key(&did); let new_info = match db.ks.get(&repo_key)? { Some(b) => { let rk = String::from_utf8_lossy(&repo_key).into_owned(); let mut info = repo::decode_repo_info(&b, &rk)?; info.state = repo::RepoState::Resyncing; info.error = None; info } None => { warn!( did = did.as_str(), "claiming resync job for did with no repo record; inserting as resyncing/active" ); repo::RepoInfo { state: repo::RepoState::Resyncing, status: repo::AccountStatus::Active, error: None, } } };
// Atomically: remove from queue + write state=Resyncing. let mut batch = db.database.batch(); batch.remove(&db.ks, key_bytes); batch.insert(&db.ks, &repo_key, repo::encode_repo_info(&new_info)); batch.commit()?;
db.stats.resync_queue_depth.fetch_sub(1, Ordering::Relaxed);
{ let mut idx = db.host_index.lock().expect("host_index poisoned"); idx.decrement(host); }
trace!( did = item.did.as_str(), host, reason = %item.retry_reason, retry = item.retry_count, "claimed resync job from host queue" ); return Ok(Some(ClaimedItem { item, host: host.to_owned(), })); }
Ok(None)}
// ---------------------------------------------------------------------------// Migration + rebuild// ---------------------------------------------------------------------------
/// Migrate old-format `"rsq"` entries to new `"rsH"` format with empty host.////// Called once on startup. Returns the number of migrated entries.pub fn migrate_legacy(db: &DbRef) -> StorageResult<u64> { let prefix = PREFIX_RESYNC_QUEUE.to_vec(); let mut count = 0u64;
for guard in db.ks.prefix(&prefix) { let (key_slice, val_slice) = guard.into_inner()?; let old_key = key_slice.as_ref(); let (ts, did) = legacy_key_parse(old_key)?;
// Write new-format key with empty host, remove old key. let new_key = host_key("", ts, &did); let mut batch = db.database.batch(); batch.insert(&db.ks, &new_key, val_slice.as_ref()); batch.remove(&db.ks, old_key); batch.commit()?;
count += 1; }
if count > 0 { debug!(count, "migrated legacy rsq entries to rsH format"); } Ok(count)}
/// Rebuild the in-memory [`HostIndex`] from a full scan of the queue.////// Called on startup after migration.pub fn rebuild_host_index(db: &DbRef) -> StorageResult<()> { let prefix = key_prefix_all(); let mut hosts = BTreeSet::new(); let mut host_counts: HashMap<String, u64> = HashMap::new();
for guard in db.ks.prefix(&prefix) { let (key_slice, _) = guard.into_inner()?; let (host, _, _) = host_key_parse(key_slice.as_ref())?; *host_counts.entry(host.clone()).or_insert(0) += 1; hosts.insert(host); }
let mut idx = db.host_index.lock().expect("host_index poisoned"); idx.hosts = hosts; idx.host_counts = host_counts; idx.rr_cursor = 0;
Ok(())}
// ---------------------------------------------------------------------------// Tests// ---------------------------------------------------------------------------
#[cfg(test)]mod tests { use std::collections::HashSet;
use super::*; use crate::storage::{open_temporary, repo};
fn did(s: &str) -> Did<'static> { Did::new_owned(s).unwrap() }
fn item(did_str: &str, retry_count: u16, reason: &str, cbor: &[u8]) -> ResyncItem { ResyncItem { did: did(did_str), retry_count, retry_reason: reason.to_owned(), commit_cbor: cbor.to_vec(), } }
fn pending_repo(db: &DbRef, did_str: &str) { repo::put_info( db, &did(did_str), &repo::RepoInfo { state: repo::RepoState::Pending, status: repo::AccountStatus::Active, error: None, }, ) .unwrap(); }
// --- host_key structure and sorting ---
#[test] fn host_key_structure() { let k = host_key("pds1.bsky.network", 0x0102030405060708, &did("did:web:example.com")); let mut expected = b"rsH".to_vec(); let host_bytes = b"pds1.bsky.network"; expected.push(host_bytes.len() as u8); expected.extend_from_slice(host_bytes); expected.push(b'\0'); expected.extend_from_slice(&0x0102030405060708u64.to_be_bytes()); expected.push(b'\0'); expected.extend_from_slice(b"did:web:example.com"); assert_eq!(k, expected); }
#[test] fn same_host_sorts_by_timestamp() { let earlier = host_key("pds.example.com", 100, &did("did:web:a.com")); let later = host_key("pds.example.com", 200, &did("did:web:a.com")); assert!(earlier < later); }
#[test] fn same_host_same_ts_sorts_by_did() { let a = host_key("pds.example.com", 100, &did("did:web:a.com")); let b = host_key("pds.example.com", 100, &did("did:web:b.com")); assert!(a < b); }
// --- host_key_parse roundtrips ---
#[test] fn host_key_parse_roundtrips() { let ts = 0xdeadbeefcafe1234u64; let d = did("did:web:example.com"); let host = "pds1.bsky.network"; let k = host_key(host, ts, &d); let (parsed_host, parsed_ts, parsed_did) = host_key_parse(&k).unwrap(); assert_eq!(parsed_host, host); assert_eq!(parsed_ts, ts); assert_eq!(parsed_did, d); }
#[test] fn host_key_parse_empty_host() { let ts = 42u64; let d = did("did:web:unknown.com"); let k = host_key("", ts, &d); let (parsed_host, parsed_ts, parsed_did) = host_key_parse(&k).unwrap(); assert_eq!(parsed_host, ""); assert_eq!(parsed_ts, ts); assert_eq!(parsed_did, d); }
#[test] fn host_key_parse_rejects_truncated() { assert!(host_key_parse(b"rsH").is_err()); }
// --- encode / decode ---
#[test] fn encode_decode_roundtrips() { let original = item("did:web:example.com", 3, "detected gap", &[0xAB, 0xCD]); let bytes = encode(&original); let decoded = decode(&bytes, "test-key", did("did:web:example.com")).unwrap(); assert_eq!(decoded.retry_count, 3); assert_eq!(decoded.retry_reason, "detected gap"); assert_eq!(decoded.commit_cbor, vec![0xAB, 0xCD]); }
#[test] fn encode_decode_empty_commit_cbor() { let original = item("did:web:example.com", 0, "first attempt", &[]); let bytes = encode(&original); let decoded = decode(&bytes, "test-key", did("did:web:example.com")).unwrap(); assert_eq!(decoded.retry_count, 0); assert_eq!(decoded.retry_reason, "first attempt"); assert!(decoded.commit_cbor.is_empty()); }
#[test] fn decode_rejects_truncated_header() { assert!(decode(&[0, 1, 2], "k", did("did:web:example.com")).is_err()); }
// --- enqueue and claim_from_host basic flow ---
#[test] fn enqueue_and_claim_basic() { let db = open_temporary().unwrap(); pending_repo(&db, "did:web:a.com"); enqueue( &db, 100, &item("did:web:a.com", 0, "backfill", &[1, 2, 3]), Some("pds.example.com"), ) .unwrap();
let claimed = claim_from_host(&db, "pds.example.com", 101, &HashSet::new()) .unwrap() .unwrap(); assert_eq!(claimed.item.did.as_str(), "did:web:a.com"); assert_eq!(claimed.item.retry_reason, "backfill"); assert_eq!(claimed.item.commit_cbor, vec![1, 2, 3]); assert_eq!(claimed.host, "pds.example.com");
// Queue should now be empty for this host. assert!( claim_from_host(&db, "pds.example.com", 9999, &HashSet::new()) .unwrap() .is_none() ); }
#[test] fn enqueue_with_none_host_uses_empty_string() { let db = open_temporary().unwrap(); pending_repo(&db, "did:web:a.com"); enqueue(&db, 100, &item("did:web:a.com", 0, "test", &[]), None).unwrap();
let claimed = claim_from_host(&db, "", 101, &HashSet::new()) .unwrap() .unwrap(); assert_eq!(claimed.item.did.as_str(), "did:web:a.com"); assert_eq!(claimed.host, ""); }
// --- next_host round-robins across hosts ---
#[test] fn next_host_round_robins() { let db = open_temporary().unwrap(); enqueue( &db, 100, &item("did:web:a.com", 0, "r1", &[]), Some("alpha.example.com"), ) .unwrap(); enqueue( &db, 100, &item("did:web:b.com", 0, "r2", &[]), Some("beta.example.com"), ) .unwrap(); enqueue( &db, 100, &item("did:web:c.com", 0, "r3", &[]), Some("gamma.example.com"), ) .unwrap();
let empty_skip: HashSet<String> = HashSet::new(); let h1 = next_host(&db, &empty_skip).unwrap(); let h2 = next_host(&db, &empty_skip).unwrap(); let h3 = next_host(&db, &empty_skip).unwrap(); let h4 = next_host(&db, &empty_skip).unwrap();
// BTreeSet order is alphabetical: alpha, beta, gamma assert_eq!(h1, "alpha.example.com"); assert_eq!(h2, "beta.example.com"); assert_eq!(h3, "gamma.example.com"); // Wraps around assert_eq!(h4, "alpha.example.com"); }
// --- next_host skips hosts in skip set ---
#[test] fn next_host_skips_hosts_in_skip_set() { let db = open_temporary().unwrap(); enqueue( &db, 100, &item("did:web:a.com", 0, "r1", &[]), Some("alpha.example.com"), ) .unwrap(); enqueue( &db, 100, &item("did:web:b.com", 0, "r2", &[]), Some("beta.example.com"), ) .unwrap(); enqueue( &db, 100, &item("did:web:c.com", 0, "r3", &[]), Some("gamma.example.com"), ) .unwrap();
let mut skip = HashSet::new(); skip.insert("alpha.example.com".to_owned()); skip.insert("gamma.example.com".to_owned());
let h = next_host(&db, &skip).unwrap(); assert_eq!(h, "beta.example.com"); }
#[test] fn next_host_returns_none_when_all_skipped() { let db = open_temporary().unwrap(); enqueue( &db, 100, &item("did:web:a.com", 0, "r1", &[]), Some("alpha.example.com"), ) .unwrap();
let mut skip = HashSet::new(); skip.insert("alpha.example.com".to_owned()); assert!(next_host(&db, &skip).is_none()); }
#[test] fn next_host_returns_none_on_empty_queue() { let db = open_temporary().unwrap(); let empty_skip: HashSet<String> = HashSet::new(); assert!(next_host(&db, &empty_skip).is_none()); }
// --- claim_from_host skips busy DIDs ---
#[test] fn claim_from_host_skips_busy_dids() { let db = open_temporary().unwrap(); pending_repo(&db, "did:web:a.com"); pending_repo(&db, "did:web:b.com"); enqueue( &db, 100, &item("did:web:a.com", 0, "first", &[]), Some("pds.example.com"), ) .unwrap(); enqueue( &db, 101, &item("did:web:b.com", 0, "second", &[]), Some("pds.example.com"), ) .unwrap();
let mut busy: HashSet<Did<'static>> = HashSet::new(); busy.insert(did("did:web:a.com"));
let claimed = claim_from_host(&db, "pds.example.com", 9999, &busy) .unwrap() .unwrap(); assert_eq!(claimed.item.did.as_str(), "did:web:b.com"); }
#[test] fn claim_from_host_returns_none_when_all_busy() { let db = open_temporary().unwrap(); pending_repo(&db, "did:web:a.com"); enqueue( &db, 100, &item("did:web:a.com", 0, "only", &[]), Some("pds.example.com"), ) .unwrap();
let mut busy: HashSet<Did<'static>> = HashSet::new(); busy.insert(did("did:web:a.com"));
assert!( claim_from_host(&db, "pds.example.com", 9999, &busy) .unwrap() .is_none() ); }
// --- claim_from_host respects timestamps ---
#[test] fn claim_from_host_respects_timestamps() { let db = open_temporary().unwrap(); pending_repo(&db, "did:web:a.com"); enqueue( &db, 200, &item("did:web:a.com", 0, "future", &[]), Some("pds.example.com"), ) .unwrap();
// now=200 means ts must be < 200 to be claimable. assert!( claim_from_host(&db, "pds.example.com", 200, &HashSet::new()) .unwrap() .is_none() ); assert!( claim_from_host(&db, "pds.example.com", 199, &HashSet::new()) .unwrap() .is_none() );
// now=201 should claim it. let claimed = claim_from_host(&db, "pds.example.com", 201, &HashSet::new()) .unwrap() .unwrap(); assert_eq!(claimed.item.retry_reason, "future"); }
#[test] fn claim_from_host_preserves_account_status() { let db = open_temporary().unwrap(); repo::put_info( &db, &did("did:web:a.com"), &repo::RepoInfo { state: repo::RepoState::Pending, status: repo::AccountStatus::Suspended, error: None, }, ) .unwrap(); enqueue( &db, 100, &item("did:web:a.com", 0, "backfill", &[]), Some("pds.example.com"), ) .unwrap();
claim_from_host(&db, "pds.example.com", 101, &HashSet::new()) .unwrap() .unwrap();
let (info, _) = repo::get(&db, &did("did:web:a.com")).unwrap().unwrap(); assert_eq!(info.status, repo::AccountStatus::Suspended); assert_eq!(info.state, repo::RepoState::Resyncing); }
// --- migrate_legacy ---
#[test] fn migrate_legacy_converts_old_entries() { let db = open_temporary().unwrap();
// Manually insert an old-format "rsq" entry. let old_did = did("did:web:old.com"); let ts = 42u64; let old_item = item("did:web:old.com", 1, "gap", &[0xFF]); let mut old_key = PREFIX_RESYNC_QUEUE.to_vec(); old_key.extend_from_slice(&ts.to_be_bytes()); old_key.push(b'\0'); old_key.extend_from_slice(old_did.as_str().as_bytes()); db.ks.insert(&old_key, encode(&old_item)).unwrap();
let count = migrate_legacy(&db).unwrap(); assert_eq!(count, 1);
// Old key should be gone. assert!(db.ks.get(&old_key).unwrap().is_none());
// Rebuild host index so next_host works. rebuild_host_index(&db).unwrap();
// Should be claimable from empty-string host. pending_repo(&db, "did:web:old.com"); let claimed = claim_from_host(&db, "", 9999, &HashSet::new()) .unwrap() .unwrap(); assert_eq!(claimed.item.did.as_str(), "did:web:old.com"); assert_eq!(claimed.item.retry_count, 1); assert_eq!(claimed.item.retry_reason, "gap"); assert_eq!(claimed.item.commit_cbor, vec![0xFF]); assert_eq!(claimed.host, ""); }
#[test] fn migrate_legacy_noop_when_no_old_entries() { let db = open_temporary().unwrap(); let count = migrate_legacy(&db).unwrap(); assert_eq!(count, 0); }
// --- rebuild_host_index ---
#[test] fn rebuild_host_index_from_disk() { let db = open_temporary().unwrap(); enqueue( &db, 100, &item("did:web:a.com", 0, "r1", &[]), Some("alpha.example.com"), ) .unwrap(); enqueue( &db, 101, &item("did:web:b.com", 0, "r2", &[]), Some("alpha.example.com"), ) .unwrap(); enqueue( &db, 102, &item("did:web:c.com", 0, "r3", &[]), Some("beta.example.com"), ) .unwrap();
// Clear the in-memory index. { let mut idx = db.host_index.lock().unwrap(); *idx = HostIndex::new(); }
rebuild_host_index(&db).unwrap();
let idx = db.host_index.lock().unwrap(); assert_eq!(idx.hosts.len(), 2); assert!(idx.hosts.contains("alpha.example.com")); assert!(idx.hosts.contains("beta.example.com")); assert_eq!(idx.host_counts["alpha.example.com"], 2); assert_eq!(idx.host_counts["beta.example.com"], 1); }
// --- count_queued ---
#[test] fn count_queued_counts_all_hosts() { let db = open_temporary().unwrap(); enqueue( &db, 100, &item("did:web:a.com", 0, "r1", &[]), Some("alpha.example.com"), ) .unwrap(); enqueue( &db, 101, &item("did:web:b.com", 0, "r2", &[]), Some("beta.example.com"), ) .unwrap(); enqueue(&db, 102, &item("did:web:c.com", 0, "r3", &[]), None).unwrap();
assert_eq!(count_queued(&db), 3); }}