use crate::control::state::SharedState;
use crate::types::{DatabaseId, Lsn, TenantId, TraceId};
use nodedb_physical::physical_plan::{BspSuperstepResult, PhysicalPlan};
use super::dispatch::NodeLevelResult;
use super::fanout::gather_graph_op_all_cores;
pub(super) async fn fan_bsp_all_cores(
state: &SharedState,
tenant_id: TenantId,
database_id: DatabaseId,
plan: PhysicalPlan,
trace_id: TraceId,
) -> crate::Result<NodeLevelResult> {
let responses =
gather_graph_op_all_cores(state, tenant_id, database_id, plan, trace_id, None, "bsp")
.await?;
let mut parts: Vec<BspSuperstepResult> = Vec::with_capacity(responses.len());
for resp in responses {
let part = if resp.payload.is_empty() {
BspSuperstepResult::default()
} else {
zerompk::from_msgpack::<BspSuperstepResult>(resp.payload.as_ref()).map_err(|e| {
crate::Error::Codec {
detail: format!("bsp gather: result decode: {e}"),
}
})?
};
parts.push(part);
}
let merged = merge_bsp_results(parts);
let payload = zerompk::to_msgpack_vec(&merged).map_err(|e| crate::Error::Serialization {
format: "msgpack".into(),
detail: format!("bsp gather: merged result encode: {e}"),
})?;
Ok(NodeLevelResult {
payload,
watermark_lsn: Lsn::ZERO,
read_version_lsn: Lsn::ZERO,
})
}
fn merge_bsp_results(parts: Vec<BspSuperstepResult>) -> BspSuperstepResult {
let mut out = BspSuperstepResult::default();
for p in parts {
out.local_delta += p.local_delta;
out.vertex_count += p.vertex_count;
out.outbound.extend(p.outbound);
out.node_names.extend(p.node_names);
out.rank_vec.extend(p.rank_vec);
out.dangling_sum += p.dangling_sum;
out.seed_hits += p.seed_hits;
}
out
}