use crate::{
db::{
access::{AccessPathKind, IndexShapeDetails},
cursor::{
CursorBoundary, CursorBoundarySlot,
effective_keep_count_for_limit as continuation_keep_count_for_limit,
effective_page_offset_for_window as continuation_page_offset_for_window,
},
data::primary_key_value_from_structural_value,
direction::Direction,
executor::{
AccessScanContinuationInput, ContinuationMode, LoweredIndexPrefixSpec,
LoweredIndexRangeSpec, LoweredKey, RouteContinuationPlan,
budget::ExecutionConstructionBudget, planning::route::LoadOrderRouteMode,
route::access_order_satisfied_by_route_mode,
},
index::IndexKey,
query::{
construction::ConstructionBudget,
plan::{AccessPlannedQuery, ContinuationPolicy, DeterministicSecondaryIndexOrderMatch},
},
schema::SchemaInfo,
},
error::InternalError,
value::Value,
};
use std::{ops::Bound, rc::Rc};
#[derive(Clone, Debug, Eq, PartialEq)]
pub(in crate::db) struct ScalarContinuationContext {
cursor_boundary: Option<Rc<CursorBoundary>>,
physical_primary_key_boundary: Option<Rc<CursorBoundary>>,
}
impl ScalarContinuationContext {
#[must_use]
pub(in crate::db) const fn initial() -> Self {
Self {
cursor_boundary: None,
physical_primary_key_boundary: None,
}
}
#[must_use]
pub(in crate::db) fn resumed(cursor_boundary: CursorBoundary) -> Self {
Self {
cursor_boundary: Some(Rc::new(cursor_boundary)),
physical_primary_key_boundary: None,
}
}
#[must_use]
pub(in crate::db) fn resumed_with_primary_progress(
cursor_boundary: Option<CursorBoundary>,
physical_primary_key_boundary: CursorBoundary,
) -> Self {
Self {
cursor_boundary: cursor_boundary.map(Rc::new),
physical_primary_key_boundary: Some(Rc::new(physical_primary_key_boundary)),
}
}
#[must_use]
pub(in crate::db::executor) fn cursor_boundary(&self) -> Option<&CursorBoundary> {
self.cursor_boundary.as_deref()
}
#[must_use]
pub(in crate::db) const fn has_progress(&self) -> bool {
self.cursor_boundary.is_some() || self.physical_primary_key_boundary.is_some()
}
#[must_use]
pub(in crate::db::executor) fn can_bound_ordered_scan(
&self,
plan: &AccessPlannedQuery,
) -> bool {
!self.has_progress()
|| scalar_order_is_primary_key_only(plan)
|| scalar_secondary_index_order(plan).is_some()
}
pub(in crate::db::executor) fn secondary_index_resume_anchor(
&self,
plan: &AccessPlannedQuery,
schema: &SchemaInfo,
prefixes: &[LoweredIndexPrefixSpec],
ranges: &[LoweredIndexRangeSpec],
) -> Result<Option<LoweredKey>, InternalError> {
let Some(boundary) = self.cursor_boundary() else {
return Ok(None);
};
let Some((index, prefix_len)) = scalar_secondary_index_order(plan) else {
return Ok(None);
};
let primary_len = plan.primary_key_names()?.len();
let order = plan
.planner_route_profile()
.secondary_order_contract()
.ok_or_else(InternalError::query_executor_invariant)?;
if boundary.slots.len() != order.non_primary_key_terms().len() + primary_len {
return Err(InternalError::query_executor_invariant());
}
let budget: &dyn ConstructionBudget = &ExecutionConstructionBudget;
let mut values = budget.vec_with_capacity(boundary.slots.len())?;
for slot in &boundary.slots {
let CursorBoundarySlot::Present(value) = slot else {
return Err(InternalError::query_executor_invariant());
};
values.push(value);
}
let primary_values = values
.get(values.len().saturating_sub(primary_len)..)
.ok_or_else(InternalError::query_executor_invariant)?;
let primary_key = match primary_values {
[value] => primary_key_value_from_structural_value(value)?,
_ => primary_key_value_from_structural_value(&Value::List(
budget.copy_slice(primary_values, |value| budget.copy_value(value))?,
))?,
};
let start = secondary_index_template(prefixes, ranges)?;
if start.component_count() != index.key_arity() {
return Err(InternalError::query_executor_invariant());
}
Ok(Some(start.raw_resume_anchor_with_accepted_suffix(
schema,
index.name(),
prefix_len,
&values,
&primary_key,
budget,
)?))
}
#[must_use]
pub(in crate::db::executor) const fn route_continuation_mode(&self) -> ContinuationMode {
if self.has_progress() {
ContinuationMode::CursorBoundary
} else {
ContinuationMode::Initial
}
}
#[must_use]
pub(in crate::db::executor) fn route_continuation_plan(
&self,
plan: &AccessPlannedQuery,
continuation_policy: ContinuationPolicy,
) -> RouteContinuationPlan {
RouteContinuationPlan::from_scalar_access_window_plan(
self.route_continuation_mode(),
continuation_policy,
plan.scalar_access_window_plan(self.has_progress()),
)
}
#[must_use]
pub(in crate::db::executor) fn access_scan_input<'a>(
&'a self,
direction: Direction,
plan: &AccessPlannedQuery,
secondary_index_anchor: Option<&'a LoweredKey>,
) -> AccessScanContinuationInput<'a> {
let primary_key_ordered = scalar_order_is_primary_key_only(plan);
AccessScanContinuationInput::with_primary_key_boundary(
secondary_index_anchor,
direction,
primary_key_ordered
.then_some(
self.physical_primary_key_boundary
.as_deref()
.or_else(|| self.cursor_boundary()),
)
.flatten(),
)
}
pub(in crate::db::executor) fn debug_assert_route_continuation_invariants(
&self,
plan: &AccessPlannedQuery,
route_continuation: RouteContinuationPlan,
) {
debug_assert!(
route_continuation.strict_advance_required_when_applied(),
"route invariant: continuation executions must enforce strict advancement policy",
);
debug_assert_eq!(
route_continuation.effective_offset(),
continuation_page_offset_for_window(plan, self.has_progress()),
"route window effective offset must match logical plan offset semantics",
);
}
#[must_use]
pub(in crate::db::executor) fn keep_count_for_limit_window(
&self,
plan: &AccessPlannedQuery,
limit: u32,
) -> usize {
continuation_keep_count_for_limit(plan, self.has_progress(), limit)
}
pub(in crate::db::executor) fn validate_load_scan_budget_hint(
&self,
scan_budget_hint: Option<usize>,
load_order_route_mode: LoadOrderRouteMode,
) -> Result<(), InternalError> {
if scan_budget_hint.is_some() && self.has_progress() {
return Err(InternalError::query_executor_invariant());
}
if scan_budget_hint.is_some() && !load_order_route_mode.allows_streaming_load() {
return Err(InternalError::query_executor_invariant());
}
Ok(())
}
}
fn scalar_order_is_primary_key_only(plan: &AccessPlannedQuery) -> bool {
let Ok(primary_key_names) = plan.primary_key_names() else {
return false;
};
plan.scalar_plan().order.as_ref().is_some_and(|order| {
order
.primary_key_only_direction_fields(primary_key_names)
.is_some()
})
}
fn scalar_secondary_index_order(plan: &AccessPlannedQuery) -> Option<(IndexShapeDetails, usize)> {
if !access_order_satisfied_by_route_mode(plan) {
return None;
}
let facts = plan.access_shape_facts();
if !matches!(
facts.single_path_facts()?.kind(),
AccessPathKind::IndexPrefix | AccessPathKind::IndexRange | AccessPathKind::IndexMultiLookup
) {
return None;
}
let index = facts
.single_path_index_prefix_details()
.or_else(|| facts.single_path_index_range_details())?;
let contract = plan.planner_route_profile().secondary_order_contract()?;
if contract.non_primary_key_terms().is_empty() {
return None;
}
let prefix_len = match contract.classify_index_key_items(index.key_items(), index.slot_arity())
{
DeterministicSecondaryIndexOrderMatch::Full => 0,
DeterministicSecondaryIndexOrderMatch::Suffix => index.slot_arity(),
DeterministicSecondaryIndexOrderMatch::None => return None,
};
Some((index, prefix_len))
}
fn secondary_index_template(
prefixes: &[LoweredIndexPrefixSpec],
ranges: &[LoweredIndexRangeSpec],
) -> Result<IndexKey, InternalError> {
let (lower, upper) = if let Some(spec) = prefixes.first() {
spec.raw_bounds(&ExecutionConstructionBudget)?
} else if let [spec] = ranges {
(spec.lower(), spec.upper())
} else {
return Err(InternalError::query_executor_invariant());
};
let raw = match lower {
Bound::Included(key) | Bound::Excluded(key) => key,
Bound::Unbounded => match upper {
Bound::Included(key) | Bound::Excluded(key) => key,
Bound::Unbounded => return Err(InternalError::query_executor_invariant()),
},
};
IndexKey::try_from_raw(raw).map_err(|_| InternalError::query_executor_invariant())
}