use crate::oracle::Oracle;
#[cfg(not(target_arch = "wasm32"))]
use crate::persistence::Persistence;
use crate::queue::{Commit, Merge};
use crate::versions::Versions;
use crate::DatabaseOptions;
use bytes::Bytes;
use crossbeam_skiplist::SkipMap;
use papaya::HashSet;
use parking_lot::RwLock;
use std::sync::atomic::{fence, AtomicBool, AtomicU64, Ordering};
use std::sync::Arc;
#[cfg(not(target_arch = "wasm32"))]
use std::thread::JoinHandle;
pub(crate) const SLOT_PINNING: u64 = u64::MAX;
pub(crate) const COMMIT_ABORTED: u64 = u64::MAX;
pub(crate) struct Slot {
pub(crate) version: AtomicU64,
pub(crate) commit: AtomicU64,
}
impl Slot {
pub(crate) fn pinning() -> Self {
Self {
version: AtomicU64::new(SLOT_PINNING),
commit: AtomicU64::new(SLOT_PINNING),
}
}
}
pub struct Inner {
pub(crate) oracle: Arc<Oracle>,
pub(crate) datastore: SkipMap<Bytes, RwLock<Versions>>,
pub(crate) readers: SkipMap<u64, Arc<Slot>>,
pub(crate) reader_slot_id: AtomicU64,
pub(crate) commit_watermark: AtomicU64,
pub(crate) transaction_queue_id: AtomicU64,
pub(crate) transaction_commit_id: AtomicU64,
pub(crate) transaction_commit_queue: SkipMap<u64, Arc<Commit>>,
pub(crate) transaction_merge_queue: SkipMap<u64, Arc<Merge>>,
pub(crate) merge_retire_id: AtomicU64,
pub(crate) gc_candidates: HashSet<Bytes>,
#[cfg(not(target_arch = "wasm32"))]
pub(crate) persistence: RwLock<Option<Arc<Persistence>>>,
pub(crate) background_threads_enabled: AtomicBool,
#[cfg(not(target_arch = "wasm32"))]
pub(crate) transaction_cleanup_handle: RwLock<Option<JoinHandle<()>>>,
#[cfg(not(target_arch = "wasm32"))]
pub(crate) garbage_collection_handle: RwLock<Option<JoinHandle<()>>>,
pub(crate) reset_threshold: usize,
}
impl Inner {
pub fn new(opts: &DatabaseOptions) -> Self {
Self {
oracle: Oracle::new(),
datastore: SkipMap::new(),
readers: SkipMap::new(),
reader_slot_id: AtomicU64::new(0),
commit_watermark: AtomicU64::new(0),
transaction_queue_id: AtomicU64::new(0),
transaction_commit_id: AtomicU64::new(0),
transaction_commit_queue: SkipMap::new(),
transaction_merge_queue: SkipMap::new(),
merge_retire_id: AtomicU64::new(0),
gc_candidates: HashSet::new(),
#[cfg(not(target_arch = "wasm32"))]
persistence: RwLock::new(None),
background_threads_enabled: AtomicBool::new(true),
#[cfg(not(target_arch = "wasm32"))]
transaction_cleanup_handle: RwLock::new(None),
#[cfg(not(target_arch = "wasm32"))]
garbage_collection_handle: RwLock::new(None),
reset_threshold: opts.reset_threshold,
}
}
}
impl Inner {
#[inline]
pub(crate) fn earliest_active_version(&self, fallback: u64) -> Option<u64> {
earliest_pinned(&self.readers, |s| &s.version, fallback, None)
}
#[inline]
pub(crate) fn earliest_active_commit(&self, fallback: u64) -> Option<u64> {
earliest_pinned(&self.readers, |s| &s.commit, fallback, None)
}
pub(crate) fn cleanup_commit_queue(&self) {
self.refresh_commit_watermark();
let fallback = self.transaction_commit_id.load(Ordering::SeqCst);
if let Some(oldest) = self.earliest_active_commit(fallback) {
self.transaction_commit_queue.range(..oldest).for_each(|e| {
e.remove();
});
}
}
pub(crate) fn inline_gc_watermark(&self, own_slot: u64) -> Option<u64> {
let now = self.oracle.timestamp.load(Ordering::SeqCst);
earliest_pinned(&self.readers, |s| &s.version, now, Some(own_slot))
}
pub(crate) fn compute_cleanup_ts(&self) -> Option<u64> {
self.refresh_merge_watermark();
for _ in 0..3 {
let now = self.oracle.timestamp.load(Ordering::SeqCst);
if let Some(earliest) = self.earliest_active_version(now) {
return Some(earliest.min(now));
}
std::hint::spin_loop();
}
None
}
pub(crate) fn try_advance_commit_prefix(&self) {
let mut spins = 0;
loop {
let cur = self.transaction_commit_id.load(Ordering::SeqCst);
let next = cur + 1;
if next > self.transaction_queue_id.load(Ordering::SeqCst) {
break;
}
if self.transaction_commit_queue.get(&next).is_none() {
break;
}
if self
.transaction_commit_id
.compare_exchange_weak(cur, next, Ordering::SeqCst, Ordering::SeqCst)
.is_err()
{
crate::tx::backoff(spins);
spins += 1;
continue;
}
}
}
pub(crate) fn try_advance_merge_clock(&self) {
let mut spins = 0;
loop {
let cur = self.oracle.timestamp.load(Ordering::SeqCst);
let next = cur + 1;
if next > self.oracle.alloc.load(Ordering::SeqCst) {
break;
}
if self.transaction_merge_queue.get(&next).is_none() {
break;
}
if self
.oracle
.timestamp
.compare_exchange_weak(cur, next, Ordering::SeqCst, Ordering::SeqCst)
.is_err()
{
crate::tx::backoff(spins);
spins += 1;
continue;
}
}
}
pub(crate) fn refresh_commit_watermark(&self) {
self.try_advance_commit_prefix();
self.advance_commit_watermark();
}
pub(crate) fn refresh_merge_watermark(&self) {
self.try_advance_merge_clock();
self.advance_merge_retirement();
}
pub(crate) fn advance_commit_watermark(&self) {
loop {
let wm = self.commit_watermark.load(Ordering::SeqCst);
let next = wm + 1;
if next > self.transaction_commit_id.load(Ordering::SeqCst) {
break;
}
let complete = match self.transaction_commit_queue.get(&next) {
Some(entry) => entry.value().merge_version.load(Ordering::SeqCst) != 0,
None => true,
};
if !complete {
break;
}
let _ = self.commit_watermark.compare_exchange(
wm,
next,
Ordering::SeqCst,
Ordering::SeqCst,
);
}
}
pub(crate) fn advance_merge_retirement(&self) {
loop {
let wm = self.merge_retire_id.load(Ordering::SeqCst);
let next = wm + 1;
if next > self.oracle.timestamp.load(Ordering::SeqCst) {
break;
}
if let Some(entry) = self.transaction_merge_queue.get(&next) {
if !entry.value().applied.load(Ordering::SeqCst) {
break;
}
entry.remove();
}
let _ =
self.merge_retire_id.compare_exchange(wm, next, Ordering::SeqCst, Ordering::SeqCst);
}
}
pub(crate) fn run_gc_tracked(&self, cleanup_ts: u64) {
let candidates = self.gc_candidates.pin();
if candidates.is_empty() {
return;
}
let mut keys: Vec<Bytes> = Vec::with_capacity(candidates.len());
keys.extend(candidates.iter().cloned());
for key in keys {
candidates.remove(&key);
let Some(entry) = self.datastore.get(&key) else {
continue;
};
let mut versions = entry.value().write();
if entry.is_removed() {
continue;
}
if versions.gc_older_versions(cleanup_ts) == 0 {
entry.remove();
} else if versions.needs_gc() {
candidates.insert(key);
}
}
}
pub(crate) fn run_gc_full(&self, cleanup_ts: u64) {
for entry in self.datastore.iter() {
let mut versions = entry.value().write();
if versions.gc_older_versions(cleanup_ts) == 0 {
entry.remove();
} else if versions.needs_gc() {
self.gc_candidates.pin().insert(entry.key().clone());
}
}
}
}
#[inline]
pub(crate) fn earliest_pinned(
map: &SkipMap<u64, Arc<Slot>>,
dim: impl Fn(&Slot) -> &AtomicU64,
fallback: u64,
exclude: Option<u64>,
) -> Option<u64> {
fence(Ordering::SeqCst);
let mut min = fallback;
for entry in map.iter() {
if Some(*entry.key()) == exclude {
continue;
}
match dim(entry.value()).load(Ordering::SeqCst) {
SLOT_PINNING => return None,
v => min = min.min(v),
}
}
Some(min)
}
impl Default for Inner {
fn default() -> Self {
Self::new(&DatabaseOptions::default())
}
}