use crate::client::connection::ConnectionMetadata;
use aws_smithy_types::config_bag::{Storable, StoreReplace};
use std::fmt;
use std::sync::{Arc, Mutex, MutexGuard};
use std::time::{Duration, SystemTime};
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum ConnectionUsage {
Fresh,
Reused,
}
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub struct ConnectionAcquisitionTelemetry {
duration: Option<Duration>,
usage: ConnectionUsage,
}
impl ConnectionAcquisitionTelemetry {
pub fn new(duration: Duration, usage: ConnectionUsage) -> Self {
Self {
duration: Some(duration),
usage,
}
}
pub fn from_interval(
started_at: SystemTime,
completed_at: SystemTime,
usage: ConnectionUsage,
) -> Self {
Self {
duration: completed_at.duration_since(started_at).ok(),
usage,
}
}
pub fn duration(&self) -> Option<Duration> {
self.duration
}
pub fn usage(&self) -> ConnectionUsage {
self.usage
}
}
#[derive(Clone, Debug, Default)]
#[non_exhaustive]
pub struct HttpAttemptTelemetry {
selection: Option<ConnectionSelection>,
connector_call_duration: Option<Duration>,
}
impl HttpAttemptTelemetry {
pub fn acquisition(&self) -> Option<&ConnectionAcquisitionTelemetry> {
self.selection
.as_ref()
.map(|selection| &selection.acquisition)
}
pub fn connection(&self) -> Option<&ConnectionMetadata> {
self.selection
.as_ref()
.map(|selection| &selection.connection)
}
pub fn connector_call_duration(&self) -> Option<Duration> {
self.connector_call_duration
}
}
#[derive(Clone, Debug)]
struct ConnectionSelection {
acquisition: ConnectionAcquisitionTelemetry,
connection: ConnectionMetadata,
}
#[derive(Clone, Default)]
pub struct CaptureHttpAttemptTelemetry {
state: Arc<Mutex<HttpAttemptTelemetry>>,
}
impl CaptureHttpAttemptTelemetry {
pub fn new() -> Self {
Self::default()
}
pub fn get(&self) -> HttpAttemptTelemetry {
self.lock().clone()
}
pub fn record_connector_call_duration(&self, duration: Duration) -> bool {
let mut state = self.lock();
if state.connector_call_duration.is_some() {
return false;
}
state.connector_call_duration = Some(duration);
true
}
pub fn record_connection_selection(
&self,
acquisition: ConnectionAcquisitionTelemetry,
connection: ConnectionMetadata,
) -> bool {
let mut state = self.lock();
if state.selection.is_some() {
return false;
}
state.selection = Some(ConnectionSelection {
acquisition,
connection,
});
true
}
fn lock(&self) -> MutexGuard<'_, HttpAttemptTelemetry> {
self.state
.lock()
.unwrap_or_else(std::sync::PoisonError::into_inner)
}
}
impl fmt::Debug for CaptureHttpAttemptTelemetry {
fn fmt(&self, f: &mut fmt::Formatter<'_>) -> fmt::Result {
f.debug_struct("CaptureHttpAttemptTelemetry")
.finish_non_exhaustive()
}
}
impl Storable for CaptureHttpAttemptTelemetry {
type Storer = StoreReplace<Self>;
}
#[cfg(test)]
mod tests {
use super::*;
fn connection() -> ConnectionMetadata {
ConnectionMetadata::builder()
.proxied(false)
.poison_fn(|| {})
.build()
}
#[test]
fn records_connector_call_and_selection_independently() {
let capture = CaptureHttpAttemptTelemetry::new();
let acquisition =
ConnectionAcquisitionTelemetry::new(Duration::from_millis(3), ConnectionUsage::Fresh);
assert!(capture.record_connection_selection(acquisition, connection()));
assert!(capture.record_connector_call_duration(Duration::from_millis(7)));
let telemetry = capture.get();
assert_eq!(telemetry.acquisition(), Some(&acquisition));
assert_eq!(
telemetry.connector_call_duration(),
Some(Duration::from_millis(7))
);
assert!(telemetry.connection().is_some());
}
#[test]
fn first_recorded_value_wins() {
let capture = CaptureHttpAttemptTelemetry::new();
let first =
ConnectionAcquisitionTelemetry::new(Duration::from_millis(3), ConnectionUsage::Fresh);
let second =
ConnectionAcquisitionTelemetry::new(Duration::from_millis(9), ConnectionUsage::Reused);
assert!(capture.record_connection_selection(first, connection()));
assert!(!capture.record_connection_selection(second, connection()));
assert!(capture.record_connector_call_duration(Duration::from_millis(4)));
assert!(!capture.record_connector_call_duration(Duration::from_millis(8)));
let telemetry = capture.get();
assert_eq!(telemetry.acquisition(), Some(&first));
assert_eq!(
telemetry.connector_call_duration(),
Some(Duration::from_millis(4))
);
}
#[test]
fn backwards_acquisition_interval_retains_selection_facts() {
let started_at = SystemTime::UNIX_EPOCH + Duration::from_secs(2);
let completed_at = SystemTime::UNIX_EPOCH + Duration::from_secs(1);
let acquisition = ConnectionAcquisitionTelemetry::from_interval(
started_at,
completed_at,
ConnectionUsage::Fresh,
);
let capture = CaptureHttpAttemptTelemetry::new();
assert!(capture.record_connection_selection(acquisition, connection()));
let telemetry = capture.get();
assert_eq!(telemetry.acquisition().expect("selection").duration(), None);
assert_eq!(
telemetry.acquisition().expect("selection").usage(),
ConnectionUsage::Fresh
);
assert!(telemetry.connection().is_some());
}
}