reifydb-engine 0.9.1

Query execution and processing engine for ReifyDB
Documentation
// SPDX-License-Identifier: Apache-2.0
// Copyright (c) 2026 ReifyDB

use std::{ops::Bound, sync::Arc};

use reifydb_codec::row::{
	bytes::{EncodedBytes, read_fingerprint},
	shape::RowShape,
};
use reifydb_core::{
	interface::{
		catalog::{dictionary::Dictionary, view::ViewStorageKind},
		resolved::ResolvedView,
		store::MultiVersionRow,
	},
	internal_error,
	key::{
		any::TaggedKey,
		bound::TaggedKeyBoundRange,
		row::{PartitionedSortedViewRowKey, RowKeyRange, SortedViewRowKey, StoragePartitionedRowKey},
		series::{PartitionedSeriesRowKeyRange, SeriesRowKeyRange},
	},
	value::column::{ColumnWithName, buffer::ColumnBuffer, columns::Columns, headers::ColumnHeaders},
};
use reifydb_transaction::{multi::RangeScope, transaction::Transaction};
use reifydb_value::{
	fragment::Fragment,
	reifydb_assertions,
	value::{partition::Partition, row_number::RowNumber, system_columns::SystemColumns, value_type::ValueType},
};
use tracing::instrument;

use super::{super::decode_dictionary_columns, guard_view_read};
use crate::{
	Result,
	vm::volcano::query::{QueryContext, QueryNode},
};

type DrainedBatch = (Vec<EncodedBytes>, Vec<RowNumber>, Option<TaggedKey>, bool);

type DrainedPartitionedBatch = (Vec<EncodedBytes>, Vec<RowNumber>, Option<StoragePartitionedRowKey>, bool);

enum Resume {
	Key(Option<TaggedKey>),
	Partitioned(Option<StoragePartitionedRowKey>),
}

fn partitioned_bounds(
	partition: Option<Partition>,
	last: Option<StoragePartitionedRowKey>,
) -> (Bound<StoragePartitionedRowKey>, Bound<StoragePartitionedRowKey>) {
	let start = match (last, partition) {
		(Some(key), _) => Bound::Excluded(key),
		(None, Some(partition)) => {
			Bound::Included(StoragePartitionedRowKey::new(partition, RowNumber(u64::MAX)))
		}
		(None, None) => Bound::Unbounded,
	};
	let end = match partition {
		Some(partition) => Bound::Included(StoragePartitionedRowKey::new(partition, RowNumber(u64::MIN))),
		None => Bound::Unbounded,
	};
	(start, end)
}

pub(crate) struct ViewScanNode {
	view: ResolvedView,
	context: Option<Arc<QueryContext>>,
	headers: ColumnHeaders,
	storage_types: Vec<ValueType>,
	dictionaries: Vec<Option<Dictionary>>,
	shape: Option<RowShape>,
	resume: Resume,
	exhausted: bool,
	sorted: bool,
	partitioned: bool,
	series: bool,
	partition: Option<Partition>,
}

impl ViewScanNode {
	pub fn new(
		view: ResolvedView,
		partition: Option<Partition>,
		context: Arc<QueryContext>,
		rx: &mut Transaction<'_>,
	) -> Result<Self> {
		let mut storage_types = Vec::with_capacity(view.columns().len());
		let mut dictionaries = Vec::with_capacity(view.columns().len());

		for col in view.columns() {
			if let Some(dict_id) = col.dictionary_id {
				if let Some(dict) = context.services.catalog.find_dictionary(rx, dict_id)? {
					storage_types.push(ValueType::DictionaryId);
					dictionaries.push(Some(dict));
				} else {
					storage_types.push(col.constraint.get_type());
					dictionaries.push(None);
				}
			} else {
				storage_types.push(col.constraint.get_type());
				dictionaries.push(None);
			}
		}

		let headers = ColumnHeaders {
			columns: view.columns().iter().map(|col| Fragment::internal(&col.name)).collect(),
		};
		let series = view.def().storage_kind() == ViewStorageKind::Series;
		let sorted = !view.def().sort().is_empty() && view.def().storage_kind() == ViewStorageKind::Table;
		let partitioned = !view.def().partition_by().is_empty();

		let resume = if partitioned && !series && !sorted {
			Resume::Partitioned(None)
		} else {
			Resume::Key(None)
		};

		Ok(Self {
			view,
			context: Some(context),
			headers,
			storage_types,
			dictionaries,
			shape: None,
			resume,
			exhausted: false,
			sorted,
			partitioned,
			series,
			partition,
		})
	}

