reifydb_engine/vm/volcano/scan/
dictionary.rs1use std::{ops::Bound, sync::Arc};
5
6use postcard::from_bytes;
7use reifydb_codec::key::encoded::{EncodedKey, EncodedKeyRange};
8use reifydb_core::{
9 interface::{catalog::dictionary::Dictionary, resolved::ResolvedDictionary, store::SingleVersionRange},
10 internal_error,
11 key::{EncodableKey, dictionary::DictionaryEntryIndexKey},
12 value::column::{ColumnWithName, buffer::ColumnBuffer, columns::Columns, headers::ColumnHeaders},
13};
14use reifydb_transaction::transaction::Transaction;
15use reifydb_value::{
16 fragment::Fragment,
17 reifydb_assertions,
18 value::{Value, dictionary::DictionaryEntryId, value_type::ValueType},
19};
20use tracing::instrument;
21
22use crate::{
23 Result,
24 vm::volcano::query::{QueryContext, QueryNode},
25};
26
27pub struct DictionaryScanNode {
28 dictionary: ResolvedDictionary,
29 context: Option<Arc<QueryContext>>,
30 headers: ColumnHeaders,
31 last_key: Option<EncodedKey>,
32 exhausted: bool,
33}
34
35impl DictionaryScanNode {
36 pub fn new(dictionary: ResolvedDictionary, context: Arc<QueryContext>) -> Result<Self> {
37 let headers = ColumnHeaders {
38 columns: vec![Fragment::internal("id"), Fragment::internal("value")],
39 };
40
41 Ok(Self {
42 dictionary,
43 context: Some(context),
44 headers,
45 last_key: None,
46 exhausted: false,
47 })
48 }
49
50 #[instrument(level = "trace", skip_all, name = "volcano::scan::dictionary::drain")]
51 fn drain_batch<'a>(
52 rx: &mut Transaction<'a>,
53 range: EncodedKeyRange,
54 batch_size: u64,
55 dict_def: &Dictionary,
56 ) -> Result<(Vec<DictionaryEntryId>, Vec<Value>, Option<EncodedKey>)> {
57 let mut ids: Vec<DictionaryEntryId> = Vec::new();
58 let mut values: Vec<Value> = Vec::new();
59 let mut new_last_key = None;
60
61 let single = rx
62 .single()
63 .ok_or_else(|| internal_error!("single-version store is not available for dictionary scans"))?;
64 let store = single.read_store();
65 let batch = SingleVersionRange::range_batch(&store, range, batch_size)?;
66
67 for entry in batch.items {
68 new_last_key = Some(entry.key.clone());
69
70 if let Some(key) = DictionaryEntryIndexKey::decode(&entry.key) {
71 let entry_id = DictionaryEntryId::from_u128(key.id, dict_def.id_type.clone())?;
72
73 let value: Value = from_bytes(&entry.bytes).map_err(|e| {
74 internal_error!("Failed to deserialize dictionary value: {}", e)
75 })?;
76
77 ids.push(entry_id);
78 values.push(value);
79 }
80 }
81
82 Ok((ids, values, new_last_key))
83 }
84
85 #[instrument(level = "trace", skip_all, name = "volcano::scan::dictionary::empty_columns")]
86 fn empty_columns(dict_def: &Dictionary) -> Vec<ColumnWithName> {
87 vec![
88 ColumnWithName {
89 name: Fragment::internal("id"),
90 data: ColumnBuffer::none_typed(dict_def.id_type.clone(), 0),
91 },
92 ColumnWithName {
93 name: Fragment::internal("value"),
94 data: ColumnBuffer::none_typed(dict_def.value_type.clone(), 0),
95 },
96 ]
97 }
98
99 #[instrument(level = "trace", skip_all, name = "volcano::scan::dictionary::assemble")]
100 fn assemble(ids: &[DictionaryEntryId], values: &[Value], dict_def: &Dictionary) -> Result<Option<Columns>> {
101 let id_column = build_id_column(ids, dict_def.id_type.clone())?;
102 let value_column = build_value_column(values, dict_def.value_type.clone())?;
103
104 Ok(Some(Columns::new(vec![id_column, value_column])))
105 }
106}
107
108impl QueryNode for DictionaryScanNode {
109 #[instrument(name = "volcano::scan::dictionary::initialize", level = "trace", skip_all)]
110 fn initialize<'a>(&mut self, _rx: &mut Transaction<'a>, _ctx: &QueryContext) -> Result<()> {
111 Ok(())
112 }
113
114 #[instrument(name = "volcano::scan::dictionary::next", level = "trace", skip_all)]
115 fn next<'a>(&mut self, rx: &mut Transaction<'a>, _ctx: &mut QueryContext) -> Result<Option<Columns>> {
116 reifydb_assertions! {
117 assert!(self.context.is_some(), "DictionaryScan::next() called before initialize()");
118 }
119 let stored_ctx = self.context.as_ref().unwrap();
120
121 if self.exhausted {
122 return Ok(None);
123 }
124
125 let batch_size = stored_ctx.batch_size;
126 let dict_def = self.dictionary.def();
127
128 let full_scan = DictionaryEntryIndexKey::full_scan(dict_def.id);
129 let range = match &self.last_key {
130 None => full_scan,
131 Some(last) => EncodedKeyRange::new(Bound::Excluded(last.clone()), full_scan.end),
132 };
133
134 let (ids, values, new_last_key) = Self::drain_batch(rx, range, batch_size, dict_def)?;
135
136 if ids.is_empty() {
137 self.exhausted = true;
138 if self.last_key.is_none() {
139 return Ok(Some(Columns::new(Self::empty_columns(dict_def))));
140 }
141 return Ok(None);
142 }
143
144 self.last_key = new_last_key;
145
146 Self::assemble(&ids, &values, dict_def)
147 }
148
149 fn headers(&self) -> Option<ColumnHeaders> {
150 Some(self.headers.clone())
151 }
152}
153
154fn build_id_column(ids: &[DictionaryEntryId], id_type: ValueType) -> Result<ColumnWithName> {
155 let data = match id_type {
156 ValueType::Uint1 => {
157 let vals: Vec<u8> = ids.iter().map(|id| id.to_u128() as u8).collect();
158 ColumnBuffer::uint1(vals)
159 }
160 ValueType::Uint2 => {
161 let vals: Vec<u16> = ids.iter().map(|id| id.to_u128() as u16).collect();
162 ColumnBuffer::uint2(vals)
163 }
164 ValueType::Uint4 => {
165 let vals: Vec<u32> = ids.iter().map(|id| id.to_u128() as u32).collect();
166 ColumnBuffer::uint4(vals)
167 }
168 ValueType::Uint8 => {
169 let vals: Vec<u64> = ids.iter().map(|id| id.to_u128() as u64).collect();
170 ColumnBuffer::uint8(vals)
171 }
172 ValueType::Uint16 => {
173 let vals: Vec<u128> = ids.iter().map(|id| id.to_u128()).collect();
174 ColumnBuffer::uint16(vals)
175 }
176 _ => return Err(internal_error!("Invalid dictionary id_type: {:?}", id_type)),
177 };
178
179 Ok(ColumnWithName {
180 name: Fragment::internal("id"),
181 data,
182 })
183}
184
185fn build_value_column(values: &[Value], value_type: ValueType) -> Result<ColumnWithName> {
186 let data = match value_type {
187 ValueType::Utf8 => {
188 let vals: Vec<String> = values
189 .iter()
190 .map(|v| match v {
191 Value::Utf8(s) => s.clone(),
192 _ => format!("{:?}", v),
193 })
194 .collect();
195 ColumnBuffer::utf8(vals)
196 }
197 ValueType::Int1 => {
198 let vals: Vec<i8> = values
199 .iter()
200 .map(|v| match v {
201 Value::Int1(n) => *n,
202 _ => 0,
203 })
204 .collect();
205 ColumnBuffer::int1(vals)
206 }
207 ValueType::Int2 => {
208 let vals: Vec<i16> = values
209 .iter()
210 .map(|v| match v {
211 Value::Int2(n) => *n,
212 _ => 0,
213 })
214 .collect();
215 ColumnBuffer::int2(vals)
216 }
217 ValueType::Int4 => {
218 let vals: Vec<i32> = values
219 .iter()
220 .map(|v| match v {
221 Value::Int4(n) => *n,
222 _ => 0,
223 })
224 .collect();
225 ColumnBuffer::int4(vals)
226 }
227 ValueType::Int8 => {
228 let vals: Vec<i64> = values
229 .iter()
230 .map(|v| match v {
231 Value::Int8(n) => *n,
232 _ => 0,
233 })
234 .collect();
235 ColumnBuffer::int8(vals)
236 }
237 ValueType::Uint1 => {
238 let vals: Vec<u8> = values
239 .iter()
240 .map(|v| match v {
241 Value::Uint1(n) => *n,
242 _ => 0,
243 })
244 .collect();
245 ColumnBuffer::uint1(vals)
246 }
247 ValueType::Uint2 => {
248 let vals: Vec<u16> = values
249 .iter()
250 .map(|v| match v {
251 Value::Uint2(n) => *n,
252 _ => 0,
253 })
254 .collect();
255 ColumnBuffer::uint2(vals)
256 }
257 ValueType::Uint4 => {
258 let vals: Vec<u32> = values
259 .iter()
260 .map(|v| match v {
261 Value::Uint4(n) => *n,
262 _ => 0,
263 })
264 .collect();
265 ColumnBuffer::uint4(vals)
266 }
267 ValueType::Uint8 => {
268 let vals: Vec<u64> = values
269 .iter()
270 .map(|v| match v {
271 Value::Uint8(n) => *n,
272 _ => 0,
273 })
274 .collect();
275 ColumnBuffer::uint8(vals)
276 }
277 _ => {
278 let vals: Vec<String> = values.iter().map(|v| format!("{:?}", v)).collect();
279 ColumnBuffer::utf8(vals)
280 }
281 };
282
283 Ok(ColumnWithName {
284 name: Fragment::internal("value"),
285 data,
286 })
287}