reifydb_engine/vm/volcano/
merge.rs1use 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}