rx_rust/operators/creating/create.rs
1use crate::utils::types::MaybeSend;
2use crate::{
3 disposable::Disposable,
4 observable::{Observable, Subscription},
5 observer::{Observer, boxed_observer::BoxedObserver},
6};
7use educe::Educe;
8
9/// Creates an Observable from scratch by means of a producer function.
10/// See <https://reactivex.io/documentation/operators/create.html>
11///
12/// # Examples
13/// ```rust
14/// use rx_rust::{
15/// observable::Subscription,
16/// observable::ObservableExt,
17/// observer::{boxed_observer::BoxedObserver, Observer, Termination},
18/// operators::creating::create::Create,
19/// };
20///
21/// let mut values = Vec::new();
22/// let mut terminations = Vec::new();
23///
24/// let observable = Create::new(|mut observer: BoxedObserver<'_, i32, ()>| {
25/// observer.on_next(42);
26/// observer.on_termination(Termination::Completed);
27/// Subscription::default()
28/// });
29///
30/// observable.subscribe_with_callback(
31/// |value| values.push(value),
32/// |termination| terminations.push(termination),
33/// );
34///
35/// assert_eq!(values, vec![42]);
36/// assert_eq!(terminations, vec![Termination::Completed]);
37/// ```
38#[derive(Educe)]
39#[educe(Debug, Clone)]
40pub struct Create<F>(F);
41
42impl<F> Create<F> {
43 pub fn new<'or, T, E, D>(builder: F) -> Self
44 where
45 // Using `Subscription` instead of FnOnce() to make `Create` more easy to wrap other observables. See more in `test_unsubscribe_wrap_observable`.
46 D: Disposable,
47 F: FnOnce(BoxedObserver<'or, T, E>) -> Subscription<D>,
48 {
49 Self(builder)
50 }
51}
52
53impl<'or, T, E, F, D> Observable<'or, T, E> for Create<F>
54where
55 D: Disposable,
56 F: FnOnce(BoxedObserver<'or, T, E>) -> Subscription<D>,
57{
58 type D = D;
59
60 fn subscribe(self, observer: impl Observer<T, E> + MaybeSend + 'or) -> Subscription<Self::D> {
61 self.0(BoxedObserver::new(observer))
62 }
63}