use std::sync::atomic::{AtomicU64, Ordering};
use crate::lsn::Lsn;
#[derive(Debug, Default)]
pub struct ReplicationMetrics {
events_forwarded: AtomicU64,
feedback_sent: AtomicU64,
keepalive_replies: AtomicU64,
stall_count: AtomicU64,
stall_micros_total: AtomicU64,
last_applied_lsn: AtomicU64,
last_wal_end: AtomicU64,
}
impl ReplicationMetrics {
#[inline]
pub fn events_forwarded(&self) -> u64 {
self.events_forwarded.load(Ordering::Relaxed)
}
#[inline]
pub fn feedback_sent(&self) -> u64 {
self.feedback_sent.load(Ordering::Relaxed)
}
#[inline]
pub fn keepalive_replies(&self) -> u64 {
self.keepalive_replies.load(Ordering::Relaxed)
}
#[inline]
pub fn stall_count(&self) -> u64 {
self.stall_count.load(Ordering::Relaxed)
}
#[inline]
pub fn stall_micros_total(&self) -> u64 {
self.stall_micros_total.load(Ordering::Relaxed)
}
#[inline]
pub fn last_applied_lsn(&self) -> Lsn {
Lsn::from_u64(self.last_applied_lsn.load(Ordering::Relaxed))
}
#[inline]
pub fn last_wal_end(&self) -> Lsn {
Lsn::from_u64(self.last_wal_end.load(Ordering::Relaxed))
}
#[inline]
pub(crate) fn record_event_forwarded(&self) {
self.events_forwarded.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub(crate) fn record_feedback_sent(&self, applied: Lsn) {
self.feedback_sent.fetch_add(1, Ordering::Relaxed);
self.last_applied_lsn
.store(applied.as_u64(), Ordering::Relaxed);
}
#[inline]
pub(crate) fn record_keepalive_reply(&self) {
self.keepalive_replies.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub(crate) fn record_stall_begin(&self) {
self.stall_count.fetch_add(1, Ordering::Relaxed);
}
#[inline]
pub(crate) fn record_stall_end(&self, micros: u64) {
self.stall_micros_total.fetch_add(micros, Ordering::Relaxed);
}
#[inline]
pub(crate) fn set_last_wal_end(&self, wal_end: Lsn) {
self.last_wal_end.store(wal_end.as_u64(), Ordering::Relaxed);
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn counters_start_at_zero() {
let m = ReplicationMetrics::default();
assert_eq!(m.events_forwarded(), 0);
assert_eq!(m.feedback_sent(), 0);
assert_eq!(m.keepalive_replies(), 0);
assert_eq!(m.stall_count(), 0);
assert_eq!(m.stall_micros_total(), 0);
assert_eq!(m.last_applied_lsn(), Lsn::ZERO);
assert_eq!(m.last_wal_end(), Lsn::ZERO);
}
#[test]
fn recorders_update_their_counters() {
let m = ReplicationMetrics::default();
m.record_event_forwarded();
m.record_event_forwarded();
assert_eq!(m.events_forwarded(), 2);
m.record_feedback_sent(Lsn::from_u64(42));
assert_eq!(m.feedback_sent(), 1);
assert_eq!(m.last_applied_lsn(), Lsn::from_u64(42));
m.record_keepalive_reply();
assert_eq!(m.keepalive_replies(), 1);
m.record_stall_begin();
assert_eq!(m.stall_count(), 1);
assert_eq!(m.stall_micros_total(), 0);
m.record_stall_end(100);
m.record_stall_end(50);
assert_eq!(m.stall_micros_total(), 150);
assert_eq!(m.stall_count(), 1);
m.set_last_wal_end(Lsn::from_u64(7));
assert_eq!(m.last_wal_end(), Lsn::from_u64(7));
}
}