use tracing::warn;
use crate::conf_change::ConfChange;
use crate::forward::PlanExecutor;
use super::super::loop_core::{CommitApplier, RaftLoop};
impl<A: CommitApplier, P: PlanExecutor> RaftLoop<A, P> {
pub(super) fn apply_group_commits(&self, group_id: u64, group_ready: &nodedb_raft::Ready) {
for entry in &group_ready.committed_entries {
if let Some(cc) = ConfChange::from_entry_data(&entry.data) {
let mut mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
if let Err(e) = mr.apply_conf_change(group_id, &cc) {
warn!(group_id, error = %e, "failed to apply conf change");
}
}
}
let last_applied = if group_id == crate::metadata_group::METADATA_GROUP_ID {
let pairs: Vec<(u64, Vec<u8>)> = group_ready
.committed_entries
.iter()
.filter(|e| ConfChange::from_entry_data(&e.data).is_none())
.map(|e| (e.index, e.data.clone()))
.collect();
self.metadata_applier.apply(&pairs)
} else {
self.applier
.apply_committed(group_id, &group_ready.committed_entries)
};
if last_applied > 0 {
let mut mr = self.multi_raft.lock().unwrap_or_else(|p| p.into_inner());
if let Err(e) = mr.advance_applied(group_id, last_applied) {
warn!(group_id, error = %e, "failed to advance applied index");
} else if group_id == crate::metadata_group::METADATA_GROUP_ID {
self.group_watchers.bump(group_id, last_applied);
}
}
if group_id == crate::metadata_group::METADATA_GROUP_ID && !*self.ready_watch.borrow() {
let _ = self.ready_watch.send(true);
}
if group_id == crate::metadata_group::METADATA_GROUP_ID {
let is_leader = self
.multi_raft
.lock()
.unwrap_or_else(|p| p.into_inner())
.group_role_is_leader(group_id);
let was_leader = self
.prev_metadata_leader
.swap(is_leader, std::sync::atomic::Ordering::AcqRel);
if is_leader && !was_leader {
if let Some(catalog) = self.catalog.as_ref() {
match crate::cluster_epoch::bump_local_cluster_epoch(catalog) {
Ok(new_epoch) => tracing::info!(
node = self.node_id,
new_epoch,
"bumped cluster epoch on metadata-group leadership acquisition"
),
Err(e) => tracing::warn!(
node = self.node_id,
error = %e,
"failed to persist bumped cluster epoch (in-memory value advanced anyway)"
),
}
} else {
let _ = crate::cluster_epoch::observe_peer_cluster_epoch(
crate::cluster_epoch::current_local_cluster_epoch() + 1,
);
}
}
}
}
}