1
  2
  3
  4
  5
  6
  7
  8
  9
 10
 11
 12
 13
 14
 15
 16
 17
 18
 19
 20
 21
 22
 23
 24
 25
 26
 27
 28
 29
 30
 31
 32
 33
 34
 35
 36
 37
 38
 39
 40
 41
 42
 43
 44
 45
 46
 47
 48
 49
 50
 51
 52
 53
 54
 55
 56
 57
 58
 59
 60
 61
 62
 63
 64
 65
 66
 67
 68
 69
 70
 71
 72
 73
 74
 75
 76
 77
 78
 79
 80
 81
 82
 83
 84
 85
 86
 87
 88
 89
 90
 91
 92
 93
 94
 95
 96
 97
 98
 99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
use crate::prelude::*;
use std::marker::PhantomData;

unsafe impl<F, Item, Err> Send for FnEmitter<F, Item, Err> where F: Send {}
unsafe impl<F, Item, Err> Sync for FnEmitter<F, Item, Err> where F: Sync {}
impl<F, Item, Err> Clone for FnEmitter<F, Item, Err>
where
  F: Clone,
{
  fn clone(&self) -> Self { FnEmitter(self.0.clone(), PhantomData) }
}

/// param `subscribe`: the function that is called when the Observable is
/// initially subscribed to. This function is given a Subscriber, to which
/// new values can be `next`ed, or an `error` method can be called to raise
/// an error, or `complete` can be called to notify of a successful
/// completion.
pub fn create<F, O, U, Item, Err>(
  subscribe: F,
) -> ObservableBase<FnEmitter<F, Item, Err>>
where
  F: FnOnce(Subscriber<O, U>),
  O: Observer<Item, Err>,
  U: SubscriptionLike,
{
  ObservableBase::new(FnEmitter(subscribe, PhantomData))
}

pub struct FnEmitter<F, Item, Err>(F, PhantomData<(Item, Err)>);

impl<F, Item, Err> Emitter for FnEmitter<F, Item, Err> {
  type Item = Item;
  type Err = Err;
}

impl<'a, F, Item, Err> LocalEmitter<'a> for FnEmitter<F, Item, Err>
where
  F: FnOnce(
    Subscriber<
      Box<dyn Observer<Item, Err> + 'a>,
      Box<dyn SubscriptionLike + 'a>,
    >,
  ),
{
  fn emit<O>(self, subscriber: Subscriber<O, LocalSubscription>)
  where
    O: Observer<Self::Item, Self::Err> + 'a,
  {
    (self.0)(Subscriber {
      observer: Box::new(subscriber.observer),
      subscription: Box::new(subscriber.subscription),
    })
  }
}

impl<F, Item, Err> SharedEmitter for FnEmitter<F, Item, Err>
where
  F: FnOnce(
    Subscriber<
      Box<dyn Observer<Item, Err> + Send + Sync + 'static>,
      SharedSubscription,
    >,
  ),
{
  fn emit<O>(self, subscriber: Subscriber<O, SharedSubscription>)
  where
    O: Observer<Self::Item, Self::Err> + Send + Sync + 'static,
  {
    (self.0)(Subscriber {
      observer: Box::new(subscriber.observer),
      subscription: subscriber.subscription,
    })
  }
}

#[cfg(test)]
mod test {
  use crate::prelude::*;
  use std::sync::{Arc, Mutex};

  #[test]
  fn proxy_call() {
    let next = Arc::new(Mutex::new(0));
    let err = Arc::new(Mutex::new(0));
    let complete = Arc::new(Mutex::new(0));
    let c_next = next.clone();
    let c_err = err.clone();
    let c_complete = complete.clone();

    observable::create(|mut subscriber| {
      subscriber.next(&1);
      subscriber.next(&2);
      subscriber.next(&3);
      subscriber.complete();
      subscriber.next(&3);
      subscriber.error("never dispatch error");
    })
    .to_shared()
    .subscribe_all(
      move |_| *next.lock().unwrap() += 1,
      move |_: &str| *err.lock().unwrap() += 1,
      move || *complete.lock().unwrap() += 1,
    );

    assert_eq!(*c_next.lock().unwrap(), 3);
    assert_eq!(*c_complete.lock().unwrap(), 1);
    assert_eq!(*c_err.lock().unwrap(), 0);
  }
  #[test]
  fn support_fork() {
    let o = observable::create(|mut subscriber| {
      subscriber.next(&1);
      subscriber.next(&2);
      subscriber.next(&3);
      subscriber.next(&4);
    });
    let sum1 = Arc::new(Mutex::new(0));
    let sum2 = Arc::new(Mutex::new(0));
    let c_sum1 = sum1.clone();
    let c_sum2 = sum2.clone();
    o.clone().subscribe(move |v| *sum1.lock().unwrap() += v);
    o.clone().subscribe(move |v| *sum2.lock().unwrap() += v);

    assert_eq!(*c_sum1.lock().unwrap(), 10);
    assert_eq!(*c_sum2.lock().unwrap(), 10);
  }

  #[test]
  fn fork_and_share() {
    let observable = observable::create(|_| {});
    observable.clone().to_shared().subscribe(|_: i32| {});
    observable.clone().to_shared().subscribe(|_| {});

    let observable = observable::create(|_| {}).to_shared();
    observable.clone().subscribe(|_: i32| {});
    observable.clone().subscribe(|_| {});
  }
}