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 162x 100.0% 100.0% boost::corosio::tcp_server::idle_pop() :189 36x 100.0% 100.0% boost::corosio::tcp_server::idle_empty() const :197 42x 100.0% 100.0% boost::corosio::tcp_server::active_push(boost::corosio::tcp_server::worker_base*) :203 18x 100.0% 100.0% boost::corosio::tcp_server::active_remove(boost::corosio::tcp_server::worker_base*) :214 42x 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 18x 100.0% 100.0% boost::corosio::tcp_server::launch_wrapper<boost::corosio::io_context::executor_type>::promise_type::get_return_object() :261 18x 100.0% 100.0% boost::corosio::tcp_server::launch_wrapper<boost::corosio::io_context::executor_type>::promise_type::initial_suspend() :266 18x 100.0% 100.0% boost::corosio::tcp_server::launch_wrapper<boost::corosio::io_context::executor_type>::promise_type::final_suspend() :270 18x 100.0% 100.0% boost::corosio::tcp_server::launch_wrapper<boost::corosio::io_context::executor_type>::promise_type::return_void() :274 18x 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 18x 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 18x 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 18x 100.0% 100.0% boost::corosio::tcp_server::launch_wrapper<boost::corosio::io_context::executor_type>::~launch_wrapper() :315 18x 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 18x 100.0% 46.0% boost::corosio::tcp_server::push_awaitable::push_awaitable(boost::corosio::tcp_server&, boost::corosio::tcp_server::worker_base&) :355 38x 100.0% 100.0% boost::corosio::tcp_server::push_awaitable::await_ready() const :361 38x 100.0% 100.0% boost::corosio::tcp_server::push_awaitable::await_suspend(std::__n4861::coroutine_handle<void>, boost::capy::io_env const*) :367 38x 100.0% 100.0% boost::corosio::tcp_server::push_awaitable::await_resume() :374 38x 100.0% 100.0% boost::corosio::tcp_server::pop_awaitable::pop_awaitable(boost::corosio::tcp_server&) :400 42x 100.0% 100.0% boost::corosio::tcp_server::pop_awaitable::await_ready() const :402 42x 100.0% 100.0% boost::corosio::tcp_server::pop_awaitable::await_suspend(std::__n4861::coroutine_handle<void>, boost::capy::io_env const*) :408 6x 100.0% 100.0% boost::corosio::tcp_server::pop_awaitable::await_resume() :418 42x 100.0% 100.0% boost::corosio::tcp_server::push(boost::corosio::tcp_server::worker_base&) :427 38x 100.0% 100.0% boost::corosio::tcp_server::push_sync(boost::corosio::tcp_server::worker_base&) :434 4x 50.0% 80.0% boost::corosio::tcp_server::pop() :451 42x 100.0% 100.0% boost::corosio::tcp_server::launcher::launcher(boost::corosio::tcp_server&, boost::corosio::tcp_server::worker_base&) :516 22x 100.0% 100.0% boost::corosio::tcp_server::launcher::~launcher() :522 22x 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 20x 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 18x 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 36x 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 36x 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 36x 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 162x void idle_push(worker_base* w) noexcept
184 {
185 162x w->next_ = idle_head_;
186 162x idle_head_ = w;
187 162x }
188
189 36x worker_base* idle_pop() noexcept
190 {
191 36x auto* w = idle_head_;
192 36x if (w)
193 36x idle_head_ = w->next_;
194 36x return w;
195 }
196
197 42x bool idle_empty() const noexcept
198 {
199 42x return idle_head_ == nullptr;
200 }
201
202 // Active list (doubly linked) - push back, remove anywhere
203 18x void active_push(worker_base* w) noexcept
204 {
205 18x w->next_ = nullptr;
206 18x w->prev_ = active_tail_;
207 18x if (active_tail_)
208 4x active_tail_->next_ = w;
209 else
210 14x active_head_ = w;
211 18x active_tail_ = w;
212 18x }
213
214 42x void active_remove(worker_base* w) noexcept
215 {
216 // Skip if not in active list (e.g., after failed accept)
217 42x if (w != active_head_ && w->prev_ == nullptr)
218 24x return;
219 18x if (w->prev_)
220 4x w->prev_->next_ = w->next_;
221 else
222 14x active_head_ = w->next_;
223 18x if (w->next_)
224 2x w->next_->prev_ = w->prev_;
225 else
226 16x active_tail_ = w->prev_;
227 18x 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 18x promise_type(Closure&&, E e, S s, Args&&...)
254 18x : ex(std::move(e))
255 18x , env_{
256 18x capy::executor_ref(ex), std::move(s),
257 18x capy::get_current_frame_allocator()}
258 {
259 18x }
260
261 18x launch_wrapper get_return_object() noexcept
262 {
263 return {
264 18x std::coroutine_handle<promise_type>::from_promise(*this)};
265 }
266 18x std::suspend_always initial_suspend() noexcept
267 {
268 18x return {};
269 }
270 18x std::suspend_never final_suspend() noexcept
271 {
272 18x return {};
273 }
274 18x 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 36x 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 54x return adapter{std::forward<Awaitable>(a), &env_};
305 18x }
306 };
307
308 std::coroutine_handle<promise_type> h;
309
310 18x launch_wrapper(std::coroutine_handle<promise_type> handle) noexcept
311 18x : h(handle)
312 {
313 18x }
314
315 18x ~launch_wrapper()
316 {
317 18x if (h)
318 h.destroy();
319 18x }
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 18x 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 36x }
346 };
347
348 class push_awaitable
349 {
350 tcp_server& self_;
351 worker_base& w_;
352 capy::continuation cont_;
353
354 public:
355 38x push_awaitable(tcp_server& self, worker_base& w) noexcept
356 38x : self_(self)
357 38x , w_(w)
358 {
359 38x }
360
361 38x bool await_ready() const noexcept
362 {
363 38x return false;
364 }
365
366 std::coroutine_handle<>
367 38x await_suspend(std::coroutine_handle<> h, capy::io_env const*) noexcept
368 {
369 // Symmetric transfer to server's executor
370 38x cont_.h = h;
371 38x return self_.ex_.dispatch(cont_);
372 }
373
374 38x 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 38x self_.active_remove(&w_);
379 38x if (self_.waiters_)
380 {
381 6x auto* wait = self_.waiters_;
382 6x self_.waiters_ = wait->next;
383 6x wait->w = &w_;
384 6x wait->cont.h = wait->h;
385 6x self_.ex_.post(wait->cont);
386 }
387 else
388 {
389 32x self_.idle_push(&w_);
390 }
391 38x }
392 };
393
394 class pop_awaitable
395 {
396 tcp_server& self_;
397 waiter wait_;
398
399 public:
400 42x pop_awaitable(tcp_server& self) noexcept : self_(self), wait_{} {}
401
402 42x bool await_ready() const noexcept
403 {
404 42x return !self_.idle_empty();
405 }
406
407 bool
408 6x await_suspend(std::coroutine_handle<> h, capy::io_env const*) noexcept
409 {
410 // Running on server executor (do_accept runs there)
411 6x wait_.h = h;
412 6x wait_.w = nullptr;
413 6x wait_.next = self_.waiters_;
414 6x self_.waiters_ = &wait_;
415 6x return true;
416 }
417
418 42x worker_base& await_resume() noexcept
419 {
420 // Running on server executor
421 42x if (wait_.w)
422 6x return *wait_.w; // Woken by push_awaitable
423 36x return *self_.idle_pop();
424 }
425 };
426
427 38x push_awaitable push(worker_base& w)
428 {
429 38x return push_awaitable{*this, w};
430 }
431
432 // Synchronous version for destructor/guard paths
433 // Must be called from server executor context
434 4x void push_sync(worker_base& w) noexcept
435 {
436 4x active_remove(&w);
437 4x 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 4x idle_push(&w);
448 }
449 4x }
450
451 42x pop_awaitable pop()
452 {
453 42x 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 22x launcher(tcp_server& srv, worker_base& w) noexcept : srv_(&srv), w_(&w)
517 {
518 22x }
519
520 public:
521 /// Return the worker to the pool if not launched.
522 22x ~launcher()
523 {
524 22x if (w_)
525 4x srv_->push_sync(*w_);
526 22x }
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 20x void operator()(Executor const& ex, capy::task<void> task)
550 {
551 20x if (!w_)
552 2x detail::throw_logic_error(); // launcher already invoked
553
554 18x auto* w = std::exchange(w_, nullptr);
555
556 // Worker is being dispatched - add to active list
557 18x 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 18x ~guard_t()
565 {
566 18x if (w)
567 srv->push_sync(*w);
568 18x }
569 18x } guard{srv_, w};
570
571 // Reset worker's stop source for this connection
572 18x w->stop_ = {};
573 18x auto st = w->stop_.get_token();
574
575 18x auto wrapper =
576 18x launch_coro<Executor>{}(ex, st, srv_, std::move(task), w);
577
578 // Executor and stop token stored in promise via constructor
579 18x ex.post(std::exchange(wrapper.h, nullptr)); // Release before post
580 18x guard.w = nullptr; // Success - dismiss guard
581 18x }
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 36x tcp_server(Ctx& ctx, Ex ex) : impl_(make_impl(ctx))
603 36x , ex_(std::move(ex))
604 {
605 36x }
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 36x void set_workers(Range&& workers)
667 {
668 // Clear existing state
669 36x storage_.reset();
670 36x idle_head_ = nullptr;
671 36x active_head_ = nullptr;
672 36x active_tail_ = nullptr;
673
674 // Take ownership and populate idle list
675 using StorageType = std::decay_t<Range>;
676 36x auto* p = new StorageType(std::forward<Range>(workers));
677 36x storage_ = std::shared_ptr<void>(
678 36x p, [](void* ptr) { delete static_cast<StorageType*>(ptr); });
679 162x for (auto&& elem : *static_cast<StorageType*>(p))
680 126x idle_push(std::to_address(elem));
681 36x }
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