Skip to main content

reifydb_engine/vm/volcano/scan/
ringbuffer.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::sync::Arc;
5
6use reifydb_codec::row::{bytes::EncodedBytes, ringbuffer::EncodedRingBufferRow, shape::RowShape};
7use reifydb_core::{
8	interface::{
9		catalog::{dictionary::Dictionary, object::ObjectId, ringbuffer::PartitionedMetadata},
10		resolved::ResolvedRingBuffer,
11	},
12	internal_error,
13	key::{
14		EncodableKey,
15		partitioned_row::{PartitionedRowKey, RowLocator},
16		row::RowKey,
17	},
18	value::column::{ColumnWithName, buffer::ColumnBuffer, columns::Columns, headers::ColumnHeaders},
19};
20use reifydb_transaction::{multi::RangeScope, transaction::Transaction};
21use reifydb_value::{
22	fragment::Fragment,
23	value::{partition::Partition, row_number::RowNumber, system_columns::SystemColumns, value_type::ValueType},
24};
25use tracing::instrument;
26
27use super::super::decode_dictionary_columns;
28use crate::{
29	Result,
30	vm::volcano::query::{QueryContext, QueryNode},
31};
32
33pub struct RingBufferScan {
34	ringbuffer: ResolvedRingBuffer,
35
36	partitions: Vec<PartitionedMetadata>,
37	current_partition_index: usize,
38	headers: ColumnHeaders,
39	shape: Option<RowShape>,
40
41	storage_types: Vec<ValueType>,
42
43	dictionaries: Vec<Option<Dictionary>>,
44
45	partition_col_indices: Vec<usize>,
46	current_partition_rows: Vec<(RowNumber, EncodedBytes)>,
47	current_partition_cursor: usize,
48	current_partition_loaded: bool,
49	finished: bool,
50	context: Option<Arc<QueryContext>>,
51	initialized: bool,
52}
53
54impl RingBufferScan {
55	pub fn new(
56		ringbuffer: ResolvedRingBuffer,
57		context: Arc<QueryContext>,
58		rx: &mut Transaction<'_>,
59	) -> Result<Self> {
60		let mut storage_types = Vec::with_capacity(ringbuffer.columns().len());
61		let mut dictionaries = Vec::with_capacity(ringbuffer.columns().len());
62
63		for col in ringbuffer.columns() {
64			if let Some(dict_id) = col.dictionary_id {
65				if let Some(dict) = context.services.catalog.find_dictionary(rx, dict_id)? {
66					storage_types.push(ValueType::DictionaryId);
67					dictionaries.push(Some(dict));
68				} else {
69					storage_types.push(col.constraint.get_type());
70					dictionaries.push(None);
71				}
72			} else {
73				storage_types.push(col.constraint.get_type());
74				dictionaries.push(None);
75			}
76		}
77
78		let partition_col_indices: Vec<usize> = ringbuffer
79			.def()
80			.partition_by
81			.iter()
82			.map(|pb_col| ringbuffer.columns().iter().position(|c| c.name == *pb_col).unwrap())
83			.collect();
84
85		let headers = ColumnHeaders {
86			columns: ringbuffer.columns().iter().map(|col| Fragment::internal(&col.name)).collect(),
87		};
88
89		Ok(Self {
90			ringbuffer,
91			partitions: Vec::new(),
92			current_partition_index: 0,
93			headers,
94			shape: None,
95			storage_types,
96			dictionaries,
97			partition_col_indices,
98			current_partition_rows: Vec::new(),
99			current_partition_cursor: 0,
100			current_partition_loaded: false,
101			finished: false,
102			context: Some(context),
103			initialized: false,
104		})
105	}
106
107	fn get_or_load_shape(&mut self, rx: &mut Transaction, first: &EncodedBytes) -> Result<RowShape> {
108		if let Some(shape) = &self.shape {
109			return Ok(shape.clone());
110		}
111
112		let fingerprint = EncodedRingBufferRow::view(first).fingerprint();
113
114		let stored_ctx = self.context.as_ref().expect("RingBufferScan context not set");
115		let shape = stored_ctx.services.catalog.get_or_load_row_shape(fingerprint, rx)?.ok_or_else(|| {
116			internal_error!(
117				"RowShape with fingerprint {:?} not found for ringbuffer {}",
118				fingerprint,
119				self.ringbuffer.def().name
120			)
121		})?;
122
123		self.shape = Some(shape.clone());
124
125		Ok(shape)
126	}
127
128	#[instrument(level = "trace", skip_all, name = "volcano::scan::ringbuffer::drain")]
129	fn drain_batch(
130		&mut self,
131		txn: &mut Transaction<'_>,
132		batch_size: usize,
133		partitioned: bool,
134	) -> Result<(Vec<EncodedBytes>, Vec<RowNumber>, Vec<Partition>)> {
135		let mut batch: Vec<EncodedBytes> = Vec::new();
136		let mut row_numbers: Vec<RowNumber> = Vec::new();
137		let mut partitions_sidecar: Vec<Partition> = Vec::new();
138
139		while batch.len() < batch_size && self.current_partition_index < self.partitions.len() {
140			if !self.current_partition_loaded {
141				self.current_partition_rows =
142					self.load_partition_rows(txn, self.current_partition_index)?;
143				self.current_partition_cursor = 0;
144				self.current_partition_loaded = true;
145			}
146
147			let hash = if partitioned {
148				Some(Partition::of(&self.partitions[self.current_partition_index].partition_values))
149			} else {
150				None
151			};
152
153			while batch.len() < batch_size
154				&& self.current_partition_cursor < self.current_partition_rows.len()
155			{
156				let (rn, row) = self.current_partition_rows[self.current_partition_cursor].clone();
157				batch.push(row);
158				row_numbers.push(rn);
159				if let Some(h) = hash {
160					partitions_sidecar.push(h);
161				}
162				self.current_partition_cursor += 1;
163			}
164
165			if self.current_partition_cursor >= self.current_partition_rows.len() {
166				self.current_partition_index += 1;
167				self.current_partition_loaded = false;
168			}
169		}
170
171		Ok((batch, row_numbers, partitions_sidecar))
172	}
173
174	#[instrument(level = "trace", skip_all, name = "volcano::scan::ringbuffer::column_alloc")]
175	fn storage_columns(&self) -> Vec<ColumnWithName> {
176		self.ringbuffer
177			.columns()
178			.iter()
179			.enumerate()
180			.map(|(idx, col)| ColumnWithName {
181				name: Fragment::internal(&col.name),
182				data: ColumnBuffer::with_capacity(self.storage_types[idx].clone(), 0),
183			})
184			.collect()
185	}
186
187	#[instrument(level = "trace", skip_all, name = "volcano::scan::ringbuffer::empty_columns")]
188	fn empty_columns(&self) -> Vec<ColumnWithName> {
189		self.ringbuffer
190			.columns()
191			.iter()
192			.map(|col| ColumnWithName {
193				name: Fragment::internal(&col.name),
194				data: ColumnBuffer::none_typed(col.constraint.get_type(), 0),
195			})
196			.collect()
197	}
198
199	#[instrument(level = "trace", skip_all, name = "volcano::scan::ringbuffer::append_rows")]
200	fn append_batch(
201		&mut self,
202		txn: &mut Transaction<'_>,
203		columns: &mut Columns,
204		bytes_vec: Vec<EncodedBytes>,
205		row_numbers: Vec<RowNumber>,
206	) -> Result<()> {
207		let shape = self.get_or_load_shape(txn, &bytes_vec[0])?;
208		columns.append_rows(&shape, bytes_vec.into_iter(), row_numbers.clone())?;
209		Ok(())
210	}
211
212	#[instrument(level = "trace", skip_all, name = "volcano::scan::ringbuffer::load_partition")]
213	fn load_partition_rows(
214		&self,
215		txn: &mut Transaction<'_>,
216		partition_index: usize,
217	) -> Result<Vec<(RowNumber, EncodedBytes)>> {
218		let pm = &self.partitions[partition_index];
219		let rb_id = self.ringbuffer.def().id;
220
221		if self.partition_col_indices.is_empty() {
222			let mut out = Vec::new();
223			for rn_value in pm.metadata.head..pm.metadata.tail {
224				let rn = RowNumber(rn_value);
225				if let Some(multi) = txn.get(&RowKey::encoded(rb_id, rn))? {
226					out.push((rn, multi.bytes));
227				}
228			}
229			return Ok(out);
230		}
231
232		let hash = Partition::of(&pm.partition_values);
233		let mut out = Vec::new();
234		let mut last_key = None;
235		loop {
236			let batch: Vec<_> = txn
237				.range(
238					PartitionedRowKey::partition_scan_range(
239						ObjectId::ringbuffer(rb_id),
240						hash,
241						last_key.as_ref(),
242					),
243					RangeScope::All,
244					1024,
245				)?
246				.collect::<Result<Vec<_>>>()?;
247			if batch.is_empty() {
248				break;
249			}
250			let n = batch.len();
251			for entry in batch {
252				if let Some(RowLocator::Row(rn)) =
253					PartitionedRowKey::decode(&entry.key).map(|pk| pk.locator)
254				{
255					out.push((rn, entry.bytes));
256				}
257				last_key = Some(entry.key);
258			}
259			if n < 1024 {
260				break;
261			}
262		}
263		out.sort_by_key(|(rn, _)| rn.0);
264		Ok(out)
265	}
266}
267
268impl QueryNode for RingBufferScan {
269	#[instrument(name = "volcano::scan::ringbuffer::initialize", level = "trace", skip_all)]
270	fn initialize<'a>(&mut self, txn: &mut Transaction<'a>, ctx: &QueryContext) -> Result<()> {
271		if !self.initialized {
272			self.partitions =
273				ctx.services.catalog.list_ringbuffer_partitions(txn, self.ringbuffer.def())?;
274			self.initialized = true;
275		}
276		Ok(())
277	}
278
279	#[instrument(name = "volcano::scan::ringbuffer::next", level = "trace", skip_all)]
280	fn next<'a>(&mut self, txn: &mut Transaction<'a>, _ctx: &mut QueryContext) -> Result<Option<Columns>> {
281		if self.finished {
282			return Ok(None);
283		}
284
285		let batch_size = self.context.as_ref().expect("RingBufferScan context not set").batch_size as usize;
286		let partitioned = !self.partition_col_indices.is_empty();
287
288		let (batch_rows, row_numbers, partitions_sidecar) = self.drain_batch(txn, batch_size, partitioned)?;
289
290		if !batch_rows.is_empty() {
291			let mut columns = Columns::with_system(self.storage_columns(), SystemColumns::default());
292			self.append_batch(txn, &mut columns, batch_rows, row_numbers)?;
293			if partitioned {
294				columns.system.set_partitions(partitions_sidecar);
295			}
296
297			decode_dictionary_columns(&mut columns, &self.dictionaries, txn)?;
298
299			return Ok(Some(columns));
300		}
301
302		self.finished = true;
303		if self.partitions.is_empty() || self.partitions.iter().all(|p| p.metadata.is_empty()) {
304			return Ok(Some(Columns::new(self.empty_columns())));
305		}
306		Ok(None)
307	}
308
309	fn headers(&self) -> Option<ColumnHeaders> {
310		Some(self.headers.clone())
311	}
312}