1use 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#[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 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 #[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 #[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 #[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 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 std::mem::swap(&mut *store, &mut self.buffer);
131
132 self.buffer = AHashMap::new();
136
137 if let Some(ready_tx) = self.ready_tx.take() {
139 ready_tx.init(())
140 }
141 }
142 }
143 }
144
145 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 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#[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#[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 pub async fn wait_until_ready(&self) -> Result<(), WriterDropped> {
216 self.ready_rx.get().await.map_err(WriterDropped)
217 }
218
219 #[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 .or_else(|| {
235 store.get(&{
236 let mut cluster_key = key.clone();
237 cluster_key.namespace = None;
238 cluster_key
239 })
240 })
241 .cloned()
243 }
244
245 #[must_use]
247 pub fn state(&self) -> Vec<Arc<K>> {
248 let s = self.store.read();
249 s.values().cloned().collect()
250 }
251
252 #[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 #[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 #[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 #[must_use]
303 pub fn len(&self) -> usize {
304 self.store.read().len()
305 }
306
307 #[must_use]
309 pub fn is_empty(&self) -> bool {
310 self.store.read().is_empty()
311 }
312}
313
314#[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#[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}