1use reifydb_codec::{
5 encoded::{row::EncodedRow, shape::RowShape},
6 key::encoded::EncodedKey,
7};
8use reifydb_core::{
9 common::CommitVersion,
10 interface::{
11 catalog::{ringbuffer::RingBuffer, shape::ShapeId},
12 change::{Change, ChangeOrigin, Diff},
13 },
14 key::{
15 partitioned_row::{PartitionedRowKey, RowLocator},
16 row::RowKey,
17 },
18 row::row_shape_from_columns,
19 value::column::columns::Columns,
20};
21use reifydb_transaction::{
22 interceptor::ringbuffer_row::RingBufferRowInterceptor,
23 transaction::{Transaction, admin::AdminTransaction, command::CommandTransaction},
24};
25use reifydb_value::{
26 util::cowvec::CowVec,
27 value::{datetime::DateTime, partition::Partition, row_number::RowNumber},
28};
29use smallvec::smallvec;
30
31use crate::Result;
32
33fn ringbuffer_key(ringbuffer: &RingBuffer, partition: Option<Partition>, row_number: RowNumber) -> EncodedKey {
34 match partition {
35 None => RowKey::encoded(ringbuffer.id, row_number),
36 Some(partition) => PartitionedRowKey::encoded(
37 ShapeId::ringbuffer(ringbuffer.id),
38 partition,
39 RowLocator::Row(row_number),
40 ),
41 }
42}
43
44fn build_ringbuffer_insert_change(
45 rb: &RingBuffer,
46 shape: &RowShape,
47 row_number: RowNumber,
48 encoded: &EncodedRow,
49) -> Change {
50 let ids = [row_number];
51 let rows = [encoded.clone()];
52 Change {
53 origin: ChangeOrigin::Shape(ShapeId::ringbuffer(rb.id)),
54 version: CommitVersion(0),
55 diffs: smallvec![Diff::insert(Columns::from_encoded_rows(shape, &ids, &rows))],
56 changed_at: DateTime::default(),
57 }
58}
59
60fn build_ringbuffer_update_change(
61 rb: &RingBuffer,
62 row_number: RowNumber,
63 pre: &EncodedRow,
64 post: &EncodedRow,
65) -> Change {
66 let shape = row_shape_from_columns(&rb.columns);
67 let ids = [row_number];
68 let pres = [pre.clone()];
69 let posts = [post.clone()];
70 Change {
71 origin: ChangeOrigin::Shape(ShapeId::ringbuffer(rb.id)),
72 version: CommitVersion(0),
73 diffs: smallvec![Diff::update(
74 Columns::from_encoded_rows(&shape, &ids, &pres),
75 Columns::from_encoded_rows(&shape, &ids, &posts),
76 )],
77 changed_at: DateTime::default(),
78 }
79}
80
81fn build_ringbuffer_remove_change(rb: &RingBuffer, row_number: RowNumber, encoded: &EncodedRow) -> Change {
82 let shape = row_shape_from_columns(&rb.columns);
83 let ids = [row_number];
84 let rows = [encoded.clone()];
85 Change {
86 origin: ChangeOrigin::Shape(ShapeId::ringbuffer(rb.id)),
87 version: CommitVersion(0),
88 diffs: smallvec![Diff::remove(Columns::from_encoded_rows(&shape, &ids, &rows))],
89 changed_at: DateTime::default(),
90 }
91}
92
93pub trait RingBufferOperations {
94 fn insert_ringbuffer(&mut self, ringbuffer: RingBuffer, row: EncodedRow) -> Result<RowNumber>;
95
96 fn insert_ringbuffer_at(
97 &mut self,
98 ringbuffer: &RingBuffer,
99 shape: &RowShape,
100 partition: Option<Partition>,
101 row_number: RowNumber,
102 row: EncodedRow,
103 ) -> Result<EncodedRow>;
104
105 fn update_ringbuffer(
106 &mut self,
107 ringbuffer: RingBuffer,
108 partition: Option<Partition>,
109 id: RowNumber,
110 row: EncodedRow,
111 ) -> Result<EncodedRow>;
112
113 fn remove_from_ringbuffer(
114 &mut self,
115 ringbuffer: &RingBuffer,
116 partition: Option<Partition>,
117 id: RowNumber,
118 ) -> Result<EncodedRow>;
119}
120
121impl RingBufferOperations for CommandTransaction {
122 fn insert_ringbuffer(&mut self, _ringbuffer: RingBuffer, _row: EncodedRow) -> Result<RowNumber> {
123 unimplemented!(
124 "Ring buffer insert must be called with explicit row_number through insert_ringbuffer_at"
125 )
126 }
127
128 fn insert_ringbuffer_at(
129 &mut self,
130 ringbuffer: &RingBuffer,
131 shape: &RowShape,
132 partition: Option<Partition>,
133 row_number: RowNumber,
134 row: EncodedRow,
135 ) -> Result<EncodedRow> {
136 let key = ringbuffer_key(ringbuffer, partition, row_number);
137
138 let pre = self.get(&key)?.map(|v| v.row);
139
140 if let Some(ref existing) = pre {
141 let ids = [row_number];
142 let existing_rows = [existing.clone()];
143 RingBufferRowInterceptor::pre_delete(self, ringbuffer, &ids)?;
144 RingBufferRowInterceptor::post_delete(self, ringbuffer, &ids, &existing_rows)?;
145 }
146
147 let mut rows_buf = [row];
148 RingBufferRowInterceptor::pre_insert(self, ringbuffer, &mut rows_buf)?;
149 let [row] = rows_buf;
150
151 self.set(&key, row.clone())?;
152
153 let ids = [row_number];
154 let rows = [row.clone()];
155 RingBufferRowInterceptor::post_insert(self, ringbuffer, &ids, &rows)?;
156
157 if let Some(pre_row) = pre.as_ref() {
158 self.track_flow_change(build_ringbuffer_update_change(ringbuffer, row_number, pre_row, &row));
159 } else {
160 self.track_flow_change(build_ringbuffer_insert_change(ringbuffer, shape, row_number, &row));
161 }
162
163 Ok(row)
164 }
165
166 fn update_ringbuffer(
167 &mut self,
168 ringbuffer: RingBuffer,
169 partition: Option<Partition>,
170 id: RowNumber,
171 row: EncodedRow,
172 ) -> Result<EncodedRow> {
173 let key = ringbuffer_key(&ringbuffer, partition, id);
174
175 let pre = match self.get(&key)? {
176 Some(v) => v.row,
177 None => return Ok(row),
178 };
179
180 let mut rows_buf = [row];
181 let ids = [id];
182 RingBufferRowInterceptor::pre_update(self, &ringbuffer, &ids, &mut rows_buf)?;
183 let [row] = rows_buf;
184
185 if self.get_committed(&key)?.is_some() {
186 self.mark_preexisting(&key)?;
187 }
188 self.set(&key, row.clone())?;
189
190 let posts = [row.clone()];
191 let pres = [pre.clone()];
192 RingBufferRowInterceptor::post_update(self, &ringbuffer, &ids, &posts, &pres)?;
193
194 self.track_flow_change(build_ringbuffer_update_change(&ringbuffer, id, &pre, &row));
195
196 Ok(row)
197 }
198
199 fn remove_from_ringbuffer(
200 &mut self,
201 ringbuffer: &RingBuffer,
202 partition: Option<Partition>,
203 id: RowNumber,
204 ) -> Result<EncodedRow> {
205 let key = ringbuffer_key(ringbuffer, partition, id);
206
207 let displayed = match self.get(&key)? {
208 Some(v) => v.row,
209 None => return Ok(EncodedRow(CowVec::new(vec![]))),
210 };
211 let committed = self.get_committed(&key)?.map(|v| v.row);
212
213 let ids = [id];
214 RingBufferRowInterceptor::pre_delete(self, ringbuffer, &ids)?;
215
216 let pre_for_cdc = committed.clone().unwrap_or_else(|| displayed.clone());
217
218 if committed.is_some() {
219 self.mark_preexisting(&key)?;
220 }
221 self.unset(&key, pre_for_cdc.clone())?;
222
223 let pre_rows = [pre_for_cdc.clone()];
224 RingBufferRowInterceptor::post_delete(self, ringbuffer, &ids, &pre_rows)?;
225
226 self.track_flow_change(build_ringbuffer_remove_change(ringbuffer, id, &pre_for_cdc));
227
228 Ok(displayed)
229 }
230}
231
232impl RingBufferOperations for AdminTransaction {
233 fn insert_ringbuffer(&mut self, _ringbuffer: RingBuffer, _row: EncodedRow) -> Result<RowNumber> {
234 unimplemented!(
235 "Ring buffer insert must be called with explicit row_number through insert_ringbuffer_at"
236 )
237 }
238
239 fn insert_ringbuffer_at(
240 &mut self,
241 ringbuffer: &RingBuffer,
242 shape: &RowShape,
243 partition: Option<Partition>,
244 row_number: RowNumber,
245 row: EncodedRow,
246 ) -> Result<EncodedRow> {
247 let key = ringbuffer_key(ringbuffer, partition, row_number);
248
249 let pre = self.get(&key)?.map(|v| v.row);
250
251 if let Some(ref existing) = pre {
252 let ids = [row_number];
253 let existing_rows = [existing.clone()];
254 RingBufferRowInterceptor::pre_delete(self, ringbuffer, &ids)?;
255 RingBufferRowInterceptor::post_delete(self, ringbuffer, &ids, &existing_rows)?;
256 }
257
258 let mut rows_buf = [row];
259 RingBufferRowInterceptor::pre_insert(self, ringbuffer, &mut rows_buf)?;
260 let [row] = rows_buf;
261
262 self.set(&key, row.clone())?;
263
264 let ids = [row_number];
265 let rows = [row.clone()];
266 RingBufferRowInterceptor::post_insert(self, ringbuffer, &ids, &rows)?;
267
268 if let Some(pre_row) = pre.as_ref() {
269 self.track_flow_change(build_ringbuffer_update_change(ringbuffer, row_number, pre_row, &row));
270 } else {
271 self.track_flow_change(build_ringbuffer_insert_change(ringbuffer, shape, row_number, &row));
272 }
273
274 Ok(row)
275 }
276
277 fn update_ringbuffer(
278 &mut self,
279 ringbuffer: RingBuffer,
280 partition: Option<Partition>,
281 id: RowNumber,
282 row: EncodedRow,
283 ) -> Result<EncodedRow> {
284 let key = ringbuffer_key(&ringbuffer, partition, id);
285
286 let pre = match self.get(&key)? {
287 Some(v) => v.row,
288 None => return Ok(row),
289 };
290
291 let mut rows_buf = [row];
292 let ids = [id];
293 RingBufferRowInterceptor::pre_update(self, &ringbuffer, &ids, &mut rows_buf)?;
294 let [row] = rows_buf;
295
296 if self.get_committed(&key)?.is_some() {
297 self.mark_preexisting(&key)?;
298 }
299 self.set(&key, row.clone())?;
300
301 let posts = [row.clone()];
302 let pres = [pre.clone()];
303 RingBufferRowInterceptor::post_update(self, &ringbuffer, &ids, &posts, &pres)?;
304
305 self.track_flow_change(build_ringbuffer_update_change(&ringbuffer, id, &pre, &row));
306
307 Ok(row)
308 }
309
310 fn remove_from_ringbuffer(
311 &mut self,
312 ringbuffer: &RingBuffer,
313 partition: Option<Partition>,
314 id: RowNumber,
315 ) -> Result<EncodedRow> {
316 let key = ringbuffer_key(ringbuffer, partition, id);
317
318 let displayed = match self.get(&key)? {
319 Some(v) => v.row,
320 None => return Ok(EncodedRow(CowVec::new(vec![]))),
321 };
322 let committed = self.get_committed(&key)?.map(|v| v.row);
323
324 let ids = [id];
325 RingBufferRowInterceptor::pre_delete(self, ringbuffer, &ids)?;
326
327 let pre_for_cdc = committed.clone().unwrap_or_else(|| displayed.clone());
328
329 if committed.is_some() {
330 self.mark_preexisting(&key)?;
331 }
332 self.unset(&key, pre_for_cdc.clone())?;
333
334 let pre_rows = [pre_for_cdc.clone()];
335 RingBufferRowInterceptor::post_delete(self, ringbuffer, &ids, &pre_rows)?;
336
337 self.track_flow_change(build_ringbuffer_remove_change(ringbuffer, id, &pre_for_cdc));
338
339 Ok(displayed)
340 }
341}
342
343impl RingBufferOperations for Transaction<'_> {
344 fn insert_ringbuffer(&mut self, _ringbuffer: RingBuffer, _row: EncodedRow) -> Result<RowNumber> {
345 unimplemented!(
346 "Ring buffer insert must be called with explicit row_number through insert_ringbuffer_at"
347 )
348 }
349
350 fn insert_ringbuffer_at(
351 &mut self,
352 ringbuffer: &RingBuffer,
353 shape: &RowShape,
354 partition: Option<Partition>,
355 row_number: RowNumber,
356 row: EncodedRow,
357 ) -> Result<EncodedRow> {
358 match self {
359 Transaction::Command(txn) => {
360 txn.insert_ringbuffer_at(ringbuffer, shape, partition, row_number, row)
361 }
362 Transaction::Admin(txn) => {
363 txn.insert_ringbuffer_at(ringbuffer, shape, partition, row_number, row)
364 }
365 Transaction::Test(t) => {
366 t.inner.insert_ringbuffer_at(ringbuffer, shape, partition, row_number, row)
367 }
368 Transaction::Query(_) => panic!("Write operations not supported on Query transaction"),
369 Transaction::Replica(_) => panic!("Write operations not supported on Replica transaction"),
370 }
371 }
372
373 fn update_ringbuffer(
374 &mut self,
375 ringbuffer: RingBuffer,
376 partition: Option<Partition>,
377 id: RowNumber,
378 row: EncodedRow,
379 ) -> Result<EncodedRow> {
380 match self {
381 Transaction::Command(txn) => txn.update_ringbuffer(ringbuffer, partition, id, row),
382 Transaction::Admin(txn) => txn.update_ringbuffer(ringbuffer, partition, id, row),
383 Transaction::Test(t) => t.inner.update_ringbuffer(ringbuffer, partition, id, row),
384 Transaction::Query(_) => panic!("Write operations not supported on Query transaction"),
385 Transaction::Replica(_) => panic!("Write operations not supported on Replica transaction"),
386 }
387 }
388
389 fn remove_from_ringbuffer(
390 &mut self,
391 ringbuffer: &RingBuffer,
392 partition: Option<Partition>,
393 id: RowNumber,
394 ) -> Result<EncodedRow> {
395 match self {
396 Transaction::Command(txn) => txn.remove_from_ringbuffer(ringbuffer, partition, id),
397 Transaction::Admin(txn) => txn.remove_from_ringbuffer(ringbuffer, partition, id),
398 Transaction::Test(t) => t.inner.remove_from_ringbuffer(ringbuffer, partition, id),
399 Transaction::Query(_) => panic!("Write operations not supported on Query transaction"),
400 Transaction::Replica(_) => panic!("Write operations not supported on Replica transaction"),
401 }
402 }
403}