use loom::sync::atomic::{AtomicU64, Ordering};
use loom::sync::{Arc, Mutex, RwLock};
use super::super::internal_key::user_key_of;
use super::super::lookup_key::LookupKey;
use super::super::memtable::MemTable;
use super::super::read_horizon::ReadHorizon;
use super::{explore, memtable, probe};
const SEQ: u64 = 7;
fn snapshot(user_key: &[u8], snapshot_seq: u64) -> LookupKey {
LookupKey::new(0, user_key, snapshot_seq)
}
type InstalledTable = Vec<(Vec<u8>, Vec<u8>)>;
fn table_holds(table: &InstalledTable, user_key: &[u8]) -> Option<Vec<u8>> {
table
.iter()
.find(|(key, _)| user_key_of(key) == user_key)
.map(|(_, value)| value.clone())
}
fn install(memtable: &MemTable) -> InstalledTable {
let mut table = Vec::new();
memtable
.try_for_each_entry(|key, value| {
table.push((key.to_vec(), value.to_vec()));
Ok(())
})
.expect("collecting into a vector cannot fail");
table
}
#[derive(Debug, PartialEq, Eq)]
enum Source {
Active,
Frozen,
Version,
}
fn read(
active: &MemTable,
frozen: &RwLock<Vec<Arc<MemTable>>>,
version: &Mutex<Vec<InstalledTable>>,
user_key: &[u8],
) -> Option<(Source, Vec<u8>)> {
let lk = probe(user_key);
if let Some((_, Some(value))) = active.get(&lk) {
return Some((Source::Active, value.as_slice().to_vec()));
}
{
let frozen = frozen.read().expect("frozen");
for memtable in frozen.iter().rev() {
if let Some((_, Some(value))) = memtable.get(&lk) {
return Some((Source::Frozen, value.as_slice().to_vec()));
}
}
}
let version = version.lock().expect("version");
version
.iter()
.rev()
.find_map(|table| table_holds(table, lk.prefixed_user_key()))
.map(|value| (Source::Version, value))
}
pub fn a_flush_never_hides_a_key() {
explore("a_flush_never_hides_a_key", 32, 8, |witness| {
let (active, frozen, version) = frozen_flush_fixture();
let flush = {
let frozen = Arc::clone(&frozen);
let version = Arc::clone(&version);
loom::thread::spawn(move || {
let retiring = Arc::clone(&frozen.read().expect("frozen")[0]);
version.lock().expect("version").push(install(&retiring));
frozen.write().expect("frozen").remove(0);
})
};
let reader = {
let active = Arc::clone(&active);
let frozen = Arc::clone(&frozen);
let version = Arc::clone(&version);
let witness = witness.clone();
loom::thread::spawn(move || {
let (source, value) = read(&active, &frozen, &version, b"k")
.expect("the flush handoff hid a key that was never deleted");
assert_eq!(value, b"v");
if source == Source::Version {
witness.record();
}
})
};
flush.join().expect("flush");
reader.join().expect("reader");
});
}
pub fn a_flush_that_retires_before_it_installs_hides_a_key() {
explore(
"a_flush_that_retires_before_it_installs_hides_a_key",
4,
1,
|witness| {
let (active, frozen, version) = frozen_flush_fixture();
let flush = {
let frozen = Arc::clone(&frozen);
let version = Arc::clone(&version);
loom::thread::spawn(move || {
let retiring = frozen.write().expect("frozen").remove(0);
version.lock().expect("version").push(install(&retiring));
})
};
let reader = {
let active = Arc::clone(&active);
let frozen = Arc::clone(&frozen);
let version = Arc::clone(&version);
let witness = witness.clone();
loom::thread::spawn(move || {
witness.record();
let found = read(&active, &frozen, &version, b"k");
assert!(found.is_some(), "the key went missing");
})
};
flush.join().expect("flush");
reader.join().expect("reader");
},
);
}
#[allow(clippy::type_complexity)]
fn frozen_flush_fixture() -> (
Arc<MemTable>,
Arc<RwLock<Vec<Arc<MemTable>>>>,
Arc<Mutex<Vec<InstalledTable>>>,
) {
let retiring = memtable();
retiring.put(probe(b"k").prefixed_user_key(), b"v", SEQ);
(
Arc::new(memtable()),
Arc::new(RwLock::new(vec![Arc::new(retiring)])),
Arc::new(Mutex::new(Vec::new())),
)
}
pub fn the_read_horizon_never_outruns_the_memtable() {
explore(
"the_read_horizon_never_outruns_the_memtable",
16,
4,
|witness| {
let mt = Arc::new(memtable());
let horizon = Arc::new(ReadHorizon::new(0));
let writer = {
let mt = Arc::clone(&mt);
let horizon = Arc::clone(&horizon);
loom::thread::spawn(move || {
mt.put(probe(b"k").prefixed_user_key(), b"v", SEQ);
horizon.publish(SEQ);
})
};
let reader = {
let mt = Arc::clone(&mt);
let horizon = Arc::clone(&horizon);
let witness = witness.clone();
loom::thread::spawn(move || {
let visible = horizon.visible();
if visible >= SEQ {
witness.record();
let (seq, value) = mt
.get(&snapshot(b"k", visible))
.expect("H2: a published sequence is readable");
assert_eq!(seq, SEQ);
assert_eq!(value.expect("live value").as_slice(), b"v");
}
})
};
writer.join().expect("writer");
reader.join().expect("reader");
},
);
}
pub fn a_relaxed_read_horizon_outruns_the_memtable() {
explore(
"a_relaxed_read_horizon_outruns_the_memtable",
4,
1,
|witness| {
let mt = Arc::new(memtable());
let horizon = Arc::new(AtomicU64::new(0));
let writer = {
let mt = Arc::clone(&mt);
let horizon = Arc::clone(&horizon);
loom::thread::spawn(move || {
mt.put(probe(b"k").prefixed_user_key(), b"v", SEQ);
horizon.store(SEQ, Ordering::Relaxed);
})
};
let reader = {
let mt = Arc::clone(&mt);
let horizon = Arc::clone(&horizon);
let witness = witness.clone();
loom::thread::spawn(move || {
let visible = horizon.load(Ordering::Relaxed);
if visible >= SEQ {
witness.record();
assert!(
mt.get(&snapshot(b"k", visible)).is_some(),
"a relaxed horizon advertised a sequence the reader cannot see"
);
}
})
};
writer.join().expect("writer");
reader.join().expect("reader");
},
);
}