Skip to main content

reifydb_engine/transaction/operation/
ringbuffer.rs

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