Skip to main content

reifydb_core/interface/
store.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use reifydb_value::{Result, util::cowvec::CowVec};
5
6use crate::{
7	common::CommitVersion,
8	delta::Delta,
9	encoded::{
10		key::{EncodedKey, EncodedKeyRange},
11		row::EncodedRow,
12	},
13	interface::catalog::{flow::FlowNodeId, shape::ShapeId},
14	key::{
15		EncodableKeyRange, Key, flow_node_internal_state::FlowNodeInternalStateKeyRange,
16		flow_node_state::FlowNodeStateKeyRange, kind::KeyKind, row::RowKeyRange,
17	},
18};
19
20#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
21pub enum Tier {
22	Buffer,
23	Persistent,
24}
25
26#[derive(Debug, Clone, Copy, PartialEq, Eq, Hash)]
27pub enum EntryKind {
28	Multi,
29
30	Source(ShapeId),
31
32	Operator(FlowNodeId),
33}
34
35pub fn classify_key(key: &EncodedKey) -> EntryKind {
36	match Key::decode(key) {
37		Some(Key::Row(row_key)) => EntryKind::Source(row_key.shape),
38		Some(Key::FlowNodeState(state_key)) => EntryKind::Operator(state_key.node),
39		Some(Key::FlowNodeInternalState(internal_key)) => EntryKind::Operator(internal_key.node),
40		_ => EntryKind::Multi,
41	}
42}
43
44pub fn is_single_version_semantics_key(key: &EncodedKey) -> bool {
45	Key::kind(key).is_some_and(|kind| matches!(kind, KeyKind::FlowNodeState | KeyKind::FlowNodeInternalState))
46}
47
48pub fn classify_range(range: &EncodedKeyRange) -> Option<EntryKind> {
49	if let (Some(start), Some(_end)) = RowKeyRange::decode(range) {
50		return Some(EntryKind::Source(start.shape));
51	}
52
53	if let (Some(start), Some(_end)) = FlowNodeStateKeyRange::decode(range) {
54		return Some(EntryKind::Operator(start.node));
55	}
56
57	if let (Some(start), Some(_end)) = FlowNodeInternalStateKeyRange::decode(range) {
58		return Some(EntryKind::Operator(start.node));
59	}
60
61	None
62}
63
64#[derive(Debug, Clone)]
65pub struct MultiVersionRow {
66	pub key: EncodedKey,
67	pub row: EncodedRow,
68	pub version: CommitVersion,
69}
70
71#[derive(Debug, Clone)]
72pub struct SingleVersionRow {
73	pub key: EncodedKey,
74	pub row: EncodedRow,
75}
76
77#[derive(Debug, Clone)]
78pub struct MultiVersionBatch {
79	pub items: Vec<MultiVersionRow>,
80
81	pub has_more: bool,
82}
83
84impl MultiVersionBatch {
85	pub fn empty() -> Self {
86		Self {
87			items: Vec::new(),
88			has_more: false,
89		}
90	}
91
92	pub fn is_empty(&self) -> bool {
93		self.items.is_empty()
94	}
95}
96
97pub trait MultiVersionCommit: Send + Sync {
98	fn commit(&self, deltas: CowVec<Delta>, version: CommitVersion) -> Result<()>;
99}
100
101#[derive(Debug, Clone, Copy, Default, PartialEq, Eq)]
102pub struct ReadOptions {
103	pub bypass_buffer: bool,
104}
105
106pub trait MultiVersionGet: Send + Sync {
107	fn get(&self, key: &EncodedKey, version: CommitVersion) -> Result<Option<MultiVersionRow>>;
108
109	fn get_with_options(
110		&self,
111		key: &EncodedKey,
112		version: CommitVersion,
113		_options: ReadOptions,
114	) -> Result<Option<MultiVersionRow>> {
115		self.get(key, version)
116	}
117}
118
119pub trait MultiVersionContains: Send + Sync {
120	fn contains(&self, key: &EncodedKey, version: CommitVersion) -> Result<bool>;
121
122	fn contains_with_options(
123		&self,
124		key: &EncodedKey,
125		version: CommitVersion,
126		_options: ReadOptions,
127	) -> Result<bool> {
128		self.contains(key, version)
129	}
130}
131
132pub trait MultiVersionGetPrevious: Send + Sync {
133	fn get_previous_version(
134		&self,
135		key: &EncodedKey,
136		before_version: CommitVersion,
137	) -> Result<Option<MultiVersionRow>>;
138}
139
140pub trait MultiVersionStore:
141	Send + Sync + Clone + MultiVersionCommit + MultiVersionGet + MultiVersionGetPrevious + MultiVersionContains + 'static
142{
143}
144
145#[derive(Debug, Clone)]
146pub struct SingleVersionBatch {
147	pub items: Vec<SingleVersionRow>,
148
149	pub has_more: bool,
150}
151
152impl SingleVersionBatch {
153	pub fn empty() -> Self {
154		Self {
155			items: Vec::new(),
156			has_more: false,
157		}
158	}
159
160	pub fn is_empty(&self) -> bool {
161		self.items.is_empty()
162	}
163}
164
165pub trait SingleVersionCommit: Send + Sync {
166	fn commit(&mut self, deltas: CowVec<Delta>) -> Result<()>;
167}
168
169pub trait SingleVersionGet: Send + Sync {
170	fn get(&self, key: &EncodedKey) -> Result<Option<SingleVersionRow>>;
171}
172
173pub trait SingleVersionContains: Send + Sync {
174	fn contains(&self, key: &EncodedKey) -> Result<bool>;
175}
176
177pub trait SingleVersionSet: SingleVersionCommit {
178	fn set(&mut self, key: &EncodedKey, row: EncodedRow) -> Result<()> {
179		Self::commit(
180			self,
181			CowVec::new(vec![Delta::Set {
182				key: key.clone(),
183				row: row.clone(),
184			}]),
185		)
186	}
187}
188
189pub trait SingleVersionRemove: SingleVersionCommit {
190	fn unset(&mut self, key: &EncodedKey, row: EncodedRow) -> Result<()> {
191		Self::commit(
192			self,
193			CowVec::new(vec![Delta::Unset {
194				key: key.clone(),
195				row,
196			}]),
197		)
198	}
199
200	fn remove(&mut self, key: &EncodedKey) -> Result<()> {
201		Self::commit(
202			self,
203			CowVec::new(vec![Delta::Remove {
204				key: key.clone(),
205			}]),
206		)
207	}
208}
209
210pub trait SingleVersionRange: Send + Sync {
211	fn range_batch(&self, range: EncodedKeyRange, batch_size: u64) -> Result<SingleVersionBatch>;
212
213	fn range(&self, range: EncodedKeyRange) -> Result<SingleVersionBatch> {
214		self.range_batch(range, 1024)
215	}
216
217	fn prefix(&self, prefix: &EncodedKey) -> Result<SingleVersionBatch> {
218		self.range(EncodedKeyRange::prefix(prefix))
219	}
220}
221
222pub trait SingleVersionRangeRev: Send + Sync {
223	fn range_rev_batch(&self, range: EncodedKeyRange, batch_size: u64) -> Result<SingleVersionBatch>;
224
225	fn range_rev(&self, range: EncodedKeyRange) -> Result<SingleVersionBatch> {
226		self.range_rev_batch(range, 1024)
227	}
228
229	fn prefix_rev(&self, prefix: &EncodedKey) -> Result<SingleVersionBatch> {
230		self.range_rev(EncodedKeyRange::prefix(prefix))
231	}
232}
233
234pub trait SingleVersionStore:
235	Send
236	+ Sync
237	+ Clone
238	+ SingleVersionCommit
239	+ SingleVersionGet
240	+ SingleVersionContains
241	+ SingleVersionSet
242	+ SingleVersionRemove
243	+ SingleVersionRange
244	+ SingleVersionRangeRev
245	+ 'static
246{
247}