use tracing::{debug, error};
use crate::forward::PlanExecutor;
use super::super::loop_core::{CommitApplier, RaftLoop};
const PLACEMENT_RECONCILE_TICK_INTERVAL: u64 = 100;
const ORPHAN_PARTIAL_GC_TICK_INTERVAL: u64 = 6000;
impl<A: CommitApplier, P: PlanExecutor> RaftLoop<A, P> {
pub(in crate::raft_loop) fn do_tick(&self) {
let ready = {
let mut mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
match mr.tick() {
Ok(ready) => ready,
Err(e) => {
error!(
error = %e,
"raft tick failed to persist hard state durably; \
skipping message/vote dispatch for this tick"
);
return;
}
}
};
if !ready.is_empty() {
self.dispatch_outbound_messages(&ready.groups);
for (group_id, group_ready) in ready.groups {
if !group_ready.committed_entries.is_empty() {
self.apply_group_commits(group_id, &group_ready);
}
if !group_ready.snapshots_needed.is_empty() {
self.dispatch_group_snapshots(group_id, group_ready.snapshots_needed);
}
}
}
self.mount_entering_groups();
self.converge_entering_learners();
self.promote_ready_learners();
self.transfer_leadership_for_leaving_voters();
self.converge_leaving_voters();
self.converge_leaving_learners();
let tick = self
.tick_count
.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
if tick.is_multiple_of(PLACEMENT_RECONCILE_TICK_INTERVAL) {
self.reconcile_placement();
}
if tick.is_multiple_of(ORPHAN_PARTIAL_GC_TICK_INTERVAL)
&& let Some(ref dir) = self.data_dir
{
match crate::install_snapshot::gc::sweep_orphans(dir, self.orphan_partial_max_age_secs)
{
Ok((removed, errs)) => {
if removed > 0 {
tracing::info!(
removed,
"periodic gc: removed orphaned partial snapshot files"
);
}
for e in errs {
tracing::warn!(error = %e, "periodic gc: partial snapshot GC error");
}
}
Err(e) => {
tracing::warn!(error = %e, "periodic gc: failed to sweep partial snapshot directory");
}
}
}
}
fn promote_ready_learners(&self) {
let promotions: Vec<(u64, u64)> = {
let mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
let group_ids = mr.group_ids();
group_ids
.into_iter()
.flat_map(|gid| {
let placement: Option<Vec<u64>> = mr
.routing()
.read()
.unwrap_or_else(|p| p.into_inner())
.group_info(gid)
.and_then(|info| info.placement.clone());
mr.ready_learners(gid)
.into_iter()
.filter(move |&learner| {
should_promote_learner(placement.as_deref(), learner)
})
.map(move |learner| (gid, learner))
})
.collect()
};
for (group_id, learner_id) in promotions {
let mut mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
let change = crate::conf_change::ConfChange {
change_type: crate::conf_change::ConfChangeType::PromoteLearner,
node_id: learner_id,
};
match mr.propose_conf_change(group_id, &change) {
Ok((_gid, idx)) => {
debug!(
group_id,
learner_id,
log_index = idx,
"proposed learner promotion"
);
}
Err(e) => {
debug!(
group_id,
learner_id,
error = %e,
"learner promotion proposal deferred"
);
}
}
}
}
}
fn should_promote_learner(placement: Option<&[u64]>, learner_id: u64) -> bool {
match placement {
None => true,
Some(set) => set.contains(&learner_id),
}
}
#[cfg(test)]
mod tests {
use super::should_promote_learner;
#[test]
fn no_placement_always_promotes() {
assert!(should_promote_learner(None, 1));
assert!(should_promote_learner(None, 42));
}
#[test]
fn placement_contains_learner_promotes() {
assert!(should_promote_learner(Some(&[1, 2, 3]), 2));
}
#[test]
fn placement_missing_learner_holds() {
assert!(!should_promote_learner(Some(&[1, 2, 3]), 99));
}
#[test]
fn empty_placement_promotes_nobody() {
assert!(!should_promote_learner(Some(&[]), 1));
assert!(!should_promote_learner(Some(&[]), 42));
}
}