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

98.6% Lines (363/370) 100.0% List of functions (45/47) 77.8% Branches (168/216)
reactor_scheduler.hpp
f(x) Functions (47)
Function Calls Lines Branches Blocks
boost::corosio::detail::reactor_find_context(boost::corosio::detail::reactor_scheduler const*) :78 2818715x 85.7% 75.0% 75.0% boost::corosio::detail::reactor_scheduler::~reactor_scheduler() :108 2977x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::inline_budget_initial() const :210 9907x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::scheduler_locking_disabled() const :216 862x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::configure_threading(boost::corosio::detail::scheduler::threading_config) :221 2967x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::reactor_scheduler() :238 2977x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::task_op::task_op() :281 5954x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::task_op::~task_op() :281 5954x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::task_op::operator()() :285 – – – – boost::corosio::detail::reactor_scheduler::task_op::destroy() :286 – – – – boost::corosio::detail::reactor_thread_context_guard::reactor_thread_context_guard(boost::corosio::detail::reactor_scheduler const*) :337 19814x 100.0% 50.0% 100.0% boost::corosio::detail::reactor_thread_context_guard::~reactor_thread_context_guard() :350 19814x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler_context::reactor_scheduler_context(boost::corosio::detail::reactor_scheduler const*, boost::corosio::detail::reactor_scheduler_context*) :358 19814x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::configure_reactor(unsigned int, unsigned int, unsigned int, unsigned int) :370 45x 100.0% 83.3% 75.0% boost::corosio::detail::reactor_scheduler::reset_inline_budget() const :395 479408x 100.0% 75.0% 93.0% boost::corosio::detail::reactor_scheduler::try_consume_inline_budget() const :423 1896084x 100.0% 83.3% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__1::coroutine_handle<void>) const :439 3721x 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>) :445 7433x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::post(std::__1::coroutine_handle<void>) const::post_handler::~post_handler() :446 7433x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::post(std::__1::coroutine_handle<void>) const::post_handler::operator()() :448 3704x 100.0% 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__1::coroutine_handle<void>) const::post_handler::destroy() :455 12x 100.0% 62.5% 100.0% boost::corosio::detail::reactor_scheduler::post(boost::corosio::detail::scheduler_op*) const :480 378491x 100.0% 75.0% 71.0% boost::corosio::detail::reactor_scheduler::post(boost::capy::continuation&) const :497 32476x 100.0% 75.0% 71.0% boost::corosio::detail::reactor_scheduler::running_in_this_thread() const :514 19542x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::stop() :520 8115x 100.0% 66.7% 71.0% boost::corosio::detail::reactor_scheduler::stopped() const :532 2526x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::restart() :538 5193x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::run() :544 6504x 100.0% 92.9% 94.0% boost::corosio::detail::reactor_scheduler::run_one() :569 112x 100.0% 66.7% 60.0% boost::corosio::detail::reactor_scheduler::wait_one(long) :583 4136x 100.0% 66.7% 60.0% boost::corosio::detail::reactor_scheduler::poll() :597 49x 100.0% 64.3% 78.0% boost::corosio::detail::reactor_scheduler::poll_one() :622 11x 100.0% 66.7% 60.0% boost::corosio::detail::reactor_scheduler::work_started() :636 149965x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::work_finished() :642 184634x 100.0% 75.0% 80.0% boost::corosio::detail::reactor_scheduler::compensating_work_started() const :649 9094x 100.0% 50.0% 100.0% boost::corosio::detail::reactor_scheduler::post_deferred_completions(boost::corosio::detail::ready_queue&) const :657 115366x 70.0% 50.0% 55.0% boost::corosio::detail::reactor_scheduler::shutdown_drain() :674 2967x 100.0% 63.6% 90.0% boost::corosio::detail::reactor_scheduler::signal_all(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :702 10228x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::maybe_unlock_and_signal_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :709 24266x 100.0% 100.0% 100.0% boost::corosio::detail::reactor_scheduler::unlock_and_signal_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :722 619186x 100.0% 100.0% 100.0% boost::corosio::detail::reactor_scheduler::clear_signal() const :733 19x 100.0% – 100.0% boost::corosio::detail::reactor_scheduler::wait_for_signal(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :739 11x 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 :750 8x 100.0% 50.0% 100.0% boost::corosio::detail::reactor_scheduler::wake_one_thread_and_unlock(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :761 24266x 100.0% 100.0% 100.0% boost::corosio::detail::reactor_scheduler::work_cleanup::~work_cleanup() :778 1078089x 100.0% 87.5% 100.0% boost::corosio::detail::reactor_scheduler::task_cleanup::~task_cleanup() :795 461568x 92.3% 62.5% 100.0% boost::corosio::detail::reactor_scheduler::do_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&, long, boost::corosio::detail::reactor_scheduler_context&) :813 547173x 100.0% 95.2% 83.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, reactor_scheduler_context* n);
71 };
72
73 /// Thread-local context stack for reactor schedulers.
74 inline thread_local_ptr<reactor_scheduler_context> reactor_context_stack;
75
76 /// Find the context frame for a scheduler on this thread.
77 inline reactor_scheduler_context*
78 2818715x reactor_find_context(reactor_scheduler const* self) noexcept
79 {
80
2/2
✓ Branch 0 taken 2775514 times.
✓ Branch 1 taken 43201 times.
2818715x for (auto* c = reactor_context_stack.get(); c != nullptr; c = c->next)
81 {
82
1/2
✓ Branch 0 taken 2775514 times.
✗ Branch 1 not taken.
2775514x if (c->key == self)
83 2775514x return c;
84 ✗ }
85 43201x return nullptr;
86 2818715x }
87
88 /** Non-template base for reactor-backed scheduler implementations.
89
90 Provides the complete threading model shared by epoll, kqueue,
91 and select schedulers: signal state machine, inline completion
92 budget, work counting, run/poll methods, and the do_one event
93 loop.
94
95 Derived classes provide platform-specific hooks by overriding:
96 - `run_task(lock, ctx)` to run the reactor poll
97 - `interrupt_reactor()` to wake a blocked reactor
98
99 De-templated from the original CRTP design to eliminate
100 duplicate instantiations when multiple backends are compiled
101 into the same binary. Virtual dispatch for run_task (called
102 once per reactor cycle, before a blocking syscall) has
103 negligible overhead.
104
105 @par Thread Safety
106 All public member functions are thread-safe.
107 */
108 class reactor_scheduler
109 : public scheduler
110 {
111 public:
112 using context_type = reactor_scheduler_context;
113 using mutex_type = conditionally_enabled_mutex;
114 using lock_type = mutex_type::scoped_lock;
115 using event_type = conditionally_enabled_event;
116
117 /// Post a coroutine for deferred execution.
118 void post(std::coroutine_handle<> h) const override;
119
120 /// Post a scheduler operation for deferred execution.
121 void post(scheduler_op* h) const override;
122
123 /// Post a continuation for deferred execution.
124 void post(capy::continuation&) const override;
125
126 /// Return true if called from a thread running this scheduler.
127 bool running_in_this_thread() const noexcept override;
128
129 /// Request the scheduler to stop dispatching handlers.
130 void stop() override;
131
132 /// Return true if the scheduler has been stopped.
133 bool stopped() const noexcept override;
134
135 /// Reset the stopped state so `run()` can resume.
136 void restart() override;
137
138 /// Run the event loop until no work remains.
139 std::size_t run() override;
140
141 /// Run until one handler completes or no work remains.
142 std::size_t run_one() override;
143
144 /// Run until one handler completes or @a usec elapses.
145 std::size_t wait_one(long usec) override;
146
147 /// Run ready handlers without blocking.
148 std::size_t poll() override;
149
150 /// Run at most one ready handler without blocking.
151 std::size_t poll_one() override;
152
153 /// Increment the outstanding work count.
154 void work_started() noexcept override;
155
156 /// Decrement the outstanding work count, stopping on zero.
157 void work_finished() noexcept override;
158
159 /** Reset the thread's inline completion budget.
160
161 Called at the start of each posted completion handler to
162 grant a fresh budget for speculative inline completions.
163 */
164 void reset_inline_budget() const noexcept;
165
166 /** Consume one unit of inline budget if available.
167
168 @return True if budget was available and consumed.
169 */
170 bool try_consume_inline_budget() const noexcept;
171
172 /** Offset a forthcoming work_finished from work_cleanup.
173
174 Called by descriptor_state when all I/O returned EAGAIN and
175 no handler will be executed. Must be called from a scheduler
176 thread.
177 */
178 void compensating_work_started() const noexcept;
179
180 /** Post completed operations for deferred invocation.
181
182 If called from a thread running this scheduler, operations
183 go to the thread's private queue (fast path). Otherwise,
184 operations are added to the global queue under mutex and a
185 waiter is signaled.
186
187 @pre work_started() must have been called for each operation.
188
189 @param ops Queue of operations to post.
190 */
191 void post_deferred_completions(ready_queue& ops) const;
192
193 /** Apply runtime configuration to the scheduler.
194
195 Called by `io_context` after construction. Values that do
196 not apply to this backend are silently ignored.
197
198 @param max_events Event buffer size for epoll/kqueue.
199 @param budget_init Starting inline completion budget.
200 @param budget_max Hard ceiling on adaptive budget ramp-up.
201 @param unassisted Budget when single-threaded.
202 */
203 virtual void configure_reactor(
204 unsigned max_events,
205 unsigned budget_init,
206 unsigned budget_max,
207 unsigned unassisted);
208
209 /// Return the configured initial inline budget.
210 9907x unsigned inline_budget_initial() const noexcept
211 {
212 9907x return inline_budget_initial_;
213 }
214
215 /// Return true when scheduler locking is disabled (fully-lockless tier).
216 862x bool scheduler_locking_disabled() const noexcept override
217 {
218 862x return scheduler_locking_disabled_;
219 }
220
221 2967x void configure_threading(threading_config cfg) noexcept override
222 {
223 2967x scheduler_locking_disabled_ = !cfg.scheduler_locking;
224 // reactor_io_locking takes effect at descriptor registration (see the
225 // register_descriptor overrides), not here.
226 2967x reactor_io_locking_ = cfg.reactor_io_locking;
227 2967x one_thread_ = cfg.one_thread;
228 2967x mutex_.set_enabled(cfg.scheduler_locking);
229 2967x cond_.set_enabled(cfg.scheduler_locking);
230 2967x }
231
232 protected:
233 2977x timer_service* timer_svc_ = nullptr;
234 2977x bool scheduler_locking_disabled_ = false;
235 2977x bool reactor_io_locking_ = true;
236 2977x bool one_thread_ = false;
237
238 8931x reactor_scheduler() = default;
239
240 /** Drain completed_ops during shutdown.
241
242 Pops all operations from the global queue and destroys them,
243 skipping the task sentinel. Signals all waiting threads.
244 Derived classes call this from their shutdown() override
245 before performing platform-specific cleanup.
246 */
247 void shutdown_drain();
248
249 /// RAII guard that re-inserts the task sentinel after `run_task`.
250 struct task_cleanup
251 {
252 reactor_scheduler const* sched;
253 lock_type* lock;
254 context_type& ctx;
255 ~task_cleanup();
256 };
257
258 2977x mutable mutex_type mutex_{true};
259 2977x mutable event_type cond_{true};
260 mutable ready_queue completed_ops_;
261 2977x mutable std::atomic<std::int64_t> outstanding_work_{0};
262 2977x std::atomic<bool> stopped_{false};
263 2977x mutable std::atomic<bool> task_running_{false};
264 2977x mutable bool task_interrupted_ = false;
265
266 // Runtime-configurable reactor tuning parameters.
267 // Defaults match the library's built-in values.
268 2977x unsigned max_events_per_poll_ = 128;
269 2977x unsigned inline_budget_initial_ = 2;
270 2977x unsigned inline_budget_max_ = 16;
271 2977x unsigned unassisted_budget_ = 4;
272
273 /// Bit 0 of `state_`: set when the condvar should be signaled.
274 static constexpr std::size_t signaled_bit = 1;
275
276 /// Increment per waiting thread in `state_`.
277 static constexpr std::size_t waiter_increment = 2;
278 2977x mutable std::size_t state_ = 0;
279
280 /// Sentinel op that triggers a reactor poll when dequeued.
281 struct task_op final : scheduler_op
282 {
283 // LCOV_EXCL_START: the sentinel is intercepted by pointer
284 // identity; its virtuals exist for vtable completeness.
285 − void operator()() override {}
286 − void destroy() override {}
287 // LCOV_EXCL_STOP
288 };
289 task_op task_op_;
290
291 /** Run the platform-specific reactor poll.
292
293 @par Postconditions
294 `lock` is owned on return, however the poll ended. An
295 implementation that unlocks around the blocking call owes the
296 caller a matching re-acquire on every path out, including the
297 errors it retries rather than reports.
298 */
299 virtual void
300 run_task(lock_type& lock, context_type& ctx, long timeout_us) = 0;
301
302 /// Wake a blocked reactor (e.g. write to eventfd or pipe).
303 virtual void interrupt_reactor() const = 0;
304
305 private:
306 struct work_cleanup
307 {
308 reactor_scheduler* sched;
309 lock_type* lock;
310 context_type& ctx;
311 ~work_cleanup();
312 };
313
314 std::size_t do_one(lock_type& lock, long timeout_us, context_type& ctx);
315
316 void signal_all(lock_type& lock) const;
317 bool maybe_unlock_and_signal_one(lock_type& lock) const;
318 bool unlock_and_signal_one(lock_type& lock) const;
319 void clear_signal() const;
320 void wait_for_signal(lock_type& lock) const;
321 void wait_for_signal_for(lock_type& lock, long timeout_us) const;
322 void wake_one_thread_and_unlock(lock_type& lock) const;
323 };
324
325 /** RAII guard that pushes/pops a scheduler context frame.
326
327 On construction, pushes a new context frame onto the
328 thread-local stack. On destruction, drains any remaining
329 private queue items to the global queue and pops the frame.
330 */
331 struct reactor_thread_context_guard
332 {
333 /// The context frame managed by this guard.
334 reactor_scheduler_context frame_;
335
336 /// Construct the guard, pushing a frame for @a sched.
337 19814x explicit reactor_thread_context_guard(
338 reactor_scheduler const* sched) noexcept
339
1/2
✓ Branch 0 taken 9907 times.
✗ Branch 1 not taken.
9907x : frame_(sched, reactor_context_stack.get())
340 9907x {
341 9907x reactor_context_stack.set(&frame_);
342 19814x }
343
344 /** Destroy the guard, popping the frame.
345
346 The private queue is empty here by invariant: work_cleanup and
347 task_cleanup splice it to the global queue after every handler
348 and every reactor pass.
349 */
350 19814x ~reactor_thread_context_guard() noexcept
351 9907x {
352 9907x reactor_context_stack.set(frame_.next);
353 19814x }
354 };
355
356 // ---- Inline implementations ------------------------------------------------
357
358 29721x inline reactor_scheduler_context::reactor_scheduler_context(
359 reactor_scheduler const* k, reactor_scheduler_context* n)
360 9907x : key(k)
361 9907x , next(n)
362 9907x , private_outstanding_work(0)
363 9907x , inline_budget(0)
364 9907x , inline_budget_max(static_cast<int>(k->inline_budget_initial()))
365 9907x , unassisted(false)
366 9907x {
367 19814x }
368
369 inline void
370 45x reactor_scheduler::configure_reactor(
371 unsigned max_events,
372 unsigned budget_init,
373 unsigned budget_max,
374 unsigned unassisted)
375 {
376
2/2
✓ Branch 0 taken 43 times.
✓ Branch 1 taken 2 times.
45x if (max_events < 1 ||
377 43x max_events > static_cast<unsigned>(std::numeric_limits<int>::max()))
378
1/2
✓ Branch 0 taken 2 times.
✗ Branch 1 not taken.
2x throw std::out_of_range("max_events_per_poll must be in [1, INT_MAX]");
379
2/2
✓ Branch 0 taken 41 times.
✓ Branch 1 taken 2 times.
43x if (budget_max > static_cast<unsigned>(std::numeric_limits<int>::max()))
380
1/2
✓ Branch 0 taken 2 times.
✗ Branch 1 not taken.
2x throw std::out_of_range("inline_budget_max must be in [0, INT_MAX]");
381
382 // Clamp initial and unassisted to budget_max.
383
2/2
✓ Branch 0 taken 23 times.
✓ Branch 1 taken 18 times.
41x if (budget_init > budget_max)
384 18x budget_init = budget_max;
385
2/2
✓ Branch 0 taken 23 times.
✓ Branch 1 taken 18 times.
41x if (unassisted > budget_max)
386 18x unassisted = budget_max;
387
388 41x max_events_per_poll_ = max_events;
389 41x inline_budget_initial_ = budget_init;
390 41x inline_budget_max_ = budget_max;
391 41x unassisted_budget_ = unassisted;
392 41x }
393
394 inline void
395 479408x reactor_scheduler::reset_inline_budget() const noexcept
396 {
397 // When budget is disabled (max==0), all paths below would no-op
398 // (inline_budget stays 0). Skip the TLS lookup entirely.
399
2/2
✓ Branch 0 taken 479350 times.
✓ Branch 1 taken 58 times.
479408x if (inline_budget_max_ == 0)
400 58x return;
401
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 479350 times.
479350x if (auto* ctx = reactor_find_context(this))
402 {
403 // Cap when no other thread absorbed queued work
404
2/2
✓ Branch 0 taken 479345 times.
✓ Branch 1 taken 5 times.
479350x if (ctx->unassisted)
405 {
406 479345x ctx->inline_budget_max = static_cast<int>(unassisted_budget_);
407 479345x ctx->inline_budget = static_cast<int>(unassisted_budget_);
408 479345x return;
409 }
410 // Ramp up when previous cycle fully consumed budget.
411 // max(1, ...) ensures the doubling escapes zero.
412
2/2
✓ Branch 0 taken 2 times.
✓ Branch 1 taken 3 times.
5x if (ctx->inline_budget == 0)
413 3x ctx->inline_budget_max =
414
3/6
✓ Branch 0 taken 3 times.
✗ Branch 1 not taken.
✓ Branch 2 taken 3 times.
✗ Branch 3 not taken.
✓ Branch 4 taken 3 times.
✗ Branch 5 not taken.
3x (std::min)((std::max)(1, ctx->inline_budget_max) * 2,
415 3x static_cast<int>(inline_budget_max_));
416
2/2
✓ Branch 0 taken 1 time.
✓ Branch 1 taken 1 time.
2x else if (ctx->inline_budget < ctx->inline_budget_max)
417 1x ctx->inline_budget_max = static_cast<int>(inline_budget_initial_);
418 5x ctx->inline_budget = ctx->inline_budget_max;
419 5x }
420 479408x }
421
422 inline bool
423 1896084x reactor_scheduler::try_consume_inline_budget() const noexcept
424 {
425
2/2
✓ Branch 0 taken 1896042 times.
✓ Branch 1 taken 42 times.
1896084x if (inline_budget_max_ == 0)
426 42x return false;
427
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 1896042 times.
1896042x if (auto* ctx = reactor_find_context(this))
428 {
429
2/2
✓ Branch 0 taken 1533764 times.
✓ Branch 1 taken 362278 times.
1896042x if (ctx->inline_budget > 0)
430 {
431 1533764x --ctx->inline_budget;
432 1533764x return true;
433 }
434 362278x }
435 362278x return false;
436 1896084x }
437
438 inline void
439 3721x reactor_scheduler::post(std::coroutine_handle<> h) const
440 {
441 struct post_handler final : scheduler_op
442 {
443 std::coroutine_handle<> h_;
444
445 7433x explicit post_handler(std::coroutine_handle<> h) : h_(h) {}
446 7433x ~post_handler() override = default;
447
448 3704x void operator()() override
449 {
450 3704x auto saved = h_;
451
2/2
✓ Branch 0 taken 1 time.
✓ Branch 1 taken 3705 times.
3704x delete this;
452 3706x saved.resume();
453 3706x }
454
455 12x void destroy() override
456 {
457 12x auto saved = h_;
458
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 12 times.
12x delete this;
459 12x saved.destroy();
460 12x }
461 };
462
463 3721x auto ph = std::make_unique<post_handler>(h);
464
465
2/2
✓ Branch 0 taken 62 times.
✓ Branch 1 taken 3659 times.
3721x if (auto* ctx = reactor_find_context(this))
466 {
467 62x ++ctx->private_outstanding_work;
468 62x ctx->private_queue.push(ph.release());
469 62x return;
470 }
471
472 3659x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
473
474
1/2
✓ Branch 0 taken 3659 times.
✗ Branch 1 not taken.
3659x lock_type lock(mutex_);
475 3659x completed_ops_.push(ph.release());
476
1/2
✓ Branch 0 taken 3659 times.
✗ Branch 1 not taken.
3659x wake_one_thread_and_unlock(lock);
477 3721x }
478
479 inline void
480 378491x reactor_scheduler::post(scheduler_op* h) const
481 {
482
2/2
✓ Branch 0 taken 377332 times.
✓ Branch 1 taken 1159 times.
378491x if (auto* ctx = reactor_find_context(this))
483 {
484 377332x ++ctx->private_outstanding_work;
485 377332x ctx->private_queue.push(h);
486 377332x return;
487 }
488
489 1159x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
490
491 1159x lock_type lock(mutex_);
492 1159x completed_ops_.push(h);
493
1/2
✓ Branch 0 taken 1159 times.
✗ Branch 1 not taken.
1159x wake_one_thread_and_unlock(lock);
494 378491x }
495
496 inline void
497 32476x reactor_scheduler::post(capy::continuation& c) const
498 {
499
2/2
✓ Branch 0 taken 13029 times.
✓ Branch 1 taken 19447 times.
32476x if (auto* ctx = reactor_find_context(this))
500 {
501 13029x ++ctx->private_outstanding_work;
502 13029x ctx->private_queue.push(c);
503 13029x return;
504 }
505
506 19447x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
507
508 19447x lock_type lock(mutex_);
509 19447x completed_ops_.push(c);
510
1/2
✓ Branch 0 taken 19447 times.
✗ Branch 1 not taken.
19447x wake_one_thread_and_unlock(lock);
511 32476x }
512
513 inline bool
514 19542x reactor_scheduler::running_in_this_thread() const noexcept
515 {
516 19542x return reactor_find_context(this) != nullptr;
517 }
518
519 inline void
520 8115x reactor_scheduler::stop()
521 {
522 8115x lock_type lock(mutex_);
523
2/2
✓ Branch 0 taken 854 times.
✓ Branch 1 taken 7261 times.
8115x if (!stopped_.load(std::memory_order_acquire))
524 {
525 7261x stopped_.store(true, std::memory_order_release);
526
1/2
✓ Branch 0 taken 7261 times.
✗ Branch 1 not taken.
7261x signal_all(lock);
527
1/2
✓ Branch 0 taken 7261 times.
✗ Branch 1 not taken.
7261x interrupt_reactor();
528 7261x }
529 8115x }
530
531 inline bool
532 2526x reactor_scheduler::stopped() const noexcept
533 {
534 2526x return stopped_.load(std::memory_order_acquire);
535 }
536
537 inline void
538 5193x reactor_scheduler::restart()
539 {
540 5193x stopped_.store(false, std::memory_order_release);
541 5193x }
542
543 inline std::size_t
544 6504x reactor_scheduler::run()
545 {
546
2/2
✓ Branch 0 taken 6451 times.
✓ Branch 1 taken 53 times.
6504x if (outstanding_work_.load(std::memory_order_acquire) == 0)
547 {
548 53x stop();
549 53x return 0;
550 }
551
552 6451x reactor_thread_context_guard ctx(this);
553
1/2
✓ Branch 0 taken 6451 times.
✗ Branch 1 not taken.
6451x lock_type lock(mutex_);
554
555 6451x std::size_t n = 0;
556 543690x for (;;)
557 {
558
4/4
✓ Branch 0 taken 543680 times.
✓ Branch 1 taken 10 times.
✓ Branch 2 taken 537235 times.
✓ Branch 3 taken 6445 times.
543690x if (!do_one(lock, -1, ctx.frame_))
559 6445x break;
560
2/2
✓ Branch 0 taken 2 times.
✓ Branch 1 taken 537233 times.
537235x if (n != (std::numeric_limits<std::size_t>::max)())
561 537233x ++n;
562
2/2
✓ Branch 0 taken 157474 times.
✓ Branch 1 taken 379757 times.
537235x if (!lock.owns_lock())
563
2/2
✓ Branch 0 taken 157482 times.
✓ Branch 1 taken 8 times.
157474x lock.lock();
564 }
565 6445x return n;
566 6516x }
567
568 inline std::size_t
569 112x reactor_scheduler::run_one()
570 {
571
2/2
✓ Branch 0 taken 109 times.
✓ Branch 1 taken 3 times.
112x if (outstanding_work_.load(std::memory_order_acquire) == 0)
572 {
573 3x stop();
574 3x return 0;
575 }
576
577 109x reactor_thread_context_guard ctx(this);
578
1/2
✓ Branch 0 taken 109 times.
✗ Branch 1 not taken.
109x lock_type lock(mutex_);
579
1/2
✓ Branch 0 taken 109 times.
✗ Branch 1 not taken.
109x return do_one(lock, -1, ctx.frame_);
580 112x }
581
582 inline std::size_t
583 4136x reactor_scheduler::wait_one(long usec)
584 {
585
2/2
✓ Branch 0 taken 3311 times.
✓ Branch 1 taken 825 times.
4136x if (outstanding_work_.load(std::memory_order_acquire) == 0)
586 {
587 825x stop();
588 825x return 0;
589 }
590
591 3311x reactor_thread_context_guard ctx(this);
592
1/2
✓ Branch 0 taken 3311 times.
✗ Branch 1 not taken.
3311x lock_type lock(mutex_);
593
1/2
✓ Branch 0 taken 3311 times.
✗ Branch 1 not taken.
3311x return do_one(lock, usec, ctx.frame_);
594 4136x }
595
596 inline std::size_t
597 49x reactor_scheduler::poll()
598 {
599
2/2
✓ Branch 0 taken 34 times.
✓ Branch 1 taken 15 times.
49x if (outstanding_work_.load(std::memory_order_acquire) == 0)
600 {
601 15x stop();
602 15x return 0;
603 }
604
605 34x reactor_thread_context_guard ctx(this);
606
1/2
✓ Branch 0 taken 34 times.
✗ Branch 1 not taken.
34x lock_type lock(mutex_);
607
608 34x std::size_t n = 0;
609 74x for (;;)
610 {
611
3/4
✓ Branch 0 taken 74 times.
✗ Branch 1 not taken.
✓ Branch 2 taken 40 times.
✓ Branch 3 taken 34 times.
74x if (!do_one(lock, 0, ctx.frame_))
612 34x break;
613
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 40 times.
40x if (n != (std::numeric_limits<std::size_t>::max)())
614 40x ++n;
615
1/2
✓ Branch 0 taken 40 times.
✗ Branch 1 not taken.
40x if (!lock.owns_lock())
616
1/2
✓ Branch 0 taken 40 times.
✗ Branch 1 not taken.
40x lock.lock();
617 }
618 34x return n;
619 49x }
620
621 inline std::size_t
622 11x reactor_scheduler::poll_one()
623 {
624
2/2
✓ Branch 0 taken 6 times.
✓ Branch 1 taken 5 times.
11x if (outstanding_work_.load(std::memory_order_acquire) == 0)
625 {
626 5x stop();
627 5x return 0;
628 }
629
630 6x reactor_thread_context_guard ctx(this);
631
1/2
✓ Branch 0 taken 6 times.
✗ Branch 1 not taken.
6x lock_type lock(mutex_);
632
1/2
✓ Branch 0 taken 6 times.
✗ Branch 1 not taken.
6x return do_one(lock, 0, ctx.frame_);
633 11x }
634
635 inline void
636 149965x reactor_scheduler::work_started() noexcept
637 {
638 149965x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
639 149965x }
640
641 inline void
642 184634x reactor_scheduler::work_finished() noexcept
643 {
644
2/2
✓ Branch 0 taken 177437 times.
✓ Branch 1 taken 7197 times.
184634x if (outstanding_work_.fetch_sub(1, std::memory_order_acq_rel) == 1)
645
1/2
✓ Branch 0 taken 7197 times.
✗ Branch 1 not taken.
7197x stop();
646 184634x }
647
648 inline void
649 9094x reactor_scheduler::compensating_work_started() const noexcept
650 {
651 9094x auto* ctx = reactor_find_context(this);
652
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 9094 times.
9094x if (ctx)
653 9094x ++ctx->private_outstanding_work;
654 9094x }
655
656 inline void
657 115366x reactor_scheduler::post_deferred_completions(ready_queue& ops) const
658 {
659
2/2
✓ Branch 0 taken 115364 times.
✓ Branch 1 taken 2 times.
115366x if (ops.empty())
660 115364x return;
661
662
1/2
✓ Branch 0 taken 2 times.
✗ Branch 1 not taken.
2x if (auto* ctx = reactor_find_context(this))
663 {
664 2x ctx->private_queue.splice(ops);
665 2x return;
666 }
667
668 ✗ lock_type lock(mutex_);
669 ✗ completed_ops_.splice(ops);
670 ✗ wake_one_thread_and_unlock(lock);
671 115366x }
672
673 inline void
674 2967x reactor_scheduler::shutdown_drain()
675 {
676 2967x lock_type lock(mutex_);
677
678
2/2
✓ Branch 0 taken 2967 times.
✓ Branch 1 taken 3192 times.
6159x while (auto e = completed_ops_.pop())
679 {
680
2/2
✓ Branch 0 taken 6 times.
✓ Branch 1 taken 3186 times.
3192x if (ready_is_continuation(e))
681 {
682
1/2
✓ Branch 0 taken 6 times.
✗ Branch 1 not taken.
6x lock.unlock();
683
1/2
✓ Branch 0 taken 6 times.
✗ Branch 1 not taken.
6x if (auto h = ready_as_cont(e)->h)
684
1/2
✓ Branch 0 taken 6 times.
✗ Branch 1 not taken.
6x h.destroy();
685
1/2
✓ Branch 0 taken 6 times.
✗ Branch 1 not taken.
6x lock.lock();
686 6x }
687 else
688 {
689 3186x auto* op = ready_as_op(e);
690
2/2
✓ Branch 0 taken 221 times.
✓ Branch 1 taken 2965 times.
3186x if (op == &task_op_)
691 2965x continue;
692
1/2
✓ Branch 0 taken 221 times.
✗ Branch 1 not taken.
221x lock.unlock();
693
1/2
✓ Branch 0 taken 221 times.
✗ Branch 1 not taken.
221x op->destroy();
694
1/2
✓ Branch 0 taken 221 times.
✗ Branch 1 not taken.
221x lock.lock();
695 }
696 }
697
698
1/2
✓ Branch 0 taken 2967 times.
✗ Branch 1 not taken.
2967x signal_all(lock);
699 2967x }
700
701 inline void
702 10228x reactor_scheduler::signal_all(lock_type&) const
703 {
704 10228x state_ |= signaled_bit;
705 10228x cond_.notify_all();
706 10228x }
707
708 inline bool
709 24266x reactor_scheduler::maybe_unlock_and_signal_one(lock_type& lock) const
710 {
711 24266x state_ |= signaled_bit;
712
2/2
✓ Branch 0 taken 27 times.
✓ Branch 1 taken 24239 times.
24266x if (state_ > signaled_bit)
713 {
714 27x lock.unlock();
715 27x cond_.notify_one();
716 27x return true;
717 }
718 24239x return false;
719 24266x }
720
721 inline bool
722 619186x reactor_scheduler::unlock_and_signal_one(lock_type& lock) const
723 {
724 619186x state_ |= signaled_bit;
725 619186x bool have_waiters = state_ > signaled_bit;
726 619186x lock.unlock();
727
2/2
✓ Branch 0 taken 619167 times.
✓ Branch 1 taken 19 times.
619186x if (have_waiters)
728 19x cond_.notify_one();
729 619186x return have_waiters;
730 }
731
732 inline void
733 19x reactor_scheduler::clear_signal() const
734 {
735 19x state_ &= ~signaled_bit;
736 19x }
737
738 inline void
739 11x reactor_scheduler::wait_for_signal(lock_type& lock) const
740 {
741
2/2
✓ Branch 0 taken 13 times.
✓ Branch 1 taken 11 times.
24x while ((state_ & signaled_bit) == 0)
742 {
743 13x state_ += waiter_increment;
744 13x cond_.wait(lock);
745 13x state_ -= waiter_increment;
746 }
747 11x }
748
749 inline void
750 8x reactor_scheduler::wait_for_signal_for(lock_type& lock, long timeout_us) const
751 {
752
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 8 times.
8x if ((state_ & signaled_bit) == 0)
753 {
754 8x state_ += waiter_increment;
755 8x cond_.wait_for(lock, std::chrono::microseconds(timeout_us));
756 8x state_ -= waiter_increment;
757 8x }
758 8x }
759
760 inline void
761 24266x reactor_scheduler::wake_one_thread_and_unlock(lock_type& lock) const
762 {
763
2/2
✓ Branch 0 taken 27 times.
✓ Branch 1 taken 24239 times.
24266x if (maybe_unlock_and_signal_one(lock))
764 27x return;
765
766
4/4
✓ Branch 0 taken 1207 times.
✓ Branch 1 taken 23032 times.
✓ Branch 2 taken 36 times.
✓ Branch 3 taken 1171 times.
24239x if (task_running_.load(std::memory_order_relaxed) && !task_interrupted_)
767 {
768 1171x task_interrupted_ = true;
769 1171x lock.unlock();
770 1171x interrupt_reactor();
771 1171x }
772 else
773 {
774 23068x lock.unlock();
775 }
776 24266x }
777
778 1078089x inline reactor_scheduler::work_cleanup::~work_cleanup()
779 539032x {
780 539057x std::int64_t produced = ctx.private_outstanding_work;
781
2/2
✓ Branch 0 taken 347 times.
✓ Branch 1 taken 538710 times.
539057x if (produced > 1)
782 694x sched->outstanding_work_.fetch_add(
783 347x produced - 1, std::memory_order_relaxed);
784
2/2
✓ Branch 0 taken 388488 times.
✓ Branch 1 taken 150222 times.
538710x else if (produced < 1)
785 150222x sched->work_finished();
786 539057x ctx.private_outstanding_work = 0;
787
788
2/2
✓ Branch 0 taken 159010 times.
✓ Branch 1 taken 380047 times.
539057x if (!ctx.private_queue.empty())
789 {
790
1/2
✓ Branch 0 taken 380047 times.
✗ Branch 1 not taken.
380047x lock->lock();
791 380047x sched->completed_ops_.splice(ctx.private_queue);
792 380047x }
793 1078089x }
794
795 461568x inline reactor_scheduler::task_cleanup::~task_cleanup()
796 230784x {
797
2/2
✓ Branch 0 taken 223316 times.
✓ Branch 1 taken 7468 times.
230784x if (ctx.private_outstanding_work > 0)
798 {
799 14936x sched->outstanding_work_.fetch_add(
800 7468x ctx.private_outstanding_work, std::memory_order_relaxed);
801 7468x ctx.private_outstanding_work = 0;
802 7468x }
803
804
2/2
✓ Branch 0 taken 223316 times.
✓ Branch 1 taken 7468 times.
230784x if (!ctx.private_queue.empty())
805 {
806
1/2
✓ Branch 0 taken 7468 times.
✗ Branch 1 not taken.
7468x if (!lock->owns_lock())
807 ✗ lock->lock();
808 7468x sched->completed_ops_.splice(ctx.private_queue);
809 7468x }
810 461568x }
811
812 inline std::size_t
813 547173x reactor_scheduler::do_one(lock_type& lock, long timeout_us, context_type& ctx)
814 {
815 547192x for (;;)
816 {
817
2/2
✓ Branch 0 taken 769891 times.
✓ Branch 1 taken 6438 times.
776329x if (stopped_.load(std::memory_order_acquire))
818 6438x return 0;
819
820 769891x std::uintptr_t e = completed_ops_.pop();
821
2/2
✓ Branch 0 taken 32457 times.
✓ Branch 1 taken 737434 times.
769891x scheduler_op* op = ready_is_continuation(e) ? nullptr : ready_as_op(e);
822
823 // Handle reactor sentinel — time to poll for I/O
824
2/2
✓ Branch 0 taken 539068 times.
✓ Branch 1 taken 230823 times.
769891x if (op == &task_op_)
825 {
826 230823x bool more_handlers = !completed_ops_.empty();
827
828
4/4
✓ Branch 0 taken 150665 times.
✓ Branch 1 taken 80158 times.
✓ Branch 2 taken 150626 times.
✓ Branch 3 taken 32 times.
381481x if (!more_handlers &&
829
2/2
✓ Branch 0 taken 150658 times.
✓ Branch 1 taken 7 times.
150665x (outstanding_work_.load(std::memory_order_acquire) == 0 ||
830 150658x timeout_us == 0))
831 {
832 39x completed_ops_.push(&task_op_);
833 39x return 0;
834 }
835
836
2/2
✓ Branch 0 taken 80158 times.
✓ Branch 1 taken 150626 times.
230784x long task_timeout_us = more_handlers ? 0 : timeout_us;
837 230784x task_interrupted_ = task_timeout_us == 0;
838 230784x task_running_.store(true, std::memory_order_release);
839
840 // Wake a peer to take the pending handlers while this thread
841 // polls the reactor; skipped when one_thread_ (no peer exists).
842
4/4
✓ Branch 0 taken 80158 times.
✓ Branch 1 taken 150626 times.
✓ Branch 2 taken 7 times.
✓ Branch 3 taken 80151 times.
230784x if (more_handlers && !one_thread_)
843 80151x unlock_and_signal_one(lock);
844
845 try
846 {
847
2/2
✓ Branch 0 taken 2 times.
✓ Branch 1 taken 230782 times.
230784x run_task(lock, ctx, task_timeout_us);
848 230784x }
849 catch (...)
850 {
851 2x task_running_.store(false, std::memory_order_relaxed);
852
1/2
✗ Branch 0 not taken.
✓ Branch 1 taken 2 times.
2x throw;
853 2x }
854
855 230782x task_running_.store(false, std::memory_order_relaxed);
856 230782x completed_ops_.push(&task_op_);
857
2/2
✓ Branch 0 taken 1645 times.
✓ Branch 1 taken 229137 times.
230782x if (timeout_us > 0)
858 1645x return 0;
859 229137x continue;
860 }
861
862 // Handle ready entry (op or continuation)
863
2/2
✓ Branch 0 taken 23 times.
✓ Branch 1 taken 539045 times.
539068x if (e != 0)
864 {
865 539045x bool more = !completed_ops_.empty();
866
867
4/4
✓ Branch 0 taken 539043 times.
✓ Branch 1 taken 2 times.
✓ Branch 2 taken 8 times.
✓ Branch 3 taken 539035 times.
539045x if (more && !one_thread_)
868 {
869 // Wake a peer for the remaining work; unassisted if none
870 // was parked to take it.
871 539035x ctx.unassisted = !unlock_and_signal_one(lock);
872 539035x }
873 else
874 {
875 // No peer to wake (one_thread_, or nothing more queued).
876 10x ctx.unassisted = more;
877 10x lock.unlock();
878 }
879
880 539041x [[maybe_unused]] work_cleanup on_exit{this, &lock, ctx};
881
882
2/2
✓ Branch 0 taken 506584 times.
✓ Branch 1 taken 32457 times.
539041x if (ready_is_continuation(e))
883
2/2
✓ Branch 0 taken 32455 times.
✓ Branch 1 taken 2 times.
32457x ready_as_cont(e)->h.resume();
884 else
885
2/2
✓ Branch 0 taken 506582 times.
✓ Branch 1 taken 2 times.
506584x (*op)();
886 539037x return 1;
887 539041x }
888
889
3/4
✓ Branch 0 taken 19 times.
✓ Branch 1 taken 4 times.
✗ Branch 2 not taken.
✓ Branch 3 taken 19 times.
23x if (outstanding_work_.load(std::memory_order_acquire) == 0 ||
890 19x timeout_us == 0)
891 4x return 0;
892
893 19x clear_signal();
894
2/2
✓ Branch 0 taken 8 times.
✓ Branch 1 taken 11 times.
19x if (timeout_us < 0)
895 11x wait_for_signal(lock);
896 else
897 8x wait_for_signal_for(lock, timeout_us);
898 }
899 547171x }
900
901 } // namespace boost::corosio::detail
902
903 #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_SCHEDULER_HPP
904