TLA Line data Source code
1 : //
2 : // Copyright (c) 2026 Steve Gerbino
3 : // Copyright (c) 2026 Michael Vandeberg
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_DELAY_HPP
12 : #define BOOST_COROSIO_DELAY_HPP
13 :
14 : #include <boost/corosio/detail/config.hpp>
15 : #include <boost/corosio/detail/except.hpp>
16 : #include <boost/corosio/detail/timer.hpp>
17 : #include <boost/corosio/wait_traits.hpp>
18 : #include <boost/capy/error.hpp>
19 : #include <boost/capy/ex/io_env.hpp>
20 : #include <boost/capy/io_result.hpp>
21 :
22 : #include <chrono>
23 : #include <concepts>
24 : #include <coroutine>
25 : #include <exception>
26 : #include <optional>
27 : #include <stdexcept>
28 : #include <system_error>
29 : #include <type_traits>
30 :
31 : namespace boost::corosio {
32 :
33 : namespace detail {
34 :
35 : // Narrow reps wrap if nanoseconds::max() is converted into them;
36 : // a double comparison clamps safely in both directions.
37 : template<typename Rep, typename Period>
38 : std::chrono::nanoseconds
39 HIT 20838 : clamp_to_ns(std::chrono::duration<Rep, Period> dur) noexcept
40 : {
41 : using namespace std::chrono;
42 : using dsec = duration<double>;
43 : if constexpr (std::is_floating_point_v<Rep>)
44 : {
45 : // NaN fails both clamp comparisons and would reach the
46 : // cast; treat it as no wait rather than undefined behavior.
47 2 : if (dur != dur)
48 2 : return nanoseconds::zero();
49 : }
50 20836 : return dsec(dur) >= dsec((nanoseconds::max)()) ? (nanoseconds::max)()
51 41670 : : dsec(dur) <= dsec((nanoseconds::min)())
52 20834 : ? (nanoseconds::min)()
53 20836 : : duration_cast<nanoseconds>(dur);
54 : }
55 :
56 : // A non-io_context executor cannot supply a timer service, and
57 : // await_suspend is driven through a noexcept wrapper, so translate
58 : // the service-lookup failure into a clear terminate.
59 : inline void
60 13035 : emplace_delay_timer(std::optional<timer>& t, capy::execution_context& ctx)
61 : {
62 : try
63 : {
64 13035 : t.emplace(ctx);
65 : }
66 2 : catch (std::logic_error const&)
67 : {
68 2 : throw_logic_error("delay requires an io_context-backed executor");
69 2 : }
70 MIS 0 : catch (std::exception const& e)
71 : {
72 0 : throw_logic_error(e.what());
73 0 : }
74 HIT 13033 : }
75 :
76 : } // namespace detail
77 :
78 : /** Suspends the calling coroutine until the deadline elapses or
79 : the environment's stop token is activated, whichever comes
80 : first. A deadline already elapsed at suspension, or a stop
81 : token already active, resumes the coroutine inline, without
82 : starting a timer (see Cancellation below). Otherwise the
83 : coroutine resumes through the executor once the timer fires
84 : or a mid-wait cancellation arrives.
85 :
86 : Not intended to be named directly. Use the @ref delay factory
87 : overloads instead.
88 :
89 : @pre The awaiting coroutine's executor must belong to an
90 : `io_context`. Any other execution context terminates with a
91 : diagnostic, because silently running without a timer would
92 : drop the requested delay.
93 :
94 : @par Cancellation
95 : If stop is already requested before suspension, the coroutine
96 : resumes immediately with `error::canceled`. If stop is
97 : requested while suspended, the pending wait is cancelled and
98 : the coroutine resumes with `error::canceled`. Requesting stop from
99 : another thread while the `io_context` runs in `single_threaded` mode is
100 : not permitted by `io_context`'s threading rules. That mode is
101 : auto-enabled at `concurrency_hint` == 1. Cross-thread cancellation
102 : requires a multi-threaded-capable context.
103 :
104 : @see delay
105 : */
106 : class delay_awaitable
107 : {
108 : // wait() names timer's private awaitable type; decltype is
109 : // the only way to store it here.
110 : using wait_type = decltype(std::declval<detail::timer&>().wait());
111 :
112 : std::chrono::steady_clock::time_point deadline_{};
113 : std::chrono::nanoseconds dur_{};
114 : bool has_deadline_ = false;
115 : bool canceled_ = false;
116 : std::optional<detail::timer> timer_;
117 : std::optional<wait_type> wait_;
118 :
119 : public:
120 : /// Construct an awaitable that waits for `dur` nanoseconds.
121 16456 : explicit delay_awaitable(std::chrono::nanoseconds dur) noexcept : dur_(dur)
122 : {
123 16456 : }
124 :
125 : /// Construct an awaitable that waits until `tp`.
126 16 : explicit delay_awaitable(std::chrono::steady_clock::time_point tp) noexcept
127 16 : : deadline_(tp)
128 16 : , has_deadline_(true)
129 : {
130 16 : }
131 :
132 : // Only moved before await_suspend; wait_ is engaged after.
133 : /// Construct by transferring state from `other`.
134 18500 : delay_awaitable(delay_awaitable&&) = default;
135 :
136 : /// Copy construction is disabled; an awaitable owns its timer.
137 : delay_awaitable(delay_awaitable const&) = delete;
138 : /// Copy assignment is disabled; an awaitable owns its timer.
139 : delay_awaitable& operator=(delay_awaitable const&) = delete;
140 : /// Move assignment is disabled; an awaitable is moved only before it is awaited.
141 : delay_awaitable& operator=(delay_awaitable&&) = delete;
142 :
143 : /// Return false unconditionally; see await_suspend.
144 : // The elapsed-deadline fast path must run after the stop-token
145 : // check, and only await_suspend receives the env carrying it.
146 16470 : bool await_ready() const noexcept
147 : {
148 16470 : return false;
149 : }
150 :
151 : /** Resume inline if stopped or elapsed; else wait on a timer.
152 :
153 : @param h Coroutine handle to resume on completion.
154 : @param env The I/O environment, carrying the executor, stop token
155 : and frame allocator.
156 :
157 : @return The handle to resume immediately, or `noop_coroutine()` when
158 : the wait was published to the timer service.
159 : */
160 : std::coroutine_handle<>
161 16472 : await_suspend(std::coroutine_handle<> h, capy::io_env const* env)
162 : {
163 16472 : if (env->stop_token.stop_requested())
164 : {
165 3495 : canceled_ = true;
166 3495 : return h;
167 : }
168 :
169 : // Elapsed deadlines complete synchronously, but only once a
170 : // pending stop request has already been ruled out above.
171 25940 : if (has_deadline_ ? deadline_ <= std::chrono::steady_clock::now()
172 12963 : : dur_.count() <= 0)
173 183 : return h;
174 :
175 12794 : detail::emplace_delay_timer(timer_, env->executor.context());
176 :
177 12792 : if (has_deadline_)
178 12 : timer_->expires_at(deadline_);
179 : else
180 12780 : timer_->expires_after(dur_);
181 :
182 12792 : wait_.emplace(timer_->wait());
183 12792 : return wait_->await_suspend(h, env);
184 : }
185 :
186 : /// Return empty on expiry, `error::canceled` if stop won.
187 16445 : [[nodiscard]] capy::io_result<> await_resume() noexcept
188 : {
189 16445 : if (canceled_)
190 3495 : return {capy::error::canceled};
191 12950 : if (wait_)
192 12767 : return wait_->await_resume();
193 183 : return {};
194 : }
195 : };
196 :
197 : /** Suspends the calling coroutine until `Clock::now()` reaches the
198 : deadline or the environment's stop token is activated. The wait is a
199 : sequence of steady-clock timer waits. After each expiry the clock is
200 : re-read. If the deadline is unreached, the same frame-embedded
201 : waiter is re-published for the next `Traits::to_wait_duration` cap.
202 : That re-publish neither resumes the coroutine nor allocates.
203 :
204 : Not intended to be named directly. Use the @ref delay factory
205 : overloads instead.
206 :
207 : @tparam Clock The clock the deadline is expressed in.
208 : @tparam Traits The wait-traits policy bounding each steady-clock wait.
209 :
210 : @pre The awaiting coroutine's executor must belong to an
211 : `io_context`. Any other execution context terminates with a
212 : diagnostic, because silently running without a timer would
213 : drop the requested delay.
214 :
215 : @par Cancellation
216 : Identical to @ref delay_awaitable: stop already requested
217 : resumes inline with `error::canceled`; stop while suspended
218 : cancels the pending wait, including between re-arms.
219 :
220 : @see delay, wait_traits
221 : */
222 : template<class Clock, class Traits>
223 : class clock_delay_awaitable
224 : {
225 : typename Clock::time_point deadline_{};
226 : bool canceled_ = false;
227 : std::optional<detail::timer> timer_;
228 : detail::waiter_node w_;
229 :
230 : std::chrono::nanoseconds
231 4384 : next_wait(typename Clock::time_point now) const noexcept
232 : {
233 4384 : return detail::clamp_to_ns(Traits::to_wait_duration(deadline_ - now));
234 : }
235 :
236 : // Runs on the scheduler thread executing the completion op,
237 : // before the continuation is posted, so the frame cannot die
238 : // concurrently.
239 4382 : static bool on_fire(void* ctx) noexcept
240 : {
241 4382 : auto* self = static_cast<clock_delay_awaitable*>(ctx);
242 : // Canceled: resume and surface the error
243 4382 : if (self->w_.ec_)
244 2 : return false;
245 4380 : auto now = Clock::now();
246 4380 : if (now >= self->deadline_)
247 237 : return false;
248 : // Re-publish and return without touching the node again:
249 : // the wait may complete on another thread immediately after.
250 4143 : if (self->timer_->rearm_wait(self->w_, self->next_wait(now)))
251 4143 : return true;
252 : // Heap growth failed; finish the wait with an error rather
253 : // than strand the frame with an unbalanced work count.
254 MIS 0 : self->w_.ec_ = std::make_error_code(std::errc::not_enough_memory);
255 0 : return false;
256 : }
257 :
258 : public:
259 : /// Construct an awaitable that waits until `tp` on `Clock`.
260 HIT 1247 : explicit clock_delay_awaitable(typename Clock::time_point tp) noexcept
261 1247 : : deadline_(tp)
262 : {
263 1247 : }
264 :
265 : // Only moved before await_suspend; w_ is quiescent until then.
266 : /// Construct by transferring the deadline from `other`.
267 1247 : clock_delay_awaitable(clock_delay_awaitable&& other) noexcept
268 1247 : : deadline_(other.deadline_)
269 : {
270 1247 : }
271 :
272 : /// Copy construction is disabled; an awaitable owns its timer.
273 : clock_delay_awaitable(clock_delay_awaitable const&) = delete;
274 : /// Copy assignment is disabled; an awaitable owns its timer.
275 : clock_delay_awaitable& operator=(clock_delay_awaitable const&) = delete;
276 : /// Move assignment is disabled; an awaitable is moved only before it is awaited.
277 : clock_delay_awaitable& operator=(clock_delay_awaitable&&) = delete;
278 :
279 : /// Return false unconditionally; see await_suspend.
280 : // The elapsed-deadline fast path must run after the stop-token
281 : // check, and only await_suspend receives the env carrying it.
282 1247 : bool await_ready() const noexcept
283 : {
284 1247 : return false;
285 : }
286 :
287 : /** Resume inline if stopped or reached; else wait on a timer.
288 :
289 : @param h Coroutine handle to resume on completion.
290 : @param env The I/O environment, carrying the executor, stop token
291 : and frame allocator.
292 :
293 : @return The handle to resume immediately, or `noop_coroutine()` when
294 : the wait was published to the timer service.
295 : */
296 : std::coroutine_handle<>
297 1247 : await_suspend(std::coroutine_handle<> h, capy::io_env const* env)
298 : {
299 1247 : if (env->stop_token.stop_requested())
300 : {
301 1004 : canceled_ = true;
302 1004 : return h;
303 : }
304 :
305 243 : auto now = Clock::now();
306 243 : if (now >= deadline_)
307 2 : return h;
308 :
309 241 : detail::emplace_delay_timer(timer_, env->executor.context());
310 :
311 241 : timer_->expires_after(next_wait(now));
312 :
313 241 : w_.bind(h, *env);
314 241 : w_.on_fire_ = &on_fire;
315 241 : w_.on_fire_ctx_ = this;
316 : // Never the elapsed fast path: a capped expiry that elapses
317 : // before publication must still reach on_fire, not complete
318 : // the clock wait early.
319 241 : return timer_->publish_wait(w_);
320 : }
321 :
322 : /// Return empty on deadline, `error::canceled` if stop won.
323 1245 : [[nodiscard]] capy::io_result<> await_resume() noexcept
324 : {
325 1245 : if (canceled_)
326 1004 : return {capy::error::canceled};
327 241 : if (timer_)
328 239 : return {w_.ec_};
329 2 : return {};
330 : }
331 : };
332 :
333 : /** Suspend the current coroutine for a duration.
334 :
335 : Returns an IoAwaitable that completes at or after the
336 : specified duration, or earlier if the environment's stop
337 : token is activated. Zero or negative durations complete
338 : synchronously.
339 :
340 : @par Example
341 : @par !example duration
342 :
343 : @param dur The duration to wait.
344 :
345 : @return A @ref delay_awaitable yielding `io_result<>`.
346 : */
347 : template<typename Rep, typename Period>
348 : [[nodiscard]] delay_awaitable
349 16454 : delay(std::chrono::duration<Rep, Period> dur) noexcept
350 : {
351 16454 : return delay_awaitable(detail::clamp_to_ns(dur));
352 : }
353 :
354 : /** Suspend the current coroutine until a time point.
355 :
356 : Returns an IoAwaitable that completes at or after `tp`, or
357 : earlier if the environment's stop token is activated. Time
358 : points already reached complete synchronously.
359 :
360 : @param tp The steady-clock time point to wait until.
361 :
362 : @return A @ref delay_awaitable yielding `io_result<>`.
363 : */
364 : [[nodiscard]] inline delay_awaitable
365 16 : delay(std::chrono::steady_clock::time_point tp) noexcept
366 : {
367 16 : return delay_awaitable(tp);
368 : }
369 :
370 : /** Suspend the current coroutine until a time point on `Clock`.
371 :
372 : Returns an IoAwaitable that completes at or after the first
373 : observation of `Clock::now() >= tp`, or earlier if the
374 : environment's stop token is activated. The wait is one or more
375 : bounded steady-clock waits, re-reading `Clock::now()` after
376 : each. `Traits::to_wait_duration` bounds each one. With the default
377 : @ref wait_traits a single full-length wait is used. An adjustment of
378 : `Clock` mid-wait is therefore observed only at natural wakeup.
379 : Supply capping traits to bound that latency. Time
380 : points already reached complete synchronously.
381 :
382 : @note `Clock::now()` and `Traits::to_wait_duration` are invoked
383 : on the `io_context`'s run thread and must not throw or block.
384 :
385 : @par Example
386 : @par !example system_clock_deadline
387 :
388 : @tparam Traits The wait-traits policy; `void` selects
389 : @ref wait_traits.
390 :
391 : @tparam Clock The clock type. This overload does not participate
392 : when `Clock` is `std::chrono::steady_clock`. The dedicated
393 : @ref delay overload taking a `steady_clock::time_point` handles
394 : that case.
395 :
396 : @param tp The time point to wait until.
397 :
398 : @return A @ref clock_delay_awaitable yielding `io_result<>`.
399 : */
400 : template<class Traits = void, class Clock, class Duration>
401 : requires(!std::same_as<Clock, std::chrono::steady_clock>) &&
402 : (std::is_void_v<Traits> || WaitTraits<Traits, Clock>)
403 : [[nodiscard]] auto
404 1247 : delay(std::chrono::time_point<Clock, Duration> tp) noexcept
405 : {
406 : using traits_type =
407 : std::conditional_t<std::is_void_v<Traits>, wait_traits<Clock>, Traits>;
408 : // ceil preserves completes-at-or-after when Duration is coarser
409 : // than the clock's native duration
410 : return clock_delay_awaitable<Clock, traits_type>(
411 1247 : std::chrono::ceil<typename Clock::duration>(tp));
412 : }
413 :
414 : } // namespace boost::corosio
415 :
416 : #endif
|