mirror of
https://github.com/rustmailer/bichon.git
synced 2026-08-03 07:48:34 +02:00
update
This commit is contained in:
+20
-5
@@ -33,13 +33,13 @@ magic(4) crc32(4) flags(1) codec(1) key(32) raw_size(4) data_size(4) data
|
|||||||
- `raw_size`: original uncompressed size
|
- `raw_size`: original uncompressed size
|
||||||
- `data_size`: on-disk data size (after compression)
|
- `data_size`: on-disk data size (after compression)
|
||||||
|
|
||||||
### Bucket index record format (52 bytes)
|
### Bucket index record format (56 bytes)
|
||||||
|
|
||||||
```
|
```
|
||||||
key(32) segment_id(4) offset(8) data_size(4) flags(1) _pad(3)
|
key(32) segment_id(4) offset(8) data_size(4) flags(1) _pad(3) crc32(4)
|
||||||
```
|
```
|
||||||
|
|
||||||
Index 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.
|
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.
|
||||||
|
|
||||||
## Read path
|
## Read path
|
||||||
|
|
||||||
@@ -48,11 +48,13 @@ get(key)
|
|||||||
→ bucket_id = (key[0..2] as u16) % 256
|
→ bucket_id = (key[0..2] as u16) % 256
|
||||||
→ check pending HashMap (most recent wins)
|
→ check pending HashMap (most recent wins)
|
||||||
→ binary search mmap'd bucket file
|
→ binary search mmap'd bucket file
|
||||||
→ IndexRecord → (segment_id, offset, data_size)
|
→ IndexRecord CRC32 verify → (segment_id, offset, data_size)
|
||||||
→ pread entry from segment file
|
→ pread entry from segment file
|
||||||
→ CRC32 verify → decompress → return value
|
→ 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.
|
||||||
|
|
||||||
## Write path
|
## Write path
|
||||||
|
|
||||||
```
|
```
|
||||||
@@ -74,10 +76,23 @@ Deletes are **tombstones** — an entry with `flags=1` and empty data is appende
|
|||||||
|
|
||||||
**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.
|
**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.
|
||||||
|
|
||||||
|
## Data integrity
|
||||||
|
|
||||||
|
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 |
|
||||||
|
| 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`.
|
||||||
|
|
||||||
## Crash recovery
|
## Crash recovery
|
||||||
|
|
||||||
- Temp files from interrupted GC are cleaned up on open.
|
- 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.
|
- 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.
|
||||||
- Partial writes at the tail of a segment (detected via CRC32 mismatch near EOF) are truncated.
|
- 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`).
|
- Buckets are always repairable by re-scanning segments (`rebuild_from_segments`).
|
||||||
|
|
||||||
|
|||||||
+57
-15
@@ -41,23 +41,36 @@ impl IndexRecord {
|
|||||||
buf[36..44].copy_from_slice(&self.offset.to_le_bytes());
|
buf[36..44].copy_from_slice(&self.offset.to_le_bytes());
|
||||||
buf[44..48].copy_from_slice(&self.data_size.to_le_bytes());
|
buf[44..48].copy_from_slice(&self.data_size.to_le_bytes());
|
||||||
buf[48] = self.flags;
|
buf[48] = self.flags;
|
||||||
|
let crc = crate::checksum::crc32(&buf[..52]);
|
||||||
|
buf[52..56].copy_from_slice(&crc.to_le_bytes());
|
||||||
buf
|
buf
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn decode(buf: &[u8; INDEX_RECORD_SIZE]) -> Self {
|
pub fn decode(buf: &[u8; INDEX_RECORD_SIZE]) -> crate::error::Result<Self> {
|
||||||
|
let stored_crc = u32::from_le_bytes(buf[52..56].try_into().unwrap());
|
||||||
|
let computed = crate::checksum::crc32(&buf[..52]);
|
||||||
|
if stored_crc != computed {
|
||||||
|
return Err(crate::error::Error::BucketIndexCorrupt {
|
||||||
|
path: std::path::PathBuf::new(),
|
||||||
|
reason: format!(
|
||||||
|
"CRC mismatch: stored=0x{:08X} computed=0x{:08X}",
|
||||||
|
stored_crc, computed
|
||||||
|
),
|
||||||
|
});
|
||||||
|
}
|
||||||
let mut key = [0u8; 32];
|
let mut key = [0u8; 32];
|
||||||
key.copy_from_slice(&buf[0..32]);
|
key.copy_from_slice(&buf[0..32]);
|
||||||
let segment_id = u32::from_le_bytes(buf[32..36].try_into().unwrap());
|
let segment_id = u32::from_le_bytes(buf[32..36].try_into().unwrap());
|
||||||
let offset = u64::from_le_bytes(buf[36..44].try_into().unwrap());
|
let offset = u64::from_le_bytes(buf[36..44].try_into().unwrap());
|
||||||
let data_size = u32::from_le_bytes(buf[44..48].try_into().unwrap());
|
let data_size = u32::from_le_bytes(buf[44..48].try_into().unwrap());
|
||||||
let flags = buf[48];
|
let flags = buf[48];
|
||||||
Self {
|
Ok(Self {
|
||||||
key,
|
key,
|
||||||
segment_id,
|
segment_id,
|
||||||
offset,
|
offset,
|
||||||
data_size,
|
data_size,
|
||||||
flags,
|
flags,
|
||||||
}
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -159,7 +172,7 @@ impl BucketStore {
|
|||||||
let result = binary_search_records(bytes, key, count);
|
let result = binary_search_records(bytes, key, count);
|
||||||
match result {
|
match result {
|
||||||
Some(idx) => {
|
Some(idx) => {
|
||||||
let rec = read_record_at(bytes, idx);
|
let rec = read_record_at(bytes, idx)?;
|
||||||
Ok(if rec.is_tombstone() { None } else { Some(rec) })
|
Ok(if rec.is_tombstone() { None } else { Some(rec) })
|
||||||
}
|
}
|
||||||
None => Ok(None),
|
None => Ok(None),
|
||||||
@@ -239,13 +252,14 @@ impl BucketStore {
|
|||||||
let bytes: &[u8] = &s.mmap;
|
let bytes: &[u8] = &s.mmap;
|
||||||
let n = bytes.len() / INDEX_RECORD_SIZE;
|
let n = bytes.len() / INDEX_RECORD_SIZE;
|
||||||
for i in 0..n {
|
for i in 0..n {
|
||||||
let rec = read_record_at(bytes, i);
|
if let Ok(rec) = read_record_at(bytes, i) {
|
||||||
// Skip keys that are overridden in pending
|
// Skip keys that are overridden in pending
|
||||||
if s.pending.contains_key(&rec.key) {
|
if s.pending.contains_key(&rec.key) {
|
||||||
continue;
|
continue;
|
||||||
}
|
}
|
||||||
if !rec.is_tombstone() {
|
if !rec.is_tombstone() {
|
||||||
count += 1;
|
count += 1;
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -268,7 +282,15 @@ impl BucketStore {
|
|||||||
let n = bytes.len() / INDEX_RECORD_SIZE;
|
let n = bytes.len() / INDEX_RECORD_SIZE;
|
||||||
all.reserve(n + state.pending.len());
|
all.reserve(n + state.pending.len());
|
||||||
for i in 0..n {
|
for i in 0..n {
|
||||||
all.push(read_record_at(bytes, i));
|
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() {
|
for r in state.pending.values() {
|
||||||
all.push(r.clone());
|
all.push(r.clone());
|
||||||
@@ -400,7 +422,15 @@ fn load_records_from_file(path: &Path) -> Result<Vec<IndexRecord>> {
|
|||||||
reason: "unexpected file size".into(),
|
reason: "unexpected file size".into(),
|
||||||
}
|
}
|
||||||
})?;
|
})?;
|
||||||
records.push(IndexRecord::decode(buf));
|
match IndexRecord::decode(buf) {
|
||||||
|
Ok(rec) => records.push(rec),
|
||||||
|
Err(_) => {
|
||||||
|
tracing::warn!(
|
||||||
|
"Bucket file {:?} record {} CRC mismatch, skipping",
|
||||||
|
path, i
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
}
|
}
|
||||||
if remainder > 0 {
|
if remainder > 0 {
|
||||||
tracing::warn!(
|
tracing::warn!(
|
||||||
@@ -458,7 +488,7 @@ fn read_key_at(bytes: &[u8], idx: usize) -> &[u8; 32] {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/// Read a full IndexRecord at index `idx`.
|
/// Read a full IndexRecord at index `idx`.
|
||||||
fn read_record_at(bytes: &[u8], idx: usize) -> IndexRecord {
|
fn read_record_at(bytes: &[u8], idx: usize) -> crate::error::Result<IndexRecord> {
|
||||||
let start = idx * INDEX_RECORD_SIZE;
|
let start = idx * INDEX_RECORD_SIZE;
|
||||||
let buf: &[u8; INDEX_RECORD_SIZE] = bytes[start..start + INDEX_RECORD_SIZE]
|
let buf: &[u8; INDEX_RECORD_SIZE] = bytes[start..start + INDEX_RECORD_SIZE]
|
||||||
.try_into()
|
.try_into()
|
||||||
@@ -486,10 +516,22 @@ mod tests {
|
|||||||
key[0..4].copy_from_slice(&[1, 2, 3, 4]);
|
key[0..4].copy_from_slice(&[1, 2, 3, 4]);
|
||||||
let rec = IndexRecord::new(key, 5, 12345, 500, 0);
|
let rec = IndexRecord::new(key, 5, 12345, 500, 0);
|
||||||
let encoded = rec.encode();
|
let encoded = rec.encode();
|
||||||
let decoded = IndexRecord::decode(&encoded);
|
assert_eq!(encoded.len(), INDEX_RECORD_SIZE);
|
||||||
|
let decoded = IndexRecord::decode(&encoded).unwrap();
|
||||||
assert_eq!(rec, decoded);
|
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]
|
#[test]
|
||||||
fn test_bucket_id_deterministic() {
|
fn test_bucket_id_deterministic() {
|
||||||
let mut key = [0u8; 32];
|
let mut key = [0u8; 32];
|
||||||
|
|||||||
@@ -100,6 +100,8 @@ impl Engine {
|
|||||||
|
|
||||||
crate::recovery::cleanup_temp_files(path)?;
|
crate::recovery::cleanup_temp_files(path)?;
|
||||||
crate::recovery::recover(path, &mut meta)?;
|
crate::recovery::recover(path, &mut meta)?;
|
||||||
|
// Recovery may have appended new records to bucket files — reload mmaps.
|
||||||
|
bucket_store.reload_all()?;
|
||||||
|
|
||||||
let seg_path = path
|
let seg_path = path
|
||||||
.join("segments")
|
.join("segments")
|
||||||
|
|||||||
@@ -6,8 +6,8 @@ pub const ENTRY_MAGIC: u32 = 0xB3DB_0001;
|
|||||||
/// Fixed header size: magic(4) + crc32(4) + flags(1) + codec(1) + key(32) + raw_size(4) + data_size(4)
|
/// Fixed header size: magic(4) + crc32(4) + flags(1) + codec(1) + key(32) + raw_size(4) + data_size(4)
|
||||||
pub const ENTRY_HEADER_SIZE: usize = 50;
|
pub const ENTRY_HEADER_SIZE: usize = 50;
|
||||||
|
|
||||||
/// Index record size: key(32) + segment_id(4) + offset(8) + data_size(4) + flags(1) + _pad(3)
|
/// 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 = 52;
|
pub const INDEX_RECORD_SIZE: usize = 56;
|
||||||
|
|
||||||
/// Maximum segment size (256 MB)
|
/// Maximum segment size (256 MB)
|
||||||
pub const SEGMENT_MAX_SIZE: u64 = 256 * 1024 * 1024;
|
pub const SEGMENT_MAX_SIZE: u64 = 256 * 1024 * 1024;
|
||||||
|
|||||||
@@ -453,6 +453,50 @@ fn test_global_dedup_after_crash() {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
// 8. Recovery + mmap consistency: entries re-indexed by recovery are readable
|
||||||
|
// ---------------------------------------------------------------------------
|
||||||
|
|
||||||
|
#[test]
|
||||||
|
fn test_recovery_reloads_bucket_mmaps() {
|
||||||
|
use bichon_blob::meta::GlobalMeta;
|
||||||
|
|
||||||
|
let dir = TempDir::new().unwrap();
|
||||||
|
let value = make_value(8192);
|
||||||
|
let n = 50u64;
|
||||||
|
|
||||||
|
// Phase 1: write data with clean shutdown, then corrupt meta to force
|
||||||
|
// re-indexing on next open (simulates crash where mark_indexed didn't run).
|
||||||
|
{
|
||||||
|
let engine = Engine::open(dir.path(), Config::default()).unwrap();
|
||||||
|
for i in 0..n {
|
||||||
|
engine.put(make_key(i), &value, Codec::None).unwrap();
|
||||||
|
}
|
||||||
|
} // clean shutdown — everything is indexed and fsynced
|
||||||
|
|
||||||
|
// Corrupt meta: zero out indexed_up_to_offset so recovery re-scans.
|
||||||
|
let mut meta = GlobalMeta::load(dir.path()).expect("failed to load meta");
|
||||||
|
for seg in meta.segments.values_mut() {
|
||||||
|
seg.indexed_up_to_offset = 0;
|
||||||
|
}
|
||||||
|
meta.save(dir.path()).expect("failed to save corrupted meta");
|
||||||
|
|
||||||
|
// Phase 2: reopen — recovery must re-scan the segment and append to bucket
|
||||||
|
// files. After the fix, reload_all() ensures the mmaps include recovered data.
|
||||||
|
{
|
||||||
|
let engine = Engine::open(dir.path(), Config::default()).unwrap();
|
||||||
|
for i in 0..n {
|
||||||
|
let result = engine.get(&make_key(i)).unwrap();
|
||||||
|
assert_eq!(
|
||||||
|
result,
|
||||||
|
Some(value.clone()),
|
||||||
|
"key {} should be readable after recovery reload",
|
||||||
|
i
|
||||||
|
);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
// Helpers
|
// Helpers
|
||||||
// ---------------------------------------------------------------------------
|
// ---------------------------------------------------------------------------
|
||||||
|
|||||||
Reference in New Issue
Block a user