use crate::bridge::envelope::{ErrorCode, PhysicalPlan, Status};
use crate::control::server::dispatch_utils::dispatch_to_data_plane_with_txn;
use crate::control::server::payload_merge::{encode_msgpack_array, extract_msgpack_elements};
use crate::control::state::SharedState;
use crate::types::{DatabaseId, TenantId, TraceId, TxnId, VShardId};
use super::gather::{GatherOutcome, gather_all_cores};
pub async fn gather_single_node(
state: &SharedState,
tenant_id: TenantId,
database_id: DatabaseId,
plan: PhysicalPlan,
trace_id: TraceId,
txn_id: Option<TxnId>,
) -> crate::Result<GatherOutcome> {
if nodedb_physical::physical_plan::plan_contains_cluster_partitioned_leaf(&plan) {
return gather_all_cores(state, tenant_id, database_id, plan, trace_id, txn_id).await;
}
if let Some(collection) = plan.collection() {
let vshard_id = VShardId::from_collection_in_database(database_id, collection);
return gather_single_owning_core(
state,
tenant_id,
database_id,
plan,
vshard_id,
trace_id,
txn_id,
)
.await;
}
gather_all_cores(state, tenant_id, database_id, plan, trace_id, txn_id).await
}
pub async fn gather_single_owning_core(
state: &SharedState,
tenant_id: TenantId,
database_id: DatabaseId,
plan: PhysicalPlan,
vshard_id: VShardId,
trace_id: TraceId,
txn_id: Option<TxnId>,
) -> crate::Result<GatherOutcome> {
let resp = Box::pin(dispatch_to_data_plane_with_txn(
state,
tenant_id,
database_id,
vshard_id,
plan,
trace_id,
txn_id,
))
.await?;
if resp.status == Status::Error
&& let Some(ec) = resp.error_code.as_deref()
&& !matches!(ec, ErrorCode::NotFound)
{
return Err(crate::Error::Dispatch {
detail: format!("{ec:?}"),
});
}
let payload_bytes: &[u8] = resp.payload.as_ref();
let all_elements = extract_msgpack_elements(payload_bytes);
let merged_array = encode_msgpack_array(&all_elements);
Ok(GatherOutcome {
raw: payload_bytes.to_vec(),
merged_array,
watermark_lsn: resp.watermark_lsn,
read_version_lsn: resp.read_version_lsn,
shard_watermarks: vec![(vshard_id, resp.watermark_lsn)],
})
}