Skip to main content

reifydb_engine/vm/volcano/
merge.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use reifydb_core::value::column::{columns::Columns, headers::ColumnHeaders};
5use reifydb_transaction::transaction::Transaction;
6use tracing::instrument;
7
8use crate::{
9	Result,
10	vm::volcano::query::{QueryContext, QueryNode},
11};
12
13pub struct DeltaMergeNode {
14	inputs: Vec<Box<dyn QueryNode>>,
15	cursor: usize,
16}
17
18impl DeltaMergeNode {
19	pub fn new(inputs: Vec<Box<dyn QueryNode>>) -> Self {
20		Self {
21			inputs,
22			cursor: 0,
23		}
24	}
25}
26
27impl QueryNode for DeltaMergeNode {
28	#[instrument(name = "volcano::merge::initialize", level = "trace", skip_all)]
29	fn initialize<'a>(&mut self, rx: &mut Transaction<'a>, ctx: &QueryContext) -> Result<()> {
30		for input in &mut self.inputs {
31			input.initialize(rx, ctx)?;
32		}
33		Ok(())
34	}
35
36	#[instrument(name = "volcano::merge::next", level = "trace", skip_all)]
37	fn next<'a>(&mut self, rx: &mut Transaction<'a>, ctx: &mut QueryContext) -> Result<Option<Columns>> {
38		while self.cursor < self.inputs.len() {
39			match self.inputs[self.cursor].next(rx, ctx)? {
40				Some(columns) => return Ok(Some(columns)),
41				None => self.cursor += 1,
42			}
43		}
44		Ok(None)
45	}
46
47	fn headers(&self) -> Option<ColumnHeaders> {
48		self.inputs.first().and_then(|input| input.headers())
49	}
50}
51
52#[cfg(test)]
53mod tests {
54	use std::collections::VecDeque;
55
56	use reifydb_core::value::column::{ColumnWithName, buffer::ColumnBuffer, headers::ColumnHeaders};
57	use reifydb_evaluate::stack::SymbolTable;
58	use reifydb_test_harness::engine::create_test_admin_transaction;
59	use reifydb_value::{
60		fragment::Fragment,
61		params::Params,
62		value::{Value, identity::IdentityId},
63	};
64
65	use super::*;
66	use crate::vm::{services::Services, volcano::query::query_budget};
67
68	struct StubNode {
69		batches: VecDeque<Columns>,
70		headers: Option<ColumnHeaders>,
71		init_count: usize,
72	}
73
74	impl StubNode {
75		fn new(batches: Vec<Columns>, headers: Option<ColumnHeaders>) -> Self {
76			Self {
77				batches: batches.into(),
78				headers,
79				init_count: 0,
80			}
81		}
82	}
83
84	impl QueryNode for StubNode {
85		fn initialize<'a>(&mut self, _rx: &mut Transaction<'a>, _ctx: &QueryContext) -> Result<()> {
86			self.init_count += 1;
87			Ok(())
88		}
89
90		fn next<'a>(&mut self, _rx: &mut Transaction<'a>, _ctx: &mut QueryContext) -> Result<Option<Columns>> {
91			Ok(self.batches.pop_front())
92		}
93
94		fn headers(&self) -> Option<ColumnHeaders> {
95			self.headers.clone()
96		}
97	}
98
99	fn batch(name: &str, vals: Vec<i32>) -> Columns {
100		Columns::new(vec![ColumnWithName {
101			name: Fragment::internal(name),
102			data: ColumnBuffer::int4(vals),
103		}])
104	}
105
106	fn first_int4(columns: &Columns) -> Vec<i32> {
107		let buf = &columns.columns[0];
108		(0..buf.len())
109			.map(|i| match buf.get_value(i) {
110				Value::Int4(v) => v,
111				other => panic!("expected Int4, got {other:?}"),
112			})
113			.collect()
114	}
115
116	fn make_ctx() -> QueryContext {
117		let services = Services::testing();
118		let memory = query_budget(&services);
119		QueryContext {
120			services,
121			source: None,
122			batch_size: 1024,
123			params: Params::None,
124			symbols: SymbolTable::new(),
125			identity: IdentityId::system(),
126			memory,
127		}
128	}
129
130	fn header(names: &[&str]) -> Option<ColumnHeaders> {
131		Some(ColumnHeaders {
132			columns: names.iter().map(|n| Fragment::internal(*n)).collect(),
133		})
134	}
135
136	#[test]
137	fn concatenates_two_inputs_in_order() {
138		let mut admin = create_test_admin_transaction();
139		let mut tx: Transaction<'_> = (&mut admin).into();
140		let mut ctx = make_ctx();
141
142		let h = header(&["v"]);
143		let a = StubNode::new(vec![batch("v", vec![1, 2]), batch("v", vec![3, 4])], h.clone());
144		let b = StubNode::new(vec![batch("v", vec![5, 6]), batch("v", vec![7, 8])], h);
145		let mut node = DeltaMergeNode::new(vec![Box::new(a), Box::new(b)]);
146		node.initialize(&mut tx, &ctx).unwrap();
147
148		let mut values = Vec::new();
149		while let Some(b) = node.next(&mut tx, &mut ctx).unwrap() {
150			values.extend(first_int4(&b));
151		}
152		assert_eq!(values, vec![1, 2, 3, 4, 5, 6, 7, 8]);
153	}
154
155	#[test]
156	fn empty_input_followed_by_nonempty_yields_nonempty() {
157		let mut admin = create_test_admin_transaction();
158		let mut tx: Transaction<'_> = (&mut admin).into();
159		let mut ctx = make_ctx();
160
161		let h = header(&["v"]);
162		let a = StubNode::new(vec![], h.clone());
163		let b = StubNode::new(vec![batch("v", vec![10, 20])], h);
164		let mut node = DeltaMergeNode::new(vec![Box::new(a), Box::new(b)]);
165		node.initialize(&mut tx, &ctx).unwrap();
166
167		let mut values = Vec::new();
168		while let Some(b) = node.next(&mut tx, &mut ctx).unwrap() {
169			values.extend(first_int4(&b));
170		}
171		assert_eq!(values, vec![10, 20]);
172	}
173
174	#[test]
175	fn all_empty_inputs_return_none() {
176		let mut admin = create_test_admin_transaction();
177		let mut tx: Transaction<'_> = (&mut admin).into();
178		let mut ctx = make_ctx();
179
180		let h = header(&["v"]);
181		let a = StubNode::new(vec![], h.clone());
182		let b = StubNode::new(vec![], h);
183		let mut node = DeltaMergeNode::new(vec![Box::new(a), Box::new(b)]);
184		node.initialize(&mut tx, &ctx).unwrap();
185		assert!(node.next(&mut tx, &mut ctx).unwrap().is_none());
186	}
187
188	#[test]
189	fn no_inputs_returns_none() {
190		let mut admin = create_test_admin_transaction();
191		let mut tx: Transaction<'_> = (&mut admin).into();
192		let mut ctx = make_ctx();
193
194		let mut node = DeltaMergeNode::new(vec![]);
195		node.initialize(&mut tx, &ctx).unwrap();
196		assert!(node.next(&mut tx, &mut ctx).unwrap().is_none());
197		assert!(node.headers().is_none());
198	}
199
200	#[test]
201	fn headers_match_first_input() {
202		let h0 = header(&["a", "b"]);
203		let h1 = header(&["x"]);
204		let a = StubNode::new(vec![], h0.clone());
205		let b = StubNode::new(vec![], h1);
206		let node = DeltaMergeNode::new(vec![Box::new(a), Box::new(b)]);
207		assert_eq!(node.headers().map(|h| h.columns), h0.map(|h| h.columns));
208	}
209}