use std::sync::Arc;
use crate::bridge::envelope::{Payload, PhysicalPlan};
use crate::control::gateway::dispatcher::{DispatchRouteParams, dispatch_route};
use crate::control::gateway::router::resolve_decision;
use crate::control::gateway::version_set::GatewayVersionSet;
use crate::control::gateway::{RouteDecision, TaskRoute};
use crate::control::server::exchange::execute_plan_all_local_cores;
use crate::control::state::SharedState;
use crate::types::{DatabaseId, TenantId, TraceId};
pub(crate) fn resolve_for_vshard(state: &SharedState, vshard_id: u32) -> RouteDecision {
let routing_guard = state
.cluster_routing
.as_ref()
.map(|rw| rw.read().unwrap_or_else(|p| p.into_inner()));
let raft_snapshot: Vec<nodedb_cluster::GroupStatus> =
state.raft_status_fn.get().map(|f| f()).unwrap_or_default();
let live_leader = move |group_id: u64| -> u64 {
raft_snapshot
.iter()
.find(|gs| gs.group_id == group_id)
.map(|gs| gs.leader_id)
.unwrap_or(0)
};
let live_lookup: Option<&dyn Fn(u64) -> u64> = if state.raft_status_fn.get().is_some() {
Some(&live_leader)
} else {
None
};
resolve_decision(
vshard_id,
state.node_id,
routing_guard.as_deref(),
live_lookup,
)
}
pub(in crate::control::server::graph_dispatch) struct DispatchSuperstepParams<'a> {
pub(in crate::control::server::graph_dispatch) tenant_id: TenantId,
pub(in crate::control::server::graph_dispatch) database_id: DatabaseId,
pub(in crate::control::server::graph_dispatch) deadline_ms: u64,
pub(in crate::control::server::graph_dispatch) node_id: u64,
pub(in crate::control::server::graph_dispatch) is_local: bool,
pub(in crate::control::server::graph_dispatch) route_vshard: u32,
pub(in crate::control::server::graph_dispatch) plan: PhysicalPlan,
pub(in crate::control::server::graph_dispatch) version_set: &'a GatewayVersionSet,
}
pub(in crate::control::server::graph_dispatch) async fn dispatch_superstep_to_node(
shared_arc: &Arc<SharedState>,
args: DispatchSuperstepParams<'_>,
) -> crate::Result<Payload> {
let DispatchSuperstepParams {
tenant_id,
database_id,
deadline_ms,
node_id,
is_local,
route_vshard,
plan,
version_set,
} = args;
if is_local {
let node_result = execute_plan_all_local_cores(
shared_arc.as_ref(),
tenant_id,
database_id,
plan,
TraceId::ZERO,
None,
)
.await?;
Ok(Payload::from_vec(node_result.payload))
} else {
let route = TaskRoute {
plan,
decision: RouteDecision::Remote {
node_id,
vshard_id: route_vshard as u64,
},
vshard_id: route_vshard,
};
let payloads = dispatch_route(DispatchRouteParams {
route,
shared: shared_arc,
tenant_id,
database_id,
trace_id: TraceId::ZERO,
deadline_ms,
version_set,
txn_id: None,
})
.await?
.payloads;
payloads
.into_iter()
.next()
.map(Payload::from_vec)
.ok_or_else(|| crate::Error::Internal {
detail: format!("graph superstep: node={node_id} returned no payload"),
})
}
}
pub(crate) fn gateway_shared(state: &SharedState) -> crate::Result<Arc<SharedState>> {
let gateway = state.gateway.get().ok_or_else(|| crate::Error::Internal {
detail: "graph scatter: cluster routing present but gateway unavailable for \
remote dispatch"
.into(),
})?;
gateway.shared()
}