1use std::cell::{Cell, RefCell};
9use std::collections::BTreeMap;
10use std::rc::Rc;
11use std::sync::Arc;
12use std::task::Poll;
13use std::time::Instant;
14
15use crate::metrics::Counters;
16
17pub(crate) struct Heap {
19 queue: BTreeMap<(Instant, u64), Rc<Slot>>,
20 seq: u64,
21 metrics: Arc<Counters>,
22}
23
24struct Slot {
26 key: Cell<Option<(Instant, u64)>>,
28 elapsed: Cell<bool>,
30 waiters: RefCell<kio::WaiterList>,
31}
32
33impl Heap {
34 pub fn new(metrics: Arc<Counters>) -> Self {
35 Self {
36 queue: BTreeMap::new(),
37 seq: 0,
38 metrics,
39 }
40 }
41
42 pub fn fire(&mut self, now: Instant) -> bool {
44 let mut fired = false;
45 while let Some(entry) = self.queue.first_entry() {
46 if entry.key().0 > now {
47 break;
48 }
49 let slot = entry.remove();
50 slot.key.set(None);
51 slot.elapsed.set(true);
52 slot.waiters.borrow_mut().wake();
53 self.metrics.timers_fired.add(1);
54 fired = true;
55 }
56 fired
57 }
58
59 pub fn next(&self) -> Option<Instant> {
61 self.queue.first_key_value().map(|(key, _)| key.0)
62 }
63
64 fn insert(&mut self, at: Instant, slot: Rc<Slot>) -> (Instant, u64) {
65 self.seq += 1;
66 let key = (at, self.seq);
67 self.queue.insert(key, slot);
68 self.metrics.timers_armed.add(1);
69 key
70 }
71
72 fn cancel(&mut self, key: (Instant, u64)) {
75 if self.queue.remove(&key).is_some() {
76 self.metrics.timers_cancelled.add(1);
77 }
78 }
79
80 fn fire_one(&mut self, key: (Instant, u64)) {
82 if self.queue.remove(&key).is_some() {
83 self.metrics.timers_fired.add(1);
84 }
85 }
86}
87
88pub struct Timer {
93 at: Option<Instant>,
94 heap: Rc<RefCell<Heap>>,
95 slot: Rc<Slot>,
96}
97
98impl Timer {
99 pub(crate) fn from_heap(heap: Rc<RefCell<Heap>>) -> Self {
100 Self {
101 at: None,
102 heap,
103 slot: Rc::new(Slot {
104 key: Cell::new(None),
105 elapsed: Cell::new(false),
106 waiters: RefCell::new(kio::WaiterList::new()),
107 }),
108 }
109 }
110}
111
112impl Timer {
113 pub fn set(&mut self, at: Option<Instant>) {
115 if self.at == at {
116 return;
117 }
118 self.at = at;
119 let mut heap = self.heap.borrow_mut();
120 if let Some(key) = self.slot.key.take() {
121 heap.cancel(key);
122 }
123 self.slot.elapsed.set(false);
124 if let Some(at) = at {
125 let key = heap.insert(at, self.slot.clone());
129 self.slot.key.set(Some(key));
130 }
131 }
132
133 pub fn poll(&mut self, waiter: &kio::Waiter) -> Poll<()> {
135 if self.slot.elapsed.get() {
136 return Poll::Ready(());
137 }
138 let Some((at, _)) = self.slot.key.get() else {
139 return Poll::Pending;
140 };
141 if at <= Instant::now() {
142 self.heap.borrow_mut().fire_one(self.slot.key.take().expect("armed"));
145 self.slot.elapsed.set(true);
146 return Poll::Ready(());
147 }
148 waiter.register(&mut self.slot.waiters.borrow_mut());
149 Poll::Pending
150 }
151}
152
153impl Drop for Timer {
154 fn drop(&mut self) {
155 if let Some(key) = self.slot.key.take() {
156 self.heap.borrow_mut().cancel(key);
157 }
158 }
159}
160
161impl std::fmt::Debug for Timer {
162 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
163 f.debug_struct("Timer")
164 .field("at", &self.slot.key.get().map(|key| key.0))
165 .field("elapsed", &self.slot.elapsed.get())
166 .finish()
167 }
168}
169
170impl Timer {
171 pub fn new(handle: &crate::Handle) -> Self {
173 handle.timer()
174 }
175 pub fn after(handle: &crate::Handle, duration: std::time::Duration) -> Self {
177 let mut timer = handle.timer();
178 timer.set(Instant::now().checked_add(duration));
179 timer
180 }
181 pub async fn wait(&mut self) {
183 kio::wait(|waiter| self.poll(waiter)).await
184 }
185}