mz_orchestratord/
metrics.rs1use std::collections::BTreeMap;
11use std::sync::Arc;
12use std::time::Duration;
13
14use axum::{Extension, Router, body::Body, routing::get};
15use http::{HeaderMap, Method, Request, Response, StatusCode};
16use k8s_controller::PrometheusMetrics;
17use prometheus::{Encoder, TextEncoder};
18use tower_http::{classify::ServerErrorsFailureClass, trace::TraceLayer};
19use tracing::{Level, Span};
20
21use mz_ore::metric;
22use mz_ore::metrics::{MetricsRegistry, UIntGauge};
23
24#[derive(Debug)]
25pub struct Metrics {
26 pub is_leader: UIntGauge,
27 pub environmentd_needs_update: UIntGauge,
28 pub reconcile: Arc<PrometheusMetrics>,
32}
33
34impl Metrics {
35 pub fn register_into(registry: &MetricsRegistry) -> Self {
36 Self {
37 is_leader: registry.register(
38 metric! {
39 name: "orchestratord_is_leader",
40 help: "Whether this operator replica holds the controller leadership lease, and is therefore the replica reconciling. Summed across the replicas this should be 1. A sustained 0 means no replica can take the lease, for instance because the service account lacks permission on leases, and the operator is reconciling nothing.",
41 }),
42 environmentd_needs_update: registry.register(
43 metric! {
44 name: "environmentd_needs_update",
45 help: "Count of organizations in this cluster which are running outdated pod templates. Only the operator replica holding the leadership lease reconciles, so the others report zero.",
46 }),
47 reconcile: {
48 let reconcile = PrometheusMetrics::new("orchestratord")
49 .expect("valid reconciliation metric definitions");
50 registry.register_collector(reconcile.clone());
51 Arc::new(reconcile)
52 },
53 }
54 }
55
56 pub fn leadership_acquired(&self) {
58 self.is_leader.set(1);
59 }
60
61 pub fn leadership_lost(&self) {
70 self.is_leader.set(0);
71 self.environmentd_needs_update.set(0);
72 }
73}
74
75pub fn router(registry: MetricsRegistry) -> Router {
76 add_tracing_layer(
77 Router::new()
78 .route("/metrics", get(metrics))
79 .layer(Extension(registry)),
80 )
81}
82
83#[allow(clippy::unused_async)]
84async fn metrics(Extension(registry): Extension<MetricsRegistry>) -> (StatusCode, Vec<u8>) {
85 let mut buf = vec![];
86 let encoder = TextEncoder::new();
87 let metric_families = registry.gather();
88 encoder.encode(&metric_families, &mut buf).unwrap();
89 (StatusCode::OK, buf)
90}
91
92fn add_tracing_layer<S>(router: Router<S>) -> Router<S>
105where
106 S: Clone + Send + Sync + 'static,
107{
108 router.layer(
109 TraceLayer::new_for_http()
110 .make_span_with(|request: &Request<Body>| {
111 macro_rules! make_span {
117 ($level:expr) => {
118 tracing::span!(
119 $level,
120 "HTTP request",
121 "request.uri" = %request.uri(),
122 "request.version" = ?request.version(),
123 "request.method" = %request.method(),
124 "request.headers" = tracing::field::Empty,
125 "response.status" = tracing::field::Empty,
126 "response.status_code" = tracing::field::Empty,
127 "response.headers" = tracing::field::Empty,
128 )
129 }
130 }
131 let span = if ["/api/health", "/metrics"].contains(&request.uri().path())
132 || request.method() == Method::OPTIONS
133 {
134 make_span!(Level::DEBUG)
135 } else {
136 make_span!(Level::INFO)
137 };
138
139 if let Ok(s) = serde_json::to_string(&display_headers(request.headers().clone())) {
140 span.record("request.headers", s);
141 }
142
143 span
144 })
145 .on_response(|response: &Response<Body>, _latency, span: &Span| {
146 span.record(
147 "response.status",
148 &tracing::field::display(response.status()),
149 );
150 span.record("response.status_code", response.status().as_u16());
151 if let Ok(s) = serde_json::to_string(&display_headers(response.headers().clone())) {
152 span.record("response.headers", s);
153 }
154
155 if span
159 .metadata()
160 .and_then(|m| Some(m.level()))
161 .unwrap_or(&Level::DEBUG)
162 == &Level::DEBUG
163 {
164 tracing::debug!("HTTP response generated");
165 } else {
166 tracing::info!("HTTP response generated");
167 }
168 })
169 .on_failure(
170 |error: ServerErrorsFailureClass, _latency: Duration, _span: &Span| {
171 tracing::warn!(error = ?error, "HTTP request handling error");
172 },
173 ),
174 )
175}
176
177fn display_headers(mut headers: HeaderMap) -> BTreeMap<String, String> {
178 _ = headers.remove(http::header::AUTHORIZATION);
180
181 headers
182 .into_iter()
183 .filter_map(|(k, v)| {
184 k.map(|k| {
185 (
186 k.to_string(),
187 String::from_utf8_lossy(v.as_bytes()).to_string(),
188 )
189 })
190 })
191 .collect()
192}