include/boost/corosio/detail/timer.hpp

98.3% Lines (58 / 59) 94.1% Functions (16 / 17)
timer.hpp
f(x) Functions (17)
Function Calls Lines Blocks
boost::corosio::detail::timer::implementation::implementation(boost::corosio::detail::timer_service&) :126 1806x 100.0% 100.0% boost::corosio::detail::timer::implementation::already_expired() const :135 28837x 100.0% 79.0% boost::corosio::detail::timer::timer(boost::corosio::detail::timer&&) :209 0 0.0% 0.0% boost::corosio::detail::timer::expires_at(std::chrono::time_point<std::chrono::_V2::steady_clock, std::chrono::duration<long, std::ratio<1l, 1000000000l> > >) :239 16x 100.0% 67.0% boost::corosio::detail::timer::expires_after(std::chrono::duration<long, std::ratio<1l, 1000000000l> >) :254 19213x 100.0% 72.0% boost::corosio::detail::timer::get() const :347 38458x 100.0% 100.0% boost::corosio::detail::waiter_node::completion_op::completion_op() :384 32184x 100.0% 100.0% boost::corosio::detail::waiter_node::waiter_node() :447 32184x 100.0% 100.0% boost::corosio::detail::waiter_node::bind(std::__n4861::coroutine_handle<void>, boost::capy::io_env const&) :467 15086x 100.0% 100.0% boost::corosio::detail::waiter_node::arm_stop_cb() :479 1785x 100.0% 100.0% boost::corosio::detail::waiter_node::reset_stop_cb() :486 14233x 100.0% 100.0% boost::corosio::detail::wait_awaitable::wait_awaitable(boost::corosio::detail::timer&) :509 14845x 100.0% 100.0% boost::corosio::detail::wait_awaitable::wait_awaitable(boost::corosio::detail::wait_awaitable&&) :511 14845x 100.0% 100.0% boost::corosio::detail::wait_awaitable::await_ready() const :517 2053x 100.0% 100.0% boost::corosio::detail::wait_awaitable::await_resume() const :525 14816x 100.0% 100.0% boost::corosio::detail::wait_awaitable::await_suspend(std::__n4861::coroutine_handle<void>, boost::capy::io_env const*) :530 14845x 100.0% 100.0% boost::corosio::detail::timer::wait() :551 14845x 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 #ifndef BOOST_COROSIO_DETAIL_TIMER_HPP
12 #define BOOST_COROSIO_DETAIL_TIMER_HPP
13
14 #include <boost/corosio/detail/config.hpp>
15 #include <boost/corosio/detail/intrusive.hpp>
16 #include <boost/corosio/detail/scheduler_op.hpp>
17 #include <boost/corosio/io/io_object.hpp>
18 #include <boost/capy/continuation.hpp>
19 #include <boost/capy/io_result.hpp>
20 #include <boost/capy/error.hpp>
21 #include <boost/capy/ex/executor_ref.hpp>
22 #include <boost/capy/ex/execution_context.hpp>
23 #include <boost/capy/ex/io_env.hpp>
24
25 #include <atomic>
26 #include <chrono>
27 #include <coroutine>
28 #include <cstddef>
29 #include <limits>
30 #include <new>
31 #include <stop_token>
32 #include <system_error>
33
34 namespace boost::corosio::detail {
35
36 // timer_service is defined in timer_service.hpp, which includes this
37 // header. waiter_node and wait_awaitable are defined below the timer
38 // class: waiter_node stores a timer::implementation*, which cannot be
39 // forward-declared as a nested type. implementation stores only a
40 // waiter_node pointer, so this forward declaration suffices for its
41 // data layout.
42 class timer_service;
43 struct waiter_node;
44 struct wait_awaitable;
45
46 /** An asynchronous timer for coroutine I/O.
47
48 This class provides asynchronous timer operations that return
49 awaitable types. The timer can be used to schedule operations
50 to occur after a specified duration or at a specific time point.
51
52 Each timer carries at most one wait: `delay` and `timeout` own a
53 private timer per `co_await`. When the timer expires the waiter
54 completes with success; a cancelled wait completes with an error
55 that compares equal to `capy::cond::canceled`.
56
57 Each timer operation participates in the affine awaitable protocol,
58 ensuring coroutines resume on the correct executor.
59
60 @par Thread Safety
61 Distinct objects: Safe.@n
62 Shared objects: Unsafe.
63
64 @par Semantics
65 Timers are not backed by per-timer kernel objects. The io_context's
66 timer service keeps a process-side min-heap of pending expirations;
67 the nearest expiry drives the reactor's poll timeout, and expirations
68 are processed in the run loop.
69 */
70 class BOOST_COROSIO_DECL timer : public io_object
71 {
72 friend struct wait_awaitable;
73
74 public:
75 /** Backend state and wait entry point for a timer.
76
77 Holds per-timer state ( expiry, heap position, the single waiter ) and
78 the `wait` entry point used by the awaitable returned from
79 @ref timer::wait. There is exactly one concrete timer backend,
80 so `wait` is a plain member function rather than a virtual
81 dispatch point.
82 */
83 struct implementation : io_object::implementation
84 {
85 /// Sentinel value indicating the timer is not in the heap.
86 static constexpr std::size_t npos =
87 (std::numeric_limits<std::size_t>::max)();
88
89 // Only mutated by the owning thread (expires_at/expires_after)
90 // before a wait is published; cross-thread consumers read the
91 // heap entry's copied time_, never this field, so it needs no
92 // atomicity.
93 /// The absolute expiry time point.
94 std::chrono::steady_clock::time_point expiry_{};
95
96 // heap_index_ and might_have_pending_waits_ are cross-thread
97 // hints, not authoritative state: the real state lives in the
98 // heap and the published waiter under timer_service::mutex_. Every
99 // unlocked fast-out that reads them is either re-validated under
100 // the mutex or safe under a stale value in both directions, and
101 // any locked writer / locked reader pair is already ordered by
102 // the mutex. All accesses therefore use memory_order_relaxed,
103 // which keeps the lock-free fast paths fence-free while making
104 // the concurrent reads well-defined.
105 /// Index in the timer service's min-heap, or `npos`.
106 std::atomic<std::size_t> heap_index_{npos};
107
108 // false implies waiter_ is null: both are cleared together
109 // under the service mutex.
110 /// True if `wait()` has been called since last cancel.
111 std::atomic<bool> might_have_pending_waits_{false};
112
113 /// The timer service that owns this implementation.
114 timer_service* svc_ = nullptr;
115
116 // Exactly one wait may be outstanding: delay and timeout own
117 // a private timer per co_await, and the service's drains rely
118 // on the one-to-one pairing.
119 /// The waiter published on this timer, or `nullptr`.
120 waiter_node* waiter_ = nullptr;
121
122 /// Free list linkage, reused when this impl is recycled.
123 implementation* next_free_ = nullptr;
124
125 /// Construct bound to the given timer service.
126 1806x explicit implementation(timer_service& svc) noexcept : svc_(&svc) {}
127
128 /** Check whether the timer is expired and absent from the heap.
129
130 The single definition of the already-expired fast-path
131 predicate: `await_suspend` tests it inline and `wait()`
132 re-tests it because the expiry can elapse between the two
133 reads.
134 */
135 28837x bool already_expired() const noexcept
136 {
137 86511x return heap_index_.load(std::memory_order_relaxed) == npos &&
138 57010x (expiry_ == (std::chrono::steady_clock::time_point::min)() ||
139 57010x expiry_ <= std::chrono::steady_clock::now());
140 }
141
142 /** Asynchronously wait for the timer to expire.
143
144 Publishes the waiter into the service's heap and the
145 timer's waiter slot, after which it may complete on any
146 thread. If the timer is already expired and not in the
147 heap, completes by posting the continuation without
148 publishing.
149
150 @pre @p w is fully initialized, and its storage (the awaitable
151 on the suspended coroutine's frame) outlives the wait.
152
153 @param w The waiter to publish.
154 */
155 // Exported at member level: dllexport on the enclosing timer
156 // class does not extend to nested classes, and header-inline
157 // callers (wait_awaitable::await_suspend) reference this
158 // symbol from outside the corosio DLL.
159 BOOST_COROSIO_DECL
160 std::coroutine_handle<> wait(waiter_node& w);
161
162 /** Publish a waiter unconditionally.
163
164 Like `wait`, but never takes the elapsed fast path. The
165 fast path posts the continuation directly, bypassing the
166 embedded op; hook-driven waits must observe every
167 completion through the op, where the re-arm hook runs.
168
169 @pre Same as `wait`.
170
171 @param w The waiter to publish.
172 */
173 std::coroutine_handle<> publish(waiter_node& w);
174 };
175
176 /// The clock type used for time operations.
177 using clock_type = std::chrono::steady_clock;
178
179 /// The time point type for absolute expiry times.
180 using time_point = clock_type::time_point;
181
182 /// The duration type for relative expiry times.
183 using duration = clock_type::duration;
184
185 /** Destructor.
186
187 Cancels any pending operations and releases timer resources.
188 */
189 ~timer() override;
190
191 /** Construct a timer from an execution context.
192
193 @param ctx The execution context that will own this timer. It
194 must be a corosio io_context; otherwise the constructor
195 throws (a timer service is required).
196
197 @throws std::logic_error if @p ctx is not an io_context.
198 */
199 explicit timer(capy::execution_context& ctx);
200
201 /** Move constructor.
202
203 Transfers ownership of the timer resources. Required so a
204 disengaged `std::optional<timer>` is movable; a timer is never
205 moved while a wait is published.
206
207 @pre No awaitables returned by @p other's methods exist.
208 */
209 ✗ timer(timer&&) noexcept = default;
210
211 /** Move assignment operator.
212
213 Closes any existing timer and transfers ownership.
214
215 @pre No awaitables returned by either `*this` or @p other's
216 methods exist.
217 */
218 timer& operator=(timer&&) noexcept = default;
219
220 timer(timer const&) = delete;
221 timer& operator=(timer const&) = delete;
222
223 /** Return the timer's expiry time as an absolute time.
224
225 @return The expiry time point. If no expiry has been set,
226 returns a default-constructed time_point.
227 */
228 time_point expiry() const noexcept
229 {
230 return get().expiry_;
231 }
232
233 /** Set the timer's expiry time as an absolute time.
234
235 @pre No wait is published on this timer.
236
237 @param t The expiry time to be used for the timer.
238 */
239 16x void expires_at(time_point t)
240 {
241 16x auto& impl = get();
242 32x BOOST_COROSIO_ASSERT(
243 impl.heap_index_.load(std::memory_order_relaxed) ==
244 implementation::npos);
245 16x impl.expiry_ = t;
246 16x }
247
248 /** Set the timer's expiry time relative to now.
249
250 @pre No wait is published on this timer.
251
252 @param d The expiry time relative to now.
253 */
254 19213x void expires_after(duration d)
255 {
256 19213x auto& impl = get();
257 38426x BOOST_COROSIO_ASSERT(
258 impl.heap_index_.load(std::memory_order_relaxed) ==
259 implementation::npos);
260 19213x if (d <= duration::zero())
261 680x impl.expiry_ = (time_point::min)();
262 else
263 {
264 // Saturate rather than overflow: a clamped near-max duration
265 // (e.g. delay(hours::max())) would wrap now() + d past the
266 // clock's range and appear already elapsed.
267 18533x auto const now = clock_type::now();
268 18533x impl.expiry_ =
269 18533x ((time_point::max)() - now < d) ? (time_point::max)() : now + d;
270 }
271 19213x }
272
273 /** Set the timer's expiry time relative to now.
274
275 This is a convenience overload that accepts any duration type
276 and converts it to the timer's native duration type.
277
278 @param d The expiry time relative to now.
279 */
280 template<class Rep, class Period>
281 void expires_after(std::chrono::duration<Rep, Period> d)
282 {
283 expires_after(std::chrono::duration_cast<duration>(d));
284 }
285
286 /** Wait for the timer to expire.
287
288 At most one wait may be outstanding at a time.
289
290 The operation supports cancellation via `std::stop_token` through
291 the affine awaitable protocol. If the associated stop token is
292 triggered, only that waiter completes with an error that
293 compares equal to `capy::cond::canceled`.
294
295 This timer must outlive the returned awaitable.
296
297 @return An awaitable that completes with `io_result<>`.
298 */
299 // Defined below wait_awaitable, which needs timer complete.
300 wait_awaitable wait();
301
302 /** Publish a hook-driven wait.
303
304 Bypasses the elapsed fast path so every completion is
305 delivered through the waiter's embedded op, where the
306 re-arm hook is consulted. Used by awaitables that
307 re-publish the waiter to continue a logical wait across
308 several timer expirations.
309
310 @pre @p w is fully initialized ( handle, executor, stop token,
311 hook fields ) and its storage outlives the wait.
312
313 @param w The waiter to publish.
314
315 @return `std::noop_coroutine()`.
316 */
317 std::coroutine_handle<> publish_wait(waiter_node& w);
318
319 /** Re-arm an already-fired waiter with a new relative expiry.
320
321 Stores the ( saturated ) expiry and re-publishes @p w. The
322 waiter's original work count and stop callback remain in
323 effect. Must only be called from the waiter's re-arm hook,
324 where the waiter has been popped from the service but not
325 yet resumed.
326
327 @pre The timer has no other waiters — this is what makes the
328 unlocked expiry write race-free.
329
330 Re-publication needs heap capacity and can fail under
331 allocation pressure. On failure the waiter is left exactly as
332 the hook received it, so the caller completes the wait through
333 the normal resume path instead of re-arming.
334
335 @param w The waiter to re-publish.
336 @param d The next expiry relative to now.
337
338 @return `true` if re-published; `false` if allocation failed.
339 */
340 [[nodiscard]] bool rearm_wait(waiter_node& w, duration d) noexcept;
341
342 protected:
343 explicit timer(handle h) noexcept : io_object(std::move(h)) {}
344
345 private:
346 /// Return the underlying implementation.
347 38458x implementation& get() const noexcept
348 {
349 38458x return *static_cast<implementation*>(h_.get());
350 }
351 };
352
353 /** Frame-resident per-wait state for a timer wait.
354
355 One node exists per `co_await` on a timer, embedded in the
356 awaitable on the suspended coroutine's frame — never allocated.
357 Once published by `implementation::wait()` the node may be
358 completed from any thread; every completion path finishes
359 touching the node before resuming or destroying the coroutine,
360 because either act may end the node's storage.
361
362 The node owns no resources: the stop token is borrowed from the
363 awaiting chain's `io_env` (which outlives the suspension) and
364 the stop callback is managed manually in `cb_buf_`, destroyed on
365 every completion path before the frame can die.
366 */
367 struct BOOST_COROSIO_SYMBOL_VISIBLE waiter_node
368 : intrusive_list<waiter_node>::node
369 {
370 // Embedded completion op — avoids heap allocation per fire/cancel.
371 // Members are exported and defined non-inline in timer.cpp: the
372 // inline waiter_node constructor references do_complete and the
373 // vtable from translation units that reach this header through
374 // delay.hpp without ever including timer_service.hpp, so the one
375 // strong definition must live in a TU that is always linked.
376 struct BOOST_COROSIO_SYMBOL_VISIBLE completion_op final : scheduler_op
377 {
378 waiter_node* waiter_ = nullptr;
379
380 BOOST_COROSIO_DECL
381 static void do_complete(
382 void* owner, scheduler_op* base, std::uint32_t, std::uint32_t);
383
384 32184x completion_op() noexcept : scheduler_op(&do_complete) {}
385
386 BOOST_COROSIO_DECL void operator()() override;
387 BOOST_COROSIO_DECL void destroy() override;
388 };
389
390 // Per-waiter stop_token cancellation
391 struct canceller
392 {
393 waiter_node* waiter_;
394 BOOST_COROSIO_DECL void operator()() const;
395 };
396
397 using stop_cb_type = std::stop_callback<canceller>;
398
399 // nullptr once unpublished from the timer ( concurrency marker )
400 /// The timer this waiter is published on, or `nullptr`.
401 timer::implementation* impl_ = nullptr;
402
403 /// The timer service that completes this waiter.
404 timer_service* svc_ = nullptr;
405
406 /// The suspended coroutine, destroyed by the shutdown drains.
407 std::coroutine_handle<> h_;
408
409 /// The continuation posted to resume the coroutine.
410 capy::continuation cont_;
411
412 /// The executor the continuation is posted through.
413 capy::executor_ref d_;
414
415 // Borrowed from the awaiting chain's io_env, which outlives the
416 // suspension; the node holds no owning state.
417 /// The stop token observed for cancellation.
418 std::stop_token const* token_ = nullptr;
419
420 /// The completion result read by `await_resume`.
421 std::error_code ec_;
422
423 // Consulted by the completion op before resuming; lets a
424 // clock-facade wait re-publish itself instead of completing.
425 // Never consulted on the shutdown destroy path. Consulted on
426 // every completion, including cancellation ( `ec_` set ) — the
427 // hook must inspect `w`'s `ec_` and must not re-arm a canceled
428 // waiter. Runs inside the completion path; must not throw.
429 /// Re-arm hook: return true to skip resumption ( wait continues ).
430 bool (*on_fire_)(void*) noexcept = nullptr;
431
432 /// Context passed to `on_fire_` ( the owning awaitable ).
433 void* on_fire_ctx_ = nullptr;
434
435 /// The embedded completion op posted to the scheduler.
436 completion_op op_;
437
438 // stop_callback is neither movable nor assignable; construct it
439 // in place once the node is pinned on the coroutine frame, and
440 // destroy it manually on every completion path.
441 /// Storage for the armed stop callback.
442 alignas(stop_cb_type) unsigned char cb_buf_[sizeof(stop_cb_type)];
443
444 /// True while `cb_buf_` holds a live stop callback.
445 bool cb_active_ = false;
446
447 32184x waiter_node() noexcept
448 32184x {
449 32184x op_.waiter_ = this;
450 32184x }
451
452 // The embedded op self-points and the list hooks are published
453 // to other threads; the node never moves.
454 waiter_node(waiter_node const&) = delete;
455 waiter_node& operator=(waiter_node const&) = delete;
456
457 /** Bind the coroutine and its environment before publication.
458
459 The single definition of the fields every wait must populate
460 before the node is published; hook-driven waits additionally
461 set `on_fire_` / `on_fire_ctx_`.
462
463 @param h The coroutine to resume on completion.
464 @param env The awaiting chain's environment; must outlive
465 the suspension.
466 */
467 15086x void bind(std::coroutine_handle<> h, capy::io_env const& env) noexcept
468 {
469 15086x h_ = h;
470 15086x cont_.h = h;
471 15086x d_ = env.executor;
472 15086x token_ = &env.stop_token;
473 15086x }
474
475 /** Arm the stop callback.
476
477 @pre `token_` is set.
478 */
479 1785x void arm_stop_cb()
480 {
481 1785x new (cb_buf_) stop_cb_type(*token_, canceller{this});
482 1785x cb_active_ = true;
483 1785x }
484
485 /// Destroy the stop callback if armed.
486 14233x void reset_stop_cb() noexcept
487 {
488 14233x if (cb_active_)
489 {
490 1785x std::launder(reinterpret_cast<stop_cb_type*>(cb_buf_))
491 1785x ->~stop_cb_type();
492 1785x cb_active_ = false;
493 }
494 14233x }
495 };
496
497 /** Awaitable returned by `timer::wait()`.
498
499 Carries the waiter node so a wait performs no allocation. The
500 awaitable is movable only before `await_suspend` publishes the
501 node (a move builds a fresh, quiescent node); afterwards it is
502 pinned on the coroutine frame until the wait completes.
503 */
504 struct wait_awaitable
505 {
506 timer& t_;
507 waiter_node w_;
508
509 14845x explicit wait_awaitable(timer& t) noexcept : t_(t) {}
510
511 14845x wait_awaitable(wait_awaitable&& o) noexcept : t_(o.t_) {}
512
513 wait_awaitable(wait_awaitable const&) = delete;
514 wait_awaitable& operator=(wait_awaitable const&) = delete;
515 wait_awaitable& operator=(wait_awaitable&&) = delete;
516
517 2053x bool await_ready() const noexcept
518 {
519 2053x return false;
520 }
521
522 // Cancellation surfaces through w_.ec_: the stop_token path in
523 // wait() completes the waiter with error::canceled written to
524 // it, so there is no separate token to consult here.
525 14816x [[nodiscard]] capy::io_result<> await_resume() const noexcept
526 {
527 14816x return {w_.ec_};
528 }
529
530 14845x auto await_suspend(std::coroutine_handle<> h, capy::io_env const* env)
531 -> std::coroutine_handle<>
532 {
533 14845x auto& impl = t_.get();
534 14845x w_.bind(h, *env);
535
536 // Inline fast path: already expired and not in the heap.
537 // Post instead of dispatch so the coroutine yields to the
538 // scheduler, allowing other queued work to run.
539 14845x if (impl.already_expired())
540 {
541 853x w_.ec_ = {};
542 853x w_.d_.post(w_.cont_);
543 853x return std::noop_coroutine();
544 }
545
546 13992x return impl.wait(w_);
547 }
548 };
549
550 inline wait_awaitable
551 14845x timer::wait()
552 {
553 14845x return wait_awaitable(*this);
554 }
555
556 } // namespace boost::corosio::detail
557
558 #endif
559