use std::collections::HashSet;
use tracing::debug;
use crate::conf_change::{ConfChange, ConfChangeType};
use crate::forward::PlanExecutor;
use super::loop_core::{CommitApplier, RaftLoop};
pub(super) fn plan_entering_learners(
actual_voters: &[u64],
actual_learners: &[u64],
placement: &[u64],
) -> Vec<u64> {
let mut out: Vec<u64> = placement
.iter()
.copied()
.filter(|n| !actual_voters.contains(n) && !actual_learners.contains(n))
.collect();
out.sort_unstable();
out.dedup();
out
}
pub(super) fn plan_entering_mounts(
self_id: u64,
hosted: &HashSet<u64>,
routing: &crate::routing::RoutingTable,
) -> Vec<u64> {
let mut out: Vec<u64> = routing
.group_members()
.iter()
.filter_map(|(&gid, info)| {
if gid == crate::metadata_group::METADATA_GROUP_ID
|| gid == crate::calvin::sequencer::SEQUENCER_GROUP_ID
{
return None;
}
if hosted.contains(&gid) {
return None;
}
match &info.placement {
Some(p) if p.contains(&self_id) => Some(gid),
_ => None,
}
})
.collect();
out.sort_unstable();
out
}
pub(super) fn plan_leaving_learners(actual_learners: &[u64], placement: &[u64]) -> Vec<u64> {
let mut leaving: Vec<u64> = actual_learners
.iter()
.copied()
.filter(|n| !placement.contains(n))
.collect();
leaving.sort_unstable();
leaving
}
pub(super) fn plan_leaving_voters(actual_voters: &[u64], placement: &[u64], rf: usize) -> Vec<u64> {
let removable = actual_voters.len().saturating_sub(rf);
if removable == 0 {
return Vec::new();
}
let mut leaving: Vec<u64> = actual_voters
.iter()
.copied()
.filter(|n| !placement.contains(n))
.collect();
leaving.sort_unstable();
leaving.truncate(removable);
leaving
}
impl<A: CommitApplier, P: PlanExecutor> RaftLoop<A, P> {
pub(super) fn mount_entering_groups(&self) {
let mount: Option<(u64, Vec<u64>, Vec<u64>)> = {
let mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
let hosted: HashSet<u64> = mr.group_ids().into_iter().collect();
let routing = mr.routing();
let routing = routing.read().unwrap_or_else(|p| p.into_inner());
plan_entering_mounts(self.node_id, &hosted, &routing)
.into_iter()
.next()
.and_then(|gid| {
routing.group_info(gid).map(|info| {
let voters: Vec<u64> = info
.members
.iter()
.copied()
.filter(|&id| id != self.node_id)
.collect();
let other_learners: Vec<u64> = info
.learners
.iter()
.copied()
.filter(|&id| id != self.node_id)
.collect();
(gid, voters, other_learners)
})
})
};
if let Some((gid, voters, other_learners)) = mount {
let mut mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
if mr.contains_group(gid) {
return;
}
match mr.add_group_as_learner(gid, voters, other_learners) {
Ok(()) => {
debug!(
group_id = gid,
node_id = self.node_id,
"mount: added local learner replica for placed group"
);
}
Err(e) => {
debug!(
group_id = gid,
node_id = self.node_id,
error = %e,
"mount: add_group_as_learner failed; retrying next tick"
);
}
}
}
}
pub(super) fn converge_entering_learners(&self) {
let additions: Vec<(u64, u64)> = {
let mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
let group_ids = mr.group_ids();
let mut out = Vec::new();
for gid in group_ids {
if gid == crate::metadata_group::METADATA_GROUP_ID
|| gid == crate::calvin::sequencer::SEQUENCER_GROUP_ID
{
continue;
}
if !mr.group_role_is_leader(gid) {
continue;
}
let placement: Option<Vec<u64>> = mr
.routing()
.read()
.unwrap_or_else(|p| p.into_inner())
.group_info(gid)
.and_then(|info| info.placement.clone());
let Some(placement) = placement else {
continue;
};
let Some(m) = mr.group_membership(gid) else {
continue;
};
for node_id in plan_entering_learners(&m.voters, &m.learners, &placement) {
out.push((gid, node_id));
}
}
out
};
for (group_id, node_id) in additions {
let mut mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
let change = ConfChange {
change_type: ConfChangeType::AddLearner,
node_id,
};
match mr.propose_conf_change(group_id, &change) {
Ok((_gid, idx)) => {
debug!(
group_id,
node_id,
log_index = idx,
"convergence: proposed AddLearner"
);
}
Err(e) => {
debug!(
group_id,
node_id,
error = %e,
"convergence: AddLearner deferred"
);
}
}
}
}
pub(super) fn converge_leaving_voters(&self) {
let rf = self.replication_factor() as usize;
let removals: Vec<(u64, u64)> = {
let mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
let group_ids = mr.group_ids();
let mut out = Vec::new();
for gid in group_ids {
if gid == crate::metadata_group::METADATA_GROUP_ID
|| gid == crate::calvin::sequencer::SEQUENCER_GROUP_ID
{
continue;
}
if !mr.group_role_is_leader(gid) {
continue;
}
let placement: Option<Vec<u64>> = mr
.routing()
.read()
.unwrap_or_else(|p| p.into_inner())
.group_info(gid)
.and_then(|info| info.placement.clone());
let Some(placement) = placement else {
continue;
};
let Some(m) = mr.group_membership(gid) else {
continue;
};
for node_id in plan_leaving_voters(&m.voters, &placement, rf) {
if node_id == m.leader_id {
debug!(
group_id = gid,
node_id,
"convergence: leaving voter is group leader; \
deferring removal (needs leadership transfer)"
);
continue;
}
out.push((gid, node_id));
break;
}
}
out
};
for (group_id, node_id) in removals {
let mut mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
let change = ConfChange {
change_type: ConfChangeType::RemoveNode,
node_id,
};
match mr.propose_conf_change(group_id, &change) {
Ok((_gid, idx)) => {
debug!(
group_id,
node_id,
log_index = idx,
"convergence: proposed RemoveNode"
);
}
Err(e) => {
debug!(
group_id,
node_id,
error = %e,
"convergence: RemoveNode deferred"
);
}
}
}
}
pub(super) fn converge_leaving_learners(&self) {
let removals: Vec<(u64, u64)> = {
let mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
let group_ids = mr.group_ids();
let mut out = Vec::new();
for gid in group_ids {
if gid == crate::metadata_group::METADATA_GROUP_ID
|| gid == crate::calvin::sequencer::SEQUENCER_GROUP_ID
{
continue;
}
if !mr.group_role_is_leader(gid) {
continue;
}
let placement: Option<Vec<u64>> = mr
.routing()
.read()
.unwrap_or_else(|p| p.into_inner())
.group_info(gid)
.and_then(|info| info.placement.clone());
let Some(placement) = placement else {
continue;
};
let Some(m) = mr.group_membership(gid) else {
continue;
};
for node_id in plan_leaving_learners(&m.learners, &placement) {
out.push((gid, node_id));
}
}
out
};
for (group_id, node_id) in removals {
let mut mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
let change = ConfChange {
change_type: ConfChangeType::RemoveLearner,
node_id,
};
match mr.propose_conf_change(group_id, &change) {
Ok((_gid, idx)) => {
debug!(
group_id,
node_id,
log_index = idx,
"convergence: proposed RemoveLearner"
);
}
Err(e) => {
debug!(
group_id,
node_id,
error = %e,
"convergence: RemoveLearner deferred"
);
}
}
}
}
}
#[cfg(test)]
mod tests {
use super::{
plan_entering_learners, plan_entering_mounts, plan_leaving_learners, plan_leaving_voters,
};
use crate::routing::{GroupInfo, RoutingTable};
use std::collections::HashSet;
fn group(members: Vec<u64>, learners: Vec<u64>, placement: Option<Vec<u64>>) -> GroupInfo {
GroupInfo {
leader: members.first().copied().unwrap_or(0),
members,
learners,
placement,
}
}
fn routing(groups: Vec<(u64, GroupInfo)>) -> RoutingTable {
RoutingTable::from_parts(Vec::new(), groups.into_iter().collect())
}
fn hosted(ids: &[u64]) -> HashSet<u64> {
ids.iter().copied().collect()
}
#[test]
fn mount_when_placement_contains_self_and_not_hosted() {
let rt = routing(vec![(5, group(vec![1, 3], vec![], Some(vec![1, 3, 4])))]);
assert_eq!(plan_entering_mounts(4, &hosted(&[]), &rt), vec![5]);
}
#[test]
fn skip_already_hosted_group() {
let rt = routing(vec![(5, group(vec![1, 3], vec![], Some(vec![1, 3, 4])))]);
assert_eq!(
plan_entering_mounts(4, &hosted(&[5]), &rt),
Vec::<u64>::new()
);
}
#[test]
fn skip_metadata_and_sequencer_groups() {
let meta = crate::metadata_group::METADATA_GROUP_ID;
let seq = crate::calvin::sequencer::SEQUENCER_GROUP_ID;
let rt = routing(vec![
(meta, group(vec![1], vec![], Some(vec![1, 4]))),
(seq, group(vec![1], vec![], Some(vec![1, 4]))),
]);
assert_eq!(
plan_entering_mounts(4, &hosted(&[]), &rt),
Vec::<u64>::new()
);
}
#[test]
fn skip_when_placement_none() {
let rt = routing(vec![(5, group(vec![1, 3], vec![], None))]);
assert_eq!(
plan_entering_mounts(4, &hosted(&[]), &rt),
Vec::<u64>::new()
);
}
#[test]
fn skip_when_placement_excludes_self() {
let rt = routing(vec![(5, group(vec![1, 3], vec![], Some(vec![1, 3, 2])))]);
assert_eq!(
plan_entering_mounts(4, &hosted(&[]), &rt),
Vec::<u64>::new()
);
}
#[test]
fn mount_self_only_placement() {
let rt = routing(vec![(7, group(vec![], vec![], Some(vec![4])))]);
assert_eq!(plan_entering_mounts(4, &hosted(&[]), &rt), vec![7]);
}
#[test]
fn mount_output_is_sorted() {
let rt = routing(vec![
(9, group(vec![1], vec![], Some(vec![1, 4]))),
(3, group(vec![1], vec![], Some(vec![1, 4]))),
(6, group(vec![1], vec![], Some(vec![1, 4]))),
]);
assert_eq!(plan_entering_mounts(4, &hosted(&[]), &rt), vec![3, 6, 9]);
}
#[test]
fn entering_nodes_not_yet_in_membership() {
assert_eq!(
plan_entering_learners(&[1, 2], &[], &[1, 2, 3, 4]),
vec![3, 4]
);
}
#[test]
fn existing_learner_not_re_added() {
assert_eq!(
plan_entering_learners(&[1, 2], &[3], &[1, 2, 3]),
Vec::<u64>::new()
);
}
#[test]
fn no_entering_when_placement_subset_of_members() {
assert_eq!(
plan_entering_learners(&[1, 2, 3], &[4], &[1, 2, 3, 4]),
Vec::<u64>::new()
);
}
#[test]
fn result_is_sorted_and_deduped() {
assert_eq!(
plan_entering_learners(&[1], &[], &[4, 2, 4, 3]),
vec![2, 3, 4]
);
}
#[test]
fn empty_placement_returns_empty() {
assert_eq!(
plan_entering_learners(&[1, 2], &[3], &[]),
Vec::<u64>::new()
);
}
#[test]
fn leaving_never_below_rf_cap() {
assert_eq!(plan_leaving_voters(&[1, 2, 3, 4], &[1, 2, 3], 3), vec![4]);
}
#[test]
fn leaving_blocked_when_at_rf() {
assert_eq!(
plan_leaving_voters(&[1, 2, 3], &[1, 2], 3),
Vec::<u64>::new()
);
}
#[test]
fn leaving_empty_when_placement_superset() {
assert_eq!(
plan_leaving_voters(&[1, 2, 3], &[1, 2, 3, 4], 3),
Vec::<u64>::new()
);
}
#[test]
fn leaving_multiple_capped_by_rf() {
assert_eq!(
plan_leaving_voters(&[1, 2, 3, 4, 5], &[1, 2, 3], 3),
vec![4, 5]
);
}
#[test]
fn leaving_only_returns_voters_not_in_placement() {
assert_eq!(
plan_leaving_voters(&[1, 2, 3, 4, 5], &[1, 2], 3),
vec![3, 4]
);
}
#[test]
fn leaving_output_is_sorted() {
assert_eq!(plan_leaving_voters(&[9, 2, 7, 4], &[], 0), vec![2, 4, 7, 9]);
}
#[test]
fn learner_not_in_placement_is_removed() {
assert_eq!(plan_leaving_learners(&[4], &[1, 2, 3]), vec![4]);
}
#[test]
fn learner_in_placement_is_kept() {
assert_eq!(plan_leaving_learners(&[3], &[1, 2, 3]), Vec::<u64>::new());
}
#[test]
fn mixed_learners_only_out_of_placement_removed() {
assert_eq!(plan_leaving_learners(&[3, 4], &[1, 2, 3]), vec![4]);
}
#[test]
fn empty_placement_superset_removes_nothing() {
assert_eq!(
plan_leaving_learners(&[2, 3], &[1, 2, 3]),
Vec::<u64>::new()
);
}
#[test]
fn leaving_learners_output_is_sorted() {
assert_eq!(plan_leaving_learners(&[9, 2, 7], &[]), vec![2, 7, 9]);
}
#[test]
fn no_learners_returns_empty() {
assert_eq!(plan_leaving_learners(&[], &[1, 2, 3]), Vec::<u64>::new());
}
}