include/boost/corosio/detail/timer_service.hpp

96.8% Lines (215 / 222) 100.0% Functions (28 / 28)
timer_service.hpp
f(x) Functions (28)
Function Calls Lines Blocks
boost::corosio::detail::timer_service::callback::callback() :122 2310x 100.0% 100.0% boost::corosio::detail::timer_service::callback::callback(void*, void (*)(void*)) :125 2310x 100.0% 100.0% boost::corosio::detail::timer_service::callback::operator()() const :134 7496x 100.0% 100.0% boost::corosio::detail::timer_service::timer_service(boost::capy::execution_context&, boost::corosio::detail::scheduler&) :163 2310x 100.0% 100.0% boost::corosio::detail::timer_service::get_scheduler() :169 32737x 100.0% 100.0% boost::corosio::detail::timer_service::~timer_service() :175 4620x 100.0% 100.0% boost::corosio::detail::timer_service::set_on_earliest_changed(boost::corosio::detail::timer_service::callback) :181 2310x 100.0% 100.0% boost::corosio::detail::timer_service::nearest_expiry() const :194 265754x 100.0% 73.0% boost::corosio::detail::timer_service::refresh_cached_nearest() :225 313523x 100.0% 70.0% boost::corosio::detail::tl_cache_owner::~tl_cache_owner() :253 64x 100.0% 100.0% boost::corosio::detail::arm_tl_cache_cleanup() :261 14815x 100.0% 100.0% boost::corosio::detail::try_pop_tl_cache(boost::corosio::detail::timer_service*) :267 17235x 87.5% 78.0% boost::corosio::detail::try_push_tl_cache(boost::corosio::detail::timer::implementation*) :282 17206x 100.0% 100.0% boost::corosio::detail::timer_service_invalidate_cache() :294 2310x 100.0% 100.0% boost::corosio::detail::timer_service::shutdown() :303 2310x 100.0% 72.0% boost::corosio::detail::timer_service::shutdown()::{lambda(boost::corosio::detail::timer::implementation*)#1}::operator()(boost::corosio::detail::timer::implementation*) const :309 29x 100.0% 100.0% boost::corosio::detail::timer_service::construct() :357 17235x 100.0% 71.0% boost::corosio::detail::timer_service::destroy(boost::corosio::io_object::implementation*) :376 17235x 100.0% 100.0% boost::corosio::detail::timer_service::destroy_impl(boost::corosio::detail::timer::implementation&) :390 17206x 66.7% 61.0% boost::corosio::detail::timer_service::insert_waiter(boost::corosio::detail::timer::implementation&, boost::corosio::detail::waiter_node*) :427 22780x 100.0% 73.0% boost::corosio::detail::timer_service::cancel_timer(boost::corosio::detail::timer::implementation&) :480 17206x 100.0% 89.0% boost::corosio::detail::timer_service::cancel_waiter(boost::corosio::detail::waiter_node*) :514 2689x 100.0% 88.0% boost::corosio::detail::timer_service::process_expired() :535 288081x 100.0% 91.0% boost::corosio::detail::timer_service::remove_timer_impl(boost::corosio::detail::timer::implementation&) :571 22751x 92.3% 71.0% boost::corosio::detail::timer_service::up_heap(unsigned long) :602 22783x 100.0% 100.0% boost::corosio::detail::timer_service::down_heap(unsigned long) :615 19515x 100.0% 100.0% boost::corosio::detail::timer_service::swap_heap(unsigned long, unsigned long) :635 53892x 100.0% 66.0% boost::corosio::detail::get_timer_service(boost::capy::execution_context&, boost::corosio::detail::scheduler&) :661 2310x 100.0% 100.0%
Line TLA Hits 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 2310x callback() = default;
123
124 /// Construct a callback with the given context and function.
125 2310x 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 7496x void operator()() const
135 {
136 7496x if (fn_)
137 7496x fn_(ctx_);
138 7496x }
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 2310x inline timer_service(capy::execution_context&, scheduler& sched)
164 2310x : sched_(&sched)
165 {
166 2310x }
167
168 /// Return the associated scheduler.
169 32737x inline scheduler& get_scheduler() noexcept
170 {
171 32737x return *sched_;
172 }
173
174 /// Destroy the timer service.
175 4620x ~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 2310x inline void set_on_earliest_changed(callback cb)
182 {
183 2310x on_earliest_changed_ = cb;
184 2310x }
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 265754x inline time_point nearest_expiry() const noexcept
195 {
196 265754x auto ns = cached_nearest_ns_.load(std::memory_order_acquire);
197 265754x 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 313523x inline void refresh_cached_nearest() noexcept
226 {
227 313523x auto ns = heap_.empty() ? (std::numeric_limits<std::int64_t>::max)()
228 306716x : heap_[0].time_.time_since_epoch().count();
229 313523x cached_nearest_ns_.store(ns, std::memory_order_release);
230 313523x }
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 64x ~tl_cache_owner()
254 {
255 64x delete tl_cached_impl.get();
256 64x tl_cached_impl.set(nullptr);
257 64x }
258 };
259
260 inline void
261 14815x arm_tl_cache_cleanup() noexcept
262 {
263 14815x [[maybe_unused]] thread_local tl_cache_owner owner;
264 14815x }
265
266 inline timer::implementation*
267 17235x try_pop_tl_cache(timer_service* svc) noexcept
268 {
269 17235x auto* impl = tl_cached_impl.get();
270 17235x if (impl)
271 {
272 14357x tl_cached_impl.set(nullptr);
273 14357x if (impl->svc_ == svc)
274 14357x return impl;
275 // Stale impl from a destroyed service
276 ✗ delete impl;
277 }
278 2878x return nullptr;
279 }
280
281 inline bool
282 17206x try_push_tl_cache(timer::implementation* impl) noexcept
283 {
284 17206x if (!tl_cached_impl.get())
285 {
286 14815x arm_tl_cache_cleanup();
287 14815x tl_cached_impl.set(impl);
288 14815x return true;
289 }
290 2391x return false;
291 }
292
293 inline void
294 2310x timer_service_invalidate_cache() noexcept
295 {
296 2310x delete tl_cached_impl.get();
297 2310x tl_cached_impl.set(nullptr);
298 2310x }
299
300 // timer_service out-of-class member function definitions
301
302 inline void
303 2310x timer_service::shutdown()
304 {
305 2310x timer_service_invalidate_cache();
306 2310x 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 2310x 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 2310x std::vector<timer::implementation*> impls;
316 2310x impls.reserve(heap_.size());
317 2339x for (auto& entry : heap_)
318 {
319 29x entry.timer_->heap_index_.store(
320 (std::numeric_limits<std::size_t>::max)(),
321 std::memory_order_relaxed);
322 29x impls.push_back(entry.timer_);
323 }
324 2310x heap_.clear();
325 2310x 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 2339x for (auto* impl : impls)
334 {
335 29x if (auto* w = std::exchange(impl->waiter_, nullptr))
336 {
337 29x w->reset_stop_cb();
338 29x auto h = std::exchange(w->h_, {});
339 29x sched_->work_finished();
340 // Destroying the frame also ends the node's storage
341 29x if (h)
342 29x 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 29x [[maybe_unused]] bool const removed = pool_.remove(impl);
348 29x BOOST_COROSIO_ASSERT(removed);
349 29x 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 2310x }
355
356 inline io_object::implementation*
357 17235x timer_service::construct()
358 {
359 17235x timer::implementation* impl = try_pop_tl_cache(this);
360 17235x if (impl)
361 {
362 14357x impl->svc_ = this;
363 14357x impl->reuse();
364 14357x 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 14357x pool_.adopt(impl);
369 14357x return impl;
370 }
371
372 2878x return pool_.acquire(*this);
373 }
374
375 inline void
376 17235x 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 17235x if (shutting_down_)
385 29x return;
386 17206x destroy_impl(static_cast<timer::implementation&>(*p));
387 }
388
389 inline void
390 17206x 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 17206x if (shutting_down_)
397 ✗ return;
398
399 17206x cancel_timer(impl);
400
401 34412x if (impl.heap_index_.load(std::memory_order_relaxed) !=
402 17206x (std::numeric_limits<std::size_t>::max)())
403 {
404 ✗ std::lock_guard lock(mutex_);
405 ✗ remove_timer_impl(impl);
406 ✗ refresh_cached_nearest();
407 ✗ }
408
409 17206x 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 14815x [[maybe_unused]] bool const removed = pool_.remove(&impl);
417 14815x BOOST_COROSIO_ASSERT(removed);
418 14815x 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 2391x release(&impl);
424 }
425
426 inline void
427 22780x timer_service::insert_waiter(timer::implementation& impl, waiter_node* w)
428 {
429 22780x bool notify = false;
430 22780x bool lost_cancel = false;
431 {
432 22780x 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 22780x if (impl.heap_index_.load(std::memory_order_relaxed) ==
437 45560x (std::numeric_limits<std::size_t>::max)() &&
438 22780x heap_.size() == heap_.capacity())
439 499x 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 22780x w->impl_ = &impl;
443 45560x if (impl.heap_index_.load(std::memory_order_relaxed) ==
444 22780x (std::numeric_limits<std::size_t>::max)())
445 {
446 22780x impl.heap_index_.store(heap_.size(), std::memory_order_relaxed);
447 22780x heap_.push_back({impl.expiry_, &impl});
448 22780x up_heap(heap_.size() - 1);
449 22780x notify = (impl.heap_index_.load(std::memory_order_relaxed) == 0);
450 22780x refresh_cached_nearest();
451 }
452 22780x BOOST_COROSIO_ASSERT(impl.waiter_ == nullptr);
453 22780x 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 22780x if (w->token_->stop_requested())
459 {
460 477x w->impl_ = nullptr;
461 477x impl.waiter_ = nullptr;
462 477x remove_timer_impl(impl);
463 477x impl.might_have_pending_waits_.store(
464 false, std::memory_order_relaxed);
465 477x refresh_cached_nearest();
466 477x lost_cancel = true;
467 477x notify = false; // insertion undone; nearest unchanged
468 }
469 22780x }
470 22780x if (notify)
471 7496x on_earliest_changed_();
472 22780x if (lost_cancel)
473 {
474 477x w->ec_ = make_error_code(capy::error::canceled);
475 477x sched_->post(&w->op_);
476 }
477 22780x }
478
479 inline void
480 17206x timer_service::cancel_timer(timer::implementation& impl)
481 {
482 17206x if (!impl.might_have_pending_waits_.load(std::memory_order_relaxed))
483 17204x 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 2x waiter_node* canceled = nullptr;
493
494 {
495 2x std::lock_guard lock(mutex_);
496 2x remove_timer_impl(impl);
497 2x canceled = std::exchange(impl.waiter_, nullptr);
498 2x if (canceled)
499 2x 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 2x impl.might_have_pending_waits_.store(false, std::memory_order_relaxed);
503 2x refresh_cached_nearest();
504 2x }
505
506 2x if (canceled)
507 {
508 2x canceled->ec_ = make_error_code(capy::error::canceled);
509 2x sched_->post(&canceled->op_);
510 }
511 }
512
513 inline void
514 2689x timer_service::cancel_waiter(waiter_node* w)
515 {
516 {
517 2689x std::lock_guard lock(mutex_);
518 // Already removed by another drain: cancel_timer,
519 // process_expired, or insert_waiter's lost-cancel recheck
520 2689x if (!w->impl_)
521 506x return;
522 2183x auto* impl = w->impl_;
523 2183x w->impl_ = nullptr;
524 2183x impl->waiter_ = nullptr;
525 2183x remove_timer_impl(*impl);
526 2183x impl->might_have_pending_waits_.store(false, std::memory_order_relaxed);
527 2183x refresh_cached_nearest();
528 2689x }
529
530 2183x w->ec_ = make_error_code(capy::error::canceled);
531 2183x sched_->post(&w->op_);
532 }
533
534 inline std::size_t
535 288081x timer_service::process_expired()
536 {
537 288081x intrusive_list<waiter_node> expired;
538
539 {
540 288081x std::lock_guard lock(mutex_);
541 288081x auto now = clock_type::now();
542
543 308170x while (!heap_.empty() && heap_[0].time_ <= now)
544 {
545 20089x timer::implementation* t = heap_[0].timer_;
546 20089x remove_timer_impl(*t);
547 20089x if (auto* w = std::exchange(t->waiter_, nullptr))
548 {
549 20089x w->impl_ = nullptr;
550 20089x w->ec_ = {};
551 20089x expired.push_back(w);
552 }
553 20089x t->might_have_pending_waits_.store(
554 false, std::memory_order_relaxed);
555 }
556
557 288081x refresh_cached_nearest();
558 288081x }
559
560 288081x std::size_t count = 0;
561 308170x while (auto* w = expired.pop_front())
562 {
563 20089x sched_->post(&w->op_);
564 20089x ++count;
565 20089x }
566
567 288081x return count;
568 }
569
570 inline void
571 22751x timer_service::remove_timer_impl(timer::implementation& impl)
572 {
573 22751x std::size_t index = impl.heap_index_.load(std::memory_order_relaxed);
574 22751x if (index >= heap_.size())
575 ✗ return; // Not in heap
576
577 22751x if (index == heap_.size() - 1)
578 {
579 // Last element, just pop
580 3233x impl.heap_index_.store(
581 (std::numeric_limits<std::size_t>::max)(),
582 std::memory_order_relaxed);
583 3233x heap_.pop_back();
584 }
585 else
586 {
587 // Swap with last and reheapify
588 19518x swap_heap(index, heap_.size() - 1);
589 19518x impl.heap_index_.store(
590 (std::numeric_limits<std::size_t>::max)(),
591 std::memory_order_relaxed);
592 19518x heap_.pop_back();
593
594 19518x if (index > 0 && heap_[index].time_ < heap_[(index - 1) / 2].time_)
595 3x up_heap(index);
596 else
597 19515x down_heap(index);
598 }
599 }
600
601 inline void
602 22783x timer_service::up_heap(std::size_t index)
603 {
604 30601x while (index > 0)
605 {
606 22738x std::size_t parent = (index - 1) / 2;
607 22738x if (!(heap_[index].time_ < heap_[parent].time_))
608 14920x break;
609 7818x swap_heap(index, parent);
610 7818x index = parent;
611 }
612 22783x }
613
614 inline void
615 19515x timer_service::down_heap(std::size_t index)
616 {
617 19515x std::size_t child = index * 2 + 1;
618 46071x while (child < heap_.size())
619 {
620 29555x std::size_t min_child = (child + 1 == heap_.size() ||
621 25385x heap_[child].time_ < heap_[child + 1].time_)
622 54940x ? child
623 29555x : child + 1;
624
625 29555x if (heap_[index].time_ < heap_[min_child].time_)
626 2999x break;
627
628 26556x swap_heap(index, min_child);
629 26556x index = min_child;
630 26556x child = index * 2 + 1;
631 }
632 19515x }
633
634 inline void
635 53892x timer_service::swap_heap(std::size_t i1, std::size_t i2)
636 {
637 53892x heap_entry tmp = heap_[i1];
638 53892x heap_[i1] = heap_[i2];
639 53892x heap_[i2] = tmp;
640 53892x heap_[i1].timer_->heap_index_.store(i1, std::memory_order_relaxed);
641 53892x heap_[i2].timer_->heap_index_.store(i2, std::memory_order_relaxed);
642 53892x }
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 2310x get_timer_service(capy::execution_context& ctx, scheduler& sched)
662 {
663 2310x return ctx.make_service<timer_service>(sched);
664 }
665
666 } // namespace boost::corosio::detail
667
668 #endif
669