1use crate::db::{QueryError, query::preparation::PreparationWork};
7use icydb_diagnostic_code::DiagnosticExecutionBudgetResource as Resource;
8
9use crate::db::predicate::MissingRowPolicy;
10use crate::{
11 db::{
12 access::{AccessPlan, SemanticIndexKeyItemRef},
13 predicate::{IndexCompileTarget, IndexCompileTargetKind, Predicate, PredicateProgram},
14 query::plan::{
15 AccessPlannedQuery, ContinuationPolicy, DistinctExecutionStrategy,
16 EffectiveRuntimeFilterProgram, ExecutionShapeSignature, GroupPlan,
17 GroupedAggregateExecutionSpec, GroupedDistinctExecutionStrategy, GroupedPlanStrategy,
18 LogicalPlan, PlannerRouteProfile, PredicatePushdownDiagnostics, QueryMode,
19 ResidualFilterContract, ResidualFilterShape, ResolvedOrder, ResolvedOrderField,
20 ResolvedOrderValueSource, ScalarPlan, StaticExecutionPlanningContract,
21 derive_logical_pushdown_eligibility,
22 expr::{
23 CompiledExpr, Expr, ProjectionSpec, compile_scalar_projection_expr_with_schema,
24 compile_scalar_projection_plan_with_schema,
25 },
26 extend_unique_grouped_aggregate_specs_from_expr, grouped_aggregate_execution_specs,
27 grouped_aggregate_specs_from_projection_spec, grouped_cursor_policy_violation,
28 grouped_plan_strategy, lower_direct_projection_layouts_with_schema,
29 lower_projection_identity, lower_projection_intent_with_schema,
30 residual_query_predicate_after_access_path_bounds,
31 residual_query_predicate_after_filtered_access_contract,
32 resolved_grouped_distinct_execution_strategy_with_schema_info,
33 },
34 schema::SchemaInfo,
35 },
36 error::InternalError,
37 value::Value,
38};
39
40impl QueryMode {
41 #[must_use]
43 pub const fn is_load(&self) -> bool {
44 match self {
45 Self::Load(_) => true,
46 Self::Delete(_) => false,
47 }
48 }
49
50 #[must_use]
52 pub const fn is_delete(&self) -> bool {
53 match self {
54 Self::Delete(_) => true,
55 Self::Load(_) => false,
56 }
57 }
58}
59
60impl LogicalPlan {
61 #[must_use]
63 pub(in crate::db) const fn scalar_semantics(&self) -> &ScalarPlan {
64 match self {
65 Self::Scalar(plan) => plan,
66 Self::Grouped(plan) => &plan.scalar,
67 }
68 }
69}
70
71impl AccessPlannedQuery {
72 #[must_use]
74 pub(in crate::db) const fn scalar_plan(&self) -> &ScalarPlan {
75 self.logical.scalar_semantics()
76 }
77
78 #[must_use]
81 pub(in crate::db) fn scalar_consistency(&self) -> MissingRowPolicy {
82 if self.access.has_selected_index_access_path() {
83 MissingRowPolicy::Error
88 } else {
89 self.scalar_plan().consistency
90 }
91 }
92
93 #[must_use]
95 pub(in crate::db) const fn grouped_plan(&self) -> Option<&GroupPlan> {
96 match &self.logical {
97 LogicalPlan::Scalar(_) => None,
98 LogicalPlan::Grouped(plan) => Some(plan),
99 }
100 }
101
102 pub(in crate::db) fn projection_spec(&self) -> Result<&ProjectionSpec, QueryError> {
105 self.static_execution_planning_contract
106 .as_ref()
107 .map(|contract| &contract.projection_spec)
108 .ok_or_else(QueryError::invariant)
109 }
110
111 pub(in crate::db) fn prepare_projection(
113 &self,
114 schema: &SchemaInfo,
115 work: &PreparationWork<'_>,
116 ) -> Result<ProjectionSpec, QueryError> {
117 lower_projection_intent_with_schema(schema, &self.logical, &self.projection_selection, work)
118 }
119
120 #[must_use]
122 pub(in crate::db::query) fn projection_spec_for_identity(&self) -> ProjectionSpec {
123 lower_projection_identity(&self.logical, &self.projection_selection)
124 }
125
126 pub(in crate::db) fn execution_preparation_predicate(
133 &self,
134 work: &PreparationWork<'_>,
135 ) -> Result<Option<Predicate>, QueryError> {
136 if let Some(static_contract) = self.static_execution_planning_contract.as_ref() {
137 return static_contract
138 .execution_preparation_predicate
139 .as_ref()
140 .map(|predicate| work.copy_predicate(predicate))
141 .transpose();
142 }
143
144 let predicate = self
145 .scalar_plan()
146 .predicate
147 .as_ref()
148 .map(|predicate| work.copy_predicate(predicate))
149 .transpose()?;
150 derive_execution_preparation_predicate(&self.access, predicate, work)
151 .map_err(QueryError::execute)
152 }
153
154 pub(in crate::db) fn residual_filter_contract(
156 &self,
157 ) -> Result<&ResidualFilterContract, InternalError> {
158 Ok(&self
159 .require_static_execution_planning_contract()?
160 .residual_filter_contract)
161 }
162
163 pub(in crate::db) fn effective_execution_predicate(
165 &self,
166 ) -> Result<Option<&Predicate>, InternalError> {
167 Ok(self.residual_filter_contract()?.residual_filter_predicate())
168 }
169
170 pub(in crate::db) fn has_residual_filter_predicate(&self) -> Result<bool, InternalError> {
172 Ok(self.effective_execution_predicate()?.is_some())
173 }
174
175 pub(in crate::db) fn residual_filter_expr(&self) -> Result<Option<&Expr>, InternalError> {
177 Ok(self.residual_filter_contract()?.residual_filter_expr())
178 }
179
180 pub(in crate::db) fn residual_filter_shape(
182 &self,
183 ) -> Result<ResidualFilterShape, InternalError> {
184 Ok(self.residual_filter_contract()?.shape())
185 }
186
187 pub(in crate::db) fn prepare_residual_filter_shape(
189 &self,
190 budget: &dyn crate::db::query::construction::ConstructionBudget,
191 ) -> Result<ResidualFilterShape, InternalError> {
192 if self.has_static_execution_planning_contract() {
193 return self.residual_filter_shape();
194 }
195 let predicate = self
196 .scalar_plan()
197 .predicate
198 .as_ref()
199 .map(|predicate| budget.copy_predicate(predicate))
200 .transpose()?;
201 Ok(
202 residual_filter_facts_for_access(self.scalar_plan(), &self.access, predicate, budget)?
203 .0,
204 )
205 }
206
207 #[cfg(feature = "sql")]
209 pub(in crate::db) fn predicate_pushdown_diagnostics(
210 &self,
211 ) -> Result<PredicatePushdownDiagnostics, InternalError> {
212 Ok(self
213 .require_static_execution_planning_contract()?
214 .predicate_pushdown_diagnostics)
215 }
216
217 #[must_use]
219 pub(in crate::db) fn execution_preparation_compiled_predicate(
220 &self,
221 ) -> Option<&PredicateProgram> {
222 self.static_execution_planning_contract()?
223 .execution_preparation_compiled_predicate
224 .as_ref()
225 }
226
227 #[must_use]
229 pub(in crate::db) fn effective_runtime_compiled_predicate(&self) -> Option<&PredicateProgram> {
230 match self
231 .static_execution_planning_contract()?
232 .residual_filter_contract
233 .effective_runtime_filter_program()
234 {
235 Some(program) => program.predicate_program(),
236 None => None,
237 }
238 }
239
240 #[must_use]
242 pub(in crate::db) fn effective_runtime_filter_program(
243 &self,
244 ) -> Option<&EffectiveRuntimeFilterProgram> {
245 self.static_execution_planning_contract()?
246 .residual_filter_contract
247 .effective_runtime_filter_program()
248 }
249
250 #[must_use]
252 pub(in crate::db) fn distinct_execution_strategy(&self) -> DistinctExecutionStrategy {
253 if !self.scalar_plan().distinct {
254 return DistinctExecutionStrategy::None;
255 }
256
257 match distinct_runtime_dedup_strategy(&self.access) {
261 Some(strategy) => strategy,
262 None => DistinctExecutionStrategy::None,
263 }
264 }
265
266 pub(in crate::db) fn finalize_planner_route_profile_for_model_with_schema(
268 &mut self,
269 schema_info: &SchemaInfo,
270 work: &PreparationWork<'_>,
271 ) -> Result<(), InternalError> {
272 self.set_planner_route_profile(project_planner_route_profile_for_schema(
273 schema_info,
274 self,
275 work,
276 )?);
277 Ok(())
278 }
279
280 pub(in crate::db) fn finalize_static_execution_planning_contract_with_schema(
284 &mut self,
285 schema_info: &SchemaInfo,
286 projection: ProjectionSpec,
287 work: &PreparationWork<'_>,
288 ) -> Result<(), QueryError> {
289 self.bind_group_field_slots_to_schema(schema_info, work)?;
290 self.static_execution_planning_contract =
291 Some(project_static_execution_planning_contract_with_schema(
292 schema_info,
293 self,
294 projection,
295 work,
296 )?);
297
298 Ok(())
299 }
300
301 fn bind_group_field_slots_to_schema(
304 &mut self,
305 schema_info: &SchemaInfo,
306 work: &PreparationWork<'_>,
307 ) -> Result<(), QueryError> {
308 let LogicalPlan::Grouped(grouped) = &mut self.logical else {
309 return Ok(());
310 };
311
312 let accepted_fields = grouped
313 .group
314 .group_fields
315 .resolve_with_schema(schema_info, work)?
316 .ok_or_else(|| QueryError::execute(InternalError::planner_executor_invariant()))?;
317 grouped.group.group_fields = accepted_fields;
318
319 Ok(())
320 }
321
322 pub(in crate::db) fn execution_shape_signature(
324 &self,
325 entity_path: &str,
326 budget: &dyn crate::db::query::construction::ConstructionBudget,
327 ) -> Result<ExecutionShapeSignature, InternalError> {
328 Ok(ExecutionShapeSignature::new(
329 self.continuation_signature(entity_path, budget)?,
330 ))
331 }
332
333 #[must_use]
335 pub(in crate::db) fn scalar_projection_plan(&self) -> Option<&[CompiledExpr]> {
336 self.static_execution_planning_contract()?
337 .scalar_projection_plan
338 .as_deref()
339 }
340
341 #[must_use]
343 pub(in crate::db) const fn has_static_execution_planning_contract(&self) -> bool {
344 self.static_execution_planning_contract.is_some()
345 }
346
347 pub(in crate::db) fn primary_key_names(&self) -> Result<&[String], InternalError> {
349 Ok(&self
350 .require_static_execution_planning_contract()?
351 .primary_key_names)
352 }
353
354 pub(in crate::db) fn projection_referenced_slots(&self) -> Result<&[usize], InternalError> {
356 Ok(self
357 .require_static_execution_planning_contract()?
358 .projection_referenced_slots
359 .as_slice())
360 }
361
362 pub(in crate::db) fn projection_is_model_identity(&self) -> Result<bool, InternalError> {
364 Ok(self
365 .require_static_execution_planning_contract()?
366 .projection_is_model_identity)
367 }
368
369 #[must_use]
371 pub(in crate::db) fn order_referenced_slots(&self) -> Option<&[usize]> {
372 self.static_execution_planning_contract()?
373 .order_referenced_slots
374 .as_deref()
375 }
376
377 #[must_use]
379 pub(in crate::db) fn resolved_order(&self) -> Option<&ResolvedOrder> {
380 self.static_execution_planning_contract()?
381 .resolved_order
382 .as_ref()
383 }
384
385 #[must_use]
387 pub(in crate::db) fn slot_map(&self) -> Option<&[usize]> {
388 self.static_execution_planning_contract()?
389 .slot_map
390 .as_deref()
391 }
392
393 #[must_use]
395 pub(in crate::db) fn grouped_aggregate_execution_specs(
396 &self,
397 ) -> Option<&[GroupedAggregateExecutionSpec]> {
398 self.static_execution_planning_contract()?
399 .grouped_aggregate_execution_specs
400 .as_deref()
401 }
402
403 #[must_use]
405 pub(in crate::db) fn grouped_distinct_execution_strategy(
406 &self,
407 ) -> Option<&GroupedDistinctExecutionStrategy> {
408 self.static_execution_planning_contract()?
409 .grouped_distinct_execution_strategy
410 .as_ref()
411 }
412
413 pub(in crate::db) fn frozen_projection_spec(&self) -> Result<&ProjectionSpec, InternalError> {
415 Ok(&self
416 .require_static_execution_planning_contract()?
417 .projection_spec)
418 }
419
420 #[must_use]
422 pub(in crate::db) fn frozen_direct_projection_slots(&self) -> Option<&[usize]> {
423 self.static_execution_planning_contract()?
424 .projection_direct_slots
425 .as_deref()
426 }
427
428 #[must_use]
430 pub(in crate::db) fn frozen_data_row_direct_projection_slots(&self) -> Option<&[usize]> {
431 self.static_execution_planning_contract()?
432 .projection_data_row_direct_slots
433 .as_deref()
434 }
435
436 #[must_use]
438 pub(in crate::db) fn index_compile_targets(&self) -> Option<&[IndexCompileTarget]> {
439 self.static_execution_planning_contract()?
440 .index_compile_targets
441 .as_deref()
442 }
443
444 const fn static_execution_planning_contract(&self) -> Option<&StaticExecutionPlanningContract> {
445 self.static_execution_planning_contract.as_ref()
446 }
447
448 fn require_static_execution_planning_contract(
449 &self,
450 ) -> Result<&StaticExecutionPlanningContract, InternalError> {
451 self.static_execution_planning_contract
452 .as_ref()
453 .ok_or_else(InternalError::query_executor_invariant)
454 }
455}
456
457fn distinct_runtime_dedup_strategy<K>(access: &AccessPlan<K>) -> Option<DistinctExecutionStrategy> {
458 match access {
459 AccessPlan::Union(_) | AccessPlan::Intersection(_) => {
460 Some(DistinctExecutionStrategy::PreOrdered)
461 }
462 AccessPlan::Path(path) if path.as_ref().is_index_multi_lookup() => {
463 Some(DistinctExecutionStrategy::HashMaterialize)
464 }
465 AccessPlan::Path(_) => None,
466 }
467}
468
469fn derive_continuation_policy_validated(plan: &AccessPlannedQuery) -> ContinuationPolicy {
470 let is_grouped_safe = plan
471 .grouped_plan()
472 .is_none_or(|grouped| grouped_cursor_policy_violation(grouped, true).is_none());
473
474 ContinuationPolicy::new(
475 true, true, is_grouped_safe,
478 )
479}
480
481pub(in crate::db) fn project_planner_route_profile_for_schema(
483 schema_info: &SchemaInfo,
484 plan: &AccessPlannedQuery,
485 work: &PreparationWork<'_>,
486) -> Result<PlannerRouteProfile, InternalError> {
487 let secondary_order_contract = plan
488 .scalar_plan()
489 .order
490 .as_ref()
491 .map(|order| {
492 order.deterministic_secondary_order_contract_fields(
493 schema_info.shared_primary_key_names(),
494 work,
495 )
496 })
497 .transpose()?
498 .flatten();
499
500 Ok(PlannerRouteProfile::new(
501 derive_continuation_policy_validated(plan),
502 derive_logical_pushdown_eligibility(plan, secondary_order_contract.as_ref()),
503 secondary_order_contract,
504 ))
505}
506
507fn project_static_execution_planning_contract_with_schema(
508 schema_info: &SchemaInfo,
509 plan: &AccessPlannedQuery,
510 projection_spec: ProjectionSpec,
511 work: &PreparationWork<'_>,
512) -> Result<StaticExecutionPlanningContract, QueryError> {
513 let execution_preparation_predicate = plan.execution_preparation_predicate(work)?;
514 let residual_input = execution_preparation_predicate
517 .as_ref()
518 .map(|predicate| work.copy_predicate(predicate))
519 .transpose()?;
520 let (residual_filter_predicate, access_satisfied) =
521 derive_residual_filter_predicate_from_preparation(
522 plan.scalar_plan(),
523 &plan.access,
524 residual_input,
525 work,
526 )
527 .map_err(QueryError::execute)?;
528 let residual_filter_expr = derive_residual_filter_expr(plan, access_satisfied);
529 let effective_runtime_filter_program = compile_effective_runtime_filter_program(
530 schema_info,
531 residual_filter_expr.as_ref(),
532 residual_filter_predicate.as_ref(),
533 work,
534 )
535 .map_err(QueryError::execute)?;
536 let residual_filter_contract = ResidualFilterContract::new(
537 residual_filter_expr,
538 residual_filter_predicate,
539 effective_runtime_filter_program,
540 );
541 let residual_filter_shape = residual_filter_contract.shape();
542 let execution_preparation_compiled_predicate =
543 (should_compile_execution_preparation_predicate(residual_filter_shape)
544 && !planner_predicate_requires_expression_runtime(plan.scalar_plan()))
545 .then(|| compile_optional_predicate(schema_info, execution_preparation_predicate.as_ref()))
546 .flatten();
547 let predicate_pushdown_diagnostics =
548 derive_predicate_pushdown_diagnostics(plan, residual_filter_shape);
549 let scalar_projection_plan = if plan.grouped_plan().is_none() {
550 Some(
551 compile_scalar_projection_plan_with_schema(schema_info, &projection_spec, work)
552 .map_err(QueryError::execute)?
553 .ok_or_else(|| QueryError::execute(InternalError::query_executor_invariant()))?,
554 )
555 } else {
556 None
557 };
558 let (grouped_aggregate_execution_specs, grouped_distinct_execution_strategy) =
559 resolve_grouped_static_planning_semantics(schema_info, plan, &projection_spec, work)
560 .map_err(QueryError::execute)?;
561 let (projection_direct_slots, projection_data_row_direct_slots) =
562 lower_direct_projection_layouts_with_schema(
563 schema_info,
564 &plan.logical,
565 &projection_spec,
566 work,
567 )?;
568 let projection_referenced_slots =
569 projection_spec.referenced_slots_for_schema(schema_info, work)?;
570 let projection_is_model_identity = projection_spec.is_schema_identity_for(schema_info, work)?;
571 let resolved_order = resolved_order_for_plan(schema_info, plan, work)?;
572 let order_referenced_slots = resolved_order
573 .as_ref()
574 .map(|order| order.referenced_slots(work))
575 .transpose()?;
576 let (slot_map, index_compile_targets) =
577 index_execution_metadata_for_schema_plan(schema_info, plan, work)?
578 .map_or((None, None), |(slots, targets)| {
579 (Some(slots), Some(targets))
580 });
581
582 Ok(StaticExecutionPlanningContract {
583 primary_key_names: schema_info.shared_primary_key_names(),
584 projection_spec,
585 execution_preparation_predicate,
586 execution_preparation_compiled_predicate,
587 residual_filter_contract,
588 predicate_pushdown_diagnostics,
589 scalar_projection_plan,
590 grouped_aggregate_execution_specs,
591 grouped_distinct_execution_strategy,
592 projection_direct_slots,
593 projection_data_row_direct_slots,
594 projection_referenced_slots,
595 projection_is_model_identity,
596 resolved_order,
597 order_referenced_slots,
598 slot_map,
599 index_compile_targets,
600 })
601}
602
603fn compile_effective_runtime_filter_program(
607 schema_info: &SchemaInfo,
608 residual_filter_expr: Option<&Expr>,
609 residual_filter_predicate: Option<&Predicate>,
610 work: &PreparationWork<'_>,
611) -> Result<Option<EffectiveRuntimeFilterProgram>, InternalError> {
612 if let Some(predicate) = residual_filter_predicate {
617 return Ok(Some(EffectiveRuntimeFilterProgram::predicate(
618 PredicateProgram::compile_with_schema_info(schema_info, predicate),
619 )));
620 }
621
622 if let Some(filter_expr) = residual_filter_expr {
623 let compiled = compile_scalar_projection_expr_with_schema(schema_info, filter_expr, work)?
624 .ok_or_else(InternalError::query_invalid_logical_plan)?;
625
626 return Ok(Some(EffectiveRuntimeFilterProgram::expression(compiled)));
627 }
628
629 Ok(None)
630}
631
632fn derive_execution_preparation_predicate(
636 access: &AccessPlan<Value>,
637 query_predicate: Option<Predicate>,
638 budget: &dyn crate::db::query::construction::ConstructionBudget,
639) -> Result<Option<Predicate>, InternalError> {
640 let Some(query_predicate) = query_predicate else {
641 return Ok(None);
642 };
643 match access.selected_index_contract() {
644 Some(index) => {
645 residual_query_predicate_after_filtered_access_contract(index, query_predicate, budget)
646 }
647 None => Ok(Some(query_predicate)),
648 }
649}
650
651fn derive_residual_filter_predicate(
655 scalar: &ScalarPlan,
656 access: &AccessPlan<Value>,
657 query_predicate: Option<Predicate>,
658 budget: &dyn crate::db::query::construction::ConstructionBudget,
659) -> Result<Option<Predicate>, InternalError> {
660 let filtered = derive_execution_preparation_predicate(access, query_predicate, budget)?;
661 Ok(derive_residual_filter_predicate_from_preparation(scalar, access, filtered, budget)?.0)
662}
663
664fn derive_residual_filter_predicate_from_preparation(
665 scalar: &ScalarPlan,
666 access: &AccessPlan<Value>,
667 execution_preparation_predicate: Option<Predicate>,
668 budget: &dyn crate::db::query::construction::ConstructionBudget,
669) -> Result<(Option<Predicate>, bool), InternalError> {
670 let Some(execution_preparation_predicate) = execution_preparation_predicate else {
671 return Ok((None, false));
672 };
673
674 let residual = residual_query_predicate_after_access_path_bounds(
675 access.as_path(),
676 execution_preparation_predicate,
677 budget,
678 )?;
679 if residual.is_some() && planner_predicate_requires_expression_runtime(scalar) {
680 return Ok((None, false));
681 }
682
683 let access_satisfied = residual.is_none();
684 Ok((residual, access_satisfied))
685}
686
687fn derive_residual_filter_expr(plan: &AccessPlannedQuery, access_satisfied: bool) -> Option<Expr> {
691 let filter_expr = plan.scalar_plan().filter_expr.as_ref()?;
692 if derive_semantic_filter_fully_satisfied_by_access_contract(plan.scalar_plan())
693 && (!planner_predicate_requires_expression_runtime(plan.scalar_plan()) || access_satisfied)
694 {
695 return None;
696 }
697
698 Some(filter_expr.clone())
699}
700
701fn planner_predicate_requires_expression_runtime(scalar: &ScalarPlan) -> bool {
706 scalar.predicate_covers_filter_expr
707 && scalar
708 .filter_expr
709 .as_ref()
710 .is_some_and(Expr::contains_field_path)
711}
712
713fn derive_predicate_pushdown_diagnostics(
716 plan: &AccessPlannedQuery,
717 residual_filter_shape: ResidualFilterShape,
718) -> PredicatePushdownDiagnostics {
719 PredicatePushdownDiagnostics::from_plan(
720 plan.scalar_plan().filter_expr.is_some(),
721 plan.scalar_plan().predicate_covers_filter_expr,
722 plan.scalar_plan().predicate.as_ref(),
723 &plan.access,
724 residual_filter_shape,
725 )
726}
727
728const fn derive_semantic_filter_fully_satisfied_by_access_contract(scalar: &ScalarPlan) -> bool {
732 scalar.filter_expr.is_some()
733 && scalar.predicate.is_some()
734 && scalar.predicate_covers_filter_expr
735}
736
737pub(in crate::db::query) fn residual_filter_facts_for_access(
741 scalar: &ScalarPlan,
742 access: &AccessPlan<Value>,
743 query_predicate: Option<Predicate>,
744 budget: &dyn crate::db::query::construction::ConstructionBudget,
745) -> Result<(ResidualFilterShape, Option<Predicate>), InternalError> {
746 let predicate = derive_residual_filter_predicate(scalar, access, query_predicate, budget)?;
747 let expression_required = scalar.filter_expr.is_some()
750 && !derive_semantic_filter_fully_satisfied_by_access_contract(scalar);
751
752 Ok((
753 ResidualFilterShape::from_presence(expression_required, predicate.is_some()),
754 predicate,
755 ))
756}
757
758fn compile_optional_predicate(
761 schema_info: &SchemaInfo,
762 predicate: Option<&Predicate>,
763) -> Option<PredicateProgram> {
764 predicate.map(|predicate| PredicateProgram::compile_with_schema_info(schema_info, predicate))
765}
766
767const fn should_compile_execution_preparation_predicate(
772 residual_filter_shape: ResidualFilterShape,
773) -> bool {
774 !residual_filter_shape.is_absent()
775}
776
777fn resolve_grouped_static_planning_semantics(
781 schema_info: &SchemaInfo,
782 plan: &AccessPlannedQuery,
783 projection_spec: &ProjectionSpec,
784 work: &PreparationWork<'_>,
785) -> Result<
786 (
787 Option<Vec<GroupedAggregateExecutionSpec>>,
788 Option<GroupedDistinctExecutionStrategy>,
789 ),
790 InternalError,
791> {
792 let Some(grouped) = plan.grouped_plan() else {
793 return Ok((None, None));
794 };
795
796 let mut aggregate_specs =
797 grouped_aggregate_specs_from_projection_spec(projection_spec, &grouped.group.group_fields)?;
798 extend_grouped_having_aggregate_specs(&mut aggregate_specs, grouped)?;
799
800 let grouped_aggregate_execution_specs = Some(grouped_aggregate_execution_specs(
801 schema_info,
802 aggregate_specs,
803 work,
804 )?);
805 let grouped_distinct_execution_strategy = Some(
806 resolved_grouped_distinct_execution_strategy_with_schema_info(
807 schema_info,
808 &grouped.group.group_fields,
809 grouped.group.aggregates.as_slice(),
810 grouped.having_expr.as_ref(),
811 )?,
812 );
813
814 Ok((
815 grouped_aggregate_execution_specs,
816 grouped_distinct_execution_strategy,
817 ))
818}
819
820fn extend_grouped_having_aggregate_specs(
821 aggregate_specs: &mut Vec<GroupedAggregateExecutionSpec>,
822 grouped: &GroupPlan,
823) -> Result<(), InternalError> {
824 if let Some(having_expr) = grouped.having_expr.as_ref() {
825 extend_unique_grouped_aggregate_specs_from_expr(aggregate_specs, having_expr)?;
826 }
827
828 Ok(())
829}
830
831fn resolved_order_for_plan(
832 schema_info: &SchemaInfo,
833 plan: &AccessPlannedQuery,
834 work: &PreparationWork<'_>,
835) -> Result<Option<ResolvedOrder>, QueryError> {
836 if grouped_plan_strategy(plan, || plan.prepare_residual_filter_shape(work))
837 .map_err(QueryError::execute)?
838 .is_some_and(GroupedPlanStrategy::is_top_k_group)
839 {
840 return Ok(None);
841 }
842
843 let Some(order) = plan.scalar_plan().order.as_ref() else {
844 return Ok(None);
845 };
846
847 let mut fields = work.vec_with_capacity(order.fields.len())?;
848 for term in &order.fields {
849 fields.push(ResolvedOrderField::new(
850 resolved_order_value_source_for_term(schema_info, term, work)?,
851 term.direction(),
852 ));
853 }
854
855 Ok(Some(ResolvedOrder::new(fields)))
856}
857
858fn resolved_order_value_source_for_term(
859 schema_info: &SchemaInfo,
860 term: &crate::db::query::plan::OrderTerm,
861 work: &PreparationWork<'_>,
862) -> Result<ResolvedOrderValueSource, QueryError> {
863 if let Some(field) = term.direct_field() {
864 work.charge(Resource::PredicateExpressionSteps, 1 + field.len() as u64)?;
865 let slot = schema_info
866 .field_slot_index(field)
867 .ok_or_else(|| QueryError::execute(InternalError::query_invalid_logical_plan()))?;
868
869 return Ok(ResolvedOrderValueSource::direct_field(slot));
870 }
871
872 validate_resolved_order_scalar_seam(term.expr(), work)?;
873 let compiled = compile_scalar_projection_expr_with_schema(schema_info, term.expr(), work)
876 .map_err(QueryError::execute)?
877 .ok_or_else(|| QueryError::execute(InternalError::query_invalid_logical_plan()))?;
878
879 Ok(ResolvedOrderValueSource::expression(compiled))
880}
881
882fn validate_resolved_order_scalar_seam(
885 expr: &Expr,
886 work: &PreparationWork<'_>,
887) -> Result<(), QueryError> {
888 expr.try_for_each_tree_expr(&mut |node| {
889 work.charge(Resource::PredicateExpressionSteps, 1)?;
890 match node {
891 Expr::Aggregate(_) | Expr::Unary { .. } => Err(QueryError::execute(
892 InternalError::query_invalid_logical_plan(),
893 )),
894 #[cfg(test)]
895 Expr::Alias { .. } => Err(QueryError::execute(
896 InternalError::query_invalid_logical_plan(),
897 )),
898 _ => Ok(()),
899 }
900 })
901}
902
903type IndexExecutionMetadata = (Vec<usize>, Vec<IndexCompileTarget>);
904
905fn index_execution_metadata_for_schema_plan(
909 schema_info: &SchemaInfo,
910 plan: &AccessPlannedQuery,
911 work: &PreparationWork<'_>,
912) -> Result<Option<IndexExecutionMetadata>, QueryError> {
913 let executable = plan.access.executable_contract();
914 let Some(path) = executable.as_path() else {
915 return Ok(None);
916 };
917 let Some(key_items) = path.shape_facts().index_key_items_for_slot_map() else {
918 return Ok(None);
919 };
920 let mut slots = work.vec_with_capacity(key_items.key_arity())?;
921 let mut targets = work.vec_with_capacity(key_items.key_arity())?;
922
923 for (component_index, key_item) in key_items.key_items().iter().enumerate() {
924 let key_item = key_item.as_ref();
925 let field = key_item.field();
926 work.charge(Resource::PredicateExpressionSteps, 1 + field.len() as u64)?;
927 let root = field.split_once('.').map_or(field, |(root, _)| root);
928 let Some(field_slot) = schema_info.field_slot_index(root) else {
929 return Ok(None);
930 };
931 slots.push(field_slot);
932 targets.push(IndexCompileTarget {
933 component_index,
934 field_slot,
935 kind: match key_item {
936 SemanticIndexKeyItemRef::Field(_) => IndexCompileTargetKind::Field,
937 SemanticIndexKeyItemRef::AcceptedExpression(expression) => {
938 IndexCompileTargetKind::Expression(expression.op())
939 }
940 },
941 });
942 }
943
944 Ok(Some((slots, targets)))
945}