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 fn inspect(path: &Path, mode: IoMode) -> Result<Scan> {
243 let _ = mode;
244 let (file, _) = open_file(path, IoMode::Buffered)?;
245 Self::scan(&*file)
246 }
247
248 pub fn end_offset(&self) -> u64 { self.end }
249
250 fn frame_crc(hdr: &[u8; HDR], payload: &[u8]) -> u32 {
251 let mut h = [0u8; HDR];
252 h.copy_from_slice(hdr);
253 h[16..20].fill(0);
254 crc32c::crc32c_append(crc32c::crc32c(&h), payload)
255 }
256
257 fn frame_crc_from_file(
258 file: &dyn FileIo,
259 hdr: &[u8; HDR],
260 payload_at: u64,
261 payload_len: u64,
262 scratch: &mut [u8],
263 ) -> Result<u32> {
264 let mut h = *hdr;
265 h[16..20].fill(0);
266 let mut crc = crc32c::crc32c(&h);
267 let mut at = payload_at;
268 let end = payload_at.checked_add(payload_len).ok_or(Error::CorruptWal {
269 offset: payload_at,
270 why: "payload boundary overflows the WAL address space",
271 })?;
272 while at < end {
273 let n = std::cmp::min(scratch.len() as u64, end - at) as usize;
274 read_exact(file, &mut scratch[..n], at)?;
275 crc = crc32c::crc32c_append(crc, &scratch[..n]);
276 at += n as u64;
277 }
278 Ok(crc)
279 }
280
281 fn scan(file: &dyn FileIo) -> Result<Scan> {
294 let len = file.len()?;
295 let mut off = 0u64;
296 let mut next_lsn = 1u64;
297 let mut crc_scratch = vec![0u8; CRC_CHUNK];
298 let stop = loop {
299 let header_end = off.checked_add(HDR as u64).ok_or(Error::CorruptWal {
300 offset: off,
301 why: "header boundary overflows the WAL address space",
302 })?;
303 if header_end > len {
304 if off == len {
305 break Stop::End("the last verified frame reaches physical EOF");
306 }
307 if Self::is_all_zeros(file, off, len)? {
308 break Stop::End("nothing past the last good frame but zeros");
309 }
310 break Stop::Damaged {
311 offset: off,
312 why: "non-zero bytes remain but there is not a complete frame header",
313 };
314 }
315 let mut hdr = [0u8; HDR];
316 read_exact(file, &mut hdr, off)?;
317 let plen = u32::from_le_bytes(hdr[0..4].try_into().unwrap()) as u64;
318 let lsn = u64::from_le_bytes(hdr[4..12].try_into().unwrap());
319 if plen > MAX_PAYLOAD_BYTES {
320 break Stop::Damaged {
321 offset: off,
322 why: "frame payload length exceeds the writer maximum",
323 };
324 }
325 let frame_end = header_end.checked_add(plen).ok_or(Error::CorruptWal {
326 offset: off,
327 why: "frame boundary overflows the WAL address space",
328 })?;
329 if frame_end > len {
330 if !Self::has_writer_header_shape(&hdr) {
335 break Stop::Damaged {
336 offset: off,
337 why: "incomplete-looking frame header was not emitted by this writer",
338 };
339 }
340 if Self::has_verified_frame_after(
345 file,
346 off.checked_add(1).ok_or(Error::CorruptWal {
347 offset: off,
348 why: "WAL resynchronisation offset overflow",
349 })?,
350 len,
351 &mut crc_scratch,
352 )? {
353 break Stop::Damaged {
354 offset: off,
355 why: "incomplete-looking frame has verified log frames behind it",
356 };
357 }
358 break Stop::End("a bounded frame header crosses physical EOF");
359 }
360 let want = u32::from_le_bytes(hdr[16..20].try_into().unwrap());
361 if Self::frame_crc_from_file(file, &hdr, header_end, plen, &mut crc_scratch)? != want {
362 if Self::is_all_zeros(file, off, len)? {
363 break Stop::End("nothing past the last good frame but zeros");
364 }
365 break Stop::Damaged {
366 offset: off,
367 why: "complete frame fails its checksum",
368 };
369 }
370 if RecKind::from_u8(hdr[12]).is_none() {
372 break Stop::Damaged {
373 offset: off,
374 why: "frame verifies but names a record kind this build does not know",
375 };
376 }
377 next_lsn = lsn.checked_add(1).ok_or(Error::CorruptWal {
378 offset: off,
379 why: "verified frame exhausts the log sequence number space",
380 })?;
381 off = frame_end;
382 };
383 Ok(Scan { end: off, next_lsn, stop })
384 }
385
386 fn is_all_zeros(file: &dyn FileIo, off: u64, len: u64) -> Result<bool> {
392 const CHUNK: usize = 64 * 1024;
393 let mut buf = vec![0u8; CHUNK];
394 let mut at = off;
395 while at < len {
396 let n = std::cmp::min(CHUNK as u64, len - at) as usize;
397 read_exact(file, &mut buf[..n], at)?;
398 if buf[..n].iter().any(|&b| b != 0) { return Ok(false); }
399 at += n as u64;
400 }
401 Ok(true)
402 }
403
404 pub fn append(&mut self, kind: RecKind, payload: &[u8]) -> Result<u64> {
405 if payload.len() as u64 > MAX_PAYLOAD_BYTES {
406 return Err(Error::TooLarge);
407 }
408 let lsn = self.next_lsn;
409 let next_lsn = self.next_lsn.checked_add(1).ok_or(Error::CorruptWal {
410 offset: self.end,
411 why: "log sequence number space is exhausted",
412 })?;
413 let frame_len = (HDR as u64).checked_add(payload.len() as u64).ok_or(Error::TooLarge)?;
414 if frame_len > MAX_FRAME_BYTES { return Err(Error::TooLarge); }
415 let new_end = self.end.checked_add(frame_len).ok_or(Error::TooLarge)?;
416 if self.limit.is_some_and(|cap| new_end > cap) {
417 return Err(Error::ResourceLimit("WAL allowance full; reduce transaction size"));
418 }
419 let mut hdr = [0u8; HDR];
420 hdr[0..4].copy_from_slice(&(payload.len() as u32).to_le_bytes());
421 hdr[4..12].copy_from_slice(&lsn.to_le_bytes());
422 hdr[12] = kind as u8;
423 #[cfg(feature = "write-trace")]
424 let crc_started = crate::write_trace::active().then(std::time::Instant::now);
425 let c = Self::frame_crc(&hdr, payload);
426 #[cfg(feature = "write-trace")]
427 if let Some(started) = crc_started {
428 crate::write_trace::add(crate::write_trace::Field::WalCrc, started.elapsed());
429 }
430 hdr[16..20].copy_from_slice(&c.to_le_bytes());
431
432 #[cfg(feature = "write-trace")]
437 let copy_started = crate::write_trace::active().then(std::time::Instant::now);
438 self.buf.extend_from_slice(&hdr);
439 self.buf.extend_from_slice(payload);
440 #[cfg(feature = "write-trace")]
441 if let Some(started) = copy_started {
442 crate::write_trace::add(crate::write_trace::Field::WalBufferCopy, started.elapsed());
443 crate::write_trace::value_copy();
444 }
445 self.next_lsn = next_lsn;
446 self.end = new_end;
447 if self.buf.len() >= WAL_BUF {
448 #[cfg(feature = "write-trace")]
449 let flush_started = crate::write_trace::active().then(std::time::Instant::now);
450 self.flush()?;
451 #[cfg(feature = "write-trace")]
452 if let Some(started) = flush_started {
453 crate::write_trace::add(crate::write_trace::Field::WalFlush, started.elapsed());
454 }
455 }
456 Ok(lsn)
457 }
458
459 pub fn flush(&mut self) -> Result<()> {
464 if self.buf.is_empty() { return Ok(()); }
465 write_all(&*self.file, &self.buf, self.flushed)?;
466 crate::write_stats::add(crate::write_stats::Phase::Wal, self.buf.len() as u64);
467 self.flushed += self.buf.len() as u64;
468 self.buf.clear();
469 debug_assert_eq!(self.flushed, self.end);
470 Ok(())
471 }
472
473 pub fn sync_data(&mut self) -> Result<()> { self.flush()?; self.file.sync_data() }
474
475 pub fn sync_full(&mut self) -> Result<()> { self.flush()?; self.file.sync_full() }
477
478 pub fn sync_full_primitive(&self) -> &'static str { self.file.sync_full_primitive() }
481
482 pub(crate) fn committed_end(&self) -> Result<u64> {
486 assert!(self.buf.is_empty(), "recovery with {} bytes buffered", self.buf.len());
490 let mut off = 0u64;
491 let mut committed = 0u64;
492 while off < self.end {
493 let mut hdr = [0u8; HDR];
494 read_exact(&*self.file, &mut hdr, off)?;
495 let plen = u32::from_le_bytes(hdr[0..4].try_into().unwrap()) as u64;
496 if plen > MAX_PAYLOAD_BYTES {
497 return Err(Error::CorruptWal { offset: off, why: "recovery payload exceeds the writer maximum" });
498 }
499 let next = off.checked_add(HDR as u64).and_then(|n| n.checked_add(plen))
500 .ok_or(Error::CorruptWal { offset: off, why: "frame boundary overflow during recovery" })?;
501 if next > self.end {
502 return Err(Error::CorruptWal { offset: off, why: "frame crosses verified WAL boundary during recovery" });
503 }
504 let kind = RecKind::from_u8(hdr[12]).ok_or(Error::CorruptWal {
505 offset: off, why: "unknown record kind inside verified WAL",
506 })?;
507 if kind == RecKind::Commit { committed = next; }
508 off = next;
509 }
510 Ok(committed)
511 }
512
513 pub(crate) fn record_at(&self, off: u64, committed_end: u64)
516 -> Result<Option<(u64, RecKind, Vec<u8>, u64)>>
517 {
518 if off == committed_end { return Ok(None); }
519 if off > committed_end || committed_end > self.end {
520 return Err(Error::CorruptWal { offset: off, why: "recovery cursor outside committed WAL boundary" });
521 }
522 let header_end = off.checked_add(HDR as u64)
523 .ok_or(Error::CorruptWal { offset: off, why: "header boundary overflow during recovery" })?;
524 if header_end > committed_end {
525 return Err(Error::CorruptWal { offset: off, why: "truncated header inside committed WAL boundary" });
526 }
527 let mut hdr = [0u8; HDR];
528 read_exact(&*self.file, &mut hdr, off)?;
529 let plen = u32::from_le_bytes(hdr[0..4].try_into().unwrap()) as usize;
530 if plen as u64 > MAX_PAYLOAD_BYTES {
531 return Err(Error::CorruptWal { offset: off, why: "recovery payload exceeds the writer maximum" });
532 }
533 let next = header_end.checked_add(plen as u64)
534 .ok_or(Error::CorruptWal { offset: off, why: "payload boundary overflow during recovery" })?;
535 if next > committed_end {
536 return Err(Error::CorruptWal { offset: off, why: "frame crosses committed WAL boundary" });
537 }
538 let mut payload = vec![0u8; plen];
539 read_exact(&*self.file, &mut payload, header_end)?;
540 let want = u32::from_le_bytes(hdr[16..20].try_into().unwrap());
541 if Self::frame_crc(&hdr, &payload) != want {
542 return Err(Error::CorruptWal { offset: off, why: "recovery frame changed after the opening CRC walk" });
543 }
544 let lsn = u64::from_le_bytes(hdr[4..12].try_into().unwrap());
546 let kind = RecKind::from_u8(hdr[12]).ok_or(Error::CorruptWal {
547 offset: off, why: "unknown record kind inside committed WAL",
548 })?;
549 Ok(Some((lsn, kind, payload, next)))
550 }
551
552 #[cfg(test)]
553 fn replay(&self) -> Result<Vec<(u64, RecKind, Vec<u8>)>> {
554 let end = self.end;
557 let mut out = Vec::new();
558 let mut off = 0;
559 while let Some((lsn, kind, payload, next)) = self.record_at(off, end)? {
560 out.push((lsn, kind, payload));
561 off = next;
562 }
563 Ok(out)
564 }
565
566 pub fn set_lsn_floor(&mut self, floor: u64) {
571 if floor > self.next_lsn { self.next_lsn = floor; }
572 }
573
574 pub fn next_lsn(&self) -> u64 { self.next_lsn }
575
576 pub(crate) fn salvage_committed(src: &Path, dst: &Path) -> Result<SalvagedWal> {
584 use std::io::{Seek, SeekFrom};
585
586 let (file, _) = open_file(src, IoMode::Buffered)?;
587 let len = file.len()?;
588 let mut out = std::fs::File::create(dst)?;
589 let mut src_off = 0u64;
590 let mut out_end = 0u64;
591 let mut last_commit = 0u64;
592 let mut resynchronising = false;
593 let mut scratch = vec![0u8; CRC_CHUNK];
594
595 while src_off < len {
596 match Self::candidate_at(&*file, src_off, len, &mut scratch)? {
597 Candidate::Valid { next, kind } => {
598 if resynchronising {
599 src_off = next;
600 if kind == RecKind::Commit {
601 resynchronising = false;
602 }
603 continue;
604 }
605 copy_range(&*file, &mut out, src_off, next - src_off, &mut scratch)?;
606 out_end = out_end.checked_add(next - src_off)
607 .ok_or(Error::CorruptWal {
608 offset: src_off,
609 why: "salvaged WAL output length overflow",
610 })?;
611 if kind == RecKind::Commit { last_commit = out_end; }
612 src_off = next;
613 }
614 Candidate::Invalid => {
615 out.set_len(last_commit)?;
616 out.seek(SeekFrom::Start(last_commit))?;
617 out_end = last_commit;
618 resynchronising = true;
619 src_off = src_off.checked_add(1).ok_or(Error::CorruptWal {
620 offset: src_off,
621 why: "WAL resynchronisation offset overflow",
622 })?;
623 }
624 }
625 }
626 out.set_len(last_commit)?;
627 out.sync_all()?;
628 drop(out);
629
630 let hash = hash_prefix(dst, last_commit)?;
631 Ok(SalvagedWal { bytes: last_commit, hash })
632 }
633
634 fn candidate_at(
635 file: &dyn FileIo,
636 off: u64,
637 len: u64,
638 scratch: &mut [u8],
639 ) -> Result<Candidate> {
640 let Some(header_end) = off.checked_add(HDR as u64) else { return Ok(Candidate::Invalid) };
641 if header_end > len { return Ok(Candidate::Invalid); }
642 let mut hdr = [0u8; HDR];
643 read_exact(file, &mut hdr, off)?;
644 if !Self::has_writer_header_shape(&hdr) {
648 return Ok(Candidate::Invalid);
649 }
650 let plen = u32::from_le_bytes(hdr[0..4].try_into().unwrap()) as u64;
651 if plen > MAX_PAYLOAD_BYTES { return Ok(Candidate::Invalid); }
652 let Some(next) = header_end.checked_add(plen) else { return Ok(Candidate::Invalid) };
653 if next > len { return Ok(Candidate::Invalid); }
654 let want = u32::from_le_bytes(hdr[16..20].try_into().unwrap());
655 if Self::frame_crc_from_file(file, &hdr, header_end, plen, scratch)? != want {
656 return Ok(Candidate::Invalid);
657 }
658 let kind = RecKind::from_u8(hdr[12]).expect("kind was prefiltered above");
659 let lsn = u64::from_le_bytes(hdr[4..12].try_into().unwrap());
660 if lsn.checked_add(1).is_none() { return Ok(Candidate::Invalid); }
661 Ok(Candidate::Valid { next, kind })
662 }
663
664 fn has_writer_header_shape(hdr: &[u8; HDR]) -> bool {
665 hdr[13..16] == [0, 0, 0]
666 && RecKind::from_u8(hdr[12]).is_some()
667 && u64::from_le_bytes(hdr[4..12].try_into().unwrap()).checked_add(1).is_some()
668 }
669
670 fn has_verified_frame_after(
671 file: &dyn FileIo,
672 mut off: u64,
673 len: u64,
674 scratch: &mut [u8],
675 ) -> Result<bool> {
676 while off < len {
677 if matches!(Self::candidate_at(file, off, len, scratch)?, Candidate::Valid { .. }) {
678 return Ok(true);
679 }
680 off = off.checked_add(1).ok_or(Error::CorruptWal {
681 offset: off,
682 why: "WAL resynchronisation offset overflow",
683 })?;
684 }
685 Ok(false)
686 }
687
688 pub fn rotate(&mut self) -> Result<()> {
691 self.buf.clear();
693 self.flushed = 0;
694 self.file.set_len(0)?;
695 self.file.sync_data()?;
696 self.file.sync_dir()?;
697 self.end = 0;
698 Ok(())
699 }
700
701 pub(crate) fn rotate_published(&mut self) -> Result<()> {
705 self.buf.clear();
706 self.flushed = 0;
707 self.file.set_len(0)?;
708 self.file.sync_data()?;
709 self.end = 0;
710 Ok(())
711 }
712}
713
714fn read_exact(f: &dyn FileIo, buf: &mut [u8], off: u64) -> Result<()> { f.read_at(buf, off) }
717fn write_all(f: &dyn FileIo, buf: &[u8], off: u64) -> Result<()> { f.write_at(buf, off) }
718
719fn copy_range(
720 src: &dyn FileIo,
721 dst: &mut std::fs::File,
722 mut at: u64,
723 mut n: u64,
724 scratch: &mut [u8],
725) -> Result<()> {
726 use std::io::Write;
727 while n > 0 {
728 let take = std::cmp::min(n, scratch.len() as u64) as usize;
729 read_exact(src, &mut scratch[..take], at)?;
730 dst.write_all(&scratch[..take])?;
731 at += take as u64;
732 n -= take as u64;
733 }
734 Ok(())
735}
736
737pub(crate) fn hash_prefix(path: &Path, n: u64) -> Result<u32> {
738 use std::io::Read;
739 const CHUNK: usize = 64 * 1024;
740 let mut file = std::fs::File::open(path)?;
741 let mut buf = vec![0u8; CHUNK];
742 let mut left = n;
743 let mut hash = 0u32;
744 let mut first = true;
745 while left > 0 {
746 let want = std::cmp::min(left, CHUNK as u64) as usize;
747 file.read_exact(&mut buf[..want])?;
748 hash = if first {
749 first = false;
750 crc32c::crc32c(&buf[..want])
751 } else {
752 crc32c::crc32c_append(hash, &buf[..want])
753 };
754 left -= want as u64;
755 }
756 Ok(hash)
757}
758
759#[cfg(test)]
760mod tests {
761 use super::*;
762
763 fn wal(dir: &std::path::Path) -> Wal { Wal::open(&dir.join("wal"), crate::io::IoMode::Buffered).unwrap() }
764
765 #[test]
766 fn records_replay_in_order() {
767 let d = tempfile::tempdir().unwrap();
768 let mut w = wal(d.path());
769 w.append(RecKind::Put, b"one").unwrap();
770 w.append(RecKind::Put, b"two").unwrap();
771 w.append(RecKind::Commit, b"").unwrap();
772 w.sync_data().unwrap();
773 drop(w);
774
775 let got = wal(d.path()).replay().unwrap();
776 assert_eq!(got.len(), 3);
777 assert_eq!(got[0].2, b"one");
778 assert_eq!(got[1].2, b"two");
779 assert_eq!(got[2].1, RecKind::Commit);
780 }
781
782 #[test]
783 fn a_torn_tail_stops_replay_at_the_last_good_frame() {
784 let d = tempfile::tempdir().unwrap();
785 let mut w = wal(d.path());
786 w.append(RecKind::Put, b"good").unwrap();
787 w.sync_data().unwrap();
788 let n = w.end_offset();
789 drop(w);
790
791 let mut f = std::fs::OpenOptions::new().write(true).open(d.path().join("wal")).unwrap();
793 use std::io::{Seek, SeekFrom, Write};
794 f.seek(SeekFrom::Start(n)).unwrap();
795 f.write_all(&[0u8; 12]).unwrap(); f.sync_all().unwrap();
797
798 let got = wal(d.path()).replay().unwrap();
799 assert_eq!(got.len(), 1, "only the intact prefix replays");
800 }
801
802 #[test]
806 fn a_write_after_a_torn_tail_survives_the_next_open() {
807 let d = tempfile::tempdir().unwrap();
808 let mut w = wal(d.path());
809 w.append(RecKind::Put, b"before").unwrap();
810 w.sync_data().unwrap();
811 let n = w.end_offset();
812 drop(w);
813
814 let mut f = std::fs::OpenOptions::new().write(true).open(d.path().join("wal")).unwrap();
815 use std::io::{Seek, SeekFrom, Write};
816 f.seek(SeekFrom::Start(n)).unwrap();
817 let mut hdr = [0u8; HDR];
818 hdr[0..4].copy_from_slice(&100u32.to_le_bytes());
819 hdr[12] = RecKind::Put as u8;
820 f.write_all(&hdr).unwrap();
821 f.write_all(&[0xAA; 9]).unwrap();
822 f.sync_all().unwrap();
823
824 let mut w2 = wal(d.path());
826 w2.append(RecKind::Put, b"after").unwrap();
827 w2.sync_data().unwrap();
828 drop(w2);
829
830 let got = wal(d.path()).replay().unwrap();
831 let payloads: Vec<&[u8]> = got.iter().map(|r| r.2.as_slice()).collect();
832 assert_eq!(payloads, vec![b"before".as_ref(), b"after".as_ref()]);
833 }
834
835 #[test]
847 fn an_all_zero_extension_is_a_clean_ending_and_is_truncated() {
848 let d = tempfile::tempdir().unwrap();
849 let mut w = wal(d.path());
850 w.append(RecKind::Put, b"one").unwrap();
851 w.sync_data().unwrap();
852 let end = w.end_offset();
853 drop(w);
854
855 let f = std::fs::OpenOptions::new().write(true)
856 .open(d.path().join("wal")).unwrap();
857 f.set_len(end + 4096).unwrap(); f.sync_all().unwrap();
859 drop(f);
860
861 let w2 = wal(d.path());
862 assert_eq!(w2.end_offset(), end);
863 assert_eq!(
864 std::fs::metadata(d.path().join("wal")).unwrap().len(), end,
865 "open must not leave orphaned bytes on disk"
866 );
867 }
868
869 #[test]
872 fn append_lands_at_the_scanned_end_not_at_physical_eof() {
873 let d = tempfile::tempdir().unwrap();
874 let mut w = wal(d.path());
875 w.append(RecKind::Put, b"one").unwrap();
876 w.sync_data().unwrap();
877 let end = w.end_offset();
878
879 {
880 let f = std::fs::OpenOptions::new().write(true)
881 .open(d.path().join("wal")).unwrap();
882 f.set_len(end + 4096).unwrap();
883 f.sync_all().unwrap();
884 }
885
886 w.append(RecKind::Put, b"two").unwrap();
887 w.sync_data().unwrap();
888
889 let raw = std::fs::read(d.path().join("wal")).unwrap();
895 let hdr = &raw[end as usize..end as usize + HDR];
896 assert_eq!(
897 u32::from_le_bytes(hdr[0..4].try_into().unwrap()), 3,
898 "the second frame must physically start at the scanned end"
899 );
900 assert_eq!(hdr[12], RecKind::Put as u8);
901
902 let got = wal(d.path()).replay().unwrap();
903 let payloads: Vec<&[u8]> = got.iter().map(|r| r.2.as_slice()).collect();
904 assert_eq!(payloads, vec![b"one".as_ref(), b"two".as_ref()]);
905 }
906
907 #[test]
908 fn non_zero_garbage_shorter_than_a_frame_is_damage_and_is_preserved() {
909 for len in [1usize, 7, 19, 20, 31, 100] {
913 let d = tempfile::tempdir().unwrap();
914 let stub: Vec<u8> = (0..len).map(|i| (i + 1) as u8).collect();
915 std::fs::write(d.path().join("wal"), &stub).unwrap();
916 assert!(matches!(
917 Wal::open(&d.path().join("wal"), IoMode::Buffered),
918 Err(crate::Error::CorruptWal { offset: 0, .. })
919 ));
920 assert_eq!(std::fs::read(d.path().join("wal")).unwrap(), stub,
921 "a refused {len}-byte tail must remain byte-for-byte intact");
922 }
923 }
924
925 #[test]
926 fn a_log_payload_larger_than_the_writer_max_is_refused_before_allocation() {
927 use std::sync::atomic::{AtomicBool, Ordering};
928 use std::sync::Arc;
929
930 struct HeaderOnly {
931 hdr: [u8; HDR],
932 len: u64,
933 payload_read: Arc<AtomicBool>,
934 }
935 impl FileIo for HeaderOnly {
936 fn requires_alignment(&self) -> bool { false }
937 fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> {
938 if off == 0 && buf.len() == HDR {
939 buf.copy_from_slice(&self.hdr);
940 return Ok(());
941 }
942 self.payload_read.store(true, Ordering::Relaxed);
943 Err(std::io::Error::other("payload read must not happen").into())
944 }
945 fn write_at(&self, _: &[u8], _: u64) -> Result<()> { Ok(()) }
946 fn sync_data(&self) -> Result<()> { Ok(()) }
947 fn sync_full(&self) -> Result<()> { Ok(()) }
948 fn sync_full_primitive(&self) -> &'static str { "test" }
949 fn sync_dir(&self) -> Result<()> { Ok(()) }
950 fn len(&self) -> Result<u64> { Ok(self.len) }
951 fn set_len(&self, _: u64) -> Result<()> { Ok(()) }
952 }
953
954 let plen = u32::MAX;
955 let mut hdr = [0u8; HDR];
956 hdr[0..4].copy_from_slice(&plen.to_le_bytes());
957 hdr[12] = RecKind::Put as u8;
958 let payload_read = Arc::new(AtomicBool::new(false));
959 let f = HeaderOnly {
960 hdr,
961 len: HDR as u64 + plen as u64,
962 payload_read: payload_read.clone(),
963 };
964 assert!(matches!(Wal::open_on(Box::new(f)), Err(crate::Error::CorruptWal { offset: 0, .. })));
965 assert!(!payload_read.load(Ordering::Relaxed),
966 "an off-disk length above the writer maximum must be rejected before allocation/read");
967 }
968
969 #[test]
970 fn an_exhausted_lsn_is_refused_with_checked_arithmetic() {
971 let d = tempfile::tempdir().unwrap();
972 let mut hdr = [0u8; HDR];
973 hdr[4..12].copy_from_slice(&u64::MAX.to_le_bytes());
974 hdr[12] = RecKind::Commit as u8;
975 let crc = Wal::frame_crc(&hdr, &[]);
976 hdr[16..20].copy_from_slice(&crc.to_le_bytes());
977 std::fs::write(d.path().join("wal"), hdr).unwrap();
978 assert!(matches!(
979 Wal::open(&d.path().join("wal"), IoMode::Buffered),
980 Err(crate::Error::CorruptWal { offset: 0, .. })
981 ));
982 }
983
984 struct FailReadAt { inner: Box<dyn FileIo>, fail_at: u64 }
993 impl FileIo for FailReadAt {
994 fn requires_alignment(&self) -> bool { self.inner.requires_alignment() }
995 fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> {
996 if off == self.fail_at {
997 return Err(std::io::Error::other("injected read failure").into());
998 }
999 self.inner.read_at(buf, off)
1000 }
1001 fn write_at(&self, buf: &[u8], off: u64) -> Result<()> { self.inner.write_at(buf, off) }
1002 fn sync_data(&self) -> Result<()> { self.inner.sync_data() }
1003 fn sync_full(&self) -> Result<()> { self.inner.sync_full() }
1004 fn sync_full_primitive(&self) -> &'static str { self.inner.sync_full_primitive() }
1005 fn sync_dir(&self) -> Result<()> { self.inner.sync_dir() }
1006 fn len(&self) -> Result<u64> { self.inner.len() }
1007 fn set_len(&self, n: u64) -> Result<()> { self.inner.set_len(n) }
1008 }
1009
1010 #[test]
1016 fn an_io_error_on_a_header_read_is_an_error_not_a_clean_end() {
1017 let d = tempfile::tempdir().unwrap();
1018 let mut w = wal(d.path());
1019 let p = vec![b'A'; 100];
1020 for _ in 0..10 { w.append(RecKind::Put, &p).unwrap(); }
1021 w.sync_data().unwrap();
1022 let full = w.end_offset();
1023 drop(w);
1024 let frame = HDR as u64 + p.len() as u64;
1025 let midpoint = frame * 5;
1026
1027 let (real, _) = open_file(&d.path().join("wal"), IoMode::Buffered).unwrap();
1028 let failing = Box::new(FailReadAt { inner: real, fail_at: midpoint });
1029 match Wal::open_on(failing) {
1030 Err(crate::Error::Io(_)) => {}
1031 Err(e) => panic!("an unreadable header must surface as an I/O error, got {e:?}"),
1032 Ok(_) => panic!("an unreadable header must not be reported as a clean end of log"),
1033 }
1034 assert_eq!(
1035 std::fs::metadata(d.path().join("wal")).unwrap().len(), full,
1036 "a read that FAILED must not be able to shorten the log -- {full} bytes were on \
1037 disk and the reader could not read one header of them"
1038 );
1039 }
1040
1041 #[test]
1048 fn damage_with_more_log_behind_it_refuses_the_open_and_keeps_every_byte() {
1049 let d = tempfile::tempdir().unwrap();
1050 let mut w = wal(d.path());
1051 let p = vec![b'A'; 100];
1052 for _ in 0..100 { w.append(RecKind::Put, &p).unwrap(); }
1053 w.sync_data().unwrap();
1054 drop(w);
1055
1056 let path = d.path().join("wal");
1057 let frame = HDR + p.len();
1058 let victim = frame * 50; let mut bytes = std::fs::read(&path).unwrap();
1060 bytes[victim + 4] ^= 0x01; std::fs::write(&path, &bytes).unwrap();
1062
1063 match Wal::open(&path, IoMode::Buffered) {
1064 Err(crate::Error::CorruptWal { offset, .. }) => {
1065 assert_eq!(offset as usize, victim, "the refusal must name where it stopped");
1066 }
1067 Err(e) => panic!("expected CorruptWal, got {e:?}"),
1068 Ok(_) => panic!("damage with committed frames behind it must not open"),
1069 }
1070 assert_eq!(
1071 std::fs::read(&path).unwrap(), bytes,
1072 "a refused open must leave the log byte for byte as it found it"
1073 );
1074 }
1075
1076 #[test]
1083 fn a_verified_frame_naming_an_unknown_kind_is_damage() {
1084 let d = tempfile::tempdir().unwrap();
1085 let mut w = wal(d.path());
1086 w.append(RecKind::Put, b"good").unwrap();
1087 w.sync_data().unwrap();
1088 let end = w.end_offset();
1089 drop(w);
1090
1091 let mut hdr = [0u8; HDR];
1093 let payload = b"future".to_vec();
1094 hdr[0..4].copy_from_slice(&(payload.len() as u32).to_le_bytes());
1095 hdr[4..12].copy_from_slice(&7u64.to_le_bytes());
1096 hdr[12] = 9;
1097 let c = Wal::frame_crc(&hdr, &payload);
1098 hdr[16..20].copy_from_slice(&c.to_le_bytes());
1099 let mut bytes = std::fs::read(d.path().join("wal")).unwrap();
1100 bytes.extend_from_slice(&hdr);
1101 bytes.extend_from_slice(&payload);
1102 std::fs::write(d.path().join("wal"), &bytes).unwrap();
1103
1104 match Wal::open(&d.path().join("wal"), IoMode::Buffered) {
1105 Err(crate::Error::CorruptWal { offset, .. }) => assert_eq!(offset, end),
1106 Err(e) => panic!("expected CorruptWal, got {e:?}"),
1107 Ok(_) => panic!("a verified frame naming an unknown kind must not be walked past"),
1108 }
1109 assert_eq!(std::fs::read(d.path().join("wal")).unwrap(), bytes,
1110 "and it must not have been truncated on the way out");
1111 }
1112
1113 #[test]
1116 fn replay_refuses_a_byte_changed_after_the_opening_crc_walk() {
1117 let d = tempfile::tempdir().unwrap();
1118 let mut w = wal(d.path());
1119 w.append(RecKind::Put, b"\x01\0key-value").unwrap();
1120 w.append(RecKind::Commit, b"").unwrap();
1121 w.sync_data().unwrap();
1122 drop(w);
1123
1124 let w = wal(d.path());
1125 let committed = w.committed_end().unwrap();
1126 let path = d.path().join("wal");
1127 let mut bytes = std::fs::read(&path).unwrap();
1128 bytes[HDR + 3] ^= 1;
1129 std::fs::write(path, bytes).unwrap();
1130
1131 assert!(matches!(
1132 w.record_at(0, committed),
1133 Err(crate::Error::CorruptWal { offset: 0, .. })
1134 ));
1135 }
1136
1137 #[test]
1145 fn writers_cannot_emit_a_frame_larger_than_max_frame_bytes() {
1146 let widest_put = MAX_PAYLOAD_BYTES as usize;
1147 let widest_put_empty_batch = MAX_PAYLOAD_BYTES as usize;
1148 let widest_delete = MAX_PAYLOAD_BYTES as usize;
1149 let widest_commit = 0usize;
1150 for (kind, payload) in [
1151 ("Put", widest_put),
1152 ("PutEmptyBatch", widest_put_empty_batch),
1153 ("Delete", widest_delete),
1154 ("Commit", widest_commit),
1155 ] {
1156 assert!(
1157 (HDR + payload) as u64 <= MAX_FRAME_BYTES,
1158 "a {kind} frame can reach {} bytes, past MAX_FRAME_BYTES ({MAX_FRAME_BYTES})",
1159 HDR + payload
1160 );
1161 }
1162 }
1163
1164 #[test]
1170 fn a_partial_frame_at_the_end_is_an_ending_not_damage() {
1171 let d = tempfile::tempdir().unwrap();
1172 let mut w = wal(d.path());
1173 w.append(RecKind::Put, b"committed").unwrap();
1174 w.sync_data().unwrap();
1175 let end = w.end_offset();
1176 drop(w);
1177
1178 let mut hdr = [0u8; HDR];
1179 hdr[0..4].copy_from_slice(&1000u32.to_le_bytes());
1180 hdr[12] = RecKind::Put as u8;
1181 let mut bytes = std::fs::read(d.path().join("wal")).unwrap();
1182 bytes.extend_from_slice(&hdr);
1183 bytes.extend_from_slice(&vec![b'x'; 280]); std::fs::write(d.path().join("wal"), &bytes).unwrap();
1185
1186 let w2 = Wal::open(&d.path().join("wal"), IoMode::Buffered)
1187 .expect("an interrupted write is how this log ends, not damage");
1188 assert_eq!(w2.end_offset(), end);
1189 assert_eq!(w2.replay().unwrap().len(), 1);
1190 assert_eq!(std::fs::metadata(d.path().join("wal")).unwrap().len(), end,
1191 "only the incomplete physical frame is truncated");
1192 }
1193}