Skip to main content

reifydb_engine/vm/volcano/scan/
series.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::sync::Arc;
5
6use reifydb_codec::key::encoded::EncodedKey;
7use reifydb_core::{
8	common::CommitVersion,
9	interface::resolved::ResolvedSeries,
10	key::{
11		EncodableKey,
12		series_row::{SeriesRowKey, SeriesRowKeyRange},
13	},
14	value::column::{ColumnWithName, buffer::ColumnBuffer, columns::Columns, headers::ColumnHeaders},
15};
16use reifydb_transaction::{multi::RangeScope, transaction::Transaction};
17use reifydb_value::{
18	fragment::Fragment,
19	reifydb_assertions,
20	value::{
21		Value, datetime::DateTime, dictionary::DictionaryEntryId, row_number::RowNumber, value_type::ValueType,
22	},
23};
24use tracing::instrument;
25
26use crate::{
27	Result,
28	transaction::operation::dictionary::DictionaryOperations,
29	vm::{
30		instruction::dml::shape::get_or_create_series_shape,
31		volcano::query::{QueryContext, QueryNode},
32	},
33};
34
35pub struct SeriesScanNode {
36	series: ResolvedSeries,
37	key_range_start: Option<u64>,
38	key_range_end: Option<u64>,
39	variant_tag: Option<u8>,
40	context: Option<Arc<QueryContext>>,
41	headers: ColumnHeaders,
42	last_key: Option<EncodedKey>,
43	exhausted: bool,
44
45	min_commit_version: Option<CommitVersion>,
46}
47
48impl SeriesScanNode {
49	pub fn with_min_commit_version(mut self, min_commit_version: Option<CommitVersion>) -> Self {
50		self.min_commit_version = min_commit_version;
51		self
52	}
53
54	pub fn new(
55		series: ResolvedSeries,
56		key_range_start: Option<u64>,
57		key_range_end: Option<u64>,
58		variant_tag: Option<u8>,
59		context: Arc<QueryContext>,
60	) -> Result<Self> {
61		let mut columns = vec![Fragment::internal(series.def().key.column())];
62		if series.def().tag.is_some() {
63			columns.push(Fragment::internal("tag"));
64		}
65		for col in series.columns() {
66			columns.push(Fragment::internal(&col.name));
67		}
68		let headers = ColumnHeaders {
69			columns,
70		};
71
72		Ok(Self {
73			series,
74			key_range_start,
75			key_range_end,
76			variant_tag,
77			context: Some(context),
78			headers,
79			last_key: None,
80			exhausted: false,
81			min_commit_version: None,
82		})
83	}
84}
85
86impl QueryNode for SeriesScanNode {
87	#[instrument(name = "volcano::scan::series::initialize", level = "trace", skip_all)]
88	fn initialize<'a>(&mut self, _rx: &mut Transaction<'a>, _ctx: &QueryContext) -> Result<()> {
89		Ok(())
90	}
91
92	#[instrument(name = "volcano::scan::series::next", level = "trace", skip_all)]
93	fn next<'a>(&mut self, rx: &mut Transaction<'a>, _ctx: &mut QueryContext) -> Result<Option<Columns>> {
94		reifydb_assertions! {
95			assert!(self.context.is_some(), "SeriesScanNode::next() called before initialize()");
96		}
97		let stored_ctx = self.context.as_ref().unwrap();
98
99		if self.exhausted {
100			return Ok(None);
101		}
102
103		let batch_size = stored_ctx.batch_size;
104		let series = self.series.def();
105		let has_tag = series.tag.is_some();
106
107		let range = SeriesRowKeyRange::scan_range(
108			series.id,
109			self.variant_tag,
110			self.key_range_start,
111			self.key_range_end,
112			self.last_key.as_ref(),
113		);
114
115		let mut key_values: Vec<u64> = Vec::new();
116		let mut tags: Vec<u8> = Vec::new();
117		let mut sequences: Vec<u64> = Vec::new();
118		let mut created_at_values: Vec<DateTime> = Vec::new();
119		let mut updated_at_values: Vec<DateTime> = Vec::new();
120		let mut data_rows: Vec<Vec<Value>> = Vec::new();
121		let mut new_last_key = None;
122
123		let read_shape = get_or_create_series_shape(&stored_ctx.services.catalog, self.series.def(), rx)?;
124
125		let scope = match self.min_commit_version {
126			Some(v) => RangeScope::After(v),
127			None => RangeScope::All,
128		};
129
130		let mut stream = rx.range(range, scope, batch_size as usize)?;
131		let mut count = 0;
132
133		for entry in stream.by_ref() {
134			let entry = entry?;
135
136			if let Some(key) = SeriesRowKey::decode(&entry.key) {
137				key_values.push(key.key);
138				sequences.push(key.sequence);
139				created_at_values.push(DateTime::from_nanos(entry.row.created_at_nanos()));
140				updated_at_values.push(DateTime::from_nanos(entry.row.updated_at_nanos()));
141				if has_tag {
142					tags.push(key.variant_tag.unwrap_or(0));
143				}
144
145				let mut values = Vec::with_capacity(series.data_columns().count());
146				for (i, _) in series.data_columns().enumerate() {
147					values.push(read_shape.get_value(&entry.row, i + 1));
148				}
149				data_rows.push(values);
150
151				new_last_key = Some(entry.key);
152				count += 1;
153				if count >= batch_size as usize {
154					break;
155				}
156			}
157		}
158
159		drop(stream);
160
161		if key_values.is_empty() {
162			self.exhausted = true;
163			if self.last_key.is_none() {
164				let key_type = series
165					.columns
166					.iter()
167					.find(|c| c.name == series.key.column())
168					.map(|c| c.constraint.get_type())
169					.unwrap_or(ValueType::Int8);
170				let mut result_columns = Vec::new();
171				result_columns.push(ColumnWithName {
172					name: Fragment::internal(series.key.column()),
173					data: ColumnBuffer::none_typed(key_type, 0),
174				});
175				if has_tag {
176					result_columns.push(ColumnWithName {
177						name: Fragment::internal("tag"),
178						data: ColumnBuffer::none_typed(ValueType::Uint1, 0),
179					});
180				}
181				for col_def in series.data_columns() {
182					result_columns.push(ColumnWithName {
183						name: Fragment::internal(&col_def.name),
184						data: ColumnBuffer::none_typed(col_def.constraint.get_type(), 0),
185					});
186				}
187				return Ok(Some(Columns::new(result_columns)));
188			}
189			return Ok(None);
190		}
191
192		self.last_key = new_last_key;
193
194		let mut result_columns = Vec::new();
195
196		result_columns.push(ColumnWithName::new(
197			Fragment::internal(series.key.column()),
198			series.key_column_data(key_values),
199		));
200
201		if has_tag {
202			result_columns.push(ColumnWithName::new(Fragment::internal("tag"), ColumnBuffer::uint1(tags)));
203		}
204
205		for (col_idx, col_def) in series.data_columns().enumerate() {
206			let col_type = col_def.constraint.get_type();
207			let mut col_values: Vec<Value> = data_rows
208				.iter()
209				.map(|row| row.get(col_idx).cloned().unwrap_or(Value::none()))
210				.collect();
211
212			if let Some(dict_id) = col_def.dictionary_id
213				&& let Some(dictionary) = stored_ctx.services.catalog.find_dictionary(rx, dict_id)?
214			{
215				for value in col_values.iter_mut() {
216					if let Some(entry_id) = DictionaryEntryId::from_value(value) {
217						*value = rx
218							.get_from_dictionary(&dictionary, entry_id)?
219							.unwrap_or_else(Value::none);
220					}
221				}
222			}
223
224			result_columns.push(build_data_column(&col_def.name, &col_values, col_type)?);
225		}
226
227		let row_numbers: Vec<RowNumber> = sequences.into_iter().map(RowNumber::from).collect();
228		Ok(Some(Columns::with_system_columns(
229			result_columns,
230			row_numbers,
231			created_at_values,
232			updated_at_values,
233		)))
234	}
235
236	fn headers(&self) -> Option<ColumnHeaders> {
237		Some(self.headers.clone())
238	}
239}
240
241pub(crate) fn build_data_column(name: &str, values: &[Value], col_type: ValueType) -> Result<ColumnWithName> {
242	let data = match col_type {
243		ValueType::Boolean => {
244			let vals: Vec<bool> = values
245				.iter()
246				.map(|v| match v {
247					Value::Boolean(b) => *b,
248					_ => false,
249				})
250				.collect();
251			ColumnBuffer::bool(vals)
252		}
253		ValueType::Int1 => {
254			let vals: Vec<i8> = values
255				.iter()
256				.map(|v| match v {
257					Value::Int1(n) => *n,
258					_ => 0,
259				})
260				.collect();
261			ColumnBuffer::int1(vals)
262		}
263		ValueType::Int2 => {
264			let vals: Vec<i16> = values
265				.iter()
266				.map(|v| match v {
267					Value::Int2(n) => *n,
268					_ => 0,
269				})
270				.collect();
271			ColumnBuffer::int2(vals)
272		}
273		ValueType::Int4 => {
274			let vals: Vec<i32> = values
275				.iter()
276				.map(|v| match v {
277					Value::Int4(n) => *n,
278					_ => 0,
279				})
280				.collect();
281			ColumnBuffer::int4(vals)
282		}
283		ValueType::Int8 => {
284			let vals: Vec<i64> = values
285				.iter()
286				.map(|v| match v {
287					Value::Int8(n) => *n,
288					_ => 0,
289				})
290				.collect();
291			ColumnBuffer::int8(vals)
292		}
293		ValueType::Uint1 => {
294			let vals: Vec<u8> = values
295				.iter()
296				.map(|v| match v {
297					Value::Uint1(n) => *n,
298					_ => 0,
299				})
300				.collect();
301			ColumnBuffer::uint1(vals)
302		}
303		ValueType::Uint2 => {
304			let vals: Vec<u16> = values
305				.iter()
306				.map(|v| match v {
307					Value::Uint2(n) => *n,
308					_ => 0,
309				})
310				.collect();
311			ColumnBuffer::uint2(vals)
312		}
313		ValueType::Uint4 => {
314			let vals: Vec<u32> = values
315				.iter()
316				.map(|v| match v {
317					Value::Uint4(n) => *n,
318					_ => 0,
319				})
320				.collect();
321			ColumnBuffer::uint4(vals)
322		}
323		ValueType::Uint8 => {
324			let vals: Vec<u64> = values
325				.iter()
326				.map(|v| match v {
327					Value::Uint8(n) => *n,
328					_ => 0,
329				})
330				.collect();
331			ColumnBuffer::uint8(vals)
332		}
333		ValueType::Float4 => {
334			let vals: Vec<f32> = values
335				.iter()
336				.map(|v| match v {
337					Value::Float4(n) => n.value(),
338					_ => 0.0,
339				})
340				.collect();
341			ColumnBuffer::float4(vals)
342		}
343		ValueType::Float8 => {
344			let vals: Vec<f64> = values
345				.iter()
346				.map(|v| match v {
347					Value::Float8(n) => n.value(),
348					_ => 0.0,
349				})
350				.collect();
351			ColumnBuffer::float8(vals)
352		}
353		ValueType::Utf8 => {
354			let vals: Vec<String> = values
355				.iter()
356				.map(|v| match v {
357					Value::Utf8(s) => s.clone(),
358					_ => String::new(),
359				})
360				.collect();
361			ColumnBuffer::utf8(vals)
362		}
363		_ => {
364			let vals: Vec<String> = values.iter().map(|v| format!("{:?}", v)).collect();
365			ColumnBuffer::utf8(vals)
366		}
367	};
368
369	Ok(ColumnWithName {
370		name: Fragment::internal(name),
371		data,
372	})
373}