Skip to main content

icydb_core/db/session/query/
dynamic.rs

1//! Module: db::session::query::dynamic
2//! Responsibility: lower and execute public dynamic reads against accepted schema.
3//! Does not own: query planning, accepted schema construction, or row projection.
4//! Boundary: entity-name requests converge on the shared structural read lane.
5
6use crate::{
7    db::{
8        DbSession, DynamicQuery, DynamicTypedEntityBinding, ExhaustiveQueryPageOutput,
9        ExhaustiveReadError, GroupedQueryOutput, LiveQueryPageOutput, MissingRowPolicy, QueryError,
10        ReadSetRevisionError, ReadSetRevisionProof, ScalarPageWork,
11        codec::{finalize_hash_sha256, new_hash_sha256_prefixed},
12        commit::{cursor_authentication_key, database_incarnation_id},
13        cursor::{
14            CursorBoundary, CursorBoundarySlot, CursorPlanError, ScalarOrderTermContract,
15            ScalarPageMode, ScalarPageToken, ScalarPageTokenAuthority, ScalarPageTokenProgress,
16            ScalarPageTokenWindow, decode_optional_cursor_token, encode_cursor,
17        },
18        data::{DecodedDataStoreKey, RawDataStoreKey},
19        executor::{
20            PageWorkEnvelope, ScalarContinuationContext, StructuralProjectionRequest,
21            execute_structural_projection_page,
22        },
23        query::{
24            admission::{QueryAdmissionPolicy, QueryAdmissionSummary},
25            expr::{CompareOperator, FilterExpr, OrderTerm as FluentOrderTerm, SetOperator},
26            intent::{IntentError, StructuralQuery},
27            plan::CardinalityTiebreakRoutePin,
28        },
29        session::AcceptedSchemaCatalogContext,
30    },
31    traits::CanisterKind,
32};
33use icydb_diagnostic_code::{
34    DiagnosticDecodeReason, DiagnosticExecutionBudgetResource, DiagnosticExecutionLane,
35    QueryReadAdmissionCode,
36};
37use sha2::Digest;
38#[cfg(test)]
39use std::cell::Cell;
40
41#[cfg(not(test))]
42const SCALAR_PAGE_OUTPUT_ROWS: usize = 1_024;
43#[cfg(test)]
44const SCALAR_PAGE_OUTPUT_ROWS: usize = 2;
45#[cfg(test)]
46const SCALAR_PAGE_KEY_ENTRIES: u64 = 4;
47
48#[cfg(test)]
49std::thread_local! {
50    static SCALAR_PAGE_RESULT_BYTES_LIMIT_OVERRIDE: Cell<Option<u64>> = const { Cell::new(None) };
51}
52
53#[cfg(test)]
54struct ScalarPageResultBytesLimitGuard(Option<u64>);
55
56#[cfg(test)]
57impl Drop for ScalarPageResultBytesLimitGuard {
58    fn drop(&mut self) {
59        SCALAR_PAGE_RESULT_BYTES_LIMIT_OVERRIDE.with(|limit| limit.set(self.0));
60    }
61}
62
63#[derive(Clone, Copy)]
64enum DynamicReadLane {
65    Public,
66    Trusted,
67}
68
69struct ScalarCursorContract {
70    signature: crate::db::cursor::ContinuationSignature,
71    authority: ScalarPageTokenAuthority,
72    route_pin: Option<CardinalityTiebreakRoutePin>,
73    window: ScalarPageTokenWindow,
74    order_terms: Vec<ScalarOrderTermContract>,
75}
76
77impl<C: CanisterKind> DbSession<C> {
78    fn may_select_exact_single_primary_key(
79        request: &DynamicQuery,
80        catalog: &AcceptedSchemaCatalogContext,
81    ) -> bool {
82        let [primary_key] = catalog.accepted_schema_info().primary_key_names() else {
83            return false;
84        };
85        matches!(
86            request.filter_expr(),
87            Some(
88                FilterExpr::Compare {
89                    operator: CompareOperator::Eq,
90                    field,
91                    ..
92                } | FilterExpr::Set {
93                    operator: SetOperator::In,
94                    field,
95                    ..
96                }
97            )
98                if field.eq_ignore_ascii_case(primary_key)
99        )
100    }
101
102    fn exact_primary_key_candidate_bound(
103        prepared_plan: &crate::db::executor::SharedPreparedExecutionPlan,
104    ) -> Option<usize> {
105        let access = &prepared_plan.logical_plan().access;
106        if access.as_by_key_path().is_some() {
107            return Some(1);
108        }
109
110        access.as_by_keys_path().map(<[crate::value::Value]>::len)
111    }
112
113    pub(super) fn structural_query_from_dynamic_request(
114        request: &DynamicQuery,
115        catalog: &AcceptedSchemaCatalogContext,
116    ) -> Result<StructuralQuery, QueryError> {
117        Self::structural_query_from_dynamic_request_with_page_limit(request, catalog, None, false)
118    }
119
120    fn structural_query_from_dynamic_request_with_page_limit(
121        request: &DynamicQuery,
122        catalog: &AcceptedSchemaCatalogContext,
123        page_limit: Option<u32>,
124        require_total_order: bool,
125    ) -> Result<StructuralQuery, QueryError> {
126        let schema = catalog.accepted_schema_info();
127        let mut query = StructuralQuery::new(MissingRowPolicy::Ignore);
128        if let Some(filter) = request.filter_expr() {
129            query = query.filter_for_schema(schema, filter.clone());
130        }
131        for order in request.order_terms() {
132            query = query.order_term(order.clone());
133        }
134        if require_total_order && request.order_terms().is_empty() {
135            for primary_key in schema.primary_key_names() {
136                query = query.order_term(FluentOrderTerm::asc(primary_key.clone()));
137            }
138        }
139        if !request.selected_fields().is_empty() {
140            query = query.select_fields(request.selected_fields().iter().cloned());
141        }
142        #[cfg(test)]
143        if request.projection_is_distinct() {
144            query = query.distinct();
145        }
146        if let Some(limit) = page_limit.or_else(|| request.row_limit()) {
147            query = query.limit(limit);
148        }
149        for field in request.group_fields() {
150            query = query.group_by_with_schema(field, schema)?;
151        }
152        for aggregate in request.aggregates() {
153            query = query.aggregate(aggregate.clone());
154        }
155        if let Some((max_groups, max_group_bytes)) = request.grouped_execution_limits() {
156            if max_groups == 0 || max_group_bytes == 0 {
157                return Err(QueryReadAdmissionCode::GroupedQueryRequiresLimits.into());
158            }
159            query = query.grouped_limits(u64::from(max_groups), u64::from(max_group_bytes));
160        }
161
162        Ok(query)
163    }
164
165    fn scalar_page_cursor_error() -> QueryError {
166        QueryError::from_cursor_plan_error(CursorPlanError::invalid_continuation_cursor_payload(
167            DiagnosticDecodeReason::CursorTokenDecode,
168        ))
169    }
170
171    fn scalar_cursor_contract(
172        request: &DynamicQuery,
173        catalog: &AcceptedSchemaCatalogContext,
174        envelope: PageWorkEnvelope,
175        prepared_plan: &crate::db::executor::SharedPreparedExecutionPlan,
176        mode: ScalarPageMode,
177        proof: Option<&ReadSetRevisionProof>,
178    ) -> Result<ScalarCursorContract, QueryError> {
179        let mut signature = prepared_plan
180            .continuation_signature_for_runtime()
181            .map_err(QueryError::execute)?;
182        match (mode, proof) {
183            (ScalarPageMode::Live, None) => {}
184            (ScalarPageMode::Exhaustive, Some(proof)) => {
185                let mut hasher = new_hash_sha256_prefixed(b"icydb.exhaustive-cursor-proof.v1");
186                hasher.update(signature.into_bytes());
187                hasher.update(proof.signature_bytes());
188                signature = crate::db::cursor::ContinuationSignature::from_bytes(
189                    finalize_hash_sha256(hasher),
190                );
191            }
192            _ => return Err(Self::scalar_page_cursor_error()),
193        }
194        let root_identity = catalog.runtime_root_identity();
195        let (root_fingerprint_method, root_fingerprint) = root_identity.fingerprint();
196        let authority = ScalarPageTokenAuthority::new(
197            database_incarnation_id()
198                .map_err(QueryError::execute)?
199                .to_bytes(),
200            root_identity.accepted_root_revision().get(),
201            root_fingerprint_method,
202            root_fingerprint,
203            catalog.fingerprint(),
204            prepared_plan.authority_ref().entity_tag(),
205        );
206        let window =
207            ScalarPageTokenWindow::new(0, request.row_limit(), envelope.profile_identity());
208        let canonical_order = prepared_plan
209            .logical_plan()
210            .scalar_plan()
211            .order
212            .as_ref()
213            .ok_or_else(Self::scalar_page_cursor_error)?;
214        let order_terms = canonical_order
215            .fields
216            .iter()
217            .map(|term| ScalarOrderTermContract::new(term.rendered_label(), term.direction()))
218            .collect::<Vec<_>>();
219
220        Ok(ScalarCursorContract {
221            signature,
222            authority,
223            route_pin: prepared_plan
224                .logical_plan()
225                .cardinality_tiebreak_route_pin(),
226            window,
227            order_terms,
228        })
229    }
230
231    fn validate_scalar_page_token(
232        token: &ScalarPageToken,
233        mode: ScalarPageMode,
234        contract: &ScalarCursorContract,
235        entity: &str,
236    ) -> Result<(), QueryError> {
237        if token.signature() != contract.signature {
238            return Err(QueryError::from_cursor_plan_error(
239                CursorPlanError::continuation_cursor_signature_mismatch(
240                    entity,
241                    &contract.signature,
242                    &token.signature(),
243                ),
244            ));
245        }
246        if token.mode() != mode
247            || token.authority() != contract.authority
248            || token.route_pin() != contract.route_pin
249            || token.window() != contract.window
250            || token.order_terms() != contract.order_terms.as_slice()
251        {
252            return Err(Self::scalar_page_cursor_error());
253        }
254
255        Ok(())
256    }
257
258    fn physical_primary_key_boundary(
259        bytes: &[u8],
260        catalog: &AcceptedSchemaCatalogContext,
261    ) -> Result<CursorBoundary, QueryError> {
262        let raw = RawDataStoreKey::from_persisted_bytes(bytes.to_vec());
263        let key = DecodedDataStoreKey::try_from_raw(&raw)
264            .map_err(|_| Self::scalar_page_cursor_error())?;
265        if key.entity_tag() != catalog.accepted_entity_authority().entity_tag() {
266            return Err(Self::scalar_page_cursor_error());
267        }
268        let primary_key_arity = catalog.accepted_schema_info().primary_key_names().len();
269        let mut slots = Vec::with_capacity(primary_key_arity);
270        for component_index in 0..primary_key_arity {
271            slots.push(CursorBoundarySlot::Present(
272                key.primary_key_component_runtime_value(component_index)
273                    .map_err(|_| Self::scalar_page_cursor_error())?,
274            ));
275        }
276
277        Ok(CursorBoundary { slots })
278    }
279
280    fn execute_dynamic_grouped_query_against_catalog(
281        &self,
282        request: &DynamicQuery,
283        lane: DynamicReadLane,
284        catalog: AcceptedSchemaCatalogContext,
285    ) -> Result<GroupedQueryOutput, QueryError> {
286        if !request.has_grouping() {
287            return Err(QueryError::intent(
288                IntentError::grouped_terminal_requires_grouped_query(),
289            ));
290        }
291        if request.grouped_execution_limits().is_none() {
292            return Err(QueryReadAdmissionCode::GroupedQueryRequiresLimits.into());
293        }
294        if !request.selected_fields().is_empty() {
295            return Err(QueryError::intent(
296                IntentError::grouped_output_defined_by_group_and_aggregates(),
297            ));
298        }
299        let query = Self::structural_query_from_dynamic_request(request, &catalog)?;
300        let public_admission = match lane {
301            DynamicReadLane::Public => Some(QueryAdmissionPolicy::default_bounded_read()),
302            DynamicReadLane::Trusted => None,
303        };
304
305        self.execute_structural_grouped_from_query(
306            &query,
307            &catalog,
308            public_admission.as_ref(),
309            request.continuation_cursor(),
310        )
311    }
312
313    #[expect(
314        clippy::too_many_lines,
315        reason = "live-page orchestration keeps planning, cursor validation, execution, and response proof in one auditable boundary"
316    )]
317    fn execute_scalar_page_against_catalog(
318        &self,
319        request: &DynamicQuery,
320        continuation: Option<&str>,
321        lane: DynamicReadLane,
322        catalog: AcceptedSchemaCatalogContext,
323        mode: ScalarPageMode,
324        supplied_proof: Option<&ReadSetRevisionProof>,
325    ) -> Result<(LiveQueryPageOutput, Option<ReadSetRevisionProof>), ExhaustiveReadError> {
326        if request.has_grouping()
327            || request.grouped_execution_limits().is_some()
328            || request.continuation_cursor().is_some()
329        {
330            return Err(
331                QueryError::intent(IntentError::scalar_terminal_requires_scalar_query()).into(),
332            );
333        }
334
335        let exhaustive_proof = match mode {
336            ScalarPageMode::Live => None,
337            ScalarPageMode::Exhaustive => {
338                if continuation.is_some() && supplied_proof.is_none() {
339                    return Err(ReadSetRevisionError::ResumeProofRequired.into());
340                }
341                let proof = supplied_proof.cloned().map_or_else(
342                    || self.capture_entity_read_set_revision_proof(catalog.identity().store_path()),
343                    Ok,
344                )?;
345                Self::ensure_read_set_contains_store(&proof, catalog.identity().store_path())?;
346                self.verify_read_set_revision_proof(&proof)?;
347                Some(proof)
348            }
349        };
350
351        let envelope = match lane {
352            DynamicReadLane::Public => PageWorkEnvelope::public_scalar(),
353            DynamicReadLane::Trusted => PageWorkEnvelope::default_scalar(),
354        };
355        #[cfg(test)]
356        let envelope = SCALAR_PAGE_RESULT_BYTES_LIMIT_OVERRIDE.with(|limit| {
357            limit.get().map_or(envelope, |limit| {
358                envelope.with_limit_for_tests(DiagnosticExecutionBudgetResource::ResultBytes, limit)
359            })
360        });
361        #[cfg(test)]
362        let envelope = envelope.with_limit_for_tests(
363            DiagnosticExecutionBudgetResource::KeyIndexEntriesVisited,
364            SCALAR_PAGE_KEY_ENTRIES,
365        );
366        let page_row_limit = envelope
367            .limit(DiagnosticExecutionBudgetResource::ResultRows)
368            .and_then(|limit| usize::try_from(limit).ok())
369            .unwrap_or(SCALAR_PAGE_OUTPUT_ROWS)
370            .min(SCALAR_PAGE_OUTPUT_ROWS);
371        let decoded_token = decode_optional_cursor_token(continuation)
372            .map_err(QueryError::from_cursor_plan_error)?
373            .map(|bytes| {
374                ScalarPageToken::decode(
375                    bytes.as_slice(),
376                    &cursor_authentication_key().map_err(QueryError::execute)?,
377                )
378                .map_err(|error| {
379                    QueryError::from_cursor_plan_error(CursorPlanError::from_token_wire_error(
380                        error,
381                    ))
382                })
383            })
384            .transpose()?;
385        let prior_rows_emitted = decoded_token
386            .as_ref()
387            .map_or(0, |token| token.progress().rows_emitted());
388        let remaining_limit = request
389            .row_limit()
390            .map(|limit| u64::from(limit).saturating_sub(prior_rows_emitted));
391        let page_output_limit = remaining_limit
392            .unwrap_or(page_row_limit as u64)
393            .min(page_row_limit as u64);
394        let page_output_limit = usize::try_from(page_output_limit).unwrap_or(page_row_limit);
395        let execution_limit = u32::try_from(page_row_limit).unwrap_or(u32::MAX);
396        let execution_lane = match lane {
397            DynamicReadLane::Public => DiagnosticExecutionLane::PublicRead,
398            DynamicReadLane::Trusted => DiagnosticExecutionLane::TrustedRead,
399        };
400        let exact_candidate =
401            decoded_token.is_none() && Self::may_select_exact_single_primary_key(request, &catalog);
402        let initial_plan = if exact_candidate {
403            let query = Self::structural_query_from_dynamic_request(request, &catalog)?;
404            Some(
405                self.structural_projection_prepared_plan_for_accepted_authority(
406                    &query,
407                    catalog.accepted_entity_authority(),
408                    catalog.snapshot(),
409                    execution_lane,
410                )?,
411            )
412        } else {
413            None
414        };
415        let initial_is_exact_exhaustion =
416            initial_plan.as_ref().is_some_and(|(prepared_plan, _)| {
417                Self::exact_primary_key_candidate_bound(prepared_plan)
418                    .is_some_and(|bound| bound == 1 && bound <= page_output_limit)
419            });
420        let (prepared_plan, projection) = if initial_is_exact_exhaustion {
421            initial_plan.ok_or_else(Self::scalar_page_cursor_error)?
422        } else {
423            let query = Self::structural_query_from_dynamic_request_with_page_limit(
424                request,
425                &catalog,
426                Some(execution_limit),
427                true,
428            )?;
429            if let Some(route_pin) = decoded_token.as_ref().and_then(ScalarPageToken::route_pin) {
430                self.structural_projection_prepared_plan_for_accepted_authority_with_route_pin(
431                    &query,
432                    catalog.accepted_entity_authority(),
433                    catalog.snapshot(),
434                    execution_lane,
435                    route_pin,
436                )?
437                .ok_or_else(Self::scalar_page_cursor_error)?
438            } else {
439                self.structural_projection_prepared_plan_for_accepted_authority(
440                    &query,
441                    catalog.accepted_entity_authority(),
442                    catalog.snapshot(),
443                    execution_lane,
444                )?
445            }
446        };
447        if matches!(lane, DynamicReadLane::Public) {
448            let policy = QueryAdmissionPolicy::default_bounded_read();
449            let summary = policy.evaluate(QueryAdmissionSummary::from_plan(
450                policy.lane(),
451                prepared_plan.logical_plan(),
452            ));
453            if let Some(rejection) = summary.rejection() {
454                return Err(QueryError::from(rejection.code()).into());
455            }
456        }
457        let exact_initial_exhaustion = initial_is_exact_exhaustion;
458        let cursor_contract = decoded_token
459            .as_ref()
460            .map(|token| {
461                let contract = Self::scalar_cursor_contract(
462                    request,
463                    &catalog,
464                    envelope,
465                    &prepared_plan,
466                    mode,
467                    exhaustive_proof.as_ref(),
468                )?;
469                Self::validate_scalar_page_token(token, mode, &contract, request.entity())?;
470                Ok::<_, QueryError>(contract)
471            })
472            .transpose()?;
473        let deferred_cursor_plan =
474            (!exact_initial_exhaustion && decoded_token.is_none()).then(|| prepared_plan.clone());
475        let continuation_context = match decoded_token.as_ref() {
476            None => ScalarContinuationContext::initial(),
477            Some(token) if token.progress().unconsumed_lookahead().is_some() => {
478                return Err(Self::scalar_page_cursor_error().into());
479            }
480            Some(token) => {
481                let logical = token.progress().last_emitted_logical().cloned();
482                match token.progress().last_consumed_physical() {
483                    Some(physical) => ScalarContinuationContext::resumed_with_primary_progress(
484                        logical,
485                        Self::physical_primary_key_boundary(physical, &catalog)?,
486                    ),
487                    None => logical.map_or_else(
488                        ScalarContinuationContext::initial,
489                        ScalarContinuationContext::resumed,
490                    ),
491                }
492            }
493        };
494        if decoded_token.is_some() && !continuation_context.has_progress() {
495            return Err(Self::scalar_page_cursor_error().into());
496        }
497
498        let value_catalog = prepared_plan
499            .authority_ref()
500            .accepted_schema_info()
501            .map(crate::db::schema::SchemaInfo::value_catalog_handle)
502            .cloned()
503            .ok_or_else(QueryError::invariant)?;
504        let (columns, _fixed_scales) = projection.into_components();
505        let projection_request = StructuralProjectionRequest::new(prepared_plan, execution_lane)
506            .with_distinct_output_offset(usize::try_from(prior_rows_emitted).unwrap_or(usize::MAX))
507            .with_page_work_envelope(envelope);
508        let projection_request = if exact_initial_exhaustion {
509            projection_request
510        } else {
511            projection_request
512                .with_continuation(continuation_context)
513                .with_cursor_emission(page_output_limit)
514        };
515        let page = execute_structural_projection_page(&self.db, projection_request)
516            .map_err(QueryError::execute)?;
517        let row_count = page.rows.row_count();
518        let rows = page
519            .rows
520            .into_value_rows()
521            .into_iter()
522            .map(|row| {
523                row.iter()
524                    .map(|value| {
525                        crate::db::schema::output_value_from_runtime(
526                            value_catalog.enum_catalog(),
527                            value,
528                        )
529                        .map_err(|_| QueryError::invariant())
530                    })
531                    .collect::<Result<Vec<_>, _>>()
532            })
533            .collect::<Result<Vec<_>, _>>()?;
534        let rows_emitted = prior_rows_emitted.saturating_add(u64::from(row_count));
535        let total_limit_reached = request
536            .row_limit()
537            .is_some_and(|limit| rows_emitted >= u64::from(limit));
538        let continuation = if page.has_more && !total_limit_reached {
539            if page.last_emitted_logical.is_none() && page.last_consumed_physical.is_none() {
540                return Err(Self::scalar_page_cursor_error().into());
541            }
542            let cursor_contract = if let Some(contract) = cursor_contract {
543                contract
544            } else {
545                let prepared_plan = deferred_cursor_plan
546                    .as_ref()
547                    .ok_or_else(Self::scalar_page_cursor_error)?;
548                Self::scalar_cursor_contract(
549                    request,
550                    &catalog,
551                    envelope,
552                    prepared_plan,
553                    mode,
554                    exhaustive_proof.as_ref(),
555                )?
556            };
557            let token = ScalarPageToken::new(
558                mode,
559                cursor_contract.signature,
560                cursor_contract.authority,
561                cursor_contract.route_pin,
562                cursor_contract.window,
563                cursor_contract.order_terms,
564                ScalarPageTokenProgress::new(
565                    page.last_emitted_logical,
566                    page.last_consumed_physical,
567                    None,
568                    decoded_token
569                        .as_ref()
570                        .map_or(0, |token| token.progress().matching_rows_skipped()),
571                    rows_emitted,
572                ),
573            );
574            Some(encode_cursor(
575                token
576                    .encode(&cursor_authentication_key().map_err(QueryError::execute)?)
577                    .map_err(|error| {
578                        QueryError::from_cursor_plan_error(CursorPlanError::from_token_wire_error(
579                            error,
580                        ))
581                    })?
582                    .as_slice(),
583            ))
584        } else {
585            None
586        };
587
588        if let Some(proof) = exhaustive_proof.as_ref() {
589            self.verify_read_set_revision_proof(proof)?;
590        }
591
592        Ok((
593            LiveQueryPageOutput {
594                entity: catalog.snapshot().entity_name().to_string(),
595                columns,
596                rows,
597                row_count,
598                continuation,
599                work: ScalarPageWork {
600                    envelope_identity: envelope.identity(),
601                    entries_visited: page.scanned_keys as u64,
602                    result_rows: row_count,
603                },
604            },
605            exhaustive_proof,
606        ))
607    }
608
609    /// Execute one revision-tolerant bounded scalar page.
610    pub fn execute_public_live_page(
611        &self,
612        request: &DynamicQuery,
613        continuation: Option<&str>,
614    ) -> Result<LiveQueryPageOutput, QueryError> {
615        let catalog = self
616            .accepted_schema_catalog_context_for_entity_name(Some(request.entity()))
617            .map_err(QueryError::execute)?;
618        self.execute_scalar_page_against_catalog(
619            request,
620            continuation,
621            DynamicReadLane::Public,
622            catalog,
623            ScalarPageMode::Live,
624            None,
625        )
626        .map(|(page, _)| page)
627        .map_err(Self::live_page_error)
628    }
629
630    /// Execute one live page through a typed binding's immutable accepted
631    /// entity identity. `None` means the opaque binding is stale.
632    #[doc(hidden)]
633    pub fn execute_public_live_page_for_typed_binding(
634        &self,
635        binding: &DynamicTypedEntityBinding,
636        request: &DynamicQuery,
637        continuation: Option<&str>,
638    ) -> Result<Option<LiveQueryPageOutput>, QueryError> {
639        let Some(catalog) = self
640            .current_typed_entity_binding_catalog(binding)
641            .map_err(QueryError::execute)?
642        else {
643            return Ok(None);
644        };
645        self.execute_scalar_page_against_catalog(
646            request,
647            continuation,
648            DynamicReadLane::Public,
649            catalog,
650            ScalarPageMode::Live,
651            None,
652        )
653        .map(|(page, _)| Some(page))
654        .map_err(Self::live_page_error)
655    }
656
657    /// Execute one ordinary entity-name-driven bounded grouped read.
658    pub fn execute_public_dynamic_grouped_query(
659        &self,
660        request: &DynamicQuery,
661    ) -> Result<GroupedQueryOutput, QueryError> {
662        let catalog = self
663            .accepted_schema_catalog_context_for_entity_name(Some(request.entity()))
664            .map_err(QueryError::execute)?;
665        self.execute_dynamic_grouped_query_against_catalog(
666            request,
667            DynamicReadLane::Public,
668            catalog,
669        )
670    }
671
672    /// Execute one grouped typed read through the binding's immutable accepted
673    /// entity identity. `None` means the opaque binding is stale.
674    #[doc(hidden)]
675    pub fn execute_public_dynamic_grouped_query_for_typed_binding(
676        &self,
677        binding: &DynamicTypedEntityBinding,
678        request: &DynamicQuery,
679    ) -> Result<Option<GroupedQueryOutput>, QueryError> {
680        let Some(catalog) = self
681            .current_typed_entity_binding_catalog(binding)
682            .map_err(QueryError::execute)?
683        else {
684            return Ok(None);
685        };
686        self.execute_dynamic_grouped_query_against_catalog(
687            request,
688            DynamicReadLane::Public,
689            catalog,
690        )
691        .map(Some)
692    }
693
694    /// Execute one trusted entity-name-driven grouped read.
695    ///
696    /// This bypasses ordinary public admission but retains accepted-schema
697    /// planning, explicit grouped limits, cursor validation, and execution.
698    pub fn execute_trusted_dynamic_grouped_query(
699        &self,
700        request: &DynamicQuery,
701    ) -> Result<GroupedQueryOutput, QueryError> {
702        let catalog = self
703            .accepted_schema_catalog_context_for_entity_name(Some(request.entity()))
704            .map_err(QueryError::execute)?;
705        self.execute_dynamic_grouped_query_against_catalog(
706            request,
707            DynamicReadLane::Trusted,
708            catalog,
709        )
710    }
711
712    /// Execute one trusted revision-tolerant bounded dynamic page.
713    ///
714    /// Trusted execution bypasses public admission but retains the same
715    /// physical and aggregate request budgets as every other read lane.
716    pub fn execute_trusted_live_page(
717        &self,
718        request: &DynamicQuery,
719        continuation: Option<&str>,
720    ) -> Result<LiveQueryPageOutput, QueryError> {
721        let catalog = self
722            .accepted_schema_catalog_context_for_entity_name(Some(request.entity()))
723            .map_err(QueryError::execute)?;
724        self.execute_scalar_page_against_catalog(
725            request,
726            continuation,
727            DynamicReadLane::Trusted,
728            catalog,
729            ScalarPageMode::Live,
730            None,
731        )
732        .map(|(page, _)| page)
733        .map_err(Self::live_page_error)
734    }
735
736    #[cfg(test)]
737    pub(in crate::db) fn execute_trusted_live_page_with_result_bytes_limit_for_tests(
738        &self,
739        request: &DynamicQuery,
740        continuation: Option<&str>,
741        result_bytes_limit: u64,
742    ) -> Result<LiveQueryPageOutput, QueryError> {
743        let previous = SCALAR_PAGE_RESULT_BYTES_LIMIT_OVERRIDE
744            .with(|limit| limit.replace(Some(result_bytes_limit)));
745        let _guard = ScalarPageResultBytesLimitGuard(previous);
746        self.execute_trusted_live_page(request, continuation)
747    }
748
749    /// Execute one revision-strict bounded dynamic page.
750    pub fn execute_public_exhaustive_page(
751        &self,
752        request: &DynamicQuery,
753        continuation: Option<&str>,
754        proof: Option<&ReadSetRevisionProof>,
755    ) -> Result<ExhaustiveQueryPageOutput, ExhaustiveReadError> {
756        let catalog =
757            self.accepted_schema_catalog_context_for_entity_name(Some(request.entity()))?;
758        let (page, proof) = self.execute_scalar_page_against_catalog(
759            request,
760            continuation,
761            DynamicReadLane::Public,
762            catalog,
763            ScalarPageMode::Exhaustive,
764            proof,
765        )?;
766        let proof = proof.ok_or(ReadSetRevisionError::NonCanonical)?;
767        Ok(ExhaustiveQueryPageOutput::from_live_page(page, proof))
768    }
769
770    /// Execute one exhaustive page through a typed binding's accepted identity.
771    #[doc(hidden)]
772    pub fn execute_public_exhaustive_page_for_typed_binding(
773        &self,
774        binding: &DynamicTypedEntityBinding,
775        request: &DynamicQuery,
776        continuation: Option<&str>,
777        proof: Option<&ReadSetRevisionProof>,
778    ) -> Result<Option<ExhaustiveQueryPageOutput>, ExhaustiveReadError> {
779        let Some(catalog) = self.current_typed_entity_binding_catalog(binding)? else {
780            return Ok(None);
781        };
782        let (page, proof) = self.execute_scalar_page_against_catalog(
783            request,
784            continuation,
785            DynamicReadLane::Public,
786            catalog,
787            ScalarPageMode::Exhaustive,
788            proof,
789        )?;
790        let proof = proof.ok_or(ReadSetRevisionError::NonCanonical)?;
791        Ok(Some(ExhaustiveQueryPageOutput::from_live_page(page, proof)))
792    }
793
794    /// Execute one trusted revision-strict bounded dynamic page.
795    pub fn execute_trusted_exhaustive_page(
796        &self,
797        request: &DynamicQuery,
798        continuation: Option<&str>,
799        proof: Option<&ReadSetRevisionProof>,
800    ) -> Result<ExhaustiveQueryPageOutput, ExhaustiveReadError> {
801        let catalog =
802            self.accepted_schema_catalog_context_for_entity_name(Some(request.entity()))?;
803        let (page, proof) = self.execute_scalar_page_against_catalog(
804            request,
805            continuation,
806            DynamicReadLane::Trusted,
807            catalog,
808            ScalarPageMode::Exhaustive,
809            proof,
810        )?;
811        let proof = proof.ok_or(ReadSetRevisionError::NonCanonical)?;
812        Ok(ExhaustiveQueryPageOutput::from_live_page(page, proof))
813    }
814
815    fn live_page_error(error: ExhaustiveReadError) -> QueryError {
816        match error {
817            ExhaustiveReadError::Query(error) => error,
818            ExhaustiveReadError::Revision(_) => QueryError::invariant(),
819        }
820    }
821}