Skip to main content

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}