include/boost/corosio/detail/timer.hpp

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