diff --git a/src/lib.rs b/src/lib.rs index 6aeea5e..056ed34 100644 --- a/src/lib.rs +++ b/src/lib.rs @@ -19,7 +19,14 @@ pub use error::*; use index::{ProbableIndex, ZeboIndex}; use page::{DOCUMENT_INDEX_OFFSET, ZeboPage, ZeboPageHeader}; -pub use crate::page::ZeboPageReservedSpace; +pub use crate::page::{CompactPageStats, ZeboPageReservedSpace}; + +#[derive(Debug)] +pub struct CompactStats { + pub pages_compacted: u32, + pub pages_skipped: u32, + pub total_bytes_reclaimed: u64, +} #[derive(Debug, PartialEq)] pub struct ZeboInfo { @@ -167,6 +174,92 @@ impl Ok(removed) } + /// Compacts all pages except the current one by rewriting them to reclaim + /// space from deleted documents. The data region is compacted while header + /// slots (including deleted sentinel entries) are preserved. + pub fn compact(&mut self) -> Result { + let page_ids = self.index.get_page_ids()?; + let current_page_id = if self.current_page.is_some() { + Some(PageId(self.next_page_id - 1)) + } else { + None + }; + + let mut stats = CompactStats { + pages_compacted: 0, + pages_skipped: 0, + total_bytes_reclaimed: 0, + }; + + let mut buf = Vec::new(); + + for page_id in page_ids { + if Some(page_id) == current_page_id { + stats.pages_skipped += 1; + continue; + } + + let page = load_page(&self.base_dir, page_id, Mode::Read)?; + + // Check if compaction would reclaim space: compare actual file size + // against header size + total live document data + let file_size = + std::fs::metadata(self.base_dir.join(format!("page_{}.zebo", page_id.0))) + .map_err(ZeboError::OperationError)? + .len(); + let header_size = ZeboPage::header_size(page.document_limit()); + let live_data_size = page.live_data_size()?; + let compacted_size = header_size + live_data_size; + if file_size <= compacted_size { + stats.pages_skipped += 1; + continue; + } + + let temp_path = self + .base_dir + .join(format!("page_{}.zebo.compact_tmp", page_id.0)); + let original_path = self.base_dir.join(format!("page_{}.zebo", page_id.0)); + + let result = (|| -> Result { + let mut temp_file = std::fs::File::options() + .create(true) + .write(true) + .read(true) + .truncate(true) + .open(&temp_path) + .map_err(ZeboError::OperationError)?; + + let page_stats = page.compact_to_file(&mut temp_file, &mut buf)?; + Ok(page_stats) + })(); + + match result { + Ok(page_stats) => { + std::fs::rename(&temp_path, &original_path) + .map_err(ZeboError::OperationError)?; + + #[cfg(unix)] + { + let dir = std::fs::File::open(&self.base_dir) + .map_err(ZeboError::OperationError)?; + dir.sync_all().map_err(ZeboError::OperationError)?; + } + + stats.total_bytes_reclaimed += page_stats + .bytes_before + .saturating_sub(page_stats.bytes_after); + stats.pages_compacted += 1; + } + Err(e) => { + let _ = std::fs::remove_file(&temp_path); + return Err(e); + } + } + } + + Ok(stats) + } + /// Returns an iterator with the order guarantees pub fn get_documents + Clone>( &self, @@ -1406,6 +1499,216 @@ mod tests { assert_eq!(result, Some(6)); } + #[test] + fn test_compact_basic() { + let test_dir = prepare_test_dir(); + + let mut zebo: Zebo<3, 2048, u32> = Zebo::<3, 2048, _>::try_new(test_dir.clone()).unwrap(); + + // Fill page 0 (3 docs), then page 1 (3 docs), then page 2 (current) + zebo.reserve_space_for(&[(1, "aaa"), (2, "bbb"), (3, "ccc")]) + .unwrap() + .write_all() + .unwrap(); + zebo.reserve_space_for(&[(4, "ddd"), (5, "eee"), (6, "fff")]) + .unwrap() + .write_all() + .unwrap(); + zebo.reserve_space_for(&[(7, "ggg")]) + .unwrap() + .write_all() + .unwrap(); + + // Delete some docs from page 0 and page 1 + zebo.remove_documents(vec![1, 2, 5], true).unwrap(); + + // Get file sizes before compaction + let page0_size_before = std::fs::metadata(test_dir.join("page_0.zebo")) + .unwrap() + .len(); + let page2_size_before = std::fs::metadata(test_dir.join("page_2.zebo")) + .unwrap() + .len(); + + let stats = zebo.compact().unwrap(); + + assert_eq!(stats.pages_compacted, 2); + // current page (page 2) should be skipped + assert_eq!(stats.pages_skipped, 1); + assert!(stats.total_bytes_reclaimed > 0); + + // Page 0 should be smaller (2 docs deleted out of 3) + let page0_size_after = std::fs::metadata(test_dir.join("page_0.zebo")) + .unwrap() + .len(); + assert!(page0_size_after < page0_size_before); + + // Current page (page 2) should be unchanged + let page2_size_after = std::fs::metadata(test_dir.join("page_2.zebo")) + .unwrap() + .len(); + assert_eq!(page2_size_after, page2_size_before); + + // All live docs should still be retrievable + assert!(zebo.get_document(1).unwrap().is_none()); + assert!(zebo.get_document(2).unwrap().is_none()); + assert_eq!(zebo.get_document(3).unwrap().unwrap(), b"ccc"); + assert_eq!(zebo.get_document(4).unwrap().unwrap(), b"ddd"); + assert!(zebo.get_document(5).unwrap().is_none()); + assert_eq!(zebo.get_document(6).unwrap().unwrap(), b"fff"); + assert_eq!(zebo.get_document(7).unwrap().unwrap(), b"ggg"); + } + + #[test] + fn test_compact_idempotent() { + let test_dir = prepare_test_dir(); + + let mut zebo: Zebo<3, 2048, u32> = Zebo::<3, 2048, _>::try_new(test_dir.clone()).unwrap(); + + zebo.reserve_space_for(&[(1, "aaa"), (2, "bbb"), (3, "ccc")]) + .unwrap() + .write_all() + .unwrap(); + zebo.reserve_space_for(&[(4, "ddd")]) + .unwrap() + .write_all() + .unwrap(); + + zebo.remove_documents(vec![2], true).unwrap(); + + let stats1 = zebo.compact().unwrap(); + assert_eq!(stats1.pages_compacted, 1); + + // Second compact should find nothing to do + let stats2 = zebo.compact().unwrap(); + assert_eq!(stats2.pages_compacted, 0); + assert_eq!(stats2.total_bytes_reclaimed, 0); + + // Docs still accessible + assert_eq!(zebo.get_document(1).unwrap().unwrap(), b"aaa"); + assert_eq!(zebo.get_document(3).unwrap().unwrap(), b"ccc"); + } + + #[test] + fn test_compact_all_deleted_page() { + let test_dir = prepare_test_dir(); + + let mut zebo: Zebo<3, 2048, u32> = Zebo::<3, 2048, _>::try_new(test_dir.clone()).unwrap(); + + zebo.reserve_space_for(&[(1, "aaa"), (2, "bbb"), (3, "ccc")]) + .unwrap() + .write_all() + .unwrap(); + zebo.reserve_space_for(&[(4, "ddd")]) + .unwrap() + .write_all() + .unwrap(); + + // Delete all docs from page 0 + zebo.remove_documents(vec![1, 2, 3], true).unwrap(); + + let page0_size_before = std::fs::metadata(test_dir.join("page_0.zebo")) + .unwrap() + .len(); + + let stats = zebo.compact().unwrap(); + assert_eq!(stats.pages_compacted, 1); + + // Page should now be header-only (no data region) + let page0_size_after = std::fs::metadata(test_dir.join("page_0.zebo")) + .unwrap() + .len(); + assert!(page0_size_after < page0_size_before); + + // Page 1 (current) doc should still work + assert_eq!(zebo.get_document(4).unwrap().unwrap(), b"ddd"); + } + + #[test] + fn test_compact_then_insert() { + let test_dir = prepare_test_dir(); + + let mut zebo: Zebo<3, 2048, u32> = Zebo::<3, 2048, _>::try_new(test_dir.clone()).unwrap(); + + zebo.reserve_space_for(&[(1, "aaa"), (2, "bbb"), (3, "ccc")]) + .unwrap() + .write_all() + .unwrap(); + zebo.reserve_space_for(&[(4, "ddd")]) + .unwrap() + .write_all() + .unwrap(); + + zebo.remove_documents(vec![2], true).unwrap(); + zebo.compact().unwrap(); + + // Insert more docs after compaction + zebo.reserve_space_for(&[(5, "eee"), (6, "fff")]) + .unwrap() + .write_all() + .unwrap(); + + assert_eq!(zebo.get_document(1).unwrap().unwrap(), b"aaa"); + assert!(zebo.get_document(2).unwrap().is_none()); + assert_eq!(zebo.get_document(3).unwrap().unwrap(), b"ccc"); + assert_eq!(zebo.get_document(4).unwrap().unwrap(), b"ddd"); + assert_eq!(zebo.get_document(5).unwrap().unwrap(), b"eee"); + assert_eq!(zebo.get_document(6).unwrap().unwrap(), b"fff"); + } + + #[test] + fn test_compact_reload() { + let test_dir = prepare_test_dir(); + + let mut zebo: Zebo<3, 2048, u32> = Zebo::<3, 2048, _>::try_new(test_dir.clone()).unwrap(); + + zebo.reserve_space_for(&[(1, "aaa"), (2, "bbb"), (3, "ccc")]) + .unwrap() + .write_all() + .unwrap(); + zebo.reserve_space_for(&[(4, "ddd"), (5, "eee"), (6, "fff")]) + .unwrap() + .write_all() + .unwrap(); + zebo.reserve_space_for(&[(7, "ggg")]) + .unwrap() + .write_all() + .unwrap(); + + zebo.remove_documents(vec![1, 2, 5], true).unwrap(); + zebo.compact().unwrap(); + + // Drop and reload from disk + drop(zebo); + let mut zebo: Zebo<3, 2048, u32> = Zebo::<3, 2048, _>::try_new(test_dir.clone()).unwrap(); + + // All live docs should be retrievable after reload + assert!(zebo.get_document(1).unwrap().is_none()); + assert!(zebo.get_document(2).unwrap().is_none()); + assert_eq!(zebo.get_document(3).unwrap().unwrap(), b"ccc"); + assert_eq!(zebo.get_document(4).unwrap().unwrap(), b"ddd"); + assert!(zebo.get_document(5).unwrap().is_none()); + assert_eq!(zebo.get_document(6).unwrap().unwrap(), b"fff"); + assert_eq!(zebo.get_document(7).unwrap().unwrap(), b"ggg"); + + // Info should reflect correct counts + let info = zebo.get_info().unwrap(); + assert_eq!(info.document_count, 4); + + // New inserts after reload should work + zebo.reserve_space_for(&[(8, "hhh"), (9, "iii")]) + .unwrap() + .write_all() + .unwrap(); + + assert_eq!(zebo.get_document(8).unwrap().unwrap(), b"hhh"); + assert_eq!(zebo.get_document(9).unwrap().unwrap(), b"iii"); + + // Second compact after reload should be idempotent + let stats = zebo.compact().unwrap(); + assert_eq!(stats.pages_compacted, 0); + } + struct MyDoc { id: String, text: String, diff --git a/src/page.rs b/src/page.rs index 97f558b..3ed99b7 100644 --- a/src/page.rs +++ b/src/page.rs @@ -53,19 +53,25 @@ pub struct ZeboPage { } impl ZeboPage { + #[inline] + pub fn header_size(document_limit: u32) -> u64 { + // 8 bytes: doc_id.as_u64() + // 4 bytes: starting offset + // 4 bytes: bytes length + // 8 + 4 + 4 = 16 bytes per document header entry + DOCUMENT_INDEX_OFFSET + (document_limit as u64) * 16 + } + pub fn try_new( document_limit: u32, starting_document_id: u64, mut page_file: std::fs::File, ) -> Result { - // 8 bytes: doc_id.as_u64() - // 4 bytes: starting offset - // 4 bytes: bytes length - let document_header_size = (4 + 4 + 8) * (document_limit as u64); + let header_size = Self::header_size(document_limit); // We shrink the file to contain at least the document header // this because we store documents *after* the header page_file - .set_len(DOCUMENT_INDEX_OFFSET + document_header_size) + .set_len(header_size) .map_err(ZeboError::OperationError)?; // Version on first byte @@ -81,7 +87,7 @@ impl ZeboPage { .write_all_at(&[0; 4], DOCUMENT_COUNT_OFFSET) .map_err(ZeboError::OperationError)?; // Next available offset - let initial_available_offset = (DOCUMENT_INDEX_OFFSET + document_header_size) as u32; + let initial_available_offset = header_size as u32; page_file .write_all_at( &initial_available_offset.to_be_bytes(), @@ -148,6 +154,24 @@ impl ZeboPage { }) } + pub fn document_limit(&self) -> u32 { + self.document_limit + } + + pub fn live_data_size(&self) -> Result { + let mut total: u64 = 0; + for i in 0..self.next_available_header_offset as u64 { + if let Some((doc_id, offset, length)) = self.get_at(i)? { + if Self::is_uninitialized_entry(offset) || Self::is_deleted(doc_id, offset, length) + { + continue; + } + total += length as u64; + } + } + Ok(total) + } + pub fn get_document_count(&self) -> Result { let mut buf = [0; 4]; self.page_file @@ -801,6 +825,130 @@ impl ZeboPage { Ok(None) } + pub fn compact_to_file( + &self, + target_file: &mut File, + buf: &mut Vec, + ) -> Result { + let bytes_before = self + .page_file + .metadata() + .map_err(ZeboError::OperationError)? + .len(); + + let header_size = Self::header_size(self.document_limit); + target_file + .set_len(header_size) + .map_err(ZeboError::OperationError)?; + + // Write version + target_file + .write_all_at(&[Version::V1.into()], VERSION_OFFSET) + .map_err(ZeboError::OperationError)?; + // Write document limit + target_file + .write_all_at( + &self.document_limit.to_be_bytes(), + DOCUMENT_COUNT_LIMIT_OFFSET, + ) + .map_err(ZeboError::OperationError)?; + // Write starting document id + target_file + .write_all_at( + &self.starting_document_id.to_be_bytes(), + STARTING_DOCUMENT_ID_OFFSET, + ) + .map_err(ZeboError::OperationError)?; + + let mut write_cursor = header_size as u32; + let mut entry_buf = [0u8; 16]; + + for i in 0..self.next_available_header_offset as u64 { + let (doc_id, offset, length) = match self.get_at(i)? { + Some(v) => v, + None => continue, + }; + + let index_pos = DOCUMENT_INDEX_OFFSET + i * 16; + + if Self::is_uninitialized_entry(offset) { + // Write zeroed entry + entry_buf = [0u8; 16]; + target_file + .write_all_at(&entry_buf, index_pos) + .map_err(ZeboError::OperationError)?; + } else if Self::is_deleted(doc_id, offset, length) { + // Preserve doc_id, write sentinel offset/length + entry_buf[0..8].copy_from_slice(&doc_id.to_be_bytes()); + entry_buf[8..12].copy_from_slice(&u32::MAX.to_be_bytes()); + entry_buf[12..16].copy_from_slice(&u32::MAX.to_be_bytes()); + target_file + .write_all_at(&entry_buf, index_pos) + .map_err(ZeboError::OperationError)?; + } else { + // Live document: read data from source, write to new location + if length > 0 { + let len = length as usize; + if buf.len() < len { + buf.resize(len, 0); + } + self.page_file + .read_exact_at(&mut buf[..len], offset as u64) + .map_err(ZeboError::OperationError)?; + target_file + .write_all_at(&buf[..len], write_cursor as u64) + .map_err(ZeboError::OperationError)?; + } + + entry_buf[0..8].copy_from_slice(&doc_id.to_be_bytes()); + entry_buf[8..12].copy_from_slice(&write_cursor.to_be_bytes()); + entry_buf[12..16].copy_from_slice(&length.to_be_bytes()); + target_file + .write_all_at(&entry_buf, index_pos) + .map_err(ZeboError::OperationError)?; + + write_cursor += length; + } + } + + // Zero-fill remaining slots + let zero_entry = [0u8; 16]; + for i in self.next_available_header_offset as u64..self.document_limit as u64 { + let index_pos = DOCUMENT_INDEX_OFFSET + i * 16; + target_file + .write_all_at(&zero_entry, index_pos) + .map_err(ZeboError::OperationError)?; + } + + // Write metadata + let document_count = self.get_document_count()?; + target_file + .write_all_at(&document_count.to_be_bytes(), DOCUMENT_COUNT_OFFSET) + .map_err(ZeboError::OperationError)?; + target_file + .write_all_at(&write_cursor.to_be_bytes(), NEXT_AVAILABLE_OFFSET) + .map_err(ZeboError::OperationError)?; + target_file + .write_all_at( + &self.next_available_header_offset.to_be_bytes(), + NEXT_AVAILABLE_HEADER_OFFSET, + ) + .map_err(ZeboError::OperationError)?; + + // Truncate to actual size + target_file + .set_len(write_cursor as u64) + .map_err(ZeboError::OperationError)?; + + target_file.flush().map_err(ZeboError::OperationError)?; + target_file.sync_all().map_err(ZeboError::OperationError)?; + + Ok(CompactPageStats { + bytes_before, + bytes_after: write_cursor as u64, + }) + } + pub fn close(&mut self) -> Result<()> { self.page_file.flush().map_err(ZeboError::OperationError)?; self.page_file @@ -1063,6 +1211,12 @@ impl<'docs, DocId: DocumentId, Doc: Document> ZeboPageReservedSpace<'docs, DocId } } +#[derive(Debug)] +pub struct CompactPageStats { + pub bytes_before: u64, + pub bytes_after: u64, +} + #[derive(Debug, PartialEq)] pub struct ZeboPageHeader { pub document_limit: u32,