use nodedb_types::protocol::NativeResponse;
use crate::control::planner::calvin::{
dispatch_dependent_edge_recon, plan_needs_implicit_edge_recon,
};
use crate::control::server::shared::session::TransactionState;
use nodedb_physical::physical_task::PhysicalTask;
use super::{DispatchCtx, SqlOutcome, error_to_native};
pub(super) async fn try_edge_recon_dispatch(
ctx: &DispatchCtx<'_>,
seq: u64,
tasks: Vec<PhysicalTask>,
) -> EdgeReconResult {
if ctx.sessions.transaction_state(ctx.peer_addr) == TransactionState::InBlock {
return EdgeReconResult::NotFired(tasks);
}
if ctx.state.calvin_completion_registry.get().is_none() {
return EdgeReconResult::NotFired(tasks);
}
let (_coll, database_id) =
match plan_needs_implicit_edge_recon(ctx.state, &tasks, ctx.tenant_id()) {
Err(e) => return EdgeReconResult::Outcome(resp(error_to_native(seq, &e))),
Ok(None) => return EdgeReconResult::NotFired(tasks),
Ok(Some(pair)) => pair,
};
let returning_plan = tasks
.iter()
.find(|t| {
matches!(
crate::control::server::response_shape::types::describe_plan(&t.plan),
crate::control::server::response_shape::types::PlanKind::ReturningRows
)
})
.map(|t| t.plan.clone());
let outcome =
dispatch_dependent_edge_recon(ctx.state, tasks, ctx.tenant_id(), database_id, false).await;
EdgeReconResult::Outcome(match outcome {
Ok(recon) => {
resp(super::conversion::calvin_native_response(
seq,
recon.apply_result,
returning_plan.as_ref(),
recon.tasks_dispatched,
ctx.state,
database_id,
ctx.tenant_id(),
))
}
Err(e) => resp(error_to_native(seq, &e)),
})
}
pub(super) enum EdgeReconResult {
NotFired(Vec<PhysicalTask>),
Outcome(SqlOutcome),
}
fn resp(r: NativeResponse) -> SqlOutcome {
SqlOutcome::Response(Box::new(r))
}