use std::sync::atomic::{AtomicU64, Ordering};
#[derive(Debug, Default)]
pub struct UdpMetrics {
pub associations_active: AtomicU64,
pub associations_total: AtomicU64,
pub association_failures: AtomicU64,
pub association_timeouts: AtomicU64,
pub packets_up: AtomicU64,
pub packets_down: AtomicU64,
pub bytes_up: AtomicU64,
pub bytes_down: AtomicU64,
pub dropped_packets: AtomicU64,
pub dropped_encode_errors: AtomicU64,
pub dropped_send_errors: AtomicU64,
pub dropped_response_channel_full: AtomicU64,
pub target_flows_active: AtomicU64,
pub target_flows_total: AtomicU64,
pub decode_errors: AtomicU64,
pub upstream_associations_total: AtomicU64,
pub upstream_associations_active: AtomicU64,
pub upstream_packets_up: AtomicU64,
pub upstream_packets_down: AtomicU64,
pub upstream_bytes_up: AtomicU64,
pub upstream_bytes_down: AtomicU64,
pub upstream_failures: AtomicU64,
pub unsupported_upstream_total: AtomicU64,
pub standalone_flows_active: AtomicU64,
pub standalone_flows_total: AtomicU64,
pub standalone_packets_in: AtomicU64,
pub standalone_packets_out: AtomicU64,
pub standalone_bytes_in: AtomicU64,
pub standalone_bytes_out: AtomicU64,
pub standalone_malformed_datagrams: AtomicU64,
pub standalone_rejected_datagrams: AtomicU64,
pub standalone_flow_reaps: AtomicU64,
}
impl UdpMetrics {
pub fn new() -> Self {
Self::default()
}
pub fn record_association_created(&self) {
self.associations_active.fetch_add(1, Ordering::Relaxed);
self.associations_total.fetch_add(1, Ordering::Relaxed);
}
pub fn record_association_closed(&self) {
Self::decr_active(&self.associations_active);
}
pub fn record_association_failure(&self) {
self.association_failures.fetch_add(1, Ordering::Relaxed);
}
pub fn record_packet_up(&self, bytes: u64) {
self.packets_up.fetch_add(1, Ordering::Relaxed);
self.bytes_up.fetch_add(bytes, Ordering::Relaxed);
}
pub fn record_packet_down(&self, bytes: u64) {
self.packets_down.fetch_add(1, Ordering::Relaxed);
self.bytes_down.fetch_add(bytes, Ordering::Relaxed);
}
pub fn record_dropped(&self) {
self.dropped_packets.fetch_add(1, Ordering::Relaxed);
}
pub fn record_dropped_encode_error(&self) {
self.dropped_encode_errors.fetch_add(1, Ordering::Relaxed);
self.record_dropped();
}
pub fn record_dropped_send_error(&self) {
self.dropped_send_errors.fetch_add(1, Ordering::Relaxed);
self.record_dropped();
}
pub fn record_dropped_response_channel_full(&self) {
self.dropped_response_channel_full
.fetch_add(1, Ordering::Relaxed);
self.record_dropped();
}
pub fn record_target_flow_created(&self) {
self.target_flows_active.fetch_add(1, Ordering::Relaxed);
self.target_flows_total.fetch_add(1, Ordering::Relaxed);
}
pub fn record_target_flow_closed(&self) {
Self::decr_active(&self.target_flows_active);
}
pub fn record_decode_error(&self) {
self.decode_errors.fetch_add(1, Ordering::Relaxed);
}
pub fn record_association_timeout(&self) {
self.association_timeouts.fetch_add(1, Ordering::Relaxed);
}
pub fn record_target_flow_timeout(&self) {
Self::decr_active(&self.target_flows_active);
}
pub fn record_upstream_association_created(&self) {
self.upstream_associations_active
.fetch_add(1, Ordering::Relaxed);
self.upstream_associations_total
.fetch_add(1, Ordering::Relaxed);
}
pub fn record_upstream_association_closed(&self) {
Self::decr_active(&self.upstream_associations_active);
}
pub fn record_upstream_failure(&self) {
self.upstream_failures.fetch_add(1, Ordering::Relaxed);
}
pub fn record_upstream_packet_up(&self, bytes: u64) {
self.upstream_packets_up.fetch_add(1, Ordering::Relaxed);
self.upstream_bytes_up.fetch_add(bytes, Ordering::Relaxed);
}
pub fn record_upstream_packet_down(&self, bytes: u64) {
self.upstream_packets_down.fetch_add(1, Ordering::Relaxed);
self.upstream_bytes_down.fetch_add(bytes, Ordering::Relaxed);
}
pub fn record_unsupported_upstream(&self) {
self.unsupported_upstream_total
.fetch_add(1, Ordering::Relaxed);
}
pub fn record_standalone_flow_created(&self) {
self.standalone_flows_active.fetch_add(1, Ordering::Relaxed);
self.standalone_flows_total.fetch_add(1, Ordering::Relaxed);
}
pub fn record_standalone_flow_closed(&self) {
Self::decr_active(&self.standalone_flows_active);
}
pub fn record_standalone_packet_in(&self, bytes: u64) {
self.standalone_packets_in.fetch_add(1, Ordering::Relaxed);
self.standalone_bytes_in.fetch_add(bytes, Ordering::Relaxed);
}
pub fn record_standalone_packet_out(&self, bytes: u64) {
self.standalone_packets_out.fetch_add(1, Ordering::Relaxed);
self.standalone_bytes_out
.fetch_add(bytes, Ordering::Relaxed);
}
pub fn record_standalone_malformed(&self) {
self.standalone_malformed_datagrams
.fetch_add(1, Ordering::Relaxed);
}
pub fn record_standalone_rejected(&self) {
self.standalone_rejected_datagrams
.fetch_add(1, Ordering::Relaxed);
}
pub fn record_standalone_flow_reap(&self) {
Self::decr_active(&self.standalone_flows_active);
self.standalone_flow_reaps.fetch_add(1, Ordering::Relaxed);
}
fn decr_active(counter: &AtomicU64) {
let _ = counter.fetch_update(Ordering::Relaxed, Ordering::Relaxed, |v| v.checked_sub(1));
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn default_metrics_are_zero() {
let metrics = UdpMetrics::new();
assert_eq!(metrics.associations_active.load(Ordering::Relaxed), 0);
assert_eq!(metrics.associations_total.load(Ordering::Relaxed), 0);
assert_eq!(metrics.packets_up.load(Ordering::Relaxed), 0);
assert_eq!(metrics.packets_down.load(Ordering::Relaxed), 0);
assert_eq!(metrics.bytes_up.load(Ordering::Relaxed), 0);
assert_eq!(metrics.bytes_down.load(Ordering::Relaxed), 0);
assert_eq!(metrics.dropped_packets.load(Ordering::Relaxed), 0);
assert_eq!(metrics.target_flows_active.load(Ordering::Relaxed), 0);
assert_eq!(metrics.target_flows_total.load(Ordering::Relaxed), 0);
assert_eq!(metrics.decode_errors.load(Ordering::Relaxed), 0);
assert_eq!(
metrics.upstream_associations_total.load(Ordering::Relaxed),
0
);
assert_eq!(
metrics.upstream_associations_active.load(Ordering::Relaxed),
0
);
assert_eq!(metrics.upstream_packets_up.load(Ordering::Relaxed), 0);
assert_eq!(metrics.upstream_packets_down.load(Ordering::Relaxed), 0);
assert_eq!(metrics.upstream_bytes_up.load(Ordering::Relaxed), 0);
assert_eq!(metrics.upstream_bytes_down.load(Ordering::Relaxed), 0);
assert_eq!(metrics.upstream_failures.load(Ordering::Relaxed), 0);
assert_eq!(
metrics.unsupported_upstream_total.load(Ordering::Relaxed),
0
);
}
#[test]
fn association_metrics() {
let metrics = UdpMetrics::new();
metrics.record_association_created();
assert_eq!(metrics.associations_active.load(Ordering::Relaxed), 1);
assert_eq!(metrics.associations_total.load(Ordering::Relaxed), 1);
metrics.record_association_created();
assert_eq!(metrics.associations_active.load(Ordering::Relaxed), 2);
assert_eq!(metrics.associations_total.load(Ordering::Relaxed), 2);
metrics.record_association_closed();
assert_eq!(metrics.associations_active.load(Ordering::Relaxed), 1);
assert_eq!(metrics.associations_total.load(Ordering::Relaxed), 2);
}
#[test]
fn association_failure() {
let metrics = UdpMetrics::new();
metrics.record_association_failure();
metrics.record_association_failure();
assert_eq!(metrics.association_failures.load(Ordering::Relaxed), 2);
}
#[test]
fn packet_metrics() {
let metrics = UdpMetrics::new();
metrics.record_packet_up(100);
metrics.record_packet_up(200);
assert_eq!(metrics.packets_up.load(Ordering::Relaxed), 2);
assert_eq!(metrics.bytes_up.load(Ordering::Relaxed), 300);
metrics.record_packet_down(50);
assert_eq!(metrics.packets_down.load(Ordering::Relaxed), 1);
assert_eq!(metrics.bytes_down.load(Ordering::Relaxed), 50);
}
#[test]
fn dropped_packets() {
let metrics = UdpMetrics::new();
metrics.record_dropped();
metrics.record_dropped();
metrics.record_dropped();
assert_eq!(metrics.dropped_packets.load(Ordering::Relaxed), 3);
}
#[test]
fn response_channel_overflow_counts_as_drop() {
let metrics = UdpMetrics::new();
metrics.record_dropped_response_channel_full();
assert_eq!(
metrics
.dropped_response_channel_full
.load(Ordering::Relaxed),
1
);
assert_eq!(metrics.dropped_packets.load(Ordering::Relaxed), 1);
}
#[test]
fn target_flow_metrics() {
let metrics = UdpMetrics::new();
metrics.record_target_flow_created();
metrics.record_target_flow_created();
assert_eq!(metrics.target_flows_active.load(Ordering::Relaxed), 2);
assert_eq!(metrics.target_flows_total.load(Ordering::Relaxed), 2);
metrics.record_target_flow_closed();
assert_eq!(metrics.target_flows_active.load(Ordering::Relaxed), 1);
assert_eq!(metrics.target_flows_total.load(Ordering::Relaxed), 2);
}
#[test]
fn decode_errors() {
let metrics = UdpMetrics::new();
metrics.record_decode_error();
assert_eq!(metrics.decode_errors.load(Ordering::Relaxed), 1);
}
#[test]
fn association_timeout_metric() {
let metrics = UdpMetrics::new();
metrics.record_association_created();
assert_eq!(metrics.associations_active.load(Ordering::Relaxed), 1);
metrics.record_association_timeout();
assert_eq!(metrics.association_timeouts.load(Ordering::Relaxed), 1);
assert_eq!(metrics.associations_active.load(Ordering::Relaxed), 1);
}
#[test]
fn target_flow_timeout_metric() {
let metrics = UdpMetrics::new();
metrics.record_target_flow_created();
assert_eq!(metrics.target_flows_active.load(Ordering::Relaxed), 1);
metrics.record_target_flow_timeout();
assert_eq!(metrics.target_flows_active.load(Ordering::Relaxed), 0);
}
#[test]
fn upstream_association_metrics() {
let metrics = UdpMetrics::new();
metrics.record_upstream_association_created();
assert_eq!(
metrics.upstream_associations_active.load(Ordering::Relaxed),
1
);
assert_eq!(
metrics.upstream_associations_total.load(Ordering::Relaxed),
1
);
metrics.record_upstream_association_created();
assert_eq!(
metrics.upstream_associations_active.load(Ordering::Relaxed),
2
);
assert_eq!(
metrics.upstream_associations_total.load(Ordering::Relaxed),
2
);
metrics.record_upstream_association_closed();
assert_eq!(
metrics.upstream_associations_active.load(Ordering::Relaxed),
1
);
assert_eq!(
metrics.upstream_associations_total.load(Ordering::Relaxed),
2
);
}
#[test]
fn upstream_failure_metric() {
let metrics = UdpMetrics::new();
metrics.record_upstream_failure();
metrics.record_upstream_failure();
assert_eq!(metrics.upstream_failures.load(Ordering::Relaxed), 2);
}
#[test]
fn upstream_packet_metrics() {
let metrics = UdpMetrics::new();
metrics.record_upstream_packet_up(100);
metrics.record_upstream_packet_up(200);
assert_eq!(metrics.upstream_packets_up.load(Ordering::Relaxed), 2);
assert_eq!(metrics.upstream_bytes_up.load(Ordering::Relaxed), 300);
metrics.record_upstream_packet_down(50);
assert_eq!(metrics.upstream_packets_down.load(Ordering::Relaxed), 1);
assert_eq!(metrics.upstream_bytes_down.load(Ordering::Relaxed), 50);
}
#[test]
fn unsupported_upstream_metric() {
let metrics = UdpMetrics::new();
metrics.record_unsupported_upstream();
metrics.record_unsupported_upstream();
assert_eq!(
metrics.unsupported_upstream_total.load(Ordering::Relaxed),
2
);
}
#[test]
fn default_standalone_metrics_are_zero() {
let metrics = UdpMetrics::new();
assert_eq!(metrics.standalone_flows_active.load(Ordering::Relaxed), 0);
assert_eq!(metrics.standalone_flows_total.load(Ordering::Relaxed), 0);
assert_eq!(metrics.standalone_packets_in.load(Ordering::Relaxed), 0);
assert_eq!(metrics.standalone_packets_out.load(Ordering::Relaxed), 0);
assert_eq!(metrics.standalone_bytes_in.load(Ordering::Relaxed), 0);
assert_eq!(metrics.standalone_bytes_out.load(Ordering::Relaxed), 0);
assert_eq!(
metrics
.standalone_malformed_datagrams
.load(Ordering::Relaxed),
0
);
assert_eq!(
metrics
.standalone_rejected_datagrams
.load(Ordering::Relaxed),
0
);
assert_eq!(metrics.standalone_flow_reaps.load(Ordering::Relaxed), 0);
}
#[test]
fn standalone_flow_metrics() {
let metrics = UdpMetrics::new();
metrics.record_standalone_flow_created();
assert_eq!(metrics.standalone_flows_active.load(Ordering::Relaxed), 1);
assert_eq!(metrics.standalone_flows_total.load(Ordering::Relaxed), 1);
metrics.record_standalone_flow_created();
assert_eq!(metrics.standalone_flows_active.load(Ordering::Relaxed), 2);
assert_eq!(metrics.standalone_flows_total.load(Ordering::Relaxed), 2);
metrics.record_standalone_flow_closed();
assert_eq!(metrics.standalone_flows_active.load(Ordering::Relaxed), 1);
assert_eq!(metrics.standalone_flows_total.load(Ordering::Relaxed), 2);
}
#[test]
fn standalone_packet_metrics() {
let metrics = UdpMetrics::new();
metrics.record_standalone_packet_in(100);
metrics.record_standalone_packet_in(200);
assert_eq!(metrics.standalone_packets_in.load(Ordering::Relaxed), 2);
assert_eq!(metrics.standalone_bytes_in.load(Ordering::Relaxed), 300);
metrics.record_standalone_packet_out(50);
assert_eq!(metrics.standalone_packets_out.load(Ordering::Relaxed), 1);
assert_eq!(metrics.standalone_bytes_out.load(Ordering::Relaxed), 50);
}
#[test]
fn standalone_malformed_metric() {
let metrics = UdpMetrics::new();
metrics.record_standalone_malformed();
metrics.record_standalone_malformed();
assert_eq!(
metrics
.standalone_malformed_datagrams
.load(Ordering::Relaxed),
2
);
}
#[test]
fn standalone_rejected_metric() {
let metrics = UdpMetrics::new();
metrics.record_standalone_rejected();
assert_eq!(
metrics
.standalone_rejected_datagrams
.load(Ordering::Relaxed),
1
);
}
#[test]
fn standalone_flow_reap_metric() {
let metrics = UdpMetrics::new();
metrics.record_standalone_flow_created();
metrics.record_standalone_flow_created();
assert_eq!(metrics.standalone_flows_active.load(Ordering::Relaxed), 2);
metrics.record_standalone_flow_reap();
assert_eq!(metrics.standalone_flows_active.load(Ordering::Relaxed), 1);
assert_eq!(metrics.standalone_flow_reaps.load(Ordering::Relaxed), 1);
}
#[test]
fn active_gauges_do_not_underflow() {
let metrics = UdpMetrics::new();
metrics.record_association_closed();
metrics.record_target_flow_closed();
metrics.record_target_flow_timeout();
metrics.record_upstream_association_closed();
metrics.record_standalone_flow_closed();
metrics.record_standalone_flow_reap();
assert_eq!(metrics.associations_active.load(Ordering::Relaxed), 0);
assert_eq!(metrics.target_flows_active.load(Ordering::Relaxed), 0);
assert_eq!(
metrics.upstream_associations_active.load(Ordering::Relaxed),
0
);
assert_eq!(metrics.standalone_flows_active.load(Ordering::Relaxed), 0);
}
}