Skip to main content

radixdb_executor/operators/
index_nested_loop.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//! Index Nested Loop Join Operator.
16//!
17//! This operator implements index nested loop join with O(N * log M) complexity
18//! by using indexes on the inner table for lookups. It's optimal when:
19//! - The inner (right) table has an index on the join key column
20//! - The outer (left) table is small or has good selectivity
21//!
22//! For each row in the outer table, we use the index/PK to find matching rows
23//! in the inner table, avoiding a full scan of the inner table.
24
25use 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/// Index Nested Loop Join lookup strategy.
88/// Determines how to find matching rows in the inner (right) table.
89#[derive(Clone)]
90pub enum IndexLookupStrategy {
91    /// Use a secondary index for lookups (index.get_row_ids_equal)
92    SecondaryIndex(Arc<dyn Index>),
93    /// Use the table's complete hot+cold exact candidate contract. The stored
94    /// column name lets a segmented table resolve immutable postings without
95    /// exposing them as a mutable hot Index handle.
96    SegmentedSecondaryIndex {
97        column_name: String,
98        index_name: String,
99    },
100    /// Use primary key lookup (direct row_id = value)
101    /// In radixdb, PRIMARY KEY INTEGER values ARE the row_ids
102    PrimaryKey,
103}
104
105/// Index Nested Loop Join Operator.
106///
107/// For each row in the outer input, looks up matching rows in the inner table
108/// using an index. This avoids full table scans of the inner table.
109pub struct IndexNestedLoopJoinOperator {
110    // Outer input operator
111    outer: Box<dyn Operator>,
112
113    // Inner table (accessed via index)
114    inner_table: Box<dyn Table>,
115
116    // Join configuration
117    join_type: JoinType,
118    outer_key_idx: usize,
119    lookup_strategy: IndexLookupStrategy,
120    residual_filter: Option<JoinFilter>,
121    cancellation: Option<CancellationHandle>,
122
123    // Output schema
124    schema: Vec<ColumnInfo>,
125    inner_col_count: usize,
126
127    // Optional projection pushdown: create projected rows during combine
128    projection: Option<JoinProjection>,
129    // Inner table columns fetched when projection is safe to push below row-id lookup.
130    // Projection ColumnSource::Inner indices are remapped to these positions.
131    inner_projection_indices: Option<Vec<usize>>,
132    // Active-row position of the INTEGER PK used by direct row-id lookup. The
133    // position is remapped when inner projection pushdown is enabled.
134    inner_lookup_key_idx: Option<usize>,
135    projection_error: Option<String>,
136    output_ordering: OrderingProperty,
137
138    // Current state
139    current_outer_row: Option<Row>,
140    // Optimization: Store (id, row) to verify specific inner rows if needed
141    current_inner_rows: RowVec,
142    current_inner_idx: usize,
143    outer_had_match: bool,
144
145    // Optimization: Reusable buffer for row IDs to avoid allocation per outer row
146    row_id_buffer: Vec<i64>,
147
148    // Optimization: Reusable row buffer to avoid allocation per join output
149    row_buffer: Row,
150
151    // Expression for fetching rows (always true - we apply residual separately)
152    true_expr: ConstBoolExpr,
153
154    // State tracking
155    opened: bool,
156    outer_exhausted: bool,
157
158    // Bounded operator-local observability, published once by close().
159    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    /// Create a new index nested loop join operator.
173    ///
174    /// # Arguments
175    /// * `outer` - Outer input operator
176    /// * `inner_table` - Inner table to lookup from
177    /// * `inner_schema` - Schema of the inner table
178    /// * `join_type` - Type of join (INNER or LEFT)
179    /// * `outer_key_idx` - Column index of the join key in outer rows
180    /// * `lookup_strategy` - How to find matching inner rows
181    /// * `residual_filter` - Optional additional filter after key match
182    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        // Build combined schema
216        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        // Pre-allocate row buffer for typical join output size
223        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            // Pre-allocate buffer for typical number of matches (small)
246            row_id_buffer: Vec::with_capacity(16),
247            // Pre-allocate row buffer to avoid per-row allocation
248            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    /// Set projection pushdown configuration.
312    ///
313    /// When set, the operator creates projected rows directly during the combine step
314    /// instead of creating full combined rows. This avoids creating large intermediate
315    /// rows that would be immediately projected down to fewer columns.
316    ///
317    /// # Arguments
318    /// * `columns` - Column sources in SELECT order (preserves original column ordering)
319    /// * `projected_schema` - Schema for the projected output
320    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    /// Create a NULL row for the inner side.
362    /// Creates the active inner row shape: projected when safe, full-width otherwise.
363    #[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    /// Combine outer and inner rows into reusable buffer (owns both).
374    /// OPTIMIZATION: Moves both outer and inner values without cloning.
375    /// Use when outer row is no longer needed (last match for this outer).
376    ///
377    /// Projection indices are validated at query planning time against the
378    /// table schemas. Checked indexing keeps a violated contract fail-closed.
379    #[inline]
380    fn combine_owned_into_buffer(&mut self, outer: Row, inner: Row) {
381        match &self.projection {
382            Some(proj) => {
383                // Fused projection into buffer - columns are in SELECT order
384                self.row_buffer.clear();
385                self.row_buffer.reserve(proj.columns.len());
386
387                // OPTIMIZATION: Move values instead of cloning (we own both rows)
388                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                // Combine both rows into buffer
406                self.row_buffer.combine_into_owned(outer, inner);
407            }
408        }
409    }
410
411    /// Take the combined row from the buffer, preserving buffer capacity.
412    /// OPTIMIZATION: Uses take_and_clear to keep buffer capacity for next iteration,
413    /// avoiding reallocation in combine_into_owned.
414    #[inline]
415    fn take_from_buffer(&mut self) -> Row {
416        self.row_buffer.take_and_clear()
417    }
418
419    /// Create combined row directly (for cases where buffer can't be used).
420    /// When projection is set, creates projected row directly.
421    ///
422    /// Projection indices are validated at query planning time against the
423    /// table schemas. Checked indexing keeps a violated contract fail-closed.
424    #[inline]
425    fn create_combined_row(&self, outer: &Row, inner: Row) -> Row {
426        match &self.projection {
427            Some(proj) => {
428                // Fused projection: create only the columns we need, in SELECT order
429                let mut values: CompactVec<Value> = CompactVec::with_capacity(proj.columns.len());
430
431                // We need to move from inner but clone from outer (we don't own outer)
432                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    /// Look up matching inner rows for the current outer row.
453    /// Uses internal buffers to avoid allocations.
454    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        // Clear buffers for reuse
458        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        // Keep the direct physical lookup for the common committed-only path.
479        // A table-level lookup is required only when the current transaction
480        // owns an unpublished delta, for segmented cold postings, or when a
481        // non-INTEGER primary key cannot be mapped directly to row_id.
482        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        // Fetch matching rows directly into inner_rows buffer. When the join has
523        // no residual ON predicate, projection can cross the row-id lookup
524        // boundary safely: the index lookup already proved the key match and the
525        // output projection is remapped to the projected inner row shape.
526        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        // Transaction-local secondary-index changes are not published to the
538        // shared index before commit. The table lookup therefore admits the
539        // complete local delta as candidates. Recheck the authoritative row
540        // value here to remove unrelated inserts, deletes and old-key entries.
541        // The projected lookup shape always retains this key column.
542        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    /// Advance to the next outer row and lookup matching inner rows.
559    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                // Get the join key value from the outer row
567                let key_value = match outer_row.get(self.outer_key_idx) {
568                    Some(v) if !v.is_null() => v.clone(),
569                    _ => {
570                        // NULL key - no match possible (NULL != NULL in SQL)
571                        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                // Lookup matching inner rows
581                self.lookup_inner_rows(&key_value)?;
582
583                self.current_outer_row = Some(outer_row);
584                // self.current_inner_rows is already populated by lookup_inner_rows
585                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        // Get first outer row
622        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            // Check if outer is exhausted
643            if self.outer_exhausted {
644                return Ok(None);
645            }
646
647            // Ensure we have an outer row
648            if self.current_outer_row.is_none() && !self.advance_outer()? {
649                return Ok(None);
650            }
651
652            // Try to find a match in current inner rows
653            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                // Check filter with borrowed references first
662                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                    // OPTIMIZATION: Check if this is the last inner row to check
677                    // If so, we can take ownership of outer_row and move its values
678                    let is_last_inner = self.current_inner_idx >= inner_len;
679                    if is_last_inner {
680                        // No more inner rows - take ownership of outer and move values
681                        let outer_row = self.current_outer_row.take().unwrap();
682                        // Pre-advance to next outer for next call
683                        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                        // More inner rows to check - create combined row directly
689                        // (Can't use buffer here due to borrow of current_outer_row)
690                        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            // Exhausted inner rows for current outer row
699            // Handle LEFT OUTER: emit outer row with NULLs if no match
700            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                // Use buffer-based combine since we own outer_row
705                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            // Move to next outer row
711            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        // Rough estimate based on outer side
730        let outer_est = self.outer.estimated_rows()?;
731        Some(match self.join_type {
732            JoinType::Inner => outer_est, // Assume most outer rows match
733            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
751// Align the logical lookup window with one artifact-backed row group. A 256-row outer
752// window made a 10k/100k fan-out reopen the table fetch cache and decode the
753// same immutable column blocks tens or hundreds of times. The byte ceiling
754// below keeps wide rows bounded independently of this cardinality ceiling.
755// Permit two physical row groups when the retained graph remains inside the
756// explicit byte budget. A 64K row-only ceiling split a common 100K fan-out
757// into two batches even though its projected/deferred graph was small enough;
758// the next unique edge then fetched the same dictionary keys twice. The byte
759// guard remains authoritative for wide values, so this does not turn the
760// operator into an unbounded collector.
761const BATCH_INDEX_NL_OUTER_ROWS: usize = 131_072;
762const BATCH_INDEX_NL_OUTER_BYTES: usize = 64 * 1024 * 1024;
763
764/// Bounded batch Index-NL operator.
765///
766/// A fixed-size outer chunk is collected, lookup keys are deduplicated, and
767/// matching inner rows are fetched once for the whole chunk. Results remain
768/// streamed in outer-row order; the operator never retains the complete outer
769/// input or complete join output.
770pub 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    /// Create a new batch index nested loop join operator.
820    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    /// Set projection pushdown configuration.
916    ///
917    /// When set, the operator creates projected rows directly during the combine step
918    /// instead of creating full combined rows. This avoids creating large intermediate
919    /// rows that would be immediately projected down to fewer columns.
920    ///
921    /// # Arguments
922    /// * `columns` - Column sources in SELECT order (preserves original column ordering)
923    /// * `projected_schema` - Schema for the projected output
924    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                    // Probe the set by reference first. Repeated UUID/reference
1087                    // keys dominate fan-out joins; cloning before dedup paid an
1088                    // Arc operation for every outer row instead of every
1089                    // distinct lookup key.
1090                    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}