reifydb_engine/vm/volcano/scan/
queue.rs1use std::sync::Arc;
5
6use reifydb_codec::row::{bytes::EncodedBytes, queue::EncodedQueueRow, shape::RowShape};
7use reifydb_core::{
8 interface::{catalog::dictionary::Dictionary, resolved::ResolvedQueue, store::MultiVersionRow},
9 internal_error,
10 key::{any::TaggedKey, bound::TaggedKeyBoundRange, row::RowKeyRange},
11 value::column::{ColumnWithName, buffer::ColumnBuffer, columns::Columns, headers::ColumnHeaders},
12};
13use reifydb_transaction::{multi::RangeScope, transaction::Transaction};
14use reifydb_value::{
15 fragment::Fragment,
16 value::{row_number::RowNumber, system_columns::SystemColumns, value_type::ValueType},
17};
18use tracing::instrument;
19
20use super::super::decode_dictionary_columns;
21use crate::{
22 Result,
23 vm::volcano::query::{QueryContext, QueryNode},
24};
25
26type DrainedBatch = (Vec<EncodedBytes>, Vec<RowNumber>, Option<TaggedKey>, bool);
27
28pub struct QueueScan {
29 queue: ResolvedQueue,
30 headers: ColumnHeaders,
31 shape: Option<RowShape>,
32 storage_types: Vec<ValueType>,
33 dictionaries: Vec<Option<Dictionary>>,
34 last_key: Option<TaggedKey>,
35 exhausted: bool,
36 context: Option<Arc<QueryContext>>,
37}
38
39impl QueueScan {
40 pub fn new(queue: ResolvedQueue, context: Arc<QueryContext>, rx: &mut Transaction<'_>) -> Result<Self> {
41 let mut storage_types = Vec::with_capacity(queue.columns().len());
42 let mut dictionaries = Vec::with_capacity(queue.columns().len());
43
44 for col in queue.columns() {
45 if let Some(dict_id) = col.dictionary_id
46 && let Some(dict) = context.services.catalog.find_dictionary(rx, dict_id)?
47 {
48 storage_types.push(ValueType::DictionaryId);
49 dictionaries.push(Some(dict));
50 continue;
51 }
52 storage_types.push(col.constraint.get_type());
53 dictionaries.push(None);
54 }
55
56 let headers = ColumnHeaders {
57 columns: queue.columns().iter().map(|col| Fragment::internal(&col.name)).collect(),
58 };
59
60 Ok(Self {
61 queue,
62 headers,
63 shape: None,
64 storage_types,
65 dictionaries,
66 last_key: None,
67 exhausted: false,
68 context: Some(context),
69 })
70 }
71
72 fn get_or_load_shape(&mut self, rx: &mut Transaction, first: &EncodedBytes) -> Result<RowShape> {
73 if let Some(shape) = &self.shape {
74 return Ok(shape.clone());
75 }
76
77 let fingerprint = EncodedQueueRow::view(first).fingerprint();
78 let stored_ctx = self.context.as_ref().expect("QueueScan context not set");
79 let shape = stored_ctx.services.catalog.get_or_load_row_shape(fingerprint, rx)?.ok_or_else(|| {
80 internal_error!(
81 "RowShape with fingerprint {:?} not found for queue {}",
82 fingerprint,
83 self.queue.def().name
84 )
85 })?;
86
87 self.shape = Some(shape.clone());
88
89 Ok(shape)
90 }
91
92 #[instrument(level = "trace", skip_all, name = "volcano::scan::queue::range_open")]
93 fn open_range<'rx, 'tx>(
94 rx: &'rx mut Transaction<'tx>,
95 range: TaggedKeyBoundRange,
96 batch_size: u64,
97 ) -> Result<Box<dyn Iterator<Item = Result<MultiVersionRow<TaggedKey>>> + Send + 'rx>> {
98 rx.range_rev(range, RangeScope::All, batch_size as usize)
99 }
100
101 #[instrument(level = "trace", skip_all, name = "volcano::scan::queue::drain")]
102 fn drain_batch(
103 stream: &mut dyn Iterator<Item = Result<MultiVersionRow<TaggedKey>>>,
104 batch_size: u64,
105 ) -> Result<DrainedBatch> {
106 let mut batch: Vec<EncodedBytes> = Vec::new();
107 let mut row_numbers: Vec<RowNumber> = Vec::new();
108 let mut new_last_key = None;
109 let mut drained = false;
110
111 for _ in 0..batch_size {
112 match stream.next() {
113 Some(Ok(multi)) => {
114 if let TaggedKey::Row(key) = &multi.key {
115 row_numbers.push(key.row);
116 batch.push(multi.bytes);
117 new_last_key = Some(multi.key);
118 }
119 }
120 Some(Err(e)) => return Err(e),
121 None => {
122 drained = true;
123 break;
124 }
125 }
126 }
127
128 Ok((batch, row_numbers, new_last_key, drained))
129 }
130
131 #[instrument(level = "trace", skip_all, name = "volcano::scan::queue::column_alloc")]
132 fn storage_columns(&self, shape: &RowShape, declared: usize) -> Result<Vec<ColumnWithName>> {
133 let mut storage_columns: Vec<ColumnWithName> = self
134 .queue
135 .columns()
136 .iter()
137 .enumerate()
138 .map(|(idx, col)| ColumnWithName {
139 name: Fragment::internal(&col.name),
140 data: ColumnBuffer::with_capacity(self.storage_types[idx].clone(), 0),
141 })
142 .collect();
143
144 for index in declared..shape.field_count() {
145 let field = shape.get_field(index).ok_or_else(|| {
146 internal_error!("queue {} shape lost field {}", self.queue.def().name, index)
147 })?;
148 storage_columns.push(ColumnWithName {
149 name: Fragment::internal(field.name.clone()),
150 data: ColumnBuffer::with_capacity(field.constraint.get_type(), 0),
151 });
152 }
153
154 Ok(storage_columns)
155 }
156
157 #[instrument(level = "trace", skip_all, name = "volcano::scan::queue::append_rows")]
158 fn append_batch(
159 shape: &RowShape,
160 columns: &mut Columns,
161 bytes_vec: Vec<EncodedBytes>,
162 row_numbers: Vec<RowNumber>,
163 ) -> Result<()> {
164 columns.append_rows(shape, bytes_vec.into_iter(), row_numbers)?;
165 Ok(())
166 }
167
168 fn enqueue_order_range(&self) -> TaggedKeyBoundRange {
169 RowKeyRange::scan_range(self.queue.def().id.into(), None).resume_before(self.last_key.as_ref())
170 }
171
172 fn empty_declared_columns(&self) -> Columns {
173 Columns::new(
174 self.queue
175 .columns()
176 .iter()
177 .map(|col| ColumnWithName {
178 name: Fragment::internal(&col.name),
179 data: ColumnBuffer::none_typed(col.constraint.get_type(), 0),
180 })
181 .collect(),
182 )
183 }
184}
185
186impl QueryNode for QueueScan {
187 #[instrument(level = "trace", skip_all, name = "volcano::scan::queue::initialize")]
188 fn initialize<'a>(&mut self, _rx: &mut Transaction<'a>, _ctx: &QueryContext) -> Result<()> {
189 Ok(())
190 }
191
192 #[instrument(level = "trace", skip_all, name = "volcano::scan::queue::next")]
193 fn next<'a>(&mut self, rx: &mut Transaction<'a>, _ctx: &mut QueryContext) -> Result<Option<Columns>> {
194 if self.exhausted {
195 return Ok(None);
196 }
197
198 let batch_size = self.context.as_ref().expect("QueueScan context not set").batch_size;
199 let range = self.enqueue_order_range();
200
201 let (batch_rows, row_numbers, new_last_key, drained) = {
202 let mut stream = Self::open_range(rx, range, batch_size)?;
203 Self::drain_batch(&mut stream, batch_size)?
204 };
205
206 if drained {
207 self.exhausted = true;
208 }
209
210 if batch_rows.is_empty() {
211 self.exhausted = true;
212 if self.last_key.is_none() {
213 return Ok(Some(self.empty_declared_columns()));
214 }
215 return Ok(None);
216 }
217
218 self.last_key = new_last_key;
219
220 let shape = self.get_or_load_shape(rx, &batch_rows[0])?;
221 let declared = self.queue.columns().len();
222
223 let storage_columns = self.storage_columns(&shape, declared)?;
224
225 let mut columns = Columns::with_system(storage_columns, SystemColumns::default());
226 Self::append_batch(&shape, &mut columns, batch_rows, row_numbers)?;
227
228 decode_dictionary_columns(&mut columns, &self.dictionaries, rx)?;
229
230 columns.columns.truncate(declared);
231 columns.names.truncate(declared);
232
233 Ok(Some(columns))
234 }
235
236 fn headers(&self) -> Option<ColumnHeaders> {
237 Some(self.headers.clone())
238 }
239}