100.00% Lines (62/62) 100.00% Functions (9/9)
TLA Baseline Branch
Line Hits Code Line Hits Code
1   // 1   //
2   // Copyright (c) 2026 Steve Gerbino 2   // Copyright (c) 2026 Steve Gerbino
3   // 3   //
4   // Distributed under the Boost Software License, Version 1.0. (See accompanying 4   // Distributed under the Boost Software License, Version 1.0. (See accompanying
5   // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt) 5   // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
6   // 6   //
7   // Official repository: https://github.com/cppalliance/corosio 7   // Official repository: https://github.com/cppalliance/corosio
8   // 8   //
9   9  
10   #ifndef BOOST_COROSIO_DETAIL_TIMEOUT_AWAITABLE_HPP 10   #ifndef BOOST_COROSIO_DETAIL_TIMEOUT_AWAITABLE_HPP
11   #define BOOST_COROSIO_DETAIL_TIMEOUT_AWAITABLE_HPP 11   #define BOOST_COROSIO_DETAIL_TIMEOUT_AWAITABLE_HPP
12   12  
13   #include <boost/corosio/io_context.hpp> 13   #include <boost/corosio/io_context.hpp>
14   #include <boost/corosio/detail/timeout_coro.hpp> 14   #include <boost/corosio/detail/timeout_coro.hpp>
15   #include <boost/corosio/detail/timer.hpp> 15   #include <boost/corosio/detail/timer.hpp>
16   #include <boost/corosio/detail/except.hpp> 16   #include <boost/corosio/detail/except.hpp>
17   #include <boost/capy/cond.hpp> 17   #include <boost/capy/cond.hpp>
18   #include <boost/capy/error.hpp> 18   #include <boost/capy/error.hpp>
19   #include <boost/capy/ex/io_env.hpp> 19   #include <boost/capy/ex/io_env.hpp>
20   #include <boost/capy/io_result.hpp> 20   #include <boost/capy/io_result.hpp>
21   21  
22   #include <chrono> 22   #include <chrono>
23   #include <coroutine> 23   #include <coroutine>
24   #include <new> 24   #include <new>
25   #include <optional> 25   #include <optional>
26   #include <stdexcept> 26   #include <stdexcept>
27   #include <stop_token> 27   #include <stop_token>
28   #include <type_traits> 28   #include <type_traits>
29   #include <utility> 29   #include <utility>
30   30  
31   /* Races an inner IoAwaitable against a timer via a shared 31   /* Races an inner IoAwaitable against a timer via a shared
32   stop_source. await_suspend arms the timer by launching a 32   stop_source. await_suspend arms the timer by launching a
33   fire-and-forget timeout_coro, then starts the inner op with 33   fire-and-forget timeout_coro, then starts the inner op with
34   an interposed stop_token. Whichever completes first signals 34   an interposed stop_token. Whichever completes first signals
35   the stop_source, cancelling the other. 35   the stop_source, cancelling the other.
36   36  
37   Parent cancellation is forwarded through a stop_callback 37   Parent cancellation is forwarded through a stop_callback
38   stored in a placement-new buffer (stop_callback is not 38   stored in a placement-new buffer (stop_callback is not
39   movable, but the awaitable must be movable for 39   movable, but the awaitable must be movable for
40   transform_awaiter). The buffer is inert during moves 40   transform_awaiter). The buffer is inert during moves
41   (before await_suspend) and constructed in-place once the 41   (before await_suspend) and constructed in-place once the
42   awaitable is pinned on the coroutine frame. 42   awaitable is pinned on the coroutine frame.
43   43  
44   The timeout_coro can outlive this awaitable — it owns its 44   The timeout_coro can outlive this awaitable — it owns its
45   env and self-destroys via suspend_never. The timer lives in 45   env and self-destroys via suspend_never. The timer lives in
46   std::optional and is constructed lazily in await_suspend, 46   std::optional and is constructed lazily in await_suspend,
47   once the awaiting coroutine's executor context is known. */ 47   once the awaiting coroutine's executor context is known. */
48   48  
49   namespace boost::corosio::detail { 49   namespace boost::corosio::detail {
50   50  
51   // Local stand-in for capy::detail's io_result trait: corosio must not 51   // Local stand-in for capy::detail's io_result trait: corosio must not
52   // reach into capy::detail, but the result-mapping switch in 52   // reach into capy::detail, but the result-mapping switch in
53   // await_resume needs to distinguish io_result from other return types. 53   // await_resume needs to distinguish io_result from other return types.
54   template<typename T> 54   template<typename T>
55   struct is_io_result : std::false_type 55   struct is_io_result : std::false_type
56   { 56   {
57   }; 57   };
58   58  
59   template<typename... Ts> 59   template<typename... Ts>
60   struct is_io_result<capy::io_result<Ts...>> : std::true_type 60   struct is_io_result<capy::io_result<Ts...>> : std::true_type
61   { 61   {
62   }; 62   };
63   63  
64   template<typename T> 64   template<typename T>
65   inline constexpr bool is_io_result_v = is_io_result<T>::value; 65   inline constexpr bool is_io_result_v = is_io_result<T>::value;
66   66  
67   /** Awaitable adapter that cancels an inner operation after a deadline. 67   /** Awaitable adapter that cancels an inner operation after a deadline.
68   68  
69   Races the inner awaitable against a timer. A shared stop_source 69   Races the inner awaitable against a timer. A shared stop_source
70   ties them together: whichever completes first cancels the other. 70   ties them together: whichever completes first cancels the other.
71   Parent cancellation is forwarded via stop_callback. 71   Parent cancellation is forwarded via stop_callback.
72   72  
73   The timer is constructed internally in `await_suspend` from the 73   The timer is constructed internally in `await_suspend` from the
74   execution context in `io_env`. 74   execution context in `io_env`.
75   75  
76   @tparam A The inner IoAwaitable type (decayed). 76   @tparam A The inner IoAwaitable type (decayed).
77   */ 77   */
78   template<typename A> 78   template<typename A>
79   struct timeout_awaitable 79   struct timeout_awaitable
80   { 80   {
81   struct stop_forwarder 81   struct stop_forwarder
82   { 82   {
83   std::stop_source* src_; 83   std::stop_source* src_;
HITCBC 84   2005 void operator()() const noexcept 84   1937 void operator()() const noexcept
85   { 85   {
HITCBC 86   2005 src_->request_stop(); 86   1937 src_->request_stop();
HITCBC 87   2005 } 87   1937 }
88   }; 88   };
89   89  
90   using time_point = std::chrono::steady_clock::time_point; 90   using time_point = std::chrono::steady_clock::time_point;
91   using stop_cb_type = std::stop_callback<stop_forwarder>; 91   using stop_cb_type = std::stop_callback<stop_forwarder>;
92   92  
93   A inner_; 93   A inner_;
94   std::optional<timer> timer_; 94   std::optional<timer> timer_;
95   time_point deadline_; 95   time_point deadline_;
96   std::chrono::nanoseconds dur_{}; 96   std::chrono::nanoseconds dur_{};
97   bool has_deadline_ = true; 97   bool has_deadline_ = true;
98   std::stop_source stop_src_; 98   std::stop_source stop_src_;
99   std::stop_token parent_token_; 99   std::stop_token parent_token_;
100   capy::io_env inner_env_; 100   capy::io_env inner_env_;
101   alignas(stop_cb_type) unsigned char cb_buf_[sizeof(stop_cb_type)]; 101   alignas(stop_cb_type) unsigned char cb_buf_[sizeof(stop_cb_type)];
102   bool cb_active_ = false; 102   bool cb_active_ = false;
103   103  
104   /// Construct without a timer, deadline given as an absolute time. 104   /// Construct without a timer, deadline given as an absolute time.
HITCBC 105   4 timeout_awaitable(A&& inner, time_point deadline) 105   4 timeout_awaitable(A&& inner, time_point deadline)
HITCBC 106   4 : inner_(std::move(inner)) 106   4 : inner_(std::move(inner))
HITCBC 107   4 , deadline_(deadline) 107   4 , deadline_(deadline)
108   { 108   {
HITCBC 109   4 } 109   4 }
110   110  
111   /// Construct without a timer, deadline measured from suspension. 111   /// Construct without a timer, deadline measured from suspension.
HITCBC 112   2041 timeout_awaitable(A&& inner, std::chrono::nanoseconds dur) 112   2052 timeout_awaitable(A&& inner, std::chrono::nanoseconds dur)
HITCBC 113   2041 : inner_(std::move(inner)) 113   2052 : inner_(std::move(inner))
HITCBC 114   2041 , dur_(dur) 114   2052 , dur_(dur)
HITCBC 115   2041 , has_deadline_(false) 115   2052 , has_deadline_(false)
116   { 116   {
HITCBC 117   2041 } 117   2052 }
118   118  
HITCBC 119   4090 ~timeout_awaitable() 119   4114 ~timeout_awaitable()
120   { 120   {
HITCBC 121   4090 destroy_parent_cb(); 121   4114 destroy_parent_cb();
HITCBC 122   4090 } 122   4114 }
123   123  
124   // Only moved before await_suspend, when cb_active_ is false 124   // Only moved before await_suspend, when cb_active_ is false
HITCBC 125   2045 timeout_awaitable(timeout_awaitable&& o) noexcept( 125   2058 timeout_awaitable(timeout_awaitable&& o) noexcept(
126   std::is_nothrow_move_constructible_v<A>) 126   std::is_nothrow_move_constructible_v<A>)
HITCBC 127   2045 : inner_(std::move(o.inner_)) 127   2058 : inner_(std::move(o.inner_))
HITCBC 128   2045 , timer_(std::move(o.timer_)) 128   2058 , timer_(std::move(o.timer_))
HITCBC 129   2045 , deadline_(o.deadline_) 129   2058 , deadline_(o.deadline_)
HITCBC 130   2045 , dur_(o.dur_) 130   2058 , dur_(o.dur_)
HITCBC 131   2045 , has_deadline_(o.has_deadline_) 131   2058 , has_deadline_(o.has_deadline_)
HITCBC 132   2045 , stop_src_(std::move(o.stop_src_)) 132   2058 , stop_src_(std::move(o.stop_src_))
133   { 133   {
HITCBC 134   2045 } 134   2058 }
135   135  
136   timeout_awaitable(timeout_awaitable const&) = delete; 136   timeout_awaitable(timeout_awaitable const&) = delete;
137   timeout_awaitable& operator=(timeout_awaitable const&) = delete; 137   timeout_awaitable& operator=(timeout_awaitable const&) = delete;
138   timeout_awaitable& operator=(timeout_awaitable&&) = delete; 138   timeout_awaitable& operator=(timeout_awaitable&&) = delete;
139   139  
ECB 140 - 2041 bool await_ready() const noexcept 140 + // Forwarding here is load-bearing, not an optimization: awaitables
  141 + // may perform setup in await_ready (type-erased stream wrappers
  142 + // construct their cached inner op there), so the full awaiter
  143 + // protocol must reach inner_ before await_suspend is driven. An
  144 + // already-ready inner op also skips arming the timer entirely.
HITGNC   145 + 2054 bool await_ready()
141   { 146   {
HITCBC 142 - 2041 return false; 147 + 2054 return inner_.await_ready();
143   } 148   }
144   149  
HITCBC 145   2045 auto await_suspend(std::coroutine_handle<> h, capy::io_env const* env) 150   2052 auto await_suspend(std::coroutine_handle<> h, capy::io_env const* env)
146   { 151   {
HITCBC 147   2045 parent_token_ = env->stop_token; 152   2052 parent_token_ = env->stop_token;
148   153  
149   // The deadline timer is built here from the awaiting 154   // The deadline timer is built here from the awaiting
150   // coroutine's executor context, the first point at which it 155   // coroutine's executor context, the first point at which it
151   // is known. await_suspend is driven through a noexcept 156   // is known. await_suspend is driven through a noexcept
152   // wrapper, so a failure cannot be surfaced as a catchable 157   // wrapper, so a failure cannot be surfaced as a catchable
153   // exception. An executor whose context is not an io_context 158   // exception. An executor whose context is not an io_context
154   // cannot supply a timer service; silently running the 159   // cannot supply a timer service; silently running the
155   // operation with no deadline would be a worse failure than 160   // operation with no deadline would be a worse failure than
156   // aborting, so translate the service-lookup error into a 161   // aborting, so translate the service-lookup error into a
157   // clear precondition diagnostic. This terminates by design 162   // clear precondition diagnostic. This terminates by design
158   // (a usage error) rather than dropping the requested timeout. 163   // (a usage error) rather than dropping the requested timeout.
159   // The detached timeout coroutine must own its executor by 164   // The detached timeout coroutine must own its executor by
160   // value (see timeout_coro::set_env_owned); io_env carries 165   // value (see timeout_coro::set_env_owned); io_env carries
161   // only a non-owning executor_ref. Recover the concrete 166   // only a non-owning executor_ref. Recover the concrete
162   // executor from the context rather than the executor_ref: 167   // executor from the context rather than the executor_ref:
163   // wrapped executors (a strand over the io_context) satisfy 168   // wrapped executors (a strand over the io_context) satisfy
164   // the documented precondition but do not expose the io 169   // the documented precondition but do not expose the io
165   // executor as their target. The timer construction below 170   // executor as their target. The timer construction below
166   // validates the context is an io_context, and the detached 171   // validates the context is an io_context, and the detached
167   // coroutine shares only the thread-safe stop_source with 172   // coroutine shares only the thread-safe stop_source with
168   // the caller, so resuming it on the raw io executor instead 173   // the caller, so resuming it on the raw io executor instead
169   // of the caller's wrapper is safe. 174   // of the caller's wrapper is safe.
170   try 175   try
171   { 176   {
HITCBC 172   2045 timer_.emplace(env->executor.context()); 177   2052 timer_.emplace(env->executor.context());
173   } 178   }
HITCBC 174   4 catch (std::logic_error const&) 179   4 catch (std::logic_error const&)
175   { 180   {
HITCBC 176   2 throw_logic_error( 181   2 throw_logic_error(
177   "timeout requires an io_context-backed executor"); 182   "timeout requires an io_context-backed executor");
178   } 183   }
179   auto ex = static_cast<io_context&>( 184   auto ex = static_cast<io_context&>(
HITCBC 180   2043 env->executor.context()).get_executor(); 185   2050 env->executor.context()).get_executor();
181   186  
HITCBC 182   2043 if (has_deadline_) 187   2050 if (has_deadline_)
HITCBC 183   4 timer_->expires_at(deadline_); 188   4 timer_->expires_at(deadline_);
184   else 189   else
HITCBC 185   2039 timer_->expires_after(dur_); 190   2046 timer_->expires_after(dur_);
186   191  
187   // Launch fire-and-forget timeout (starts suspended) 192   // Launch fire-and-forget timeout (starts suspended)
HITCBC 188   2043 auto timeout = make_timeout(*timer_, stop_src_); 193   2050 auto timeout = make_timeout(*timer_, stop_src_);
HITCBC 189   4086 timeout.h_.promise().set_env_owned( 194   4100 timeout.h_.promise().set_env_owned(
HITCBC 190   2043 ex, stop_src_.get_token(), env->frame_allocator); 195   2050 ex, stop_src_.get_token(), env->frame_allocator);
191   // Runs synchronously until timer.wait() suspends 196   // Runs synchronously until timer.wait() suspends
HITCBC 192   2043 timeout.h_.resume(); 197   2050 timeout.h_.resume();
193   // timeout goes out of scope; destructor is a no-op, 198   // timeout goes out of scope; destructor is a no-op,
194   // the coroutine self-destroys via suspend_never 199   // the coroutine self-destroys via suspend_never
195   200  
196   // Forward parent cancellation 201   // Forward parent cancellation
HITCBC 197   2043 new (cb_buf_) stop_cb_type(env->stop_token, stop_forwarder{&stop_src_}); 202   2050 new (cb_buf_) stop_cb_type(env->stop_token, stop_forwarder{&stop_src_});
HITCBC 198   2043 cb_active_ = true; 203   2050 cb_active_ = true;
199   204  
200   // Start the inner op with our interposed stop_token 205   // Start the inner op with our interposed stop_token
HITCBC 201   2043 inner_env_ = { 206   2050 inner_env_ = {
HITCBC 202   2043 env->executor, stop_src_.get_token(), env->frame_allocator}; 207   2050 env->executor, stop_src_.get_token(), env->frame_allocator};
HITCBC 203   4086 return inner_.await_suspend(h, &inner_env_); 208   4100 return inner_.await_suspend(h, &inner_env_);
HITCBC 204   2043 } 209   2050 }
205   210  
HITCBC 206   2039 decltype(auto) await_resume() 211   2050 decltype(auto) await_resume()
207   { 212   {
208   // Read before request_stop: afterwards stop_requested() 213   // Read before request_stop: afterwards stop_requested()
209   // can no longer distinguish who fired first. This must also 214   // can no longer distinguish who fired first. This must also
210   // happen before inner_.await_resume() rather than after: when 215   // happen before inner_.await_resume() rather than after: when
211   // the inner awaitable is itself a timeout_awaitable (nested 216   // the inner awaitable is itself a timeout_awaitable (nested
212   // timeout()), our own request_stop() below is visible through 217   // timeout()), our own request_stop() below is visible through
213   // its parent_token_ (aliasing our stop_src_), and would 218   // its parent_token_ (aliasing our stop_src_), and would
214   // otherwise make its read of "parent" look like a 219   // otherwise make its read of "parent" look like a
215   // cancellation that never happened. 220   // cancellation that never happened.
HITCBC 216   2039 bool const parent = parent_token_.stop_requested(); 221   2050 bool const parent = parent_token_.stop_requested();
HITCBC 217   2039 bool const fired = stop_src_.stop_requested(); 222   2050 bool const fired = stop_src_.stop_requested();
218   223  
219   // If inner_.await_resume() throws below, request_stop() is 224   // If inner_.await_resume() throws below, request_stop() is
220   // skipped; the still-armed timeout coroutine is then drained 225   // skipped; the still-armed timeout coroutine is then drained
221   // by timer_'s destructor rather than by us. 226   // by timer_'s destructor rather than by us.
HITCBC 222   2039 auto r = inner_.await_resume(); 227   2050 auto r = inner_.await_resume();
223   228  
224   // Cancel whichever is still pending (idempotent) 229   // Cancel whichever is still pending (idempotent)
HITCBC 225   2037 stop_src_.request_stop(); 230   2048 stop_src_.request_stop();
HITCBC 226   2037 destroy_parent_cb(); 231   2048 destroy_parent_cb();
227   232  
228   // Deadline won: stop_src_ is assumed to be the only 233   // Deadline won: stop_src_ is assumed to be the only
229   // cancellation source, whose only writers are the timer 234   // cancellation source, whose only writers are the timer
230   // coroutine and the parent forwarder, so fired && !parent 235   // coroutine and the parent forwarder, so fired && !parent
231   // identifies a timeout. A third-party cancellation of the 236   // identifies a timeout. A third-party cancellation of the
232   // inner op (e.g. a socket cancel issued from elsewhere) 237   // inner op (e.g. a socket cancel issued from elsewhere)
233   // landing in the same window as the deadline firing is 238   // landing in the same window as the deadline firing is
234   // reported as a timeout. 239   // reported as a timeout.
HITCBC 235   2059 if (fired && !parent && 240   2070 if (fired && !parent &&
HITCBC 236   2059 r.ec == capy::cond::canceled) 241   2070 r.ec == capy::cond::canceled)
237   { 242   {
HITCBC 238   22 std::remove_cvref_t<decltype(r)> t{}; 243   22 std::remove_cvref_t<decltype(r)> t{};
HITCBC 239   22 t.ec = make_error_code(capy::error::timeout); 244   22 t.ec = make_error_code(capy::error::timeout);
HITCBC 240   22 return t; 245   22 return t;
241   } 246   }
HITCBC 242   2015 return r; 247   2026 return r;
243   } 248   }
244   249  
HITCBC 245   6127 void destroy_parent_cb() noexcept 250   6162 void destroy_parent_cb() noexcept
246   { 251   {
HITCBC 247   6127 if (cb_active_) 252   6162 if (cb_active_)
248   { 253   {
HITCBC 249   2043 std::launder(reinterpret_cast<stop_cb_type*>(cb_buf_)) 254   2050 std::launder(reinterpret_cast<stop_cb_type*>(cb_buf_))
HITCBC 250   2043 ->~stop_cb_type(); 255   2050 ->~stop_cb_type();
HITCBC 251   2043 cb_active_ = false; 256   2050 cb_active_ = false;
252   } 257   }
HITCBC 253   6127 } 258   6162 }
254   }; 259   };
255   260  
256   } // namespace boost::corosio::detail 261   } // namespace boost::corosio::detail
257   262  
258   #endif 263   #endif