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

93.2% Lines (193/207) 100.0% List of functions (16/16) 81.0% Branches (132/163)
win_wait_reactor.hpp
f(x) Functions (16)
Function Calls Lines Branches Blocks
boost::corosio::detail::win_wait_reactor::last_error() :146 137x 100.0% 50.0% 83.3% boost::corosio::detail::win_wait_reactor::events_for_wait(boost::corosio::wait_type) :152 46x 100.0% 100.0% 100.0% boost::corosio::detail::win_wait_reactor::ready_for_wait(boost::corosio::wait_type, short) :171 29x 100.0% 100.0% 100.0% boost::corosio::detail::win_wait_reactor::win_wait_reactor(boost::corosio::detail::win_scheduler&) :214 4581x 100.0% 100.0% 100.0% boost::corosio::detail::win_wait_reactor::~win_wait_reactor() :229 4428x 100.0% – 100.0% boost::corosio::detail::win_wait_reactor::make_wakeup_pair() :236 4564x 100.0% 100.0% 100.0% boost::corosio::detail::win_wait_reactor::close_wakeup_pair() :308 4445x 100.0% 50.0% 100.0% boost::corosio::detail::win_wait_reactor::wake_self() :323 33181x 100.0% 100.0% 100.0% boost::corosio::detail::win_wait_reactor::register_wait(unsigned long long, boost::corosio::wait_type, boost::corosio::detail::overlapped_op*) :345 48x 75.0% 80.0% 81.8% boost::corosio::detail::win_wait_reactor::queue_register(boost::corosio::detail::win_wait_reactor::entry const&) :370 48x 69.2% 87.5% 68.0% boost::corosio::detail::win_wait_reactor::queue_register(boost::corosio::detail::win_wait_reactor::entry const&)::{lambda()#1}::operator()() const :400 39x 100.0% 100.0% 100.0% boost::corosio::detail::win_wait_reactor::cancel_wait(boost::corosio::detail::overlapped_op*) :413 28718x 100.0% 100.0% 89.5% boost::corosio::detail::win_wait_reactor::stop() :431 13284x 100.0% 100.0% 84.2% boost::corosio::detail::win_wait_reactor::drop_refused_entries() :450 1x 72.7% 57.1% 52.4% boost::corosio::detail::win_wait_reactor::run() :475 39x 93.0% 75.0% 76.6% boost::corosio::detail::win_wait_reactor::run()::{lambda(boost::corosio::detail::win_wait_reactor::entry const&)#1}::operator()(boost::corosio::detail::win_wait_reactor::entry const&) const :509 225x 100.0% – 100.0%
Line Branch TLA Hits Source Code
1 //
2 // Copyright (c) 2026 Michael Vandeberg
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_NATIVE_DETAIL_IOCP_WIN_WAIT_REACTOR_HPP
11 #define BOOST_COROSIO_NATIVE_DETAIL_IOCP_WIN_WAIT_REACTOR_HPP
12
13 #include <boost/corosio/detail/platform.hpp>
14
15 #if BOOST_COROSIO_HAS_IOCP
16
17 // This header is included from the bottom of win_scheduler.hpp after
18 // the scheduler class is fully defined. Including it directly would
19 // circle back into a still-incomplete win_scheduler when the dtor's
20 // unique_ptr<win_wait_reactor>::reset() is parsed. Diagnose that
21 // rather than emitting a confusing "incomplete type" error far away.
22 #ifndef BOOST_COROSIO_DETAIL_IOCP_WIN_SCHEDULER_BODY_DONE
23 #error "Include <boost/corosio/native/detail/iocp/win_scheduler.hpp> \
24 instead of including this header directly."
25 #endif
26
27 #include <boost/corosio/wait_type.hpp>
28 #include <boost/corosio/native/detail/iocp/win_overlapped_op.hpp>
29 #include <boost/corosio/native/detail/iocp/win_scheduler.hpp>
30 #include <boost/corosio/native/detail/iocp/win_windows.hpp>
31 #include <boost/corosio/native/detail/iocp/win_wsa_init.hpp>
32
33 #include <ws2tcpip.h>
34
35 #include <algorithm>
36 #include <atomic>
37 #include <cstddef>
38 #include <mutex>
39 #include <thread>
40 #include <vector>
41
42 namespace boost::corosio::detail {
43
44 /** Auxiliary select-based reactor for IOCP wait operations.
45
46 IOCP has no native primitive for socket readiness without I/O.
47 For cases where a zero-byte WSARecv won't work (datagram-read,
48 acceptor-read, write-wait, error-wait), this reactor runs a
49 dedicated thread using WSAPoll to detect readiness and posts a
50 synthetic completion to the owning IOCP scheduler via
51 win_scheduler::on_completion().
52
53 The same dispatch path used by overlapped I/O then delivers the
54 completion to the user's coroutine, so the public API is uniform
55 across backends.
56
57 Per-op lifecycle:
58 1. Caller sets up an overlapped_op (h, ex, ec_out, cancelled flag).
59 2. Caller calls register_wait(fd, w, op) and returns
60 std::noop_coroutine. The op is parked in the reactor's table.
61 3. Reactor thread polls. When the fd is ready, the op is removed
62 from the table and posted to the scheduler. A read or write
63 wait always completes with ec={}: error revents mean "ready",
64 and the next read or write names the error. A wait_type::error
65 wait completes with the SO_ERROR value when it is set, or
66 WSAECONNABORTED as a synthesized fallback when it is zero.
67 4. On socket cancel(), the user's thread calls cancel_wait(op),
68 which queues a cancel request. The reactor thread removes the
69 op from the table and posts a completion; invoke_handler sees
70 op.cancelled==true and yields capy::cond::canceled.
71
72 A poll the provider refuses on account of one descriptor answers
73 that descriptor's op with the refusal and leaves the rest of the
74 table polling. Such an error reaches the caller only where nothing
75 flagged the op first: close() and cancel() set the cancelled flag
76 before the handle goes, and a flagged op yields canceled whatever
77 the completion carries.
78
79 A refusal none of the entries accounts for is the end of the
80 reactor: the error is latched, the ops it still holds are completed
81 with it, and every later register_wait is completed with it too,
82 since nothing would drain an op parked after that point. Both
83 answers name the reason waits can no longer be served; canceled is
84 reserved for ops close() or cancel() flagged, and for the ops a
85 stop() drain finds parked.
86
87 The constructor builds the wakeup channel and throws if it cannot:
88 a reactor that cannot be woken can never report readiness, so
89 there is no reactor worth handing back. The polling thread is a
90 separate cost, paid by the first register_wait, so a context that
91 never waits never carries one. A thread the system refuses costs
92 only the wait that asked for it: that wait completes with
93 `resource_unavailable_try_again` and the next one tries again.
94
95 Thread-safe: register_wait, cancel_wait, and stop may be called
96 from any thread.
97 */
98 class win_wait_reactor : private win_wsa_init
99 {
100 public:
101 /** Construct the reactor and its wakeup channel.
102
103 @par Exception Safety
104 Strong guarantee. A channel that cannot be formed leaves no
105 socket open.
106
107 @param sched The scheduler synthetic completions are posted to.
108
109 @throws std::system_error If Winsock could not be started or
110 the wakeup socket pair could not be built.
111 */
112 explicit win_wait_reactor(win_scheduler& sched);
113
114 /// Stop the reactor thread and close the wakeup channel.
115 ~win_wait_reactor();
116
117 win_wait_reactor(win_wait_reactor const&) = delete;
118 win_wait_reactor& operator=(win_wait_reactor const&) = delete;
119
120 /// Park an overlapped_op until @p fd is ready for @p w.
121 void register_wait(SOCKET fd, wait_type w, overlapped_op* op);
122
123 /// Remove a parked op and post a completion. Idempotent.
124 void cancel_wait(overlapped_op* op);
125
126 /// Stop the reactor thread and drain remaining ops as cancelled.
127 void stop();
128
129 private:
130 struct entry
131 {
132 SOCKET fd = INVALID_SOCKET;
133 wait_type w = wait_type::read;
134 overlapped_op* op = nullptr;
135 };
136
137 void run();
138 bool drop_refused_entries();
139 DWORD queue_register(entry const& e);
140 void wake_self() noexcept;
141 DWORD make_wakeup_pair() noexcept;
142 void close_wakeup_pair() noexcept;
143
144 // A failed call that left a zero last error would answer "no
145 // error" and put the reactor straight back on the silent path.
146 137x static DWORD last_error() noexcept
147 {
148 137x DWORD const err = ::WSAGetLastError();
149
1/2
✓ Branch 3 → 4 taken 137 times.
✗ Branch 3 → 5 not taken.
137x return err != 0 ? err : static_cast<DWORD>(WSAEINVAL);
150 }
151
152 46x static SHORT events_for_wait(wait_type w) noexcept
153 {
154
3/3
✓ Branch 2 → 3 taken 22 times.
✓ Branch 2 → 4 taken 14 times.
✓ Branch 2 → 5 taken 10 times.
46x switch (w)
155 {
156 22x case wait_type::read:
157 22x return POLLRDNORM;
158 14x case wait_type::write:
159 14x return POLLWRNORM;
160 // The Microsoft provider does not implement POLLPRI and
161 // refuses the whole call when it is asked for, which would
162 // take every other registration in the set with it. The band
163 // it does implement carries the same out-of-band meaning, and
164 // the error conditions an error wait is really after arrive in
165 // revents whether or not they were asked for.
166 10x default:
167 10x return POLLRDBAND;
168 }
169 }
170
171 29x static bool ready_for_wait(wait_type w, SHORT revents) noexcept
172 {
173 29x constexpr SHORT err_bits = POLLERR | POLLHUP | POLLNVAL;
174
3/3
✓ Branch 2 → 3 taken 11 times.
✓ Branch 2 → 4 taken 14 times.
✓ Branch 2 → 5 taken 4 times.
29x switch (w)
175 {
176 11x case wait_type::read:
177 11x return (revents & (POLLRDNORM | POLLRDBAND | err_bits)) != 0;
178 14x case wait_type::write:
179 14x return (revents & (POLLWRNORM | POLLWRBAND | err_bits)) != 0;
180 4x default:
181 4x return (revents & (POLLRDBAND | err_bits)) != 0;
182 }
183 }
184
185 win_scheduler& sched_;
186
187 // Built by the constructor and closed by the destructor, so every
188 // other member function can assume a usable channel.
189 SOCKET wakeup_read_ = INVALID_SOCKET;
190 SOCKET wakeup_write_ = INVALID_SOCKET;
191
192 // Also guards thread_ against a start racing the stop that joins it.
193 std::mutex mutex_;
194 std::vector<entry> pending_register_;
195 std::vector<overlapped_op*> pending_cancel_;
196 std::atomic<bool> stop_{false};
197 std::atomic<bool> wake_pending_{false};
198
199 // What the exit told the ops it was still holding, kept so a later
200 // register_wait can be told the same thing instead of inventing a
201 // cancellation nobody asked for. Non-zero is also what says the
202 // polling thread left on its own after a WSAPoll error: stop_ must
203 // not be used for that exit -- stop() reads it as "already
204 // stopped" and would skip the join that keeps the thread from
205 // being destroyed joinable. Nothing latches a zero here, because
206 // a failed call that left a zero last error is substituted for.
207 DWORD dead_err_ = 0;
208
209 std::vector<entry> registered_; // reactor-thread-only
210
211 std::thread thread_;
212 };
213
214 4581x inline win_wait_reactor::win_wait_reactor(win_scheduler& sched) : sched_(sched)
215 {
216 // The win_wsa_init base is what makes the sockets below legal, and
217 // holding the reference rather than borrowing someone else's is
218 // what keeps them legal to the end: a base is constructed before
219 // this body and released after ~win_wait_reactor has closed the
220 // pair, so WSACleanup can never land between the two.
221 //
222 // A reactor that cannot be woken would park every op it is handed
223 // forever, so it refuses to exist rather than being handed out
224 // broken. The polling thread waits for the first register_wait.
225
2/2
✓ Branch 11 → 12 taken 136 times.
✓ Branch 11 → 14 taken 4428 times.
4564x if (DWORD const err = make_wakeup_pair(); err != 0)
226 136x detail::throw_system_error(make_err(err), "win_wait_reactor");
227 5244x }
228
229 4428x inline win_wait_reactor::~win_wait_reactor()
230 {
231 4428x stop();
232 4428x close_wakeup_pair();
233 4428x }
234
235 inline DWORD
236 4564x win_wait_reactor::make_wakeup_pair() noexcept
237 {
238 // Build a pair of connected loopback sockets to use as a wakeup
239 // channel. Winsock has no socketpair(2), so we listen on
240 // 127.0.0.1:0, connect a peer, then accept it.
241 //
242 // Every failure path reads the last error before closing
243 // anything: closesocket() overwrites it.
244 4564x SOCKET listener = ::socket(AF_INET, SOCK_STREAM, IPPROTO_TCP);
245
2/2
✓ Branch 3 → 4 taken 17 times.
✓ Branch 3 → 5 taken 4547 times.
4564x if (listener == INVALID_SOCKET)
246 17x return last_error();
247
248 4547x sockaddr_in addr{};
249 4547x addr.sin_family = AF_INET;
250 4547x addr.sin_addr.s_addr = ::htonl(INADDR_LOOPBACK);
251 4547x addr.sin_port = 0;
252
253 4547x int len = sizeof(addr);
254 4547x if (::bind(listener, reinterpret_cast<sockaddr*>(&addr), len) ==
255 4530x SOCKET_ERROR ||
256
6/6
✓ Branch 7 → 8 taken 4530 times.
✓ Branch 7 → 12 taken 17 times.
✓ Branch 9 → 10 taken 4513 times.
✓ Branch 9 → 12 taken 17 times.
✓ Branch 14 → 15 taken 51 times.
✓ Branch 14 → 18 taken 4496 times.
9060x ::listen(listener, 1) == SOCKET_ERROR ||
257
2/2
✓ Branch 11 → 12 taken 17 times.
✓ Branch 11 → 13 taken 4496 times.
4513x ::getsockname(listener, reinterpret_cast<sockaddr*>(&addr), &len) ==
258 SOCKET_ERROR)
259 {
260 51x DWORD const err = last_error();
261 51x ::closesocket(listener);
262 51x return err;
263 }
264
265 4496x wakeup_write_ = ::socket(AF_INET, SOCK_STREAM, IPPROTO_TCP);
266
2/2
✓ Branch 19 → 20 taken 17 times.
✓ Branch 19 → 23 taken 4479 times.
4496x if (wakeup_write_ == INVALID_SOCKET)
267 {
268 17x DWORD const err = last_error();
269 17x ::closesocket(listener);
270 17x return err;
271 }
272
273
2/2
✓ Branch 24 → 25 taken 17 times.
✓ Branch 24 → 29 taken 4462 times.
4479x if (::connect(wakeup_write_, reinterpret_cast<sockaddr*>(&addr), len) ==
274 SOCKET_ERROR)
275 {
276 17x DWORD const err = last_error();
277 17x ::closesocket(wakeup_write_);
278 17x wakeup_write_ = INVALID_SOCKET;
279 17x ::closesocket(listener);
280 17x return err;
281 }
282
283 4462x wakeup_read_ = ::accept(listener, nullptr, nullptr);
284
2/2
✓ Branch 30 → 31 taken 17 times.
✓ Branch 30 → 35 taken 4445 times.
4462x if (wakeup_read_ == INVALID_SOCKET)
285 {
286 17x DWORD const err = last_error();
287 17x ::closesocket(listener);
288 17x ::closesocket(wakeup_write_);
289 17x wakeup_write_ = INVALID_SOCKET;
290 17x return err;
291 }
292 4445x ::closesocket(listener);
293
294 // The drain loop in run() calls recv() until it returns <= 0.
295 // With a blocking socket that second recv() would block instead
296 // of returning WSAEWOULDBLOCK, deadlocking the reactor thread.
297 4445x u_long non_blocking = 1;
298
2/2
✓ Branch 37 → 38 taken 17 times.
✓ Branch 37 → 41 taken 4428 times.
4445x if (::ioctlsocket(wakeup_read_, FIONBIO, &non_blocking) == SOCKET_ERROR)
299 {
300 17x DWORD const err = last_error();
301 17x close_wakeup_pair();
302 17x return err;
303 }
304 4428x return 0;
305 }
306
307 inline void
308 4445x win_wait_reactor::close_wakeup_pair() noexcept
309 {
310
1/2
✓ Branch 2 → 3 taken 4445 times.
✗ Branch 2 → 5 not taken.
4445x if (wakeup_read_ != INVALID_SOCKET)
311 {
312 4445x ::closesocket(wakeup_read_);
313 4445x wakeup_read_ = INVALID_SOCKET;
314 }
315
1/2
✓ Branch 5 → 6 taken 4445 times.
✗ Branch 5 → 8 not taken.
4445x if (wakeup_write_ != INVALID_SOCKET)
316 {
317 4445x ::closesocket(wakeup_write_);
318 4445x wakeup_write_ = INVALID_SOCKET;
319 }
320 4445x }
321
322 inline void
323 33181x win_wait_reactor::wake_self() noexcept
324 {
325 // Coalesce wakes: only send a byte if no wake is already pending.
326 33181x bool expected = false;
327
2/2
✓ Branch 3 → 4 taken 28547 times.
✓ Branch 3 → 5 taken 4634 times.
33181x if (!wake_pending_.compare_exchange_strong(
328 expected, true, std::memory_order_acq_rel))
329 28547x return;
330
331 4634x char b = 0;
332
2/2
✓ Branch 6 → 7 taken 2 times.
✓ Branch 6 → 8 taken 4632 times.
4634x if (::send(wakeup_write_, &b, 1, 0) == SOCKET_ERROR)
333 {
334 // The flag is what coalesces later wakes into a byte already
335 // in the channel; a send that failed put no byte there, so
336 // leaving it latched would swallow every wake that follows.
337 // Disarming keeps the cost to the wakes already in flight: the
338 // next register, cancel or stop sends its own byte and the
339 // reactor learns about both.
340 2x wake_pending_.store(false, std::memory_order_release);
341 }
342 }
343
344 inline void
345 48x win_wait_reactor::register_wait(SOCKET fd, wait_type w, overlapped_op* op)
346 {
347 // If the op was already cancelled (e.g. pre-cancelled stop_token
348 // fired synchronously before this call), complete immediately
349 // instead of registering. Otherwise the reactor would park the
350 // op forever because the earlier cancel_wait() found nothing to
351 // cancel in registered_.
352
1/2
✗ Branch 3 → 4 not taken.
✓ Branch 3 → 6 taken 48 times.
48x if (op->cancelled.load(std::memory_order_acquire))
353 {
354 ✗ sched_.on_completion(op, 0, 0);
355 ✗ return;
356 }
357
358
3/3
✓ Branch 6 → 7 taken 48 times.
✓ Branch 7 → 8 taken 1 time.
✓ Branch 7 → 10 taken 47 times.
48x if (DWORD const err = queue_register(entry{fd, w, op}); err != 0)
359 {
360 // Nothing would ever drain a parked op, and the refusal knows
361 // why: the abort a stop drains with, the error the polling
362 // thread died of, or the system declining a thread.
363 1x sched_.on_completion(op, err, 0);
364 1x return;
365 }
366 47x wake_self();
367 }
368
369 inline DWORD
370 48x win_wait_reactor::queue_register(entry const& e)
371 {
372
1/1
✓ Branch 2 → 3 taken 48 times.
48x std::lock_guard lock(mutex_);
373 // stop() sets the flag before it takes the thread out from under
374 // this lock, and never joins again. Checking the flag and queueing
375 // in one critical section is what keeps both halves honest: a
376 // thread started after that join would be destroyed still joinable,
377 // which ends the process, and an op queued after it would have no
378 // drainer. Queueing under the flag instead leaves the op for the
379 // drain run() performs on its way out.
380 // dead_err_ says the same thing for the other exit: a thread that
381 // left on a WSAPoll error drained what it held and will not poll
382 // again, so a register queued after it would wait on nobody. It
383 // answers with what killed it, so the caller learns why its waits
384 // stopped being served instead of hearing that something cancelled
385 // them.
386
1/2
✗ Branch 4 → 5 not taken.
✓ Branch 4 → 6 taken 48 times.
48x if (stop_.load(std::memory_order_acquire))
387 ✗ return ERROR_OPERATION_ABORTED;
388
2/2
✓ Branch 6 → 7 taken 1 time.
✓ Branch 6 → 8 taken 47 times.
48x if (dead_err_)
389 1x return dead_err_;
390
391 // A polling thread costs a thread per context, and a context that
392 // never waits never pays for one; the first wait is what starts it.
393 // A system that will not give one refuses this wait alone: nothing
394 // has been queued yet, so the reactor is left exactly as it was and
395 // the next register_wait asks again.
396
2/2
✓ Branch 9 → 10 taken 39 times.
✓ Branch 9 → 14 taken 8 times.
47x if (!thread_.joinable())
397 {
398 try
399 {
400
1/1
✓ Branch 10 → 11 taken 39 times.
78x thread_ = std::thread([this] { run(); });
401 }
402 ✗ catch (...)
403 {
404 ✗ return ERROR_MAX_THRDS_REACHED;
405 ✗ }
406 }
407
408
1/1
✓ Branch 14 → 15 taken 47 times.
47x pending_register_.push_back(e);
409 47x return 0;
410 48x }
411
412 inline void
413 28718x win_wait_reactor::cancel_wait(overlapped_op* op)
414 {
415 {
416
1/1
✓ Branch 2 → 3 taken 28718 times.
28718x std::lock_guard lock(mutex_);
417 // Same refusal queue_register makes, for the same reason: once
418 // the polling thread is gone nothing reads pending_cancel_
419 // again, so a cancel queued here would sit there for good. It
420 // has nothing to cancel either -- the drain those exits run
421 // already answered whatever was parked, and a register that
422 // arrived after them was refused at its caller.
423
6/6
✓ Branch 4 → 5 taken 28712 times.
✓ Branch 4 → 6 taken 6 times.
✓ Branch 5 → 6 taken 6 times.
✓ Branch 5 → 7 taken 28706 times.
✓ Branch 8 → 9 taken 12 times.
✓ Branch 8 → 10 taken 28706 times.
28718x if (stop_.load(std::memory_order_acquire) || dead_err_)
424 12x return;
425
1/1
✓ Branch 10 → 11 taken 28706 times.
28706x pending_cancel_.push_back(op);
426 28718x }
427 28706x wake_self();
428 }
429
430 inline void
431 13284x win_wait_reactor::stop()
432 {
433
2/2
✓ Branch 3 → 4 taken 8856 times.
✓ Branch 3 → 5 taken 4428 times.
13284x if (stop_.exchange(true, std::memory_order_acq_rel))
434 8856x return;
435 4428x wake_self();
436 // Moved out under the lock, then joined without it: the reactor
437 // thread takes the same lock on every pass, so joining while
438 // holding it would deadlock. A context that never waited has no
439 // thread here at all.
440 4428x std::thread t;
441 {
442
1/1
✓ Branch 7 → 8 taken 4428 times.
4428x std::lock_guard lock(mutex_);
443 4428x t = std::move(thread_);
444 4428x }
445
2/2
✓ Branch 12 → 13 taken 39 times.
✓ Branch 12 → 14 taken 4389 times.
4428x if (t.joinable())
446
1/1
✓ Branch 13 → 14 taken 39 times.
39x t.join();
447 4428x }
448
449 inline bool
450 1x win_wait_reactor::drop_refused_entries()
451 {
452 // A descriptor the provider no longer recognises either comes back
453 // POLLNVAL on its own entry, which the revents walk already
454 // answers, or refuses the whole call and reaches here. Asking
455 // about each entry alone covers both without depending on which
456 // one this provider does. A socket closed under a parked wait is
457 // the ordinary way to get here: the reactor is told to drop the
458 // entry, but the handle can go before that ask is read.
459 1x bool dropped = false;
460
2/2
✓ Branch 16 → 3 taken 1 time.
✓ Branch 16 → 17 taken 1 time.
2x for (std::size_t i = registered_.size(); i > 0; --i)
461 {
462 1x auto const& e = registered_[i - 1];
463 1x WSAPOLLFD pfd{e.fd, events_for_wait(e.w), 0};
464
2/3
✓ Branch 5 → 6 taken 1 time.
✓ Branch 6 → 7 taken 1 time.
✗ Branch 6 → 8 not taken.
1x if (::WSAPoll(&pfd, 1, 0) != SOCKET_ERROR)
465 1x continue;
466
467 ✗ sched_.on_completion(e.op, last_error(), 0);
468 ✗ registered_.erase(registered_.begin() + (i - 1));
469 ✗ dropped = true;
470 }
471 1x return dropped;
472 }
473
474 inline void
475 39x win_wait_reactor::run()
476 {
477 39x std::vector<WSAPOLLFD> pollfds;
478
479 // What the ops still parked here are told on the way out, and what
480 // every later register_wait is told too. Ending on a poll the
481 // provider refused is not a cancellation, and an op that reports
482 // one hides the reason its wait could not be kept.
483 39x DWORD drain_err = ERROR_OPERATION_ABORTED;
484
485
4/4
✓ Branch 108 → 3 taken 239 times.
✓ Branch 108 → 109 taken 37 times.
✓ Branch 110 → 3 taken 12 times.
✓ Branch 110 → 111 taken 1 time.
289x while (!stop_.load(std::memory_order_acquire))
486 {
487 // Drain pending register/cancel under the lock.
488 251x std::vector<entry> to_add;
489 251x std::vector<overlapped_op*> to_cancel;
490 {
491
1/1
✓ Branch 3 → 4 taken 251 times.
251x std::lock_guard lock(mutex_);
492 251x to_add.swap(pending_register_);
493 251x to_cancel.swap(pending_cancel_);
494 251x }
495
496
2/2
✓ Branch 13 → 9 taken 47 times.
✓ Branch 13 → 14 taken 251 times.
298x for (auto& e : to_add)
497
1/1
✓ Branch 10 → 11 taken 47 times.
47x registered_.push_back(e);
498
499
2/2
✓ Branch 33 → 16 taken 485 times.
✓ Branch 33 → 34 taken 251 times.
736x for (auto* op : to_cancel)
500 {
501 // A socket's close() and cancel() ask the reactor to drop
502 // their wait op whether or not one is parked -- open() goes
503 // through close_socket() before it has a socket at all.
504 // Only a parked op pins its owner, so an op that is not in
505 // registered_ may belong to an object already destroyed:
506 // look it up before touching it.
507
1/1
✓ Branch 19 → 20 taken 485 times.
485x auto it = std::find_if(
508 registered_.begin(), registered_.end(),
509 225x [op](entry const& e) { return e.op == op; });
510
2/2
✓ Branch 22 → 23 taken 452 times.
✓ Branch 22 → 24 taken 33 times.
485x if (it == registered_.end())
511 452x continue; // Already fired, or never parked.
512
513 // Such an ask can also outlive the op's next reset(). Acting
514 // on it then would complete the wait that reset started,
515 // the moment it is registered. Every real cancel flags the
516 // op before queueing the ask, and only reset() clears the
517 // flag, so an unflagged op is one of those stale asks.
518
2/2
✓ Branch 25 → 26 taken 17 times.
✓ Branch 25 → 30 taken 16 times.
33x if (op->cancelled.load(std::memory_order_acquire))
519 {
520 // The op's cancelled flag has already been set by
521 // request_cancel; invoke_handler will translate it.
522
1/1
✓ Branch 26 → 27 taken 17 times.
17x sched_.on_completion(op, 0, 0);
523
1/1
✓ Branch 28 → 29 taken 17 times.
17x registered_.erase(it);
524 }
525 }
526
527 // Build the poll set. Slot 0 is the wakeup socket.
528 251x pollfds.clear();
529
1/1
✓ Branch 36 → 37 taken 251 times.
251x pollfds.reserve(registered_.size() + 1);
530
1/1
✓ Branch 37 → 38 taken 251 times.
251x pollfds.push_back({wakeup_read_, POLLRDNORM, 0});
531
2/2
✓ Branch 45 → 40 taken 45 times.
✓ Branch 45 → 46 taken 251 times.
296x for (auto& e : registered_)
532
1/1
✓ Branch 42 → 43 taken 45 times.
45x pollfds.push_back({e.fd, events_for_wait(e.w), 0});
533
534 // Block until the self-pipe (slot 0) is poked by a register,
535 // cancel, or stop, or a watched socket becomes ready. No
536 // periodic timeout, so an idle reactor consumes no CPU: the
537 // self-pipe is the only thing that ends this wait, and
538 // wake_self() leaves the channel free for the next poke when
539 // its own send fails, so a lost wake costs the wakes already in
540 // flight rather than every wake after it.
541
1/1
✓ Branch 48 → 49 taken 251 times.
251x int n = ::WSAPoll(
542 251x pollfds.data(), static_cast<ULONG>(pollfds.size()),
543 -1 /* infinite */);
544
2/2
✓ Branch 49 → 50 taken 1 time.
✓ Branch 49 → 60 taken 250 times.
251x if (n == SOCKET_ERROR)
545 {
546 // A refusal one entry accounts for costs that entry its
547 // wait and nothing else; the rest of the set keeps being
548 // polled, and a later register still finds a live reactor.
549 // Only a failure no entry explains ends the loop.
550 1x DWORD const err = last_error();
551
4/7
✓ Branch 52 → 53 taken 1 time.
✗ Branch 52 → 55 not taken.
✓ Branch 53 → 54 taken 1 time.
✓ Branch 54 → 55 taken 1 time.
✗ Branch 54 → 56 not taken.
✓ Branch 57 → 58 taken 1 time.
✗ Branch 57 → 59 not taken.
1x if (registered_.empty() || !drop_refused_entries())
552 {
553 1x drain_err = err;
554 1x break;
555 }
556 ✗ continue;
557 ✗ }
558
559 // Drain the wakeup socket so it stops reporting readable.
560
3/3
✓ Branch 61 → 62 taken 242 times.
✓ Branch 61 → 68 taken 7 times.
✓ Branch 61 → 69 taken 1 time.
250x if (pollfds[0].revents != 0)
561 {
562 char buf[64];
563 for (;;)
564 {
565
2/2
✓ Branch 62 → 63 taken 462 times.
✓ Branch 63 → 64 taken 22 times.
484x int r = ::recv(wakeup_read_, buf, sizeof(buf), 0);
566
4/4
✓ Branch 63 → 64 taken 231 times.
✓ Branch 63 → 65 taken 231 times.
✓ Branch 64 → 65 taken 11 times.
✓ Branch 64 → 66 taken 11 times.
484x if (r <= 0)
567 242x break;
568 242x }
569 242x wake_pending_.store(false, std::memory_order_release);
570 }
571
572 // Walk events in reverse so erases don't invalidate later indices.
573
4/4
✓ Branch 94 → 69 taken 42 times.
✓ Branch 94 → 95 taken 238 times.
✓ Branch 95 → 70 taken 2 times.
✓ Branch 95 → 96 taken 12 times.
294x for (std::size_t i = pollfds.size(); i > 1; --i)
574 {
575 44x auto const& pfd = pollfds[i - 1];
576
4/4
✓ Branch 70 → 71 taken 14 times.
✓ Branch 70 → 72 taken 28 times.
✓ Branch 71 → 72 taken 1 time.
✓ Branch 71 → 73 taken 1 time.
44x if (pfd.revents == 0)
577 15x continue;
578
579 29x auto const& e = registered_[i - 2];
580
2/4
✗ Branch 74 → 75 not taken.
✓ Branch 74 → 76 taken 28 times.
✗ Branch 75 → 76 not taken.
✓ Branch 75 → 77 taken 1 time.
29x if (!ready_for_wait(e.w, pfd.revents))
581 ✗ continue;
582
583 29x DWORD err = 0;
584 29x constexpr SHORT err_bits = POLLERR | POLLHUP | POLLNVAL;
585 // A read or write wait reports readiness: the next read or
586 // write names the error. SO_ERROR is read only for an error
587 // wait, because reading it resets it.
588
4/7
✓ Branch 76 → 77 taken 4 times.
✓ Branch 76 → 87 taken 24 times.
✓ Branch 77 → 78 taken 4 times.
✗ Branch 77 → 87 not taken.
✓ Branch 77 → 88 taken 1 time.
✗ Branch 78 → 79 not taken.
✗ Branch 78 → 88 not taken.
29x if (e.w == wait_type::error && (pfd.revents & err_bits))
589 {
590 4x int so_err = 0;
591 4x int sz = sizeof(so_err);
592 12x if (::getsockopt(
593
1/2
✓ Branch 78 → 79 taken 4 times.
✗ Branch 79 → 80 not taken.
4x e.fd, SOL_SOCKET, SO_ERROR,
594
3/8
✓ Branch 79 → 80 taken 2 times.
✓ Branch 79 → 82 taken 2 times.
✗ Branch 80 → 81 not taken.
✗ Branch 80 → 83 not taken.
✗ Branch 83 → 84 not taken.
✓ Branch 83 → 85 taken 4 times.
✗ Branch 84 → 85 not taken.
✗ Branch 84 → 86 not taken.
6x reinterpret_cast<char*>(&so_err), &sz) == 0 &&
595
1/4
✗ Branch 80 → 81 not taken.
✓ Branch 80 → 82 taken 2 times.
✗ Branch 81 → 82 not taken.
✗ Branch 81 → 83 not taken.
2x so_err != 0)
596 ✗ err = static_cast<DWORD>(so_err);
597 else
598 // The contract is to report a non-zero error_code.
599 4x err = WSAECONNABORTED;
600 }
601
602
2/2
✓ Branch 87 → 88 taken 28 times.
✓ Branch 88 → 89 taken 1 time.
29x sched_.on_completion(e.op, err, 0);
603
2/2
✓ Branch 91 → 92 taken 28 times.
✓ Branch 92 → 93 taken 1 time.
29x registered_.erase(registered_.begin() + (i - 2));
604 }
605 252x }
606
607 // Drain remaining ops on the way out. This must cover both the
608 // active set and anything still queued by user threads that hasn't
609 // been moved into registered_ yet, otherwise those ops leak
610 // work_started credit and stall scheduler shutdown.
611 {
612
2/2
✓ Branch 109 → 110 taken 38 times.
✓ Branch 111 → 112 taken 1 time.
39x std::lock_guard lock(mutex_);
613 // Closing the door and taking what is behind it in one critical
614 // section is what leaves no register in between: one that got
615 // in is drained here, one that arrives after is refused by
616 // queue_register and completes with the same code at its
617 // caller.
618 39x dead_err_ = drain_err;
619
2/4
✗ Branch 116 → 112 not taken.
✓ Branch 116 → 117 taken 38 times.
✗ Branch 118 → 114 not taken.
✓ Branch 118 → 119 taken 1 time.
39x for (auto& e : pending_register_)
620 ✗ registered_.push_back(e);
621 39x pending_register_.clear();
622 39x pending_cancel_.clear();
623 39x }
624
3/4
✓ Branch 126 → 122 taken 1 time.
✓ Branch 126 → 127 taken 38 times.
✗ Branch 128 → 124 not taken.
✓ Branch 128 → 129 taken 1 time.
40x for (auto& e : registered_)
625
1/2
✓ Branch 123 → 124 taken 1 time.
✗ Branch 125 → 126 not taken.
1x sched_.on_completion(e.op, drain_err, 0);
626 39x registered_.clear();
627 39x }
628
629 } // namespace boost::corosio::detail
630
631 #endif // BOOST_COROSIO_HAS_IOCP
632
633 #endif // BOOST_COROSIO_NATIVE_DETAIL_IOCP_WIN_WAIT_REACTOR_HPP
634