use std::sync::atomic::{AtomicU64, Ordering};
use std::sync::Arc;
use crossbeam_skiplist::SkipSet;
pub(crate) struct ActiveTxnTracker {
seqs: Arc<SkipSet<(u64, u64)>>,
next_id: AtomicU64,
}
impl Default for ActiveTxnTracker {
fn default() -> Self {
Self::new()
}
}
impl ActiveTxnTracker {
pub(crate) fn new() -> Self {
Self {
seqs: Arc::new(SkipSet::new()),
next_id: AtomicU64::new(0),
}
}
pub(crate) fn register(self: &Arc<Self>, start_seq: u64) -> ActiveTxnGuard {
let id = self.next_id.fetch_add(1, Ordering::Relaxed);
let entry = (start_seq, id);
self.seqs.insert(entry);
ActiveTxnGuard {
tracker: Arc::clone(self),
entry,
released: false,
}
}
pub(crate) fn oldest(&self) -> Option<u64> {
self.seqs.front().map(|e| e.value().0)
}
#[cfg(test)]
pub(crate) fn len(&self) -> usize {
self.seqs.len()
}
}
pub(crate) struct ActiveTxnGuard {
tracker: Arc<ActiveTxnTracker>,
entry: (u64, u64),
released: bool,
}
impl ActiveTxnGuard {
pub(crate) fn release(&mut self) {
if !self.released {
self.tracker.seqs.remove(&self.entry);
self.released = true;
}
}
}
impl Drop for ActiveTxnGuard {
fn drop(&mut self) {
self.release();
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn empty_tracker_has_no_oldest() {
let t = Arc::new(ActiveTxnTracker::new());
assert_eq!(t.oldest(), None);
assert_eq!(t.len(), 0);
}
#[test]
fn registers_and_drops() {
let t = Arc::new(ActiveTxnTracker::new());
{
let _g = t.register(10);
assert_eq!(t.oldest(), Some(10));
assert_eq!(t.len(), 1);
}
assert_eq!(t.oldest(), None);
assert_eq!(t.len(), 0);
}
#[test]
fn oldest_picks_min() {
let t = Arc::new(ActiveTxnTracker::new());
let _g1 = t.register(20);
let _g2 = t.register(10);
let _g3 = t.register(15);
assert_eq!(t.oldest(), Some(10));
}
#[test]
fn duplicate_start_seqs_dont_collide() {
let t = Arc::new(ActiveTxnTracker::new());
let g1 = t.register(5);
let g2 = t.register(5);
assert_eq!(t.len(), 2);
drop(g1);
assert_eq!(t.oldest(), Some(5));
assert_eq!(t.len(), 1);
drop(g2);
assert_eq!(t.oldest(), None);
}
#[test]
fn explicit_release_is_idempotent() {
let t = Arc::new(ActiveTxnTracker::new());
let mut g = t.register(7);
g.release();
assert_eq!(t.oldest(), None);
g.release();
drop(g);
assert_eq!(t.oldest(), None);
}
}