pgwire_replication/client/
metrics.rs1use std::sync::atomic::{AtomicU64, Ordering};
4
5use crate::lsn::Lsn;
6
7#[derive(Debug, Default)]
16pub struct ReplicationMetrics {
17 events_forwarded: AtomicU64,
18 feedback_sent: AtomicU64,
19 keepalive_replies: AtomicU64,
20 stall_count: AtomicU64,
21 stall_micros_total: AtomicU64,
22 last_applied_lsn: AtomicU64,
23 last_wal_end: AtomicU64,
24}
25
26impl ReplicationMetrics {
27 #[inline]
29 pub fn events_forwarded(&self) -> u64 {
30 self.events_forwarded.load(Ordering::Relaxed)
31 }
32
33 #[inline]
35 pub fn feedback_sent(&self) -> u64 {
36 self.feedback_sent.load(Ordering::Relaxed)
37 }
38
39 #[inline]
42 pub fn keepalive_replies(&self) -> u64 {
43 self.keepalive_replies.load(Ordering::Relaxed)
44 }
45
46 #[inline]
49 pub fn stall_count(&self) -> u64 {
50 self.stall_count.load(Ordering::Relaxed)
51 }
52
53 #[inline]
56 pub fn stall_micros_total(&self) -> u64 {
57 self.stall_micros_total.load(Ordering::Relaxed)
58 }
59
60 #[inline]
62 pub fn last_applied_lsn(&self) -> Lsn {
63 Lsn::from_u64(self.last_applied_lsn.load(Ordering::Relaxed))
64 }
65
66 #[inline]
68 pub fn last_wal_end(&self) -> Lsn {
69 Lsn::from_u64(self.last_wal_end.load(Ordering::Relaxed))
70 }
71
72 #[inline]
75 pub(crate) fn record_event_forwarded(&self) {
76 self.events_forwarded.fetch_add(1, Ordering::Relaxed);
77 }
78
79 #[inline]
80 pub(crate) fn record_feedback_sent(&self, applied: Lsn) {
81 self.feedback_sent.fetch_add(1, Ordering::Relaxed);
82 self.last_applied_lsn
83 .store(applied.as_u64(), Ordering::Relaxed);
84 }
85
86 #[inline]
87 pub(crate) fn record_keepalive_reply(&self) {
88 self.keepalive_replies.fetch_add(1, Ordering::Relaxed);
89 }
90
91 #[inline]
95 pub(crate) fn record_stall_begin(&self) {
96 self.stall_count.fetch_add(1, Ordering::Relaxed);
97 }
98
99 #[inline]
101 pub(crate) fn record_stall_end(&self, micros: u64) {
102 self.stall_micros_total.fetch_add(micros, Ordering::Relaxed);
103 }
104
105 #[inline]
106 pub(crate) fn set_last_wal_end(&self, wal_end: Lsn) {
107 self.last_wal_end.store(wal_end.as_u64(), Ordering::Relaxed);
108 }
109}
110
111#[cfg(test)]
112mod tests {
113 use super::*;
114
115 #[test]
116 fn counters_start_at_zero() {
117 let m = ReplicationMetrics::default();
118 assert_eq!(m.events_forwarded(), 0);
119 assert_eq!(m.feedback_sent(), 0);
120 assert_eq!(m.keepalive_replies(), 0);
121 assert_eq!(m.stall_count(), 0);
122 assert_eq!(m.stall_micros_total(), 0);
123 assert_eq!(m.last_applied_lsn(), Lsn::ZERO);
124 assert_eq!(m.last_wal_end(), Lsn::ZERO);
125 }
126
127 #[test]
128 fn recorders_update_their_counters() {
129 let m = ReplicationMetrics::default();
130
131 m.record_event_forwarded();
132 m.record_event_forwarded();
133 assert_eq!(m.events_forwarded(), 2);
134
135 m.record_feedback_sent(Lsn::from_u64(42));
136 assert_eq!(m.feedback_sent(), 1);
137 assert_eq!(m.last_applied_lsn(), Lsn::from_u64(42));
138
139 m.record_keepalive_reply();
140 assert_eq!(m.keepalive_replies(), 1);
141
142 m.record_stall_begin();
143 assert_eq!(m.stall_count(), 1);
144 assert_eq!(m.stall_micros_total(), 0);
145
146 m.record_stall_end(100);
147 m.record_stall_end(50);
148 assert_eq!(m.stall_micros_total(), 150);
149 assert_eq!(m.stall_count(), 1);
151
152 m.set_last_wal_end(Lsn::from_u64(7));
153 assert_eq!(m.last_wal_end(), Lsn::from_u64(7));
154 }
155}