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