use std::collections::HashMap;
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex};
use fsqlite_error::{FrankenError, Result};
use fsqlite_types::cx::Cx;
use fsqlite_types::flags::{AccessFlags, SyncFlags, VfsOpenFlags};
use fsqlite_types::{
BTreePageHeader, CommitSeq, DATABASE_HEADER_SIZE, DatabaseHeader,
FRANKENSQLITE_SQLITE_VERSION_NUMBER, PageData, PageNumber, PageSize,
};
use fsqlite_vfs::{Vfs, VfsFile};
use crate::journal::{JournalHeader, JournalPageRecord};
use crate::page_buf::{PageBuf, PageBufPool};
use crate::page_cache::{PageCache, PageCacheMetricsSnapshot};
use crate::traits::{self, JournalMode, MvccPager, TransactionHandle, TransactionMode, WalBackend};
pub(crate) struct PagerInner<F: VfsFile> {
db_file: F,
cache: PageCache,
page_size: PageSize,
db_size: u32,
next_page: u32,
writer_active: bool,
active_transactions: u32,
checkpoint_active: bool,
freelist: Vec<PageNumber>,
journal_mode: JournalMode,
wal_backend: Option<Box<dyn WalBackend>>,
commit_seq: CommitSeq,
}
impl<F: VfsFile> PagerInner<F> {
fn read_page_copy(&mut self, cx: &Cx, page_no: PageNumber) -> Result<Vec<u8>> {
if self.journal_mode == JournalMode::Wal {
if let Some(ref mut wal) = self.wal_backend {
if let Some(wal_data) = wal.read_page(cx, page_no.get())? {
return Ok(wal_data);
}
}
}
if let Some(data) = self.cache.get(page_no) {
return Ok(data.to_vec());
}
if page_no.get() > self.db_size {
return Ok(vec![0_u8; self.page_size.as_usize()]);
}
let mut buf = match self.cache.pool().clone().acquire() {
Ok(buf) => buf,
Err(FrankenError::OutOfMemory) => {
if self.cache.evict_any() {
self.cache.pool().clone().acquire()?
} else {
return Err(FrankenError::OutOfMemory);
}
}
Err(err) => return Err(err),
};
let page_size = self.page_size.as_usize();
let offset = u64::from(page_no.get() - 1) * page_size as u64;
let bytes_read = self.db_file.read(cx, buf.as_mut_slice(), offset)?;
if bytes_read < page_size {
return Err(FrankenError::DatabaseCorrupt {
detail: format!(
"short read fetching page {page}: got {bytes_read} of {page_size}",
page = page_no.get()
),
});
}
let out = buf.as_slice()[..page_size].to_vec();
let fresh = match self.cache.insert_fresh(page_no) {
Ok(fresh) => fresh,
Err(FrankenError::OutOfMemory) => {
if self.cache.evict_any() {
self.cache.insert_fresh(page_no)?
} else {
return Ok(out);
}
}
Err(err) => return Err(err),
};
fresh.copy_from_slice(&out);
Ok(out)
}
fn flush_page(&mut self, cx: &Cx, page_no: PageNumber, data: &[u8]) -> Result<()> {
if let Some(cached) = self.cache.get_mut(page_no) {
let len = cached.len().min(data.len());
cached[..len].copy_from_slice(&data[..len]);
} else {
match self.cache.insert_fresh(page_no) {
Ok(fresh) => {
let len = fresh.len().min(data.len());
fresh[..len].copy_from_slice(&data[..len]);
}
Err(FrankenError::OutOfMemory) => {
if self.cache.evict_any() {
if let Ok(fresh) = self.cache.insert_fresh(page_no) {
let len = fresh.len().min(data.len());
fresh[..len].copy_from_slice(&data[..len]);
}
}
}
Err(err) => return Err(err),
}
}
let page_size = self.page_size.as_usize();
let offset = u64::from(page_no.get() - 1) * page_size as u64;
self.db_file.write(cx, data, offset)?;
Ok(())
}
}
pub struct SimplePager<V: Vfs> {
vfs: Arc<V>,
db_path: PathBuf,
inner: Arc<Mutex<PagerInner<V::File>>>,
}
impl<V: Vfs> traits::sealed::Sealed for SimplePager<V> {}
impl<V> MvccPager for SimplePager<V>
where
V: Vfs + Send + Sync,
V::File: Send + Sync,
{
type Txn = SimpleTransaction<V>;
fn begin(&self, cx: &Cx, mode: TransactionMode) -> Result<Self::Txn> {
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimplePager lock poisoned"))?;
if inner.checkpoint_active {
return Err(FrankenError::Busy);
}
if inner.journal_mode == JournalMode::Wal {
let wal = inner.wal_backend.as_mut().ok_or_else(|| {
FrankenError::internal("WAL mode active but no WAL backend installed")
})?;
wal.begin_transaction(cx)?;
}
let eager_writer = matches!(
mode,
TransactionMode::Immediate | TransactionMode::Exclusive
);
if eager_writer && inner.writer_active {
return Err(FrankenError::Busy);
}
if eager_writer {
inner.writer_active = true;
}
inner.active_transactions = inner.active_transactions.saturating_add(1);
let original_db_size = inner.db_size;
let journal_mode = inner.journal_mode;
let pool = inner.cache.pool().clone();
drop(inner);
Ok(SimpleTransaction {
vfs: Arc::clone(&self.vfs),
journal_path: Self::journal_path(&self.db_path),
inner: Arc::clone(&self.inner),
write_set: HashMap::new(),
freed_pages: Vec::new(),
allocated_from_freelist: Vec::new(),
mode,
is_writer: eager_writer,
committed: false,
finished: false,
original_db_size,
savepoint_stack: Vec::new(),
journal_mode,
pool,
})
}
fn journal_mode(&self) -> JournalMode {
self.inner
.lock()
.map(|inner| inner.journal_mode)
.unwrap_or_default()
}
fn set_journal_mode(&self, _cx: &Cx, mode: JournalMode) -> Result<JournalMode> {
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimplePager lock poisoned"))?;
if inner.checkpoint_active {
return Err(FrankenError::Busy);
}
if inner.writer_active {
return Err(FrankenError::Busy);
}
if mode == JournalMode::Wal && inner.wal_backend.is_none() {
return Err(FrankenError::Unsupported);
}
inner.journal_mode = mode;
drop(inner);
Ok(mode)
}
fn set_wal_backend(&self, backend: Box<dyn WalBackend>) -> Result<()> {
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimplePager lock poisoned"))?;
if inner.checkpoint_active {
return Err(FrankenError::Busy);
}
inner.wal_backend = Some(backend);
drop(inner);
Ok(())
}
}
impl<V: Vfs> SimplePager<V>
where
V::File: Send + Sync,
{
pub fn cache_metrics_snapshot(&self) -> Result<PageCacheMetricsSnapshot> {
let inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimplePager lock poisoned"))?;
Ok(inner.cache.metrics_snapshot())
}
pub fn reset_cache_metrics(&self) -> Result<()> {
self.inner
.lock()
.map_err(|_| FrankenError::internal("SimplePager lock poisoned"))?
.cache
.reset_metrics();
Ok(())
}
fn journal_path(db_path: &Path) -> PathBuf {
let mut jp = db_path.as_os_str().to_owned();
jp.push("-journal");
PathBuf::from(jp)
}
#[allow(clippy::too_many_lines)]
pub fn open(vfs: V, path: &Path, page_size: PageSize) -> Result<Self> {
let cx = Cx::new();
let vfs = Arc::new(vfs);
let flags = VfsOpenFlags::CREATE | VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_DB;
let (mut db_file, _actual_flags) = vfs.open(&cx, Some(path), flags)?;
let journal_path = Self::journal_path(path);
if vfs.access(&cx, &journal_path, AccessFlags::EXISTS)? {
Self::replay_journal(&cx, &*vfs, &mut db_file, &journal_path, page_size)?;
let _ = vfs.delete(&cx, &journal_path, true);
}
let mut file_size = db_file.file_size(&cx)?;
if file_size == 0 {
let page_len = page_size.as_usize();
let mut page1 = vec![0u8; page_len];
let header = DatabaseHeader {
page_size,
page_count: 1,
sqlite_version: FRANKENSQLITE_SQLITE_VERSION_NUMBER,
..DatabaseHeader::default()
};
let hdr_bytes = header.to_bytes().map_err(|err| {
FrankenError::internal(format!("failed to encode new database header: {err}"))
})?;
page1[..DATABASE_HEADER_SIZE].copy_from_slice(&hdr_bytes);
let usable = page_size.usable(header.reserved_per_page);
BTreePageHeader::write_empty_leaf_table(&mut page1, DATABASE_HEADER_SIZE, usable);
db_file.write(&cx, &page1, 0)?;
db_file.sync(&cx, SyncFlags::NORMAL)?;
file_size = db_file.file_size(&cx)?;
} else {
if file_size < DATABASE_HEADER_SIZE as u64 {
return Err(FrankenError::DatabaseCorrupt {
detail: format!(
"database file too small for header: {file_size} bytes (< {DATABASE_HEADER_SIZE})"
),
});
}
let mut header_bytes = [0u8; DATABASE_HEADER_SIZE];
let header_read = db_file.read(&cx, &mut header_bytes, 0)?;
if header_read < DATABASE_HEADER_SIZE {
return Err(FrankenError::DatabaseCorrupt {
detail: format!(
"short read fetching database header: got {header_read} of {DATABASE_HEADER_SIZE}"
),
});
}
let header = DatabaseHeader::from_bytes(&header_bytes).map_err(|error| {
FrankenError::DatabaseCorrupt {
detail: format!("invalid database header: {error}"),
}
})?;
if header.page_size != page_size {
return Err(FrankenError::DatabaseCorrupt {
detail: format!(
"database page size mismatch: header={} requested={}",
header.page_size.get(),
page_size.get()
),
});
}
}
let page_size_u64 = page_size.as_usize() as u64;
if file_size % page_size_u64 != 0 {
return Err(FrankenError::DatabaseCorrupt {
detail: format!(
"database file size {file_size} is not aligned to page size {}",
page_size.get()
),
});
}
let db_pages = file_size
.checked_div(page_size_u64)
.ok_or_else(|| FrankenError::internal("page size must be non-zero"))?;
let db_size = u32::try_from(db_pages).map_err(|_| FrankenError::OutOfRange {
what: "database page count".to_owned(),
value: db_pages.to_string(),
})?;
let next_page = if db_size >= 2 { db_size + 1 } else { 2 };
Ok(Self {
vfs,
db_path: path.to_owned(),
inner: Arc::new(Mutex::new(PagerInner {
db_file,
cache: PageCache::new(page_size, 256),
page_size,
db_size,
next_page,
writer_active: false,
active_transactions: 0,
checkpoint_active: false,
freelist: Vec::new(),
journal_mode: JournalMode::Delete,
wal_backend: None,
commit_seq: CommitSeq::ZERO,
})),
})
}
fn replay_journal(
cx: &Cx,
vfs: &V,
db_file: &mut V::File,
journal_path: &Path,
page_size: PageSize,
) -> Result<()> {
let jrnl_flags = VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_JOURNAL;
let Ok((mut jrnl_file, _)) = vfs.open(cx, Some(journal_path), jrnl_flags) else {
return Ok(()); };
let jrnl_size = jrnl_file.file_size(cx)?;
if jrnl_size < crate::journal::JOURNAL_HEADER_SIZE as u64 {
return Ok(()); }
let mut hdr_buf = vec![0u8; crate::journal::JOURNAL_HEADER_SIZE];
let _ = jrnl_file.read(cx, &mut hdr_buf, 0)?;
let Ok(header) = JournalHeader::decode(&hdr_buf) else {
return Ok(()); };
if header.page_size != page_size.get() {
return Err(FrankenError::DatabaseCorrupt {
detail: format!(
"hot journal page size mismatch: header={} expected={}",
header.page_size,
page_size.get()
),
});
}
let page_count = if header.page_count < 0 {
header.compute_page_count_from_file_size(jrnl_size)
} else {
#[allow(clippy::cast_sign_loss)]
let c = header.page_count as u32;
c
};
let header_size = u64::try_from(crate::journal::JOURNAL_HEADER_SIZE)
.expect("journal header size should fit in u64");
let hdr_padded = u64::from(header.sector_size).max(header_size);
let ps = page_size.as_usize();
let record_size = 4 + ps + 4;
let mut offset = hdr_padded;
for _ in 0..page_count {
let mut rec_buf = vec![0u8; record_size];
let bytes_read = jrnl_file.read(cx, &mut rec_buf, offset)?;
if bytes_read < record_size {
break; }
#[allow(clippy::cast_possible_truncation)]
let Ok(record) = JournalPageRecord::decode(&rec_buf, ps as u32) else {
break; };
if record.verify_checksum(header.nonce).is_err() {
break; }
let Some(page_no) = PageNumber::new(record.page_number) else {
return Err(FrankenError::DatabaseCorrupt {
detail: format!(
"hot journal contains invalid page number {}",
record.page_number
),
});
};
let page_offset = u64::from(page_no.get() - 1) * ps as u64;
db_file.write(cx, &record.content, page_offset)?;
offset += record_size as u64;
}
db_file.sync(cx, SyncFlags::NORMAL)?;
if header.initial_db_size > 0 {
let target_size = u64::from(header.initial_db_size) * ps as u64;
let current_size = db_file.file_size(cx)?;
if current_size > target_size {
db_file.truncate(cx, target_size)?;
}
}
Ok(())
}
}
struct SavepointEntry {
name: String,
write_set_snapshot: HashMap<PageNumber, Vec<u8>>,
freed_pages_snapshot: Vec<PageNumber>,
next_page_snapshot: u32,
freelist_snapshot: Vec<PageNumber>,
allocated_from_freelist_snapshot: Vec<PageNumber>,
}
pub struct SimpleTransaction<V: Vfs> {
vfs: Arc<V>,
journal_path: PathBuf,
inner: Arc<Mutex<PagerInner<V::File>>>,
write_set: HashMap<PageNumber, PageBuf>,
freed_pages: Vec<PageNumber>,
allocated_from_freelist: Vec<PageNumber>,
mode: TransactionMode,
is_writer: bool,
committed: bool,
finished: bool,
original_db_size: u32,
savepoint_stack: Vec<SavepointEntry>,
journal_mode: JournalMode,
pool: PageBufPool,
}
impl<V: Vfs> traits::sealed::Sealed for SimpleTransaction<V> {}
impl<V: Vfs> SimpleTransaction<V> {
#[must_use]
pub fn is_writer(&self) -> bool {
self.is_writer
}
}
impl<V> SimpleTransaction<V>
where
V: Vfs + Send,
V::File: Send + Sync,
{
#[allow(clippy::too_many_lines)]
fn commit_journal(
cx: &Cx,
vfs: &Arc<V>,
journal_path: &Path,
inner: &mut PagerInner<V::File>,
write_set: &HashMap<PageNumber, PageBuf>,
freed_pages: &mut Vec<PageNumber>,
original_db_size: u32,
) -> Result<()> {
if !write_set.is_empty() {
let nonce = 0x4652_414E; let page_size = inner.page_size;
let ps = page_size.as_usize();
let jrnl_flags =
VfsOpenFlags::CREATE | VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_JOURNAL;
let (mut jrnl_file, _) = vfs.open(cx, Some(journal_path), jrnl_flags)?;
let header = JournalHeader {
#[allow(clippy::cast_possible_truncation, clippy::cast_possible_wrap)]
page_count: write_set.len() as i32,
nonce,
initial_db_size: original_db_size,
sector_size: 512,
page_size: page_size.get(),
};
let hdr_bytes = header.encode_padded();
jrnl_file.write(cx, &hdr_bytes, 0)?;
let mut jrnl_offset = hdr_bytes.len() as u64;
for &page_no in write_set.keys() {
let mut pre_image = vec![0u8; ps];
if page_no.get() <= inner.db_size {
let disk_offset = u64::from(page_no.get() - 1) * ps as u64;
let bytes_read = inner.db_file.read(cx, &mut pre_image, disk_offset)?;
if bytes_read < ps {
return Err(FrankenError::DatabaseCorrupt {
detail: format!(
"short read while journaling pre-image for page {}: got {bytes_read} of {ps}",
page_no.get()
),
});
}
}
let record = JournalPageRecord::new(page_no.get(), pre_image, nonce);
let rec_bytes = record.encode();
jrnl_file.write(cx, &rec_bytes, jrnl_offset)?;
jrnl_offset += rec_bytes.len() as u64;
}
jrnl_file.sync(cx, SyncFlags::NORMAL)?;
for (page_no, data) in write_set {
inner.flush_page(cx, *page_no, data)?;
inner.db_size = inner.db_size.max(page_no.get());
}
inner.db_file.sync(cx, SyncFlags::NORMAL)?;
let _ = vfs.delete(cx, journal_path, true);
}
for page_no in freed_pages.drain(..) {
inner.freelist.push(page_no);
}
Ok(())
}
fn commit_wal(
cx: &Cx,
inner: &mut PagerInner<V::File>,
write_set: &HashMap<PageNumber, PageBuf>,
freed_pages: &mut Vec<PageNumber>,
) -> Result<()> {
if !write_set.is_empty() {
let wal = inner.wal_backend.as_mut().ok_or_else(|| {
FrankenError::internal("WAL mode active but no WAL backend installed")
})?;
let page_count = write_set.len();
let mut written = 0_usize;
for (page_no, data) in write_set {
written += 1;
let db_size_if_commit = if written == page_count {
let max_written = write_set.keys().map(|p| p.get()).max().unwrap_or(0);
inner.db_size.max(max_written)
} else {
0
};
wal.append_frame(cx, page_no.get(), data, db_size_if_commit)?;
}
let wal = inner
.wal_backend
.as_mut()
.ok_or_else(|| FrankenError::internal("WAL backend disappeared during commit"))?;
wal.sync(cx)?;
for page_no in write_set.keys() {
inner.db_size = inner.db_size.max(page_no.get());
}
}
for page_no in freed_pages.drain(..) {
inner.freelist.push(page_no);
}
Ok(())
}
fn ensure_writer(&mut self) -> Result<()> {
if self.is_writer {
return Ok(());
}
match self.mode {
TransactionMode::ReadOnly => Err(FrankenError::ReadOnly),
TransactionMode::Deferred | TransactionMode::Concurrent => {
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimpleTransaction lock poisoned"))?;
if inner.checkpoint_active {
return Err(FrankenError::Busy);
}
if inner.writer_active {
return Err(FrankenError::Busy);
}
inner.writer_active = true;
drop(inner);
self.is_writer = true;
Ok(())
}
TransactionMode::Immediate | TransactionMode::Exclusive => Err(FrankenError::internal(
"writer transaction lost writer role",
)),
}
}
}
impl<V> TransactionHandle for SimpleTransaction<V>
where
V: Vfs + Send,
V::File: Send + Sync,
{
fn get_page(&self, cx: &Cx, page_no: PageNumber) -> Result<PageData> {
if let Some(data) = self.write_set.get(&page_no) {
return Ok(PageData::from_vec(data.to_vec()));
}
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimpleTransaction lock poisoned"))?;
let data = inner.read_page_copy(cx, page_no)?;
drop(inner);
Ok(PageData::from_vec(data))
}
fn write_page(&mut self, _cx: &Cx, page_no: PageNumber, data: &[u8]) -> Result<()> {
self.ensure_writer()?;
if let Some(pos) = self.freed_pages.iter().position(|&p| p == page_no) {
self.freed_pages.swap_remove(pos);
}
let mut buf = self.pool.acquire()?;
let len = buf.len().min(data.len());
buf[..len].copy_from_slice(&data[..len]);
if len < buf.len() {
buf[len..].fill(0);
}
self.write_set.insert(page_no, buf);
Ok(())
}
fn allocate_page(&mut self, _cx: &Cx) -> Result<PageNumber> {
self.ensure_writer()?;
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimpleTransaction lock poisoned"))?;
if let Some(page) = inner.freelist.pop() {
self.allocated_from_freelist.push(page);
return Ok(page);
}
let raw = inner.next_page;
inner.next_page = inner.next_page.saturating_add(1);
drop(inner);
PageNumber::new(raw).ok_or_else(|| FrankenError::OutOfRange {
what: "allocated page number".to_owned(),
value: raw.to_string(),
})
}
fn free_page(&mut self, _cx: &Cx, page_no: PageNumber) -> Result<()> {
self.ensure_writer()?;
if page_no == PageNumber::ONE {
return Err(FrankenError::OutOfRange {
what: "free page number".to_owned(),
value: page_no.get().to_string(),
});
}
if !self.freed_pages.contains(&page_no) {
self.freed_pages.push(page_no);
}
self.write_set.remove(&page_no);
Ok(())
}
#[allow(clippy::too_many_lines)]
fn commit(&mut self, cx: &Cx) -> Result<()> {
if self.finished {
return Ok(());
}
if !self.is_writer {
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimpleTransaction lock poisoned"))?;
inner.active_transactions = inner.active_transactions.saturating_sub(1);
drop(inner);
self.committed = true;
self.finished = true;
return Ok(());
}
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimpleTransaction lock poisoned"))?;
let commit_result = if self.journal_mode == JournalMode::Wal {
Self::commit_wal(cx, &mut inner, &self.write_set, &mut self.freed_pages)
} else {
Self::commit_journal(
cx,
&self.vfs,
&self.journal_path,
&mut inner,
&self.write_set,
&mut self.freed_pages,
self.original_db_size,
)
};
if commit_result.is_ok() {
inner.commit_seq = inner.commit_seq.next();
inner.active_transactions = inner.active_transactions.saturating_sub(1);
inner.writer_active = false;
drop(inner);
self.write_set.clear();
self.committed = true;
self.finished = true;
} else {
drop(inner);
}
commit_result
}
fn rollback(&mut self, cx: &Cx) -> Result<()> {
if self.finished {
return Ok(());
}
self.write_set.clear();
self.freed_pages.clear();
self.savepoint_stack.clear();
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimpleTransaction lock poisoned"))?;
for page in self.allocated_from_freelist.drain(..) {
inner.freelist.push(page);
}
if self.is_writer {
inner.db_size = self.original_db_size;
let db_size = inner.db_size;
inner.next_page = if db_size >= 2 { db_size + 1 } else { 2 };
inner.writer_active = false;
}
inner.active_transactions = inner.active_transactions.saturating_sub(1);
drop(inner);
if self.is_writer {
let _ = self.vfs.delete(cx, &self.journal_path, true);
}
self.committed = false;
self.finished = true;
Ok(())
}
fn savepoint(&mut self, _cx: &Cx, name: &str) -> Result<()> {
let inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimpleTransaction lock poisoned"))?;
self.savepoint_stack.push(SavepointEntry {
name: name.to_owned(),
write_set_snapshot: self
.write_set
.iter()
.map(|(&k, v)| (k, v.to_vec()))
.collect(),
freed_pages_snapshot: self.freed_pages.clone(),
next_page_snapshot: inner.next_page,
freelist_snapshot: inner.freelist.clone(),
allocated_from_freelist_snapshot: self.allocated_from_freelist.clone(),
});
drop(inner);
Ok(())
}
fn release_savepoint(&mut self, _cx: &Cx, name: &str) -> Result<()> {
let pos = self
.savepoint_stack
.iter()
.rposition(|sp| sp.name == name)
.ok_or_else(|| FrankenError::internal(format!("no savepoint named '{name}'")))?;
self.savepoint_stack.truncate(pos);
Ok(())
}
fn rollback_to_savepoint(&mut self, _cx: &Cx, name: &str) -> Result<()> {
let pos = self
.savepoint_stack
.iter()
.rposition(|sp| sp.name == name)
.ok_or_else(|| FrankenError::internal(format!("no savepoint named '{name}'")))?;
let entry = &self.savepoint_stack[pos];
self.write_set = entry
.write_set_snapshot
.iter()
.map(|(&k, v)| -> Result<(PageNumber, PageBuf)> {
let mut buf = self.pool.acquire()?;
let len = buf.len().min(v.len());
buf.as_mut_slice()[..len].copy_from_slice(&v[..len]);
if len < buf.len() {
buf[len..].fill(0);
}
Ok((k, buf))
})
.collect::<Result<HashMap<_, _>>>()?;
self.freed_pages = entry.freed_pages_snapshot.clone();
self.allocated_from_freelist = entry.allocated_from_freelist_snapshot.clone();
{
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimpleTransaction lock poisoned"))?;
inner.next_page = entry.next_page_snapshot;
inner.freelist.clone_from(&entry.freelist_snapshot);
}
self.savepoint_stack.truncate(pos + 1);
Ok(())
}
}
impl<V: Vfs> Drop for SimpleTransaction<V> {
fn drop(&mut self) {
if self.finished {
return;
}
if let Ok(mut inner) = self.inner.lock() {
if self.is_writer {
inner.writer_active = false;
}
inner.active_transactions = inner.active_transactions.saturating_sub(1);
}
self.finished = true;
}
}
pub struct SimplePagerCheckpointWriter<V: Vfs>
where
V::File: Send + Sync,
{
inner: Arc<Mutex<PagerInner<V::File>>>,
}
impl<V: Vfs> traits::sealed::Sealed for SimplePagerCheckpointWriter<V> where V::File: Send + Sync {}
impl<V> traits::CheckpointPageWriter for SimplePagerCheckpointWriter<V>
where
V: Vfs + Send + Sync,
V::File: Send + Sync,
{
fn write_page(&mut self, cx: &Cx, page_no: PageNumber, data: &[u8]) -> Result<()> {
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimplePagerCheckpointWriter lock poisoned"))?;
let page_size = inner.page_size.as_usize();
let offset = u64::from(page_no.get() - 1) * page_size as u64;
inner.db_file.write(cx, data, offset)?;
inner.cache.evict(page_no);
inner.db_size = inner.db_size.max(page_no.get());
drop(inner);
Ok(())
}
fn truncate(&mut self, cx: &Cx, n_pages: u32) -> Result<()> {
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimplePagerCheckpointWriter lock poisoned"))?;
let old_db_size = inner.db_size;
let page_size = inner.page_size.as_usize();
let target_size = u64::from(n_pages) * page_size as u64;
inner.db_file.truncate(cx, target_size)?;
inner.db_size = n_pages;
for pgno in (n_pages + 1)..=old_db_size {
if let Some(page_no) = PageNumber::new(pgno) {
inner.cache.evict(page_no);
}
}
drop(inner);
Ok(())
}
fn sync(&mut self, cx: &Cx) -> Result<()> {
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimplePagerCheckpointWriter lock poisoned"))?;
inner.db_file.sync(cx, SyncFlags::NORMAL)
}
}
impl<V: Vfs> SimplePager<V>
where
V::File: Send + Sync,
{
#[must_use]
pub fn checkpoint_writer(&self) -> SimplePagerCheckpointWriter<V> {
SimplePagerCheckpointWriter {
inner: Arc::clone(&self.inner),
}
}
pub fn checkpoint(
&self,
cx: &Cx,
mode: traits::CheckpointMode,
) -> Result<traits::CheckpointResult> {
let mut wal = {
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimplePager lock poisoned"))?;
if inner.journal_mode != JournalMode::Wal {
return Err(FrankenError::Unsupported);
}
if inner.checkpoint_active {
return Err(FrankenError::Busy);
}
if inner.active_transactions > 0 {
return Err(FrankenError::Busy);
}
inner.checkpoint_active = true;
let Some(wal) = inner.wal_backend.take() else {
inner.checkpoint_active = false;
return Err(FrankenError::internal(
"WAL mode active but no WAL backend installed",
));
};
drop(inner);
wal
};
let mut writer = self.checkpoint_writer();
let result = wal.checkpoint(cx, mode, &mut writer, 0, None);
{
let mut inner = self
.inner
.lock()
.map_err(|_| FrankenError::internal("SimplePager lock poisoned"))?;
inner.wal_backend = Some(wal);
inner.checkpoint_active = false;
}
result
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::traits::{MvccPager, TransactionHandle, TransactionMode};
use fsqlite_types::PageSize;
use fsqlite_types::flags::AccessFlags;
use fsqlite_types::{BTreePageHeader, DatabaseHeader};
use fsqlite_vfs::MemoryVfs;
use std::path::PathBuf;
const BEAD_ID: &str = "bd-bca.1";
fn test_pager() -> (SimplePager<MemoryVfs>, PathBuf) {
let vfs = MemoryVfs::new();
let path = PathBuf::from("/test.db");
let pager = SimplePager::open(vfs, &path, PageSize::DEFAULT).unwrap();
(pager, path)
}
#[test]
fn test_open_empty_database() {
let (pager, _) = test_pager();
let inner = pager.inner.lock().unwrap();
assert_eq!(inner.db_size, 1, "bead_id={BEAD_ID} case=empty_db_size");
assert_eq!(
inner.page_size,
PageSize::DEFAULT,
"bead_id={BEAD_ID} case=page_size_default"
);
drop(inner);
let cx = Cx::new();
let txn = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let raw_page = txn.get_page(&cx, PageNumber::ONE).unwrap().into_vec();
let hdr: [u8; DATABASE_HEADER_SIZE] = raw_page[..DATABASE_HEADER_SIZE]
.try_into()
.expect("page 1 must contain database header");
let parsed = DatabaseHeader::from_bytes(&hdr).expect("header should parse");
assert_eq!(
parsed.page_size,
PageSize::DEFAULT,
"bead_id={BEAD_ID} case=page1_header_page_size"
);
assert_eq!(
parsed.page_count, 1,
"bead_id={BEAD_ID} case=page1_header_page_count"
);
let btree_hdr =
BTreePageHeader::parse(&raw_page, PageSize::DEFAULT, 0, true).expect("btree header");
assert_eq!(
btree_hdr.cell_count, 0,
"bead_id={BEAD_ID} case=sqlite_master_initially_empty"
);
}
#[test]
fn test_open_existing_database_rejects_page_size_mismatch() {
let vfs = MemoryVfs::new();
let path = PathBuf::from("/page_size_mismatch.db");
let _pager = SimplePager::open(vfs.clone(), &path, PageSize::DEFAULT).unwrap();
let wrong_page_size = PageSize::new(8192).unwrap();
let Err(err) = SimplePager::open(vfs, &path, wrong_page_size) else {
panic!("expected page size mismatch error");
};
assert!(
matches!(err, FrankenError::DatabaseCorrupt { .. }),
"bead_id={BEAD_ID} case=reject_page_size_mismatch"
);
}
#[test]
fn test_open_existing_database_rejects_non_page_aligned_size() {
let vfs = MemoryVfs::new();
let path = PathBuf::from("/misaligned.db");
let cx = Cx::new();
let _pager = SimplePager::open(vfs.clone(), &path, PageSize::DEFAULT).unwrap();
let flags = VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_DB;
let (mut db_file, _) = vfs.open(&cx, Some(&path), flags).unwrap();
let file_size = db_file.file_size(&cx).unwrap();
db_file.write(&cx, &[0xAB], file_size).unwrap();
let Err(err) = SimplePager::open(vfs, &path, PageSize::DEFAULT) else {
panic!("expected non-page-aligned file size error");
};
assert!(
matches!(err, FrankenError::DatabaseCorrupt { .. }),
"bead_id={BEAD_ID} case=reject_non_page_aligned_file_size"
);
}
#[test]
fn test_begin_readonly_transaction() {
let (pager, _) = test_pager();
let cx = Cx::new();
let txn = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
assert!(!txn.is_writer, "bead_id={BEAD_ID} case=readonly_not_writer");
}
#[test]
fn test_begin_write_transaction() {
let (pager, _) = test_pager();
let cx = Cx::new();
let txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
assert!(txn.is_writer, "bead_id={BEAD_ID} case=immediate_is_writer");
}
#[test]
fn test_begin_deferred_transaction_starts_reader() {
let (pager, _) = test_pager();
let cx = Cx::new();
let txn = pager.begin(&cx, TransactionMode::Deferred).unwrap();
assert!(
!txn.is_writer,
"bead_id={BEAD_ID} case=deferred_starts_readonly"
);
}
#[test]
fn test_begin_concurrent_transaction_starts_reader() {
let (pager, _) = test_pager();
let cx = Cx::new();
let txn = pager.begin(&cx, TransactionMode::Concurrent).unwrap();
assert!(
!txn.is_writer,
"bead_id={BEAD_ID} case=concurrent_starts_readonly"
);
}
#[test]
fn test_deferred_upgrades_on_first_write_intent() {
let (pager, _) = test_pager();
let cx = Cx::new();
let mut deferred = pager.begin(&cx, TransactionMode::Deferred).unwrap();
assert!(
!deferred.is_writer,
"bead_id={BEAD_ID} case=deferred_pre_upgrade"
);
let _page = deferred.allocate_page(&cx).unwrap();
assert!(
deferred.is_writer,
"bead_id={BEAD_ID} case=deferred_upgraded_to_writer"
);
}
#[test]
fn test_deferred_upgrade_busy_when_writer_active() {
let (pager, _) = test_pager();
let cx = Cx::new();
let _writer = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let mut deferred = pager.begin(&cx, TransactionMode::Deferred).unwrap();
let err = deferred.allocate_page(&cx).unwrap_err();
assert!(matches!(err, FrankenError::Busy));
}
#[test]
fn test_concurrent_writer_blocked() {
let (pager, _) = test_pager();
let cx = Cx::new();
let _txn1 = pager.begin(&cx, TransactionMode::Exclusive).unwrap();
let result = pager.begin(&cx, TransactionMode::Immediate);
assert!(
result.is_err(),
"bead_id={BEAD_ID} case=concurrent_writer_busy"
);
}
#[test]
fn test_multiple_readers_allowed() {
let (pager, _) = test_pager();
let cx = Cx::new();
let _r1 = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let _r2 = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
}
#[test]
fn test_write_page_and_read_back() {
let (pager, _) = test_pager();
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let page_no = txn.allocate_page(&cx).unwrap();
let page_size = PageSize::DEFAULT.as_usize();
let mut data = vec![0_u8; page_size];
data[0] = 0xDE;
data[1] = 0xAD;
txn.write_page(&cx, page_no, &data).unwrap();
let read_back = txn.get_page(&cx, page_no).unwrap();
assert_eq!(
read_back.as_ref()[0],
0xDE,
"bead_id={BEAD_ID} case=read_back_byte0"
);
assert_eq!(
read_back.as_ref()[1],
0xAD,
"bead_id={BEAD_ID} case=read_back_byte1"
);
}
#[test]
fn test_commit_persists_pages() {
let (pager, _) = test_pager();
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let page_no = txn.allocate_page(&cx).unwrap();
let page_size = PageSize::DEFAULT.as_usize();
let mut data = vec![0_u8; page_size];
data[0..4].copy_from_slice(&[0xCA, 0xFE, 0xBA, 0xBE]);
txn.write_page(&cx, page_no, &data).unwrap();
txn.commit(&cx).unwrap();
let txn2 = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let read_back = txn2.get_page(&cx, page_no).unwrap();
assert_eq!(
&read_back.as_ref()[0..4],
&[0xCA, 0xFE, 0xBA, 0xBE],
"bead_id={BEAD_ID} case=commit_persists"
);
}
#[test]
fn test_rollback_discards_writes() {
let (pager, _) = test_pager();
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let page_no = txn.allocate_page(&cx).unwrap();
let page_size = PageSize::DEFAULT.as_usize();
let original = vec![0x11_u8; page_size];
txn.write_page(&cx, page_no, &original).unwrap();
txn.commit(&cx).unwrap();
let mut txn2 = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let modified = vec![0x99_u8; page_size];
txn2.write_page(&cx, page_no, &modified).unwrap();
txn2.rollback(&cx).unwrap();
let txn3 = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let read_back = txn3.get_page(&cx, page_no).unwrap();
assert_eq!(
read_back.as_ref()[0],
0x11,
"bead_id={BEAD_ID} case=rollback_restores"
);
}
#[test]
fn test_allocate_returns_sequential_pages() {
let (pager, _) = test_pager();
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
let p2 = txn.allocate_page(&cx).unwrap();
assert!(
p2.get() > p1.get(),
"bead_id={BEAD_ID} case=sequential_alloc p1={} p2={}",
p1.get(),
p2.get()
);
}
#[test]
fn test_free_page_reuses_on_next_alloc() {
let (pager, _) = test_pager();
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
let p2 = txn.allocate_page(&cx).unwrap();
let page_size = PageSize::DEFAULT.as_usize();
txn.write_page(&cx, p1, &vec![1_u8; page_size]).unwrap();
txn.write_page(&cx, p2, &vec![2_u8; page_size]).unwrap();
txn.commit(&cx).unwrap();
let mut txn2 = pager.begin(&cx, TransactionMode::Immediate).unwrap();
txn2.free_page(&cx, p1).unwrap();
txn2.commit(&cx).unwrap();
let mut txn3 = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p3 = txn3.allocate_page(&cx).unwrap();
assert_eq!(
p3,
p1,
"bead_id={BEAD_ID} case=freelist_reuse p3={} p1={}",
p3.get(),
p1.get()
);
}
#[test]
fn test_cannot_free_page_one() {
let (pager, _) = test_pager();
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let result = txn.free_page(&cx, PageNumber::ONE);
assert!(
result.is_err(),
"bead_id={BEAD_ID} case=cannot_free_page_one"
);
}
#[test]
fn test_readonly_cannot_write() {
let (pager, _) = test_pager();
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let result = txn.write_page(&cx, PageNumber::ONE, &[0_u8; 4096]);
assert!(
result.is_err(),
"bead_id={BEAD_ID} case=readonly_cannot_write"
);
}
#[test]
fn test_readonly_cannot_allocate() {
let (pager, _) = test_pager();
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let result = txn.allocate_page(&cx);
assert!(
result.is_err(),
"bead_id={BEAD_ID} case=readonly_cannot_allocate"
);
}
#[test]
fn test_drop_uncommitted_writer_releases_lock() {
let (pager, _) = test_pager();
let cx = Cx::new();
{
let _txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
}
let txn2 = pager.begin(&cx, TransactionMode::Immediate);
assert!(
txn2.is_ok(),
"bead_id={BEAD_ID} case=drop_releases_writer_lock"
);
}
#[test]
fn test_commit_then_drop_no_double_release() {
let (pager, _) = test_pager();
let cx = Cx::new();
{
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
txn.commit(&cx).unwrap();
}
let txn2 = pager.begin(&cx, TransactionMode::Immediate);
assert!(
txn2.is_ok(),
"bead_id={BEAD_ID} case=commit_releases_writer"
);
}
#[test]
fn test_double_commit_is_idempotent() {
let (pager, _) = test_pager();
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
txn.commit(&cx).unwrap();
txn.commit(&cx).unwrap();
}
#[test]
fn test_multi_page_write_commit_read() {
let (pager, _) = test_pager();
let cx = Cx::new();
let page_size = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let mut allocated_pages = Vec::new();
for i in 0_u8..5 {
let p = txn.allocate_page(&cx).unwrap();
let data = vec![i; page_size];
txn.write_page(&cx, p, &data).unwrap();
allocated_pages.push(p);
}
txn.commit(&cx).unwrap();
let txn2 = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
for (i, p) in allocated_pages.iter().enumerate() {
let data = txn2.get_page(&cx, *p).unwrap();
#[allow(clippy::cast_possible_truncation)]
let expected = i as u8;
assert_eq!(
data.as_ref()[0],
expected,
"bead_id={BEAD_ID} case=multi_page idx={i}"
);
}
}
#[test]
fn test_commit_journal_short_preimage_read_errors() {
let vfs = MemoryVfs::new();
let path = PathBuf::from("/short_preimage.db");
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let pager = SimplePager::open(vfs.clone(), &path, PageSize::DEFAULT).unwrap();
let page_two = {
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p, &vec![0x11; ps]).unwrap();
txn.commit(&cx).unwrap();
p
};
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
txn.write_page(&cx, page_two, &vec![0x22; ps]).unwrap();
let flags = VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_DB;
let (mut db_file, _) = vfs.open(&cx, Some(&path), flags).unwrap();
db_file
.truncate(&cx, PageSize::DEFAULT.as_usize() as u64)
.unwrap();
let err = txn.commit(&cx).unwrap_err();
assert!(
matches!(err, FrankenError::DatabaseCorrupt { .. }),
"bead_id={BEAD_ID} case=short_preimage_read_is_corruption"
);
let Err(busy) = pager.begin(&cx, TransactionMode::Immediate) else {
panic!("expected begin to fail while writer lock is still held");
};
assert!(
matches!(busy, FrankenError::Busy),
"bead_id={BEAD_ID} case=commit_error_keeps_writer_lock"
);
txn.rollback(&cx).unwrap();
let _next_writer = pager.begin(&cx, TransactionMode::Immediate).unwrap();
}
#[test]
fn test_commit_creates_and_deletes_journal() {
let vfs = MemoryVfs::new();
let path = PathBuf::from("/jrnl_test.db");
let pager = SimplePager::open(vfs.clone(), &path, PageSize::DEFAULT).unwrap();
let cx = Cx::new();
let journal_path = SimplePager::<MemoryVfs>::journal_path(&path);
assert!(
!vfs.access(&cx, &journal_path, AccessFlags::EXISTS).unwrap(),
"bead_id={BEAD_ID} case=no_journal_before_commit"
);
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p, &vec![0xAA; 4096]).unwrap();
txn.commit(&cx).unwrap();
assert!(
!vfs.access(&cx, &journal_path, AccessFlags::EXISTS).unwrap(),
"bead_id={BEAD_ID} case=journal_deleted_after_commit"
);
}
#[test]
fn test_hot_journal_recovery_restores_original_data() {
let vfs = MemoryVfs::new();
let path = PathBuf::from("/crash_test.db");
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
{
let pager = SimplePager::open(vfs.clone(), &path, PageSize::DEFAULT).unwrap();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
assert_eq!(p1.get(), 2);
txn.write_page(&cx, p1, &vec![0x11; ps]).unwrap();
txn.commit(&cx).unwrap();
}
{
let flags = VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_DB;
let (mut db_file, _) = vfs.open(&cx, Some(&path), flags).unwrap();
let corrupt_data = vec![0x99; ps];
let offset = u64::from(2_u32 - 1) * ps as u64;
db_file.write(&cx, &corrupt_data, offset).unwrap();
}
{
let journal_path = SimplePager::<MemoryVfs>::journal_path(&path);
let jrnl_flags =
VfsOpenFlags::CREATE | VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_JOURNAL;
let (mut jrnl, _) = vfs.open(&cx, Some(&journal_path), jrnl_flags).unwrap();
let nonce = 42;
let header = JournalHeader {
page_count: 1,
nonce,
initial_db_size: 2,
sector_size: 512,
page_size: 4096,
};
let hdr_bytes = header.encode_padded();
jrnl.write(&cx, &hdr_bytes, 0).unwrap();
let record = JournalPageRecord::new(2, vec![0x11; ps], nonce);
let rec_bytes = record.encode();
jrnl.write(&cx, &rec_bytes, hdr_bytes.len() as u64).unwrap();
}
{
let pager = SimplePager::open(vfs.clone(), &path, PageSize::DEFAULT).unwrap();
let txn = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let page_no_2 = PageNumber::new(2).unwrap();
let data = txn.get_page(&cx, page_no_2).unwrap();
assert_eq!(
data.as_ref()[0],
0x11,
"bead_id={BEAD_ID} case=journal_recovery_restores"
);
assert_eq!(
data.as_ref()[ps - 1],
0x11,
"bead_id={BEAD_ID} case=journal_recovery_restores_last_byte"
);
}
let journal_path = SimplePager::<MemoryVfs>::journal_path(&path);
assert!(
!vfs.access(&cx, &journal_path, AccessFlags::EXISTS).unwrap(),
"bead_id={BEAD_ID} case=journal_deleted_after_recovery"
);
}
#[test]
fn test_hot_journal_truncated_record_stops_replay() {
let vfs = MemoryVfs::new();
let path = PathBuf::from("/trunc_jrnl.db");
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
{
let pager = SimplePager::open(vfs.clone(), &path, PageSize::DEFAULT).unwrap();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
let p2 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0xAA; ps]).unwrap();
txn.write_page(&cx, p2, &vec![0xBB; ps]).unwrap();
txn.commit(&cx).unwrap();
}
{
let flags = VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_DB;
let (mut db_file, _) = vfs.open(&cx, Some(&path), flags).unwrap();
db_file
.write(&cx, &vec![0xFF; ps], u64::from(2_u32 - 1) * ps as u64)
.unwrap();
}
{
let journal_path = SimplePager::<MemoryVfs>::journal_path(&path);
let jrnl_flags =
VfsOpenFlags::CREATE | VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_JOURNAL;
let (mut jrnl, _) = vfs.open(&cx, Some(&journal_path), jrnl_flags).unwrap();
let nonce = 7;
let header = JournalHeader {
page_count: 2,
nonce,
initial_db_size: 3,
sector_size: 512,
page_size: 4096,
};
let hdr_bytes = header.encode_padded();
jrnl.write(&cx, &hdr_bytes, 0).unwrap();
let rec1 = JournalPageRecord::new(3, vec![0xCC; ps], nonce);
let rec1_bytes = rec1.encode();
jrnl.write(&cx, &rec1_bytes, hdr_bytes.len() as u64)
.unwrap();
let rec2 = JournalPageRecord::new(2, vec![0xBB; ps], nonce);
let rec2_bytes = rec2.encode();
let trunc_len = rec2_bytes.len() / 2;
let offset = hdr_bytes.len() as u64 + rec1_bytes.len() as u64;
jrnl.write(&cx, &rec2_bytes[..trunc_len], offset).unwrap();
}
{
let pager = SimplePager::open(vfs, &path, PageSize::DEFAULT).unwrap();
let txn = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let page_no_2 = PageNumber::new(2).unwrap();
let data2 = txn.get_page(&cx, page_no_2).unwrap();
assert_eq!(
data2.as_ref()[0],
0xFF,
"bead_id={BEAD_ID} case=truncated_journal_page2_not_restored"
);
}
}
#[test]
fn test_hot_journal_checksum_mismatch_stops_replay() {
let vfs = MemoryVfs::new();
let path = PathBuf::from("/cksum_jrnl.db");
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
{
let pager = SimplePager::open(vfs.clone(), &path, PageSize::DEFAULT).unwrap();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p, &vec![0x55; ps]).unwrap();
txn.commit(&cx).unwrap();
}
{
let flags = VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_DB;
let (mut db_file, _) = vfs.open(&cx, Some(&path), flags).unwrap();
db_file
.write(&cx, &vec![0xEE; ps], u64::from(2_u32 - 1) * ps as u64)
.unwrap();
}
{
let journal_path = SimplePager::<MemoryVfs>::journal_path(&path);
let jrnl_flags =
VfsOpenFlags::CREATE | VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_JOURNAL;
let (mut jrnl, _) = vfs.open(&cx, Some(&journal_path), jrnl_flags).unwrap();
let nonce = 99;
let header = JournalHeader {
page_count: 1,
nonce,
initial_db_size: 2,
sector_size: 512,
page_size: 4096,
};
let hdr_bytes = header.encode_padded();
jrnl.write(&cx, &hdr_bytes, 0).unwrap();
let record = JournalPageRecord::new(2, vec![0x55; ps], nonce + 1);
let rec_bytes = record.encode();
jrnl.write(&cx, &rec_bytes, hdr_bytes.len() as u64).unwrap();
}
{
let pager = SimplePager::open(vfs, &path, PageSize::DEFAULT).unwrap();
let txn = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let page_no_2 = PageNumber::new(2).unwrap();
let data = txn.get_page(&cx, page_no_2).unwrap();
assert_eq!(
data.as_ref()[0],
0xEE,
"bead_id={BEAD_ID} case=bad_checksum_stops_replay"
);
}
}
#[test]
fn test_hot_journal_invalid_page_number_errors() {
let vfs = MemoryVfs::new();
let path = PathBuf::from("/bad_pgno_jrnl.db");
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
{
let _pager = SimplePager::open(vfs.clone(), &path, PageSize::DEFAULT).unwrap();
let journal_path = SimplePager::<MemoryVfs>::journal_path(&path);
let jrnl_flags =
VfsOpenFlags::CREATE | VfsOpenFlags::READWRITE | VfsOpenFlags::MAIN_JOURNAL;
let (mut jrnl, _) = vfs.open(&cx, Some(&journal_path), jrnl_flags).unwrap();
let nonce = 321;
let header = JournalHeader {
page_count: 1,
nonce,
initial_db_size: 1,
sector_size: 512,
page_size: 4096,
};
let hdr_bytes = header.encode_padded();
jrnl.write(&cx, &hdr_bytes, 0).unwrap();
let record = JournalPageRecord::new(0, vec![0xAA; ps], nonce);
let rec_bytes = record.encode();
jrnl.write(&cx, &rec_bytes, hdr_bytes.len() as u64).unwrap();
}
let Err(err) = SimplePager::open(vfs, &path, PageSize::DEFAULT) else {
panic!("expected invalid journal page number error");
};
assert!(
matches!(err, FrankenError::DatabaseCorrupt { .. }),
"bead_id={BEAD_ID} case=invalid_journal_page_number_rejected"
);
}
#[test]
fn test_journal_not_created_for_readonly_commit() {
let vfs = MemoryVfs::new();
let path = PathBuf::from("/readonly_jrnl.db");
let pager = SimplePager::open(vfs.clone(), &path, PageSize::DEFAULT).unwrap();
let cx = Cx::new();
let journal_path = SimplePager::<MemoryVfs>::journal_path(&path);
let mut txn = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
txn.commit(&cx).unwrap();
assert!(
!vfs.access(&cx, &journal_path, AccessFlags::EXISTS).unwrap(),
"bead_id={BEAD_ID} case=no_journal_for_readonly"
);
}
#[test]
fn test_rollback_deletes_journal() {
let vfs = MemoryVfs::new();
let path = PathBuf::from("/rollback_jrnl.db");
let pager = SimplePager::open(vfs.clone(), &path, PageSize::DEFAULT).unwrap();
let cx = Cx::new();
let journal_path = SimplePager::<MemoryVfs>::journal_path(&path);
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p, &vec![0xDD; 4096]).unwrap();
txn.rollback(&cx).unwrap();
assert!(
!vfs.access(&cx, &journal_path, AccessFlags::EXISTS).unwrap(),
"bead_id={BEAD_ID} case=journal_deleted_on_rollback"
);
}
#[test]
fn test_savepoint_basic_rollback_to() {
let (pager, _) = test_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0x11; ps]).unwrap();
txn.savepoint(&cx, "sp1").unwrap();
let p2 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p2, &vec![0x22; ps]).unwrap();
txn.write_page(&cx, p1, &vec![0x33; ps]).unwrap();
txn.rollback_to_savepoint(&cx, "sp1").unwrap();
let data = txn.get_page(&cx, p1).unwrap();
assert_eq!(
data.as_ref()[0],
0x11,
"bead_id={BEAD_ID} case=savepoint_rollback_restores_p1"
);
let data2 = txn.get_page(&cx, p2).unwrap();
assert_eq!(
data2.as_ref()[0],
0x00,
"bead_id={BEAD_ID} case=savepoint_rollback_removes_p2"
);
txn.commit(&cx).unwrap();
}
#[test]
fn test_savepoint_release_keeps_changes() {
let (pager, _) = test_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0xAA; ps]).unwrap();
txn.savepoint(&cx, "sp1").unwrap();
txn.write_page(&cx, p1, &vec![0xBB; ps]).unwrap();
txn.release_savepoint(&cx, "sp1").unwrap();
let data = txn.get_page(&cx, p1).unwrap();
assert_eq!(
data.as_ref()[0],
0xBB,
"bead_id={BEAD_ID} case=release_keeps_changes"
);
txn.commit(&cx).unwrap();
}
#[test]
fn test_savepoint_nested_rollback_to_inner() {
let (pager, _) = test_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0x11; ps]).unwrap();
txn.savepoint(&cx, "outer").unwrap();
txn.write_page(&cx, p1, &vec![0x22; ps]).unwrap();
txn.savepoint(&cx, "inner").unwrap();
txn.write_page(&cx, p1, &vec![0x33; ps]).unwrap();
txn.rollback_to_savepoint(&cx, "inner").unwrap();
let data = txn.get_page(&cx, p1).unwrap();
assert_eq!(
data.as_ref()[0],
0x22,
"bead_id={BEAD_ID} case=nested_rollback_inner"
);
txn.commit(&cx).unwrap();
}
#[test]
fn test_savepoint_nested_rollback_to_outer() {
let (pager, _) = test_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0x11; ps]).unwrap();
txn.savepoint(&cx, "outer").unwrap();
txn.write_page(&cx, p1, &vec![0x22; ps]).unwrap();
txn.savepoint(&cx, "inner").unwrap();
txn.write_page(&cx, p1, &vec![0x33; ps]).unwrap();
txn.rollback_to_savepoint(&cx, "outer").unwrap();
let data = txn.get_page(&cx, p1).unwrap();
assert_eq!(
data.as_ref()[0],
0x11,
"bead_id={BEAD_ID} case=nested_rollback_outer"
);
txn.commit(&cx).unwrap();
}
#[test]
fn test_savepoint_rollback_to_preserves_savepoint() {
let (pager, _) = test_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0x11; ps]).unwrap();
txn.savepoint(&cx, "sp1").unwrap();
txn.write_page(&cx, p1, &vec![0x22; ps]).unwrap();
txn.rollback_to_savepoint(&cx, "sp1").unwrap();
let data = txn.get_page(&cx, p1).unwrap();
assert_eq!(data.as_ref()[0], 0x11);
txn.write_page(&cx, p1, &vec![0x33; ps]).unwrap();
txn.commit(&cx).unwrap();
let txn2 = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let data2 = txn2.get_page(&cx, p1).unwrap();
assert_eq!(data2.as_ref()[0], 0x33);
}
#[test]
fn test_savepoint_rollback_reclaims_allocated_pages() {
let (pager, _) = test_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap(); txn.write_page(&cx, p1, &vec![0x11; ps]).unwrap();
txn.savepoint(&cx, "sp1").unwrap();
let p2 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p2, &vec![0x22; ps]).unwrap();
assert_eq!(p2.get(), p1.get() + 1, "Expected sequential allocation");
txn.rollback_to_savepoint(&cx, "sp1").unwrap();
let p3 = txn.allocate_page(&cx).unwrap();
assert_eq!(
p3.get(),
p2.get(),
"bead_id={BEAD_ID} case=rollback_reclaims_allocation: expected page {} but got {}",
p2.get(),
p3.get()
);
txn.commit(&cx).unwrap();
}
#[test]
fn test_savepoint_rollback_to_preserves_savepoint_multi() {
let (pager, _) = test_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0x11; ps]).unwrap();
txn.savepoint(&cx, "sp1").unwrap();
txn.write_page(&cx, p1, &vec![0x22; ps]).unwrap();
txn.rollback_to_savepoint(&cx, "sp1").unwrap();
txn.write_page(&cx, p1, &vec![0x33; ps]).unwrap();
txn.rollback_to_savepoint(&cx, "sp1").unwrap();
let data = txn.get_page(&cx, p1).unwrap();
assert_eq!(
data.as_ref()[0],
0x11,
"bead_id={BEAD_ID} case=rollback_to_preserves_savepoint"
);
txn.commit(&cx).unwrap();
}
#[test]
fn test_savepoint_freed_pages_restored_on_rollback() {
let (pager, _) = test_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
let p2 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0xAA; ps]).unwrap();
txn.write_page(&cx, p2, &vec![0xBB; ps]).unwrap();
txn.savepoint(&cx, "sp1").unwrap();
txn.free_page(&cx, p2).unwrap();
txn.rollback_to_savepoint(&cx, "sp1").unwrap();
let data = txn.get_page(&cx, p2).unwrap();
assert_eq!(
data.as_ref()[0],
0xBB,
"bead_id={BEAD_ID} case=freed_pages_restored"
);
txn.commit(&cx).unwrap();
}
#[test]
fn test_savepoint_unknown_name_errors() {
let (pager, _) = test_pager();
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let result = txn.rollback_to_savepoint(&cx, "nonexistent");
assert!(
result.is_err(),
"bead_id={BEAD_ID} case=rollback_to_unknown_savepoint_errors"
);
let result = txn.release_savepoint(&cx, "nonexistent");
assert!(
result.is_err(),
"bead_id={BEAD_ID} case=release_unknown_savepoint_errors"
);
}
#[test]
fn test_savepoint_release_then_rollback_to_outer() {
let (pager, _) = test_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0x11; ps]).unwrap();
txn.savepoint(&cx, "outer").unwrap();
txn.write_page(&cx, p1, &vec![0x22; ps]).unwrap();
txn.savepoint(&cx, "inner").unwrap();
txn.write_page(&cx, p1, &vec![0x33; ps]).unwrap();
txn.release_savepoint(&cx, "inner").unwrap();
txn.rollback_to_savepoint(&cx, "outer").unwrap();
let data = txn.get_page(&cx, p1).unwrap();
assert_eq!(
data.as_ref()[0],
0x11,
"bead_id={BEAD_ID} case=release_inner_then_rollback_outer"
);
txn.commit(&cx).unwrap();
}
#[test]
fn test_savepoint_commit_with_active_savepoints() {
let (pager, _) = test_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0x11; ps]).unwrap();
txn.savepoint(&cx, "sp1").unwrap();
txn.write_page(&cx, p1, &vec![0x22; ps]).unwrap();
txn.commit(&cx).unwrap();
let txn2 = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let data = txn2.get_page(&cx, p1).unwrap();
assert_eq!(
data.as_ref()[0],
0x22,
"bead_id={BEAD_ID} case=commit_with_savepoints_persists_all"
);
}
#[test]
fn test_savepoint_full_rollback_clears_savepoints() {
let (pager, _) = test_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0x11; ps]).unwrap();
txn.savepoint(&cx, "sp1").unwrap();
txn.savepoint(&cx, "sp2").unwrap();
txn.rollback(&cx).unwrap();
let result = txn.rollback_to_savepoint(&cx, "sp1");
assert!(
result.is_err(),
"bead_id={BEAD_ID} case=full_rollback_clears_savepoints"
);
}
#[test]
fn test_savepoint_three_levels_deep() {
let (pager, _) = test_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0x00; ps]).unwrap();
txn.savepoint(&cx, "L0").unwrap();
txn.write_page(&cx, p1, &vec![0x11; ps]).unwrap();
txn.savepoint(&cx, "L1").unwrap();
txn.write_page(&cx, p1, &vec![0x22; ps]).unwrap();
txn.savepoint(&cx, "L2").unwrap();
txn.write_page(&cx, p1, &vec![0x33; ps]).unwrap();
let data = txn.get_page(&cx, p1).unwrap();
assert_eq!(
data.as_ref()[0],
0x33,
"bead_id={BEAD_ID} case=3level_current"
);
txn.rollback_to_savepoint(&cx, "L2").unwrap();
let data = txn.get_page(&cx, p1).unwrap();
assert_eq!(data.as_ref()[0], 0x22, "bead_id={BEAD_ID} case=3level_L2");
txn.rollback_to_savepoint(&cx, "L1").unwrap();
let data = txn.get_page(&cx, p1).unwrap();
assert_eq!(data.as_ref()[0], 0x11, "bead_id={BEAD_ID} case=3level_L1");
txn.rollback_to_savepoint(&cx, "L0").unwrap();
let data = txn.get_page(&cx, p1).unwrap();
assert_eq!(data.as_ref()[0], 0x00, "bead_id={BEAD_ID} case=3level_L0");
txn.commit(&cx).unwrap();
let txn2 = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let data = txn2.get_page(&cx, p1).unwrap();
assert_eq!(
data.as_ref()[0],
0x00,
"bead_id={BEAD_ID} case=3level_committed"
);
}
use std::sync::{Arc as StdArc, Mutex as StdMutex};
type WalFrame = (u32, Vec<u8>, u32);
type SharedFrames = StdArc<StdMutex<Vec<WalFrame>>>;
struct MockWalBackend {
frames: SharedFrames,
}
impl MockWalBackend {
fn new() -> (Self, SharedFrames) {
let frames: SharedFrames = StdArc::new(StdMutex::new(Vec::new()));
(
Self {
frames: StdArc::clone(&frames),
},
frames,
)
}
}
impl crate::traits::WalBackend for MockWalBackend {
fn append_frame(
&mut self,
_cx: &Cx,
page_number: u32,
page_data: &[u8],
db_size_if_commit: u32,
) -> fsqlite_error::Result<()> {
self.frames
.lock()
.unwrap()
.push((page_number, page_data.to_vec(), db_size_if_commit));
Ok(())
}
fn read_page(
&mut self,
_cx: &Cx,
page_number: u32,
) -> fsqlite_error::Result<Option<Vec<u8>>> {
let frames = self.frames.lock().unwrap();
let result = frames
.iter()
.rev()
.find(|(pn, _, _)| *pn == page_number)
.map(|(_, data, _)| data.clone());
drop(frames);
Ok(result)
}
fn sync(&mut self, _cx: &Cx) -> fsqlite_error::Result<()> {
Ok(())
}
fn frame_count(&self) -> usize {
self.frames.lock().unwrap().len()
}
fn checkpoint(
&mut self,
_cx: &Cx,
_mode: crate::traits::CheckpointMode,
_writer: &mut dyn crate::traits::CheckpointPageWriter,
_backfilled_frames: u32,
_oldest_reader_frame: Option<u32>,
) -> fsqlite_error::Result<crate::traits::CheckpointResult> {
let total_frames = u32::try_from(self.frames.lock().unwrap().len()).map_err(|_| {
fsqlite_error::FrankenError::internal("mock wal frame count exceeds u32")
})?;
Ok(crate::traits::CheckpointResult {
total_frames,
frames_backfilled: total_frames,
completed: true,
wal_was_reset: false,
})
}
}
fn wal_pager() -> (SimplePager<MemoryVfs>, SharedFrames) {
let vfs = MemoryVfs::new();
let path = PathBuf::from("/wal_test.db");
let pager = SimplePager::open(vfs, &path, PageSize::DEFAULT).unwrap();
let cx = Cx::new();
let (backend, frames) = MockWalBackend::new();
pager.set_wal_backend(Box::new(backend)).unwrap();
pager.set_journal_mode(&cx, JournalMode::Wal).unwrap();
(pager, frames)
}
#[test]
fn test_journal_mode_default_is_delete() {
let (pager, _) = test_pager();
assert_eq!(
pager.journal_mode(),
JournalMode::Delete,
"bead_id={BEAD_ID} case=default_journal_mode"
);
}
#[test]
fn test_set_journal_mode_wal_requires_backend() {
let (pager, _) = test_pager();
let cx = Cx::new();
let result = pager.set_journal_mode(&cx, JournalMode::Wal);
assert!(
result.is_err(),
"bead_id={BEAD_ID} case=wal_requires_backend"
);
}
#[test]
fn test_set_journal_mode_wal_with_backend() {
let (pager, _) = test_pager();
let cx = Cx::new();
let (backend, _frames) = MockWalBackend::new();
pager.set_wal_backend(Box::new(backend)).unwrap();
let mode = pager.set_journal_mode(&cx, JournalMode::Wal).unwrap();
assert_eq!(
mode,
JournalMode::Wal,
"bead_id={BEAD_ID} case=wal_mode_set"
);
assert_eq!(
pager.journal_mode(),
JournalMode::Wal,
"bead_id={BEAD_ID} case=wal_mode_persisted"
);
}
#[test]
fn test_set_journal_mode_blocked_during_write() {
let (pager, _) = test_pager();
let cx = Cx::new();
let _writer = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let (backend, _frames) = MockWalBackend::new();
pager.set_wal_backend(Box::new(backend)).unwrap();
let result = pager.set_journal_mode(&cx, JournalMode::Wal);
assert!(
result.is_err(),
"bead_id={BEAD_ID} case=mode_switch_blocked_during_write"
);
}
#[test]
fn test_checkpoint_busy_with_active_reader() {
let (pager, _frames) = wal_pager();
let cx = Cx::new();
let reader = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let err = pager
.checkpoint(&cx, crate::traits::CheckpointMode::Passive)
.expect_err("checkpoint should be blocked by active reader");
assert!(matches!(err, FrankenError::Busy));
drop(reader);
let result = pager
.checkpoint(&cx, crate::traits::CheckpointMode::Passive)
.expect("checkpoint should succeed after reader closes");
assert_eq!(result.total_frames, 0);
}
#[test]
fn test_checkpoint_busy_with_active_writer() {
let (pager, _frames) = wal_pager();
let cx = Cx::new();
let _writer = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let err = pager
.checkpoint(&cx, crate::traits::CheckpointMode::Passive)
.expect_err("checkpoint should be blocked by active writer");
assert!(matches!(err, FrankenError::Busy));
}
#[test]
fn test_wal_commit_appends_frames() {
let (pager, frames) = wal_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
let data = vec![0xAA_u8; ps];
txn.write_page(&cx, p1, &data).unwrap();
txn.commit(&cx).unwrap();
let locked_frames = frames.lock().unwrap();
assert_eq!(
locked_frames.len(),
1,
"bead_id={BEAD_ID} case=wal_one_frame_appended"
);
assert_eq!(
locked_frames[0].0,
p1.get(),
"bead_id={BEAD_ID} case=wal_frame_page_number"
);
assert_eq!(
locked_frames[0].1[0], 0xAA,
"bead_id={BEAD_ID} case=wal_frame_data"
);
assert!(
locked_frames[0].2 > 0,
"bead_id={BEAD_ID} case=wal_commit_marker"
);
drop(locked_frames);
}
#[test]
fn test_wal_commit_multi_page() {
let (pager, frames) = wal_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
let p2 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0x11; ps]).unwrap();
txn.write_page(&cx, p2, &vec![0x22; ps]).unwrap();
txn.commit(&cx).unwrap();
let locked_frames = frames.lock().unwrap();
assert_eq!(
locked_frames.len(),
2,
"bead_id={BEAD_ID} case=wal_multi_page_count"
);
let commit_count = locked_frames.iter().filter(|f| f.2 > 0).count();
drop(locked_frames);
assert_eq!(
commit_count, 1,
"bead_id={BEAD_ID} case=wal_exactly_one_commit_marker"
);
}
#[test]
fn test_wal_read_page_from_wal() {
let (pager, _frames) = wal_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
let data = vec![0xBB_u8; ps];
txn.write_page(&cx, p1, &data).unwrap();
txn.commit(&cx).unwrap();
let txn2 = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let read_back = txn2.get_page(&cx, p1).unwrap();
assert_eq!(
read_back.as_ref()[0],
0xBB,
"bead_id={BEAD_ID} case=wal_read_back_from_wal"
);
}
#[test]
fn test_wal_no_journal_file_created() {
let vfs = MemoryVfs::new();
let path = PathBuf::from("/wal_no_jrnl.db");
let pager = SimplePager::open(vfs.clone(), &path, PageSize::DEFAULT).unwrap();
let cx = Cx::new();
let (backend, _frames) = MockWalBackend::new();
pager.set_wal_backend(Box::new(backend)).unwrap();
pager.set_journal_mode(&cx, JournalMode::Wal).unwrap();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0xFF; 4096]).unwrap();
txn.commit(&cx).unwrap();
let journal_path = SimplePager::<MemoryVfs>::journal_path(&path);
assert!(
!vfs.access(&cx, &journal_path, AccessFlags::EXISTS).unwrap(),
"bead_id={BEAD_ID} case=wal_no_journal_created"
);
}
#[test]
fn test_wal_mode_switch_back_to_delete() {
let (pager, _frames) = wal_pager();
let cx = Cx::new();
assert_eq!(pager.journal_mode(), JournalMode::Wal);
let mode = pager.set_journal_mode(&cx, JournalMode::Delete).unwrap();
assert_eq!(
mode,
JournalMode::Delete,
"bead_id={BEAD_ID} case=switch_back_to_delete"
);
}
#[test]
fn test_wal_overwrite_page_reads_latest() {
let (pager, _frames) = wal_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0x11; ps]).unwrap();
txn.commit(&cx).unwrap();
let mut txn2 = pager.begin(&cx, TransactionMode::Immediate).unwrap();
txn2.write_page(&cx, p1, &vec![0x22; ps]).unwrap();
txn2.commit(&cx).unwrap();
let txn3 = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let data = txn3.get_page(&cx, p1).unwrap();
assert_eq!(
data.as_ref()[0],
0x22,
"bead_id={BEAD_ID} case=wal_latest_version"
);
}
#[test]
fn test_wal_rollback_does_not_append_frames() {
let (pager, frames) = wal_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = txn.allocate_page(&cx).unwrap();
txn.write_page(&cx, p1, &vec![0xDD; ps]).unwrap();
txn.rollback(&cx).unwrap();
assert_eq!(
frames.lock().unwrap().len(),
0,
"bead_id={BEAD_ID} case=wal_rollback_no_frames"
);
}
const BEAD_5A1: &str = "bd-2yy6";
#[test]
fn test_page1_database_header_all_fields() {
let (pager, _) = test_pager();
let cx = Cx::new();
let txn = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let raw = txn.get_page(&cx, PageNumber::ONE).unwrap().into_vec();
eprintln!(
"[5A1][test=page1_database_header_all_fields][step=parse] page_len={}",
raw.len()
);
let hdr_bytes: [u8; DATABASE_HEADER_SIZE] = raw[..DATABASE_HEADER_SIZE]
.try_into()
.expect("page 1 must have 100-byte header");
let hdr = DatabaseHeader::from_bytes(&hdr_bytes).expect("header must parse");
assert_eq!(
hdr.page_size,
PageSize::DEFAULT,
"bead_id={BEAD_5A1} case=page_size"
);
assert_eq!(hdr.page_count, 1, "bead_id={BEAD_5A1} case=page_count");
assert_eq!(
hdr.sqlite_version, FRANKENSQLITE_SQLITE_VERSION_NUMBER,
"bead_id={BEAD_5A1} case=sqlite_version"
);
assert_eq!(
hdr.schema_format, 4,
"bead_id={BEAD_5A1} case=schema_format"
);
assert_eq!(
hdr.freelist_trunk, 0,
"bead_id={BEAD_5A1} case=freelist_trunk"
);
assert_eq!(
hdr.freelist_count, 0,
"bead_id={BEAD_5A1} case=freelist_count"
);
assert_eq!(
hdr.schema_cookie, 0,
"bead_id={BEAD_5A1} case=schema_cookie"
);
assert_eq!(
hdr.text_encoding,
fsqlite_types::TextEncoding::Utf8,
"bead_id={BEAD_5A1} case=text_encoding"
);
assert_eq!(hdr.user_version, 0, "bead_id={BEAD_5A1} case=user_version");
assert_eq!(
hdr.application_id, 0,
"bead_id={BEAD_5A1} case=application_id"
);
assert_eq!(
hdr.change_counter, 0,
"bead_id={BEAD_5A1} case=change_counter"
);
assert_eq!(
&raw[..16],
b"SQLite format 3\0",
"bead_id={BEAD_5A1} case=magic_string"
);
assert_eq!(raw[21], 64, "bead_id={BEAD_5A1} case=max_payload_fraction");
assert_eq!(raw[22], 32, "bead_id={BEAD_5A1} case=min_payload_fraction");
assert_eq!(raw[23], 32, "bead_id={BEAD_5A1} case=leaf_payload_fraction");
eprintln!(
"[5A1][test=page1_database_header_all_fields][step=verify] \
page_size={} page_count={} schema_format={} encoding=UTF8 \u{2713}",
hdr.page_size.get(),
hdr.page_count,
hdr.schema_format
);
}
#[test]
fn test_page1_btree_header_is_valid_empty_leaf_table() {
let (pager, _) = test_pager();
let cx = Cx::new();
let txn = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let raw = txn.get_page(&cx, PageNumber::ONE).unwrap().into_vec();
let btree = BTreePageHeader::parse(&raw, PageSize::DEFAULT, 0, true)
.expect("page 1 must parse as B-tree page");
eprintln!(
"[5A1][test=page1_btree_header][step=parse] \
page_type={:?} cell_count={} content_start={} freeblock={} frag={}",
btree.page_type,
btree.cell_count,
btree.cell_content_start,
btree.first_freeblock,
btree.fragmented_free_bytes
);
assert_eq!(
btree.page_type,
fsqlite_types::BTreePageType::LeafTable,
"bead_id={BEAD_5A1} case=btree_page_type"
);
assert_eq!(
btree.cell_count, 0,
"bead_id={BEAD_5A1} case=btree_cell_count"
);
assert_eq!(
btree.cell_content_start,
PageSize::DEFAULT.get(),
"bead_id={BEAD_5A1} case=btree_content_start"
);
assert_eq!(
btree.first_freeblock, 0,
"bead_id={BEAD_5A1} case=btree_first_freeblock"
);
assert_eq!(
btree.fragmented_free_bytes, 0,
"bead_id={BEAD_5A1} case=btree_fragmented_free"
);
assert_eq!(
btree.header_offset, DATABASE_HEADER_SIZE,
"bead_id={BEAD_5A1} case=btree_header_offset"
);
assert!(
btree.right_most_child.is_none(),
"bead_id={BEAD_5A1} case=leaf_no_child"
);
eprintln!("[5A1][test=page1_btree_header][step=verify] empty_leaf_table valid \u{2713}");
}
#[test]
fn test_page1_rest_is_zeroed() {
let (pager, _) = test_pager();
let cx = Cx::new();
let txn = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let raw = txn.get_page(&cx, PageNumber::ONE).unwrap().into_vec();
let btree_header_end = DATABASE_HEADER_SIZE + 8;
let trailing = &raw[btree_header_end..];
let non_zero_count = trailing.iter().filter(|&&b| b != 0).count();
assert_eq!(
non_zero_count, 0,
"bead_id={BEAD_5A1} case=trailing_bytes_zeroed non_zero_count={non_zero_count}"
);
eprintln!(
"[5A1][test=page1_rest_is_zeroed][step=verify] \
trailing_bytes={} all_zero=true \u{2713}",
trailing.len()
);
}
#[test]
fn test_page1_various_page_sizes() {
for &ps_val in &[512u32, 1024, 2048, 4096, 8192, 16384, 32768, 65536] {
let page_size = PageSize::new(ps_val).unwrap();
let vfs = MemoryVfs::new();
let path = PathBuf::from(format!("/test_{ps_val}.db"));
let pager = SimplePager::open(vfs, &path, page_size).unwrap();
let cx = Cx::new();
let txn = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let raw = txn.get_page(&cx, PageNumber::ONE).unwrap().into_vec();
eprintln!(
"[5A1][test=page1_various_page_sizes][step=open] page_size={ps_val} page_len={}",
raw.len()
);
assert_eq!(
raw.len(),
ps_val as usize,
"bead_id={BEAD_5A1} case=page_len ps={ps_val}"
);
let hdr_bytes: [u8; DATABASE_HEADER_SIZE] =
raw[..DATABASE_HEADER_SIZE].try_into().unwrap();
let hdr = DatabaseHeader::from_bytes(&hdr_bytes).unwrap_or_else(|e| {
panic!("bead_id={BEAD_5A1} case=hdr_parse ps={ps_val} err={e}")
});
assert_eq!(
hdr.page_size, page_size,
"bead_id={BEAD_5A1} case=hdr_page_size ps={ps_val}"
);
let btree = BTreePageHeader::parse(&raw, page_size, 0, true).unwrap_or_else(|e| {
panic!("bead_id={BEAD_5A1} case=btree_parse ps={ps_val} err={e}")
});
assert_eq!(
btree.cell_count, 0,
"bead_id={BEAD_5A1} case=empty_cells ps={ps_val}"
);
let expected_content = ps_val;
assert_eq!(
btree.cell_content_start, expected_content,
"bead_id={BEAD_5A1} case=content_start ps={ps_val}"
);
eprintln!(
"[5A1][test=page1_various_page_sizes][step=verify] \
page_size={ps_val} content_start={} \u{2713}",
btree.cell_content_start
);
}
}
#[test]
fn test_write_empty_leaf_table_roundtrip() {
let page_size = PageSize::DEFAULT;
let mut page = vec![0u8; page_size.as_usize()];
BTreePageHeader::write_empty_leaf_table(&mut page, 0, page_size.get());
let parsed = BTreePageHeader::parse(&page, page_size, 0, false)
.expect("bead_id=bd-2yy6 written page must parse");
assert_eq!(parsed.page_type, fsqlite_types::BTreePageType::LeafTable);
assert_eq!(parsed.cell_count, 0);
assert_eq!(parsed.first_freeblock, 0);
assert_eq!(parsed.fragmented_free_bytes, 0);
assert_eq!(parsed.cell_content_start, page_size.get());
assert_eq!(parsed.header_offset, 0);
eprintln!(
"[5A1][test=write_empty_leaf_roundtrip][step=verify] \
non_page1 roundtrip \u{2713}"
);
let mut page1 = vec![0u8; page_size.as_usize()];
BTreePageHeader::write_empty_leaf_table(&mut page1, DATABASE_HEADER_SIZE, page_size.get());
let hdr = DatabaseHeader {
page_size,
page_count: 1,
sqlite_version: FRANKENSQLITE_SQLITE_VERSION_NUMBER,
..DatabaseHeader::default()
};
let hdr_bytes = hdr.to_bytes().unwrap();
page1[..DATABASE_HEADER_SIZE].copy_from_slice(&hdr_bytes);
let parsed1 = BTreePageHeader::parse(&page1, page_size, 0, true)
.expect("bead_id=bd-2yy6 page1 written page must parse");
assert_eq!(parsed1.page_type, fsqlite_types::BTreePageType::LeafTable);
assert_eq!(parsed1.cell_count, 0);
assert_eq!(parsed1.header_offset, DATABASE_HEADER_SIZE);
eprintln!(
"[5A1][test=write_empty_leaf_roundtrip][step=verify] \
page1 roundtrip \u{2713}"
);
}
#[test]
fn test_write_empty_leaf_table_65536_page_size() {
let page_size = PageSize::new(65536).unwrap();
let mut page = vec![0u8; page_size.as_usize()];
BTreePageHeader::write_empty_leaf_table(&mut page, 0, page_size.get());
assert_eq!(page[5], 0x00);
assert_eq!(page[6], 0x00);
let parsed =
BTreePageHeader::parse(&page, page_size, 0, false).expect("65536 page must parse");
assert_eq!(parsed.cell_content_start, 65536);
eprintln!(
"[5A1][test=write_empty_leaf_65536][step=verify] \
content_start=65536 encoding=0x0000 \u{2713}"
);
}
#[test]
fn test_freelist_leak_on_rollback() {
let (pager, _) = test_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn1 = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p = txn1.allocate_page(&cx).unwrap();
txn1.write_page(&cx, p, &vec![0xAA; ps]).unwrap();
txn1.commit(&cx).unwrap();
let mut txn2 = pager.begin(&cx, TransactionMode::Immediate).unwrap();
txn2.free_page(&cx, p).unwrap();
txn2.commit(&cx).unwrap();
{
let inner = pager.inner.lock().unwrap();
assert_eq!(inner.freelist.len(), 1);
assert_eq!(inner.freelist[0], p);
drop(inner);
}
let mut txn3 = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p2 = txn3.allocate_page(&cx).unwrap();
assert_eq!(p2, p, "should reuse freed page");
{
let inner = pager.inner.lock().unwrap();
assert!(inner.freelist.is_empty());
drop(inner);
}
txn3.rollback(&cx).unwrap();
{
let inner = pager.inner.lock().unwrap();
assert_eq!(
inner.freelist.len(),
1,
"bead_id={BEAD_ID} case=freelist_leak_on_rollback"
);
assert_eq!(inner.freelist[0], p);
drop(inner);
}
}
#[test]
#[allow(clippy::similar_names, clippy::cast_possible_truncation)]
fn test_cache_eviction_under_pressure() {
let (pager, _) = test_pager();
let cx = Cx::new();
let ps = PageSize::DEFAULT.as_usize();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let mut pages = Vec::new();
for i in 0..300u32 {
let p = txn.allocate_page(&cx).unwrap();
pages.push(p);
let byte = (i % 256) as u8;
let data = vec![byte; ps];
txn.write_page(&cx, p, &data).unwrap();
}
txn.commit(&cx).unwrap();
let txn = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
for (i, &p) in pages.iter().enumerate() {
let data = txn.get_page(&cx, p).unwrap();
let expected_byte = (i % 256) as u8;
assert_eq!(
data.as_ref()[0],
expected_byte,
"bead_id={BEAD_ID} case=cache_pressure page={p}"
);
}
}
}