	fn get_or_load_shape<'a>(&mut self, rx: &mut Transaction<'a>, first: &EncodedBytes) -> Result<RowShape> {
		if let Some(shape) = &self.shape {
			return Ok(shape.clone());
		}

		let fingerprint = read_fingerprint(first);

		let stored_ctx = self.context.as_ref().expect("ViewScanNode context not set");
		let shape = stored_ctx.services.catalog.get_or_load_row_shape(fingerprint, rx)?.ok_or_else(|| {
			internal_error!(
				"RowShape with fingerprint {:?} not found for view {}",
				fingerprint,
				self.view.def().name()
			)
		})?;

		self.shape = Some(shape.clone());

		Ok(shape)
	}

	#[instrument(level = "trace", skip_all, name = "volcano::scan::view::range_open")]
	fn open_range<'rx, 'tx>(
		rx: &'rx mut Transaction<'tx>,
		range: TaggedKeyBoundRange,
		batch_size: u64,
	) -> Result<Box<dyn Iterator<Item = Result<MultiVersionRow<TaggedKey>>> + Send + 'rx>> {
		rx.range(range, RangeScope::All, batch_size as usize)
	}

	#[instrument(level = "trace", skip_all, name = "volcano::scan::view::drain")]
	fn drain_batch(
		&self,
		stream: &mut dyn Iterator<Item = Result<MultiVersionRow<TaggedKey>>>,
		batch_size: u64,
	) -> Result<DrainedBatch> {
		let mut batch = Vec::new();
		let mut row_numbers = Vec::new();
		let mut new_last_key = None;
		let mut drained = false;

		for _ in 0..batch_size {
			match stream.next() {
				Some(Ok(multi)) => {
					let row = if self.series {
						if self.partitioned {
							match &multi.key {
								TaggedKey::PartitionedSeriesRow(key) => {
									RowNumber(key.sequence)
								}
								_ => continue,
							}
						} else {
							match &multi.key {
								TaggedKey::SeriesRow(key) => RowNumber(key.sequence),
								_ => continue,
							}
						}
					} else if self.sorted {
						let row = if self.partitioned {
							match &multi.key {
								TaggedKey::PartitionedSortedViewRow(key) => {
									Some(key.row.0)
								}
								_ => None,
							}
						} else {
							match &multi.key {
								TaggedKey::SortedViewRow(key) => Some(key.row.0),
								_ => None,
							}
						};
						match row {
							Some(row) => row,
							None => continue,
						}
					} else if let TaggedKey::Row(key) = &multi.key {
						key.row
					} else {
						continue;
					};
					batch.push(multi.bytes);
					row_numbers.push(row);
					new_last_key = Some(multi.key);
				}
				Some(Err(e)) => return Err(e),
				None => {
					drained = true;
					break;
				}
			}
		}

		Ok((batch, row_numbers, new_last_key, drained))
	}

	#[instrument(level = "trace", skip_all, name = "volcano::scan::view::drain_partitioned")]
	fn drain_batch_partitioned(
		stream: &mut dyn Iterator<Item = Result<MultiVersionRow<StoragePartitionedRowKey>>>,
		batch_size: u64,
	) -> Result<DrainedPartitionedBatch> {
		let mut batch = Vec::new();
		let mut row_numbers = Vec::new();
		let mut new_last_key = None;
		let mut drained = false;

		for _ in 0..batch_size {
			match stream.next() {
				Some(Ok(multi)) => {
					batch.push(multi.bytes);
					row_numbers.push(multi.key.row());
					new_last_key = Some(multi.key);
				}
				Some(Err(e)) => return Err(e),
				None => {
					drained = true;
					break;
				}
			}
		}

		Ok((batch, row_numbers, new_last_key, drained))
	}

	#[instrument(level = "trace", skip_all, name = "volcano::scan::view::column_alloc")]
	fn storage_columns(&self) -> Vec<ColumnWithName> {
		self.view
			.columns()
			.iter()
			.enumerate()
			.map(|(idx, col)| ColumnWithName {
				name: Fragment::internal(&col.name),
				data: ColumnBuffer::with_capacity(self.storage_types[idx].clone(), 0),
			})
			.collect()
	}

	#[instrument(level = "trace", skip_all, name = "volcano::scan::view::append_rows")]
	fn append_batch<'a>(
		&mut self,
		rx: &mut Transaction<'a>,
		columns: &mut Columns,
		bytes_vec: Vec<EncodedBytes>,
		row_numbers: Vec<RowNumber>,
	) -> Result<()> {
		let shape = self.get_or_load_shape(rx, &bytes_vec[0])?;
		columns.append_rows(&shape, bytes_vec.into_iter(), row_numbers)?;
		Ok(())
	}
}

