include/boost/corosio/tcp_server.hpp

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