include/boost/corosio/native/detail/kqueue/kqueue_scheduler.hpp

98.2% Lines (166/169) 100.0% List of functions (18/18) 75.4% Branches (86/114)
kqueue_scheduler.hpp
f(x) Functions (18)
Function Calls Lines Branches Blocks
boost::corosio::detail::kqueue_set(kevent64_s&, int, short, unsigned short, unsigned int, void*) :61 68278x 100.0% – 100.0% boost::corosio::detail::kqueue_udata(kevent64_s const&) :80 76044x 100.0% – 100.0% boost::corosio::detail::kqueue_call(int, kevent64_s const*, int, kevent64_s*, int, timespec const*) :91 165304x 100.0% 87.5% 85.0% boost::corosio::detail::kqueue_scheduler::register_signal_reader(int) :228 78x 100.0% – 100.0% boost::corosio::detail::kqueue_scheduler::kqueue_scheduler(boost::capy::execution_context&, int) :251 3612x 100.0% 58.8% 100.0% boost::corosio::detail::kqueue_scheduler::kqueue_scheduler(boost::capy::execution_context&, int)::'lambda'(void*)::__invoke(void*) :277 7182x 100.0% – 100.0% boost::corosio::detail::kqueue_scheduler::kqueue_scheduler(boost::capy::execution_context&, int)::'lambda'(void*)::operator void (*)(void*)() const :277 1803x 100.0% – 100.0% boost::corosio::detail::kqueue_scheduler::kqueue_scheduler(boost::capy::execution_context&, int)::'lambda'(void*)::operator()(void*) const :277 7182x 100.0% – 100.0% boost::corosio::detail::kqueue_scheduler::~kqueue_scheduler() :284 5409x 100.0% 50.0% 100.0% boost::corosio::detail::kqueue_scheduler::shutdown() :291 1803x 100.0% 50.0% 100.0% boost::corosio::detail::kqueue_scheduler::configure_reactor(unsigned int, unsigned int, unsigned int, unsigned int) :300 27x 100.0% – 100.0% boost::corosio::detail::kqueue_add_filter(int, int, short, void*) :313 22157x 100.0% 66.7% 77.0% boost::corosio::detail::kqueue_scheduler::register_descriptor(int, boost::corosio::detail::reactor_descriptor_state*) const :324 15637x 100.0% 100.0% 100.0% boost::corosio::detail::kqueue_scheduler::ensure_write_registered(int, boost::corosio::detail::reactor_descriptor_state*) const :348 7023x 100.0% 100.0% 83.0% boost::corosio::detail::kqueue_scheduler::deregister_descriptor(int) const :364 15550x 100.0% – 100.0% boost::corosio::detail::kqueue_scheduler::interrupt_reactor() const :377 15277x 100.0% 100.0% 100.0% boost::corosio::detail::kqueue_scheduler::calculate_timeout(long) const :401 57527x 95.0% 87.5% 90.0% boost::corosio::detail::kqueue_scheduler::run_task(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&, boost::corosio::detail::reactor_scheduler_context&, long) :431 112075x 96.0% 81.6% 92.0%
Line Branch TLA Hits Source Code
1 //
2 // Copyright (c) 2026 Steve Gerbino
3 // Copyright (c) 2026 Michael Vandeberg
4 //
5 // Distributed under the Boost Software License, Version 1.0. (See accompanying
6 // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
7 //
8 // Official repository: https://github.com/cppalliance/corosio
9 //
10
11 #ifndef BOOST_COROSIO_NATIVE_DETAIL_KQUEUE_KQUEUE_SCHEDULER_HPP
12 #define BOOST_COROSIO_NATIVE_DETAIL_KQUEUE_KQUEUE_SCHEDULER_HPP
13
14 #include <boost/corosio/detail/platform.hpp>
15
16 #if BOOST_COROSIO_HAS_KQUEUE
17
18 #include <boost/corosio/detail/config.hpp>
19 #include <boost/capy/ex/execution_context.hpp>
20
21 #include <boost/corosio/native/detail/reactor/reactor_scheduler.hpp>
22 #include <boost/corosio/native/detail/reactor/reactor_signal_pipe.hpp>
23
24 #include <boost/corosio/native/detail/kqueue/kqueue_traits.hpp>
25 #include <boost/corosio/detail/timer_service.hpp>
26 #include <boost/corosio/native/detail/make_err.hpp>
27
28 #include <boost/corosio/detail/except.hpp>
29
30 #include <atomic>
31 #include <chrono>
32 #include <cstdint>
33 #include <limits>
34 #include <mutex>
35 #include <vector>
36
37 #include <errno.h>
38 #include <fcntl.h>
39 #include <sys/event.h>
40 #include <sys/time.h>
41 #include <unistd.h>
42
43 namespace boost::corosio::detail {
44
45 struct kqueue_op;
46
47 /* Every call on a kqueue goes through kqueue_call, because Darwin
48 needs kevent64 for one case: a poll with a zero timeout that finds
49 nothing ready. kevent with a zero timespec takes ~12us there; only
50 KEVENT_FLAG_IMMEDIATE returns at once. A kqueue accepts a single
51 flavor of call, so the registrations use kevent64 as well.
52 */
53 #if defined(__APPLE__)
54 using kqueue_event = struct kevent64_s;
55 #else
56 using kqueue_event = struct kevent;
57 #endif
58
59 /// Fill @p ev, as EV_SET does.
60 inline void
61 68278x kqueue_set(
62 kqueue_event& ev,
63 int ident,
64 short filter,
65 unsigned short flags,
66 unsigned fflags,
67 void* udata) noexcept
68 {
69 #if defined(__APPLE__)
70 68278x EV_SET64(
71 &ev, static_cast<std::uint64_t>(ident), filter, flags, fflags, 0,
72 reinterpret_cast<std::uint64_t>(udata), 0, 0);
73 #else
74 EV_SET(&ev, static_cast<uintptr_t>(ident), filter, flags, fflags, 0, udata);
75 #endif
76 68278x }
77
78 /// The udata stored by kqueue_set.
79 inline void*
80 76044x kqueue_udata(kqueue_event const& ev) noexcept
81 {
82 #if defined(__APPLE__)
83 76044x return reinterpret_cast<void*>(static_cast<std::uintptr_t>(ev.udata));
84 #else
85 return ev.udata;
86 #endif
87 }
88
89 /// Call kevent on @p kq; a zero @p ts polls without waiting.
90 inline int
91 165304x kqueue_call(
92 int kq,
93 kqueue_event const* changes,
94 int nchanges,
95 kqueue_event* events,
96 int nevents,
97 struct timespec const* ts) noexcept
98 {
99 #if defined(__APPLE__)
100 165304x unsigned flags = 0;
101
6/6
✓ Branch 0 taken 107592 times.
✓ Branch 1 taken 57712 times.
✓ Branch 2 taken 106650 times.
✓ Branch 3 taken 942 times.
✓ Branch 4 taken 51445 times.
✓ Branch 5 taken 55205 times.
165304x if (ts && ts->tv_sec == 0 && ts->tv_nsec == 0)
102 {
103 55205x flags = KEVENT_FLAG_IMMEDIATE;
104 55205x ts = nullptr;
105 55205x }
106
1/2
✓ Branch 0 taken 165304 times.
✗ Branch 1 not taken.
165304x return ::kevent64(kq, changes, nchanges, events, nevents, flags, ts);
107 #else
108 return ::kevent(kq, changes, nchanges, events, nevents, ts);
109 #endif
110 }
111
112 /** macOS/BSD scheduler using kqueue for I/O multiplexing.
113
114 This scheduler implements the scheduler interface using the BSD kqueue
115 API for efficient I/O event notification. It uses a single reactor model
116 where one thread runs kevent() while other threads
117 wait on a condition variable for handler work. This design provides:
118
119 - Handler parallelism: N posted handlers can execute on N threads
120 - No thundering herd: condition_variable wakes exactly one thread
121 - IOCP parity: Behavior matches Windows I/O completion port semantics
122
123 When threads call run(), they first try to execute queued handlers.
124 If the queue is empty and no reactor is running, one thread becomes
125 the reactor and runs kevent(). Other threads wait on a condition
126 variable until handlers are available.
127
128 kqueue uses EV_CLEAR for edge-triggered semantics (equivalent to
129 epoll's EPOLLET). A descriptor gets EVFILT_READ at registration and
130 EVFILT_WRITE when a write-direction operation first parks, as asio
131 does; both stay registered until the descriptor is closed.
132
133 @par Thread Safety
134 All public member functions are thread-safe.
135 */
136 class BOOST_COROSIO_DECL kqueue_scheduler final : public reactor_scheduler
137 {
138 public:
139 /** Construct the scheduler.
140
141 Creates a kqueue file descriptor via kqueue(), sets
142 close-on-exec, and registers EVFILT_USER for reactor
143 interruption. On failure the kqueue fd is closed before
144 throwing.
145
146 @param ctx Reference to the owning execution_context.
147 @param concurrency_hint Hint for expected thread count (unused).
148
149 @throws std::system_error if kqueue() fails, if setting
150 FD_CLOEXEC on the kqueue fd fails, or if registering
151 the EVFILT_USER event fails. The error code contains
152 the errno from the failed syscall.
153 */
154 kqueue_scheduler(capy::execution_context& ctx, int concurrency_hint = -1);
155
156 /** Destructor.
157
158 Closes the kqueue file descriptor if valid. Does not throw.
159 */
160 ~kqueue_scheduler();
161
162 kqueue_scheduler(kqueue_scheduler const&) = delete;
163 kqueue_scheduler& operator=(kqueue_scheduler const&) = delete;
164
165 /// Shut down the scheduler, draining pending operations.
166 void shutdown() override;
167
168 /// Apply runtime configuration, resizing the event buffer.
169 void configure_reactor(
170 unsigned max_events,
171 unsigned budget_init,
172 unsigned budget_max,
173 unsigned unassisted) override;
174
175 /** Return the kqueue file descriptor.
176
177 Used by socket services to register file descriptors
178 for I/O event notification.
179
180 @return The kqueue file descriptor.
181 */
182 int kq_fd() const noexcept
183 {
184 return kq_fd_;
185 }
186
187 /** Register a descriptor for persistent monitoring.
188
189 Adds EVFILT_READ (EV_CLEAR) for @a fd and stores @a desc in the
190 kevent udata field so that the reactor can dispatch events to
191 the correct reactor_descriptor_state. EVFILT_WRITE is added later
192 by ensure_write_registered.
193
194 @param fd The file descriptor to register.
195 @param desc Pointer to the caller-owned reactor_descriptor_state.
196
197 @return The error if kevent(EV_ADD) fails, otherwise a default
198 constructed error code.
199 */
200 std::error_code
201 register_descriptor(int fd, reactor_descriptor_state* desc) const;
202
203 /** Add EVFILT_WRITE for @a fd if it is not registered yet.
204
205 Called with `desc->mutex` held, when a write-direction
206 operation is about to park.
207
208 @param fd The registered file descriptor.
209 @param desc Its reactor_descriptor_state.
210
211 @return The kernel's refusal of the filter, otherwise a default
212 constructed error code.
213 */
214 std::error_code ensure_write_registered(
215 int fd, reactor_descriptor_state* desc) const noexcept;
216
217 /** Deregister a persistently registered descriptor.
218
219 Issues kevent(EV_DELETE) for both EVFILT_READ and EVFILT_WRITE.
220 Errors are silently ignored because the fd may already be
221 closed and kqueue automatically removes closed descriptors.
222
223 @param fd The file descriptor to deregister.
224 */
225 void deregister_descriptor(int fd) const;
226
227 /// Watch the read end of the POSIX signal self-pipe (see scheduler.hpp).
228 78x [[nodiscard]] std::error_code register_signal_reader(int read_fd) override
229 {
230 78x return register_descriptor(read_fd, signal_pipe_reader_.arm());
231 }
232
233 private:
234 void run_task(lock_type& lock, context_type& ctx, long timeout_us) override;
235 void interrupt_reactor() const override;
236 long calculate_timeout(long requested_timeout_us) const;
237
238 int kq_fd_;
239
240 // Watches the global signal self-pipe's read end (armed lazily by
241 // register_signal_reader on the first signal registration).
242 reactor_signal_pipe_reader signal_pipe_reader_;
243
244 // EVFILT_USER idempotency
245 1806x mutable std::atomic<bool> user_event_armed_{false};
246
247 // Event buffer sized from max_events_per_poll_.
248 std::vector<kqueue_event> event_buffer_;
249 };
250
251 5418x inline kqueue_scheduler::kqueue_scheduler(capy::execution_context& ctx, int)
252 1806x : kq_fd_(-1)
253
1/2
✓ Branch 0 taken 1806 times.
✗ Branch 1 not taken.
1806x , event_buffer_(max_events_per_poll_)
254 3612x {
255
1/2
✓ Branch 0 taken 1806 times.
✗ Branch 1 not taken.
1806x kq_fd_ = ::kqueue();
256
2/2
✓ Branch 0 taken 1 time.
✓ Branch 1 taken 1805 times.
1806x if (kq_fd_ < 0)
257
2/4
✓ Branch 0 taken 1 time.
✗ Branch 1 not taken.
✗ Branch 2 not taken.
✓ Branch 3 taken 1 time.
1x detail::throw_system_error(make_err(errno), "kqueue");
258
259
3/4
✓ Branch 0 taken 1805 times.
✗ Branch 1 not taken.
✓ Branch 2 taken 1 time.
✓ Branch 3 taken 1804 times.
1805x if (::fcntl(kq_fd_, F_SETFD, FD_CLOEXEC) == -1)
260 {
261
1/2
✓ Branch 0 taken 1 time.
✗ Branch 1 not taken.
1x int errn = errno;
262
1/2
✓ Branch 0 taken 1 time.
✗ Branch 1 not taken.
1x ::close(kq_fd_);
263
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 1 time.
1x detail::throw_system_error(make_err(errn), "fcntl (kqueue FD_CLOEXEC)");
264 }
265
266 kqueue_event ev;
267 1804x kqueue_set(ev, 0, EVFILT_USER, EV_ADD | EV_CLEAR, 0, nullptr);
268
2/2
✓ Branch 0 taken 1 time.
✓ Branch 1 taken 1803 times.
1804x if (kqueue_call(kq_fd_, &ev, 1, nullptr, 0, nullptr) < 0)
269 {
270
1/2
✓ Branch 0 taken 1 time.
✗ Branch 1 not taken.
1x int errn = errno;
271
1/2
✓ Branch 0 taken 1 time.
✗ Branch 1 not taken.
1x ::close(kq_fd_);
272
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 1 time.
1x detail::throw_system_error(make_err(errn), "kevent (EVFILT_USER)");
273 }
274
275
1/2
✓ Branch 0 taken 1803 times.
✗ Branch 1 not taken.
1803x timer_svc_ = &get_timer_service(ctx, *this);
276
2/4
✓ Branch 0 taken 1680 times.
✗ Branch 1 not taken.
✓ Branch 2 taken 1680 times.
✗ Branch 3 not taken.
3606x timer_svc_->set_on_earliest_changed(
277 8985x timer_service::callback(this, [](void* p) {
278 7182x static_cast<kqueue_scheduler*>(p)->interrupt_reactor();
279 7182x }));
280
281 1803x completed_ops_.push(&task_op_);
282 3612x }
283
284 5409x inline kqueue_scheduler::~kqueue_scheduler()
285 3606x {
286
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 1803 times.
1803x if (kq_fd_ >= 0)
287
1/2
✓ Branch 0 taken 1803 times.
✗ Branch 1 not taken.
1803x ::close(kq_fd_);
288 5409x }
289
290 inline void
291 1803x kqueue_scheduler::shutdown()
292 {
293 1803x shutdown_drain();
294
295
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 1803 times.
1803x if (kq_fd_ >= 0)
296 1803x interrupt_reactor();
297 1803x }
298
299 inline void
300 27x kqueue_scheduler::configure_reactor(
301 unsigned max_events,
302 unsigned budget_init,
303 unsigned budget_max,
304 unsigned unassisted)
305 {
306 27x reactor_scheduler::configure_reactor(
307 27x max_events, budget_init, budget_max, unassisted);
308 27x event_buffer_.resize(max_events_per_poll_);
309 27x }
310
311 /// Add one edge-triggered filter for @p fd; returns 0 or the errno.
312 inline int
313 22157x kqueue_add_filter(int kq, int fd, short filter, void* udata) noexcept
314 {
315 kqueue_event ch;
316 22157x kqueue_set(ch, fd, filter, EV_ADD | EV_CLEAR | EV_RECEIPT, 0, udata);
317 kqueue_event receipt;
318
2/2
✓ Branch 0 taken 22148 times.
✓ Branch 1 taken 9 times.
22157x if (kqueue_call(kq, &ch, 1, &receipt, 1, nullptr) < 0)
319
1/2
✓ Branch 0 taken 9 times.
✗ Branch 1 not taken.
9x return errno;
320
1/2
✓ Branch 0 taken 22148 times.
✗ Branch 1 not taken.
22148x return (receipt.flags & EV_ERROR) ? static_cast<int>(receipt.data) : 0;
321 22157x }
322
323 inline std::error_code
324 15637x kqueue_scheduler::register_descriptor(
325 int fd, reactor_descriptor_state* desc) const
326 {
327 15637x int const err = kqueue_add_filter(kq_fd_, fd, EVFILT_READ, desc);
328 // EINVAL/ENODEV: a device with no kqfilter. As on epoll's EPERM,
329 // adopt it unwatched.
330
2/2
✓ Branch 0 taken 3 times.
✓ Branch 1 taken 15634 times.
15637x bool const unpollable = err == EINVAL || err == ENODEV;
331
4/4
✓ Branch 0 taken 11 times.
✓ Branch 1 taken 15626 times.
✓ Branch 2 taken 3 times.
✓ Branch 3 taken 8 times.
15637x if (err != 0 && !unpollable)
332 8x return make_err(err);
333 15629x desc->registered_events = unpollable ? 0 : reactor_event_read;
334 15629x desc->unpollable = unpollable;
335 15629x desc->fd = fd;
336 15629x desc->scheduler_ = this;
337 15629x desc->mutex.set_enabled(reactor_io_locking_);
338 15629x desc->ready_events_.store(0, std::memory_order_relaxed);
339
340 15629x conditionally_enabled_mutex::scoped_lock lock(desc->mutex);
341 15629x desc->impl_ref_.reset();
342 15629x desc->read_ready = false;
343 15629x desc->write_ready = false;
344 15629x return {};
345 15637x }
346
347 inline std::error_code
348 7023x kqueue_scheduler::ensure_write_registered(
349 int fd, reactor_descriptor_state* desc) const noexcept
350 {
351 // Added on first need, as asio does: adoption never asks for a
352 // write filter a read-only device would refuse. Adding a filter
353 // reports the current state, so an fd already writable fires at
354 // once and no edge is lost.
355
2/2
✓ Branch 0 taken 506 times.
✓ Branch 1 taken 6517 times.
7023x if (desc->registered_events & reactor_event_write)
356 506x return {};
357
2/2
✓ Branch 0 taken 1 time.
✓ Branch 1 taken 6516 times.
6517x if (int err = kqueue_add_filter(kq_fd_, fd, EVFILT_WRITE, desc))
358 1x return make_err(err);
359 6516x desc->registered_events |= reactor_event_write;
360 6516x return {};
361 7023x }
362
363 inline void
364 15550x kqueue_scheduler::deregister_descriptor(int fd) const
365 {
366 // EV_RECEIPT reports each change separately, so a never-added
367 // EVFILT_WRITE (ENOENT) does not stop the EVFILT_READ delete.
368 kqueue_event changes[2];
369 15550x kqueue_set(changes[0], fd, EVFILT_READ, EV_DELETE | EV_RECEIPT, 0, nullptr);
370 15550x kqueue_set(
371 15550x changes[1], fd, EVFILT_WRITE, EV_DELETE | EV_RECEIPT, 0, nullptr);
372 kqueue_event receipts[2];
373 15550x kqueue_call(kq_fd_, changes, 2, receipts, 2, nullptr);
374 15550x }
375
376 inline void
377 15277x kqueue_scheduler::interrupt_reactor() const
378 {
379 15277x bool expected = false;
380
2/2
✓ Branch 0 taken 2063 times.
✓ Branch 1 taken 13214 times.
15277x if (user_event_armed_.compare_exchange_strong(
381 expected, true, std::memory_order_acq_rel,
382 std::memory_order_acquire))
383 {
384 kqueue_event ev;
385 13214x kqueue_set(ev, 0, EVFILT_USER, 0, NOTE_TRIGGER, nullptr);
386
2/2
✓ Branch 0 taken 13212 times.
✓ Branch 1 taken 2 times.
13214x if (kqueue_call(kq_fd_, &ev, 1, nullptr, 0, nullptr) < 0)
387 {
388 // The flag is what coalesces later interrupts into a
389 // trigger already queued on the kqueue; a kevent that
390 // failed queued nothing, so leaving it armed would swallow
391 // every interrupt that follows. Disarming keeps the cost to
392 // the interrupts already in flight -- the next one arms and
393 // triggers again, instead of every one after this
394 // coalescing into a trigger that does not exist.
395 2x user_event_armed_.store(false, std::memory_order_release);
396 2x }
397 13214x }
398 15277x }
399
400 inline long
401 57527x kqueue_scheduler::calculate_timeout(long requested_timeout_us) const
402 {
403
1/2
✓ Branch 0 taken 57527 times.
✗ Branch 1 not taken.
57527x if (requested_timeout_us == 0)
404 ✗ return 0;
405
406 57527x auto nearest = timer_svc_->nearest_expiry();
407
2/2
✓ Branch 0 taken 5152 times.
✓ Branch 1 taken 52375 times.
57527x if (nearest == timer_service::time_point::max())
408 5152x return requested_timeout_us;
409
410 52375x auto now = std::chrono::steady_clock::now();
411
2/2
✓ Branch 0 taken 154 times.
✓ Branch 1 taken 52221 times.
52375x if (nearest <= now)
412 154x return 0;
413
414 52221x auto timer_timeout_us =
415 52221x std::chrono::duration_cast<std::chrono::microseconds>(nearest - now)
416 52221x .count();
417
418 52221x constexpr auto long_max =
419 static_cast<long long>((std::numeric_limits<long>::max)());
420 52221x auto capped_timer_us = std::min(
421 52221x std::max(timer_timeout_us, static_cast<long long>(0)), long_max);
422
423
2/2
✓ Branch 0 taken 2 times.
✓ Branch 1 taken 52219 times.
52221x if (requested_timeout_us < 0)
424 52219x return static_cast<long>(capped_timer_us);
425
426 2x return static_cast<long>(std::min(
427 2x static_cast<long long>(requested_timeout_us), capped_timer_us));
428 57527x }
429
430 inline void
431 112075x kqueue_scheduler::run_task(lock_type& lock, context_type& ctx, long timeout_us)
432 {
433 112075x long effective_timeout_us =
434
2/2
✓ Branch 0 taken 54548 times.
✓ Branch 1 taken 57527 times.
112075x task_interrupted_ ? 0 : calculate_timeout(timeout_us);
435
436
2/2
✓ Branch 0 taken 54546 times.
✓ Branch 1 taken 57529 times.
112075x if (lock.owns_lock())
437 57529x lock.unlock();
438
439 112075x task_cleanup on_exit{this, &lock, ctx};
440
441 struct timespec ts;
442 112075x struct timespec* ts_ptr = nullptr;
443
2/2
✓ Branch 0 taken 4984 times.
✓ Branch 1 taken 107091 times.
112075x if (effective_timeout_us >= 0)
444 {
445 107091x ts.tv_sec = effective_timeout_us / 1000000;
446 107091x ts.tv_nsec = (effective_timeout_us % 1000000) * 1000;
447 107091x ts_ptr = &ts;
448 107091x }
449
450 112075x int nev = kqueue_call(
451 112075x kq_fd_, nullptr, 0, event_buffer_.data(),
452 112075x static_cast<int>(event_buffer_.size()), ts_ptr);
453
1/2
✓ Branch 0 taken 112075 times.
✗ Branch 1 not taken.
112075x int saved_errno = errno;
454
455
4/4
✓ Branch 0 taken 2 times.
✓ Branch 1 taken 112073 times.
✓ Branch 2 taken 1 time.
✓ Branch 3 taken 1 time.
112075x if (nev < 0 && saved_errno != EINTR)
456
1/2
✓ Branch 0 taken 1 time.
✗ Branch 1 not taken.
1x detail::throw_system_error(make_err(saved_errno), "kevent");
457
458 112074x ready_queue local_ops;
459
460
2/2
✓ Branch 0 taken 112074 times.
✓ Branch 1 taken 87453 times.
199527x for (int i = 0; i < nev; ++i)
461 {
462
2/2
✓ Branch 0 taken 11409 times.
✓ Branch 1 taken 76044 times.
87453x if (event_buffer_[i].filter == EVFILT_USER)
463 {
464 11409x user_event_armed_.store(false, std::memory_order_release);
465 11409x continue;
466 }
467
468 76044x auto* desc =
469 static_cast<reactor_descriptor_state*>(
470 76044x kqueue_udata(event_buffer_[i]));
471
1/2
✓ Branch 0 taken 76044 times.
✗ Branch 1 not taken.
76044x if (!desc)
472 ✗ continue;
473
474 76044x std::uint32_t ready = 0;
475
476
2/2
✓ Branch 0 taken 8147 times.
✓ Branch 1 taken 67897 times.
76044x if (event_buffer_[i].filter == EVFILT_READ)
477 67897x ready |= reactor_event_read;
478
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 8147 times.
8147x else if (event_buffer_[i].filter == EVFILT_WRITE)
479 8147x ready |= reactor_event_write;
480
481
1/2
✓ Branch 0 taken 76044 times.
✗ Branch 1 not taken.
76044x if (event_buffer_[i].flags & EV_ERROR)
482 ✗ ready |= reactor_event_error;
483
484
2/2
✓ Branch 0 taken 75879 times.
✓ Branch 1 taken 165 times.
76044x if (event_buffer_[i].flags & EV_EOF)
485 {
486
2/2
✓ Branch 0 taken 32 times.
✓ Branch 1 taken 133 times.
165x if (event_buffer_[i].filter == EVFILT_READ)
487 133x ready |= reactor_event_read;
488
2/2
✓ Branch 0 taken 102 times.
✓ Branch 1 taken 63 times.
165x if (event_buffer_[i].fflags != 0)
489 63x ready |= reactor_event_error;
490 165x }
491
492 76044x desc->add_ready_events(ready);
493
494 76044x bool expected = false;
495
2/2
✓ Branch 0 taken 75 times.
✓ Branch 1 taken 75969 times.
76044x if (desc->is_enqueued_.compare_exchange_strong(
496 expected, true, std::memory_order_acq_rel,
497 std::memory_order_acquire))
498 {
499 75969x local_ops.push(desc);
500 75969x }
501 76044x }
502
503
1/2
✓ Branch 0 taken 112074 times.
✗ Branch 1 not taken.
112074x timer_svc_->process_expired();
504
505
1/2
✓ Branch 0 taken 112074 times.
✗ Branch 1 not taken.
112074x lock.lock();
506
507 112074x completed_ops_.splice(local_ops);
508 112075x }
509
510 } // namespace boost::corosio::detail
511
512 #endif // BOOST_COROSIO_HAS_KQUEUE
513
514 #endif // BOOST_COROSIO_NATIVE_DETAIL_KQUEUE_KQUEUE_SCHEDULER_HPP
515