Skip to main content

reifydb_engine/vm/volcano/scan/
dictionary.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use 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}