rx_rust/operators/conditional_boolean/
default_if_empty.rs1use crate::utils::types::MaybeSend;
2use crate::{
3 observable::{Observable, Subscription},
4 observer::{Flow, Observer, Termination},
5};
6use educe::Educe;
7
8#[derive(Educe)]
35#[educe(Debug, Clone)]
36pub struct DefaultIfEmpty<T, OE> {
37 source: OE,
38 default_value: T,
39}
40
41impl<T, OE> DefaultIfEmpty<T, OE> {
42 pub fn new<'or, E>(source: OE, item: T) -> Self
43 where
44 OE: Observable<'or, T, E>,
45 {
46 Self {
47 source,
48 default_value: item,
49 }
50 }
51}
52
53impl<'or, T, E, OE> Observable<'or, T, E> for DefaultIfEmpty<T, OE>
54where
55 OE: Observable<'or, T, E>,
56 T: MaybeSend + 'or,
57{
58 type D = OE::D;
59
60 fn subscribe(self, observer: impl Observer<T, E> + MaybeSend + 'or) -> Subscription<Self::D> {
61 let observer = DefaultIfEmptyObserver {
62 observer,
63 default_value: Some(self.default_value),
64 };
65 self.source.subscribe(observer)
66 }
67}
68
69struct DefaultIfEmptyObserver<T, OR> {
70 observer: OR,
71 default_value: Option<T>, }
73
74impl<T, E, OR> Observer<T, E> for DefaultIfEmptyObserver<T, OR>
75where
76 OR: Observer<T, E>,
77{
78 fn on_next(&mut self, value: T) -> Flow {
79 self.default_value = None;
80 self.observer.on_next(value)
81 }
82
83 fn on_termination(mut self, termination: Termination<E>) {
84 if matches!(termination, Termination::Completed)
87 && let Some(default_value) = self.default_value.take()
88 && self.observer.on_next(default_value).is_stop()
89 {
90 return;
91 }
92 self.observer.on_termination(termination);
93 }
94}