use crate::db::executor::{SharedPreparedExecutionPlan, StructuralGroupedProjectionResult};
use crate::db::registry::StoreHandle;
use crate::{
db::{
cursor::ValidatedGroupedCursor,
executor::{
EntityAuthority, ExecutionPreparation, PreparedGroupedRuntimeResidents,
PreparedLoadPlan, RetainedSlotLayout,
aggregate::runtime::{
build_grouped_stream_with_runtime, execute_group_fold_stage,
finalize_grouped_output, try_execute_grouped_count_metadata,
},
budget::ExecutionConstructionBudget,
budget::{
charge_current_execution_budget, charge_runtime_grouped_rows,
prepared_read_execution_context, runtime_value_work, with_read_execution_budget,
},
pipeline::contracts::{ExecutionRuntimeAdapter, GroupedCursorPage, GroupedRouteStage},
pipeline::grouped_runtime::resolve_grouped_route_for_plan,
pipeline::runtime::{
GroupedStreamStage, StructuralGroupedRowRuntime,
compile_grouped_row_slot_layout_from_inputs,
},
stream::access::TraversalRuntime,
},
schema::cardinality_generation::CardinalityAcceptedRootIdentity,
},
error::InternalError,
metrics::EntityMetricsSpan,
traits::CanisterKind,
};
use icydb_diagnostic_code::{DiagnosticExecutionBudgetResource, DiagnosticExecutionLane};
pub(in crate::db) fn execute_shared_grouped_plan_for_canister<C>(
db: &crate::db::Db<C>,
plan: SharedPreparedExecutionPlan,
cursor: ValidatedGroupedCursor,
execution_lane: DiagnosticExecutionLane,
) -> Result<StructuralGroupedProjectionResult, InternalError>
where
C: CanisterKind,
{
let context = prepared_read_execution_context(&plan, execution_lane);
with_read_execution_budget(db.request_execution_scope(), context, || {
execute_shared_grouped_plan_for_canister_inner(db, plan, cursor)
})
}
fn execute_shared_grouped_plan_for_canister_inner<C>(
db: &crate::db::Db<C>,
plan: SharedPreparedExecutionPlan,
cursor: ValidatedGroupedCursor,
) -> Result<StructuralGroupedProjectionResult, InternalError>
where
C: CanisterKind,
{
let entity_path = plan.authority_ref().entity_path_handle();
let _metrics_span = EntityMetricsSpan::new(entity_path.as_ref());
charge_grouped_cursor_input(&cursor)?;
let value_catalog = plan
.authority_ref()
.accepted_schema_info()
.map(crate::db::schema::SchemaInfo::value_catalog_handle)
.cloned()
.ok_or_else(InternalError::query_executor_invariant)?;
let prepared =
prepare_grouped_route_runtime_for_load_plan(db, plan.into_prepared_load_plan(), cursor)?;
let page = execute_prepared_grouped_route_runtime(prepared)?;
charge_grouped_page_result(&page)?;
Ok(StructuralGroupedProjectionResult::from_page(
page,
value_catalog,
))
}
fn charge_grouped_page_result(page: &GroupedCursorPage) -> Result<(), InternalError> {
charge_runtime_grouped_rows(&page.rows)?;
if let Some(cursor) = page.next_cursor.as_ref() {
let encoded = cursor
.encode()
.map_err(|_| InternalError::query_executor_invariant())?;
charge_current_execution_budget(DiagnosticExecutionBudgetResource::CursorSteps, 1)?;
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::TemporaryBytes,
u64::try_from(encoded.len()).unwrap_or(u64::MAX),
)?;
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::ResultBytes,
u64::try_from(encoded.len().saturating_mul(2)).unwrap_or(u64::MAX),
)?;
}
Ok(())
}
fn charge_grouped_cursor_input(cursor: &ValidatedGroupedCursor) -> Result<(), InternalError> {
let Some(group_key) = cursor.last_group_key() else {
return Ok(());
};
let (bytes, nested_steps) = group_key.iter().fold((0_u64, 0_u64), |total, value| {
let value_work = runtime_value_work(value);
(
total.0.saturating_add(value_work.0),
total.1.saturating_add(value_work.1),
)
});
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::CursorSteps,
u64::try_from(group_key.len()).unwrap_or(u64::MAX),
)?;
charge_current_execution_budget(
DiagnosticExecutionBudgetResource::NestedValueSteps,
nested_steps,
)?;
charge_current_execution_budget(DiagnosticExecutionBudgetResource::TemporaryBytes, bytes)
}
struct GroupedPathRuntimeContext {
traversal_runtime: TraversalRuntime,
row_store: StoreHandle,
authority: EntityAuthority,
}
pub(in crate::db::executor) struct PreparedGroupedRouteRuntime {
route: GroupedRouteStage,
runtime: GroupedPathRuntimeContext,
execution_preparation: ExecutionPreparation,
grouped_slot_layout: RetainedSlotLayout,
}
impl GroupedPathRuntimeContext {
fn from_store(store: StoreHandle, authority: EntityAuthority) -> Result<Self, InternalError> {
let entity_tag = authority.entity_tag();
let accepted_schema = authority.accepted_schema_authority()?;
let accepted_root = CardinalityAcceptedRootIdentity::new(
accepted_schema.revision(),
accepted_schema.fingerprint(),
)?;
Ok(Self {
traversal_runtime: TraversalRuntime::new(
store,
entity_tag,
authority
.accepted_runtime_root_identity()
.database_incarnation(),
accepted_root,
),
row_store: store,
authority,
})
}
fn build_grouped_stream(
&self,
route: &GroupedRouteStage,
execution_preparation: ExecutionPreparation,
grouped_slot_layout: RetainedSlotLayout,
) -> Result<GroupedStreamStage, InternalError> {
let runtime = ExecutionRuntimeAdapter::from_stream_runtime(self.traversal_runtime);
let single_grouped_path = if execution_preparation
.effective_runtime_filter_program()
.is_none()
&& matches!(route.grouped_aggregate_execution_specs(), [aggregate] if aggregate.admits_count_rows_dedicated_fold())
{
route
.group_fields()
.as_path_aware()
.and_then(|fields| match fields {
[field] => field.as_scalar_path(),
_ => None,
})
} else {
None
};
build_grouped_stream_with_runtime(
route,
&runtime,
execution_preparation,
StructuralGroupedRowRuntime::new(
self.row_store,
self.authority.row_layout()?,
grouped_slot_layout,
single_grouped_path,
),
)
}
}
impl PreparedGroupedRouteRuntime {
fn new(
route: GroupedRouteStage,
runtime: GroupedPathRuntimeContext,
prepared_residents: Option<PreparedGroupedRuntimeResidents>,
) -> Result<Self, InternalError> {
let residents = if let Some(residents) = prepared_residents {
residents
} else {
let execution_preparation = ExecutionPreparation::from_runtime_plan(
route.plan(),
route.plan().slot_map().map(<[usize]>::to_vec),
&ExecutionConstructionBudget,
)?;
let grouped_slot_layout = compile_grouped_row_slot_layout_from_inputs(
runtime.authority.row_layout()?,
route.group_fields(),
route.grouped_aggregate_execution_specs(),
route.grouped_distinct_execution_strategy(),
execution_preparation.effective_runtime_filter_program(),
);
PreparedGroupedRuntimeResidents::new(execution_preparation, grouped_slot_layout)
};
let (execution_preparation, grouped_slot_layout) = residents.into_parts();
Ok(Self {
route,
runtime,
execution_preparation,
grouped_slot_layout,
})
}
}
pub(in crate::db::executor) fn prepare_grouped_route_runtime_for_load_plan<C>(
db: &crate::db::Db<C>,
plan: PreparedLoadPlan,
cursor: ValidatedGroupedCursor,
) -> Result<PreparedGroupedRouteRuntime, InternalError>
where
C: CanisterKind,
{
let authority = plan.authority();
let prepared_residents = plan.cloned_grouped_runtime_residents()?;
let route = resolve_grouped_route_for_plan(plan, cursor)?;
let store = db.recovered_store(authority.store_path())?;
PreparedGroupedRouteRuntime::new(
route,
GroupedPathRuntimeContext::from_store(store, authority)?,
prepared_residents,
)
}
fn execute_grouped_route_path(
runtime: &GroupedPathRuntimeContext,
route: GroupedRouteStage,
execution_preparation: ExecutionPreparation,
grouped_slot_layout: RetainedSlotLayout,
) -> Result<GroupedCursorPage, InternalError> {
if let Some(page) = try_execute_grouped_count_metadata(
runtime.row_store,
runtime.authority.entity_tag(),
&route,
)? {
return Ok(page);
}
let stream =
runtime.build_grouped_stream(&route, execution_preparation, grouped_slot_layout)?;
let folded = execute_group_fold_stage(&route, stream)?;
Ok(finalize_grouped_output(folded))
}
pub(in crate::db::executor) fn execute_prepared_grouped_route_runtime(
prepared: PreparedGroupedRouteRuntime,
) -> Result<GroupedCursorPage, InternalError> {
let PreparedGroupedRouteRuntime {
route,
runtime,
execution_preparation,
grouped_slot_layout,
} = prepared;
execute_grouped_route_path(&runtime, route, execution_preparation, grouped_slot_layout)
}