include/boost/corosio/tcp_server.hpp

93.5% Lines (130/0/139) 97.1% List of functions (33/0/34)
tcp_server.hpp
f(x) Functions (34)
Function Calls Lines Blocks
boost::corosio::tcp_server::idle_push(boost::corosio::tcp_server::worker_base*) :183 239x 100.0% 100.0% boost::corosio::tcp_server::idle_pop() :189 54x 100.0% 100.0% boost::corosio::tcp_server::idle_empty() const :197 63x 100.0% 100.0% boost::corosio::tcp_server::active_push(boost::corosio::tcp_server::worker_base*) :203 27x 100.0% 100.0% boost::corosio::tcp_server::active_remove(boost::corosio::tcp_server::worker_base*) :214 63x 100.0% 91.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 27x 100.0% 100.0% boost::corosio::tcp_server::launch_wrapper<boost::corosio::io_context::executor_type>::promise_type::get_return_object() :261 27x 100.0% 100.0% boost::corosio::tcp_server::launch_wrapper<boost::corosio::io_context::executor_type>::promise_type::initial_suspend() :266 27x 100.0% 100.0% boost::corosio::tcp_server::launch_wrapper<boost::corosio::io_context::executor_type>::promise_type::final_suspend() :270 27x 100.0% 100.0% boost::corosio::tcp_server::launch_wrapper<boost::corosio::io_context::executor_type>::promise_type::return_void() :274 27x 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>&&) :282 27x 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&&) :282 27x 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>) :310 27x 100.0% 100.0% boost::corosio::tcp_server::launch_wrapper<boost::corosio::io_context::executor_type>::~launch_wrapper() :315 27x 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*) :335 27x 100.0% 46.0% boost::corosio::tcp_server::push_awaitable::push_awaitable(boost::corosio::tcp_server&, boost::corosio::tcp_server::worker_base&) :355 57x 100.0% 100.0% boost::corosio::tcp_server::push_awaitable::await_ready() const :361 57x 100.0% 100.0% boost::corosio::tcp_server::push_awaitable::await_suspend(std::__n4861::coroutine_handle<void>, boost::capy::io_env const*) :367 57x 100.0% 100.0% boost::corosio::tcp_server::push_awaitable::await_resume() :374 57x 100.0% 100.0% boost::corosio::tcp_server::pop_awaitable::pop_awaitable(boost::corosio::tcp_server&) :400 63x 100.0% 100.0% boost::corosio::tcp_server::pop_awaitable::await_ready() const :402 63x 100.0% 100.0% boost::corosio::tcp_server::pop_awaitable::await_suspend(std::__n4861::coroutine_handle<void>, boost::capy::io_env const*) :408 9x 100.0% 100.0% boost::corosio::tcp_server::pop_awaitable::await_resume() :418 63x 100.0% 100.0% boost::corosio::tcp_server::push(boost::corosio::tcp_server::worker_base&) :427 57x 100.0% 100.0% boost::corosio::tcp_server::push_sync(boost::corosio::tcp_server::worker_base&) :434 6x 50.0% 80.0% boost::corosio::tcp_server::pop() :451 63x 100.0% 100.0% boost::corosio::tcp_server::launcher::launcher(boost::corosio::tcp_server&, boost::corosio::tcp_server::worker_base&) :516 33x 100.0% 100.0% boost::corosio::tcp_server::launcher::~launcher() :522 33x 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>) :549 30x 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() :564 27x 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) :602 53x 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> > > >&&) :666 53x 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 :678 53x 100.0% 100.0%
Line TLA Hits Source Code
1 //
2 // Copyright (c) 2026 Vinnie Falco ([email protected])
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 239x void idle_push(worker_base* w) noexcept
184 {
185 239x w->next_ = idle_head_;
186 239x idle_head_ = w;
187 239x }
188
189 54x worker_base* idle_pop() noexcept
190 {
191 54x auto* w = idle_head_;
192 54x if (w)
193 54x idle_head_ = w->next_;
194 54x return w;
195 }
196
197 63x bool idle_empty() const noexcept
198 {
199 63x return idle_head_ == nullptr;
200 }
201
202 // Active list (doubly linked) - push back, remove anywhere
203 27x void active_push(worker_base* w) noexcept
204 {
205 27x w->next_ = nullptr;
206 27x w->prev_ = active_tail_;
207 27x if (active_tail_)
208 6x active_tail_->next_ = w;
209 else
210 21x active_head_ = w;
211 27x active_tail_ = w;
212 27x }
213
214 63x void active_remove(worker_base* w) noexcept
215 {
216 // Skip if not in active list (e.g., after failed accept)
217 63x if (w != active_head_ && w->prev_ == nullptr)
218 36x return;
219 27x if (w->prev_)
220 6x w->prev_->next_ = w->next_;
221 else
222 21x active_head_ = w->next_;
223 27x if (w->next_)
224 3x w->next_->prev_ = w->prev_;
225 else
226 24x active_tail_ = w->prev_;
227 27x 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 27x promise_type(Closure&&, E e, S s, Args&&...)
254 27x : ex(std::move(e))
255 27x , env_{
256 27x capy::executor_ref(ex), std::move(s),
257 27x capy::get_current_frame_allocator()}
258 {
259 27x }
260
261 27x launch_wrapper get_return_object() noexcept
262 {
263 return {
264 27x std::coroutine_handle<promise_type>::from_promise(*this)};
265 }
266 27x std::suspend_always initial_suspend() noexcept
267 {
268 27x return {};
269 }
270 27x std::suspend_never final_suspend() noexcept
271 {
272 27x return {};
273 }
274 27x void return_void() noexcept {}
275 void unhandled_exception()
276 {
277 std::terminate();
278 }
279
280 // Inject io_env for IoAwaitable
281 template<capy::IoAwaitable Awaitable>
282 54x auto await_transform(Awaitable&& a)
283 {
284 using AwaitableT = std::decay_t<Awaitable>;
285 struct adapter
286 {
287 AwaitableT aw;
288 capy::io_env const* env;
289
290 bool await_ready()
291 {
292 return aw.await_ready();
293 }
294 decltype(auto) await_resume()
295 {
296 return aw.await_resume();
297 }
298
299 auto await_suspend(std::coroutine_handle<promise_type> h)
300 {
301 return aw.await_suspend(h, env);
302 }
303 };
304 81x return adapter{std::forward<Awaitable>(a), &env_};
305 27x }
306 };
307
308 std::coroutine_handle<promise_type> h;
309
310 27x launch_wrapper(std::coroutine_handle<promise_type> handle) noexcept
311 27x : h(handle)
312 {
313 27x }
314
315 27x ~launch_wrapper()
316 {
317 27x if (h)
318 h.destroy();
319 27x }
320
321 launch_wrapper(launch_wrapper&& o) noexcept
322 : h(std::exchange(o.h, nullptr))
323 {
324 }
325
326 launch_wrapper(launch_wrapper const&) = delete;
327 launch_wrapper& operator=(launch_wrapper const&) = delete;
328 launch_wrapper& operator=(launch_wrapper&&) = delete;
329 };
330
331 // Named functor to avoid incomplete lambda type in coroutine promise
332 template<class Executor>
333 struct launch_coro
334 {
335 27x launch_wrapper<Executor> operator()(
336 Executor,
337 std::stop_token,
338 tcp_server* self,
339 capy::task<void> t,
340 worker_base* wp)
341 {
342 // Executor and stop token stored in promise via constructor
343 co_await std::move(t);
344 co_await self->push(*wp); // worker goes back to idle list
345 54x }
346 };
347
348 class push_awaitable
349 {
350 tcp_server& self_;
351 worker_base& w_;
352 capy::continuation cont_;
353
354 public:
355 57x push_awaitable(tcp_server& self, worker_base& w) noexcept
356 57x : self_(self)
357 57x , w_(w)
358 {
359 57x }
360
361 57x bool await_ready() const noexcept
362 {
363 57x return false;
364 }
365
366 std::coroutine_handle<>
367 57x await_suspend(std::coroutine_handle<> h, capy::io_env const*) noexcept
368 {
369 // Symmetric transfer to server's executor
370 57x cont_.h = h;
371 57x return self_.ex_.dispatch(cont_);
372 }
373
374 57x void await_resume() noexcept
375 {
376 // Running on server executor - safe to modify lists
377 // Remove from active (if present), then wake waiter or add to idle
378 57x self_.active_remove(&w_);
379 57x if (self_.waiters_)
380 {
381 9x auto* wait = self_.waiters_;
382 9x self_.waiters_ = wait->next;
383 9x wait->w = &w_;
384 9x wait->cont.h = wait->h;
385 9x self_.ex_.post(wait->cont);
386 }
387 else
388 {
389 48x self_.idle_push(&w_);
390 }
391 57x }
392 };
393
394 class pop_awaitable
395 {
396 tcp_server& self_;
397 waiter wait_;
398
399 public:
400 63x pop_awaitable(tcp_server& self) noexcept : self_(self), wait_{} {}
401
402 63x bool await_ready() const noexcept
403 {
404 63x return !self_.idle_empty();
405 }
406
407 bool
408 9x await_suspend(std::coroutine_handle<> h, capy::io_env const*) noexcept
409 {
410 // Running on server executor (do_accept runs there)
411 9x wait_.h = h;
412 9x wait_.w = nullptr;
413 9x wait_.next = self_.waiters_;
414 9x self_.waiters_ = &wait_;
415 9x return true;
416 }
417
418 63x worker_base& await_resume() noexcept
419 {
420 // Running on server executor
421 63x if (wait_.w)
422 9x return *wait_.w; // Woken by push_awaitable
423 54x return *self_.idle_pop();
424 }
425 };
426
427 57x push_awaitable push(worker_base& w)
428 {
429 57x return push_awaitable{*this, w};
430 }
431
432 // Synchronous version for destructor/guard paths
433 // Must be called from server executor context
434 6x void push_sync(worker_base& w) noexcept
435 {
436 6x active_remove(&w);
437 6x if (waiters_)
438 {
439 auto* wait = waiters_;
440 waiters_ = wait->next;
441 wait->w = &w;
442 wait->cont.h = wait->h;
443 ex_.post(wait->cont);
444 }
445 else
446 {
447 6x idle_push(&w);
448 }
449 6x }
450
451 63x pop_awaitable pop()
452 {
453 63x return pop_awaitable{*this};
454 }
455
456 capy::task<void> do_accept(tcp_acceptor& acc);
457
458 public:
459 /** Abstract base class for connection handlers.
460
461 Derive from this class to implement custom connection handling.
462 Each worker owns a socket and is reused across multiple
463 connections to avoid per-connection allocation.
464
465 @see tcp_server, launcher
466 */
467 class BOOST_COROSIO_DECL worker_base
468 {
469 // Ordered largest to smallest for optimal packing
470 std::stop_source stop_; // ~16 bytes
471 worker_base* next_ = nullptr; // 8 bytes - used by idle and active lists
472 worker_base* prev_ = nullptr; // 8 bytes - used only by active list
473
474 friend class tcp_server;
475
476 public:
477 /// Construct a worker.
478 worker_base();
479
480 /// Destroy the worker.
481 virtual ~worker_base();
482
483 /** Handle an accepted connection.
484
485 Called when this worker is dispatched to handle a new
486 connection. The implementation must invoke the launcher
487 exactly once to start the handling coroutine.
488
489 @param launch Handle to launch the connection coroutine.
490 */
491 virtual void run(launcher launch) = 0;
492
493 /// Return the socket used for connections.
494 virtual corosio::tcp_socket& socket() = 0;
495 };
496
497 /** Move-only handle to launch a worker coroutine.
498
499 Passed to @ref worker_base::run to start the connection-handling
500 coroutine. The launcher ensures the worker returns to the idle
501 pool when the coroutine completes or if launching fails.
502
503 The launcher must be invoked exactly once via `operator()`.
504 If destroyed without invoking, the worker is returned to the
505 idle pool automatically.
506
507 @see worker_base::run
508 */
509 class BOOST_COROSIO_DECL launcher
510 {
511 tcp_server* srv_;
512 worker_base* w_;
513
514 friend class tcp_server;
515
516 33x launcher(tcp_server& srv, worker_base& w) noexcept : srv_(&srv), w_(&w)
517 {
518 33x }
519
520 public:
521 /// Return the worker to the pool if not launched.
522 33x ~launcher()
523 {
524 33x if (w_)
525 6x srv_->push_sync(*w_);
526 33x }
527
528 launcher(launcher&& o) noexcept
529 : srv_(o.srv_)
530 , w_(std::exchange(o.w_, nullptr))
531 {
532 }
533 launcher(launcher const&) = delete;
534 launcher& operator=(launcher const&) = delete;
535 launcher& operator=(launcher&&) = delete;
536
537 /** Launch the connection-handling coroutine.
538
539 Starts the given coroutine on the specified executor. When
540 the coroutine completes, the worker is automatically returned
541 to the idle pool.
542
543 @param ex The executor to run the coroutine on.
544 @param task The coroutine to execute.
545
546 @throws std::logic_error If this launcher was already invoked.
547 */
548 template<class Executor>
549 30x void operator()(Executor const& ex, capy::task<void> task)
550 {
551 30x if (!w_)
552 3x detail::throw_logic_error(); // launcher already invoked
553
554 27x auto* w = std::exchange(w_, nullptr);
555
556 // Worker is being dispatched - add to active list
557 27x srv_->active_push(w);
558
559 // Return worker to pool if coroutine setup throws
560 struct guard_t
561 {
562 tcp_server* srv;
563 worker_base* w;
564 27x ~guard_t()
565 {
566 27x if (w)
567 srv->push_sync(*w);
568 27x }
569 27x } guard{srv_, w};
570
571 // Reset worker's stop source for this connection
572 27x w->stop_ = {};
573 27x auto st = w->stop_.get_token();
574
575 27x auto wrapper =
576 27x launch_coro<Executor>{}(ex, st, srv_, std::move(task), w);
577
578 // Executor and stop token stored in promise via constructor
579 27x ex.post(std::exchange(wrapper.h, nullptr)); // Release before post
580 27x guard.w = nullptr; // Success - dismiss guard
581 27x }
582 };
583
584 /** Construct a TCP server.
585
586 @tparam Ctx Execution context type satisfying ExecutionContext.
587 @tparam Ex Executor type satisfying Executor.
588
589 @param ctx The execution context for socket operations.
590 @param ex The executor for dispatching coroutines.
591
592 @par Example
593 @code
594 tcp_server srv(ctx, ctx.get_executor());
595 srv.set_workers(make_workers(ctx, 100));
596 if (auto ec = srv.bind(endpoint{...}))
597 return;
598 srv.start();
599 @endcode
600 */
601 template<capy::ExecutionContext Ctx, capy::Executor Ex>
602 53x tcp_server(Ctx& ctx, Ex ex) : impl_(make_impl(ctx))
603 53x , ex_(std::move(ex))
604 {
605 53x }
606
607 public:
608 /// Destroy the server, stopping all accept loops.
609 ~tcp_server();
610
611 tcp_server(tcp_server const&) = delete;
612 tcp_server& operator=(tcp_server const&) = delete;
613
614 /** Move construct from another server.
615
616 @param o The source server. After the move, @p o is
617 in a valid but unspecified state.
618 */
619 tcp_server(tcp_server&& o) noexcept;
620
621 /** Move assign from another server.
622
623 @param o The source server. After the move, @p o is
624 in a valid but unspecified state.
625
626 @return `*this`.
627 */
628 tcp_server& operator=(tcp_server&& o) noexcept;
629
630 /** Bind to a local endpoint.
631
632 Creates an acceptor listening on the specified endpoint.
633 Multiple endpoints can be bound by calling this method
634 multiple times before @ref start.
635
636 @param ep The local endpoint to bind to.
637
638 @return The error code if binding fails.
639 */
640 [[nodiscard]] std::error_code bind(endpoint ep);
641
642 /** Set the worker pool.
643
644 Replaces any existing workers with the given range. Any
645 previous workers are released and the idle/active lists
646 are cleared before populating with new workers.
647
648 @tparam Range Forward range of pointer-like objects to worker_base.
649
650 @param workers Range of workers to manage. Each element must
651 support `std::to_address()` yielding `worker_base*`.
652
653 @par Example
654 @code
655 std::vector<std::unique_ptr<my_worker>> workers;
656 for(int i = 0; i < 100; ++i)
657 workers.push_back(std::make_unique<my_worker>(ctx));
658 srv.set_workers(std::move(workers));
659 @endcode
660 */
661 template<std::ranges::forward_range Range>
662 requires std::convertible_to<
663 decltype(std::to_address(
664 std::declval<std::ranges::range_value_t<Range>&>())),
665 worker_base*>
666 53x void set_workers(Range&& workers)
667 {
668 // Clear existing state
669 53x storage_.reset();
670 53x idle_head_ = nullptr;
671 53x active_head_ = nullptr;
672 53x active_tail_ = nullptr;
673
674 // Take ownership and populate idle list
675 using StorageType = std::decay_t<Range>;
676 53x auto* p = new StorageType(std::forward<Range>(workers));
677 53x storage_ = std::shared_ptr<void>(
678 53x p, [](void* ptr) { delete static_cast<StorageType*>(ptr); });
679 238x for (auto&& elem : *static_cast<StorageType*>(p))
680 185x idle_push(std::to_address(elem));
681 53x }
682
683 /** Start accepting connections.
684
685 Launches accept loops for all bound endpoints. Incoming
686 connections are dispatched to idle workers from the pool.
687
688 Calling `start()` on an already-running server has no effect.
689
690 @par Preconditions
691 - At least one endpoint bound via @ref bind.
692 - Workers provided via @ref set_workers.
693 - If restarting, @ref join must have completed first.
694
695 @par Effects
696 Creates one accept coroutine per bound endpoint. Each coroutine
697 runs on the server's executor, waiting for connections and
698 dispatching them to idle workers.
699
700 @par Restart Sequence
701 To restart after stopping, complete the full shutdown cycle:
702 @code
703 srv.start();
704 ioc.run_for( 1s );
705 srv.stop(); // 1. Signal shutdown
706 ioc.run(); // 2. Drain remaining completions
707 srv.join(); // 3. Wait for accept loops
708
709 // Now safe to restart
710 srv.start();
711 ioc.run();
712 @endcode
713
714 @par Thread Safety
715 Not thread safe.
716
717 @throws std::logic_error If a previous session has not been
718 joined (accept loops still active).
719 */
720 void start();
721
722 /** Return the local endpoint for the i-th bound port.
723
724 @param index Zero-based index into the list of bound ports.
725
726 @return The local endpoint, or a default-constructed endpoint
727 if @p index is out of range or the acceptor is not open.
728 */
729 endpoint local_endpoint(std::size_t index = 0) const noexcept;
730
731 /** Stop accepting connections.
732
733 Requests the accept loops' stop token and requests cancellation
734 of active workers via their stop tokens. The acceptors are not
735 closed; a suspended accept completes once more before its loop
736 observes the stop token and ends.
737
738 This function returns immediately; it does not wait for workers
739 to finish. Pending I/O operations complete asynchronously.
740
741 Calling `stop()` on a non-running server has no effect.
742
743 @par Effects
744 - Requests stop on the accept loops' stop token. The acceptors
745 are not closed; a pending accept completes once more before
746 the accept loop ends.
747 - Requests stop on each active worker's stop token.
748 - Workers observing their stop token should exit promptly.
749
750 @par Postconditions
751 No new connections will be accepted. Active workers continue
752 until they observe their stop token or complete naturally.
753
754 @par What Happens Next
755 After calling `stop()`:
756 1. Let `ioc.run()` return (drains pending completions).
757 2. Call @ref join to wait for accept loops to finish.
758 3. Only then is it safe to restart or destroy the server.
759
760 @par Thread Safety
761 Not thread safe.
762
763 @see join, start
764 */
765 void stop();
766
767 /** Block until all accept loops complete.
768
769 Blocks the calling thread until all accept coroutines launched
770 by @ref start have finished executing. This synchronizes the
771 shutdown sequence, ensuring the server is fully stopped before
772 restarting or destroying it.
773
774 @par Preconditions
775 @ref stop has been called and `ioc.run()` has returned.
776
777 @par Postconditions
778 All accept loops have completed. The server is in the stopped
779 state and may be restarted via @ref start.
780
781 @par Example (Correct Usage)
782 @code
783 // main thread
784 srv.start();
785 ioc.run(); // Blocks until work completes
786 srv.join(); // Safe: called after ioc.run() returns
787 @endcode
788
789 @par WARNING: Deadlock Scenarios
790 Calling `join()` from the wrong context causes deadlock:
791
792 @code
793 // WRONG: calling join() from inside a worker coroutine
794 void run( launcher launch ) override
795 {
796 launch( ex, [this]() -> capy::task<>
797 {
798 srv_.join(); // DEADLOCK: blocks the executor
799 co_return;
800 }());
801 }
802
803 // WRONG: calling join() while ioc.run() is still active
804 std::thread t( [&]{ ioc.run(); } );
805 srv.stop();
806 srv.join(); // DEADLOCK: ioc.run() still running in thread t
807 @endcode
808
809 @par Thread Safety
810 May be called from any thread, but will deadlock if called
811 from within the io_context event loop or from a worker coroutine.
812
813 @see stop, start
814 */
815 void join();
816
817 private:
818 capy::task<> do_stop();
819 };
820
821 #ifdef _MSC_VER
822 #pragma warning(pop)
823 #endif
824
825 } // namespace boost::corosio
826
827 #endif
828