reifydb_engine/vm/volcano/scan/
series.rs1use 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}