#[cfg(test)]
mod tests;
use crate::{
db::{
query::{
builder::AggregateExpr,
plan::{
AccessPlannedQuery, AggregateIdentity, AggregateKind, AggregateSemanticKeyRef,
FieldSlot, GlobalDistinctAggregateKind, GroupAggregateSpec,
GroupDistinctAdmissibility, GroupDistinctPolicyReason, GroupedExecutionConfig,
GroupedPlanStrategy,
expr::{
CompiledExpr, Expr, ProjectionSpec, compile_scalar_projection_expr_with_schema,
},
grouped_distinct_admissibility, grouped_plan_strategy,
resolve_global_distinct_field_aggregate, validate_grouped_projection_layout,
},
},
schema::SchemaInfo,
},
error::InternalError,
};
use icydb_diagnostic_code::DiagnosticExecutionBudgetResource as Resource;
#[derive(Clone, Debug, Eq, PartialEq)]
pub(in crate::db) struct PlannedProjectionLayout {
pub(in crate::db) group_field_positions: Vec<usize>,
pub(in crate::db) aggregate_positions: Vec<usize>,
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub(in crate::db) struct GroupedAggregateExecutionSpec {
identity: AggregateIdentity,
target_slot: Option<FieldSlot>,
filter_expr: Option<Expr>,
compiled_input_expr: Option<CompiledExpr>,
compiled_filter_expr: Option<CompiledExpr>,
}
impl GroupedAggregateExecutionSpec {
fn compile_attached_scalar_expr(
schema_info: &SchemaInfo,
expr: &Expr,
budget: &dyn crate::db::query::construction::ConstructionBudget,
) -> Result<CompiledExpr, InternalError> {
let scalar = compile_scalar_projection_expr_with_schema(schema_info, expr, budget)?
.ok_or_else(InternalError::planner_executor_invariant)?;
Ok(scalar)
}
#[must_use]
pub(in crate::db) fn from_aggregate_expr(aggregate_expr: &AggregateExpr) -> Self {
Self {
identity: AggregateIdentity::from_aggregate_expr(aggregate_expr),
target_slot: None,
filter_expr: aggregate_expr.filter_expr().cloned(),
compiled_input_expr: None,
compiled_filter_expr: None,
}
}
#[must_use]
pub(in crate::db) fn from_uncompiled_inputs(
kind: AggregateKind,
target_slot: Option<FieldSlot>,
input_expr: Option<Expr>,
filter_expr: Option<Expr>,
distinct: bool,
) -> Self {
Self {
identity: AggregateIdentity::from_kind_input_and_distinct(kind, input_expr, distinct),
target_slot,
filter_expr,
compiled_input_expr: None,
compiled_filter_expr: None,
}
}
pub(in crate::db) fn resolve_with_schema_info(
&mut self,
schema_info: &SchemaInfo,
budget: &dyn crate::db::query::construction::ConstructionBudget,
) -> Result<(), InternalError> {
budget.charge(Resource::PredicateExpressionSteps, 1)?;
let compiled_input_expr = self
.input_expr()
.map(|expr| Self::compile_attached_scalar_expr(schema_info, expr, budget))
.transpose()?;
let compiled_filter_expr = self
.filter_expr()
.map(|expr| Self::compile_attached_scalar_expr(schema_info, expr, budget))
.transpose()?;
let target_slot = self
.target_field()
.map(|field| {
resolve_aggregate_target_field_slot_from_schema(schema_info, field)
.ok_or_else(InternalError::planner_executor_invariant)
})
.transpose()?;
self.target_slot = target_slot;
self.compiled_input_expr = compiled_input_expr;
self.compiled_filter_expr = compiled_filter_expr;
Ok(())
}
#[must_use]
pub(in crate::db) const fn kind(&self) -> AggregateKind {
self.identity.kind()
}
#[must_use]
pub(in crate::db) const fn target_field(&self) -> Option<&str> {
match self.input_expr() {
Some(Expr::Field(field_id)) => Some(field_id.as_str()),
_ => None,
}
}
#[must_use]
pub(in crate::db) const fn target_slot(&self) -> Option<&FieldSlot> {
self.target_slot.as_ref()
}
#[must_use]
pub(in crate::db) const fn input_expr(&self) -> Option<&Expr> {
self.identity.input_expr()
}
#[must_use]
pub(in crate::db) const fn filter_expr(&self) -> Option<&Expr> {
self.filter_expr.as_ref()
}
#[must_use]
pub(in crate::db) fn semantic_key(&self) -> AggregateSemanticKeyRef<'_> {
AggregateSemanticKeyRef::new(
self.kind(),
self.input_expr(),
self.filter_expr(),
self.distinct(),
)
}
#[must_use]
pub(in crate::db) const fn distinct(&self) -> bool {
self.identity.distinct()
}
#[must_use]
pub(in crate::db) const fn uses_grouped_distinct_value_dedup(&self) -> bool {
self.identity.uses_grouped_distinct_value_dedup()
}
#[must_use]
pub(in crate::db) const fn compiled_input_expr(&self) -> Option<&CompiledExpr> {
self.compiled_input_expr.as_ref()
}
#[must_use]
pub(in crate::db) const fn compiled_filter_expr(&self) -> Option<&CompiledExpr> {
self.compiled_filter_expr.as_ref()
}
#[must_use]
pub(in crate::db) const fn admits_count_rows_dedicated_fold(&self) -> bool {
matches!(self.kind(), AggregateKind::Count)
&& self.input_expr().is_none()
&& self.filter_expr().is_none()
&& !self.distinct()
}
#[cfg(test)]
#[must_use]
pub(in crate::db) fn from_test_inputs(
kind: AggregateKind,
target_slot: Option<FieldSlot>,
target_field: Option<&str>,
distinct: bool,
) -> Self {
Self {
identity: AggregateIdentity::from_kind_input_and_distinct(
kind,
target_field
.map(|field| Expr::Field(crate::db::query::plan::expr::FieldId::new(field))),
distinct,
),
target_slot,
filter_expr: None,
compiled_input_expr: None,
compiled_filter_expr: None,
}
}
}
impl PlannedProjectionLayout {
#[must_use]
pub(in crate::db) const fn group_field_positions(&self) -> &[usize] {
self.group_field_positions.as_slice()
}
#[must_use]
pub(in crate::db) const fn aggregate_positions(&self) -> &[usize] {
self.aggregate_positions.as_slice()
}
pub(in crate::db) fn group_field_positions_not_strictly_increasing() -> InternalError {
InternalError::planner_executor_invariant()
}
pub(in crate::db) fn aggregate_positions_not_strictly_increasing() -> InternalError {
InternalError::planner_executor_invariant()
}
pub(in crate::db) fn group_fields_must_precede_aggregates() -> InternalError {
InternalError::planner_executor_invariant()
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(in crate::db) enum GroupedFoldPath {
CountRowsDedicated,
GenericReducers,
}
impl GroupedFoldPath {
#[must_use]
pub(in crate::db) fn from_plan_strategy(
strategy: GroupedPlanStrategy,
aggregate_specs: &[GroupedAggregateExecutionSpec],
) -> Self {
if strategy.is_single_count_rows()
&& aggregate_specs
.iter()
.all(GroupedAggregateExecutionSpec::admits_count_rows_dedicated_fold)
{
Self::CountRowsDedicated
} else {
Self::GenericReducers
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(in crate::db) enum GroupedExecutionRoute {
GlobalDistinctTopK,
GlobalDistinctFull,
CountRowsDedicated,
GenericTopK,
GenericFull,
}
impl GroupedExecutionRoute {
#[must_use]
pub(in crate::db) const fn from_planner_strategy(
strategy: GroupedPlanStrategy,
fold_path: GroupedFoldPath,
distinct_strategy: &GroupedDistinctExecutionStrategy,
) -> Self {
let uses_global_distinct = distinct_strategy.global_distinct_target_slot().is_some();
let uses_top_k = strategy.is_top_k_group();
match (uses_global_distinct, uses_top_k, fold_path) {
(true, true, _) => Self::GlobalDistinctTopK,
(true, false, _) => Self::GlobalDistinctFull,
(false, true, _) => Self::GenericTopK,
(false, false, GroupedFoldPath::CountRowsDedicated) => Self::CountRowsDedicated,
(false, false, GroupedFoldPath::GenericReducers) => Self::GenericFull,
}
}
#[must_use]
pub(in crate::db) const fn uses_global_distinct_fold(self) -> bool {
matches!(self, Self::GlobalDistinctTopK | Self::GlobalDistinctFull)
}
#[must_use]
pub(in crate::db) const fn uses_top_k_group_selection(self) -> bool {
matches!(self, Self::GlobalDistinctTopK | Self::GenericTopK)
}
#[must_use]
pub(in crate::db) const fn uses_count_rows_dedicated_fold(self) -> bool {
matches!(self, Self::CountRowsDedicated)
}
}
#[derive(Clone)]
pub(in crate::db) struct GroupedExecutorHandoff<'a> {
base: &'a AccessPlannedQuery,
group_fields: &'a crate::db::query::plan::GroupFieldSet,
grouped_aggregate_execution_specs: &'a [GroupedAggregateExecutionSpec],
projection_layout: PlannedProjectionLayout,
projection_is_identity: bool,
grouped_plan_strategy: GroupedPlanStrategy,
grouped_execution_route: GroupedExecutionRoute,
grouped_distinct_policy_contract: GroupedDistinctPolicyContract<'a>,
execution: GroupedExecutionConfig,
}
impl<'a> GroupedExecutorHandoff<'a> {
#[must_use]
pub(in crate::db) const fn base(&self) -> &'a AccessPlannedQuery {
self.base
}
#[must_use]
pub(in crate::db) const fn group_fields(&self) -> &'a crate::db::query::plan::GroupFieldSet {
self.group_fields
}
#[must_use]
pub(in crate::db) const fn projection_is_identity(&self) -> bool {
self.projection_is_identity
}
#[must_use]
pub(in crate::db) const fn grouped_plan_strategy(&self) -> GroupedPlanStrategy {
self.grouped_plan_strategy
}
#[must_use]
pub(in crate::db) const fn grouped_execution_route(&self) -> GroupedExecutionRoute {
self.grouped_execution_route
}
#[must_use]
pub(in crate::db) fn into_route_stage_residents(
self,
) -> (
Vec<GroupedAggregateExecutionSpec>,
PlannedProjectionLayout,
GroupedDistinctExecutionStrategy,
) {
(
self.grouped_aggregate_execution_specs.to_vec(),
self.projection_layout,
self.grouped_distinct_policy_contract
.execution_strategy()
.clone(),
)
}
#[must_use]
pub(in crate::db) const fn distinct_policy_violation_for_executor(
&self,
) -> Option<GroupDistinctPolicyReason> {
self.grouped_distinct_policy_contract
.violation_for_executor()
}
#[must_use]
pub(in crate::db) const fn execution(&self) -> GroupedExecutionConfig {
self.execution
}
}
pub(in crate::db) fn grouped_executor_handoff(
plan: &AccessPlannedQuery,
) -> Result<GroupedExecutorHandoff<'_>, InternalError> {
let Some(grouped) = plan.grouped_plan() else {
return Err(InternalError::planner_executor_invariant());
};
let projection_spec = plan.frozen_projection_spec()?;
let grouped_aggregate_execution_specs = plan
.grouped_aggregate_execution_specs()
.ok_or_else(InternalError::planner_executor_invariant)?;
let (projection_layout, projection_is_identity) = planned_projection_layout_from_spec(
projection_spec,
&grouped.group.group_fields,
grouped.group.aggregates.as_slice(),
grouped_aggregate_execution_specs,
)?;
validate_grouped_projection_layout(&projection_layout)?;
let grouped_plan_strategy = grouped_plan_strategy(plan, || plan.residual_filter_shape())?
.ok_or_else(InternalError::planner_executor_invariant)?;
let grouped_fold_path = if grouped.group.group_fields.as_path_aware().is_some() {
GroupedFoldPath::GenericReducers
} else {
GroupedFoldPath::from_plan_strategy(
grouped_plan_strategy,
grouped_aggregate_execution_specs,
)
};
let grouped_distinct_policy_contract = GroupedDistinctPolicyContract::new(
match grouped_distinct_admissibility(grouped.scalar.distinct, grouped.having_expr.is_some())
{
GroupDistinctAdmissibility::Allowed => None,
GroupDistinctAdmissibility::Disallowed(reason) => Some(reason),
},
plan.grouped_distinct_execution_strategy()
.ok_or_else(InternalError::planner_executor_invariant)?,
);
let grouped_execution_route = GroupedExecutionRoute::from_planner_strategy(
grouped_plan_strategy,
grouped_fold_path,
grouped_distinct_policy_contract.execution_strategy(),
);
Ok(GroupedExecutorHandoff {
base: plan,
group_fields: &grouped.group.group_fields,
grouped_aggregate_execution_specs,
projection_layout,
projection_is_identity,
grouped_plan_strategy,
grouped_execution_route,
grouped_distinct_policy_contract,
execution: grouped.group.execution,
})
}
pub(in crate::db) fn grouped_aggregate_execution_specs(
schema_info: &SchemaInfo,
mut aggregate_specs: Vec<GroupedAggregateExecutionSpec>,
budget: &dyn crate::db::query::construction::ConstructionBudget,
) -> Result<Vec<GroupedAggregateExecutionSpec>, InternalError> {
for aggregate_spec in &mut aggregate_specs {
aggregate_spec.resolve_with_schema_info(schema_info, budget)?;
}
Ok(aggregate_specs)
}
pub(in crate::db) fn grouped_aggregate_specs_from_projection_spec(
projection_spec: &ProjectionSpec,
group_fields: &crate::db::query::plan::GroupFieldSet,
) -> Result<Vec<GroupedAggregateExecutionSpec>, InternalError> {
#[cfg(not(test))]
let _ = group_fields;
#[cfg(test)]
validate_grouped_projection_references(projection_spec, group_fields)?;
let mut aggregate_specs = Vec::new();
for field in projection_spec.fields() {
extend_unique_grouped_aggregate_specs_from_expr(&mut aggregate_specs, field.expr())?;
}
Ok(aggregate_specs)
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub(in crate::db) struct GroupedDistinctPolicyContract<'a> {
violation_for_executor: Option<GroupDistinctPolicyReason>,
execution_strategy: &'a GroupedDistinctExecutionStrategy,
}
impl<'a> GroupedDistinctPolicyContract<'a> {
#[must_use]
const fn new(
violation_for_executor: Option<GroupDistinctPolicyReason>,
execution_strategy: &'a GroupedDistinctExecutionStrategy,
) -> Self {
Self {
violation_for_executor,
execution_strategy,
}
}
#[must_use]
pub(in crate::db) const fn execution_strategy(&self) -> &GroupedDistinctExecutionStrategy {
self.execution_strategy
}
#[must_use]
pub(in crate::db) const fn violation_for_executor(&self) -> Option<GroupDistinctPolicyReason> {
self.violation_for_executor
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
pub(in crate::db) enum GroupedDistinctExecutionStrategy {
None,
GlobalDistinctFieldCount { target_slot: FieldSlot },
GlobalDistinctFieldSum { target_slot: FieldSlot },
GlobalDistinctFieldAvg { target_slot: FieldSlot },
}
impl GroupedDistinctExecutionStrategy {
#[must_use]
pub(in crate::db) const fn global_distinct_target_slot(&self) -> Option<&FieldSlot> {
match self {
Self::None => None,
Self::GlobalDistinctFieldCount { target_slot }
| Self::GlobalDistinctFieldSum { target_slot }
| Self::GlobalDistinctFieldAvg { target_slot } => Some(target_slot),
}
}
#[must_use]
pub(in crate::db) const fn global_distinct_aggregate_kind(&self) -> Option<AggregateKind> {
match self {
Self::None => None,
Self::GlobalDistinctFieldCount { .. } => Some(AggregateKind::Count),
Self::GlobalDistinctFieldSum { .. } => Some(AggregateKind::Sum),
Self::GlobalDistinctFieldAvg { .. } => Some(AggregateKind::Avg),
}
}
#[must_use]
pub(in crate::db) const fn from_supported_global_distinct(
kind: GlobalDistinctAggregateKind,
target_slot: FieldSlot,
) -> Self {
match kind {
GlobalDistinctAggregateKind::Count => Self::GlobalDistinctFieldCount { target_slot },
GlobalDistinctAggregateKind::Sum => Self::GlobalDistinctFieldSum { target_slot },
GlobalDistinctAggregateKind::Avg => Self::GlobalDistinctFieldAvg { target_slot },
}
}
}
pub(in crate::db) fn resolved_grouped_distinct_execution_strategy_with_schema_info(
schema_info: &SchemaInfo,
group_fields: &crate::db::query::plan::GroupFieldSet,
aggregates: &[GroupAggregateSpec],
having_expr: Option<&crate::db::query::plan::expr::Expr>,
) -> Result<GroupedDistinctExecutionStrategy, InternalError> {
match resolve_global_distinct_field_aggregate(group_fields, aggregates, having_expr) {
Ok(Some(aggregate)) => {
let target_slot = resolve_aggregate_target_field_slot_from_schema(
schema_info,
aggregate.target_field(),
)
.ok_or_else(InternalError::planner_executor_invariant)?;
let distinct_kind = aggregate
.kind()
.global_distinct_kind()
.ok_or_else(InternalError::planner_executor_invariant)?;
Ok(
GroupedDistinctExecutionStrategy::from_supported_global_distinct(
distinct_kind,
target_slot,
),
)
}
Ok(None) => Ok(GroupedDistinctExecutionStrategy::None),
Err(reason) => Err(reason.into_planner_handoff_internal_error()),
}
}
fn resolve_aggregate_target_field_slot_from_schema(
schema_info: &SchemaInfo,
field: &str,
) -> Option<FieldSlot> {
FieldSlot::resolve_with_schema(schema_info, field)
}
fn planned_projection_layout_from_spec(
projection_spec: &ProjectionSpec,
group_fields: &crate::db::query::plan::GroupFieldSet,
aggregates: &[GroupAggregateSpec],
aggregate_specs: &[GroupedAggregateExecutionSpec],
) -> Result<(PlannedProjectionLayout, bool), InternalError> {
#[cfg(test)]
validate_grouped_projection_references(projection_spec, group_fields)?;
let mut group_field_positions = Vec::new();
let mut aggregate_positions = Vec::new();
let mut seen_aggregates = vec![false; aggregate_specs.len()];
let mut projection_is_identity =
projection_spec.len() == group_fields.len().saturating_add(aggregates.len());
let mut next_group_field_index = 0usize;
let mut next_aggregate_index = 0usize;
for (index, field) in projection_spec.fields().enumerate() {
let root_expr = expression_without_alias(field.expr());
let mut contains_aggregate = false;
let mut introduced_aggregate_count = 0usize;
root_expr.try_for_each_tree_aggregate(&mut |aggregate| {
contains_aggregate = true;
let key = AggregateSemanticKeyRef::from_aggregate_expr(aggregate);
let slot = aggregate_specs
.iter()
.position(|spec| spec.semantic_key() == key)
.ok_or_else(InternalError::planner_executor_invariant)?;
if !seen_aggregates[slot] {
seen_aggregates[slot] = true;
introduced_aggregate_count = introduced_aggregate_count.saturating_add(1);
}
Ok::<(), InternalError>(())
})?;
match root_expr {
Expr::Field(_) | Expr::FieldPath(_) => {
group_field_positions.push(index);
projection_is_identity &= next_aggregate_index == 0
&& group_fields
.get(next_group_field_index)
.is_some_and(|group_field| group_field.matches_expr(root_expr));
next_group_field_index = next_group_field_index.saturating_add(1);
}
Expr::Aggregate(aggregate_expr) => {
aggregate_positions.push(index);
let aggregate_key = AggregateSemanticKeyRef::from_aggregate_expr(aggregate_expr);
projection_is_identity &= next_group_field_index == group_fields.len()
&& aggregates
.get(next_aggregate_index)
.is_some_and(|aggregate| aggregate_key == aggregate.semantic_key());
next_aggregate_index =
next_aggregate_index.saturating_add(introduced_aggregate_count);
}
_ if contains_aggregate => {
aggregate_positions.push(index);
projection_is_identity = false;
next_aggregate_index =
next_aggregate_index.saturating_add(introduced_aggregate_count);
}
_ => {
group_field_positions.push(index);
projection_is_identity = false;
next_group_field_index = next_group_field_index.saturating_add(1);
}
}
}
projection_is_identity &=
next_group_field_index == group_fields.len() && next_aggregate_index == aggregates.len();
Ok((
PlannedProjectionLayout {
group_field_positions,
aggregate_positions,
},
projection_is_identity,
))
}
#[cfg(test)]
fn validate_grouped_projection_references(
projection_spec: &ProjectionSpec,
group_fields: &crate::db::query::plan::GroupFieldSet,
) -> Result<(), InternalError> {
crate::db::query::preparation::with_preparation_work(|work| {
for field in projection_spec.fields() {
let root_expr = expression_without_alias(field.expr());
if !group_fields.try_contains_all_expr_references(root_expr, &mut |steps| {
crate::db::query::construction::ConstructionBudget::charge(
work,
Resource::PredicateExpressionSteps,
steps,
)
})? {
return Err(InternalError::planner_executor_invariant());
}
}
Ok(())
})
}
pub(in crate::db::query::plan) fn extend_unique_grouped_aggregate_specs_from_expr(
aggregate_specs: &mut Vec<GroupedAggregateExecutionSpec>,
expr: &Expr,
) -> Result<(), InternalError> {
expr.try_for_each_tree_aggregate(&mut |aggregate_expr| {
push_unique_grouped_aggregate_spec(aggregate_specs, aggregate_expr);
Ok::<(), InternalError>(())
})
}
fn push_unique_grouped_aggregate_spec(
aggregate_specs: &mut Vec<GroupedAggregateExecutionSpec>,
aggregate_expr: &AggregateExpr,
) {
let aggregate_key = AggregateSemanticKeyRef::from_aggregate_expr(aggregate_expr);
if aggregate_specs
.iter()
.all(|current| current.semantic_key() != aggregate_key)
{
aggregate_specs.push(GroupedAggregateExecutionSpec::from_aggregate_expr(
aggregate_expr,
));
}
}
#[cfg_attr(
not(test),
expect(
clippy::missing_const_for_fn,
reason = "test-only alias stripping keeps the shared helper non-const across the full target matrix"
)
)]
fn expression_without_alias(expr: &Expr) -> &Expr {
#[cfg(test)]
{
let mut current = expr;
while let Expr::Alias { expr: inner, .. } = current {
current = inner.as_ref();
}
current
}
#[cfg(not(test))]
expr
}
crate::retained::retained_fields!(GroupedAggregateExecutionSpec {
Self{identity,target_slot,filter_expr,compiled_input_expr,compiled_filter_expr} => [identity,target_slot,filter_expr,compiled_input_expr,compiled_filter_expr],
});
crate::retained::retained_fields!(GroupedDistinctExecutionStrategy {
Self::None => [],
Self::GlobalDistinctFieldCount{target_slot} => [target_slot],
Self::GlobalDistinctFieldSum{target_slot} => [target_slot],
Self::GlobalDistinctFieldAvg{target_slot} => [target_slot],
});