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_value::{
58		fragment::Fragment,
59		params::Params,
60		value::{Value, identity::IdentityId},
61	};
62
63	use super::*;
64	use crate::{
65		test_harness::create_test_admin_transaction,
66		vm::{services::Services, stack::SymbolTable},
67	};
68
69	struct StubNode {
70		batches: VecDeque<Columns>,
71		headers: Option<ColumnHeaders>,
72		init_count: usize,
73	}
74
75	impl StubNode {
76		fn new(batches: Vec<Columns>, headers: Option<ColumnHeaders>) -> Self {
77			Self {
78				batches: batches.into(),
79				headers,
80				init_count: 0,
81			}
82		}
83	}
84
85	impl QueryNode for StubNode {
86		fn initialize<'a>(&mut self, _rx: &mut Transaction<'a>, _ctx: &QueryContext) -> Result<()> {
87			self.init_count += 1;
88			Ok(())
89		}
90
91		fn next<'a>(&mut self, _rx: &mut Transaction<'a>, _ctx: &mut QueryContext) -> Result<Option<Columns>> {
92			Ok(self.batches.pop_front())
93		}
94
95		fn headers(&self) -> Option<ColumnHeaders> {
96			self.headers.clone()
97		}
98	}
99
100	fn batch(name: &str, vals: Vec<i32>) -> Columns {
101		Columns::new(vec![ColumnWithName {
102			name: Fragment::internal(name),
103			data: ColumnBuffer::int4(vals),
104		}])
105	}
106
107	fn first_int4(columns: &Columns) -> Vec<i32> {
108		let buf = &columns.columns[0];
109		(0..buf.len())
110			.map(|i| match buf.get_value(i) {
111				Value::Int4(v) => v,
112				other => panic!("expected Int4, got {other:?}"),
113			})
114			.collect()
115	}
116
117	fn make_ctx() -> QueryContext {
118		QueryContext {
119			services: Services::testing(),
120			source: None,
121			batch_size: 1024,
122			params: Params::None,
123			symbols: SymbolTable::new(),
124			identity: IdentityId::system(),
125		}
126	}
127
128	fn header(names: &[&str]) -> Option<ColumnHeaders> {
129		Some(ColumnHeaders {
130			columns: names.iter().map(|n| Fragment::internal(*n)).collect(),
131		})
132	}
133
134	#[test]
135	fn concatenates_two_inputs_in_order() {
136		let mut admin = create_test_admin_transaction();
137		let mut tx: Transaction<'_> = (&mut admin).into();
138		let mut ctx = make_ctx();
139
140		let h = header(&["v"]);
141		let a = StubNode::new(vec![batch("v", vec![1, 2]), batch("v", vec![3, 4])], h.clone());
142		let b = StubNode::new(vec![batch("v", vec![5, 6]), batch("v", vec![7, 8])], h);
143		let mut node = DeltaMergeNode::new(vec![Box::new(a), Box::new(b)]);
144		node.initialize(&mut tx, &ctx).unwrap();
145
146		let mut values = Vec::new();
147		while let Some(b) = node.next(&mut tx, &mut ctx).unwrap() {
148			values.extend(first_int4(&b));
149		}
150		assert_eq!(values, vec![1, 2, 3, 4, 5, 6, 7, 8]);
151	}
152
153	#[test]
154	fn empty_input_followed_by_nonempty_yields_nonempty() {
155		let mut admin = create_test_admin_transaction();
156		let mut tx: Transaction<'_> = (&mut admin).into();
157		let mut ctx = make_ctx();
158
159		let h = header(&["v"]);
160		let a = StubNode::new(vec![], h.clone());
161		let b = StubNode::new(vec![batch("v", vec![10, 20])], h);
162		let mut node = DeltaMergeNode::new(vec![Box::new(a), Box::new(b)]);
163		node.initialize(&mut tx, &ctx).unwrap();
164
165		let mut values = Vec::new();
166		while let Some(b) = node.next(&mut tx, &mut ctx).unwrap() {
167			values.extend(first_int4(&b));
168		}
169		assert_eq!(values, vec![10, 20]);
170	}
171
172	#[test]
173	fn all_empty_inputs_return_none() {
174		let mut admin = create_test_admin_transaction();
175		let mut tx: Transaction<'_> = (&mut admin).into();
176		let mut ctx = make_ctx();
177
178		let h = header(&["v"]);
179		let a = StubNode::new(vec![], h.clone());
180		let b = StubNode::new(vec![], h);
181		let mut node = DeltaMergeNode::new(vec![Box::new(a), Box::new(b)]);
182		node.initialize(&mut tx, &ctx).unwrap();
183		assert!(node.next(&mut tx, &mut ctx).unwrap().is_none());
184	}
185
186	#[test]
187	fn no_inputs_returns_none() {
188		let mut admin = create_test_admin_transaction();
189		let mut tx: Transaction<'_> = (&mut admin).into();
190		let mut ctx = make_ctx();
191
192		let mut node = DeltaMergeNode::new(vec![]);
193		node.initialize(&mut tx, &ctx).unwrap();
194		assert!(node.next(&mut tx, &mut ctx).unwrap().is_none());
195		assert!(node.headers().is_none());
196	}
197
198	#[test]
199	fn headers_match_first_input() {
200		let h0 = header(&["a", "b"]);
201		let h1 = header(&["x"]);
202		let a = StubNode::new(vec![], h0.clone());
203		let b = StubNode::new(vec![], h1);
204		let node = DeltaMergeNode::new(vec![Box::new(a), Box::new(b)]);
205		assert_eq!(node.headers().map(|h| h.columns), h0.map(|h| h.columns));
206	}
207}