use parking_lot::{Mutex, RwLock};
use std::any::{Any, TypeId};
use std::collections::{BTreeMap, BTreeSet, HashMap};
use std::ops::Bound;
use std::sync::Arc;
use std::sync::atomic::{AtomicU16, AtomicU64, AtomicUsize, Ordering};
use crate::core::{
Error, IndexKey, IsolationLevel, Result, Snapshot as SnapshotLevel, TableId, Timestamp, TxnId,
Versioned, Visibility,
};
use crossbeam_epoch::{Atomic, Guard, Shared};
use crate::engine::oracle::{Oracle, OracleConfig};
use crate::engine::slotmap::SlotMap;
use crate::engine::ssi::{Readers, TxnState};
use crate::engine::txn::Transaction;
pub(crate) struct Version<T> {
pub(crate) begin: AtomicU64,
pub(crate) end: AtomicU64,
pub(crate) prev: Atomic<Version<T>>,
pub(crate) value: Option<T>,
pub(crate) writer: Option<Arc<TxnState>>,
}
impl<T> Version<T> {
pub(crate) fn visible_to(&self, snapshot: Timestamp, reader: TxnId) -> bool {
let begin = Visibility::decode(self.begin.load(Ordering::Acquire));
let end = Visibility::decode(self.end.load(Ordering::Acquire));
begin.reached(snapshot, reader) && !end.reached(snapshot, reader)
}
}
#[repr(align(64))] pub(crate) struct Slot<T> {
pub(crate) latest: Atomic<Version<T>>,
pub(crate) lock: AtomicU64,
pub(crate) readers: Mutex<Readers>,
}
impl<T> Slot<T> {
pub(crate) fn new() -> Self {
Slot {
latest: Atomic::null(),
lock: AtomicU64::new(0),
readers: Mutex::new(Readers::default()),
}
}
pub(crate) fn read<'g>(
&self,
snapshot: Timestamp,
reader: TxnId,
guard: &'g Guard,
) -> Option<&'g Version<T>> {
let mut cur = self.latest.load(Ordering::Acquire, guard);
loop {
let v = unsafe { cur.as_ref() }?;
if v.visible_to(snapshot, reader) {
return Some(v);
}
cur = v.prev.load(Ordering::Acquire, guard);
}
}
pub(crate) fn read_pair<'g>(
&self,
snapshot: Timestamp,
reader: TxnId,
guard: &'g Guard,
) -> (Option<&'g Version<T>>, Option<&'g Version<T>>) {
let mut seen = None;
let mut cur = self.latest.load(Ordering::Acquire, guard);
loop {
let Some(v) = (unsafe { cur.as_ref() }) else {
return (seen, None);
};
if seen.is_none() && v.visible_to(snapshot, reader) {
seen = Some(v);
}
if v.visible_to(snapshot, TxnId::NONE) {
return (seen, Some(v));
}
cur = v.prev.load(Ordering::Acquire, guard);
}
}
pub(crate) fn read_committed_now<'g>(
&self,
reader: TxnId,
guard: &'g Guard,
) -> Option<&'g Version<T>> {
let mut cur = self.latest.load(Ordering::Acquire, guard);
loop {
let v = unsafe { cur.as_ref() }?;
match Visibility::decode(v.begin.load(Ordering::Acquire)) {
Visibility::CommittedAt(_) => return Some(v),
Visibility::InFlight(id) if id == reader => return Some(v),
Visibility::InFlight(_) => {}
}
cur = v.prev.load(Ordering::Acquire, guard);
}
}
pub(crate) fn try_lock(&self, txn: TxnId) -> bool {
self.lock
.compare_exchange(0, txn.0, Ordering::AcqRel, Ordering::Acquire)
.is_ok()
|| self.lock.load(Ordering::Acquire) == txn.0 }
pub(crate) fn unlock(&self, txn: TxnId) {
let _ = self
.lock
.compare_exchange(txn.0, 0, Ordering::AcqRel, Ordering::Acquire);
}
pub(crate) fn prune(&self, gc: Timestamp, guard: &Guard) {
let head = self.latest.load(Ordering::Acquire, guard);
if let Some(v) = unsafe { head.as_ref() }
&& v.value.is_none()
&& Self::began_at_or_before(v, gc)
{
self.latest.store(Shared::null(), Ordering::Release);
Self::retire_chain(head, guard);
return;
}
let mut cur = head;
for _ in 0..Self::PRUNE_PROBE {
let Some(v) = (unsafe { cur.as_ref() }) else {
return;
};
let next = v.prev.load(Ordering::Acquire, guard);
let Some(n) = (unsafe { next.as_ref() }) else {
return;
};
if !Self::ended_at_or_before(n, gc) {
cur = next;
continue;
}
v.prev.store(Shared::null(), Ordering::Release);
Self::retire_chain(next, guard);
return;
}
}
fn needs_prune(&self, gc: Timestamp, guard: &Guard) -> bool {
let head = self.latest.load(Ordering::Acquire, guard);
let Some(first) = (unsafe { head.as_ref() }) else {
return false;
};
if first.value.is_none() && Self::began_at_or_before(first, gc) {
return true;
}
let mut cur = head;
for _ in 0..Self::PRUNE_PROBE {
let Some(v) = (unsafe { cur.as_ref() }) else {
return false;
};
let next = v.prev.load(Ordering::Acquire, guard);
let Some(n) = (unsafe { next.as_ref() }) else {
return false;
};
if Self::ended_at_or_before(n, gc) {
return true;
}
cur = next;
}
false
}
fn retire_chain(head: Shared<'_, Version<T>>, guard: &Guard) {
let mut doomed = head;
while let Some(d) = unsafe { doomed.as_ref() } {
let following = d.prev.load(Ordering::Acquire, guard);
unsafe { guard.defer_destroy(doomed) };
doomed = following;
}
}
fn began_at_or_before(v: &Version<T>, gc: Timestamp) -> bool {
match Visibility::decode(v.begin.load(Ordering::Acquire)) {
Visibility::CommittedAt(ts) => ts <= gc,
Visibility::InFlight(_) => false,
}
}
const PRUNE_PROBE: usize = 8;
fn ended_at_or_before(v: &Version<T>, gc: Timestamp) -> bool {
match Visibility::decode(v.end.load(Ordering::Acquire)) {
Visibility::CommittedAt(ts) => ts <= gc,
Visibility::InFlight(_) => false,
}
}
}
type Matches<'g, T> = Vec<(<T as Versioned>::Key, &'g Version<T>)>;
pub(crate) trait Sweep: Send + Sync {
fn sweep_shard(&self, round: usize, gc: Timestamp, sweeper: TxnId);
fn compact(&self) -> usize;
}
pub(crate) struct Table<T: Versioned> {
slots: SlotMap<T>,
secondary: Vec<RwLock<BTreeMap<IndexKey, BTreeSet<T::Key>>>>,
predicate_locks: Mutex<Vec<PredicateLock<T>>>,
unique_claims: Vec<Mutex<HashMap<IndexKey, TxnId>>>,
}
#[derive(Clone, Copy, PartialEq, Eq, Debug)]
pub(crate) enum Claim {
Acquired,
Held,
Contended,
}
impl<T: Versioned> Drop for Table<T> {
fn drop(&mut self) {
let guard = unsafe { crossbeam_epoch::unprotected() };
self.slots.for_each(guard, |record| {
let mut cur = record
.slot
.latest
.swap(Shared::null(), Ordering::Relaxed, guard);
while !cur.is_null() {
let owned = unsafe { cur.into_owned() };
cur = owned.prev.load(Ordering::Relaxed, guard);
drop(owned);
}
});
}
}
struct PredicateLock<T> {
state: Arc<TxnState>,
predicate: Arc<dyn Fn(&T) -> bool + Send + Sync>,
}
impl<T: Versioned> Sweep for Table<T> {
fn sweep_shard(&self, round: usize, gc: Timestamp, sweeper: TxnId) {
let guard = crossbeam_epoch::pin();
self.slots.for_each_in_shard(round, &guard, |record| {
if record.slot.needs_prune(gc, &guard) && record.slot.try_lock(sweeper) {
record.slot.prune(gc, &guard);
record.slot.unlock(sweeper);
}
});
}
fn compact(&self) -> usize {
let reclaimed = self.slots.compact();
{
let guard = unsafe { crossbeam_epoch::unprotected() };
for index in &self.secondary {
let mut index = index.write();
index.retain(|_, candidates| {
candidates.retain(|key| self.slots.get(key, guard).is_some());
!candidates.is_empty()
});
}
}
self.predicate_locks.lock().clear();
for claims in &self.unique_claims {
claims.lock().clear();
}
reclaimed
}
}
impl<T: Versioned> Table<T> {
fn new() -> Self {
Table {
slots: SlotMap::new(),
secondary: T::indexes()
.iter()
.map(|_| RwLock::new(BTreeMap::new()))
.collect(),
predicate_locks: Mutex::new(Vec::new()),
unique_claims: T::indexes()
.iter()
.map(|_| Mutex::new(HashMap::new()))
.collect(),
}
}
pub(crate) fn slot(&self, key: &T::Key, guard: &Guard) -> Option<&Slot<T>> {
self.slots.get(key, guard)
}
pub(crate) fn slot_or_create(&self, key: &T::Key, guard: &Guard) -> &Slot<T> {
self.slots.get_or_create(key, guard)
}
pub(crate) fn index_record(&self, record: &T, previous: Option<&T>) {
let key = record.key();
for (desc, map) in T::indexes().iter().zip(&self.secondary) {
let index_key = (desc.extract)(record);
if previous.is_some_and(|prev| (desc.extract)(prev) == index_key) {
continue;
}
map.write()
.entry(index_key)
.or_default()
.insert(key.clone());
}
}
pub(crate) fn index_candidates(
&self,
position: usize,
lo: Bound<IndexKey>,
hi: Bound<IndexKey>,
) -> Vec<T::Key> {
let map = self.secondary[position].read();
map.range((lo, hi))
.flat_map(|(_, keys)| keys.iter().cloned())
.collect()
}
fn visit_slots<'s>(&'s self, guard: &Guard, mut f: impl FnMut(&'s Slot<T>)) {
self.slots.for_each(guard, |record| f(&record.slot));
}
pub(crate) fn matching<'g>(
&self,
snapshot: Timestamp,
reader: TxnId,
predicate: &dyn Fn(&T) -> bool,
guard: &'g Guard,
) -> Matches<'g, T> {
let mut out = Vec::new();
self.visit_slots(guard, |slot| {
if let Some(version) = slot.read(snapshot, reader, guard)
&& let Some(value) = version.value.as_ref()
&& predicate(value)
{
out.push((value.key(), version));
}
});
out.sort_unstable_by(|a, b| a.0.cmp(&b.0));
out
}
pub(crate) fn matching_pair<'g>(
&self,
snapshot: Timestamp,
reader: TxnId,
predicate: &dyn Fn(&T) -> bool,
guard: &'g Guard,
) -> (Matches<'g, T>, Matches<'g, T>) {
let mut seen = Vec::new();
let mut committed = Vec::new();
self.visit_slots(guard, |slot| {
let (a, b) = slot.read_pair(snapshot, reader, guard);
let same = match (a, b) {
(Some(a), Some(b)) => std::ptr::eq(a, b),
(None, None) => true,
_ => false,
};
if let Some(version) = a
&& let Some(value) = version.value.as_ref()
&& predicate(value)
{
let key = value.key();
if same {
committed.push((key.clone(), version));
}
seen.push((key, version));
}
if !same
&& let Some(version) = b
&& let Some(value) = version.value.as_ref()
&& predicate(value)
{
committed.push((value.key(), version));
}
});
seen.sort_unstable_by(|a, b| a.0.cmp(&b.0));
committed.sort_unstable_by(|a, b| a.0.cmp(&b.0));
(seen, committed)
}
pub(crate) fn matching_in_index<'g>(
&self,
position: usize,
lo: &Bound<IndexKey>,
hi: &Bound<IndexKey>,
snapshot: Timestamp,
reader: TxnId,
guard: &'g Guard,
) -> Matches<'g, T> {
let desc = &T::indexes()[position];
let in_range = |k: &IndexKey| {
let lo_ok = match lo {
Bound::Included(b) => k >= b,
Bound::Excluded(b) => k > b,
Bound::Unbounded => true,
};
let hi_ok = match hi {
Bound::Included(b) => k <= b,
Bound::Excluded(b) => k < b,
Bound::Unbounded => true,
};
lo_ok && hi_ok
};
let mut out: Vec<(IndexKey, T::Key, &'g Version<T>)> = Vec::new();
for key in self.index_candidates(position, lo.clone(), hi.clone()) {
let Some(slot) = self.slot(&key, guard) else {
continue;
};
let Some(version) = slot.read(snapshot, reader, guard) else {
continue;
};
let Some(value) = version.value.as_ref() else {
continue;
};
let actual = (desc.extract)(value);
if in_range(&actual) {
out.push((actual, key, version));
}
}
out.sort_unstable_by(|a, b| a.0.cmp(&b.0).then_with(|| a.1.cmp(&b.1)));
out.into_iter().map(|(_, k, v)| (k, v)).collect()
}
pub(crate) fn register_predicate(
&self,
state: &Arc<TxnState>,
predicate: Arc<dyn Fn(&T) -> bool + Send + Sync>,
) {
self.predicate_locks.lock().push(PredicateLock {
state: Arc::clone(state),
predicate,
});
}
pub(crate) fn predicate_readers_of(
&self,
record: &T,
writer: TxnId,
gc_watermark: Timestamp,
) -> Vec<Arc<TxnState>> {
let mut locks = self.predicate_locks.lock();
locks.retain(|l| !l.state.is_expired(gc_watermark));
locks
.iter()
.filter(|l| l.state.id() != writer && (l.predicate)(record))
.map(|l| Arc::clone(&l.state))
.collect()
}
pub(crate) fn try_claim_unique(
&self,
position: usize,
index_key: &IndexKey,
txn: TxnId,
) -> Claim {
let mut claims = self.unique_claims[position].lock();
match claims.get(index_key) {
Some(&owner) if owner == txn => Claim::Held,
Some(_) => Claim::Contended,
None => {
claims.insert(index_key.clone(), txn);
Claim::Acquired
}
}
}
pub(crate) fn release_unique(&self, position: usize, index_key: &IndexKey, txn: TxnId) {
let mut claims = self.unique_claims[position].lock();
if claims.get(index_key) == Some(&txn) {
claims.remove(index_key);
}
}
pub(crate) fn unique_key_taken(
&self,
position: usize,
index_key: &IndexKey,
own_key: &T::Key,
reader: TxnId,
guard: &Guard,
) -> bool {
let extract = T::indexes()[position].extract;
let candidates = {
let map = self.secondary[position].read();
map.get(index_key).cloned().unwrap_or_default()
};
candidates.into_iter().any(|candidate| {
if candidate == *own_key {
return false;
}
self.slot(&candidate, guard)
.and_then(|slot| slot.read_committed_now(reader, guard))
.and_then(|version| version.value.as_ref())
.is_some_and(|value| extract(value) == *index_key)
})
}
}
#[derive(Debug, Default)]
pub struct Config {
pub oracle: OracleConfig,
}
impl Config {
pub fn in_memory() -> Self {
Config::default()
}
}
pub struct Database {
oracle: Oracle,
tables: RwLock<HashMap<TypeId, Arc<dyn Any + Send + Sync>>>,
next_table_id: AtomicU16,
commit_lock: Mutex<()>,
gc_hint: AtomicU64,
sweepable: RwLock<Vec<Arc<dyn Sweep>>>,
sweep_round: AtomicUsize,
}
impl Database {
#[allow(clippy::needless_pass_by_value)]
pub fn open(config: Config) -> Result<Self> {
Ok(Database {
oracle: Oracle::new(config.oracle),
tables: RwLock::new(HashMap::new()),
next_table_id: AtomicU16::new(0),
commit_lock: Mutex::new(()),
gc_hint: AtomicU64::new(0),
sweepable: RwLock::new(Vec::new()),
sweep_round: AtomicUsize::new(0),
})
}
pub(crate) fn gc_hint(&self) -> Timestamp {
Timestamp(self.gc_hint.load(Ordering::Relaxed))
}
pub(crate) fn refresh_gc_hint(&self, ts: Timestamp, sweeper: TxnId) {
if ts.raw().is_multiple_of(Self::GC_HINT_INTERVAL) {
let watermark = self.oracle.gc_watermark();
self.gc_hint.fetch_max(watermark.raw(), Ordering::Relaxed);
if ts.raw().is_multiple_of(Self::SWEEP_INTERVAL) {
self.sweep(watermark, sweeper);
}
}
}
fn sweep(&self, gc: Timestamp, sweeper: TxnId) {
let round = self.sweep_round.fetch_add(1, Ordering::Relaxed);
let tables: Vec<Arc<dyn Sweep>> = self.sweepable.read().clone();
for table in tables {
table.sweep_shard(round, gc, sweeper);
}
}
pub(crate) const GC_HINT_INTERVAL: u64 = 128;
pub(crate) const SWEEP_INTERVAL: u64 = Self::GC_HINT_INTERVAL * 32;
pub(crate) fn oracle(&self) -> &Oracle {
&self.oracle
}
pub(crate) fn commit_lock(&self) -> parking_lot::MutexGuard<'_, ()> {
self.commit_lock.lock()
}
pub fn stats(&self) -> crate::engine::gc::GcStats {
crate::engine::gc::GcStats {
watermark: self.oracle.gc_watermark(),
active_transactions: self.oracle.active_count(),
}
}
pub fn compact(&mut self) -> usize {
let tables: Vec<Arc<dyn Sweep>> = self.sweepable.read().clone();
tables.iter().map(|table| table.compact()).sum()
}
pub fn register<T: Versioned>(&self) -> Result<()> {
let type_id = TypeId::of::<T>();
let mut tables = self.tables.write();
if tables.contains_key(&type_id) {
return Ok(());
}
let id = TableId(self.next_table_id.fetch_add(1, Ordering::Relaxed));
let _ = T::table_id_cell().set(id);
let table = Arc::new(Table::<T>::new());
self.sweepable.write().push(table.clone() as Arc<dyn Sweep>);
tables.insert(type_id, table);
Ok(())
}
pub(crate) fn table_erased(
&self,
type_id: TypeId,
name: &'static str,
) -> Result<&(dyn Any + Send + Sync)> {
let tables = self.tables.read();
let entry = tables
.get(&type_id)
.ok_or(Error::TableNotRegistered { table: name })?;
let erased: &(dyn Any + Send + Sync) = &**entry;
Ok(unsafe { &*(erased as *const (dyn Any + Send + Sync)) })
}
pub fn begin(&self) -> Transaction<'_, SnapshotLevel> {
self.begin_with()
}
pub fn begin_with<I: IsolationLevel>(&self) -> Transaction<'_, I> {
Transaction::new(self)
}
pub fn transaction<R, F>(&self, f: F) -> Result<R>
where
F: FnMut(&mut Transaction<'_, SnapshotLevel>) -> Result<R>,
{
self.transaction_with::<SnapshotLevel, R, F>(f)
}
pub fn transaction_with<I, R, F>(&self, mut f: F) -> Result<R>
where
I: IsolationLevel,
F: FnMut(&mut Transaction<'_, I>) -> Result<R>,
{
const MAX_ATTEMPTS: u32 = 100;
let mut attempt = 0;
loop {
attempt += 1;
let mut tx = self.begin_with::<I>();
let outcome = f(&mut tx).and_then(|r| tx.commit().map(|_| r));
match outcome {
Ok(r) => return Ok(r),
Err(e) if e.is_retriable() && attempt < MAX_ATTEMPTS => {
let backoff = 1u64 << attempt.min(10);
std::thread::sleep(std::time::Duration::from_micros(backoff));
}
Err(e) => return Err(e),
}
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::Mvcc;
#[derive(Mvcc, Clone, Debug)]
struct Counter {
#[mvcc(primary_key)]
id: u64,
hits: u64,
}
fn chain_len(db: &Database, key: u64) -> usize {
let guard = crossbeam_epoch::pin();
let table = db
.table_erased(TypeId::of::<Counter>(), Counter::TABLE_NAME)
.unwrap()
.downcast_ref::<Table<Counter>>()
.unwrap();
let slot = table.slot(&key, &guard).expect("slot exists");
let mut n = 0;
let mut cur = slot.latest.load(Ordering::Acquire, &guard);
while let Some(v) = unsafe { cur.as_ref() } {
n += 1;
cur = v.prev.load(Ordering::Acquire, &guard);
}
n
}
fn bump(db: &Database, times: u64) {
for _ in 0..times {
db.transaction(|tx| tx.update::<Counter>(&1, |c| c.hits += 1))
.unwrap();
}
}
fn open() -> Database {
let db = Database::open(Config::in_memory()).unwrap();
db.register::<Counter>().unwrap();
db.transaction(|tx| tx.insert(Counter { id: 1, hits: 0 }))
.unwrap();
db
}
fn chain_of(db: &Database, key: u64) -> Option<usize> {
let guard = crossbeam_epoch::pin();
let table = db
.table_erased(TypeId::of::<Counter>(), Counter::TABLE_NAME)
.unwrap()
.downcast_ref::<Table<Counter>>()
.unwrap();
let slot = table.slot(&key, &guard)?;
let mut n = 0;
let mut cur = slot.latest.load(Ordering::Acquire, &guard);
while let Some(v) = unsafe { cur.as_ref() } {
n += 1;
cur = v.prev.load(Ordering::Acquire, &guard);
}
Some(n)
}
const FULL_PASS: u64 = Database::SWEEP_INTERVAL * (crate::engine::slotmap::SHARDS as u64 + 2);
fn churn(db: &Database, times: u64) {
for i in 0..times {
db.transaction(|tx| {
tx.insert(Counter {
id: 1_000_000 + i,
hits: 0,
})
})
.unwrap();
}
}
#[derive(crate::Mvcc, Clone, Debug)]
#[mvcc(table = "indexed")]
struct Indexed {
#[mvcc(primary_key)]
id: u64,
#[mvcc(index)]
tag: u64,
}
fn record_count(db: &Database) -> usize {
let guard = crossbeam_epoch::pin();
let table = db
.table_erased(TypeId::of::<Counter>(), Counter::TABLE_NAME)
.unwrap()
.downcast_ref::<Table<Counter>>()
.unwrap();
let mut n = 0;
table.slots.for_each(&guard, |_| n += 1);
n
}
#[test]
fn compaction_frees_deleted_records_and_keeps_live_ones() {
let mut db = open();
for id in 2..200 {
db.transaction(|tx| tx.insert(Counter { id, hits: id }))
.unwrap();
}
for id in (2..200).filter(|id| id % 2 == 0) {
db.transaction(|tx| tx.delete::<Counter>(&id)).unwrap();
}
churn(&db, FULL_PASS);
let before = record_count(&db);
let freed = db.compact();
let after = record_count(&db);
assert!(freed > 0, "nothing was reclaimed");
assert_eq!(before - after, freed, "freed count disagrees with the map");
let mut tx = db.begin();
assert_eq!(tx.get::<Counter>(&1).unwrap().unwrap().hits, 0);
for id in 2..200 {
let got = tx.get::<Counter>(&id).unwrap().map(|c| c.hits);
if id % 2 == 0 {
assert_eq!(got, None, "deleted key {id} came back");
} else {
assert_eq!(got, Some(id), "survivor {id} was lost or corrupted");
}
}
}
#[test]
fn a_compacted_key_can_be_inserted_again() {
let mut db = open();
db.transaction(|tx| tx.delete::<Counter>(&1)).unwrap();
churn(&db, FULL_PASS);
assert!(db.compact() > 0);
db.transaction(|tx| tx.insert(Counter { id: 1, hits: 9 }))
.unwrap();
let mut tx = db.begin();
assert_eq!(tx.get::<Counter>(&1).unwrap().unwrap().hits, 9);
}
#[test]
fn compaction_drops_index_entries_that_can_no_longer_resolve() {
let mut db = Database::open(Config::in_memory()).unwrap();
db.register::<Indexed>().unwrap();
db.register::<Counter>().unwrap();
db.transaction(|tx| tx.insert(Counter { id: 1, hits: 0 }))
.unwrap();
for id in 0..100 {
db.transaction(|tx| tx.insert(Indexed { id, tag: id % 5 }))
.unwrap();
}
for id in 0..50 {
db.transaction(|tx| tx.delete::<Indexed>(&id)).unwrap();
}
churn(&db, FULL_PASS);
db.compact();
let table = db
.table_erased(TypeId::of::<Indexed>(), Indexed::TABLE_NAME)
.unwrap()
.downcast_ref::<Table<Indexed>>()
.unwrap();
let candidates: usize = table.secondary[0].read().values().map(BTreeSet::len).sum();
assert_eq!(
candidates, 50,
"index should hold only the 50 surviving keys"
);
let mut tx = db.begin();
let hits = tx.scan_index(Indexed::TAG, 0u64..=4).unwrap();
assert_eq!(hits.len(), 50);
}
#[test]
fn a_deleted_record_gives_its_versions_back() {
let db = open();
bump(&db, 64);
db.transaction(|tx| tx.delete::<Counter>(&1)).unwrap();
assert!(
chain_of(&db, 1).unwrap() > 1,
"the tombstone and its history should still be here"
);
churn(&db, FULL_PASS);
assert_eq!(
chain_of(&db, 1),
Some(0),
"a tombstone below the watermark should leave an empty slot"
);
let mut tx = db.begin();
assert!(tx.get::<Counter>(&1).unwrap().is_none());
}
#[test]
fn an_emptied_slot_can_be_refilled() {
let db = open();
db.transaction(|tx| tx.delete::<Counter>(&1)).unwrap();
churn(&db, FULL_PASS);
assert_eq!(chain_of(&db, 1), Some(0), "precondition: slot was emptied");
db.transaction(|tx| tx.insert(Counter { id: 1, hits: 7 }))
.unwrap();
let mut tx = db.begin();
assert_eq!(tx.get::<Counter>(&1).unwrap().unwrap().hits, 7);
}
#[test]
fn a_record_nobody_writes_any_more_is_still_collected() {
let db = open();
bump(&db, Database::GC_HINT_INTERVAL * 4);
let cold = chain_of(&db, 1).unwrap();
assert!(cold > 1, "precondition: key 1 has history to collect");
churn(&db, FULL_PASS);
let swept = chain_of(&db, 1).unwrap();
assert!(
swept < cold,
"cold chain stayed at {swept} (was {cold}): the sweep never reached it"
);
}
#[test]
fn repeated_updates_do_not_grow_the_chain_without_bound() {
let db = open();
let updates = Database::GC_HINT_INTERVAL * 8;
bump(&db, updates);
let len = chain_len(&db, 1) as u64;
assert!(
len < updates / 4,
"chain was {len} after {updates} updates: pruning is not keeping up"
);
}
#[test]
fn an_open_transaction_pins_the_versions_it_can_still_see() {
let db = open();
bump(&db, Database::GC_HINT_INTERVAL * 4);
let mut reader = db.begin();
let seen = reader.get::<Counter>(&1).unwrap().unwrap().hits;
let held = Database::GC_HINT_INTERVAL * 4;
bump(&db, held);
let pinned = chain_len(&db, 1) as u64;
assert_eq!(
reader.get::<Counter>(&1).unwrap().unwrap().hits,
seen,
"the snapshot moved under a live transaction"
);
assert!(
pinned >= held,
"chain was {pinned} with a reader pinning {held} versions"
);
drop(reader);
bump(&db, Database::GC_HINT_INTERVAL * 4);
let after = chain_len(&db, 1) as u64;
assert!(
after < pinned,
"chain stayed at {after} after the reader was dropped (was {pinned})"
);
}
#[test]
fn a_tombstone_survives_pruning() {
let db = open();
bump(&db, Database::GC_HINT_INTERVAL * 2);
db.transaction(|tx| tx.delete::<Counter>(&1)).unwrap();
db.transaction(|tx| tx.insert(Counter { id: 1, hits: 0 }))
.unwrap();
bump(&db, Database::GC_HINT_INTERVAL * 4);
let mut tx = db.begin();
assert!(tx.get::<Counter>(&1).unwrap().is_some());
}
}