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

96.0% Lines (316 / 329, 2 excl) 100.0% Functions (42 / 42, 2 excl)
reactor_scheduler.hpp
f(x) Functions (44)
Function Calls Lines Blocks
boost::corosio::detail::reactor_find_context(boost::corosio::detail::reactor_scheduler const*) :78 1011750x 100.0% 86.0% boost::corosio::detail::reactor_scheduler::inline_budget_initial() const :210 6283x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::scheduler_locking_disabled() const :216 581x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::configure_threading(boost::corosio::detail::scheduler::threading_config) :221 2310x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::reactor_scheduler() :238 2322x 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 6283x 100.0% 100.0% boost::corosio::detail::reactor_thread_context_guard::~reactor_thread_context_guard() :350 6283x 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 6283x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::configure_reactor(unsigned int, unsigned int, unsigned int, unsigned int) :370 44x 100.0% 87.0% boost::corosio::detail::reactor_scheduler::reset_inline_budget() const :395 102736x 53.3% 50.0% boost::corosio::detail::reactor_scheduler::try_consume_inline_budget() const :423 457251x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const :439 3668x 100.0% 84.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::post_handler(std::__n4861::coroutine_handle<void>) :445 3668x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::~post_handler() :446 7336x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::operator()() :448 3656x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::destroy() :455 12x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(boost::corosio::detail::scheduler_op*) const :480 116230x 100.0% 87.0% boost::corosio::detail::reactor_scheduler::post(boost::capy::continuation&) const :497 28879x 100.0% 87.0% boost::corosio::detail::reactor_scheduler::running_in_this_thread() const :514 11153x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::stop() :520 3857x 100.0% 82.0% boost::corosio::detail::reactor_scheduler::stopped() const :532 2693x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::restart() :538 1678x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::run() :544 2354x 100.0% 86.0% boost::corosio::detail::reactor_scheduler::run_one() :569 118x 100.0% 70.0% boost::corosio::detail::reactor_scheduler::wait_one(long) :583 4540x 100.0% 70.0% boost::corosio::detail::reactor_scheduler::poll() :597 49x 100.0% 76.0% boost::corosio::detail::reactor_scheduler::poll_one() :622 11x 100.0% 70.0% boost::corosio::detail::reactor_scheduler::work_started() :636 40295x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::work_finished() :642 77056x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::compensating_work_started() const :649 291927x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post_deferred_completions(boost::corosio::detail::ready_queue&) const :657 10474x 60.0% 59.0% boost::corosio::detail::reactor_scheduler::shutdown_drain() :674 2310x 100.0% 88.0% boost::corosio::detail::reactor_scheduler::signal_all(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :702 5427x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::maybe_unlock_and_signal_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :709 16465x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::unlock_and_signal_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :722 510458x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::clear_signal() const :733 22x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wait_for_signal(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :739 10x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wait_for_signal_for(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&, long) const :750 12x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wake_one_thread_and_unlock(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :761 16465x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::work_cleanup::~work_cleanup() :778 451096x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::task_cleanup::~task_cleanup() :795 326062x 90.0% 91.0% boost::corosio::detail::reactor_scheduler::do_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&, long, boost::corosio::detail::reactor_scheduler_context&) :813 455357x 97.8% 84.0%
Line 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 1011750x reactor_find_context(reactor_scheduler const* self) noexcept
79 {
80 1011750x for (auto* c = reactor_context_stack.get(); c != nullptr; c = c->next)
81 {
82 985011x if (c->key == self)
83 985011x return c;
84 }
85 26739x return nullptr;
86 }
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 6283x unsigned inline_budget_initial() const noexcept
211 {
212 6283x return inline_budget_initial_;
213 }
214
215 /// Return true when scheduler locking is disabled (fully-lockless tier).
216 581x bool scheduler_locking_disabled() const noexcept override
217 {
218 581x return scheduler_locking_disabled_;
219 }
220
221 2310x void configure_threading(threading_config cfg) noexcept override
222 {
223 2310x scheduler_locking_disabled_ = !cfg.scheduler_locking;
224 // reactor_io_locking takes effect at descriptor registration (see the
225 // register_descriptor overrides), not here.
226 2310x reactor_io_locking_ = cfg.reactor_io_locking;
227 2310x one_thread_ = cfg.one_thread;
228 2310x mutex_.set_enabled(cfg.scheduler_locking);
229 2310x cond_.set_enabled(cfg.scheduler_locking);
230 2310x }
231
232 protected:
233 timer_service* timer_svc_ = nullptr;
234 bool scheduler_locking_disabled_ = false;
235 bool reactor_io_locking_ = true;
236 bool one_thread_ = false;
237
238 2322x 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 mutable mutex_type mutex_{true};
259 mutable event_type cond_{true};
260 mutable ready_queue completed_ops_;
261 mutable std::atomic<std::int64_t> outstanding_work_{0};
262 std::atomic<bool> stopped_{false};
263 mutable std::atomic<bool> task_running_{false};
264 mutable bool task_interrupted_ = false;
265
266 // Runtime-configurable reactor tuning parameters.
267 // Defaults match the library's built-in values.
268 unsigned max_events_per_poll_ = 128;
269 unsigned inline_budget_initial_ = 2;
270 unsigned inline_budget_max_ = 16;
271 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 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 6283x explicit reactor_thread_context_guard(
338 reactor_scheduler const* sched) noexcept
339 6283x : frame_(sched, reactor_context_stack.get())
340 {
341 6283x reactor_context_stack.set(&frame_);
342 6283x }
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 6283x ~reactor_thread_context_guard() noexcept
351 {
352 6283x reactor_context_stack.set(frame_.next);
353 6283x }
354 };
355
356 // ---- Inline implementations ------------------------------------------------
357
358 6283x inline reactor_scheduler_context::reactor_scheduler_context(
359 6283x reactor_scheduler const* k, reactor_scheduler_context* n)
360 6283x : key(k)
361 6283x , next(n)
362 6283x , private_outstanding_work(0)
363 6283x , inline_budget(0)
364 6283x , inline_budget_max(static_cast<int>(k->inline_budget_initial()))
365 6283x , unassisted(false)
366 {
367 6283x }
368
369 inline void
370 44x reactor_scheduler::configure_reactor(
371 unsigned max_events,
372 unsigned budget_init,
373 unsigned budget_max,
374 unsigned unassisted)
375 {
376 86x if (max_events < 1 ||
377 42x max_events > static_cast<unsigned>(std::numeric_limits<int>::max()))
378 2x throw std::out_of_range("max_events_per_poll must be in [1, INT_MAX]");
379 42x if (budget_max > static_cast<unsigned>(std::numeric_limits<int>::max()))
380 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 40x if (budget_init > budget_max)
384 16x budget_init = budget_max;
385 40x if (unassisted > budget_max)
386 16x unassisted = budget_max;
387
388 40x max_events_per_poll_ = max_events;
389 40x inline_budget_initial_ = budget_init;
390 40x inline_budget_max_ = budget_max;
391 40x unassisted_budget_ = unassisted;
392 40x }
393
394 inline void
395 102736x 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 102736x if (inline_budget_max_ == 0)
400 56x return;
401 102680x if (auto* ctx = reactor_find_context(this))
402 {
403 // Cap when no other thread absorbed queued work
404 102680x if (ctx->unassisted)
405 {
406 102680x ctx->inline_budget_max = static_cast<int>(unassisted_budget_);
407 102680x ctx->inline_budget = static_cast<int>(unassisted_budget_);
408 102680x return;
409 }
410 // Ramp up when previous cycle fully consumed budget.
411 // max(1, ...) ensures the doubling escapes zero.
412 ✗ if (ctx->inline_budget == 0)
413 ✗ ctx->inline_budget_max =
414 ✗ (std::min)((std::max)(1, ctx->inline_budget_max) * 2,
415 ✗ static_cast<int>(inline_budget_max_));
416 ✗ else if (ctx->inline_budget < ctx->inline_budget_max)
417 ✗ ctx->inline_budget_max = static_cast<int>(inline_budget_initial_);
418 ✗ ctx->inline_budget = ctx->inline_budget_max;
419 }
420 }
421
422 inline bool
423 457251x reactor_scheduler::try_consume_inline_budget() const noexcept
424 {
425 457251x if (inline_budget_max_ == 0)
426 40x return false;
427 457211x if (auto* ctx = reactor_find_context(this))
428 {
429 457211x if (ctx->inline_budget > 0)
430 {
431 365586x --ctx->inline_budget;
432 365586x return true;
433 }
434 }
435 91625x return false;
436 }
437
438 inline void
439 3668x reactor_scheduler::post(std::coroutine_handle<> h) const
440 {
441 struct post_handler final : scheduler_op
442 {
443 std::coroutine_handle<> h_;
444
445 3668x explicit post_handler(std::coroutine_handle<> h) : h_(h) {}
446 7336x ~post_handler() override = default;
447
448 3656x void operator()() override
449 {
450 3656x auto saved = h_;
451 3656x delete this;
452 3656x saved.resume();
453 3656x }
454
455 12x void destroy() override
456 {
457 12x auto saved = h_;
458 12x delete this;
459 12x saved.destroy();
460 12x }
461 };
462
463 3668x auto ph = std::make_unique<post_handler>(h);
464
465 3668x if (auto* ctx = reactor_find_context(this))
466 {
467 8x ++ctx->private_outstanding_work;
468 8x ctx->private_queue.push(ph.release());
469 8x return;
470 }
471
472 3660x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
473
474 3660x lock_type lock(mutex_);
475 3660x completed_ops_.push(ph.release());
476 3660x wake_one_thread_and_unlock(lock);
477 3668x }
478
479 inline void
480 116230x reactor_scheduler::post(scheduler_op* h) const
481 {
482 116230x if (auto* ctx = reactor_find_context(this))
483 {
484 114425x ++ctx->private_outstanding_work;
485 114425x ctx->private_queue.push(h);
486 114425x return;
487 }
488
489 1805x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
490
491 1805x lock_type lock(mutex_);
492 1805x completed_ops_.push(h);
493 1805x wake_one_thread_and_unlock(lock);
494 1805x }
495
496 inline void
497 28879x reactor_scheduler::post(capy::continuation& c) const
498 {
499 28879x if (auto* ctx = reactor_find_context(this))
500 {
501 17879x ++ctx->private_outstanding_work;
502 17879x ctx->private_queue.push(c);
503 17879x return;
504 }
505
506 11000x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
507
508 11000x lock_type lock(mutex_);
509 11000x completed_ops_.push(c);
510 11000x wake_one_thread_and_unlock(lock);
511 11000x }
512
513 inline bool
514 11153x reactor_scheduler::running_in_this_thread() const noexcept
515 {
516 11153x return reactor_find_context(this) != nullptr;
517 }
518
519 inline void
520 3857x reactor_scheduler::stop()
521 {
522 3857x lock_type lock(mutex_);
523 3857x if (!stopped_.load(std::memory_order_acquire))
524 {
525 3117x stopped_.store(true, std::memory_order_release);
526 3117x signal_all(lock);
527 3117x interrupt_reactor();
528 }
529 3857x }
530
531 inline bool
532 2693x reactor_scheduler::stopped() const noexcept
533 {
534 2693x return stopped_.load(std::memory_order_acquire);
535 }
536
537 inline void
538 1678x reactor_scheduler::restart()
539 {
540 1678x stopped_.store(false, std::memory_order_release);
541 1678x }
542
543 inline std::size_t
544 2354x reactor_scheduler::run()
545 {
546 4708x if (outstanding_work_.load(std::memory_order_acquire) == 0)
547 {
548 96x stop();
549 96x return 0;
550 }
551
552 2258x reactor_thread_context_guard ctx(this);
553 2258x lock_type lock(mutex_);
554
555 2258x std::size_t n = 0;
556 for (;;)
557 {
558 451291x if (!do_one(lock, -1, ctx.frame_))
559 2255x break;
560 449033x if (n != (std::numeric_limits<std::size_t>::max)())
561 449033x ++n;
562 449033x if (!lock.owns_lock())
563 337010x lock.lock();
564 }
565 2255x return n;
566 2261x }
567
568 inline std::size_t
569 118x reactor_scheduler::run_one()
570 {
571 236x if (outstanding_work_.load(std::memory_order_acquire) == 0)
572 {
573 3x stop();
574 3x return 0;
575 }
576
577 115x reactor_thread_context_guard ctx(this);
578 115x lock_type lock(mutex_);
579 115x return do_one(lock, -1, ctx.frame_);
580 115x }
581
582 inline std::size_t
583 4540x reactor_scheduler::wait_one(long usec)
584 {
585 9080x if (outstanding_work_.load(std::memory_order_acquire) == 0)
586 {
587 670x stop();
588 670x return 0;
589 }
590
591 3870x reactor_thread_context_guard ctx(this);
592 3870x lock_type lock(mutex_);
593 3870x return do_one(lock, usec, ctx.frame_);
594 3870x }
595
596 inline std::size_t
597 49x reactor_scheduler::poll()
598 {
599 98x 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 34x lock_type lock(mutex_);
607
608 34x std::size_t n = 0;
609 for (;;)
610 {
611 75x if (!do_one(lock, 0, ctx.frame_))
612 34x break;
613 41x if (n != (std::numeric_limits<std::size_t>::max)())
614 41x ++n;
615 41x if (!lock.owns_lock())
616 41x lock.lock();
617 }
618 34x return n;
619 34x }
620
621 inline std::size_t
622 11x reactor_scheduler::poll_one()
623 {
624 22x 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 6x lock_type lock(mutex_);
632 6x return do_one(lock, 0, ctx.frame_);
633 6x }
634
635 inline void
636 40295x reactor_scheduler::work_started() noexcept
637 {
638 40295x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
639 40295x }
640
641 inline void
642 77056x reactor_scheduler::work_finished() noexcept
643 {
644 154112x if (outstanding_work_.fetch_sub(1, std::memory_order_acq_rel) == 1)
645 3053x stop();
646 77056x }
647
648 inline void
649 291927x reactor_scheduler::compensating_work_started() const noexcept
650 {
651 291927x auto* ctx = reactor_find_context(this);
652 291927x if (ctx)
653 291927x ++ctx->private_outstanding_work;
654 291927x }
655
656 inline void
657 10474x reactor_scheduler::post_deferred_completions(ready_queue& ops) const
658 {
659 10474x if (ops.empty())
660 10474x return;
661
662 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 ✗ }
672
673 inline void
674 2310x reactor_scheduler::shutdown_drain()
675 {
676 2310x lock_type lock(mutex_);
677
678 5004x while (auto e = completed_ops_.pop())
679 {
680 2694x if (ready_is_continuation(e))
681 {
682 8x lock.unlock();
683 8x if (auto h = ready_as_cont(e)->h)
684 8x h.destroy();
685 8x lock.lock();
686 }
687 else
688 {
689 2686x auto* op = ready_as_op(e);
690 2686x if (op == &task_op_)
691 2307x continue;
692 379x lock.unlock();
693 379x op->destroy();
694 379x lock.lock();
695 }
696 2694x }
697
698 2310x signal_all(lock);
699 2310x }
700
701 inline void
702 5427x reactor_scheduler::signal_all(lock_type&) const
703 {
704 5427x state_ |= signaled_bit;
705 5427x cond_.notify_all();
706 5427x }
707
708 inline bool
709 16465x reactor_scheduler::maybe_unlock_and_signal_one(lock_type& lock) const
710 {
711 16465x state_ |= signaled_bit;
712 16465x if (state_ > signaled_bit)
713 {
714 5x lock.unlock();
715 5x cond_.notify_one();
716 5x return true;
717 }
718 16460x return false;
719 }
720
721 inline bool
722 510458x reactor_scheduler::unlock_and_signal_one(lock_type& lock) const
723 {
724 510458x state_ |= signaled_bit;
725 510458x bool have_waiters = state_ > signaled_bit;
726 510458x lock.unlock();
727 510458x if (have_waiters)
728 308x cond_.notify_one();
729 510458x return have_waiters;
730 }
731
732 inline void
733 22x reactor_scheduler::clear_signal() const
734 {
735 22x state_ &= ~signaled_bit;
736 22x }
737
738 inline void
739 10x reactor_scheduler::wait_for_signal(lock_type& lock) const
740 {
741 22x while ((state_ & signaled_bit) == 0)
742 {
743 12x state_ += waiter_increment;
744 12x cond_.wait(lock);
745 12x state_ -= waiter_increment;
746 }
747 10x }
748
749 inline void
750 12x reactor_scheduler::wait_for_signal_for(lock_type& lock, long timeout_us) const
751 {
752 12x if ((state_ & signaled_bit) == 0)
753 {
754 12x state_ += waiter_increment;
755 12x cond_.wait_for(lock, std::chrono::microseconds(timeout_us));
756 12x state_ -= waiter_increment;
757 }
758 12x }
759
760 inline void
761 16465x reactor_scheduler::wake_one_thread_and_unlock(lock_type& lock) const
762 {
763 16465x if (maybe_unlock_and_signal_one(lock))
764 5x return;
765
766 16460x if (task_running_.load(std::memory_order_relaxed) && !task_interrupted_)
767 {
768 1180x task_interrupted_ = true;
769 1180x lock.unlock();
770 1180x interrupt_reactor();
771 }
772 else
773 {
774 15280x lock.unlock();
775 }
776 }
777
778 451096x inline reactor_scheduler::work_cleanup::~work_cleanup()
779 {
780 451096x std::int64_t produced = ctx.private_outstanding_work;
781 451096x if (produced > 1)
782 347x sched->outstanding_work_.fetch_add(
783 produced - 1, std::memory_order_relaxed);
784 450749x else if (produced < 1)
785 47305x sched->work_finished();
786 451096x ctx.private_outstanding_work = 0;
787
788 451096x if (!ctx.private_queue.empty())
789 {
790 112185x lock->lock();
791 112185x sched->completed_ops_.splice(ctx.private_queue);
792 }
793 451096x }
794
795 326062x inline reactor_scheduler::task_cleanup::~task_cleanup()
796 {
797 326062x if (ctx.private_outstanding_work > 0)
798 {
799 10433x sched->outstanding_work_.fetch_add(
800 10433x ctx.private_outstanding_work, std::memory_order_relaxed);
801 10433x ctx.private_outstanding_work = 0;
802 }
803
804 326062x if (!ctx.private_queue.empty())
805 {
806 10433x if (!lock->owns_lock())
807 ✗ lock->lock();
808 10433x sched->completed_ops_.splice(ctx.private_queue);
809 }
810 326062x }
811
812 inline std::size_t
813 455357x reactor_scheduler::do_one(lock_type& lock, long timeout_us, context_type& ctx)
814 {
815 for (;;)
816 {
817 779471x if (stopped_.load(std::memory_order_acquire))
818 2259x return 0;
819
820 777212x std::uintptr_t e = completed_ops_.pop();
821 777212x scheduler_op* op = ready_is_continuation(e) ? nullptr : ready_as_op(e);
822
823 // Handle reactor sentinel — time to poll for I/O
824 777212x if (op == &task_op_)
825 {
826 326094x bool more_handlers = !completed_ops_.empty();
827
828 592765x if (!more_handlers &&
829 533342x (outstanding_work_.load(std::memory_order_acquire) == 0 ||
830 timeout_us == 0))
831 {
832 32x completed_ops_.push(&task_op_);
833 32x return 0;
834 }
835
836 326062x long task_timeout_us = more_handlers ? 0 : timeout_us;
837 326062x task_interrupted_ = task_timeout_us == 0;
838 326062x 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 326062x if (more_handlers && !one_thread_)
843 59416x unlock_and_signal_one(lock);
844
845 try
846 {
847 326062x run_task(lock, ctx, task_timeout_us);
848 }
849 3x catch (...)
850 {
851 3x task_running_.store(false, std::memory_order_relaxed);
852 3x throw;
853 3x }
854
855 326059x task_running_.store(false, std::memory_order_relaxed);
856 326059x completed_ops_.push(&task_op_);
857 326059x if (timeout_us > 0)
858 1967x return 0;
859 324092x continue;
860 324092x }
861
862 // Handle ready entry (op or continuation)
863 451118x if (e != 0)
864 {
865 451096x bool more = !completed_ops_.empty();
866
867 451096x if (more && !one_thread_)
868 {
869 // Wake a peer for the remaining work; unassisted if none
870 // was parked to take it.
871 451042x ctx.unassisted = !unlock_and_signal_one(lock);
872 }
873 else
874 {
875 // No peer to wake (one_thread_, or nothing more queued).
876 54x ctx.unassisted = more;
877 54x lock.unlock();
878 }
879
880 451096x [[maybe_unused]] work_cleanup on_exit{this, &lock, ctx};
881
882 451096x if (ready_is_continuation(e))
883 28871x ready_as_cont(e)->h.resume();
884 else
885 422225x (*op)();
886 451096x return 1;
887 451096x }
888
889 44x if (outstanding_work_.load(std::memory_order_acquire) == 0 ||
890 timeout_us == 0)
891 ✗ return 0;
892
893 22x clear_signal();
894 22x if (timeout_us < 0)
895 10x wait_for_signal(lock);
896 else
897 12x wait_for_signal_for(lock, timeout_us);
898 324114x }
899 }
900
901 } // namespace boost::corosio::detail
902
903 #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_SCHEDULER_HPP
904