use std::sync::atomic::{AtomicU64, AtomicUsize, Ordering};
use std::sync::{Arc, Mutex};
#[derive(Debug, Default)]
pub(crate) struct ReplicaProgress {
loading: AtomicUsize,
control_epoch: Arc<AtomicU64>,
primary_offset: AtomicU64,
applied_offset: AtomicU64,
last_ping_ms: AtomicU64,
runner_offsets: Mutex<Box<[AtomicU64]>>,
upstream_gens: Mutex<Box<[u64]>>,
}
impl ReplicaProgress {
pub(super) fn with_epoch(control_epoch: Arc<AtomicU64>) -> Self {
Self { control_epoch, ..Self::default() }
}
pub(crate) fn begin_loading(&self) {
if self.loading.fetch_add(1, Ordering::AcqRel) == 0 {
self.control_epoch.fetch_add(1, Ordering::Release);
}
}
pub(crate) fn end_loading(&self) {
if self.loading.fetch_sub(1, Ordering::AcqRel) == 1 {
self.control_epoch.fetch_add(1, Ordering::Release);
}
}
pub(crate) fn loading(&self) -> bool {
self.loading.load(Ordering::Acquire) > 0
}
pub(crate) fn record_ping(
&self,
runner_slot: usize,
generation: u64,
primary_offset: u64,
applied: u64,
) {
self.primary_offset.fetch_max(primary_offset, Ordering::Relaxed);
self.applied_offset.fetch_max(applied, Ordering::Relaxed);
self.last_ping_ms.store(epoch_ms(), Ordering::Relaxed);
self.record_applied(runner_slot, applied);
let mut gens = self.upstream_gens.lock().expect("upstream_gens poisoned");
if let Some(slot) = gens.get_mut(runner_slot) {
*slot = generation;
}
}
pub(crate) fn record_applied(&self, runner_slot: usize, applied: u64) {
let slots = self.runner_offsets.lock().expect("runner_offsets poisoned");
if let Some(slot) = slots.get(runner_slot) {
slot.store(applied, Ordering::Relaxed);
}
}
pub(super) fn size_runner_slots(&self, n: usize) {
*self.runner_offsets.lock().expect("runner_offsets poisoned") =
(0..n).map(|_| AtomicU64::new(0)).collect();
*self.upstream_gens.lock().expect("upstream_gens poisoned") = vec![0; n].into_boxed_slice();
}
pub(super) fn clear_runner_slots(&self) {
self.size_runner_slots(0);
}
pub(super) fn runner_offsets(&self) -> Vec<u64> {
self.runner_offsets
.lock()
.expect("runner_offsets poisoned")
.iter()
.map(|a| a.load(Ordering::Relaxed))
.collect()
}
pub(super) fn upstream_gens(&self) -> Vec<u64> {
self.upstream_gens.lock().expect("upstream_gens poisoned").to_vec()
}
pub(super) fn applied_offset_sum(&self) -> u64 {
self.runner_offsets
.lock()
.expect("runner_offsets poisoned")
.iter()
.fold(0u64, |acc, v| acc.saturating_add(v.load(Ordering::Relaxed)))
}
pub(super) fn last_ping_ms(&self) -> u64 {
self.last_ping_ms.load(Ordering::Relaxed)
}
pub(super) fn link_view(&self) -> (bool, u64, u64, u64) {
let last = self.last_ping_ms();
let age_ms = epoch_ms().saturating_sub(last);
let up = last != 0 && age_ms < 3_000;
let primary = self.primary_offset.load(Ordering::Relaxed);
let applied = self.applied_offset.load(Ordering::Relaxed);
(up, applied, primary.saturating_sub(applied), age_ms / 1000)
}
}
pub(super) fn epoch_ms() -> u64 {
std::time::SystemTime::now()
.duration_since(std::time::UNIX_EPOCH)
.map(|d| d.as_millis() as u64)
.unwrap_or(0)
}