use std::collections::HashMap;
use crate::pager::SimpleTransaction;
use fsqlite_error::{FrankenError, Result};
use fsqlite_types::cx::Cx;
use fsqlite_types::{PageData, PageNumber, PageSize};
#[cfg(all(feature = "native", target_os = "linux"))]
use fsqlite_vfs::IoUringVfs;
use fsqlite_vfs::MemoryVfs;
#[cfg(all(feature = "native", unix))]
use fsqlite_vfs::UnixVfs;
#[cfg(all(feature = "native", target_os = "windows"))]
use fsqlite_vfs::WindowsVfs;
use fsqlite_wal::{
TransactionConflictSnapshot, WalGenerationIdentity, checksum::WalChecksumTransform,
};
pub(crate) mod sealed {
pub trait Sealed {}
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
pub enum JournalMode {
#[default]
Delete,
Wal,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub enum CheckpointMode {
#[default]
Passive,
Full,
Restart,
Truncate,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct CheckpointResult {
pub total_frames: u32,
pub frames_backfilled: u32,
pub completed: bool,
pub wal_was_reset: bool,
pub requested_mode: CheckpointMode,
pub effective_mode: CheckpointMode,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct WalPublicationSnapshot {
pub publication_seq: u64,
pub generation: WalGenerationIdentity,
pub last_commit_frame: Option<usize>,
pub commit_count: u64,
pub latest_frame_entries: usize,
pub index_is_partial: bool,
}
impl WalPublicationSnapshot {
#[must_use]
pub const fn lookup_contract_is_authoritative(self) -> bool {
!self.index_is_partial
}
}
pub trait WalBackend: Send + Sync {
fn begin_transaction(&mut self, _cx: &Cx) -> Result<()> {
Ok(())
}
#[must_use]
fn published_snapshot(&self) -> Option<WalPublicationSnapshot> {
None
}
#[must_use]
fn pinned_read_snapshot(&self) -> Option<WalPublicationSnapshot> {
None
}
fn refresh_published_snapshot(&mut self, _cx: &Cx) -> Result<Option<WalPublicationSnapshot>> {
Ok(self.published_snapshot())
}
fn append_frame(
&mut self,
cx: &Cx,
page_number: u32,
page_data: &[u8],
db_size_if_commit: u32,
) -> Result<()>;
fn append_frames(&mut self, cx: &Cx, frames: &[WalFrameRef<'_>]) -> Result<()> {
for frame in frames {
self.append_frame(
cx,
frame.page_number,
frame.page_data,
frame.db_size_if_commit,
)?;
}
Ok(())
}
fn prepare_append_frames(
&self,
_frames: &[WalFrameRef<'_>],
) -> Result<Option<PreparedWalFrameBatch>> {
Ok(None)
}
fn finalize_prepared_frames(
&self,
_cx: &Cx,
_prepared: &mut PreparedWalFrameBatch,
) -> Result<()> {
Ok(())
}
fn append_prepared_frames(
&mut self,
cx: &Cx,
prepared: &mut PreparedWalFrameBatch,
) -> Result<()> {
let frame_refs = prepared.frame_refs();
self.append_frames(cx, &frame_refs)
}
fn read_page(&mut self, cx: &Cx, page_number: u32) -> Result<Option<Vec<u8>>>;
fn read_page_pinned(&self, _cx: &Cx, _page_number: u32) -> Result<Option<Vec<u8>>> {
Err(FrankenError::internal(
"read_page_pinned not supported by this WalBackend; use read_page",
))
}
fn supports_pinned_reads(&self) -> bool {
false
}
fn committed_txns_since_page(&mut self, _cx: &Cx, _page_number: u32) -> Result<u64> {
Ok(0)
}
fn conflicting_pages_since_snapshot(
&mut self,
_cx: &Cx,
_snapshot: TransactionConflictSnapshot,
_page_numbers: &[u32],
) -> Result<Vec<u32>> {
Ok(Vec::new())
}
fn committed_txn_count(&mut self, _cx: &Cx) -> Result<u64> {
Ok(0)
}
fn sync(&mut self, cx: &Cx) -> Result<()>;
fn frame_count(&self) -> usize;
fn checkpoint(
&mut self,
cx: &Cx,
mode: CheckpointMode,
writer: &mut dyn CheckpointPageWriter,
backfilled_frames: u32,
oldest_reader_frame: Option<u32>,
) -> Result<CheckpointResult>;
}
#[derive(Debug, Clone, Copy)]
pub struct WalFrameRef<'a> {
pub page_number: u32,
pub page_data: &'a [u8],
pub db_size_if_commit: u32,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub struct PreparedWalFrameMeta {
pub page_number: u32,
pub db_size_if_commit: u32,
}
pub type PreparedWalChecksumTransform = WalChecksumTransform;
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct PreparedWalChecksumSeed {
pub s1: u32,
pub s2: u32,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct PreparedWalFinalizationState {
pub checkpoint_seq: u32,
pub salt1: u32,
pub salt2: u32,
pub start_frame_index: usize,
pub seed: PreparedWalChecksumSeed,
}
#[derive(Debug, Clone, PartialEq, Eq)]
pub struct PreparedWalFrameBatch {
pub frame_size: usize,
pub page_data_offset: usize,
pub big_endian_checksum: bool,
pub frame_metas: Vec<PreparedWalFrameMeta>,
pub checksum_transforms: Vec<PreparedWalChecksumTransform>,
pub frame_bytes: Vec<u8>,
pub last_commit_frame_offset: Option<usize>,
pub finalized_for: Option<PreparedWalFinalizationState>,
pub finalized_running_checksum: Option<PreparedWalChecksumSeed>,
}
impl PreparedWalFrameBatch {
#[must_use]
pub fn frame_count(&self) -> usize {
self.frame_metas.len()
}
#[must_use]
pub fn page_size(&self) -> usize {
self.frame_size.saturating_sub(self.page_data_offset)
}
#[must_use]
pub fn frame_refs(&self) -> Vec<WalFrameRef<'_>> {
self.frame_metas
.iter()
.enumerate()
.map(|(index, meta)| {
let frame_start = index * self.frame_size;
let page_start = frame_start + self.page_data_offset;
let page_end = frame_start + self.frame_size;
WalFrameRef {
page_number: meta.page_number,
page_data: &self.frame_bytes[page_start..page_end],
db_size_if_commit: meta.db_size_if_commit,
}
})
.collect()
}
#[must_use]
pub fn page_data(&self, index: usize) -> &[u8] {
let frame_start = index * self.frame_size;
let page_start = frame_start + self.page_data_offset;
let page_end = frame_start + self.frame_size;
&self.frame_bytes[page_start..page_end]
}
#[must_use]
pub fn frame_slice(&self, index: usize) -> &[u8] {
let frame_start = index * self.frame_size;
let frame_end = frame_start + self.frame_size;
&self.frame_bytes[frame_start..frame_end]
}
pub fn set_db_size_if_commit(&mut self, index: usize, db_size_if_commit: u32) {
self.frame_metas[index].db_size_if_commit = db_size_if_commit;
let frame_start = index * self.frame_size;
let db_size_offset = frame_start + 4;
self.frame_bytes[db_size_offset..db_size_offset + 4]
.copy_from_slice(&db_size_if_commit.to_be_bytes());
self.finalized_for = None;
self.finalized_running_checksum = None;
}
pub fn recompute_checksum_transforms(&mut self) -> Result<()> {
let page_size = self.page_size();
self.checksum_transforms = (0..self.frame_count())
.map(|index| {
WalChecksumTransform::for_wal_frame(
self.frame_slice(index),
page_size,
self.big_endian_checksum,
)
})
.collect::<Result<Vec<_>>>()?;
self.finalized_for = None;
self.finalized_running_checksum = None;
Ok(())
}
}
#[derive(Debug, Default, Clone, Copy, PartialEq, Eq)]
pub enum TransactionMode {
#[default]
Deferred,
Immediate,
Exclusive,
Concurrent,
ReadOnly,
}
pub trait MvccPager: sealed::Sealed + Send + Sync {
type Txn: TransactionHandle;
fn begin(&self, cx: &Cx, mode: TransactionMode) -> Result<Self::Txn>;
fn journal_mode(&self) -> JournalMode;
fn is_readonly(&self) -> bool;
fn set_journal_mode(&self, cx: &Cx, mode: JournalMode) -> Result<JournalMode>;
fn set_wal_backend(&self, backend: Box<dyn WalBackend>) -> Result<()>;
}
pub trait TransactionHandle: sealed::Sealed + Send {
fn get_page(&self, cx: &Cx, page_no: PageNumber) -> Result<PageData>;
fn prefetch_page_hint(&self, _cx: &Cx, _page_no: PageNumber) {}
fn write_page(&mut self, cx: &Cx, page_no: PageNumber, data: &[u8]) -> Result<()>;
fn write_page_data(&mut self, cx: &Cx, page_no: PageNumber, data: PageData) -> Result<()> {
self.write_page(cx, page_no, data.as_bytes())
}
fn try_take_staged_page_data(&mut self, _page_no: PageNumber) -> Option<PageData> {
None
}
fn try_mutate_staged_page_data(
&mut self,
_page_no: PageNumber,
_f: &mut dyn FnMut(&mut PageData),
) -> bool {
false
}
fn restore_staged_page_data(
&mut self,
cx: &Cx,
page_no: PageNumber,
data: PageData,
) -> Result<()> {
self.write_page_data(cx, page_no, data)
}
fn allocate_page(&mut self, cx: &Cx) -> Result<PageNumber>;
fn free_page(&mut self, cx: &Cx, page_no: PageNumber) -> Result<()>;
fn commit(&mut self, cx: &Cx) -> Result<()>;
fn commit_and_retain(&mut self, cx: &Cx) -> Result<bool> {
self.commit(cx)?;
Ok(false)
}
fn is_writer(&self) -> bool;
fn has_pending_writes(&self) -> bool;
fn published_visible_commit_seq_hint(&self) -> Option<fsqlite_types::CommitSeq> {
None
}
fn pending_commit_pages(&self) -> Result<Vec<PageNumber>> {
Ok(Vec::new())
}
fn pending_conflict_pages(&self) -> Result<Vec<PageNumber>> {
self.pending_commit_pages()
}
fn pending_conflict_pages_conservative(&self) -> Vec<PageNumber> {
self.write_set_page_numbers()
}
fn write_set_page_numbers(&self) -> Vec<PageNumber> {
Vec::new()
}
fn page_one_in_pending_commit_surface(&self) -> Result<bool> {
Ok(self.pending_commit_pages()?.contains(&PageNumber::ONE))
}
fn page_size(&self) -> PageSize {
PageSize::default()
}
fn allocate_page_requires_page_one_conflict_tracking(&self) -> Result<bool> {
Ok(true)
}
fn free_page_requires_page_one_conflict_tracking(&self, _page_no: PageNumber) -> Result<bool> {
Ok(true)
}
fn write_page_requires_page_one_conflict_tracking(&self, _page_no: PageNumber) -> Result<bool> {
Ok(true)
}
fn rollback(&mut self, cx: &Cx) -> Result<()>;
fn record_write_witness(&mut self, _cx: &Cx, _key: fsqlite_types::WitnessKey) {}
fn savepoint(&mut self, cx: &Cx, name: &str) -> Result<()>;
fn release_savepoint(&mut self, cx: &Cx, name: &str) -> Result<()>;
fn rollback_to_savepoint(&mut self, cx: &Cx, name: &str) -> Result<()>;
}
pub trait CheckpointPageWriter: sealed::Sealed + Send {
fn write_page(&mut self, cx: &Cx, page_no: PageNumber, data: &[u8]) -> Result<()>;
fn truncate(&mut self, cx: &Cx, n_pages: u32) -> Result<()>;
fn sync(&mut self, cx: &Cx) -> Result<()>;
}
#[derive(Debug, Default, Clone, Copy)]
pub struct MockMvccPager;
impl sealed::Sealed for MockMvccPager {}
impl MvccPager for MockMvccPager {
type Txn = MockTransaction;
fn begin(&self, _cx: &Cx, _mode: TransactionMode) -> Result<Self::Txn> {
Ok(MockTransaction {
committed: false,
next_page: 2,
savepoint_names: Vec::new(),
})
}
fn journal_mode(&self) -> JournalMode {
JournalMode::Delete
}
fn is_readonly(&self) -> bool {
false
}
fn set_journal_mode(&self, _cx: &Cx, mode: JournalMode) -> Result<JournalMode> {
Ok(mode)
}
fn set_wal_backend(&self, _backend: Box<dyn WalBackend>) -> Result<()> {
Ok(())
}
}
#[derive(Debug, Clone)]
pub struct MockTransaction {
committed: bool,
next_page: u32,
savepoint_names: Vec<String>,
}
impl sealed::Sealed for MockTransaction {}
impl TransactionHandle for MockTransaction {
fn get_page(&self, _cx: &Cx, page_no: PageNumber) -> Result<PageData> {
let size = fsqlite_types::PageSize::default();
let mut data = PageData::zeroed(size);
data.as_bytes_mut()[..4].copy_from_slice(&page_no.get().to_le_bytes());
Ok(data)
}
fn write_page(&mut self, _cx: &Cx, _page_no: PageNumber, _data: &[u8]) -> Result<()> {
Ok(())
}
fn allocate_page(&mut self, _cx: &Cx) -> Result<PageNumber> {
let page = PageNumber::new(self.next_page)
.expect("mock allocator must always produce non-zero page numbers");
self.next_page += 1;
Ok(page)
}
fn free_page(&mut self, _cx: &Cx, _page_no: PageNumber) -> Result<()> {
Ok(())
}
fn commit(&mut self, _cx: &Cx) -> Result<()> {
self.committed = true;
Ok(())
}
fn is_writer(&self) -> bool {
false
}
fn has_pending_writes(&self) -> bool {
false
}
fn pending_commit_pages(&self) -> Result<Vec<PageNumber>> {
Ok(Vec::new())
}
fn rollback(&mut self, _cx: &Cx) -> Result<()> {
Ok(())
}
fn record_write_witness(&mut self, _cx: &Cx, _key: fsqlite_types::WitnessKey) {}
fn savepoint(&mut self, _cx: &Cx, name: &str) -> Result<()> {
self.savepoint_names.push(name.to_owned());
Ok(())
}
fn release_savepoint(&mut self, _cx: &Cx, name: &str) -> Result<()> {
if let Some(pos) = self.savepoint_names.iter().rposition(|n| n == name) {
self.savepoint_names.truncate(pos);
Ok(())
} else {
Err(fsqlite_error::FrankenError::internal(format!(
"no savepoint named '{name}'"
)))
}
}
fn rollback_to_savepoint(&mut self, _cx: &Cx, name: &str) -> Result<()> {
if let Some(pos) = self.savepoint_names.iter().rposition(|n| n == name) {
self.savepoint_names.truncate(pos + 1);
Ok(())
} else {
Err(fsqlite_error::FrankenError::internal(format!(
"no savepoint named '{name}'"
)))
}
}
}
#[derive(Debug, Default, Clone, Copy)]
pub struct MemoryMockMvccPager;
impl sealed::Sealed for MemoryMockMvccPager {}
impl MvccPager for MemoryMockMvccPager {
type Txn = MemoryMockTransaction;
fn begin(&self, _cx: &Cx, _mode: TransactionMode) -> Result<Self::Txn> {
Ok(MemoryMockTransaction {
committed: false,
next_page: 2,
pages: HashMap::new(),
savepoints: Vec::new(),
})
}
fn journal_mode(&self) -> JournalMode {
JournalMode::Delete
}
fn is_readonly(&self) -> bool {
false
}
fn set_journal_mode(&self, _cx: &Cx, mode: JournalMode) -> Result<JournalMode> {
Ok(mode)
}
fn set_wal_backend(&self, _backend: Box<dyn WalBackend>) -> Result<()> {
Ok(())
}
}
#[derive(Debug, Clone)]
struct MemoryMockSavepoint {
name: String,
next_page: u32,
pages: HashMap<PageNumber, PageData>,
}
#[derive(Debug, Clone)]
pub struct MemoryMockTransaction {
committed: bool,
next_page: u32,
pages: HashMap<PageNumber, PageData>,
savepoints: Vec<MemoryMockSavepoint>,
}
impl sealed::Sealed for MemoryMockTransaction {}
impl TransactionHandle for MemoryMockTransaction {
fn get_page(&self, _cx: &Cx, page_no: PageNumber) -> Result<PageData> {
Ok(self
.pages
.get(&page_no)
.cloned()
.unwrap_or_else(|| PageData::zeroed(fsqlite_types::PageSize::default())))
}
fn write_page(&mut self, _cx: &Cx, page_no: PageNumber, data: &[u8]) -> Result<()> {
self.committed = false;
let page_size = fsqlite_types::PageSize::default().as_usize();
let mut page = vec![0_u8; page_size];
let copy_len = data.len().min(page_size);
page[..copy_len].copy_from_slice(&data[..copy_len]);
self.pages.insert(page_no, PageData::from_vec(page));
Ok(())
}
fn write_page_data(&mut self, _cx: &Cx, page_no: PageNumber, data: PageData) -> Result<()> {
self.committed = false;
let page_size = fsqlite_types::PageSize::default().as_usize();
let mut page = vec![0_u8; page_size];
let copy_len = data.len().min(page_size);
page[..copy_len].copy_from_slice(&data.as_bytes()[..copy_len]);
self.pages.insert(page_no, PageData::from_vec(page));
Ok(())
}
fn allocate_page(&mut self, _cx: &Cx) -> Result<PageNumber> {
self.committed = false;
let page = PageNumber::new(self.next_page)
.expect("mock allocator must always produce non-zero page numbers");
self.next_page += 1;
self.pages
.entry(page)
.or_insert_with(|| PageData::zeroed(fsqlite_types::PageSize::default()));
Ok(page)
}
fn free_page(&mut self, _cx: &Cx, page_no: PageNumber) -> Result<()> {
self.committed = false;
self.pages.remove(&page_no);
Ok(())
}
fn commit(&mut self, _cx: &Cx) -> Result<()> {
self.committed = true;
Ok(())
}
fn is_writer(&self) -> bool {
!self.pages.is_empty()
}
fn has_pending_writes(&self) -> bool {
!self.committed && !self.pages.is_empty()
}
fn pending_commit_pages(&self) -> Result<Vec<PageNumber>> {
let mut pages = self.pages.keys().copied().collect::<Vec<_>>();
pages.sort_unstable();
Ok(pages)
}
fn rollback(&mut self, _cx: &Cx) -> Result<()> {
self.committed = false;
self.next_page = 2;
self.pages.clear();
self.savepoints.clear();
Ok(())
}
fn record_write_witness(&mut self, _cx: &Cx, _key: fsqlite_types::WitnessKey) {}
fn savepoint(&mut self, _cx: &Cx, name: &str) -> Result<()> {
self.savepoints.push(MemoryMockSavepoint {
name: name.to_owned(),
next_page: self.next_page,
pages: self.pages.clone(),
});
Ok(())
}
fn release_savepoint(&mut self, _cx: &Cx, name: &str) -> Result<()> {
if let Some(pos) = self.savepoints.iter().rposition(|sp| sp.name == name) {
self.savepoints.truncate(pos);
Ok(())
} else {
Err(fsqlite_error::FrankenError::internal(format!(
"no savepoint named '{name}'"
)))
}
}
fn rollback_to_savepoint(&mut self, _cx: &Cx, name: &str) -> Result<()> {
if let Some(pos) = self.savepoints.iter().rposition(|sp| sp.name == name) {
let snapshot = self.savepoints[pos].clone();
self.next_page = snapshot.next_page;
self.pages = snapshot.pages;
self.savepoints.truncate(pos + 1);
Ok(())
} else {
Err(fsqlite_error::FrankenError::internal(format!(
"no savepoint named '{name}'"
)))
}
}
}
pub enum TransactionKind {
Memory(SimpleTransaction<MemoryVfs>),
#[cfg(all(feature = "native", target_os = "linux"))]
IoUring(SimpleTransaction<IoUringVfs>),
#[cfg(all(feature = "native", unix))]
Unix(SimpleTransaction<UnixVfs>),
#[cfg(all(feature = "native", target_os = "windows"))]
Windows(SimpleTransaction<WindowsVfs>),
Mock(MockTransaction),
MemoryMock(MemoryMockTransaction),
Drained,
}
impl std::fmt::Debug for TransactionKind {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
match self {
Self::Memory(_) => f.write_str("TransactionKind::Memory"),
#[cfg(all(feature = "native", target_os = "linux"))]
Self::IoUring(_) => f.write_str("TransactionKind::IoUring"),
#[cfg(all(feature = "native", unix))]
Self::Unix(_) => f.write_str("TransactionKind::Unix"),
#[cfg(all(feature = "native", target_os = "windows"))]
Self::Windows(_) => f.write_str("TransactionKind::Windows"),
Self::Mock(_) => f.write_str("TransactionKind::Mock"),
Self::MemoryMock(_) => f.write_str("TransactionKind::MemoryMock"),
Self::Drained => f.write_str("TransactionKind::Drained"),
}
}
}
impl TransactionKind {
fn with_handle<R>(&self, f: impl FnOnce(&dyn TransactionHandle) -> R) -> R {
match self {
Self::Memory(txn) => f(txn),
#[cfg(all(feature = "native", target_os = "linux"))]
Self::IoUring(txn) => f(txn),
#[cfg(all(feature = "native", unix))]
Self::Unix(txn) => f(txn),
#[cfg(all(feature = "native", target_os = "windows"))]
Self::Windows(txn) => f(txn),
Self::Mock(txn) => f(txn),
Self::MemoryMock(txn) => f(txn),
Self::Drained => panic!(
"BUG: TransactionKind::Drained accessed — a retained cursor tried to \
read/write pages while the transaction was extracted. This sentinel \
only exists between engine.drain_transaction() and the next \
engine.refill_transaction(). If you see this, a cursor was accessed \
outside the VDBE execution window."
),
}
}
fn with_handle_mut<R>(&mut self, f: impl FnOnce(&mut dyn TransactionHandle) -> R) -> R {
match self {
Self::Memory(txn) => f(txn),
#[cfg(all(feature = "native", target_os = "linux"))]
Self::IoUring(txn) => f(txn),
#[cfg(all(feature = "native", unix))]
Self::Unix(txn) => f(txn),
#[cfg(all(feature = "native", target_os = "windows"))]
Self::Windows(txn) => f(txn),
Self::Mock(txn) => f(txn),
Self::MemoryMock(txn) => f(txn),
Self::Drained => panic!(
"BUG: TransactionKind::Drained accessed — a retained cursor tried to \
read/write pages while the transaction was extracted. This sentinel \
only exists between engine.drain_transaction() and the next \
engine.refill_transaction(). If you see this, a cursor was accessed \
outside the VDBE execution window."
),
}
}
}
impl From<SimpleTransaction<MemoryVfs>> for TransactionKind {
fn from(txn: SimpleTransaction<MemoryVfs>) -> Self {
Self::Memory(txn)
}
}
#[cfg(all(feature = "native", target_os = "linux"))]
impl From<SimpleTransaction<IoUringVfs>> for TransactionKind {
fn from(txn: SimpleTransaction<IoUringVfs>) -> Self {
Self::IoUring(txn)
}
}
#[cfg(all(feature = "native", unix))]
impl From<SimpleTransaction<UnixVfs>> for TransactionKind {
fn from(txn: SimpleTransaction<UnixVfs>) -> Self {
Self::Unix(txn)
}
}
#[cfg(all(feature = "native", target_os = "windows"))]
impl From<SimpleTransaction<WindowsVfs>> for TransactionKind {
fn from(txn: SimpleTransaction<WindowsVfs>) -> Self {
Self::Windows(txn)
}
}
impl From<MockTransaction> for TransactionKind {
fn from(txn: MockTransaction) -> Self {
Self::Mock(txn)
}
}
impl From<MemoryMockTransaction> for TransactionKind {
fn from(txn: MemoryMockTransaction) -> Self {
Self::MemoryMock(txn)
}
}
impl sealed::Sealed for TransactionKind {}
impl TransactionHandle for TransactionKind {
fn get_page(&self, cx: &Cx, page_no: PageNumber) -> Result<PageData> {
match self {
Self::Memory(txn) => txn.get_page(cx, page_no),
#[cfg(all(feature = "native", target_os = "linux"))]
Self::IoUring(txn) => txn.get_page(cx, page_no),
#[cfg(all(feature = "native", unix))]
Self::Unix(txn) => txn.get_page(cx, page_no),
#[cfg(all(feature = "native", target_os = "windows"))]
Self::Windows(txn) => txn.get_page(cx, page_no),
Self::Mock(txn) => txn.get_page(cx, page_no),
Self::MemoryMock(txn) => txn.get_page(cx, page_no),
Self::Drained => panic!(
"BUG: TransactionKind::Drained accessed in get_page — a retained \
cursor tried to read pages while the transaction was extracted."
),
}
}
fn prefetch_page_hint(&self, cx: &Cx, page_no: PageNumber) {
self.with_handle(|txn| txn.prefetch_page_hint(cx, page_no));
}
fn write_page(&mut self, cx: &Cx, page_no: PageNumber, data: &[u8]) -> Result<()> {
self.with_handle_mut(|txn| txn.write_page(cx, page_no, data))
}
fn write_page_data(&mut self, cx: &Cx, page_no: PageNumber, data: PageData) -> Result<()> {
match self {
Self::Memory(txn) => txn.write_page_data(cx, page_no, data),
#[cfg(all(feature = "native", target_os = "linux"))]
Self::IoUring(txn) => txn.write_page_data(cx, page_no, data),
#[cfg(all(feature = "native", unix))]
Self::Unix(txn) => txn.write_page_data(cx, page_no, data),
#[cfg(all(feature = "native", target_os = "windows"))]
Self::Windows(txn) => txn.write_page_data(cx, page_no, data),
Self::Mock(txn) => txn.write_page_data(cx, page_no, data),
Self::MemoryMock(txn) => txn.write_page_data(cx, page_no, data),
Self::Drained => panic!(
"BUG: TransactionKind::Drained accessed in write_page_data — a \
retained cursor tried to write pages while the transaction was \
extracted."
),
}
}
fn try_mutate_staged_page_data(
&mut self,
page_no: PageNumber,
f: &mut dyn FnMut(&mut PageData),
) -> bool {
self.with_handle_mut(|txn| txn.try_mutate_staged_page_data(page_no, f))
}
fn allocate_page(&mut self, cx: &Cx) -> Result<PageNumber> {
self.with_handle_mut(|txn| txn.allocate_page(cx))
}
fn free_page(&mut self, cx: &Cx, page_no: PageNumber) -> Result<()> {
match self {
Self::Memory(txn) => txn.free_page(cx, page_no),
#[cfg(all(feature = "native", target_os = "linux"))]
Self::IoUring(txn) => txn.free_page(cx, page_no),
#[cfg(all(feature = "native", unix))]
Self::Unix(txn) => txn.free_page(cx, page_no),
#[cfg(all(feature = "native", target_os = "windows"))]
Self::Windows(txn) => txn.free_page(cx, page_no),
Self::Mock(txn) => txn.free_page(cx, page_no),
Self::MemoryMock(txn) => txn.free_page(cx, page_no),
Self::Drained => panic!(
"BUG: TransactionKind::Drained accessed in free_page — a retained \
cursor tried to free pages while the transaction was extracted."
),
}
}
fn commit(&mut self, cx: &Cx) -> Result<()> {
self.with_handle_mut(|txn| txn.commit(cx))
}
fn commit_and_retain(&mut self, cx: &Cx) -> Result<bool> {
self.with_handle_mut(|txn| txn.commit_and_retain(cx))
}
fn is_writer(&self) -> bool {
self.with_handle(|txn| txn.is_writer())
}
fn has_pending_writes(&self) -> bool {
self.with_handle(|txn| txn.has_pending_writes())
}
fn published_visible_commit_seq_hint(&self) -> Option<fsqlite_types::CommitSeq> {
self.with_handle(|txn| txn.published_visible_commit_seq_hint())
}
fn pending_commit_pages(&self) -> Result<Vec<PageNumber>> {
self.with_handle(|txn| txn.pending_commit_pages())
}
fn pending_conflict_pages(&self) -> Result<Vec<PageNumber>> {
self.with_handle(|txn| txn.pending_conflict_pages())
}
fn pending_conflict_pages_conservative(&self) -> Vec<PageNumber> {
self.with_handle(|txn| txn.pending_conflict_pages_conservative())
}
fn write_set_page_numbers(&self) -> Vec<PageNumber> {
self.with_handle(|txn| txn.write_set_page_numbers())
}
fn page_one_in_pending_commit_surface(&self) -> Result<bool> {
self.with_handle(|txn| txn.page_one_in_pending_commit_surface())
}
fn page_size(&self) -> PageSize {
self.with_handle(|txn| txn.page_size())
}
fn allocate_page_requires_page_one_conflict_tracking(&self) -> Result<bool> {
self.with_handle(|txn| txn.allocate_page_requires_page_one_conflict_tracking())
}
fn free_page_requires_page_one_conflict_tracking(&self, page_no: PageNumber) -> Result<bool> {
self.with_handle(|txn| txn.free_page_requires_page_one_conflict_tracking(page_no))
}
fn write_page_requires_page_one_conflict_tracking(&self, page_no: PageNumber) -> Result<bool> {
self.with_handle(|txn| txn.write_page_requires_page_one_conflict_tracking(page_no))
}
fn rollback(&mut self, cx: &Cx) -> Result<()> {
self.with_handle_mut(|txn| txn.rollback(cx))
}
fn record_write_witness(&mut self, cx: &Cx, key: fsqlite_types::WitnessKey) {
self.with_handle_mut(|txn| txn.record_write_witness(cx, key));
}
fn savepoint(&mut self, cx: &Cx, name: &str) -> Result<()> {
self.with_handle_mut(|txn| txn.savepoint(cx, name))
}
fn release_savepoint(&mut self, cx: &Cx, name: &str) -> Result<()> {
self.with_handle_mut(|txn| txn.release_savepoint(cx, name))
}
fn rollback_to_savepoint(&mut self, cx: &Cx, name: &str) -> Result<()> {
self.with_handle_mut(|txn| txn.rollback_to_savepoint(cx, name))
}
}
#[derive(Debug, Default, Clone, Copy)]
pub struct MockCheckpointPageWriter;
impl sealed::Sealed for MockCheckpointPageWriter {}
impl CheckpointPageWriter for MockCheckpointPageWriter {
fn write_page(&mut self, _cx: &Cx, _page_no: PageNumber, _data: &[u8]) -> Result<()> {
Ok(())
}
fn truncate(&mut self, _cx: &Cx, _n_pages: u32) -> Result<()> {
Ok(())
}
fn sync(&mut self, _cx: &Cx) -> Result<()> {
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
const fn test_wal_generation_identity() -> WalGenerationIdentity {
WalGenerationIdentity {
checkpoint_seq: 0,
salts: fsqlite_wal::checksum::WalSalts { salt1: 0, salt2: 0 },
}
}
#[test]
fn test_pager_trait_is_sealed_mock_impl() {
let pager = MockMvccPager;
let cx = Cx::new();
let _txn = pager.begin(&cx, TransactionMode::Deferred).unwrap();
}
#[test]
fn test_mvccpager_begin_commit_rollback_signatures() {
let pager = MockMvccPager;
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::ReadOnly).unwrap();
let page_no = PageNumber::new(1).unwrap();
let data = txn.get_page(&cx, page_no).unwrap();
assert_eq!(
u32::from_le_bytes(data.as_bytes()[..4].try_into().unwrap()),
1
);
txn.write_page(&cx, page_no, &[0u8; 4096]).unwrap();
let new_page = txn.allocate_page(&cx).unwrap();
assert_eq!(new_page.get(), 2);
txn.free_page(&cx, new_page).unwrap();
txn.commit(&cx).unwrap();
}
#[test]
fn test_transaction_rollback_is_infallible() {
let pager = MockMvccPager;
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::Deferred).unwrap();
txn.rollback(&cx).unwrap();
}
#[test]
fn test_checkpoint_page_writer_signatures() {
let mut writer = MockCheckpointPageWriter;
let cx = Cx::new();
let page1 = PageNumber::new(1).unwrap();
writer.write_page(&cx, page1, &[0u8; 4096]).unwrap();
writer.truncate(&cx, 10).unwrap();
writer.sync(&cx).unwrap();
}
#[test]
fn test_transaction_mode_default_is_deferred() {
assert_eq!(TransactionMode::default(), TransactionMode::Deferred);
}
#[test]
fn test_open_traits_are_extensible() {
let pager = MockMvccPager;
let _: &dyn MvccPager<Txn = MockTransaction> = &pager;
}
#[test]
fn test_memory_mock_transaction_persists_writes() {
let pager = MemoryMockMvccPager;
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let page_no = PageNumber::new(256).unwrap();
let mut bytes = vec![0_u8; fsqlite_types::PageSize::default().as_usize()];
bytes[0] = 0x0A;
txn.write_page(&cx, page_no, &bytes).unwrap();
let page = txn.get_page(&cx, page_no).unwrap();
assert_eq!(page.as_bytes()[0], 0x0A);
assert!(txn.has_pending_writes());
assert!(txn.is_writer());
}
#[test]
fn test_memory_mock_transaction_commit_clears_pending_writes() {
let pager = MemoryMockMvccPager;
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let page_no = PageNumber::new(2).unwrap();
txn.write_page(&cx, page_no, &[1_u8; 4096]).unwrap();
assert!(txn.has_pending_writes());
txn.commit(&cx).unwrap();
assert!(
!txn.has_pending_writes(),
"committed mock transactions must not report pending writes"
);
}
#[test]
fn test_memory_mock_transaction_rollback_resets_allocator() {
let pager = MemoryMockMvccPager;
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
assert_eq!(txn.allocate_page(&cx).unwrap().get(), 2);
assert_eq!(txn.allocate_page(&cx).unwrap().get(), 3);
txn.rollback(&cx).unwrap();
assert_eq!(
txn.allocate_page(&cx).unwrap().get(),
2,
"rollback should restore the mock allocator to its initial state"
);
}
#[test]
fn test_checkpoint_mode_default_is_passive() {
assert_eq!(CheckpointMode::default(), CheckpointMode::Passive);
}
#[test]
fn test_journal_mode_default_is_delete() {
assert_eq!(JournalMode::default(), JournalMode::Delete);
}
#[test]
fn test_wal_publication_snapshot_authoritative_when_index_full() {
let snap = WalPublicationSnapshot {
publication_seq: 1,
generation: test_wal_generation_identity(),
last_commit_frame: Some(10),
commit_count: 5,
latest_frame_entries: 10,
index_is_partial: false,
};
assert!(
snap.lookup_contract_is_authoritative(),
"full index must be authoritative"
);
}
#[test]
fn test_wal_publication_snapshot_not_authoritative_when_partial() {
let snap = WalPublicationSnapshot {
publication_seq: 1,
generation: test_wal_generation_identity(),
last_commit_frame: None,
commit_count: 0,
latest_frame_entries: 0,
index_is_partial: true,
};
assert!(
!snap.lookup_contract_is_authoritative(),
"partial index must not be authoritative"
);
}
#[test]
fn test_prepared_wal_frame_batch_frame_count_and_page_size() {
let batch = PreparedWalFrameBatch {
frame_size: 4120,
page_data_offset: 24,
big_endian_checksum: false,
frame_metas: vec![
PreparedWalFrameMeta {
page_number: 1,
db_size_if_commit: 0,
},
PreparedWalFrameMeta {
page_number: 2,
db_size_if_commit: 10,
},
],
checksum_transforms: Vec::new(),
frame_bytes: vec![0u8; 4120 * 2],
last_commit_frame_offset: Some(4120),
finalized_for: None,
finalized_running_checksum: None,
};
assert_eq!(batch.frame_count(), 2);
assert_eq!(batch.page_size(), 4096);
}
#[test]
fn test_prepared_wal_frame_batch_set_db_size_clears_finalized() {
let mut batch = PreparedWalFrameBatch {
frame_size: 32,
page_data_offset: 8,
big_endian_checksum: false,
frame_metas: vec![PreparedWalFrameMeta {
page_number: 1,
db_size_if_commit: 0,
}],
checksum_transforms: Vec::new(),
frame_bytes: vec![0u8; 32],
last_commit_frame_offset: None,
finalized_for: Some(PreparedWalFinalizationState {
checkpoint_seq: 1,
salt1: 0xAA,
salt2: 0xBB,
start_frame_index: 0,
seed: PreparedWalChecksumSeed::default(),
}),
finalized_running_checksum: Some(PreparedWalChecksumSeed { s1: 1, s2: 2 }),
};
batch.set_db_size_if_commit(0, 42);
assert_eq!(batch.frame_metas[0].db_size_if_commit, 42);
assert!(
batch.finalized_for.is_none(),
"set_db_size_if_commit must invalidate finalized_for"
);
assert!(
batch.finalized_running_checksum.is_none(),
"set_db_size_if_commit must invalidate finalized_running_checksum"
);
let db_bytes = &batch.frame_bytes[4..8];
assert_eq!(u32::from_be_bytes(db_bytes.try_into().unwrap()), 42);
}
#[test]
fn test_mock_release_savepoint_unknown_name_returns_error() {
let pager = MockMvccPager;
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::Deferred).unwrap();
let result = txn.release_savepoint(&cx, "nonexistent");
assert!(result.is_err(), "releasing unknown savepoint must fail");
}
#[test]
fn test_memory_mock_savepoint_rollback_restores_pages() {
let pager = MemoryMockMvccPager;
let cx = Cx::new();
let mut txn = pager.begin(&cx, TransactionMode::Immediate).unwrap();
let p1 = PageNumber::new(1).unwrap();
let page_size = fsqlite_types::PageSize::default().as_usize();
let mut data_a = vec![0u8; page_size];
data_a[0] = 0xAA;
txn.write_page(&cx, p1, &data_a).unwrap();
txn.savepoint(&cx, "sp1").unwrap();
let mut data_b = vec![0u8; page_size];
data_b[0] = 0xBB;
txn.write_page(&cx, p1, &data_b).unwrap();
assert_eq!(txn.get_page(&cx, p1).unwrap().as_bytes()[0], 0xBB);
txn.rollback_to_savepoint(&cx, "sp1").unwrap();
assert_eq!(
txn.get_page(&cx, p1).unwrap().as_bytes()[0],
0xAA,
"rollback_to_savepoint must restore page state"
);
}
#[test]
fn test_transaction_mode_default_trait_contract_is_deferred() {
assert_eq!(TransactionMode::default(), TransactionMode::Deferred);
}
#[test]
fn test_checkpoint_result_fields() {
let result = CheckpointResult {
total_frames: 100,
frames_backfilled: 80,
completed: false,
wal_was_reset: false,
requested_mode: CheckpointMode::Full,
effective_mode: CheckpointMode::Passive,
};
assert_eq!(result.total_frames, 100);
assert_eq!(result.frames_backfilled, 80);
assert!(!result.completed);
assert_ne!(result.requested_mode, result.effective_mode);
}
#[test]
fn test_journal_mode_debug_clone_copy_eq() {
let a = JournalMode::Wal;
let b = a;
assert_eq!(a, b);
assert_ne!(JournalMode::Delete, JournalMode::Wal);
let dbg = format!("{a:?}");
assert!(dbg.contains("Wal"));
}
#[test]
fn test_checkpoint_result_clone_debug() {
let result = CheckpointResult {
total_frames: 50,
frames_backfilled: 50,
completed: true,
wal_was_reset: true,
requested_mode: CheckpointMode::Truncate,
effective_mode: CheckpointMode::Truncate,
};
let cloned = result.clone();
assert_eq!(result, cloned);
let dbg = format!("{result:?}");
assert!(dbg.contains("CheckpointResult"));
assert!(dbg.contains("Truncate"));
assert!(dbg.contains("wal_was_reset"));
}
#[test]
fn test_wal_publication_snapshot_clone_copy_debug() {
let snap = WalPublicationSnapshot {
publication_seq: 42,
generation: test_wal_generation_identity(),
last_commit_frame: Some(100),
commit_count: 7,
latest_frame_entries: 50,
index_is_partial: false,
};
let copied = snap;
assert_eq!(copied, snap);
let dbg = format!("{snap:?}");
assert!(dbg.contains("WalPublicationSnapshot"));
assert!(dbg.contains("publication_seq"));
assert!(dbg.contains("42"));
}
#[test]
fn test_checkpoint_mode_all_variants_debug() {
for (mode, expected) in [
(CheckpointMode::Passive, "Passive"),
(CheckpointMode::Full, "Full"),
(CheckpointMode::Restart, "Restart"),
(CheckpointMode::Truncate, "Truncate"),
] {
let dbg = format!("{mode:?}");
assert!(dbg.contains(expected), "expected {expected} in {dbg}");
let copy = mode;
assert_eq!(mode, copy);
}
}
#[test]
fn test_prepared_wal_frame_batch_page_data_and_frame_slice() {
let frame_size = 32;
let page_data_offset = 8;
let mut frame_bytes = vec![0u8; frame_size * 2];
frame_bytes[8] = 0xAA;
frame_bytes[frame_size + 8] = 0xBB;
let batch = PreparedWalFrameBatch {
frame_size,
page_data_offset,
big_endian_checksum: false,
frame_metas: vec![
PreparedWalFrameMeta {
page_number: 1,
db_size_if_commit: 0,
},
PreparedWalFrameMeta {
page_number: 2,
db_size_if_commit: 5,
},
],
checksum_transforms: Vec::new(),
frame_bytes,
last_commit_frame_offset: None,
finalized_for: None,
finalized_running_checksum: None,
};
assert_eq!(batch.page_data(0)[0], 0xAA);
assert_eq!(batch.page_data(1)[0], 0xBB);
assert_eq!(batch.frame_slice(0).len(), frame_size);
assert_eq!(batch.frame_slice(1).len(), frame_size);
let refs = batch.frame_refs();
assert_eq!(refs.len(), 2);
assert_eq!(refs[0].page_number, 1);
assert_eq!(refs[1].db_size_if_commit, 5);
assert_eq!(refs[0].page_data[0], 0xAA);
assert_eq!(refs[1].page_data[0], 0xBB);
}
#[test]
fn prepared_wal_frame_meta_debug_clone_copy_eq() {
let a = PreparedWalFrameMeta {
page_number: 5,
db_size_if_commit: 0,
};
let b = PreparedWalFrameMeta {
page_number: 5,
db_size_if_commit: 10,
};
let copied = a;
assert_eq!(copied, a);
assert_ne!(a, b);
let dbg = format!("{a:?}");
assert!(dbg.contains("PreparedWalFrameMeta"));
assert!(dbg.contains("5"));
}
#[test]
fn prepared_wal_checksum_seed_default_and_eq() {
let def = PreparedWalChecksumSeed::default();
assert_eq!(def.s1, 0);
assert_eq!(def.s2, 0);
let other = PreparedWalChecksumSeed { s1: 1, s2: 2 };
assert_ne!(def, other);
let copied = other;
assert_eq!(copied, other);
let dbg = format!("{def:?}");
assert!(dbg.contains("PreparedWalChecksumSeed"));
}
#[test]
fn prepared_wal_finalization_state_default_and_eq() {
let def = PreparedWalFinalizationState::default();
assert_eq!(def.checkpoint_seq, 0);
assert_eq!(def.salt1, 0);
assert_eq!(def.salt2, 0);
assert_eq!(def.start_frame_index, 0);
assert_eq!(def.seed, PreparedWalChecksumSeed::default());
let other = PreparedWalFinalizationState {
checkpoint_seq: 1,
salt1: 0xAA,
salt2: 0xBB,
start_frame_index: 42,
seed: PreparedWalChecksumSeed { s1: 10, s2: 20 },
};
assert_ne!(def, other);
let copied = other;
assert_eq!(copied, other);
let dbg = format!("{other:?}");
assert!(dbg.contains("PreparedWalFinalizationState"));
}
#[test]
fn transaction_mode_all_variants_debug_copy_eq() {
let variants = [
(TransactionMode::Deferred, "Deferred"),
(TransactionMode::Immediate, "Immediate"),
(TransactionMode::Exclusive, "Exclusive"),
(TransactionMode::Concurrent, "Concurrent"),
(TransactionMode::ReadOnly, "ReadOnly"),
];
for (mode, expected) in &variants {
let dbg = format!("{mode:?}");
assert!(dbg.contains(expected), "expected {expected} in {dbg}");
let copied = *mode;
assert_eq!(copied, *mode);
}
assert_ne!(TransactionMode::Deferred, TransactionMode::Concurrent);
}
#[test]
fn wal_frame_ref_debug_clone_copy() {
let data = [0xABu8; 16];
let frame = WalFrameRef {
page_number: 3,
page_data: &data,
db_size_if_commit: 0,
};
let copied = frame;
assert_eq!(copied.page_number, 3);
assert_eq!(copied.page_data.len(), 16);
assert_eq!(copied.db_size_if_commit, 0);
let dbg = format!("{frame:?}");
assert!(dbg.contains("WalFrameRef"));
}
#[test]
fn mock_checkpoint_page_writer_default_and_trait_methods() {
let mut writer = MockCheckpointPageWriter;
let cx = Cx::new();
let page = PageNumber::new(1).unwrap();
writer.write_page(&cx, page, &[0u8; 4096]).unwrap();
writer.truncate(&cx, 10).unwrap();
writer.sync(&cx).unwrap();
let dbg = format!("{writer:?}");
assert!(dbg.contains("MockCheckpointPageWriter"));
}
#[test]
fn transaction_kind_drained_debug() {
let kind = TransactionKind::Drained;
let dbg = format!("{kind:?}");
assert!(dbg.contains("Drained"));
}
#[test]
fn wal_publication_snapshot_authoritative_boundary() {
let base = WalPublicationSnapshot {
publication_seq: 1,
generation: test_wal_generation_identity(),
last_commit_frame: Some(10),
commit_count: 5,
latest_frame_entries: 10,
index_is_partial: false,
};
assert!(base.lookup_contract_is_authoritative());
let partial = WalPublicationSnapshot {
index_is_partial: true,
..base
};
assert!(!partial.lookup_contract_is_authoritative());
}
}