Skip to main content

rx_rust/operators/transforming/
buffer_with_time.rs

1use crate::disposable::{Disposable, bound_drop_disposal::BoundDropDisposal};
2use crate::utils::serialized_delivery::UpdateOutcome;
3use crate::utils::subscribe_with_context::{
4    self, SubscriptionContext, subscribe_with_context_owning_source,
5};
6use crate::utils::types::{MarkerType, MaybeSend};
7use crate::{
8    observable::Observable,
9    observable::Subscription,
10    observer::{Flow, Observer, Termination},
11    scheduler::Scheduler,
12};
13use educe::Educe;
14use std::time::Duration;
15
16/// Periodically gathers items from an Observable into bundles and emits these bundles as `Vec<T>`, after a specified time interval.
17/// See <https://reactivex.io/documentation/operators/buffer.html>
18///
19/// # Examples
20/// ```rust
21/// # #[cfg(not(feature = "tokio-scheduler"))]
22/// # fn main() {}
23/// # #[cfg(feature = "tokio-scheduler")]
24/// #[tokio::main]
25/// async fn main() {
26///     use rx_rust::{
27///         observable::ObservableExt,
28///         observer::Termination,
29///         operators::{
30///             creating::from_iter::FromIter,
31///             transforming::buffer_with_time::BufferWithTime,
32///         },
33///     };
34///     use std::sync::{Arc, Mutex};
35///     use std::time::Duration;
36///     use tokio::time::sleep;
37///
38///     let handle = tokio::runtime::Handle::current();
39///     let values = Arc::new(Mutex::new(Vec::new()));
40///     let terminations = Arc::new(Mutex::new(Vec::new()));
41///     let values_observer = Arc::clone(&values);
42///     let terminations_observer = Arc::clone(&terminations);
43///
44///     let subscription = BufferWithTime::new(
45///         FromIter::new(vec![1, 2, 3]),
46///         Duration::from_millis(5),
47///         handle.clone(),
48///         None,
49///     )
50///     .subscribe_with_callback(
51///         move |value| values_observer.lock().unwrap().push(value),
52///         move |termination| terminations_observer
53///             .lock()
54///             .unwrap()
55///             .push(termination),
56///     );
57///
58///     sleep(Duration::from_millis(10)).await;
59///     drop(subscription);
60///
61///     assert_eq!(&*values.lock().unwrap(), &[vec![1, 2, 3]]);
62///     assert_eq!(
63///         &*terminations.lock().unwrap(),
64///         &[Termination::Completed]
65///     );
66/// }
67/// ```
68#[derive(Educe)]
69#[educe(Debug, Clone)]
70pub struct BufferWithTime<'or, OE, S> {
71    source: OE,
72    time_span: Duration,
73    scheduler: S,
74    delay: Option<Duration>,
75    _marker: MarkerType<&'or ()>,
76}
77
78impl<'or, OE, S> BufferWithTime<'or, OE, S> {
79    pub fn new(source: OE, time_span: Duration, scheduler: S, delay: Option<Duration>) -> Self {
80        Self {
81            source,
82            time_span,
83            scheduler,
84            delay,
85            _marker: Default::default(),
86        }
87    }
88}
89
90impl<'or, T, E, OE, S> Observable<'static, Vec<T>, E> for BufferWithTime<'or, OE, S>
91where
92    T: MaybeSend + 'static,
93    E: MaybeSend + 'static,
94    OE: Observable<'or, T, E>,
95    OE::D: MaybeSend + 'static,
96    S: Scheduler + Clone + MaybeSend + 'static,
97{
98    type D = subscribe_with_context::OwningDisposal<'or>;
99
100    fn subscribe(
101        self,
102        observer: impl Observer<Vec<T>, E> + MaybeSend + 'static,
103    ) -> Subscription<Self::D> {
104        subscribe_with_context_owning_source(observer, Vec::new(), |context| {
105            let sub = self
106                .source
107                .subscribe(BufferWithTimeObserver(context.clone()));
108            let disposal = setup_emit_timer(context, self.scheduler, self.time_span, self.delay);
109            sub.preceded_by_bound(disposal)
110        })
111    }
112}
113
114struct BufferWithTimeObserver<T, E, OR, D: Disposable>(
115    SubscriptionContext<Vec<T>, E, OR, Vec<T>, D>,
116);
117
118impl<T, E, OR, D> Observer<T, E> for BufferWithTimeObserver<T, E, OR, D>
119where
120    OR: Observer<Vec<T>, E>,
121    D: Disposable + MaybeSend + 'static,
122{
123    fn on_next(&mut self, value: T) -> Flow {
124        self.0.update_flow(|values| {
125            values.push(value);
126            UpdateOutcome::empty().without_events()
127        })
128    }
129
130    fn on_termination(self, termination: Termination<E>) {
131        match termination {
132            completion @ Termination::Completed => {
133                let _ = self.0.update(|values| {
134                    if !values.is_empty() {
135                        UpdateOutcome::empty()
136                            .with_next_and_termination_events(std::mem::take(values), completion)
137                    } else {
138                        UpdateOutcome::empty().with_termination_event(completion)
139                    }
140                });
141            }
142            error @ Termination::Error(_) => {
143                self.0.send_termination(error);
144            }
145        }
146    }
147}
148
149fn setup_emit_timer<T, E, OR, D, S>(
150    context: SubscriptionContext<Vec<T>, E, OR, Vec<T>, D>,
151    scheduler: S,
152    time_span: Duration,
153    delay: Option<Duration>,
154) -> BoundDropDisposal<S::D>
155where
156    T: MaybeSend + 'static,
157    E: MaybeSend + 'static,
158    OR: Observer<Vec<T>, E> + MaybeSend + 'static,
159    D: Disposable + MaybeSend + 'static,
160    S: Scheduler + Clone + MaybeSend + 'static,
161{
162    let weak_context = context.downgrade();
163    scheduler.schedule_periodically(
164        move |_| {
165            let Some(context) = weak_context.upgrade() else {
166                return false;
167            };
168            context
169                .update(|values| {
170                    UpdateOutcome::new(true).with_next_event(std::mem::replace(
171                        values,
172                        Vec::with_capacity(values.len()),
173                    ))
174                })
175                .unwrap_or(false)
176        },
177        time_span,
178        delay,
179    )
180}