use nodedb_physical::physical_plan::{ExchangeMode, ExchangeOp, PhysicalPlan, QueryOp};
use crate::control::state::SharedState;
use crate::data::executor::response_codec::flatten_to_relational_rows;
use crate::types::{DatabaseId, TenantId, TraceId, TxnId};
use crate::control::server::exchange::full_scan::full_scan_plan_for_collection;
use crate::control::server::exchange::gather::{
finalize_aggregate, gather_all_cores, gather_all_vshards,
};
use super::capture::DistributedReadCapture;
pub(super) async fn resolve_join_input(
state: &SharedState,
database_id: DatabaseId,
tenant_id: TenantId,
input: Option<Box<PhysicalPlan>>,
trace_id: TraceId,
txn_id: Option<TxnId>,
captures: &mut Vec<DistributedReadCapture>,
) -> crate::Result<Option<Box<PhysicalPlan>>> {
let Some(boxed) = input else {
return Ok(None);
};
match *boxed {
PhysicalPlan::Query(QueryOp::Exchange(ExchangeOp {
child,
mode: ExchangeMode::Broadcast,
})) => {
let outcome =
gather_all_cores(state, tenant_id, database_id, *child, trace_id, txn_id).await?;
let provider_scan = PhysicalPlan::Query(QueryOp::ProviderScan {
provider: None,
rows: flatten_to_relational_rows(&outcome.merged_array),
filters: Vec::new(),
projection: Vec::new(),
sort_keys: Vec::new(),
limit: None,
offset: 0,
distinct: false,
});
Ok(Some(Box::new(provider_scan)))
}
PhysicalPlan::Query(QueryOp::Exchange(ExchangeOp {
mode: ExchangeMode::Shuffle { .. },
..
})) => Err(crate::Error::Internal {
detail: "ExchangeMode::Shuffle is only valid wrapping a complete hash join, \
not as a join input"
.into(),
}),
PhysicalPlan::Query(QueryOp::Exchange(ExchangeOp {
mode: ExchangeMode::ShuffleAggregate { .. },
..
})) => Err(crate::Error::Internal {
detail: "ExchangeMode::ShuffleAggregate is only valid wrapping a complete root \
aggregate, not as a join input"
.into(),
}),
PhysicalPlan::Query(QueryOp::Exchange(ExchangeOp {
child,
mode: ExchangeMode::Gather { as_aggregate },
})) => {
let child_collection: Option<String> = if txn_id.is_some() {
child.collection().map(str::to_owned)
} else {
None
};
let outcome =
gather_all_cores(state, tenant_id, database_id, *child, trace_id, txn_id).await?;
if let Some(coll) = child_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 merged = if as_aggregate {
finalize_aggregate(&outcome.merged_array)
} else {
outcome.merged_array
};
let provider_scan = PhysicalPlan::Query(QueryOp::ProviderScan {
provider: None,
rows: flatten_to_relational_rows(&merged),
filters: Vec::new(),
projection: Vec::new(),
sort_keys: Vec::new(),
limit: None,
offset: 0,
distinct: false,
});
Ok(Some(Box::new(provider_scan)))
}
other => Ok(Some(Box::new(other))),
}
}
pub(super) async fn gather_join_build_side(
state: &SharedState,
database_id: DatabaseId,
tenant_id: TenantId,
collection: &str,
trace_id: TraceId,
txn_id: Option<TxnId>,
captures: &mut Vec<DistributedReadCapture>,
) -> crate::Result<Option<Box<PhysicalPlan>>> {
let Some(scan_plan) =
crate::control::server::exchange::full_scan::full_scan_plan_for_collection(
state,
database_id,
tenant_id,
collection,
)?
else {
return Ok(None);
};
let capture_plan = if txn_id.is_some() {
Some(scan_plan.clone())
} else {
None
};
let outcome = Box::pin(gather_all_vshards(
state,
tenant_id,
database_id,
scan_plan,
trace_id,
txn_id,
))
.await?;
if let Some(scan_plan) = capture_plan {
captures.push(DistributedReadCapture {
scan_plan,
read_version_lsn: outcome.read_version_lsn,
});
}
Ok(Some(Box::new(PhysicalPlan::Query(QueryOp::ProviderScan {
provider: None,
rows: flatten_to_relational_rows(&outcome.merged_array),
filters: Vec::new(),
projection: Vec::new(),
sort_keys: Vec::new(),
limit: None,
offset: 0,
distinct: false,
}))))
}