lattice_mode/
idle_gate.rs1use std::sync::Arc;
50use std::sync::Mutex;
51
52use lattice_grammar::effect::Effect;
53use tokio::time::Instant;
54
55pub type IdleGateHandler = Box<dyn FnMut() -> Vec<Effect> + Send + 'static>;
61
62pub type IdleGateRegistryHandle = Arc<IdleGateRegistry>;
65
66struct Gate {
67 id: u64,
68 name: &'static str,
70 deadline: Option<Instant>,
73 handler: IdleGateHandler,
74}
75
76struct Inner {
77 next_id: u64,
78 gates: Vec<Gate>,
79}
80
81pub struct IdleGateRegistry {
83 inner: Mutex<Inner>,
84}
85
86impl IdleGateRegistry {
87 pub fn new() -> Self {
111 Self {
112 inner: Mutex::new(Inner {
113 next_id: 0,
114 gates: Vec::new(),
115 }),
116 }
117 }
118
119 pub fn register(
122 self: &Arc<Self>,
123 name: &'static str,
124 handler: IdleGateHandler,
125 ) -> IdleGateHandle {
126 let id = {
127 let mut g = self.lock();
128 let id = g.next_id;
129 g.next_id += 1;
130 g.gates.push(Gate {
131 id,
132 name,
133 deadline: None,
134 handler,
135 });
136 id
137 };
138 IdleGateHandle {
139 registry: Arc::clone(self),
140 id,
141 }
142 }
143
144 pub fn earliest(&self) -> Option<Instant> {
148 self.lock().gates.iter().filter_map(|g| g.deadline).min()
149 }
150
151 pub fn fire_elapsed(&self, now: Instant) -> Vec<Effect> {
158 let mut g = self.lock();
159 let mut effects = Vec::new();
160 for gate in g.gates.iter_mut() {
161 if gate.deadline.is_some_and(|d| d <= now) {
162 gate.deadline = None;
163 tracing::debug!(gate = gate.name, "idle gate fired");
164 effects.extend((gate.handler)());
165 }
166 }
167 effects
168 }
169
170 fn arm(&self, id: u64, at: Instant) {
171 if let Some(gate) = self.lock().gates.iter_mut().find(|g| g.id == id) {
172 gate.deadline = Some(at);
173 }
174 }
175
176 fn disarm(&self, id: u64) {
177 if let Some(gate) = self.lock().gates.iter_mut().find(|g| g.id == id) {
178 gate.deadline = None;
179 }
180 }
181
182 fn unregister(&self, id: u64) {
183 self.lock().gates.retain(|g| g.id != id);
184 }
185
186 #[doc(hidden)]
188 pub fn registered_count(&self) -> usize {
189 self.lock().gates.len()
190 }
191
192 fn lock(&self) -> std::sync::MutexGuard<'_, Inner> {
196 self.inner.lock().unwrap_or_else(|e| e.into_inner())
197 }
198}
199
200impl Default for IdleGateRegistry {
201 fn default() -> Self {
202 Self::new()
203 }
204}
205
206impl std::fmt::Debug for IdleGateRegistry {
207 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
208 f.debug_struct("IdleGateRegistry")
209 .field("registered_count", &self.registered_count())
210 .finish_non_exhaustive()
211 }
212}
213
214pub struct IdleGateHandle {
218 registry: Arc<IdleGateRegistry>,
219 id: u64,
220}
221
222impl IdleGateHandle {
223 pub fn arm(&self, at: Instant) {
227 self.registry.arm(self.id, at);
228 }
229
230 pub fn disarm(&self) {
232 self.registry.disarm(self.id);
233 }
234}
235
236impl Drop for IdleGateHandle {
237 fn drop(&mut self) {
238 self.registry.unregister(self.id);
239 }
240}
241
242impl std::fmt::Debug for IdleGateHandle {
243 fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
244 f.debug_struct("IdleGateHandle")
245 .field("id", &self.id)
246 .finish_non_exhaustive()
247 }
248}
249
250#[cfg(test)]
251mod tests {
252 use super::*;
253 use std::sync::atomic::{AtomicUsize, Ordering};
254 use std::time::Duration;
255
256 fn counting_gate(count: Arc<AtomicUsize>) -> IdleGateHandler {
257 Box::new(move || {
258 count.fetch_add(1, Ordering::SeqCst);
259 Vec::new()
260 })
261 }
262
263 #[tokio::test]
264 async fn an_unarmed_registry_parks_the_actor_sleep() {
265 let r = Arc::new(IdleGateRegistry::new());
266 let _h = r.register("test", counting_gate(Arc::new(AtomicUsize::new(0))));
267 assert!(
268 r.earliest().is_none(),
269 "a registered-but-disarmed gate contributes no deadline"
270 );
271 }
272
273 #[tokio::test]
274 async fn the_earlier_of_two_armed_gates_is_the_one_the_actor_sleeps_to() {
275 let r = Arc::new(IdleGateRegistry::new());
276 let early_count = Arc::new(AtomicUsize::new(0));
277 let late_count = Arc::new(AtomicUsize::new(0));
278 let early = r.register("early", counting_gate(Arc::clone(&early_count)));
279 let late = r.register("late", counting_gate(Arc::clone(&late_count)));
280
281 let now = Instant::now();
282 late.arm(now + Duration::from_millis(500));
283 early.arm(now + Duration::from_millis(100));
284 assert_eq!(
285 r.earliest(),
286 Some(now + Duration::from_millis(100)),
287 "the minimum across gates, not registration order"
288 );
289
290 r.fire_elapsed(now + Duration::from_millis(200));
292 assert_eq!(early_count.load(Ordering::SeqCst), 1);
293 assert_eq!(late_count.load(Ordering::SeqCst), 0, "not yet due");
294 assert_eq!(
295 r.earliest(),
296 Some(now + Duration::from_millis(500)),
297 "a fired gate disarms itself; the later one remains"
298 );
299
300 r.fire_elapsed(now + Duration::from_millis(600));
301 assert_eq!(late_count.load(Ordering::SeqCst), 1);
302 assert!(r.earliest().is_none(), "both fired and disarmed");
303 }
304
305 #[tokio::test]
306 async fn a_fired_gate_does_not_fire_again() {
307 let r = Arc::new(IdleGateRegistry::new());
308 let count = Arc::new(AtomicUsize::new(0));
309 let h = r.register("once", counting_gate(Arc::clone(&count)));
310 let now = Instant::now();
311 h.arm(now);
312 r.fire_elapsed(now);
313 r.fire_elapsed(now + Duration::from_secs(1));
314 assert_eq!(
315 count.load(Ordering::SeqCst),
316 1,
317 "firing disarms — otherwise a popup would reopen every loop \
318 iteration forever"
319 );
320 }
321
322 #[tokio::test]
323 async fn disarm_cancels_a_pending_fire() {
324 let r = Arc::new(IdleGateRegistry::new());
325 let count = Arc::new(AtomicUsize::new(0));
326 let h = r.register("cancelled", counting_gate(Arc::clone(&count)));
327 let now = Instant::now();
328 h.arm(now + Duration::from_millis(50));
329 h.disarm();
330 assert!(r.earliest().is_none());
331 r.fire_elapsed(now + Duration::from_secs(1));
332 assert_eq!(
333 count.load(Ordering::SeqCst),
334 0,
335 "the chord resolved before the delay elapsed — no popup"
336 );
337 }
338
339 #[tokio::test]
340 async fn re_arming_replaces_the_deadline() {
341 let r = Arc::new(IdleGateRegistry::new());
342 let h = r.register("regrow", counting_gate(Arc::new(AtomicUsize::new(0))));
343 let now = Instant::now();
344 h.arm(now + Duration::from_millis(100));
345 h.arm(now + Duration::from_millis(300));
346 assert_eq!(
347 r.earliest(),
348 Some(now + Duration::from_millis(300)),
349 "re-arm replaces rather than adding a second deadline"
350 );
351 }
352
353 #[tokio::test]
354 async fn dropping_the_handle_deregisters_the_gate() {
355 let r = Arc::new(IdleGateRegistry::new());
356 let count = Arc::new(AtomicUsize::new(0));
357 let now = Instant::now();
358 {
359 let h = r.register("scoped", counting_gate(Arc::clone(&count)));
360 h.arm(now);
361 assert_eq!(r.registered_count(), 1);
362 }
363 assert_eq!(
364 r.registered_count(),
365 0,
366 "a deactivated subsystem contributes no timer"
367 );
368 r.fire_elapsed(now + Duration::from_secs(1));
369 assert_eq!(count.load(Ordering::SeqCst), 0);
370 }
371
372 #[tokio::test]
373 async fn only_due_gates_fire_when_several_are_armed() {
374 let r = Arc::new(IdleGateRegistry::new());
375 let now = Instant::now();
376 let counts: Vec<Arc<AtomicUsize>> = (0..3).map(|_| Arc::new(AtomicUsize::new(0))).collect();
377 let handles: Vec<IdleGateHandle> = counts
378 .iter()
379 .enumerate()
380 .map(|(i, c)| {
381 let h = r.register("multi", counting_gate(Arc::clone(c)));
382 h.arm(now + Duration::from_millis(100 * (i as u64 + 1)));
383 h
384 })
385 .collect();
386
387 r.fire_elapsed(now + Duration::from_millis(250));
388 assert_eq!(
389 counts
390 .iter()
391 .map(|c| c.load(Ordering::SeqCst))
392 .collect::<Vec<_>>(),
393 vec![1, 1, 0],
394 "the 100ms and 200ms gates fired; the 300ms one is still armed"
395 );
396 drop(handles);
397 }
398}