include/boost/corosio/native/detail/iocp/win_scheduler.hpp

94.6% Lines (295/317) 100.0% List of functions (33/33) 73.9% Branches (161/218)
win_scheduler.hpp
f(x) Functions (33)
Function Calls Lines Branches Blocks
boost::corosio::detail::win_scheduler::native_handle() const :77 1618x 100.0% – 100.0% boost::corosio::detail::win_scheduler::configure_threading(boost::corosio::detail::scheduler::threading_config) :87 4428x 100.0% – 100.0% boost::corosio::detail::win_scheduler::scheduler_locking_disabled() const :94 2824x 100.0% – 100.0% boost::corosio::detail::iocp::thread_context_guard::thread_context_guard(boost::corosio::detail::win_scheduler const*) :197 8061x 100.0% – 100.0% boost::corosio::detail::iocp::thread_context_guard::~thread_context_guard() :203 8061x 100.0% – 100.0% boost::corosio::detail::win_scheduler::post(std::__n4861::coroutine_handle<void>) const :218 1845x 100.0% 100.0% 100.0% boost::corosio::detail::win_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::do_complete(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int) :224 1845x 100.0% 70.0% 86.7% boost::corosio::detail::win_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::post_handler(std::__n4861::coroutine_handle<void>) :259 1845x 100.0% – 100.0% boost::corosio::detail::win_scheduler::post(boost::corosio::detail::scheduler_op*) const :279 2753x 100.0% 100.0% 100.0% boost::corosio::detail::win_scheduler::post(boost::capy::continuation&) const :293 21752x 100.0% 100.0% 100.0% boost::corosio::detail::win_scheduler::running_in_this_thread() const :309 21421x 100.0% 75.0% 85.7% boost::corosio::detail::win_scheduler::work_started() :318 451197x 100.0% – 100.0% boost::corosio::detail::win_scheduler::work_finished() :324 477447x 100.0% 100.0% 100.0% boost::corosio::detail::win_scheduler::on_pending(boost::corosio::detail::overlapped_op*) const :331 423575x 41.7% 14.3% 37.9% boost::corosio::detail::win_scheduler::on_completion(boost::corosio::detail::overlapped_op*, unsigned long, unsigned long) const :367 3422x 100.0% 75.0% 71.4% boost::corosio::detail::win_scheduler::stop() :390 15416x 100.0% 85.7% 100.0% boost::corosio::detail::win_scheduler::stopped() const :410 4429x 100.0% – 100.0% boost::corosio::detail::win_scheduler::restart() :417 3894x 100.0% – 100.0% boost::corosio::detail::win_scheduler::run() :424 7432x 100.0% 90.9% 100.0% boost::corosio::detail::win_scheduler::run_one() :451 66x 100.0% 100.0% 71.4% boost::corosio::detail::win_scheduler::wait_one(long) :464 829x 100.0% 75.0% 84.2% boost::corosio::detail::win_scheduler::poll() :484 19x 100.0% 70.0% 88.9% boost::corosio::detail::win_scheduler::poll_one() :502 8x 100.0% 100.0% 85.7% boost::corosio::detail::win_scheduler::post_deferred_completions(boost::corosio::detail::intrusive_queue<boost::corosio::detail::scheduler_op>&) :515 896x 100.0% 66.7% 100.0% boost::corosio::detail::win_scheduler::do_one(unsigned long) :533 453541x 94.8% 78.4% 84.3% boost::corosio::detail::win_scheduler::on_timer_changed(void*) :667 2459x 100.0% – 100.0% boost::corosio::detail::win_scheduler::set_timer_service(boost::corosio::detail::timer_service*) :673 4581x 100.0% 50.0% 100.0% boost::corosio::detail::win_scheduler::update_timeout() :684 3355x 100.0% 50.0% 90.0% boost::corosio::detail::win_scheduler::win_scheduler(boost::capy::execution_context&, int) :705 4582x 100.0% 91.7% 95.0% boost::corosio::detail::win_scheduler::shutdown() :760 4428x 83.7% 61.1% 83.9% boost::corosio::detail::win_scheduler::~win_scheduler() :860 8856x 100.0% – 100.0% boost::corosio::detail::win_scheduler::wait_reactor() :870 48x 100.0% – 100.0% boost::corosio::detail::win_scheduler::cancel_wait(boost::corosio::detail::overlapped_op*) :876 28718x 100.0% – 100.0%
Line Branch TLA Hits Source Code
1 //
2 // Copyright (c) 2025 Vinnie Falco ([email protected])
3 // Copyright (c) 2026 Steve Gerbino
4 // Copyright (c) 2026 Michael Vandeberg
5 //
6 // Distributed under the Boost Software License, Version 1.0. (See accompanying
7 // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
8 //
9 // Official repository: https://github.com/cppalliance/corosio
10 //
11
12 #ifndef BOOST_COROSIO_NATIVE_DETAIL_IOCP_WIN_SCHEDULER_HPP
13 #define BOOST_COROSIO_NATIVE_DETAIL_IOCP_WIN_SCHEDULER_HPP
14
15 #include <boost/corosio/detail/platform.hpp>
16
17 #if BOOST_COROSIO_HAS_IOCP
18
19 #include <boost/corosio/detail/config.hpp>
20 #include <boost/capy/ex/execution_context.hpp>
21
22 #include <boost/corosio/detail/scheduler.hpp>
23 #include <system_error>
24
25 #include <boost/corosio/detail/scheduler_op.hpp>
26 #include <boost/capy/continuation.hpp>
27 #include <boost/corosio/native/detail/iocp/win_completion_key.hpp>
28 #include <boost/corosio/native/detail/iocp/win_mutex.hpp>
29
30 #include <boost/corosio/native/detail/iocp/win_overlapped_op.hpp>
31 #include <boost/corosio/native/detail/iocp/win_timers.hpp>
32 #include <boost/corosio/detail/timer_service.hpp>
33 #include <boost/corosio/native/detail/make_err.hpp>
34 #include <boost/corosio/detail/except.hpp>
35 #include <boost/corosio/detail/thread_local_ptr.hpp>
36
37 #include <atomic>
38 #include <chrono>
39 #include <cstdint>
40 #include <limits>
41 #include <memory>
42 #include <mutex>
43
44 #include <boost/corosio/native/detail/iocp/win_windows.hpp>
45
46 namespace boost::corosio::detail {
47
48 // Forward declarations
49 struct overlapped_op;
50 class win_timers;
51 class win_wait_reactor;
52
53 class BOOST_COROSIO_DECL win_scheduler final
54 : public scheduler
55 {
56 public:
57
58 win_scheduler(capy::execution_context& ctx, int concurrency_hint = -1);
59 ~win_scheduler();
60 win_scheduler(win_scheduler const&) = delete;
61 win_scheduler& operator=(win_scheduler const&) = delete;
62
63 void shutdown() override;
64 void post(std::coroutine_handle<> h) const override;
65 void post(scheduler_op* h) const override;
66 void post(capy::continuation&) const override;
67 bool running_in_this_thread() const noexcept override;
68 void stop() override;
69 bool stopped() const noexcept override;
70 void restart() override;
71 std::size_t run() override;
72 std::size_t run_one() override;
73 std::size_t wait_one(long usec) override;
74 std::size_t poll() override;
75 std::size_t poll_one() override;
76
77 1618x void* native_handle() const noexcept
78 {
79 1618x return iocp_;
80 }
81
82 void work_started() noexcept override;
83 void work_finished() noexcept override;
84
85 // IOCP has only the dispatch mutex; reactor_io_locking and one_thread do
86 // not apply (the completion port provides its own synchronization).
87 4428x void configure_threading(threading_config cfg) noexcept override
88 {
89 4428x scheduler_locking_disabled_ = !cfg.scheduler_locking;
90 4428x dispatch_mutex_.set_enabled(cfg.scheduler_locking);
91 4428x }
92
93 /// Return true when scheduler locking is disabled (fully-lockless tier).
94 2824x bool scheduler_locking_disabled() const noexcept override
95 {
96 // LCOV_EXCL_START: consulted only by the POSIX pool-backed
97 // services; the IOCP services do not read it yet.
98 − return scheduler_locking_disabled_;
99 // LCOV_EXCL_STOP
100 }
101
102 /** Signal that an overlapped I/O operation is now pending.
103 Coordinates with do_one() via the ready_ CAS protocol. */
104 void on_pending(overlapped_op* op) const;
105
106 /** Post an immediate completion with pre-stored results.
107 Used for sync errors and noop paths. */
108 void on_completion(overlapped_op* op, DWORD error, DWORD bytes) const;
109
110 // Timer service integration
111 void set_timer_service(timer_service* svc);
112 void update_timeout();
113
114 private:
115 static void on_timer_changed(void* ctx);
116 void post_deferred_completions(op_queue& ops);
117 std::size_t do_one(unsigned long timeout_ms);
118
119 timer_service* timer_svc_ = nullptr;
120 void* iocp_;
121 mutable long outstanding_work_;
122
123 // Packets in flight to the completion port that reference
124 // overlapped-op memory: kernel completions owed after a pending
125 // submission, plus successful stored-result posts. Shutdown reaps
126 // until this is zero before the services free op storage. The
127 // run-loop counter cannot serve that role: frames abandoned at
128 // teardown never return their work-guard credits.
129 mutable long pending_io_ = 0;
130 mutable long stopped_;
131 long stop_event_posted_;
132 mutable long dispatch_required_;
133 bool scheduler_locking_disabled_ = false;
134
135 BOOST_COROSIO_MSVC_WARNING_PUSH
136 BOOST_COROSIO_MSVC_WARNING_DISABLE(
137 4251) // std::/detail:: members, dll-interface
138 mutable win_mutex dispatch_mutex_;
139 mutable op_queue completed_ops_;
140 std::unique_ptr<win_timers> timers_;
141 std::unique_ptr<win_wait_reactor> wait_reactor_;
142 BOOST_COROSIO_MSVC_WARNING_POP
143
144 public:
145 /** Return the auxiliary select-based reactor for wait operations.
146
147 Built with the scheduler and stopped+joined in ~win_scheduler,
148 so it is always there to hand a wait to. Used by socket and
149 acceptor wait() implementations whose readiness cannot be
150 expressed natively in IOCP (datagram-read, acceptor-read,
151 error-wait).
152 */
153 win_wait_reactor& wait_reactor();
154
155 /// Cancel a parked wait op. Safe to call from any thread.
156 void cancel_wait(overlapped_op* op) noexcept;
157 };
158
159 /*
160 ARCHITECTURE NOTE: Function Pointer Dispatch
161
162 All I/O handles are registered with the IOCP using key_io (0).
163 Dispatch happens via the function pointer stored in each scheduler_op.
164
165 When GQCS returns with an OVERLAPPED*, we cast it to scheduler_op*
166 and call the function pointer directly - no virtual dispatch.
167
168 The completion_key enum values are used only for internal signals:
169 - key_io (0): Normal I/O completion, dispatch via func_
170 - key_wake_dispatch (1): Timer wakeup, check dispatch_required_
171 - key_shutdown (2): Stop signal
172 - key_result_stored (3): Results pre-stored in OVERLAPPED
173 - key_posted: Carries a scheduler_op* in the OVERLAPPED pointer
174 - key_continuation: Carries a capy::continuation* in the OVERLAPPED pointer
175 */
176
177 namespace iocp {
178
179 // Poll interval (ms) for the one-time shutdown drain loop. Unlike the
180 // steady-state loop (which blocks indefinitely until a real wake-up),
181 // the drain must keep re-checking the work count and the deferred
182 // completion queue to make progress, so it uses a finite wait.
183 inline constexpr unsigned long shutdown_drain_timeout_ms = 500;
184
185 struct BOOST_COROSIO_SYMBOL_VISIBLE scheduler_context
186 {
187 win_scheduler const* key;
188 scheduler_context* next;
189 };
190
191 inline thread_local_ptr<scheduler_context> context_stack;
192
193 struct thread_context_guard
194 {
195 scheduler_context frame_;
196
197 8061x explicit thread_context_guard(win_scheduler const* ctx) noexcept
198 8061x : frame_{ctx, context_stack.get()}
199 {
200 8061x context_stack.set(&frame_);
201 8061x }
202
203 8061x ~thread_context_guard() noexcept
204 {
205 8061x context_stack.set(frame_.next);
206 8061x }
207 };
208
209 } // namespace iocp
210
211 // The constructor, ~win_scheduler(), shutdown(), wait_reactor() and
212 // cancel_wait() are defined at the bottom of this header so the
213 // unique_ptr<win_wait_reactor>'s deleter, its operator* and
214 // wait_reactor_->stop() see the type complete. The constructor needs
215 // it to build the reactor, and its unwind path to destroy it.
216
217 inline void
218 1845x win_scheduler::post(std::coroutine_handle<> h) const
219 {
220 struct post_handler final : scheduler_op
221 {
222 std::coroutine_handle<> h_;
223
224 1845x static void do_complete(
225 void* owner, scheduler_op* base, std::uint32_t, std::uint32_t)
226 {
227 1845x auto* self = static_cast<post_handler*>(base);
228
2/2
✓ Branch 2 → 3 taken 6 times.
✓ Branch 2 → 10 taken 1839 times.
1845x if (!owner)
229 {
230 // Shutdown path: destroy the coroutine frame synchronously.
231 //
232 // Bounded destruction invariant: the chain triggered by
233 // coro.destroy() is at most two levels deep:
234 // 1. task frame destroyed → ~io_awaitable_promise_base()
235 // destroys stored continuation (if != noop_coroutine)
236 // 2. continuation (trampoline) destroyed → final_suspend
237 // returns suspend_never, no further continuation
238 //
239 // If a future refactor adds deeper continuation chains,
240 // this would reintroduce re-entrant stack overflow risk.
241 #ifndef NDEBUG
242 static thread_local int destroy_depth = 0;
243 6x ++destroy_depth;
244
1/2
✗ Branch 3 → 4 not taken.
✓ Branch 3 → 5 taken 6 times.
6x BOOST_COROSIO_ASSERT(destroy_depth <= 2);
245 #endif
246 6x auto coro = self->h_;
247
1/2
✓ Branch 6 → 7 taken 6 times.
✗ Branch 6 → 8 not taken.
6x delete self;
248
1/1
✓ Branch 8 → 9 taken 6 times.
6x coro.destroy();
249 #ifndef NDEBUG
250 6x --destroy_depth;
251 #endif
252 6x return;
253 }
254 1839x auto coro = self->h_;
255
1/2
✓ Branch 10 → 11 taken 1839 times.
✗ Branch 10 → 12 not taken.
1839x delete self;
256
1/1
✓ Branch 12 → 13 taken 1839 times.
1839x coro.resume();
257 }
258
259 1845x explicit post_handler(std::coroutine_handle<> coro)
260 1845x : scheduler_op(&do_complete)
261 1845x , h_(coro)
262 {
263 1845x }
264 };
265
266 1845x auto* ph = new post_handler(h);
267 1845x ::InterlockedIncrement(&outstanding_work_);
268
269
2/2
✓ Branch 7 → 8 taken 1 time.
✓ Branch 7 → 14 taken 1844 times.
1845x if (!::PostQueuedCompletionStatus(
270 1845x iocp_, 0, key_posted, reinterpret_cast<LPOVERLAPPED>(ph)))
271 {
272 1x std::lock_guard<win_mutex> lock(dispatch_mutex_);
273 1x completed_ops_.push(ph);
274 1x ::InterlockedExchange(&dispatch_required_, 1);
275 1x }
276 1845x }
277
278 inline void
279 2753x win_scheduler::post(scheduler_op* h) const
280 {
281 2753x ::InterlockedIncrement(&outstanding_work_);
282
283
2/2
✓ Branch 5 → 6 taken 1 time.
✓ Branch 5 → 12 taken 2752 times.
2753x if (!::PostQueuedCompletionStatus(
284 2753x iocp_, 0, key_posted, reinterpret_cast<LPOVERLAPPED>(h)))
285 {
286 1x std::lock_guard<win_mutex> lock(dispatch_mutex_);
287 1x completed_ops_.push(h);
288 1x ::InterlockedExchange(&dispatch_required_, 1);
289 1x }
290 2753x }
291
292 inline void
293 21752x win_scheduler::post(capy::continuation& c) const
294 {
295 21752x ::InterlockedIncrement(&outstanding_work_);
296
297
2/2
✓ Branch 5 → 6 taken 2 times.
✓ Branch 5 → 9 taken 21750 times.
21752x if (!::PostQueuedCompletionStatus(
298 21752x iocp_, 0, key_continuation, reinterpret_cast<LPOVERLAPPED>(&c)))
299 {
300 // completed_ops_ is an op_queue and cannot carry a raw continuation,
301 // so on the rare PQCS failure fall back to the allocating handle
302 // path. Drop the increment first; post(c.h) does its own accounting.
303 2x ::InterlockedDecrement(&outstanding_work_);
304 2x post(c.h);
305 }
306 21752x }
307
308 inline bool
309 21421x win_scheduler::running_in_this_thread() const noexcept
310 {
311
2/2
✓ Branch 6 → 3 taken 3353 times.
✓ Branch 6 → 7 taken 18068 times.
21421x for (auto* c = iocp::context_stack.get(); c != nullptr; c = c->next)
312
1/2
✓ Branch 3 → 4 taken 3353 times.
✗ Branch 3 → 5 not taken.
3353x if (c->key == this)
313 3353x return true;
314 18068x return false;
315 }
316
317 inline void
318 451197x win_scheduler::work_started() noexcept
319 {
320 451197x ::InterlockedIncrement(&outstanding_work_);
321 451197x }
322
323 inline void
324 477447x win_scheduler::work_finished() noexcept
325 {
326
2/2
✓ Branch 4 → 5 taken 7767 times.
✓ Branch 4 → 6 taken 469680 times.
954894x if (::InterlockedDecrement(&outstanding_work_) == 0)
327 7767x stop();
328 477447x }
329
330 inline void
331 423575x win_scheduler::on_pending(overlapped_op* op) const
332 {
333 // If the CAS fails (ready_ was already 1), the completer got here first
334 // and stored the results — re-post so do_one() can dispatch. The acquire
335 // on failure makes those payload writes visible, so the re-posted op
336 // carries valid dwError / bytes_transferred.
337 //
338 // pending_io_ counts the packet that will dispatch this op: on CAS
339 // success the kernel's own completion will find ready_ == 1 and
340 // dispatch; on CAS failure the kernel's packet was consumed as a
341 // skip (uncounted) and the re-post is the dispatching packet. A
342 // failed re-post falls back to the deferred queue, which holds the
343 // op memory itself — no packet, no count.
344 423575x long expected = 0;
345
1/2
✓ Branch 11 → 12 taken 423575 times.
✗ Branch 11 → 15 not taken.
847150x if (op->ready_.compare_exchange_strong(
346 expected, 1, std::memory_order_acq_rel, std::memory_order_acquire))
347 {
348 423575x ::InterlockedIncrement(&pending_io_);
349 }
350 else
351 {
352 ✗ if (::PostQueuedCompletionStatus(
353 ✗ iocp_, 0, key_result_stored, static_cast<LPOVERLAPPED>(op)))
354 {
355 ✗ ::InterlockedIncrement(&pending_io_);
356 }
357 else
358 {
359 ✗ std::lock_guard<win_mutex> lock(dispatch_mutex_);
360 ✗ completed_ops_.push(op);
361 ✗ ::InterlockedExchange(&dispatch_required_, 1);
362 ✗ }
363 }
364 423575x }
365
366 inline void
367 3422x win_scheduler::on_completion(overlapped_op* op, DWORD error, DWORD bytes) const
368 {
369 // Synchronous-completion path. Write the payload before the release store
370 // to ready_ so the GQCS thread that dequeues the key_result_stored post
371 // sees both fields once it observes ready_ == 1.
372 3422x op->dwError = error;
373 3422x op->bytes_transferred = bytes;
374 3422x op->ready_.store(1, std::memory_order_release);
375
376
3/4
✓ Branch 22 → 23 taken 3422 times.
✗ Branch 22 → 24 not taken.
✓ Branch 26 → 27 taken 3421 times.
✓ Branch 26 → 30 taken 1 time.
6844x if (::PostQueuedCompletionStatus(
377 3422x iocp_, 0, key_result_stored, static_cast<LPOVERLAPPED>(op)))
378 {
379 3421x ::InterlockedIncrement(&pending_io_);
380 }
381 else
382 {
383 1x std::lock_guard<win_mutex> lock(dispatch_mutex_);
384 1x completed_ops_.push(op);
385 1x ::InterlockedExchange(&dispatch_required_, 1);
386 1x }
387 3422x }
388
389 inline void
390 15416x win_scheduler::stop()
391 {
392
2/2
✓ Branch 4 → 5 taken 7807 times.
✓ Branch 4 → 13 taken 7609 times.
30832x if (::InterlockedExchange(&stopped_, 1) == 0)
393 {
394
1/2
✓ Branch 7 → 8 taken 7807 times.
✗ Branch 7 → 13 not taken.
15614x if (::InterlockedExchange(&stop_event_posted_, 1) == 0)
395 {
396
2/2
✓ Branch 9 → 10 taken 1 time.
✓ Branch 9 → 13 taken 7806 times.
7807x if (!::PostQueuedCompletionStatus(iocp_, 0, key_shutdown, nullptr))
397 {
398 // The shutdown post is the only thing that wakes a
399 // run()/run_one() thread blocked indefinitely in GQCS.
400 // With no periodic timeout there is no fallback, so a
401 // failed post is fatal (matches Asio). It can only fail
402 // under resource exhaustion (ERROR_NO_SYSTEM_RESOURCES).
403
1/1
✓ Branch 10 → 11 taken 1 time.
1x detail::throw_system_error(make_err(::GetLastError()));
404 }
405 }
406 }
407 15415x }
408
409 inline bool
410 4429x win_scheduler::stopped() const noexcept
411 {
412 // equivalent to atomic read
413 8858x return ::InterlockedExchangeAdd(&stopped_, 0) != 0;
414 }
415
416 inline void
417 3894x win_scheduler::restart()
418 {
419 3894x ::InterlockedExchange(&stopped_, 0);
420 3894x ::InterlockedExchange(&stop_event_posted_, 0);
421 3894x }
422
423 inline std::size_t
424 7432x win_scheduler::run()
425 {
426
2/2
✓ Branch 4 → 5 taken 48 times.
✓ Branch 4 → 7 taken 7384 times.
14864x if (::InterlockedExchangeAdd(&outstanding_work_, 0) == 0)
427 {
428
1/1
✓ Branch 5 → 6 taken 48 times.
48x stop();
429 48x return 0;
430 }
431
432 7384x iocp::thread_context_guard ctx(this);
433
434 7384x std::size_t n = 0;
435 for (;;)
436 {
437
3/3
✓ Branch 9 → 10 taken 452843 times.
✓ Branch 10 → 11 taken 53 times.
✓ Branch 10 → 12 taken 452790 times.
452844x if (!do_one(INFINITE))
438 53x break;
439
1/2
✓ Branch 13 → 14 taken 452790 times.
✗ Branch 13 → 15 not taken.
452790x if (n != (std::numeric_limits<std::size_t>::max)())
440 452790x ++n;
441
2/2
✓ Branch 17 → 18 taken 7330 times.
✓ Branch 17 → 20 taken 445460 times.
905580x if (::InterlockedExchangeAdd(&outstanding_work_, 0) == 0)
442 {
443
1/1
✓ Branch 18 → 19 taken 7330 times.
7330x stop();
444 7330x break;
445 }
446 }
447 7383x return n;
448 7384x }
449
450 inline std::size_t
451 66x win_scheduler::run_one()
452 {
453
2/2
✓ Branch 4 → 5 taken 1 time.
✓ Branch 4 → 7 taken 65 times.
132x if (::InterlockedExchangeAdd(&outstanding_work_, 0) == 0)
454 {
455
1/1
✓ Branch 5 → 6 taken 1 time.
1x stop();
456 1x return 0;
457 }
458
459 65x iocp::thread_context_guard ctx(this);
460
1/1
✓ Branch 8 → 9 taken 65 times.
65x return do_one(INFINITE);
461 65x }
462
463 inline std::size_t
464 829x win_scheduler::wait_one(long usec)
465 {
466
2/2
✓ Branch 4 → 5 taken 233 times.
✓ Branch 4 → 7 taken 596 times.
1658x if (::InterlockedExchangeAdd(&outstanding_work_, 0) == 0)
467 {
468
1/1
✓ Branch 5 → 6 taken 233 times.
233x stop();
469 233x return 0;
470 }
471
472 596x iocp::thread_context_guard ctx(this);
473 596x unsigned long timeout_ms = INFINITE;
474
1/2
✓ Branch 8 → 9 taken 596 times.
✗ Branch 8 → 13 not taken.
596x if (usec >= 0)
475 {
476 596x auto ms = (static_cast<long long>(usec) + 999) / 1000;
477
1/2
✓ Branch 9 → 10 taken 596 times.
✗ Branch 9 → 11 not taken.
596x timeout_ms = ms >= 0xFFFFFFFELL ? static_cast<unsigned long>(0xFFFFFFFE)
478 : static_cast<unsigned long>(ms);
479 }
480
1/1
✓ Branch 13 → 14 taken 596 times.
596x return do_one(timeout_ms);
481 596x }
482
483 inline std::size_t
484 19x win_scheduler::poll()
485 {
486
2/2
✓ Branch 4 → 5 taken 8 times.
✓ Branch 4 → 7 taken 11 times.
38x if (::InterlockedExchangeAdd(&outstanding_work_, 0) == 0)
487 {
488
1/1
✓ Branch 5 → 6 taken 8 times.
8x stop();
489 8x return 0;
490 }
491
492 11x iocp::thread_context_guard ctx(this);
493
494 11x std::size_t n = 0;
495
3/5
✓ Branch 12 → 13 taken 31 times.
✓ Branch 13 → 9 taken 20 times.
✓ Branch 13 → 14 taken 11 times.
✗ Branch 14 → 9 not taken.
✗ Branch 14 → 15 not taken.
31x while (do_one(0))
496
1/2
✓ Branch 10 → 11 taken 20 times.
✗ Branch 10 → 12 not taken.
20x if (n != (std::numeric_limits<std::size_t>::max)())
497 20x ++n;
498 11x return n;
499 11x }
500
501 inline std::size_t
502 8x win_scheduler::poll_one()
503 {
504
2/2
✓ Branch 4 → 5 taken 3 times.
✓ Branch 4 → 7 taken 5 times.
16x if (::InterlockedExchangeAdd(&outstanding_work_, 0) == 0)
505 {
506
1/1
✓ Branch 5 → 6 taken 3 times.
3x stop();
507 3x return 0;
508 }
509
510 5x iocp::thread_context_guard ctx(this);
511
1/1
✓ Branch 8 → 9 taken 5 times.
5x return do_one(0);
512 5x }
513
514 inline void
515 899x win_scheduler::post_deferred_completions(op_queue& ops)
516 {
517
3/4
✓ Branch 3 → 4 taken 4 times.
✓ Branch 3 → 15 taken 878 times.
✗ Branch 4 → 5 not taken.
✓ Branch 4 → 16 taken 17 times.
899x while (auto h = ops.pop())
518 {
519
3/5
✓ Branch 4 → 5 taken 4 times.
✓ Branch 5 → 6 taken 3 times.
✓ Branch 5 → 7 taken 1 time.
✗ Branch 6 → 7 not taken.
✗ Branch 6 → 8 not taken.
4x if (::PostQueuedCompletionStatus(
520 iocp_, 0, key_posted, reinterpret_cast<LPOVERLAPPED>(h)))
521 3x continue;
522
523 // Out of resources, put the failed op and remaining ops back
524 1x ops.push(h);
525 1x std::lock_guard<win_mutex> lock(dispatch_mutex_);
526 1x completed_ops_.splice(ops);
527 1x ::InterlockedExchange(&dispatch_required_, 1);
528 1x return;
529 4x }
530 }
531
532 inline std::size_t
533 458109x win_scheduler::do_one(unsigned long timeout_ms)
534 {
535 for (;;)
536 {
537 // Check if we need to process timers or deferred ops
538
4/4
✓ Branch 4 → 5 taken 879 times.
✓ Branch 4 → 13 taken 456699 times.
✓ Branch 5 → 6 taken 17 times.
✓ Branch 5 → 14 taken 585 times.
916360x if (::InterlockedCompareExchange(&dispatch_required_, 0, 1) == 1)
539 {
540 896x op_queue local_ops;
541 {
542 896x std::lock_guard<win_mutex> lock(dispatch_mutex_);
543 896x local_ops.splice(completed_ops_);
544 896x }
545
2/2
✓ Branch 8 → 9 taken 879 times.
✓ Branch 9 → 10 taken 17 times.
896x post_deferred_completions(local_ops);
546
547
2/4
✓ Branch 9 → 10 taken 879 times.
✗ Branch 9 → 11 not taken.
✓ Branch 10 → 11 taken 17 times.
✗ Branch 10 → 12 not taken.
896x if (timer_svc_)
548
2/2
✓ Branch 10 → 11 taken 879 times.
✓ Branch 11 → 12 taken 17 times.
896x timer_svc_->process_expired();
549
550
2/2
✓ Branch 11 → 12 taken 879 times.
✓ Branch 12 → 13 taken 17 times.
896x update_timeout();
551 }
552
553 458180x DWORD bytes = 0;
554 458180x ULONG_PTR key = 0;
555 458180x LPOVERLAPPED overlapped = nullptr;
556
2/2
✓ Branch 13 → 14 taken 457578 times.
✓ Branch 14 → 15 taken 602 times.
458180x ::SetLastError(0);
557
558
2/2
✓ Branch 14 → 15 taken 457578 times.
✓ Branch 15 → 16 taken 602 times.
458180x BOOL result = ::GetQueuedCompletionStatus(
559 iocp_, &bytes, &key, &overlapped, timeout_ms);
560
2/2
✓ Branch 15 → 16 taken 457578 times.
✓ Branch 16 → 17 taken 602 times.
458180x DWORD dwError = ::GetLastError();
561
562 // Handle based on completion key
563
4/4
✓ Branch 16 → 17 taken 452761 times.
✓ Branch 16 → 47 taken 4817 times.
✓ Branch 17 → 18 taken 531 times.
✓ Branch 17 → 48 taken 71 times.
458180x if (overlapped)
564 {
565
4/4
✓ Branch 17 → 18 taken 2145 times.
✓ Branch 17 → 19 taken 450616 times.
✓ Branch 18 → 19 taken 12 times.
✓ Branch 18 → 20 taken 519 times.
453292x DWORD err = result ? 0 : dwError;
566
567
6/8
✓ Branch 20 → 21 taken 426713 times.
✓ Branch 20 → 40 taken 4566 times.
✓ Branch 20 → 43 taken 21482 times.
✗ Branch 20 → 46 not taken.
✓ Branch 21 → 22 taken 246 times.
✓ Branch 21 → 41 taken 20 times.
✓ Branch 21 → 44 taken 265 times.
✗ Branch 21 → 47 not taken.
453292x switch (key)
568 {
569 426959x case key_io:
570 case key_result_stored:
571 {
572 426959x auto* ov_op = overlapped_to_op(overlapped);
573
574 // key_io carries fresh kernel results — publish them before
575 // the CAS so that losing the race (old value 0) still leaves
576 // valid data for the on_pending() re-post. For key_result_stored
577 // the payload was already written and released (by the completer
578 // in on_pending, or by on_completion), so we must not overwrite
579 // it.
580
4/4
✓ Branch 22 → 23 taken 423304 times.
✓ Branch 22 → 24 taken 3409 times.
✓ Branch 23 → 24 taken 241 times.
✓ Branch 23 → 25 taken 5 times.
426959x if (key == key_io)
581 423545x ov_op->store_result(bytes, err);
582
583 // If old value was 1 the initiator already returned — dispatch.
584 // The acquire pairs with the publisher's release so the payload
585 // reads in complete() are ordered after the store that produced
586 // them. If old value was 0 the initiator hasn't returned yet;
587 // skip and let on_pending() re-post.
588 426959x long expected = 0;
589
2/4
✓ Branch 33 → 34 taken 426713 times.
✗ Branch 33 → 39 not taken.
✓ Branch 34 → 35 taken 246 times.
✗ Branch 34 → 40 not taken.
853918x if (!ov_op->ready_.compare_exchange_strong(
590 expected, 1, std::memory_order_acq_rel,
591 std::memory_order_acquire))
592 {
593 426959x ::InterlockedDecrement(&pending_io_);
594 426959x ov_op->complete(
595
2/2
✓ Branch 36 → 37 taken 426713 times.
✓ Branch 37 → 38 taken 246 times.
426959x this, ov_op->bytes_transferred, ov_op->dwError);
596 426959x work_finished();
597 426959x return 1;
598 }
599 ✗ continue;
600 ✗ }
601
602 4586x case key_posted:
603 {
604 // Posted scheduler_op*: overlapped is actually a scheduler_op*
605 4586x auto* op = reinterpret_cast<scheduler_op*>(overlapped);
606
2/2
✓ Branch 40 → 41 taken 4566 times.
✓ Branch 41 → 42 taken 20 times.
4586x op->complete(this, bytes, err);
607 4586x work_finished();
608 4586x return 1;
609 }
610
611 21747x case key_continuation:
612 {
613 // Posted continuation: overlapped is actually a continuation*
614 21747x auto* c = reinterpret_cast<capy::continuation*>(overlapped);
615
2/2
✓ Branch 43 → 44 taken 21482 times.
✓ Branch 44 → 45 taken 265 times.
21747x c->h.resume();
616 21747x work_finished();
617 21747x return 1;
618 }
619
620 − default: // LCOV_EXCL_LINE unreachable: closed key set
621 − continue; // LCOV_EXCL_LINE unreachable: closed key set
622 ✗ }
623 }
624
625 // Signal completions (no OVERLAPPED)
626
3/4
✓ Branch 47 → 48 taken 4803 times.
✓ Branch 47 → 61 taken 14 times.
✓ Branch 48 → 49 taken 71 times.
✗ Branch 48 → 62 not taken.
4888x if (result)
627 {
628
4/6
✓ Branch 48 → 49 taken 875 times.
✓ Branch 48 → 50 taken 3928 times.
✗ Branch 48 → 60 not taken.
✓ Branch 49 → 50 taken 17 times.
✓ Branch 49 → 51 taken 54 times.
✗ Branch 49 → 61 not taken.
4874x switch (key)
629 {
630 892x case key_wake_dispatch:
631 // Timer wakeup - loop to check dispatch_required_
632 892x continue;
633
634 3982x case key_shutdown:
635 3982x ::InterlockedExchange(&stop_event_posted_, 0);
636
3/4
✓ Branch 53 → 54 taken 235 times.
✓ Branch 53 → 59 taken 3693 times.
✗ Branch 54 → 55 not taken.
✓ Branch 54 → 60 taken 54 times.
3982x if (stopped())
637 {
638 // Re-post for other waiting threads
639
1/4
✓ Branch 56 → 57 taken 235 times.
✗ Branch 56 → 58 not taken.
✗ Branch 57 → 58 not taken.
✗ Branch 57 → 59 not taken.
470x if (::InterlockedExchange(&stop_event_posted_, 1) == 0)
640 {
641
1/2
✓ Branch 57 → 58 taken 235 times.
✗ Branch 58 → 59 not taken.
235x ::PostQueuedCompletionStatus(
642 iocp_, 0, key_shutdown, nullptr);
643 }
644 235x return 0;
645 }
646 3747x continue;
647
648 // A key outside the closed set reaches here only if a
649 // third party posts to the port.
650 − default: // LCOV_EXCL_LINE unreachable: closed key set
651 − continue; // LCOV_EXCL_LINE unreachable: closed key set
652 }
653 }
654
655 // Timeout or error. INFINITE never times out, so a WAIT_TIMEOUT
656 // can only be a finite caller timeout (wait_one/poll) elapsing
657 // with no work. run()/run_one() pass INFINITE and block until a
658 // real wake-up; they exit only when stop() posts key_shutdown
659 // (a failed post is fatal in stop()).
660
2/4
✓ Branch 61 → 62 taken 1 time.
✓ Branch 61 → 64 taken 13 times.
✗ Branch 62 → 63 not taken.
✗ Branch 62 → 65 not taken.
14x if (dwError != WAIT_TIMEOUT)
661 1x detail::throw_system_error(make_err(dwError));
662 13x return 0;
663 4639x }
664 }
665
666 inline void
667 2459x win_scheduler::on_timer_changed(void* ctx)
668 {
669 2459x static_cast<win_scheduler*>(ctx)->update_timeout();
670 2459x }
671
672 inline void
673 4581x win_scheduler::set_timer_service(timer_service* svc)
674 {
675 4581x timer_svc_ = svc;
676 // Pass 'this' as context - callback routes to correct instance
677 4581x svc->set_on_earliest_changed(
678 4581x timer_service::callback{this, &on_timer_changed});
679
1/2
✓ Branch 5 → 6 taken 4581 times.
✗ Branch 5 → 8 not taken.
4581x if (timers_)
680 4581x timers_->start();
681 4581x }
682
683 inline void
684 3355x win_scheduler::update_timeout()
685 {
686
3/6
✓ Branch 2 → 3 taken 3355 times.
✗ Branch 2 → 6 not taken.
✓ Branch 4 → 5 taken 3355 times.
✗ Branch 4 → 6 not taken.
✓ Branch 7 → 8 taken 3355 times.
✗ Branch 7 → 11 not taken.
3355x if (timer_svc_ && timers_)
687 3355x timers_->update_timeout(timer_svc_->nearest_expiry());
688 3355x }
689
690 } // namespace boost::corosio::detail
691
692 // Defer including the auxiliary wait reactor until the scheduler is
693 // fully defined, since the reactor's inline methods call back into
694 // win_scheduler. This also gives the ctor, dtor and wait_reactor()
695 // below a complete win_wait_reactor type: the constructor builds it,
696 // the destructor and unique_ptr's deleter tear it down.
697 //
698 // The macro lets win_wait_reactor.hpp diagnose direct inclusion
699 // (which would land it here with win_scheduler still incomplete).
700 #define BOOST_COROSIO_DETAIL_IOCP_WIN_SCHEDULER_BODY_DONE
701 #include <boost/corosio/native/detail/iocp/win_wait_reactor.hpp>
702
703 namespace boost::corosio::detail {
704
705 4582x inline win_scheduler::win_scheduler(
706 4582x capy::execution_context& ctx, int concurrency_hint)
707 4582x : iocp_(nullptr)
708 4582x , outstanding_work_(0)
709 4582x , stopped_(0)
710 4582x , stop_event_posted_(0)
711
1/1
✓ Branch 3 → 4 taken 4582 times.
4582x , dispatch_required_(0)
712 {
713 // concurrency_hint < 0 means use system default (DWORD(~0) = max)
714
2/3
✓ Branch 7 → 8 taken 4582 times.
✗ Branch 7 → 9 not taken.
✓ Branch 10 → 11 taken 4582 times.
4582x iocp_ = ::CreateIoCompletionPort(
715 INVALID_HANDLE_VALUE, nullptr, 0,
716 static_cast<DWORD>(
717 concurrency_hint >= 0 ? concurrency_hint : DWORD(~0)));
718
719
2/2
✓ Branch 11 → 12 taken 1 time.
✓ Branch 11 → 15 taken 4581 times.
4582x if (iocp_ == nullptr)
720
1/1
✓ Branch 12 → 13 taken 1 time.
1x detail::throw_system_error(make_err(::GetLastError()));
721
722 try
723 {
724
1/1
✓ Branch 15 → 16 taken 4581 times.
4581x timers_ = make_win_timers(iocp_, &dispatch_required_);
725
2/2
✓ Branch 18 → 19 taken 4581 times.
✓ Branch 19 → 20 taken 4581 times.
4581x set_timer_service(&get_timer_service(ctx, *this));
726
727 // A scheduler whose wait reactor could not be built would
728 // answer every wait with a parked op, so it refuses to exist
729 // instead. Last, so the catch below is the whole cleanup:
730 // the reactor holds its own Winsock reference and needs no
731 // help from the order.
732
1/1
✓ Branch 20 → 21 taken 4428 times.
4581x wait_reactor_ = std::make_unique<win_wait_reactor>(*this);
733 }
734 153x catch (...)
735 {
736 // ~win_scheduler never runs for a constructor that throws, and
737 // the port is a raw handle nothing else owns. The timer thread
738 // is stopped first because it posts to that port.
739 //
740 // The services registered above are not unregistered here, and
741 // cannot be: the context owns them and offers no way to take
742 // one back. Each holds a scheduler reference this unwind is
743 // about to invalidate, and each survives to be shut down and
744 // destroyed with the context. What makes that safe is that
745 // none of them touches its scheduler while it holds nothing:
746 // a scheduler that never finished constructing handed out no
747 // timer, resolver or file, so every one of those shutdowns
748 // walks an empty list. Registering anything here that would
749 // reach back into the scheduler on an empty shutdown breaks
750 // that, and the fix would have to be a rollback in the
751 // context, not an ordering trick here.
752 153x timers_.reset();
753
1/1
✓ Branch 30 → 31 taken 153 times.
153x ::CloseHandle(iocp_);
754 153x iocp_ = nullptr;
755 153x throw;
756 153x }
757 5044x }
758
759 inline void
760 4428x win_scheduler::shutdown()
761 {
762
1/2
✓ Branch 3 → 4 taken 4428 times.
✗ Branch 3 → 6 not taken.
4428x if (timers_)
763 4428x timers_->stop();
764
765 // Drain timer heap before the work-counting loop. The timer_service
766 // was registered after this scheduler (nested make_service from our
767 // constructor), so execution_context::shutdown() calls us first.
768 // Asio avoids this by owning timer queues directly inside the
769 // scheduler; we bridge the gap by shutting down the timer service
770 // early. The subsequent call from execution_context is a no-op.
771
1/2
✓ Branch 6 → 7 taken 4428 times.
✗ Branch 6 → 8 not taken.
4428x if (timer_svc_)
772 4428x timer_svc_->shutdown();
773
774 // Same problem for the auxiliary wait reactor: ops parked in it
775 // owe completion packets. Stop the reactor early so its loop
776 // posts them as cancelled and the pending count can reach zero.
777 4428x wait_reactor_->stop();
778
779 // Reap every packet still owed to the port before the services
780 // free the op memory those packets reference. Work-guard credits,
781 // posted handlers, and queued continuations have no bearing here.
782
2/2
✓ Branch 34 → 11 taken 37 times.
✓ Branch 34 → 35 taken 4428 times.
8930x while (::InterlockedExchangeAdd(&pending_io_, 0) > 0)
783 {
784 37x op_queue ops;
785 {
786 37x std::lock_guard<win_mutex> lock(dispatch_mutex_);
787 37x ops.splice(completed_ops_);
788 37x }
789
790 // Deferred-queue entries are process-owned (failed-post
791 // fallbacks and posted handlers); no packet references them.
792
1/2
✗ Branch 16 → 17 not taken.
✓ Branch 16 → 19 taken 37 times.
37x while (auto* h = ops.pop())
793 ✗ h->destroy();
794
795 DWORD bytes;
796 ULONG_PTR key;
797 LPOVERLAPPED overlapped;
798
1/1
✓ Branch 19 → 20 taken 37 times.
37x ::GetQueuedCompletionStatus(
799 iocp_, &bytes, &key, &overlapped, iocp::shutdown_drain_timeout_ms);
800
1/2
✓ Branch 20 → 21 taken 37 times.
✗ Branch 20 → 31 not taken.
37x if (overlapped)
801 {
802
1/2
✗ Branch 21 → 22 not taken.
✓ Branch 21 → 23 taken 37 times.
37x if (key == key_posted)
803 {
804 ✗ auto* op = reinterpret_cast<scheduler_op*>(overlapped);
805 ✗ op->destroy();
806 }
807
1/2
✗ Branch 23 → 24 not taken.
✓ Branch 23 → 27 taken 37 times.
37x else if (key == key_continuation)
808 {
809 // Drain without resuming: destroy the parked frame.
810 ✗ auto* c = reinterpret_cast<capy::continuation*>(overlapped);
811 ✗ if (c->h)
812 ✗ c->h.destroy();
813 }
814 else
815 {
816 37x ::InterlockedDecrement(&pending_io_);
817 37x auto* op = overlapped_to_op(overlapped);
818
1/1
✓ Branch 30 → 31 taken 37 times.
37x op->destroy();
819 }
820 }
821 }
822
823 // Final sweep: packets can sit in the port or deferred queue
824 // without a pending_io_ count (posted directly against the
825 // handle). Destroy them so service teardown does not free state
826 // they still reference.
827 for (;;)
828 {
829 4448x op_queue ops;
830 {
831 4448x std::lock_guard<win_mutex> lock(dispatch_mutex_);
832 4448x ops.splice(completed_ops_);
833 4448x }
834
1/2
✗ Branch 40 → 41 not taken.
✓ Branch 40 → 43 taken 4448 times.
4448x while (auto* h = ops.pop())
835 ✗ h->destroy();
836
837 DWORD bytes;
838 ULONG_PTR key;
839 LPOVERLAPPED overlapped;
840
1/1
✓ Branch 43 → 44 taken 4448 times.
4448x ::GetQueuedCompletionStatus(iocp_, &bytes, &key, &overlapped, 0);
841
2/2
✓ Branch 44 → 45 taken 4428 times.
✓ Branch 44 → 46 taken 20 times.
4448x if (!overlapped)
842 4428x break;
843
2/2
✓ Branch 46 → 47 taken 13 times.
✓ Branch 46 → 48 taken 7 times.
20x if (key == key_posted)
844 {
845
1/1
✓ Branch 47 → 54 taken 13 times.
13x reinterpret_cast<scheduler_op*>(overlapped)->destroy();
846 }
847
2/2
✓ Branch 48 → 49 taken 3 times.
✓ Branch 48 → 52 taken 4 times.
7x else if (key == key_continuation)
848 {
849 3x auto* c = reinterpret_cast<capy::continuation*>(overlapped);
850
1/2
✓ Branch 50 → 51 taken 3 times.
✗ Branch 50 → 54 not taken.
3x if (c->h)
851
1/1
✓ Branch 51 → 54 taken 3 times.
3x c->h.destroy();
852 }
853 else
854 {
855
1/1
✓ Branch 53 → 54 taken 4 times.
4x overlapped_to_op(overlapped)->destroy();
856 }
857 20x }
858 4428x }
859
860 8856x inline win_scheduler::~win_scheduler()
861 {
862 4428x wait_reactor_->stop();
863 4428x wait_reactor_.reset();
864
865
1/2
✓ Branch 5 → 6 taken 4428 times.
✗ Branch 5 → 7 not taken.
4428x if (iocp_ != nullptr)
866 4428x ::CloseHandle(iocp_);
867 8856x }
868
869 inline win_wait_reactor&
870 48x win_scheduler::wait_reactor()
871 {
872 48x return *wait_reactor_;
873 }
874
875 inline void
876 28718x win_scheduler::cancel_wait(overlapped_op* op) noexcept
877 {
878 28718x wait_reactor_->cancel_wait(op);
879 28718x }
880
881 } // namespace boost::corosio::detail
882
883 #endif // BOOST_COROSIO_HAS_IOCP
884
885 #endif // BOOST_COROSIO_NATIVE_DETAIL_IOCP_WIN_SCHEDULER_HPP
886