use nodedb_sql::types::{EngineType, Filter, SqlValue};
use nodedb_types::Surrogate;
use crate::bridge::envelope::PhysicalPlan;
use crate::types::{TenantId, VShardId};
use nodedb_physical::physical_plan::*;
use crate::control::planner::sql_plan_convert::convert::ConvertContext;
use crate::control::planner::sql_plan_convert::filter::serialize_filters;
use crate::control::planner::sql_plan_convert::value::{sql_value_to_bytes, sql_value_to_string};
use nodedb_physical::physical_task::{PhysicalTask, PostSetOp};
use super::shared::{
delete_effective_filter, document_collection_is_edge_bearing, pk_effective_filter,
};
pub(in crate::control::planner::sql_plan_convert) fn convert_delete(
collection: &str,
engine: &EngineType,
filters: &[Filter],
target_keys: &[SqlValue],
tenant_id: TenantId,
ctx: &ConvertContext,
) -> crate::Result<Vec<PhysicalTask>> {
let coll_qualified = crate::control::planner::sql_plan_convert::convert::db_qualified(
ctx.database_id,
collection,
);
let collection = coll_qualified.as_str();
let vshard = VShardId::from_collection_in_database(ctx.database_id, collection);
if matches!(engine, EngineType::KeyValue) && !target_keys.is_empty() {
let keys: Vec<Vec<u8>> = target_keys.iter().map(sql_value_to_bytes).collect();
return Ok(vec![PhysicalTask {
tenant_id,
vshard_id: vshard,
database_id: ctx.database_id,
plan: PhysicalPlan::Kv(KvOp::Delete {
collection: collection.into(),
keys,
}),
post_set_op: PostSetOp::None,
txn_id: None,
}]);
}
if matches!(engine, EngineType::Columnar | EngineType::Spatial) {
let filter_bytes = serialize_filters(filters)?;
let effective_filter = pk_effective_filter(filter_bytes, target_keys)?;
return Ok(vec![PhysicalTask {
tenant_id,
vshard_id: vshard,
database_id: ctx.database_id,
plan: PhysicalPlan::Columnar(ColumnarOp::Delete {
collection: collection.into(),
filters: effective_filter,
}),
post_set_op: PostSetOp::None,
txn_id: None,
}]);
}
let is_crdt = super::super::crdt_gate::document_collection_is_crdt(ctx, collection)?;
if is_crdt && target_keys.is_empty() {
return Err(crate::Error::BadRequest {
detail: format!(
"predicate (non-primary-key) DELETE on CRDT collection '{collection}' is not \
supported; target rows by primary key"
),
});
}
if !is_crdt && !target_keys.is_empty() && document_collection_is_edge_bearing(ctx, collection)?
{
let effective_filter = delete_effective_filter(filters, target_keys)?;
return Ok(vec![PhysicalTask {
tenant_id,
vshard_id: vshard,
database_id: ctx.database_id,
plan: PhysicalPlan::Document(DocumentOp::BulkDelete {
collection: collection.into(),
filters: effective_filter,
returning: None,
ollp_predicted_surrogates: None,
ollp_predicted_edges: None,
}),
post_set_op: PostSetOp::None,
txn_id: None,
}]);
}
if !target_keys.is_empty() {
let mut tasks = Vec::new();
for key in target_keys {
let pk_string = sql_value_to_string(key);
let pk_bytes = pk_string.clone().into_bytes();
let surrogate = match ctx.surrogate_assigner.as_ref() {
Some(a) => match a.lookup(ctx.database_id, ctx.tenant_id, collection, &pk_bytes)? {
Some(s) => s,
None => continue,
},
None => Surrogate::ZERO,
};
let plan = if is_crdt {
PhysicalPlan::Crdt(CrdtOp::DocDelete {
collection: collection.into(),
document_id: pk_string,
surrogate,
returning: None,
})
} else {
PhysicalPlan::Document(DocumentOp::PointDelete {
collection: collection.into(),
document_id: pk_string,
surrogate,
pk_bytes,
returning: None,
})
};
tasks.push(PhysicalTask {
tenant_id,
vshard_id: vshard,
database_id: ctx.database_id,
plan,
post_set_op: PostSetOp::None,
txn_id: None,
});
}
Ok(tasks)
} else {
let filter_bytes = serialize_filters(filters)?;
Ok(vec![PhysicalTask {
tenant_id,
vshard_id: vshard,
database_id: ctx.database_id,
plan: PhysicalPlan::Document(DocumentOp::BulkDelete {
collection: collection.into(),
filters: filter_bytes,
returning: None,
ollp_predicted_surrogates: None,
ollp_predicted_edges: None,
}),
post_set_op: PostSetOp::None,
txn_id: None,
}])
}
}