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