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

86.2% Lines (313/0/363) 89.4% List of functions (42/0/47)
reactor_scheduler.hpp
f(x) Functions (47)
Function Calls Lines Blocks
boost::corosio::detail::reactor_find_context(boost::corosio::detail::reactor_scheduler const*) :79 1321334x 100.0% 86.0% boost::corosio::detail::reactor_flush_private_work(boost::corosio::detail::reactor_scheduler_context*, std::atomic<long>&) :91 0 0.0% 0.0% boost::corosio::detail::reactor_drain_private_queue(boost::corosio::detail::reactor_scheduler_context*, std::atomic<long>&, boost::corosio::detail::ready_queue&) :108 11x 50.0% 64.0% boost::corosio::detail::reactor_scheduler::inline_budget_initial() const :256 5044x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::scheduler_locking_disabled() const :262 190x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::configure_threading(boost::corosio::detail::scheduler::threading_config) :267 1948x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::reactor_scheduler() :284 1948x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::task_op::operator()() :329 0 0.0% 0.0% boost::corosio::detail::reactor_scheduler::task_op::destroy() :330 0 0.0% 0.0% boost::corosio::detail::reactor_thread_context_guard::reactor_thread_context_guard(boost::corosio::detail::reactor_scheduler const*) :376 5044x 100.0% 100.0% boost::corosio::detail::reactor_thread_context_guard::~reactor_thread_context_guard() :384 5044x 66.7% 80.0% boost::corosio::detail::reactor_scheduler_context::reactor_scheduler_context(boost::corosio::detail::reactor_scheduler const*, boost::corosio::detail::reactor_scheduler_context*) :396 5044x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::configure_reactor(unsigned int, unsigned int, unsigned int, unsigned int) :410 30x 93.3% 74.0% boost::corosio::detail::reactor_scheduler::reset_inline_budget() const :437 124734x 94.4% 93.0% boost::corosio::detail::reactor_scheduler::try_consume_inline_budget() const :468 511379x 87.5% 88.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const :484 3686x 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>) :490 3686x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::~post_handler() :491 7372x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::operator()() :493 3674x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::destroy() :500 12x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(boost::corosio::detail::scheduler_op*) const :525 116299x 100.0% 87.0% boost::corosio::detail::reactor_scheduler::post(boost::capy::continuation&) const :542 25585x 100.0% 87.0% boost::corosio::detail::reactor_scheduler::running_in_this_thread() const :559 15928x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::stop() :565 5088x 100.0% 82.0% boost::corosio::detail::reactor_scheduler::stopped() const :577 101x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::restart() :583 3777x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::run() :589 5044x 100.0% 76.0% boost::corosio::detail::reactor_scheduler::run_one() :614 29x 100.0% 64.0% boost::corosio::detail::reactor_scheduler::wait_one(long) :628 65x 100.0% 70.0% boost::corosio::detail::reactor_scheduler::poll() :642 33x 100.0% 76.0% boost::corosio::detail::reactor_scheduler::poll_one() :667 9x 100.0% 70.0% boost::corosio::detail::reactor_scheduler::work_started() :681 45634x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::work_finished() :687 71914x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::compensating_work_started() const :694 523723x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::drain_thread_queue(boost::corosio::detail::ready_queue&, long) const :702 0 0.0% 0.0% boost::corosio::detail::reactor_scheduler::post_deferred_completions(boost::corosio::detail::ready_queue&) const :715 18327x 30.0% 35.0% boost::corosio::detail::reactor_scheduler::shutdown_drain() :732 1948x 100.0% 88.0% boost::corosio::detail::reactor_scheduler::signal_all(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :760 6948x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::maybe_unlock_and_signal_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :767 19296x 57.1% 50.0% boost::corosio::detail::reactor_scheduler::unlock_and_signal_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :781 759676x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::clear_signal() const :793 11x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wait_for_signal(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :799 11x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wait_for_signal_for(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&, long) const :811 0 0.0% 0.0% boost::corosio::detail::reactor_scheduler::wake_one_thread_and_unlock(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :823 19296x 87.5% 92.0% boost::corosio::detail::reactor_scheduler::work_cleanup::~work_cleanup() :841 687586x 92.3% 92.0% boost::corosio::detail::reactor_scheduler::task_cleanup::~task_cleanup() :865 570770x 83.3% 86.0% boost::corosio::detail::reactor_scheduler::do_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&, long, boost::corosio::detail::reactor_scheduler_context*) :886 692574x 85.4% 78.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,
71 reactor_scheduler_context* n);
72 };
73
74 /// Thread-local context stack for reactor schedulers.
75 inline thread_local_ptr<reactor_scheduler_context> reactor_context_stack;
76
77 /// Find the context frame for a scheduler on this thread.
78 inline reactor_scheduler_context*
79 1321334x reactor_find_context(reactor_scheduler const* self) noexcept
80 {
81 1321334x for (auto* c = reactor_context_stack.get(); c != nullptr; c = c->next)
82 {
83 1286818x if (c->key == self)
84 1286818x return c;
85 }
86 34516x return nullptr;
87 }
88
89 /// Flush private work count to global counter.
90 inline void
91 reactor_flush_private_work(
92 reactor_scheduler_context* ctx,
93 std::atomic<std::int64_t>& outstanding_work) noexcept
94 {
95 if (ctx && ctx->private_outstanding_work > 0)
96 {
97 outstanding_work.fetch_add(
98 ctx->private_outstanding_work, std::memory_order_relaxed);
99 ctx->private_outstanding_work = 0;
100 }
101 }
102
103 /** Drain private queue to global queue, flushing work count first.
104
105 @return True if any ops were drained.
106 */
107 inline bool
108 11x reactor_drain_private_queue(
109 reactor_scheduler_context* ctx,
110 std::atomic<std::int64_t>& outstanding_work,
111 ready_queue& completed_ops) noexcept
112 {
113 11x if (!ctx || ctx->private_queue.empty())
114 11x return false;
115
116 reactor_flush_private_work(ctx, outstanding_work);
117 completed_ops.splice(ctx->private_queue);
118 return true;
119 }
120
121 /** Non-template base for reactor-backed scheduler implementations.
122
123 Provides the complete threading model shared by epoll, kqueue,
124 and select schedulers: signal state machine, inline completion
125 budget, work counting, run/poll methods, and the do_one event
126 loop.
127
128 Derived classes provide platform-specific hooks by overriding:
129 - `run_task(lock, ctx)` to run the reactor poll
130 - `interrupt_reactor()` to wake a blocked reactor
131
132 De-templated from the original CRTP design to eliminate
133 duplicate instantiations when multiple backends are compiled
134 into the same binary. Virtual dispatch for run_task (called
135 once per reactor cycle, before a blocking syscall) has
136 negligible overhead.
137
138 @par Thread Safety
139 All public member functions are thread-safe.
140 */
141 class reactor_scheduler
142 : public scheduler
143 , public capy::execution_context::service
144 {
145 public:
146 using key_type = scheduler;
147 using context_type = reactor_scheduler_context;
148 using mutex_type = conditionally_enabled_mutex;
149 using lock_type = mutex_type::scoped_lock;
150 using event_type = conditionally_enabled_event;
151
152 /// Post a coroutine for deferred execution.
153 void post(std::coroutine_handle<> h) const override;
154
155 /// Post a scheduler operation for deferred execution.
156 void post(scheduler_op* h) const override;
157
158 /// Post a continuation for deferred execution.
159 void post(capy::continuation&) const override;
160
161 /// Return true if called from a thread running this scheduler.
162 bool running_in_this_thread() const noexcept override;
163
164 /// Request the scheduler to stop dispatching handlers.
165 void stop() override;
166
167 /// Return true if the scheduler has been stopped.
168 bool stopped() const noexcept override;
169
170 /// Reset the stopped state so `run()` can resume.
171 void restart() override;
172
173 /// Run the event loop until no work remains.
174 std::size_t run() override;
175
176 /// Run until one handler completes or no work remains.
177 std::size_t run_one() override;
178
179 /// Run until one handler completes or @a usec elapses.
180 std::size_t wait_one(long usec) override;
181
182 /// Run ready handlers without blocking.
183 std::size_t poll() override;
184
185 /// Run at most one ready handler without blocking.
186 std::size_t poll_one() override;
187
188 /// Increment the outstanding work count.
189 void work_started() noexcept override;
190
191 /// Decrement the outstanding work count, stopping on zero.
192 void work_finished() noexcept override;
193
194 /** Reset the thread's inline completion budget.
195
196 Called at the start of each posted completion handler to
197 grant a fresh budget for speculative inline completions.
198 */
199 void reset_inline_budget() const noexcept;
200
201 /** Consume one unit of inline budget if available.
202
203 @return True if budget was available and consumed.
204 */
205 bool try_consume_inline_budget() const noexcept;
206
207 /** Offset a forthcoming work_finished from work_cleanup.
208
209 Called by descriptor_state when all I/O returned EAGAIN and
210 no handler will be executed. Must be called from a scheduler
211 thread.
212 */
213 void compensating_work_started() const noexcept;
214
215 /** Drain work from thread context's private queue to global queue.
216
217 Flushes private work count to the global counter, then
218 transfers the queue under mutex protection.
219
220 @param queue The private queue to drain.
221 @param count Private work count to flush before draining.
222 */
223 void drain_thread_queue(ready_queue& queue, std::int64_t count) const;
224
225 /** Post completed operations for deferred invocation.
226
227 If called from a thread running this scheduler, operations
228 go to the thread's private queue (fast path). Otherwise,
229 operations are added to the global queue under mutex and a
230 waiter is signaled.
231
232 @par Preconditions
233 work_started() must have been called for each operation.
234
235 @param ops Queue of operations to post.
236 */
237 void post_deferred_completions(ready_queue& ops) const;
238
239 /** Apply runtime configuration to the scheduler.
240
241 Called by `io_context` after construction. Values that do
242 not apply to this backend are silently ignored.
243
244 @param max_events Event buffer size for epoll/kqueue.
245 @param budget_init Starting inline completion budget.
246 @param budget_max Hard ceiling on adaptive budget ramp-up.
247 @param unassisted Budget when single-threaded.
248 */
249 virtual void configure_reactor(
250 unsigned max_events,
251 unsigned budget_init,
252 unsigned budget_max,
253 unsigned unassisted);
254
255 /// Return the configured initial inline budget.
256 5044x unsigned inline_budget_initial() const noexcept
257 {
258 5044x return inline_budget_initial_;
259 }
260
261 /// Return true when scheduler locking is disabled (fully-lockless tier).
262 190x bool scheduler_locking_disabled() const noexcept override
263 {
264 190x return scheduler_locking_disabled_;
265 }
266
267 1948x void configure_threading(threading_config cfg) noexcept override
268 {
269 1948x scheduler_locking_disabled_ = !cfg.scheduler_locking;
270 // reactor_io_locking takes effect at descriptor registration (see the
271 // register_descriptor overrides), not here.
272 1948x reactor_io_locking_ = cfg.reactor_io_locking;
273 1948x one_thread_ = cfg.one_thread;
274 1948x mutex_.set_enabled(cfg.scheduler_locking);
275 1948x cond_.set_enabled(cfg.scheduler_locking);
276 1948x }
277
278 protected:
279 timer_service* timer_svc_ = nullptr;
280 bool scheduler_locking_disabled_ = false;
281 bool reactor_io_locking_ = true;
282 bool one_thread_ = false;
283
284 1948x reactor_scheduler() = default;
285
286 /** Drain completed_ops during shutdown.
287
288 Pops all operations from the global queue and destroys them,
289 skipping the task sentinel. Signals all waiting threads.
290 Derived classes call this from their shutdown() override
291 before performing platform-specific cleanup.
292 */
293 void shutdown_drain();
294
295 /// RAII guard that re-inserts the task sentinel after `run_task`.
296 struct task_cleanup
297 {
298 reactor_scheduler const* sched;
299 lock_type* lock;
300 context_type* ctx;
301 ~task_cleanup();
302 };
303
304 mutable mutex_type mutex_{true};
305 mutable event_type cond_{true};
306 mutable ready_queue completed_ops_;
307 mutable std::atomic<std::int64_t> outstanding_work_{0};
308 std::atomic<bool> stopped_{false};
309 mutable std::atomic<bool> task_running_{false};
310 mutable bool task_interrupted_ = false;
311
312 // Runtime-configurable reactor tuning parameters.
313 // Defaults match the library's built-in values.
314 unsigned max_events_per_poll_ = 128;
315 unsigned inline_budget_initial_ = 2;
316 unsigned inline_budget_max_ = 16;
317 unsigned unassisted_budget_ = 4;
318
319 /// Bit 0 of `state_`: set when the condvar should be signaled.
320 static constexpr std::size_t signaled_bit = 1;
321
322 /// Increment per waiting thread in `state_`.
323 static constexpr std::size_t waiter_increment = 2;
324 mutable std::size_t state_ = 0;
325
326 /// Sentinel op that triggers a reactor poll when dequeued.
327 struct task_op final : scheduler_op
328 {
329 void operator()() override {}
330 void destroy() override {}
331 };
332 task_op task_op_;
333
334 /// Run the platform-specific reactor poll.
335 virtual void
336 run_task(lock_type& lock, context_type* ctx,
337 long timeout_us) = 0;
338
339 /// Wake a blocked reactor (e.g. write to eventfd or pipe).
340 virtual void interrupt_reactor() const = 0;
341
342 private:
343 struct work_cleanup
344 {
345 reactor_scheduler* sched;
346 lock_type* lock;
347 context_type* ctx;
348 ~work_cleanup();
349 };
350
351 std::size_t do_one(
352 lock_type& lock, long timeout_us, context_type* ctx);
353
354 void signal_all(lock_type& lock) const;
355 bool maybe_unlock_and_signal_one(lock_type& lock) const;
356 bool unlock_and_signal_one(lock_type& lock) const;
357 void clear_signal() const;
358 void wait_for_signal(lock_type& lock) const;
359 void wait_for_signal_for(
360 lock_type& lock, long timeout_us) const;
361 void wake_one_thread_and_unlock(lock_type& lock) const;
362 };
363
364 /** RAII guard that pushes/pops a scheduler context frame.
365
366 On construction, pushes a new context frame onto the
367 thread-local stack. On destruction, drains any remaining
368 private queue items to the global queue and pops the frame.
369 */
370 struct reactor_thread_context_guard
371 {
372 /// The context frame managed by this guard.
373 reactor_scheduler_context frame_;
374
375 /// Construct the guard, pushing a frame for @a sched.
376 5044x explicit reactor_thread_context_guard(
377 reactor_scheduler const* sched) noexcept
378 5044x : frame_(sched, reactor_context_stack.get())
379 {
380 5044x reactor_context_stack.set(&frame_);
381 5044x }
382
383 /// Destroy the guard, draining private work and popping the frame.
384 5044x ~reactor_thread_context_guard() noexcept
385 {
386 5044x if (!frame_.private_queue.empty())
387 frame_.key->drain_thread_queue(
388 frame_.private_queue, frame_.private_outstanding_work);
389 5044x reactor_context_stack.set(frame_.next);
390 5044x }
391 };
392
393 // ---- Inline implementations ------------------------------------------------
394
395 inline
396 5044x reactor_scheduler_context::reactor_scheduler_context(
397 reactor_scheduler const* k,
398 5044x reactor_scheduler_context* n)
399 5044x : key(k)
400 5044x , next(n)
401 5044x , private_outstanding_work(0)
402 5044x , inline_budget(0)
403 5044x , inline_budget_max(
404 5044x static_cast<int>(k->inline_budget_initial()))
405 5044x , unassisted(false)
406 {
407 5044x }
408
409 inline void
410 30x reactor_scheduler::configure_reactor(
411 unsigned max_events,
412 unsigned budget_init,
413 unsigned budget_max,
414 unsigned unassisted)
415 {
416 58x if (max_events < 1 ||
417 28x max_events > static_cast<unsigned>(std::numeric_limits<int>::max()))
418 throw std::out_of_range(
419 2x "max_events_per_poll must be in [1, INT_MAX]");
420 28x if (budget_max > static_cast<unsigned>(std::numeric_limits<int>::max()))
421 throw std::out_of_range(
422 "inline_budget_max must be in [0, INT_MAX]");
423
424 // Clamp initial and unassisted to budget_max.
425 28x if (budget_init > budget_max)
426 2x budget_init = budget_max;
427 28x if (unassisted > budget_max)
428 2x unassisted = budget_max;
429
430 28x max_events_per_poll_ = max_events;
431 28x inline_budget_initial_ = budget_init;
432 28x inline_budget_max_ = budget_max;
433 28x unassisted_budget_ = unassisted;
434 28x }
435
436 inline void
437 124734x reactor_scheduler::reset_inline_budget() const noexcept
438 {
439 // When budget is disabled (max==0), all paths below would no-op
440 // (inline_budget stays 0). Skip the TLS lookup entirely.
441 124734x if (inline_budget_max_ == 0)
442 return;
443 124734x if (auto* ctx = reactor_find_context(this))
444 {
445 // Cap when no other thread absorbed queued work
446 124734x if (ctx->unassisted)
447 {
448 124732x ctx->inline_budget_max =
449 124732x static_cast<int>(unassisted_budget_);
450 124732x ctx->inline_budget =
451 124732x static_cast<int>(unassisted_budget_);
452 124732x return;
453 }
454 // Ramp up when previous cycle fully consumed budget.
455 // max(1, ...) ensures the doubling escapes zero.
456 2x if (ctx->inline_budget == 0)
457 2x ctx->inline_budget_max = (std::min)(
458 1x (std::max)(1, ctx->inline_budget_max) * 2,
459 2x static_cast<int>(inline_budget_max_));
460 1x else if (ctx->inline_budget < ctx->inline_budget_max)
461 1x ctx->inline_budget_max =
462 1x static_cast<int>(inline_budget_initial_);
463 2x ctx->inline_budget = ctx->inline_budget_max;
464 }
465 }
466
467 inline bool
468 511379x reactor_scheduler::try_consume_inline_budget() const noexcept
469 {
470 511379x if (inline_budget_max_ == 0)
471 return false;
472 511379x if (auto* ctx = reactor_find_context(this))
473 {
474 511379x if (ctx->inline_budget > 0)
475 {
476 406667x --ctx->inline_budget;
477 406667x return true;
478 }
479 }
480 104712x return false;
481 }
482
483 inline void
484 3686x reactor_scheduler::post(std::coroutine_handle<> h) const
485 {
486 struct post_handler final : scheduler_op
487 {
488 std::coroutine_handle<> h_;
489
490 3686x explicit post_handler(std::coroutine_handle<> h) : h_(h) {}
491 7372x ~post_handler() override = default;
492
493 3674x void operator()() override
494 {
495 3674x auto saved = h_;
496 3674x delete this;
497 3674x saved.resume();
498 3674x }
499
500 12x void destroy() override
501 {
502 12x auto saved = h_;
503 12x delete this;
504 12x saved.destroy();
505 12x }
506 };
507
508 3686x auto ph = std::make_unique<post_handler>(h);
509
510 3686x if (auto* ctx = reactor_find_context(this))
511 {
512 26x ++ctx->private_outstanding_work;
513 26x ctx->private_queue.push(ph.release());
514 26x return;
515 }
516
517 3660x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
518
519 3660x lock_type lock(mutex_);
520 3660x completed_ops_.push(ph.release());
521 3660x wake_one_thread_and_unlock(lock);
522 3686x }
523
524 inline void
525 116299x reactor_scheduler::post(scheduler_op* h) const
526 {
527 116299x if (auto* ctx = reactor_find_context(this))
528 {
529 115888x ++ctx->private_outstanding_work;
530 115888x ctx->private_queue.push(h);
531 115888x return;
532 }
533
534 411x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
535
536 411x lock_type lock(mutex_);
537 411x completed_ops_.push(h);
538 411x wake_one_thread_and_unlock(lock);
539 411x }
540
541 inline void
542 25585x reactor_scheduler::post(capy::continuation& c) const
543 {
544 25585x if (auto* ctx = reactor_find_context(this))
545 {
546 10360x ++ctx->private_outstanding_work;
547 10360x ctx->private_queue.push(c);
548 10360x return;
549 }
550
551 15225x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
552
553 15225x lock_type lock(mutex_);
554 15225x completed_ops_.push(c);
555 15225x wake_one_thread_and_unlock(lock);
556 15225x }
557
558 inline bool
559 15928x reactor_scheduler::running_in_this_thread() const noexcept
560 {
561 15928x return reactor_find_context(this) != nullptr;
562 }
563
564 inline void
565 5088x reactor_scheduler::stop()
566 {
567 5088x lock_type lock(mutex_);
568 5088x if (!stopped_.load(std::memory_order_acquire))
569 {
570 5000x stopped_.store(true, std::memory_order_release);
571 5000x signal_all(lock);
572 5000x interrupt_reactor();
573 }
574 5088x }
575
576 inline bool
577 101x reactor_scheduler::stopped() const noexcept
578 {
579 101x return stopped_.load(std::memory_order_acquire);
580 }
581
582 inline void
583 3777x reactor_scheduler::restart()
584 {
585 3777x stopped_.store(false, std::memory_order_release);
586 3777x }
587
588 inline std::size_t
589 5044x reactor_scheduler::run()
590 {
591 10088x if (outstanding_work_.load(std::memory_order_acquire) == 0)
592 {
593 92x stop();
594 92x return 0;
595 }
596
597 4952x reactor_thread_context_guard ctx(this);
598 4952x lock_type lock(mutex_);
599
600 4952x std::size_t n = 0;
601 for (;;)
602 {
603 692446x if (!do_one(lock, -1, &ctx.frame_))
604 4952x break;
605 687494x if (n != (std::numeric_limits<std::size_t>::max)())
606 687494x ++n;
607 687494x if (!lock.owns_lock())
608 567952x lock.lock();
609 }
610 4952x return n;
611 4952x }
612
613 inline std::size_t
614 29x reactor_scheduler::run_one()
615 {
616 58x if (outstanding_work_.load(std::memory_order_acquire) == 0)
617 {
618 1x stop();
619 1x return 0;
620 }
621
622 28x reactor_thread_context_guard ctx(this);
623 28x lock_type lock(mutex_);
624 28x return do_one(lock, -1, &ctx.frame_);
625 28x }
626
627 inline std::size_t
628 65x reactor_scheduler::wait_one(long usec)
629 {
630 130x if (outstanding_work_.load(std::memory_order_acquire) == 0)
631 {
632 23x stop();
633 23x return 0;
634 }
635
636 42x reactor_thread_context_guard ctx(this);
637 42x lock_type lock(mutex_);
638 42x return do_one(lock, usec, &ctx.frame_);
639 42x }
640
641 inline std::size_t
642 33x reactor_scheduler::poll()
643 {
644 66x if (outstanding_work_.load(std::memory_order_acquire) == 0)
645 {
646 15x stop();
647 15x return 0;
648 }
649
650 18x reactor_thread_context_guard ctx(this);
651 18x lock_type lock(mutex_);
652
653 18x std::size_t n = 0;
654 for (;;)
655 {
656 54x if (!do_one(lock, 0, &ctx.frame_))
657 18x break;
658 36x if (n != (std::numeric_limits<std::size_t>::max)())
659 36x ++n;
660 36x if (!lock.owns_lock())
661 36x lock.lock();
662 }
663 18x return n;
664 18x }
665
666 inline std::size_t
667 9x reactor_scheduler::poll_one()
668 {
669 18x if (outstanding_work_.load(std::memory_order_acquire) == 0)
670 {
671 5x stop();
672 5x return 0;
673 }
674
675 4x reactor_thread_context_guard ctx(this);
676 4x lock_type lock(mutex_);
677 4x return do_one(lock, 0, &ctx.frame_);
678 4x }
679
680 inline void
681 45634x reactor_scheduler::work_started() noexcept
682 {
683 45634x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
684 45634x }
685
686 inline void
687 71914x reactor_scheduler::work_finished() noexcept
688 {
689 143828x if (outstanding_work_.fetch_sub(1, std::memory_order_acq_rel) == 1)
690 4943x stop();
691 71914x }
692
693 inline void
694 523723x reactor_scheduler::compensating_work_started() const noexcept
695 {
696 523723x auto* ctx = reactor_find_context(this);
697 523723x if (ctx)
698 523723x ++ctx->private_outstanding_work;
699 523723x }
700
701 inline void
702 reactor_scheduler::drain_thread_queue(
703 ready_queue& queue, std::int64_t count) const
704 {
705 if (count > 0)
706 outstanding_work_.fetch_add(count, std::memory_order_relaxed);
707
708 lock_type lock(mutex_);
709 completed_ops_.splice(queue);
710 if (count > 0)
711 maybe_unlock_and_signal_one(lock);
712 }
713
714 inline void
715 18327x reactor_scheduler::post_deferred_completions(ready_queue& ops) const
716 {
717 18327x if (ops.empty())
718 18327x return;
719
720 if (auto* ctx = reactor_find_context(this))
721 {
722 ctx->private_queue.splice(ops);
723 return;
724 }
725
726 lock_type lock(mutex_);
727 completed_ops_.splice(ops);
728 wake_one_thread_and_unlock(lock);
729 }
730
731 inline void
732 1948x reactor_scheduler::shutdown_drain()
733 {
734 1948x lock_type lock(mutex_);
735
736 4246x while (auto e = completed_ops_.pop())
737 {
738 2298x if (ready_is_continuation(e))
739 {
740 4x lock.unlock();
741 4x if (auto h = ready_as_cont(e)->h)
742 4x h.destroy();
743 4x lock.lock();
744 }
745 else
746 {
747 2294x auto* op = ready_as_op(e);
748 2294x if (op == &task_op_)
749 1948x continue;
750 346x lock.unlock();
751 346x op->destroy();
752 346x lock.lock();
753 }
754 2298x }
755
756 1948x signal_all(lock);
757 1948x }
758
759 inline void
760 6948x reactor_scheduler::signal_all(lock_type&) const
761 {
762 6948x state_ |= signaled_bit;
763 6948x cond_.notify_all();
764 6948x }
765
766 inline bool
767 19296x reactor_scheduler::maybe_unlock_and_signal_one(
768 lock_type& lock) const
769 {
770 19296x state_ |= signaled_bit;
771 19296x if (state_ > signaled_bit)
772 {
773 lock.unlock();
774 cond_.notify_one();
775 return true;
776 }
777 19296x return false;
778 }
779
780 inline bool
781 759676x reactor_scheduler::unlock_and_signal_one(
782 lock_type& lock) const
783 {
784 759676x state_ |= signaled_bit;
785 759676x bool have_waiters = state_ > signaled_bit;
786 759676x lock.unlock();
787 759676x if (have_waiters)
788 7x cond_.notify_one();
789 759676x return have_waiters;
790 }
791
792 inline void
793 11x reactor_scheduler::clear_signal() const
794 {
795 11x state_ &= ~signaled_bit;
796 11x }
797
798 inline void
799 11x reactor_scheduler::wait_for_signal(
800 lock_type& lock) const
801 {
802 22x while ((state_ & signaled_bit) == 0)
803 {
804 11x state_ += waiter_increment;
805 11x cond_.wait(lock);
806 11x state_ -= waiter_increment;
807 }
808 11x }
809
810 inline void
811 reactor_scheduler::wait_for_signal_for(
812 lock_type& lock, long timeout_us) const
813 {
814 if ((state_ & signaled_bit) == 0)
815 {
816 state_ += waiter_increment;
817 cond_.wait_for(lock, std::chrono::microseconds(timeout_us));
818 state_ -= waiter_increment;
819 }
820 }
821
822 inline void
823 19296x reactor_scheduler::wake_one_thread_and_unlock(
824 lock_type& lock) const
825 {
826 19296x if (maybe_unlock_and_signal_one(lock))
827 return;
828
829 19296x if (task_running_.load(std::memory_order_relaxed) && !task_interrupted_)
830 {
831 126x task_interrupted_ = true;
832 126x lock.unlock();
833 126x interrupt_reactor();
834 }
835 else
836 {
837 19170x lock.unlock();
838 }
839 }
840
841 687586x inline reactor_scheduler::work_cleanup::~work_cleanup()
842 {
843 687586x if (ctx)
844 {
845 687586x std::int64_t produced = ctx->private_outstanding_work;
846 687586x if (produced > 1)
847 326x sched->outstanding_work_.fetch_add(
848 produced - 1, std::memory_order_relaxed);
849 687260x else if (produced < 1)
850 44623x sched->work_finished();
851 687586x ctx->private_outstanding_work = 0;
852
853 687586x if (!ctx->private_queue.empty())
854 {
855 119542x lock->lock();
856 119542x sched->completed_ops_.splice(ctx->private_queue);
857 }
858 }
859 else
860 {
861 sched->work_finished();
862 }
863 687586x }
864
865 1141540x inline reactor_scheduler::task_cleanup::~task_cleanup()
866 {
867 570770x if (!ctx)
868 return;
869
870 570770x if (ctx->private_outstanding_work > 0)
871 {
872 6694x sched->outstanding_work_.fetch_add(
873 6694x ctx->private_outstanding_work, std::memory_order_relaxed);
874 6694x ctx->private_outstanding_work = 0;
875 }
876
877 570770x if (!ctx->private_queue.empty())
878 {
879 6694x if (!lock->owns_lock())
880 lock->lock();
881 6694x sched->completed_ops_.splice(ctx->private_queue);
882 }
883 570770x }
884
885 inline std::size_t
886 692574x reactor_scheduler::do_one(
887 lock_type& lock, long timeout_us, context_type* ctx)
888 {
889 for (;;)
890 {
891 1263338x if (stopped_.load(std::memory_order_acquire))
892 4954x return 0;
893
894 1258384x std::uintptr_t e = completed_ops_.pop();
895 1258384x scheduler_op* op = ready_is_continuation(e) ? nullptr : ready_as_op(e);
896
897 // Handle reactor sentinel — time to poll for I/O
898 1258384x if (op == &task_op_)
899 {
900 bool more_handlers =
901 570787x !completed_ops_.empty() || (ctx && !ctx->private_queue.empty());
902
903 1069458x if (!more_handlers &&
904 997342x (outstanding_work_.load(std::memory_order_acquire) == 0 ||
905 timeout_us == 0))
906 {
907 17x completed_ops_.push(&task_op_);
908 17x return 0;
909 }
910
911 570770x long task_timeout_us = more_handlers ? 0 : timeout_us;
912 570770x task_interrupted_ = task_timeout_us == 0;
913 570770x task_running_.store(true, std::memory_order_release);
914
915 // Wake a peer to take the pending handlers while this thread
916 // polls the reactor; skipped when one_thread_ (no peer exists).
917 570770x if (more_handlers && !one_thread_)
918 72109x unlock_and_signal_one(lock);
919
920 try
921 {
922 570770x run_task(lock, ctx, task_timeout_us);
923 }
924 catch (...)
925 {
926 task_running_.store(false, std::memory_order_relaxed);
927 throw;
928 }
929
930 570770x task_running_.store(false, std::memory_order_relaxed);
931 570770x completed_ops_.push(&task_op_);
932 570770x if (timeout_us > 0)
933 17x return 0;
934 570753x continue;
935 570753x }
936
937 // Handle ready entry (op or continuation)
938 687597x if (e != 0)
939 {
940 687586x bool more = !completed_ops_.empty();
941
942 687586x if (more && !one_thread_)
943 {
944 // Wake a peer for the remaining work; unassisted if none
945 // was parked to take it.
946 687567x ctx->unassisted = !unlock_and_signal_one(lock);
947 }
948 else
949 {
950 // No peer to wake (one_thread_, or nothing more queued).
951 19x ctx->unassisted = more;
952 19x lock.unlock();
953 }
954
955 687586x [[maybe_unused]] work_cleanup on_exit{this, &lock, ctx};
956
957 687586x if (ready_is_continuation(e))
958 25581x ready_as_cont(e)->h.resume();
959 else
960 662005x (*op)();
961 687586x return 1;
962 687586x }
963
964 // Try private queue before blocking
965 11x if (reactor_drain_private_queue(ctx, outstanding_work_, completed_ops_))
966 continue;
967
968 22x if (outstanding_work_.load(std::memory_order_acquire) == 0 ||
969 timeout_us == 0)
970 return 0;
971
972 11x clear_signal();
973 11x if (timeout_us < 0)
974 11x wait_for_signal(lock);
975 else
976 wait_for_signal_for(lock, timeout_us);
977 570764x }
978 }
979
980 } // namespace boost::corosio::detail
981
982 #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_SCHEDULER_HPP
983