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

96.4% Lines (317 / 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 878613x 100.0% 86.0% boost::corosio::detail::reactor_scheduler::inline_budget_initial() const :212 5451x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::scheduler_locking_disabled() const :218 502x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::configure_threading(boost::corosio::detail::scheduler::threading_config) :223 2241x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::reactor_scheduler() :240 2253x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::task_op::operator()() :287 – – – boost::corosio::detail::reactor_scheduler::task_op::destroy() :288 – – – boost::corosio::detail::reactor_thread_context_guard::reactor_thread_context_guard(boost::corosio::detail::reactor_scheduler const*) :339 5451x 100.0% 100.0% boost::corosio::detail::reactor_thread_context_guard::~reactor_thread_context_guard() :352 5451x 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*) :360 5451x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::configure_reactor(unsigned int, unsigned int, unsigned int, unsigned int) :372 44x 100.0% 87.0% boost::corosio::detail::reactor_scheduler::reset_inline_budget() const :397 86044x 53.3% 50.0% boost::corosio::detail::reactor_scheduler::try_consume_inline_budget() const :425 395233x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const :441 3756x 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>) :447 3756x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::~post_handler() :448 7512x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::operator()() :450 3744x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::destroy() :457 12x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(boost::corosio::detail::scheduler_op*) const :482 97697x 100.0% 87.0% boost::corosio::detail::reactor_scheduler::post(boost::capy::continuation&) const :499 24257x 100.0% 87.0% boost::corosio::detail::reactor_scheduler::running_in_this_thread() const :516 10717x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::stop() :522 3733x 100.0% 82.0% boost::corosio::detail::reactor_scheduler::stopped() const :534 2506x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::restart() :540 1439x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::run() :546 2102x 100.0% 86.0% boost::corosio::detail::reactor_scheduler::run_one() :571 112x 100.0% 70.0% boost::corosio::detail::reactor_scheduler::wait_one(long) :585 4093x 100.0% 70.0% boost::corosio::detail::reactor_scheduler::poll() :599 49x 100.0% 76.0% boost::corosio::detail::reactor_scheduler::poll_one() :624 11x 100.0% 70.0% boost::corosio::detail::reactor_scheduler::work_started() :638 31259x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::work_finished() :644 61730x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::compensating_work_started() const :651 261003x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post_deferred_completions(boost::corosio::detail::ready_queue&) const :659 6263x 60.0% 59.0% boost::corosio::detail::reactor_scheduler::shutdown_drain() :676 2241x 100.0% 88.0% boost::corosio::detail::reactor_scheduler::signal_all(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :704 5106x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::maybe_unlock_and_signal_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :711 15039x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::unlock_and_signal_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :724 445819x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::clear_signal() const :735 50x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wait_for_signal(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :741 6x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wait_for_signal_for(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&, long) const :752 44x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wake_one_thread_and_unlock(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :763 15039x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::work_cleanup::~work_cleanup() :780 392894x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::task_cleanup::~task_cleanup() :797 284131x 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&) :815 396589x 100.0% 86.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 878613x reactor_find_context(reactor_scheduler const* self) noexcept
79 {
80 878613x for (auto* c = reactor_context_stack.get(); c != nullptr; c = c->next)
81 {
82 853720x if (c->key == self)
83 853720x return c;
84 }
85 24893x 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 , 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 @pre work_started() must have been called for each operation.
190
191 @param ops Queue of operations to post.
192 */
193 void post_deferred_completions(ready_queue& ops) const;
194
195 /** Apply runtime configuration to the scheduler.
196
197 Called by `io_context` after construction. Values that do
198 not apply to this backend are silently ignored.
199
200 @param max_events Event buffer size for epoll/kqueue.
201 @param budget_init Starting inline completion budget.
202 @param budget_max Hard ceiling on adaptive budget ramp-up.
203 @param unassisted Budget when single-threaded.
204 */
205 virtual void configure_reactor(
206 unsigned max_events,
207 unsigned budget_init,
208 unsigned budget_max,
209 unsigned unassisted);
210
211 /// Return the configured initial inline budget.
212 5451x unsigned inline_budget_initial() const noexcept
213 {
214 5451x return inline_budget_initial_;
215 }
216
217 /// Return true when scheduler locking is disabled (fully-lockless tier).
218 502x bool scheduler_locking_disabled() const noexcept override
219 {
220 502x return scheduler_locking_disabled_;
221 }
222
223 2241x void configure_threading(threading_config cfg) noexcept override
224 {
225 2241x scheduler_locking_disabled_ = !cfg.scheduler_locking;
226 // reactor_io_locking takes effect at descriptor registration (see the
227 // register_descriptor overrides), not here.
228 2241x reactor_io_locking_ = cfg.reactor_io_locking;
229 2241x one_thread_ = cfg.one_thread;
230 2241x mutex_.set_enabled(cfg.scheduler_locking);
231 2241x cond_.set_enabled(cfg.scheduler_locking);
232 2241x }
233
234 protected:
235 timer_service* timer_svc_ = nullptr;
236 bool scheduler_locking_disabled_ = false;
237 bool reactor_io_locking_ = true;
238 bool one_thread_ = false;
239
240 2253x reactor_scheduler() = default;
241
242 /** Drain completed_ops during shutdown.
243
244 Pops all operations from the global queue and destroys them,
245 skipping the task sentinel. Signals all waiting threads.
246 Derived classes call this from their shutdown() override
247 before performing platform-specific cleanup.
248 */
249 void shutdown_drain();
250
251 /// RAII guard that re-inserts the task sentinel after `run_task`.
252 struct task_cleanup
253 {
254 reactor_scheduler const* sched;
255 lock_type* lock;
256 context_type& ctx;
257 ~task_cleanup();
258 };
259
260 mutable mutex_type mutex_{true};
261 mutable event_type cond_{true};
262 mutable ready_queue completed_ops_;
263 mutable std::atomic<std::int64_t> outstanding_work_{0};
264 std::atomic<bool> stopped_{false};
265 mutable std::atomic<bool> task_running_{false};
266 mutable bool task_interrupted_ = false;
267
268 // Runtime-configurable reactor tuning parameters.
269 // Defaults match the library's built-in values.
270 unsigned max_events_per_poll_ = 128;
271 unsigned inline_budget_initial_ = 2;
272 unsigned inline_budget_max_ = 16;
273 unsigned unassisted_budget_ = 4;
274
275 /// Bit 0 of `state_`: set when the condvar should be signaled.
276 static constexpr std::size_t signaled_bit = 1;
277
278 /// Increment per waiting thread in `state_`.
279 static constexpr std::size_t waiter_increment = 2;
280 mutable std::size_t state_ = 0;
281
282 /// Sentinel op that triggers a reactor poll when dequeued.
283 struct task_op final : scheduler_op
284 {
285 // LCOV_EXCL_START: the sentinel is intercepted by pointer
286 // identity; its virtuals exist for vtable completeness.
287 − void operator()() override {}
288 − void destroy() override {}
289 // LCOV_EXCL_STOP
290 };
291 task_op task_op_;
292
293 /** Run the platform-specific reactor poll.
294
295 @par Postconditions
296 `lock` is owned on return, however the poll ended. An
297 implementation that unlocks around the blocking call owes the
298 caller a matching re-acquire on every path out, including the
299 errors it retries rather than reports.
300 */
301 virtual void
302 run_task(lock_type& lock, context_type& ctx, long timeout_us) = 0;
303
304 /// Wake a blocked reactor (e.g. write to eventfd or pipe).
305 virtual void interrupt_reactor() const = 0;
306
307 private:
308 struct work_cleanup
309 {
310 reactor_scheduler* sched;
311 lock_type* lock;
312 context_type& ctx;
313 ~work_cleanup();
314 };
315
316 std::size_t do_one(lock_type& lock, long timeout_us, context_type& ctx);
317
318 void signal_all(lock_type& lock) const;
319 bool maybe_unlock_and_signal_one(lock_type& lock) const;
320 bool unlock_and_signal_one(lock_type& lock) const;
321 void clear_signal() const;
322 void wait_for_signal(lock_type& lock) const;
323 void wait_for_signal_for(lock_type& lock, long timeout_us) const;
324 void wake_one_thread_and_unlock(lock_type& lock) const;
325 };
326
327 /** RAII guard that pushes/pops a scheduler context frame.
328
329 On construction, pushes a new context frame onto the
330 thread-local stack. On destruction, drains any remaining
331 private queue items to the global queue and pops the frame.
332 */
333 struct reactor_thread_context_guard
334 {
335 /// The context frame managed by this guard.
336 reactor_scheduler_context frame_;
337
338 /// Construct the guard, pushing a frame for @a sched.
339 5451x explicit reactor_thread_context_guard(
340 reactor_scheduler const* sched) noexcept
341 5451x : frame_(sched, reactor_context_stack.get())
342 {
343 5451x reactor_context_stack.set(&frame_);
344 5451x }
345
346 /** Destroy the guard, popping the frame.
347
348 The private queue is empty here by invariant: work_cleanup and
349 task_cleanup splice it to the global queue after every handler
350 and every reactor pass.
351 */
352 5451x ~reactor_thread_context_guard() noexcept
353 {
354 5451x reactor_context_stack.set(frame_.next);
355 5451x }
356 };
357
358 // ---- Inline implementations ------------------------------------------------
359
360 5451x inline reactor_scheduler_context::reactor_scheduler_context(
361 5451x reactor_scheduler const* k, reactor_scheduler_context* n)
362 5451x : key(k)
363 5451x , next(n)
364 5451x , private_outstanding_work(0)
365 5451x , inline_budget(0)
366 5451x , inline_budget_max(static_cast<int>(k->inline_budget_initial()))
367 5451x , unassisted(false)
368 {
369 5451x }
370
371 inline void
372 44x reactor_scheduler::configure_reactor(
373 unsigned max_events,
374 unsigned budget_init,
375 unsigned budget_max,
376 unsigned unassisted)
377 {
378 86x if (max_events < 1 ||
379 42x max_events > static_cast<unsigned>(std::numeric_limits<int>::max()))
380 2x throw std::out_of_range("max_events_per_poll must be in [1, INT_MAX]");
381 42x if (budget_max > static_cast<unsigned>(std::numeric_limits<int>::max()))
382 2x throw std::out_of_range("inline_budget_max must be in [0, INT_MAX]");
383
384 // Clamp initial and unassisted to budget_max.
385 40x if (budget_init > budget_max)
386 16x budget_init = budget_max;
387 40x if (unassisted > budget_max)
388 16x unassisted = budget_max;
389
390 40x max_events_per_poll_ = max_events;
391 40x inline_budget_initial_ = budget_init;
392 40x inline_budget_max_ = budget_max;
393 40x unassisted_budget_ = unassisted;
394 40x }
395
396 inline void
397 86044x reactor_scheduler::reset_inline_budget() const noexcept
398 {
399 // When budget is disabled (max==0), all paths below would no-op
400 // (inline_budget stays 0). Skip the TLS lookup entirely.
401 86044x if (inline_budget_max_ == 0)
402 56x return;
403 85988x if (auto* ctx = reactor_find_context(this))
404 {
405 // Cap when no other thread absorbed queued work
406 85988x if (ctx->unassisted)
407 {
408 85988x ctx->inline_budget_max = static_cast<int>(unassisted_budget_);
409 85988x ctx->inline_budget = static_cast<int>(unassisted_budget_);
410 85988x return;
411 }
412 // Ramp up when previous cycle fully consumed budget.
413 // max(1, ...) ensures the doubling escapes zero.
414 ✗ if (ctx->inline_budget == 0)
415 ✗ ctx->inline_budget_max =
416 ✗ (std::min)((std::max)(1, ctx->inline_budget_max) * 2,
417 ✗ static_cast<int>(inline_budget_max_));
418 ✗ else if (ctx->inline_budget < ctx->inline_budget_max)
419 ✗ ctx->inline_budget_max = static_cast<int>(inline_budget_initial_);
420 ✗ ctx->inline_budget = ctx->inline_budget_max;
421 }
422 }
423
424 inline bool
425 395233x reactor_scheduler::try_consume_inline_budget() const noexcept
426 {
427 395233x if (inline_budget_max_ == 0)
428 40x return false;
429 395193x if (auto* ctx = reactor_find_context(this))
430 {
431 395193x if (ctx->inline_budget > 0)
432 {
433 315977x --ctx->inline_budget;
434 315977x return true;
435 }
436 }
437 79216x return false;
438 }
439
440 inline void
441 3756x reactor_scheduler::post(std::coroutine_handle<> h) const
442 {
443 struct post_handler final : scheduler_op
444 {
445 std::coroutine_handle<> h_;
446
447 3756x explicit post_handler(std::coroutine_handle<> h) : h_(h) {}
448 7512x ~post_handler() override = default;
449
450 3744x void operator()() override
451 {
452 3744x auto saved = h_;
453 3744x delete this;
454 3744x saved.resume();
455 3744x }
456
457 12x void destroy() override
458 {
459 12x auto saved = h_;
460 12x delete this;
461 12x saved.destroy();
462 12x }
463 };
464
465 3756x auto ph = std::make_unique<post_handler>(h);
466
467 3756x if (auto* ctx = reactor_find_context(this))
468 {
469 96x ++ctx->private_outstanding_work;
470 96x ctx->private_queue.push(ph.release());
471 96x return;
472 }
473
474 3660x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
475
476 3660x lock_type lock(mutex_);
477 3660x completed_ops_.push(ph.release());
478 3660x wake_one_thread_and_unlock(lock);
479 3756x }
480
481 inline void
482 97697x reactor_scheduler::post(scheduler_op* h) const
483 {
484 97697x if (auto* ctx = reactor_find_context(this))
485 {
486 96698x ++ctx->private_outstanding_work;
487 96698x ctx->private_queue.push(h);
488 96698x return;
489 }
490
491 999x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
492
493 999x lock_type lock(mutex_);
494 999x completed_ops_.push(h);
495 999x wake_one_thread_and_unlock(lock);
496 999x }
497
498 inline void
499 24257x reactor_scheduler::post(capy::continuation& c) const
500 {
501 24257x if (auto* ctx = reactor_find_context(this))
502 {
503 13877x ++ctx->private_outstanding_work;
504 13877x ctx->private_queue.push(c);
505 13877x return;
506 }
507
508 10380x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
509
510 10380x lock_type lock(mutex_);
511 10380x completed_ops_.push(c);
512 10380x wake_one_thread_and_unlock(lock);
513 10380x }
514
515 inline bool
516 10717x reactor_scheduler::running_in_this_thread() const noexcept
517 {
518 10717x return reactor_find_context(this) != nullptr;
519 }
520
521 inline void
522 3733x reactor_scheduler::stop()
523 {
524 3733x lock_type lock(mutex_);
525 3733x if (!stopped_.load(std::memory_order_acquire))
526 {
527 2865x stopped_.store(true, std::memory_order_release);
528 2865x signal_all(lock);
529 2865x interrupt_reactor();
530 }
531 3733x }
532
533 inline bool
534 2506x reactor_scheduler::stopped() const noexcept
535 {
536 2506x return stopped_.load(std::memory_order_acquire);
537 }
538
539 inline void
540 1439x reactor_scheduler::restart()
541 {
542 1439x stopped_.store(false, std::memory_order_release);
543 1439x }
544
545 inline std::size_t
546 2102x reactor_scheduler::run()
547 {
548 4204x if (outstanding_work_.load(std::memory_order_acquire) == 0)
549 {
550 101x stop();
551 101x return 0;
552 }
553
554 2001x reactor_thread_context_guard ctx(this);
555 2001x lock_type lock(mutex_);
556
557 2001x std::size_t n = 0;
558 for (;;)
559 {
560 393098x if (!do_one(lock, -1, ctx.frame_))
561 1998x break;
562 391097x if (n != (std::numeric_limits<std::size_t>::max)())
563 391097x ++n;
564 391097x if (!lock.owns_lock())
565 295953x lock.lock();
566 }
567 1998x return n;
568 2004x }
569
570 inline std::size_t
571 112x reactor_scheduler::run_one()
572 {
573 224x if (outstanding_work_.load(std::memory_order_acquire) == 0)
574 {
575 3x stop();
576 3x return 0;
577 }
578
579 109x reactor_thread_context_guard ctx(this);
580 109x lock_type lock(mutex_);
581 109x return do_one(lock, -1, ctx.frame_);
582 109x }
583
584 inline std::size_t
585 4093x reactor_scheduler::wait_one(long usec)
586 {
587 8186x if (outstanding_work_.load(std::memory_order_acquire) == 0)
588 {
589 792x stop();
590 792x return 0;
591 }
592
593 3301x reactor_thread_context_guard ctx(this);
594 3301x lock_type lock(mutex_);
595 3301x return do_one(lock, usec, ctx.frame_);
596 3301x }
597
598 inline std::size_t
599 49x reactor_scheduler::poll()
600 {
601 98x if (outstanding_work_.load(std::memory_order_acquire) == 0)
602 {
603 15x stop();
604 15x return 0;
605 }
606
607 34x reactor_thread_context_guard ctx(this);
608 34x lock_type lock(mutex_);
609
610 34x std::size_t n = 0;
611 for (;;)
612 {
613 75x if (!do_one(lock, 0, ctx.frame_))
614 34x break;
615 41x if (n != (std::numeric_limits<std::size_t>::max)())
616 41x ++n;
617 41x if (!lock.owns_lock())
618 41x lock.lock();
619 }
620 34x return n;
621 34x }
622
623 inline std::size_t
624 11x reactor_scheduler::poll_one()
625 {
626 22x if (outstanding_work_.load(std::memory_order_acquire) == 0)
627 {
628 5x stop();
629 5x return 0;
630 }
631
632 6x reactor_thread_context_guard ctx(this);
633 6x lock_type lock(mutex_);
634 6x return do_one(lock, 0, ctx.frame_);
635 6x }
636
637 inline void
638 31259x reactor_scheduler::work_started() noexcept
639 {
640 31259x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
641 31259x }
642
643 inline void
644 61730x reactor_scheduler::work_finished() noexcept
645 {
646 123460x if (outstanding_work_.fetch_sub(1, std::memory_order_acq_rel) == 1)
647 2802x stop();
648 61730x }
649
650 inline void
651 261003x reactor_scheduler::compensating_work_started() const noexcept
652 {
653 261003x auto* ctx = reactor_find_context(this);
654 261003x if (ctx)
655 261003x ++ctx->private_outstanding_work;
656 261003x }
657
658 inline void
659 6263x reactor_scheduler::post_deferred_completions(ready_queue& ops) const
660 {
661 6263x if (ops.empty())
662 6263x return;
663
664 2x if (auto* ctx = reactor_find_context(this))
665 {
666 2x ctx->private_queue.splice(ops);
667 2x return;
668 }
669
670 ✗ lock_type lock(mutex_);
671 ✗ completed_ops_.splice(ops);
672 ✗ wake_one_thread_and_unlock(lock);
673 ✗ }
674
675 inline void
676 2241x reactor_scheduler::shutdown_drain()
677 {
678 2241x lock_type lock(mutex_);
679
680 4862x while (auto e = completed_ops_.pop())
681 {
682 2621x if (ready_is_continuation(e))
683 {
684 8x lock.unlock();
685 8x if (auto h = ready_as_cont(e)->h)
686 8x h.destroy();
687 8x lock.lock();
688 }
689 else
690 {
691 2613x auto* op = ready_as_op(e);
692 2613x if (op == &task_op_)
693 2238x continue;
694 375x lock.unlock();
695 375x op->destroy();
696 375x lock.lock();
697 }
698 2621x }
699
700 2241x signal_all(lock);
701 2241x }
702
703 inline void
704 5106x reactor_scheduler::signal_all(lock_type&) const
705 {
706 5106x state_ |= signaled_bit;
707 5106x cond_.notify_all();
708 5106x }
709
710 inline bool
711 15039x reactor_scheduler::maybe_unlock_and_signal_one(lock_type& lock) const
712 {
713 15039x state_ |= signaled_bit;
714 15039x if (state_ > signaled_bit)
715 {
716 31x lock.unlock();
717 31x cond_.notify_one();
718 31x return true;
719 }
720 15008x return false;
721 }
722
723 inline bool
724 445819x reactor_scheduler::unlock_and_signal_one(lock_type& lock) const
725 {
726 445819x state_ |= signaled_bit;
727 445819x bool have_waiters = state_ > signaled_bit;
728 445819x lock.unlock();
729 445819x if (have_waiters)
730 6x cond_.notify_one();
731 445819x return have_waiters;
732 }
733
734 inline void
735 50x reactor_scheduler::clear_signal() const
736 {
737 50x state_ &= ~signaled_bit;
738 50x }
739
740 inline void
741 6x reactor_scheduler::wait_for_signal(lock_type& lock) const
742 {
743 14x while ((state_ & signaled_bit) == 0)
744 {
745 8x state_ += waiter_increment;
746 8x cond_.wait(lock);
747 8x state_ -= waiter_increment;
748 }
749 6x }
750
751 inline void
752 44x reactor_scheduler::wait_for_signal_for(lock_type& lock, long timeout_us) const
753 {
754 44x if ((state_ & signaled_bit) == 0)
755 {
756 44x state_ += waiter_increment;
757 44x cond_.wait_for(lock, std::chrono::microseconds(timeout_us));
758 44x state_ -= waiter_increment;
759 }
760 44x }
761
762 inline void
763 15039x reactor_scheduler::wake_one_thread_and_unlock(lock_type& lock) const
764 {
765 15039x if (maybe_unlock_and_signal_one(lock))
766 31x return;
767
768 15008x if (task_running_.load(std::memory_order_relaxed) && !task_interrupted_)
769 {
770 891x task_interrupted_ = true;
771 891x lock.unlock();
772 891x interrupt_reactor();
773 }
774 else
775 {
776 14117x lock.unlock();
777 }
778 }
779
780 392894x inline reactor_scheduler::work_cleanup::~work_cleanup()
781 {
782 392894x std::int64_t produced = ctx.private_outstanding_work;
783 392894x if (produced > 1)
784 345x sched->outstanding_work_.fetch_add(
785 produced - 1, std::memory_order_relaxed);
786 392549x else if (produced < 1)
787 36804x sched->work_finished();
788 392894x ctx.private_outstanding_work = 0;
789
790 392894x if (!ctx.private_queue.empty())
791 {
792 95406x lock->lock();
793 95406x sched->completed_ops_.splice(ctx.private_queue);
794 }
795 392894x }
796
797 284131x inline reactor_scheduler::task_cleanup::~task_cleanup()
798 {
799 284131x if (ctx.private_outstanding_work > 0)
800 {
801 7849x sched->outstanding_work_.fetch_add(
802 7849x ctx.private_outstanding_work, std::memory_order_relaxed);
803 7849x ctx.private_outstanding_work = 0;
804 }
805
806 284131x if (!ctx.private_queue.empty())
807 {
808 7849x if (!lock->owns_lock())
809 ✗ lock->lock();
810 7849x sched->completed_ops_.splice(ctx.private_queue);
811 }
812 284131x }
813
814 inline std::size_t
815 396589x reactor_scheduler::do_one(lock_type& lock, long timeout_us, context_type& ctx)
816 {
817 for (;;)
818 {
819 679109x if (stopped_.load(std::memory_order_acquire))
820 2002x return 0;
821
822 677107x std::uintptr_t e = completed_ops_.pop();
823 677107x scheduler_op* op = ready_is_continuation(e) ? nullptr : ready_as_op(e);
824
825 // Handle reactor sentinel — time to poll for I/O
826 677107x if (op == &task_op_)
827 {
828 284162x bool more_handlers = !completed_ops_.empty();
829
830 515335x if (!more_handlers &&
831 462346x (outstanding_work_.load(std::memory_order_acquire) == 0 ||
832 timeout_us == 0))
833 {
834 31x completed_ops_.push(&task_op_);
835 31x return 0;
836 }
837
838 284131x long task_timeout_us = more_handlers ? 0 : timeout_us;
839 284131x task_interrupted_ = task_timeout_us == 0;
840 284131x task_running_.store(true, std::memory_order_release);
841
842 // Wake a peer to take the pending handlers while this thread
843 // polls the reactor; skipped when one_thread_ (no peer exists).
844 284131x if (more_handlers && !one_thread_)
845 52982x unlock_and_signal_one(lock);
846
847 try
848 {
849 284131x run_task(lock, ctx, task_timeout_us);
850 }
851 3x catch (...)
852 {
853 3x task_running_.store(false, std::memory_order_relaxed);
854 3x throw;
855 3x }
856
857 284128x task_running_.store(false, std::memory_order_relaxed);
858 284128x completed_ops_.push(&task_op_);
859 284128x if (timeout_us > 0)
860 1658x return 0;
861 282470x continue;
862 282470x }
863
864 // Handle ready entry (op or continuation)
865 392945x if (e != 0)
866 {
867 392894x bool more = !completed_ops_.empty();
868
869 392894x if (more && !one_thread_)
870 {
871 // Wake a peer for the remaining work; unassisted if none
872 // was parked to take it.
873 392837x ctx.unassisted = !unlock_and_signal_one(lock);
874 }
875 else
876 {
877 // No peer to wake (one_thread_, or nothing more queued).
878 57x ctx.unassisted = more;
879 57x lock.unlock();
880 }
881
882 392894x [[maybe_unused]] work_cleanup on_exit{this, &lock, ctx};
883
884 392894x if (ready_is_continuation(e))
885 24249x ready_as_cont(e)->h.resume();
886 else
887 368645x (*op)();
888 392894x return 1;
889 392894x }
890
891 102x if (outstanding_work_.load(std::memory_order_acquire) == 0 ||
892 timeout_us == 0)
893 1x return 0;
894
895 50x clear_signal();
896 50x if (timeout_us < 0)
897 6x wait_for_signal(lock);
898 else
899 44x wait_for_signal_for(lock, timeout_us);
900 282520x }
901 }
902
903 } // namespace boost::corosio::detail
904
905 #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_SCHEDULER_HPP
906