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::{
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 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 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 #[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 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 #[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 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 #[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 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 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 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 #[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 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}