icydb-core 0.213.35

IcyDB — A schema-first typed query engine and persistence runtime for Internet Computer canisters
Documentation
//! Module: executor::pipeline::runtime::adapter
//! Responsibility: runtime adapters for stream resolution and scalar materialization.
//! Does not own: execution-input DTO construction or planning semantics.
//! Boundary: executes already-assembled execution contracts through runtime owners.

#[cfg(feature = "sql")]
use crate::db::executor::{
    pipeline::contracts::KernelRowsExecutionAttempt,
    terminal::page::materialize_key_stream_into_kernel_rows,
};
use crate::{
    db::{
        access::ExecutableAccessPlan,
        direction::Direction,
        executor::{
            AccessScanContinuationInput, AccessStreamBindings, AccessStreamExecutionPolicy,
            EntityAuthority, ExecutableAccess, ExecutionKernel, LoweredIndexRangeSpec,
            OrderedKeyStreamBox, ScalarContinuationContext,
            pipeline::contracts::{
                CursorEmissionMode, FastPathKeyResult, FastStreamRouteKind, FastStreamRouteRequest,
                KernelPageMaterializationRequest, RowCollectorMaterializationRequest,
                ScalarMaterializationCapabilities, StructuralCursorPage,
            },
            projection::PreparedProjectionContract,
            route::LoadOrderRouteMode,
            scan::execute_fast_stream_route,
            stream::access::TraversalRuntime,
            terminal::page::{
                ScalarRowRuntimeHandle, ScalarRowRuntimeState,
                materialize_key_stream_into_execution_payload,
            },
        },
        index::predicate::IndexPredicateExecution,
        predicate::MissingRowPolicy,
        query::plan::{AccessPlannedQuery, EffectiveRuntimeFilterProgram},
        registry::StoreHandle,
    },
    error::InternalError,
    value::Value,
};

type MaterializedExecutionPayloadResult = (StructuralCursorPage, usize, usize);

///
/// ExecutionMaterializationContract
///
/// ExecutionMaterializationContract captures the execution-input fields shared
/// by the row-collector and kernel-page materialization requests.
/// Runtime materialization consumes this once so the two terminal request
/// shapes do not re-spell predicate/projection/retained-slot wiring.
///

#[derive(Clone, Copy)]
pub(in crate::db::executor) struct ExecutionMaterializationContract<'a> {
    pub(in crate::db::executor) plan: &'a AccessPlannedQuery,
    pub(in crate::db::executor) residual_filter_program: Option<&'a EffectiveRuntimeFilterProgram>,
    pub(in crate::db::executor) scan_budget_hint: Option<usize>,
    pub(in crate::db::executor) load_order_route_mode: LoadOrderRouteMode,
    pub(in crate::db::executor) validate_projection: bool,
    pub(in crate::db::executor) retain_slot_rows: bool,
    pub(in crate::db::executor) retained_slot_layout: Option<&'a RetainedSlotLayout>,
    pub(in crate::db::executor) prepared_projection_validation:
        Option<&'a PreparedProjectionContract>,
}

impl<'a> ExecutionMaterializationContract<'a> {
    // Project the shared predicate/projection contract into one terminal
    // capability bundle without introducing another request DTO.
    const fn capabilities(
        &self,
        cursor_emission: CursorEmissionMode,
    ) -> ScalarMaterializationCapabilities<'a> {
        ScalarMaterializationCapabilities {
            residual_filter_program: self.residual_filter_program,
            validate_projection: self.validate_projection,
            retain_slot_rows: self.retain_slot_rows,
            retained_slot_layout: self.retained_slot_layout,
            prepared_projection_validation: self.prepared_projection_validation,
            cursor_emission,
        }
    }

    // Materialize one resolved scalar key stream through the aligned
    // row-collector or canonical page runtime lane without rebuilding the
    // shared predicate/projection/retained-slot contract twice.
    pub(in crate::db::executor) fn materialize_resolved_execution_stream(
        &self,
        runtime: &'a ExecutionRuntimeAdapter,
        emit_cursor: bool,
        consistency: MissingRowPolicy,
        continuation: ScalarContinuationContext,
        key_stream: &'a mut OrderedKeyStreamBox,
    ) -> Result<MaterializedExecutionPayloadResult, InternalError> {
        runtime.materialize_resolved_execution_stream(
            self,
            emit_cursor,
            consistency,
            continuation,
            key_stream,
        )
    }

    // Materialize one resolved scalar key stream through post-access/window
    // processing while stopping before structural page payload construction.
    #[cfg(feature = "sql")]
    pub(in crate::db::executor) fn materialize_resolved_execution_stream_to_kernel_rows(
        &self,
        runtime: &'a ExecutionRuntimeAdapter,
        consistency: MissingRowPolicy,
        continuation: ScalarContinuationContext,
        key_stream: &'a mut OrderedKeyStreamBox,
    ) -> Result<KernelRowsExecutionAttempt, InternalError> {
        runtime.materialize_resolved_execution_stream_to_kernel_rows(
            self,
            consistency,
            continuation,
            key_stream,
        )
    }

    // Build the cursorless row-collector materialization request from one
    // already-aligned scalar materialization contract.
    const fn row_collector_request(
        &self,
        continuation: ScalarContinuationContext,
        consistency: MissingRowPolicy,
        key_stream: &'a mut OrderedKeyStreamBox,
    ) -> RowCollectorMaterializationRequest<'a> {
        RowCollectorMaterializationRequest {
            plan: self.plan,
            scan_budget_hint: self.scan_budget_hint,
            load_order_route_mode: self.load_order_route_mode,
            continuation,
            cursor_boundary: continuation.cursor_boundary(),
            capabilities: self.capabilities(CursorEmissionMode::Suppress),
            consistency,
            key_stream,
        }
    }
}

