Skip to main content

radixdb_executor/operators/
merge_join.rs

1// Copyright 2026 RadixDB Contributors
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15//! Merge Join Operator for pre-sorted inputs.
16//!
17//! This operator implements merge join with O(N+M) complexity when both
18//! inputs are already sorted on the join keys. It's optimal for:
19//! - Joining tables that are physically sorted (e.g., clustered index)
20//! - Joining results of ORDER BY queries
21//! - Self-joins on sorted columns
22
23use std::cmp::Ordering;
24
25use crate::context::check_current_query_cancelled;
26use crate::operator::{ColumnInfo, Operator, OrderingProperty, RowRef};
27use crate::utils::compare_values;
28use radixdb_core::value::NULL_VALUE;
29use radixdb_core::{Result, Row, Value};
30
31use super::hash_join::JoinType;
32
33fn compare_rows_on_keys(
34    left: &Row,
35    right: &Row,
36    left_key_indices: &[usize],
37    right_key_indices: &[usize],
38) -> Ordering {
39    for (li, ri) in left_key_indices.iter().zip(right_key_indices.iter()) {
40        let lv = left.get(*li).cloned().unwrap_or(NULL_VALUE);
41        let rv = right.get(*ri).cloned().unwrap_or(NULL_VALUE);
42
43        // SQL equality join keys containing NULL never match. Returning a
44        // stable ordering still lets the merge cursors make progress.
45        if lv.is_null() && rv.is_null() {
46            return Ordering::Less;
47        }
48        if lv.is_null() {
49            return Ordering::Greater;
50        }
51        if rv.is_null() {
52            return Ordering::Less;
53        }
54
55        let cmp = compare_values(&lv, &rv);
56        if cmp != Ordering::Equal {
57            return cmp;
58        }
59    }
60    Ordering::Equal
61}
62
63fn compare_same_side_keys(row1: &Row, row2: &Row, key_indices: &[usize]) -> Ordering {
64    for &idx in key_indices {
65        let v1 = row1.get(idx).cloned().unwrap_or(NULL_VALUE);
66        let v2 = row2.get(idx).cloned().unwrap_or(NULL_VALUE);
67
68        if v1.is_null() && v2.is_null() {
69            continue;
70        }
71        if v1.is_null() {
72            return Ordering::Greater;
73        }
74        if v2.is_null() {
75            return Ordering::Less;
76        }
77
78        let cmp = compare_values(&v1, &v2);
79        if cmp != Ordering::Equal {
80            return cmp;
81        }
82    }
83    Ordering::Equal
84}
85
86/// Conservative upper bound for the two duplicate-group Vec allocations.
87/// Row payloads are moved from the already materialized inputs, not cloned;
88/// the blocking state adds only Row slots. The factor of two covers geometric
89/// Vec capacity growth.
90fn merge_group_slot_bytes(rows: usize) -> Option<usize> {
91    rows.checked_mul(std::mem::size_of::<Row>())?.checked_mul(2)
92}
93
94/// Merge Join Operator for pre-sorted inputs.
95///
96/// Both inputs must be sorted on their respective join keys.
97/// The operator performs a single pass through both inputs,
98/// producing matches as they're found.
99pub struct MergeJoinOperator {
100    // Input operators
101    left: Box<dyn Operator>,
102    right: Box<dyn Operator>,
103
104    // Join configuration
105    join_type: JoinType,
106    left_key_indices: Vec<usize>,
107    right_key_indices: Vec<usize>,
108
109    // Output schema
110    schema: Vec<ColumnInfo>,
111    left_col_count: usize,
112    right_col_count: usize,
113
114    // Ordered cursors. Only the current row and one duplicate-key group are
115    // retained; complete inputs are never materialized by open().
116    left_current: Option<Row>,
117    right_current: Option<Row>,
118    left_group: Vec<Row>,
119    right_group: Vec<Row>,
120    group_left_idx: usize,
121    group_right_idx: usize,
122    group_retained_rows: usize,
123    max_group_rows: usize,
124    max_group_bytes: usize,
125
126    // Cached null rows for OUTER joins (avoid repeated allocation)
127    cached_null_left: Vec<Value>,
128    cached_null_right: Vec<Value>,
129
130    // Physical property contract. Merge is invalid without both certificates.
131    input_ordering_certified: bool,
132    output_ordering: OrderingProperty,
133
134    // State
135    opened: bool,
136}
137
138impl MergeJoinOperator {
139    /// Create a new merge join operator.
140    ///
141    /// # Arguments
142    /// * `left` - Left input operator (must be sorted on left_key_indices)
143    /// * `right` - Right input operator (must be sorted on right_key_indices)
144    /// * `join_type` - Type of join (INNER, LEFT, RIGHT, FULL - NOT Cross)
145    /// * `left_key_indices` - Column indices for left join keys
146    /// * `right_key_indices` - Column indices for right join keys
147    ///
148    /// # Panics
149    /// Debug builds will panic if `join_type` is `JoinType::Cross`.
150    /// Cross joins should use NestedLoopJoinOperator instead.
151    pub fn new(
152        left: Box<dyn Operator>,
153        right: Box<dyn Operator>,
154        join_type: JoinType,
155        left_key_indices: Vec<usize>,
156        right_key_indices: Vec<usize>,
157    ) -> Self {
158        // Cross joins should use NestedLoop, not MergeJoin (no keys to merge on)
159        debug_assert!(
160            !matches!(join_type, JoinType::Cross),
161            "MergeJoin cannot be used for CROSS JOIN - use NestedLoopJoin instead"
162        );
163
164        let left_ordering = left.ordering();
165        let right_ordering = right.ordering();
166        let input_ordering_certified = left_ordering.proves_ascending_nulls_last(&left_key_indices)
167            && right_ordering.proves_ascending_nulls_last(&right_key_indices);
168        let output_ordering =
169            if input_ordering_certified && matches!(join_type, JoinType::Inner | JoinType::Left) {
170                left_ordering
171            } else {
172                OrderingProperty::Unknown
173            };
174
175        // Build combined schema
176        let mut schema = Vec::new();
177        schema.extend(left.schema().iter().cloned());
178        schema.extend(right.schema().iter().cloned());
179
180        let left_col_count = left.schema().len();
181        let right_col_count = right.schema().len();
182
183        Self {
184            left,
185            right,
186            join_type,
187            left_key_indices,
188            right_key_indices,
189            schema,
190            left_col_count,
191            right_col_count,
192            left_current: None,
193            right_current: None,
194            left_group: Vec::new(),
195            right_group: Vec::new(),
196            group_left_idx: 0,
197            group_right_idx: 0,
198            group_retained_rows: 0,
199            max_group_rows: crate::utils::RetainedRowsBudget::DEFAULT_MAX_ROWS,
200            max_group_bytes: crate::utils::RetainedRowsBudget::DEFAULT_MAX_BYTES,
201            cached_null_left: Vec::new(),  // Initialized in open()
202            cached_null_right: Vec::new(), // Initialized in open()
203            input_ordering_certified,
204            output_ordering,
205            opened: false,
206        }
207    }
208
209    /// Bound the only blocking state retained by streaming merge: the two
210    /// matching duplicate-key groups. Production planning preflights the same
211    /// limits and falls back to hash before execution when they cannot fit.
212    pub(crate) fn with_group_budget(mut self, max_rows: usize, max_bytes: usize) -> Self {
213        self.max_group_rows = max_rows;
214        self.max_group_bytes = max_bytes;
215        self
216    }
217
218    /// Prove that every *matching* duplicate-key group fits the merge owner.
219    /// Small relations are admitted from cardinality alone in O(1). Only an
220    /// oversized merge candidate pays the linear key-only safety scan.
221    pub(crate) fn matching_groups_retained_bytes(
222        left_rows: &[Row],
223        right_rows: &[Row],
224        left_key_indices: &[usize],
225        right_key_indices: &[usize],
226        max_rows: usize,
227        max_bytes: usize,
228    ) -> Result<Option<usize>> {
229        if left_key_indices.is_empty() || left_key_indices.len() != right_key_indices.len() {
230            return Ok(None);
231        }
232
233        let total_rows = left_rows.len().checked_add(right_rows.len());
234        if let Some(total_rows) = total_rows {
235            if total_rows <= max_rows {
236                if let Some(bytes) = merge_group_slot_bytes(total_rows) {
237                    if bytes <= max_bytes {
238                        return Ok(Some(bytes));
239                    }
240                }
241            }
242        }
243
244        let mut left_idx = 0;
245        let mut right_idx = 0;
246        let mut peak_group_bytes = 0;
247        while left_idx < left_rows.len() && right_idx < right_rows.len() {
248            if (left_idx.wrapping_add(right_idx)) & 0xff == 0 {
249                check_current_query_cancelled()?;
250            }
251            match compare_rows_on_keys(
252                &left_rows[left_idx],
253                &right_rows[right_idx],
254                left_key_indices,
255                right_key_indices,
256            ) {
257                Ordering::Less => left_idx += 1,
258                Ordering::Greater => right_idx += 1,
259                Ordering::Equal => {
260                    let left_start = left_idx;
261                    let right_start = right_idx;
262                    left_idx += 1;
263                    while left_idx < left_rows.len()
264                        && compare_same_side_keys(
265                            &left_rows[left_start],
266                            &left_rows[left_idx],
267                            left_key_indices,
268                        ) == Ordering::Equal
269                    {
270                        left_idx += 1;
271                    }
272                    right_idx += 1;
273                    while right_idx < right_rows.len()
274                        && compare_same_side_keys(
275                            &right_rows[right_start],
276                            &right_rows[right_idx],
277                            right_key_indices,
278                        ) == Ordering::Equal
279                    {
280                        right_idx += 1;
281                    }
282
283                    let group_rows =
284                        (left_idx - left_start).saturating_add(right_idx - right_start);
285                    if group_rows > max_rows
286                        || merge_group_slot_bytes(group_rows).is_none_or(|bytes| bytes > max_bytes)
287                    {
288                        return Ok(None);
289                    }
290                    peak_group_bytes = peak_group_bytes.max(
291                        merge_group_slot_bytes(group_rows)
292                            .expect("admitted merge group must have addressable slot bytes"),
293                    );
294                }
295            }
296        }
297        Ok(Some(peak_group_bytes))
298    }
299
300    #[cfg(test)]
301    pub(crate) fn matching_groups_fit_budget(
302        left_rows: &[Row],
303        right_rows: &[Row],
304        left_key_indices: &[usize],
305        right_key_indices: &[usize],
306        max_rows: usize,
307        max_bytes: usize,
308    ) -> Result<bool> {
309        Self::matching_groups_retained_bytes(
310            left_rows,
311            right_rows,
312            left_key_indices,
313            right_key_indices,
314            max_rows,
315            max_bytes,
316        )
317        .map(|bytes| bytes.is_some())
318    }
319
320    /// Compare two rows on their respective join keys.
321    fn compare_on_keys(&self, left: &Row, right: &Row) -> Ordering {
322        compare_rows_on_keys(left, right, &self.left_key_indices, &self.right_key_indices)
323    }
324
325    /// Compare two rows from the same side on their keys.
326    fn compare_same_side(&self, row1: &Row, row2: &Row, key_indices: &[usize]) -> Ordering {
327        compare_same_side_keys(row1, row2, key_indices)
328    }
329
330    fn admit_group_row(&mut self) -> Result<()> {
331        let next_rows = self.group_retained_rows.saturating_add(1);
332        let next_bytes = merge_group_slot_bytes(next_rows).unwrap_or(usize::MAX);
333        if next_rows > self.max_group_rows || next_bytes > self.max_group_bytes {
334            return Err(radixdb_core::Error::invalid_argument(format!(
335                "merge join duplicate group exceeds memory budget (rows {next_rows}/{}, slot bytes {next_bytes}/{})",
336                self.max_group_rows, self.max_group_bytes
337            )));
338        }
339        self.group_retained_rows = next_rows;
340        Ok(())
341    }
342
343    /// Create a NULL row for the left side (uses cached values).
344    #[inline]
345    fn null_left_row(&self) -> Row {
346        Row::from_values(self.cached_null_left.clone())
347    }
348
349    /// Create a NULL row for the right side (uses cached values).
350    #[inline]
351    fn null_right_row(&self) -> Row {
352        Row::from_values(self.cached_null_right.clone())
353    }
354
355    /// Combine left and right rows into output row.
356    #[inline]
357    fn combine(&self, left: &Row, right: &Row) -> Row {
358        Row::from_combined(left, right)
359    }
360
361    fn pull_left(&mut self) -> Result<Option<Row>> {
362        self.left.next().map(|row| row.map(RowRef::into_owned))
363    }
364
365    fn pull_right(&mut self) -> Result<Option<Row>> {
366        self.right.next().map(|row| row.map(RowRef::into_owned))
367    }
368
369    /// Consume the complete duplicate-key groups starting at both current
370    /// cursors. The first non-matching row remains as the next cursor value.
371    fn collect_equal_groups(&mut self) -> Result<()> {
372        self.left_group.clear();
373        self.right_group.clear();
374        self.group_left_idx = 0;
375        self.group_right_idx = 0;
376        self.group_retained_rows = 0;
377
378        self.admit_group_row()?;
379        self.left_group.push(
380            self.left_current
381                .take()
382                .expect("equal merge key without left cursor"),
383        );
384        self.admit_group_row()?;
385        self.right_group.push(
386            self.right_current
387                .take()
388                .expect("equal merge key without right cursor"),
389        );
390
391        loop {
392            if self.left_group.len() & 0xff == 0 {
393                check_current_query_cancelled()?;
394            }
395            match self.pull_left()? {
396                Some(row)
397                    if self.compare_same_side(
398                        &row,
399                        &self.left_group[0],
400                        &self.left_key_indices,
401                    ) == Ordering::Equal =>
402                {
403                    self.admit_group_row()?;
404                    self.left_group.push(row);
405                }
406                next => {
407                    self.left_current = next;
408                    break;
409                }
410            }
411        }
412
413        loop {
414            if self.right_group.len() & 0xff == 0 {
415                check_current_query_cancelled()?;
416            }
417            match self.pull_right()? {
418                Some(row)
419                    if self.compare_same_side(
420                        &row,
421                        &self.right_group[0],
422                        &self.right_key_indices,
423                    ) == Ordering::Equal =>
424                {
425                    self.admit_group_row()?;
426                    self.right_group.push(row);
427                }
428                next => {
429                    self.right_current = next;
430                    break;
431                }
432            }
433        }
434        Ok(())
435    }
436
437    fn emit_group_match(&mut self) -> Option<RowRef> {
438        if self.group_left_idx >= self.left_group.len()
439            || self.group_right_idx >= self.right_group.len()
440        {
441            return None;
442        }
443
444        let row = self.combine(
445            &self.left_group[self.group_left_idx],
446            &self.right_group[self.group_right_idx],
447        );
448        self.group_right_idx += 1;
449        if self.group_right_idx == self.right_group.len() {
450            self.group_right_idx = 0;
451            self.group_left_idx += 1;
452            if self.group_left_idx == self.left_group.len() {
453                self.left_group.clear();
454                self.right_group.clear();
455                self.group_left_idx = 0;
456                self.group_retained_rows = 0;
457            }
458        }
459        Some(RowRef::Owned(row))
460    }
461}
462
463impl Operator for MergeJoinOperator {
464    fn open(&mut self) -> Result<()> {
465        if !self.input_ordering_certified {
466            return Err(radixdb_core::Error::invalid_argument(
467                "merge join requires explicit ascending NULLS LAST ordering certificates",
468            ));
469        }
470        if matches!(
471            self.join_type,
472            JoinType::Cross | JoinType::Semi | JoinType::Anti
473        ) {
474            return Err(radixdb_core::Error::invalid_argument(
475                "merge join supports only INNER, LEFT, RIGHT and FULL joins",
476            ));
477        }
478
479        if let Err(error) = self.left.open() {
480            let _ = self.left.close();
481            return Err(error);
482        }
483        if let Err(error) = self.right.open() {
484            let _ = self.right.close();
485            let _ = self.left.close();
486            return Err(error);
487        }
488
489        let open_result = (|| {
490            check_current_query_cancelled()?;
491
492            // Pre-cache null rows for OUTER joins (avoids repeated allocation)
493            // NULL_VALUE is a static constant, cloning Vec is just memcpy
494            if matches!(
495                self.join_type,
496                JoinType::Left | JoinType::Right | JoinType::Full
497            ) {
498                self.cached_null_left = vec![NULL_VALUE; self.left_col_count];
499                self.cached_null_right = vec![NULL_VALUE; self.right_col_count];
500            }
501
502            self.left_group.clear();
503            self.right_group.clear();
504            self.group_left_idx = 0;
505            self.group_right_idx = 0;
506            self.group_retained_rows = 0;
507            self.left_current = self.pull_left()?;
508            self.right_current = self.pull_right()?;
509
510            self.opened = true;
511            Ok(())
512        })();
513        if let Err(error) = open_result {
514            let _ = self.close();
515            return Err(error);
516        }
517        Ok(())
518    }
519
520    fn next(&mut self) -> Result<Option<RowRef>> {
521        check_current_query_cancelled()?;
522        if !self.opened {
523            return Err(radixdb_core::Error::internal(
524                "MergeJoinOperator::next called before open",
525            ));
526        }
527
528        let is_left_outer = matches!(self.join_type, JoinType::Left | JoinType::Full);
529        let is_right_outer = matches!(self.join_type, JoinType::Right | JoinType::Full);
530
531        if let Some(row) = self.emit_group_match() {
532            return Ok(Some(row));
533        }
534
535        loop {
536            check_current_query_cancelled()?;
537
538            match (&self.left_current, &self.right_current) {
539                (None, None) => return Ok(None),
540                (None, Some(_)) => {
541                    if !is_right_outer {
542                        return Ok(None);
543                    }
544                    let right = self
545                        .right_current
546                        .take()
547                        .expect("right cursor disappeared during outer emission");
548                    self.right_current = self.pull_right()?;
549                    let null_left = self.null_left_row();
550                    return Ok(Some(RowRef::Owned(self.combine(&null_left, &right))));
551                }
552                (Some(_), None) => {
553                    if !is_left_outer {
554                        return Ok(None);
555                    }
556                    let left = self
557                        .left_current
558                        .take()
559                        .expect("left cursor disappeared during outer emission");
560                    self.left_current = self.pull_left()?;
561                    let null_right = self.null_right_row();
562                    return Ok(Some(RowRef::Owned(self.combine(&left, &null_right))));
563                }
564                (Some(_), Some(_)) => {}
565            }
566
567            let ordering = self.compare_on_keys(
568                self.left_current.as_ref().expect("left cursor missing"),
569                self.right_current.as_ref().expect("right cursor missing"),
570            );
571            match ordering {
572                Ordering::Less => {
573                    let left = self.left_current.take().expect("left cursor missing");
574                    self.left_current = self.pull_left()?;
575                    if is_left_outer {
576                        let null_right = self.null_right_row();
577                        return Ok(Some(RowRef::Owned(self.combine(&left, &null_right))));
578                    }
579                }
580                Ordering::Greater => {
581                    let right = self.right_current.take().expect("right cursor missing");
582                    self.right_current = self.pull_right()?;
583                    if is_right_outer {
584                        let null_left = self.null_left_row();
585                        return Ok(Some(RowRef::Owned(self.combine(&null_left, &right))));
586                    }
587                }
588                Ordering::Equal => {
589                    self.collect_equal_groups()?;
590                    return Ok(self.emit_group_match());
591                }
592            }
593        }
594    }
595
596    fn close(&mut self) -> Result<()> {
597        self.left_current = None;
598        self.right_current = None;
599        self.left_group.clear();
600        self.right_group.clear();
601        self.group_retained_rows = 0;
602        self.opened = false;
603        let left = self.left.close();
604        let right = self.right.close();
605        left.and(right)
606    }
607
608    fn schema(&self) -> &[ColumnInfo] {
609        &self.schema
610    }
611
612    fn estimated_rows(&self) -> Option<usize> {
613        let left_est = self.left.estimated_rows()?;
614        let right_est = self.right.estimated_rows()?;
615
616        Some(match self.join_type {
617            JoinType::Inner => left_est.min(right_est),
618            JoinType::Left => left_est,
619            JoinType::Right => right_est,
620            JoinType::Full => left_est + right_est,
621            JoinType::Cross => left_est * right_est,
622            JoinType::Semi => left_est.min(right_est),
623            JoinType::Anti => left_est,
624        })
625    }
626
627    fn ordering(&self) -> OrderingProperty {
628        self.output_ordering.clone()
629    }
630
631    fn name(&self) -> &str {
632        match self.join_type {
633            JoinType::Inner => "MergeJoin (INNER)",
634            JoinType::Left => "MergeJoin (LEFT)",
635            JoinType::Right => "MergeJoin (RIGHT)",
636            JoinType::Full => "MergeJoin (FULL)",
637            JoinType::Cross => "MergeJoin (CROSS)",
638            JoinType::Semi => "MergeJoin (SEMI)",
639            JoinType::Anti => "MergeJoin (ANTI)",
640        }
641    }
642}
643
644#[cfg(test)]
645mod tests {
646    use super::*;
647    use crate::operator::MaterializedOperator;
648    use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering};
649    use std::sync::Arc;
650
651    struct CountingOrderedOperator {
652        rows: Vec<Option<Row>>,
653        schema: Vec<ColumnInfo>,
654        next_index: usize,
655        next_calls: Arc<AtomicUsize>,
656    }
657
658    impl CountingOrderedOperator {
659        fn new(rows: Vec<Row>, next_calls: Arc<AtomicUsize>) -> Self {
660            Self {
661                rows: rows.into_iter().map(Some).collect(),
662                schema: vec![ColumnInfo::new("id")],
663                next_index: 0,
664                next_calls,
665            }
666        }
667    }
668
669    impl Operator for CountingOrderedOperator {
670        fn open(&mut self) -> Result<()> {
671            self.next_index = 0;
672            Ok(())
673        }
674
675        fn next(&mut self) -> Result<Option<RowRef>> {
676            self.next_calls.fetch_add(1, AtomicOrdering::Relaxed);
677            let Some(row) = self.rows.get_mut(self.next_index) else {
678                return Ok(None);
679            };
680            self.next_index += 1;
681            Ok(row.take().map(RowRef::Owned))
682        }
683
684        fn close(&mut self) -> Result<()> {
685            Ok(())
686        }
687
688        fn schema(&self) -> &[ColumnInfo] {
689            &self.schema
690        }
691
692        fn ordering(&self) -> OrderingProperty {
693            OrderingProperty::ascending_nulls_last(vec![0])
694        }
695
696        fn name(&self) -> &str {
697            "CountingOrdered"
698        }
699    }
700
701    fn make_rows(data: Vec<Vec<i64>>) -> Vec<Row> {
702        data.into_iter()
703            .map(|vals| Row::from_values(vals.into_iter().map(Value::integer).collect()))
704            .collect()
705    }
706
707    fn make_operator(data: Vec<Vec<i64>>, cols: Vec<&str>) -> Box<dyn Operator> {
708        let rows = make_rows(data);
709        let schema = cols.into_iter().map(ColumnInfo::new).collect();
710        Box::new(
711            MaterializedOperator::new(rows, schema)
712                .with_ordering(OrderingProperty::ascending_nulls_last(vec![0])),
713        )
714    }
715
716    fn make_value_operator(rows: Vec<Row>, cols: Vec<&str>) -> Box<dyn Operator> {
717        let schema = cols.into_iter().map(ColumnInfo::new).collect();
718        Box::new(
719            MaterializedOperator::new(rows, schema)
720                .with_ordering(OrderingProperty::ascending_nulls_last(vec![0])),
721        )
722    }
723
724    fn collect_results(op: &mut dyn Operator) -> Result<Vec<Row>> {
725        let mut results = Vec::new();
726        op.open()?;
727        while let Some(row_ref) = op.next()? {
728            results.push(row_ref.into_owned());
729        }
730        op.close()?;
731        Ok(results)
732    }
733
734    #[test]
735    fn test_inner_merge_join() {
736        // Both sides sorted on id
737        let left = make_operator(
738            vec![vec![1, 10], vec![2, 20], vec![3, 30]],
739            vec!["id", "value"],
740        );
741        let right = make_operator(vec![vec![1, 100], vec![3, 300]], vec!["id", "data"]);
742
743        let mut join = MergeJoinOperator::new(
744            left,
745            right,
746            JoinType::Inner,
747            vec![0], // left key: id
748            vec![0], // right key: id
749        );
750        assert_eq!(
751            join.ordering(),
752            OrderingProperty::ascending_nulls_last(vec![0])
753        );
754
755        let results = collect_results(&mut join).unwrap();
756
757        // Should have 2 matches: id=1 and id=3
758        assert_eq!(results.len(), 2);
759    }
760
761    #[test]
762    fn open_prefetches_only_one_row_per_ordered_input() {
763        let left_calls = Arc::new(AtomicUsize::new(0));
764        let right_calls = Arc::new(AtomicUsize::new(0));
765        let rows = (0..10_000)
766            .map(|value| Row::from_values(vec![Value::Integer(value)]))
767            .collect::<Vec<_>>();
768        let left = Box::new(CountingOrderedOperator::new(
769            rows.clone(),
770            Arc::clone(&left_calls),
771        ));
772        let right = Box::new(CountingOrderedOperator::new(rows, Arc::clone(&right_calls)));
773        let mut join = MergeJoinOperator::new(left, right, JoinType::Inner, vec![0], vec![0]);
774
775        join.open().unwrap();
776        assert_eq!(left_calls.load(AtomicOrdering::Relaxed), 1);
777        assert_eq!(right_calls.load(AtomicOrdering::Relaxed), 1);
778        assert!(join.next().unwrap().is_some());
779        assert_eq!(left_calls.load(AtomicOrdering::Relaxed), 2);
780        assert_eq!(right_calls.load(AtomicOrdering::Relaxed), 2);
781        join.close().unwrap();
782    }
783
784    #[test]
785    fn merge_join_rejects_inputs_without_physical_ordering_certificate() {
786        let schema = vec![ColumnInfo::new("id")];
787        let left = Box::new(MaterializedOperator::new(
788            make_rows(vec![vec![1], vec![2]]),
789            schema.clone(),
790        ));
791        let right = Box::new(MaterializedOperator::new(
792            make_rows(vec![vec![1], vec![2]]),
793            schema,
794        ));
795        let mut join = MergeJoinOperator::new(left, right, JoinType::Inner, vec![0], vec![0]);
796
797        let error = join.open().unwrap_err();
798        assert!(error.to_string().contains("ordering certificates"));
799    }
800
801    #[test]
802    fn test_left_merge_join() {
803        let left = make_operator(
804            vec![vec![1, 10], vec![2, 20], vec![3, 30]],
805            vec!["id", "value"],
806        );
807        let right = make_operator(vec![vec![1, 100]], vec!["id", "data"]);
808
809        let mut join = MergeJoinOperator::new(left, right, JoinType::Left, vec![0], vec![0]);
810
811        let results = collect_results(&mut join).unwrap();
812
813        // All 3 left rows should be preserved
814        assert_eq!(results.len(), 3);
815
816        // Check that id=2 and id=3 have NULLs on right side
817        let row2 = results
818            .iter()
819            .find(|r| r.get(0) == Some(&Value::integer(2)))
820            .unwrap();
821        assert!(row2.get(2).unwrap().is_null());
822    }
823
824    #[test]
825    fn streaming_right_and_full_merge_emit_unmatched_rows_once() {
826        let make_left = || make_operator(vec![vec![1, 10], vec![3, 30]], vec!["id", "value"]);
827        let make_right = || {
828            make_operator(
829                vec![vec![2, 200], vec![3, 300], vec![4, 400]],
830                vec!["id", "data"],
831            )
832        };
833
834        let mut right =
835            MergeJoinOperator::new(make_left(), make_right(), JoinType::Right, vec![0], vec![0]);
836        let right_rows = collect_results(&mut right).unwrap();
837        assert_eq!(right_rows.len(), 3);
838        assert_eq!(
839            right_rows
840                .iter()
841                .filter(|row| row.get(0).is_some_and(Value::is_null))
842                .count(),
843            2
844        );
845
846        let mut full =
847            MergeJoinOperator::new(make_left(), make_right(), JoinType::Full, vec![0], vec![0]);
848        let full_rows = collect_results(&mut full).unwrap();
849        assert_eq!(full_rows.len(), 4);
850        assert_eq!(
851            full_rows
852                .iter()
853                .filter(|row| row.get(0).is_some_and(Value::is_null))
854                .count(),
855            2
856        );
857        assert_eq!(
858            full_rows
859                .iter()
860                .filter(|row| row.get(2).is_some_and(Value::is_null))
861                .count(),
862            1
863        );
864    }
865
866    #[test]
867    fn streaming_merge_never_matches_null_key_groups() {
868        let left = make_value_operator(
869            vec![
870                Row::from_values(vec![Value::Integer(1)]),
871                Row::from_values(vec![Value::Null(radixdb_core::DataType::Integer)]),
872            ],
873            vec!["key"],
874        );
875        let right = make_value_operator(
876            vec![
877                Row::from_values(vec![Value::Integer(1)]),
878                Row::from_values(vec![Value::Null(radixdb_core::DataType::Integer)]),
879            ],
880            vec!["key"],
881        );
882        let mut join = MergeJoinOperator::new(left, right, JoinType::Inner, vec![0], vec![0]);
883
884        let rows = collect_results(&mut join).unwrap();
885        assert_eq!(rows.len(), 1);
886        assert_eq!(rows[0].get(0), Some(&Value::Integer(1)));
887    }
888
889    #[test]
890    fn test_merge_join_with_duplicates() {
891        // Both sides have duplicate keys
892        let left = make_operator(
893            vec![vec![1, 10], vec![1, 11], vec![2, 20]],
894            vec!["id", "value"],
895        );
896        let right = make_operator(
897            vec![vec![1, 100], vec![1, 101], vec![2, 200]],
898            vec!["id", "data"],
899        );
900
901        let mut join = MergeJoinOperator::new(left, right, JoinType::Inner, vec![0], vec![0]);
902
903        let results = collect_results(&mut join).unwrap();
904
905        // id=1: 2 left x 2 right = 4 matches
906        // id=2: 1 left x 1 right = 1 match
907        // Total = 5
908        assert_eq!(results.len(), 5);
909    }
910
911    #[test]
912    fn matching_group_preflight_rejects_only_oversized_matching_groups() {
913        let matching_left = make_rows((0..10).map(|_| vec![1]).collect());
914        let matching_right = make_rows((0..10).map(|_| vec![1]).collect());
915        let bytes = merge_group_slot_bytes(20).unwrap();
916        assert!(!MergeJoinOperator::matching_groups_fit_budget(
917            &matching_left,
918            &matching_right,
919            &[0],
920            &[0],
921            19,
922            bytes,
923        )
924        .unwrap());
925        assert!(MergeJoinOperator::matching_groups_fit_budget(
926            &matching_left,
927            &matching_right,
928            &[0],
929            &[0],
930            20,
931            bytes,
932        )
933        .unwrap());
934
935        let unmatched_left = make_rows((0..100).map(|_| vec![1]).collect());
936        let unmatched_right = make_rows((0..100).map(|_| vec![2]).collect());
937        assert!(MergeJoinOperator::matching_groups_fit_budget(
938            &unmatched_left,
939            &unmatched_right,
940            &[0],
941            &[0],
942            1,
943            merge_group_slot_bytes(1).unwrap(),
944        )
945        .unwrap());
946    }
947
948    #[test]
949    fn direct_merge_owner_fails_before_duplicate_group_can_exceed_budget() {
950        let left = make_operator(vec![vec![1], vec![1]], vec!["id"]);
951        let right = make_operator(vec![vec![1], vec![1]], vec!["id"]);
952        let mut join = MergeJoinOperator::new(left, right, JoinType::Inner, vec![0], vec![0])
953            .with_group_budget(3, usize::MAX);
954
955        join.open().unwrap();
956        let error = join.next().unwrap_err();
957        assert!(error
958            .to_string()
959            .contains("duplicate group exceeds memory budget"));
960        join.close().unwrap();
961    }
962
963    #[test]
964    fn test_merge_join_uses_exact_numeric_and_nan_ordering() {
965        const EXACT: i64 = 1_i64 << 53;
966        let left = make_value_operator(
967            vec![
968                Row::from_values(vec![Value::Float(-0.0), Value::text("left-zero")]),
969                Row::from_values(vec![Value::Integer(EXACT), Value::text("left-exact")]),
970                Row::from_values(vec![
971                    Value::Integer(EXACT + 1),
972                    Value::text("left-neighbor"),
973                ]),
974                Row::from_values(vec![Value::Float(f64::NAN), Value::text("left-nan")]),
975            ],
976            vec!["key", "value"],
977        );
978        let right = make_value_operator(
979            vec![
980                Row::from_values(vec![Value::Integer(0), Value::text("right-zero")]),
981                Row::from_values(vec![Value::Float(EXACT as f64), Value::text("right-exact")]),
982                Row::from_values(vec![
983                    Value::Float(f64::from_bits(0x7ff8_0000_0000_0042)),
984                    Value::text("right-nan"),
985                ]),
986            ],
987            vec!["key", "value"],
988        );
989
990        let mut join = MergeJoinOperator::new(left, right, JoinType::Inner, vec![0], vec![0]);
991        let results = collect_results(&mut join).unwrap();
992
993        assert_eq!(results.len(), 3);
994        assert!(results
995            .iter()
996            .any(|row| row.get(1) == Some(&Value::text("left-zero"))));
997        assert!(results
998            .iter()
999            .any(|row| row.get(1) == Some(&Value::text("left-exact"))));
1000        assert!(results
1001            .iter()
1002            .any(|row| row.get(1) == Some(&Value::text("left-nan"))));
1003        assert!(!results
1004            .iter()
1005            .any(|row| row.get(1) == Some(&Value::text("left-neighbor"))));
1006    }
1007}