use std::collections::BTreeSet;
use crate::Result;
use crate::errors::PagedbError;
use crate::pager::freelist::{self, CHAIN_METADATA_CID};
use crate::vfs::types::OpenMode;
use crate::vfs::{Vfs, VfsFile};
use super::core::{Db, WriterState};
pub(super) struct StagedFreeList {
pub free_list_root_page_id: u64,
pub next_page_id: u64,
}
impl<V: Vfs + Clone> Db<V> {
pub(super) async fn stage_reclaimed_free_list(
&self,
state: &WriterState,
reclaimed: &BTreeSet<u64>,
target_page_ids: &BTreeSet<u64>,
freeing_commit_id: u64,
alloc_cursor: u64,
) -> Result<StagedFreeList> {
let (existing, old_chain_pages) =
freelist::read_chain(&self.pager, self.realm_id, state.free_list_root_page_id).await?;
let reallocated = existing
.iter()
.any(|(_, page_id)| target_page_ids.contains(page_id));
let hosted_on_target_pages = old_chain_pages
.iter()
.any(|page_id| target_page_ids.contains(page_id));
if reclaimed.is_empty() && !reallocated && !hosted_on_target_pages {
return Ok(StagedFreeList {
free_list_root_page_id: state.free_list_root_page_id,
next_page_id: alloc_cursor,
});
}
let mut entries: Vec<(u64, u64)> =
Vec::with_capacity(existing.len() + reclaimed.len() + old_chain_pages.len());
for (freed_at, page_id) in existing {
if target_page_ids.contains(&page_id) {
continue;
}
if reclaimed.contains(&page_id) {
return Err(PagedbError::corruption(
crate::errors::CorruptionDetail::CatalogRowInvalid {
field: "free-list entry names a page the base commit still reaches",
},
));
}
entries.push((freed_at, page_id));
}
let floor = self.reclamation_floor(state).await?;
let host_candidates: Vec<u64> = entries
.iter()
.filter(|(freed_at, _)| *freed_at < floor)
.map(|(_, page_id)| *page_id)
.collect();
entries.extend(
reclaimed
.iter()
.map(|&page_id| (freeing_commit_id, page_id)),
);
entries.extend(
old_chain_pages
.iter()
.filter(|page_id| !target_page_ids.contains(page_id))
.map(|&page_id| (CHAIN_METADATA_CID, page_id)),
);
let (free_list_root_page_id, next_page_id) = freelist::rewrite_chain(
&self.pager,
self.realm_id,
self.page_size,
entries,
host_candidates,
alloc_cursor,
0,
)
.await?;
Ok(StagedFreeList {
free_list_root_page_id,
next_page_id,
})
}
pub(super) async fn ensure_image_covers_cursor(
&self,
image_path: &str,
next_page_id: u64,
) -> Result<()> {
let page_size = u64::try_from(self.page_size)
.map_err(|_| PagedbError::arithmetic_overflow("page size"))?;
let required = next_page_id
.checked_mul(page_size)
.ok_or_else(|| PagedbError::arithmetic_overflow("main.db extent"))?;
let mut file = self.vfs.open(image_path, OpenMode::CreateOrOpen).await?;
if file.len().await? >= required {
return Ok(());
}
file.set_len(required).await?;
file.sync().await
}
async fn reclamation_floor(&self, state: &WriterState) -> Result<u64> {
let min_reader = {
let readers = self.tracked_readers.lock();
readers.iter().map(|reader| reader.commit_id.0).min()
};
let history_floor = self
.oldest_retained_history_commit(state.commit_history_root_page_id, state.next_page_id)
.await?;
Ok(min_reader
.unwrap_or(u64::MAX)
.min(history_floor.map_or(u64::MAX, |commit| commit.saturating_add(1))))
}
}