1use 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 VerifiedSubset,
26 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 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 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 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 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 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 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
271fn verify_overflow(file: &dyn FileIo, marker: &[u8]) -> Result<()> {
274 super::reader::visit_overflow(file, marker, |_| {})
275}
276
277pub 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 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)?; 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 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 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}