icydb 0.252.11

IcyDB — A schema-first typed query engine and persistence runtime for Internet Computer canisters
Documentation
//! Module: db::query::typed
//!
//! Responsibility: typed read ergonomics over the accepted dynamic-query lane.
//! Does not own: schema identity, planning, admission, execution, or row decoding.
//! Boundary: generated adapters supply an opaque binding; accepted authority
//! resolves and executes the structural query before generated output decoding.

use crate::{
    db::{
        DbSession, DynamicQuery, ExhaustiveReadError, GroupedQueryOutput, PreparedLivePageOutput,
        TypedEntityAdapter, TypedEntityBinding, TypedOperationError,
    },
    traits::{CanisterKind, EntityKey},
    types::Id,
};
use candid::CandidType;
use icydb_core::db::{AggregateExpr, FilterExpr, OrderTerm, PrimaryKeyEncode, PrimaryKeyValue};
use serde::Deserialize;
use std::{error::Error as StdError, fmt, marker::PhantomData};

/// One revision-tolerant bounded typed page.
#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
pub struct LivePage<Row> {
    /// Decoded typed rows returned by this page.
    pub rows: Vec<Row>,
    /// Authenticated continuation, or `None` after proven exhaustion.
    pub continuation: Option<String>,
    /// Bounded work observed while producing the page.
    pub work: crate::db::ScalarPageWork,
}

/// One revision-strict bounded typed page.
#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
pub struct ExhaustivePage<Row> {
    /// Decoded typed rows returned by this page.
    pub rows: Vec<Row>,
    /// Authenticated continuation, or `None` after proof-bound exhaustion.
    pub continuation: Option<String>,
    /// Bounded work observed while producing this page.
    pub work: crate::db::ScalarPageWork,
    /// Complete source proof that must accompany the next resume.
    pub proof: crate::db::ReadSetRevisionProof,
}

impl<Row> LivePage<Row> {
    /// Return the number of decoded rows in this page.
    #[must_use]
    pub const fn len(&self) -> usize {
        self.rows.len()
    }

    /// Return whether this page contains no decoded rows.
    #[must_use]
    pub const fn is_empty(&self) -> bool {
        self.rows.is_empty()
    }
}

impl<Row> ExhaustivePage<Row> {
    /// Return the number of decoded rows in this page.
    #[must_use]
    pub const fn len(&self) -> usize {
        self.rows.len()
    }

    /// Return whether this page contains no decoded rows.
    #[must_use]
    pub const fn is_empty(&self) -> bool {
        self.rows.is_empty()
    }
}

/// Failure while decoding or executing one typed exhaustive page.
#[derive(Debug)]
pub enum TypedExhaustiveQueryError {
    /// The accepted exhaustive-read boundary or its source proof failed.
    Exhaustive(ExhaustiveReadError),
    /// Typed binding, decoding, or ordinary database execution failed.
    Operation(TypedOperationError),
}

impl fmt::Display for TypedExhaustiveQueryError {
    fn fmt(&self, formatter: &mut fmt::Formatter<'_>) -> fmt::Result {
        match self {
            Self::Exhaustive(error) => error.fmt(formatter),
            Self::Operation(error) => error.fmt(formatter),
        }
    }
}

impl StdError for TypedExhaustiveQueryError {}

impl From<TypedOperationError> for TypedExhaustiveQueryError {
    fn from(error: TypedOperationError) -> Self {
        Self::Operation(error)
    }
}

fn encode_typed_exact_keys<K>(keys: &[K]) -> Result<Vec<PrimaryKeyValue>, TypedOperationError>
where
    K: PrimaryKeyEncode,
{
    keys.iter()
        .map(|key| {
            key.to_primary_key_value().map_err(|error| {
                TypedOperationError::Database(crate::Error::from(
                    icydb_core::error::InternalError::from(error),
                ))
            })
        })
        .collect()
}

///
/// Query
///
/// Typed application projection over one accepted-schema-driven dynamic read.
/// Query planning, admission, and execution never consume generated model
/// metadata. The generated type participates only in binding and output decode.
///
pub struct Query<'session, C, E>
where
    C: CanisterKind,
    E: TypedEntityAdapter,
{
    session: &'session DbSession<C>,
    binding: TypedEntityBinding,
    request: DynamicQuery,
    entity: PhantomData<fn() -> E>,
}

