1use crate::btree::{BTree, RangeIter};
11use crate::budget::MemoryBudget;
12use crate::io::{open_file_writer, Barrier, FileIo, IoMode};
13#[cfg(test)]
14use crate::io::open_file;
15use crate::meta::Meta;
16use crate::page::PAGE_SIZE;
17use crate::pool::BufferPool;
18use crate::wal::{RecKind, Wal};
19use crate::{Error, Result};
20use std::cell::Cell;
21use std::path::Path;
22use std::sync::Arc;
23
24#[cfg(feature = "test-support")]
29#[derive(Debug, Clone, Copy, PartialEq, Eq)]
30pub enum TestFaultKind { Read, Write, Commit }
31
32#[cfg(feature = "test-support")]
33#[derive(Debug, Default)]
34pub struct TestFaultInjector {
35 armed: std::sync::Mutex<Option<(TestFaultKind, usize)>>,
36}
37
38#[cfg(feature = "test-support")]
39impl TestFaultInjector {
40 pub fn arm(&self, kind: TestFaultKind, successful_matches: usize) {
42 *self.armed.lock().unwrap() = Some((kind, successful_matches));
43 }
44
45 fn check(&self, kind: TestFaultKind) -> Result<()> {
46 let mut armed = self.armed.lock().unwrap();
47 let Some((wanted, remaining)) = armed.as_mut() else { return Ok(()) };
48 if *wanted != kind { return Ok(()) }
49 if *remaining > 0 {
50 *remaining -= 1;
51 return Ok(());
52 }
53 *armed = None;
54 Err(std::io::Error::other(format!("injected {kind:?} failure at storage boundary")).into())
55 }
56}
57
58pub const PUT_EMPTY_BATCH_MAX_KEYS: usize = 64;
62
63#[derive(Debug, Clone, Copy, PartialEq, Eq)]
64pub enum SyncMode { Full, Normal, Off }
65
66#[derive(Debug, Clone, Copy)]
67pub struct Config { pub budget_bytes: usize, pub io: IoMode, pub sync: SyncMode }
68
69#[derive(Clone, Debug, PartialEq, Eq)]
76pub struct PreparedGraft {
77 pub base_generation: u64,
78 pub write_generation: u64,
79 pub root: u32,
80 pub rows: u64,
81 pub min: Vec<u8>,
82 pub max: Vec<u8>,
83 pub last_next: u32,
84 pub inserted_min: Vec<u8>,
85 pub inserted_max: Vec<u8>,
86 pub inserted_rows: u64,
87}
88
89impl PreparedGraft {
90 const MAGIC: &'static [u8; 8] = b"KGRFT01\0";
91
92 pub fn encode(&self) -> Result<Vec<u8>> {
93 let mut out = Vec::new();
94 out.extend_from_slice(Self::MAGIC);
95 out.extend_from_slice(&self.base_generation.to_le_bytes());
96 out.extend_from_slice(&self.write_generation.to_le_bytes());
97 out.extend_from_slice(&self.root.to_le_bytes());
98 out.extend_from_slice(&self.rows.to_le_bytes());
99 out.extend_from_slice(&self.last_next.to_le_bytes());
100 out.extend_from_slice(&self.inserted_rows.to_le_bytes());
101 for bytes in [&self.min, &self.max, &self.inserted_min, &self.inserted_max] {
102 let len = u32::try_from(bytes.len()).map_err(|_| Error::TooLarge)?;
103 out.extend_from_slice(&len.to_le_bytes());
104 out.extend_from_slice(bytes);
105 }
106 let crc = crc32c::crc32c(&out);
107 out.extend_from_slice(&crc.to_le_bytes());
108 Ok(out)
109 }
110
111 pub fn decode(bytes: &[u8]) -> Result<Self> {
112 const FIXED: usize = 8 + 8 + 8 + 4 + 8 + 4 + 8 + 4 * 4 + 4;
113 let invalid = || Error::Io(std::io::Error::new(
114 std::io::ErrorKind::InvalidData, "prepared graft manifest is invalid"));
115 if bytes.len() < FIXED || bytes.get(..8) != Some(Self::MAGIC) { return Err(invalid()); }
116 let crc_at = bytes.len() - 4;
117 let want = u32::from_le_bytes(bytes[crc_at..].try_into().unwrap());
118 if crc32c::crc32c(&bytes[..crc_at]) != want { return Err(invalid()); }
119 let mut pos = 8usize;
120 let mut u64_at = || {
121 let end = pos.checked_add(8).ok_or_else(invalid)?;
122 let value = u64::from_le_bytes(bytes.get(pos..end).ok_or_else(invalid)?.try_into().unwrap());
123 pos = end; Ok::<_, Error>(value)
124 };
125 let base_generation = u64_at()?;
126 let write_generation = u64_at()?;
127 drop(u64_at);
128 let root = u32::from_le_bytes(bytes.get(pos..pos + 4).ok_or_else(invalid)?.try_into().unwrap());
129 pos += 4;
130 let rows = u64::from_le_bytes(bytes.get(pos..pos + 8).ok_or_else(invalid)?.try_into().unwrap());
131 pos += 8;
132 let last_next = u32::from_le_bytes(bytes.get(pos..pos + 4).ok_or_else(invalid)?.try_into().unwrap());
133 pos += 4;
134 let inserted_rows = u64::from_le_bytes(bytes.get(pos..pos + 8).ok_or_else(invalid)?.try_into().unwrap());
135 pos += 8;
136 let mut fields = Vec::with_capacity(4);
137 for _ in 0..4 {
138 let raw = bytes.get(pos..pos + 4).ok_or_else(invalid)?;
139 pos += 4;
140 let len = u32::from_le_bytes(raw.try_into().unwrap()) as usize;
141 let end = pos.checked_add(len).ok_or_else(invalid)?;
142 if end > crc_at { return Err(invalid()); }
143 fields.push(bytes[pos..end].to_vec());
144 pos = end;
145 }
146 if pos != crc_at || fields.iter().any(Vec::is_empty) { return Err(invalid()); }
147 Ok(Self { base_generation, write_generation, root, rows,
148 min: fields.remove(0), max: fields.remove(0), inserted_min: fields.remove(0),
149 inserted_max: fields.remove(0), inserted_rows, last_next })
150 }
151
152 pub fn write_manifest(&self, path: &Path) -> Result<()> {
155 let parent = path.parent().unwrap_or_else(|| Path::new("."));
156 std::fs::create_dir_all(parent)?;
157 let tmp = path.with_extension("tmp");
158 {
159 use std::io::Write;
160 let mut file = std::fs::File::create(&tmp)?;
161 let bytes = self.encode()?;
162 file.write_all(&bytes)?;
163 crate::write_stats::add(
164 crate::write_stats::Phase::Manifest,
165 bytes.len() as u64,
166 );
167 file.sync_all()?;
168 }
169 std::fs::rename(&tmp, path)?;
170 std::fs::File::open(parent)?.sync_all()?;
171 Ok(())
172 }
173
174 pub fn read_manifest(path: &Path) -> Result<Self> {
175 let metadata = std::fs::metadata(path)?;
176 if metadata.len() > (1 << 20) {
177 return Err(Error::Io(std::io::Error::new(
178 std::io::ErrorKind::InvalidData, "prepared graft manifest is too large")));
179 }
180 Self::decode(&std::fs::read(path)?)
181 }
182}
183
184impl Default for Config {
190 fn default() -> Config {
191 Config { budget_bytes: 64 << 20, io: IoMode::Buffered, sync: SyncMode::Normal }
192 }
193}
194
195fn sync_barrier(sync: SyncMode) -> Barrier {
199 match sync {
200 SyncMode::Full => Barrier::Full,
201 SyncMode::Normal => Barrier::Data,
202 SyncMode::Off => Barrier::None,
203 }
204}
205
206pub struct Store {
207 pool: BufferPool,
208 wal: Option<Wal>,
211 root: u32,
212 dir: std::path::PathBuf,
215 generation: u64,
218 reader_slot: Option<crate::readers::ReaderSlot>,
221 tree_id: u16,
222 sync: SyncMode,
223 io_mode: IoMode,
224 #[cfg(feature = "test-support")]
225 test_faults: Arc<TestFaultInjector>,
226 poisoned: bool,
236 last_leaf: Cell<Option<u32>>,
249 fast_path_hits: Cell<u64>,
254 fast_path_attempts: Cell<u64>,
262 tag_hints: crate::btree::TagHints,
269 format_version: u16,
273 #[cfg(test)]
274 trace: Vec<&'static str>,
275 #[cfg(test)]
276 barriers: Vec<&'static str>,
277}
278
279#[cfg(test)]
280#[path = "store_reuse_probe.rs"]
281mod reuse_probe;
282
283impl Store {
284 fn build(dir: &Path, cfg: Config, fresh: bool) -> Result<Store> {
285 Self::build_limited(dir, cfg, fresh, None)
286 }
287 fn build_limited(dir: &Path, cfg: Config, fresh: bool, limits: Option<crate::limits::ResourceLimits>) -> Result<Store> {
288 if fresh && dir.join("data").exists() {
289 return Err(std::io::Error::new(std::io::ErrorKind::AlreadyExists, "create refuses existing data; use open").into());
290 }
291 std::fs::create_dir_all(dir)?;
292 if !fresh && crate::io::writer_owned_by_this_process(&dir.join("data")) {
299 if let crate::wal::Stop::Damaged { offset, why } =
300 crate::wal::Wal::inspect(&dir.join("wal"), cfg.io)?.stop
301 {
302 return Err(crate::Error::CorruptWal { offset, why });
303 }
304 }
305 let (file, io_mode) = open_file_writer(&dir.join("data"), cfg.io)?;
306 Self::build_on_limited(dir, cfg, fresh, file.into(), io_mode, limits)
307 }
308
309 #[cfg(test)]
317 fn build_on(dir: &Path, cfg: Config, fresh: bool, file: Arc<dyn FileIo>, io_mode: IoMode)
318 -> Result<Store> {
319 Self::build_on_limited(dir, cfg, fresh, file, io_mode, None)
320 }
321 fn build_on_limited(dir: &Path, cfg: Config, fresh: bool, file: Arc<dyn FileIo>, io_mode: IoMode,
322 requested: Option<crate::limits::ResourceLimits>) -> Result<Store> {
323 let dir_owned = dir.to_path_buf();
324 let frames = (cfg.budget_bytes / 3 * 2) / PAGE_SIZE;
326 let budget = Arc::new(MemoryBudget::new(cfg.budget_bytes));
327 let pool = BufferPool::new(file, budget, frames.max(16))?;
328 let limits = if fresh { requested } else { Meta::read_limits(&pool)? };
329 if let Some(l) = limits { pool.set_resource_limits(l)?; }
330 let mut wal = Wal::open_limited(&dir.join("wal"), cfg.io, limits.map(|l| l.wal_bytes))?;
331
332 let last_leaf = Cell::new(None);
336 let fast_path_hits = Cell::new(0);
337 let fast_path_attempts = Cell::new(0);
338 let tag_hints = crate::btree::TagHints::default();
339
340 let (root, generation, format_version) = if fresh {
341 let _meta_page = pool.allocate()?; drop(_meta_page);
343 let _slot_b = pool.allocate()?; drop(_slot_b);
345 Meta::init_slot_b(&pool)?;
346 pool.set_compact_cells(crate::meta::FORMAT_VERSION == 2);
347 let t = BTree::create(&pool, 1, &last_leaf, &fast_path_hits, &fast_path_attempts)?;
348 let r = t.root();
349 Meta { format_version: crate::meta::FORMAT_VERSION, roots: [r, 0, 0, 0, 0, 0, 0, 0], next_lsn: wal.next_lsn(),
350 generation: 0 }
351 .write(&pool)?;
352 pool.flush_all(sync_barrier(cfg.sync))?;
353 pool.sync_dir()?;
358 (r, 0, crate::meta::FORMAT_VERSION)
359 } else {
360 let meta = Meta::read_latest(&pool)?;
367 let format_version = meta.format_version & !crate::meta::LIMITED;
368 pool.set_compact_cells(format_version == 2);
369 wal.set_lsn_floor(meta.next_lsn);
370 pool.sync_dir()?;
374 (meta.roots[0], meta.generation, format_version)
375 };
376 pool.set_frozen_boundary();
381 let write_generation = generation.checked_add(1).ok_or(crate::Error::Corrupt {
382 page_no: if generation % 2 == 0 { crate::meta::META_PAGE } else { crate::meta::META_PAGE_B },
383 why: "published generation is exhausted",
384 })?;
385 pool.set_stamp_gen(write_generation);
386 if !fresh {
387 let cap = limits.map_or(28 + 24 * pool.page_count() as u64, |l| l.freelist_bytes());
390 if let Ok(f) = std::fs::File::open(dir.join("free")) {
391 use std::io::Read;
392 if f.metadata()?.len() <= cap {
393 let mut b = Vec::new();
394 f.take(cap + 1).read_to_end(&mut b)?;
395 if b.len() as u64 <= cap { pool.import_free(&b, generation); }
396 }
397 }
398 }
399 pool.set_reuse_limit(
400 generation.saturating_sub(1)
401 .min(crate::readers::oldest_live_reader(dir)));
402
403 let mut s = Store { pool, wal: Some(wal), root, generation, reader_slot: None, dir: dir_owned, tree_id: 1, sync: cfg.sync, io_mode,
404 #[cfg(feature = "test-support")] test_faults: Arc::new(TestFaultInjector::default()),
405 poisoned: false, last_leaf, fast_path_hits, fast_path_attempts,
406 tag_hints,
407 format_version,
408 #[cfg(test)] trace: Vec::new(),
409 #[cfg(test)] barriers: Vec::new() };
410 if !fresh {
411 if limits.is_some() {
412 if s.wal.as_ref().unwrap().committed_end()? != 0 {
415 return Err(Error::ResourceLimit("unexpected committed WAL in constrained store; preserve and inspect"));
416 }
417 s.wal_mut()?.rotate()?;
418 } else { s.recover_from_log()?; }
419 }
420 Ok(s)
421 }
422
423 pub fn open_snapshot(dir: &Path, cfg: Config) -> Result<Store> {
432 Self::open_snapshot_with_after_meta(dir, cfg, |_| Ok(()))
433 }
434
435 fn open_snapshot_with_after_meta<F>(dir: &Path, cfg: Config, after_meta: F) -> Result<Store>
440 where
441 F: FnOnce(u64) -> Result<()>,
442 {
443 let mut slot = crate::readers::ReaderSlot::reserve(dir)?;
447 let file = crate::io::open_file_readonly(&dir.join("data"))?;
448 let frames = (cfg.budget_bytes / 3 * 2) / PAGE_SIZE;
449 let budget = Arc::new(MemoryBudget::new(cfg.budget_bytes));
450 let pool = BufferPool::new(file.into(), budget, frames.max(16))?;
451 let meta = Meta::read_latest(&pool)?;
452 let format_version = meta.format_version & !crate::meta::LIMITED;
453 pool.set_compact_cells(format_version == 2);
454 if let Some(l) = Meta::read_limits(&pool)? { pool.set_resource_limits(l)?; }
455 after_meta(meta.generation)?;
456 slot.set_generation(meta.generation)?;
457 let s = Store {
458 pool, wal: None, root: meta.roots[0], generation: meta.generation,
459 reader_slot: Some(slot),
460 dir: dir.to_path_buf(),
461 tree_id: 1, sync: cfg.sync, io_mode: IoMode::Buffered,
462 #[cfg(feature = "test-support")]
463 test_faults: Arc::new(TestFaultInjector::default()),
464 poisoned: false,
465 last_leaf: Cell::new(None),
466 fast_path_hits: Cell::new(0), fast_path_attempts: Cell::new(0),
467 tag_hints: crate::btree::TagHints::default(),
468 format_version,
469 #[cfg(test)] trace: Vec::new(),
470 #[cfg(test)] barriers: Vec::new(),
471 };
472 Ok(s)
473 }
474
475 pub fn pool_ref(&self) -> &BufferPool { &self.pool }
477
478 pub fn create_limited(dir: &Path, cfg: Config, limits: crate::limits::ResourceLimits) -> Result<Store> {
482 let limits = limits.validate()?;
483 std::fs::create_dir(dir)?;
484 Self::build_limited(dir, cfg, true, Some(limits))
485 }
486 pub fn resource_limits(&self) -> Option<crate::limits::ResourceLimits> { self.pool.resource_limits() }
487 fn refuse_external_workspace(&self) -> Result<()> {
488 if self.resource_limits().is_some() {
489 return Err(Error::ResourceLimit("external sort/graft requires a separately budgeted destination; use batched put"));
490 }
491 Ok(())
492 }
493
494 pub fn create(dir: &Path, cfg: Config) -> Result<Store> { Self::build(dir, cfg, true) }
495
496 #[cfg(test)]
500 fn create_on(dir: &Path, cfg: Config, file: Arc<dyn FileIo>) -> Result<Store> {
501 std::fs::create_dir_all(dir)?;
502 Self::build_on(dir, cfg, true, file, IoMode::Buffered)
503 }
504 pub fn open(dir: &Path, cfg: Config) -> Result<Store> {
505 if dir.join("data").exists() { Self::build(dir, cfg, false) }
506 else { Self::build(dir, cfg, true) }
507 }
508
509 pub fn io_mode(&self) -> IoMode { self.io_mode }
510
511 #[cfg(feature = "test-support")]
512 #[doc(hidden)]
513 pub fn test_fault_injector(&self) -> Arc<TestFaultInjector> {
514 self.test_faults.clone()
515 }
516
517 fn recover_from_log(&mut self) -> Result<()> {
519 let committed_end = self.wal.as_ref().ok_or(crate::Error::ReadOnly)?.committed_end()?;
523 let mut off = 0u64;
524 while off < committed_end {
525 let (_, kind, payload, next) = self.wal.as_ref()
526 .ok_or(crate::Error::ReadOnly)?
527 .record_at(off, committed_end)?
528 .ok_or(crate::Error::CorruptWal {
529 offset: off, why: "committed recovery prefix ended early",
530 })?;
531 self.apply(kind, &payload, off)?;
532 off = next;
533 }
534 Ok(())
535 }
536
537 fn apply(&mut self, kind: RecKind, payload: &[u8], wal_offset: u64) -> Result<()> {
538 let corrupt = |why| crate::Error::CorruptWal { offset: wal_offset, why };
539 match kind {
540 RecKind::Put => {
541 let klen = payload.get(..2)
542 .map(|b| u16::from_le_bytes([b[0], b[1]]) as usize)
543 .ok_or_else(|| corrupt("put payload has no key length"))?;
544 let key_end = 2usize.checked_add(klen)
545 .ok_or_else(|| corrupt("put key boundary overflow"))?;
546 let key = payload.get(2..key_end)
547 .ok_or_else(|| corrupt("put key crosses its WAL payload"))?;
548 let val = payload.get(key_end..)
549 .ok_or_else(|| corrupt("put value boundary is invalid"))?;
550 if self.generation == 0 && crate::keys::is_field_aggregate_key(key) {
555 if Meta::is_salvaged(&self.pool)? { return Ok(()); }
556 }
557 let mut t = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
558 t.insert(key, val)?;
559 self.root = t.root();
560 }
561 RecKind::PutEmptyBatch => {
562 if payload.len() > crate::page::MAX_RECORD_LEN {
563 return Err(corrupt("put-empty-batch exceeds the WAL frame bound"));
564 }
565 let count = payload.get(..2)
566 .map(|b| u16::from_le_bytes([b[0], b[1]]) as usize)
567 .ok_or_else(|| corrupt("put-empty-batch payload has no key count"))?;
568 if count == 0 || count > PUT_EMPTY_BATCH_MAX_KEYS {
569 return Err(corrupt("put-empty-batch key count is outside its fixed bound"));
570 }
571
572 let mut at = 2usize;
576 for _ in 0..count {
577 let len_end = at.checked_add(2)
578 .ok_or_else(|| corrupt("put-empty-batch length boundary overflow"))?;
579 let len_bytes = payload.get(at..len_end)
580 .ok_or_else(|| corrupt("put-empty-batch key has no length"))?;
581 let key_len = u16::from_le_bytes([len_bytes[0], len_bytes[1]]) as usize;
582 let key_end = len_end.checked_add(key_len)
583 .ok_or_else(|| corrupt("put-empty-batch key boundary overflow"))?;
584 payload.get(len_end..key_end)
585 .ok_or_else(|| corrupt("put-empty-batch key crosses its WAL payload"))?;
586 if 4 + key_len + 12 > crate::page::MAX_RECORD_LEN {
587 return Err(corrupt("put-empty-batch key cannot fit a leaf record"));
588 }
589 at = key_end;
590 }
591 if at != payload.len() {
592 return Err(corrupt("put-empty-batch payload has trailing bytes"));
593 }
594
595 let mut t = BTree::open(&self.pool, self.tree_id, self.root,
596 &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
597 at = 2;
598 for _ in 0..count {
599 let key_len = u16::from_le_bytes([payload[at], payload[at + 1]]) as usize;
600 at += 2;
601 t.insert(&payload[at..at + key_len], &[])?;
602 at += key_len;
603 }
604 self.root = t.root();
605 }
606 RecKind::Delete => {
607 let klen = payload.get(..2)
608 .map(|b| u16::from_le_bytes([b[0], b[1]]) as usize)
609 .ok_or_else(|| corrupt("delete payload has no key length"))?;
610 let key_end = 2usize.checked_add(klen)
611 .ok_or_else(|| corrupt("delete key boundary overflow"))?;
612 if key_end != payload.len() {
613 return Err(corrupt("delete key length does not match its WAL payload"));
614 }
615 let mut t = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
616 t.delete(&payload[2..key_end])?;
617 self.root = t.root();
618 }
619 RecKind::DeletePrefix => {
620 if payload.is_empty() {
621 return Err(corrupt("delete-prefix WAL payload is empty"));
622 }
623 let mut t = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
624 t.delete_prefix(payload)?;
625 self.root = t.root();
626 }
627 RecKind::Commit => {
628 if !payload.is_empty() {
629 return Err(corrupt("commit WAL payload is not empty"));
630 }
631 }
632 RecKind::PageImage => {
633 return Err(corrupt("page-image WAL records have no recovery implementation"));
634 }
635 }
636 Ok(())
637 }
638
639 fn frame(k: &[u8], v: &[u8]) -> Vec<u8> {
645 let mut b = Vec::with_capacity(2 + k.len() + v.len());
646 b.extend_from_slice(&(k.len() as u16).to_le_bytes());
647 b.extend_from_slice(k);
648 b.extend_from_slice(v);
649 b
650 }
651
652 fn wal_mut(&mut self) -> Result<&mut Wal> {
654 self.wal.as_mut().ok_or(crate::Error::ReadOnly)
655 }
656
657 pub fn put(&mut self, k: &[u8], v: &[u8]) -> Result<()> {
658 #[cfg(feature = "test-support")]
659 self.test_faults.check(TestFaultKind::Write)?;
660 if self.poisoned { return Err(crate::Error::StorePoisoned); }
661 if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
662 let wal_payload_len = 2usize.checked_add(k.len())
663 .and_then(|n| n.checked_add(v.len()))
664 .ok_or(crate::Error::TooLarge)?;
665 if self.resource_limits().is_some_and(|l| v.len() > l.record_bytes as usize) {
666 return Err(Error::ResourceLimit("record exceeds configured maximum"));
667 }
668 if wal_payload_len as u64 > crate::wal::MAX_PAYLOAD_BYTES {
672 return Err(crate::Error::TooLarge);
673 }
674 if k.len() > crate::page::MAX_RECORD_LEN - 16 || v.len() > u32::MAX as usize {
675 return Err(crate::Error::TooLarge);
676 }
677 #[cfg(feature = "write-trace")]
678 let frame_started = crate::write_trace::active().then(std::time::Instant::now);
679 let payload = Self::frame(k, v);
680 #[cfg(feature = "write-trace")]
681 if let Some(started) = frame_started {
682 crate::write_trace::add(crate::write_trace::Field::FrameEncode, started.elapsed());
683 crate::write_trace::value_copy();
684 }
685 if 4 + k.len() + 12 > crate::page::MAX_RECORD_LEN || v.len() > u32::MAX as usize {
692 return Err(crate::Error::TooLarge);
693 }
694 if let Err(e) = self.wal_mut()?.append(RecKind::Put, &payload) {
695 self.poisoned = true;
699 return Err(e);
700 }
701 let mut t = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
702 #[cfg(feature = "write-trace")]
703 let btree_started = crate::write_trace::active().then(std::time::Instant::now);
704 let inserted = t.insert(k, v);
705 #[cfg(feature = "write-trace")]
706 if let Some(started) = btree_started {
707 crate::write_trace::add(crate::write_trace::Field::BtreeTotal, started.elapsed());
708 }
709 let root = t.root();
710 drop(t);
711 if let Err(e) = inserted {
712 self.poisoned = true;
713 return Err(e);
714 }
715 self.root = root;
716 Ok(())
717 }
718
719 #[cfg(feature = "write-trace")]
722 pub fn put_profiled(&mut self, k: &[u8], v: &[u8]) -> (Result<()>, crate::write_trace::PutTrace) {
723 let trace_started = crate::write_trace::begin();
724 let result = self.put(k, v);
725 let trace = crate::write_trace::finish(trace_started);
726 (result, trace)
727 }
728
729 pub fn put_empty_batch(&mut self, keys: &[Vec<u8>]) -> Result<()> {
734 if keys.is_empty() { return Ok(()); }
735 #[cfg(feature = "test-support")]
736 self.test_faults.check(TestFaultKind::Write)?;
737 if self.poisoned { return Err(crate::Error::StorePoisoned); }
738 if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
739 if keys.len() > PUT_EMPTY_BATCH_MAX_KEYS { return Err(crate::Error::TooLarge); }
740
741 let payload_len = keys.iter().try_fold(2usize, |total, key| {
742 if 4 + key.len() + 12 > crate::page::MAX_RECORD_LEN || key.len() > u16::MAX as usize {
743 return None;
744 }
745 total.checked_add(2 + key.len())
746 }).ok_or(crate::Error::TooLarge)?;
747 if payload_len > crate::page::MAX_RECORD_LEN { return Err(crate::Error::TooLarge); }
748 let mut payload = Vec::with_capacity(payload_len);
749 payload.extend_from_slice(&(keys.len() as u16).to_le_bytes());
750 for key in keys {
751 payload.extend_from_slice(&(key.len() as u16).to_le_bytes());
752 payload.extend_from_slice(key);
753 }
754 if let Err(e) = self.wal_mut()?.append(RecKind::PutEmptyBatch, &payload) {
755 self.poisoned = true;
756 return Err(e);
757 }
758 let mut t = BTree::open(&self.pool, self.tree_id, self.root,
759 &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
760 let inserted = keys.iter().try_for_each(|key| t.insert(key, &[]));
761 let root = t.root();
762 drop(t);
763 if let Err(e) = inserted {
764 self.poisoned = true;
765 return Err(e);
766 }
767 self.root = root;
768 Ok(())
769 }
770
771 pub fn delete(&mut self, k: &[u8]) -> Result<bool> {
772 #[cfg(feature = "test-support")]
773 self.test_faults.check(TestFaultKind::Write)?;
774 if self.poisoned { return Err(crate::Error::StorePoisoned); }
775 if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
776 if k.len() > crate::page::MAX_RECORD_LEN - 2 { return Ok(false); }
777 let mut payload = (k.len() as u16).to_le_bytes().to_vec();
778 payload.extend_from_slice(k);
779 if payload.len() > crate::page::MAX_RECORD_LEN { return Ok(false); }
788 if let Err(e) = self.wal_mut()?.append(RecKind::Delete, &payload) {
789 self.poisoned = true;
790 return Err(e);
791 }
792 let mut t = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
793 let deleted = t.delete(k);
794 let root = t.root();
795 drop(t);
796 let hit = match deleted {
797 Ok(hit) => hit,
798 Err(e) => {
799 self.poisoned = true;
800 return Err(e);
801 }
802 };
803 self.root = root;
804 Ok(hit)
805 }
806
807 pub fn delete_prefix(&mut self, prefix: &[u8]) -> Result<u64> {
813 #[cfg(feature = "test-support")]
814 self.test_faults.check(TestFaultKind::Write)?;
815 if self.poisoned { return Err(crate::Error::StorePoisoned); }
816 if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
817 if prefix.is_empty() || prefix.len() > crate::page::MAX_RECORD_LEN { return Err(crate::Error::TooLarge); }
818 if let Err(e) = self.wal_mut()?.append(RecKind::DeletePrefix, prefix) {
819 self.poisoned = true;
820 return Err(e);
821 }
822 let mut t = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
823 let deleted = t.delete_prefix(prefix);
824 let root = t.root();
825 drop(t);
826 let n = match deleted {
827 Ok(n) => n,
828 Err(e) => {
829 self.poisoned = true;
830 return Err(e);
831 }
832 };
833 self.root = root;
834 Ok(n)
835 }
836
837 pub fn get(&self, k: &[u8]) -> Result<Option<Vec<u8>>> {
838 if self.poisoned && self.resource_limits().is_some() {
839 return Err(crate::Error::StorePoisoned);
840 }
841 #[cfg(feature = "test-support")]
842 self.test_faults.check(TestFaultKind::Read)?;
843 BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints).get(k)
844 }
845
846 pub fn scan(&self, from: &[u8]) -> Result<RangeIter<'_>> {
847 if self.poisoned && self.resource_limits().is_some() {
848 return Err(crate::Error::StorePoisoned);
849 }
850 BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints).range(from)
851 }
852
853 pub fn scan_reverse(&self, to: &[u8]) -> Result<crate::btree::ReverseRangeIter<'_>> {
855 if self.poisoned && self.resource_limits().is_some() {
856 return Err(crate::Error::StorePoisoned);
857 }
858 BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf,
859 &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints).range_reverse(to)
860 }
861
862 pub fn commit(&mut self) -> Result<()> {
876 if self.resource_limits().is_some() { return self.checkpoint(); }
877 #[cfg(feature = "test-support")]
878 self.test_faults.check(TestFaultKind::Commit)?;
879 if self.poisoned { return Err(crate::Error::StorePoisoned); }
880 let wal = self.wal.as_mut().ok_or(crate::Error::ReadOnly)?;
881 wal.append(RecKind::Commit, &[])?;
882 wal.flush()?;
886 match self.sync {
887 SyncMode::Full => {
888 #[cfg(test)] self.barriers.push("sync_full");
889 wal.sync_full()?
890 }
891 SyncMode::Normal => {
892 #[cfg(test)] self.barriers.push("sync_data");
893 wal.sync_data()?
894 }
895 SyncMode::Off => {}
896 }
897 let readers = crate::readers::live_generations(&self.dir);
903 self.pool.refresh_reuse(self.generation, readers.as_deref());
904 Ok(())
905 }
906
907 pub fn commit_with_checkpoint(&mut self, wal_bytes: u64, page_bytes: u64) -> Result<bool> {
917 if self.resource_limits().is_some() { self.commit()?; return Ok(true); }
918 if self.poisoned { return Err(crate::Error::StorePoisoned); }
919 let end = self.wal.as_ref().ok_or(crate::Error::ReadOnly)?.end_offset();
920 let publish = (wal_bytes > 0 && end >= wal_bytes)
921 || (page_bytes > 0 && self.pool.epoch_allocated_bytes() >= page_bytes);
922 if !publish { self.commit()?; return Ok(false); }
923 #[cfg(feature = "test-support")]
924 self.test_faults.check(TestFaultKind::Commit)?;
925 let wal = self.wal_mut()?;
926 if let Err(e) = wal.append(RecKind::Commit, &[]).and_then(|_| wal.flush()) {
927 self.poisoned = true;
928 return Err(e);
929 }
930 if let Err(e) = self.checkpoint() {
931 self.poisoned = true;
932 return Err(e);
933 }
934 Ok(true)
935 }
936
937 #[cfg(test)]
938 fn barriers(&self) -> Vec<&'static str> { self.barriers.clone() }
939
940 pub fn dir(&self) -> &Path { &self.dir }
943 pub fn published_root(&self) -> u32 { self.root }
946 pub fn main_tree_id(&self) -> u16 { self.tree_id }
948
949 pub fn sync_full_primitive(&self) -> &'static str {
950 self.wal.as_ref().map_or("none (snapshot reader)", |w| w.sync_full_primitive())
951 }
952
953 pub fn pool_stats(&self) -> crate::pool::PoolStats { self.pool.stats() }
960 pub fn tag_hints(&self) -> &crate::btree::TagHints { &self.tag_hints }
963 pub fn io_stats(&self) -> Option<&crate::io::IoStats> { self.pool.io_stats() }
965 pub fn sweep_steps(&self) -> u64 { self.pool.sweep_steps() }
966
967 fn barrier(&self) -> Barrier { sync_barrier(self.sync) }
974
975 pub fn set_sync(&mut self, s: SyncMode) { self.sync = s; }
978
979 pub fn sync_mode(&self) -> SyncMode { self.sync }
984 pub fn generation(&self) -> u64 { self.generation }
986
987 pub fn checkpoint(&mut self) -> Result<()> {
988 let result = self.checkpoint_inner();
989 if result.is_err() && self.wal.is_some() { self.poisoned = true; }
990 result
991 }
992 fn checkpoint_inner(&mut self) -> Result<()> {
993 if self.poisoned { return Err(crate::Error::StorePoisoned); }
994 if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
995 let gen = self.generation.checked_add(1).ok_or(crate::Error::Corrupt {
1000 page_no: if self.generation % 2 == 0 { crate::meta::META_PAGE } else { crate::meta::META_PAGE_B },
1001 why: "published generation is exhausted",
1002 })?;
1003 let next_write_generation = gen.checked_add(1).ok_or(crate::Error::Corrupt {
1004 page_no: if self.generation % 2 == 0 { crate::meta::META_PAGE } else { crate::meta::META_PAGE_B },
1005 why: "no generation remains for the next write epoch",
1006 })?;
1007 #[cfg(test)] self.trace.clear();
1012 if let Err(e) = self.pool.flush_all(self.barrier()) {
1019 self.poisoned = true;
1020 return Err(e);
1021 }
1022 #[cfg(test)] { self.trace.push("flush_pages"); self.trace.push("sync_file"); }
1023 Meta {
1029 format_version: self.format_version,
1030 roots: [self.root, 0, 0, 0, 0, 0, 0, 0],
1031 next_lsn: self.wal.as_ref().ok_or(crate::Error::ReadOnly)?.next_lsn(),
1032 generation: gen,
1033 }.write_slot(&self.pool)?;
1034 if let Err(e) = self.pool.flush_all(self.barrier()) {
1035 self.poisoned = true;
1036 return Err(e);
1037 }
1038 self.generation = gen;
1039 self.pool.set_stamp_gen(next_write_generation);
1040 let readers = crate::readers::live_generations(&self.dir);
1044 self.pool.refresh_reuse(gen, readers.as_deref());
1045 let free_bytes = self.pool.export_free(gen);
1046 let _ = crate::verify::persist_checkpoint_freelist(
1047 &self.dir,
1048 &free_bytes,
1049 self.pool.file_ref(),
1050 );
1051 #[cfg(test)] self.trace.push("flip_meta");
1057 self.pool.set_frozen_boundary();
1060 self.wal_mut()?.rotate_published()?;
1063 #[cfg(test)] self.trace.push("rotate_wal");
1064 Ok(())
1065 }
1066
1067 #[cfg(test)]
1068 fn checkpoint_trace(&self) -> Vec<&'static str> { self.trace.clone() }
1069
1070 #[cfg(test)]
1077 fn next_lsn(&self) -> u64 { self.wal.as_ref().unwrap().next_lsn() }
1078
1079 pub fn bulk_load<I>(&mut self, items: I) -> Result<()>
1105 where I: Iterator<Item = (Vec<u8>, Vec<u8>)> {
1106 self.bulk_load_with_before_publish(items, |_, _| Ok(()))
1107 }
1108
1109 fn bulk_load_with_before_publish<I, F>(&mut self, items: I, before_publish: F) -> Result<()>
1112 where
1113 I: Iterator<Item = (Vec<u8>, Vec<u8>)>,
1114 F: FnOnce(&Path, u32) -> Result<()>,
1115 {
1116 self.refuse_external_workspace()?;
1121 if self.poisoned { return Err(crate::Error::StorePoisoned); }
1122 if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
1123 static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1137 let seq = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1138 let tmp = std::env::temp_dir()
1139 .join(format!("kernel-sort-{}-{}", std::process::id(), seq));
1140 let mut s = crate::bulk::ExternalSort::new(&tmp, 64 << 20)?;
1141 let mut expected_rows = 0u64;
1142 for (k, v) in items {
1143 expected_rows = expected_rows.checked_add(1).ok_or(crate::Error::TooLarge)?;
1144 if 4 + k.len() + v.len() > crate::page::MAX_RECORD_LEN {
1149 let (head, crc) = crate::btree::write_overflow(&self.pool, &v)?;
1150 let m = crate::btree::enc_marker(v.len() as u32, head, crc);
1151 s.push_flagged(k, m.to_vec(), true)?;
1152 } else {
1153 s.push(k, v)?;
1154 }
1155 }
1156 let mut runs = s.finish()?;
1157 let root = crate::bulk::pack_tree(&self.pool, self.tree_id, runs.iter()?, 0.9, &tmp)?;
1162 if let Err(e) = self.pool.flush_all(Barrier::None) {
1165 self.poisoned = true;
1166 return Err(e);
1167 }
1168 let data = self.dir.join("data");
1169 before_publish(&data, root)?;
1170 let verified = crate::verify::verify_file(
1171 &data,
1172 self.io_mode,
1173 root,
1174 self.tree_id,
1175 expected_rows,
1176 )?;
1177 debug_assert_eq!(verified.rows, expected_rows);
1178 debug_assert!(verified.pages > 0);
1179 self.root = root;
1180 self.last_leaf.set(None);
1188 self.tag_hints.clear();
1189 self.checkpoint()
1190 }
1191
1192 pub fn graft_range<I>(&mut self, items: I) -> Result<()>
1203 where
1204 I: Iterator<Item = (Vec<u8>, Vec<u8>)>,
1205 {
1206 self.graft_range_with_before_publish(items, |_, _| Ok(()))
1207 }
1208
1209 fn graft_range_with_before_publish<I, F>(
1212 &mut self,
1213 items: I,
1214 before_publish: F,
1215 ) -> Result<()>
1216 where
1217 I: Iterator<Item = (Vec<u8>, Vec<u8>)>,
1218 F: FnOnce(&Path, &crate::bulk::PackedRange) -> Result<()>,
1219 {
1220 self.refuse_external_workspace()?;
1221 if self.poisoned { return Err(crate::Error::StorePoisoned); }
1222 if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
1223
1224 static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1225 let seq = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1226 let tmp = std::env::temp_dir()
1227 .join(format!("kernel-graft-{}-{}", std::process::id(), seq));
1228 let mut sort = crate::bulk::ExternalSort::new(&tmp, 64 << 20)?;
1229 let mut expected_rows = 0u64;
1230 let mut min: Option<Vec<u8>> = None;
1231 let mut max: Option<Vec<u8>> = None;
1232
1233 for (key, value) in items {
1234 expected_rows = expected_rows.checked_add(1).ok_or(crate::Error::TooLarge)?;
1235 if min.as_ref().is_none_or(|current| key.as_slice() < current.as_slice()) {
1236 min = Some(key.clone());
1237 }
1238 if max.as_ref().is_none_or(|current| key.as_slice() > current.as_slice()) {
1239 max = Some(key.clone());
1240 }
1241 if 4 + key.len() + value.len() > crate::page::MAX_RECORD_LEN {
1242 if value.len() > u32::MAX as usize || 4 + key.len() + 12 > crate::page::MAX_RECORD_LEN {
1243 return Err(crate::Error::TooLarge);
1244 }
1245 let (head, crc) = crate::btree::write_overflow(&self.pool, &value)?;
1246 let marker = crate::btree::enc_marker(value.len() as u32, head, crc);
1247 sort.push_flagged(key, marker.to_vec(), true)?;
1248 } else {
1249 sort.push(key, value)?;
1250 }
1251 }
1252 let (Some(min), Some(max)) = (min, max) else {
1253 return Ok(());
1256 };
1257
1258 let mut runs = sort.finish()?;
1259 self.graft_sorted_range_with_before_publish(
1260 runs.iter()?, expected_rows, min, max, &tmp, true, before_publish)
1261 }
1262
1263 pub fn graft_sorted_range<I>(
1272 &mut self,
1273 sorted: I,
1274 expected_rows: u64,
1275 min: Vec<u8>,
1276 max: Vec<u8>,
1277 scratch_dir: &Path,
1278 ) -> Result<()>
1279 where
1280 I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>,
1281 {
1282 self.graft_sorted_range_with_before_publish(
1283 sorted, expected_rows, min, max, scratch_dir, true, |_, _| Ok(()))
1284 }
1285
1286 pub fn graft_sorted_range_deferred<I>(
1292 &mut self,
1293 sorted: I,
1294 expected_rows: u64,
1295 min: Vec<u8>,
1296 max: Vec<u8>,
1297 scratch_dir: &Path,
1298 ) -> Result<()>
1299 where
1300 I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>,
1301 {
1302 self.graft_sorted_range_with_before_publish(
1303 sorted, expected_rows, min, max, scratch_dir, false, |_, _| Ok(()))
1304 }
1305
1306 pub fn prepare_graft_candidate<I>(
1311 &mut self,
1312 sorted: I,
1313 expected_rows: u64,
1314 min: Vec<u8>,
1315 max: Vec<u8>,
1316 scratch_dir: &Path,
1317 ) -> Result<PreparedGraft>
1318 where
1319 I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>,
1320 {
1321 self.refuse_external_workspace()?;
1322 if self.poisoned { return Err(Error::StorePoisoned); }
1323 if self.wal.is_none() { return Err(Error::ReadOnly); }
1324 if expected_rows == 0 || min > max { return Err(Error::TooLarge); }
1325 if let Some(row) = self.scan(&min)?.next() {
1326 let (key, _) = row?;
1327 if key <= max { return Err(Error::RangeNotEmpty); }
1328 }
1329 let tree = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf,
1330 &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
1331 let boundary = tree.plan_graft(&min)?;
1332 drop(tree);
1333 let last_next = boundary.right_page.unwrap_or(boundary.old_next);
1334 let packed = crate::bulk::pack_range(
1335 &self.pool, self.tree_id, sorted, 0.9, scratch_dir, last_next)?;
1336 if packed.rows != expected_rows || packed.min.as_deref() != Some(min.as_slice())
1337 || packed.max.as_deref() != Some(max.as_slice()) {
1338 return Err(Error::Corrupt { page_no: packed.root,
1339 why: "prepared range disagrees with its sorted-stream manifest" });
1340 }
1341 let tree = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf,
1342 &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
1343 let candidate = tree.build_graft_candidate(&boundary, &packed)?;
1344 drop(tree);
1345 if let Err(error) = self.pool.flush_all(Barrier::None) {
1346 self.poisoned = true;
1347 return Err(error);
1348 }
1349 let write_generation = self.generation.checked_add(1).ok_or(Error::Corrupt {
1350 page_no: 0, why: "prepared graft generation is exhausted" })?;
1351 crate::verify::verify_range_file_generation(
1352 &self.dir.join("data"), self.io_mode, candidate.root, self.tree_id,
1353 candidate.rows, &candidate.min, &candidate.max, candidate.last_next,
1354 Some(write_generation))?;
1355 Ok(PreparedGraft {
1356 base_generation: self.generation,
1357 write_generation,
1358 root: candidate.root,
1359 rows: candidate.rows,
1360 min: candidate.min,
1361 max: candidate.max,
1362 last_next: candidate.last_next,
1363 inserted_min: min,
1364 inserted_max: max,
1365 inserted_rows: expected_rows,
1366 })
1367 }
1368
1369 pub fn publish_existing_candidate(&mut self, prepared: &PreparedGraft) -> Result<()> {
1374 self.refuse_external_workspace()?;
1375 if self.poisoned { return Err(Error::StorePoisoned); }
1376 if self.wal.is_none() { return Err(Error::ReadOnly); }
1377 if prepared.write_generation != prepared.base_generation.checked_add(1)
1378 .ok_or(Error::TooLarge)? || prepared.inserted_min > prepared.inserted_max
1379 || prepared.inserted_rows == 0 {
1380 return Err(Error::Corrupt { page_no: prepared.root,
1381 why: "prepared graft has an invalid generation or interval" });
1382 }
1383 let data = self.dir.join("data");
1384 crate::verify::verify_range_file_generation(
1385 &data, self.io_mode, prepared.root, self.tree_id, prepared.rows,
1386 &prepared.min, &prepared.max, prepared.last_next,
1387 Some(prepared.write_generation))?;
1388
1389 if self.generation == prepared.write_generation {
1390 let mut rows = 0u64;
1391 self.scan(&prepared.inserted_min)?.for_each_ref(|key, _| {
1392 if key > prepared.inserted_max.as_slice() { return false; }
1393 rows = rows.saturating_add(1);
1394 true
1395 })?;
1396 if rows != prepared.inserted_rows {
1397 return Err(Error::Corrupt { page_no: prepared.root,
1398 why: "published graft interval disagrees with its manifest" });
1399 }
1400 return Ok(());
1401 }
1402 if self.generation != prepared.base_generation {
1403 return Err(Error::Corrupt { page_no: prepared.root,
1404 why: "prepared graft belongs to a stale base generation" });
1405 }
1406 if let Some(row) = self.scan(&prepared.inserted_min)?.next() {
1407 let (key, _) = row?;
1408 if key <= prepared.inserted_max { return Err(Error::RangeNotEmpty); }
1409 }
1410 let mut tree = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf,
1411 &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
1412 let boundary = tree.plan_existing_graft(&prepared.inserted_min)?;
1413 let retired = tree.install_graft(&boundary, prepared.root, &prepared.min)?;
1414 self.root = tree.root();
1415 drop(tree);
1416 for page in retired { self.pool.free_page(page)?; }
1417 self.last_leaf.set(None);
1418 self.tag_hints.clear();
1419 self.checkpoint()
1420 }
1421
1422 fn graft_sorted_range_with_before_publish<I, F>(
1423 &mut self,
1424 sorted: I,
1425 expected_rows: u64,
1426 min: Vec<u8>,
1427 max: Vec<u8>,
1428 scratch_dir: &Path,
1429 publish: bool,
1430 before_publish: F,
1431 ) -> Result<()>
1432 where
1433 I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>,
1434 F: FnOnce(&Path, &crate::bulk::PackedRange) -> Result<()>,
1435 {
1436 let trace = std::env::var_os("SEKEJAP_LOAD_BREAKDOWN").is_some();
1437 let total_started = trace.then(std::time::Instant::now);
1438 self.refuse_external_workspace()?;
1439 if self.poisoned {
1440 return Err(crate::Error::StorePoisoned);
1441 }
1442 if self.wal.is_none() {
1443 return Err(crate::Error::ReadOnly);
1444 }
1445 if expected_rows == 0 {
1446 return Ok(());
1447 }
1448 if min > max {
1449 return Err(crate::Error::TooLarge);
1450 }
1451
1452 let preflight_started = trace.then(std::time::Instant::now);
1453 if let Some(row) = self.scan(&min)?.next() {
1454 let (key, _) = row?;
1455 if key <= max { return Err(crate::Error::RangeNotEmpty); }
1456 }
1457
1458 let tree = BTree::open(
1459 &self.pool,
1460 self.tree_id,
1461 self.root,
1462 &self.last_leaf,
1463 &self.fast_path_hits,
1464 &self.fast_path_attempts,
1465 ).with_tags(&self.tag_hints);
1466 let boundary = tree.plan_graft(&min)?;
1467 drop(tree);
1468 let last_next = boundary.right_page.unwrap_or(boundary.old_next);
1469 let preflight =
1470 preflight_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1471
1472 let pack_started = trace.then(std::time::Instant::now);
1473 let packed = crate::bulk::pack_range(
1474 &self.pool,
1475 self.tree_id,
1476 sorted,
1477 0.9,
1478 scratch_dir,
1479 last_next,
1480 )?;
1481 let pack = pack_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1482 if packed.rows != expected_rows
1483 || packed.min.as_deref() != Some(min.as_slice())
1484 || packed.max.as_deref() != Some(max.as_slice())
1485 {
1486 return Err(crate::Error::Corrupt {
1487 page_no: packed.root,
1488 why: "packed range disagrees with its sorted-stream manifest",
1489 });
1490 }
1491 if packed.last_leaf >= self.pool.page_count() {
1492 return Err(crate::Error::Corrupt {
1493 page_no: packed.last_leaf,
1494 why: "packed range returned a last leaf outside the data file",
1495 });
1496 }
1497
1498 let tree = BTree::open(
1499 &self.pool,
1500 self.tree_id,
1501 self.root,
1502 &self.last_leaf,
1503 &self.fast_path_hits,
1504 &self.fast_path_attempts,
1505 ).with_tags(&self.tag_hints);
1506 let candidate = tree.build_graft_candidate(&boundary, &packed)?;
1507 drop(tree);
1508
1509 let flush_started = trace.then(std::time::Instant::now);
1512 if let Err(error) = self.pool.flush_all(Barrier::None) {
1513 self.poisoned = true;
1514 return Err(error);
1515 }
1516 let flush = flush_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1517 let data = self.dir.join("data");
1518 before_publish(&data, &packed)?;
1519 let verify_started = trace.then(std::time::Instant::now);
1520 let verified = crate::verify::verify_range_file(
1521 &data,
1522 self.io_mode,
1523 candidate.root,
1524 self.tree_id,
1525 candidate.rows,
1526 &candidate.min,
1527 &candidate.max,
1528 candidate.last_next,
1529 )?;
1530 let verify = verify_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1531 debug_assert_eq!(verified.rows, candidate.rows);
1532 debug_assert!(verified.pages > 0);
1533
1534 let install_started = trace.then(std::time::Instant::now);
1539 let mut tree = BTree::open(
1540 &self.pool,
1541 self.tree_id,
1542 self.root,
1543 &self.last_leaf,
1544 &self.fast_path_hits,
1545 &self.fast_path_attempts,
1546 ).with_tags(&self.tag_hints);
1547 let retired = tree.install_graft(&boundary, candidate.root, &candidate.min)?;
1548 self.root = tree.root();
1549 drop(tree);
1550 for page in retired { self.pool.free_page(page)?; }
1551 self.last_leaf.set(None);
1552 self.tag_hints.clear();
1553 let install =
1554 install_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1555 let publish_started = trace.then(std::time::Instant::now);
1556 let result = if publish { self.checkpoint() } else { Ok(()) };
1557 let publication =
1558 publish_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1559 if trace {
1560 eprintln!("graft detail: rows={} pages={} bytes={} preflight={:.6}s pack={:.6}s flush={:.6}s verify={:.6}s graft={:.6}s publish={:.6}s total={:.6}s deferred={}",
1561 expected_rows, verified.pages, verified.pages as u64 * crate::page::PAGE_SIZE as u64,
1562 preflight.as_secs_f64(), pack.as_secs_f64(), flush.as_secs_f64(),
1563 verify.as_secs_f64(), install.as_secs_f64(), publication.as_secs_f64(),
1564 total_started.unwrap().elapsed().as_secs_f64(), !publish);
1565 }
1566 result
1567 }
1568}
1569
1570#[cfg(test)]
1578mod tests {
1579 use super::*;
1580
1581 fn cfg() -> Config { Config { budget_bytes: 32 << 20, io: IoMode::Buffered, sync: SyncMode::Full } }
1582
1583 #[test]
1589 fn a_store_checkpoint_preserves_the_file_own_format_version() {
1590 let d = tempfile::tempdir().unwrap();
1591 let other = if crate::meta::FORMAT_VERSION == 2 { 1 } else { 2 };
1592 {
1593 let s = Store::create(d.path(), cfg()).unwrap();
1594 let meta = Meta::read_latest(&s.pool).unwrap();
1595 Meta {
1596 format_version: other,
1597 roots: meta.roots,
1598 next_lsn: meta.next_lsn,
1599 generation: meta.generation,
1600 }
1601 .write(&s.pool)
1602 .unwrap();
1603 s.pool.flush_all(crate::io::Barrier::Data).unwrap();
1604 }
1605 let mut s = Store::open(d.path(), cfg()).unwrap_or_else(|e| {
1606 panic!("supported superblock version {other} must open in this build: {e}")
1607 });
1608 s.put(b"k", b"v").unwrap();
1609 s.commit().unwrap();
1610 s.checkpoint().unwrap();
1611 drop(s);
1612 let s = Store::open(d.path(), cfg()).unwrap();
1613 let got = Meta::read_latest(&s.pool).unwrap().format_version & !crate::meta::LIMITED;
1614 assert_eq!(
1615 got, other,
1616 "checkpoint must write the file's format version, not the build's FORMAT_VERSION"
1617 );
1618 assert_eq!(s.get(b"k").unwrap().as_deref(), Some(&b"v"[..]));
1619 }
1620
1621 #[test]
1622 fn byte_policy_publishes_exact_rows_without_a_redundant_wal_barrier() {
1623 let d = tempfile::tempdir().unwrap();
1624 let mut s = Store::create(d.path(), cfg()).unwrap();
1625 s.put(b"old", b"durable").unwrap(); s.commit().unwrap(); s.checkpoint().unwrap();
1626 let reader = Store::open_snapshot(d.path(), cfg()).unwrap();
1627 s.barriers.clear();
1628 s.put(b"new", &[7; 9000]).unwrap();
1629 let before = s.pool_stats().sync_full_calls;
1630 assert!(s.commit_with_checkpoint(u64::MAX, 4096).unwrap());
1631 assert_eq!(s.pool_stats().sync_full_calls - before, 2);
1632 assert!(s.barriers.is_empty(), "publication already supplies the durability barriers");
1633 assert_eq!(s.wal.as_ref().unwrap().end_offset(), 0);
1634 assert_eq!(reader.get(b"new").unwrap(), None);
1635 drop(s);
1636 let s = Store::open(d.path(), cfg()).unwrap();
1637 assert_eq!(s.get(b"old").unwrap().as_deref(), Some(&b"durable"[..]));
1638 assert_eq!(s.get(b"new").unwrap(), Some(vec![7; 9000]));
1639 }
1640
1641 #[test]
1642 fn byte_policy_keeps_wal_and_old_snapshot_after_data_barrier_failure() {
1643 let d = tempfile::tempdir().unwrap();
1644 let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1645 let fio = Arc::new(FailingBarrier { inner: real, fail: false.into() });
1646 let mut s = Store::create_on(d.path(), cfg(), fio.clone()).unwrap();
1647 s.put(b"old", b"committed").unwrap(); s.commit().unwrap(); s.checkpoint().unwrap();
1648 let reader = Store::open_snapshot(d.path(), cfg()).unwrap();
1649 s.put(b"new", &[8; 9000]).unwrap();
1650 fio.fail.store(true, std::sync::atomic::Ordering::Relaxed);
1651 assert!(s.commit_with_checkpoint(1,1).is_err());
1652 assert!(s.wal.as_ref().unwrap().end_offset() > 0);
1653 assert!(matches!(s.commit_with_checkpoint(1,1), Err(crate::Error::StorePoisoned)));
1654 assert_eq!(reader.get(b"old").unwrap().as_deref(), Some(&b"committed"[..]));
1655 assert_eq!(reader.get(b"new").unwrap(), None);
1656 fio.fail.store(false, std::sync::atomic::Ordering::Relaxed);
1657 drop(s); drop(fio);
1658 let s = Store::open(d.path(), cfg()).unwrap();
1659 assert_eq!(s.get(b"old").unwrap().as_deref(), Some(&b"committed"[..]));
1660 assert_eq!(s.get(b"new").unwrap(), Some(vec![8; 9000]));
1661 }
1662
1663 #[test]
1664 fn constrained_commit_failure_preserves_published_state_and_snapshot() {
1665 for barrier in [false, true] {
1666 let d = tempfile::tempdir().unwrap();
1667 let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1668 let fio = Arc::new(FailingBarrier { inner: real, fail: false.into() });
1669 let limits = crate::limits::ResourceLimits { data_bytes: 1 << 20, wal_bytes: 64 << 10,
1670 tracked_pages: 256, readers: 2, record_bytes: 16000, recovery_bytes: 65536 };
1671 let mut s = Store::build_on_limited(d.path(), cfg(), true, fio.clone(), IoMode::Buffered, Some(limits)).unwrap();
1672 s.put(b"old", b"durable").unwrap(); s.commit().unwrap();
1673 let reader = Store::open_snapshot(d.path(), cfg()).unwrap();
1674 s.put(b"new", &[8;9000]).unwrap();
1675 if barrier { fio.fail.store(true, std::sync::atomic::Ordering::Relaxed); }
1676 if barrier {
1677 assert!(s.commit().is_err());
1678 assert!(matches!(s.checkpoint(), Err(Error::StorePoisoned)));
1679 } assert_eq!(reader.get(b"old").unwrap(), Some(b"durable".to_vec()));
1681 assert_eq!(reader.get(b"new").unwrap(), None);
1682 fio.fail.store(false, std::sync::atomic::Ordering::Relaxed);
1683 drop(s); drop(fio);
1684 let s = Store::open(d.path(), cfg()).unwrap();
1685 assert_eq!(s.get(b"old").unwrap(), Some(b"durable".to_vec()));
1686 assert_eq!(s.get(b"new").unwrap(), None);
1687 }
1688 }
1689
1690 #[test]
1691 fn constrained_commit_data_write_failures_do_not_publish_partial_rows() {
1692 for at in 0..3 {
1693 let d = tempfile::tempdir().unwrap();
1694 let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1695 let fio = Arc::new(FailingWrite { inner: real, armed: std::sync::Mutex::new(None) });
1696 let limits = crate::limits::ResourceLimits { data_bytes: 1 << 20, wal_bytes: 64 << 10,
1697 tracked_pages: 256, readers: 2, record_bytes: 16000, recovery_bytes: 65536 };
1698 let mut s = Store::build_on_limited(d.path(), cfg(), true, fio.clone(), IoMode::Buffered, Some(limits)).unwrap();
1699 s.put(b"old", b"durable").unwrap(); s.commit().unwrap();
1700 let reader = Store::open_snapshot(d.path(), cfg()).unwrap();
1701 s.put(b"new", &[8;12000]).unwrap();
1702 *fio.armed.lock().unwrap() = Some(at);
1703 assert!(s.commit().is_err());
1704 assert!(matches!(s.put(b"later", b"no"), Err(Error::StorePoisoned)));
1705 assert_eq!(reader.get(b"old").unwrap(), Some(b"durable".to_vec()));
1706 *fio.armed.lock().unwrap() = None;
1707 drop(s); drop(fio);
1708 let s = Store::open(d.path(), cfg()).unwrap();
1709 assert_eq!(s.get(b"old").unwrap(), Some(b"durable".to_vec()));
1710 assert_eq!(s.get(b"new").unwrap(), None);
1711 }
1712 }
1713
1714 #[test]
1722 fn a_snapshot_is_registered_before_its_generation_can_be_recycled() {
1723 let d = tempfile::tempdir().unwrap();
1724 let tiny = Config { budget_bytes: 1 << 16, io: IoMode::Buffered, sync: SyncMode::Off };
1725 let mut writer = Store::create(d.path(), tiny).unwrap();
1726 for i in 0..2_000u64 {
1727 writer.put(&i.to_be_bytes(), format!("generation one row {i}").as_bytes()).unwrap();
1728 }
1729 writer.commit().unwrap();
1730 writer.checkpoint().unwrap();
1731
1732 let (selected_tx, selected_rx) = std::sync::mpsc::channel();
1733 let (churned_tx, churned_rx) = std::sync::mpsc::channel();
1734 let result = std::thread::scope(|scope| {
1735 let dir = d.path();
1736 let reader = scope.spawn(move || {
1737 let snapshot = Store::open_snapshot_with_after_meta(dir, tiny, |generation| {
1738 selected_tx.send(generation).unwrap();
1739 churned_rx.recv().unwrap();
1740 Ok(())
1741 })?;
1742 let iter = snapshot.scan(&[])?;
1743 iter.collect::<Result<Vec<_>>>()
1744 });
1745 let writer_thread = scope.spawn(move || {
1746 assert_eq!(selected_rx.recv().unwrap(), 1, "fixture must stop after selecting generation 1");
1747 for round in 0..4u64 {
1748 for i in 0..2_000u64 {
1749 writer.put(&i.to_be_bytes(),
1750 format!("writer round {round} row {i} is different").as_bytes()).unwrap();
1751 }
1752 writer.commit().unwrap();
1753 writer.checkpoint().unwrap();
1754 }
1755 churned_tx.send(()).unwrap();
1756 });
1757 writer_thread.join().unwrap();
1758 reader.join().unwrap()
1759 });
1760
1761 let rows = result.expect("a snapshot must not reach recycled pages while it is opening");
1762 assert_eq!(rows.len(), 2_000);
1763 for (i, (key, value)) in rows.iter().enumerate() {
1764 assert_eq!(key.as_slice(), &(i as u64).to_be_bytes());
1765 assert_eq!(value.as_slice(), format!("generation one row {i}").as_bytes(),
1766 "the opening snapshot observed a recycled page at row {i}");
1767 }
1768 }
1769
1770 struct CrashDirectory {
1771 inner: Box<dyn FileIo>,
1772 directory_durable: std::sync::atomic::AtomicBool,
1773 }
1774
1775 impl FileIo for CrashDirectory {
1776 fn requires_alignment(&self) -> bool { self.inner.requires_alignment() }
1777 fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> { self.inner.read_at(buf, off) }
1778 fn write_at(&self, buf: &[u8], off: u64) -> Result<()> { self.inner.write_at(buf, off) }
1779 fn sync_data(&self) -> Result<()> { self.inner.sync_data() }
1780 fn sync_full(&self) -> Result<()> { self.inner.sync_full() }
1781 fn sync_full_primitive(&self) -> &'static str { self.inner.sync_full_primitive() }
1782 fn sync_dir(&self) -> Result<()> {
1783 self.inner.sync_dir()?;
1784 self.directory_durable.store(true, std::sync::atomic::Ordering::Release);
1785 Ok(())
1786 }
1787 fn len(&self) -> Result<u64> { self.inner.len() }
1788 fn set_len(&self, n: u64) -> Result<()> { self.inner.set_len(n) }
1789 }
1790
1791 #[test]
1796 fn a_commit_before_the_first_checkpoint_survives_a_directory_crash() {
1797 let d = tempfile::tempdir().unwrap();
1798 let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1799 let fio = Arc::new(CrashDirectory {
1800 inner: real,
1801 directory_durable: std::sync::atomic::AtomicBool::new(false),
1802 });
1803 {
1804 let mut s = Store::create_on(d.path(), cfg(), fio.clone()).unwrap();
1805 s.put(b"acknowledged", b"must survive power loss").unwrap();
1806 s.commit().unwrap();
1807 }
1808
1809 if !fio.directory_durable.load(std::sync::atomic::Ordering::Acquire) {
1810 std::fs::remove_file(d.path().join("data")).unwrap();
1811 std::fs::remove_file(d.path().join("wal")).unwrap();
1812 }
1813 let reopened = Store::open(d.path(), cfg()).unwrap();
1814 assert_eq!(reopened.get(b"acknowledged").unwrap().as_deref(),
1815 Some(&b"must survive power loss"[..]),
1816 "creation must make file names durable before any commit can be acknowledged");
1817 }
1818
1819 #[test]
1820 fn recreated_wal_name_survives_a_commit_before_checkpoint() {
1821 let d = tempfile::tempdir().unwrap();
1822 {
1823 let mut s = Store::create(d.path(), cfg()).unwrap();
1824 s.put(b"old", b"published").unwrap();
1825 s.checkpoint().unwrap();
1826 }
1827 std::fs::remove_file(d.path().join("wal")).unwrap();
1828 let (real, mode) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1829 let fio = Arc::new(CrashDirectory {
1830 inner: real,
1831 directory_durable: std::sync::atomic::AtomicBool::new(false),
1832 });
1833 {
1834 let mut s = Store::build_on(d.path(), cfg(), false, fio.clone(), mode).unwrap();
1835 s.put(b"new", b"acknowledged").unwrap();
1836 s.commit().unwrap();
1837 }
1838 if !fio.directory_durable.load(std::sync::atomic::Ordering::Acquire) {
1839 std::fs::remove_file(d.path().join("wal")).unwrap();
1840 }
1841 drop(fio);
1842 let s = Store::open(d.path(), cfg()).unwrap();
1843 assert_eq!(s.get(b"old").unwrap().as_deref(), Some(&b"published"[..]));
1844 assert_eq!(s.get(b"new").unwrap().as_deref(), Some(&b"acknowledged"[..]));
1845 }
1846
1847 #[test]
1848 fn a_maximum_generation_read_from_disk_is_refused_not_wrapped() {
1849 let d = tempfile::tempdir().unwrap();
1850 {
1851 let s = Store::create(d.path(), cfg()).unwrap();
1852 Meta {
1853 format_version: crate::meta::FORMAT_VERSION,
1854 roots: [s.root, 0, 0, 0, 0, 0, 0, 0],
1855 next_lsn: s.wal.as_ref().unwrap().next_lsn(),
1856 generation: u64::MAX,
1857 }.write_slot(&s.pool).unwrap();
1858 s.pool.flush_all(Barrier::None).unwrap();
1859 }
1860
1861 assert!(matches!(Store::open(d.path(), cfg()), Err(crate::Error::Corrupt { .. })),
1862 "the generation after u64::MAX does not exist and must not become zero");
1863 }
1864
1865 #[test]
1866 fn a_checkpoint_refuses_when_no_later_page_generation_exists() {
1867 let d = tempfile::tempdir().unwrap();
1868 let mut s = Store::create(d.path(), cfg()).unwrap();
1869 s.generation = u64::MAX - 1;
1870 s.put(b"pending", b"kept in the log").unwrap();
1871 s.commit().unwrap();
1872
1873 assert!(matches!(s.checkpoint(), Err(crate::Error::Corrupt { .. })),
1874 "publishing the final generation would leave the next epoch wrapping to zero");
1875 assert!(std::fs::metadata(d.path().join("wal")).unwrap().len() > 0,
1876 "refusing exhaustion must preserve the committed log");
1877 }
1878
1879 #[test]
1880 fn a_checkpoint_rotates_the_log_only_after_the_pages_are_durable() {
1881 let d = tempfile::tempdir().unwrap();
1882 let mut s = Store::create(d.path(), cfg()).unwrap();
1883 for i in 0..2000u64 { s.put(&i.to_be_bytes(), b"v").unwrap(); }
1884 s.commit().unwrap();
1885 s.checkpoint().unwrap();
1886 let order = s.checkpoint_trace();
1887 assert_eq!(order, vec!["flush_pages", "sync_file", "flip_meta", "rotate_wal"],
1892 "checkpoint order must be: data durable, then flip, then drop the log");
1893 assert!(s.get(&1999u64.to_be_bytes()).unwrap().is_some());
1894 }
1895
1896 #[test]
1906 fn a_rotation_followed_by_a_reopen_does_not_reissue_lsn_1() {
1907 let d = tempfile::tempdir().unwrap();
1908 let lsn_at_checkpoint = {
1909 let mut s = Store::create(d.path(), cfg()).unwrap();
1910 for i in 0..500u64 { s.put(&i.to_be_bytes(), b"v").unwrap(); }
1911 s.commit().unwrap();
1912 s.checkpoint().unwrap(); s.next_lsn()
1914 };
1915 assert!(lsn_at_checkpoint > 1, "sanity: many records were appended before the rotation");
1916
1917 let s2 = Store::open(d.path(), cfg()).unwrap();
1918 assert_eq!(
1919 s2.next_lsn(), lsn_at_checkpoint,
1920 "a reopen after rotation must not renumber LSNs from 1"
1921 );
1922 }
1923
1924 #[test]
1930 fn the_three_durability_modes_issue_different_barriers() {
1931 for (mode, want) in [
1932 (SyncMode::Full, vec!["sync_full"]),
1933 (SyncMode::Normal, vec!["sync_data"]),
1934 (SyncMode::Off, vec![]),
1935 ] {
1936 let d = tempfile::tempdir().unwrap();
1937 let cfg = Config { budget_bytes: 16 << 20, io: IoMode::Buffered, sync: mode };
1938 let mut s = Store::create(d.path(), cfg).unwrap();
1939 s.put(b"k", b"v").unwrap();
1940 s.commit().unwrap();
1941 assert_eq!(s.barriers(), want, "{mode:?} issued the wrong barrier");
1942 }
1943 }
1944
1945 fn armed(s: &Store) -> (u64, u64) {
1960 (s.fast_path_hits.get() + s.tag_hints.hits(),
1961 s.fast_path_attempts.get() + s.tag_hints.attempts())
1962 }
1963
1964 #[test]
1980 fn store_put_ascending_uses_fast_path() {
1981 let d = tempfile::tempdir().unwrap();
1982 let mut s = Store::create(d.path(), cfg()).unwrap();
1983 let n = 100u64;
1984 for i in 0..n { s.put(&i.to_be_bytes(), b"v").unwrap(); }
1985 assert_eq!(
1986 armed(&s).0, n - 1,
1987 "every Store::put but the first must hit the append fast path"
1988 );
1989 }
1990
1991 #[test]
1998 fn store_put_survives_bulk_load() {
1999 let d = tempfile::tempdir().unwrap();
2000 let mut s = Store::create(d.path(), cfg()).unwrap();
2001 let n = 1_000u64;
2002 let items = (0..n).map(|i| (i.to_be_bytes().to_vec(), b"bulk".to_vec()));
2003 s.bulk_load(items).unwrap();
2004 assert_eq!(s.last_leaf.get(), None, "bulk_load must clear a hint it just made meaningless");
2005
2006 let tail_key = n.to_be_bytes();
2008 s.put(&tail_key, b"tail").unwrap();
2009
2010 assert_eq!(
2011 s.get(&tail_key).unwrap().as_deref(), Some(&b"tail"[..]),
2012 "a put right after bulk_load must be found"
2013 );
2014
2015 let scanned: Vec<Vec<u8>> = s.scan(&[]).unwrap().map(|r| r.unwrap().0).collect();
2016 let mut sorted = scanned.clone();
2017 sorted.sort();
2018 assert_eq!(scanned.len(), n as usize + 1);
2019 assert_eq!(scanned, sorted, "a full scan after bulk_load + put must stay in sorted order");
2020 }
2021
2022 #[test]
2029 fn store_random_order_matches_sequential() {
2030 let n = 20_000u64;
2031 let scatter = |i: u64| i.wrapping_mul(0x9E37_79B9_7F4A_7C15);
2032
2033 let d1 = tempfile::tempdir().unwrap();
2034 let mut s1 = Store::create(d1.path(), cfg()).unwrap();
2035 for i in 0..n { s1.put(&i.to_be_bytes(), &i.to_le_bytes()).unwrap(); }
2036
2037 let d2 = tempfile::tempdir().unwrap();
2038 let mut s2 = Store::create(d2.path(), cfg()).unwrap();
2039 let mut order: Vec<u64> = (0..n).collect();
2040 order.sort_by_key(|&i| scatter(i));
2041 for &i in &order { s2.put(&i.to_be_bytes(), &i.to_le_bytes()).unwrap(); }
2042
2043 let seq1: Vec<(Vec<u8>, Vec<u8>)> = s1.scan(&[]).unwrap().map(|r| r.unwrap()).collect();
2044 let seq2: Vec<(Vec<u8>, Vec<u8>)> = s2.scan(&[]).unwrap().map(|r| r.unwrap()).collect();
2045 assert_eq!(seq1.len(), n as usize);
2046 assert_eq!(
2047 seq1, seq2,
2048 "ascending vs scattered insertion order through Store::put must produce identical scans"
2049 );
2050 }
2051
2052 #[test]
2079 fn random_order_disarms_fast_path() {
2080 let d = tempfile::tempdir().unwrap();
2081 let mut s = Store::create(d.path(), cfg()).unwrap();
2082 let n = 20_000u64;
2083 let scatter = |i: u64| i.wrapping_mul(0x9E37_79B9_7F4A_7C15);
2084 let payload = vec![b'x'; 200]; for i in 0..n { s.put(&scatter(i).to_be_bytes(), &payload).unwrap(); }
2086 #[cfg(feature = "sqlite-balance")]
2090 assert!(s.fast_path_attempts.get() <= 117);
2091 assert!(
2099 s.tag_hints.attempts() <= 2_000,
2100 "a scattered workload attempted the per-keyspace fast path {} times in {n} \
2101 inserts; arming belongs to appends, not to every descent",
2102 s.tag_hints.attempts()
2103 );
2104 #[cfg(not(feature = "sqlite-balance"))]
2105 assert_eq!(
2106 s.fast_path_attempts.get(), 117,
2107 "a disarming hint must attempt the fast path a number of times bounded by \
2108 leaf capacity and log(n), not by n -- a bound proportional to n would not \
2109 have caught the pre-Task-18 defect this test exists for"
2110 );
2111 }
2112
2113 #[test]
2123 fn ascending_keeps_fast_path_armed() {
2124 let d = tempfile::tempdir().unwrap();
2125 let mut s = Store::create(d.path(), cfg()).unwrap();
2126 let n = 100u64;
2127 for i in 0..n { s.put(&i.to_be_bytes(), b"v").unwrap(); }
2128 assert_eq!(
2129 armed(&s).0, n - 1,
2130 "an unbroken ascending run must still hit the fast path on every insert but the first"
2131 );
2132 }
2133
2134 #[test]
2167 fn mixed_workload_rearms() {
2168 let d = tempfile::tempdir().unwrap();
2169 let mut s = Store::create(d.path(), cfg()).unwrap();
2170 let m = 300u64;
2171 let k = 30u64;
2172
2173 for i in 1..=m { s.put(&i.to_be_bytes(), b"v").unwrap(); }
2174 assert_eq!(armed(&s).1, m - 1, "one attempt per insert but the first");
2175 assert_eq!(
2176 armed(&s).0, m - 2,
2177 "one miss expected: the insert that fills the leaf and forces the one split \
2178 this run crosses"
2179 );
2180
2181 s.put(&0u64.to_be_bytes(), b"v").unwrap();
2182 assert_eq!(armed(&s).1, m, "the out-of-order key is one more attempt");
2183 assert_eq!(armed(&s).0, m - 2, "the out-of-order key must not hit");
2184 assert_eq!(
2185 s.last_leaf.get(), None,
2186 "landing on a non-rightmost leaf must leave the hint disarmed, not re-armed \
2187 to the wrong leaf"
2188 );
2189
2190 let (hits, attempts) = armed(&s);
2191 for i in (m + 1)..(m + 1 + k) { s.put(&i.to_be_bytes(), b"v").unwrap(); }
2192 let (hits, attempts) = (armed(&s).0 - hits, armed(&s).1 - attempts);
2193 assert_eq!(
2194 attempts, hits,
2195 "an unbroken ascending run must hit every probe it makes"
2196 );
2197 assert!(
2205 k - hits <= 1,
2206 "{} of the {k} ascending inserts after the disarm paid a descent",
2207 k - hits
2208 );
2209 }
2210
2211 struct FailingBarrier { inner: Box<dyn FileIo>, fail: std::sync::atomic::AtomicBool }
2220 impl FailingBarrier {
2221 fn failing(&self) -> bool { self.fail.load(std::sync::atomic::Ordering::Relaxed) }
2222 fn eio() -> crate::Error {
2223 std::io::Error::other("injected barrier failure").into()
2224 }
2225 }
2226 impl FileIo for FailingBarrier {
2227 fn requires_alignment(&self) -> bool { self.inner.requires_alignment() }
2228 fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> { self.inner.read_at(buf, off) }
2229 fn write_at(&self, buf: &[u8], off: u64) -> Result<()> { self.inner.write_at(buf, off) }
2230 fn sync_data(&self) -> Result<()> {
2231 if self.failing() { return Err(Self::eio()); }
2232 self.inner.sync_data()
2233 }
2234 fn sync_full(&self) -> Result<()> {
2235 if self.failing() { return Err(Self::eio()); }
2236 self.inner.sync_full()
2237 }
2238 fn sync_full_primitive(&self) -> &'static str { self.inner.sync_full_primitive() }
2239 fn sync_dir(&self) -> Result<()> { self.inner.sync_dir() }
2240 fn len(&self) -> Result<u64> { self.inner.len() }
2241 fn set_len(&self, n: u64) -> Result<()> { self.inner.set_len(n) }
2242 }
2243
2244 struct FailingWrite {
2248 inner: Box<dyn FileIo>,
2249 armed: std::sync::Mutex<Option<usize>>,
2250 }
2251 impl FailingWrite {
2252 fn arm(&self, successful_writes: usize) {
2253 *self.armed.lock().unwrap() = Some(successful_writes);
2254 }
2255 }
2256 impl FileIo for FailingWrite {
2257 fn requires_alignment(&self) -> bool { self.inner.requires_alignment() }
2258 fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> { self.inner.read_at(buf, off) }
2259 fn write_at(&self, buf: &[u8], off: u64) -> Result<()> {
2260 let mut armed = self.armed.lock().unwrap();
2261 if let Some(remaining) = armed.as_mut() {
2262 if *remaining == 0 {
2263 *armed = None;
2264 return Err(std::io::Error::other(
2265 "injected page-write failure during a tree mutation",
2266 ).into());
2267 }
2268 *remaining -= 1;
2269 }
2270 self.inner.write_at(buf, off)
2271 }
2272 fn sync_data(&self) -> Result<()> { self.inner.sync_data() }
2273 fn sync_full(&self) -> Result<()> { self.inner.sync_full() }
2274 fn sync_full_primitive(&self) -> &'static str { self.inner.sync_full_primitive() }
2275 fn sync_dir(&self) -> Result<()> { self.inner.sync_dir() }
2276 fn len(&self) -> Result<u64> { self.inner.len() }
2277 fn set_len(&self, n: u64) -> Result<()> { self.inner.set_len(n) }
2278 }
2279
2280 #[test]
2281 fn deletion_merge_eviction_failures_cannot_publish_partial_changes() {
2282 for fail_after in 0..12 {
2283 let d=tempfile::tempdir().unwrap();
2284 let cfg=Config{budget_bytes:64<<10,io:IoMode::Buffered,sync:SyncMode::Full};
2285 let (real,_)=open_file(&d.path().join("data"),IoMode::Buffered).unwrap();
2286 let fio=Arc::new(FailingWrite{inner:real,armed:std::sync::Mutex::new(None)});
2287 let mut s=Store::create_on(d.path(),cfg,fio.clone()).unwrap();
2288 for i in 0u64..1024 {s.put(&i.to_be_bytes(),&vec![7;240]).unwrap();}
2289 for i in 0u64..1024 {if i%15>=6 {s.delete(&i.to_be_bytes()).unwrap();}}
2292 s.checkpoint().unwrap();
2293 let old=Store::open_snapshot(d.path(),cfg).unwrap();
2294 fio.arm(fail_after);let mut failed=false;
2295 for i in 0u64..1024 {
2296 if let Err(e)=s.delete(&((i*71)%1024).to_be_bytes()) {
2297 assert!(matches!(e,crate::Error::Io(_)));failed=true;break;
2298 }
2299 }
2300 assert!(failed,"fault {fail_after} was not reached");
2301 assert!(s.commit().is_err());assert!(s.checkpoint().is_err());
2302 for i in 0u64..1024 {assert_eq!(old.get(&i.to_be_bytes()).unwrap(),(i%15<6).then(||vec![7;240]));}
2303 drop(old);drop(s);
2304 let reopened=Store::open(d.path(),cfg).unwrap();
2305 for i in 0u64..1024 {assert_eq!(reopened.get(&i.to_be_bytes()).unwrap(),(i%15<6).then(||vec![7;240]));}
2306 assert_eq!(crate::verify::verify_published_tree(&d.path().join("data"),IoMode::Buffered,reopened.root,1).unwrap().0,412);
2307 }
2308 }
2309
2310 #[test]
2316 fn a_partial_tree_mutation_poisons_and_reopen_replays_the_retained_log() {
2317 let d = tempfile::tempdir().unwrap();
2318 let tiny = Config {
2319 budget_bytes: 16 * crate::page::PAGE_SIZE,
2320 io: IoMode::Buffered,
2321 sync: SyncMode::Off,
2322 };
2323 let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
2324 let fio = Arc::new(FailingWrite {
2325 inner: real,
2326 armed: std::sync::Mutex::new(None),
2327 });
2328 let mut s = Store::create_on(d.path(), tiny, fio.clone()).unwrap();
2329
2330 let value = vec![b'v'; 256];
2331 let record_bytes = 4 + 8 + value.len() + 4;
2332 let mut rows = 0u64;
2333 loop {
2334 let room = {
2335 let r = s.pool.get(s.root).unwrap();
2336 crate::page::PageRef::open_resident(&r, s.root).unwrap().free_space()
2337 };
2338 if room < record_bytes { break; }
2339 s.put(&(rows * 2).to_be_bytes(), &value).unwrap();
2340 s.commit().unwrap();
2341 rows += 1;
2342 }
2343 assert!(rows > 4, "fixture must have committed rows to move across the split");
2344
2345 while s.pool.page_count() < 16 {
2350 drop(s.pool.allocate().unwrap());
2351 }
2352 {
2353 let root_pin = s.pool.get(s.root).unwrap();
2354 for _ in 0..15 { drop(s.pool.allocate().unwrap()); }
2355 drop(root_pin);
2356 }
2357
2358 let wal_before = std::fs::metadata(d.path().join("wal")).unwrap().len();
2359 fio.arm(1);
2360 let failed = s.put(&1u64.to_be_bytes(), &value);
2361 assert!(matches!(failed, Err(crate::Error::Io(_))),
2362 "the injected eviction failure must escape the mutating insert, got {failed:?}");
2363 assert_eq!(s.get(&((rows - 1) * 2).to_be_bytes()).unwrap(), None,
2364 "fixture must prove the insert failed after committed keys moved off the root");
2365
2366 let commit = s.commit();
2370 let checkpoint = s.checkpoint();
2371 let wal_after = std::fs::metadata(d.path().join("wal")).unwrap().len();
2372 assert!(matches!(commit, Err(crate::Error::StorePoisoned)),
2373 "a partial tree mutation must make commit refuse, got {commit:?}");
2374 assert!(matches!(checkpoint, Err(crate::Error::StorePoisoned)),
2375 "a partial tree mutation must make checkpoint refuse, got {checkpoint:?}");
2376 assert_eq!(wal_after, wal_before,
2377 "a poisoned store must retain the committed recovery log byte-for-byte");
2378
2379 drop(s);
2380 drop(fio);
2381 let reopened = Store::open(d.path(), tiny)
2382 .expect("reopen must rebuild a poisoned handle from the retained log");
2383 for i in 0..rows {
2384 assert_eq!(reopened.get(&(i * 2).to_be_bytes()).unwrap().as_deref(), Some(value.as_slice()));
2385 }
2386 assert_eq!(reopened.get(&1u64.to_be_bytes()).unwrap(), None,
2387 "the failed, uncommitted insert must not be replayed");
2388 }
2389
2390 #[test]
2394 fn a_pre_mutation_validation_error_does_not_poison_the_store() {
2395 let d = tempfile::tempdir().unwrap();
2396 let mut s = Store::create(d.path(), cfg()).unwrap();
2397 let oversized_key = vec![b'k'; crate::page::MAX_RECORD_LEN];
2398 assert!(matches!(s.put(&oversized_key, b"v"), Err(crate::Error::TooLarge)));
2399 s.put(b"usable", b"still").unwrap();
2400 s.commit().unwrap();
2401 s.checkpoint().unwrap();
2402 assert_eq!(s.get(b"usable").unwrap().as_deref(), Some(&b"still"[..]));
2403 }
2404
2405 #[test]
2414 fn a_failed_checkpoint_barrier_poisons_every_writer() {
2415 let d = tempfile::tempdir().unwrap();
2416 let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
2417 let fio = Arc::new(FailingBarrier { inner: real, fail: false.into() });
2418 let mut s = Store::create_on(d.path(), cfg(), fio.clone()).unwrap();
2419 s.put(b"a", b"1").unwrap();
2420 s.commit().unwrap();
2421 s.checkpoint().expect("sanity: the checkpoint works while the disk does");
2422
2423 s.put(b"b", b"2").unwrap();
2424 s.commit().unwrap();
2425 fio.fail.store(true, std::sync::atomic::Ordering::Relaxed);
2426
2427 match s.checkpoint() {
2428 Err(crate::Error::Io(_)) => {}
2429 other => panic!("a failing barrier must surface, got {other:?}"),
2430 }
2431 fio.fail.store(false, std::sync::atomic::Ordering::Relaxed); assert!(matches!(s.put(b"c", b"3"), Err(crate::Error::StorePoisoned)));
2434 assert!(matches!(s.delete(b"a"), Err(crate::Error::StorePoisoned)));
2435 assert!(matches!(s.commit(), Err(crate::Error::StorePoisoned)));
2436 assert!(matches!(s.checkpoint(), Err(crate::Error::StorePoisoned)),
2437 "a retried checkpoint would flush nothing (those frames are marked clean \
2438 already) and then rotate the log away");
2439 let items = vec![(b"z".to_vec(), b"9".to_vec())].into_iter();
2440 assert!(matches!(s.bulk_load(items), Err(crate::Error::StorePoisoned)),
2441 "bulk_load replaces the whole tree AND discards the log");
2442 assert_eq!(s.get(b"a").unwrap().as_deref(), Some(&b"1"[..]),
2448 "a refused bulk_load must not have replaced the tree");
2449 assert_eq!(s.get(b"z").unwrap(), None);
2450 }
2451
2452 #[test]
2459 fn reopening_clears_poisoning_and_the_committed_data_is_still_there() {
2460 let d = tempfile::tempdir().unwrap();
2461 let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
2462 let fio = Arc::new(FailingBarrier { inner: real, fail: false.into() });
2463 {
2464 let mut s = Store::create_on(d.path(), cfg(), fio.clone()).unwrap();
2465 s.put(b"a", b"1").unwrap();
2466 s.commit().unwrap();
2467 fio.fail.store(true, std::sync::atomic::Ordering::Relaxed);
2468 assert!(s.checkpoint().is_err());
2469 assert!(matches!(s.put(b"b", b"2"), Err(crate::Error::StorePoisoned)));
2470 }
2471 fio.fail.store(false, std::sync::atomic::Ordering::Relaxed);
2472
2473 let mut s2 = Store::open(d.path(), cfg()).expect("a poisoned store must not poison the DIRECTORY");
2474 assert_eq!(s2.get(b"a").unwrap().as_deref(), Some(&b"1"[..]),
2475 "the committed row is in the log, which the failed checkpoint never rotated");
2476 s2.put(b"b", b"2").unwrap();
2477 s2.commit().unwrap();
2478 }
2479
2480 #[test]
2484 fn malformed_wal_payloads_are_bounded_and_do_not_mutate_the_tree() {
2485 let d = tempfile::tempdir().unwrap();
2486 let mut s = Store::create(d.path(), cfg()).unwrap();
2487 s.put(b"kept", b"value").unwrap();
2488 s.commit().unwrap();
2489 let root = s.root;
2490
2491 let malformed: &[(RecKind, &[u8])] = &[
2492 (RecKind::Put, &[]),
2493 (RecKind::Put, &[4, 0, b'a']),
2494 (RecKind::Delete, &[]),
2495 (RecKind::Delete, &[2, 0, b'a']),
2496 (RecKind::Delete, &[0, 0, b'x']),
2497 (RecKind::DeletePrefix, &[]),
2498 (RecKind::PutEmptyBatch, &[]),
2499 (RecKind::PutEmptyBatch, &[1, 0, 4, 0, b'a']),
2500 (RecKind::PutEmptyBatch, &[65, 0]),
2501 (RecKind::PutEmptyBatch, &[1, 0, 1, 0, b'a', b'x']),
2502 (RecKind::Commit, b"not empty"),
2503 (RecKind::PageImage, b"unsupported"),
2504 ];
2505 for (n, &(kind, payload)) in malformed.iter().enumerate() {
2506 assert!(matches!(
2507 s.apply(kind, payload, n as u64),
2508 Err(crate::Error::CorruptWal { offset, .. }) if offset == n as u64
2509 ));
2510 assert_eq!(s.root, root, "malformed frame {n} changed the root");
2511 assert_eq!(s.get(b"kept").unwrap().as_deref(), Some(&b"value"[..]),
2512 "malformed frame {n} changed existing data");
2513 }
2514 }
2515
2516 #[test]
2517 fn bounded_empty_key_batch_replays_as_one_committed_wal_unit() {
2518 let d = tempfile::tempdir().unwrap();
2519 {
2520 let mut s = Store::create(d.path(), cfg()).unwrap();
2521 s.put_empty_batch(&[b"alpha".to_vec(), b"beta".to_vec()]).unwrap();
2522 s.commit().unwrap();
2523 }
2524 let s = Store::open(d.path(), cfg()).unwrap();
2525 assert_eq!(s.get(b"alpha").unwrap().as_deref(), Some(&b""[..]));
2526 assert_eq!(s.get(b"beta").unwrap().as_deref(), Some(&b""[..]));
2527 }
2528
2529 #[test]
2530 fn a_corrupted_packed_tree_never_becomes_authoritative() {
2531 let d = tempfile::tempdir().unwrap();
2532 let mut s = Store::create(d.path(), cfg()).unwrap();
2533 s.put(b"old", b"authoritative").unwrap();
2534 s.commit().unwrap();
2535 s.checkpoint().unwrap();
2536 let old_root = s.root;
2537
2538 let result = s.bulk_load_with_before_publish(
2539 (0..2_000u64).map(|i| (i.to_be_bytes().to_vec(), b"new".to_vec())),
2540 |data, root| {
2541 use std::io::{Seek, SeekFrom, Write};
2542 let mut file = std::fs::OpenOptions::new().write(true).open(data)?;
2543 file.seek(SeekFrom::Start(root as u64 * PAGE_SIZE as u64))?;
2544 file.write_all(&[0u8])?;
2545 file.sync_all()?;
2546 Ok(())
2547 },
2548 );
2549
2550 assert!(result.is_err(), "a packed tree corrupted before publication must be refused");
2551 assert_eq!(s.root, old_root, "the old root must stay authoritative after refusal");
2552 assert_eq!(s.get(b"old").unwrap().as_deref(), Some(&b"authoritative"[..]));
2553 }
2554
2555 fn graft_key(space: u8, i: u32) -> Vec<u8> {
2556 let mut key = vec![space];
2557 key.extend_from_slice(&i.to_be_bytes());
2558 key
2559 }
2560
2561 fn seed_graft_base(store: &mut Store) -> Vec<(Vec<u8>, Vec<u8>)> {
2562 let mut rows = Vec::new();
2563 for space in [0x10, 0x30] {
2564 for i in 0..1_500u32 {
2565 let row = (graft_key(space, i), format!("base-{space:02x}-{i}").into_bytes());
2566 store.put(&row.0, &row.1).unwrap();
2567 rows.push(row);
2568 }
2569 }
2570 store.commit().unwrap();
2571 store.checkpoint().unwrap();
2572 rows.sort_by(|a, b| a.0.cmp(&b.0));
2573 rows
2574 }
2575
2576 fn graft_rows() -> Vec<(Vec<u8>, Vec<u8>)> {
2577 (0..2_000u32)
2578 .rev()
2579 .map(|i| (graft_key(0x20, i), format!("graft-{i}").into_bytes()))
2580 .collect()
2581 }
2582
2583 fn collect_from(store: &Store, from: &[u8]) -> Vec<(Vec<u8>, Vec<u8>)> {
2584 store.scan(from).unwrap().map(Result::unwrap).collect()
2585 }
2586
2587 fn collect_below(store: &Store, to: &[u8]) -> Vec<(Vec<u8>, Vec<u8>)> {
2588 let mut rows = Vec::new();
2589 store.scan_reverse(to).unwrap().for_each_ref(|key, value| {
2590 rows.push((key.to_vec(), value.to_vec()));
2591 true
2592 }).unwrap();
2593 rows
2594 }
2595
2596 #[test]
2600 fn graft_preserves_every_preexisting_key() {
2601 let d = tempfile::tempdir().unwrap();
2602 let mut s = Store::create(d.path(), cfg()).unwrap();
2603 let base = seed_graft_base(&mut s);
2604 let pinned = Store::open_snapshot(d.path(), cfg()).unwrap();
2605
2606 s.graft_range(graft_rows().into_iter()).unwrap();
2607
2608 for (key, value) in &base {
2609 assert_eq!(s.get(key).unwrap().as_deref(), Some(value.as_slice()),
2610 "graft lost pre-existing key {key:?}");
2611 }
2612 assert_eq!(s.scan(&[]).unwrap().count(), base.len() + 2_000);
2613 assert!(pinned.get(&graft_key(0x20, 17)).unwrap().is_none(),
2614 "a reader pinned before publication must remain on the old generation");
2615 assert_eq!(pinned.scan(&[]).unwrap().count(), base.len());
2616 let fresh = Store::open_snapshot(d.path(), cfg()).unwrap();
2617 assert_eq!(fresh.get(&graft_key(0x20, 17)).unwrap().as_deref(), Some(&b"graft-17"[..]));
2618 }
2619
2620 #[test]
2624 fn graft_and_individual_inserts_answer_every_kernel_query_identically() {
2625 let dg = tempfile::tempdir().unwrap();
2626 let di = tempfile::tempdir().unwrap();
2627 let mut grafted = Store::create(dg.path(), cfg()).unwrap();
2628 let mut inserted = Store::create(di.path(), cfg()).unwrap();
2629 seed_graft_base(&mut grafted);
2630 seed_graft_base(&mut inserted);
2631 let rows = graft_rows();
2632
2633 grafted.graft_range(rows.clone().into_iter()).unwrap();
2634 for (key, value) in &rows { inserted.put(key, value).unwrap(); }
2635 inserted.commit().unwrap();
2636 inserted.checkpoint().unwrap();
2637
2638 for (key, _) in seed_query_keys(&rows) {
2639 assert_eq!(grafted.get(&key).unwrap(), inserted.get(&key).unwrap(),
2640 "point query disagreed at {key:?}");
2641 }
2642 for from in [vec![], graft_key(0x10, 777), graft_key(0x20, 0),
2643 graft_key(0x20, 999), graft_key(0x30, 0), vec![0xff]] {
2644 assert_eq!(collect_from(&grafted, &from), collect_from(&inserted, &from),
2645 "forward range disagreed from {from:?}");
2646 }
2647 for to in [graft_key(0x10, 0), graft_key(0x20, 0), graft_key(0x20, 999),
2648 graft_key(0x30, 0), vec![0xff]] {
2649 assert_eq!(collect_below(&grafted, &to), collect_below(&inserted, &to),
2650 "reverse range disagreed below {to:?}");
2651 }
2652
2653 let live_key = graft_key(0x20, 2_500);
2657 grafted.put(&live_key, b"later-live-write").unwrap();
2658 inserted.put(&live_key, b"later-live-write").unwrap();
2659 grafted.commit().unwrap();
2660 inserted.commit().unwrap();
2661 grafted.checkpoint().unwrap();
2662 inserted.checkpoint().unwrap();
2663 assert_eq!(collect_from(&grafted, &graft_key(0x20, 1_900)),
2664 collect_from(&inserted, &graft_key(0x20, 1_900)));
2665 }
2666
2667 #[test]
2668 fn graft_refuses_a_nonempty_range_without_changing_the_tree() {
2669 let d = tempfile::tempdir().unwrap();
2670 let mut s = Store::create(d.path(), cfg()).unwrap();
2671 let base = seed_graft_base(&mut s);
2672 let old_root = s.root;
2673 let rows = vec![
2674 (graft_key(0x0f, 0), b"before".to_vec()),
2675 (graft_key(0x10, 10), b"overlap".to_vec()),
2676 ];
2677 assert!(matches!(s.graft_range(rows.into_iter()), Err(crate::Error::RangeNotEmpty)));
2678 assert_eq!(s.root, old_root);
2679 assert_eq!(s.scan(&[]).unwrap().count(), base.len());
2680 for (key, value) in base {
2681 assert_eq!(s.get(&key).unwrap().as_deref(), Some(value.as_slice()));
2682 }
2683 }
2684
2685 #[test]
2686 fn graft_into_an_empty_tree_becomes_the_tree() {
2687 let d = tempfile::tempdir().unwrap();
2688 let mut s = Store::create(d.path(), cfg()).unwrap();
2689 let rows = graft_rows();
2690 s.graft_range(rows.clone().into_iter()).unwrap();
2691 assert_eq!(s.scan(&[]).unwrap().count(), rows.len());
2692 for (key, value) in rows.iter().step_by(97) {
2693 assert_eq!(s.get(key).unwrap().as_deref(), Some(value.as_slice()));
2694 }
2695 }
2696
2697 fn seed_query_keys(rows: &[(Vec<u8>, Vec<u8>)]) -> Vec<(Vec<u8>, Vec<u8>)> {
2698 let mut keys = rows.to_vec();
2699 keys.push((graft_key(0x10, 0), Vec::new()));
2700 keys.push((graft_key(0x10, 1_499), Vec::new()));
2701 keys.push((graft_key(0x30, 0), Vec::new()));
2702 keys.push((graft_key(0x30, 1_499), Vec::new()));
2703 keys.push((graft_key(0x20, 2_001), Vec::new()));
2704 keys
2705 }
2706
2707 #[test]
2711 fn graft_crash_boundary_is_old_before_and_durable_after() {
2712 let before_dir = tempfile::tempdir().unwrap();
2713 {
2714 let mut s = Store::create(before_dir.path(), cfg()).unwrap();
2715 seed_graft_base(&mut s);
2716 let old_root = s.root;
2717 let result = s.graft_range_with_before_publish(
2718 graft_rows().into_iter(),
2719 |_, _| Err(std::io::Error::other("crash before graft").into()),
2720 );
2721 assert!(result.is_err());
2722 assert_eq!(s.root, old_root, "a pre-publication failure changed the live root");
2723 }
2724 let before = Store::open(before_dir.path(), cfg()).unwrap();
2725 assert!(before.get(&graft_key(0x20, 17)).unwrap().is_none());
2726 assert_eq!(before.scan(&[]).unwrap().count(), 3_000);
2727
2728 let after_dir = tempfile::tempdir().unwrap();
2729 {
2730 let mut s = Store::create(after_dir.path(), cfg()).unwrap();
2731 seed_graft_base(&mut s);
2732 s.graft_range(graft_rows().into_iter()).unwrap();
2733 }
2735 let after = Store::open(after_dir.path(), cfg()).unwrap();
2736 assert_eq!(after.get(&graft_key(0x20, 17)).unwrap().as_deref(), Some(&b"graft-17"[..]));
2737 assert_eq!(after.scan(&[]).unwrap().count(), 5_000);
2738 }
2739
2740 #[test]
2744 fn prepared_graft_publishes_after_reopen_and_resume_of_resume_is_idempotent() {
2745 let d = tempfile::tempdir().unwrap();
2746 let scratch = d.path().join("prepared-graft-scratch");
2747 let manifest = d.path().join("prepared-graft");
2748 {
2749 let mut s = Store::create(d.path(), cfg()).unwrap();
2750 seed_graft_base(&mut s);
2751 let mut rows = graft_rows();
2752 rows.sort_by(|a, b| a.0.cmp(&b.0));
2753 let min = rows.first().unwrap().0.clone();
2754 let max = rows.last().unwrap().0.clone();
2755 let prepared = s.prepare_graft_candidate(
2756 rows.into_iter().map(|(k, v)| Ok((k, v, false))),
2757 2_000, min, max, &scratch).unwrap();
2758 prepared.write_manifest(&manifest).unwrap();
2759 assert!(s.get(&graft_key(0x20, 17)).unwrap().is_none(),
2760 "preparing a candidate must not publish it in the live handle");
2761 }
2763 std::fs::write(manifest.with_extension("tmp"), b"torn next candidate").unwrap();
2764 let prepared = PreparedGraft::read_manifest(&manifest).unwrap();
2765 {
2766 let mut resumed = Store::open(d.path(), cfg()).unwrap();
2767 assert!(resumed.get(&graft_key(0x20, 17)).unwrap().is_none());
2768 resumed.publish_existing_candidate(&prepared).unwrap();
2769 assert_eq!(resumed.get(&graft_key(0x20, 17)).unwrap().as_deref(),
2770 Some(&b"graft-17"[..]));
2771 }
2773 {
2774 let mut resumed_again = Store::open(d.path(), cfg()).unwrap();
2775 let generation = resumed_again.generation;
2776 resumed_again.publish_existing_candidate(&prepared).unwrap();
2777 assert_eq!(resumed_again.generation, generation,
2778 "resume-of-resume must not publish another generation");
2779 assert_eq!(resumed_again.scan(&[]).unwrap().count(), 5_000);
2780 }
2781 }
2782
2783 #[test]
2784 fn prepared_graft_manifest_and_generation_are_both_enforced() {
2785 let d = tempfile::tempdir().unwrap();
2786 let scratch = d.path().join("prepared-graft-scratch");
2787 let mut s = Store::create(d.path(), cfg()).unwrap();
2788 seed_graft_base(&mut s);
2789 let mut rows = graft_rows();
2790 rows.sort_by(|a, b| a.0.cmp(&b.0));
2791 let prepared = s.prepare_graft_candidate(
2792 rows.clone().into_iter().map(|(k, v)| Ok((k, v, false))), 2_000,
2793 rows.first().unwrap().0.clone(), rows.last().unwrap().0.clone(), &scratch).unwrap();
2794 let mut damaged = prepared.encode().unwrap();
2795 damaged[20] ^= 0x80;
2796 assert!(PreparedGraft::decode(&damaged).is_err());
2797
2798 let mut wrong_generation = prepared.clone();
2799 wrong_generation.base_generation += 1;
2800 wrong_generation.write_generation += 1;
2801 assert!(s.publish_existing_candidate(&wrong_generation).is_err(),
2802 "candidate pages stamped for one generation must not publish in another");
2803 assert!(s.get(&graft_key(0x20, 17)).unwrap().is_none());
2804 }
2805
2806 #[test]
2810 fn corrupt_packed_range_is_caught_and_never_published() {
2811 let d = tempfile::tempdir().unwrap();
2812 let mut s = Store::create(d.path(), cfg()).unwrap();
2813 seed_graft_base(&mut s);
2814 let old_root = s.root;
2815
2816 let result = s.graft_range_with_before_publish(
2817 graft_rows().into_iter(),
2818 |data, packed| {
2819 use std::io::{Seek, SeekFrom, Write};
2820 let mut file = std::fs::OpenOptions::new().write(true).open(data)?;
2821 file.seek(SeekFrom::Start(packed.root as u64 * PAGE_SIZE as u64 + 80))?;
2822 file.write_all(&[0xa5])?;
2823 file.sync_all()?;
2824 Ok(())
2825 },
2826 );
2827
2828 assert!(matches!(result, Err(crate::Error::Corrupt { .. })),
2829 "a corrupt packed page must be refused, got {result:?}");
2830 assert_eq!(s.root, old_root);
2831 assert_eq!(s.get(&graft_key(0x10, 17)).unwrap().as_deref(), Some(&b"base-10-17"[..]));
2832 assert!(s.get(&graft_key(0x20, 17)).unwrap().is_none());
2833 }
2834
2835 #[test]
2842 fn an_oversized_record_never_reaches_the_log() {
2843 let d = tempfile::tempdir().unwrap();
2844 let mut s = Store::create(d.path(), cfg()).unwrap();
2845 s.put(b"a", b"1").unwrap();
2846 s.commit().unwrap();
2847 let before = s.wal.as_ref().unwrap().end_offset();
2848
2849 let max = crate::page::MAX_RECORD_LEN;
2850 let k = vec![b'k'; max];
2854 assert!(matches!(s.put(&k, b"v"), Err(crate::Error::TooLarge)));
2855 let k = vec![b'k'; max - 1];
2858 assert!(!s.delete(&k).unwrap(),
2859 "a key that long can never have been inserted, so `not found` is the truth");
2860 let half = vec![b'b'; max / 2];
2863 assert!(matches!(
2864 s.put_empty_batch(&[half.clone(), half]),
2865 Err(crate::Error::TooLarge)
2866 ));
2867 assert_eq!(
2868 s.wal.as_ref().unwrap().end_offset(), before,
2869 "neither refusal may append a single byte to the log"
2870 );
2871 }
2872}