LCOV - code coverage report
Current view: top level - corosio/detail - timer_service.hpp (source / functions) Coverage Total Hit Missed
Test: coverage_remapped.info Lines: 96.8 % 222 215 7
Test Date: 2026-10-08 18:13:32 Functions: 100.0 % 29 29

           TLA  Line data    Source code
       1                 : //
       2                 : // Copyright (c) 2025 Vinnie Falco (vinnie.falco@gmail.com)
       3                 : // Copyright (c) 2026 Steve Gerbino
       4                 : //
       5                 : // Distributed under the Boost Software License, Version 1.0. (See accompanying
       6                 : // file LICENSE_1_0.txt or copy at http://www.boost.org/LICENSE_1_0.txt)
       7                 : //
       8                 : // Official repository: https://github.com/cppalliance/corosio
       9                 : //
      10                 : 
      11                 : #ifndef BOOST_COROSIO_DETAIL_TIMER_SERVICE_HPP
      12                 : #define BOOST_COROSIO_DETAIL_TIMER_SERVICE_HPP
      13                 : 
      14                 : #include <boost/corosio/detail/timer.hpp>
      15                 : #include <boost/corosio/detail/scheduler.hpp>
      16                 : #include <boost/corosio/detail/scheduler_op.hpp>
      17                 : #include <boost/corosio/detail/object_pool.hpp>
      18                 : #include <boost/corosio/detail/object_ref.hpp>
      19                 : #include <boost/corosio/detail/intrusive.hpp>
      20                 : #include <boost/corosio/detail/thread_local_ptr.hpp>
      21                 : #include <boost/capy/error.hpp>
      22                 : #include <boost/capy/ex/execution_context.hpp>
      23                 : #include <boost/capy/ex/executor_ref.hpp>
      24                 : #include <system_error>
      25                 : 
      26                 : #include <atomic>
      27                 : #include <chrono>
      28                 : #include <coroutine>
      29                 : #include <cstddef>
      30                 : #include <limits>
      31                 : #include <mutex>
      32                 : #include <stop_token>
      33                 : #include <utility>
      34                 : #include <vector>
      35                 : 
      36                 : namespace boost::corosio::detail {
      37                 : 
      38                 : struct scheduler;
      39                 : 
      40                 : /*
      41                 :     Timer Service
      42                 :     =============
      43                 : 
      44                 :     Data Structures
      45                 :     ---------------
      46                 :     waiter_node (defined in timer.hpp) holds per-waiter state:
      47                 :     coroutine handle, executor, error output, embedded
      48                 :     completion_op. Each concurrent co_await t.wait() embeds one
      49                 :     waiter_node in the awaitable on the suspended coroutine's
      50                 :     frame — waits perform no allocation.
      51                 : 
      52                 :     timer::implementation holds per-timer state: expiry, heap
      53                 :     index, and the single published waiter. Each timer holds
      54                 :     at most one waiter; process_expired's local cross-timer drain
      55                 :     list still threads waiters through their intrusive hooks when
      56                 :     collecting several timers' waiters past the lock.
      57                 : 
      58                 :     timer_service owns a min-heap of active timers and recycles
      59                 :     retired impls through a per-service object_pool, fronted by a
      60                 :     thread-local single-slot cache. The heap is ordered by expiry
      61                 :     time; the scheduler queries nearest_expiry() to set the
      62                 :     epoll/timerfd timeout.
      63                 : 
      64                 :     Optimization Strategy
      65                 :     ---------------------
      66                 :     1. Deferred heap insertion — expires_after() stores the expiry
      67                 :        but does not insert into the heap. Insertion happens in wait().
      68                 :     2. Thread-local impl cache — single-slot per-thread cache, in
      69                 :        front of the object_pool fallback.
      70                 :     3. Frame-resident waiter_node with embedded completion_op —
      71                 :        eliminates heap allocation per wait/fire/cancel.
      72                 :     4. Cached nearest expiry — atomic avoids mutex in nearest_expiry().
      73                 :     5. might_have_pending_waits_ flag — skips lock when no wait issued.
      74                 : 
      75                 :     Concurrency
      76                 :     -----------
      77                 :     stop_token callbacks can fire from any thread. The impl_
      78                 :     pointer on waiter_node is used as a "still in list" marker.
      79                 :     A waiter_node's storage is the suspended coroutine's frame:
      80                 :     every completion path must finish touching the node before
      81                 :     posting the continuation or destroying the handle.
      82                 : */
      83                 : 
      84                 : inline void timer_service_invalidate_cache() noexcept;
      85                 : 
      86                 : /** The timer service's recycling pool.
      87                 : 
      88                 :     Adds no behavior of its own — it exists only to re-expose
      89                 :     `adopt()` and `remove()` as public, so the timer service's
      90                 :     thread-local cache can transfer a cached impl's `live_` tracking
      91                 :     on both sides of the handoff: `construct()` re-adopts a
      92                 :     TL-cached impl as live, `shutdown()` detaches one before its
      93                 :     direct delete.
      94                 : */
      95                 : class timer_object_pool : public object_pool<timer::implementation>
      96                 : {
      97                 : public:
      98                 :     using object_pool::adopt;
      99                 :     using object_pool::remove;
     100                 : };
     101                 : 
     102                 : // timer_service class body — member function definitions are
     103                 : // out-of-class (after implementation and waiter_node are complete)
     104                 : class BOOST_COROSIO_DECL timer_service final
     105                 :     : public capy::execution_context::service
     106                 :     , public io_object::io_service
     107                 : {
     108                 :     friend struct timer::implementation;
     109                 : 
     110                 : public:
     111                 :     using clock_type = std::chrono::steady_clock;
     112                 :     using time_point = clock_type::time_point;
     113                 : 
     114                 :     /// Type-erased callback for earliest-expiry-changed notifications.
     115                 :     class callback
     116                 :     {
     117                 :         void* ctx_         = nullptr;
     118                 :         void (*fn_)(void*) = nullptr;
     119                 : 
     120                 :     public:
     121                 :         /// Construct an empty callback.
     122 HIT        2310 :         callback() = default;
     123                 : 
     124                 :         /// Construct a callback with the given context and function.
     125            2310 :         callback(void* ctx, void (*fn)(void*)) noexcept : ctx_(ctx), fn_(fn) {}
     126                 : 
     127                 :         /// Return true if the callback is non-empty.
     128                 :         explicit operator bool() const noexcept
     129                 :         {
     130                 :             return fn_ != nullptr;
     131                 :         }
     132                 : 
     133                 :         /// Invoke the callback.
     134            7496 :         void operator()() const
     135                 :         {
     136            7496 :             if (fn_)
     137            7496 :                 fn_(ctx_);
     138            7496 :         }
     139                 :     };
     140                 : 
     141                 : private:
     142                 :     struct heap_entry
     143                 :     {
     144                 :         time_point time_;
     145                 :         timer::implementation* timer_;
     146                 :     };
     147                 : 
     148                 :     scheduler* sched_ = nullptr;
     149                 :     BOOST_COROSIO_MSVC_WARNING_PUSH
     150                 :     BOOST_COROSIO_MSVC_WARNING_DISABLE(4251) // std:: members, dll-interface
     151                 :     mutable std::mutex mutex_;
     152                 :     std::vector<heap_entry> heap_;
     153                 :     timer_object_pool pool_;
     154                 :     callback on_earliest_changed_;
     155                 :     bool shutting_down_ = false;
     156                 :     // Avoids mutex in nearest_expiry() and empty()
     157                 :     mutable std::atomic<std::int64_t> cached_nearest_ns_{
     158                 :         (std::numeric_limits<std::int64_t>::max)()};
     159                 :     BOOST_COROSIO_MSVC_WARNING_POP
     160                 : 
     161                 : public:
     162                 :     /// Construct the timer service bound to a scheduler.
     163            2310 :     inline timer_service(capy::execution_context&, scheduler& sched)
     164            2310 :         : sched_(&sched)
     165                 :     {
     166            2310 :     }
     167                 : 
     168                 :     /// Return the associated scheduler.
     169           32737 :     inline scheduler& get_scheduler() noexcept
     170                 :     {
     171           32737 :         return *sched_;
     172                 :     }
     173                 : 
     174                 :     /// Destroy the timer service.
     175            4620 :     ~timer_service() override = default;
     176                 : 
     177                 :     timer_service(timer_service const&)            = delete;
     178                 :     timer_service& operator=(timer_service const&) = delete;
     179                 : 
     180                 :     /// Register a callback invoked when the earliest expiry changes.
     181            2310 :     inline void set_on_earliest_changed(callback cb)
     182                 :     {
     183            2310 :         on_earliest_changed_ = cb;
     184            2310 :     }
     185                 : 
     186                 :     /// Return true if no timers are in the heap.
     187                 :     inline bool empty() const noexcept
     188                 :     {
     189                 :         return cached_nearest_ns_.load(std::memory_order_acquire) ==
     190                 :             (std::numeric_limits<std::int64_t>::max)();
     191                 :     }
     192                 : 
     193                 :     /// Return the nearest timer expiry without acquiring the mutex.
     194          265754 :     inline time_point nearest_expiry() const noexcept
     195                 :     {
     196          265754 :         auto ns = cached_nearest_ns_.load(std::memory_order_acquire);
     197          265754 :         return time_point(time_point::duration(ns));
     198                 :     }
     199                 : 
     200                 :     /// Cancel all pending timers and free cached resources.
     201                 :     inline void shutdown() override;
     202                 : 
     203                 :     /// Construct a new timer implementation.
     204                 :     inline io_object::implementation* construct() override;
     205                 : 
     206                 :     /// Destroy a timer implementation, cancelling pending waiters.
     207                 :     inline void destroy(io_object::implementation* p) override;
     208                 : 
     209                 :     /// Cancel and recycle a timer implementation.
     210                 :     inline void destroy_impl(timer::implementation& impl);
     211                 : 
     212                 :     /// Publish the timer's waiter and insert the timer into the heap.
     213                 :     inline void insert_waiter(timer::implementation& impl, waiter_node* w);
     214                 : 
     215                 :     /// Cancel the timer's published waiter, if any.
     216                 :     inline void cancel_timer(timer::implementation& impl);
     217                 : 
     218                 :     /// Cancel one specific waiter ( stop_token callback path ).
     219                 :     inline void cancel_waiter(waiter_node* w);
     220                 : 
     221                 :     /// Complete all waiters whose timers have expired.
     222                 :     inline std::size_t process_expired();
     223                 : 
     224                 : private:
     225          313523 :     inline void refresh_cached_nearest() noexcept
     226                 :     {
     227          313523 :         auto ns = heap_.empty() ? (std::numeric_limits<std::int64_t>::max)()
     228          306716 :                                 : heap_[0].time_.time_since_epoch().count();
     229          313523 :         cached_nearest_ns_.store(ns, std::memory_order_release);
     230          313523 :     }
     231                 : 
     232                 :     inline void remove_timer_impl(timer::implementation& impl);
     233                 :     inline void up_heap(std::size_t index);
     234                 :     inline void down_heap(std::size_t index);
     235                 :     inline void swap_heap(std::size_t i1, std::size_t i2);
     236                 : };
     237                 : 
     238                 : // Thread-local cache avoids hot-path mutex acquisitions:
     239                 : // single-slot impl cache, validated by comparing svc_. Cleared by
     240                 : // timer_service_invalidate_cache() during shutdown.
     241                 : 
     242                 : inline thread_local_ptr<timer::implementation> tl_cached_impl;
     243                 : 
     244                 : // The POD TLS slot above never runs destructors, so a short-lived
     245                 : // run() thread would leak its cached impl. Each push arms this
     246                 : // owner, whose destructor frees the slot at thread exit. A cached
     247                 : // entry is a quiescent heap object (nothing in the heap or free
     248                 : // list) and deletion touches no service state, so it is safe after
     249                 : // the owning service is gone (the stale-entry path in
     250                 : // try_pop_tl_cache deletes the same way).
     251                 : struct tl_cache_owner
     252                 : {
     253              64 :     ~tl_cache_owner()
     254                 :     {
     255              64 :         delete tl_cached_impl.get();
     256              64 :         tl_cached_impl.set(nullptr);
     257              64 :     }
     258                 : };
     259                 : 
     260                 : inline void
     261           14815 : arm_tl_cache_cleanup() noexcept
     262                 : {
     263           14815 :     [[maybe_unused]] thread_local tl_cache_owner owner;
     264           14815 : }
     265                 : 
     266                 : inline timer::implementation*
     267           17235 : try_pop_tl_cache(timer_service* svc) noexcept
     268                 : {
     269           17235 :     auto* impl = tl_cached_impl.get();
     270           17235 :     if (impl)
     271                 :     {
     272           14357 :         tl_cached_impl.set(nullptr);
     273           14357 :         if (impl->svc_ == svc)
     274           14357 :             return impl;
     275                 :         // Stale impl from a destroyed service
     276 MIS           0 :         delete impl;
     277                 :     }
     278 HIT        2878 :     return nullptr;
     279                 : }
     280                 : 
     281                 : inline bool
     282           17206 : try_push_tl_cache(timer::implementation* impl) noexcept
     283                 : {
     284           17206 :     if (!tl_cached_impl.get())
     285                 :     {
     286           14815 :         arm_tl_cache_cleanup();
     287           14815 :         tl_cached_impl.set(impl);
     288           14815 :         return true;
     289                 :     }
     290            2391 :     return false;
     291                 : }
     292                 : 
     293                 : inline void
     294            2310 : timer_service_invalidate_cache() noexcept
     295                 : {
     296            2310 :     delete tl_cached_impl.get();
     297            2310 :     tl_cached_impl.set(nullptr);
     298            2310 : }
     299                 : 
     300                 : // timer_service out-of-class member function definitions
     301                 : 
     302                 : inline void
     303            2310 : timer_service::shutdown()
     304                 : {
     305            2310 :     timer_service_invalidate_cache();
     306            2310 :     shutting_down_ = true;
     307                 :     // Flag only; the real drain walks the expiry heap below, not
     308                 :     // object_pool's live_ list, so the visit callback is a no-op.
     309            2310 :     pool_.shutdown([](timer::implementation*) {});
     310                 : 
     311                 :     // Snapshot impls and detach them from the heap so that
     312                 :     // coroutine-owned timer destructors (triggered by h.destroy()
     313                 :     // below) cannot re-enter remove_timer_impl() and mutate the
     314                 :     // vector during iteration.
     315            2310 :     std::vector<timer::implementation*> impls;
     316            2310 :     impls.reserve(heap_.size());
     317            2339 :     for (auto& entry : heap_)
     318                 :     {
     319              29 :         entry.timer_->heap_index_.store(
     320                 :             (std::numeric_limits<std::size_t>::max)(),
     321                 :             std::memory_order_relaxed);
     322              29 :         impls.push_back(entry.timer_);
     323                 :     }
     324            2310 :     heap_.clear();
     325            2310 :     cached_nearest_ns_.store(
     326                 :         (std::numeric_limits<std::int64_t>::max)(), std::memory_order_release);
     327                 : 
     328                 :     // Cancel waiting timers. Each waiter called work_started()
     329                 :     // in implementation::wait(). On IOCP the scheduler shutdown
     330                 :     // loop exits when outstanding_work_ reaches zero, so we must
     331                 :     // call work_finished() here to balance it. On other backends
     332                 :     // this is harmless.
     333            2339 :     for (auto* impl : impls)
     334                 :     {
     335              29 :         if (auto* w = std::exchange(impl->waiter_, nullptr))
     336                 :         {
     337              29 :             w->reset_stop_cb();
     338              29 :             auto h = std::exchange(w->h_, {});
     339              29 :             sched_->work_finished();
     340                 :             // Destroying the frame also ends the node's storage
     341              29 :             if (h)
     342              29 :                 h.destroy();
     343                 :         }
     344                 :         // Unlink from the pool's live_ list before the direct delete
     345                 :         // below, or ~object_pool()'s unconditional sweep would delete
     346                 :         // this impl a second time.
     347              29 :         [[maybe_unused]] bool const removed = pool_.remove(impl);
     348              29 :         BOOST_COROSIO_ASSERT(removed);
     349              29 :         delete impl;
     350                 :     }
     351                 : 
     352                 :     // Anything still parked in pool_ (free-listed or live with no
     353                 :     // pending wait at shutdown) is freed by object_pool's own destructor.
     354            2310 : }
     355                 : 
     356                 : inline io_object::implementation*
     357           17235 : timer_service::construct()
     358                 : {
     359           17235 :     timer::implementation* impl = try_pop_tl_cache(this);
     360           17235 :     if (impl)
     361                 :     {
     362           14357 :         impl->svc_ = this;
     363           14357 :         impl->reuse();
     364           14357 :         impl->refs_.store(1, std::memory_order_relaxed);
     365                 :         // The TL push that parked impl here also removed it from
     366                 :         // live_ (see destroy_impl) — re-adopt it now that it is live
     367                 :         // again.
     368           14357 :         pool_.adopt(impl);
     369           14357 :         return impl;
     370                 :     }
     371                 : 
     372            2878 :     return pool_.acquire(*this);
     373                 : }
     374                 : 
     375                 : inline void
     376           17235 : timer_service::destroy(io_object::implementation* p)
     377                 : {
     378                 :     // During shutdown the drain loop owns every impl and deletes
     379                 :     // them directly. A frame destroyed by that loop can unwind a
     380                 :     // handle whose impl was freed in an earlier iteration (a
     381                 :     // timeout's parent frame owns the timeout timer while
     382                 :     // suspended on the inner delay's timer), so bail out before
     383                 :     // even downcasting the pointer.
     384           17235 :     if (shutting_down_)
     385              29 :         return;
     386           17206 :     destroy_impl(static_cast<timer::implementation&>(*p));
     387                 : }
     388                 : 
     389                 : inline void
     390           17206 : timer_service::destroy_impl(timer::implementation& impl)
     391                 : {
     392                 :     // During shutdown the impl is owned by the shutdown loop.
     393                 :     // Re-entering here (from a coroutine-owned timer destructor
     394                 :     // triggered by h.destroy()) must not modify the heap or
     395                 :     // recycle the impl — shutdown deletes it directly.
     396           17206 :     if (shutting_down_)
     397 MIS           0 :         return;
     398                 : 
     399 HIT       17206 :     cancel_timer(impl);
     400                 : 
     401           34412 :     if (impl.heap_index_.load(std::memory_order_relaxed) !=
     402           17206 :         (std::numeric_limits<std::size_t>::max)())
     403                 :     {
     404 MIS           0 :         std::lock_guard lock(mutex_);
     405               0 :         remove_timer_impl(impl);
     406               0 :         refresh_cached_nearest();
     407               0 :     }
     408                 : 
     409 HIT       17206 :     if (try_push_tl_cache(&impl))
     410                 :     {
     411                 :         // The TL slot takes sole ownership: it is deleted directly by
     412                 :         // tl_cache_owner / the stale-service path, never through this
     413                 :         // pool, so it must leave live_ now or ~object_pool() would see
     414                 :         // it as abandoned (and, at the wrong moment, double-delete it
     415                 :         // after a later construct() re-adopts it).
     416           14815 :         [[maybe_unused]] bool const removed = pool_.remove(&impl);
     417           14815 :         BOOST_COROSIO_ASSERT(removed);
     418           14815 :         return;
     419                 :     }
     420                 : 
     421                 :     // No op keepalives on a timer impl: this is the only reference,
     422                 :     // so release() always recycles immediately via retire().
     423            2391 :     release(&impl);
     424                 : }
     425                 : 
     426                 : inline void
     427           22780 : timer_service::insert_waiter(timer::implementation& impl, waiter_node* w)
     428                 : {
     429           22780 :     bool notify      = false;
     430           22780 :     bool lost_cancel = false;
     431                 :     {
     432           22780 :         std::lock_guard lock(mutex_);
     433                 :         // Grow before publishing anything, so the push_back below
     434                 :         // cannot throw: a failure here leaves the waiter untouched,
     435                 :         // the strong guarantee rearm_wait's recovery relies on.
     436           22780 :         if (impl.heap_index_.load(std::memory_order_relaxed) ==
     437           45560 :                 (std::numeric_limits<std::size_t>::max)() &&
     438           22780 :             heap_.size() == heap_.capacity())
     439             499 :             heap_.reserve(heap_.capacity() == 0 ? 16 : 2 * heap_.capacity());
     440                 :         // Publish: from here the waiter is visible to the fire path and
     441                 :         // to its own stop callback (impl_ non-null enables cancel_waiter).
     442           22780 :         w->impl_ = &impl;
     443           45560 :         if (impl.heap_index_.load(std::memory_order_relaxed) ==
     444           22780 :             (std::numeric_limits<std::size_t>::max)())
     445                 :         {
     446           22780 :             impl.heap_index_.store(heap_.size(), std::memory_order_relaxed);
     447           22780 :             heap_.push_back({impl.expiry_, &impl});
     448           22780 :             up_heap(heap_.size() - 1);
     449           22780 :             notify = (impl.heap_index_.load(std::memory_order_relaxed) == 0);
     450           22780 :             refresh_cached_nearest();
     451                 :         }
     452           22780 :         BOOST_COROSIO_ASSERT(impl.waiter_ == nullptr);
     453           22780 :         impl.waiter_ = w;
     454                 : 
     455                 :         // Lost-cancel re-check: a stop requested after the canceller was
     456                 :         // armed in wait() but before this publication found impl_ null
     457                 :         // and returned a no-op. Observe it now and undo the insertion.
     458           22780 :         if (w->token_->stop_requested())
     459                 :         {
     460             477 :             w->impl_     = nullptr;
     461             477 :             impl.waiter_ = nullptr;
     462             477 :             remove_timer_impl(impl);
     463             477 :             impl.might_have_pending_waits_.store(
     464                 :                 false, std::memory_order_relaxed);
     465             477 :             refresh_cached_nearest();
     466             477 :             lost_cancel = true;
     467             477 :             notify      = false; // insertion undone; nearest unchanged
     468                 :         }
     469           22780 :     }
     470           22780 :     if (notify)
     471            7496 :         on_earliest_changed_();
     472           22780 :     if (lost_cancel)
     473                 :     {
     474             477 :         w->ec_ = make_error_code(capy::error::canceled);
     475             477 :         sched_->post(&w->op_);
     476                 :     }
     477           22780 : }
     478                 : 
     479                 : inline void
     480           17206 : timer_service::cancel_timer(timer::implementation& impl)
     481                 : {
     482           17206 :     if (!impl.might_have_pending_waits_.load(std::memory_order_relaxed))
     483           17204 :         return;
     484                 : 
     485                 :     // No unlocked already-done fast-out here: it would need the
     486                 :     // non-atomic waiter_ (a race with concurrent drains), and an
     487                 :     // index-only check is lifetime-unsafe because npos is stored
     488                 :     // before the drain finishes touching the impl. A stale-true
     489                 :     // flag is rare with the stateless API; the locked path below
     490                 :     // re-validates.
     491                 : 
     492               2 :     waiter_node* canceled = nullptr;
     493                 : 
     494                 :     {
     495               2 :         std::lock_guard lock(mutex_);
     496               2 :         remove_timer_impl(impl);
     497               2 :         canceled = std::exchange(impl.waiter_, nullptr);
     498               2 :         if (canceled)
     499               2 :             canceled->impl_ = nullptr;
     500                 :         // Store false as the final touch of the impl under the lock so
     501                 :         // a pre-lock false-flag check trusts it unqualified.
     502               2 :         impl.might_have_pending_waits_.store(false, std::memory_order_relaxed);
     503               2 :         refresh_cached_nearest();
     504               2 :     }
     505                 : 
     506               2 :     if (canceled)
     507                 :     {
     508               2 :         canceled->ec_ = make_error_code(capy::error::canceled);
     509               2 :         sched_->post(&canceled->op_);
     510                 :     }
     511                 : }
     512                 : 
     513                 : inline void
     514            2689 : timer_service::cancel_waiter(waiter_node* w)
     515                 : {
     516                 :     {
     517            2689 :         std::lock_guard lock(mutex_);
     518                 :         // Already removed by another drain: cancel_timer,
     519                 :         // process_expired, or insert_waiter's lost-cancel recheck
     520            2689 :         if (!w->impl_)
     521             506 :             return;
     522            2183 :         auto* impl    = w->impl_;
     523            2183 :         w->impl_      = nullptr;
     524            2183 :         impl->waiter_ = nullptr;
     525            2183 :         remove_timer_impl(*impl);
     526            2183 :         impl->might_have_pending_waits_.store(false, std::memory_order_relaxed);
     527            2183 :         refresh_cached_nearest();
     528            2689 :     }
     529                 : 
     530            2183 :     w->ec_ = make_error_code(capy::error::canceled);
     531            2183 :     sched_->post(&w->op_);
     532                 : }
     533                 : 
     534                 : inline std::size_t
     535          288081 : timer_service::process_expired()
     536                 : {
     537          288081 :     intrusive_list<waiter_node> expired;
     538                 : 
     539                 :     {
     540          288081 :         std::lock_guard lock(mutex_);
     541          288081 :         auto now = clock_type::now();
     542                 : 
     543          308170 :         while (!heap_.empty() && heap_[0].time_ <= now)
     544                 :         {
     545           20089 :             timer::implementation* t = heap_[0].timer_;
     546           20089 :             remove_timer_impl(*t);
     547           20089 :             if (auto* w = std::exchange(t->waiter_, nullptr))
     548                 :             {
     549           20089 :                 w->impl_ = nullptr;
     550           20089 :                 w->ec_   = {};
     551           20089 :                 expired.push_back(w);
     552                 :             }
     553           20089 :             t->might_have_pending_waits_.store(
     554                 :                 false, std::memory_order_relaxed);
     555                 :         }
     556                 : 
     557          288081 :         refresh_cached_nearest();
     558          288081 :     }
     559                 : 
     560          288081 :     std::size_t count = 0;
     561          308170 :     while (auto* w = expired.pop_front())
     562                 :     {
     563           20089 :         sched_->post(&w->op_);
     564           20089 :         ++count;
     565           20089 :     }
     566                 : 
     567          288081 :     return count;
     568                 : }
     569                 : 
     570                 : inline void
     571           22751 : timer_service::remove_timer_impl(timer::implementation& impl)
     572                 : {
     573           22751 :     std::size_t index = impl.heap_index_.load(std::memory_order_relaxed);
     574           22751 :     if (index >= heap_.size())
     575 MIS           0 :         return; // Not in heap
     576                 : 
     577 HIT       22751 :     if (index == heap_.size() - 1)
     578                 :     {
     579                 :         // Last element, just pop
     580            3233 :         impl.heap_index_.store(
     581                 :             (std::numeric_limits<std::size_t>::max)(),
     582                 :             std::memory_order_relaxed);
     583            3233 :         heap_.pop_back();
     584                 :     }
     585                 :     else
     586                 :     {
     587                 :         // Swap with last and reheapify
     588           19518 :         swap_heap(index, heap_.size() - 1);
     589           19518 :         impl.heap_index_.store(
     590                 :             (std::numeric_limits<std::size_t>::max)(),
     591                 :             std::memory_order_relaxed);
     592           19518 :         heap_.pop_back();
     593                 : 
     594           19518 :         if (index > 0 && heap_[index].time_ < heap_[(index - 1) / 2].time_)
     595               3 :             up_heap(index);
     596                 :         else
     597           19515 :             down_heap(index);
     598                 :     }
     599                 : }
     600                 : 
     601                 : inline void
     602           22783 : timer_service::up_heap(std::size_t index)
     603                 : {
     604           30601 :     while (index > 0)
     605                 :     {
     606           22738 :         std::size_t parent = (index - 1) / 2;
     607           22738 :         if (!(heap_[index].time_ < heap_[parent].time_))
     608           14920 :             break;
     609            7818 :         swap_heap(index, parent);
     610            7818 :         index = parent;
     611                 :     }
     612           22783 : }
     613                 : 
     614                 : inline void
     615           19515 : timer_service::down_heap(std::size_t index)
     616                 : {
     617           19515 :     std::size_t child = index * 2 + 1;
     618           46071 :     while (child < heap_.size())
     619                 :     {
     620           29555 :         std::size_t min_child = (child + 1 == heap_.size() ||
     621           25385 :                                  heap_[child].time_ < heap_[child + 1].time_)
     622           54940 :             ? child
     623           29555 :             : child + 1;
     624                 : 
     625           29555 :         if (heap_[index].time_ < heap_[min_child].time_)
     626            2999 :             break;
     627                 : 
     628           26556 :         swap_heap(index, min_child);
     629           26556 :         index = min_child;
     630           26556 :         child = index * 2 + 1;
     631                 :     }
     632           19515 : }
     633                 : 
     634                 : inline void
     635           53892 : timer_service::swap_heap(std::size_t i1, std::size_t i2)
     636                 : {
     637           53892 :     heap_entry tmp = heap_[i1];
     638           53892 :     heap_[i1]      = heap_[i2];
     639           53892 :     heap_[i2]      = tmp;
     640           53892 :     heap_[i1].timer_->heap_index_.store(i1, std::memory_order_relaxed);
     641           53892 :     heap_[i2].timer_->heap_index_.store(i2, std::memory_order_relaxed);
     642           53892 : }
     643                 : 
     644                 : // waiter_node's completion_op and canceller members are defined in
     645                 : // timer.cpp alongside implementation::wait(), for the same reason
     646                 : // wait() lives there (see below).
     647                 : 
     648                 : // timer::implementation::wait() is defined in timer.cpp, not here.
     649                 : // It must be a non-inline definition in a translation unit that is
     650                 : // always pulled into the link whenever detail::timer is used (every
     651                 : // consumer needs timer's constructors from that same object file).
     652                 : // An inline definition in this header would only be emitted in
     653                 : // translation units that happen to also include this header, which
     654                 : // is not guaranteed for every caller of wait_awaitable::await_suspend
     655                 : // in timer.hpp (e.g. code that only reaches timer.hpp through
     656                 : // delay.hpp, without transitively including a scheduler header).
     657                 : 
     658                 : // Free functions
     659                 : 
     660                 : inline timer_service&
     661            2310 : get_timer_service(capy::execution_context& ctx, scheduler& sched)
     662                 : {
     663            2310 :     return ctx.make_service<timer_service>(sched);
     664                 : }
     665                 : 
     666                 : } // namespace boost::corosio::detail
     667                 : 
     668                 : #endif
        

Generated by: LCOV version 2.3