reifydb_engine/vm/volcano/scan/
ringbuffer.rs1use std::sync::Arc;
5
6use reifydb_codec::row::{bytes::EncodedBytes, ringbuffer::EncodedRingBufferRow, shape::RowShape};
7use reifydb_core::{
8 interface::{
9 catalog::{dictionary::Dictionary, ringbuffer::PartitionedMetadata},
10 resolved::ResolvedRingBuffer,
11 },
12 internal_error,
13 key::{
14 any::TaggedKey,
15 row::{PartitionedRowKey, RowKey},
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::{partition::Partition, 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
32pub struct RingBufferScan {
33 ringbuffer: ResolvedRingBuffer,
34
35 partitions: Vec<PartitionedMetadata>,
36 current_partition_index: usize,
37 headers: ColumnHeaders,
38 shape: Option<RowShape>,
39
40 storage_types: Vec<ValueType>,
41
42 dictionaries: Vec<Option<Dictionary>>,
43
44 partition_col_indices: Vec<usize>,
45 current_partition_rows: Vec<(RowNumber, EncodedBytes)>,
46 current_partition_cursor: usize,
47 current_partition_loaded: bool,
48 finished: bool,
49 context: Option<Arc<QueryContext>>,
50 initialized: bool,
51}
52
53impl RingBufferScan {
54 pub fn new(
55 ringbuffer: ResolvedRingBuffer,
56 context: Arc<QueryContext>,
57 rx: &mut Transaction<'_>,
58 ) -> Result<Self> {
59 let mut storage_types = Vec::with_capacity(ringbuffer.columns().len());
60 let mut dictionaries = Vec::with_capacity(ringbuffer.columns().len());
61
62 for col in ringbuffer.columns() {
63 if let Some(dict_id) = col.dictionary_id {
64 if let Some(dict) = context.services.catalog.find_dictionary(rx, dict_id)? {
65 storage_types.push(ValueType::DictionaryId);
66 dictionaries.push(Some(dict));
67 } else {
68 storage_types.push(col.constraint.get_type());
69 dictionaries.push(None);
70 }
71 } else {
72 storage_types.push(col.constraint.get_type());
73 dictionaries.push(None);
74 }
75 }
76
77 let partition_col_indices: Vec<usize> = ringbuffer
78 .def()
79 .partition_by
80 .iter()
81 .map(|pb_col| ringbuffer.columns().iter().position(|c| c.name == *pb_col).unwrap())
82 .collect();
83
84 let headers = ColumnHeaders {
85 columns: ringbuffer.columns().iter().map(|col| Fragment::internal(&col.name)).collect(),
86 };
87
88 Ok(Self {
89 ringbuffer,
90 partitions: Vec::new(),
91 current_partition_index: 0,
92 headers,
93 shape: None,
94 storage_types,
95 dictionaries,
96 partition_col_indices,
97 current_partition_rows: Vec::new(),
98 current_partition_cursor: 0,
99 current_partition_loaded: false,
100 finished: false,
101 context: Some(context),
102 initialized: false,
103 })
104 }
105
106 fn get_or_load_shape(&mut self, rx: &mut Transaction, first: &EncodedBytes) -> Result<RowShape> {
107 if let Some(shape) = &self.shape {
108 return Ok(shape.clone());
109 }
110
111 let fingerprint = EncodedRingBufferRow::view(first).fingerprint();
112
113 let stored_ctx = self.context.as_ref().expect("RingBufferScan context not set");
114 let shape = stored_ctx.services.catalog.get_or_load_row_shape(fingerprint, rx)?.ok_or_else(|| {
115 internal_error!(
116 "RowShape with fingerprint {:?} not found for ringbuffer {}",
117 fingerprint,
118 self.ringbuffer.def().name
119 )
120 })?;
121
122 self.shape = Some(shape.clone());
123
124 Ok(shape)
125 }
126
127 #[instrument(level = "trace", skip_all, name = "volcano::scan::ringbuffer::drain")]
128 fn drain_batch(
129 &mut self,
130 txn: &mut Transaction<'_>,
131 batch_size: usize,
132 partitioned: bool,
133 ) -> Result<(Vec<EncodedBytes>, Vec<RowNumber>, Vec<Partition>)> {
134 let mut batch: Vec<EncodedBytes> = Vec::new();
135 let mut row_numbers: Vec<RowNumber> = Vec::new();
136 let mut partitions_sidecar: Vec<Partition> = Vec::new();
137
138 while batch.len() < batch_size && self.current_partition_index < self.partitions.len() {
139 if !self.current_partition_loaded {
140 self.current_partition_rows =
141 self.load_partition_rows(txn, self.current_partition_index)?;
142 self.current_partition_cursor = 0;
143 self.current_partition_loaded = true;
144 }
145
146 let hash = if partitioned {
147 Some(Partition::of(&self.partitions[self.current_partition_index].partition_values))
148 } else {
149 None
150 };
151
152 while batch.len() < batch_size
153 && self.current_partition_cursor < self.current_partition_rows.len()
154 {
155 let (rn, row) = self.current_partition_rows[self.current_partition_cursor].clone();
156 batch.push(row);
157 row_numbers.push(rn);
158 if let Some(h) = hash {
159 partitions_sidecar.push(h);
160 }
161 self.current_partition_cursor += 1;
162 }
163
164 if self.current_partition_cursor >= self.current_partition_rows.len() {
165 self.current_partition_index += 1;
166 self.current_partition_loaded = false;
167 }
168 }
169
170 Ok((batch, row_numbers, partitions_sidecar))
171 }
172
173 #[instrument(level = "trace", skip_all, name = "volcano::scan::ringbuffer::column_alloc")]
174 fn storage_columns(&self) -> Vec<ColumnWithName> {
175 self.ringbuffer
176 .columns()
177 .iter()
178 .enumerate()
179 .map(|(idx, col)| ColumnWithName {
180 name: Fragment::internal(&col.name),
181 data: ColumnBuffer::with_capacity(self.storage_types[idx].clone(), 0),
182 })
183 .collect()
184 }
185
186 #[instrument(level = "trace", skip_all, name = "volcano::scan::ringbuffer::empty_columns")]
187 fn empty_columns(&self) -> Vec<ColumnWithName> {
188 self.ringbuffer
189 .columns()
190 .iter()
191 .map(|col| ColumnWithName {
192 name: Fragment::internal(&col.name),
193 data: ColumnBuffer::none_typed(col.constraint.get_type(), 0),
194 })
195 .collect()
196 }
197
198 #[instrument(level = "trace", skip_all, name = "volcano::scan::ringbuffer::append_rows")]
199 fn append_batch(
200 &mut self,
201 txn: &mut Transaction<'_>,
202 columns: &mut Columns,
203 bytes_vec: Vec<EncodedBytes>,
204 row_numbers: Vec<RowNumber>,
205 ) -> Result<()> {
206 let shape = self.get_or_load_shape(txn, &bytes_vec[0])?;
207 columns.append_rows(&shape, bytes_vec.into_iter(), row_numbers.clone())?;
208 Ok(())
209 }
210
211 #[instrument(level = "trace", skip_all, name = "volcano::scan::ringbuffer::load_partition")]
212 fn load_partition_rows(
213 &self,
214 txn: &mut Transaction<'_>,
215 partition_index: usize,
216 ) -> Result<Vec<(RowNumber, EncodedBytes)>> {
217 let pm = &self.partitions[partition_index];
218 let rb_id = self.ringbuffer.def().id;
219
220 if self.partition_col_indices.is_empty() {
221 let mut out = Vec::new();
222 for rn_value in pm.metadata.head..pm.metadata.tail {
223 let rn = RowNumber(rn_value);
224 if let Some(multi) = txn.get(&RowKey::new(rb_id, rn))? {
225 out.push((rn, multi.bytes));
226 }
227 }
228 return Ok(out);
229 }
230
231 let hash = Partition::of(&pm.partition_values);
232 let mut out = Vec::new();
233 let mut last_key = None;
234 loop {
235 let batch: Vec<_> = txn
236 .range(
237 PartitionedRowKey::partition_scan_range(rb_id, hash, last_key.as_ref()),
238 RangeScope::All,
239 1024,
240 )?
241 .collect::<Result<Vec<_>>>()?;
242 if batch.is_empty() {
243 break;
244 }
245 let n = batch.len();
246 for entry in batch {
247 if let TaggedKey::PartitionedRow(pk) = &entry.key {
248 out.push((pk.row, entry.bytes));
249 }
250 last_key = Some(entry.key.clone());
251 }
252 if n < 1024 {
253 break;
254 }
255 }
256 out.sort_by_key(|(rn, _)| rn.0);
257 Ok(out)
258 }
259}
260
261impl QueryNode for RingBufferScan {
262 #[instrument(name = "volcano::scan::ringbuffer::initialize", level = "trace", skip_all)]
263 fn initialize<'a>(&mut self, txn: &mut Transaction<'a>, ctx: &QueryContext) -> Result<()> {
264 if !self.initialized {
265 self.partitions =
266 ctx.services.catalog.list_ringbuffer_partitions(txn, self.ringbuffer.def())?;
267 self.initialized = true;
268 }
269 Ok(())
270 }
271
272 #[instrument(name = "volcano::scan::ringbuffer::next", level = "trace", skip_all)]
273 fn next<'a>(&mut self, txn: &mut Transaction<'a>, _ctx: &mut QueryContext) -> Result<Option<Columns>> {
274 if self.finished {
275 return Ok(None);
276 }
277
278 let batch_size = self.context.as_ref().expect("RingBufferScan context not set").batch_size as usize;
279 let partitioned = !self.partition_col_indices.is_empty();
280
281 let (batch_rows, row_numbers, partitions_sidecar) = self.drain_batch(txn, batch_size, partitioned)?;
282
283 if !batch_rows.is_empty() {
284 let mut columns = Columns::with_system(self.storage_columns(), SystemColumns::default());
285 self.append_batch(txn, &mut columns, batch_rows, row_numbers)?;
286 if partitioned {
287 columns.system.set_partitions(partitions_sidecar);
288 }
289
290 decode_dictionary_columns(&mut columns, &self.dictionaries, txn)?;
291
292 return Ok(Some(columns));
293 }
294
295 self.finished = true;
296 if self.partitions.is_empty() || self.partitions.iter().all(|p| p.metadata.is_empty()) {
297 return Ok(Some(Columns::new(self.empty_columns())));
298 }
299 Ok(None)
300 }
301
302 fn headers(&self) -> Option<ColumnHeaders> {
303 Some(self.headers.clone())
304 }
305}