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::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}