icydb 0.252.1

IcyDB — A schema-first typed query engine and persistence runtime for Internet Computer canisters
Documentation
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
419
420
421
422
423
424
425
426
427
428
429
430
431
432
433
434
435
436
437
438
439
440
441
442
443
444
445
446
447
448
449
450
451
452
453
454
455
456
457
458
459
460
461
462
463
464
465
466
467
468
469
470
471
472
473
474
475
476
477
//! 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::{
        AttributedRead, DbSession, DynamicQuery, ExhaustiveReadError, GroupedQueryOutput,
        PreparedLivePageOutput, TypedBindingError, TypedEntityAdapter, TypedEntityBinding,
        TypedRowError,
    },
    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};

/// Failure while executing one accepted-schema-bound typed query.
#[derive(Debug)]
pub enum TypedQueryError {
    /// The accepted dynamic read rejected or failed.
    Database(crate::Error),
    /// A returned accepted row could not be projected through the binding.
    Row(TypedRowError),
}

/// 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 {
    Exhaustive(ExhaustiveReadError),
    Row(TypedRowError),
}

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

impl StdError for TypedExhaustiveQueryError {}

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

impl StdError for TypedQueryError {}

fn typed_query_error_from_binding(error: TypedBindingError) -> TypedQueryError {
    match error {
        TypedBindingError::Adapter(error) => TypedQueryError::Row(TypedRowError::Adapter(error)),
        TypedBindingError::Database(error) => TypedQueryError::Database(error),
    }
}

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

#[must_use]
#[cfg(target_arch = "wasm32")]
fn read_operation_local_instruction_counter() -> u64 {
    ic_cdk::api::performance_counter(1)
}

#[must_use]
#[cfg(not(target_arch = "wasm32"))]
const fn read_operation_local_instruction_counter() -> u64 {
    0
}

///
/// 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, TypedBindingError> {
        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, TypedQueryError> {
        self.session
            .execute_public_typed_exact_count(&self.binding, &self.request)
            .map_err(TypedQueryError::Database)?
            .ok_or({
                TypedQueryError::Row(TypedRowError::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>, TypedQueryError> {
        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)
    }

    /// Execute one revision-tolerant typed page with bounded per-call cost attribution.
    ///
    /// The operation follows the same accepted dynamic execution and typed row
    /// decoding as [`Self::execute_live_page`]. Its fixed attribution envelope
    /// is returned beside the page without entering retained metrics or
    /// exposing query, entity, index, literal or caller identity.
    pub fn execute_live_page_with_attribution(
        self,
        continuation: Option<&str>,
    ) -> Result<AttributedRead<LivePage<E::Row>>, TypedQueryError> {
        let start = read_operation_local_instruction_counter();
        let cursor = self
            .session
            .prepare_live_page_cursor(self.binding, self.request);
        let attributed = cursor.execute_public_page_with_attribution(continuation)?;
        let decode_start = read_operation_local_instruction_counter();
        let result = Self::decode_live_page(cursor.binding(), attributed.result)?;
        let response_decode_local_instructions =
            read_operation_local_instruction_counter().saturating_sub(decode_start);
        let total_local_instructions =
            read_operation_local_instruction_counter().saturating_sub(start);
        let mut attribution = attributed.attribution;
        attribution.response_decode_local_instructions = attribution
            .response_decode_local_instructions
            .saturating_add(response_decode_local_instructions);
        attribution.total_local_instructions = total_local_instructions;

        Ok(AttributedRead {
            result,
            attribution,
        })
    }

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

        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::Row(TypedRowError::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::Row)?;
        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::Row(TypedRowError::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, TypedQueryError> {
        self.session
            .execute_public_typed_dynamic_grouped_query(&self.binding, &self.request)
            .map_err(TypedQueryError::Database)?
            .ok_or({
                TypedQueryError::Row(TypedRowError::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>, TypedBindingError>
    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>, TypedQueryError>
    where
        E: EntityKey + TypedEntityAdapter,
        E::Row: Clone,
    {
        let mut rows = self.get_many::<E>(&[id])?;
        rows.pop().ok_or({
            TypedQueryError::Row(TypedRowError::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>>, TypedQueryError>
    where
        E: EntityKey + TypedEntityAdapter,
        E::Row: Clone,
    {
        let binding = E::typed_binding(self).map_err(typed_query_error_from_binding)?;
        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)
                        .map_err(|error| TypedQueryError::Row(TypedRowError::Adapter(error)))?,
                ),
                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(|_| {
                TypedQueryError::Row(TypedRowError::Adapter(
                    crate::db::TypedAdapterError::RowShapeMismatch,
                ))
            })?;
            rows.push(distinct_rows.get(index).cloned().ok_or({
                TypedQueryError::Row(TypedRowError::Adapter(
                    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;