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