1use std::{
5 cmp::Ordering,
6 collections, iter,
7 ops::{
8 Bound::{Excluded, Included, Unbounded},
9 RangeBounds,
10 },
11 vec,
12};
13
14use collections::BTreeMap;
15use iter::Peekable;
16use reifydb_codec::{
17 encoded::row::EncodedRow,
18 key::encoded::{EncodedKey, EncodedKeyRange},
19};
20use reifydb_core::{
21 actors::pending::PendingWrite,
22 common::CommitVersion,
23 interface::store::{MultiVersionBatch, MultiVersionRow},
24 key::{Key, kind::KeyKind},
25};
26use reifydb_transaction::multi::RangeScope;
27use reifydb_value::Result;
28use vec::IntoIter;
29
30use super::FlowTransaction;
31
32pub(crate) enum ReadFrom {
33 StateQuery,
34
35 Query,
36
37 DictionaryQuery,
38}
39
40impl FlowTransaction {
41 pub fn get(&mut self, key: &EncodedKey) -> Result<Option<EncodedRow>> {
42 let inner = self.inner();
43 if inner.pending.is_removed(key) {
44 return Ok(None);
45 }
46 if let Some(value) = inner.pending.get(key) {
47 return Ok(Some(value.clone()));
48 }
49
50 if let Self::Transactional {
51 base_pending,
52 ..
53 } = self
54 {
55 if base_pending.is_removed(key) {
56 return Ok(None);
57 }
58 if let Some(value) = base_pending.get(key) {
59 return Ok(Some(value.clone()));
60 }
61 }
62
63 if let Self::Ephemeral {
64 inner,
65 state,
66 } = self
67 {
68 return match Self::read_from(key) {
69 ReadFrom::StateQuery => Ok(state.get(key).cloned()),
70 ReadFrom::Query | ReadFrom::DictionaryQuery => {
71 match inner.dictionary_query.as_ref().unwrap_or(&inner.query).get(key)? {
72 Some(multi) => Ok(Some(multi.row().clone())),
73 None => Ok(None),
74 }
75 }
76 };
77 }
78
79 if let Some(cached) = self.inner().prefetch.get(key) {
80 return Ok(cached.clone());
81 }
82
83 let inner = self.inner_mut();
84 let query = match Self::read_from(key) {
85 ReadFrom::StateQuery => inner.state_query.as_ref().unwrap(),
86 ReadFrom::Query => &inner.query,
87 ReadFrom::DictionaryQuery => inner.dictionary_query.as_ref().unwrap_or(&inner.query),
88 };
89 match query.get(key)? {
90 Some(multi) => Ok(Some(multi.row().clone())),
91 None => Ok(None),
92 }
93 }
94
95 pub fn contains_key(&mut self, key: &EncodedKey) -> Result<bool> {
96 let inner = self.inner();
97 if inner.pending.is_removed(key) {
98 return Ok(false);
99 }
100 if inner.pending.get(key).is_some() {
101 return Ok(true);
102 }
103
104 if let Self::Transactional {
105 base_pending,
106 ..
107 } = self
108 {
109 if base_pending.is_removed(key) {
110 return Ok(false);
111 }
112 if base_pending.get(key).is_some() {
113 return Ok(true);
114 }
115 }
116
117 if let Self::Ephemeral {
118 inner,
119 state,
120 } = self
121 {
122 return match Self::read_from(key) {
123 ReadFrom::StateQuery => Ok(state.contains_key(key)),
124 ReadFrom::Query | ReadFrom::DictionaryQuery => {
125 inner.dictionary_query.as_ref().unwrap_or(&inner.query).contains_key(key)
126 }
127 };
128 }
129
130 let inner = self.inner_mut();
131 let query = match Self::read_from(key) {
132 ReadFrom::StateQuery => inner.state_query.as_ref().unwrap(),
133 ReadFrom::Query => &inner.query,
134 ReadFrom::DictionaryQuery => inner.dictionary_query.as_ref().unwrap_or(&inner.query),
135 };
136 query.contains_key(key)
137 }
138
139 pub fn prefix(&mut self, prefix: &EncodedKey) -> Result<MultiVersionBatch> {
140 let range = EncodedKeyRange::prefix(prefix);
141 let items = self.range(range, RangeScope::All, 1024).collect::<Result<Vec<_>>>()?;
142 Ok(MultiVersionBatch {
143 items,
144 has_more: false,
145 })
146 }
147
148 pub(crate) fn read_from(key: &EncodedKey) -> ReadFrom {
149 match Key::kind(key) {
150 None => ReadFrom::Query,
151 Some(kind) => match kind {
152 KeyKind::FlowNodeState => ReadFrom::StateQuery,
153 KeyKind::FlowNodeInternalState => ReadFrom::StateQuery,
154 KeyKind::RingBufferMetadata => ReadFrom::StateQuery,
155 KeyKind::SeriesMetadata => ReadFrom::StateQuery,
156
157 KeyKind::Row => ReadFrom::Query,
158
159 KeyKind::Namespace => ReadFrom::Query,
160 KeyKind::Table => ReadFrom::Query,
161 KeyKind::NamespaceTable => ReadFrom::Query,
162 KeyKind::SystemSequence => ReadFrom::Query,
163 KeyKind::Columns => ReadFrom::Query,
164 KeyKind::Column => ReadFrom::Query,
165 KeyKind::RowSequence => ReadFrom::Query,
166 KeyKind::ColumnProperty => ReadFrom::Query,
167 KeyKind::SystemVersion => ReadFrom::Query,
168 KeyKind::TransactionVersion => ReadFrom::Query,
169 KeyKind::Index => ReadFrom::Query,
170 KeyKind::IndexEntry => ReadFrom::Query,
171 KeyKind::ColumnSequence => ReadFrom::Query,
172 KeyKind::CdcConsumer => ReadFrom::Query,
173 KeyKind::View => ReadFrom::Query,
174 KeyKind::NamespaceView => ReadFrom::Query,
175 KeyKind::PrimaryKey => ReadFrom::Query,
176 KeyKind::RingBuffer => ReadFrom::Query,
177 KeyKind::NamespaceRingBuffer => ReadFrom::Query,
178 KeyKind::ShapeRetentionStrategy => ReadFrom::Query,
179 KeyKind::OperatorRetentionStrategy => ReadFrom::Query,
180 KeyKind::Flow => ReadFrom::Query,
181 KeyKind::NamespaceFlow => ReadFrom::Query,
182 KeyKind::FlowNode => ReadFrom::Query,
183 KeyKind::FlowNodeByFlow => ReadFrom::Query,
184 KeyKind::FlowEdge => ReadFrom::Query,
185 KeyKind::FlowEdgeByFlow => ReadFrom::Query,
186 KeyKind::Dictionary => ReadFrom::Query,
187 KeyKind::DictionaryEntry => ReadFrom::DictionaryQuery,
188 KeyKind::DictionaryEntryIndex => ReadFrom::DictionaryQuery,
189 KeyKind::NamespaceDictionary => ReadFrom::Query,
190 KeyKind::Metric => ReadFrom::Query,
191 KeyKind::FlowVersion => ReadFrom::Query,
192 KeyKind::Subscription => ReadFrom::Query,
193 KeyKind::SubscriptionRow => ReadFrom::Query,
194 KeyKind::SubscriptionColumn => ReadFrom::Query,
195 KeyKind::Shape => ReadFrom::Query,
196 KeyKind::RowShapeField => ReadFrom::Query,
197 KeyKind::SumType => ReadFrom::Query,
198 KeyKind::NamespaceSumType => ReadFrom::Query,
199 KeyKind::Handler => ReadFrom::Query,
200 KeyKind::NamespaceHandler => ReadFrom::Query,
201 KeyKind::VariantHandler => ReadFrom::Query,
202 KeyKind::Series => ReadFrom::Query,
203 KeyKind::NamespaceSeries => ReadFrom::Query,
204 KeyKind::Identity => ReadFrom::Query,
205 KeyKind::Role => ReadFrom::Query,
206 KeyKind::GrantedRole => ReadFrom::Query,
207 KeyKind::Policy => ReadFrom::Query,
208 KeyKind::PolicyOp => ReadFrom::Query,
209 KeyKind::Migration => ReadFrom::Query,
210 KeyKind::MigrationEvent => ReadFrom::Query,
211 KeyKind::Authentication => ReadFrom::Query,
212 KeyKind::ConfigStorage => ReadFrom::Query,
213 KeyKind::Token => ReadFrom::Query,
214 KeyKind::Source => ReadFrom::Query,
215 KeyKind::NamespaceSource => ReadFrom::Query,
216 KeyKind::Sink => ReadFrom::Query,
217 KeyKind::NamespaceSink => ReadFrom::Query,
218 KeyKind::SourceCheckpoint => ReadFrom::Query,
219 KeyKind::RowSettings => ReadFrom::Query,
220 KeyKind::OperatorSettings => ReadFrom::Query,
221 KeyKind::Procedure => ReadFrom::Query,
222 KeyKind::NamespaceProcedure => ReadFrom::Query,
223 KeyKind::ProcedureParam => ReadFrom::Query,
224 KeyKind::Binding => ReadFrom::Query,
225 KeyKind::NamespaceBinding => ReadFrom::Query,
226 KeyKind::ColumnSnapshot => ReadFrom::Query,
227 KeyKind::SeriesColumnSnapshot => ReadFrom::Query,
228 KeyKind::TableColumnSnapshot => ReadFrom::Query,
229 KeyKind::VersionEpoch => ReadFrom::Query,
230 },
231 }
232 }
233
234 pub fn range(
235 &mut self,
236 range: EncodedKeyRange,
237 scope: RangeScope,
238 batch_size: usize,
239 ) -> Box<dyn Iterator<Item = Result<MultiVersionRow>> + Send + '_> {
240 match self {
241 Self::Deferred {
242 inner,
243 ..
244 }
245 | Self::Committing {
246 inner,
247 ..
248 } => {
249 let merged: BTreeMap<EncodedKey, PendingWrite> = inner
250 .pending
251 .range((range.start.as_ref(), range.end.as_ref()))
252 .map(|(k, v)| (k.clone(), v.clone()))
253 .collect();
254 let pending_vec: Vec<(EncodedKey, PendingWrite)> = merged.into_iter().collect();
255
256 let query = match range.start.as_ref() {
257 Included(start) | Excluded(start) => match Self::read_from(start) {
258 ReadFrom::StateQuery => inner.state_query.as_ref().unwrap(),
259 ReadFrom::Query => &inner.query,
260 ReadFrom::DictionaryQuery => {
261 inner.dictionary_query.as_ref().unwrap_or(&inner.query)
262 }
263 },
264 Unbounded => &inner.query,
265 };
266
267 let storage_iter = query.range(range, scope, batch_size);
268 let v = inner.version;
269 Box::new(flow_merge_pending_iterator(pending_vec, storage_iter, v))
270 }
271 Self::Transactional {
272 inner,
273 base_pending,
274 ..
275 } => {
276 let mut merged: BTreeMap<EncodedKey, PendingWrite> = base_pending
277 .range((range.start.as_ref(), range.end.as_ref()))
278 .map(|(k, v)| (k.clone(), v.clone()))
279 .collect();
280 for (k, v) in inner.pending.range((range.start.as_ref(), range.end.as_ref())) {
281 merged.insert(k.clone(), v.clone());
282 }
283 let pending_vec: Vec<(EncodedKey, PendingWrite)> = merged.into_iter().collect();
284
285 let query = match range.start.as_ref() {
286 Included(start) | Excluded(start) => match Self::read_from(start) {
287 ReadFrom::StateQuery => inner.state_query.as_ref().unwrap(),
288 ReadFrom::Query => &inner.query,
289 ReadFrom::DictionaryQuery => {
290 inner.dictionary_query.as_ref().unwrap_or(&inner.query)
291 }
292 },
293 Unbounded => &inner.query,
294 };
295
296 let storage_iter = query.range(range, scope, batch_size);
297 let v = inner.version;
298 Box::new(flow_merge_pending_iterator(pending_vec, storage_iter, v))
299 }
300 Self::Ephemeral {
301 inner,
302 state,
303 } => {
304 let is_state_range = match range.start.as_ref() {
305 Included(start) | Excluded(start) => {
306 matches!(Self::read_from(start), ReadFrom::StateQuery)
307 }
308 Unbounded => false,
309 };
310
311 let merged: BTreeMap<EncodedKey, PendingWrite> = inner
312 .pending
313 .range((range.start.as_ref(), range.end.as_ref()))
314 .map(|(k, v)| (k.clone(), v.clone()))
315 .collect();
316 let pending_vec: Vec<(EncodedKey, PendingWrite)> = merged.into_iter().collect();
317
318 if is_state_range {
319 let state_items: Vec<Result<MultiVersionRow>> = state
320 .iter()
321 .filter(|(k, _)| range.contains(k))
322 .map(|(k, v)| {
323 Ok(MultiVersionRow {
324 key: k.clone(),
325 row: v.clone(),
326 version: inner.version,
327 })
328 })
329 .collect();
330 let v = inner.version;
331
332 let mut sorted_items = state_items;
333 sorted_items.sort_by(|a, b| match (a, b) {
334 (Ok(a), Ok(b)) => a.key.cmp(&b.key),
335 _ => Ordering::Equal,
336 });
337 Box::new(flow_merge_pending_iterator(pending_vec, sorted_items.into_iter(), v))
338 } else {
339 let storage_iter = inner.query.range(range, scope, batch_size);
340 let v = inner.version;
341 Box::new(flow_merge_pending_iterator(pending_vec, storage_iter, v))
342 }
343 }
344 }
345 }
346
347 pub fn range_rev(
348 &mut self,
349 range: EncodedKeyRange,
350 scope: RangeScope,
351 batch_size: usize,
352 ) -> Box<dyn Iterator<Item = Result<MultiVersionRow>> + Send + '_> {
353 match self {
354 Self::Deferred {
355 inner,
356 ..
357 }
358 | Self::Committing {
359 inner,
360 ..
361 } => {
362 let merged: BTreeMap<EncodedKey, PendingWrite> = inner
363 .pending
364 .range((range.start.as_ref(), range.end.as_ref()))
365 .map(|(k, v)| (k.clone(), v.clone()))
366 .collect();
367 let pending_vec: Vec<(EncodedKey, PendingWrite)> = merged.into_iter().rev().collect();
368
369 let query = match range.start.as_ref() {
370 Included(start) | Excluded(start) => match Self::read_from(start) {
371 ReadFrom::StateQuery => inner.state_query.as_ref().unwrap(),
372 ReadFrom::Query => &inner.query,
373 ReadFrom::DictionaryQuery => {
374 inner.dictionary_query.as_ref().unwrap_or(&inner.query)
375 }
376 },
377 Unbounded => &inner.query,
378 };
379
380 let storage_iter = query.range_rev(range, scope, batch_size);
381 let v = inner.version;
382 Box::new(flow_merge_pending_iterator_rev(pending_vec, storage_iter, v))
383 }
384 Self::Transactional {
385 inner,
386 base_pending,
387 ..
388 } => {
389 let mut merged: BTreeMap<EncodedKey, PendingWrite> = base_pending
390 .range((range.start.as_ref(), range.end.as_ref()))
391 .map(|(k, v)| (k.clone(), v.clone()))
392 .collect();
393 for (k, v) in inner.pending.range((range.start.as_ref(), range.end.as_ref())) {
394 merged.insert(k.clone(), v.clone());
395 }
396 let pending_vec: Vec<(EncodedKey, PendingWrite)> = merged.into_iter().rev().collect();
397
398 let query = match range.start.as_ref() {
399 Included(start) | Excluded(start) => match Self::read_from(start) {
400 ReadFrom::StateQuery => inner.state_query.as_ref().unwrap(),
401 ReadFrom::Query => &inner.query,
402 ReadFrom::DictionaryQuery => {
403 inner.dictionary_query.as_ref().unwrap_or(&inner.query)
404 }
405 },
406 Unbounded => &inner.query,
407 };
408
409 let storage_iter = query.range_rev(range, scope, batch_size);
410 let v = inner.version;
411 Box::new(flow_merge_pending_iterator_rev(pending_vec, storage_iter, v))
412 }
413 Self::Ephemeral {
414 inner,
415 state,
416 } => {
417 let is_state_range = match range.start.as_ref() {
418 Included(start) | Excluded(start) => {
419 matches!(Self::read_from(start), ReadFrom::StateQuery)
420 }
421 Unbounded => false,
422 };
423
424 let merged: BTreeMap<EncodedKey, PendingWrite> = inner
425 .pending
426 .range((range.start.as_ref(), range.end.as_ref()))
427 .map(|(k, v)| (k.clone(), v.clone()))
428 .collect();
429 let pending_vec: Vec<(EncodedKey, PendingWrite)> = merged.into_iter().rev().collect();
430
431 if is_state_range {
432 let mut state_items: Vec<Result<MultiVersionRow>> = state
433 .iter()
434 .filter(|(k, _)| range.contains(k))
435 .map(|(k, v)| {
436 Ok(MultiVersionRow {
437 key: k.clone(),
438 row: v.clone(),
439 version: inner.version,
440 })
441 })
442 .collect();
443 let v = inner.version;
444
445 state_items.sort_by(|a, b| match (a, b) {
446 (Ok(a), Ok(b)) => b.key.cmp(&a.key),
447 _ => Ordering::Equal,
448 });
449 Box::new(flow_merge_pending_iterator_rev(
450 pending_vec,
451 state_items.into_iter(),
452 v,
453 ))
454 } else {
455 let storage_iter = inner.query.range_rev(range, scope, batch_size);
456 let v = inner.version;
457 Box::new(flow_merge_pending_iterator_rev(pending_vec, storage_iter, v))
458 }
459 }
460 }
461 }
462}
463
464struct FlowMergePendingIterator<I>
465where
466 I: Iterator<Item = Result<MultiVersionRow>>,
467{
468 storage_iter: Peekable<I>,
469 pending_iter: Peekable<IntoIter<(EncodedKey, PendingWrite)>>,
470 version: CommitVersion,
471}
472
473impl<I> Iterator for FlowMergePendingIterator<I>
474where
475 I: Iterator<Item = Result<MultiVersionRow>>,
476{
477 type Item = Result<MultiVersionRow>;
478
479 fn next(&mut self) -> Option<Self::Item> {
480 loop {
481 let next_storage = self.storage_iter.peek();
482
483 match (self.pending_iter.peek(), next_storage) {
484 (Some((pending_key, _)), Some(storage_result)) => {
485 let storage_val = match storage_result {
486 Ok(v) => v,
487 Err(_) => {
488 let err = self.storage_iter.next().unwrap();
489 return Some(err);
490 }
491 };
492 let cmp = pending_key.cmp(&storage_val.key);
493
494 if matches!(cmp, Ordering::Less) {
495 let (key, value) = self.pending_iter.next().unwrap();
496 if let PendingWrite::Set(row) = value {
497 return Some(Ok(MultiVersionRow {
498 key,
499 row,
500 version: self.version,
501 }));
502 }
503 } else if matches!(cmp, Ordering::Equal) {
504 let (key, value) = self.pending_iter.next().unwrap();
505 self.storage_iter.next();
506 if let PendingWrite::Set(row) = value {
507 return Some(Ok(MultiVersionRow {
508 key,
509 row,
510 version: self.version,
511 }));
512 }
513 } else {
514 return Some(self.storage_iter.next().unwrap());
515 }
516 }
517 (Some(_), None) => {
518 let (key, value) = self.pending_iter.next().unwrap();
519 if let PendingWrite::Set(row) = value {
520 return Some(Ok(MultiVersionRow {
521 key,
522 row,
523 version: self.version,
524 }));
525 }
526 }
527 (None, Some(_)) => {
528 return Some(self.storage_iter.next().unwrap());
529 }
530 (None, None) => return None,
531 }
532 }
533 }
534}
535
536fn flow_merge_pending_iterator<I>(
537 pending: Vec<(EncodedKey, PendingWrite)>,
538 storage_iter: I,
539 version: CommitVersion,
540) -> FlowMergePendingIterator<I>
541where
542 I: Iterator<Item = Result<MultiVersionRow>>,
543{
544 FlowMergePendingIterator {
545 storage_iter: storage_iter.peekable(),
546 pending_iter: pending.into_iter().peekable(),
547 version,
548 }
549}
550
551struct FlowMergePendingIteratorRev<I>
552where
553 I: Iterator<Item = Result<MultiVersionRow>>,
554{
555 storage_iter: Peekable<I>,
556 pending_iter: Peekable<IntoIter<(EncodedKey, PendingWrite)>>,
557 version: CommitVersion,
558}
559
560impl<I> Iterator for FlowMergePendingIteratorRev<I>
561where
562 I: Iterator<Item = Result<MultiVersionRow>>,
563{
564 type Item = Result<MultiVersionRow>;
565
566 fn next(&mut self) -> Option<Self::Item> {
567 loop {
568 let next_storage = self.storage_iter.peek();
569
570 match (self.pending_iter.peek(), next_storage) {
571 (Some((pending_key, _)), Some(storage_result)) => {
572 let storage_val = match storage_result {
573 Ok(v) => v,
574 Err(_) => {
575 let err = self.storage_iter.next().unwrap();
576 return Some(err);
577 }
578 };
579 let cmp = pending_key.cmp(&storage_val.key);
580
581 if matches!(cmp, Ordering::Greater) {
582 let (key, value) = self.pending_iter.next().unwrap();
583 if let PendingWrite::Set(row) = value {
584 return Some(Ok(MultiVersionRow {
585 key,
586 row,
587 version: self.version,
588 }));
589 }
590 } else if matches!(cmp, Ordering::Equal) {
591 let (key, value) = self.pending_iter.next().unwrap();
592 self.storage_iter.next();
593 if let PendingWrite::Set(row) = value {
594 return Some(Ok(MultiVersionRow {
595 key,
596 row,
597 version: self.version,
598 }));
599 }
600 } else {
601 return Some(self.storage_iter.next().unwrap());
602 }
603 }
604 (Some(_), None) => {
605 let (key, value) = self.pending_iter.next().unwrap();
606 if let PendingWrite::Set(row) = value {
607 return Some(Ok(MultiVersionRow {
608 key,
609 row,
610 version: self.version,
611 }));
612 }
613 }
614 (None, Some(_)) => {
615 return Some(self.storage_iter.next().unwrap());
616 }
617 (None, None) => return None,
618 }
619 }
620 }
621}
622
623fn flow_merge_pending_iterator_rev<I>(
624 pending: Vec<(EncodedKey, PendingWrite)>,
625 storage_iter: I,
626 version: CommitVersion,
627) -> FlowMergePendingIteratorRev<I>
628where
629 I: Iterator<Item = Result<MultiVersionRow>>,
630{
631 FlowMergePendingIteratorRev {
632 storage_iter: storage_iter.peekable(),
633 pending_iter: pending.into_iter().peekable(),
634 version,
635 }
636}
637
638#[cfg(test)]
639pub mod tests {
640 use reifydb_catalog::catalog::Catalog;
641 use reifydb_codec::{
642 encoded::row::EncodedRow,
643 key::encoded::{EncodedKey, EncodedKeyRange},
644 };
645 use reifydb_engine::test_harness::TestEngine;
646 use reifydb_runtime::context::clock::{Clock, MockClock};
647 use reifydb_transaction::interceptor::interceptors::Interceptors;
648 use reifydb_value::{util::cowvec::CowVec, value::identity::IdentityId};
649
650 use super::*;
651 use crate::operator::stateful::test_utils::test::create_test_transaction;
652
653 fn make_key(s: &str) -> EncodedKey {
654 EncodedKey::new(s.as_bytes().to_vec())
655 }
656
657 fn make_value(s: &str) -> EncodedRow {
658 EncodedRow(CowVec::new(s.as_bytes().to_vec()))
659 }
660
661 #[test]
662 fn test_get_from_pending() {
663 let parent = create_test_transaction();
664 let mut txn = FlowTransaction::deferred(
665 &parent,
666 CommitVersion(1),
667 Catalog::testing(),
668 Interceptors::new(),
669 Clock::Mock(MockClock::from_millis(1000)),
670 );
671
672 let key = make_key("key1");
673 let value = make_value("value1");
674
675 txn.set(&key, value.clone()).unwrap();
676
677 let result = txn.get(&key).unwrap();
679 assert_eq!(result, Some(value));
680 }
681
682 #[test]
683 fn test_get_from_committed() {
684 let t = TestEngine::new();
685
686 let key = make_key("key1");
687 let value = make_value("value1");
688
689 {
691 let mut cmd_txn = t.begin_admin(IdentityId::system()).unwrap();
692 cmd_txn.set(&key, value.clone()).unwrap();
693 cmd_txn.commit().unwrap();
694 }
695
696 let parent = t.begin_admin(IdentityId::system()).unwrap();
698 let version = parent.version();
699
700 let mut txn = FlowTransaction::deferred(
702 &parent,
703 version,
704 Catalog::testing(),
705 Interceptors::new(),
706 Clock::Mock(MockClock::from_millis(1000)),
707 );
708
709 let result = txn.get(&key).unwrap();
711 assert_eq!(result, Some(value));
712 }
713
714 #[test]
715 fn test_get_pending_shadows_committed() {
716 let mut parent = create_test_transaction();
717
718 let key = make_key("key1");
719 parent.set(&key, make_value("old")).unwrap();
720 let version = parent.version();
721
722 let mut txn = FlowTransaction::deferred(
723 &parent,
724 version,
725 Catalog::testing(),
726 Interceptors::new(),
727 Clock::Mock(MockClock::from_millis(1000)),
728 );
729
730 let new_value = make_value("new");
732 txn.set(&key, new_value.clone()).unwrap();
733
734 let result = txn.get(&key).unwrap();
736 assert_eq!(result, Some(new_value));
737 }
738
739 #[test]
740 fn test_get_removed_returns_none() {
741 let mut parent = create_test_transaction();
742
743 let key = make_key("key1");
744 parent.set(&key, make_value("value1")).unwrap();
745 let version = parent.version();
746
747 let mut txn = FlowTransaction::deferred(
748 &parent,
749 version,
750 Catalog::testing(),
751 Interceptors::new(),
752 Clock::Mock(MockClock::from_millis(1000)),
753 );
754
755 txn.remove(&key).unwrap();
757
758 let result = txn.get(&key).unwrap();
760 assert_eq!(result, None);
761 }
762
763 #[test]
764 fn test_get_nonexistent_key() {
765 let parent = create_test_transaction();
766 let mut txn = FlowTransaction::deferred(
767 &parent,
768 CommitVersion(1),
769 Catalog::testing(),
770 Interceptors::new(),
771 Clock::Mock(MockClock::from_millis(1000)),
772 );
773
774 let result = txn.get(&make_key("missing")).unwrap();
775 assert_eq!(result, None);
776 }
777
778 #[test]
779 fn test_contains_key_pending() {
780 let parent = create_test_transaction();
781 let mut txn = FlowTransaction::deferred(
782 &parent,
783 CommitVersion(1),
784 Catalog::testing(),
785 Interceptors::new(),
786 Clock::Mock(MockClock::from_millis(1000)),
787 );
788
789 let key = make_key("key1");
790 txn.set(&key, make_value("value1")).unwrap();
791
792 assert!(txn.contains_key(&key).unwrap());
793 }
794
795 #[test]
796 fn test_contains_key_committed() {
797 let t = TestEngine::new();
798
799 let key = make_key("key1");
800
801 {
803 let mut cmd_txn = t.begin_admin(IdentityId::system()).unwrap();
804 cmd_txn.set(&key, make_value("value1")).unwrap();
805 cmd_txn.commit().unwrap();
806 }
807
808 let parent = t.begin_admin(IdentityId::system()).unwrap();
810 let version = parent.version();
811 let mut txn = FlowTransaction::deferred(
812 &parent,
813 version,
814 Catalog::testing(),
815 Interceptors::new(),
816 Clock::Mock(MockClock::from_millis(1000)),
817 );
818
819 assert!(txn.contains_key(&key).unwrap());
820 }
821
822 #[test]
823 fn test_contains_key_removed_returns_false() {
824 let mut parent = create_test_transaction();
825
826 let key = make_key("key1");
827 parent.set(&key, make_value("value1")).unwrap();
828 let version = parent.version();
829
830 let mut txn = FlowTransaction::deferred(
831 &parent,
832 version,
833 Catalog::testing(),
834 Interceptors::new(),
835 Clock::Mock(MockClock::from_millis(1000)),
836 );
837 txn.remove(&key).unwrap();
838
839 assert!(!txn.contains_key(&key).unwrap());
840 }
841
842 #[test]
843 fn test_contains_key_nonexistent() {
844 let parent = create_test_transaction();
845 let mut txn = FlowTransaction::deferred(
846 &parent,
847 CommitVersion(1),
848 Catalog::testing(),
849 Interceptors::new(),
850 Clock::Mock(MockClock::from_millis(1000)),
851 );
852
853 assert!(!txn.contains_key(&make_key("missing")).unwrap());
854 }
855
856 #[test]
857 fn test_scan_empty() {
858 let parent = create_test_transaction();
859 let mut txn = FlowTransaction::deferred(
860 &parent,
861 CommitVersion(1),
862 Catalog::testing(),
863 Interceptors::new(),
864 Clock::Mock(MockClock::from_millis(1000)),
865 );
866
867 let mut iter = txn.range(EncodedKeyRange::all(), RangeScope::All, 1024);
868 assert!(iter.next().is_none());
869 }
870
871 #[test]
872 fn test_scan_only_pending() {
873 let parent = create_test_transaction();
874 let mut txn = FlowTransaction::deferred(
875 &parent,
876 CommitVersion(1),
877 Catalog::testing(),
878 Interceptors::new(),
879 Clock::Mock(MockClock::from_millis(1000)),
880 );
881
882 txn.set(&make_key("b"), make_value("2")).unwrap();
883 txn.set(&make_key("a"), make_value("1")).unwrap();
884 txn.set(&make_key("c"), make_value("3")).unwrap();
885
886 let items: Vec<_> =
887 txn.range(EncodedKeyRange::all(), RangeScope::All, 1024).collect::<Result<Vec<_>>>().unwrap();
888
889 assert_eq!(items.len(), 3);
891 assert_eq!(items[0].key, make_key("a"));
892 assert_eq!(items[1].key, make_key("b"));
893 assert_eq!(items[2].key, make_key("c"));
894 }
895
896 #[test]
897 fn test_scan_filters_removes() {
898 let parent = create_test_transaction();
899 let mut txn = FlowTransaction::deferred(
900 &parent,
901 CommitVersion(1),
902 Catalog::testing(),
903 Interceptors::new(),
904 Clock::Mock(MockClock::from_millis(1000)),
905 );
906
907 txn.set(&make_key("a"), make_value("1")).unwrap();
908 txn.remove(&make_key("b")).unwrap();
909 txn.set(&make_key("c"), make_value("3")).unwrap();
910
911 let items: Vec<_> =
912 txn.range(EncodedKeyRange::all(), RangeScope::All, 1024).collect::<Result<Vec<_>>>().unwrap();
913
914 assert_eq!(items.len(), 2);
916 assert_eq!(items[0].key, make_key("a"));
917 assert_eq!(items[1].key, make_key("c"));
918 }
919
920 #[test]
921 fn test_range_empty() {
922 let parent = create_test_transaction();
923 let mut txn = FlowTransaction::deferred(
924 &parent,
925 CommitVersion(1),
926 Catalog::testing(),
927 Interceptors::new(),
928 Clock::Mock(MockClock::from_millis(1000)),
929 );
930
931 let range = EncodedKeyRange::start_end(Some(make_key("a")), Some(make_key("z")));
932 let mut iter = txn.range(range, RangeScope::All, 1024);
933 assert!(iter.next().is_none());
934 }
935
936 #[test]
937 fn test_range_only_pending() {
938 let parent = create_test_transaction();
939 let mut txn = FlowTransaction::deferred(
940 &parent,
941 CommitVersion(1),
942 Catalog::testing(),
943 Interceptors::new(),
944 Clock::Mock(MockClock::from_millis(1000)),
945 );
946
947 txn.set(&make_key("a"), make_value("1")).unwrap();
948 txn.set(&make_key("b"), make_value("2")).unwrap();
949 txn.set(&make_key("c"), make_value("3")).unwrap();
950 txn.set(&make_key("d"), make_value("4")).unwrap();
951
952 let range = EncodedKeyRange::new(Included(make_key("b")), Excluded(make_key("d")));
953 let items: Vec<_> = txn.range(range, RangeScope::All, 1024).collect::<Result<Vec<_>>>().unwrap();
954
955 assert_eq!(items.len(), 2);
957 assert_eq!(items[0].key, make_key("b"));
958 assert_eq!(items[1].key, make_key("c"));
959 }
960
961 #[test]
962 fn test_prefix_empty() {
963 let parent = create_test_transaction();
964 let mut txn = FlowTransaction::deferred(
965 &parent,
966 CommitVersion(1),
967 Catalog::testing(),
968 Interceptors::new(),
969 Clock::Mock(MockClock::from_millis(1000)),
970 );
971
972 let prefix = make_key("test_");
973 let iter = txn.prefix(&prefix).unwrap();
974 assert!(iter.items.into_iter().next().is_none());
975 }
976
977 #[test]
978 fn test_prefix_only_pending() {
979 let parent = create_test_transaction();
980 let mut txn = FlowTransaction::deferred(
981 &parent,
982 CommitVersion(1),
983 Catalog::testing(),
984 Interceptors::new(),
985 Clock::Mock(MockClock::from_millis(1000)),
986 );
987
988 txn.set(&make_key("test_a"), make_value("1")).unwrap();
989 txn.set(&make_key("test_b"), make_value("2")).unwrap();
990 txn.set(&make_key("other_c"), make_value("3")).unwrap();
991
992 let prefix = make_key("test_");
993 let iter = txn.prefix(&prefix).unwrap();
994 let items: Vec<_> = iter.items.into_iter().collect();
995
996 assert_eq!(items.len(), 2);
998 assert_eq!(items[0].key, make_key("test_a"));
999 assert_eq!(items[1].key, make_key("test_b"));
1000 }
1001}