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 1806 : 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 28837 : bool already_expired() const noexcept
136 : {
137 86511 : return heap_index_.load(std::memory_order_relaxed) == npos &&
138 57010 : (expiry_ == (std::chrono::steady_clock::time_point::min)() ||
139 57010 : 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 MIS 0 : 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 HIT 16 : void expires_at(time_point t)
240 : {
241 16 : auto& impl = get();
242 32 : BOOST_COROSIO_ASSERT(
243 : impl.heap_index_.load(std::memory_order_relaxed) ==
244 : implementation::npos);
245 16 : impl.expiry_ = t;
246 16 : }
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 19213 : void expires_after(duration d)
255 : {
256 19213 : auto& impl = get();
257 38426 : BOOST_COROSIO_ASSERT(
258 : impl.heap_index_.load(std::memory_order_relaxed) ==
259 : implementation::npos);
260 19213 : if (d <= duration::zero())
261 680 : 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 18533 : auto const now = clock_type::now();
268 18533 : impl.expiry_ =
269 18533 : ((time_point::max)() - now < d) ? (time_point::max)() : now + d;
270 : }
271 19213 : }
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 38458 : implementation& get() const noexcept
348 : {
349 38458 : 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 32184 : 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 32184 : waiter_node() noexcept
448 32184 : {
449 32184 : op_.waiter_ = this;
450 32184 : }
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 15086 : void bind(std::coroutine_handle<> h, capy::io_env const& env) noexcept
468 : {
469 15086 : h_ = h;
470 15086 : cont_.h = h;
471 15086 : d_ = env.executor;
472 15086 : token_ = &env.stop_token;
473 15086 : }
474 :
475 : /** Arm the stop callback.
476 :
477 : @pre `token_` is set.
478 : */
479 1785 : void arm_stop_cb()
480 : {
481 1785 : new (cb_buf_) stop_cb_type(*token_, canceller{this});
482 1785 : cb_active_ = true;
483 1785 : }
484 :
485 : /// Destroy the stop callback if armed.
486 14233 : void reset_stop_cb() noexcept
487 : {
488 14233 : if (cb_active_)
489 : {
490 1785 : std::launder(reinterpret_cast<stop_cb_type*>(cb_buf_))
491 1785 : ->~stop_cb_type();
492 1785 : cb_active_ = false;
493 : }
494 14233 : }
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 14845 : explicit wait_awaitable(timer& t) noexcept : t_(t) {}
510 :
511 14845 : 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 2053 : bool await_ready() const noexcept
518 : {
519 2053 : 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 14816 : [[nodiscard]] capy::io_result<> await_resume() const noexcept
526 : {
527 14816 : return {w_.ec_};
528 : }
529 :
530 14845 : auto await_suspend(std::coroutine_handle<> h, capy::io_env const* env)
531 : -> std::coroutine_handle<>
532 : {
533 14845 : auto& impl = t_.get();
534 14845 : 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 14845 : if (impl.already_expired())
540 : {
541 853 : w_.ec_ = {};
542 853 : w_.d_.post(w_.cont_);
543 853 : return std::noop_coroutine();
544 : }
545 :
546 13992 : return impl.wait(w_);
547 : }
548 : };
549 :
550 : inline wait_awaitable
551 14845 : timer::wait()
552 : {
553 14845 : return wait_awaitable(*this);
554 : }
555 :
556 : } // namespace boost::corosio::detail
557 :
558 : #endif
|