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