1use std::time::Duration;
6
7pub trait Telemetry: Send + Sync + 'static {
12 fn on_publish(&self, event_name: &str, receivers: usize);
14
15 fn on_publish_complete(&self, event_name: &str, elapsed: Duration);
17
18 fn on_subscribe(&self, event_name: &str, sub_id: usize);
20
21 fn on_handler_start(&self, event_name: &str, sub_id: usize);
23
24 fn on_handler_complete(
26 &self,
27 event_name: &str,
28 sub_id: usize,
29 elapsed: Duration,
30 error: Option<&str>,
31 );
32
33 fn on_handler_lagged(&self, event_name: &str, sub_id: usize, lagged_count: usize);
35}
36
37pub struct TracingTelemetry;
43
44impl Telemetry for TracingTelemetry {
45 fn on_publish(&self, event_name: &str, receivers: usize) {
46 tracing::debug!(
47 event = event_name,
48 receivers,
49 "telemetry: event publish started"
50 );
51 }
52
53 fn on_publish_complete(&self, event_name: &str, elapsed: Duration) {
54 tracing::debug!(
55 event = event_name,
56 elapsed_ms = elapsed.as_secs_f64() * 1000.0,
57 "telemetry: event publish complete"
58 );
59 }
60
61 fn on_subscribe(&self, event_name: &str, sub_id: usize) {
62 tracing::debug!(
63 event = event_name,
64 sub_id,
65 "telemetry: subscriber registered"
66 );
67 }
68
69 fn on_handler_start(&self, event_name: &str, sub_id: usize) {
70 tracing::debug!(
71 event = event_name,
72 sub_id,
73 "telemetry: handler started"
74 );
75 }
76
77 fn on_handler_complete(
78 &self,
79 event_name: &str,
80 sub_id: usize,
81 elapsed: Duration,
82 error: Option<&str>,
83 ) {
84 if let Some(err) = error {
85 tracing::warn!(
86 event = event_name,
87 sub_id,
88 elapsed_ms = elapsed.as_secs_f64() * 1000.0,
89 error = err,
90 "telemetry: handler completed with error"
91 );
92 } else {
93 tracing::debug!(
94 event = event_name,
95 sub_id,
96 elapsed_ms = elapsed.as_secs_f64() * 1000.0,
97 "telemetry: handler completed"
98 );
99 }
100 }
101
102 fn on_handler_lagged(&self, event_name: &str, sub_id: usize, lagged_count: usize) {
103 tracing::warn!(
104 event = event_name,
105 sub_id,
106 lagged = lagged_count,
107 "telemetry: handler lagged"
108 );
109 }
110}
111
112pub struct NoopTelemetry;
118
119impl Telemetry for NoopTelemetry {
120 fn on_publish(&self, _event_name: &str, _receivers: usize) {}
121 fn on_publish_complete(&self, _event_name: &str, _elapsed: Duration) {}
122 fn on_subscribe(&self, _event_name: &str, _sub_id: usize) {}
123 fn on_handler_start(&self, _event_name: &str, _sub_id: usize) {}
124 fn on_handler_complete(
125 &self,
126 _event_name: &str,
127 _sub_id: usize,
128 _elapsed: Duration,
129 _error: Option<&str>,
130 ) {
131 }
132 fn on_handler_lagged(&self, _event_name: &str, _sub_id: usize, _lagged_count: usize) {}
133}
134
135#[cfg(test)]
136mod tests {
137 use super::*;
138 use std::time::Duration;
139
140 #[test]
141 fn noop_telemetry_does_not_panic() {
142 let tel = NoopTelemetry;
143 tel.on_publish("test.event", 3);
144 tel.on_publish_complete("test.event", Duration::from_millis(10));
145 tel.on_subscribe("test.event", 1);
146 tel.on_handler_start("test.event", 1);
147 tel.on_handler_complete("test.event", 1, Duration::from_millis(5), None);
148 tel.on_handler_complete("test.event", 1, Duration::from_millis(5), Some("err"));
149 tel.on_handler_lagged("test.event", 1, 42);
150 }
151
152 #[test]
153 fn tracing_telemetry_does_not_panic() {
154 let tel = TracingTelemetry;
155 tel.on_publish("test.event", 3);
156 tel.on_publish_complete("test.event", Duration::from_millis(10));
157 tel.on_subscribe("test.event", 1);
158 tel.on_handler_start("test.event", 1);
159 tel.on_handler_complete("test.event", 1, Duration::from_millis(5), None);
160 tel.on_handler_complete("test.event", 1, Duration::from_millis(5), Some("err"));
161 tel.on_handler_lagged("test.event", 1, 42);
162 }
163}