use std::time::Duration;
use tokio::time::Instant;
#[derive(Debug, Clone, Copy)]
pub struct NodeReadCounts {
pub valid_samples: u64,
pub invalid_samples: u64,
}
pub trait ReadAnchor: 'static {
fn elapsed(self, now: Instant) -> Duration;
}
impl ReadAnchor for Instant {
fn elapsed(self, now: Instant) -> Duration {
now.duration_since(self)
}
}
#[derive(Debug)]
pub struct PollLoopMetricsLogger {
flush_interval: Duration,
window: MetricsWindow,
}
#[derive(Debug)]
struct MetricsWindow {
pub window_start: Instant,
pub iterations: u64,
pub ok_count: u64,
pub err_count: u64,
pub read_duration_sum: Duration,
pub read_duration_min: Option<Duration>,
pub read_duration_max: Option<Duration>,
pub valid_samples: u64,
pub invalid_samples: u64,
}
impl MetricsWindow {
fn new() -> Self {
Self {
window_start: Instant::now(),
iterations: 0,
ok_count: 0,
err_count: 0,
read_duration_sum: Duration::ZERO,
read_duration_min: None,
read_duration_max: None,
valid_samples: 0,
invalid_samples: 0,
}
}
fn update(&mut self, read_duration: Duration, counts: Option<NodeReadCounts>) {
self.iterations = self.iterations.saturating_add(1);
self.read_duration_sum = self.read_duration_sum.saturating_add(read_duration);
if let Some(min) = self.read_duration_min {
self.read_duration_min = Some(min.min(read_duration));
} else {
self.read_duration_min = Some(read_duration);
}
if let Some(max) = self.read_duration_max {
self.read_duration_max = Some(max.max(read_duration));
} else {
self.read_duration_max = Some(read_duration);
}
match counts {
Some(NodeReadCounts {
valid_samples,
invalid_samples,
}) => {
self.ok_count = self.ok_count.saturating_add(1);
self.valid_samples = self.valid_samples.saturating_add(valid_samples);
self.invalid_samples = self.invalid_samples.saturating_add(invalid_samples);
}
None => {
self.err_count = self.err_count.saturating_add(1);
}
}
}
}
impl PollLoopMetricsLogger {
#[must_use]
pub fn new(flush_interval: Duration) -> Self {
Self {
flush_interval,
window: MetricsWindow::new(),
}
}
#[must_use]
pub fn start_read(&self) -> impl ReadAnchor {
Instant::now()
}
pub fn finish_read(&mut self, anchor: impl ReadAnchor, counts: Option<NodeReadCounts>) {
let duration = anchor.elapsed(Instant::now());
self.window.update(duration, counts);
let now = Instant::now();
if now.duration_since(self.window.window_start) >= self.flush_interval {
self.flush(now);
self.window = MetricsWindow::new()
}
}
fn flush(&self, now: Instant) {
if self.window.iterations < 1 {
return;
}
let window_secs = now.duration_since(self.window.window_start).as_secs_f64();
let loop_rate_hz = if window_secs > 0.0 {
self.window.iterations as f64 / window_secs
} else {
0.0
};
let read_total_count = self.window.ok_count.saturating_add(self.window.err_count);
let read_avg_ms = if read_total_count > 0 {
self.window.read_duration_sum.as_secs_f64() * 1000.0 / read_total_count as f64
} else {
0.0
};
tracing::debug!(
target: "opcua::client::poll_loop",
valid_samples = self.window.valid_samples,
invalid_samples = self.window.invalid_samples,
successful_read_count = self.window.ok_count,
failed_read_count = self.window.err_count,
loop_rate_hz,
read_avg_ms,
);
}
}
impl Drop for PollLoopMetricsLogger {
fn drop(&mut self) {
tracing::info!(target: "opcua::client::poll_loop", "final metrics flush");
self.flush(Instant::now());
}
}
#[cfg(test)]
mod tests {
use tokio::time::Duration;
use tokio::time::Instant;
use super::NodeReadCounts;
use super::PollLoopMetricsLogger;
use super::ReadAnchor;
#[derive(Clone, Copy, Debug)]
struct MockAnchor {
pub start: Instant,
pub end: Instant,
}
impl ReadAnchor for MockAnchor {
fn elapsed(self, _now: Instant) -> Duration {
self.end.duration_since(self.start)
}
}
impl MockAnchor {
fn new(start_offset: Duration, dur: Duration) -> MockAnchor {
let base = Instant::now();
#[expect(clippy::arithmetic_side_effects, reason = "test code")]
MockAnchor {
start: base + start_offset,
end: base + start_offset + dur,
}
}
}
#[test]
fn from_first_read_ok_seeds_window_with_that_read() {
let anchor = MockAnchor::new(Duration::ZERO, Duration::from_millis(10));
let mut m = PollLoopMetricsLogger::new(Duration::from_secs(5));
m.finish_read(
anchor,
Some(NodeReadCounts {
valid_samples: 3,
invalid_samples: 1,
}),
);
assert_eq!(m.window.iterations, 1);
assert_eq!(m.window.ok_count, 1);
assert_eq!(m.window.err_count, 0);
assert_eq!(m.window.valid_samples, 3);
assert_eq!(m.window.invalid_samples, 1);
assert_eq!(m.window.read_duration_sum, Duration::from_millis(10));
assert_eq!(m.window.read_duration_min, Some(Duration::from_millis(10)));
assert_eq!(m.window.read_duration_max, Some(Duration::from_millis(10)));
}
#[test]
fn from_first_read_err_seeds_err_counter() {
let anchor = MockAnchor::new(Duration::ZERO, Duration::from_millis(7));
let mut m = PollLoopMetricsLogger::new(Duration::from_secs(5));
m.finish_read(anchor, None);
assert_eq!(m.window.iterations, 1);
assert_eq!(m.window.ok_count, 0);
assert_eq!(m.window.err_count, 1);
assert_eq!(m.window.valid_samples, 0);
assert_eq!(m.window.invalid_samples, 0);
assert_eq!(m.window.read_duration_sum, Duration::from_millis(7));
assert_eq!(m.window.read_duration_min, Some(Duration::from_millis(7)));
assert_eq!(m.window.read_duration_max, Some(Duration::from_millis(7)));
}
#[test]
fn record_read_ok_accumulates_counts_and_tracks_min_max_sum() {
let a1 = MockAnchor::new(Duration::from_secs(1), Duration::from_millis(10));
let a2 = MockAnchor::new(Duration::from_secs(2), Duration::from_millis(15));
let a3 = MockAnchor::new(Duration::from_secs(3), Duration::from_millis(8));
let mut m = PollLoopMetricsLogger::new(Duration::from_secs(5));
m.finish_read(
a1,
Some(NodeReadCounts {
valid_samples: 3,
invalid_samples: 1,
}),
);
m.finish_read(
a2,
Some(NodeReadCounts {
valid_samples: 2,
invalid_samples: 0,
}),
);
m.finish_read(
a3,
Some(NodeReadCounts {
valid_samples: 4,
invalid_samples: 2,
}),
);
assert_eq!(m.window.iterations, 3);
assert_eq!(m.window.ok_count, 3);
assert_eq!(m.window.err_count, 0);
assert_eq!(m.window.valid_samples, 9);
assert_eq!(m.window.invalid_samples, 3);
assert_eq!(m.window.read_duration_min, Some(Duration::from_millis(8)));
assert_eq!(m.window.read_duration_max, Some(Duration::from_millis(15)));
assert_eq!(m.window.read_duration_sum, Duration::from_millis(33));
}
#[test]
fn record_read_err_increments_err_count_and_iteration() {
let a1 = MockAnchor::new(Duration::from_secs(1), Duration::from_millis(15));
let a2 = MockAnchor::new(Duration::from_secs(2), Duration::from_millis(25));
let mut m = PollLoopMetricsLogger::new(Duration::from_secs(5));
m.finish_read(a1, None);
m.finish_read(a2, None);
assert_eq!(m.window.iterations, 2);
assert_eq!(m.window.ok_count, 0);
assert_eq!(m.window.err_count, 2);
assert_eq!(m.window.valid_samples, 0);
assert_eq!(m.window.invalid_samples, 0);
assert_eq!(m.window.read_duration_sum, Duration::from_millis(40));
assert_eq!(m.window.read_duration_min, Some(Duration::from_millis(15)));
assert_eq!(m.window.read_duration_max, Some(Duration::from_millis(25)));
}
#[test]
fn record_read_mixed_ok_and_err() {
let a1 = MockAnchor::new(Duration::from_secs(1), Duration::from_millis(5));
let a2 = MockAnchor::new(Duration::from_secs(2), Duration::from_millis(50));
let mut m = PollLoopMetricsLogger::new(Duration::from_secs(5));
m.finish_read(
a1,
Some(NodeReadCounts {
valid_samples: 1,
invalid_samples: 0,
}),
);
m.finish_read(a2, None);
assert_eq!(m.window.iterations, 2);
assert_eq!(m.window.ok_count, 1);
assert_eq!(m.window.err_count, 1);
assert_eq!(m.window.valid_samples, 1);
assert_eq!(m.window.invalid_samples, 0);
assert_eq!(m.window.read_duration_sum, Duration::from_millis(55));
assert_eq!(m.window.read_duration_min, Some(Duration::from_millis(5)));
assert_eq!(m.window.read_duration_max, Some(Duration::from_millis(50)));
}
#[test]
fn auto_flush_resets_window() {
let anchor = MockAnchor::new(Duration::ZERO, Duration::from_millis(1));
let mut m = PollLoopMetricsLogger::new(Duration::ZERO);
m.finish_read(
anchor,
Some(NodeReadCounts {
valid_samples: 1,
invalid_samples: 0,
}),
);
assert_eq!(m.window.iterations, 0);
assert_eq!(m.window.ok_count, 0);
assert_eq!(m.window.err_count, 0);
assert_eq!(m.window.valid_samples, 0);
assert_eq!(m.window.invalid_samples, 0);
assert_eq!(m.window.read_duration_sum, Duration::ZERO);
assert_eq!(m.window.read_duration_min, None);
assert_eq!(m.window.read_duration_max, None);
}
#[test]
fn flush_after_only_errors_does_not_panic() {
let anchor = MockAnchor::new(Duration::ZERO, Duration::from_millis(10));
let mut m = PollLoopMetricsLogger::new(Duration::from_secs(5));
m.finish_read(anchor, None);
let flush_at = m.window.window_start + Duration::from_secs(5);
m.flush(flush_at);
}
#[test]
fn instant_now_is_monotonic_relative_to_window_start() {
let anchor = MockAnchor::new(Duration::ZERO, Duration::from_millis(0));
let mut m = PollLoopMetricsLogger::new(Duration::from_secs(5));
m.finish_read(anchor, None);
let later = MockAnchor::new(Duration::from_secs(1), Duration::from_millis(0));
assert!(later.start >= m.window.window_start);
}
}