///
/// ExecutionRuntimeAdapter
///
/// Typed runtime adapter that captures recovered context plus structural
/// runtime helpers once at the execution boundary and exposes one monomorphic
/// runtime surface to shared executor code.
///

pub(in crate::db::executor) struct ExecutionRuntimeAdapter {
    runtime: TraversalRuntime,
    scalar_row_runtime: Option<ScalarRowRuntimeState>,
}

impl ExecutionRuntimeAdapter {
    /// Build one structural runtime adapter for scalar execution paths.
    pub(in crate::db::executor) fn from_scalar_runtime(
        runtime: TraversalRuntime,
        store: StoreHandle,
        authority: EntityAuthority,
    ) -> Result<Self, InternalError> {
        let row_layout = authority.row_layout()?;

        Ok(Self {
            runtime,
            scalar_row_runtime: Some(ScalarRowRuntimeState::new(store, row_layout)),
        })
    }

    /// Build one stream-only runtime adapter for key-stream resolution paths
    /// that never materialize scalar rows.
    pub(in crate::db::executor) const fn from_stream_runtime(runtime: TraversalRuntime) -> Self {
        Self {
            runtime,
            scalar_row_runtime: None,
        }
    }

    // Require the scalar materialization runtime when the caller enters one
    // scalar-only row materialization path through the shared execution spine.
    fn scalar_row_runtime(&self) -> Result<&ScalarRowRuntimeState, InternalError> {
        self.scalar_row_runtime
            .as_ref()
            .ok_or_else(InternalError::query_executor_invariant)
    }

    // Reuse the adapter-owned scalar row runtime for one materialization call
    // so callers do not each rebuild the same borrowed runtime-handle shell.
    fn with_scalar_row_runtime_handle<'a, T>(
        &'a self,
        run: impl FnOnce(&mut ScalarRowRuntimeHandle<'a>) -> Result<T, InternalError>,
    ) -> Result<T, InternalError> {
        let scalar_row_runtime = self.scalar_row_runtime()?;
        let mut row_runtime = ScalarRowRuntimeHandle::from_borrowed(scalar_row_runtime);

        run(&mut row_runtime)
    }