impl<'session, C, E> Query<'session, C, E>
where
    C: CanisterKind,
    E: TypedEntityAdapter,
{
    pub(crate) fn new(session: &'session DbSession<C>) -> Result<Self, TypedOperationError> {
        let binding = E::typed_binding(session)?;
        let request = DynamicQuery::new(binding.entity());
        Ok(Self {
            session,
            binding,
            request,
            entity: PhantomData,
        })
    }

    /// Add one accepted-field filter expression, joined with prior filters by `AND`.
    #[must_use]
    pub fn filter(mut self, filter: impl Into<FilterExpr>) -> Self {
        self.request = self.request.filter(filter);
        self
    }

    /// Append one deterministic accepted-field ordering term.
    #[must_use]
    pub fn order_by(mut self, order: OrderTerm) -> Self {
        self.request = self.request.order_by(order);
        self
    }

    /// Select explicit accepted fields in scalar output order.
    ///
    /// Grouped execution rejects an explicit scalar selection because group
    /// keys and aggregates define its output contract.
    #[must_use]
    pub fn select<I, S>(mut self, fields: I) -> Self
    where
        I: IntoIterator<Item = S>,
        S: Into<String>,
    {
        self.request = self.request.select(fields);
        self
    }

    /// Bound the maximum number of returned rows.
    #[must_use]
    pub fn limit(mut self, limit: u32) -> Self {
        self.request = self.request.limit(limit);
        self
    }

    /// Append one accepted field to the grouped key in declaration order.
    #[must_use]
    pub fn group_by(mut self, field: impl Into<String>) -> Self {
        self.request = self.request.group_by(field);
        self
    }

    /// Append one grouped aggregate in declaration order.
    #[must_use]
    pub fn aggregate(mut self, aggregate: AggregateExpr) -> Self {
        self.request = self.request.aggregate(aggregate);
        self
    }

    /// Set explicit hard limits for grouped execution.
    #[must_use]
    pub fn grouped_limits(mut self, max_groups: u32, max_group_bytes: u32) -> Self {
        self.request = self.request.grouped_limits(max_groups, max_group_bytes);
        self
    }

    /// Continue from one opaque grouped cursor returned by IcyDB.
    #[must_use]
    pub fn cursor(mut self, cursor: impl Into<String>) -> Self {
        self.request = self.request.cursor(cursor);
        self
    }

    /// Return exact visible cardinality without scanning rows.
    ///
    /// This terminal accepts a bare entity query or one strict equality or
    /// bounded `IN` filter over the leading field of an accepted unfiltered
    /// field-path user index. The index may have trailing fields. Other shapes
    /// and unavailable exact-cardinality metadata fail closed.
    pub fn execute_exact_count(self) -> Result<u64, TypedOperationError> {
        self.session
            .execute_public_typed_exact_count(&self.binding, &self.request)?
            .ok_or(TypedOperationError::Adapter(
                crate::db::TypedAdapterError::StaleBinding,
            ))
    }

    /// Execute one revision-tolerant bounded page and decode its typed rows.
    ///
    /// Pass the prior page's opaque continuation to resume. A non-null
    /// continuation means traversal has not yet been proven exhausted.
    pub fn execute_live_page(
        self,
        continuation: Option<&str>,
    ) -> Result<LivePage<E::Row>, TypedOperationError> {
        let cursor = self
            .session
            .prepare_live_page_cursor(self.binding, self.request);
        let result = cursor.execute_page(continuation)?;
        Self::decode_live_page(cursor.binding(), result)
    }

    fn decode_live_page(
        binding: &TypedEntityBinding,
        prepared: PreparedLivePageOutput,
    ) -> Result<LivePage<E::Row>, TypedOperationError> {
        let mut rows = Vec::with_capacity(prepared.rows.len());
        for row in prepared.rows {
            rows.push(E::decode_row(binding, row)?);
        }

        Ok(LivePage {
            rows,
            continuation: prepared.continuation,
            work: prepared.work,
        })
    }

    /// Execute one revision-strict bounded page and decode its typed rows.
    ///
    /// Supply the prior page's continuation and proof together. Omitting both
    /// starts a single-store proof; a pre-captured multi-store proof may be
    /// supplied on the first call for a larger validation job.
    pub fn execute_exhaustive_page(
        self,
        continuation: Option<&str>,
        proof: Option<&crate::db::ReadSetRevisionProof>,
    ) -> Result<ExhaustivePage<E::Row>, TypedExhaustiveQueryError> {
        let result = self
            .session
            .execute_public_typed_exhaustive_page(&self.binding, &self.request, continuation, proof)
            .map_err(TypedExhaustiveQueryError::Exhaustive)?
            .ok_or({
                TypedExhaustiveQueryError::Operation(TypedOperationError::Adapter(
                    crate::db::TypedAdapterError::StaleBinding,
                ))
            })?;
        let crate::db::ExhaustiveQueryPageOutput {
            entity,
            columns,
            rows: output_rows,
            row_count: _,
            continuation,
            work,
            proof,
        } = result;
        let output_rows = self
            .session
            .prepare_typed_output_rows(&self.binding, entity, columns, output_rows)
            .map_err(TypedExhaustiveQueryError::Operation)?;
        let mut rows = Vec::with_capacity(output_rows.len());
        for row in output_rows {
            rows.push(E::decode_row(&self.binding, row).map_err(|error| {
                TypedExhaustiveQueryError::Operation(TypedOperationError::Adapter(error))
            })?);
        }

        Ok(ExhaustivePage {
            rows,
            continuation,
            work,
            proof,
        })
    }

    /// Execute through ordinary bounded grouped-read admission.
    ///
    /// Group keys and aggregate outputs preserve their declaration order. The
    /// accepted schema, shared query planner, and grouped executor remain the
    /// sole runtime authorities; `E` supplies only the source-bound entity
    /// binding used to reject stale adapters.
    pub fn execute_grouped(self) -> Result<GroupedQueryOutput, TypedOperationError> {
        self.session
            .execute_public_typed_dynamic_grouped_query(&self.binding, &self.request)?
            .ok_or(TypedOperationError::Adapter(
                crate::db::TypedAdapterError::StaleBinding,
            ))
    }
}

