Skip to main content

reifydb_engine/vm/volcano/join/
common.rs

1// SPDX-License-Identifier: Apache-2.0
2// Copyright (c) 2026 ReifyDB
3
4use std::sync::Arc;
5
6use postcard::to_stdvec;
7use reifydb_core::value::column::{ColumnWithName, buffer::ColumnBuffer, columns::Columns};
8use reifydb_evaluate::expression::compile::CompiledExpr;
9use reifydb_transaction::transaction::Transaction;
10use reifydb_value::{
11	fragment::Fragment,
12	util::hash::{Hash128, xxh3_128},
13	value::Value,
14};
15
16use crate::{
17	Result,
18	vm::volcano::query::{QueryContext, QueryNode, charge_query_memory, eval_context_from_query},
19};
20
21pub(crate) fn load_and_merge_all<'a>(
22	node: &mut Box<dyn QueryNode>,
23	rx: &mut Transaction<'a>,
24	ctx: &mut QueryContext,
25) -> Result<Columns> {
26	let mut result: Option<Columns> = None;
27	let mut charged = 0usize;
28
29	while let Some(columns) = node.next(rx, ctx)? {
30		if let Some(mut acc) = result.take() {
31			acc.append_columns(columns)?;
32			result = Some(acc);
33		} else {
34			result = Some(columns);
35		}
36		if let Some(acc) = &result {
37			charge_query_memory(&ctx.memory, &mut charged, acc)?;
38		}
39	}
40	let result = result.unwrap_or_else(Columns::empty);
41	Ok(result)
42}
43
44pub struct ResolvedColumnNames {
45	pub qualified_names: Vec<String>,
46}
47
48pub fn resolve_column_names(
49	left_columns: &Columns,
50	right_columns: &Columns,
51	alias: &Option<Fragment>,
52	excluded_right_indices: Option<&[usize]>,
53) -> ResolvedColumnNames {
54	let mut qualified_names = Vec::new();
55
56	for col in left_columns.iter() {
57		qualified_names.push(col.name().text().to_string());
58	}
59
60	for (idx, col) in right_columns.iter().enumerate() {
61		if let Some(excluded) = excluded_right_indices
62			&& excluded.contains(&idx)
63		{
64			continue;
65		}
66
67		let col_name = col.name().text();
68
69		let alias_text = alias.as_ref().map(|a| a.text()).unwrap_or("other");
70		let prefixed_name = format!("{}_{}", alias_text, col_name);
71
72		let mut final_name = prefixed_name.clone();
73		if qualified_names.contains(&final_name) {
74			let mut counter = 2;
75			loop {
76				let candidate = format!("{}_{}", prefixed_name, counter);
77				if !qualified_names.contains(&candidate) {
78					final_name = candidate;
79					break;
80				}
81				counter += 1;
82			}
83		}
84
85		qualified_names.push(final_name);
86	}
87
88	ResolvedColumnNames {
89		qualified_names,
90	}
91}
92
93pub fn build_eval_columns(
94	left_columns: &Columns,
95	right_columns: &Columns,
96	left_row: &[Value],
97	right_row: &[Value],
98	alias: &Option<Fragment>,
99) -> Vec<ColumnWithName> {
100	let mut eval_columns = Vec::new();
101
102	for (idx, col) in left_columns.iter().enumerate() {
103		let data = match &left_row[idx] {
104			Value::None {
105				..
106			} => ColumnBuffer::typed_none(&col.get_type()),
107			value => ColumnBuffer::from(value.clone()),
108		};
109		eval_columns.push(ColumnWithName::new(col.name().clone(), data));
110	}
111
112	for (idx, col) in right_columns.iter().enumerate() {
113		let data = match &right_row[idx] {
114			Value::None {
115				..
116			} => ColumnBuffer::typed_none(&col.get_type()),
117			value => ColumnBuffer::from(value.clone()),
118		};
119		if let Some(alias) = alias {
120			let aliased_name = Fragment::internal(format!("{}.{}", alias.text(), col.name().text()));
121			eval_columns.push(ColumnWithName {
122				name: aliased_name,
123				data,
124			});
125		} else {
126			eval_columns.push(ColumnWithName::new(col.name().clone(), data));
127		}
128	}
129
130	eval_columns
131}
132
133pub struct JoinContext {
134	pub context: Option<Arc<QueryContext>>,
135	pub compiled: Vec<CompiledExpr>,
136}
137
138impl Default for JoinContext {
139	fn default() -> Self {
140		Self::new()
141	}
142}
143
144impl JoinContext {
145	pub fn new() -> Self {
146		Self {
147			context: None,
148			compiled: vec![],
149		}
150	}
151
152	pub fn set(&mut self, ctx: &QueryContext) {
153		self.context = Some(Arc::new(ctx.clone()));
154	}
155
156	pub fn get(&self) -> &Arc<QueryContext> {
157		self.context.as_ref().expect("Join context not initialized")
158	}
159
160	pub fn is_initialized(&self) -> bool {
161		self.context.is_some()
162	}
163}
164
165pub(crate) fn compute_join_hash(
166	columns: &Columns,
167	col_indices: &[usize],
168	row_idx: usize,
169	buf: &mut Vec<u8>,
170) -> Option<Hash128> {
171	buf.clear();
172	for &idx in col_indices {
173		let value = columns[idx].get_value(row_idx);
174		if matches!(value, Value::None { .. }) {
175			return None;
176		}
177		let bytes = to_stdvec(&value).ok()?;
178		buf.extend_from_slice(&bytes);
179	}
180	Some(xxh3_128(buf))
181}
182
183pub(crate) fn keys_equal_by_index(
184	left: &Columns,
185	left_row: usize,
186	left_indices: &[usize],
187	right: &Columns,
188	right_row: usize,
189	right_indices: &[usize],
190) -> bool {
191	for (&li, &ri) in left_indices.iter().zip(right_indices.iter()) {
192		let lv = left[li].get_value(left_row);
193		let rv = right[ri].get_value(right_row);
194		if lv != rv {
195			return false;
196		}
197	}
198	true
199}
200
201pub(crate) fn eval_join_condition(
202	compiled: &[CompiledExpr],
203	left_columns: &Columns,
204	right_columns: &Columns,
205	left_row: &[Value],
206	right_row: &[Value],
207	alias: &Option<Fragment>,
208	ctx: &QueryContext,
209) -> bool {
210	if compiled.is_empty() {
211		return true;
212	}
213	let eval_columns = build_eval_columns(left_columns, right_columns, left_row, right_row, alias);
214	let session = eval_context_from_query(ctx);
215	let exec_ctx = session.with_eval_join(Columns::new(eval_columns));
216	compiled.iter().all(|compiled_expr| {
217		let col = compiled_expr.execute(&exec_ctx).unwrap();
218		matches!(col.data().get_value(0), Value::Boolean(true))
219	})
220}