use std::sync::Arc;
use apollo_opentelemetry::default_instrumentation_scope;
use apollo_opentelemetry::metrics::{
Clock, HistogramExt, RecordDurationGuard, TrackGuard, UpDownCounterExt,
};
use http::{Method, StatusCode};
use opentelemetry::KeyValue;
use opentelemetry::metrics::{Histogram, UpDownCounter};
use opentelemetry_semantic_conventions::{attribute as semconv, metric as metric_semconv};
use apollo_http_shared::body::BodySizeGuard;
use apollo_http_shared::method::normalize_method;
use crate::protocol::{Origin, ProtocolVersionCell, version_str};
const REQUEST_DURATION_BOUNDARIES: &[f64] = &[
0.005, 0.01, 0.025, 0.05, 0.075, 0.1, 0.25, 0.5, 0.75, 1.0, 2.5, 5.0, 7.5, 10.0,
];
const CONNECTION_DURATION_BOUNDARIES: &[f64] = &[
0.01, 0.02, 0.05, 0.1, 0.2, 0.5, 1.0, 2.0, 5.0, 10.0, 30.0, 60.0, 120.0, 300.0,
];
const BODY_SIZE_BOUNDARIES: &[f64] = &[
0.0,
1_024.0,
4_096.0,
16_384.0,
65_536.0,
262_144.0,
1_048_576.0,
4_194_304.0,
16_777_216.0,
];
pub(crate) const STATE_IDLE: &str = "idle";
pub(crate) const STATE_ACTIVE: &str = "active";
pub(crate) const STATE_CUSTOM_CONNECTING: &str = "connecting";
#[derive(Clone, Debug)]
pub(crate) struct HttpMetrics {
active_requests: UpDownCounter<i64>,
request_duration: Histogram<f64>,
request_body_size: Histogram<f64>,
response_body_size: Histogram<f64>,
connection_duration: Histogram<f64>,
open_connections: UpDownCounter<i64>,
clock: Clock,
record_request_body_size: bool,
record_response_body_size: bool,
record_url_scheme: bool,
}
impl HttpMetrics {
pub(crate) fn new(
clock: Clock,
record_request_body_size: bool,
record_response_body_size: bool,
record_url_scheme: bool,
) -> Self {
let meter =
opentelemetry::global::meter_with_scope(default_instrumentation_scope!().clone());
let active_requests = meter
.i64_up_down_counter(metric_semconv::HTTP_CLIENT_ACTIVE_REQUESTS)
.with_description("Number of HTTP requests currently in flight")
.with_unit("{request}")
.build();
let request_duration = meter
.f64_histogram(metric_semconv::HTTP_CLIENT_REQUEST_DURATION)
.with_description("Duration of HTTP client requests")
.with_unit("s")
.with_boundaries(REQUEST_DURATION_BOUNDARIES.to_vec())
.build();
let request_body_size = meter
.f64_histogram(metric_semconv::HTTP_CLIENT_REQUEST_BODY_SIZE)
.with_description("Size of HTTP request bodies sent by the client")
.with_unit("By")
.with_boundaries(BODY_SIZE_BOUNDARIES.to_vec())
.build();
let response_body_size = meter
.f64_histogram(metric_semconv::HTTP_CLIENT_RESPONSE_BODY_SIZE)
.with_description("Size of HTTP response bodies received by the client")
.with_unit("By")
.with_boundaries(BODY_SIZE_BOUNDARIES.to_vec())
.build();
let connection_duration = meter
.f64_histogram(metric_semconv::HTTP_CLIENT_CONNECTION_DURATION)
.with_description("Duration of outbound HTTP connections from establishment to close")
.with_unit("s")
.with_boundaries(CONNECTION_DURATION_BOUNDARIES.to_vec())
.build();
let open_connections = meter
.i64_up_down_counter(metric_semconv::HTTP_CLIENT_OPEN_CONNECTIONS)
.with_description("Number of open connections in the HTTP client pool")
.with_unit("{connection}")
.build();
Self {
active_requests,
request_duration,
request_body_size,
response_body_size,
connection_duration,
open_connections,
clock,
record_request_body_size,
record_response_body_size,
record_url_scheme,
}
}
pub(crate) fn begin_request(
&self,
method: &Method,
origin: &Origin,
version: ProtocolVersionCell,
) -> (RequestMetrics, BodySizeGuard) {
let method = normalize_method(method);
let server_address = origin.address();
let server_port = origin.port();
let mut base = vec![
KeyValue::new(semconv::HTTP_REQUEST_METHOD, method),
KeyValue::new(
semconv::NETWORK_PROTOCOL_VERSION,
version_str(version.get()),
),
KeyValue::new(semconv::SERVER_ADDRESS, server_address.clone()),
];
let mut active_attrs = vec![
KeyValue::new(semconv::HTTP_REQUEST_METHOD, method),
KeyValue::new(semconv::SERVER_ADDRESS, server_address),
];
if let Some(port) = server_port {
base.push(KeyValue::new(semconv::SERVER_PORT, port));
active_attrs.push(KeyValue::new(semconv::SERVER_PORT, port));
}
if self.record_url_scheme {
base.push(KeyValue::new(
semconv::URL_SCHEME,
origin.scheme.to_string(),
));
active_attrs.push(KeyValue::new(
semconv::URL_SCHEME,
origin.scheme.to_string(),
));
}
let req_metrics = RequestMetrics {
_active: self.active_requests.track(active_attrs),
timer: self
.request_duration
.record_duration_on_drop_with_clock(self.clock.clone(), base.clone()),
version,
};
let body_size = if self.record_request_body_size {
BodySizeGuard::new(self.request_body_size.record_on_drop(0.0, base))
} else {
BodySizeGuard::noop()
};
(req_metrics, body_size)
}
pub(crate) fn begin_response(
&self,
method: &Method,
origin: &Origin,
status_code: StatusCode,
) -> BodySizeGuard {
let method = normalize_method(method);
if self.record_response_body_size {
let mut attrs = vec![
KeyValue::new(semconv::HTTP_REQUEST_METHOD, method),
KeyValue::new(
semconv::HTTP_RESPONSE_STATUS_CODE,
i64::from(status_code.as_u16()),
),
KeyValue::new(semconv::SERVER_ADDRESS, origin.address()),
];
if let Some(port) = origin.port() {
attrs.push(KeyValue::new(semconv::SERVER_PORT, port));
}
if self.record_url_scheme {
attrs.push(KeyValue::new(
semconv::URL_SCHEME,
origin.scheme.to_string(),
));
}
BodySizeGuard::new(self.response_body_size.record_on_drop(0.0, attrs))
} else {
BodySizeGuard::noop()
}
}
pub(crate) fn connection_metrics(&self) -> ConnectionMetrics {
ConnectionMetrics {
counter: self.open_connections.clone(),
duration: self.connection_duration.clone(),
record_url_scheme: self.record_url_scheme,
peer: None,
}
}
}
#[derive(Clone)]
struct Peer {
address: Arc<str>,
port: u16,
}
#[derive(Clone)]
pub(crate) struct ConnectionMetrics {
counter: UpDownCounter<i64>,
duration: Histogram<f64>,
record_url_scheme: bool,
peer: Option<Peer>,
}
impl ConnectionMetrics {
pub(crate) fn with_peer(mut self, host: impl Into<Arc<str>>, port: u16) -> Self {
self.peer = Some(Peer {
address: host.into(),
port,
});
self
}
pub(crate) fn connection_attrs(
&self,
protocol_version: Option<http::Version>,
origin: &Origin,
) -> Vec<KeyValue> {
let mut attrs = vec![KeyValue::new(semconv::SERVER_ADDRESS, origin.address())];
if let Some(port) = origin.port() {
attrs.push(KeyValue::new(semconv::SERVER_PORT, port));
}
if let Some(protocol_version) = protocol_version {
attrs.push(KeyValue::new(
semconv::NETWORK_PROTOCOL_VERSION,
version_str(protocol_version),
));
}
if self.record_url_scheme {
attrs.push(KeyValue::new(
semconv::URL_SCHEME,
origin.scheme.to_string(),
));
}
if let Some(peer) = &self.peer {
attrs.push(KeyValue::new(
semconv::NETWORK_PEER_ADDRESS,
peer.address.to_string(),
));
attrs.push(KeyValue::new(
semconv::NETWORK_PEER_PORT,
i64::from(peer.port),
));
}
attrs
}
pub(crate) fn counter(&self) -> &UpDownCounter<i64> {
&self.counter
}
pub(crate) fn duration_histogram(&self) -> &Histogram<f64> {
&self.duration
}
}
pub(crate) struct RequestMetrics {
_active: TrackGuard<i64>,
timer: RecordDurationGuard,
version: ProtocolVersionCell,
}
impl Drop for RequestMetrics {
fn drop(&mut self) {
self.timer.set(KeyValue::new(
semconv::NETWORK_PROTOCOL_VERSION,
version_str(self.version.get()),
));
}
}
impl RequestMetrics {
pub(crate) fn on_success(&mut self, status: http::StatusCode) {
let status_code = status.as_u16();
self.timer.set(KeyValue::new(
semconv::HTTP_RESPONSE_STATUS_CODE,
i64::from(status_code),
));
if status.is_client_error() || status.is_server_error() {
self.timer
.set(KeyValue::new(semconv::ERROR_TYPE, status_code.to_string()));
}
}
pub(crate) fn on_error(&mut self) {
self.timer.set(KeyValue::new(semconv::ERROR_TYPE, "_OTHER"));
}
}
#[cfg(test)]
mod tests {
use apollo_opentelemetry::metrics::Clock;
use apollo_opentelemetry_test::TelemetryContext;
use http::uri::{Authority, Scheme};
use opentelemetry::{KeyValue, Value};
fn make_conn_metrics(record_url_scheme: bool) -> super::ConnectionMetrics {
let (clock, _mock) = Clock::mock();
super::HttpMetrics::new(clock, false, false, record_url_scheme).connection_metrics()
}
fn find<'a>(attrs: &'a [KeyValue], key: &str) -> Option<&'a Value> {
attrs
.iter()
.find(|kv| kv.key.as_str() == key)
.map(|kv| &kv.value)
}
fn pool_key(scheme: Scheme, authority: Authority) -> super::Origin {
super::Origin { scheme, authority }
}
#[test]
fn standard_fields_always_present() {
let _ctx = TelemetryContext::new();
let metrics = make_conn_metrics(false);
let attrs = metrics.connection_attrs(
Some(http::Version::HTTP_11),
&pool_key(Scheme::HTTP, Authority::from_static("example.com:8080")),
);
assert_eq!(
find(&attrs, "network.protocol.version"),
Some(&Value::from("1.1"))
);
assert_eq!(
find(&attrs, "server.address"),
Some(&Value::from("example.com"))
);
assert_eq!(find(&attrs, "server.port"), Some(&Value::I64(8080)));
}
#[test]
fn url_scheme_absent_by_default() {
let _ctx = TelemetryContext::new();
let metrics = make_conn_metrics(false);
let attrs = metrics.connection_attrs(
Some(http::Version::HTTP_11),
&pool_key(Scheme::HTTP, Authority::from_static("example.com:80")),
);
assert!(find(&attrs, "url.scheme").is_none());
}
#[test]
fn url_scheme_present_when_enabled() {
let _ctx = TelemetryContext::new();
let metrics = make_conn_metrics(true);
let attrs = metrics.connection_attrs(
Some(http::Version::HTTP_11),
&pool_key(Scheme::HTTP, Authority::from_static("example.com:80")),
);
assert_eq!(find(&attrs, "url.scheme"), Some(&Value::from("http")));
}
#[cfg(unix)]
#[test]
fn unix_attrs_set_address_and_omit_port() {
let _ctx = TelemetryContext::new();
let metrics = make_conn_metrics(false);
let authority = crate::protocol::unix::path_to_authority("/var/run/app.sock").unwrap();
let attrs = metrics.connection_attrs(
Some(http::Version::HTTP_11),
&pool_key(Scheme::try_from("unix").unwrap(), authority),
);
assert_eq!(
find(&attrs, "server.address"),
Some(&Value::from("/var/run/app.sock"))
);
assert!(find(&attrs, "server.port").is_none());
assert_eq!(
find(&attrs, "network.protocol.version"),
Some(&Value::from("1.1"))
);
}
#[cfg(unix)]
#[test]
fn unix_attrs_record_http2_version() {
let _ctx = TelemetryContext::new();
let metrics = make_conn_metrics(false);
let authority = crate::protocol::unix::path_to_authority("/sock").unwrap();
let attrs = metrics.connection_attrs(
Some(http::Version::HTTP_2),
&pool_key(Scheme::try_from("unix").unwrap(), authority),
);
assert_eq!(
find(&attrs, "network.protocol.version"),
Some(&Value::from("2"))
);
}
#[cfg(unix)]
#[test]
fn unix_attrs_omit_version_when_none() {
let _ctx = TelemetryContext::new();
let metrics = make_conn_metrics(false);
let authority = crate::protocol::unix::path_to_authority("/sock").unwrap();
let attrs = metrics.connection_attrs(
None,
&pool_key(Scheme::try_from("unix").unwrap(), authority),
);
assert!(find(&attrs, "network.protocol.version").is_none());
}
#[cfg(unix)]
#[test]
fn unix_url_scheme_respects_flag() {
let _ctx = TelemetryContext::new();
let authority = crate::protocol::unix::path_to_authority("/sock").unwrap();
let key = pool_key(Scheme::try_from("unix").unwrap(), authority);
let off = make_conn_metrics(false).connection_attrs(None, &key);
assert!(find(&off, "url.scheme").is_none());
let on = make_conn_metrics(true).connection_attrs(None, &key);
assert_eq!(find(&on, "url.scheme"), Some(&Value::from("unix")));
}
#[test]
fn network_peer_absent_when_unset() {
let _ctx = TelemetryContext::new();
let metrics = make_conn_metrics(false);
let attrs = metrics.connection_attrs(
Some(http::Version::HTTP_11),
&pool_key(Scheme::HTTP, Authority::from_static("example.com:80")),
);
assert!(find(&attrs, "network.peer.address").is_none());
assert!(find(&attrs, "network.peer.port").is_none());
}
#[test]
fn network_peer_present_when_set() {
let _ctx = TelemetryContext::new();
let metrics = make_conn_metrics(false).with_peer("proxy.example.com", 3128);
let attrs = metrics.connection_attrs(
Some(http::Version::HTTP_11),
&pool_key(Scheme::HTTP, Authority::from_static("example.com:80")),
);
assert_eq!(
find(&attrs, "network.peer.address"),
Some(&Value::from("proxy.example.com")),
);
assert_eq!(find(&attrs, "network.peer.port"), Some(&Value::I64(3128)));
assert_eq!(
find(&attrs, "server.address"),
Some(&Value::from("example.com"))
);
assert_eq!(find(&attrs, "server.port"), Some(&Value::I64(80)));
}
}