include/boost/corosio/detail/timer.hpp

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