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