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