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

97.6% Lines (163/167) 100.0% List of functions (28/28) 77.8% Branches (49/63)
win_object_handle_service.hpp
f(x) Functions (28)
Function Calls Lines Branches Blocks
boost::corosio::detail::win_object_handle_state::win_object_handle_state(boost::corosio::detail::win_scheduler&, boost::corosio::detail::win_state_list<boost::corosio::detail::win_object_handle_state>&) :59 2824x 100.0% – 100.0% boost::corosio::detail::win_object_handle_state::~win_object_handle_state() :69 2824x 100.0% 100.0% 100.0% boost::corosio::detail::win_object_handle_state::native_handle() const :83 8261x 100.0% – 100.0% boost::corosio::detail::win_object_handle_state::wait(boost::capy::continuation&, boost::capy::executor_ref, std::stop_token, std::error_code*) :93 3222x 89.2% 68.2% 54.1% boost::corosio::detail::win_object_handle_state::cancel() :160 2201x 100.0% – 100.0% boost::corosio::detail::win_object_handle_state::close_handle() :165 8264x 100.0% 100.0% 100.0% boost::corosio::detail::win_object_handle_state::release() :175 201x 100.0% – 100.0% boost::corosio::detail::win_object_handle_state::assign(unsigned long long) :183 2823x 100.0% 100.0% 100.0% boost::corosio::detail::win_object_handle_state::wait_op::wait_op() :204 2824x 100.0% – 100.0% boost::corosio::detail::win_object_handle_state::wait_op::do_cancel_impl(boost::corosio::detail::overlapped_op*) :209 2x 100.0% – 100.0% boost::corosio::detail::win_object_handle_state::wait_op::do_complete(void*, boost::corosio::detail::scheduler_op*, unsigned int, unsigned int) :214 3021x 100.0% 90.0% 89.7% boost::corosio::detail::win_object_handle_state::set_idle() :256 3021x 100.0% – 58.3% boost::corosio::detail::win_object_handle_state::on_signaled(_TP_CALLBACK_INSTANCE*, void*, _TP_WAIT*, unsigned long) :263 1885x 100.0% 50.0% 69.6% boost::corosio::detail::win_object_handle_state::claim(unsigned long long) :272 3020x 100.0% – 71.4% boost::corosio::detail::win_object_handle_state::rundown() :282 10668x 100.0% 87.5% 73.3% boost::corosio::detail::win_object_handle_impl::win_object_handle_impl(std::shared_ptr<boost::corosio::detail::win_object_handle_state>) :319 2824x 100.0% – 100.0% boost::corosio::detail::win_object_handle_impl::close_internal() :325 2823x 100.0% 50.0% 100.0% boost::corosio::detail::win_object_handle_impl::get_internal() const :334 8263x 100.0% – 100.0% boost::corosio::detail::win_object_handle_impl::wait(boost::capy::continuation&, boost::capy::executor_ref, std::stop_token, std::error_code*) :339 3222x 100.0% 100.0% 81.8% boost::corosio::detail::win_object_handle_impl::native_handle() const :348 8261x 100.0% – 100.0% boost::corosio::detail::win_object_handle_impl::release_handle() :353 201x 100.0% – 100.0% boost::corosio::detail::win_object_handle_impl::cancel() :358 2201x 100.0% – 100.0% boost::corosio::detail::win_object_handle_service::win_object_handle_service(boost::capy::execution_context&) :369 2824x 100.0% 100.0% 83.3% boost::corosio::detail::win_object_handle_service::construct() :374 2824x 100.0% 100.0% 85.7% boost::corosio::detail::win_object_handle_service::destroy(boost::corosio::io_object::implementation*) :381 2823x 100.0% 50.0% 100.0% boost::corosio::detail::win_object_handle_service::close(boost::corosio::io_object::handle&) :387 5440x 100.0% – 100.0% boost::corosio::detail::win_object_handle_service::shutdown() :394 2824x 100.0% – 100.0% boost::corosio::detail::win_object_handle_service::assign_object_handle(boost::corosio::win_object_handle::implementation&, unsigned long long) :399 2824x 100.0% 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_OBJECT_HANDLE_SERVICE_HPP
11 #define BOOST_COROSIO_NATIVE_DETAIL_IOCP_WIN_OBJECT_HANDLE_SERVICE_HPP
12
13 #include <boost/corosio/detail/platform.hpp>
14
15 #if BOOST_COROSIO_HAS_IOCP
16
17 #include <boost/corosio/detail/config.hpp>
18 #include <boost/corosio/detail/dispatch_coro.hpp>
19 #include <boost/corosio/detail/win_handle_service.hpp>
20 #include <boost/corosio/native/detail/iocp/win_overlapped_handle.hpp>
21 #include <boost/corosio/native/detail/iocp/win_validate_handle.hpp>
22 #include <boost/capy/error.hpp>
23
24 #include <atomic>
25 #include <cstdint>
26 #include <memory>
27 #include <mutex>
28
29 /* win_object_handle on the Windows thread pool.
30
31 One PTP_WAIT per state, created at the first wait() and re-armed per
32 wait. The claim word (phase in the low two bits, a generation above)
33 decides which of the pool callback and a cancellation delivers the
34 one op: whoever moves it from armed to claimed posts the op.
35
36 Cancellation never cancels a queued callback. By the time a callback
37 is queued the kernel has already satisfied the wait -- consumed an
38 auto-reset event's signal, a semaphore unit -- and the callback is
39 the only witness to that. So every rundown (cancel, stop token,
40 close, release, assign, destroy, shutdown) drains instead:
41 SetThreadpoolWait(nullptr) stops an unsatisfied wait, then
42 WaitForThreadpoolWaitCallbacks(FALSE) lets a queued callback finish
43 and claim, and only then does the rundown try to claim for
44 cancellation.
45
46 arm_mutex_ orders arming against the disarm step, and disarmed_gen_
47 records which generation a rundown disarmed, so a rundown that
48 disarmed before wait() armed is never followed by that arm. The
49 callback never takes the mutex, which keeps the drain deadlock-free.
50 */
51
52 namespace boost::corosio::detail {
53
54 class win_object_handle_state
55 : public intrusive_list<win_object_handle_state>::node
56 , public std::enable_shared_from_this<win_object_handle_state>
57 {
58 public:
59 2824x win_object_handle_state(
60 win_scheduler& sched,
61 win_state_list<win_object_handle_state>& list) noexcept
62 5648x : sched_(sched)
63 2824x , list_(list)
64 {
65 2824x op_.self = this;
66 2824x list_.add(*this);
67 2824x }
68
69 2824x ~win_object_handle_state()
70 {
71
2/2
✓ Branch 2 → 3 taken 2818 times.
✓ Branch 2 → 6 taken 6 times.
2824x if (tp_wait_)
72 {
73 2818x ::SetThreadpoolWait(tp_wait_, nullptr, nullptr);
74 2818x ::WaitForThreadpoolWaitCallbacks(tp_wait_, FALSE);
75 2818x ::CloseThreadpoolWait(tp_wait_);
76 }
77 2824x list_.remove(*this);
78 2824x }
79
80 win_object_handle_state(win_object_handle_state const&) = delete;
81 win_object_handle_state& operator=(win_object_handle_state const&) = delete;
82
83 8261x HANDLE native_handle() const noexcept
84 {
85 8261x return handle_;
86 }
87
88 bool is_open() const noexcept
89 {
90 return handle_ != INVALID_HANDLE_VALUE;
91 }
92
93 3222x std::coroutine_handle<> wait(
94 capy::continuation& cont,
95 capy::executor_ref ex,
96 std::stop_token token,
97 std::error_code* ec)
98 {
99 3222x std::uint64_t const cur = state_.load(std::memory_order_acquire);
100
2/2
✓ Branch 17 → 18 taken 201 times.
✓ Branch 17 → 21 taken 3021 times.
3222x if ((cur & phase_mask) != phase_idle)
101 {
102 // Complete immediately: the returned handle is resumed by
103 // symmetric transfer on the caller's own executor.
104
1/2
✓ Branch 18 → 19 taken 201 times.
✗ Branch 18 → 20 not taken.
201x if (ec)
105 201x *ec = std::make_error_code(std::errc::operation_in_progress);
106 201x return cont.h;
107 }
108
109 3021x std::uint64_t const gen = (cur >> 2) + 1;
110 3021x std::uint64_t const armed = (gen << 2) | phase_armed;
111
112 3021x op_.reset();
113
1/1
✓ Branch 22 → 23 taken 3021 times.
3021x op_.keep_alive = shared_from_this();
114 // The op may complete and the object be destroyed on another
115 // thread before wait() returns; this pin outlives both.
116 3021x auto const pin = op_.keep_alive;
117 3021x op_.user_cont = &cont;
118 3021x op_.h = cont.h;
119 3021x op_.ex = ex;
120 3021x op_.ec_out = ec;
121 3021x op_.bytes_out = nullptr;
122 3021x sched_.work_started();
123
124
2/2
✓ Branch 27 → 28 taken 1 time.
✓ Branch 27 → 52 taken 3020 times.
3021x if (handle_ == INVALID_HANDLE_VALUE)
125 {
126 1x state_.store((gen << 2) | phase_claimed, std::memory_order_release);
127
1/1
✓ Branch 48 → 49 taken 1 time.
1x sched_.on_completion(&op_, ERROR_INVALID_HANDLE, 0);
128 1x return std::noop_coroutine();
129 }
130
131
2/2
✓ Branch 52 → 53 taken 2818 times.
✓ Branch 52 → 80 taken 202 times.
3020x if (!tp_wait_)
132 {
133
1/1
✓ Branch 53 → 54 taken 2818 times.
2818x tp_wait_ = ::CreateThreadpoolWait(&on_signaled, this, nullptr);
134
1/2
✗ Branch 54 → 55 not taken.
✓ Branch 54 → 80 taken 2818 times.
2818x if (!tp_wait_)
135 {
136 ✗ state_.store(
137 ✗ (gen << 2) | phase_claimed, std::memory_order_release);
138 ✗ sched_.on_completion(&op_, ::GetLastError(), 0);
139 ✗ return std::noop_coroutine();
140 }
141 }
142
143 3020x state_.store(armed, std::memory_order_release);
144 // May run the rundown synchronously if the token is already
145 // stopped; that claims the op before it is ever armed.
146 3020x op_.start(token);
147 {
148 // A rundown that disarmed this generation first (from a
149 // stop request on another thread) has claimed, or will
150 // claim, the op; arming now would let a later signal be
151 // consumed with nobody left to report it.
152 3020x std::lock_guard<win_mutex> lock(arm_mutex_);
153
2/4
✓ Branch 117 → 118 taken 3020 times.
✗ Branch 117 → 120 not taken.
✓ Branch 121 → 122 taken 3020 times.
✗ Branch 121 → 123 not taken.
9060x if (state_.load(std::memory_order_acquire) == armed &&
154
1/2
✓ Branch 118 → 119 taken 3020 times.
✗ Branch 118 → 120 not taken.
3020x disarmed_gen_ != gen)
155
1/1
✓ Branch 122 → 123 taken 3020 times.
3020x ::SetThreadpoolWait(tp_wait_, handle_, nullptr);
156 3020x }
157 3020x return std::noop_coroutine();
158 3021x }
159
160 2201x void cancel() noexcept
161 {
162 2201x rundown();
163 2201x }
164
165 8264x void close_handle() noexcept
166 {
167 8264x rundown();
168
2/2
✓ Branch 3 → 4 taken 2618 times.
✓ Branch 3 → 6 taken 5646 times.
8264x if (handle_ != INVALID_HANDLE_VALUE)
169 {
170 2618x ::CloseHandle(handle_);
171 2618x handle_ = INVALID_HANDLE_VALUE;
172 }
173 8264x }
174
175 201x native_handle_type release() noexcept
176 {
177 201x rundown();
178 201x HANDLE h = handle_;
179 201x handle_ = INVALID_HANDLE_VALUE;
180 201x return reinterpret_cast<native_handle_type>(h);
181 }
182
183 2823x std::error_code assign(native_handle_type nh) noexcept
184 {
185 2823x HANDLE h = reinterpret_cast<HANDLE>(nh);
186
2/2
✓ Branch 4 → 5 taken 4 times.
✓ Branch 4 → 6 taken 2819 times.
2823x if (auto ec = validate_object_handle(h))
187 4x return ec;
188 2819x handle_ = h;
189 2819x return {};
190 }
191
192 private:
193 static constexpr std::uint64_t phase_idle = 0;
194 static constexpr std::uint64_t phase_armed = 1;
195 static constexpr std::uint64_t phase_claimed = 2;
196 static constexpr std::uint64_t phase_mask = 3;
197
198 struct wait_op : overlapped_op
199 {
200 win_object_handle_state* self = nullptr;
201 std::shared_ptr<win_object_handle_state> keep_alive;
202 capy::continuation* user_cont = nullptr;
203
204 2824x wait_op() noexcept : overlapped_op(&do_complete)
205 {
206 2824x cancel_func_ = &do_cancel_impl;
207 2824x }
208
209 2x static void do_cancel_impl(overlapped_op* base) noexcept
210 {
211 2x static_cast<wait_op*>(base)->self->rundown();
212 2x }
213
214 3021x static void do_complete(
215 void* owner,
216 scheduler_op* base,
217 std::uint32_t /*bytes*/,
218 std::uint32_t /*error*/)
219 {
220 3021x auto* op = static_cast<wait_op*>(base);
221 3021x auto* self = op->self;
222 3021x auto prevent_premature_destruction = std::move(op->keep_alive);
223
224
2/2
✓ Branch 4 → 5 taken 1 time.
✓ Branch 4 → 8 taken 3020 times.
3021x if (!owner)
225 {
226 1x self->set_idle();
227 1x op->cleanup_only();
228 1x return;
229 }
230
231 3020x op->stop_cb.reset();
232 // A claimed signal is success even if a cancel request raced
233 // it and set the cancelled flag: completion wins.
234
1/2
✓ Branch 9 → 10 taken 3020 times.
✗ Branch 9 → 17 not taken.
3020x if (op->ec_out)
235 {
236
2/2
✓ Branch 10 → 11 taken 1885 times.
✓ Branch 10 → 13 taken 1135 times.
3020x if (op->dwError == 0)
237 1885x *op->ec_out = {};
238
2/2
✓ Branch 13 → 14 taken 1134 times.
✓ Branch 13 → 16 taken 1 time.
1135x else if (op->dwError == ERROR_OPERATION_ABORTED)
239 1134x *op->ec_out = capy::error::canceled;
240 else
241 1x *op->ec_out =
242 1x iocp_make_err(op->dwError, /*accept_path=*/false);
243 }
244
245 3020x capy::continuation* cont = op->user_cont;
246 3020x cont->h = op->h;
247 3020x capy::executor_ref ex = op->ex;
248
249 // Last touch of op_: from here a new wait() may reset it.
250 3020x self->set_idle();
251
2/2
✓ Branch 18 → 19 taken 3020 times.
✓ Branch 19 → 20 taken 3020 times.
3020x dispatch_coro(ex, *cont).resume();
252 3021x }
253 };
254
255 // Back to idle, keeping the generation.
256 3021x void set_idle() noexcept
257 {
258 6042x state_.store(
259 6042x state_.load(std::memory_order_relaxed) & ~phase_mask,
260 std::memory_order_release);
261 3021x }
262
263 1885x static void CALLBACK on_signaled(
264 PTP_CALLBACK_INSTANCE, void* ctx, PTP_WAIT, TP_WAIT_RESULT) noexcept
265 {
266 1885x auto* self = static_cast<win_object_handle_state*>(ctx);
267 1885x std::uint64_t cur = self->state_.load(std::memory_order_acquire);
268
3/6
✓ Branch 17 → 18 taken 1885 times.
✗ Branch 17 → 21 not taken.
✓ Branch 19 → 20 taken 1885 times.
✗ Branch 19 → 21 not taken.
✓ Branch 22 → 23 taken 1885 times.
✗ Branch 22 → 24 not taken.
1885x if ((cur & phase_mask) == phase_armed && self->claim(cur))
269 1885x self->sched_.on_completion(&self->op_, 0, 0);
270 1885x }
271
272 3020x bool claim(std::uint64_t armed) noexcept
273 {
274 6040x return state_.compare_exchange_strong(
275 3020x armed, (armed & ~phase_mask) | phase_claimed,
276 3020x std::memory_order_acq_rel);
277 }
278
279 // Stop an unsatisfied wait, let a queued callback finish, then claim
280 // for cancellation only if nothing else did. Never called from the
281 // pool callback.
282 10668x void rundown() noexcept
283 {
284
2/2
✓ Branch 2 → 3 taken 13 times.
✓ Branch 2 → 4 taken 10655 times.
10668x if (!tp_wait_)
285 13x return;
286 {
287 10655x std::lock_guard<win_mutex> lock(arm_mutex_);
288 10655x ::SetThreadpoolWait(tp_wait_, nullptr, nullptr);
289 10655x disarmed_gen_ = state_.load(std::memory_order_acquire) >> 2;
290 10655x }
291 10655x ::WaitForThreadpoolWaitCallbacks(tp_wait_, FALSE);
292
293 10655x std::uint64_t cur = state_.load(std::memory_order_acquire);
294
5/6
✓ Branch 38 → 39 taken 1135 times.
✓ Branch 38 → 42 taken 9520 times.
✓ Branch 40 → 41 taken 1135 times.
✗ Branch 40 → 42 not taken.
✓ Branch 43 → 44 taken 1135 times.
✓ Branch 43 → 46 taken 9520 times.
10655x if ((cur & phase_mask) == phase_armed && claim(cur))
295 {
296 1135x op_.request_cancel();
297 1135x sched_.on_completion(&op_, ERROR_OPERATION_ABORTED, 0);
298 }
299 }
300
301 win_scheduler& sched_;
302 win_state_list<win_object_handle_state>& list_;
303 HANDLE handle_ = INVALID_HANDLE_VALUE;
304 PTP_WAIT tp_wait_ = nullptr;
305 std::atomic<std::uint64_t> state_{phase_idle};
306 win_mutex arm_mutex_;
307 std::uint64_t disarmed_gen_ = 0; // guarded by arm_mutex_
308 wait_op op_;
309 };
310
311 /** IOCP implementation of @ref win_object_handle. */
312 class win_object_handle_impl final
313 : public win_object_handle::implementation
314 , public intrusive_list<win_object_handle_impl>::node
315 {
316 std::shared_ptr<win_object_handle_state> internal_;
317
318 public:
319 2824x explicit win_object_handle_impl(
320 std::shared_ptr<win_object_handle_state> internal) noexcept
321 2824x : internal_(std::move(internal))
322 {
323 2824x }
324
325 2823x void close_internal() noexcept
326 {
327
1/2
✓ Branch 3 → 4 taken 2823 times.
✗ Branch 3 → 7 not taken.
2823x if (internal_)
328 {
329 2823x internal_->close_handle();
330 2823x internal_.reset();
331 }
332 2823x }
333
334 8263x win_object_handle_state* get_internal() const noexcept
335 {
336 8263x return internal_.get();
337 }
338
339 3222x std::coroutine_handle<> wait(
340 capy::continuation& cont,
341 capy::executor_ref ex,
342 std::stop_token token,
343 std::error_code* ec) override
344 {
345
1/1
✓ Branch 5 → 6 taken 3222 times.
3222x return internal_->wait(cont, ex, std::move(token), ec);
346 }
347
348 8261x native_handle_type native_handle() const noexcept override
349 {
350 8261x return reinterpret_cast<native_handle_type>(internal_->native_handle());
351 }
352
353 201x native_handle_type release_handle() noexcept override
354 {
355 201x return internal_->release();
356 }
357
358 2201x void cancel() noexcept override
359 {
360 2201x internal_->cancel();
361 2201x }
362 };
363
364 /** IOCP service that owns the @ref win_object_handle implementations. */
365 class BOOST_COROSIO_DECL win_object_handle_service final
366 : public object_handle_service
367 {
368 public:
369 2824x explicit win_object_handle_service(capy::execution_context& ctx)
370
2/2
✓ Branch 3 → 4 taken 2824 times.
✓ Branch 4 → 5 taken 2824 times.
2824x : sched_(ctx.use_service<win_scheduler>())
371 {
372 2824x }
373
374 2824x io_object::implementation* construct() override
375 {
376 auto internal =
377
1/1
✓ Branch 3 → 4 taken 2824 times.
2824x std::make_shared<win_object_handle_state>(sched_, reg_.states());
378
1/1
✓ Branch 4 → 5 taken 2824 times.
5648x return reg_.add(new win_object_handle_impl(std::move(internal)));
379 2824x }
380
381 2823x void destroy(io_object::implementation* p) override
382 {
383
1/2
✓ Branch 2 → 3 taken 2823 times.
✗ Branch 2 → 4 not taken.
2823x if (p)
384 2823x reg_.destroy(static_cast<win_object_handle_impl&>(*p));
385 2823x }
386
387 5440x void close(io_object::handle& h) override
388 {
389 5440x static_cast<win_object_handle_impl&>(*h.get())
390 .get_internal()
391 5440x ->close_handle();
392 5440x }
393
394 2824x void shutdown() override
395 {
396 2824x reg_.shutdown();
397 2824x }
398
399 2824x std::error_code assign_object_handle(
400 win_object_handle::implementation& impl,
401 native_handle_type h) override
402 {
403 // The thread-pool callback completes from a foreign thread.
404
2/2
✓ Branch 3 → 4 taken 1 time.
✓ Branch 3 → 5 taken 2823 times.
2824x if (sched_.scheduler_locking_disabled())
405 1x return std::make_error_code(std::errc::operation_not_supported);
406 return static_cast<win_object_handle_impl&>(impl)
407 .get_internal()
408 2823x ->assign(h);
409 }
410
411 private:
412 win_scheduler& sched_;
413 BOOST_COROSIO_MSVC_WARNING_PUSH
414 BOOST_COROSIO_MSVC_WARNING_DISABLE(4251) // detail:: members, dll-interface
415 win_handle_registry<win_object_handle_impl, win_object_handle_state> reg_;
416 BOOST_COROSIO_MSVC_WARNING_POP
417 };
418
419 } // namespace boost::corosio::detail
420
421 #endif // BOOST_COROSIO_HAS_IOCP
422
423 #endif
424