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#[derive(Debug, thiserror::Error)]
24pub enum Error<E: std::error::Error + 'static> {
25 #[error("{0}")]
28 ControllerError(#[source] E),
29 #[error("{0}")]
35 FinalizerError(#[source] kube_runtime::finalizer::Error<E>),
36}
37
38impl<E: std::error::Error + 'static> Error<E> {
39 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 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
59struct Instrumentation {
61 name: Arc<str>,
62 observer: Option<Arc<dyn ReconcileObserver>>,
63 events: Option<Arc<EventRecorder>>,
64}
65
66pub 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 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 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 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 pub fn with_name(mut self, name: impl Into<String>) -> Self {
191 self.name = Some(name.into().into());
192 self
193 }
194
195 pub fn with_observer(mut self, observer: Arc<dyn ReconcileObserver>) -> Self {
199 self.observer = Some(observer);
200 self
201 }
202
203 pub fn with_event_recorder(mut self, events: Arc<EventRecorder>) -> Self {
214 self.events = Some(events);
215 self
216 }
217
218 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 if let Err(e) = res
283 && !matches!(e, kube_runtime::controller::Error::ReconcilerFailed(..))
284 {
285 warn!(
288 error = %e,
289 source = e.source(),
290 "internal kube controller error",
291 );
292 }
293 })
294 .await
295 }
296
297 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#[cfg_attr(not(docsrs), async_trait::async_trait)]
315pub trait Context {
316 type Resource: Resource + Send + Sync + 'static;
319 type Error: std::error::Error;
322
323 const FINALIZER_NAME: Option<&'static str> = None;
330
331 async fn apply(
348 &self,
349 client: Client,
350 resource: &Self::Resource,
351 metadata: &mut TraceMetadata,
352 ) -> Result<Option<Action>, Self::Error>;
353
354 async fn cleanup(
367 &self,
368 client: Client,
369 resource: &Self::Resource,
370 metadata: &mut TraceMetadata,
371 ) -> Result<Option<Action>, Self::Error> {
372 let _client = client;
374 let _resource = resource;
375 let _metadata = metadata;
376
377 Ok(Some(Action::await_change()))
378 }
379
380 fn success_action(&self, resource: &Self::Resource) -> Action {
386 let _resource = resource;
388
389 Action::requeue(Duration::from_secs(rng().random_range(2400..3600)))
390 }
391
392 fn error_action(
400 self: Arc<Self>,
401 resource: Arc<Self::Resource>,
402 err: &Error<Self::Error>,
403 consecutive_errors: u32,
404 ) -> Action {
405 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 fn failure_event(
430 &self,
431 resource: &Self::Resource,
432 phase: Phase,
433 err: &Error<Self::Error>,
434 ) -> Option<Event> {
435 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
452struct 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 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}