use nodedb_types::TraceId;
use nodedb_types::protocol::NativeResponse;
use nodedb_types::value::Value;
use crate::bridge::envelope::{Response, Status};
use crate::control::server::response_shape::compose::{ShapeOutcome, shape_response_materialized};
use crate::control::server::response_shape::schema::OutputSchema;
use crate::control::server::response_shape::types::describe_plan;
use crate::control::server::shared::ddl::sqlstate::error_code_to_sqlstate;
use crate::control::server::shared::session::expander_stage::{
ExpanderOutcome, route_in_tx_expander,
};
use crate::control::server::shared::session::staging_gate::{
InTxnRoute, StagedTagKind, StagingGateError, route_in_tx_write,
};
use crate::types::{DatabaseId, Lsn, VShardId};
use nodedb_physical::physical_task::PhysicalTask;
use super::sql_gateway::dispatch_task_via_gateway;
use super::streaming::SqlOutcome;
use super::{DispatchCtx, error_to_native, shape_error_to_native, to_native_columns_rows};
use crate::control::server::broadcast::broadcast_count_to_all_cores;
use crate::control::server::exchange::DistributedReadCapture;
use crate::control::server::exchange::resolve::{Resolved, resolve_and_materialize};
#[inline]
fn resp(r: NativeResponse) -> SqlOutcome {
SqlOutcome::Response(Box::new(r))
}
pub(super) async fn run_dispatch_loop(
ctx: &DispatchCtx<'_>,
seq: u64,
tasks: Vec<PhysicalTask>,
output_schema: Option<&OutputSchema>,
database_id: DatabaseId,
) -> SqlOutcome {
let mut all_columns: Option<Vec<String>> = None;
let mut all_rows: Vec<Vec<Value>> = Vec::new();
let mut warnings: Vec<String> = Vec::new();
let mut last_lsn = 0u64;
let mut total_affected = 0u64;
for task in tasks {
if task.tenant_id != ctx.tenant_id() {
return resp(NativeResponse::error(
seq,
"42501",
"tenant isolation violation",
));
}
let plan_for_staged_response = task.plan.clone();
let routed = match route_in_tx_expander(
ctx.state,
ctx.sessions,
ctx.peer_addr,
task,
|stage_task| async move {
dispatch_task(ctx, stage_task)
.await
.map(|(resp, _, _)| resp)
},
)
.await
{
Ok(ExpanderOutcome::Handled(route)) => Ok(route),
Ok(ExpanderOutcome::Passthrough(task)) => {
route_in_tx_write(
ctx.state,
ctx.sessions,
ctx.peer_addr,
*task,
|stage_task| async move {
dispatch_task(ctx, stage_task)
.await
.map(|(resp, _, _)| resp)
},
)
.await
}
Err(e) => Err(e),
};
let task = match routed {
Ok(InTxnRoute::Read(routed_task)) => *routed_task,
Ok(InTxnRoute::Buffered) => {
total_affected += 1;
continue;
}
Ok(InTxnRoute::Staged(outcome)) => {
if matches!(outcome.kind, StagedTagKind::RawPayload) && !outcome.payload.is_empty()
{
let plan_kind = describe_plan(&plan_for_staged_response);
match shape_response_materialized(
&outcome.payload,
&plan_for_staged_response,
plan_kind,
output_schema,
ctx.state,
database_id,
ctx.tenant_id(),
) {
Ok(ShapeOutcome::Rows(mut shaped)) => {
if let Some(notice) = shaped.notice.take() {
warnings.push(notice);
}
let (cols, rows) = to_native_columns_rows(&shaped);
if !cols.is_empty() && all_columns.is_none() {
all_columns = Some(cols);
}
all_rows.extend(rows);
}
Ok(ShapeOutcome::Passthrough) => {
total_affected += 1;
}
Err(e) => return resp(shape_error_to_native(seq, &e)),
}
} else {
total_affected += outcome.affected as u64;
}
continue;
}
Err(StagingGateError::Dispatch(e)) => return resp(error_to_native(seq, &e)),
Err(StagingGateError::Rejected { code }) => {
let (_, sqlstate, message) = match code {
Some(code) => error_code_to_sqlstate(&code),
None => ("ERROR", "XX000", "unknown data plane error".to_owned()),
};
return resp(NativeResponse::error(seq, sqlstate, message));
}
};
let plan_for_response = task.plan.clone();
let task_vshard = task.vshard_id;
let (task_resp, shard_watermarks, dist_reads) = match dispatch_task(ctx, task).await {
Ok(r) => r,
Err(e) => return resp(error_to_native(seq, &e)),
};
let records_read = task_resp.status == Status::Ok
|| task_resp.error_code.as_deref()
== Some(&crate::bridge::envelope::ErrorCode::NotFound);
if records_read
&& ctx.sessions.transaction_state(ctx.peer_addr)
== crate::control::server::shared::session::TransactionState::InBlock
{
let watermarks = if shard_watermarks.is_empty() {
vec![(task_vshard, task_resp.watermark_lsn)]
} else {
shard_watermarks
};
crate::control::server::shared::session::record_reads_for_response(
ctx.state,
ctx.sessions,
ctx.peer_addr,
ctx.tenant_id(),
crate::control::server::shared::session::ResponseReads {
plan: &plan_for_response,
watermarks: &watermarks,
read_version_lsn: task_resp.read_version_lsn,
found: task_resp.status == Status::Ok,
distributed_reads: &dist_reads,
read_lsn_vshard: task_vshard,
},
)
.await;
}
if task_resp.status == Status::Error {
let msg = if task_resp.payload.is_empty() {
task_resp
.error_code
.as_ref()
.map(|c| format!("{c:?}"))
.unwrap_or_else(|| "unknown error".into())
} else {
String::from_utf8_lossy(&task_resp.payload).into_owned()
};
return resp(NativeResponse::error(seq, "XX000", msg));
}
last_lsn = task_resp.watermark_lsn.as_u64();
if task_resp.payload.is_empty() {
total_affected += 1;
} else {
let plan_kind = describe_plan(&plan_for_response);
match shape_response_materialized(
&task_resp.payload,
&plan_for_response,
plan_kind,
output_schema,
ctx.state,
database_id,
ctx.tenant_id(),
) {
Ok(ShapeOutcome::Rows(mut shaped)) => {
if let Some(notice) = shaped.notice.take() {
warnings.push(notice);
}
let (cols, rows) = to_native_columns_rows(&shaped);
if !cols.is_empty() && all_columns.is_none() {
all_columns = Some(cols);
}
all_rows.extend(rows);
}
Ok(ShapeOutcome::Passthrough) => {
total_affected += 1;
}
Err(e) => return resp(shape_error_to_native(seq, &e)),
}
}
}
if all_rows.is_empty() {
let mut r = NativeResponse::ok(seq);
r.rows_affected = Some(total_affected);
r.watermark_lsn = last_lsn;
r.warnings = warnings;
resp(r)
} else {
resp(NativeResponse {
seq,
status: nodedb_types::protocol::ResponseStatus::Ok,
columns: all_columns,
rows: Some(all_rows),
rows_affected: Some(total_affected),
watermark_lsn: last_lsn,
error: None,
auth: None,
warnings,
})
}
}
async fn dispatch_task(
ctx: &DispatchCtx<'_>,
mut task: PhysicalTask,
) -> crate::Result<(Response, Vec<(VShardId, Lsn)>, Vec<DistributedReadCapture>)> {
if let crate::bridge::envelope::PhysicalPlan::Document(
nodedb_physical::physical_plan::DocumentOp::InsertSelect {
target_collection,
source_collection,
source_filters,
source_limit,
},
) = &task.plan
{
let resp = crate::control::insert_select::run_insert_select(
ctx.state,
task.tenant_id,
task.database_id,
target_collection,
source_collection,
source_filters,
*source_limit,
)
.await?;
return Ok((resp, Vec::new(), Vec::new()));
}
if let crate::bridge::envelope::PhysicalPlan::Document(
nodedb_physical::physical_plan::DocumentOp::Merge {
target_collection,
source_collection,
source_alias,
target_join_col,
source_join_col,
clauses,
returning: _,
resolve_only: false,
resolved_inserts: None,
source_rows: _,
},
) = &task.plan
{
let resp = crate::control::merge_orchestrator::run_merge(
ctx.state,
crate::control::merge_orchestrator::MergeArgs {
tenant_id: task.tenant_id,
database_id: task.database_id,
target_collection,
source_collection,
source_alias,
target_join_col,
source_join_col,
clauses,
},
)
.await?;
return Ok((resp, Vec::new(), Vec::new()));
}
if let crate::bridge::envelope::PhysicalPlan::Document(
nodedb_physical::physical_plan::DocumentOp::UpdateFromJoin {
target_collection,
source_collection,
source_alias,
target_join_col,
source_join_col,
updates,
target_filters,
returning,
resolve_only: false,
source_rows: None,
},
) = &task.plan
{
let resp = crate::control::update_from_join_orchestrator::run_update_from_join(
ctx.state,
crate::control::update_from_join_orchestrator::UpdateFromJoinArgs {
tenant_id: task.tenant_id,
database_id: task.database_id,
target_collection,
source_collection,
source_alias,
target_join_col,
source_join_col,
updates,
target_filters,
returning: returning.as_ref(),
},
)
.await?;
return Ok((resp, Vec::new(), Vec::new()));
}
if matches!(
task.plan,
crate::bridge::envelope::PhysicalPlan::Array(
nodedb_physical::physical_plan::ArrayOp::DropArray { .. }
)
) {
let resp = broadcast_count_to_all_cores(
ctx.state,
task.tenant_id,
task.database_id,
task.plan,
TraceId::ZERO,
"dropped",
)
.await?;
return Ok((resp, Vec::new(), Vec::new()));
}
match resolve_and_materialize(
ctx.state,
ctx.identity,
task.database_id,
task.tenant_id,
task.plan,
TraceId::ZERO,
task.txn_id,
)
.await?
{
Resolved::Gathered(resp, shard_watermarks, dist_reads) => {
return Ok((resp, shard_watermarks, dist_reads));
}
Resolved::Plan(resolved_plan) => {
task.plan = resolved_plan;
}
Resolved::Stream(s) => {
let resp = crate::control::server::exchange::gather::stream_to_response(s).await?;
return Ok((resp, Vec::new(), Vec::new()));
}
}
let resp = dispatch_task_via_gateway(ctx, task).await?;
Ok((resp, Vec::new(), Vec::new()))
}