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