src/corosio/src/timer.cpp
76.2% Lines (48 / 63)
90.9% Functions (10 / 11)
Functions (11)
Function
Calls
Lines
Blocks
boost::corosio::detail::timer::~timer()
:16
17235x
100.0%
100.0%
boost::corosio::detail::timer::timer(boost::capy::execution_context&)
:18
17239x
100.0%
100.0%
boost::corosio::detail::timer::implementation::retire()
:27
2391x
100.0%
100.0%
boost::corosio::detail::timer::implementation::wait(boost::corosio::detail::waiter_node&)
:38
15395x
50.0%
44.0%
boost::corosio::detail::timer::implementation::publish(boost::corosio::detail::waiter_node&)
:53
16383x
64.3%
63.0%
boost::corosio::detail::timer::publish_wait(boost::corosio::detail::waiter_node&)
:94
988x
100.0%
100.0%
boost::corosio::detail::timer::rearm_wait(boost::corosio::detail::waiter_node&, std::chrono::duration<long, std::ratio<1l, 1000000000l> >)
:100
6397x
66.7%
70.0%
boost::corosio::detail::waiter_node::canceller::operator()() const
:129
2689x
100.0%
100.0%
boost::corosio::detail::waiter_node::completion_op::do_complete(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int)
:135
0
0.0%
0.0%
boost::corosio::detail::waiter_node::completion_op::operator()()
:149
22749x
100.0%
100.0%
boost::corosio::detail::waiter_node::completion_op::destroy()
:169
2x
100.0%
100.0%
| Line | TLA | Hits | Source Code |
|---|---|---|---|
| 1 | // | ||
| 2 | // Copyright (c) 2025 Vinnie Falco (vinnie.falco@gmail.com) | ||
| 3 | // Copyright (c) 2026 Steve Gerbino | ||
| 4 | // | ||
| 5 | // Distributed under the Boost Software License, Version 1.0. (See accompanying | ||
| 6 | // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt) | ||
| 7 | // | ||
| 8 | // Official repository: https://github.com/cppalliance/corosio | ||
| 9 | // | ||
| 10 | |||
| 11 | #include <boost/corosio/detail/timer.hpp> | ||
| 12 | #include <boost/corosio/detail/timer_service.hpp> | ||
| 13 | |||
| 14 | namespace boost::corosio::detail { | ||
| 15 | |||
| 16 | 17235x | timer::~timer() = default; | |
| 17 | |||
| 18 | 17239x | timer::timer(capy::execution_context& ctx) | |
| 19 | 17239x | : io_object(create_handle<detail::timer_service>(ctx)) | |
| 20 | { | ||
| 21 | 17235x | } | |
| 22 | |||
| 23 | // Not inline for the same reason as wait() below: retire() needs | ||
| 24 | // timer_service's complete type (to reach pool_), which timer.hpp only | ||
| 25 | // forward-declares. | ||
| 26 | void | ||
| 27 | 2391x | timer::implementation::retire() noexcept | |
| 28 | { | ||
| 29 | 2391x | svc_->pool_.recycle(this); | |
| 30 | 2391x | } | |
| 31 | |||
| 32 | // Not inline: wait_awaitable::await_suspend (defined in timer.hpp) calls | ||
| 33 | // this from translation units that may never include timer_service.hpp, | ||
| 34 | // so this must be the one strong definition the linker can always find | ||
| 35 | // wherever a detail::timer is used (every such user also needs timer's | ||
| 36 | // constructors, defined in this same translation unit). | ||
| 37 | std::coroutine_handle<> | ||
| 38 | 15395x | timer::implementation::wait(waiter_node& w) | |
| 39 | { | ||
| 40 | // Already-expired fast path — no publication, no mutex. | ||
| 41 | // Post instead of dispatch so the coroutine yields to the | ||
| 42 | // scheduler, allowing other queued work to run. | ||
| 43 | 15395x | if (already_expired()) | |
| 44 | { | ||
| 45 | ✗ | w.ec_ = {}; | |
| 46 | ✗ | w.d_.post(w.cont_); | |
| 47 | ✗ | return std::noop_coroutine(); | |
| 48 | } | ||
| 49 | 15395x | return publish(w); | |
| 50 | } | ||
| 51 | |||
| 52 | std::coroutine_handle<> | ||
| 53 | 16383x | timer::implementation::publish(waiter_node& w) | |
| 54 | { | ||
| 55 | // Publication-last invariant: fully initialize the waiter, count | ||
| 56 | // its work, and arm cancellation BEFORE insert_waiter() publishes | ||
| 57 | // it into the heap/waiter slot where a concurrent run() thread can fire | ||
| 58 | // it. impl_ stays null until insert_waiter() sets it under the | ||
| 59 | // mutex, so a stop callback that fires early (cancel_waiter) sees a | ||
| 60 | // null impl_ and is a safe no-op. To avoid losing such an early | ||
| 61 | // cancel, insert_waiter() re-checks stop_requested() under the lock | ||
| 62 | // and completes as canceled if it fires in this window. | ||
| 63 | 16383x | w.impl_ = nullptr; | |
| 64 | 16383x | w.svc_ = svc_; | |
| 65 | |||
| 66 | 16383x | might_have_pending_waits_.store(true, std::memory_order_relaxed); | |
| 67 | 16383x | svc_->get_scheduler().work_started(); | |
| 68 | |||
| 69 | 16383x | if (w.token_->stop_possible()) | |
| 70 | 2802x | w.arm_stop_cb(); | |
| 71 | |||
| 72 | // insert_waiter() grows the heap before publishing anything (see | ||
| 73 | // its own comment), so a throw out of it leaves w.impl_ null and | ||
| 74 | // the waiter never visible to the fire/cancel paths above -- but | ||
| 75 | // work_started() and arm_stop_cb() already ran and must be | ||
| 76 | // unwound here, or the escaping exception leaves a dangling | ||
| 77 | // stop-callback registration (UAF on a later cancel) and a work | ||
| 78 | // count that never balances. | ||
| 79 | try | ||
| 80 | { | ||
| 81 | 16383x | svc_->insert_waiter(*this, &w); | |
| 82 | } | ||
| 83 | ✗ | catch (...) | |
| 84 | { | ||
| 85 | ✗ | w.reset_stop_cb(); | |
| 86 | ✗ | svc_->get_scheduler().work_finished(); | |
| 87 | ✗ | throw; | |
| 88 | ✗ | } | |
| 89 | |||
| 90 | 16383x | return std::noop_coroutine(); | |
| 91 | } | ||
| 92 | |||
| 93 | std::coroutine_handle<> | ||
| 94 | 988x | timer::publish_wait(waiter_node& w) | |
| 95 | { | ||
| 96 | 988x | return get().publish(w); | |
| 97 | } | ||
| 98 | |||
| 99 | bool | ||
| 100 | 6397x | timer::rearm_wait(waiter_node& w, duration d) noexcept | |
| 101 | { | ||
| 102 | // The single waiter was popped before its op ran, so the impl is | ||
| 103 | // out of the heap with no published waiters: expires_after only | ||
| 104 | // stores the saturated expiry, and writing it is race-free. | ||
| 105 | 6397x | expires_after(d); | |
| 106 | 6397x | auto& impl = get(); | |
| 107 | // The drain that popped the waiter cleared the flag. | ||
| 108 | 6397x | impl.might_have_pending_waits_.store(true, std::memory_order_relaxed); | |
| 109 | try | ||
| 110 | { | ||
| 111 | 6397x | impl.svc_->insert_waiter(impl, &w); | |
| 112 | } | ||
| 113 | ✗ | catch (std::bad_alloc const&) | |
| 114 | { | ||
| 115 | // insert_waiter grows the heap before publishing anything, | ||
| 116 | // so the waiter is untouched and the caller can complete | ||
| 117 | // the wait through the normal resume path. | ||
| 118 | ✗ | return false; | |
| 119 | ✗ | } | |
| 120 | 6397x | return true; | |
| 121 | } | ||
| 122 | |||
| 123 | // completion_op and canceller definitions live here, non-inline, for | ||
| 124 | // the same reason wait() does: the inline waiter_node constructor in | ||
| 125 | // timer.hpp references do_complete and the vtable from translation | ||
| 126 | // units that never include timer_service.hpp. | ||
| 127 | |||
| 128 | void | ||
| 129 | 2689x | waiter_node::canceller::operator()() const | |
| 130 | { | ||
| 131 | 2689x | waiter_->svc_->cancel_waiter(waiter_); | |
| 132 | 2689x | } | |
| 133 | |||
| 134 | void | ||
| 135 | ✗ | waiter_node::completion_op::do_complete( | |
| 136 | [[maybe_unused]] void* owner, | ||
| 137 | scheduler_op* base, | ||
| 138 | std::uint32_t, | ||
| 139 | std::uint32_t) | ||
| 140 | { | ||
| 141 | // owner is always non-null here. The destroy path (owner == nullptr) | ||
| 142 | // is unreachable because completion_op overrides destroy() directly, | ||
| 143 | // bypassing scheduler_op::destroy() which would call func_(nullptr, ...). | ||
| 144 | ✗ | BOOST_COROSIO_ASSERT(owner); | |
| 145 | ✗ | static_cast<completion_op*>(base)->operator()(); | |
| 146 | ✗ | } | |
| 147 | |||
| 148 | void | ||
| 149 | 22749x | waiter_node::completion_op::operator()() | |
| 150 | { | ||
| 151 | // The node lives in the resuming coroutine's frame: posting the | ||
| 152 | // continuation is the last access, since the frame (and node) | ||
| 153 | // may complete and die on another thread immediately after. | ||
| 154 | 22749x | auto* w = waiter_; | |
| 155 | // A true return means the waiter re-published itself: the frame | ||
| 156 | // stays suspended, the wait's work count stays live, and the | ||
| 157 | // node may already be firing on another thread — no access past | ||
| 158 | // this point. | ||
| 159 | 22749x | if (w->on_fire_ && w->on_fire_(w->on_fire_ctx_)) | |
| 160 | 6397x | return; | |
| 161 | 16352x | w->reset_stop_cb(); | |
| 162 | 16352x | auto d = w->d_; | |
| 163 | 16352x | auto& sched = w->svc_->get_scheduler(); | |
| 164 | 16352x | d.post(w->cont_); | |
| 165 | 16352x | sched.work_finished(); | |
| 166 | } | ||
| 167 | |||
| 168 | void | ||
| 169 | 2x | waiter_node::completion_op::destroy() | |
| 170 | { | ||
| 171 | // Called during scheduler shutdown drain when this completion_op is | ||
| 172 | // in the scheduler's ready queue (posted by cancel_timer() or | ||
| 173 | // process_expired()). Balances the work_started() from | ||
| 174 | // implementation::wait(), keeping the run-loop counter sane; no | ||
| 175 | // shutdown path waits on that counter. | ||
| 176 | // | ||
| 177 | // This override also prevents scheduler_op::destroy() from calling | ||
| 178 | // do_complete(nullptr, ...). See also: timer_service::shutdown() | ||
| 179 | // which drains waiters still in the timer heap (the other path). | ||
| 180 | // Destroying the frame also ends the node's storage, so it is | ||
| 181 | // the last access. | ||
| 182 | 2x | auto* w = waiter_; | |
| 183 | 2x | w->reset_stop_cb(); | |
| 184 | 2x | auto h = std::exchange(w->h_, {}); | |
| 185 | 2x | auto& sched = w->svc_->get_scheduler(); | |
| 186 | 2x | sched.work_finished(); | |
| 187 | 2x | if (h) | |
| 188 | 2x | h.destroy(); | |
| 189 | 2x | } | |
| 190 | |||
| 191 | } // namespace boost::corosio::detail | ||
| 192 |