rx_rust/operators/creating/from_result.rs
1use crate::utils::types::MaybeSend;
2use crate::{
3 observable::{Observable, Subscription},
4 observer::{Observer, Termination},
5};
6use educe::Educe;
7
8/// Converts a `Result` into an Observable.
9/// See <https://reactivex.io/documentation/operators/from.html>
10///
11/// # Examples
12/// ```rust
13/// use rx_rust::{
14/// observable::ObservableExt,
15/// observer::Termination,
16/// operators::creating::from_result::FromResult,
17/// };
18///
19/// let mut values = Vec::new();
20/// let mut terminations = Vec::new();
21///
22/// FromResult::new(Ok::<i32, &str>(10)).subscribe_with_callback(
23/// |value| values.push(value),
24/// |termination| terminations.push(termination),
25/// );
26///
27/// assert_eq!(values, vec![10]);
28/// assert_eq!(terminations, vec![Termination::Completed]);
29/// ```
30#[derive(Educe)]
31#[educe(Debug, Clone)]
32pub struct FromResult<T, E>(Result<T, E>);
33
34impl<T, E> FromResult<T, E> {
35 pub fn new(result: Result<T, E>) -> Self {
36 Self(result)
37 }
38}
39
40impl<'or, T, E> Observable<'or, T, E> for FromResult<T, E> {
41 type D = ();
42
43 fn subscribe(
44 self,
45 mut observer: impl Observer<T, E> + MaybeSend + 'or,
46 ) -> Subscription<Self::D> {
47 match self.0 {
48 Ok(value) => {
49 if observer.on_next(value).is_continue() {
50 observer.on_termination(Termination::Completed);
51 }
52 }
53 Err(error) => observer.on_termination(Termination::Error(error)),
54 }
55 Subscription::default()
56 }
57}