include/boost/corosio/tcp_server.hpp

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