Skip to main content

k8s_controller/
prometheus.rs

1use prometheus::core::{Collector, Desc};
2use prometheus::proto::MetricFamily;
3use prometheus::{HistogramOpts, HistogramVec, IntCounterVec, Opts};
4
5use crate::observe::{ReconcileObserver, ReconcileRecord, StepRecord};
6
7/// A [`ReconcileObserver`] that exports Prometheus metrics. Requires the
8/// `prometheus` feature.
9///
10/// The metrics, each prefixed with the namespace given to
11/// [`new`](PrometheusMetrics::new), are:
12///
13/// * `<namespace>_reconciliations_total{controller, phase, outcome}`: count
14///   of reconciliation passes. An `outcome` of `failed` is the signal to
15///   alert on. `waiting` is normal while a resource converges.
16/// * `<namespace>_reconciliation_duration_seconds{controller, phase}`:
17///   histogram of time spent in each pass. Reconcilers that wait by
18///   requeueing, rather than by blocking, return promptly, so this measures
19///   work done rather than time to converge.
20/// * `<namespace>_reconciliation_steps_total{controller, step, outcome}`:
21///   count of [steps](crate::Step). Since a step that an error propagates
22///   out of is reported as `abandoned`, as is one in a pass that was
23///   cancelled, `abandoned` locates where passes stop, but is not itself a
24///   failure signal.
25/// * `<namespace>_reconciliation_step_duration_seconds{controller, step}`:
26///   histogram of time spent in each step.
27///
28/// Label values are the [controller name](crate::Controller::with_name),
29/// [`Phase::as_str`](crate::Phase::as_str),
30/// [`Outcome::as_str`](crate::Outcome::as_str), and the step name.
31///
32/// This is a [`Collector`], so it must be registered with a registry to be
33/// exported, and then shared by every controller that reports to it:
34///
35/// ```no_run
36/// # use std::sync::Arc;
37/// # fn f() -> Result<(), prometheus::Error> {
38/// let metrics = k8s_controller::PrometheusMetrics::new("my_operator")?;
39/// let registry = prometheus::Registry::new();
40/// registry.register(Box::new(metrics.clone()))?;
41/// let observer: Arc<dyn k8s_controller::ReconcileObserver> = Arc::new(metrics);
42/// // pass `Arc::clone(&observer)` to each `Controller::with_observer`
43/// # Ok(())
44/// # }
45/// ```
46#[derive(Clone, Debug)]
47pub struct PrometheusMetrics {
48    reconciliations: IntCounterVec,
49    reconciliation_duration: HistogramVec,
50    steps: IntCounterVec,
51    step_duration: HistogramVec,
52}
53
54impl PrometheusMetrics {
55    /// The default histogram buckets, in seconds.
56    pub const DEFAULT_BUCKETS: &[f64] = &[0.01, 0.05, 0.25, 1.0, 5.0, 30.0];
57
58    /// Creates the metrics, with names prefixed by `namespace` (which may be
59    /// empty, for no prefix) and histograms using
60    /// [`DEFAULT_BUCKETS`](Self::DEFAULT_BUCKETS).
61    pub fn new(namespace: &str) -> Result<Self, prometheus::Error> {
62        Self::with_buckets(namespace, Self::DEFAULT_BUCKETS.to_vec())
63    }
64
65    /// Creates the metrics, with names prefixed by `namespace` (which may be
66    /// empty, for no prefix) and histograms using the given `buckets`, in
67    /// seconds.
68    pub fn with_buckets(namespace: &str, buckets: Vec<f64>) -> Result<Self, prometheus::Error> {
69        Ok(Self {
70            reconciliations: IntCounterVec::new(
71                Opts::new(
72                    "reconciliations_total",
73                    "Count of reconciliation passes, by controller, by the phase of the resource's \
74                     lifecycle handled, and by what the pass concluded. An outcome of `failed` \
75                     means the reconciler returned an error.",
76                )
77                .namespace(namespace),
78                &["controller", "phase", "outcome"],
79            )?,
80            reconciliation_duration: HistogramVec::new(
81                HistogramOpts::new(
82                    "reconciliation_duration_seconds",
83                    "Time spent in one reconciliation pass.",
84                )
85                .namespace(namespace)
86                .buckets(buckets.clone()),
87                &["controller", "phase"],
88            )?,
89            steps: IntCounterVec::new(
90                Opts::new(
91                    "reconciliation_steps_total",
92                    "Count of reconciliation steps, by controller, by step, and by what the step \
93                     concluded. An outcome of `abandoned` means the step did not conclude, either \
94                     because an error propagated out of it or because the pass was cancelled.",
95                )
96                .namespace(namespace),
97                &["controller", "step", "outcome"],
98            )?,
99            step_duration: HistogramVec::new(
100                HistogramOpts::new(
101                    "reconciliation_step_duration_seconds",
102                    "Time spent in one reconciliation step.",
103                )
104                .namespace(namespace)
105                .buckets(buckets),
106                &["controller", "step"],
107            )?,
108        })
109    }
110}
111
112impl ReconcileObserver for PrometheusMetrics {
113    fn reconciled(&self, record: &ReconcileRecord<'_>) {
114        let phase = record.phase.as_str();
115        self.reconciliations
116            .with_label_values(&[record.controller, phase, record.outcome.as_str()])
117            .inc();
118        self.reconciliation_duration
119            .with_label_values(&[record.controller, phase])
120            .observe(record.duration.as_secs_f64());
121    }
122
123    fn step_finished(&self, record: &StepRecord<'_>) {
124        self.steps
125            .with_label_values(&[record.controller, record.step, record.outcome.as_str()])
126            .inc();
127        self.step_duration
128            .with_label_values(&[record.controller, record.step])
129            .observe(record.duration.as_secs_f64());
130    }
131}
132
133impl PrometheusMetrics {
134    fn collectors(&self) -> [&dyn Collector; 4] {
135        let Self {
136            reconciliations,
137            reconciliation_duration,
138            steps,
139            step_duration,
140        } = self;
141        [
142            reconciliations,
143            reconciliation_duration,
144            steps,
145            step_duration,
146        ]
147    }
148}
149
150impl Collector for PrometheusMetrics {
151    fn desc(&self) -> Vec<&Desc> {
152        self.collectors()
153            .into_iter()
154            .flat_map(Collector::desc)
155            .collect()
156    }
157
158    fn collect(&self) -> Vec<MetricFamily> {
159        self.collectors()
160            .into_iter()
161            .flat_map(Collector::collect)
162            .collect()
163    }
164}
165
166#[cfg(test)]
167mod tests {
168    use std::time::Duration;
169
170    use super::*;
171    use crate::{Outcome, Phase};
172
173    #[test]
174    fn exports_records() {
175        let metrics = PrometheusMetrics::new("test").unwrap();
176        let registry = prometheus::Registry::new();
177        registry.register(Box::new(metrics.clone())).unwrap();
178
179        metrics.reconciled(&ReconcileRecord {
180            controller: "widgets",
181            kind: "Widget",
182            namespace: Some("ns"),
183            name: "w",
184            phase: Phase::Apply,
185            outcome: Outcome::Failed,
186            duration: Duration::from_millis(20),
187        });
188        metrics.step_finished(&StepRecord {
189            controller: "widgets",
190            kind: "Widget",
191            namespace: Some("ns"),
192            name: "w",
193            step: "deployment",
194            outcome: Outcome::Abandoned,
195            duration: Duration::from_millis(10),
196        });
197
198        let families = registry.gather();
199        let names: Vec<_> = families.iter().map(|f| f.name()).collect();
200        assert_eq!(
201            names,
202            [
203                "test_reconciliation_duration_seconds",
204                "test_reconciliation_step_duration_seconds",
205                "test_reconciliation_steps_total",
206                "test_reconciliations_total",
207            ]
208        );
209        let labels = |name: &str| -> Vec<(String, String)> {
210            let family = families.iter().find(|f| f.name() == name).unwrap();
211            family.get_metric()[0]
212                .get_label()
213                .iter()
214                .map(|l| (l.name().to_owned(), l.value().to_owned()))
215                .collect()
216        };
217        let owned = |pairs: &[(&str, &str)]| -> Vec<(String, String)> {
218            pairs
219                .iter()
220                .map(|(k, v)| (k.to_string(), v.to_string()))
221                .collect()
222        };
223        assert_eq!(
224            labels("test_reconciliations_total"),
225            owned(&[
226                ("controller", "widgets"),
227                ("outcome", "failed"),
228                ("phase", "apply")
229            ])
230        );
231        assert_eq!(
232            labels("test_reconciliation_steps_total"),
233            owned(&[
234                ("controller", "widgets"),
235                ("outcome", "abandoned"),
236                ("step", "deployment")
237            ])
238        );
239    }
240}