reifydb_engine/vm/volcano/scan/
table.rs1use std::sync::Arc;
5
6use reifydb_codec::{
7 key::encoded::{EncodedKey, EncodedKeyRange},
8 row::{bytes::EncodedBytes, shape::RowShape, table::EncodedTableRow},
9};
10use reifydb_core::{
11 common::CommitVersion,
12 error::diagnostic,
13 interface::{catalog::dictionary::Dictionary, resolved::ResolvedTable, store::MultiVersionRow},
14 key::{
15 EncodableKey,
16 partitioned_row::{PartitionedRowKey, RowLocator},
17 row::{RowKey, RowKeyRange},
18 },
19 value::column::{ColumnWithName, buffer::ColumnBuffer, columns::Columns, headers::ColumnHeaders},
20};
21use reifydb_transaction::{multi::RangeScope, transaction::Transaction};
22use reifydb_value::{
23 error,
24 fragment::Fragment,
25 reifydb_assertions,
26 value::{partition::Partition, row_number::RowNumber, system_columns::SystemColumns, value_type::ValueType},
27};
28use tracing::instrument;
29
30use super::super::decode_dictionary_columns;
31use crate::{
32 Result,
33 vm::volcano::query::{QueryContext, QueryNode},
34};
35
36pub struct TableScanNode {
37 table: ResolvedTable,
38 context: Option<Arc<QueryContext>>,
39 headers: ColumnHeaders,
40
41 storage_types: Vec<ValueType>,
42
43 dictionaries: Vec<Option<Dictionary>>,
44
45 shape: Option<RowShape>,
46 last_key: Option<EncodedKey>,
47 exhausted: bool,
48
49 partition: Option<Partition>,
50
51 min_commit_version: Option<CommitVersion>,
52}
53
54impl TableScanNode {
55 pub fn with_min_commit_version(mut self, min_commit_version: Option<CommitVersion>) -> Self {
56 self.min_commit_version = min_commit_version;
57 self
58 }
59
60 pub fn new(
61 table: ResolvedTable,
62 partition: Option<Partition>,
63 context: Arc<QueryContext>,
64 rx: &mut Transaction<'_>,
65 ) -> Result<Self> {
66 let mut storage_types = Vec::with_capacity(table.columns().len());
67 let mut dictionaries = Vec::with_capacity(table.columns().len());
68
69 for col in table.columns() {
70 if let Some(dict_id) = col.dictionary_id {
71 if let Some(dict) = context.services.catalog.find_dictionary(rx, dict_id)? {
72 storage_types.push(ValueType::DictionaryId);
73 dictionaries.push(Some(dict));
74 } else {
75 storage_types.push(col.constraint.get_type());
76 dictionaries.push(None);
77 }
78 } else {
79 storage_types.push(col.constraint.get_type());
80 dictionaries.push(None);
81 }
82 }
83
84 let headers = ColumnHeaders {
85 columns: table.columns().iter().map(|col| Fragment::internal(&col.name)).collect(),
86 };
87
88 Ok(Self {
89 table,
90 context: Some(context),
91 headers,
92 storage_types,
93 dictionaries,
94 shape: None,
95 last_key: None,
96 exhausted: false,
97 partition,
98 min_commit_version: None,
99 })
100 }
101
102 fn get_or_load_shape<'a>(&mut self, rx: &mut Transaction<'a>, first: &EncodedBytes) -> Result<RowShape> {
103 if let Some(shape) = &self.shape {
104 return Ok(shape.clone());
105 }
106
107 let fingerprint = EncodedTableRow::view(first).fingerprint();
108
109 let stored_ctx = self.context.as_ref().expect("TableScanNode context not set");
110 let shape = stored_ctx.services.catalog.get_or_load_row_shape(fingerprint, rx)?.ok_or_else(|| {
111 error!(diagnostic::internal::internal(format!(
112 "RowShape with fingerprint {:?} not found for table {}",
113 fingerprint,
114 self.table.def().name
115 )))
116 })?;
117
118 self.shape = Some(shape.clone());
119
120 Ok(shape)
121 }
122
123 #[instrument(level = "trace", skip_all, name = "volcano::scan::table::range_open")]
124 fn open_range<'rx, 'tx>(
125 rx: &'rx mut Transaction<'tx>,
126 range: EncodedKeyRange,
127 scope: RangeScope,
128 batch_size: u64,
129 ) -> Result<Box<dyn Iterator<Item = Result<MultiVersionRow>> + Send + 'rx>> {
130 rx.range(range, scope, batch_size as usize)
131 }
132
133 #[instrument(level = "trace", skip_all, name = "volcano::scan::table::drain")]
134 fn drain_batch(
135 stream: &mut dyn Iterator<Item = Result<MultiVersionRow>>,
136 batch_size: u64,
137 partitioned: bool,
138 ) -> Result<ScannedBatch> {
139 let mut batch = ScannedBatch::default();
140
141 for _ in 0..batch_size {
142 match stream.next() {
143 Some(Ok(multi)) => {
144 let decoded = if partitioned {
145 PartitionedRowKey::decode(&multi.key).and_then(|k| match k.locator {
146 RowLocator::Row(rn) => Some((rn, Some(k.partition))),
147 _ => None,
148 })
149 } else {
150 RowKey::decode(&multi.key).map(|k| (k.row, None))
151 };
152 if let Some((rn, partition)) = decoded {
153 batch.rows.push(multi.bytes);
154 batch.row_numbers.push(rn);
155 if let Some(p) = partition {
156 batch.partitions.push(p);
157 }
158 batch.last_key = Some(multi.key);
159 }
160 }
161 Some(Err(e)) => return Err(e),
162 None => {
163 batch.exhausted = true;
164 break;
165 }
166 }
167 }
168
169 Ok(batch)
170 }
171
172 #[instrument(level = "trace", skip_all, name = "volcano::scan::table::column_alloc")]
173 fn storage_columns(&self) -> Vec<ColumnWithName> {
174 self.table
175 .columns()
176 .iter()
177 .enumerate()
178 .map(|(idx, col)| ColumnWithName {
179 name: Fragment::internal(&col.name),
180 data: ColumnBuffer::with_capacity(self.storage_types[idx].clone(), 0),
181 })
182 .collect()
183 }
184
185 #[instrument(level = "trace", skip_all, name = "volcano::scan::table::empty_columns")]
186 fn empty_columns(&self) -> Vec<ColumnWithName> {
187 self.table
188 .columns()
189 .iter()
190 .map(|col| ColumnWithName {
191 name: Fragment::internal(&col.name),
192 data: ColumnBuffer::none_typed(col.constraint.get_type(), 0),
193 })
194 .collect()
195 }
196
197 #[instrument(level = "trace", skip_all, name = "volcano::scan::table::append_rows")]
198 fn append_batch<'a>(
199 &mut self,
200 rx: &mut Transaction<'a>,
201 columns: &mut Columns,
202 bytes_vec: Vec<EncodedBytes>,
203 row_numbers: Vec<RowNumber>,
204 ) -> Result<()> {
205 let shape = self.get_or_load_shape(rx, &bytes_vec[0])?;
206 columns.append_rows(&shape, bytes_vec.into_iter(), row_numbers)?;
207 Ok(())
208 }
209}
210
211#[derive(Default)]
212struct ScannedBatch {
213 rows: Vec<EncodedBytes>,
214 row_numbers: Vec<RowNumber>,
215 partitions: Vec<Partition>,
216 last_key: Option<EncodedKey>,
217 exhausted: bool,
218}
219
220impl QueryNode for TableScanNode {
221 #[instrument(level = "trace", skip_all, name = "volcano::scan::table::initialize")]
222 fn initialize<'a>(&mut self, _rx: &mut Transaction<'a>, _ctx: &QueryContext) -> Result<()> {
223 Ok(())
224 }
225
226 #[instrument(level = "trace", skip_all, name = "volcano::scan::table::next")]
227 fn next<'a>(&mut self, rx: &mut Transaction<'a>, _ctx: &mut QueryContext) -> Result<Option<Columns>> {
228 reifydb_assertions! {
229 assert!(self.context.is_some(), "TableScanNode::next() called before initialize()");
230 }
231 let stored_ctx = self.context.as_ref().unwrap();
232
233 if self.exhausted {
234 return Ok(None);
235 }
236
237 let batch_size = stored_ctx.batch_size;
238
239 let partitioned = !self.table.def().partition_by.is_empty();
240 let range = if partitioned {
241 match self.partition {
242 Some(partition) => PartitionedRowKey::partition_scan_range(
243 self.table.def().id,
244 partition,
245 self.last_key.as_ref(),
246 ),
247 None => PartitionedRowKey::scan_range(self.table.def().id, self.last_key.as_ref()),
248 }
249 } else {
250 RowKeyRange::scan_range(self.table.def().id.into(), self.last_key.as_ref())
251 };
252
253 let scope = match self.min_commit_version {
254 Some(v) => RangeScope::After(v),
255 None => RangeScope::All,
256 };
257
258 let batch = {
259 let mut stream = Self::open_range(rx, range, scope, batch_size)?;
260 Self::drain_batch(&mut stream, batch_size, partitioned)?
261 };
262
263 if batch.exhausted {
264 self.exhausted = true;
265 }
266
267 if batch.rows.is_empty() {
268 self.exhausted = true;
269 if self.last_key.is_none() {
270 return Ok(Some(Columns::new(self.empty_columns())));
271 }
272 return Ok(None);
273 }
274
275 self.last_key = batch.last_key;
276
277 let mut columns = Columns::with_system(self.storage_columns(), SystemColumns::default());
278 self.append_batch(rx, &mut columns, batch.rows, batch.row_numbers)?;
279
280 if partitioned {
281 columns.system.set_partitions(batch.partitions);
282 }
283
284 decode_dictionary_columns(&mut columns, &self.dictionaries, rx)?;
285
286 Ok(Some(columns))
287 }
288
289 fn headers(&self) -> Option<ColumnHeaders> {
290 Some(self.headers.clone())
291 }
292}