use std::cmp;
use self::inflights::Inflights;
use crate::raft::INVALID_INDEX;
pub mod inflights;
pub mod progress_set;
#[derive(Debug, PartialEq, Clone, Copy)]
pub enum ProgressState {
Probe,
Replicate,
Snapshot,
}
impl Default for ProgressState {
fn default() -> ProgressState {
ProgressState::Probe
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct Progress {
pub matched: u64,
pub next_idx: u64,
pub state: ProgressState,
pub paused: bool,
pub pending_snapshot: u64,
pub pending_request_snapshot: u64,
pub recent_active: bool,
pub ins: Inflights,
}
impl Progress {
pub fn new(next_idx: u64, ins_size: usize) -> Self {
Progress {
matched: 0,
next_idx,
state: ProgressState::default(),
paused: false,
pending_snapshot: 0,
pending_request_snapshot: 0,
recent_active: false,
ins: Inflights::new(ins_size),
}
}
fn reset_state(&mut self, state: ProgressState) {
self.paused = false;
self.pending_snapshot = 0;
self.state = state;
self.ins.reset();
}
pub(crate) fn reset(&mut self, next_idx: u64) {
self.matched = 0;
self.next_idx = next_idx;
self.state = ProgressState::default();
self.paused = false;
self.pending_snapshot = 0;
self.pending_request_snapshot = INVALID_INDEX;
self.recent_active = false;
debug_assert!(self.ins.cap() != 0);
self.ins.reset();
}
pub fn become_probe(&mut self) {
if self.state == ProgressState::Snapshot {
let pending_snapshot = self.pending_snapshot;
self.reset_state(ProgressState::Probe);
self.next_idx = cmp::max(self.matched + 1, pending_snapshot + 1);
} else {
self.reset_state(ProgressState::Probe);
self.next_idx = self.matched + 1;
}
}
#[inline]
pub fn become_replicate(&mut self) {
self.reset_state(ProgressState::Replicate);
self.next_idx = self.matched + 1;
}
#[inline]
pub fn become_snapshot(&mut self, snapshot_idx: u64) {
self.reset_state(ProgressState::Snapshot);
self.pending_snapshot = snapshot_idx;
}
#[inline]
pub fn snapshot_failure(&mut self) {
self.pending_snapshot = 0;
}
#[inline]
pub fn maybe_snapshot_abort(&self) -> bool {
self.state == ProgressState::Snapshot && self.matched >= self.pending_snapshot
}
pub fn maybe_update(&mut self, n: u64) -> bool {
let need_update = self.matched < n;
if need_update {
self.matched = n;
self.resume();
};
if self.next_idx < n + 1 {
self.next_idx = n + 1
}
need_update
}
#[inline]
pub fn optimistic_update(&mut self, n: u64) {
self.next_idx = n + 1;
}
pub fn maybe_decr_to(&mut self, rejected: u64, last: u64, request_snapshot: u64) -> bool {
if self.state == ProgressState::Replicate {
if rejected < self.matched
|| (rejected == self.matched && request_snapshot == INVALID_INDEX)
{
return false;
}
if request_snapshot == INVALID_INDEX {
self.next_idx = self.matched + 1;
} else {
self.pending_request_snapshot = request_snapshot;
}
return true;
}
if (self.next_idx == 0 || self.next_idx - 1 != rejected)
&& request_snapshot == INVALID_INDEX
{
return false;
}
if request_snapshot == INVALID_INDEX {
self.next_idx = cmp::min(rejected, last + 1);
if self.next_idx < 1 {
self.next_idx = 1;
}
} else if self.pending_request_snapshot == INVALID_INDEX {
self.pending_request_snapshot = request_snapshot;
}
self.resume();
true
}
#[inline]
pub fn is_paused(&self) -> bool {
match self.state {
ProgressState::Probe => self.paused,
ProgressState::Replicate => self.ins.full(),
ProgressState::Snapshot => true,
}
}
#[inline]
pub fn resume(&mut self) {
self.paused = false;
}
#[inline]
pub fn pause(&mut self) {
self.paused = true;
}
pub fn update_state(&mut self, last: u64) {
match self.state {
ProgressState::Replicate => {
self.optimistic_update(last);
self.ins.add(last);
}
ProgressState::Probe => self.pause(),
ProgressState::Snapshot => panic!(
"updating progress state in unhandled state {:?}",
self.state
),
}
}
}