Skip to main content

synapse/
telemetry.rs

1//! OpenTelemetry metrics. Every instrument the gateway emits lives on one
2//! [`GatewayMetrics`], passed explicitly (there is no global recorder). The
3//! Prometheus and OTLP exporters sit behind the `server` feature.
4
5use std::sync::Arc;
6
7use opentelemetry::metrics::{
8    Counter, Gauge, Histogram, Meter, MeterProvider as _, NoopMeterProvider,
9};
10use opentelemetry::KeyValue;
11
12use crate::observability::GenAiSpan;
13use crate::routing::classify::Lane;
14
15#[cfg(feature = "server")]
16pub use exporter::{install, metrics_router, scrape, test_metrics, MetricsExporter};
17
18/// Second-scaled histogram boundaries; OTel's defaults are millisecond-scaled.
19pub const SECONDS_BUCKETS: [f64; 14] = [
20    0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0, 30.0, 60.0, 120.0,
21];
22
23/// A duration histogram in seconds with [`SECONDS_BUCKETS`] boundaries.
24pub fn seconds_histogram(meter: &Meter, name: &'static str) -> Histogram<f64> {
25    meter
26        .f64_histogram(name)
27        .with_boundaries(SECONDS_BUCKETS.to_vec())
28        .build()
29}
30
31/// Every instrument the gateway records. Build once, share behind an `Arc`.
32pub struct GatewayMetrics {
33    requests: Counter<u64>,
34    request_duration: Histogram<f64>,
35    input_tokens: Counter<u64>,
36    output_tokens: Counter<u64>,
37    embeddings: Counter<u64>,
38    embedding_duration: Histogram<f64>,
39    passthrough: Counter<u64>,
40    passthrough_fallback: Counter<u64>,
41    jev_extraction: Counter<u64>,
42    ledger_errors: Counter<u64>,
43    ledger_dropped: Counter<u64>,
44    retry_attempts: Counter<u64>,
45    resilience_calls: Counter<u64>,
46    resilience_call_duration: Histogram<f64>,
47    breaker_transitions: Counter<u64>,
48    breaker_state: Gauge<f64>,
49    guard_scans: Counter<u64>,
50    guard_matches: Counter<u64>,
51    guard_scan_duration: Histogram<f64>,
52    routing_decisions: Counter<u64>,
53    routing_decision_duration: Histogram<f64>,
54}
55
56impl Default for GatewayMetrics {
57    fn default() -> Self {
58        Self::new(&NoopMeterProvider::new().meter("synapse-gateway"))
59    }
60}
61
62impl GatewayMetrics {
63    pub fn new(meter: &Meter) -> Self {
64        Self {
65            requests: meter.u64_counter("synapse_requests_total").build(),
66            request_duration: seconds_histogram(meter, "synapse_request_duration_seconds"),
67            input_tokens: meter.u64_counter("synapse_input_tokens_total").build(),
68            output_tokens: meter.u64_counter("synapse_output_tokens_total").build(),
69            embeddings: meter.u64_counter("synapse_embeddings_total").build(),
70            embedding_duration: seconds_histogram(meter, "synapse_embedding_duration_seconds"),
71            passthrough: meter.u64_counter("synapse_passthrough_total").build(),
72            passthrough_fallback: meter
73                .u64_counter("synapse_passthrough_fallback_total")
74                .build(),
75            jev_extraction: meter.u64_counter("synapse_jev_extraction_total").build(),
76            ledger_errors: meter.u64_counter("synapse_ledger_errors_total").build(),
77            ledger_dropped: meter.u64_counter("synapse_ledger_dropped_total").build(),
78            retry_attempts: meter
79                .u64_counter("synapse_resilience_retry_attempts_total")
80                .build(),
81            resilience_calls: meter.u64_counter("synapse_resilience_calls_total").build(),
82            resilience_call_duration: seconds_histogram(
83                meter,
84                "synapse_resilience_call_duration_seconds",
85            ),
86            breaker_transitions: meter
87                .u64_counter("synapse_resilience_breaker_transitions_total")
88                .build(),
89            breaker_state: meter.f64_gauge("synapse_resilience_breaker_state").build(),
90            guard_scans: meter.u64_counter("synapse_guard_scans_total").build(),
91            guard_matches: meter.u64_counter("synapse_guard_matches_total").build(),
92            guard_scan_duration: seconds_histogram(meter, "synapse_guard_scan_duration_seconds"),
93            routing_decisions: meter.u64_counter("synapse_routing_decisions_total").build(),
94            routing_decision_duration: seconds_histogram(
95                meter,
96                "synapse_routing_decision_duration_seconds",
97            ),
98        }
99    }
100
101    /// Instruments that record nothing, for embedders and tests that never
102    /// read metrics.
103    pub fn noop() -> Arc<Self> {
104        Arc::new(Self::default())
105    }
106
107    /// One served chat request: count, latency, and token totals.
108    pub fn request(&self, span: &GenAiSpan, latency_secs: f64) {
109        let labels = [
110            KeyValue::new("route", span.route.clone()),
111            KeyValue::new("model", span.response_model.clone()),
112            KeyValue::new("system", span.system),
113            KeyValue::new("lane", lane_label(&span.lane)),
114        ];
115        self.requests.add(1, &labels);
116        self.request_duration.record(latency_secs, &labels);
117        self.input_tokens.add(span.input_tokens, &labels);
118        self.output_tokens.add(span.output_tokens, &labels);
119    }
120
121    pub fn embedding(&self, route: &str, model: &str, provider: &str, secs: f64) {
122        let labels = [
123            KeyValue::new("route", route.to_string()),
124            KeyValue::new("model", model.to_string()),
125            KeyValue::new("provider", provider.to_string()),
126        ];
127        self.embeddings.add(1, &labels);
128        self.embedding_duration.record(secs, &labels);
129    }
130
131    pub fn passthrough(&self, provider: &'static str, model: &str, action: &str, ok: bool) {
132        self.passthrough.add(
133            1,
134            &[
135                KeyValue::new("provider", provider),
136                KeyValue::new("model", model.to_string()),
137                KeyValue::new("action", action.to_string()),
138                KeyValue::new("status", if ok { "ok" } else { "error" }),
139            ],
140        );
141    }
142
143    pub fn passthrough_fallback(&self, from_model: &str, to_model: &str) {
144        self.passthrough_fallback.add(
145            1,
146            &[
147                KeyValue::new("from_model", from_model.to_string()),
148                KeyValue::new("to_model", to_model.to_string()),
149            ],
150        );
151    }
152
153    pub fn jev_extraction(&self, route: &str, degraded: bool) {
154        self.jev_extraction.add(
155            1,
156            &[
157                KeyValue::new("route", route.to_string()),
158                KeyValue::new("degraded", degraded.to_string()),
159            ],
160        );
161    }
162
163    pub fn ledger_error(&self, backend: &'static str) {
164        self.ledger_errors
165            .add(1, &[KeyValue::new("backend", backend)]);
166    }
167
168    pub fn ledger_dropped(&self) {
169        self.ledger_dropped.add(1, &[]);
170    }
171
172    pub fn retry_attempt(&self, label: &'static str) {
173        self.retry_attempts.add(1, &[KeyValue::new("label", label)]);
174    }
175
176    pub fn resilience_call(&self, label: &'static str, outcome: &'static str, secs: f64) {
177        let labels = [
178            KeyValue::new("label", label),
179            KeyValue::new("outcome", outcome),
180        ];
181        self.resilience_calls.add(1, &labels);
182        self.resilience_call_duration.record(secs, &labels);
183    }
184
185    /// A breaker transition plus its new state (0 closed, 1 open, 2 half-open).
186    pub fn breaker_transition(&self, name: &'static str, transition: &'static str, state: u8) {
187        self.breaker_transitions.add(
188            1,
189            &[
190                KeyValue::new("name", name),
191                KeyValue::new("transition", transition),
192            ],
193        );
194        self.breaker_state
195            .record(f64::from(state), &[KeyValue::new("name", name)]);
196    }
197
198    pub fn guard_scan(&self, policy: &str, outcome: &'static str, secs: f64) {
199        self.guard_scans.add(
200            1,
201            &[
202                KeyValue::new("policy", policy.to_string()),
203                KeyValue::new("outcome", outcome),
204            ],
205        );
206        self.guard_scan_duration
207            .record(secs, &[KeyValue::new("policy", policy.to_string())]);
208    }
209
210    pub fn guard_match(&self, policy: &str, scanner: &str, severity: &'static str) {
211        self.guard_matches.add(
212            1,
213            &[
214                KeyValue::new("policy", policy.to_string()),
215                KeyValue::new("scanner", scanner.to_string()),
216                KeyValue::new("severity", severity),
217            ],
218        );
219    }
220
221    /// One per request to a `jev` route; `tier` is the decided tier.
222    pub fn routing_decision(&self, route: &str, tier: &str, outcome: &'static str) {
223        self.routing_decisions.add(
224            1,
225            &[
226                KeyValue::new("route", route.to_string()),
227                KeyValue::new("tier", tier.to_string()),
228                KeyValue::new("outcome", outcome),
229            ],
230        );
231    }
232
233    /// Latency of one Jev decision call.
234    pub fn routing_decision_duration(&self, route: &str, secs: f64) {
235        self.routing_decision_duration
236            .record(secs, &[KeyValue::new("route", route.to_string())]);
237    }
238}
239
240fn lane_label(lane: &Lane) -> &'static str {
241    match lane {
242        Lane::Standard => "standard",
243        Lane::NativeVertex => "native",
244        Lane::Jev => "jev",
245    }
246}
247
248#[cfg(feature = "server")]
249mod exporter {
250    use std::sync::Arc;
251
252    use axum::extract::State;
253    use axum::http::{header, StatusCode};
254    use axum::response::{IntoResponse, Response};
255    use axum::routing::get;
256    use axum::Router;
257    use opentelemetry::metrics::MeterProvider as _;
258    use opentelemetry_otlp::WithExportConfig;
259    use opentelemetry_sdk::metrics::{PeriodicReader, SdkMeterProvider};
260    use opentelemetry_sdk::Resource;
261    use prometheus::{Registry, TextEncoder, TEXT_FORMAT};
262    use tap::Pipe;
263
264    use super::GatewayMetrics;
265
266    /// The Prometheus registry plus the provider feeding it. Keep it alive for
267    /// the process lifetime: dropping the last provider handle shuts every
268    /// reader down, and scrapes come back empty.
269    #[derive(Clone, Debug)]
270    pub struct MetricsExporter {
271        pub registry: Registry,
272        _provider: SdkMeterProvider,
273    }
274
275    /// Build the meter provider and the gateway's instruments. Prometheus is
276    /// always on; OTLP/HTTP is added when `otlp_endpoint` (a base collector URL
277    /// such as `http://collector:4318`) is set. opentelemetry-otlp 0.32 uses
278    /// programmatic endpoints verbatim, so `/v1/metrics` is appended here.
279    pub fn install(
280        otlp_endpoint: Option<&str>,
281        service_name: &str,
282    ) -> anyhow::Result<(Arc<GatewayMetrics>, MetricsExporter)> {
283        let registry = Registry::new();
284        let prometheus = opentelemetry_prometheus::exporter()
285            .with_registry(registry.clone())
286            .without_scope_info()
287            .without_target_info()
288            .without_counter_suffixes()
289            .without_units()
290            .build()?;
291        let otlp = otlp_endpoint
292            .map(|endpoint| {
293                opentelemetry_otlp::MetricExporter::builder()
294                    .with_http()
295                    .with_endpoint(format!("{}/v1/metrics", endpoint.trim_end_matches('/')))
296                    .build()
297            })
298            .transpose()?;
299        let provider = SdkMeterProvider::builder()
300            .with_reader(prometheus)
301            .with_resource(
302                Resource::builder()
303                    .with_service_name(service_name.to_string())
304                    .build(),
305            )
306            .pipe(|builder| match otlp {
307                Some(exporter) => builder.with_reader(PeriodicReader::builder(exporter).build()),
308                None => builder,
309            })
310            .build();
311        let metrics = Arc::new(GatewayMetrics::new(&provider.meter("synapse-gateway")));
312        Ok((
313            metrics,
314            MetricsExporter {
315                registry,
316                _provider: provider,
317            },
318        ))
319    }
320
321    /// A private Prometheus-only exporter for tests.
322    #[doc(hidden)]
323    pub fn test_metrics() -> (Arc<GatewayMetrics>, MetricsExporter) {
324        install(None, "synapse-gateway-test").expect("prometheus exporter builds")
325    }
326
327    /// Prometheus text exposition of everything recorded so far.
328    #[doc(hidden)]
329    pub fn scrape(exporter: &MetricsExporter) -> String {
330        encode(exporter).unwrap_or_default()
331    }
332
333    fn encode(exporter: &MetricsExporter) -> prometheus::Result<String> {
334        TextEncoder::new().encode_to_string(&exporter.registry.gather())
335    }
336
337    /// `GET /metrics` (and `GET /`) in Prometheus text format.
338    pub fn metrics_router(exporter: MetricsExporter) -> Router {
339        Router::new()
340            .route("/metrics", get(serve))
341            .route("/", get(serve))
342            .with_state(exporter)
343    }
344
345    async fn serve(State(exporter): State<MetricsExporter>) -> Response {
346        match encode(&exporter) {
347            Ok(text) => ([(header::CONTENT_TYPE, TEXT_FORMAT)], text).into_response(),
348            Err(e) => {
349                tracing::error!(error = %e, "prometheus metrics encoding failed");
350                StatusCode::INTERNAL_SERVER_ERROR.into_response()
351            }
352        }
353    }
354}
355
356#[cfg(test)]
357mod tests {
358    use super::*;
359    use crate::observability::GenAiSpan;
360    use crate::routing::classify::Lane;
361    use crate::routing::executor::Completion;
362    use crate::routing::stream::FinishReason;
363
364    fn span() -> GenAiSpan {
365        let c = Completion {
366            provider: "qwen".into(),
367            model: "qwen-max".into(),
368            content: String::new(),
369            tool_calls: Vec::new(),
370            finish_reason: FinishReason::Stop,
371            input_tokens: 3,
372            output_tokens: 5,
373        };
374        GenAiSpan::from_completion(&c, Lane::Standard, "fast", "acme", None, 1, false)
375    }
376
377    #[test]
378    fn noop_accepts_every_measurement() {
379        let m = GatewayMetrics::noop();
380        m.request(&span(), 0.1);
381        m.embedding("embed", "text-embedding-3-small", "openai", 0.05);
382        m.passthrough("vertex", "gemini-2.5-flash", "countTokens", true);
383        m.passthrough_fallback("gemini-2.5-pro", "gemini-2.5-flash");
384        m.jev_extraction("extract", false);
385        m.ledger_error("writer");
386        m.ledger_dropped();
387        m.retry_attempt("qwen");
388        m.resilience_call("qwen", "success", 0.3);
389        m.breaker_transition("qwen", "open", 1);
390        m.guard_scan("strict", "block", 0.001);
391        m.guard_match("strict", "ban_substrings", "block");
392        m.routing_decision("auto", "hard", "decided");
393        m.routing_decision_duration("auto", 0.12);
394    }
395
396    #[cfg(feature = "server")]
397    #[test]
398    fn routing_decision_instruments_export() {
399        let (m, exporter) = test_metrics();
400        m.routing_decision("auto", "hard", "decided");
401        m.routing_decision_duration("auto", 0.12);
402        let text = scrape(&exporter);
403        assert!(text.contains("synapse_routing_decisions_total"), "{text}");
404        ["route=\"auto\"", "tier=\"hard\"", "outcome=\"decided\""]
405            .iter()
406            .for_each(|label| assert!(text.contains(label), "{label} missing in {text}"));
407        assert!(
408            !text.contains("synapse_routing_decisions_total_total"),
409            "{text}"
410        );
411        assert!(
412            text.contains("synapse_routing_decision_duration_seconds_bucket"),
413            "{text}"
414        );
415        assert!(text.contains("le=\"0.25\""), "{text}");
416    }
417
418    #[cfg(feature = "server")]
419    #[test]
420    fn every_instrument_exports_with_its_labels() {
421        let (m, exporter) = test_metrics();
422        m.request(&span(), 0.2);
423        m.embedding("embed", "text-embedding-3-small", "openai", 0.05);
424        m.passthrough("vertex", "gemini-2.5-flash", "countTokens", true);
425        m.passthrough_fallback("gemini-2.5-pro", "gemini-2.5-flash");
426        m.jev_extraction("extract", false);
427        m.ledger_error("writer");
428        m.ledger_dropped();
429        m.retry_attempt("qwen");
430        m.resilience_call("qwen", "success", 0.3);
431        m.breaker_transition("qwen", "open", 1);
432        m.guard_scan("strict", "block", 0.001);
433        m.guard_match("strict", "ban_substrings", "block");
434        m.routing_decision("auto", "hard", "decided");
435        m.routing_decision_duration("auto", 0.12);
436        let text = scrape(&exporter);
437        for line in [
438            r#"synapse_requests_total{lane="standard",model="qwen-max",route="fast",system="dashscope"} 1"#,
439            r#"synapse_input_tokens_total{lane="standard",model="qwen-max",route="fast",system="dashscope"} 3"#,
440            r#"synapse_output_tokens_total{lane="standard",model="qwen-max",route="fast",system="dashscope"} 5"#,
441            r#"synapse_request_duration_seconds_count{lane="standard",model="qwen-max",route="fast",system="dashscope"} 1"#,
442            r#"synapse_embeddings_total{model="text-embedding-3-small",provider="openai",route="embed"} 1"#,
443            r#"synapse_embedding_duration_seconds_count{model="text-embedding-3-small",provider="openai",route="embed"} 1"#,
444            r#"synapse_passthrough_total{action="countTokens",model="gemini-2.5-flash",provider="vertex",status="ok"} 1"#,
445            r#"synapse_passthrough_fallback_total{from_model="gemini-2.5-pro",to_model="gemini-2.5-flash"} 1"#,
446            r#"synapse_jev_extraction_total{degraded="false",route="extract"} 1"#,
447            r#"synapse_ledger_errors_total{backend="writer"} 1"#,
448            "synapse_ledger_dropped_total 1",
449            r#"synapse_resilience_retry_attempts_total{label="qwen"} 1"#,
450            r#"synapse_resilience_calls_total{label="qwen",outcome="success"} 1"#,
451            r#"synapse_resilience_call_duration_seconds_count{label="qwen",outcome="success"} 1"#,
452            r#"synapse_resilience_breaker_transitions_total{name="qwen",transition="open"} 1"#,
453            r#"synapse_resilience_breaker_state{name="qwen"} 1"#,
454            r#"synapse_guard_scans_total{outcome="block",policy="strict"} 1"#,
455            r#"synapse_guard_matches_total{policy="strict",scanner="ban_substrings",severity="block"} 1"#,
456            r#"synapse_guard_scan_duration_seconds_count{policy="strict"} 1"#,
457            r#"synapse_routing_decisions_total{outcome="decided",route="auto",tier="hard"} 1"#,
458            r#"synapse_routing_decision_duration_seconds_count{route="auto"} 1"#,
459        ] {
460            assert!(
461                text.lines().any(|l| l == line),
462                "missing `{line}` in:\n{text}"
463            );
464        }
465    }
466
467    #[cfg(feature = "server")]
468    #[test]
469    fn failed_passthrough_and_latest_breaker_state_export() {
470        let (m, exporter) = test_metrics();
471        m.passthrough("vertex", "gemini-2.5-flash", "countTokens", false);
472        m.breaker_transition("qwen", "open", 1);
473        m.breaker_transition("qwen", "half_open", 2);
474        let text = scrape(&exporter);
475        for line in [
476            r#"synapse_passthrough_total{action="countTokens",model="gemini-2.5-flash",provider="vertex",status="error"} 1"#,
477            r#"synapse_resilience_breaker_state{name="qwen"} 2"#,
478        ] {
479            assert!(
480                text.lines().any(|l| l == line),
481                "missing `{line}` in:\n{text}"
482            );
483        }
484    }
485
486    #[cfg(feature = "server")]
487    #[test]
488    fn exposition_has_seconds_buckets_and_no_otel_artifacts() {
489        let (m, exporter) = test_metrics();
490        m.request(&span(), 0.2);
491        let text = scrape(&exporter);
492        assert!(text.contains(
493            r#"synapse_request_duration_seconds_bucket{lane="standard",model="qwen-max",route="fast",system="dashscope",le="0.25"} 1"#
494        ));
495        assert!(!text.contains("otel_scope"));
496        assert!(!text.contains("target_info"));
497        assert!(!text.contains("_total_total"));
498    }
499
500    /// Accept one OTLP request, answer 200, and return its request line.
501    #[cfg(feature = "server")]
502    fn serve_one_otlp_request(listener: std::net::TcpListener) -> std::io::Result<String> {
503        use std::io::{BufRead, BufReader, Read, Write};
504
505        let (mut stream, _) = listener.accept()?;
506        stream.set_read_timeout(Some(std::time::Duration::from_secs(5)))?;
507        let mut reader = BufReader::new(stream.try_clone()?);
508        let head: Vec<String> = std::iter::from_fn(|| {
509            let mut line = String::new();
510            match reader.read_line(&mut line) {
511                Ok(n) if n > 0 && line != "\r\n" => Some(line),
512                _ => None,
513            }
514        })
515        .collect();
516        let body_len = head
517            .iter()
518            .find_map(|l| {
519                l.to_ascii_lowercase()
520                    .strip_prefix("content-length:")
521                    .and_then(|v| v.trim().parse::<usize>().ok())
522            })
523            .unwrap_or(0);
524        reader.read_exact(&mut vec![0; body_len])?;
525        stream.write_all(b"HTTP/1.1 200 OK\r\ncontent-length: 0\r\n\r\n")?;
526        Ok(head.into_iter().next().unwrap_or_default())
527    }
528
529    /// Install with a local collector, check Prometheus still scrapes, then
530    /// drop everything so shutdown flushes one OTLP export to the collector.
531    #[cfg(feature = "server")]
532    fn assert_otlp_exports_alongside_prometheus() {
533        use std::sync::mpsc;
534        use std::time::Duration;
535
536        let listener = std::net::TcpListener::bind("127.0.0.1:0").unwrap();
537        let addr = listener.local_addr().unwrap();
538        let (tx, rx) = mpsc::channel();
539        std::thread::spawn(move || tx.send(serve_one_otlp_request(listener)));
540
541        let (m, exporter) = install(Some(&format!("http://{addr}/")), "synapse-gateway").unwrap();
542        m.ledger_dropped();
543        assert!(scrape(&exporter).contains("synapse_ledger_dropped_total 1"));
544        drop(m);
545        drop(exporter);
546
547        let request_line = rx
548            .recv_timeout(Duration::from_secs(5))
549            .expect("collector received no OTLP export")
550            .unwrap();
551        assert!(
552            request_line.starts_with("POST /v1/metrics"),
553            "unexpected request line: {request_line}"
554        );
555    }
556
557    #[cfg(feature = "server")]
558    #[test]
559    fn otlp_and_prometheus_readers_coexist() {
560        assert_otlp_exports_alongside_prometheus();
561    }
562
563    #[cfg(feature = "server")]
564    #[tokio::test(flavor = "multi_thread")]
565    async fn otlp_exports_when_installed_inside_a_tokio_runtime() {
566        assert_otlp_exports_alongside_prometheus();
567    }
568
569    #[cfg(feature = "server")]
570    #[tokio::test]
571    async fn metrics_router_serves_prometheus_text_on_metrics_and_root() {
572        use axum::body::Body;
573        use axum::http::{Request, StatusCode};
574        use http_body_util::BodyExt;
575        use tower::ServiceExt;
576
577        let (m, exporter) = test_metrics();
578        m.ledger_dropped();
579        for uri in ["/metrics", "/"] {
580            let resp = metrics_router(exporter.clone())
581                .oneshot(Request::get(uri).body(Body::empty()).unwrap())
582                .await
583                .unwrap();
584            assert_eq!(resp.status(), StatusCode::OK);
585            assert_eq!(resp.headers()["content-type"], "text/plain; version=0.0.4");
586            let body = resp.into_body().collect().await.unwrap().to_bytes();
587            assert!(String::from_utf8_lossy(&body).contains("synapse_ledger_dropped_total 1"));
588        }
589    }
590}