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