1use 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
63pub const SERIES_WINDOW: Duration = Duration::from_secs(10 * 60);
68
69pub const MAX_NOTE_BYTES: usize = 1024;
72
73pub const DEFAULT_TIMEOUT: Duration = Duration::from_secs(5);
76
77const CLUSTER_EVENT_NAMESPACE: &str = "default";
80
81#[derive(Clone, Debug, PartialEq)]
86pub struct Event {
87 pub type_: EventType,
90 pub reason: String,
95 pub action: String,
99 pub note: Option<String>,
102 pub related: Option<ObjectReference>,
104}
105
106#[derive(Debug, thiserror::Error)]
108#[non_exhaustive]
109pub enum PublishError {
110 #[error(transparent)]
112 Kube(#[from] kube::Error),
113 #[error("timed out after {0:?} publishing event")]
116 Timeout(Duration),
117}
118
119pub struct EventRecorder {
126 client: Client,
127 reporter: Reporter,
128 timeout: Duration,
129 series: Mutex<HashMap<SeriesKey, Slot>>,
130}
131
132type Slot = Arc<tokio::sync::Mutex<Option<Series>>>;
136
137#[derive(Clone, Debug, PartialEq, Eq, Hash)]
139struct SeriesKey {
140 failure: bool,
145 uid: String,
146 reason: String,
147 action: String,
148}
149
150struct 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 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 pub fn with_timeout(mut self, timeout: Duration) -> Self {
193 self.timeout = timeout;
194 self
195 }
196
197 pub fn reporter(&self) -> &Reporter {
199 &self.reporter
200 }
201
202 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 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 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 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 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 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
395fn 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}", ¬e[..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}