97.97% Lines (145/148) 97.30% Functions (36/37)
TLA Baseline Branch
Line Hits Code Line Hits Code
1   // 1   //
2   // Copyright (c) 2026 Vinnie Falco (vinnie.falco@gmail.com) 2   // Copyright (c) 2026 Vinnie Falco (vinnie.falco@gmail.com)
3   // Copyright (c) 2026 Michael Vandeberg 3   // Copyright (c) 2026 Michael Vandeberg
4   // 4   //
5   // Distributed under the Boost Software License, Version 1.0. (See accompanying 5   // Distributed under the Boost Software License, Version 1.0. (See accompanying
6   // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt) 6   // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
7   // 7   //
8   // Official repository: https://github.com/cppalliance/corosio 8   // Official repository: https://github.com/cppalliance/corosio
9   // 9   //
10   10  
11   #ifndef BOOST_COROSIO_TCP_SERVER_HPP 11   #ifndef BOOST_COROSIO_TCP_SERVER_HPP
12   #define BOOST_COROSIO_TCP_SERVER_HPP 12   #define BOOST_COROSIO_TCP_SERVER_HPP
13   13  
14   #include <boost/corosio/detail/config.hpp> 14   #include <boost/corosio/detail/config.hpp>
15   #include <boost/corosio/detail/except.hpp> 15   #include <boost/corosio/detail/except.hpp>
16   #include <boost/corosio/tcp_acceptor.hpp> 16   #include <boost/corosio/tcp_acceptor.hpp>
17   #include <boost/corosio/tcp_socket.hpp> 17   #include <boost/corosio/tcp_socket.hpp>
18   #include <boost/corosio/io_context.hpp> 18   #include <boost/corosio/io_context.hpp>
19   #include <boost/corosio/endpoint.hpp> 19   #include <boost/corosio/endpoint.hpp>
20   #include <boost/capy/task.hpp> 20   #include <boost/capy/task.hpp>
21   #include <boost/capy/concept/execution_context.hpp> 21   #include <boost/capy/concept/execution_context.hpp>
22   #include <boost/capy/concept/io_awaitable.hpp> 22   #include <boost/capy/concept/io_awaitable.hpp>
23   #include <boost/capy/concept/executor.hpp> 23   #include <boost/capy/concept/executor.hpp>
24   #include <boost/capy/ex/any_executor.hpp> 24   #include <boost/capy/ex/any_executor.hpp>
25   #include <boost/capy/ex/frame_allocator.hpp> 25   #include <boost/capy/ex/frame_allocator.hpp>
26   #include <boost/capy/ex/io_env.hpp> 26   #include <boost/capy/ex/io_env.hpp>
27   #include <boost/capy/ex/run_async.hpp> 27   #include <boost/capy/ex/run_async.hpp>
28   28  
29   #include <coroutine> 29   #include <coroutine>
30   #include <memory> 30   #include <memory>
31   #include <ranges> 31   #include <ranges>
32   #include <vector> 32   #include <vector>
33   33  
34   namespace boost::corosio { 34   namespace boost::corosio {
35   35  
36   #ifdef _MSC_VER 36   #ifdef _MSC_VER
37   #pragma warning(push) 37   #pragma warning(push)
38   #pragma warning(disable : 4251) // class needs to have dll-interface 38   #pragma warning(disable : 4251) // class needs to have dll-interface
39   #endif 39   #endif
40   40  
41   /** Manages a pool of reusable workers that handle incoming TCP connections. 41   /** Manages a pool of reusable workers that handle incoming TCP connections.
42   42  
43   This class manages a pool of reusable worker objects that handle 43   This class manages a pool of reusable worker objects that handle
44   incoming connections. When a connection arrives, an idle worker 44   incoming connections. When a connection arrives, an idle worker
45   is dispatched to handle it. After the connection completes, the 45   is dispatched to handle it. After the connection completes, the
46   worker returns to the pool for reuse, avoiding allocation overhead 46   worker returns to the pool for reuse, avoiding allocation overhead
47   per connection. 47   per connection.
48   48  
49   Workers are set via @ref set_workers as a forward range of 49   Workers are set via @ref set_workers as a forward range of
50   pointer-like objects (e.g., `unique_ptr<worker_base>`). The server 50   pointer-like objects (e.g., `unique_ptr<worker_base>`). The server
51   takes ownership of the container via type erasure. 51   takes ownership of the container via type erasure.
52   52  
53   @par Thread Safety 53   @par Thread Safety
54   Distinct objects: Safe. 54   Distinct objects: Safe.
55   Shared objects: Unsafe. 55   Shared objects: Unsafe.
56   56  
57   @par Lifecycle 57   @par Lifecycle
58   The server operates in three states: 58   The server operates in three states:
59   59  
60   - **Stopped**: Initial state, or after @ref join completes. 60   - **Stopped**: Initial state, or after @ref join completes.
61   - **Running**: After @ref start, actively accepting connections. 61   - **Running**: After @ref start, actively accepting connections.
62   - **Stopping**: After @ref stop, draining active work. 62   - **Stopping**: After @ref stop, draining active work.
63   63  
64   State transitions: 64   State transitions:
65   @code 65   @code
66   [Stopped] --start()--> [Running] --stop()--> [Stopping] --join()--> [Stopped] 66   [Stopped] --start()--> [Running] --stop()--> [Stopping] --join()--> [Stopped]
67   @endcode 67   @endcode
68   68  
69   @par Running the Server 69   @par Running the Server
70   @par !example running_the_server 70   @par !example running_the_server
71   71  
72   @par Graceful Shutdown 72   @par Graceful Shutdown
73   To shut down gracefully, call @ref stop then drain the `io_context`: 73   To shut down gracefully, call @ref stop then drain the `io_context`:
74   @par !example graceful_shutdown 74   @par !example graceful_shutdown
75   75  
76   @par Restart After Stop 76   @par Restart After Stop
77   The server can be restarted after a complete shutdown cycle. 77   The server can be restarted after a complete shutdown cycle.
78   You must drain the `io_context`, call @ref join, and restart the 78   You must drain the `io_context`, call @ref join, and restart the
79   `io_context` itself (`ioc.restart()`) before restarting: 79   `io_context` itself (`ioc.restart()`) before restarting:
80   @par !example restart_after_stop 80   @par !example restart_after_stop
81   81  
82   @par WARNING: What NOT to Do 82   @par WARNING: What NOT to Do
83   - Do NOT call @ref join from inside a worker coroutine (deadlock). 83   - Do NOT call @ref join from inside a worker coroutine (deadlock).
84   - Do NOT call @ref join from a thread running `ioc.run()` (deadlock). 84   - Do NOT call @ref join from a thread running `ioc.run()` (deadlock).
85   - Do NOT call @ref start without completing @ref join after @ref stop. 85   - Do NOT call @ref start without completing @ref join after @ref stop.
86   - Do NOT call `ioc.stop()` for graceful shutdown; use @ref stop instead. 86   - Do NOT call `ioc.stop()` for graceful shutdown; use @ref stop instead.
87   87  
88   @par Example 88   @par Example
89   @par !example custom_worker 89   @par !example custom_worker
90   90  
91   @see worker_base, set_workers, launcher 91   @see worker_base, set_workers, launcher
92   */ 92   */
93   class BOOST_COROSIO_DECL tcp_server 93   class BOOST_COROSIO_DECL tcp_server
94   { 94   {
95   public: 95   public:
96   class worker_base; ///< Abstract base for connection handlers. 96   class worker_base; ///< Abstract base for connection handlers.
97   class launcher; ///< Move-only handle to launch worker coroutines. 97   class launcher; ///< Move-only handle to launch worker coroutines.
98   98  
99   private: 99   private:
100   struct waiter 100   struct waiter
101   { 101   {
102   waiter* next; 102   waiter* next;
103   std::coroutine_handle<> h; 103   std::coroutine_handle<> h;
104   capy::continuation cont; 104   capy::continuation cont;
105   worker_base* w; 105   worker_base* w;
106   }; 106   };
107   107  
108   struct impl; 108   struct impl;
109   109  
110   static impl* make_impl(capy::execution_context& ctx); 110   static impl* make_impl(capy::execution_context& ctx);
111   111  
112   impl* impl_; 112   impl* impl_;
113   capy::any_executor ex_; 113   capy::any_executor ex_;
114   waiter* waiters_ = nullptr; 114   waiter* waiters_ = nullptr;
115   worker_base* idle_head_ = nullptr; // Forward list: available workers 115   worker_base* idle_head_ = nullptr; // Forward list: available workers
116   worker_base* active_head_ = 116   worker_base* active_head_ =
117   nullptr; // Doubly linked: workers handling connections 117   nullptr; // Doubly linked: workers handling connections
118   worker_base* active_tail_ = nullptr; // Tail for O(1) push_back 118   worker_base* active_tail_ = nullptr; // Tail for O(1) push_back
119   std::size_t active_accepts_ = 0; // Number of active do_accept coroutines 119   std::size_t active_accepts_ = 0; // Number of active do_accept coroutines
120   std::shared_ptr<void> storage_; // Owns the worker container (type-erased) 120   std::shared_ptr<void> storage_; // Owns the worker container (type-erased)
121   bool running_ = false; 121   bool running_ = false;
122   122  
123   // Idle list (forward/singly linked) - push front, pop front 123   // Idle list (forward/singly linked) - push front, pop front
HITCBC 124   238 void idle_push(worker_base* w) noexcept 124   238 void idle_push(worker_base* w) noexcept
125   { 125   {
HITCBC 126   238 w->next_ = idle_head_; 126   238 w->next_ = idle_head_;
HITCBC 127   238 idle_head_ = w; 127   238 idle_head_ = w;
HITCBC 128   238 } 128   238 }
129   129  
HITCBC 130   75 worker_base* idle_pop() noexcept 130   75 worker_base* idle_pop() noexcept
131   { 131   {
HITCBC 132   75 auto* w = idle_head_; 132   75 auto* w = idle_head_;
HITCBC 133   75 if (w) 133   75 if (w)
HITCBC 134   75 idle_head_ = w->next_; 134   75 idle_head_ = w->next_;
HITCBC 135   75 return w; 135   75 return w;
136   } 136   }
137   137  
HITCBC 138   153 bool idle_empty() const noexcept 138   153 bool idle_empty() const noexcept
139   { 139   {
HITCBC 140   153 return idle_head_ == nullptr; 140   153 return idle_head_ == nullptr;
141   } 141   }
142   142  
143   // Active list (doubly linked) - push back, remove anywhere 143   // Active list (doubly linked) - push back, remove anywhere
HITCBC 144   88 void active_push(worker_base* w) noexcept 144   88 void active_push(worker_base* w) noexcept
145   { 145   {
HITCBC 146   88 w->next_ = nullptr; 146   88 w->next_ = nullptr;
HITCBC 147   88 w->prev_ = active_tail_; 147   88 w->prev_ = active_tail_;
HITCBC 148   88 if (active_tail_) 148   88 if (active_tail_)
HITCBC 149   4 active_tail_->next_ = w; 149   4 active_tail_->next_ = w;
150   else 150   else
HITCBC 151   84 active_head_ = w; 151   84 active_head_ = w;
HITCBC 152   88 active_tail_ = w; 152   88 active_tail_ = w;
HITCBC 153   88 } 153   88 }
154   154  
HITCBC 155   153 void active_remove(worker_base* w) noexcept 155   153 void active_remove(worker_base* w) noexcept
156   { 156   {
157   // Skip if not in active list (e.g., after failed accept) 157   // Skip if not in active list (e.g., after failed accept)
HITCBC 158   153 if (w != active_head_ && w->prev_ == nullptr) 158   153 if (w != active_head_ && w->prev_ == nullptr)
HITCBC 159   65 return; 159   65 return;
HITCBC 160   88 if (w->prev_) 160   88 if (w->prev_)
HITCBC 161   4 w->prev_->next_ = w->next_; 161   4 w->prev_->next_ = w->next_;
162   else 162   else
HITCBC 163   84 active_head_ = w->next_; 163   84 active_head_ = w->next_;
HITCBC 164   88 if (w->next_) 164   88 if (w->next_)
HITCBC 165   2 w->next_->prev_ = w->prev_; 165   2 w->next_->prev_ = w->prev_;
166   else 166   else
HITCBC 167   86 active_tail_ = w->prev_; 167   86 active_tail_ = w->prev_;
HITCBC 168   88 w->prev_ = nullptr; // Mark as not in active list 168   88 w->prev_ = nullptr; // Mark as not in active list
169   } 169   }
170   170  
171   template<capy::Executor Ex> 171   template<capy::Executor Ex>
172   struct launch_wrapper 172   struct launch_wrapper
173   { 173   {
174   struct promise_type 174   struct promise_type
175   { 175   {
176   Ex ex; // Executor stored directly in frame (outlives child tasks) 176   Ex ex; // Executor stored directly in frame (outlives child tasks)
177   capy::io_env env_; 177   capy::io_env env_;
178   178  
179   // For regular coroutines: first arg is executor, second is stop token 179   // For regular coroutines: first arg is executor, second is stop token
180   template<class E, class S, class... Args> 180   template<class E, class S, class... Args>
181   requires capy::Executor<std::decay_t<E>> 181   requires capy::Executor<std::decay_t<E>>
182   promise_type(E e, S s, Args&&...) 182   promise_type(E e, S s, Args&&...)
183   : ex(std::move(e)) 183   : ex(std::move(e))
184   , env_{ 184   , env_{
185   capy::executor_ref(ex), std::move(s), 185   capy::executor_ref(ex), std::move(s),
186   capy::get_current_frame_allocator()} 186   capy::get_current_frame_allocator()}
187   { 187   {
188   } 188   }
189   189  
190   // For lambda coroutines: first arg is closure, second is executor, third is stop token 190   // For lambda coroutines: first arg is closure, second is executor, third is stop token
191   template<class Closure, class E, class S, class... Args> 191   template<class Closure, class E, class S, class... Args>
192   requires(!capy::Executor<std::decay_t<Closure>> && 192   requires(!capy::Executor<std::decay_t<Closure>> &&
193   capy::Executor<std::decay_t<E>>) 193   capy::Executor<std::decay_t<E>>)
HITCBC 194   88 promise_type(Closure&&, E e, S s, Args&&...) 194   88 promise_type(Closure&&, E e, S s, Args&&...)
HITCBC 195   88 : ex(std::move(e)) 195   88 : ex(std::move(e))
HITCBC 196   88 , env_{ 196   88 , env_{
HITCBC 197   88 capy::executor_ref(ex), std::move(s), 197   88 capy::executor_ref(ex), std::move(s),
HITCBC 198   88 capy::get_current_frame_allocator()} 198   88 capy::get_current_frame_allocator()}
199   { 199   {
HITCBC 200   88 } 200   88 }
201   201  
HITCBC 202   88 launch_wrapper get_return_object() noexcept 202   88 launch_wrapper get_return_object() noexcept
203   { 203   {
204   return { 204   return {
HITCBC 205   88 std::coroutine_handle<promise_type>::from_promise(*this)}; 205   88 std::coroutine_handle<promise_type>::from_promise(*this)};
206   } 206   }
HITCBC 207   88 std::suspend_always initial_suspend() noexcept 207   88 std::suspend_always initial_suspend() noexcept
208   { 208   {
HITCBC 209   88 return {}; 209   88 return {};
210   } 210   }
HITCBC 211   88 std::suspend_never final_suspend() noexcept 211   88 std::suspend_never final_suspend() noexcept
212   { 212   {
HITCBC 213   88 return {}; 213   88 return {};
214   } 214   }
HITCBC 215   88 void return_void() noexcept {} 215   88 void return_void() noexcept {}
MISUBC 216   ✗ void unhandled_exception() 216   ✗ void unhandled_exception()
217   { 217   {
218   // LCOV_EXCL_START: terminating by contract is not a 218   // LCOV_EXCL_START: terminating by contract is not a
219   // coverable outcome. 219   // coverable outcome.
220   std::terminate(); 220   std::terminate();
221   // LCOV_EXCL_STOP 221   // LCOV_EXCL_STOP
222   } 222   }
223   223  
224   // Inject io_env for IoAwaitable 224   // Inject io_env for IoAwaitable
225   template<capy::IoAwaitable Awaitable> 225   template<capy::IoAwaitable Awaitable>
HITCBC 226   176 auto await_transform(Awaitable&& a) 226   176 auto await_transform(Awaitable&& a)
227   { 227   {
228   using AwaitableT = std::decay_t<Awaitable>; 228   using AwaitableT = std::decay_t<Awaitable>;
229   struct adapter 229   struct adapter
230   { 230   {
231   AwaitableT aw; 231   AwaitableT aw;
232   capy::io_env const* env; 232   capy::io_env const* env;
233   233  
HITCBC 234   176 bool await_ready() 234   176 bool await_ready()
235   { 235   {
HITCBC 236   176 return aw.await_ready(); 236   176 return aw.await_ready();
237   } 237   }
HITCBC 238   176 decltype(auto) await_resume() 238   176 decltype(auto) await_resume()
239   { 239   {
HITCBC 240   176 return aw.await_resume(); 240   176 return aw.await_resume();
241   } 241   }
242   242  
HITCBC 243   176 auto await_suspend(std::coroutine_handle<promise_type> h) 243   176 auto await_suspend(std::coroutine_handle<promise_type> h)
244   { 244   {
HITCBC 245   176 return aw.await_suspend(h, env); 245   176 return aw.await_suspend(h, env);
246   } 246   }
247   }; 247   };
HITCBC 248   264 return adapter{std::forward<Awaitable>(a), &env_}; 248   264 return adapter{std::forward<Awaitable>(a), &env_};
HITCBC 249   88 } 249   88 }
250   }; 250   };
251   251  
252   std::coroutine_handle<promise_type> h; 252   std::coroutine_handle<promise_type> h;
253   253  
HITCBC 254   88 launch_wrapper(std::coroutine_handle<promise_type> handle) noexcept 254   88 launch_wrapper(std::coroutine_handle<promise_type> handle) noexcept
HITCBC 255   88 : h(handle) 255   88 : h(handle)
256   { 256   {
HITCBC 257   88 } 257   88 }
258   258  
HITCBC 259   88 ~launch_wrapper() 259   88 ~launch_wrapper()
260   { 260   {
HITCBC 261   88 if (h) 261   88 if (h)
MISUBC 262   ✗ h.destroy(); 262   ✗ h.destroy();
HITCBC 263   88 } 263   88 }
264   264  
265   launch_wrapper(launch_wrapper&& o) noexcept 265   launch_wrapper(launch_wrapper&& o) noexcept
266   : h(std::exchange(o.h, nullptr)) 266   : h(std::exchange(o.h, nullptr))
267   { 267   {
268   } 268   }
269   269  
270   launch_wrapper(launch_wrapper const&) = delete; 270   launch_wrapper(launch_wrapper const&) = delete;
271   launch_wrapper& operator=(launch_wrapper const&) = delete; 271   launch_wrapper& operator=(launch_wrapper const&) = delete;
272   launch_wrapper& operator=(launch_wrapper&&) = delete; 272   launch_wrapper& operator=(launch_wrapper&&) = delete;
273   }; 273   };
274   274  
275   // Named functor to avoid incomplete lambda type in coroutine promise 275   // Named functor to avoid incomplete lambda type in coroutine promise
276   template<class Executor> 276   template<class Executor>
277   struct launch_coro 277   struct launch_coro
278   { 278   {
HITCBC 279   88 launch_wrapper<Executor> operator()( 279   88 launch_wrapper<Executor> operator()(
280   Executor, 280   Executor,
281   std::stop_token, 281   std::stop_token,
282   tcp_server* self, 282   tcp_server* self,
283   capy::task<void> t, 283   capy::task<void> t,
284   worker_base* wp) 284   worker_base* wp)
285   { 285   {
286   // Executor and stop token stored in promise via constructor 286   // Executor and stop token stored in promise via constructor
287   co_await std::move(t); 287   co_await std::move(t);
288   co_await self->push(*wp); // worker goes back to idle list 288   co_await self->push(*wp); // worker goes back to idle list
HITCBC 289   176 } 289   176 }
290   }; 290   };
291   291  
292   class push_awaitable 292   class push_awaitable
293   { 293   {
294   tcp_server& self_; 294   tcp_server& self_;
295   worker_base& w_; 295   worker_base& w_;
296   capy::continuation cont_; 296   capy::continuation cont_;
297   297  
298   public: 298   public:
HITCBC 299   145 push_awaitable(tcp_server& self, worker_base& w) noexcept 299   145 push_awaitable(tcp_server& self, worker_base& w) noexcept
HITCBC 300   145 : self_(self) 300   145 : self_(self)
HITCBC 301   145 , w_(w) 301   145 , w_(w)
302   { 302   {
HITCBC 303   145 } 303   145 }
304   304  
HITCBC 305   145 bool await_ready() const noexcept 305   145 bool await_ready() const noexcept
306   { 306   {
HITCBC 307   145 return false; 307   145 return false;
308   } 308   }
309   309  
310   std::coroutine_handle<> 310   std::coroutine_handle<>
HITCBC 311   145 await_suspend(std::coroutine_handle<> h, capy::io_env const*) noexcept 311   145 await_suspend(std::coroutine_handle<> h, capy::io_env const*) noexcept
312   { 312   {
313   // Symmetric transfer to server's executor 313   // Symmetric transfer to server's executor
HITCBC 314   145 cont_.h = h; 314   145 cont_.h = h;
HITCBC 315   145 return self_.ex_.dispatch(cont_); 315   145 return self_.ex_.dispatch(cont_);
316   } 316   }
317   317  
HITCBC 318   145 void await_resume() noexcept 318   145 void await_resume() noexcept
319   { 319   {
320   // Running on server executor - safe to modify lists 320   // Running on server executor - safe to modify lists
321   // Remove from active (if present), then wake waiter or add to idle 321   // Remove from active (if present), then wake waiter or add to idle
HITCBC 322   145 self_.active_remove(&w_); 322   145 self_.active_remove(&w_);
HITCBC 323   145 if (self_.waiters_) 323   145 if (self_.waiters_)
324   { 324   {
HITCBC 325   76 auto* wait = self_.waiters_; 325   76 auto* wait = self_.waiters_;
HITCBC 326   76 self_.waiters_ = wait->next; 326   76 self_.waiters_ = wait->next;
HITCBC 327   76 wait->w = &w_; 327   76 wait->w = &w_;
HITCBC 328   76 wait->cont.h = wait->h; 328   76 wait->cont.h = wait->h;
HITCBC 329   76 self_.ex_.post(wait->cont); 329   76 self_.ex_.post(wait->cont);
330   } 330   }
331   else 331   else
332   { 332   {
HITCBC 333   69 self_.idle_push(&w_); 333   69 self_.idle_push(&w_);
334   } 334   }
HITCBC 335   145 } 335   145 }
336   }; 336   };
337   337  
338   class pop_awaitable 338   class pop_awaitable
339   { 339   {
340   tcp_server& self_; 340   tcp_server& self_;
341   waiter wait_; 341   waiter wait_;
342   342  
343   public: 343   public:
HITCBC 344   153 pop_awaitable(tcp_server& self) noexcept : self_(self), wait_{} {} 344   153 pop_awaitable(tcp_server& self) noexcept : self_(self), wait_{} {}
345   345  
HITCBC 346   153 bool await_ready() const noexcept 346   153 bool await_ready() const noexcept
347   { 347   {
HITCBC 348   153 return !self_.idle_empty(); 348   153 return !self_.idle_empty();
349   } 349   }
350   350  
351   bool 351   bool
HITCBC 352   78 await_suspend(std::coroutine_handle<> h, capy::io_env const*) noexcept 352   78 await_suspend(std::coroutine_handle<> h, capy::io_env const*) noexcept
353   { 353   {
354   // Running on server executor (do_accept runs there) 354   // Running on server executor (do_accept runs there)
HITCBC 355   78 wait_.h = h; 355   78 wait_.h = h;
HITCBC 356   78 wait_.w = nullptr; 356   78 wait_.w = nullptr;
HITCBC 357   78 wait_.next = self_.waiters_; 357   78 wait_.next = self_.waiters_;
HITCBC 358   78 self_.waiters_ = &wait_; 358   78 self_.waiters_ = &wait_;
HITCBC 359   78 return true; 359   78 return true;
360   } 360   }
361   361  
HITCBC 362   153 worker_base& await_resume() noexcept 362   153 worker_base& await_resume() noexcept
363   { 363   {
364   // Running on server executor 364   // Running on server executor
HITCBC 365   153 if (wait_.w) 365   153 if (wait_.w)
HITCBC 366   78 return *wait_.w; // Woken by push_awaitable 366   78 return *wait_.w; // Woken by push_awaitable
HITCBC 367   75 return *self_.idle_pop(); 367   75 return *self_.idle_pop();
368   } 368   }
369   }; 369   };
370   370  
HITCBC 371   145 push_awaitable push(worker_base& w) 371   145 push_awaitable push(worker_base& w)
372   { 372   {
HITCBC 373   145 return push_awaitable{*this, w}; 373   145 return push_awaitable{*this, w};
374   } 374   }
375   375  
376   // Synchronous version for destructor/guard paths 376   // Synchronous version for destructor/guard paths
377   // Must be called from server executor context 377   // Must be called from server executor context
HITCBC 378   8 void push_sync(worker_base& w) noexcept 378   8 void push_sync(worker_base& w) noexcept
379   { 379   {
HITCBC 380   8 active_remove(&w); 380   8 active_remove(&w);
HITCBC 381   8 if (waiters_) 381   8 if (waiters_)
382   { 382   {
HITCBC 383   2 auto* wait = waiters_; 383   2 auto* wait = waiters_;
HITCBC 384   2 waiters_ = wait->next; 384   2 waiters_ = wait->next;
HITCBC 385   2 wait->w = &w; 385   2 wait->w = &w;
HITCBC 386   2 wait->cont.h = wait->h; 386   2 wait->cont.h = wait->h;
HITCBC 387   2 ex_.post(wait->cont); 387   2 ex_.post(wait->cont);
388   } 388   }
389   else 389   else
390   { 390   {
HITCBC 391   6 idle_push(&w); 391   6 idle_push(&w);
392   } 392   }
HITCBC 393   8 } 393   8 }
394   394  
HITCBC 395   153 pop_awaitable pop() 395   153 pop_awaitable pop()
396   { 396   {
HITCBC 397   153 return pop_awaitable{*this}; 397   153 return pop_awaitable{*this};
398   } 398   }
399   399  
400   capy::task<void> do_accept(tcp_acceptor& acc); 400   capy::task<void> do_accept(tcp_acceptor& acc);
401   401  
402   public: 402   public:
403   /** Handles one accepted connection using a socket the derived class owns. 403   /** Handles one accepted connection using a socket the derived class owns.
404   404  
405   Derive from this class to implement custom connection handling. 405   Derive from this class to implement custom connection handling.
406   Each worker owns a socket and is reused across multiple 406   Each worker owns a socket and is reused across multiple
407   connections to avoid per-connection allocation. 407   connections to avoid per-connection allocation.
408   408  
409   @par Thread Safety 409   @par Thread Safety
410   run() and socket() execute on the server's executor. 410   run() and socket() execute on the server's executor.
411   411  
412   @see tcp_server, launcher 412   @see tcp_server, launcher
413   */ 413   */
414   class BOOST_COROSIO_DECL worker_base 414   class BOOST_COROSIO_DECL worker_base
415   { 415   {
416   // Ordered largest to smallest for optimal packing 416   // Ordered largest to smallest for optimal packing
417   std::stop_source stop_; // ~16 bytes 417   std::stop_source stop_; // ~16 bytes
418   worker_base* next_ = nullptr; // 8 bytes - used by idle and active lists 418   worker_base* next_ = nullptr; // 8 bytes - used by idle and active lists
419   worker_base* prev_ = nullptr; // 8 bytes - used only by active list 419   worker_base* prev_ = nullptr; // 8 bytes - used only by active list
420   420  
421   friend class tcp_server; 421   friend class tcp_server;
422   422  
423   public: 423   public:
424   /// Construct a worker. 424   /// Construct a worker.
425   worker_base(); 425   worker_base();
426   426  
427   /// Destroy the worker. 427   /// Destroy the worker.
428   virtual ~worker_base(); 428   virtual ~worker_base();
429   429  
430   /** Handle an accepted connection. 430   /** Handle an accepted connection.
431   431  
432   Called when this worker is dispatched to handle a new 432   Called when this worker is dispatched to handle a new
433   connection. The implementation must invoke the launcher 433   connection. The implementation must invoke the launcher
434   exactly once to start the handling coroutine. 434   exactly once to start the handling coroutine.
435   435  
436   @param launch Handle to start the connection coroutine. 436   @param launch Handle to start the connection coroutine.
437   */ 437   */
438   virtual void run(launcher launch) = 0; 438   virtual void run(launcher launch) = 0;
439   439  
440   /// Return the socket used for connections. 440   /// Return the socket used for connections.
441   virtual corosio::tcp_socket& socket() = 0; 441   virtual corosio::tcp_socket& socket() = 0;
442   }; 442   };
443   443  
444   /** Starts a worker's connection-handling coroutine and returns the 444   /** Starts a worker's connection-handling coroutine and returns the
445   worker to the idle pool automatically. 445   worker to the idle pool automatically.
446   446  
447   Passed to @ref worker_base::run to start the connection-handling 447   Passed to @ref worker_base::run to start the connection-handling
448   coroutine. The launcher ensures the worker returns to the idle 448   coroutine. The launcher ensures the worker returns to the idle
449   pool when the coroutine completes or if starting fails. 449   pool when the coroutine completes or if starting fails.
450   450  
451   The launcher must be invoked exactly once via `operator()`. 451   The launcher must be invoked exactly once via `operator()`.
452   If destroyed without invoking, the worker is returned to the 452   If destroyed without invoking, the worker is returned to the
453   idle pool automatically. 453   idle pool automatically.
454   454  
455   @see worker_base::run 455   @see worker_base::run
456   */ 456   */
457   class BOOST_COROSIO_DECL launcher 457   class BOOST_COROSIO_DECL launcher
458   { 458   {
459   tcp_server* srv_; 459   tcp_server* srv_;
460   worker_base* w_; 460   worker_base* w_;
461   461  
462   friend class tcp_server; 462   friend class tcp_server;
463   463  
HITCBC 464   96 launcher(tcp_server& srv, worker_base& w) noexcept : srv_(&srv), w_(&w) 464   96 launcher(tcp_server& srv, worker_base& w) noexcept : srv_(&srv), w_(&w)
465   { 465   {
HITCBC 466   96 } 466   96 }
467   467  
468   public: 468   public:
469   /// Return the worker to the pool if not started. 469   /// Return the worker to the pool if not started.
HITCBC 470   98 ~launcher() 470   98 ~launcher()
471   { 471   {
HITCBC 472   98 if (w_) 472   98 if (w_)
HITCBC 473   8 srv_->push_sync(*w_); 473   8 srv_->push_sync(*w_);
HITCBC 474   98 } 474   98 }
475   475  
476   /** Move construct, transferring the borrowed worker. 476   /** Move construct, transferring the borrowed worker.
477   477  
478   @param o The launcher to take the worker from. It is left 478   @param o The launcher to take the worker from. It is left
479   holding none, so only one of the two returns it. 479   holding none, so only one of the two returns it.
480   */ 480   */
HITCBC 481   2 launcher(launcher&& o) noexcept 481   2 launcher(launcher&& o) noexcept
HITCBC 482   2 : srv_(o.srv_) 482   2 : srv_(o.srv_)
HITCBC 483   2 , w_(std::exchange(o.w_, nullptr)) 483   2 , w_(std::exchange(o.w_, nullptr))
484   { 484   {
HITCBC 485   2 } 485   2 }
486   /// Copy construction is disabled; a launcher holds a borrowed worker it must return exactly once. 486   /// Copy construction is disabled; a launcher holds a borrowed worker it must return exactly once.
487   launcher(launcher const&) = delete; 487   launcher(launcher const&) = delete;
488   /// Copy assignment is disabled; a launcher holds a borrowed worker it must return exactly once. 488   /// Copy assignment is disabled; a launcher holds a borrowed worker it must return exactly once.
489   launcher& operator=(launcher const&) = delete; 489   launcher& operator=(launcher const&) = delete;
490   /// Move assignment is disabled; a launcher is moved, never reassigned. 490   /// Move assignment is disabled; a launcher is moved, never reassigned.
491   launcher& operator=(launcher&&) = delete; 491   launcher& operator=(launcher&&) = delete;
492   492  
493   /** Start the connection-handling coroutine. 493   /** Start the connection-handling coroutine.
494   494  
495   Starts the given coroutine on the specified executor. When 495   Starts the given coroutine on the specified executor. When
496   the coroutine completes, the worker is automatically returned 496   the coroutine completes, the worker is automatically returned
497   to the idle pool. 497   to the idle pool.
498   498  
499   @tparam Executor Executor type satisfying capy::Executor. 499   @tparam Executor Executor type satisfying capy::Executor.
500   500  
501   @param ex The executor to run the coroutine on. 501   @param ex The executor to run the coroutine on.
502   @param task The coroutine to execute. 502   @param task The coroutine to execute.
503   503  
504   @throws std::logic_error If this launcher was already invoked. 504   @throws std::logic_error If this launcher was already invoked.
505   */ 505   */
506   template<class Executor> 506   template<class Executor>
HITCBC 507   90 void operator()(Executor const& ex, capy::task<void> task) 507   90 void operator()(Executor const& ex, capy::task<void> task)
508   { 508   {
HITCBC 509   90 if (!w_) 509   90 if (!w_)
HITCBC 510   2 detail::throw_logic_error(); // launcher already invoked 510   2 detail::throw_logic_error(); // launcher already invoked
511   511  
HITCBC 512   88 auto* w = std::exchange(w_, nullptr); 512   88 auto* w = std::exchange(w_, nullptr);
513   513  
514   // Worker is being dispatched - add to active list 514   // Worker is being dispatched - add to active list
HITCBC 515   88 srv_->active_push(w); 515   88 srv_->active_push(w);
516   516  
517   // Return worker to pool if coroutine setup throws 517   // Return worker to pool if coroutine setup throws
518   struct guard_t 518   struct guard_t
519   { 519   {
520   tcp_server* srv; 520   tcp_server* srv;
521   worker_base* w; 521   worker_base* w;
HITCBC 522   88 ~guard_t() 522   88 ~guard_t()
523   { 523   {
HITCBC 524   88 if (w) 524   88 if (w)
MISUBC 525   ✗ srv->push_sync(*w); 525   ✗ srv->push_sync(*w);
HITCBC 526   88 } 526   88 }
HITCBC 527   88 } guard{srv_, w}; 527   88 } guard{srv_, w};
528   528  
529   // Reset worker's stop source for this connection 529   // Reset worker's stop source for this connection
HITCBC 530   88 w->stop_ = {}; 530   88 w->stop_ = {};
HITCBC 531   88 auto st = w->stop_.get_token(); 531   88 auto st = w->stop_.get_token();
532   532  
HITCBC 533   88 auto wrapper = 533   88 auto wrapper =
HITCBC 534   88 launch_coro<Executor>{}(ex, st, srv_, std::move(task), w); 534   88 launch_coro<Executor>{}(ex, st, srv_, std::move(task), w);
535   535  
536   // Executor and stop token stored in promise via constructor 536   // Executor and stop token stored in promise via constructor
HITCBC 537   88 ex.post(std::exchange(wrapper.h, nullptr)); // Release before post 537   88 ex.post(std::exchange(wrapper.h, nullptr)); // Release before post
HITCBC 538   88 guard.w = nullptr; // Success - dismiss guard 538   88 guard.w = nullptr; // Success - dismiss guard
HITCBC 539   88 } 539   88 }
540   }; 540   };
541   541  
542   /** Construct a TCP server. 542   /** Construct a TCP server.
543   543  
544   @tparam Ctx Execution context type satisfying ExecutionContext. 544   @tparam Ctx Execution context type satisfying ExecutionContext.
545   @tparam Ex Executor type satisfying Executor. 545   @tparam Ex Executor type satisfying Executor.
546   546  
547   @param ctx The execution context for socket operations. 547   @param ctx The execution context for socket operations.
548   @param ex The executor for dispatching coroutines. 548   @param ex The executor for dispatching coroutines.
549   549  
550   @par Example 550   @par Example
551   @par !example tcp_server 551   @par !example tcp_server
552   */ 552   */
553   template<capy::ExecutionContext Ctx, capy::Executor Ex> 553   template<capy::ExecutionContext Ctx, capy::Executor Ex>
HITCBC 554   73 tcp_server(Ctx& ctx, Ex ex) : impl_(make_impl(ctx)) 554   73 tcp_server(Ctx& ctx, Ex ex) : impl_(make_impl(ctx))
HITCBC 555   73 , ex_(std::move(ex)) 555   73 , ex_(std::move(ex))
556   { 556   {
HITCBC 557   73 } 557   73 }
558   558  
559   public: 559   public:
560   /// Destroy the server, stopping all accept loops. 560   /// Destroy the server, stopping all accept loops.
561   ~tcp_server(); 561   ~tcp_server();
562   562  
563   /// Copy construction is disabled; the server owns its worker storage. 563   /// Copy construction is disabled; the server owns its worker storage.
564   tcp_server(tcp_server const&) = delete; 564   tcp_server(tcp_server const&) = delete;
565   /// Copy assignment is disabled; the server owns its worker storage. 565   /// Copy assignment is disabled; the server owns its worker storage.
566   tcp_server& operator=(tcp_server const&) = delete; 566   tcp_server& operator=(tcp_server const&) = delete;
567   567  
568   /** Move construct from another server. 568   /** Move construct from another server.
569   569  
570   @param o The source server. After the move, @p o is 570   @param o The source server. After the move, @p o is
571   in a valid but unspecified state. 571   in a valid but unspecified state.
572   */ 572   */
573   tcp_server(tcp_server&& o) noexcept; 573   tcp_server(tcp_server&& o) noexcept;
574   574  
575   /** Move assign from another server. 575   /** Move assign from another server.
576   576  
577   @param o The source server. After the move, @p o is 577   @param o The source server. After the move, @p o is
578   in a valid but unspecified state. 578   in a valid but unspecified state.
579   579  
580   @return `*this`. 580   @return `*this`.
581   */ 581   */
582   tcp_server& operator=(tcp_server&& o) noexcept; 582   tcp_server& operator=(tcp_server&& o) noexcept;
583   583  
584   /** Bind to a local endpoint. 584   /** Bind to a local endpoint.
585   585  
586   Creates an acceptor listening on the specified endpoint. 586   Creates an acceptor listening on the specified endpoint.
587   Multiple endpoints can be bound by calling this method 587   Multiple endpoints can be bound by calling this method
588   multiple times before @ref start. 588   multiple times before @ref start.
589   589  
590   @param ep The local endpoint to bind to. 590   @param ep The local endpoint to bind to.
591   591  
592   @return An error code indicating success, or the reason binding 592   @return An error code indicating success, or the reason binding
593   failed. 593   failed.
594   */ 594   */
595   [[nodiscard]] std::error_code bind(endpoint ep); 595   [[nodiscard]] std::error_code bind(endpoint ep);
596   596  
597   /** Set the worker pool. 597   /** Set the worker pool.
598   598  
599   Replaces any existing workers with the given range. Any 599   Replaces any existing workers with the given range. Any
600   previous workers are released and the idle/active lists 600   previous workers are released and the idle/active lists
601   are cleared before populating with new workers. 601   are cleared before populating with new workers.
602   602  
603   @tparam Range Forward range of pointer-like objects to worker_base. 603   @tparam Range Forward range of pointer-like objects to worker_base.
604   604  
605   @param workers Range of workers to manage. Each element must 605   @param workers Range of workers to manage. Each element must
606   support `std::to_address()` yielding `worker_base*`. 606   support `std::to_address()` yielding `worker_base*`.
607   607  
608   @par Example 608   @par Example
609   @par !example set_workers 609   @par !example set_workers
610   */ 610   */
611   template<std::ranges::forward_range Range> 611   template<std::ranges::forward_range Range>
612   requires std::convertible_to< 612   requires std::convertible_to<
613   decltype(std::to_address( 613   decltype(std::to_address(
614   std::declval<std::ranges::range_value_t<Range>&>())), 614   std::declval<std::ranges::range_value_t<Range>&>())),
615   worker_base*> 615   worker_base*>
HITCBC 616   73 void set_workers(Range&& workers) 616   73 void set_workers(Range&& workers)
617   { 617   {
618   // Clear existing state 618   // Clear existing state
HITCBC 619   73 storage_.reset(); 619   73 storage_.reset();
HITCBC 620   73 idle_head_ = nullptr; 620   73 idle_head_ = nullptr;
HITCBC 621   73 active_head_ = nullptr; 621   73 active_head_ = nullptr;
HITCBC 622   73 active_tail_ = nullptr; 622   73 active_tail_ = nullptr;
623   623  
624   // Take ownership and populate idle list 624   // Take ownership and populate idle list
625   using StorageType = std::decay_t<Range>; 625   using StorageType = std::decay_t<Range>;
HITCBC 626   73 auto* p = new StorageType(std::forward<Range>(workers)); 626   73 auto* p = new StorageType(std::forward<Range>(workers));
HITCBC 627   73 storage_ = std::shared_ptr<void>( 627   73 storage_ = std::shared_ptr<void>(
HITCBC 628   73 p, [](void* ptr) { delete static_cast<StorageType*>(ptr); }); 628   73 p, [](void* ptr) { delete static_cast<StorageType*>(ptr); });
HITCBC 629   236 for (auto&& elem : *static_cast<StorageType*>(p)) 629   236 for (auto&& elem : *static_cast<StorageType*>(p))
HITCBC 630   163 idle_push(std::to_address(elem)); 630   163 idle_push(std::to_address(elem));
HITCBC 631   73 } 631   73 }
632   632  
633   /** Start accepting connections. 633   /** Start accepting connections.
634   634  
635   Starts accept loops for all bound endpoints. Incoming 635   Starts accept loops for all bound endpoints. Incoming
636   connections are dispatched to idle workers from the pool. 636   connections are dispatched to idle workers from the pool.
637   637  
638   Calling `start()` on an already-running server has no effect. 638   Calling `start()` on an already-running server has no effect.
639   639  
640   @pre At least one endpoint bound via @ref bind. 640   @pre At least one endpoint bound via @ref bind.
641   @pre Workers provided via @ref set_workers. 641   @pre Workers provided via @ref set_workers.
642   @pre If restarting, @ref join must have completed first, and the 642   @pre If restarting, @ref join must have completed first, and the
643   `io_context` must be restarted (`ioc.restart()`). 643   `io_context` must be restarted (`ioc.restart()`).
644   644  
645   @par Effects 645   @par Effects
646   Creates one accept coroutine per bound endpoint. Each coroutine 646   Creates one accept coroutine per bound endpoint. Each coroutine
647   runs on the server's executor, waiting for connections and 647   runs on the server's executor, waiting for connections and
648   dispatching them to idle workers. 648   dispatching them to idle workers.
649   649  
650   @par Restart Sequence 650   @par Restart Sequence
651   To restart after stopping, complete the full shutdown cycle: 651   To restart after stopping, complete the full shutdown cycle:
652   @par !example start 652   @par !example start
653   653  
654   @par Thread Safety 654   @par Thread Safety
655   Not thread safe. 655   Not thread safe.
656   656  
657   @throws std::logic_error If a previous session has not been 657   @throws std::logic_error If a previous session has not been
658   joined (accept loops still active). 658   joined (accept loops still active).
659   */ 659   */
660   void start(); 660   void start();
661   661  
662   /** Return the local endpoint for the i-th bound port. 662   /** Return the local endpoint for the i-th bound port.
663   663  
664   @param index Zero-based index into the list of bound ports. 664   @param index Zero-based index into the list of bound ports.
665   665  
666   @return The local endpoint, or a default-constructed endpoint 666   @return The local endpoint, or a default-constructed endpoint
667   if @p index is out of range or the acceptor is not open. 667   if @p index is out of range or the acceptor is not open.
668   */ 668   */
669   endpoint local_endpoint(std::size_t index = 0) const noexcept; 669   endpoint local_endpoint(std::size_t index = 0) const noexcept;
670   670  
671   /** Stop accepting connections. 671   /** Stop accepting connections.
672   672  
673   Requests the accept loops' stop token and requests cancellation 673   Requests the accept loops' stop token and requests cancellation
674   of active workers via their stop tokens. The acceptors are not 674   of active workers via their stop tokens. The acceptors are not
675   closed. A suspended accept completes once more before its loop 675   closed. A suspended accept completes once more before its loop
676   observes the stop token and ends. 676   observes the stop token and ends.
677   677  
678   This function returns immediately; it does not wait for workers 678   This function returns immediately; it does not wait for workers
679   to finish. Pending I/O operations complete asynchronously. 679   to finish. Pending I/O operations complete asynchronously.
680   680  
681   Calling `stop()` on a non-running server has no effect. 681   Calling `stop()` on a non-running server has no effect.
682   682  
683   @par Effects 683   @par Effects
684   - Requests stop on the accept loops' stop token. The acceptors 684   - Requests stop on the accept loops' stop token. The acceptors
685   are not closed; a pending accept completes once more before 685   are not closed; a pending accept completes once more before
686   the accept loop ends. 686   the accept loop ends.
687   - Requests stop on each active worker's stop token. 687   - Requests stop on each active worker's stop token.
688   - Workers observing their stop token should exit promptly. 688   - Workers observing their stop token should exit promptly.
689   689  
690   @par Postconditions 690   @par Postconditions
691   The server accepts no new connections. Active workers continue 691   The server accepts no new connections. Active workers continue
692   until they observe their stop token or complete naturally. 692   until they observe their stop token or complete naturally.
693   693  
694   @par What Happens Next 694   @par What Happens Next
695   After calling `stop()`: 695   After calling `stop()`:
696   1. Let `ioc.run()` return (drains pending completions). 696   1. Let `ioc.run()` return (drains pending completions).
697   2. Call @ref join to wait for accept loops to finish. 697   2. Call @ref join to wait for accept loops to finish.
698   3. Only then is it safe to restart or destroy the server. 698   3. Only then is it safe to restart or destroy the server.
699   699  
700   @par Thread Safety 700   @par Thread Safety
701   Not thread safe. 701   Not thread safe.
702   702  
703   @see join, start 703   @see join, start
704   */ 704   */
705   void stop(); 705   void stop();
706   706  
707   /** Block until all accept loops complete. 707   /** Block until all accept loops complete.
708   708  
709   Blocks the calling thread until all accept coroutines started 709   Blocks the calling thread until all accept coroutines started
710   by @ref start have finished executing. This synchronizes the 710   by @ref start have finished executing. This synchronizes the
711   shutdown sequence, ensuring the server is fully stopped before 711   shutdown sequence, ensuring the server is fully stopped before
712   restarting or destroying it. 712   restarting or destroying it.
713   713  
714   @pre @ref stop was called and `ioc.run()` returned. 714   @pre @ref stop was called and `ioc.run()` returned.
715   715  
716   @par Postconditions 716   @par Postconditions
717   All accept loops have completed. The server is in the stopped 717   All accept loops have completed. The server is in the stopped
718   state and may be restarted via @ref start. 718   state and may be restarted via @ref start.
719   719  
720   @par Example (Correct Usage) 720   @par Example (Correct Usage)
721   @par !example correct_usage 721   @par !example correct_usage
722   722  
723   @par WARNING: Deadlock Scenario 723   @par WARNING: Deadlock Scenario
724   Calling `join()` from inside a worker coroutine deadlocks: 724   Calling `join()` from inside a worker coroutine deadlocks:
725   725  
726   @par !example deadlock_scenarios 726   @par !example deadlock_scenarios
727   727  
728   @par Thread Safety 728   @par Thread Safety
729   May be called from any thread. It deadlocks if called 729   May be called from any thread. It deadlocks if called
730   from within the `io_context` event loop or from a worker coroutine. 730   from within the `io_context` event loop or from a worker coroutine.
731   731  
732   @see stop, start 732   @see stop, start
733   */ 733   */
734   void join(); 734   void join();
735   735  
736   private: 736   private:
737   capy::task<> do_stop(); 737   capy::task<> do_stop();
738   }; 738   };
739   739  
740   #ifdef _MSC_VER 740   #ifdef _MSC_VER
741   #pragma warning(pop) 741   #pragma warning(pop)
742   #endif 742   #endif
743   743  
744   } // namespace boost::corosio 744   } // namespace boost::corosio
745   745  
746   #endif 746   #endif