Skip to main content

reifydb_engine/vm/volcano/scan/
table.rs

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