1use 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}