#[cfg(test)]
mod tests;
use crate::db::{
QueryError,
access::AccessPlan,
query::plan::{
AccessPlannedQuery, GroupAggregateSpec, GroupFieldSet, GroupPlan,
GroupedPlanAggregateFamily, OrderSpec, ResidualFilterShape,
expr::{
GroupedOrderTermAdmissibility, GroupedTopKOrderTermAdmissibility,
try_classify_grouped_order_term_for_field, try_classify_grouped_top_k_order_term,
try_grouped_top_k_order_term_requires_heap,
},
},
query::preparation::PreparationWork,
};
use crate::error::InternalError;
use icydb_diagnostic_code::DiagnosticExecutionBudgetResource as Resource;
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
enum GroupedPlanFamily {
Hash,
Ordered,
TopK,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(in crate::db) enum GroupedPlanFallbackReason {
DistinctGroupingNotAdmitted,
ResidualFilterBlocksGroupedOrder,
AggregateStreamingNotSupported,
HavingBlocksGroupedOrder,
GroupKeyOrderPrefixMismatch,
GroupKeyOrderDirectionMismatch,
GroupKeyOrderExpressionNotAdmissible,
GroupKeyOrderUnavailable,
}
impl GroupedPlanFallbackReason {
#[must_use]
pub(in crate::db) const fn code(self) -> &'static str {
match self {
Self::DistinctGroupingNotAdmitted => "distinct_grouping_not_admitted",
Self::ResidualFilterBlocksGroupedOrder => "residual_filter_blocks_grouped_order",
Self::AggregateStreamingNotSupported => "aggregate_streaming_not_supported",
Self::HavingBlocksGroupedOrder => "having_blocks_grouped_order",
Self::GroupKeyOrderPrefixMismatch => "group_key_order_prefix_mismatch",
Self::GroupKeyOrderDirectionMismatch => "group_key_order_direction_mismatch",
Self::GroupKeyOrderExpressionNotAdmissible => {
"group_key_order_expression_not_admissible"
}
Self::GroupKeyOrderUnavailable => "group_key_order_unavailable",
}
}
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
pub(in crate::db) struct GroupedPlanStrategy {
family: GroupedPlanFamily,
aggregate_family: GroupedPlanAggregateFamily,
fallback_reason: Option<GroupedPlanFallbackReason>,
}
impl GroupedPlanStrategy {
#[must_use]
pub(in crate::db) const fn code(self) -> &'static str {
match self.family {
GroupedPlanFamily::Hash => "hash_group",
GroupedPlanFamily::Ordered => "ordered_group",
GroupedPlanFamily::TopK => "top_k_group",
}
}
#[must_use]
pub(in crate::db) const fn hash_group_with_aggregate_family(
reason: GroupedPlanFallbackReason,
aggregate_family: GroupedPlanAggregateFamily,
) -> Self {
Self {
family: GroupedPlanFamily::Hash,
aggregate_family,
fallback_reason: Some(reason),
}
}
#[must_use]
pub(in crate::db) const fn ordered_group_with_aggregate_family(
aggregate_family: GroupedPlanAggregateFamily,
) -> Self {
Self {
family: GroupedPlanFamily::Ordered,
aggregate_family,
fallback_reason: None,
}
}
#[must_use]
pub(in crate::db) const fn top_k_group_with_aggregate_family(
aggregate_family: GroupedPlanAggregateFamily,
) -> Self {
Self {
family: GroupedPlanFamily::TopK,
aggregate_family,
fallback_reason: None,
}
}
#[must_use]
pub(in crate::db) const fn is_ordered_group(self) -> bool {
matches!(self.family, GroupedPlanFamily::Ordered)
}
#[must_use]
pub(in crate::db) const fn is_top_k_group(self) -> bool {
matches!(self.family, GroupedPlanFamily::TopK)
}
#[must_use]
pub(in crate::db) const fn ordered_group_admitted(self) -> bool {
self.is_ordered_group()
}
#[must_use]
pub(in crate::db) const fn aggregate_family(self) -> GroupedPlanAggregateFamily {
self.aggregate_family
}
#[must_use]
pub(in crate::db) const fn is_single_count_rows(self) -> bool {
matches!(
self.aggregate_family,
GroupedPlanAggregateFamily::CountRowsOnly
)
}
#[must_use]
pub(in crate::db) const fn fallback_reason(self) -> Option<GroupedPlanFallbackReason> {
self.fallback_reason
}
}
pub(in crate::db) fn grouped_plan_strategy(
plan: &AccessPlannedQuery,
residual: impl FnOnce() -> Result<ResidualFilterShape, InternalError>,
) -> Result<Option<GroupedPlanStrategy>, InternalError> {
plan.grouped_plan()
.map(|grouped| derive_grouped_plan_strategy(plan, grouped, residual, &mut |_| Ok(())))
.transpose()
}
pub(in crate::db) fn grouped_plan_strategy_for_explain(
plan: &AccessPlannedQuery,
grouped: &GroupPlan,
work: &PreparationWork<'_>,
) -> Result<GroupedPlanStrategy, QueryError> {
derive_grouped_plan_strategy(
plan,
grouped,
|| {
plan.prepare_residual_filter_shape(work)
.map_err(QueryError::execute)
},
&mut |steps| work.charge(Resource::PredicateExpressionSteps, steps),
)
}
fn derive_grouped_plan_strategy<E>(
plan: &AccessPlannedQuery,
grouped: &GroupPlan,
residual: impl FnOnce() -> Result<ResidualFilterShape, E>,
observe: &mut impl FnMut(u64) -> Result<(), E>,
) -> Result<GroupedPlanStrategy, E> {
observe(1)?;
let aggregate_family = GroupedPlanAggregateFamily::try_from_grouped_aggregates(
grouped.group.aggregates.as_slice(),
observe,
)?;
let order_strategy_projection = grouped_order_strategy_projection(
grouped.scalar.order.as_ref(),
&grouped.group.group_fields,
observe,
)?;
if grouped.scalar.distinct {
return Ok(hash_group_fallback_strategy(
GroupedPlanFallbackReason::DistinctGroupingNotAdmitted,
aggregate_family,
));
}
if matches!(
order_strategy_projection,
GroupedOrderStrategyProjection::TopK
) {
return Ok(GroupedPlanStrategy::top_k_group_with_aggregate_family(
aggregate_family,
));
}
if !residual()?.is_absent() && grouped.group.group_fields.as_path_aware().is_none() {
return Ok(hash_group_fallback_strategy(
GroupedPlanFallbackReason::ResidualFilterBlocksGroupedOrder,
aggregate_family,
));
}
if !matches!(
order_strategy_projection,
GroupedOrderStrategyProjection::TopK
) && !grouped_aggregates_streaming_compatible(grouped.group.aggregates.as_slice(), observe)?
{
return Ok(hash_group_fallback_strategy(
GroupedPlanFallbackReason::AggregateStreamingNotSupported,
aggregate_family,
));
}
if !crate::db::query::plan::semantics::group_having::grouped_having_streaming_compatible(
grouped.having_expr.as_ref(),
observe,
)? {
return Ok(hash_group_fallback_strategy(
GroupedPlanFallbackReason::HavingBlocksGroupedOrder,
aggregate_family,
));
}
match order_strategy_projection {
GroupedOrderStrategyProjection::Canonical => {}
GroupedOrderStrategyProjection::TopK => {
return Ok(GroupedPlanStrategy::top_k_group_with_aggregate_family(
aggregate_family,
));
}
GroupedOrderStrategyProjection::HashFallback(reason) => {
return Ok(hash_group_fallback_strategy(reason, aggregate_family));
}
}
if grouped_access_path_proves_group_order(&grouped.group.group_fields, &plan.access, observe)? {
return Ok(GroupedPlanStrategy::ordered_group_with_aggregate_family(
aggregate_family,
));
}
Ok(hash_group_fallback_strategy(
GroupedPlanFallbackReason::GroupKeyOrderUnavailable,
aggregate_family,
))
}
fn grouped_aggregates_streaming_compatible<E>(
aggregates: &[GroupAggregateSpec],
observe: &mut impl FnMut(u64) -> Result<(), E>,
) -> Result<bool, E> {
for aggregate in aggregates {
observe(1)?;
if !aggregate.streaming_compatible() {
return Ok(false);
}
}
Ok(true)
}
const fn hash_group_fallback_strategy(
reason: GroupedPlanFallbackReason,
aggregate_family: GroupedPlanAggregateFamily,
) -> GroupedPlanStrategy {
GroupedPlanStrategy::hash_group_with_aggregate_family(reason, aggregate_family)
}
enum GroupedOrderStrategyProjection {
Canonical,
TopK,
HashFallback(GroupedPlanFallbackReason),
}
fn grouped_order_strategy_projection<E>(
order: Option<&OrderSpec>,
group_fields: &GroupFieldSet,
observe: &mut impl FnMut(u64) -> Result<(), E>,
) -> Result<GroupedOrderStrategyProjection, E> {
let Some(order) = order else {
return Ok(GroupedOrderStrategyProjection::Canonical);
};
for term in &order.fields {
observe(1)?;
if try_grouped_top_k_order_term_requires_heap(term.expr(), observe)? {
return grouped_top_k_strategy_projection(order, group_fields, observe);
}
}
grouped_canonical_order_strategy_projection(order, group_fields, observe)
}
fn grouped_canonical_order_strategy_projection<E>(
order: &OrderSpec,
group_fields: &GroupFieldSet,
observe: &mut impl FnMut(u64) -> Result<(), E>,
) -> Result<GroupedOrderStrategyProjection, E> {
observe(1)?;
if order.fields.len() < group_fields.len() {
return Ok(GroupedOrderStrategyProjection::HashFallback(
GroupedPlanFallbackReason::GroupKeyOrderPrefixMismatch,
));
}
let mut canonical_direction = None;
for (index, term) in order.fields.iter().take(group_fields.len()).enumerate() {
observe(1)?;
let direction = term.direction();
if canonical_direction.is_some_and(|expected| expected != direction) {
return Ok(GroupedOrderStrategyProjection::HashFallback(
GroupedPlanFallbackReason::GroupKeyOrderDirectionMismatch,
));
}
canonical_direction.get_or_insert(direction);
let Some(group_field) = group_fields.get(index) else {
return Ok(GroupedOrderStrategyProjection::HashFallback(
GroupedPlanFallbackReason::GroupKeyOrderPrefixMismatch,
));
};
match try_classify_grouped_order_term_for_field(term.expr(), group_field, observe)? {
GroupedOrderTermAdmissibility::Preserves(_) => {}
GroupedOrderTermAdmissibility::PrefixMismatch => {
return Ok(GroupedOrderStrategyProjection::HashFallback(
GroupedPlanFallbackReason::GroupKeyOrderPrefixMismatch,
));
}
GroupedOrderTermAdmissibility::UnsupportedExpression => {
return Ok(GroupedOrderStrategyProjection::HashFallback(
GroupedPlanFallbackReason::GroupKeyOrderExpressionNotAdmissible,
));
}
}
}
Ok(GroupedOrderStrategyProjection::Canonical)
}
fn grouped_top_k_strategy_projection<E>(
order: &OrderSpec,
group_fields: &GroupFieldSet,
observe: &mut impl FnMut(u64) -> Result<(), E>,
) -> Result<GroupedOrderStrategyProjection, E> {
for term in &order.fields {
observe(1)?;
match try_classify_grouped_top_k_order_term(term.expr(), group_fields, observe)? {
GroupedTopKOrderTermAdmissibility::Admissible => {}
GroupedTopKOrderTermAdmissibility::NonGroupFieldReference => {
return Ok(GroupedOrderStrategyProjection::HashFallback(
GroupedPlanFallbackReason::GroupKeyOrderPrefixMismatch,
));
}
GroupedTopKOrderTermAdmissibility::UnsupportedExpression => {
return Ok(GroupedOrderStrategyProjection::HashFallback(
GroupedPlanFallbackReason::GroupKeyOrderExpressionNotAdmissible,
));
}
}
}
Ok(GroupedOrderStrategyProjection::TopK)
}
fn grouped_access_path_proves_group_order<K, E>(
group_fields: &GroupFieldSet,
access: &AccessPlan<K>,
observe: &mut impl FnMut(u64) -> Result<(), E>,
) -> Result<bool, E> {
observe(1)?;
let executable = access.executable_contract();
let Some(path) = executable.as_path() else {
return Ok(false);
};
let Some(details) = path
.index_prefix_details()
.or_else(|| path.index_range_details())
else {
return Ok(false);
};
let prefix_len = details.slot_arity();
let mut cursor = 0usize;
for group_field in group_fields.iter() {
observe(1)?;
let comparison_steps = 1_u64.saturating_add(group_field.field().len() as u64);
while cursor < prefix_len && cursor < details.key_arity() {
observe(comparison_steps)?;
if details.key_field_at(cursor) == Some(group_field.field()) {
break;
}
cursor = cursor.saturating_add(1);
}
if cursor >= details.key_arity() {
return Ok(false);
}
observe(comparison_steps)?;
if details.key_field_at(cursor) != Some(group_field.field()) {
return Ok(false);
}
cursor = cursor.saturating_add(1);
}
Ok(true)
}
crate::retained::retained_copy!(GroupedPlanStrategy);