1use 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 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 #[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 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 #[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 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 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 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 #[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 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}