Skip to main content

reifydb_sub_flow/transaction/
read.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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		// Should get value from pending buffer
678		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		// Set value in first transaction and commit
690		{
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		// Create new command transaction to read committed data
697		let parent = t.begin_admin(IdentityId::system()).unwrap();
698		let version = parent.version();
699
700		// Create FlowTransaction - should see committed value
701		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		// Should get value from query transaction
710		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		// Override with new value in pending
731		let new_value = make_value("new");
732		txn.set(&key, new_value.clone()).unwrap();
733
734		// Should get new value from pending, not old value from committed
735		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		// Remove in pending
756		txn.remove(&key).unwrap();
757
758		// Should return None even though it exists in committed
759		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		// Set value in first transaction and commit
802		{
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		// Create new command transaction
809		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		// Should be in sorted order
890		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		// Should only have 2 items (remove filtered out)
915		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		// Should only include b and c (not d, exclusive end)
956		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		// Should only include keys with prefix "test_"
997		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}