use loom::sync::{Arc, Mutex, RwLock};
use super::explore;
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
struct Table {
id: u64,
key: u8,
}
const INPUT_A: Table = Table { id: 1, key: b'k' };
const INPUT_B: Table = Table { id: 2, key: b'k' };
const OUTPUT: Table = Table { id: 3, key: b'k' };
const FLUSHED: Table = Table { id: 4, key: b'w' };
#[derive(Clone)]
struct Version {
levels: [Vec<Table>; 2],
}
impl Version {
fn find(&self, key: u8) -> Option<Table> {
self.levels[0]
.iter()
.rev()
.chain(self.levels[1].iter())
.copied()
.find(|table| table.key == key)
}
}
enum Edit {
Add(usize, Table),
Remove(usize, u64),
}
struct VersionSet {
current: RwLock<Arc<Version>>,
}
impl VersionSet {
fn new(levels: [Vec<Table>; 2]) -> Self {
Self {
current: RwLock::new(Arc::new(Version { levels })),
}
}
fn current(&self) -> Arc<Version> {
Arc::clone(&self.current.read().expect("current"))
}
fn apply(&self, edits: &[Edit]) {
let mut version = (*self.current()).clone();
for edit in edits {
match edit {
Edit::Add(level, table) => version.levels[*level].push(*table),
Edit::Remove(level, id) => version.levels[*level].retain(|t| t.id != *id),
}
}
*self.current.write().expect("current") = Arc::new(version);
}
}
fn apply_locked(pipeline: &Mutex<()>, versions: &VersionSet, edits: &[Edit]) {
let _lock = pipeline.lock().expect("pipeline");
versions.apply(edits);
}
fn current_locked(pipeline: &Mutex<()>, versions: &VersionSet) -> Arc<Version> {
let _lock = pipeline.lock().expect("pipeline");
versions.current()
}
fn compaction_edits() -> [Edit; 3] {
[
Edit::Remove(0, INPUT_A.id),
Edit::Remove(0, INPUT_B.id),
Edit::Add(1, OUTPUT),
]
}
fn before_compaction() -> Arc<VersionSet> {
Arc::new(VersionSet::new([vec![INPUT_A, INPUT_B], Vec::new()]))
}
pub fn a_reader_pins_one_version_across_a_compaction() {
explore(
"a_reader_pins_one_version_across_a_compaction",
64,
8,
|witness| {
let versions = before_compaction();
let pipeline = Arc::new(Mutex::new(()));
let compaction = {
let versions = Arc::clone(&versions);
let pipeline = Arc::clone(&pipeline);
loom::thread::spawn(move || {
apply_locked(&pipeline, &versions, &compaction_edits());
})
};
let reader = {
let versions = Arc::clone(&versions);
let pipeline = Arc::clone(&pipeline);
let witness = witness.clone();
loom::thread::spawn(move || {
let version = current_locked(&pipeline, &versions);
let found = version
.find(b'k')
.expect("the compaction handoff hid a key that was never deleted");
if found.id == OUTPUT.id {
witness.record();
}
})
};
compaction.join().expect("compaction");
reader.join().expect("reader");
let version = versions.current();
assert!(version.levels[0].is_empty(), "both inputs were removed");
assert_eq!(version.levels[1], vec![OUTPUT]);
},
);
}
pub fn a_split_compaction_apply_hides_a_key() {
explore("a_split_compaction_apply_hides_a_key", 4, 1, |witness| {
let versions = before_compaction();
let pipeline = Arc::new(Mutex::new(()));
let compaction = {
let versions = Arc::clone(&versions);
let pipeline = Arc::clone(&pipeline);
loom::thread::spawn(move || {
apply_locked(
&pipeline,
&versions,
&[Edit::Remove(0, INPUT_A.id), Edit::Remove(0, INPUT_B.id)],
);
apply_locked(&pipeline, &versions, &[Edit::Add(1, OUTPUT)]);
})
};
let reader = {
let versions = Arc::clone(&versions);
let pipeline = Arc::clone(&pipeline);
let witness = witness.clone();
loom::thread::spawn(move || {
let version = current_locked(&pipeline, &versions);
witness.record();
assert!(version.find(b'k').is_some(), "the split apply hid `k`");
})
};
compaction.join().expect("compaction");
reader.join().expect("reader");
});
}
pub fn a_flush_and_a_compaction_cannot_lose_each_other() {
explore(
"a_flush_and_a_compaction_cannot_lose_each_other",
64,
8,
|witness| {
let versions = before_compaction();
let pipeline = Arc::new(Mutex::new(()));
let flush = {
let versions = Arc::clone(&versions);
let pipeline = Arc::clone(&pipeline);
let witness = witness.clone();
loom::thread::spawn(move || {
if current_locked(&pipeline, &versions).levels[1].contains(&OUTPUT) {
witness.record();
}
apply_locked(&pipeline, &versions, &[Edit::Add(0, FLUSHED)]);
})
};
let compaction = {
let versions = Arc::clone(&versions);
let pipeline = Arc::clone(&pipeline);
loom::thread::spawn(move || {
apply_locked(&pipeline, &versions, &compaction_edits());
})
};
flush.join().expect("flush");
compaction.join().expect("compaction");
let version = versions.current();
assert_eq!(
version.levels[0],
vec![FLUSHED],
"the flush's table was lost by the compaction's version swap"
);
assert_eq!(
version.levels[1],
vec![OUTPUT],
"the compaction's output was lost by the flush's version swap"
);
},
);
}
pub fn an_unserialized_version_swap_loses_an_edit() {
explore(
"an_unserialized_version_swap_loses_an_edit",
4,
1,
|witness| {
let versions = before_compaction();
let flush = {
let versions = Arc::clone(&versions);
loom::thread::spawn(move || {
versions.apply(&[Edit::Add(0, FLUSHED)]);
})
};
let compaction = {
let versions = Arc::clone(&versions);
loom::thread::spawn(move || {
versions.apply(&compaction_edits());
})
};
flush.join().expect("flush");
compaction.join().expect("compaction");
witness.record();
let version = versions.current();
assert!(
version.levels[0] == vec![FLUSHED] && version.levels[1] == vec![OUTPUT],
"an unserialized version swap lost an edit"
);
},
);
}