From 54244075ccce33fcb18eb9019eec11147d4718dd Mon Sep 17 00:00:00 2001 From: Lewis Date: Tue, 2 Jun 2026 11:51:20 +0300 Subject: [PATCH] test(store): cross-store, firehose, read-validation coverage w/ faults Lewis: May this revision serve well! --- crates/tranquil-store/tests/gauntlet_smoke.rs | 80 +++ .../tranquil-store/tests/sim_cross_store.rs | 507 +++++++++++++++++- crates/tranquil-store/tests/sim_firehose.rs | 196 +++++++ 3 files changed, 778 insertions(+), 5 deletions(-) create mode 100644 crates/tranquil-store/tests/sim_firehose.rs diff --git a/crates/tranquil-store/tests/gauntlet_smoke.rs b/crates/tranquil-store/tests/gauntlet_smoke.rs index b0a7b9d..e5055d6 100644 --- a/crates/tranquil-store/tests/gauntlet_smoke.rs +++ b/crates/tranquil-store/tests/gauntlet_smoke.rs @@ -109,6 +109,86 @@ async fn full_stack_restart_port() { ); } +#[tokio::test] +async fn read_corruption_detects_no_silent_mismatch() { + let reports = farm::run_many( + |seed| { + let mut cfg = config_for(Scenario::ReadCorruption, seed); + ConfigOverrides { + op_count: Some(4_000), + ..ConfigOverrides::default() + } + .apply_to(&mut cfg); + cfg + }, + (0..4).map(Seed), + ); + let failures: Vec = reports + .iter() + .filter(|r| !r.is_clean()) + .map(|r| { + format!( + "seed {}: {}", + r.seed.0, + r.violations + .iter() + .map(|v| format!("{}: {}", v.invariant, v.detail)) + .collect::>() + .join("; ") + ) + }) + .collect(); + assert!(failures.is_empty(), "{}", failures.join("\n---\n")); +} + +#[tokio::test] +async fn mst_list_walk_validates_under_recoverable_faults() { + let reports = farm::run_many( + |seed| { + let mut cfg = config_for(Scenario::BlockChurnRecoverable, seed); + ConfigOverrides { + op_count: Some(4_000), + ..ConfigOverrides::default() + } + .apply_to(&mut cfg); + cfg + }, + (0..6).map(Seed), + ); + let failures: Vec = reports + .iter() + .filter(|r| !r.is_clean()) + .map(|r| { + format!( + "seed {}: {}", + r.seed.0, + r.violations + .iter() + .map(|v| format!("{}: {}", v.invariant, v.detail)) + .collect::>() + .join("; ") + ) + }) + .collect(); + assert!(failures.is_empty(), "{}", failures.join("\n---\n")); +} + +#[tokio::test] +async fn inline_commit_synchronous_path_sanity() { + let mut cfg = config_for(Scenario::InlineCommit, Seed(5)); + ConfigOverrides { + op_count: Some(2_000), + ..ConfigOverrides::default() + } + .apply_to(&mut cfg); + let report = Gauntlet::new(cfg).expect("build gauntlet").run().await; + assert_clean(&report); + assert!( + report.restarts.0 >= 1, + "InlineCommit must exercise at least one restart of the synchronous commit path" + ); +} + #[tokio::test] async fn compaction_idempotent_sanity() { let cfg = GauntletConfig { diff --git a/crates/tranquil-store/tests/sim_cross_store.rs b/crates/tranquil-store/tests/sim_cross_store.rs index a9cb6c4..b6801ba 100644 --- a/crates/tranquil-store/tests/sim_cross_store.rs +++ b/crates/tranquil-store/tests/sim_cross_store.rs @@ -4,21 +4,24 @@ use std::ops::ControlFlow; use std::sync::atomic::{AtomicBool, Ordering}; use rayon::prelude::*; -use tranquil_store::RealIO; use tranquil_store::backup::{BackupCoordinator, restore_from_backup, verify_backup}; use tranquil_store::blockstore::{ BlockStoreConfig, CidBytes, DEFAULT_MAX_FILE_SIZE, GroupCommitConfig, TranquilBlockStore, }; -use tranquil_store::eventlog::{EventLog, EventLogConfig, EventSequence}; +use tranquil_store::consistency::{ConsistencyCheckOptions, verify_store_consistency_with_options}; +use tranquil_store::eventlog::{EventLog, EventLogBridge, EventLogConfig, EventSequence}; +use tranquil_store::metastore::record_ops::RecordWrite; use tranquil_store::metastore::{Metastore, MetastoreConfig}; -use tranquil_store::{sim_seed_range, sim_single_seed}; +use tranquil_store::{RealIO, SimClock, SimulatedIO, SystemClock, sim_seed_range, sim_single_seed}; use common::{ TestStores, assert_store_consistent, block_data, open_test_stores, test_cid, test_cid_link, test_did, test_handle, test_uuid, with_runtime, }; -use tranquil_db_traits::{RepoEventType, SequenceNumber, SequencedEvent}; -use tranquil_types::Did; +use tranquil_db_traits::{ + ImportBlock, ImportRecord, RepoEventType, SequenceNumber, SequencedEvent, +}; +use tranquil_types::{CidLink, Did, Nsid, Rkey}; use uuid::Uuid; const CACHE_SIZE: u64 = 16 * 1024 * 1024; @@ -584,3 +587,497 @@ fn sim_backup_during_concurrent_block_and_event_writes() { }); }); } + +fn import_blockstore(dir: &std::path::Path) -> TranquilBlockStore { + TranquilBlockStore::open(BlockStoreConfig { + data_dir: dir.join("blockstore/data"), + index_dir: dir.join("blockstore/index"), + max_file_size: DEFAULT_MAX_FILE_SIZE, + group_commit: GroupCommitConfig::default(), + shard_count: 1, + }) + .unwrap() +} + +fn import_metastore(dir: &std::path::Path) -> Metastore { + Metastore::open( + &dir.join("metastore"), + MetastoreConfig { + cache_size_bytes: CACHE_SIZE, + }, + ) + .unwrap() +} + +fn import_eventlog(dir: &std::path::Path) -> std::sync::Arc> { + std::sync::Arc::new( + EventLog::open( + EventLogConfig { + segments_dir: dir.join("eventlog/segments"), + ..EventLogConfig::default() + }, + RealIO::new(), + ) + .unwrap(), + ) +} + +fn create_import_dirs(dir: &std::path::Path) { + [ + dir.join("blockstore/data"), + dir.join("blockstore/index"), + dir.join("eventlog/segments"), + dir.join("metastore"), + ] + .iter() + .for_each(|d| std::fs::create_dir_all(d).unwrap()); +} + +fn import_fixture(seed: u64) -> (Vec, Vec) { + let block_count = ((seed % 12) + 4) as u32; + let collection = Nsid::from("app.bsky.feed.post".to_string()); + let blocks: Vec = (0..block_count) + .map(|i| ImportBlock { + cid_bytes: test_cid(i).to_vec(), + data: block_data(i), + }) + .collect(); + let records: Vec = (0..block_count) + .map(|i| { + let cid = cid::Cid::try_from(&test_cid(i)[..]).unwrap(); + ImportRecord { + collection: collection.clone(), + rkey: Rkey::from(format!("3kimport{i:04}")), + record_cid: CidLink::from_cid(&cid), + } + }) + .collect(); + (blocks, records) +} + +fn assert_import_complete( + dir: &std::path::Path, + user_id: Uuid, + blocks: &[ImportBlock], + records: &[ImportRecord], + ctx: &str, +) { + let bs = import_blockstore(dir); + blocks.iter().for_each(|b| { + let cid: CidBytes = b.cid_bytes.as_slice().try_into().unwrap(); + assert!( + bs.get_block_sync(&cid).unwrap().is_some(), + "{ctx}: imported block must be durable" + ); + }); + let ms = import_metastore(dir); + records.iter().enumerate().for_each(|(i, r)| { + let got = ms + .record_ops() + .get_record_cid(user_id, &r.collection, &r.rkey) + .unwrap(); + assert_eq!( + got.as_ref(), + Some(&r.record_cid), + "{ctx}: imported record {i} must reference its block after recovery" + ); + }); + let el = import_eventlog(dir); + let report = verify_store_consistency_with_options( + &bs, + &ms, + &el, + ConsistencyCheckOptions { + check_block_references: true, + ..ConsistencyCheckOptions::default() + }, + ); + assert!( + report.dangling_record_cids.is_empty(), + "{ctx}: every imported record must reference a present block: {report}" + ); +} + +fn run_import_crash_scenario(seed: u64) { + let dir = tempfile::TempDir::new().unwrap(); + create_import_dirs(dir.path()); + let user_id = test_uuid(seed); + let did = test_did(seed); + let handle = test_handle(seed); + let root_cid = test_cid_link((seed & 0xFF) as u8); + let (blocks, records) = import_fixture(seed); + + { + let ms = import_metastore(dir.path()); + ms.repo_ops() + .create_repo(ms.database(), user_id, &did, &handle, &root_cid, "rev0") + .unwrap(); + ms.persist().unwrap(); + } + + if seed.is_multiple_of(2) { + { + let bs = import_blockstore(dir.path()); + let pairs: Vec<(CidBytes, Vec)> = blocks + .iter() + .map(|b| (b.cid_bytes.as_slice().try_into().unwrap(), b.data.clone())) + .collect(); + bs.put_blocks_blocking(pairs).unwrap(); + } + + { + let bs = import_blockstore(dir.path()); + blocks.iter().for_each(|b| { + let cid: CidBytes = b.cid_bytes.as_slice().try_into().unwrap(); + assert!( + bs.get_block_sync(&cid).unwrap().is_some(), + "seed={seed} blocks are durable after put, before metastore commit" + ); + }); + let ms = import_metastore(dir.path()); + records.iter().for_each(|r| { + assert!( + ms.record_ops() + .get_record_cid(user_id, &r.collection, &r.rkey) + .unwrap() + .is_none(), + "seed={seed} no import record may exist before the metastore commit lands" + ); + }); + let bridge = EventLogBridge::new(import_eventlog(dir.path())); + let commit_ops = ms + .commit_ops(std::sync::Arc::new(bridge)) + .with_blockstore(bs); + commit_ops + .import_repo_data(user_id, &blocks, &records, Some(&root_cid)) + .unwrap(); + ms.persist().unwrap(); + } + } else { + let ms = import_metastore(dir.path()); + let bs = import_blockstore(dir.path()); + let bridge = EventLogBridge::new(import_eventlog(dir.path())); + let commit_ops = ms + .commit_ops(std::sync::Arc::new(bridge)) + .with_blockstore(bs); + commit_ops + .import_repo_data(user_id, &blocks, &records, Some(&root_cid)) + .unwrap(); + ms.persist().unwrap(); + } + + assert_import_complete( + dir.path(), + user_id, + &blocks, + &records, + &format!("seed={seed} import recovery"), + ); +} + +#[test] +fn sim_import_repo_data_crash_during_import_recovers() { + with_runtime(|| { + sim_seed_range().into_par_iter().for_each(|seed| { + run_import_crash_scenario(seed); + }); + }); +} + +struct SimStores { + blockstore: TranquilBlockStore, SimClock>, + eventlog: std::sync::Arc>>, + metastore: Metastore, +} + +fn open_sim_stores( + dir: &std::path::Path, + sim: &std::sync::Arc, + clock: &SimClock, +) -> SimStores { + let blockstore = { + let factory = std::sync::Arc::clone(sim); + TranquilBlockStore::, SimClock>::open_with_io( + BlockStoreConfig { + data_dir: dir.join("blockstore/data"), + index_dir: dir.join("blockstore/index"), + max_file_size: 16 * 1024, + group_commit: GroupCommitConfig { + synchronous: true, + verify_persisted_blocks: true, + ..GroupCommitConfig::default() + }, + shard_count: 1, + }, + move || std::sync::Arc::clone(&factory), + clock.clone(), + ) + .unwrap() + }; + let eventlog = std::sync::Arc::new( + EventLog::open( + EventLogConfig { + segments_dir: dir.join("eventlog/segments"), + ..EventLogConfig::default() + }, + std::sync::Arc::clone(sim), + ) + .unwrap(), + ); + let metastore = import_metastore(dir); + SimStores { + blockstore, + eventlog, + metastore, + } +} + +fn sim_fault_for(seed: u64) -> tranquil_store::FaultConfig { + match seed % 3 { + 0 => tranquil_store::FaultConfig::recoverable(), + 1 => tranquil_store::FaultConfig::torn_pages_only(), + _ => tranquil_store::FaultConfig::fsyncgate_only(), + } +} + +fn block_cid_link(i: u32) -> CidLink { + let cid = cid::Cid::try_from(&test_cid(i)[..]).expect("test_cid is a valid cid"); + CidLink::from_cid(&cid) +} + +fn commit_block_records(stores: &SimStores, uid: Uuid, collection: &Nsid, block_ids: &[u32]) { + let Some(user_hash) = stores.metastore.user_hashes().get(&uid) else { + return; + }; + let rkeys: Vec = block_ids + .iter() + .map(|i| Rkey::from(format!("3k{i:08}"))) + .collect(); + let links: Vec = block_ids.iter().map(|&i| block_cid_link(i)).collect(); + let writes: Vec> = rkeys + .iter() + .zip(links.iter()) + .map(|(rkey, cid)| RecordWrite { + collection, + rkey, + cid, + }) + .collect(); + let mut batch = stores.metastore.database().batch(); + stores + .metastore + .record_ops() + .upsert_records(&mut batch, user_hash, &writes) + .unwrap(); + batch.commit().unwrap(); + stores.metastore.persist().unwrap(); +} + +fn run_cross_store_fault_scenario(seed: u64) { + let dir = tempfile::TempDir::new().unwrap(); + create_import_dirs(dir.path()); + let fault = sim_fault_for(seed); + let sim = std::sync::Arc::new(SimulatedIO::new(seed, fault)); + let clock = sim.clock(); + + let repo_count = ((seed % 4) + 1) as u32; + let blocks_per_round = ((seed % 6) + 3) as u32; + let rounds = ((seed % 4) + 2) as u32; + let collection = Nsid::from("app.bsky.feed.post".to_string()); + let root_base = (seed as u32).wrapping_mul(100_000); + + let mut acked_blocks: Vec = Vec::new(); + let mut committed_repos: Vec<(u64, Uuid)> = Vec::new(); + + sim.set_pristine_mode(true); + let stores = open_sim_stores(dir.path(), &sim, &clock); + + let repo_uids: Vec = (0..repo_count) + .map(|i| test_uuid(seed * 100 + i as u64)) + .collect(); + + let root_pairs: Vec<(CidBytes, Vec)> = (0..repo_count) + .map(|i| (test_cid(root_base + i), block_data(root_base + i))) + .collect(); + stores.blockstore.put_blocks_blocking(root_pairs).unwrap(); + acked_blocks.extend((0..repo_count).map(|i| root_base + i)); + + (0..repo_count).for_each(|i| { + let idx = seed * 100 + i as u64; + let uid = repo_uids[i as usize]; + let did = test_did(idx); + let handle = test_handle(idx); + let root = block_cid_link(root_base + i); + stores + .metastore + .repo_ops() + .create_repo( + stores.metastore.database(), + uid, + &did, + &handle, + &root, + &format!("rev{idx}"), + ) + .unwrap(); + committed_repos.push((idx, uid)); + }); + stores.metastore.persist().unwrap(); + + sim.set_pristine_mode(false); + let mut next_block: u32 = root_base.wrapping_add(1000); + (0..rounds).for_each(|round| { + let start = next_block; + next_block = next_block.wrapping_add(blocks_per_round); + let block_ids: Vec = (start..start + blocks_per_round).collect(); + let pairs: Vec<(CidBytes, Vec)> = block_ids + .iter() + .map(|&i| (test_cid(i), block_data(i))) + .collect(); + if stores.blockstore.put_blocks_blocking(pairs).is_ok() { + acked_blocks.extend(block_ids.iter().copied()); + let uid = repo_uids[(round % repo_count) as usize]; + commit_block_records(&stores, uid, &collection, &block_ids); + } + }); + + sim.set_pristine_mode(true); + sim.crash(); + drop(stores); + + let stores = open_sim_stores(dir.path(), &sim, &clock); + + acked_blocks.iter().for_each(|&i| { + let got = stores.blockstore.get_block_sync(&test_cid(i)).unwrap(); + assert!( + got.is_some(), + "seed={seed} fault={fault:?}: acked block {i} must survive fault+crash" + ); + assert_eq!( + &got.unwrap()[..4], + &i.to_le_bytes(), + "seed={seed} fault={fault:?}: acked block {i} content mismatch after recovery" + ); + }); + + committed_repos.iter().for_each(|&(idx, uid)| { + assert!( + stores + .metastore + .repo_ops() + .get_repo_meta(uid) + .unwrap() + .is_some(), + "seed={seed} fault={fault:?}: persisted repo {idx} must survive" + ); + }); + + let report = verify_store_consistency_with_options( + &stores.blockstore, + &stores.metastore, + &stores.eventlog, + ConsistencyCheckOptions::default(), + ); + assert!( + report.is_consistent(), + "seed={seed} fault={fault:?}: every surviving record and repo root must reference a present block after fault+crash: {report}" + ); +} + +#[test] +fn sim_cross_store_coordinated_commit_under_faults() { + sim_seed_range().into_par_iter().for_each(|seed| { + run_cross_store_fault_scenario(seed); + }); +} + +fn open_sim_blockstore( + dir: &std::path::Path, + sim: &std::sync::Arc, + clock: &SimClock, +) -> TranquilBlockStore, SimClock> { + let factory = std::sync::Arc::clone(sim); + TranquilBlockStore::, SimClock>::open_with_io( + BlockStoreConfig { + data_dir: dir.join("blockstore/data"), + index_dir: dir.join("blockstore/index"), + max_file_size: 64 * 1024, + group_commit: GroupCommitConfig { + synchronous: true, + verify_persisted_blocks: true, + ..GroupCommitConfig::default() + }, + shard_count: 1, + }, + move || std::sync::Arc::clone(&factory), + clock.clone(), + ) + .unwrap() +} + +fn run_midcommit_crash_atomicity(seed: u64) -> bool { + let dir = tempfile::TempDir::new().unwrap(); + create_import_dirs(dir.path()); + let sim = std::sync::Arc::new(SimulatedIO::new(seed, tranquil_store::FaultConfig::none())); + let clock = sim.clock(); + let base = (seed as u32).wrapping_mul(10_000); + let baseline: Vec = (base..base + 5).collect(); + let victim: Vec = (base + 5..base + 5 + 12).collect(); + let after_writes = ((seed % 8) + 1) as i64; + + { + let bs = open_sim_blockstore(dir.path(), &sim, &clock); + bs.put_blocks_blocking( + baseline + .iter() + .map(|&i| (test_cid(i), block_data(i))) + .collect(), + ) + .unwrap(); + sim.arm_write_crash(after_writes); + let _ = bs.put_blocks_blocking( + victim + .iter() + .map(|&i| (test_cid(i), block_data(i))) + .collect(), + ); + sim.crash(); + drop(bs); + } + + let bs = open_sim_blockstore(dir.path(), &sim, &clock); + baseline.iter().for_each(|&i| { + assert!( + bs.get_block_sync(&test_cid(i)).unwrap().is_some(), + "seed={seed} baseline block {i} committed before the mid-commit crash must survive" + ); + }); + let present = victim + .iter() + .filter(|&&i| bs.get_block_sync(&test_cid(i)).unwrap().is_some()) + .count(); + assert!( + present == 0 || present == victim.len(), + "seed={seed} multi-block commit must be atomic across a mid-commit crash: {present}/{} victim blocks present", + victim.len() + ); + present == 0 +} + +#[test] +fn sim_blockstore_atomic_commit_under_midcommit_crash() { + let interrupted = std::sync::atomic::AtomicU32::new(0); + let total = std::sync::atomic::AtomicU32::new(0); + sim_seed_range().into_par_iter().for_each(|seed| { + total.fetch_add(1, Ordering::Relaxed); + if run_midcommit_crash_atomicity(seed) { + interrupted.fetch_add(1, Ordering::Relaxed); + } + }); + let interrupted = interrupted.load(Ordering::Relaxed); + let total = total.load(Ordering::Relaxed); + assert!( + interrupted * 2 >= total, + "mid-commit crash must actually interrupt the commit in most seeds, otherwise this test has no teeth: interrupted {interrupted}/{total}" + ); +} diff --git a/crates/tranquil-store/tests/sim_firehose.rs b/crates/tranquil-store/tests/sim_firehose.rs new file mode 100644 index 0000000..c1a8e21 --- /dev/null +++ b/crates/tranquil-store/tests/sim_firehose.rs @@ -0,0 +1,196 @@ +mod common; + +use std::sync::Arc; +use std::sync::atomic::{AtomicBool, AtomicU32, Ordering}; + +use rayon::prelude::*; +use tranquil_store::eventlog::{EventLog, EventLogConfig, EventSequence}; +use tranquil_store::{FaultConfig, SimulatedIO, sim_seed_range}; + +use common::with_runtime; +use tranquil_db_traits::{RepoEventType, SequenceNumber, SequencedEvent}; +use tranquil_types::Did; + +fn open_sim_eventlog( + dir: &std::path::Path, + sim: &Arc, +) -> Arc>> { + Arc::new( + EventLog::open( + EventLogConfig { + segments_dir: dir.join("segments"), + max_segment_size: 4096, + ..EventLogConfig::default() + }, + Arc::clone(sim), + ) + .unwrap(), + ) +} + +fn append_seq(el: &EventLog>, idx: u32) { + let did = Did::from(format!("did:plc:firehose{}", idx % 16)); + let event = SequencedEvent { + seq: SequenceNumber::from_raw(0), + did: did.clone(), + created_at: chrono::Utc::now(), + event_type: RepoEventType::Commit, + commit_cid: None, + prev_cid: None, + prev_data_cid: None, + ops: None, + blobs: None, + blocks: None, + handle: None, + active: None, + status: None, + rev: Some(format!("rev{idx}")), + }; + el.append_event(&did, RepoEventType::Commit, &event) + .unwrap(); +} + +fn drain_replay(el: &EventLog>, batch: usize) -> Result, String> { + let mut cursor = EventSequence::BEFORE_ALL; + let mut seqs: Vec = Vec::new(); + loop { + let events = el + .get_events_since(cursor, batch) + .map_err(|e| e.to_string())?; + if events.is_empty() { + return Ok(seqs); + } + events.iter().for_each(|e| { + let raw = e.seq.as_u64().expect("event seq is non-negative"); + seqs.push(raw); + cursor = EventSequence::new(raw); + }); + } +} + +#[test] +fn sim_firehose_replay_recovers_after_crash() { + sim_seed_range().into_par_iter().for_each(|seed| { + let dir = tempfile::TempDir::new().unwrap(); + std::fs::create_dir_all(dir.path().join("segments")).unwrap(); + let n = ((seed % 400) + 50) as u32; + + let sim = Arc::new(SimulatedIO::new(seed, FaultConfig::read_corruption())); + sim.set_pristine_mode(true); + + let synced: u64 = { + let el = open_sim_eventlog(dir.path(), &sim); + (1..=n).for_each(|i| append_seq(&el, i)); + let result = el.sync().unwrap(); + let synced = result.synced_through.raw(); + el.shutdown().unwrap(); + synced + }; + assert_eq!(synced, n as u64, "seed={seed} all appends must sync clean"); + + sim.crash(); + + sim.set_pristine_mode(true); + let el = open_sim_eventlog(dir.path(), &sim); + let expected: Vec = (1..=synced).collect(); + + let clean = drain_replay(&el, 64).unwrap(); + assert_eq!( + clean, expected, + "seed={seed} clean replay after crash must return every synced event in order" + ); + + sim.set_pristine_mode(false); + if let Ok(faulted) = drain_replay(&el, 32) { + assert!( + faulted.windows(2).all(|w| w[1] > w[0]), + "seed={seed} replayed seqs must be strictly increasing under read corruption" + ); + assert!( + faulted.iter().all(|&s| s >= 1 && s <= synced), + "seed={seed} replay must never fabricate a seq outside the synced range" + ); + assert!( + faulted.len() as u64 <= synced, + "seed={seed} replay under corruption must never return more events than were synced" + ); + } + + sim.set_pristine_mode(true); + let recovered = drain_replay(&el, 64).unwrap(); + assert_eq!( + recovered, expected, + "seed={seed} after read faults clear, full replay must be intact" + ); + el.shutdown().unwrap(); + }); +} + +#[test] +fn sim_blockstore_quiesce_under_concurrent_writers_and_reopen() { + use tranquil_store::blockstore::{CidBytes, TranquilBlockStore}; + + with_runtime(|| { + let seed_range = match tranquil_store::sim_single_seed() { + Some(s) => s..s + 1, + None => 0..64u64, + }; + seed_range.into_par_iter().for_each(|seed| { + let dir = tempfile::TempDir::new().unwrap(); + let cfg = common::default_blockstore_config(dir.path()); + let acked = Arc::new(AtomicU32::new(0)); + let base = (seed as u32).wrapping_mul(100_000); + + { + let store = TranquilBlockStore::open(cfg.clone()).unwrap(); + let stop = AtomicBool::new(false); + + std::thread::scope(|s| { + let writer = s.spawn(|| { + let mut next = base; + while !stop.load(Ordering::Relaxed) { + let batch: Vec<(CidBytes, Vec)> = (next..next + 4) + .map(|i| (common::test_cid(i), common::block_data(i))) + .collect(); + if store.put_blocks_blocking(batch).is_ok() { + acked.fetch_max(next + 4, Ordering::Relaxed); + } + next += 4; + } + }); + + let mut waited_ms = 0; + while acked.load(Ordering::Relaxed) <= base && waited_ms < 5_000 { + std::thread::sleep(std::time::Duration::from_millis(1)); + waited_ms += 1; + } + let (_snapshot, guard) = store.quiesce().unwrap(); + let acked_at_quiesce = acked.load(Ordering::Relaxed); + assert!( + acked_at_quiesce > base, + "seed={seed} writer must ack at least one batch before quiesce, otherwise the snapshot assertion has no teeth" + ); + (base..acked_at_quiesce).for_each(|i| { + assert!( + store.get_block_sync(&common::test_cid(i)).unwrap().is_some(), + "seed={seed} block {i} acked before quiesce must be readable in the quiesced snapshot" + ); + }); + drop(guard); + std::thread::sleep(std::time::Duration::from_millis(20)); + stop.store(true, Ordering::Relaxed); + writer.join().unwrap(); + }); + } + + let store = TranquilBlockStore::open(cfg).unwrap(); + let acked_count = acked.load(Ordering::Relaxed); + (base..acked_count).for_each(|i| { + assert!( + store.get_block_sync(&common::test_cid(i)).unwrap().is_some(), + "seed={seed} acked block {i} must survive quiesce + clean reopen" + ); + }); + }); + }); +}