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_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}