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