ReactivePlusPlus
ReactiveX implementation for C++20
Toggle main menu visibility
Loading...
Searching...
No Matches
take_until.hpp
1
// ReactivePlusPlus library
2
//
3
// Copyright Aleksey Loginov 2023 - present.
4
// Distributed under the Boost Software License, Version 1.0.
5
// (See accompanying file LICENSE_1_0.txt or copy at
6
// https://www.boost.org/LICENSE_1_0.txt)
7
//
8
// Project home: https://github.com/AlexInLog/ReactivePlusPlus
9
//
10
11
#pragma once
12
13
#include <rpp/operators/fwd.hpp>
14
15
#include <rpp/defs.hpp>
16
#include <rpp/disposables/composite_disposable.hpp>
17
#include <rpp/schedulers/current_thread.hpp>
18
#include <rpp/utils/utils.hpp>
19
20
namespace
rpp::operators::details
21
{
22
template
<rpp::constra
int
::observer TObserver>
23
class
take_until_disposable final :
public
rpp::composite_disposable
24
{
25
public
:
26
take_until_disposable(TObserver&&
observer
)
27
: m_observer_with_mutex(std::move(
observer
))
28
{
29
}
30
31
take_until_disposable(
const
TObserver&
observer
)
32
: m_observer_with_mutex(
observer
)
33
{
34
}
35
36
bool
is_stopped()
const
{
return
m_stopped; }
37
bool
stop_return_was_stopped() {
return
m_stopped.exchange(
true
); }
38
39
rpp::utils::pointer_under_lock<TObserver> get_observer() {
return
m_observer_with_mutex; }
40
41
private
:
42
rpp::utils::value_with_mutex<TObserver>
m_observer_with_mutex{};
43
std::atomic_bool m_stopped{};
44
};
45
46
template
<rpp::constra
int
::observer TObserver>
47
struct
take_until_observer_strategy_base
48
{
49
static
constexpr
auto
preferred_disposables_mode = rpp::details::observers::disposables_mode::Auto;
50
51
std::shared_ptr<take_until_disposable<TObserver>> state;
52
53
void
on_error(
const
std::exception_ptr& err)
const
54
{
55
if
(!state->stop_return_was_stopped())
56
state->get_observer()->on_error(err);
57
}
58
59
void
on_completed()
const
60
{
61
if
(!state->stop_return_was_stopped())
62
state->get_observer()->on_completed();
63
}
64
65
void
set_upstream(
const
disposable_wrapper
& d) { state->add(d); }
66
67
bool
is_disposed()
const
{
return
state->is_disposed(); }
68
};
69
70
template
<rpp::constra
int
::observer TObserver>
71
struct
take_until_throttle_observer_strategy
:
public
take_until_observer_strategy_base
<TObserver>
72
{
73
template
<
typename
T>
74
void
on_next(
const
T&)
const
75
{
76
if
(!take_until_observer_strategy_base<TObserver>::state->stop_return_was_stopped())
77
take_until_observer_strategy_base<TObserver>::state->get_observer()->on_completed();
78
}
79
};
80
81
template
<rpp::constra
int
::observer TObserver>
82
struct
take_until_observer_strategy
:
public
take_until_observer_strategy_base
<TObserver>
83
{
84
template
<
typename
T>
85
void
on_next(T&& v)
const
86
{
87
if
(!take_until_observer_strategy_base<TObserver>::state->is_stopped())
88
take_until_observer_strategy_base<TObserver>::state->get_observer()->on_next(std::forward<T>(v));
89
}
90
};
91
92
template
<rpp::constra
int
::observable TObservable>
93
struct
take_until_t
94
{
95
RPP_NO_UNIQUE_ADDRESS TObservable observable{};
96
97
template
<rpp::constra
int
::decayed_type T>
98
struct
operator_traits
99
{
100
using
result_type = T;
101
102
constexpr
static
bool
own_current_queue =
true
;
103
};
104
105
template
<rpp::details::observables::constra
int
::disposables_strategy Prev>
106
using
updated_optimal_disposables_strategy =
rpp::details::observables::fixed_disposables_strategy<1>
;
107
108
template
<rpp::constra
int
::decayed_type Type, rpp::constra
int
::observer Observer>
109
auto
lift(Observer&&
observer
)
const
110
{
111
const
auto
d =
disposable_wrapper_impl<take_until_disposable<std::decay_t<Observer>
>>::make(std::forward<Observer>(
observer
));
112
auto
ptr = d.lock();
113
ptr->get_observer()->set_upstream(d.as_weak());
114
115
observable.subscribe(
take_until_throttle_observer_strategy
<std::decay_t<Observer>>{ptr});
116
return
rpp::observer<Type, take_until_observer_strategy<std::decay_t<Observer>
>>(std::move(ptr));
117
}
118
};
119
}
// namespace rpp::operators::details
120
121
namespace
rpp::operators
122
{
142
*
143
* @ingroup conditional_operators
144
* @see https://reactivex.io/documentation/operators/takeuntil.html
145
*/
146
template
<rpp::constra
int
::observable TObservable>
147
auto
take_until
(TObservable&& until_observable)
148
{
149
return
details::take_until_t<std::decay_t<TObservable>
>{std::forward<TObservable>(until_observable)};
150
}
151
}
// namespace rpp::operators
rpp::composite_disposable
Disposable which can keep some other sub-disposables. When this root disposable is disposed,...
Definition
composite_disposable.hpp:175
rpp::disposable_wrapper_impl
Main RPP wrapper over disposables.
Definition
disposable_wrapper.hpp:142
rpp::observer
Base class for any observer used in RPP. It handles core callbacks of observers. Objects of this clas...
Definition
observer.hpp:172
rpp::utils::value_with_mutex
Definition
utils.hpp:260
rpp::operators::take_until
auto take_until(TObservable &&until_observable)
Discard any items emitted by an Observable after a second Observable emits an item or terminates.
Definition
take_until.hpp:142
rpp::disposable_wrapper
disposable_wrapper_impl< interface_disposable > disposable_wrapper
Wrapper to keep "simple" disposable. Specialization of rpp::disposable_wrapper_impl.
Definition
fwd.hpp:34
rpp::details::observables::fixed_disposables_strategy
Definition
disposables_strategy.hpp:29
rpp::operators::details::take_until_observer_strategy_base
Definition
take_until.hpp:48
rpp::operators::details::take_until_observer_strategy
Definition
take_until.hpp:83
rpp::operators::details::take_until_t::operator_traits
Definition
take_until.hpp:99
rpp::operators::details::take_until_t
Definition
take_until.hpp:94
rpp::operators::details::take_until_throttle_observer_strategy
Definition
take_until.hpp:72
src
rpp
rpp
operators
take_until.hpp
Generated by
1.17.0