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 self.wal
537 .as_mut()
538 .ok_or(crate::Error::ReadOnly)?
539 .cut_to_committed(committed_end)
540 }
541
542 fn apply(&mut self, kind: RecKind, payload: &[u8], wal_offset: u64) -> Result<()> {
543 let corrupt = |why| crate::Error::CorruptWal { offset: wal_offset, why };
544 match kind {
545 RecKind::Put => {
546 let klen = payload.get(..2)
547 .map(|b| u16::from_le_bytes([b[0], b[1]]) as usize)
548 .ok_or_else(|| corrupt("put payload has no key length"))?;
549 let key_end = 2usize.checked_add(klen)
550 .ok_or_else(|| corrupt("put key boundary overflow"))?;
551 let key = payload.get(2..key_end)
552 .ok_or_else(|| corrupt("put key crosses its WAL payload"))?;
553 let val = payload.get(key_end..)
554 .ok_or_else(|| corrupt("put value boundary is invalid"))?;
555 if self.generation == 0 && crate::keys::is_field_aggregate_key(key) {
560 if Meta::is_salvaged(&self.pool)? { return Ok(()); }
561 }
562 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);
563 t.insert(key, val)?;
564 self.root = t.root();
565 }
566 RecKind::PutEmptyBatch => {
567 if payload.len() > crate::page::MAX_RECORD_LEN {
568 return Err(corrupt("put-empty-batch exceeds the WAL frame bound"));
569 }
570 let count = payload.get(..2)
571 .map(|b| u16::from_le_bytes([b[0], b[1]]) as usize)
572 .ok_or_else(|| corrupt("put-empty-batch payload has no key count"))?;
573 if count == 0 || count > PUT_EMPTY_BATCH_MAX_KEYS {
574 return Err(corrupt("put-empty-batch key count is outside its fixed bound"));
575 }
576
577 let mut at = 2usize;
581 for _ in 0..count {
582 let len_end = at.checked_add(2)
583 .ok_or_else(|| corrupt("put-empty-batch length boundary overflow"))?;
584 let len_bytes = payload.get(at..len_end)
585 .ok_or_else(|| corrupt("put-empty-batch key has no length"))?;
586 let key_len = u16::from_le_bytes([len_bytes[0], len_bytes[1]]) as usize;
587 let key_end = len_end.checked_add(key_len)
588 .ok_or_else(|| corrupt("put-empty-batch key boundary overflow"))?;
589 payload.get(len_end..key_end)
590 .ok_or_else(|| corrupt("put-empty-batch key crosses its WAL payload"))?;
591 if 4 + key_len + 12 > crate::page::MAX_RECORD_LEN {
592 return Err(corrupt("put-empty-batch key cannot fit a leaf record"));
593 }
594 at = key_end;
595 }
596 if at != payload.len() {
597 return Err(corrupt("put-empty-batch payload has trailing bytes"));
598 }
599
600 let mut t = BTree::open(&self.pool, self.tree_id, self.root,
601 &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
602 at = 2;
603 for _ in 0..count {
604 let key_len = u16::from_le_bytes([payload[at], payload[at + 1]]) as usize;
605 at += 2;
606 t.insert(&payload[at..at + key_len], &[])?;
607 at += key_len;
608 }
609 self.root = t.root();
610 }
611 RecKind::Delete => {
612 let klen = payload.get(..2)
613 .map(|b| u16::from_le_bytes([b[0], b[1]]) as usize)
614 .ok_or_else(|| corrupt("delete payload has no key length"))?;
615 let key_end = 2usize.checked_add(klen)
616 .ok_or_else(|| corrupt("delete key boundary overflow"))?;
617 if key_end != payload.len() {
618 return Err(corrupt("delete key length does not match its WAL payload"));
619 }
620 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);
621 t.delete(&payload[2..key_end])?;
622 self.root = t.root();
623 }
624 RecKind::DeletePrefix => {
625 if payload.is_empty() {
626 return Err(corrupt("delete-prefix WAL payload is empty"));
627 }
628 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);
629 t.delete_prefix(payload)?;
630 self.root = t.root();
631 }
632 RecKind::Commit => {
633 if !payload.is_empty() {
634 return Err(corrupt("commit WAL payload is not empty"));
635 }
636 }
637 RecKind::PageImage => {
638 return Err(corrupt("page-image WAL records have no recovery implementation"));
639 }
640 }
641 Ok(())
642 }
643
644 fn frame(k: &[u8], v: &[u8]) -> Vec<u8> {
650 let mut b = Vec::with_capacity(2 + k.len() + v.len());
651 b.extend_from_slice(&(k.len() as u16).to_le_bytes());
652 b.extend_from_slice(k);
653 b.extend_from_slice(v);
654 b
655 }
656
657 fn wal_mut(&mut self) -> Result<&mut Wal> {
659 self.wal.as_mut().ok_or(crate::Error::ReadOnly)
660 }
661
662 pub fn put(&mut self, k: &[u8], v: &[u8]) -> Result<()> {
663 #[cfg(feature = "test-support")]
664 self.test_faults.check(TestFaultKind::Write)?;
665 if self.poisoned { return Err(crate::Error::StorePoisoned); }
666 if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
667 let wal_payload_len = 2usize.checked_add(k.len())
668 .and_then(|n| n.checked_add(v.len()))
669 .ok_or(crate::Error::TooLarge)?;
670 if self.resource_limits().is_some_and(|l| v.len() > l.record_bytes as usize) {
671 return Err(Error::ResourceLimit("record exceeds configured maximum"));
672 }
673 if wal_payload_len as u64 > crate::wal::MAX_PAYLOAD_BYTES {
677 return Err(crate::Error::TooLarge);
678 }
679 if k.len() > crate::page::MAX_RECORD_LEN - 16 || v.len() > u32::MAX as usize {
680 return Err(crate::Error::TooLarge);
681 }
682 #[cfg(feature = "write-trace")]
683 let frame_started = crate::write_trace::active().then(std::time::Instant::now);
684 let payload = Self::frame(k, v);
685 #[cfg(feature = "write-trace")]
686 if let Some(started) = frame_started {
687 crate::write_trace::add(crate::write_trace::Field::FrameEncode, started.elapsed());
688 crate::write_trace::value_copy();
689 }
690 if 4 + k.len() + 12 > crate::page::MAX_RECORD_LEN || v.len() > u32::MAX as usize {
697 return Err(crate::Error::TooLarge);
698 }
699 if let Err(e) = self.wal_mut()?.append(RecKind::Put, &payload) {
700 self.poisoned = true;
704 return Err(e);
705 }
706 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);
707 #[cfg(feature = "write-trace")]
708 let btree_started = crate::write_trace::active().then(std::time::Instant::now);
709 let inserted = t.insert(k, v);
710 #[cfg(feature = "write-trace")]
711 if let Some(started) = btree_started {
712 crate::write_trace::add(crate::write_trace::Field::BtreeTotal, started.elapsed());
713 }
714 let root = t.root();
715 drop(t);
716 if let Err(e) = inserted {
717 self.poisoned = true;
718 return Err(e);
719 }
720 self.root = root;
721 Ok(())
722 }
723
724 #[cfg(feature = "write-trace")]
727 pub fn put_profiled(&mut self, k: &[u8], v: &[u8]) -> (Result<()>, crate::write_trace::PutTrace) {
728 let trace_started = crate::write_trace::begin();
729 let result = self.put(k, v);
730 let trace = crate::write_trace::finish(trace_started);
731 (result, trace)
732 }
733
734 pub fn put_empty_batch(&mut self, keys: &[Vec<u8>]) -> Result<()> {
739 if keys.is_empty() { return Ok(()); }
740 #[cfg(feature = "test-support")]
741 self.test_faults.check(TestFaultKind::Write)?;
742 if self.poisoned { return Err(crate::Error::StorePoisoned); }
743 if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
744 if keys.len() > PUT_EMPTY_BATCH_MAX_KEYS { return Err(crate::Error::TooLarge); }
745
746 let payload_len = keys.iter().try_fold(2usize, |total, key| {
747 if 4 + key.len() + 12 > crate::page::MAX_RECORD_LEN || key.len() > u16::MAX as usize {
748 return None;
749 }
750 total.checked_add(2 + key.len())
751 }).ok_or(crate::Error::TooLarge)?;
752 if payload_len > crate::page::MAX_RECORD_LEN { return Err(crate::Error::TooLarge); }
753 let mut payload = Vec::with_capacity(payload_len);
754 payload.extend_from_slice(&(keys.len() as u16).to_le_bytes());
755 for key in keys {
756 payload.extend_from_slice(&(key.len() as u16).to_le_bytes());
757 payload.extend_from_slice(key);
758 }
759 if let Err(e) = self.wal_mut()?.append(RecKind::PutEmptyBatch, &payload) {
760 self.poisoned = true;
761 return Err(e);
762 }
763 let mut t = BTree::open(&self.pool, self.tree_id, self.root,
764 &self.last_leaf, &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
765 let inserted = keys.iter().try_for_each(|key| t.insert(key, &[]));
766 let root = t.root();
767 drop(t);
768 if let Err(e) = inserted {
769 self.poisoned = true;
770 return Err(e);
771 }
772 self.root = root;
773 Ok(())
774 }
775
776 pub fn delete(&mut self, k: &[u8]) -> Result<bool> {
777 #[cfg(feature = "test-support")]
778 self.test_faults.check(TestFaultKind::Write)?;
779 if self.poisoned { return Err(crate::Error::StorePoisoned); }
780 if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
781 if k.len() > crate::page::MAX_RECORD_LEN - 2 { return Ok(false); }
782 let mut payload = (k.len() as u16).to_le_bytes().to_vec();
783 payload.extend_from_slice(k);
784 if payload.len() > crate::page::MAX_RECORD_LEN { return Ok(false); }
793 if let Err(e) = self.wal_mut()?.append(RecKind::Delete, &payload) {
794 self.poisoned = true;
795 return Err(e);
796 }
797 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);
798 let deleted = t.delete(k);
799 let root = t.root();
800 drop(t);
801 let hit = match deleted {
802 Ok(hit) => hit,
803 Err(e) => {
804 self.poisoned = true;
805 return Err(e);
806 }
807 };
808 self.root = root;
809 Ok(hit)
810 }
811
812 pub fn delete_prefix(&mut self, prefix: &[u8]) -> Result<u64> {
818 #[cfg(feature = "test-support")]
819 self.test_faults.check(TestFaultKind::Write)?;
820 if self.poisoned { return Err(crate::Error::StorePoisoned); }
821 if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
822 if prefix.is_empty() || prefix.len() > crate::page::MAX_RECORD_LEN { return Err(crate::Error::TooLarge); }
823 if let Err(e) = self.wal_mut()?.append(RecKind::DeletePrefix, prefix) {
824 self.poisoned = true;
825 return Err(e);
826 }
827 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);
828 let deleted = t.delete_prefix(prefix);
829 let root = t.root();
830 drop(t);
831 let n = match deleted {
832 Ok(n) => n,
833 Err(e) => {
834 self.poisoned = true;
835 return Err(e);
836 }
837 };
838 self.root = root;
839 Ok(n)
840 }
841
842 pub fn get(&self, k: &[u8]) -> Result<Option<Vec<u8>>> {
843 if self.poisoned && self.resource_limits().is_some() {
844 return Err(crate::Error::StorePoisoned);
845 }
846 #[cfg(feature = "test-support")]
847 self.test_faults.check(TestFaultKind::Read)?;
848 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)
849 }
850
851 pub fn scan(&self, from: &[u8]) -> Result<RangeIter<'_>> {
852 if self.poisoned && self.resource_limits().is_some() {
853 return Err(crate::Error::StorePoisoned);
854 }
855 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)
856 }
857
858 pub fn scan_reverse(&self, to: &[u8]) -> Result<crate::btree::ReverseRangeIter<'_>> {
860 if self.poisoned && self.resource_limits().is_some() {
861 return Err(crate::Error::StorePoisoned);
862 }
863 BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf,
864 &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints).range_reverse(to)
865 }
866
867 pub fn commit(&mut self) -> Result<()> {
881 if self.resource_limits().is_some() { return self.checkpoint(); }
882 #[cfg(feature = "test-support")]
883 self.test_faults.check(TestFaultKind::Commit)?;
884 if self.poisoned { return Err(crate::Error::StorePoisoned); }
885 let wal = self.wal.as_mut().ok_or(crate::Error::ReadOnly)?;
886 wal.append(RecKind::Commit, &[])?;
887 wal.flush()?;
891 match self.sync {
892 SyncMode::Full => {
893 #[cfg(test)] self.barriers.push("sync_full");
894 wal.sync_full()?
895 }
896 SyncMode::Normal => {
897 #[cfg(test)] self.barriers.push("sync_data");
898 wal.sync_data()?
899 }
900 SyncMode::Off => {}
901 }
902 let readers = crate::readers::live_generations(&self.dir);
908 self.pool.refresh_reuse(self.generation, readers.as_deref());
909 Ok(())
910 }
911
912 pub fn commit_with_checkpoint(&mut self, wal_bytes: u64, page_bytes: u64) -> Result<bool> {
922 if self.resource_limits().is_some() { self.commit()?; return Ok(true); }
923 if self.poisoned { return Err(crate::Error::StorePoisoned); }
924 let end = self.wal.as_ref().ok_or(crate::Error::ReadOnly)?.end_offset();
925 let publish = (wal_bytes > 0 && end >= wal_bytes)
926 || (page_bytes > 0 && self.pool.epoch_allocated_bytes() >= page_bytes);
927 if !publish { self.commit()?; return Ok(false); }
928 #[cfg(feature = "test-support")]
929 self.test_faults.check(TestFaultKind::Commit)?;
930 let wal = self.wal_mut()?;
931 if let Err(e) = wal.append(RecKind::Commit, &[]).and_then(|_| wal.flush()) {
932 self.poisoned = true;
933 return Err(e);
934 }
935 if let Err(e) = self.checkpoint() {
936 self.poisoned = true;
937 return Err(e);
938 }
939 Ok(true)
940 }
941
942 #[cfg(test)]
943 fn barriers(&self) -> Vec<&'static str> { self.barriers.clone() }
944
945 pub fn dir(&self) -> &Path { &self.dir }
948 pub fn published_root(&self) -> u32 { self.root }
951 pub fn main_tree_id(&self) -> u16 { self.tree_id }
953
954 pub fn sync_full_primitive(&self) -> &'static str {
955 self.wal.as_ref().map_or("none (snapshot reader)", |w| w.sync_full_primitive())
956 }
957
958 pub fn pool_stats(&self) -> crate::pool::PoolStats { self.pool.stats() }
965 pub fn tag_hints(&self) -> &crate::btree::TagHints { &self.tag_hints }
968 pub fn io_stats(&self) -> Option<&crate::io::IoStats> { self.pool.io_stats() }
970 pub fn sweep_steps(&self) -> u64 { self.pool.sweep_steps() }
971
972 fn barrier(&self) -> Barrier { sync_barrier(self.sync) }
979
980 pub fn set_sync(&mut self, s: SyncMode) { self.sync = s; }
983
984 pub fn sync_mode(&self) -> SyncMode { self.sync }
989 pub fn generation(&self) -> u64 { self.generation }
991
992 pub fn checkpoint(&mut self) -> Result<()> {
993 let result = self.checkpoint_inner();
994 if result.is_err() && self.wal.is_some() { self.poisoned = true; }
995 result
996 }
997 fn checkpoint_inner(&mut self) -> Result<()> {
998 if self.poisoned { return Err(crate::Error::StorePoisoned); }
999 if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
1000 let gen = self.generation.checked_add(1).ok_or(crate::Error::Corrupt {
1005 page_no: if self.generation % 2 == 0 { crate::meta::META_PAGE } else { crate::meta::META_PAGE_B },
1006 why: "published generation is exhausted",
1007 })?;
1008 let next_write_generation = gen.checked_add(1).ok_or(crate::Error::Corrupt {
1009 page_no: if self.generation % 2 == 0 { crate::meta::META_PAGE } else { crate::meta::META_PAGE_B },
1010 why: "no generation remains for the next write epoch",
1011 })?;
1012 #[cfg(test)] self.trace.clear();
1017 if let Err(e) = self.pool.flush_all(self.barrier()) {
1024 self.poisoned = true;
1025 return Err(e);
1026 }
1027 #[cfg(test)] { self.trace.push("flush_pages"); self.trace.push("sync_file"); }
1028 Meta {
1034 format_version: self.format_version,
1035 roots: [self.root, 0, 0, 0, 0, 0, 0, 0],
1036 next_lsn: self.wal.as_ref().ok_or(crate::Error::ReadOnly)?.next_lsn(),
1037 generation: gen,
1038 }.write_slot(&self.pool)?;
1039 if let Err(e) = self.pool.flush_all(self.barrier()) {
1040 self.poisoned = true;
1041 return Err(e);
1042 }
1043 self.generation = gen;
1044 self.pool.set_stamp_gen(next_write_generation);
1045 let readers = crate::readers::live_generations(&self.dir);
1049 self.pool.refresh_reuse(gen, readers.as_deref());
1050 let free_bytes = self.pool.export_free(gen);
1051 let _ = crate::verify::persist_checkpoint_freelist(
1052 &self.dir,
1053 &free_bytes,
1054 self.pool.file_ref(),
1055 );
1056 #[cfg(test)] self.trace.push("flip_meta");
1062 self.pool.set_frozen_boundary();
1065 self.wal_mut()?.rotate_published()?;
1068 #[cfg(test)] self.trace.push("rotate_wal");
1069 Ok(())
1070 }
1071
1072 #[cfg(test)]
1073 fn checkpoint_trace(&self) -> Vec<&'static str> { self.trace.clone() }
1074
1075 #[cfg(test)]
1082 fn next_lsn(&self) -> u64 { self.wal.as_ref().unwrap().next_lsn() }
1083
1084 pub fn bulk_load<I>(&mut self, items: I) -> Result<()>
1110 where I: Iterator<Item = (Vec<u8>, Vec<u8>)> {
1111 self.bulk_load_with_before_publish(items, |_, _| Ok(()))
1112 }
1113
1114 fn bulk_load_with_before_publish<I, F>(&mut self, items: I, before_publish: F) -> Result<()>
1117 where
1118 I: Iterator<Item = (Vec<u8>, Vec<u8>)>,
1119 F: FnOnce(&Path, u32) -> Result<()>,
1120 {
1121 self.refuse_external_workspace()?;
1126 if self.poisoned { return Err(crate::Error::StorePoisoned); }
1127 if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
1128 static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1142 let seq = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1143 let tmp = std::env::temp_dir()
1144 .join(format!("kernel-sort-{}-{}", std::process::id(), seq));
1145 let mut s = crate::bulk::ExternalSort::new(&tmp, 64 << 20)?;
1146 let mut expected_rows = 0u64;
1147 for (k, v) in items {
1148 expected_rows = expected_rows.checked_add(1).ok_or(crate::Error::TooLarge)?;
1149 if 4 + k.len() + v.len() > crate::page::MAX_RECORD_LEN {
1154 let (head, crc) = crate::btree::write_overflow(&self.pool, &v)?;
1155 let m = crate::btree::enc_marker(v.len() as u32, head, crc);
1156 s.push_flagged(k, m.to_vec(), true)?;
1157 } else {
1158 s.push(k, v)?;
1159 }
1160 }
1161 let mut runs = s.finish()?;
1162 let root = crate::bulk::pack_tree(&self.pool, self.tree_id, runs.iter()?, 0.9, &tmp)?;
1167 if let Err(e) = self.pool.flush_all(Barrier::None) {
1170 self.poisoned = true;
1171 return Err(e);
1172 }
1173 let data = self.dir.join("data");
1174 before_publish(&data, root)?;
1175 let verified = crate::verify::verify_file(
1176 &data,
1177 self.io_mode,
1178 root,
1179 self.tree_id,
1180 expected_rows,
1181 )?;
1182 debug_assert_eq!(verified.rows, expected_rows);
1183 debug_assert!(verified.pages > 0);
1184 self.root = root;
1185 self.last_leaf.set(None);
1193 self.tag_hints.clear();
1194 self.checkpoint()
1195 }
1196
1197 pub fn graft_range<I>(&mut self, items: I) -> Result<()>
1208 where
1209 I: Iterator<Item = (Vec<u8>, Vec<u8>)>,
1210 {
1211 self.graft_range_with_before_publish(items, |_, _| Ok(()))
1212 }
1213
1214 fn graft_range_with_before_publish<I, F>(
1217 &mut self,
1218 items: I,
1219 before_publish: F,
1220 ) -> Result<()>
1221 where
1222 I: Iterator<Item = (Vec<u8>, Vec<u8>)>,
1223 F: FnOnce(&Path, &crate::bulk::PackedRange) -> Result<()>,
1224 {
1225 self.refuse_external_workspace()?;
1226 if self.poisoned { return Err(crate::Error::StorePoisoned); }
1227 if self.wal.is_none() { return Err(crate::Error::ReadOnly); }
1228
1229 static SEQ: std::sync::atomic::AtomicU64 = std::sync::atomic::AtomicU64::new(0);
1230 let seq = SEQ.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
1231 let tmp = std::env::temp_dir()
1232 .join(format!("kernel-graft-{}-{}", std::process::id(), seq));
1233 let mut sort = crate::bulk::ExternalSort::new(&tmp, 64 << 20)?;
1234 let mut expected_rows = 0u64;
1235 let mut min: Option<Vec<u8>> = None;
1236 let mut max: Option<Vec<u8>> = None;
1237
1238 for (key, value) in items {
1239 expected_rows = expected_rows.checked_add(1).ok_or(crate::Error::TooLarge)?;
1240 if min.as_ref().is_none_or(|current| key.as_slice() < current.as_slice()) {
1241 min = Some(key.clone());
1242 }
1243 if max.as_ref().is_none_or(|current| key.as_slice() > current.as_slice()) {
1244 max = Some(key.clone());
1245 }
1246 if 4 + key.len() + value.len() > crate::page::MAX_RECORD_LEN {
1247 if value.len() > u32::MAX as usize || 4 + key.len() + 12 > crate::page::MAX_RECORD_LEN {
1248 return Err(crate::Error::TooLarge);
1249 }
1250 let (head, crc) = crate::btree::write_overflow(&self.pool, &value)?;
1251 let marker = crate::btree::enc_marker(value.len() as u32, head, crc);
1252 sort.push_flagged(key, marker.to_vec(), true)?;
1253 } else {
1254 sort.push(key, value)?;
1255 }
1256 }
1257 let (Some(min), Some(max)) = (min, max) else {
1258 return Ok(());
1261 };
1262
1263 let mut runs = sort.finish()?;
1264 self.graft_sorted_range_with_before_publish(
1265 runs.iter()?, expected_rows, min, max, &tmp, true, before_publish)
1266 }
1267
1268 pub fn graft_sorted_range<I>(
1277 &mut self,
1278 sorted: I,
1279 expected_rows: u64,
1280 min: Vec<u8>,
1281 max: Vec<u8>,
1282 scratch_dir: &Path,
1283 ) -> Result<()>
1284 where
1285 I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>,
1286 {
1287 self.graft_sorted_range_with_before_publish(
1288 sorted, expected_rows, min, max, scratch_dir, true, |_, _| Ok(()))
1289 }
1290
1291 pub fn graft_sorted_range_deferred<I>(
1297 &mut self,
1298 sorted: I,
1299 expected_rows: u64,
1300 min: Vec<u8>,
1301 max: Vec<u8>,
1302 scratch_dir: &Path,
1303 ) -> Result<()>
1304 where
1305 I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>,
1306 {
1307 self.graft_sorted_range_with_before_publish(
1308 sorted, expected_rows, min, max, scratch_dir, false, |_, _| Ok(()))
1309 }
1310
1311 pub fn prepare_graft_candidate<I>(
1316 &mut self,
1317 sorted: I,
1318 expected_rows: u64,
1319 min: Vec<u8>,
1320 max: Vec<u8>,
1321 scratch_dir: &Path,
1322 ) -> Result<PreparedGraft>
1323 where
1324 I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>,
1325 {
1326 self.refuse_external_workspace()?;
1327 if self.poisoned { return Err(Error::StorePoisoned); }
1328 if self.wal.is_none() { return Err(Error::ReadOnly); }
1329 if expected_rows == 0 || min > max { return Err(Error::TooLarge); }
1330 if let Some(row) = self.scan(&min)?.next() {
1331 let (key, _) = row?;
1332 if key <= max { return Err(Error::RangeNotEmpty); }
1333 }
1334 let tree = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf,
1335 &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
1336 let boundary = tree.plan_graft(&min)?;
1337 drop(tree);
1338 let last_next = boundary.right_page.unwrap_or(boundary.old_next);
1339 let packed = crate::bulk::pack_range(
1340 &self.pool, self.tree_id, sorted, 0.9, scratch_dir, last_next)?;
1341 if packed.rows != expected_rows || packed.min.as_deref() != Some(min.as_slice())
1342 || packed.max.as_deref() != Some(max.as_slice()) {
1343 return Err(Error::Corrupt { page_no: packed.root,
1344 why: "prepared range disagrees with its sorted-stream manifest" });
1345 }
1346 let tree = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf,
1347 &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
1348 let candidate = tree.build_graft_candidate(&boundary, &packed)?;
1349 drop(tree);
1350 if let Err(error) = self.pool.flush_all(Barrier::None) {
1351 self.poisoned = true;
1352 return Err(error);
1353 }
1354 let write_generation = self.generation.checked_add(1).ok_or(Error::Corrupt {
1355 page_no: 0, why: "prepared graft generation is exhausted" })?;
1356 crate::verify::verify_range_file_generation(
1357 &self.dir.join("data"), self.io_mode, candidate.root, self.tree_id,
1358 candidate.rows, &candidate.min, &candidate.max, candidate.last_next,
1359 Some(write_generation))?;
1360 Ok(PreparedGraft {
1361 base_generation: self.generation,
1362 write_generation,
1363 root: candidate.root,
1364 rows: candidate.rows,
1365 min: candidate.min,
1366 max: candidate.max,
1367 last_next: candidate.last_next,
1368 inserted_min: min,
1369 inserted_max: max,
1370 inserted_rows: expected_rows,
1371 })
1372 }
1373
1374 pub fn publish_existing_candidate(&mut self, prepared: &PreparedGraft) -> Result<()> {
1379 self.refuse_external_workspace()?;
1380 if self.poisoned { return Err(Error::StorePoisoned); }
1381 if self.wal.is_none() { return Err(Error::ReadOnly); }
1382 if prepared.write_generation != prepared.base_generation.checked_add(1)
1383 .ok_or(Error::TooLarge)? || prepared.inserted_min > prepared.inserted_max
1384 || prepared.inserted_rows == 0 {
1385 return Err(Error::Corrupt { page_no: prepared.root,
1386 why: "prepared graft has an invalid generation or interval" });
1387 }
1388 let data = self.dir.join("data");
1389 crate::verify::verify_range_file_generation(
1390 &data, self.io_mode, prepared.root, self.tree_id, prepared.rows,
1391 &prepared.min, &prepared.max, prepared.last_next,
1392 Some(prepared.write_generation))?;
1393
1394 if self.generation == prepared.write_generation {
1395 let mut rows = 0u64;
1396 self.scan(&prepared.inserted_min)?.for_each_ref(|key, _| {
1397 if key > prepared.inserted_max.as_slice() { return false; }
1398 rows = rows.saturating_add(1);
1399 true
1400 })?;
1401 if rows != prepared.inserted_rows {
1402 return Err(Error::Corrupt { page_no: prepared.root,
1403 why: "published graft interval disagrees with its manifest" });
1404 }
1405 return Ok(());
1406 }
1407 if self.generation != prepared.base_generation {
1408 return Err(Error::Corrupt { page_no: prepared.root,
1409 why: "prepared graft belongs to a stale base generation" });
1410 }
1411 if let Some(row) = self.scan(&prepared.inserted_min)?.next() {
1412 let (key, _) = row?;
1413 if key <= prepared.inserted_max { return Err(Error::RangeNotEmpty); }
1414 }
1415 let mut tree = BTree::open(&self.pool, self.tree_id, self.root, &self.last_leaf,
1416 &self.fast_path_hits, &self.fast_path_attempts).with_tags(&self.tag_hints);
1417 let boundary = tree.plan_existing_graft(&prepared.inserted_min)?;
1418 let retired = tree.install_graft(&boundary, prepared.root, &prepared.min)?;
1419 self.root = tree.root();
1420 drop(tree);
1421 for page in retired { self.pool.free_page(page)?; }
1422 self.last_leaf.set(None);
1423 self.tag_hints.clear();
1424 self.checkpoint()
1425 }
1426
1427 fn graft_sorted_range_with_before_publish<I, F>(
1428 &mut self,
1429 sorted: I,
1430 expected_rows: u64,
1431 min: Vec<u8>,
1432 max: Vec<u8>,
1433 scratch_dir: &Path,
1434 publish: bool,
1435 before_publish: F,
1436 ) -> Result<()>
1437 where
1438 I: Iterator<Item = Result<(Vec<u8>, Vec<u8>, bool)>>,
1439 F: FnOnce(&Path, &crate::bulk::PackedRange) -> Result<()>,
1440 {
1441 let trace = std::env::var_os("SEKEJAP_LOAD_BREAKDOWN").is_some();
1442 let total_started = trace.then(std::time::Instant::now);
1443 self.refuse_external_workspace()?;
1444 if self.poisoned {
1445 return Err(crate::Error::StorePoisoned);
1446 }
1447 if self.wal.is_none() {
1448 return Err(crate::Error::ReadOnly);
1449 }
1450 if expected_rows == 0 {
1451 return Ok(());
1452 }
1453 if min > max {
1454 return Err(crate::Error::TooLarge);
1455 }
1456
1457 let preflight_started = trace.then(std::time::Instant::now);
1458 if let Some(row) = self.scan(&min)?.next() {
1459 let (key, _) = row?;
1460 if key <= max { return Err(crate::Error::RangeNotEmpty); }
1461 }
1462
1463 let tree = BTree::open(
1464 &self.pool,
1465 self.tree_id,
1466 self.root,
1467 &self.last_leaf,
1468 &self.fast_path_hits,
1469 &self.fast_path_attempts,
1470 ).with_tags(&self.tag_hints);
1471 let boundary = tree.plan_graft(&min)?;
1472 drop(tree);
1473 let last_next = boundary.right_page.unwrap_or(boundary.old_next);
1474 let preflight =
1475 preflight_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1476
1477 let pack_started = trace.then(std::time::Instant::now);
1478 let packed = crate::bulk::pack_range(
1479 &self.pool,
1480 self.tree_id,
1481 sorted,
1482 0.9,
1483 scratch_dir,
1484 last_next,
1485 )?;
1486 let pack = pack_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1487 if packed.rows != expected_rows
1488 || packed.min.as_deref() != Some(min.as_slice())
1489 || packed.max.as_deref() != Some(max.as_slice())
1490 {
1491 return Err(crate::Error::Corrupt {
1492 page_no: packed.root,
1493 why: "packed range disagrees with its sorted-stream manifest",
1494 });
1495 }
1496 if packed.last_leaf >= self.pool.page_count() {
1497 return Err(crate::Error::Corrupt {
1498 page_no: packed.last_leaf,
1499 why: "packed range returned a last leaf outside the data file",
1500 });
1501 }
1502
1503 let tree = BTree::open(
1504 &self.pool,
1505 self.tree_id,
1506 self.root,
1507 &self.last_leaf,
1508 &self.fast_path_hits,
1509 &self.fast_path_attempts,
1510 ).with_tags(&self.tag_hints);
1511 let candidate = tree.build_graft_candidate(&boundary, &packed)?;
1512 drop(tree);
1513
1514 let flush_started = trace.then(std::time::Instant::now);
1517 if let Err(error) = self.pool.flush_all(Barrier::None) {
1518 self.poisoned = true;
1519 return Err(error);
1520 }
1521 let flush = flush_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1522 let data = self.dir.join("data");
1523 before_publish(&data, &packed)?;
1524 let verify_started = trace.then(std::time::Instant::now);
1525 let verified = crate::verify::verify_range_file(
1526 &data,
1527 self.io_mode,
1528 candidate.root,
1529 self.tree_id,
1530 candidate.rows,
1531 &candidate.min,
1532 &candidate.max,
1533 candidate.last_next,
1534 )?;
1535 let verify = verify_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1536 debug_assert_eq!(verified.rows, candidate.rows);
1537 debug_assert!(verified.pages > 0);
1538
1539 let install_started = trace.then(std::time::Instant::now);
1544 let mut tree = BTree::open(
1545 &self.pool,
1546 self.tree_id,
1547 self.root,
1548 &self.last_leaf,
1549 &self.fast_path_hits,
1550 &self.fast_path_attempts,
1551 ).with_tags(&self.tag_hints);
1552 let retired = tree.install_graft(&boundary, candidate.root, &candidate.min)?;
1553 self.root = tree.root();
1554 drop(tree);
1555 for page in retired { self.pool.free_page(page)?; }
1556 self.last_leaf.set(None);
1557 self.tag_hints.clear();
1558 let install =
1559 install_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1560 let publish_started = trace.then(std::time::Instant::now);
1561 let result = if publish { self.checkpoint() } else { Ok(()) };
1562 let publication =
1563 publish_started.map_or(std::time::Duration::ZERO, |started| started.elapsed());
1564 if trace {
1565 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={}",
1566 expected_rows, verified.pages, verified.pages as u64 * crate::page::PAGE_SIZE as u64,
1567 preflight.as_secs_f64(), pack.as_secs_f64(), flush.as_secs_f64(),
1568 verify.as_secs_f64(), install.as_secs_f64(), publication.as_secs_f64(),
1569 total_started.unwrap().elapsed().as_secs_f64(), !publish);
1570 }
1571 result
1572 }
1573}
1574
1575#[cfg(test)]
1583mod tests {
1584 use super::*;
1585
1586 fn cfg() -> Config { Config { budget_bytes: 32 << 20, io: IoMode::Buffered, sync: SyncMode::Full } }
1587
1588 #[test]
1594 fn a_store_checkpoint_preserves_the_file_own_format_version() {
1595 let d = tempfile::tempdir().unwrap();
1596 let other = if crate::meta::FORMAT_VERSION == 2 { 1 } else { 2 };
1597 {
1598 let s = Store::create(d.path(), cfg()).unwrap();
1599 let meta = Meta::read_latest(&s.pool).unwrap();
1600 Meta {
1601 format_version: other,
1602 roots: meta.roots,
1603 next_lsn: meta.next_lsn,
1604 generation: meta.generation,
1605 }
1606 .write(&s.pool)
1607 .unwrap();
1608 s.pool.flush_all(crate::io::Barrier::Data).unwrap();
1609 }
1610 let mut s = Store::open(d.path(), cfg()).unwrap_or_else(|e| {
1611 panic!("supported superblock version {other} must open in this build: {e}")
1612 });
1613 s.put(b"k", b"v").unwrap();
1614 s.commit().unwrap();
1615 s.checkpoint().unwrap();
1616 drop(s);
1617 let s = Store::open(d.path(), cfg()).unwrap();
1618 let got = Meta::read_latest(&s.pool).unwrap().format_version & !crate::meta::LIMITED;
1619 assert_eq!(
1620 got, other,
1621 "checkpoint must write the file's format version, not the build's FORMAT_VERSION"
1622 );
1623 assert_eq!(s.get(b"k").unwrap().as_deref(), Some(&b"v"[..]));
1624 }
1625
1626 #[test]
1627 fn byte_policy_publishes_exact_rows_without_a_redundant_wal_barrier() {
1628 let d = tempfile::tempdir().unwrap();
1629 let mut s = Store::create(d.path(), cfg()).unwrap();
1630 s.put(b"old", b"durable").unwrap(); s.commit().unwrap(); s.checkpoint().unwrap();
1631 let reader = Store::open_snapshot(d.path(), cfg()).unwrap();
1632 s.barriers.clear();
1633 s.put(b"new", &[7; 9000]).unwrap();
1634 let before = s.pool_stats().sync_full_calls;
1635 assert!(s.commit_with_checkpoint(u64::MAX, 4096).unwrap());
1636 assert_eq!(s.pool_stats().sync_full_calls - before, 2);
1637 assert!(s.barriers.is_empty(), "publication already supplies the durability barriers");
1638 assert_eq!(s.wal.as_ref().unwrap().end_offset(), 0);
1639 assert_eq!(reader.get(b"new").unwrap(), None);
1640 drop(s);
1641 let s = Store::open(d.path(), cfg()).unwrap();
1642 assert_eq!(s.get(b"old").unwrap().as_deref(), Some(&b"durable"[..]));
1643 assert_eq!(s.get(b"new").unwrap(), Some(vec![7; 9000]));
1644 }
1645
1646 #[test]
1647 fn byte_policy_keeps_wal_and_old_snapshot_after_data_barrier_failure() {
1648 let d = tempfile::tempdir().unwrap();
1649 let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1650 let fio = Arc::new(FailingBarrier { inner: real, fail: false.into() });
1651 let mut s = Store::create_on(d.path(), cfg(), fio.clone()).unwrap();
1652 s.put(b"old", b"committed").unwrap(); s.commit().unwrap(); s.checkpoint().unwrap();
1653 let reader = Store::open_snapshot(d.path(), cfg()).unwrap();
1654 s.put(b"new", &[8; 9000]).unwrap();
1655 fio.fail.store(true, std::sync::atomic::Ordering::Relaxed);
1656 assert!(s.commit_with_checkpoint(1,1).is_err());
1657 assert!(s.wal.as_ref().unwrap().end_offset() > 0);
1658 assert!(matches!(s.commit_with_checkpoint(1,1), Err(crate::Error::StorePoisoned)));
1659 assert_eq!(reader.get(b"old").unwrap().as_deref(), Some(&b"committed"[..]));
1660 assert_eq!(reader.get(b"new").unwrap(), None);
1661 fio.fail.store(false, std::sync::atomic::Ordering::Relaxed);
1662 drop(s); drop(fio);
1663 let s = Store::open(d.path(), cfg()).unwrap();
1664 assert_eq!(s.get(b"old").unwrap().as_deref(), Some(&b"committed"[..]));
1665 assert_eq!(s.get(b"new").unwrap(), Some(vec![8; 9000]));
1666 }
1667
1668 #[test]
1669 fn constrained_commit_failure_preserves_published_state_and_snapshot() {
1670 for barrier in [false, true] {
1671 let d = tempfile::tempdir().unwrap();
1672 let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1673 let fio = Arc::new(FailingBarrier { inner: real, fail: false.into() });
1674 let limits = crate::limits::ResourceLimits { data_bytes: 1 << 20, wal_bytes: 64 << 10,
1675 tracked_pages: 256, readers: 2, record_bytes: 16000, recovery_bytes: 65536 };
1676 let mut s = Store::build_on_limited(d.path(), cfg(), true, fio.clone(), IoMode::Buffered, Some(limits)).unwrap();
1677 s.put(b"old", b"durable").unwrap(); s.commit().unwrap();
1678 let reader = Store::open_snapshot(d.path(), cfg()).unwrap();
1679 s.put(b"new", &[8;9000]).unwrap();
1680 if barrier { fio.fail.store(true, std::sync::atomic::Ordering::Relaxed); }
1681 if barrier {
1682 assert!(s.commit().is_err());
1683 assert!(matches!(s.checkpoint(), Err(Error::StorePoisoned)));
1684 } assert_eq!(reader.get(b"old").unwrap(), Some(b"durable".to_vec()));
1686 assert_eq!(reader.get(b"new").unwrap(), None);
1687 fio.fail.store(false, std::sync::atomic::Ordering::Relaxed);
1688 drop(s); drop(fio);
1689 let s = Store::open(d.path(), cfg()).unwrap();
1690 assert_eq!(s.get(b"old").unwrap(), Some(b"durable".to_vec()));
1691 assert_eq!(s.get(b"new").unwrap(), None);
1692 }
1693 }
1694
1695 #[test]
1696 fn constrained_commit_data_write_failures_do_not_publish_partial_rows() {
1697 for at in 0..3 {
1698 let d = tempfile::tempdir().unwrap();
1699 let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1700 let fio = Arc::new(FailingWrite { inner: real, armed: std::sync::Mutex::new(None) });
1701 let limits = crate::limits::ResourceLimits { data_bytes: 1 << 20, wal_bytes: 64 << 10,
1702 tracked_pages: 256, readers: 2, record_bytes: 16000, recovery_bytes: 65536 };
1703 let mut s = Store::build_on_limited(d.path(), cfg(), true, fio.clone(), IoMode::Buffered, Some(limits)).unwrap();
1704 s.put(b"old", b"durable").unwrap(); s.commit().unwrap();
1705 let reader = Store::open_snapshot(d.path(), cfg()).unwrap();
1706 s.put(b"new", &[8;12000]).unwrap();
1707 *fio.armed.lock().unwrap() = Some(at);
1708 assert!(s.commit().is_err());
1709 assert!(matches!(s.put(b"later", b"no"), Err(Error::StorePoisoned)));
1710 assert_eq!(reader.get(b"old").unwrap(), Some(b"durable".to_vec()));
1711 *fio.armed.lock().unwrap() = None;
1712 drop(s); drop(fio);
1713 let s = Store::open(d.path(), cfg()).unwrap();
1714 assert_eq!(s.get(b"old").unwrap(), Some(b"durable".to_vec()));
1715 assert_eq!(s.get(b"new").unwrap(), None);
1716 }
1717 }
1718
1719 #[test]
1727 fn a_snapshot_is_registered_before_its_generation_can_be_recycled() {
1728 let d = tempfile::tempdir().unwrap();
1729 let tiny = Config { budget_bytes: 1 << 16, io: IoMode::Buffered, sync: SyncMode::Off };
1730 let mut writer = Store::create(d.path(), tiny).unwrap();
1731 for i in 0..2_000u64 {
1732 writer.put(&i.to_be_bytes(), format!("generation one row {i}").as_bytes()).unwrap();
1733 }
1734 writer.commit().unwrap();
1735 writer.checkpoint().unwrap();
1736
1737 let (selected_tx, selected_rx) = std::sync::mpsc::channel();
1738 let (churned_tx, churned_rx) = std::sync::mpsc::channel();
1739 let result = std::thread::scope(|scope| {
1740 let dir = d.path();
1741 let reader = scope.spawn(move || {
1742 let snapshot = Store::open_snapshot_with_after_meta(dir, tiny, |generation| {
1743 selected_tx.send(generation).unwrap();
1744 churned_rx.recv().unwrap();
1745 Ok(())
1746 })?;
1747 let iter = snapshot.scan(&[])?;
1748 iter.collect::<Result<Vec<_>>>()
1749 });
1750 let writer_thread = scope.spawn(move || {
1751 assert_eq!(selected_rx.recv().unwrap(), 1, "fixture must stop after selecting generation 1");
1752 for round in 0..4u64 {
1753 for i in 0..2_000u64 {
1754 writer.put(&i.to_be_bytes(),
1755 format!("writer round {round} row {i} is different").as_bytes()).unwrap();
1756 }
1757 writer.commit().unwrap();
1758 writer.checkpoint().unwrap();
1759 }
1760 churned_tx.send(()).unwrap();
1761 });
1762 writer_thread.join().unwrap();
1763 reader.join().unwrap()
1764 });
1765
1766 let rows = result.expect("a snapshot must not reach recycled pages while it is opening");
1767 assert_eq!(rows.len(), 2_000);
1768 for (i, (key, value)) in rows.iter().enumerate() {
1769 assert_eq!(key.as_slice(), &(i as u64).to_be_bytes());
1770 assert_eq!(value.as_slice(), format!("generation one row {i}").as_bytes(),
1771 "the opening snapshot observed a recycled page at row {i}");
1772 }
1773 }
1774
1775 struct CrashDirectory {
1776 inner: Box<dyn FileIo>,
1777 directory_durable: std::sync::atomic::AtomicBool,
1778 }
1779
1780 impl FileIo for CrashDirectory {
1781 fn requires_alignment(&self) -> bool { self.inner.requires_alignment() }
1782 fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> { self.inner.read_at(buf, off) }
1783 fn write_at(&self, buf: &[u8], off: u64) -> Result<()> { self.inner.write_at(buf, off) }
1784 fn sync_data(&self) -> Result<()> { self.inner.sync_data() }
1785 fn sync_full(&self) -> Result<()> { self.inner.sync_full() }
1786 fn sync_full_primitive(&self) -> &'static str { self.inner.sync_full_primitive() }
1787 fn sync_dir(&self) -> Result<()> {
1788 self.inner.sync_dir()?;
1789 self.directory_durable.store(true, std::sync::atomic::Ordering::Release);
1790 Ok(())
1791 }
1792 fn len(&self) -> Result<u64> { self.inner.len() }
1793 fn set_len(&self, n: u64) -> Result<()> { self.inner.set_len(n) }
1794 }
1795
1796 #[test]
1801 fn a_commit_before_the_first_checkpoint_survives_a_directory_crash() {
1802 let d = tempfile::tempdir().unwrap();
1803 let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1804 let fio = Arc::new(CrashDirectory {
1805 inner: real,
1806 directory_durable: std::sync::atomic::AtomicBool::new(false),
1807 });
1808 {
1809 let mut s = Store::create_on(d.path(), cfg(), fio.clone()).unwrap();
1810 s.put(b"acknowledged", b"must survive power loss").unwrap();
1811 s.commit().unwrap();
1812 }
1813
1814 if !fio.directory_durable.load(std::sync::atomic::Ordering::Acquire) {
1815 std::fs::remove_file(d.path().join("data")).unwrap();
1816 std::fs::remove_file(d.path().join("wal")).unwrap();
1817 }
1818 let reopened = Store::open(d.path(), cfg()).unwrap();
1819 assert_eq!(reopened.get(b"acknowledged").unwrap().as_deref(),
1820 Some(&b"must survive power loss"[..]),
1821 "creation must make file names durable before any commit can be acknowledged");
1822 }
1823
1824 #[test]
1825 fn recreated_wal_name_survives_a_commit_before_checkpoint() {
1826 let d = tempfile::tempdir().unwrap();
1827 {
1828 let mut s = Store::create(d.path(), cfg()).unwrap();
1829 s.put(b"old", b"published").unwrap();
1830 s.checkpoint().unwrap();
1831 }
1832 std::fs::remove_file(d.path().join("wal")).unwrap();
1833 let (real, mode) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
1834 let fio = Arc::new(CrashDirectory {
1835 inner: real,
1836 directory_durable: std::sync::atomic::AtomicBool::new(false),
1837 });
1838 {
1839 let mut s = Store::build_on(d.path(), cfg(), false, fio.clone(), mode).unwrap();
1840 s.put(b"new", b"acknowledged").unwrap();
1841 s.commit().unwrap();
1842 }
1843 if !fio.directory_durable.load(std::sync::atomic::Ordering::Acquire) {
1844 std::fs::remove_file(d.path().join("wal")).unwrap();
1845 }
1846 drop(fio);
1847 let s = Store::open(d.path(), cfg()).unwrap();
1848 assert_eq!(s.get(b"old").unwrap().as_deref(), Some(&b"published"[..]));
1849 assert_eq!(s.get(b"new").unwrap().as_deref(), Some(&b"acknowledged"[..]));
1850 }
1851
1852 #[test]
1853 fn a_maximum_generation_read_from_disk_is_refused_not_wrapped() {
1854 let d = tempfile::tempdir().unwrap();
1855 {
1856 let s = Store::create(d.path(), cfg()).unwrap();
1857 Meta {
1858 format_version: crate::meta::FORMAT_VERSION,
1859 roots: [s.root, 0, 0, 0, 0, 0, 0, 0],
1860 next_lsn: s.wal.as_ref().unwrap().next_lsn(),
1861 generation: u64::MAX,
1862 }.write_slot(&s.pool).unwrap();
1863 s.pool.flush_all(Barrier::None).unwrap();
1864 }
1865
1866 assert!(matches!(Store::open(d.path(), cfg()), Err(crate::Error::Corrupt { .. })),
1867 "the generation after u64::MAX does not exist and must not become zero");
1868 }
1869
1870 #[test]
1871 fn a_checkpoint_refuses_when_no_later_page_generation_exists() {
1872 let d = tempfile::tempdir().unwrap();
1873 let mut s = Store::create(d.path(), cfg()).unwrap();
1874 s.generation = u64::MAX - 1;
1875 s.put(b"pending", b"kept in the log").unwrap();
1876 s.commit().unwrap();
1877
1878 assert!(matches!(s.checkpoint(), Err(crate::Error::Corrupt { .. })),
1879 "publishing the final generation would leave the next epoch wrapping to zero");
1880 assert!(std::fs::metadata(d.path().join("wal")).unwrap().len() > 0,
1881 "refusing exhaustion must preserve the committed log");
1882 }
1883
1884 #[test]
1885 fn a_checkpoint_rotates_the_log_only_after_the_pages_are_durable() {
1886 let d = tempfile::tempdir().unwrap();
1887 let mut s = Store::create(d.path(), cfg()).unwrap();
1888 for i in 0..2000u64 { s.put(&i.to_be_bytes(), b"v").unwrap(); }
1889 s.commit().unwrap();
1890 s.checkpoint().unwrap();
1891 let order = s.checkpoint_trace();
1892 assert_eq!(order, vec!["flush_pages", "sync_file", "flip_meta", "rotate_wal"],
1897 "checkpoint order must be: data durable, then flip, then drop the log");
1898 assert!(s.get(&1999u64.to_be_bytes()).unwrap().is_some());
1899 }
1900
1901 #[test]
1911 fn a_rotation_followed_by_a_reopen_does_not_reissue_lsn_1() {
1912 let d = tempfile::tempdir().unwrap();
1913 let lsn_at_checkpoint = {
1914 let mut s = Store::create(d.path(), cfg()).unwrap();
1915 for i in 0..500u64 { s.put(&i.to_be_bytes(), b"v").unwrap(); }
1916 s.commit().unwrap();
1917 s.checkpoint().unwrap(); s.next_lsn()
1919 };
1920 assert!(lsn_at_checkpoint > 1, "sanity: many records were appended before the rotation");
1921
1922 let s2 = Store::open(d.path(), cfg()).unwrap();
1923 assert_eq!(
1924 s2.next_lsn(), lsn_at_checkpoint,
1925 "a reopen after rotation must not renumber LSNs from 1"
1926 );
1927 }
1928
1929 #[test]
1935 fn the_three_durability_modes_issue_different_barriers() {
1936 for (mode, want) in [
1937 (SyncMode::Full, vec!["sync_full"]),
1938 (SyncMode::Normal, vec!["sync_data"]),
1939 (SyncMode::Off, vec![]),
1940 ] {
1941 let d = tempfile::tempdir().unwrap();
1942 let cfg = Config { budget_bytes: 16 << 20, io: IoMode::Buffered, sync: mode };
1943 let mut s = Store::create(d.path(), cfg).unwrap();
1944 s.put(b"k", b"v").unwrap();
1945 s.commit().unwrap();
1946 assert_eq!(s.barriers(), want, "{mode:?} issued the wrong barrier");
1947 }
1948 }
1949
1950 fn armed(s: &Store) -> (u64, u64) {
1965 (s.fast_path_hits.get() + s.tag_hints.hits(),
1966 s.fast_path_attempts.get() + s.tag_hints.attempts())
1967 }
1968
1969 #[test]
1985 fn store_put_ascending_uses_fast_path() {
1986 let d = tempfile::tempdir().unwrap();
1987 let mut s = Store::create(d.path(), cfg()).unwrap();
1988 let n = 100u64;
1989 for i in 0..n { s.put(&i.to_be_bytes(), b"v").unwrap(); }
1990 assert_eq!(
1991 armed(&s).0, n - 1,
1992 "every Store::put but the first must hit the append fast path"
1993 );
1994 }
1995
1996 #[test]
2003 fn store_put_survives_bulk_load() {
2004 let d = tempfile::tempdir().unwrap();
2005 let mut s = Store::create(d.path(), cfg()).unwrap();
2006 let n = 1_000u64;
2007 let items = (0..n).map(|i| (i.to_be_bytes().to_vec(), b"bulk".to_vec()));
2008 s.bulk_load(items).unwrap();
2009 assert_eq!(s.last_leaf.get(), None, "bulk_load must clear a hint it just made meaningless");
2010
2011 let tail_key = n.to_be_bytes();
2013 s.put(&tail_key, b"tail").unwrap();
2014
2015 assert_eq!(
2016 s.get(&tail_key).unwrap().as_deref(), Some(&b"tail"[..]),
2017 "a put right after bulk_load must be found"
2018 );
2019
2020 let scanned: Vec<Vec<u8>> = s.scan(&[]).unwrap().map(|r| r.unwrap().0).collect();
2021 let mut sorted = scanned.clone();
2022 sorted.sort();
2023 assert_eq!(scanned.len(), n as usize + 1);
2024 assert_eq!(scanned, sorted, "a full scan after bulk_load + put must stay in sorted order");
2025 }
2026
2027 #[test]
2034 fn store_random_order_matches_sequential() {
2035 let n = 20_000u64;
2036 let scatter = |i: u64| i.wrapping_mul(0x9E37_79B9_7F4A_7C15);
2037
2038 let d1 = tempfile::tempdir().unwrap();
2039 let mut s1 = Store::create(d1.path(), cfg()).unwrap();
2040 for i in 0..n { s1.put(&i.to_be_bytes(), &i.to_le_bytes()).unwrap(); }
2041
2042 let d2 = tempfile::tempdir().unwrap();
2043 let mut s2 = Store::create(d2.path(), cfg()).unwrap();
2044 let mut order: Vec<u64> = (0..n).collect();
2045 order.sort_by_key(|&i| scatter(i));
2046 for &i in &order { s2.put(&i.to_be_bytes(), &i.to_le_bytes()).unwrap(); }
2047
2048 let seq1: Vec<(Vec<u8>, Vec<u8>)> = s1.scan(&[]).unwrap().map(|r| r.unwrap()).collect();
2049 let seq2: Vec<(Vec<u8>, Vec<u8>)> = s2.scan(&[]).unwrap().map(|r| r.unwrap()).collect();
2050 assert_eq!(seq1.len(), n as usize);
2051 assert_eq!(
2052 seq1, seq2,
2053 "ascending vs scattered insertion order through Store::put must produce identical scans"
2054 );
2055 }
2056
2057 #[test]
2084 fn random_order_disarms_fast_path() {
2085 let d = tempfile::tempdir().unwrap();
2086 let mut s = Store::create(d.path(), cfg()).unwrap();
2087 let n = 20_000u64;
2088 let scatter = |i: u64| i.wrapping_mul(0x9E37_79B9_7F4A_7C15);
2089 let payload = vec![b'x'; 200]; for i in 0..n { s.put(&scatter(i).to_be_bytes(), &payload).unwrap(); }
2091 #[cfg(feature = "sqlite-balance")]
2095 assert!(s.fast_path_attempts.get() <= 117);
2096 assert!(
2104 s.tag_hints.attempts() <= 2_000,
2105 "a scattered workload attempted the per-keyspace fast path {} times in {n} \
2106 inserts; arming belongs to appends, not to every descent",
2107 s.tag_hints.attempts()
2108 );
2109 #[cfg(not(feature = "sqlite-balance"))]
2110 assert_eq!(
2111 s.fast_path_attempts.get(), 117,
2112 "a disarming hint must attempt the fast path a number of times bounded by \
2113 leaf capacity and log(n), not by n -- a bound proportional to n would not \
2114 have caught the pre-Task-18 defect this test exists for"
2115 );
2116 }
2117
2118 #[test]
2128 fn ascending_keeps_fast_path_armed() {
2129 let d = tempfile::tempdir().unwrap();
2130 let mut s = Store::create(d.path(), cfg()).unwrap();
2131 let n = 100u64;
2132 for i in 0..n { s.put(&i.to_be_bytes(), b"v").unwrap(); }
2133 assert_eq!(
2134 armed(&s).0, n - 1,
2135 "an unbroken ascending run must still hit the fast path on every insert but the first"
2136 );
2137 }
2138
2139 #[test]
2172 fn mixed_workload_rearms() {
2173 let d = tempfile::tempdir().unwrap();
2174 let mut s = Store::create(d.path(), cfg()).unwrap();
2175 let m = 300u64;
2176 let k = 30u64;
2177
2178 for i in 1..=m { s.put(&i.to_be_bytes(), b"v").unwrap(); }
2179 assert_eq!(armed(&s).1, m - 1, "one attempt per insert but the first");
2180 assert_eq!(
2181 armed(&s).0, m - 2,
2182 "one miss expected: the insert that fills the leaf and forces the one split \
2183 this run crosses"
2184 );
2185
2186 s.put(&0u64.to_be_bytes(), b"v").unwrap();
2187 assert_eq!(armed(&s).1, m, "the out-of-order key is one more attempt");
2188 assert_eq!(armed(&s).0, m - 2, "the out-of-order key must not hit");
2189 assert_eq!(
2190 s.last_leaf.get(), None,
2191 "landing on a non-rightmost leaf must leave the hint disarmed, not re-armed \
2192 to the wrong leaf"
2193 );
2194
2195 let (hits, attempts) = armed(&s);
2196 for i in (m + 1)..(m + 1 + k) { s.put(&i.to_be_bytes(), b"v").unwrap(); }
2197 let (hits, attempts) = (armed(&s).0 - hits, armed(&s).1 - attempts);
2198 assert_eq!(
2199 attempts, hits,
2200 "an unbroken ascending run must hit every probe it makes"
2201 );
2202 assert!(
2210 k - hits <= 1,
2211 "{} of the {k} ascending inserts after the disarm paid a descent",
2212 k - hits
2213 );
2214 }
2215
2216 struct FailingBarrier { inner: Box<dyn FileIo>, fail: std::sync::atomic::AtomicBool }
2225 impl FailingBarrier {
2226 fn failing(&self) -> bool { self.fail.load(std::sync::atomic::Ordering::Relaxed) }
2227 fn eio() -> crate::Error {
2228 std::io::Error::other("injected barrier failure").into()
2229 }
2230 }
2231 impl FileIo for FailingBarrier {
2232 fn requires_alignment(&self) -> bool { self.inner.requires_alignment() }
2233 fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> { self.inner.read_at(buf, off) }
2234 fn write_at(&self, buf: &[u8], off: u64) -> Result<()> { self.inner.write_at(buf, off) }
2235 fn sync_data(&self) -> Result<()> {
2236 if self.failing() { return Err(Self::eio()); }
2237 self.inner.sync_data()
2238 }
2239 fn sync_full(&self) -> Result<()> {
2240 if self.failing() { return Err(Self::eio()); }
2241 self.inner.sync_full()
2242 }
2243 fn sync_full_primitive(&self) -> &'static str { self.inner.sync_full_primitive() }
2244 fn sync_dir(&self) -> Result<()> { self.inner.sync_dir() }
2245 fn len(&self) -> Result<u64> { self.inner.len() }
2246 fn set_len(&self, n: u64) -> Result<()> { self.inner.set_len(n) }
2247 }
2248
2249 struct FailingWrite {
2253 inner: Box<dyn FileIo>,
2254 armed: std::sync::Mutex<Option<usize>>,
2255 }
2256 impl FailingWrite {
2257 fn arm(&self, successful_writes: usize) {
2258 *self.armed.lock().unwrap() = Some(successful_writes);
2259 }
2260 }
2261 impl FileIo for FailingWrite {
2262 fn requires_alignment(&self) -> bool { self.inner.requires_alignment() }
2263 fn read_at(&self, buf: &mut [u8], off: u64) -> Result<()> { self.inner.read_at(buf, off) }
2264 fn write_at(&self, buf: &[u8], off: u64) -> Result<()> {
2265 let mut armed = self.armed.lock().unwrap();
2266 if let Some(remaining) = armed.as_mut() {
2267 if *remaining == 0 {
2268 *armed = None;
2269 return Err(std::io::Error::other(
2270 "injected page-write failure during a tree mutation",
2271 ).into());
2272 }
2273 *remaining -= 1;
2274 }
2275 self.inner.write_at(buf, off)
2276 }
2277 fn sync_data(&self) -> Result<()> { self.inner.sync_data() }
2278 fn sync_full(&self) -> Result<()> { self.inner.sync_full() }
2279 fn sync_full_primitive(&self) -> &'static str { self.inner.sync_full_primitive() }
2280 fn sync_dir(&self) -> Result<()> { self.inner.sync_dir() }
2281 fn len(&self) -> Result<u64> { self.inner.len() }
2282 fn set_len(&self, n: u64) -> Result<()> { self.inner.set_len(n) }
2283 }
2284
2285 #[test]
2286 fn deletion_merge_eviction_failures_cannot_publish_partial_changes() {
2287 for fail_after in 0..12 {
2288 let d=tempfile::tempdir().unwrap();
2289 let cfg=Config{budget_bytes:64<<10,io:IoMode::Buffered,sync:SyncMode::Full};
2290 let (real,_)=open_file(&d.path().join("data"),IoMode::Buffered).unwrap();
2291 let fio=Arc::new(FailingWrite{inner:real,armed:std::sync::Mutex::new(None)});
2292 let mut s=Store::create_on(d.path(),cfg,fio.clone()).unwrap();
2293 for i in 0u64..1024 {s.put(&i.to_be_bytes(),&vec![7;240]).unwrap();}
2294 for i in 0u64..1024 {if i%15>=6 {s.delete(&i.to_be_bytes()).unwrap();}}
2297 s.checkpoint().unwrap();
2298 let old=Store::open_snapshot(d.path(),cfg).unwrap();
2299 fio.arm(fail_after);let mut failed=false;
2300 for i in 0u64..1024 {
2301 if let Err(e)=s.delete(&((i*71)%1024).to_be_bytes()) {
2302 assert!(matches!(e,crate::Error::Io(_)));failed=true;break;
2303 }
2304 }
2305 assert!(failed,"fault {fail_after} was not reached");
2306 assert!(s.commit().is_err());assert!(s.checkpoint().is_err());
2307 for i in 0u64..1024 {assert_eq!(old.get(&i.to_be_bytes()).unwrap(),(i%15<6).then(||vec![7;240]));}
2308 drop(old);drop(s);
2309 let reopened=Store::open(d.path(),cfg).unwrap();
2310 for i in 0u64..1024 {assert_eq!(reopened.get(&i.to_be_bytes()).unwrap(),(i%15<6).then(||vec![7;240]));}
2311 assert_eq!(crate::verify::verify_published_tree(&d.path().join("data"),IoMode::Buffered,reopened.root,1).unwrap().0,412);
2312 }
2313 }
2314
2315 #[test]
2321 fn a_partial_tree_mutation_poisons_and_reopen_replays_the_retained_log() {
2322 let d = tempfile::tempdir().unwrap();
2323 let tiny = Config {
2324 budget_bytes: 16 * crate::page::PAGE_SIZE,
2325 io: IoMode::Buffered,
2326 sync: SyncMode::Off,
2327 };
2328 let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
2329 let fio = Arc::new(FailingWrite {
2330 inner: real,
2331 armed: std::sync::Mutex::new(None),
2332 });
2333 let mut s = Store::create_on(d.path(), tiny, fio.clone()).unwrap();
2334
2335 let value = vec![b'v'; 256];
2336 let record_bytes = 4 + 8 + value.len() + 4;
2337 let mut rows = 0u64;
2338 loop {
2339 let room = {
2340 let r = s.pool.get(s.root).unwrap();
2341 crate::page::PageRef::open_resident(&r, s.root).unwrap().free_space()
2342 };
2343 if room < record_bytes { break; }
2344 s.put(&(rows * 2).to_be_bytes(), &value).unwrap();
2345 s.commit().unwrap();
2346 rows += 1;
2347 }
2348 assert!(rows > 4, "fixture must have committed rows to move across the split");
2349
2350 while s.pool.page_count() < 16 {
2355 drop(s.pool.allocate().unwrap());
2356 }
2357 {
2358 let root_pin = s.pool.get(s.root).unwrap();
2359 for _ in 0..15 { drop(s.pool.allocate().unwrap()); }
2360 drop(root_pin);
2361 }
2362
2363 let wal_before = std::fs::metadata(d.path().join("wal")).unwrap().len();
2364 fio.arm(1);
2365 let failed = s.put(&1u64.to_be_bytes(), &value);
2366 assert!(matches!(failed, Err(crate::Error::Io(_))),
2367 "the injected eviction failure must escape the mutating insert, got {failed:?}");
2368 assert_eq!(s.get(&((rows - 1) * 2).to_be_bytes()).unwrap(), None,
2369 "fixture must prove the insert failed after committed keys moved off the root");
2370
2371 let commit = s.commit();
2375 let checkpoint = s.checkpoint();
2376 let wal_after = std::fs::metadata(d.path().join("wal")).unwrap().len();
2377 assert!(matches!(commit, Err(crate::Error::StorePoisoned)),
2378 "a partial tree mutation must make commit refuse, got {commit:?}");
2379 assert!(matches!(checkpoint, Err(crate::Error::StorePoisoned)),
2380 "a partial tree mutation must make checkpoint refuse, got {checkpoint:?}");
2381 assert_eq!(wal_after, wal_before,
2382 "a poisoned store must retain the committed recovery log byte-for-byte");
2383
2384 drop(s);
2385 drop(fio);
2386 let reopened = Store::open(d.path(), tiny)
2387 .expect("reopen must rebuild a poisoned handle from the retained log");
2388 for i in 0..rows {
2389 assert_eq!(reopened.get(&(i * 2).to_be_bytes()).unwrap().as_deref(), Some(value.as_slice()));
2390 }
2391 assert_eq!(reopened.get(&1u64.to_be_bytes()).unwrap(), None,
2392 "the failed, uncommitted insert must not be replayed");
2393 }
2394
2395 #[test]
2399 fn a_pre_mutation_validation_error_does_not_poison_the_store() {
2400 let d = tempfile::tempdir().unwrap();
2401 let mut s = Store::create(d.path(), cfg()).unwrap();
2402 let oversized_key = vec![b'k'; crate::page::MAX_RECORD_LEN];
2403 assert!(matches!(s.put(&oversized_key, b"v"), Err(crate::Error::TooLarge)));
2404 s.put(b"usable", b"still").unwrap();
2405 s.commit().unwrap();
2406 s.checkpoint().unwrap();
2407 assert_eq!(s.get(b"usable").unwrap().as_deref(), Some(&b"still"[..]));
2408 }
2409
2410 #[test]
2419 fn a_failed_checkpoint_barrier_poisons_every_writer() {
2420 let d = tempfile::tempdir().unwrap();
2421 let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
2422 let fio = Arc::new(FailingBarrier { inner: real, fail: false.into() });
2423 let mut s = Store::create_on(d.path(), cfg(), fio.clone()).unwrap();
2424 s.put(b"a", b"1").unwrap();
2425 s.commit().unwrap();
2426 s.checkpoint().expect("sanity: the checkpoint works while the disk does");
2427
2428 s.put(b"b", b"2").unwrap();
2429 s.commit().unwrap();
2430 fio.fail.store(true, std::sync::atomic::Ordering::Relaxed);
2431
2432 match s.checkpoint() {
2433 Err(crate::Error::Io(_)) => {}
2434 other => panic!("a failing barrier must surface, got {other:?}"),
2435 }
2436 fio.fail.store(false, std::sync::atomic::Ordering::Relaxed); assert!(matches!(s.put(b"c", b"3"), Err(crate::Error::StorePoisoned)));
2439 assert!(matches!(s.delete(b"a"), Err(crate::Error::StorePoisoned)));
2440 assert!(matches!(s.commit(), Err(crate::Error::StorePoisoned)));
2441 assert!(matches!(s.checkpoint(), Err(crate::Error::StorePoisoned)),
2442 "a retried checkpoint would flush nothing (those frames are marked clean \
2443 already) and then rotate the log away");
2444 let items = vec![(b"z".to_vec(), b"9".to_vec())].into_iter();
2445 assert!(matches!(s.bulk_load(items), Err(crate::Error::StorePoisoned)),
2446 "bulk_load replaces the whole tree AND discards the log");
2447 assert_eq!(s.get(b"a").unwrap().as_deref(), Some(&b"1"[..]),
2453 "a refused bulk_load must not have replaced the tree");
2454 assert_eq!(s.get(b"z").unwrap(), None);
2455 }
2456
2457 #[test]
2464 fn reopening_clears_poisoning_and_the_committed_data_is_still_there() {
2465 let d = tempfile::tempdir().unwrap();
2466 let (real, _) = open_file(&d.path().join("data"), IoMode::Buffered).unwrap();
2467 let fio = Arc::new(FailingBarrier { inner: real, fail: false.into() });
2468 {
2469 let mut s = Store::create_on(d.path(), cfg(), fio.clone()).unwrap();
2470 s.put(b"a", b"1").unwrap();
2471 s.commit().unwrap();
2472 fio.fail.store(true, std::sync::atomic::Ordering::Relaxed);
2473 assert!(s.checkpoint().is_err());
2474 assert!(matches!(s.put(b"b", b"2"), Err(crate::Error::StorePoisoned)));
2475 }
2476 fio.fail.store(false, std::sync::atomic::Ordering::Relaxed);
2477
2478 let mut s2 = Store::open(d.path(), cfg()).expect("a poisoned store must not poison the DIRECTORY");
2479 assert_eq!(s2.get(b"a").unwrap().as_deref(), Some(&b"1"[..]),
2480 "the committed row is in the log, which the failed checkpoint never rotated");
2481 s2.put(b"b", b"2").unwrap();
2482 s2.commit().unwrap();
2483 }
2484
2485 #[test]
2489 fn malformed_wal_payloads_are_bounded_and_do_not_mutate_the_tree() {
2490 let d = tempfile::tempdir().unwrap();
2491 let mut s = Store::create(d.path(), cfg()).unwrap();
2492 s.put(b"kept", b"value").unwrap();
2493 s.commit().unwrap();
2494 let root = s.root;
2495
2496 let malformed: &[(RecKind, &[u8])] = &[
2497 (RecKind::Put, &[]),
2498 (RecKind::Put, &[4, 0, b'a']),
2499 (RecKind::Delete, &[]),
2500 (RecKind::Delete, &[2, 0, b'a']),
2501 (RecKind::Delete, &[0, 0, b'x']),
2502 (RecKind::DeletePrefix, &[]),
2503 (RecKind::PutEmptyBatch, &[]),
2504 (RecKind::PutEmptyBatch, &[1, 0, 4, 0, b'a']),
2505 (RecKind::PutEmptyBatch, &[65, 0]),
2506 (RecKind::PutEmptyBatch, &[1, 0, 1, 0, b'a', b'x']),
2507 (RecKind::Commit, b"not empty"),
2508 (RecKind::PageImage, b"unsupported"),
2509 ];
2510 for (n, &(kind, payload)) in malformed.iter().enumerate() {
2511 assert!(matches!(
2512 s.apply(kind, payload, n as u64),
2513 Err(crate::Error::CorruptWal { offset, .. }) if offset == n as u64
2514 ));
2515 assert_eq!(s.root, root, "malformed frame {n} changed the root");
2516 assert_eq!(s.get(b"kept").unwrap().as_deref(), Some(&b"value"[..]),
2517 "malformed frame {n} changed existing data");
2518 }
2519 }
2520
2521 #[test]
2522 fn bounded_empty_key_batch_replays_as_one_committed_wal_unit() {
2523 let d = tempfile::tempdir().unwrap();
2524 {
2525 let mut s = Store::create(d.path(), cfg()).unwrap();
2526 s.put_empty_batch(&[b"alpha".to_vec(), b"beta".to_vec()]).unwrap();
2527 s.commit().unwrap();
2528 }
2529 let s = Store::open(d.path(), cfg()).unwrap();
2530 assert_eq!(s.get(b"alpha").unwrap().as_deref(), Some(&b""[..]));
2531 assert_eq!(s.get(b"beta").unwrap().as_deref(), Some(&b""[..]));
2532 }
2533
2534 #[test]
2535 fn a_corrupted_packed_tree_never_becomes_authoritative() {
2536 let d = tempfile::tempdir().unwrap();
2537 let mut s = Store::create(d.path(), cfg()).unwrap();
2538 s.put(b"old", b"authoritative").unwrap();
2539 s.commit().unwrap();
2540 s.checkpoint().unwrap();
2541 let old_root = s.root;
2542
2543 let result = s.bulk_load_with_before_publish(
2544 (0..2_000u64).map(|i| (i.to_be_bytes().to_vec(), b"new".to_vec())),
2545 |data, root| {
2546 use std::io::{Seek, SeekFrom, Write};
2547 let mut file = std::fs::OpenOptions::new().write(true).open(data)?;
2548 file.seek(SeekFrom::Start(root as u64 * PAGE_SIZE as u64))?;
2549 file.write_all(&[0u8])?;
2550 file.sync_all()?;
2551 Ok(())
2552 },
2553 );
2554
2555 assert!(result.is_err(), "a packed tree corrupted before publication must be refused");
2556 assert_eq!(s.root, old_root, "the old root must stay authoritative after refusal");
2557 assert_eq!(s.get(b"old").unwrap().as_deref(), Some(&b"authoritative"[..]));
2558 }
2559
2560 fn graft_key(space: u8, i: u32) -> Vec<u8> {
2561 let mut key = vec![space];
2562 key.extend_from_slice(&i.to_be_bytes());
2563 key
2564 }
2565
2566 fn seed_graft_base(store: &mut Store) -> Vec<(Vec<u8>, Vec<u8>)> {
2567 let mut rows = Vec::new();
2568 for space in [0x10, 0x30] {
2569 for i in 0..1_500u32 {
2570 let row = (graft_key(space, i), format!("base-{space:02x}-{i}").into_bytes());
2571 store.put(&row.0, &row.1).unwrap();
2572 rows.push(row);
2573 }
2574 }
2575 store.commit().unwrap();
2576 store.checkpoint().unwrap();
2577 rows.sort_by(|a, b| a.0.cmp(&b.0));
2578 rows
2579 }
2580
2581 fn graft_rows() -> Vec<(Vec<u8>, Vec<u8>)> {
2582 (0..2_000u32)
2583 .rev()
2584 .map(|i| (graft_key(0x20, i), format!("graft-{i}").into_bytes()))
2585 .collect()
2586 }
2587
2588 fn collect_from(store: &Store, from: &[u8]) -> Vec<(Vec<u8>, Vec<u8>)> {
2589 store.scan(from).unwrap().map(Result::unwrap).collect()
2590 }
2591
2592 fn collect_below(store: &Store, to: &[u8]) -> Vec<(Vec<u8>, Vec<u8>)> {
2593 let mut rows = Vec::new();
2594 store.scan_reverse(to).unwrap().for_each_ref(|key, value| {
2595 rows.push((key.to_vec(), value.to_vec()));
2596 true
2597 }).unwrap();
2598 rows
2599 }
2600
2601 #[test]
2605 fn graft_preserves_every_preexisting_key() {
2606 let d = tempfile::tempdir().unwrap();
2607 let mut s = Store::create(d.path(), cfg()).unwrap();
2608 let base = seed_graft_base(&mut s);
2609 let pinned = Store::open_snapshot(d.path(), cfg()).unwrap();
2610
2611 s.graft_range(graft_rows().into_iter()).unwrap();
2612
2613 for (key, value) in &base {
2614 assert_eq!(s.get(key).unwrap().as_deref(), Some(value.as_slice()),
2615 "graft lost pre-existing key {key:?}");
2616 }
2617 assert_eq!(s.scan(&[]).unwrap().count(), base.len() + 2_000);
2618 assert!(pinned.get(&graft_key(0x20, 17)).unwrap().is_none(),
2619 "a reader pinned before publication must remain on the old generation");
2620 assert_eq!(pinned.scan(&[]).unwrap().count(), base.len());
2621 let fresh = Store::open_snapshot(d.path(), cfg()).unwrap();
2622 assert_eq!(fresh.get(&graft_key(0x20, 17)).unwrap().as_deref(), Some(&b"graft-17"[..]));
2623 }
2624
2625 #[test]
2629 fn graft_and_individual_inserts_answer_every_kernel_query_identically() {
2630 let dg = tempfile::tempdir().unwrap();
2631 let di = tempfile::tempdir().unwrap();
2632 let mut grafted = Store::create(dg.path(), cfg()).unwrap();
2633 let mut inserted = Store::create(di.path(), cfg()).unwrap();
2634 seed_graft_base(&mut grafted);
2635 seed_graft_base(&mut inserted);
2636 let rows = graft_rows();
2637
2638 grafted.graft_range(rows.clone().into_iter()).unwrap();
2639 for (key, value) in &rows { inserted.put(key, value).unwrap(); }
2640 inserted.commit().unwrap();
2641 inserted.checkpoint().unwrap();
2642
2643 for (key, _) in seed_query_keys(&rows) {
2644 assert_eq!(grafted.get(&key).unwrap(), inserted.get(&key).unwrap(),
2645 "point query disagreed at {key:?}");
2646 }
2647 for from in [vec![], graft_key(0x10, 777), graft_key(0x20, 0),
2648 graft_key(0x20, 999), graft_key(0x30, 0), vec![0xff]] {
2649 assert_eq!(collect_from(&grafted, &from), collect_from(&inserted, &from),
2650 "forward range disagreed from {from:?}");
2651 }
2652 for to in [graft_key(0x10, 0), graft_key(0x20, 0), graft_key(0x20, 999),
2653 graft_key(0x30, 0), vec![0xff]] {
2654 assert_eq!(collect_below(&grafted, &to), collect_below(&inserted, &to),
2655 "reverse range disagreed below {to:?}");
2656 }
2657
2658 let live_key = graft_key(0x20, 2_500);
2662 grafted.put(&live_key, b"later-live-write").unwrap();
2663 inserted.put(&live_key, b"later-live-write").unwrap();
2664 grafted.commit().unwrap();
2665 inserted.commit().unwrap();
2666 grafted.checkpoint().unwrap();
2667 inserted.checkpoint().unwrap();
2668 assert_eq!(collect_from(&grafted, &graft_key(0x20, 1_900)),
2669 collect_from(&inserted, &graft_key(0x20, 1_900)));
2670 }
2671
2672 #[test]
2673 fn graft_refuses_a_nonempty_range_without_changing_the_tree() {
2674 let d = tempfile::tempdir().unwrap();
2675 let mut s = Store::create(d.path(), cfg()).unwrap();
2676 let base = seed_graft_base(&mut s);
2677 let old_root = s.root;
2678 let rows = vec![
2679 (graft_key(0x0f, 0), b"before".to_vec()),
2680 (graft_key(0x10, 10), b"overlap".to_vec()),
2681 ];
2682 assert!(matches!(s.graft_range(rows.into_iter()), Err(crate::Error::RangeNotEmpty)));
2683 assert_eq!(s.root, old_root);
2684 assert_eq!(s.scan(&[]).unwrap().count(), base.len());
2685 for (key, value) in base {
2686 assert_eq!(s.get(&key).unwrap().as_deref(), Some(value.as_slice()));
2687 }
2688 }
2689
2690 #[test]
2691 fn graft_into_an_empty_tree_becomes_the_tree() {
2692 let d = tempfile::tempdir().unwrap();
2693 let mut s = Store::create(d.path(), cfg()).unwrap();
2694 let rows = graft_rows();
2695 s.graft_range(rows.clone().into_iter()).unwrap();
2696 assert_eq!(s.scan(&[]).unwrap().count(), rows.len());
2697 for (key, value) in rows.iter().step_by(97) {
2698 assert_eq!(s.get(key).unwrap().as_deref(), Some(value.as_slice()));
2699 }
2700 }
2701
2702 fn seed_query_keys(rows: &[(Vec<u8>, Vec<u8>)]) -> Vec<(Vec<u8>, Vec<u8>)> {
2703 let mut keys = rows.to_vec();
2704 keys.push((graft_key(0x10, 0), Vec::new()));
2705 keys.push((graft_key(0x10, 1_499), Vec::new()));
2706 keys.push((graft_key(0x30, 0), Vec::new()));
2707 keys.push((graft_key(0x30, 1_499), Vec::new()));
2708 keys.push((graft_key(0x20, 2_001), Vec::new()));
2709 keys
2710 }
2711
2712 #[test]
2716 fn graft_crash_boundary_is_old_before_and_durable_after() {
2717 let before_dir = tempfile::tempdir().unwrap();
2718 {
2719 let mut s = Store::create(before_dir.path(), cfg()).unwrap();
2720 seed_graft_base(&mut s);
2721 let old_root = s.root;
2722 let result = s.graft_range_with_before_publish(
2723 graft_rows().into_iter(),
2724 |_, _| Err(std::io::Error::other("crash before graft").into()),
2725 );
2726 assert!(result.is_err());
2727 assert_eq!(s.root, old_root, "a pre-publication failure changed the live root");
2728 }
2729 let before = Store::open(before_dir.path(), cfg()).unwrap();
2730 assert!(before.get(&graft_key(0x20, 17)).unwrap().is_none());
2731 assert_eq!(before.scan(&[]).unwrap().count(), 3_000);
2732
2733 let after_dir = tempfile::tempdir().unwrap();
2734 {
2735 let mut s = Store::create(after_dir.path(), cfg()).unwrap();
2736 seed_graft_base(&mut s);
2737 s.graft_range(graft_rows().into_iter()).unwrap();
2738 }
2740 let after = Store::open(after_dir.path(), cfg()).unwrap();
2741 assert_eq!(after.get(&graft_key(0x20, 17)).unwrap().as_deref(), Some(&b"graft-17"[..]));
2742 assert_eq!(after.scan(&[]).unwrap().count(), 5_000);
2743 }
2744
2745 #[test]
2749 fn prepared_graft_publishes_after_reopen_and_resume_of_resume_is_idempotent() {
2750 let d = tempfile::tempdir().unwrap();
2751 let scratch = d.path().join("prepared-graft-scratch");
2752 let manifest = d.path().join("prepared-graft");
2753 {
2754 let mut s = Store::create(d.path(), cfg()).unwrap();
2755 seed_graft_base(&mut s);
2756 let mut rows = graft_rows();
2757 rows.sort_by(|a, b| a.0.cmp(&b.0));
2758 let min = rows.first().unwrap().0.clone();
2759 let max = rows.last().unwrap().0.clone();
2760 let prepared = s.prepare_graft_candidate(
2761 rows.into_iter().map(|(k, v)| Ok((k, v, false))),
2762 2_000, min, max, &scratch).unwrap();
2763 prepared.write_manifest(&manifest).unwrap();
2764 assert!(s.get(&graft_key(0x20, 17)).unwrap().is_none(),
2765 "preparing a candidate must not publish it in the live handle");
2766 }
2768 std::fs::write(manifest.with_extension("tmp"), b"torn next candidate").unwrap();
2769 let prepared = PreparedGraft::read_manifest(&manifest).unwrap();
2770 {
2771 let mut resumed = Store::open(d.path(), cfg()).unwrap();
2772 assert!(resumed.get(&graft_key(0x20, 17)).unwrap().is_none());
2773 resumed.publish_existing_candidate(&prepared).unwrap();
2774 assert_eq!(resumed.get(&graft_key(0x20, 17)).unwrap().as_deref(),
2775 Some(&b"graft-17"[..]));
2776 }
2778 {
2779 let mut resumed_again = Store::open(d.path(), cfg()).unwrap();
2780 let generation = resumed_again.generation;
2781 resumed_again.publish_existing_candidate(&prepared).unwrap();
2782 assert_eq!(resumed_again.generation, generation,
2783 "resume-of-resume must not publish another generation");
2784 assert_eq!(resumed_again.scan(&[]).unwrap().count(), 5_000);
2785 }
2786 }
2787
2788 #[test]
2789 fn prepared_graft_manifest_and_generation_are_both_enforced() {
2790 let d = tempfile::tempdir().unwrap();
2791 let scratch = d.path().join("prepared-graft-scratch");
2792 let mut s = Store::create(d.path(), cfg()).unwrap();
2793 seed_graft_base(&mut s);
2794 let mut rows = graft_rows();
2795 rows.sort_by(|a, b| a.0.cmp(&b.0));
2796 let prepared = s.prepare_graft_candidate(
2797 rows.clone().into_iter().map(|(k, v)| Ok((k, v, false))), 2_000,
2798 rows.first().unwrap().0.clone(), rows.last().unwrap().0.clone(), &scratch).unwrap();
2799 let mut damaged = prepared.encode().unwrap();
2800 damaged[20] ^= 0x80;
2801 assert!(PreparedGraft::decode(&damaged).is_err());
2802
2803 let mut wrong_generation = prepared.clone();
2804 wrong_generation.base_generation += 1;
2805 wrong_generation.write_generation += 1;
2806 assert!(s.publish_existing_candidate(&wrong_generation).is_err(),
2807 "candidate pages stamped for one generation must not publish in another");
2808 assert!(s.get(&graft_key(0x20, 17)).unwrap().is_none());
2809 }
2810
2811 #[test]
2815 fn corrupt_packed_range_is_caught_and_never_published() {
2816 let d = tempfile::tempdir().unwrap();
2817 let mut s = Store::create(d.path(), cfg()).unwrap();
2818 seed_graft_base(&mut s);
2819 let old_root = s.root;
2820
2821 let result = s.graft_range_with_before_publish(
2822 graft_rows().into_iter(),
2823 |data, packed| {
2824 use std::io::{Seek, SeekFrom, Write};
2825 let mut file = std::fs::OpenOptions::new().write(true).open(data)?;
2826 file.seek(SeekFrom::Start(packed.root as u64 * PAGE_SIZE as u64 + 80))?;
2827 file.write_all(&[0xa5])?;
2828 file.sync_all()?;
2829 Ok(())
2830 },
2831 );
2832
2833 assert!(matches!(result, Err(crate::Error::Corrupt { .. })),
2834 "a corrupt packed page must be refused, got {result:?}");
2835 assert_eq!(s.root, old_root);
2836 assert_eq!(s.get(&graft_key(0x10, 17)).unwrap().as_deref(), Some(&b"base-10-17"[..]));
2837 assert!(s.get(&graft_key(0x20, 17)).unwrap().is_none());
2838 }
2839
2840 #[test]
2847 fn an_oversized_record_never_reaches_the_log() {
2848 let d = tempfile::tempdir().unwrap();
2849 let mut s = Store::create(d.path(), cfg()).unwrap();
2850 s.put(b"a", b"1").unwrap();
2851 s.commit().unwrap();
2852 let before = s.wal.as_ref().unwrap().end_offset();
2853
2854 let max = crate::page::MAX_RECORD_LEN;
2855 let k = vec![b'k'; max];
2859 assert!(matches!(s.put(&k, b"v"), Err(crate::Error::TooLarge)));
2860 let k = vec![b'k'; max - 1];
2863 assert!(!s.delete(&k).unwrap(),
2864 "a key that long can never have been inserted, so `not found` is the truth");
2865 let half = vec![b'b'; max / 2];
2868 assert!(matches!(
2869 s.put_empty_batch(&[half.clone(), half]),
2870 Err(crate::Error::TooLarge)
2871 ));
2872 assert_eq!(
2873 s.wal.as_ref().unwrap().end_offset(), before,
2874 "neither refusal may append a single byte to the log"
2875 );
2876 }
2877}