Skip to main content

kernel/recover/
safe.rs

1//! Source-preserving salvage. Only verified parent links establish membership.
2//! Unreachable leaves are exported as raw evidence, never merged into current rows.
3use super::{copy_overflow, decode_leaf_record, quarantine_damaged_wal, RecoveryReport};
4use crate::{
5    budget::MemoryBudget,
6    bulk::{pack_tree, ExternalSort},
7    io::{open_file, open_recovery_source, Barrier, FileIo, IoMode},
8    meta::Meta,
9    page::{PageKind, PageRef, PAGE_SIZE},
10    pool::BufferPool,
11    store::{Config, Store, SyncMode},
12    Error, Result,
13};
14use std::{
15    fs::{self, File, OpenOptions},
16    io::{BufWriter, Write},
17    path::{Path, PathBuf},
18    sync::Arc,
19};
20
21#[derive(Debug, Clone, Copy, PartialEq, Eq)]
22pub enum RecoveryClass {
23    /// Surviving rows were reached through verified published links and WAL.
24    /// This does not assert completeness: consult losses and unknown extents.
25    VerifiedSubset,
26    /// Meta/WAL evidence is missing: even membership in the latest state is uncertain.
27    MembershipUncertain,
28}
29
30#[derive(Debug)]
31pub struct SafeRecoveryReport {
32    pub version: u32,
33    pub source: PathBuf,
34    pub destination: PathBuf,
35    pub database: PathBuf,
36    pub class: RecoveryClass,
37    pub entries_recovered: u64,
38    pub entries_before_wal: u64,
39    pub known_value_losses: u64,
40    pub unknown_extents: u64,
41    pub pages_read: u64,
42    pub raw_candidate_records: u64,
43    pub loss_journal: PathBuf,
44    /// Framed raw cells; overflow cells retain SOURCE markers, not literal values.
45    pub candidate_archive: Option<PathBuf>,
46    pub wal: RecoveryReport,
47}
48
49fn bad(no: u32, why: &'static str) -> Error {
50    Error::Corrupt { page_no: no, why }
51}
52fn fresh(path: &Path) -> Result<File> {
53    Ok(OpenOptions::new().write(true).create_new(true).open(path)?)
54}
55fn hex(w: &mut impl Write, bytes: &[u8]) -> Result<()> {
56    for b in bytes {
57        write!(w, "{b:02x}")?;
58    }
59    Ok(())
60}
61
62struct Scan<'a> {
63    file: &'a dyn FileIo,
64    pages: u64,
65    sort: ExternalSort,
66    journal: BufWriter<File>,
67    rep: SafeRecoveryReport,
68    need_candidates: bool,
69    max_lsn: u64,
70}
71impl Scan<'_> {
72    fn loss(
73        &mut self,
74        no: u32,
75        key: Option<&[u8]>,
76        lower: Option<&[u8]>,
77        upper: Option<&[u8]>,
78        why: &'static str,
79    ) -> Result<()> {
80        if key.is_some() {
81            self.rep.known_value_losses += 1;
82        } else {
83            self.rep.unknown_extents += 1;
84        }
85        write!(
86            self.journal,
87            "{}\t{no}\t",
88            if key.is_some() { "key" } else { "extent" }
89        )?;
90        if let Some(k) = key {
91            hex(&mut self.journal, k)?;
92        }
93        write!(self.journal, "\t")?;
94        if let Some(k) = lower {
95            hex(&mut self.journal, k)?;
96        }
97        write!(self.journal, "\t")?;
98        if let Some(k) = upper {
99            hex(&mut self.journal, k)?;
100        }
101        writeln!(self.journal, "\t{why}")?;
102        Ok(())
103    }
104
105    // One 4KiB page per level, depth <= 64; no visited set proportional to file.
106    // Strict parent bounds and leaf ordering make shared/cyclic children fail.
107    fn walk(
108        &mut self,
109        no: u32,
110        lower: Option<&[u8]>,
111        upper: Option<&[u8]>,
112        depth: usize,
113    ) -> Result<()> {
114        if depth >= 64 || no < 2 || no as u64 >= self.pages {
115            self.need_candidates = true;
116            return self.loss(no, None, lower, upper, "invalid child or depth");
117        }
118        let mut buf = [0u8; PAGE_SIZE];
119        self.rep.pages_read += 1;
120        // Actual I/O errors abort; corruption/truncation are named data losses.
121        self.file.read_at(&mut buf, no as u64 * PAGE_SIZE as u64)?;
122        let page = match PageRef::open(&buf, no) {
123            Ok(p) if p.tree_id() == 1 => p,
124            _ => {
125                self.need_candidates = true;
126                return self.loss(no, None, lower, upper, "page integrity or tree identity");
127            }
128        };
129        self.max_lsn = self.max_lsn.max(page.lsn());
130        match page.kind() {
131            PageKind::Leaf => {
132                // Validate the whole leaf before admitting any record. A malformed
133                // cell/order makes its extent uncertain; do not partially trust it.
134                let mut previous: Option<&[u8]> = None;
135                for i in 0..page.nentries() {
136                    let k = match decode_leaf_record(page.slot(i), no) {
137                        Ok((k, _, _)) => k,
138                        Err(_) => {
139                            self.need_candidates = true;
140                            return self.loss(no, None, lower, upper, "malformed leaf cell");
141                        }
142                    };
143                    if previous.is_some_and(|p| k <= p)
144                        || lower.is_some_and(|p| k < p)
145                        || upper.is_some_and(|p| k >= p)
146                    {
147                        self.need_candidates = true;
148                        return self.loss(no, None, lower, upper, "leaf ordering or parent bounds");
149                    }
150                    previous = Some(k);
151                }
152                for i in 0..page.nentries() {
153                    let (k, v, overflow) = decode_leaf_record(page.slot(i), no)?;
154                    if crate::keys::is_field_aggregate_key(k) {
155                        continue;
156                    }
157                    if overflow {
158                        match verify_overflow(self.file, v) {
159                            Ok(()) => (),
160                            Err(Error::Corrupt { why, .. }) => {
161                                self.loss(no, Some(k), None, None, why)?;
162                                continue;
163                            }
164                            Err(e) => return Err(e),
165                        }
166                    }
167                    self.sort.push_flagged(k.to_vec(), v.to_vec(), overflow)?;
168                    self.rep.entries_before_wal += 1;
169                }
170            }
171            PageKind::Interior => {
172                // Decode every separator before following ANY child of this page.
173                let mut separators = Vec::new();
174                for i in 0..page.nentries() {
175                    let (k, child) =
176                        match crate::verify::decode_record(page.slot(i), no, PageKind::Interior) {
177                            Ok(crate::verify::DecodedRecord::Interior { key, child }) => {
178                                (key, child)
179                            }
180                            _ => {
181                                self.need_candidates = true;
182                                return self.loss(
183                                    no,
184                                    None,
185                                    lower,
186                                    upper,
187                                    "malformed interior cell",
188                                );
189                            }
190                        };
191                    if separators
192                        .last()
193                        .is_some_and(|(p, _): &(&[u8], u32)| k <= *p)
194                        || lower.is_some_and(|p| k < p)
195                        || upper.is_some_and(|p| k >= p)
196                    {
197                        self.need_candidates = true;
198                        return self.loss(
199                            no,
200                            None,
201                            lower,
202                            upper,
203                            "interior ordering or parent bounds",
204                        );
205                    }
206                    separators.push((k, child));
207                }
208                for i in 0..=separators.len() {
209                    let child = if i == 0 {
210                        page.child0()
211                    } else {
212                        separators[i - 1].1
213                    };
214                    let lo = if i == 0 {
215                        lower
216                    } else {
217                        Some(separators[i - 1].0)
218                    };
219                    let hi = if i == separators.len() {
220                        upper
221                    } else {
222                        Some(separators[i].0)
223                    };
224                    self.walk(child, lo, hi, depth + 1)?;
225                }
226            }
227            _ => {
228                self.need_candidates = true;
229                self.loss(no, None, lower, upper, "unexpected tree page kind")?;
230            }
231        }
232        Ok(())
233    }
234
235    fn archive_candidates(&mut self, path: &Path) -> Result<()> {
236        let mut out = BufWriter::new(fresh(path)?);
237        out.write_all(b"E4RAW001")?;
238        let mut buf = [0u8; PAGE_SIZE];
239        for no in 2..self.pages {
240            self.rep.pages_read += 1;
241            self.file.read_at(&mut buf, no * PAGE_SIZE as u64)?;
242            let page = match PageRef::open(&buf, no as u32) {
243                Ok(p) if p.kind() == PageKind::Leaf && p.tree_id() == 1 => p,
244                _ => continue,
245            };
246            for i in 0..page.nentries() {
247                let cell = page.slot(i);
248                if decode_leaf_record(cell, no as u32).is_err() {
249                    continue;
250                }
251                // The raw source cell and origin remain intact; no version winner
252                // or currentness is inferred, and marker bytes are never values.
253                let mut header = Vec::with_capacity(16);
254                header.extend_from_slice(&(no as u32).to_le_bytes());
255                header.extend_from_slice(&page.lsn().to_le_bytes());
256                header.extend_from_slice(&(cell.len() as u32).to_le_bytes());
257                let crc = crc32c::crc32c_append(crc32c::crc32c(&header), cell);
258                out.write_all(&header)?;
259                out.write_all(cell)?;
260                out.write_all(&crc.to_le_bytes())?;
261                self.rep.raw_candidate_records += 1;
262            }
263        }
264        out.flush()?;
265        out.get_ref().sync_all()?;
266        self.rep.candidate_archive = Some(path.to_path_buf());
267        Ok(())
268    }
269}
270
271// Read-only preflight prevents orphan destination allocations for damaged values.
272// Whole-value CRC rejects individually valid pages crossed between chains.
273fn verify_overflow(file: &dyn FileIo, marker: &[u8]) -> Result<()> {
274    super::reader::visit_overflow(file, marker, |_| {})
275}
276
277/// Recover into a NEW directory, leaving all source files unchanged on every
278/// result. `COMPLETE` is written last; failed attempts retain evidence and must
279/// be retried with another destination. This API never publishes over source.
280///
281/// Costs: tree walk + two reads of each healthy overflow + fresh build/verification;
282/// damage to ancestry adds a sequential raw-cell archive pass. Minimum budget
283/// 8MiB. The report streams losses rather than keeping O(losses) collections.
284pub fn recover_to(source: &Path, destination: &Path, cfg: Config) -> Result<SafeRecoveryReport> {
285    recover_with_hook(source, destination, cfg, |_| Ok(()))
286}
287
288fn recover_with_hook(
289    source: &Path,
290    destination: &Path,
291    mut cfg: Config,
292    mut phase: impl FnMut(&'static str) -> Result<()>,
293) -> Result<SafeRecoveryReport> {
294    if cfg.budget_bytes < 8 << 20 {
295        return Err(Error::OutOfBudget);
296    }
297    cfg.sync = SyncMode::Full;
298    cfg.io = IoMode::Buffered;
299    let source = fs::canonicalize(source)?;
300    let parent = fs::canonicalize(destination.parent().ok_or(Error::TooLarge)?)?;
301    let destination = parent.join(destination.file_name().ok_or(Error::TooLarge)?);
302    if destination.starts_with(&source) || source.starts_with(&destination) {
303        return Err(bad(0, "recovery destination aliases or contains source"));
304    }
305    // Do not let open_file_writer's create semantics create a missing source.
306    if !fs::metadata(source.join("data"))?.is_file() {
307        return Err(bad(0, "source data is not a file"));
308    }
309    let file = open_recovery_source(&source.join("data"))?;
310    phase("locked")?;
311    fs::create_dir(&destination)?; // exclusive; never delete or reuse an attempt
312    crate::io::sync_directory(&parent)?;
313    phase("destination-created")?;
314    let database = destination.join("database");
315    fs::create_dir(&database)?;
316    let journal_path = destination.join("losses.tsv");
317    let mut journal = BufWriter::new(fresh(&journal_path)?);
318    writeln!(
319        journal,
320        "kind\tpage\tkey_hex\tlower_inclusive_hex\tupper_exclusive_hex\treason"
321    )?;
322    let rep = SafeRecoveryReport {
323        version: 1,
324        source: source.clone(),
325        destination: destination.clone(),
326        database: database.clone(),
327        class: RecoveryClass::VerifiedSubset,
328        entries_recovered: 0,
329        entries_before_wal: 0,
330        known_value_losses: 0,
331        unknown_extents: 0,
332        pages_read: 0,
333        raw_candidate_records: 0,
334        loss_journal: journal_path,
335        candidate_archive: None,
336        wal: RecoveryReport::default(),
337    };
338    let scratch = destination.join("scratch");
339    let mut scan = Scan {
340        file: &*file,
341        pages: file.len()? / PAGE_SIZE as u64,
342        sort: ExternalSort::new(&scratch, cfg.budget_bytes / 4)?,
343        journal,
344        rep,
345        need_candidates: false,
346        max_lsn: 0,
347    };
348    let mut chosen: Option<Meta> = None;
349    let mut buf = [0u8; PAGE_SIZE];
350    for no in 0..2u32 {
351        scan.rep.pages_read += 1;
352        let m = if no as u64 >= scan.pages {
353            Err(bad(no, "missing meta"))
354        } else {
355            file.read_at(&mut buf, no as u64 * PAGE_SIZE as u64)?;
356            PageRef::open(&buf, no).and_then(|p| {
357                if p.kind() != PageKind::Meta || p.tree_id() != 0 {
358                    return Err(bad(no, "meta identity"));
359                }
360                // Empty slot B is the valid unused initial slot, not lost state.
361                if no == 1 && p.nentries() == 0 {
362                    return Ok(None);
363                }
364                Meta::from_page(&p).map(Some)
365            })
366        };
367        match m {
368            Ok(Some(m)) => {
369                if m.roots[1..].iter().any(|r| *r != 0) {
370                    return Err(bad(no, "multiple trees require explicit salvage mapping"));
371                }
372                scan.max_lsn = scan.max_lsn.max(m.next_lsn);
373                if chosen
374                    .as_ref()
375                    .is_none_or(|old| m.generation > old.generation)
376                {
377                    chosen = Some(m);
378                }
379            }
380            Ok(None) => (),
381            Err(_) => {
382                scan.rep.class = RecoveryClass::MembershipUncertain;
383                scan.need_candidates = true;
384                scan.loss(no, None, None, None, "meta evidence missing")?;
385            }
386        }
387    }
388    if let Some(meta) = chosen {
389        scan.walk(meta.roots[0], None, None, 0)?;
390    } else {
391        scan.need_candidates = true;
392        scan.rep.class = RecoveryClass::MembershipUncertain;
393    }
394    if file.len()? % PAGE_SIZE as u64 != 0 {
395        scan.need_candidates = true;
396        scan.loss(scan.pages as u32, None, None, None, "truncated file tail")?;
397    }
398    if scan.need_candidates {
399        scan.archive_candidates(&destination.join("candidates.raw"))?;
400    }
401    scan.journal.flush()?;
402    scan.journal.get_ref().sync_all()?;
403
404    phase("scanned")?;
405    let mut runs = scan.sort.finish()?;
406    let data = database.join("data");
407    let (nf, _) = open_file(&data, cfg.io)?;
408    let pool = BufferPool::new(
409        nf.into(),
410        Arc::new(MemoryBudget::new(cfg.budget_bytes)),
411        cfg.budget_bytes / 2 / PAGE_SIZE,
412    )?;
413    drop(pool.allocate()?);
414    drop(pool.allocate()?);
415    Meta::init_slot_b(&pool)?;
416    let iter = runs.iter()?.map(|item| {
417        let (k, v, overflow) = item?;
418        let v = if overflow {
419            copy_overflow(&*file, &pool, &v)?
420        } else {
421            v
422        };
423        Ok((k, v, overflow))
424    });
425    let root = pack_tree(&pool, 1, iter, 0.9, &scratch)?;
426    let mut roots = [0; crate::meta::MAX_TREES];
427    roots[0] = root;
428    let meta = Meta {
429        format_version: crate::meta::FORMAT_VERSION,
430        roots,
431        generation: 0,
432        next_lsn: scan.max_lsn.checked_add(1).ok_or(Error::TooLarge)?,
433    };
434    meta.write(&pool)?;
435    Meta::mark_salvaged(&pool)?;
436    pool.flush_all(Barrier::Full)?;
437    drop(pool);
438    crate::verify::verify_rebuild(&data, cfg.io, &meta, scan.rep.entries_before_wal)?;
439
440    phase("build-verified")?;
441    let wal_source = source.join("wal");
442    match fs::metadata(&wal_source) {
443        Ok(m) => {
444            if !m.is_file() {
445                return Err(bad(0, "source WAL is not a file"));
446            }
447            let wal_dest = database.join("wal");
448            let want = crate::wal::hash_prefix(&wal_source, m.len())?;
449            fs::copy(&wal_source, &wal_dest)?;
450            File::open(&wal_dest)?.sync_all()?;
451            if fs::metadata(&wal_dest)?.len() != m.len()
452                || crate::wal::hash_prefix(&wal_dest, m.len())? != want
453            {
454                return Err(bad(0, "copied WAL failed verification"));
455            }
456        }
457        Err(e) if e.kind() == std::io::ErrorKind::NotFound => (),
458        Err(e) => return Err(e.into()),
459    }
460    quarantine_damaged_wal(&database, cfg, &mut scan.rep.wal)?;
461    if scan.rep.wal.wal_quarantined.is_some() {
462        scan.rep.class = RecoveryClass::MembershipUncertain;
463    }
464    // Replay only into the new database. An independent reopen verifies values,
465    // including overflow; counts after replay are reported separately.
466    phase("wal-preserved")?;
467    let mut store = Store::open(&database, cfg)?;
468    store.checkpoint()?;
469    drop(store);
470    let store = Store::open(&database, cfg)?;
471    scan.rep.entries_recovered =
472        crate::verify::verify_published_tree(&data, cfg.io, store.published_root(), 1)?.0;
473    drop(store);
474    let (durable, _) = open_file(&data, cfg.io)?;
475    durable.sync_dir()?;
476    phase("replay-verified")?;
477    let mut complete = fresh(&destination.join("COMPLETE"))?;
478    writeln!(complete, "e4-recovery-v1\nclass={:?}\nrows={}\nknown_value_losses={}\nunknown_extents={}\nraw_candidates={}",
479        scan.rep.class, scan.rep.entries_recovered, scan.rep.known_value_losses,
480        scan.rep.unknown_extents, scan.rep.raw_candidate_records)?;
481    complete.sync_all()?;
482    let (marker, _) = open_file(&destination.join("COMPLETE"), cfg.io)?;
483    marker.sync_dir()?;
484    Ok(scan.rep)
485}
486
487#[cfg(test)]
488mod tests {
489    use super::*;
490
491    #[test]
492    fn interrupted_repair_is_safe_to_retry() {
493        let cfg = Config {
494            budget_bytes: 8 << 20,
495            io: IoMode::Buffered,
496            sync: SyncMode::Full,
497        };
498        for stop in [
499            "locked",
500            "destination-created",
501            "scanned",
502            "build-verified",
503            "wal-preserved",
504            "replay-verified",
505        ] {
506            let d = tempfile::tempdir().unwrap();
507            let source = d.path().join("source");
508            let mut s = Store::create(&source, cfg).unwrap();
509            s.put(b"keep", &vec![9; 9000]).unwrap();
510            s.commit().unwrap();
511            s.checkpoint().unwrap();
512            s.put(b"wal", b"also keep").unwrap();
513            s.commit().unwrap();
514            drop(s);
515            let before: Vec<_> = ["data", "wal", "free"]
516                .map(|f| fs::read(source.join(f)).ok())
517                .into();
518            let dest = d.path().join("interrupted");
519            let result = recover_with_hook(&source, &dest, cfg, |phase| {
520                if phase == stop {
521                    Err(std::io::Error::other("injected repair interruption").into())
522                } else {
523                    Ok(())
524                }
525            });
526            assert!(result.is_err(), "{stop}");
527            assert!(!dest.join("COMPLETE").exists(), "{stop}");
528            let after: Vec<_> = ["data", "wal", "free"]
529                .map(|f| fs::read(source.join(f)).ok())
530                .into();
531            assert_eq!(before, after, "{stop}");
532            if dest.exists() {
533                assert!(recover_to(&source, &dest, cfg).is_err());
534            }
535            let report = recover_to(&source, &d.path().join("retry"), cfg).unwrap();
536            assert_eq!(report.entries_recovered, 2, "{stop}");
537            let s = Store::open(&report.database, cfg).unwrap();
538            assert_eq!(s.get(b"keep").unwrap(), Some(vec![9; 9000]));
539            assert_eq!(s.get(b"wal").unwrap(), Some(b"also keep".to_vec()));
540        }
541    }
542}