From 06261515f36c9422adb184fd58ba1d0a70d0469a Mon Sep 17 00:00:00 2001 From: Kim Altintop Date: Wed, 23 Sep 2026 13:44:31 +0200 Subject: [PATCH] Base machinery to replace `Pages` with `PageSet` --- Cargo.lock | 2 + crates/sats/src/layout.rs | 9 + crates/table/Cargo.toml | 2 + crates/table/src/lib.rs | 1 + crates/table/src/page.rs | 194 ++++++++-- crates/table/src/tiered/budget.rs | 104 ++++++ crates/table/src/tiered/mod.rs | 8 + crates/table/src/tiered/page_manager.rs | 477 ++++++++++++++++++++++++ crates/table/src/tiered/page_set.rs | 255 +++++++++++++ 9 files changed, 1020 insertions(+), 32 deletions(-) create mode 100644 crates/table/src/tiered/budget.rs create mode 100644 crates/table/src/tiered/mod.rs create mode 100644 crates/table/src/tiered/page_manager.rs create mode 100644 crates/table/src/tiered/page_set.rs diff --git a/Cargo.lock b/Cargo.lock index 6a114150c4d..8a22b997fce 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -8926,9 +8926,11 @@ dependencies = [ "foldhash 0.2.0", "hashbrown 0.16.1", "itertools 0.12.1", + "parking_lot 0.12.5", "proptest", "proptest-derive", "rand 0.9.2", + "slab", "smallvec", "spacetimedb-data-structures", "spacetimedb-lib", diff --git a/crates/sats/src/layout.rs b/crates/sats/src/layout.rs index b9e8f919899..524c278c6b1 100644 --- a/crates/sats/src/layout.rs +++ b/crates/sats/src/layout.rs @@ -60,6 +60,15 @@ impl Size { pub const fn len(self) -> usize { self.0 as usize } + + /// Computes `self - rhs`, returning `None` if underflow occurred. + #[inline] + pub const fn checked_sub(self, rhs: Self) -> Option { + match self.0.checked_sub(rhs.0) { + Some(v) => Some(Self(v)), + None => None, + } + } } impl Mul for Size { diff --git a/crates/table/Cargo.toml b/crates/table/Cargo.toml index 1b5d9684480..55e5af355be 100644 --- a/crates/table/Cargo.toml +++ b/crates/table/Cargo.toml @@ -46,6 +46,8 @@ enum-as-inner.workspace = true foldhash.workspace = true hashbrown.workspace = true itertools.workspace = true +parking_lot.workspace = true +slab.workspace = true smallvec.workspace = true thiserror.workspace = true diff --git a/crates/table/src/lib.rs b/crates/table/src/lib.rs index 9aa10f04f1b..bd27732e4d8 100644 --- a/crates/table/src/lib.rs +++ b/crates/table/src/lib.rs @@ -24,6 +24,7 @@ pub mod static_bsatn_validator; pub mod static_layout; pub mod table; pub mod table_index; +pub mod tiered; pub mod var_len; #[doc(hidden)] // Used in tests and benchmarks. diff --git a/crates/table/src/page.rs b/crates/table/src/page.rs index cfb66d97f10..a705b053b82 100644 --- a/crates/table/src/page.rs +++ b/crates/table/src/page.rs @@ -1346,38 +1346,13 @@ impl Page { /// where the fixed size part is `fixed_row_size` bytes large, /// and the variable part requires `num_granules`. pub fn has_space_for_row(&self, fixed_row_size: Size, num_granules: usize) -> bool { - let has_fixed_free = self.header.fixed.next_free.has(); - let var_first = self.header.var.first; - let fixed_last = self.header.fixed.last; - - if num_granules == 0 { - // No granules needed. Just verify that there's space for the fixed part. - return has_fixed_free || gap_enough_size_for_row(var_first, fixed_last, fixed_row_size); - } - - // Determine the gap remaining after allocating for the fixed part. - let gap_remaining = gap_remaining_size(var_first, fixed_last); - let gap_avail_for_granules = if has_fixed_free { - // If we have a free fixed length block, then we can use the whole gap for var-len granules. - gap_remaining - } else { - // If we need to grow the fixed-length store into the gap, - if gap_remaining < fixed_row_size { - // If the gap is too small for fixed-length row, fail. - return false; - } - - // Otherwise, the space available in the gap for var-len granules - // is the current gap size less the fixed-len row size. - gap_remaining - fixed_row_size - }; - - // Convert the gap size to granules. - let gap_in_granules = VarLenGranule::space_to_granules(gap_avail_for_granules); - // Account for granules available in the freelist. - let needed_granules_after_freelist = num_granules.saturating_sub(self.header.var.freelist_len as usize); - - gap_in_granules >= needed_granules_after_freelist + has_space_for_row( + self.header.fixed.next_free.has(), + gap_remaining_size(self.header.var.first, self.header.fixed.last), + self.available_var_len_granules() as _, + fixed_row_size, + num_granules, + ) } /// Returns whether the row is full with respect to storing a fixed row with `fixed_row_size` @@ -1938,6 +1913,30 @@ impl Page { pub fn unmodified_hash(&self) -> Option<&blake3::Hash> { self.header.unmodified_hash.as_ref() } + + pub fn metadata(&self, fixed_row_size: Size) -> PageMetadata { + PageMetadata { + num_rows: self.num_rows() as _, + bytes_used_by_rows: self.bytes_used_by_rows(fixed_row_size) as _, + has_free_fixed_slot: self.header.fixed.next_free.has(), + gap_bytes: gap_remaining_size(self.header.var.first, self.header.fixed.last).0 as _, + available_granules: self.available_var_len_granules() as _, + } + } + + pub fn capacity(&self, fixed_row_size: Size) -> PageCapacity { + let allocated_fixed_slots = self.header.fixed.last / fixed_row_size; + let free_fixed_slots = allocated_fixed_slots + .checked_sub(self.header.fixed.num_rows as usize) + .expect("live row count exceeds allocated fixed slots"); + + PageCapacity { + num_rows: self.num_rows(), + gap_size: gap_remaining_size(self.header.var.first, self.header.fixed.last), + free_fixed_slots, + available_granules: self.available_var_len_granules(), + } + } } /// An iterator over the `PageOffset`s of all present fixed-length rows in a [`Page`]. @@ -1996,6 +1995,137 @@ impl<'page> Iterator for VarLenGranulesIter<'page> { } } +#[derive(Clone, Copy, Serialize, Deserialize)] +pub struct PageMetadata { + pub num_rows: u16, + pub bytes_used_by_rows: u32, + has_free_fixed_slot: bool, + gap_bytes: u16, + available_granules: u16, +} + +impl PageMetadata { + pub fn has_space_for_row(&self, fixed_row_size: Size, num_var_len_granules: usize) -> bool { + has_space_for_row( + self.has_free_fixed_slot, + Size(self.gap_bytes), + self.available_granules as _, + fixed_row_size, + num_var_len_granules, + ) + } + + pub fn is_full(&self, fixed_row_size: Size) -> bool { + !self.has_space_for_row(fixed_row_size, 0) + } + + pub fn available_var_len_granules(&self) -> usize { + self.available_granules as _ + } +} + +fn has_space_for_row( + has_fixed_free: bool, + gap_size: Size, + available_granules: usize, + fixed_row_size: Size, + num_granules: usize, +) -> bool { + if num_granules == 0 { + // No granules needed. Just verify that there's space for the fixed part. + return has_fixed_free || gap_size >= fixed_row_size; + } + + // Determine the gap remaining after allocating for the fixed part. + let gap_avail_for_granules = if has_fixed_free { + // If we have a free fixed length block, then we can use the whole gap for var-len granules. + gap_size + } else { + // If we need to grow the fixed-length store into the gap, + if gap_size < fixed_row_size { + // If the gap is too small for fixed-length row, fail. + return false; + } + + // Otherwise, the space available in the gap for var-len granules + // is the current gap size less the fixed-len row size. + gap_size - fixed_row_size + }; + + // Convert the gap size to granules. + let gap_in_granules = VarLenGranule::space_to_granules(gap_avail_for_granules); + // Account for granules available in the freelist. + let freelist_len = available_granules + .checked_sub(VarLenGranule::space_to_granules(gap_size)) + .expect("available granules must include the gap"); + let needed_granules_after_freelist = num_granules.saturating_sub(freelist_len); + + gap_in_granules >= needed_granules_after_freelist +} + +pub struct PageCapacity { + pub num_rows: usize, + gap_size: Size, + free_fixed_slots: usize, + available_granules: usize, +} + +impl PageCapacity { + pub fn empty(fixed_row_size: Size) -> Self { + let header = PageHeader::new(max_rows_in_page(fixed_row_size)); + Self { + num_rows: 0, + gap_size: gap_remaining_size(header.var.first, header.fixed.last), + free_fixed_slots: 0, + available_granules: header.available_var_len_granules(), + } + } + + pub fn has_space_for_row(&self, fixed_row_size: Size, num_granules: usize) -> bool { + has_space_for_row( + self.free_fixed_slots != 0, + self.gap_size, + self.available_granules as _, + fixed_row_size, + num_granules, + ) + } + + pub fn available_var_len_granules(&self) -> usize { + self.available_granules + } + + pub fn release_row(&mut self, row_granules: usize) { + assert!(self.num_rows != 0); + + self.num_rows -= 1; + self.free_fixed_slots += 1; + self.available_granules += row_granules; + } + + pub fn reserve_row(&mut self, fixed_row_size: Size, required_granules: usize) { + assert!(self.has_space_for_row(fixed_row_size, required_granules)); + + let gap_granules = VarLenGranule::space_to_granules(self.gap_size); + let freelist_granules = self + .available_granules + .checked_sub(gap_granules) + .expect("available granules must include the gap"); + + if self.free_fixed_slots != 0 { + self.free_fixed_slots -= 1; + } else { + self.gap_size = self.gap_size.checked_sub(fixed_row_size).unwrap(); + } + + let from_gap = required_granules.saturating_sub(freelist_granules); + self.gap_size = self.gap_size.checked_sub(VarLenGranule::SIZE * from_gap).unwrap(); + self.available_granules = + freelist_granules.saturating_sub(required_granules) + VarLenGranule::space_to_granules(self.gap_size); + self.num_rows += 1; + } +} + #[cfg(test)] pub(crate) mod tests { use super::*; diff --git a/crates/table/src/tiered/budget.rs b/crates/table/src/tiered/budget.rs new file mode 100644 index 00000000000..e820f2e2037 --- /dev/null +++ b/crates/table/src/tiered/budget.rs @@ -0,0 +1,104 @@ +use std::sync::{Arc, Mutex}; + +#[derive(Debug, thiserror::Error)] +#[error("memory limit exceeded")] +pub struct BudgetExceeded { + pub requested_bytes: u64, + pub accounted_bytes: u64, + pub hard_limit_bytes: u64, +} + +#[derive(Debug, thiserror::Error)] +pub enum ConfigError { + #[error("invalid configuration")] + InvalidBudgetOrder(ByteBudgetConfig), +} + +#[derive(Clone)] +pub struct ByteBudget { + state: Arc>, + config: ByteBudgetConfig, +} + +#[derive(Clone, Copy, Debug)] +pub struct ByteBudgetConfig { + pub low_water_bytes: u64, + pub soft_limit_bytes: u64, + pub hard_limit_bytes: u64, +} + +impl ByteBudgetConfig { + fn validate_then(self, f: impl FnOnce(Self) -> T) -> Result { + if self.low_water_bytes <= self.soft_limit_bytes && self.soft_limit_bytes <= self.hard_limit_bytes { + Ok(f(self)) + } else { + Err(ConfigError::InvalidBudgetOrder(self)) + } + } +} + +#[derive(Debug)] +pub struct ByteBudgetState { + accounted_bytes: u64, +} + +#[derive(Debug)] +pub struct BudgetPermit { + state: Arc>, + bytes: u64, +} + +impl Drop for BudgetPermit { + fn drop(&mut self) { + let mut state = self.state.lock().unwrap(); + state.accounted_bytes = state + .accounted_bytes + .checked_sub(self.bytes) + .expect("accounted bytes underflow"); + } +} + +pub struct ByteBudgetUsage { + pub accounted_bytes: u64, +} + +impl ByteBudget { + pub fn new(config: ByteBudgetConfig) -> Result { + config.validate_then(|config| Self { + state: Arc::new(Mutex::new(ByteBudgetState { accounted_bytes: 0 })), + config, + }) + } + + pub fn usage(&self) -> ByteBudgetUsage { + ByteBudgetUsage { + accounted_bytes: self.state.lock().unwrap().accounted_bytes, + } + } + + pub fn acquire(&self, bytes: u64) -> Result { + let mut state = self.state.lock().unwrap(); + if state.accounted_bytes + bytes <= self.config.hard_limit_bytes { + state.accounted_bytes += bytes; + Ok(BudgetPermit { + state: self.state.clone(), + bytes, + }) + } else { + Err(BudgetExceeded { + requested_bytes: bytes, + accounted_bytes: state.accounted_bytes, + hard_limit_bytes: self.config.hard_limit_bytes, + }) + } + } + + pub(super) fn force_acquire(&self, bytes: u64) -> BudgetPermit { + let mut state = self.state.lock().unwrap(); + state.accounted_bytes += bytes; + BudgetPermit { + state: self.state.clone(), + bytes, + } + } +} diff --git a/crates/table/src/tiered/mod.rs b/crates/table/src/tiered/mod.rs new file mode 100644 index 00000000000..c289c4e6544 --- /dev/null +++ b/crates/table/src/tiered/mod.rs @@ -0,0 +1,8 @@ +mod budget; +pub use budget::{BudgetExceeded, BudgetPermit, ByteBudget, ByteBudgetConfig, ByteBudgetUsage}; + +mod page_manager; +pub use page_manager::PageManager; + +mod page_set; +pub use page_set::PageSet; diff --git a/crates/table/src/tiered/page_manager.rs b/crates/table/src/tiered/page_manager.rs new file mode 100644 index 00000000000..c4d0d109088 --- /dev/null +++ b/crates/table/src/tiered/page_manager.rs @@ -0,0 +1,477 @@ +use std::{ + io, + sync::{ + atomic::{AtomicU64, Ordering}, + Arc, Mutex, Weak, + }, +}; + +use parking_lot::{ArcRwLockReadGuard, RawRwLock, RwLock, RwLockWriteGuard}; +use slab::Slab; +use spacetimedb_lib::bsatn::DecodeError; +use spacetimedb_sats::layout::Size; + +use crate::{ + indexes::{PageIndex, PAGE_SIZE}, + page::{self, Page, PageMetadata}, + page_pool::PagePool, + tiered::{BudgetExceeded, BudgetPermit, ByteBudget}, +}; + +pub type PageFrameReadGuard = ArcRwLockReadGuard>; + +pub trait PageBackingStore: Send + Sync + 'static { + /// Load a [Page] by its content hash from backing storage . + fn load_page(&self, hash: blake3::Hash) -> Result, PageIoError>; +} + +impl PageBackingStore for () { + fn load_page(&self, _: blake3::Hash) -> Result, PageIoError> { + unimplemented!("no page backing store configured") + } +} + +#[derive(Debug, thiserror::Error)] +pub enum PageError { + #[error("maximum number of pages exceeded")] + TooManyPages, + #[error(transparent)] + MemoryLimitExceeded(#[from] BudgetExceeded), + #[error(transparent)] + Page(page::Error), + #[error("page at index {0:?} is missing")] + MissingPage(PageIndex), + #[error("object {0} is missing")] + MissingObject(blake3::Hash), + #[error(transparent)] + Io(#[from] PageIoError), + #[error("error decoding page from bsatn")] + Deserialize(DecodeError), +} + +#[derive(Debug, thiserror::Error)] +pub enum PageIoError { + #[error("expected page hash {expected} doesn't match computed page hash {computed}")] + HashMismatch { + expected: blake3::Hash, + computed: blake3::Hash, + }, + #[error(transparent)] + Io(#[from] io::Error), +} + +#[derive(Clone, Copy)] +pub enum PageEvictionPolicy { + Evictable, + NeverEvict, +} + +#[derive(Clone)] +pub struct PageSlotHandle { + slot: Arc>, +} + +impl PageSlotHandle { + pub fn is_resident(&self) -> bool { + self.slot.lock().unwrap().is_resident() + } + + pub fn is_absent(&self) -> bool { + self.slot.lock().unwrap().is_absent() + } + + pub fn has_space_for_row(&self, fixed_row_size: Size, num_var_len_granules: usize) -> Option { + self.slot + .lock() + .unwrap() + .has_space_for_row(fixed_row_size, num_var_len_granules) + } + + pub fn is_full(&self, fixed_row_size: Size) -> Option { + self.slot.lock().unwrap().is_full(fixed_row_size) + } + + pub fn available_var_len_granules(&self) -> Option { + self.slot.lock().unwrap().available_var_len_granules() + } + + pub fn bytes_used_by_rows(&self, fixed_row_size: Size) -> usize { + self.slot.lock().unwrap().bytes_used_by_rows(fixed_row_size) + } + + pub fn metadata(&self, fixed_row_size: Size) -> Option { + self.slot.lock().unwrap().metadata(fixed_row_size) + } + + pub fn page(&self) -> Option { + self.slot.lock().unwrap().page().cloned() + } + + pub(super) fn free(&self) { + let mut slot = self.slot.lock().unwrap(); + *slot = PageSlot::Absent; + } +} + +pub enum PageSlot { + Absent, + Resident { + handle: PageHandle, + #[allow(unused)] + state: ResidentPageState, + }, + #[allow(unused)] + NonResident { + hash: blake3::Hash, + metadata: PageMetadata, + }, +} + +impl PageSlot { + pub fn is_resident(&self) -> bool { + matches!(self, Self::Resident { .. }) + } + + pub fn is_absent(&self) -> bool { + matches!(self, Self::Absent) + } + + pub fn has_space_for_row(&self, fixed_row_size: Size, num_var_len_granules: usize) -> Option { + match self { + PageSlot::Absent => None, + PageSlot::Resident { handle, .. } => { + Some(handle.read().has_space_for_row(fixed_row_size, num_var_len_granules)) + } + PageSlot::NonResident { metadata, .. } => { + Some(metadata.has_space_for_row(fixed_row_size, num_var_len_granules)) + } + } + } + + pub fn is_full(&self, fixed_row_size: Size) -> Option { + match self { + PageSlot::Absent => None, + PageSlot::Resident { handle, .. } => Some(handle.read().is_full(fixed_row_size)), + PageSlot::NonResident { metadata, .. } => Some(metadata.is_full(fixed_row_size)), + } + } + + pub fn available_var_len_granules(&self) -> Option { + match self { + PageSlot::Absent => None, + PageSlot::Resident { handle, .. } => Some(handle.read().available_var_len_granules()), + PageSlot::NonResident { metadata, .. } => Some(metadata.available_var_len_granules()), + } + } + + pub fn bytes_used_by_rows(&self, fixed_row_size: Size) -> usize { + match self { + PageSlot::Absent => 0, + PageSlot::Resident { handle, .. } => handle.read().bytes_used_by_rows(fixed_row_size), + PageSlot::NonResident { metadata, .. } => metadata.bytes_used_by_rows as _, + } + } + + pub fn metadata(&self, fixed_row_size: Size) -> Option { + match self { + PageSlot::Absent => None, + PageSlot::Resident { handle, .. } => Some(handle.read().metadata(fixed_row_size)), + PageSlot::NonResident { metadata, .. } => Some(*metadata), + } + } + + pub fn page(&self) -> Option<&PageHandle> { + match self { + PageSlot::Absent | PageSlot::NonResident { .. } => None, + PageSlot::Resident { handle, .. } => Some(handle), + } + } +} + +#[derive(Clone)] +pub struct PageHandle { + frame: Arc, +} + +impl PageHandle { + pub fn read(&self) -> PageFrameReadGuard { + self.frame.read() + } + + pub fn with_page_mut(&mut self, f: impl FnOnce(&mut Page) -> T) -> T { + let mut guard = self.frame.write(); + f(&mut guard) + } +} + +#[derive(Debug)] +pub struct PageFrame { + page: Arc>>, + #[allow(unused)] + permit: BudgetPermit, + access: Arc, +} + +impl PageFrame { + fn new(permit: BudgetPermit, page: Box, epoch: u64) -> Self { + Self { + page: RwLock::new(page).into(), + permit, + access: FrameAccess::new(epoch).into(), + } + } + + pub fn read(&self) -> PageFrameReadGuard { + RwLock::read_arc(&self.page) + } + + fn write(&self) -> RwLockWriteGuard<'_, Box> { + self.page.write() + } + + fn touch(&self, epoch: u64) { + self.access.touch(epoch); + } +} + +#[allow(unused)] +pub enum ResidentPageState { + Clean { hash: Option }, + Dirty { hash: Option }, +} + +pub struct ReservedPage { + permit: BudgetPermit, + page: Box, +} + +pub struct PageManager { + frames: RwLock, + pool: PagePool, + store: Arc, + memory: ByteBudget, + access_epoch: AtomicU64, +} + +impl PageManager { + pub fn new(pool: PagePool, store: Arc, memory: ByteBudget) -> Self { + Self { + frames: <_>::default(), + pool, + store, + memory, + access_epoch: <_>::default(), + } + } + + pub fn get( + &self, + slot: &PageSlotHandle, + eviction_policy: PageEvictionPolicy, + ) -> Result, PageError> { + Ok(self.may_fault(slot, eviction_policy)?.map(|frame| PageHandle { frame })) + } + + pub fn with_page_mut( + &self, + slot: &PageSlotHandle, + eviction_policy: PageEvictionPolicy, + f: impl FnOnce(&mut Page) -> T, + ) -> Result { + let frame = self + .may_fault(slot, eviction_policy)? + .expect("page requested for mutation to be present"); + let res = { + let mut page = frame.write(); + f(&mut page) + }; + let mut slot = slot.slot.lock().unwrap(); + *slot = PageSlot::Resident { + handle: PageHandle { frame }, + state: ResidentPageState::Dirty { hash: None }, + }; + Ok(res) + } + + fn may_fault( + &self, + slot: &PageSlotHandle, + eviction_policy: PageEvictionPolicy, + ) -> Result>, PageError> { + let mut slot_guard = slot.slot.lock().unwrap(); + match *slot_guard { + PageSlot::Absent => Ok(None), + PageSlot::Resident { ref handle, .. } => { + let frame = handle.frame.clone(); + frame.touch(self.access_epoch.fetch_add(1, Ordering::Relaxed)); + Ok(Some(frame)) + } + PageSlot::NonResident { hash, .. } => { + let permit = self.acquire_memory_budget_permit()?; + let page = self.store.load_page(hash)?; + let frame = self.frames.write().register( + permit, + page, + self.access_epoch.fetch_add(1, Ordering::Relaxed), + |frame| { + let registry_entry = FrameRegistryEntry::new(&frame, &slot.slot, eviction_policy); + let resident = PageSlot::Resident { + handle: PageHandle { frame }, + state: ResidentPageState::Clean { hash: Some(hash) }, + }; + *slot_guard = resident; + registry_entry + }, + ); + + Ok(Some(frame)) + } + } + } + + pub fn allocate( + &self, + fixed_row_size: Size, + eviction_policy: PageEvictionPolicy, + ) -> Result { + let reservation = self.reserve(fixed_row_size)?; + Ok(self.redeem(reservation, eviction_policy)) + } + + pub fn reserve(&self, fixed_row_size: Size) -> Result { + let permit = self.acquire_memory_budget_permit()?; + let page = self.pool.take_with_fixed_row_size(fixed_row_size); + + Ok(ReservedPage { permit, page }) + } + + pub fn redeem( + &self, + ReservedPage { permit, page }: ReservedPage, + eviction_policy: PageEvictionPolicy, + ) -> PageSlotHandle { + let slot = Arc::new(Mutex::new(PageSlot::Absent)); + { + let mut slot_guard = slot.lock().unwrap(); + self.frames.write().register( + permit, + page, + self.access_epoch.fetch_add(1, Ordering::Relaxed), + |frame| { + let registry_entry = FrameRegistryEntry::new(&frame, &slot, eviction_policy); + *slot_guard = PageSlot::Resident { + handle: PageHandle { frame }, + state: ResidentPageState::Clean { hash: None }, + }; + registry_entry + }, + ); + } + + PageSlotHandle { slot } + } + + pub(super) fn register( + &self, + eviction_policy: PageEvictionPolicy, + pages: impl IntoIterator>>, + ) -> Vec { + let mut handles = Vec::new(); + for page in pages { + let slot = Arc::new(Mutex::new(PageSlot::Absent)); + if let Some(page) = page { + let permit = self.force_acquire_memory_budget_limit(); + let mut slot_gard = slot.lock().unwrap(); + self.frames.write().register( + permit, + page, + self.access_epoch.fetch_add(1, Ordering::Relaxed), + |frame| { + let registry_entry = FrameRegistryEntry::new(&frame, &slot, eviction_policy); + *slot_gard = PageSlot::Resident { + handle: PageHandle { frame }, + state: ResidentPageState::Clean { hash: None }, + }; + registry_entry + }, + ); + } + + handles.push(PageSlotHandle { slot }); + } + + handles + } + + fn force_acquire_memory_budget_limit(&self) -> BudgetPermit { + self.memory.force_acquire(PAGE_SIZE as _) + } + + fn acquire_memory_budget_permit(&self) -> Result { + // TODO: Try to evict pages if acquisition fails. + self.memory.acquire(PAGE_SIZE as _) + } +} + +#[derive(Default)] +struct FrameRegistry { + frames: Slab, +} + +impl FrameRegistry { + pub fn register( + &mut self, + permit: BudgetPermit, + page: Box, + epoch: u64, + mk_entry: impl FnOnce(Arc) -> FrameRegistryEntry, + ) -> Arc { + let entry = self.frames.vacant_entry(); + let frame = Arc::new(PageFrame::new(permit, page, epoch)); + entry.insert(mk_entry(frame.clone())); + frame + } +} + +/// Access counters for cache eviction. +/// +/// Shared between [FrameRegistryEntry] and [PageFrame], to avoid lock +/// contention on the [FrameRegistry] for updates. +#[derive(Debug)] +struct FrameAccess { + last_access_epoch: AtomicU64, + last_access_count: AtomicU64, +} + +impl FrameAccess { + fn new(epoch: u64) -> Self { + Self { + last_access_epoch: AtomicU64::new(epoch), + last_access_count: <_>::default(), + } + } + + fn touch(&self, epoch: u64) { + self.last_access_epoch.store(epoch, Ordering::Relaxed); + self.last_access_count.fetch_add(1, Ordering::Relaxed); + } +} + +#[allow(unused)] +pub struct FrameRegistryEntry { + frame: Weak, + slot: Weak>, + eviction_policy: PageEvictionPolicy, + access: Arc, +} + +impl FrameRegistryEntry { + pub fn new(frame: &Arc, slot: &Arc>, eviction_policy: PageEvictionPolicy) -> Self { + Self { + access: Arc::clone(&frame.access), + frame: Arc::downgrade(frame), + slot: Arc::downgrade(slot), + eviction_policy, + } + } +} diff --git a/crates/table/src/tiered/page_set.rs b/crates/table/src/tiered/page_set.rs new file mode 100644 index 00000000000..b6daccbf9f0 --- /dev/null +++ b/crates/table/src/tiered/page_set.rs @@ -0,0 +1,255 @@ +use std::{collections::BTreeSet, sync::Arc}; + +use spacetimedb_sats::layout::Size; + +use crate::{ + blob_store::BlobStore, + indexes::{PageIndex, RowPointer}, + page::Page, + table::BlobNumBytes, + tiered::page_manager::{PageEvictionPolicy, PageHandle, PageManager, PageSlotHandle, ReservedPage}, + var_len::VarLenMembers, +}; + +pub use crate::tiered::page_manager::PageError; + +pub struct PageSet { + eviction_policy: PageEvictionPolicy, + slots: Vec, + free_page_slots: BTreeSet, + non_full_pages: BTreeSet<(usize, PageIndex)>, + manager: Arc, +} + +impl PageSet { + /// Create a fresh, empty page set. + pub fn new(manager: Arc, eviction_policy: PageEvictionPolicy) -> Self { + Self { + eviction_policy, + slots: <_>::default(), + free_page_slots: <_>::default(), + non_full_pages: <_>::default(), + manager, + } + } + + /// Populate this [PageSet] with the `pages` acquired externally. + /// + /// This method is provided for compatibility. It is used when restoring + /// from a snapshot. + pub fn set_contents(&mut self, pages: impl IntoIterator>>, fixed_row_size: Size) { + assert!(self.slots.is_empty()); + + self.slots = self.manager.register(self.eviction_policy, pages); + self.non_full_pages = self + .slots + .iter() + .enumerate() + .filter_map(|(idx, page)| { + if !page.is_full(fixed_row_size)? { + page.available_var_len_granules() + .map(|granules| (granules, PageIndex(idx as _))) + } else { + None + } + }) + .collect(); + self.free_page_slots = self + .slots + .iter() + .enumerate() + .filter_map(|(idx, page)| page.is_absent().then_some(PageIndex(idx as _))) + .collect(); + } + + pub fn get_page(&self, index: PageIndex) -> Result, PageError> { + let Some(slot) = self.slots.get(index.idx()) else { + return Ok(None); + }; + self.manager.get(slot, self.eviction_policy) + } + + pub fn with_page_mut( + &mut self, + index: PageIndex, + fixed_row_size: Size, + f: impl FnOnce(&mut Page) -> T, + ) -> Result { + let Some(slot) = self.slots.get(index.idx()) else { + return Err(PageError::MissingPage(index)); + }; + let mut is_empty = false; + let mut available_granules = None; + let ret = self.manager.with_page_mut(slot, self.eviction_policy, |page| { + self.non_full_pages.remove(&(page.available_var_len_granules(), index)); + let ret = f(page); + is_empty = page.num_rows() == 0; + if !is_empty && !page.is_full(fixed_row_size) { + available_granules = Some(page.available_var_len_granules()); + } + ret + })?; + + if is_empty { + slot.free(); + self.free_page_slots.insert(index); + } else if let Some(available_granules) = available_granules { + self.non_full_pages.insert((available_granules, index)); + } + Ok(ret) + } + + pub fn with_page_to_insert_row( + &mut self, + fixed_row_size: Size, + num_var_len_granules: usize, + reservation: Option<(PageIndex, ReservedPage)>, + f: impl FnOnce(&mut Page) -> T, + ) -> Result<(PageIndex, T), PageError> { + let index = match reservation { + None => { + //eprintln!("finding page for insert"); + self.find_page_with_space_for_row(fixed_row_size, num_var_len_granules)? + } + Some((index, page)) => { + //eprintln!("registering reservation at {index:?} for insert"); + self.register(index, page); + index + } + }; + self.with_page_mut(index, fixed_row_size, f).map(|res| (index, res)) + } + + /// Free the row that is pointed to by `row_ptr`, + /// marking its fixed-len storage + /// and var-len storage granules as available for re-use. + /// + /// # Safety + /// + /// The `row_ptr` must point to a valid row in this page manager, + /// of `fixed_row_size` bytes for the fixed part. + /// + /// The `fixed_row_size` must be consistent + /// with what has been passed to the manager in all other operations + /// and must be consistent with the `var_len_visitor` the manager was made with. + pub unsafe fn delete_row( + &mut self, + var_len_visitor: &impl VarLenMembers, + fixed_row_size: Size, + row_ptr: RowPointer, + blob_store: &mut dyn BlobStore, + ) -> Result { + let page_index = row_ptr.page_index(); + self.with_page_mut(page_index, fixed_row_size, |page| { + // SAFETY: + // - `row_ptr.page_offset()` does point to a valid row in this page + // as the caller promised that `row_ptr` points to a valid row in `self`. + // + // - `fixed_row_size` is consistent with the size in bytes of the fixed part of the row. + // The size is also conistent with `var_len_visitor`. + unsafe { page.delete_row(row_ptr.page_offset(), fixed_row_size, var_len_visitor, blob_store) } + }) + } + + pub fn is_resident(&self, index: PageIndex) -> bool { + self.slots[index.idx()].is_resident() + } + + /// Find a page with sufficient available space to store a row of size `fixed_row_size` + /// containing `num_var_len_granules` granules of var-len data. + /// + /// Retrieving a page in this way will remove it from the non-full set. + /// After performing an insertion, the caller should use [`Self::record_page_non_full`] + /// to restore the page to the non-full set. + fn find_page_with_space_for_row( + &mut self, + fixed_row_size: Size, + num_var_len_granules: usize, + ) -> Result { + if let Some((page_num_free_granules, page_idx)) = self + .non_full_pages + .range((num_var_len_granules, PageIndex(0))..) + .copied() + .find(|(_, page_index)| { + self.slots[page_index.idx()] + .has_space_for_row(fixed_row_size, num_var_len_granules) + .expect("page in `self.non_full_pages` to be present in `self.pages`") + }) + { + self.non_full_pages.remove(&(page_num_free_granules, page_idx)); + return Ok(page_idx); + } + + self.allocate_new_page(fixed_row_size) + } + + /// Allocates one additional page, + /// returning an error if the new number of pages would overflow `PageIndex::MAX`. + /// + /// The new page is initially empty, but is not added to the non-full set. + /// Callers should call [`Pages::record_page_non_full`] after operating on the new page. + fn allocate_new_page(&mut self, fixed_row_size: Size) -> Result { + let page_index = self + .free_page_slots + .pop_first() + .map(Ok) + .unwrap_or_else(|| self.can_allocate_new_page())?; + let slot = self.manager.allocate(fixed_row_size, self.eviction_policy)?; + // SAFETY: The page is resident and the index is from free slots or + // fresh. + unsafe { self.install_page_slot(page_index, slot) }; + + Ok(page_index) + } + + #[allow(unused)] + pub(crate) fn register(&mut self, index: PageIndex, reservation: ReservedPage) -> PageHandle { + let slot = self.manager.redeem(reservation, self.eviction_policy); + // SAFETY: The page is resident. The index was obtained during commit + // planning. + unsafe { self.install_page_slot(index, slot) } + } + + /// SAFETY: + /// - The [PageSlotHandle] must contain a resident page. + /// - If there already exists an entry at `index`, the existing slot must be + /// absent. + unsafe fn install_page_slot(&mut self, index: PageIndex, slot: PageSlotHandle) -> PageHandle { + let page = slot.page().expect("page must be resident"); + + let idx = index.idx(); + if idx == self.slots.len() { + self.slots.push(slot); + } else { + assert!(self.slots[idx].is_absent()); + self.free_page_slots.remove(&index); + self.slots[idx] = slot; + } + + self.non_full_pages + .insert((page.read().available_var_len_granules(), index)); + + page + } + + /// Is there space to allocate another page? + pub fn can_allocate_new_page(&self) -> Result { + let new_idx = self.slots.len(); + if new_idx <= PageIndex::MAX.idx() { + Ok(PageIndex(new_idx as _)) + } else { + Err(PageError::TooManyPages) + } + } + + pub fn iter_present_pages(&self) -> impl Iterator { + self.slots.iter().filter(|slot| !slot.is_absent()) + } + + pub fn iter_present_pages_with_page_index(&self) -> impl Iterator { + self.slots + .iter() + .enumerate() + .filter_map(|(idx, slot)| (!slot.is_absent()).then_some((PageIndex(idx as _), slot))) + } +}