include/boost/corosio/tcp_server.hpp

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