Skip to main content

kernel/
store.rs

1//! The store ties the pool, the log and the tree together and owns the two
2//! orderings that durability depends on.
3//!
4//! COMMIT    append -> fsync -> apply to pages -> publish
5//! CHECKPOINT flush pages -> fsync file -> fsync dir -> rotate log
6//!
7//! Moving the rotation earlier changes nothing observable until the machine loses
8//! power, which is why it has an explicit test rather than a review comment.
9
10use crate::btree::{BTree, RangeIter};
11use crate::budget::MemoryBudget;
12use crate::io::{open_file_writer, Barrier, FileIo, IoMode};
13#[cfg(test)]
14use crate::io::open_file;
15use crate::meta::Meta;
16use crate::page::PAGE_SIZE;
17use crate::pool::BufferPool;
18use crate::wal::{RecKind, Wal};
19use crate::{Error, Result};
20use std::cell::Cell;
21use std::path::Path;
22use std::sync::Arc;
23
24/// Deterministic failures at the SQL/native-to-storage boundary. This is
25/// compiled only for tests: production stores contain neither the branch nor
26/// the counter, while callers can prove error handling without depending on a
27/// full disk, permissions, or a particular filesystem.
28#[cfg(feature = "test-support")]
29#[derive(Debug, Clone, Copy, PartialEq, Eq)]
30pub enum TestFaultKind { Read, Write, Commit }
31
32#[cfg(feature = "test-support")]
33#[derive(Debug, Default)]
34pub struct TestFaultInjector {
35    armed: std::sync::Mutex<Option<(TestFaultKind, usize)>>,
36}
37
38#[cfg(feature = "test-support")]
39impl TestFaultInjector {
40    /// Fail the matching operation after `successful_matches` matching calls.
41    pub fn arm(&self, kind: TestFaultKind, successful_matches: usize) {
42        *self.armed.lock().unwrap() = Some((kind, successful_matches));
43    }
44
45    fn check(&self, kind: TestFaultKind) -> Result<()> {
46        let mut armed = self.armed.lock().unwrap();
47        let Some((wanted, remaining)) = armed.as_mut() else { return Ok(()) };
48        if *wanted != kind { return Ok(()) }
49        if *remaining > 0 {
50            *remaining -= 1;
51            return Ok(());
52        }
53        *armed = None;
54        Err(std::io::Error::other(format!("injected {kind:?} failure at storage boundary")).into())
55    }
56}
57
58/// One bounded logical batch for empty-valued index rows: at most 64 keys and
59/// one page-sized WAL payload. Ordinary GIN batches are ~2 KiB. Both bounds
60/// are format-validated during recovery.
61pub const PUT_EMPTY_BATCH_MAX_KEYS: usize = 64;
62
63#[derive(Debug, Clone, Copy, PartialEq, Eq)]
64pub enum SyncMode { Full, Normal, Off }
65
66#[derive(Debug, Clone, Copy)]
67pub struct Config { pub budget_bytes: usize, pub io: IoMode, pub sync: SyncMode }
68
69/// Durable description of an unreachable, independently verified graft.
70/// Persist this beside a resumable build before acknowledging its final
71/// sorter watermark. On reopen, [`Store::publish_existing_candidate`] checks
72/// both generations, re-verifies every reachable candidate page through a
73/// fresh handle, reconstructs the standing parent path, and only then flips
74/// the root.
75#[derive(Clone, Debug, PartialEq, Eq)]
76pub struct PreparedGraft {
77    pub base_generation: u64,
78    pub write_generation: u64,
79    pub root: u32,
80    pub rows: u64,
81    pub min: Vec<u8>,
82    pub max: Vec<u8>,
83    pub last_next: u32,
84    pub inserted_min: Vec<u8>,
85    pub inserted_max: Vec<u8>,
86    pub inserted_rows: u64,
87}
88
89impl PreparedGraft {
90    const MAGIC: &'static [u8; 8] = b"KGRFT01\0";
91
92    pub fn encode(&self) -> Result<Vec<u8>> {
93        let mut out = Vec::new();
94        out.extend_from_slice(Self::MAGIC);
95        out.extend_from_slice(&self.base_generation.to_le_bytes());
96        out.extend_from_slice(&self.write_generation.to_le_bytes());
97        out.extend_from_slice(&self.root.to_le_bytes());
98        out.extend_from_slice(&self.rows.to_le_bytes());
99        out.extend_from_slice(&self.last_next.to_le_bytes());
100        out.extend_from_slice(&self.inserted_rows.to_le_bytes());
101        for bytes in [&self.min, &self.max, &self.inserted_min, &self.inserted_max] {
102            let len = u32::try_from(bytes.len()).map_err(|_| Error::TooLarge)?;
103            out.extend_from_slice(&len.to_le_bytes());
104            out.extend_from_slice(bytes);
105        }
106        let crc = crc32c::crc32c(&out);
107        out.extend_from_slice(&crc.to_le_bytes());
108        Ok(out)
109    }
110
111    pub fn decode(bytes: &[u8]) -> Result<Self> {
112        const FIXED: usize = 8 + 8 + 8 + 4 + 8 + 4 + 8 + 4 * 4 + 4;
113        let invalid = || Error::Io(std::io::Error::new(
114            std::io::ErrorKind::InvalidData, "prepared graft manifest is invalid"));
115        if bytes.len() < FIXED || bytes.get(..8) != Some(Self::MAGIC) { return Err(invalid()); }
116        let crc_at = bytes.len() - 4;
117        let want = u32::from_le_bytes(bytes[crc_at..].try_into().unwrap());
118        if crc32c::crc32c(&bytes[..crc_at]) != want { return Err(invalid()); }
119        let mut pos = 8usize;
120        let mut u64_at = || {
121            let end = pos.checked_add(8).ok_or_else(invalid)?;
122            let value = u64::from_le_bytes(bytes.get(pos..end).ok_or_else(invalid)?.try_into().unwrap());
123            pos = end; Ok::<_, Error>(value)
124        };
125        let base_generation = u64_at()?;
126        let write_generation = u64_at()?;
127        drop(u64_at);
128        let root = u32::from_le_bytes(bytes.get(pos..pos + 4).ok_or_else(invalid)?.try_into().unwrap());
129        pos += 4;
130        let rows = u64::from_le_bytes(bytes.get(pos..pos + 8).ok_or_else(invalid)?.try_into().unwrap());
131        pos += 8;
132        let last_next = u32::from_le_bytes(bytes.get(pos..pos + 4).ok_or_else(invalid)?.try_into().unwrap());
133        pos += 4;
134        let inserted_rows = u64::from_le_bytes(bytes.get(pos..pos + 8).ok_or_else(invalid)?.try_into().unwrap());
135        pos += 8;
136        let mut fields = Vec::with_capacity(4);
137        for _ in 0..4 {
138            let raw = bytes.get(pos..pos + 4).ok_or_else(invalid)?;
139            pos += 4;
140            let len = u32::from_le_bytes(raw.try_into().unwrap()) as usize;
141            let end = pos.checked_add(len).ok_or_else(invalid)?;
142            if end > crc_at { return Err(invalid()); }
143            fields.push(bytes[pos..end].to_vec());
144            pos = end;
145        }
146        if pos != crc_at || fields.iter().any(Vec::is_empty) { return Err(invalid()); }
147        Ok(Self { base_generation, write_generation, root, rows,
148            min: fields.remove(0), max: fields.remove(0), inserted_min: fields.remove(0),
149            inserted_max: fields.remove(0), inserted_rows, last_next })
150    }
151
152    /// Durably replace a candidate descriptor. A torn next descriptor remains
153    /// at `.tmp`; reopen reads only the last fsynced-and-renamed manifest.
154    pub fn write_manifest(&self, path: &Path) -> Result<()> {
155        let parent = path.parent().unwrap_or_else(|| Path::new("."));
156        std::fs::create_dir_all(parent)?;
157        let tmp = path.with_extension("tmp");
158        {
159            use std::io::Write;
160            let mut file = std::fs::File::create(&tmp)?;
161            let bytes = self.encode()?;
162            file.write_all(&bytes)?;
163            crate::write_stats::add(
164                crate::write_stats::Phase::Manifest,
165                bytes.len() as u64,
166            );
167            file.sync_all()?;
168        }
169        std::fs::rename(&tmp, path)?;
170        std::fs::File::open(parent)?.sync_all()?;
171        Ok(())
172    }
173
174    pub fn read_manifest(path: &Path) -> Result<Self> {
175        let metadata = std::fs::metadata(path)?;
176        if metadata.len() > (1 << 20) {
177            return Err(Error::Io(std::io::Error::new(
178                std::io::ErrorKind::InvalidData, "prepared graft manifest is too large")));
179        }
180        Self::decode(&std::fs::read(path)?)
181    }
182}
183
184/// The blessed default (D24): 64 MiB pool, buffered I/O, Normal durability.
185/// 64 MiB is the WEIGHTED-TRAINING budget (D25): every gate and ablation is
186/// measured at this budget or smaller, so optimization pressure lands on the
187/// disk path -- layout, clustering, read counts -- never on cache flatter.
188/// A server that can spend 4 GB sets it; the layout never changes.
189impl Default for Config {
190    fn default() -> Config {
191        Config { budget_bytes: 64 << 20, io: IoMode::Buffered, sync: SyncMode::Normal }
192    }
193}
194
195/// The `Barrier` a `SyncMode` maps to for the buffer pool's own flush. A free
196/// function, not just an inherent method on `Store`, because `build`'s fresh
197/// path needs it before a `Store` exists to call a method on.
198fn sync_barrier(sync: SyncMode) -> Barrier {
199    match sync {
200        SyncMode::Full => Barrier::Full,
201        SyncMode::Normal => Barrier::Data,
202        SyncMode::Off => Barrier::None,
203    }
204}
205
206pub struct Store {
207    pool: BufferPool,
208    /// None = snapshot reader (2f): opened read-only at a published
209    /// generation; every mutating method refuses with Error::ReadOnly.
210    wal: Option<Wal>,
211    root: u32,
212    /// Directory this store lives in -- lets `nearest_par` (2g.2) open
213    /// sibling snapshot readers on the same files.
214    dir: std::path::PathBuf,
215    /// Publication counter (2f): the generation of the newest meta slot on
216    /// disk. The next checkpoint publishes generation + 1.
217    generation: u64,
218    /// A snapshot reader's registration in the reader table (2n): present
219    /// iff this store is a reader; dropping it releases the pin.
220    reader_slot: Option<crate::readers::ReaderSlot>,
221    tree_id: u16,
222    sync: SyncMode,
223    io_mode: IoMode,
224    #[cfg(feature = "test-support")]
225    test_faults: Arc<TestFaultInjector>,
226    /// Set once, never cleared within this instance's life, when a logged
227    /// write or `checkpoint` reaches a fallible durability step and fails.
228    /// Checked at the top of EVERY writer -- `put`, `delete`, `commit`,
229    /// `checkpoint`, `bulk_load` -- because the one that matters most is
230    /// `checkpoint`: it ends in `Wal::rotate`, which deletes the log. See
231    /// `Error::StorePoisoned`'s doc comment for why a failed barrier leaves
232    /// the store unable to say what is durable, and why the remedy is to
233    /// drop this instance and reopen (which clears the flag by rebuilding
234    /// every belief from disk) rather than to retry in place.
235    poisoned: bool,
236    /// The append hint, hoisted here from `BTree` (Task 16). `BTree::open` is
237    /// called fresh on every `put`/`delete`/`get`/`scan` and dropped
238    /// immediately, so a hint owned by that transient `BTree` would reset to
239    /// `None` before the next call ever saw it -- the fast path added in
240    /// Task 15 exercised only a long-lived `BTree` in its own tests and
241    /// bought nothing through `Store`, the path production and the
242    /// benchmark actually use. `Store` outlives every one of those calls, so
243    /// the hint lives here instead -- PostgreSQL's `rel->rd_targblock` on the
244    /// `Relation`, not on a per-statement scan. `Cell`, not a plain field,
245    /// because `BTree<'p>` needs `&'p Cell<Option<u32>>` and `put`/`delete`
246    /// already borrow `self` mutably elsewhere; interior mutability avoids
247    /// fighting the borrow checker over a single `u32` hint.
248    last_leaf: Cell<Option<u32>>,
249    /// Times `insert` used `last_leaf` instead of a full descent, borrowed
250    /// into each transient `BTree` the same way `last_leaf` is so the count
251    /// survives across `Store::put` calls instead of resetting with every
252    /// fresh `BTree`.
253    fast_path_hits: Cell<u64>,
254    /// Times `insert` found the hint armed and tried `fast_path_leaf` at
255    /// all, hit or miss (Task 18). Same borrowed-`Cell` reasoning as
256    /// `fast_path_hits`: a counter living on the transient `BTree` would
257    /// reset every call, and the disarm regression test needs a count of
258    /// attempts, not hits, that survives across `Store::put` calls to prove
259    /// a random-order workload pays for one wasted probe and then stops --
260    /// not one wasted probe per row forever.
261    fast_path_attempts: Cell<u64>,
262    /// The per-keyspace append hints (`TagHints`), owned here for exactly the
263    /// reason `last_leaf` is: every `Store` operation builds a transient
264    /// `BTree` and drops it, so a cache living on the tree would be empty on
265    /// arrival every time. Every `BTree` this `Store` opens gets it -- the
266    /// ones that write, so they can use it, and the ones that delete or graft,
267    /// so they can clear it.
268    tag_hints: crate::btree::TagHints,
269    /// Superblock logical version this file declared (1 = plain cells, 2 =
270    /// compact cells). Independent of this build's `compact-cells` feature;
271    /// checkpoint writes this value, never the compile-time default.
272    format_version: u16,
273    #[cfg(test)]
274    trace: Vec<&'static str>,
275    #[cfg(test)]
276    barriers: Vec<&'static str>,
277}
278
279#[cfg(test)]
280#[path = "store_reuse_probe.rs"]
281mod reuse_probe;
282
283impl Store {
284    fn build(dir: &Path, cfg: Config, fresh: bool) -> Result<Store> {
285        Self::build_limited(dir, cfg, fresh, None)
286    }
287    fn build_limited(dir: &Path, cfg: Config, fresh: bool, limits: Option<crate::limits::ResourceLimits>) -> Result<Store> {
288        if fresh && dir.join("data").exists() {
289            return Err(std::io::Error::new(std::io::ErrorKind::AlreadyExists, "create refuses existing data; use open").into());
290        }
291        std::fs::create_dir_all(dir)?;
292        // A refused same-process open normally reports the writer lock. If an
293        // external fault has already made that process's WAL demonstrably
294        // corrupt, preserve the more specific pre-existing corruption
295        // refusal. Inspect only in this already-refused local case, never on
296        // a successful-open path and never for another process whose WAL may
297        // still be changing.
298        if !fresh && crate::io::writer_owned_by_this_process(&dir.join("data")) {
299            if let crate::wal::Stop::Damaged { offset, why } =
300                crate::wal::Wal::inspect(&dir.join("wal"), cfg.io)?.stop
301            {
302                return Err(crate::Error::CorruptWal { offset, why });
303            }
304        }
305        let (file, io_mode) = open_file_writer(&dir.join("data"), cfg.io)?;
306        Self::build_on_limited(dir, cfg, fresh, file.into(), io_mode, limits)
307    }
308
309    /// `build` with the data file handed in rather than opened here. The
310    /// split exists for the same reason `recover_impl`'s does: a test needs
311    /// to drive a store whose disk fails in a specific way, and without a
312    /// seam the poisoning TRIGGER is unreachable from any test -- which is
313    /// exactly how it came to be shipped unpinned (Task 17 final review,
314    /// F6). Nothing but the `open_file` call moves out of the path a test
315    /// can reach.
316    #[cfg(test)]
317    fn build_on(dir: &Path, cfg: Config, fresh: bool, file: Arc<dyn FileIo>, io_mode: IoMode)
318        -> Result<Store> {
319        Self::build_on_limited(dir, cfg, fresh, file, io_mode, None)
320    }
321    fn build_on_limited(dir: &Path, cfg: Config, fresh: bool, file: Arc<dyn FileIo>, io_mode: IoMode,
322        requested: Option<crate::limits::ResourceLimits>) -> Result<Store> {
323        let dir_owned = dir.to_path_buf();
324        // Two thirds of the budget to the pool; the rest is scratch and log.
325        let frames = (cfg.budget_bytes / 3 * 2) / PAGE_SIZE;
326        let budget = Arc::new(MemoryBudget::new(cfg.budget_bytes));
327        let pool = BufferPool::new(file, budget, frames.max(16))?;
328        let limits = if fresh { requested } else { Meta::read_limits(&pool)? };
329        if let Some(l) = limits { pool.set_resource_limits(l)?; }
330        let mut wal = Wal::open_limited(&dir.join("wal"), cfg.io, limits.map(|l| l.wal_bytes))?;
331
332        // Declared before the `fresh`/reopen branch because `BTree::create`
333        // below borrows them for its lifetime, and they must still be here,
334        // outside that borrow, to move into the `Store` literal afterward.
335        let last_leaf = Cell::new(None);
336        let fast_path_hits = Cell::new(0);
337        let fast_path_attempts = Cell::new(0);
338        let tag_hints = crate::btree::TagHints::default();
339
340        let (root, generation, format_version) = if fresh {
341            let _meta_page = pool.allocate()?;     // page 0 is the superblock
342            drop(_meta_page);
343            let _slot_b = pool.allocate()?;        // page 1 is meta slot B (2f)
344            drop(_slot_b);
345            Meta::init_slot_b(&pool)?;
346            pool.set_compact_cells(crate::meta::FORMAT_VERSION == 2);
347            let t = BTree::create(&pool, 1, &last_leaf, &fast_path_hits, &fast_path_attempts)?;
348            let r = t.root();
349            Meta { format_version: crate::meta::FORMAT_VERSION, roots: [r, 0, 0, 0, 0, 0, 0, 0], next_lsn: wal.next_lsn(),
350                   generation: 0 }
351                .write(&pool)?;
352            pool.flush_all(sync_barrier(cfg.sync))?;
353            // The data and WAL contents now have the requested barriers; make
354            // their new directory entries durable before create returns and a
355            // later commit can be acknowledged. This is unconditional because
356            // even SyncMode::Off does not ask for files to vanish by name.
357            pool.sync_dir()?;
358            (r, 0, crate::meta::FORMAT_VERSION)
359        } else {
360            // Re-supply the LSN high-water mark BEFORE anything else touches
361            // the log: `Wal::open` above has already derived `next_lsn` from
362            // whatever the log file itself contains, which after a rotation
363            // is empty and would renumber from 1. `set_lsn_floor` is a no-op
364            // if the log turned out to know a higher number already (i.e. no
365            // rotation happened since the last checkpoint).
366            let meta = Meta::read_latest(&pool)?;
367            let format_version = meta.format_version & !crate::meta::LIMITED;
368            pool.set_compact_cells(format_version == 2);
369            wal.set_lsn_floor(meta.next_lsn);
370            // Opening can recreate a missing WAL. Publish that filename once
371            // before any later acknowledgement, rather than relying on a
372            // recurring ordinary-checkpoint directory barrier.
373            pool.sync_dir()?;
374            (meta.roots[0], meta.generation, format_version)
375        };
376        // 2f: everything on disk up to here is the published state a snapshot
377        // reader may be standing on. Freeze it BEFORE recovery replays the
378        // WAL tail -- the replay goes through the ordinary shadowed write
379        // path, so it relocates instead of tearing the published pages.
380        pool.set_frozen_boundary();
381        let write_generation = generation.checked_add(1).ok_or(crate::Error::Corrupt {
382            page_no: if generation % 2 == 0 { crate::meta::META_PAGE } else { crate::meta::META_PAGE_B },
383            why: "published generation is exhausted",
384        })?;
385        pool.set_stamp_gen(write_generation);
386        if !fresh {
387            // Derived state may be lost, but never allocate from an unchecked
388            // file length. At most one retirement per physical page.
389            let cap = limits.map_or(28 + 24 * pool.page_count() as u64, |l| l.freelist_bytes());
390            if let Ok(f) = std::fs::File::open(dir.join("free")) {
391                use std::io::Read;
392                if f.metadata()?.len() <= cap {
393                    let mut b = Vec::new();
394                    f.take(cap + 1).read_to_end(&mut b)?;
395                    if b.len() as u64 <= cap { pool.import_free(&b, generation); }
396                }
397            }
398        }
399        pool.set_reuse_limit(
400            generation.saturating_sub(1)
401                .min(crate::readers::oldest_live_reader(dir)));
402
403        let mut s = Store { pool, wal: Some(wal), root, generation, reader_slot: None, dir: dir_owned, tree_id: 1, sync: cfg.sync, io_mode,
404                            #[cfg(feature = "test-support")] test_faults: Arc::new(TestFaultInjector::default()),
405                            poisoned: false, last_leaf, fast_path_hits, fast_path_attempts,
406                            tag_hints,
407                            format_version,
408                            #[cfg(test)] trace: Vec::new(),
409                            #[cfg(test)] barriers: Vec::new() };
410        if !fresh {
411            if limits.is_some() {
412                // Limited commits publish their root without a WAL Commit.
413                // Reject a foreign committed tail before applying anything.
414                if s.wal.as_ref().unwrap().committed_end()? != 0 {
415                    return Err(Error::ResourceLimit("unexpected committed WAL in constrained store; preserve and inspect"));
416                }
417                s.wal_mut()?.rotate()?;
418            } else { s.recover_from_log()?; }
419        }
420        Ok(s)
421    }
422
423    /// 2f: open a SNAPSHOT READER on `dir`. Serves the newest PUBLISHED
424    /// generation (the last checkpoint's roots) and nothing later: the WAL
425    /// tail is deliberately not replayed -- replay writes, and this store
426    /// cannot write. Correct beside a live writer with zero coordination:
427    /// published pages are immutable (the writer shadows instead of editing,
428    /// and page numbers are never reused), so everything reachable from a
429    /// published root stays byte-identical for as long as this reader lives.
430    /// Sacrifice (Law 4): staleness up to one checkpoint cadence.
431    pub fn open_snapshot(dir: &Path, cfg: Config) -> Result<Store> {
432        Self::open_snapshot_with_after_meta(dir, cfg, |_| Ok(()))
433    }
434
435    /// Test seam inside snapshot open: metadata has been selected while a
436    /// conservative generation-zero reader slot is already live. Production
437    /// supplies an empty closure; the concurrency regression test holds this
438    /// point open to prove the writer cannot recycle the selected generation.
439    fn open_snapshot_with_after_meta<F>(dir: &Path, cfg: Config, after_meta: F) -> Result<Store>
440    where
441        F: FnOnce(u64) -> Result<()>,
442    {
443        // Register BEFORE selecting metadata. Zero is deliberately
444        // conservative: until the exact published generation is known, the
445        // writer must assume this reader can need every historical page.
446        let mut slot = crate::readers::ReaderSlot::reserve(dir)?;
447        let file = crate::io::open_file_readonly(&dir.join("data"))?;
448        let frames = (cfg.budget_bytes / 3 * 2) / PAGE_SIZE;
449        let budget = Arc::new(MemoryBudget::new(cfg.budget_bytes));
450        let pool = BufferPool::new(file.into(), budget, frames.max(16))?;
451        let meta = Meta::read_latest(&pool)?;
452        let format_version = meta.format_version & !crate::meta::LIMITED;
453        pool.set_compact_cells(format_version == 2);
454        if let Some(l) = Meta::read_limits(&pool)? { pool.set_resource_limits(l)?; }
455        after_meta(meta.generation)?;
456        slot.set_generation(meta.generation)?;
457        let s = Store {
458            pool, wal: None, root: meta.roots[0], generation: meta.generation,
459            reader_slot: Some(slot),
460            dir: dir.to_path_buf(),
461            tree_id: 1, sync: cfg.sync, io_mode: IoMode::Buffered,
462            #[cfg(feature = "test-support")]
463            test_faults: Arc::new(TestFaultInjector::default()),
464            poisoned: false,
465            last_leaf: Cell::new(None),
466            fast_path_hits: Cell::new(0), fast_path_attempts: Cell::new(0),
467            tag_hints: crate::btree::TagHints::default(),
468            format_version,
469            #[cfg(test)] trace: Vec::new(),
470            #[cfg(test)] barriers: Vec::new(),
471        };
472        Ok(s)
473    }
474
475    /// The pool, for probes and tests (2n).
476    pub fn pool_ref(&self) -> &BufferPool { &self.pool }
477
478    /// Create a new constrained entry store. Commit publishes durable metadata;
479    /// external-sort/graft and in-place repair require a separate workspace.
480    /// Existing directories are refused, so policy cannot be applied halfway.
481    pub fn create_limited(dir: &Path, cfg: Config, limits: crate::limits::ResourceLimits) -> Result<Store> {
482        let limits = limits.validate()?;
483        std::fs::create_dir(dir)?;
484        Self::build_limited(dir, cfg, true, Some(limits))
485    }
486    pub fn resource_limits(&self) -> Option<crate::limits::ResourceLimits> { self.pool.resource_limits() }
487    fn refuse_external_workspace(&self) -> Result<()> {
488        if self.resource_limits().is_some() {
489            return Err(Error::ResourceLimit("external sort/graft requires a separately budgeted destination; use batched put"));
490        }
491        Ok(())
492    }
493
494    pub fn create(dir: &Path, cfg: Config) -> Result<Store> { Self::build(dir, cfg, true) }
495
496    /// `create`, but on a caller-supplied data file. Test-only: the one
497    /// caller is the poisoning-trigger test, which needs a `FileIo` whose
498    /// barrier can be made to fail on demand.
499    #[cfg(test)]
500    fn create_on(dir: &Path, cfg: Config, file: Arc<dyn FileIo>) -> Result<Store> {
501        std::fs::create_dir_all(dir)?;
502        Self::build_on(dir, cfg, true, file, IoMode::Buffered)
503    }
504    pub fn open(dir: &Path, cfg: Config) -> Result<Store> {
505        if dir.join("data").exists() { Self::build(dir, cfg, false) }
506        else { Self::build(dir, cfg, true) }
507    }
508
509    pub fn io_mode(&self) -> IoMode { self.io_mode }
510
511    #[cfg(feature = "test-support")]
512    #[doc(hidden)]
513    pub fn test_fault_injector(&self) -> Arc<TestFaultInjector> {
514        self.test_faults.clone()
515    }
516
517    /// Replay committed transactions the pages have not yet absorbed.
518    fn recover_from_log(&mut self) -> Result<()> {
519        // Pass 1 names the final committed byte. Pass 2 applies one frame at a
520        // time through that boundary. Nothing after it happened, and neither a
521        // whole WAL nor a whole transaction is ever resident (Law 1).
522        let committed_end = self.wal.as_ref().ok_or(crate::Error::ReadOnly)?.committed_end()?;
523        let mut off = 0u64;
524        while off < committed_end {
525            let (_, kind, payload, next) = self.wal.as_ref()
526                .ok_or(crate::Error::ReadOnly)?
527                .record_at(off, committed_end)?
528                .ok_or(crate::Error::CorruptWal {
529                    offset: off, why: "committed recovery prefix ended early",
530                })?;
531            self.apply(kind, &payload, off)?;
532            off = next;
533        }
534        Ok(())
535    }
536
537    fn apply(&mut self, kind: RecKind, payload: &[u8], wal_offset: u64) -> Result<()> {
538        let corrupt = |why| crate::Error::CorruptWal { offset: wal_offset, why };
539        match kind {
540            RecKind::Put => {
541                let klen = payload.get(..2)
542                    .map(|b| u16::from_le_bytes([b[0], b[1]]) as usize)
543                    .ok_or_else(|| corrupt("put payload has no key length"))?;
544                let key_end = 2usize.checked_add(klen)
545                    .ok_or_else(|| corrupt("put key boundary overflow"))?;
546                let key = payload.get(2..key_end)
547                    .ok_or_else(|| corrupt("put key crosses its WAL payload"))?;
548                let val = payload.get(key_end..)
549                    .ok_or_else(|| corrupt("put value boundary is invalid"))?;
550                // Salvage publishes a marker at generation zero and retains
551                // the logical WAL. Replay repairs source rows, but old derived
552                // counts cannot be trusted beside a potentially damaged source.
553                // Fresh never-checkpointed databases have no marker.
554                if self.generation == 0 && crate::keys::is_field_aggregate_key(key) {
555                    if Meta::is_salvaged(&self.pool)? { return Ok(()); }
556                }
557                let mut t = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
558                t.insert(key, val)?;
559                self.root = t.root();
560            }
561            RecKind::PutEmptyBatch => {
562                if payload.len() > crate::page::MAX_RECORD_LEN {
563                    return Err(corrupt("put-empty-batch exceeds the WAL frame bound"));
564                }
565                let count = payload.get(..2)
566                    .map(|b| u16::from_le_bytes([b[0], b[1]]) as usize)
567                    .ok_or_else(|| corrupt("put-empty-batch payload has no key count"))?;
568                if count == 0 || count > PUT_EMPTY_BATCH_MAX_KEYS {
569                    return Err(corrupt("put-empty-batch key count is outside its fixed bound"));
570                }
571
572                // Validate every boundary and key-size invariant before the
573                // first tree mutation. A CRC-valid malformed frame must not
574                // apply a valid prefix and fail halfway through.
575                let mut at = 2usize;
576                for _ in 0..count {
577                    let len_end = at.checked_add(2)
578                        .ok_or_else(|| corrupt("put-empty-batch length boundary overflow"))?;
579                    let len_bytes = payload.get(at..len_end)
580                        .ok_or_else(|| corrupt("put-empty-batch key has no length"))?;
581                    let key_len = u16::from_le_bytes([len_bytes[0], len_bytes[1]]) as usize;
582                    let key_end = len_end.checked_add(key_len)
583                        .ok_or_else(|| corrupt("put-empty-batch key boundary overflow"))?;
584                    payload.get(len_end..key_end)
585                        .ok_or_else(|| corrupt("put-empty-batch key crosses its WAL payload"))?;
586                    if 4 + key_len + 12 > crate::page::MAX_RECORD_LEN {
587                        return Err(corrupt("put-empty-batch key cannot fit a leaf record"));
588                    }
589                    at = key_end;
590                }
591                if at != payload.len() {
592                    return Err(corrupt("put-empty-batch payload has trailing bytes"));
593                }
594
595                let mut t = BTree::open(&self.pool, self.tree_id, self.root,
596                    &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
597                at = 2;
598                for _ in 0..count {
599                    let key_len = u16::from_le_bytes([payload[at], payload[at + 1]]) as usize;
600                    at += 2;
601                    t.insert(&payload[at..at + key_len], &[])?;
602                    at += key_len;
603                }
604                self.root = t.root();
605            }
606            RecKind::Delete => {
607                let klen = payload.get(..2)
608                    .map(|b| u16::from_le_bytes([b[0], b[1]]) as usize)
609                    .ok_or_else(|| corrupt("delete payload has no key length"))?;
610                let key_end = 2usize.checked_add(klen)
611                    .ok_or_else(|| corrupt("delete key boundary overflow"))?;
612                if key_end != payload.len() {
613                    return Err(corrupt("delete key length does not match its WAL payload"));
614                }
615                let mut t = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
616                t.delete(&payload[2..key_end])?;
617                self.root = t.root();
618            }
619            RecKind::DeletePrefix => {
620                if payload.is_empty() {
621                    return Err(corrupt("delete-prefix WAL payload is empty"));
622                }
623                let mut t = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
624                t.delete_prefix(payload)?;
625                self.root = t.root();
626            }
627            RecKind::Commit => {
628                if !payload.is_empty() {
629                    return Err(corrupt("commit WAL payload is not empty"));
630                }
631            }
632            RecKind::PageImage => {
633                return Err(corrupt("page-image WAL records have no recovery implementation"));
634            }
635        }
636        Ok(())
637    }
638
639    /// [klen u16][key][value...to end]. The value's length is implicit -- the
640    /// WAL frame already carries its own length. The old format stored vlen as
641    /// u16, which silently TRUNCATED any value over 64KB on its way into the
642    /// log (v.len() as u16); with overflow chains making big values legal,
643    /// that was a data-loss bug waiting on the replay path.
644    fn frame(k: &[u8], v: &[u8]) -> Vec<u8> {
645        let mut b = Vec::with_capacity(2 + k.len() + v.len());
646        b.extend_from_slice(&(k.len() as u16).to_le_bytes());
647        b.extend_from_slice(k);
648        b.extend_from_slice(v);
649        b
650    }
651
652    /// The writer's log, or ReadOnly for a snapshot reader (2f).
653    fn wal_mut(&mut self) -> Result<&mut Wal> {
654        self.wal.as_mut().ok_or(crate::Error::ReadOnly)
655    }
656
657    pub fn put(&mut self, k: &[u8], v: &[u8]) -> Result<()> {
658        #[cfg(feature = "test-support")]
659        self.test_faults.check(TestFaultKind::Write)?;
660        if self.poisoned { return Err(crate::Error::StorePoisoned); }
661        if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
662        let wal_payload_len = 2usize.checked_add(k.len())
663            .and_then(|n| n.checked_add(v.len()))
664            .ok_or(crate::Error::TooLarge)?;
665        if self.resource_limits().is_some_and(|l| v.len() > l.record_bytes as usize) {
666            return Err(Error::ResourceLimit("record exceeds configured maximum"));
667        }
668        // Reject before `frame` copies the change. D7's large overflow values
669        // retain the encoded format's nearly 4 GiB range, while no disk length
670        // may exceed the exact maximum any writer can emit.
671        if wal_payload_len as u64 > crate::wal::MAX_PAYLOAD_BYTES {
672            return Err(crate::Error::TooLarge);
673        }
674        if k.len() > crate::page::MAX_RECORD_LEN - 16 || v.len() > u32::MAX as usize {
675            return Err(crate::Error::TooLarge);
676        }
677        #[cfg(feature = "write-trace")]
678        let frame_started = crate::write_trace::active().then(std::time::Instant::now);
679        let payload = Self::frame(k, v);
680        #[cfg(feature = "write-trace")]
681        if let Some(started) = frame_started {
682            crate::write_trace::add(crate::write_trace::Field::FrameEncode, started.elapsed());
683            crate::write_trace::value_copy();
684        }
685        // Reject BEFORE the WAL ever sees it (Task 17 re-review, R1): a frame
686        // the insert below would refuse must never enter the log, or replay
687        // meets a record nothing can apply. With overflow chains the value can
688        // be any size; what must still fit in a page is the KEY (a marker
689        // record is 4 + klen + 12 bytes), and the value length must fit the
690        // marker's u32.
691        if 4 + k.len() + 12 > crate::page::MAX_RECORD_LEN || v.len() > u32::MAX as usize {
692            return Err(crate::Error::TooLarge);
693        }
694        if let Err(e) = self.wal_mut()?.append(RecKind::Put, &payload) {
695            // An append error can leave a partial frame in the log buffer or
696            // file. The tree is still intact, but continuing could place a
697            // later commit beyond an unreadable tail, so this is not benign.
698            self.poisoned = true;
699            return Err(e);
700        }
701        let mut t = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
702        #[cfg(feature = "write-trace")]
703        let btree_started = crate::write_trace::active().then(std::time::Instant::now);
704        let inserted = t.insert(k, v);
705        #[cfg(feature = "write-trace")]
706        if let Some(started) = btree_started {
707            crate::write_trace::add(crate::write_trace::Field::BtreeTotal, started.elapsed());
708        }
709        let root = t.root();
710        drop(t);
711        if let Err(e) = inserted {
712            self.poisoned = true;
713            return Err(e);
714        }
715        self.root = root;
716        Ok(())
717    }
718
719    /// The exact ordinary put path, with feature-gated clocks enabled for the
720    /// load diagnostic. This does not exist in production builds.
721    #[cfg(feature = "write-trace")]
722    pub fn put_profiled(&mut self, k: &[u8], v: &[u8]) -> (Result<()>, crate::write_trace::PutTrace) {
723        let trace_started = crate::write_trace::begin();
724        let result = self.put(k, v);
725        let trace = crate::write_trace::finish(trace_started);
726        (result, trace)
727    }
728
729    /// Write up to 64 empty-valued index rows under one WAL frame and one
730    /// B-tree handle. Key order is preserved exactly. This is not bulk-load:
731    /// it appends to the existing shared tree and participates in the caller's
732    /// ordinary commit/recovery boundary.
733    pub fn put_empty_batch(&mut self, keys: &[Vec<u8>]) -> Result<()> {
734        if keys.is_empty() { return Ok(()); }
735        #[cfg(feature = "test-support")]
736        self.test_faults.check(TestFaultKind::Write)?;
737        if self.poisoned { return Err(crate::Error::StorePoisoned); }
738        if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
739        if keys.len() > PUT_EMPTY_BATCH_MAX_KEYS { return Err(crate::Error::TooLarge); }
740
741        let payload_len = keys.iter().try_fold(2usize, |total, key| {
742            if 4 + key.len() + 12 > crate::page::MAX_RECORD_LEN || key.len() > u16::MAX as usize {
743                return None;
744            }
745            total.checked_add(2 + key.len())
746        }).ok_or(crate::Error::TooLarge)?;
747        if payload_len > crate::page::MAX_RECORD_LEN { return Err(crate::Error::TooLarge); }
748        let mut payload = Vec::with_capacity(payload_len);
749        payload.extend_from_slice(&(keys.len() as u16).to_le_bytes());
750        for key in keys {
751            payload.extend_from_slice(&(key.len() as u16).to_le_bytes());
752            payload.extend_from_slice(key);
753        }
754        if let Err(e) = self.wal_mut()?.append(RecKind::PutEmptyBatch, &payload) {
755            self.poisoned = true;
756            return Err(e);
757        }
758        let mut t = BTree::open(&self.pool, self.tree_id, self.root,
759            &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
760        let inserted = keys.iter().try_for_each(|key| t.insert(key, &[]));
761        let root = t.root();
762        drop(t);
763        if let Err(e) = inserted {
764            self.poisoned = true;
765            return Err(e);
766        }
767        self.root = root;
768        Ok(())
769    }
770
771    pub fn delete(&mut self, k: &[u8]) -> Result<bool> {
772        #[cfg(feature = "test-support")]
773        self.test_faults.check(TestFaultKind::Write)?;
774        if self.poisoned { return Err(crate::Error::StorePoisoned); }
775        if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
776        if k.len() > crate::page::MAX_RECORD_LEN - 2 { return Ok(false); }
777        let mut payload = (k.len() as u16).to_le_bytes().to_vec();
778        payload.extend_from_slice(k);
779        // The same pre-WAL bound as `put`, on the PAYLOAD -- `2 + k.len()`
780        // -- not on the key (Task 17 final review, F1: a guard written
781        // against the key while the append writes the payload is a guard
782        // that is two bytes wrong, and a bound that is two bytes wrong is
783        // the whole defect). No key this long can be present in the first
784        // place: `put` above refuses any record whose key alone would reach
785        // this size, so `Ok(false)` -- `delete`'s ordinary "not found"
786        // answer -- is not merely convenient here, it is the true one.
787        if payload.len() > crate::page::MAX_RECORD_LEN { return Ok(false); }
788        if let Err(e) = self.wal_mut()?.append(RecKind::Delete, &payload) {
789            self.poisoned = true;
790            return Err(e);
791        }
792        let mut t = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
793        let deleted = t.delete(k);
794        let root = t.root();
795        drop(t);
796        let hit = match deleted {
797            Ok(hit) => hit,
798            Err(e) => {
799                self.poisoned = true;
800                return Err(e);
801            }
802        };
803        self.root = root;
804        Ok(hit)
805    }
806
807    /// Delete every key starting with `prefix` (2h A4). One WAL record,
808    /// idempotent on replay; leaves wholly inside the range are cleared in
809    /// ONE page write instead of per-slot removals -- the fold's head erase
810    /// was 8M row deletes (~20 minutes at 1M docs) and is now ~one write
811    /// per leaf. Returns the number of keys removed.
812    pub fn delete_prefix(&mut self, prefix: &[u8]) -> Result<u64> {
813        #[cfg(feature = "test-support")]
814        self.test_faults.check(TestFaultKind::Write)?;
815        if self.poisoned { return Err(crate::Error::StorePoisoned); }
816        if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
817        if prefix.is_empty() || prefix.len() > crate::page::MAX_RECORD_LEN { return Err(crate::Error::TooLarge); }
818        if let Err(e) = self.wal_mut()?.append(RecKind::DeletePrefix, prefix) {
819            self.poisoned = true;
820            return Err(e);
821        }
822        let mut t = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
823        let deleted = t.delete_prefix(prefix);
824        let root = t.root();
825        drop(t);
826        let n = match deleted {
827            Ok(n) => n,
828            Err(e) => {
829                self.poisoned = true;
830                return Err(e);
831            }
832        };
833        self.root = root;
834        Ok(n)
835    }
836
837    pub fn get(&self, k: &[u8]) -> Result<Option<Vec<u8>>> {
838        if self.poisoned && self.resource_limits().is_some() {
839            return Err(crate::Error::StorePoisoned);
840        }
841        #[cfg(feature = "test-support")]
842        self.test_faults.check(TestFaultKind::Read)?;
843        BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints).get(k)
844    }
845
846    pub fn scan(&self, from: &[u8]) -> Result<RangeIter<'_>> {
847        if self.poisoned && self.resource_limits().is_some() {
848            return Err(crate::Error::StorePoisoned);
849        }
850        BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints).range(from)
851    }
852
853    /// Descending scan of keys strictly below `to`.
854    pub fn scan_reverse(&self, to: &[u8]) -> Result<crate::btree::ReverseRangeIter<'_>> {
855        if self.poisoned && self.resource_limits().is_some() {
856            return Err(crate::Error::StorePoisoned);
857        }
858        BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf,
859                    &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints).range_reverse(to)
860    }
861
862    /// `SyncMode` is a promise about what reached the medium, so the three modes
863    /// must issue three different things. An earlier draft had `Full` and
864    /// `Normal` both call one `sync()` -- which made the label decorative, and a
865    /// decorative durability label is worse than none, because every benchmark
866    /// carrying it becomes unattributable. It is worse than that on macOS
867    /// specifically: `std::fs::File::sync_data` (what the single `sync()` used)
868    /// issues `fcntl(F_FULLFSYNC)` there, so `Normal` would have silently been
869    /// `Full`'s ~65x-costlier barrier on this machine while meaning something
870    /// cheaper on Linux -- the same code looking wildly different speeds for
871    /// reasons the label never states.
872    ///
873    /// Nothing in the committed suite failed if this collapsed back into one
874    /// call for both arms -- `barriers` exists so something does.
875    pub fn commit(&mut self) -> Result<()> {
876        if self.resource_limits().is_some() { return self.checkpoint(); }
877        #[cfg(feature = "test-support")]
878        self.test_faults.check(TestFaultKind::Commit)?;
879        if self.poisoned { return Err(crate::Error::StorePoisoned); }
880        let wal = self.wal.as_mut().ok_or(crate::Error::ReadOnly)?;
881        wal.append(RecKind::Commit, &[])?;
882        // Buffering does not redefine commit: Off promises no BARRIER, not
883        // that the record stays in process RAM. Free under Normal/Full (their
884        // barriers flush anyway).
885        wal.flush()?;
886        match self.sync {
887            SyncMode::Full => {
888                #[cfg(test)] self.barriers.push("sync_full");
889                wal.sync_full()?
890            }
891            SyncMode::Normal => {
892                #[cfg(test)] self.barriers.push("sync_data");
893                wal.sync_data()?
894            }
895            SyncMode::Off => {}
896        }
897        // A snapshot can close between checkpoints. Refresh at the existing
898        // transaction boundary so the next batch can use newly eligible pages
899        // instead of growing for the rest of this epoch. No publication or
900        // page scan: only the live reader table, with the same conservative
901        // fallback/ambiguous-reader protection as checkpoint and reopen.
902        let readers = crate::readers::live_generations(&self.dir);
903        self.pool.refresh_reuse(self.generation, readers.as_deref());
904        Ok(())
905    }
906
907    /// Commit with an explicit byte-based publication policy. Limits are
908    /// triggers checked at the transaction boundary, NOT disk quotas. Both
909    /// WAL bytes and allocated/shadow pages matter: logical WAL records are
910    /// much smaller than the pages scattered updates cause us to copy.
911    ///
912    /// On the checkpoint branch, the data + metadata barriers establish the
913    /// commit's durability; a redundant WAL barrier is avoided. Checkpoint
914    /// errors remain errors and preserve the WAL for reopening. This opt-in
915    /// API does not change ordinary commit() or snapshot staleness defaults.
916    pub fn commit_with_checkpoint(&mut self, wal_bytes: u64, page_bytes: u64) -> Result<bool> {
917        if self.resource_limits().is_some() { self.commit()?; return Ok(true); }
918        if self.poisoned { return Err(crate::Error::StorePoisoned); }
919        let end = self.wal.as_ref().ok_or(crate::Error::ReadOnly)?.end_offset();
920        let publish = (wal_bytes > 0 && end >= wal_bytes)
921            || (page_bytes > 0 && self.pool.epoch_allocated_bytes() >= page_bytes);
922        if !publish { self.commit()?; return Ok(false); }
923        #[cfg(feature = "test-support")]
924        self.test_faults.check(TestFaultKind::Commit)?;
925        let wal = self.wal_mut()?;
926        if let Err(e) = wal.append(RecKind::Commit, &[]).and_then(|_| wal.flush()) {
927            self.poisoned = true;
928            return Err(e);
929        }
930        if let Err(e) = self.checkpoint() {
931            self.poisoned = true;
932            return Err(e);
933        }
934        Ok(true)
935    }
936
937    #[cfg(test)]
938    fn barriers(&self) -> Vec<&'static str> { self.barriers.clone() }
939
940    /// The exact primitive `SyncMode::Full` issues on this platform, so a
941    /// measurement can name it instead of implying it.
942    pub fn dir(&self) -> &Path { &self.dir }
943    /// The published root of the main tree. For structural verification of a
944    /// live database — see `verify::verify_published_tree`.
945    pub fn published_root(&self) -> u32 { self.root }
946    /// The main tree's identity, so a verifier can prove every page belongs to it.
947    pub fn main_tree_id(&self) -> u16 { self.tree_id }
948
949    pub fn sync_full_primitive(&self) -> &'static str {
950        self.wal.as_ref().map_or("none (snapshot reader)", |w| w.sync_full_primitive())
951    }
952
953    /// The pool's own counters, including how many `sync_data`/`sync_full`
954    /// barriers it has actually issued. Real observability API, not test
955    /// instrumentation -- `PoolStats` and `BufferPool::stats` are already
956    /// public and ungated -- so this needs no `#[cfg(test)]` and the test
957    /// that reads it can live in `kernel/tests/durability.rs` like any other
958    /// public-API test.
959    pub fn pool_stats(&self) -> crate::pool::PoolStats { self.pool.stats() }
960    /// The per-keyspace append hints, for diagnostics and for the oracle that
961    /// runs one workload with them and once without.
962    pub fn tag_hints(&self) -> &crate::btree::TagHints { &self.tag_hints }
963    /// Data-file counters only; the WAL keeps its own.
964    pub fn io_stats(&self) -> Option<&crate::io::IoStats> { self.pool.io_stats() }
965    pub fn sweep_steps(&self) -> u64 { self.pool.sweep_steps() }
966
967    /// The `Barrier` `checkpoint`'s page flush should use, derived from this
968    /// store's `SyncMode`. Before this, the pool's checkpoint flush went
969    /// through `FileIo::sync` -- `File::sync_data`, `F_FULLFSYNC` on macOS,
970    /// `fdatasync` on Linux -- ungoverned by `SyncMode` at all: the exact
971    /// ambiguity `commit` was fixed to no longer have, left standing in the
972    /// quieter of the two places that used to share it.
973    fn barrier(&self) -> Barrier { sync_barrier(self.sync) }
974
975    /// Runtime durability change (SQL SET WAL_SYNC): applies to every
976    /// subsequent commit/checkpoint barrier.
977    pub fn set_sync(&mut self, s: SyncMode) { self.sync = s; }
978
979    /// The level in force right now. A caller that raises durability for one
980    /// operation has to be able to put back what it found -- without this it
981    /// could only guess, and guessing wrong silently re-levels every later
982    /// commit.
983    pub fn sync_mode(&self) -> SyncMode { self.sync }
984    /// The published generation this handle currently serves.
985    pub fn generation(&self) -> u64 { self.generation }
986
987    pub fn checkpoint(&mut self) -> Result<()> {
988        let result = self.checkpoint_inner();
989        if result.is_err() && self.wal.is_some() { self.poisoned = true; }
990        result
991    }
992    fn checkpoint_inner(&mut self) -> Result<()> {
993        if self.poisoned { return Err(crate::Error::StorePoisoned); }
994        if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
995        // Preflight both numbers before flushing any page or touching either
996        // meta slot. Publishing `u64::MAX` would leave the following epoch
997        // with no distinct stamp; wrapping it to zero would make recycled
998        // pages look older than every live snapshot.
999        let gen = self.generation.checked_add(1).ok_or(crate::Error::Corrupt {
1000            page_no: if self.generation % 2 == 0 { crate::meta::META_PAGE } else { crate::meta::META_PAGE_B },
1001            why: "published generation is exhausted",
1002        })?;
1003        let next_write_generation = gen.checked_add(1).ok_or(crate::Error::Corrupt {
1004            page_no: if self.generation % 2 == 0 { crate::meta::META_PAGE } else { crate::meta::META_PAGE_B },
1005            why: "no generation remains for the next write epoch",
1006        })?;
1007        // The LSN high-water mark goes into the same page as the roots, so it
1008        // is flushed atomically with them below, and a reopen after this
1009        // checkpoint's rotation re-supplies it via `set_lsn_floor` instead of
1010        // renumbering from 1.
1011        #[cfg(test)] self.trace.clear();
1012        // Poison on `Err` rather than merely propagating it (Task 17
1013        // re-review, R5): `flush_all` clears each frame's dirty bit BEFORE
1014        // issuing its barrier, so a failure here leaves the store believing
1015        // pages are durable that may not be -- and the very next line
1016        // deletes the log those pages' contents are otherwise only recorded
1017        // in. See `Error::StorePoisoned`.
1018        if let Err(e) = self.pool.flush_all(self.barrier()) {
1019            self.poisoned = true;
1020            return Err(e);
1021        }
1022        #[cfg(test)] { self.trace.push("flush_pages"); self.trace.push("sync_file"); }
1023        // 2f dual-slot flip, strictly AFTER the data barrier above: the slot
1024        // write must not be able to reach the medium before the pages its
1025        // roots name. A crash between the two leaves the previous generation
1026        // standing and the WAL tail replayable -- exactly a missed
1027        // checkpoint, never a torn one.
1028        Meta {
1029            format_version: self.format_version,
1030            roots: [self.root, 0, 0, 0, 0, 0, 0, 0],
1031            next_lsn: self.wal.as_ref().ok_or(crate::Error::ReadOnly)?.next_lsn(),
1032            generation: gen,
1033        }.write_slot(&self.pool)?;
1034        if let Err(e) = self.pool.flush_all(self.barrier()) {
1035            self.poisoned = true;
1036            return Err(e);
1037        }
1038        self.generation = gen;
1039        self.pool.set_stamp_gen(next_write_generation);
1040        // The new root is durable, making the preceding generation's hint
1041        // stale. Reuse its existing file: a torn overwrite loses only derived
1042        // reuse knowledge. Generation and whole-body CRC are checked on reopen.
1043        let readers = crate::readers::live_generations(&self.dir);
1044        self.pool.refresh_reuse(gen, readers.as_deref());
1045        let free_bytes = self.pool.export_free(gen);
1046        let _ = crate::verify::persist_checkpoint_freelist(
1047            &self.dir,
1048            &free_bytes,
1049            self.pool.file_ref(),
1050        );
1051        // 2n recycling horizon: pages freed at generations <= this are safe.
1052        // published - 1 keeps the dual-slot fallback tree intact; the oldest
1053        // live snapshot reader caps it further. An unreadable reader table
1054        // reports 0 -- recycling halts rather than guesses.
1055
1056        #[cfg(test)] self.trace.push("flip_meta");
1057        // Everything just published is now immutable to writers (2f): the
1058        // next epoch's writes shadow instead of editing in place.
1059        self.pool.set_frozen_boundary();
1060        // Data/WAL names were published at create/open; no name is replaced at
1061        // ordinary checkpoint. Newly created free hints sync their own name.
1062        self.wal_mut()?.rotate_published()?;
1063        #[cfg(test)] self.trace.push("rotate_wal");
1064        Ok(())
1065    }
1066
1067    #[cfg(test)]
1068    fn checkpoint_trace(&self) -> Vec<&'static str> { self.trace.clone() }
1069
1070    /// The log's own high-water mark, test-only. Exists to pin the LSN
1071    /// persistence obligation: without `checkpoint` writing `next_lsn` into
1072    /// `Meta` and `open` re-supplying it via `set_lsn_floor`, a rotation
1073    /// followed by a reopen would renumber records from 1 -- silently
1074    /// harmless today because nothing compares a page's `lsn` against a
1075    /// record's yet, but wrong the moment that check is added.
1076    #[cfg(test)]
1077    fn next_lsn(&self) -> u64 { self.wal.as_ref().unwrap().next_lsn() }
1078
1079    /// Build the tree from scratch by external sort and pack.
1080    ///
1081    /// The new tree replaces the old one, but the old one's PAGES are not
1082    /// reclaimed -- page reclamation is deferred in Phase 1, so calling this
1083    /// on a non-empty store leaves the previous tree's pages allocated and
1084    /// unreachable. Intended for loading into a fresh store; on a populated
1085    /// one the file grows by the size of both trees.
1086    ///
1087    /// Unlike `put`/`delete`, this writes pages straight into the pool and
1088    /// never through the WAL -- logging every bulk-loaded record before
1089    /// packing it would reintroduce the per-record write this task exists to
1090    /// remove. That means there is nothing in the log for recovery to replay,
1091    /// so this makes itself durable before returning rather than leaving that
1092    /// to a caller who may reasonably assume `bulk_load` behaves like any
1093    /// other write the engine accepted: pack, then `checkpoint()`. Without
1094    /// this, `bulk_load` followed only by `commit()` followed by a crash
1095    /// loses everything -- the superblock still names the OLD root, and the
1096    /// log holds a commit record describing nothing -- while the caller was
1097    /// told `Ok` twice. `checkpoint()` also rotates the log, so this discards
1098    /// prior log state; correct for a whole-tree replacement, not something
1099    /// an incremental write may do.
1100    ///
1101    /// SACRIFICE (Law 4): publication adds one sequential read traversal of
1102    /// the packed tree. It retains only 16 verification pages and tree-height
1103    /// state; ordinary queries and writes do not use this path.
1104    pub fn bulk_load<I>(&mut self, items: I) -> Result<()>
1105    where I: Iterator<Item = (Vec<u8>, Vec<u8>)> {
1106        self.bulk_load_with_before_publish(items, |_, _| Ok(()))
1107    }
1108
1109    /// Test seam at the exact Law-3 boundary: packing has returned a fresh
1110    /// root, but that root is not yet authoritative.
1111    fn bulk_load_with_before_publish<I, F>(&mut self, items: I, before_publish: F) -> Result<()>
1112    where
1113        I: Iterator<Item = (Vec<u8>, Vec<u8>)>,
1114        F: FnOnce(&Path, u32) -> Result<()>,
1115    {
1116        // Guarded like every other writer (Task 17 final review, F5): this
1117        // is the write that replaces the whole tree AND, via `checkpoint`,
1118        // discards the log -- the last thing a store with an ambiguous
1119        // checkpoint behind it should be allowed to do.
1120        self.refuse_external_workspace()?;
1121        if self.poisoned { return Err(crate::Error::StorePoisoned); }
1122        if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
1123        // A directory keyed only on the process id would be shared by every
1124        // `bulk_load` call live in this process at once -- and the default
1125        // test harness runs the tests in `kernel/tests/bulk.rs` concurrently
1126        // on separate threads of the SAME process, each calling this method.
1127        // Two concurrent `ExternalSort`s in one directory collide on their
1128        // `run-NNNNN.tmp` names, and `SortedRuns::drop`'s `remove_dir_all`
1129        // would delete run files a sibling call is still merging from. The
1130        // sequence number makes each call's scratch directory its own.
1131        //
1132        // `pack_tree` spills each tree level's separators into this same
1133        // directory rather than inventing a second temp-file scheme with its
1134        // own lifetime and its own collision rules, so the guarantee this
1135        // sequence number provides covers the pack's level files too.
1136        static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1137        let seq = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1138        let tmp = std::env::temp_dir()
1139            .join(format!("kernel-sort-{}-{}", std::process::id(), seq));
1140        let mut s = crate::bulk::ExternalSort::new(&tmp, 64 << 20)?;
1141        let mut expected_rows = 0u64;
1142        for (k, v) in items {
1143            expected_rows = expected_rows.checked_add(1).ok_or(crate::Error::TooLarge)?;
1144            // Oversized values spill to overflow chains HERE, before the sort,
1145            // so the sorted stream and the packed leaves only ever carry
1146            // page-sized records. Chain pages are allocated from the same pool
1147            // the pack writes through; the marker is what gets sorted.
1148            if 4 + k.len() + v.len() > crate::page::MAX_RECORD_LEN {
1149                let (head, crc) = crate::btree::write_overflow(&self.pool, &v)?;
1150                let m = crate::btree::enc_marker(v.len() as u32, head, crc);
1151                s.push_flagged(k, m.to_vec(), true)?;
1152            } else {
1153                s.push(k, v)?;
1154            }
1155        }
1156        let mut runs = s.finish()?;
1157        // `&tmp` is the sort's own scratch directory: `pack_tree` spills its
1158        // level separators there and removes them as it consumes them, and
1159        // `SortedRuns::drop` takes the directory itself when `runs` falls out
1160        // of scope below.
1161        let root = crate::bulk::pack_tree(&self.pool, self.tree_id, runs.iter()?, 0.9, &tmp)?;
1162        // Make the candidate a physical file, then reopen it through a fresh,
1163        // fixed-size pool. The old root remains authoritative throughout.
1164        if let Err(e) = self.pool.flush_all(Barrier::None) {
1165            self.poisoned = true;
1166            return Err(e);
1167        }
1168        let data = self.dir.join("data");
1169        before_publish(&data, root)?;
1170        let verified = crate::verify::verify_file(
1171            &data,
1172            self.io_mode,
1173            root,
1174            self.tree_id,
1175            expected_rows,
1176        )?;
1177        debug_assert_eq!(verified.rows, expected_rows);
1178        debug_assert!(verified.pages > 0);
1179        self.root = root;
1180        // The old hint names a leaf number that means nothing in the
1181        // repacked tree -- it may not even be allocated, or may now belong
1182        // to an unrelated page. `fast_path_leaf`'s checks (tree_id, kind,
1183        // rightmost, room, ordering) would probably catch a reused page
1184        // number too, but "probably" is not the standard: a stale hint here
1185        // has no reason to survive `bulk_load`, so it is cleared rather than
1186        // left to be caught.
1187        self.last_leaf.set(None);
1188        self.tag_hints.clear();
1189        self.checkpoint()
1190    }
1191
1192    /// Sort and pack an empty key interval, independently reopen and verify
1193    /// it, then splice it into the shared tree through a copy-on-write parent
1194    /// path. Packed pages bypass the per-record WAL because they are
1195    /// unreachable until the checkpoint publishes the new root.
1196    ///
1197    /// The sort arena is fixed at 64 MiB. Inputs larger than that spill
1198    /// checksummed runs, so RAM is independent of corpus size; the named cost
1199    /// is scratch space approximately twice the index size plus one readback
1200    /// verification pass. Packed leaves are 90% full, leaving room for the
1201    /// first later live write before ordinary split policy takes over.
1202    pub fn graft_range<I>(&mut self, items: I) -> Result<()>
1203    where
1204        I: Iterator<Item = (Vec<u8>, Vec<u8>)>,
1205    {
1206        self.graft_range_with_before_publish(items, |_, _| Ok(()))
1207    }
1208
1209    /// Test seam after candidate pages have reached the data file but before
1210    /// the independent reopen and the one logical child replacement.
1211    fn graft_range_with_before_publish<I, F>(
1212        &mut self,
1213        items: I,
1214        before_publish: F,
1215    ) -> Result<()>
1216    where
1217        I: Iterator<Item = (Vec<u8>, Vec<u8>)>,
1218        F: FnOnce(&Path, &crate::bulk::PackedRange) -> Result<()>,
1219    {
1220        self.refuse_external_workspace()?;
1221        if self.poisoned { return Err(crate::Error::StorePoisoned); }
1222        if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
1223
1224        static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1225        let seq = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1226        let tmp = std::env::temp_dir()
1227            .join(format!("kernel-graft-{}-{}", std::process::id(), seq));
1228        let mut sort = crate::bulk::ExternalSort::new(&tmp, 64 << 20)?;
1229        let mut expected_rows = 0u64;
1230        let mut min: Option<Vec<u8>> = None;
1231        let mut max: Option<Vec<u8>> = None;
1232
1233        for (key, value) in items {
1234            expected_rows = expected_rows.checked_add(1).ok_or(crate::Error::TooLarge)?;
1235            if min.as_ref().is_none_or(|current| key.as_slice() < current.as_slice()) {
1236                min = Some(key.clone());
1237            }
1238            if max.as_ref().is_none_or(|current| key.as_slice() > current.as_slice()) {
1239                max = Some(key.clone());
1240            }
1241            if 4 + key.len() + value.len() > crate::page::MAX_RECORD_LEN {
1242                if value.len() > u32::MAX as usize || 4 + key.len() + 12 > crate::page::MAX_RECORD_LEN {
1243                    return Err(crate::Error::TooLarge);
1244                }
1245                let (head, crc) = crate::btree::write_overflow(&self.pool, &value)?;
1246                let marker = crate::btree::enc_marker(value.len() as u32, head, crc);
1247                sort.push_flagged(key, marker.to_vec(), true)?;
1248            } else {
1249                sort.push(key, value)?;
1250            }
1251        }
1252        let (Some(min), Some(max)) = (min, max) else {
1253            // An empty graft changes nothing and must not manufacture an empty
1254            // child that the verifier correctly refuses below a parent.
1255            return Ok(());
1256        };
1257
1258        let mut runs = sort.finish()?;
1259        self.graft_sorted_range_with_before_publish(
1260            runs.iter()?, expected_rows, min, max, &tmp, true, before_publish)
1261    }
1262
1263    /// Publish a caller's already-sorted, checksummed run stream. This is the
1264    /// lower half of [`graft_range`]: late-index builders that already paid for
1265    /// an external sort use it directly instead of materialising or sorting the
1266    /// same keys a second time.
1267    ///
1268    /// `expected_rows`, `min`, and `max` are trusted only for the preflight
1269    /// overlap probe. `pack_range` recomputes all three and this method refuses
1270    /// a mismatch before publication.
1271    pub fn graft_sorted_range<I>(
1272        &mut self,
1273        sorted: I,
1274        expected_rows: u64,
1275        min: Vec<u8>,
1276        max: Vec<u8>,
1277        scratch_dir: &Path,
1278    ) -> Result<()>
1279    where
1280        I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>,
1281    {
1282        self.graft_sorted_range_with_before_publish(
1283            sorted, expected_rows, min, max, scratch_dir, true, |_, _| Ok(()))
1284    }
1285
1286    /// Install an independently verified packed range in this writer's
1287    /// unpublished root. The caller must finish its other namespace changes
1288    /// and issue one checkpoint; until then snapshot readers continue to use
1289    /// the previous generation. This is the transaction form of
1290    /// [`graft_sorted_range`](Self::graft_sorted_range).
1291    pub fn graft_sorted_range_deferred<I>(
1292        &mut self,
1293        sorted: I,
1294        expected_rows: u64,
1295        min: Vec<u8>,
1296        max: Vec<u8>,
1297        scratch_dir: &Path,
1298    ) -> Result<()>
1299    where
1300        I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>,
1301    {
1302        self.graft_sorted_range_with_before_publish(
1303            sorted, expected_rows, min, max, scratch_dir, false, |_, _| Ok(()))
1304    }
1305
1306    /// Pack and independently verify a candidate without publishing it. The
1307    /// returned descriptor is self-checksummable and sufficient for a later
1308    /// process to publish the already-written pages without sorting or packing
1309    /// again. The standing root is unchanged on every return path.
1310    pub fn prepare_graft_candidate<I>(
1311        &mut self,
1312        sorted: I,
1313        expected_rows: u64,
1314        min: Vec<u8>,
1315        max: Vec<u8>,
1316        scratch_dir: &Path,
1317    ) -> Result<PreparedGraft>
1318    where
1319        I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>,
1320    {
1321        self.refuse_external_workspace()?;
1322        if self.poisoned { return Err(Error::StorePoisoned); }
1323        if self.wal.is_none() { return Err(Error::ReadOnly); }
1324        if expected_rows == 0 || min > max { return Err(Error::TooLarge); }
1325        if let Some(row) = self.scan(&min)?.next() {
1326            let (key, _) = row?;
1327            if key <= max { return Err(Error::RangeNotEmpty); }
1328        }
1329        let tree = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf,
1330            &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
1331        let boundary = tree.plan_graft(&min)?;
1332        drop(tree);
1333        let last_next = boundary.right_page.unwrap_or(boundary.old_next);
1334        let packed = crate::bulk::pack_range(
1335            &self.pool, self.tree_id, sorted, 0.9, scratch_dir, last_next)?;
1336        if packed.rows != expected_rows || packed.min.as_deref() != Some(min.as_slice())
1337            || packed.max.as_deref() != Some(max.as_slice()) {
1338            return Err(Error::Corrupt { page_no: packed.root,
1339                why: "prepared range disagrees with its sorted-stream manifest" });
1340        }
1341        let tree = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf,
1342            &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
1343        let candidate = tree.build_graft_candidate(&boundary, &packed)?;
1344        drop(tree);
1345        if let Err(error) = self.pool.flush_all(Barrier::None) {
1346            self.poisoned = true;
1347            return Err(error);
1348        }
1349        let write_generation = self.generation.checked_add(1).ok_or(Error::Corrupt {
1350            page_no: 0, why: "prepared graft generation is exhausted" })?;
1351        crate::verify::verify_range_file_generation(
1352            &self.dir.join("data"), self.io_mode, candidate.root, self.tree_id,
1353            candidate.rows, &candidate.min, &candidate.max, candidate.last_next,
1354            Some(write_generation))?;
1355        Ok(PreparedGraft {
1356            base_generation: self.generation,
1357            write_generation,
1358            root: candidate.root,
1359            rows: candidate.rows,
1360            min: candidate.min,
1361            max: candidate.max,
1362            last_next: candidate.last_next,
1363            inserted_min: min,
1364            inserted_max: max,
1365            inserted_rows: expected_rows,
1366        })
1367    }
1368
1369    /// Verify and publish an already-packed candidate discovered on reopen.
1370    /// Calling it after the candidate's checkpoint but before scratch cleanup
1371    /// is idempotent: the published generation and exact inserted interval are
1372    /// checked, then no second root flip occurs.
1373    pub fn publish_existing_candidate(&mut self, prepared: &PreparedGraft) -> Result<()> {
1374        self.refuse_external_workspace()?;
1375        if self.poisoned { return Err(Error::StorePoisoned); }
1376        if self.wal.is_none() { return Err(Error::ReadOnly); }
1377        if prepared.write_generation != prepared.base_generation.checked_add(1)
1378            .ok_or(Error::TooLarge)? || prepared.inserted_min > prepared.inserted_max
1379            || prepared.inserted_rows == 0 {
1380            return Err(Error::Corrupt { page_no: prepared.root,
1381                why: "prepared graft has an invalid generation or interval" });
1382        }
1383        let data = self.dir.join("data");
1384        crate::verify::verify_range_file_generation(
1385            &data, self.io_mode, prepared.root, self.tree_id, prepared.rows,
1386            &prepared.min, &prepared.max, prepared.last_next,
1387            Some(prepared.write_generation))?;
1388
1389        if self.generation == prepared.write_generation {
1390            let mut rows = 0u64;
1391            self.scan(&prepared.inserted_min)?.for_each_ref(|key, _| {
1392                if key > prepared.inserted_max.as_slice() { return false; }
1393                rows = rows.saturating_add(1);
1394                true
1395            })?;
1396            if rows != prepared.inserted_rows {
1397                return Err(Error::Corrupt { page_no: prepared.root,
1398                    why: "published graft interval disagrees with its manifest" });
1399            }
1400            return Ok(());
1401        }
1402        if self.generation != prepared.base_generation {
1403            return Err(Error::Corrupt { page_no: prepared.root,
1404                why: "prepared graft belongs to a stale base generation" });
1405        }
1406        if let Some(row) = self.scan(&prepared.inserted_min)?.next() {
1407            let (key, _) = row?;
1408            if key <= prepared.inserted_max { return Err(Error::RangeNotEmpty); }
1409        }
1410        let mut tree = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf,
1411            &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
1412        let boundary = tree.plan_existing_graft(&prepared.inserted_min)?;
1413        let retired = tree.install_graft(&boundary, prepared.root, &prepared.min)?;
1414        self.root = tree.root();
1415        drop(tree);
1416        for page in retired { self.pool.free_page(page)?; }
1417        self.last_leaf.set(None);
1418        self.tag_hints.clear();
1419        self.checkpoint()
1420    }
1421
1422    fn graft_sorted_range_with_before_publish<I, F>(
1423        &mut self,
1424        sorted: I,
1425        expected_rows: u64,
1426        min: Vec<u8>,
1427        max: Vec<u8>,
1428        scratch_dir: &Path,
1429        publish: bool,
1430        before_publish: F,
1431    ) -> Result<()>
1432    where
1433        I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>,
1434        F: FnOnce(&Path, &crate::bulk::PackedRange) -> Result<()>,
1435    {
1436        let trace = std::env::var_os("SEKEJAP_LOAD_BREAKDOWN").is_some();
1437        let total_started = trace.then(std::time::Instant::now);
1438        self.refuse_external_workspace()?;
1439        if self.poisoned {
1440            return Err(crate::Error::StorePoisoned);
1441        }
1442        if self.wal.is_none() {
1443            return Err(crate::Error::ReadOnly);
1444        }
1445        if expected_rows == 0 {
1446            return Ok(());
1447        }
1448        if min > max {
1449            return Err(crate::Error::TooLarge);
1450        }
1451
1452        let preflight_started = trace.then(std::time::Instant::now);
1453        if let Some(row) = self.scan(&min)?.next() {
1454            let (key, _) = row?;
1455            if key <= max { return Err(crate::Error::RangeNotEmpty); }
1456        }
1457
1458        let tree = BTree::open(
1459            &self.pool,
1460            self.tree_id,
1461            self.root,
1462            &self.last_leaf,
1463            &self.fast_path_hits,
1464            &self.fast_path_attempts,
1465        ).with_tags(&self.tag_hints);
1466        let boundary = tree.plan_graft(&min)?;
1467        drop(tree);
1468        let last_next = boundary.right_page.unwrap_or(boundary.old_next);
1469        let preflight =
1470            preflight_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1471
1472        let pack_started = trace.then(std::time::Instant::now);
1473        let packed = crate::bulk::pack_range(
1474            &self.pool,
1475            self.tree_id,
1476            sorted,
1477            0.9,
1478            scratch_dir,
1479            last_next,
1480        )?;
1481        let pack = pack_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1482        if packed.rows != expected_rows
1483            || packed.min.as_deref() != Some(min.as_slice())
1484            || packed.max.as_deref() != Some(max.as_slice())
1485        {
1486            return Err(crate::Error::Corrupt {
1487                page_no: packed.root,
1488                why: "packed range disagrees with its sorted-stream manifest",
1489            });
1490        }
1491        if packed.last_leaf >= self.pool.page_count() {
1492            return Err(crate::Error::Corrupt {
1493                page_no: packed.last_leaf,
1494                why: "packed range returned a last leaf outside the data file",
1495            });
1496        }
1497
1498        let tree = BTree::open(
1499            &self.pool,
1500            self.tree_id,
1501            self.root,
1502            &self.last_leaf,
1503            &self.fast_path_hits,
1504            &self.fast_path_attempts,
1505        ).with_tags(&self.tag_hints);
1506        let candidate = tree.build_graft_candidate(&boundary, &packed)?;
1507        drop(tree);
1508
1509        // Materialise candidate bytes, then verify through a different file
1510        // handle and a 16-page pool. The authoritative root is still unchanged.
1511        let flush_started = trace.then(std::time::Instant::now);
1512        if let Err(error) = self.pool.flush_all(Barrier::None) {
1513            self.poisoned = true;
1514            return Err(error);
1515        }
1516        let flush = flush_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1517        let data = self.dir.join("data");
1518        before_publish(&data, &packed)?;
1519        let verify_started = trace.then(std::time::Instant::now);
1520        let verified = crate::verify::verify_range_file(
1521            &data,
1522            self.io_mode,
1523            candidate.root,
1524            self.tree_id,
1525            candidate.rows,
1526            &candidate.min,
1527            &candidate.max,
1528            candidate.last_next,
1529        )?;
1530        let verify = verify_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1531        debug_assert_eq!(verified.rows, candidate.rows);
1532        debug_assert!(verified.pages > 0);
1533
1534        // Only after verification: build fresh copies of the O(height) parent
1535        // path and publish their root. Old pages are merely retired for a
1536        // future safe generation; snapshot readers retain their byte-stable
1537        // versions throughout.
1538        let install_started = trace.then(std::time::Instant::now);
1539        let mut tree = BTree::open(
1540            &self.pool,
1541            self.tree_id,
1542            self.root,
1543            &self.last_leaf,
1544            &self.fast_path_hits,
1545            &self.fast_path_attempts,
1546        ).with_tags(&self.tag_hints);
1547        let retired = tree.install_graft(&boundary, candidate.root, &candidate.min)?;
1548        self.root = tree.root();
1549        drop(tree);
1550        for page in retired { self.pool.free_page(page)?; }
1551        self.last_leaf.set(None);
1552        self.tag_hints.clear();
1553        let install =
1554            install_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1555        let publish_started = trace.then(std::time::Instant::now);
1556        let result = if publish { self.checkpoint() } else { Ok(()) };
1557        let publication =
1558            publish_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1559        if trace {
1560            eprintln!("graft detail: rows={} pages={} bytes={} preflight={:.6}s pack={:.6}s flush={:.6}s verify={:.6}s graft={:.6}s publish={:.6}s total={:.6}s deferred={}",
1561                expected_rows, verified.pages, verified.pages as u64 * crate::page::PAGE_SIZE as u64,
1562                preflight.as_secs_f64(), pack.as_secs_f64(), flush.as_secs_f64(),
1563                verify.as_secs_f64(), install.as_secs_f64(), publication.as_secs_f64(),
1564                total_started.unwrap().elapsed().as_secs_f64(), !publish);
1565        }
1566        result
1567    }
1568}
1569
1570// Tests that need `checkpoint_trace`/`barriers`/`next_lsn` live here, inside
1571// the crate, where `#[cfg(test)]` actually applies. An integration test in
1572// kernel/tests/ links the library built WITHOUT `--cfg test` -- that cfg only
1573// governs the crate's own unit-test binary -- so a test-only method gated
1574// this way would compile fine and simply not exist from there. The
1575// public-API durability tests (kernel/tests/durability.rs) need nothing
1576// test-only and stay where a real consumer of this crate would exercise them.
1577#[cfg(test)]
1578mod tests {
1579    use super::*;
1580
1581    fn cfg() -> Config { Config { budget_bytes: 32 << 20, io: IoMode::Buffered, sync: SyncMode::Full } }
1582
1583    /// The inherited kernel Store superblock still stamps `FORMAT_VERSION` from
1584    /// the `compact-cells` cargo feature. A release must accept both supported
1585    /// versions and write the file's own on checkpoint, not the build's.
1586    /// (This path is not the typed-collection release backend — that is
1587    /// PageWalStore with per-database feature bits — but the same defect.)
1588    #[test]
1589    fn a_store_checkpoint_preserves_the_file_own_format_version() {
1590        let d = tempfile::tempdir().unwrap();
1591        let other = if crate::meta::FORMAT_VERSION == 2 { 1 } else { 2 };
1592        {
1593            let s = Store::create(d.path(), cfg()).unwrap();
1594            let meta = Meta::read_latest(&s.pool).unwrap();
1595            Meta {
1596                format_version: other,
1597                roots: meta.roots,
1598                next_lsn: meta.next_lsn,
1599                generation: meta.generation,
1600            }
1601            .write(&s.pool)
1602            .unwrap();
1603            s.pool.flush_all(crate::io::Barrier::Data).unwrap();
1604        }
1605        let mut s = Store::open(d.path(), cfg()).unwrap_or_else(|e| {
1606            panic!("supported superblock version {other} must open in this build: {e}")
1607        });
1608        s.put(b"k", b"v").unwrap();
1609        s.commit().unwrap();
1610        s.checkpoint().unwrap();
1611        drop(s);
1612        let s = Store::open(d.path(), cfg()).unwrap();
1613        let got = Meta::read_latest(&s.pool).unwrap().format_version & !crate::meta::LIMITED;
1614        assert_eq!(
1615            got, other,
1616            "checkpoint must write the file's format version, not the build's FORMAT_VERSION"
1617        );
1618        assert_eq!(s.get(b"k").unwrap().as_deref(), Some(&b"v"[..]));
1619    }
1620
1621    #[test]
1622    fn byte_policy_publishes_exact_rows_without_a_redundant_wal_barrier() {
1623        let d = tempfile::tempdir().unwrap();
1624        let mut s = Store::create(d.path(), cfg()).unwrap();
1625        s.put(b"old", b"durable").unwrap(); s.commit().unwrap(); s.checkpoint().unwrap();
1626        let reader = Store::open_snapshot(d.path(), cfg()).unwrap();
1627        s.barriers.clear();
1628        s.put(b"new", &[7; 9000]).unwrap();
1629        let before = s.pool_stats().sync_full_calls;
1630        assert!(s.commit_with_checkpoint(u64::MAX, 4096).unwrap());
1631        assert_eq!(s.pool_stats().sync_full_calls - before, 2);
1632        assert!(s.barriers.is_empty(), "publication already supplies the durability barriers");
1633        assert_eq!(s.wal.as_ref().unwrap().end_offset(), 0);
1634        assert_eq!(reader.get(b"new").unwrap(), None);
1635        drop(s);
1636        let s = Store::open(d.path(), cfg()).unwrap();
1637        assert_eq!(s.get(b"old").unwrap().as_deref(), Some(&b"durable"[..]));
1638        assert_eq!(s.get(b"new").unwrap(), Some(vec![7; 9000]));
1639    }
1640
1641    #[test]
1642    fn byte_policy_keeps_wal_and_old_snapshot_after_data_barrier_failure() {
1643        let d = tempfile::tempdir().unwrap();
1644        let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1645        let fio = Arc::new(FailingBarrier { inner: real, fail: false.into() });
1646        let mut s = Store::create_on(d.path(), cfg(), fio.clone()).unwrap();
1647        s.put(b"old", b"committed").unwrap(); s.commit().unwrap(); s.checkpoint().unwrap();
1648        let reader = Store::open_snapshot(d.path(), cfg()).unwrap();
1649        s.put(b"new", &[8; 9000]).unwrap();
1650        fio.fail.store(true, std::sync::atomic::Ordering::Relaxed);
1651        assert!(s.commit_with_checkpoint(1,1).is_err());
1652        assert!(s.wal.as_ref().unwrap().end_offset() > 0);
1653        assert!(matches!(s.commit_with_checkpoint(1,1), Err(crate::Error::StorePoisoned)));
1654        assert_eq!(reader.get(b"old").unwrap().as_deref(), Some(&b"committed"[..]));
1655        assert_eq!(reader.get(b"new").unwrap(), None);
1656        fio.fail.store(false, std::sync::atomic::Ordering::Relaxed);
1657        drop(s); drop(fio);
1658        let s = Store::open(d.path(), cfg()).unwrap();
1659        assert_eq!(s.get(b"old").unwrap().as_deref(), Some(&b"committed"[..]));
1660        assert_eq!(s.get(b"new").unwrap(), Some(vec![8; 9000]));
1661    }
1662
1663    #[test]
1664    fn constrained_commit_failure_preserves_published_state_and_snapshot() {
1665        for barrier in [false, true] {
1666            let d = tempfile::tempdir().unwrap();
1667            let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1668            let fio = Arc::new(FailingBarrier { inner: real, fail: false.into() });
1669            let limits = crate::limits::ResourceLimits { data_bytes: 1 << 20, wal_bytes: 64 << 10,
1670                tracked_pages: 256, readers: 2, record_bytes: 16000, recovery_bytes: 65536 };
1671            let mut s = Store::build_on_limited(d.path(), cfg(), true, fio.clone(), IoMode::Buffered, Some(limits)).unwrap();
1672            s.put(b"old", b"durable").unwrap(); s.commit().unwrap();
1673            let reader = Store::open_snapshot(d.path(), cfg()).unwrap();
1674            s.put(b"new", &[8;9000]).unwrap();
1675            if barrier { fio.fail.store(true, std::sync::atomic::Ordering::Relaxed); }
1676            if barrier {
1677                assert!(s.commit().is_err());
1678                assert!(matches!(s.checkpoint(), Err(Error::StorePoisoned)));
1679            } // otherwise simulate process exit before commit
1680            assert_eq!(reader.get(b"old").unwrap(), Some(b"durable".to_vec()));
1681            assert_eq!(reader.get(b"new").unwrap(), None);
1682            fio.fail.store(false, std::sync::atomic::Ordering::Relaxed);
1683            drop(s); drop(fio);
1684            let s = Store::open(d.path(), cfg()).unwrap();
1685            assert_eq!(s.get(b"old").unwrap(), Some(b"durable".to_vec()));
1686            assert_eq!(s.get(b"new").unwrap(), None);
1687        }
1688    }
1689
1690    #[test]
1691    fn constrained_commit_data_write_failures_do_not_publish_partial_rows() {
1692        for at in 0..3 {
1693            let d = tempfile::tempdir().unwrap();
1694            let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1695            let fio = Arc::new(FailingWrite { inner: real, armed: std::sync::Mutex::new(None) });
1696            let limits = crate::limits::ResourceLimits { data_bytes: 1 << 20, wal_bytes: 64 << 10,
1697                tracked_pages: 256, readers: 2, record_bytes: 16000, recovery_bytes: 65536 };
1698            let mut s = Store::build_on_limited(d.path(), cfg(), true, fio.clone(), IoMode::Buffered, Some(limits)).unwrap();
1699            s.put(b"old", b"durable").unwrap(); s.commit().unwrap();
1700            let reader = Store::open_snapshot(d.path(), cfg()).unwrap();
1701            s.put(b"new", &[8;12000]).unwrap();
1702            *fio.armed.lock().unwrap() = Some(at);
1703            assert!(s.commit().is_err());
1704            assert!(matches!(s.put(b"later", b"no"), Err(Error::StorePoisoned)));
1705            assert_eq!(reader.get(b"old").unwrap(), Some(b"durable".to_vec()));
1706            *fio.armed.lock().unwrap() = None;
1707            drop(s); drop(fio);
1708            let s = Store::open(d.path(), cfg()).unwrap();
1709            assert_eq!(s.get(b"old").unwrap(), Some(b"durable".to_vec()));
1710            assert_eq!(s.get(b"new").unwrap(), None);
1711        }
1712    }
1713
1714    /// F5c: this is a real reader/writer race with a handshake, not a reader
1715    /// registered before single-threaded churn. The reader has selected
1716    /// generation 1 and is stopped inside snapshot open. While it is stopped,
1717    /// the writer publishes enough epochs to make generation-1 pages eligible
1718    /// for recycling, then overwrites them. A correct open registers a
1719    /// conservative pin before selecting metadata, so the writer can never
1720    /// advance past the reader while the hook is held.
1721    #[test]
1722    fn a_snapshot_is_registered_before_its_generation_can_be_recycled() {
1723        let d = tempfile::tempdir().unwrap();
1724        let tiny = Config { budget_bytes: 1 << 16, io: IoMode::Buffered, sync: SyncMode::Off };
1725        let mut writer = Store::create(d.path(), tiny).unwrap();
1726        for i in 0..2_000u64 {
1727            writer.put(&i.to_be_bytes(), format!("generation one row {i}").as_bytes()).unwrap();
1728        }
1729        writer.commit().unwrap();
1730        writer.checkpoint().unwrap();
1731
1732        let (selected_tx, selected_rx) = std::sync::mpsc::channel();
1733        let (churned_tx, churned_rx) = std::sync::mpsc::channel();
1734        let result = std::thread::scope(|scope| {
1735            let dir = d.path();
1736            let reader = scope.spawn(move || {
1737                let snapshot = Store::open_snapshot_with_after_meta(dir, tiny, |generation| {
1738                    selected_tx.send(generation).unwrap();
1739                    churned_rx.recv().unwrap();
1740                    Ok(())
1741                })?;
1742                let iter = snapshot.scan(&[])?;
1743                iter.collect::<Result<Vec<_>>>()
1744            });
1745            let writer_thread = scope.spawn(move || {
1746                assert_eq!(selected_rx.recv().unwrap(), 1, "fixture must stop after selecting generation 1");
1747                for round in 0..4u64 {
1748                    for i in 0..2_000u64 {
1749                        writer.put(&i.to_be_bytes(),
1750                                   format!("writer round {round} row {i} is different").as_bytes()).unwrap();
1751                    }
1752                    writer.commit().unwrap();
1753                    writer.checkpoint().unwrap();
1754                }
1755                churned_tx.send(()).unwrap();
1756            });
1757            writer_thread.join().unwrap();
1758            reader.join().unwrap()
1759        });
1760
1761        let rows = result.expect("a snapshot must not reach recycled pages while it is opening");
1762        assert_eq!(rows.len(), 2_000);
1763        for (i, (key, value)) in rows.iter().enumerate() {
1764            assert_eq!(key.as_slice(), &(i as u64).to_be_bytes());
1765            assert_eq!(value.as_slice(), format!("generation one row {i}").as_bytes(),
1766                       "the opening snapshot observed a recycled page at row {i}");
1767        }
1768    }
1769
1770    struct CrashDirectory {
1771        inner: Box<dyn FileIo>,
1772        directory_durable: std::sync::atomic::AtomicBool,
1773    }
1774
1775    impl FileIo for CrashDirectory {
1776        fn requires_alignment(&self) -> bool { self.inner.requires_alignment() }
1777        fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> { self.inner.read_at(buf, off) }
1778        fn write_at(&self, buf: &[u8], off: u64) -> Result<()> { self.inner.write_at(buf, off) }
1779        fn sync_data(&self) -> Result<()> { self.inner.sync_data() }
1780        fn sync_full(&self) -> Result<()> { self.inner.sync_full() }
1781        fn sync_full_primitive(&self) -> &'static str { self.inner.sync_full_primitive() }
1782        fn sync_dir(&self) -> Result<()> {
1783            self.inner.sync_dir()?;
1784            self.directory_durable.store(true, std::sync::atomic::Ordering::Release);
1785            Ok(())
1786        }
1787        fn len(&self) -> Result<u64> { self.inner.len() }
1788        fn set_len(&self, n: u64) -> Result<()> { self.inner.set_len(n) }
1789    }
1790
1791    /// F5d: model power loss after an acknowledged commit but before the first
1792    /// checkpoint. File barriers preserve contents; only a directory barrier
1793    /// preserves the new data/WAL names. If creation did not issue one, the
1794    /// crash model drops those volatile names and the acknowledged row vanishes.
1795    #[test]
1796    fn a_commit_before_the_first_checkpoint_survives_a_directory_crash() {
1797        let d = tempfile::tempdir().unwrap();
1798        let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1799        let fio = Arc::new(CrashDirectory {
1800            inner: real,
1801            directory_durable: std::sync::atomic::AtomicBool::new(false),
1802        });
1803        {
1804            let mut s = Store::create_on(d.path(), cfg(), fio.clone()).unwrap();
1805            s.put(b"acknowledged", b"must survive power loss").unwrap();
1806            s.commit().unwrap();
1807        }
1808
1809        if !fio.directory_durable.load(std::sync::atomic::Ordering::Acquire) {
1810            std::fs::remove_file(d.path().join("data")).unwrap();
1811            std::fs::remove_file(d.path().join("wal")).unwrap();
1812        }
1813        let reopened = Store::open(d.path(), cfg()).unwrap();
1814        assert_eq!(reopened.get(b"acknowledged").unwrap().as_deref(),
1815                   Some(&b"must survive power loss"[..]),
1816                   "creation must make file names durable before any commit can be acknowledged");
1817    }
1818
1819    #[test]
1820    fn recreated_wal_name_survives_a_commit_before_checkpoint() {
1821        let d = tempfile::tempdir().unwrap();
1822        {
1823            let mut s = Store::create(d.path(), cfg()).unwrap();
1824            s.put(b"old", b"published").unwrap();
1825            s.checkpoint().unwrap();
1826        }
1827        std::fs::remove_file(d.path().join("wal")).unwrap();
1828        let (real, mode) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1829        let fio = Arc::new(CrashDirectory {
1830            inner: real,
1831            directory_durable: std::sync::atomic::AtomicBool::new(false),
1832        });
1833        {
1834            let mut s = Store::build_on(d.path(), cfg(), false, fio.clone(), mode).unwrap();
1835            s.put(b"new", b"acknowledged").unwrap();
1836            s.commit().unwrap();
1837        }
1838        if !fio.directory_durable.load(std::sync::atomic::Ordering::Acquire) {
1839            std::fs::remove_file(d.path().join("wal")).unwrap();
1840        }
1841        drop(fio);
1842        let s = Store::open(d.path(), cfg()).unwrap();
1843        assert_eq!(s.get(b"old").unwrap().as_deref(), Some(&b"published"[..]));
1844        assert_eq!(s.get(b"new").unwrap().as_deref(), Some(&b"acknowledged"[..]));
1845    }
1846
1847    #[test]
1848    fn a_maximum_generation_read_from_disk_is_refused_not_wrapped() {
1849        let d = tempfile::tempdir().unwrap();
1850        {
1851            let s = Store::create(d.path(), cfg()).unwrap();
1852            Meta {
1853                format_version: crate::meta::FORMAT_VERSION,
1854                roots: [s.root, 0, 0, 0, 0, 0, 0, 0],
1855                next_lsn: s.wal.as_ref().unwrap().next_lsn(),
1856                generation: u64::MAX,
1857            }.write_slot(&s.pool).unwrap();
1858            s.pool.flush_all(Barrier::None).unwrap();
1859        }
1860
1861        assert!(matches!(Store::open(d.path(), cfg()), Err(crate::Error::Corrupt { .. })),
1862                "the generation after u64::MAX does not exist and must not become zero");
1863    }
1864
1865    #[test]
1866    fn a_checkpoint_refuses_when_no_later_page_generation_exists() {
1867        let d = tempfile::tempdir().unwrap();
1868        let mut s = Store::create(d.path(), cfg()).unwrap();
1869        s.generation = u64::MAX - 1;
1870        s.put(b"pending", b"kept in the log").unwrap();
1871        s.commit().unwrap();
1872
1873        assert!(matches!(s.checkpoint(), Err(crate::Error::Corrupt { .. })),
1874                "publishing the final generation would leave the next epoch wrapping to zero");
1875        assert!(std::fs::metadata(d.path().join("wal")).unwrap().len() > 0,
1876                "refusing exhaustion must preserve the committed log");
1877    }
1878
1879    #[test]
1880    fn a_checkpoint_rotates_the_log_only_after_the_pages_are_durable() {
1881        let d = tempfile::tempdir().unwrap();
1882        let mut s = Store::create(d.path(), cfg()).unwrap();
1883        for i in 0..2000u64 { s.put(&i.to_be_bytes(), b"v").unwrap(); }
1884        s.commit().unwrap();
1885        s.checkpoint().unwrap();
1886        let order = s.checkpoint_trace();
1887        // 2f: the dual-slot flip sits strictly between the data barrier and
1888        // the WAL rotation -- a flip before the pages are durable can
1889        // publish roots whose pages never landed; a rotation before the
1890        // flip is durable drops the only other copy of this epoch.
1891        assert_eq!(order, vec!["flush_pages", "sync_file", "flip_meta", "rotate_wal"],
1892                   "checkpoint order must be: data durable, then flip, then drop the log");
1893        assert!(s.get(&1999u64.to_be_bytes()).unwrap().is_some());
1894    }
1895
1896    /// `Wal::scan` derives `next_lsn` from the file, so once `rotate()` empties
1897    /// it a restart would otherwise renumber from 1 and the LSN space would
1898    /// repeat. Nothing compares a page's `lsn` against a record's today, so a
1899    /// repeat is currently harmless -- but it becomes silently wrong the
1900    /// moment anyone adds the obvious "skip this record, the page already has
1901    /// a higher lsn" check. `checkpoint` persists `wal.next_lsn()` into
1902    /// `Meta`, and `open` re-supplies it via `wal.set_lsn_floor`; this test is
1903    /// the only thing in the suite that would notice if that wiring were
1904    /// dropped -- `wal.rs`'s own unit tests never call `set_lsn_floor`.
1905    #[test]
1906    fn a_rotation_followed_by_a_reopen_does_not_reissue_lsn_1() {
1907        let d = tempfile::tempdir().unwrap();
1908        let lsn_at_checkpoint = {
1909            let mut s = Store::create(d.path(), cfg()).unwrap();
1910            for i in 0..500u64 { s.put(&i.to_be_bytes(), b"v").unwrap(); }
1911            s.commit().unwrap();
1912            s.checkpoint().unwrap();   // rotates the log -- would reset next_lsn to 1 without the floor
1913            s.next_lsn()
1914        };
1915        assert!(lsn_at_checkpoint > 1, "sanity: many records were appended before the rotation");
1916
1917        let s2 = Store::open(d.path(), cfg()).unwrap();
1918        assert_eq!(
1919            s2.next_lsn(), lsn_at_checkpoint,
1920            "a reopen after rotation must not renumber LSNs from 1"
1921        );
1922    }
1923
1924    /// The three modes must issue three DIFFERENT things. Without this test,
1925    /// collapsing `Full` and `Normal` back into one call would break nothing --
1926    /// and a durability label no test defends is decorative, which makes every
1927    /// benchmark carrying it unattributable. This is the gap the round-1 fix
1928    /// shipped with: the reviewer found it independently of the coordinator.
1929    #[test]
1930    fn the_three_durability_modes_issue_different_barriers() {
1931        for (mode, want) in [
1932            (SyncMode::Full,   vec!["sync_full"]),
1933            (SyncMode::Normal, vec!["sync_data"]),
1934            (SyncMode::Off,    vec![]),
1935        ] {
1936            let d = tempfile::tempdir().unwrap();
1937            let cfg = Config { budget_bytes: 16 << 20, io: IoMode::Buffered, sync: mode };
1938            let mut s = Store::create(d.path(), cfg).unwrap();
1939            s.put(b"k", b"v").unwrap();
1940            s.commit().unwrap();
1941            assert_eq!(s.barriers(), want, "{mode:?} issued the wrong barrier");
1942        }
1943    }
1944
1945    // -- Task 16: the append hint lives on `Store`, not on a throwaway `BTree` --
1946
1947    /// Inserts that reached their leaf without a descent, across BOTH append
1948    /// caches, and the same for attempts.
1949    ///
1950    /// K1 added the per-keyspace cache, and the whole-tree hint these tests
1951    /// were written for is now one case of it: a leaf that is rightmost of the
1952    /// tree is also rightmost of its own tag. Whichever cache owns a key is
1953    /// the only one that probes it -- probing both would make two interleaved
1954    /// ascending runs knock each other's hint out -- so an ascending run's
1955    /// hits move between the two as splits move which leaf is the tree's last.
1956    /// The PROPERTY these tests pin is unchanged and is about the sum: an
1957    /// unbroken ascending run pays no descent, an out-of-order key disarms,
1958    /// and the cache comes back on its own.
1959    fn armed(s: &Store) -> (u64, u64) {
1960        (s.fast_path_hits.get() + s.tag_hints.hits(),
1961         s.fast_path_attempts.get() + s.tag_hints.attempts())
1962    }
1963
1964    /// The test Task 15 needed but did not get. Task 15's own fast-path test
1965    /// (`btree::tests::insert_ascending_uses_fast_path`) holds one long-lived
1966    /// `BTree` for its whole run and passes -- but that is not what
1967    /// production, or the benchmark, actually does: `Store::put` built a
1968    /// fresh `BTree` on every call and dropped it immediately, so
1969    /// `last_leaf` was `None` on entry to every single insert and the fast
1970    /// path never fired through `Store` at all. Before Task 16's fix this
1971    /// assertion reads 0, not `n - 1`; asserting the exact count, not just
1972    /// `> 0`, is what would have caught a path that only fires sometimes.
1973    ///
1974    /// `n` is small enough (as in `btree::tests::insert_ascending_uses_fast_path`)
1975    /// that these tiny records never fill a single leaf, so no split ever
1976    /// falls back to a full descent -- a split is a legitimate miss (the
1977    /// leaf genuinely has no room), and mixing that into this count would
1978    /// make the assertion about page capacity instead of about the hint.
1979    #[test]
1980    fn store_put_ascending_uses_fast_path() {
1981        let d = tempfile::tempdir().unwrap();
1982        let mut s = Store::create(d.path(), cfg()).unwrap();
1983        let n = 100u64;
1984        for i in 0..n { s.put(&i.to_be_bytes(), b"v").unwrap(); }
1985        assert_eq!(
1986            armed(&s).0, n - 1,
1987            "every Store::put but the first must hit the append fast path"
1988        );
1989    }
1990
1991    /// `bulk_load` repacks the whole tree by external sort, so any leaf
1992    /// number the hint named before the call means nothing afterward -- it
1993    /// may not even be allocated, or may now hold an unrelated page. A
1994    /// `Store::put` right after `bulk_load` must still land correctly (via
1995    /// a real descent, since the hint was cleared) and a full scan must stay
1996    /// in order.
1997    #[test]
1998    fn store_put_survives_bulk_load() {
1999        let d = tempfile::tempdir().unwrap();
2000        let mut s = Store::create(d.path(), cfg()).unwrap();
2001        let n = 1_000u64;
2002        let items = (0..n).map(|i| (i.to_be_bytes().to_vec(), b"bulk".to_vec()));
2003        s.bulk_load(items).unwrap();
2004        assert_eq!(s.last_leaf.get(), None, "bulk_load must clear a hint it just made meaningless");
2005
2006        // Sorts after every bulk-loaded key.
2007        let tail_key = n.to_be_bytes();
2008        s.put(&tail_key, b"tail").unwrap();
2009
2010        assert_eq!(
2011            s.get(&tail_key).unwrap().as_deref(), Some(&b"tail"[..]),
2012            "a put right after bulk_load must be found"
2013        );
2014
2015        let scanned: Vec<Vec<u8>> = s.scan(&[]).unwrap().map(|r| r.unwrap().0).collect();
2016        let mut sorted = scanned.clone();
2017        sorted.sort();
2018        assert_eq!(scanned.len(), n as usize + 1);
2019        assert_eq!(scanned, sorted, "a full scan after bulk_load + put must stay in sorted order");
2020    }
2021
2022    /// The correctness net for the hint's new, longer-lived home: the same
2023    /// 20k keys through `Store::put`, once in ascending order (exercising
2024    /// the fast path constantly) and once in scattered order (exercising
2025    /// `descend_for_write` and splits constantly, since a scattered key is
2026    /// essentially never greater than the rightmost leaf's last key), must
2027    /// produce byte-identical full scans out of two independent stores.
2028    #[test]
2029    fn store_random_order_matches_sequential() {
2030        let n = 20_000u64;
2031        let scatter = |i: u64| i.wrapping_mul(0x9E37_79B9_7F4A_7C15);
2032
2033        let d1 = tempfile::tempdir().unwrap();
2034        let mut s1 = Store::create(d1.path(), cfg()).unwrap();
2035        for i in 0..n { s1.put(&i.to_be_bytes(), &i.to_le_bytes()).unwrap(); }
2036
2037        let d2 = tempfile::tempdir().unwrap();
2038        let mut s2 = Store::create(d2.path(), cfg()).unwrap();
2039        let mut order: Vec<u64> = (0..n).collect();
2040        order.sort_by_key(|&i| scatter(i));
2041        for &i in &order { s2.put(&i.to_be_bytes(), &i.to_le_bytes()).unwrap(); }
2042
2043        let seq1: Vec<(Vec<u8>, Vec<u8>)> = s1.scan(&[]).unwrap().map(|r| r.unwrap()).collect();
2044        let seq2: Vec<(Vec<u8>, Vec<u8>)> = s2.scan(&[]).unwrap().map(|r| r.unwrap()).collect();
2045        assert_eq!(seq1.len(), n as usize);
2046        assert_eq!(
2047            seq1, seq2,
2048            "ascending vs scattered insertion order through Store::put must produce identical scans"
2049        );
2050    }
2051
2052    // -- Task 18: the append hint disarms itself when it stops paying --
2053
2054    /// Before this task, `insert_into_leaf` re-armed `last_leaf` after EVERY
2055    /// descent insert whether or not the leaf it wrote was the rightmost
2056    /// one, and nothing ever cleared a hint that had just failed. On the
2057    /// probe's own scattered-key workload (MODE=incr: key =
2058    /// `i.wrapping_mul(0x9E37_79B9_7F4A_7C15)`, a 200-byte payload -- both
2059    /// matched exactly here, not approximated, because this is the workload
2060    /// the project is measured on) that meant a
2061    /// `pool.get_mut` plus a full CRC-verifying `PageRef::open` thrown away
2062    /// on very nearly every single row, forever: attempts grew with `n`.
2063    ///
2064    /// After the fix, an attempt is followed by the hint disarming unless it
2065    /// actually paid off, so most inserts see `last_leaf.get() == None` and
2066    /// never call `fast_path_leaf` at all. What is left is (a) the leaf's
2067    /// worth of inserts before the very first split -- every write lands on
2068    /// the sole, trivially-rightmost leaf, so the hint stays armed and is
2069    /// tried again next time regardless of key order -- plus (b) one more
2070    /// attempt each time a later scattered key happens to be a new running
2071    /// maximum, which is what actually re-arms the hint after (a) ends. Both
2072    /// are bounded by leaf capacity and by how many new maxima a random
2073    /// sequence produces (~log n), not by `n` -- measured at 85/97/107/117/126
2074    /// attempts for n=2000/5000/10000/20000/40000, a ~1.5x change over a 20x
2075    /// change in `n`. 20000 is asserted here exactly (117): a bound written
2076    /// as `< n` would still pass with the pre-fix O(n) behavior and catch
2077    /// nothing.
2078    #[test]
2079    fn random_order_disarms_fast_path() {
2080        let d = tempfile::tempdir().unwrap();
2081        let mut s = Store::create(d.path(), cfg()).unwrap();
2082        let n = 20_000u64;
2083        let scatter = |i: u64| i.wrapping_mul(0x9E37_79B9_7F4A_7C15);
2084        let payload = vec![b'x'; 200]; // same shape as probe's MODE=incr payload
2085        for i in 0..n { s.put(&scatter(i).to_be_bytes(), &payload).unwrap(); }
2086        // Neighbor redistribution explicitly disarms the append hint, reducing
2087        // attempts (101 for this workload). Preserve the original upper bound;
2088        // the non-redistributing path still has its exact historical count.
2089        #[cfg(feature = "sqlite-balance")]
2090        assert!(s.fast_path_attempts.get() <= 117);
2091        // And the per-keyspace cache must be bounded the same way. It arms
2092        // only when a record lands at the END of its leaf, which a scattered
2093        // key does about once per leaf-full, so this is bounded by leaf
2094        // capacity and the number of leaves -- not by n. An implementation
2095        // that re-armed on every descent would read one extra page per insert
2096        // here, 20,000 of them, which is precisely the pre-Task-18 defect in
2097        // a new place.
2098        assert!(
2099            s.tag_hints.attempts() <= 2_000,
2100            "a scattered workload attempted the per-keyspace fast path {} times in {n} \
2101             inserts; arming belongs to appends, not to every descent",
2102            s.tag_hints.attempts()
2103        );
2104        #[cfg(not(feature = "sqlite-balance"))]
2105        assert_eq!(
2106            s.fast_path_attempts.get(), 117,
2107            "a disarming hint must attempt the fast path a number of times bounded by \
2108             leaf capacity and log(n), not by n -- a bound proportional to n would not \
2109             have caught the pre-Task-18 defect this test exists for"
2110        );
2111    }
2112
2113    /// The regression guard test 1 needs: a too-aggressive disarm -- say, one
2114    /// that clears the hint on any descent instead of only on a failed
2115    /// fast-path attempt -- would silently destroy the append optimisation
2116    /// for the common ascending-key case (autoincrementing IDs, time-ordered
2117    /// writes) without any test noticing, since test 1 only bounds an upper
2118    /// limit. `n` is small enough that these tiny records never fill one
2119    /// leaf, so no split's legitimate room-check miss dilutes the count: the
2120    /// only insert that does not hit is the very first, before the hint is
2121    /// ever armed.
2122    #[test]
2123    fn ascending_keeps_fast_path_armed() {
2124        let d = tempfile::tempdir().unwrap();
2125        let mut s = Store::create(d.path(), cfg()).unwrap();
2126        let n = 100u64;
2127        for i in 0..n { s.put(&i.to_be_bytes(), b"v").unwrap(); }
2128        assert_eq!(
2129            armed(&s).0, n - 1,
2130            "an unbroken ascending run must still hit the fast path on every insert but the first"
2131        );
2132    }
2133
2134    /// Ascending, then one out-of-order key, then ascending again above the
2135    /// prior high-water mark: the fast path must disarm on the out-of-order
2136    /// key and come back on its own once a genuine append re-establishes it.
2137    ///
2138    /// `m = 300` ascending keys is deliberately past this leaf's capacity
2139    /// (~238 of these 13-byte records), so the run crosses exactly one
2140    /// split: the insert that fills the leaf is a legitimate room-check miss
2141    /// (attempted, since the hint was armed, but correctly rejected -- there
2142    /// is no space), immediately followed by the split path re-arming the
2143    /// hint onto the new right leaf. That costs one hit out of `m - 1`
2144    /// attempts, not `m - 1` hits, and is why the first assertion is `m - 2`
2145    /// rather than `m - 1`. Splitting first also means the tree has more
2146    /// than one leaf before the out-of-order key lands, so its failure is
2147    /// the real disarm this test exists to check -- with only one leaf in
2148    /// the tree that leaf is always "rightmost" and the hint would re-arm to
2149    /// it immediately regardless of this task's fix, proving nothing.
2150    ///
2151    /// Key `0` sorts before everything already inserted, so it lands on the
2152    /// leftmost leaf -- not the rightmost one the hint pointed at -- and
2153    /// that leaf's own last key is 300-ish, not `0`, so the ordering check
2154    /// (check 5) rejects it too if the wrong-leaf check somehow didn't.
2155    /// Either way it is one attempt, one miss, and (since the leftmost leaf
2156    /// is not rightmost) the hint stays cleared afterward: `last_leaf` reads
2157    /// `None`, not just "hits did not increase".
2158    ///
2159    /// The next ascending key, `m + 1`, finds the hint disarmed and pays for
2160    /// one descent -- no attempt is counted for it -- but that descent lands
2161    /// on the genuine rightmost leaf (it is the new global maximum) and
2162    /// re-arms the hint there. Every ascending key after that hits again:
2163    /// `k - 1` more hits and attempts, for a workload small enough (`k =
2164    /// 30`, well under the ~119 records of headroom left in that leaf) that
2165    /// no second split dilutes the count.
2166    #[test]
2167    fn mixed_workload_rearms() {
2168        let d = tempfile::tempdir().unwrap();
2169        let mut s = Store::create(d.path(), cfg()).unwrap();
2170        let m = 300u64;
2171        let k = 30u64;
2172
2173        for i in 1..=m { s.put(&i.to_be_bytes(), b"v").unwrap(); }
2174        assert_eq!(armed(&s).1, m - 1, "one attempt per insert but the first");
2175        assert_eq!(
2176            armed(&s).0, m - 2,
2177            "one miss expected: the insert that fills the leaf and forces the one split \
2178             this run crosses"
2179        );
2180
2181        s.put(&0u64.to_be_bytes(), b"v").unwrap();
2182        assert_eq!(armed(&s).1, m, "the out-of-order key is one more attempt");
2183        assert_eq!(armed(&s).0, m - 2, "the out-of-order key must not hit");
2184        assert_eq!(
2185            s.last_leaf.get(), None,
2186            "landing on a non-rightmost leaf must leave the hint disarmed, not re-armed \
2187             to the wrong leaf"
2188        );
2189
2190        let (hits, attempts) = armed(&s);
2191        for i in (m + 1)..(m + 1 + k) { s.put(&i.to_be_bytes(), b"v").unwrap(); }
2192        let (hits, attempts) = (armed(&s).0 - hits, armed(&s).1 - attempts);
2193        assert_eq!(
2194            attempts, hits,
2195            "an unbroken ascending run must hit every probe it makes"
2196        );
2197        // AT MOST ONE descent in the whole run: the fast path must come back
2198        // on its own. Exactly which insert pays it depends on how the build
2199        // splits. Without `sqlite-balance` the split re-arms the per-keyspace
2200        // hint on the page it left the record in, so the run is never
2201        // interrupted at all and k of k hit. With it, the leaf fill goes
2202        // through `redistribute_neighbors`, which has to forget the fences of
2203        // every page in its window, and the insert after it descends once.
2204        assert!(
2205            k - hits <= 1,
2206            "{} of the {k} ascending inserts after the disarm paid a descent",
2207            k - hits
2208        );
2209    }
2210
2211    // -- Task 19: poisoning, both halves. The CHECK (every writer refuses)
2212    // and the TRIGGER (a failed checkpoint barrier is what sets it), plus
2213    // the way out, which Law 5 requires any refusal to have.
2214
2215    /// A `FileIo` whose barrier can be made to fail on demand. `Store` had
2216    /// no fault-injection seam at all, which is why the poisoning trigger
2217    /// shipped with zero executable evidence (Task 17 final review, F6):
2218    /// deleting the `self.poisoned = true` line left the whole suite green.
2219    struct FailingBarrier { inner: Box<dyn FileIo>, fail: std::sync::atomic::AtomicBool }
2220    impl FailingBarrier {
2221        fn failing(&self) -> bool { self.fail.load(std::sync::atomic::Ordering::Relaxed) }
2222        fn eio() -> crate::Error {
2223            std::io::Error::other("injected barrier failure").into()
2224        }
2225    }
2226    impl FileIo for FailingBarrier {
2227        fn requires_alignment(&self) -> bool { self.inner.requires_alignment() }
2228        fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> { self.inner.read_at(buf, off) }
2229        fn write_at(&self, buf: &[u8], off: u64) -> Result<()> { self.inner.write_at(buf, off) }
2230        fn sync_data(&self) -> Result<()> {
2231            if self.failing() { return Err(Self::eio()); }
2232            self.inner.sync_data()
2233        }
2234        fn sync_full(&self) -> Result<()> {
2235            if self.failing() { return Err(Self::eio()); }
2236            self.inner.sync_full()
2237        }
2238        fn sync_full_primitive(&self) -> &'static str { self.inner.sync_full_primitive() }
2239        fn sync_dir(&self) -> Result<()> { self.inner.sync_dir() }
2240        fn len(&self) -> Result<u64> { self.inner.len() }
2241        fn set_len(&self, n: u64) -> Result<()> { self.inner.set_len(n) }
2242    }
2243
2244    /// Fail one selected data-page write, after a chosen number of matching
2245    /// writes have succeeded. The WAL is a separate file, so this reaches the
2246    /// exact pool eviction boundary without making log append fail first.
2247    struct FailingWrite {
2248        inner: Box<dyn FileIo>,
2249        armed: std::sync::Mutex<Option<usize>>,
2250    }
2251    impl FailingWrite {
2252        fn arm(&self, successful_writes: usize) {
2253            *self.armed.lock().unwrap() = Some(successful_writes);
2254        }
2255    }
2256    impl FileIo for FailingWrite {
2257        fn requires_alignment(&self) -> bool { self.inner.requires_alignment() }
2258        fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> { self.inner.read_at(buf, off) }
2259        fn write_at(&self, buf: &[u8], off: u64) -> Result<()> {
2260            let mut armed = self.armed.lock().unwrap();
2261            if let Some(remaining) = armed.as_mut() {
2262                if *remaining == 0 {
2263                    *armed = None;
2264                    return Err(std::io::Error::other(
2265                        "injected page-write failure during a tree mutation",
2266                    ).into());
2267                }
2268                *remaining -= 1;
2269            }
2270            self.inner.write_at(buf, off)
2271        }
2272        fn sync_data(&self) -> Result<()> { self.inner.sync_data() }
2273        fn sync_full(&self) -> Result<()> { self.inner.sync_full() }
2274        fn sync_full_primitive(&self) -> &'static str { self.inner.sync_full_primitive() }
2275        fn sync_dir(&self) -> Result<()> { self.inner.sync_dir() }
2276        fn len(&self) -> Result<u64> { self.inner.len() }
2277        fn set_len(&self, n: u64) -> Result<()> { self.inner.set_len(n) }
2278    }
2279
2280    #[test]
2281    fn deletion_merge_eviction_failures_cannot_publish_partial_changes() {
2282        for fail_after in 0..12 {
2283            let d=tempfile::tempdir().unwrap();
2284            let cfg=Config{budget_bytes:64<<10,io:IoMode::Buffered,sync:SyncMode::Full};
2285            let (real,_)=open_file(&d.path().join("data"),IoMode::Buffered).unwrap();
2286            let fio=Arc::new(FailingWrite{inner:real,armed:std::sync::Mutex::new(None)});
2287            let mut s=Store::create_on(d.path(),cfg,fio.clone()).unwrap();
2288            for i in 0u64..1024 {s.put(&i.to_be_bytes(),&vec![7;240]).unwrap();}
2289            // Fifteen 256-byte cells fit each full leaf. Keep six, just
2290            // above the maintenance threshold; subsequent hits force merges.
2291            for i in 0u64..1024 {if i%15>=6 {s.delete(&i.to_be_bytes()).unwrap();}}
2292            s.checkpoint().unwrap();
2293            let old=Store::open_snapshot(d.path(),cfg).unwrap();
2294            fio.arm(fail_after);let mut failed=false;
2295            for i in 0u64..1024 {
2296                if let Err(e)=s.delete(&((i*71)%1024).to_be_bytes()) {
2297                    assert!(matches!(e,crate::Error::Io(_)));failed=true;break;
2298                }
2299            }
2300            assert!(failed,"fault {fail_after} was not reached");
2301            assert!(s.commit().is_err());assert!(s.checkpoint().is_err());
2302            for i in 0u64..1024 {assert_eq!(old.get(&i.to_be_bytes()).unwrap(),(i%15<6).then(||vec![7;240]));}
2303            drop(old);drop(s);
2304            let reopened=Store::open(d.path(),cfg).unwrap();
2305            for i in 0u64..1024 {assert_eq!(reopened.get(&i.to_be_bytes()).unwrap(),(i%15<6).then(||vec![7;240]));}
2306            assert_eq!(crate::verify::verify_published_tree(&d.path().join("data"),IoMode::Buffered,reopened.root,1).unwrap().0,412);
2307        }
2308    }
2309
2310    /// F7's reachable tear, not merely an error before a mutation starts.
2311    /// Earlier writes are committed in the same unpublished epoch. The next
2312    /// scattered insert splits the root leaf, rewrites its left half, then
2313    /// fails while allocating the new parent; committed keys moved right are
2314    /// no longer reachable from the still-leaf root.
2315    #[test]
2316    fn a_partial_tree_mutation_poisons_and_reopen_replays_the_retained_log() {
2317        let d = tempfile::tempdir().unwrap();
2318        let tiny = Config {
2319            budget_bytes: 16 * crate::page::PAGE_SIZE,
2320            io: IoMode::Buffered,
2321            sync: SyncMode::Off,
2322        };
2323        let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
2324        let fio = Arc::new(FailingWrite {
2325            inner: real,
2326            armed: std::sync::Mutex::new(None),
2327        });
2328        let mut s = Store::create_on(d.path(), tiny, fio.clone()).unwrap();
2329
2330        let value = vec![b'v'; 256];
2331        let record_bytes = 4 + 8 + value.len() + 4;
2332        let mut rows = 0u64;
2333        loop {
2334            let room = {
2335                let r = s.pool.get(s.root).unwrap();
2336                crate::page::PageRef::open_resident(&r, s.root).unwrap().free_space()
2337            };
2338            if room < record_bytes { break; }
2339            s.put(&(rows * 2).to_be_bytes(), &value).unwrap();
2340            s.commit().unwrap();
2341            rows += 1;
2342        }
2343        assert!(rows > 4, "fixture must have committed rows to move across the split");
2344
2345        // Fill every initially-unused frame, then pin the root while replacing
2346        // all other residents with dirty throwaway pages. The split's first
2347        // allocation can therefore evict once; its second allocation (the new
2348        // root, after the old leaf was rewritten) reaches the injected failure.
2349        while s.pool.page_count() < 16 {
2350            drop(s.pool.allocate().unwrap());
2351        }
2352        {
2353            let root_pin = s.pool.get(s.root).unwrap();
2354            for _ in 0..15 { drop(s.pool.allocate().unwrap()); }
2355            drop(root_pin);
2356        }
2357
2358        let wal_before = std::fs::metadata(d.path().join("wal")).unwrap().len();
2359        fio.arm(1);
2360        let failed = s.put(&1u64.to_be_bytes(), &value);
2361        assert!(matches!(failed, Err(crate::Error::Io(_))),
2362                "the injected eviction failure must escape the mutating insert, got {failed:?}");
2363        assert_eq!(s.get(&((rows - 1) * 2).to_be_bytes()).unwrap(), None,
2364                   "fixture must prove the insert failed after committed keys moved off the root");
2365
2366        // Exercise both dangerous retries before asserting: on the old code
2367        // commit succeeds and checkpoint then publishes the tear and rotates
2368        // away the only log that could rebuild it.
2369        let commit = s.commit();
2370        let checkpoint = s.checkpoint();
2371        let wal_after = std::fs::metadata(d.path().join("wal")).unwrap().len();
2372        assert!(matches!(commit, Err(crate::Error::StorePoisoned)),
2373                "a partial tree mutation must make commit refuse, got {commit:?}");
2374        assert!(matches!(checkpoint, Err(crate::Error::StorePoisoned)),
2375                "a partial tree mutation must make checkpoint refuse, got {checkpoint:?}");
2376        assert_eq!(wal_after, wal_before,
2377                   "a poisoned store must retain the committed recovery log byte-for-byte");
2378
2379        drop(s);
2380        drop(fio);
2381        let reopened = Store::open(d.path(), tiny)
2382            .expect("reopen must rebuild a poisoned handle from the retained log");
2383        for i in 0..rows {
2384            assert_eq!(reopened.get(&(i * 2).to_be_bytes()).unwrap().as_deref(), Some(value.as_slice()));
2385        }
2386        assert_eq!(reopened.get(&1u64.to_be_bytes()).unwrap(), None,
2387                   "the failed, uncommitted insert must not be replayed");
2388    }
2389
2390    /// A validation refusal happens before the WAL and before any tree page is
2391    /// touched. It is a caller-correctable input error, not evidence of a torn
2392    /// tree, and must not turn the handle into a denial of service.
2393    #[test]
2394    fn a_pre_mutation_validation_error_does_not_poison_the_store() {
2395        let d = tempfile::tempdir().unwrap();
2396        let mut s = Store::create(d.path(), cfg()).unwrap();
2397        let oversized_key = vec![b'k'; crate::page::MAX_RECORD_LEN];
2398        assert!(matches!(s.put(&oversized_key, b"v"), Err(crate::Error::TooLarge)));
2399        s.put(b"usable", b"still").unwrap();
2400        s.commit().unwrap();
2401        s.checkpoint().unwrap();
2402        assert_eq!(s.get(b"usable").unwrap().as_deref(), Some(&b"still"[..]));
2403    }
2404
2405    /// THE TRIGGER, executed rather than asserted by inspection. A
2406    /// checkpoint whose barrier fails must (a) surface the error and (b)
2407    /// leave the store refusing every subsequent write -- above all
2408    /// `checkpoint` itself, whose last step deletes the log (see
2409    /// `Error::StorePoisoned`). Deleting `checkpoint`'s `self.poisoned =
2410    /// true` line makes every assertion after the first one fail: the store
2411    /// happily accepts a `put`, a `commit`, and another `checkpoint` -- and
2412    /// that second checkpoint rotates the log away.
2413    #[test]
2414    fn a_failed_checkpoint_barrier_poisons_every_writer() {
2415        let d = tempfile::tempdir().unwrap();
2416        let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
2417        let fio = Arc::new(FailingBarrier { inner: real, fail: false.into() });
2418        let mut s = Store::create_on(d.path(), cfg(), fio.clone()).unwrap();
2419        s.put(b"a", b"1").unwrap();
2420        s.commit().unwrap();
2421        s.checkpoint().expect("sanity: the checkpoint works while the disk does");
2422
2423        s.put(b"b", b"2").unwrap();
2424        s.commit().unwrap();
2425        fio.fail.store(true, std::sync::atomic::Ordering::Relaxed);
2426
2427        match s.checkpoint() {
2428            Err(crate::Error::Io(_)) => {}
2429            other => panic!("a failing barrier must surface, got {other:?}"),
2430        }
2431        fio.fail.store(false, std::sync::atomic::Ordering::Relaxed);   // the disk "recovers"
2432
2433        assert!(matches!(s.put(b"c", b"3"), Err(crate::Error::StorePoisoned)));
2434        assert!(matches!(s.delete(b"a"), Err(crate::Error::StorePoisoned)));
2435        assert!(matches!(s.commit(), Err(crate::Error::StorePoisoned)));
2436        assert!(matches!(s.checkpoint(), Err(crate::Error::StorePoisoned)),
2437                "a retried checkpoint would flush nothing (those frames are marked clean \
2438                 already) and then rotate the log away");
2439        let items = vec![(b"z".to_vec(), b"9".to_vec())].into_iter();
2440        assert!(matches!(s.bulk_load(items), Err(crate::Error::StorePoisoned)),
2441                "bulk_load replaces the whole tree AND discards the log");
2442        // And it must refuse BEFORE doing any of that, not fail on its way
2443        // out. Without its own guard, `bulk_load` still returns
2444        // `StorePoisoned` -- from the `checkpoint()` it ends with -- having
2445        // already repacked the tree and republished the root, which is the
2446        // assertion below (and only the assertion below) that notices.
2447        assert_eq!(s.get(b"a").unwrap().as_deref(), Some(&b"1"[..]),
2448                   "a refused bulk_load must not have replaced the tree");
2449        assert_eq!(s.get(b"z").unwrap(), None);
2450    }
2451
2452    /// The other half of the same guarantee, and the one Law 5 demands: a
2453    /// refusal with no way out is itself an unrecoverable state. Poisoning
2454    /// is per-instance and never persisted, so dropping the store and
2455    /// reopening the directory clears it -- and that reopen is not a way of
2456    /// ignoring the problem, it is what re-derives every belief from disk
2457    /// and replays the log that was never rotated away.
2458    #[test]
2459    fn reopening_clears_poisoning_and_the_committed_data_is_still_there() {
2460        let d = tempfile::tempdir().unwrap();
2461        let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
2462        let fio = Arc::new(FailingBarrier { inner: real, fail: false.into() });
2463        {
2464            let mut s = Store::create_on(d.path(), cfg(), fio.clone()).unwrap();
2465            s.put(b"a", b"1").unwrap();
2466            s.commit().unwrap();
2467            fio.fail.store(true, std::sync::atomic::Ordering::Relaxed);
2468            assert!(s.checkpoint().is_err());
2469            assert!(matches!(s.put(b"b", b"2"), Err(crate::Error::StorePoisoned)));
2470        }
2471        fio.fail.store(false, std::sync::atomic::Ordering::Relaxed);
2472
2473        let mut s2 = Store::open(d.path(), cfg()).expect("a poisoned store must not poison the DIRECTORY");
2474        assert_eq!(s2.get(b"a").unwrap().as_deref(), Some(&b"1"[..]),
2475                   "the committed row is in the log, which the failed checkpoint never rotated");
2476        s2.put(b"b", b"2").unwrap();
2477        s2.commit().unwrap();
2478    }
2479
2480    /// CRC-valid does not mean structurally valid. Recovery must reject a
2481    /// malformed frame without indexing through unchecked bytes and without
2482    /// letting any part of that frame reach the tree.
2483    #[test]
2484    fn malformed_wal_payloads_are_bounded_and_do_not_mutate_the_tree() {
2485        let d = tempfile::tempdir().unwrap();
2486        let mut s = Store::create(d.path(), cfg()).unwrap();
2487        s.put(b"kept", b"value").unwrap();
2488        s.commit().unwrap();
2489        let root = s.root;
2490
2491        let malformed: &[(RecKind, &[u8])] = &[
2492            (RecKind::Put, &[]),
2493            (RecKind::Put, &[4, 0, b'a']),
2494            (RecKind::Delete, &[]),
2495            (RecKind::Delete, &[2, 0, b'a']),
2496            (RecKind::Delete, &[0, 0, b'x']),
2497            (RecKind::DeletePrefix, &[]),
2498            (RecKind::PutEmptyBatch, &[]),
2499            (RecKind::PutEmptyBatch, &[1, 0, 4, 0, b'a']),
2500            (RecKind::PutEmptyBatch, &[65, 0]),
2501            (RecKind::PutEmptyBatch, &[1, 0, 1, 0, b'a', b'x']),
2502            (RecKind::Commit, b"not empty"),
2503            (RecKind::PageImage, b"unsupported"),
2504        ];
2505        for (n, &(kind, payload)) in malformed.iter().enumerate() {
2506            assert!(matches!(
2507                s.apply(kind, payload, n as u64),
2508                Err(crate::Error::CorruptWal { offset, .. }) if offset == n as u64
2509            ));
2510            assert_eq!(s.root, root, "malformed frame {n} changed the root");
2511            assert_eq!(s.get(b"kept").unwrap().as_deref(), Some(&b"value"[..]),
2512                "malformed frame {n} changed existing data");
2513        }
2514    }
2515
2516    #[test]
2517    fn bounded_empty_key_batch_replays_as_one_committed_wal_unit() {
2518        let d = tempfile::tempdir().unwrap();
2519        {
2520            let mut s = Store::create(d.path(), cfg()).unwrap();
2521            s.put_empty_batch(&[b"alpha".to_vec(), b"beta".to_vec()]).unwrap();
2522            s.commit().unwrap();
2523        }
2524        let s = Store::open(d.path(), cfg()).unwrap();
2525        assert_eq!(s.get(b"alpha").unwrap().as_deref(), Some(&b""[..]));
2526        assert_eq!(s.get(b"beta").unwrap().as_deref(), Some(&b""[..]));
2527    }
2528
2529    #[test]
2530    fn a_corrupted_packed_tree_never_becomes_authoritative() {
2531        let d = tempfile::tempdir().unwrap();
2532        let mut s = Store::create(d.path(), cfg()).unwrap();
2533        s.put(b"old", b"authoritative").unwrap();
2534        s.commit().unwrap();
2535        s.checkpoint().unwrap();
2536        let old_root = s.root;
2537
2538        let result = s.bulk_load_with_before_publish(
2539            (0..2_000u64).map(|i| (i.to_be_bytes().to_vec(), b"new".to_vec())),
2540            |data, root| {
2541                use std::io::{Seek, SeekFrom, Write};
2542                let mut file = std::fs::OpenOptions::new().write(true).open(data)?;
2543                file.seek(SeekFrom::Start(root as u64 * PAGE_SIZE as u64))?;
2544                file.write_all(&[0u8])?;
2545                file.sync_all()?;
2546                Ok(())
2547            },
2548        );
2549
2550        assert!(result.is_err(), "a packed tree corrupted before publication must be refused");
2551        assert_eq!(s.root, old_root, "the old root must stay authoritative after refusal");
2552        assert_eq!(s.get(b"old").unwrap().as_deref(), Some(&b"authoritative"[..]));
2553    }
2554
2555    fn graft_key(space: u8, i: u32) -> Vec<u8> {
2556        let mut key = vec![space];
2557        key.extend_from_slice(&i.to_be_bytes());
2558        key
2559    }
2560
2561    fn seed_graft_base(store: &mut Store) -> Vec<(Vec<u8>, Vec<u8>)> {
2562        let mut rows = Vec::new();
2563        for space in [0x10, 0x30] {
2564            for i in 0..1_500u32 {
2565                let row = (graft_key(space, i), format!("base-{space:02x}-{i}").into_bytes());
2566                store.put(&row.0, &row.1).unwrap();
2567                rows.push(row);
2568            }
2569        }
2570        store.commit().unwrap();
2571        store.checkpoint().unwrap();
2572        rows.sort_by(|a, b| a.0.cmp(&b.0));
2573        rows
2574    }
2575
2576    fn graft_rows() -> Vec<(Vec<u8>, Vec<u8>)> {
2577        (0..2_000u32)
2578            .rev()
2579            .map(|i| (graft_key(0x20, i), format!("graft-{i}").into_bytes()))
2580            .collect()
2581    }
2582
2583    fn collect_from(store: &Store, from: &[u8]) -> Vec<(Vec<u8>, Vec<u8>)> {
2584        store.scan(from).unwrap().map(Result::unwrap).collect()
2585    }
2586
2587    fn collect_below(store: &Store, to: &[u8]) -> Vec<(Vec<u8>, Vec<u8>)> {
2588        let mut rows = Vec::new();
2589        store.scan_reverse(to).unwrap().for_each_ref(|key, value| {
2590            rows.push((key.to_vec(), value.to_vec()));
2591            true
2592        }).unwrap();
2593        rows
2594    }
2595
2596    /// Stage 1's first obligation: installing an empty key range may not
2597    /// replace the shared tree. Keys on both sides exercise a true mid-tree
2598    /// graft rather than the easier leftmost/rightmost append cases.
2599    #[test]
2600    fn graft_preserves_every_preexisting_key() {
2601        let d = tempfile::tempdir().unwrap();
2602        let mut s = Store::create(d.path(), cfg()).unwrap();
2603        let base = seed_graft_base(&mut s);
2604        let pinned = Store::open_snapshot(d.path(), cfg()).unwrap();
2605
2606        s.graft_range(graft_rows().into_iter()).unwrap();
2607
2608        for (key, value) in &base {
2609            assert_eq!(s.get(key).unwrap().as_deref(), Some(value.as_slice()),
2610                "graft lost pre-existing key {key:?}");
2611        }
2612        assert_eq!(s.scan(&[]).unwrap().count(), base.len() + 2_000);
2613        assert!(pinned.get(&graft_key(0x20, 17)).unwrap().is_none(),
2614            "a reader pinned before publication must remain on the old generation");
2615        assert_eq!(pinned.scan(&[]).unwrap().count(), base.len());
2616        let fresh = Store::open_snapshot(d.path(), cfg()).unwrap();
2617        assert_eq!(fresh.get(&graft_key(0x20, 17)).unwrap().as_deref(), Some(&b"graft-17"[..]));
2618    }
2619
2620    /// THE EQUIVALENCE ORACLE: construction is allowed to change, answers are
2621    /// not. Point, forward-range, and reverse-range queries must be byte-for-
2622    /// byte identical to the ordinary insertion path over the same rows.
2623    #[test]
2624    fn graft_and_individual_inserts_answer_every_kernel_query_identically() {
2625        let dg = tempfile::tempdir().unwrap();
2626        let di = tempfile::tempdir().unwrap();
2627        let mut grafted = Store::create(dg.path(), cfg()).unwrap();
2628        let mut inserted = Store::create(di.path(), cfg()).unwrap();
2629        seed_graft_base(&mut grafted);
2630        seed_graft_base(&mut inserted);
2631        let rows = graft_rows();
2632
2633        grafted.graft_range(rows.clone().into_iter()).unwrap();
2634        for (key, value) in &rows { inserted.put(key, value).unwrap(); }
2635        inserted.commit().unwrap();
2636        inserted.checkpoint().unwrap();
2637
2638        for (key, _) in seed_query_keys(&rows) {
2639            assert_eq!(grafted.get(&key).unwrap(), inserted.get(&key).unwrap(),
2640                "point query disagreed at {key:?}");
2641        }
2642        for from in [vec![], graft_key(0x10, 777), graft_key(0x20, 0),
2643                     graft_key(0x20, 999), graft_key(0x30, 0), vec![0xff]] {
2644            assert_eq!(collect_from(&grafted, &from), collect_from(&inserted, &from),
2645                "forward range disagreed from {from:?}");
2646        }
2647        for to in [graft_key(0x10, 0), graft_key(0x20, 0), graft_key(0x20, 999),
2648                   graft_key(0x30, 0), vec![0xff]] {
2649            assert_eq!(collect_below(&grafted, &to), collect_below(&inserted, &to),
2650                "reverse range disagreed below {to:?}");
2651        }
2652
2653        // LIVE WRITES stay on the existing path after publication. Exercise a
2654        // later key in the grafted namespace, including its first packed-page
2655        // split, and compare again.
2656        let live_key = graft_key(0x20, 2_500);
2657        grafted.put(&live_key, b"later-live-write").unwrap();
2658        inserted.put(&live_key, b"later-live-write").unwrap();
2659        grafted.commit().unwrap();
2660        inserted.commit().unwrap();
2661        grafted.checkpoint().unwrap();
2662        inserted.checkpoint().unwrap();
2663        assert_eq!(collect_from(&grafted, &graft_key(0x20, 1_900)),
2664                   collect_from(&inserted, &graft_key(0x20, 1_900)));
2665    }
2666
2667    #[test]
2668    fn graft_refuses_a_nonempty_range_without_changing_the_tree() {
2669        let d = tempfile::tempdir().unwrap();
2670        let mut s = Store::create(d.path(), cfg()).unwrap();
2671        let base = seed_graft_base(&mut s);
2672        let old_root = s.root;
2673        let rows = vec![
2674            (graft_key(0x0f, 0), b"before".to_vec()),
2675            (graft_key(0x10, 10), b"overlap".to_vec()),
2676        ];
2677        assert!(matches!(s.graft_range(rows.into_iter()), Err(crate::Error::RangeNotEmpty)));
2678        assert_eq!(s.root, old_root);
2679        assert_eq!(s.scan(&[]).unwrap().count(), base.len());
2680        for (key, value) in base {
2681            assert_eq!(s.get(&key).unwrap().as_deref(), Some(value.as_slice()));
2682        }
2683    }
2684
2685    #[test]
2686    fn graft_into_an_empty_tree_becomes_the_tree() {
2687        let d = tempfile::tempdir().unwrap();
2688        let mut s = Store::create(d.path(), cfg()).unwrap();
2689        let rows = graft_rows();
2690        s.graft_range(rows.clone().into_iter()).unwrap();
2691        assert_eq!(s.scan(&[]).unwrap().count(), rows.len());
2692        for (key, value) in rows.iter().step_by(97) {
2693            assert_eq!(s.get(key).unwrap().as_deref(), Some(value.as_slice()));
2694        }
2695    }
2696
2697    fn seed_query_keys(rows: &[(Vec<u8>, Vec<u8>)]) -> Vec<(Vec<u8>, Vec<u8>)> {
2698        let mut keys = rows.to_vec();
2699        keys.push((graft_key(0x10, 0), Vec::new()));
2700        keys.push((graft_key(0x10, 1_499), Vec::new()));
2701        keys.push((graft_key(0x30, 0), Vec::new()));
2702        keys.push((graft_key(0x30, 1_499), Vec::new()));
2703        keys.push((graft_key(0x20, 2_001), Vec::new()));
2704        keys
2705    }
2706
2707    /// Packed bytes have no WAL. Before the verified subtree is linked, a
2708    /// crash must reopen the old generation; after the link returns, the new
2709    /// generation must already be durable without a caller checkpoint.
2710    #[test]
2711    fn graft_crash_boundary_is_old_before_and_durable_after() {
2712        let before_dir = tempfile::tempdir().unwrap();
2713        {
2714            let mut s = Store::create(before_dir.path(), cfg()).unwrap();
2715            seed_graft_base(&mut s);
2716            let old_root = s.root;
2717            let result = s.graft_range_with_before_publish(
2718                graft_rows().into_iter(),
2719                |_, _| Err(std::io::Error::other("crash before graft").into()),
2720            );
2721            assert!(result.is_err());
2722            assert_eq!(s.root, old_root, "a pre-publication failure changed the live root");
2723        }
2724        let before = Store::open(before_dir.path(), cfg()).unwrap();
2725        assert!(before.get(&graft_key(0x20, 17)).unwrap().is_none());
2726        assert_eq!(before.scan(&[]).unwrap().count(), 3_000);
2727
2728        let after_dir = tempfile::tempdir().unwrap();
2729        {
2730            let mut s = Store::create(after_dir.path(), cfg()).unwrap();
2731            seed_graft_base(&mut s);
2732            s.graft_range(graft_rows().into_iter()).unwrap();
2733            // No commit and no extra checkpoint.
2734        }
2735        let after = Store::open(after_dir.path(), cfg()).unwrap();
2736        assert_eq!(after.get(&graft_key(0x20, 17)).unwrap().as_deref(), Some(&b"graft-17"[..]));
2737        assert_eq!(after.scan(&[]).unwrap().count(), 5_000);
2738    }
2739
2740    /// Kill after packing, resume and publish, kill after publication, then
2741    /// resume again before cleanup. The second resume must recognize that the
2742    /// exact generation is already authoritative rather than grafting twice.
2743    #[test]
2744    fn prepared_graft_publishes_after_reopen_and_resume_of_resume_is_idempotent() {
2745        let d = tempfile::tempdir().unwrap();
2746        let scratch = d.path().join("prepared-graft-scratch");
2747        let manifest = d.path().join("prepared-graft");
2748        {
2749            let mut s = Store::create(d.path(), cfg()).unwrap();
2750            seed_graft_base(&mut s);
2751            let mut rows = graft_rows();
2752            rows.sort_by(|a, b| a.0.cmp(&b.0));
2753            let min = rows.first().unwrap().0.clone();
2754            let max = rows.last().unwrap().0.clone();
2755            let prepared = s.prepare_graft_candidate(
2756                rows.into_iter().map(|(k, v)| Ok((k, v, false))),
2757                2_000, min, max, &scratch).unwrap();
2758            prepared.write_manifest(&manifest).unwrap();
2759            assert!(s.get(&graft_key(0x20, 17)).unwrap().is_none(),
2760                "preparing a candidate must not publish it in the live handle");
2761            // Simulated kill before root publication.
2762        }
2763        std::fs::write(manifest.with_extension("tmp"), b"torn next candidate").unwrap();
2764        let prepared = PreparedGraft::read_manifest(&manifest).unwrap();
2765        {
2766            let mut resumed = Store::open(d.path(), cfg()).unwrap();
2767            assert!(resumed.get(&graft_key(0x20, 17)).unwrap().is_none());
2768            resumed.publish_existing_candidate(&prepared).unwrap();
2769            assert_eq!(resumed.get(&graft_key(0x20, 17)).unwrap().as_deref(),
2770                Some(&b"graft-17"[..]));
2771            // Simulated kill after checkpoint, before manifest cleanup.
2772        }
2773        {
2774            let mut resumed_again = Store::open(d.path(), cfg()).unwrap();
2775            let generation = resumed_again.generation;
2776            resumed_again.publish_existing_candidate(&prepared).unwrap();
2777            assert_eq!(resumed_again.generation, generation,
2778                "resume-of-resume must not publish another generation");
2779            assert_eq!(resumed_again.scan(&[]).unwrap().count(), 5_000);
2780        }
2781    }
2782
2783    #[test]
2784    fn prepared_graft_manifest_and_generation_are_both_enforced() {
2785        let d = tempfile::tempdir().unwrap();
2786        let scratch = d.path().join("prepared-graft-scratch");
2787        let mut s = Store::create(d.path(), cfg()).unwrap();
2788        seed_graft_base(&mut s);
2789        let mut rows = graft_rows();
2790        rows.sort_by(|a, b| a.0.cmp(&b.0));
2791        let prepared = s.prepare_graft_candidate(
2792            rows.clone().into_iter().map(|(k, v)| Ok((k, v, false))), 2_000,
2793            rows.first().unwrap().0.clone(), rows.last().unwrap().0.clone(), &scratch).unwrap();
2794        let mut damaged = prepared.encode().unwrap();
2795        damaged[20] ^= 0x80;
2796        assert!(PreparedGraft::decode(&damaged).is_err());
2797
2798        let mut wrong_generation = prepared.clone();
2799        wrong_generation.base_generation += 1;
2800        wrong_generation.write_generation += 1;
2801        assert!(s.publish_existing_candidate(&wrong_generation).is_err(),
2802            "candidate pages stamped for one generation must not publish in another");
2803        assert!(s.get(&graft_key(0x20, 17)).unwrap().is_none());
2804    }
2805
2806    /// Verification is through an independent reopen. Damage a packed page
2807    /// after the pool flush but before publication: the old root must remain
2808    /// authoritative and the damaged candidate must never become reachable.
2809    #[test]
2810    fn corrupt_packed_range_is_caught_and_never_published() {
2811        let d = tempfile::tempdir().unwrap();
2812        let mut s = Store::create(d.path(), cfg()).unwrap();
2813        seed_graft_base(&mut s);
2814        let old_root = s.root;
2815
2816        let result = s.graft_range_with_before_publish(
2817            graft_rows().into_iter(),
2818            |data, packed| {
2819                use std::io::{Seek, SeekFrom, Write};
2820                let mut file = std::fs::OpenOptions::new().write(true).open(data)?;
2821                file.seek(SeekFrom::Start(packed.root as u64 * PAGE_SIZE as u64 + 80))?;
2822                file.write_all(&[0xa5])?;
2823                file.sync_all()?;
2824                Ok(())
2825            },
2826        );
2827
2828        assert!(matches!(result, Err(crate::Error::Corrupt { .. })),
2829            "a corrupt packed page must be refused, got {result:?}");
2830        assert_eq!(s.root, old_root);
2831        assert_eq!(s.get(&graft_key(0x10, 17)).unwrap().as_deref(), Some(&b"base-10-17"[..]));
2832        assert!(s.get(&graft_key(0x20, 17)).unwrap().is_none());
2833    }
2834
2835    /// An oversized record must be refused BEFORE it reaches the log, and
2836    /// the refusal must cost the caller nothing else -- neither `put` nor
2837    /// `delete` may leave a frame behind that nothing will ever apply.
2838    /// `delete`'s bound is on its PAYLOAD (`2 + k.len()`), not its key: the
2839    /// two differ by exactly the two bytes that made a 4,052-byte key
2840    /// produce a 4,054-byte frame (Task 17 final review, F1).
2841    #[test]
2842    fn an_oversized_record_never_reaches_the_log() {
2843        let d = tempfile::tempdir().unwrap();
2844        let mut s = Store::create(d.path(), cfg()).unwrap();
2845        s.put(b"a", b"1").unwrap();
2846        s.commit().unwrap();
2847        let before = s.wal.as_ref().unwrap().end_offset();
2848
2849        let max = crate::page::MAX_RECORD_LEN;
2850        // A huge VALUE is legal now (overflow chains) -- what must never reach
2851        // the log is a record nothing could apply: a KEY too big for a leaf
2852        // even as a marker record.
2853        let k = vec![b'k'; max];
2854        assert!(matches!(s.put(&k, b"v"), Err(crate::Error::TooLarge)));
2855        // A delete whose PAYLOAD is over the bound while its KEY is not:
2856        // the exact two-byte band that bricked a store in round 3.
2857        let k = vec![b'k'; max - 1];
2858        assert!(!s.delete(&k).unwrap(),
2859                "a key that long can never have been inserted, so `not found` is the truth");
2860        // Each key fits alone, but their encoded batch would exceed the
2861        // classifier's one-frame bound. Refuse the whole unit before WAL.
2862        let half = vec![b'b'; max / 2];
2863        assert!(matches!(
2864            s.put_empty_batch(&[half.clone(), half]),
2865            Err(crate::Error::TooLarge)
2866        ));
2867        assert_eq!(
2868            s.wal.as_ref().unwrap().end_offset(), before,
2869            "neither refusal may append a single byte to the log"
2870        );
2871    }
2872}