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}