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