include/boost/corosio/native/detail/reactor/reactor_scheduler.hpp

85.5% Lines (348/407) 90.0% List of functions (45/50) 64.8% Branches (158/244)
reactor_scheduler.hpp
f(x) Functions (50)
Function Calls Lines Branches Blocks
boost::corosio::detail::reactor_find_context(boost::corosio::detail::reactor_scheduler const*) :79 2800563x 85.7% 75.0% 75.0% boost::corosio::detail::reactor_flush_private_work(boost::corosio::detail::reactor_scheduler_context*, std::__1::atomic<long long>&) :91 0 0.0% 0.0% 0.0% boost::corosio::detail::reactor_drain_private_queue(boost::corosio::detail::reactor_scheduler_context*, std::__1::atomic<long long>&, boost::corosio::detail::ready_queue&) :108 8x 57.1% 50.0% 80.0% boost::corosio::detail::reactor_scheduler::~reactor_scheduler() :141 1331x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::inline_budget_initial() const :256 3194x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::scheduler_locking_disabled() const :262 162x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::configure_threading(boost::corosio::detail::scheduler::threading_config) :267 1331x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::reactor_scheduler() :284 1331x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::task_op::task_op() :327 2662x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::task_op::~task_op() :327 2662x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::task_op::operator()() :329 0 0.0% 0.0% boost::corosio::detail::reactor_scheduler::task_op::destroy() :330 0 0.0% 0.0% boost::corosio::detail::reactor_thread_context_guard::reactor_thread_context_guard(boost::corosio::detail::reactor_scheduler const*) :376 6388x 100.0% 50.0% 100.0% boost::corosio::detail::reactor_thread_context_guard::~reactor_thread_context_guard() :384 6388x 71.4% 16.7% 100.0% boost::corosio::detail::reactor_scheduler_context::reactor_scheduler_context(boost::corosio::detail::reactor_scheduler const*, boost::corosio::detail::reactor_scheduler_context*) :396 6388x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::configure_reactor(unsigned int, unsigned int, unsigned int, unsigned int) :410 25x 93.3% 66.7% 62.0% boost::corosio::detail::reactor_scheduler::reset_inline_budget() const :437 310835x 50.0% 21.4% 33.0% boost::corosio::detail::reactor_scheduler::try_consume_inline_budget() const :468 1142736x 90.0% 66.7% 87.0% boost::corosio::detail::reactor_scheduler::post(std::__1::coroutine_handle<void>) const :484 3684x 100.0% 53.0% boost::corosio::detail::reactor_scheduler::post(std::__1::coroutine_handle<void>) const::post_handler::post_handler(std::__1::coroutine_handle<void>) :490 7359x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__1::coroutine_handle<void>) const::post_handler::~post_handler() :491 7359x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__1::coroutine_handle<void>) const::post_handler::operator()() :493 3666x 100.0% 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__1::coroutine_handle<void>) const::post_handler::destroy() :500 12x 100.0% 62.5% 100.0% boost::corosio::detail::reactor_scheduler::post(boost::corosio::detail::scheduler_op*) const :525 222591x 100.0% 75.0% 71.0% boost::corosio::detail::reactor_scheduler::post(boost::capy::continuation&) const :542 19523x 100.0% 75.0% 71.0% boost::corosio::detail::reactor_scheduler::running_in_this_thread() const :559 10701x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::stop() :565 3178x 100.0% 66.7% 71.0% boost::corosio::detail::reactor_scheduler::stopped() const :577 90x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::restart() :583 2353x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::run() :589 3186x 100.0% 92.9% 78.0% boost::corosio::detail::reactor_scheduler::run_one() :614 22x 75.0% 50.0% 50.0% boost::corosio::detail::reactor_scheduler::wait_one(long) :628 62x 100.0% 66.7% 60.0% boost::corosio::detail::reactor_scheduler::poll() :642 18x 100.0% 64.3% 78.0% boost::corosio::detail::reactor_scheduler::poll_one() :667 8x 100.0% 66.7% 60.0% boost::corosio::detail::reactor_scheduler::work_started() :681 117720x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::work_finished() :687 138502x 100.0% 75.0% 80.0% boost::corosio::detail::reactor_scheduler::compensating_work_started() const :694 1090499x 100.0% 50.0% 100.0% boost::corosio::detail::reactor_scheduler::drain_thread_queue(boost::corosio::detail::ready_queue&, long long) const :702 0 0.0% 0.0% 0.0% boost::corosio::detail::reactor_scheduler::post_deferred_completions(boost::corosio::detail::ready_queue&) const :715 97200x 70.0% 50.0% 55.0% boost::corosio::detail::reactor_scheduler::shutdown_drain() :732 1331x 100.0% 63.6% 90.0% boost::corosio::detail::reactor_scheduler::signal_all(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :760 4449x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::maybe_unlock_and_signal_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :767 14289x 62.5% 50.0% 75.0% boost::corosio::detail::reactor_scheduler::unlock_and_signal_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :781 1487119x 100.0% 100.0% 100.0% boost::corosio::detail::reactor_scheduler::clear_signal() const :793 7x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wait_for_signal(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :799 7x 100.0% 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wait_for_signal_for(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&, long) const :811 0 0.0% 0.0% 0.0% boost::corosio::detail::reactor_scheduler::wake_one_thread_and_unlock(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :823 14289x 90.0% 83.3% 85.0% boost::corosio::detail::reactor_scheduler::work_cleanup::~work_cleanup() :841 2866914x 94.1% 80.0% 100.0% boost::corosio::detail::reactor_scheduler::task_cleanup::~task_cleanup() :865 2255968x 86.7% 60.0% 100.0% boost::corosio::detail::reactor_scheduler::do_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&, long, boost::corosio::detail::reactor_scheduler_context*) :886 1436597x 90.4% 85.4% 81.0%
Line Branch TLA Hits Source Code
1 //
2 // Copyright (c) 2026 Steve Gerbino
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_REACTOR_REACTOR_SCHEDULER_HPP
11 #define BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_SCHEDULER_HPP
12
13 #include <boost/corosio/detail/config.hpp>
14 #include <boost/capy/ex/execution_context.hpp>
15
16 #include <boost/corosio/detail/ready_queue.hpp>
17 #include <boost/corosio/detail/scheduler.hpp>
18 #include <boost/corosio/detail/scheduler_op.hpp>
19 #include <boost/corosio/detail/thread_local_ptr.hpp>
20
21 #include <atomic>
22 #include <chrono>
23 #include <coroutine>
24 #include <cstddef>
25 #include <cstdint>
26 #include <limits>
27 #include <memory>
28 #include <stdexcept>
29
30 #include <boost/corosio/detail/conditionally_enabled_mutex.hpp>
31 #include <boost/corosio/detail/conditionally_enabled_event.hpp>
32
33 namespace boost::corosio::detail {
34
35 // Forward declarations
36 class reactor_scheduler;
37 class timer_service;
38
39 /** Per-thread state for a reactor scheduler.
40
41 Each thread running a scheduler's event loop has one of these
42 on a thread-local stack. It holds a private work queue and
43 inline completion budget for speculative I/O fast paths.
44 */
45 struct BOOST_COROSIO_SYMBOL_VISIBLE reactor_scheduler_context
46 {
47 /// Scheduler this context belongs to.
48 reactor_scheduler const* key;
49
50 /// Next context frame on this thread's stack.
51 reactor_scheduler_context* next;
52
53 /// Private work queue for reduced contention.
54 ready_queue private_queue;
55
56 /// Unflushed work count for the private queue.
57 std::int64_t private_outstanding_work;
58
59 /// Remaining inline completions allowed this cycle.
60 int inline_budget;
61
62 /// Maximum inline budget (adaptive, 2-16).
63 int inline_budget_max;
64
65 /// True if no other thread absorbed queued work last cycle.
66 bool unassisted;
67
68 /// Construct a context frame linked to @a n.
69 reactor_scheduler_context(
70 reactor_scheduler const* k,
71 reactor_scheduler_context* n);
72 };
73
74 /// Thread-local context stack for reactor schedulers.
75 inline thread_local_ptr<reactor_scheduler_context> reactor_context_stack;
76
77 /// Find the context frame for a scheduler on this thread.
78 inline reactor_scheduler_context*
79 2800563x reactor_find_context(reactor_scheduler const* self) noexcept
80 {
81
2/2
✓ Branch 0 taken 2776044 times.
✓ Branch 1 taken 24519 times.
2800563x for (auto* c = reactor_context_stack.get(); c != nullptr; c = c->next)
82 {
83
1/2
✓ Branch 0 taken 2776044 times.
✗ Branch 1 not taken.
2776044x if (c->key == self)
84 2776044x return c;
85 }
86 24519x return nullptr;
87 2800563x }
88
89 /// Flush private work count to global counter.
90 inline void
91 reactor_flush_private_work(
92 reactor_scheduler_context* ctx,
93 std::atomic<std::int64_t>& outstanding_work) noexcept
94 {
95 if (ctx && ctx->private_outstanding_work > 0)
96 {
97 outstanding_work.fetch_add(
98 ctx->private_outstanding_work, std::memory_order_relaxed);
99 ctx->private_outstanding_work = 0;
100 }
101 }
102
103 /** Drain private queue to global queue, flushing work count first.
104
105 @return True if any ops were drained.
106 */
107 inline bool
108 8x reactor_drain_private_queue(
109 reactor_scheduler_context* ctx,
110 std::atomic<std::int64_t>& outstanding_work,
111 ready_queue& completed_ops) noexcept
112 {
113
2/4
✓ Branch 0 taken 8 times.
✗ Branch 1 not taken.
✓ Branch 2 taken 8 times.
✗ Branch 3 not taken.
8x if (!ctx || ctx->private_queue.empty())
114 8x return false;
115
116 reactor_flush_private_work(ctx, outstanding_work);
117 completed_ops.splice(ctx->private_queue);
118 return true;
119 8x }
120
121 /** Non-template base for reactor-backed scheduler implementations.
122
123 Provides the complete threading model shared by epoll, kqueue,
124 and select schedulers: signal state machine, inline completion
125 budget, work counting, run/poll methods, and the do_one event
126 loop.
127
128 Derived classes provide platform-specific hooks by overriding:
129 - `run_task(lock, ctx)` to run the reactor poll
130 - `interrupt_reactor()` to wake a blocked reactor
131
132 De-templated from the original CRTP design to eliminate
133 duplicate instantiations when multiple backends are compiled
134 into the same binary. Virtual dispatch for run_task (called
135 once per reactor cycle, before a blocking syscall) has
136 negligible overhead.
137
138 @par Thread Safety
139 All public member functions are thread-safe.
140 */
141 class reactor_scheduler
142 : public scheduler
143 , public capy::execution_context::service
144 {
145 public:
146 using key_type = scheduler;
147 using context_type = reactor_scheduler_context;
148 using mutex_type = conditionally_enabled_mutex;
149 using lock_type = mutex_type::scoped_lock;
150 using event_type = conditionally_enabled_event;
151
152 /// Post a coroutine for deferred execution.
153 void post(std::coroutine_handle<> h) const override;
154
155 /// Post a scheduler operation for deferred execution.
156 void post(scheduler_op* h) const override;
157
158 /// Post a continuation for deferred execution.
159 void post(capy::continuation&) const override;
160
161 /// Return true if called from a thread running this scheduler.
162 bool running_in_this_thread() const noexcept override;
163
164 /// Request the scheduler to stop dispatching handlers.
165 void stop() override;
166
167 /// Return true if the scheduler has been stopped.
168 bool stopped() const noexcept override;
169
170 /// Reset the stopped state so `run()` can resume.
171 void restart() override;
172
173 /// Run the event loop until no work remains.
174 std::size_t run() override;
175
176 /// Run until one handler completes or no work remains.
177 std::size_t run_one() override;
178
179 /// Run until one handler completes or @a usec elapses.
180 std::size_t wait_one(long usec) override;
181
182 /// Run ready handlers without blocking.
183 std::size_t poll() override;
184
185 /// Run at most one ready handler without blocking.
186 std::size_t poll_one() override;
187
188 /// Increment the outstanding work count.
189 void work_started() noexcept override;
190
191 /// Decrement the outstanding work count, stopping on zero.
192 void work_finished() noexcept override;
193
194 /** Reset the thread's inline completion budget.
195
196 Called at the start of each posted completion handler to
197 grant a fresh budget for speculative inline completions.
198 */
199 void reset_inline_budget() const noexcept;
200
201 /** Consume one unit of inline budget if available.
202
203 @return True if budget was available and consumed.
204 */
205 bool try_consume_inline_budget() const noexcept;
206
207 /** Offset a forthcoming work_finished from work_cleanup.
208
209 Called by descriptor_state when all I/O returned EAGAIN and
210 no handler will be executed. Must be called from a scheduler
211 thread.
212 */
213 void compensating_work_started() const noexcept;
214
215 /** Drain work from thread context's private queue to global queue.
216
217 Flushes private work count to the global counter, then
218 transfers the queue under mutex protection.
219
220 @param queue The private queue to drain.
221 @param count Private work count to flush before draining.
222 */
223 void drain_thread_queue(ready_queue& queue, std::int64_t count) const;
224
225 /** Post completed operations for deferred invocation.
226
227 If called from a thread running this scheduler, operations
228 go to the thread's private queue (fast path). Otherwise,
229 operations are added to the global queue under mutex and a
230 waiter is signaled.
231
232 @par Preconditions
233 work_started() must have been called for each operation.
234
235 @param ops Queue of operations to post.
236 */
237 void post_deferred_completions(ready_queue& ops) const;
238
239 /** Apply runtime configuration to the scheduler.
240
241 Called by `io_context` after construction. Values that do
242 not apply to this backend are silently ignored.
243
244 @param max_events Event buffer size for epoll/kqueue.
245 @param budget_init Starting inline completion budget.
246 @param budget_max Hard ceiling on adaptive budget ramp-up.
247 @param unassisted Budget when single-threaded.
248 */
249 virtual void configure_reactor(
250 unsigned max_events,
251 unsigned budget_init,
252 unsigned budget_max,
253 unsigned unassisted);
254
255 /// Return the configured initial inline budget.
256 3194x unsigned inline_budget_initial() const noexcept
257 {
258 3194x return inline_budget_initial_;
259 }
260
261 /// Return true when scheduler locking is disabled (fully-lockless tier).
262 162x bool scheduler_locking_disabled() const noexcept override
263 {
264 162x return scheduler_locking_disabled_;
265 }
266
267 1331x void configure_threading(threading_config cfg) noexcept override
268 {
269 1331x scheduler_locking_disabled_ = !cfg.scheduler_locking;
270 // reactor_io_locking takes effect at descriptor registration (see the
271 // register_descriptor overrides), not here.
272 1331x reactor_io_locking_ = cfg.reactor_io_locking;
273 1331x one_thread_ = cfg.one_thread;
274 1331x mutex_.set_enabled(cfg.scheduler_locking);
275 1331x cond_.set_enabled(cfg.scheduler_locking);
276 1331x }
277
278 protected:
279 1331x timer_service* timer_svc_ = nullptr;
280 1331x bool scheduler_locking_disabled_ = false;
281 1331x bool reactor_io_locking_ = true;
282 1331x bool one_thread_ = false;
283
284 3993x reactor_scheduler() = default;
285
286 /** Drain completed_ops during shutdown.
287
288 Pops all operations from the global queue and destroys them,
289 skipping the task sentinel. Signals all waiting threads.
290 Derived classes call this from their shutdown() override
291 before performing platform-specific cleanup.
292 */
293 void shutdown_drain();
294
295 /// RAII guard that re-inserts the task sentinel after `run_task`.
296 struct task_cleanup
297 {
298 reactor_scheduler const* sched;
299 lock_type* lock;
300 context_type* ctx;
301 ~task_cleanup();
302 };
303
304 1331x mutable mutex_type mutex_{true};
305 1331x mutable event_type cond_{true};
306 mutable ready_queue completed_ops_;
307 1331x mutable std::atomic<std::int64_t> outstanding_work_{0};
308 1331x std::atomic<bool> stopped_{false};
309 1331x mutable std::atomic<bool> task_running_{false};
310 1331x mutable bool task_interrupted_ = false;
311
312 // Runtime-configurable reactor tuning parameters.
313 // Defaults match the library's built-in values.
314 1331x unsigned max_events_per_poll_ = 128;
315 1331x unsigned inline_budget_initial_ = 2;
316 1331x unsigned inline_budget_max_ = 16;
317 1331x unsigned unassisted_budget_ = 4;
318
319 /// Bit 0 of `state_`: set when the condvar should be signaled.
320 static constexpr std::size_t signaled_bit = 1;
321
322 /// Increment per waiting thread in `state_`.
323 static constexpr std::size_t waiter_increment = 2;
324 1331x mutable std::size_t state_ = 0;
325
326 /// Sentinel op that triggers a reactor poll when dequeued.
327 struct task_op final : scheduler_op
328 {
329 void operator()() override {}
330 void destroy() override {}
331 };
332 task_op task_op_;
333
334 /// Run the platform-specific reactor poll.
335 virtual void
336 run_task(lock_type& lock, context_type* ctx,
337 long timeout_us) = 0;
338
339 /// Wake a blocked reactor (e.g. write to eventfd or pipe).
340 virtual void interrupt_reactor() const = 0;
341
342 private:
343 struct work_cleanup
344 {
345 reactor_scheduler* sched;
346 lock_type* lock;
347 context_type* ctx;
348 ~work_cleanup();
349 };
350
351 std::size_t do_one(
352 lock_type& lock, long timeout_us, context_type* ctx);
353
354 void signal_all(lock_type& lock) const;
355 bool maybe_unlock_and_signal_one(lock_type& lock) const;
356 bool unlock_and_signal_one(lock_type& lock) const;
357 void clear_signal() const;
358 void wait_for_signal(lock_type& lock) const;
359 void wait_for_signal_for(
360 lock_type& lock, long timeout_us) const;
361 void wake_one_thread_and_unlock(lock_type& lock) const;
362 };
363
364 /** RAII guard that pushes/pops a scheduler context frame.
365
366 On construction, pushes a new context frame onto the
367 thread-local stack. On destruction, drains any remaining
368 private queue items to the global queue and pops the frame.
369 */
370 struct reactor_thread_context_guard
371 {
372 /// The context frame managed by this guard.
373 reactor_scheduler_context frame_;
374
375 /// Construct the guard, pushing a frame for @a sched.
376 6388x explicit reactor_thread_context_guard(
377 reactor_scheduler const* sched) noexcept
378
1/2
✓ Branch 0 taken 3194 times.
✗ Branch 1 not taken.
3194x : frame_(sched, reactor_context_stack.get())
379 3194x {
380 3194x reactor_context_stack.set(&frame_);
381 6388x }
382
383 /// Destroy the guard, draining private work and popping the frame.
384 6388x ~reactor_thread_context_guard() noexcept
385 3194x {
386
1/2
✓ Branch 0 taken 3194 times.
✗ Branch 1 not taken.
3194x if (!frame_.private_queue.empty())
387 frame_.key->drain_thread_queue(
388 frame_.private_queue, frame_.private_outstanding_work);
389 3194x reactor_context_stack.set(frame_.next);
390 6388x }
391 };
392
393 // ---- Inline implementations ------------------------------------------------
394
395 inline
396 9582x reactor_scheduler_context::reactor_scheduler_context(
397 reactor_scheduler const* k,
398 reactor_scheduler_context* n)
399 3194x : key(k)
400 3194x , next(n)
401 3194x , private_outstanding_work(0)
402 3194x , inline_budget(0)
403 6388x , inline_budget_max(
404 3194x static_cast<int>(k->inline_budget_initial()))
405 3194x , unassisted(false)
406 3194x {
407 6388x }
408
409 inline void
410 25x reactor_scheduler::configure_reactor(
411 unsigned max_events,
412 unsigned budget_init,
413 unsigned budget_max,
414 unsigned unassisted)
415 {
416
2/2
✓ Branch 0 taken 23 times.
✓ Branch 1 taken 2 times.
25x if (max_events < 1 ||
417 23x max_events > static_cast<unsigned>(std::numeric_limits<int>::max()))
418
1/2
✓ Branch 0 taken 2 times.
✗ Branch 1 not taken.
2x throw std::out_of_range(
419 "max_events_per_poll must be in [1, INT_MAX]");
420
1/2
✓ Branch 0 taken 23 times.
✗ Branch 1 not taken.
23x if (budget_max > static_cast<unsigned>(std::numeric_limits<int>::max()))
421 throw std::out_of_range(
422 "inline_budget_max must be in [0, INT_MAX]");
423
424 // Clamp initial and unassisted to budget_max.
425
2/2
✓ Branch 0 taken 21 times.
✓ Branch 1 taken 2 times.
23x if (budget_init > budget_max)
426 2x budget_init = budget_max;
427
2/2
✓ Branch 0 taken 21 times.
✓ Branch 1 taken 2 times.
23x if (unassisted > budget_max)
428 2x unassisted = budget_max;
429
430 23x max_events_per_poll_ = max_events;
431 23x inline_budget_initial_ = budget_init;
432 23x inline_budget_max_ = budget_max;
433 23x unassisted_budget_ = unassisted;
434 23x }
435
436 inline void
437 310835x reactor_scheduler::reset_inline_budget() const noexcept
438 {
439 // When budget is disabled (max==0), all paths below would no-op
440 // (inline_budget stays 0). Skip the TLS lookup entirely.
441
1/2
✓ Branch 0 taken 310835 times.
✗ Branch 1 not taken.
310835x if (inline_budget_max_ == 0)
442 return;
443
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 310835 times.
310835x if (auto* ctx = reactor_find_context(this))
444 {
445 // Cap when no other thread absorbed queued work
446
1/2
✓ Branch 0 taken 310835 times.
✗ Branch 1 not taken.
310835x if (ctx->unassisted)
447 {
448 310835x ctx->inline_budget_max =
449 310835x static_cast<int>(unassisted_budget_);
450 310835x ctx->inline_budget =
451 310835x static_cast<int>(unassisted_budget_);
452 310835x return;
453 }
454 // Ramp up when previous cycle fully consumed budget.
455 // max(1, ...) ensures the doubling escapes zero.
456 if (ctx->inline_budget == 0)
457 ctx->inline_budget_max = (std::min)(
458 (std::max)(1, ctx->inline_budget_max) * 2,
459 static_cast<int>(inline_budget_max_));
460 else if (ctx->inline_budget < ctx->inline_budget_max)
461 ctx->inline_budget_max =
462 static_cast<int>(inline_budget_initial_);
463 ctx->inline_budget = ctx->inline_budget_max;
464 }
465 310835x }
466
467 inline bool
468 1142736x reactor_scheduler::try_consume_inline_budget() const noexcept
469 {
470
1/2
✓ Branch 0 taken 1142736 times.
✗ Branch 1 not taken.
1142736x if (inline_budget_max_ == 0)
471 return false;
472
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 1142736 times.
1142736x if (auto* ctx = reactor_find_context(this))
473 {
474
2/2
✓ Branch 0 taken 930030 times.
✓ Branch 1 taken 212706 times.
1142736x if (ctx->inline_budget > 0)
475 {
476 930030x --ctx->inline_budget;
477 930030x return true;
478 }
479 212706x }
480 212706x return false;
481 1142736x }
482
483 inline void
484 3684x reactor_scheduler::post(std::coroutine_handle<> h) const
485 {
486 struct post_handler final : scheduler_op
487 {
488 std::coroutine_handle<> h_;
489
490 7359x explicit post_handler(std::coroutine_handle<> h) : h_(h) {}
491 7359x ~post_handler() override = default;
492
493 3666x void operator()() override
494 {
495 3666x auto saved = h_;
496
2/2
✓ Branch 0 taken 1 time.
✓ Branch 1 taken 3667 times.
3666x delete this;
497 3668x saved.resume();
498 3668x }
499
500 12x void destroy() override
501 {
502 12x auto saved = h_;
503
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 12 times.
12x delete this;
504 12x saved.destroy();
505 12x }
506 };
507
508 3684x auto ph = std::make_unique<post_handler>(h);
509
510
2/2
✓ Branch 0 taken 26 times.
✓ Branch 1 taken 3658 times.
3684x if (auto* ctx = reactor_find_context(this))
511 {
512 26x ++ctx->private_outstanding_work;
513 26x ctx->private_queue.push(ph.release());
514 26x return;
515 }
516
517 3658x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
518
519
1/2
✓ Branch 0 taken 3658 times.
✗ Branch 1 not taken.
3658x lock_type lock(mutex_);
520 3658x completed_ops_.push(ph.release());
521
1/2
✓ Branch 0 taken 3658 times.
✗ Branch 1 not taken.
3658x wake_one_thread_and_unlock(lock);
522 3684x }
523
524 inline void
525 222591x reactor_scheduler::post(scheduler_op* h) const
526 {
527
2/2
✓ Branch 0 taken 222203 times.
✓ Branch 1 taken 388 times.
222591x if (auto* ctx = reactor_find_context(this))
528 {
529 222203x ++ctx->private_outstanding_work;
530 222203x ctx->private_queue.push(h);
531 222203x return;
532 }
533
534 388x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
535
536 388x lock_type lock(mutex_);
537 388x completed_ops_.push(h);
538
1/2
✓ Branch 0 taken 388 times.
✗ Branch 1 not taken.
388x wake_one_thread_and_unlock(lock);
539 222591x }
540
541 inline void
542 19523x reactor_scheduler::post(capy::continuation& c) const
543 {
544
2/2
✓ Branch 0 taken 9282 times.
✓ Branch 1 taken 10241 times.
19523x if (auto* ctx = reactor_find_context(this))
545 {
546 9282x ++ctx->private_outstanding_work;
547 9282x ctx->private_queue.push(c);
548 9282x return;
549 }
550
551 10241x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
552
553 10241x lock_type lock(mutex_);
554 10241x completed_ops_.push(c);
555
1/2
✓ Branch 0 taken 10241 times.
✗ Branch 1 not taken.
10241x wake_one_thread_and_unlock(lock);
556 19523x }
557
558 inline bool
559 10701x reactor_scheduler::running_in_this_thread() const noexcept
560 {
561 10701x return reactor_find_context(this) != nullptr;
562 }
563
564 inline void
565 3178x reactor_scheduler::stop()
566 {
567 3178x lock_type lock(mutex_);
568
2/2
✓ Branch 0 taken 60 times.
✓ Branch 1 taken 3118 times.
3178x if (!stopped_.load(std::memory_order_acquire))
569 {
570 3118x stopped_.store(true, std::memory_order_release);
571
1/2
✓ Branch 0 taken 3118 times.
✗ Branch 1 not taken.
3118x signal_all(lock);
572
1/2
✓ Branch 0 taken 3118 times.
✗ Branch 1 not taken.
3118x interrupt_reactor();
573 3118x }
574 3178x }
575
576 inline bool
577 90x reactor_scheduler::stopped() const noexcept
578 {
579 90x return stopped_.load(std::memory_order_acquire);
580 }
581
582 inline void
583 2353x reactor_scheduler::restart()
584 {
585 2353x stopped_.store(false, std::memory_order_release);
586 2353x }
587
588 inline std::size_t
589 3186x reactor_scheduler::run()
590 {
591
2/2
✓ Branch 0 taken 3134 times.
✓ Branch 1 taken 52 times.
3186x if (outstanding_work_.load(std::memory_order_acquire) == 0)
592 {
593 52x stop();
594 52x return 0;
595 }
596
597 3134x reactor_thread_context_guard ctx(this);
598
1/2
✓ Branch 0 taken 3134 times.
✗ Branch 1 not taken.
3134x lock_type lock(mutex_);
599
600 3134x std::size_t n = 0;
601 1436521x for (;;)
602 {
603
4/4
✓ Branch 0 taken 1436498 times.
✓ Branch 1 taken 23 times.
✓ Branch 2 taken 1433388 times.
✓ Branch 3 taken 3110 times.
1436521x if (!do_one(lock, -1, &ctx.frame_))
604 3110x break;
605
2/2
✓ Branch 0 taken 12 times.
✓ Branch 1 taken 1433376 times.
1433388x if (n != (std::numeric_limits<std::size_t>::max)())
606 1433376x ++n;
607
2/2
✓ Branch 0 taken 1208165 times.
✓ Branch 1 taken 225199 times.
1433388x if (!lock.owns_lock())
608
2/2
✓ Branch 0 taken 1208188 times.
✓ Branch 1 taken 23 times.
1208165x lock.lock();
609 }
610 3110x return n;
611 3208x }
612
613 inline std::size_t
614 22x reactor_scheduler::run_one()
615 {
616
1/2
✓ Branch 0 taken 22 times.
✗ Branch 1 not taken.
22x if (outstanding_work_.load(std::memory_order_acquire) == 0)
617 {
618 stop();
619 return 0;
620 }
621
622 22x reactor_thread_context_guard ctx(this);
623
1/2
✓ Branch 0 taken 22 times.
✗ Branch 1 not taken.
22x lock_type lock(mutex_);
624
1/2
✓ Branch 0 taken 22 times.
✗ Branch 1 not taken.
22x return do_one(lock, -1, &ctx.frame_);
625 22x }
626
627 inline std::size_t
628 62x reactor_scheduler::wait_one(long usec)
629 {
630
2/2
✓ Branch 0 taken 42 times.
✓ Branch 1 taken 20 times.
62x if (outstanding_work_.load(std::memory_order_acquire) == 0)
631 {
632 20x stop();
633 20x return 0;
634 }
635
636 42x reactor_thread_context_guard ctx(this);
637
1/2
✓ Branch 0 taken 42 times.
✗ Branch 1 not taken.
42x lock_type lock(mutex_);
638
1/2
✓ Branch 0 taken 42 times.
✗ Branch 1 not taken.
42x return do_one(lock, usec, &ctx.frame_);
639 62x }
640
641 inline std::size_t
642 18x reactor_scheduler::poll()
643 {
644
2/2
✓ Branch 0 taken 16 times.
✓ Branch 1 taken 2 times.
18x if (outstanding_work_.load(std::memory_order_acquire) == 0)
645 {
646 2x stop();
647 2x return 0;
648 }
649
650 16x reactor_thread_context_guard ctx(this);
651
1/2
✓ Branch 0 taken 16 times.
✗ Branch 1 not taken.
16x lock_type lock(mutex_);
652
653 16x std::size_t n = 0;
654 50x for (;;)
655 {
656
3/4
✓ Branch 0 taken 50 times.
✗ Branch 1 not taken.
✓ Branch 2 taken 34 times.
✓ Branch 3 taken 16 times.
50x if (!do_one(lock, 0, &ctx.frame_))
657 16x break;
658
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 34 times.
34x if (n != (std::numeric_limits<std::size_t>::max)())
659 34x ++n;
660
1/2
✓ Branch 0 taken 34 times.
✗ Branch 1 not taken.
34x if (!lock.owns_lock())
661
1/2
✓ Branch 0 taken 34 times.
✗ Branch 1 not taken.
34x lock.lock();
662 }
663 16x return n;
664 18x }
665
666 inline std::size_t
667 8x reactor_scheduler::poll_one()
668 {
669
2/2
✓ Branch 0 taken 4 times.
✓ Branch 1 taken 4 times.
8x if (outstanding_work_.load(std::memory_order_acquire) == 0)
670 {
671 4x stop();
672 4x return 0;
673 }
674
675 4x reactor_thread_context_guard ctx(this);
676
1/2
✓ Branch 0 taken 4 times.
✗ Branch 1 not taken.
4x lock_type lock(mutex_);
677
1/2
✓ Branch 0 taken 4 times.
✗ Branch 1 not taken.
4x return do_one(lock, 0, &ctx.frame_);
678 8x }
679
680 inline void
681 117720x reactor_scheduler::work_started() noexcept
682 {
683 117720x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
684 117720x }
685
686 inline void
687 138502x reactor_scheduler::work_finished() noexcept
688 {
689
2/2
✓ Branch 0 taken 135408 times.
✓ Branch 1 taken 3094 times.
138502x if (outstanding_work_.fetch_sub(1, std::memory_order_acq_rel) == 1)
690
1/2
✓ Branch 0 taken 3094 times.
✗ Branch 1 not taken.
3094x stop();
691 138502x }
692
693 inline void
694 1090499x reactor_scheduler::compensating_work_started() const noexcept
695 {
696 1090499x auto* ctx = reactor_find_context(this);
697
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 1090499 times.
1090499x if (ctx)
698 1090499x ++ctx->private_outstanding_work;
699 1090499x }
700
701 inline void
702 reactor_scheduler::drain_thread_queue(
703 ready_queue& queue, std::int64_t count) const
704 {
705 if (count > 0)
706 outstanding_work_.fetch_add(count, std::memory_order_relaxed);
707
708 lock_type lock(mutex_);
709 completed_ops_.splice(queue);
710 if (count > 0)
711 maybe_unlock_and_signal_one(lock);
712 }
713
714 inline void
715 97200x reactor_scheduler::post_deferred_completions(ready_queue& ops) const
716 {
717
2/2
✓ Branch 0 taken 97199 times.
✓ Branch 1 taken 1 time.
97200x if (ops.empty())
718 97199x return;
719
720
1/2
✓ Branch 0 taken 1 time.
✗ Branch 1 not taken.
1x if (auto* ctx = reactor_find_context(this))
721 {
722 1x ctx->private_queue.splice(ops);
723 1x return;
724 }
725
726 lock_type lock(mutex_);
727 completed_ops_.splice(ops);
728 wake_one_thread_and_unlock(lock);
729 97200x }
730
731 inline void
732 1331x reactor_scheduler::shutdown_drain()
733 {
734 1331x lock_type lock(mutex_);
735
736
2/2
✓ Branch 0 taken 1331 times.
✓ Branch 1 taken 1543 times.
2874x while (auto e = completed_ops_.pop())
737 {
738
2/2
✓ Branch 0 taken 4 times.
✓ Branch 1 taken 1539 times.
1543x if (ready_is_continuation(e))
739 {
740
1/2
✓ Branch 0 taken 4 times.
✗ Branch 1 not taken.
4x lock.unlock();
741
1/2
✓ Branch 0 taken 4 times.
✗ Branch 1 not taken.
4x if (auto h = ready_as_cont(e)->h)
742
1/2
✓ Branch 0 taken 4 times.
✗ Branch 1 not taken.
4x h.destroy();
743
1/2
✓ Branch 0 taken 4 times.
✗ Branch 1 not taken.
4x lock.lock();
744 4x }
745 else
746 {
747 1539x auto* op = ready_as_op(e);
748
2/2
✓ Branch 0 taken 208 times.
✓ Branch 1 taken 1331 times.
1539x if (op == &task_op_)
749 1331x continue;
750
1/2
✓ Branch 0 taken 208 times.
✗ Branch 1 not taken.
208x lock.unlock();
751
1/2
✓ Branch 0 taken 208 times.
✗ Branch 1 not taken.
208x op->destroy();
752
1/2
✓ Branch 0 taken 208 times.
✗ Branch 1 not taken.
208x lock.lock();
753 }
754 }
755
756
1/2
✓ Branch 0 taken 1331 times.
✗ Branch 1 not taken.
1331x signal_all(lock);
757 1331x }
758
759 inline void
760 4449x reactor_scheduler::signal_all(lock_type&) const
761 {
762 4449x state_ |= signaled_bit;
763 4449x cond_.notify_all();
764 4449x }
765
766 inline bool
767 14289x reactor_scheduler::maybe_unlock_and_signal_one(
768 lock_type& lock) const
769 {
770 14289x state_ |= signaled_bit;
771
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 14289 times.
14289x if (state_ > signaled_bit)
772 {
773 lock.unlock();
774 cond_.notify_one();
775 return true;
776 }
777 14289x return false;
778 14289x }
779
780 inline bool
781 1487119x reactor_scheduler::unlock_and_signal_one(
782 lock_type& lock) const
783 {
784 1487119x state_ |= signaled_bit;
785 1487119x bool have_waiters = state_ > signaled_bit;
786 1487119x lock.unlock();
787
2/2
✓ Branch 0 taken 1487114 times.
✓ Branch 1 taken 5 times.
1487119x if (have_waiters)
788 5x cond_.notify_one();
789 1487119x return have_waiters;
790 }
791
792 inline void
793 7x reactor_scheduler::clear_signal() const
794 {
795 7x state_ &= ~signaled_bit;
796 7x }
797
798 inline void
799 7x reactor_scheduler::wait_for_signal(
800 lock_type& lock) const
801 {
802
2/2
✓ Branch 0 taken 8 times.
✓ Branch 1 taken 7 times.
15x while ((state_ & signaled_bit) == 0)
803 {
804 8x state_ += waiter_increment;
805 8x cond_.wait(lock);
806 8x state_ -= waiter_increment;
807 }
808 7x }
809
810 inline void
811 reactor_scheduler::wait_for_signal_for(
812 lock_type& lock, long timeout_us) const
813 {
814 if ((state_ & signaled_bit) == 0)
815 {
816 state_ += waiter_increment;
817 cond_.wait_for(lock, std::chrono::microseconds(timeout_us));
818 state_ -= waiter_increment;
819 }
820 }
821
822 inline void
823 14289x reactor_scheduler::wake_one_thread_and_unlock(
824 lock_type& lock) const
825 {
826
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 14289 times.
14289x if (maybe_unlock_and_signal_one(lock))
827 return;
828
829
4/4
✓ Branch 0 taken 206 times.
✓ Branch 1 taken 14083 times.
✓ Branch 2 taken 90 times.
✓ Branch 3 taken 116 times.
14289x if (task_running_.load(std::memory_order_relaxed) && !task_interrupted_)
830 {
831 116x task_interrupted_ = true;
832 116x lock.unlock();
833 116x interrupt_reactor();
834 116x }
835 else
836 {
837 14173x lock.unlock();
838 }
839 14289x }
840
841 2866914x inline reactor_scheduler::work_cleanup::~work_cleanup()
842 1433446x {
843
1/2
✓ Branch 0 taken 1433468 times.
✗ Branch 1 not taken.
1433468x if (ctx)
844 {
845 1433468x std::int64_t produced = ctx->private_outstanding_work;
846
2/2
✓ Branch 0 taken 306 times.
✓ Branch 1 taken 1433162 times.
1433468x if (produced > 1)
847 612x sched->outstanding_work_.fetch_add(
848 306x produced - 1, std::memory_order_relaxed);
849
2/2
✓ Branch 0 taken 1315109 times.
✓ Branch 1 taken 118053 times.
1433162x else if (produced < 1)
850 118053x sched->work_finished();
851 1433468x ctx->private_outstanding_work = 0;
852
853
2/2
✓ Branch 0 taken 1208269 times.
✓ Branch 1 taken 225199 times.
1433468x if (!ctx->private_queue.empty())
854 {
855
1/2
✓ Branch 0 taken 225199 times.
✗ Branch 1 not taken.
225199x lock->lock();
856 225199x sched->completed_ops_.splice(ctx->private_queue);
857 225199x }
858 1433468x }
859 else
860 {
861 sched->work_finished();
862 }
863 2866914x }
864
865 2255968x inline reactor_scheduler::task_cleanup::~task_cleanup()
866 1127984x {
867
1/2
✓ Branch 0 taken 1127984 times.
✗ Branch 1 not taken.
1127984x if (!ctx)
868 return;
869
870
2/2
✓ Branch 0 taken 1121689 times.
✓ Branch 1 taken 6295 times.
1127984x if (ctx->private_outstanding_work > 0)
871 {
872 12590x sched->outstanding_work_.fetch_add(
873 6295x ctx->private_outstanding_work, std::memory_order_relaxed);
874 6295x ctx->private_outstanding_work = 0;
875 6295x }
876
877
2/2
✓ Branch 0 taken 1121689 times.
✓ Branch 1 taken 6295 times.
1127984x if (!ctx->private_queue.empty())
878 {
879
1/2
✓ Branch 0 taken 6295 times.
✗ Branch 1 not taken.
6295x if (!lock->owns_lock())
880 lock->lock();
881 6295x sched->completed_ops_.splice(ctx->private_queue);
882 6295x }
883 2255968x }
884
885 inline std::size_t
886 1436597x reactor_scheduler::do_one(
887 lock_type& lock, long timeout_us, context_type* ctx)
888 {
889 1436604x for (;;)
890 {
891
2/2
✓ Branch 0 taken 2561458 times.
✓ Branch 1 taken 3112 times.
2564570x if (stopped_.load(std::memory_order_acquire))
892 3112x return 0;
893
894 2561458x std::uintptr_t e = completed_ops_.pop();
895
2/2
✓ Branch 0 taken 19502 times.
✓ Branch 1 taken 2541956 times.
2561458x scheduler_op* op = ready_is_continuation(e) ? nullptr : ready_as_op(e);
896
897 // Handle reactor sentinel — time to poll for I/O
898
2/2
✓ Branch 0 taken 1433461 times.
✓ Branch 1 taken 1127997 times.
2561458x if (op == &task_op_)
899 {
900 1127997x bool more_handlers =
901
3/4
✓ Branch 0 taken 53683 times.
✓ Branch 1 taken 1074314 times.
✗ Branch 2 not taken.
✓ Branch 3 taken 1074314 times.
1127997x !completed_ops_.empty() || (ctx && !ctx->private_queue.empty());
902
903
4/4
✓ Branch 0 taken 1074314 times.
✓ Branch 1 taken 53683 times.
✓ Branch 2 taken 1074301 times.
✓ Branch 3 taken 12 times.
2202310x if (!more_handlers &&
904
2/2
✓ Branch 0 taken 1074313 times.
✓ Branch 1 taken 1 time.
1074314x (outstanding_work_.load(std::memory_order_acquire) == 0 ||
905 1074313x timeout_us == 0))
906 {
907 13x completed_ops_.push(&task_op_);
908 13x return 0;
909 }
910
911
2/2
✓ Branch 0 taken 53683 times.
✓ Branch 1 taken 1074301 times.
1127984x long task_timeout_us = more_handlers ? 0 : timeout_us;
912 1127984x task_interrupted_ = task_timeout_us == 0;
913 1127984x task_running_.store(true, std::memory_order_release);
914
915 // Wake a peer to take the pending handlers while this thread
916 // polls the reactor; skipped when one_thread_ (no peer exists).
917
4/4
✓ Branch 0 taken 53683 times.
✓ Branch 1 taken 1074301 times.
✓ Branch 2 taken 7 times.
✓ Branch 3 taken 53676 times.
1127984x if (more_handlers && !one_thread_)
918 53676x unlock_and_signal_one(lock);
919
920 try
921 {
922
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 1127984 times.
1127984x run_task(lock, ctx, task_timeout_us);
923 1127984x }
924 catch (...)
925 {
926 task_running_.store(false, std::memory_order_relaxed);
927 throw;
928 }
929
930 1127984x task_running_.store(false, std::memory_order_relaxed);
931 1127984x completed_ops_.push(&task_op_);
932
2/2
✓ Branch 0 taken 18 times.
✓ Branch 1 taken 1127966 times.
1127984x if (timeout_us > 0)
933 18x return 0;
934 1127966x continue;
935 }
936
937 // Handle ready entry (op or continuation)
938
2/2
✓ Branch 0 taken 8 times.
✓ Branch 1 taken 1433453 times.
1433461x if (e != 0)
939 {
940 1433453x bool more = !completed_ops_.empty();
941
942
4/4
✓ Branch 0 taken 1433451 times.
✓ Branch 1 taken 2 times.
✓ Branch 2 taken 8 times.
✓ Branch 3 taken 1433443 times.
1433453x if (more && !one_thread_)
943 {
944 // Wake a peer for the remaining work; unassisted if none
945 // was parked to take it.
946 1433443x ctx->unassisted = !unlock_and_signal_one(lock);
947 1433443x }
948 else
949 {
950 // No peer to wake (one_thread_, or nothing more queued).
951 10x ctx->unassisted = more;
952 10x lock.unlock();
953 }
954
955 1433453x work_cleanup on_exit{this, &lock, ctx};
956 (void)on_exit;
957
958
2/2
✓ Branch 0 taken 1413946 times.
✓ Branch 1 taken 19507 times.
1433453x if (ready_is_continuation(e))
959
2/2
✓ Branch 0 taken 19504 times.
✓ Branch 1 taken 3 times.
19507x ready_as_cont(e)->h.resume();
960 else
961
2/2
✓ Branch 0 taken 1413949 times.
✓ Branch 1 taken 3 times.
1413946x (*op)();
962 1433453x return 1;
963 1433459x }
964
965 // Try private queue before blocking
966
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 8 times.
8x if (reactor_drain_private_queue(ctx, outstanding_work_, completed_ops_))
967 continue;
968
969
3/4
✓ Branch 0 taken 7 times.
✓ Branch 1 taken 1 time.
✗ Branch 2 not taken.
✓ Branch 3 taken 7 times.
8x if (outstanding_work_.load(std::memory_order_acquire) == 0 ||
970 7x timeout_us == 0)
971 1x return 0;
972
973 7x clear_signal();
974
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 7 times.
7x if (timeout_us < 0)
975 7x wait_for_signal(lock);
976 else
977 wait_for_signal_for(lock, timeout_us);
978 }
979 1436603x }
980
981 } // namespace boost::corosio::detail
982
983 #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_SCHEDULER_HPP
984