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