100.00% Lines (83/83) 100.00% Functions (25/25)
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   // Copyright (c) 2026 Michael Vandeberg 4   // Copyright (c) 2026 Michael Vandeberg
5   // 5   //
6   // Distributed under the Boost Software License, Version 1.0. (See accompanying 6   // Distributed under the Boost Software License, Version 1.0. (See accompanying
7   // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt) 7   // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
8   // 8   //
9   // Official repository: https://github.com/cppalliance/corosio 9   // Official repository: https://github.com/cppalliance/corosio
10   // 10   //
11   11  
12   #ifndef BOOST_COROSIO_IO_CONTEXT_HPP 12   #ifndef BOOST_COROSIO_IO_CONTEXT_HPP
13   #define BOOST_COROSIO_IO_CONTEXT_HPP 13   #define BOOST_COROSIO_IO_CONTEXT_HPP
14   14  
15   #include <boost/corosio/detail/config.hpp> 15   #include <boost/corosio/detail/config.hpp>
16   #include <boost/corosio/detail/platform.hpp> 16   #include <boost/corosio/detail/platform.hpp>
17   #include <boost/corosio/detail/scheduler.hpp> 17   #include <boost/corosio/detail/scheduler.hpp>
18   #include <boost/capy/continuation.hpp> 18   #include <boost/capy/continuation.hpp>
19   #include <boost/capy/ex/execution_context.hpp> 19   #include <boost/capy/ex/execution_context.hpp>
20   20  
21   #include <chrono> 21   #include <chrono>
22   #include <coroutine> 22   #include <coroutine>
23   #include <cstddef> 23   #include <cstddef>
24   #include <limits> 24   #include <limits>
25   #include <thread> 25   #include <thread>
26   26  
27   namespace boost::corosio { 27   namespace boost::corosio {
28   28  
29   /** Locking-safety tier for an @ref io_context. 29   /** Locking-safety tier for an @ref io_context.
30   30  
31   Selects which internal locks the scheduler and reactor elide, trading 31   Selects which internal locks the scheduler and reactor elide, trading
32   thread-safety guarantees for reduced synchronization overhead. This is 32   thread-safety guarantees for reduced synchronization overhead. This is
33   the analog of Boost.Asio's `SAFE` / `UNSAFE_IO` / `UNSAFE` concurrency 33   the analog of Boost.Asio's `SAFE` / `UNSAFE_IO` / `UNSAFE` concurrency
34   hint constants. The tier is chosen explicitly, not derived from the 34   hint constants. The tier is chosen explicitly, not derived from the
35   `concurrency_hint`. (The reverse does apply: a lockless tier reduces the 35   `concurrency_hint`. (The reverse does apply: a lockless tier reduces the
36   effective hint used for performance tuning to 1.) 36   effective hint used for performance tuning to 1.)
37   37  
38   @see io_context_options::locking 38   @see io_context_options::locking
39   */ 39   */
40   enum class locking_mode 40   enum class locking_mode
41   { 41   {
42   /** Full thread safety (default). All locks enabled; equivalent to 42   /** Full thread safety (default). All locks enabled; equivalent to
43   Boost.Asio's `SAFE`/`DEFAULT`. Any thread may use the context. */ 43   Boost.Asio's `SAFE`/`DEFAULT`. Any thread may use the context. */
44   safe, 44   safe,
45   45  
46   /** Disable only the per-descriptor I/O locks; keep scheduler locking. 46   /** Disable only the per-descriptor I/O locks; keep scheduler locking.
47   Equivalent to Boost.Asio's `UNSAFE_IO`. The context must be run 47   Equivalent to Boost.Asio's `UNSAFE_IO`. The context must be run
48   and driven by a single thread, but resolver and POSIX file 48   and driven by a single thread, but resolver and POSIX file
49   services remain available (they rely on scheduler locking, which 49   services remain available (they rely on scheduler locking, which
50   stays on). */ 50   stays on). */
51   unsafe_io, 51   unsafe_io,
52   52  
53   /** Disable all locking (fully lockless). Equivalent to Boost.Asio's 53   /** Disable all locking (fully lockless). Equivalent to Boost.Asio's
54   `UNSAFE`. 54   `UNSAFE`.
55   55  
56   @par Restrictions 56   @par Restrictions
57   - Only one thread may call `run()` (or any run variant). 57   - Only one thread may call `run()` (or any run variant).
58   - Posting work from another thread is undefined behavior. 58   - Posting work from another thread is undefined behavior.
59   - DNS resolution returns `operation_not_supported`. 59   - DNS resolution returns `operation_not_supported`.
60   - POSIX file I/O returns `operation_not_supported`. 60   - POSIX file I/O returns `operation_not_supported`.
61   - Signal sets should not be shared across contexts. */ 61   - Signal sets should not be shared across contexts. */
62   unsafe 62   unsafe
63   }; 63   };
64   64  
65   /** Runtime tuning options for @ref io_context. 65   /** Runtime tuning options for @ref io_context.
66   66  
67   All fields have defaults that match the library's built-in 67   All fields have defaults that match the library's built-in
68   values, so constructing a default `io_context_options` produces 68   values, so constructing a default `io_context_options` produces
69   identical behavior to an unconfigured context. 69   identical behavior to an unconfigured context.
70   70  
71   Options that apply only to a specific backend family are 71   Options that apply only to a specific backend family are
72   silently ignored when the active backend does not support them. 72   silently ignored when the active backend does not support them.
73   73  
74   @par Example 74   @par Example
75   @code 75   @code
76   io_context_options opts; 76   io_context_options opts;
77   opts.max_events_per_poll = 256; // larger batch per syscall 77   opts.max_events_per_poll = 256; // larger batch per syscall
78   opts.inline_budget_max = 32; // more speculative completions 78   opts.inline_budget_max = 32; // more speculative completions
79   opts.thread_pool_size = 4; // more file-I/O workers 79   opts.thread_pool_size = 4; // more file-I/O workers
80   80  
81   io_context ioc(opts); 81   io_context ioc(opts);
82   @endcode 82   @endcode
83   83  
84   @see io_context, native_io_context 84   @see io_context, native_io_context
85   */ 85   */
86   struct io_context_options 86   struct io_context_options
87   { 87   {
88   /** Maximum events fetched per reactor poll call. 88   /** Maximum events fetched per reactor poll call.
89   89  
90   Controls the buffer size passed to `epoll_wait()` or 90   Controls the buffer size passed to `epoll_wait()` or
91   `kevent()`. Larger values reduce syscall frequency under 91   `kevent()`. Larger values reduce syscall frequency under
92   high load; smaller values improve fairness between 92   high load; smaller values improve fairness between
93   connections. Ignored on IOCP and select backends. 93   connections. Ignored on IOCP and select backends.
94   */ 94   */
95   unsigned max_events_per_poll = 128; 95   unsigned max_events_per_poll = 128;
96   96  
97   /** Starting inline completion budget per handler chain. 97   /** Starting inline completion budget per handler chain.
98   98  
99   After a posted handler executes, the reactor grants this 99   After a posted handler executes, the reactor grants this
100   many speculative inline completions before forcing a 100   many speculative inline completions before forcing a
101   re-queue. Applies to reactor backends only. 101   re-queue. Applies to reactor backends only.
102   102  
103   @note Constructing an `io_context` with `concurrency_hint > 1` 103   @note Constructing an `io_context` with `concurrency_hint > 1`
104   and all three budget fields at their defaults overrides 104   and all three budget fields at their defaults overrides
105   them to disable inline completion (post-everything mode), 105   them to disable inline completion (post-everything mode),
106   since multi-thread workloads benefit from cross-thread 106   since multi-thread workloads benefit from cross-thread
107   work-stealing. Setting any budget field to a non-default 107   work-stealing. Setting any budget field to a non-default
108   value disables the override. 108   value disables the override.
109   */ 109   */
110   unsigned inline_budget_initial = 2; 110   unsigned inline_budget_initial = 2;
111   111  
112   /** Hard ceiling on adaptive inline budget ramp-up. 112   /** Hard ceiling on adaptive inline budget ramp-up.
113   113  
114   The budget doubles each cycle it is fully consumed, up to 114   The budget doubles each cycle it is fully consumed, up to
115   this limit. Applies to reactor backends only. 115   this limit. Applies to reactor backends only.
116   */ 116   */
117   unsigned inline_budget_max = 16; 117   unsigned inline_budget_max = 16;
118   118  
119   /** Inline budget when no other thread assists the reactor. 119   /** Inline budget when no other thread assists the reactor.
120   120  
121   When only one thread is running the event loop, this 121   When only one thread is running the event loop, this
122   value caps the inline budget to preserve fairness. 122   value caps the inline budget to preserve fairness.
123   Applies to reactor backends only. 123   Applies to reactor backends only.
124   */ 124   */
125   unsigned unassisted_budget = 4; 125   unsigned unassisted_budget = 4;
126   126  
127   /** Thread pool size for blocking I/O (file I/O, DNS resolution). 127   /** Thread pool size for blocking I/O (file I/O, DNS resolution).
128   128  
129   Sets the number of worker threads in the shared thread pool 129   Sets the number of worker threads in the shared thread pool
130   used by POSIX file services and DNS resolution. Must be at 130   used by POSIX file services and DNS resolution. Must be at
131   least 1. Applies to POSIX backends only; ignored on IOCP 131   least 1. Applies to POSIX backends only; ignored on IOCP
132   where file I/O uses native overlapped I/O. 132   where file I/O uses native overlapped I/O.
133   */ 133   */
134   unsigned thread_pool_size = 1; 134   unsigned thread_pool_size = 1;
135   135  
136   /** Thread-safety tier. See @ref locking_mode for the tiers and their 136   /** Thread-safety tier. See @ref locking_mode for the tiers and their
137   restrictions. 137   restrictions.
138   */ 138   */
139   locking_mode locking = locking_mode::safe; 139   locking_mode locking = locking_mode::safe;
140   140  
141   /** Enable IORING_SETUP_SQPOLL on the io_uring backend. 141   /** Enable IORING_SETUP_SQPOLL on the io_uring backend.
142   142  
143   With SQPOLL, the kernel forks a thread that busy-polls the 143   With SQPOLL, the kernel forks a thread that busy-polls the
144   submission ring; submission becomes a userspace-only memory 144   submission ring; submission becomes a userspace-only memory
145   store, eliminating the io_uring_enter syscall on the submit 145   store, eliminating the io_uring_enter syscall on the submit
146   path. Most useful for sustained traffic. Idle thread parks 146   path. Most useful for sustained traffic. Idle thread parks
147   after `sq_thread_idle_ms` of no activity. 147   after `sq_thread_idle_ms` of no activity.
148   148  
149   Independent of `locking`. Default: off. 149   Independent of `locking`. Default: off.
150   150  
151   Ignored on non-io_uring backends. 151   Ignored on non-io_uring backends.
152   */ 152   */
153   bool enable_sqpoll = false; 153   bool enable_sqpoll = false;
154   154  
155   /** SQ-poll idle timeout in milliseconds. 155   /** SQ-poll idle timeout in milliseconds.
156   156  
157   After this many ms of no submissions, the kernel polling 157   After this many ms of no submissions, the kernel polling
158   thread sleeps; next submit re-wakes it via SQ_WAKEUP. 0 158   thread sleeps; next submit re-wakes it via SQ_WAKEUP. 0
159   means use the kernel default (1ms). Recommended for bursty 159   means use the kernel default (1ms). Recommended for bursty
160   workloads: 100-1000ms (avoids park/unpark thrash). 160   workloads: 100-1000ms (avoids park/unpark thrash).
161   161  
162   Ignored unless `enable_sqpoll` is true. Ignored on 162   Ignored unless `enable_sqpoll` is true. Ignored on
163   non-io_uring backends. 163   non-io_uring backends.
164   */ 164   */
165   unsigned sq_thread_idle_ms = 0; 165   unsigned sq_thread_idle_ms = 0;
166   166  
167   /** Pin the SQ-poll kernel thread to this CPU. 167   /** Pin the SQ-poll kernel thread to this CPU.
168   168  
169   -1 means do not pin (kernel scheduler picks). Pinning off 169   -1 means do not pin (kernel scheduler picks). Pinning off
170   the dispatch core is recommended on latency-sensitive 170   the dispatch core is recommended on latency-sensitive
171   deployments to avoid cache contention. 171   deployments to avoid cache contention.
172   172  
173   Ignored unless `enable_sqpoll` is true. Ignored on 173   Ignored unless `enable_sqpoll` is true. Ignored on
174   non-io_uring backends. 174   non-io_uring backends.
175   */ 175   */
176   int sq_thread_cpu = -1; 176   int sq_thread_cpu = -1;
177   }; 177   };
178   178  
179   namespace detail { 179   namespace detail {
180   class timer_service; 180   class timer_service;
181   181  
182   /** Return the hint used for performance tuning: the lockless tiers are 182   /** Return the hint used for performance tuning: the lockless tiers are
183   single-threaded, so their effective hint is 1 whatever the caller passed. 183   single-threaded, so their effective hint is 1 whatever the caller passed.
184   */ 184   */
185   inline unsigned 185   inline unsigned
HITCBC 186   36 effective_concurrency_hint( 186   36 effective_concurrency_hint(
187   io_context_options const& opts, unsigned hint) noexcept 187   io_context_options const& opts, unsigned hint) noexcept
188   { 188   {
HITCBC 189   36 return opts.locking == locking_mode::safe ? hint : 1u; 189   36 return opts.locking == locking_mode::safe ? hint : 1u;
190   } 190   }
191   } // namespace detail 191   } // namespace detail
192   192  
193   /** An I/O context for running asynchronous operations. 193   /** An I/O context for running asynchronous operations.
194   194  
195   The io_context provides an execution environment for async 195   The io_context provides an execution environment for async
196   operations. It maintains a queue of pending work items and 196   operations. It maintains a queue of pending work items and
197   processes them when `run()` is called. 197   processes them when `run()` is called.
198   198  
199   The default and unsigned constructors select the platform's 199   The default and unsigned constructors select the platform's
200   native backend: 200   native backend:
201   - Windows: IOCP 201   - Windows: IOCP
202   - Linux: epoll 202   - Linux: epoll
203   - BSD/macOS: kqueue 203   - BSD/macOS: kqueue
204   - Other POSIX: select 204   - Other POSIX: select
205   205  
206   The template constructor accepts a backend tag value to 206   The template constructor accepts a backend tag value to
207   choose a specific backend at compile time: 207   choose a specific backend at compile time:
208   208  
209   @par Example 209   @par Example
210   @code 210   @code
211   io_context ioc; // platform default 211   io_context ioc; // platform default
212   io_context ioc2(corosio::epoll); // explicit backend 212   io_context ioc2(corosio::epoll); // explicit backend
213   @endcode 213   @endcode
214   214  
215   @par Preconditions 215   @par Preconditions
216   The context must outlive every operation posted or dispatched 216   The context must outlive every operation posted or dispatched
217   through its executor, and no thread may be executing a run 217   through its executor, and no thread may be executing a run
218   variant when the context is destroyed. Posting to the context 218   variant when the context is destroyed. Posting to the context
219   concurrently with, or after, its destruction is undefined 219   concurrently with, or after, its destruction is undefined
220   behavior. The safe teardown pattern is to stop submitting new 220   behavior. The safe teardown pattern is to stop submitting new
221   work, let every `run()` call return (each returns once no 221   work, let every `run()` call return (each returns once no
222   outstanding work remains), and join the threads that ran the 222   outstanding work remains), and join the threads that ran the
223   loop before destroying the context. Work launched with 223   loop before destroying the context. Work launched with
224   `capy::run` / `capy::run_async` is work-tracked, so a normal 224   `capy::run` / `capy::run_async` is work-tracked, so a normal
225   `run()` completion already waits for it. 225   `run()` completion already waits for it.
226   226  
  227 + @par Exception Safety
  228 + A context that constructs is usable. The infrastructure its
  229 + backend needs — the completion port, the ring, the reactor's
  230 + wakeup channel — is created during construction, so a system that
  231 + refuses it throws from the constructor rather than from the first
  232 + operation, and the failed construction leaves nothing open.
  233 +
227   @par Thread Safety 234   @par Thread Safety
228   Distinct objects: Safe.@n 235   Distinct objects: Safe.@n
229   Shared objects: Safe, unless the context was constructed with a 236   Shared objects: Safe, unless the context was constructed with a
230   lockless @ref io_context_options::locking tier (`unsafe_io` or 237   lockless @ref io_context_options::locking tier (`unsafe_io` or
231   `unsafe`), in which case a single thread must drive it. 238   `unsafe`), in which case a single thread must drive it.
232   239  
233   @see epoll_t, select_t, kqueue_t, iocp_t 240   @see epoll_t, select_t, kqueue_t, iocp_t
234   */ 241   */
235   class BOOST_COROSIO_DECL io_context : public capy::execution_context 242   class BOOST_COROSIO_DECL io_context : public capy::execution_context
236   { 243   {
237 - /// Pre-create services that depend on options (before construct). 244 + /// Reject invalid options before the backend is constructed.
238   void apply_options_pre_(io_context_options const& opts); 245   void apply_options_pre_(io_context_options const& opts);
239   246  
240 - /// Apply runtime tuning to the scheduler (after construct). 247 + /** Create the blocking-I/O thread pool, apply runtime tuning to the
  248 + scheduler and finish bringing the backend up. The tail of every
  249 + options constructor: the backend infrastructure whose setup reads
  250 + these options is created here, so a failure to create it throws
  251 + from the constructor. */
241   void apply_options_post_( 252   void apply_options_post_(
242   io_context_options const& opts, 253   io_context_options const& opts,
243   unsigned concurrency_hint); 254   unsigned concurrency_hint);
244   255  
245 - /** Apply only the decomposed threading configuration (locking tiers). 256 + /** Create the blocking-I/O thread pool and apply only the decomposed
246 - Used by the plain constructors, which — unlike the options 257 + threading configuration (locking tiers), then finish bringing the
247 - constructors — deliberately leave the reactor budget at its defaults 258 + backend up. The tail of every plain constructor, which — unlike
248 - rather than engaging the multi-thread post-everything heuristic. */ 259 + the options constructors — deliberately leaves the reactor budget
  260 + at its defaults rather than engaging the multi-thread
  261 + post-everything heuristic. */
249   void apply_threading_(io_context_options const& opts); 262   void apply_threading_(io_context_options const& opts);
250   263  
251   protected: 264   protected:
252   detail::scheduler* sched_; 265   detail::scheduler* sched_;
253   266  
254   public: 267   public:
255   /** The executor type for this context. */ 268   /** The executor type for this context. */
256   class executor_type; 269   class executor_type;
257   270  
258   /** Construct with default concurrency and platform backend. 271   /** Construct with default concurrency and platform backend.
259   272  
260   Uses `std::thread::hardware_concurrency()` (floored to 1, in 273   Uses `std::thread::hardware_concurrency()` (floored to 1, in
261   case it reports 0) as the concurrency hint, and the default 274   case it reports 0) as the concurrency hint, and the default
262   @ref locking_mode::safe tier. Select a lockless tier via 275   @ref locking_mode::safe tier. Select a lockless tier via
263   @ref io_context_options::locking. 276   @ref io_context_options::locking.
  277 +
  278 + @throws std::system_error If the backend's infrastructure
  279 + could not be created.
264   */ 280   */
265   io_context(); 281   io_context();
266   282  
267   /** Construct with a concurrency hint and platform backend. 283   /** Construct with a concurrency hint and platform backend.
268   284  
269   @param concurrency_hint Hint for the number of threads 285   @param concurrency_hint Hint for the number of threads
270   that will call `run()`. 286   that will call `run()`.
  287 +
  288 + @throws std::system_error If the backend's infrastructure
  289 + could not be created.
271   */ 290   */
272   explicit io_context(unsigned concurrency_hint); 291   explicit io_context(unsigned concurrency_hint);
273   292  
274   /** Construct with runtime tuning options and platform backend. 293   /** Construct with runtime tuning options and platform backend.
275   294  
276   @param opts Runtime options controlling scheduler and 295   @param opts Runtime options controlling scheduler and
277   service behavior. 296   service behavior.
278   @param concurrency_hint Hint for the number of threads 297   @param concurrency_hint Hint for the number of threads
279   that will call `run()`. 298   that will call `run()`.
280   299  
281   @throws std::invalid_argument If `opts.thread_pool_size` is 300   @throws std::invalid_argument If `opts.thread_pool_size` is
282   less than 1 (POSIX). 301   less than 1 (POSIX).
  302 +
  303 + @throws std::system_error If the backend's infrastructure
  304 + could not be created.
283   */ 305   */
284   explicit io_context( 306   explicit io_context(
285   io_context_options const& opts, 307   io_context_options const& opts,
286   unsigned concurrency_hint = std::thread::hardware_concurrency()); 308   unsigned concurrency_hint = std::thread::hardware_concurrency());
287   309  
288   /** Construct with an explicit backend tag. 310   /** Construct with an explicit backend tag.
289   311  
290   @param backend The backend tag value selecting the I/O 312   @param backend The backend tag value selecting the I/O
291   multiplexer (e.g. `corosio::epoll`). 313   multiplexer (e.g. `corosio::epoll`).
292   @param concurrency_hint Hint for the number of threads 314   @param concurrency_hint Hint for the number of threads
293   that will call `run()`. 315   that will call `run()`.
  316 +
  317 + @throws std::system_error If the backend's infrastructure
  318 + could not be created.
294   */ 319   */
295   template<class Backend> 320   template<class Backend>
296   requires requires { Backend::construct; } 321   requires requires { Backend::construct; }
HITCBC 297   1386 explicit io_context( 322   1581 explicit io_context(
298   [[maybe_unused]] Backend backend, 323   [[maybe_unused]] Backend backend,
299   unsigned concurrency_hint = std::thread::hardware_concurrency()) 324   unsigned concurrency_hint = std::thread::hardware_concurrency())
300   : capy::execution_context(this) 325   : capy::execution_context(this)
HITCBC 301   1386 , sched_(nullptr) 326   1581 , sched_(nullptr)
302   { 327   {
HITCBC 303   1386 sched_ = &Backend::construct(*this, concurrency_hint); 328   1581 sched_ = &Backend::construct(*this, concurrency_hint);
304   // Apply threading config only (locking tier). Unlike the options 329   // Apply threading config only (locking tier). Unlike the options
305   // ctor, the plain path leaves the reactor budget at its defaults. 330   // ctor, the plain path leaves the reactor budget at its defaults.
HITCBC 306   1386 apply_threading_(io_context_options{}); 331   1569 apply_threading_(io_context_options{});
HITCBC 307   1386 } 332   1581 }
308   333  
309   /** Construct with an explicit backend tag and runtime options. 334   /** Construct with an explicit backend tag and runtime options.
310   335  
311   @param backend The backend tag value selecting the I/O 336   @param backend The backend tag value selecting the I/O
312   multiplexer (e.g. `corosio::epoll`). 337   multiplexer (e.g. `corosio::epoll`).
313   @param opts Runtime options controlling scheduler and 338   @param opts Runtime options controlling scheduler and
314   service behavior. 339   service behavior.
315   @param concurrency_hint Hint for the number of threads 340   @param concurrency_hint Hint for the number of threads
316   that will call `run()`. 341   that will call `run()`.
317   342  
318   @throws std::invalid_argument If `opts.thread_pool_size` is 343   @throws std::invalid_argument If `opts.thread_pool_size` is
319   less than 1 (POSIX). 344   less than 1 (POSIX).
  345 +
  346 + @throws std::system_error If the backend's infrastructure
  347 + could not be created.
320   */ 348   */
321   template<class Backend> 349   template<class Backend>
322   requires requires { Backend::construct; } 350   requires requires { Backend::construct; }
HITCBC 323   19 explicit io_context( 351   19 explicit io_context(
324   [[maybe_unused]] Backend backend, 352   [[maybe_unused]] Backend backend,
325   io_context_options const& opts, 353   io_context_options const& opts,
326   unsigned concurrency_hint = std::thread::hardware_concurrency()) 354   unsigned concurrency_hint = std::thread::hardware_concurrency())
327   : capy::execution_context(this) 355   : capy::execution_context(this)
HITCBC 328   19 , sched_(nullptr) 356   19 , sched_(nullptr)
329   { 357   {
HITCBC 330   19 apply_options_pre_(opts); 358   19 apply_options_pre_(opts);
331   // Effective hint (1 for lockless tiers); see effective_concurrency_hint. 359   // Effective hint (1 for lockless tiers); see effective_concurrency_hint.
332   unsigned const eff = 360   unsigned const eff =
HITCBC 333   19 detail::effective_concurrency_hint(opts, concurrency_hint); 361   19 detail::effective_concurrency_hint(opts, concurrency_hint);
HITCBC 334   19 sched_ = &Backend::construct(*this, eff); 362   19 sched_ = &Backend::construct(*this, eff);
HITCBC 335   19 apply_options_post_(opts, eff); 363   19 apply_options_post_(opts, eff);
HITCBC 336   19 } 364   19 }
337   365  
338   ~io_context(); 366   ~io_context();
339   367  
340   io_context(io_context const&) = delete; 368   io_context(io_context const&) = delete;
341   io_context& operator=(io_context const&) = delete; 369   io_context& operator=(io_context const&) = delete;
342   370  
343   /** Return an executor for this context. 371   /** Return an executor for this context.
344   372  
345   The returned executor can be used to dispatch coroutines 373   The returned executor can be used to dispatch coroutines
346   and post work items to this context. 374   and post work items to this context.
347   375  
348   @return An executor associated with this context. 376   @return An executor associated with this context.
349   */ 377   */
350   executor_type get_executor() const noexcept; 378   executor_type get_executor() const noexcept;
351   379  
352   /** Signal the context to stop processing. 380   /** Signal the context to stop processing.
353   381  
354   This causes `run()` to return as soon as possible. Any pending 382   This causes `run()` to return as soon as possible. Any pending
355   work items remain queued. 383   work items remain queued.
356   */ 384   */
HITCBC 357   7 void stop() 385   13 void stop()
358   { 386   {
HITCBC 359   7 sched_->stop(); 387   13 sched_->stop();
HITCBC 360   7 } 388   13 }
361   389  
362   /** Return whether the context has been stopped. 390   /** Return whether the context has been stopped.
363   391  
364   @return `true` if `stop()` has been called and `restart()` 392   @return `true` if `stop()` has been called and `restart()`
365   has not been called since. 393   has not been called since.
366   */ 394   */
HITCBC 367   74 bool stopped() const noexcept 395   74 bool stopped() const noexcept
368   { 396   {
HITCBC 369   74 return sched_->stopped(); 397   74 return sched_->stopped();
370   } 398   }
371   399  
372   /** Restart the context after being stopped. 400   /** Restart the context after being stopped.
373   401  
374   This function must be called before `run()` can be called 402   This function must be called before `run()` can be called
375   again after `stop()` has been called. 403   again after `stop()` has been called.
376   */ 404   */
HITCBC 377   311 void restart() 405   357 void restart()
378   { 406   {
HITCBC 379   311 sched_->restart(); 407   357 sched_->restart();
HITCBC 380   311 } 408   357 }
381   409  
382   /** Process all pending work items. 410   /** Process all pending work items.
383   411  
384   This function blocks until all pending work items have been 412   This function blocks until all pending work items have been
385   executed or `stop()` is called. The context is stopped 413   executed or `stop()` is called. The context is stopped
386   when there is no more outstanding work. 414   when there is no more outstanding work.
387   415  
388   @note The context must be restarted with `restart()` before 416   @note The context must be restarted with `restart()` before
389   calling this function again after it returns. 417   calling this function again after it returns.
390   418  
391   @return The number of handlers executed. 419   @return The number of handlers executed.
392   */ 420   */
HITCBC 393   1316 std::size_t run() 421   1419 std::size_t run()
394   { 422   {
HITCBC 395   1316 return sched_->run(); 423   1419 return sched_->run();
396   } 424   }
397   425  
398   /** Process at most one pending work item. 426   /** Process at most one pending work item.
399   427  
400   This function blocks until one work item has been executed 428   This function blocks until one work item has been executed
401   or `stop()` is called. The context is stopped when there 429   or `stop()` is called. The context is stopped when there
402   is no more outstanding work. 430   is no more outstanding work.
403   431  
404   @note The context must be restarted with `restart()` before 432   @note The context must be restarted with `restart()` before
405   calling this function again after it returns. 433   calling this function again after it returns.
406   434  
407   @return The number of handlers executed (0 or 1). 435   @return The number of handlers executed (0 or 1).
408   */ 436   */
HITCBC 409   29 std::size_t run_one() 437   110 std::size_t run_one()
410   { 438   {
HITCBC 411   29 return sched_->run_one(); 439   110 return sched_->run_one();
412   } 440   }
413   441  
414   /** Process work items for the specified duration. 442   /** Process work items for the specified duration.
415   443  
416   This function blocks until work items have been executed for 444   This function blocks until work items have been executed for
417   the specified duration, or `stop()` is called. The context 445   the specified duration, or `stop()` is called. The context
418   is stopped when there is no more outstanding work. 446   is stopped when there is no more outstanding work.
419   447  
420   @note The context must be restarted with `restart()` before 448   @note The context must be restarted with `restart()` before
421   calling this function again after it returns. 449   calling this function again after it returns.
422   450  
423   @param rel_time The duration for which to process work. 451   @param rel_time The duration for which to process work.
424   452  
425   @return The number of handlers executed. 453   @return The number of handlers executed.
426   */ 454   */
427   template<class Rep, class Period> 455   template<class Rep, class Period>
HITCBC 428   11 std::size_t run_for(std::chrono::duration<Rep, Period> const& rel_time) 456   11 std::size_t run_for(std::chrono::duration<Rep, Period> const& rel_time)
429   { 457   {
HITCBC 430   11 return run_until(std::chrono::steady_clock::now() + rel_time); 458   11 return run_until(std::chrono::steady_clock::now() + rel_time);
431   } 459   }
432   460  
433   /** Process work items until the specified time. 461   /** Process work items until the specified time.
434   462  
435   This function blocks until the specified time is reached 463   This function blocks until the specified time is reached
436   or `stop()` is called. The context is stopped when there 464   or `stop()` is called. The context is stopped when there
437   is no more outstanding work. 465   is no more outstanding work.
438   466  
439   @note The context must be restarted with `restart()` before 467   @note The context must be restarted with `restart()` before
440   calling this function again after it returns. 468   calling this function again after it returns.
441   469  
442   @param abs_time The time point until which to process work. 470   @param abs_time The time point until which to process work.
443   471  
444   @return The number of handlers executed. 472   @return The number of handlers executed.
445   */ 473   */
446   template<class Clock, class Duration> 474   template<class Clock, class Duration>
447   std::size_t 475   std::size_t
HITCBC 448   12 run_until(std::chrono::time_point<Clock, Duration> const& abs_time) 476   12 run_until(std::chrono::time_point<Clock, Duration> const& abs_time)
449   { 477   {
HITCBC 450   12 std::size_t n = 0; 478   12 std::size_t n = 0;
HITCBC 451   30 while (run_one_until(abs_time)) 479   30 while (run_one_until(abs_time))
HITCBC 452   18 if (n != (std::numeric_limits<std::size_t>::max)()) 480   18 if (n != (std::numeric_limits<std::size_t>::max)())
HITCBC 453   18 ++n; 481   18 ++n;
HITCBC 454   12 return n; 482   12 return n;
455   } 483   }
456   484  
457   /** Process at most one work item for the specified duration. 485   /** Process at most one work item for the specified duration.
458   486  
459   This function blocks until one work item has been executed, 487   This function blocks until one work item has been executed,
460   the specified duration has elapsed, or `stop()` is called. 488   the specified duration has elapsed, or `stop()` is called.
461   The context is stopped when there is no more outstanding work. 489   The context is stopped when there is no more outstanding work.
462   490  
463   @note The context must be restarted with `restart()` before 491   @note The context must be restarted with `restart()` before
464   calling this function again after it returns. 492   calling this function again after it returns.
465   493  
466   @param rel_time The duration for which the call may block. 494   @param rel_time The duration for which the call may block.
467   495  
468   @return The number of handlers executed (0 or 1). 496   @return The number of handlers executed (0 or 1).
469   */ 497   */
470   template<class Rep, class Period> 498   template<class Rep, class Period>
HITCBC 471   6 std::size_t run_one_for(std::chrono::duration<Rep, Period> const& rel_time) 499   6 std::size_t run_one_for(std::chrono::duration<Rep, Period> const& rel_time)
472   { 500   {
HITCBC 473   6 return run_one_until(std::chrono::steady_clock::now() + rel_time); 501   6 return run_one_until(std::chrono::steady_clock::now() + rel_time);
474   } 502   }
475   503  
476   /** Process at most one work item until the specified time. 504   /** Process at most one work item until the specified time.
477   505  
478   This function blocks until one work item has been executed, 506   This function blocks until one work item has been executed,
479   the specified time is reached, or `stop()` is called. 507   the specified time is reached, or `stop()` is called.
480   The context is stopped when there is no more outstanding work. 508   The context is stopped when there is no more outstanding work.
481   509  
482   @note The context must be restarted with `restart()` before 510   @note The context must be restarted with `restart()` before
483   calling this function again after it returns. 511   calling this function again after it returns.
484   512  
485   @param abs_time The time point until which the call may block. 513   @param abs_time The time point until which the call may block.
486   514  
487   @return The number of handlers executed (0 or 1). 515   @return The number of handlers executed (0 or 1).
488   */ 516   */
489   template<class Clock, class Duration> 517   template<class Clock, class Duration>
490   std::size_t 518   std::size_t
HITCBC 491   44 run_one_until(std::chrono::time_point<Clock, Duration> const& abs_time) 519   44 run_one_until(std::chrono::time_point<Clock, Duration> const& abs_time)
492   { 520   {
HITCBC 493   44 typename Clock::time_point now = Clock::now(); 521   44 typename Clock::time_point now = Clock::now();
HITCBC 494   8 for (;;) 522   8 for (;;)
495   { 523   {
HITCBC 496   52 auto rel_time = abs_time - now; 524   52 auto rel_time = abs_time - now;
497   using rel_type = decltype(rel_time); 525   using rel_type = decltype(rel_time);
HITCBC 498   52 if (rel_time < rel_type::zero()) 526   52 if (rel_time < rel_type::zero())
HITCBC 499   5 rel_time = rel_type::zero(); 527   5 rel_time = rel_type::zero();
HITCBC 500   47 else if (rel_time > std::chrono::seconds(1)) 528   47 else if (rel_time > std::chrono::seconds(1))
HITCBC 501   22 rel_time = std::chrono::seconds(1); 529   22 rel_time = std::chrono::seconds(1);
502   530  
HITCBC 503   52 std::size_t s = sched_->wait_one( 531   52 std::size_t s = sched_->wait_one(
504   static_cast<long>( 532   static_cast<long>(
HITCBC 505   52 std::chrono::duration_cast<std::chrono::microseconds>( 533   52 std::chrono::duration_cast<std::chrono::microseconds>(
506   rel_time) 534   rel_time)
HITCBC 507   52 .count())); 535   52 .count()));
508   536  
HITCBC 509   52 if (s || stopped()) 537   52 if (s || stopped())
HITCBC 510   44 return s; 538   44 return s;
511   539  
HITCBC 512   12 now = Clock::now(); 540   12 now = Clock::now();
HITCBC 513   12 if (now >= abs_time) 541   12 if (now >= abs_time)
HITCBC 514   4 return 0; 542   4 return 0;
515   } 543   }
516   } 544   }
517   545  
518   /** Process all ready work items without blocking. 546   /** Process all ready work items without blocking.
519   547  
520   This function executes all work items that are ready to run 548   This function executes all work items that are ready to run
521   without blocking for more work. The context is stopped 549   without blocking for more work. The context is stopped
522   when there is no more outstanding work. 550   when there is no more outstanding work.
523   551  
524   @note The context must be restarted with `restart()` before 552   @note The context must be restarted with `restart()` before
525   calling this function again after it returns. 553   calling this function again after it returns.
526   554  
527   @return The number of handlers executed. 555   @return The number of handlers executed.
528   */ 556   */
HITCBC 529   31 std::size_t poll() 557   31 std::size_t poll()
530   { 558   {
HITCBC 531   31 return sched_->poll(); 559   31 return sched_->poll();
532   } 560   }
533   561  
534   /** Process at most one ready work item without blocking. 562   /** Process at most one ready work item without blocking.
535   563  
536   This function executes at most one work item that is ready 564   This function executes at most one work item that is ready
537   to run without blocking for more work. The context is 565   to run without blocking for more work. The context is
538   stopped when there is no more outstanding work. 566   stopped when there is no more outstanding work.
539   567  
540   @note The context must be restarted with `restart()` before 568   @note The context must be restarted with `restart()` before
541   calling this function again after it returns. 569   calling this function again after it returns.
542   570  
543   @return The number of handlers executed (0 or 1). 571   @return The number of handlers executed (0 or 1).
544   */ 572   */
HITCBC 545   9 std::size_t poll_one() 573   9 std::size_t poll_one()
546   { 574   {
HITCBC 547   9 return sched_->poll_one(); 575   9 return sched_->poll_one();
548   } 576   }
549   }; 577   };
550   578  
551   /** An executor for dispatching work to an I/O context. 579   /** An executor for dispatching work to an I/O context.
552   580  
553   The executor provides the interface for posting work items and 581   The executor provides the interface for posting work items and
554   dispatching coroutines to the associated context. It satisfies 582   dispatching coroutines to the associated context. It satisfies
555   the `capy::Executor` concept. 583   the `capy::Executor` concept.
556   584  
557   Executors are lightweight handles that can be copied and compared 585   Executors are lightweight handles that can be copied and compared
558   for equality. Two executors compare equal if they refer to the 586   for equality. Two executors compare equal if they refer to the
559   same context. 587   same context.
560   588  
561   @par Thread Safety 589   @par Thread Safety
562   Distinct objects: Safe.@n 590   Distinct objects: Safe.@n
563   Shared objects: Safe. 591   Shared objects: Safe.
564   */ 592   */
565   class io_context::executor_type 593   class io_context::executor_type
566   { 594   {
567   io_context* ctx_ = nullptr; 595   io_context* ctx_ = nullptr;
568   596  
569   public: 597   public:
570   /** Default constructor. 598   /** Default constructor.
571   599  
572   Constructs an executor not associated with any context. 600   Constructs an executor not associated with any context.
573   */ 601   */
HITCBC 574   2053 executor_type() = default; 602   2053 executor_type() = default;
575   603  
576   /** Construct an executor from a context. 604   /** Construct an executor from a context.
577   605  
578   @param ctx The context to associate with this executor. 606   @param ctx The context to associate with this executor.
579   */ 607   */
HITCBC 580   3662 explicit executor_type(io_context& ctx) noexcept : ctx_(&ctx) {} 608   3830 explicit executor_type(io_context& ctx) noexcept : ctx_(&ctx) {}
581   609  
582   /** Return a reference to the associated execution context. 610   /** Return a reference to the associated execution context.
583   611  
584   @return Reference to the context. 612   @return Reference to the context.
585   */ 613   */
HITCBC 586   16597 io_context& context() const noexcept 614   17777 io_context& context() const noexcept
587   { 615   {
HITCBC 588   16597 return *ctx_; 616   17777 return *ctx_;
589   } 617   }
590   618  
591   /** Check if the current thread is running this executor's context. 619   /** Check if the current thread is running this executor's context.
592   620  
593   @return `true` if `run()` is being called on this thread. 621   @return `true` if `run()` is being called on this thread.
594   */ 622   */
HITCBC 595   7724 bool running_in_this_thread() const noexcept 623   7955 bool running_in_this_thread() const noexcept
596   { 624   {
HITCBC 597   7724 return ctx_->sched_->running_in_this_thread(); 625   7955 return ctx_->sched_->running_in_this_thread();
598   } 626   }
599   627  
600   /** Informs the executor that work is beginning. 628   /** Informs the executor that work is beginning.
601   629  
602   Must be paired with `on_work_finished()`. 630   Must be paired with `on_work_finished()`.
603   */ 631   */
HITCBC 604   8023 void on_work_started() const noexcept 632   8310 void on_work_started() const noexcept
605   { 633   {
HITCBC 606   8023 ctx_->sched_->work_started(); 634   8310 ctx_->sched_->work_started();
HITCBC 607   8023 } 635   8310 }
608   636  
609   /** Informs the executor that work has completed. 637   /** Informs the executor that work has completed.
610   638  
611   @par Preconditions 639   @par Preconditions
612   A preceding call to `on_work_started()` on an equal executor. 640   A preceding call to `on_work_started()` on an equal executor.
613   */ 641   */
HITCBC 614   7971 void on_work_finished() const noexcept 642   8248 void on_work_finished() const noexcept
615   { 643   {
HITCBC 616   7971 ctx_->sched_->work_finished(); 644   8248 ctx_->sched_->work_finished();
HITCBC 617   7971 } 645   8248 }
618   646  
619   /** Dispatch a continuation. 647   /** Dispatch a continuation.
620   648  
621   Returns a handle for symmetric transfer. If called from 649   Returns a handle for symmetric transfer. If called from
622   within `run()`, returns `c.h`. Otherwise posts `c` for 650   within `run()`, returns `c.h`. Otherwise posts `c` for
623   later execution and returns `std::noop_coroutine()`. 651   later execution and returns `std::noop_coroutine()`.
624   652  
625   @param c The continuation to dispatch. 653   @param c The continuation to dispatch.
626   654  
627   @return A handle for symmetric transfer or `std::noop_coroutine()`. 655   @return A handle for symmetric transfer or `std::noop_coroutine()`.
628   656  
629   @par Preconditions 657   @par Preconditions
630   The associated context must outlive this call. Dispatching 658   The associated context must outlive this call. Dispatching
631   concurrently with, or after, the context's destruction is 659   concurrently with, or after, the context's destruction is
632   undefined behavior. 660   undefined behavior.
633   */ 661   */
HITCBC 634   7719 std::coroutine_handle<> dispatch(capy::continuation& c) const 662   7950 std::coroutine_handle<> dispatch(capy::continuation& c) const
635   { 663   {
HITCBC 636   7719 if (running_in_this_thread()) 664   7950 if (running_in_this_thread())
HITCBC 637   677 return c.h; 665   683 return c.h;
HITCBC 638   7042 post(c); 666   7267 post(c);
HITCBC 639   7042 return std::noop_coroutine(); 667   7267 return std::noop_coroutine();
640   } 668   }
641   669  
642   /** Post a continuation for deferred execution. 670   /** Post a continuation for deferred execution.
643   671  
644   Enqueues `c` directly on the scheduler's ready queue. 672   Enqueues `c` directly on the scheduler's ready queue.
645   No heap allocation occurs. 673   No heap allocation occurs.
646   674  
647   @par Preconditions 675   @par Preconditions
648   The associated context must outlive this call. Posting 676   The associated context must outlive this call. Posting
649   concurrently with, or after, the context's destruction is 677   concurrently with, or after, the context's destruction is
650   undefined behavior. 678   undefined behavior.
651   */ 679   */
HITCBC 652   14240 void post(capy::continuation& c) const 680   15417 void post(capy::continuation& c) const
653   { 681   {
HITCBC 654   14240 ctx_->sched_->post(c); 682   15417 ctx_->sched_->post(c);
HITCBC 655   14240 } 683   15417 }
656   684  
657   /** Post a bare coroutine handle for deferred execution. 685   /** Post a bare coroutine handle for deferred execution.
658   686  
659   Heap-allocates a scheduler_op to wrap the handle. A caller 687   Heap-allocates a scheduler_op to wrap the handle. A caller
660   that already owns a `scheduler_op` can post it directly via 688   that already owns a `scheduler_op` can post it directly via
661   the `post(scheduler_op*)` overload to avoid the allocation. 689   the `post(scheduler_op*)` overload to avoid the allocation.
662   690  
663   @param h The coroutine handle to post. 691   @param h The coroutine handle to post.
664   692  
665   @par Preconditions 693   @par Preconditions
666   The associated context must outlive this call. Posting 694   The associated context must outlive this call. Posting
667   concurrently with, or after, the context's destruction is 695   concurrently with, or after, the context's destruction is
668   undefined behavior. 696   undefined behavior.
669   */ 697   */
HITCBC 670   3686 void post(std::coroutine_handle<> h) const 698   3686 void post(std::coroutine_handle<> h) const
671   { 699   {
HITCBC 672   3686 ctx_->sched_->post(h); 700   3686 ctx_->sched_->post(h);
HITCBC 673   3686 } 701   3686 }
674   702  
675   /** Compare two executors for equality. 703   /** Compare two executors for equality.
676   704  
677   @return `true` if both executors refer to the same context. 705   @return `true` if both executors refer to the same context.
678   */ 706   */
HITCBC 679   2 bool operator==(executor_type const& other) const noexcept 707   2 bool operator==(executor_type const& other) const noexcept
680   { 708   {
HITCBC 681   2 return ctx_ == other.ctx_; 709   2 return ctx_ == other.ctx_;
682   } 710   }
683   711  
684   /** Compare two executors for inequality. 712   /** Compare two executors for inequality.
685   713  
686   @return `true` if the executors refer to different contexts. 714   @return `true` if the executors refer to different contexts.
687   */ 715   */
688   bool operator!=(executor_type const& other) const noexcept 716   bool operator!=(executor_type const& other) const noexcept
689   { 717   {
690   return ctx_ != other.ctx_; 718   return ctx_ != other.ctx_;
691   } 719   }
692   }; 720   };
693   721  
694   inline io_context::executor_type 722   inline io_context::executor_type
HITCBC 695   3662 io_context::get_executor() const noexcept 723   3830 io_context::get_executor() const noexcept
696   { 724   {
HITCBC 697   3662 return executor_type(const_cast<io_context&>(*this)); 725   3830 return executor_type(const_cast<io_context&>(*this));
698   } 726   }
699   727  
700   } // namespace boost::corosio 728   } // namespace boost::corosio
701   729  
702   #endif // BOOST_COROSIO_IO_CONTEXT_HPP 730   #endif // BOOST_COROSIO_IO_CONTEXT_HPP