1use std::sync::Arc;
26use std::time::Instant;
27
28#[cfg(test)]
29use std::sync::atomic::{AtomicUsize, Ordering};
30
31use crate::context::{current_query_is_cancelled, CancellationHandle};
32use crate::expression::JoinFilter;
33use crate::lookup_key::exact_integer_pk_value;
34use crate::operator::{
35 ColumnInfo, ColumnSource, JoinProjection, Operator, OrderingProperty, RowRef,
36};
37use radixdb_core::value::NULL_VALUE;
38use radixdb_core::{CompactArc, CompactVec};
39use radixdb_core::{Result, Row, RowVec, Value, ValueMap, ValueSet};
40use radixdb_storage::expression::ConstBoolExpr;
41use radixdb_storage::instrumentation::{self, JoinExecutionKind, JoinExecutionRecord};
42use radixdb_storage::traits::{Index, Table};
43
44use super::hash_join::JoinType;
45use super::reference_unique_lookup::{
46 execute_unique_lookup_join_batch, lookup_edge_batch, LookupEdgeCardinality, LookupEdgeFallback,
47 SharedUniqueLookupRows, UniqueLookupIntegrity,
48};
49
50#[cfg(test)]
51pub(crate) static INDEX_NL_CANCELLATION_OBSERVED: AtomicUsize = AtomicUsize::new(0);
52
53#[inline]
54fn check_query_cancelled(cancellation: Option<&CancellationHandle>) -> Result<()> {
55 if cancellation.map_or_else(current_query_is_cancelled, CancellationHandle::is_cancelled) {
56 #[cfg(test)]
57 INDEX_NL_CANCELLATION_OBSERVED.fetch_add(1, Ordering::Relaxed);
58 Err(radixdb_core::Error::QueryCancelled)
59 } else {
60 Ok(())
61 }
62}
63
64fn remap_inner_projection_columns(columns: &[ColumnSource]) -> (Vec<ColumnSource>, Vec<usize>) {
65 let mut inner_indices = Vec::new();
66 let mut remapped = Vec::with_capacity(columns.len());
67
68 for col in columns {
69 match col {
70 ColumnSource::Outer(idx) => remapped.push(ColumnSource::Outer(*idx)),
71 ColumnSource::Inner(idx) => {
72 let position = if let Some(pos) = inner_indices.iter().position(|seen| seen == idx)
73 {
74 pos
75 } else {
76 inner_indices.push(*idx);
77 inner_indices.len() - 1
78 };
79 remapped.push(ColumnSource::Inner(position));
80 }
81 }
82 }
83
84 (remapped, inner_indices)
85}
86
87#[derive(Clone)]
90pub enum IndexLookupStrategy {
91 SecondaryIndex(Arc<dyn Index>),
93 SegmentedSecondaryIndex {
97 column_name: String,
98 index_name: String,
99 },
100 PrimaryKey,
103}
104
105pub struct IndexNestedLoopJoinOperator {
110 outer: Box<dyn Operator>,
112
113 inner_table: Box<dyn Table>,
115
116 join_type: JoinType,
118 outer_key_idx: usize,
119 lookup_strategy: IndexLookupStrategy,
120 residual_filter: Option<JoinFilter>,
121 cancellation: Option<CancellationHandle>,
122
123 schema: Vec<ColumnInfo>,
125 inner_col_count: usize,
126
127 projection: Option<JoinProjection>,
129 inner_projection_indices: Option<Vec<usize>>,
132 inner_lookup_key_idx: Option<usize>,
135 projection_error: Option<String>,
136 output_ordering: OrderingProperty,
137
138 current_outer_row: Option<Row>,
140 current_inner_rows: RowVec,
142 current_inner_idx: usize,
143 outer_had_match: bool,
144
145 row_id_buffer: Vec<i64>,
147
148 row_buffer: Row,
150
151 true_expr: ConstBoolExpr,
153
154 opened: bool,
156 outer_exhausted: bool,
157
158 execution_kind: JoinExecutionKind,
160 metrics_started: Option<Instant>,
161 metrics_recorded: bool,
162 observed_outer_rows: u64,
163 observed_key_rows: u64,
164 observed_lookup_calls: u64,
165 observed_lookup_candidates: u64,
166 observed_output_rows: u64,
167 outer_width: u64,
168 inner_width: u64,
169}
170
171impl IndexNestedLoopJoinOperator {
172 pub fn new(
183 outer: Box<dyn Operator>,
184 inner_table: Box<dyn Table>,
185 inner_schema: Vec<ColumnInfo>,
186 join_type: JoinType,
187 outer_key_idx: usize,
188 lookup_strategy: IndexLookupStrategy,
189 residual_filter: Option<JoinFilter>,
190 ) -> Self {
191 let lookup_is_unique = match &lookup_strategy {
192 IndexLookupStrategy::PrimaryKey => true,
193 IndexLookupStrategy::SecondaryIndex(index) => index.is_unique(),
194 IndexLookupStrategy::SegmentedSecondaryIndex { column_name, .. } => inner_table
195 .get_index_on_column(column_name)
196 .is_some_and(|index| index.is_unique()),
197 };
198 let output_ordering = if lookup_is_unique {
199 outer.ordering()
200 } else {
201 OrderingProperty::Unknown
202 };
203 let inner_table_schema = inner_table.schema();
204 let inner_lookup_key_idx = match &lookup_strategy {
205 IndexLookupStrategy::PrimaryKey => inner_table_schema.pk_column_index(),
206 IndexLookupStrategy::SecondaryIndex(index) => index
207 .column_names()
208 .first()
209 .and_then(|column_name| inner_table_schema.find_column(column_name))
210 .map(|(column_index, _)| column_index),
211 IndexLookupStrategy::SegmentedSecondaryIndex { column_name, .. } => inner_table_schema
212 .find_column(column_name)
213 .map(|(column_index, _)| column_index),
214 };
215 let mut schema = Vec::new();
217 schema.extend(outer.schema().iter().cloned());
218 schema.extend(inner_schema.iter().cloned());
219
220 let inner_col_count = inner_schema.len();
221
222 let outer_col_count = outer.schema().len();
224 let total_cols = outer_col_count + inner_col_count;
225
226 Self {
227 outer,
228 inner_table,
229 join_type,
230 outer_key_idx,
231 lookup_strategy,
232 residual_filter,
233 cancellation: None,
234 schema,
235 inner_col_count,
236 projection: None,
237 inner_projection_indices: None,
238 inner_lookup_key_idx,
239 projection_error: None,
240 output_ordering,
241 current_outer_row: None,
242 current_inner_rows: RowVec::new(),
243 current_inner_idx: 0,
244 outer_had_match: false,
245 row_id_buffer: Vec::with_capacity(16),
247 row_buffer: Row::with_capacity(total_cols),
249 true_expr: ConstBoolExpr::true_expr(),
250 opened: false,
251 outer_exhausted: false,
252 execution_kind: JoinExecutionKind::IndexNestedLoop,
253 metrics_started: None,
254 metrics_recorded: false,
255 observed_outer_rows: 0,
256 observed_key_rows: 0,
257 observed_lookup_calls: 0,
258 observed_lookup_candidates: 0,
259 observed_output_rows: 0,
260 outer_width: outer_col_count as u64,
261 inner_width: inner_col_count as u64,
262 }
263 }
264
265 fn reset_metrics(&mut self) {
266 self.metrics_started = Some(Instant::now());
267 self.metrics_recorded = false;
268 self.observed_outer_rows = 0;
269 self.observed_key_rows = 0;
270 self.observed_lookup_calls = 0;
271 self.observed_lookup_candidates = 0;
272 self.observed_output_rows = 0;
273 }
274
275 fn publish_metrics(&mut self) {
276 if self.metrics_recorded {
277 return;
278 }
279 let Some(started) = self.metrics_started.take() else {
280 return;
281 };
282 instrumentation::record_join_outer_rows(self.observed_outer_rows, self.observed_key_rows);
283 instrumentation::record_join_rows_constructed(self.observed_output_rows);
284 instrumentation::record_join_execution(
285 self.execution_kind,
286 JoinExecutionRecord {
287 left_rows: self.observed_outer_rows,
288 right_rows: self.observed_lookup_candidates,
289 output_rows: self.observed_output_rows,
290 left_width: self.outer_width,
291 right_width: self.inner_width,
292 output_width: self.schema.len() as u64,
293 candidate_pairs: self.observed_lookup_candidates,
294 lookup_calls: self.observed_lookup_calls,
295 lookup_candidate_rows: self.observed_lookup_candidates,
296 wall_nanos: started.elapsed().as_nanos().min(u128::from(u64::MAX)) as u64,
297 outer_pull_nanos: 0,
298 key_prepare_nanos: 0,
299 lookup_nanos: 0,
300 candidate_map_nanos: 0,
301 },
302 );
303 self.metrics_recorded = true;
304 }
305
306 pub fn with_cancellation(mut self, cancellation: CancellationHandle) -> Self {
307 self.cancellation = Some(cancellation);
308 self
309 }
310
311 pub fn with_projection(
321 mut self,
322 columns: Vec<ColumnSource>,
323 projected_schema: Vec<ColumnInfo>,
324 ) -> Self {
325 self.output_ordering = self.output_ordering.remap_outer_projection(&columns);
326 self.projection_error = JoinProjection {
327 columns: columns.clone(),
328 }
329 .validate(
330 self.outer.schema().len(),
331 self.inner_col_count,
332 projected_schema.len(),
333 )
334 .err()
335 .map(|error| error.to_string());
336 let (columns, inner_projection_indices) = if self.residual_filter.is_none() {
337 let (remapped, mut inner_indices) = remap_inner_projection_columns(&columns);
338 if let Some(key_idx) = self.inner_lookup_key_idx {
339 let active_idx = inner_indices
340 .iter()
341 .position(|index| *index == key_idx)
342 .unwrap_or_else(|| {
343 inner_indices.push(key_idx);
344 inner_indices.len() - 1
345 });
346 self.inner_lookup_key_idx = Some(active_idx);
347 }
348 (remapped, Some(inner_indices))
349 } else {
350 (columns, None)
351 };
352 self.projection = Some(JoinProjection { columns });
353 self.inner_projection_indices = inner_projection_indices;
354 if let Some(indices) = self.inner_projection_indices.as_ref() {
355 self.inner_width = indices.len() as u64;
356 }
357 self.schema = projected_schema;
358 self
359 }
360
361 #[inline]
364 fn null_inner_row(&self) -> Row {
365 let inner_col_count = self
366 .inner_projection_indices
367 .as_ref()
368 .map_or(self.inner_col_count, Vec::len);
369 let null_values: Vec<Value> = (0..inner_col_count).map(|_| NULL_VALUE).collect();
370 Row::from_values(null_values)
371 }
372
373 #[inline]
380 fn combine_owned_into_buffer(&mut self, outer: Row, inner: Row) {
381 match &self.projection {
382 Some(proj) => {
383 self.row_buffer.clear();
385 self.row_buffer.reserve(proj.columns.len());
386
387 let mut outer_values = outer.into_values();
389 let mut inner_values = inner.into_values();
390
391 for col_source in &proj.columns {
392 match col_source {
393 ColumnSource::Outer(idx) => {
394 self.row_buffer
395 .push(std::mem::take(&mut outer_values[*idx]));
396 }
397 ColumnSource::Inner(idx) => {
398 self.row_buffer
399 .push(std::mem::take(&mut inner_values[*idx]));
400 }
401 }
402 }
403 }
404 None => {
405 self.row_buffer.combine_into_owned(outer, inner);
407 }
408 }
409 }
410
411 #[inline]
415 fn take_from_buffer(&mut self) -> Row {
416 self.row_buffer.take_and_clear()
417 }
418
419 #[inline]
425 fn create_combined_row(&self, outer: &Row, inner: Row) -> Row {
426 match &self.projection {
427 Some(proj) => {
428 let mut values: CompactVec<Value> = CompactVec::with_capacity(proj.columns.len());
430
431 let outer_slice = outer.as_slice();
433 let mut inner_values = inner.into_values();
434
435 for col_source in &proj.columns {
436 match col_source {
437 ColumnSource::Outer(idx) => {
438 values.push(outer_slice[*idx].clone());
439 }
440 ColumnSource::Inner(idx) => {
441 values.push(std::mem::take(&mut inner_values[*idx]));
442 }
443 }
444 }
445
446 Row::from_compact_vec(values)
447 }
448 None => Row::from_combined_clone_move(outer, inner),
449 }
450 }
451
452 fn lookup_inner_rows(&mut self, key_value: &Value) -> Result<()> {
455 check_query_cancelled(self.cancellation.as_ref())?;
456 self.observed_lookup_calls = self.observed_lookup_calls.saturating_add(1);
457 self.row_id_buffer.clear();
459 self.current_inner_rows.clear();
460
461 let column_name = match &self.lookup_strategy {
462 IndexLookupStrategy::PrimaryKey => self
463 .inner_table
464 .schema()
465 .pk_column_index()
466 .and_then(|column_index| self.inner_table.schema().columns.get(column_index))
467 .map(|column| column.name.as_str()),
468 IndexLookupStrategy::SecondaryIndex(index) => {
469 index.column_names().first().map(String::as_str)
470 }
471 IndexLookupStrategy::SegmentedSecondaryIndex { column_name, .. } => {
472 Some(column_name.as_str())
473 }
474 }
475 .ok_or_else(|| {
476 radixdb_core::Error::internal("index nested-loop lookup has no key column")
477 })?;
478 let requires_table_lookup = self.inner_table.has_local_changes()
483 || matches!(
484 &self.lookup_strategy,
485 IndexLookupStrategy::SegmentedSecondaryIndex { .. }
486 )
487 || matches!(&self.lookup_strategy, IndexLookupStrategy::PrimaryKey)
488 && exact_integer_pk_value(key_value).is_none();
489 if requires_table_lookup {
490 let row_ids = self
491 .inner_table
492 .collect_row_ids_by_index_values(column_name, std::slice::from_ref(key_value))
493 .ok_or_else(|| {
494 radixdb_core::Error::internal(format!(
495 "transactional join index coverage disappeared for column {column_name}",
496 ))
497 })??;
498 self.row_id_buffer.extend(row_ids);
499 } else {
500 match &self.lookup_strategy {
501 IndexLookupStrategy::SecondaryIndex(index) => {
502 index.get_row_ids_equal_into(
503 std::slice::from_ref(key_value),
504 &mut self.row_id_buffer,
505 )?;
506 }
507 IndexLookupStrategy::SegmentedSecondaryIndex { .. } => unreachable!(
508 "segmented join lookup must use the table-level exact-index contract"
509 ),
510 IndexLookupStrategy::PrimaryKey => {
511 if let Some(id) = exact_integer_pk_value(key_value) {
512 self.row_id_buffer.push(id);
513 }
514 }
515 }
516 }
517
518 if self.row_id_buffer.is_empty() {
519 return Ok(());
520 }
521
522 if let Some(indices) = self.inner_projection_indices.as_ref() {
527 self.current_inner_rows = self
528 .inner_table
529 .collect_rows_by_ids_projected(&self.row_id_buffer, indices)?;
530 } else {
531 self.inner_table.fetch_rows_by_ids_into(
532 &self.row_id_buffer,
533 &self.true_expr,
534 &mut self.current_inner_rows,
535 )?;
536 }
537 let key_idx = self.inner_lookup_key_idx.ok_or_else(|| {
543 radixdb_core::Error::invalid_argument(
544 "index nested-loop lookup cannot resolve its inner key column",
545 )
546 })?;
547 self.current_inner_rows.retain(|(_, row)| {
548 row.get(key_idx)
549 .is_some_and(|inner_key| !inner_key.is_null() && inner_key == key_value)
550 });
551 self.observed_lookup_candidates = self
552 .observed_lookup_candidates
553 .saturating_add(self.current_inner_rows.len() as u64);
554 check_query_cancelled(self.cancellation.as_ref())?;
555 Ok(())
556 }
557
558 fn advance_outer(&mut self) -> Result<bool> {
560 check_query_cancelled(self.cancellation.as_ref())?;
561 match self.outer.next()? {
562 Some(row_ref) => {
563 self.observed_outer_rows = self.observed_outer_rows.saturating_add(1);
564 let outer_row = row_ref.into_owned();
565
566 let key_value = match outer_row.get(self.outer_key_idx) {
568 Some(v) if !v.is_null() => v.clone(),
569 _ => {
570 self.current_outer_row = Some(outer_row);
572 self.current_inner_rows.clear();
573 self.current_inner_idx = 0;
574 self.outer_had_match = false;
575 return Ok(true);
576 }
577 };
578 self.observed_key_rows = self.observed_key_rows.saturating_add(1);
579
580 self.lookup_inner_rows(&key_value)?;
582
583 self.current_outer_row = Some(outer_row);
584 self.current_inner_idx = 0;
586 self.outer_had_match = false;
587 Ok(true)
588 }
589 None => {
590 self.outer_exhausted = true;
591 Ok(false)
592 }
593 }
594 }
595}
596
597impl Operator for IndexNestedLoopJoinOperator {
598 fn open(&mut self) -> Result<()> {
599 self.reset_metrics();
600 check_query_cancelled(self.cancellation.as_ref())?;
601 if let Some(message) = self.projection_error.take() {
602 return Err(radixdb_core::Error::invalid_argument(message));
603 }
604 if let Some(projection) = &self.projection {
605 JoinProjection {
606 columns: projection.columns.clone(),
607 }
608 .validate(
609 self.outer.schema().len(),
610 self.inner_projection_indices
611 .as_ref()
612 .map_or(self.inner_col_count, Vec::len),
613 self.schema.len(),
614 )?;
615 }
616 if let Err(error) = self.outer.open() {
617 let _ = self.outer.close();
618 return Err(error);
619 }
620
621 if let Err(error) = self.advance_outer() {
623 let _ = self.outer.close();
624 return Err(error);
625 }
626
627 self.opened = true;
628 Ok(())
629 }
630
631 fn next(&mut self) -> Result<Option<RowRef>> {
632 check_query_cancelled(self.cancellation.as_ref())?;
633 if !self.opened {
634 return Err(radixdb_core::Error::internal(
635 "IndexNestedLoopJoinOperator::next called before open",
636 ));
637 }
638
639 let is_left_outer = matches!(self.join_type, JoinType::Left | JoinType::Full);
640
641 loop {
642 if self.outer_exhausted {
644 return Ok(None);
645 }
646
647 if self.current_outer_row.is_none() && !self.advance_outer()? {
649 return Ok(None);
650 }
651
652 let inner_len = self.current_inner_rows.len();
654 while self.current_inner_idx < inner_len {
655 if self.current_inner_idx & 0xff == 0 {
656 check_query_cancelled(self.cancellation.as_ref())?;
657 }
658 let inner_idx = self.current_inner_idx;
659 self.current_inner_idx += 1;
660
661 let passes_filter = {
663 let outer_row = self.current_outer_row.as_ref().unwrap();
664 let inner_entry = &self.current_inner_rows[inner_idx];
665 if let Some(ref filter) = self.residual_filter {
666 filter.matches_checked(outer_row, &inner_entry.1)?
667 } else {
668 true
669 }
670 };
671
672 if passes_filter {
673 self.outer_had_match = true;
674 let inner_row = std::mem::take(&mut self.current_inner_rows[inner_idx].1);
675
676 let is_last_inner = self.current_inner_idx >= inner_len;
679 if is_last_inner {
680 let outer_row = self.current_outer_row.take().unwrap();
682 self.advance_outer()?;
684 self.combine_owned_into_buffer(outer_row, inner_row);
685 self.observed_output_rows = self.observed_output_rows.saturating_add(1);
686 return Ok(Some(RowRef::Owned(self.take_from_buffer())));
687 } else {
688 let outer_row = self.current_outer_row.as_ref().unwrap();
691 let combined = self.create_combined_row(outer_row, inner_row);
692 self.observed_output_rows = self.observed_output_rows.saturating_add(1);
693 return Ok(Some(RowRef::Owned(combined)));
694 }
695 }
696 }
697
698 if is_left_outer && !self.outer_had_match {
701 let outer_row = self.current_outer_row.take().unwrap();
702 self.advance_outer()?;
703 let null_inner = self.null_inner_row();
704 self.combine_owned_into_buffer(outer_row, null_inner);
706 self.observed_output_rows = self.observed_output_rows.saturating_add(1);
707 return Ok(Some(RowRef::Owned(self.take_from_buffer())));
708 }
709
710 if !self.advance_outer()? {
712 return Ok(None);
713 }
714 }
715 }
716
717 fn close(&mut self) -> Result<()> {
718 let result = self.outer.close();
719 self.publish_metrics();
720 self.opened = false;
721 result
722 }
723
724 fn schema(&self) -> &[ColumnInfo] {
725 &self.schema
726 }
727
728 fn estimated_rows(&self) -> Option<usize> {
729 let outer_est = self.outer.estimated_rows()?;
731 Some(match self.join_type {
732 JoinType::Inner => outer_est, JoinType::Left | JoinType::Full => outer_est,
734 _ => outer_est,
735 })
736 }
737
738 fn ordering(&self) -> OrderingProperty {
739 self.output_ordering.clone()
740 }
741
742 fn name(&self) -> &str {
743 match self.join_type {
744 JoinType::Inner => "IndexNL (INNER)",
745 JoinType::Left => "IndexNL (LEFT)",
746 _ => "IndexNL",
747 }
748 }
749}
750
751const BATCH_INDEX_NL_OUTER_ROWS: usize = 131_072;
762const BATCH_INDEX_NL_OUTER_BYTES: usize = 64 * 1024 * 1024;
763
764pub struct BatchIndexNestedLoopJoinOperator {
771 outer: Box<dyn Operator>,
772 inner_table: Box<dyn Table>,
773 join_type: JoinType,
774 outer_key_idx: usize,
775 lookup_strategy: IndexLookupStrategy,
776 residual_filter: Option<JoinFilter>,
777 cancellation: Option<CancellationHandle>,
778 schema: Vec<ColumnInfo>,
779 inner_col_count: usize,
780 inner_lookup_column: Option<String>,
781 inner_lookup_key_idx: Option<usize>,
782 lookup_is_unique: bool,
783 output_ordering: OrderingProperty,
784 active_inner_key_idx: Option<usize>,
785 projection: Option<JoinProjection>,
786 projection_columns: Option<CompactArc<[ColumnSource]>>,
787 inner_projection_indices: Option<Vec<usize>>,
788 projection_error: Option<String>,
789 outer_batch: Vec<RowRef>,
790 outer_index: usize,
791 current_inner_index: usize,
792 current_outer_had_match: bool,
793 current_non_unique_indices: Option<CompactArc<[usize]>>,
794 unique_inner_rows_by_key: Option<SharedUniqueLookupRows>,
795 inner_rows: Option<CompactArc<Vec<Row>>>,
796 inner_row_indices_by_key: ValueMap<CompactArc<[usize]>>,
797 opened: bool,
798 outer_exhausted: bool,
799 metrics_started: Option<Instant>,
800 metrics_recorded: bool,
801 observed_outer_rows: u64,
802 observed_key_rows: u64,
803 observed_lookup_calls: u64,
804 observed_lookup_candidates: u64,
805 observed_lookup_key_rows: u64,
806 observed_lookup_distinct_keys: u64,
807 observed_output_rows: u64,
808 observed_deferred_outer_rows: u64,
809 observed_deferred_output_rows: u64,
810 observed_outer_pull_nanos: u64,
811 observed_key_prepare_nanos: u64,
812 observed_lookup_nanos: u64,
813 observed_candidate_map_nanos: u64,
814 outer_width: u64,
815 inner_width: u64,
816}
817
818impl BatchIndexNestedLoopJoinOperator {
819 pub fn new(
821 outer: Box<dyn Operator>,
822 inner_table: Box<dyn Table>,
823 inner_schema: Vec<ColumnInfo>,
824 join_type: JoinType,
825 outer_key_idx: usize,
826 lookup_strategy: IndexLookupStrategy,
827 residual_filter: Option<JoinFilter>,
828 ) -> Self {
829 let inner_table_schema = inner_table.schema();
830 let inner_lookup_column = match &lookup_strategy {
831 IndexLookupStrategy::PrimaryKey => inner_table_schema
832 .pk_column_index()
833 .and_then(|index| inner_table_schema.columns.get(index))
834 .map(|column| column.name.clone()),
835 IndexLookupStrategy::SecondaryIndex(index) => index.column_names().first().cloned(),
836 IndexLookupStrategy::SegmentedSecondaryIndex { column_name, .. } => {
837 Some(column_name.clone())
838 }
839 };
840 let inner_lookup_key_idx = inner_lookup_column.as_ref().and_then(|column_name| {
841 inner_table_schema
842 .find_column(column_name)
843 .map(|(index, _)| index)
844 });
845 let lookup_is_unique = match &lookup_strategy {
846 IndexLookupStrategy::PrimaryKey => true,
847 IndexLookupStrategy::SecondaryIndex(index) => index.is_unique(),
848 IndexLookupStrategy::SegmentedSecondaryIndex { column_name, .. } => inner_table
849 .get_index_on_column(column_name)
850 .is_some_and(|index| index.is_unique()),
851 };
852 let output_ordering = if lookup_is_unique {
853 outer.ordering()
854 } else {
855 OrderingProperty::Unknown
856 };
857 let mut schema = outer.schema().to_vec();
858 schema.extend(inner_schema.iter().cloned());
859 let outer_width = outer.schema().len() as u64;
860 let inner_width = inner_schema.len() as u64;
861 Self {
862 outer,
863 inner_table,
864 join_type,
865 outer_key_idx,
866 lookup_strategy,
867 residual_filter,
868 cancellation: None,
869 schema,
870 inner_col_count: inner_schema.len(),
871 inner_lookup_column,
872 inner_lookup_key_idx,
873 lookup_is_unique,
874 output_ordering,
875 active_inner_key_idx: inner_lookup_key_idx,
876 projection: None,
877 projection_columns: None,
878 inner_projection_indices: None,
879 projection_error: None,
880 outer_batch: Vec::with_capacity(4_096),
881 outer_index: 0,
882 current_inner_index: 0,
883 current_outer_had_match: false,
884 current_non_unique_indices: None,
885 unique_inner_rows_by_key: None,
886 inner_rows: None,
887 inner_row_indices_by_key: ValueMap::default(),
888 opened: false,
889 outer_exhausted: false,
890 metrics_started: None,
891 metrics_recorded: false,
892 observed_outer_rows: 0,
893 observed_key_rows: 0,
894 observed_lookup_calls: 0,
895 observed_lookup_candidates: 0,
896 observed_lookup_key_rows: 0,
897 observed_lookup_distinct_keys: 0,
898 observed_output_rows: 0,
899 observed_deferred_outer_rows: 0,
900 observed_deferred_output_rows: 0,
901 observed_outer_pull_nanos: 0,
902 observed_key_prepare_nanos: 0,
903 observed_lookup_nanos: 0,
904 observed_candidate_map_nanos: 0,
905 outer_width,
906 inner_width,
907 }
908 }
909
910 pub fn with_cancellation(mut self, cancellation: CancellationHandle) -> Self {
911 self.cancellation = Some(cancellation);
912 self
913 }
914
915 pub fn with_projection(
925 mut self,
926 columns: Vec<ColumnSource>,
927 projected_schema: Vec<ColumnInfo>,
928 ) -> Self {
929 self.output_ordering = self.output_ordering.remap_outer_projection(&columns);
930 self.projection_error = JoinProjection {
931 columns: columns.clone(),
932 }
933 .validate(
934 self.outer.schema().len(),
935 self.inner_col_count,
936 projected_schema.len(),
937 )
938 .err()
939 .map(|error| error.to_string());
940 let (columns, inner_projection_indices) = if self.residual_filter.is_none() {
941 let (remapped, mut inner_indices) = remap_inner_projection_columns(&columns);
942 if let Some(key_idx) = self.inner_lookup_key_idx {
943 let active_idx = inner_indices
944 .iter()
945 .position(|index| *index == key_idx)
946 .unwrap_or_else(|| {
947 inner_indices.push(key_idx);
948 inner_indices.len() - 1
949 });
950 self.active_inner_key_idx = Some(active_idx);
951 }
952 (remapped, Some(inner_indices))
953 } else {
954 (columns, None)
955 };
956 self.projection_columns = Some(CompactArc::from(columns.clone()));
957 self.projection = Some(JoinProjection { columns });
958 self.inner_projection_indices = inner_projection_indices;
959 if let Some(indices) = self.inner_projection_indices.as_ref() {
960 self.inner_width = indices.len() as u64;
961 }
962 self.schema = projected_schema;
963 self
964 }
965
966 fn reset_metrics(&mut self) {
967 self.metrics_started = Some(Instant::now());
968 self.metrics_recorded = false;
969 self.observed_outer_rows = 0;
970 self.observed_key_rows = 0;
971 self.observed_lookup_calls = 0;
972 self.observed_lookup_candidates = 0;
973 self.observed_lookup_key_rows = 0;
974 self.observed_lookup_distinct_keys = 0;
975 self.observed_output_rows = 0;
976 self.observed_deferred_outer_rows = 0;
977 self.observed_deferred_output_rows = 0;
978 self.observed_outer_pull_nanos = 0;
979 self.observed_key_prepare_nanos = 0;
980 self.observed_lookup_nanos = 0;
981 self.observed_candidate_map_nanos = 0;
982 }
983
984 fn publish_metrics(&mut self) {
985 if self.metrics_recorded {
986 return;
987 }
988 let Some(started) = self.metrics_started.take() else {
989 return;
990 };
991 instrumentation::record_join_outer_rows(self.observed_outer_rows, self.observed_key_rows);
992 instrumentation::record_join_lookup_key_batch(
993 self.observed_lookup_key_rows,
994 self.observed_lookup_distinct_keys,
995 );
996 instrumentation::record_join_rows_constructed(
997 self.observed_output_rows
998 .saturating_sub(self.observed_deferred_output_rows),
999 );
1000 instrumentation::record_join_deferred_rows(
1001 self.observed_deferred_output_rows,
1002 self.observed_deferred_outer_rows,
1003 );
1004 instrumentation::record_join_execution(
1005 JoinExecutionKind::BatchIndexNestedLoop,
1006 JoinExecutionRecord {
1007 left_rows: self.observed_outer_rows,
1008 right_rows: self.observed_lookup_candidates,
1009 output_rows: self.observed_output_rows,
1010 left_width: self.outer_width,
1011 right_width: self.inner_width,
1012 output_width: self.schema.len() as u64,
1013 candidate_pairs: self.observed_lookup_candidates,
1014 lookup_calls: self.observed_lookup_calls,
1015 lookup_candidate_rows: self.observed_lookup_candidates,
1016 wall_nanos: started.elapsed().as_nanos().min(u128::from(u64::MAX)) as u64,
1017 outer_pull_nanos: self.observed_outer_pull_nanos,
1018 key_prepare_nanos: self.observed_key_prepare_nanos,
1019 lookup_nanos: self.observed_lookup_nanos,
1020 candidate_map_nanos: self.observed_candidate_map_nanos,
1021 },
1022 );
1023 self.metrics_recorded = true;
1024 }
1025
1026 fn load_outer_batch(&mut self) -> Result<bool> {
1027 check_query_cancelled(self.cancellation.as_ref())?;
1028 self.outer_batch.clear();
1029 self.unique_inner_rows_by_key = None;
1030 self.inner_rows = None;
1031 self.inner_row_indices_by_key.clear();
1032 self.outer_index = 0;
1033 self.current_inner_index = 0;
1034 self.current_outer_had_match = false;
1035 self.current_non_unique_indices = None;
1036
1037 let outer_pull_started = Instant::now();
1038 let mut retained_bytes = 0usize;
1039 while self.outer_batch.len() < BATCH_INDEX_NL_OUTER_ROWS {
1040 let Some(row_ref) = self.outer.next()? else {
1041 self.outer_exhausted = true;
1042 break;
1043 };
1044 retained_bytes = retained_bytes.saturating_add(row_ref.estimated_retained_bytes());
1045 self.observed_deferred_outer_rows = self
1046 .observed_deferred_outer_rows
1047 .saturating_add(u64::from(row_ref.is_deferred()));
1048 self.outer_batch.push(row_ref);
1049 if retained_bytes >= BATCH_INDEX_NL_OUTER_BYTES {
1050 break;
1051 }
1052 }
1053 self.observed_outer_pull_nanos = self.observed_outer_pull_nanos.saturating_add(
1054 outer_pull_started
1055 .elapsed()
1056 .as_nanos()
1057 .min(u128::from(u64::MAX)) as u64,
1058 );
1059 if self.outer_batch.is_empty() {
1060 return Ok(false);
1061 }
1062
1063 self.observed_outer_rows = self
1064 .observed_outer_rows
1065 .saturating_add(self.outer_batch.len() as u64);
1066 let key_prepare_started = Instant::now();
1067 let mut seen = ValueSet::with_capacity(self.outer_batch.len());
1068 let mut keys = Vec::with_capacity(self.outer_batch.len());
1069 let mut key_rows = 0u64;
1070 for row in &self.outer_batch {
1071 let Some(value) = row.get(self.outer_key_idx).filter(|value| !value.is_null()) else {
1072 continue;
1073 };
1074 self.observed_key_rows = self.observed_key_rows.saturating_add(1);
1075 key_rows = key_rows.saturating_add(1);
1076 match self.lookup_strategy {
1077 IndexLookupStrategy::PrimaryKey => {
1078 let Some(key) = exact_integer_pk_value(value).map(Value::Integer) else {
1079 continue;
1080 };
1081 if seen.insert(key.clone()) {
1082 keys.push(key);
1083 }
1084 }
1085 _ => {
1086 if !seen.contains(value) {
1091 let key = value.clone();
1092 seen.insert(key.clone());
1093 keys.push(key);
1094 }
1095 }
1096 }
1097 }
1098 self.observed_lookup_key_rows = self.observed_lookup_key_rows.saturating_add(key_rows);
1099 self.observed_lookup_distinct_keys = self
1100 .observed_lookup_distinct_keys
1101 .saturating_add(keys.len() as u64);
1102 self.observed_key_prepare_nanos = self.observed_key_prepare_nanos.saturating_add(
1103 key_prepare_started
1104 .elapsed()
1105 .as_nanos()
1106 .min(u128::from(u64::MAX)) as u64,
1107 );
1108 if keys.is_empty() {
1109 return Ok(true);
1110 }
1111
1112 let column_name = self.inner_lookup_column.as_deref().ok_or_else(|| {
1113 radixdb_core::Error::invalid_argument(
1114 "batch index nested-loop lookup has no inner key column",
1115 )
1116 })?;
1117 let fallback = match &self.lookup_strategy {
1118 IndexLookupStrategy::PrimaryKey => LookupEdgeFallback::IntegerPrimaryKey,
1119 IndexLookupStrategy::SecondaryIndex(index) => {
1120 LookupEdgeFallback::SecondaryIndex(index.as_ref())
1121 }
1122 IndexLookupStrategy::SegmentedSecondaryIndex { .. } => LookupEdgeFallback::None,
1123 };
1124 let cancellation = self.cancellation.clone();
1125 let key_idx = self.active_inner_key_idx.ok_or_else(|| {
1126 radixdb_core::Error::invalid_argument(
1127 "batch index nested-loop lookup cannot resolve active inner key",
1128 )
1129 })?;
1130
1131 if self.lookup_is_unique {
1132 let lookup_started = Instant::now();
1133 let lookup = execute_unique_lookup_join_batch(
1134 self.inner_table.as_ref(),
1135 column_name,
1136 &keys,
1137 self.inner_projection_indices.as_deref(),
1138 key_idx,
1139 LookupEdgeCardinality::AtMostOne,
1140 fallback,
1141 |_| Ok(()),
1142 || check_query_cancelled(cancellation.as_ref()),
1143 )?
1144 .ok_or_else(|| lookup_disappeared_error(&self.lookup_strategy))?;
1145 if lookup.integrity != UniqueLookupIntegrity::Complete {
1146 return Err(radixdb_core::Error::internal(
1147 "proven unique batch join lookup returned duplicate visible rows",
1148 ));
1149 }
1150 self.observed_lookup_calls = self
1151 .observed_lookup_calls
1152 .saturating_add(lookup.lookup_calls);
1153 self.observed_lookup_candidates = self
1154 .observed_lookup_candidates
1155 .saturating_add(lookup.candidate_rows as u64);
1156 self.unique_inner_rows_by_key = Some(lookup.rows_by_key.into_shared());
1157 self.observed_lookup_nanos = self.observed_lookup_nanos.saturating_add(
1158 lookup_started
1159 .elapsed()
1160 .as_nanos()
1161 .min(u128::from(u64::MAX)) as u64,
1162 );
1163 } else {
1164 let lookup_started = Instant::now();
1165 let lookup = lookup_edge_batch(
1166 self.inner_table.as_ref(),
1167 column_name,
1168 &keys,
1169 self.inner_projection_indices.as_deref(),
1170 LookupEdgeCardinality::Unbounded,
1171 fallback,
1172 || check_query_cancelled(cancellation.as_ref()),
1173 )?
1174 .ok_or_else(|| lookup_disappeared_error(&self.lookup_strategy))?;
1175 self.observed_lookup_nanos = self.observed_lookup_nanos.saturating_add(
1176 lookup_started
1177 .elapsed()
1178 .as_nanos()
1179 .min(u128::from(u64::MAX)) as u64,
1180 );
1181 self.observed_lookup_calls = self
1182 .observed_lookup_calls
1183 .saturating_add(lookup.lookup_calls);
1184 self.observed_lookup_candidates = self
1185 .observed_lookup_candidates
1186 .saturating_add(lookup.rows.len() as u64);
1187 let candidate_map_started = Instant::now();
1188 let mut inner_rows = Vec::with_capacity(lookup.rows.len());
1189 let mut row_indices_by_key: ValueMap<Vec<usize>> = ValueMap::default();
1190 for (_, row) in lookup.rows {
1191 if let Some(key) = row.get(key_idx).cloned() {
1192 row_indices_by_key
1193 .entry(key)
1194 .or_default()
1195 .push(inner_rows.len());
1196 }
1197 inner_rows.push(row);
1198 }
1199 for (key, indices) in row_indices_by_key {
1200 self.inner_row_indices_by_key
1201 .insert(key, CompactArc::from(indices));
1202 }
1203 self.inner_rows = Some(CompactArc::new(inner_rows));
1204 self.observed_candidate_map_nanos = self.observed_candidate_map_nanos.saturating_add(
1205 candidate_map_started
1206 .elapsed()
1207 .as_nanos()
1208 .min(u128::from(u64::MAX)) as u64,
1209 );
1210 }
1211 check_query_cancelled(self.cancellation.as_ref())?;
1212 Ok(true)
1213 }
1214
1215 fn combine_rows(&self, outer: RowRef, inner: RowRef) -> RowRef {
1216 if let Some(columns) = &self.projection_columns {
1217 RowRef::projected(outer, inner, CompactArc::clone(columns))
1218 } else {
1219 RowRef::Owned(Row::from_combined_owned(
1220 outer.into_owned(),
1221 inner.into_owned(),
1222 ))
1223 }
1224 }
1225
1226 #[inline]
1227 fn take_current_outer(&mut self) -> RowRef {
1228 std::mem::replace(
1229 &mut self.outer_batch[self.outer_index],
1230 RowRef::Owned(Row::new()),
1231 )
1232 }
1233
1234 fn null_inner_row(&self) -> Row {
1235 let width = self
1236 .inner_projection_indices
1237 .as_ref()
1238 .map_or(self.inner_col_count, Vec::len);
1239 Row::from_values((0..width).map(|_| NULL_VALUE).collect())
1240 }
1241
1242 fn advance_outer(&mut self) {
1243 self.outer_index += 1;
1244 self.current_inner_index = 0;
1245 self.current_outer_had_match = false;
1246 self.current_non_unique_indices = None;
1247 }
1248
1249 fn lookup_candidate(&mut self) -> (usize, Option<RowRef>) {
1250 if let Some(rows) = self.unique_inner_rows_by_key.as_ref() {
1251 let value = self.outer_batch[self.outer_index].get(self.outer_key_idx);
1252 let row = match self.lookup_strategy {
1253 IndexLookupStrategy::PrimaryKey => value
1254 .and_then(exact_integer_pk_value)
1255 .map(Value::Integer)
1256 .as_ref()
1257 .and_then(|key| rows.get(key)),
1258 _ => value
1259 .filter(|value| !value.is_null())
1260 .and_then(|key| rows.get(key)),
1261 };
1262 return (
1263 usize::from(row.is_some()),
1264 if self.current_inner_index == 0 {
1265 row.map(|(rows, row_idx)| RowRef::shared(rows, row_idx))
1266 } else {
1267 None
1268 },
1269 );
1270 }
1271 if self.current_non_unique_indices.is_none() {
1272 let value = self.outer_batch[self.outer_index].get(self.outer_key_idx);
1273 self.current_non_unique_indices = match self.lookup_strategy {
1274 IndexLookupStrategy::PrimaryKey => value
1275 .and_then(exact_integer_pk_value)
1276 .map(Value::Integer)
1277 .as_ref()
1278 .and_then(|key| self.inner_row_indices_by_key.get(key)),
1279 _ => value
1280 .filter(|value| !value.is_null())
1281 .and_then(|key| self.inner_row_indices_by_key.get(key)),
1282 }
1283 .cloned();
1284 }
1285 let candidates = self.current_non_unique_indices.as_ref();
1286 (
1287 candidates.map_or(0, |indices| indices.len()),
1288 candidates.and_then(|indices| {
1289 let row_idx = *indices.get(self.current_inner_index)?;
1290 Some(RowRef::shared(
1291 CompactArc::clone(self.inner_rows.as_ref()?),
1292 row_idx,
1293 ))
1294 }),
1295 )
1296 }
1297}
1298
1299fn lookup_disappeared_error(lookup_strategy: &IndexLookupStrategy) -> radixdb_core::Error {
1300 match lookup_strategy {
1301 IndexLookupStrategy::SegmentedSecondaryIndex { column_name, .. } => {
1302 radixdb_core::Error::internal(format!(
1303 "segmented batch join index coverage disappeared for column {column_name}"
1304 ))
1305 }
1306 _ => radixdb_core::Error::internal(
1307 "batch index nested-loop selected lookup disappeared during execution",
1308 ),
1309 }
1310}
1311
1312impl Operator for BatchIndexNestedLoopJoinOperator {
1313 fn open(&mut self) -> Result<()> {
1314 self.reset_metrics();
1315 self.outer_exhausted = false;
1316 check_query_cancelled(self.cancellation.as_ref())?;
1317 if let Some(message) = self.projection_error.take() {
1318 return Err(radixdb_core::Error::invalid_argument(message));
1319 }
1320 if let Err(error) = self.outer.open() {
1321 let _ = self.outer.close();
1322 return Err(error);
1323 }
1324 if let Err(error) = self.load_outer_batch() {
1325 let _ = self.outer.close();
1326 return Err(error);
1327 }
1328 self.opened = true;
1329 Ok(())
1330 }
1331
1332 fn next(&mut self) -> Result<Option<RowRef>> {
1333 check_query_cancelled(self.cancellation.as_ref())?;
1334 if !self.opened {
1335 return Err(radixdb_core::Error::internal(
1336 "BatchIndexNestedLoopJoinOperator::next called before open",
1337 ));
1338 }
1339 let is_left_outer = matches!(self.join_type, JoinType::Left | JoinType::Full);
1340
1341 loop {
1342 if self.outer_index >= self.outer_batch.len()
1343 && (self.outer_exhausted || !self.load_outer_batch()?)
1344 {
1345 return Ok(None);
1346 }
1347
1348 let (candidate_count, candidate) = self.lookup_candidate();
1349
1350 if let Some(inner) = candidate {
1351 self.current_inner_index += 1;
1352 let passes = self.residual_filter.as_ref().map_or(Ok(true), |filter| {
1353 filter.matches_row_refs_checked(&self.outer_batch[self.outer_index], &inner)
1354 })?;
1355 if passes {
1356 self.current_outer_had_match = true;
1357 let outer = if self.current_inner_index >= candidate_count {
1358 let outer = self.take_current_outer();
1359 self.advance_outer();
1360 outer
1361 } else {
1362 self.outer_batch[self.outer_index].clone()
1363 };
1364 let output = self.combine_rows(outer, inner);
1365 self.observed_output_rows = self.observed_output_rows.saturating_add(1);
1366 self.observed_deferred_output_rows = self
1367 .observed_deferred_output_rows
1368 .saturating_add(u64::from(output.is_deferred()));
1369 return Ok(Some(output));
1370 }
1371 continue;
1372 }
1373
1374 if is_left_outer && !self.current_outer_had_match {
1375 let outer = self.take_current_outer();
1376 let output = self.combine_rows(outer, RowRef::Owned(self.null_inner_row()));
1377 self.advance_outer();
1378 self.observed_output_rows = self.observed_output_rows.saturating_add(1);
1379 self.observed_deferred_output_rows = self
1380 .observed_deferred_output_rows
1381 .saturating_add(u64::from(output.is_deferred()));
1382 return Ok(Some(output));
1383 }
1384 self.advance_outer();
1385 }
1386 }
1387
1388 fn close(&mut self) -> Result<()> {
1389 let result = self.outer.close();
1390 self.publish_metrics();
1391 self.outer_batch.clear();
1392 self.current_non_unique_indices = None;
1393 self.unique_inner_rows_by_key = None;
1394 self.inner_rows = None;
1395 self.inner_row_indices_by_key.clear();
1396 self.opened = false;
1397 result
1398 }
1399
1400 fn schema(&self) -> &[ColumnInfo] {
1401 &self.schema
1402 }
1403
1404 fn estimated_rows(&self) -> Option<usize> {
1405 self.outer.estimated_rows()
1406 }
1407
1408 fn ordering(&self) -> OrderingProperty {
1409 self.output_ordering.clone()
1410 }
1411
1412 fn name(&self) -> &str {
1413 "BatchIndexNL"
1414 }
1415}
1416
1417#[cfg(test)]
1418mod tests {
1419 use super::*;
1420 use crate::operator::MaterializedOperator;
1421 use radixdb_storage::mvcc::engine::MVCCEngine;
1422 use radixdb_storage::traits::Engine;
1423 use std::sync::atomic::{AtomicUsize, Ordering as AtomicOrdering};
1424 use std::sync::Arc;
1425
1426 struct CountingOuter {
1427 rows: Vec<Row>,
1428 index: usize,
1429 consumed: Arc<AtomicUsize>,
1430 schema: Vec<ColumnInfo>,
1431 }
1432
1433 impl Operator for CountingOuter {
1434 fn open(&mut self) -> Result<()> {
1435 self.index = 0;
1436 Ok(())
1437 }
1438
1439 fn next(&mut self) -> Result<Option<RowRef>> {
1440 let Some(row) = self.rows.get(self.index).cloned() else {
1441 return Ok(None);
1442 };
1443 self.index += 1;
1444 self.consumed.fetch_add(1, AtomicOrdering::Relaxed);
1445 Ok(Some(RowRef::Owned(row)))
1446 }
1447
1448 fn close(&mut self) -> Result<()> {
1449 Ok(())
1450 }
1451
1452 fn schema(&self) -> &[ColumnInfo] {
1453 &self.schema
1454 }
1455
1456 fn name(&self) -> &str {
1457 "CountingOuter"
1458 }
1459 }
1460
1461 #[test]
1462 fn remap_inner_projection_columns_keeps_outer_and_compacts_inner_columns() {
1463 let columns = vec![
1464 ColumnSource::Outer(1),
1465 ColumnSource::Inner(3),
1466 ColumnSource::Inner(1),
1467 ColumnSource::Outer(0),
1468 ColumnSource::Inner(3),
1469 ];
1470
1471 let (remapped, inner_indices) = remap_inner_projection_columns(&columns);
1472
1473 assert_eq!(
1474 remapped,
1475 vec![
1476 ColumnSource::Outer(1),
1477 ColumnSource::Inner(0),
1478 ColumnSource::Inner(1),
1479 ColumnSource::Outer(0),
1480 ColumnSource::Inner(0),
1481 ]
1482 );
1483 assert_eq!(inner_indices, vec![3, 1]);
1484 }
1485
1486 fn make_outer_operator() -> Box<dyn Operator> {
1487 Box::new(MaterializedOperator::new(
1488 vec![
1489 Row::from_values(vec![Value::integer(1), Value::integer(10)]),
1490 Row::from_values(vec![Value::integer(3), Value::integer(30)]),
1491 ],
1492 vec![ColumnInfo::new("outer_id"), ColumnInfo::new("outer_value")],
1493 ))
1494 }
1495
1496 fn make_inner_table() -> (Arc<MVCCEngine>, Box<dyn Table>) {
1497 let engine = Arc::new(MVCCEngine::in_memory());
1498 engine.open_engine().unwrap();
1499 let executor = crate::Executor::new(Arc::clone(&engine));
1500 executor
1501 .execute(
1502 "CREATE TABLE inner_join_target (id INTEGER PRIMARY KEY, data INTEGER NOT NULL, unused INTEGER NOT NULL)",
1503 )
1504 .unwrap();
1505 drop(executor);
1506
1507 {
1508 let mut tx = engine.begin_transaction().unwrap();
1509 let mut table = tx.get_table("inner_join_target").unwrap();
1510 table
1511 .insert(Row::from_values(vec![
1512 Value::integer(1),
1513 Value::integer(100),
1514 Value::integer(1000),
1515 ]))
1516 .unwrap();
1517 table
1518 .insert(Row::from_values(vec![
1519 Value::integer(2),
1520 Value::integer(200),
1521 Value::integer(2000),
1522 ]))
1523 .unwrap();
1524 table
1525 .insert(Row::from_values(vec![
1526 Value::integer(3),
1527 Value::integer(300),
1528 Value::integer(3000),
1529 ]))
1530 .unwrap();
1531 tx.commit().unwrap();
1532 }
1533
1534 let tx = engine.begin_transaction().unwrap();
1535 let table = tx.get_table("inner_join_target").unwrap();
1536 (engine, table)
1537 }
1538
1539 fn projected_schema() -> Vec<ColumnInfo> {
1540 vec![ColumnInfo::new("outer_value"), ColumnInfo::new("data")]
1541 }
1542
1543 fn collect_operator_rows(op: &mut dyn Operator) -> Vec<Row> {
1544 let mut rows = Vec::new();
1545 op.open().unwrap();
1546 while let Some(row_ref) = op.next().unwrap() {
1547 rows.push(row_ref.into_owned());
1548 }
1549 op.close().unwrap();
1550 rows
1551 }
1552
1553 #[test]
1554 fn streaming_index_nested_loop_projection_returns_requested_columns_only() {
1555 let (_engine, inner_table) = make_inner_table();
1556 let inner_schema = vec![
1557 ColumnInfo::new("id"),
1558 ColumnInfo::new("data"),
1559 ColumnInfo::new("unused"),
1560 ];
1561 let mut join = IndexNestedLoopJoinOperator::new(
1562 make_outer_operator(),
1563 inner_table,
1564 inner_schema,
1565 JoinType::Inner,
1566 0,
1567 IndexLookupStrategy::PrimaryKey,
1568 None,
1569 )
1570 .with_projection(
1571 vec![ColumnSource::Outer(1), ColumnSource::Inner(1)],
1572 projected_schema(),
1573 );
1574
1575 let before = instrumentation::snapshot();
1576 let rows = collect_operator_rows(&mut join);
1577 let after = instrumentation::snapshot();
1578
1579 assert_eq!(rows.len(), 2);
1580 assert!(rows.iter().all(|row| row.len() == 2));
1581 assert_eq!(rows[0].get(0), Some(&Value::integer(10)));
1582 assert_eq!(rows[0].get(1), Some(&Value::integer(100)));
1583 assert_eq!(rows[1].get(0), Some(&Value::integer(30)));
1584 assert_eq!(rows[1].get(1), Some(&Value::integer(300)));
1585 assert!(
1586 after.join_index_nested_loop_calls
1587 >= before.join_index_nested_loop_calls.saturating_add(1)
1588 );
1589 assert!(after.join_lookup_calls >= before.join_lookup_calls.saturating_add(2));
1590 assert!(
1591 after.join_lookup_candidate_rows >= before.join_lookup_candidate_rows.saturating_add(2)
1592 );
1593 assert!(after.join_output_rows >= before.join_output_rows.saturating_add(2));
1594 }
1595
1596 #[test]
1597 fn streaming_index_nested_loop_observes_a_pre_cancelled_request() {
1598 let (_engine, inner_table) = make_inner_table();
1599 let context = crate::context::ExecutionContext::new();
1600 let cancellation = context.cancellation_handle();
1601 cancellation.cancel();
1602 let mut join = IndexNestedLoopJoinOperator::new(
1603 make_outer_operator(),
1604 inner_table,
1605 vec![
1606 ColumnInfo::new("id"),
1607 ColumnInfo::new("data"),
1608 ColumnInfo::new("unused"),
1609 ],
1610 JoinType::Inner,
1611 0,
1612 IndexLookupStrategy::PrimaryKey,
1613 None,
1614 )
1615 .with_cancellation(cancellation);
1616
1617 let before = INDEX_NL_CANCELLATION_OBSERVED.load(AtomicOrdering::Relaxed);
1618 assert!(matches!(
1619 join.open(),
1620 Err(radixdb_core::Error::QueryCancelled)
1621 ));
1622 assert!(
1623 INDEX_NL_CANCELLATION_OBSERVED.load(AtomicOrdering::Relaxed) > before,
1624 "index nested-loop operator did not observe its cancellation handle"
1625 );
1626 }
1627
1628 #[test]
1629 fn primary_key_lookup_admits_only_exact_integer_domain_keys() {
1630 let (_engine, inner_table) = make_inner_table();
1631 let outer = Box::new(MaterializedOperator::new(
1632 vec![
1633 Row::from_values(vec![Value::Float(0.5)]),
1634 Row::from_values(vec![Value::Float(1.0)]),
1635 Row::from_values(vec![Value::Integer(3)]),
1636 ],
1637 vec![ColumnInfo::new("outer_key")],
1638 ));
1639 let mut join = IndexNestedLoopJoinOperator::new(
1640 outer,
1641 inner_table,
1642 vec![
1643 ColumnInfo::new("id"),
1644 ColumnInfo::new("data"),
1645 ColumnInfo::new("unused"),
1646 ],
1647 JoinType::Inner,
1648 0,
1649 IndexLookupStrategy::PrimaryKey,
1650 None,
1651 )
1652 .with_projection(
1653 vec![ColumnSource::Outer(0), ColumnSource::Inner(1)],
1654 vec![ColumnInfo::new("outer_key"), ColumnInfo::new("data")],
1655 );
1656
1657 let rows = collect_operator_rows(&mut join);
1658 assert_eq!(rows.len(), 2);
1659 assert_eq!(rows[0].get(0), Some(&Value::Float(1.0)));
1660 assert_eq!(rows[0].get(1), Some(&Value::Integer(100)));
1661 assert_eq!(rows[1].get(0), Some(&Value::Integer(3)));
1662 assert_eq!(rows[1].get(1), Some(&Value::Integer(300)));
1663 }
1664
1665 #[test]
1666 fn public_index_nested_loop_rejects_invalid_projection_before_lookup() {
1667 let (_engine, inner_table) = make_inner_table();
1668 let mut join = IndexNestedLoopJoinOperator::new(
1669 make_outer_operator(),
1670 inner_table,
1671 vec![
1672 ColumnInfo::new("id"),
1673 ColumnInfo::new("data"),
1674 ColumnInfo::new("unused"),
1675 ],
1676 JoinType::Inner,
1677 0,
1678 IndexLookupStrategy::PrimaryKey,
1679 None,
1680 )
1681 .with_projection(
1682 vec![ColumnSource::Inner(3)],
1683 vec![ColumnInfo::new("invalid")],
1684 );
1685
1686 assert!(join.open().is_err());
1687 }
1688
1689 #[test]
1690 fn batch_index_nested_loop_fetches_multiple_terminal_columns_once() {
1691 let (_engine, inner_table) = make_inner_table();
1692 let inner_schema = vec![
1693 ColumnInfo::new("id"),
1694 ColumnInfo::new("data"),
1695 ColumnInfo::new("unused"),
1696 ];
1697 let mut join = BatchIndexNestedLoopJoinOperator::new(
1698 make_outer_operator(),
1699 inner_table,
1700 inner_schema,
1701 JoinType::Inner,
1702 0,
1703 IndexLookupStrategy::PrimaryKey,
1704 None,
1705 )
1706 .with_projection(
1707 vec![
1708 ColumnSource::Outer(1),
1709 ColumnSource::Inner(1),
1710 ColumnSource::Inner(2),
1711 ],
1712 vec![
1713 ColumnInfo::new("outer_value"),
1714 ColumnInfo::new("data"),
1715 ColumnInfo::new("unused"),
1716 ],
1717 );
1718
1719 assert_eq!(join.inner_projection_indices, Some(vec![1, 2, 0]));
1720 instrumentation::begin_join_execution_probe();
1721 let mut rows = collect_operator_rows(&mut join);
1722 let probe = instrumentation::end_join_execution_probe();
1723 rows.sort_by_key(|row| row.get(0).and_then(Value::as_int64).unwrap_or_default());
1724
1725 assert_eq!(rows.len(), 2);
1726 assert!(rows.iter().all(|row| row.len() == 3));
1727 assert_eq!(rows[0].get(0), Some(&Value::integer(10)));
1728 assert_eq!(rows[0].get(1), Some(&Value::integer(100)));
1729 assert_eq!(rows[0].get(2), Some(&Value::integer(1000)));
1730 assert_eq!(rows[1].get(0), Some(&Value::integer(30)));
1731 assert_eq!(rows[1].get(1), Some(&Value::integer(300)));
1732 assert_eq!(rows[1].get(2), Some(&Value::integer(3000)));
1733 assert_eq!(probe.lookup_calls, 1);
1734 assert_eq!(probe.lookup_candidate_rows, 2);
1735 }
1736
1737 #[test]
1738 fn batch_index_nested_loop_keeps_projection_deferred_across_edges() {
1739 let (engine, first_inner) = make_inner_table();
1740 let first = BatchIndexNestedLoopJoinOperator::new(
1741 make_outer_operator(),
1742 first_inner,
1743 vec![
1744 ColumnInfo::new("id"),
1745 ColumnInfo::new("data"),
1746 ColumnInfo::new("unused"),
1747 ],
1748 JoinType::Inner,
1749 0,
1750 IndexLookupStrategy::PrimaryKey,
1751 None,
1752 )
1753 .with_projection(
1754 vec![
1755 ColumnSource::Outer(0),
1756 ColumnSource::Outer(1),
1757 ColumnSource::Inner(1),
1758 ],
1759 vec![
1760 ColumnInfo::new("id"),
1761 ColumnInfo::new("outer_value"),
1762 ColumnInfo::new("first_data"),
1763 ],
1764 );
1765
1766 let tx = engine.begin_transaction().unwrap();
1767 let second_inner = tx.get_table("inner_join_target").unwrap();
1768 let mut second = BatchIndexNestedLoopJoinOperator::new(
1769 Box::new(first),
1770 second_inner,
1771 vec![
1772 ColumnInfo::new("id"),
1773 ColumnInfo::new("data"),
1774 ColumnInfo::new("unused"),
1775 ],
1776 JoinType::Inner,
1777 0,
1778 IndexLookupStrategy::PrimaryKey,
1779 None,
1780 )
1781 .with_projection(
1782 vec![
1783 ColumnSource::Outer(1),
1784 ColumnSource::Outer(2),
1785 ColumnSource::Inner(2),
1786 ],
1787 vec![
1788 ColumnInfo::new("outer_value"),
1789 ColumnInfo::new("first_data"),
1790 ColumnInfo::new("second_unused"),
1791 ],
1792 );
1793
1794 second.open().unwrap();
1795 let first_output = second.next().unwrap().unwrap();
1796 assert!(first_output.is_deferred());
1797 assert_eq!(first_output.get(0), Some(&Value::integer(10)));
1798 assert_eq!(first_output.get(1), Some(&Value::integer(100)));
1799 assert_eq!(first_output.get(2), Some(&Value::integer(1000)));
1800 assert_eq!(first_output.into_owned().len(), 3);
1801
1802 let second_output = second.next().unwrap().unwrap();
1803 assert!(second_output.is_deferred());
1804 assert_eq!(second_output.get(0), Some(&Value::integer(30)));
1805 assert_eq!(second_output.get(1), Some(&Value::integer(300)));
1806 assert_eq!(second_output.get(2), Some(&Value::integer(3000)));
1807 assert!(second.next().unwrap().is_none());
1808 second.close().unwrap();
1809 assert_eq!(
1810 second.observed_deferred_outer_rows, 2,
1811 "the second edge must consume both deferred rows from the first edge"
1812 );
1813 assert_eq!(
1814 second.observed_deferred_output_rows, 2,
1815 "the second edge must preserve both outputs as deferred rows"
1816 );
1817 assert_eq!(
1818 second
1819 .observed_output_rows
1820 .saturating_sub(second.observed_deferred_output_rows),
1821 0,
1822 "the second edge must not construct an owned output row"
1823 );
1824 }
1825
1826 #[test]
1827 fn r8_l01_batch_h_batch_index_nl_does_not_drain_outer_in_open() {
1828 let (_engine, inner_table) = make_inner_table();
1829 let consumed = Arc::new(AtomicUsize::new(0));
1830 let outer = CountingOuter {
1831 rows: (0..200_000)
1832 .map(|index| Row::from_values(vec![Value::Integer(index)]))
1833 .collect(),
1834 index: 0,
1835 consumed: Arc::clone(&consumed),
1836 schema: vec![ColumnInfo::new("id")],
1837 };
1838 let mut join = BatchIndexNestedLoopJoinOperator::new(
1839 Box::new(outer),
1840 inner_table,
1841 vec![
1842 ColumnInfo::new("id"),
1843 ColumnInfo::new("data"),
1844 ColumnInfo::new("unused"),
1845 ],
1846 JoinType::Left,
1847 0,
1848 IndexLookupStrategy::PrimaryKey,
1849 None,
1850 );
1851
1852 join.open().unwrap();
1853 assert_eq!(
1854 consumed.load(AtomicOrdering::Relaxed),
1855 BATCH_INDEX_NL_OUTER_ROWS
1856 );
1857 assert!(join.next().unwrap().is_some());
1858 assert!(consumed.load(AtomicOrdering::Relaxed) < 200_000);
1859 join.close().unwrap();
1860 }
1861
1862 #[test]
1863 fn batch_index_nested_loop_bounds_wide_outer_rows_by_bytes() {
1864 let (_engine, inner_table) = make_inner_table();
1865 let consumed = Arc::new(AtomicUsize::new(0));
1866 let payload = "x".repeat(1024 * 1024);
1867 let outer = CountingOuter {
1868 rows: (0..128)
1869 .map(|index| {
1870 Row::from_values(vec![Value::Integer(index), Value::text(payload.clone())])
1871 })
1872 .collect(),
1873 index: 0,
1874 consumed: Arc::clone(&consumed),
1875 schema: vec![ColumnInfo::new("id"), ColumnInfo::new("payload")],
1876 };
1877 let mut join = BatchIndexNestedLoopJoinOperator::new(
1878 Box::new(outer),
1879 inner_table,
1880 vec![
1881 ColumnInfo::new("id"),
1882 ColumnInfo::new("data"),
1883 ColumnInfo::new("unused"),
1884 ],
1885 JoinType::Left,
1886 0,
1887 IndexLookupStrategy::PrimaryKey,
1888 None,
1889 );
1890
1891 join.open().unwrap();
1892 let retained: usize = join
1893 .outer_batch
1894 .iter()
1895 .map(RowRef::estimated_retained_bytes)
1896 .sum();
1897 assert!(retained >= BATCH_INDEX_NL_OUTER_BYTES);
1898 assert!(join.outer_batch.len() < 128);
1899 assert_eq!(
1900 consumed.load(AtomicOrdering::Relaxed),
1901 join.outer_batch.len()
1902 );
1903 join.close().unwrap();
1904 }
1905
1906 #[test]
1907 fn batch_index_nested_loop_deduplicates_keys_inside_bounded_chunks() {
1908 let (_engine, inner_table) = make_inner_table();
1909 let outer = Box::new(MaterializedOperator::new(
1910 (0..1_000)
1911 .map(|_| Row::from_values(vec![Value::Integer(1)]))
1912 .collect(),
1913 vec![ColumnInfo::new("id")],
1914 ));
1915 let mut join = BatchIndexNestedLoopJoinOperator::new(
1916 outer,
1917 inner_table,
1918 vec![
1919 ColumnInfo::new("id"),
1920 ColumnInfo::new("data"),
1921 ColumnInfo::new("unused"),
1922 ],
1923 JoinType::Inner,
1924 0,
1925 IndexLookupStrategy::PrimaryKey,
1926 None,
1927 );
1928
1929 instrumentation::begin_join_execution_probe();
1930 let rows = collect_operator_rows(&mut join);
1931 let probe = instrumentation::end_join_execution_probe();
1932
1933 assert_eq!(rows.len(), 1_000);
1934 assert_eq!(probe.operator_calls, 1);
1935 assert_eq!(probe.batch_index_nested_loop_calls, 1);
1936 assert_eq!(probe.lookup_calls, 1);
1937 assert_eq!(probe.lookup_candidate_rows, 1);
1938 assert_eq!(probe.output_rows, 1_000);
1939 }
1940}