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