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