Skip to main content

k8s_controller/
observe.rs

1//! Observing what reconciliation does, for metrics.
2//!
3//! A [`Controller`](crate::Controller) given a [`ReconcileObserver`] via
4//! [`with_observer`](crate::Controller::with_observer) reports every
5//! reconciliation pass to it as a [`ReconcileRecord`]. Within a pass,
6//! reconcilers can additionally divide their work into named [`Step`]s,
7//! started with [`TraceMetadata::step`], each reported as a [`StepRecord`].
8//! Steps are what attribute a slow or failing pass to the part of the
9//! reconciler responsible, and are the only way to report that a part of it
10//! was [skipped](Outcome::Skipped).
11//!
12//! With the `prometheus` feature enabled,
13//! [`PrometheusMetrics`](crate::PrometheusMetrics) implements
14//! [`ReconcileObserver`] by exporting Prometheus metrics. Implement the trait
15//! directly to export to another metrics system.
16
17use std::collections::BTreeMap;
18use std::fmt::Display;
19use std::sync::Arc;
20use std::time::{Duration, Instant};
21
22use k8s_openapi::api::core::v1::ObjectReference;
23use kube_runtime::controller::Action;
24use tracing::warn;
25
26use crate::events::{Event, EventRecorder};
27
28/// What a reconciliation pass, or a step of one, concluded.
29///
30/// [`as_str`](Outcome::as_str) gives the value used as a metric label, which
31/// is stable.
32#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
33#[non_exhaustive]
34pub enum Outcome {
35    /// Brought what it manages to the desired state, or found it already
36    /// there.
37    Completed,
38    /// Made progress, but the desired state has not been reached yet, and
39    /// asked to be run again to continue.
40    Waiting,
41    /// Had nothing to do, for instance because what it manages is disabled
42    /// by configuration. Never inferred; reconcilers must report it
43    /// explicitly.
44    Skipped,
45    /// Returned an error.
46    Failed,
47    /// Stopped without reaching a conclusion: either the reconciliation was
48    /// cancelled (for instance because leadership was lost, or the process
49    /// is shutting down), or, for a [`Step`], an error propagated out of it
50    /// before it was finished.
51    Abandoned,
52}
53
54impl Outcome {
55    /// The label value for this outcome.
56    pub fn as_str(self) -> &'static str {
57        match self {
58            Outcome::Completed => "completed",
59            Outcome::Waiting => "waiting",
60            Outcome::Skipped => "skipped",
61            Outcome::Failed => "failed",
62            Outcome::Abandoned => "abandoned",
63        }
64    }
65
66    /// Classifies a result in the form returned by
67    /// [`Context::apply`](crate::Context::apply) and
68    /// [`Context::cleanup`](crate::Context::cleanup):
69    ///
70    /// * `Ok(None)` and `Ok(Some(Action::await_change()))` are
71    ///   [`Completed`](Outcome::Completed), since there is nothing more to
72    ///   do until something changes.
73    /// * Any other `Ok(Some(action))` is [`Waiting`](Outcome::Waiting),
74    ///   since the reconciler asked to be run again.
75    /// * `Err(_)` is [`Failed`](Outcome::Failed).
76    ///
77    /// A reconciler that returns a requeue action as a periodic resync after
78    /// having fully converged will therefore be reported as waiting, and
79    /// should override that with [`TraceMetadata::set_outcome`].
80    pub fn of_result<E>(result: &Result<Option<Action>, E>) -> Self {
81        match result {
82            Ok(None) => Outcome::Completed,
83            Ok(Some(action)) if *action == Action::await_change() => Outcome::Completed,
84            Ok(Some(_)) => Outcome::Waiting,
85            Err(_) => Outcome::Failed,
86        }
87    }
88}
89
90impl Display for Outcome {
91    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
92        f.write_str(self.as_str())
93    }
94}
95
96/// Which part of the resource's lifecycle a reconciliation pass handled.
97///
98/// [`as_str`](Phase::as_str) gives the value used as a metric label, which is
99/// stable, and matches the `event_type` field of the `reconcile` tracing
100/// span.
101#[derive(Clone, Copy, Debug, PartialEq, Eq, Hash)]
102#[non_exhaustive]
103pub enum Phase {
104    /// The controller added its finalizer to the resource. Neither
105    /// [`apply`](crate::Context::apply) nor
106    /// [`cleanup`](crate::Context::cleanup) was called; the resource will be
107    /// reconciled again once the finalizer is in place.
108    Init,
109    /// [`Context::apply`](crate::Context::apply) was called.
110    Apply,
111    /// [`Context::cleanup`](crate::Context::cleanup) was called.
112    Cleanup,
113    /// The resource is being deleted, and there was no cleanup to run:
114    /// either the context has no finalizer, or the controller's finalizer
115    /// had already been removed.
116    Delete,
117}
118
119impl Phase {
120    /// The label value for this phase.
121    pub fn as_str(self) -> &'static str {
122        match self {
123            Phase::Init => "init",
124            Phase::Apply => "apply",
125            Phase::Cleanup => "cleanup",
126            Phase::Delete => "delete",
127        }
128    }
129}
130
131impl Display for Phase {
132    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
133        f.write_str(self.as_str())
134    }
135}
136
137/// A completed (or [abandoned](Outcome::Abandoned)) reconciliation pass, as
138/// reported to [`ReconcileObserver::reconciled`].
139#[derive(Clone, Copy, Debug)]
140#[non_exhaustive]
141pub struct ReconcileRecord<'a> {
142    /// The [name](crate::Controller::with_name) of the controller.
143    pub controller: &'a str,
144    /// The kind of the reconciled resource.
145    pub kind: &'a str,
146    /// The namespace of the reconciled resource, if it is namespaced.
147    pub namespace: Option<&'a str>,
148    /// The name of the reconciled resource.
149    pub name: &'a str,
150    pub phase: Phase,
151    pub outcome: Outcome,
152    /// Time spent in the pass, including the controller's finalizer
153    /// bookkeeping but excluding publishing a failure event.
154    pub duration: Duration,
155}
156
157/// A finished (or [abandoned](Outcome::Abandoned)) [`Step`], as reported to
158/// [`ReconcileObserver::step_finished`].
159#[derive(Clone, Copy, Debug)]
160#[non_exhaustive]
161pub struct StepRecord<'a> {
162    /// The [name](crate::Controller::with_name) of the controller.
163    pub controller: &'a str,
164    /// The kind of the reconciled resource.
165    pub kind: &'a str,
166    /// The namespace of the reconciled resource, if it is namespaced.
167    pub namespace: Option<&'a str>,
168    /// The name of the reconciled resource.
169    pub name: &'a str,
170    /// The name the step was started with.
171    pub step: &'static str,
172    pub outcome: Outcome,
173    pub duration: Duration,
174}
175
176/// Receives reports of what reconciliation did.
177///
178/// Methods are called synchronously from the reconciliation (including from
179/// [`Step`]'s `Drop` implementation), so they should be cheap and must not
180/// block.
181///
182/// Records identify the reconciled resource, but its namespace and name
183/// usually make poor metric labels: they are unbounded, and a label set is
184/// kept for as long as the process lives, even after the resource is
185/// deleted.
186pub trait ReconcileObserver: Send + Sync {
187    /// Called when a reconciliation pass ends, including when it is
188    /// cancelled.
189    fn reconciled(&self, record: &ReconcileRecord<'_>) {
190        let _record = record;
191    }
192
193    /// Called when a [`Step`] is finished or dropped.
194    fn step_finished(&self, record: &StepRecord<'_>) {
195        let _record = record;
196    }
197}
198
199/// The context shared by everything recorded during one reconciliation pass.
200pub(crate) struct Pass {
201    pub(crate) controller: Arc<str>,
202    pub(crate) reference: ObjectReference,
203    pub(crate) observer: Option<Arc<dyn ReconcileObserver>>,
204    pub(crate) events: Option<Arc<EventRecorder>>,
205}
206
207impl Pass {
208    fn kind(&self) -> &str {
209        self.reference.kind.as_deref().unwrap_or_default()
210    }
211
212    fn name(&self) -> &str {
213        self.reference.name.as_deref().unwrap_or_default()
214    }
215
216    pub(crate) fn record(&self, phase: Phase, outcome: Outcome, duration: Duration) {
217        if let Some(observer) = &self.observer {
218            observer.reconciled(&ReconcileRecord {
219                controller: &self.controller,
220                kind: self.kind(),
221                namespace: self.reference.namespace.as_deref(),
222                name: self.name(),
223                phase,
224                outcome,
225                duration,
226            });
227        }
228    }
229}
230
231/// Times one named step of a reconciliation pass, reporting it to the
232/// controller's [`ReconcileObserver`] when finished or dropped.
233///
234/// Created by [`TraceMetadata::step`]. A step dropped without being finished
235/// is reported as [`Outcome::Abandoned`], so that `?` propagating an error
236/// out of a step identifies the step the pass stopped in, without any
237/// bookkeeping at the early return. The flip side is that every path out of
238/// a step that does reach a conclusion must finish it explicitly:
239///
240/// ```no_run
241/// # use k8s_controller::{Outcome, TraceMetadata};
242/// # async fn create_certificate() -> Result<(), kube::Error> { Ok(()) }
243/// # async fn f(metadata: &TraceMetadata, tls_enabled: bool) -> Result<(), kube::Error> {
244/// let step = metadata.step("certificate");
245/// if tls_enabled {
246///     create_certificate().await?;
247///     step.finish(Outcome::Completed);
248/// } else {
249///     step.finish(Outcome::Skipped);
250/// }
251/// # Ok(())
252/// # }
253/// ```
254#[must_use = "a step that is never finished is reported as abandoned"]
255pub struct Step {
256    pass: Option<Arc<Pass>>,
257    name: &'static str,
258    start: Instant,
259    outcome: Option<Outcome>,
260}
261
262impl Step {
263    /// Finishes the step with the given outcome.
264    pub fn finish(mut self, outcome: Outcome) {
265        self.outcome = Some(outcome);
266    }
267
268    /// Finishes the step with the outcome [classified](Outcome::of_result)
269    /// from a reconciler-style result.
270    ///
271    /// Prefer this to propagating an error out of the step with `?` when the
272    /// result is in hand, since that reports the step as
273    /// [abandoned](Outcome::Abandoned) rather than [failed](Outcome::Failed).
274    pub fn finish_with<E>(self, result: &Result<Option<Action>, E>) {
275        self.finish(Outcome::of_result(result));
276    }
277}
278
279impl Drop for Step {
280    fn drop(&mut self) {
281        let Some(pass) = &self.pass else { return };
282        let Some(observer) = &pass.observer else {
283            return;
284        };
285        observer.step_finished(&StepRecord {
286            controller: &pass.controller,
287            kind: pass.kind(),
288            namespace: pass.reference.namespace.as_deref(),
289            name: pass.name(),
290            step: self.name,
291            outcome: self.outcome.unwrap_or(Outcome::Abandoned),
292            duration: self.start.elapsed(),
293        });
294    }
295}
296
297/// Per-pass state handed to [`Context::apply`](crate::Context::apply) and
298/// [`Context::cleanup`](crate::Context::cleanup).
299///
300/// It collects annotations for the pass's tracing span, times
301/// [steps](TraceMetadata::step), overrides the pass's reported
302/// [outcome](TraceMetadata::set_outcome), and
303/// [publishes events](TraceMetadata::publish_event) about the resource being
304/// reconciled.
305///
306/// A `TraceMetadata` created with [`Default`] (for instance, to call a
307/// reconciler from a unit test) accepts all of these, but reports steps and
308/// publishes events nowhere.
309#[derive(Default)]
310pub struct TraceMetadata {
311    pub(crate) annotations: BTreeMap<String, String>,
312    pub(crate) outcome: Option<Outcome>,
313    pass: Option<Arc<Pass>>,
314}
315
316impl std::fmt::Debug for TraceMetadata {
317    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
318        f.debug_struct("TraceMetadata")
319            .field("annotations", &self.annotations)
320            .field("outcome", &self.outcome)
321            .finish_non_exhaustive()
322    }
323}
324
325impl TraceMetadata {
326    pub(crate) fn for_pass(pass: Arc<Pass>) -> Self {
327        Self {
328            pass: Some(pass),
329            ..Default::default()
330        }
331    }
332
333    /// Adds a key/value pair to the `metadata` field of the pass's
334    /// `reconcile` tracing span, replacing any earlier value for `key`.
335    pub fn annotate<K: Display, V: Display>(&mut self, key: K, val: V) {
336        self.annotations.insert(key.to_string(), val.to_string());
337    }
338
339    /// Starts timing the step named `name`. See [`Step`].
340    ///
341    /// Step names become metric labels, so they should come from a small
342    /// fixed set.
343    pub fn step(&self, name: &'static str) -> Step {
344        Step {
345            pass: self.pass.clone(),
346            name,
347            start: Instant::now(),
348            outcome: None,
349        }
350    }
351
352    /// Overrides the outcome reported for this pass, if the reconciler
353    /// returns `Ok`. By default the outcome is
354    /// [classified](Outcome::of_result) from what the reconciler returns;
355    /// an `Err` is always reported as [`Outcome::Failed`].
356    pub fn set_outcome(&mut self, outcome: Outcome) {
357        self.outcome = Some(outcome);
358    }
359
360    /// Publishes `event` about the resource being reconciled, using the
361    /// controller's [event recorder](crate::Controller::with_event_recorder).
362    /// Does nothing if the controller has none.
363    ///
364    /// Failure to publish is logged rather than returned, since an event
365    /// only reports on reconciliation and should not change its course.
366    /// Use [`EventRecorder::publish`] directly to handle errors.
367    pub async fn publish_event(&self, event: Event) {
368        let Some(pass) = &self.pass else { return };
369        let Some(events) = &pass.events else { return };
370        if let Err(e) = events.publish_to(&pass.reference, &event, false).await {
371            warn!(
372                error = %e,
373                reason = %event.reason,
374                controller = %pass.controller,
375                "failed to publish event",
376            );
377        }
378    }
379}