//! In-memory throttle for Popfeed sync emissions. //! //! KOReader pings progress every few pages — sometimes seconds apart — which //! would flood a user's PDS with append-only `book.progress` log records. The //! throttle implements leading-edge emit with a trailing collapse: the first //! sync in an idle window fires immediately, additional syncs within the //! window are buffered (only the latest percent is kept), and a single //! trailing emit fires once the window closes. //! //! State lives in memory only. On process restart, all slots clear; the worst //! case is that one pending trailing emit per active document is dropped, but //! the next incoming sync will be "first after idle window" and emit //! immediately, so the eventual recorded progress is correct. use std::collections::HashMap; use std::sync::{Arc, Mutex}; use std::time::{Duration, Instant}; use tokio::task::JoinHandle; pub type DocKey = (i64, String); /// What to do with an incoming sync request. #[derive(Debug, PartialEq)] pub enum Decision { /// Emit this sync now and update the slot's last-emit timestamp. EmitNow, /// Buffer the percent and let an already-scheduled trailing task handle it. Buffered, /// Buffer the percent and schedule a trailing emit at the returned wake time. /// The caller is responsible for spawning the wake task and calling /// [`Throttle::take_pending`] when it fires. ScheduleTrailing { fire_at: Instant }, } struct Slot { last_emit_at: Instant, pending: Option, trailing_scheduled: bool, } pub struct Throttle { window: Duration, slots: Mutex>, } impl Throttle { pub fn new(window: Duration) -> Self { Self { window, slots: Mutex::new(HashMap::new()) } } /// Decide what to do with an incoming sync of `percent` at `now`. /// /// - If no prior emission exists for `key`, or the prior emission was /// longer than `window` ago: `EmitNow` (and the slot is updated as if /// the emit happened at `now`). /// - If a prior emission is within `window` and a trailing task is /// already scheduled: `Buffered` (the existing trailing task will /// pick up the new `percent`). /// - If a prior emission is within `window` and no trailing task is /// scheduled: `ScheduleTrailing { fire_at }` and the slot is marked /// as having a trailing task scheduled. pub fn decide(&self, key: &DocKey, percent: f64, now: Instant) -> Decision { let mut slots = self.slots.lock().expect("throttle mutex poisoned"); match slots.get_mut(key) { None => { slots.insert( key.clone(), Slot { last_emit_at: now, pending: None, trailing_scheduled: false }, ); Decision::EmitNow } Some(slot) => { let elapsed = now.saturating_duration_since(slot.last_emit_at); if elapsed >= self.window { slot.last_emit_at = now; slot.pending = None; Decision::EmitNow } else if slot.trailing_scheduled { slot.pending = Some(percent); Decision::Buffered } else { slot.pending = Some(percent); slot.trailing_scheduled = true; Decision::ScheduleTrailing { fire_at: slot.last_emit_at + self.window } } } } } /// Called when a trailing-emit task fires. Returns the buffered percent /// (if any) and marks the slot as having emitted at `now`. pub fn take_pending(&self, key: &DocKey, now: Instant) -> Option { let mut slots = self.slots.lock().expect("throttle mutex poisoned"); let slot = slots.get_mut(key)?; slot.trailing_scheduled = false; let pending = slot.pending.take(); if pending.is_some() { slot.last_emit_at = now; } pending } } /// Spawn a tokio task that sleeps until `fire_at`, then calls `on_fire` with /// the percent the throttle had buffered when the wake fires. If the buffered /// slot is empty by the time the task wakes (rare race), nothing happens. pub fn spawn_trailing_emit( throttle: Arc, key: DocKey, fire_at: Instant, on_fire: F, ) -> JoinHandle<()> where F: FnOnce(f64) -> Fut + Send + 'static, Fut: std::future::Future + Send + 'static, { tokio::spawn(async move { let now = Instant::now(); if fire_at > now { tokio::time::sleep(fire_at - now).await; } if let Some(percent) = throttle.take_pending(&key, Instant::now()) { on_fire(percent).await; } }) } #[cfg(test)] mod tests { use super::*; fn key() -> DocKey { (1, "doc-a".to_owned()) } #[test] fn first_sync_emits_immediately() { let t = Throttle::new(Duration::from_secs(300)); let now = Instant::now(); assert_eq!(t.decide(&key(), 10.0, now), Decision::EmitNow); } #[test] fn second_sync_within_window_schedules_trailing() { let t = Throttle::new(Duration::from_secs(300)); let now = Instant::now(); let _ = t.decide(&key(), 10.0, now); match t.decide(&key(), 20.0, now + Duration::from_secs(30)) { Decision::ScheduleTrailing { fire_at } => { assert_eq!(fire_at, now + Duration::from_secs(300)); } other => panic!("expected ScheduleTrailing, got {other:?}"), } } #[test] fn third_sync_with_trailing_scheduled_is_buffered() { let t = Throttle::new(Duration::from_secs(300)); let now = Instant::now(); let _ = t.decide(&key(), 10.0, now); let _ = t.decide(&key(), 20.0, now + Duration::from_secs(30)); assert_eq!( t.decide(&key(), 30.0, now + Duration::from_secs(60)), Decision::Buffered ); } #[test] fn take_pending_returns_latest_buffered_percent() { let t = Throttle::new(Duration::from_secs(300)); let now = Instant::now(); let _ = t.decide(&key(), 10.0, now); let _ = t.decide(&key(), 20.0, now + Duration::from_secs(30)); let _ = t.decide(&key(), 30.0, now + Duration::from_secs(60)); // The latest buffered value is 30.0. assert_eq!( t.take_pending(&key(), now + Duration::from_secs(300)), Some(30.0) ); } #[test] fn sync_after_window_emits_immediately_again() { let t = Throttle::new(Duration::from_secs(300)); let now = Instant::now(); let _ = t.decide(&key(), 10.0, now); assert_eq!( t.decide(&key(), 50.0, now + Duration::from_secs(301)), Decision::EmitNow ); } #[test] fn sync_after_trailing_fired_emits_immediately() { let t = Throttle::new(Duration::from_secs(300)); let now = Instant::now(); let _ = t.decide(&key(), 10.0, now); let _ = t.decide(&key(), 20.0, now + Duration::from_secs(30)); let trailing_time = now + Duration::from_secs(300); assert_eq!(t.take_pending(&key(), trailing_time), Some(20.0)); // After the trailing fires, last_emit_at = trailing_time. A new sync // within the next window should schedule another trailing. match t.decide(&key(), 40.0, trailing_time + Duration::from_secs(10)) { Decision::ScheduleTrailing { .. } => {} other => panic!("expected ScheduleTrailing, got {other:?}"), } } #[test] fn take_pending_with_no_buffer_is_noop() { let t = Throttle::new(Duration::from_secs(300)); let now = Instant::now(); let _ = t.decide(&key(), 10.0, now); // No second decide call → nothing buffered. assert_eq!(t.take_pending(&key(), now + Duration::from_secs(300)), None); } #[test] fn distinct_keys_do_not_interfere() { let t = Throttle::new(Duration::from_secs(300)); let now = Instant::now(); let k1: DocKey = (1, "a".to_owned()); let k2: DocKey = (2, "b".to_owned()); assert_eq!(t.decide(&k1, 10.0, now), Decision::EmitNow); assert_eq!(t.decide(&k2, 50.0, now), Decision::EmitNow); } }