Skip to main content

reifydb_engine/vm/volcano/scan/
queue.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::{collections::Bound, sync::Arc};
5
6use reifydb_codec::{
7	key::encoded::{EncodedKey, EncodedKeyRange},
8	row::{bytes::EncodedBytes, queue::EncodedQueueRow, shape::RowShape},
9};
10use reifydb_core::{
11	interface::{catalog::dictionary::Dictionary, resolved::ResolvedQueue, store::MultiVersionRow},
12	internal_error,
13	key::{
14		EncodableKey,
15		row::{RowKey, RowKeyRange},
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::{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
32type DrainedBatch = (Vec<EncodedBytes>, Vec<RowNumber>, Option<EncodedKey>, bool);
33
34pub struct QueueScan {
35	queue: ResolvedQueue,
36	headers: ColumnHeaders,
37	shape: Option<RowShape>,
38	storage_types: Vec<ValueType>,
39	dictionaries: Vec<Option<Dictionary>>,
40	last_key: Option<EncodedKey>,
41	exhausted: bool,
42	context: Option<Arc<QueryContext>>,
43}
44
45impl QueueScan {
46	pub fn new(queue: ResolvedQueue, context: Arc<QueryContext>, rx: &mut Transaction<'_>) -> Result<Self> {
47		let mut storage_types = Vec::with_capacity(queue.columns().len());
48		let mut dictionaries = Vec::with_capacity(queue.columns().len());
49
50		for col in queue.columns() {
51			if let Some(dict_id) = col.dictionary_id
52				&& let Some(dict) = context.services.catalog.find_dictionary(rx, dict_id)?
53			{
54				storage_types.push(ValueType::DictionaryId);
55				dictionaries.push(Some(dict));
56				continue;
57			}
58			storage_types.push(col.constraint.get_type());
59			dictionaries.push(None);
60		}
61
62		let headers = ColumnHeaders {
63			columns: queue.columns().iter().map(|col| Fragment::internal(&col.name)).collect(),
64		};
65
66		Ok(Self {
67			queue,
68			headers,
69			shape: None,
70			storage_types,
71			dictionaries,
72			last_key: None,
73			exhausted: false,
74			context: Some(context),
75		})
76	}
77
78	fn get_or_load_shape(&mut self, rx: &mut Transaction, first: &EncodedBytes) -> Result<RowShape> {
79		if let Some(shape) = &self.shape {
80			return Ok(shape.clone());
81		}
82
83		let fingerprint = EncodedQueueRow::view(first).fingerprint();
84		let stored_ctx = self.context.as_ref().expect("QueueScan context not set");
85		let shape = stored_ctx.services.catalog.get_or_load_row_shape(fingerprint, rx)?.ok_or_else(|| {
86			internal_error!(
87				"RowShape with fingerprint {:?} not found for queue {}",
88				fingerprint,
89				self.queue.def().name
90			)
91		})?;
92
93		self.shape = Some(shape.clone());
94
95		Ok(shape)
96	}
97
98	#[instrument(level = "trace", skip_all, name = "volcano::scan::queue::range_open")]
99	fn open_range<'rx, 'tx>(
100		rx: &'rx mut Transaction<'tx>,
101		range: EncodedKeyRange,
102		batch_size: u64,
103	) -> Result<Box<dyn Iterator<Item = Result<MultiVersionRow>> + Send + 'rx>> {
104		rx.range_rev(range, RangeScope::All, batch_size as usize)
105	}
106
107	#[instrument(level = "trace", skip_all, name = "volcano::scan::queue::drain")]
108	fn drain_batch(
109		stream: &mut dyn Iterator<Item = Result<MultiVersionRow>>,
110		batch_size: u64,
111	) -> Result<DrainedBatch> {
112		let mut batch: Vec<EncodedBytes> = Vec::new();
113		let mut row_numbers: Vec<RowNumber> = Vec::new();
114		let mut new_last_key = None;
115		let mut drained = false;
116
117		for _ in 0..batch_size {
118			match stream.next() {
119				Some(Ok(multi)) => {
120					if let Some(key) = RowKey::decode(&multi.key) {
121						batch.push(multi.bytes);
122						row_numbers.push(key.row);
123						new_last_key = Some(multi.key);
124					}
125				}
126				Some(Err(e)) => return Err(e),
127				None => {
128					drained = true;
129					break;
130				}
131			}
132		}
133
134		Ok((batch, row_numbers, new_last_key, drained))
135	}
136
137	#[instrument(level = "trace", skip_all, name = "volcano::scan::queue::column_alloc")]
138	fn storage_columns(&self, shape: &RowShape, declared: usize) -> Result<Vec<ColumnWithName>> {
139		let mut storage_columns: Vec<ColumnWithName> = self
140			.queue
141			.columns()
142			.iter()
143			.enumerate()
144			.map(|(idx, col)| ColumnWithName {
145				name: Fragment::internal(&col.name),
146				data: ColumnBuffer::with_capacity(self.storage_types[idx].clone(), 0),
147			})
148			.collect();
149
150		for index in declared..shape.field_count() {
151			let field = shape.get_field(index).ok_or_else(|| {
152				internal_error!("queue {} shape lost field {}", self.queue.def().name, index)
153			})?;
154			storage_columns.push(ColumnWithName {
155				name: Fragment::internal(field.name.clone()),
156				data: ColumnBuffer::with_capacity(field.constraint.get_type(), 0),
157			});
158		}
159
160		Ok(storage_columns)
161	}
162
163	#[instrument(level = "trace", skip_all, name = "volcano::scan::queue::append_rows")]
164	fn append_batch(
165		shape: &RowShape,
166		columns: &mut Columns,
167		bytes_vec: Vec<EncodedBytes>,
168		row_numbers: Vec<RowNumber>,
169	) -> Result<()> {
170		columns.append_rows(shape, bytes_vec.into_iter(), row_numbers)?;
171		Ok(())
172	}
173
174	fn enqueue_order_range(&self) -> EncodedKeyRange {
175		let full = RowKeyRange::scan_range(self.queue.def().id.into(), None);
176		match &self.last_key {
177			Some(last_key) => EncodedKeyRange::new(full.start.clone(), Bound::Excluded(last_key.clone())),
178			None => full,
179		}
180	}
181
182	fn empty_declared_columns(&self) -> Columns {
183		Columns::new(
184			self.queue
185				.columns()
186				.iter()
187				.map(|col| ColumnWithName {
188					name: Fragment::internal(&col.name),
189					data: ColumnBuffer::none_typed(col.constraint.get_type(), 0),
190				})
191				.collect(),
192		)
193	}
194}
195
196impl QueryNode for QueueScan {
197	#[instrument(level = "trace", skip_all, name = "volcano::scan::queue::initialize")]
198	fn initialize<'a>(&mut self, _rx: &mut Transaction<'a>, _ctx: &QueryContext) -> Result<()> {
199		Ok(())
200	}
201
202	#[instrument(level = "trace", skip_all, name = "volcano::scan::queue::next")]
203	fn next<'a>(&mut self, rx: &mut Transaction<'a>, _ctx: &mut QueryContext) -> Result<Option<Columns>> {
204		if self.exhausted {
205			return Ok(None);
206		}
207
208		let batch_size = self.context.as_ref().expect("QueueScan context not set").batch_size;
209		let range = self.enqueue_order_range();
210
211		let (batch_rows, row_numbers, new_last_key, drained) = {
212			let mut stream = Self::open_range(rx, range, batch_size)?;
213			Self::drain_batch(&mut stream, batch_size)?
214		};
215
216		if drained {
217			self.exhausted = true;
218		}
219
220		if batch_rows.is_empty() {
221			self.exhausted = true;
222			if self.last_key.is_none() {
223				return Ok(Some(self.empty_declared_columns()));
224			}
225			return Ok(None);
226		}
227
228		self.last_key = new_last_key;
229
230		let shape = self.get_or_load_shape(rx, &batch_rows[0])?;
231		let declared = self.queue.columns().len();
232
233		let storage_columns = self.storage_columns(&shape, declared)?;
234
235		let mut columns = Columns::with_system(storage_columns, SystemColumns::default());
236		Self::append_batch(&shape, &mut columns, batch_rows, row_numbers)?;
237
238		decode_dictionary_columns(&mut columns, &self.dictionaries, rx)?;
239
240		columns.columns.truncate(declared);
241		columns.names.truncate(declared);
242
243		Ok(Some(columns))
244	}
245
246	fn headers(&self) -> Option<ColumnHeaders> {
247		Some(self.headers.clone())
248	}
249}