1use 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 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
86fn merge_group_slot_bytes(rows: usize) -> Option<usize> {
91 rows.checked_mul(std::mem::size_of::<Row>())?.checked_mul(2)
92}
93
94pub struct MergeJoinOperator {
100 left: Box<dyn Operator>,
102 right: Box<dyn Operator>,
103
104 join_type: JoinType,
106 left_key_indices: Vec<usize>,
107 right_key_indices: Vec<usize>,
108
109 schema: Vec<ColumnInfo>,
111 left_col_count: usize,
112 right_col_count: usize,
113
114 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_left: Vec<Value>,
128 cached_null_right: Vec<Value>,
129
130 input_ordering_certified: bool,
132 output_ordering: OrderingProperty,
133
134 opened: bool,
136}
137
138impl MergeJoinOperator {
139 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 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 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(), cached_null_right: Vec::new(), input_ordering_certified,
204 output_ordering,
205 opened: false,
206 }
207 }
208
209 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 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 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 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 #[inline]
345 fn null_left_row(&self) -> Row {
346 Row::from_values(self.cached_null_left.clone())
347 }
348
349 #[inline]
351 fn null_right_row(&self) -> Row {
352 Row::from_values(self.cached_null_right.clone())
353 }
354
355 #[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 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 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 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], vec![0], );
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 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 assert_eq!(results.len(), 3);
815
816 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 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 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}