mirror of
https://tangled.org/tranquil.farm/tranquil-pds
synced 2026-07-20 15:02:37 +00:00
test(store): cross-store, firehose, read-validation coverage w/ faults
Lewis: May this revision serve well! <lu5a@proton.me>
This commit is contained in:
@@ -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<String> = 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::<Vec<_>>()
|
||||
.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<String> = 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::<Vec<_>>()
|
||||
.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 {
|
||||
|
||||
@@ -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<RealIO, SystemClock> {
|
||||
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<EventLog<RealIO>> {
|
||||
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<ImportBlock>, Vec<ImportRecord>) {
|
||||
let block_count = ((seed % 12) + 4) as u32;
|
||||
let collection = Nsid::from("app.bsky.feed.post".to_string());
|
||||
let blocks: Vec<ImportBlock> = (0..block_count)
|
||||
.map(|i| ImportBlock {
|
||||
cid_bytes: test_cid(i).to_vec(),
|
||||
data: block_data(i),
|
||||
})
|
||||
.collect();
|
||||
let records: Vec<ImportRecord> = (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<u8>)> = 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<std::sync::Arc<SimulatedIO>, SimClock>,
|
||||
eventlog: std::sync::Arc<EventLog<std::sync::Arc<SimulatedIO>>>,
|
||||
metastore: Metastore,
|
||||
}
|
||||
|
||||
fn open_sim_stores(
|
||||
dir: &std::path::Path,
|
||||
sim: &std::sync::Arc<SimulatedIO>,
|
||||
clock: &SimClock,
|
||||
) -> SimStores {
|
||||
let blockstore = {
|
||||
let factory = std::sync::Arc::clone(sim);
|
||||
TranquilBlockStore::<std::sync::Arc<SimulatedIO>, 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<Rkey> = block_ids
|
||||
.iter()
|
||||
.map(|i| Rkey::from(format!("3k{i:08}")))
|
||||
.collect();
|
||||
let links: Vec<CidLink> = block_ids.iter().map(|&i| block_cid_link(i)).collect();
|
||||
let writes: Vec<RecordWrite<'_>> = 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<u32> = 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<Uuid> = (0..repo_count)
|
||||
.map(|i| test_uuid(seed * 100 + i as u64))
|
||||
.collect();
|
||||
|
||||
let root_pairs: Vec<(CidBytes, Vec<u8>)> = (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<u32> = (start..start + blocks_per_round).collect();
|
||||
let pairs: Vec<(CidBytes, Vec<u8>)> = 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<SimulatedIO>,
|
||||
clock: &SimClock,
|
||||
) -> TranquilBlockStore<std::sync::Arc<SimulatedIO>, SimClock> {
|
||||
let factory = std::sync::Arc::clone(sim);
|
||||
TranquilBlockStore::<std::sync::Arc<SimulatedIO>, 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<u32> = (base..base + 5).collect();
|
||||
let victim: Vec<u32> = (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}"
|
||||
);
|
||||
}
|
||||
|
||||
@@ -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<SimulatedIO>,
|
||||
) -> Arc<EventLog<Arc<SimulatedIO>>> {
|
||||
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<Arc<SimulatedIO>>, 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<Arc<SimulatedIO>>, batch: usize) -> Result<Vec<u64>, String> {
|
||||
let mut cursor = EventSequence::BEFORE_ALL;
|
||||
let mut seqs: Vec<u64> = 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<u64> = (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<u8>)> = (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"
|
||||
);
|
||||
});
|
||||
});
|
||||
});
|
||||
}
|
||||
Reference in New Issue
Block a user