1use std::{
4 sync::{Arc, Mutex},
5 time::Duration,
6};
7
8use opentelemetry::{KeyValue, global, trace::TracerProvider};
9use opentelemetry_otlp::WithExportConfig;
10use opentelemetry_sdk::{Resource, metrics::SdkMeterProvider, trace::SdkTracerProvider};
11use serde::{Deserialize, Serialize};
12use tracing_subscriber::{EnvFilter, layer::SubscriberExt, util::SubscriberInitExt};
13
14use crate::{env::TelemetryEnvOverrides, project::PROJECT_NAME};
15
16const DEFAULT_OTEL_ENDPOINT: &str = "http://127.0.0.1:4318";
17const DEFAULT_EXPORT_TIMEOUT_MS: u64 = 5_000;
18const OTLP_TRACE_PATH: &str = "/v1/traces";
19const OTLP_METRIC_PATH: &str = "/v1/metrics";
20
21#[derive(Debug, Clone, PartialEq, Eq)]
23pub struct TelemetryConfig {
24 pub otel_endpoint: String,
25 pub traces_enabled: bool,
26 pub metrics_enabled: bool,
27 pub export_timeout: Duration,
28 pub service_environment: String,
29}
30
31impl TelemetryConfig {
32 pub fn from_environment(environment: &TelemetryEnvOverrides) -> Self {
34 Self {
35 otel_endpoint: environment
36 .otel_endpoint
37 .clone()
38 .unwrap_or_else(|| DEFAULT_OTEL_ENDPOINT.to_owned()),
39 traces_enabled: environment.otel_traces.unwrap_or(false),
40 metrics_enabled: environment.otel_metrics.unwrap_or(false),
41 export_timeout: Duration::from_millis(
42 environment
43 .export_timeout_ms
44 .unwrap_or(DEFAULT_EXPORT_TIMEOUT_MS),
45 ),
46 service_environment: environment
47 .service_environment
48 .clone()
49 .unwrap_or_else(|| "local".to_owned()),
50 }
51 }
52
53 fn trace_endpoint(&self) -> String {
54 signal_endpoint(&self.otel_endpoint, OTLP_TRACE_PATH)
55 }
56
57 fn metric_endpoint(&self) -> String {
58 signal_endpoint(&self.otel_endpoint, OTLP_METRIC_PATH)
59 }
60}
61
62#[derive(Debug, Clone)]
64pub struct ObservabilityRuntime {
65 config: TelemetryConfig,
66 state: Arc<Mutex<ObservabilityState>>,
67 metrics: AgentProtocolMetrics,
68}
69
70#[derive(Debug, Default)]
71struct ObservabilityState {
72 trace_initialized: bool,
73 metrics_initialized: bool,
74 trace_provider: Option<SdkTracerProvider>,
75 metrics_provider: Option<SdkMeterProvider>,
76 last_error: Option<String>,
77}
78
79impl ObservabilityRuntime {
80 pub fn new(config: TelemetryConfig) -> Self {
82 Self {
83 config,
84 state: Arc::new(Mutex::new(ObservabilityState::default())),
85 metrics: AgentProtocolMetrics::default(),
86 }
87 }
88
89 pub fn initialize(&self) {
91 let initialized = self.try_initialize();
92 let mut state = self
93 .state
94 .lock()
95 .unwrap_or_else(|poisoned| poisoned.into_inner());
96 state.trace_initialized = initialized.trace_initialized;
97 state.metrics_initialized = initialized.metrics_initialized;
98 state.trace_provider = initialized.trace_provider;
99 state.metrics_provider = initialized.metrics_provider;
100 state.last_error = initialized.last_error;
101 }
102
103 pub fn agent_metrics(&self) -> AgentProtocolMetrics {
105 self.metrics.clone()
106 }
107
108 pub fn status(&self) -> TelemetryStatus {
110 let state = self
111 .state
112 .lock()
113 .unwrap_or_else(|poisoned| poisoned.into_inner());
114 TelemetryStatus {
115 otlp_endpoint_configured: self.config.otel_endpoint != DEFAULT_OTEL_ENDPOINT,
116 traces_enabled: self.config.traces_enabled,
117 metrics_enabled: self.config.metrics_enabled,
118 trace_exporter_initialized: state.trace_initialized,
119 metrics_exporter_initialized: state.metrics_initialized,
120 export_timeout_ms: duration_millis(self.config.export_timeout),
121 service_environment: self.config.service_environment.clone(),
122 last_error: state.last_error.clone(),
123 agent_protocol: self.metrics.snapshot(),
124 }
125 }
126
127 pub fn shutdown(&self) {
129 let mut state = self
130 .state
131 .lock()
132 .unwrap_or_else(|poisoned| poisoned.into_inner());
133 if let Some(provider) = state.trace_provider.take() {
134 if let Err(error) = provider.shutdown_with_timeout(self.config.export_timeout) {
135 state.push_error(format!("trace shutdown: {error}"));
136 }
137 state.trace_initialized = false;
138 }
139 if let Some(provider) = state.metrics_provider.take() {
140 if let Err(error) = provider.shutdown_with_timeout(self.config.export_timeout) {
141 state.push_error(format!("metrics shutdown: {error}"));
142 }
143 state.metrics_initialized = false;
144 }
145 }
146
147 fn try_initialize(&self) -> InitializedTelemetry {
148 let resource = Resource::builder()
149 .with_service_name(PROJECT_NAME.to_owned())
150 .with_attribute(KeyValue::new(
151 "deployment.environment",
152 self.config.service_environment.clone(),
153 ))
154 .build();
155 let mut initialized = InitializedTelemetry::default();
156
157 if self.config.metrics_enabled {
158 match opentelemetry_otlp::MetricExporter::builder()
159 .with_http()
160 .with_endpoint(self.config.metric_endpoint())
161 .with_timeout(self.config.export_timeout)
162 .build()
163 {
164 Ok(exporter) => {
165 let reader = opentelemetry_sdk::metrics::PeriodicReader::builder(exporter)
166 .with_interval(Duration::from_secs(5))
167 .build();
168 let provider = SdkMeterProvider::builder()
169 .with_resource(resource.clone())
170 .with_reader(reader)
171 .build();
172 global::set_meter_provider(provider.clone());
173 initialized.metrics_provider = Some(provider);
174 initialized.metrics_initialized = true;
175 }
176 Err(error) => initialized.push_error(format!("metrics exporter: {error}")),
177 }
178 }
179
180 if self.config.traces_enabled {
181 match opentelemetry_otlp::SpanExporter::builder()
182 .with_http()
183 .with_endpoint(self.config.trace_endpoint())
184 .with_timeout(self.config.export_timeout)
185 .build()
186 {
187 Ok(exporter) => {
188 let provider = SdkTracerProvider::builder()
189 .with_resource(resource)
190 .with_batch_exporter(exporter)
191 .build();
192 let tracer = provider.tracer(PROJECT_NAME.to_owned());
193 global::set_tracer_provider(provider.clone());
194 match install_otel_subscriber(tracer) {
195 Ok(()) => {
196 initialized.trace_provider = Some(provider);
197 initialized.trace_initialized = true;
198 }
199 Err(error) => initialized.push_error(format!("trace subscriber: {error}")),
200 }
201 }
202 Err(error) => {
203 initialized.push_error(format!("trace exporter: {error}"));
204 install_fallback_subscriber(&mut initialized);
205 }
206 }
207 } else {
208 install_fallback_subscriber(&mut initialized);
209 }
210
211 initialized
212 }
213}
214
215#[derive(Default)]
216struct InitializedTelemetry {
217 trace_initialized: bool,
218 metrics_initialized: bool,
219 trace_provider: Option<SdkTracerProvider>,
220 metrics_provider: Option<SdkMeterProvider>,
221 last_error: Option<String>,
222}
223
224impl InitializedTelemetry {
225 fn push_error(&mut self, error: String) {
226 match &mut self.last_error {
227 Some(existing) => {
228 existing.push_str("; ");
229 existing.push_str(&error);
230 }
231 None => self.last_error = Some(error),
232 }
233 }
234}
235
236impl ObservabilityState {
237 fn push_error(&mut self, error: String) {
238 match &mut self.last_error {
239 Some(existing) => {
240 existing.push_str("; ");
241 existing.push_str(&error);
242 }
243 None => self.last_error = Some(error),
244 }
245 }
246}
247
248fn install_otel_subscriber(
249 tracer: opentelemetry_sdk::trace::SdkTracer,
250) -> Result<(), tracing_subscriber::util::TryInitError> {
251 let filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info"));
252 let otel_layer = tracing_opentelemetry::layer().with_tracer(tracer);
253 tracing_subscriber::registry()
254 .with(filter)
255 .with(tracing_subscriber::fmt::layer())
256 .with(otel_layer)
257 .try_init()
258}
259
260fn install_fallback_subscriber(initialized: &mut InitializedTelemetry) {
261 let filter = EnvFilter::try_from_default_env().unwrap_or_else(|_| EnvFilter::new("info"));
262 if let Err(error) = tracing_subscriber::registry()
263 .with(filter)
264 .with(tracing_subscriber::fmt::layer())
265 .try_init()
266 {
267 initialized.push_error(format!("fallback subscriber: {error}"));
268 }
269}
270
271#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
273pub struct TelemetryStatus {
274 pub otlp_endpoint_configured: bool,
275 pub traces_enabled: bool,
276 pub metrics_enabled: bool,
277 pub trace_exporter_initialized: bool,
278 pub metrics_exporter_initialized: bool,
279 pub export_timeout_ms: u64,
280 pub service_environment: String,
281 #[serde(skip_serializing_if = "Option::is_none")]
282 pub last_error: Option<String>,
283 pub agent_protocol: AgentProtocolMetricsSnapshot,
284}
285
286#[derive(Debug, Clone, Default)]
288pub struct AgentProtocolMetrics {
289 inner: Arc<Mutex<AgentProtocolMetricsSnapshot>>,
290}
291
292impl AgentProtocolMetrics {
293 pub fn record_request(
295 &self,
296 protocol: &str,
297 operation: &str,
298 status: &str,
299 duration_ms: u64,
300 truncated: bool,
301 ) {
302 {
303 let mut inner = self
304 .inner
305 .lock()
306 .unwrap_or_else(|poisoned| poisoned.into_inner());
307 inner.requests_total = inner.requests_total.saturating_add(1);
308 inner.request_duration_ms_total =
309 inner.request_duration_ms_total.saturating_add(duration_ms);
310 if truncated {
311 inner.context_truncated_total = inner.context_truncated_total.saturating_add(1);
312 }
313 }
314
315 let meter = global::meter(PROJECT_NAME);
316 meter
317 .u64_counter("relay_agent_protocol_requests_total")
318 .build()
319 .add(
320 1,
321 &[
322 KeyValue::new("protocol", protocol.to_owned()),
323 KeyValue::new("operation", operation.to_owned()),
324 KeyValue::new("status", status.to_owned()),
325 ],
326 );
327 meter
328 .u64_histogram("relay_agent_protocol_request_duration_ms")
329 .build()
330 .record(
331 duration_ms,
332 &[
333 KeyValue::new("protocol", protocol.to_owned()),
334 KeyValue::new("operation", operation.to_owned()),
335 ],
336 );
337 if truncated {
338 meter
339 .u64_counter("relay_agent_context_truncated_total")
340 .build()
341 .add(
342 1,
343 &[
344 KeyValue::new("protocol", protocol.to_owned()),
345 KeyValue::new("reason", "budget".to_owned()),
346 ],
347 );
348 }
349 }
350
351 pub fn record_rejection(&self, protocol: &str, reason: &str) {
353 {
354 let mut inner = self
355 .inner
356 .lock()
357 .unwrap_or_else(|poisoned| poisoned.into_inner());
358 inner.rejections_total = inner.rejections_total.saturating_add(1);
359 }
360 global::meter(PROJECT_NAME)
361 .u64_counter("relay_agent_protocol_rejections_total")
362 .build()
363 .add(
364 1,
365 &[
366 KeyValue::new("protocol", protocol.to_owned()),
367 KeyValue::new("reason", reason.to_owned()),
368 ],
369 );
370 }
371
372 pub fn record_cancelled(&self, protocol: &str) {
374 {
375 let mut inner = self
376 .inner
377 .lock()
378 .unwrap_or_else(|poisoned| poisoned.into_inner());
379 inner.cancelled_total = inner.cancelled_total.saturating_add(1);
380 }
381 global::meter(PROJECT_NAME)
382 .u64_counter("relay_agent_retrieval_cancelled_total")
383 .build()
384 .add(1, &[KeyValue::new("protocol", protocol.to_owned())]);
385 }
386
387 pub fn snapshot(&self) -> AgentProtocolMetricsSnapshot {
389 self.inner
390 .lock()
391 .unwrap_or_else(|poisoned| poisoned.into_inner())
392 .clone()
393 }
394}
395
396#[derive(Debug, Clone, Default, PartialEq, Eq, Serialize, Deserialize)]
398pub struct AgentProtocolMetricsSnapshot {
399 pub requests_total: u64,
400 pub request_duration_ms_total: u64,
401 pub rejections_total: u64,
402 pub cancelled_total: u64,
403 pub context_truncated_total: u64,
404}
405
406fn signal_endpoint(base: &str, path: &str) -> String {
407 let trimmed = base.trim_end_matches('/');
408 if let Some(prefix) = trimmed.strip_suffix(OTLP_TRACE_PATH) {
409 format!("{prefix}{path}")
410 } else if let Some(prefix) = trimmed.strip_suffix(OTLP_METRIC_PATH) {
411 format!("{prefix}{path}")
412 } else {
413 format!("{trimmed}{path}")
414 }
415}
416
417fn duration_millis(duration: Duration) -> u64 {
418 u64::try_from(duration.as_millis()).unwrap_or(u64::MAX)
419}
420
421#[cfg(test)]
422mod tests {
423 use super::*;
424
425 #[test]
426 fn telemetry_config_applies_documented_defaults() {
427 let config = TelemetryConfig::from_environment(&TelemetryEnvOverrides::default());
428
429 assert_eq!(config.otel_endpoint, DEFAULT_OTEL_ENDPOINT);
430 assert!(!config.traces_enabled);
431 assert!(!config.metrics_enabled);
432 assert_eq!(config.export_timeout, Duration::from_millis(5_000));
433 assert_eq!(config.service_environment, "local");
434 assert_eq!(config.trace_endpoint(), "http://127.0.0.1:4318/v1/traces");
435 assert_eq!(config.metric_endpoint(), "http://127.0.0.1:4318/v1/metrics");
436 }
437
438 #[test]
439 fn telemetry_config_uses_validated_environment_overrides() {
440 let config = TelemetryConfig::from_environment(&TelemetryEnvOverrides {
441 otel_endpoint: Some("http://collector:4318/".to_owned()),
442 otel_traces: Some(true),
443 otel_metrics: Some(true),
444 export_timeout_ms: Some(250),
445 service_environment: Some("ci".to_owned()),
446 });
447
448 assert_eq!(config.otel_endpoint, "http://collector:4318/");
449 assert!(config.traces_enabled);
450 assert!(config.metrics_enabled);
451 assert_eq!(config.export_timeout, Duration::from_millis(250));
452 assert_eq!(config.service_environment, "ci");
453 assert_eq!(config.trace_endpoint(), "http://collector:4318/v1/traces");
454 assert_eq!(config.metric_endpoint(), "http://collector:4318/v1/metrics");
455 }
456
457 #[test]
458 fn signal_endpoint_preserves_signal_specific_paths() {
459 assert_eq!(
460 signal_endpoint("http://collector:4318/v1/traces", OTLP_TRACE_PATH),
461 "http://collector:4318/v1/traces"
462 );
463 assert_eq!(
464 signal_endpoint("http://collector:4318/v1/metrics", OTLP_METRIC_PATH),
465 "http://collector:4318/v1/metrics"
466 );
467 }
468
469 #[test]
470 fn signal_endpoint_routes_sibling_signals_from_specific_paths() {
471 assert_eq!(
472 signal_endpoint("http://collector:4318/v1/traces", OTLP_METRIC_PATH),
473 "http://collector:4318/v1/metrics"
474 );
475 assert_eq!(
476 signal_endpoint("http://collector:4318/v1/metrics", OTLP_TRACE_PATH),
477 "http://collector:4318/v1/traces"
478 );
479 }
480
481 #[test]
482 fn disabled_exporters_still_report_runtime_status() {
483 let runtime =
484 ObservabilityRuntime::new(TelemetryConfig::from_environment(&TelemetryEnvOverrides {
485 otel_endpoint: Some("http://collector:4318".to_owned()),
486 service_environment: Some("test".to_owned()),
487 ..TelemetryEnvOverrides::default()
488 }));
489
490 runtime.initialize();
491 let status = runtime.status();
492
493 assert!(status.otlp_endpoint_configured);
494 assert!(!status.traces_enabled);
495 assert!(!status.metrics_enabled);
496 assert!(!status.trace_exporter_initialized);
497 assert!(!status.metrics_exporter_initialized);
498 assert_eq!(status.export_timeout_ms, 5_000);
499 assert_eq!(status.service_environment, "test");
500 assert_eq!(status.last_error, None);
501 }
502
503 #[test]
504 fn agent_protocol_metrics_snapshot_records_all_event_types() {
505 let metrics = AgentProtocolMetrics::default();
506
507 metrics.record_request("mcp", "tools/call", "ok", 12, false);
508 metrics.record_request("mcp", "resources/read", "ok", 34, true);
509 metrics.record_rejection("mcp", "qos");
510 metrics.record_cancelled("acp");
511
512 assert_eq!(
513 metrics.snapshot(),
514 AgentProtocolMetricsSnapshot {
515 requests_total: 2,
516 request_duration_ms_total: 46,
517 rejections_total: 1,
518 cancelled_total: 1,
519 context_truncated_total: 1,
520 }
521 );
522 }
523
524 #[test]
525 fn agent_protocol_metrics_saturate_instead_of_wrapping() {
526 let metrics = AgentProtocolMetrics::default();
527 {
528 let mut snapshot = metrics.inner.lock().expect("metrics mutex");
529 snapshot.requests_total = u64::MAX;
530 snapshot.request_duration_ms_total = u64::MAX - 1;
531 snapshot.rejections_total = u64::MAX;
532 snapshot.cancelled_total = u64::MAX;
533 snapshot.context_truncated_total = u64::MAX;
534 }
535
536 metrics.record_request("mcp", "tools/call", "ok", 50, true);
537 metrics.record_rejection("mcp", "budget");
538 metrics.record_cancelled("mcp");
539
540 let snapshot = metrics.snapshot();
541 assert_eq!(snapshot.requests_total, u64::MAX);
542 assert_eq!(snapshot.request_duration_ms_total, u64::MAX);
543 assert_eq!(snapshot.rejections_total, u64::MAX);
544 assert_eq!(snapshot.cancelled_total, u64::MAX);
545 assert_eq!(snapshot.context_truncated_total, u64::MAX);
546 }
547
548 #[test]
549 fn duration_millis_saturates_large_durations() {
550 let large = Duration::from_millis(u64::MAX) + Duration::from_millis(1);
551
552 assert_eq!(duration_millis(large), u64::MAX);
553 }
554
555 #[test]
556 fn initialized_telemetry_accumulates_exporter_errors() {
557 let mut initialized = InitializedTelemetry::default();
558
559 initialized.push_error("metrics exporter: invalid endpoint".to_owned());
560 initialized.push_error("trace exporter: invalid endpoint".to_owned());
561
562 assert_eq!(
563 initialized.last_error.as_deref(),
564 Some("metrics exporter: invalid endpoint; trace exporter: invalid endpoint")
565 );
566 }
567
568 #[test]
569 fn observability_state_accumulates_shutdown_errors() {
570 let mut state = ObservabilityState::default();
571
572 state.push_error("trace shutdown: timed out".to_owned());
573 state.push_error("metrics shutdown: already shut down".to_owned());
574
575 assert_eq!(
576 state.last_error.as_deref(),
577 Some("trace shutdown: timed out; metrics shutdown: already shut down")
578 );
579 }
580}