From e394bbc43bf054e0b64c6f1e8a3a8b25460c9008 Mon Sep 17 00:00:00 2001 From: Phoebe Goldman Date: Thu, 20 Aug 2026 11:01:35 -0400 Subject: [PATCH 01/10] Represent absent pages in `Table` and snapshots Preparation for freeing empty pages. Because the snapshot format depends on the density of page vectors and didn't previously reserve a sentinel, to preserve rollback safety we have to do this preparation step before actually implementing freeing pages as part of the row delete operation. In this PR, the `Table`/`Pages` switches to a `Vec>>`, with pages allowed to be absent. However, until a later patch, outside of tests, no page entry will ever be `None`. The table code is still able to use and reason about `None` page entries, as they may arise if we deploy said later patch, free a page, capture a snapshot, then roll back to this version. In the snapshot format, absent pages are recorded in the pages vec as the all-zeroes hash. Page objects are not written or read in this case; the all-zeroes hash does not correspond to an actual object on disk. When allocating a new page, we attempt to fill the lowest empty slot. We do this in log time by storing a `BTreeSet` of the empty slots, and popping the lowest value from it to use as the slot for the newly allocated page. I believe that for at least some access patterns, this should allow us to gradually converge on a dense array of pages in the case where rows are deleted at a higher rate than new inserts. --- crates/snapshot/src/lib.rs | 76 +++++++----- crates/table/src/pages.rs | 241 +++++++++++++++++++++++++++---------- crates/table/src/table.rs | 239 +++++++++++++++++++++++++++++++----- 3 files changed, 429 insertions(+), 127 deletions(-) diff --git a/crates/snapshot/src/lib.rs b/crates/snapshot/src/lib.rs index b04ee20216f..f8b2eb57f64 100644 --- a/crates/snapshot/src/lib.rs +++ b/crates/snapshot/src/lib.rs @@ -187,6 +187,8 @@ pub const INVALID_SNAPSHOT_DIR_EXT: &str = "invalid_snapshot"; /// File extension of snapshots which have been archived pub const ARCHIVED_SNAPSHOT_EXT: &str = "archived_snapshot"; +pub const ZERO_HASH_DENOTING_ABSENT_PAGE: blake3::Hash = blake3::Hash::from_bytes([0; _]); + #[derive(Clone, Serialize, Deserialize)] /// The hash and refcount of a single blob in the blob store. struct BlobEntry { @@ -428,7 +430,13 @@ impl Snapshot { ) -> Result<(), SnapshotError> { let pages = table .iter_pages_with_hashes() - .map(|(hash, page)| Self::write_page(object_repo, page, hash, prev_snapshot, counter)) + .map(|option| { + if let Some((hash, page)) = option { + Self::write_page(object_repo, page, hash, prev_snapshot, counter) + } else { + Ok(ZERO_HASH_DENOTING_ABSENT_PAGE) + } + }) .collect::, SnapshotError>>()?; self.tables.push(TableEntry { @@ -592,40 +600,44 @@ impl Snapshot { pages: &[blake3::Hash], page_pool: &PagePool, metrics: &mut SnapshotReadKindMetrics, - ) -> Result>, SnapshotError> { + ) -> Result>>, SnapshotError> { pages .iter() .map(|hash| { - // Read the BSATN bytes of the on-disk page object. - let (buf, disk_bytes) = - Self::read_object_with_disk_bytes(object_repo, hash.as_bytes(), ObjectType::Page(*hash))?; - metrics.disk_bytes += disk_bytes; - - // Deserialize the bytes into a `Page`. - let page = page_pool.take_deserialize_from(&buf); - let page = page.map_err(|cause| SnapshotError::Deserialize { - ty: ObjectType::Page(*hash), - source_repo: object_repo.root().to_path_buf(), - cause, - })?; - - // Compute the hash of the page. - let hash_start = Instant::now(); - let computed_hash = page.content_hash(); - metrics.hash_time += hash_start.elapsed(); - - // Compare the computed hash to the one recorded in the `Snapshot`, - // and fail if they do not match. - if *hash != computed_hash { - return Err(SnapshotError::HashMismatch { + if *hash == ZERO_HASH_DENOTING_ABSENT_PAGE { + Ok(None) + } else { + // Read the BSATN bytes of the on-disk page object. + let (buf, disk_bytes) = + Self::read_object_with_disk_bytes(object_repo, hash.as_bytes(), ObjectType::Page(*hash))?; + metrics.disk_bytes += disk_bytes; + + // Deserialize the bytes into a `Page`. + let page = page_pool.take_deserialize_from(&buf); + let page = page.map_err(|cause| SnapshotError::Deserialize { ty: ObjectType::Page(*hash), - expected: *hash.as_bytes(), - computed: *computed_hash.as_bytes(), source_repo: object_repo.root().to_path_buf(), - }); - } + cause, + })?; + + // Compute the hash of the page. + let hash_start = Instant::now(); + let computed_hash = page.content_hash(); + metrics.hash_time += hash_start.elapsed(); + + // Compare the computed hash to the one recorded in the `Snapshot`, + // and fail if they do not match. + if *hash != computed_hash { + return Err(SnapshotError::HashMismatch { + ty: ObjectType::Page(*hash), + expected: *hash.as_bytes(), + computed: *computed_hash.as_bytes(), + source_repo: object_repo.root().to_path_buf(), + }); + } - Ok::, SnapshotError>(page) + Ok::>, SnapshotError>(Some(page)) + } }) .collect() } @@ -635,7 +647,7 @@ impl Snapshot { TableEntry { table_id, pages }: &TableEntry, page_pool: &PagePool, metrics: &mut SnapshotReadKindMetrics, - ) -> Result<(TableId, Vec>), SnapshotError> { + ) -> Result<(TableId, Vec>>), SnapshotError> { Ok(( *table_id, Self::reconstruct_one_table_pages(object_repo, pages, page_pool, metrics)?, @@ -657,7 +669,7 @@ impl Snapshot { object_repo: &DirTrie, page_pool: &PagePool, metrics: &mut SnapshotReadKindMetrics, - ) -> Result>>, SnapshotError> { + ) -> Result>>>, SnapshotError> { self.tables .iter() .map(|tbl| Self::reconstruct_one_table(object_repo, tbl, page_pool, metrics)) @@ -1550,7 +1562,7 @@ pub struct ReconstructedSnapshot { /// This includes the system tables, /// so the schema of user-defined tables can be recovered /// given knowledge of the schema of `st_table` and `st_column`. - pub tables: BTreeMap>>, + pub tables: BTreeMap>>>, /// If the snapshot was compressed or not. pub compress_type: CompressType, /// Metrics collected while reading this snapshot from disk. diff --git a/crates/table/src/pages.rs b/crates/table/src/pages.rs index b815b6415e2..5949419deea 100644 --- a/crates/table/src/pages.rs +++ b/crates/table/src/pages.rs @@ -6,7 +6,7 @@ use super::page::Page; use super::page_pool::PagePool; use super::table::BlobNumBytes; use super::var_len::VarLenMembers; -use core::ops::{ControlFlow, Deref, Index, IndexMut}; +use core::ops::{ControlFlow, Deref}; use spacetimedb_sats::layout::Size; use spacetimedb_sats::memory_usage::MemoryUsage; use std::collections::BTreeSet; @@ -21,25 +21,24 @@ pub enum Error { Page(#[from] super::page::Error), } -impl Index for Pages { - type Output = Page; - - fn index(&self, pi: PageIndex) -> &Self::Output { - &self.pages[pi.idx()] - } -} - -impl IndexMut for Pages { - fn index_mut(&mut self, pi: PageIndex) -> &mut Self::Output { - &mut self.pages[pi.idx()] - } -} - /// A manager of [`Page`]s. #[derive(Default, Debug, PartialEq, Eq)] pub struct Pages { /// The collection of pages under management. - pages: Vec>, + /// + /// A `None` here is a page that was previously allocated and then became empty during [`Self::delete_row`]. + pages: Vec>>, + /// Indexes into [`self.pages`](Self::pages) which hold `None`, where newly-allocated pages can be inserted. + /// + /// When freeing a page other than the last one, resulting in a `None` in [`self.pages`](Self::pages), + /// we'll insert the freed page's [`PageIndex`] into this set. + /// When allocating a new page, we'll use [`BTreeSet::pop_first`] to insert it into the lowest-available [`PageIndex`]. + /// + /// Using a `BTreeSet` and popping the lowest item means that, over time, + /// datastores can converge on a dense sequence of pages. + /// This assumes that deletes are distributed randomly. + /// Low-indexed pages will be quickly replaced, while high-indexed pages will remain vacant for longer. + free_page_slots: BTreeSet, /// The set of pages that aren't yet full, /// sorted by the number of var-len granules available in each page. /// @@ -59,12 +58,28 @@ pub struct Pages { impl MemoryUsage for Pages { fn heap_usage(&self) -> usize { - let Self { pages, non_full_pages } = self; - pages.heap_usage() + non_full_pages.heap_usage() + let Self { + pages, + free_page_slots, + non_full_pages, + } = self; + pages.heap_usage() + free_page_slots.heap_usage() + non_full_pages.heap_usage() } } impl Pages { + pub fn get(&self, page_index: PageIndex) -> Option<&Page> { + self.pages + .get(page_index.idx()) + .and_then(|page_slot| page_slot.as_deref()) + } + + pub fn get_mut(&mut self, page_index: PageIndex) -> Option<&mut Page> { + self.pages + .get_mut(page_index.idx()) + .and_then(|page_slot| page_slot.as_deref_mut()) + } + #[cfg(test)] pub(crate) fn assert_non_full_pages_consistent(&self, fixed_row_size: Size) { let mut seen_page_indexes = BTreeSet::new(); @@ -78,30 +93,45 @@ impl Pages { for (idx, page) in self.pages.iter().enumerate() { let page_index = PageIndex(idx as u64); - let is_full = page.is_full(fixed_row_size); - let available_granules = page.available_var_len_granules(); - let entries_for_page: Vec<_> = self - .non_full_pages - .iter() - .copied() - .filter(|&(_, idx)| idx == page_index) - .collect(); - - if is_full { + if let Some(page) = page { + let is_full = page.is_full(fixed_row_size); + let available_granules = page.available_var_len_granules(); + let entries_for_page: Vec<_> = self + .non_full_pages + .iter() + .copied() + .filter(|&(_, idx)| idx == page_index) + .collect(); + + if is_full { + assert!( + entries_for_page.is_empty(), + "page {:?} has 0 available var-len granules but appears in non_full_pages as {:?}", + page_index, + entries_for_page + ); + } else { + assert_eq!( + entries_for_page, + vec![(available_granules, page_index)], + "page {:?} has {} available var-len granules but non_full_pages has {:?}", + page_index, + available_granules, + entries_for_page + ); + } + } else { + let entries_for_page: Vec<_> = self + .non_full_pages + .iter() + .copied() + .filter(|&(_free_granules, idx)| idx == page_index) + .collect(); assert!( entries_for_page.is_empty(), - "page {:?} has 0 available var-len granules but appears in non_full_pages as {:?}", + "page slot {:?} is None, but appears in non_full_pages as {:?}", page_index, - entries_for_page - ); - } else { - assert_eq!( entries_for_page, - vec![(available_granules, page_index)], - "page {:?} has {} available var-len granules but non_full_pages has {:?}", - page_index, - available_granules, - entries_for_page ); } } @@ -117,21 +147,45 @@ impl Pages { } } - /// Get a mutable reference to a `Page`. - /// - /// Used in benchmarks. Internal operators will prefer directly indexing into `self.pages`, - /// as that allows split borrows. - pub fn get_page_mut(&mut self, page: PageIndex) -> &mut Page { - &mut self.pages[page.idx()] + /// The number of present pages in `self`, i.e. those which have been allocated and not since freed. + pub fn num_present_pages(&self) -> usize { + self.pages + .len() + .checked_sub(self.free_page_slots.len()) + .expect("pages len to be greater than number of free slots") + } + + #[cfg(test)] + pub fn free_empty_page(&mut self, page_index: PageIndex) { + let page = self.get(page_index).expect("page to free to have been present"); + + assert_eq!(page.num_rows(), 0); + + let free_granules = page.available_var_len_granules(); + + let removed_from_non_full = self.non_full_pages.remove(&(free_granules, page_index)); + assert!(removed_from_non_full); + + self.pages[page_index.idx()] = None; + + let newly_inserted_into_free_pages_set = self.free_page_slots.insert(page_index); + + assert!(newly_inserted_into_free_pages_set) } /// Make all pages within `self` clear, /// deleting all rows. + // + // TODO(delete-free-page): Determine what to do with this method. + // It doesn't really make sense given that it clears all pages but doesn't delete them, + // but it's only used in benchmarks, and those benchmarks are using it specifically to bypass allocating new pages. #[doc(hidden)] // Used in benchmarks. pub fn clear(&mut self) { // Clear every page. for page in &mut self.pages { - page.clear(); + if let Some(page) = page { + page.clear(); + } } // Mark every page non-full. self.non_full_pages = (0..self.pages.len()) @@ -140,7 +194,10 @@ impl Pages { // but we'd have to do some amount of reasoning to demonstrate it was correct // based on the definition of `Page::clear`, // and why bother? - .map(|idx| (self.pages[idx].available_var_len_granules(), PageIndex(idx as u64))) + .filter_map(|idx| { + let idx = PageIndex(idx as u64); + self.get(idx).map(|page| (page.available_var_len_granules(), idx)) + }) .collect(); } @@ -150,7 +207,9 @@ impl Pages { /// Higher-level code paths are expected to go through [`super::de::read_row_from_pages`]. #[doc(hidden)] // Used in benchmarks. pub fn get_fixed_len_row(&self, row: RowPointer, fixed_row_size: Size) -> &Bytes { - self[row.page_index()].get_row_data(row.page_offset(), fixed_row_size) + self.get(row.page_index()) + .expect("`get_fixed_len_row` of row in not-present page") + .get_row_data(row.page_offset(), fixed_row_size) } /// Allocates one additional page, @@ -159,15 +218,27 @@ impl Pages { /// 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, pool: &PagePool, fixed_row_size: Size) -> Result { - let new_idx = self.can_allocate_new_page()?; + if let Some(idx) = self.free_page_slots.pop_first() { + // If `self` contains holes from previously-freed pages, + // fill the lowest-`PageIndex` such hole. + let page = pool.take_with_fixed_row_size(fixed_row_size); + self.pages[idx.idx()] = Some(page); + Ok(idx) + } else { + // If `self` is currently dense, try to put the new page at the end. + let new_idx = self.can_allocate_new_page()?; - let page = pool.take_with_fixed_row_size(fixed_row_size); - self.pages.push(page); + let page = pool.take_with_fixed_row_size(fixed_row_size); + self.pages.push(Some(page)); - Ok(new_idx) + Ok(new_idx) + } } /// Reserve a new, initially empty page. + // TODO(delete-free-page): Determine what to do with this method. + // It doesn't really make sense in a world where `Pages` doesn't contain empty `Page`s, + // but it's only used for tests and benches. pub fn reserve_empty_page(&mut self, pool: &PagePool, fixed_row_size: Size) -> Result { let idx = self.allocate_new_page(pool, fixed_row_size)?; self.record_page_non_full(idx, fixed_row_size); @@ -184,7 +255,9 @@ impl Pages { f: impl FnOnce(&mut Page) -> Res, ) -> Result<(PageIndex, Res), Error> { let page_index = self.find_page_with_space_for_row(pool, fixed_row_size, num_var_len_granules)?; - let res = f(&mut self[page_index]); + let res = f(&mut self + .get_mut(page_index) + .expect("page returned by `find_page_with_space_for_row` to be present in `self.pages`")); self.record_page_non_full(page_index, fixed_row_size); Ok((page_index, res)) } @@ -205,7 +278,11 @@ impl Pages { .non_full_pages .range((num_var_len_granules, PageIndex(0))..) .copied() - .find(|(_, page_idx)| self[*page_idx].has_space_for_row(fixed_row_size, num_var_len_granules)) + .find(|(_, page_idx)| { + self.get(*page_idx) + .expect("page in `self.non_full_pages` to be present in `self.pages`") + .has_space_for_row(fixed_row_size, num_var_len_granules) + }) { self.non_full_pages.remove(&(page_num_free_granules, page_idx)); return Ok(page_idx); @@ -284,7 +361,9 @@ impl Pages { let page_index = row_ptr.page_index(); self.with_updating_non_full_pages(page_index, fixed_row_size, |this| { - let page = &mut this[page_index]; + let page = this + .get_mut(page_index) + .expect("page containing `row_ptr` with validity safety invariant to be present"); // SAFETY: // - `row_ptr.page_offset()` does point to a valid row in this page @@ -300,13 +379,18 @@ impl Pages { /// then run `body` to update the page, and finally update [`Self::non_full_pages`] for its new fullness and capacity. /// /// `body` should not update any pages other than the one identified by `page_index`. + /// + /// If `page_index` does not refer to a present page, i.e. `self.pages[page_index.idx]` is `None`, + /// this method will panic. fn with_updating_non_full_pages( &mut self, page_index: PageIndex, fixed_row_size: Size, body: impl FnOnce(&mut Self) -> Ret, ) -> Ret { - let page = &self[page_index]; + let page = self + .get(page_index) + .expect("page to be present in `with_updating_non_full_pages`"); let full_before = page.is_full(fixed_row_size); let available_granules_before = page.available_var_len_granules(); @@ -352,10 +436,15 @@ impl Pages { /// Record the number of available var-len granules in the page at `self[page_index]` into [`Self::non_full_pages`]. /// /// Prior to calling this function, there must not be an entry for `page_index` in [`Self::non_full_pages`]. + /// + /// The page must still be present, that is, `self.pages[page_index.idx()]` must be `Some`. + /// Otherwise this method will panic. fn record_page_non_full(&mut self, page_index: PageIndex, fixed_row_size: Size) { debug_assert!(!self.non_full_pages.iter().any(|(_, idx)| *idx == page_index)); - let page = &self[page_index]; + let page = self + .get(page_index) + .expect("page to be present when recording its fullness"); let available_granules = page.available_var_len_granules(); if !page.is_full(fixed_row_size) { @@ -372,6 +461,8 @@ impl Pages { /// /// - The `fixed_row_size` is consistent with the `var_len_visitor` /// and is equal to the value provided to all other methods on `self`. + // FIXME: this method appears not to correctly set `non_full_pages` on the result. + // It is also unused except for benchmarks, so it may be best to just remove it. pub unsafe fn copy_filter( &self, var_len_visitor: &impl VarLenMembers, @@ -388,7 +479,7 @@ impl Pages { let mut partial_page = None; // Copy each page. - for from_page in &self.pages { + for from_page in self.pages.iter().filter_map(|page| page.as_deref()) { // You may require multiple calls to `Page::copy_starting_from` // if `partial_page` fills up; // the first call starts from 0. @@ -447,7 +538,7 @@ impl Pages { if copy_starting_from.is_none() { partial_page = Some(to_page); } else { - partial_copied_pages.pages.push(to_page); + partial_copied_pages.pages.push(Some(to_page)); } } } @@ -465,26 +556,52 @@ impl Pages { /// Should only ever be called when `self.is_empty()`. /// /// Also populates `self.non_full_pages`. - pub fn set_contents(&mut self, pages: Vec>, fixed_row_size: Size) { + pub fn set_contents(&mut self, pages: Vec>>, fixed_row_size: Size) { debug_assert!(self.is_empty()); self.non_full_pages = pages .iter() .enumerate() .filter_map(|(idx, page)| { - (!page.is_full(fixed_row_size)).then_some((page.available_var_len_granules(), PageIndex(idx as _))) + page.as_ref().and_then(|page| { + (!page.is_full(fixed_row_size)).then_some((page.available_var_len_granules(), PageIndex(idx as _))) + }) }) .collect(); + self.free_page_slots = pages + .iter() + .enumerate() + .filter_map(|(idx, page)| page.is_none().then_some(PageIndex(idx as _))) + .collect(); self.pages = pages; } /// Consumes the page manager, returning all the pages it held. + /// + /// Indexes in the iterator will not necessarily correspond to the pages' `PageIndex`es. pub fn into_page_iter(self) -> impl Iterator> { - self.pages.into_iter() + self.pages.into_iter().filter_map(|page| page) + } + + /// Iterate over only those pages in `self` that are present. + /// + /// Indexes in the iterator will not necessarily correspond to the pages' `PageIndex`es. + /// `pages.iter_present_pages().enumerate()` does not yield the correct `PageIndex`es for the yielded pages. + /// Use [`Self::iter_present_pages_with_page_index`] for a version that also yields correct `PageIndex`es. + pub fn iter_present_pages(&self) -> impl Iterator { + self.pages.iter().filter_map(|page| page.as_deref()) + } + + /// Iterate over only those pages in `self` that are present, paired with their `PageIndex`es. + pub fn iter_present_pages_with_page_index(&self) -> impl Iterator { + self.pages + .iter() + .enumerate() + .filter_map(|(idx, page)| page.as_deref().map(|page| (PageIndex(idx as u64), page))) } } impl Deref for Pages { - type Target = [Box]; + type Target = [Option>]; fn deref(&self) -> &Self::Target { &self.pages diff --git a/crates/table/src/table.rs b/crates/table/src/table.rs index 7fb63a2e429..75c864d5e49 100644 --- a/crates/table/src/table.rs +++ b/crates/table/src/table.rs @@ -191,7 +191,7 @@ impl TableInner { } fn try_page_and_offset(&self, ptr: RowPointer) -> Option<(&Page, PageOffset)> { - (ptr.page_index().idx() < self.pages.len()).then(|| (&self.pages[ptr.page_index()], ptr.page_offset())) + self.pages.get(ptr.page_index()).map(|page| (page, ptr.page_offset())) } /// Returns the page and page offset that `ptr` points to. @@ -200,7 +200,7 @@ impl TableInner { } } -static_assert_size!(Table, 264); +static_assert_size!(Table, 288); impl MemoryUsage for Table { fn heap_usage(&self) -> usize { @@ -783,7 +783,12 @@ impl Table { }; let fixed_row_size = self.inner.row_layout.size(); - let fixed_buf = self.inner.pages[ptr.page_index()].get_fixed_row_data_mut(ptr.page_offset(), fixed_row_size); + let fixed_buf = self + .inner + .pages + .get_mut(ptr.page_index()) + .expect("page for row with validity safety constraint to be present") + .get_fixed_row_data_mut(ptr.page_offset(), fixed_row_size); fn write(dst: &mut [u8], offset: u16, bytes: [u8; N]) { let offset = offset as usize; @@ -1687,7 +1692,7 @@ impl Table { /// # Safety /// /// The schema of rows stored in the `pages` must exactly match `self.schema` and `self.inner.row_layout`. - pub unsafe fn set_pages(&mut self, pages: Vec>, blob_store: &dyn BlobStore) { + pub unsafe fn set_pages(&mut self, pages: Vec>>, blob_store: &dyn BlobStore) { self.inner.pages.set_contents(pages, self.inner.row_layout.size()); // Recompute table metadata based on the new pages. @@ -1697,6 +1702,9 @@ impl Table { } /// Consumes the table, returning some constituents needed for merge. + /// + /// The returned iterator of `Page`s will not necessarily yield pages in their `PageIndex` order. + /// It is intended for reclaiming pages to a pool, not for reading data out of the pages. pub fn consume_for_merge( self, ) -> ( @@ -1716,7 +1724,10 @@ impl Table { #[cfg(test)] fn reconstruct_num_rows(&self) -> u64 { - self.pages().iter().map(|page| page.reconstruct_num_rows() as u64).sum() + self.pages() + .iter_present_pages() + .map(|page| page.reconstruct_num_rows() as u64) + .sum() } /// Returns the number of bytes used by rows resident in this table. @@ -1739,7 +1750,7 @@ impl Table { // so that this runs in constant time rather than O(|Pages|). pub fn bytes_used_by_rows(&self) -> u64 { self.pages() - .iter() + .iter_present_pages() .map(|page| page.bytes_used_by_rows(self.inner.row_layout.size()) as u64) .sum() } @@ -1747,7 +1758,7 @@ impl Table { #[cfg(test)] fn reconstruct_bytes_used_by_rows(&self) -> u64 { self.pages() - .iter() + .iter_present_pages() .map(|page| unsafe { // Safety: `page` is in `self`, and was constructed using `self.innser.row_layout` and `self.inner.visitor_prog`, // so the three are mutually consistent. @@ -2136,8 +2147,9 @@ pub struct TableScanIter<'table> { /// The current page we're yielding rows from. /// When `None`, the iterator will attempt to advance to the next page, if any. current_page: Option>, - /// The current page index we are or will visit. + /// The current page index we visiting. current_page_idx: PageIndex, + /// The table the iterator is yielding rows from. pub(crate) table: &'table Table, /// The `BlobStore` that row references may refer into. @@ -2177,11 +2189,27 @@ impl<'a> Iterator for TableScanIter<'a> { // already incremented `self.current_page_idx`, // or we're just beginning and so it was initialized as 0. None => { + // Search for another page to yield from until we run past `self.pages.len`. // If there's another page, set `self.current_page` to it, // and go to the `Some` case in the match. - let next_page = self.table.pages().get(self.current_page_idx.idx())?; - let iter = next_page.iter_fixed_len(self.table.row_size()); - self.current_page = Some(iter); + + 'find_next_page: loop { + if self.current_page_idx.idx() >= self.table.pages().len() { + // We're past the end of the pages. + return None; + } + + if let Some(next_page) = self.table.pages().get(self.current_page_idx) { + // There's another page, so start yielding from it. + let iter = next_page.iter_fixed_len(self.table.row_size()); + self.current_page = Some(iter); + break 'find_next_page; + } else { + // The next page slot is not occupied by a page, + // so continue past it to the next page slot. + self.current_page_idx.0 += 1; + } + } } } } @@ -2457,19 +2485,26 @@ impl Table { &self.inner.pages } + #[cfg(test)] + fn pages_mut(&mut self) -> &mut Pages { + &mut self.inner.pages + } + /// Iterates over each [`Page`] in this table, ensuring that its hash is computed before yielding it. /// /// Used when capturing a snapshot. - pub fn iter_pages_with_hashes(&mut self) -> impl Iterator { + pub fn iter_pages_with_hashes(&mut self) -> impl Iterator> { self.inner.pages.iter_mut().map(|page| { - let hash = page.save_or_get_content_hash(); - (hash, &**page) + page.as_mut().map(|page| { + let hash = page.save_or_get_content_hash(); + (hash, &**page) + }) }) } /// Returns the number of pages storing the physical rows of this table. fn num_pages(&self) -> usize { - self.inner.pages.len() + self.inner.pages.num_present_pages() } /// Returns the [`StaticLayout`] for this table, @@ -2563,7 +2598,8 @@ pub(crate) mod test { // Reserve a page so that we can check the hash. let pi = table.inner.pages.reserve_empty_page(&pool, table.row_size()).unwrap(); - let hash_pre_ins = hash_unmodified_save_get(&mut table.inner.pages[pi]); + let hash_pre_ins = + hash_unmodified_save_get(table.inner.pages.get_mut(pi).expect("reserved page to be present")); // Insert the row (0, 0). table @@ -2571,7 +2607,8 @@ pub(crate) mod test { .expect("Initial insert failed"); // Inserting cleared the hash. - let hash_post_ins = hash_unmodified_save_get(&mut table.inner.pages[pi]); + let hash_post_ins = + hash_unmodified_save_get(table.inner.pages.get_mut(pi).expect("reserved page to be present")); assert_ne!(hash_pre_ins, hash_post_ins); // Try to insert the row (0, 1), and assert that we get the expected error. @@ -2593,7 +2630,15 @@ pub(crate) mod test { // Second insert did clear the hash while we had a constraint violation, // as constraint checking is done after insertion and then rolled back. - assert_eq!(table.inner.pages[pi].unmodified_hash(), None); + assert_eq!( + table + .inner + .pages + .get(pi) + .expect("reserved page to be present") + .unmodified_hash(), + None + ); } fn insert_retrieve_body(ty: impl Into, val: impl Into) -> TestCaseResult { @@ -2608,7 +2653,7 @@ pub(crate) mod test { prop_assert_eq!(table.pointers_for(hash), &[ptr]); prop_assert_eq!(table.inner.pages.len(), 1); - prop_assert_eq!(table.inner.pages[PageIndex(0)].num_rows(), 1); + prop_assert_eq!(table.inner.pages.get(PageIndex(0)).map(|page| page.num_rows()), Some(1)); let row_ref = table.get_row_ref(&blob_store, ptr).unwrap(); prop_assert_eq!(row_ref.to_product_value(), val.clone()); @@ -2738,21 +2783,23 @@ pub(crate) mod test { prop_assert_eq!(table.pointers_for(hash), &[ptr]); prop_assert_eq!(table.inner.pages.len(), 1); - prop_assert_eq!(table.inner.pages[PageIndex(0)].num_rows(), 1); + prop_assert_eq!(table.inner.pages.get(PageIndex(0)).map(|page| page.num_rows()), Some(1)); prop_assert_eq!(&table.scan_rows(&blob_store).map(|r| r.pointer()).collect::>(), &[ptr]); prop_assert_eq!(table.row_count, 1); - let hash_pre_del = hash_unmodified_save_get(&mut table.inner.pages[ptr.page_index()]); + let hash_pre_del = hash_unmodified_save_get(table.inner.pages.get_mut(ptr.page_index()).expect("page containing row to be present")); table.delete(&mut blob_store, ptr, |_| ()); - let hash_post_del = hash_unmodified_save_get(&mut table.inner.pages[ptr.page_index()]); + // FIXME(delete-free-page): Page will no longer be present here after deleting its only row. + // Amend this test to insert two rows and only delete one of them. + let hash_post_del = hash_unmodified_save_get(table.inner.pages.get_mut(ptr.page_index()).expect("page to remain present after delete, until we move to freeing empty pages")); assert_ne!(hash_pre_del, hash_post_del); prop_assert_eq!(table.pointers_for(hash), &[]); prop_assert_eq!(table.inner.pages.len(), 1); - prop_assert_eq!(table.inner.pages[PageIndex(0)].num_rows(), 0); + prop_assert_eq!(table.inner.pages.get(PageIndex(0)).map(|page| page.num_rows()), Some(0)); prop_assert_eq!(table.row_count, 0); prop_assert!(&table.scan_rows(&blob_store).next().is_none()); @@ -2775,12 +2822,12 @@ pub(crate) mod test { let blob_uses = blob_store.usage_counter(); - let hash_pre_ins = hash_unmodified_save_get(&mut table.inner.pages[ptr.page_index()]); + let hash_pre_ins = hash_unmodified_save_get(table.inner.pages.get_mut(ptr.page_index()).expect("page containing row to be present")); prop_assert!(table.insert(&pool, &mut blob_store, &val).is_err()); // Hash was cleared and is different despite failure to insert. - let hash_post_ins = hash_unmodified_save_get(&mut table.inner.pages[ptr.page_index()]); + let hash_post_ins = hash_unmodified_save_get(table.inner.pages.get_mut(ptr.page_index()).expect("page containing row to be present")); assert_ne!(hash_pre_ins, hash_post_ins); prop_assert_eq!(table.row_count, 1); @@ -2790,7 +2837,7 @@ pub(crate) mod test { let blob_uses_after = blob_store.usage_counter(); prop_assert_eq!(blob_uses_after, blob_uses); - prop_assert_eq!(table.inner.pages[PageIndex(0)].num_rows(), 1); + prop_assert_eq!(table.inner.pages.get(PageIndex(0)).map(|page| page.num_rows()), Some(1)); prop_assert_eq!(&table.scan_rows(&blob_store).map(|r| r.pointer()).collect::>(), &[ptr]); } @@ -2916,9 +2963,8 @@ pub(crate) mod test { let simple = table .inner .pages - .iter() - .zip((0..).map(PageIndex)) - .flat_map(|(page, pi)| { + .iter_present_pages_with_page_index() + .flat_map(|(pi, page)| { page.iter_fixed_len(table.row_size()) .map(move |po| RowPointer::new(false, pi, po, table.squashed_offset)) }); @@ -2950,8 +2996,16 @@ pub(crate) mod test { next_value += 1; } - let first_page = &table.inner.pages[PageIndex(0)]; - let second_page = &table.inner.pages[PageIndex(1)]; + let first_page = &table + .inner + .pages + .get(PageIndex(0)) + .expect("page zero to be present after inserting many rows"); + let second_page = &table + .inner + .pages + .get(PageIndex(1)) + .expect("page one to be present after inserting many rows"); assert!(first_page.is_full(table.row_size())); assert_eq!(first_page.available_var_len_granules(), 0); assert!(!second_page.is_full(table.row_size())); @@ -2960,7 +3014,11 @@ pub(crate) mod test { let first_ptr = inserted_ptrs[0]; table.delete(&mut blob_store, first_ptr, |_| ()); - let first_page = &table.inner.pages[PageIndex(0)]; + let first_page = &table + .inner + .pages + .get(PageIndex(0)) + .expect("page zero to still be present after a delete that does not empty it"); assert!(!first_page.is_full(table.row_size())); assert_eq!(first_page.available_var_len_granules(), 0); @@ -2969,7 +3027,11 @@ pub(crate) mod test { assert_eq!(new_ptr.page_index(), first_ptr.page_index()); assert_eq!(new_ptr.page_offset(), first_ptr.page_offset()); - let first_page = &table.inner.pages[PageIndex(0)]; + let first_page = &table + .inner + .pages + .get(PageIndex(0)) + .expect("page zero to still be present after an insert"); assert!(first_page.is_full(table.row_size())); assert_eq!(first_page.available_var_len_granules(), 0); } @@ -3069,4 +3131,115 @@ pub(crate) mod test { ) .is_none()); } + + #[test] + fn table_iter_skips_absent_pages() { + let pool = PagePool::new_for_test(); + let blob_store = &mut NullBlobStore; + let mut table = table([AlgebraicType::I32].into()); + + let mut next_value = 0i32; + let mut row_ptrs = vec![]; + loop { + let (_, row_ref) = table.insert(&pool, blob_store, &product![next_value]).unwrap(); + next_value += 1; + let pointer = row_ref.pointer(); + row_ptrs.push(pointer); + if pointer.page_index() == PageIndex(1) { + break; + } + } + assert_eq!(table.pages().len(), 2); + + let scanned_rows = table + .scan_rows(blob_store) + .map(|row| row.read_col::(0).unwrap()) + .collect::>(); + let expected_contents = (0..next_value).collect::>(); + assert_eq!(scanned_rows, expected_contents); + + // Delete all the rows from page 0 before freeing it, to avoid leaving stale entries in the pointer map. + for ptr in row_ptrs.into_iter().take_while(|ptr| ptr.page_index() == PageIndex(0)) { + table + .delete(blob_store, ptr, |_| ()) + .expect("deleted row to have been present"); + } + + table.pages_mut().free_empty_page(PageIndex(0)); + + let scanned_rows = table + .scan_rows(blob_store) + .map(|row| row.read_col::(0).unwrap()) + .collect::>(); + assert_eq!(scanned_rows.len(), 1); + assert_eq!(scanned_rows[0], next_value - 1); + } + + #[test] + fn alloc_new_page_fills_hole() { + let pool = PagePool::new_for_test(); + let blob_store = &mut NullBlobStore; + let mut table = table([AlgebraicType::I32].into()); + + fn has_two_full_pages(table: &Table) -> bool { + table.pages().len() >= 2 + && table + .pages() + .get(PageIndex(0)) + .map(|page| page.is_full(table.inner.row_layout.size())) + .unwrap_or(false) + && table + .pages() + .get(PageIndex(1)) + .map(|page| page.is_full(table.inner.row_layout.size())) + .unwrap_or(false) + } + + let mut next_value = 0i32; + let mut row_ptrs = vec![]; + loop { + let (_, row_ref) = table.insert(&pool, blob_store, &product![next_value]).unwrap(); + next_value += 1; + row_ptrs.push(row_ref.pointer()); + if has_two_full_pages(&table) { + break; + } + } + + // Delete all the rows from page 0 before freeing it, to avoid leaving stale entries in the pointer map. + for ptr in row_ptrs.into_iter().take_while(|ptr| ptr.page_index() == PageIndex(0)) { + table + .delete(blob_store, ptr, |_| ()) + .expect("deleted row to have been present"); + } + + assert_eq!( + table + .pages() + .get(PageIndex(0)) + .expect("empty page to still be present") + .num_rows(), + 0 + ); + assert!(table + .pages() + .get(PageIndex(1)) + .expect("full page to still be present") + .is_full(table.inner.row_layout.size())); + + assert_eq!(table.num_pages(), 2); + + table.pages_mut().free_empty_page(PageIndex(0)); + + assert_eq!(table.num_pages(), 1); + assert_eq!(table.pages().len(), 2); + + let (_, row_ref) = table.insert(&pool, blob_store, &product![next_value]).unwrap(); + let ptr = row_ref.pointer(); + assert_eq!(ptr.page_index(), PageIndex(0)); + + assert_eq!(table.num_pages(), 2); + assert_eq!(table.pages().len(), 2); + assert!(table.pages().get(PageIndex(0)).is_some()); + } } From da456e7af9fc6a69a28f807b785b552d76cb95c8 Mon Sep 17 00:00:00 2001 From: Phoebe Goldman Date: Thu, 20 Aug 2026 14:34:19 -0400 Subject: [PATCH 02/10] Free pages which become empty during a delete This commit expands upon PR 5770, causing empty pages to transition to `None` or the all-zeroes hash in two circumstances: 1. When a delete causes a page to become empty. 2. When a page is empty while loading from a snapshot. Historical snapshots are not rewritten; all-zeroes page hashes will appear only in newly-created snapshots. The diff here is relatively inflated due to now needing a `PagePool` in delete operations, in order to return freed pages to the pool. Many `Table` operations which one might not initially expect to perform a delete require this, as many operations are implemented in terms of transient inserts which are quickly deleted. Also, during implementation, I removed some `Pages` or `Table` operators which no longer make sense, as they cause a `Pages` to contain an empty `Page` at rest. These were only used by tests, which I have rewritten to not require them, and benchmarks, which had already bitrotted so significantly as to not be worth maintaining. I removed benchmarks sufficient to get `cargo check --tests --benches` passing, but did not attempt to repair `cargo test --benches`, as that appears to have been broken prior to this change. --- crates/sats/src/proptest.rs | 6 + crates/table/benches/page_manager.rs | 215 ++------------------------- crates/table/src/eq.rs | 4 +- crates/table/src/page.rs | 7 + crates/table/src/pages.rs | 118 ++++++--------- crates/table/src/table.rs | 150 ++++++++++--------- 6 files changed, 154 insertions(+), 346 deletions(-) diff --git a/crates/sats/src/proptest.rs b/crates/sats/src/proptest.rs index 503d804f6a6..b291ae906ba 100644 --- a/crates/sats/src/proptest.rs +++ b/crates/sats/src/proptest.rs @@ -224,6 +224,12 @@ pub fn generate_typed_row() -> impl Strategy impl Strategy { + gen_with(generate_row_type(0..=SIZE), |ty| { + (generate_product_value(ty.clone()), generate_product_value(ty)) + }) +} + pub fn generate_typed_row_vec( size: impl Into, num_rows_min: usize, diff --git a/crates/table/benches/page_manager.rs b/crates/table/benches/page_manager.rs index 10fa708e8d2..f8f22178531 100644 --- a/crates/table/benches/page_manager.rs +++ b/crates/table/benches/page_manager.rs @@ -10,7 +10,7 @@ use rand::{Rng, SeedableRng}; use spacetimedb_lib::db::raw_def::v9::RawIndexAlgorithm; use spacetimedb_lib::db::raw_def::v9::RawModuleDefV9Builder; use spacetimedb_primitives::{ColList, IndexId, TableId}; -use spacetimedb_sats::layout::{row_size_for_bytes, row_size_for_type, Size}; +use spacetimedb_sats::layout::row_size_for_type; use spacetimedb_sats::raw_identifier::RawIdentifier; use spacetimedb_sats::{AlgebraicType, AlgebraicValue, ProductType, ProductValue}; use spacetimedb_schema::def::BTreeAlgorithm; @@ -18,7 +18,7 @@ use spacetimedb_schema::def::ModuleDef; use spacetimedb_schema::schema::TableSchema; use spacetimedb_schema::table_name::TableName; use spacetimedb_table::blob_store::NullBlobStore; -use spacetimedb_table::indexes::{Byte, Bytes, PageOffset, RowPointer, SquashedOffset, PAGE_DATA_SIZE}; +use spacetimedb_table::indexes::{Byte, Bytes, PageOffset, RowPointer, SquashedOffset}; use spacetimedb_table::page_pool::PagePool; use spacetimedb_table::pages::Pages; use spacetimedb_table::row_type_visitor::{row_type_visitor, VarLenVisitorProgram}; @@ -176,75 +176,8 @@ fn var_len_rows_per_page(data_size_in_bytes: usize) -> usize { PageOffset::PAGE_END.idx() / var_object_size } -fn reserve_empty_page(c: &mut Criterion) { - const RESERVE_SIZE: Size = row_size_for_bytes(8); - - let mut group = c.benchmark_group("reserve_empty_page"); - group.throughput(Throughput::Bytes(PAGE_DATA_SIZE as _)); - group.bench_function("leave_uninit", |b| { - let pool = PagePool::new_for_test(); - let mut pages = Pages::default(); - b.iter(|| { - let _ = black_box(pages.reserve_empty_page(&pool, RESERVE_SIZE)); - }); - }); - - let fill_with_zeros = |_, _, pages: &mut Pages| { - let pool = PagePool::new_for_test(); - let page = pages.reserve_empty_page(&pool, RESERVE_SIZE).unwrap(); - let page = pages.get_page_mut(page); - unsafe { page.zero_data() }; - }; - group.bench_function("fill_with_zeros", |b| { - iter_time_with(b, &mut Pages::default(), |_, _| (), fill_with_zeros) - }); -} - -fn insert_one_page_worth_fixed_len( - pool: &PagePool, - pages: &mut Pages, - visitor: &impl VarLenMembers, - val: &R, -) { - let size = row_size_for_type::(); - for _ in 0..rows_per_page::() { - let _ = black_box(unsafe { - black_box(&mut *pages).insert_row(pool, visitor, size, val.as_bytes(), &[], &mut NullBlobStore) - }); - } -} - type Group<'a, 'b> = &'a mut BenchmarkGroup<'b, WallTime>; -// time to insert a whole bunch of rows -fn insert_one_page_fixed_len(c: &mut Criterion) { - fn bench_insert_one_page_fixed_len(group: Group<'_, '_>, visitor: &impl VarLenMembers, name: &str) { - group.throughput(Throughput::Bytes( - rows_per_page::() as u64 * mem::size_of::() as u64, - )); - group.bench_function(name, |b| { - let pool = PagePool::new_for_test(); - let mut pages = Pages::default(); - // `0xa5` is the alternating bit pattern, which makes incorrect accesses obvious. - insert_one_page_worth_fixed_len(&pool, &mut pages, visitor, &R::from_u64(0xa5a5a5a5_a5a5a5a5)); - let pre = |_, pages: &mut Pages| pages.clear(); - iter_time_with(b, &mut pages, pre, |_, _, pages| { - insert_one_page_worth_fixed_len(&pool, pages, visitor, &R::from_u64(0xdeadbeef_0badbeef)) - }); - }); - } - - let mut group = c.benchmark_group("insert_one_page_fixed_len"); - bench_insert_one_page_fixed_len::(&mut group, &NullVarLenVisitor, "u64/NullVarLenVisitor"); - bench_insert_one_page_fixed_len::(&mut group, &u64::var_len_visitor(), "u64/VarLenVisitorProgram"); - - bench_insert_one_page_fixed_len::(&mut group, &NullVarLenVisitor, "U32x8/NullVarLenVisitor"); - bench_insert_one_page_fixed_len::(&mut group, &U32x8::var_len_visitor(), "U32x8/VarLenVisitorProgram"); - - bench_insert_one_page_fixed_len::(&mut group, &NullVarLenVisitor, "U32x64/NullVarLenVisitor"); - bench_insert_one_page_fixed_len::(&mut group, &U32x64::var_len_visitor(), "U32x64/VarLenVisitorProgram"); -} - fn fill_page_with_fixed_len_collect_row_pointers( pool: &PagePool, pages: &mut Pages, @@ -270,45 +203,6 @@ fn fill_page_with_fixed_len_collect_row_pointers( ptrs } -// insert a whole bunch of rows, then time to delete them all -fn delete_one_page_fixed_len(c: &mut Criterion) { - fn bench_delete_one_page_fixed_len(group: Group<'_, '_>, visitor: &impl VarLenMembers, name: &str) { - let rows_per_page = rows_per_page::(); - - group.throughput(Throughput::Bytes(rows_per_page as u64 * mem::size_of::() as u64)); - - group.bench_function(name, |b| { - let pre = |i, (pages, pool): &mut _| { - let val = R::from_u64(i); - fill_page_with_fixed_len_collect_row_pointers::(pool, pages, visitor, &val) - }; - iter_time_with( - b, - &mut (Pages::default(), PagePool::new_for_test()), - pre, - |ptrs, _, (pages, _)| { - for ptr in ptrs { - unsafe { - pages.delete_row(visitor, row_size_for_type::(), black_box(ptr), &mut NullBlobStore) - }; - } - }, - ); - }); - } - - let mut group = c.benchmark_group("delete_one_page_fixed_len"); - - bench_delete_one_page_fixed_len::(&mut group, &NullVarLenVisitor, "u64/NullVarLenVisitor"); - bench_delete_one_page_fixed_len::(&mut group, &u64::var_len_visitor(), "u64/VarLenVisitorProgram"); - - bench_delete_one_page_fixed_len::(&mut group, &NullVarLenVisitor, "U32x8/NullVarLenVisitor"); - bench_delete_one_page_fixed_len::(&mut group, &U32x8::var_len_visitor(), "U32x8/VarLenVisitorProgram"); - - bench_delete_one_page_fixed_len::(&mut group, &NullVarLenVisitor, "U32x64/NullVarLenVisitor"); - bench_delete_one_page_fixed_len::(&mut group, &U32x64::var_len_visitor(), "U32x64/VarLenVisitorProgram"); -} - // insert a whole bunch of rows, then time to access them fn retrieve_one_page_fixed_len(c: &mut Criterion) { fn bench_retrieve_one_page(group: Group<'_, '_>, visitor: &impl VarLenMembers, name: &str) { @@ -349,85 +243,6 @@ fn retrieve_one_page_fixed_len(c: &mut Criterion) { bench_retrieve_one_page::(&mut group, &U32x64::var_len_visitor(), "U32x64/VarLenVisitorProgram"); } -// insert a bunch of rows, -// delete some fraction of them to create holes in multiple pages, -// then time to insert into those holes -fn insert_with_holes_fixed_len(c: &mut Criterion) { - fn bench_insert_with_holes(c: &mut Criterion, var_len_visitor: &impl VarLenMembers, name: &str) { - let mut group = c.benchmark_group(format!("insert_with_holes_fixed_len/{name}")); - let val = R::from_u64(0xdeadbeef_0badbeef); - for delete_ratio in [0.1f64, 0.25, 0.5, 0.75, 0.9, 1.0] { - let num_pages = 16; - let total_num_rows = rows_per_page::() * num_pages; - let num_to_delete = (total_num_rows as f64 * delete_ratio) as usize; - - let num_to_delete_in_bytes = num_to_delete * mem::size_of::(); - - let row_size = row_size_for_type::(); - - group.throughput(Throughput::Bytes(num_to_delete_in_bytes as u64)); - - group.bench_function(delete_ratio.to_string(), |b| { - let pool = PagePool::new_for_test(); - let mut pages = Pages::default(); - - let mut rng = StdRng::seed_from_u64(0xa5a5a5a5_a5a5a5a5); - - for _ in 0..num_pages { - let page = pages.reserve_empty_page(&pool, row_size).unwrap(); - let page = pages.get_page_mut(page); - - unsafe { page.zero_data() }; - } - - let pre = |_, (pages, pool): &mut (Pages, PagePool)| { - pages.clear(); - let mut ptrs_to_delete = Vec::with_capacity(num_to_delete); - for _ in 0..total_num_rows { - let (page_idx, offset) = unsafe { - pages.insert_row(pool, var_len_visitor, row_size, val.as_bytes(), &[], &mut NullBlobStore) - } - .unwrap(); - - if rng.random_bool(delete_ratio) { - ptrs_to_delete.push(RowPointer::new( - false, - page_idx, - offset, - SquashedOffset::COMMITTED_STATE, - )); - } - } - let actual_num_deleted = ptrs_to_delete.len(); - for ptr in ptrs_to_delete { - unsafe { - pages.delete_row(var_len_visitor, row_size, ptr, &mut NullBlobStore); - } - } - actual_num_deleted - }; - let body = |actual_num_deleted, _, (pages, pool): &mut (Pages, PagePool)| { - for _ in 0..actual_num_deleted { - let _ = black_box(unsafe { - pages.insert_row(pool, var_len_visitor, row_size, val.as_bytes(), &[], &mut NullBlobStore) - }); - } - }; - iter_time_with(b, &mut (pages, pool), pre, body); - }); - } - } - - bench_insert_with_holes::(c, &NullVarLenVisitor, "u64/NullVarLenVisitor"); - bench_insert_with_holes::(c, &u64::var_len_visitor(), "u64/VarLenVisitorProgram"); - - bench_insert_with_holes::(c, &NullVarLenVisitor, "U32x8/NullVarLenVisitor"); - bench_insert_with_holes::(c, &U32x8::var_len_visitor(), "U32x8/VarLenVisitorProgram"); - - bench_insert_with_holes::(c, &NullVarLenVisitor, "U32x64/NullVarLenVisitor"); - bench_insert_with_holes::(c, &U32x64::var_len_visitor(), "U32x64/VarLenVisitorProgram"); -} - // insert a whole bunch of rows, then time to copy_filter materialize a view fn copy_filter_fixed_len(c: &mut Criterion) { @@ -480,15 +295,7 @@ fn copy_filter_fixed_len(c: &mut Criterion) { // - In the insert-with-holes benchmark, randomize size of each row to simulate fragmentation. // - Extend above benchmarks to go through `Table` with `AlgebraicValue`. -criterion_group!( - pages, - reserve_empty_page, - insert_one_page_fixed_len, - delete_one_page_fixed_len, - retrieve_one_page_fixed_len, - insert_with_holes_fixed_len, - copy_filter_fixed_len, -); +criterion_group!(pages, retrieve_one_page_fixed_len, copy_filter_fixed_len,); fn schema_from_ty(ty: ProductType, name: &str) -> TableSchema { let mut result = TableSchema::from_product_type(ty); @@ -542,7 +349,7 @@ fn table_insert_one_row(c: &mut Criterion) { let mut ctx = (table, NullBlobStore); let ptr = ctx.0.insert(&pool, &mut ctx.1, &val).unwrap().1.pointer(); let pre = |_, (table, bs): &mut (Table, NullBlobStore)| { - table.delete(bs, ptr, |_| ()).unwrap(); + table.delete(&pool, bs, ptr, |_| ()).unwrap(); }; group.bench_function(name, |b| { iter_time_with(b, &mut ctx, pre, |_, _, (table, bs)| { @@ -595,8 +402,8 @@ fn table_delete_one_row(c: &mut Criterion) { }; group.bench_function(name, |b| { - iter_time_with(b, &mut ctx, insert, |row, _, (table, bs, _)| { - table.delete(bs, row, |_| ()) + iter_time_with(b, &mut ctx, insert, |row, _, (table, bs, pool)| { + table.delete(pool, bs, row, |_| ()) }); }); } @@ -797,13 +604,13 @@ fn insert_num_same( .flatten() } -fn clear_all_same(tbl: &mut Table, index_id: IndexId, val_same: u64) { +fn clear_all_same(pool: &PagePool, tbl: &mut Table, index_id: IndexId, val_same: u64) { let index = tbl.get_index_by_id(index_id).unwrap(); let key = R::column_value_from_u64(val_same); let key = index.key_from_algebraic_value(&key); let ptrs = index.seek_point(&key).collect::>(); for ptr in ptrs { - tbl.delete(&mut NullBlobStore, ptr, |_| ()).unwrap(); + tbl.delete(pool, &mut NullBlobStore, ptr, |_| ()).unwrap(); } } @@ -857,7 +664,7 @@ fn index_insert(c: &mut Criterion) { &num_rows, |b, &num_rows| { let pre = |_, (tbl, _, pool): &mut (Table, NullBlobStore, PagePool)| { - clear_all_same::(tbl, index_id, num_rows); + clear_all_same::(pool, tbl, index_id, num_rows); insert_num_same(pool, tbl, || make_row(num_rows), num_same - 1); make_row(num_rows).to_product() }; @@ -980,11 +787,11 @@ fn index_delete(c: &mut Criterion) { &num_rows, |b, &num_rows| { let pre = |_, tbl: &mut Table| { - clear_all_same::(tbl, index_id, num_rows); + clear_all_same::(&pool, tbl, index_id, num_rows); insert_num_same(&pool, tbl, || make_row(num_rows), num_same).unwrap() }; iter_time_with(b, &mut tbl, pre, |ptr, _, tbl| { - tbl.delete(&mut NullBlobStore, ptr, |_| ()) + tbl.delete(&pool, &mut NullBlobStore, ptr, |_| ()) }); }, ); diff --git a/crates/table/src/eq.rs b/crates/table/src/eq.rs index 173f07a8cce..51b6ab3c83b 100644 --- a/crates/table/src/eq.rs +++ b/crates/table/src/eq.rs @@ -250,14 +250,14 @@ mod test { let a0 = product![AlgebraicValue::sum(0, u64::MAX.into())]; let (_, a0_rr) = table_a.insert(&pool, bs, &a0).unwrap(); let a0_ptr = a0_rr.pointer(); - assert!(table_a.delete(bs, a0_ptr, |_| {}).is_some()); + assert!(table_a.delete(&pool, bs, a0_ptr, |_| {}).is_some()); // Insert u64::ALTERNATING_BIT_PATTERN with tag 0 and then delete it. let b0 = 0b01010101_01010101_01010101_01010101_01010101_01010101_01010101_01010101u64; let b0 = product![AlgebraicValue::sum(0, b0.into())]; let (_, b0_rr) = table_b.insert(&pool, bs, &b0).unwrap(); let b0_ptr = b0_rr.pointer(); - assert!(table_b.delete(bs, b0_ptr, |_| {}).is_some()); + assert!(table_b.delete(&pool, bs, b0_ptr, |_| {}).is_some()); // Insert two identical rows `a1` and `b2` into the tables. // They should occupy the spaces of the previous rows. diff --git a/crates/table/src/page.rs b/crates/table/src/page.rs index cfb66d97f10..3ef134da395 100644 --- a/crates/table/src/page.rs +++ b/crates/table/src/page.rs @@ -1185,6 +1185,13 @@ impl Page { self.header.fixed.num_rows as usize } + /// Is this page empty, that is, does it contain zero rows? + /// + /// This method runs in constant time. + pub fn is_empty(&self) -> bool { + self.num_rows() == 0 + } + #[cfg(test)] /// Use this page's present rows bitvec to compute the number of present rows. /// diff --git a/crates/table/src/pages.rs b/crates/table/src/pages.rs index 5949419deea..ce475778044 100644 --- a/crates/table/src/pages.rs +++ b/crates/table/src/pages.rs @@ -155,50 +155,40 @@ impl Pages { .expect("pages len to be greater than number of free slots") } - #[cfg(test)] - pub fn free_empty_page(&mut self, page_index: PageIndex) { + /// Free the page at `page_index`, leaving `self.pages[page_index.idx()]` as `None`. + /// + /// Prior to calling this method, `self.pages[page_index.idx()]` must be: + /// - Present, i.e. `Some` and not `None`. + /// - Empty, i.e. have `page.num_rows() == 0`. + /// + /// Prior to calling this method, the page's entry should already have been deleted from `self.non_full_pages`. + /// + /// Calling when in violation of these invariants may result in a panic and/or unexpected behavior. + /// If a non-empty page is freed, + /// future `unsafe` operations which attempt to read from rows that were previously in that page + /// may result in Undefined Behavior. + /// + /// The freed page will be returned to the `pool`. + fn free_empty_page(&mut self, pool: &PagePool, page_index: PageIndex) { let page = self.get(page_index).expect("page to free to have been present"); - assert_eq!(page.num_rows(), 0); + debug_assert_eq!(page.num_rows(), 0); + + // The page should already have been removed from `self.non_full_pages`. + debug_assert!({ + let free_granules = page.available_var_len_granules(); - let free_granules = page.available_var_len_granules(); + !self.non_full_pages.remove(&(free_granules, page_index)) + }); - let removed_from_non_full = self.non_full_pages.remove(&(free_granules, page_index)); - assert!(removed_from_non_full); + let page = std::mem::replace(&mut self.pages[page_index.idx()], None) + .expect("freed page to have been present after we already checked its presence"); - self.pages[page_index.idx()] = None; + pool.put(page); let newly_inserted_into_free_pages_set = self.free_page_slots.insert(page_index); - assert!(newly_inserted_into_free_pages_set) - } - - /// Make all pages within `self` clear, - /// deleting all rows. - // - // TODO(delete-free-page): Determine what to do with this method. - // It doesn't really make sense given that it clears all pages but doesn't delete them, - // but it's only used in benchmarks, and those benchmarks are using it specifically to bypass allocating new pages. - #[doc(hidden)] // Used in benchmarks. - pub fn clear(&mut self) { - // Clear every page. - for page in &mut self.pages { - if let Some(page) = page { - page.clear(); - } - } - // Mark every page non-full. - self.non_full_pages = (0..self.pages.len()) - // We could probably compute the number of available granules once and use it for all pages, - // rather than calling the method on each page, - // but we'd have to do some amount of reasoning to demonstrate it was correct - // based on the definition of `Page::clear`, - // and why bother? - .filter_map(|idx| { - let idx = PageIndex(idx as u64); - self.get(idx).map(|page| (page.available_var_len_granules(), idx)) - }) - .collect(); + debug_assert!(newly_inserted_into_free_pages_set) } /// Get a reference to fixed-len row data. @@ -235,16 +225,6 @@ impl Pages { } } - /// Reserve a new, initially empty page. - // TODO(delete-free-page): Determine what to do with this method. - // It doesn't really make sense in a world where `Pages` doesn't contain empty `Page`s, - // but it's only used for tests and benches. - pub fn reserve_empty_page(&mut self, pool: &PagePool, fixed_row_size: Size) -> Result { - let idx = self.allocate_new_page(pool, fixed_row_size)?; - self.record_page_non_full(idx, fixed_row_size); - Ok(idx) - } - /// Call `f` with a reference to a page which satisfies /// `page.has_space_for_row(fixed_row_size, num_var_len_granules)`. pub fn with_page_to_insert_row( @@ -356,11 +336,12 @@ impl Pages { var_len_visitor: &impl VarLenMembers, fixed_row_size: Size, row_ptr: RowPointer, + page_pool: &PagePool, blob_store: &mut dyn BlobStore, ) -> BlobNumBytes { let page_index = row_ptr.page_index(); - self.with_updating_non_full_pages(page_index, fixed_row_size, |this| { + self.with_updating_non_full_pages_and_maybe_freeing_page(page_pool, page_index, fixed_row_size, |this| { let page = this .get_mut(page_index) .expect("page containing `row_ptr` with validity safety invariant to be present"); @@ -377,60 +358,51 @@ impl Pages { /// Collect information about the page `self[page_index]` sufficient to update [`Self::non_full_pages`], /// then run `body` to update the page, and finally update [`Self::non_full_pages`] for its new fullness and capacity. + /// Free the page and return it to the `pool` if it is empty after the `body`. /// /// `body` should not update any pages other than the one identified by `page_index`. /// /// If `page_index` does not refer to a present page, i.e. `self.pages[page_index.idx]` is `None`, /// this method will panic. - fn with_updating_non_full_pages( + fn with_updating_non_full_pages_and_maybe_freeing_page( &mut self, + pool: &PagePool, page_index: PageIndex, fixed_row_size: Size, body: impl FnOnce(&mut Self) -> Ret, ) -> Ret { let page = self .get(page_index) - .expect("page to be present in `with_updating_non_full_pages`"); + .expect("page to be present in `with_updating_non_full_pages_and_maybe_freeing_page`"); let full_before = page.is_full(fixed_row_size); let available_granules_before = page.available_var_len_granules(); let ret = body(self); - self.update_page_non_full(available_granules_before, full_before, page_index, fixed_row_size); + let page = self + .get(page_index) + .expect("page to remain present in `with_updating_non_full_pages_and_maybe_freeing_page` after it was checked earlier"); + + if page.is_empty() { + self.remove_non_full_marker(available_granules_before, full_before, page_index); + self.free_empty_page(pool, page_index); + } else { + self.remove_non_full_marker(available_granules_before, full_before, page_index); + self.record_page_non_full(page_index, fixed_row_size); + } ret } - /// Update [`Self::non_full_pages`] to change the number of var-len granules available in the page at `self[page_index]`, - /// first deleting any old entry and then re-inserting the new entry. - /// - /// The entry for `page` in `self.non_full_granules` should not have been deleted prior to calling this method. - /// If the entry has already been deleted or was never present, instead use [`Self::record_page_non_full`]. - /// - /// `available_granules_before` should be the previous count from [`Page::available_var_len_granules`], - /// prior to whatever operation made space available in the page. - /// This is necessary because `non_full_pages` is a `BTreeSet` sorted by `(available_granules, page_index)`, - /// so locating the `page_index` without the `available_granules` would be slow. - /// - /// `full_before` should be the result of [`Page::is_full`] prior to whatever operation made space available in the page. - /// This is necessary because `non_full_pages` does not store full pages (as the name implies), - /// so we should not attempt to delete the previous entry if the page was previously full. - fn update_page_non_full( - &mut self, - available_granules_before: usize, - full_before: bool, - page_index: PageIndex, - fixed_row_size: Size, - ) { + /// Remove a page's pre-existing entry in `self.non_full_pages`. + fn remove_non_full_marker(&mut self, available_granules_before: usize, full_before: bool, page_index: PageIndex) { if full_before { debug_assert!(!self.non_full_pages.remove(&(available_granules_before, page_index))); } else { let _prev = self.non_full_pages.remove(&(available_granules_before, page_index)); debug_assert!(_prev); } - - self.record_page_non_full(page_index, fixed_row_size); } /// Record the number of available var-len granules in the page at `self[page_index]` into [`Self::non_full_pages`]. diff --git a/crates/table/src/table.rs b/crates/table/src/table.rs index 75c864d5e49..3d2bce7ef0b 100644 --- a/crates/table/src/table.rs +++ b/crates/table/src/table.rs @@ -627,7 +627,7 @@ impl Table { // where `insert` is called, we are not dealing with transactions, // and we already know there cannot be a duplicate row error, // but we check just in case it isn't. - let (hash, row_ptr) = unsafe { self.confirm_insertion::(blob_store, row_ptr, blob_bytes) }?; + let (hash, row_ptr) = unsafe { self.confirm_insertion::(pool, blob_store, row_ptr, blob_bytes) }?; // SAFETY: Per post-condition of `confirm_insertion`, `row_ptr` refers to a valid row. let row_ref = unsafe { self.get_row_ref_unchecked(blob_store, row_ptr) }; Ok((hash, row_ref)) @@ -829,14 +829,15 @@ impl Table { /// `self.is_row_present(row)` must hold. pub unsafe fn confirm_insertion<'a, const CHECK_SAME_ROW: bool>( &'a mut self, + page_pool: &PagePool, blob_store: &'a mut dyn BlobStore, ptr: RowPointer, blob_bytes: BlobNumBytes, ) -> Result<(Option, RowPointer), InsertError> { // SAFETY: Caller promised that `self.is_row_present(ptr)` holds. - let hash = unsafe { self.insert_into_pointer_map(blob_store, ptr) }?; + let hash = unsafe { self.insert_into_pointer_map(page_pool, blob_store, ptr) }?; // SAFETY: Caller promised that `self.is_row_present(ptr)` holds. - unsafe { self.insert_into_indices::(blob_store, ptr) }?; + unsafe { self.insert_into_indices::(page_pool, blob_store, ptr) }?; self.update_statistics_added_row(blob_bytes); Ok((hash, ptr)) @@ -857,6 +858,7 @@ impl Table { /// `self.is_row_present(new_row)` and `self.is_row_present(old_row)` must hold. pub unsafe fn confirm_update<'a>( &'a mut self, + page_pool: &PagePool, blob_store: &'a mut dyn BlobStore, new_ptr: RowPointer, old_ptr: RowPointer, @@ -868,17 +870,17 @@ impl Table { // Insert new row into indices. // SAFETY: Caller promised that `self.is_row_present(ptr)` holds. - let res = unsafe { self.insert_into_indices::(blob_store, new_ptr) }; + let res = unsafe { self.insert_into_indices::(page_pool, blob_store, new_ptr) }; if let Err(e) = res { // Undo (1). - unsafe { self.insert_into_indices::(blob_store, old_ptr) } + unsafe { self.insert_into_indices::(page_pool, blob_store, old_ptr) } .expect("re-inserting the old row into indices should always work"); return Err(e); } // Remove the old row physically. // SAFETY: The physical `old_ptr` still exists. - let blob_bytes_removed = unsafe { self.delete_internal_skip_pointer_map(blob_store, old_ptr) }; + let blob_bytes_removed = unsafe { self.delete_internal_skip_pointer_map(page_pool, blob_store, old_ptr) }; self.update_statistics_deleted_row(blob_bytes_removed); // Update statistics. @@ -912,6 +914,7 @@ impl Table { /// Post-condition: If this method returns `Ok(_)`, the row still exists. unsafe fn insert_into_indices<'a, const CHECK_SAME_ROW: bool>( &'a mut self, + pool: &PagePool, blob_store: &'a mut dyn BlobStore, new: RowPointer, ) -> Result<(), InsertError> { @@ -955,7 +958,7 @@ impl Table { // Cleanup, undo the row insertion of `new`s. // SAFETY: We just inserted `new`, so it must be present. - unsafe { self.delete_internal(blob_store, new) }; + unsafe { self.delete_internal(pool, blob_store, new) }; error }) @@ -1018,6 +1021,7 @@ impl Table { /// Post-condition: If this method returns `Ok(_)`, the row still exists. unsafe fn insert_into_pointer_map<'a>( &'a mut self, + page_pool: &PagePool, blob_store: &'a mut dyn BlobStore, ptr: RowPointer, ) -> Result, DuplicateError> { @@ -1039,7 +1043,7 @@ impl Table { unsafe { self.inner .pages - .delete_row(&self.inner.visitor_prog, self.row_size(), ptr, blob_store) + .delete_row(&self.inner.visitor_prog, self.row_size(), ptr, page_pool, blob_store) }; return Err(DuplicateError(existing_row)); } @@ -1215,6 +1219,7 @@ impl Table { /// `ptr` must point to a valid, live row in this table. pub unsafe fn delete_internal_skip_pointer_map( &mut self, + page_pool: &PagePool, blob_store: &mut dyn BlobStore, ptr: RowPointer, ) -> BlobNumBytes { @@ -1228,7 +1233,7 @@ impl Table { unsafe { self.inner .pages - .delete_row(&self.inner.visitor_prog, self.row_size(), ptr, blob_store) + .delete_row(&self.inner.visitor_prog, self.row_size(), ptr, page_pool, blob_store) } } @@ -1240,7 +1245,12 @@ impl Table { /// Use `delete_unchecked` or `delete` to delete a row with index updating. /// /// SAFETY: `self.is_row_present(row)` must hold. - unsafe fn delete_internal(&mut self, blob_store: &mut dyn BlobStore, ptr: RowPointer) -> BlobNumBytes { + unsafe fn delete_internal( + &mut self, + page_pool: &PagePool, + blob_store: &mut dyn BlobStore, + ptr: RowPointer, + ) -> BlobNumBytes { // Remove the set semantic association. if let Some(pointer_map) = &mut self.pointer_map { // SAFETY: `self.is_row_present(row)` holds. @@ -1252,7 +1262,7 @@ impl Table { // Delete the physical row. // SAFETY: `ptr` points to a valid row in this table as `self.is_row_present(row)` holds. - unsafe { self.delete_internal_skip_pointer_map(blob_store, ptr) } + unsafe { self.delete_internal_skip_pointer_map(page_pool, blob_store, ptr) } } /// Deletes the row identified by `ptr` from the table. @@ -1260,7 +1270,7 @@ impl Table { /// This method does update statistics. /// /// SAFETY: `self.is_row_present(row)` must hold. - unsafe fn delete_unchecked(&mut self, blob_store: &mut dyn BlobStore, ptr: RowPointer) { + unsafe fn delete_unchecked(&mut self, page_pool: &PagePool, blob_store: &mut dyn BlobStore, ptr: RowPointer) { // Delete row from indices. // Do this before the actual deletion, as `index.delete` needs a `RowRef` // so it can extract the appropriate value. @@ -1268,7 +1278,7 @@ impl Table { unsafe { self.delete_from_indices(blob_store, ptr) }; // SAFETY: Caller promised that `self.is_row_present(row)` holds. - let blob_bytes_deleted = unsafe { self.delete_internal(blob_store, ptr) }; + let blob_bytes_deleted = unsafe { self.delete_internal(page_pool, blob_store, ptr) }; self.update_statistics_deleted_row(blob_bytes_deleted); } @@ -1310,6 +1320,7 @@ impl Table { /// so that the resulting `ProductValue`s can be passed to the subscription evaluator. pub fn delete<'a, R>( &'a mut self, + page_pool: &PagePool, blob_store: &'a mut dyn BlobStore, ptr: RowPointer, before: impl for<'b> FnOnce(RowRef<'b>) -> R, @@ -1324,7 +1335,7 @@ impl Table { let ret = before(row_ref); // SAFETY: We've checked above that `self.is_row_present(ptr)`. - unsafe { self.delete_unchecked(blob_store, ptr) }; + unsafe { self.delete_unchecked(page_pool, blob_store, ptr) }; Some(ret) } @@ -1365,29 +1376,29 @@ impl Table { // If an equal row was present, delete it. if let Some(existing_row_ptr) = existing_row_ptr { // SAFETY: `find_same_row` ensures that the pointer is valid. - unsafe { self.delete_unchecked(blob_store, existing_row_ptr) }; + unsafe { self.delete_unchecked(pool, blob_store, existing_row_ptr) }; } // Remove the temporary row we inserted in the beginning. // Avoid the pointer map, since we don't want to delete it twice. // SAFETY: `ptr` is valid as we just inserted it. unsafe { - self.delete_internal_skip_pointer_map(blob_store, temp_ptr); + self.delete_internal_skip_pointer_map(pool, blob_store, temp_ptr); } Ok(existing_row_ptr) } - /// Clears this table, removing all present rows from it. - pub fn clear(&mut self, blob_store: &mut dyn BlobStore) -> u64 { - let ptrs = self.scan_all_row_ptrs(); - let len = ptrs.len() as u64; - for ptr in ptrs { - // SAFETY: `ptr` came rom `self.scan_rows(...)`, so it's present. - unsafe { self.delete_unchecked(blob_store, ptr) }; - } - len - } + // /// Clears this table, removing all present rows from it. + // pub fn clear(&mut self, pool: &PagePool, blob_store: &mut dyn BlobStore) -> u64 { + // let ptrs = self.scan_all_row_ptrs(); + // let len = ptrs.len() as u64; + // for ptr in ptrs { + // // SAFETY: `ptr` came rom `self.scan_rows(...)`, so it's present. + // unsafe { self.delete_unchecked(pool, blob_store, ptr) }; + // } + // len + // } /// Returns the row type for rows in this table. pub fn get_row_type(&self) -> &ProductType { @@ -1693,6 +1704,13 @@ impl Table { /// /// The schema of rows stored in the `pages` must exactly match `self.schema` and `self.inner.row_layout`. pub unsafe fn set_pages(&mut self, pages: Vec>>, blob_store: &dyn BlobStore) { + // If the snapshot contains any present but empty pages, remove them and replace them with gaps. + // This may be the case for historical snapshots which predate our support for empty pages. + let pages = pages + .into_iter() + .map(|page| page.and_then(|page| (!page.is_empty()).then_some(page))) + .collect(); + self.inner.pages.set_contents(pages, self.inner.row_layout.size()); // Recompute table metadata based on the new pages. @@ -2485,11 +2503,6 @@ impl Table { &self.inner.pages } - #[cfg(test)] - fn pages_mut(&mut self) -> &mut Pages { - &mut self.inner.pages - } - /// Iterates over each [`Page`] in this table, ensuring that its hash is computed before yielding it. /// /// Used when capturing a snapshot. @@ -2554,7 +2567,7 @@ pub(crate) mod test { use spacetimedb_lib::db::raw_def::v9::{btree, RawModuleDefV9Builder}; use spacetimedb_primitives::TableId; use spacetimedb_sats::bsatn::to_vec; - use spacetimedb_sats::proptest::{generate_typed_row, generate_typed_row_vec, SIZE}; + use spacetimedb_sats::proptest::{generate_two_typed_rows, generate_typed_row, generate_typed_row_vec, SIZE}; use spacetimedb_sats::{product, AlgebraicType, ArrayValue}; use spacetimedb_schema::def::{BTreeAlgorithm, ModuleDef}; use spacetimedb_schema::schema::Schema as _; @@ -2589,21 +2602,24 @@ pub(crate) mod test { let mut table = Table::new(schema.into(), SquashedOffset::COMMITTED_STATE); let pool = PagePool::new_for_test(); + let blob_store = &mut NullBlobStore; let cols = ColList::new(0.into()); let algo = BTreeAlgorithm { columns: cols.clone() }.into(); let index = table.new_index(&algo, true).unwrap(); // SAFETY: Index was derived from `table`. - unsafe { table.insert_index(&NullBlobStore, index_schema.index_id, index) }.unwrap(); + unsafe { table.insert_index(blob_store, index_schema.index_id, index) }.unwrap(); // Reserve a page so that we can check the hash. - let pi = table.inner.pages.reserve_empty_page(&pool, table.row_size()).unwrap(); + let (_, row_ref) = table.insert(&pool, blob_store, &product![i32::MAX, i32::MAX]).unwrap(); + let pi = row_ref.pointer().page_index(); + let hash_pre_ins = hash_unmodified_save_get(table.inner.pages.get_mut(pi).expect("reserved page to be present")); // Insert the row (0, 0). table - .insert(&pool, &mut NullBlobStore, &product![0i32, 0i32]) + .insert(&pool, blob_store, &product![0i32, 0i32]) .expect("Initial insert failed"); // Inserting cleared the hash. @@ -2612,7 +2628,7 @@ pub(crate) mod test { assert_ne!(hash_pre_ins, hash_post_ins); // Try to insert the row (0, 1), and assert that we get the expected error. - match table.insert(&pool, &mut NullBlobStore, &product![0i32, 1i32]) { + match table.insert(&pool, blob_store, &product![0i32, 1i32]) { Ok(_) => panic!("Second insert with same unique value succeeded"), Err(InsertError::IndexError(UniqueConstraintViolation { constraint_name, @@ -2772,10 +2788,23 @@ pub(crate) mod test { } #[test] - fn insert_delete_removed_from_pointer_map((ty, val) in generate_typed_row()) { + fn insert_delete_removed_from_pointer_map((ty, (anchor, val)) in generate_two_typed_rows()) { + // We need two distinct values so that one can remain resident in the page + // while the other gets inserted and deleted. + // Without the "anchor" remaining present, the page will be deleted, + // so we'll be unable to make assertions about its hash. + prop_assume!(anchor != val); + let pool = PagePool::new_for_test(); let mut blob_store = HashMapBlobStore::default(); let mut table = table(ty); + + let (anchor_hash, row) = table.insert(&pool, &mut blob_store, &anchor).unwrap(); + let anchor_hash = anchor_hash.unwrap(); + prop_assert_eq!(row.row_hash(), anchor_hash); + let anchor_ptr = row.pointer(); + prop_assert_eq!(table.pointers_for(anchor_hash), &[anchor_ptr]); + let (hash, row) = table.insert(&pool, &mut blob_store, &val).unwrap(); let hash = hash.unwrap(); prop_assert_eq!(row.row_hash(), hash); @@ -2783,26 +2812,26 @@ pub(crate) mod test { prop_assert_eq!(table.pointers_for(hash), &[ptr]); prop_assert_eq!(table.inner.pages.len(), 1); - prop_assert_eq!(table.inner.pages.get(PageIndex(0)).map(|page| page.num_rows()), Some(1)); - prop_assert_eq!(&table.scan_rows(&blob_store).map(|r| r.pointer()).collect::>(), &[ptr]); - prop_assert_eq!(table.row_count, 1); + prop_assert_eq!(table.inner.pages.get(PageIndex(0)).map(|page| page.num_rows()), Some(2)); + prop_assert_eq!(&table.scan_rows(&blob_store).map(|r| r.pointer()).collect::>(), &[anchor_ptr, ptr]); + prop_assert_eq!(table.row_count, 2); let hash_pre_del = hash_unmodified_save_get(table.inner.pages.get_mut(ptr.page_index()).expect("page containing row to be present")); - table.delete(&mut blob_store, ptr, |_| ()); + table.delete(&pool, &mut blob_store, ptr, |_| ()); // FIXME(delete-free-page): Page will no longer be present here after deleting its only row. // Amend this test to insert two rows and only delete one of them. - let hash_post_del = hash_unmodified_save_get(table.inner.pages.get_mut(ptr.page_index()).expect("page to remain present after delete, until we move to freeing empty pages")); + let hash_post_del = hash_unmodified_save_get(table.inner.pages.get_mut(ptr.page_index()).expect("page to remain present after delete because anchor is still present")); assert_ne!(hash_pre_del, hash_post_del); prop_assert_eq!(table.pointers_for(hash), &[]); prop_assert_eq!(table.inner.pages.len(), 1); - prop_assert_eq!(table.inner.pages.get(PageIndex(0)).map(|page| page.num_rows()), Some(0)); - prop_assert_eq!(table.row_count, 0); + prop_assert_eq!(table.inner.pages.get(PageIndex(0)).map(|page| page.num_rows()), Some(1)); + prop_assert_eq!(table.row_count, 1); - prop_assert!(&table.scan_rows(&blob_store).next().is_none()); + prop_assert_eq!(&table.scan_rows(&blob_store).map(|r| r.pointer()).collect::>(), &[anchor_ptr]); } #[test] @@ -2903,7 +2932,7 @@ pub(crate) mod test { } if let Some(row_ptr) = inserted_row_ptrs.pop() { - table.delete(&mut blob_store, row_ptr, |_| ()); + table.delete(&pool, &mut blob_store, row_ptr, |_| ()); table.inner.pages.assert_non_full_pages_consistent(table.inner.row_layout.size()); } } @@ -2943,7 +2972,7 @@ pub(crate) mod test { // Confirm the insertion, checking any constraints, removing the physical row on error. // SAFETY: We just inserted `ptr`, so it must be present. - let (hash, row_ptr) = unsafe { table.confirm_insertion::(blob_store, row_ptr, blob_bytes) }?; + let (hash, row_ptr) = unsafe { table.confirm_insertion::(&pool, blob_store, row_ptr, blob_bytes) }?; // SAFETY: Per post-condition of `confirm_insertion`, `row_ptr` refers to a valid row. let row_ref = unsafe { table.get_row_ref_unchecked(blob_store, row_ptr) }; Ok((hash, row_ref)) @@ -3012,7 +3041,7 @@ pub(crate) mod test { assert!(second_page.available_var_len_granules() > 0); let first_ptr = inserted_ptrs[0]; - table.delete(&mut blob_store, first_ptr, |_| ()); + table.delete(&pool, &mut blob_store, first_ptr, |_| ()); let first_page = &table .inner @@ -3096,15 +3125,15 @@ pub(crate) mod test { assert_eq!(table2.blob_store_bytes, BLOB_OBJ_LEN); // Delete `short_str` row. This should not affect the byte count. - table1.delete(blob_store, short_row_ptr, |_| ()).unwrap(); + table1.delete(&pool, blob_store, short_row_ptr, |_| ()).unwrap(); assert_eq!(table1.blob_store_bytes, BLOB_OBJ_LEN_2X); // Delete the first long string row. This gets us down to `BLOB_OBJ_LEN` (we had 2x before). - table1.delete(blob_store, long_row_ptr, |_| ()).unwrap(); + table1.delete(&pool, blob_store, long_row_ptr, |_| ()).unwrap(); assert_eq!(table1.blob_store_bytes, BLOB_OBJ_LEN); // Delete the first long string row. This gets us down to 0 (we've now deleted 2x). - table1.delete(blob_store, long_row_ptr2, |_| ()).unwrap(); + table1.delete(&pool, blob_store, long_row_ptr2, |_| ()).unwrap(); assert_eq!(table1.blob_store_bytes, 0.into()); } @@ -3161,12 +3190,10 @@ pub(crate) mod test { // Delete all the rows from page 0 before freeing it, to avoid leaving stale entries in the pointer map. for ptr in row_ptrs.into_iter().take_while(|ptr| ptr.page_index() == PageIndex(0)) { table - .delete(blob_store, ptr, |_| ()) + .delete(&pool, blob_store, ptr, |_| ()) .expect("deleted row to have been present"); } - table.pages_mut().free_empty_page(PageIndex(0)); - let scanned_rows = table .scan_rows(blob_store) .map(|row| row.read_col::(0).unwrap()) @@ -3209,28 +3236,17 @@ pub(crate) mod test { // Delete all the rows from page 0 before freeing it, to avoid leaving stale entries in the pointer map. for ptr in row_ptrs.into_iter().take_while(|ptr| ptr.page_index() == PageIndex(0)) { table - .delete(blob_store, ptr, |_| ()) + .delete(&pool, blob_store, ptr, |_| ()) .expect("deleted row to have been present"); } - assert_eq!( - table - .pages() - .get(PageIndex(0)) - .expect("empty page to still be present") - .num_rows(), - 0 - ); + assert!(table.pages().get(PageIndex(0)).is_none()); assert!(table .pages() .get(PageIndex(1)) .expect("full page to still be present") .is_full(table.inner.row_layout.size())); - assert_eq!(table.num_pages(), 2); - - table.pages_mut().free_empty_page(PageIndex(0)); - assert_eq!(table.num_pages(), 1); assert_eq!(table.pages().len(), 2); From 592ea3bf68ac1b8417a7382a99a0eba877149827 Mon Sep 17 00:00:00 2001 From: Phoebe Goldman Date: Fri, 21 Aug 2026 10:43:15 -0400 Subject: [PATCH 03/10] clippy, and fix (remove) broken benchmarks As discussed in the follow-up PR, the benchmarks in the `table` crate appear to have bit-rotted significantly, and fail when run under `cargo test --benches`. I have not endeavored to fix those runtime failures, as I believe them to have been pre-existing, but I have removed any of the benchmarks which no longer compile due to the changes in this branch. --- crates/snapshot/src/lib.rs | 10 ++- crates/table/benches/page_manager.rs | 109 +-------------------------- crates/table/src/pages.rs | 10 +-- 3 files changed, 12 insertions(+), 117 deletions(-) diff --git a/crates/snapshot/src/lib.rs b/crates/snapshot/src/lib.rs index f8b2eb57f64..b5922653e2b 100644 --- a/crates/snapshot/src/lib.rs +++ b/crates/snapshot/src/lib.rs @@ -99,6 +99,8 @@ pub struct SnapshotReadMetrics { pub blob: SnapshotReadKindMetrics, } +pub type TablePages = Vec>>; + impl SnapshotReadMetrics { pub fn iter(&self) -> impl Iterator + '_ { [("metadata", &self.metadata), ("page", &self.page), ("blob", &self.blob)].into_iter() @@ -600,7 +602,7 @@ impl Snapshot { pages: &[blake3::Hash], page_pool: &PagePool, metrics: &mut SnapshotReadKindMetrics, - ) -> Result>>, SnapshotError> { + ) -> Result { pages .iter() .map(|hash| { @@ -647,7 +649,7 @@ impl Snapshot { TableEntry { table_id, pages }: &TableEntry, page_pool: &PagePool, metrics: &mut SnapshotReadKindMetrics, - ) -> Result<(TableId, Vec>>), SnapshotError> { + ) -> Result<(TableId, TablePages), SnapshotError> { Ok(( *table_id, Self::reconstruct_one_table_pages(object_repo, pages, page_pool, metrics)?, @@ -669,7 +671,7 @@ impl Snapshot { object_repo: &DirTrie, page_pool: &PagePool, metrics: &mut SnapshotReadKindMetrics, - ) -> Result>>>, SnapshotError> { + ) -> Result, SnapshotError> { self.tables .iter() .map(|tbl| Self::reconstruct_one_table(object_repo, tbl, page_pool, metrics)) @@ -1562,7 +1564,7 @@ pub struct ReconstructedSnapshot { /// This includes the system tables, /// so the schema of user-defined tables can be recovered /// given knowledge of the schema of `st_table` and `st_column`. - pub tables: BTreeMap>>>, + pub tables: BTreeMap, /// If the snapshot was compressed or not. pub compress_type: CompressType, /// Metrics collected while reading this snapshot from disk. diff --git a/crates/table/benches/page_manager.rs b/crates/table/benches/page_manager.rs index 10fa708e8d2..36892640604 100644 --- a/crates/table/benches/page_manager.rs +++ b/crates/table/benches/page_manager.rs @@ -10,7 +10,7 @@ use rand::{Rng, SeedableRng}; use spacetimedb_lib::db::raw_def::v9::RawIndexAlgorithm; use spacetimedb_lib::db::raw_def::v9::RawModuleDefV9Builder; use spacetimedb_primitives::{ColList, IndexId, TableId}; -use spacetimedb_sats::layout::{row_size_for_bytes, row_size_for_type, Size}; +use spacetimedb_sats::layout::row_size_for_type; use spacetimedb_sats::raw_identifier::RawIdentifier; use spacetimedb_sats::{AlgebraicType, AlgebraicValue, ProductType, ProductValue}; use spacetimedb_schema::def::BTreeAlgorithm; @@ -18,7 +18,7 @@ use spacetimedb_schema::def::ModuleDef; use spacetimedb_schema::schema::TableSchema; use spacetimedb_schema::table_name::TableName; use spacetimedb_table::blob_store::NullBlobStore; -use spacetimedb_table::indexes::{Byte, Bytes, PageOffset, RowPointer, SquashedOffset, PAGE_DATA_SIZE}; +use spacetimedb_table::indexes::{Byte, Bytes, PageOffset, RowPointer, SquashedOffset}; use spacetimedb_table::page_pool::PagePool; use spacetimedb_table::pages::Pages; use spacetimedb_table::row_type_visitor::{row_type_visitor, VarLenVisitorProgram}; @@ -176,30 +176,6 @@ fn var_len_rows_per_page(data_size_in_bytes: usize) -> usize { PageOffset::PAGE_END.idx() / var_object_size } -fn reserve_empty_page(c: &mut Criterion) { - const RESERVE_SIZE: Size = row_size_for_bytes(8); - - let mut group = c.benchmark_group("reserve_empty_page"); - group.throughput(Throughput::Bytes(PAGE_DATA_SIZE as _)); - group.bench_function("leave_uninit", |b| { - let pool = PagePool::new_for_test(); - let mut pages = Pages::default(); - b.iter(|| { - let _ = black_box(pages.reserve_empty_page(&pool, RESERVE_SIZE)); - }); - }); - - let fill_with_zeros = |_, _, pages: &mut Pages| { - let pool = PagePool::new_for_test(); - let page = pages.reserve_empty_page(&pool, RESERVE_SIZE).unwrap(); - let page = pages.get_page_mut(page); - unsafe { page.zero_data() }; - }; - group.bench_function("fill_with_zeros", |b| { - iter_time_with(b, &mut Pages::default(), |_, _| (), fill_with_zeros) - }); -} - fn insert_one_page_worth_fixed_len( pool: &PagePool, pages: &mut Pages, @@ -349,85 +325,6 @@ fn retrieve_one_page_fixed_len(c: &mut Criterion) { bench_retrieve_one_page::(&mut group, &U32x64::var_len_visitor(), "U32x64/VarLenVisitorProgram"); } -// insert a bunch of rows, -// delete some fraction of them to create holes in multiple pages, -// then time to insert into those holes -fn insert_with_holes_fixed_len(c: &mut Criterion) { - fn bench_insert_with_holes(c: &mut Criterion, var_len_visitor: &impl VarLenMembers, name: &str) { - let mut group = c.benchmark_group(format!("insert_with_holes_fixed_len/{name}")); - let val = R::from_u64(0xdeadbeef_0badbeef); - for delete_ratio in [0.1f64, 0.25, 0.5, 0.75, 0.9, 1.0] { - let num_pages = 16; - let total_num_rows = rows_per_page::() * num_pages; - let num_to_delete = (total_num_rows as f64 * delete_ratio) as usize; - - let num_to_delete_in_bytes = num_to_delete * mem::size_of::(); - - let row_size = row_size_for_type::(); - - group.throughput(Throughput::Bytes(num_to_delete_in_bytes as u64)); - - group.bench_function(delete_ratio.to_string(), |b| { - let pool = PagePool::new_for_test(); - let mut pages = Pages::default(); - - let mut rng = StdRng::seed_from_u64(0xa5a5a5a5_a5a5a5a5); - - for _ in 0..num_pages { - let page = pages.reserve_empty_page(&pool, row_size).unwrap(); - let page = pages.get_page_mut(page); - - unsafe { page.zero_data() }; - } - - let pre = |_, (pages, pool): &mut (Pages, PagePool)| { - pages.clear(); - let mut ptrs_to_delete = Vec::with_capacity(num_to_delete); - for _ in 0..total_num_rows { - let (page_idx, offset) = unsafe { - pages.insert_row(pool, var_len_visitor, row_size, val.as_bytes(), &[], &mut NullBlobStore) - } - .unwrap(); - - if rng.random_bool(delete_ratio) { - ptrs_to_delete.push(RowPointer::new( - false, - page_idx, - offset, - SquashedOffset::COMMITTED_STATE, - )); - } - } - let actual_num_deleted = ptrs_to_delete.len(); - for ptr in ptrs_to_delete { - unsafe { - pages.delete_row(var_len_visitor, row_size, ptr, &mut NullBlobStore); - } - } - actual_num_deleted - }; - let body = |actual_num_deleted, _, (pages, pool): &mut (Pages, PagePool)| { - for _ in 0..actual_num_deleted { - let _ = black_box(unsafe { - pages.insert_row(pool, var_len_visitor, row_size, val.as_bytes(), &[], &mut NullBlobStore) - }); - } - }; - iter_time_with(b, &mut (pages, pool), pre, body); - }); - } - } - - bench_insert_with_holes::(c, &NullVarLenVisitor, "u64/NullVarLenVisitor"); - bench_insert_with_holes::(c, &u64::var_len_visitor(), "u64/VarLenVisitorProgram"); - - bench_insert_with_holes::(c, &NullVarLenVisitor, "U32x8/NullVarLenVisitor"); - bench_insert_with_holes::(c, &U32x8::var_len_visitor(), "U32x8/VarLenVisitorProgram"); - - bench_insert_with_holes::(c, &NullVarLenVisitor, "U32x64/NullVarLenVisitor"); - bench_insert_with_holes::(c, &U32x64::var_len_visitor(), "U32x64/VarLenVisitorProgram"); -} - // insert a whole bunch of rows, then time to copy_filter materialize a view fn copy_filter_fixed_len(c: &mut Criterion) { @@ -482,11 +379,9 @@ fn copy_filter_fixed_len(c: &mut Criterion) { criterion_group!( pages, - reserve_empty_page, insert_one_page_fixed_len, delete_one_page_fixed_len, retrieve_one_page_fixed_len, - insert_with_holes_fixed_len, copy_filter_fixed_len, ); diff --git a/crates/table/src/pages.rs b/crates/table/src/pages.rs index 5949419deea..04969fe1321 100644 --- a/crates/table/src/pages.rs +++ b/crates/table/src/pages.rs @@ -182,10 +182,8 @@ impl Pages { #[doc(hidden)] // Used in benchmarks. pub fn clear(&mut self) { // Clear every page. - for page in &mut self.pages { - if let Some(page) = page { - page.clear(); - } + for page in self.pages.iter_mut().flatten() { + page.clear(); } // Mark every page non-full. self.non_full_pages = (0..self.pages.len()) @@ -255,7 +253,7 @@ impl Pages { f: impl FnOnce(&mut Page) -> Res, ) -> Result<(PageIndex, Res), Error> { let page_index = self.find_page_with_space_for_row(pool, fixed_row_size, num_var_len_granules)?; - let res = f(&mut self + let res = f(self .get_mut(page_index) .expect("page returned by `find_page_with_space_for_row` to be present in `self.pages`")); self.record_page_non_full(page_index, fixed_row_size); @@ -579,7 +577,7 @@ impl Pages { /// /// Indexes in the iterator will not necessarily correspond to the pages' `PageIndex`es. pub fn into_page_iter(self) -> impl Iterator> { - self.pages.into_iter().filter_map(|page| page) + self.pages.into_iter().flatten() } /// Iterate over only those pages in `self` that are present. From 159dd5dbe0846a695587b6d936fed5426e0d0798 Mon Sep 17 00:00:00 2001 From: Phoebe Goldman Date: Fri, 21 Aug 2026 11:04:26 -0400 Subject: [PATCH 04/10] fix compliation errors in datastore crate --- .../locking_tx_datastore/committed_state.rs | 7 +++-- .../src/locking_tx_datastore/mut_tx.rs | 28 +++++++++++-------- .../src/locking_tx_datastore/replay.rs | 12 ++++---- crates/table/src/table.rs | 22 ++++++++------- 4 files changed, 40 insertions(+), 29 deletions(-) diff --git a/crates/datastore/src/locking_tx_datastore/committed_state.rs b/crates/datastore/src/locking_tx_datastore/committed_state.rs index 0ce629524c4..88cadf976e0 100644 --- a/crates/datastore/src/locking_tx_datastore/committed_state.rs +++ b/crates/datastore/src/locking_tx_datastore/committed_state.rs @@ -617,6 +617,7 @@ impl CommittedState { tx_data: &mut TxData, table_id: TableId, table: &mut Table, + page_pool: &PagePool, blob_store: &mut dyn BlobStore, row_ptrs_len: usize, row_ptrs: impl Iterator, @@ -633,7 +634,7 @@ impl CommittedState { // TODO: re-write `TxData` to remove `ProductValue`s let pv = table - .delete(blob_store, row_ptr, |row| row.to_product_value()) + .delete(page_pool, blob_store, row_ptr, |row| row.to_product_value()) .expect("Delete for non-existent row!"); deletes.push(pv); } @@ -650,10 +651,11 @@ impl CommittedState { for (table_id, row_ptrs) in delete_tables { match self.get_table_and_blob_store_mut(table_id) { - Ok((table, blob_store, ..)) => delete_rows( + Ok((table, blob_store, _index_map, page_pool)) => delete_rows( tx_data, table_id, table, + page_pool, blob_store, row_ptrs.len(), row_ptrs.iter(), @@ -675,6 +677,7 @@ impl CommittedState { tx_data, table_id, &mut table, + &self.page_pool, &mut self.blob_store, row_ptrs.len(), row_ptrs.into_iter(), diff --git a/crates/datastore/src/locking_tx_datastore/mut_tx.rs b/crates/datastore/src/locking_tx_datastore/mut_tx.rs index fec1d4804dd..5a6b9f8a77e 100644 --- a/crates/datastore/src/locking_tx_datastore/mut_tx.rs +++ b/crates/datastore/src/locking_tx_datastore/mut_tx.rs @@ -3451,11 +3451,12 @@ pub(super) fn insert<'a, const GENERATE: bool>( is_scheduler_table: tx_table.is_scheduler(), }; let ok = |row_ref| Ok((gen_cols, row_ref, insert_flags)); + let page_pool = &committed_state.page_pool; // `CHECK_SAME_ROW = true`, as there might be an identical row already in the tx state. // SAFETY: `tx_table.is_row_present(row)` holds as we still haven't deleted the row, // in particular, the `write_gen_val_to_col` call does not remove the row. - let res = unsafe { tx_table.confirm_insertion::(tx_blob_store, tx_row_ptr, blob_bytes) }; + let res = unsafe { tx_table.confirm_insertion::(page_pool, tx_blob_store, tx_row_ptr, blob_bytes) }; match res { Ok((tx_row_hash, tx_row_ptr)) => { @@ -3499,7 +3500,7 @@ pub(super) fn insert<'a, const GENERATE: bool>( // - Insert Row A // This is impossible to recover if `Running 2` elides its insert. tx_table - .delete(tx_blob_store, tx_row_ptr, |_| ()) + .delete(page_pool, tx_blob_store, tx_row_ptr, |_| ()) .expect("Failed to delete a row we just inserted"); // It's possible that `row` appears in the committed state, @@ -3528,7 +3529,7 @@ pub(super) fn insert<'a, const GENERATE: bool>( let res = unsafe { commit_table.check_unique_constraints(tx_row_ref, |ixs| ixs, is_deleted) }; if let Err(e) = res { // There was a constraint violation, so undo the insertion. - tx_table.delete(tx_blob_store, tx_row_ptr, |_| {}); + tx_table.delete(&page_pool, tx_blob_store, tx_row_ptr, |_| {}); return Err(IndexError::from(e).into()); } @@ -3601,6 +3602,8 @@ impl MutTxId { // SAFETY: `tx_table.is_row_present(tx_row_ptr)` holds as we just inserted it. let tx_row_ref = unsafe { tx_table.get_row_ref_unchecked(tx_blob_store, tx_row_ptr) }; + let page_pool = &self.committed_state_write_lock.page_pool; + let err = 'error: { // This macros can be thought of as a `throw $e` within `'error`. // TODO(centril): Get rid of this once we have stable `try` blocks or polonius. @@ -3665,7 +3668,7 @@ impl MutTxId { // 3. we just inserted `tx_row_ptr` into `tx_table`, so we know it is valid. if unsafe { Table::eq_row_in_page(commit_table, old_ptr, tx_table, tx_row_ptr) } { // SAFETY: `tx_table.is_row_present(tx_row_ptr)` holds, as noted in 3. - unsafe { tx_table.delete_internal_skip_pointer_map(tx_blob_store, tx_row_ptr) }; + unsafe { tx_table.delete_internal_skip_pointer_map(page_pool, tx_blob_store, tx_row_ptr) }; // SAFETY: `commit_table.is_row_present(old_ptr)` holds, as noted in 2. let row_ref = unsafe { commit_table.get_row_ref_unchecked(commit_blob_store, old_ptr) }; return ok(RowRefInsertion::Existed(row_ref)); @@ -3685,7 +3688,7 @@ impl MutTxId { // in particular, the `write_gen_val_to_col` call does not remove the row. // On error, `tx_row_ptr` has already been removed, so don't do it again. let (_, tx_row_ptr) = - unsafe { tx_table.confirm_insertion::(tx_blob_store, tx_row_ptr, blob_bytes) }?; + unsafe { tx_table.confirm_insertion::(page_pool, tx_blob_store, tx_row_ptr, blob_bytes) }?; // Delete the old row. del_table.insert(old_ptr); @@ -3703,7 +3706,8 @@ impl MutTxId { // SAFETY: `tx_table.is_row_present(tx_row_ptr)` and `tx_table.is_row_present(old_ptr)` both hold // as we've deleted neither. // In particular, the `write_gen_val_to_col` call does not remove the row. - let tx_row_ptr = unsafe { tx_table.confirm_update(tx_blob_store, tx_row_ptr, old_ptr, blob_bytes) }?; + let tx_row_ptr = + unsafe { tx_table.confirm_update(page_pool, tx_blob_store, tx_row_ptr, old_ptr, blob_bytes) }?; if let Some(old_commit_del_ptr) = old_commit_del_ptr { // If we have an identical deleted row in the committed state, @@ -3718,7 +3722,7 @@ impl MutTxId { // It is important that we `confirm_update` first, // as we must ensure that undeleting the row causes no tx state conflict. tx_table - .delete(tx_blob_store, tx_row_ptr, |_| ()) + .delete(page_pool, tx_blob_store, tx_row_ptr, |_| ()) .expect("Failed to delete a row we just inserted"); // Undelete. @@ -3748,7 +3752,7 @@ impl MutTxId { // When we reach here, we had an error and we need to revert the insertion of `tx_row_ref`. // SAFETY: `tx_table.is_row_present(tx_row_ptr)` holds, // as we still haven't deleted the row physically. - unsafe { tx_table.delete_internal_skip_pointer_map(tx_blob_store, tx_row_ptr) }; + unsafe { tx_table.delete_internal_skip_pointer_map(page_pool, tx_blob_store, tx_row_ptr) }; Err(err) } @@ -3770,7 +3774,7 @@ impl MutTxId { let (tx_table, tx_blob_store, delete_table) = self .tx_state .get_table_and_blob_store_or_create_from(table_id, commit_table); - let mut rows_removed = tx_table.clear(tx_blob_store); + let mut rows_removed = tx_table.clear(&self.committed_state_write_lock.page_pool, tx_blob_store); // Mark every row in the committed state as deleted. for row in commit_table.scan_rows(commit_bs) { @@ -3796,7 +3800,9 @@ pub(super) fn delete( let (table, blob_store) = tx_state .get_table_and_blob_store(table_id) .ok_or(TableError::IdNotFoundState(table_id))?; - Ok(table.delete(blob_store, row_pointer, |_| ()).is_some()) + Ok(table + .delete(&committed_state.page_pool, blob_store, row_pointer, |_| ()) + .is_some()) } SquashedOffset::COMMITTED_STATE => { let commit_table = committed_state @@ -3855,7 +3861,7 @@ impl MutTxId { // Do this before actually deleting to drop the borrows on the table. // SAFETY: `temp_ptr` is valid because we just inserted it and haven't deleted it since. unsafe { - tx_table.delete_internal_skip_pointer_map(tx_blob_store, temp_ptr); + tx_table.delete_internal_skip_pointer_map(page_pool, tx_blob_store, temp_ptr); } // Delete the found row either by marking (commit table) diff --git a/crates/datastore/src/locking_tx_datastore/replay.rs b/crates/datastore/src/locking_tx_datastore/replay.rs index c42185e2bb2..b4b6abddd3f 100644 --- a/crates/datastore/src/locking_tx_datastore/replay.rs +++ b/crates/datastore/src/locking_tx_datastore/replay.rs @@ -570,7 +570,7 @@ impl<'cs> ReplayCommittedState<'cs> { }) .collect::>(); - let (st_sequence, blob_store, ..) = self + let (st_sequence, blob_store, _index_map, page_pool) = self .get_table_and_blob_store_mut(ST_SEQUENCE_ID) .expect("`st_sequence` should exist"); @@ -602,7 +602,7 @@ impl<'cs> ReplayCommittedState<'cs> { prev_row_pointer }; - st_sequence.delete(blob_store, row_pointer_to_delete, |_| ()) + st_sequence.delete(page_pool, blob_store, row_pointer_to_delete, |_| ()) .expect("Duplicated `st_sequence` row at `row_pointer_to_delete` should be present in `st_sequence` during fixup"); } } @@ -646,13 +646,13 @@ impl<'cs> ReplayCommittedState<'cs> { .map(|row_ref| row_ref.pointer()) .collect(); - let (st_event_table, blob_store, ..) = self + let (st_event_table, blob_store, _index_map, page_pool) = self .get_table_and_blob_store_mut(ST_EVENT_TABLE_ID) .expect("`st_event_table` was found above"); for ptr in orphaned_rows { st_event_table - .delete(blob_store, ptr, |_| ()) + .delete(page_pool, blob_store, ptr, |_| ()) .expect("Orphaned `st_event_table` row at `ptr` should be present in `st_event_table` during fixup"); } } @@ -1039,12 +1039,12 @@ impl<'cs> ReplayCommittedState<'cs> { } // Get the table for mutation. - let (table, blob_store, ..) = self.get_table_and_blob_store_mut(table_id)?; + let (table, blob_store, _index_map, page_pool) = self.get_table_and_blob_store_mut(table_id)?; // We do not need to consider a truncation of `st_table` itself, // as if that happens, the database is bricked. - table.clear(blob_store); + table.clear(page_pool, blob_store); Ok(()) } diff --git a/crates/table/src/table.rs b/crates/table/src/table.rs index 3d2bce7ef0b..689789654ec 100644 --- a/crates/table/src/table.rs +++ b/crates/table/src/table.rs @@ -1389,16 +1389,18 @@ impl Table { Ok(existing_row_ptr) } - // /// Clears this table, removing all present rows from it. - // pub fn clear(&mut self, pool: &PagePool, blob_store: &mut dyn BlobStore) -> u64 { - // let ptrs = self.scan_all_row_ptrs(); - // let len = ptrs.len() as u64; - // for ptr in ptrs { - // // SAFETY: `ptr` came rom `self.scan_rows(...)`, so it's present. - // unsafe { self.delete_unchecked(pool, blob_store, ptr) }; - // } - // len - // } + /// Clears this table, removing all present rows from it. + /// + /// This implements the `truncate` operation. + pub fn clear(&mut self, pool: &PagePool, blob_store: &mut dyn BlobStore) -> u64 { + let ptrs = self.scan_all_row_ptrs(); + let len = ptrs.len() as u64; + for ptr in ptrs { + // SAFETY: `ptr` came rom `self.scan_rows(...)`, so it's present. + unsafe { self.delete_unchecked(pool, blob_store, ptr) }; + } + len + } /// Returns the row type for rows in this table. pub fn get_row_type(&self) -> &ProductType { From f0cc92ace23033cc5a52bb2507c364327c750759 Mon Sep 17 00:00:00 2001 From: Phoebe Goldman Date: Fri, 21 Aug 2026 11:54:06 -0400 Subject: [PATCH 05/10] Don't attempt to fsync or count non-pages in snapshot --- crates/snapshot/src/lib.rs | 20 ++++++++++++++++++-- 1 file changed, 18 insertions(+), 2 deletions(-) diff --git a/crates/snapshot/src/lib.rs b/crates/snapshot/src/lib.rs index b5922653e2b..d4673276620 100644 --- a/crates/snapshot/src/lib.rs +++ b/crates/snapshot/src/lib.rs @@ -680,7 +680,18 @@ impl Snapshot { /// The number of objects in this snapshot, both blobs and pages. pub fn total_objects(&self) -> usize { - self.blobs.len() + self.tables.iter().map(|table| table.pages.len()).sum::() + self.blobs.len() + + self + .tables + .iter() + .map(|table| { + table + .pages + .iter() + .filter(|hash| **hash != ZERO_HASH_DENOTING_ABSENT_PAGE) + .count() + }) + .sum::() } /// Obtain an iterator over the [`blake3::Hash`]es of all objects @@ -689,7 +700,12 @@ impl Snapshot { self.blobs .iter() .map(|b| blake3::Hash::from_bytes(b.hash.data)) - .chain(self.tables.iter().flat_map(|t| t.pages.iter().copied())) + .chain(self.tables.iter().flat_map(|t| { + t.pages + .iter() + .filter(|hash| **hash != ZERO_HASH_DENOTING_ABSENT_PAGE) + .copied() + })) } /// Obtain an iterator over the [`PathBuf`]s of all objects From 78dfcfcc0474745d15023dc22c5304bbc4940487 Mon Sep 17 00:00:00 2001 From: Phoebe Goldman Date: Fri, 21 Aug 2026 12:34:30 -0400 Subject: [PATCH 06/10] Clippy --- .../datastore/src/locking_tx_datastore/committed_state.rs | 1 + crates/datastore/src/locking_tx_datastore/mut_tx.rs | 2 +- crates/table/src/pages.rs | 3 ++- crates/table/src/table.rs | 6 +++--- 4 files changed, 7 insertions(+), 5 deletions(-) diff --git a/crates/datastore/src/locking_tx_datastore/committed_state.rs b/crates/datastore/src/locking_tx_datastore/committed_state.rs index 88cadf976e0..71f0d3f0931 100644 --- a/crates/datastore/src/locking_tx_datastore/committed_state.rs +++ b/crates/datastore/src/locking_tx_datastore/committed_state.rs @@ -613,6 +613,7 @@ impl CommittedState { pending_schema_changes: ThinVec, truncates: &mut IntSet, ) { + #[allow(clippy::too_many_arguments)] fn delete_rows( tx_data: &mut TxData, table_id: TableId, diff --git a/crates/datastore/src/locking_tx_datastore/mut_tx.rs b/crates/datastore/src/locking_tx_datastore/mut_tx.rs index 5a6b9f8a77e..ba0cda8e616 100644 --- a/crates/datastore/src/locking_tx_datastore/mut_tx.rs +++ b/crates/datastore/src/locking_tx_datastore/mut_tx.rs @@ -3529,7 +3529,7 @@ pub(super) fn insert<'a, const GENERATE: bool>( let res = unsafe { commit_table.check_unique_constraints(tx_row_ref, |ixs| ixs, is_deleted) }; if let Err(e) = res { // There was a constraint violation, so undo the insertion. - tx_table.delete(&page_pool, tx_blob_store, tx_row_ptr, |_| {}); + tx_table.delete(page_pool, tx_blob_store, tx_row_ptr, |_| {}); return Err(IndexError::from(e).into()); } diff --git a/crates/table/src/pages.rs b/crates/table/src/pages.rs index daf3989a6bf..d39bfeb0f9f 100644 --- a/crates/table/src/pages.rs +++ b/crates/table/src/pages.rs @@ -181,7 +181,8 @@ impl Pages { !self.non_full_pages.remove(&(free_granules, page_index)) }); - let page = std::mem::replace(&mut self.pages[page_index.idx()], None) + let page = self.pages[page_index.idx()] + .take() .expect("freed page to have been present after we already checked its presence"); pool.put(page); diff --git a/crates/table/src/table.rs b/crates/table/src/table.rs index 689789654ec..0a439a572bd 100644 --- a/crates/table/src/table.rs +++ b/crates/table/src/table.rs @@ -3127,15 +3127,15 @@ pub(crate) mod test { assert_eq!(table2.blob_store_bytes, BLOB_OBJ_LEN); // Delete `short_str` row. This should not affect the byte count. - table1.delete(&pool, blob_store, short_row_ptr, |_| ()).unwrap(); + table1.delete(pool, blob_store, short_row_ptr, |_| ()).unwrap(); assert_eq!(table1.blob_store_bytes, BLOB_OBJ_LEN_2X); // Delete the first long string row. This gets us down to `BLOB_OBJ_LEN` (we had 2x before). - table1.delete(&pool, blob_store, long_row_ptr, |_| ()).unwrap(); + table1.delete(pool, blob_store, long_row_ptr, |_| ()).unwrap(); assert_eq!(table1.blob_store_bytes, BLOB_OBJ_LEN); // Delete the first long string row. This gets us down to 0 (we've now deleted 2x). - table1.delete(&pool, blob_store, long_row_ptr2, |_| ()).unwrap(); + table1.delete(pool, blob_store, long_row_ptr2, |_| ()).unwrap(); assert_eq!(table1.blob_store_bytes, 0.into()); } From bda0156a9291e438fc2e8a923a37c484656fba1e Mon Sep 17 00:00:00 2001 From: Phoebe Goldman Date: Fri, 21 Aug 2026 15:19:34 -0400 Subject: [PATCH 07/10] Remove methods that will be removed from next patch Per Joshua's request, remove these methods (and the stuff that depends on them) now, rather than leaving them in place here to be cleaned up by the next PR. --- crates/sats/src/proptest.rs | 6 ++ crates/table/benches/page_manager.rs | 139 +-------------------------- crates/table/src/pages.rs | 134 +------------------------- crates/table/src/table.rs | 10 +- 4 files changed, 15 insertions(+), 274 deletions(-) diff --git a/crates/sats/src/proptest.rs b/crates/sats/src/proptest.rs index 503d804f6a6..b291ae906ba 100644 --- a/crates/sats/src/proptest.rs +++ b/crates/sats/src/proptest.rs @@ -224,6 +224,12 @@ pub fn generate_typed_row() -> impl Strategy impl Strategy { + gen_with(generate_row_type(0..=SIZE), |ty| { + (generate_product_value(ty.clone()), generate_product_value(ty)) + }) +} + pub fn generate_typed_row_vec( size: impl Into, num_rows_min: usize, diff --git a/crates/table/benches/page_manager.rs b/crates/table/benches/page_manager.rs index 36892640604..b35b4e54c61 100644 --- a/crates/table/benches/page_manager.rs +++ b/crates/table/benches/page_manager.rs @@ -5,8 +5,6 @@ use criterion::measurement::{Measurement, WallTime}; use criterion::{ black_box, criterion_group, criterion_main, Bencher, BenchmarkGroup, BenchmarkId, Criterion, Throughput, }; -use rand::rngs::StdRng; -use rand::{Rng, SeedableRng}; use spacetimedb_lib::db::raw_def::v9::RawIndexAlgorithm; use spacetimedb_lib::db::raw_def::v9::RawModuleDefV9Builder; use spacetimedb_primitives::{ColList, IndexId, TableId}; @@ -176,51 +174,8 @@ fn var_len_rows_per_page(data_size_in_bytes: usize) -> usize { PageOffset::PAGE_END.idx() / var_object_size } -fn insert_one_page_worth_fixed_len( - pool: &PagePool, - pages: &mut Pages, - visitor: &impl VarLenMembers, - val: &R, -) { - let size = row_size_for_type::(); - for _ in 0..rows_per_page::() { - let _ = black_box(unsafe { - black_box(&mut *pages).insert_row(pool, visitor, size, val.as_bytes(), &[], &mut NullBlobStore) - }); - } -} - type Group<'a, 'b> = &'a mut BenchmarkGroup<'b, WallTime>; -// time to insert a whole bunch of rows -fn insert_one_page_fixed_len(c: &mut Criterion) { - fn bench_insert_one_page_fixed_len(group: Group<'_, '_>, visitor: &impl VarLenMembers, name: &str) { - group.throughput(Throughput::Bytes( - rows_per_page::() as u64 * mem::size_of::() as u64, - )); - group.bench_function(name, |b| { - let pool = PagePool::new_for_test(); - let mut pages = Pages::default(); - // `0xa5` is the alternating bit pattern, which makes incorrect accesses obvious. - insert_one_page_worth_fixed_len(&pool, &mut pages, visitor, &R::from_u64(0xa5a5a5a5_a5a5a5a5)); - let pre = |_, pages: &mut Pages| pages.clear(); - iter_time_with(b, &mut pages, pre, |_, _, pages| { - insert_one_page_worth_fixed_len(&pool, pages, visitor, &R::from_u64(0xdeadbeef_0badbeef)) - }); - }); - } - - let mut group = c.benchmark_group("insert_one_page_fixed_len"); - bench_insert_one_page_fixed_len::(&mut group, &NullVarLenVisitor, "u64/NullVarLenVisitor"); - bench_insert_one_page_fixed_len::(&mut group, &u64::var_len_visitor(), "u64/VarLenVisitorProgram"); - - bench_insert_one_page_fixed_len::(&mut group, &NullVarLenVisitor, "U32x8/NullVarLenVisitor"); - bench_insert_one_page_fixed_len::(&mut group, &U32x8::var_len_visitor(), "U32x8/VarLenVisitorProgram"); - - bench_insert_one_page_fixed_len::(&mut group, &NullVarLenVisitor, "U32x64/NullVarLenVisitor"); - bench_insert_one_page_fixed_len::(&mut group, &U32x64::var_len_visitor(), "U32x64/VarLenVisitorProgram"); -} - fn fill_page_with_fixed_len_collect_row_pointers( pool: &PagePool, pages: &mut Pages, @@ -246,45 +201,6 @@ fn fill_page_with_fixed_len_collect_row_pointers( ptrs } -// insert a whole bunch of rows, then time to delete them all -fn delete_one_page_fixed_len(c: &mut Criterion) { - fn bench_delete_one_page_fixed_len(group: Group<'_, '_>, visitor: &impl VarLenMembers, name: &str) { - let rows_per_page = rows_per_page::(); - - group.throughput(Throughput::Bytes(rows_per_page as u64 * mem::size_of::() as u64)); - - group.bench_function(name, |b| { - let pre = |i, (pages, pool): &mut _| { - let val = R::from_u64(i); - fill_page_with_fixed_len_collect_row_pointers::(pool, pages, visitor, &val) - }; - iter_time_with( - b, - &mut (Pages::default(), PagePool::new_for_test()), - pre, - |ptrs, _, (pages, _)| { - for ptr in ptrs { - unsafe { - pages.delete_row(visitor, row_size_for_type::(), black_box(ptr), &mut NullBlobStore) - }; - } - }, - ); - }); - } - - let mut group = c.benchmark_group("delete_one_page_fixed_len"); - - bench_delete_one_page_fixed_len::(&mut group, &NullVarLenVisitor, "u64/NullVarLenVisitor"); - bench_delete_one_page_fixed_len::(&mut group, &u64::var_len_visitor(), "u64/VarLenVisitorProgram"); - - bench_delete_one_page_fixed_len::(&mut group, &NullVarLenVisitor, "U32x8/NullVarLenVisitor"); - bench_delete_one_page_fixed_len::(&mut group, &U32x8::var_len_visitor(), "U32x8/VarLenVisitorProgram"); - - bench_delete_one_page_fixed_len::(&mut group, &NullVarLenVisitor, "U32x64/NullVarLenVisitor"); - bench_delete_one_page_fixed_len::(&mut group, &U32x64::var_len_visitor(), "U32x64/VarLenVisitorProgram"); -} - // insert a whole bunch of rows, then time to access them fn retrieve_one_page_fixed_len(c: &mut Criterion) { fn bench_retrieve_one_page(group: Group<'_, '_>, visitor: &impl VarLenMembers, name: &str) { @@ -325,65 +241,12 @@ fn retrieve_one_page_fixed_len(c: &mut Criterion) { bench_retrieve_one_page::(&mut group, &U32x64::var_len_visitor(), "U32x64/VarLenVisitorProgram"); } -// insert a whole bunch of rows, then time to copy_filter materialize a view - -fn copy_filter_fixed_len(c: &mut Criterion) { - fn bench_copy_filter(c: &mut Criterion, name: &str) { - let mut group = c.benchmark_group(format!("copy_filter_fixed_len/{name}")); - let row_size = black_box(row_size_for_type::()); - - let val = R::from_u64(0xdeadbeef_0badbeef); - for keep_ratio in [0.1, 0.25, 0.5, 0.75, 0.9, 1.0] { - let visitor = &NullVarLenVisitor; - let pool = PagePool::new_for_test(); - let mut pages = Pages::default(); - - let num_pages = 16; - let total_num_rows = rows_per_page::() * num_pages; - - for _ in 0..total_num_rows { - unsafe { pages.insert_row(&pool, visitor, row_size, val.as_bytes(), &[], &mut NullBlobStore) }.unwrap(); - } - - let num_to_keep = (total_num_rows as f64 * keep_ratio) as usize; - let num_to_keep_bytes = num_to_keep * mem::size_of::(); - - group.throughput(Throughput::Bytes(num_to_keep_bytes as u64)); - - let mut rng = StdRng::seed_from_u64(0xa5a5a5a5_a5a5a5a5); - - // To avoid advancing RNG in the benchmark, - // precompute a big vec of bools, with one bool for each value that we may or may not keep. - let keep_seq: Vec = (0..total_num_rows).map(|_| rng.random_bool(keep_ratio)).collect(); - - group.bench_function(keep_ratio.to_string(), |b| { - b.iter_with_large_drop(|| unsafe { - let mut keep_iter = keep_seq.iter().copied(); - black_box(&pages).copy_filter(visitor, row_size, None::<&mut Box>, |_, _| { - black_box(keep_iter.next().unwrap_or_default()) - }) - }); - }); - } - } - - bench_copy_filter::(c, "u64"); - bench_copy_filter::(c, "U32x8"); - bench_copy_filter::(c, "U32x64"); -} - // TODO(bench): // - Duplicate above benchmarks with var-len rows of various sizes // - In the insert-with-holes benchmark, randomize size of each row to simulate fragmentation. // - Extend above benchmarks to go through `Table` with `AlgebraicValue`. -criterion_group!( - pages, - insert_one_page_fixed_len, - delete_one_page_fixed_len, - retrieve_one_page_fixed_len, - copy_filter_fixed_len, -); +criterion_group!(pages, retrieve_one_page_fixed_len); fn schema_from_ty(ty: ProductType, name: &str) -> TableSchema { let mut result = TableSchema::from_product_type(ty); diff --git a/crates/table/src/pages.rs b/crates/table/src/pages.rs index 04969fe1321..451d00a6fb5 100644 --- a/crates/table/src/pages.rs +++ b/crates/table/src/pages.rs @@ -1,12 +1,12 @@ //! Provides [`Pages`], a page manager dealing with [`Page`]s as a collection. -use super::blob_store::{BlobHash, BlobStore}; +use super::blob_store::BlobStore; use super::indexes::{Bytes, PageIndex, PageOffset, RowPointer}; use super::page::Page; use super::page_pool::PagePool; use super::table::BlobNumBytes; use super::var_len::VarLenMembers; -use core::ops::{ControlFlow, Deref}; +use core::ops::Deref; use spacetimedb_sats::layout::Size; use spacetimedb_sats::memory_usage::MemoryUsage; use std::collections::BTreeSet; @@ -173,32 +173,6 @@ impl Pages { assert!(newly_inserted_into_free_pages_set) } - /// Make all pages within `self` clear, - /// deleting all rows. - // - // TODO(delete-free-page): Determine what to do with this method. - // It doesn't really make sense given that it clears all pages but doesn't delete them, - // but it's only used in benchmarks, and those benchmarks are using it specifically to bypass allocating new pages. - #[doc(hidden)] // Used in benchmarks. - pub fn clear(&mut self) { - // Clear every page. - for page in self.pages.iter_mut().flatten() { - page.clear(); - } - // Mark every page non-full. - self.non_full_pages = (0..self.pages.len()) - // We could probably compute the number of available granules once and use it for all pages, - // rather than calling the method on each page, - // but we'd have to do some amount of reasoning to demonstrate it was correct - // based on the definition of `Page::clear`, - // and why bother? - .filter_map(|idx| { - let idx = PageIndex(idx as u64); - self.get(idx).map(|page| (page.available_var_len_granules(), idx)) - }) - .collect(); - } - /// Get a reference to fixed-len row data. /// /// Used in benchmarks. @@ -233,16 +207,6 @@ impl Pages { } } - /// Reserve a new, initially empty page. - // TODO(delete-free-page): Determine what to do with this method. - // It doesn't really make sense in a world where `Pages` doesn't contain empty `Page`s, - // but it's only used for tests and benches. - pub fn reserve_empty_page(&mut self, pool: &PagePool, fixed_row_size: Size) -> Result { - let idx = self.allocate_new_page(pool, fixed_row_size)?; - self.record_page_non_full(idx, fixed_row_size); - Ok(idx) - } - /// Call `f` with a reference to a page which satisfies /// `page.has_space_for_row(fixed_row_size, num_var_len_granules)`. pub fn with_page_to_insert_row( @@ -450,100 +414,6 @@ impl Pages { } } - /// Materialize a view of rows in `self` for which the `filter` returns `true`. - /// - /// # Safety - /// - /// - The `var_len_visitor` will visit the same set of `VarLenRef`s in the row - /// as the visitor provided to all other methods on `self`. - /// - /// - The `fixed_row_size` is consistent with the `var_len_visitor` - /// and is equal to the value provided to all other methods on `self`. - // FIXME: this method appears not to correctly set `non_full_pages` on the result. - // It is also unused except for benchmarks, so it may be best to just remove it. - pub unsafe fn copy_filter( - &self, - var_len_visitor: &impl VarLenMembers, - fixed_row_size: Size, - mut blob_policy: Option<&mut impl FnMut(BlobHash)>, - mut filter: impl FnMut(&Page, PageOffset) -> bool, - ) -> Self { - // Build a new container to hold the materialized view. - // Push pages into it later. - let mut partial_copied_pages = Self::default(); - - // A destination page that was not filled entirely, - // or `None` if it's time to allocate a new destination page. - let mut partial_page = None; - - // Copy each page. - for from_page in self.pages.iter().filter_map(|page| page.as_deref()) { - // You may require multiple calls to `Page::copy_starting_from` - // if `partial_page` fills up; - // the first call starts from 0. - let mut copy_starting_from = Some(PageOffset(0)); - - // While there are unprocessed rows in `from_page`, - while let Some(next_offset) = copy_starting_from.take() { - // Grab the `partial_page` or allocate a new one. - let mut to_page = partial_page.take().unwrap_or_else(|| Page::new(fixed_row_size)); - - // Copy as many rows as will fit in `to_page`. - // - // SAFETY: - // - // - The `var_len_visitor` will visit the same set of `VarLenRef`s in the row - // as the visitor provided to all other methods on `self`. - // The `to_page` uses the same visitor as the `from_page`. - // - // - The `fixed_row_size` is consistent with the `var_len_visitor` - // and is equal to the value provided to all other methods on `self`, - // as promised by the caller. - // The newly made `to_page` uses the same `fixed_row_size` as the `from_page`. - // - // - The `next_offset` is either 0, - // which is always a valid starting offset for any row size, - // or it came from `copy_filter_into` in a previous iteration, - // which, given that `fixed_row_size` was valid, - // always returns a valid starting offset in case of `Continue(_)`. - let cfi_ret = unsafe { - from_page.copy_filter_into( - next_offset, - &mut to_page, - fixed_row_size, - var_len_visitor, - blob_policy.as_mut(), - &mut filter, - ) - }; - copy_starting_from = if let ControlFlow::Continue(continue_point) = cfi_ret { - // If `to_page` couldn't fit all of `from_page`, - // repeat the `while_let` loop to copy the rest. - Some(continue_point) - } else { - // If `to_page` fit all of `from_page`, we can move on. - None - }; - - // If `from_page` finished copying into `to_page`, then `to_page` may have extra room. - // - // If `copy_filtered_into` returns `Some`, - // that means at least one row didn't have space in `to_page`, - // so we must consider `to_page` full. - // - // Note that this is distinct from `Page::is_full`, - // as that method considers the optimistic case of a row with no var-len members. - if copy_starting_from.is_none() { - partial_page = Some(to_page); - } else { - partial_copied_pages.pages.push(Some(to_page)); - } - } - } - - partial_copied_pages - } - /// Set this [`Pages`]' contents to be the `pages`. /// /// Used when restoring from a snapshot. diff --git a/crates/table/src/table.rs b/crates/table/src/table.rs index 75c864d5e49..a3f10125620 100644 --- a/crates/table/src/table.rs +++ b/crates/table/src/table.rs @@ -2589,21 +2589,23 @@ pub(crate) mod test { let mut table = Table::new(schema.into(), SquashedOffset::COMMITTED_STATE); let pool = PagePool::new_for_test(); + let blob_store = &mut NullBlobStore; let cols = ColList::new(0.into()); let algo = BTreeAlgorithm { columns: cols.clone() }.into(); let index = table.new_index(&algo, true).unwrap(); // SAFETY: Index was derived from `table`. - unsafe { table.insert_index(&NullBlobStore, index_schema.index_id, index) }.unwrap(); + unsafe { table.insert_index(blob_store, index_schema.index_id, index) }.unwrap(); // Reserve a page so that we can check the hash. - let pi = table.inner.pages.reserve_empty_page(&pool, table.row_size()).unwrap(); + let (_, row_ref) = table.insert(&pool, blob_store, &product![i32::MAX, i32::MAX]).unwrap(); + let pi = row_ref.pointer().page_index(); let hash_pre_ins = hash_unmodified_save_get(table.inner.pages.get_mut(pi).expect("reserved page to be present")); // Insert the row (0, 0). table - .insert(&pool, &mut NullBlobStore, &product![0i32, 0i32]) + .insert(&pool, blob_store, &product![0i32, 0i32]) .expect("Initial insert failed"); // Inserting cleared the hash. @@ -2612,7 +2614,7 @@ pub(crate) mod test { assert_ne!(hash_pre_ins, hash_post_ins); // Try to insert the row (0, 1), and assert that we get the expected error. - match table.insert(&pool, &mut NullBlobStore, &product![0i32, 1i32]) { + match table.insert(&pool, blob_store, &product![0i32, 1i32]) { Ok(_) => panic!("Second insert with same unique value succeeded"), Err(InsertError::IndexError(UniqueConstraintViolation { constraint_name, From 524f6902ba3b16ff4370f9b085bf0f7978c78ecd Mon Sep 17 00:00:00 2001 From: Phoebe Goldman Date: Fri, 21 Aug 2026 21:10:22 -0400 Subject: [PATCH 08/10] Count absent pages during snapshot replay towards a different metric Per Joshua's review, this commit adds a new metric, `spacetime_replay_snapshot_num_absent_pages`. Pages which aren't read from files due to having the all-zeroes hash are counted towards that metric and not towards the existing `spacetime_replay_snapshot_num_objects_read`. --- crates/engine/src/metrics.rs | 6 ++++ crates/engine/src/relational_db.rs | 5 ++++ crates/snapshot/src/lib.rs | 47 +++++++++++++++++++++--------- 3 files changed, 44 insertions(+), 14 deletions(-) diff --git a/crates/engine/src/metrics.rs b/crates/engine/src/metrics.rs index dd57244e9a4..d2e09428c1b 100644 --- a/crates/engine/src/metrics.rs +++ b/crates/engine/src/metrics.rs @@ -28,6 +28,12 @@ metrics_group!( #[labels(db: Identity, kind: str)] pub replay_snapshot_num_objects_read: IntGaugeVec, + #[name = spacetime_replay_snapshot_num_absent_pages] + #[help = "Number of pages marked absent and skipped during snapshot replay"] + // Not labeled with `kind` as blobs and metadata will never have absent files. + #[labels(db: Identity)] + pub replay_snapshot_num_absent_pages: IntGaugeVec, + #[name = spacetime_replay_snapshot_bytes_read_from_disk] #[help = "Number of snapshot bytes read from disk during replay"] #[labels(db: Identity, kind: str)] diff --git a/crates/engine/src/relational_db.rs b/crates/engine/src/relational_db.rs index 933949b666b..491b203803c 100644 --- a/crates/engine/src/relational_db.rs +++ b/crates/engine/src/relational_db.rs @@ -534,6 +534,11 @@ impl RelationalDB { let elapsed_time = start.elapsed(); + ENGINE_METRICS + .replay_snapshot_num_absent_pages + .with_label_values(database_identity) + .set(u64_to_i64(snapshot.read_metrics.absent_pages)); + for (kind, metrics) in snapshot.read_metrics.iter() { ENGINE_METRICS .replay_snapshot_read_time_seconds diff --git a/crates/snapshot/src/lib.rs b/crates/snapshot/src/lib.rs index d4673276620..76df6e08f6e 100644 --- a/crates/snapshot/src/lib.rs +++ b/crates/snapshot/src/lib.rs @@ -97,10 +97,25 @@ pub struct SnapshotReadMetrics { pub metadata: SnapshotReadKindMetrics, pub page: SnapshotReadKindMetrics, pub blob: SnapshotReadKindMetrics, + /// Number of files skipped due to having [`ZERO_HASH_DENOTING_ABSENT_PAGE`] as their hash. + /// + /// This is only meaningful for pages, so is stored separate from the other metrics. + pub absent_pages: u64, } pub type TablePages = Vec>>; +fn table_present_pages(pages: &[blake3::Hash]) -> impl Iterator { + pages + .iter() + .copied() + .filter(|hash| *hash != ZERO_HASH_DENOTING_ABSENT_PAGE) +} + +fn table_num_present_pages(pages: &[blake3::Hash]) -> u64 { + table_present_pages(pages).count() as u64 +} + impl SnapshotReadMetrics { pub fn iter(&self) -> impl Iterator + '_ { [("metadata", &self.metadata), ("page", &self.page), ("blob", &self.blob)].into_iter() @@ -684,13 +699,7 @@ impl Snapshot { + self .tables .iter() - .map(|table| { - table - .pages - .iter() - .filter(|hash| **hash != ZERO_HASH_DENOTING_ABSENT_PAGE) - .count() - }) + .map(|table| table_num_present_pages(&table.pages) as usize) .sum::() } @@ -700,12 +709,7 @@ impl Snapshot { self.blobs .iter() .map(|b| blake3::Hash::from_bytes(b.hash.data)) - .chain(self.tables.iter().flat_map(|t| { - t.pages - .iter() - .filter(|hash| **hash != ZERO_HASH_DENOTING_ABSENT_PAGE) - .copied() - })) + .chain(self.tables.iter().flat_map(|t| table_present_pages(&t.pages))) } /// Obtain an iterator over the [`PathBuf`]s of all objects @@ -1071,7 +1075,22 @@ impl SnapshotRepository { ..Default::default() }; read_metrics.blob.files = snapshot.blobs.len() as u64; - read_metrics.page.files = snapshot.tables.iter().map(|table| table.pages.len() as u64).sum(); + read_metrics.page.files = snapshot + .tables + .iter() + .map(|table| table_num_present_pages(&table.pages)) + .sum(); + read_metrics.absent_pages = snapshot + .tables + .iter() + .map(|table| { + table + .pages + .iter() + .filter(|page| **page == ZERO_HASH_DENOTING_ABSENT_PAGE) + .count() as u64 + }) + .sum(); if snapshot.magic != MAGIC { return Err(SnapshotError::BadMagic { From 93902957295f25d3418dc1657654efaa17b95677 Mon Sep 17 00:00:00 2001 From: Phoebe Goldman Date: Fri, 21 Aug 2026 21:10:22 -0400 Subject: [PATCH 09/10] Count absent pages during snapshot replay towards a different metric Per Joshua's review, this commit adds a new metric, `spacetime_replay_snapshot_num_absent_pages`. Pages which aren't read from files due to having the all-zeroes hash are counted towards that metric and not towards the existing `spacetime_replay_snapshot_num_objects_read`. --- crates/engine/src/metrics.rs | 6 ++++ crates/engine/src/relational_db.rs | 5 ++++ crates/snapshot/src/lib.rs | 47 +++++++++++++++++++++--------- 3 files changed, 44 insertions(+), 14 deletions(-) diff --git a/crates/engine/src/metrics.rs b/crates/engine/src/metrics.rs index dd57244e9a4..d2e09428c1b 100644 --- a/crates/engine/src/metrics.rs +++ b/crates/engine/src/metrics.rs @@ -28,6 +28,12 @@ metrics_group!( #[labels(db: Identity, kind: str)] pub replay_snapshot_num_objects_read: IntGaugeVec, + #[name = spacetime_replay_snapshot_num_absent_pages] + #[help = "Number of pages marked absent and skipped during snapshot replay"] + // Not labeled with `kind` as blobs and metadata will never have absent files. + #[labels(db: Identity)] + pub replay_snapshot_num_absent_pages: IntGaugeVec, + #[name = spacetime_replay_snapshot_bytes_read_from_disk] #[help = "Number of snapshot bytes read from disk during replay"] #[labels(db: Identity, kind: str)] diff --git a/crates/engine/src/relational_db.rs b/crates/engine/src/relational_db.rs index 933949b666b..491b203803c 100644 --- a/crates/engine/src/relational_db.rs +++ b/crates/engine/src/relational_db.rs @@ -534,6 +534,11 @@ impl RelationalDB { let elapsed_time = start.elapsed(); + ENGINE_METRICS + .replay_snapshot_num_absent_pages + .with_label_values(database_identity) + .set(u64_to_i64(snapshot.read_metrics.absent_pages)); + for (kind, metrics) in snapshot.read_metrics.iter() { ENGINE_METRICS .replay_snapshot_read_time_seconds diff --git a/crates/snapshot/src/lib.rs b/crates/snapshot/src/lib.rs index d4673276620..76df6e08f6e 100644 --- a/crates/snapshot/src/lib.rs +++ b/crates/snapshot/src/lib.rs @@ -97,10 +97,25 @@ pub struct SnapshotReadMetrics { pub metadata: SnapshotReadKindMetrics, pub page: SnapshotReadKindMetrics, pub blob: SnapshotReadKindMetrics, + /// Number of files skipped due to having [`ZERO_HASH_DENOTING_ABSENT_PAGE`] as their hash. + /// + /// This is only meaningful for pages, so is stored separate from the other metrics. + pub absent_pages: u64, } pub type TablePages = Vec>>; +fn table_present_pages(pages: &[blake3::Hash]) -> impl Iterator { + pages + .iter() + .copied() + .filter(|hash| *hash != ZERO_HASH_DENOTING_ABSENT_PAGE) +} + +fn table_num_present_pages(pages: &[blake3::Hash]) -> u64 { + table_present_pages(pages).count() as u64 +} + impl SnapshotReadMetrics { pub fn iter(&self) -> impl Iterator + '_ { [("metadata", &self.metadata), ("page", &self.page), ("blob", &self.blob)].into_iter() @@ -684,13 +699,7 @@ impl Snapshot { + self .tables .iter() - .map(|table| { - table - .pages - .iter() - .filter(|hash| **hash != ZERO_HASH_DENOTING_ABSENT_PAGE) - .count() - }) + .map(|table| table_num_present_pages(&table.pages) as usize) .sum::() } @@ -700,12 +709,7 @@ impl Snapshot { self.blobs .iter() .map(|b| blake3::Hash::from_bytes(b.hash.data)) - .chain(self.tables.iter().flat_map(|t| { - t.pages - .iter() - .filter(|hash| **hash != ZERO_HASH_DENOTING_ABSENT_PAGE) - .copied() - })) + .chain(self.tables.iter().flat_map(|t| table_present_pages(&t.pages))) } /// Obtain an iterator over the [`PathBuf`]s of all objects @@ -1071,7 +1075,22 @@ impl SnapshotRepository { ..Default::default() }; read_metrics.blob.files = snapshot.blobs.len() as u64; - read_metrics.page.files = snapshot.tables.iter().map(|table| table.pages.len() as u64).sum(); + read_metrics.page.files = snapshot + .tables + .iter() + .map(|table| table_num_present_pages(&table.pages)) + .sum(); + read_metrics.absent_pages = snapshot + .tables + .iter() + .map(|table| { + table + .pages + .iter() + .filter(|page| **page == ZERO_HASH_DENOTING_ABSENT_PAGE) + .count() as u64 + }) + .sum(); if snapshot.magic != MAGIC { return Err(SnapshotError::BadMagic { From 1d758fbde708475c478bdca61653e778608557ef Mon Sep 17 00:00:00 2001 From: Phoebe Goldman Date: Wed, 2 Sep 2026 13:48:18 -0400 Subject: [PATCH 10/10] Delete unused method --- crates/table/src/table.rs | 5 ----- 1 file changed, 5 deletions(-) diff --git a/crates/table/src/table.rs b/crates/table/src/table.rs index 39d5aa1a50c..97b6bbde683 100644 --- a/crates/table/src/table.rs +++ b/crates/table/src/table.rs @@ -2505,11 +2505,6 @@ impl Table { &self.inner.pages } - #[cfg(test)] - fn pages_mut(&mut self) -> &mut Pages { - &mut self.inner.pages - } - /// Iterates over each [`Page`] in this table, ensuring that its hash is computed before yielding it. /// /// Used when capturing a snapshot.