use std::pin::Pin;
use std::sync::{Arc, Mutex, RwLock};
use nodedb_cluster::calvin::{CalvinCompletionRegistry, SEQUENCER_GROUP_ID, SequencerStateMachine};
use nodedb_cluster::distributed_array::{ArrayLocalExecutor, handle_array_shard_rpc};
use nodedb_cluster::vshard_handler::{DispatchTarget, dispatch_by_type};
use nodedb_cluster::wire::VShardEnvelope;
use crate::control::cluster::calvin::scheduler::metrics::SchedulerMetrics;
use crate::control::cluster::calvin::scheduler::read_applied_recovery;
use crate::control::cluster::calvin::{
ReadResultEvent, Scheduler, SchedulerConfig, SchedulerParams,
};
use crate::control::cluster::handle::ClusterHandle;
use crate::control::state::SharedState;
use crate::event::cross_shard::CrossShardReceiver;
pub(super) fn build_vshard_handler(
array_executor: Arc<dyn ArrayLocalExecutor>,
cross_shard_receiver: Arc<CrossShardReceiver>,
) -> nodedb_cluster::VShardEnvelopeHandler {
Arc::new(move |bytes: Vec<u8>| {
let executor = array_executor.clone();
let receiver = Arc::clone(&cross_shard_receiver);
let fut: Pin<
Box<dyn std::future::Future<Output = nodedb_cluster::error::Result<Vec<u8>>> + Send>,
> = Box::pin(async move {
let envelope = VShardEnvelope::from_bytes(&bytes).ok_or_else(|| {
nodedb_cluster::error::ClusterError::Codec {
detail: "vshard_handler: failed to deserialize VShardEnvelope".into(),
}
})?;
let target = dispatch_by_type(&envelope);
match target {
DispatchTarget::ArrayShard => {
let opcode = envelope.msg_type as u32;
let resp_payload = handle_array_shard_rpc(
opcode,
envelope.vshard_id,
&envelope.payload,
&executor,
)
.await?;
let resp_opcode = opcode + 1;
let resp_msg_type = resolve_vshard_msg_type(resp_opcode)?;
let resp_envelope = VShardEnvelope::new(
resp_msg_type,
envelope.target_node,
envelope.source_node,
envelope.vshard_id,
resp_payload,
);
Ok(resp_envelope.to_bytes())
}
DispatchTarget::EventPlane => Ok(receiver.handle_envelope(bytes).await),
other => Err(nodedb_cluster::error::ClusterError::Transport {
detail: format!(
"vshard_handler: no handler registered for dispatch target {other:?}"
),
}),
}
});
fut
})
}
type ReadResultSenders =
Arc<Mutex<std::collections::BTreeMap<u32, tokio::sync::mpsc::Sender<ReadResultEvent>>>>;
fn hosted_vshards(routing: &RwLock<nodedb_cluster::RoutingTable>, node_id: u64) -> Vec<u32> {
let routing = routing.read().unwrap_or_else(|p| p.into_inner());
let mut vshards = Vec::new();
for (group_id, info) in routing.group_members() {
if info.members.contains(&node_id) {
vshards.extend(routing.vshards_for_group(*group_id));
}
}
vshards.sort_unstable();
vshards.dedup();
vshards
}
struct ReconcileSchedulersParams<'a> {
node_id: u64,
routing: &'a Arc<RwLock<nodedb_cluster::RoutingTable>>,
shared: &'a Arc<SharedState>,
raft_loop_handle: &'a Arc<Mutex<nodedb_cluster::multi_raft::MultiRaft>>,
sequencer_state_machine: &'a Arc<Mutex<SequencerStateMachine>>,
calvin_read_result_senders: &'a ReadResultSenders,
calvin_completion_registry: &'a Arc<CalvinCompletionRegistry>,
scheduler_config: &'a SchedulerConfig,
}
fn reconcile_vshard_schedulers(params: ReconcileSchedulersParams<'_>) -> crate::Result<usize> {
let ReconcileSchedulersParams {
node_id,
routing,
shared,
raft_loop_handle,
sequencer_state_machine,
calvin_read_result_senders,
calvin_completion_registry,
scheduler_config,
} = params;
let mut spawned = 0usize;
for vshard_id in hosted_vshards(routing, node_id) {
if calvin_read_result_senders
.lock()
.unwrap_or_else(|p| p.into_inner())
.contains_key(&vshard_id)
{
continue;
}
let recovery = read_applied_recovery(&shared.wal, vshard_id)?;
let (sequenced_tx, sequenced_rx) =
tokio::sync::mpsc::channel(scheduler_config.channel_capacity);
let first_available = raft_loop_handle
.lock()
.unwrap_or_else(|p| p.into_inner())
.first_available_index(SEQUENCER_GROUP_ID)
.unwrap_or(1);
{
let mut sm = sequencer_state_machine
.lock()
.unwrap_or_else(|p| p.into_inner());
sm.set_vshard_sender(vshard_id, sequenced_tx);
sm.arm_catch_up_from(vshard_id, first_available);
}
let (read_result_tx, read_result_rx) =
tokio::sync::mpsc::channel(scheduler_config.channel_capacity);
calvin_read_result_senders
.lock()
.unwrap_or_else(|p| p.into_inner())
.insert(vshard_id, read_result_tx);
let lock_manager = Arc::new(Mutex::new(
crate::control::cluster::calvin::scheduler::lock_manager::LockManager::new(),
));
shared
.calvin_lock_managers
.lock()
.unwrap_or_else(|p| p.into_inner())
.insert(vshard_id, Arc::clone(&lock_manager));
let (promotion_tx, promotion_rx) = tokio::sync::mpsc::unbounded_channel();
shared
.calvin_promotion_senders
.lock()
.unwrap_or_else(|p| p.into_inner())
.insert(vshard_id, promotion_tx);
let (verdict_tx, verdict_rx) =
tokio::sync::mpsc::channel(scheduler_config.channel_capacity);
calvin_completion_registry.register_verdict_signal_sender(vshard_id, verdict_tx);
let scheduler = Scheduler::new(SchedulerParams {
vshard_id,
receiver: sequenced_rx,
shared: Arc::clone(shared),
multi_raft: raft_loop_handle.clone(),
sequencer_state_machine: Arc::clone(sequencer_state_machine),
fully_applied_epoch: recovery.fully_applied_epoch,
applied_tail: recovery.applied_tail,
rebuild_target_epoch: recovery.max_applied_epoch,
config: scheduler_config.clone(),
metrics: SchedulerMetrics::new(),
read_result_rx,
lock_manager,
promotion_rx,
registry: Arc::clone(calvin_completion_registry),
verdict_rx,
});
crate::control::shutdown::spawn_loop_no_abort(
&shared.loop_registry,
&shared.shutdown,
"calvin_scheduler",
move |shutdown| async move {
scheduler.run(shutdown).await;
},
);
spawned += 1;
}
Ok(spawned)
}
pub(super) struct SpawnVshardSchedulersParams<'a> {
pub(super) handle: &'a ClusterHandle,
pub(super) shared: &'a Arc<SharedState>,
pub(super) raft_loop_handle: Arc<Mutex<nodedb_cluster::multi_raft::MultiRaft>>,
pub(super) sequencer_state_machine: &'a Arc<Mutex<SequencerStateMachine>>,
pub(super) calvin_read_result_senders: &'a ReadResultSenders,
pub(super) calvin_completion_registry: &'a Arc<CalvinCompletionRegistry>,
pub(super) scheduler_config: &'a SchedulerConfig,
}
pub(super) fn spawn_vshard_schedulers(
params: SpawnVshardSchedulersParams<'_>,
) -> crate::Result<()> {
let SpawnVshardSchedulersParams {
handle,
shared,
raft_loop_handle,
sequencer_state_machine,
calvin_read_result_senders,
calvin_completion_registry,
scheduler_config,
} = params;
let node_id = handle.node_id;
let routing = Arc::clone(&handle.routing);
reconcile_vshard_schedulers(ReconcileSchedulersParams {
node_id,
routing: &routing,
shared,
raft_loop_handle: &raft_loop_handle,
sequencer_state_machine,
calvin_read_result_senders,
calvin_completion_registry,
scheduler_config,
})?;
let shared_task = Arc::clone(shared);
let sm_task = Arc::clone(sequencer_state_machine);
let rr_task = Arc::clone(calvin_read_result_senders);
let registry_task = Arc::clone(calvin_completion_registry);
let cfg_task = scheduler_config.clone();
crate::control::shutdown::spawn_loop(
&shared.loop_registry,
&shared.shutdown,
"calvin_vshard_reconcile",
move |mut shutdown| async move {
let mut tick = tokio::time::interval(std::time::Duration::from_millis(500));
tick.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Skip);
loop {
tokio::select! {
biased;
_ = shutdown.wait_cancelled() => break,
_ = tick.tick() => {
if let Err(e) = reconcile_vshard_schedulers(ReconcileSchedulersParams {
node_id,
routing: &routing,
shared: &shared_task,
raft_loop_handle: &raft_loop_handle,
sequencer_state_machine: &sm_task,
calvin_read_result_senders: &rr_task,
calvin_completion_registry: ®istry_task,
scheduler_config: &cfg_task,
}) {
tracing::warn!(node_id, error = %e, "calvin scheduler reconcile pass failed");
}
}
}
}
},
);
Ok(())
}
pub(super) fn resolve_vshard_msg_type(
opcode: u32,
) -> nodedb_cluster::error::Result<nodedb_cluster::wire::VShardMessageType> {
let mut scratch = [0u8; 26];
scratch[0..2].copy_from_slice(&1u16.to_le_bytes()); scratch[2..4].copy_from_slice(&(opcode as u16).to_le_bytes());
VShardEnvelope::from_bytes(&scratch)
.map(|e| e.msg_type)
.ok_or_else(|| nodedb_cluster::error::ClusterError::Codec {
detail: format!("resolve_vshard_msg_type: unknown opcode {opcode}"),
})
}