1use std::collections::{BTreeMap, HashSet};
5
6use postcard::{from_bytes, to_stdvec};
7use reifydb_abi::operator::capabilities::OperatorCapability;
8use reifydb_catalog::store::ringbuffer::update::{decode_ringbuffer_metadata, encode_ringbuffer_metadata};
9use reifydb_codec::{
10 encoded::{row::EncodedRow, shape::RowShape},
11 key::encoded::EncodedKey,
12};
13use reifydb_core::{
14 interface::{
15 catalog::{
16 flow::FlowNodeId, id::RingBufferId, ringbuffer::RingBufferMetadata, shape::ShapeId, view::View,
17 },
18 change::{Change, ChangeOrigin, Diff},
19 resolved::ResolvedView,
20 },
21 key::{ringbuffer::RingBufferMetadataKey, row::RowKey},
22 row::row_shape_from_columns,
23 value::column::columns::Columns,
24};
25use reifydb_value::{
26 Result,
27 error::Error,
28 value::{blob::Blob, datetime::DateTime, row_number::RowNumber},
29};
30use serde::{Deserialize, Serialize};
31use smallvec::smallvec;
32
33use super::{coerce_columns, encode_row_at_index, shape_field_columns, view::dictionary_encode_view_columns};
34use crate::{
35 Operator,
36 error::FlowStateError,
37 operator::{
38 OperatorCell,
39 stateful::{raw::RawStatefulOperator, single::SingleStateful},
40 },
41 transaction::FlowTransaction,
42};
43
44#[derive(Debug, Clone, Serialize, Deserialize, Default)]
45struct RingBufferState {
46 forward: BTreeMap<RowNumber, RowNumber>,
47 reverse: BTreeMap<RowNumber, RowNumber>,
48}
49
50pub struct SinkRingBufferViewOperator {
51 #[allow(dead_code)]
52 parent: OperatorCell,
53 node: FlowNodeId,
54 view: ResolvedView,
55 ringbuffer_id: RingBufferId,
56 capacity: u64,
57 propagate_evictions: bool,
58 state_shape: RowShape,
59}
60
61impl SinkRingBufferViewOperator {
62 pub fn new(
63 parent: OperatorCell,
64 node: FlowNodeId,
65 view: ResolvedView,
66 ringbuffer_id: RingBufferId,
67 capacity: u64,
68 propagate_evictions: bool,
69 ) -> Self {
70 Self {
71 parent,
72 node,
73 view,
74 ringbuffer_id,
75 capacity,
76 propagate_evictions,
77 state_shape: RowShape::operator_state(),
78 }
79 }
80
81 fn read_metadata(&self, txn: &mut FlowTransaction) -> Result<RingBufferMetadata> {
82 let key = RingBufferMetadataKey::encoded(self.ringbuffer_id);
83 match txn.get(&key)? {
84 Some(row) => Ok(decode_ringbuffer_metadata(&row)),
85 None => Ok(RingBufferMetadata::new(self.ringbuffer_id, self.capacity)),
86 }
87 }
88
89 fn write_metadata(&self, txn: &mut FlowTransaction, metadata: &RingBufferMetadata) -> Result<()> {
90 let key = RingBufferMetadataKey::encoded(self.ringbuffer_id);
91 let row = encode_ringbuffer_metadata(metadata);
92 txn.set(&key, row)
93 }
94
95 fn load(&self, txn: &mut FlowTransaction) -> Result<RingBufferState> {
96 let state_row = self.load_state(txn)?;
97
98 if state_row.is_empty() || !state_row.is_defined(0) {
99 return Ok(RingBufferState::default());
100 }
101
102 let blob = self.state_shape.get_blob(&state_row, 0);
103 if blob.is_empty() {
104 return Ok(RingBufferState::default());
105 }
106
107 from_bytes(blob.as_ref()).map_err(|e| {
108 Error::from(FlowStateError::Decode {
109 state: "RingBufferState",
110 cause: e.to_string(),
111 })
112 })
113 }
114
115 fn save(&self, txn: &mut FlowTransaction, state: &RingBufferState) -> Result<()> {
116 let serialized = to_stdvec(state).map_err(|e| {
117 Error::from(FlowStateError::Encode {
118 state: "RingBufferState",
119 cause: e.to_string(),
120 })
121 })?;
122 let blob = Blob::from(serialized);
123
124 self.update_state(txn, |shape, row| {
125 shape.set_blob(row, 0, &blob);
126 Ok(())
127 })?;
128 Ok(())
129 }
130}
131
132impl RawStatefulOperator for SinkRingBufferViewOperator {}
133
134impl SingleStateful for SinkRingBufferViewOperator {
135 fn layout(&self) -> RowShape {
136 self.state_shape.clone()
137 }
138}
139
140impl Operator for SinkRingBufferViewOperator {
141 fn id(&self) -> FlowNodeId {
142 self.node
143 }
144
145 fn capabilities(&self) -> &[OperatorCapability] {
146 OperatorCapability::STANDARD
147 }
148
149 fn apply(&self, txn: &mut FlowTransaction, change: Change) -> Result<Change> {
150 let view = self.view.def().clone();
151 let shape = row_shape_from_columns(view.columns());
152 let object_id = ShapeId::ringbuffer(self.ringbuffer_id);
153 let mut metadata = self.read_metadata(txn)?;
154 let mut state = self.load(txn)?;
155
156 for diff in change.diffs.iter() {
157 match diff {
158 Diff::Insert {
159 post,
160 ..
161 } => self.apply_ringbuffer_insert(
162 txn,
163 &view,
164 &shape,
165 object_id,
166 &mut metadata,
167 &mut state,
168 post,
169 )?,
170 Diff::Update {
171 pre,
172 post,
173 ..
174 } => self.apply_ringbuffer_update(txn, &view, &shape, object_id, &state, pre, post)?,
175 Diff::Remove {
176 pre,
177 ..
178 } => self.apply_ringbuffer_remove(txn, &view, object_id, &mut state, pre)?,
179 }
180 }
181
182 self.write_metadata(txn, &metadata)?;
183 self.save(txn, &state)?;
184
185 Ok(Change::from_flow(self.node, change.version, Vec::new(), change.changed_at))
186 }
187}
188
189impl SinkRingBufferViewOperator {
190 #[inline]
191 #[allow(clippy::too_many_arguments)]
192 fn apply_ringbuffer_insert(
193 &self,
194 txn: &mut FlowTransaction,
195 view: &View,
196 shape: &RowShape,
197 object_id: ShapeId,
198 metadata: &mut RingBufferMetadata,
199 state: &mut RingBufferState,
200 post: &Columns,
201 ) -> Result<()> {
202 let coerced = coerce_columns(post, view.columns())?;
203 let dict_encoded = dictionary_encode_view_columns(txn, view, &coerced)?;
204 let source = dict_encoded.as_ref().unwrap_or(&coerced);
205 let row_count = source.row_count();
206 let field_columns = shape_field_columns(source, shape);
207 let mut assigned_ids: Vec<RowNumber> = Vec::with_capacity(row_count);
208 let mut encoded_rows: Vec<EncodedRow> = Vec::with_capacity(row_count);
209 let mut evicted_in_batch: HashSet<RowNumber> = HashSet::new();
210 for row_idx in 0..row_count {
211 if metadata.is_full() {
212 let oldest_rn = RowNumber(metadata.head);
213 let pre_key = RowKey::encoded(object_id, oldest_rn);
214 txn.remove(&pre_key)?;
215 metadata.head += 1;
216 metadata.count -= 1;
217 evicted_in_batch.insert(oldest_rn);
218
219 if let Some(source_rn) = state.reverse.remove(&oldest_rn) {
220 state.forward.remove(&source_rn);
221 }
222
223 if self.propagate_evictions {}
224 }
225
226 let source_rn = source.row_numbers[row_idx];
227 let assigned_rn = RowNumber(metadata.tail);
228 let (_, encoded) = encode_row_at_index(source, row_idx, shape, assigned_rn, &field_columns)?;
229
230 if source_rn != assigned_rn {
231 state.forward.insert(source_rn, assigned_rn);
232 state.reverse.insert(assigned_rn, source_rn);
233 }
234
235 assigned_ids.push(assigned_rn);
236 encoded_rows.push(encoded);
237
238 if metadata.is_empty() {
239 metadata.head = assigned_rn.0;
240 }
241 metadata.count += 1;
242 metadata.tail = assigned_rn.0 + 1;
243 }
244
245 let surviving: Vec<usize> =
246 (0..assigned_ids.len()).filter(|&i| !evicted_in_batch.contains(&assigned_ids[i])).collect();
247 let final_ids: Vec<RowNumber> = surviving.iter().map(|&i| assigned_ids[i]).collect();
248 let final_rows: Vec<EncodedRow> = surviving.iter().map(|&i| encoded_rows[i].clone()).collect();
249
250 for (assigned_rn, encoded) in final_ids.iter().zip(final_rows.iter()) {
251 let key = RowKey::encoded(object_id, *assigned_rn);
252 txn.set(&key, encoded.clone())?;
253 }
254 emit_view_change(txn, view, Diff::insert(coerced));
255 Ok(())
256 }
257
258 #[inline]
259 #[allow(clippy::too_many_arguments)]
260 fn apply_ringbuffer_update(
261 &self,
262 txn: &mut FlowTransaction,
263 view: &View,
264 shape: &RowShape,
265 object_id: ShapeId,
266 state: &RingBufferState,
267 pre: &Columns,
268 post: &Columns,
269 ) -> Result<()> {
270 let coerced_pre = coerce_columns(pre, view.columns())?;
271 let coerced_post = coerce_columns(post, view.columns())?;
272 let dict_pre = dictionary_encode_view_columns(txn, view, &coerced_pre)?;
273 let dict_post = dictionary_encode_view_columns(txn, view, &coerced_post)?;
274 let source_pre = dict_pre.as_ref().unwrap_or(&coerced_pre);
275 let source_post = dict_post.as_ref().unwrap_or(&coerced_post);
276 let row_count = source_post.row_count();
277 let field_columns = shape_field_columns(source_post, shape);
278 let mut pre_keys: Vec<EncodedKey> = Vec::with_capacity(row_count);
279 let mut post_keys: Vec<EncodedKey> = Vec::with_capacity(row_count);
280 let mut post_encoded_rows: Vec<EncodedRow> = Vec::with_capacity(row_count);
281 for row_idx in 0..row_count {
282 let pre_source_rn = source_pre.row_numbers[row_idx];
283 let post_source_rn = source_post.row_numbers[row_idx];
284 let pre_storage_rn = state.forward.get(&pre_source_rn).copied().unwrap_or(pre_source_rn);
285 let post_storage_rn = state.forward.get(&post_source_rn).copied().unwrap_or(post_source_rn);
286 let (_, post_encoded) =
287 encode_row_at_index(source_post, row_idx, shape, post_storage_rn, &field_columns)?;
288
289 pre_keys.push(RowKey::encoded(object_id, pre_storage_rn));
290 post_keys.push(RowKey::encoded(object_id, post_storage_rn));
291 post_encoded_rows.push(post_encoded);
292 }
293
294 for ((pre_key, post_key), post_encoded) in
295 pre_keys.iter().zip(post_keys.iter()).zip(post_encoded_rows.iter())
296 {
297 txn.remove(pre_key)?;
298 txn.set(post_key, post_encoded.clone())?;
299 }
300 emit_view_change(txn, view, Diff::update(coerced_pre, coerced_post));
301 Ok(())
302 }
303
304 #[inline]
305 fn apply_ringbuffer_remove(
306 &self,
307 txn: &mut FlowTransaction,
308 view: &View,
309 object_id: ShapeId,
310 state: &mut RingBufferState,
311 pre: &Columns,
312 ) -> Result<()> {
313 let coerced = coerce_columns(pre, view.columns())?;
314 let row_count = coerced.row_count();
315 let mut storage_ids: Vec<RowNumber> = Vec::with_capacity(row_count);
316 for row_idx in 0..row_count {
317 let source_rn = coerced.row_numbers[row_idx];
318 let storage_rn = state.forward.remove(&source_rn).unwrap_or(source_rn);
319 state.reverse.remove(&storage_rn);
320 storage_ids.push(storage_rn);
321 }
322 for storage_rn in storage_ids.iter() {
323 let key = RowKey::encoded(object_id, *storage_rn);
324 txn.remove(&key)?;
325 }
326 emit_view_change(txn, view, Diff::remove(coerced));
327 Ok(())
328 }
329}
330
331#[inline]
332fn emit_view_change(txn: &mut FlowTransaction, view: &View, diff: Diff) {
333 let version = txn.version();
334 let changed_at = DateTime::from_nanos(txn.clock().now_nanos());
335 txn.track_flow_change(Change {
336 origin: ChangeOrigin::Shape(ShapeId::view(view.id())),
337 version,
338 diffs: smallvec![diff],
339 changed_at,
340 });
341}