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}