MCPcopy Create free account
hub / github.com/AlexInLog/ReactivePlusPlus / main

Function main

src/examples/rpp/doxygen/observe_on.cpp:10–49  ·  view source on GitHub ↗

* @example observe_on.cpp **/

Source from the content-addressed store, hash-verified

8 * @example observe_on.cpp
9 **/
10int main() // NOLINT(bugprone-exception-escape)
11{
12 //! [observe_on]
13
14 auto start = rpp::schedulers::clock_type::now();
15
16 rpp::source::create<int>([&start](const auto& obs) {
17 for (int i = 0; i < 3; ++i)
18 {
19 auto emitting_time = rpp::schedulers::clock_type::now();
20 std::cout << "emit " << i << " in thread{" << std::this_thread::get_id() << "} duration since start " << std::chrono::duration_cast<std::chrono::seconds>(emitting_time - start).count() << "s" << std::endl;
21
22 obs.on_next(i);
23 std::this_thread::sleep_for(std::chrono::seconds{1});
24 }
25 auto emitting_time = rpp::schedulers::clock_type::now();
26 std::cout << "emit error in thread{" << std::this_thread::get_id() << "} duration since start " << std::chrono::duration_cast<std::chrono::seconds>(emitting_time - start).count() << "s" << std::endl;
27
28 obs.on_error({});
29 })
30 | rpp::operators::observe_on(rpp::schedulers::new_thread{}, std::chrono::seconds{3})
31 | rpp::operators::as_blocking()
32 | rpp::operators::subscribe([&](int v) {
33 auto observing_time = rpp::schedulers::clock_type::now();
34 std::cout << "observe " << v << " in thread{" << std::this_thread::get_id() << "} duration since start " << std::chrono::duration_cast<std::chrono::seconds>(observing_time - start).count() <<"s" << std::endl; },
35 [&](const std::exception_ptr&) {
36 auto observing_time = rpp::schedulers::clock_type::now();
37 std::cout << "observe error in thread{" << std::this_thread::get_id() << "} duration since start " << std::chrono::duration_cast<std::chrono::seconds>(observing_time - start).count() << "s" << std::endl;
38 });
39
40 // Template for output:
41 // emit 0 in thread{139800298538880} duration since start 0s
42 // emit 1 in thread{139800298538880} duration since start 1s
43 // emit 2 in thread{139800298538880} duration since start 2s
44 // observe 0 in thread{139800298534464} duration since start 3s
45 // emit error in thread{139800298538880} duration since start 3s
46 // observe error in thread{139800298538880} duration since start 3s
47 //! [observe_on]
48 return 0;
49}

Callers

nothing calls this directly

Calls 6

nowFunction · 0.85
observe_onFunction · 0.85
as_blockingFunction · 0.85
subscribeFunction · 0.85
on_nextMethod · 0.45
on_errorMethod · 0.45

Tested by

no test coverage detected