Skip to main content

rx_rust/operators/creating/
defer.rs

1use crate::utils::types::MaybeSend;
2use crate::{
3    observable::{Observable, Subscription},
4    observer::Observer,
5};
6use educe::Educe;
7
8/// Do not create the Observable until a Observer subscribes, and create a fresh Observable for each Observer.
9/// See <https://reactivex.io/documentation/operators/defer.html>
10///
11/// # Examples
12/// ```rust
13/// use rx_rust::{
14///     observable::ObservableExt,
15///     observer::Termination,
16///     operators::creating::{defer::Defer, just::Just},
17/// };
18///
19/// let mut values = Vec::new();
20/// let mut terminations = Vec::new();
21///
22/// let observable = Defer::new(|| Just::new(5));
23/// observable.subscribe_with_callback(
24///     |value| values.push(value),
25///     |termination| terminations.push(termination),
26/// );
27///
28/// assert_eq!(values, vec![5]);
29/// assert_eq!(terminations, vec![Termination::Completed]);
30/// ```
31#[derive(Educe)]
32#[educe(Debug, Clone)]
33pub struct Defer<F>(F);
34
35impl<F> Defer<F> {
36    pub fn new<OE>(builder: F) -> Self
37    where
38        F: FnOnce() -> OE,
39    {
40        Self(builder)
41    }
42}
43
44impl<'or, T, E, OE, F> Observable<'or, T, E> for Defer<F>
45where
46    F: FnOnce() -> OE,
47    OE: Observable<'or, T, E>,
48{
49    type D = OE::D;
50
51    fn subscribe(self, observer: impl Observer<T, E> + MaybeSend + 'or) -> Subscription<Self::D> {
52        let observable = self.0();
53        observable.subscribe(observer)
54    }
55}