TLA Line data 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 HIT 986551 : reactor_find_context(reactor_scheduler const* self) noexcept
79 : {
80 986551 : for (auto* c = reactor_context_stack.get(); c != nullptr; c = c->next)
81 : {
82 961413 : if (c->key == self)
83 961413 : return c;
84 : }
85 25138 : return nullptr;
86 : }
87 :
88 : /** Non-template base for reactor-backed scheduler implementations.
89 :
90 : Provides the complete threading model shared by epoll, kqueue,
91 : and select schedulers: signal state machine, inline completion
92 : budget, work counting, run/poll methods, and the do_one event
93 : loop.
94 :
95 : Derived classes provide platform-specific hooks by overriding:
96 : - `run_task(lock, ctx)` to run the reactor poll
97 : - `interrupt_reactor()` to wake a blocked reactor
98 :
99 : De-templated from the original CRTP design to eliminate
100 : duplicate instantiations when multiple backends are compiled
101 : into the same binary. Virtual dispatch for run_task (called
102 : once per reactor cycle, before a blocking syscall) has
103 : negligible overhead.
104 :
105 : @par Thread Safety
106 : All public member functions are thread-safe.
107 : */
108 : class reactor_scheduler
109 : : public scheduler
110 : {
111 : public:
112 : using context_type = reactor_scheduler_context;
113 : using mutex_type = conditionally_enabled_mutex;
114 : using lock_type = mutex_type::scoped_lock;
115 : using event_type = conditionally_enabled_event;
116 :
117 : /// Post a coroutine for deferred execution.
118 : void post(std::coroutine_handle<> h) const override;
119 :
120 : /// Post a scheduler operation for deferred execution.
121 : void post(scheduler_op* h) const override;
122 :
123 : /// Post a continuation for deferred execution.
124 : void post(capy::continuation&) const override;
125 :
126 : /// Return true if called from a thread running this scheduler.
127 : bool running_in_this_thread() const noexcept override;
128 :
129 : /// Request the scheduler to stop dispatching handlers.
130 : void stop() override;
131 :
132 : /// Return true if the scheduler has been stopped.
133 : bool stopped() const noexcept override;
134 :
135 : /// Reset the stopped state so `run()` can resume.
136 : void restart() override;
137 :
138 : /// Run the event loop until no work remains.
139 : std::size_t run() override;
140 :
141 : /// Run until one handler completes or no work remains.
142 : std::size_t run_one() override;
143 :
144 : /// Run until one handler completes or @a usec elapses.
145 : std::size_t wait_one(long usec) override;
146 :
147 : /// Run ready handlers without blocking.
148 : std::size_t poll() override;
149 :
150 : /// Run at most one ready handler without blocking.
151 : std::size_t poll_one() override;
152 :
153 : /// Increment the outstanding work count.
154 : void work_started() noexcept override;
155 :
156 : /// Decrement the outstanding work count, stopping on zero.
157 : void work_finished() noexcept override;
158 :
159 : /** Reset the thread's inline completion budget.
160 :
161 : Called at the start of each posted completion handler to
162 : grant a fresh budget for speculative inline completions.
163 : */
164 : void reset_inline_budget() const noexcept;
165 :
166 : /** Consume one unit of inline budget if available.
167 :
168 : @return True if budget was available and consumed.
169 : */
170 : bool try_consume_inline_budget() const noexcept;
171 :
172 : /** Offset a forthcoming work_finished from work_cleanup.
173 :
174 : Called by descriptor_state when all I/O returned EAGAIN and
175 : no handler will be executed. Must be called from a scheduler
176 : thread.
177 : */
178 : void compensating_work_started() const noexcept;
179 :
180 : /** Post completed operations for deferred invocation.
181 :
182 : If called from a thread running this scheduler, operations
183 : go to the thread's private queue (fast path). Otherwise,
184 : operations are added to the global queue under mutex and a
185 : waiter is signaled.
186 :
187 : @pre work_started() must have been called for each operation.
188 :
189 : @param ops Queue of operations to post.
190 : */
191 : void post_deferred_completions(ready_queue& ops) const;
192 :
193 : /** Apply runtime configuration to the scheduler.
194 :
195 : Called by `io_context` after construction. Values that do
196 : not apply to this backend are silently ignored.
197 :
198 : @param max_events Event buffer size for epoll/kqueue.
199 : @param budget_init Starting inline completion budget.
200 : @param budget_max Hard ceiling on adaptive budget ramp-up.
201 : @param unassisted Budget when single-threaded.
202 : */
203 : virtual void configure_reactor(
204 : unsigned max_events,
205 : unsigned budget_init,
206 : unsigned budget_max,
207 : unsigned unassisted);
208 :
209 : /// Return the configured initial inline budget.
210 5453 : unsigned inline_budget_initial() const noexcept
211 : {
212 5453 : return inline_budget_initial_;
213 : }
214 :
215 : /// Return true when scheduler locking is disabled (fully-lockless tier).
216 514 : bool scheduler_locking_disabled() const noexcept override
217 : {
218 514 : return scheduler_locking_disabled_;
219 : }
220 :
221 2310 : void configure_threading(threading_config cfg) noexcept override
222 : {
223 2310 : scheduler_locking_disabled_ = !cfg.scheduler_locking;
224 : // reactor_io_locking takes effect at descriptor registration (see the
225 : // register_descriptor overrides), not here.
226 2310 : reactor_io_locking_ = cfg.reactor_io_locking;
227 2310 : one_thread_ = cfg.one_thread;
228 2310 : mutex_.set_enabled(cfg.scheduler_locking);
229 2310 : cond_.set_enabled(cfg.scheduler_locking);
230 2310 : }
231 :
232 : protected:
233 : timer_service* timer_svc_ = nullptr;
234 : bool scheduler_locking_disabled_ = false;
235 : bool reactor_io_locking_ = true;
236 : bool one_thread_ = false;
237 :
238 2322 : reactor_scheduler() = default;
239 :
240 : /** Drain completed_ops during shutdown.
241 :
242 : Pops all operations from the global queue and destroys them,
243 : skipping the task sentinel. Signals all waiting threads.
244 : Derived classes call this from their shutdown() override
245 : before performing platform-specific cleanup.
246 : */
247 : void shutdown_drain();
248 :
249 : /// RAII guard that re-inserts the task sentinel after `run_task`.
250 : struct task_cleanup
251 : {
252 : reactor_scheduler const* sched;
253 : lock_type* lock;
254 : context_type& ctx;
255 : ~task_cleanup();
256 : };
257 :
258 : mutable mutex_type mutex_{true};
259 : mutable event_type cond_{true};
260 : mutable ready_queue completed_ops_;
261 : mutable std::atomic<std::int64_t> outstanding_work_{0};
262 : std::atomic<bool> stopped_{false};
263 : mutable std::atomic<bool> task_running_{false};
264 : mutable bool task_interrupted_ = false;
265 :
266 : // Runtime-configurable reactor tuning parameters.
267 : // Defaults match the library's built-in values.
268 : unsigned max_events_per_poll_ = 128;
269 : unsigned inline_budget_initial_ = 2;
270 : unsigned inline_budget_max_ = 16;
271 : unsigned unassisted_budget_ = 4;
272 :
273 : /// Bit 0 of `state_`: set when the condvar should be signaled.
274 : static constexpr std::size_t signaled_bit = 1;
275 :
276 : /// Increment per waiting thread in `state_`.
277 : static constexpr std::size_t waiter_increment = 2;
278 : mutable std::size_t state_ = 0;
279 :
280 : /// Sentinel op that triggers a reactor poll when dequeued.
281 : struct task_op final : scheduler_op
282 : {
283 : // LCOV_EXCL_START: the sentinel is intercepted by pointer
284 : // identity; its virtuals exist for vtable completeness.
285 : void operator()() override {}
286 : void destroy() override {}
287 : // LCOV_EXCL_STOP
288 : };
289 : task_op task_op_;
290 :
291 : /** Run the platform-specific reactor poll.
292 :
293 : @par Postconditions
294 : `lock` is owned on return, however the poll ended. An
295 : implementation that unlocks around the blocking call owes the
296 : caller a matching re-acquire on every path out, including the
297 : errors it retries rather than reports.
298 : */
299 : virtual void
300 : run_task(lock_type& lock, context_type& ctx, long timeout_us) = 0;
301 :
302 : /// Wake a blocked reactor (e.g. write to eventfd or pipe).
303 : virtual void interrupt_reactor() const = 0;
304 :
305 : private:
306 : struct work_cleanup
307 : {
308 : reactor_scheduler* sched;
309 : lock_type* lock;
310 : context_type& ctx;
311 : ~work_cleanup();
312 : };
313 :
314 : std::size_t do_one(lock_type& lock, long timeout_us, context_type& ctx);
315 :
316 : void signal_all(lock_type& lock) const;
317 : bool maybe_unlock_and_signal_one(lock_type& lock) const;
318 : bool unlock_and_signal_one(lock_type& lock) const;
319 : void clear_signal() const;
320 : void wait_for_signal(lock_type& lock) const;
321 : void wait_for_signal_for(lock_type& lock, long timeout_us) const;
322 : void wake_one_thread_and_unlock(lock_type& lock) const;
323 : };
324 :
325 : /** RAII guard that pushes/pops a scheduler context frame.
326 :
327 : On construction, pushes a new context frame onto the
328 : thread-local stack. On destruction, drains any remaining
329 : private queue items to the global queue and pops the frame.
330 : */
331 : struct reactor_thread_context_guard
332 : {
333 : /// The context frame managed by this guard.
334 : reactor_scheduler_context frame_;
335 :
336 : /// Construct the guard, pushing a frame for @a sched.
337 5453 : explicit reactor_thread_context_guard(
338 : reactor_scheduler const* sched) noexcept
339 5453 : : frame_(sched, reactor_context_stack.get())
340 : {
341 5453 : reactor_context_stack.set(&frame_);
342 5453 : }
343 :
344 : /** Destroy the guard, popping the frame.
345 :
346 : The private queue is empty here by invariant: work_cleanup and
347 : task_cleanup splice it to the global queue after every handler
348 : and every reactor pass.
349 : */
350 5453 : ~reactor_thread_context_guard() noexcept
351 : {
352 5453 : reactor_context_stack.set(frame_.next);
353 5453 : }
354 : };
355 :
356 : // ---- Inline implementations ------------------------------------------------
357 :
358 5453 : inline reactor_scheduler_context::reactor_scheduler_context(
359 5453 : reactor_scheduler const* k, reactor_scheduler_context* n)
360 5453 : : key(k)
361 5453 : , next(n)
362 5453 : , private_outstanding_work(0)
363 5453 : , inline_budget(0)
364 5453 : , inline_budget_max(static_cast<int>(k->inline_budget_initial()))
365 5453 : , unassisted(false)
366 : {
367 5453 : }
368 :
369 : inline void
370 46 : reactor_scheduler::configure_reactor(
371 : unsigned max_events,
372 : unsigned budget_init,
373 : unsigned budget_max,
374 : unsigned unassisted)
375 : {
376 90 : if (max_events < 1 ||
377 44 : max_events > static_cast<unsigned>(std::numeric_limits<int>::max()))
378 2 : throw std::out_of_range("max_events_per_poll must be in [1, INT_MAX]");
379 44 : if (budget_max > static_cast<unsigned>(std::numeric_limits<int>::max()))
380 2 : throw std::out_of_range("inline_budget_max must be in [0, INT_MAX]");
381 :
382 : // Clamp initial and unassisted to budget_max.
383 42 : if (budget_init > budget_max)
384 18 : budget_init = budget_max;
385 42 : if (unassisted > budget_max)
386 18 : unassisted = budget_max;
387 :
388 42 : max_events_per_poll_ = max_events;
389 42 : inline_budget_initial_ = budget_init;
390 42 : inline_budget_max_ = budget_max;
391 42 : unassisted_budget_ = unassisted;
392 42 : }
393 :
394 : inline void
395 92473 : reactor_scheduler::reset_inline_budget() const noexcept
396 : {
397 : // When budget is disabled (max==0), all paths below would no-op
398 : // (inline_budget stays 0). Skip the TLS lookup entirely.
399 92473 : if (inline_budget_max_ == 0)
400 58 : return;
401 92415 : if (auto* ctx = reactor_find_context(this))
402 : {
403 : // Cap when no other thread absorbed queued work
404 92415 : if (ctx->unassisted)
405 : {
406 92415 : ctx->inline_budget_max = static_cast<int>(unassisted_budget_);
407 92415 : ctx->inline_budget = static_cast<int>(unassisted_budget_);
408 92415 : return;
409 : }
410 : // Ramp up when previous cycle fully consumed budget.
411 : // max(1, ...) ensures the doubling escapes zero.
412 MIS 0 : if (ctx->inline_budget == 0)
413 0 : ctx->inline_budget_max =
414 0 : (std::min)((std::max)(1, ctx->inline_budget_max) * 2,
415 0 : static_cast<int>(inline_budget_max_));
416 0 : else if (ctx->inline_budget < ctx->inline_budget_max)
417 0 : ctx->inline_budget_max = static_cast<int>(inline_budget_initial_);
418 0 : ctx->inline_budget = ctx->inline_budget_max;
419 : }
420 : }
421 :
422 : inline bool
423 HIT 409870 : reactor_scheduler::try_consume_inline_budget() const noexcept
424 : {
425 409870 : if (inline_budget_max_ == 0)
426 42 : return false;
427 409828 : if (auto* ctx = reactor_find_context(this))
428 : {
429 409828 : if (ctx->inline_budget > 0)
430 : {
431 327680 : --ctx->inline_budget;
432 327680 : return true;
433 : }
434 : }
435 82148 : return false;
436 : }
437 :
438 : inline void
439 3756 : reactor_scheduler::post(std::coroutine_handle<> h) const
440 : {
441 : struct post_handler final : scheduler_op
442 : {
443 : std::coroutine_handle<> h_;
444 :
445 3756 : explicit post_handler(std::coroutine_handle<> h) : h_(h) {}
446 7512 : ~post_handler() override = default;
447 :
448 3744 : void operator()() override
449 : {
450 3744 : auto saved = h_;
451 3744 : delete this;
452 3744 : saved.resume();
453 3744 : }
454 :
455 12 : void destroy() override
456 : {
457 12 : auto saved = h_;
458 12 : delete this;
459 12 : saved.destroy();
460 12 : }
461 : };
462 :
463 3756 : auto ph = std::make_unique<post_handler>(h);
464 :
465 3756 : if (auto* ctx = reactor_find_context(this))
466 : {
467 96 : ++ctx->private_outstanding_work;
468 96 : ctx->private_queue.push(ph.release());
469 96 : return;
470 : }
471 :
472 3660 : outstanding_work_.fetch_add(1, std::memory_order_relaxed);
473 :
474 3660 : lock_type lock(mutex_);
475 3660 : completed_ops_.push(ph.release());
476 3660 : wake_one_thread_and_unlock(lock);
477 3756 : }
478 :
479 : inline void
480 102216 : reactor_scheduler::post(scheduler_op* h) const
481 : {
482 102216 : if (auto* ctx = reactor_find_context(this))
483 : {
484 100949 : ++ctx->private_outstanding_work;
485 100949 : ctx->private_queue.push(h);
486 100949 : return;
487 : }
488 :
489 1267 : outstanding_work_.fetch_add(1, std::memory_order_relaxed);
490 :
491 1267 : lock_type lock(mutex_);
492 1267 : completed_ops_.push(h);
493 1267 : wake_one_thread_and_unlock(lock);
494 1267 : }
495 :
496 : inline void
497 26092 : reactor_scheduler::post(capy::continuation& c) const
498 : {
499 26092 : if (auto* ctx = reactor_find_context(this))
500 : {
501 15730 : ++ctx->private_outstanding_work;
502 15730 : ctx->private_queue.push(c);
503 15730 : return;
504 : }
505 :
506 10362 : outstanding_work_.fetch_add(1, std::memory_order_relaxed);
507 :
508 10362 : lock_type lock(mutex_);
509 10362 : completed_ops_.push(c);
510 10362 : wake_one_thread_and_unlock(lock);
511 10362 : }
512 :
513 : inline bool
514 10795 : reactor_scheduler::running_in_this_thread() const noexcept
515 : {
516 10795 : return reactor_find_context(this) != nullptr;
517 : }
518 :
519 : inline void
520 3736 : reactor_scheduler::stop()
521 : {
522 3736 : lock_type lock(mutex_);
523 3736 : if (!stopped_.load(std::memory_order_acquire))
524 : {
525 2868 : stopped_.store(true, std::memory_order_release);
526 2868 : signal_all(lock);
527 2868 : interrupt_reactor();
528 : }
529 3736 : }
530 :
531 : inline bool
532 2502 : reactor_scheduler::stopped() const noexcept
533 : {
534 2502 : return stopped_.load(std::memory_order_acquire);
535 : }
536 :
537 : inline void
538 1435 : reactor_scheduler::restart()
539 : {
540 1435 : stopped_.store(false, std::memory_order_release);
541 1435 : }
542 :
543 : inline std::size_t
544 2105 : reactor_scheduler::run()
545 : {
546 4210 : if (outstanding_work_.load(std::memory_order_acquire) == 0)
547 : {
548 108 : stop();
549 108 : return 0;
550 : }
551 :
552 1997 : reactor_thread_context_guard ctx(this);
553 1997 : lock_type lock(mutex_);
554 :
555 1997 : std::size_t n = 0;
556 : for (;;)
557 : {
558 483336 : if (!do_one(lock, -1, ctx.frame_))
559 1994 : break;
560 481339 : if (n != (std::numeric_limits<std::size_t>::max)())
561 481339 : ++n;
562 481339 : if (!lock.owns_lock())
563 381478 : lock.lock();
564 : }
565 1994 : return n;
566 2000 : }
567 :
568 : inline std::size_t
569 112 : reactor_scheduler::run_one()
570 : {
571 224 : if (outstanding_work_.load(std::memory_order_acquire) == 0)
572 : {
573 3 : stop();
574 3 : return 0;
575 : }
576 :
577 109 : reactor_thread_context_guard ctx(this);
578 109 : lock_type lock(mutex_);
579 109 : return do_one(lock, -1, ctx.frame_);
580 109 : }
581 :
582 : inline std::size_t
583 4092 : reactor_scheduler::wait_one(long usec)
584 : {
585 8184 : if (outstanding_work_.load(std::memory_order_acquire) == 0)
586 : {
587 785 : stop();
588 785 : return 0;
589 : }
590 :
591 3307 : reactor_thread_context_guard ctx(this);
592 3307 : lock_type lock(mutex_);
593 3307 : return do_one(lock, usec, ctx.frame_);
594 3307 : }
595 :
596 : inline std::size_t
597 49 : reactor_scheduler::poll()
598 : {
599 98 : if (outstanding_work_.load(std::memory_order_acquire) == 0)
600 : {
601 15 : stop();
602 15 : return 0;
603 : }
604 :
605 34 : reactor_thread_context_guard ctx(this);
606 34 : lock_type lock(mutex_);
607 :
608 34 : std::size_t n = 0;
609 : for (;;)
610 : {
611 75 : if (!do_one(lock, 0, ctx.frame_))
612 34 : break;
613 41 : if (n != (std::numeric_limits<std::size_t>::max)())
614 41 : ++n;
615 41 : if (!lock.owns_lock())
616 41 : lock.lock();
617 : }
618 34 : return n;
619 34 : }
620 :
621 : inline std::size_t
622 11 : reactor_scheduler::poll_one()
623 : {
624 22 : if (outstanding_work_.load(std::memory_order_acquire) == 0)
625 : {
626 5 : stop();
627 5 : return 0;
628 : }
629 :
630 6 : reactor_thread_context_guard ctx(this);
631 6 : lock_type lock(mutex_);
632 6 : return do_one(lock, 0, ctx.frame_);
633 6 : }
634 :
635 : inline void
636 36602 : reactor_scheduler::work_started() noexcept
637 : {
638 36602 : outstanding_work_.fetch_add(1, std::memory_order_relaxed);
639 36602 : }
640 :
641 : inline void
642 68694 : reactor_scheduler::work_finished() noexcept
643 : {
644 137388 : if (outstanding_work_.fetch_sub(1, std::memory_order_acq_rel) == 1)
645 2805 : stop();
646 68694 : }
647 :
648 : inline void
649 341447 : reactor_scheduler::compensating_work_started() const noexcept
650 : {
651 341447 : auto* ctx = reactor_find_context(this);
652 341447 : if (ctx)
653 341447 : ++ctx->private_outstanding_work;
654 341447 : }
655 :
656 : inline void
657 9710 : reactor_scheduler::post_deferred_completions(ready_queue& ops) const
658 : {
659 9710 : if (ops.empty())
660 9710 : return;
661 :
662 2 : if (auto* ctx = reactor_find_context(this))
663 : {
664 2 : ctx->private_queue.splice(ops);
665 2 : return;
666 : }
667 :
668 MIS 0 : lock_type lock(mutex_);
669 0 : completed_ops_.splice(ops);
670 0 : wake_one_thread_and_unlock(lock);
671 0 : }
672 :
673 : inline void
674 HIT 2310 : reactor_scheduler::shutdown_drain()
675 : {
676 2310 : lock_type lock(mutex_);
677 :
678 5010 : while (auto e = completed_ops_.pop())
679 : {
680 2700 : if (ready_is_continuation(e))
681 : {
682 8 : lock.unlock();
683 8 : if (auto h = ready_as_cont(e)->h)
684 8 : h.destroy();
685 8 : lock.lock();
686 : }
687 : else
688 : {
689 2692 : auto* op = ready_as_op(e);
690 2692 : if (op == &task_op_)
691 2307 : continue;
692 385 : lock.unlock();
693 385 : op->destroy();
694 385 : lock.lock();
695 : }
696 2700 : }
697 :
698 2310 : signal_all(lock);
699 2310 : }
700 :
701 : inline void
702 5178 : reactor_scheduler::signal_all(lock_type&) const
703 : {
704 5178 : state_ |= signaled_bit;
705 5178 : cond_.notify_all();
706 5178 : }
707 :
708 : inline bool
709 15289 : reactor_scheduler::maybe_unlock_and_signal_one(lock_type& lock) const
710 : {
711 15289 : state_ |= signaled_bit;
712 15289 : if (state_ > signaled_bit)
713 : {
714 12 : lock.unlock();
715 12 : cond_.notify_one();
716 12 : return true;
717 : }
718 15277 : return false;
719 : }
720 :
721 : inline bool
722 534159 : reactor_scheduler::unlock_and_signal_one(lock_type& lock) const
723 : {
724 534159 : state_ |= signaled_bit;
725 534159 : bool have_waiters = state_ > signaled_bit;
726 534159 : lock.unlock();
727 534159 : if (have_waiters)
728 6 : cond_.notify_one();
729 534159 : return have_waiters;
730 : }
731 :
732 : inline void
733 20 : reactor_scheduler::clear_signal() const
734 : {
735 20 : state_ &= ~signaled_bit;
736 20 : }
737 :
738 : inline void
739 6 : reactor_scheduler::wait_for_signal(lock_type& lock) const
740 : {
741 14 : while ((state_ & signaled_bit) == 0)
742 : {
743 8 : state_ += waiter_increment;
744 8 : cond_.wait(lock);
745 8 : state_ -= waiter_increment;
746 : }
747 6 : }
748 :
749 : inline void
750 14 : reactor_scheduler::wait_for_signal_for(lock_type& lock, long timeout_us) const
751 : {
752 14 : if ((state_ & signaled_bit) == 0)
753 : {
754 14 : state_ += waiter_increment;
755 14 : cond_.wait_for(lock, std::chrono::microseconds(timeout_us));
756 14 : state_ -= waiter_increment;
757 : }
758 14 : }
759 :
760 : inline void
761 15289 : reactor_scheduler::wake_one_thread_and_unlock(lock_type& lock) const
762 : {
763 15289 : if (maybe_unlock_and_signal_one(lock))
764 12 : return;
765 :
766 15277 : if (task_running_.load(std::memory_order_relaxed) && !task_interrupted_)
767 : {
768 790 : task_interrupted_ = true;
769 790 : lock.unlock();
770 790 : interrupt_reactor();
771 : }
772 : else
773 : {
774 14487 : lock.unlock();
775 : }
776 : }
777 :
778 483139 : inline reactor_scheduler::work_cleanup::~work_cleanup()
779 : {
780 483139 : std::int64_t produced = ctx.private_outstanding_work;
781 483139 : if (produced > 1)
782 346 : sched->outstanding_work_.fetch_add(
783 : produced - 1, std::memory_order_relaxed);
784 482793 : else if (produced < 1)
785 41872 : sched->work_finished();
786 483139 : ctx.private_outstanding_work = 0;
787 :
788 483139 : if (!ctx.private_queue.empty())
789 : {
790 100139 : lock->lock();
791 100139 : sched->completed_ops_.splice(ctx.private_queue);
792 : }
793 483139 : }
794 :
795 370997 : inline reactor_scheduler::task_cleanup::~task_cleanup()
796 : {
797 370997 : if (ctx.private_outstanding_work > 0)
798 : {
799 9448 : sched->outstanding_work_.fetch_add(
800 9448 : ctx.private_outstanding_work, std::memory_order_relaxed);
801 9448 : ctx.private_outstanding_work = 0;
802 : }
803 :
804 370997 : if (!ctx.private_queue.empty())
805 : {
806 9448 : if (!lock->owns_lock())
807 MIS 0 : lock->lock();
808 HIT 9448 : sched->completed_ops_.splice(ctx.private_queue);
809 : }
810 370997 : }
811 :
812 : inline std::size_t
813 486833 : reactor_scheduler::do_one(lock_type& lock, long timeout_us, context_type& ctx)
814 : {
815 : for (;;)
816 : {
817 856186 : if (stopped_.load(std::memory_order_acquire))
818 1998 : return 0;
819 :
820 854188 : std::uintptr_t e = completed_ops_.pop();
821 854188 : scheduler_op* op = ready_is_continuation(e) ? nullptr : ready_as_op(e);
822 :
823 : // Handle reactor sentinel — time to poll for I/O
824 854188 : if (op == &task_op_)
825 : {
826 371029 : bool more_handlers = !completed_ops_.empty();
827 :
828 691001 : if (!more_handlers &&
829 639944 : (outstanding_work_.load(std::memory_order_acquire) == 0 ||
830 : timeout_us == 0))
831 : {
832 32 : completed_ops_.push(&task_op_);
833 32 : return 0;
834 : }
835 :
836 370997 : long task_timeout_us = more_handlers ? 0 : timeout_us;
837 370997 : task_interrupted_ = task_timeout_us == 0;
838 370997 : task_running_.store(true, std::memory_order_release);
839 :
840 : // Wake a peer to take the pending handlers while this thread
841 : // polls the reactor; skipped when one_thread_ (no peer exists).
842 370997 : if (more_handlers && !one_thread_)
843 51050 : unlock_and_signal_one(lock);
844 :
845 : try
846 : {
847 370997 : run_task(lock, ctx, task_timeout_us);
848 : }
849 3 : catch (...)
850 : {
851 3 : task_running_.store(false, std::memory_order_relaxed);
852 3 : throw;
853 3 : }
854 :
855 370994 : task_running_.store(false, std::memory_order_relaxed);
856 370994 : completed_ops_.push(&task_op_);
857 370994 : if (timeout_us > 0)
858 1661 : return 0;
859 369333 : continue;
860 369333 : }
861 :
862 : // Handle ready entry (op or continuation)
863 483159 : if (e != 0)
864 : {
865 483139 : bool more = !completed_ops_.empty();
866 :
867 483139 : if (more && !one_thread_)
868 : {
869 : // Wake a peer for the remaining work; unassisted if none
870 : // was parked to take it.
871 483109 : ctx.unassisted = !unlock_and_signal_one(lock);
872 : }
873 : else
874 : {
875 : // No peer to wake (one_thread_, or nothing more queued).
876 30 : ctx.unassisted = more;
877 30 : lock.unlock();
878 : }
879 :
880 483139 : [[maybe_unused]] work_cleanup on_exit{this, &lock, ctx};
881 :
882 483139 : if (ready_is_continuation(e))
883 26084 : ready_as_cont(e)->h.resume();
884 : else
885 457055 : (*op)();
886 483139 : return 1;
887 483139 : }
888 :
889 40 : if (outstanding_work_.load(std::memory_order_acquire) == 0 ||
890 : timeout_us == 0)
891 MIS 0 : return 0;
892 :
893 HIT 20 : clear_signal();
894 20 : if (timeout_us < 0)
895 6 : wait_for_signal(lock);
896 : else
897 14 : wait_for_signal_for(lock, timeout_us);
898 369353 : }
899 : }
900 :
901 : } // namespace boost::corosio::detail
902 :
903 : #endif // BOOST_COROSIO_NATIVE_DETAIL_REACTOR_REACTOR_SCHEDULER_HPP
|