use tracing::{debug, warn};
use crate::forward::PlanExecutor;
use super::loop_core::{CommitApplier, RaftLoop};
impl<A: CommitApplier, P: PlanExecutor> RaftLoop<A, P> {
pub(super) fn reconcile_placement(&self) {
let changes: Vec<(u64, Vec<u64>)> = {
let mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
if !mr.group_role_is_leader(crate::metadata_group::METADATA_GROUP_ID) {
return;
}
let active_nodes: Vec<u64> = {
let topo = self.topology.read().unwrap_or_else(|p| p.into_inner());
topo.active_nodes().iter().map(|n| n.node_id).collect()
};
if active_nodes.is_empty() {
return;
}
let data_group_ids: Vec<u64> = mr
.group_ids()
.into_iter()
.filter(|g| {
*g != crate::metadata_group::METADATA_GROUP_ID
&& *g != crate::calvin::sequencer::SEQUENCER_GROUP_ID
})
.collect();
let target = crate::rebalancer::placement::compute_placement(
&active_nodes,
&data_group_ids,
self.replication_factor(),
);
let routing = mr.routing();
let routing = routing.read().unwrap_or_else(|p| p.into_inner());
crate::rebalancer::placement::compute_placement_changes(&routing, &target)
};
for (group_id, placement) in changes {
let entry = crate::metadata_group::entry::MetadataEntry::RoutingChange(
crate::metadata_group::entry::RoutingChange::SetPlacement {
group_id,
placement,
},
);
let bytes = match crate::metadata_group::codec::encode_entry(&entry) {
Ok(b) => b,
Err(e) => {
warn!(group_id, error = %e, "placement: encode SetPlacement failed");
continue;
}
};
match self.propose_to_metadata_group(bytes) {
Ok(idx) => {
debug!(
group_id,
log_index = idx,
"placement: proposed SetPlacement"
)
}
Err(e) => {
debug!(group_id, error = %e, "placement: SetPlacement proposal deferred")
}
}
}
}
}