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}