impl<C: CanisterKind> DbSession<C> {
    /// Start one typed read bound to current accepted schema authority.
    pub fn query<E>(&self) -> Result<Query<'_, C, E>, TypedOperationError>
    where
        E: TypedEntityAdapter,
    {
        Query::new(self)
    }

    /// Read one generated entity directly by its typed entity identifier.
    ///
    /// This path validates one current accepted binding and performs one
    /// bounded store lookup without constructing or caching a dynamic plan.
    pub fn get<E>(&self, id: Id<E>) -> Result<Option<E::Row>, TypedOperationError>
    where
        E: EntityKey + TypedEntityAdapter,
        E::Row: Clone,
    {
        let mut rows = self.get_many::<E>(&[id])?;
        rows.pop().ok_or(TypedOperationError::Adapter(
            crate::db::TypedAdapterError::RowShapeMismatch,
        ))
    }

    /// Read generated entities directly by typed entity identifier.
    ///
    /// The result has exactly one position per input identifier in request
    /// order. Missing identifiers produce `None`; duplicates preserve positions
    /// while sharing one physical lookup and persisted-row decode. One batch is
    /// bounded by [`MAX_TYPED_EXACT_KEY_BATCH_ITEMS`], encoded key bytes,
    /// distinct stored-row bytes, and logical result bytes.
    pub fn get_many<E>(&self, ids: &[Id<E>]) -> Result<Vec<Option<E::Row>>, TypedOperationError>
    where
        E: EntityKey + TypedEntityAdapter,
        E::Row: Clone,
    {
        let binding = E::typed_binding(self)?;
        let keys = encode_typed_exact_keys(ids)?;
        let prepared = self.execute_public_prepared_exact_key_batch(&binding, &keys)?;
        let mut distinct_rows = Vec::with_capacity(prepared.distinct_rows.len());
        for row in prepared.distinct_rows {
            let row = match row {
                Some(row) => Some(E::decode_row(&binding, row)?),
                None => None,
            };
            distinct_rows.push(row);
        }

        let mut rows = Vec::with_capacity(prepared.positions.len());
        for position in prepared.positions {
            let index = usize::try_from(position)
                .map_err(|_| crate::db::TypedAdapterError::RowShapeMismatch)?;
            rows.push(
                distinct_rows
                    .get(index)
                    .cloned()
                    .ok_or(crate::db::TypedAdapterError::RowShapeMismatch)?,
            );
        }
        Ok(rows)
    }
}

/// Maximum input positions admitted by one typed exact-key batch.
pub const MAX_TYPED_EXACT_KEY_BATCH_ITEMS: usize = icydb_core::db::MAX_TYPED_EXACT_KEY_BATCH_ITEMS;

/// Maximum encoded stored-key bytes admitted before input deduplication.
pub const MAX_TYPED_EXACT_KEY_BATCH_INPUT_BYTES: usize =
    icydb_core::db::MAX_TYPED_EXACT_KEY_BATCH_INPUT_BYTES;

/// Maximum raw row bytes admitted across distinct stored keys.
pub const MAX_TYPED_EXACT_KEY_BATCH_STORED_BYTES: usize =
    icydb_core::db::MAX_TYPED_EXACT_KEY_BATCH_STORED_BYTES;

/// Maximum logical projection bytes admitted across original input positions.
pub const MAX_TYPED_EXACT_KEY_BATCH_RESULT_BYTES: usize =
    icydb_core::db::MAX_TYPED_EXACT_KEY_BATCH_RESULT_BYTES;