use std::cell::Cell;
use fsqlite_error::{FrankenError, Result};
use fsqlite_types::cx::Cx;
use fsqlite_types::sync_primitives::Mutex;
use fsqlite_types::{PageNumber, PageSize};
use fsqlite_vfs::VfsFile;
use crate::page_buf::{PageBuf, PageBufPool};
#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
pub struct PageCacheMetricsSnapshot {
pub hits: u64,
pub misses: u64,
pub admits: u64,
pub evictions: u64,
pub cached_pages: usize,
pub pool_capacity: usize,
pub dirty_ratio_pct: u64,
pub t1_size: usize,
pub t2_size: usize,
pub b1_size: usize,
pub b2_size: usize,
pub p_target: usize,
pub mvcc_multi_version_pages: usize,
}
impl PageCacheMetricsSnapshot {
#[must_use]
pub fn total_accesses(self) -> u64 {
self.hits.saturating_add(self.misses)
}
#[must_use]
pub fn hit_rate_percent(self) -> f64 {
let total = self.total_accesses();
if total == 0 {
0.0
} else {
(self.hits as f64 * 100.0) / total as f64
}
}
}
pub struct PageCache {
pool: PageBufPool,
pages: std::collections::HashMap<PageNumber, PageBuf, foldhash::fast::FixedState>,
page_size: PageSize,
hits: Cell<u64>,
misses: Cell<u64>,
admits: Cell<u64>,
evictions: Cell<u64>,
}
impl PageCache {
pub fn new(page_size: PageSize) -> Self {
Self::with_pool(PageBufPool::new(page_size, 65_536), page_size)
}
pub fn with_pool(pool: PageBufPool, page_size: PageSize) -> Self {
Self {
pool,
pages: std::collections::HashMap::with_hasher(foldhash::fast::FixedState::default()),
page_size,
hits: Cell::new(0),
misses: Cell::new(0),
admits: Cell::new(0),
evictions: Cell::new(0),
}
}
pub fn pool(&self) -> &PageBufPool {
&self.pool
}
pub fn len(&self) -> usize {
self.pages.len()
}
pub fn is_empty(&self) -> bool {
self.pages.is_empty()
}
pub fn get(&self, page_no: PageNumber) -> Option<&[u8]> {
if let Some(page) = self.pages.get(&page_no) {
self.hits.set(self.hits.get().saturating_add(1));
Some(page.as_slice())
} else {
self.misses.set(self.misses.get().saturating_add(1));
None
}
}
#[inline]
pub fn get_mut(&mut self, page_no: PageNumber) -> Option<&mut [u8]> {
if let Some(page) = self.pages.get_mut(&page_no) {
self.hits.set(self.hits.get().saturating_add(1));
Some(page.as_mut_slice())
} else {
self.misses.set(self.misses.get().saturating_add(1));
None
}
}
#[inline]
#[must_use]
pub fn contains(&self, page_no: PageNumber) -> bool {
self.pages.contains_key(&page_no)
}
pub fn read_page(
&mut self,
cx: &Cx,
file: &mut impl VfsFile,
page_no: PageNumber,
) -> Result<&[u8]> {
if !self.contains(page_no) {
let mut buf = self.pool.acquire()?;
let offset = page_offset(page_no, self.page_size);
let bytes_read = file.read(cx, buf.as_mut_slice(), offset)?;
if bytes_read < self.page_size.as_usize() {
return Err(fsqlite_error::FrankenError::DatabaseCorrupt {
detail: format!(
"short read fetching page {page}: got {bytes_read} of {page_size}",
page = page_no.get(),
page_size = self.page_size.as_usize()
),
});
}
self.pages.insert(page_no, buf);
self.admits.set(self.admits.get().saturating_add(1));
}
Ok(self.pages.get(&page_no).expect("just inserted").as_slice())
}
pub fn write_page(&self, cx: &Cx, file: &mut impl VfsFile, page_no: PageNumber) -> Result<()> {
let Some(buf) = self.pages.get(&page_no) else {
self.misses.set(self.misses.get().saturating_add(1));
return Err(fsqlite_error::FrankenError::internal(format!(
"page {} not in cache",
page_no
)));
};
self.hits.set(self.hits.get().saturating_add(1));
let offset = page_offset(page_no, self.page_size);
file.write(cx, buf.as_slice(), offset)?;
Ok(())
}
pub fn insert_fresh(&mut self, page_no: PageNumber) -> Result<&mut [u8]> {
let mut buf = self.pool.acquire()?;
buf.as_mut_slice().fill(0);
let (out, admitted_new) = match self.pages.entry(page_no) {
std::collections::hash_map::Entry::Occupied(mut entry) => {
entry.insert(buf);
(entry.into_mut().as_mut_slice(), false)
}
std::collections::hash_map::Entry::Vacant(entry) => {
(entry.insert(buf).as_mut_slice(), true)
}
};
if admitted_new {
self.admits.set(self.admits.get().saturating_add(1));
}
Ok(out)
}
pub fn insert_buffer(&mut self, page_no: PageNumber, buf: PageBuf) {
let admitted_new = match self.pages.entry(page_no) {
std::collections::hash_map::Entry::Occupied(mut entry) => {
entry.insert(buf);
false
}
std::collections::hash_map::Entry::Vacant(entry) => {
entry.insert(buf);
true
}
};
if admitted_new {
self.admits.set(self.admits.get().saturating_add(1));
}
}
pub fn evict(&mut self, page_no: PageNumber) -> bool {
let removed = self.pages.remove(&page_no).is_some();
if removed {
self.evictions.set(self.evictions.get().saturating_add(1));
}
removed
}
pub fn evict_any(&mut self) -> bool {
let key = self.pages.keys().next().copied();
if let Some(key) = key {
self.pages.remove(&key);
self.evictions.set(self.evictions.get().saturating_add(1));
true
} else {
false
}
}
pub fn clear(&mut self) {
let removed = self.pages.len();
let removed_u64 = u64::try_from(removed).unwrap_or(u64::MAX);
self.evictions
.set(self.evictions.get().saturating_add(removed_u64));
self.pages.clear();
}
#[must_use]
pub fn metrics_snapshot(&self) -> PageCacheMetricsSnapshot {
let cached_pages = self.pages.len();
PageCacheMetricsSnapshot {
hits: self.hits.get(),
misses: self.misses.get(),
admits: self.admits.get(),
evictions: self.evictions.get(),
cached_pages,
pool_capacity: self.pool.capacity(),
dirty_ratio_pct: 0,
t1_size: cached_pages,
t2_size: 0,
b1_size: 0,
b2_size: 0,
p_target: cached_pages,
mvcc_multi_version_pages: 0,
}
}
pub fn reset_metrics(&mut self) {
self.hits.set(0);
self.misses.set(0);
self.admits.set(0);
self.evictions.set(0);
}
}
impl std::fmt::Debug for PageCache {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PageCache")
.field("page_size", &self.page_size)
.field("cached_pages", &self.pages.len())
.field("pool", &self.pool)
.field("hits", &self.hits)
.field("misses", &self.misses)
.field("admits", &self.admits)
.field("evictions", &self.evictions)
.field("metrics", &self.metrics_snapshot())
.finish()
}
}
const SHARD_COUNT: usize = 128;
const SHARD_MASK: usize = SHARD_COUNT - 1;
const GOLDEN_RATIO_32: u32 = 2_654_435_769;
#[repr(align(64))]
struct PageCacheShard {
pages: std::collections::HashMap<PageNumber, PageBuf, foldhash::fast::FixedState>,
hits: u64,
misses: u64,
admits: u64,
evictions: u64,
}
impl PageCacheShard {
fn new() -> Self {
Self {
pages: std::collections::HashMap::with_hasher(foldhash::fast::FixedState::default()),
hits: 0,
misses: 0,
admits: 0,
evictions: 0,
}
}
#[inline]
fn len(&self) -> usize {
self.pages.len()
}
#[inline]
fn contains(&self, page_no: PageNumber) -> bool {
self.pages.contains_key(&page_no)
}
#[inline]
fn get(&mut self, page_no: PageNumber) -> Option<&[u8]> {
if let Some(page) = self.pages.get(&page_no) {
self.hits = self.hits.saturating_add(1);
Some(page.as_slice())
} else {
self.misses = self.misses.saturating_add(1);
None
}
}
#[inline]
fn get_mut(&mut self, page_no: PageNumber) -> Option<&mut [u8]> {
if let Some(page) = self.pages.get_mut(&page_no) {
self.hits = self.hits.saturating_add(1);
Some(page.as_mut_slice())
} else {
self.misses = self.misses.saturating_add(1);
None
}
}
fn insert(&mut self, page_no: PageNumber, buf: PageBuf) -> bool {
let admitted_new = match self.pages.entry(page_no) {
std::collections::hash_map::Entry::Occupied(mut entry) => {
entry.insert(buf);
false
}
std::collections::hash_map::Entry::Vacant(entry) => {
entry.insert(buf);
true
}
};
if admitted_new {
self.admits = self.admits.saturating_add(1);
}
admitted_new
}
fn remove(&mut self, page_no: PageNumber) -> bool {
let removed = self.pages.remove(&page_no).is_some();
if removed {
self.evictions = self.evictions.saturating_add(1);
}
removed
}
fn remove_any(&mut self) -> Option<PageNumber> {
let key = self.pages.keys().next().copied();
if let Some(k) = key {
self.pages.remove(&k);
self.evictions = self.evictions.saturating_add(1);
}
key
}
fn clear(&mut self) -> usize {
let removed = self.pages.len();
self.evictions = self.evictions.saturating_add(removed as u64);
self.pages.clear();
removed
}
fn reset_metrics(&mut self) {
self.hits = 0;
self.misses = 0;
self.admits = 0;
self.evictions = 0;
}
}
impl std::fmt::Debug for PageCacheShard {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("PageCacheShard")
.field("pages", &self.pages.len())
.field("hits", &self.hits)
.field("misses", &self.misses)
.field("admits", &self.admits)
.field("evictions", &self.evictions)
.finish()
}
}
pub struct ShardedPageCache {
shards: Box<[Mutex<PageCacheShard>; SHARD_COUNT]>,
pool: PageBufPool,
page_size: PageSize,
}
impl ShardedPageCache {
pub fn new(page_size: PageSize) -> Self {
Self::with_pool(PageBufPool::new(page_size, 65_536), page_size)
}
pub fn with_pool(pool: PageBufPool, page_size: PageSize) -> Self {
let shards: Box<[Mutex<PageCacheShard>; SHARD_COUNT]> =
Box::new(std::array::from_fn(|_| Mutex::new(PageCacheShard::new())));
Self {
shards,
pool,
page_size,
}
}
#[inline]
fn shard_index(page_no: PageNumber) -> usize {
let hash = page_no.get().wrapping_mul(GOLDEN_RATIO_32);
(hash >> 25) as usize
}
pub fn pool(&self) -> &PageBufPool {
&self.pool
}
pub fn len(&self) -> usize {
self.shards.iter().map(|s| s.lock().len()).sum()
}
pub fn is_empty(&self) -> bool {
self.shards.iter().all(|s| s.lock().pages.is_empty())
}
#[inline]
pub fn contains(&self, page_no: PageNumber) -> bool {
let idx = Self::shard_index(page_no);
self.shards[idx].lock().contains(page_no)
}
#[inline]
pub fn get(&self, page_no: PageNumber) -> Option<Vec<u8>> {
let idx = Self::shard_index(page_no);
let mut shard = self.shards[idx].lock();
shard.get(page_no).map(|slice| slice.to_vec())
}
#[inline]
pub fn with_page<R>(&self, page_no: PageNumber, f: impl FnOnce(&[u8]) -> R) -> Option<R> {
let idx = Self::shard_index(page_no);
let mut shard = self.shards[idx].lock();
shard.get(page_no).map(f)
}
#[inline]
pub fn with_page_mut<R>(
&self,
page_no: PageNumber,
f: impl FnOnce(&mut [u8]) -> R,
) -> Option<R> {
let idx = Self::shard_index(page_no);
let mut shard = self.shards[idx].lock();
shard.get_mut(page_no).map(f)
}
pub fn read_page<R>(
&self,
cx: &Cx,
file: &mut impl VfsFile,
page_no: PageNumber,
f: impl FnOnce(&[u8]) -> R,
) -> Result<R> {
let idx = Self::shard_index(page_no);
let mut shard = self.shards[idx].lock();
if shard.pages.contains_key(&page_no) {
shard.hits = shard.hits.saturating_add(1);
let data = shard.pages.get(&page_no).unwrap();
return Ok(f(data.as_slice()));
}
shard.misses = shard.misses.saturating_add(1);
let mut buf = self.pool.acquire()?;
let offset = page_offset(page_no, self.page_size);
let bytes_read = file.read(cx, buf.as_mut_slice(), offset)?;
if bytes_read < self.page_size.as_usize() {
return Err(FrankenError::DatabaseCorrupt {
detail: format!(
"short read fetching page {page}: got {bytes_read} of {page_size}",
page = page_no.get(),
page_size = self.page_size.as_usize()
),
});
}
let result = f(buf.as_slice());
shard.insert(page_no, buf);
Ok(result)
}
pub fn write_page(&self, cx: &Cx, file: &mut impl VfsFile, page_no: PageNumber) -> Result<()> {
let idx = Self::shard_index(page_no);
let mut shard = self.shards[idx].lock();
if !shard.pages.contains_key(&page_no) {
shard.misses = shard.misses.saturating_add(1);
return Err(FrankenError::internal(format!(
"page {} not in cache",
page_no
)));
}
shard.hits = shard.hits.saturating_add(1);
let buf = shard.pages.get(&page_no).unwrap();
let offset = page_offset(page_no, self.page_size);
file.write(cx, buf.as_slice(), offset)?;
Ok(())
}
pub fn insert_fresh<R>(
&self,
page_no: PageNumber,
f: impl FnOnce(&mut [u8]) -> R,
) -> Result<R> {
let idx = Self::shard_index(page_no);
let mut shard = self.shards[idx].lock();
let mut buf = self.pool.acquire()?;
buf.as_mut_slice().fill(0);
let result = f(buf.as_mut_slice());
shard.insert(page_no, buf);
Ok(result)
}
pub fn insert_buffer(&self, page_no: PageNumber, buf: PageBuf) {
let idx = Self::shard_index(page_no);
let mut shard = self.shards[idx].lock();
shard.insert(page_no, buf);
}
pub fn evict(&self, page_no: PageNumber) -> bool {
let idx = Self::shard_index(page_no);
let mut shard = self.shards[idx].lock();
shard.remove(page_no)
}
pub fn evict_any(&self) -> bool {
let start = (std::time::Instant::now().elapsed().as_nanos() as usize) & SHARD_MASK;
for i in 0..SHARD_COUNT {
let idx = (start + i) & SHARD_MASK;
let mut shard = self.shards[idx].lock();
if shard.remove_any().is_some() {
return true;
}
}
false
}
pub fn clear(&self) {
for shard in self.shards.iter() {
shard.lock().clear();
}
}
#[must_use]
pub fn metrics_snapshot(&self) -> PageCacheMetricsSnapshot {
let mut total_hits = 0_u64;
let mut total_misses = 0_u64;
let mut total_admits = 0_u64;
let mut total_evictions = 0_u64;
let mut total_pages = 0_usize;
for shard in self.shards.iter() {
let s = shard.lock();
total_hits = total_hits.saturating_add(s.hits);
total_misses = total_misses.saturating_add(s.misses);
total_admits = total_admits.saturating_add(s.admits);
total_evictions = total_evictions.saturating_add(s.evictions);
total_pages += s.len();
}
PageCacheMetricsSnapshot {
hits: total_hits,
misses: total_misses,
admits: total_admits,
evictions: total_evictions,
cached_pages: total_pages,
pool_capacity: self.pool.capacity(),
dirty_ratio_pct: 0,
t1_size: total_pages,
t2_size: 0,
b1_size: 0,
b2_size: 0,
p_target: total_pages,
mvcc_multi_version_pages: 0,
}
}
pub fn reset_metrics(&self) {
for shard in self.shards.iter() {
shard.lock().reset_metrics();
}
}
#[must_use]
pub fn page_size(&self) -> PageSize {
self.page_size
}
#[must_use]
pub fn shard_distribution(&self) -> Vec<usize> {
self.shards.iter().map(|s| s.lock().len()).collect()
}
pub fn read_page_copy(
&self,
cx: &Cx,
file: &mut impl VfsFile,
page_no: PageNumber,
) -> Result<Vec<u8>> {
self.read_page(cx, file, page_no, |data| data.to_vec())
}
#[inline]
pub fn get_copy(&self, page_no: PageNumber) -> Option<Vec<u8>> {
let idx = Self::shard_index(page_no);
let mut shard = self.shards[idx].lock();
shard.get(page_no).map(|data| data.to_vec())
}
}
impl std::fmt::Debug for ShardedPageCache {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
let metrics = self.metrics_snapshot();
f.debug_struct("ShardedPageCache")
.field("shard_count", &SHARD_COUNT)
.field("page_size", &self.page_size)
.field("cached_pages", &metrics.cached_pages)
.field("hits", &metrics.hits)
.field("misses", &metrics.misses)
.field("admits", &metrics.admits)
.field("evictions", &metrics.evictions)
.finish_non_exhaustive()
}
}
#[inline]
fn page_offset(page_no: PageNumber, page_size: PageSize) -> u64 {
u64::from(page_no.get() - 1) * u64::from(page_size.get())
}
pub fn read_db_header(cx: &Cx, file: &mut impl VfsFile) -> Result<[u8; 100]> {
let mut header = [0u8; 100];
let bytes_read = file.read(cx, &mut header, 0)?;
if bytes_read < 100 {
return Err(FrankenError::DatabaseCorrupt {
detail: format!("database header short read: expected 100 bytes, got {bytes_read}"),
});
}
Ok(header)
}
#[cfg(test)]
#[allow(clippy::cast_possible_truncation)]
mod tests {
use super::*;
use fsqlite_types::flags::VfsOpenFlags;
use fsqlite_vfs::{MemoryVfs, Vfs};
use std::path::Path;
const BEAD_ID: &str = "bd-22n.2";
fn setup() -> (Cx, impl VfsFile) {
let cx = Cx::new();
let vfs = MemoryVfs::new();
let flags = VfsOpenFlags::MAIN_DB | VfsOpenFlags::CREATE | VfsOpenFlags::READWRITE;
let (file, _) = vfs.open(&cx, Some(Path::new("test.db")), flags).unwrap();
(cx, file)
}
#[cfg(unix)]
#[test]
fn test_spawn_blocking_io_read_page() {
use asupersync::runtime::{RuntimeBuilder, spawn_blocking_io};
use std::io::{ErrorKind, Write as _};
use std::os::unix::fs::FileExt as _;
use std::sync::Arc;
use tempfile::NamedTempFile;
fn read_exact_at(file: &std::fs::File, buf: &mut [u8], offset: u64) -> std::io::Result<()> {
let mut total = 0_usize;
while total < buf.len() {
#[allow(clippy::cast_possible_truncation)]
let off = offset + total as u64;
let n = file.read_at(&mut buf[total..], off)?;
if n == 0 {
return Err(std::io::Error::new(ErrorKind::UnexpectedEof, "short read"));
}
total += n;
}
Ok(())
}
let mut tmp = NamedTempFile::new().unwrap();
let page_data: Vec<u8> = (0..4096u16)
.map(|i| u8::try_from(i % 256).expect("i % 256 fits in u8"))
.collect();
tmp.as_file_mut().write_all(&page_data).unwrap();
tmp.as_file_mut().flush().unwrap();
let file = Arc::new(tmp.reopen().unwrap());
let pool = PageBufPool::new(PageSize::DEFAULT, 1);
let rt = RuntimeBuilder::low_latency()
.worker_threads(1)
.blocking_threads(1, 1)
.build()
.unwrap();
let join = rt.handle().spawn(async move {
let worker_tid = std::thread::current().id();
let mut buf = pool.acquire().unwrap();
let file2 = Arc::clone(&file);
let (buf, io_tid) = spawn_blocking_io(move || {
let io_tid = std::thread::current().id();
read_exact_at(file2.as_ref(), buf.as_mut_slice(), 0)?;
Ok::<_, std::io::Error>((buf, io_tid))
})
.await
.unwrap();
assert_ne!(
io_tid, worker_tid,
"spawn_blocking_io must dispatch work to a blocking thread"
);
assert_eq!(
buf.as_slice(),
page_data.as_slice(),
"bead_id={BEAD_ID} case=spawn_blocking_io_read_page data mismatch"
);
drop(buf);
assert_eq!(
pool.available(),
1,
"bead_id={BEAD_ID} case=spawn_blocking_io_read_page buf must return to pool"
);
});
rt.block_on(join);
}
#[test]
fn test_spawn_blocking_io_no_unsafe() {
let manifest = include_str!("../../../Cargo.toml");
assert!(
manifest.contains(r#"unsafe_code = "forbid""#),
"workspace must keep unsafe_code=forbid for IO dispatch paths"
);
}
#[test]
fn test_blocking_pool_lab_mode_inline() {
use asupersync::lab::{LabConfig, LabRuntime};
use asupersync::runtime::spawn_blocking_io;
use asupersync::types::Budget;
use std::sync::Arc;
use std::sync::atomic::{AtomicBool, Ordering};
let mut rt = LabRuntime::new(LabConfig::new(42));
let region = rt.state.create_root_region(Budget::INFINITE);
let ok = Arc::new(AtomicBool::new(false));
let ok_task = Arc::clone(&ok);
let (task_id, _handle) = rt
.state
.create_task(region, Budget::INFINITE, async move {
let worker_tid = std::thread::current().id();
let io_tid =
spawn_blocking_io(|| Ok::<_, std::io::Error>(std::thread::current().id()))
.await
.unwrap();
ok_task.store(worker_tid == io_tid, Ordering::Release);
})
.unwrap();
rt.scheduler.lock().schedule(task_id, 0);
rt.run_until_quiescent();
assert!(
ok.load(Ordering::Acquire),
"spawn_blocking_io must execute inline when no blocking pool exists (lab determinism)"
);
}
#[test]
fn test_cancel_mid_io_returns_buf_to_pool() {
use asupersync::runtime::{RuntimeBuilder, spawn_blocking_io, yield_now};
use std::future::poll_fn;
use std::task::Poll;
use std::time::Duration;
let rt = RuntimeBuilder::low_latency()
.worker_threads(1)
.blocking_threads(1, 1)
.build()
.unwrap();
let pool = PageBufPool::new(PageSize::DEFAULT, 1);
let join = rt.handle().spawn(async move {
let buf = pool.acquire().unwrap();
let mut fut = Box::pin(spawn_blocking_io(move || {
std::thread::sleep(Duration::from_millis(20));
Ok::<_, std::io::Error>(buf)
}));
let mut polled = false;
poll_fn(|cx| {
if !polled {
polled = true;
let _ = fut.as_mut().poll(cx);
}
Poll::Ready(())
})
.await;
drop(fut);
for _ in 0..200u32 {
if pool.available() == 1 {
break;
}
std::thread::sleep(Duration::from_millis(1));
yield_now().await;
}
assert_eq!(
pool.available(),
1,
"bead_id={BEAD_ID} case=cancel_mid_io_returns_buf_to_pool"
);
});
rt.block_on(join);
}
#[test]
fn test_pager_reads_pages_via_pool() {
let (cx, mut file) = setup();
let page_data = vec![0xAB_u8; 4096];
file.write(&cx, &page_data, 0).unwrap();
let pool = PageBufPool::new(PageSize::DEFAULT, 4);
let mut cache = PageCache::with_pool(pool.clone(), PageSize::DEFAULT);
let read = cache.read_page(&cx, &mut file, PageNumber::ONE).unwrap();
assert_eq!(read, page_data.as_slice());
assert_eq!(pool.available(), 0, "cached page still holds the buffer");
assert!(cache.evict(PageNumber::ONE));
assert_eq!(
pool.available(),
1,
"evicting a cached page should return its buffer to the pool"
);
}
#[test]
fn test_vfs_read_no_intermediate_alloc() {
let (cx, mut file) = setup();
let pattern: Vec<u8> = (0..4096u16)
.map(|i| u8::try_from(i % 256).expect("i % 256 fits in u8"))
.collect();
file.write(&cx, &pattern, 0).unwrap();
let pool = PageBufPool::new(PageSize::DEFAULT, 4);
let mut buf = pool.acquire().unwrap();
let ptr_before = buf.as_ptr();
file.read(&cx, buf.as_mut_slice(), 0).unwrap();
let ptr_after = buf.as_ptr();
assert_eq!(
ptr_before, ptr_after,
"bead_id={BEAD_ID} case=vfs_read_no_intermediate_alloc \
pointer must not change — read goes directly into PageBuf"
);
assert_eq!(
buf.as_slice(),
pattern.as_slice(),
"bead_id={BEAD_ID} case=vfs_read_data_correct"
);
}
#[test]
fn test_vfs_write_no_intermediate_alloc() {
let (cx, mut file) = setup();
let pool = PageBufPool::new(PageSize::DEFAULT, 4);
let mut buf = pool.acquire().unwrap();
for (i, b) in buf.as_mut_slice().iter_mut().enumerate() {
*b = u8::try_from(i % 251).expect("i % 251 fits in u8"); }
let ptr_before = buf.as_ptr();
file.write(&cx, buf.as_slice(), 0).unwrap();
let ptr_after = buf.as_ptr();
assert_eq!(
ptr_before, ptr_after,
"bead_id={BEAD_ID} case=vfs_write_no_intermediate_alloc \
PageBuf pointer must be stable through write"
);
let mut verify = vec![0u8; 4096];
file.read(&cx, &mut verify, 0).unwrap();
assert_eq!(
verify.as_slice(),
buf.as_slice(),
"bead_id={BEAD_ID} case=vfs_write_data_roundtrip"
);
}
#[test]
fn test_pager_returns_ref_not_copy() {
let (cx, mut file) = setup();
let data = vec![0xAB_u8; 4096];
file.write(&cx, &data, 0).unwrap();
let mut cache = PageCache::new(PageSize::DEFAULT);
let page1 = PageNumber::ONE;
let ref1 = cache.read_page(&cx, &mut file, page1).unwrap();
let ref1_ptr = ref1.as_ptr();
assert_eq!(
&ref1[..4096],
data.as_slice(),
"bead_id={BEAD_ID} case=pager_ref_data_correct"
);
let ref2 = cache.get(page1).unwrap();
let ref2_ptr = ref2.as_ptr();
assert_eq!(
ref1_ptr, ref2_ptr,
"bead_id={BEAD_ID} case=pager_returns_ref_not_copy \
get() must return reference to same memory as read_page()"
);
}
#[test]
fn test_wal_uses_buffered_io_compat() {
let wal_header_size: u64 = 24;
for &size in &[512u32, 1024, 2048, 4096, 8192, 16384, 32768, 65536] {
let frame_size = wal_header_size + u64::from(size);
let wal_header_bytes: u64 = 32; let frame2_offset = wal_header_bytes + frame_size;
let _sector_4k_aligned = frame2_offset % 4096 == 0;
let sector_512_aligned = frame2_offset % 512 == 0;
assert!(
!sector_512_aligned,
"bead_id={BEAD_ID} case=wal_frame_not_512_aligned \
WAL frame 2 at offset {frame2_offset} should NOT be 512-byte aligned \
for page_size={size}"
);
}
}
#[test]
fn test_small_header_stack_buffer_ok() {
let (cx, mut file) = setup();
let mut header_data = [0u8; 100];
header_data[..16].copy_from_slice(b"SQLite format 3\0");
header_data[16..18].copy_from_slice(&4096u16.to_be_bytes()); file.write(&cx, &header_data, 0).unwrap();
let header = read_db_header(&cx, &mut file).unwrap();
assert_eq!(
&header[..16],
b"SQLite format 3\0",
"bead_id={BEAD_ID} case=small_header_stack_buffer_ok"
);
let page_size = u16::from_be_bytes([header[16], header[17]]);
assert_eq!(
page_size, 4096,
"bead_id={BEAD_ID} case=header_page_size_correct"
);
}
#[test]
fn test_page_decode_bounds_checked() {
let (cx, mut file) = setup();
let mut page_data = vec![0u8; 4096];
page_data[0] = 0x0D; page_data[3..5].copy_from_slice(&10u16.to_be_bytes()); page_data[5..7].copy_from_slice(&100u16.to_be_bytes()); file.write(&cx, &page_data, 0).unwrap();
let mut cache = PageCache::new(PageSize::DEFAULT);
let page = cache.read_page(&cx, &mut file, PageNumber::ONE).unwrap();
let page_type = page[0];
assert_eq!(page_type, 0x0D, "bead_id={BEAD_ID} case=page_decode_type");
let cell_count = u16::from_be_bytes([page[3], page[4]]);
assert_eq!(
cell_count, 10,
"bead_id={BEAD_ID} case=page_decode_cell_count"
);
let content_offset = u16::from_be_bytes([page[5], page[6]]);
assert_eq!(
content_offset, 100,
"bead_id={BEAD_ID} case=page_decode_content_offset"
);
assert_eq!(
page.len(),
4096,
"bead_id={BEAD_ID} case=page_decode_bounds_checked"
);
}
#[test]
fn test_cache_insert_fresh_zeroed() {
let mut cache = PageCache::new(PageSize::DEFAULT);
let page1 = PageNumber::ONE;
let data = cache.insert_fresh(page1).unwrap();
assert!(
data.iter().all(|&b| b == 0),
"bead_id={BEAD_ID} case=insert_fresh_zeroed"
);
assert_eq!(data.len(), 4096);
assert_eq!(cache.len(), 1);
}
#[test]
fn test_cache_get_mut_modifies_in_place() {
let mut cache = PageCache::new(PageSize::DEFAULT);
let page1 = PageNumber::ONE;
cache.insert_fresh(page1).unwrap();
let data = cache.get_mut(page1).unwrap();
data[0] = 0xFF;
data[4095] = 0xEE;
let read_back = cache.get(page1).unwrap();
assert_eq!(read_back[0], 0xFF);
assert_eq!(read_back[4095], 0xEE);
}
#[test]
fn test_cache_evict_returns_to_pool() {
let mut cache = PageCache::new(PageSize::DEFAULT);
let page1 = PageNumber::ONE;
assert_eq!(cache.pool().available(), 0);
cache.insert_fresh(page1).unwrap();
assert_eq!(cache.pool().available(), 0);
assert!(cache.evict(page1));
assert!(!cache.contains(page1));
assert_eq!(
cache.pool().available(),
1,
"bead_id={BEAD_ID} case=evict_returns_to_pool"
);
}
#[test]
fn test_cache_evict_nonexistent() {
let mut cache = PageCache::new(PageSize::DEFAULT);
let page1 = PageNumber::ONE;
assert!(!cache.evict(page1));
}
#[test]
fn test_cache_clear_returns_all_to_pool() {
let mut cache = PageCache::new(PageSize::DEFAULT);
for i in 1..=5u32 {
let pn = PageNumber::new(i).unwrap();
cache.insert_fresh(pn).unwrap();
}
assert_eq!(cache.len(), 5);
assert_eq!(cache.pool().available(), 0);
cache.clear();
assert_eq!(cache.len(), 0);
assert_eq!(
cache.pool().available(),
5,
"bead_id={BEAD_ID} case=clear_returns_all_to_pool"
);
}
#[test]
fn test_cache_multiple_pages() {
let (cx, mut file) = setup();
for i in 0..3u32 {
let seed = u8::try_from(i).expect("i <= 2");
let data = vec![(seed + 1) * 0x11; 4096];
let offset = u64::from(i) * 4096;
file.write(&cx, &data, offset).unwrap();
}
let mut cache = PageCache::new(PageSize::DEFAULT);
for i in 1..=3u32 {
let pn = PageNumber::new(i).unwrap();
let page = cache.read_page(&cx, &mut file, pn).unwrap();
let expected = u8::try_from(i).expect("i <= 3") * 0x11;
assert!(
page.iter().all(|&b| b == expected),
"bead_id={BEAD_ID} case=multiple_pages page={i} expected={expected:#x}"
);
}
assert_eq!(cache.len(), 3);
}
#[test]
fn test_cache_write_page_roundtrip() {
let (cx, mut file) = setup();
let mut cache = PageCache::new(PageSize::DEFAULT);
let page1 = PageNumber::ONE;
let data = cache.insert_fresh(page1).unwrap();
data.fill(0xCD);
cache.write_page(&cx, &mut file, page1).unwrap();
let mut verify = vec![0u8; 4096];
file.read(&cx, &mut verify, 0).unwrap();
assert!(
verify.iter().all(|&b| b == 0xCD),
"bead_id={BEAD_ID} case=write_page_roundtrip"
);
}
#[test]
fn test_page_offset_calculation() {
assert_eq!(
page_offset(PageNumber::ONE, PageSize::DEFAULT),
0,
"bead_id={BEAD_ID} case=page_offset_page1"
);
let p2 = PageNumber::new(2).unwrap();
assert_eq!(
page_offset(p2, PageSize::DEFAULT),
4096,
"bead_id={BEAD_ID} case=page_offset_page2"
);
let p100 = PageNumber::new(100).unwrap();
let ps512 = PageSize::new(512).unwrap();
assert_eq!(
page_offset(p100, ps512),
50688,
"bead_id={BEAD_ID} case=page_offset_page100_512"
);
}
#[test]
fn test_e2e_zero_copy_io_no_allocations() {
let (cx, mut file) = setup();
let num_pages: u32 = 10;
for i in 0..num_pages {
let byte = u8::try_from(i).expect("i <= 9").wrapping_add(0x10);
let data = vec![byte; 4096];
file.write(&cx, &data, u64::from(i) * 4096).unwrap();
}
let mut cache = PageCache::new(PageSize::DEFAULT);
let mut ptrs: Vec<usize> = Vec::with_capacity(num_pages as usize);
for i in 1..=num_pages {
let pn = PageNumber::new(i).unwrap();
let page = cache.read_page(&cx, &mut file, pn).unwrap();
ptrs.push(page.as_ptr() as usize);
}
for round in 0..5u32 {
for i in 1..=num_pages {
let pn = PageNumber::new(i).unwrap();
let page = cache.get(pn).unwrap();
let ptr = page.as_ptr() as usize;
assert_eq!(
ptr,
ptrs[(i - 1) as usize],
"bead_id={BEAD_ID} case=e2e_pointer_stable \
round={round} page={i}"
);
let expected = u8::try_from(i - 1).expect("i - 1 <= 9").wrapping_add(0x10);
assert!(
page.iter().all(|&b| b == expected),
"bead_id={BEAD_ID} case=e2e_data_correct \
round={round} page={i}"
);
}
}
let pool_available_before = cache.pool().available();
let old_ptr = ptrs[0];
cache.evict(PageNumber::ONE);
assert_eq!(
cache.pool().available(),
pool_available_before + 1,
"bead_id={BEAD_ID} case=e2e_evict_returns_to_pool"
);
let page1_reread = cache.read_page(&cx, &mut file, PageNumber::ONE).unwrap();
let new_ptr = page1_reread.as_ptr() as usize;
assert_eq!(
new_ptr, old_ptr,
"bead_id={BEAD_ID} case=e2e_pool_reuse_after_evict \
Expected recycled buffer at {old_ptr:#x}, got {new_ptr:#x}"
);
eprintln!("pages_cached={}", cache.len());
eprintln!("pool_available={}", cache.pool().available());
eprintln!("pointer_checks_passed={}", num_pages * 5 + 1);
}
#[test]
fn test_page_cache_debug() {
let cache = PageCache::new(PageSize::DEFAULT);
let debug = format!("{cache:?}");
assert!(
debug.contains("PageCache"),
"bead_id={BEAD_ID} case=debug_format"
);
}
#[test]
fn test_metrics_snapshot_and_reset() {
let mut cache = PageCache::new(PageSize::DEFAULT);
let page1 = PageNumber::ONE;
assert!(cache.get(page1).is_none());
let fresh = cache.insert_fresh(page1).unwrap();
fresh[0] = 7;
assert!(cache.get(page1).is_some());
assert!(cache.evict(page1));
let snapshot = cache.metrics_snapshot();
assert_eq!(snapshot.hits, 1, "bead_id={BEAD_ID} case=metrics_hits");
assert_eq!(snapshot.misses, 1, "bead_id={BEAD_ID} case=metrics_misses");
assert_eq!(snapshot.admits, 1, "bead_id={BEAD_ID} case=metrics_admits");
assert_eq!(
snapshot.evictions, 1,
"bead_id={BEAD_ID} case=metrics_evictions"
);
assert_eq!(
snapshot.total_accesses(),
2,
"bead_id={BEAD_ID} case=metrics_total_accesses"
);
assert!(
(snapshot.hit_rate_percent() - 50.0).abs() < f64::EPSILON,
"bead_id={BEAD_ID} case=metrics_hit_rate"
);
cache.reset_metrics();
let reset = cache.metrics_snapshot();
assert_eq!(reset.hits, 0, "bead_id={BEAD_ID} case=reset_hits");
assert_eq!(reset.misses, 0, "bead_id={BEAD_ID} case=reset_misses");
assert_eq!(reset.admits, 0, "bead_id={BEAD_ID} case=reset_admits");
assert_eq!(reset.evictions, 0, "bead_id={BEAD_ID} case=reset_evictions");
}
const BEAD_22N8: &str = "bd-22n.8";
#[test]
fn test_cache_lookup_no_alloc() {
let (cx, mut file) = setup();
let data = vec![0xBE_u8; 4096];
file.write(&cx, &data, 0).unwrap();
let mut cache = PageCache::new(PageSize::DEFAULT);
let page1 = PageNumber::ONE;
let initial = cache.read_page(&cx, &mut file, page1).unwrap();
let initial_ptr = initial.as_ptr();
for round in 0..100u32 {
let cached = cache.get(page1).unwrap();
assert_eq!(
cached.as_ptr(),
initial_ptr,
"bead_id={BEAD_22N8} case=cache_lookup_no_alloc \
round={round} pointer must be stable (no realloc)"
);
}
}
#[test]
fn test_cache_lookup_hit_returns_reference() {
let mut cache = PageCache::new(PageSize::DEFAULT);
let page1 = PageNumber::ONE;
let fresh = cache.insert_fresh(page1).unwrap();
fresh[0] = 0xAA;
let ptr_after_insert = cache.get(page1).unwrap().as_ptr();
let mutref = cache.get_mut(page1).unwrap();
mutref[1] = 0xBB;
let read_back = cache.get(page1).unwrap();
assert_eq!(
read_back.as_ptr(),
ptr_after_insert,
"bead_id={BEAD_22N8} case=cache_lookup_returns_reference \
pointer must be stable through mutation"
);
assert_eq!(read_back[0], 0xAA);
assert_eq!(read_back[1], 0xBB);
}
#[test]
fn test_pool_reuse_avoids_alloc_on_reread() {
let (cx, mut file) = setup();
let data = vec![0xDD_u8; 4096];
file.write(&cx, &data, 0).unwrap();
let mut cache = PageCache::new(PageSize::DEFAULT);
let page1 = PageNumber::ONE;
let _ = cache.read_page(&cx, &mut file, page1).unwrap();
assert_eq!(cache.pool().available(), 0);
cache.evict(page1);
assert_eq!(
cache.pool().available(),
1,
"bead_id={BEAD_22N8} case=evicted_buffer_returned_to_pool"
);
let reread = cache.read_page(&cx, &mut file, page1).unwrap();
assert_eq!(
reread,
data.as_slice(),
"bead_id={BEAD_22N8} case=pool_reuse_data_correct"
);
assert_eq!(
cache.pool().available(),
0,
"bead_id={BEAD_22N8} case=pool_buffer_consumed_on_reread"
);
}
const BEAD_3WOP3_2: &str = "bd-3wop3.2";
#[test]
fn test_sharded_cache_basic_operations() {
let cache = ShardedPageCache::new(PageSize::DEFAULT);
let p1 = PageNumber::ONE;
let p2 = PageNumber::new(2).unwrap();
cache.insert_fresh(p1, |data| data[0] = 0xAA).unwrap();
cache.insert_fresh(p2, |data| data[0] = 0xBB).unwrap();
assert_eq!(cache.len(), 2);
assert!(cache.contains(p1));
assert!(cache.contains(p2));
cache.with_page(p1, |data| assert_eq!(data[0], 0xAA));
cache.with_page(p2, |data| assert_eq!(data[0], 0xBB));
assert!(cache.evict(p1));
assert!(!cache.contains(p1));
assert!(cache.contains(p2));
assert_eq!(cache.len(), 1);
let m = cache.metrics_snapshot();
assert_eq!(m.admits, 2, "bead_id={BEAD_3WOP3_2} case=basic_admits");
assert_eq!(
m.evictions, 1,
"bead_id={BEAD_3WOP3_2} case=basic_evictions"
);
}
#[test]
fn test_sharded_cache_shard_distribution() {
let cache = ShardedPageCache::new(PageSize::DEFAULT);
for i in 1..=256u32 {
let pn = PageNumber::new(i).unwrap();
cache.insert_fresh(pn, |_| {}).unwrap();
}
let dist = cache.shard_distribution();
assert_eq!(dist.len(), 128);
let non_empty = dist.iter().filter(|&&n| n > 0).count();
assert!(
non_empty >= 64,
"bead_id={BEAD_3WOP3_2} case=shard_distribution \
expected at least 64 non-empty shards, got {non_empty}"
);
let max_per_shard = *dist.iter().max().unwrap();
assert!(
max_per_shard <= 16,
"bead_id={BEAD_3WOP3_2} case=shard_balance \
expected max 16 pages per shard, got {max_per_shard}"
);
}
#[test]
fn test_sharded_cache_cross_shard_eviction() {
let cache = ShardedPageCache::new(PageSize::DEFAULT);
for i in 1..=16u32 {
let pn = PageNumber::new(i * 100).unwrap();
cache.insert_fresh(pn, |_| {}).unwrap();
}
assert_eq!(cache.len(), 16);
let mut evicted = 0;
while cache.evict_any() {
evicted += 1;
if evicted > 100 {
panic!("bead_id={BEAD_3WOP3_2} case=cross_shard_eviction infinite loop");
}
}
assert_eq!(
evicted, 16,
"bead_id={BEAD_3WOP3_2} case=cross_shard_eviction_count"
);
assert!(cache.is_empty());
}
#[test]
fn test_sharded_cache_clear() {
let cache = ShardedPageCache::new(PageSize::DEFAULT);
for i in 1..=100u32 {
let pn = PageNumber::new(i).unwrap();
cache.insert_fresh(pn, |_| {}).unwrap();
}
assert_eq!(cache.len(), 100);
cache.clear();
assert!(cache.is_empty());
assert_eq!(cache.len(), 0);
let m = cache.metrics_snapshot();
assert_eq!(
m.evictions, 100,
"bead_id={BEAD_3WOP3_2} case=clear_evictions"
);
}
#[test]
fn test_sharded_cache_metrics_aggregation() {
let cache = ShardedPageCache::new(PageSize::DEFAULT);
for i in 1..=50u32 {
let pn = PageNumber::new(i).unwrap();
cache.insert_fresh(pn, |_| {}).unwrap();
}
for i in 1..=50u32 {
let pn = PageNumber::new(i).unwrap();
cache.with_page(pn, |_| {});
}
for i in 51..=100u32 {
let pn = PageNumber::new(i).unwrap();
cache.with_page(pn, |_| {});
}
let m = cache.metrics_snapshot();
assert_eq!(m.admits, 50, "bead_id={BEAD_3WOP3_2} case=metrics_admits");
assert_eq!(m.hits, 50, "bead_id={BEAD_3WOP3_2} case=metrics_hits");
assert_eq!(m.misses, 50, "bead_id={BEAD_3WOP3_2} case=metrics_misses");
assert_eq!(
m.cached_pages, 50,
"bead_id={BEAD_3WOP3_2} case=metrics_cached_pages"
);
cache.reset_metrics();
let reset = cache.metrics_snapshot();
assert_eq!(reset.hits, 0, "bead_id={BEAD_3WOP3_2} case=reset_metrics");
assert_eq!(reset.misses, 0);
assert_eq!(reset.admits, 0);
assert_eq!(reset.cached_pages, 50);
}
#[test]
fn test_sharded_cache_shard_padding_alignment() {
let shard_size = std::mem::size_of::<PageCacheShard>();
assert!(
shard_size >= 64,
"bead_id={BEAD_3WOP3_2} case=shard_padding \
PageCacheShard size {shard_size} should be >= 64 bytes"
);
assert_eq!(
shard_size % 64,
0,
"bead_id={BEAD_3WOP3_2} case=shard_alignment \
PageCacheShard size {shard_size} must be multiple of 64"
);
let shard_align = std::mem::align_of::<PageCacheShard>();
assert_eq!(
shard_align, 64,
"bead_id={BEAD_3WOP3_2} case=shard_align_req \
PageCacheShard alignment should be 64, got {shard_align}"
);
}
#[test]
fn test_sharded_cache_with_page_mut() {
let cache = ShardedPageCache::new(PageSize::DEFAULT);
let p1 = PageNumber::ONE;
cache.insert_fresh(p1, |data| data.fill(0)).unwrap();
cache.with_page_mut(p1, |data| {
data[0] = 0x12;
data[1] = 0x34;
});
cache.with_page(p1, |data| {
assert_eq!(data[0], 0x12, "bead_id={BEAD_3WOP3_2} case=with_page_mut_0");
assert_eq!(data[1], 0x34, "bead_id={BEAD_3WOP3_2} case=with_page_mut_1");
});
}
#[test]
fn test_sharded_cache_insert_buffer() {
let cache = ShardedPageCache::new(PageSize::DEFAULT);
let p1 = PageNumber::ONE;
let mut buf = cache.pool().acquire().unwrap();
buf.as_mut_slice().fill(0xEE);
cache.insert_buffer(p1, buf);
assert!(cache.contains(p1));
cache.with_page(p1, |data| {
assert!(
data.iter().all(|&b| b == 0xEE),
"bead_id={BEAD_3WOP3_2} case=insert_buffer_data"
);
});
}
#[test]
fn test_sharded_cache_vfs_read_write() {
let (cx, mut file) = setup();
let test_data = vec![0xAB_u8; 4096];
file.write(&cx, &test_data, 0).unwrap();
let cache = ShardedPageCache::new(PageSize::DEFAULT);
let p1 = PageNumber::ONE;
let result = cache.read_page(&cx, &mut file, p1, |data| {
assert_eq!(
data,
test_data.as_slice(),
"bead_id={BEAD_3WOP3_2} case=vfs_read_data"
);
data[0]
});
assert_eq!(result.unwrap(), 0xAB);
cache.with_page_mut(p1, |data| data[0] = 0xCD);
cache.write_page(&cx, &mut file, p1).unwrap();
let mut verify = vec![0u8; 4096];
file.read(&cx, &mut verify, 0).unwrap();
assert_eq!(
verify[0], 0xCD,
"bead_id={BEAD_3WOP3_2} case=vfs_write_verify"
);
}
#[test]
fn test_sharded_cache_8_threads_no_deadlock() {
use std::sync::Arc;
use std::thread;
let cache = Arc::new(ShardedPageCache::new(PageSize::DEFAULT));
let num_threads = 8;
let ops_per_thread = 1000;
let handles: Vec<_> = (0..num_threads)
.map(|tid| {
let c = Arc::clone(&cache);
thread::spawn(move || {
for i in 0..ops_per_thread {
let base = tid * 10000 + i;
let pn = PageNumber::new(base as u32 + 1).unwrap();
c.insert_fresh(pn, |data| data[0] = (tid & 0xFF) as u8)
.unwrap();
c.with_page(pn, |data| {
assert_eq!(data[0], (tid & 0xFF) as u8);
});
if i % 10 == 0 {
c.evict(pn);
}
}
})
})
.collect();
for h in handles {
h.join()
.expect("bead_id={BEAD_3WOP3_2} case=8t_no_deadlock thread panic");
}
let m = cache.metrics_snapshot();
assert!(
m.admits >= (num_threads * ops_per_thread) as u64,
"bead_id={BEAD_3WOP3_2} case=8t_admits"
);
}
#[test]
fn test_sharded_cache_16_threads_no_deadlock() {
use std::sync::Arc;
use std::thread;
let cache = Arc::new(ShardedPageCache::new(PageSize::DEFAULT));
let num_threads = 16;
let ops_per_thread = 500;
let handles: Vec<_> = (0..num_threads)
.map(|tid| {
let c = Arc::clone(&cache);
thread::spawn(move || {
for i in 0..ops_per_thread {
let base = tid * 10000 + i;
let pn = PageNumber::new(base as u32 + 1).unwrap();
c.insert_fresh(pn, |data| data[0] = ((tid * 7) & 0xFF) as u8)
.unwrap();
c.with_page(pn, |data| {
assert_eq!(data[0], ((tid * 7) & 0xFF) as u8);
});
if i % 5 == 0 {
c.evict(pn);
}
}
})
})
.collect();
for h in handles {
h.join()
.expect("bead_id={BEAD_3WOP3_2} case=16t_no_deadlock thread panic");
}
let m = cache.metrics_snapshot();
assert!(
m.admits >= (num_threads * ops_per_thread) as u64,
"bead_id={BEAD_3WOP3_2} case=16t_admits"
);
}
#[test]
fn test_sharded_cache_throughput_vs_single() {
use std::time::Instant;
let iterations = 10_000;
let mut single = PageCache::new(PageSize::DEFAULT);
let start = Instant::now();
for i in 1..=iterations {
let pn = PageNumber::new(i).unwrap();
single.insert_fresh(pn).unwrap();
let _ = single.get(pn);
}
let single_elapsed = start.elapsed();
let sharded = ShardedPageCache::new(PageSize::DEFAULT);
let start = Instant::now();
for i in 1..=iterations {
let pn = PageNumber::new(i).unwrap();
sharded.insert_fresh(pn, |_| {}).unwrap();
sharded.with_page(pn, |_| {});
}
let sharded_elapsed = start.elapsed();
let ratio = sharded_elapsed.as_nanos() as f64 / single_elapsed.as_nanos() as f64;
assert!(
ratio < 3.0,
"bead_id={BEAD_3WOP3_2} case=throughput_overhead \
sharded cache is {ratio:.2}x slower than single (max 3x allowed)"
);
eprintln!(
"bead_id={BEAD_3WOP3_2} throughput_ratio={ratio:.2}x \
single={:?} sharded={:?}",
single_elapsed, sharded_elapsed
);
}
#[test]
fn test_sharded_cache_concurrent_same_shard() {
use std::sync::Arc;
use std::thread;
let cache = Arc::new(ShardedPageCache::new(PageSize::DEFAULT));
let num_threads = 4;
let ops_per_thread = 500;
let base_page = PageNumber::ONE;
let base_shard = ShardedPageCache::shard_index(base_page);
let mut same_shard_pages = vec![1u32];
for i in 2..10000u32 {
let pn = PageNumber::new(i).unwrap();
if ShardedPageCache::shard_index(pn) == base_shard {
same_shard_pages.push(i);
if same_shard_pages.len() >= (num_threads * ops_per_thread) {
break;
}
}
}
let pages = Arc::new(same_shard_pages);
let handles: Vec<_> = (0..num_threads)
.map(|tid| {
let c = Arc::clone(&cache);
let p = Arc::clone(&pages);
thread::spawn(move || {
let start = tid * ops_per_thread;
for i in 0..ops_per_thread {
let idx = start + i;
if idx >= p.len() {
break;
}
let pn = PageNumber::new(p[idx]).unwrap();
c.insert_fresh(pn, |data| data[0] = (tid & 0xFF) as u8)
.unwrap();
c.with_page(pn, |_| {});
}
})
})
.collect();
for h in handles {
h.join()
.expect("bead_id={BEAD_3WOP3_2} case=concurrent_same_shard panic");
}
let dist = cache.shard_distribution();
assert!(
dist[base_shard] > 0,
"bead_id={BEAD_3WOP3_2} case=same_shard_populated"
);
}
}