include/boost/corosio/tcp_server.hpp

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