systemprompt_api/services/server/
metrics.rs1use std::sync::{Mutex, OnceLock};
7use std::time::Instant;
8
9use axum::extract::{MatchedPath, Request};
10use axum::http::header::CONTENT_TYPE;
11use axum::middleware::Next;
12use axum::response::{IntoResponse, Response};
13use metrics_exporter_prometheus::{Matcher, PrometheusBuilder, PrometheusHandle};
14use systemprompt_events::{
15 A2A_BROADCASTER, AGUI_BROADCASTER, ANALYTICS_BROADCASTER, Broadcaster, CONTEXT_BROADCASTER,
16};
17use systemprompt_identifiers::InstanceId;
18use systemprompt_traits::OwnedTask;
19
20const METRICS_CONTENT_TYPE: &str = "text/plain; version=0.0.4; charset=utf-8";
21
22const HTTP_REQUESTS_TOTAL: &str = "http_requests_total";
23pub const HTTP_REQUEST_DURATION_SECONDS: &str = "http_request_duration_seconds";
24pub const HTTP_BUCKETS: [f64; 11] = [
25 0.005, 0.01, 0.025, 0.05, 0.1, 0.25, 0.5, 1.0, 2.5, 5.0, 10.0,
26];
27const HTTP_REQUESTS_IN_FLIGHT: &str = "http_requests_in_flight";
28const SSE_CONNECTIONS: &str = "sse_active_connections";
29
30static RECORDER: OnceLock<PrometheusHandle> = OnceLock::new();
33static RECORDER_INIT: Mutex<()> = Mutex::new(());
34
35pub fn install_recorder(instance_id: &InstanceId) -> anyhow::Result<PrometheusHandle> {
36 if let Some(handle) = RECORDER.get() {
37 return Ok(handle.clone());
38 }
39 let _guard = RECORDER_INIT
40 .lock()
41 .unwrap_or_else(std::sync::PoisonError::into_inner);
42 if let Some(handle) = RECORDER.get() {
43 return Ok(handle.clone());
44 }
45 let handle = PrometheusBuilder::new()
46 .add_global_label("instance", instance_id.as_str())
47 .set_buckets_for_metric(
48 Matcher::Full(systemprompt_gateway::GATEWAY_OVERHEAD_SECONDS.to_owned()),
49 &systemprompt_gateway::OVERHEAD_BUCKETS,
50 )?
51 .set_buckets_for_metric(
52 Matcher::Full(systemprompt_gateway::GATEWAY_UPSTREAM_DURATION_SECONDS.to_owned()),
53 &systemprompt_gateway::UPSTREAM_BUCKETS,
54 )?
55 .set_buckets_for_metric(
56 Matcher::Full(HTTP_REQUEST_DURATION_SECONDS.to_owned()),
57 &HTTP_BUCKETS,
58 )?
59 .install_recorder()
60 .map_err(|e| anyhow::anyhow!("failed to install Prometheus recorder: {e}"))?;
61 describe_metrics();
62 Ok(RECORDER.get_or_init(|| handle).clone())
63}
64
65fn describe_metrics() {
66 use crate::services::middleware::load_shed::{
67 HTTP_IN_FLIGHT_LIMIT, HTTP_IN_FLIGHT_SATURATION, HTTP_LOAD_SHED_TOTAL,
68 };
69 metrics::describe_counter!(
70 HTTP_LOAD_SHED_TOTAL,
71 "Requests refused with 503 because the in-flight ceiling was reached"
72 );
73 metrics::describe_gauge!(
74 HTTP_IN_FLIGHT_LIMIT,
75 "Configured server.max_in_flight ceiling"
76 );
77 metrics::describe_gauge!(
78 HTTP_IN_FLIGHT_SATURATION,
79 "Fraction of the in-flight ceiling in use (1.0 = shedding)"
80 );
81 metrics::describe_histogram!(
82 HTTP_REQUEST_DURATION_SECONDS,
83 metrics::Unit::Seconds,
84 "HTTP request duration by method, matched path and status"
85 );
86 metrics::describe_histogram!(
87 systemprompt_gateway::GATEWAY_OVERHEAD_SECONDS,
88 metrics::Unit::Seconds,
89 "Time the gateway added to a completed inference request, excluding the upstream call"
90 );
91 metrics::describe_histogram!(
92 systemprompt_gateway::GATEWAY_UPSTREAM_DURATION_SECONDS,
93 metrics::Unit::Seconds,
94 "Upstream provider call duration for a completed inference request"
95 );
96}
97
98pub fn metrics_router(handle: PrometheusHandle) -> axum::Router {
99 axum::Router::new()
100 .route("/metrics", axum::routing::get(handle_metrics))
101 .with_state(handle)
102}
103
104pub async fn serve_metrics_listener(
105 addr: std::net::SocketAddr,
106 handle: PrometheusHandle,
107) -> anyhow::Result<OwnedTask<()>> {
108 let listener = tokio::net::TcpListener::bind(addr).await?;
109 let mut readiness = super::readiness::get_readiness_receiver();
110 let shutdown = async move {
111 loop {
112 match readiness.recv().await {
113 Ok(super::readiness::ReadinessEvent::ApiShuttingDown) | Err(_) => break,
114 Ok(_) => {},
115 }
116 }
117 };
118 tracing::info!(%addr, "metrics listener bound");
119 Ok(OwnedTask::spawn("metrics_listener", async move {
120 if let Err(error) = axum::serve(listener, metrics_router(handle))
121 .with_graceful_shutdown(shutdown)
122 .await
123 {
124 tracing::warn!(error = %error, "metrics listener exited with error");
125 }
126 }))
127}
128
129pub async fn handle_metrics(
130 axum::extract::State(handle): axum::extract::State<PrometheusHandle>,
131) -> Response {
132 refresh_connection_gauges().await;
133 let body = handle.render();
134 ([(CONTENT_TYPE, METRICS_CONTENT_TYPE)], body).into_response()
135}
136
137async fn refresh_connection_gauges() {
138 let context = CONTEXT_BROADCASTER.total_connections().await;
139 let agui = AGUI_BROADCASTER.total_connections().await;
140 let a2a = A2A_BROADCASTER.total_connections().await;
141 let analytics = ANALYTICS_BROADCASTER.total_connections().await;
142
143 metrics::gauge!(SSE_CONNECTIONS, "channel" => "context").set(context as f64);
144 metrics::gauge!(SSE_CONNECTIONS, "channel" => "agui").set(agui as f64);
145 metrics::gauge!(SSE_CONNECTIONS, "channel" => "a2a").set(a2a as f64);
146 metrics::gauge!(SSE_CONNECTIONS, "channel" => "analytics").set(analytics as f64);
147}
148
149pub async fn track_metrics(req: Request, next: Next) -> Response {
150 let method = req.method().clone();
151 let path = req
152 .extensions()
153 .get::<MatchedPath>()
154 .map_or_else(|| req.uri().path().to_owned(), |m| m.as_str().to_owned());
155
156 let in_flight = metrics::gauge!(HTTP_REQUESTS_IN_FLIGHT);
157 in_flight.increment(1.0);
158
159 let start = Instant::now();
160 let response = next.run(req).await;
161 let elapsed = start.elapsed().as_secs_f64();
162
163 in_flight.decrement(1.0);
164
165 let status = response.status().as_u16().to_string();
166 let method = method.to_string();
167
168 metrics::counter!(
169 HTTP_REQUESTS_TOTAL,
170 "method" => method.clone(),
171 "path" => path.clone(),
172 "status" => status.clone(),
173 )
174 .increment(1);
175 metrics::histogram!(
176 HTTP_REQUEST_DURATION_SECONDS,
177 "method" => method,
178 "path" => path,
179 "status" => status,
180 )
181 .record(elapsed);
182
183 response
184}