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