//! Two-level host-partitioned resync queue. //! //! Keys: `"rsH" \0 \0 ` //! 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, host_counts: HashMap, 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" \0 \0 ` fn host_key(host: &str, ts: u64, did: &Did<'_>) -> Vec { 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" \0` fn host_prefix(host: &str) -> Vec { 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 { PREFIX_RESYNC_HOST_QUEUE.to_vec() } /// Upper bound for timestamp within a host prefix (exclusive): /// `"rsH" \0 \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 { 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, } fn encode(item: &ResyncItem) -> Vec { 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 { 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) -> Option { 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>, ) -> StorageResult> { 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 { 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 = 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 = 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 = 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> = 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> = 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); } }