Skip to main content

systemprompt_api/services/server/
metrics.rs

1//! Prometheus recorder installation and HTTP metrics rendering.
2//!
3//! Copyright (c) systemprompt.io — Business Source License 1.1.
4//! See <https://systemprompt.io> for licensing details.
5
6use 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
30// Why: metrics::set_global_recorder rejects a second installation in the same
31// process.
32static 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}