Skip to main content

kube_runtime/reflector/
store.rs

1//! A reader/writer split store for reflectors
2use super::{Lookup, ObjectRef, dispatcher::Dispatcher};
3#[cfg(feature = "unstable-runtime-subscribe")]
4use crate::reflector::ReflectHandle;
5use crate::{
6    utils::delayed_init::{self, DelayedInit},
7    watcher,
8};
9use ahash::AHashMap;
10use educe::Educe;
11use kube_client::{
12    ResourceExt,
13    core::{Selector, SelectorExt},
14};
15use parking_lot::RwLock;
16use std::{fmt::Debug, hash::Hash, sync::Arc};
17use thiserror::Error;
18
19type Cache<K> = Arc<RwLock<AHashMap<ObjectRef<K>, Arc<K>>>>;
20
21/// A writable Store handle
22///
23/// This is exclusive since it's not safe to share a single `Store` between multiple reflectors.
24/// In particular, `Restarted` events will clobber the state of other connected reflectors.
25#[derive(Debug)]
26pub struct Writer<K: 'static + Lookup + Clone>
27where
28    K::DynamicType: Eq + Hash + Clone,
29{
30    store: Cache<K>,
31    buffer: AHashMap<ObjectRef<K>, Arc<K>>,
32    dyntype: K::DynamicType,
33    ready_tx: Option<delayed_init::Initializer<()>>,
34    ready_rx: Arc<DelayedInit<()>>,
35    dispatcher: Option<Dispatcher<K>>,
36}
37
38impl<K: 'static + Lookup + Clone> Writer<K>
39where
40    K::DynamicType: Eq + Hash + Clone,
41{
42    /// Creates a new Writer with the specified dynamic type.
43    ///
44    /// If the dynamic type is default-able (for example when writer is used with
45    /// `k8s_openapi` types) you can use `Default` instead.
46    pub fn new(dyntype: K::DynamicType) -> Self {
47        let (ready_tx, ready_rx) = DelayedInit::new();
48        Writer {
49            store: Default::default(),
50            buffer: Default::default(),
51            dyntype,
52            ready_tx: Some(ready_tx),
53            ready_rx: Arc::new(ready_rx),
54            dispatcher: None,
55        }
56    }
57
58    /// Creates a new Writer with the specified dynamic type and buffer size.
59    ///
60    /// When the Writer is created through `new_shared`, it will be able to
61    /// be subscribed. Stored objects will be propagated to all subscribers. The
62    /// buffer size is used for the underlying channel. An object is cleared
63    /// from the buffer only when all subscribers have seen it.
64    ///
65    /// If the dynamic type is default-able (for example when writer is used with
66    /// `k8s_openapi` types) you can use `Default` instead.
67    #[cfg(feature = "unstable-runtime-subscribe")]
68    pub fn new_shared(buf_size: usize, dyntype: K::DynamicType) -> Self {
69        let (ready_tx, ready_rx) = DelayedInit::new();
70        Writer {
71            store: Default::default(),
72            buffer: Default::default(),
73            dyntype,
74            ready_tx: Some(ready_tx),
75            ready_rx: Arc::new(ready_rx),
76            dispatcher: Some(Dispatcher::new(buf_size)),
77        }
78    }
79
80    /// Return a read handle to the store
81    ///
82    /// Multiple read handles may be obtained, by either calling `as_reader` multiple times,
83    /// or by calling `Store::clone()` afterwards.
84    #[must_use]
85    pub fn as_reader(&self) -> Store<K> {
86        Store {
87            store: self.store.clone(),
88            ready_rx: self.ready_rx.clone(),
89        }
90    }
91
92    /// Return a handle to a subscriber
93    ///
94    /// Multiple subscribe handles may be obtained, by either calling
95    /// `subscribe` multiple times, or by calling `clone()`
96    ///
97    /// This function returns a `Some` when the [`Writer`] is constructed through
98    /// [`Writer::new_shared`] or [`store_shared`], and a `None` otherwise.
99    #[cfg(feature = "unstable-runtime-subscribe")]
100    pub fn subscribe(&self) -> Option<ReflectHandle<K>> {
101        self.dispatcher
102            .as_ref()
103            .map(|dispatcher| dispatcher.subscribe(self.as_reader()))
104    }
105
106    /// Applies a single watcher event to the store
107    pub fn apply_watcher_event(&mut self, event: &watcher::Event<K>) {
108        match event {
109            watcher::Event::Apply(obj) => {
110                let key = obj.to_object_ref(self.dyntype.clone());
111                let obj = Arc::new(obj.clone());
112                self.store.write().insert(key, obj);
113            }
114            watcher::Event::Delete(obj) => {
115                let key = obj.to_object_ref(self.dyntype.clone());
116                self.store.write().remove(&key);
117            }
118            watcher::Event::Init => {
119                self.buffer = AHashMap::new();
120            }
121            watcher::Event::InitApply(obj) => {
122                let key = obj.to_object_ref(self.dyntype.clone());
123                let obj = Arc::new(obj.clone());
124                self.buffer.insert(key, obj);
125            }
126            watcher::Event::InitDone => {
127                let mut store = self.store.write();
128
129                // Swap the buffer into the store
130                std::mem::swap(&mut *store, &mut self.buffer);
131
132                // Clear the buffer
133                // This is preferred over self.buffer.clear(), as clear() will keep the allocated memory for reuse.
134                // This way, the old buffer is dropped.
135                self.buffer = AHashMap::new();
136
137                // Mark as ready after the Restart, "releasing" any calls to Store::wait_until_ready()
138                if let Some(ready_tx) = self.ready_tx.take() {
139                    ready_tx.init(())
140                }
141            }
142        }
143    }
144
145    /// Broadcast an event to any downstream listeners subscribed on the store
146    pub(crate) async fn dispatch_event(&mut self, event: &watcher::Event<K>) {
147        if let Some(ref mut dispatcher) = self.dispatcher {
148            match event {
149                watcher::Event::Apply(obj) => {
150                    let obj_ref = obj.to_object_ref(self.dyntype.clone());
151                    // TODO (matei): should this take a timeout to log when backpressure has
152                    // been applied for too long, e.g. 10s
153                    dispatcher.broadcast(obj_ref).await;
154                }
155
156                watcher::Event::InitDone => {
157                    let obj_refs: Vec<_> = {
158                        let store = self.store.read();
159                        store.keys().cloned().collect()
160                    };
161
162                    for obj_ref in obj_refs {
163                        dispatcher.broadcast(obj_ref).await;
164                    }
165                }
166
167                _ => {}
168            }
169        }
170    }
171}
172
173impl<K> Default for Writer<K>
174where
175    K: Lookup + Clone + 'static,
176    K::DynamicType: Default + Eq + Hash + Clone,
177{
178    fn default() -> Self {
179        Self::new(K::DynamicType::default())
180    }
181}
182
183/// A readable cache of Kubernetes objects of kind `K`
184///
185/// Cloning will produce a new reference to the same backing store.
186///
187/// Cannot be constructed directly since one writer handle is required,
188/// use `Writer::as_reader()` instead.
189#[derive(Educe)]
190#[educe(Debug(bound("K: Debug, K::DynamicType: Debug")), Clone)]
191pub struct Store<K: 'static + Lookup>
192where
193    K::DynamicType: Hash + Eq,
194{
195    store: Cache<K>,
196    ready_rx: Arc<DelayedInit<()>>,
197}
198
199/// The error returned by `Store::wait_until_ready`
200#[derive(Debug, Error)]
201#[error("writer was dropped before store became ready")]
202pub struct WriterDropped(delayed_init::InitDropped);
203
204impl<K: 'static + Clone + Lookup> Store<K>
205where
206    K::DynamicType: Eq + Hash + Clone,
207{
208    /// Wait for the store to be populated by Kubernetes.
209    ///
210    /// Note that polling this will _not_ await the source of the stream that populates the [`Writer`].
211    /// The [`reflector`](crate::reflector()) stream must be awaited separately.
212    ///
213    /// # Errors
214    /// Returns an error if the [`Writer`] was dropped before any value was written.
215    pub async fn wait_until_ready(&self) -> Result<(), WriterDropped> {
216        self.ready_rx.get().await.map_err(WriterDropped)
217    }
218
219    /// Retrieve a `clone()` of the entry referred to by `key`, if it is in the cache.
220    ///
221    /// `key.namespace` is ignored for cluster-scoped resources.
222    ///
223    /// Note that this is a cache and may be stale. Deleted objects may still exist in the cache
224    /// despite having been deleted in the cluster, and new objects may not yet exist in the cache.
225    /// If any of these are a problem for you then you should abort your reconciler and retry later.
226    /// If you use `kube_rt::controller` then you can do this by returning an error and specifying a
227    /// reasonable `error_policy`.
228    #[must_use]
229    pub fn get(&self, key: &ObjectRef<K>) -> Option<Arc<K>> {
230        let store = self.store.read();
231        store
232            .get(key)
233            // Try to erase the namespace and try again, in case the object is cluster-scoped
234            .or_else(|| {
235                store.get(&{
236                    let mut cluster_key = key.clone();
237                    cluster_key.namespace = None;
238                    cluster_key
239                })
240            })
241            // Clone to let go of the entry lock ASAP
242            .cloned()
243    }
244
245    /// Return a full snapshot of the current values
246    #[must_use]
247    pub fn state(&self) -> Vec<Arc<K>> {
248        let s = self.store.read();
249        s.values().cloned().collect()
250    }
251
252    /// Retrieve a `clone()` of the entry found by the given predicate
253    #[must_use]
254    pub fn find<P>(&self, predicate: P) -> Option<Arc<K>>
255    where
256        P: Fn(&K) -> bool,
257    {
258        self.store
259            .read()
260            .values()
261            .find(|k| predicate(k.as_ref()))
262            .cloned()
263    }
264
265    /// Return a filtered snapshot of the current values, retaining only objects matching `predicate`
266    ///
267    /// Note: the store read lock is held for the entire duration of predicate evaluation.
268    /// Avoid blocking or expensive operations inside the predicate.
269    #[must_use]
270    pub fn state_filter<P>(&self, predicate: P) -> Vec<Arc<K>>
271    where
272        P: Fn(&K) -> bool,
273    {
274        self.store
275            .read()
276            .values()
277            .filter(|k| predicate(k.as_ref()))
278            .cloned()
279            .collect()
280    }
281
282    /// Return a filtered snapshot of the current values, retaining only objects whose labels match `selector`
283    ///
284    /// # Example
285    ///
286    /// ```
287    /// # use k8s_openapi::api::core::v1::ConfigMap;
288    /// # use kube_client::core::{Expression, Selector};
289    /// # let (reader, _writer) = kube_runtime::reflector::store::<ConfigMap>();
290    /// let selector: Selector = Expression::Equal("app".into(), "nginx".into()).into();
291    /// let result = reader.state_filter_selector(&selector);
292    /// ```
293    #[must_use]
294    pub fn state_filter_selector(&self, selector: &Selector) -> Vec<Arc<K>>
295    where
296        K: ResourceExt,
297    {
298        self.state_filter(|k| selector.matches(k.labels()))
299    }
300
301    /// Return the number of elements in the store
302    #[must_use]
303    pub fn len(&self) -> usize {
304        self.store.read().len()
305    }
306
307    /// Return whether the store is empty
308    #[must_use]
309    pub fn is_empty(&self) -> bool {
310        self.store.read().is_empty()
311    }
312}
313
314/// Create a (Reader, Writer) for a `Store<K>` for a typed resource `K`
315///
316/// The `Writer` should be passed to a [`reflector`](crate::reflector()),
317/// and the [`Store`] is a read-only handle.
318#[must_use]
319pub fn store<K>() -> (Store<K>, Writer<K>)
320where
321    K: Lookup + Clone + 'static,
322    K::DynamicType: Eq + Hash + Clone + Default,
323{
324    let w = Writer::<K>::default();
325    let r = w.as_reader();
326    (r, w)
327}
328
329/// Create a (Reader, Writer) for a `Store<K>` for a typed resource `K`
330///
331/// The resulting `Writer` can be subscribed on in order to fan out events from
332/// a watcher. The `Writer` should be passed to a [`reflector`](crate::reflector()),
333/// and the [`Store`] is a read-only handle.
334///
335/// A buffer size is used for the underlying message channel. When the buffer is
336/// full, backpressure will be applied by waiting for capacity.
337#[must_use]
338#[cfg(feature = "unstable-runtime-subscribe")]
339pub fn store_shared<K>(buf_size: usize) -> (Store<K>, Writer<K>)
340where
341    K: Lookup + Clone + 'static,
342    K::DynamicType: Eq + Hash + Clone + Default,
343{
344    let w = Writer::<K>::new_shared(buf_size, Default::default());
345    let r = w.as_reader();
346    (r, w)
347}
348
349#[cfg(test)]
350mod tests {
351    use super::{Writer, store};
352    use crate::{reflector::ObjectRef, watcher};
353    use k8s_openapi::api::core::v1::ConfigMap;
354    use kube_client::api::ObjectMeta;
355
356    #[test]
357    fn should_allow_getting_namespaced_object_by_namespaced_ref() {
358        let cm = ConfigMap {
359            metadata: ObjectMeta {
360                name: Some("obj".to_string()),
361                namespace: Some("ns".to_string()),
362                ..ObjectMeta::default()
363            },
364            ..ConfigMap::default()
365        };
366        let mut store_w = Writer::default();
367        store_w.apply_watcher_event(&watcher::Event::Apply(cm.clone()));
368        let store = store_w.as_reader();
369        assert_eq!(store.get(&ObjectRef::from_obj(&cm)).as_deref(), Some(&cm));
370    }
371
372    #[test]
373    fn should_not_allow_getting_namespaced_object_by_clusterscoped_ref() {
374        let cm = ConfigMap {
375            metadata: ObjectMeta {
376                name: Some("obj".to_string()),
377                namespace: Some("ns".to_string()),
378                ..ObjectMeta::default()
379            },
380            ..ConfigMap::default()
381        };
382        let mut cluster_cm = cm.clone();
383        cluster_cm.metadata.namespace = None;
384        let mut store_w = Writer::default();
385        store_w.apply_watcher_event(&watcher::Event::Apply(cm));
386        let store = store_w.as_reader();
387        assert_eq!(store.get(&ObjectRef::from_obj(&cluster_cm)), None);
388    }
389
390    #[test]
391    fn should_allow_getting_clusterscoped_object_by_clusterscoped_ref() {
392        let cm = ConfigMap {
393            metadata: ObjectMeta {
394                name: Some("obj".to_string()),
395                namespace: None,
396                ..ObjectMeta::default()
397            },
398            ..ConfigMap::default()
399        };
400        let (store, mut writer) = store();
401        writer.apply_watcher_event(&watcher::Event::Apply(cm.clone()));
402        assert_eq!(store.get(&ObjectRef::from_obj(&cm)).as_deref(), Some(&cm));
403    }
404
405    #[test]
406    fn should_allow_getting_clusterscoped_object_by_namespaced_ref() {
407        let cm = ConfigMap {
408            metadata: ObjectMeta {
409                name: Some("obj".to_string()),
410                namespace: None,
411                ..ObjectMeta::default()
412            },
413            ..ConfigMap::default()
414        };
415        let mut nsed_cm = cm.clone();
416        nsed_cm.metadata.namespace = Some("ns".to_string());
417        let mut store_w = Writer::default();
418        store_w.apply_watcher_event(&watcher::Event::Apply(cm.clone()));
419        let store = store_w.as_reader();
420        assert_eq!(store.get(&ObjectRef::from_obj(&nsed_cm)).as_deref(), Some(&cm));
421    }
422
423    #[test]
424    fn state_filter_filters_by_predicate() {
425        let (reader, mut writer) = store::<ConfigMap>();
426
427        let cm1 = ConfigMap {
428            metadata: ObjectMeta {
429                name: Some("cm1".to_string()),
430                namespace: Some("ns".to_string()),
431                labels: Some([("app".to_string(), "nginx".to_string())].into()),
432                ..ObjectMeta::default()
433            },
434            ..ConfigMap::default()
435        };
436        let cm2 = ConfigMap {
437            metadata: ObjectMeta {
438                name: Some("cm2".to_string()),
439                namespace: Some("ns".to_string()),
440                labels: Some([("app".to_string(), "postgres".to_string())].into()),
441                ..ObjectMeta::default()
442            },
443            ..ConfigMap::default()
444        };
445
446        writer.apply_watcher_event(&watcher::Event::Apply(cm1.clone()));
447        writer.apply_watcher_event(&watcher::Event::Apply(cm2));
448
449        let result = reader.state_filter(|k| {
450            k.metadata
451                .labels
452                .as_ref()
453                .and_then(|l| l.get("app"))
454                .is_some_and(|v| v == "nginx")
455        });
456        assert_eq!(result.len(), 1);
457        assert_eq!(result[0].as_ref(), &cm1);
458    }
459
460    #[test]
461    fn state_filter_selector_filters_by_label_selector() {
462        use kube_client::core::{Expression, Selector};
463
464        let (reader, mut writer) = store::<ConfigMap>();
465
466        let cm1 = ConfigMap {
467            metadata: ObjectMeta {
468                name: Some("cm1".to_string()),
469                namespace: Some("ns".to_string()),
470                labels: Some([("app".to_string(), "nginx".to_string())].into()),
471                ..ObjectMeta::default()
472            },
473            ..ConfigMap::default()
474        };
475        let cm2 = ConfigMap {
476            metadata: ObjectMeta {
477                name: Some("cm2".to_string()),
478                namespace: Some("ns".to_string()),
479                labels: Some([("app".to_string(), "postgres".to_string())].into()),
480                ..ObjectMeta::default()
481            },
482            ..ConfigMap::default()
483        };
484
485        writer.apply_watcher_event(&watcher::Event::Apply(cm1.clone()));
486        writer.apply_watcher_event(&watcher::Event::Apply(cm2));
487
488        let selector: Selector = Expression::Equal("app".into(), "nginx".into()).into();
489        let result = reader.state_filter_selector(&selector);
490        assert_eq!(result.len(), 1);
491        assert_eq!(result[0].as_ref(), &cm1);
492    }
493
494    #[test]
495    fn find_element_in_store() {
496        let cm = ConfigMap {
497            metadata: ObjectMeta {
498                name: Some("obj".to_string()),
499                namespace: None,
500                ..ObjectMeta::default()
501            },
502            ..ConfigMap::default()
503        };
504        let mut target_cm = cm.clone();
505
506        let (reader, mut writer) = store::<ConfigMap>();
507        assert!(reader.is_empty());
508        writer.apply_watcher_event(&watcher::Event::Apply(cm));
509
510        assert_eq!(reader.len(), 1);
511        assert!(reader.find(|k| k.metadata.generation == Some(1234)).is_none());
512
513        target_cm.metadata.name = Some("obj1".to_string());
514        target_cm.metadata.generation = Some(1234);
515        writer.apply_watcher_event(&watcher::Event::Apply(target_cm.clone()));
516        assert!(!reader.is_empty());
517        assert_eq!(reader.len(), 2);
518        let found = reader.find(|k| k.metadata.generation == Some(1234));
519        assert_eq!(found.as_deref(), Some(&target_cm));
520    }
521}