rx_rust/operators/others/
hook_on_next.rs1use crate::utils::types::MaybeSend;
2use crate::{
3 observable::{Observable, Subscription},
4 observer::{Flow, Observer, Termination},
5};
6use educe::Educe;
7
8#[derive(Educe)]
36#[educe(Debug, Clone)]
37pub struct HookOnNext<OE, F> {
38 source: OE,
39 callback: F,
40}
41
42impl<OE, F> HookOnNext<OE, F> {
43 pub fn new<'or, T, E>(source: OE, callback: F) -> Self
44 where
45 OE: Observable<'or, T, E>,
46 F: FnMut(&mut dyn Observer<T, E>, T) -> Flow,
47 {
48 Self { source, callback }
49 }
50}
51
52impl<'or, T, E, OE, F> Observable<'or, T, E> for HookOnNext<OE, F>
53where
54 OE: Observable<'or, T, E>,
55 F: FnMut(&mut dyn Observer<T, E>, T) -> Flow + MaybeSend + 'or,
56{
57 type D = OE::D;
58
59 fn subscribe(self, observer: impl Observer<T, E> + MaybeSend + 'or) -> Subscription<Self::D> {
60 let observer = HookOnNextObserver {
61 observer,
62 callback: self.callback,
63 };
64 self.source.subscribe(observer)
65 }
66}
67
68struct HookOnNextObserver<OR, F> {
69 observer: OR,
70 callback: F,
71}
72
73impl<T, E, OR, F> Observer<T, E> for HookOnNextObserver<OR, F>
74where
75 OR: Observer<T, E>,
76 F: FnMut(&mut dyn Observer<T, E>, T) -> Flow,
77{
78 fn on_next(&mut self, value: T) -> Flow {
79 (self.callback)(&mut self.observer, value)
80 }
81
82 fn on_termination(self, termination: Termination<E>) {
83 self.observer.on_termination(termination);
84 }
85}