1use std::mem::size_of;
2use std::ops::Bound;
3use std::sync::Arc;
4
5use crate::Key;
6
7use crate::batch::{Applied, BatchWrite};
8use crate::byte_view::ByteView;
9use crate::cache::{BlockCache, BlockKey};
10use crate::compaction::{CompactionIndex, compact_shard};
11use crate::config::Config;
12use crate::disk_loc::DiskLoc;
13use crate::engine::Engine;
14use crate::error::{DbError, DbResult};
15use crate::hook::{NoHook, WriteHook};
16use crate::io::aligned_buf::AlignedBuf;
17use crate::recovery::recover_var_tree;
18use crate::shard::ShardInner;
19use crate::skiplist::node::{SkipNode, VarNode, random_height};
20use crate::skiplist::{InsertResult, SkipList};
21use crate::sync::MutexGuard;
22use crate::value_cache::{ValueCache, ValueKey};
23
24const MAX_STALE_RETRIES: usize = 3;
25
26enum UpdateOutcome {
28 Missing,
30 Unchanged(ByteView),
32 Updated { old: ByteView, new: ByteView },
34}
35
36pub struct VarTree<K: Key, H: WriteHook<K> = NoHook> {
77 index: SkipList<VarNode<K>>,
78 engine: Engine,
79 block_cache: BlockCache,
80 value_cache: ValueCache,
81 compaction_threshold: f64,
82 shard_prefix_bits: usize,
83 reversed: bool,
84 hook: H,
85}
86
87impl<K: Key> VarTree<K> {
88 #[doc = include_str!("../docs/exclusive-access.md")]
92 pub fn open(path: impl AsRef<std::path::Path>, config: Config) -> DbResult<Self> {
93 Self::open_hooked(path, config, NoHook)
94 }
95}
96
97impl<K: Key, H: WriteHook<K>> VarTree<K, H> {
98 #[doc = include_str!("../docs/exclusive-access.md")]
101 pub fn open_hooked(
102 path: impl AsRef<std::path::Path>,
103 config: Config,
104 hook: H,
105 ) -> DbResult<Self> {
106 let config = config.with_resolved_hints(true);
107 if config.var_cache_warn_needed() {
108 tracing::warn!(
109 "Var collection opened with no cache (block_cache=0 and value_cache=0): every get hits disk"
110 );
111 }
112 if config.app_cache_without_direct_io() {
113 tracing::warn!(
114 "Var collection has an app cache enabled (block_cache/value_cache) but direct_io is off: \
115 values are cached twice (app cache + OS page cache). Enable direct_io to avoid redundant \
116 RAM — exception: keeping the app cache for plaintext caching under encryption."
117 );
118 }
119 Self::open_inner(path, config, hook)
120 }
121
122 fn open_inner(path: impl AsRef<std::path::Path>, config: Config, hook: H) -> DbResult<Self> {
123 let compaction_threshold = config.compaction_threshold;
124 let shard_prefix_bits = config.shard_prefix_bits;
125 let reversed = config.reversed;
126 let block_cache = BlockCache::new(&config.block_cache);
127 let value_cache = ValueCache::new(&config.value_cache);
128 let engine = Engine::open(path, config)?;
129
130 let tree = Self {
131 index: SkipList::new(reversed),
132 engine,
133 block_cache,
134 value_cache,
135 compaction_threshold,
136 shard_prefix_bits,
137 reversed,
138 hook,
139 };
140
141 let shard_dirs = tree.engine.shard_dirs();
143 let shard_dir_refs = Engine::shard_dir_refs(&shard_dirs);
144 let shard_ids = tree.engine.shard_ids();
145
146 let hints = tree.engine.hints();
147 let outcome = recover_var_tree::<K>(
148 &shard_dir_refs,
149 &shard_ids,
150 tree.index(),
151 hints,
152 #[cfg(feature = "encryption")]
153 tree.engine.shard_ciphers(),
154 )?;
155 for tail in &outcome.active_tails {
156 tree.engine.shards()[tail.shard_idx].apply_recovery_tail(tail)?;
157 }
158 for (shard_idx, dead) in outcome.shard_dead_bytes {
159 tree.engine.shards()[shard_idx].install_dead_bytes(dead);
160 }
161 #[cfg(feature = "replication")]
162 for (shard_idx, max_gsn) in &outcome.shard_max_gsns {
163 tree.engine.shards()[*shard_idx].set_durable_recovered_gsn(*max_gsn);
164 }
165 let max_gsn = outcome.max_gsn;
166
167 tree.engine
168 .gsn()
169 .fetch_max(max_gsn + 1, std::sync::atomic::Ordering::Relaxed);
170 if hints {
171 for shard in tree.engine.shards().iter() {
172 shard.set_key_len(size_of::<K>());
173 }
174 }
175 tracing::info!(
176 key_size = size_of::<K>(),
177 entries = tree.len(),
178 "var_tree recovered"
179 );
180
181 Ok(tree)
182 }
183
184 pub fn clean_shutdown(&self) -> DbResult<()> {
186 if self.engine.hints() {
187 self.sync_hints()?;
188 }
189 self.engine.flush()
190 }
191
192 pub fn close(self) -> DbResult<()> {
194 self.clean_shutdown()
195 }
196
197 pub fn flush_buffers(&self) -> DbResult<()> {
199 self.engine.flush_buffers()
200 }
201
202 pub(crate) fn flush_durable(&self) -> DbResult<()> {
205 self.engine.flush()
206 }
207
208 pub fn config(&self) -> &Config {
210 self.engine.config()
211 }
212}
213
214impl<K: Key, H: WriteHook<K>> CompactionIndex<K> for VarTree<K, H> {
215 fn update_if_match(&self, key: &K, old_loc: DiskLoc, new_loc: DiskLoc) -> bool {
216 let guard = self.index.collector().enter();
217 if let Some(node) = self.index.get(key.as_bytes(), &guard)
218 && node.read_loc() == old_loc
219 {
220 node.write_loc(new_loc);
221 return true;
222 }
223 false
224 }
225
226 fn invalidate_blocks(&self, shard_id: u8, file_id: u32, total_bytes: u64) {
227 self.block_cache
228 .invalidate_file(shard_id, file_id, total_bytes);
229 }
230
231 fn contains_key(&self, key: &K) -> bool {
232 self.contains(key)
233 }
234
235 fn is_live(&self, _shard_id: u8, key: &K, loc: DiskLoc) -> bool {
236 let guard = self.index.collector().enter();
237 self.index
238 .get(key.as_bytes(), &guard)
239 .is_some_and(|node| node.read_loc() == loc)
240 }
241}
242
243impl<K: Key, H: WriteHook<K>> VarTree<K, H> {
244 pub fn compact(&self) -> DbResult<usize> {
246 let mut total_compacted = 0;
247 for shard in self.engine.shards().iter() {
248 total_compacted += compact_shard(shard, self, self.compaction_threshold)?;
249 }
250 Ok(total_compacted)
251 }
252
253 pub fn get(&self, key: &K) -> Option<ByteView> {
258 metrics::counter!("armdb.ops", "op" => "get", "tree" => "var_tree").increment(1);
259 #[cfg(feature = "hot-path-tracing")]
260 tracing::trace!("var_tree.get");
261 let guard = self.index.collector().enter();
262 let node = match self.index.get(key.as_bytes(), &guard) {
263 Some(n) => n,
264 None => {
265 #[cfg(feature = "hot-path-tracing")]
266 tracing::error!(
267 "VarTree get error: index.get returned None for key {:?}",
268 key.as_bytes()
269 );
270 return None;
271 }
272 };
273 self.read_value_cached(node, &guard)
274 }
275
276 pub fn get_or_err(&self, key: &K) -> DbResult<ByteView> {
279 let guard = self.index.collector().enter();
280 let node = self
281 .index
282 .get(key.as_bytes(), &guard)
283 .ok_or(DbError::KeyNotFound)?;
284 self.read_value_cached_result(node, &guard)
285 }
286
287 pub fn try_get(&self, key: &K) -> DbResult<Option<ByteView>> {
290 let guard = self.index.collector().enter();
291 match self.index.get(key.as_bytes(), &guard) {
292 Some(node) => self.read_value_cached_result(node, &guard).map(Some),
293 None => Ok(None),
294 }
295 }
296
297 pub fn try_first(&self) -> DbResult<Option<(K, ByteView)>> {
300 let guard = self.index.collector().enter();
301 let mut ptr = crate::skiplist::strip_mark(unsafe {
302 (*self.index.head_ptr())
303 .tower(0)
304 .load(std::sync::atomic::Ordering::Acquire)
305 });
306 while !ptr.is_null() {
307 let node = unsafe { &*ptr };
308 if !node.is_marked() {
309 return self
310 .read_value_cached_result(node, &guard)
311 .map(|v| Some((node.key, v)));
312 }
313 ptr = crate::skiplist::strip_mark(
314 node.tower(0).load(std::sync::atomic::Ordering::Acquire),
315 );
316 }
317 Ok(None)
318 }
319
320 pub fn try_last(&self) -> DbResult<Option<(K, ByteView)>> {
322 let guard = self.index.collector().enter();
323 let mut ptr = self.index.find_last(&guard);
324 while !ptr.is_null() {
325 let node = unsafe { &*ptr };
326 if !node.is_marked() {
327 return self
328 .read_value_cached_result(node, &guard)
329 .map(|v| Some((node.key, v)));
330 }
331 ptr = self.index.find_last_lt(node.key_bytes(), &guard);
332 }
333 Ok(None)
334 }
335
336 pub fn put(&self, key: &K, value: &[u8]) -> DbResult<bool> {
342 metrics::counter!("armdb.ops", "op" => "put", "tree" => "var_tree").increment(1);
343 #[cfg(feature = "hot-path-tracing")]
344 tracing::trace!("var_tree.put");
345 let shard_id = self.shard_for(key);
346 let mut inner = self.engine.shards()[shard_id].lock();
347 let guard = self.index.collector().enter();
348 let old_value = if H::NEEDS_OLD_VALUE {
349 if let Some(node) = self.index.get(key.as_bytes(), &guard) {
350 let disk = node.read_loc();
351 Some(self.read_value_locked_result(&disk, shard_id as u8, &inner)?)
352 } else {
353 None
354 }
355 } else {
356 None
357 };
358 let existed = self.put_locked(shard_id, &mut inner, &guard, key, value)?;
359 drop(inner);
360 self.hook.on_write(key, old_value.as_deref(), Some(value));
361 Ok(existed)
362 }
363
364 pub fn insert(&self, key: &K, value: &[u8]) -> DbResult<()> {
367 metrics::counter!("armdb.ops", "op" => "insert", "tree" => "var_tree").increment(1);
368 #[cfg(feature = "hot-path-tracing")]
369 tracing::trace!("var_tree.insert");
370 let shard_id = self.shard_for(key);
371 let mut inner = self.engine.shards()[shard_id].lock();
372 let guard = self.index.collector().enter();
373 self.insert_locked(shard_id, &mut inner, &guard, key, value)?;
374 drop(inner);
375 self.hook.on_write(key, None, Some(value));
376 Ok(())
377 }
378
379 pub fn delete(&self, key: &K) -> DbResult<bool> {
381 metrics::counter!("armdb.ops", "op" => "delete", "tree" => "var_tree").increment(1);
382 #[cfg(feature = "hot-path-tracing")]
383 tracing::trace!("var_tree.delete");
384 let shard_id = self.shard_for(key);
385 let mut inner = self.engine.shards()[shard_id].lock();
386 let guard = self.index.collector().enter();
387 let old_value = if H::NEEDS_OLD_VALUE {
388 if let Some(node) = self.index.get(key.as_bytes(), &guard) {
389 let disk = node.read_loc();
390 Some(self.read_value_locked_result(&disk, shard_id as u8, &inner)?)
391 } else {
392 None
393 }
394 } else {
395 None
396 };
397 let existed = self.delete_locked(shard_id, &mut inner, &guard, key)?;
398 drop(inner);
399 if existed {
400 self.hook.on_write(key, old_value.as_deref(), None);
401 }
402 Ok(existed)
403 }
404
405 pub fn atomic<R>(
409 &self,
410 shard_key: &K,
411 f: impl FnOnce(&mut VarShard<'_, K, H>) -> DbResult<R>,
412 ) -> DbResult<R> {
413 let shard_id = self.shard_for(shard_key);
414 let inner = self.engine.shards()[shard_id].lock();
415 let guard = self.index.collector().enter();
416 let mut shard = VarShard {
417 tree: self,
418 inner,
419 shard_id,
420 guard,
421 events: Vec::new(),
422 };
423 let result = f(&mut shard);
424 let VarShard {
425 inner,
426 guard,
427 events,
428 ..
429 } = shard;
430 drop(inner);
431 if H::NEEDS_WRITE {
432 for (k, old, new) in &events {
433 self.hook.on_write(k, old.as_deref(), new.as_deref());
434 }
435 }
436 drop(guard);
437 result
438 }
439
440 fn put_locked(
442 &self,
443 shard_id: usize,
444 inner: &mut ShardInner,
445 guard: &seize::LocalGuard<'_>,
446 key: &K,
447 value: &[u8],
448 ) -> DbResult<bool> {
449 let (disk_loc, _gsn) = inner.append_entry(shard_id as u8, key.as_bytes(), value, false)?;
450
451 if let Some(existing) = self.index.get(key.as_bytes(), guard) {
453 let old_disk = existing.read_loc();
454 inner.add_dead_bytes(
455 old_disk.file_id,
456 crate::entry::entry_size(size_of::<K>(), old_disk.len),
457 );
458 existing.write_loc(disk_loc);
459 return Ok(true);
460 }
461
462 let height = random_height();
464 let node_ptr = VarNode::alloc(*key, disk_loc, height);
465
466 match self.index.insert(node_ptr, guard) {
467 InsertResult::Inserted => Ok(false),
468 InsertResult::Exists(existing) => {
469 let old_disk = existing.read_loc();
471 inner.add_dead_bytes(
472 old_disk.file_id,
473 crate::entry::entry_size(size_of::<K>(), old_disk.len),
474 );
475 existing.write_loc(disk_loc);
476 unsafe {
477 VarNode::<K>::dealloc_node(node_ptr);
478 }
479 Ok(true)
480 }
481 }
482 }
483
484 fn insert_locked(
485 &self,
486 shard_id: usize,
487 inner: &mut ShardInner,
488 guard: &seize::LocalGuard<'_>,
489 key: &K,
490 value: &[u8],
491 ) -> DbResult<()> {
492 if self.index.get(key.as_bytes(), guard).is_some() {
493 return Err(DbError::KeyExists);
494 }
495
496 let (disk_loc, _gsn) = inner.append_entry(shard_id as u8, key.as_bytes(), value, false)?;
497 let height = random_height();
498 let node_ptr = VarNode::alloc(*key, disk_loc, height);
499
500 match self.index.insert(node_ptr, guard) {
501 InsertResult::Inserted => Ok(()),
502 InsertResult::Exists(_existing) => {
503 inner.add_dead_bytes(
507 disk_loc.file_id,
508 crate::entry::entry_size(size_of::<K>(), disk_loc.len),
509 );
510 unsafe { VarNode::<K>::dealloc_node(node_ptr) };
514 Err(DbError::KeyExists)
515 }
516 }
517 }
518
519 fn delete_locked(
520 &self,
521 shard_id: usize,
522 inner: &mut ShardInner,
523 guard: &seize::LocalGuard<'_>,
524 key: &K,
525 ) -> DbResult<bool> {
526 if self.index.get(key.as_bytes(), guard).is_none() {
527 return Ok(false);
528 }
529
530 inner.append_entry(shard_id as u8, key.as_bytes(), &[], true)?;
531
532 let removed = self.index.remove(key.as_bytes(), guard);
533
534 if let Some(node_ptr) = removed {
535 let disk = unsafe { &*node_ptr }.read_loc();
536 inner.add_dead_bytes(
537 disk.file_id,
538 crate::entry::entry_size(size_of::<K>(), disk.len),
539 );
540 }
541
542 Ok(removed.is_some())
543 }
544
545 pub fn contains(&self, key: &K) -> bool {
547 let guard = self.index.collector().enter();
548 self.index.get(key.as_bytes(), &guard).is_some()
549 }
550
551 pub fn entry_len(&self, key: &K) -> Option<u32> {
554 let guard = self.index.collector().enter();
555 self.index
556 .get(key.as_bytes(), &guard)
557 .map(|node| node.read_loc().len)
558 }
559
560 pub fn first(&self) -> Option<(K, ByteView)> {
564 let guard = self.index.collector().enter();
565 let mut ptr = crate::skiplist::strip_mark(unsafe {
566 (*self.index.head_ptr())
567 .tower(0)
568 .load(std::sync::atomic::Ordering::Acquire)
569 });
570 while !ptr.is_null() {
571 let node = unsafe { &*ptr };
572 if !node.is_marked()
573 && let Some(v) = self.read_value_cached(node, &guard)
574 {
575 return Some((node.key, v));
576 }
577 ptr = crate::skiplist::strip_mark(
578 node.tower(0).load(std::sync::atomic::Ordering::Acquire),
579 );
580 }
581 None
582 }
583
584 pub fn last(&self) -> Option<(K, ByteView)> {
588 self.iter().next_back()
589 }
590
591 fn resolve_front(&self, bound: &Bound<&K>, guard: &seize::LocalGuard<'_>) -> *mut VarNode<K> {
594 match bound {
595 Bound::Included(k) => self.index.find_first_ge(k.as_bytes(), guard),
596 Bound::Excluded(k) => {
597 let ge = self.index.find_first_ge(k.as_bytes(), guard);
598 if !ge.is_null()
599 && !unsafe { &*ge }.is_marked()
600 && unsafe { &*ge }.key_bytes() == k.as_bytes()
601 {
602 crate::skiplist::strip_mark(unsafe {
603 (*ge).tower(0).load(std::sync::atomic::Ordering::Acquire)
604 })
605 } else {
606 ge
607 }
608 }
609 Bound::Unbounded => crate::skiplist::strip_mark(unsafe {
610 (*self.index.head_ptr())
611 .tower(0)
612 .load(std::sync::atomic::Ordering::Acquire)
613 }),
614 }
615 }
616
617 fn prefix_bounds(&self, prefix: &[u8]) -> (K, Bound<K>) {
618 if self.reversed {
619 let mut search = K::zeroed();
620 search.as_bytes_mut().fill(0xFF);
621 search.as_bytes_mut()[..prefix.len()].copy_from_slice(prefix);
622 let mut end_key = K::zeroed();
623 end_key.as_bytes_mut()[..prefix.len()].copy_from_slice(prefix);
624 (search, Bound::Included(end_key))
625 } else {
626 let mut search = K::zeroed();
627 search.as_bytes_mut()[..prefix.len()].copy_from_slice(prefix);
628 let end = prefix_to_end_bound::<K>(prefix);
629 (search, end)
630 }
631 }
632
633 pub fn prefix_iter(&self, prefix: &[u8]) -> VarIter<'_, K, H> {
638 let guard = self.index.collector().enter();
639 let (search_key, end) = self.prefix_bounds(prefix);
640 let front = self.index.find_first_ge(search_key.as_bytes(), &guard);
641 VarIter {
642 tree: self,
643 front,
644 back: None,
645 end,
646 start: Bound::Included(search_key),
647 reversed: self.reversed,
648 done: false,
649 _guard: guard,
650 }
651 }
652
653 pub fn iter(&self) -> VarIter<'_, K, H> {
658 let guard = self.index.collector().enter();
659 let front = crate::skiplist::strip_mark(unsafe {
660 (*self.index.head_ptr())
661 .tower(0)
662 .load(std::sync::atomic::Ordering::Acquire)
663 });
664 VarIter {
665 tree: self,
666 front,
667 back: None,
668 end: Bound::Unbounded,
669 start: Bound::Unbounded,
670 reversed: self.reversed,
671 done: false,
672 _guard: guard,
673 }
674 }
675
676 pub fn range(&self, start: &K, end: &K) -> VarIter<'_, K, H> {
681 self.range_bounds(Bound::Included(start), Bound::Excluded(end))
682 }
683
684 pub fn range_bounds(&self, start: Bound<&K>, end: Bound<&K>) -> VarIter<'_, K, H> {
692 let guard = self.index.collector().enter();
693 if self.reversed {
694 let front = self.resolve_front(&end, &guard);
695 VarIter {
696 tree: self,
697 front,
698 back: None,
699 end: bound_owned(&start),
700 start: bound_owned(&end),
701 reversed: true,
702 done: false,
703 _guard: guard,
704 }
705 } else {
706 let front = self.resolve_front(&start, &guard);
707 VarIter {
708 tree: self,
709 front,
710 back: None,
711 end: bound_owned(&end),
712 start: bound_owned(&start),
713 reversed: false,
714 done: false,
715 _guard: guard,
716 }
717 }
718 }
719
720 pub fn len(&self) -> usize {
721 self.index.len()
722 }
723
724 pub fn is_empty(&self) -> bool {
725 self.index.is_empty()
726 }
727
728 pub fn sync_hints(&self) -> DbResult<()> {
730 for shard in self.engine.shards().iter() {
731 shard.write_active_hint(size_of::<K>())?;
732 }
733 Ok(())
734 }
735
736 pub fn warmup(&self) -> DbResult<()> {
742 use std::collections::BTreeSet;
743
744 let guard = self.index.collector().enter();
745
746 let mut blocks: BTreeSet<(u8, u32, u64)> = BTreeSet::new();
748 let mut current = crate::skiplist::strip_mark(unsafe {
749 (*self.index.head_ptr())
750 .tower(0)
751 .load(std::sync::atomic::Ordering::Acquire)
752 });
753 while !current.is_null() {
754 let node = unsafe { &*current };
755 current = crate::skiplist::strip_mark(
756 node.tower(0).load(std::sync::atomic::Ordering::Acquire),
757 );
758 if node.is_marked() {
759 continue;
760 }
761 let disk = node.read_loc();
762 if disk.is_value_cache_routed() {
765 continue;
766 }
767 let shard_id = self.shard_for(&node.key) as u8;
768 let block_offset = disk.offset as u64 & !4095;
769 blocks.insert((shard_id, disk.file_id, block_offset));
770 }
771 drop(guard);
772
773 for (shard_id, file_id, block_offset) in &blocks {
775 let key = BlockKey {
776 shard_id: *shard_id,
777 file_id: *file_id,
778 block_offset: *block_offset,
779 };
780 if self.block_cache.get(&key).is_some() {
781 continue;
782 }
783 let shard = &self.engine.shards()[*shard_id as usize];
784 let (buf, is_full_block) = shard.read_block(*file_id, *block_offset)?;
785 if is_full_block {
786 self.block_cache.insert(key, Arc::new(buf));
787 }
788 }
789
790 Ok(())
791 }
792
793 pub(crate) fn index(&self) -> &SkipList<VarNode<K>> {
794 &self.index
795 }
796
797 pub fn migrate(
813 &self,
814 f: impl Fn(&K, &[u8]) -> crate::MigrateAction<ByteView>,
815 ) -> DbResult<usize> {
816 self.migrate_inner(f, true)
817 }
818
819 pub(crate) fn migrate_inner(
820 &self,
821 f: impl Fn(&K, &[u8]) -> crate::MigrateAction<ByteView>,
822 fire_init: bool,
823 ) -> DbResult<usize> {
824 use crate::MigrateAction;
825
826 let guard = self.index.collector().enter();
827 let mut current = crate::skiplist::strip_mark(unsafe {
828 (*self.index.head_ptr())
829 .tower(0)
830 .load(std::sync::atomic::Ordering::Acquire)
831 });
832 let mut count = 0;
833 while !current.is_null() {
834 let node = unsafe { &*current };
835 current = crate::skiplist::strip_mark(
836 node.tower(0).load(std::sync::atomic::Ordering::Acquire),
837 );
838 if node.is_marked() {
839 continue;
840 }
841 let value = match self.read_value_cached_result(node, &guard) {
842 Ok(v) => v,
843 Err(e) => {
844 tracing::error!(
845 key = ?node.key.as_bytes(),
846 error = %e,
847 "var_tree migrate: value read failed — aborting migration step"
848 );
849 return Err(e);
850 }
851 };
852 match f(&node.key, &value) {
853 MigrateAction::Keep => {
854 if fire_init && H::NEEDS_INIT {
855 self.hook.on_init(&node.key, &value);
856 }
857 }
858 MigrateAction::Update(new_value) => {
859 let shard_id = self.shard_for(&node.key);
860 {
861 let mut inner = self.engine.shards()[shard_id].lock();
862 self.put_locked(shard_id, &mut inner, &guard, &node.key, &new_value)?;
863 }
864 if fire_init && H::NEEDS_INIT {
865 self.hook.on_init(&node.key, &new_value);
866 }
867 count += 1;
868 }
869 MigrateAction::Delete => {
870 let shard_id = self.shard_for(&node.key);
871 let mut inner = self.engine.shards()[shard_id].lock();
872 self.delete_locked(shard_id, &mut inner, &guard, &node.key)?;
873 count += 1;
874 }
875 }
876 }
877
878 tracing::info!(mutations = count, "var_tree migration complete");
879 Ok(count)
880 }
881
882 pub(crate) fn replay_init(&self) {
886 if !H::NEEDS_INIT {
887 return;
888 }
889 let guard = self.index.collector().enter();
890 let mut current = crate::skiplist::strip_mark(unsafe {
891 (*self.index.head_ptr())
892 .tower(0)
893 .load(std::sync::atomic::Ordering::Acquire)
894 });
895 let mut count = 0usize;
896 while !current.is_null() {
897 let node = unsafe { &*current };
898 current = crate::skiplist::strip_mark(
899 node.tower(0).load(std::sync::atomic::Ordering::Acquire),
900 );
901 if node.is_marked() {
902 continue;
903 }
904 let value = match self.read_value_cached(node, &guard) {
905 Some(v) => v,
906 None => {
907 tracing::warn!(
908 key = ?node.key.as_bytes(),
909 "var_tree replay_init: skipping entry — value read failed"
910 );
911 continue;
912 }
913 };
914 self.hook.on_init(&node.key, &value);
915 count += 1;
916 }
917 tracing::debug!(replayed = count, "var_tree replay_init complete");
918 }
919
920 fn read_value_cached_inner(&self, disk: &DiskLoc, shard_id: u8) -> DbResult<ByteView> {
925 let len = disk.len as usize;
926 let start = (disk.offset & 4095) as usize;
927
928 if start + len > 8192 {
933 {
934 let shard = &self.engine.shards()[shard_id as usize];
935 let inner = shard.lock();
936 if inner.active.file_id == disk.file_id {
937 if let Some(bytes) = inner.write_buf.read(disk.offset as u64, len) {
938 return Ok(ByteView::new(bytes));
939 }
940 if inner.straddles_write_buf(disk) {
943 return Ok(ByteView::from_vec(
944 inner.read_straddling_value_locked(disk)?,
945 ));
946 }
947 }
948 }
949 let shard = &self.engine.shards()[shard_id as usize];
950 return Ok(ByteView::from_vec(shard.read_value(
951 disk.file_id,
952 disk.offset,
953 len,
954 )?));
955 }
956
957 let block_offset = disk.offset as u64 & !4095;
958 let cache_key = BlockKey {
959 shard_id,
960 file_id: disk.file_id,
961 block_offset,
962 };
963
964 if let Some(block) = self.block_cache.get(&cache_key) {
966 metrics::counter!("armdb.block_cache.hit").increment(1);
967 if start + len > 4096 {
974 let shard = &self.engine.shards()[shard_id as usize];
975 let inner = shard.lock();
976 if inner.active.file_id == disk.file_id && inner.straddles_write_buf(disk) {
977 return Ok(ByteView::from_vec(
978 inner.read_straddling_value_locked(disk)?,
979 ));
980 }
981 }
982 return Self::extract_from_block(&block, start, len, || {
983 self.get_or_read_block(shard_id, disk.file_id, block_offset + 4096)
984 });
985 }
986
987 {
989 let shard = &self.engine.shards()[shard_id as usize];
990 let inner = shard.lock();
991 if inner.active.file_id == disk.file_id {
992 if let Some(bytes) = inner.write_buf.read(disk.offset as u64, len) {
993 return Ok(ByteView::new(bytes));
994 }
995 if inner.straddles_write_buf(disk) {
999 return Ok(ByteView::from_vec(
1000 inner.read_straddling_value_locked(disk)?,
1001 ));
1002 }
1003 }
1004 }
1005
1006 metrics::counter!("armdb.block_cache.miss").increment(1);
1008 let block = self.get_or_read_block(shard_id, disk.file_id, block_offset)?;
1009 Self::extract_from_block(&block, start, len, || {
1010 self.get_or_read_block(shard_id, disk.file_id, block_offset + 4096)
1011 })
1012 }
1013
1014 fn read_value_cached_result(
1020 &self,
1021 node: &VarNode<K>,
1022 _guard: &seize::LocalGuard<'_>,
1023 ) -> DbResult<ByteView> {
1024 let shard_id = self.shard_for(&node.key) as u8;
1025 for _ in 0..MAX_STALE_RETRIES {
1026 let disk = node.read_loc();
1027 let is_large = disk.is_value_cache_routed();
1028 if is_large {
1029 let vkey = ValueKey {
1030 shard_id,
1031 file_id: disk.file_id,
1032 offset: disk.offset,
1033 };
1034 if let Some(v) = self.value_cache.get(&vkey) {
1035 let now = node.read_loc();
1037 if now.file_id == disk.file_id && now.offset == disk.offset {
1041 metrics::counter!("armdb.value_cache.hit").increment(1);
1042 return Ok(v);
1043 }
1044 continue; }
1046 metrics::counter!("armdb.value_cache.miss").increment(1);
1047 match self.read_value_cached_inner(&disk, shard_id) {
1048 Ok(v) => {
1049 self.value_cache.insert(vkey, v.clone());
1050 return Ok(v);
1051 }
1052 Err(DbError::StaleDiskLoc) => {
1053 metrics::counter!("armdb.read.stale_retry", "tree" => "var_tree")
1054 .increment(1);
1055 continue;
1056 }
1057 Err(e) => return Err(e),
1058 }
1059 }
1060 match self.read_value_cached_inner(&disk, shard_id) {
1061 Ok(v) => return Ok(v),
1062 Err(DbError::StaleDiskLoc) => {
1063 metrics::counter!("armdb.read.stale_retry", "tree" => "var_tree").increment(1);
1064 continue;
1065 }
1066 Err(e) => return Err(e),
1067 }
1068 }
1069 Err(DbError::StaleDiskLoc)
1070 }
1071
1072 fn read_value_cached(
1075 &self,
1076 node: &VarNode<K>,
1077 guard: &seize::LocalGuard<'_>,
1078 ) -> Option<ByteView> {
1079 let result = self.read_value_cached_result(node, guard);
1080 #[cfg(feature = "hot-path-tracing")]
1081 if let Err(ref _e) = result {
1082 tracing::error!("VarTree read_value_cached error: {:?}", _e);
1083 }
1084 result.ok()
1085 }
1086
1087 fn extract_from_block(
1088 block: &AlignedBuf,
1089 start: usize,
1090 len: usize,
1091 next_block: impl FnOnce() -> DbResult<Arc<AlignedBuf>>,
1092 ) -> DbResult<ByteView> {
1093 debug_assert!(
1094 start + len <= 8192,
1095 "extract_from_block supports at most 2 blocks (8192 bytes)"
1096 );
1097 if start + len <= 4096 {
1098 Ok(ByteView::new(&block[start..start + len]))
1099 } else {
1100 let next = next_block()?;
1101 let first_part = &block[start..];
1102 let second_len = len - first_part.len();
1103 let mut combined = Vec::with_capacity(len);
1104 combined.extend_from_slice(first_part);
1105 combined.extend_from_slice(&next[..second_len]);
1106 Ok(ByteView::from_vec(combined))
1107 }
1108 }
1109
1110 fn get_or_read_block(
1112 &self,
1113 shard_id: u8,
1114 file_id: u32,
1115 block_offset: u64,
1116 ) -> DbResult<Arc<AlignedBuf>> {
1117 let key = BlockKey {
1118 shard_id,
1119 file_id,
1120 block_offset,
1121 };
1122 if let Some(cached) = self.block_cache.get(&key) {
1123 return Ok(cached);
1124 }
1125 let shard = &self.engine.shards()[shard_id as usize];
1126 let (buf, is_full_block) = shard.read_block(file_id, block_offset)?;
1127 let arc = Arc::new(buf);
1128 if is_full_block {
1132 self.block_cache.insert(key, arc.clone());
1133 }
1134 Ok(arc)
1135 }
1136
1137 pub fn cas(&self, key: &K, expected: &[u8], new_value: &[u8]) -> DbResult<()> {
1145 metrics::counter!("armdb.ops", "op" => "cas", "tree" => "var_tree").increment(1);
1146 #[cfg(feature = "hot-path-tracing")]
1147 tracing::trace!("var_tree.cas");
1148 let shard_id = self.shard_for(key);
1149 let shard = &self.engine.shards()[shard_id];
1150 let mut inner = shard.lock();
1151
1152 let guard = self.index.collector().enter();
1153 let existing = self
1154 .index
1155 .get(key.as_bytes(), &guard)
1156 .ok_or(DbError::KeyNotFound)?;
1157
1158 let disk = existing.read_loc();
1159 let current = self.read_value_locked_result(&disk, shard_id as u8, &inner)?;
1160 if current.as_ref() != expected {
1161 return Err(DbError::CasMismatch);
1162 }
1163
1164 let (new_disk_loc, _gsn) =
1165 inner.append_entry(shard_id as u8, key.as_bytes(), new_value, false)?;
1166
1167 let old_disk = existing.read_loc();
1168 inner.add_dead_bytes(
1169 old_disk.file_id,
1170 crate::entry::entry_size(size_of::<K>(), old_disk.len),
1171 );
1172 existing.write_loc(new_disk_loc);
1173
1174 drop(inner);
1175 self.hook.on_write(
1176 key,
1177 if H::NEEDS_OLD_VALUE {
1178 Some(&*current)
1179 } else {
1180 None
1181 },
1182 Some(new_value),
1183 );
1184 Ok(())
1185 }
1186
1187 pub fn compare_delete(&self, key: &K, expected: &[u8]) -> DbResult<()> {
1195 metrics::counter!("armdb.ops", "op" => "compare_delete", "tree" => "var_tree").increment(1);
1196 #[cfg(feature = "hot-path-tracing")]
1197 tracing::trace!("var_tree.compare_delete");
1198 let shard_id = self.shard_for(key);
1199 let shard = &self.engine.shards()[shard_id];
1200 let mut inner = shard.lock();
1201
1202 let guard = self.index.collector().enter();
1203 let existing = self
1204 .index
1205 .get(key.as_bytes(), &guard)
1206 .ok_or(DbError::KeyNotFound)?;
1207
1208 let disk = existing.read_loc();
1209 let current = self.read_value_locked_result(&disk, shard_id as u8, &inner)?;
1210 if current.as_ref() != expected {
1211 return Err(DbError::CasMismatch);
1212 }
1213
1214 let removed = self.delete_locked(shard_id, &mut inner, &guard, key)?;
1216 debug_assert!(
1217 removed,
1218 "compare_delete: key must exist after a successful match"
1219 );
1220
1221 drop(inner);
1222 self.hook.on_write(
1223 key,
1224 if H::NEEDS_OLD_VALUE {
1225 Some(&*current)
1226 } else {
1227 None
1228 },
1229 None,
1230 );
1231 Ok(())
1232 }
1233
1234 pub fn update(&self, key: &K, f: impl FnOnce(&[u8]) -> ByteView) -> DbResult<Option<ByteView>> {
1242 self.update_inner(key, f, false)
1243 }
1244
1245 pub fn fetch_update(
1247 &self,
1248 key: &K,
1249 f: impl FnOnce(&[u8]) -> ByteView,
1250 ) -> DbResult<Option<ByteView>> {
1251 self.update_inner(key, f, true)
1252 }
1253
1254 fn try_update_locked(
1258 &self,
1259 shard_id: usize,
1260 inner: &mut ShardInner,
1261 guard: &seize::LocalGuard<'_>,
1262 key: &K,
1263 f: impl FnOnce(&[u8]) -> DbResult<Option<ByteView>>,
1264 ) -> DbResult<UpdateOutcome> {
1265 let existing = match self.index.get(key.as_bytes(), guard) {
1266 Some(n) => n,
1267 None => return Ok(UpdateOutcome::Missing),
1268 };
1269 let disk = existing.read_loc();
1270 let current = self.read_value_locked_result(&disk, shard_id as u8, inner)?;
1271
1272 let new_value = match f(¤t)? {
1273 Some(v) => v,
1274 None => return Ok(UpdateOutcome::Unchanged(current)),
1275 };
1276
1277 let (new_disk_loc, _gsn) =
1278 inner.append_entry(shard_id as u8, key.as_bytes(), &new_value, false)?;
1279 let old_disk = existing.read_loc();
1280 inner.add_dead_bytes(
1281 old_disk.file_id,
1282 crate::entry::entry_size(size_of::<K>(), old_disk.len),
1283 );
1284 existing.write_loc(new_disk_loc);
1285
1286 Ok(UpdateOutcome::Updated {
1287 old: current,
1288 new: new_value,
1289 })
1290 }
1291
1292 pub(crate) fn try_update_inner(
1293 &self,
1294 key: &K,
1295 f: impl FnOnce(&[u8]) -> DbResult<Option<ByteView>>,
1296 return_old: bool,
1297 ) -> DbResult<Option<ByteView>> {
1298 metrics::counter!("armdb.ops", "op" => "update", "tree" => "var_tree").increment(1);
1299 #[cfg(feature = "hot-path-tracing")]
1300 tracing::trace!("var_tree.update");
1301 let shard_id = self.shard_for(key);
1302 let shard = &self.engine.shards()[shard_id];
1303 let mut inner = shard.lock();
1304 let guard = self.index.collector().enter();
1305
1306 match self.try_update_locked(shard_id, &mut inner, &guard, key, f)? {
1307 UpdateOutcome::Missing => Ok(None),
1308 UpdateOutcome::Unchanged(current) => Ok(Some(current)),
1309 UpdateOutcome::Updated { old, new } => {
1310 drop(inner);
1311 self.hook.on_write(
1312 key,
1313 if H::NEEDS_OLD_VALUE {
1314 Some(&*old)
1315 } else {
1316 None
1317 },
1318 Some(&*new),
1319 );
1320 Ok(Some(if return_old { old } else { new }))
1321 }
1322 }
1323 }
1324
1325 fn update_inner(
1326 &self,
1327 key: &K,
1328 f: impl FnOnce(&[u8]) -> ByteView,
1329 return_old: bool,
1330 ) -> DbResult<Option<ByteView>> {
1331 self.try_update_inner(key, |bytes| Ok(Some(f(bytes))), return_old)
1332 }
1333
1334 fn read_value_locked_result(
1338 &self,
1339 disk: &DiskLoc,
1340 shard_id: u8,
1341 inner: &ShardInner,
1342 ) -> DbResult<ByteView> {
1343 let len = disk.len as usize;
1344
1345 if inner.active.file_id == disk.file_id {
1349 if let Some(bytes) = inner.write_buf.read(disk.offset as u64, len) {
1350 return Ok(ByteView::new(bytes));
1351 }
1352 if inner.straddles_write_buf(disk) {
1353 return Ok(ByteView::from_vec(
1354 inner.read_straddling_value_locked(disk)?,
1355 ));
1356 }
1357 }
1358
1359 let block_offset = disk.offset as u64 & !4095;
1361 let start = (disk.offset & 4095) as usize;
1362 if start + len <= 4096 {
1363 let cache_key = BlockKey {
1364 shard_id,
1365 file_id: disk.file_id,
1366 block_offset,
1367 };
1368 if let Some(block) = self.block_cache.get(&cache_key) {
1369 return Ok(ByteView::new(&block[start..start + len]));
1370 }
1371
1372 let (buf, is_full_block) = inner.read_block_locked(disk.file_id, block_offset)?;
1374 let arc = Arc::new(buf);
1375 if is_full_block {
1376 self.block_cache.insert(cache_key, arc.clone());
1377 }
1378 return Ok(ByteView::new(&arc[start..start + len]));
1379 }
1380
1381 let bytes = inner.read_value_from_disk_locked(disk)?;
1383 Ok(ByteView::new(&bytes))
1384 }
1385
1386 pub fn shard_for(&self, key: &K) -> usize {
1387 crate::shard_for_key(key, self.shard_prefix_bits, self.engine.shards().len())
1388 }
1389
1390 pub fn get_many(&self, keys: &[K]) -> Vec<Option<ByteView>> {
1395 let n = keys.len();
1396 let mut out: Vec<Option<ByteView>> = (0..n).map(|_| None).collect();
1397 if n == 0 {
1398 return out;
1399 }
1400 let mut order: Vec<usize> = (0..n).collect();
1401 order.sort_by(|&a, &b| {
1402 let ord = keys[a].as_bytes().cmp(keys[b].as_bytes());
1403 if self.reversed { ord.reverse() } else { ord }
1404 });
1405 let guard = self.index.collector().enter();
1406 let mut finger = self.index.finger();
1407 for &i in &order {
1408 let kb = keys[i].as_bytes();
1409 if let Some(node) = self.index.finger_seek(&mut finger, kb, &guard)
1410 && node.key_bytes() == kb
1411 {
1412 out[i] = self.read_value_cached(node, &guard);
1413 }
1414 }
1415 out
1416 }
1417
1418 pub fn update_many<P>(
1422 &self,
1423 items: Vec<(K, P)>,
1424 f: impl Fn(&K, Option<&[u8]>, &P) -> BatchWrite<ByteView>,
1425 ) -> DbResult<Vec<(K, Applied<ByteView>)>> {
1426 let n = items.len();
1427 let mut by_shard: std::collections::BTreeMap<usize, Vec<usize>> =
1428 std::collections::BTreeMap::new();
1429 for (i, (k, _)) in items.iter().enumerate() {
1430 by_shard.entry(self.shard_for(k)).or_default().push(i);
1431 }
1432
1433 let mut out: Vec<Option<(K, Applied<ByteView>)>> = (0..n).map(|_| None).collect();
1434 let mut events: Vec<(K, Option<ByteView>, Option<ByteView>)> = Vec::new();
1435 let mut err: Option<DbError> = None;
1436
1437 'shards: for (shard_id, idxs) in by_shard {
1438 let mut inner = self.engine.shards()[shard_id].lock();
1439 let guard = self.index.collector().enter();
1440 for &i in &idxs {
1441 let (k, p) = &items[i];
1442 let cur: Option<ByteView> = match self.index.get(k.as_bytes(), &guard) {
1446 Some(node) => {
1447 let disk = node.read_loc();
1448 match self.read_value_locked_result(&disk, shard_id as u8, &inner) {
1449 Ok(v) => Some(v),
1450 Err(e) => {
1451 err = Some(e);
1452 drop(inner);
1453 drop(guard);
1454 break 'shards;
1455 }
1456 }
1457 }
1458 None => None,
1459 };
1460 let applied = match f(k, cur.as_deref(), p) {
1461 BatchWrite::Set(v) => {
1462 match self.put_locked(shard_id, &mut inner, &guard, k, &v) {
1463 Ok(_existed) => {
1464 if H::NEEDS_WRITE {
1465 events.push((*k, cur.clone(), Some(v.clone())));
1466 }
1467 Applied::Written { old: cur, new: v }
1468 }
1469 Err(e) => {
1470 err = Some(e);
1471 drop(inner);
1472 drop(guard);
1473 break 'shards;
1474 }
1475 }
1476 }
1477 BatchWrite::Keep => Applied::Kept,
1478 BatchWrite::Delete => {
1479 match self.delete_locked(shard_id, &mut inner, &guard, k) {
1480 Ok(true) => {
1481 let old = cur.expect("existing key has an old value");
1484 if H::NEEDS_WRITE {
1485 events.push((*k, Some(old.clone()), None));
1486 }
1487 Applied::Deleted(old)
1488 }
1489 Ok(false) => Applied::Kept, Err(e) => {
1491 err = Some(e);
1492 drop(inner);
1493 drop(guard);
1494 break 'shards;
1495 }
1496 }
1497 }
1498 };
1499 out[i] = Some((*k, applied));
1500 }
1501 drop(inner);
1502 drop(guard);
1503 }
1504
1505 if H::NEEDS_WRITE {
1506 for (k, old, new) in &events {
1507 self.hook.on_write(k, old.as_deref(), new.as_deref());
1508 }
1509 }
1510
1511 if let Some(e) = err {
1512 return Err(e);
1513 }
1514 Ok(out
1515 .into_iter()
1516 .map(|slot| slot.expect("every index assigned on the success path"))
1517 .collect())
1518 }
1519}
1520
1521#[cfg(feature = "replication")]
1522impl<K: Key, H: WriteHook<K>> crate::replication::ReplicationTarget for VarTree<K, H> {
1523 fn apply_entry(
1524 &self,
1525 _shard_inner: &mut crate::shard::ShardInner,
1526 _shard_id: u8,
1527 file_id: u32,
1528 entry_offset: u64,
1529 header: &crate::entry::EntryHeader,
1530 key: &[u8],
1531 _value: &[u8],
1532 ) -> DbResult<crate::replication::ApplyOutcome> {
1533 use crate::replication::ApplyOutcome;
1534
1535 let key: K = K::from_bytes(key);
1536
1537 let value_offset =
1538 entry_offset + size_of::<crate::entry::EntryHeader>() as u64 + size_of::<K>() as u64;
1539 let disk = DiskLoc::new(file_id, value_offset as u32, header.value_len);
1540
1541 if header.is_tombstone() {
1542 let guard = self.index.collector().enter();
1543 let removed = self.index.remove(key.as_bytes(), &guard);
1544 match removed {
1545 Some(node_ptr) => {
1546 let old_disk = unsafe { &*node_ptr }.read_loc();
1547 Ok(ApplyOutcome::TombstoneRemoved(old_disk))
1548 }
1549 None => Ok(ApplyOutcome::Inserted), }
1551 } else {
1552 let guard = self.index.collector().enter();
1553 let height = random_height();
1554 let node_ptr = VarNode::alloc(key, disk, height);
1555 match self.index.insert(node_ptr, &guard) {
1556 InsertResult::Inserted => Ok(ApplyOutcome::Inserted),
1557 InsertResult::Exists(existing) => {
1558 let old_disk = existing.read_loc();
1559 existing.write_loc(disk);
1560 unsafe {
1561 VarNode::<K>::dealloc_node(node_ptr);
1562 }
1563 Ok(ApplyOutcome::Replaced(old_disk))
1564 }
1565 }
1566 }
1567 }
1568
1569 fn try_apply_entry(
1570 &self,
1571 shard_inner: &mut crate::shard::ShardInner,
1572 shard_id: u8,
1573 file_id: u32,
1574 entry_offset: u64,
1575 header: &crate::entry::EntryHeader,
1576 raw_after_header: &[u8],
1577 ) -> DbResult<crate::replication::ApplyOutcome> {
1578 use crate::replication::ApplyOutcome;
1579
1580 if raw_after_header.len() < size_of::<K>() + header.value_len as usize {
1581 return Ok(ApplyOutcome::NotMatched);
1582 }
1583 let key = &raw_after_header[..size_of::<K>()];
1584 let value = &raw_after_header[size_of::<K>()..size_of::<K>() + header.value_len as usize];
1585 let crc = crate::entry::compute_crc32(header.gsn, header.value_len, key, value);
1586 if crc != header.crc32 {
1587 return Ok(ApplyOutcome::NotMatched);
1588 }
1589 self.apply_entry(
1590 shard_inner,
1591 shard_id,
1592 file_id,
1593 entry_offset,
1594 header,
1595 key,
1596 value,
1597 )
1598 }
1599
1600 fn key_len(&self) -> usize {
1601 size_of::<K>()
1602 }
1603}
1604
1605#[cfg(feature = "replication")]
1606impl<K: Key, H: WriteHook<K>> VarTree<K, H> {
1607 pub fn start_replication_server(
1618 &self,
1619 bind_addr: std::net::SocketAddr,
1620 signal: crate::shutdown::ShutdownSignal,
1621 ) -> crate::error::DbResult<crate::replication::ReplicationServer> {
1622 let consumers = self.install_replication_producers()?;
1623 crate::replication::ReplicationServer::start(
1624 bind_addr,
1625 self.engine.shards().clone(),
1626 consumers,
1627 self.engine.config().max_file_size,
1628 signal,
1629 )
1630 }
1631
1632 pub fn start_replication_server_with_options(
1633 &self,
1634 bind_addr: std::net::SocketAddr,
1635 signal: crate::shutdown::ShutdownSignal,
1636 options: crate::replication::ReplicationServerOptions,
1637 ) -> crate::error::DbResult<crate::replication::ReplicationServer> {
1638 let consumers = self.install_replication_producers()?;
1639 crate::replication::ReplicationServer::start_with_options(
1640 bind_addr,
1641 self.engine.shards().clone(),
1642 consumers,
1643 self.engine.config().max_file_size,
1644 signal,
1645 options,
1646 )
1647 }
1648
1649 fn install_replication_producers(
1650 &self,
1651 ) -> crate::error::DbResult<Vec<rtrb::Consumer<crate::replication::ReplicationEntry>>> {
1652 const SPSC_CAPACITY: usize = 4096;
1653 let shards = self.engine.shards();
1654 let mut consumers = Vec::with_capacity(shards.len());
1655 for shard in shards.iter() {
1656 let (p, c) = rtrb::RingBuffer::new(SPSC_CAPACITY);
1657 shard.set_replication_producer(p);
1658 consumers.push(c);
1659 }
1660 Ok(consumers)
1661 }
1662
1663 pub fn start_replication_client(
1668 &self,
1669 leader_addr: std::net::SocketAddr,
1670 registry: std::sync::Arc<crate::replication::ReplicationRegistry>,
1671 signal: crate::shutdown::ShutdownSignal,
1672 ) -> crate::error::DbResult<crate::replication::ReplicationClient> {
1673 crate::replication::ReplicationClient::start(
1674 leader_addr,
1675 self.engine.shards().clone(),
1676 registry,
1677 size_of::<K>() as u16,
1678 signal,
1679 )
1680 }
1681
1682 pub fn start_replication_client_with_options(
1683 &self,
1684 leader_addr: std::net::SocketAddr,
1685 registry: std::sync::Arc<crate::replication::ReplicationRegistry>,
1686 signal: crate::shutdown::ShutdownSignal,
1687 options: crate::replication::ReplicationClientOptions,
1688 ) -> crate::error::DbResult<crate::replication::ReplicationClient> {
1689 crate::replication::ReplicationClient::start_with_options(
1690 leader_addr,
1691 self.engine.shards().clone(),
1692 registry,
1693 size_of::<K>() as u16,
1694 signal,
1695 options,
1696 )
1697 }
1698}
1699
1700#[cfg(feature = "replication")]
1701impl<K, H> VarTree<K, H>
1702where
1703 K: Key + Send + Sync + 'static,
1704 H: WriteHook<K> + Send + Sync + 'static,
1705{
1706 pub fn as_replication_target(
1718 self: &std::sync::Arc<Self>,
1719 ) -> Box<dyn crate::replication::ReplicationTarget> {
1720 Box::new(std::sync::Arc::clone(self))
1721 }
1722}
1723
1724pub struct VarShard<'a, K: Key, H: WriteHook<K> = NoHook> {
1728 tree: &'a VarTree<K, H>,
1729 inner: MutexGuard<'a, ShardInner>,
1730 shard_id: usize,
1731 guard: seize::LocalGuard<'a>,
1732 events: Vec<(K, Option<ByteView>, Option<ByteView>)>,
1733}
1734
1735impl<K: Key, H: WriteHook<K>> VarShard<'_, K, H> {
1736 pub fn put(&mut self, key: &K, value: &[u8]) -> DbResult<bool> {
1739 self.check_shard(key)?;
1740 let old = if H::NEEDS_OLD_VALUE && H::NEEDS_WRITE {
1741 if let Some(node) = self.tree.index.get(key.as_bytes(), &self.guard) {
1742 let disk = node.read_loc();
1743 Some(
1744 self.tree
1745 .read_value_locked_result(&disk, self.shard_id as u8, &self.inner)?,
1746 )
1747 } else {
1748 None
1749 }
1750 } else {
1751 None
1752 };
1753 let existed =
1754 self.tree
1755 .put_locked(self.shard_id, &mut self.inner, &self.guard, key, value)?;
1756 if H::NEEDS_WRITE {
1757 self.events.push((*key, old, Some(ByteView::from(value))));
1758 }
1759 Ok(existed)
1760 }
1761
1762 pub fn insert(&mut self, key: &K, value: &[u8]) -> DbResult<()> {
1763 self.check_shard(key)?;
1764 self.tree
1765 .insert_locked(self.shard_id, &mut self.inner, &self.guard, key, value)?;
1766 if H::NEEDS_WRITE {
1767 self.events.push((*key, None, Some(ByteView::from(value))));
1768 }
1769 Ok(())
1770 }
1771
1772 pub fn delete(&mut self, key: &K) -> DbResult<bool> {
1773 self.check_shard(key)?;
1774 let old = if H::NEEDS_OLD_VALUE && H::NEEDS_WRITE {
1775 if let Some(node) = self.tree.index.get(key.as_bytes(), &self.guard) {
1776 let disk = node.read_loc();
1777 Some(
1778 self.tree
1779 .read_value_locked_result(&disk, self.shard_id as u8, &self.inner)?,
1780 )
1781 } else {
1782 None
1783 }
1784 } else {
1785 None
1786 };
1787 let existed = self
1788 .tree
1789 .delete_locked(self.shard_id, &mut self.inner, &self.guard, key)?;
1790 if existed && H::NEEDS_WRITE {
1791 self.events.push((*key, old, None));
1792 }
1793 Ok(existed)
1794 }
1795
1796 pub fn try_get(&self, key: &K) -> DbResult<Option<ByteView>> {
1797 self.check_shard(key)?;
1798 match self.tree.index.get(key.as_bytes(), &self.guard) {
1799 Some(node) => {
1800 let disk = node.read_loc();
1801 let view =
1802 self.tree
1803 .read_value_locked_result(&disk, self.shard_id as u8, &self.inner)?;
1804 Ok(Some(view))
1805 }
1806 None => Ok(None),
1807 }
1808 }
1809
1810 pub fn get_or_err(&self, key: &K) -> DbResult<ByteView> {
1811 self.check_shard(key)?;
1812 let node = self
1813 .tree
1814 .index
1815 .get(key.as_bytes(), &self.guard)
1816 .ok_or(DbError::KeyNotFound)?;
1817 let disk = node.read_loc();
1818 self.tree
1819 .read_value_locked_result(&disk, self.shard_id as u8, &self.inner)
1820 }
1821
1822 pub fn try_contains(&self, key: &K) -> DbResult<bool> {
1823 self.check_shard(key)?;
1824 Ok(self.tree.index.get(key.as_bytes(), &self.guard).is_some())
1825 }
1826
1827 pub fn update(
1831 &mut self,
1832 key: &K,
1833 f: impl FnOnce(&[u8]) -> ByteView,
1834 ) -> DbResult<Option<ByteView>> {
1835 self.update_inner(key, |b| Ok(Some(f(b))), false)
1836 }
1837
1838 pub fn fetch_update(
1840 &mut self,
1841 key: &K,
1842 f: impl FnOnce(&[u8]) -> ByteView,
1843 ) -> DbResult<Option<ByteView>> {
1844 self.update_inner(key, |b| Ok(Some(f(b))), true)
1845 }
1846
1847 pub(crate) fn update_inner(
1850 &mut self,
1851 key: &K,
1852 f: impl FnOnce(&[u8]) -> DbResult<Option<ByteView>>,
1853 return_old: bool,
1854 ) -> DbResult<Option<ByteView>> {
1855 self.check_shard(key)?;
1856 match self
1857 .tree
1858 .try_update_locked(self.shard_id, &mut self.inner, &self.guard, key, f)?
1859 {
1860 UpdateOutcome::Missing => Ok(None),
1861 UpdateOutcome::Unchanged(current) => Ok(Some(current)),
1862 UpdateOutcome::Updated { old, new } => {
1863 if H::NEEDS_WRITE {
1864 let event_old = if H::NEEDS_OLD_VALUE {
1865 Some(old.clone())
1866 } else {
1867 None
1868 };
1869 self.events.push((*key, event_old, Some(new.clone())));
1870 }
1871 Ok(Some(if return_old { old } else { new }))
1872 }
1873 }
1874 }
1875
1876 fn check_shard(&self, key: &K) -> DbResult<()> {
1877 if self.tree.shard_for(key) != self.shard_id {
1878 return Err(DbError::ShardMismatch);
1879 }
1880 Ok(())
1881 }
1882}
1883
1884fn bound_owned<K: Copy>(b: &Bound<&K>) -> Bound<K> {
1885 match b {
1886 Bound::Included(k) => Bound::Included(**k),
1887 Bound::Excluded(k) => Bound::Excluded(**k),
1888 Bound::Unbounded => Bound::Unbounded,
1889 }
1890}
1891
1892fn prefix_to_end_bound<K: Key>(prefix: &[u8]) -> Bound<K> {
1893 let mut incremented = prefix.to_vec();
1894 let mut carry = true;
1895 for byte in incremented.iter_mut().rev() {
1896 if carry {
1897 if *byte == 0xFF {
1898 *byte = 0x00;
1899 } else {
1900 *byte += 1;
1901 carry = false;
1902 break;
1903 }
1904 }
1905 }
1906 if carry {
1907 Bound::Unbounded
1908 } else {
1909 let mut end = K::zeroed();
1910 end.as_bytes_mut()[..incremented.len()].copy_from_slice(&incremented);
1911 Bound::Excluded(end)
1912 }
1913}
1914
1915pub struct VarIter<'a, K: Key, H: WriteHook<K> = NoHook> {
1921 tree: &'a VarTree<K, H>,
1922 front: *mut VarNode<K>,
1923 back: Option<*mut VarNode<K>>,
1925 end: Bound<K>,
1926 start: Bound<K>,
1927 reversed: bool,
1928 done: bool,
1929 _guard: seize::LocalGuard<'a>,
1930}
1931
1932impl<K: Key, H: WriteHook<K>> Iterator for VarIter<'_, K, H> {
1933 type Item = (K, ByteView);
1934
1935 fn next(&mut self) -> Option<Self::Item> {
1936 loop {
1937 if self.done || self.front.is_null() {
1938 return None;
1939 }
1940 let node = unsafe { &*self.front };
1941 let converged = self.back.is_some_and(|back| std::ptr::eq(self.front, back));
1942 self.front = crate::skiplist::strip_mark(
1943 node.tower(0).load(std::sync::atomic::Ordering::Acquire),
1944 );
1945 if converged {
1946 self.done = true;
1947 }
1948 if node.is_marked() {
1949 if converged {
1950 return None;
1951 }
1952 continue;
1953 }
1954 if !self.check_end(&node.key) {
1955 self.done = true;
1956 return None;
1957 }
1958 match self.tree.read_value_cached(node, &self._guard) {
1959 Some(value) => return Some((node.key, value)),
1960 None => {
1961 if converged {
1962 return None;
1963 }
1964 continue;
1965 }
1966 }
1967 }
1968 }
1969}
1970
1971impl<K: Key, H: WriteHook<K>> DoubleEndedIterator for VarIter<'_, K, H> {
1972 fn next_back(&mut self) -> Option<Self::Item> {
1973 if self.back.is_none() {
1974 self.back = Some(self.resolve_back());
1975 if self.front.is_null() {
1976 self.done = true;
1977 }
1978 }
1979 loop {
1980 let back = self.back.unwrap_or(std::ptr::null_mut());
1981 if self.done || back.is_null() {
1982 return None;
1983 }
1984 let node = unsafe { &*back };
1985 let key = node.key;
1986 let converged = std::ptr::eq(self.front, back);
1987 self.back = Some(self.tree.index().find_last_lt(key.as_bytes(), &self._guard));
1988 if converged {
1989 self.done = true;
1990 }
1991 if node.is_marked() {
1992 if converged {
1993 return None;
1994 }
1995 continue;
1996 }
1997 if !self.check_start(&key) {
1998 self.done = true;
1999 return None;
2000 }
2001 match self.tree.read_value_cached(node, &self._guard) {
2002 Some(value) => return Some((key, value)),
2003 None => {
2004 if converged {
2005 return None;
2006 }
2007 continue;
2008 }
2009 }
2010 }
2011 }
2012}
2013
2014impl<K: Key, H: WriteHook<K>> VarIter<'_, K, H> {
2015 fn resolve_back(&self) -> *mut VarNode<K> {
2017 let index = self.tree.index();
2018 match &self.end {
2019 Bound::Unbounded => index.find_last(&self._guard),
2020 Bound::Excluded(k) => index.find_last_lt(k.as_bytes(), &self._guard),
2021 Bound::Included(k) => {
2022 let ge = index.find_first_ge(k.as_bytes(), &self._guard);
2023 if !ge.is_null()
2024 && !unsafe { &*ge }.is_marked()
2025 && unsafe { &*ge }.key_bytes() == k.as_bytes()
2026 {
2027 ge
2028 } else {
2029 index.find_last_lt(k.as_bytes(), &self._guard)
2030 }
2031 }
2032 }
2033 }
2034
2035 #[inline(always)]
2036 fn check_end(&self, key: &K) -> bool {
2037 match &self.end {
2038 Bound::Unbounded => true,
2039 Bound::Excluded(end) => {
2040 if self.reversed {
2041 key.as_bytes() > end.as_bytes()
2042 } else {
2043 key.as_bytes() < end.as_bytes()
2044 }
2045 }
2046 Bound::Included(end) => {
2047 if self.reversed {
2048 key.as_bytes() >= end.as_bytes()
2049 } else {
2050 key.as_bytes() <= end.as_bytes()
2051 }
2052 }
2053 }
2054 }
2055
2056 #[inline(always)]
2057 fn check_start(&self, key: &K) -> bool {
2058 match &self.start {
2059 Bound::Unbounded => true,
2060 Bound::Excluded(s) => {
2061 if self.reversed {
2062 key.as_bytes() < s.as_bytes()
2063 } else {
2064 key.as_bytes() > s.as_bytes()
2065 }
2066 }
2067 Bound::Included(s) => {
2068 if self.reversed {
2069 key.as_bytes() <= s.as_bytes()
2070 } else {
2071 key.as_bytes() >= s.as_bytes()
2072 }
2073 }
2074 }
2075 }
2076 pub fn collect_keys(&mut self) -> Vec<K> {
2078 self.map(|(k, _)| k).collect()
2079 }
2080
2081 pub fn collect_entries(&mut self) -> Vec<(K, ByteView)> {
2083 self.collect()
2084 }
2085}
2086
2087#[cfg(feature = "armour")]
2097pub struct VarTx<'a, K: Key, H: WriteHook<K> = NoHook> {
2098 tree: &'a VarTree<K, H>,
2099 inners: Vec<(usize, MutexGuard<'a, ShardInner>)>,
2100 seize: seize::LocalGuard<'a>,
2101 log: Vec<(K, Option<ByteView>, Option<ByteView>)>,
2102}
2103
2104#[cfg(feature = "armour")]
2105impl<K: Key, H: WriteHook<K>> VarTx<'_, K, H> {
2106 fn position(&self, key: &K) -> DbResult<usize> {
2107 let sid = self.tree.shard_for(key);
2108 self.inners
2109 .iter()
2110 .position(|(s, _)| *s == sid)
2111 .ok_or(DbError::ShardMismatch)
2112 }
2113
2114 pub fn try_get(&self, key: &K) -> DbResult<Option<ByteView>> {
2115 let i = self.position(key)?;
2116 let (sid, inner) = &self.inners[i];
2117 let Some(node) = self.tree.index.get(key.as_bytes(), &self.seize) else {
2118 return Ok(None);
2119 };
2120 let disk = node.read_loc();
2121 Ok(Some(
2122 self.tree
2123 .read_value_locked_result(&disk, *sid as u8, inner)?,
2124 ))
2125 }
2126
2127 pub fn try_contains(&self, key: &K) -> DbResult<bool> {
2128 self.position(key)?;
2129 Ok(self.tree.index.get(key.as_bytes(), &self.seize).is_some())
2130 }
2131
2132 pub fn get_or_err(&self, key: &K) -> DbResult<ByteView> {
2133 self.try_get(key)?.ok_or(DbError::KeyNotFound)
2134 }
2135
2136 pub fn put(&mut self, key: &K, value: &[u8]) -> DbResult<()> {
2137 let i = self.position(key)?;
2138 let old = if H::NEEDS_OLD_VALUE && H::NEEDS_WRITE {
2139 let (sid, inner) = &self.inners[i];
2140 if let Some(node) = self.tree.index.get(key.as_bytes(), &self.seize) {
2141 let disk = node.read_loc();
2142 Some(
2143 self.tree
2144 .read_value_locked_result(&disk, *sid as u8, inner)?,
2145 )
2146 } else {
2147 None
2148 }
2149 } else {
2150 None
2151 };
2152 let (sid, inner) = &mut self.inners[i];
2153 self.tree.put_locked(*sid, inner, &self.seize, key, value)?;
2154 if H::NEEDS_WRITE {
2155 self.log.push((*key, old, Some(ByteView::from(value))));
2156 }
2157 Ok(())
2158 }
2159
2160 pub fn insert(&mut self, key: &K, value: &[u8]) -> DbResult<()> {
2161 let i = self.position(key)?;
2162 let (sid, inner) = &mut self.inners[i];
2163 self.tree
2164 .insert_locked(*sid, inner, &self.seize, key, value)?;
2165 if H::NEEDS_WRITE {
2166 self.log.push((*key, None, Some(ByteView::from(value))));
2167 }
2168 Ok(())
2169 }
2170
2171 pub fn delete(&mut self, key: &K) -> DbResult<bool> {
2172 let i = self.position(key)?;
2173 let old = if H::NEEDS_OLD_VALUE && H::NEEDS_WRITE {
2174 let (sid, inner) = &self.inners[i];
2175 if let Some(node) = self.tree.index.get(key.as_bytes(), &self.seize) {
2176 let disk = node.read_loc();
2177 Some(
2178 self.tree
2179 .read_value_locked_result(&disk, *sid as u8, inner)?,
2180 )
2181 } else {
2182 None
2183 }
2184 } else {
2185 None
2186 };
2187 let (sid, inner) = &mut self.inners[i];
2188 let existed = self.tree.delete_locked(*sid, inner, &self.seize, key)?;
2189 if existed && H::NEEDS_WRITE {
2190 self.log.push((*key, old, None));
2191 }
2192 Ok(existed)
2193 }
2194}
2195
2196#[cfg(feature = "armour")]
2197impl<K: Key, H: WriteHook<K>> crate::armour::multi_tx::MultiTx for VarTree<K, H> {
2198 type Key = K;
2199 type Tx<'a>
2200 = VarTx<'a, K, H>
2201 where
2202 Self: 'a;
2203
2204 fn shard_for_key(&self, key: &K) -> usize {
2205 self.shard_for(key)
2206 }
2207
2208 fn begin_tx(&self) -> VarTx<'_, K, H> {
2209 VarTx {
2210 tree: self,
2211 inners: Vec::new(),
2212 seize: self.index.collector().enter(),
2213 log: Vec::new(),
2214 }
2215 }
2216
2217 fn lock_shard_into<'a>(&'a self, shard_id: usize, tx: &mut VarTx<'a, K, H>) {
2218 tx.inners
2219 .push((shard_id, self.engine.shards()[shard_id].lock()));
2220 }
2221
2222 fn release_locks(&self, tx: &mut VarTx<'_, K, H>) -> crate::armour::multi_tx::SyncNeeds {
2223 tx.inners.clear(); crate::armour::multi_tx::SyncNeeds::none()
2225 }
2226
2227 fn run_sync(&self, _needs: crate::armour::multi_tx::SyncNeeds) -> DbResult<()> {
2228 Ok(())
2229 }
2230
2231 fn replay_hooks(&self, tx: VarTx<'_, K, H>) {
2232 if H::NEEDS_WRITE {
2233 for (k, old, new) in &tx.log {
2234 self.hook.on_write(k, old.as_deref(), new.as_deref());
2235 }
2236 }
2237 }
2238}
2239
2240#[cfg(test)]
2241mod tests {
2242 use super::*;
2243 use crate::Config;
2244 use crate::compaction::compact_shard;
2245 use tempfile::tempdir;
2246
2247 use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering};
2248
2249 #[derive(Default)]
2252 struct CountingHook<const NEEDS_INIT: bool, const NEEDS_OLD: bool> {
2253 writes: AtomicUsize,
2254 writes_with_old: AtomicUsize,
2255 inits: AtomicUsize,
2256 last_write_new: crate::sync::Mutex<Option<Vec<u8>>>,
2257 last_init_value: crate::sync::Mutex<Option<Vec<u8>>>,
2258 }
2259
2260 impl<const NEEDS_INIT: bool, const NEEDS_OLD: bool> WriteHook<[u8; 8]>
2261 for CountingHook<NEEDS_INIT, NEEDS_OLD>
2262 {
2263 const NEEDS_OLD_VALUE: bool = NEEDS_OLD;
2264 const NEEDS_INIT: bool = NEEDS_INIT;
2265
2266 fn on_write(&self, _key: &[u8; 8], old: Option<&[u8]>, new: Option<&[u8]>) {
2267 self.writes.fetch_add(1, AtomicOrdering::Relaxed);
2268 if old.is_some() {
2269 self.writes_with_old.fetch_add(1, AtomicOrdering::Relaxed);
2270 }
2271 *crate::sync::lock(&self.last_write_new) = new.map(<[u8]>::to_vec);
2272 }
2273
2274 fn on_init(&self, _key: &[u8; 8], value: &[u8]) {
2275 self.inits.fetch_add(1, AtomicOrdering::Relaxed);
2276 *crate::sync::lock(&self.last_init_value) = Some(value.to_vec());
2277 }
2278 }
2279
2280 fn open_test_tree(dir: &std::path::Path) -> VarTree<[u8; 8]> {
2281 let mut cfg = Config::test();
2282 cfg.shard_count = 1;
2283 cfg.max_file_size = 8192;
2284 cfg.write_buffer_size = 8192;
2285 cfg.compaction_threshold = 0.0;
2286 VarTree::open(dir, cfg).expect("open test tree")
2287 }
2288
2289 fn open_test_tree_hooked<const NEEDS_INIT: bool, const NEEDS_OLD: bool>(
2290 dir: &std::path::Path,
2291 hook: CountingHook<NEEDS_INIT, NEEDS_OLD>,
2292 ) -> VarTree<[u8; 8], CountingHook<NEEDS_INIT, NEEDS_OLD>> {
2293 let mut cfg = Config::test();
2294 cfg.shard_count = 1;
2295 cfg.max_file_size = 8192;
2296 cfg.write_buffer_size = 8192;
2297 cfg.compaction_threshold = 0.0;
2298 VarTree::open_hooked(dir, cfg, hook).expect("open hooked test tree")
2299 }
2300
2301 #[test]
2306 fn var_compaction_reclaims_dead_bytes_under_overwrite_churn() {
2307 let dir = tempdir().unwrap();
2308 let mut cfg = Config::test();
2309 cfg.shard_count = 1;
2310 cfg.max_file_size = 64 * 1024;
2311 cfg.write_buffer_size = 16 * 1024;
2312 cfg.compaction_threshold = 0.30;
2313 let max_file_size = cfg.max_file_size;
2314 let tree = VarTree::<[u8; 8]>::open(dir.path(), cfg).unwrap();
2315
2316 const N: u64 = 500;
2317 const M: u64 = 40;
2318 let entry_sz: u64 = 16 + 8 + 64;
2319 let live_set = N * entry_sz;
2320 let total_written = N * (M + 1) * entry_sz;
2321
2322 for k in 0..N {
2323 tree.put(&k.to_be_bytes(), &[0u8; 64]).unwrap();
2324 }
2325 for round in 1..=M {
2326 let mut v = [0u8; 64];
2327 v[0] = round as u8;
2328 for k in 0..N {
2329 tree.put(&k.to_be_bytes(), &v).unwrap();
2330 }
2331 }
2332 tree.flush_buffers().unwrap();
2333
2334 let shard = &tree.engine.shards()[0];
2335 let disk = || {
2336 let inner = shard.lock();
2337 inner.active.write_offset + inner.immutable.iter().map(|f| f.total_bytes).sum::<u64>()
2338 };
2339 assert!(disk() > total_written / 2, "expected dead-byte buildup");
2340
2341 for _ in 0..200 {
2342 let before = disk();
2343 tree.compact().unwrap();
2344 if disk() == before {
2345 break;
2346 }
2347 }
2348 assert!(
2349 disk() <= live_set + max_file_size,
2350 "VarTree compaction failed to reclaim: disk={} live_set={live_set}",
2351 disk()
2352 );
2353
2354 let mut v = [0u8; 64];
2356 v[0] = M as u8;
2357 assert_eq!(tree.get(&0u64.to_be_bytes()).as_deref(), Some(&v[..]));
2358 assert_eq!(tree.get(&(N - 1).to_be_bytes()).as_deref(), Some(&v[..]));
2359 assert_eq!(tree.len(), N as usize);
2360 }
2361
2362 fn put_until_compactable(tree: &VarTree<[u8; 8]>, key: [u8; 8]) -> DiskLoc {
2363 tree.put(&key, &[0u8; 256]).expect("first put");
2367 let snap = {
2368 let guard = tree.index.collector().enter();
2369 let node = tree.index.get(key.as_bytes(), &guard).expect("indexed");
2370 node.read_loc()
2371 };
2372 for i in 1..65u8 {
2376 tree.put(&key, &[i; 256]).expect("put");
2377 }
2378 tree.put(&key, b"final-value-payload-XX")
2379 .expect("final put");
2380 snap
2381 }
2382
2383 #[test]
2384 fn put_reports_whether_key_existed() {
2385 let dir = tempdir().unwrap();
2386 let tree = open_test_tree(dir.path());
2387 let key = 42u64.to_be_bytes();
2388
2389 assert!(!tree.put(&key, b"v1").expect("put"));
2391 assert!(tree.put(&key, b"v2").expect("put"));
2393 assert!(tree.put(&key, b"v3").expect("put"));
2394 tree.flush_buffers().expect("flush");
2396 assert!(tree.put(&key, b"v4").expect("put"));
2397 assert!(!tree.put(&7u64.to_be_bytes(), b"x").expect("put"));
2399 }
2400
2401 #[test]
2402 fn warmup_threads_correct_shard_across_shards() {
2403 let dir = tempdir().unwrap();
2404 let tree = open_test_tree(dir.path());
2405 for i in 0..256u64 {
2407 tree.put(&i.to_be_bytes(), &[i as u8; 64]).expect("put");
2408 }
2409 tree.flush_buffers().expect("flush");
2410 tree.warmup().expect("warmup");
2412 for i in 0..256u64 {
2414 assert_eq!(
2415 tree.get(&i.to_be_bytes()).as_deref(),
2416 Some(&[i as u8; 64][..])
2417 );
2418 }
2419 }
2420
2421 #[test]
2422 fn read_value_cached_inner_returns_stale_after_compaction() {
2423 let dir = tempdir().unwrap();
2424 let tree = open_test_tree(dir.path());
2425
2426 let key = 7u64.to_be_bytes();
2427 let snap = put_until_compactable(&tree, key);
2428 let shard_id = tree.shard_for(&key) as u8;
2429
2430 let shard = &tree.engine.shards()[shard_id as usize];
2431 let _ = compact_shard(shard, &tree, 0.0).expect("compaction");
2432
2433 match tree.read_value_cached_inner(&snap, shard_id) {
2434 Err(DbError::StaleDiskLoc) => {}
2435 Ok(v) => panic!("expected StaleDiskLoc, got Ok({:?})", v.as_bytes()),
2436 Err(e) => panic!("expected StaleDiskLoc, got Err({e})"),
2437 }
2438 }
2439
2440 #[test]
2441 fn read_value_cached_returns_some_after_compaction() {
2442 let dir = tempdir().unwrap();
2443 let tree = open_test_tree(dir.path());
2444
2445 let key = 11u64.to_be_bytes();
2446 let _snap = put_until_compactable(&tree, key);
2447 let shard = &tree.engine.shards()[0];
2448 let _ = compact_shard(shard, &tree, 0.0).expect("compaction");
2449
2450 let guard = tree.index.collector().enter();
2451 let node = tree.index.get(key.as_bytes(), &guard).expect("indexed");
2452 let v = tree
2453 .read_value_cached(node, &guard)
2454 .expect("post-compaction value must be readable");
2455 assert_eq!(v.as_bytes(), b"final-value-payload-XX");
2456 }
2457
2458 #[test]
2459 fn get_during_compaction_returns_some() {
2460 let dir = tempdir().unwrap();
2461 let tree = open_test_tree(dir.path());
2462
2463 let key = 13u64.to_be_bytes();
2464 let _snap = put_until_compactable(&tree, key);
2465 let shard = &tree.engine.shards()[0];
2466 let _ = compact_shard(shard, &tree, 0.0).expect("compaction");
2467
2468 let v = tree.get(&key).expect("post-compaction get");
2469 assert_eq!(v.as_bytes(), b"final-value-payload-XX");
2470 }
2471
2472 #[test]
2473 fn iter_during_compaction_yields_all_live_keys() {
2474 let dir = tempdir().unwrap();
2475 let tree = open_test_tree(dir.path());
2476
2477 for k in 1u64..=3 {
2478 put_until_compactable(&tree, k.to_be_bytes());
2479 }
2480 let shard = &tree.engine.shards()[0];
2481 let _ = compact_shard(shard, &tree, 0.0).expect("compaction");
2482
2483 let collected: std::collections::BTreeMap<[u8; 8], Vec<u8>> = tree
2484 .iter()
2485 .map(|(k, v)| (k, v.as_bytes().to_vec()))
2486 .collect();
2487 assert_eq!(collected.len(), 3);
2488 for k in 1u64..=3 {
2489 let bytes = collected
2490 .get(&k.to_be_bytes())
2491 .expect("every original key must remain");
2492 assert_eq!(bytes.as_slice(), b"final-value-payload-XX");
2493 }
2494 }
2495
2496 #[test]
2497 fn get_or_read_block_returns_stale_for_unknown_file_id() {
2498 let dir = tempdir().unwrap();
2499 let tree = open_test_tree(dir.path());
2500
2501 match tree.get_or_read_block(0, 9999, 0) {
2502 Err(DbError::StaleDiskLoc) => {}
2503 Ok(_) => panic!("expected StaleDiskLoc, got Ok"),
2504 Err(e) => panic!("expected StaleDiskLoc, got Err({e})"),
2505 }
2506 }
2507
2508 #[test]
2509 fn extract_from_block_propagates_next_block_error() {
2510 let block = AlignedBuf::zeroed(4096);
2511 let start = 4090;
2512 let len = 32;
2513 let result: DbResult<ByteView> =
2514 VarTree::<[u8; 8]>::extract_from_block(&block, start, len, || {
2515 Err(DbError::StaleDiskLoc)
2516 });
2517 match result {
2518 Err(DbError::StaleDiskLoc) => {}
2519 Ok(_) => panic!("expected StaleDiskLoc, got Ok"),
2520 Err(e) => panic!("expected StaleDiskLoc, got Err({e})"),
2521 }
2522 }
2523
2524 #[test]
2525 fn extract_from_block_multi_block_first_cached_second_stale() {
2526 let dir = tempdir().unwrap();
2527 let tree = open_test_tree(dir.path());
2528
2529 let key = 21u64.to_be_bytes();
2530 let value = vec![0xCDu8; 4073];
2531 tree.put(&key, &value).expect("initial put");
2532
2533 let snap = {
2535 let guard = tree.index.collector().enter();
2536 let node = tree.index.get(key.as_bytes(), &guard).expect("indexed");
2537 node.read_loc()
2538 };
2539
2540 for i in 0..65u8 {
2545 tree.put(&key, &[i; 256]).expect("overwrite");
2546 }
2547 tree.put(&key, b"final").expect("final");
2548
2549 let shard = &tree.engine.shards()[0];
2550 let _ = compact_shard(shard, &tree, 0.0).expect("compaction");
2551
2552 let block_offset = snap.offset as u64 & !4095;
2553 let shard_id = tree.shard_for(&key) as u8;
2554 let first_block = AlignedBuf::zeroed(4096);
2555 let cache_key = BlockKey {
2556 shard_id,
2557 file_id: snap.file_id,
2558 block_offset,
2559 };
2560 tree.block_cache.insert(cache_key, Arc::new(first_block));
2561
2562 match tree.read_value_cached_inner(&snap, shard_id) {
2563 Err(DbError::StaleDiskLoc) => {}
2564 Ok(_) => panic!("expected StaleDiskLoc from second block, got Ok"),
2565 Err(e) => panic!("expected StaleDiskLoc, got Err({e})"),
2566 }
2567 }
2568
2569 #[test]
2570 fn retry_limit_returns_none_on_persistent_stale() {
2571 let dir = tempdir().unwrap();
2572 let tree = open_test_tree(dir.path());
2573
2574 let key = 99u64.to_be_bytes();
2575 tree.put(&key, b"payload").expect("put");
2576
2577 let guard = tree.index.collector().enter();
2583 assert!(tree.index.get(key.as_bytes(), &guard).is_some());
2584 drop(guard);
2585
2586 {
2587 let shard_id = tree.shard_for(&key) as u8;
2588 let shard = &tree.engine.shards()[shard_id as usize];
2589 let mut inner = shard.lock();
2590 inner.immutable = Vec::new();
2591 inner.active.file_id = u32::MAX;
2595 }
2596
2597 let guard = tree.index.collector().enter();
2598 let node = tree
2599 .index
2600 .get(key.as_bytes(), &guard)
2601 .expect("still indexed");
2602 assert!(
2603 tree.read_value_cached(node, &guard).is_none(),
2604 "MAX_STALE_RETRIES must terminate the loop and return None"
2605 );
2606 }
2607
2608 #[test]
2609 fn var_tree_replay_init_fires_on_init_per_live_key_raw() {
2610 let dir = tempdir().unwrap();
2611 let tree = open_test_tree_hooked::<true, false>(dir.path(), CountingHook::default());
2612
2613 for i in 0u64..5 {
2614 tree.put(&i.to_be_bytes(), &[i as u8; 16]).expect("put");
2615 }
2616 tree.delete(&3u64.to_be_bytes()).expect("delete");
2618
2619 tree.hook.writes.store(0, AtomicOrdering::Relaxed);
2621 tree.hook.inits.store(0, AtomicOrdering::Relaxed);
2622
2623 tree.replay_init();
2624
2625 assert_eq!(
2626 tree.hook.inits.load(AtomicOrdering::Relaxed),
2627 4,
2628 "4 live keys"
2629 );
2630 assert_eq!(
2631 tree.hook.writes.load(AtomicOrdering::Relaxed),
2632 0,
2633 "no on_write"
2634 );
2635 }
2636
2637 #[test]
2638 fn var_tree_replay_init_no_hook_is_noop() {
2639 let dir = tempdir().unwrap();
2640 let tree = open_test_tree(dir.path());
2641
2642 for i in 0u64..3 {
2643 tree.put(&i.to_be_bytes(), &[i as u8; 8]).expect("put");
2644 }
2645 tree.replay_init();
2647 assert!(tree.get(&0u64.to_be_bytes()).is_some());
2649 }
2650
2651 #[test]
2652 fn var_tree_migrate_keep_fires_on_init_not_on_write_raw() {
2653 use crate::MigrateAction;
2654 let dir = tempdir().unwrap();
2655 let tree = open_test_tree_hooked::<true, false>(dir.path(), CountingHook::default());
2656
2657 for i in 0u64..4 {
2658 tree.put(&i.to_be_bytes(), &[i as u8; 16]).expect("put");
2659 }
2660 tree.hook.writes.store(0, AtomicOrdering::Relaxed);
2661 tree.hook.inits.store(0, AtomicOrdering::Relaxed);
2662
2663 let mutated = tree.migrate(|_, _| MigrateAction::Keep).expect("migrate");
2664
2665 assert_eq!(mutated, 0);
2666 assert_eq!(
2667 tree.hook.inits.load(AtomicOrdering::Relaxed),
2668 4,
2669 "4 keeps -> 4 on_init"
2670 );
2671 assert_eq!(
2672 tree.hook.writes.load(AtomicOrdering::Relaxed),
2673 0,
2674 "Keep must not fire on_write"
2675 );
2676 }
2677
2678 #[test]
2679 fn var_tree_migrate_update_fires_on_init_with_new_value_raw() {
2680 use crate::MigrateAction;
2681 let dir = tempdir().unwrap();
2682 let tree = open_test_tree_hooked::<true, false>(dir.path(), CountingHook::default());
2683
2684 let key = 42u64.to_be_bytes();
2685 tree.put(&key, b"old-value").expect("put");
2686 tree.hook.writes.store(0, AtomicOrdering::Relaxed);
2687 tree.hook.inits.store(0, AtomicOrdering::Relaxed);
2688
2689 let new = ByteView::new(b"new-value");
2690 let mutated = tree
2691 .migrate(move |_, _| MigrateAction::Update(new.clone()))
2692 .expect("migrate");
2693
2694 assert_eq!(mutated, 1);
2695 assert_eq!(tree.hook.inits.load(AtomicOrdering::Relaxed), 1);
2696 assert_eq!(
2697 crate::sync::lock(&tree.hook.last_init_value).as_deref(),
2698 Some(b"new-value".as_ref()),
2699 "on_init must receive the NEW value"
2700 );
2701 assert_eq!(
2702 tree.hook.writes.load(AtomicOrdering::Relaxed),
2703 0,
2704 "Update must NOT fire on_write (was double-firing through self.put)"
2705 );
2706 assert_eq!(tree.get(&key).unwrap().as_bytes(), b"new-value");
2707 }
2708
2709 #[test]
2710 fn var_tree_migrate_delete_fires_no_hooks_raw() {
2711 use crate::MigrateAction;
2712 let dir = tempdir().unwrap();
2713 let tree = open_test_tree_hooked::<true, false>(dir.path(), CountingHook::default());
2714
2715 let key = 7u64.to_be_bytes();
2716 tree.put(&key, b"x").expect("put");
2717 tree.hook.writes.store(0, AtomicOrdering::Relaxed);
2718 tree.hook.inits.store(0, AtomicOrdering::Relaxed);
2719
2720 let mutated = tree.migrate(|_, _| MigrateAction::Delete).expect("migrate");
2721
2722 assert_eq!(mutated, 1);
2723 assert_eq!(tree.hook.inits.load(AtomicOrdering::Relaxed), 0);
2724 assert_eq!(tree.hook.writes.load(AtomicOrdering::Relaxed), 0);
2725 assert!(tree.get(&key).is_none());
2726 }
2727
2728 #[test]
2729 fn var_tree_migrate_no_init_hook_is_silent_for_keep_and_update() {
2730 use crate::MigrateAction;
2731 let dir = tempdir().unwrap();
2732 let tree = open_test_tree_hooked::<false, false>(dir.path(), CountingHook::default());
2734
2735 for i in 0u64..3 {
2736 tree.put(&i.to_be_bytes(), &[i as u8; 16]).expect("put");
2737 }
2738 tree.hook.writes.store(0, AtomicOrdering::Relaxed);
2739 tree.hook.inits.store(0, AtomicOrdering::Relaxed);
2740
2741 tree.migrate(|_, _| MigrateAction::Keep)
2743 .expect("migrate keep");
2744 assert_eq!(
2745 tree.hook.inits.load(AtomicOrdering::Relaxed),
2746 0,
2747 "Keep with NEEDS_INIT=false"
2748 );
2749 assert_eq!(tree.hook.writes.load(AtomicOrdering::Relaxed), 0);
2750
2751 let new = ByteView::new(b"new");
2753 tree.migrate(move |_, _| MigrateAction::Update(new.clone()))
2754 .expect("migrate update");
2755 assert_eq!(
2756 tree.hook.inits.load(AtomicOrdering::Relaxed),
2757 0,
2758 "Update with NEEDS_INIT=false"
2759 );
2760 assert_eq!(
2761 tree.hook.writes.load(AtomicOrdering::Relaxed),
2762 0,
2763 "Update must not fire on_write either"
2764 );
2765 }
2766
2767 #[test]
2768 fn var_tree_public_put_still_fires_on_write_once() {
2769 let dir = tempdir().unwrap();
2770 let tree = open_test_tree_hooked::<true, false>(dir.path(), CountingHook::default());
2771
2772 tree.put(&1u64.to_be_bytes(), b"v").expect("put");
2773 assert_eq!(tree.hook.writes.load(AtomicOrdering::Relaxed), 1);
2774 }
2775
2776 #[test]
2777 fn var_tree_atomic_fires_hooks() {
2778 let dir = tempdir().unwrap();
2779 let tree = open_test_tree_hooked::<true, true>(dir.path(), CountingHook::default());
2780 let key = 1u64.to_be_bytes();
2781 tree.atomic(&key, |shard| {
2782 shard.put(&key, b"a")?; shard.delete(&key)?; Ok(())
2785 })
2786 .expect("atomic");
2787 assert_eq!(tree.hook.writes.load(AtomicOrdering::Relaxed), 2);
2788 assert_eq!(tree.hook.writes_with_old.load(AtomicOrdering::Relaxed), 1); }
2790
2791 #[test]
2792 fn var_tree_atomic_fires_for_applied_on_err() {
2793 let dir = tempdir().unwrap();
2794 let tree = open_test_tree_hooked::<true, true>(dir.path(), CountingHook::default());
2795 let key = 1u64.to_be_bytes();
2796 let r: DbResult<()> = tree.atomic(&key, |shard| {
2797 shard.put(&key, b"x")?;
2798 Err(DbError::KeyNotFound)
2799 });
2800 assert!(r.is_err());
2801 assert_eq!(tree.hook.writes.load(AtomicOrdering::Relaxed), 1);
2802 }
2803
2804 #[test]
2805 fn var_tree_atomic_same_key_old_from_disk() {
2806 let dir = tempdir().unwrap();
2809 let tree = open_test_tree_hooked::<false, true>(dir.path(), CountingHook::default());
2810 let key = 1u64.to_be_bytes();
2811 tree.atomic(&key, |shard| {
2812 shard.put(&key, b"first")?;
2813 shard.put(&key, b"second")?;
2814 Ok(())
2815 })
2816 .expect("atomic");
2817 assert_eq!(tree.hook.writes.load(AtomicOrdering::Relaxed), 2);
2818 assert_eq!(tree.hook.writes_with_old.load(AtomicOrdering::Relaxed), 1); assert_eq!(
2820 crate::sync::lock(&tree.hook.last_write_new).clone(),
2821 Some(b"second".to_vec())
2822 );
2823 }
2824
2825 #[test]
2826 fn var_tree_shard_put_returns_existed() {
2827 let dir = tempdir().unwrap();
2828 let tree = open_test_tree_hooked::<true, true>(dir.path(), CountingHook::default());
2829 let key = 1u64.to_be_bytes();
2830 let (fresh, overwrite) = tree
2831 .atomic(&key, |shard| {
2832 let fresh = shard.put(&key, b"a")?; let overwrite = shard.put(&key, b"b")?; Ok((fresh, overwrite))
2835 })
2836 .expect("atomic");
2837 assert!(!fresh, "fresh insert must return false");
2838 assert!(overwrite, "overwrite must return true");
2839 }
2840
2841 #[test]
2842 fn compare_delete_match_mismatch_absent() {
2843 let dir = tempdir().unwrap();
2844 let tree = open_test_tree(dir.path());
2845
2846 let k = [1u8; 8];
2847 tree.put(&k, b"hello").unwrap();
2848
2849 assert!(matches!(
2850 tree.compare_delete(&k, b"WRONG"),
2851 Err(DbError::CasMismatch)
2852 ));
2853 assert_eq!(tree.get(&k).as_deref(), Some(&b"hello"[..]));
2854
2855 assert!(tree.compare_delete(&k, b"hello").is_ok());
2856 assert!(tree.get(&k).is_none());
2857
2858 assert!(matches!(
2859 tree.compare_delete(&k, b"hello"),
2860 Err(DbError::KeyNotFound)
2861 ));
2862 }
2863
2864 #[test]
2865 fn compare_delete_hook_needs_old_true_gets_old_value() {
2866 let dir = tempdir().unwrap();
2867 let tree = open_test_tree_hooked::<false, true>(dir.path(), CountingHook::default());
2868
2869 let k = [2u8; 8];
2870 tree.put(&k, b"payload").unwrap();
2871 let with_old_before = tree.hook.writes_with_old.load(AtomicOrdering::Relaxed);
2872
2873 tree.compare_delete(&k, b"payload").unwrap();
2874
2875 assert_eq!(
2877 tree.hook.writes_with_old.load(AtomicOrdering::Relaxed) - with_old_before,
2878 1
2879 );
2880 assert_eq!(*crate::sync::lock(&tree.hook.last_write_new), None);
2881 }
2882
2883 #[test]
2884 fn compare_delete_hook_needs_old_false_gets_none() {
2885 let dir = tempdir().unwrap();
2886 let tree = open_test_tree_hooked::<false, false>(dir.path(), CountingHook::default());
2887
2888 let k = [3u8; 8];
2889 tree.put(&k, b"payload").unwrap();
2890 let writes_before = tree.hook.writes.load(AtomicOrdering::Relaxed);
2891 let with_old_before = tree.hook.writes_with_old.load(AtomicOrdering::Relaxed);
2892
2893 tree.compare_delete(&k, b"payload").unwrap();
2894
2895 assert_eq!(
2897 tree.hook.writes.load(AtomicOrdering::Relaxed) - writes_before,
2898 1
2899 );
2900 assert_eq!(
2901 tree.hook.writes_with_old.load(AtomicOrdering::Relaxed) - with_old_before,
2902 0
2903 );
2904 }
2905
2906 #[test]
2908 fn read_value_locked_result_ok_from_write_buf() {
2909 let dir = tempdir().unwrap();
2910 let tree = open_test_tree(dir.path());
2911
2912 let key = 1u64.to_be_bytes();
2913 let payload = b"in-write-buffer-value";
2914 tree.put(&key, payload).expect("put");
2915
2916 let guard = tree.index.collector().enter();
2917 let node = tree.index.get(key.as_bytes(), &guard).expect("indexed");
2918 let disk = node.read_loc();
2919 drop(guard);
2920
2921 let shard_id = tree.shard_for(&key) as u8;
2922 let shard = &tree.engine.shards()[shard_id as usize];
2923 let inner = shard.lock();
2924 let v = tree
2925 .read_value_locked_result(&disk, shard_id, &inner)
2926 .expect("write-buf read must succeed");
2927 assert_eq!(v.as_bytes(), payload);
2928 }
2929
2930 #[test]
2933 fn read_value_locked_result_ok_from_disk_immutable() {
2934 let dir = tempdir().unwrap();
2935 let tree = open_test_tree(dir.path());
2936
2937 let key = 2u64.to_be_bytes();
2938 let payload = b"single-block-immutable";
2939 tree.put(&key, payload).expect("put");
2940 for i in 100u64..135 {
2944 tree.put(&i.to_be_bytes(), &[i as u8; 256])
2945 .expect("rotator");
2946 }
2947
2948 let guard = tree.index.collector().enter();
2949 let node = tree.index.get(key.as_bytes(), &guard).expect("indexed");
2950 let disk = node.read_loc();
2951 drop(guard);
2952
2953 let shard_id = tree.shard_for(&key) as u8;
2958 let shard = &tree.engine.shards()[shard_id as usize];
2959 let inner = shard.lock();
2960 assert_ne!(
2963 disk.file_id, inner.active.file_id,
2964 "test setup failed: key=2 entry is still in the active file's write buffer",
2965 );
2966 let v = tree
2967 .read_value_locked_result(&disk, shard_id, &inner)
2968 .expect("disk read must succeed");
2969 assert_eq!(v.as_bytes(), payload);
2970 }
2971
2972 #[test]
2974 fn read_value_locked_result_propagates_stale_disk_loc() {
2975 let dir = tempdir().unwrap();
2976 let tree = open_test_tree(dir.path());
2977
2978 let key = 3u64.to_be_bytes();
2979 let snap = put_until_compactable(&tree, key);
2980
2981 let shard_id = tree.shard_for(&key) as u8;
2982 let shard = &tree.engine.shards()[shard_id as usize];
2983 let _ = compact_shard(shard, &tree, 0.0).expect("compaction");
2984
2985 let inner = shard.lock();
2986 match tree.read_value_locked_result(&snap, shard_id, &inner) {
2987 Err(DbError::StaleDiskLoc) => {}
2988 Ok(v) => panic!("expected StaleDiskLoc, got Ok({:?})", v.as_bytes()),
2989 Err(e) => panic!("expected StaleDiskLoc, got Err({e})"),
2990 }
2991 }
2992
2993 #[test]
3003 fn large_value_with_first_block_cached_uses_fallback() {
3004 let dir = tempdir().unwrap();
3005 let mut cfg = Config::test();
3009 cfg.shard_count = 1;
3010 cfg.max_file_size = 128 * 1024;
3011 cfg.write_buffer_size = 128 * 1024;
3012 cfg.compaction_threshold = 0.0;
3013 cfg.block_cache.max_size = 1 << 20;
3017 let tree: VarTree<[u8; 8]> = VarTree::open(dir.path(), cfg).expect("open");
3018
3019 let key = 42u64.to_be_bytes();
3020 let payload: Vec<u8> = (0..20_000u32).map(|i| i as u8).collect();
3022 tree.put(&key, &payload).expect("put large");
3023 tree.engine.shards()[0]
3025 .rotate_active_for_test(8)
3026 .expect("rotate");
3027
3028 let guard = tree.index.collector().enter();
3029 let node = tree.index.get(key.as_bytes(), &guard).expect("indexed");
3030 let disk = node.read_loc();
3031 drop(guard);
3032
3033 let start = (disk.offset & 4095) as usize;
3035 assert!(
3036 start + disk.len as usize > 8192,
3037 "test precondition: large value must span >2 blocks",
3038 );
3039
3040 let block_offset = disk.offset as u64 & !4095;
3044 let shard_id = tree.shard_for(&key) as u8;
3045 let cache_key = BlockKey {
3046 shard_id,
3047 file_id: disk.file_id,
3048 block_offset,
3049 };
3050 tree.block_cache
3051 .insert(cache_key, Arc::new(AlignedBuf::zeroed(4096)));
3052 assert!(
3054 tree.block_cache.get(&cache_key).is_some(),
3055 "cache must be enabled for warm-cache scenario",
3056 );
3057
3058 let v = tree
3059 .read_value_cached_inner(&disk, shard_id)
3060 .expect("large value must read via locked fallback");
3061 assert_eq!(v.as_bytes(), payload.as_slice());
3062 }
3063
3064 #[test]
3066 fn extract_from_block_single_block() {
3067 let mut block = AlignedBuf::zeroed(4096);
3068 for (i, byte) in block.iter_mut().enumerate() {
3069 *byte = i as u8;
3070 }
3071 let v = VarTree::<[u8; 8]>::extract_from_block(&block, 100, 50, || {
3072 panic!("next_block must not be called for single-block reads")
3073 })
3074 .expect("ok");
3075 let expected: Vec<u8> = (100u8..150u8).collect();
3076 assert_eq!(v.as_bytes(), expected.as_slice());
3077 }
3078
3079 #[test]
3081 fn extract_from_block_two_blocks_exact() {
3082 let mut first = AlignedBuf::zeroed(4096);
3083 for byte in first.iter_mut() {
3084 *byte = 0xAA;
3085 }
3086 let mut second = AlignedBuf::zeroed(4096);
3087 for byte in second.iter_mut() {
3088 *byte = 0xBB;
3089 }
3090 let v = VarTree::<[u8; 8]>::extract_from_block(&first, 4095, 4097, || Ok(Arc::new(second)))
3093 .expect("ok");
3094 let bytes = v.as_bytes();
3095 assert_eq!(bytes.len(), 4097);
3096 assert_eq!(bytes[0], 0xAA);
3097 assert_eq!(bytes[1], 0xBB);
3098 assert_eq!(bytes[4096], 0xBB);
3099 }
3100
3101 #[test]
3103 fn extract_from_block_two_blocks_partial() {
3104 let mut first = AlignedBuf::zeroed(4096);
3105 for byte in first.iter_mut() {
3106 *byte = 0x11;
3107 }
3108 let mut second = AlignedBuf::zeroed(4096);
3109 for byte in second.iter_mut() {
3110 *byte = 0x22;
3111 }
3112 let v = VarTree::<[u8; 8]>::extract_from_block(&first, 4000, 200, || Ok(Arc::new(second)))
3114 .expect("ok");
3115 let bytes = v.as_bytes();
3116 assert_eq!(bytes.len(), 200);
3117 assert!(bytes[..96].iter().all(|b| *b == 0x11));
3118 assert!(bytes[96..].iter().all(|b| *b == 0x22));
3119 }
3120
3121 fn open_large_value_tree(dir: &std::path::Path) -> VarTree<[u8; 8]> {
3126 let mut cfg = Config::test();
3127 cfg.shard_count = 1;
3128 cfg.max_file_size = 128 * 1024;
3129 cfg.write_buffer_size = 128 * 1024;
3130 cfg.compaction_threshold = 0.0;
3131 VarTree::open(dir, cfg).expect("open large-value test tree")
3132 }
3133
3134 fn build_large_payload(seed: u8) -> Vec<u8> {
3135 (0..50_000u32)
3136 .map(|i| (i as u8).wrapping_add(seed))
3137 .collect()
3138 }
3139
3140 fn open_value_cache_tree(dir: &std::path::Path) -> VarTree<[u8; 8]> {
3141 let mut cfg = Config::test();
3142 cfg.shard_count = 1;
3143 cfg.max_file_size = 128 * 1024;
3144 cfg.write_buffer_size = 128 * 1024;
3145 cfg.compaction_threshold = 0.0;
3146 cfg.value_cache.max_size = 4 << 20; VarTree::open(dir, cfg).expect("open value-cache test tree")
3148 }
3149
3150 #[test]
3151 fn large_value_miss_then_cached() {
3152 let dir = tempdir().unwrap();
3153 let tree = open_value_cache_tree(dir.path());
3154 let key = 1u64.to_be_bytes();
3155 let payload = build_large_payload(0x33);
3156 tree.put(&key, &payload).unwrap();
3157
3158 assert_eq!(tree.get(&key).unwrap().as_bytes(), &payload[..]);
3160
3161 let guard = tree.index.collector().enter();
3163 let node = tree.index.get(key.as_bytes(), &guard).unwrap();
3164 let disk = node.read_loc();
3165 let vkey = crate::value_cache::ValueKey {
3166 shard_id: 0,
3167 file_id: disk.file_id,
3168 offset: disk.offset,
3169 };
3170 assert!(
3171 tree.value_cache.get(&vkey).is_some(),
3172 "value cached after read"
3173 );
3174 drop(guard);
3175
3176 assert_eq!(tree.get(&key).unwrap().as_bytes(), &payload[..]);
3178 }
3179
3180 #[test]
3181 fn large_value_read_then_edit_returns_new() {
3182 let dir = tempdir().unwrap();
3183 let tree = open_value_cache_tree(dir.path());
3184 let key = 7u64.to_be_bytes();
3185 let v1 = build_large_payload(0x11);
3186 let v2 = build_large_payload(0x22);
3187 tree.put(&key, &v1).unwrap();
3188 assert_eq!(tree.get(&key).unwrap().as_bytes(), &v1[..]); tree.put(&key, &v2).unwrap(); assert_eq!(tree.get(&key).unwrap().as_bytes(), &v2[..]);
3192 }
3193
3194 #[test]
3195 fn large_value_read_then_compact_returns_live() {
3196 let dir = tempdir().unwrap();
3197 let tree = open_value_cache_tree(dir.path());
3198 let key = 102u64.to_be_bytes();
3199 let payload = build_large_payload(0x42);
3200 tree.put(&key, &payload).unwrap();
3201 assert_eq!(tree.get(&key).unwrap().as_bytes(), &payload[..]); tree.engine.shards()[0]
3204 .rotate_active_for_test(8)
3205 .expect("rotate");
3206 for i in 1..20u8 {
3207 tree.put(&key, &[i; 256]).unwrap();
3208 }
3209 tree.put(&key, b"live-after-compaction").unwrap();
3210 let shard = &tree.engine.shards()[0];
3211 let _ = compact_shard(shard, &tree, 0.0).expect("compaction");
3212
3213 assert_eq!(tree.get(&key).unwrap().as_bytes(), b"live-after-compaction");
3215 }
3216
3217 #[test]
3218 fn large_value_disabled_cache_reads_correctly() {
3219 let dir = tempdir().unwrap();
3220 let tree = open_large_value_tree(dir.path()); let key = 3u64.to_be_bytes();
3222 let payload = build_large_payload(0x55);
3223 tree.put(&key, &payload).unwrap();
3224 assert_eq!(tree.get(&key).unwrap().as_bytes(), &payload[..]);
3225 assert_eq!(tree.get(&key).unwrap().as_bytes(), &payload[..]);
3226 }
3227
3228 #[test]
3235 fn large_value_read_via_locked_fallback() {
3236 let dir = tempdir().unwrap();
3237 let tree = open_large_value_tree(dir.path());
3238
3239 let key = 100u64.to_be_bytes();
3240 let payload = build_large_payload(0);
3241 tree.put(&key, &payload).expect("put large");
3242 tree.engine.shards()[0]
3243 .rotate_active_for_test(8)
3244 .expect("rotate");
3245
3246 let guard = tree.index.collector().enter();
3247 let node = tree.index.get(key.as_bytes(), &guard).expect("indexed");
3248 let v = tree
3249 .read_value_cached(node, &guard)
3250 .expect("read must succeed via locked fallback");
3251 assert_eq!(v.as_bytes(), payload.as_slice());
3252 }
3253
3254 #[cfg(feature = "encryption")]
3255 #[test]
3256 fn get_or_err_and_try_get_surface_read_fault() {
3257 use crate::test_faults::{big_value, corrupt_tags};
3258 let dir = tempfile::tempdir().unwrap();
3259 let mut cfg = Config::test();
3260 cfg.shard_count = 1;
3261 cfg.max_file_size = 128 * 1024;
3262 cfg.write_buffer_size = 128 * 1024;
3263 cfg.compaction_threshold = 0.0;
3264 cfg.encryption_key = Some([7u8; 32]);
3265 let tree: VarTree<[u8; 8]> = VarTree::open(dir.path(), cfg).expect("open enc");
3266
3267 let key = 1u64.to_be_bytes();
3268 tree.put(&key, &big_value(0xAB)).expect("put");
3269 tree.engine.shards()[0]
3270 .rotate_active_for_test(8)
3271 .expect("rotate");
3272 corrupt_tags(&tree.engine.shard_dirs()[0]);
3273
3274 assert!(tree.get(&key).is_none());
3276 assert!(matches!(
3278 tree.get_or_err(&key),
3279 Err(DbError::EncryptionError(_))
3280 ));
3281 assert!(matches!(
3282 tree.try_get(&key),
3283 Err(DbError::EncryptionError(_))
3284 ));
3285 let absent = 999u64.to_be_bytes();
3287 assert!(matches!(
3288 tree.get_or_err(&absent),
3289 Err(DbError::KeyNotFound)
3290 ));
3291 assert!(matches!(tree.try_get(&absent), Ok(None)));
3292 }
3293
3294 #[cfg(feature = "encryption")]
3295 #[test]
3296 fn mutators_surface_read_fault_not_keynotfound() {
3297 use crate::test_faults::{big_value, corrupt_tags};
3298 let dir = tempfile::tempdir().unwrap();
3299 let mut cfg = Config::test();
3300 cfg.shard_count = 1;
3301 cfg.max_file_size = 128 * 1024;
3302 cfg.write_buffer_size = 128 * 1024;
3303 cfg.compaction_threshold = 0.0;
3304 cfg.encryption_key = Some([7u8; 32]);
3305 let tree: VarTree<[u8; 8]> = VarTree::open(dir.path(), cfg).expect("open enc");
3306
3307 let key = 1u64.to_be_bytes();
3308 tree.put(&key, &big_value(0xAB)).unwrap();
3309 tree.engine.shards()[0]
3310 .rotate_active_for_test(8)
3311 .expect("rotate");
3312 corrupt_tags(&tree.engine.shard_dirs()[0]);
3313
3314 assert!(matches!(
3315 tree.cas(&key, b"x", b"y"),
3316 Err(DbError::EncryptionError(_))
3317 ));
3318 assert!(matches!(
3319 tree.compare_delete(&key, b"x"),
3320 Err(DbError::EncryptionError(_))
3321 ));
3322 assert!(matches!(
3323 tree.update(&key, |_| ByteView::new(b"z")),
3324 Err(DbError::EncryptionError(_))
3325 ));
3326 assert!(matches!(
3327 tree.fetch_update(&key, |_| ByteView::new(b"z")),
3328 Err(DbError::EncryptionError(_))
3329 ));
3330 }
3331
3332 #[cfg(feature = "encryption")]
3333 #[test]
3334 fn first_skips_unreadable_boundary_try_first_surfaces_it() {
3335 use crate::test_faults::{big_value, corrupt_tags};
3336 let dir = tempfile::tempdir().unwrap();
3337 let mut cfg = Config::test();
3338 cfg.shard_count = 1;
3339 cfg.max_file_size = 128 * 1024;
3340 cfg.write_buffer_size = 128 * 1024;
3341 cfg.compaction_threshold = 0.0;
3342 cfg.encryption_key = Some([7u8; 32]);
3343 let tree: VarTree<[u8; 8]> = VarTree::open(dir.path(), cfg).expect("open enc");
3345
3346 let k1 = 1u64.to_be_bytes();
3347 let k2 = 2u64.to_be_bytes();
3348 tree.put(&k1, &big_value(0x11)).unwrap();
3349 tree.put(&k2, &big_value(0x22)).unwrap();
3350 tree.engine.shards()[0]
3351 .rotate_active_for_test(8)
3352 .expect("rotate");
3353 corrupt_tags(&tree.engine.shard_dirs()[0]); assert!(tree.first().is_none());
3357 assert_eq!(
3358 tree.first().map(|(k, _)| k),
3359 tree.iter().next().map(|(k, _)| k)
3360 );
3361 assert!(matches!(tree.try_first(), Err(DbError::EncryptionError(_))));
3363 assert!(matches!(tree.try_last(), Err(DbError::EncryptionError(_))));
3364 }
3365
3366 #[cfg(feature = "encryption")]
3367 #[test]
3368 fn large_value_read_encrypted() {
3369 let dir = tempdir().unwrap();
3370 let mut cfg = Config::test();
3371 cfg.shard_count = 1;
3372 cfg.max_file_size = 128 * 1024;
3373 cfg.write_buffer_size = 128 * 1024;
3374 cfg.compaction_threshold = 0.0;
3375 cfg.encryption_key = Some([7u8; 32]);
3376 let tree: VarTree<[u8; 8]> = VarTree::open(dir.path(), cfg).expect("open enc");
3380
3381 let key = 101u64.to_be_bytes();
3382 let payload = build_large_payload(0xAB);
3383 tree.put(&key, &payload).expect("put encrypted large");
3384 tree.engine.shards()[0]
3385 .rotate_active_for_test(8)
3386 .expect("rotate");
3387
3388 let guard = tree.index.collector().enter();
3389 let node = tree.index.get(key.as_bytes(), &guard).expect("indexed");
3390 let v = tree
3391 .read_value_cached(node, &guard)
3392 .expect("encrypted large read must succeed via pread_value_encrypted");
3393 assert_eq!(v.as_bytes(), payload.as_slice());
3394 }
3395
3396 #[test]
3402 fn large_value_stale_disk_loc_deterministic() {
3403 let dir = tempdir().unwrap();
3404 let tree = open_large_value_tree(dir.path());
3405
3406 let key = 102u64.to_be_bytes();
3407 let payload = build_large_payload(0x42);
3408 tree.put(&key, &payload).expect("first large put");
3409
3410 let snap = {
3413 let guard = tree.index.collector().enter();
3414 let node = tree.index.get(key.as_bytes(), &guard).expect("indexed");
3415 node.read_loc()
3416 };
3417
3418 tree.engine.shards()[0]
3420 .rotate_active_for_test(8)
3421 .expect("rotate after large put");
3422
3423 for i in 1..20u8 {
3426 tree.put(&key, &[i; 256]).expect("overwrite");
3427 }
3428 tree.put(&key, b"live-after-compaction").expect("final put");
3429
3430 let shard_id = tree.shard_for(&key) as u8;
3431 let shard = &tree.engine.shards()[shard_id as usize];
3432 let _ = compact_shard(shard, &tree, 0.0).expect("compaction");
3433
3434 match tree.read_value_cached_inner(&snap, shard_id) {
3437 Err(DbError::StaleDiskLoc) => {}
3438 Ok(v) => panic!("expected StaleDiskLoc, got Ok({:?})", v.as_bytes()),
3439 Err(e) => panic!("expected StaleDiskLoc, got Err({e})"),
3440 }
3441
3442 let guard = tree.index.collector().enter();
3444 let node = tree.index.get(key.as_bytes(), &guard).expect("indexed");
3445 let v = tree
3446 .read_value_cached(node, &guard)
3447 .expect("public path must retry and return live value");
3448 assert_eq!(v.as_bytes(), b"live-after-compaction");
3449 }
3450
3451 #[test]
3456 fn read_value_locked_result_ok_from_cache_single_block() {
3457 let dir = tempdir().unwrap();
3458 let mut cfg = Config::test();
3460 cfg.shard_count = 1;
3461 cfg.max_file_size = 8192;
3462 cfg.write_buffer_size = 8192;
3463 cfg.compaction_threshold = 0.0;
3464 cfg.block_cache.max_size = 1 << 20;
3465 let tree: VarTree<[u8; 8]> = VarTree::open(dir.path(), cfg).expect("open");
3466
3467 let key = 9u64.to_be_bytes();
3468 let payload = b"small-single-block-value";
3469 tree.put(&key, payload).expect("put");
3470 for i in 100u64..135 {
3474 tree.put(&i.to_be_bytes(), &[i as u8; 256])
3475 .expect("rotator");
3476 }
3477
3478 let guard = tree.index.collector().enter();
3479 let node = tree.index.get(key.as_bytes(), &guard).expect("indexed");
3480 let disk = node.read_loc();
3481 drop(guard);
3482
3483 let start = (disk.offset & 4095) as usize;
3486 let len = disk.len as usize;
3487 assert!(
3488 start + len <= 4096,
3489 "test precondition: value must fit in a single block",
3490 );
3491
3492 let block_offset = disk.offset as u64 & !4095;
3496 let shard_id = tree.shard_for(&key) as u8;
3497 let cache_key = BlockKey {
3498 shard_id,
3499 file_id: disk.file_id,
3500 block_offset,
3501 };
3502 let block = tree
3503 .get_or_read_block(shard_id, disk.file_id, block_offset)
3504 .expect("read block");
3505 tree.block_cache.insert(cache_key, block);
3506 assert!(
3507 tree.block_cache.get(&cache_key).is_some(),
3508 "cache must contain the block for step 2 to fire",
3509 );
3510
3511 let shard = &tree.engine.shards()[shard_id as usize];
3512 let inner = shard.lock();
3513 assert_ne!(
3514 disk.file_id, inner.active.file_id,
3515 "test setup failed: key entry is still in the active write buffer",
3516 );
3517
3518 let v = tree
3519 .read_value_locked_result(&disk, shard_id, &inner)
3520 .expect("cache-hit read must succeed");
3521 assert_eq!(v.as_bytes(), payload);
3522 }
3523
3524 #[test]
3535 fn var_tree_get_with_file_id_above_u16() {
3536 let dir = tempdir().unwrap();
3537 let mut cfg = Config::test();
3538 cfg.shard_count = 1;
3539 cfg.max_file_size = 128 * 1024;
3540 cfg.write_buffer_size = 128 * 1024;
3541 cfg.compaction_threshold = 0.0;
3542 cfg.block_cache.max_size = 1 << 20;
3543 let tree: VarTree<[u8; 8]> = VarTree::open(dir.path(), cfg).expect("open");
3544
3545 let shard = &tree.engine.shards()[0];
3546
3547 shard.set_next_file_id(70_000);
3550 shard.rotate_active_for_test(8).expect("first rotate");
3551 assert!(
3552 shard.active_file_id() >= 70_000,
3553 "active_file_id should be >= 70_000 after rotation"
3554 );
3555
3556 let key = 42u64.to_be_bytes();
3557 let value = vec![0xC3u8; 512];
3558 tree.put(&key, &value).expect("put");
3559
3560 {
3562 let guard = tree.index.collector().enter();
3563 let node = tree.index.get(key.as_bytes(), &guard).expect("indexed");
3564 let disk = node.read_loc();
3565 assert!(
3566 disk.file_id > u16::MAX as u32,
3567 "DiskLoc.file_id must be above u16::MAX, got {}",
3568 disk.file_id,
3569 );
3570 }
3571
3572 shard.flush().expect("flush");
3575 shard.rotate_active_for_test(8).expect("second rotate");
3576
3577 let got = tree.get(&key).expect("get must return Some");
3578 assert_eq!(got.as_bytes(), value.as_slice());
3579 }
3580
3581 #[test]
3594 fn recovery_handles_file_id_above_u16() {
3595 let dir = tempdir().unwrap();
3596
3597 let mut cfg = Config::test();
3598 cfg.shard_count = 1;
3599 cfg.max_file_size = 128 * 1024;
3600 cfg.write_buffer_size = 128 * 1024;
3601 cfg.compaction_threshold = 0.0;
3602 cfg.block_cache.max_size = 1 << 20;
3603
3604 {
3606 let tree: VarTree<[u8; 8]> = VarTree::open(dir.path(), cfg.clone()).expect("open A");
3607 let shard = &tree.engine.shards()[0];
3608
3609 shard.set_next_file_id(70_000);
3611 shard
3612 .rotate_active_for_test(8)
3613 .expect("rotate to file_id 70_000");
3614 assert!(
3615 shard.active_file_id() >= 70_000,
3616 "active_file_id should be >= 70_000 after rotation"
3617 );
3618
3619 for i in 0u64..4 {
3620 let key = i.to_be_bytes();
3621 let value = vec![i as u8; 200];
3622 tree.put(&key, &value).expect("put phase A");
3623 }
3624
3625 tree.close().expect("close phase A");
3626 }
3627
3628 {
3630 let tree: VarTree<[u8; 8]> = VarTree::open(dir.path(), cfg).expect("open B");
3631
3632 assert_eq!(tree.len(), 4, "all 4 entries must survive recovery");
3633
3634 for i in 0u64..4 {
3635 let key = i.to_be_bytes();
3636 let expected = vec![i as u8; 200];
3637 let got = tree
3638 .get(&key)
3639 .unwrap_or_else(|| panic!("key {i} not found after recovery"));
3640 assert_eq!(
3641 got.as_bytes(),
3642 expected.as_slice(),
3643 "value mismatch for key {i} after recovery"
3644 );
3645 }
3646
3647 let shard = &tree.engine.shards()[0];
3650 let max_fid = shard.file_ids().into_iter().max().expect("non-empty");
3651 assert!(
3652 max_fid > u16::MAX as u32,
3653 "max file_id should exceed u16::MAX after recovery (got {})",
3654 max_fid
3655 );
3656 }
3657 }
3658
3659 #[test]
3660 fn warmup_skips_large_values_block_cache() {
3661 let dir = tempdir().unwrap();
3662 let mut cfg = Config::test();
3663 cfg.shard_count = 1;
3664 cfg.max_file_size = 128 * 1024;
3665 cfg.write_buffer_size = 128 * 1024;
3666 cfg.compaction_threshold = 0.0;
3667 cfg.block_cache.max_size = 1 << 20; let tree = VarTree::<[u8; 8]>::open(dir.path(), cfg).unwrap();
3669
3670 let key = 5u64.to_be_bytes();
3671 let payload = build_large_payload(0x42); tree.put(&key, &payload).unwrap();
3673 tree.engine.shards()[0]
3676 .rotate_active_for_test(8)
3677 .expect("rotate");
3678
3679 let guard = tree.index.collector().enter();
3681 let node = tree.index.get(key.as_bytes(), &guard).unwrap();
3682 let disk = node.read_loc();
3683 assert!(
3684 disk.is_value_cache_routed(),
3685 "test precondition: value must be large"
3686 );
3687 drop(guard);
3688
3689 tree.warmup().unwrap();
3690
3691 let guard = tree.index.collector().enter();
3693 let node = tree.index.get(key.as_bytes(), &guard).unwrap();
3694 let disk = node.read_loc();
3695 let bkey = BlockKey {
3696 shard_id: 0,
3697 file_id: disk.file_id,
3698 block_offset: disk.offset as u64 & !4095,
3699 };
3700 assert!(
3701 tree.block_cache.get(&bkey).is_none(),
3702 "warmup must not cache a large value's first block"
3703 );
3704 }
3705
3706 #[test]
3711 fn var_get_many_lossy_aligned() {
3712 let dir = tempdir().unwrap();
3713 let cfg = Config::test();
3714 let tree = VarTree::<[u8; 8]>::open(dir.path(), cfg).unwrap();
3715 tree.put(&1u64.to_be_bytes(), b"alpha").unwrap();
3716 tree.put(&2u64.to_be_bytes(), b"beta").unwrap();
3717 let keys: Vec<[u8; 8]> = vec![2u64.to_be_bytes(), 9u64.to_be_bytes(), 1u64.to_be_bytes()];
3718 let got = tree.get_many(&keys);
3719 assert_eq!(got.len(), 3);
3720 assert_eq!(got[0].as_deref(), Some(&b"beta"[..]));
3721 assert_eq!(got[1], None);
3722 assert_eq!(got[2].as_deref(), Some(&b"alpha"[..]));
3723 }
3724
3725 #[test]
3726 fn var_get_many_large_value_cache_path() {
3727 let dir = tempdir().unwrap();
3728 let mut cfg = Config::test();
3729 cfg.value_cache.max_size = 4 << 20;
3732 let tree = VarTree::<[u8; 8]>::open(dir.path(), cfg).unwrap();
3733 let big = vec![7u8; 16 * 1024]; tree.put(&1u64.to_be_bytes(), &big).unwrap();
3735 let got1 = tree.get_many(&[1u64.to_be_bytes()]);
3738 assert_eq!(got1[0].as_deref(), Some(&big[..]));
3739 let got2 = tree.get_many(&[1u64.to_be_bytes()]);
3740 assert_eq!(got2[0].as_deref(), Some(&big[..]));
3741 }
3742
3743 #[test]
3744 fn var_update_many_set_keep_delete() {
3745 use crate::{Applied, BatchWrite};
3746 let dir = tempdir().unwrap();
3747 let cfg = Config::test();
3748 let tree = VarTree::<[u8; 8]>::open(dir.path(), cfg).unwrap();
3749 tree.put(&1u64.to_be_bytes(), b"old").unwrap();
3750 tree.put(&2u64.to_be_bytes(), b"gone").unwrap();
3751
3752 let items: Vec<([u8; 8], &[u8])> = vec![
3753 (1u64.to_be_bytes(), b"new"),
3754 (2u64.to_be_bytes(), b""), (3u64.to_be_bytes(), b"ins"), (4u64.to_be_bytes(), b""), ];
3758 let out = tree
3759 .update_many(items, |k, _cur, p| {
3760 let k2 = 2u64.to_be_bytes();
3761 let k4 = 4u64.to_be_bytes();
3762 if *k == k2 || *k == k4 {
3763 BatchWrite::Delete
3764 } else {
3765 BatchWrite::Set(ByteView::new(p))
3766 }
3767 })
3768 .unwrap();
3769 assert!(matches!(out[0].1, Applied::Written { .. }));
3770 assert!(matches!(out[1].1, Applied::Deleted(_)));
3771 assert!(matches!(out[2].1, Applied::Written { old: None, .. }));
3772 assert_eq!(out[3].1, Applied::Kept);
3773
3774 assert_eq!(tree.get(&1u64.to_be_bytes()).as_deref(), Some(&b"new"[..]));
3775 assert!(tree.get(&2u64.to_be_bytes()).is_none());
3776 assert_eq!(tree.get(&3u64.to_be_bytes()).as_deref(), Some(&b"ins"[..]));
3777 }
3778
3779 #[test]
3780 fn var_tree_shard_update_and_fetch_update() {
3781 let dir = tempdir().unwrap();
3782 let tree = open_test_tree_hooked::<true, true>(dir.path(), CountingHook::default());
3783 let key = 1u64.to_be_bytes();
3784 tree.put(&key, b"a").unwrap();
3785
3786 let (upd, fetched, missing) = tree
3787 .atomic(&key, |shard| {
3788 let upd = shard.update(&key, |old| {
3789 let mut v = old.to_vec();
3790 v.push(b'b');
3791 ByteView::from(v.as_slice())
3792 })?;
3793 let fetched = shard.fetch_update(&key, |old| {
3794 let mut v = old.to_vec();
3795 v.push(b'c');
3796 ByteView::from(v.as_slice())
3797 })?;
3798 let missing = shard.update(&2u64.to_be_bytes(), |old| ByteView::from(old))?;
3799 Ok((upd, fetched, missing))
3800 })
3801 .expect("atomic");
3802
3803 assert_eq!(upd.as_deref(), Some(&b"ab"[..])); assert_eq!(fetched.as_deref(), Some(&b"ab"[..])); assert!(missing.is_none());
3806 assert_eq!(tree.get(&key).as_deref(), Some(&b"abc"[..]));
3807 }
3808
3809 #[test]
3810 fn var_tree_atomic_update_fires_hook() {
3811 let dir = tempdir().unwrap();
3812 let tree = open_test_tree_hooked::<true, true>(dir.path(), CountingHook::default());
3813 let key = 1u64.to_be_bytes();
3814
3815 tree.atomic(&key, |shard| {
3816 shard.put(&key, b"init")?; shard.update(&key, |old| {
3818 let mut v = old.to_vec();
3819 v.push(b'!');
3820 ByteView::from(v.as_slice())
3821 })?; Ok(())
3823 })
3824 .expect("atomic");
3825
3826 assert_eq!(
3828 tree.hook.writes.load(AtomicOrdering::Relaxed),
3829 2,
3830 "hook must fire once for put and once for update"
3831 );
3832 assert_eq!(
3834 tree.hook.writes_with_old.load(AtomicOrdering::Relaxed),
3835 1,
3836 "only the update event carries an old value"
3837 );
3838 assert_eq!(tree.get(&key).as_deref(), Some(&b"init!"[..]));
3840 }
3841
3842 #[test]
3849 fn straddle_read_2block_direct_io_returns_full_value() {
3850 let dir = tempdir().unwrap();
3851 let mut cfg = Config::test();
3852 cfg.shard_count = 1;
3853 cfg.max_file_size = 1 << 20;
3854 cfg.write_buffer_size = 8192;
3855 cfg.direct_io = true;
3856 cfg.compaction_threshold = 0.0;
3857 let tree = VarTree::<[u8; 8]>::open(dir.path(), cfg).expect("open");
3858
3859 let v1: Vec<u8> = (0..3000).map(|i| ((i % 251) + 1) as u8).collect();
3862 let v2: Vec<u8> = (0..3000).map(|i| ((i % 241) + 2) as u8).collect();
3863 let k1 = 1u64.to_be_bytes();
3864 let k2 = 2u64.to_be_bytes();
3865 tree.put(&k1, &v1).expect("put v1");
3866 tree.put(&k2, &v2).expect("put v2");
3867
3868 tree.flush_buffers().expect("flush");
3871
3872 assert_eq!(
3873 tree.get(&k2).as_deref(),
3874 Some(&v2[..]),
3875 "straddling value must read back byte-for-byte (tail was zeroed before the fix)"
3876 );
3877 assert_eq!(tree.get(&k1).as_deref(), Some(&v1[..]));
3879 }
3880
3881 #[test]
3887 fn straddle_read_large_value_direct_io_try_get_ok() {
3888 let dir = tempdir().unwrap();
3889 let mut cfg = Config::test();
3890 cfg.shard_count = 1;
3891 cfg.max_file_size = 1 << 20;
3892 cfg.write_buffer_size = 16 * 1024;
3893 cfg.direct_io = true;
3894 cfg.compaction_threshold = 0.0;
3895 let tree = VarTree::<[u8; 8]>::open(dir.path(), cfg).expect("open");
3896
3897 let v1: Vec<u8> = (0..3000).map(|i| ((i % 251) + 1) as u8).collect();
3900 let v2: Vec<u8> = (0..9000).map(|i| ((i % 239) + 3) as u8).collect();
3901 let k1 = 1u64.to_be_bytes();
3902 let k2 = 2u64.to_be_bytes();
3903 tree.put(&k1, &v1).expect("put v1");
3904 tree.put(&k2, &v2).expect("put v2");
3905
3906 tree.flush_buffers().expect("flush");
3909
3910 let got = tree
3911 .try_get(&k2)
3912 .expect("try_get must not error for a live key");
3913 assert_eq!(
3914 got.as_deref(),
3915 Some(&v2[..]),
3916 "straddling large value must read back byte-for-byte"
3917 );
3918 }
3919
3920 #[cfg(feature = "encryption")]
3925 #[test]
3926 fn straddle_read_2block_encrypted_returns_full_value() {
3927 let dir = tempdir().unwrap();
3928 let mut cfg = Config::test();
3929 cfg.shard_count = 1;
3930 cfg.max_file_size = 1 << 20;
3931 cfg.write_buffer_size = 8192;
3932 cfg.encryption_key = Some([7u8; 32]);
3933 cfg.compaction_threshold = 0.0;
3934 let tree = VarTree::<[u8; 8]>::open(dir.path(), cfg).expect("open enc");
3935
3936 let v1: Vec<u8> = (0..3000).map(|i| ((i % 251) + 1) as u8).collect();
3937 let v2: Vec<u8> = (0..3000).map(|i| ((i % 241) + 2) as u8).collect();
3938 let k1 = 1u64.to_be_bytes();
3939 let k2 = 2u64.to_be_bytes();
3940 tree.put(&k1, &v1).expect("put v1");
3941 tree.put(&k2, &v2).expect("put v2");
3942 tree.flush_buffers().expect("flush");
3943
3944 let got = tree
3945 .try_get(&k2)
3946 .expect("encrypted straddle read must not error");
3947 assert_eq!(got.as_deref(), Some(&v2[..]));
3948 assert_eq!(tree.get(&k1).as_deref(), Some(&v1[..]));
3949 }
3950
3951 #[test]
3961 fn straddle_read_2block_cache_hit_returns_full_value() {
3962 let dir = tempdir().unwrap();
3963 let mut cfg = Config::test();
3964 cfg.shard_count = 1;
3965 cfg.max_file_size = 1 << 20;
3966 cfg.write_buffer_size = 8192;
3967 cfg.direct_io = true;
3968 cfg.block_cache.max_size = 1 << 20;
3969 cfg.compaction_threshold = 0.0;
3970 let tree = VarTree::<[u8; 8]>::open(dir.path(), cfg).expect("open");
3971
3972 let v1: Vec<u8> = (0..3000).map(|i| ((i % 251) + 1) as u8).collect();
3973 let v2: Vec<u8> = (0..3000).map(|i| ((i % 241) + 2) as u8).collect();
3974 let k1 = 1u64.to_be_bytes();
3975 let k2 = 2u64.to_be_bytes();
3976 tree.put(&k1, &v1).expect("put v1");
3977 tree.put(&k2, &v2).expect("put v2");
3978 tree.flush_buffers().expect("flush");
3980
3981 let shard_id = tree.shard_for(&k2) as u8;
3984 let key0 = BlockKey {
3985 shard_id,
3986 file_id: 1,
3987 block_offset: 0,
3988 };
3989 let (blk0, _is_full) = tree.engine.shards()[shard_id as usize]
3990 .read_block(1, 0)
3991 .expect("read block 0");
3992 tree.block_cache.insert(key0, Arc::new(blk0));
3993
3994 assert_eq!(
3995 tree.get(&k2).as_deref(),
3996 Some(&v2[..]),
3997 "straddling value via cache-hit path must read back byte-for-byte (tail was zeroed before the fix)"
3998 );
3999 }
4000}