use super::sql_write_candidate_bounds;
use crate::db::query::preparation::PreparationWork;
use crate::{
db::{
DbSession, MissingRowPolicy, QueryError,
data::{AcceptedMutationIntentPatch, DecodedDataStoreKey},
query::intent::StructuralQuery,
schema::SchemaInfo,
session::{
AcceptedSchemaCatalogContext,
sql::{
SqlDeleteExposurePolicy, SqlDeletePolicyContext, SqlPublicBoundedDeletePlan,
SqlPublicPrimaryKeyDeletePlan, SqlStatementDispatch, SqlStatementResult,
SqlValidatedDeletePlan, classify_sql_delete_statement_policy,
execute::write_returning::{
projection_labels_from_accepted_write_descriptor,
sql_returning_statement_projection, validate_sql_materialized_returning_bounds,
},
write_policy::SqlWriteExecutionBounds,
},
},
sql::{
lowering::bind_sql_delete_statement_structural_with_schema,
parser::{SqlDeleteStatement, SqlReturningProjection},
},
},
traits::CanisterKind,
};
impl<C: CanisterKind> DbSession<C> {
pub(in crate::db::session::sql::execute) fn execute_sql_delete_statement(
&self,
query: &StructuralQuery,
returning: Option<&SqlReturningProjection>,
catalog: Option<&AcceptedSchemaCatalogContext>,
) -> Result<SqlStatementResult, QueryError> {
self.execute_sql_delete_statement_with_execution_bounds(query, returning, catalog, None)
}
fn execute_sql_delete_statement_with_execution_bounds(
&self,
query: &StructuralQuery,
returning: Option<&SqlReturningProjection>,
catalog: Option<&AcceptedSchemaCatalogContext>,
execution_bounds: Option<SqlWriteExecutionBounds>,
) -> Result<SqlStatementResult, QueryError> {
let entity_name = catalog
.map(|catalog| {
catalog
.snapshot()
.persisted_snapshot()
.entity_name()
.to_string()
})
.ok_or_else(QueryError::invariant)?;
self.with_checked_accepted_write_descriptor_for_returning(
catalog,
Some(entity_name.as_str()),
returning,
|catalog, descriptor| {
let (authority, _) = Self::accepted_sql_write_authority_schema_info(catalog);
let primary_names = descriptor.primary_key_names();
let selector = query
.clone()
.into_load_selection()
.select_fields(primary_names.iter().copied());
let bounds = sql_write_candidate_bounds(execution_bounds);
let entity_tag = catalog.identity().entity_tag();
let rows = self.collect_bounded_sql_write_mutation_batch_from_structural_query(
catalog.snapshot(),
authority,
&selector,
bounds,
None,
|row| {
let key =
DecodedDataStoreKey::try_from_structural_key_values(entity_tag, row)
.map_err(QueryError::execute)?;
Ok((key, AcceptedMutationIntentPatch::new()))
},
)?;
let keys = rows.into_rows().into_iter().map(|(key, _)| key).collect();
let columns = projection_labels_from_accepted_write_descriptor(&descriptor);
let rows = self
.execute_accepted_structural_delete_batch(
catalog,
returning.is_some(),
keys,
|rows| {
let Some(returning) = returning else {
return Ok(());
};
validate_sql_materialized_returning_bounds(
entity_name.as_str(),
columns.as_slice(),
rows,
u32::try_from(rows.len()).unwrap_or(u32::MAX),
returning,
catalog.enum_catalog(),
execution_bounds.map(|bounds| bounds.returning),
)
},
)
.map_err(QueryError::execute)?;
let row_count = u32::try_from(rows.len()).unwrap_or(u32::MAX);
match returning {
None => Ok(SqlStatementResult::Count { row_count }),
Some(returning) => sql_returning_statement_projection(
catalog.enum_catalog(),
columns,
rows,
row_count,
returning,
),
}
},
)
}
fn sql_delete_query_from_statement(
&self,
schema_info: &SchemaInfo,
statement: &SqlDeleteStatement,
) -> Result<StructuralQuery, QueryError> {
PreparationWork::run(
self.db.request_execution_scope(),
icydb_diagnostic_code::DiagnosticExecutionLane::Mutation,
|work| {
bind_sql_delete_statement_structural_with_schema(
statement.clone(),
MissingRowPolicy::Ignore,
schema_info,
work,
)
.map_err(QueryError::from_sql_lowering_error)
},
)
}
fn schema_derived_sql_delete_plan(
&self,
dispatch: &SqlStatementDispatch<'_>,
policy: SqlDeleteExposurePolicy,
) -> Result<SqlValidatedDeletePlan, QueryError> {
let entity_name = dispatch.entity_name();
self.with_checked_accepted_write_descriptor_for_returning(
None,
entity_name,
None,
|_catalog, descriptor| {
let context =
SqlDeletePolicyContext::public_generated(descriptor.primary_key_names());
let result =
classify_sql_delete_statement_policy(dispatch.statement(), policy, context);
result.map_err(|_| QueryError::unsupported_query())
},
)
}
fn execute_validated_sql_delete_statement(
&self,
statement: &SqlDeleteStatement,
execution_bounds: SqlWriteExecutionBounds,
) -> Result<SqlStatementResult, QueryError> {
self.with_checked_accepted_write_descriptor_for_returning(
None,
Some(statement.entity.as_str()),
statement.returning.as_ref(),
|catalog, _descriptor| {
let (_authority, schema_info) =
Self::accepted_sql_write_authority_schema_info(catalog);
let query = self.sql_delete_query_from_statement(&schema_info, statement)?;
self.execute_sql_delete_statement_with_execution_bounds(
&query,
statement.returning.as_ref(),
Some(catalog),
Some(execution_bounds),
)
},
)
}
#[doc(hidden)]
pub(in crate::db) fn execute_validated_sql_public_primary_key_delete(
&self,
plan: &SqlPublicPrimaryKeyDeletePlan,
) -> Result<SqlStatementResult, QueryError> {
self.execute_validated_sql_delete_statement(plan.statement(), plan.execution_bounds())
}
#[doc(hidden)]
pub(in crate::db) fn execute_validated_sql_public_bounded_delete(
&self,
plan: &SqlPublicBoundedDeletePlan,
) -> Result<SqlStatementResult, QueryError> {
self.execute_validated_sql_delete_statement(plan.statement(), plan.execution_bounds())
}
#[doc(hidden)]
pub fn execute_sql_public_primary_key_delete(
&self,
dispatch: &SqlStatementDispatch<'_>,
) -> Result<SqlStatementResult, QueryError> {
let plan = self.schema_derived_sql_delete_plan(
dispatch,
SqlDeleteExposurePolicy::PublicPrimaryKeyOnly,
)?;
let SqlValidatedDeletePlan::PublicPrimaryKeyOnly(plan) = plan else {
return Err(QueryError::invariant());
};
self.execute_validated_sql_public_primary_key_delete(&plan)
}
#[doc(hidden)]
pub fn execute_sql_public_bounded_delete(
&self,
dispatch: &SqlStatementDispatch<'_>,
) -> Result<SqlStatementResult, QueryError> {
let plan = self.schema_derived_sql_delete_plan(
dispatch,
SqlDeleteExposurePolicy::PublicBoundedDeterministic,
)?;
let SqlValidatedDeletePlan::PublicBoundedDeterministic(plan) = plan else {
return Err(QueryError::invariant());
};
self.execute_validated_sql_public_bounded_delete(&plan)
}
}