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 : 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 HIT 341 : 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 20604 : bool already_expired() const noexcept
136 : {
137 61812 : return heap_index_.load(std::memory_order_relaxed) == npos &&
138 20604 : (expiry_ ==
139 40544 : (std::chrono::steady_clock::time_point::min)() ||
140 40544 : 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 MIS 0 : 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 HIT 16 : void expires_at(time_point t)
244 : {
245 16 : auto& impl = get();
246 32 : BOOST_COROSIO_ASSERT(
247 : impl.heap_index_.load(std::memory_order_relaxed) ==
248 : implementation::npos);
249 16 : impl.expiry_ = t;
250 16 : }
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 10738 : void expires_after(duration d)
260 : {
261 10738 : auto& impl = get();
262 21476 : BOOST_COROSIO_ASSERT(
263 : impl.heap_index_.load(std::memory_order_relaxed) ==
264 : implementation::npos);
265 10738 : if (d <= duration::zero())
266 684 : 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 10054 : auto const now = clock_type::now();
273 10054 : impl.expiry_ = ((time_point::max)() - now < d)
274 20104 : ? (time_point::max)()
275 10050 : : now + d;
276 : }
277 10738 : }
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 21508 : implementation& get() const noexcept
356 : {
357 21508 : 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 23488 : 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 23488 : waiter_node() noexcept
456 23488 : {
457 23488 : op_.waiter_ = this;
458 23488 : }
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 10742 : void bind(std::coroutine_handle<> h, capy::io_env const& env) noexcept
476 : {
477 10742 : h_ = h;
478 10742 : cont_.h = h;
479 10742 : d_ = env.executor;
480 10742 : token_ = &env.stop_token;
481 10742 : }
482 :
483 : /** Arm the stop callback.
484 :
485 : @par Preconditions
486 : `token_` is set.
487 : */
488 1436 : void arm_stop_cb()
489 : {
490 1436 : new (cb_buf_) stop_cb_type(*token_, canceller{this});
491 1436 : cb_active_ = true;
492 1436 : }
493 :
494 : /// Destroy the stop callback if armed.
495 9889 : void reset_stop_cb() noexcept
496 : {
497 9889 : if (cb_active_)
498 : {
499 1436 : std::launder(reinterpret_cast<stop_cb_type*>(cb_buf_))
500 1436 : ->~stop_cb_type();
501 1436 : cb_active_ = false;
502 : }
503 9889 : }
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 10728 : explicit wait_awaitable(timer& t) noexcept : t_(t) {}
519 :
520 10728 : 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 2053 : bool await_ready() const noexcept
527 : {
528 2053 : 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 10700 : capy::io_result<> await_resume() const noexcept
535 : {
536 10700 : return {w_.ec_};
537 : }
538 :
539 10728 : auto await_suspend(std::coroutine_handle<> h, capy::io_env const* env)
540 : -> std::coroutine_handle<>
541 : {
542 10728 : auto& impl = t_.get();
543 10728 : 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 10728 : if (impl.already_expired())
549 : {
550 852 : w_.ec_ = {};
551 852 : w_.d_.post(w_.cont_);
552 852 : return std::noop_coroutine();
553 : }
554 :
555 9876 : return impl.wait(w_);
556 : }
557 : };
558 :
559 : inline wait_awaitable
560 10728 : timer::wait()
561 : {
562 10728 : return wait_awaitable(*this);
563 : }
564 :
565 : } // namespace boost::corosio::detail
566 :
567 : #endif
|