use crate::{
db::{
data::{DataRow, DecodedDataStoreKey, RawRow, StoreVisit},
executor::{
BoundedOrderWindow, DataRowOrderWindow, OrderedKeyStreamBox, PendingOrderRows,
ScalarContinuationContext, begin_production_scalar_page_unit,
exact_output_key_count_hint, finish_production_scalar_page_unit,
key_stream_budget_is_redundant, measure_execution_stats_phase,
production_scalar_page_work_is_active, record_key_stream_micros,
record_key_stream_yield,
route::LoadOrderRouteMode,
terminal::page::{
KernelRow, KernelRowOrderWindow, KernelRowScanStrategy, RetainedSlotLayout,
ScalarRowRuntimeHandle,
},
},
predicate::MissingRowPolicy,
query::plan::{EffectiveRuntimeFilterProgram, ResolvedOrder},
},
error::InternalError,
};
use icydb_diagnostic_code::DiagnosticExecutionBudgetResource;
#[cfg(feature = "diagnostics")]
use super::metrics::{
measure_direct_data_row_phase, record_direct_data_row_key_stream_local_instructions,
record_direct_data_row_peak_retained_backing_bytes,
record_direct_data_row_peak_retained_candidates,
record_direct_data_row_row_read_local_instructions,
};
#[cfg(feature = "diagnostics")]
use super::metrics::{
measure_kernel_row_phase, record_kernel_retained_slot_layout,
record_kernel_row_key_stream_local_instructions, record_kernel_row_peak_retained_backing_bytes,
record_kernel_row_peak_retained_candidates, record_kernel_row_row_read_local_instructions,
record_kernel_row_scan_local_instructions,
};
pub(super) struct RowScanResult<T> {
pub(super) rows: Vec<T>,
pub(super) rows_scanned: usize,
pub(super) rows_matched: usize,
}
pub(super) struct DataRowOrderScanResult<'a> {
pub(super) window: DataRowOrderWindow<'a>,
pub(super) rows_scanned: usize,
pub(super) rows_matched: usize,
}
#[derive(Clone, Copy)]
struct KernelRowScanBounds<'a> {
row_keep_cap: Option<usize>,
row_skip_count: usize,
order_window: Option<KernelRowOrderWindow<'a>>,
}
impl<'a> KernelRowScanBounds<'a> {
const fn new(
row_keep_cap: Option<usize>,
row_skip_count: usize,
order_window: Option<KernelRowOrderWindow<'a>>,
) -> Self {
Self {
row_keep_cap,
row_skip_count,
order_window,
}
}
}
pub(super) struct ScalarPageKernelRequest<'a, 'r> {
pub(super) key_stream: &'a mut OrderedKeyStreamBox,
pub(super) scan_budget_hint: Option<usize>,
pub(super) row_keep_cap: Option<usize>,
pub(super) order_window: Option<KernelRowOrderWindow<'a>>,
pub(super) load_order_route_mode: LoadOrderRouteMode,
pub(super) consistency: MissingRowPolicy,
pub(super) scan_strategy: KernelRowScanStrategy<'a>,
pub(super) continuation: ScalarContinuationContext,
pub(super) row_runtime: &'r mut ScalarRowRuntimeHandle<'a>,
}
pub(in crate::db::executor) struct KernelRowScanRequest<'a, 'r> {
pub(in crate::db::executor) key_stream: &'a mut OrderedKeyStreamBox,
pub(in crate::db::executor) scan_budget_hint: Option<usize>,
pub(in crate::db::executor) consistency: MissingRowPolicy,
pub(in crate::db::executor) scan_strategy: KernelRowScanStrategy<'a>,
pub(in crate::db::executor) row_keep_cap: Option<usize>,
pub(in crate::db::executor) row_skip_count: usize,
pub(in crate::db::executor) order_window: Option<KernelRowOrderWindow<'a>>,
pub(in crate::db::executor) row_runtime: &'r mut ScalarRowRuntimeHandle<'a>,
}
pub(in crate::db::executor) fn execute_kernel_row_scan(
request: KernelRowScanRequest<'_, '_>,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError> {
#[cfg(feature = "diagnostics")]
{
let (scan_local_instructions, result) =
measure_kernel_row_phase(|| execute_kernel_row_scan_inner(request));
record_kernel_row_scan_local_instructions(scan_local_instructions);
let result = result?;
record_kernel_row_peak_retained_candidates(result.0.retained_count());
record_kernel_row_peak_retained_backing_bytes(result.0.retained_backing_bytes());
Ok(result)
}
#[cfg(not(feature = "diagnostics"))]
execute_kernel_row_scan_inner(request)
}
#[expect(clippy::too_many_lines)]
fn execute_kernel_row_scan_inner(
request: KernelRowScanRequest<'_, '_>,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError> {
let KernelRowScanRequest {
key_stream,
scan_budget_hint,
consistency,
scan_strategy,
row_keep_cap,
row_skip_count,
order_window,
row_runtime,
} = request;
let scan_bounds = KernelRowScanBounds::new(row_keep_cap, row_skip_count, order_window);
match scan_strategy {
KernelRowScanStrategy::DataRows => {
execute_scalar_read_loop(key_stream, scan_budget_hint, |key_stream| {
scan_data_rows_only_into_kernel(key_stream, consistency, scan_bounds, row_runtime)
})
}
KernelRowScanStrategy::DataRowsFiltered { filter_program } => {
execute_scalar_read_loop(key_stream, scan_budget_hint, |key_stream| {
scan_data_rows_only_into_kernel_with_filter_program(
key_stream,
consistency,
filter_program,
scan_bounds,
row_runtime,
)
})
}
KernelRowScanStrategy::RetainedFullRows {
retained_slot_layout,
} => {
#[cfg(feature = "diagnostics")]
record_kernel_retained_slot_layout(retained_slot_layout);
execute_retained_kernel_scan(
key_stream,
scan_budget_hint,
Some(retained_slot_layout),
|key_stream, retained_slot_layout| {
scan_full_retained_rows_into_kernel(
key_stream,
consistency,
retained_slot_layout,
scan_bounds,
row_runtime,
)
},
)
}
KernelRowScanStrategy::RetainedFullRowsFiltered {
filter_program,
retained_slot_layout,
} => {
#[cfg(feature = "diagnostics")]
record_kernel_retained_slot_layout(retained_slot_layout);
execute_retained_kernel_scan(
key_stream,
scan_budget_hint,
Some(retained_slot_layout),
|key_stream, retained_slot_layout| {
scan_full_retained_rows_into_kernel_with_filter_program(
key_stream,
consistency,
filter_program,
retained_slot_layout,
scan_bounds,
row_runtime,
)
},
)
}
KernelRowScanStrategy::SlotOnlyRows {
retained_slot_layout,
} => {
#[cfg(feature = "diagnostics")]
record_kernel_retained_slot_layout(retained_slot_layout);
execute_retained_kernel_scan(
key_stream,
scan_budget_hint,
Some(retained_slot_layout),
|key_stream, retained_slot_layout| {
scan_slot_rows_into_kernel(
key_stream,
consistency,
retained_slot_layout,
scan_bounds,
row_runtime,
)
},
)
}
KernelRowScanStrategy::SlotOnlyRowsFiltered {
filter_program,
retained_slot_layout,
} => {
#[cfg(feature = "diagnostics")]
record_kernel_retained_slot_layout(retained_slot_layout);
execute_retained_kernel_scan(
key_stream,
scan_budget_hint,
Some(retained_slot_layout),
|key_stream, retained_slot_layout| {
scan_slot_rows_into_kernel_with_filter_program(
key_stream,
consistency,
filter_program,
retained_slot_layout,
scan_bounds,
row_runtime,
)
},
)
}
}
}
fn execute_retained_kernel_scan(
key_stream: &mut OrderedKeyStreamBox,
scan_budget_hint: Option<usize>,
retained_slot_layout: Option<&RetainedSlotLayout>,
mut scan_rows: impl FnMut(
&mut OrderedKeyStreamBox,
&RetainedSlotLayout,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError>,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError> {
let retained_slot_layout =
retained_slot_layout.ok_or_else(InternalError::query_executor_invariant)?;
execute_scalar_read_loop(key_stream, scan_budget_hint, |key_stream| {
scan_rows(key_stream, retained_slot_layout)
})
}
pub(super) fn execute_scalar_page_kernel_dyn(
request: ScalarPageKernelRequest<'_, '_>,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError> {
let ScalarPageKernelRequest {
key_stream,
scan_budget_hint,
row_keep_cap,
order_window,
load_order_route_mode,
consistency,
scan_strategy,
continuation,
row_runtime,
} = request;
continuation.validate_load_scan_budget_hint(scan_budget_hint, load_order_route_mode)?;
execute_kernel_row_scan(KernelRowScanRequest {
key_stream,
scan_budget_hint,
consistency,
scan_strategy,
row_keep_cap,
row_skip_count: 0,
order_window,
row_runtime,
})
}
fn execute_scalar_read_loop<T>(
key_stream: &mut OrderedKeyStreamBox,
scan_budget_hint: Option<usize>,
mut scan_rows: impl FnMut(&mut OrderedKeyStreamBox) -> Result<T, InternalError>,
) -> Result<T, InternalError> {
if let Some(scan_budget) = scan_budget_hint
&& !key_stream_budget_is_redundant(key_stream, scan_budget)
{
let inner = std::mem::replace(key_stream, OrderedKeyStreamBox::empty());
*key_stream = OrderedKeyStreamBox::budgeted(inner, scan_budget);
return scan_rows(key_stream);
}
scan_rows(key_stream)
}
fn scan_kernel_rows_with(
key_stream: &mut OrderedKeyStreamBox,
bounds: KernelRowScanBounds<'_>,
mut read_row: impl FnMut(DecodedDataStoreKey) -> Result<Option<KernelRow>, InternalError>,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError> {
if let Some(order_window) = bounds.order_window {
return scan_kernel_rows_with_bounded_order_window(
key_stream,
bounds,
order_window,
read_row,
);
}
let result = scan_rows_with(
key_stream,
bounds.row_keep_cap,
bounds.row_skip_count,
next_kernel_scan_key,
|_key_stream, key| read_kernel_scan_row(key, &mut read_row),
)?;
Ok((PendingOrderRows::plain(result.rows), result.rows_scanned))
}
fn scan_kernel_rows_with_bounded_order_window(
key_stream: &mut OrderedKeyStreamBox,
bounds: KernelRowScanBounds<'_>,
order_window: KernelRowOrderWindow<'_>,
mut read_row: impl FnMut(DecodedDataStoreKey) -> Result<Option<KernelRow>, InternalError>,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError> {
if bounds.row_keep_cap.is_some() || bounds.row_skip_count != 0 {
return Err(InternalError::query_executor_invariant());
}
if order_window.keep_count == 0 {
return Ok((PendingOrderRows::plain(Vec::new()), 0));
}
let mut rows_scanned = 0usize;
let mut window = BoundedOrderWindow::new(order_window.keep_count, order_window.resolved_order);
loop {
let page_unit = begin_scan_page_unit(key_stream)?;
if matches!(page_unit, ScanPageUnit::EnvelopeFull) {
break;
}
let key = next_kernel_scan_key(key_stream)?;
let Some(key) = key else {
finish_scan_page_unit(page_unit)?;
break;
};
record_key_stream_yield();
rows_scanned = rows_scanned.saturating_add(1);
let row = read_kernel_scan_row(key, &mut read_row)?;
finish_scan_page_unit(page_unit)?;
let Some(row) = row else {
continue;
};
if !row.has_materialized_slots() {
return Err(InternalError::query_executor_invariant());
}
window.push(row)?;
}
#[cfg(feature = "diagnostics")]
record_kernel_row_peak_retained_backing_bytes(window.peak_retained_backing_bytes());
Ok((window.into_pending_rows(), rows_scanned))
}
fn try_scan_borrowed_primary_rows_with_bounded_order_window(
key_stream: &mut OrderedKeyStreamBox,
bounds: KernelRowScanBounds<'_>,
mut read_row: impl FnMut(&DecodedDataStoreKey, &RawRow) -> Result<Option<KernelRow>, InternalError>,
) -> Result<Option<(PendingOrderRows<KernelRow>, usize)>, InternalError> {
let Some(order_window) = bounds.order_window else {
return Ok(None);
};
if bounds.row_keep_cap.is_some() || bounds.row_skip_count != 0 {
return Err(InternalError::query_executor_invariant());
}
if order_window.keep_count == 0 {
return Ok(Some((PendingOrderRows::plain(Vec::new()), 0)));
}
let rows_scanned = std::cell::Cell::new(0usize);
let envelope_stopped = std::cell::Cell::new(false);
let active_unit = std::cell::Cell::new(ScanPageUnit::Untracked);
let mut window = BoundedOrderWindow::new(order_window.keep_count, order_window.resolved_order);
let mut begin_row = || {
let page_unit = begin_direct_row_scan_page_unit()?;
if matches!(page_unit, ScanPageUnit::EnvelopeFull) {
envelope_stopped.set(true);
return Ok(false);
}
active_unit.set(page_unit);
Ok(true)
};
let mut visit_row = |key: DecodedDataStoreKey, row: &RawRow| {
record_key_stream_yield();
rows_scanned.set(rows_scanned.get().saturating_add(1));
let decoded = read_borrowed_kernel_scan_row(&key, row, &mut read_row)?;
finish_scan_page_unit(active_unit.replace(ScanPageUnit::Untracked))?;
if let Some(decoded) = decoded {
if !decoded.has_materialized_slots() {
return Err(InternalError::query_executor_invariant());
}
window.push(decoded)?;
}
Ok(StoreVisit::Continue)
};
let Some(()) = key_stream.try_visit_primary_rows_direct(&mut begin_row, &mut visit_row)? else {
return Ok(None);
};
if envelope_stopped.get() {
let observed = u64::try_from(rows_scanned.get()).unwrap_or(u64::MAX);
return Err(InternalError::page_unit_too_large(
DiagnosticExecutionBudgetResource::KeyIndexEntriesVisited,
observed,
observed.saturating_add(1),
));
}
#[cfg(feature = "diagnostics")]
record_kernel_row_peak_retained_backing_bytes(window.peak_retained_backing_bytes());
Ok(Some((window.into_pending_rows(), rows_scanned.get())))
}
fn scan_rows_with<T>(
key_stream: &mut OrderedKeyStreamBox,
row_keep_cap: Option<usize>,
row_skip_count: usize,
mut next_key: impl FnMut(
&mut OrderedKeyStreamBox,
) -> Result<Option<DecodedDataStoreKey>, InternalError>,
mut read_row: impl FnMut(
&mut OrderedKeyStreamBox,
DecodedDataStoreKey,
) -> Result<Option<T>, InternalError>,
) -> Result<RowScanResult<T>, InternalError> {
if row_keep_cap == Some(0) {
return Ok(RowScanResult {
rows: Vec::new(),
rows_scanned: 0,
rows_matched: 0,
});
}
let mut rows_scanned = 0usize;
let staged_capacity = staged_row_capacity(key_stream, row_keep_cap, row_skip_count);
let mut rows = Vec::with_capacity(staged_capacity);
let mut rows_matched = 0usize;
let Some(row_keep_cap) = row_keep_cap else {
loop {
let page_unit = begin_scan_page_unit(key_stream)?;
if matches!(page_unit, ScanPageUnit::EnvelopeFull) {
break;
}
let key = next_key(key_stream)?;
let Some(key) = key else {
finish_scan_page_unit(page_unit)?;
break;
};
record_key_stream_yield();
rows_scanned = rows_scanned.saturating_add(1);
let row = read_row(key_stream, key)?;
finish_scan_page_unit(page_unit)?;
let Some(row) = row else {
continue;
};
retain_scanned_row(row, row_skip_count, rows_matched, &mut rows);
rows_matched = rows_matched.saturating_add(1);
}
return Ok(RowScanResult {
rows,
rows_scanned,
rows_matched,
});
};
loop {
let page_unit = begin_scan_page_unit(key_stream)?;
if matches!(page_unit, ScanPageUnit::EnvelopeFull) {
break;
}
let key = next_key(key_stream)?;
let Some(key) = key else {
finish_scan_page_unit(page_unit)?;
break;
};
record_key_stream_yield();
rows_scanned = rows_scanned.saturating_add(1);
let row = read_row(key_stream, key)?;
finish_scan_page_unit(page_unit)?;
let Some(row) = row else {
continue;
};
retain_scanned_row(row, row_skip_count, rows_matched, &mut rows);
rows_matched = rows_matched.saturating_add(1);
if rows_matched >= row_keep_cap {
break;
}
}
Ok(RowScanResult {
rows,
rows_scanned,
rows_matched,
})
}
#[derive(Clone, Copy)]
enum ScanPageUnit {
Untracked,
Started,
EnvelopeFull,
}
fn begin_scan_page_unit(key_stream: &OrderedKeyStreamBox) -> Result<ScanPageUnit, InternalError> {
let Some(access_entry_bound) = key_stream.page_access_entry_bound() else {
if production_scalar_page_work_is_active()? {
return Err(InternalError::query_executor_invariant());
}
return Ok(ScanPageUnit::Untracked);
};
if begin_production_scalar_page_unit(access_entry_bound)? {
Ok(ScanPageUnit::Started)
} else {
Ok(ScanPageUnit::EnvelopeFull)
}
}
fn begin_direct_row_scan_page_unit() -> Result<ScanPageUnit, InternalError> {
if !production_scalar_page_work_is_active()? {
return Ok(ScanPageUnit::Untracked);
}
if begin_production_scalar_page_unit(1)? {
Ok(ScanPageUnit::Started)
} else {
Ok(ScanPageUnit::EnvelopeFull)
}
}
fn finish_scan_page_unit(unit: ScanPageUnit) -> Result<(), InternalError> {
if matches!(unit, ScanPageUnit::Started) {
finish_production_scalar_page_unit()?;
}
Ok(())
}
fn next_kernel_scan_key(
key_stream: &mut OrderedKeyStreamBox,
) -> Result<Option<DecodedDataStoreKey>, InternalError> {
#[cfg(feature = "diagnostics")]
let ((key_stream_local_instructions, next_key), key_stream_micros) =
measure_execution_stats_phase(|| measure_kernel_row_phase(|| key_stream.next_key()));
#[cfg(not(feature = "diagnostics"))]
let (next_key, key_stream_micros) = measure_execution_stats_phase(|| key_stream.next_key());
record_key_stream_micros(key_stream_micros);
#[cfg(feature = "diagnostics")]
record_kernel_row_key_stream_local_instructions(key_stream_local_instructions);
next_key
}
fn read_kernel_scan_row(
key: DecodedDataStoreKey,
read_row: &mut impl FnMut(DecodedDataStoreKey) -> Result<Option<KernelRow>, InternalError>,
) -> Result<Option<KernelRow>, InternalError> {
#[cfg(feature = "diagnostics")]
let (row_read_local_instructions, row) = measure_kernel_row_phase(|| read_row(key));
#[cfg(not(feature = "diagnostics"))]
let row = read_row(key);
#[cfg(feature = "diagnostics")]
record_kernel_row_row_read_local_instructions(row_read_local_instructions);
row
}
fn read_borrowed_kernel_scan_row(
key: &DecodedDataStoreKey,
row: &RawRow,
read_row: &mut impl FnMut(&DecodedDataStoreKey, &RawRow) -> Result<Option<KernelRow>, InternalError>,
) -> Result<Option<KernelRow>, InternalError> {
#[cfg(feature = "diagnostics")]
let (row_read_local_instructions, decoded) = measure_kernel_row_phase(|| read_row(key, row));
#[cfg(not(feature = "diagnostics"))]
let decoded = read_row(key, row);
#[cfg(feature = "diagnostics")]
record_kernel_row_row_read_local_instructions(row_read_local_instructions);
decoded
}
fn staged_row_capacity(
key_stream: &OrderedKeyStreamBox,
row_keep_cap: Option<usize>,
row_skip_count: usize,
) -> usize {
row_keep_cap
.map(|row_keep_cap| row_keep_cap.saturating_sub(row_skip_count))
.or_else(|| {
exact_output_key_count_hint(key_stream, row_keep_cap)
.map(|hint| hint.saturating_sub(row_skip_count))
})
.unwrap_or(0)
}
fn retain_scanned_row<T>(row: T, row_skip_count: usize, rows_matched: usize, rows: &mut Vec<T>) {
if rows_matched >= row_skip_count {
rows.push(row);
}
}
fn scan_data_rows_direct(
key_stream: &mut OrderedKeyStreamBox,
consistency: MissingRowPolicy,
row_keep_cap: Option<usize>,
row_skip_count: usize,
row_runtime: &ScalarRowRuntimeHandle<'_>,
) -> Result<RowScanResult<DataRow>, InternalError> {
scan_data_rows_direct_with_reader(key_stream, row_keep_cap, row_skip_count, |key| {
row_runtime.read_data_row(consistency, key)
})
}
fn scan_data_rows_direct_with_reader(
key_stream: &mut OrderedKeyStreamBox,
row_keep_cap: Option<usize>,
row_skip_count: usize,
mut read_data_row: impl FnMut(DecodedDataStoreKey) -> Result<Option<DataRow>, InternalError>,
) -> Result<RowScanResult<DataRow>, InternalError> {
scan_rows_with(
key_stream,
row_keep_cap,
row_skip_count,
next_direct_data_row_scan_key,
|_key_stream, key| read_direct_data_row_scan_row(key, &mut read_data_row),
)
}
fn next_direct_data_row_scan_key(
key_stream: &mut OrderedKeyStreamBox,
) -> Result<Option<DecodedDataStoreKey>, InternalError> {
#[cfg(feature = "diagnostics")]
let ((key_stream_local_instructions, read_result), key_stream_micros) =
measure_execution_stats_phase(|| measure_direct_data_row_phase(|| key_stream.next_key()));
#[cfg(not(feature = "diagnostics"))]
let (read_result, key_stream_micros) = measure_execution_stats_phase(|| key_stream.next_key());
record_key_stream_micros(key_stream_micros);
#[cfg(feature = "diagnostics")]
record_direct_data_row_key_stream_local_instructions(key_stream_local_instructions);
read_result
}
fn read_direct_data_row_scan_row(
key: DecodedDataStoreKey,
read_data_row: &mut impl FnMut(DecodedDataStoreKey) -> Result<Option<DataRow>, InternalError>,
) -> Result<Option<DataRow>, InternalError> {
#[cfg(feature = "diagnostics")]
let (row_read_local_instructions, row_read_result) =
measure_direct_data_row_phase(|| read_data_row(key));
#[cfg(not(feature = "diagnostics"))]
let row_read_result = read_data_row(key);
#[cfg(feature = "diagnostics")]
record_direct_data_row_row_read_local_instructions(row_read_local_instructions);
row_read_result
}
fn scan_data_rows_direct_with_filter_program(
key_stream: &mut OrderedKeyStreamBox,
consistency: MissingRowPolicy,
row_keep_cap: Option<usize>,
row_skip_count: usize,
row_runtime: &ScalarRowRuntimeHandle<'_>,
filter_program: &EffectiveRuntimeFilterProgram,
) -> Result<RowScanResult<DataRow>, InternalError> {
scan_data_rows_direct_with_reader(key_stream, row_keep_cap, row_skip_count, |key| {
row_runtime.read_data_row_with_filter_program(consistency, key, filter_program)
})
}
pub(super) fn scan_materialized_order_direct_data_rows<'a>(
key_stream: &mut OrderedKeyStreamBox,
scan_budget_hint: Option<usize>,
consistency: MissingRowPolicy,
row_runtime: &ScalarRowRuntimeHandle<'_>,
residual_filter_program: Option<&EffectiveRuntimeFilterProgram>,
resolved_order: &'a ResolvedOrder,
keep_count: Option<usize>,
) -> Result<DataRowOrderScanResult<'a>, InternalError> {
execute_scalar_read_loop(key_stream, scan_budget_hint, |key_stream| {
let mut rows_scanned = 0usize;
let mut rows_matched = 0usize;
let mut window =
DataRowOrderWindow::new(row_runtime.row_layout(), resolved_order, keep_count);
if keep_count == Some(0) {
return Ok(DataRowOrderScanResult {
window,
rows_scanned,
rows_matched,
});
}
while let Some(key) = next_direct_data_row_scan_key(key_stream)? {
record_key_stream_yield();
rows_scanned = rows_scanned.saturating_add(1);
let row =
read_direct_data_row_scan_row(key, &mut |key| match residual_filter_program {
None => row_runtime.read_data_row(consistency, key),
Some(filter_program) => row_runtime.read_data_row_with_filter_program(
consistency,
key,
filter_program,
),
})?;
let Some(row) = row else {
continue;
};
window.push(row)?;
rows_matched = rows_matched.saturating_add(1);
}
#[cfg(feature = "diagnostics")]
{
record_direct_data_row_peak_retained_candidates(window.retained_count());
record_direct_data_row_peak_retained_backing_bytes(
window.peak_retained_backing_bytes(),
);
}
Ok(DataRowOrderScanResult {
window,
rows_scanned,
rows_matched,
})
})
}
pub(super) fn scan_direct_data_rows_with_residual_policy(
key_stream: &mut OrderedKeyStreamBox,
scan_budget_hint: Option<usize>,
row_keep_cap: Option<usize>,
row_skip_count: usize,
consistency: MissingRowPolicy,
row_runtime: &ScalarRowRuntimeHandle<'_>,
residual_filter_program: Option<&EffectiveRuntimeFilterProgram>,
) -> Result<RowScanResult<DataRow>, InternalError> {
if row_keep_cap == Some(0) {
return Ok(RowScanResult {
rows: Vec::new(),
rows_scanned: 0,
rows_matched: 0,
});
}
execute_scalar_read_loop(key_stream, scan_budget_hint, |key_stream| {
match residual_filter_program {
None => scan_data_rows_direct(
key_stream,
consistency,
row_keep_cap,
row_skip_count,
row_runtime,
),
Some(filter_program) => scan_data_rows_direct_with_filter_program(
key_stream,
consistency,
row_keep_cap,
row_skip_count,
row_runtime,
filter_program,
),
}
})
}
fn scan_data_rows_only_into_kernel(
key_stream: &mut OrderedKeyStreamBox,
consistency: MissingRowPolicy,
bounds: KernelRowScanBounds<'_>,
row_runtime: &ScalarRowRuntimeHandle<'_>,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError> {
scan_kernel_rows_with(key_stream, bounds, |key| {
row_runtime.read_data_row_only(consistency, key)
})
}
fn scan_data_rows_only_into_kernel_with_filter_program(
key_stream: &mut OrderedKeyStreamBox,
consistency: MissingRowPolicy,
filter_program: &EffectiveRuntimeFilterProgram,
bounds: KernelRowScanBounds<'_>,
row_runtime: &ScalarRowRuntimeHandle<'_>,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError> {
scan_kernel_rows_with(key_stream, bounds, |key| {
row_runtime
.read_data_row_with_filter_program(consistency, key, filter_program)
.map(|row| row.map(KernelRow::new_data_row_only))
})
}
fn scan_full_retained_rows_into_kernel(
key_stream: &mut OrderedKeyStreamBox,
consistency: MissingRowPolicy,
retained_slot_layout: &RetainedSlotLayout,
bounds: KernelRowScanBounds<'_>,
row_runtime: &ScalarRowRuntimeHandle<'_>,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError> {
scan_full_retained_rows_into_kernel_with_reader(key_stream, bounds, |key| {
row_runtime.read_full_row_retained(consistency, key, retained_slot_layout)
})
}
fn scan_full_retained_rows_into_kernel_with_reader(
key_stream: &mut OrderedKeyStreamBox,
bounds: KernelRowScanBounds<'_>,
read_row: impl FnMut(DecodedDataStoreKey) -> Result<Option<KernelRow>, InternalError>,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError> {
scan_kernel_rows_with(key_stream, bounds, read_row)
}
fn scan_full_retained_rows_into_kernel_with_filter_program(
key_stream: &mut OrderedKeyStreamBox,
consistency: MissingRowPolicy,
filter_program: &EffectiveRuntimeFilterProgram,
retained_slot_layout: &RetainedSlotLayout,
bounds: KernelRowScanBounds<'_>,
row_runtime: &ScalarRowRuntimeHandle<'_>,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError> {
scan_full_retained_rows_into_kernel_with_reader(key_stream, bounds, |key| {
row_runtime.read_full_row_retained_with_filter_program(
consistency,
key,
filter_program,
retained_slot_layout,
)
})
}
fn scan_slot_rows_into_kernel(
key_stream: &mut OrderedKeyStreamBox,
consistency: MissingRowPolicy,
retained_slot_layout: &RetainedSlotLayout,
bounds: KernelRowScanBounds<'_>,
row_runtime: &ScalarRowRuntimeHandle<'_>,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError> {
if let Some(rows) =
try_scan_borrowed_primary_rows_with_bounded_order_window(key_stream, bounds, |key, row| {
row_runtime
.read_borrowed_slot_only(key, row, retained_slot_layout)
.map(Some)
})?
{
return Ok(rows);
}
scan_slot_rows_into_kernel_with_reader(key_stream, bounds, |key| {
row_runtime.read_slot_only(consistency, &key, retained_slot_layout)
})
}
fn scan_slot_rows_into_kernel_with_reader(
key_stream: &mut OrderedKeyStreamBox,
bounds: KernelRowScanBounds<'_>,
read_row: impl FnMut(DecodedDataStoreKey) -> Result<Option<KernelRow>, InternalError>,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError> {
scan_kernel_rows_with(key_stream, bounds, read_row)
}
fn scan_slot_rows_into_kernel_with_filter_program(
key_stream: &mut OrderedKeyStreamBox,
consistency: MissingRowPolicy,
filter_program: &EffectiveRuntimeFilterProgram,
retained_slot_layout: &RetainedSlotLayout,
bounds: KernelRowScanBounds<'_>,
row_runtime: &ScalarRowRuntimeHandle<'_>,
) -> Result<(PendingOrderRows<KernelRow>, usize), InternalError> {
if let Some(rows) =
try_scan_borrowed_primary_rows_with_bounded_order_window(key_stream, bounds, |key, row| {
row_runtime.read_borrowed_slot_only_with_filter_program(
key,
row,
filter_program,
retained_slot_layout,
)
})?
{
return Ok(rows);
}
scan_slot_rows_into_kernel_with_reader(key_stream, bounds, |key| {
row_runtime.read_slot_only_with_filter_program(
consistency,
&key,
filter_program,
retained_slot_layout,
)
})
}