Skip to main content

kernel/
bulk.rs

1//! Building a large tree by repeated insertion means one descent and one random
2//! write per key. Sorting first and packing bottom-up means one sequential pass.
3//!
4//! Peak RAM is the arena in phase 1 and `MAX_FANOUT x buffer` in phase 2. Both
5//! are chosen numbers, neither is proportional to the input. That is Law 1 by
6//! construction rather than by care: without a fanout cap, `iter` would open
7//! every run at once and merge memory would be `run_count * buffer`, and
8//! `run_count` is `input / arena` -- linear in the input, correct at the test
9//! scale and wrong at the target one. `SortedRuns::merge_down` bounds it by
10//! merging runs down to at most `MAX_FANOUT` before the final pass.
11//!
12//! Phase 3, `pack_tree`, is bounded the same way and for the same reason. It
13//! used to accumulate one `(first_key, page_no)` pair per page of the level it
14//! was building into a `Vec`, which is one entry per page and therefore linear
15//! in the rows -- correct at the test scale, wrong at the target one. Measured
16//! on a 209M-row / 50.08 GiB load under a 1 GiB cap, that vector was several
17//! hundred MiB of the 539 MiB `rss_anon` the process peaked at, against a
18//! 113 MiB pool. Each level is now spilled to a scratch file as it is built and
19//! streamed back to build the level above it, so peak pack memory is one write
20//! buffer plus one read buffer whatever the input.
21//!
22//! SACRIFICE (Law 4): the input is written to temporary files and read back, so
23//! a bulk load costs roughly 3x the data in sequential I/O and needs scratch disk
24//! comparable to the input, plus one more read-and-rewrite pass of the whole
25//! dataset for every factor of `MAX_FANOUT` the run count exceeds it. Bought:
26//! linear build time instead of a curve, with merge memory that stays fixed
27//! regardless of how large the input grows.
28//!
29//! SACRIFICE (Law 4), scratch integrity: every temporary record carries four
30//! checksum bytes and is checksummed once when written and once when read.
31//! Extra merge passes repeat that work. Bought: a changed scratch byte cannot
32//! become a checksummed page and then pass the publication row-count check.
33//!
34//! SACRIFICE (Law 4), spilled separators: each level's separators are written
35//! once and read once instead of being held. That is `16 + key_len` bytes per
36//! PAGE of the level, not per row -- for 8-byte keys, 24 bytes per 4096-byte
37//! page, so the level-0 file is about 0.59% of the tree it describes and every
38//! level above it is another ~1/200th of that. Measured at 0.586% of the tree
39//! (265,080 B of scratch against a 45.2 MB tree) by
40//! `pack_tree_spill_is_a_small_fraction_of_the_tree` in
41//! kernel/tests/pack_shape.rs, which polls the scratch directory while the pack
42//! runs rather than deriving the figure from the record layout. Scaled to the
43//! 50.08 GiB reference load that is ~300 MiB of extra sequential I/O. Bought: pack memory that does not grow with rows
44//! at all, which is the whole premise of indexing 50 GB under a 1 GiB cap.
45
46use crate::page::{PageKind, PageMut, HEADER_LEN, PAGE_SIZE};
47use crate::pool::BufferPool;
48use crate::io::AlignedRegion;
49use crate::{Error, Result};
50use std::collections::BinaryHeap;
51use std::fs::File;
52use std::io::{BufReader, BufWriter, Read, Write};
53use std::path::{Path, PathBuf};
54
55/// Run-file framing. The top bit of vlen carries "this value is an overflow
56/// MARKER" -- explicit in the framing, never inside the value bytes, because an
57/// in-band tag was tried and ate the first byte of every value on recover's
58/// path within the hour. Real vlen is far below 2^31.
59const MARKER_BIT: u32 = 1 << 31;
60const SCRATCH_HEADER_LEN: usize = 12;
61
62fn invalid_scratch(why: &'static str) -> std::io::Error {
63    std::io::Error::new(std::io::ErrorKind::InvalidData, why)
64}
65
66fn write_item(w: &mut impl Write, k: &[u8], v: &[u8], marker: bool) -> std::io::Result<()> {
67    let kl = u32::try_from(k.len()).map_err(|_| invalid_scratch("scratch key is too long"))?;
68    let vl = u32::try_from(v.len()).map_err(|_| invalid_scratch("scratch value is too long"))?;
69    if vl & MARKER_BIT != 0 {
70        return Err(invalid_scratch("scratch value is too long for the marker bit"));
71    }
72    let raw_vl = vl | if marker { MARKER_BIT } else { 0 };
73    let mut header = [0u8; SCRATCH_HEADER_LEN];
74    header[0..4].copy_from_slice(&kl.to_le_bytes());
75    header[4..8].copy_from_slice(&raw_vl.to_le_bytes());
76    let mut crc = crc32c::crc32c(&header[..8]);
77    crc = crc32c::crc32c_append(crc, k);
78    crc = crc32c::crc32c_append(crc, v);
79    header[8..12].copy_from_slice(&crc.to_le_bytes());
80    w.write_all(&header)?;
81    w.write_all(k)?;
82    w.write_all(v)
83}
84
85fn read_item(r: &mut impl Read, max_key_len: usize, max_val_len: usize)
86    -> std::io::Result<Option<(Vec<u8>, Vec<u8>, bool)>>
87{
88    let mut h = [0u8; SCRATCH_HEADER_LEN];
89    match r.read_exact(&mut h[..1]) {
90        Ok(()) => {}
91        Err(e) if e.kind() == std::io::ErrorKind::UnexpectedEof => return Ok(None),
92        Err(e) => return Err(e),
93    }
94    // Once even one header byte exists, every other byte is mandatory.  This
95    // is what distinguishes an exact record boundary from a torn header.
96    r.read_exact(&mut h[1..])?;
97    let kl = u32::from_le_bytes(h[0..4].try_into().unwrap()) as usize;
98    let raw = u32::from_le_bytes(h[4..8].try_into().unwrap());
99    let marker = raw & MARKER_BIT != 0;
100    let vl = (raw & !MARKER_BIT) as usize;
101    // These maxima came from the in-memory records handed to the writer, not
102    // from this file.  Validate both disk lengths before either controls an
103    // allocation; a corrupted u32 must not request gigabytes of RAM.
104    if kl > max_key_len || vl > max_val_len {
105        return Err(invalid_scratch("scratch record length exceeds its writer's bound"));
106    }
107    let mut k = Vec::new();
108    k.try_reserve_exact(kl)
109        .map_err(|_| invalid_scratch("scratch key allocation exceeds address space"))?;
110    k.resize(kl, 0);
111    r.read_exact(&mut k)?;
112    let mut v = Vec::new();
113    v.try_reserve_exact(vl)
114        .map_err(|_| invalid_scratch("scratch value allocation exceeds address space"))?;
115    v.resize(vl, 0);
116    r.read_exact(&mut v)?;
117    let want_crc = u32::from_le_bytes(h[8..12].try_into().unwrap());
118    let mut crc = crc32c::crc32c(&h[..8]);
119    crc = crc32c::crc32c_append(crc, &k);
120    crc = crc32c::crc32c_append(crc, &v);
121    if crc != want_crc {
122        return Err(invalid_scratch("scratch record checksum mismatch"));
123    }
124    Ok(Some((k, v, marker)))
125}
126
127const SORT_MANIFEST_MAGIC: &[u8; 8] = b"KSRUN01\0";
128const SORT_MANIFEST_MAX: u64 = 16 << 20;
129
130struct SortManifest {
131    watermark: u64,
132    records: u64,
133    framed_bytes: u64,
134    max_key_len: usize,
135    max_val_len: usize,
136    runs: Vec<PathBuf>,
137    run_counts: Vec<u64>,
138}
139
140fn manifest_u64(bytes: &[u8], pos: &mut usize) -> std::io::Result<u64> {
141    let end = pos.checked_add(8).ok_or_else(|| invalid_scratch("sort manifest offset overflow"))?;
142    let raw = bytes.get(*pos..end).ok_or_else(|| invalid_scratch("sort manifest is truncated"))?;
143    *pos = end;
144    Ok(u64::from_le_bytes(raw.try_into().unwrap()))
145}
146
147fn write_sort_manifest(sort: &ExternalSort, generation: u64, watermark: u64) -> Result<()> {
148    let mut body = Vec::new();
149    body.extend_from_slice(SORT_MANIFEST_MAGIC);
150    for n in [generation, watermark, sort.records, sort.framed_bytes,
151              sort.max_key_len as u64, sort.max_val_len as u64, sort.runs.len() as u64] {
152        body.extend_from_slice(&n.to_le_bytes());
153    }
154    for (path, &count) in sort.runs.iter().zip(&sort.run_counts) {
155        let name = path.file_name().and_then(|s| s.to_str())
156            .ok_or_else(|| Error::Io(invalid_scratch("durable run has no UTF-8 file name")))?;
157        let name = name.as_bytes();
158        let len = u16::try_from(name.len())
159            .map_err(|_| Error::Io(invalid_scratch("durable run file name is too long")))?;
160        body.extend_from_slice(&len.to_le_bytes());
161        body.extend_from_slice(name);
162        body.extend_from_slice(&count.to_le_bytes());
163    }
164    let crc = crc32c::crc32c(&body);
165    body.extend_from_slice(&crc.to_le_bytes());
166    let tmp = sort.dir.join("manifest.tmp");
167    let published = sort.dir.join("manifest");
168    {
169        let mut file = File::create(&tmp)?;
170        file.write_all(&body)?;
171        crate::write_stats::add(
172            crate::write_stats::Phase::Manifest,
173            body.len() as u64,
174        );
175        file.sync_all()?;
176    }
177    std::fs::rename(&tmp, &published)?;
178    File::open(&sort.dir)?.sync_all()?;
179    Ok(())
180}
181
182fn read_sort_manifest(dir: &Path, expected_generation: u64) -> Result<SortManifest> {
183    let path = dir.join("manifest");
184    let len = std::fs::metadata(&path)?.len();
185    if len < (SORT_MANIFEST_MAGIC.len() + 7 * 8 + 4) as u64 || len > SORT_MANIFEST_MAX {
186        return Err(Error::Io(invalid_scratch("sort manifest length is invalid")));
187    }
188    let bytes = std::fs::read(&path)?;
189    let split = bytes.len() - 4;
190    let want = u32::from_le_bytes(bytes[split..].try_into().unwrap());
191    if crc32c::crc32c(&bytes[..split]) != want {
192        return Err(Error::Io(invalid_scratch("sort manifest checksum mismatch")));
193    }
194    if bytes.get(..8) != Some(SORT_MANIFEST_MAGIC) {
195        return Err(Error::Io(invalid_scratch("sort manifest magic mismatch")));
196    }
197    let mut pos = 8;
198    let generation = manifest_u64(&bytes[..split], &mut pos)?;
199    if generation != expected_generation {
200        return Err(Error::Io(invalid_scratch("sort manifest generation mismatch")));
201    }
202    let watermark = manifest_u64(&bytes[..split], &mut pos)?;
203    let records = manifest_u64(&bytes[..split], &mut pos)?;
204    let framed_bytes = manifest_u64(&bytes[..split], &mut pos)?;
205    let max_key_len = usize::try_from(manifest_u64(&bytes[..split], &mut pos)?)
206        .map_err(|_| Error::TooLarge)?;
207    let max_val_len = usize::try_from(manifest_u64(&bytes[..split], &mut pos)?)
208        .map_err(|_| Error::TooLarge)?;
209    let run_count = usize::try_from(manifest_u64(&bytes[..split], &mut pos)?)
210        .map_err(|_| Error::TooLarge)?;
211    // Every entry needs at least a u16 name length and u64 row count.
212    if run_count > (split.saturating_sub(pos)) / 10 {
213        return Err(Error::Io(invalid_scratch("sort manifest run count exceeds its bytes")));
214    }
215    let mut runs = Vec::with_capacity(run_count);
216    let mut run_counts = Vec::with_capacity(run_count);
217    for _ in 0..run_count {
218        let end = pos.checked_add(2).ok_or(Error::TooLarge)?;
219        let raw = bytes.get(pos..end)
220            .ok_or_else(|| Error::Io(invalid_scratch("sort manifest run name is truncated")))?;
221        pos = end;
222        let name_len = u16::from_le_bytes(raw.try_into().unwrap()) as usize;
223        let end = pos.checked_add(name_len).ok_or(Error::TooLarge)?;
224        let raw_name = bytes.get(pos..end)
225            .ok_or_else(|| Error::Io(invalid_scratch("sort manifest run name is truncated")))?;
226        pos = end;
227        let name = std::str::from_utf8(raw_name)
228            .map_err(|_| Error::Io(invalid_scratch("sort manifest run name is not UTF-8")))?;
229        if name.contains('/') || name.contains('\\') || name == "." || name == ".." {
230            return Err(Error::Io(invalid_scratch("sort manifest run name escapes its directory")));
231        }
232        let path = dir.join(name);
233        if !path.is_file() {
234            return Err(Error::Io(invalid_scratch("sort manifest names a missing run")));
235        }
236        runs.push(path);
237        run_counts.push(manifest_u64(&bytes[..split], &mut pos)?);
238    }
239    if pos != split {
240        return Err(Error::Io(invalid_scratch("sort manifest has trailing bytes")));
241    }
242    if run_counts.iter().try_fold(0u64, |sum, &n| sum.checked_add(n)) != Some(records) {
243        return Err(Error::Io(invalid_scratch("sort manifest row counts disagree")));
244    }
245    Ok(SortManifest { watermark, records, framed_bytes, max_key_len, max_val_len,
246        runs, run_counts })
247}
248
249/// Removes its scratch directory if it is dropped without `finish()`.
250///
251/// `push` can fail, and `Store::bulk_load` returns through `?` when it does
252/// -- before `finish()` has produced the `SortedRuns` whose own `Drop` would
253/// have cleaned up. Without this, every failed bulk load leaves a temp
254/// directory behind for the life of the machine.
255pub struct ExternalSort {
256    finished: bool,
257    dir: PathBuf,
258    arena: Vec<(Vec<u8>, Vec<u8>, bool)>,
259    arena_bytes: usize,
260    used: usize,
261    runs: Vec<PathBuf>,
262    run_counts: Vec<u64>,
263    max_key_len: usize,
264    max_val_len: usize,
265    records: u64,
266    framed_bytes: u64,
267    /// A durable build generation retains runs on Drop and publishes a
268    /// reopenable manifest after each explicit watermark checkpoint.
269    durable_generation: Option<u64>,
270}
271
272impl ExternalSort {
273    pub fn new(dir: &Path, arena_bytes: usize) -> Result<Self> {
274        Self::new_inner(dir, arena_bytes, None)
275    }
276
277    /// Create a sorter whose runs survive process death. `generation` is the
278    /// build generation, not the store page generation; reopen refuses a
279    /// manifest from a different build so stale scratch cannot be attached to
280    /// a later CREATE INDEX using the same directory.
281    pub fn new_durable(dir: &Path, arena_bytes: usize, generation: u64) -> Result<Self> {
282        Self::new_inner(dir, arena_bytes, Some(generation))
283    }
284
285    fn new_inner(dir: &Path, arena_bytes: usize, durable_generation: Option<u64>) -> Result<Self> {
286        std::fs::create_dir_all(dir)?;
287        Ok(ExternalSort {
288            finished: false, dir: dir.to_path_buf(), arena: Vec::new(), arena_bytes, used: 0,
289            runs: Vec::new(), run_counts: Vec::new(), max_key_len: 0, max_val_len: 0,
290            records: 0, framed_bytes: 0, durable_generation,
291        })
292    }
293
294    /// Push with the overflow-marker flag carried explicitly through the run
295    /// files. `val` is whatever the caller's pipeline stores (recover tags it
296    /// with a page number); the flag survives sort and merge untouched.
297    pub fn push_flagged(&mut self, key: Vec<u8>, val: Vec<u8>, marker: bool) -> Result<()> {
298        self.push_inner(key, val, marker)
299    }
300
301    pub fn push(&mut self, key: Vec<u8>, val: Vec<u8>) -> Result<()> {
302        self.push_inner(key, val, false)
303    }
304
305    fn push_inner(&mut self, key: Vec<u8>, val: Vec<u8>, marker: bool) -> Result<()> {
306        self.records = self.records.checked_add(1).ok_or(Error::TooLarge)?;
307        self.framed_bytes = self.framed_bytes
308            .checked_add((SCRATCH_HEADER_LEN + key.len() + val.len()) as u64)
309            .ok_or(Error::TooLarge)?;
310        self.max_key_len = self.max_key_len.max(key.len());
311        self.max_val_len = self.max_val_len.max(val.len());
312        self.used += key.len() + val.len() + 48;   // 48 = two Vec headers, approx
313        self.arena.push((key, val, marker));
314        if self.used >= self.arena_bytes { self.spill()?; }
315        Ok(())
316    }
317
318    fn spill(&mut self) -> Result<()> {
319        if self.arena.is_empty() { return Ok(()); }
320        self.arena.sort_by(|a, b| a.0.cmp(&b.0));
321        let path = match self.durable_generation {
322            Some(generation) => self.dir.join(format!("g{generation:016x}-run-{:05}.dat", self.runs.len())),
323            None => self.dir.join(format!("run-{:05}.tmp", self.runs.len())),
324        };
325        let mut w = BufWriter::new(File::create(&path)?);
326        let count = u64::try_from(self.arena.len()).map_err(|_| Error::TooLarge)?;
327        let mut written = 0u64;
328        for (k, v, m) in self.arena.drain(..) {
329            written = written.checked_add(
330                (SCRATCH_HEADER_LEN + k.len() + v.len()) as u64,
331            ).ok_or(Error::TooLarge)?;
332            write_item(&mut w, &k, &v, m)?;
333        }
334        w.flush()?;
335        crate::write_stats::add(crate::write_stats::Phase::SortScratch, written);
336        if self.durable_generation.is_some() { w.get_ref().sync_all()?; }
337        self.runs.push(path);
338        self.run_counts.push(count);
339        self.used = 0;
340        Ok(())
341    }
342
343    /// Seal the current arena as one checksummed run without finishing the
344    /// sorter. Chunked SQL builders call this at their scan watermark so peak
345    /// live heap is bounded by the fixed chunk even when the configured arena
346    /// is larger (the arena remains the hard upper ceiling).
347    pub fn flush_run(&mut self) -> Result<()> { self.spill() }
348
349    /// Mechanism counters for load profiling.  `framed_bytes` is the exact
350    /// first-pass scratch payload (including each record header), excluding
351    /// any extra merge-down pass.  The pending arena counts as one future run.
352    pub fn profile(&self) -> (u64, u64, usize) {
353        (self.records, self.framed_bytes,
354         self.runs.len() + usize::from(!self.arena.is_empty()))
355    }
356
357    /// Flush one scan watermark and atomically publish the complete run list.
358    /// The returned cost includes the run fsync, manifest fsync, rename and
359    /// directory fsync: exactly the durability tax paid at this interval.
360    pub fn checkpoint(&mut self, watermark: u64) -> Result<std::time::Duration> {
361        let started = std::time::Instant::now();
362        let generation = self.durable_generation.ok_or_else(|| Error::Io(
363            invalid_scratch("checkpoint requested for a non-durable sorter")))?;
364        self.spill()?;
365        write_sort_manifest(self, generation, watermark)?;
366        Ok(started.elapsed())
367    }
368
369    /// Reopen a durable sorter at its last completely published watermark.
370    /// Torn `.tmp` manifests are ignored; every listed run is later checked by
371    /// the ordinary framed CRC and exact record-count reader.
372    pub fn reopen_durable(dir: &Path, arena_bytes: usize, expected_generation: u64)
373        -> Result<(Self, u64)>
374    {
375        let manifest = read_sort_manifest(dir, expected_generation)?;
376        let mut sort = Self::new_inner(dir, arena_bytes, Some(expected_generation))?;
377        sort.runs = manifest.runs;
378        sort.run_counts = manifest.run_counts;
379        sort.max_key_len = manifest.max_key_len;
380        sort.max_val_len = manifest.max_val_len;
381        sort.records = manifest.records;
382        sort.framed_bytes = manifest.framed_bytes;
383        Ok((sort, manifest.watermark))
384    }
385
386    /// Successful publication owns the cleanup decision. Until this is
387    /// called, dropping the sorter deliberately leaves its checkpoint intact.
388    pub fn discard_durable(mut self) -> Result<()> {
389        self.finished = true;
390        std::fs::remove_dir_all(&self.dir)?;
391        Ok(())
392    }
393
394    pub fn finish(mut self) -> Result<SortedRuns> {
395        self.spill()?;
396        self.finished = true;
397        Ok(SortedRuns {
398            dir: std::mem::take(&mut self.dir), runs: std::mem::take(&mut self.runs),
399            owns_dir: self.durable_generation.is_none(),
400            run_counts: std::mem::take(&mut self.run_counts),
401            max_key_len: self.max_key_len, max_val_len: self.max_val_len,
402        })
403    }
404}
405
406impl Drop for ExternalSort {
407    fn drop(&mut self) {
408        if !self.finished && self.durable_generation.is_none() {
409            let _ = std::fs::remove_dir_all(&self.dir);
410        }
411    }
412}
413
414pub struct SortedRuns {
415    dir: PathBuf,
416    runs: Vec<PathBuf>,
417    /// Trusted record count beside each run, so truncation exactly at a
418    /// record boundary cannot masquerade as the run ending normally.
419    run_counts: Vec<u64>,
420    owns_dir: bool,
421    /// Trusted allocation bounds retained from the caller's in-memory input.
422    max_key_len: usize,
423    max_val_len: usize,
424}
425
426impl SortedRuns {
427    pub fn run_count(&self) -> usize { self.runs.len() }
428
429    /// Remove a durable sort checkpoint after its graft and final metadata
430    /// publication are known durable. Ordinary temporary runs already clean
431    /// themselves on Drop; accepting both shapes keeps caller cleanup simple.
432    pub fn discard(mut self) -> Result<()> {
433        self.owns_dir = false;
434        if self.dir.exists() { std::fs::remove_dir_all(&self.dir)?; }
435        Ok(())
436    }
437
438    /// The most runs `iter` will hold open at once.
439    ///
440    /// Without a cap, `iter` opens EVERY run with a read buffer each, so merge
441    /// memory is `run_count * buffer` -- and `run_count` is `input / arena`,
442    /// which makes it linear in the input. Law 1 forbids that: it happens to
443    /// fit at moderate scale and stops fitting above it, which is precisely
444    /// the shape this project exists to avoid -- correct at the test scale,
445    /// wrong at the target one, invisible until someone's data grows.
446    pub const MAX_FANOUT: usize = 64;
447
448    /// Merge groups of runs into intermediate runs until at most `MAX_FANOUT`
449    /// remain, so the final merge's memory is bounded by a chosen number
450    /// rather than by the input.
451    ///
452    /// SACRIFICE (Law 4): each extra pass reads and rewrites the whole
453    /// dataset once more. Bought: merge memory that does not grow with the
454    /// input at all.
455    fn merge_down(&mut self) -> Result<()> {
456        // A monotonic counter across the WHOLE call, not just within one
457        // pass. `next.len()`/`gi` alone repeat every pass (`next` always
458        // starts at 0), so a name built only from them collides with a
459        // SURVIVOR of the previous pass sitting at the same position in
460        // `self.runs` -- concretely, pass 2's group 0 output and pass 1's
461        // group 0 output are both named `pass-0-00000.tmp`, and pass 1's
462        // output is typically pass 2's group 0's OWN first input. Creating
463        // `out` then truncates that input before it is read (silently
464        // losing every record it held), and the later cleanup below deletes
465        // the file `out` itself, since it was also listed as one of
466        // `group`'s members -- so the very output just written vanishes and
467        // the next `File::open` on it panics with "No such file or
468        // directory". Reachable on any input needing more than one
469        // `merge_down` pass (run_count > MAX_FANOUT^2), which is exactly the
470        // scale this bound exists to make safe -- found by forcing a
471        // two-pass `merge_down` in a scratch test, not by inspection alone.
472        let mut pass_no = 0usize;
473        while self.runs.len() > Self::MAX_FANOUT {
474            let mut next: Vec<PathBuf> = Vec::new();
475            let mut next_counts: Vec<u64> = Vec::new();
476            for (gi, start) in (0..self.runs.len()).step_by(Self::MAX_FANOUT).enumerate() {
477                let end = (start + Self::MAX_FANOUT).min(self.runs.len());
478                let group = &self.runs[start..end];
479                let group_count = self.run_counts[start..end].iter().try_fold(0u64, |total, &n| {
480                    total.checked_add(n).ok_or(Error::TooLarge)
481                })?;
482                if group.len() == 1 {
483                    next.push(group[0].clone());
484                    next_counts.push(group_count);
485                    continue;
486                }
487                let out = self.dir.join(format!("pass-{}-{}-{:05}.tmp", pass_no, next.len(), gi));
488                let mut w = BufWriter::new(File::create(&out)?);
489                // `owns_dir: false` -- this SortedRuns shares `self.dir` with
490                // its parent. If its Drop removed that directory the way the
491                // owning one does, finishing this group's merge would delete
492                // run files sibling groups (and the parent) still need. Same
493                // shape as the temp-directory race `Store::bulk_load` guards
494                // against with its per-call sequence number, one level in:
495                // here the collision is between a SortedRuns and the parent
496                // it was carved out of, not between two unrelated calls.
497                let mut part = SortedRuns {
498                    dir: self.dir.clone(), runs: group.to_vec(), owns_dir: false,
499                    run_counts: self.run_counts[start..end].to_vec(),
500                    max_key_len: self.max_key_len, max_val_len: self.max_val_len,
501                };
502                let mut written = 0u64;
503                for item in part.iter_unbounded()? {
504                    let (k, v, m) = item?;
505                    write_item(&mut w, &k, &v, m)?;
506                    written = written.checked_add(
507                        (SCRATCH_HEADER_LEN + k.len() + v.len()) as u64,
508                    ).ok_or(Error::TooLarge)?;
509                }
510                w.flush()?;
511                crate::write_stats::add(crate::write_stats::Phase::SortScratch, written);
512                for p in group { let _ = std::fs::remove_file(p); }
513                next.push(out);
514                next_counts.push(group_count);
515            }
516            self.runs = next;
517            self.run_counts = next_counts;
518            pass_no += 1;
519        }
520        Ok(())
521    }
522
523    /// K-way merge over at most `MAX_FANOUT` runs. RAM is one buffered
524    /// reader per open run, a chosen bound rather than an input-shaped one.
525    pub fn iter(&mut self) -> Result<MergeIter> {
526        self.merge_down()?;
527        self.iter_unbounded()
528    }
529
530    /// Open every remaining run and build the merge over them.
531    ///
532    /// Returns a `Result` because opening a file can fail for reasons that
533    /// have nothing to do with the caller: a permission change, a full disk,
534    /// a future regression that reintroduces a name collision. An earlier
535    /// version used `File::open(p).unwrap()` -- the filename-collision fix
536    /// removed the one *trigger* this project had found for that panic and
537    /// left the panic itself standing for every other cause, contradicting
538    /// the `pending_err`/`done` machinery two paragraphs below, which exists
539    /// specifically because nothing else in this crate aborts on a
540    /// recoverable I/O error. `iter`, and `merge_down`'s own internal use of
541    /// `iter_unbounded`, both propagate this with `?` rather than unwrap it.
542    ///
543    /// The actual k-way merge over whatever is currently in `self.runs`, with
544    /// no fanout bound of its own -- callers (`iter`, after `merge_down`, and
545    /// `merge_down` itself over one `<= MAX_FANOUT` group) are what keep the
546    /// run count this opens bounded.
547    fn iter_unbounded(&mut self) -> Result<MergeIter> {
548        let mut readers = Vec::with_capacity(self.runs.len());
549        // One fixed aggregate read arena, divided across the active fanout.
550        // A buffer per run with a fixed per-buffer size makes live heap grow
551        // with the number of corpus runs until MAX_FANOUT; bounded is not the
552        // same as independent of corpus size. Four KiB is one page at the
553        // maximum 64-way fanout, 256 KiB for a single-run replay.
554        let per_reader = (256 * 1024 / self.runs.len().max(1)).max(4096);
555        for p in &self.runs {
556            readers.push(BufReader::with_capacity(per_reader, File::open(p)?));
557        }
558        let mut m = MergeIter {
559            readers, heap: BinaryHeap::new(), pending_err: None, done: false,
560            max_key_len: self.max_key_len, max_val_len: self.max_val_len,
561            expected: self.run_counts.iter().try_fold(0u64, |total, &n| {
562                total.checked_add(n).ok_or(Error::TooLarge)
563            })?,
564            yielded: 0,
565        };
566        for i in 0..m.readers.len() {
567            // An I/O error this early (priming a run's first record) must
568            // not be swallowed as if the run were simply empty -- that would
569            // drop every key still sitting in it with no signal to the
570            // caller. Keep priming the rest so their handles stay primed,
571            // and hand back the first error once iteration starts.
572            if let Err(e) = m.pull(i) {
573                if m.pending_err.is_none() { m.pending_err = Some(e); }
574            }
575        }
576        Ok(m)
577    }
578}
579
580impl Drop for SortedRuns {
581    fn drop(&mut self) {
582        // Only the owner removes the directory. A temporary SortedRuns built
583        // over a subset of runs during `merge_down` must not delete the
584        // scratch its parent (or sibling groups) is still using.
585        if self.owns_dir { let _ = std::fs::remove_dir_all(&self.dir); }
586    }
587}
588
589/// Reversed ordering so BinaryHeap (a max-heap) yields the smallest key.
590struct Head { key: Vec<u8>, val: Vec<u8>, marker: bool, from: usize }
591impl PartialEq for Head { fn eq(&self, o: &Self) -> bool { self.key == o.key } }
592impl Eq for Head {}
593impl Ord for Head {
594    fn cmp(&self, o: &Self) -> std::cmp::Ordering { o.key.cmp(&self.key).then(o.from.cmp(&self.from)) }
595}
596impl PartialOrd for Head { fn partial_cmp(&self, o: &Self) -> Option<std::cmp::Ordering> { Some(self.cmp(o)) } }
597
598pub struct MergeIter {
599    readers: Vec<BufReader<File>>,
600    heap: BinaryHeap<Head>,
601    /// An I/O error observed while refilling from a run. `read_item` cannot
602    /// tell corruption/truncation apart from "the run legitimately ended" by
603    /// itself, so `pull` must not collapse a real error into `None` the way
604    /// the reference sketch did -- that reads as the run finishing early and
605    /// silently drops every key still sitting in it, with the merge and its
606    /// caller both reporting success. Carried here because `pull` runs both
607    /// from `next` (which can return it immediately) and from `SortedRuns::iter`
608    /// while priming (which cannot: the `Iterator` isn't constructed yet).
609    pending_err: Option<Error>,
610    /// Set once `pending_err` has been handed back. Same discipline as
611    /// `RangeIter::fail` in btree.rs: after an error, stop for good rather
612    /// than keep popping the heap and serving items from runs the error
613    /// didn't touch, which would look like the merge completed when an
614    /// unknown number of keys from the broken run were never emitted.
615    done: bool,
616    max_key_len: usize,
617    max_val_len: usize,
618    expected: u64,
619    yielded: u64,
620}
621
622impl MergeIter {
623    fn pull(&mut self, i: usize) -> Result<()> {
624        match read_item(&mut self.readers[i], self.max_key_len, self.max_val_len) {
625            Ok(Some((k, v, m))) => { self.heap.push(Head { key: k, val: v, marker: m, from: i }); Ok(()) }
626            Ok(None) => Ok(()),
627            Err(e) => Err(e.into()),
628        }
629    }
630}
631
632impl Iterator for MergeIter {
633    type Item = Result<(Vec<u8>, Vec<u8>, bool)>;
634    fn next(&mut self) -> Option<Self::Item> {
635        if self.done { return None; }
636        if let Some(e) = self.pending_err.take() {
637            self.done = true;
638            return Some(Err(e));
639        }
640        let Some(h) = self.heap.pop() else {
641            self.done = true;
642            if self.yielded == self.expected { return None; }
643            return Some(Err(Error::Io(invalid_scratch(
644                "scratch run ended before its recorded item count",
645            ))));
646        };
647        if self.yielded >= self.expected {
648            self.done = true;
649            return Some(Err(Error::Io(invalid_scratch(
650                "scratch run exceeded its recorded item count",
651            ))));
652        }
653        self.yielded += 1;
654        if let Err(e) = self.pull(h.from) { self.pending_err = Some(e); }
655        Some(Ok((h.key, h.val, h.marker)))
656    }
657}
658
659fn enc_leaf(key: &[u8], val: &[u8], compact: bool) -> Vec<u8> {
660    crate::btree::enc_leaf(key,val,compact)
661}
662
663fn enc_interior(key: &[u8], child: u32) -> Vec<u8> {
664    let mut r = Vec::with_capacity(6 + key.len());
665    r.extend_from_slice(&(key.len() as u16).to_le_bytes());
666    r.extend_from_slice(key);
667    r.extend_from_slice(&child.to_le_bytes());
668    r
669}
670
671/// Build one page's full contents in a scratch buffer, then return it whole.
672///
673/// Mirrors `btree.rs`'s `build_page`: writing straight into a live pinned
674/// frame and then hitting a fallible `insert_slot` would leave that frame
675/// wiped, dirty and never finalised -- a stale checksum over a blank page,
676/// which the next read refuses, losing every key that was meant to land on
677/// it. Building in scratch first removes the class rather than relying on
678/// the packing arithmetic below never being wrong.
679fn build_scratch(kind: PageKind, tree_id: u16, page_no: u32, recs: &[Vec<u8>]) -> Result<Vec<u8>> {
680    let mut scratch = vec![0u8; PAGE_SIZE];
681    {
682        let mut p = PageMut::init(&mut scratch, kind, tree_id, page_no);
683        for r in recs {
684            let at = p.nentries_pub();
685            p.insert_slot(at, r)?;
686        }
687        p.finalise(0);
688    }
689    Ok(scratch)
690}
691
692/// Fixed 256 KiB write-behind used only for unreachable packed candidates.
693/// Page numbers from the append allocator are contiguous in the common case;
694/// a freelist discontinuity flushes the current run and starts another.
695struct PackedPageWriter<'a> {
696    pool: &'a BufferPool,
697    region: AlignedRegion,
698    first: Option<u32>,
699    pages: usize,
700}
701
702impl<'a> PackedPageWriter<'a> {
703    // The 64-page write aggregation did not move the 250k end-to-end wall
704    // (9.21s baseline, 9.20s one-page ablation) and the combined server 1M
705    // result remained flat. Keep the uncached candidate-page path, which
706    // avoids polluting the shared pool, but omit the unearned aggregation.
707    const CAPACITY: usize = 1;
708
709    fn new(pool: &'a BufferPool) -> Result<Self> {
710        Ok(Self {
711            pool,
712            region: AlignedRegion::new(Self::CAPACITY * PAGE_SIZE)?,
713            first: None,
714            pages: 0,
715        })
716    }
717
718    fn push(&mut self, page_no: u32, scratch: &[u8]) -> Result<()> {
719        if scratch.len() != PAGE_SIZE { return Err(Error::TooLarge); }
720        if self.pages == Self::CAPACITY
721            || self.first.is_some_and(|first| first + self.pages as u32 != page_no)
722        {
723            self.flush()?;
724        }
725        if self.first.is_none() { self.first = Some(page_no); }
726        // SAFETY: this writer exclusively owns the region and no slice escapes.
727        let page = unsafe { self.region.page_mut(self.pages) };
728        page.copy_from_slice(scratch);
729        crate::page::seal(page, self.pool.stamp_generation());
730        self.pages += 1;
731        Ok(())
732    }
733
734    fn flush(&mut self) -> Result<()> {
735        if self.pages == 0 { return Ok(()); }
736        let len = self.pages * PAGE_SIZE;
737        // SAFETY: no mutable page slice is live; the write completes before
738        // the region can be reused.
739        let bytes = unsafe { self.region.prefix(len) };
740        self.pool.write_unpooled_run(self.first.unwrap(), bytes)?;
741        self.first = None;
742        self.pages = 0;
743        Ok(())
744    }
745}
746
747/// Where a packed page goes, and therefore which durability rule it obeys.
748///
749/// `Direct` reserves an unpooled page number and writes the finished, sealed
750/// page straight into the data file. Those pages are unreachable until the
751/// root swap publishes them, which is what lets `Store::graft_range` skip the
752/// record WAL for them entirely (D11).
753///
754/// `Pooled` allocates ordinary pool pages instead. A page-WAL database then
755/// logs every packed page as a normal frame and publishes the whole graft
756/// with the ordinary commit, so there is no root swap, no skipped log and no
757/// second durability rule to reason about: Law 3 and Law 5 hold exactly as
758/// they do for a single `insert`. The price is one WAL frame per packed page
759/// -- the same frame an ordinary insert would have paid for that leaf anyway,
760/// but paid once instead of once per commit the leaf is dirty in.
761///
762/// A page is RESERVED before it is filled, because a leaf's `next_leaf`
763/// pointer is only known once the following leaf has a number. The pooled
764/// sink therefore holds the write guard between `reserve` and `push` (at most
765/// two at a time: the leaf being finished and the leaf that follows it),
766/// which keeps every packed page a single pooled write with no read-back.
767pub(crate) enum PageSink<'a> {
768    Direct(PackedPageWriter<'a>),
769    Pooled { pool: &'a BufferPool, held: Vec<crate::pool::PinnedWrite<'a>> },
770}
771
772impl<'a> PageSink<'a> {
773    pub(crate) fn direct(pool: &'a BufferPool) -> Result<Self> {
774        Ok(PageSink::Direct(PackedPageWriter::new(pool)?))
775    }
776
777    pub(crate) fn pooled(pool: &'a BufferPool) -> Self {
778        PageSink::Pooled { pool, held: Vec::new() }
779    }
780
781    fn reserve(&mut self) -> Result<u32> {
782        match self {
783            PageSink::Direct(writer) => writer.pool.allocate_unpooled(),
784            PageSink::Pooled { pool, held } => {
785                let guard = pool.allocate()?;
786                let page_no = guard.page_no();
787                held.push(guard);
788                Ok(page_no)
789            }
790        }
791    }
792
793    fn push(&mut self, page_no: u32, scratch: &[u8]) -> Result<()> {
794        match self {
795            PageSink::Direct(writer) => writer.push(page_no, scratch),
796            PageSink::Pooled { held, .. } => {
797                if scratch.len() != PAGE_SIZE { return Err(Error::TooLarge); }
798                // A push for a number this sink never reserved would write a
799                // page the allocator still considers free: refuse instead of
800                // guessing, the same way the direct writer refuses a run that
801                // exceeds its allocation.
802                let at = held.iter().position(|guard| guard.page_no() == page_no)
803                    .ok_or(Error::Corrupt { page_no, why: "packed page was never reserved" })?;
804                let mut guard = held.remove(at);
805                guard.bytes_mut().copy_from_slice(scratch);
806                Ok(())
807            }
808        }
809    }
810
811    /// No page may stay pinned past the pack: `flush_all` asserts that a
812    /// dirty frame has no live guard, so a guard held across the caller's
813    /// commit would be a panic, not a slow path. Guards still held here
814    /// belong to reserved-but-unfilled pages on an error path; dropping them
815    /// leaves each as the valid empty Free page `allocate` initialised.
816    fn finish(&mut self) -> Result<()> {
817        match self {
818            PageSink::Direct(writer) => writer.flush(),
819            PageSink::Pooled { held, .. } => { held.clear(); Ok(()) }
820        }
821    }
822}
823
824/// One level's separators, on disk instead of in RAM.
825///
826/// `pack_tree` builds a level of pages and needs, for each page, the pair
827/// `(first_key, page_no)` to build the level above it. Held in a `Vec`, that is
828/// one entry per page and so linear in the rows -- the same shape
829/// `MAX_FANOUT` above exists to forbid, one phase later. This holds the pairs
830/// in a file and hands back a reader, so the pack's memory is a write buffer
831/// plus a read buffer regardless of how many pages the level has.
832///
833/// `Drop` closes the writer and unlinks the file unconditionally. That is what
834/// makes the scratch directory empty after a FAILED pack as well as a
835/// successful one: every early return in `pack_tree` -- a duplicate key, an
836/// oversized record, a pool that cannot allocate -- drops the live
837/// `Separators` on its way out.
838struct Separators {
839    path: PathBuf,
840    /// `None` once `seal` has flushed and closed it.
841    w: Option<BufWriter<File>>,
842    /// Entries written. The pack needs only "is this level one page yet?",
843    /// which is `count == 1`, and "was there any input at all?", `count == 0`.
844    count: u64,
845    /// Trusted upper bound for a key length read back from this file.
846    max_key_len: usize,
847    framed_bytes: u64,
848}
849
850impl Separators {
851    /// `seq` is a counter that runs across the WHOLE pack and is never reset
852    /// per level, and `pack` is a per-call number from a process-wide atomic.
853    ///
854    /// A name built from the level number alone would happen to be unique
855    /// today, and that is exactly the reasoning that produced the
856    /// `merge_down` collision documented above: a name derived only from
857    /// position repeats the moment the loop producing it runs more than once,
858    /// and the file it collides with is typically one still being read, so
859    /// `File::create` truncates live data with no error anywhere. This
860    /// creates temp files across LEVELS, which is structurally the same
861    /// hazard, so the counter is monotonic across the whole call. `pack` and
862    /// the pid do for concurrent packs sharing a scratch directory what
863    /// `Store::bulk_load`'s per-call sequence number does for the sort, and
864    /// keep these names disjoint from `run-*.tmp` and `pass-*.tmp` besides.
865    fn create(dir: &Path, pack: u64, seq: &mut u64) -> Result<Self> {
866        let n = *seq;
867        *seq += 1;
868        let path = dir.join(format!("sep-{}-{}-{:05}.tmp", std::process::id(), pack, n));
869        let w = BufWriter::with_capacity(256 * 1024, File::create(&path)?);
870        Ok(Separators {
871            path,
872            w: Some(w),
873            count: 0,
874            max_key_len: 0,
875            framed_bytes: 0,
876        })
877    }
878
879    fn push(&mut self, key: &[u8], page_no: u32) -> Result<()> {
880        let w = self.w.as_mut().expect("separators pushed after seal");
881        write_item(w, key, &page_no.to_le_bytes(), false)?;
882        self.max_key_len = self.max_key_len.max(key.len());
883        self.count += 1;
884        self.framed_bytes = self.framed_bytes
885            .checked_add((SCRATCH_HEADER_LEN + key.len() + 4) as u64)
886            .ok_or(Error::TooLarge)?;
887        Ok(())
888    }
889
890    /// Flush and close. Explicit rather than left to `Drop`, because a
891    /// `BufWriter` dropped with bytes still buffered discards the write error
892    /// and the level would simply be short -- pages silently absent from the
893    /// tree, which is the failure mode this crate refuses everywhere else.
894    fn seal(&mut self) -> Result<()> {
895        if let Some(mut w) = self.w.take() {
896            w.flush()?;
897            crate::write_stats::add(
898                crate::write_stats::Phase::PackScratch,
899                self.framed_bytes,
900            );
901        }
902        Ok(())
903    }
904
905    fn reader(&self) -> Result<BufReader<File>> {
906        // `push` after `seal` is already a hard error; this is the mirror.
907        // Reading before `seal` would open the path and get whatever had
908        // happened to reach the filesystem -- silently short by up to a whole
909        // write buffer, which is the exact failure `short_read` below exists
910        // to catch and the exact reason `seal` is explicit at all. No current
911        // path does it; the seal discipline is load-bearing enough to say so.
912        debug_assert!(self.w.is_none(), "separators read before seal");
913        Ok(BufReader::with_capacity(256 * 1024, File::open(&self.path)?))
914    }
915}
916
917impl Drop for Separators {
918    fn drop(&mut self) {
919        // Writer first, then unlink -- and unlink even when `seal` already
920        // ran, since the file is scratch either way.
921        self.w = None;
922        let _ = std::fs::remove_file(&self.path);
923    }
924}
925
926/// A level file yielded fewer separators than were written to it.
927///
928/// `write_item` uses `write_all` into a `BufWriter`, `count` is incremented
929/// only after that returns `Ok`, and `seal` flushes explicitly so a buffered
930/// write error surfaces rather than leaving the file short -- so the write
931/// side is defended. The READ side was not, and `read_item` cannot tell the
932/// difference by itself: EOF at a record boundary reads as "the level ended"
933/// and EOF mid-record as an error, so a file short by a whole number of
934/// records is indistinguishable from a complete one.
935///
936/// What that costs is not a crash. Each missing separator is a page the level
937/// above never points at, so an entire subtree becomes unreachable by descent
938/// while the leaf sibling chain -- stitched independently of this file --
939/// still walks every leaf. `get` misses keys that `range` returns: exactly the
940/// point-lookup/scan divergence `Error::DuplicateKey` refuses inputs to
941/// prevent, arriving silently, with `bulk_load` returning `Ok`. Verified by
942/// injection: truncating the level-0 file by one whole record left 144 of
943/// 50,000 keys unreachable by `get` while `scan` still counted all 50,000.
944/// A full-scan row count -- the check the 209M-row load was verified with --
945/// cannot see it.
946///
947/// This file's own recorded history is a temp-name collision that truncated a
948/// file still being read, so "nothing can truncate it" is not a claim this
949/// code gets to make about itself. One `u64` and one comparison per level
950/// converts the whole class from silent to loud.
951fn short_level(seen: u64, want: u64) -> Error {
952    Error::Io(std::io::Error::new(
953        std::io::ErrorKind::InvalidData,
954        format!("a separator level came back with {seen} of {want} entries"),
955    ))
956}
957
958/// Read one `(first_key, page_no)` back off a level file.
959fn read_sep(r: &mut BufReader<File>, max_key_len: usize) -> Result<Option<(Vec<u8>, u32)>> {
960    match read_item(r, max_key_len, 4)? {
961        None => Ok(None),
962        Some((k, v, _)) => {
963            // A value that is not exactly four bytes means the level file was
964            // truncated or scribbled on. Taking a page number out of the
965            // wrong bytes would build interior pages pointing at arbitrary
966            // pages -- a structurally valid tree serving wrong answers -- so
967            // refuse instead. Not `Corrupt`: that carries a page number, and
968            // page 0 is the superblock; a scratch file has no page at all, so
969            // naming one would read as damage to the store itself.
970            if v.len() != 4 {
971                return Err(Error::Io(std::io::Error::new(
972                    std::io::ErrorKind::InvalidData,
973                    "separator spill record is not a 4-byte page number",
974                )));
975            }
976            Ok(Some((k, u32::from_le_bytes(v[..4].try_into().unwrap()))))
977        }
978    }
979}
980
981/// The physical result of packing one sorted key range.  `first_leaf` and
982/// `last_leaf` let a caller stitch the range into an existing leaf sequence
983/// without walking the packed tree or retaining one separator per page.
984#[derive(Debug, Clone)]
985pub(crate) struct PackedRange {
986    pub(crate) root: u32,
987    pub(crate) first_leaf: u32,
988    pub(crate) last_leaf: u32,
989    pub(crate) rows: u64,
990    pub(crate) min: Option<Vec<u8>>,
991    pub(crate) max: Option<Vec<u8>>,
992}
993
994/// Pack a sorted stream into a fresh tree, bottom up.
995///
996/// `last_next` is the already-known leaf immediately to the right of this
997/// range.  Whole-tree builds pass zero.  A graft passes the copied boundary
998/// leaf (or the standing right sibling), which means the final packed page is
999/// complete before its first and only write to disk.  No post-verification
1000/// pointer patch is needed.
1001///
1002/// `scratch_dir` is where each level's separators are spilled. Callers pass
1003/// the directory they already made for the sort (`Store::bulk_load`,
1004/// `recover`), rather than this inventing a second temp-file scheme with its
1005/// own lifetime and its own collision rules; the files created here are
1006/// removed as each level is consumed, and on every error path.
1007pub(crate) fn pack_range<I>(
1008    pool: &BufferPool,
1009    tree_id: u16,
1010    sorted: I,
1011    fill: f32,
1012    scratch_dir: &Path,
1013    last_next: u32,
1014) -> Result<PackedRange>
1015where I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>> {
1016    let mut sink = PageSink::direct(pool)?;
1017    pack_range_into(pool, tree_id, sorted, fill, scratch_dir, last_next, &mut sink)
1018}
1019
1020/// The same pack, with every page allocated and written through the ordinary
1021/// buffer pool. A page-WAL store uses this so a graft is an ordinary logged
1022/// transaction; see [`PageSink`].
1023pub(crate) fn pack_range_pooled<I>(
1024    pool: &BufferPool,
1025    tree_id: u16,
1026    sorted: I,
1027    fill: f32,
1028    scratch_dir: &Path,
1029    last_next: u32,
1030) -> Result<PackedRange>
1031where I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>> {
1032    let mut sink = PageSink::pooled(pool);
1033    let packed = pack_range_into(pool, tree_id, sorted, fill, scratch_dir, last_next, &mut sink);
1034    // Drop any guard a failed pack still holds before the error reaches a
1035    // caller that may roll back or commit.
1036    sink.finish()?;
1037    packed
1038}
1039
1040fn pack_range_into<I>(
1041    pool: &BufferPool,
1042    tree_id: u16,
1043    sorted: I,
1044    fill: f32,
1045    scratch_dir: &Path,
1046    last_next: u32,
1047    sink: &mut PageSink<'_>,
1048) -> Result<PackedRange>
1049where I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>> {
1050    let usable = ((PAGE_SIZE - HEADER_LEN) as f32 * fill) as usize;
1051    let capacity = PAGE_SIZE - HEADER_LEN;
1052
1053    std::fs::create_dir_all(scratch_dir)?;
1054    // One number per call, so two packs sharing a scratch directory cannot
1055    // name the same file. Same discipline as `Store::bulk_load`'s per-call
1056    // sequence number for the sort directory.
1057    static PACK_SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1058    let pack = PACK_SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1059    let mut seq = 0u64;
1060
1061    // Level 0: leaves. Their (first_key, page_no) pairs are spilled, not held.
1062    let mut level = Separators::create(scratch_dir, pack, &mut seq)?;
1063    let mut cur: Option<(u32, Vec<Vec<u8>>, Vec<u8>, usize)> = None; // (page, recs, first_key, used)
1064    let mut first_leaf: Option<u32> = None;
1065    let mut last_leaf: Option<u32> = None;
1066    let mut rows = 0u64;
1067    let mut range_min: Option<Vec<u8>> = None;
1068    let mut range_max: Option<Vec<u8>> = None;
1069
1070    let flush_leaf = |cur: &mut Option<(u32, Vec<Vec<u8>>, Vec<u8>, usize)>,
1071                          level: &mut Separators,
1072                          next_leaf: u32,
1073                          sink: &mut PageSink<'_>| -> Result<()> {
1074        if let Some((no, recs, first, _)) = cur.take() {
1075            let mut scratch = build_scratch(PageKind::Leaf, tree_id, no, &recs)?;
1076            {
1077                let mut page = PageMut::reopen(&mut scratch);
1078                page.set_next_leaf(next_leaf);
1079                page.finalise(0);
1080            }
1081            sink.push(no, &scratch)?;
1082            level.push(&first, no)?;
1083        }
1084        Ok(())
1085    };
1086
1087    let mut prev_key: Option<Vec<u8>> = None;
1088    for item in sorted {
1089        let (k, v, is_marker) = item?;
1090        // Refuse duplicates rather than pack them.
1091        //
1092        // Two entries with the same key can straddle a leaf boundary, and
1093        // the separator promoted into the parent then equals a key in the
1094        // leaf to its left -- so a descent routes past one of them and `get`
1095        // cannot reach it, while `range` still returns it. A divergence
1096        // between point lookup and scan is the worst kind of wrong answer:
1097        // both paths look correct in isolation. Refusing is honest, and a
1098        // caller who wants last-wins can deduplicate before calling.
1099        if prev_key.as_deref() == Some(k.as_slice()) {
1100            // Not `Corrupt { page_no: 0 }` -- page 0 is the superblock, so
1101            // that would name a real, meaningful page for a condition that
1102            // has nothing to do with one, and read as structural damage in
1103            // a log when it is really just an input the caller must
1104            // deduplicate.
1105            return Err(Error::DuplicateKey);
1106        }
1107        rows = rows.checked_add(1).ok_or(Error::TooLarge)?;
1108        if range_min.is_none() { range_min = Some(k.clone()); }
1109        range_max = Some(k.clone());
1110        prev_key = Some(k.clone());
1111        let rec = if is_marker {
1112            let m: [u8; 12] = v.as_slice().try_into().map_err(|_| {
1113                Error::Io(invalid_scratch("overflow marker scratch value is not exactly 12 bytes"))
1114            })?;
1115            crate::btree::enc_leaf_marker(&k, &m)
1116        } else {
1117            enc_leaf(&k, &v, pool.compact_cells())
1118        };
1119        let need = rec.len() + 4;
1120        // A record that cannot fit an empty page at all can never be packed,
1121        // bulk or otherwise -- `BTree::insert` refuses it before it ever
1122        // touches a page (`kernel/src/btree.rs`); refuse it here too, before
1123        // any page is allocated for it, rather than letting `insert_slot`
1124        // discover it deep inside `build_scratch`.
1125        if need > capacity {
1126            return Err(Error::TooLarge);
1127        }
1128        let fits = matches!(&cur, Some((_, _, _, used)) if used + need <= usable);
1129        if !fits {
1130            let no = sink.reserve()?;
1131            flush_leaf(&mut cur, &mut level, no, sink)?;
1132            if first_leaf.is_none() { first_leaf = Some(no); }
1133            last_leaf = Some(no);
1134            cur = Some((no, Vec::new(), k.clone(), 0));
1135        }
1136        if let Some((_, recs, _, used)) = cur.as_mut() { recs.push(rec); *used += need; }
1137    }
1138    flush_leaf(&mut cur, &mut level, last_next, sink)?;
1139    level.seal()?;
1140
1141    if level.count == 0 {
1142        let no = sink.reserve()?;
1143        let scratch = build_scratch(PageKind::Leaf, tree_id, no, &[])?;
1144        sink.push(no, &scratch)?;
1145        sink.finish()?;
1146        return Ok(PackedRange {
1147            root: no,
1148            first_leaf: no,
1149            last_leaf: no,
1150            rows: 0,
1151            min: None,
1152            max: None,
1153        });
1154    }
1155
1156    // Build interior levels the same way until one page remains. Same
1157    // convention as btree.rs: leftmost child in the header, slot array holds
1158    // strictly sorted (min_key, child) pairs.
1159    //
1160    // The old loop indexed `level[i]` and looked one entry ahead to decide
1161    // where a page ends. Streaming has no index, so the lookahead becomes an
1162    // explicit `pending`: the separator that did not fit the page just
1163    // finished is the first child of the next one. The page boundaries this
1164    // produces are identical to the indexed version's, which is what keeps
1165    // the resulting tree byte-for-byte the same.
1166    while level.count > 1 {
1167        let mut up = Separators::create(scratch_dir, pack, &mut seq)?;
1168        {
1169            let mut r = level.reader()?;
1170            // Counted against `level.count` once the file is consumed. See
1171            // `short_level`.
1172            let mut seen = 0u64;
1173            let mut pending = std::collections::VecDeque::new();
1174            if let Some(first) = read_sep(&mut r, level.max_key_len)? {
1175                seen += 1;
1176                pending.push_back(first);
1177            }
1178            while let Some((first, child0)) = pending.pop_front() {
1179                let no = sink.reserve()?;
1180                // Keep the separator beside its encoded form until this page
1181                // is sealed. If the very last child would strand itself on a
1182                // one-child page, we may have to move this page's final
1183                // separator over to become that page's child0.
1184                let mut recs: Vec<(Vec<u8>, u32, Vec<u8>)> = Vec::new();
1185                let mut used = 0usize;
1186                loop {
1187                    let next = if let Some(queued) = pending.pop_front() {
1188                        Some(queued)
1189                    } else {
1190                        let read = read_sep(&mut r, level.max_key_len)?;
1191                        if read.is_some() { seen += 1; }
1192                        read
1193                    };
1194                    let Some((k, child)) = next else { break };
1195                    // Counted where it is READ, not where it is used: the
1196                    // entry that does not fit is handed to the next page
1197                    // through `pending` and must not be counted twice.
1198                    let rec = enc_interior(&k, child);
1199                    let need = rec.len() + 4;
1200                    if need > capacity { return Err(Error::TooLarge); }
1201                    if used + need > usable && !recs.is_empty() {
1202                        // Never strand the final child in a one-child
1203                        // interior page. Let the penultimate page consume the
1204                        // last separator up to physical capacity; its slight
1205                        // overfill is bounded by the normal 10% reserve.
1206                        if seen == level.count && used + need <= capacity {
1207                            used += need;
1208                            recs.push((k, child, rec));
1209                            continue;
1210                        }
1211                        if seen == level.count {
1212                            // The reserve is not large enough for the final
1213                            // separator. Rebalance one separator from this
1214                            // page: it becomes child0 of the final page, and
1215                            // the separator just read becomes that page's one
1216                            // slot. Both pages therefore retain at least two
1217                            // children. With only one existing slot there is
1218                            // no valid two-page partition at this level.
1219                            if recs.len() == 1 { return Err(Error::TooLarge); }
1220                            let (moved_key, moved_child, _) = recs.pop().unwrap();
1221                            pending.push_back((moved_key, moved_child));
1222                            pending.push_back((k, child));
1223                            break;
1224                        }
1225                        pending.push_back((k, child));
1226                        break;
1227                    }
1228                    used += need;
1229                    recs.push((k, child, rec));
1230                }
1231                let encoded: Vec<Vec<u8>> = recs.into_iter().map(|(_, _, rec)| rec).collect();
1232                let mut scratch = build_scratch(PageKind::Interior, tree_id, no, &encoded)?;
1233                // `set_child0` and `set_next_leaf` are the same header field
1234                // (page.rs); set it in the scratch buffer and re-finalise so the
1235                // page reaches disk complete in its one batched write.
1236                {
1237                    let mut p = PageMut::reopen(&mut scratch);
1238                    p.set_child0(child0);
1239                    p.finalise(0);
1240                }
1241                sink.push(no, &scratch)?;
1242                up.push(&first, no)?;
1243            }
1244            // The outer loop only ends on `pending == None`, which only
1245            // happens at EOF, so the whole file has been consumed here.
1246            if seen != level.count { return Err(short_level(seen, level.count)); }
1247        }
1248        up.seal()?;
1249        // Assigning drops the level just consumed, which unlinks its file --
1250        // the cleanup happens level by level, so a deep tree never has more
1251        // than two separator files on disk at once.
1252        level = up;
1253    }
1254
1255    // Exactly one entry left: that page is the root. Checked both ways --
1256    // an empty file, and a file with anything after the root -- so the last
1257    // level gets the same treatment as every level below it.
1258    let mut r = level.reader()?;
1259    let root = match read_sep(&mut r, level.max_key_len)? {
1260        Some((_, no)) => no,
1261        None => return Err(short_level(0, level.count)),
1262    };
1263    if read_sep(&mut r, level.max_key_len)?.is_some() {
1264        return Err(Error::Io(std::io::Error::new(
1265            std::io::ErrorKind::InvalidData,
1266            "the final separator level held more than the one entry it counted",
1267        )));
1268    }
1269    // The root is about to become readable through the ordinary pool and,
1270    // after graft publication, reachable from committed metadata. Do not let
1271    // either happen while any candidate pages remain only in this buffer.
1272    sink.finish()?;
1273    Ok(PackedRange {
1274        root,
1275        first_leaf: first_leaf.expect("a non-empty level has a first leaf"),
1276        last_leaf: last_leaf.expect("a non-empty level has a last leaf"),
1277        rows,
1278        min: range_min,
1279        max: range_max,
1280    })
1281}
1282
1283/// Pack a complete sorted tree.  Kept as the stable whole-tree primitive used
1284/// by bulk load and recovery; range grafting extends it through `pack_range`
1285/// instead of duplicating its page builder.
1286pub fn pack_tree<I>(pool: &BufferPool, tree_id: u16, sorted: I, fill: f32, scratch_dir: &Path) -> Result<u32>
1287where I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>> {
1288    Ok(pack_range(pool, tree_id, sorted, fill, scratch_dir, 0)?.root)
1289}
1290
1291/// Pack a complete sorted tree through the ORDINARY buffer pool.
1292///
1293/// Same bottom-up pack as [`pack_tree`], but every page is an ordinary pooled
1294/// page, so a page-WAL database logs each one as a normal frame and publishes
1295/// the whole tree with the caller's commit (`PageSink::Pooled`). There is no
1296/// root swap, no skipped log and no second durability rule: Law 3 and Law 5
1297/// hold exactly as they do for a single `insert`, and a crash before the
1298/// caller's commit leaves nothing reachable.
1299///
1300/// This is what a per-index tree is built with: the tree is EMPTY, so there is
1301/// no standing content to graft beside and no boundary to plan -- the packed
1302/// root simply becomes the index descriptor's root in the same transaction.
1303/// Returns `(root, rows)`; an empty stream returns `(0, 0)`, the descriptor's
1304/// encoding of an empty tree.
1305pub fn pack_tree_pooled<I>(pool: &BufferPool, tree_id: u16, sorted: I, fill: f32, scratch_dir: &Path)
1306    -> Result<(u32, u64)>
1307where I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>> {
1308    let mut peek = sorted.peekable();
1309    if peek.peek().is_none() { return Ok((0, 0)); }
1310    let packed = pack_range_pooled(pool, tree_id, peek, fill, scratch_dir, 0)?;
1311    Ok((packed.root, packed.rows))
1312}