Skip to main content

reifydb_transaction/multi/transaction/
read.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}