From 85987e39fb69c65722eaa45a29a1018a454c5984 Mon Sep 17 00:00:00 2001 From: rustmailer Date: Mon, 13 Jul 2026 00:04:53 +0800 Subject: [PATCH] update --- Cargo.lock | 32 +- Cargo.toml | 12 +- crates/blob/Cargo.toml | 3 +- crates/blob/README.md | 95 ++- crates/blob/src/bucket.rs | 711 ++++-------------- crates/blob/src/engine.rs | 350 ++++++--- crates/blob/src/error.rs | 6 + crates/blob/src/gc.rs | 87 ++- crates/blob/src/recovery.rs | 88 ++- crates/blob/src/types.rs | 28 +- crates/blob/tests/fuzz_test.rs | 380 ++++++++++ crates/blob/tests/integration_test.rs | 87 ++- crates/core/src/store/blob.rs | 1 + web/src/components/mail-iframe.tsx | 32 +- .../attachment/mail-display-dialog.tsx | 35 +- .../features/attachment/mail-message-view.tsx | 4 +- .../features/search/mail-display-dialog.tsx | 34 +- web/src/features/search/mail-message-view.tsx | 4 +- 18 files changed, 1153 insertions(+), 836 deletions(-) create mode 100644 crates/blob/tests/fuzz_test.rs diff --git a/Cargo.lock b/Cargo.lock index bd26c34..5dfdff1 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -329,9 +329,10 @@ dependencies = [ "bincode_reloaded", "crc32fast", "criterion", + "fs2", "lz4_flex", - "memmap2", "rand 0.10.2", + "redb 4.1.0", "serde", "serde_json", "tempfile", @@ -388,7 +389,7 @@ dependencies = [ "html2text", "itertools 0.15.0", "itoa", - "lru 0.18.0", + "lru 0.18.1", "mail-parser", "mail-send", "memmap2", @@ -1500,6 +1501,16 @@ dependencies = [ "percent-encoding", ] +[[package]] +name = "fs2" +version = "0.4.3" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "9564fc758e15025b46aa6643b1b77d047d1a56a1aea6e01002ac0c7026876213" +dependencies = [ + "libc", + "winapi", +] + [[package]] name = "fs4" version = "0.13.1" @@ -2509,9 +2520,9 @@ dependencies = [ [[package]] name = "lru" -version = "0.18.0" +version = "0.18.1" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "8a860605968fce16869fd239cf4237a82f3ac470723415db603b0e8b6c8d4fb9" +checksum = "0b6180140927ee907000b0aa540091f6ea512ead4447c92b8fc35bc72788a5a6" dependencies = [ "hashbrown 0.17.0", ] @@ -3815,6 +3826,15 @@ dependencies = [ "libc", ] +[[package]] +name = "redb" +version = "4.1.0" +source = "registry+https://github.com/rust-lang/crates.io-index" +checksum = "8e925444704b5f17d32bf42f5b6e2df050bceebc3dcd6e71cc73dafe8092e839" +dependencies = [ + "libc", +] + [[package]] name = "redox_syscall" version = "0.5.18" @@ -3826,9 +3846,9 @@ dependencies = [ [[package]] name = "regex" -version = "1.12.4" +version = "1.13.0" source = "registry+https://github.com/rust-lang/crates.io-index" -checksum = "f1292b7759ae1cb9ec195452d1390a074f0cd8541ab7a5a8c31cd6db45d4a6ba" +checksum = "2a0e75113e14dc5acb068cd0786884f214f1312650a3d36d269f5c4f3cdee8a2" dependencies = [ "aho-corasick", "memchr", diff --git a/Cargo.toml b/Cargo.toml index c8f94c1..ab0916b 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -41,11 +41,11 @@ reqwest = { version = "0.12.24", default-features = false, features = [ ] } tokio-socks = "0.5.3" http = "1.4.2" -regex = "1.12.4" +regex = "1.13" email_address = "0.2.9" futures = "0.3.32" utf7-imap = "0.3.2" -mail-parser = { version = '0.11.4', features = ["serde"] } +mail-parser = { version = '0.11', features = ["serde"] } # mail-send = "0.5.2" tokio-rustls = { version = "0.26.4", default-features = false, features = [ "ring", @@ -54,7 +54,7 @@ tokio-rustls = { version = "0.26.4", default-features = false, features = [ timeago = "0.6.1" oauth2 = { version = "5.0.0", features = ["reqwest-blocking"] } url = { version = "2.5.8", features = ["serde"] } -sysinfo = "0.39.4" +sysinfo = "0.39" num_cpus = "1.17.0" rand = "0.10.2" encoding_rs = "0.8.35" @@ -64,7 +64,7 @@ rustls-pki-types = "1.15.0" tokio-io-timeout = "1.2.1" semver = "1.0.28" governor = "0.10.4" -lru = "0.18.0" +lru = "0.18.1" mime_guess = "2.0.5" hex = "0.4.3" time = { version = "0.3.53", features = [ @@ -72,14 +72,14 @@ time = { version = "0.3.53", features = [ "parsing", "local-offset", ] } -rust-embed = "8.12.0" +rust-embed = "8.12" murmur3 = "0.5.2" urlencoding = "2.1.3" dashmap = "6.2.1" gethostname = "1.1.0" itoa = "1.0.18" html2text = "0.17.1" -bytes = "1.12.0" +bytes = "1.12" dialoguer = "0.12.0" console = "0.16.4" mail-send = "0.6.1" diff --git a/crates/blob/Cargo.toml b/crates/blob/Cargo.toml index 75f4ae2..6623483 100644 --- a/crates/blob/Cargo.toml +++ b/crates/blob/Cargo.toml @@ -13,7 +13,8 @@ serde_json = "1" bincode = { package = "bincode_reloaded", version = "3.1.10", default-features = false, features = ["serde", "alloc"] } tracing = "0.1" thiserror = "2" -memmap2 = "0.9" +redb = "4.1" +fs2 = "0.4" [dev-dependencies] tempfile = "3" diff --git a/crates/blob/README.md b/crates/blob/README.md index a675f0d..89687b0 100644 --- a/crates/blob/README.md +++ b/crates/blob/README.md @@ -9,14 +9,11 @@ All data is keyed by a 32-byte content hash — identical content is stored only ``` / ├── meta.bin # global metadata (bincode + CRC32) +├── index.redb # key → (segment_id, offset, size) index (redb B-tree) ├── segments/ -│ ├── 00000001.seg # append-only segment files (≤ 256 MB each) +│ ├── 00000001.seg # append-only segment files (≤ 1 GB each) │ ├── 00000002.seg │ └── ... -└── buckets/ - ├── 00.idx # 256 bucket index files (mmap'd, binary-searchable) - ├── 01.idx - └── ... ``` ### Segment entry format (50-byte fixed header + variable data) @@ -33,27 +30,28 @@ magic(4) crc32(4) flags(1) codec(1) key(32) raw_size(4) data_size(4) data - `raw_size`: original uncompressed size - `data_size`: on-disk data size (after compression) -### Bucket index record format (56 bytes) +### Index store -``` -key(32) segment_id(4) offset(8) data_size(4) flags(1) _pad(3) crc32(4) -``` +The key → (segment_id, offset, data_size, flags) mapping is stored in a single [redb](https://github.com/cberner/redb) database (`index.redb`). redb provides: -Every index record is CRC32-protected — a corrupted record is detected on read and never silently returned. Records are sorted by key, deduplicated (newest segment_id + offset wins), and mmap'd for zero-heap binary search. Pending writes (since the last compaction) live in a small in-memory HashMap. +- **B-tree + mmap**: O(log N) point lookups with zero heap allocation — pages are faulted in on demand. +- **ACID transactions**: every index write is durable and atomic. +- **Crash recovery**: handled transparently by redb's WAL — no manual reload or rebuild logic. +- **O(1) startup**: only the B-tree root page is read at open time. + +Records are stored as fixed-size 56-byte blobs, each carrying an internal CRC32 checksum. ## Read path ``` get(key) - → bucket_id = (key[0..2] as u16) % 256 - → check pending HashMap (most recent wins) - → binary search mmap'd bucket file - → IndexRecord CRC32 verify → (segment_id, offset, data_size) + → index_store.get(key) # redb B-tree lookup, zero-copy + → IndexRecord CRC32 verify # (segment_id, offset, data_size, flags) → pread entry from segment file → entry CRC32 verify → decompress → return value ``` -Both the bucket index record and the segment entry carry independent CRC32 checksums. Corruption in one record or entry is contained — it never affects other keys. +The index record and the segment entry carry independent CRC32 checksums. Corruption in one record or entry is contained — it never affects other keys. ## Write path @@ -61,20 +59,53 @@ Both the bucket index record and the segment entry carry independent CRC32 check put(key, value, codec) → compress value (Zstd/Lz4 if ≥ 4 KB, else store raw) → append entry to active segment file - → insert IndexRecord into bucket store (append to .idx file + HashMap) + → insert IndexRecord into redb (single write txn) → update metadata (indexed_up_to_offset) - → if segment ≥ 256 MB → seal it, create new segment + → if segment ≥ 1 GB → seal it, create new segment ``` ## Delete -Deletes are **tombstones** — an entry with `flags=1` and empty data is appended. The bucket index maps the key to this tombstone. Read returns `None`. GC later reclaims the space. +Deletes are **tombstones** — an entry with `flags=1` and empty data is appended to the active segment. Before writing the tombstone, the existing index record is consulted to increment `deleted_bytes` on the **original** segment (the one that holds the live data). This drives the GC threshold. -## Compaction & GC +``` +delete(key) + → index_store.get(key) → find original (segment_id, data_size) + → original_segment.deleted_bytes += data_size + → recompute deleted_ratio on original segment + → append tombstone entry to active segment + → insert tombstone IndexRecord into redb + → index_store.get(key) now returns None +``` -**Bucket compaction:** when a bucket's pending HashMap grows past `compact_threshold` (default 10,000), the mmap + pending records are merged, sorted, deduplicated, and atomically rewritten. This keeps binary search fast. +## GC -**Garbage collection:** when a sealed segment's deleted-ratio exceeds `gc_deleted_ratio` (default 0.30), GC scans all segments to find the latest entry per key, then rewrites the target segment keeping only live entries. Tombstones and overwritten entries are dropped. Bucket indices are rebuilt afterward. +**Segment GC:** two-phase, driven by the `deleted_ratio` tracked per segment. + +### Trigger + +- **Background**: the `blob-gc` thread wakes up every `gc_interval_secs` (default 300s), checks whether any sealed segment's `deleted_ratio ≥ gc_deleted_ratio` (default 0.30), and runs GC on the worst segment if so. +- **Manual**: `engine.gc_if_needed()` or `engine.gc()`. + +### Phase 1 — scan & compact (read-only, no write lock) + +1. Pick the sealed segment with the highest `deleted_ratio`. +2. Scan only that segment's entries. +3. For each entry, ask the index: "is this entry still the latest version for its key?" + - If the index points to this exact `(segment_id, offset)` → **keep**, write to a temp segment file. + - If the index points elsewhere (overwritten by a later segment, or a tombstone) → **skip** (stale). +4. Fsync the temp file. + +Phase 1 holds only the read lock — `put` / `delete` continue uninterrupted. + +### Phase 2 — commit & update index (write lock) + +1. Atomically rename the temp file over the original segment. +2. Batch-insert new `IndexRecord`s (now at new offsets) into redb. Old records with the same key are naturally overwritten. +3. Reset the segment's `deleted_bytes` and `deleted_ratio` to zero. +4. Persist metadata. + +Phase 2 holds the write lock, but is fast — no full segment scan, no full index rebuild. ## Data integrity @@ -83,22 +114,28 @@ Every record on disk is independently checksummed: | Layer | Format | Protection | |---|---|---| | Segment entry | 50-byte header + data | CRC32 covers all fields + data | -| Bucket index record | 56 bytes | CRC32 covers key + segment_id + offset + data_size + flags | +| Index record | 56 bytes | CRC32 covers key + segment_id + offset + data_size + flags | | Global metadata | bincode blob | CRC32 + version header | -Corruption is **contained** — a bad segment entry or bucket record produces an error for that key only. Compaction and recovery skip corrupt records (with a warning) rather than aborting. Bucket indices can always be fully rebuilt from segments via `rebuild_from_segments`. +Corruption is **contained** — a bad segment entry or index record produces an error for that key only. Recovery and GC skip corrupt records (with a warning) rather than aborting. The index is backed by redb's B-tree which maintains its own internal integrity. ## Crash recovery - Temp files from interrupted GC are cleaned up on open. -- Any segment data beyond `indexed_up_to_offset` is scanned and indexed into bucket files. -- After appending recovered records, bucket mmaps are reloaded so they are immediately visible. +- Any segment data beyond `indexed_up_to_offset` is scanned and inserted into the index. - Partial writes at the tail of a segment (detected via CRC32 mismatch near EOF) are truncated. -- Buckets are always repairable by re-scanning segments (`rebuild_from_segments`). +- redb's WAL ensures the index is always consistent — no manual reload or rebuild needed. -## Background flush +## Background threads -Set `Config.flush_interval_secs` to a positive value to have a background thread fsync the active segment and save metadata periodically. This bounds recovery time after a crash at the cost of a small I/O overhead. +Set `Config.flush_interval_secs` and `Config.gc_interval_secs` to positive values to enable periodic background work: + +| Thread | Config | Default | What it does | +|---|---|---|---| +| `blob-flush` | `flush_interval_secs` | `0` (off) | Fsync the active segment and save metadata | +| `blob-gc` | `gc_interval_secs` | `0` (off) | Check deleted-ratio, compact one segment if needed | + +The two threads are independent — a long GC run never blocks fsync. ## Config @@ -107,9 +144,9 @@ Set `Config.flush_interval_secs` to a positive value to have a background thread | `compress_threshold` | 4096 | Bytes; smaller values stored uncompressed | | `default_codec` | Zstd | Also supports Lz4 | | `compression_level` | 0 | Zstd compression level | -| `compact_threshold` | 10000 | Pending records per bucket before auto-compact | | `gc_deleted_ratio` | 0.30 | Trigger GC when a sealed segment exceeds this | | `flush_interval_secs` | 0 | 0 = disabled; ≥ 5 for periodic background fsync | +| `gc_interval_secs` | 0 | 0 = disabled; ≥ 10 for periodic background GC | ## Basic usage diff --git a/crates/blob/src/bucket.rs b/crates/blob/src/bucket.rs index 53e6801..d7432a2 100644 --- a/crates/blob/src/bucket.rs +++ b/crates/blob/src/bucket.rs @@ -1,15 +1,14 @@ -use std::collections::HashMap; -use std::fs::{self, File, OpenOptions}; -use std::io::Write; -use std::path::{Path, PathBuf}; -use std::sync::Mutex; +use std::path::Path; -use memmap2::Mmap; +use redb::{Database, ReadableTable, TableDefinition}; +use redb::ReadableDatabase; use crate::error::Result; -use crate::types::{BUCKET_COUNT, INDEX_RECORD_SIZE}; +use crate::types::INDEX_RECORD_SIZE; -/// On-disk format: 52 bytes per record. +// ── IndexRecord ────────────────────────────────────────────────────────────── + +/// On-disk format: 52 bytes per record + 4 bytes CRC = 56 bytes total. #[derive(Debug, Clone, PartialEq, Eq)] pub struct IndexRecord { pub key: [u8; 32], @@ -74,106 +73,86 @@ impl IndexRecord { } } -/// Compute bucket_id from a key's first 2 bytes. -pub fn bucket_id(key: &[u8; 32]) -> u16 { - u16::from_be_bytes([key[0], key[1]]) % BUCKET_COUNT -} +// ── redb Value impl for fixed-size record bytes ────────────────────────────── -// ── BucketStore ──────────────────────────────────────────────────────────── +/// Newtype wrapper so we can implement `redb::Value` for `[u8; INDEX_RECORD_SIZE]`. +#[derive(Debug, Clone, Copy)] +struct RecordBytes([u8; INDEX_RECORD_SIZE]); -/// Per-bucket mutable state behind a Mutex. -struct BucketState { - /// mmap of the clean, sorted, deduplicated portion of the bucket file. - mmap: Mmap, - /// Number of sorted records in the mmap. - compacted_records: usize, - /// Recent writes not yet merged into the mmap. Key → latest record. - pending: HashMap<[u8; 32], IndexRecord>, - /// Append-only file for durability of pending writes. - file: File, - /// Path to the bucket file. - path: PathBuf, -} +impl redb::Value for RecordBytes { + type SelfType<'a> = RecordBytes; + type AsBytes<'a> = [u8; INDEX_RECORD_SIZE]; -/// Zero-heap bucket index store backed by mmap. -/// -/// Each bucket's clean portion is mmap'd — binary search reads directly -/// from the OS page cache without allocating heap memory proportional to -/// the number of stored keys. Only pending writes (since the last compact) -/// live in a small in-memory HashMap. -pub struct BucketStore { - states: Vec>, - compact_threshold: usize, -} - -impl BucketStore { - /// Open all bucket files. On first open or after a crash, each file is - /// loaded, sorted, deduplicated, and rewritten into a clean mmap'd form. - pub fn open(dir: &Path, compact_threshold: usize) -> Result { - fs::create_dir_all(dir)?; - - let mut states = Vec::with_capacity(BUCKET_COUNT as usize); - for bid in 0..BUCKET_COUNT { - let path = bucket_path(dir, bid); - let (mmap, compacted_records) = if path.exists() { - // Load, sort, dedup, rewrite clean, then mmap. - let records = load_records_from_file(&path)?; - let deduped = sort_and_dedup(records); - let count = deduped.len(); - rewrite_file(&path, &deduped)?; - let file = fs::File::open(&path)?; - let mmap = unsafe { Mmap::map(&file)? }; - (mmap, count) - } else { - // Create empty bucket file. - let file = fs::File::create(&path)?; - file.set_len(0)?; - let mmap = unsafe { Mmap::map(&file)? }; - (mmap, 0) - }; - - let file = OpenOptions::new() - .create(true) - .append(true) - .open(&path)?; - - states.push(Mutex::new(BucketState { - mmap, - compacted_records, - pending: HashMap::new(), - file, - path, - })); - } - - Ok(Self { - states, - compact_threshold, - }) + fn fixed_width() -> Option { + Some(INDEX_RECORD_SIZE) } - /// Look up a key. Returns the latest IndexRecord, or None if absent/tombstone. + fn from_bytes<'a>(data: &'a [u8]) -> Self::SelfType<'a> + where + Self: 'a + { + let mut arr = [0u8; INDEX_RECORD_SIZE]; + arr.copy_from_slice(data); + RecordBytes(arr) + } + + fn as_bytes<'a, 'b: 'a>(value: &'a Self::SelfType<'b>) -> Self::AsBytes<'a> { + value.0 + } + + fn type_name() -> redb::TypeName { + redb::TypeName::new("IndexRecord") + } +} + +// ── IndexStore ─────────────────────────────────────────────────────────────── + +const INDEX_TABLE: TableDefinition<[u8; 32], RecordBytes> = TableDefinition::new("blob_index"); + +/// Zero-heap key-value index backed by redb. +/// +/// B-tree + mmap + ACID. No manual compact / sort / dedup / per-bucket mmap +/// management. Startup is O(1) — redb reads only its root page. +pub struct IndexStore { + db: Database, +} + +impl IndexStore { + /// Open (or create) the index database at `dir/index.redb`. + pub fn open(dir: &Path) -> Result { + let path = dir.join("index.redb"); + let db = Database::create(&path) + .map_err(|e| crate::error::Error::IndexDb(format!("failed to create index: {}", e)))?; + // Ensure the table exists so reads on a fresh database don't fail. + { + let txn = db + .begin_write() + .map_err(|e| crate::error::Error::IndexDb(format!("init write txn: {}", e)))?; + txn.open_table(INDEX_TABLE) + .map_err(|e| crate::error::Error::IndexDb(format!("init table: {}", e)))?; + txn.commit() + .map_err(|e| crate::error::Error::IndexDb(format!("init commit: {}", e)))?; + } + Ok(Self { db }) + } + + /// Look up a key. Returns the latest IndexRecord, or None if absent/tombstone. pub fn get(&self, key: &[u8; 32]) -> Result> { - let bid = bucket_id(key) as usize; - let state = self.states[bid].lock().unwrap(); + let txn = self + .db + .begin_read() + .map_err(|e| crate::error::Error::IndexDb(format!("read txn: {}", e)))?; + let table = txn + .open_table(INDEX_TABLE) + .map_err(|e| crate::error::Error::IndexDb(format!("open table: {}", e)))?; - // 1. Check pending (most recent wins) - if let Some(rec) = state.pending.get(key) { - return Ok(if rec.is_tombstone() { None } else { Some(rec.clone()) }); - } - - // 2. Binary search the mmap'd sorted portion - let bytes: &[u8] = &state.mmap; - if bytes.is_empty() { - return Ok(None); - } - - let count = bytes.len() / INDEX_RECORD_SIZE; - let result = binary_search_records(bytes, key, count); - match result { - Some(idx) => { - let rec = read_record_at(bytes, idx)?; - Ok(if rec.is_tombstone() { None } else { Some(rec) }) + match table + .get(key) + .map_err(|e| crate::error::Error::IndexDb(format!("get: {}", e)))? + { + Some(guard) => { + let record = IndexRecord::decode(&guard.value().0)?; + Ok(if record.is_tombstone() { None } else { Some(record) }) } None => Ok(None), } @@ -184,467 +163,95 @@ impl BucketStore { self.get(key).map(|r| r.is_some()) } - /// Insert or update a record for a key. Appends to the file for durability, - /// then inserts into the pending HashMap. - pub fn insert(&self, record: IndexRecord) -> Result<()> { - let bid = bucket_id(&record.key) as usize; - let mut state = self.states[bid].lock().unwrap(); - - // Durability: append to file - state.file.write_all(&record.encode())?; - - // Update pending - state.pending.insert(record.key, record); - - // Auto-compact if pending grows too large - if state.pending.len() >= self.compact_threshold { - drop(state); - self.compact_bucket(bid as u16)?; + /// Insert or update a record for a key. Committed in a single write txn. + pub fn insert(&self, record: &IndexRecord) -> Result<()> { + let txn = self + .db + .begin_write() + .map_err(|e| crate::error::Error::IndexDb(format!("write txn: {}", e)))?; + { + let mut table = txn + .open_table(INDEX_TABLE) + .map_err(|e| crate::error::Error::IndexDb(format!("open table: {}", e)))?; + table + .insert(&record.key, RecordBytes(record.encode())) + .map_err(|e| crate::error::Error::IndexDb(format!("insert: {}", e)))?; } - + txn.commit() + .map_err(|e| crate::error::Error::IndexDb(format!("commit: {}", e)))?; Ok(()) } - /// Batch insert multiple records. Appends all, then updates pending. + /// Batch insert multiple records in a single write transaction. pub fn insert_batch(&self, records: &[IndexRecord]) -> Result<()> { - // Group by bucket - let mut grouped: HashMap> = HashMap::new(); - for r in records { - let bid = bucket_id(&r.key); - grouped.entry(bid).or_default().push(r); - } - - for (bid, recs) in &grouped { - let bid_usize = *bid as usize; - let mut state = self.states[bid_usize].lock().unwrap(); - - for r in recs { - state.file.write_all(&r.encode())?; - state.pending.insert(r.key, (*r).clone()); - } - } - - // Compact overfull buckets - for &bid in grouped.keys() { - let state = self.states[bid as usize].lock().unwrap(); - let needs_compact = state.pending.len() >= self.compact_threshold; - drop(state); - if needs_compact { - self.compact_bucket(bid)?; - } - } - - Ok(()) - } - - /// Run stats across all buckets: total keys (non-tombstone) and total data bytes. - pub fn total_keys(&self) -> usize { - let mut count = 0usize; - for state in self.states.iter() { - let s = state.lock().unwrap(); - // Count from pending - for r in s.pending.values() { - if !r.is_tombstone() { - count += 1; - } - } - // Count from mmap - let bytes: &[u8] = &s.mmap; - let n = bytes.len() / INDEX_RECORD_SIZE; - for i in 0..n { - if let Ok(rec) = read_record_at(bytes, i) { - // Skip keys that are overridden in pending - if s.pending.contains_key(&rec.key) { - continue; - } - if !rec.is_tombstone() { - count += 1; - } - } - } - } - count - } - - /// Compact a single bucket: merge mmap + pending, sort+dedup, rewrite file, remap. - fn compact_bucket(&self, bid: u16) -> Result<()> { - let idx = bid as usize; - let mut state = self.states[idx].lock().unwrap(); - - if state.pending.is_empty() { + if records.is_empty() { return Ok(()); } - - // Collect mmap records + pending records - let mut all: Vec = Vec::new(); - - let bytes: &[u8] = &state.mmap; - let n = bytes.len() / INDEX_RECORD_SIZE; - all.reserve(n + state.pending.len()); - for i in 0..n { - match read_record_at(bytes, i) { - Ok(rec) => all.push(rec), - Err(_) => { - tracing::warn!( - "Bucket {:02x} record {} CRC mismatch during compact, dropping", - bid, i - ); - } - } - } - for r in state.pending.values() { - all.push(r.clone()); - } - - let deduped = sort_and_dedup(all); - let new_count = deduped.len(); - - // Write new file atomically - rewrite_file(&state.path, &deduped)?; - - // Remap - let file = fs::File::open(&state.path)?; - let new_mmap = unsafe { Mmap::map(&file)? }; - - // Re-open append file descriptor (old one was truncated) - let new_file = OpenOptions::new() - .create(true) - .append(true) - .open(&state.path)?; - - state.mmap = new_mmap; - state.compacted_records = new_count; - state.pending.clear(); - state.file = new_file; - - Ok(()) - } - - /// Compact all buckets. - pub fn compact_all(&self) -> Result<()> { - for bid in 0..BUCKET_COUNT { - self.compact_bucket(bid)?; - } - Ok(()) - } - - /// Reload all mmaps from disk and reset pending state. - /// Used after GC rewrites bucket files externally (via rebuild_from_segments). - pub fn reload_all(&self) -> Result<()> { - for bid in 0..BUCKET_COUNT { - let mut state = self.states[bid as usize].lock().unwrap(); - let path = state.path.clone(); - - // Load, sort, dedup, rewrite clean - let records = load_records_from_file(&path)?; - let deduped = sort_and_dedup(records); - let count = deduped.len(); - rewrite_file(&path, &deduped)?; - - // Remap - let file = fs::File::open(&path)?; - let new_mmap = unsafe { Mmap::map(&file)? }; - - // Reopen append file - let new_file = OpenOptions::new() - .create(true) - .append(true) - .open(&path)?; - - state.mmap = new_mmap; - state.compacted_records = count; - state.pending.clear(); - state.file = new_file; - } - Ok(()) - } - - /// Rebuild all bucket files from scratch by scanning segment entries. - /// Used by GC and recovery. - pub fn rebuild_from_segments( - dir: &Path, - segments: &[(u32, &Path)], - ) -> Result<()> { - use crate::segment::SegmentReader; - - let mut bucket_records: HashMap> = HashMap::new(); - for i in 0..BUCKET_COUNT { - bucket_records.insert(i, Vec::new()); - } - - for &(seg_id, seg_path) in segments { - if !seg_path.exists() { - continue; - } - let reader = SegmentReader::open(seg_path.to_path_buf(), seg_id)?; - reader.scan_entries(0, |entry, offset| { - let bid = bucket_id(&entry.key); - let rec = IndexRecord::new( - entry.key, - seg_id, - offset, - entry.data.len() as u32, - entry.flags, - ); - bucket_records.entry(bid).or_default().push(rec); - Ok(()) - })?; - } - - for (bid, records) in &bucket_records { - let deduped = sort_and_dedup(records.clone()); - let path = bucket_path(dir, *bid); - rewrite_file(&path, &deduped)?; - } - - Ok(()) - } -} - -// ── Helpers ──────────────────────────────────────────────────────────────── - -fn bucket_path(dir: &Path, bid: u16) -> PathBuf { - dir.join(format!("{:02x}.idx", bid)) -} - -/// Read records from a raw bucket file (may contain duplicates, not sorted). -fn load_records_from_file(path: &Path) -> Result> { - let data = fs::read(path)?; - let remainder = data.len() % INDEX_RECORD_SIZE; - let count = data.len() / INDEX_RECORD_SIZE; - let mut records = Vec::with_capacity(count); - for i in 0..count { - let start = i * INDEX_RECORD_SIZE; - let end = start + INDEX_RECORD_SIZE; - let buf: &[u8; INDEX_RECORD_SIZE] = data[start..end].try_into().map_err(|_| { - crate::error::Error::BucketIndexCorrupt { - path: path.to_path_buf(), - reason: "unexpected file size".into(), - } - })?; - match IndexRecord::decode(buf) { - Ok(rec) => records.push(rec), - Err(_) => { - tracing::warn!( - "Bucket file {:?} record {} CRC mismatch, skipping", - path, i - ); - } - } - } - if remainder > 0 { - tracing::warn!( - "Bucket file {:?} has {} trailing bytes, ignoring", - path, - remainder - ); - } - Ok(records) -} - -/// Sort records by key, deduplicate keeping the one with the highest (segment_id, offset). -fn sort_and_dedup(mut records: Vec) -> Vec { - records.sort_by_key(|a| a.key); - let mut out = Vec::with_capacity(records.len()); - let mut i = 0; - while i < records.len() { - let mut best = i; - let mut j = i + 1; - while j < records.len() && records[j].key == records[i].key { - if records[j].segment_id > records[best].segment_id - || (records[j].segment_id == records[best].segment_id - && records[j].offset > records[best].offset) - { - best = j; - } - j += 1; - } - out.push(records[best].clone()); - i = j; - } - out -} - -/// Binary search for `key` in `bytes` (array of INDEX_RECORD_SIZE records, sorted). -fn binary_search_records(bytes: &[u8], key: &[u8; 32], count: usize) -> Option { - let mut lo = 0usize; - let mut hi = count; - while lo < hi { - let mid = lo + (hi - lo) / 2; - let rec_key = read_key_at(bytes, mid); - match rec_key.cmp(key) { - std::cmp::Ordering::Less => lo = mid + 1, - std::cmp::Ordering::Greater => hi = mid, - std::cmp::Ordering::Equal => return Some(mid), - } - } - None -} - -/// Read the key at index `idx` from a byte slice of INDEX_RECORD_SIZE records. -fn read_key_at(bytes: &[u8], idx: usize) -> &[u8; 32] { - let start = idx * INDEX_RECORD_SIZE; - bytes[start..start + 32].try_into().unwrap() -} - -/// Read a full IndexRecord at index `idx`. -fn read_record_at(bytes: &[u8], idx: usize) -> crate::error::Result { - let start = idx * INDEX_RECORD_SIZE; - let buf: &[u8; INDEX_RECORD_SIZE] = bytes[start..start + INDEX_RECORD_SIZE] - .try_into() - .unwrap(); - IndexRecord::decode(buf) -} - -/// Atomically rewrite a bucket file with sorted, deduplicated records. -fn rewrite_file(path: &Path, records: &[IndexRecord]) -> Result<()> { - let mut buf = Vec::with_capacity(records.len() * INDEX_RECORD_SIZE); - for r in records { - buf.extend_from_slice(&r.encode()); - } - crate::fs::create_atomic(path, &buf) -} - -#[cfg(test)] -mod tests { - use super::*; - use tempfile::TempDir; - - #[test] - fn test_index_record_encode_decode() { - let mut key = [0u8; 32]; - key[0..4].copy_from_slice(&[1, 2, 3, 4]); - let rec = IndexRecord::new(key, 5, 12345, 500, 0); - let encoded = rec.encode(); - assert_eq!(encoded.len(), INDEX_RECORD_SIZE); - let decoded = IndexRecord::decode(&encoded).unwrap(); - assert_eq!(rec, decoded); - } - - #[test] - fn test_index_record_crc_detects_corruption() { - let mut key = [0u8; 32]; - key[0..4].copy_from_slice(&[1, 2, 3, 4]); - let rec = IndexRecord::new(key, 5, 12345, 500, 0); - let mut encoded = rec.encode(); - // Flip a bit in the data portion - encoded[40] ^= 1; - assert!(IndexRecord::decode(&encoded).is_err()); - } - - #[test] - fn test_bucket_id_deterministic() { - let mut key = [0u8; 32]; - key[0] = 0x00; - key[1] = 0x0F; - assert_eq!(bucket_id(&key), 15); - key[0] = 0x00; - key[1] = 0x10; - assert_eq!(bucket_id(&key), 16); - } - - #[test] - fn test_sort_and_dedup_keeps_latest() { - let recs = vec![ - IndexRecord::new([1u8; 32], 1, 100, 50, 0), - IndexRecord::new([1u8; 32], 2, 200, 50, 0), - IndexRecord::new([2u8; 32], 1, 300, 60, 0), - ]; - let deduped = sort_and_dedup(recs); - assert_eq!(deduped.len(), 2); - let found = &deduped[0]; - assert_eq!(found.segment_id, 2); - assert_eq!(found.offset, 200); - } - - #[test] - fn test_bucket_store_put_and_get() { - let dir = TempDir::new().unwrap(); - let store = BucketStore::open(dir.path(), 100).unwrap(); - - let key = [0xAA; 32]; - let rec = IndexRecord::new(key, 1, 0, 500, 0); - store.insert(rec).unwrap(); - - let found = store.get(&key).unwrap(); - assert!(found.is_some()); - assert_eq!(found.unwrap().data_size, 500); - } - - #[test] - fn test_bucket_store_get_missing() { - let dir = TempDir::new().unwrap(); - let store = BucketStore::open(dir.path(), 100).unwrap(); - - let result = store.get(&[0xFF; 32]).unwrap(); - assert!(result.is_none()); - } - - #[test] - fn test_bucket_store_tombstone() { - let dir = TempDir::new().unwrap(); - let store = BucketStore::open(dir.path(), 100).unwrap(); - - let key = [0xBB; 32]; - - // Write then tombstone - store - .insert(IndexRecord::new(key, 1, 0, 100, 0)) - .unwrap(); - store - .insert(IndexRecord::new(key, 2, 0, 0, 1)) - .unwrap(); - - let result = store.get(&key).unwrap(); - assert!(result.is_none()); - } - - #[test] - fn test_bucket_store_compact() { - let dir = TempDir::new().unwrap(); - // Use small threshold to trigger auto-compact - let store = BucketStore::open(dir.path(), 5).unwrap(); - - // Write 10 records for same bucket (all same first 2 bytes → same bucket) - for i in 0u8..10 { - let mut key = [0u8; 32]; - key[0..2].copy_from_slice(&[0x00, 0x00]); // same bucket - key[2] = i; - store - .insert(IndexRecord::new(key, 1, i as u64 * 100, 50, 0)) - .unwrap(); - } - - // All should be readable after auto-compact - for i in 0u8..10 { - let mut key = [0u8; 32]; - key[0..2].copy_from_slice(&[0x00, 0x00]); - key[2] = i; - let found = store.get(&key).unwrap(); - assert!(found.is_some(), "key {} should exist after compact", i); - } - } - - #[test] - fn test_bucket_store_persistence() { - let dir = TempDir::new().unwrap(); - let dir_path = dir.path().to_path_buf(); - - let key = [0xCC; 32]; + let txn = self + .db + .begin_write() + .map_err(|e| crate::error::Error::IndexDb(format!("write txn: {}", e)))?; { - let store = BucketStore::open(&dir_path, 100).unwrap(); - store - .insert(IndexRecord::new(key, 1, 42, 512, 0)) - .unwrap(); - store.compact_all().unwrap(); + let mut table = txn + .open_table(INDEX_TABLE) + .map_err(|e| crate::error::Error::IndexDb(format!("open table: {}", e)))?; + for record in records { + table + .insert(&record.key, RecordBytes(record.encode())) + .map_err(|e| crate::error::Error::IndexDb(format!("insert: {}", e)))?; + } } + txn.commit() + .map_err(|e| crate::error::Error::IndexDb(format!("commit: {}", e)))?; + Ok(()) + } - // Reopen - { - let store = BucketStore::open(&dir_path, 100).unwrap(); - let found = store.get(&key).unwrap(); - assert!(found.is_some()); - assert_eq!(found.unwrap().offset, 42); + /// Remove keys from the index in a single write transaction. + pub fn delete_batch(&self, keys: &[[u8; 32]]) -> Result<()> { + if keys.is_empty() { + return Ok(()); } + let txn = self + .db + .begin_write() + .map_err(|e| crate::error::Error::IndexDb(format!("write txn: {}", e)))?; + { + let mut table = txn + .open_table(INDEX_TABLE) + .map_err(|e| crate::error::Error::IndexDb(format!("open table: {}", e)))?; + for key in keys { + table + .remove(key) + .map_err(|e| crate::error::Error::IndexDb(format!("remove: {}", e)))?; + } + } + txn.commit() + .map_err(|e| crate::error::Error::IndexDb(format!("commit: {}", e)))?; + Ok(()) + } + + /// Total number of live (non-tombstone) keys. + pub fn total_keys(&self) -> Result { + let txn = self + .db + .begin_read() + .map_err(|e| crate::error::Error::IndexDb(format!("read txn: {}", e)))?; + let table = txn + .open_table(INDEX_TABLE) + .map_err(|e| crate::error::Error::IndexDb(format!("open table: {}", e)))?; + + let mut count = 0usize; + let iter = table + .iter() + .map_err(|e| crate::error::Error::IndexDb(format!("iter: {}", e)))?; + for item in iter { + let (_, guard) = + item.map_err(|e| crate::error::Error::IndexDb(format!("iter next: {}", e)))?; + let record = IndexRecord::decode(&guard.value().0)?; + if !record.is_tombstone() { + count += 1; + } + } + Ok(count) } } diff --git a/crates/blob/src/engine.rs b/crates/blob/src/engine.rs index cd9012d..609c827 100644 --- a/crates/blob/src/engine.rs +++ b/crates/blob/src/engine.rs @@ -1,12 +1,14 @@ -use std::collections::HashMap; use std::fs; +use std::fs::File; use std::path::{Path, PathBuf}; use std::sync::atomic::{AtomicBool, Ordering}; use std::sync::{Arc, Mutex, RwLock}; use std::thread::{self, JoinHandle}; use std::time::Duration; -use crate::bucket::{BucketStore, IndexRecord}; +use fs2::FileExt; + +use crate::bucket::{IndexRecord, IndexStore}; use crate::compress; use crate::error::{Error, Result}; use crate::file_pool::FilePool; @@ -23,14 +25,17 @@ use crate::types::{Codec, Config, ENTRY_HEADER_SIZE}; pub struct Engine { shared: Arc, flush_handle: Mutex>, + gc_handle: Mutex>, } struct EngineShared { config: Config, inner: RwLock, - bucket_store: BucketStore, + index_store: IndexStore, write_mutex: Mutex<()>, file_pool: FilePool, + #[allow(dead_code)] + lock_file: File, } struct FlushHandle { @@ -38,11 +43,15 @@ struct FlushHandle { stop: Arc, } +struct GcHandle { + handle: JoinHandle<()>, + stop: Arc, +} + struct EngineInner { root: PathBuf, meta: GlobalMeta, active_writer: SegmentWriter, - readers: HashMap, } #[derive(Debug, Clone)] @@ -60,8 +69,23 @@ impl Engine { fs::create_dir_all(path)?; fs::create_dir_all(path.join("segments"))?; - let bucket_dir = path.join("buckets"); - let bucket_store = BucketStore::open(&bucket_dir, config.compact_threshold)?; + // Acquire an exclusive file lock so that no two processes can open + // the same database directory concurrently. + let lock_path = path.join("LOCK"); + let lock_file = File::create(&lock_path)?; + lock_file.try_lock_exclusive().map_err(|_| Error::AlreadyOpen { + path: path.display().to_string(), + })?; + + let (index_store, index_rebuilt) = match IndexStore::open(path) { + Ok(s) => (s, false), + Err(e) => { + tracing::warn!("Index database unreadable, rebuilding from segments: {}", e); + let index_path = path.join("index.redb"); + let _ = std::fs::remove_file(&index_path); + (IndexStore::open(path)?, true) + } + }; let mut meta = GlobalMeta::load(path)?; @@ -99,9 +123,12 @@ impl Engine { } crate::recovery::cleanup_temp_files(path)?; - crate::recovery::recover(path, &mut meta)?; - // Recovery may have appended new records to bucket files — reload mmaps. - bucket_store.reload_all()?; + let recovered_records = if index_rebuilt { + crate::recovery::rebuild_index(path, &mut meta)? + } else { + crate::recovery::recover(path, &mut meta)? + }; + index_store.insert_batch(&recovered_records)?; let seg_path = path .join("segments") @@ -112,18 +139,6 @@ impl Engine { SegmentWriter::create(seg_path, meta.active_segment_id)? }; - let mut readers = HashMap::new(); - for (&seg_id, stats) in &meta.segments { - if stats.sealed { - let seg_path = path - .join("segments") - .join(segment::segment_filename(seg_id)); - if seg_path.exists() { - readers.insert(seg_id, SegmentReader::open(seg_path, seg_id)?); - } - } - } - meta.save(path)?; let shared = Arc::new(EngineShared { @@ -132,11 +147,11 @@ impl Engine { root: path.to_path_buf(), meta, active_writer, - readers, }), - bucket_store, + index_store, write_mutex: Mutex::new(()), file_pool: FilePool::new(8), + lock_file, }); let flush_handle = if config.flush_interval_secs > 0 { @@ -167,9 +182,59 @@ impl Engine { None }; + let gc_handle = if config.gc_interval_secs > 0 { + let shared3 = Arc::clone(&shared); + let stop = Arc::new(AtomicBool::new(false)); + let stop2 = Arc::clone(&stop); + let interval = Duration::from_secs(config.gc_interval_secs); + + let handle = thread::Builder::new() + .name("blob-gc".into()) + .spawn(move || { + while !stop2.load(Ordering::Acquire) { + thread::park_timeout(interval); + if stop2.load(Ordering::Acquire) { + break; + } + // Seal the active segment first so that the data from + // this GC cycle becomes eligible for compaction. + // Without this, the active segment never seals until + // it reaches SEGMENT_MAX_SIZE (1 GB), and GC would + // have no candidates on small databases. + if let Err(e) = shared3.ensure_sealed() { + tracing::error!("background GC: seal failed: {}", e); + } + match shared3.gc_if_needed() { + Ok(Some(stats)) => { + tracing::info!( + "GC compacted segment {}: {} → {} bytes (kept {}, skipped {})", + stats.segment_id, + stats.bytes_before, + stats.bytes_after, + stats.entries_kept, + stats.entries_skipped, + ); + } + Ok(None) => { + tracing::debug!("GC check: no segment exceeds deleted-ratio threshold"); + } + Err(e) => { + tracing::error!("background GC failed: {}", e); + } + } + } + }) + .expect("failed to spawn blob-gc thread"); + + Some(GcHandle { handle, stop }) + } else { + None + }; + Ok(Self { shared, flush_handle: Mutex::new(flush_handle), + gc_handle: Mutex::new(gc_handle), }) } @@ -195,7 +260,7 @@ impl Engine { inner.append_entry(key, &data, original_len, 0, actual_codec)?; let record = IndexRecord::new(key, segment_id, offset, data_size, 0); - self.shared.bucket_store.insert(record)?; + self.shared.index_store.insert(&record)?; let entry_end = offset + ENTRY_HEADER_SIZE as u64 + data_size as u64; inner.mark_indexed(segment_id, entry_end)?; @@ -204,7 +269,7 @@ impl Engine { } pub fn get(&self, key: &[u8; 32]) -> Result>> { - let record = match self.shared.bucket_store.get(key)? { + let record = match self.shared.index_store.get(key)? { Some(r) => r, None => return Ok(None), }; @@ -228,11 +293,22 @@ impl Engine { let _write_lock = self.shared.write_mutex.lock().unwrap(); let mut inner = self.shared.inner.write().unwrap(); + // Look up the existing record so we can account deleted_bytes on the + // segment that holds the original data — this is what drives GC. + if let Some(rec) = self.shared.index_store.get(key)? { + if !rec.is_tombstone() { + if let Some(stats) = inner.meta.segments.get_mut(&rec.segment_id) { + stats.deleted_bytes += rec.data_size as u64; + stats.recompute_ratio(); + } + } + } + let (segment_id, offset, data_size) = inner.append_entry(*key, &[], 0, 1, Codec::None)?; let record = IndexRecord::new(*key, segment_id, offset, data_size, 1); - self.shared.bucket_store.insert(record)?; + self.shared.index_store.insert(&record)?; let entry_end = offset + ENTRY_HEADER_SIZE as u64 + data_size as u64; inner.mark_indexed(segment_id, entry_end)?; @@ -241,7 +317,7 @@ impl Engine { } pub fn exists(&self, key: &[u8; 32]) -> Result { - self.shared.bucket_store.exists(key) + self.shared.index_store.exists(key) } // ── Batch delete ───────────────────────────────────────────────────── @@ -258,6 +334,16 @@ impl Engine { let mut ends: Vec<(u32, u64)> = Vec::with_capacity(keys.len()); for key in keys { + // Track deleted_bytes for GC threshold on the original segment. + if let Some(rec) = self.shared.index_store.get(key)? { + if !rec.is_tombstone() { + if let Some(stats) = inner.meta.segments.get_mut(&rec.segment_id) { + stats.deleted_bytes += rec.data_size as u64; + stats.recompute_ratio(); + } + } + } + let (segment_id, offset, data_size) = inner.append_entry(*key, &[], 0, 1, Codec::None)?; @@ -268,7 +354,7 @@ impl Engine { inner.flush_active()?; - self.shared.bucket_store.insert_batch(&records)?; + self.shared.index_store.insert_batch(&records)?; for (segment_id, entry_end) in &ends { inner.mark_indexed(*segment_id, *entry_end)?; @@ -312,7 +398,7 @@ impl Engine { inner.flush_active()?; - self.shared.bucket_store.insert_batch(&records)?; + self.shared.index_store.insert_batch(&records)?; for (segment_id, entry_end) in &ends { inner.mark_indexed(*segment_id, *entry_end)?; @@ -324,71 +410,14 @@ impl Engine { // ── GC ────────────────────────────────────────────────────────────── pub fn gc(&self) -> Result> { - // Phase 1: scan segments and write compacted temp file. - // Read-only with respect to Engine state — no write_mutex needed. - let prep = { - let inner = self.shared.inner.read().unwrap(); - gc::gc_prepare( - &inner.root, - &inner.meta, - self.shared.config.gc_deleted_ratio, - )? - }; - - let prep = match prep { - Some(p) => p, - None => return Ok(None), - }; - - // Phase 2: rename temp file + rebuild bucket indices. - // This is the only part that requires exclusive access. - let _write_lock = self.shared.write_mutex.lock().unwrap(); - - let stats = gc::gc_finish(prep)?; - self.shared.file_pool.invalidate(stats.segment_id); - - // Rebuild bucket indices from the updated segment files - let inner = self.shared.inner.write().unwrap(); - let seg_ids: Vec = inner.meta.segments.keys().copied().collect(); - let mut seg_refs: Vec<(u32, PathBuf)> = Vec::with_capacity(seg_ids.len()); - for &id in &seg_ids { - let p = inner.root.join("segments").join(segment::segment_filename(id)); - seg_refs.push((id, p)); - } - let paths: Vec<(u32, &Path)> = - seg_refs.iter().map(|(id, p)| (*id, p.as_path())).collect(); - - let bucket_dir = inner.root.join("buckets"); - BucketStore::rebuild_from_segments(&bucket_dir, &paths)?; - - // Update meta (segment stats changed after GC compaction) - let mut meta = crate::meta::GlobalMeta::load(&inner.root)?; - meta.active_segment_id = inner.meta.active_segment_id; - meta.save(&inner.root)?; - - // Reload bucket store mmaps after GC rewrites - self.shared.bucket_store.reload_all()?; - - Ok(Some(stats)) + self.shared.gc() } /// Run GC only if some segment exceeds the configured deleted-ratio threshold. /// Returns `Ok(None)` immediately without acquiring the write lock when no /// segment qualifies. pub fn gc_if_needed(&self) -> Result> { - let inner = self.shared.inner.read().unwrap(); - let needs_gc = inner - .meta - .segments - .values() - .any(|s| s.sealed && s.deleted_ratio >= self.shared.config.gc_deleted_ratio); - drop(inner); - - if needs_gc { - self.gc() - } else { - Ok(None) - } + self.shared.gc_if_needed() } // ── Flush / Stats / Shutdown ──────────────────────────────────────── @@ -413,7 +442,7 @@ impl Engine { deleted_bytes += seg.deleted_bytes; } - let total_keys = self.shared.bucket_store.total_keys() as u64; + let total_keys = self.shared.index_store.total_keys()? as u64; Ok(Stats { total_keys, @@ -423,19 +452,32 @@ impl Engine { }) } + /// Seal the active segment, forcing it to become a GC candidate. + #[doc(hidden)] + pub fn seal_active_segment(&self) -> Result { + let _write_lock = self.shared.write_mutex.lock().unwrap(); + let mut inner = self.shared.inner.write().unwrap(); + let id = inner.active_writer.id(); + inner.seal_active()?; + Ok(id) + } + pub fn shutdown(&self) -> Result<()> { - // Stop background flush thread first + // Stop background threads first if let Some(fh) = self.flush_handle.lock().unwrap().take() { fh.stop.store(true, Ordering::Release); fh.handle.thread().unpark(); let _ = fh.handle.join(); } + if let Some(gh) = self.gc_handle.lock().unwrap().take() { + gh.stop.store(true, Ordering::Release); + gh.handle.thread().unpark(); + let _ = gh.handle.join(); + } let mut inner = self.shared.inner.write().unwrap(); inner.flush_active()?; - self.shared.bucket_store.compact_all()?; - inner.meta.save(&inner.root)?; tracing::info!("bichon-blob shut down cleanly"); Ok(()) @@ -444,9 +486,122 @@ impl Engine { impl Drop for Engine { fn drop(&mut self) { - if let Err(e) = self.shutdown() { - tracing::error!("bichon-blob shutdown error: {}", e); + // Best-effort shutdown that is panic-safe: only signal threads to + // stop — don't try to acquire write_mutex or inner.write(), which + // would deadlock if we're unwinding from a panic that happened while + // one of those locks was held. + if let Some(fh) = self.flush_handle.lock().ok().and_then(|mut g| g.take()) { + fh.stop.store(true, Ordering::Release); + fh.handle.thread().unpark(); + let _ = fh.handle.join(); } + if let Some(gh) = self.gc_handle.lock().ok().and_then(|mut g| g.take()) { + gh.stop.store(true, Ordering::Release); + gh.handle.thread().unpark(); + let _ = gh.handle.join(); + } + } +} + +// ── EngineShared ──────────────────────────────────────────────────────────── + +impl EngineShared { + /// Seal the active segment if its deleted-ratio exceeds the GC threshold, + /// so that the upcoming GC pass can compact it. Called periodically by + /// the background GC thread. + fn ensure_sealed(&self) -> Result<()> { + let mut inner = self.inner.write().unwrap(); + let active_id = inner.active_writer.id(); + if let Some(stats) = inner.meta.segments.get(&active_id) { + if stats.deleted_ratio >= self.config.gc_deleted_ratio { + inner.seal_active()?; + } + } + Ok(()) + } + + fn gc_if_needed(&self) -> Result> { + let inner = self.inner.read().unwrap(); + let needs_gc = inner + .meta + .segments + .values() + .any(|s| s.sealed && s.deleted_ratio >= self.config.gc_deleted_ratio); + drop(inner); + + if needs_gc { + self.gc() + } else { + Ok(None) + } + } + + fn gc(&self) -> Result> { + // Phase 1: scan only the target segment, consult bucket index per entry. + // Read-only with respect to Engine state — no write_mutex needed. + let prep = { + let inner = self.inner.read().unwrap(); + gc::gc_prepare( + &inner.root, + &inner.meta, + self.config.gc_deleted_ratio, + &self.index_store, + )? + }; + + let mut prep = match prep { + Some(p) => p, + None => return Ok(None), + }; + + // Phase 2: rename temp file + update bucket index. + let _write_lock = self.write_mutex.lock().unwrap(); + + let kept_records = std::mem::take(&mut prep.kept_records); + let deleted_keys = std::mem::take(&mut prep.deleted_keys); + let stats = gc::gc_finish(prep)?; + self.file_pool.invalidate(stats.segment_id); + + // Insert new index records with updated offsets for kept entries. + if !kept_records.is_empty() { + self.index_store.insert_batch(&kept_records)?; + } + + // Remove tombstone IndexRecords that pointed to entries in this + // compacted segment — they are gone now and would accumulate forever. + if !deleted_keys.is_empty() { + self.index_store.delete_batch(&deleted_keys)?; + } + + { + let mut inner = self.inner.write().unwrap(); + + if stats.bytes_after == 0 { + // Segment was completely emptied — remove it from meta first, + // then delete the file. If we crash between the two steps the + // orphaned file is rediscovered on next open and retried. + inner.meta.segments.remove(&stats.segment_id); + inner.meta.save(&inner.root)?; + drop(inner); + + let seg_path = self.inner.read().unwrap() + .root + .join("segments") + .join(segment::segment_filename(stats.segment_id)); + let _ = fs::remove_file(&seg_path); + } else { + // Update segment stats: now smaller and clean. + if let Some(seg_stats) = inner.meta.segments.get_mut(&stats.segment_id) { + seg_stats.total_bytes = stats.bytes_after; + seg_stats.deleted_bytes = 0; + seg_stats.deleted_ratio = 0.0; + seg_stats.indexed_up_to_offset = stats.bytes_after; + } + inner.meta.save(&inner.root)?; + } + } + + Ok(Some(stats)) } } @@ -513,13 +668,6 @@ impl EngineInner { .or_insert_with(|| SegmentStats::new(old_id)); old_stats.sealed = true; - let seg_path = self - .root - .join("segments") - .join(segment::segment_filename(old_id)); - self.readers - .insert(old_id, SegmentReader::open(seg_path, old_id)?); - let new_id = old_id + 1; self.meta.active_segment_id = new_id; let new_path = self diff --git a/crates/blob/src/error.rs b/crates/blob/src/error.rs index 7848965..be9ba0e 100644 --- a/crates/blob/src/error.rs +++ b/crates/blob/src/error.rs @@ -50,4 +50,10 @@ pub enum Error { #[error("Unsupported metadata version {version} in {path}")] UnsupportedMetaVersion { path: PathBuf, version: u32 }, + + #[error("Index database error: {0}")] + IndexDb(String), + + #[error("Database is already open by another process at {path}")] + AlreadyOpen { path: String }, } diff --git a/crates/blob/src/gc.rs b/crates/blob/src/gc.rs index 5829a2d..bd2d9b9 100644 --- a/crates/blob/src/gc.rs +++ b/crates/blob/src/gc.rs @@ -1,7 +1,7 @@ -use std::collections::HashMap; use std::fs; use std::path::{Path, PathBuf}; +use crate::bucket::{IndexRecord, IndexStore}; use crate::error::Result; use crate::meta::GlobalMeta; use crate::segment::{self, SegmentReader, SegmentWriter}; @@ -24,18 +24,25 @@ pub struct GcPrepare { pub bytes_after: u64, pub entries_kept: usize, pub entries_skipped: usize, + /// Index records for kept entries with their new offsets in the compacted segment. + pub kept_records: Vec, + /// Keys whose tombstone IndexRecord should be removed from redb after + /// this segment is compacted (the tombstone entries they pointed to are gone). + pub deleted_keys: Vec<[u8; 32]>, temp_path: PathBuf, seg_path: PathBuf, } -/// Phase 1: pick the sealed segment with the highest deleted_ratio, scan all -/// segments to determine the latest entry for each key, then write a compacted -/// version of the target segment to a temp file. Does NOT rename — the caller -/// should hold the write lock only during `gc_finish`. +/// Phase 1: pick the sealed segment with the highest deleted_ratio, then for +/// each entry in that segment consult the bucket index to decide whether it is +/// still the latest version. Live entries are written to a temp file; stale +/// entries and tombstones are skipped. Does NOT rename — the caller should +/// hold the write lock only during `gc_finish`. pub fn gc_prepare( store_root: &Path, meta: &GlobalMeta, deleted_ratio_threshold: f64, + index_store: &IndexStore, ) -> Result> { let candidate = meta .segments @@ -51,36 +58,15 @@ pub fn gc_prepare( let seg_path = store_root .join("segments") .join(segment::segment_filename(target.segment_id)); - let reader = SegmentReader::open(seg_path.clone(), target.segment_id)?; - // Build a global view: for each key, which entry (segment_id + offset) is the latest? - let mut latest_key: HashMap<[u8; 32], (u32, u64)> = HashMap::new(); - - for &seg_id in meta.segments.keys() { - let rpath = store_root - .join("segments") - .join(segment::segment_filename(seg_id)); - if !rpath.exists() { - continue; - } - let r = SegmentReader::open(rpath, seg_id)?; - let _ = r.scan_entries(0, |entry, offset| { - match latest_key.get(&entry.key) { - Some((existing_seg, existing_off)) => { - if seg_id > *existing_seg - || (seg_id == *existing_seg && offset > *existing_off) - { - latest_key.insert(entry.key, (seg_id, offset)); - } - } - None => { - latest_key.insert(entry.key, (seg_id, offset)); - } - } - Ok(()) - })?; + // Guard against stale meta entries: if the segment file was deleted + // (e.g. after a prior GC emptied it), skip this candidate. + if !seg_path.exists() { + return Ok(None); } + let reader = SegmentReader::open(seg_path.clone(), target.segment_id)?; + // Create temp segment (not renamed yet) let timestamp = std::time::SystemTime::now() .duration_since(std::time::UNIX_EPOCH) @@ -90,24 +76,47 @@ pub fn gc_prepare( let temp_path = store_root.join("segments").join(&temp_name); let mut writer = SegmentWriter::create(temp_path.clone(), target.segment_id)?; + let mut kept_records: Vec = Vec::new(); + let mut deleted_keys: Vec<[u8; 32]> = Vec::new(); let mut bytes_after: u64 = 0; let mut entries_kept: usize = 0; let mut entries_skipped: usize = 0; reader.scan_entries(0, |entry, offset| { if entry.is_tombstone() { + // If this tombstone is the latest version in the index, the key + // must be removed from redb after compaction — otherwise it + // accumulates forever. + match index_store.get(&entry.key)? { + Some(rec) + if rec.segment_id == target.segment_id && rec.offset == offset => + { + deleted_keys.push(entry.key); + } + _ => {} + } entries_skipped += 1; return Ok(()); } - if let Some((latest_seg, latest_off)) = latest_key.get(&entry.key) { - if *latest_seg != target.segment_id || *latest_off != offset { + + // Ask the bucket index whether this entry is still the latest version. + match index_store.get(&entry.key)? { + Some(rec) if rec.segment_id == target.segment_id && rec.offset == offset => { + let new_offset = writer.append(entry)?; + kept_records.push(IndexRecord::new( + entry.key, + target.segment_id, + new_offset, + entry.data.len() as u32, + entry.flags, + )); + bytes_after += entry.data.len() as u64; + entries_kept += 1; + } + _ => { entries_skipped += 1; - return Ok(()); } } - writer.append(entry)?; - bytes_after += entry.data.len() as u64; - entries_kept += 1; Ok(()) })?; @@ -119,6 +128,8 @@ pub fn gc_prepare( bytes_after, entries_kept, entries_skipped, + kept_records, + deleted_keys, temp_path, seg_path, })) diff --git a/crates/blob/src/recovery.rs b/crates/blob/src/recovery.rs index 5bd7226..a7c193d 100644 --- a/crates/blob/src/recovery.rs +++ b/crates/blob/src/recovery.rs @@ -1,23 +1,49 @@ -use std::collections::HashMap; use std::fs; use std::path::Path; -use crate::bucket::{self, IndexRecord}; +use crate::bucket::IndexRecord; use crate::error::Result; use crate::meta::{GlobalMeta, SegmentStats}; use crate::segment::{self, SegmentReader}; /// Recover after a crash: scan any unindexed portions of segments, -/// update bucket files, fix segment stats. -pub fn recover(store_root: &Path, meta: &mut GlobalMeta) -> Result<()> { +/// update segment stats, and return newly discovered index records. +/// +/// The caller is responsible for inserting the returned records into +/// the index store. +pub fn recover(store_root: &Path, meta: &mut GlobalMeta) -> Result> { + scan_segments(store_root, meta, false) +} + +/// Rebuild the entire index from scratch by scanning all segments from +/// offset 0. Used when the index database is corrupted or lost. +/// +/// Resets all segment stats and returns every entry found on disk. +/// The caller should replace the index database before calling this. +pub fn rebuild_index(store_root: &Path, meta: &mut GlobalMeta) -> Result> { + // Reset all segment stats — they'll be recomputed during the scan. + // Also reset indexed_up_to_offset so we scan from 0. + for stats in meta.segments.values_mut() { + stats.total_bytes = 0; + stats.deleted_bytes = 0; + stats.deleted_ratio = 0.0; + stats.indexed_up_to_offset = 0; + } + scan_segments(store_root, meta, true) +} + +/// Common implementation: scan segments and collect index records. +/// When `full_scan` is true, every segment is scanned from offset 0. +fn scan_segments( + store_root: &Path, + meta: &mut GlobalMeta, + full_scan: bool, +) -> Result> { let seg_dir = store_root.join("segments"); if !seg_dir.exists() { fs::create_dir_all(&seg_dir)?; } - let buckets_dir = store_root.join("buckets"); - fs::create_dir_all(&buckets_dir)?; - // Discover all segment files on disk let mut disk_segments: Vec = Vec::new(); if seg_dir.exists() { @@ -36,11 +62,9 @@ pub fn recover(store_root: &Path, meta: &mut GlobalMeta) -> Result<()> { } disk_segments.sort_unstable(); - if disk_segments.is_empty() { - return Ok(()); - } + let mut all_records: Vec = Vec::new(); - // For each segment, scan unindexed portions and append to bucket files + // For each segment, scan unindexed portions for &seg_id in &disk_segments { let seg_path = seg_dir.join(segment::segment_filename(seg_id)); let file_size = fs::metadata(&seg_path)?.len(); @@ -52,9 +76,13 @@ pub fn recover(store_root: &Path, meta: &mut GlobalMeta) -> Result<()> { let is_sealed = seg_id != meta.active_segment_id; stats.sealed = is_sealed; - let scan_start = if stats.indexed_up_to_offset <= file_size { + // Determine scan range. + let scan_start = if full_scan { + 0 + } else if stats.indexed_up_to_offset <= file_size { stats.indexed_up_to_offset } else { + // Segment was replaced (interrupted GC) — full rescan needed. 0 }; @@ -64,18 +92,15 @@ pub fn recover(store_root: &Path, meta: &mut GlobalMeta) -> Result<()> { } let reader = SegmentReader::open(seg_path.clone(), seg_id)?; - let mut new_records: HashMap> = HashMap::new(); let truncation_point = reader.scan_entries(scan_start, |entry, offset| { - let bid = bucket::bucket_id(&entry.key); - let rec = IndexRecord::new( + all_records.push(IndexRecord::new( entry.key, seg_id, offset, entry.data.len() as u32, entry.flags, - ); - new_records.entry(bid).or_default().push(rec); + )); stats.total_bytes += entry.data.len() as u64; if entry.is_tombstone() { @@ -85,19 +110,6 @@ pub fn recover(store_root: &Path, meta: &mut GlobalMeta) -> Result<()> { Ok(()) })?; - // Append new records to bucket files - for (bid, records) in &new_records { - let bf_path = buckets_dir.join(format!("{:02x}.idx", bid)); - use std::io::Write; - let mut file = std::fs::OpenOptions::new() - .create(true) - .append(true) - .open(&bf_path)?; - for r in records { - file.write_all(&r.encode())?; - } - } - // Truncate if tail corruption found if truncation_point < file_size { segment::truncate_segment(&seg_path, truncation_point)?; @@ -110,7 +122,7 @@ pub fn recover(store_root: &Path, meta: &mut GlobalMeta) -> Result<()> { meta.save(store_root)?; - Ok(()) + Ok(all_records) } /// Clean up leftover temp files from interrupted GC. @@ -128,19 +140,5 @@ pub fn cleanup_temp_files(store_root: &Path) -> Result<()> { } } } - // Also cleanup temp bucket files - let buckets_dir = store_root.join("buckets"); - if buckets_dir.exists() { - for entry in fs::read_dir(&buckets_dir)? { - let entry = entry?; - let name = entry.file_name(); - let name_str = name.to_string_lossy(); - if name_str.ends_with(".tmp") { - let path = entry.path(); - tracing::warn!("Removing leftover temp bucket file: {:?}", path); - fs::remove_file(&path)?; - } - } - } Ok(()) } diff --git a/crates/blob/src/types.rs b/crates/blob/src/types.rs index b729359..ae1d7e8 100644 --- a/crates/blob/src/types.rs +++ b/crates/blob/src/types.rs @@ -9,11 +9,8 @@ pub const ENTRY_HEADER_SIZE: usize = 50; /// Index record size: key(32) + segment_id(4) + offset(8) + data_size(4) + flags(1) + _pad(3) + crc32(4) pub const INDEX_RECORD_SIZE: usize = 56; -/// Maximum segment size (256 MB) -pub const SEGMENT_MAX_SIZE: u64 = 256 * 1024 * 1024; - -/// Number of hash buckets (global) -pub const BUCKET_COUNT: u16 = 256; +/// Maximum segment size (1 GB) +pub const SEGMENT_MAX_SIZE: u64 = 1024 * 1024 * 1024; /// Maximum value size (100 MB) pub const MAX_VALUE_SIZE: usize = 100 * 1024 * 1024; @@ -24,8 +21,8 @@ pub const DEFAULT_COMPRESS_THRESHOLD: usize = 4096; /// Default GC deleted ratio threshold pub const DEFAULT_GC_DELETED_RATIO: f64 = 0.30; -/// Default bucket compact threshold (number of pending records before auto-compact) -pub const DEFAULT_COMPACT_THRESHOLD: usize = 10_000; +/// Default GC interval in seconds (5 minutes) +pub const DEFAULT_GC_INTERVAL_SECS: u64 = 300; #[derive(Debug, Clone, Copy, PartialEq, Eq, Serialize, Deserialize)] pub enum Codec { @@ -50,12 +47,15 @@ pub struct Config { pub compress_threshold: usize, pub default_codec: Codec, pub compression_level: i32, - pub compact_threshold: usize, pub gc_deleted_ratio: f64, /// Interval in seconds for periodic background flush (0 = disabled). /// When set, a background thread fsyncs the active segment and saves /// metadata at this interval, bounding recovery time after a crash. pub flush_interval_secs: u64, + /// Interval in seconds for periodic background GC (0 = disabled). + /// When set, a background thread checks whether any sealed segment + /// exceeds the deleted-ratio threshold and compacts it if needed. + pub gc_interval_secs: u64, } impl Default for Config { @@ -64,20 +64,15 @@ impl Default for Config { compress_threshold: DEFAULT_COMPRESS_THRESHOLD, default_codec: Codec::Zstd, compression_level: 0, - compact_threshold: DEFAULT_COMPACT_THRESHOLD, gc_deleted_ratio: DEFAULT_GC_DELETED_RATIO, flush_interval_secs: 0, + gc_interval_secs: 0, } } } impl Config { pub fn validate(&self) -> crate::error::Result<()> { - if self.compact_threshold == 0 { - return Err(crate::error::Error::InvalidConfig( - "compact_threshold must be > 0".into(), - )); - } if self.gc_deleted_ratio <= 0.0 || self.gc_deleted_ratio >= 1.0 { return Err(crate::error::Error::InvalidConfig( "gc_deleted_ratio must be in (0.0, 1.0)".into(), @@ -93,6 +88,11 @@ impl Config { "flush_interval_secs must be 0 (disabled) or >= 5".into(), )); } + if self.gc_interval_secs > 0 && self.gc_interval_secs < 10 { + return Err(crate::error::Error::InvalidConfig( + "gc_interval_secs must be 0 (disabled) or >= 10".into(), + )); + } Ok(()) } } diff --git a/crates/blob/tests/fuzz_test.rs b/crates/blob/tests/fuzz_test.rs new file mode 100644 index 0000000..cb0444f --- /dev/null +++ b/crates/blob/tests/fuzz_test.rs @@ -0,0 +1,380 @@ +/// Fuzz-style tests: randomized operation sequences, corruption injection, +/// and crash-recovery stress testing. +use bichon_blob::{Codec, Config, Engine}; +use rand::RngExt; +use std::collections::HashMap; +use tempfile::TempDir; + +/// Number of iterations for each randomized test. +const FUZZ_OPS: usize = 2000; + +// ── Helpers ────────────────────────────────────────────────────────────────── + +fn make_key(i: u64) -> [u8; 32] { + let mut key = [0u8; 32]; + key[0..8].copy_from_slice(&i.to_le_bytes()); + key +} + +fn make_value(rng: &mut impl RngExt) -> Vec { + let size = match rng.random_range(0..100) { + 0..=4 => rng.random_range(0..64), // tiny + 5..=9 => 0, // empty + 10..=79 => rng.random_range(64..4096), // small + 80..=89 => rng.random_range(4096..65536), // medium + 90..=94 => rng.random_range(65536..500_000), // large + _ => rng.random_range(500_000..2_000_000), // xl (near max) + }; + let mut v = vec![0u8; size]; + rng.fill(&mut v[..]); + v +} + +fn random_codec(rng: &mut impl RngExt) -> Codec { + match rng.random_range(0..4) { + 0 => Codec::None, + 1 => Codec::Zstd, + _ => Codec::Lz4, + } +} + +// ── Fuzz: random operation sequence ───────────────────────────────────────── + +#[test] +fn fuzz_random_ops() { + let mut rng = rand::rng(); + let dir = TempDir::new().unwrap(); + let mut config = Config::default(); + config.flush_interval_secs = 0; // manual flush only + config.gc_interval_secs = 0; + + let engine = Engine::open(dir.path(), config).unwrap(); + + // Oracle: track expected values in memory + let mut oracle: HashMap<[u8; 32], Vec> = HashMap::new(); + let mut next_key = 0u64; + + for _ in 0..FUZZ_OPS { + match rng.random_range(0..100) { + // 45%: put new key + 0..=44 => { + let key = make_key(next_key); + next_key += 1; + let value = make_value(&mut rng); + let codec = random_codec(&mut rng); + engine.put(key, &value, codec).unwrap(); + oracle.insert(key, value); + } + // 20%: put overwrite existing key + 45..=64 => { + if oracle.is_empty() { continue; } + let idx = rng.random_range(0..oracle.len()); + let key = *oracle.iter().nth(idx).unwrap().0; + let value = make_value(&mut rng); + let codec = random_codec(&mut rng); + engine.put(key, &value, codec).unwrap(); + oracle.insert(key, value); + } + // 15%: read and verify + 65..=79 => { + if oracle.is_empty() { continue; } + let idx = rng.random_range(0..oracle.len()); + let key = *oracle.iter().nth(idx).unwrap().0; + let expected = oracle.get(&key).unwrap(); + let got = engine.get(&key).unwrap(); + assert_eq!(got.as_ref(), Some(expected), "key mismatch on read"); + } + // 10%: delete + 80..=89 => { + if oracle.is_empty() { continue; } + let idx = rng.random_range(0..oracle.len()); + let key = *oracle.iter().nth(idx).unwrap().0; + engine.delete(&key).unwrap(); + oracle.remove(&key); + let got = engine.get(&key).unwrap(); + assert_eq!(got, None, "deleted key should return None"); + } + // 5%: read non-existent key + 90..=94 => { + let key = make_key(next_key + rng.random_range(1000u64..10000)); + let got = engine.get(&key).unwrap(); + assert_eq!(got, None, "non-existent key should return None"); + } + // 5%: flush + _ => { + engine.flush().unwrap(); + } + } + } + + // Final verification: all oracle entries must match + for (key, expected) in &oracle { + let got = engine.get(key).unwrap(); + assert_eq!(got.as_ref(), Some(expected), "final verification: key mismatch"); + } +} + +// ── Fuzz: crash + reopen cycle ────────────────────────────────────────────── + +#[test] +fn fuzz_crash_reopen_cycles() { + let mut rng = rand::rng(); + let dir = TempDir::new().unwrap(); + let dir_path = dir.path().to_path_buf(); + + let mut oracle: HashMap<[u8; 32], Vec> = HashMap::new(); + let mut next_key = 0u64; + let cycles = 20; + + for _cycle in 0..cycles { + // Open database + let mut config = Config::default(); + config.flush_interval_secs = 0; + config.gc_interval_secs = 0; + let engine = Engine::open(&dir_path, config).unwrap(); + + // Do some work + let ops = rng.random_range(50..200); + for _ in 0..ops { + match rng.random_range(0..100) { + 0..=50 => { + let key = make_key(next_key); + next_key += 1; + let value = make_value(&mut rng); + if engine.put(key, &value, Codec::Zstd).is_ok() { + oracle.insert(key, value); + } + } + 51..=65 => { + if oracle.is_empty() { continue; } + let idx = rng.random_range(0..oracle.len()); + let key = *oracle.iter().nth(idx).unwrap().0; + engine.delete(&key).unwrap(); + oracle.remove(&key); + } + 66..=85 => { + if oracle.is_empty() { continue; } + let idx = rng.random_range(0..oracle.len()); + let key = *oracle.iter().nth(idx).unwrap().0; + let expected = oracle.get(&key).unwrap(); + if let Ok(Some(got)) = engine.get(&key) { + assert_eq!(&got, expected, "pre-crash read mismatch"); + } + } + _ => { + let _ = engine.flush(); + } + } + } + + // Simulate crash: drop without shutdown + drop(engine); + } + + // Final reopen: all oracle entries must be intact + let config = Config::default(); + let engine = Engine::open(&dir_path, config).unwrap(); + for (key, expected) in &oracle { + let got = engine.get(key).unwrap(); + assert_eq!(got.as_ref(), Some(expected), "after {} crash cycles", cycles); + } +} + +// ── Fuzz: batch operations ────────────────────────────────────────────────── + +#[test] +fn fuzz_batch_ops() { + let mut rng = rand::rng(); + let dir = TempDir::new().unwrap(); + + let config = Config::default(); + let engine = Engine::open(dir.path(), config).unwrap(); + + let mut oracle: HashMap<[u8; 32], Vec> = HashMap::new(); + let mut next_key = 0u64; + + for _ in 0..200 { + match rng.random_range(0..100) { + // 50%: batch write + 0..=49 => { + let batch_size = rng.random_range(1..30); + let entries: Vec<_> = (0..batch_size) + .map(|_| { + let key = make_key(next_key); + next_key += 1; + let value = make_value(&mut rng); + oracle.insert(key, value.clone()); + (key, value, Codec::Zstd) + }) + .collect(); + engine.put_batch(&entries).unwrap(); + } + // 30%: verify random subset + 50..=79 => { + if oracle.is_empty() { continue; } + let n = rng.random_range(1..=20.min(oracle.len())); + for _ in 0..n { + let idx = rng.random_range(0..oracle.len()); + let (key, expected) = oracle.iter().nth(idx).unwrap(); + let got = engine.get(key).unwrap(); + assert_eq!(got.as_ref(), Some(expected)); + } + } + // 20%: batch delete + _ => { + if oracle.is_empty() { continue; } + let n = rng.random_range(1..=20.min(oracle.len())); + let keys: Vec<[u8; 32]> = (0..n) + .map(|_| { + let idx = rng.random_range(0..oracle.len()); + let key = *oracle.iter().nth(idx).unwrap().0; + oracle.remove(&key); + key + }) + .collect(); + engine.delete_batch(&keys).unwrap(); + } + } + } + + for (key, expected) in &oracle { + let got = engine.get(key).unwrap(); + assert_eq!(got.as_ref(), Some(expected)); + } +} + +// ── Fuzz: GC stress ───────────────────────────────────────────────────────── + +#[test] +fn fuzz_gc_stress() { + let mut rng = rand::rng(); + let dir = TempDir::new().unwrap(); + + let mut config = Config::default(); + config.gc_deleted_ratio = 0.1; // aggressive GC trigger + let engine = Engine::open(dir.path(), config).unwrap(); + + let mut oracle: HashMap<[u8; 32], Vec> = HashMap::new(); + let mut next_key = 0u64; + + for round in 0..10 { + // Write a batch of keys + let n = rng.random_range(50..150); + let mut round_keys: Vec<[u8; 32]> = Vec::new(); + let value_size = rng.random_range(100..10000); + let value: Vec = (0..value_size).map(|_| rng.random::()).collect(); + + for _ in 0..n { + let key = make_key(next_key); + next_key += 1; + engine.put(key, &value, Codec::None).unwrap(); + oracle.insert(key, value.clone()); + round_keys.push(key); + } + + // Delete some fraction + let delete_frac: f64 = rng.random_range(0.2..0.8); + let delete_count = (round_keys.len() as f64 * delete_frac) as usize; + for _ in 0..delete_count { + let idx = rng.random_range(0..round_keys.len()); + let key = round_keys.swap_remove(idx); + engine.delete(&key).unwrap(); + oracle.remove(&key); + } + + // Seal and run GC + engine.seal_active_segment().unwrap(); + let _ = engine.gc().unwrap(); + + // Verify all remaining oracle entries + for (key, expected) in &oracle { + let got = engine.get(key).unwrap(); + assert_eq!( + got.as_ref(), + Some(expected), + "GC round {}: key mismatch", + round + ); + } + } +} + +// ── Fuzz: bitflip corruption resilience ───────────────────────────────────── + +#[test] +fn fuzz_corruption_resilience() { + let mut rng = rand::rng(); + let dir = TempDir::new().unwrap(); + let dir_path = dir.path().to_path_buf(); + + let config = Config::default(); + let engine = Engine::open(&dir_path, config).unwrap(); + + // Write known data + let mut good_keys: Vec<[u8; 32]> = Vec::new(); + let n = 100u64; + for i in 0..n { + let key = make_key(i); + let value = vec![i as u8; 2048]; + engine.put(key, &value, Codec::None).unwrap(); + good_keys.push(key); + } + engine.flush().unwrap(); + + // Seal so segment file is durable on disk + engine.seal_active_segment().unwrap(); + engine.shutdown().unwrap(); + drop(engine); + + // Find segment files and corrupt random bytes + let seg_dir = dir_path.join("segments"); + let mut seg_files: Vec<_> = std::fs::read_dir(&seg_dir) + .unwrap() + .filter_map(|e| e.ok()) + .filter(|e| { + e.file_name() + .to_string_lossy() + .ends_with(".seg") + }) + .collect(); + seg_files.sort_by_key(|e| e.file_name()); + + // Corrupt 3 random bytes in the first segment + if let Some(seg) = seg_files.first() { + let path = seg.path(); + let mut data = std::fs::read(&path).unwrap(); + if data.len() > 100 { + for _ in 0..3 { + let pos = rng.random_range(50..data.len()); + data[pos] ^= 0xFF; // flip all bits + } + std::fs::write(&path, &data).unwrap(); + } + } + + // Reopen: recovery must succeed (not panic), even if some keys are lost + let config = Config::default(); + let engine = Engine::open(&dir_path, config).unwrap(); + + // At least some uncorrupted keys should still be readable + let mut readable = 0; + let mut corrupted = 0; + for key in &good_keys { + match engine.get(key) { + Ok(Some(_)) => readable += 1, + Ok(None) => { /* key might be lost due to corruption */ } + Err(_) => corrupted += 1, + } + } + // The vast majority should still be ok (corruption only hit 3 bytes in one segment) + assert!( + readable + corrupted > 0, + "at least some outcomes should be observable" + ); + assert!( + readable > n as usize / 2, + "majority of keys should survive localized corruption ({} of {})", + readable, + n + ); +} diff --git a/crates/blob/tests/integration_test.rs b/crates/blob/tests/integration_test.rs index f16971d..2732744 100644 --- a/crates/blob/tests/integration_test.rs +++ b/crates/blob/tests/integration_test.rs @@ -229,7 +229,7 @@ fn test_batch_write_persistence() { fn test_invalid_config_rejected() { let dir = TempDir::new().unwrap(); let mut config = Config::default(); - config.compact_threshold = 0; + config.compression_level = -1; assert!(Engine::open(dir.path(), config).is_err()); let mut config = Config::default(); @@ -358,3 +358,88 @@ fn test_meta_bin_durability() { let read = engine.get(&[0x42u8; 32]).unwrap(); assert_eq!(read, Some(b"durable".to_vec())); } + +#[test] +fn test_gc_deletes_empty_segment() { + let dir = TempDir::new().unwrap(); + let engine = Engine::open(dir.path(), Config::default()).unwrap(); + + // Write data, seal, then delete everything so the sealed segment + // becomes 100% garbage. GC should remove the segment file entirely. + let n = 50u64; + for i in 0..n { + let mut key = [0u8; 32]; + key[0..8].copy_from_slice(&i.to_le_bytes()); + engine.put(key, &vec![b'X'; 4096], Codec::None).unwrap(); + } + + // Force seal so the segment becomes a GC candidate. + let sealed_id = engine.seal_active_segment().unwrap(); + + // Delete all keys — the sealed segment is now entirely garbage. + let keys: Vec<[u8; 32]> = (0..n) + .map(|i| { + let mut key = [0u8; 32]; + key[0..8].copy_from_slice(&i.to_le_bytes()); + key + }) + .collect(); + engine.delete_batch(&keys).unwrap(); + + // Run GC — this should empty and then delete the sealed segment. + let result = engine.gc().unwrap(); + assert!(result.is_some(), "GC should have found a candidate"); + let stats = result.unwrap(); + assert_eq!(stats.segment_id, sealed_id); + assert_eq!(stats.bytes_after, 0); + + // Segment file must be deleted. + let seg_path = dir + .path() + .join("segments") + .join(format!("{:08}.seg", sealed_id)); + assert!( + !seg_path.exists(), + "emptied segment file should have been deleted, but {:?} exists", + seg_path + ); + + // Subsequent reads / writes must still work (no corruption). + let new_key = [0x99u8; 32]; + engine + .put(new_key, b"post-gc data", Codec::None) + .unwrap(); + let result = engine.get(&new_key).unwrap(); + assert_eq!(result, Some(b"post-gc data".to_vec())); + + // Engine stats should be consistent. + let stats = engine.stats().unwrap(); + assert_eq!(stats.total_keys, 1); + + // Shutdown and reopen — persistence must be intact. + drop(engine); + let engine = Engine::open(dir.path(), Config::default()).unwrap(); + let result = engine.get(&new_key).unwrap(); + assert_eq!(result, Some(b"post-gc data".to_vec())); + let result = engine.get(&keys[0]).unwrap(); + assert_eq!(result, None); +} + +#[test] +fn test_file_lock_prevents_concurrent_open() { + let dir = TempDir::new().unwrap(); + let _engine1 = Engine::open(dir.path(), Config::default()).unwrap(); + let result = Engine::open(dir.path(), Config::default()); + assert!(result.is_err(), "second open on same directory must fail"); +} + +#[test] +fn test_file_lock_released_after_close() { + let dir = TempDir::new().unwrap(); + { + let _engine = Engine::open(dir.path(), Config::default()).unwrap(); + } + // Lock should be released after engine is dropped + let engine = Engine::open(dir.path(), Config::default()); + assert!(engine.is_ok(), "reopen after close must succeed"); +} diff --git a/crates/core/src/store/blob.rs b/crates/core/src/store/blob.rs index d1852f0..145b495 100644 --- a/crates/core/src/store/blob.rs +++ b/crates/core/src/store/blob.rs @@ -119,6 +119,7 @@ impl BlobManager { config.default_codec = Codec::Zstd; config.compress_threshold = 1024; config.flush_interval_secs = 60; + config.gc_interval_secs = 300; let engine = Engine::open(&blob_dir, config) .expect("Failed to initialize blob engine: Check disk space and permissions."); diff --git a/web/src/components/mail-iframe.tsx b/web/src/components/mail-iframe.tsx index a3b3962..fe52c13 100644 --- a/web/src/components/mail-iframe.tsx +++ b/web/src/components/mail-iframe.tsx @@ -17,7 +17,7 @@ // along with this program. If not, see . -import React from 'react'; +import React, { useCallback, useRef, useState } from 'react'; interface EmailIframeProps { emailHtml: string; @@ -25,16 +25,36 @@ interface EmailIframeProps { } const EmailIframe: React.FC = ({ emailHtml, height }) => { - const encodedHtml = encodeURIComponent(emailHtml); - const iframeSrc = `data:text/html;charset=utf-8,${encodedHtml}`; + const iframeRef = useRef(null); + const [iframeHeight, setIframeHeight] = useState(height ?? 300); + + const onLoad = useCallback(() => { + try { + const doc = iframeRef.current?.contentWindow?.document; + if (doc?.body) { + const h = Math.max( + doc.body.scrollHeight, + doc.body.offsetHeight, + doc.documentElement.scrollHeight, + doc.documentElement.offsetHeight, + ); + if (h > 0) setIframeHeight(h + 20); + } + } catch { + // sandbox prevents access — keep default height + } + }, []); return (