1use crate::io::{open_file, FileIo, IoMode};
80use crate::{Error, Result};
81use std::path::Path;
82
83const HDR: usize = 20;
84
85const MAX_FRAME_BYTES: u64 = u32::MAX as u64;
94pub const MAX_PAYLOAD_BYTES: u64 = MAX_FRAME_BYTES - HDR as u64;
95const CRC_CHUNK: usize = 64 * 1024;
96
97#[derive(Debug, Clone, Copy, PartialEq, Eq)]
118pub enum Stop {
119 End(&'static str),
123 Damaged { offset: u64, why: &'static str },
127}
128
129#[derive(Debug, Clone, Copy)]
132pub struct Scan {
133 pub end: u64,
134 pub next_lsn: u64,
135 pub stop: Stop,
136}
137
138#[derive(Debug, Clone, Copy, PartialEq, Eq)]
139#[repr(u8)]
140pub enum RecKind {
141 Commit = 1,
142 PageImage = 2,
143 Put = 3,
144 Delete = 4,
145 DeletePrefix = 5,
146 PutEmptyBatch = 6,
147}
148
149impl RecKind {
150 fn from_u8(v: u8) -> Option<Self> {
151 Some(match v {
152 1 => Self::Commit,
153 2 => Self::PageImage,
154 3 => Self::Put,
155 4 => Self::Delete,
156 5 => Self::DeletePrefix,
157 6 => Self::PutEmptyBatch,
158 _ => return None,
159 })
160 }
161}
162
163const WAL_BUF: usize = 256 * 1024;
171
172pub struct Wal {
173 file: Box<dyn FileIo>,
174 end: u64,
175 next_lsn: u64,
176 buf: Vec<u8>,
177 flushed: u64,
178 limit: Option<u64>,
179}
180
181pub(crate) struct SalvagedWal {
182 pub bytes: u64,
183 pub hash: u32,
184}
185
186enum Candidate {
187 Valid { next: u64, kind: RecKind },
188 Invalid,
189}
190
191impl Wal {
192 pub fn open(path: &Path, mode: IoMode) -> Result<Wal> {
193 Self::open_limited(path, mode, None)
194 }
195 pub(crate) fn open_limited(path: &Path, mode: IoMode, limit: Option<u64>) -> Result<Wal> {
196 if let Some(cap) = limit {
197 match std::fs::metadata(path) {
198 Ok(m) if m.len() > cap => return Err(Error::ResourceLimit("existing WAL exceeds allowance")),
199 Ok(_) => (), Err(e) if e.kind() == std::io::ErrorKind::NotFound => (), Err(e) => return Err(e.into()),
200 }
201 }
202 let _ = mode;
205 let (file, _) = open_file(path, IoMode::Buffered)?;
206 let mut wal = Self::open_on(file)?;
207 wal.limit = limit;
208 Ok(wal)
209 }
210
211 fn open_on(file: Box<dyn FileIo>) -> Result<Wal> {
217 let scan = Self::scan(&*file)?;
218 match scan.stop {
219 Stop::End(_) => {
222 file.set_len(scan.end)?;
229 file.sync_data()?;
230 }
231 Stop::Damaged { offset, why } => return Err(Error::CorruptWal { offset, why }),
232 }
233 Ok(Wal { file, end: scan.end, next_lsn: scan.next_lsn,
234 buf: Vec::with_capacity(WAL_BUF), flushed: scan.end, limit: None })
235 }
236
237 pub(crate) fn cut_to_committed(&mut self, end: u64) -> Result<()> {
242 if end < self.end {
243 self.buf.clear();
244 self.file.set_len(end)?;
245 self.file.sync_data()?;
246 self.end = end;
247 self.flushed = end;
248 }
249 Ok(())
250 }
251
252 pub fn inspect(path: &Path, mode: IoMode) -> Result<Scan> {
258 let _ = mode;
259 let (file, _) = open_file(path, IoMode::Buffered)?;
260 Self::scan(&*file)
261 }
262
263 pub fn end_offset(&self) -> u64 { self.end }
264
265 fn frame_crc(hdr: &[u8; HDR], payload: &[u8]) -> u32 {
266 let mut h = [0u8; HDR];
267 h.copy_from_slice(hdr);
268 h[16..20].fill(0);
269 crc32c::crc32c_append(crc32c::crc32c(&h), payload)
270 }
271
272 fn frame_crc_from_file(
273 file: &dyn FileIo,
274 hdr: &[u8; HDR],
275 payload_at: u64,
276 payload_len: u64,
277 scratch: &mut [u8],
278 ) -> Result<u32> {
279 let mut h = *hdr;
280 h[16..20].fill(0);
281 let mut crc = crc32c::crc32c(&h);
282 let mut at = payload_at;
283 let end = payload_at.checked_add(payload_len).ok_or(Error::CorruptWal {
284 offset: payload_at,
285 why: "payload boundary overflows the WAL address space",
286 })?;
287 while at < end {
288 let n = std::cmp::min(scratch.len() as u64, end - at) as usize;
289 read_exact(file, &mut scratch[..n], at)?;
290 crc = crc32c::crc32c_append(crc, &scratch[..n]);
291 at += n as u64;
292 }
293 Ok(crc)
294 }
295
296 fn scan(file: &dyn FileIo) -> Result<Scan> {
309 let len = file.len()?;
310 let mut off = 0u64;
311 let mut next_lsn = 1u64;
312 let mut crc_scratch = vec![0u8; CRC_CHUNK];
313 let stop = loop {
314 let header_end = off.checked_add(HDR as u64).ok_or(Error::CorruptWal {
315 offset: off,
316 why: "header boundary overflows the WAL address space",
317 })?;
318 if header_end > len {
319 if off == len {
320 break Stop::End("the last verified frame reaches physical EOF");
321 }
322 if Self::is_all_zeros(file, off, len)? {
323 break Stop::End("nothing past the last good frame but zeros");
324 }
325 break Stop::Damaged {
326 offset: off,
327 why: "non-zero bytes remain but there is not a complete frame header",
328 };
329 }
330 let mut hdr = [0u8; HDR];
331 read_exact(file, &mut hdr, off)?;
332 let plen = u32::from_le_bytes(hdr[0..4].try_into().unwrap()) as u64;
333 let lsn = u64::from_le_bytes(hdr[4..12].try_into().unwrap());
334 if plen > MAX_PAYLOAD_BYTES {
335 break Stop::Damaged {
336 offset: off,
337 why: "frame payload length exceeds the writer maximum",
338 };
339 }
340 let frame_end = header_end.checked_add(plen).ok_or(Error::CorruptWal {
341 offset: off,
342 why: "frame boundary overflows the WAL address space",
343 })?;
344 if frame_end > len {
345 if !Self::has_writer_header_shape(&hdr) {
350 break Stop::Damaged {
351 offset: off,
352 why: "incomplete-looking frame header was not emitted by this writer",
353 };
354 }
355 if Self::has_verified_frame_after(
360 file,
361 off.checked_add(1).ok_or(Error::CorruptWal {
362 offset: off,
363 why: "WAL resynchronisation offset overflow",
364 })?,
365 len,
366 &mut crc_scratch,
367 )? {
368 break Stop::Damaged {
369 offset: off,
370 why: "incomplete-looking frame has verified log frames behind it",
371 };
372 }
373 break Stop::End("a bounded frame header crosses physical EOF");
374 }
375 let want = u32::from_le_bytes(hdr[16..20].try_into().unwrap());
376 if Self::frame_crc_from_file(file, &hdr, header_end, plen, &mut crc_scratch)? != want {
377 if Self::is_all_zeros(file, off, len)? {
378 break Stop::End("nothing past the last good frame but zeros");
379 }
380 break Stop::Damaged {
381 offset: off,
382 why: "complete frame fails its checksum",
383 };
384 }
385 if RecKind::from_u8(hdr[12]).is_none() {
387 break Stop::Damaged {
388 offset: off,
389 why: "frame verifies but names a record kind this build does not know",
390 };
391 }
392 next_lsn = lsn.checked_add(1).ok_or(Error::CorruptWal {
393 offset: off,
394 why: "verified frame exhausts the log sequence number space",
395 })?;
396 off = frame_end;
397 };
398 Ok(Scan { end: off, next_lsn, stop })
399 }
400
401 fn is_all_zeros(file: &dyn FileIo, off: u64, len: u64) -> Result<bool> {
407 const CHUNK: usize = 64 * 1024;
408 let mut buf = vec![0u8; CHUNK];
409 let mut at = off;
410 while at < len {
411 let n = std::cmp::min(CHUNK as u64, len - at) as usize;
412 read_exact(file, &mut buf[..n], at)?;
413 if buf[..n].iter().any(|&b| b != 0) { return Ok(false); }
414 at += n as u64;
415 }
416 Ok(true)
417 }
418
419 pub fn append(&mut self, kind: RecKind, payload: &[u8]) -> Result<u64> {
420 if payload.len() as u64 > MAX_PAYLOAD_BYTES {
421 return Err(Error::TooLarge);
422 }
423 let lsn = self.next_lsn;
424 let next_lsn = self.next_lsn.checked_add(1).ok_or(Error::CorruptWal {
425 offset: self.end,
426 why: "log sequence number space is exhausted",
427 })?;
428 let frame_len = (HDR as u64).checked_add(payload.len() as u64).ok_or(Error::TooLarge)?;
429 if frame_len > MAX_FRAME_BYTES { return Err(Error::TooLarge); }
430 let new_end = self.end.checked_add(frame_len).ok_or(Error::TooLarge)?;
431 if self.limit.is_some_and(|cap| new_end > cap) {
432 return Err(Error::ResourceLimit("WAL allowance full; reduce transaction size"));
433 }
434 let mut hdr = [0u8; HDR];
435 hdr[0..4].copy_from_slice(&(payload.len() as u32).to_le_bytes());
436 hdr[4..12].copy_from_slice(&lsn.to_le_bytes());
437 hdr[12] = kind as u8;
438 #[cfg(feature = "write-trace")]
439 let crc_started = crate::write_trace::active().then(std::time::Instant::now);
440 let c = Self::frame_crc(&hdr, payload);
441 #[cfg(feature = "write-trace")]
442 if let Some(started) = crc_started {
443 crate::write_trace::add(crate::write_trace::Field::WalCrc, started.elapsed());
444 }
445 hdr[16..20].copy_from_slice(&c.to_le_bytes());
446
447 #[cfg(feature = "write-trace")]
452 let copy_started = crate::write_trace::active().then(std::time::Instant::now);
453 self.buf.extend_from_slice(&hdr);
454 self.buf.extend_from_slice(payload);
455 #[cfg(feature = "write-trace")]
456 if let Some(started) = copy_started {
457 crate::write_trace::add(crate::write_trace::Field::WalBufferCopy, started.elapsed());
458 crate::write_trace::value_copy();
459 }
460 self.next_lsn = next_lsn;
461 self.end = new_end;
462 if self.buf.len() >= WAL_BUF {
463 #[cfg(feature = "write-trace")]
464 let flush_started = crate::write_trace::active().then(std::time::Instant::now);
465 self.flush()?;
466 #[cfg(feature = "write-trace")]
467 if let Some(started) = flush_started {
468 crate::write_trace::add(crate::write_trace::Field::WalFlush, started.elapsed());
469 }
470 }
471 Ok(lsn)
472 }
473
474 pub fn flush(&mut self) -> Result<()> {
479 if self.buf.is_empty() { return Ok(()); }
480 write_all(&*self.file, &self.buf, self.flushed)?;
481 crate::write_stats::add(crate::write_stats::Phase::Wal, self.buf.len() as u64);
482 self.flushed += self.buf.len() as u64;
483 self.buf.clear();
484 debug_assert_eq!(self.flushed, self.end);
485 Ok(())
486 }
487
488 pub fn sync_data(&mut self) -> Result<()> { self.flush()?; self.file.sync_data() }
489
490 pub fn sync_full(&mut self) -> Result<()> { self.flush()?; self.file.sync_full() }
492
493 pub fn sync_full_primitive(&self) -> &'static str { self.file.sync_full_primitive() }
496
497 pub(crate) fn committed_end(&self) -> Result<u64> {
501 assert!(self.buf.is_empty(), "recovery with {} bytes buffered", self.buf.len());
505 let mut off = 0u64;
506 let mut committed = 0u64;
507 while off < self.end {
508 let mut hdr = [0u8; HDR];
509 read_exact(&*self.file, &mut hdr, off)?;
510 let plen = u32::from_le_bytes(hdr[0..4].try_into().unwrap()) as u64;
511 if plen > MAX_PAYLOAD_BYTES {
512 return Err(Error::CorruptWal { offset: off, why: "recovery payload exceeds the writer maximum" });
513 }
514 let next = off.checked_add(HDR as u64).and_then(|n| n.checked_add(plen))
515 .ok_or(Error::CorruptWal { offset: off, why: "frame boundary overflow during recovery" })?;
516 if next > self.end {
517 return Err(Error::CorruptWal { offset: off, why: "frame crosses verified WAL boundary during recovery" });
518 }
519 let kind = RecKind::from_u8(hdr[12]).ok_or(Error::CorruptWal {
520 offset: off, why: "unknown record kind inside verified WAL",
521 })?;
522 if kind == RecKind::Commit { committed = next; }
523 off = next;
524 }
525 Ok(committed)
526 }
527
528 pub(crate) fn record_at(&self, off: u64, committed_end: u64)
531 -> Result<Option<(u64, RecKind, Vec<u8>, u64)>>
532 {
533 if off == committed_end { return Ok(None); }
534 if off > committed_end || committed_end > self.end {
535 return Err(Error::CorruptWal { offset: off, why: "recovery cursor outside committed WAL boundary" });
536 }
537 let header_end = off.checked_add(HDR as u64)
538 .ok_or(Error::CorruptWal { offset: off, why: "header boundary overflow during recovery" })?;
539 if header_end > committed_end {
540 return Err(Error::CorruptWal { offset: off, why: "truncated header inside committed WAL boundary" });
541 }
542 let mut hdr = [0u8; HDR];
543 read_exact(&*self.file, &mut hdr, off)?;
544 let plen = u32::from_le_bytes(hdr[0..4].try_into().unwrap()) as usize;
545 if plen as u64 > MAX_PAYLOAD_BYTES {
546 return Err(Error::CorruptWal { offset: off, why: "recovery payload exceeds the writer maximum" });
547 }
548 let next = header_end.checked_add(plen as u64)
549 .ok_or(Error::CorruptWal { offset: off, why: "payload boundary overflow during recovery" })?;
550 if next > committed_end {
551 return Err(Error::CorruptWal { offset: off, why: "frame crosses committed WAL boundary" });
552 }
553 let mut payload = vec![0u8; plen];
554 read_exact(&*self.file, &mut payload, header_end)?;
555 let want = u32::from_le_bytes(hdr[16..20].try_into().unwrap());
556 if Self::frame_crc(&hdr, &payload) != want {
557 return Err(Error::CorruptWal { offset: off, why: "recovery frame changed after the opening CRC walk" });
558 }
559 let lsn = u64::from_le_bytes(hdr[4..12].try_into().unwrap());
561 let kind = RecKind::from_u8(hdr[12]).ok_or(Error::CorruptWal {
562 offset: off, why: "unknown record kind inside committed WAL",
563 })?;
564 Ok(Some((lsn, kind, payload, next)))
565 }
566
567 #[cfg(test)]
568 fn replay(&self) -> Result<Vec<(u64, RecKind, Vec<u8>)>> {
569 let end = self.end;
572 let mut out = Vec::new();
573 let mut off = 0;
574 while let Some((lsn, kind, payload, next)) = self.record_at(off, end)? {
575 out.push((lsn, kind, payload));
576 off = next;
577 }
578 Ok(out)
579 }
580
581 pub fn set_lsn_floor(&mut self, floor: u64) {
586 if floor > self.next_lsn { self.next_lsn = floor; }
587 }
588
589 pub fn next_lsn(&self) -> u64 { self.next_lsn }
590
591 pub(crate) fn salvage_committed(src: &Path, dst: &Path) -> Result<SalvagedWal> {
599 use std::io::{Seek, SeekFrom};
600
601 let (file, _) = open_file(src, IoMode::Buffered)?;
602 let len = file.len()?;
603 let mut out = std::fs::File::create(dst)?;
604 let mut src_off = 0u64;
605 let mut out_end = 0u64;
606 let mut last_commit = 0u64;
607 let mut resynchronising = false;
608 let mut scratch = vec![0u8; CRC_CHUNK];
609
610 while src_off < len {
611 match Self::candidate_at(&*file, src_off, len, &mut scratch)? {
612 Candidate::Valid { next, kind } => {
613 if resynchronising {
614 src_off = next;
615 if kind == RecKind::Commit {
616 resynchronising = false;
617 }
618 continue;
619 }
620 copy_range(&*file, &mut out, src_off, next - src_off, &mut scratch)?;
621 out_end = out_end.checked_add(next - src_off)
622 .ok_or(Error::CorruptWal {
623 offset: src_off,
624 why: "salvaged WAL output length overflow",
625 })?;
626 if kind == RecKind::Commit { last_commit = out_end; }
627 src_off = next;
628 }
629 Candidate::Invalid => {
630 out.set_len(last_commit)?;
631 out.seek(SeekFrom::Start(last_commit))?;
632 out_end = last_commit;
633 resynchronising = true;
634 src_off = src_off.checked_add(1).ok_or(Error::CorruptWal {
635 offset: src_off,
636 why: "WAL resynchronisation offset overflow",
637 })?;
638 }
639 }
640 }
641 out.set_len(last_commit)?;
642 out.sync_all()?;
643 drop(out);
644
645 let hash = hash_prefix(dst, last_commit)?;
646 Ok(SalvagedWal { bytes: last_commit, hash })
647 }
648
649 fn candidate_at(
650 file: &dyn FileIo,
651 off: u64,
652 len: u64,
653 scratch: &mut [u8],
654 ) -> Result<Candidate> {
655 let Some(header_end) = off.checked_add(HDR as u64) else { return Ok(Candidate::Invalid) };
656 if header_end > len { return Ok(Candidate::Invalid); }
657 let mut hdr = [0u8; HDR];
658 read_exact(file, &mut hdr, off)?;
659 if !Self::has_writer_header_shape(&hdr) {
663 return Ok(Candidate::Invalid);
664 }
665 let plen = u32::from_le_bytes(hdr[0..4].try_into().unwrap()) as u64;
666 if plen > MAX_PAYLOAD_BYTES { return Ok(Candidate::Invalid); }
667 let Some(next) = header_end.checked_add(plen) else { return Ok(Candidate::Invalid) };
668 if next > len { return Ok(Candidate::Invalid); }
669 let want = u32::from_le_bytes(hdr[16..20].try_into().unwrap());
670 if Self::frame_crc_from_file(file, &hdr, header_end, plen, scratch)? != want {
671 return Ok(Candidate::Invalid);
672 }
673 let kind = RecKind::from_u8(hdr[12]).expect("kind was prefiltered above");
674 let lsn = u64::from_le_bytes(hdr[4..12].try_into().unwrap());
675 if lsn.checked_add(1).is_none() { return Ok(Candidate::Invalid); }
676 Ok(Candidate::Valid { next, kind })
677 }
678
679 fn has_writer_header_shape(hdr: &[u8; HDR]) -> bool {
680 hdr[13..16] == [0, 0, 0]
681 && RecKind::from_u8(hdr[12]).is_some()
682 && u64::from_le_bytes(hdr[4..12].try_into().unwrap()).checked_add(1).is_some()
683 }
684
685 fn has_verified_frame_after(
686 file: &dyn FileIo,
687 mut off: u64,
688 len: u64,
689 scratch: &mut [u8],
690 ) -> Result<bool> {
691 while off < len {
692 if matches!(Self::candidate_at(file, off, len, scratch)?, Candidate::Valid { .. }) {
693 return Ok(true);
694 }
695 off = off.checked_add(1).ok_or(Error::CorruptWal {
696 offset: off,
697 why: "WAL resynchronisation offset overflow",
698 })?;
699 }
700 Ok(false)
701 }
702
703 pub fn rotate(&mut self) -> Result<()> {
706 self.buf.clear();
708 self.flushed = 0;
709 self.file.set_len(0)?;
710 self.file.sync_data()?;
711 self.file.sync_dir()?;
712 self.end = 0;
713 Ok(())
714 }
715
716 pub(crate) fn rotate_published(&mut self) -> Result<()> {
720 self.buf.clear();
721 self.flushed = 0;
722 self.file.set_len(0)?;
723 self.file.sync_data()?;
724 self.end = 0;
725 Ok(())
726 }
727}
728
729fn read_exact(f: &dyn FileIo, buf: &mut [u8], off: u64) -> Result<()> { f.read_at(buf, off) }
732fn write_all(f: &dyn FileIo, buf: &[u8], off: u64) -> Result<()> { f.write_at(buf, off) }
733
734fn copy_range(
735 src: &dyn FileIo,
736 dst: &mut std::fs::File,
737 mut at: u64,
738 mut n: u64,
739 scratch: &mut [u8],
740) -> Result<()> {
741 use std::io::Write;
742 while n > 0 {
743 let take = std::cmp::min(n, scratch.len() as u64) as usize;
744 read_exact(src, &mut scratch[..take], at)?;
745 dst.write_all(&scratch[..take])?;
746 at += take as u64;
747 n -= take as u64;
748 }
749 Ok(())
750}
751
752pub(crate) fn hash_prefix(path: &Path, n: u64) -> Result<u32> {
753 use std::io::Read;
754 const CHUNK: usize = 64 * 1024;
755 let mut file = std::fs::File::open(path)?;
756 let mut buf = vec![0u8; CHUNK];
757 let mut left = n;
758 let mut hash = 0u32;
759 let mut first = true;
760 while left > 0 {
761 let want = std::cmp::min(left, CHUNK as u64) as usize;
762 file.read_exact(&mut buf[..want])?;
763 hash = if first {
764 first = false;
765 crc32c::crc32c(&buf[..want])
766 } else {
767 crc32c::crc32c_append(hash, &buf[..want])
768 };
769 left -= want as u64;
770 }
771 Ok(hash)
772}
773
774#[cfg(test)]
775mod tests {
776 use super::*;
777
778 fn wal(dir: &std::path::Path) -> Wal { Wal::open(&dir.join("wal"), crate::io::IoMode::Buffered).unwrap() }
779
780 #[test]
781 fn records_replay_in_order() {
782 let d = tempfile::tempdir().unwrap();
783 let mut w = wal(d.path());
784 w.append(RecKind::Put, b"one").unwrap();
785 w.append(RecKind::Put, b"two").unwrap();
786 w.append(RecKind::Commit, b"").unwrap();
787 w.sync_data().unwrap();
788 drop(w);
789
790 let got = wal(d.path()).replay().unwrap();
791 assert_eq!(got.len(), 3);
792 assert_eq!(got[0].2, b"one");
793 assert_eq!(got[1].2, b"two");
794 assert_eq!(got[2].1, RecKind::Commit);
795 }
796
797 #[test]
798 fn a_torn_tail_stops_replay_at_the_last_good_frame() {
799 let d = tempfile::tempdir().unwrap();
800 let mut w = wal(d.path());
801 w.append(RecKind::Put, b"good").unwrap();
802 w.sync_data().unwrap();
803 let n = w.end_offset();
804 drop(w);
805
806 let mut f = std::fs::OpenOptions::new().write(true).open(d.path().join("wal")).unwrap();
808 use std::io::{Seek, SeekFrom, Write};
809 f.seek(SeekFrom::Start(n)).unwrap();
810 f.write_all(&[0u8; 12]).unwrap(); f.sync_all().unwrap();
812
813 let got = wal(d.path()).replay().unwrap();
814 assert_eq!(got.len(), 1, "only the intact prefix replays");
815 }
816
817 #[test]
821 fn a_write_after_a_torn_tail_survives_the_next_open() {
822 let d = tempfile::tempdir().unwrap();
823 let mut w = wal(d.path());
824 w.append(RecKind::Put, b"before").unwrap();
825 w.sync_data().unwrap();
826 let n = w.end_offset();
827 drop(w);
828
829 let mut f = std::fs::OpenOptions::new().write(true).open(d.path().join("wal")).unwrap();
830 use std::io::{Seek, SeekFrom, Write};
831 f.seek(SeekFrom::Start(n)).unwrap();
832 let mut hdr = [0u8; HDR];
833 hdr[0..4].copy_from_slice(&100u32.to_le_bytes());
834 hdr[12] = RecKind::Put as u8;
835 f.write_all(&hdr).unwrap();
836 f.write_all(&[0xAA; 9]).unwrap();
837 f.sync_all().unwrap();
838
839 let mut w2 = wal(d.path());
841 w2.append(RecKind::Put, b"after").unwrap();
842 w2.sync_data().unwrap();
843 drop(w2);
844
845 let got = wal(d.path()).replay().unwrap();
846 let payloads: Vec<&[u8]> = got.iter().map(|r| r.2.as_slice()).collect();
847 assert_eq!(payloads, vec![b"before".as_ref(), b"after".as_ref()]);
848 }
849
850 #[test]
862 fn an_all_zero_extension_is_a_clean_ending_and_is_truncated() {
863 let d = tempfile::tempdir().unwrap();
864 let mut w = wal(d.path());
865 w.append(RecKind::Put, b"one").unwrap();
866 w.sync_data().unwrap();
867 let end = w.end_offset();
868 drop(w);
869
870 let f = std::fs::OpenOptions::new().write(true)
871 .open(d.path().join("wal")).unwrap();
872 f.set_len(end + 4096).unwrap(); f.sync_all().unwrap();
874 drop(f);
875
876 let w2 = wal(d.path());
877 assert_eq!(w2.end_offset(), end);
878 assert_eq!(
879 std::fs::metadata(d.path().join("wal")).unwrap().len(), end,
880 "open must not leave orphaned bytes on disk"
881 );
882 }
883
884 #[test]
887 fn append_lands_at_the_scanned_end_not_at_physical_eof() {
888 let d = tempfile::tempdir().unwrap();
889 let mut w = wal(d.path());
890 w.append(RecKind::Put, b"one").unwrap();
891 w.sync_data().unwrap();
892 let end = w.end_offset();
893
894 {
895 let f = std::fs::OpenOptions::new().write(true)
896 .open(d.path().join("wal")).unwrap();
897 f.set_len(end + 4096).unwrap();
898 f.sync_all().unwrap();
899 }
900
901 w.append(RecKind::Put, b"two").unwrap();
902 w.sync_data().unwrap();
903
904 let raw = std::fs::read(d.path().join("wal")).unwrap();
910 let hdr = &raw[end as usize..end as usize + HDR];
911 assert_eq!(
912 u32::from_le_bytes(hdr[0..4].try_into().unwrap()), 3,
913 "the second frame must physically start at the scanned end"
914 );
915 assert_eq!(hdr[12], RecKind::Put as u8);
916
917 let got = wal(d.path()).replay().unwrap();
918 let payloads: Vec<&[u8]> = got.iter().map(|r| r.2.as_slice()).collect();
919 assert_eq!(payloads, vec![b"one".as_ref(), b"two".as_ref()]);
920 }
921
922 #[test]
923 fn non_zero_garbage_shorter_than_a_frame_is_damage_and_is_preserved() {
924 for len in [1usize, 7, 19, 20, 31, 100] {
928 let d = tempfile::tempdir().unwrap();
929 let stub: Vec<u8> = (0..len).map(|i| (i + 1) as u8).collect();
930 std::fs::write(d.path().join("wal"), &stub).unwrap();
931 assert!(matches!(
932 Wal::open(&d.path().join("wal"), IoMode::Buffered),
933 Err(crate::Error::CorruptWal { offset: 0, .. })
934 ));
935 assert_eq!(std::fs::read(d.path().join("wal")).unwrap(), stub,
936 "a refused {len}-byte tail must remain byte-for-byte intact");
937 }
938 }
939
940 #[test]
941 fn a_log_payload_larger_than_the_writer_max_is_refused_before_allocation() {
942 use std::sync::atomic::{AtomicBool, Ordering};
943 use std::sync::Arc;
944
945 struct HeaderOnly {
946 hdr: [u8; HDR],
947 len: u64,
948 payload_read: Arc<AtomicBool>,
949 }
950 impl FileIo for HeaderOnly {
951 fn requires_alignment(&self) -> bool { false }
952 fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> {
953 if off == 0 && buf.len() == HDR {
954 buf.copy_from_slice(&self.hdr);
955 return Ok(());
956 }
957 self.payload_read.store(true, Ordering::Relaxed);
958 Err(std::io::Error::other("payload read must not happen").into())
959 }
960 fn write_at(&self, _: &[u8], _: u64) -> Result<()> { Ok(()) }
961 fn sync_data(&self) -> Result<()> { Ok(()) }
962 fn sync_full(&self) -> Result<()> { Ok(()) }
963 fn sync_full_primitive(&self) -> &'static str { "test" }
964 fn sync_dir(&self) -> Result<()> { Ok(()) }
965 fn len(&self) -> Result<u64> { Ok(self.len) }
966 fn set_len(&self, _: u64) -> Result<()> { Ok(()) }
967 }
968
969 let plen = u32::MAX;
970 let mut hdr = [0u8; HDR];
971 hdr[0..4].copy_from_slice(&plen.to_le_bytes());
972 hdr[12] = RecKind::Put as u8;
973 let payload_read = Arc::new(AtomicBool::new(false));
974 let f = HeaderOnly {
975 hdr,
976 len: HDR as u64 + plen as u64,
977 payload_read: payload_read.clone(),
978 };
979 assert!(matches!(Wal::open_on(Box::new(f)), Err(crate::Error::CorruptWal { offset: 0, .. })));
980 assert!(!payload_read.load(Ordering::Relaxed),
981 "an off-disk length above the writer maximum must be rejected before allocation/read");
982 }
983
984 #[test]
985 fn an_exhausted_lsn_is_refused_with_checked_arithmetic() {
986 let d = tempfile::tempdir().unwrap();
987 let mut hdr = [0u8; HDR];
988 hdr[4..12].copy_from_slice(&u64::MAX.to_le_bytes());
989 hdr[12] = RecKind::Commit as u8;
990 let crc = Wal::frame_crc(&hdr, &[]);
991 hdr[16..20].copy_from_slice(&crc.to_le_bytes());
992 std::fs::write(d.path().join("wal"), hdr).unwrap();
993 assert!(matches!(
994 Wal::open(&d.path().join("wal"), IoMode::Buffered),
995 Err(crate::Error::CorruptWal { offset: 0, .. })
996 ));
997 }
998
999 struct FailReadAt { inner: Box<dyn FileIo>, fail_at: u64 }
1008 impl FileIo for FailReadAt {
1009 fn requires_alignment(&self) -> bool { self.inner.requires_alignment() }
1010 fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> {
1011 if off == self.fail_at {
1012 return Err(std::io::Error::other("injected read failure").into());
1013 }
1014 self.inner.read_at(buf, off)
1015 }
1016 fn write_at(&self, buf: &[u8], off: u64) -> Result<()> { self.inner.write_at(buf, off) }
1017 fn sync_data(&self) -> Result<()> { self.inner.sync_data() }
1018 fn sync_full(&self) -> Result<()> { self.inner.sync_full() }
1019 fn sync_full_primitive(&self) -> &'static str { self.inner.sync_full_primitive() }
1020 fn sync_dir(&self) -> Result<()> { self.inner.sync_dir() }
1021 fn len(&self) -> Result<u64> { self.inner.len() }
1022 fn set_len(&self, n: u64) -> Result<()> { self.inner.set_len(n) }
1023 }
1024
1025 #[test]
1031 fn an_io_error_on_a_header_read_is_an_error_not_a_clean_end() {
1032 let d = tempfile::tempdir().unwrap();
1033 let mut w = wal(d.path());
1034 let p = vec![b'A'; 100];
1035 for _ in 0..10 { w.append(RecKind::Put, &p).unwrap(); }
1036 w.sync_data().unwrap();
1037 let full = w.end_offset();
1038 drop(w);
1039 let frame = HDR as u64 + p.len() as u64;
1040 let midpoint = frame * 5;
1041
1042 let (real, _) = open_file(&d.path().join("wal"), IoMode::Buffered).unwrap();
1043 let failing = Box::new(FailReadAt { inner: real, fail_at: midpoint });
1044 match Wal::open_on(failing) {
1045 Err(crate::Error::Io(_)) => {}
1046 Err(e) => panic!("an unreadable header must surface as an I/O error, got {e:?}"),
1047 Ok(_) => panic!("an unreadable header must not be reported as a clean end of log"),
1048 }
1049 assert_eq!(
1050 std::fs::metadata(d.path().join("wal")).unwrap().len(), full,
1051 "a read that FAILED must not be able to shorten the log -- {full} bytes were on \
1052 disk and the reader could not read one header of them"
1053 );
1054 }
1055
1056 #[test]
1063 fn damage_with_more_log_behind_it_refuses_the_open_and_keeps_every_byte() {
1064 let d = tempfile::tempdir().unwrap();
1065 let mut w = wal(d.path());
1066 let p = vec![b'A'; 100];
1067 for _ in 0..100 { w.append(RecKind::Put, &p).unwrap(); }
1068 w.sync_data().unwrap();
1069 drop(w);
1070
1071 let path = d.path().join("wal");
1072 let frame = HDR + p.len();
1073 let victim = frame * 50; let mut bytes = std::fs::read(&path).unwrap();
1075 bytes[victim + 4] ^= 0x01; std::fs::write(&path, &bytes).unwrap();
1077
1078 match Wal::open(&path, IoMode::Buffered) {
1079 Err(crate::Error::CorruptWal { offset, .. }) => {
1080 assert_eq!(offset as usize, victim, "the refusal must name where it stopped");
1081 }
1082 Err(e) => panic!("expected CorruptWal, got {e:?}"),
1083 Ok(_) => panic!("damage with committed frames behind it must not open"),
1084 }
1085 assert_eq!(
1086 std::fs::read(&path).unwrap(), bytes,
1087 "a refused open must leave the log byte for byte as it found it"
1088 );
1089 }
1090
1091 #[test]
1098 fn a_verified_frame_naming_an_unknown_kind_is_damage() {
1099 let d = tempfile::tempdir().unwrap();
1100 let mut w = wal(d.path());
1101 w.append(RecKind::Put, b"good").unwrap();
1102 w.sync_data().unwrap();
1103 let end = w.end_offset();
1104 drop(w);
1105
1106 let mut hdr = [0u8; HDR];
1108 let payload = b"future".to_vec();
1109 hdr[0..4].copy_from_slice(&(payload.len() as u32).to_le_bytes());
1110 hdr[4..12].copy_from_slice(&7u64.to_le_bytes());
1111 hdr[12] = 9;
1112 let c = Wal::frame_crc(&hdr, &payload);
1113 hdr[16..20].copy_from_slice(&c.to_le_bytes());
1114 let mut bytes = std::fs::read(d.path().join("wal")).unwrap();
1115 bytes.extend_from_slice(&hdr);
1116 bytes.extend_from_slice(&payload);
1117 std::fs::write(d.path().join("wal"), &bytes).unwrap();
1118
1119 match Wal::open(&d.path().join("wal"), IoMode::Buffered) {
1120 Err(crate::Error::CorruptWal { offset, .. }) => assert_eq!(offset, end),
1121 Err(e) => panic!("expected CorruptWal, got {e:?}"),
1122 Ok(_) => panic!("a verified frame naming an unknown kind must not be walked past"),
1123 }
1124 assert_eq!(std::fs::read(d.path().join("wal")).unwrap(), bytes,
1125 "and it must not have been truncated on the way out");
1126 }
1127
1128 #[test]
1131 fn replay_refuses_a_byte_changed_after_the_opening_crc_walk() {
1132 let d = tempfile::tempdir().unwrap();
1133 let mut w = wal(d.path());
1134 w.append(RecKind::Put, b"\x01\0key-value").unwrap();
1135 w.append(RecKind::Commit, b"").unwrap();
1136 w.sync_data().unwrap();
1137 drop(w);
1138
1139 let w = wal(d.path());
1140 let committed = w.committed_end().unwrap();
1141 let path = d.path().join("wal");
1142 let mut bytes = std::fs::read(&path).unwrap();
1143 bytes[HDR + 3] ^= 1;
1144 std::fs::write(path, bytes).unwrap();
1145
1146 assert!(matches!(
1147 w.record_at(0, committed),
1148 Err(crate::Error::CorruptWal { offset: 0, .. })
1149 ));
1150 }
1151
1152 #[test]
1160 fn writers_cannot_emit_a_frame_larger_than_max_frame_bytes() {
1161 let widest_put = MAX_PAYLOAD_BYTES as usize;
1162 let widest_put_empty_batch = MAX_PAYLOAD_BYTES as usize;
1163 let widest_delete = MAX_PAYLOAD_BYTES as usize;
1164 let widest_commit = 0usize;
1165 for (kind, payload) in [
1166 ("Put", widest_put),
1167 ("PutEmptyBatch", widest_put_empty_batch),
1168 ("Delete", widest_delete),
1169 ("Commit", widest_commit),
1170 ] {
1171 assert!(
1172 (HDR + payload) as u64 <= MAX_FRAME_BYTES,
1173 "a {kind} frame can reach {} bytes, past MAX_FRAME_BYTES ({MAX_FRAME_BYTES})",
1174 HDR + payload
1175 );
1176 }
1177 }
1178
1179 #[test]
1185 fn a_partial_frame_at_the_end_is_an_ending_not_damage() {
1186 let d = tempfile::tempdir().unwrap();
1187 let mut w = wal(d.path());
1188 w.append(RecKind::Put, b"committed").unwrap();
1189 w.sync_data().unwrap();
1190 let end = w.end_offset();
1191 drop(w);
1192
1193 let mut hdr = [0u8; HDR];
1194 hdr[0..4].copy_from_slice(&1000u32.to_le_bytes());
1195 hdr[12] = RecKind::Put as u8;
1196 let mut bytes = std::fs::read(d.path().join("wal")).unwrap();
1197 bytes.extend_from_slice(&hdr);
1198 bytes.extend_from_slice(&vec![b'x'; 280]); std::fs::write(d.path().join("wal"), &bytes).unwrap();
1200
1201 let w2 = Wal::open(&d.path().join("wal"), IoMode::Buffered)
1202 .expect("an interrupted write is how this log ends, not damage");
1203 assert_eq!(w2.end_offset(), end);
1204 assert_eq!(w2.replay().unwrap().len(), 1);
1205 assert_eq!(std::fs::metadata(d.path().join("wal")).unwrap().len(), end,
1206 "only the incomplete physical frame is truncated");
1207 }
1208}