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};
#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
pub struct LivePage<Row> {
pub rows: Vec<Row>,
pub continuation: Option<String>,
pub work: crate::db::ScalarPageWork,
}
#[derive(CandidType, Clone, Debug, Deserialize, Eq, PartialEq)]
pub struct ExhaustivePage<Row> {
pub rows: Vec<Row>,
pub continuation: Option<String>,
pub work: crate::db::ScalarPageWork,
pub proof: crate::db::ReadSetRevisionProof,
}
impl<Row> LivePage<Row> {
#[must_use]
pub const fn len(&self) -> usize {
self.rows.len()
}
#[must_use]
pub const fn is_empty(&self) -> bool {
self.rows.is_empty()
}
}
impl<Row> ExhaustivePage<Row> {
#[must_use]
pub const fn len(&self) -> usize {
self.rows.len()
}
#[must_use]
pub const fn is_empty(&self) -> bool {
self.rows.is_empty()
}
}
#[derive(Debug)]
pub enum TypedExhaustiveQueryError {
Exhaustive(ExhaustiveReadError),
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()
}
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,
})
}
#[must_use]
pub fn filter(mut self, filter: impl Into<FilterExpr>) -> Self {
self.request = self.request.filter(filter);
self
}
#[must_use]
pub fn order_by(mut self, order: OrderTerm) -> Self {
self.request = self.request.order_by(order);
self
}
#[must_use]
pub fn limit(mut self, limit: u32) -> Self {
self.request = self.request.limit(limit);
self
}
#[must_use]
pub fn group_by(mut self, field: impl Into<String>) -> Self {
self.request = self.request.group_by(field);
self
}
#[must_use]
pub fn aggregate(mut self, aggregate: AggregateExpr) -> Self {
self.request = self.request.aggregate(aggregate);
self
}
#[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
}
#[must_use]
pub fn cursor(mut self, cursor: impl Into<String>) -> Self {
self.request = self.request.cursor(cursor);
self
}
pub fn explain(self) -> Result<crate::db::query::ExplainPlan, TypedOperationError> {
self.session
.explain_typed_query(&self.binding, &self.request)?
.ok_or(TypedOperationError::Adapter(
crate::db::TypedAdapterError::StaleBinding,
))
}
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,
))
}
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,
})
}
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,
})
}
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> {
pub fn query<E>(&self) -> Result<Query<'_, C, E>, TypedOperationError>
where
E: TypedEntityAdapter,
{
Query::new(self)
}
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,
))
}
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((0, row));
}
for (output_position, position) in prepared.positions.iter().enumerate() {
let index = usize::try_from(*position)
.map_err(|_| crate::db::TypedAdapterError::RowShapeMismatch)?;
let (last_use, _) = distinct_rows
.get_mut(index)
.ok_or(crate::db::TypedAdapterError::RowShapeMismatch)?;
*last_use = output_position;
}
let mut rows = Vec::with_capacity(prepared.positions.len());
for (output_position, position) in prepared.positions.into_iter().enumerate() {
let index = usize::try_from(position)
.map_err(|_| crate::db::TypedAdapterError::RowShapeMismatch)?;
let (last_use, row) = distinct_rows
.get_mut(index)
.ok_or(crate::db::TypedAdapterError::RowShapeMismatch)?;
rows.push(if *last_use == output_position {
row.take()
} else {
row.clone()
});
}
Ok(rows)
}
}
pub const MAX_TYPED_EXACT_KEY_BATCH_ITEMS: usize = icydb_core::db::MAX_TYPED_EXACT_KEY_BATCH_ITEMS;
pub const MAX_TYPED_EXACT_KEY_BATCH_INPUT_BYTES: usize =
icydb_core::db::MAX_TYPED_EXACT_KEY_BATCH_INPUT_BYTES;
pub const MAX_TYPED_EXACT_KEY_BATCH_STORED_BYTES: usize =
icydb_core::db::MAX_TYPED_EXACT_KEY_BATCH_STORED_BYTES;
pub const MAX_TYPED_EXACT_KEY_BATCH_RESULT_BYTES: usize =
icydb_core::db::MAX_TYPED_EXACT_KEY_BATCH_RESULT_BYTES;