Skip to main content

kernel/
wal.rs

1//! Write-ahead log. Frame layout, little-endian:
2//!   0  len     u32   payload length
3//!   4  lsn     u64
4//!   12 kind    u8
5//!   13 pad     u8 x3
6//!   16 crc32c  u32   over the 16-byte header with crc zeroed, then payload
7//!   20 payload
8//!
9//! THE TORN TAIL, and the two mechanisms that hold the property.
10//!
11//! The predecessor engine lost five rows that had been written AND fsynced. A
12//! crash left a complete header with a short payload; the next writer appended
13//! at PHYSICAL END OF FILE, so those bytes survived and the new records landed
14//! after them. On replay the torn header's declared length swallowed the new
15//! records as its own payload, failed CRC, and stopped replay before reaching
16//! them. Every torn-tail test in that engine passed, because none of them wrote
17//! anything AFTER recovering.
18//!
19//! Here the property is held REDUNDANTLY, by two independent mechanisms, and
20//! **either one alone is sufficient**:
21//!
22//!   1. `open` truncates to the end of the last CRC-valid frame only when the
23//!      next header declares a frame that physically crosses EOF, or when
24//!      every remaining byte is zero. A complete CRC-failed frame and a
25//!      non-zero fragment shorter than a header are damage, never endings.
26//!   2. `append` writes at `self.end` — the offset `scan`'s CRC walk
27//!      established — and never at physical EOF, so a new frame is laid down ON
28//!      TOP of any torn bytes and there is nothing left to swallow it.
29//!
30//! **This redundancy is a maintenance hazard and must be treated as one.**
31//! Removing either mechanism alone breaks NO test, because the other still holds
32//! the property. Someone deleting one will see a green suite and reasonably
33//! conclude it was dead code; someone later deleting the other will see a green
34//! suite too, right up until a power cut. Reproducing the original bug requires
35//! removing BOTH — which is exactly what the falsification in this task does.
36//!
37//! Because the integration test cannot distinguish the two, each mechanism is
38//! additionally pinned by its own direct test:
39//! `an_all_zero_extension_is_a_clean_ending_and_is_truncated` and
40//! `append_lands_at_the_scanned_end_not_at_physical_eof`.
41//!
42//! WHY A WALK MUST SAY WHY IT STOPPED (Task 19; Task 17 final review, F4).
43//!
44//! `scan` walks frames until one cannot be accepted. There are two utterly
45//! different reasons that can happen, and for most of this engine's life
46//! they left through the same door:
47//!
48//!   - THE LOG ENDS HERE. A writer was interrupted; what follows is a
49//!     partial frame, or nothing, or zeros. Discarding it costs nothing,
50//!     because no completed write is in it.
51//!   - SOMETHING IS WRONG HERE. A bit rotted, a sector went bad, a header
52//!     verifies but names a record kind this build does not know. What
53//!     follows may be hundreds of committed frames.
54//!
55//! Treating the second as the first is how one flipped bit at the midpoint
56//! of a log destroyed 1,499 of 3,000 committed rows while `open` returned
57//! `Ok`, and how a single injected EIO on a header read erased half a log.
58//! Both measured. The distinction is now carried by `Stop`, a type -- not by
59//! a comment asserting that every `break` means the same thing, which is
60//! what the previous version of this file did, incorrectly, in four places.
61//!
62//! Three rules follow, and all three are load-bearing:
63//!   1. Only `Stop::End` may truncate.
64//!   2. A read that FAILED is not an answer about content: it leaves as
65//!      `Err`. "I could not read this" is not "the log ends here".
66//!   3. Before CRC, the payload bound and physical extent limit reads. The
67//!      unverified writer shape may only move classification toward DAMAGE;
68//!      kind and LSN are trusted only after CRC verifies them.
69//!
70//! And a fourth rule, about the refusal itself: `Stop::Damaged` makes
71//! `open` refuse, which is only half a design. A refusal nothing can clear
72//! is as unrecoverable as a deletion, and Law 5 forbids both. `recover()`
73//! is the other half -- it copies the whole log aside intact,
74//! independently re-reads the copy, then reconstructs committed regions that
75//! can be resynchronised as the live log -- and
76//! `a_damaged_log_refuses_to_open_and_recover_clears_it` pins the two
77//! halves together, because either alone is a defect.
78
79use crate::io::{open_file, FileIo, IoMode};
80use crate::{Error, Result};
81use std::path::Path;
82
83const HDR: usize = 20;
84
85/// The largest number of bytes ONE frame this crate writes can occupy on
86/// disk. The whole encoded frame must fit in the format's `u32` size bound;
87/// D7's value-to-end frames therefore retain essentially the whole 4 GiB
88/// payload range instead of returning to the old 64 KiB truncation.
89///
90/// This is a validity bound on every disk frame as well as every writer. The
91/// length is checked before allocation: a larger value read off disk is
92/// damage even when the file is large enough to satisfy it.
93const 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/// Why `scan` stopped walking, as a type rather than as a comment.
98///
99/// `scan` had EIGHT `break` statements, four of which meant "the log
100/// genuinely ends here" and four of which meant "something is wrong here" --
101/// and all eight left through the same `Ok((end, ..))` return, which
102/// `open_impl` then handed to `set_len(end)`. So an injected EIO on a single
103/// header read returned a clean end and erased 640 of 1,280 bytes, and one
104/// flipped bit at the midpoint of a 3,000-row log destroyed 1,499 committed
105/// rows while `open()` returned `Ok`. Both measured (Task 17 final review,
106/// F4).
107///
108/// The distinction is now in the type system, so no exit can join the wrong
109/// one by accident:
110///   - `End` is the ONLY variant `open_impl` will truncate on.
111///   - `Damaged` refuses the open, touching nothing -- and, because a
112///     refusal with no way out is its own Law 5 violation, `recover()`
113///     clears it by setting the whole log aside intact and reconstructing
114///     independently verified committed regions.
115///   - An I/O error is neither: it is not an answer about the log's content
116///     at all, so it leaves `scan` as `Err` and never reaches this enum.
117#[derive(Debug, Clone, Copy, PartialEq, Eq)]
118pub enum Stop {
119    /// The log genuinely ends at `end`. Everything from there to EOF is
120    /// either a writer-shaped header whose bounded frame crosses physical
121    /// EOF, or an all-zero extension. `why` names which proof was established.
122    End(&'static str),
123    /// Something is wrong at `offset`. The remaining bytes may include
124    /// committed frames regardless of how short they are. Nothing here may
125    /// be discarded.
126    Damaged { offset: u64, why: &'static str },
127}
128
129/// What a walk of the log found: how far it got, the LSN to carry on from,
130/// and -- the part that used to be missing -- WHY it stopped.
131#[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
163/// 256 KiB append buffer. Appending was one pwrite + one Vec per record --
164/// measured 59.4% of a load's wall time inside pwrite, 2.1M syscalls per 1M
165/// rows. No engine does that (PostgreSQL: wal_buffers). Constant, so Law 1
166/// holds. `end` = logical end; `flushed` = what the file has; invariant
167/// flushed + buf.len() == end. Every barrier flushes first, so no durability
168/// promise changes; commit() must flush in every mode -- Off promises no
169/// BARRIER, it never promised the record stays in process RAM.
170const 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        // The WAL is byte-addressed, not page-addressed, so it always uses
203        // buffered I/O regardless of the store's mode.
204        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    /// `open` with the file handed in, so a test can drive the whole open
212    /// path -- classification, refusal, truncation -- against a `FileIo`
213    /// that fails a specific read. Without this seam the "an I/O error is
214    /// not the end of the log" property can only be argued, and it was
215    /// argued wrongly for three rounds.
216    fn open_on(file: Box<dyn FileIo>) -> Result<Wal> {
217        let scan = Self::scan(&*file)?;
218        match scan.stop {
219            // Only ever on `End`. See `Stop`, and `set_len`'s own note
220            // below on what this is and is not buying now.
221            Stop::End(_) => {
222                // Discard anything after the last good frame before we
223                // append to it -- mechanism 1 of the two the module doc
224                // comment names. What makes this safe, which it was not
225                // before, is that `End` is now a CLASSIFIED answer: the
226                // bytes past `end` have been shown to be a physically
227                // incomplete bounded frame, or zeros.
228                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    /// Walk the log without opening it for use and without changing a byte
238    /// of it -- no `set_len`, no `sync_data`, and no refusal either: the
239    /// caller gets the `Stop` and decides. `recover()` is the caller, and it
240    /// needs all three of those properties, because it is the path that has
241    /// to still work on exactly the images `open` refuses.
242    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    /// Walk frames from offset 0 until one of them cannot be accepted, and
282    /// say WHY the walk stopped (see `Stop`).
283    ///
284    /// Two rules govern the order of the checks:
285    ///
286    /// 1. Before CRC, payload bounds and physical extents limit reads. An
287    ///    impossible writer shape may only reject toward `Damaged`; semantic
288    ///    fields are trusted for acceptance only after verification.
289    /// 2. A read that FAILS is not an answer about the log's content. EIO on
290    ///    a header read means "I could not read", not "the log ends here",
291    ///    so it leaves as `Err` -- measured: laundering it into a clean end
292    ///    made `set_len` erase half the log on one transient error.
293    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                // An incomplete physical write can leave a complete header,
331                // but that header must still have the shape every writer
332                // emits. Rejecting an implausible unverified header is the
333                // conservative direction: it preserves bytes and refuses.
334                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                // A corrupt length byte in a COMPLETE frame can also point
341                // past EOF. If an independently checksummed frame exists
342                // later, this cannot be a lone interrupted physical write.
343                // Preserve the whole non-zero tail and refuse ordinary open.
344                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            // Verified. Only now may anything in the header be believed.
371            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    /// Read `[off, len)` in fixed-size chunks looking for a non-zero byte.
387    /// Chunked, not slurped: the remainder can be as large as the log, and
388    /// nothing in this engine is allowed to allocate in proportion to the
389    /// store (Law 1). It runs only after frame verification has already
390    /// stopped, so ordinary valid frames never pay for it.
391    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        // Into the buffer, no per-record allocation. Frames still land at
433        // self.end (scan's CRC-walk offset, never physical EOF), so a torn
434        // tail is still overwritten, not landed behind: flush writes at
435        // self.flushed which starts at the same scan.end.
436        #[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    /// `SyncMode::Normal`'s barrier: durable against an OS crash, not
460    /// necessarily against a loss of power to the drive.
461    /// Push the buffer to the file. Not a barrier. Every path that reads the
462    /// file or promises durability calls this first.
463    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    /// `SyncMode::Full`'s barrier: the strongest this platform can issue.
476    pub fn sync_full(&mut self) -> Result<()> { self.flush()?; self.file.sync_full() }
477
478    /// The exact primitive `sync_full` issues here, so a measurement can name
479    /// what it did rather than imply it.
480    pub fn sync_full_primitive(&self) -> &'static str { self.file.sync_full_primitive() }
481
482    /// First recovery pass: find the byte immediately after the last Commit.
483    /// Headers are fixed-size and payloads are skipped, so retained RAM is one
484    /// header regardless of the number or size of committed transactions.
485    pub(crate) fn committed_end(&self) -> Result<u64> {
486        // Reads the FILE; anything buffered would be silently missing -- a
487        // replay that drops committed records. Only caller is open, where
488        // nothing is buffered; checked, not trusted.
489        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    /// Second recovery pass: read exactly one frame from the committed prefix.
514    /// The returned payload is the only transaction data retained by replay.
515    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        // Only verified fields may control replay.
545        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        // Unit tests for the scanner inspect every verified frame, including
555        // an uncommitted tail; production recovery uses `committed_end` above.
556        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    /// Raise the LSN floor. `scan` derives `next_lsn` from the file, so after a
567    /// rotation leaves the log empty, a restart would otherwise renumber from 1.
568    /// The store persists the high-water mark in the superblock and re-supplies
569    /// it here at open.
570    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    /// Reconstruct committed regions from a damaged log into a fresh file.
577    ///
578    /// A damaged region invalidates the transaction around it. Resync scans
579    /// byte-by-byte for a complete independently checksummed frame, discards
580    /// through the first verified Commit boundary, then copies later complete
581    /// transactions. Any later damage repeats that rule. The source is never
582    /// changed; recovery has already made and verified its quarantine copy.
583    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        // Cheap necessary writer invariants keep byte-wise resynchronisation
645        // linear on arbitrary input. They never make a candidate valid; a
646        // surviving candidate still has to pass its full CRC independently.
647        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    /// Discard the log. MUST be called only after the pages it protects are
689    /// durably on disk and the directory has been fsynced.
690    pub fn rotate(&mut self) -> Result<()> {
691        // Discard, not flush: the log is being thrown away.
692        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    /// Ordinary checkpoint after durable data/root publication, with already
702    /// durable filenames. Truncation still needs its file barrier, but creates
703    /// no directory-entry obligation. Create/recovery retain `rotate` above.
704    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
714// The FileIo trait is page-aligned by contract; the WAL needs byte access, so it
715// wraps unaligned reads/writes here rather than weakening the trait.
716fn 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        // Simulate a crash mid-frame: a complete header, a short payload.
792        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();   // header claiming a payload that isn't there
796        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    /// THE test. sekejap lost five fsynced rows to exactly this: a writer that
803    /// appended after a torn header, so replay swallowed the new records as the
804    /// torn frame's payload.
805    #[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        // Reopen and write MORE. This is the half every torn-tail test forgets.
825        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    /// Mechanism 1, pinned directly. The integration test above cannot see
836    /// this one, because mechanism 2 would hold the property without it --
837    /// and that is still exactly true after Task 19: deleting the
838    /// `set_len(scan.end)` in `open_on` fails THIS test and nothing else in
839    /// the suite (checked: 73 lib tests, the 200-trial corruption sweep,
840    /// durability, recovery and bulk all stay green). It is kept as the
841    /// documented redundancy this module exists to explain, and it is safe
842    /// to keep now in a way it was not before: `End` is a classified
843    /// answer: a bounded frame physically crosses EOF, or every remaining
844    /// byte is zero. Here it is the zeros branch -- only "these cannot be a
845    /// record anyone wrote" makes it an ending.
846    #[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();   // orphaned bytes past the last frame
858        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    /// Mechanism 2, pinned directly. Bytes are added past the log's logical end
870    /// AFTER open, so mechanism 1 has already run and cannot mask this.
871    #[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        // Read the bytes off disk, bypassing the Wal's own bookkeeping.
890        // Asserting on `end_offset()` here would be tautological: `self.end` is
891        // incremented unconditionally by `append`, whatever offset the write
892        // actually landed at, so that assertion passes identically for a broken
893        // implementation. Only the file can say where the frame really went.
894        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        // The whole 1..=7 range, not a single representative. Short is not
910        // evidence of an interrupted write: every non-zero byte remains
911        // evidence until a complete physical frame can be ruled out.
912        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    // -- Task 19: the exit taxonomy (Task 17 final review, F4) --
985
986    /// A `FileIo` that fails one specific read. `scan` used to swallow this:
987    /// `read_exact(..).is_err()` was one of its eight `break`s, so an EIO
988    /// came back as a clean `end` and `open`'s `set_len` then erased
989    /// everything behind it (measured: 640 of 1,280 bytes on a single
990    /// injected failure). An I/O error is not a statement about what the
991    /// log contains.
992    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    /// Ten equal frames; the header read for frame 5 fails. The open must
1011    /// come back `Err`, and -- the half that actually cost data -- the file
1012    /// must still be ten frames long. With the `read_exact(..).is_err() =>
1013    /// break` this replaces, `scan` returns `end = 640` and `set_len`
1014    /// truncates the last five frames away, `open` reporting `Ok`.
1015    #[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    /// One flipped bit in the middle of a log, with committed frames behind
1042    /// it. This is the measured 1,499-rows-for-one-byte case at its own
1043    /// level: the frame fails CRC, but there are thousands of bytes behind
1044    /// it -- far more than one interrupted write could have left -- so it is
1045    /// damage, not an ending. The open must refuse, and every byte must
1046    /// still be there afterwards for `recover()` to work with.
1047    #[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;          // the header of frame 50
1059        let mut bytes = std::fs::read(&path).unwrap();
1060        bytes[victim + 4] ^= 0x01;        // one bit of the LSN: only the CRC can see it
1061        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    /// The other half of the CRC-first rule. A frame that VERIFIES cannot
1077    /// have come from an interrupted write, so an unknown kind byte on it is
1078    /// damage however small the remainder is -- and, conversely, the kind
1079    /// byte is never consulted before the CRC has spoken, so a torn tail
1080    /// with a garbage kind byte is still an ordinary ending (the test above
1081    /// this one and `a_write_after_a_torn_tail_survives_the_next_open`).
1082    #[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        // A complete, correctly checksummed frame whose kind byte is 9.
1092        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    /// The recovery read is independently checksummed. A byte changing after
1114    /// `open` performed its scan must be refused before the payload is used.
1115    #[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    /// `MAX_FRAME_BYTES` is only as good as the claim that no writer exceeds
1138    /// it, and the last two bounds in this file were wrong by two bytes and
1139    /// by a whole record respectively. So assert it per `RecKind`, against
1140    /// the widest payload each writer can actually construct, rather than
1141    /// against a comment. Every current logical writer is stopped by
1142    /// `Wal::append`'s payload bound; `Commit` is empty; `PageImage` is never
1143    /// appended by anything, and if that changes this test is what fails.
1144    #[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    /// A real, partial frame at the physical end -- the case the whole
1165    /// module exists for -- must still be an ordinary ending: verification
1166    /// fails because the bounded header declares bytes beyond physical EOF,
1167    /// so nothing behind it can be lost by truncating. 300 bytes of a
1168    /// 1,000-byte frame.
1169    #[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]);   // 300 of the 1,020 bytes landed
1184        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}