Skip to main content

k8s_controller/
events.rs

1//! Publishing Kubernetes [events](https://kubernetes.io/docs/reference/kubernetes-api/cluster-resources/event-v1/)
2//! about the resources a controller reconciles.
3//!
4//! Events are what `kubectl describe` shows beneath an object, which makes
5//! them the place to explain to someone without access to the controller's
6//! logs why an object is not converging.
7//!
8//! An [`EventRecorder`] can be given to a [`Controller`](crate::Controller)
9//! via [`with_event_recorder`](crate::Controller::with_event_recorder), in
10//! which case every failed reconciliation is published as an event on the
11//! resource being reconciled (see [`Context::failure_event`](crate::Context::failure_event)).
12//! Reconcilers can publish events of their own through
13//! [`TraceMetadata::publish_event`](crate::TraceMetadata::publish_event), or
14//! through [`EventRecorder::publish`] directly, which also works outside of
15//! reconciliation (for instance, from a background task).
16//!
17//! Each controller should have its own recorder, whose [`Reporter`] names
18//! that controller, even when several controllers in a process reconcile
19//! the same kind of resource. The reporter is the only thing in an event
20//! that identifies which controller published it.
21//!
22//! # Aggregation
23//!
24//! Publishing an event identical to the one last published for the same
25//! object, reason, and action, within [`SERIES_WINDOW`], does not create a
26//! new event object. Instead it increments the `series.count` of the existing
27//! one, so that a reconciliation retrying on a backoff does not bury its
28//! resource in identical events. Any difference, including in the note,
29//! starts a new event, so that an event always describes the most recent
30//! occurrence accurately.
31//!
32//! Failure events published by a controller additionally start afresh after
33//! the resource next reconciles successfully, so a failure that recurs after
34//! a recovery is reported as a new event rather than as a continuation of
35//! the old one.
36//!
37//! # Timeouts
38//!
39//! Each publish is bounded by the recorder's [timeout](EventRecorder::with_timeout),
40//! so that an unresponsive API server delays the reconciliation publishing
41//! an event by at most that long.
42//!
43//! # RBAC
44//!
45//! Publishing requires `create` and `patch` permissions on `events` in the
46//! `events.k8s.io` API group, in every namespace the controller reconciles
47//! resources in. Events about cluster-scoped resources are created in the
48//! `default` namespace.
49
50use std::collections::HashMap;
51use std::sync::{Arc, Mutex};
52use std::time::{Duration, Instant};
53
54use k8s_openapi::api::core::v1::ObjectReference;
55use k8s_openapi::api::events::v1::Event as KubeEvent;
56use k8s_openapi::apimachinery::pkg::apis::meta::v1::MicroTime;
57use k8s_openapi::jiff::Timestamp;
58use kube::api::{Api, ObjectMeta, Patch, PatchParams, PostParams};
59use kube::{Client, Resource, ResourceExt};
60
61pub use kube_runtime::events::{EventType, Reporter};
62
63/// How long after an event was last published that publishing it again
64/// still aggregates into it, rather than creating a new event.
65///
66/// This matches the aggregation window of client-go's event correlator.
67pub const SERIES_WINDOW: Duration = Duration::from_secs(10 * 60);
68
69/// The Kubernetes API server rejects events whose note is longer than this
70/// many bytes.
71pub const MAX_NOTE_BYTES: usize = 1024;
72
73/// The default [timeout](EventRecorder::with_timeout) for publishing an
74/// event.
75pub const DEFAULT_TIMEOUT: Duration = Duration::from_secs(5);
76
77/// Namespace used for events about cluster-scoped resources, which have no
78/// namespace of their own.
79const CLUSTER_EVENT_NAMESPACE: &str = "default";
80
81/// An event to publish about a resource.
82///
83/// `reason` and `action` must each be at most 128 characters, or the API
84/// server will reject the event.
85#[derive(Clone, Debug, PartialEq)]
86pub struct Event {
87    /// Whether this is an ordinary event or a problem. Shown as `Type` by
88    /// `kubectl describe`.
89    pub type_: EventType,
90    /// Why the event happened, as a machine-readable `PascalCase` word (for
91    /// instance `ReconcileFailed`). Shown as `Reason` by `kubectl describe`,
92    /// and what tools and dashboards group events by, so it should be
93    /// treated as a stable interface.
94    pub reason: String,
95    /// What was being done when the event happened, as a machine-readable
96    /// `PascalCase` word (for instance `Reconcile`). Not shown by
97    /// `kubectl describe`.
98    pub action: String,
99    /// A human-readable description of what happened. Shown as `Message` by
100    /// `kubectl describe`. Truncated to [`MAX_NOTE_BYTES`] if longer.
101    pub note: Option<String>,
102    /// Another object involved in the event, if any.
103    pub related: Option<ObjectReference>,
104}
105
106/// An error publishing an event.
107#[derive(Debug, thiserror::Error)]
108#[non_exhaustive]
109pub enum PublishError {
110    /// The API server rejected the event, or could not be reached.
111    #[error(transparent)]
112    Kube(#[from] kube::Error),
113    /// Publishing did not complete within the recorder's
114    /// [timeout](EventRecorder::with_timeout).
115    #[error("timed out after {0:?} publishing event")]
116    Timeout(Duration),
117}
118
119/// Publishes Kubernetes events on resources on behalf of one controller,
120/// aggregating repeats as described in the [module documentation](self).
121///
122/// Share a recorder (via [`Arc`]) between a [`Controller`](crate::Controller)
123/// and whatever else publishes events on that controller's behalf, but not
124/// between controllers.
125pub struct EventRecorder {
126    client: Client,
127    reporter: Reporter,
128    timeout: Duration,
129    series: Mutex<HashMap<SeriesKey, Slot>>,
130}
131
132/// The series published under one [`SeriesKey`], locked for the whole of a
133/// publish so that concurrent publishes of the same event cannot both
134/// create it, or both patch it to the same count.
135type Slot = Arc<tokio::sync::Mutex<Option<Series>>>;
136
137/// Identifies the events that can aggregate into one another.
138#[derive(Clone, Debug, PartialEq, Eq, Hash)]
139struct SeriesKey {
140    /// Whether these are the failure events the controller publishes when
141    /// reconciling the resource fails, as opposed to events published
142    /// through [`EventRecorder::publish`] or
143    /// [`TraceMetadata::publish_event`](crate::TraceMetadata::publish_event).
144    failure: bool,
145    uid: String,
146    reason: String,
147    action: String,
148}
149
150/// The last event published under a [`SeriesKey`].
151struct Series {
152    event: Event,
153    namespace: String,
154    name: String,
155    count: i32,
156    last_published: Instant,
157}
158
159impl std::fmt::Debug for EventRecorder {
160    fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
161        f.debug_struct("EventRecorder")
162            .field("reporter", &self.reporter)
163            .field("timeout", &self.timeout)
164            .finish_non_exhaustive()
165    }
166}
167
168impl EventRecorder {
169    /// Creates a recorder that publishes events as `reporter`, with the
170    /// [default timeout](DEFAULT_TIMEOUT).
171    ///
172    /// [`Reporter::controller`] becomes each event's `reportingController`,
173    /// shown by `kubectl describe` as where the event came from. It should
174    /// name the individual controller, and must be a qualified name (for
175    /// instance `"my-operator.example.com/widgets"`), or the API server
176    /// rejects the event. [`Reporter::instance`] becomes each event's
177    /// `reportingInstance`, and should identify the replica, for which the
178    /// pod name is a good choice; if it is `None`, the controller name is
179    /// used instead.
180    pub fn new(client: Client, reporter: Reporter) -> Self {
181        Self {
182            client,
183            reporter,
184            timeout: DEFAULT_TIMEOUT,
185            series: Mutex::new(HashMap::new()),
186        }
187    }
188
189    /// Sets how long publishing an event may take, including any wait for
190    /// a concurrent publish of the same event, before giving up with
191    /// [`PublishError::Timeout`].
192    pub fn with_timeout(mut self, timeout: Duration) -> Self {
193        self.timeout = timeout;
194        self
195    }
196
197    /// The reporter this recorder publishes events as.
198    pub fn reporter(&self) -> &Reporter {
199        &self.reporter
200    }
201
202    /// Publishes `event` about `resource`.
203    ///
204    /// A note longer than [`MAX_NOTE_BYTES`] is truncated to fit rather than
205    /// having the whole event rejected.
206    ///
207    /// Events are informational, so callers should usually log a failure to
208    /// publish one rather than let it change what they do next.
209    pub async fn publish<K>(&self, resource: &K, event: &Event) -> Result<(), PublishError>
210    where
211        K: Resource,
212        K::DynamicType: Default,
213    {
214        let reference = resource.object_ref(&Default::default());
215        self.publish_to(&reference, event, false).await
216    }
217
218    /// Forgets every event published about `resource`, so that the next
219    /// event published about it creates a new event object even if it is
220    /// identical to the last one.
221    ///
222    /// This is useful when an event reports a state that the resource can
223    /// leave and later re-enter, and each entry should be reported as a
224    /// distinct occurrence.
225    pub fn forget<K: Resource>(&self, resource: &K) {
226        if let Some(uid) = resource.meta().uid.as_deref() {
227            self.series.lock().unwrap().retain(|key, _| key.uid != uid);
228        }
229    }
230
231    /// Forgets the failure events published about the object with the given
232    /// `uid`.
233    pub(crate) fn forget_failures(&self, uid: &str) {
234        self.series
235            .lock()
236            .unwrap()
237            .retain(|key, _| key.uid != uid || !key.failure);
238    }
239
240    /// Publishes `event` about the object `reference` refers to, as a
241    /// failure event if `failure` is set.
242    pub(crate) async fn publish_to(
243        &self,
244        reference: &ObjectReference,
245        event: &Event,
246        failure: bool,
247    ) -> Result<(), PublishError> {
248        tokio::time::timeout(
249            self.timeout,
250            self.publish_unbounded(reference, event, failure),
251        )
252        .await
253        .map_err(|_| PublishError::Timeout(self.timeout))?
254        .map_err(PublishError::Kube)
255    }
256
257    async fn publish_unbounded(
258        &self,
259        reference: &ObjectReference,
260        event: &Event,
261        failure: bool,
262    ) -> Result<(), kube::Error> {
263        let event = Event {
264            note: event.note.as_deref().map(truncate_note),
265            ..(*event).clone()
266        };
267        let namespace = reference
268            .namespace
269            .clone()
270            .unwrap_or_else(|| CLUSTER_EVENT_NAMESPACE.to_owned());
271        let Some(uid) = reference.uid.clone() else {
272            self.create(&namespace, reference, &event).await?;
273            return Ok(());
274        };
275        let slot = self.slot(SeriesKey {
276            failure,
277            uid,
278            reason: event.reason.clone(),
279            action: event.action.clone(),
280        });
281        let mut series = slot.lock().await;
282
283        if let Some(series) = series
284            .as_mut()
285            .filter(|s| s.event == event && s.last_published.elapsed() < SERIES_WINDOW)
286        {
287            let count = series.count.saturating_add(1);
288            match self
289                .patch_series(&series.namespace, &series.name, count)
290                .await
291            {
292                Ok(()) => {
293                    series.count = count;
294                    series.last_published = Instant::now();
295                    return Ok(());
296                }
297                // The API server deletes events some time after their last
298                // update (an hour, by default), which a long-running series
299                // can outlive. Start a new one in that case.
300                Err(kube::Error::Api(e)) if e.code == 404 => {}
301                Err(e) => return Err(e),
302            }
303        }
304
305        let created = self.create(&namespace, reference, &event).await?;
306        *series = Some(Series {
307            event,
308            namespace,
309            name: created.name_any(),
310            count: 1,
311            last_published: Instant::now(),
312        });
313        Ok(())
314    }
315
316    /// Returns the slot for `key`, evicting the slots of series that can no
317    /// longer be aggregated into.
318    fn slot(&self, key: SeriesKey) -> Slot {
319        let mut slots = self.series.lock().unwrap();
320        let now = Instant::now();
321        slots.retain(|k, slot| {
322            *k == key
323                || slot.try_lock().map_or(true, |series| {
324                    series
325                        .as_ref()
326                        .is_some_and(|s| now.duration_since(s.last_published) < SERIES_WINDOW)
327                })
328        });
329        Arc::clone(slots.entry(key).or_default())
330    }
331
332    async fn create(
333        &self,
334        namespace: &str,
335        reference: &ObjectReference,
336        event: &Event,
337    ) -> Result<KubeEvent, kube::Error> {
338        let api: Api<KubeEvent> = Api::namespaced(self.client.clone(), namespace);
339        let event = KubeEvent {
340            metadata: ObjectMeta {
341                generate_name: Some(format!(
342                    "{}.",
343                    reference
344                        .name
345                        .as_deref()
346                        .unwrap_or(&self.reporter.controller)
347                )),
348                namespace: Some(namespace.to_owned()),
349                ..Default::default()
350            },
351            action: Some(event.action.clone()),
352            reason: Some(event.reason.clone()),
353            note: event.note.clone(),
354            type_: Some(
355                match event.type_ {
356                    EventType::Normal => "Normal",
357                    EventType::Warning => "Warning",
358                }
359                .to_owned(),
360            ),
361            event_time: Some(MicroTime(Timestamp::now())),
362            regarding: Some(reference.clone()),
363            related: event.related.clone(),
364            reporting_controller: Some(self.reporter.controller.clone()),
365            reporting_instance: Some(
366                self.reporter
367                    .instance
368                    .clone()
369                    .unwrap_or_else(|| self.reporter.controller.clone()),
370            ),
371            ..Default::default()
372        };
373        api.create(&PostParams::default(), &event).await
374    }
375
376    async fn patch_series(
377        &self,
378        namespace: &str,
379        name: &str,
380        count: i32,
381    ) -> Result<(), kube::Error> {
382        let api: Api<KubeEvent> = Api::namespaced(self.client.clone(), namespace);
383        let patch = serde_json::json!({
384            "series": {
385                "count": count,
386                "lastObservedTime": MicroTime(Timestamp::now()),
387            },
388        });
389        api.patch(name, &PatchParams::default(), &Patch::Merge(&patch))
390            .await?;
391        Ok(())
392    }
393}
394
395/// Shortens `note` to at most [`MAX_NOTE_BYTES`], on a character boundary.
396fn truncate_note(note: &str) -> String {
397    if note.len() <= MAX_NOTE_BYTES {
398        return note.to_owned();
399    }
400    const ELLIPSIS: &str = "...";
401    let mut end = MAX_NOTE_BYTES - ELLIPSIS.len();
402    while !note.is_char_boundary(end) {
403        end -= 1;
404    }
405    format!("{}{ELLIPSIS}", &note[..end])
406}
407
408#[cfg(test)]
409mod tests {
410    use super::*;
411    use crate::test_util::{MockApiServer, block_on};
412
413    use k8s_openapi::api::core::v1::ConfigMap;
414    use serde_json::json;
415
416    fn config_map(uid: &str) -> ConfigMap {
417        ConfigMap {
418            metadata: ObjectMeta {
419                name: Some("cm".to_owned()),
420                namespace: Some("ns".to_owned()),
421                uid: Some(uid.to_owned()),
422                ..Default::default()
423            },
424            ..Default::default()
425        }
426    }
427
428    fn event(note: &str) -> Event {
429        Event {
430            type_: EventType::Warning,
431            reason: "Broken".to_owned(),
432            action: "Reconcile".to_owned(),
433            note: Some(note.to_owned()),
434            related: None,
435        }
436    }
437
438    fn recorder(server: &MockApiServer) -> EventRecorder {
439        EventRecorder::new(
440            server.client(),
441            Reporter {
442                controller: "test.example.com".to_owned(),
443                instance: Some("pod-0".to_owned()),
444            },
445        )
446    }
447
448    #[test]
449    fn test_truncate_note() {
450        assert_eq!(truncate_note("short"), "short");
451
452        let exact = "a".repeat(MAX_NOTE_BYTES);
453        assert_eq!(truncate_note(&exact), exact);
454
455        let long = "a".repeat(MAX_NOTE_BYTES + 1);
456        let truncated = truncate_note(&long);
457        assert_eq!(truncated.len(), MAX_NOTE_BYTES);
458        assert!(truncated.ends_with("..."));
459
460        let multibyte = "é".repeat(MAX_NOTE_BYTES);
461        let truncated = truncate_note(&multibyte);
462        assert!(truncated.len() <= MAX_NOTE_BYTES);
463        assert!(truncated.ends_with("..."));
464    }
465
466    #[test]
467    fn creates_event_with_reporter_and_reference() {
468        block_on(async {
469            let server = MockApiServer::new();
470            let recorder = recorder(&server);
471            recorder
472                .publish(&config_map("uid-1"), &event("it broke"))
473                .await
474                .unwrap();
475
476            let requests = server.requests();
477            assert_eq!(requests.len(), 1);
478            let req = &requests[0];
479            assert_eq!(req.method, "POST");
480            assert_eq!(req.path, "/apis/events.k8s.io/v1/namespaces/ns/events");
481            assert_eq!(req.body["type"], "Warning");
482            assert_eq!(req.body["reason"], "Broken");
483            assert_eq!(req.body["action"], "Reconcile");
484            assert_eq!(req.body["note"], "it broke");
485            assert_eq!(req.body["reportingController"], "test.example.com");
486            assert_eq!(req.body["reportingInstance"], "pod-0");
487            assert_eq!(req.body["regarding"]["kind"], "ConfigMap");
488            assert_eq!(req.body["regarding"]["uid"], "uid-1");
489            assert_eq!(req.body["metadata"]["generateName"], "cm.");
490        });
491    }
492
493    #[test]
494    fn aggregates_identical_events() {
495        block_on(async {
496            let server = MockApiServer::new();
497            let recorder = recorder(&server);
498            let cm = config_map("uid-1");
499            for _ in 0..3 {
500                recorder.publish(&cm, &event("it broke")).await.unwrap();
501            }
502
503            let requests = server.requests();
504            assert_eq!(requests.len(), 3);
505            assert_eq!(requests[0].method, "POST");
506            let name = requests[0].created_name();
507            for (req, count) in requests[1..].iter().zip([2, 3]) {
508                assert_eq!(req.method, "PATCH");
509                assert_eq!(
510                    req.path,
511                    format!("/apis/events.k8s.io/v1/namespaces/ns/events/{name}")
512                );
513                assert_eq!(req.body["series"]["count"], json!(count));
514                assert!(req.body["series"]["lastObservedTime"].is_string());
515            }
516        });
517    }
518
519    #[test]
520    fn concurrent_identical_events_aggregate() {
521        block_on(async {
522            let server = MockApiServer::new();
523            let recorder = recorder(&server);
524            let cm = config_map("uid-1");
525            let ev = event("it broke");
526            let results =
527                futures::future::join_all((0..3).map(|_| recorder.publish(&cm, &ev))).await;
528            assert!(results.iter().all(Result::is_ok));
529
530            let requests = server.requests();
531            let methods: Vec<_> = requests.iter().map(|r| r.method.as_str()).collect();
532            assert_eq!(methods, ["POST", "PATCH", "PATCH"]);
533            assert_eq!(requests[1].body["series"]["count"], json!(2));
534            assert_eq!(requests[2].body["series"]["count"], json!(3));
535        });
536    }
537
538    #[test]
539    fn changed_note_starts_new_event() {
540        block_on(async {
541            let server = MockApiServer::new();
542            let recorder = recorder(&server);
543            let cm = config_map("uid-1");
544            recorder.publish(&cm, &event("first cause")).await.unwrap();
545            recorder.publish(&cm, &event("second cause")).await.unwrap();
546            recorder.publish(&cm, &event("second cause")).await.unwrap();
547
548            let methods: Vec<_> = server.requests().into_iter().map(|r| r.method).collect();
549            assert_eq!(methods, ["POST", "POST", "PATCH"]);
550        });
551    }
552
553    #[test]
554    fn forget_starts_new_event() {
555        block_on(async {
556            let server = MockApiServer::new();
557            let recorder = recorder(&server);
558            let cm = config_map("uid-1");
559            recorder.publish(&cm, &event("it broke")).await.unwrap();
560            recorder.forget(&cm);
561            recorder.publish(&cm, &event("it broke")).await.unwrap();
562
563            let methods: Vec<_> = server.requests().into_iter().map(|r| r.method).collect();
564            assert_eq!(methods, ["POST", "POST"]);
565        });
566    }
567
568    #[test]
569    fn forget_failures_keeps_other_events() {
570        block_on(async {
571            let server = MockApiServer::new();
572            let recorder = recorder(&server);
573            let cm = config_map("uid-1");
574            let reference = cm.object_ref(&());
575            recorder
576                .publish_to(&reference, &event("it broke"), true)
577                .await
578                .unwrap();
579            recorder.publish(&cm, &event("it broke")).await.unwrap();
580            recorder.forget_failures("uid-1");
581            recorder
582                .publish_to(&reference, &event("it broke"), true)
583                .await
584                .unwrap();
585            recorder.publish(&cm, &event("it broke")).await.unwrap();
586
587            let methods: Vec<_> = server.requests().into_iter().map(|r| r.method).collect();
588            assert_eq!(methods, ["POST", "POST", "POST", "PATCH"]);
589        });
590    }
591
592    #[test]
593    fn times_out_and_recovers() {
594        block_on(async {
595            let server = MockApiServer::new();
596            let recorder = recorder(&server).with_timeout(Duration::from_millis(50));
597            let cm = config_map("uid-1");
598            server.hang_next_request();
599            let err = recorder.publish(&cm, &event("it broke")).await.unwrap_err();
600            assert!(matches!(err, PublishError::Timeout(_)), "{err:?}");
601
602            recorder.publish(&cm, &event("it broke")).await.unwrap();
603            let methods: Vec<_> = server.requests().into_iter().map(|r| r.method).collect();
604            assert_eq!(methods, ["POST"]);
605        });
606    }
607
608    #[test]
609    fn expired_event_is_recreated() {
610        block_on(async {
611            let server = MockApiServer::new();
612            let recorder = recorder(&server);
613            let cm = config_map("uid-1");
614            recorder.publish(&cm, &event("it broke")).await.unwrap();
615            server.fail_next_patch(404);
616            recorder.publish(&cm, &event("it broke")).await.unwrap();
617            recorder.publish(&cm, &event("it broke")).await.unwrap();
618
619            let requests = server.requests();
620            let methods: Vec<_> = requests.iter().map(|r| r.method.clone()).collect();
621            assert_eq!(methods, ["POST", "PATCH", "POST", "PATCH"]);
622            assert_eq!(requests[3].body["series"]["count"], json!(2));
623        });
624    }
625
626    #[test]
627    fn cluster_scoped_events_go_to_default_namespace() {
628        block_on(async {
629            let server = MockApiServer::new();
630            let recorder = recorder(&server);
631            let ns = k8s_openapi::api::core::v1::Namespace {
632                metadata: ObjectMeta {
633                    name: Some("some-namespace".to_owned()),
634                    uid: Some("uid-1".to_owned()),
635                    ..Default::default()
636                },
637                ..Default::default()
638            };
639            recorder.publish(&ns, &event("it broke")).await.unwrap();
640
641            let requests = server.requests();
642            assert_eq!(
643                requests[0].path,
644                "/apis/events.k8s.io/v1/namespaces/default/events"
645            );
646        });
647    }
648}