use hashbrown::HashMap;
use std::panic::{AssertUnwindSafe, catch_unwind, resume_unwind};
use std::sync::Arc;
use std::sync::atomic::{AtomicU32, Ordering};
use fsqlite_types::sync_primitives::{Condvar, Mutex};
use fsqlite_types::{CommitSeq, PageNumber};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct CacheMetricsSnapshot {
pub hits: u64,
pub misses: u64,
pub ghost_hits_b1: u64,
pub ghost_hits_b2: u64,
pub evictions_t1: u64,
pub evictions_t2: u64,
pub version_coalesce_count: u64,
pub admits: u64,
pub t1_len: usize,
pub t2_len: usize,
pub b1_len: usize,
pub b2_len: usize,
pub p: usize,
pub capacity: usize,
pub total_bytes: usize,
pub max_bytes: usize,
pub multi_version_pages: usize,
pub capacity_overflow_events: usize,
}
impl CacheMetricsSnapshot {
#[must_use]
pub fn hit_rate_pct(&self) -> f64 {
let total = self.hits + self.misses + self.ghost_hits_b1 + self.ghost_hits_b2;
if total == 0 {
return 0.0;
}
(self.hits as f64 / total as f64) * 100.0
}
#[must_use]
pub fn total_accesses(&self) -> u64 {
self.hits + self.misses + self.ghost_hits_b1 + self.ghost_hits_b2
}
#[must_use]
pub fn resident_pages(&self) -> usize {
self.t1_len + self.t2_len
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
pub struct CacheKey {
pub pgno: PageNumber,
pub commit_seq: CommitSeq,
}
impl CacheKey {
#[inline]
#[must_use]
pub const fn new(pgno: PageNumber, commit_seq: CommitSeq) -> Self {
Self { pgno, commit_seq }
}
}
pub struct CachedPage {
pub key: CacheKey,
pub data: fsqlite_types::PageData,
pub ref_count: AtomicU32,
pub xxh3: u64,
pub byte_size: usize,
pub wal_frame: Option<u32>,
}
impl CachedPage {
#[must_use]
pub fn new(
key: CacheKey,
data: fsqlite_types::PageData,
xxh3: u64,
wal_frame: Option<u32>,
) -> Self {
let byte_size = data.len();
Self {
key,
data,
ref_count: AtomicU32::new(0),
xxh3,
byte_size,
wal_frame,
}
}
#[inline]
pub fn pin(&self) {
self.ref_count.fetch_add(1, Ordering::Acquire);
}
#[inline]
pub fn unpin(&self) {
let prev = self.ref_count.fetch_sub(1, Ordering::Release);
assert!(prev > 0, "unpin on page with ref_count 0");
}
#[inline]
#[must_use]
pub fn is_pinned(&self) -> bool {
self.ref_count.load(Ordering::Acquire) > 0
}
}
impl std::fmt::Debug for CachedPage {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("CachedPage")
.field("key", &self.key)
.field("data_len", &self.data.len())
.field("ref_count", &self.ref_count.load(Ordering::Relaxed))
.field("xxh3", &format_args!("{:#018x}", self.xxh3))
.field("byte_size", &self.byte_size)
.field("wal_frame", &self.wal_frame)
.finish()
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
struct SlabIdx(u32);
struct SlabNode<T> {
value: T,
prev: Option<SlabIdx>,
next: Option<SlabIdx>,
}
struct IntrusiveList<T> {
slots: Vec<Option<SlabNode<T>>>,
free_indices: Vec<u32>,
head: Option<SlabIdx>,
tail: Option<SlabIdx>,
len: usize,
}
impl<T> IntrusiveList<T> {
fn new() -> Self {
Self {
slots: Vec::new(),
free_indices: Vec::new(),
head: None,
tail: None,
len: 0,
}
}
#[inline]
fn len(&self) -> usize {
self.len
}
#[inline]
fn is_empty(&self) -> bool {
self.len == 0
}
fn push_back(&mut self, value: T) -> SlabIdx {
let idx = self.alloc_slot(value);
if let Some(old_tail) = self.tail {
self.node_mut(old_tail).next = Some(idx);
self.node_mut(idx).prev = Some(old_tail);
} else {
self.head = Some(idx);
}
self.tail = Some(idx);
self.len += 1;
idx
}
fn pop_front(&mut self) -> Option<T> {
let head = self.head?;
Some(self.remove(head))
}
fn remove(&mut self, idx: SlabIdx) -> T {
let node = self.slots[idx.0 as usize]
.take()
.expect("IntrusiveList::remove on vacant slot");
match (node.prev, node.next) {
(Some(p), Some(n)) => {
self.node_mut(p).next = Some(n);
self.node_mut(n).prev = Some(p);
}
(None, Some(n)) => {
self.node_mut(n).prev = None;
self.head = Some(n);
}
(Some(p), None) => {
self.node_mut(p).next = None;
self.tail = Some(p);
}
(None, None) => {
self.head = None;
self.tail = None;
}
}
self.free_indices.push(idx.0);
self.len -= 1;
node.value
}
fn move_to_back(&mut self, idx: SlabIdx) {
if self.tail == Some(idx) {
return;
}
let (prev, next) = {
let n = self.node_ref(idx);
(n.prev, n.next)
};
match (prev, next) {
(Some(p), Some(n)) => {
self.node_mut(p).next = Some(n);
self.node_mut(n).prev = Some(p);
}
(None, Some(n)) => {
self.node_mut(n).prev = None;
self.head = Some(n);
}
_ => return, }
if let Some(old_tail) = self.tail {
self.node_mut(old_tail).next = Some(idx);
}
let old_tail = self.tail;
let node = self.node_mut(idx);
node.prev = old_tail;
node.next = None;
self.tail = Some(idx);
}
fn get(&self, idx: SlabIdx) -> Option<&T> {
self.slots.get(idx.0 as usize)?.as_ref().map(|n| &n.value)
}
#[cfg(test)]
fn front(&self) -> Option<&T> {
let idx = self.head?;
self.get(idx)
}
#[cfg(test)]
fn back(&self) -> Option<&T> {
let idx = self.tail?;
self.get(idx)
}
fn iter(&self) -> IntrusiveListIter<'_, T> {
IntrusiveListIter {
list: self,
current: self.head,
}
}
fn alloc_slot(&mut self, value: T) -> SlabIdx {
if let Some(free) = self.free_indices.pop() {
let idx = SlabIdx(free);
self.slots[free as usize] = Some(SlabNode {
value,
prev: None,
next: None,
});
idx
} else {
let raw = u32::try_from(self.slots.len()).expect("slab overflow");
let idx = SlabIdx(raw);
self.slots.push(Some(SlabNode {
value,
prev: None,
next: None,
}));
idx
}
}
#[inline]
fn node_ref(&self, idx: SlabIdx) -> &SlabNode<T> {
self.slots[idx.0 as usize]
.as_ref()
.expect("dangling SlabIdx")
}
#[inline]
fn node_mut(&mut self, idx: SlabIdx) -> &mut SlabNode<T> {
self.slots[idx.0 as usize]
.as_mut()
.expect("dangling SlabIdx")
}
}
impl<T> Default for IntrusiveList<T> {
fn default() -> Self {
Self::new()
}
}
struct IntrusiveListIter<'a, T> {
list: &'a IntrusiveList<T>,
current: Option<SlabIdx>,
}
impl<'a, T> Iterator for IntrusiveListIter<'a, T> {
type Item = (SlabIdx, &'a T);
fn next(&mut self) -> Option<Self::Item> {
let idx = self.current?;
let node = self.list.slots[idx.0 as usize].as_ref()?;
self.current = node.next;
Some((idx, &node.value))
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Location {
T1(SlabIdx),
T2(SlabIdx),
B1(SlabIdx),
B2(SlabIdx),
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum CacheLookup {
Hit,
GhostHitB1,
GhostHitB2,
Miss,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum AsyncLookup {
Hit,
Loaded,
WaitedForPeerHit,
WaitedForPeerMiss,
}
#[derive(Debug)]
struct InflightLoad {
state: Mutex<InflightState>,
cv: Condvar,
}
impl InflightLoad {
fn new_loading() -> Self {
Self {
state: Mutex::new(InflightState {
loading: true,
waiters: 0,
}),
cv: Condvar::new(),
}
}
}
#[derive(Debug, Clone, Copy)]
struct InflightState {
loading: bool,
waiters: usize,
}
enum InflightRole {
Leader(Arc<InflightLoad>),
Waiter(Arc<InflightLoad>),
}
pub struct ArcCacheInner {
t1: IntrusiveList<Arc<CachedPage>>,
t2: IntrusiveList<Arc<CachedPage>>,
b1: IntrusiveList<CacheKey>,
b2: IntrusiveList<CacheKey>,
p: usize,
capacity: usize,
total_bytes: usize,
max_bytes: usize,
directory: HashMap<CacheKey, Location>,
page_versions: HashMap<PageNumber, Vec<CommitSeq>>,
gc_horizon: CommitSeq,
capacity_overflow_events: usize,
hits: u64,
misses: u64,
ghost_hits_b1: u64,
ghost_hits_b2: u64,
evictions_t1: u64,
evictions_t2: u64,
version_coalesce_count: u64,
admits: u64,
}
impl ArcCacheInner {
#[must_use]
pub fn new(capacity: usize, max_bytes: usize) -> Self {
Self {
t1: IntrusiveList::new(),
t2: IntrusiveList::new(),
b1: IntrusiveList::new(),
b2: IntrusiveList::new(),
p: 0,
capacity,
total_bytes: 0,
max_bytes,
directory: HashMap::with_capacity(capacity * 2),
page_versions: HashMap::new(),
gc_horizon: CommitSeq::ZERO,
capacity_overflow_events: 0,
hits: 0,
misses: 0,
ghost_hits_b1: 0,
ghost_hits_b2: 0,
evictions_t1: 0,
evictions_t2: 0,
version_coalesce_count: 0,
admits: 0,
}
}
#[inline]
#[must_use]
pub fn p(&self) -> usize {
self.p
}
#[inline]
#[must_use]
pub fn capacity(&self) -> usize {
self.capacity
}
#[inline]
#[must_use]
pub fn max_bytes(&self) -> usize {
self.max_bytes
}
#[inline]
#[must_use]
pub fn len(&self) -> usize {
self.t1.len() + self.t2.len()
}
#[inline]
#[must_use]
pub fn is_empty(&self) -> bool {
self.t1.is_empty() && self.t2.is_empty()
}
#[inline]
#[must_use]
pub fn total_bytes(&self) -> usize {
self.total_bytes
}
#[inline]
#[must_use]
pub fn t1_len(&self) -> usize {
self.t1.len()
}
#[inline]
#[must_use]
pub fn t2_len(&self) -> usize {
self.t2.len()
}
#[inline]
#[must_use]
pub fn b1_len(&self) -> usize {
self.b1.len()
}
#[inline]
#[must_use]
pub fn b2_len(&self) -> usize {
self.b2.len()
}
#[inline]
#[must_use]
pub fn capacity_overflow_events(&self) -> usize {
self.capacity_overflow_events
}
#[must_use]
pub fn metrics_snapshot(&self) -> CacheMetricsSnapshot {
let multi_version_pages = self
.page_versions
.values()
.filter(|&versions| versions.len() > 1)
.count();
CacheMetricsSnapshot {
hits: self.hits,
misses: self.misses,
ghost_hits_b1: self.ghost_hits_b1,
ghost_hits_b2: self.ghost_hits_b2,
evictions_t1: self.evictions_t1,
evictions_t2: self.evictions_t2,
version_coalesce_count: self.version_coalesce_count,
admits: self.admits,
t1_len: self.t1.len(),
t2_len: self.t2.len(),
b1_len: self.b1.len(),
b2_len: self.b2.len(),
p: self.p,
capacity: self.capacity,
total_bytes: self.total_bytes,
max_bytes: self.max_bytes,
multi_version_pages,
capacity_overflow_events: self.capacity_overflow_events,
}
}
pub fn reset_metrics(&mut self) {
self.hits = 0;
self.misses = 0;
self.ghost_hits_b1 = 0;
self.ghost_hits_b2 = 0;
self.evictions_t1 = 0;
self.evictions_t2 = 0;
self.version_coalesce_count = 0;
self.admits = 0;
}
#[inline]
#[must_use]
pub fn is_visible(version_commit_seq: CommitSeq, snapshot_high: CommitSeq) -> bool {
version_commit_seq != CommitSeq::ZERO && version_commit_seq <= snapshot_high
}
#[inline]
#[must_use]
pub fn is_visible_or_self(
version_commit_seq: CommitSeq,
snapshot_high: CommitSeq,
in_write_set: bool,
) -> bool {
if version_commit_seq == CommitSeq::ZERO {
return in_write_set;
}
version_commit_seq <= snapshot_high
}
pub fn set_gc_horizon(&mut self, horizon: CommitSeq) {
self.gc_horizon = horizon;
self.coalesce_all_versions();
self.prune_ghosts_below_horizon();
}
fn prune_ghosts_below_horizon(&mut self) {
let mut b1_victims = Vec::new();
for (idx, key) in self.b1.iter() {
if key.commit_seq != CommitSeq::ZERO && key.commit_seq <= self.gc_horizon {
b1_victims.push(idx);
}
}
for idx in b1_victims {
let key = self.b1.remove(idx);
self.directory.remove(&key);
}
let mut b2_victims = Vec::new();
for (idx, key) in self.b2.iter() {
if key.commit_seq != CommitSeq::ZERO && key.commit_seq <= self.gc_horizon {
b2_victims.push(idx);
}
}
for idx in b2_victims {
let key = self.b2.remove(idx);
self.directory.remove(&key);
}
}
pub fn apply_pragma_cache_size(&mut self, n: i32, page_size: usize) {
assert!(page_size > 0, "page_size must be > 0");
let (new_capacity, new_max_bytes) = match n.cmp(&0) {
std::cmp::Ordering::Greater => {
let cap = usize::try_from(n).expect("positive cache_size fits usize");
(cap, cap.saturating_mul(page_size))
}
std::cmp::Ordering::Less => {
let kib =
usize::try_from(n.unsigned_abs()).expect("cache_size magnitude fits usize");
let max_bytes = kib.saturating_mul(1024);
(max_bytes / page_size, max_bytes)
}
std::cmp::Ordering::Equal => (0, 0),
};
self.resize(new_capacity, new_max_bytes);
}
pub fn resize(&mut self, new_capacity: usize, new_max_bytes: usize) {
self.capacity = new_capacity;
self.max_bytes = new_max_bytes;
self.enforce_limits();
self.trim_ghosts();
self.p = self.p.min(self.capacity);
}
#[inline]
#[must_use]
pub fn gc_horizon(&self) -> CommitSeq {
self.gc_horizon
}
#[cfg(test)]
fn set_p_for_tests(&mut self, p: usize) {
self.p = std::cmp::min(self.capacity, p);
}
#[cfg(test)]
fn in_t1(&self, key: &CacheKey) -> bool {
matches!(self.directory.get(key), Some(Location::T1(_)))
}
#[cfg(test)]
fn in_b1(&self, key: &CacheKey) -> bool {
matches!(self.directory.get(key), Some(Location::B1(_)))
}
#[cfg(test)]
fn in_b2(&self, key: &CacheKey) -> bool {
matches!(self.directory.get(key), Some(Location::B2(_)))
}
#[cfg(test)]
fn in_t2(&self, key: &CacheKey) -> bool {
matches!(self.directory.get(key), Some(Location::T2(_)))
}
#[cfg(test)]
fn t2_lru_key(&self) -> Option<CacheKey> {
self.t2.front().map(|page| page.key)
}
#[cfg(test)]
fn t2_mru_key(&self) -> Option<CacheKey> {
self.t2.back().map(|page| page.key)
}
#[cfg(test)]
fn t1_lru_key(&self) -> Option<CacheKey> {
self.t1.front().map(|page| page.key)
}
#[cfg(test)]
#[allow(dead_code)]
fn t1_mru_key(&self) -> Option<CacheKey> {
self.t1.back().map(|page| page.key)
}
#[must_use]
pub fn get(&self, key: &CacheKey) -> Option<Arc<CachedPage>> {
match self.directory.get(key)? {
Location::T1(idx) => self.t1.get(*idx).cloned(),
Location::T2(idx) => self.t2.get(*idx).cloned(),
Location::B1(_) | Location::B2(_) => None,
}
}
pub fn unpin(&mut self, key: &CacheKey) -> bool {
let was_pinned = if let Some(page) = self.get(key) {
if page.is_pinned() {
page.unpin();
true
} else {
false
}
} else {
false
};
if was_pinned {
self.reclaim_one_overflow_slot();
}
was_pinned
}
pub fn request(&mut self, key: &CacheKey) -> CacheLookup {
let location = self.directory.get(key).copied();
match location {
Some(Location::T1(idx)) => {
let page = self.t1.remove(idx);
let new_idx = self.t2.push_back(page);
self.directory.insert(*key, Location::T2(new_idx));
self.hits += 1;
CacheLookup::Hit
}
Some(Location::T2(idx)) => {
self.t2.move_to_back(idx);
self.hits += 1;
CacheLookup::Hit
}
Some(Location::B1(idx)) => {
let delta = std::cmp::max(self.b2.len() / std::cmp::max(self.b1.len(), 1), 1);
self.p = std::cmp::min(self.capacity, self.p.saturating_add(delta));
self.b1.remove(idx);
self.directory.remove(key);
self.ghost_hits_b1 += 1;
CacheLookup::GhostHitB1
}
Some(Location::B2(idx)) => {
let delta = std::cmp::max(self.b1.len() / std::cmp::max(self.b2.len(), 1), 1);
self.p = self.p.saturating_sub(delta);
self.b2.remove(idx);
self.directory.remove(key);
self.ghost_hits_b2 += 1;
CacheLookup::GhostHitB2
}
None => {
self.misses += 1;
CacheLookup::Miss
}
}
}
pub fn admit(&mut self, key: CacheKey, page: CachedPage, lookup: CacheLookup) {
debug_assert!(
!matches!(lookup, CacheLookup::Hit),
"admit called after a cache hit"
);
let page = Arc::new(page);
self.admits += 1;
if self.capacity == 0 {
self.trim_ghosts();
self.p = 0;
return;
}
let byte_size = page.byte_size;
let l1_len = self.t1.len() + self.b1.len();
if l1_len >= self.capacity {
if self.t1.len() < self.capacity {
if let Some(ghost) = self.b1.pop_front() {
self.directory.remove(&ghost);
}
if !self.replace(&key, matches!(lookup, CacheLookup::GhostHitB2)) {
self.capacity_overflow_events = self.capacity_overflow_events.saturating_add(1);
}
} else {
if !self.delete_lru_from_t1() && !self.delete_lru_from_t2() {
self.capacity_overflow_events = self.capacity_overflow_events.saturating_add(1);
}
}
} else {
let total_dir = l1_len + self.t2.len() + self.b2.len();
if total_dir >= self.capacity * 2 {
if let Some(ghost) = self.b2.pop_front() {
self.directory.remove(&ghost);
}
}
if self.t1.len() + self.t2.len() >= self.capacity
&& !self.replace(&key, matches!(lookup, CacheLookup::GhostHitB2))
{
self.capacity_overflow_events = self.capacity_overflow_events.saturating_add(1);
}
}
match lookup {
CacheLookup::GhostHitB1 | CacheLookup::GhostHitB2 => {
let t2_idx = self.t2.push_back(page.clone());
self.directory.insert(key, Location::T2(t2_idx));
}
CacheLookup::Miss => {
let t1_idx = self.t1.push_back(page.clone());
self.directory.insert(key, Location::T1(t1_idx));
}
CacheLookup::Hit => {
unreachable!("admit called after a cache hit");
}
}
self.total_bytes += byte_size;
self.page_versions
.entry(key.pgno)
.or_default()
.push(key.commit_seq);
while self.max_bytes > 0 && self.total_bytes > self.max_bytes && self.len() > 1 {
if !self.evict_one_preferred() {
self.capacity_overflow_events = self.capacity_overflow_events.saturating_add(1);
break; }
}
self.trim_ghosts();
}
fn replace(&mut self, incoming: &CacheKey, from_b2: bool) -> bool {
if self.coalesce_one_superseded_for_replace() {
return true;
}
let incoming_in_b2 = matches!(self.directory.get(incoming), Some(Location::B2(_)));
let from_b2_bias = from_b2 || incoming_in_b2;
let t1_len = self.t1.len();
let prefer_t1 = t1_len > 0 && (t1_len > self.p || (from_b2_bias && t1_len == self.p));
if prefer_t1 {
if self.evict_lru_from_t1() {
return true;
}
return self.evict_lru_from_t2();
}
if self.evict_lru_from_t2() {
return true;
}
self.evict_lru_from_t1()
}
fn evict_lru_from_t1(&mut self) -> bool {
if let Some(victim_idx) = Self::find_unpinned_victim(&self.t1) {
let page = self.t1.remove(victim_idx);
let key = page.key;
self.total_bytes -= page.byte_size;
self.remove_page_version(key.pgno, key.commit_seq);
self.directory.remove(&key);
let ghost_idx = self.b1.push_back(key);
self.directory.insert(key, Location::B1(ghost_idx));
drop(page);
self.evictions_t1 += 1;
true
} else {
false
}
}
fn evict_lru_from_t2(&mut self) -> bool {
if let Some(victim_idx) = Self::find_unpinned_victim(&self.t2) {
let page = self.t2.remove(victim_idx);
let key = page.key;
self.total_bytes -= page.byte_size;
self.remove_page_version(key.pgno, key.commit_seq);
self.directory.remove(&key);
let ghost_idx = self.b2.push_back(key);
self.directory.insert(key, Location::B2(ghost_idx));
drop(page);
self.evictions_t2 += 1;
true
} else {
false
}
}
fn delete_lru_from_t1(&mut self) -> bool {
if let Some(victim_idx) = Self::find_unpinned_victim(&self.t1) {
let page = self.t1.remove(victim_idx);
let key = page.key;
self.total_bytes -= page.byte_size;
self.remove_page_version(key.pgno, key.commit_seq);
self.directory.remove(&key);
drop(page);
self.evictions_t1 += 1;
true
} else {
false
}
}
fn delete_lru_from_t2(&mut self) -> bool {
if let Some(victim_idx) = Self::find_unpinned_victim(&self.t2) {
let page = self.t2.remove(victim_idx);
let key = page.key;
self.total_bytes -= page.byte_size;
self.remove_page_version(key.pgno, key.commit_seq);
self.directory.remove(&key);
drop(page);
self.evictions_t2 += 1;
true
} else {
false
}
}
fn evict_one_preferred(&mut self) -> bool {
if self.coalesce_one_superseded_for_replace() {
return true;
}
let t1_len = self.t1.len();
let prefer_t1 = t1_len > 0 && t1_len > self.p;
if prefer_t1 {
if self.evict_lru_from_t1() {
return true;
}
return self.evict_lru_from_t2();
}
if self.evict_lru_from_t2() {
return true;
}
self.evict_lru_from_t1()
}
fn find_unpinned_victim(list: &IntrusiveList<Arc<CachedPage>>) -> Option<SlabIdx> {
for (idx, page) in list.iter() {
if !page.is_pinned() {
return Some(idx);
}
}
None
}
fn find_superseded_victim(&self, list: &IntrusiveList<Arc<CachedPage>>) -> Option<SlabIdx> {
for (idx, page) in list.iter().take(64) {
if page.is_pinned() {
continue;
}
let versions = self
.page_versions
.get(&page.key.pgno)
.map(|v| v.as_slice())
.unwrap_or(&[]);
let seq = page.key.commit_seq;
if seq != CommitSeq::ZERO && seq <= self.gc_horizon {
let is_superseded = versions.iter().any(|&v| v > seq && v <= self.gc_horizon);
if is_superseded {
return Some(idx);
}
}
}
None
}
fn enforce_limits(&mut self) {
while self.len() > self.capacity
|| (self.max_bytes > 0 && self.total_bytes > self.max_bytes && self.len() > 1)
{
if !self.evict_one_preferred() {
self.capacity_overflow_events = self.capacity_overflow_events.saturating_add(1);
break;
}
}
}
fn reclaim_one_overflow_slot(&mut self) {
if self.capacity_overflow_events == 0 {
return;
}
let over_capacity = self.len() > self.capacity;
let over_bytes = self.max_bytes > 0 && self.total_bytes > self.max_bytes;
if !over_capacity && !over_bytes {
self.capacity_overflow_events = self.capacity_overflow_events.saturating_sub(1);
return;
}
if self.evict_one_preferred() {
self.capacity_overflow_events = self.capacity_overflow_events.saturating_sub(1);
}
}
fn coalesce_one_superseded_for_replace(&mut self) -> bool {
if let Some(idx) = self.find_superseded_victim(&self.t1) {
let page = self.t1.remove(idx);
let key = page.key;
self.total_bytes -= page.byte_size;
self.remove_page_version(key.pgno, key.commit_seq);
self.directory.remove(&key);
self.version_coalesce_count = self.version_coalesce_count.saturating_add(1);
return true;
}
if let Some(idx) = self.find_superseded_victim(&self.t2) {
let page = self.t2.remove(idx);
let key = page.key;
self.total_bytes -= page.byte_size;
self.remove_page_version(key.pgno, key.commit_seq);
self.directory.remove(&key);
self.version_coalesce_count = self.version_coalesce_count.saturating_add(1);
return true;
}
false
}
fn coalesce_all_versions(&mut self) {
let mut candidates: HashMap<PageNumber, Vec<(u64, CacheKey, SlabIdx, bool)>> =
HashMap::new();
for (idx, page) in self.t1.iter() {
if page.key.commit_seq != CommitSeq::ZERO && page.key.commit_seq <= self.gc_horizon {
candidates.entry(page.key.pgno).or_default().push((
page.key.commit_seq.get(),
page.key,
idx,
true, ));
}
}
for (idx, page) in self.t2.iter() {
if page.key.commit_seq != CommitSeq::ZERO && page.key.commit_seq <= self.gc_horizon {
candidates.entry(page.key.pgno).or_default().push((
page.key.commit_seq.get(),
page.key,
idx,
false, ));
}
}
let mut removed = 0;
for (_, mut versions) in candidates {
if versions.len() > 1 {
versions.sort_by_key(|entry| std::cmp::Reverse(entry.0));
for (_, key, idx, is_t1) in versions.into_iter().skip(1) {
let is_pinned = if is_t1 {
self.t1.get(idx).is_some_and(|p| p.is_pinned())
} else {
self.t2.get(idx).is_some_and(|p| p.is_pinned())
};
if is_pinned {
continue;
}
let page = if is_t1 {
self.t1.remove(idx)
} else {
self.t2.remove(idx)
};
self.total_bytes = self.total_bytes.saturating_sub(page.byte_size);
self.remove_page_version(key.pgno, key.commit_seq);
self.directory.remove(&key);
removed += 1;
}
}
}
let removed_u64 = u64::try_from(removed).unwrap_or(u64::MAX);
self.version_coalesce_count = self.version_coalesce_count.saturating_add(removed_u64);
}
#[allow(dead_code)]
fn coalesce_versions_for_pgno(&mut self, pgno: PageNumber) -> usize {
#[derive(Clone, Copy)]
enum ResidentList {
T1,
T2,
}
let mut committed: Vec<(u64, ResidentList, SlabIdx, CacheKey, bool)> = Vec::new();
for (idx, page) in self.t1.iter() {
if page.key.pgno == pgno
&& page.key.commit_seq != CommitSeq::ZERO
&& page.key.commit_seq <= self.gc_horizon
{
committed.push((
page.key.commit_seq.get(),
ResidentList::T1,
idx,
page.key,
page.is_pinned(),
));
}
}
for (idx, page) in self.t2.iter() {
if page.key.pgno == pgno
&& page.key.commit_seq != CommitSeq::ZERO
&& page.key.commit_seq <= self.gc_horizon
{
committed.push((
page.key.commit_seq.get(),
ResidentList::T2,
idx,
page.key,
page.is_pinned(),
));
}
}
if committed.len() <= 1 {
return 0;
}
committed.sort_by_key(|entry| std::cmp::Reverse(entry.0));
let mut removed = 0;
for (_, resident_list, idx, key, is_pinned) in committed.into_iter().skip(1) {
if is_pinned {
continue;
}
let page = match resident_list {
ResidentList::T1 => self.t1.remove(idx),
ResidentList::T2 => self.t2.remove(idx),
};
self.total_bytes = self.total_bytes.saturating_sub(page.byte_size);
self.remove_page_version(key.pgno, key.commit_seq);
self.directory.remove(&key);
removed += 1;
}
let removed_u64 = u64::try_from(removed).unwrap_or(u64::MAX);
self.version_coalesce_count = self.version_coalesce_count.saturating_add(removed_u64);
removed
}
fn remove_page_version(&mut self, pgno: PageNumber, seq: CommitSeq) {
if let Some(versions) = self.page_versions.get_mut(&pgno) {
versions.retain(|&v| v != seq);
if versions.is_empty() {
self.page_versions.remove(&pgno);
}
}
}
fn trim_ghosts(&mut self) {
while self.b1.len() > self.capacity {
if let Some(ghost) = self.b1.pop_front() {
self.directory.remove(&ghost);
} else {
break;
}
}
while self.b2.len() > self.capacity {
if let Some(ghost) = self.b2.pop_front() {
self.directory.remove(&ghost);
} else {
break;
}
}
}
pub fn shrink_memory(&mut self) {
self.coalesce_all_versions();
}
pub fn notify_unpin(&mut self) {
while self.capacity_overflow_events > 0
&& (self.len() > self.capacity
|| (self.max_bytes > 0 && self.total_bytes > self.max_bytes))
{
if self.evict_one_preferred() {
self.capacity_overflow_events -= 1;
} else {
break;
}
}
}
}
impl std::fmt::Debug for ArcCacheInner {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ArcCacheInner")
.field("t1_len", &self.t1.len())
.field("t2_len", &self.t2.len())
.field("b1_len", &self.b1.len())
.field("b2_len", &self.b2.len())
.field("p", &self.p)
.field("capacity", &self.capacity)
.field("total_bytes", &self.total_bytes)
.field("max_bytes", &self.max_bytes)
.field("directory_len", &self.directory.len())
.field("page_versions_len", &self.page_versions.len())
.field("gc_horizon", &self.gc_horizon)
.field("capacity_overflow_events", &self.capacity_overflow_events)
.field("hits", &self.hits)
.field("misses", &self.misses)
.field("ghost_hits_b1", &self.ghost_hits_b1)
.field("ghost_hits_b2", &self.ghost_hits_b2)
.field("evictions_t1", &self.evictions_t1)
.field("evictions_t2", &self.evictions_t2)
.field("version_coalesce_count", &self.version_coalesce_count)
.field("admits", &self.admits)
.finish()
}
}
pub struct ArcCache {
shards: Box<[Mutex<ArcCacheInner>]>,
shard_mask: usize,
inflight: Mutex<HashMap<CacheKey, Arc<InflightLoad>>>,
}
impl ArcCache {
const SHARD_COUNT: usize = 16;
#[must_use]
pub fn new(capacity: usize, max_bytes: usize) -> Self {
let shard_capacity = (capacity / Self::SHARD_COUNT).max(1);
let shard_max_bytes = max_bytes / Self::SHARD_COUNT;
let mut shards = Vec::with_capacity(Self::SHARD_COUNT);
for _ in 0..Self::SHARD_COUNT {
shards.push(Mutex::new(ArcCacheInner::new(
shard_capacity,
shard_max_bytes,
)));
}
Self {
shards: shards.into_boxed_slice(),
shard_mask: Self::SHARD_COUNT - 1,
inflight: Mutex::new(HashMap::new()),
}
}
#[inline]
fn select_shard(&self, key: &CacheKey) -> &Mutex<ArcCacheInner> {
use std::hash::{BuildHasher, Hasher};
let mut hasher = foldhash::fast::FixedState::default().build_hasher();
std::hash::Hash::hash(key, &mut hasher);
let hash = hasher.finish() as usize;
&self.shards[hash & self.shard_mask]
}
fn resident_pages_snapshot(&self) -> Option<usize> {
let mut total = 0usize;
for shard in &self.shards {
total = total.checked_add(shard.try_lock()?.len())?;
}
Some(total)
}
fn inflight_count_snapshot(&self) -> Option<usize> {
Some(self.inflight.try_lock()?.len())
}
#[must_use]
pub fn get(&self, key: &CacheKey) -> Option<Arc<CachedPage>> {
let shard = self.select_shard(key);
let inner = shard.lock();
inner.get(key)
}
pub fn request_async<F, E>(&self, key: CacheKey, loader: F) -> Result<AsyncLookup, E>
where
F: FnOnce() -> Result<CachedPage, E> + std::panic::UnwindSafe,
{
match self.claim_inflight_slot(key) {
InflightRole::Leader(slot) => self.lead_request_async(key, &slot, loader),
InflightRole::Waiter(slot) => Ok(self.wait_for_peer_load(key, &slot)),
}
}
fn claim_inflight_slot(&self, key: CacheKey) -> InflightRole {
let existing = self.inflight.lock().get(&key).cloned();
if let Some(existing) = existing {
return InflightRole::Waiter(existing);
}
let slot = Arc::new(InflightLoad::new_loading());
let mut inflight = self.inflight.lock();
if let Some(existing) = inflight.get(&key).cloned() {
return InflightRole::Waiter(existing);
}
inflight.insert(key, Arc::clone(&slot));
drop(inflight);
InflightRole::Leader(slot)
}
fn wait_for_peer_load(&self, key: CacheKey, slot: &Arc<InflightLoad>) -> AsyncLookup {
Self::wait_on_slot(slot.as_ref());
let is_hit = {
let shard = self.select_shard(&key);
let mut inner = shard.lock();
if inner.get(&key).is_some() {
let lookup = inner.request(&key);
debug_assert!(matches!(lookup, CacheLookup::Hit));
true
} else {
false
}
};
if is_hit {
AsyncLookup::WaitedForPeerHit
} else {
AsyncLookup::WaitedForPeerMiss
}
}
fn lead_request_async<F, E>(
&self,
key: CacheKey,
slot: &Arc<InflightLoad>,
loader: F,
) -> Result<AsyncLookup, E>
where
F: FnOnce() -> Result<CachedPage, E> + std::panic::UnwindSafe,
{
let is_hit = {
let shard = self.select_shard(&key);
let mut inner = shard.lock();
if inner.get(&key).is_some() {
let lookup = inner.request(&key);
debug_assert!(matches!(lookup, CacheLookup::Hit));
true
} else {
false
}
};
if is_hit {
self.release_inflight_slot(key, slot);
return Ok(AsyncLookup::Hit);
}
let load_result = catch_unwind(AssertUnwindSafe(loader));
match load_result {
Ok(Ok(page)) => {
{
let shard = self.select_shard(&key);
let mut inner = shard.lock();
if inner.get(&key).is_some() {
let _ = inner.request(&key);
} else {
let lookup = inner.request(&key);
debug_assert_eq!(page.key, key, "request_async loader returned wrong key");
inner.admit(key, page, lookup);
}
}
self.release_inflight_slot(key, slot);
Ok(AsyncLookup::Loaded)
}
Ok(Err(err)) => {
self.release_inflight_slot(key, slot);
Err(err)
}
Err(payload) => {
self.release_inflight_slot(key, slot);
resume_unwind(payload)
}
}
}
fn wait_on_slot(slot: &InflightLoad) {
let mut state = slot.state.lock();
state.waiters = state.waiters.saturating_add(1);
while state.loading {
slot.cv.wait(&mut state);
}
state.waiters = state.waiters.saturating_sub(1);
}
fn release_inflight_slot(&self, key: CacheKey, slot: &Arc<InflightLoad>) {
{
let mut inflight = self.inflight.lock();
if let Some(existing) = inflight.get(&key) {
if Arc::ptr_eq(existing, slot) {
inflight.remove(&key);
}
}
}
{
let mut state = slot.state.lock();
state.loading = false;
}
slot.cv.notify_all();
}
#[cfg(test)]
fn inflight_count(&self) -> usize {
self.inflight.lock().len()
}
}
impl std::fmt::Debug for ArcCache {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("ArcCache")
.field("shards_len", &self.shards.len())
.field("shard_mask", &self.shard_mask)
.field("resident_pages", &self.resident_pages_snapshot())
.field("inflight_len", &self.inflight_count_snapshot())
.finish_non_exhaustive()
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering};
use std::thread;
use std::time::Duration;
use fsqlite_types::PageSize;
const BEAD_ID: &str = "bd-125g";
const BEAD_ID_BD_3JK9: &str = "bd-3jk9";
const BEAD_ID_BD_1ZLA: &str = "bd-1zla";
const BEAD_ID_BD_7PU_1: &str = "bd-7pu.1";
fn key(pgno: u32, commit_seq: u64) -> CacheKey {
CacheKey::new(PageNumber::new(pgno).unwrap(), CommitSeq::new(commit_seq))
}
fn page(k: CacheKey, size: usize) -> CachedPage {
let ps = if size <= 512 {
PageSize::new(512).unwrap()
} else if size <= 1024 {
PageSize::new(1024).unwrap()
} else if size <= 2048 {
PageSize::new(2048).unwrap()
} else if size <= 4096 {
PageSize::new(4096).unwrap()
} else if size <= 8192 {
PageSize::new(8192).unwrap()
} else if size <= 16384 {
PageSize::new(16384).unwrap()
} else if size <= 32768 {
PageSize::new(32768).unwrap()
} else {
PageSize::new(65536).unwrap()
};
let mut cp = CachedPage::new(k, fsqlite_types::PageData::zeroed(ps), 0, None);
cp.byte_size = size;
cp
}
fn lock_shard_for_key<'a>(
cache: &'a ArcCache,
key: &CacheKey,
) -> fsqlite_types::sync_primitives::MutexGuard<'a, ArcCacheInner> {
cache.select_shard(key).lock()
}
fn total_cached_pages(cache: &ArcCache) -> usize {
cache.shards.iter().map(|shard| shard.lock().len()).sum()
}
#[test]
fn test_cache_key_mvcc_awareness() {
let k1 = key(1, 0);
let k2 = key(1, 1);
let k3 = key(2, 0);
assert_ne!(k1, k2, "bead_id={BEAD_ID} case=same_page_diff_seq");
assert_ne!(k1, k3, "bead_id={BEAD_ID} case=diff_page_same_seq");
assert_ne!(k2, k3, "bead_id={BEAD_ID} case=diff_page_diff_seq");
let k4 = key(1, 0);
assert_eq!(k1, k4, "bead_id={BEAD_ID} case=same_key_equal");
let mut map = HashMap::new();
map.insert(k1, "v1");
map.insert(k2, "v2");
map.insert(k3, "v3");
assert_eq!(map.len(), 3, "bead_id={BEAD_ID} case=hashmap_distinct_keys");
}
#[test]
fn test_request_t1_to_t2_promotion() {
let mut cache = ArcCacheInner::new(10, 0);
let k = key(1, 0);
let lookup = cache.request(&k);
assert_eq!(
lookup,
CacheLookup::Miss,
"bead_id={BEAD_ID} case=first_access_miss"
);
cache.admit(k, page(k, 4096), lookup);
assert_eq!(
cache.t1_len(),
1,
"bead_id={BEAD_ID} case=in_t1_after_admit"
);
assert_eq!(
cache.t2_len(),
0,
"bead_id={BEAD_ID} case=t2_empty_after_admit"
);
let lookup = cache.request(&k);
assert_eq!(
lookup,
CacheLookup::Hit,
"bead_id={BEAD_ID} case=second_access_hit"
);
assert_eq!(
cache.t1_len(),
0,
"bead_id={BEAD_ID} case=t1_empty_after_promotion"
);
assert_eq!(
cache.t2_len(),
1,
"bead_id={BEAD_ID} case=in_t2_after_promotion"
);
let lookup = cache.request(&k);
assert_eq!(
lookup,
CacheLookup::Hit,
"bead_id={BEAD_ID} case=third_access_t2_hit"
);
assert_eq!(
cache.t2_len(),
1,
"bead_id={BEAD_ID} case=stays_in_t2_after_refresh"
);
}
#[test]
fn test_request_t2_refresh() {
let mut cache = ArcCacheInner::new(4, 0);
let k1 = key(1, 0);
let k2 = key(2, 0);
let l = cache.request(&k1);
cache.admit(k1, page(k1, 4096), l);
let l = cache.request(&k2);
cache.admit(k2, page(k2, 4096), l);
assert_eq!(cache.request(&k1), CacheLookup::Hit);
assert_eq!(cache.request(&k2), CacheLookup::Hit);
assert_eq!(cache.t2_lru_key(), Some(k1));
assert_eq!(cache.t2_mru_key(), Some(k2));
assert_eq!(cache.request(&k1), CacheLookup::Hit);
assert_eq!(cache.t2_lru_key(), Some(k2));
assert_eq!(cache.t2_mru_key(), Some(k1));
}
#[test]
fn test_request_miss_inserts_t1() {
let mut cache = ArcCacheInner::new(2, 0);
let k = key(1, 0);
let lookup = cache.request(&k);
assert_eq!(lookup, CacheLookup::Miss);
cache.admit(k, page(k, 4096), lookup);
assert!(cache.in_t1(&k), "bead_id={BEAD_ID} case=miss_goes_to_t1");
assert_eq!(cache.t1_len(), 1);
assert_eq!(cache.t2_len(), 0);
}
#[test]
fn test_replace_prefers_t1_when_over_p() {
let mut cache = ArcCacheInner::new(3, 0);
let k1 = key(1, 0);
let k2 = key(2, 0);
let k3 = key(3, 0);
let k4 = key(4, 0);
let l = cache.request(&k1);
cache.admit(k1, page(k1, 4096), l);
cache.request(&k1);
for &k in &[k2, k3] {
let l = cache.request(&k);
cache.admit(k, page(k, 4096), l);
}
cache.set_p_for_tests(1);
let l = cache.request(&k4);
cache.admit(k4, page(k4, 4096), l);
assert!(cache.in_b1(&k2), "bead_id={BEAD_ID} case=prefer_t1");
assert!(cache.get(&k4).is_some());
}
#[test]
fn test_replace_b2_tiebreaker() {
let mut cache = ArcCacheInner::new(3, 0);
let k1 = key(1, 0);
let k2 = key(2, 0);
let k3 = key(3, 0);
let k4 = key(4, 0);
for &k in &[k1, k2] {
let l = cache.request(&k);
cache.admit(k, page(k, 4096), l);
cache.request(&k);
}
let l = cache.request(&k3);
cache.admit(k3, page(k3, 4096), l);
cache.set_p_for_tests(1);
let l = cache.request(&k4);
cache.admit(k4, page(k4, 4096), l);
assert!(cache.in_b2(&k1), "setup requires B2 ghost");
let t1_before = cache.t1_lru_key().expect("T1 must be non-empty"); cache.set_p_for_tests(cache.t1_len());
let l = cache.request(&k1);
assert_eq!(l, CacheLookup::GhostHitB2);
cache.admit(k1, page(k1, 4096), l);
assert!(
cache.in_b1(&t1_before),
"bead_id={BEAD_ID} case=b2_tiebreaker_prefers_t1"
);
}
#[test]
fn test_replace_fallback() {
let mut cache = ArcCacheInner::new(2, 0);
let pinned = key(1, 0);
let victim_t2 = key(2, 0);
let incoming = key(3, 0);
let l = cache.request(&pinned);
cache.admit(pinned, page(pinned, 4096), l);
let l = cache.request(&victim_t2);
cache.admit(victim_t2, page(victim_t2, 4096), l);
assert_eq!(cache.request(&victim_t2), CacheLookup::Hit);
cache.set_p_for_tests(0); cache.get(&pinned).expect("pinned page exists").pin();
let l = cache.request(&incoming);
cache.admit(incoming, page(incoming, 4096), l);
assert!(cache.get(&pinned).is_some(), "pinned T1 entry must remain");
assert!(
cache.get(&victim_t2).is_none(),
"fallback should evict from T2 when preferred T1 victim is pinned"
);
assert!(cache.get(&incoming).is_some());
let _ = cache.unpin(&pinned);
}
#[test]
fn test_replace_overflow_safety_valve() {
let mut cache = ArcCacheInner::new(1, 0);
let pinned = key(1, 0);
let incoming = key(2, 0);
let l = cache.request(&pinned);
cache.admit(pinned, page(pinned, 4096), l);
cache.get(&pinned).expect("pinned page exists").pin();
let l = cache.request(&incoming);
cache.admit(incoming, page(incoming, 4096), l);
assert_eq!(cache.capacity_overflow_events(), 1);
assert_eq!(cache.len(), 2, "safety valve allows temporary growth");
cache.get(&pinned).expect("pinned page exists").unpin();
}
#[test]
fn test_request_ghost_trim() {
let mut cache = ArcCacheInner::new(2, 0);
for pgno in 1..=10_u32 {
let k = key(pgno, 0);
let l = cache.request(&k);
if !matches!(l, CacheLookup::Hit) {
cache.admit(k, page(k, 4096), l);
}
}
assert!(cache.b1_len() <= cache.capacity);
assert!(cache.b2_len() <= cache.capacity);
}
#[test]
fn test_request_b1_ghost_increases_p() {
let mut cache = ArcCacheInner::new(2, 0);
let k1 = key(1, 0);
let k2 = key(2, 0);
let k3 = key(3, 0);
let l = cache.request(&k1);
cache.admit(k1, page(k1, 4096), l);
cache.request(&k1);
let l = cache.request(&k2);
cache.admit(k2, page(k2, 4096), l);
let l = cache.request(&k3);
cache.admit(k3, page(k3, 4096), l);
let p_before = cache.p();
let lookup_k2 = cache.request(&k2);
assert_eq!(
lookup_k2,
CacheLookup::GhostHitB1,
"bead_id={BEAD_ID} case=ghost_hit_b1"
);
assert!(
cache.p() > p_before,
"bead_id={BEAD_ID} case=p_increased_on_b1_hit p_before={p_before} p_after={}",
cache.p()
);
}
#[test]
fn test_request_b2_ghost_decreases_p() {
let mut cache = ArcCacheInner::new(3, 0);
let k1 = key(1, 0);
let k2 = key(2, 0);
let k3 = key(3, 0);
let k4 = key(4, 0);
for &k in &[k1, k2] {
let l = cache.request(&k);
cache.admit(k, page(k, 4096), l);
cache.request(&k); }
let l = cache.request(&k3);
cache.admit(k3, page(k3, 4096), l);
let l = cache.request(&k4);
cache.admit(k4, page(k4, 4096), l);
let l = cache.request(&k3);
assert_eq!(l, CacheLookup::GhostHitB1, "setup: B1 ghost hit");
cache.admit(k3, page(k3, 4096), l);
let p_before = cache.p();
assert!(p_before > 0, "setup: p > 0 before B2 ghost hit");
let lookup_k1 = cache.request(&k1);
assert_eq!(
lookup_k1,
CacheLookup::GhostHitB2,
"bead_id={BEAD_ID} case=ghost_hit_b2"
);
assert!(
cache.p() < p_before,
"bead_id={BEAD_ID} case=p_decreased_on_b2_hit p_before={p_before} p_after={}",
cache.p()
);
}
#[test]
fn test_scan_resistance() {
let mut cache = ArcCacheInner::new(10, 0);
let hot_keys: Vec<CacheKey> = (1..=5).map(|i| key(i, 0)).collect();
for &k in &hot_keys {
let l = cache.request(&k);
cache.admit(k, page(k, 4096), l);
cache.request(&k); }
assert_eq!(cache.t2_len(), 5, "bead_id={BEAD_ID} case=hot_set_in_t2");
for i in 100..120u32 {
let k = key(i, 0);
let l = cache.request(&k);
if !matches!(l, CacheLookup::Hit) {
cache.admit(k, page(k, 4096), l);
}
}
for &k in &hot_keys {
assert!(
cache.get(&k).is_some(),
"bead_id={BEAD_ID} case=hot_page_survived_scan pgno={}",
k.pgno
);
}
}
#[test]
fn test_replace_skips_pinned() {
let mut cache = ArcCacheInner::new(2, 0);
let k1 = key(1, 0);
let k2 = key(2, 0);
let l = cache.request(&k1);
cache.admit(k1, page(k1, 4096), l);
cache.get(&k1).unwrap().pin();
let l = cache.request(&k2);
cache.admit(k2, page(k2, 4096), l);
let k3 = key(3, 0);
let l = cache.request(&k3);
cache.admit(k3, page(k3, 4096), l);
assert!(
cache.get(&k1).is_some(),
"bead_id={BEAD_ID} case=pinned_page_not_evicted"
);
cache.get(&k1).unwrap().unpin();
}
#[test]
fn test_eviction_no_io() {
let mut cache = ArcCacheInner::new(3, 0);
for i in 1..=3u32 {
let k = key(i, 0);
let l = cache.request(&k);
cache.admit(k, page(k, 4096), l);
}
assert_eq!(cache.len(), 3);
let k4 = key(4, 0);
let l = cache.request(&k4);
cache.admit(k4, page(k4, 4096), l);
assert_eq!(
cache.len(),
3,
"bead_id={BEAD_ID} case=eviction_kept_capacity"
);
assert!(
cache.get(&k4).is_some(),
"bead_id={BEAD_ID} case=new_page_admitted_after_eviction"
);
}
#[test]
fn test_superseded_version_preferred() {
let mut cache = ArcCacheInner::new(3, 0);
let k_old = key(1, 1);
let l = cache.request(&k_old);
cache.admit(k_old, page(k_old, 4096), l);
let k_new = key(1, 5);
let l = cache.request(&k_new);
cache.admit(k_new, page(k_new, 4096), l);
let k_other = key(2, 0);
let l = cache.request(&k_other);
cache.admit(k_other, page(k_other, 4096), l);
cache.set_gc_horizon(CommitSeq::new(3));
let k_trigger = key(3, 0);
let l = cache.request(&k_trigger);
cache.admit(k_trigger, page(k_trigger, 4096), l);
assert!(
cache.get(&k_old).is_none(),
"bead_id={BEAD_ID} case=old_version_evicted"
);
assert!(
cache.get(&k_new).is_some(),
"bead_id={BEAD_ID} case=new_version_retained"
);
}
#[test]
fn test_memory_accounting() {
let mut cache = ArcCacheInner::new(100, 0);
let sizes = [1024_usize, 2048, 4096, 8192];
let mut expected_total = 0_usize;
for (i, &size) in sizes.iter().enumerate() {
let k = key(u32::try_from(i + 1).unwrap(), 0);
let l = cache.request(&k);
cache.admit(k, page(k, size), l);
expected_total += size;
assert_eq!(
cache.total_bytes(),
expected_total,
"bead_id={BEAD_ID} case=accounting_after_insert_{i}"
);
}
let mut small_cache = ArcCacheInner::new(2, 0);
let k1 = key(1, 0);
let k2 = key(2, 0);
let k3 = key(3, 0);
let l = small_cache.request(&k1);
small_cache.admit(k1, page(k1, 4096), l);
assert_eq!(
small_cache.total_bytes(),
4096,
"bead_id={BEAD_ID} case=accounting_single"
);
let l = small_cache.request(&k2);
small_cache.admit(k2, page(k2, 2048), l);
assert_eq!(
small_cache.total_bytes(),
6144,
"bead_id={BEAD_ID} case=accounting_two"
);
let l = small_cache.request(&k3);
small_cache.admit(k3, page(k3, 1024), l);
assert_eq!(
small_cache.len(),
2,
"bead_id={BEAD_ID} case=accounting_after_eviction_count"
);
assert_eq!(
small_cache.total_bytes(),
3072,
"bead_id={BEAD_ID} case=accounting_after_eviction_bytes"
);
}
#[test]
fn test_e2e_arc_cache_behavior_under_mixed_workload() {
let cache = ArcCache::new(128, 0);
for i in 1..=5u32 {
let k = key(i, 1);
let mut inner = lock_shard_for_key(&cache, &k);
let l = inner.request(&k);
inner.admit(k, page(k, 4096), l);
}
for i in 1..=3u32 {
let k = key(i, 1);
let mut inner = lock_shard_for_key(&cache, &k);
let l = inner.request(&k);
assert_eq!(
l,
CacheLookup::Hit,
"bead_id={BEAD_ID} case=e2e_txn2_read_hit page={i}"
);
inner.get(&k).unwrap().pin();
}
for i in 6..=8u32 {
let k = key(i, 2);
let mut inner = lock_shard_for_key(&cache, &k);
let l = inner.request(&k);
inner.admit(k, page(k, 4096), l);
}
for i in 1..=3u32 {
let k = key(i, 1);
let inner = lock_shard_for_key(&cache, &k);
inner.get(&k).unwrap().unpin();
}
assert_eq!(
total_cached_pages(&cache),
8,
"bead_id={BEAD_ID} case=e2e_total_pages"
);
for i in 1..=5u32 {
assert!(
cache.get(&key(i, 1)).is_some(),
"bead_id={BEAD_ID} case=e2e_page_{i}_v1_present"
);
}
for i in 6..=8u32 {
assert!(
cache.get(&key(i, 2)).is_some(),
"bead_id={BEAD_ID} case=e2e_page_{i}_v2_present"
);
}
}
#[test]
fn test_request_async_singleflight_duplicate_load_suppression() {
let cache = Arc::new(ArcCache::new(8, 0));
let load_calls = Arc::new(AtomicUsize::new(0));
let k = key(900, 1);
let mut handles = Vec::new();
for _ in 0..8 {
let cache_ref = Arc::clone(&cache);
let calls_ref = Arc::clone(&load_calls);
handles.push(thread::spawn(move || {
cache_ref
.request_async(k, || -> Result<CachedPage, &'static str> {
let _ = calls_ref.fetch_add(1, AtomicOrdering::SeqCst);
thread::sleep(Duration::from_millis(20));
Ok(page(k, 4096))
})
.expect("singleflight request should not fail")
}));
}
let mut waited_hits = 0usize;
for handle in handles {
let outcome = handle.join().expect("worker thread should not panic");
if matches!(outcome, AsyncLookup::WaitedForPeerHit) {
waited_hits += 1;
}
}
assert_eq!(
load_calls.load(AtomicOrdering::SeqCst),
1,
"bead_id={BEAD_ID_BD_7PU_1} case=singleflight_duplicate_load_suppression"
);
assert!(
waited_hits >= 1,
"bead_id={BEAD_ID_BD_7PU_1} case=singleflight_waiters_observed"
);
assert_eq!(
cache.inflight_count(),
0,
"bead_id={BEAD_ID_BD_7PU_1} case=singleflight_placeholder_cleared"
);
let mut inner = lock_shard_for_key(&cache, &k);
assert_eq!(
inner.request(&k),
CacheLookup::Hit,
"bead_id={BEAD_ID_BD_7PU_1} case=singleflight_page_admitted"
);
drop(inner);
}
#[test]
fn test_request_async_error_path_clears_placeholder() {
let cache = ArcCache::new(4, 0);
let k = key(901, 1);
let first = cache.request_async(k, || -> Result<CachedPage, &'static str> {
Err("loader failed")
});
assert_eq!(
first,
Err("loader failed"),
"bead_id={BEAD_ID_BD_7PU_1} case=error_propagated_to_leader"
);
assert_eq!(
cache.inflight_count(),
0,
"bead_id={BEAD_ID_BD_7PU_1} case=error_placeholder_cleared"
);
let mut inner = lock_shard_for_key(&cache, &k);
assert_eq!(
inner.request(&k),
CacheLookup::Miss,
"bead_id={BEAD_ID_BD_7PU_1} case=error_leaves_clean_miss"
);
drop(inner);
let retry = cache
.request_async(k, || -> Result<CachedPage, &'static str> {
Ok(page(k, 4096))
})
.expect("retry after error should succeed");
assert!(
matches!(retry, AsyncLookup::Loaded | AsyncLookup::Hit),
"bead_id={BEAD_ID_BD_7PU_1} case=error_retry_succeeds"
);
let mut inner = lock_shard_for_key(&cache, &k);
assert_eq!(
inner.request(&k),
CacheLookup::Hit,
"bead_id={BEAD_ID_BD_7PU_1} case=error_retry_admits_page"
);
drop(inner);
}
#[test]
fn test_request_async_hit_path_skips_loader() {
let cache = ArcCache::new(4, 0);
let k = key(903, 1);
{
let mut inner = lock_shard_for_key(&cache, &k);
let lookup = inner.request(&k);
inner.admit(k, page(k, 4096), lookup);
drop(inner);
}
let loader_calls = Arc::new(AtomicUsize::new(0));
let loader_calls_for_closure = Arc::clone(&loader_calls);
let outcome = cache
.request_async(k, move || -> Result<CachedPage, &'static str> {
loader_calls_for_closure.fetch_add(1, AtomicOrdering::SeqCst);
Err("waiter loader must not execute while placeholder is active")
})
.expect("hit path should not fail");
assert!(
matches!(outcome, AsyncLookup::Hit | AsyncLookup::Loaded),
"bead_id={BEAD_ID_BD_7PU_1} case=hit_path_skips_loader"
);
assert_eq!(
loader_calls.load(AtomicOrdering::SeqCst),
0,
"bead_id={BEAD_ID_BD_7PU_1} case=hit_path_loader_not_called"
);
}
#[test]
fn test_request_async_panic_notifies_waiters_and_allows_retry() {
let cache = Arc::new(ArcCache::new(4, 0));
let k = key(902, 1);
let cache_leader = Arc::clone(&cache);
let leader = thread::spawn(move || {
let panic_result = std::panic::catch_unwind(AssertUnwindSafe(|| {
let _ = cache_leader.request_async(k, || -> Result<CachedPage, &'static str> {
thread::sleep(Duration::from_millis(75));
resume_unwind(Box::new(
"forced loader unwind for cancellation path (test only)",
));
});
}));
panic_result.is_err()
});
while cache.inflight_count() == 0 {
thread::yield_now();
}
let waiter_loader_calls = Arc::new(AtomicUsize::new(0));
let waiter_loader_calls_for_thread = Arc::clone(&waiter_loader_calls);
let cache_waiter = Arc::clone(&cache);
let waiter = thread::spawn(move || {
cache_waiter
.request_async(k, move || -> Result<CachedPage, &'static str> {
waiter_loader_calls_for_thread.fetch_add(1, AtomicOrdering::SeqCst);
Err("waiter loader must not execute while placeholder is active")
})
.expect("waiter should observe peer outcome")
});
assert!(
leader.join().expect("leader thread join failed"),
"bead_id={BEAD_ID_BD_7PU_1} case=panic_expected"
);
let waiter_outcome = waiter.join().expect("waiter thread join failed");
assert_eq!(
waiter_outcome,
AsyncLookup::WaitedForPeerMiss,
"bead_id={BEAD_ID_BD_7PU_1} case=panic_waiter_unblocked"
);
assert_eq!(
waiter_loader_calls.load(AtomicOrdering::SeqCst),
0,
"bead_id={BEAD_ID_BD_7PU_1} case=waiter_loader_not_called"
);
assert_eq!(
cache.inflight_count(),
0,
"bead_id={BEAD_ID_BD_7PU_1} case=panic_placeholder_cleared"
);
let retry_outcome = cache
.request_async(k, || -> Result<CachedPage, &'static str> {
Ok(page(k, 4096))
})
.expect("retry after panic should succeed");
assert!(
matches!(retry_outcome, AsyncLookup::Loaded | AsyncLookup::Hit),
"bead_id={BEAD_ID_BD_7PU_1} case=panic_retry_succeeds"
);
}
#[test]
fn test_ghost_hit_exact_match() {
let mut cache = ArcCacheInner::new(2, 0);
let k_v1 = key(1, 1);
let k2 = key(2, 1);
let k3 = key(3, 1);
let lookup = cache.request(&k2);
cache.admit(k2, page(k2, 4096), lookup);
cache.request(&k2);
let lookup = cache.request(&k_v1);
cache.admit(k_v1, page(k_v1, 4096), lookup);
let lookup = cache.request(&k3);
cache.admit(k3, page(k3, 4096), lookup); assert!(cache.in_b1(&k_v1), "setup requires ghost entry in B1");
assert_eq!(
cache.request(&k_v1),
CacheLookup::GhostHitB1,
"bead_id={BEAD_ID_BD_3JK9} case=ghost_hit_exact_match"
);
}
#[test]
fn test_ghost_miss_different_version() {
let mut cache = ArcCacheInner::new(2, 0);
let k_v1 = key(1, 1);
let k_v2 = key(1, 2);
let k2 = key(2, 1);
let k3 = key(3, 1);
let lookup = cache.request(&k2);
cache.admit(k2, page(k2, 4096), lookup);
cache.request(&k2);
let lookup = cache.request(&k_v1);
cache.admit(k_v1, page(k_v1, 4096), lookup);
let lookup = cache.request(&k3);
cache.admit(k3, page(k3, 4096), lookup); assert!(cache.in_b1(&k_v1), "setup requires ghost entry in B1");
assert_eq!(
cache.request(&k_v2),
CacheLookup::Miss,
"bead_id={BEAD_ID_BD_3JK9} case=ghost_miss_different_version"
);
}
#[test]
fn test_ghost_prune_below_horizon() {
let mut cache = ArcCacheInner::new(2, 0);
let k1 = key(1, 1);
let k2 = key(2, 1);
let k3 = key(3, 1);
let lookup = cache.request(&k2);
cache.admit(k2, page(k2, 4096), lookup);
cache.request(&k2);
let lookup = cache.request(&k1);
cache.admit(k1, page(k1, 4096), lookup);
let lookup = cache.request(&k3);
cache.admit(k3, page(k3, 4096), lookup); assert!(cache.in_b1(&k1), "setup requires B1 ghost");
cache.set_gc_horizon(CommitSeq::new(2));
assert!(
!cache.in_b1(&k1),
"bead_id={BEAD_ID_BD_3JK9} case=ghost_prune_below_horizon"
);
assert!(!cache.in_b2(&k1));
}
#[test]
fn test_pinned_page_eviction_overflow() {
let mut cache = ArcCacheInner::new(1, 0);
let pinned = key(1, 1);
let incoming = key(2, 1);
let lookup = cache.request(&pinned);
cache.admit(pinned, page(pinned, 4096), lookup);
cache.get(&pinned).expect("pinned page exists").pin();
let lookup = cache.request(&incoming);
cache.admit(incoming, page(incoming, 4096), lookup);
assert_eq!(
cache.capacity_overflow_events(),
1,
"bead_id={BEAD_ID_BD_3JK9} case=pinned_page_eviction_overflow"
);
}
#[test]
fn test_overflow_decrement_on_unpin() {
let mut cache = ArcCacheInner::new(1, 0);
let pinned = key(1, 1);
let incoming = key(2, 1);
let lookup = cache.request(&pinned);
cache.admit(pinned, page(pinned, 4096), lookup);
cache.get(&pinned).expect("pinned page exists").pin();
let lookup = cache.request(&incoming);
cache.admit(incoming, page(incoming, 4096), lookup);
assert_eq!(
cache.capacity_overflow_events(),
1,
"setup requires overflow"
);
assert_eq!(cache.len(), 2, "setup requires temporary growth");
assert!(cache.unpin(&pinned), "unpin should observe prior pin");
assert_eq!(
cache.capacity_overflow_events(),
0,
"bead_id={BEAD_ID_BD_3JK9} case=overflow_decrement_on_unpin"
);
assert_eq!(cache.len(), 1);
}
#[test]
fn test_eviction_never_writes_wal() {
let mut cache = ArcCacheInner::new(2, 0);
for i in 1..=2_u32 {
let k = key(i, 1);
let lookup = cache.request(&k);
let mut p = page(k, 4096);
p.wal_frame = Some(i);
cache.admit(k, p, lookup);
}
let k3 = key(3, 1);
let lookup = cache.request(&k3);
let mut p3 = page(k3, 4096);
p3.wal_frame = Some(3);
cache.admit(k3, p3, lookup);
assert_eq!(cache.len(), 2);
assert!(
cache.t1.iter().all(|(_, page)| page.wal_frame.is_some())
&& cache.t2.iter().all(|(_, page)| page.wal_frame.is_some()),
"bead_id={BEAD_ID_BD_3JK9} case=eviction_never_writes_wal"
);
}
#[test]
fn test_uncommitted_pages_in_write_set() {
let mut cache = ArcCacheInner::new(2, 0);
let private_key = key(99, 0);
let mut write_set = HashMap::new();
write_set.insert(private_key, page(private_key, 4096));
assert_eq!(
cache.request(&private_key),
CacheLookup::Miss,
"bead_id={BEAD_ID_BD_3JK9} case=write_set_not_arc_hit"
);
assert!(cache.get(&private_key).is_none());
assert_eq!(cache.len(), 0);
assert_eq!(write_set.len(), 1);
}
#[test]
fn test_version_coalesce_keeps_newest() {
let mut cache = ArcCacheInner::new(8, 0);
let k1 = key(1, 1);
let k2 = key(1, 2);
let k3 = key(1, 3);
for &k in &[k1, k2, k3] {
let lookup = cache.request(&k);
cache.admit(k, page(k, 4096), lookup);
}
cache.set_gc_horizon(CommitSeq::new(3));
assert!(cache.get(&k3).is_some());
assert!(
cache.get(&k2).is_none() && cache.get(&k1).is_none(),
"bead_id={BEAD_ID_BD_3JK9} case=version_coalesce_keeps_newest"
);
}
#[test]
fn test_version_coalesce_removes_superseded() {
let mut cache = ArcCacheInner::new(8, 0);
let old = key(1, 1);
let new = key(1, 2);
for &k in &[old, new] {
let lookup = cache.request(&k);
cache.admit(k, page(k, 4096), lookup);
}
cache.set_gc_horizon(CommitSeq::new(2));
assert!(cache.get(&new).is_some());
assert!(
cache.get(&old).is_none(),
"bead_id={BEAD_ID_BD_3JK9} case=version_coalesce_removes_superseded"
);
}
#[test]
fn test_version_coalesce_skips_pinned() {
let mut cache = ArcCacheInner::new(8, 0);
let old = key(1, 1);
let new = key(1, 2);
for &k in &[old, new] {
let lookup = cache.request(&k);
cache.admit(k, page(k, 4096), lookup);
}
cache.get(&old).expect("old version must exist").pin();
cache.set_gc_horizon(CommitSeq::new(2));
assert!(
cache.get(&old).is_some(),
"bead_id={BEAD_ID_BD_3JK9} case=version_coalesce_skips_pinned"
);
let _ = cache.unpin(&old);
}
#[test]
fn test_coalesced_not_ghosted() {
let mut cache = ArcCacheInner::new(8, 0);
let old = key(1, 1);
let new = key(1, 2);
for &k in &[old, new] {
let lookup = cache.request(&k);
cache.admit(k, page(k, 4096), lookup);
}
cache.set_gc_horizon(CommitSeq::new(2));
assert!(cache.get(&old).is_none());
assert!(
!cache.in_b1(&old) && !cache.in_b2(&old),
"bead_id={BEAD_ID_BD_3JK9} case=coalesced_not_ghosted"
);
}
#[test]
fn test_coalesce_trigger_on_replace() {
let mut cache = ArcCacheInner::new(3, 0);
let old = key(1, 1);
let new = key(1, 2);
let other = key(2, 1);
let incoming = key(3, 1);
for &k in &[old, new, other] {
let lookup = cache.request(&k);
cache.admit(k, page(k, 4096), lookup);
}
assert_eq!(cache.request(&new), CacheLookup::Hit);
assert!(cache.in_t2(&new), "setup requires a T2 resident");
cache.get(&old).expect("old version must exist").pin();
cache.set_gc_horizon(CommitSeq::new(2));
assert!(
cache.get(&old).is_some(),
"pinned old version should survive batch"
);
let _ = cache.unpin(&old);
let lookup = cache.request(&incoming);
cache.admit(incoming, page(incoming, 4096), lookup);
assert!(
cache.get(&old).is_none(),
"bead_id={BEAD_ID_BD_3JK9} case=coalesce_trigger_on_replace"
);
assert!(!cache.in_b1(&old) && !cache.in_b2(&old));
assert!(cache.get(&incoming).is_some());
}
#[test]
fn test_coalesce_trigger_on_gc_advance() {
let mut cache = ArcCacheInner::new(8, 0);
let old = key(1, 1);
let new = key(1, 2);
for &k in &[old, new] {
let lookup = cache.request(&k);
cache.admit(k, page(k, 4096), lookup);
}
cache.set_gc_horizon(CommitSeq::new(2));
assert!(
cache.get(&old).is_none() && cache.get(&new).is_some(),
"bead_id={BEAD_ID_BD_3JK9} case=coalesce_trigger_on_gc_advance"
);
}
#[test]
fn test_e2e_bd_3jk9() {
let mut cache = ArcCacheInner::new(64, 0);
let mut commit_seq = 1_u64;
for writer in 0..8_u32 {
for op in 0..100_u32 {
let pgno = 1 + ((writer * 29 + op) % 24);
let key = key(pgno, commit_seq);
let lookup = cache.request(&key);
if !matches!(lookup, CacheLookup::Hit) {
cache.admit(key, page(key, 4096), lookup);
}
if op % 17 == 0 {
if let Some(page) = cache.get(&key) {
page.pin();
}
}
if op % 17 == 5 {
let _ = cache.unpin(&key);
}
if op % 20 == 0 {
let horizon = commit_seq.saturating_sub(10);
cache.set_gc_horizon(CommitSeq::new(horizon));
}
commit_seq = commit_seq.saturating_add(1);
}
}
assert!(cache.b1_len() <= cache.capacity);
assert!(cache.b2_len() <= cache.capacity);
assert!(cache.capacity_overflow_events() <= 200);
assert!(cache.len() <= cache.capacity + cache.capacity_overflow_events());
}
#[test]
fn test_visibility_committed_below_high() {
assert!(ArcCacheInner::is_visible(
CommitSeq::new(3),
CommitSeq::new(7)
));
}
#[test]
fn test_visibility_committed_above_high() {
assert!(!ArcCacheInner::is_visible(
CommitSeq::new(9),
CommitSeq::new(7)
));
}
#[test]
fn test_visibility_uncommitted() {
assert!(!ArcCacheInner::is_visible(
CommitSeq::ZERO,
CommitSeq::new(7)
));
}
#[test]
fn test_self_visibility_via_write_set() {
assert!(ArcCacheInner::is_visible_or_self(
CommitSeq::ZERO,
CommitSeq::new(7),
true
));
assert!(!ArcCacheInner::is_visible_or_self(
CommitSeq::ZERO,
CommitSeq::new(7),
false
));
}
#[test]
fn test_dual_eviction_by_count() {
let mut cache = ArcCacheInner::new(2, 0);
for pgno in 1..=3_u32 {
let k = key(pgno, 1);
let lookup = cache.request(&k);
cache.admit(k, page(k, 4096), lookup);
}
assert_eq!(cache.len(), 2);
}
#[test]
fn test_dual_eviction_by_bytes() {
let mut cache = ArcCacheInner::new(10, 5000);
let k1 = key(1, 1);
let k2 = key(2, 1);
let lookup = cache.request(&k1);
cache.admit(k1, page(k1, 4096), lookup);
let lookup = cache.request(&k2);
cache.admit(k2, page(k2, 4096), lookup);
assert!(
cache.total_bytes() <= 5000,
"bead_id={BEAD_ID_BD_1ZLA} case=dual_eviction_by_bytes"
);
}
#[test]
fn test_memory_accounting_delta_vs_full() {
let mut cache = ArcCacheInner::new(10, 0);
let full = key(1, 1);
let delta = key(2, 1);
let lookup = cache.request(&full);
cache.admit(full, page(full, 4096), lookup);
let lookup = cache.request(&delta);
cache.admit(delta, page(delta, 200), lookup);
assert_eq!(
cache.total_bytes(),
4296,
"bead_id={BEAD_ID_BD_1ZLA} case=memory_accounting_delta_vs_full"
);
}
#[test]
fn test_pragma_cache_size_positive() {
let mut cache = ArcCacheInner::new(1, 1);
cache.apply_pragma_cache_size(500, 4096);
assert_eq!(cache.capacity(), 500);
assert_eq!(cache.max_bytes(), 500 * 4096);
}
#[test]
fn test_pragma_cache_size_negative() {
let mut cache = ArcCacheInner::new(1, 1);
cache.apply_pragma_cache_size(-2000, 4096);
assert_eq!(cache.max_bytes(), 2_048_000);
assert_eq!(cache.capacity(), 500);
}
#[test]
fn test_pragma_cache_size_zero() {
let mut cache = ArcCacheInner::new(4, 0);
for pgno in 1..=4_u32 {
let k = key(pgno, 1);
let lookup = cache.request(&k);
cache.admit(k, page(k, 4096), lookup);
}
cache.apply_pragma_cache_size(0, 4096);
assert_eq!(cache.capacity(), 0);
assert_eq!(cache.max_bytes(), 0);
assert_eq!(cache.len(), 0);
}
#[test]
fn test_cache_resize_evicts() {
let mut cache = ArcCacheInner::new(4, 0);
for pgno in 1..=4_u32 {
let k = key(pgno, 1);
let lookup = cache.request(&k);
cache.admit(k, page(k, 4096), lookup);
}
cache.resize(2, 0);
assert!(cache.len() <= 2);
}
#[test]
fn test_cache_resize_trims_ghosts() {
let mut cache = ArcCacheInner::new(4, 0);
for pgno in 1..=12_u32 {
let k = key(pgno, 1);
let lookup = cache.request(&k);
if !matches!(lookup, CacheLookup::Hit) {
cache.admit(k, page(k, 4096), lookup);
}
}
cache.resize(1, 0);
assert!(cache.b1_len() <= 1);
assert!(cache.b2_len() <= 1);
}
#[test]
fn test_cache_resize_clamps_p() {
let mut cache = ArcCacheInner::new(8, 0);
cache.set_p_for_tests(8);
cache.resize(2, 0);
assert_eq!(cache.p(), 2);
}
#[test]
fn test_e2e_arc_scan_then_hotset() {
let mut cache = ArcCacheInner::new(16, 16 * 4096);
for i in 1_u32..=200 {
if i == 50 {
cache.apply_pragma_cache_size(64, 4096);
} else if i == 100 {
cache.apply_pragma_cache_size(-2000, 4096);
} else if i == 150 {
cache.apply_pragma_cache_size(0, 4096);
} else if i == 170 {
cache.apply_pragma_cache_size(32, 4096);
}
let pgno = 1 + ((i - 1) % 64);
let k = key(pgno, u64::from(i));
let size = if i % 3 == 0 { 200 } else { 4096 };
let lookup = cache.request(&k);
if !matches!(lookup, CacheLookup::Hit) {
cache.admit(k, page(k, size), lookup);
}
}
assert!(cache.b1_len() <= cache.capacity());
assert!(cache.b2_len() <= cache.capacity());
if cache.max_bytes() > 0 {
assert!(cache.total_bytes() <= cache.max_bytes());
}
}
#[test]
fn test_arc_cache_scan_resistance() {
test_scan_resistance();
}
#[test]
fn test_arc_hit_t1_promote_to_t2() {
test_request_t1_to_t2_promotion();
}
#[test]
fn test_arc_ghost_hit_b1_increases_p() {
test_request_b1_ghost_increases_p();
}
#[test]
fn test_arc_ghost_hit_b2_decreases_p() {
test_request_b2_ghost_decreases_p();
}
#[test]
fn test_replace_all_pinned_overflow() {
test_replace_overflow_safety_valve();
test_pinned_page_eviction_overflow();
}
#[test]
fn test_mvcc_keying_exact_match_ghosts() {
test_ghost_hit_exact_match();
test_ghost_miss_different_version();
}
#[test]
fn test_version_coalescing_drops_superseded() {
test_version_coalesce_removes_superseded();
}
#[test]
fn test_eviction_never_appends_wal() {
test_eviction_never_writes_wal();
test_eviction_no_io();
}
#[test]
fn test_arc_scan_then_hotset_acceptance() {
test_scan_resistance();
}
const BEAD_ID_BD_2ZOA: &str = "bd-2zoa";
struct Xorshift64 {
state: u64,
}
impl Xorshift64 {
fn new(seed: u64) -> Self {
Self {
state: if seed == 0 { 1 } else { seed },
}
}
fn next_u64(&mut self) -> u64 {
let mut x = self.state;
x ^= x << 13;
x ^= x >> 7;
x ^= x << 17;
self.state = x;
x
}
fn next_bounded(&mut self, bound: u64) -> u64 {
self.next_u64() % bound
}
}
struct ZipfSampler {
cdf: Vec<f64>,
}
impl ZipfSampler {
fn new(n: usize, s: f64) -> Self {
let mut weights = Vec::with_capacity(n);
let mut total = 0.0_f64;
for rank in 1..=n {
let w = 1.0 / (rank as f64).powf(s);
total += w;
weights.push(total);
}
for w in &mut weights {
*w /= total;
}
Self { cdf: weights }
}
fn sample(&self, rng: &mut Xorshift64) -> usize {
let u = (rng.next_u64() as f64) / (u64::MAX as f64);
self.cdf.partition_point(|&c| c < u)
}
}
fn run_workload(
cache: &mut ArcCacheInner,
accesses: impl Iterator<Item = CacheKey>,
) -> (usize, usize) {
let mut hits = 0_usize;
let mut total = 0_usize;
for k in accesses {
total += 1;
let lookup = cache.request(&k);
if matches!(lookup, CacheLookup::Hit) {
hits += 1;
} else {
cache.admit(k, page(k, 4096), lookup);
}
}
(hits, total)
}
#[test]
fn test_arc_oltp_hit_rate() {
let capacity = 2000_usize;
let hot_set = 500_u64;
let total_pages = 100_000_u64;
let num_accesses = 50_000_usize;
let mut cache = ArcCacheInner::new(capacity, 0);
let mut rng = Xorshift64::new(42);
#[allow(clippy::cast_possible_truncation)]
let hot_set_usize = hot_set as usize; for i in 1..=capacity.min(hot_set_usize) {
let k = key(u32::try_from(i).unwrap(), 0);
let l = cache.request(&k);
if !matches!(l, CacheLookup::Hit) {
cache.admit(k, page(k, 4096), l);
}
let _ = cache.request(&k);
}
let accesses = (0..num_accesses).map(|_| {
let pgno = if rng.next_bounded(100) < 95 {
u32::try_from(1 + rng.next_bounded(hot_set)).unwrap()
} else {
u32::try_from(1 + rng.next_bounded(total_pages)).unwrap()
};
key(pgno, 0)
});
let (hits, total) = run_workload(&mut cache, accesses);
let hit_rate = hits as f64 / total as f64;
eprintln!(
"INFO bead_id={BEAD_ID_BD_2ZOA} case=oltp_hit_rate \
hits={hits} total={total} hit_rate={hit_rate:.4} \
capacity={capacity} hot_set={hot_set}"
);
assert!(
hit_rate > 0.90,
"bead_id={BEAD_ID_BD_2ZOA} case=oltp_hit_rate \
expected>0.90 got={hit_rate:.4}"
);
}
#[test]
fn test_arc_mixed_hit_rate() {
let capacity = 2000_usize;
let hot_set = 500_u64;
let num_accesses = 40_000_usize;
let mut cache = ArcCacheInner::new(capacity, 0);
let mut rng = Xorshift64::new(7);
for i in 1..=hot_set {
let k = key(u32::try_from(i).unwrap(), 0);
let l = cache.request(&k);
if !matches!(l, CacheLookup::Hit) {
cache.admit(k, page(k, 4096), l);
}
let _ = cache.request(&k);
}
let mut scan_offset = 10_000_u32;
let accesses = (0..num_accesses).map(|_| {
if rng.next_bounded(100) < 70 {
let pgno = u32::try_from(1 + rng.next_bounded(hot_set)).unwrap();
key(pgno, 0)
} else {
scan_offset += 1;
if scan_offset > 50_000 {
scan_offset = 10_001;
}
key(scan_offset, 0)
}
});
let (hits, total) = run_workload(&mut cache, accesses);
let hit_rate = hits as f64 / total as f64;
eprintln!(
"INFO bead_id={BEAD_ID_BD_2ZOA} case=mixed_hit_rate \
hits={hits} total={total} hit_rate={hit_rate:.4}"
);
assert!(
hit_rate > 0.65,
"bead_id={BEAD_ID_BD_2ZOA} case=mixed_hit_rate \
expected>0.65 got={hit_rate:.4}"
);
}
#[test]
fn test_arc_scan_resistance_preserves_t2() {
let capacity = 100_usize;
let hot_count = 40_u32;
let mut cache = ArcCacheInner::new(capacity, 0);
for i in 1..=hot_count {
let k = key(i, 0);
let l = cache.request(&k);
cache.admit(k, page(k, 4096), l);
let _ = cache.request(&k); }
let t2_before = cache.t2_len();
assert_eq!(
t2_before, hot_count as usize,
"bead_id={BEAD_ID_BD_2ZOA} case=scan_resistance_t2_setup"
);
for i in 1000..1200_u32 {
let k = key(i, 0);
let l = cache.request(&k);
if !matches!(l, CacheLookup::Hit) {
cache.admit(k, page(k, 4096), l);
}
}
let mut survived = 0_u32;
for i in 1..=hot_count {
if cache.get(&key(i, 0)).is_some() {
survived += 1;
}
}
eprintln!(
"INFO bead_id={BEAD_ID_BD_2ZOA} case=scan_resistance \
survived={survived}/{hot_count} t2_before={t2_before} t2_after={}",
cache.t2_len()
);
assert!(
survived >= hot_count / 2,
"bead_id={BEAD_ID_BD_2ZOA} case=scan_resistance \
hot pages should survive scan: survived={survived}/{hot_count}"
);
}
#[test]
#[allow(clippy::cast_possible_truncation)]
fn test_arc_zipf_hit_rate() {
let capacity = 2000_usize;
let total_pages = 10_000_usize;
let num_accesses = 50_000_usize;
let mut cache = ArcCacheInner::new(capacity, 0);
let sampler = ZipfSampler::new(total_pages, 1.0);
let mut rng = Xorshift64::new(123);
let accesses = (0..num_accesses).map(|_| {
let rank = sampler.sample(&mut rng);
let pgno = u32::try_from(rank.min(total_pages - 1) + 1).unwrap();
key(pgno, 0)
});
let (hits, total) = run_workload(&mut cache, accesses);
let hit_rate = hits as f64 / total as f64;
eprintln!(
"INFO bead_id={BEAD_ID_BD_2ZOA} case=zipf_hit_rate \
hits={hits} total={total} hit_rate={hit_rate:.4} \
capacity={capacity} total_pages={total_pages}"
);
assert!(
hit_rate > 0.75,
"bead_id={BEAD_ID_BD_2ZOA} case=zipf_hit_rate \
expected>0.75 got={hit_rate:.4}"
);
}
#[test]
fn test_arc_mvcc_hit_rate() {
let capacity = 2000_usize;
let hot_set = 800_u64;
let num_writers = 8_u64;
let accesses_per_writer = 5_000_usize;
let mut cache = ArcCacheInner::new(capacity, 0);
let mut rng = Xorshift64::new(99);
let mut commit_seq = 1_u64;
for i in 1..=hot_set {
let k = key(i as u32, 0);
let l = cache.request(&k);
if !matches!(l, CacheLookup::Hit) {
cache.admit(k, page(k, 4096), l);
}
let _ = cache.request(&k);
}
let mut hits = 0_usize;
let mut total = 0_usize;
for _writer in 0..num_writers {
commit_seq += 1;
for _ in 0..accesses_per_writer {
total += 1;
let pgno = u32::try_from(1 + rng.next_bounded(hot_set)).unwrap();
let is_read = rng.next_bounded(100) < 80;
let k = key(pgno, commit_seq);
let lookup = cache.request(&k);
if matches!(lookup, CacheLookup::Hit) {
hits += 1;
} else if is_read {
let k0 = key(pgno, 0);
let l0 = cache.request(&k0);
if matches!(l0, CacheLookup::Hit) {
hits += 1;
} else {
cache.admit(k, page(k, 4096), lookup);
}
} else {
cache.admit(k, page(k, 4096), lookup);
}
}
}
let hit_rate = hits as f64 / total as f64;
eprintln!(
"INFO bead_id={BEAD_ID_BD_2ZOA} case=mvcc_hit_rate \
hits={hits} total={total} hit_rate={hit_rate:.4} \
writers={num_writers} hot_set={hot_set}"
);
assert!(
hit_rate > 0.30,
"bead_id={BEAD_ID_BD_2ZOA} case=mvcc_hit_rate \
expected>0.30 got={hit_rate:.4}"
);
}
#[test]
fn test_warmup_phase1_cold() {
let capacity = 100_usize;
let mut cache = ArcCacheInner::new(capacity, 0);
let half = capacity / 2;
let mut all_misses = true;
for i in 1..=half {
let k = key(u32::try_from(i).unwrap(), 0);
let lookup = cache.request(&k);
if matches!(lookup, CacheLookup::Hit) {
all_misses = false;
} else {
cache.admit(k, page(k, 4096), lookup);
}
}
eprintln!(
"INFO bead_id={BEAD_ID_BD_2ZOA} case=warmup_phase1_cold \
all_misses={all_misses} p={} len={} capacity={capacity}",
cache.p(),
cache.len()
);
assert!(
all_misses,
"bead_id={BEAD_ID_BD_2ZOA} case=warmup_phase1_cold \
first access to each unique page should be a miss"
);
assert_eq!(
cache.p(),
0,
"bead_id={BEAD_ID_BD_2ZOA} case=warmup_phase1_cold \
p should remain 0 during cold start (no ghost hits yet)"
);
}
#[test]
fn test_warmup_phase2_learning() {
let capacity = 50_usize;
let mut cache = ArcCacheInner::new(capacity, 0);
let mut rng = Xorshift64::new(77);
for i in 1..=capacity {
let k = key(u32::try_from(i).unwrap(), 0);
let l = cache.request(&k);
cache.admit(k, page(k, 4096), l);
}
let num_learning_accesses = capacity * 4;
let mut ghost_hits = 0_usize;
for _ in 0..num_learning_accesses {
let pgno = if rng.next_bounded(100) < 60 {
u32::try_from(1 + rng.next_bounded(capacity as u64)).unwrap()
} else {
u32::try_from(capacity as u64 + 1 + rng.next_bounded(capacity as u64 * 2)).unwrap()
};
let k = key(pgno, 0);
let lookup = cache.request(&k);
if matches!(lookup, CacheLookup::GhostHitB1 | CacheLookup::GhostHitB2) {
ghost_hits += 1;
cache.admit(k, page(k, 4096), lookup);
} else if !matches!(lookup, CacheLookup::Hit) {
cache.admit(k, page(k, 4096), lookup);
}
}
eprintln!(
"INFO bead_id={BEAD_ID_BD_2ZOA} case=warmup_phase2_learning \
p={} ghost_hits={ghost_hits} b1={} b2={} capacity={capacity}",
cache.p(),
cache.b1_len(),
cache.b2_len()
);
let ghosts_populated = cache.b1_len() > 0 || cache.b2_len() > 0;
assert!(
ghosts_populated,
"bead_id={BEAD_ID_BD_2ZOA} case=warmup_phase2_learning \
ghost lists should populate during learning phase"
);
}
#[test]
fn test_warmup_phase3_steady() {
let capacity = 200_usize;
let hot_set = 150_u64; let warmup_accesses = capacity * 3;
let measure_window = 2000_usize;
let mut cache = ArcCacheInner::new(capacity, 0);
let mut rng = Xorshift64::new(55);
for _ in 0..warmup_accesses {
let pgno = if rng.next_bounded(100) < 85 {
u32::try_from(1 + rng.next_bounded(hot_set)).unwrap()
} else {
u32::try_from(1 + rng.next_bounded(1000)).unwrap()
};
let k = key(pgno, 0);
let l = cache.request(&k);
if !matches!(l, CacheLookup::Hit) {
cache.admit(k, page(k, 4096), l);
}
}
let measure = |cache: &mut ArcCacheInner, rng: &mut Xorshift64, count: usize| -> f64 {
let mut hits = 0_usize;
for _ in 0..count {
let pgno = if rng.next_bounded(100) < 85 {
u32::try_from(1 + rng.next_bounded(hot_set)).unwrap()
} else {
u32::try_from(1 + rng.next_bounded(1000)).unwrap()
};
let k = key(pgno, 0);
let l = cache.request(&k);
if matches!(l, CacheLookup::Hit) {
hits += 1;
} else {
cache.admit(k, page(k, 4096), l);
}
}
hits as f64 / count as f64
};
let rate1 = measure(&mut cache, &mut rng, measure_window);
let rate2 = measure(&mut cache, &mut rng, measure_window);
let diff = (rate1 - rate2).abs();
eprintln!(
"INFO bead_id={BEAD_ID_BD_2ZOA} case=warmup_phase3_steady \
rate1={rate1:.4} rate2={rate2:.4} diff={diff:.4} p={}",
cache.p()
);
assert!(
diff < 0.10,
"bead_id={BEAD_ID_BD_2ZOA} case=warmup_phase3_steady \
hit rate should stabilize within 10%: rate1={rate1:.4} rate2={rate2:.4} diff={diff:.4}"
);
}
#[test]
fn test_prewarm_wal_index() {
let capacity = 100_usize;
let half_capacity = capacity / 2;
let wal_index_pages = 30_u32;
let mut cache = ArcCacheInner::new(capacity, 0);
let pages_to_load = wal_index_pages.min(u32::try_from(half_capacity).unwrap());
for i in 1..=pages_to_load {
let k = key(i, 0);
let l = cache.request(&k);
if !matches!(l, CacheLookup::Hit) {
cache.admit(k, page(k, 4096), l);
}
}
eprintln!(
"INFO bead_id={BEAD_ID_BD_2ZOA} case=prewarm_wal_index \
loaded={pages_to_load} t1={} t2={} capacity={capacity}",
cache.t1_len(),
cache.t2_len()
);
assert_eq!(
cache.t1_len(),
pages_to_load as usize,
"bead_id={BEAD_ID_BD_2ZOA} case=prewarm_wal_index \
WAL index pages should be in T1"
);
assert_eq!(
cache.t2_len(),
0,
"bead_id={BEAD_ID_BD_2ZOA} case=prewarm_wal_index \
T2 should be empty after pre-warming (single access)"
);
}
#[test]
fn test_prewarm_limited() {
let capacity = 100_usize;
let half_capacity = capacity / 2;
let requested_prewarm = 80_u32;
let mut cache = ArcCacheInner::new(capacity, 0);
let pages_to_load = requested_prewarm.min(u32::try_from(half_capacity).unwrap());
for i in 1..=pages_to_load {
let k = key(i, 0);
let l = cache.request(&k);
if !matches!(l, CacheLookup::Hit) {
cache.admit(k, page(k, 4096), l);
}
}
eprintln!(
"INFO bead_id={BEAD_ID_BD_2ZOA} case=prewarm_limited \
requested={requested_prewarm} loaded={pages_to_load} \
half_capacity={half_capacity}"
);
assert!(
cache.len() <= half_capacity,
"bead_id={BEAD_ID_BD_2ZOA} case=prewarm_limited \
pre-warming must load at most half capacity: loaded={} half={}",
cache.len(),
half_capacity
);
}
#[test]
fn test_prewarm_root_pages() {
let capacity = 100_usize;
let mut cache = ArcCacheInner::new(capacity, 0);
let root_pages: Vec<u32> = vec![1, 2, 5, 8];
for &pg in &root_pages {
let k = key(pg, 0);
let l = cache.request(&k);
if !matches!(l, CacheLookup::Hit) {
cache.admit(k, page(k, 4096), l);
}
}
for &pg in &root_pages {
let k = key(pg, 0);
assert!(
cache.get(&k).is_some(),
"bead_id={BEAD_ID_BD_2ZOA} case=prewarm_root_pages \
root page {pg} should be in cache"
);
}
let mut all_hits = true;
for &pg in &root_pages {
let k = key(pg, 0);
let l = cache.request(&k);
if !matches!(l, CacheLookup::Hit) {
all_hits = false;
}
}
eprintln!(
"INFO bead_id={BEAD_ID_BD_2ZOA} case=prewarm_root_pages \
root_pages_loaded={} all_hits={all_hits}",
root_pages.len()
);
assert!(
all_hits,
"bead_id={BEAD_ID_BD_2ZOA} case=prewarm_root_pages \
pre-warmed root pages should be cache hits"
);
}
#[test]
#[allow(clippy::cast_possible_truncation)]
fn test_e2e_bd_2zoa_arc_performance() {
let capacity = 500_usize;
let hot_set = 200_u64;
let mut cache = ArcCacheInner::new(capacity, 0);
let mut rng = Xorshift64::new(2026);
let mut cold_misses = 0_usize;
for i in 1..=(capacity / 2) {
let k = key(i as u32, 0);
let l = cache.request(&k);
if !matches!(l, CacheLookup::Hit) {
cold_misses += 1;
cache.admit(k, page(k, 4096), l);
}
}
assert_eq!(
cold_misses,
capacity / 2,
"bead_id={BEAD_ID_BD_2ZOA} case=e2e_phase1_all_misses"
);
for _ in 0..capacity * 2 {
let pgno = if rng.next_bounded(100) < 70 {
u32::try_from(1 + rng.next_bounded(hot_set)).unwrap()
} else {
u32::try_from(1 + rng.next_bounded(2000)).unwrap()
};
let k = key(pgno, 0);
let l = cache.request(&k);
if !matches!(l, CacheLookup::Hit) {
cache.admit(k, page(k, 4096), l);
}
}
let mut hits = 0_usize;
let measure_count = 5000_usize;
for _ in 0..measure_count {
let pgno = if rng.next_bounded(100) < 70 {
u32::try_from(1 + rng.next_bounded(hot_set)).unwrap()
} else {
u32::try_from(1 + rng.next_bounded(2000)).unwrap()
};
let k = key(pgno, 0);
let l = cache.request(&k);
if matches!(l, CacheLookup::Hit) {
hits += 1;
} else {
cache.admit(k, page(k, 4096), l);
}
}
let hit_rate = hits as f64 / measure_count as f64;
eprintln!(
"INFO bead_id={BEAD_ID_BD_2ZOA} case=e2e_arc_performance \
hit_rate={hit_rate:.4} hits={hits}/{measure_count} \
p={} t1={} t2={} b1={} b2={}",
cache.p(),
cache.t1_len(),
cache.t2_len(),
cache.b1_len(),
cache.b2_len()
);
assert!(
hit_rate > 0.50,
"bead_id={BEAD_ID_BD_2ZOA} case=e2e_arc_performance \
expected>0.50 got={hit_rate:.4}"
);
}
#[test]
fn test_evict_one_preferred_fallback() {
let mut cache = ArcCacheInner::new(10, 0);
let k1 = key(1, 0); let k2 = key(2, 0);
let l = cache.request(&k1);
cache.admit(k1, page(k1, 4096), l);
cache.request(&k1);
let l = cache.request(&k2);
cache.admit(k2, page(k2, 4096), l);
cache.get(&k2).unwrap().pin();
cache.resize(0, 1);
assert!(
cache.get(&k1).is_none(),
"k1 should be evicted from T2 because T1 was pinned"
);
assert!(cache.get(&k2).is_some(), "k2 must remain (pinned)");
assert_eq!(
cache.capacity_overflow_events(),
1,
"overflow only because we couldn't evict k2, but we DID evict k1"
);
cache.get(&k2).unwrap().unpin();
}
#[test]
fn test_cache_metrics_snapshot_hit_rate_pct() {
let snap = CacheMetricsSnapshot {
hits: 75,
misses: 25,
ghost_hits_b1: 0,
ghost_hits_b2: 0,
evictions_t1: 0,
evictions_t2: 0,
version_coalesce_count: 0,
admits: 0,
t1_len: 0,
t2_len: 0,
b1_len: 0,
b2_len: 0,
p: 0,
capacity: 0,
total_bytes: 0,
max_bytes: 0,
multi_version_pages: 0,
capacity_overflow_events: 0,
};
assert!((snap.hit_rate_pct() - 75.0).abs() < 0.01);
let zero = CacheMetricsSnapshot {
hits: 0,
misses: 0,
..snap
};
assert!((zero.hit_rate_pct() - 0.0).abs() < f64::EPSILON);
}
#[test]
fn test_cache_metrics_snapshot_total_accesses_and_resident() {
let snap = CacheMetricsSnapshot {
hits: 10,
misses: 5,
ghost_hits_b1: 3,
ghost_hits_b2: 2,
evictions_t1: 0,
evictions_t2: 0,
version_coalesce_count: 0,
admits: 0,
t1_len: 7,
t2_len: 4,
b1_len: 0,
b2_len: 0,
p: 0,
capacity: 0,
total_bytes: 0,
max_bytes: 0,
multi_version_pages: 0,
capacity_overflow_events: 0,
};
assert_eq!(snap.total_accesses(), 20);
assert_eq!(snap.resident_pages(), 11);
}
#[test]
fn test_cached_page_debug_format() {
let k = key(42, 7);
let cp = page(k, 4096);
let dbg = format!("{cp:?}");
assert!(dbg.contains("CachedPage"));
assert!(dbg.contains("data_len"));
assert!(dbg.contains("ref_count"));
}
#[test]
fn test_cached_page_pin_unpin_is_pinned() {
let k = key(99, 0);
let cp = page(k, 4096);
assert!(!cp.is_pinned());
cp.pin();
assert!(cp.is_pinned());
cp.pin();
assert!(cp.is_pinned());
cp.unpin();
assert!(cp.is_pinned());
cp.unpin();
assert!(!cp.is_pinned());
}
#[test]
fn cache_lookup_debug_clone_copy_eq() {
let variants = [
CacheLookup::Hit,
CacheLookup::GhostHitB1,
CacheLookup::GhostHitB2,
CacheLookup::Miss,
];
for v in &variants {
let copied = *v;
assert_eq!(copied, *v);
}
assert_ne!(CacheLookup::Hit, CacheLookup::Miss);
let dbg = format!("{:?}", CacheLookup::GhostHitB1);
assert!(dbg.contains("GhostHitB1"));
}
#[test]
fn async_lookup_debug_clone_copy_eq() {
let variants = [
AsyncLookup::Hit,
AsyncLookup::Loaded,
AsyncLookup::WaitedForPeerHit,
AsyncLookup::WaitedForPeerMiss,
];
for v in &variants {
let copied = *v;
assert_eq!(copied, *v);
}
assert_ne!(AsyncLookup::Hit, AsyncLookup::Loaded);
let dbg = format!("{:?}", AsyncLookup::WaitedForPeerMiss);
assert!(dbg.contains("WaitedForPeerMiss"));
}
#[test]
fn cache_key_hash_distinguishes_pgno_and_seq() {
use std::collections::HashSet;
let a = key(1, 0);
let b = key(1, 1);
let c = key(2, 0);
let mut set = HashSet::new();
set.insert(a);
set.insert(b);
set.insert(c);
assert_eq!(set.len(), 3);
assert!(set.contains(&key(1, 0)));
}
#[test]
fn cache_metrics_snapshot_debug_clone_copy() {
let snap = CacheMetricsSnapshot {
hits: 1,
misses: 2,
ghost_hits_b1: 0,
ghost_hits_b2: 0,
evictions_t1: 0,
evictions_t2: 0,
version_coalesce_count: 0,
admits: 0,
t1_len: 3,
t2_len: 4,
b1_len: 0,
b2_len: 0,
p: 5,
capacity: 10,
total_bytes: 100,
max_bytes: 200,
multi_version_pages: 0,
capacity_overflow_events: 0,
};
let dbg = format!("{snap:?}");
assert!(dbg.contains("CacheMetricsSnapshot"));
let copied = snap;
assert_eq!(copied, snap);
let cloned = snap.clone();
assert_eq!(cloned.hits, 1);
}
#[test]
fn arc_cache_inner_fresh_state_accessors() {
let inner = ArcCacheInner::new(128, 1024 * 1024);
assert_eq!(inner.capacity(), 128);
assert_eq!(inner.max_bytes(), 1024 * 1024);
assert_eq!(inner.len(), 0);
assert!(inner.is_empty());
assert_eq!(inner.total_bytes(), 0);
assert_eq!(inner.p(), 0);
assert_eq!(inner.t1_len(), 0);
assert_eq!(inner.t2_len(), 0);
assert_eq!(inner.b1_len(), 0);
assert_eq!(inner.b2_len(), 0);
assert_eq!(inner.capacity_overflow_events(), 0);
}
#[test]
fn cached_page_new_stores_fields_and_byte_size() {
let k = key(7, 3);
let data = fsqlite_types::PageData::zeroed(fsqlite_types::PageSize::default());
let page = CachedPage::new(k, data, 0xDEAD_BEEF, Some(42));
assert_eq!(page.key, k);
assert_eq!(page.xxh3, 0xDEAD_BEEF);
assert_eq!(page.wal_frame, Some(42));
assert_eq!(page.byte_size, 4096);
assert!(!page.is_pinned());
}
#[test]
fn arc_cache_new_creates_sharded_instance() {
let cache = ArcCache::new(256, 2 * 1024 * 1024);
assert!(cache.get(&key(1, 0)).is_none());
let snap = cache.resident_pages_snapshot();
assert_eq!(snap, Some(0));
}
#[test]
fn arc_cache_inner_reset_metrics_clears_counters() {
let mut inner = ArcCacheInner::new(64, 512 * 1024);
let snap_before = inner.metrics_snapshot();
assert_eq!(snap_before.hits, 0);
assert_eq!(snap_before.misses, 0);
inner.reset_metrics();
let snap_after = inner.metrics_snapshot();
assert_eq!(snap_after.hits, 0);
assert_eq!(snap_after.admits, 0);
}
}