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