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

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