rx_rust/operators/creating/repeat.rs
1use crate::operators::creating::from_iter::FromIter;
2use crate::utils::types::MaybeSend;
3use crate::{
4 observable::{Observable, Subscription},
5 observer::Observer,
6};
7use educe::Educe;
8use std::convert::Infallible;
9
10/// Creates an Observable that emits a particular item multiple times.
11/// See <https://reactivex.io/documentation/operators/repeat.html>
12///
13/// # Examples
14/// ```rust
15/// use rx_rust::{
16/// observable::ObservableExt,
17/// observer::Termination,
18/// operators::creating::repeat::Repeat,
19/// };
20///
21/// let mut values = Vec::new();
22/// let mut terminations = Vec::new();
23///
24/// Repeat::new("ping", 3).subscribe_with_callback(
25/// |value| values.push(value),
26/// |termination| terminations.push(termination),
27/// );
28///
29/// assert_eq!(values, vec!["ping", "ping", "ping"]);
30/// assert_eq!(terminations, vec![Termination::Completed]);
31/// ```
32#[derive(Educe)]
33#[educe(Debug, Clone)]
34pub struct Repeat<T> {
35 value: T,
36 n: usize,
37}
38
39impl<T> Repeat<T> {
40 pub fn new(value: T, n: usize) -> Self
41 where
42 T: Clone,
43 {
44 Self { value, n }
45 }
46}
47
48impl<'or, T> Observable<'or, T, Infallible> for Repeat<T>
49where
50 T: Clone,
51{
52 type D = ();
53
54 fn subscribe(
55 self,
56 observer: impl Observer<T, Infallible> + MaybeSend + 'or,
57 ) -> Subscription<Self::D> {
58 FromIter::new(std::iter::repeat_n(self.value, self.n)).subscribe(observer)
59 }
60}