-
Notifications
You must be signed in to change notification settings - Fork 24
Commit
This commit does not belong to any branch on this repository, and may belong to a fork outside of the repository.
- Loading branch information
1 parent
a6db1fd
commit ea672e7
Showing
10 changed files
with
192 additions
and
159 deletions.
There are no files selected for viewing
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file was deleted.
Oops, something went wrong.
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
This file contains bidirectional Unicode text that may be interpreted or compiled differently than what appears below. To review, open the file in an editor that reveals hidden Unicode characters.
Learn more about bidirectional Unicode characters
Original file line number | Diff line number | Diff line change |
---|---|---|
@@ -0,0 +1,65 @@ | ||
// ReactivePlusPlus library | ||
// | ||
// Copyright Aleksey Loginov 2022 - present. | ||
// Distributed under the Boost Software License, Version 1.0. | ||
// (See accompanying file LICENSE_1_0.txt or copy at | ||
// https://www.boost.org/LICENSE_1_0.txt) | ||
// | ||
// Project home: https://github.com/victimsnino/ReactivePlusPlus | ||
|
||
#pragma once | ||
|
||
#include <rpp/subscribers/constraints.hpp> | ||
#include <rpp/operators/details/subscriber_with_state.hpp> | ||
#include <rpp/subscriptions/composite_subscription.hpp> | ||
|
||
#include <memory> | ||
#include <mutex> | ||
|
||
namespace rpp::details | ||
{ | ||
struct forwarding_on_next_under_lock | ||
{ | ||
template<typename T, typename TSerializationPrimitive> | ||
void operator()(T&& v, const auto& subscriber, const std::shared_ptr<TSerializationPrimitive>& primitive) const | ||
{ | ||
std::lock_guard lock{*primitive}; | ||
subscriber.on_next(std::forward<T>(v)); | ||
} | ||
}; | ||
|
||
struct forwarding_on_error_under_lock | ||
{ | ||
template<typename TSerializationPrimitive> | ||
void operator()(const std::exception_ptr& err, | ||
const auto& subscriber, | ||
const std::shared_ptr<TSerializationPrimitive>& primitive) const | ||
{ | ||
std::lock_guard lock{*primitive}; | ||
subscriber.on_error(err); | ||
} | ||
}; | ||
|
||
struct forwarding_on_completed_under_lock | ||
{ | ||
template<typename TSerializationPrimitive> | ||
void operator()(const auto& subscriber, const std::shared_ptr<TSerializationPrimitive>& primitive) const | ||
{ | ||
std::lock_guard lock{*primitive}; | ||
subscriber.on_completed(); | ||
} | ||
}; | ||
|
||
template<typename TSerializationPrimitive, constraint::subscriber TSub> | ||
auto make_serialized_subscriber(TSub&& subscriber, | ||
const std::shared_ptr<TSerializationPrimitive>& primitive) | ||
{ | ||
auto sub = subscriber.get_subscription(); | ||
return create_subscriber_with_state<utils::extract_subscriber_type_t<std::decay_t<TSub>>>(std::move(sub), | ||
forwarding_on_next_under_lock{}, | ||
forwarding_on_error_under_lock{}, | ||
forwarding_on_completed_under_lock{}, | ||
std::forward<TSub>(subscriber), | ||
primitive); | ||
} | ||
} // namespace rpp::details |
Oops, something went wrong.
ea672e7
There was a problem hiding this comment.
Choose a reason for hiding this comment
The reason will be displayed to describe this comment to others. Learn more.
BENCHMARK RESULTS (AUTOGENERATED)
ci-ubuntu-clang
Observable construction
Table
Observable lift
Table
Observable subscribe
Table
Observable subscribe #2
Table
Observer construction
Table
OnNext
Table
Subscriber construction
Table
Subscription
Table
buffer
Table
chains creation test
Table
combine_latest
Table
concat
Table
distinct_until_changed
Table
first
Table
foundamental sources
Table
from
Table
immediate scheduler
Table
just
Table
last
Table
map
Table
merge
Table
observe_on
Table
publish_subject callbacks
Table
publish_subject routines
Table
repeat
Table
scan
Table
skip
Table
switch_on_next
Table
take
Table
take_last
Table
take_until
Table
trampoline scheduler
Table
window
Table
with_latest_from
Table
ci-ubuntu-gcc
Observable construction
Table
Observable lift
Table
Observable subscribe
Table
Observable subscribe #2
Table
Observer construction
Table
OnNext
Table
Subscriber construction
Table
Subscription
Table
buffer
Table
chains creation test
Table
combine_latest
Table
concat
Table
distinct_until_changed
Table
first
Table
foundamental sources
Table
from
Table
immediate scheduler
Table
just
Table
last
Table
map
Table
merge
Table
observe_on
Table
publish_subject callbacks
Table
publish_subject routines
Table
repeat
Table
scan
Table
skip
Table
switch_on_next
Table
take
Table
take_last
Table
take_until
Table
trampoline scheduler
Table
window
Table
with_latest_from
Table
ci-windows
Observable construction
Table
Observable lift
Table
Observable subscribe
Table
Observable subscribe #2
Table
Observer construction
Table
OnNext
Table
Subscriber construction
Table
Subscription
Table
buffer
Table
chains creation test
Table
combine_latest
Table
concat
Table
distinct_until_changed
Table
first
Table
foundamental sources
Table
from
Table
immediate scheduler
Table
just
Table
last
Table
map
Table
merge
Table
observe_on
Table
publish_subject callbacks
Table
publish_subject routines
Table
repeat
Table
scan
Table
skip
Table
switch_on_next
Table
take
Table
take_last
Table
take_until
Table
trampoline scheduler
Table
window
Table
with_latest_from
Table