src/corosio/src/timer.cpp

76.2% Lines (48 / 63) 90.9% Functions (10 / 11)
timer.cpp
f(x) Functions (11)
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