1use std::{
5 collections::{BTreeMap, HashMap, HashSet},
6 ops::{Bound, RangeBounds},
7};
8
9use reifydb_core::{
10 actors::drop::{DropMessage, DropRequest},
11 common::CommitVersion,
12 delta::Delta,
13 encoded::{
14 key::{EncodedKey, EncodedKeyRange},
15 row::EncodedRow,
16 },
17 event::metric::{MultiCommittedEvent, MultiDelete, MultiWrite},
18 interface::store::{
19 EntryKind, MultiVersionBatch, MultiVersionCommit, MultiVersionContains, MultiVersionGet,
20 MultiVersionGetPrevious, MultiVersionRow, MultiVersionStore,
21 },
22};
23use reifydb_value::util::{cowvec::CowVec, hex};
24use tracing::{instrument, warn};
25
26use super::{
27 StandardMultiStore,
28 router::{classify_key, classify_range, is_single_version_semantics_key},
29};
30use crate::{
31 Result,
32 tier::{
33 RangeBatch, RangeCursor, TierBatch, TierStorage, VersionedGetResult,
34 commit::buffer::MultiCommitBufferTier, persistent::MultiPersistentTier,
35 },
36};
37
38const TIER_SCAN_CHUNK_SIZE: usize = 32;
39
40impl MultiVersionGet for StandardMultiStore {
41 #[instrument(name = "store::multi::get", level = "trace", skip(self), fields(key_hex = %hex::display(key.as_ref()), version = version.0))]
42 fn get(&self, key: &EncodedKey, version: CommitVersion) -> Result<Option<MultiVersionRow>> {
43 let table = classify_key(key);
44
45 if let Some(commit) = &self.commit {
46 match commit.get(table, key.as_ref(), version)? {
47 VersionedGetResult::Value {
48 value,
49 version: v,
50 } => {
51 return Ok(Some(MultiVersionRow {
52 key: key.clone(),
53 row: EncodedRow(value),
54 version: v,
55 }));
56 }
57 VersionedGetResult::Tombstone => return Ok(None),
58 VersionedGetResult::NotFound => {}
59 }
60 }
61
62 if let Some(read) = &self.read {
63 match read.get(key, version) {
64 VersionedGetResult::Value {
65 value,
66 version: v,
67 } => {
68 return Ok(Some(MultiVersionRow {
69 key: key.clone(),
70 row: EncodedRow(value),
71 version: v,
72 }));
73 }
74 VersionedGetResult::Tombstone => return Ok(None),
75 VersionedGetResult::NotFound => {}
76 }
77 }
78
79 if let Some(persistent) = &self.persistent {
80 match persistent.get(table, key.as_ref(), version)? {
81 VersionedGetResult::Value {
82 value,
83 version: v,
84 } => {
85 if let Some(read) = &self.read {
86 read.insert(key.clone(), v, Some(value.clone()));
87 }
88 return Ok(Some(MultiVersionRow {
89 key: key.clone(),
90 row: EncodedRow(value),
91 version: v,
92 }));
93 }
94 VersionedGetResult::Tombstone => return Ok(None),
95 VersionedGetResult::NotFound => {}
96 }
97 }
98
99 Ok(None)
100 }
101}
102
103impl MultiVersionContains for StandardMultiStore {
104 #[instrument(name = "store::multi::contains", level = "trace", skip(self), fields(key_hex = %hex::display(key.as_ref()), version = version.0), ret)]
105 fn contains(&self, key: &EncodedKey, version: CommitVersion) -> Result<bool> {
106 Ok(MultiVersionGet::get(self, key, version)?.is_some())
107 }
108}
109
110impl MultiVersionCommit for StandardMultiStore {
111 #[instrument(name = "store::multi::commit", level = "debug", skip(self, deltas), fields(delta_count = deltas.len(), version = version.0))]
112 fn commit(&self, deltas: CowVec<Delta>, version: CommitVersion) -> Result<()> {
113 let classified = classify_deltas(&deltas);
114
115 let (operator_drops, source_drops): (Vec<_>, Vec<_>) = classified
116 .explicit_drops
117 .into_iter()
118 .partition(|(table, _)| matches!(table, EntryKind::Operator(_)));
119
120 let drop_batch = build_drop_batch(source_drops, &classified.pending_set_keys, version);
121 self.dispatch_drops(drop_batch);
122
123 if let Some(read) = &self.read {
124 for write in &classified.writes {
125 read.invalidate(&write.key);
126 }
127 for delete in &classified.deletes {
128 read.invalidate(&delete.key);
129 }
130 }
131
132 if let Some(commit) = &self.commit {
133 commit.set(version, classified.batches)?;
134 } else if let Some(persistent) = &self.persistent {
135 persistent.set(version, classified.batches)?;
136 } else {
137 return Ok(());
138 }
139
140 self.evict_operator_state(&operator_drops)?;
141
142 self.emit_commit_metrics(classified.writes, classified.deletes, version);
143
144 Ok(())
145 }
146}
147
148struct ClassifiedDeltas {
149 pending_set_keys: HashSet<EncodedKey>,
150 writes: Vec<MultiWrite>,
151 deletes: Vec<MultiDelete>,
152 batches: TierBatch,
153 explicit_drops: Vec<(EntryKind, EncodedKey)>,
154}
155
156#[inline]
157fn classify_deltas(deltas: &CowVec<Delta>) -> ClassifiedDeltas {
158 let mut pending_set_keys: HashSet<EncodedKey> = HashSet::new();
159 let mut writes: Vec<MultiWrite> = Vec::new();
160 let mut deletes: Vec<MultiDelete> = Vec::new();
161 let mut batches: TierBatch = HashMap::new();
162 let mut explicit_drops: Vec<(EntryKind, EncodedKey)> = Vec::new();
163
164 for delta in deltas.iter() {
165 let key = delta.key();
166 let table = classify_key(key);
167 let is_single_version = is_single_version_semantics_key(key);
168
169 match delta {
170 Delta::Set {
171 key,
172 row,
173 } => {
174 if is_single_version {
175 pending_set_keys.insert(key.clone());
176 }
177 writes.push(MultiWrite {
178 key: key.clone(),
179 value_bytes: row.len() as u64,
180 });
181 batches.entry(table).or_default().push((key.clone(), Some(row.0.clone())));
182 }
183 Delta::Unset {
184 key,
185 row,
186 } => {
187 deletes.push(MultiDelete {
188 key: key.clone(),
189 value_bytes: row.len() as u64,
190 });
191 batches.entry(table).or_default().push((key.clone(), None));
192 }
193 Delta::Remove {
194 key,
195 } => {
196 deletes.push(MultiDelete {
197 key: key.clone(),
198 value_bytes: 0,
199 });
200 batches.entry(table).or_default().push((key.clone(), None));
201 }
202 Delta::Drop {
203 key,
204 } => {
205 explicit_drops.push((table, key.clone()));
206 }
207 }
208 }
209
210 ClassifiedDeltas {
211 pending_set_keys,
212 writes,
213 deletes,
214 batches,
215 explicit_drops,
216 }
217}
218
219#[inline]
220fn build_drop_batch(
221 explicit_drops: Vec<(EntryKind, EncodedKey)>,
222 pending_set_keys: &HashSet<EncodedKey>,
223 version: CommitVersion,
224) -> Vec<DropRequest> {
225 let mut drop_batch = Vec::with_capacity(explicit_drops.len() + pending_set_keys.len());
226 for (table, key) in explicit_drops {
227 let pending_version = if pending_set_keys.contains(key.as_ref()) {
228 Some(version)
229 } else {
230 None
231 };
232 drop_batch.push(DropRequest {
233 table,
234 key,
235 commit_version: version,
236 pending_version,
237 });
238 }
239 for key in pending_set_keys.iter() {
240 let encoded = EncodedKey::new(key.to_vec());
241 let table = classify_key(&encoded);
242 drop_batch.push(DropRequest {
243 table,
244 key: encoded,
245 commit_version: version,
246 pending_version: Some(version),
247 });
248 }
249 drop_batch
250}
251
252impl StandardMultiStore {
253 pub fn get_many(
254 &self,
255 keys: &[EncodedKey],
256 version: CommitVersion,
257 ) -> Result<HashMap<EncodedKey, MultiVersionRow>> {
258 let mut by_table: HashMap<EntryKind, Vec<&EncodedKey>> = HashMap::new();
259 for key in keys {
260 by_table.entry(classify_key(key)).or_default().push(key);
261 }
262
263 let mut out: HashMap<EncodedKey, MultiVersionRow> = HashMap::new();
264
265 for (table, table_keys) in by_table {
266 let key_slices: Vec<&[u8]> = table_keys.iter().map(|k| k.as_ref()).collect();
267
268 let commit_results = match &self.commit {
269 Some(commit) => commit.get_many(table, &key_slices, version)?,
270 None => vec![VersionedGetResult::NotFound; key_slices.len()],
271 };
272
273 let mut read_aligned = vec![VersionedGetResult::NotFound; key_slices.len()];
274 let mut persistent_idx: Vec<usize> = Vec::new();
275 let mut persistent_slices: Vec<&[u8]> = Vec::new();
276 for (i, result) in commit_results.iter().enumerate() {
277 if !matches!(result, VersionedGetResult::NotFound) {
278 continue;
279 }
280 let read_hit = self
281 .read
282 .as_ref()
283 .map(|c| c.get(table_keys[i], version))
284 .unwrap_or(VersionedGetResult::NotFound);
285 match read_hit {
286 VersionedGetResult::Value {
287 value,
288 version: v,
289 } => {
290 read_aligned[i] = VersionedGetResult::Value {
291 value,
292 version: v,
293 };
294 }
295 VersionedGetResult::Tombstone => {
296 read_aligned[i] = VersionedGetResult::Tombstone;
297 }
298 VersionedGetResult::NotFound => {
299 persistent_idx.push(i);
300 persistent_slices.push(key_slices[i]);
301 }
302 }
303 }
304
305 let mut persistent_aligned = vec![VersionedGetResult::NotFound; key_slices.len()];
306 if !persistent_slices.is_empty()
307 && let Some(persistent) = &self.persistent
308 {
309 let persistent_results = persistent.get_many(table, &persistent_slices, version)?;
310 for (slot, result) in persistent_idx.into_iter().zip(persistent_results) {
311 if let (
312 Some(read),
313 VersionedGetResult::Value {
314 value,
315 version: v,
316 },
317 ) = (&self.read, &result)
318 {
319 read.insert(table_keys[slot].clone(), *v, Some(value.clone()));
320 }
321 persistent_aligned[slot] = result;
322 }
323 }
324
325 for (i, key) in table_keys.into_iter().enumerate() {
326 let resolved = match &commit_results[i] {
327 VersionedGetResult::Value {
328 value,
329 version: v,
330 } => Some((value.clone(), *v)),
331 VersionedGetResult::Tombstone => None,
332 VersionedGetResult::NotFound => match &read_aligned[i] {
333 VersionedGetResult::Value {
334 value,
335 version: v,
336 } => Some((value.clone(), *v)),
337 VersionedGetResult::Tombstone => None,
338 VersionedGetResult::NotFound => match &persistent_aligned[i] {
339 VersionedGetResult::Value {
340 value,
341 version: v,
342 } => Some((value.clone(), *v)),
343 _ => None,
344 },
345 },
346 };
347
348 if let Some((value, v)) = resolved {
349 out.insert(
350 key.clone(),
351 MultiVersionRow {
352 key: key.clone(),
353 row: EncodedRow(value),
354 version: v,
355 },
356 );
357 }
358 }
359 }
360
361 Ok(out)
362 }
363
364 #[inline]
365 fn dispatch_drops(&self, drop_batch: Vec<DropRequest>) {
366 if drop_batch.is_empty() {
367 return;
368 }
369 if let Some(actor) = &self.drop_actor
370 && actor.send_blocking(DropMessage::Batch(drop_batch)).is_err()
371 {
372 warn!("Failed to send drop batch");
373 }
374 }
375
376 fn evict_operator_state(&self, drops: &[(EntryKind, EncodedKey)]) -> Result<()> {
377 if drops.is_empty() {
378 return Ok(());
379 }
380
381 if let Some(read) = &self.read {
382 for (_, key) in drops {
383 read.invalidate(key);
384 }
385 }
386
387 if let Some(commit) = &self.commit {
388 let mut batches: HashMap<EntryKind, Vec<(EncodedKey, CommitVersion)>> = HashMap::new();
389 for (table, key) in drops {
390 for (entry_version, _) in commit.get_all_versions(*table, key.as_ref())? {
391 batches.entry(*table).or_default().push((key.clone(), entry_version));
392 }
393 }
394 if !batches.is_empty() {
395 commit.drop(batches)?;
396 }
397 }
398
399 if let Some(persistent) = &self.persistent {
400 let mut by_table: HashMap<EntryKind, Vec<EncodedKey>> = HashMap::new();
401 for (table, key) in drops {
402 by_table.entry(*table).or_default().push(key.clone());
403 }
404 for (table, keys) in by_table {
405 persistent.delete_keys(table, &keys)?;
406 }
407 }
408
409 Ok(())
410 }
411
412 #[inline]
413 fn emit_commit_metrics(&self, writes: Vec<MultiWrite>, deletes: Vec<MultiDelete>, version: CommitVersion) {
414 if writes.is_empty() && deletes.is_empty() {
415 return;
416 }
417 self.event_bus.emit(MultiCommittedEvent::new(writes, deletes, vec![], version));
418 }
419}
420
421#[derive(Debug, Clone, Default)]
422pub struct MultiVersionRangeCursor {
423 pub commit: RangeCursor,
424
425 pub persistent: RangeCursor,
426
427 pub exhausted: bool,
428}
429
430impl MultiVersionRangeCursor {
431 pub fn new() -> Self {
432 Self::default()
433 }
434
435 pub fn is_exhausted(&self) -> bool {
436 self.exhausted
437 }
438}
439
440pub struct TierScanQuery<'a> {
441 pub table: EntryKind,
442 pub start: &'a [u8],
443 pub end: &'a [u8],
444 pub version: CommitVersion,
445 pub range: &'a EncodedKeyRange,
446}
447
448pub fn scan_tier_chunk<S: TierStorage>(
449 storage: &S,
450 cursor: &mut RangeCursor,
451 scan: &TierScanQuery,
452 collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
453) -> Result<bool> {
454 let batch = storage.range_next(
455 scan.table,
456 cursor,
457 Bound::Included(scan.start),
458 Bound::Included(scan.end),
459 scan.version,
460 TIER_SCAN_CHUNK_SIZE,
461 )?;
462 merge_tier_batch(batch, scan.range, collected)
463}
464
465pub fn scan_tier_chunk_rev<S: TierStorage>(
466 storage: &S,
467 cursor: &mut RangeCursor,
468 scan: &TierScanQuery,
469 collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
470) -> Result<bool> {
471 let batch = storage.range_rev_next(
472 scan.table,
473 cursor,
474 Bound::Included(scan.start),
475 Bound::Included(scan.end),
476 scan.version,
477 TIER_SCAN_CHUNK_SIZE,
478 )?;
479 merge_tier_batch(batch, scan.range, collected)
480}
481
482#[inline]
483fn merge_tier_batch(
484 batch: RangeBatch,
485 range: &EncodedKeyRange,
486 collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
487) -> Result<bool> {
488 if batch.entries.is_empty() {
489 return Ok(false);
490 }
491
492 for entry in batch.entries {
493 let original_key = entry.key.as_slice().to_vec();
494 let entry_version = entry.version;
495
496 let original_key_encoded = EncodedKey::new(original_key.clone());
497 if !range.contains(&original_key_encoded) {
498 continue;
499 }
500
501 let should_update = match collected.get(&original_key) {
502 None => true,
503 Some((existing_version, _)) => entry_version > *existing_version,
504 };
505
506 if should_update {
507 collected.insert(original_key, (entry_version, entry.value));
508 }
509 }
510
511 Ok(true)
512}
513
514#[inline]
515pub fn collected_to_batch(
516 collected: BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
517 has_more: bool,
518) -> MultiVersionBatch {
519 let items: Vec<MultiVersionRow> = collected
520 .into_iter()
521 .filter_map(|(key_bytes, (v, value))| {
522 value.map(|val| MultiVersionRow {
523 key: EncodedKey::new(key_bytes),
524 row: EncodedRow(val),
525 version: v,
526 })
527 })
528 .collect();
529
530 MultiVersionBatch {
531 items,
532 has_more,
533 }
534}
535
536#[inline]
537fn step_all_tiers(
538 buffer: Option<&MultiCommitBufferTier>,
539 buffer_cursor: &mut RangeCursor,
540 persistent: Option<&MultiPersistentTier>,
541 persistent_cursor: &mut RangeCursor,
542 scan: &TierScanQuery,
543 collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
544) -> Result<bool> {
545 let mut any_progress = false;
546 if let Some(s) = buffer
547 && !buffer_cursor.exhausted
548 {
549 any_progress |= scan_tier_chunk(s, buffer_cursor, scan, collected)?;
550 }
551 if let Some(s) = persistent
552 && !persistent_cursor.exhausted
553 {
554 any_progress |= scan_tier_chunk(s, persistent_cursor, scan, collected)?;
555 }
556 Ok(any_progress)
557}
558
559pub fn scan_tiers_latest(
560 buffer: Option<&MultiCommitBufferTier>,
561 persistent: Option<&MultiPersistentTier>,
562 range: EncodedKeyRange,
563 version: CommitVersion,
564 max_keys: usize,
565) -> Result<MultiVersionBatch> {
566 let table = classify_key_range(&range);
567 let (start, end) = make_range_bounds(&range);
568 let scan = TierScanQuery {
569 table,
570 start: &start,
571 end: &end,
572 version,
573 range: &range,
574 };
575
576 let mut collected: BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)> = BTreeMap::new();
577 let mut buffer_cursor = RangeCursor::default();
578 let mut persistent_cursor = RangeCursor::default();
579 let mut exhausted = false;
580
581 while collected.len() < max_keys {
582 let progress = step_all_tiers(
583 buffer,
584 &mut buffer_cursor,
585 persistent,
586 &mut persistent_cursor,
587 &scan,
588 &mut collected,
589 )?;
590 if !progress {
591 exhausted = true;
592 break;
593 }
594 }
595
596 Ok(collected_to_batch(collected, !exhausted))
597}
598
599impl StandardMultiStore {
600 pub fn range_next(
601 &self,
602 cursor: &mut MultiVersionRangeCursor,
603 range: EncodedKeyRange,
604 version: CommitVersion,
605 batch_size: u64,
606 ) -> Result<MultiVersionBatch> {
607 if cursor.exhausted {
608 return Ok(MultiVersionBatch {
609 items: Vec::new(),
610 has_more: false,
611 });
612 }
613
614 mark_unconfigured_exhausted(self, cursor);
615
616 let table = classify_key_range(&range);
617 let (start, end) = make_range_bounds(&range);
618 let batch_size = batch_size as usize;
619 let scan = TierScanQuery {
620 table,
621 start: &start,
622 end: &end,
623 version,
624 range: &range,
625 };
626
627 let mut collected: BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)> = BTreeMap::new();
628
629 while collected.len() < batch_size {
630 let progress = step_all_tiers(
631 self.commit.as_ref(),
632 &mut cursor.commit,
633 self.persistent.as_ref(),
634 &mut cursor.persistent,
635 &scan,
636 &mut collected,
637 )?;
638 if !progress {
639 cursor.exhausted = true;
640 break;
641 }
642 }
643
644 apply_forward_horizon(cursor, &mut collected);
645
646 let items: Vec<MultiVersionRow> = collected
647 .into_iter()
648 .filter_map(|(key_bytes, (v, value))| {
649 value.map(|val| MultiVersionRow {
650 key: EncodedKey::new(key_bytes),
651 row: EncodedRow(val),
652 version: v,
653 })
654 })
655 .collect();
656
657 let has_more = !cursor.exhausted;
658
659 Ok(MultiVersionBatch {
660 items,
661 has_more,
662 })
663 }
664
665 pub fn range(
666 &self,
667 range: EncodedKeyRange,
668 version: CommitVersion,
669 batch_size: usize,
670 ) -> MultiVersionRangeIter {
671 MultiVersionRangeIter {
672 store: self.clone(),
673 cursor: MultiVersionRangeCursor::new(),
674 range,
675 version,
676 batch_size,
677 current_batch: Vec::new(),
678 current_index: 0,
679 }
680 }
681
682 pub fn range_rev(
683 &self,
684 range: EncodedKeyRange,
685 version: CommitVersion,
686 batch_size: usize,
687 ) -> MultiVersionRangeRevIter {
688 MultiVersionRangeRevIter {
689 store: self.clone(),
690 cursor: MultiVersionRangeCursor::new(),
691 range,
692 version,
693 batch_size,
694 current_batch: Vec::new(),
695 current_index: 0,
696 }
697 }
698
699 fn range_rev_next(
700 &self,
701 cursor: &mut MultiVersionRangeCursor,
702 range: EncodedKeyRange,
703 version: CommitVersion,
704 batch_size: u64,
705 ) -> Result<MultiVersionBatch> {
706 if cursor.exhausted {
707 return Ok(MultiVersionBatch {
708 items: Vec::new(),
709 has_more: false,
710 });
711 }
712
713 mark_unconfigured_exhausted(self, cursor);
714
715 let table = classify_key_range(&range);
716 let (start, end) = make_range_bounds(&range);
717 let batch_size = batch_size as usize;
718 let scan = TierScanQuery {
719 table,
720 start: &start,
721 end: &end,
722 version,
723 range: &range,
724 };
725
726 let mut collected: BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)> = BTreeMap::new();
727
728 while collected.len() < batch_size {
729 let mut any_progress = false;
730
731 if let Some(commit) = &self.commit
732 && !cursor.commit.exhausted
733 {
734 any_progress |= scan_tier_chunk_rev(commit, &mut cursor.commit, &scan, &mut collected)?;
735 }
736
737 if let Some(persistent) = &self.persistent
738 && !cursor.persistent.exhausted
739 {
740 any_progress |=
741 scan_tier_chunk_rev(persistent, &mut cursor.persistent, &scan, &mut collected)?;
742 }
743
744 if !any_progress {
745 cursor.exhausted = true;
746 break;
747 }
748 }
749
750 apply_reverse_horizon(cursor, &mut collected);
751
752 let items: Vec<MultiVersionRow> = collected
753 .into_iter()
754 .rev()
755 .filter_map(|(key_bytes, (v, value))| {
756 value.map(|val| MultiVersionRow {
757 key: EncodedKey::new(key_bytes),
758 row: EncodedRow(val),
759 version: v,
760 })
761 })
762 .collect();
763
764 let has_more = !cursor.exhausted;
765
766 Ok(MultiVersionBatch {
767 items,
768 has_more,
769 })
770 }
771}
772
773fn mark_unconfigured_exhausted(store: &StandardMultiStore, cursor: &mut MultiVersionRangeCursor) {
774 if store.commit.is_none() {
775 cursor.commit.exhausted = true;
776 }
777 if store.persistent.is_none() {
778 cursor.persistent.exhausted = true;
779 }
780}
781
782fn apply_forward_horizon(
783 cursor: &mut MultiVersionRangeCursor,
784 collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
785) {
786 let horizon = forward_horizon(cursor);
787 if let Some(h) = horizon {
788 collected.retain(|k, _| k.as_slice() <= h.as_slice());
789 rewind_over_advanced_forward(cursor, &h);
790 }
791}
792
793fn apply_reverse_horizon(
794 cursor: &mut MultiVersionRangeCursor,
795 collected: &mut BTreeMap<Vec<u8>, (CommitVersion, Option<CowVec<u8>>)>,
796) {
797 let horizon = reverse_horizon(cursor);
798 if let Some(h) = horizon {
799 collected.retain(|k, _| k.as_slice() >= h.as_slice());
800 rewind_over_advanced_reverse(cursor, &h);
801 }
802}
803
804fn forward_horizon(cursor: &MultiVersionRangeCursor) -> Option<EncodedKey> {
805 let mut horizon: Option<EncodedKey> = None;
806 for tier in [&cursor.commit, &cursor.persistent] {
807 if tier.exhausted {
808 continue;
809 }
810 let last = match &tier.last_key {
811 Some(k) => k.clone(),
812
813 None => return None,
814 };
815 horizon = Some(match horizon {
816 None => last,
817 Some(prev) => {
818 if last.as_slice() < prev.as_slice() {
819 last
820 } else {
821 prev
822 }
823 }
824 });
825 }
826 horizon
827}
828
829fn reverse_horizon(cursor: &MultiVersionRangeCursor) -> Option<EncodedKey> {
830 let mut horizon: Option<EncodedKey> = None;
831 for tier in [&cursor.commit, &cursor.persistent] {
832 if tier.exhausted {
833 continue;
834 }
835 let last = match &tier.last_key {
836 Some(k) => k.clone(),
837 None => return None,
838 };
839 horizon = Some(match horizon {
840 None => last,
841 Some(prev) => {
842 if last.as_slice() > prev.as_slice() {
843 last
844 } else {
845 prev
846 }
847 }
848 });
849 }
850 horizon
851}
852
853fn rewind_over_advanced_forward(cursor: &mut MultiVersionRangeCursor, horizon: &EncodedKey) {
854 for tier in [&mut cursor.commit, &mut cursor.persistent] {
855 if tier.exhausted {
856 continue;
857 }
858 if let Some(last) = &tier.last_key
859 && last.as_slice() > horizon.as_slice()
860 {
861 tier.last_key = Some(horizon.clone());
862 }
863 }
864}
865
866fn rewind_over_advanced_reverse(cursor: &mut MultiVersionRangeCursor, horizon: &EncodedKey) {
867 for tier in [&mut cursor.commit, &mut cursor.persistent] {
868 if tier.exhausted {
869 continue;
870 }
871 if let Some(last) = &tier.last_key
872 && last.as_slice() < horizon.as_slice()
873 {
874 tier.last_key = Some(horizon.clone());
875 }
876 }
877}
878
879impl MultiVersionGetPrevious for StandardMultiStore {
880 fn get_previous_version(
881 &self,
882 key: &EncodedKey,
883 before_version: CommitVersion,
884 ) -> Result<Option<MultiVersionRow>> {
885 if before_version.0 == 0 {
886 return Ok(None);
887 }
888
889 let table = classify_key(key);
890 let prev_version = CommitVersion(before_version.0 - 1);
891
892 if let Some(commit) = &self.commit {
893 match commit.get(table, key.as_ref(), prev_version)? {
894 VersionedGetResult::Value {
895 value,
896 version,
897 } => {
898 return Ok(Some(MultiVersionRow {
899 key: key.clone(),
900 row: EncodedRow(CowVec::new(value.to_vec())),
901 version,
902 }));
903 }
904 VersionedGetResult::Tombstone => return Ok(None),
905 VersionedGetResult::NotFound => {}
906 }
907 }
908
909 if let Some(read) = &self.read {
910 match read.get(key, prev_version) {
911 VersionedGetResult::Value {
912 value,
913 version,
914 } => {
915 return Ok(Some(MultiVersionRow {
916 key: key.clone(),
917 row: EncodedRow(CowVec::new(value.to_vec())),
918 version,
919 }));
920 }
921 VersionedGetResult::Tombstone => return Ok(None),
922 VersionedGetResult::NotFound => {}
923 }
924 }
925
926 if let Some(persistent) = &self.persistent {
927 match persistent.get(table, key.as_ref(), prev_version)? {
928 VersionedGetResult::Value {
929 value,
930 version,
931 } => {
932 if let Some(read) = &self.read {
933 read.insert(key.clone(), version, Some(value.clone()));
934 }
935 return Ok(Some(MultiVersionRow {
936 key: key.clone(),
937 row: EncodedRow(CowVec::new(value.to_vec())),
938 version,
939 }));
940 }
941 VersionedGetResult::Tombstone => return Ok(None),
942 VersionedGetResult::NotFound => {}
943 }
944 }
945
946 Ok(None)
947 }
948}
949
950impl MultiVersionStore for StandardMultiStore {}
951
952pub struct MultiVersionRangeIter {
953 store: StandardMultiStore,
954 cursor: MultiVersionRangeCursor,
955 range: EncodedKeyRange,
956 version: CommitVersion,
957 batch_size: usize,
958 current_batch: Vec<MultiVersionRow>,
959 current_index: usize,
960}
961
962impl Iterator for MultiVersionRangeIter {
963 type Item = Result<MultiVersionRow>;
964
965 fn next(&mut self) -> Option<Self::Item> {
966 if self.current_index < self.current_batch.len() {
967 let item = self.current_batch[self.current_index].clone();
968 self.current_index += 1;
969 return Some(Ok(item));
970 }
971
972 if self.cursor.exhausted {
973 return None;
974 }
975
976 match self.store.range_next(&mut self.cursor, self.range.clone(), self.version, self.batch_size as u64)
977 {
978 Ok(batch) => {
979 if batch.items.is_empty() {
980 if self.cursor.exhausted {
981 return None;
982 }
983 return self.next();
984 }
985 self.current_batch = batch.items;
986 self.current_index = 0;
987 self.next()
988 }
989 Err(e) => Some(Err(e)),
990 }
991 }
992}
993
994pub struct MultiVersionRangeRevIter {
995 store: StandardMultiStore,
996 cursor: MultiVersionRangeCursor,
997 range: EncodedKeyRange,
998 version: CommitVersion,
999 batch_size: usize,
1000 current_batch: Vec<MultiVersionRow>,
1001 current_index: usize,
1002}
1003
1004impl Iterator for MultiVersionRangeRevIter {
1005 type Item = Result<MultiVersionRow>;
1006
1007 fn next(&mut self) -> Option<Self::Item> {
1008 if self.current_index < self.current_batch.len() {
1009 let item = self.current_batch[self.current_index].clone();
1010 self.current_index += 1;
1011 return Some(Ok(item));
1012 }
1013
1014 if self.cursor.exhausted {
1015 return None;
1016 }
1017
1018 match self.store.range_rev_next(
1019 &mut self.cursor,
1020 self.range.clone(),
1021 self.version,
1022 self.batch_size as u64,
1023 ) {
1024 Ok(batch) => {
1025 if batch.items.is_empty() {
1026 if self.cursor.exhausted {
1027 return None;
1028 }
1029 return self.next();
1030 }
1031 self.current_batch = batch.items;
1032 self.current_index = 0;
1033 self.next()
1034 }
1035 Err(e) => Some(Err(e)),
1036 }
1037 }
1038}
1039
1040fn classify_key_range(range: &EncodedKeyRange) -> EntryKind {
1041 classify_range(range).unwrap_or(EntryKind::Multi)
1042}
1043
1044fn make_range_bounds(range: &EncodedKeyRange) -> (Vec<u8>, Vec<u8>) {
1045 let start = match &range.start {
1046 Bound::Included(key) => key.as_ref().to_vec(),
1047 Bound::Excluded(key) => key.as_ref().to_vec(),
1048 Bound::Unbounded => vec![],
1049 };
1050
1051 let end = match &range.end {
1052 Bound::Included(key) => key.as_ref().to_vec(),
1053 Bound::Excluded(key) => key.as_ref().to_vec(),
1054 Bound::Unbounded => vec![0xFFu8; 256],
1055 };
1056
1057 (start, end)
1058}