use std::sync::Arc;
use crate::btree::BTree;
use crate::catalog::codec::{Catalog, CatalogRowKind, SegmentMeta};
use crate::errors::PagedbError;
use crate::pager::anchor::HeaderCursor;
use crate::pager::header::commit_header;
use crate::pager::structural_header::MainDbHeaderFields;
use crate::txn::db::{Db, WriterState};
use crate::txn::write::SegmentSideEffect;
use crate::vfs::Vfs;
use crate::{CommitId, Result};
const REPACK_RECORD_BATCH: usize = 256;
pub(super) async fn stream_dense_tree<V: Vfs + Clone>(
db: &Db<V>,
scratch: &str,
source_root: u64,
source_next: u64,
dest: &mut BTree<V>,
keep: impl Fn(&[u8]) -> bool + Send,
) -> Result<()> {
let source = BTree::open(
db.pager.clone(),
db.realm_id,
source_root,
source_next,
db.page_size,
);
let mut loader = dest.bulk_loader()?;
let mut cursor: Vec<u8> = Vec::new();
loop {
db.pager.flush_main_to(db.realm_id, scratch).await?;
db.pager.reset_main_pages();
let batch = source
.collect_batch_from(&cursor, REPACK_RECORD_BATCH)
.await?;
let Some((last_key, _)) = batch.last() else {
break;
};
cursor.clear();
cursor.extend_from_slice(last_key);
cursor.push(0);
let exhausted = batch.len() < REPACK_RECORD_BATCH;
for (key, value) in batch {
if !keep(&key) {
continue;
}
loader.push(key.to_vec(), value).await?;
if db.pager.main_dirty_at_budget() {
db.pager.flush_main_to(db.realm_id, scratch).await?;
db.pager.reset_main_pages();
}
}
if exhausted {
break;
}
}
loader.finish().await
}
pub(super) const SEGMENT_ROW_BATCH: usize = 256;
pub(super) async fn segment_rows_from<V: Vfs + Clone>(
pager: &Arc<crate::pager::Pager<V>>,
realm_id: crate::RealmId,
state: &WriterState,
cursor: &[u8],
) -> Result<Vec<(Vec<u8>, SegmentMeta)>> {
if state.catalog_root_page_id == 0 {
return Ok(Vec::new());
}
let tree = BTree::open(
pager.clone(),
realm_id,
state.catalog_root_page_id,
state.next_page_id,
pager.page_size(),
);
let prefix = [CatalogRowKind::Segment as u8];
let rows = tree
.collect_prefix_batch_from(&prefix, cursor, SEGMENT_ROW_BATCH)
.await?;
let mut out = Vec::with_capacity(rows.len());
for (key, value) in rows {
let meta = Catalog::decode_segment_meta(&value)?;
out.push((key.to_vec(), meta));
}
Ok(out)
}
pub(super) async fn replace_segment_compact<V: Vfs + Clone>(
db: &Db<V>,
state: &mut WriterState,
visibility: &tokio::sync::RwLockWriteGuard<'_, ()>,
name: &str,
old_segment_id: &[u8; 16],
new_meta: &SegmentMeta,
) -> Result<()> {
let key = Catalog::segment_key(db.realm_id, name.as_bytes())?;
let value = Catalog::encode_segment_meta(new_meta);
let mut cat_tree = BTree::open(
db.pager.clone(),
db.realm_id,
state.catalog_root_page_id,
state.next_page_id,
db.page_size,
);
cat_tree.put(&key, &value).await?;
cat_tree.flush().await?;
let new_cat_root = cat_tree.root_page_id();
let new_next = cat_tree.next_page_id().max(state.next_page_id);
let new_commit_id = state.latest_commit_id + 1;
let header_cursor = db.pager.header_cursor()?;
let new_seq = header_cursor.next_seq()?;
let counter_anchor = db.pager.pending_anchor();
let mut catalog_root_bytes = [0u8; 16];
catalog_root_bytes[..8].copy_from_slice(&new_cat_root.to_le_bytes());
catalog_root_bytes[8..].copy_from_slice(&new_commit_id.to_le_bytes());
let fields = MainDbHeaderFields {
format_version: crate::pager::structural_header::MAIN_FORMAT_VERSION,
cipher_id: db.cipher_id.as_byte(),
page_size_log2: page_size_log2(db.page_size)?,
flags: 0,
file_id: db.file_id,
kek_salt: db.kek_salt,
mk_epoch: db.mk_epoch.load(std::sync::atomic::Ordering::SeqCst),
seq: new_seq,
active_root_page_id: state.root_page_id,
active_root_txn_id: state.latest_commit_id,
counter_anchor,
commit_id: CommitId(new_commit_id),
free_list_root: [0u8; 16],
catalog_root: catalog_root_bytes,
apply_journal_root_page_id: 0,
apply_journal_root_version: 0,
commit_history_root_page_id: 0,
commit_history_root_version: 0,
restore_mode: 0,
next_page_id: new_next,
commit_retain_policy_tag: 0,
commit_retain_policy_value: 0,
realm_id: db.realm_id,
};
let hk_clone = { db.hk.read().clone() };
let new_slot = commit_header(
&*db.vfs,
&db.main_db_path,
&hk_clone,
&fields,
header_cursor.slot,
db.page_size,
)
.await?;
db.pager.note_header_written(HeaderCursor {
slot: new_slot,
seq: new_seq,
});
state.catalog_root_page_id = new_cat_root;
state.catalog_root_txn_id = new_commit_id;
state.next_page_id = new_next;
state.latest_commit_id = new_commit_id;
let effects = [
SegmentSideEffect::Tombstone {
segment_id: *old_segment_id,
tombstone_commit_id: None,
},
SegmentSideEffect::Promote {
segment_id: new_meta.segment_id,
},
];
let _ = db
.finish_durable_commit_visible(
visibility,
state,
CommitId(new_commit_id),
counter_anchor,
&effects,
)
.await?;
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub(super) fn make_header_fields<V: Vfs + Clone>(
db: &Db<V>,
state: &WriterState,
new_commit_id: u64,
new_seq: u64,
counter_anchor: u64,
new_root: u64,
new_cat_root: u64,
new_next: u64,
free_list_root_page_id: u64,
) -> MainDbHeaderFields {
let _ = state;
let mut catalog_root_bytes = [0u8; 16];
catalog_root_bytes[..8].copy_from_slice(&new_cat_root.to_le_bytes());
catalog_root_bytes[8..].copy_from_slice(&new_commit_id.to_le_bytes());
MainDbHeaderFields {
format_version: crate::pager::structural_header::MAIN_FORMAT_VERSION,
cipher_id: db.cipher_id.as_byte(),
page_size_log2: page_size_log2(db.page_size).unwrap_or(12),
flags: 0,
file_id: db.file_id,
kek_salt: db.kek_salt,
mk_epoch: db.mk_epoch.load(std::sync::atomic::Ordering::SeqCst),
seq: new_seq,
active_root_page_id: new_root,
active_root_txn_id: new_commit_id,
counter_anchor,
commit_id: CommitId(new_commit_id),
free_list_root: crate::txn::db::encode_free_list_root(free_list_root_page_id),
catalog_root: catalog_root_bytes,
apply_journal_root_page_id: 0,
apply_journal_root_version: 0,
commit_history_root_page_id: 0,
commit_history_root_version: 0,
restore_mode: 0,
next_page_id: new_next,
commit_retain_policy_tag: 0,
commit_retain_policy_value: 0,
realm_id: db.realm_id,
}
}
pub(super) fn page_size_log2(page_size: usize) -> Result<u8> {
match page_size {
4096 => Ok(12),
8192 => Ok(13),
16384 => Ok(14),
32768 => Ok(15),
65536 => Ok(16),
_ => Err(PagedbError::Unsupported),
}
}