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 966774x 100.0% 86.0% boost::corosio::detail::reactor_scheduler::inline_budget_initial() const :213 5449x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::scheduler_locking_disabled() const :219 502x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::configure_threading(boost::corosio::detail::scheduler::threading_config) :224 2241x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::reactor_scheduler() :241 2253x 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 5449x 100.0% 100.0% boost::corosio::detail::reactor_thread_context_guard::~reactor_thread_context_guard() :353 5449x 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 5449x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::configure_reactor(unsigned int, unsigned int, unsigned int, unsigned int) :373 44x 100.0% 87.0% boost::corosio::detail::reactor_scheduler::reset_inline_budget() const :398 95878x 53.3% 50.0% boost::corosio::detail::reactor_scheduler::try_consume_inline_budget() const :426 427135x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const :442 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>) :448 3756x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::~post_handler() :449 7512x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::operator()() :451 3744x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::destroy() :458 12x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(boost::corosio::detail::scheduler_op*) const :483 105543x 100.0% 87.0% boost::corosio::detail::reactor_scheduler::post(boost::capy::continuation&) const :500 26056x 100.0% 87.0% boost::corosio::detail::reactor_scheduler::running_in_this_thread() const :517 10763x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::stop() :523 3678x 100.0% 82.0% boost::corosio::detail::reactor_scheduler::stopped() const :535 2510x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::restart() :541 1439x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::run() :547 2066x 100.0% 86.0% boost::corosio::detail::reactor_scheduler::run_one() :572 112x 100.0% 70.0% boost::corosio::detail::reactor_scheduler::wait_one(long) :586 4108x 100.0% 70.0% boost::corosio::detail::reactor_scheduler::poll() :600 49x 100.0% 76.0% boost::corosio::detail::reactor_scheduler::poll_one() :625 11x 100.0% 70.0% boost::corosio::detail::reactor_scheduler::work_started() :639 36609x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::work_finished() :645 68581x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::compensating_work_started() const :652 297737x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post_deferred_completions(boost::corosio::detail::ready_queue&) const :660 9681x 60.0% 59.0% boost::corosio::detail::reactor_scheduler::shutdown_drain() :677 2241x 100.0% 88.0% boost::corosio::detail::reactor_scheduler::signal_all(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :705 5070x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::maybe_unlock_and_signal_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :712 15108x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::unlock_and_signal_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :725 494921x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::clear_signal() const :736 42x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wait_for_signal(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :742 6x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wait_for_signal_for(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&, long) const :753 36x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wake_one_thread_and_unlock(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :764 15108x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::work_cleanup::~work_cleanup() :781 442691x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::task_cleanup::~task_cleanup() :798 325555x 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&) :816 446373x 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 966774x reactor_find_context(reactor_scheduler const* self) noexcept
79 {
80 966774x for (auto* c = reactor_context_stack.get(); c != nullptr; c = c->next)
81 {
82 941848x if (c->key == self)
83 941848x return c;
84 }
85 24926x 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 @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 5449x unsigned inline_budget_initial() const noexcept
214 {
215 5449x return inline_budget_initial_;
216 }
217
218 /// Return true when scheduler locking is disabled (fully-lockless tier).
219 502x bool scheduler_locking_disabled() const noexcept override
220 {
221 502x return scheduler_locking_disabled_;
222 }
223
224 2241x void configure_threading(threading_config cfg) noexcept override
225 {
226 2241x scheduler_locking_disabled_ = !cfg.scheduler_locking;
227 // reactor_io_locking takes effect at descriptor registration (see the
228 // register_descriptor overrides), not here.
229 2241x reactor_io_locking_ = cfg.reactor_io_locking;
230 2241x one_thread_ = cfg.one_thread;
231 2241x mutex_.set_enabled(cfg.scheduler_locking);
232 2241x cond_.set_enabled(cfg.scheduler_locking);
233 2241x }
234
235 protected:
236 timer_service* timer_svc_ = nullptr;
237 bool scheduler_locking_disabled_ = false;
238 bool reactor_io_locking_ = true;
239 bool one_thread_ = false;
240
241 2253x 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 mutable mutex_type mutex_{true};
262 mutable event_type cond_{true};
263 mutable ready_queue completed_ops_;
264 mutable std::atomic<std::int64_t> outstanding_work_{0};
265 std::atomic<bool> stopped_{false};
266 mutable std::atomic<bool> task_running_{false};
267 mutable bool task_interrupted_ = false;
268
269 // Runtime-configurable reactor tuning parameters.
270 // Defaults match the library's built-in values.
271 unsigned max_events_per_poll_ = 128;
272 unsigned inline_budget_initial_ = 2;
273 unsigned inline_budget_max_ = 16;
274 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 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 5449x explicit reactor_thread_context_guard(
341 reactor_scheduler const* sched) noexcept
342 5449x : frame_(sched, reactor_context_stack.get())
343 {
344 5449x reactor_context_stack.set(&frame_);
345 5449x }
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 5449x ~reactor_thread_context_guard() noexcept
354 {
355 5449x reactor_context_stack.set(frame_.next);
356 5449x }
357 };
358
359 // ---- Inline implementations ------------------------------------------------
360
361 5449x inline reactor_scheduler_context::reactor_scheduler_context(
362 5449x reactor_scheduler const* k, reactor_scheduler_context* n)
363 5449x : key(k)
364 5449x , next(n)
365 5449x , private_outstanding_work(0)
366 5449x , inline_budget(0)
367 5449x , inline_budget_max(static_cast<int>(k->inline_budget_initial()))
368 5449x , unassisted(false)
369 {
370 5449x }
371
372 inline void
373 44x reactor_scheduler::configure_reactor(
374 unsigned max_events,
375 unsigned budget_init,
376 unsigned budget_max,
377 unsigned unassisted)
378 {
379 86x if (max_events < 1 ||
380 42x max_events > static_cast<unsigned>(std::numeric_limits<int>::max()))
381 2x throw std::out_of_range("max_events_per_poll must be in [1, INT_MAX]");
382 42x if (budget_max > static_cast<unsigned>(std::numeric_limits<int>::max()))
383 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 40x if (budget_init > budget_max)
387 16x budget_init = budget_max;
388 40x if (unassisted > budget_max)
389 16x unassisted = budget_max;
390
391 40x max_events_per_poll_ = max_events;
392 40x inline_budget_initial_ = budget_init;
393 40x inline_budget_max_ = budget_max;
394 40x unassisted_budget_ = unassisted;
395 40x }
396
397 inline void
398 95878x 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 95878x if (inline_budget_max_ == 0)
403 56x return;
404 95822x if (auto* ctx = reactor_find_context(this))
405 {
406 // Cap when no other thread absorbed queued work
407 95822x if (ctx->unassisted)
408 {
409 95822x ctx->inline_budget_max = static_cast<int>(unassisted_budget_);
410 95822x ctx->inline_budget = static_cast<int>(unassisted_budget_);
411 95822x return;
412 }
413 // Ramp up when previous cycle fully consumed budget.
414 // max(1, ...) ensures the doubling escapes zero.
415 ✗ if (ctx->inline_budget == 0)
416 ✗ ctx->inline_budget_max =
417 ✗ (std::min)((std::max)(1, ctx->inline_budget_max) * 2,
418 ✗ static_cast<int>(inline_budget_max_));
419 ✗ else if (ctx->inline_budget < ctx->inline_budget_max)
420 ✗ ctx->inline_budget_max = static_cast<int>(inline_budget_initial_);
421 ✗ ctx->inline_budget = ctx->inline_budget_max;
422 }
423 }
424
425 inline bool
426 427135x reactor_scheduler::try_consume_inline_budget() const noexcept
427 {
428 427135x if (inline_budget_max_ == 0)
429 40x return false;
430 427095x if (auto* ctx = reactor_find_context(this))
431 {
432 427095x if (ctx->inline_budget > 0)
433 {
434 341511x --ctx->inline_budget;
435 341511x return true;
436 }
437 }
438 85584x return false;
439 }
440
441 inline void
442 3756x reactor_scheduler::post(std::coroutine_handle<> h) const
443 {
444 struct post_handler final : scheduler_op
445 {
446 std::coroutine_handle<> h_;
447
448 3756x explicit post_handler(std::coroutine_handle<> h) : h_(h) {}
449 7512x ~post_handler() override = default;
450
451 3744x void operator()() override
452 {
453 3744x auto saved = h_;
454 3744x delete this;
455 3744x saved.resume();
456 3744x }
457
458 12x void destroy() override
459 {
460 12x auto saved = h_;
461 12x delete this;
462 12x saved.destroy();
463 12x }
464 };
465
466 3756x auto ph = std::make_unique<post_handler>(h);
467
468 3756x if (auto* ctx = reactor_find_context(this))
469 {
470 96x ++ctx->private_outstanding_work;
471 96x ctx->private_queue.push(ph.release());
472 96x return;
473 }
474
475 3660x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
476
477 3660x lock_type lock(mutex_);
478 3660x completed_ops_.push(ph.release());
479 3660x wake_one_thread_and_unlock(lock);
480 3756x }
481
482 inline void
483 105543x reactor_scheduler::post(scheduler_op* h) const
484 {
485 105543x if (auto* ctx = reactor_find_context(this))
486 {
487 104518x ++ctx->private_outstanding_work;
488 104518x ctx->private_queue.push(h);
489 104518x return;
490 }
491
492 1025x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
493
494 1025x lock_type lock(mutex_);
495 1025x completed_ops_.push(h);
496 1025x wake_one_thread_and_unlock(lock);
497 1025x }
498
499 inline void
500 26056x reactor_scheduler::post(capy::continuation& c) const
501 {
502 26056x if (auto* ctx = reactor_find_context(this))
503 {
504 15633x ++ctx->private_outstanding_work;
505 15633x ctx->private_queue.push(c);
506 15633x return;
507 }
508
509 10423x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
510
511 10423x lock_type lock(mutex_);
512 10423x completed_ops_.push(c);
513 10423x wake_one_thread_and_unlock(lock);
514 10423x }
515
516 inline bool
517 10763x reactor_scheduler::running_in_this_thread() const noexcept
518 {
519 10763x return reactor_find_context(this) != nullptr;
520 }
521
522 inline void
523 3678x reactor_scheduler::stop()
524 {
525 3678x lock_type lock(mutex_);
526 3678x if (!stopped_.load(std::memory_order_acquire))
527 {
528 2829x stopped_.store(true, std::memory_order_release);
529 2829x signal_all(lock);
530 2829x interrupt_reactor();
531 }
532 3678x }
533
534 inline bool
535 2510x reactor_scheduler::stopped() const noexcept
536 {
537 2510x return stopped_.load(std::memory_order_acquire);
538 }
539
540 inline void
541 1439x reactor_scheduler::restart()
542 {
543 1439x stopped_.store(false, std::memory_order_release);
544 1439x }
545
546 inline std::size_t
547 2066x reactor_scheduler::run()
548 {
549 4132x if (outstanding_work_.load(std::memory_order_acquire) == 0)
550 {
551 110x stop();
552 110x return 0;
553 }
554
555 1956x reactor_thread_context_guard ctx(this);
556 1956x lock_type lock(mutex_);
557
558 1956x std::size_t n = 0;
559 for (;;)
560 {
561 442839x if (!do_one(lock, -1, ctx.frame_))
562 1953x break;
563 440883x if (n != (std::numeric_limits<std::size_t>::max)())
564 440883x ++n;
565 440883x if (!lock.owns_lock())
566 337527x lock.lock();
567 }
568 1953x return n;
569 1959x }
570
571 inline std::size_t
572 112x reactor_scheduler::run_one()
573 {
574 224x 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 109x lock_type lock(mutex_);
582 109x return do_one(lock, -1, ctx.frame_);
583 109x }
584
585 inline std::size_t
586 4108x reactor_scheduler::wait_one(long usec)
587 {
588 8216x if (outstanding_work_.load(std::memory_order_acquire) == 0)
589 {
590 764x stop();
591 764x return 0;
592 }
593
594 3344x reactor_thread_context_guard ctx(this);
595 3344x lock_type lock(mutex_);
596 3344x return do_one(lock, usec, ctx.frame_);
597 3344x }
598
599 inline std::size_t
600 49x reactor_scheduler::poll()
601 {
602 98x 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 34x lock_type lock(mutex_);
610
611 34x std::size_t n = 0;
612 for (;;)
613 {
614 75x if (!do_one(lock, 0, ctx.frame_))
615 34x break;
616 41x if (n != (std::numeric_limits<std::size_t>::max)())
617 41x ++n;
618 41x if (!lock.owns_lock())
619 41x lock.lock();
620 }
621 34x return n;
622 34x }
623
624 inline std::size_t
625 11x reactor_scheduler::poll_one()
626 {
627 22x 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 6x lock_type lock(mutex_);
635 6x return do_one(lock, 0, ctx.frame_);
636 6x }
637
638 inline void
639 36609x reactor_scheduler::work_started() noexcept
640 {
641 36609x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
642 36609x }
643
644 inline void
645 68581x reactor_scheduler::work_finished() noexcept
646 {
647 137162x if (outstanding_work_.fetch_sub(1, std::memory_order_acq_rel) == 1)
648 2766x stop();
649 68581x }
650
651 inline void
652 297737x reactor_scheduler::compensating_work_started() const noexcept
653 {
654 297737x auto* ctx = reactor_find_context(this);
655 297737x if (ctx)
656 297737x ++ctx->private_outstanding_work;
657 297737x }
658
659 inline void
660 9681x reactor_scheduler::post_deferred_completions(ready_queue& ops) const
661 {
662 9681x if (ops.empty())
663 9681x return;
664
665 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 ✗ }
675
676 inline void
677 2241x reactor_scheduler::shutdown_drain()
678 {
679 2241x lock_type lock(mutex_);
680
681 4861x while (auto e = completed_ops_.pop())
682 {
683 2620x if (ready_is_continuation(e))
684 {
685 8x lock.unlock();
686 8x if (auto h = ready_as_cont(e)->h)
687 8x h.destroy();
688 8x lock.lock();
689 }
690 else
691 {
692 2612x auto* op = ready_as_op(e);
693 2612x if (op == &task_op_)
694 2238x continue;
695 374x lock.unlock();
696 374x op->destroy();
697 374x lock.lock();
698 }
699 2620x }
700
701 2241x signal_all(lock);
702 2241x }
703
704 inline void
705 5070x reactor_scheduler::signal_all(lock_type&) const
706 {
707 5070x state_ |= signaled_bit;
708 5070x cond_.notify_all();
709 5070x }
710
711 inline bool
712 15108x reactor_scheduler::maybe_unlock_and_signal_one(lock_type& lock) const
713 {
714 15108x state_ |= signaled_bit;
715 15108x if (state_ > signaled_bit)
716 {
717 38x lock.unlock();
718 38x cond_.notify_one();
719 38x return true;
720 }
721 15070x return false;
722 }
723
724 inline bool
725 494921x reactor_scheduler::unlock_and_signal_one(lock_type& lock) const
726 {
727 494921x state_ |= signaled_bit;
728 494921x bool have_waiters = state_ > signaled_bit;
729 494921x lock.unlock();
730 494921x if (have_waiters)
731 6x cond_.notify_one();
732 494921x return have_waiters;
733 }
734
735 inline void
736 42x reactor_scheduler::clear_signal() const
737 {
738 42x state_ &= ~signaled_bit;
739 42x }
740
741 inline void
742 6x reactor_scheduler::wait_for_signal(lock_type& lock) const
743 {
744 14x while ((state_ & signaled_bit) == 0)
745 {
746 8x state_ += waiter_increment;
747 8x cond_.wait(lock);
748 8x state_ -= waiter_increment;
749 }
750 6x }
751
752 inline void
753 36x reactor_scheduler::wait_for_signal_for(lock_type& lock, long timeout_us) const
754 {
755 36x if ((state_ & signaled_bit) == 0)
756 {
757 36x state_ += waiter_increment;
758 36x cond_.wait_for(lock, std::chrono::microseconds(timeout_us));
759 36x state_ -= waiter_increment;
760 }
761 36x }
762
763 inline void
764 15108x reactor_scheduler::wake_one_thread_and_unlock(lock_type& lock) const
765 {
766 15108x if (maybe_unlock_and_signal_one(lock))
767 38x return;
768
769 15070x if (task_running_.load(std::memory_order_relaxed) && !task_interrupted_)
770 {
771 850x task_interrupted_ = true;
772 850x lock.unlock();
773 850x interrupt_reactor();
774 }
775 else
776 {
777 14220x lock.unlock();
778 }
779 }
780
781 442691x inline reactor_scheduler::work_cleanup::~work_cleanup()
782 {
783 442691x std::int64_t produced = ctx.private_outstanding_work;
784 442691x if (produced > 1)
785 344x sched->outstanding_work_.fetch_add(
786 produced - 1, std::memory_order_relaxed);
787 442347x else if (produced < 1)
788 41723x sched->work_finished();
789 442691x ctx.private_outstanding_work = 0;
790
791 442691x if (!ctx.private_queue.empty())
792 {
793 103550x lock->lock();
794 103550x sched->completed_ops_.splice(ctx.private_queue);
795 }
796 442691x }
797
798 325555x inline reactor_scheduler::task_cleanup::~task_cleanup()
799 {
800 325555x if (ctx.private_outstanding_work > 0)
801 {
802 9690x sched->outstanding_work_.fetch_add(
803 9690x ctx.private_outstanding_work, std::memory_order_relaxed);
804 9690x ctx.private_outstanding_work = 0;
805 }
806
807 325555x if (!ctx.private_queue.empty())
808 {
809 9690x if (!lock->owns_lock())
810 ✗ lock->lock();
811 9690x sched->completed_ops_.splice(ctx.private_queue);
812 }
813 325555x }
814
815 inline std::size_t
816 446373x reactor_scheduler::do_one(lock_type& lock, long timeout_us, context_type& ctx)
817 {
818 for (;;)
819 {
820 770277x if (stopped_.load(std::memory_order_acquire))
821 1956x return 0;
822
823 768321x std::uintptr_t e = completed_ops_.pop();
824 768321x scheduler_op* op = ready_is_continuation(e) ? nullptr : ready_as_op(e);
825
826 // Handle reactor sentinel — time to poll for I/O
827 768321x if (op == &task_op_)
828 {
829 325588x bool more_handlers = !completed_ops_.empty();
830
831 598894x if (!more_handlers &&
832 546612x (outstanding_work_.load(std::memory_order_acquire) == 0 ||
833 timeout_us == 0))
834 {
835 33x completed_ops_.push(&task_op_);
836 33x return 0;
837 }
838
839 325555x long task_timeout_us = more_handlers ? 0 : timeout_us;
840 325555x task_interrupted_ = task_timeout_us == 0;
841 325555x 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 325555x if (more_handlers && !one_thread_)
846 52275x unlock_and_signal_one(lock);
847
848 try
849 {
850 325555x run_task(lock, ctx, task_timeout_us);
851 }
852 3x catch (...)
853 {
854 3x task_running_.store(false, std::memory_order_relaxed);
855 3x throw;
856 3x }
857
858 325552x task_running_.store(false, std::memory_order_relaxed);
859 325552x completed_ops_.push(&task_op_);
860 325552x if (timeout_us > 0)
861 1690x return 0;
862 323862x continue;
863 323862x }
864
865 // Handle ready entry (op or continuation)
866 442733x if (e != 0)
867 {
868 442691x bool more = !completed_ops_.empty();
869
870 442691x if (more && !one_thread_)
871 {
872 // Wake a peer for the remaining work; unassisted if none
873 // was parked to take it.
874 442646x ctx.unassisted = !unlock_and_signal_one(lock);
875 }
876 else
877 {
878 // No peer to wake (one_thread_, or nothing more queued).
879 45x ctx.unassisted = more;
880 45x lock.unlock();
881 }
882
883 442691x [[maybe_unused]] work_cleanup on_exit{this, &lock, ctx};
884
885 442691x if (ready_is_continuation(e))
886 26048x ready_as_cont(e)->h.resume();
887 else
888 416643x (*op)();
889 442691x return 1;
890 442691x }
891
892 84x if (outstanding_work_.load(std::memory_order_acquire) == 0 ||
893 timeout_us == 0)
894 ✗ return 0;
895
896 42x clear_signal();
897 42x if (timeout_us < 0)
898 6x wait_for_signal(lock);
899 else
900 36x wait_for_signal_for(lock, timeout_us);
901 323904x }
902 }
903
904 } // namespace boost::corosio::detail
905
906 #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_SCHEDULER_HPP
907