Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
Original file line number Diff line number Diff line change
Expand Up @@ -30,7 +30,7 @@ impl Generator {
} else {
quote! {
/// Returns the physical size of this table's `.wt.data` file.
/// Persisted vacuum currently compacts logical/in-memory pages
/// Persisted vacuum makes freed pages reusable across reloads,
/// but does not truncate this file.
pub async fn persisted_data_file_size_bytes(&self) -> std::io::Result<u64> {
self.1.persisted_data_file_size_bytes().await
Expand Down
4 changes: 4 additions & 0 deletions src/persistence/engine.rs
Original file line number Diff line number Diff line change
Expand Up @@ -285,6 +285,10 @@ where
Ok(())
}

async fn reclaim_data_pages(&mut self, page_ids: Vec<data_bucket::page::PageId>) -> eyre::Result<()> {
self.data.reclaim_data_pages(page_ids).await
}

async fn ensure_schema(
&mut self,
row_schema: Vec<(String, String)>,
Expand Down
11 changes: 11 additions & 0 deletions src/persistence/mod.rs
Original file line number Diff line number Diff line change
@@ -1,5 +1,7 @@
use std::future::Future;

use data_bucket::page::PageId;

use crate::persistence::operation::BatchOperation;

pub use engine::DiskConfig;
Expand Down Expand Up @@ -57,6 +59,15 @@ pub trait PersistenceEngine<PrimaryKeyGenState, PrimaryKey, SecondaryIndexEvents
batch_op: BatchOperation<PrimaryKeyGenState, PrimaryKey, SecondaryIndexEvents, AvailableIndexes>,
) -> impl Future<Output = eyre::Result<()>> + Send;

/// Persists whole data pages made reusable by vacuum.
///
/// The persistence task invokes this only after every row move queued
/// before the reclamation barrier has reached the engine. Custom engines
/// that do not manage data pages may keep the default no-op.
fn reclaim_data_pages(&mut self, _page_ids: Vec<PageId>) -> impl Future<Output = eyre::Result<()>> + Send {
async { Ok(()) }
}

/// Installs the generated table schema and rejects a non-empty schema that
/// belongs to a different table shape.
/// Custom engines may keep the default no-op when they do not expose
Expand Down
99 changes: 99 additions & 0 deletions src/persistence/space/data.rs
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,63 @@ impl<PkGenState, const INNER_PAGE_SIZE: usize, const PAGE_SIZE: u32> SpaceData<P
self.data_file.write_all(bytes.as_ref()).await?;
Ok(())
}

/// Removes written byte ranges from the durable free-range list.
///
/// This is persisted before the corresponding row bytes. A crash between
/// those writes can leak reusable space, but can never leave a live row
/// described as free and eligible to be overwritten after reload.
fn consume_reusable_ranges(&mut self, used_links: impl IntoIterator<Item = Link>) -> bool {
let mut changed = false;

for used in used_links {
let used_start = u64::from(used.offset);
let used_end = used_start + u64::from(used.length);
let mut remaining = Vec::with_capacity(self.info.inner.empty_links_list.len() + 1);

for free in self.info.inner.empty_links_list.drain(..) {
if free.page_id != used.page_id {
remaining.push(free);
continue;
}

let free_start = u64::from(free.offset);
let free_end = free_start + u64::from(free.length);
let overlap_start = free_start.max(used_start);
let overlap_end = free_end.min(used_end);
if overlap_start >= overlap_end {
remaining.push(free);
continue;
}

changed = true;
if free_start < overlap_start {
remaining.push(Link {
page_id: free.page_id,
offset: free.offset,
length: (overlap_start - free_start) as u32,
});
}
if overlap_end < free_end {
remaining.push(Link {
page_id: free.page_id,
offset: overlap_end as u32,
length: (free_end - overlap_end) as u32,
});
}
}

self.info.inner.empty_links_list = remaining;
}

if changed {
self.info.inner.empty_links_list.sort_by_key(|link| {
let page_id: u32 = link.page_id.into();
(page_id, link.offset)
});
}
changed
}
}

impl<PkGenState, const INNER_PAGE_SIZE: usize, const PAGE_SIZE: u32> SpaceDataOps<PkGenState>
Expand Down Expand Up @@ -100,6 +157,9 @@ where
}

async fn save_data(&mut self, link: Link, bytes: &[u8]) -> eyre::Result<()> {
if self.consume_reusable_ranges([link]) {
self.save_info().await?;
}
if link.page_id > self.last_page_id.into() {
let mut page = GeneralPage {
header: GeneralHeader::new(link.page_id, PageType::Data, 0.into()),
Expand All @@ -123,6 +183,11 @@ where
}

async fn save_batch_data(&mut self, batch_data: BatchData) -> eyre::Result<()> {
let used_links = batch_data.values().flat_map(|ops| ops.iter().map(|(link, _)| *link));
if self.consume_reusable_ranges(used_links.collect::<Vec<_>>()) {
self.save_info().await?;
}

let page_ids = batch_data.keys().map(|id| (*id).into()).collect::<Vec<_>>();
let ids_to_create = page_ids
.iter()
Expand Down Expand Up @@ -182,6 +247,40 @@ where
Ok(())
}

async fn reclaim_data_pages(&mut self, page_ids: Vec<data_bucket::page::PageId>) -> eyre::Result<()> {
let mut page_ids = page_ids
.into_iter()
.filter(|page_id| {
let id: u32 = (*page_id).into();
id != 0 && id <= self.last_page_id
})
.collect::<Vec<_>>();
page_ids.sort_unstable();
page_ids.dedup();

if page_ids.is_empty() {
return Ok(());
}

self.info
.inner
.empty_links_list
.retain(|link| !page_ids.contains(&link.page_id));
self.info
.inner
.empty_links_list
.extend(page_ids.into_iter().map(|page_id| Link {
page_id,
offset: 0,
length: INNER_PAGE_SIZE as u32,
}));
self.info.inner.empty_links_list.sort_by_key(|link| {
let page_id: u32 = link.page_id.into();
(page_id, link.offset)
});
self.save_info().await
}

fn get_mut_info(&mut self) -> &mut GeneralPage<SpaceInfoPage<PkGenState>> {
&mut self.info
}
Expand Down
1 change: 1 addition & 0 deletions src/persistence/space/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -33,6 +33,7 @@ pub trait SpaceDataOps<PkGenState> {
fn bootstrap(file: &mut File, table_name: String, version: u32) -> impl Future<Output = eyre::Result<()>> + Send;
fn save_data(&mut self, link: Link, bytes: &[u8]) -> impl Future<Output = eyre::Result<()>> + Send;
fn save_batch_data(&mut self, batch_data: BatchData) -> impl Future<Output = eyre::Result<()>> + Send;
fn reclaim_data_pages(&mut self, page_ids: Vec<PageId>) -> impl Future<Output = eyre::Result<()>> + Send;
fn get_mut_info(&mut self) -> &mut GeneralPage<SpaceInfoPage<PkGenState>>;
fn save_info(&mut self) -> impl Future<Output = eyre::Result<()>> + Send;
}
Expand Down
Loading
Loading