use crate::bloom::BloomFilter;
use crate::direction::Direction;
use crate::iter::MergeQueueIter;
use crate::queue::{Commit, Merge};
use bytes::Bytes;
use papaya::HashSet;
use std::collections::BTreeMap;
use std::sync::atomic::{AtomicBool, AtomicU64};
use std::sync::Arc;
pub struct ReadsetConflictScenario {
commit: Arc<Commit>,
readset: HashSet<Bytes>,
readset_bloom: BloomFilter,
}
impl ReadsetConflictScenario {
pub fn new(writeset_keys: &[Bytes], readset_keys: &[Bytes]) -> Self {
let mut ws: Vec<Bytes> = writeset_keys.to_vec();
ws.sort();
ws.dedup();
let keys: Arc<[Bytes]> = ws.into();
let mut writeset_bloom = BloomFilter::new();
for k in keys.iter() {
writeset_bloom.insert(k);
}
let commit = Arc::new(Commit {
keys,
writeset_bloom,
merge_version: AtomicU64::new(0),
});
let readset = HashSet::new();
let mut readset_bloom = BloomFilter::new();
{
let pin = readset.pin();
for k in readset_keys {
pin.insert(k.clone());
readset_bloom.insert(k);
}
}
Self {
commit,
readset,
readset_bloom,
}
}
pub fn check_with_bloom(&self) -> bool {
self.commit.is_disjoint_readset_bloom(&self.readset, &self.readset_bloom)
}
pub fn check_without_bloom(&self) -> bool {
self.commit.is_disjoint_readset(&self.readset)
}
}
pub struct WritesetConflictScenario {
committed: Arc<Commit>,
current: Arc<Commit>,
}
impl WritesetConflictScenario {
pub fn new(committed_keys: &[Bytes], current_keys: &[Bytes]) -> Self {
Self {
committed: Arc::new(Self::build_commit(committed_keys)),
current: Arc::new(Self::build_commit(current_keys)),
}
}
pub fn check_with_bloom(&self) -> bool {
self.committed.is_disjoint_writeset_bloom(&self.current)
}
pub fn check_without_bloom(&self) -> bool {
self.committed.is_disjoint_writeset(&self.current)
}
fn build_commit(input: &[Bytes]) -> Commit {
let mut ws: Vec<Bytes> = input.to_vec();
ws.sort();
ws.dedup();
let keys: Arc<[Bytes]> = ws.into();
let mut writeset_bloom = BloomFilter::new();
for k in keys.iter() {
writeset_bloom.insert(k);
}
Commit {
keys,
writeset_bloom,
merge_version: AtomicU64::new(0),
}
}
}
pub struct MergeQueueScenario {
sources: Vec<Arc<Merge>>,
beg: Bytes,
end: Bytes,
}
impl MergeQueueScenario {
pub fn new(num_sources: usize, keys_per_source: usize, total_keys: usize) -> Self {
let mut sources = Vec::with_capacity(num_sources);
for i in 0..num_sources {
let mut ws = BTreeMap::new();
for j in 0..keys_per_source {
let key_idx = (i.wrapping_mul(31) + j.wrapping_mul(17)) % total_keys;
let key = Bytes::from(format!("key_{:08}", key_idx).into_bytes());
ws.insert(key, Some(Bytes::from_static(b"v")));
}
sources.push(Arc::new(Merge {
writeset: Arc::new(ws),
applied: AtomicBool::new(false),
}));
}
let beg = Bytes::from(b"key_00000000".to_vec());
let end = Bytes::from(format!("key_{:08}", total_keys).into_bytes());
Self {
sources,
beg,
end,
}
}
pub fn iter_forward_count(&self) -> usize {
MergeQueueIter::new(
self.sources.clone(),
self.beg.clone(),
self.end.clone(),
Direction::Forward,
)
.count()
}
pub fn iter_forward_take(&self, n: usize) -> usize {
MergeQueueIter::new(
self.sources.clone(),
self.beg.clone(),
self.end.clone(),
Direction::Forward,
)
.take(n)
.count()
}
}
pub struct WatermarkScanScenario {
db: crate::Database,
own_slot: u64,
}
impl WatermarkScanScenario {
pub fn new(num_readers: usize) -> Self {
use crate::inner::Slot;
use std::sync::atomic::Ordering;
let db = crate::Database::new_with_options(
crate::DatabaseOptions::default().with_all_workers_disabled(),
);
for i in 0..num_readers {
let id = db.reader_slot_id.fetch_add(1, Ordering::SeqCst) + 1;
let slot = Arc::new(Slot::pinning());
slot.version.store(1 + (i as u64 % 16), Ordering::SeqCst);
slot.commit.store(1 + (i as u64 % 16), Ordering::SeqCst);
db.readers.insert(id, slot);
}
let own_slot = db.reader_slot_id.fetch_add(1, Ordering::SeqCst) + 1;
let slot = Arc::new(Slot::pinning());
slot.version.store(32, Ordering::SeqCst);
slot.commit.store(32, Ordering::SeqCst);
db.readers.insert(own_slot, slot);
Self {
db,
own_slot,
}
}
pub fn scan(&self) -> Option<u64> {
self.db.inline_gc_watermark(self.own_slot)
}
}