rx_rust/operators/transforming/
buffer_with_time.rs1use 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#[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}