use crate::Result;
use crate::btree::BTree;
use crate::pager::header::{ActiveSlot, commit_header};
use crate::txn::db::{Db, WriterState};
use crate::vfs::types::OpenMode;
use crate::vfs::{Vfs, VfsFile};
use super::helpers::{collect_all_pairs, collect_catalog_split, make_header_fields};
pub(super) struct RepackOutcome {
pub pages_reclaimed: u64,
pub bytes_truncated: u64,
}
pub(super) struct PendingSwap {
new_root: u64,
new_cat_root: u64,
new_next: u64,
counter_anchor: u64,
old_next: u64,
}
fn scratch_path(db: &Db<impl Vfs + Clone>) -> String {
format!("{}.compact", db.main_db_path)
}
pub(super) async fn atomic_dense_repack<V: Vfs + Clone>(
db: &Db<V>,
state: &mut WriterState,
visibility: &tokio::sync::RwLockWriteGuard<'_, ()>,
) -> Result<RepackOutcome> {
let scratch = scratch_path(db);
db.vfs.remove(&scratch).await.ok();
match build_scratch(db, state, &scratch).await {
Ok(pending) => commit_swap(db, state, visibility, &scratch, pending).await,
Err(e) => {
db.pager.reset_main_pages();
db.vfs.remove(&scratch).await.ok();
Err(e)
}
}
}
pub(super) async fn build_scratch<V: Vfs + Clone>(
db: &Db<V>,
state: &WriterState,
scratch: &str,
) -> Result<PendingSwap> {
let old_next = state.next_page_id;
let main_pairs = if state.root_page_id != 0 {
let old = BTree::open(
db.pager.clone(),
db.realm_id,
state.root_page_id,
old_next,
db.page_size,
);
collect_all_pairs(&old).await?
} else {
Vec::new()
};
let cs_prefix = crate::catalog::codec::CatalogRowKind::CompactionState as u8;
let cat_rows: Vec<(Vec<u8>, Vec<u8>)> = collect_catalog_split(&db.pager, db.realm_id, state)
.await?
.into_iter()
.filter(|(k, _)| k.first() != Some(&cs_prefix))
.collect();
let mut new_main = BTree::open(db.pager.clone(), db.realm_id, 0, 4, db.page_size);
new_main.bulk_load(main_pairs).await?;
let new_root = new_main.root_page_id();
let after_main = new_main.next_page_id();
let mut new_cat = BTree::open(db.pager.clone(), db.realm_id, 0, after_main, db.page_size);
new_cat.bulk_load(cat_rows).await?;
let new_cat_root = new_cat.root_page_id();
let new_next = new_cat.next_page_id();
{
let mut f = db.vfs.open(scratch, OpenMode::CreateOrOpen).await?;
f.set_len(new_next.saturating_mul(db.page_size as u64))
.await?;
f.sync().await?;
}
db.pager.flush_main_to(db.realm_id, scratch).await?;
let new_commit_id = state.latest_commit_id + 1;
let new_seq = state.seq + 1;
let counter_anchor = db.pager.pending_anchor();
let fields = make_header_fields(
db,
state,
new_commit_id,
new_seq,
counter_anchor,
new_root,
new_cat_root,
new_next,
0,
);
let hk_clone = { db.hk.read().clone() };
commit_header(
&*db.vfs,
scratch,
&hk_clone,
&fields,
ActiveSlot::B,
db.page_size,
)
.await?;
Ok(PendingSwap {
new_root,
new_cat_root,
new_next,
counter_anchor,
old_next,
})
}
async fn commit_swap<V: Vfs + Clone>(
db: &Db<V>,
state: &mut WriterState,
visibility: &tokio::sync::RwLockWriteGuard<'_, ()>,
scratch: &str,
pending: PendingSwap,
) -> Result<RepackOutcome> {
db.pager.close_main_handle().await;
let new_commit_id = state.latest_commit_id + 1;
if db.vfs.rename(scratch, &db.main_db_path).await.is_err() {
let commit = crate::CommitId(new_commit_id);
db.poison(commit);
return Err(crate::errors::PagedbError::durably_committed_but_unpublished(commit));
}
state.root_page_id = pending.new_root;
state.catalog_root_page_id = pending.new_cat_root;
state.catalog_root_txn_id = new_commit_id;
state.next_page_id = pending.new_next;
state.active_slot = ActiveSlot::A;
state.seq += 1;
state.latest_commit_id = new_commit_id;
state.commit_history_root_page_id = 0;
state.commit_history_root_version = 0;
state.free_list_root_page_id = 0;
db.pager.reset_main_pages();
if db.vfs.sync_dir(db.main_db_parent_dir()).await.is_err() {
let commit = crate::CommitId(new_commit_id);
db.poison(commit);
return Err(crate::errors::PagedbError::durably_committed_but_unpublished(commit));
}
let _ = db
.finish_durable_commit_visible(
visibility,
state,
crate::CommitId(new_commit_id),
pending.counter_anchor,
&[],
)
.await?;
let pages_reclaimed = pending.old_next.saturating_sub(pending.new_next);
Ok(RepackOutcome {
pages_reclaimed,
bytes_truncated: pages_reclaimed.saturating_mul(db.page_size as u64),
})
}
#[cfg(test)]
mod tests {
use super::*;
use crate::vfs::memory::MemVfs;
use crate::{Db, RealmId};
#[tokio::test(flavor = "current_thread")]
async fn crash_before_rename_leaves_old_store_intact() {
let vfs = MemVfs::new();
let kek = [0x42u8; 32];
let realm = RealmId::new([0x7u8; 16]);
let big = vec![0xABu8; 2048]; let n = 50u32;
{
let db = Db::open_internal(vfs.clone(), kek, 4096, realm)
.await
.unwrap();
{
let mut w = db.begin_write().await.unwrap();
for i in 0..n {
w.put(format!("k-{i:04}").as_bytes(), &big).await.unwrap();
}
w.commit().await.unwrap();
}
{
let state = db.writer.lock().await;
let scratch = scratch_path(&db);
build_scratch(&db, &state, &scratch).await.unwrap();
}
}
let db2 = Db::open_existing(vfs.clone(), kek, 4096, realm)
.await
.unwrap();
let r = db2.begin_read().await.unwrap();
for i in 0..n {
let key = format!("k-{i:04}");
assert_eq!(
r.get(key.as_bytes()).await.unwrap().as_deref(),
Some(big.as_slice()),
"value {key} lost after a crash before the compaction rename"
);
}
assert!(
vfs.open("/main.db.compact", OpenMode::Read).await.is_err(),
"orphaned compaction scratch should be removed on open"
);
}
}