reifydb_transaction/multi/transaction/
read.rs1use std::{collections::HashMap, ops::Bound};
5
6use reifydb_codec::key::encoded::{EncodedKey, EncodedKeyRange};
7use reifydb_core::{
8 common::CommitVersion,
9 interface::{
10 catalog::storage::StorageId,
11 store::{MultiVersionBatch, MultiVersionRow},
12 },
13 key::{
14 any::TaggedKey,
15 bound::TaggedKeyBoundRange,
16 row::{StoragePartitionedRowKey, StorageRowKey},
17 },
18};
19use reifydb_value::Result;
20use tracing::instrument;
21
22use super::{MultiTransaction, manager::TransactionManagerQuery, version::StandardVersionProvider};
23use crate::multi::{RangeScope, lease::VersionLeaseGuard, types::TransactionValue};
24
25pub struct MultiReadTransaction {
26 pub(crate) engine: MultiTransaction,
27 pub(crate) tm: TransactionManagerQuery<StandardVersionProvider>,
28 #[allow(dead_code)]
29 pub(crate) lease: Option<VersionLeaseGuard>,
30}
31
32impl MultiReadTransaction {
33 pub fn new(engine: MultiTransaction, version: Option<CommitVersion>) -> Result<Self> {
34 let tm = engine.tm.query(version)?;
35 Ok(Self {
36 engine,
37 tm,
38 lease: None,
39 })
40 }
41
42 pub fn new_with_lease(engine: MultiTransaction, lease: VersionLeaseGuard) -> Result<Self> {
43 let version = lease.version();
44 let tm = engine.tm.query(Some(version))?;
45 Ok(Self {
46 engine,
47 tm,
48 lease: Some(lease),
49 })
50 }
51}
52
53impl MultiReadTransaction {
54 pub fn version(&self) -> CommitVersion {
55 self.tm.version()
56 }
57
58 pub fn read_as_of_version_exclusive(&mut self, version: CommitVersion) {
59 self.tm.read_as_of_version_exclusive(version);
60 }
61
62 pub fn read_as_of_version_inclusive(&mut self, version: CommitVersion) {
63 self.read_as_of_version_exclusive(CommitVersion(version.0 + 1))
64 }
65
66 pub fn get<K: Into<TaggedKey> + Clone>(&self, key: &K) -> Result<Option<TransactionValue>> {
67 let version = self.tm.version();
68 Ok(self.engine.get(&key.clone().into(), version)?.map(Into::into))
69 }
70
71 #[instrument(name = "transaction::get_many", level = "trace", skip(self, keys), fields(key_count = keys.len()))]
72 pub fn get_many(&self, keys: &[EncodedKey]) -> Result<HashMap<EncodedKey, MultiVersionRow>> {
73 let version = self.tm.version();
74 self.engine.store.get_many(keys, version)
75 }
76
77 pub fn contains<K: Into<TaggedKey> + Clone>(&self, key: &K) -> Result<bool> {
78 let version = self.tm.version();
79 self.engine.contains_key(&key.clone().into(), version)
80 }
81
82 pub fn scan(&self) -> Result<MultiVersionBatch<TaggedKey>> {
83 let items: Vec<_> = self
84 .range_encoded(EncodedKeyRange::all(), RangeScope::All, 1024)
85 .collect::<Result<Vec<_>>>()?;
86 Ok(MultiVersionBatch {
87 items,
88 has_more: false,
89 })
90 }
91
92 pub fn prefix(&self, prefix: &EncodedKey) -> Result<MultiVersionBatch<TaggedKey>> {
93 let items: Vec<_> = self
94 .range_encoded(EncodedKeyRange::prefix(prefix), RangeScope::All, 1024)
95 .collect::<Result<Vec<_>>>()?;
96 Ok(MultiVersionBatch {
97 items,
98 has_more: false,
99 })
100 }
101
102 pub fn prefix_rev(&self, prefix: &EncodedKey) -> Result<MultiVersionBatch<TaggedKey>> {
103 let items: Vec<_> = self
104 .range_rev_encoded(EncodedKeyRange::prefix(prefix), RangeScope::All, 1024)
105 .collect::<Result<Vec<_>>>()?;
106 Ok(MultiVersionBatch {
107 items,
108 has_more: false,
109 })
110 }
111
112 pub fn range_encoded(
113 &self,
114 range: EncodedKeyRange,
115 scope: RangeScope,
116 batch_size: usize,
117 ) -> Box<dyn Iterator<Item = Result<MultiVersionRow<TaggedKey>>> + Send + '_> {
118 let multi_scope = scope.into_multi(self.tm.version());
119 Box::new(self.engine.store.range(range, multi_scope, batch_size))
120 }
121
122 pub fn range_rev_encoded(
123 &self,
124 range: EncodedKeyRange,
125 scope: RangeScope,
126 batch_size: usize,
127 ) -> Box<dyn Iterator<Item = Result<MultiVersionRow<TaggedKey>>> + Send + '_> {
128 let multi_scope = scope.into_multi(self.tm.version());
129 Box::new(self.engine.store.range_rev(range, multi_scope, batch_size))
130 }
131
132 pub fn range(
133 &self,
134 range: TaggedKeyBoundRange,
135 scope: RangeScope,
136 batch_size: usize,
137 ) -> Box<dyn Iterator<Item = Result<MultiVersionRow<TaggedKey>>> + Send + '_> {
138 let multi_scope = scope.into_multi(self.tm.version());
139 Box::new(self.engine.store.range(range.encode(), multi_scope, batch_size))
140 }
141
142 pub fn range_row(
143 &self,
144 storage: StorageId,
145 start: Bound<StorageRowKey>,
146 end: Bound<StorageRowKey>,
147 scope: RangeScope,
148 batch_size: usize,
149 ) -> Box<dyn Iterator<Item = Result<MultiVersionRow<StorageRowKey>>> + Send + '_> {
150 let multi_scope = scope.into_multi(self.tm.version());
151 Box::new(self.engine.store.range_row(storage, start, end, multi_scope, batch_size))
152 }
153
154 pub fn range_partitioned_row(
155 &self,
156 storage: StorageId,
157 start: Bound<StoragePartitionedRowKey>,
158 end: Bound<StoragePartitionedRowKey>,
159 scope: RangeScope,
160 batch_size: usize,
161 ) -> Box<dyn Iterator<Item = Result<MultiVersionRow<StoragePartitionedRowKey>>> + Send + '_> {
162 let multi_scope = scope.into_multi(self.tm.version());
163 Box::new(self.engine.store.range_partitioned_row(storage, start, end, multi_scope, batch_size))
164 }
165
166 pub fn range_rev(
167 &self,
168 range: TaggedKeyBoundRange,
169 scope: RangeScope,
170 batch_size: usize,
171 ) -> Box<dyn Iterator<Item = Result<MultiVersionRow<TaggedKey>>> + Send + '_> {
172 let multi_scope = scope.into_multi(self.tm.version());
173 Box::new(self.engine.store.range_rev(range.encode(), multi_scope, batch_size))
174 }
175}
176
177impl Clone for MultiReadTransaction {
178 fn clone(&self) -> Self {
179 Self {
180 engine: self.engine.clone(),
181 tm: self.tm.clone(),
182 lease: None,
183 }
184 }
185}