impl QueryNode for ViewScanNode {
	#[instrument(name = "volcano::scan::view::initialize", level = "trace", skip_all)]
	fn initialize<'a>(&mut self, rx: &mut Transaction<'a>, ctx: &QueryContext) -> Result<()> {
		guard_view_read(&self.view, rx, &ctx.services)
	}

	#[instrument(name = "volcano::scan::view::next", level = "trace", skip_all)]
	fn next<'a>(&mut self, rx: &mut Transaction<'a>, _ctx: &mut QueryContext) -> Result<Option<Columns>> {
		reifydb_assertions! {
			assert!(self.context.is_some(), "ViewScanNode::next() called before initialize()");
		}
		let stored_ctx = self.context.as_ref().unwrap();

		if self.exhausted {
			return Ok(None);
		}

		let batch_size = stored_ctx.batch_size;
		let storage = self.view.def().storage_id();

		let (batch, row_numbers, next_resume, resumed, drained) = match &self.resume {
			Resume::Partitioned(last) => {
				let last = *last;
				let (start, end) = partitioned_bounds(self.partition, last);
				let (batch, row_numbers, new_last_key, drained) = {
					let mut stream = rx.range_partitioned_row(
						storage,
						start,
						end,
						RangeScope::All,
						batch_size as usize,
					)?;
					Self::drain_batch_partitioned(&mut stream, batch_size)?
				};
				(batch, row_numbers, Resume::Partitioned(new_last_key), last.is_some(), drained)
			}
			Resume::Key(last) => {
				let range = match (self.series, self.partitioned, self.partition) {
					(true, true, Some(partition)) => {
						PartitionedSeriesRowKeyRange::partition_scan_range(
							storage,
							partition,
							last.as_ref(),
						)
					}
					(true, true, None) => {
						PartitionedSeriesRowKeyRange::full_scan_range(storage, last.as_ref())
					}
					(true, false, _) => {
						SeriesRowKeyRange::scan_range(storage, None, None, None, last.as_ref())
					}
					(false, true, Some(partition)) if self.sorted => {
						PartitionedSortedViewRowKey::partition_scan_range(
							storage,
							partition,
							last.as_ref(),
						)
					}
					(false, true, None) if self.sorted => {
						PartitionedSortedViewRowKey::scan_range(storage, last.as_ref())
					}
					(false, false, _) if self.sorted => {
						SortedViewRowKey::scan_range(storage, last.as_ref())
					}
					(false, false, _) => RowKeyRange::scan_range(storage, last.as_ref()),
					(false, true, _) => unreachable!(
						"unsorted partitioned view rows resume through Resume::Partitioned"
					),
				};

				let resumed = last.is_some();
				let (batch, row_numbers, new_last_key, drained) = {
					let mut stream = Self::open_range(rx, range, batch_size)?;
					self.drain_batch(&mut stream, batch_size)?
				};
				(batch, row_numbers, Resume::Key(new_last_key), resumed, drained)
			}
		};

		if drained {
			self.exhausted = true;
		}

		if batch.is_empty() {
			self.exhausted = true;
			if !resumed {
				return Ok(Some(Columns::from_catalog_columns(self.view.columns())));
			}
			return Ok(None);
		}

		self.resume = next_resume;

		let mut columns = Columns::with_system(self.storage_columns(), SystemColumns::default());
		self.append_batch(rx, &mut columns, batch, row_numbers)?;

		decode_dictionary_columns(&mut columns, &self.dictionaries, rx)?;

		Ok(Some(columns))
	}

	fn headers(&self) -> Option<ColumnHeaders> {
		Some(self.headers.clone())
	}
}