1use 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}