use nodedb_physical::physical_plan::{ExchangeMode, ExchangeOp, PhysicalPlan, QueryOp};
use crate::bridge::envelope::Response;
use crate::control::security::identity::AuthenticatedIdentity;
use crate::control::state::SharedState;
use crate::types::{DatabaseId, Lsn, TenantId, TraceId, TxnId, VShardId};
use crate::control::server::exchange::gather::{
GatherOutcome, finalize_aggregate, gather_all_cores_stream, gather_all_vshards,
outcome_to_response,
};
use crate::control::server::result_stream::ResultStream;
use super::capture::DistributedReadCapture;
use super::join_input::{gather_join_build_side, resolve_join_input};
use super::materialize::materialize_providers;
use crate::control::server::exchange::full_scan::full_scan_plan_for_collection;
pub enum Resolved {
Gathered(Response, Vec<(VShardId, Lsn)>, Vec<DistributedReadCapture>),
Plan(PhysicalPlan),
Stream(ResultStream),
}
pub async fn resolve_and_materialize(
state: &SharedState,
identity: &AuthenticatedIdentity,
database_id: DatabaseId,
tenant_id: TenantId,
plan: PhysicalPlan,
trace_id: TraceId,
txn_id: Option<TxnId>,
) -> crate::Result<Resolved> {
let plan = materialize_providers(state, identity, plan).await?;
let mut captures = Vec::new();
resolve_exchange(
state,
database_id,
tenant_id,
plan,
trace_id,
txn_id,
&mut captures,
)
.await
}
pub async fn resolve_exchange_in_plan(
state: &SharedState,
database_id: DatabaseId,
tenant_id: TenantId,
plan: PhysicalPlan,
trace_id: TraceId,
txn_id: Option<TxnId>,
) -> crate::Result<Resolved> {
let mut captures = Vec::new();
resolve_exchange(
state,
database_id,
tenant_id,
plan,
trace_id,
txn_id,
&mut captures,
)
.await
}
async fn resolve_exchange(
state: &SharedState,
database_id: DatabaseId,
tenant_id: TenantId,
plan: PhysicalPlan,
trace_id: TraceId,
txn_id: Option<TxnId>,
captures: &mut Vec<DistributedReadCapture>,
) -> crate::Result<Resolved> {
match plan {
PhysicalPlan::Query(QueryOp::Exchange(ExchangeOp {
child,
mode: ExchangeMode::Gather { as_aggregate },
})) => {
let child = match Box::pin(resolve_exchange(
state,
database_id,
tenant_id,
*child,
trace_id,
txn_id,
captures,
))
.await?
{
Resolved::Plan(p) => p,
Resolved::Gathered(resp, wms, caps) => {
return Ok(Resolved::Gathered(resp, wms, caps));
}
Resolved::Stream(s) => return Ok(Resolved::Stream(s)),
};
if !as_aggregate && txn_id.is_none() && child.is_streamable_unordered_scan() {
let stream = if let Some(gw) = state.gateway.get() {
let ctx = crate::control::gateway::core::QueryContext {
tenant_id,
trace_id,
database_id,
txn_id: None,
};
gw.execute_stream(&ctx, child).await?
} else {
gather_all_cores_stream(state, tenant_id, database_id, child, trace_id, txn_id)?
};
return Ok(Resolved::Stream(stream));
}
let probe_collection: Option<String> = if txn_id.is_some() {
match &child {
PhysicalPlan::Query(QueryOp::HashJoin {
left_collection, ..
}) => Some(left_collection.clone()),
other => other.collection().map(str::to_owned),
}
} else {
None
};
let outcome: GatherOutcome =
gather_all_vshards(state, tenant_id, database_id, child, trace_id, txn_id).await?;
if let Some(coll) = probe_collection
&& let Some(scan_plan) =
full_scan_plan_for_collection(state, database_id, tenant_id, &coll)?
{
captures.push(DistributedReadCapture {
scan_plan,
read_version_lsn: outcome.read_version_lsn,
});
}
let payload = if as_aggregate {
finalize_aggregate(&outcome.merged_array)
} else {
outcome.merged_array
};
Ok(Resolved::Gathered(
outcome_to_response(payload, outcome.watermark_lsn, outcome.read_version_lsn),
outcome.shard_watermarks,
std::mem::take(captures),
))
}
PhysicalPlan::Query(QueryOp::Exchange(ExchangeOp {
child,
mode: ExchangeMode::Broadcast,
})) => {
let outcome =
gather_all_vshards(state, tenant_id, database_id, *child, trace_id, txn_id).await?;
Ok(Resolved::Gathered(
outcome_to_response(
outcome.merged_array,
outcome.watermark_lsn,
outcome.read_version_lsn,
),
outcome.shard_watermarks,
std::mem::take(captures),
))
}
PhysicalPlan::Query(QueryOp::Exchange(ExchangeOp {
child,
mode: ExchangeMode::Shuffle { keys, num_parts },
})) => {
super::shuffle::resolve_shuffle_join(
state,
database_id,
tenant_id,
*child,
keys,
num_parts,
trace_id,
)
.await
}
PhysicalPlan::Query(QueryOp::Exchange(ExchangeOp {
child,
mode: ExchangeMode::ShuffleAggregate { keys, num_parts },
})) => {
super::shuffle_aggregate::resolve_shuffle_aggregate(
state,
database_id,
tenant_id,
*child,
keys,
num_parts,
trace_id,
)
.await
}
PhysicalPlan::Query(QueryOp::HashJoin {
left_collection,
right_collection,
left_alias,
right_alias,
on,
join_type,
limit,
post_group_by,
post_aggregates,
projection,
computed_projection,
join_filters,
post_filters,
left_input,
right_input,
left_bitmap,
right_bitmap,
}) => {
let left_input = resolve_join_input(
state,
database_id,
tenant_id,
left_input,
trace_id,
txn_id,
captures,
)
.await?;
let mut right_input = resolve_join_input(
state,
database_id,
tenant_id,
right_input,
trace_id,
txn_id,
captures,
)
.await?;
if state.gateway.get().is_some()
&& right_input.is_none()
&& !right_collection.is_empty()
{
right_input = gather_join_build_side(
state,
database_id,
tenant_id,
&right_collection,
trace_id,
txn_id,
captures,
)
.await?;
}
Ok(Resolved::Plan(PhysicalPlan::Query(QueryOp::HashJoin {
left_collection,
right_collection,
left_alias,
right_alias,
on,
join_type,
limit,
post_group_by,
post_aggregates,
projection,
computed_projection,
join_filters,
post_filters,
left_input,
right_input,
left_bitmap,
right_bitmap,
})))
}
other => Ok(Resolved::Plan(other)),
}
}