use nodedb_physical::physical_task::PostSetOp;
use nodedb_types::protocol::NativeResponse;
use crate::control::gateway::core::QueryContext;
use crate::control::server::exchange::gather::gather_all_cores_stream;
use crate::control::server::exchange::streamable::streamable_gather_child;
use crate::control::server::response_shape::schema::OutputSchema;
use crate::control::server::result_stream::ResultStream;
use crate::control::server::shared::session::TransactionState;
use super::DispatchCtx;
pub(crate) enum SqlOutcome {
Response(Box<NativeResponse>),
Stream(SqlStream),
}
impl SqlOutcome {
pub(crate) fn into_response(self) -> NativeResponse {
match self {
SqlOutcome::Response(r) => *r,
SqlOutcome::Stream(s) => NativeResponse::error(
s.seq,
"XX000",
"internal error: SQL stream produced on a non-streaming path",
),
}
}
}
pub(crate) struct SqlStream {
pub seq: u64,
pub limit: usize,
pub stream: ResultStream,
pub projection: Option<OutputSchema>,
}
pub(crate) async fn try_open_sql_stream(
ctx: &DispatchCtx<'_>,
seq: u64,
tasks: &[nodedb_physical::physical_task::PhysicalTask],
database_id: crate::types::DatabaseId,
output_schema: Option<&OutputSchema>,
) -> crate::Result<Option<SqlStream>> {
let [task] = tasks else {
return Ok(None);
};
if task.post_set_op != PostSetOp::None
|| ctx.sessions.transaction_state(ctx.peer_addr) == TransactionState::InBlock
{
return Ok(None);
}
let Some((child_plan, limit)) = streamable_gather_child(&task.plan) else {
return Ok(None);
};
let stream = if let Some(gw) = ctx.state.gateway.get() {
let gw_ctx = QueryContext {
tenant_id: task.tenant_id,
trace_id: crate::types::TraceId::ZERO,
database_id,
txn_id: task.txn_id,
};
gw.execute_stream(&gw_ctx, child_plan).await?
} else {
gather_all_cores_stream(
ctx.state,
task.tenant_id,
task.database_id,
child_plan,
crate::types::TraceId::ZERO,
task.txn_id,
)?
};
Ok(Some(SqlStream {
seq,
limit,
stream,
projection: output_schema.cloned(),
}))
}