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