TLA Line data 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 HIT 2874 : 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 14361 : void reuse() noexcept
135 : {
136 14361 : expiry_ = {};
137 14361 : heap_index_.store(npos, std::memory_order_relaxed);
138 14361 : might_have_pending_waits_.store(
139 : false, std::memory_order_relaxed);
140 14361 : BOOST_COROSIO_ASSERT(waiter_ == nullptr);
141 14361 : }
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 31642 : bool already_expired() const noexcept
159 : {
160 94926 : return heap_index_.load(std::memory_order_relaxed) == npos &&
161 62620 : (expiry_ == (std::chrono::steady_clock::time_point::min)() ||
162 62620 : 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 MIS 0 : 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 HIT 16 : void expires_at(time_point t)
263 : {
264 16 : auto& impl = get();
265 32 : BOOST_COROSIO_ASSERT(
266 : impl.heap_index_.load(std::memory_order_relaxed) ==
267 : implementation::npos);
268 16 : impl.expiry_ = t;
269 16 : }
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 23616 : void expires_after(duration d)
278 : {
279 23616 : auto& impl = get();
280 47232 : BOOST_COROSIO_ASSERT(
281 : impl.heap_index_.load(std::memory_order_relaxed) ==
282 : implementation::npos);
283 23616 : if (d <= duration::zero())
284 3160 : 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 20456 : auto const now = clock_type::now();
291 20456 : impl.expiry_ =
292 20456 : ((time_point::max)() - now < d) ? (time_point::max)() : now + d;
293 : }
294 23616 : }
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 47264 : implementation& get() const noexcept
371 : {
372 47264 : 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 35042 : 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 35042 : waiter_node() noexcept
471 35042 : {
472 35042 : op_.waiter_ = this;
473 35042 : }
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 17235 : void bind(std::coroutine_handle<> h, capy::io_env const& env) noexcept
491 : {
492 17235 : h_ = h;
493 17235 : cont_.h = h;
494 17235 : d_ = env.executor;
495 17235 : token_ = &env.stop_token;
496 17235 : }
497 :
498 : /** Arm the stop callback.
499 :
500 : @pre `token_` is set.
501 : */
502 2802 : void arm_stop_cb()
503 : {
504 2802 : new (cb_buf_) stop_cb_type(*token_, canceller{this});
505 2802 : cb_active_ = true;
506 2802 : }
507 :
508 : /// Destroy the stop callback if armed.
509 16383 : void reset_stop_cb() noexcept
510 : {
511 16383 : if (cb_active_)
512 : {
513 2802 : std::launder(reinterpret_cast<stop_cb_type*>(cb_buf_))
514 2802 : ->~stop_cb_type();
515 2802 : cb_active_ = false;
516 : }
517 16383 : }
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 16247 : explicit wait_awaitable(timer& t) noexcept : t_(t) {}
533 :
534 16247 : 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 2053 : bool await_ready() const noexcept
541 : {
542 2053 : 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 16218 : [[nodiscard]] capy::io_result<> await_resume() const noexcept
549 : {
550 16218 : return {w_.ec_};
551 : }
552 :
553 16247 : auto await_suspend(std::coroutine_handle<> h, capy::io_env const* env)
554 : -> std::coroutine_handle<>
555 : {
556 16247 : auto& impl = t_.get();
557 16247 : 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 16247 : if (impl.already_expired())
563 : {
564 852 : w_.ec_ = {};
565 852 : w_.d_.post(w_.cont_);
566 852 : return std::noop_coroutine();
567 : }
568 :
569 15395 : return impl.wait(w_);
570 : }
571 : };
572 :
573 : inline wait_awaitable
574 16247 : timer::wait()
575 : {
576 16247 : return wait_awaitable(*this);
577 : }
578 :
579 : } // namespace boost::corosio::detail
580 :
581 : #endif
|