rama_http/layer/trace/
on_eos.rs1use super::{DEFAULT_MESSAGE_LEVEL, Latency};
2use crate::header::HeaderMap;
3use crate::layer::classify::grpc_errors_as_failures::ParsedGrpcStatus;
4use rama_core::telemetry::tracing::{Level, Span};
5use rama_utils::latency::LatencyUnit;
6use std::time::Duration;
7
8pub trait OnEos: Send + Sync + 'static {
15 fn on_eos(self, trailers: Option<&HeaderMap>, stream_duration: Duration, span: &Span);
27}
28
29impl OnEos for () {
30 #[inline]
31 fn on_eos(self, _: Option<&HeaderMap>, _: Duration, _: &Span) {}
32}
33
34impl<F> OnEos for F
35where
36 F: Fn(Option<&HeaderMap>, Duration, &Span) + Send + Sync + 'static,
37{
38 fn on_eos(self, trailers: Option<&HeaderMap>, stream_duration: Duration, span: &Span) {
39 self(trailers, stream_duration, span)
40 }
41}
42
43#[derive(Clone, Debug)]
47pub struct DefaultOnEos {
48 level: Level,
49 latency_unit: LatencyUnit,
50}
51
52impl Default for DefaultOnEos {
53 fn default() -> Self {
54 Self {
55 level: DEFAULT_MESSAGE_LEVEL,
56 latency_unit: LatencyUnit::Millis,
57 }
58 }
59}
60
61impl DefaultOnEos {
62 #[must_use]
64 pub fn new() -> Self {
65 Self::default()
66 }
67
68 rama_utils::macros::generate_set_and_with! {
69 pub fn level(mut self, level: Level) -> Self {
76 self.level = level;
77 self
78 }
79 }
80
81 rama_utils::macros::generate_set_and_with! {
82 pub fn latency_unit(mut self, latency_unit: LatencyUnit) -> Self {
86 self.latency_unit = latency_unit;
87 self
88 }
89 }
90}
91
92impl OnEos for DefaultOnEos {
93 fn on_eos(self, trailers: Option<&HeaderMap>, stream_duration: Duration, span: &Span) {
94 let stream_duration = Latency {
95 unit: self.latency_unit,
96 duration: stream_duration,
97 };
98 let status = trailers.and_then(|trailers| {
99 match crate::layer::classify::grpc_errors_as_failures::classify_grpc_metadata(
100 trailers,
101 crate::layer::classify::GrpcCode::Ok.into_bitmask(),
102 ) {
103 ParsedGrpcStatus::Success | ParsedGrpcStatus::HeaderNotGrpcCode => Some(0),
104 ParsedGrpcStatus::NonSuccess(status) => Some(status.code_raw()),
105 ParsedGrpcStatus::GrpcStatusHeaderMissing => None,
106 }
107 });
108
109 event_dynamic_lvl!(parent: span, self.level, %stream_duration, status, "end of stream");
110 }
111}