Skip to main content

reifydb_sub_flow/operator/sink/
ringbuffer_view.rs

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