    // Materialize one resolved scalar key stream through the aligned
    // row-collector or canonical page runtime lane owned by this runtime
    // adapter.
    fn materialize_resolved_execution_stream<'a>(
        &'a self,
        contract: &ExecutionMaterializationContract<'a>,
        emit_cursor: bool,
        consistency: MissingRowPolicy,
        continuation: ScalarContinuationContext,
        key_stream: &'a mut OrderedKeyStreamBox,
    ) -> Result<MaterializedExecutionPayloadResult, InternalError> {
        if !emit_cursor
            && let Some(materialized) = self.try_materialize_load_via_row_collector(
                contract.row_collector_request(continuation, consistency, key_stream),
            )?
        {
            return Ok(materialized);
        }

        self.materialize_key_stream_into_structural_page(
            contract,
            emit_cursor,
            consistency,
            continuation,
            key_stream,
        )
    }

    // Materialize one ordered key stream into post-access scalar kernel rows for
    // aggregate sinks that do not need an outward cursor page.
    #[cfg(feature = "sql")]
    fn materialize_resolved_execution_stream_to_kernel_rows<'a>(
        &'a self,
        contract: &ExecutionMaterializationContract<'a>,
        consistency: MissingRowPolicy,
        continuation: ScalarContinuationContext,
        key_stream: &'a mut OrderedKeyStreamBox,
    ) -> Result<KernelRowsExecutionAttempt, InternalError> {
        self.materialize_key_stream_into_kernel_rows(
            contract,
            consistency,
            continuation,
            key_stream,
        )
    }

    /// Resolve one primary-key fast path when the route is already verified.
    pub(in crate::db::executor) fn try_execute_pk_order_stream(
        &self,
        plan: &AccessPlannedQuery,
        executable_access: ExecutableAccessPlan<'_, Value>,
        direction: Direction,
        physical_fetch_hint: Option<usize>,
    ) -> Result<Option<FastPathKeyResult>, InternalError> {
        execute_fast_stream_route(
            &self.runtime,
            FastStreamRouteKind::PrimaryKey,
            FastStreamRouteRequest::PrimaryKey {
                plan,
                executable_access: &executable_access,
                stream_direction: direction,
                probe_fetch_hint: physical_fetch_hint,
            },
        )
    }

    /// Resolve one verified secondary-prefix fast path.
    pub(in crate::db::executor) fn try_execute_secondary_index_order_stream(
        &self,
        plan: &AccessPlannedQuery,
        executable_access: ExecutableAccessPlan<'_, Value>,
        bindings: AccessStreamBindings<'_>,
        physical_fetch_hint: Option<usize>,
        index_predicate_execution: Option<IndexPredicateExecution<'_>>,
    ) -> Result<Option<FastPathKeyResult>, InternalError> {
        execute_fast_stream_route(
            &self.runtime,
            FastStreamRouteKind::SecondaryIndex,
            FastStreamRouteRequest::SecondaryIndex {
                plan,
                executable_access: &executable_access,
                bindings,
                probe_fetch_hint: physical_fetch_hint,
                index_predicate_execution,
            },
        )
    }

    /// Resolve one verified index-range limit-pushdown fast path.
    pub(in crate::db::executor) fn try_execute_index_range_limit_pushdown_stream(
        &self,
        plan: &AccessPlannedQuery,
        executable_access: ExecutableAccessPlan<'_, Value>,
        index_range_spec: Option<&LoweredIndexRangeSpec>,
        continuation: AccessScanContinuationInput<'_>,
        fetch: usize,
        index_predicate_execution: Option<IndexPredicateExecution<'_>>,
    ) -> Result<Option<FastPathKeyResult>, InternalError> {
        execute_fast_stream_route(
            &self.runtime,
            FastStreamRouteKind::IndexRangeLimitPushdown,
            FastStreamRouteRequest::IndexRangeLimitPushdown {
                plan,
                executable_access: &executable_access,
                index_range_spec,
                continuation,
                effective_fetch: fetch,
                index_predicate_execution,
            },
        )
    }

    /// Resolve the canonical fallback routed key stream for this execution attempt.
    pub(in crate::db::executor) fn resolve_fallback_execution_key_stream(
        &self,
        executable_access: ExecutableAccessPlan<'_, Value>,
        bindings: AccessStreamBindings<'_>,
        execution_policy: AccessStreamExecutionPolicy,
        index_predicate_execution: Option<IndexPredicateExecution<'_>>,
    ) -> Result<OrderedKeyStreamBox, InternalError> {
        let access = ExecutableAccess::from_executable_plan_with_policy(
            executable_access,
            bindings,
            execution_policy,
            index_predicate_execution,
        );

        self.runtime.ordered_key_stream_from_runtime_access(access)
    }

    /// Attempt the cursorless row-collector short path and erase the typed page result.
    fn try_materialize_load_via_row_collector<'req>(
        &'req self,
        request: RowCollectorMaterializationRequest<'req>,
    ) -> Result<Option<MaterializedExecutionPayloadResult>, InternalError> {
        self.with_scalar_row_runtime_handle(|row_runtime| {
            ExecutionKernel::try_materialize_load_via_row_collector(request, row_runtime)
        })
    }

    /// Materialize one ordered key stream into one structural scalar page payload.
    fn materialize_key_stream_into_structural_page<'a>(
        &'a self,
        contract: &ExecutionMaterializationContract<'a>,
        emit_cursor: bool,
        consistency: MissingRowPolicy,
        continuation: ScalarContinuationContext,
        key_stream: &'a mut OrderedKeyStreamBox,
    ) -> Result<MaterializedExecutionPayloadResult, InternalError> {
        let cursor_emission = if emit_cursor {
            CursorEmissionMode::Emit
        } else {
            CursorEmissionMode::Suppress
        };

        self.with_scalar_row_runtime_handle(|row_runtime| {
            materialize_key_stream_into_execution_payload(
                KernelPageMaterializationRequest {
                    plan: contract.plan,
                    key_stream,
                    scan_budget_hint: contract.scan_budget_hint,
                    load_order_route_mode: contract.load_order_route_mode,
                    capabilities: contract.capabilities(cursor_emission),
                    consistency,
                    continuation,
                },
                row_runtime,
            )
        })
    }

    /// Materialize one ordered key stream into post-access kernel rows.
    #[cfg(feature = "sql")]
    fn materialize_key_stream_into_kernel_rows<'a>(
        &'a self,
        contract: &ExecutionMaterializationContract<'a>,
        consistency: MissingRowPolicy,
        continuation: ScalarContinuationContext,
        key_stream: &'a mut OrderedKeyStreamBox,
    ) -> Result<KernelRowsExecutionAttempt, InternalError> {
        self.with_scalar_row_runtime_handle(|row_runtime| {
            materialize_key_stream_into_kernel_rows(
                KernelPageMaterializationRequest {
                    plan: contract.plan,
                    key_stream,
                    scan_budget_hint: contract.scan_budget_hint,
                    load_order_route_mode: contract.load_order_route_mode,
                    capabilities: contract.capabilities(CursorEmissionMode::Suppress),
                    consistency,
                    continuation,
                },
                row_runtime,
            )
        })
    }
}

type RetainedSlotLayout = crate::db::executor::terminal::RetainedSlotLayout;