use std::sync::atomic::{AtomicU64, Ordering};
use sipx_sip::Method;
use sipx_sip::transaction::Timer;
use crate::target::TransportKind;
const TRANSPORTS: usize = 6;
const fn slot(transport: TransportKind) -> usize {
match transport {
TransportKind::Udp => 0,
TransportKind::Tcp => 1,
TransportKind::Tls => 2,
TransportKind::Ws => 3,
TransportKind::Wss => 4,
TransportKind::Quic => 5,
}
}
#[derive(Debug, Default)]
pub(crate) struct Shed {
pub(crate) requests: AtomicU64,
pub(crate) acks: AtomicU64,
pub(crate) unmatched: AtomicU64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct ShedCounts {
pub requests: u64,
pub acks: u64,
pub unmatched: u64,
}
impl ShedCounts {
#[must_use]
pub fn any(self) -> bool {
self.total() > 0
}
#[must_use]
pub fn total(self) -> u64 {
self.requests
.saturating_add(self.acks)
.saturating_add(self.unmatched)
}
}
#[derive(Debug, Default)]
struct TransportMeter {
requests_in: AtomicU64,
requests_out: AtomicU64,
responses_in: AtomicU64,
responses_out: AtomicU64,
parse_failures: AtomicU64,
source_refusals: AtomicU64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct TransportCounts {
pub requests_in: u64,
pub requests_out: u64,
pub responses_in: u64,
pub responses_out: u64,
pub parse_failures: u64,
pub source_refusals: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct TimeoutCounts {
pub b: u64,
pub f: u64,
pub h: u64,
}
impl TimeoutCounts {
#[must_use]
pub fn total(self) -> u64 {
self.b.saturating_add(self.f).saturating_add(self.h)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct CaptureCounts {
pub records: u64,
pub dropped: u64,
pub errors: u64,
pub hep_records: u64,
pub hep_dropped: u64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct DiscardCounts {
pub transaction_events: u64,
pub unanswered: u64,
pub no_destination: u64,
pub send_failures: u64,
pub stun_unmatched: u64,
}
impl DiscardCounts {
#[must_use]
pub fn total(self) -> u64 {
self.transaction_events
.saturating_add(self.unanswered)
.saturating_add(self.no_destination)
.saturating_add(self.send_failures)
.saturating_add(self.stun_unmatched)
}
}
#[derive(Debug, Default)]
struct Unsent {
invite: AtomicU64,
ack: AtomicU64,
bye: AtomicU64,
cancel: AtomicU64,
other: AtomicU64,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct UnsentCounts {
pub invite: u64,
pub ack: u64,
pub bye: u64,
pub cancel: u64,
pub other: u64,
}
impl UnsentCounts {
#[must_use]
pub fn total(self) -> u64 {
self.invite
.saturating_add(self.ack)
.saturating_add(self.bye)
.saturating_add(self.cancel)
.saturating_add(self.other)
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default)]
pub struct Counters {
pub shed: ShedCounts,
pub overload_rejections: u64,
pub unmatched_responses: u64,
pub retransmissions_sent: u64,
pub oversized_request_tcp_fallbacks: u64,
pub timeouts: TimeoutCounts,
pub discards: DiscardCounts,
pub unsent: UnsentCounts,
pub capture: CaptureCounts,
pub observation_dropped: u64,
per_transport: [TransportCounts; TRANSPORTS],
}
impl Counters {
#[must_use]
pub fn transport(&self, transport: TransportKind) -> TransportCounts {
self.per_transport
.get(slot(transport))
.copied()
.unwrap_or_default()
}
#[must_use]
pub fn messages_in(&self) -> u64 {
self.per_transport.iter().fold(0, |total, counts| {
total
.saturating_add(counts.requests_in)
.saturating_add(counts.responses_in)
})
}
#[must_use]
pub fn messages_out(&self) -> u64 {
self.per_transport.iter().fold(0, |total, counts| {
total
.saturating_add(counts.requests_out)
.saturating_add(counts.responses_out)
})
}
#[must_use]
pub fn parse_failures(&self) -> u64 {
self.per_transport.iter().fold(0, |total, counts| {
total.saturating_add(counts.parse_failures)
})
}
#[must_use]
pub fn any_loss(&self) -> bool {
self.shed.any()
|| self.overload_rejections > 0
|| self.discards.total() > 0
|| self.capture.dropped > 0
|| self.capture.hep_dropped > 0
|| self.observation_dropped > 0
|| self.unsent.total() > 0
}
}
#[derive(Debug, Default)]
pub(crate) struct Meters {
pub(crate) shed: Shed,
overload_rejections: AtomicU64,
per_transport: [TransportMeter; TRANSPORTS],
unmatched_responses: AtomicU64,
retransmissions: AtomicU64,
oversized_request_tcp_fallbacks: AtomicU64,
timeout_b: AtomicU64,
timeout_f: AtomicU64,
timeout_h: AtomicU64,
discard_transaction_events: AtomicU64,
discard_unanswered: AtomicU64,
discard_no_destination: AtomicU64,
discard_send_failures: AtomicU64,
discard_stun_unmatched: AtomicU64,
capture_records: AtomicU64,
capture_dropped: AtomicU64,
capture_errors: AtomicU64,
capture_hep_records: AtomicU64,
capture_hep_dropped: AtomicU64,
observation_dropped: AtomicU64,
unsent: Unsent,
}
fn bump(counter: &AtomicU64) {
counter.fetch_add(1, Ordering::Relaxed);
}
fn read(counter: &AtomicU64) -> u64 {
counter.load(Ordering::Relaxed)
}
impl Meters {
pub(crate) fn overload_rejection(&self) {
bump(&self.overload_rejections);
}
fn meter(&self, transport: TransportKind) -> Option<&TransportMeter> {
self.per_transport.get(slot(transport))
}
pub(crate) fn message_in(&self, transport: TransportKind, is_response: bool) {
if let Some(meter) = self.meter(transport) {
if is_response {
bump(&meter.responses_in);
} else {
bump(&meter.requests_in);
}
}
}
pub(crate) fn message_out(&self, transport: TransportKind, is_response: bool) {
if let Some(meter) = self.meter(transport) {
if is_response {
bump(&meter.responses_out);
} else {
bump(&meter.requests_out);
}
}
}
pub(crate) fn parse_failure(&self, transport: TransportKind) {
if let Some(meter) = self.meter(transport) {
bump(&meter.parse_failures);
}
}
pub(crate) fn source_refusal(&self, transport: TransportKind) {
if let Some(meter) = self.meter(transport) {
bump(&meter.source_refusals);
}
}
pub(crate) fn observation_drop(&self) {
bump(&self.observation_dropped);
}
pub(crate) fn unmatched_response(&self) {
bump(&self.unmatched_responses);
}
pub(crate) fn on_timer(&self, timer: Timer) {
match timer {
Timer::A | Timer::E | Timer::G => bump(&self.retransmissions),
Timer::B => bump(&self.timeout_b),
Timer::F => bump(&self.timeout_f),
Timer::H => bump(&self.timeout_h),
Timer::D | Timer::I | Timer::J | Timer::K | Timer::L | Timer::M | Timer::Trying100 => {}
}
}
pub(crate) fn oversized_request_tcp_fallback(&self) {
bump(&self.oversized_request_tcp_fallbacks);
}
pub(crate) fn discard_transaction_event(&self) {
bump(&self.discard_transaction_events);
}
pub(crate) fn discard_unanswered(&self) {
bump(&self.discard_unanswered);
}
pub(crate) fn discard_no_destination(&self) {
bump(&self.discard_no_destination);
}
pub(crate) fn discard_send_failure(&self) {
bump(&self.discard_send_failures);
}
pub(crate) fn discard_stun_unmatched(&self) {
bump(&self.discard_stun_unmatched);
}
pub(crate) fn capture_record(&self) {
bump(&self.capture_records);
}
pub(crate) fn capture_drop(&self) {
bump(&self.capture_dropped);
}
pub(crate) fn capture_error(&self) {
bump(&self.capture_errors);
}
pub(crate) fn capture_hep_record(&self) {
bump(&self.capture_hep_records);
}
pub(crate) fn capture_hep_drop(&self) {
bump(&self.capture_hep_dropped);
}
pub(crate) fn unsent(&self, method: &Method) {
bump(match method {
Method::Invite => &self.unsent.invite,
Method::Ack => &self.unsent.ack,
Method::Bye => &self.unsent.bye,
Method::Cancel => &self.unsent.cancel,
Method::Register
| Method::Options
| Method::Info
| Method::Prack
| Method::Update
| Method::Subscribe
| Method::Notify
| Method::Refer
| Method::Message
| Method::Publish
| Method::Other(_) => &self.unsent.other,
});
}
pub(crate) fn snapshot(&self) -> Counters {
let mut per_transport = [TransportCounts::default(); TRANSPORTS];
for (counts, meter) in per_transport.iter_mut().zip(self.per_transport.iter()) {
*counts = TransportCounts {
requests_in: read(&meter.requests_in),
requests_out: read(&meter.requests_out),
responses_in: read(&meter.responses_in),
responses_out: read(&meter.responses_out),
parse_failures: read(&meter.parse_failures),
source_refusals: read(&meter.source_refusals),
};
}
Counters {
shed: ShedCounts {
requests: read(&self.shed.requests),
acks: read(&self.shed.acks),
unmatched: read(&self.shed.unmatched),
},
overload_rejections: read(&self.overload_rejections),
unmatched_responses: read(&self.unmatched_responses),
retransmissions_sent: read(&self.retransmissions),
oversized_request_tcp_fallbacks: read(&self.oversized_request_tcp_fallbacks),
timeouts: TimeoutCounts {
b: read(&self.timeout_b),
f: read(&self.timeout_f),
h: read(&self.timeout_h),
},
discards: DiscardCounts {
transaction_events: read(&self.discard_transaction_events),
unanswered: read(&self.discard_unanswered),
no_destination: read(&self.discard_no_destination),
send_failures: read(&self.discard_send_failures),
stun_unmatched: read(&self.discard_stun_unmatched),
},
capture: CaptureCounts {
records: read(&self.capture_records),
dropped: read(&self.capture_dropped),
errors: read(&self.capture_errors),
hep_records: read(&self.capture_hep_records),
hep_dropped: read(&self.capture_hep_dropped),
},
observation_dropped: read(&self.observation_dropped),
unsent: UnsentCounts {
invite: read(&self.unsent.invite),
ack: read(&self.unsent.ack),
bye: read(&self.unsent.bye),
cancel: read(&self.unsent.cancel),
other: read(&self.unsent.other),
},
per_transport,
}
}
}
#[cfg(test)]
#[allow(
clippy::unwrap_used,
clippy::expect_used,
clippy::panic,
clippy::indexing_slicing
)]
mod tests {
use super::*;
#[test]
fn every_transport_has_its_own_slot() {
let kinds = [
TransportKind::Udp,
TransportKind::Tcp,
TransportKind::Tls,
TransportKind::Ws,
TransportKind::Wss,
TransportKind::Quic,
];
let mut slots: Vec<usize> = kinds.iter().copied().map(slot).collect();
slots.sort_unstable();
slots.dedup();
assert_eq!(
slots.len(),
TRANSPORTS,
"two transports share a slot, so their counts are being added together"
);
assert_eq!(
slots.last().copied(),
Some(TRANSPORTS - 1),
"a slot is out of range of the array it indexes"
);
}
#[test]
fn a_message_is_counted_against_its_own_transport_only() {
let meters = Meters::default();
meters.message_in(TransportKind::Tcp, false);
meters.message_out(TransportKind::Tcp, true);
let counters = meters.snapshot();
assert_eq!(counters.transport(TransportKind::Tcp).requests_in, 1);
assert_eq!(counters.transport(TransportKind::Tcp).responses_out, 1);
assert_eq!(
counters.transport(TransportKind::Udp),
TransportCounts::default(),
"a TCP message must not appear in UDP's counts"
);
assert_eq!(counters.messages_in(), 1);
assert_eq!(counters.messages_out(), 1);
}
#[test]
fn a_parse_failure_is_not_counted_as_a_message() {
let meters = Meters::default();
meters.parse_failure(TransportKind::Udp);
let counters = meters.snapshot();
assert_eq!(counters.parse_failures(), 1);
assert_eq!(
counters.messages_in(),
0,
"which it would have been is exactly what could not be determined"
);
}
#[test]
fn only_the_retransmission_timers_count_as_retransmissions() {
let meters = Meters::default();
for timer in [Timer::A, Timer::E, Timer::G] {
meters.on_timer(timer);
}
for timer in [
Timer::D,
Timer::I,
Timer::J,
Timer::K,
Timer::L,
Timer::M,
Timer::Trying100,
] {
meters.on_timer(timer);
}
let counters = meters.snapshot();
assert_eq!(counters.retransmissions_sent, 3);
assert_eq!(
counters.timeouts.total(),
0,
"an absorption window closing is not a transaction timing out"
);
}
#[test]
fn each_deadline_timer_is_counted_apart() {
let meters = Meters::default();
meters.on_timer(Timer::B);
meters.on_timer(Timer::H);
meters.on_timer(Timer::H);
let counters = meters.snapshot();
assert_eq!(counters.timeouts.b, 1);
assert_eq!(counters.timeouts.f, 0);
assert_eq!(counters.timeouts.h, 2);
assert_eq!(counters.timeouts.total(), 3);
}
#[test]
fn any_loss_covers_discards_and_not_arrivals() {
let meters = Meters::default();
assert!(!meters.snapshot().any_loss());
meters.parse_failure(TransportKind::Udp);
meters.on_timer(Timer::B);
assert!(
!meters.snapshot().any_loss(),
"a stranger's malformed datagram is not this endpoint losing something"
);
meters.discard_transaction_event();
assert!(meters.snapshot().any_loss());
}
#[test]
fn a_fresh_endpoint_reports_zero_everywhere() {
let counters = Meters::default().snapshot();
assert_eq!(counters, Counters::default());
assert!(!counters.any_loss());
assert_eq!(counters.capture, CaptureCounts::default());
}
}