From 005012edcf0804a3f238e45ddfbf762130b0ee65 Mon Sep 17 00:00:00 2001 From: Eliah Rusin Date: Mon, 7 Sep 2026 23:50:49 +0300 Subject: [PATCH] rust volume: quick-repair redb on durable checkpoints (#11203) * rust volume: insert redb rebuilds in needle-id order Unlink the .rdb before create (create does not truncate). Collapse last-write-wins, then insert live keys sorted so 4.2.0 packs leaves. * rust volume: rebuild redb from a BTreeMap and clear leftover keys Peak rebuild memory is one ordered map instead of HashMap + Vec + stable-sort scratch. Unlink stays best-effort: if it fails, retain clears the leftover table before sorted insert. Compute idx metrics before the write so a read error does not unlink a committed .rdb. * rust volume: drop the extra redb read transaction on put/delete put uses insert()'s previous value. delete gets then inserts the tombstone in the same write transaction. Truncate the .idx row on any failed redb write after the append. * rust volume: unpack redb blobs through packed_to_needle_value save_to_idx, ascending_visit, and collect_entries used the same length-check copy as get. Route them through the helper so a wrong-length value is absent everywhere, not a panic. * rust volume: quick-repair redb on durable checkpoints set_quick_repair(true) on the durable checkpoint transaction so an OOM-killed volume server opens without a full-file repair scan. * rust volume: reopen redb from .idx on non-poisoned commit error redb 4.2.0 can make a Durability::None commit visible before returning Err(CommitError::Storage(..)). In that state the database refuses further write transactions, so truncating the .idx row (the old behavior) would leave a redb-only put or tombstone that the stored idx_size makes the reload skip. Distinguish CommitError::TransactionPoisoned (txn rolled back, db still usable -- truncate the orphan .idx row as before) from other commit errors (change may be visible, db refuses writes -- keep the .idx row, close the database, and reopen from .idx to repair redb's internal state). db becomes Option so reopen_from_idx can drop the old file lock before load_from_idx opens the same path. rdb_path, version, and cache_bytes are stored so the reopen uses the same configuration. * rust volume: truncate .idx row when redb is closed in put put appends the .idx entry before acquiring the write transaction. When db_or_err() fails (db is None after a failed reopen), the ? returned without calling truncate_idx_to_offset, so a write reported failed remained in the authoritative .idx and was replayed on restart. Handle db_or_err() explicitly and truncate the orphan .idx row before returning the error, matching the existing handling for begin_write, open_table, and insert failures. --------- Co-authored-by: Chris Lu Co-authored-by: Chris Lu --- seaweed-volume/src/storage/needle_map.rs | 227 +++++++++++++++++++---- 1 file changed, 188 insertions(+), 39 deletions(-) diff --git a/seaweed-volume/src/storage/needle_map.rs b/seaweed-volume/src/storage/needle_map.rs index f08e896c2..a1d15111e 100644 --- a/seaweed-volume/src/storage/needle_map.rs +++ b/seaweed-volume/src/storage/needle_map.rs @@ -443,7 +443,15 @@ const REDB_CHECKPOINT_INTERVAL: u32 = 1000; /// Low memory usage — data lives on disk behind a small, bounded redb page /// cache sized by `NeedleMapKind::redb_cache_bytes`. pub struct RedbNeedleMap { - db: Database, + /// `Option` so a non-`TransactionPoisoned` commit error can drop the + /// `Database` (releasing the file lock) before reopening from `.idx`. + db: Option, + /// Path of the `.rdb` file — needed to reopen after a commit error. + rdb_path: String, + /// Volume version, for replaying `.idx` rows on reopen. + version: Version, + /// redb page-cache size, for reopening with the same configuration. + cache_bytes: usize, metric: NeedleMapMetric, idx_file: Option>, idx_file_offset: u64, @@ -464,6 +472,16 @@ impl RedbNeedleMap { Ok(txn) } + /// Return a reference to the open `Database`, or an error if a previous + /// non-poisoned commit error followed by a failed reopen left it closed. + fn db_or_err(&self) -> io::Result<&Database> { + self.db.as_ref().ok_or_else(|| { + io::Error::other( + "redb database is closed after a commit error and failed reopen from .idx", + ) + }) + } + /// True once `REDB_CHECKPOINT_INTERVAL` puts/deletes have been committed /// non-durably. The caller (the volume) must then flush the .dat and /// call [`checkpoint`](Self::checkpoint): a checkpoint makes the index @@ -498,7 +516,7 @@ impl RedbNeedleMap { if sync_idx { self.sync()?; } - let mut txn = self.db.begin_write().map_err(|e| { + let mut txn = self.db_or_err()?.begin_write().map_err(|e| { io::Error::new(io::ErrorKind::Other, format!("redb begin_write: {}", e)) })?; txn.set_quick_repair(true); @@ -538,7 +556,10 @@ impl RedbNeedleMap { .map_err(|e| io::Error::new(io::ErrorKind::Other, format!("redb commit: {}", e)))?; Ok(RedbNeedleMap { - db, + db: Some(db), + rdb_path: db_path.to_string(), + version: Version::current(), + cache_bytes, metric: NeedleMapMetric::default(), idx_file: None, idx_file_offset: 0, @@ -549,7 +570,7 @@ impl RedbNeedleMap { /// Save the .idx file size into redb metadata so we can detect whether /// the .rdb is up-to-date on the next startup. fn save_idx_size_meta(&self, idx_size: u64) -> io::Result<()> { - let txn = Self::begin_write_no_fsync(&self.db)?; + let txn = Self::begin_write_no_fsync(self.db_or_err()?)?; { let mut meta = txn.open_table(META_TABLE).map_err(|e| { io::Error::new(io::ErrorKind::Other, format!("redb open meta: {}", e)) @@ -567,7 +588,7 @@ impl RedbNeedleMap { /// Read the stored .idx file size from redb metadata. fn read_idx_size_meta(&self) -> io::Result> { let txn = self - .db + .db_or_err()? .begin_read() .map_err(|e| io::Error::new(io::ErrorKind::Other, format!("redb begin_read: {}", e)))?; let meta = txn @@ -630,7 +651,10 @@ impl RedbNeedleMap { .map_err(|e| io::Error::new(io::ErrorKind::Other, format!("redb open: {}", e)))?; let mut nm = RedbNeedleMap { - db, + db: Some(db), + rdb_path: db_path.to_string(), + version, + cache_bytes, metric: NeedleMapMetric::default(), idx_file: None, idx_file_offset: 0, @@ -657,7 +681,7 @@ impl RedbNeedleMap { if stored_idx_size < idx_size { // .idx grew — replay new entries incrementally let start_entry = stored_idx_size / NEEDLE_MAP_ENTRY_SIZE as u64; - let txn = Self::begin_write_no_fsync(&nm.db)?; + let txn = Self::begin_write_no_fsync(nm.db.as_ref().unwrap())?; { let mut table = txn.open_table(NEEDLE_TABLE).map_err(|e| { io::Error::new(io::ErrorKind::Other, format!("redb open_table: {}", e)) @@ -747,6 +771,7 @@ impl RedbNeedleMap { return Err(e); } }; + nm.version = version; let result = (|| -> io::Result<()> { // Seek independently of walk_index_file. Run before the write txn @@ -763,7 +788,7 @@ impl RedbNeedleMap { Ok(()) })?; - let txn = Self::begin_write_no_fsync(&nm.db)?; + let txn = Self::begin_write_no_fsync(nm.db.as_ref().unwrap())?; { let mut table = txn.open_table(NEEDLE_TABLE).map_err(|e| { io::Error::new(io::ErrorKind::Other, format!("redb open_table: {}", e)) @@ -853,26 +878,72 @@ impl RedbNeedleMap { let key_u64: u64 = key.into(); let packed = pack_needle_value(&NeedleValue { offset, size }); - let old = match (|| -> io::Result> { - let txn = Self::begin_write_no_fsync(&self.db)?; - let old = { - let mut table = txn.open_table(NEEDLE_TABLE).map_err(|e| { - io::Error::new(io::ErrorKind::Other, format!("redb open_table: {}", e)) - })?; - let prev = table.insert(key_u64, packed.as_slice()).map_err(|e| { - io::Error::new(io::ErrorKind::Other, format!("redb insert: {}", e)) - })?; - prev.and_then(|g| packed_to_needle_value(g.value())) + let old = { + let db = match self.db_or_err() { + Ok(db) => db, + Err(e) => { + self.truncate_idx_to_offset(); + return Err(e); + } }; - txn.commit().map_err(|e| { - io::Error::new(io::ErrorKind::Other, format!("redb commit: {}", e)) - })?; - Ok(old) - })() { - Ok(old) => old, - Err(e) => { - self.truncate_idx_to_offset(); - return Err(e); + let txn = Self::begin_write_no_fsync(db); + let txn = match txn { + Ok(txn) => txn, + Err(e) => { + self.truncate_idx_to_offset(); + return Err(e); + } + }; + // Insert and extract the old value inside a block so the table + // and AccessGuard are dropped before commit (which moves txn). + let old = { + let mut table = match txn.open_table(NEEDLE_TABLE) { + Ok(t) => t, + Err(e) => { + self.truncate_idx_to_offset(); + return Err(io::Error::new( + io::ErrorKind::Other, + format!("redb open_table: {}", e), + )); + } + }; + let result = match table.insert(key_u64, packed.as_slice()) { + Ok(prev) => prev.and_then(|g| packed_to_needle_value(g.value())), + Err(e) => { + self.truncate_idx_to_offset(); + return Err(io::Error::new( + io::ErrorKind::Other, + format!("redb insert: {}", e), + )); + } + }; + result + }; + match txn.commit() { + Ok(()) => old, + Err(redb::CommitError::TransactionPoisoned) => { + // Transaction rolled back, database still usable: + // truncate the orphan .idx row. + self.truncate_idx_to_offset(); + return Err(io::Error::new( + io::ErrorKind::Other, + "redb commit: Transaction was poisoned by a panic", + )); + } + Err(e) => { + // Non-poisoned commit error: the change may be + // visible and redb refuses further writes. Keep + // the .idx row (do NOT truncate) and reopen from + // .idx to repair redb's internal state. + let err = io::Error::new(io::ErrorKind::Other, format!("redb commit: {}", e)); + if let Err(reopen_err) = self.reopen_from_idx() { + tracing::warn!( + "redb reopen after put commit error failed: {}", + reopen_err + ); + } + return Err(err); + } } }; @@ -895,7 +966,7 @@ impl RedbNeedleMap { /// Internal get that returns io::Result for error propagation. fn get_internal(&self, key_u64: u64) -> io::Result> { let txn = self - .db + .db_or_err()? .begin_read() .map_err(|e| io::Error::new(io::ErrorKind::Other, format!("redb begin_read: {}", e)))?; let table = txn @@ -918,7 +989,7 @@ impl RedbNeedleMap { /// Mark a needle as deleted. Appends tombstone to .idx file, negates size in redb. pub fn delete(&mut self, key: NeedleId, offset: Offset) -> io::Result> { let key_u64: u64 = key.into(); - let txn = Self::begin_write_no_fsync(&self.db)?; + let txn = Self::begin_write_no_fsync(self.db_or_err()?)?; let mut table = txn.open_table(NEEDLE_TABLE).map_err(|e| { io::Error::new(io::ErrorKind::Other, format!("redb open_table: {}", e)) })?; @@ -955,12 +1026,31 @@ impl RedbNeedleMap { format!("redb insert: {}", e), )); } - if let Err(e) = txn.commit() { - self.truncate_idx_to_offset(); - return Err(io::Error::new( - io::ErrorKind::Other, - format!("redb commit: {}", e), - )); + match txn.commit() { + Ok(()) => {} + Err(redb::CommitError::TransactionPoisoned) => { + // Transaction rolled back, database still usable: + // truncate the orphan .idx row. + self.truncate_idx_to_offset(); + return Err(io::Error::new( + io::ErrorKind::Other, + "redb commit: Transaction was poisoned by a panic", + )); + } + Err(e) => { + // Non-poisoned commit error: the change may be visible + // and redb refuses further writes. Keep the .idx row + // (do NOT truncate) and reopen from .idx to repair + // redb's internal state. + let err = io::Error::new(io::ErrorKind::Other, format!("redb commit: {}", e)); + if let Err(reopen_err) = self.reopen_from_idx() { + tracing::warn!( + "redb reopen after delete commit error failed: {}", + reopen_err + ); + } + return Err(err); + } } if self.idx_file.is_some() { self.idx_file_offset += NEEDLE_MAP_ENTRY_SIZE as u64; @@ -1022,6 +1112,64 @@ impl RedbNeedleMap { } } + /// Reopen the database from `.idx` after a non-`TransactionPoisoned` + /// commit error. + /// + /// redb 4.2.0 can make a `Durability::None` commit visible *before* + /// returning `Err(CommitError::Storage(..))`. In that state the database + /// refuses further write transactions, so the map cannot be used as-is. + /// The `.idx` row that triggered the commit is **not** truncated — the + /// change may already be in the table — and the database is closed and + /// reopened from `.idx`, which repairs redb's internal state and replays + /// any tail the failed commit did not apply. + /// + /// The `.idx` writer is preserved and `idx_file_offset` is advanced to + /// the on-disk `.idx` size so the preserved row is part of the replay + /// watermark. + fn reopen_from_idx(&mut self) -> io::Result<()> { + let idx_path = self + .rdb_path + .strip_suffix(".rdb") + .map(|base| format!("{base}.idx")) + .ok_or_else(|| { + io::Error::other(format!( + "rdb path lacks .rdb suffix, cannot derive .idx path: {}", + self.rdb_path + )) + })?; + + // Close the current database first — redb holds an exclusive file + // lock, so load_from_idx cannot open the same path until we drop it. + let old_db = self.db.take(); + drop(old_db); + + let read_file = std::fs::OpenOptions::new() + .read(true) + .open(&idx_path) + .map_err(|e| { + io::Error::other(format!("reopen: open .idx {}: {}", idx_path, e)) + })?; + let actual_idx_size = read_file.metadata()?.len(); + let mut reader = io::BufReader::new(read_file); + + let reopened = Self::load_from_idx( + &self.rdb_path, + &mut reader, + self.version, + self.cache_bytes, + )?; + + // Preserve the append writer and the paths/version/cache; adopt the + // repaired database, metrics, and idx_file_offset from the reload. + self.db = reopened.db; + self.metric = reopened.metric; + self.idx_file_offset = actual_idx_size; + // The reopen replayed all rows since the last durable checkpoint + // non-durably; start the counter fresh. + self.writes_since_checkpoint = 0; + Ok(()) + } + /// Close the index file, checkpointing first so a reload starts from /// the recorded .idx size instead of replaying entries the table /// already holds. @@ -1048,7 +1196,7 @@ impl RedbNeedleMap { /// Save the redb contents to an index file, sorted by needle ID ascending. pub fn save_to_idx(&self, path: &str) -> io::Result<()> { let txn = self - .db + .db_or_err()? .begin_read() .map_err(|e| io::Error::new(io::ErrorKind::Other, format!("redb begin_read: {}", e)))?; let table = txn @@ -1088,7 +1236,8 @@ impl RedbNeedleMap { F: FnMut(NeedleId, &NeedleValue) -> Result<(), String>, { let txn = self - .db + .db_or_err() + .map_err(|e| format!("redb begin_read: {}", e))? .begin_read() .map_err(|e| format!("redb begin_read: {}", e))?; let table = txn @@ -1111,7 +1260,7 @@ impl RedbNeedleMap { pub fn collect_entries(&self) -> io::Result> { let mut result = Vec::new(); let txn: redb::ReadTransaction = self - .db + .db_or_err()? .begin_read() .map_err(|e| io::Error::other(format!("redb begin_read: {e}")))?; let table = txn @@ -1891,7 +2040,7 @@ mod tests { let db_path = dir.path().join("test.rdb"); let mut nm = RedbNeedleMap::new(db_path.to_str().unwrap(), redb_test_cache()).unwrap(); { - let txn = RedbNeedleMap::begin_write_no_fsync(&nm.db).unwrap(); + let txn = RedbNeedleMap::begin_write_no_fsync(nm.db.as_ref().unwrap()).unwrap(); { let mut table = txn.open_table(NEEDLE_TABLE).unwrap(); table.insert(1u64, &b"xx"[..]).unwrap();