Skip to main content

k8s_controller/
controller.rs

1use std::collections::BTreeMap;
2use std::error::Error as _;
3use std::sync::{Arc, Mutex};
4use std::time::{Duration, Instant};
5
6use futures::future::FutureExt;
7use futures::stream::StreamExt;
8use kube::api::Api;
9use kube::core::{ClusterResourceScope, NamespaceResourceScope};
10use kube::{Client, Resource, ResourceExt};
11use kube_runtime::controller::Action;
12use kube_runtime::finalizer::{Event as FinalizerEvent, finalizer};
13use kube_runtime::watcher;
14use rand::{Rng, rng};
15use tracing::field::Empty;
16use tracing::{Instrument, Span, error, info, info_span, trace, warn};
17
18use crate::events::{Event, EventRecorder, EventType};
19use crate::observe::{Outcome, Pass, Phase, ReconcileObserver, TraceMetadata};
20
21/// An error from a reconciliation pass, as passed to
22/// [`Context::error_action`] and [`Context::failure_event`].
23#[derive(Debug, thiserror::Error)]
24pub enum Error<E: std::error::Error + 'static> {
25    /// [`Context::apply`] returned an error, and the context has no
26    /// finalizer.
27    #[error("{0}")]
28    ControllerError(#[source] E),
29    /// The context has a finalizer, and either [`Context::apply`] or
30    /// [`Context::cleanup`] returned an error (wrapped as
31    /// [`ApplyFailed`](kube_runtime::finalizer::Error::ApplyFailed) or
32    /// [`CleanupFailed`](kube_runtime::finalizer::Error::CleanupFailed)), or
33    /// adding or removing the finalizer failed.
34    #[error("{0}")]
35    FinalizerError(#[source] kube_runtime::finalizer::Error<E>),
36}
37
38impl<E: std::error::Error + 'static> Error<E> {
39    /// Renders this error followed by each of its causes, separated by
40    /// `": "`, as used in the note of the default
41    /// [failure event](Context::failure_event).
42    pub fn display_chain(&self) -> String {
43        let mut out = self.to_string();
44        let mut source = self.source();
45        while let Some(err) = source {
46            let message = err.to_string();
47            // Error types commonly end their own message with their source's
48            // (as this one does), which a naive join would repeat.
49            if !out.ends_with(&message) {
50                out.push_str(": ");
51                out.push_str(&message);
52            }
53            source = err.source();
54        }
55        out
56    }
57}
58
59/// The observability configured on a [`Controller`], shared by its passes.
60struct Instrumentation {
61    name: Arc<str>,
62    observer: Option<Arc<dyn ReconcileObserver>>,
63    events: Option<Arc<EventRecorder>>,
64}
65
66/// The [`Controller`] watches a set of resources, calling methods on the
67/// provided [`Context`] when events occur.
68pub struct Controller<Ctx: Context>
69where
70    Ctx: Send + Sync + 'static,
71    Ctx::Error: Send + Sync + 'static,
72    Ctx::Resource: Send + Sync + 'static,
73    Ctx::Resource: Clone + std::fmt::Debug + serde::Serialize,
74    for<'de> Ctx::Resource: serde::Deserialize<'de>,
75    <Ctx::Resource as Resource>::DynamicType:
76        Eq + Clone + std::hash::Hash + std::default::Default + std::fmt::Debug + std::marker::Unpin,
77{
78    client: kube::Client,
79    make_api: Box<dyn Fn(&Ctx::Resource) -> Api<Ctx::Resource> + Sync + Send + 'static>,
80    controller: kube_runtime::controller::Controller<Ctx::Resource>,
81    context: Ctx,
82    name: Option<Arc<str>>,
83    observer: Option<Arc<dyn ReconcileObserver>>,
84    events: Option<Arc<EventRecorder>>,
85}
86
87impl<Ctx: Context> Controller<Ctx>
88where
89    Ctx: Send + Sync + 'static,
90    Ctx::Error: Send + Sync + 'static,
91    Ctx::Resource: Clone + std::fmt::Debug + serde::Serialize,
92    for<'de> Ctx::Resource: serde::Deserialize<'de>,
93    <Ctx::Resource as Resource>::DynamicType:
94        Eq + Clone + std::hash::Hash + std::default::Default + std::fmt::Debug + std::marker::Unpin,
95{
96    /// Creates a new controller for a namespaced resource using the given
97    /// `client`. The `context` given determines the type of resource
98    /// to watch (via the [`Context::Resource`] type provided as part of
99    /// the trait implementation). The resources to be watched will be
100    /// limited to resources in the given `namespace`. A [`watcher::Config`]
101    /// can be given to limit the resources watched (for instance,
102    /// `watcher::Config::default().labels("app=myapp")`).
103    pub fn namespaced(client: Client, context: Ctx, namespace: &str, wc: watcher::Config) -> Self
104    where
105        Ctx::Resource: Resource<Scope = NamespaceResourceScope>,
106    {
107        let make_api = {
108            let client = client.clone();
109            Box::new(move |resource: &Ctx::Resource| {
110                Api::<Ctx::Resource>::namespaced(client.clone(), &resource.namespace().unwrap())
111            })
112        };
113        let controller = kube_runtime::controller::Controller::new(
114            Api::<Ctx::Resource>::namespaced(client.clone(), namespace),
115            wc,
116        );
117        Self::new(client, make_api, controller, context)
118    }
119
120    /// Creates a new controller for a namespaced resource using the given
121    /// `client`. The `context` given determines the type of resource to
122    /// watch (via the [`Context::Resource`] type provided as part of the
123    /// trait implementation). The resources to be watched will not be
124    /// limited by namespace. A [`watcher::Config`] can be given to limit the
125    /// resources watched (for instance,
126    /// `watcher::Config::default().labels("app=myapp")`).
127    pub fn namespaced_all(client: Client, context: Ctx, wc: watcher::Config) -> Self
128    where
129        Ctx::Resource: Resource<Scope = NamespaceResourceScope>,
130    {
131        let make_api = {
132            let client = client.clone();
133            Box::new(move |resource: &Ctx::Resource| {
134                Api::<Ctx::Resource>::namespaced(client.clone(), &resource.namespace().unwrap())
135            })
136        };
137        let controller = kube_runtime::controller::Controller::new(
138            Api::<Ctx::Resource>::all(client.clone()),
139            wc,
140        );
141        Self::new(client, make_api, controller, context)
142    }
143
144    /// Creates a new controller for a cluster-scoped resource using the
145    /// given `client`. The `context` given determines the type of resource
146    /// to watch (via the [`Context::Resource`] type provided as part of the
147    /// trait implementation). A [`watcher::Config`] can be given to limit the
148    /// resources watched (for instance,
149    /// `watcher::Config::default().labels("app=myapp")`).
150    pub fn cluster(client: Client, context: Ctx, wc: watcher::Config) -> Self
151    where
152        Ctx::Resource: Resource<Scope = ClusterResourceScope>,
153    {
154        let make_api = {
155            let client = client.clone();
156            Box::new(move |_: &Ctx::Resource| Api::<Ctx::Resource>::all(client.clone()))
157        };
158        let controller = kube_runtime::controller::Controller::new(
159            Api::<Ctx::Resource>::all(client.clone()),
160            wc,
161        );
162        Self::new(client, make_api, controller, context)
163    }
164
165    fn new(
166        client: Client,
167        make_api: Box<dyn Fn(&Ctx::Resource) -> Api<Ctx::Resource> + Sync + Send + 'static>,
168        controller: kube_runtime::controller::Controller<Ctx::Resource>,
169        context: Ctx,
170    ) -> Self {
171        Self {
172            client,
173            make_api,
174            controller,
175            context,
176            name: None,
177            observer: None,
178            events: None,
179        }
180    }
181
182    /// Sets the name identifying this controller in its metrics (the
183    /// `controller` field of [`ReconcileRecord`](crate::ReconcileRecord) and
184    /// [`StepRecord`](crate::StepRecord)), and in the `controller` field of
185    /// its `reconcile` tracing span. Defaults to [`Context::FINALIZER_NAME`]
186    /// if set, and otherwise to the kind of the resource being watched.
187    ///
188    /// Controllers sharing an [observer](Controller::with_observer) must
189    /// have distinct names, or their metrics will be merged.
190    pub fn with_name(mut self, name: impl Into<String>) -> Self {
191        self.name = Some(name.into().into());
192        self
193    }
194
195    /// Reports every reconciliation pass, and every [step](crate::Step)
196    /// within one, to `observer`. See the [`observe`](crate::observe)
197    /// module.
198    pub fn with_observer(mut self, observer: Arc<dyn ReconcileObserver>) -> Self {
199        self.observer = Some(observer);
200        self
201    }
202
203    /// Publishes a Kubernetes event on the resource whenever reconciling it
204    /// fails, as determined by [`Context::failure_event`], and enables
205    /// [`TraceMetadata::publish_event`] for the context's own events. See
206    /// the [`events`](crate::events) module, including for the RBAC
207    /// permissions this requires.
208    ///
209    /// `events` must not be given to any other controller. Besides every
210    /// event appearing to come from the same controller, a successful pass
211    /// of one controller would reset the aggregation of another's failure
212    /// events for the same resource.
213    pub fn with_event_recorder(mut self, events: Arc<EventRecorder>) -> Self {
214        self.events = Some(events);
215        self
216    }
217
218    /// Run the controller. This method will not return. The [`Context`]
219    /// given to the constructor will have its [`apply`](Context::apply)
220    /// method called when a resource is created or updated, and its
221    /// [`cleanup`](Context::cleanup) method called when a resource is about
222    /// to be deleted.
223    ///
224    /// To run multiple replicas of a controller with only one reconciling
225    /// at a time, pass this method's future to
226    /// [`LeaderElection::with_lease`](crate::LeaderElection::with_lease).
227    pub async fn run(self) {
228        let Self {
229            client,
230            make_api,
231            controller,
232            context,
233            name,
234            observer,
235            events,
236        } = self;
237        let instrumentation = Arc::new(Instrumentation {
238            name: name.unwrap_or_else(|| match Ctx::FINALIZER_NAME {
239                Some(finalizer_name) => finalizer_name.into(),
240                None => Ctx::Resource::kind(&Default::default()).into(),
241            }),
242            observer,
243            events,
244        });
245        let instrumentation = &instrumentation;
246        let backoffs = Arc::new(Mutex::new(BTreeMap::new()));
247        let backoffs = &backoffs;
248        controller
249            .run(
250                |resource, context| {
251                    let uid = resource.uid().unwrap();
252                    let backoffs = Arc::clone(backoffs);
253                    reconcile(
254                        context,
255                        client.clone(),
256                        make_api(&resource),
257                        resource,
258                        Arc::clone(instrumentation),
259                    )
260                    .inspect(move |result| {
261                        if result.is_ok() {
262                            backoffs.lock().unwrap().remove(&uid);
263                        }
264                    })
265                },
266                |resource, err, context| {
267                    let consecutive_errors = {
268                        let uid = resource.uid().unwrap();
269                        let mut backoffs = backoffs.lock().unwrap();
270                        let consecutive_errors: u32 =
271                            backoffs.get(&uid).copied().unwrap_or_default();
272                        backoffs.insert(uid, consecutive_errors.saturating_add(1));
273                        consecutive_errors
274                    };
275                    context.error_action(resource, err, consecutive_errors)
276                },
277                Arc::new(context),
278            )
279            .for_each(|res| async {
280                // ReconcilerFailed errors will already have been reported by
281                // the _reconcile function
282                if let Err(e) = res
283                    && !matches!(e, kube_runtime::controller::Error::ReconcilerFailed(..))
284                {
285                    // warn instead of error because these kinds of errors
286                    // are almost always recoverable
287                    warn!(
288                        error = %e,
289                        source = e.source(),
290                        "internal kube controller error",
291                    );
292                }
293            })
294            .await
295    }
296
297    /// Allow configuring the underlying [`kube_runtime::Controller`]. For
298    /// example, you can use
299    /// `controller.with_controller(|controller| controller.with_config(Config::default().concurrency(10)))`
300    /// to limit the created controller to reconciling 10 resources at once.
301    pub fn with_controller<F>(mut self, f: F) -> Self
302    where
303        F: FnOnce(
304            kube_runtime::Controller<Ctx::Resource>,
305        ) -> kube_runtime::Controller<Ctx::Resource>,
306    {
307        self.controller = f(self.controller);
308        self
309    }
310}
311
312/// The [`Context`] trait should be implemented in order to provide callbacks
313/// for events that happen to resources watched by a [`Controller`].
314#[cfg_attr(not(docsrs), async_trait::async_trait)]
315pub trait Context {
316    /// The type of Kubernetes [resource](Resource) that will be watched by
317    /// the [`Controller`] this context is passed to
318    type Resource: Resource + Send + Sync + 'static;
319    /// The error type which will be returned by the [`apply`](Self::apply)
320    /// and [`cleanup`](Self::cleanup) methods
321    type Error: std::error::Error;
322
323    /// The name to use for the finalizer. This must be unique across
324    /// controllers - if multiple controllers with the same finalizer name
325    /// run against the same resource, unexpected behavior can occur.
326    ///
327    /// If this is None (the default), a finalizer will not be used, and
328    /// cleanup events will not be reported.
329    const FINALIZER_NAME: Option<&'static str> = None;
330
331    /// This method is called when a watched resource is created or updated.
332    /// The [`Client`] used by the controller is passed in to allow making
333    /// additional API requests, as is the resource which triggered this
334    /// event. If this method returns `Some(action)`, the given action will
335    /// be performed, otherwise if `None` is returned,
336    /// [`success_action`](Self::success_action) will be called to find the
337    /// action to perform.
338    ///
339    /// `metadata` is state for this reconciliation pass, which can be used
340    /// to annotate its tracing span, time [steps](TraceMetadata::step) of
341    /// the work for metrics, and [publish events](TraceMetadata::publish_event)
342    /// about the resource.
343    ///
344    /// What this returns determines the [`Outcome`] reported to the
345    /// controller's [observer](Controller::with_observer), as described in
346    /// [`Outcome::of_result`].
347    async fn apply(
348        &self,
349        client: Client,
350        resource: &Self::Resource,
351        metadata: &mut TraceMetadata,
352    ) -> Result<Option<Action>, Self::Error>;
353
354    /// This method is called when a watched resource is marked for deletion.
355    /// The [`Client`] used by the controller is passed in to allow making
356    /// additional API requests, as is the resource which triggered this
357    /// event. If this method returns `Some(action)`, the given action will
358    /// be performed, otherwise if `None` is returned,
359    /// [`success_action`](Self::success_action) will be called to find the
360    /// action to perform.
361    ///
362    /// `metadata` and the return value are treated as for
363    /// [`apply`](Self::apply).
364    ///
365    /// Note that this method will only be called if a finalizer is used.
366    async fn cleanup(
367        &self,
368        client: Client,
369        resource: &Self::Resource,
370        metadata: &mut TraceMetadata,
371    ) -> Result<Option<Action>, Self::Error> {
372        // use a better name for the parameter name in the docs
373        let _client = client;
374        let _resource = resource;
375        let _metadata = metadata;
376
377        Ok(Some(Action::await_change()))
378    }
379
380    /// This method is called when a call to [`apply`](Self::apply) or
381    /// [`cleanup`](Self::cleanup) returns `Ok(None)`. It should return the
382    /// default [`Action`] to perform. The default implementation will
383    /// requeue the event at a random time between 40 and 60 minutes in the
384    /// future.
385    fn success_action(&self, resource: &Self::Resource) -> Action {
386        // use a better name for the parameter name in the docs
387        let _resource = resource;
388
389        Action::requeue(Duration::from_secs(rng().random_range(2400..3600)))
390    }
391
392    /// This method is called when a call to [`apply`](Self::apply) or
393    /// [`cleanup`](Self::cleanup) returns `Err`. It should return the
394    /// default [`Action`] to perform. The error returned will be passed in
395    /// here, as well as a count of how many consecutive errors have happened
396    /// for this resource, to allow for an exponential backoff strategy. The
397    /// default implementation uses exponential backoff with a max of 256
398    /// seconds and some added randomization to avoid thundering herds.
399    fn error_action(
400        self: Arc<Self>,
401        resource: Arc<Self::Resource>,
402        err: &Error<Self::Error>,
403        consecutive_errors: u32,
404    ) -> Action {
405        // use a better name for the parameter name in the docs
406        let _resource = resource;
407        let _err = err;
408
409        let seconds = 2u64.pow(consecutive_errors.min(7) + 1);
410        Action::requeue(Duration::from_millis(
411            rng().random_range((seconds * 500)..(seconds * 1000)),
412        ))
413    }
414
415    /// This method is called when a reconciliation pass fails, if the
416    /// controller has an [event recorder](Controller::with_event_recorder),
417    /// to determine the Kubernetes event to publish on the resource. Return
418    /// `None` to publish nothing, for instance for errors that are an
419    /// expected part of waiting on something else.
420    ///
421    /// The default implementation publishes a `Warning` event with reason
422    /// `ReconcileFailed` (or `CleanupFailed`, if the resource is being
423    /// deleted), and the error and its causes as the note. Repeats of an
424    /// identical event are aggregated, as described in the
425    /// [`events`](crate::events) module; this means that errors which
426    /// include something that varies from one attempt to the next (such as a
427    /// timestamp or request ID) create a new event per attempt, so should be
428    /// rephrased here.
429    fn failure_event(
430        &self,
431        resource: &Self::Resource,
432        phase: Phase,
433        err: &Error<Self::Error>,
434    ) -> Option<Event> {
435        // use a better name for the parameter name in the docs
436        let _resource = resource;
437
438        let (reason, action) = match phase {
439            Phase::Cleanup | Phase::Delete => ("CleanupFailed", "Cleanup"),
440            _ => ("ReconcileFailed", "Reconcile"),
441        };
442        Some(Event {
443            type_: EventType::Warning,
444            reason: reason.to_owned(),
445            action: action.to_owned(),
446            note: Some(err.display_chain()),
447            related: None,
448        })
449    }
450}
451
452/// Records a reconciliation pass as [abandoned](Outcome::Abandoned) if it is
453/// dropped before being finished.
454struct PassGuard {
455    pass: Arc<Pass>,
456    start: Instant,
457    phase: Option<Phase>,
458    fallback_phase: Phase,
459    finished: bool,
460}
461
462impl PassGuard {
463    fn phase(&self) -> Phase {
464        self.phase.unwrap_or(self.fallback_phase)
465    }
466
467    fn finish(&mut self, outcome: Outcome) {
468        self.finished = true;
469        self.pass
470            .record(self.phase(), outcome, self.start.elapsed());
471    }
472}
473
474impl Drop for PassGuard {
475    fn drop(&mut self) {
476        if !self.finished {
477            self.pass
478                .record(self.phase(), Outcome::Abandoned, self.start.elapsed());
479        }
480    }
481}
482
483async fn reconcile<Ctx>(
484    ctx: Arc<Ctx>,
485    client: Client,
486    api: Api<Ctx::Resource>,
487    resource: Arc<Ctx::Resource>,
488    instrumentation: Arc<Instrumentation>,
489) -> Result<Action, Error<Ctx::Error>>
490where
491    Ctx: Context + Send + Sync + 'static,
492    Ctx::Error: Send + Sync + 'static,
493    Ctx::Resource: Send + Sync + 'static,
494    Ctx::Resource: Clone + std::fmt::Debug + serde::Serialize,
495    for<'de> Ctx::Resource: serde::Deserialize<'de>,
496    <Ctx::Resource as Resource>::DynamicType:
497        Eq + Clone + std::hash::Hash + std::default::Default + std::fmt::Debug + std::marker::Unpin,
498{
499    let span = info_span!(
500        "reconcile",
501        resource_type = Ctx::Resource::kind(&Default::default()).as_ref(),
502        resource_name = resource.name_unchecked().as_str(),
503        controller = &*instrumentation.name,
504        event_type = Empty,
505        outcome = Empty,
506        success = Empty,
507        duration_seconds = Empty,
508        metadata = Empty,
509    );
510    async {
511        trace!("beginning reconciliation");
512
513        let pass = Arc::new(Pass {
514            controller: Arc::clone(&instrumentation.name),
515            reference: resource.object_ref(&Default::default()),
516            observer: instrumentation.observer.clone(),
517            events: instrumentation.events.clone(),
518        });
519        let mut metadata = TraceMetadata::for_pass(Arc::clone(&pass));
520        let mut guard = PassGuard {
521            pass: Arc::clone(&pass),
522            start: Instant::now(),
523            phase: None,
524            fallback_phase: if resource.meta().deletion_timestamp.is_some() {
525                Phase::Delete
526            } else {
527                Phase::Init
528            },
529            finished: false,
530        };
531        let mut reconciler_outcome = None;
532
533        let res = if let Some(finalizer_name) = Ctx::FINALIZER_NAME {
534            finalizer(&api, finalizer_name, Arc::clone(&resource), |event| async {
535                match event {
536                    FinalizerEvent::Apply(resource) => {
537                        guard.phase = Some(Phase::Apply);
538                        let res = ctx.apply(client, &resource, &mut metadata).await;
539                        reconciler_outcome = Some(Outcome::of_result(&res));
540                        res.map(|action| action.unwrap_or_else(|| ctx.success_action(&resource)))
541                    }
542                    FinalizerEvent::Cleanup(resource) => {
543                        guard.phase = Some(Phase::Cleanup);
544                        let res = ctx.cleanup(client, &resource, &mut metadata).await;
545                        reconciler_outcome = Some(Outcome::of_result(&res));
546                        res.map(|action| action.unwrap_or_else(Action::await_change))
547                    }
548                }
549            })
550            .await
551            .map_err(Error::FinalizerError)
552        } else if resource.meta().deletion_timestamp.is_none() {
553            guard.phase = Some(Phase::Apply);
554            let res = ctx.apply(client, &resource, &mut metadata).await;
555            reconciler_outcome = Some(Outcome::of_result(&res));
556            res.map(|action| action.unwrap_or_else(|| ctx.success_action(&resource)))
557                .map_err(Error::ControllerError)
558        } else {
559            Ok(Action::await_change())
560        };
561
562        let outcome = match &res {
563            Err(_) => Outcome::Failed,
564            Ok(_) => metadata
565                .outcome
566                .or(reconciler_outcome)
567                .unwrap_or(Outcome::Completed),
568        };
569        let phase = guard.phase();
570        let duration = guard.start.elapsed();
571        guard.finish(outcome);
572
573        let span = Span::current();
574        span.record("event_type", phase.as_str());
575        span.record("outcome", outcome.as_str());
576        span.record("duration_seconds", duration.as_secs_f64());
577
578        if !metadata.annotations.is_empty()
579            && let Ok(s) = serde_json::to_string(&metadata.annotations)
580        {
581            span.record("metadata", s);
582        }
583
584        if let Err(e) = &res {
585            span.record("success", false);
586            error!(error = %e, source = e.source(), "reconcile");
587        } else {
588            span.record("success", true);
589            info!("reconcile");
590        }
591
592        if let Some(events) = &instrumentation.events {
593            match &res {
594                Err(e) => {
595                    if let Some(event) = ctx.failure_event(&resource, phase, e)
596                        && let Err(publish_err) =
597                            events.publish_to(&pass.reference, &event, true).await
598                    {
599                        warn!(
600                            error = %publish_err,
601                            reason = %event.reason,
602                            "failed to publish reconciliation failure event",
603                        );
604                    }
605                }
606                Ok(_) => {
607                    if let Some(uid) = &pass.reference.uid {
608                        events.forget_failures(uid);
609                    }
610                }
611            }
612        }
613
614        res
615    }
616    .instrument(span)
617    .await
618}
619
620#[cfg(test)]
621mod tests {
622    use std::future::pending;
623    use std::sync::Mutex;
624
625    use futures::FutureExt;
626    use k8s_openapi::api::core::v1::ConfigMap;
627    use k8s_openapi::apimachinery::pkg::apis::meta::v1::{ObjectMeta, Time};
628    use k8s_openapi::jiff::Timestamp;
629
630    use super::*;
631    use crate::events::Reporter;
632    use crate::observe::{ReconcileRecord, StepRecord};
633    use crate::test_util::{MockApiServer, block_on};
634
635    #[derive(Debug, thiserror::Error)]
636    enum TestError {
637        #[error("reconciling failed: {0}")]
638        Wrapped(#[source] std::io::Error),
639    }
640
641    #[derive(Clone, Copy)]
642    enum Behavior {
643        Done,
644        Requeue,
645        Skip,
646        Fail,
647        FailQuietly,
648        Hang,
649    }
650
651    struct TestContext {
652        behavior: Mutex<Behavior>,
653    }
654
655    impl TestContext {
656        fn new(behavior: Behavior) -> Arc<Self> {
657            Arc::new(Self {
658                behavior: Mutex::new(behavior),
659            })
660        }
661
662        fn set(&self, behavior: Behavior) {
663            *self.behavior.lock().unwrap() = behavior;
664        }
665    }
666
667    #[async_trait::async_trait]
668    impl Context for TestContext {
669        type Resource = ConfigMap;
670        type Error = TestError;
671
672        async fn apply(
673            &self,
674            _client: Client,
675            _resource: &ConfigMap,
676            metadata: &mut TraceMetadata,
677        ) -> Result<Option<Action>, TestError> {
678            let behavior = *self.behavior.lock().unwrap();
679            let step = metadata.step("work");
680            match behavior {
681                Behavior::Done => {
682                    step.finish(Outcome::Completed);
683                    Ok(None)
684                }
685                Behavior::Requeue => {
686                    step.finish(Outcome::Waiting);
687                    Ok(Some(Action::requeue(Duration::from_secs(1))))
688                }
689                Behavior::Skip => {
690                    step.finish(Outcome::Skipped);
691                    metadata.set_outcome(Outcome::Skipped);
692                    Ok(None)
693                }
694                Behavior::Fail | Behavior::FailQuietly => {
695                    Err(TestError::Wrapped(std::io::Error::other("disk on fire")))
696                }
697                Behavior::Hang => pending().await,
698            }
699        }
700
701        fn failure_event(
702            &self,
703            _resource: &ConfigMap,
704            _phase: Phase,
705            err: &Error<TestError>,
706        ) -> Option<Event> {
707            match *self.behavior.lock().unwrap() {
708                Behavior::FailQuietly => None,
709                _ => Some(Event {
710                    type_: EventType::Warning,
711                    reason: "ReconcileFailed".to_owned(),
712                    action: "Reconcile".to_owned(),
713                    note: Some(err.display_chain()),
714                    related: None,
715                }),
716            }
717        }
718    }
719
720    #[derive(Debug, Clone, PartialEq)]
721    enum Recorded {
722        Pass(String, Phase, Outcome),
723        Step(String, &'static str, Outcome),
724    }
725
726    #[derive(Default)]
727    struct TestObserver(Mutex<Vec<Recorded>>);
728
729    impl TestObserver {
730        fn take(&self) -> Vec<Recorded> {
731            std::mem::take(&mut self.0.lock().unwrap())
732        }
733    }
734
735    impl ReconcileObserver for TestObserver {
736        fn reconciled(&self, record: &ReconcileRecord<'_>) {
737            assert_eq!(record.kind, "ConfigMap");
738            assert_eq!(record.namespace, Some("ns"));
739            assert_eq!(record.name, "cm");
740            self.0.lock().unwrap().push(Recorded::Pass(
741                record.controller.to_owned(),
742                record.phase,
743                record.outcome,
744            ));
745        }
746
747        fn step_finished(&self, record: &StepRecord<'_>) {
748            self.0.lock().unwrap().push(Recorded::Step(
749                record.controller.to_owned(),
750                record.step,
751                record.outcome,
752            ));
753        }
754    }
755
756    struct Harness {
757        server: MockApiServer,
758        observer: Arc<TestObserver>,
759        instrumentation: Arc<Instrumentation>,
760    }
761
762    impl Harness {
763        fn new() -> Self {
764            let server = MockApiServer::new();
765            let observer = Arc::new(TestObserver::default());
766            let instrumentation = Arc::new(Instrumentation {
767                name: "test".into(),
768                observer: Some(Arc::<TestObserver>::clone(&observer)),
769                events: Some(Arc::new(EventRecorder::new(
770                    server.client(),
771                    Reporter {
772                        controller: "test.example.com".to_owned(),
773                        instance: None,
774                    },
775                ))),
776            });
777            Self {
778                server,
779                observer,
780                instrumentation,
781            }
782        }
783
784        fn reconcile(
785            &self,
786            ctx: &Arc<TestContext>,
787            resource: ConfigMap,
788        ) -> impl Future<Output = Result<Action, Error<TestError>>> + use<> {
789            let client = self.server.client();
790            reconcile(
791                Arc::clone(ctx),
792                client.clone(),
793                Api::namespaced(client, "ns"),
794                Arc::new(resource),
795                Arc::clone(&self.instrumentation),
796            )
797        }
798
799        fn methods(&self) -> Vec<String> {
800            self.server
801                .requests()
802                .into_iter()
803                .map(|r| r.method)
804                .collect()
805        }
806    }
807
808    fn config_map() -> ConfigMap {
809        ConfigMap {
810            metadata: ObjectMeta {
811                name: Some("cm".to_owned()),
812                namespace: Some("ns".to_owned()),
813                uid: Some("uid-1".to_owned()),
814                ..Default::default()
815            },
816            ..Default::default()
817        }
818    }
819
820    fn pass(phase: Phase, outcome: Outcome) -> Recorded {
821        Recorded::Pass("test".to_owned(), phase, outcome)
822    }
823
824    fn step(outcome: Outcome) -> Recorded {
825        Recorded::Step("test".to_owned(), "work", outcome)
826    }
827
828    #[test]
829    fn classifies_successful_passes() {
830        block_on(async {
831            let harness = Harness::new();
832            let ctx = TestContext::new(Behavior::Done);
833            harness.reconcile(&ctx, config_map()).await.unwrap();
834            assert_eq!(
835                harness.observer.take(),
836                [
837                    step(Outcome::Completed),
838                    pass(Phase::Apply, Outcome::Completed)
839                ]
840            );
841
842            ctx.set(Behavior::Requeue);
843            harness.reconcile(&ctx, config_map()).await.unwrap();
844            assert_eq!(
845                harness.observer.take(),
846                [step(Outcome::Waiting), pass(Phase::Apply, Outcome::Waiting)]
847            );
848
849            ctx.set(Behavior::Skip);
850            harness.reconcile(&ctx, config_map()).await.unwrap();
851            assert_eq!(
852                harness.observer.take(),
853                [step(Outcome::Skipped), pass(Phase::Apply, Outcome::Skipped)]
854            );
855
856            assert!(harness.server.requests().is_empty());
857        });
858    }
859
860    #[test]
861    fn publishes_and_aggregates_failure_events() {
862        block_on(async {
863            let harness = Harness::new();
864            let ctx = TestContext::new(Behavior::Fail);
865            harness.reconcile(&ctx, config_map()).await.unwrap_err();
866            assert_eq!(
867                harness.observer.take(),
868                [
869                    step(Outcome::Abandoned),
870                    pass(Phase::Apply, Outcome::Failed)
871                ]
872            );
873            let requests = harness.server.requests();
874            assert_eq!(requests[0].body["reason"], "ReconcileFailed");
875            assert_eq!(requests[0].body["note"], "reconciling failed: disk on fire");
876
877            harness.reconcile(&ctx, config_map()).await.unwrap_err();
878            assert_eq!(harness.methods(), ["POST", "PATCH"]);
879
880            // a success ends the series, so the next failure is a new event
881            ctx.set(Behavior::Done);
882            harness.reconcile(&ctx, config_map()).await.unwrap();
883            ctx.set(Behavior::Fail);
884            harness.reconcile(&ctx, config_map()).await.unwrap_err();
885            assert_eq!(harness.methods(), ["POST", "PATCH", "POST"]);
886        });
887    }
888
889    #[test]
890    fn failure_event_can_be_suppressed() {
891        block_on(async {
892            let harness = Harness::new();
893            let ctx = TestContext::new(Behavior::FailQuietly);
894            harness.reconcile(&ctx, config_map()).await.unwrap_err();
895            assert_eq!(
896                harness.observer.take(),
897                [
898                    step(Outcome::Abandoned),
899                    pass(Phase::Apply, Outcome::Failed)
900                ]
901            );
902            assert!(harness.server.requests().is_empty());
903        });
904    }
905
906    #[test]
907    fn deleted_resource_without_finalizer() {
908        block_on(async {
909            let harness = Harness::new();
910            let ctx = TestContext::new(Behavior::Fail);
911            let mut resource = config_map();
912            resource.metadata.deletion_timestamp = Some(Time(Timestamp::now()));
913            harness.reconcile(&ctx, resource).await.unwrap();
914            assert_eq!(
915                harness.observer.take(),
916                [pass(Phase::Delete, Outcome::Completed)]
917            );
918        });
919    }
920
921    #[test]
922    fn cancelled_pass_is_abandoned() {
923        block_on(async {
924            let harness = Harness::new();
925            let ctx = TestContext::new(Behavior::Hang);
926            let mut fut = Box::pin(harness.reconcile(&ctx, config_map()));
927            assert!((&mut fut).now_or_never().is_none());
928            assert!(harness.observer.take().is_empty());
929            drop(fut);
930            assert_eq!(
931                harness.observer.take(),
932                [
933                    step(Outcome::Abandoned),
934                    pass(Phase::Apply, Outcome::Abandoned)
935                ]
936            );
937        });
938    }
939
940    #[test]
941    fn display_chain_does_not_repeat_messages() {
942        let err: Error<TestError> =
943            Error::FinalizerError(kube_runtime::finalizer::Error::ApplyFailed(
944                TestError::Wrapped(std::io::Error::other("disk on fire")),
945            ));
946        assert_eq!(
947            err.display_chain(),
948            "failed to apply object: reconciling failed: disk on fire"
949        );
950
951        #[derive(Debug, thiserror::Error)]
952        #[error("outer")]
953        struct Outer(#[source] std::io::Error);
954        let err: Error<Outer> = Error::ControllerError(Outer(std::io::Error::other("inner")));
955        assert_eq!(err.display_chain(), "outer: inner");
956
957        #[derive(Debug, thiserror::Error)]
958        #[error("reading lock: timed out waiting for lock")]
959        struct Mentions(#[source] std::io::Error);
960        let err: Error<Mentions> =
961            Error::ControllerError(Mentions(std::io::Error::other("timed out")));
962        assert_eq!(
963            err.display_chain(),
964            "reading lock: timed out waiting for lock: timed out"
965        );
966    }
967}