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

85.1% Lines (309/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 992581x 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 6x 50.0% 64.0% boost::corosio::detail::reactor_scheduler::inline_budget_initial() const :256 1496x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::scheduler_locking_disabled() const :262 240x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::configure_threading(boost::corosio::detail::scheduler::threading_config) :267 1790x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::reactor_scheduler() :284 1802x 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*) :383 1496x 100.0% 100.0% boost::corosio::detail::reactor_thread_context_guard::~reactor_thread_context_guard() :391 1496x 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*) :403 1496x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::configure_reactor(unsigned int, unsigned int, unsigned int, unsigned int) :417 28x 93.3% 74.0% boost::corosio::detail::reactor_scheduler::reset_inline_budget() const :444 93273x 50.0% 43.0% boost::corosio::detail::reactor_scheduler::try_consume_inline_budget() const :475 408267x 87.5% 88.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const :491 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>) :497 3686x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::~post_handler() :498 7372x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::operator()() :500 3674x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(std::__n4861::coroutine_handle<void>) const::post_handler::destroy() :507 12x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::post(boost::corosio::detail::scheduler_op*) const :532 90086x 100.0% 87.0% boost::corosio::detail::reactor_scheduler::post(boost::capy::continuation&) const :549 15417x 100.0% 87.0% boost::corosio::detail::reactor_scheduler::running_in_this_thread() const :566 7955x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::stop() :572 1477x 100.0% 82.0% boost::corosio::detail::reactor_scheduler::stopped() const :584 96x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::restart() :590 359x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::run() :596 1422x 100.0% 86.0% boost::corosio::detail::reactor_scheduler::run_one() :621 110x 100.0% 64.0% boost::corosio::detail::reactor_scheduler::wait_one(long) :635 64x 100.0% 70.0% boost::corosio::detail::reactor_scheduler::poll() :649 33x 100.0% 76.0% boost::corosio::detail::reactor_scheduler::poll_one() :674 9x 100.0% 70.0% boost::corosio::detail::reactor_scheduler::work_started() :688 27117x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::work_finished() :694 44305x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::compensating_work_started() const :701 373897x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::drain_thread_queue(boost::corosio::detail::ready_queue&, long) const :709 0 0.0% 0.0% boost::corosio::detail::reactor_scheduler::post_deferred_completions(boost::corosio::detail::ready_queue&) const :722 10996x 30.0% 35.0% boost::corosio::detail::reactor_scheduler::shutdown_drain() :739 1790x 100.0% 88.0% boost::corosio::detail::reactor_scheduler::signal_all(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :767 3171x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::maybe_unlock_and_signal_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :774 11500x 57.1% 50.0% boost::corosio::detail::reactor_scheduler::unlock_and_signal_one(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :788 542866x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::clear_signal() const :800 6x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wait_for_signal(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :806 6x 100.0% 100.0% boost::corosio::detail::reactor_scheduler::wait_for_signal_for(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&, long) const :818 0 0.0% 0.0% boost::corosio::detail::reactor_scheduler::wake_one_thread_and_unlock(boost::corosio::detail::conditionally_enabled_mutex::scoped_lock&) const :830 11500x 87.5% 92.0% boost::corosio::detail::reactor_scheduler::work_cleanup::~work_cleanup() :848 493998x 92.3% 92.0% boost::corosio::detail::reactor_scheduler::task_cleanup::~task_cleanup() :872 398899x 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*) :893 495357x 93.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,
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 992581x reactor_find_context(reactor_scheduler const* self) noexcept
80 {
81 992581x for (auto* c = reactor_context_stack.get(); c != nullptr; c = c->next)
82 {
83 973811x if (c->key == self)
84 973811x return c;
85 }
86 18770x 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 6x 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 6x if (!ctx || ctx->private_queue.empty())
114 6x 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 1496x unsigned inline_budget_initial() const noexcept
257 {
258 1496x return inline_budget_initial_;
259 }
260
261 /// Return true when scheduler locking is disabled (fully-lockless tier).
262 240x bool scheduler_locking_disabled() const noexcept override
263 {
264 240x return scheduler_locking_disabled_;
265 }
266
267 1790x void configure_threading(threading_config cfg) noexcept override
268 {
269 1790x scheduler_locking_disabled_ = !cfg.scheduler_locking;
270 // reactor_io_locking takes effect at descriptor registration (see the
271 // register_descriptor overrides), not here.
272 1790x reactor_io_locking_ = cfg.reactor_io_locking;
273 1790x one_thread_ = cfg.one_thread;
274 1790x mutex_.set_enabled(cfg.scheduler_locking);
275 1790x cond_.set_enabled(cfg.scheduler_locking);
276 1790x }
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 1802x 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
336 @par Postconditions
337 `lock` is owned on return, however the poll ended. An
338 implementation that unlocks around the blocking call owes the
339 caller a matching re-acquire on every path out, including the
340 errors it retries rather than reports.
341 */
342 virtual void
343 run_task(lock_type& lock, context_type* ctx,
344 long timeout_us) = 0;
345
346 /// Wake a blocked reactor (e.g. write to eventfd or pipe).
347 virtual void interrupt_reactor() const = 0;
348
349 private:
350 struct work_cleanup
351 {
352 reactor_scheduler* sched;
353 lock_type* lock;
354 context_type* ctx;
355 ~work_cleanup();
356 };
357
358 std::size_t do_one(
359 lock_type& lock, long timeout_us, context_type* ctx);
360
361 void signal_all(lock_type& lock) const;
362 bool maybe_unlock_and_signal_one(lock_type& lock) const;
363 bool unlock_and_signal_one(lock_type& lock) const;
364 void clear_signal() const;
365 void wait_for_signal(lock_type& lock) const;
366 void wait_for_signal_for(
367 lock_type& lock, long timeout_us) const;
368 void wake_one_thread_and_unlock(lock_type& lock) const;
369 };
370
371 /** RAII guard that pushes/pops a scheduler context frame.
372
373 On construction, pushes a new context frame onto the
374 thread-local stack. On destruction, drains any remaining
375 private queue items to the global queue and pops the frame.
376 */
377 struct reactor_thread_context_guard
378 {
379 /// The context frame managed by this guard.
380 reactor_scheduler_context frame_;
381
382 /// Construct the guard, pushing a frame for @a sched.
383 1496x explicit reactor_thread_context_guard(
384 reactor_scheduler const* sched) noexcept
385 1496x : frame_(sched, reactor_context_stack.get())
386 {
387 1496x reactor_context_stack.set(&frame_);
388 1496x }
389
390 /// Destroy the guard, draining private work and popping the frame.
391 1496x ~reactor_thread_context_guard() noexcept
392 {
393 1496x if (!frame_.private_queue.empty())
394 frame_.key->drain_thread_queue(
395 frame_.private_queue, frame_.private_outstanding_work);
396 1496x reactor_context_stack.set(frame_.next);
397 1496x }
398 };
399
400 // ---- Inline implementations ------------------------------------------------
401
402 inline
403 1496x reactor_scheduler_context::reactor_scheduler_context(
404 reactor_scheduler const* k,
405 1496x reactor_scheduler_context* n)
406 1496x : key(k)
407 1496x , next(n)
408 1496x , private_outstanding_work(0)
409 1496x , inline_budget(0)
410 1496x , inline_budget_max(
411 1496x static_cast<int>(k->inline_budget_initial()))
412 1496x , unassisted(false)
413 {
414 1496x }
415
416 inline void
417 28x reactor_scheduler::configure_reactor(
418 unsigned max_events,
419 unsigned budget_init,
420 unsigned budget_max,
421 unsigned unassisted)
422 {
423 54x if (max_events < 1 ||
424 26x max_events > static_cast<unsigned>(std::numeric_limits<int>::max()))
425 throw std::out_of_range(
426 2x "max_events_per_poll must be in [1, INT_MAX]");
427 26x if (budget_max > static_cast<unsigned>(std::numeric_limits<int>::max()))
428 throw std::out_of_range(
429 "inline_budget_max must be in [0, INT_MAX]");
430
431 // Clamp initial and unassisted to budget_max.
432 26x if (budget_init > budget_max)
433 2x budget_init = budget_max;
434 26x if (unassisted > budget_max)
435 2x unassisted = budget_max;
436
437 26x max_events_per_poll_ = max_events;
438 26x inline_budget_initial_ = budget_init;
439 26x inline_budget_max_ = budget_max;
440 26x unassisted_budget_ = unassisted;
441 26x }
442
443 inline void
444 93273x reactor_scheduler::reset_inline_budget() const noexcept
445 {
446 // When budget is disabled (max==0), all paths below would no-op
447 // (inline_budget stays 0). Skip the TLS lookup entirely.
448 93273x if (inline_budget_max_ == 0)
449 return;
450 93273x if (auto* ctx = reactor_find_context(this))
451 {
452 // Cap when no other thread absorbed queued work
453 93273x if (ctx->unassisted)
454 {
455 93273x ctx->inline_budget_max =
456 93273x static_cast<int>(unassisted_budget_);
457 93273x ctx->inline_budget =
458 93273x static_cast<int>(unassisted_budget_);
459 93273x return;
460 }
461 // Ramp up when previous cycle fully consumed budget.
462 // max(1, ...) ensures the doubling escapes zero.
463 if (ctx->inline_budget == 0)
464 ctx->inline_budget_max = (std::min)(
465 (std::max)(1, ctx->inline_budget_max) * 2,
466 static_cast<int>(inline_budget_max_));
467 else if (ctx->inline_budget < ctx->inline_budget_max)
468 ctx->inline_budget_max =
469 static_cast<int>(inline_budget_initial_);
470 ctx->inline_budget = ctx->inline_budget_max;
471 }
472 }
473
474 inline bool
475 408267x reactor_scheduler::try_consume_inline_budget() const noexcept
476 {
477 408267x if (inline_budget_max_ == 0)
478 return false;
479 408267x if (auto* ctx = reactor_find_context(this))
480 {
481 408267x if (ctx->inline_budget > 0)
482 {
483 326456x --ctx->inline_budget;
484 326456x return true;
485 }
486 }
487 81811x return false;
488 }
489
490 inline void
491 3686x reactor_scheduler::post(std::coroutine_handle<> h) const
492 {
493 struct post_handler final : scheduler_op
494 {
495 std::coroutine_handle<> h_;
496
497 3686x explicit post_handler(std::coroutine_handle<> h) : h_(h) {}
498 7372x ~post_handler() override = default;
499
500 3674x void operator()() override
501 {
502 3674x auto saved = h_;
503 3674x delete this;
504 3674x saved.resume();
505 3674x }
506
507 12x void destroy() override
508 {
509 12x auto saved = h_;
510 12x delete this;
511 12x saved.destroy();
512 12x }
513 };
514
515 3686x auto ph = std::make_unique<post_handler>(h);
516
517 3686x if (auto* ctx = reactor_find_context(this))
518 {
519 26x ++ctx->private_outstanding_work;
520 26x ctx->private_queue.push(ph.release());
521 26x return;
522 }
523
524 3660x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
525
526 3660x lock_type lock(mutex_);
527 3660x completed_ops_.push(ph.release());
528 3660x wake_one_thread_and_unlock(lock);
529 3686x }
530
531 inline void
532 90086x reactor_scheduler::post(scheduler_op* h) const
533 {
534 90086x if (auto* ctx = reactor_find_context(this))
535 {
536 89519x ++ctx->private_outstanding_work;
537 89519x ctx->private_queue.push(h);
538 89519x return;
539 }
540
541 567x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
542
543 567x lock_type lock(mutex_);
544 567x completed_ops_.push(h);
545 567x wake_one_thread_and_unlock(lock);
546 567x }
547
548 inline void
549 15417x reactor_scheduler::post(capy::continuation& c) const
550 {
551 15417x if (auto* ctx = reactor_find_context(this))
552 {
553 8144x ++ctx->private_outstanding_work;
554 8144x ctx->private_queue.push(c);
555 8144x return;
556 }
557
558 7273x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
559
560 7273x lock_type lock(mutex_);
561 7273x completed_ops_.push(c);
562 7273x wake_one_thread_and_unlock(lock);
563 7273x }
564
565 inline bool
566 7955x reactor_scheduler::running_in_this_thread() const noexcept
567 {
568 7955x return reactor_find_context(this) != nullptr;
569 }
570
571 inline void
572 1477x reactor_scheduler::stop()
573 {
574 1477x lock_type lock(mutex_);
575 1477x if (!stopped_.load(std::memory_order_acquire))
576 {
577 1381x stopped_.store(true, std::memory_order_release);
578 1381x signal_all(lock);
579 1381x interrupt_reactor();
580 }
581 1477x }
582
583 inline bool
584 96x reactor_scheduler::stopped() const noexcept
585 {
586 96x return stopped_.load(std::memory_order_acquire);
587 }
588
589 inline void
590 359x reactor_scheduler::restart()
591 {
592 359x stopped_.store(false, std::memory_order_release);
593 359x }
594
595 inline std::size_t
596 1422x reactor_scheduler::run()
597 {
598 2844x if (outstanding_work_.load(std::memory_order_acquire) == 0)
599 {
600 99x stop();
601 99x return 0;
602 }
603
604 1323x reactor_thread_context_guard ctx(this);
605 1323x lock_type lock(mutex_);
606
607 1323x std::size_t n = 0;
608 for (;;)
609 {
610 495148x if (!do_one(lock, -1, &ctx.frame_))
611 1320x break;
612 493825x if (n != (std::numeric_limits<std::size_t>::max)())
613 493825x ++n;
614 493825x if (!lock.owns_lock())
615 401692x lock.lock();
616 }
617 1320x return n;
618 1326x }
619
620 inline std::size_t
621 110x reactor_scheduler::run_one()
622 {
623 220x if (outstanding_work_.load(std::memory_order_acquire) == 0)
624 {
625 1x stop();
626 1x return 0;
627 }
628
629 109x reactor_thread_context_guard ctx(this);
630 109x lock_type lock(mutex_);
631 109x return do_one(lock, -1, &ctx.frame_);
632 109x }
633
634 inline std::size_t
635 64x reactor_scheduler::wait_one(long usec)
636 {
637 128x if (outstanding_work_.load(std::memory_order_acquire) == 0)
638 {
639 22x stop();
640 22x return 0;
641 }
642
643 42x reactor_thread_context_guard ctx(this);
644 42x lock_type lock(mutex_);
645 42x return do_one(lock, usec, &ctx.frame_);
646 42x }
647
648 inline std::size_t
649 33x reactor_scheduler::poll()
650 {
651 66x if (outstanding_work_.load(std::memory_order_acquire) == 0)
652 {
653 15x stop();
654 15x return 0;
655 }
656
657 18x reactor_thread_context_guard ctx(this);
658 18x lock_type lock(mutex_);
659
660 18x std::size_t n = 0;
661 for (;;)
662 {
663 54x if (!do_one(lock, 0, &ctx.frame_))
664 18x break;
665 36x if (n != (std::numeric_limits<std::size_t>::max)())
666 36x ++n;
667 36x if (!lock.owns_lock())
668 36x lock.lock();
669 }
670 18x return n;
671 18x }
672
673 inline std::size_t
674 9x reactor_scheduler::poll_one()
675 {
676 18x if (outstanding_work_.load(std::memory_order_acquire) == 0)
677 {
678 5x stop();
679 5x return 0;
680 }
681
682 4x reactor_thread_context_guard ctx(this);
683 4x lock_type lock(mutex_);
684 4x return do_one(lock, 0, &ctx.frame_);
685 4x }
686
687 inline void
688 27117x reactor_scheduler::work_started() noexcept
689 {
690 27117x outstanding_work_.fetch_add(1, std::memory_order_relaxed);
691 27117x }
692
693 inline void
694 44305x reactor_scheduler::work_finished() noexcept
695 {
696 88610x if (outstanding_work_.fetch_sub(1, std::memory_order_acq_rel) == 1)
697 1320x stop();
698 44305x }
699
700 inline void
701 373897x reactor_scheduler::compensating_work_started() const noexcept
702 {
703 373897x auto* ctx = reactor_find_context(this);
704 373897x if (ctx)
705 373897x ++ctx->private_outstanding_work;
706 373897x }
707
708 inline void
709 reactor_scheduler::drain_thread_queue(
710 ready_queue& queue, std::int64_t count) const
711 {
712 if (count > 0)
713 outstanding_work_.fetch_add(count, std::memory_order_relaxed);
714
715 lock_type lock(mutex_);
716 completed_ops_.splice(queue);
717 if (count > 0)
718 maybe_unlock_and_signal_one(lock);
719 }
720
721 inline void
722 10996x reactor_scheduler::post_deferred_completions(ready_queue& ops) const
723 {
724 10996x if (ops.empty())
725 10996x return;
726
727 if (auto* ctx = reactor_find_context(this))
728 {
729 ctx->private_queue.splice(ops);
730 return;
731 }
732
733 lock_type lock(mutex_);
734 completed_ops_.splice(ops);
735 wake_one_thread_and_unlock(lock);
736 }
737
738 inline void
739 1790x reactor_scheduler::shutdown_drain()
740 {
741 1790x lock_type lock(mutex_);
742
743 3934x while (auto e = completed_ops_.pop())
744 {
745 2144x if (ready_is_continuation(e))
746 {
747 8x lock.unlock();
748 8x if (auto h = ready_as_cont(e)->h)
749 8x h.destroy();
750 8x lock.lock();
751 }
752 else
753 {
754 2136x auto* op = ready_as_op(e);
755 2136x if (op == &task_op_)
756 1787x continue;
757 349x lock.unlock();
758 349x op->destroy();
759 349x lock.lock();
760 }
761 2144x }
762
763 1790x signal_all(lock);
764 1790x }
765
766 inline void
767 3171x reactor_scheduler::signal_all(lock_type&) const
768 {
769 3171x state_ |= signaled_bit;
770 3171x cond_.notify_all();
771 3171x }
772
773 inline bool
774 11500x reactor_scheduler::maybe_unlock_and_signal_one(
775 lock_type& lock) const
776 {
777 11500x state_ |= signaled_bit;
778 11500x if (state_ > signaled_bit)
779 {
780 lock.unlock();
781 cond_.notify_one();
782 return true;
783 }
784 11500x return false;
785 }
786
787 inline bool
788 542866x reactor_scheduler::unlock_and_signal_one(
789 lock_type& lock) const
790 {
791 542866x state_ |= signaled_bit;
792 542866x bool have_waiters = state_ > signaled_bit;
793 542866x lock.unlock();
794 542866x if (have_waiters)
795 6x cond_.notify_one();
796 542866x return have_waiters;
797 }
798
799 inline void
800 6x reactor_scheduler::clear_signal() const
801 {
802 6x state_ &= ~signaled_bit;
803 6x }
804
805 inline void
806 6x reactor_scheduler::wait_for_signal(
807 lock_type& lock) const
808 {
809 14x while ((state_ & signaled_bit) == 0)
810 {
811 8x state_ += waiter_increment;
812 8x cond_.wait(lock);
813 8x state_ -= waiter_increment;
814 }
815 6x }
816
817 inline void
818 reactor_scheduler::wait_for_signal_for(
819 lock_type& lock, long timeout_us) const
820 {
821 if ((state_ & signaled_bit) == 0)
822 {
823 state_ += waiter_increment;
824 cond_.wait_for(lock, std::chrono::microseconds(timeout_us));
825 state_ -= waiter_increment;
826 }
827 }
828
829 inline void
830 11500x reactor_scheduler::wake_one_thread_and_unlock(
831 lock_type& lock) const
832 {
833 11500x if (maybe_unlock_and_signal_one(lock))
834 return;
835
836 11500x if (task_running_.load(std::memory_order_relaxed) && !task_interrupted_)
837 {
838 165x task_interrupted_ = true;
839 165x lock.unlock();
840 165x interrupt_reactor();
841 }
842 else
843 {
844 11335x lock.unlock();
845 }
846 }
847
848 493998x inline reactor_scheduler::work_cleanup::~work_cleanup()
849 {
850 493998x if (ctx)
851 {
852 493998x std::int64_t produced = ctx->private_outstanding_work;
853 493998x if (produced > 1)
854 322x sched->outstanding_work_.fetch_add(
855 produced - 1, std::memory_order_relaxed);
856 493676x else if (produced < 1)
857 28252x sched->work_finished();
858 493998x ctx->private_outstanding_work = 0;
859
860 493998x if (!ctx->private_queue.empty())
861 {
862 92155x lock->lock();
863 92155x sched->completed_ops_.splice(ctx->private_queue);
864 }
865 }
866 else
867 {
868 sched->work_finished();
869 }
870 493998x }
871
872 797798x inline reactor_scheduler::task_cleanup::~task_cleanup()
873 {
874 398899x if (!ctx)
875 return;
876
877 398899x if (ctx->private_outstanding_work > 0)
878 {
879 5505x sched->outstanding_work_.fetch_add(
880 5505x ctx->private_outstanding_work, std::memory_order_relaxed);
881 5505x ctx->private_outstanding_work = 0;
882 }
883
884 398899x if (!ctx->private_queue.empty())
885 {
886 5505x if (!lock->owns_lock())
887 lock->lock();
888 5505x sched->completed_ops_.splice(ctx->private_queue);
889 }
890 398899x }
891
892 inline std::size_t
893 495357x reactor_scheduler::do_one(
894 lock_type& lock, long timeout_us, context_type* ctx)
895 {
896 for (;;)
897 {
898 894241x if (stopped_.load(std::memory_order_acquire))
899 1322x return 0;
900
901 892919x std::uintptr_t e = completed_ops_.pop();
902 892919x scheduler_op* op = ready_is_continuation(e) ? nullptr : ready_as_op(e);
903
904 // Handle reactor sentinel — time to poll for I/O
905 892919x if (op == &task_op_)
906 {
907 bool more_handlers =
908 398915x !completed_ops_.empty() || (ctx && !ctx->private_queue.empty());
909
910 748934x if (!more_handlers &&
911 700038x (outstanding_work_.load(std::memory_order_acquire) == 0 ||
912 timeout_us == 0))
913 {
914 16x completed_ops_.push(&task_op_);
915 16x return 0;
916 }
917
918 398899x long task_timeout_us = more_handlers ? 0 : timeout_us;
919 398899x task_interrupted_ = task_timeout_us == 0;
920 398899x task_running_.store(true, std::memory_order_release);
921
922 // Wake a peer to take the pending handlers while this thread
923 // polls the reactor; skipped when one_thread_ (no peer exists).
924 398899x if (more_handlers && !one_thread_)
925 48889x unlock_and_signal_one(lock);
926
927 try
928 {
929 398899x run_task(lock, ctx, task_timeout_us);
930 }
931 3x catch (...)
932 {
933 3x task_running_.store(false, std::memory_order_relaxed);
934 3x throw;
935 3x }
936
937 398896x task_running_.store(false, std::memory_order_relaxed);
938 398896x completed_ops_.push(&task_op_);
939 398896x if (timeout_us > 0)
940 18x return 0;
941 398878x continue;
942 398878x }
943
944 // Handle ready entry (op or continuation)
945 494004x if (e != 0)
946 {
947 493998x bool more = !completed_ops_.empty();
948
949 493998x if (more && !one_thread_)
950 {
951 // Wake a peer for the remaining work; unassisted if none
952 // was parked to take it.
953 493977x ctx->unassisted = !unlock_and_signal_one(lock);
954 }
955 else
956 {
957 // No peer to wake (one_thread_, or nothing more queued).
958 21x ctx->unassisted = more;
959 21x lock.unlock();
960 }
961
962 493998x [[maybe_unused]] work_cleanup on_exit{this, &lock, ctx};
963
964 493998x if (ready_is_continuation(e))
965 15409x ready_as_cont(e)->h.resume();
966 else
967 478589x (*op)();
968 493998x return 1;
969 493998x }
970
971 // Try private queue before blocking
972 6x if (reactor_drain_private_queue(ctx, outstanding_work_, completed_ops_))
973 continue;
974
975 12x if (outstanding_work_.load(std::memory_order_acquire) == 0 ||
976 timeout_us == 0)
977 return 0;
978
979 6x clear_signal();
980 6x if (timeout_us < 0)
981 6x wait_for_signal(lock);
982 else
983 wait_for_signal_for(lock, timeout_us);
984 398884x }
985 }
986
987 } // namespace boost::corosio::detail
988
989 #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_SCHEDULER_HPP
990