ReactivePlusPlus
ReactiveX implementation for C++20
Loading...
Searching...
No Matches
share.cpp
#include <rpp/rpp.hpp>
#include <iostream>
int main() // NOLINT(bugprone-exception-escape)
{
{
auto observable = rpp::source::create<int>([&](auto&& observer) {
std::cout << "SUBSCRIBE" << std::endl;
subject.get_observable().subscribe(std::forward<decltype(observer)>(observer));
})
std::cout << "before subscriptions" << std::endl;
observable.subscribe([](int v) { std::cout << "#1 " << v << std::endl; });
observable.subscribe([](int v) { std::cout << "#2 " << v << std::endl; });
subject.get_observer().on_next(1);
subject.get_observer().on_completed();
// Output:
// before subscriptions
// SUBSCRIBE
// #1 1
// #2 1
}
return 0;
}
Subject which just multicasts values to observers subscribed on it. It contains two parts: observer a...
Definition publish_subject.hpp:81
auto share()
Shares single subscription to original observable between multiple observers.
Definition share.hpp:49
auto create(OnSubscribe &&on_subscribe)
Construct observable specialized with passed callback function. Most easiesest way to construct obser...
Definition create.hpp:57