1use 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
18pub 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
23pub 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
31pub 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 pub fn noop() -> Arc<Self> {
104 Arc::new(Self::default())
105 }
106
107 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 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 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 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 #[derive(Clone, Debug)]
270 pub struct MetricsExporter {
271 pub registry: Registry,
272 _provider: SdkMeterProvider,
273 }
274
275 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 #[doc(hidden)]
323 pub fn test_metrics() -> (Arc<GatewayMetrics>, MetricsExporter) {
324 install(None, "synapse-gateway-test").expect("prometheus exporter builds")
325 }
326
327 #[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 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 #[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 #[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}