Skip to main content

mz_orchestratord/
metrics.rs

1// Copyright Materialize, Inc. and contributors. All rights reserved.
2//
3// Use of this software is governed by the Business Source License
4// included in the LICENSE file.
5//
6// As of the Change Date specified in that file, in accordance with
7// the Business Source License, use of this software will be governed
8// by the Apache License, Version 2.0.
9
10use 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    /// The reconciliation metrics of every controller, as defined by
29    /// `k8s_controller`: counts and durations of reconciliation passes, and
30    /// of the steps within them.
31    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    /// Records that this replica has taken the leadership lease.
57    pub fn leadership_acquired(&self) {
58        self.is_leader.set(1);
59    }
60
61    /// Records that this replica no longer holds the leadership lease, and
62    /// resets the metrics that only mean anything while it is reconciling.
63    ///
64    /// Those are derived from what reconciliation observed, and the process
65    /// outlives its own leadership: it keeps serving the conversion webhook
66    /// after losing the lease. Without this, a former leader would go on
67    /// publishing its last observation forever, so summing across the
68    /// replicas would count the same organizations once per past leader.
69    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
92///   Adds a tracing layer that reports an `INFO` level span per
93///   request and reports a `WARN` event when a handler returns a
94///   server error to the given Axum Router
95///
96///   This accepts a router instead of returning a layer itself
97///   to avoid dealing with defining generics over a bunch of closures
98///   (see <https://users.rust-lang.org/t/how-to-encapsulate-a-builder-that-depends-on-a-closure/71139/6>)
99///
100///   And this also can't be returned as a Router::new()::layer(TraceLayer)...
101///   because the TraceLayer needs to be added to a Router after
102///   all routes are defined, as it won't trace any routes defined
103///   on the router after it's attached.
104fn 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                // This ugly macro is needed, unfortunately (and
112                // copied from tower-http), because
113                // `tracing::span!` required the level argument to
114                // be static. Meaning we can't just pass
115                // `self.level`.
116                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                // Emit an event at the same level as the span. For the same reason as noted in the comment
156                // above we can't use `tracing::event!(dynamic_level, ...)` since the level argument
157                // needs to be static
158                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    // Don't log Authorization headers
179    _ = 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}