use std::time::Duration;
use tokio::sync::mpsc;
use tokio::time;
use crate::watcher::ChangeEvent;
#[non_exhaustive]
#[derive(Debug, Clone, Copy)]
pub struct ScheduleConfig {
pub quiet: Duration,
pub cap: Duration,
}
impl Default for ScheduleConfig {
fn default() -> Self {
Self {
quiet: Duration::from_millis(100),
cap: Duration::from_millis(500),
}
}
}
#[non_exhaustive]
#[derive(Debug, Clone)]
pub struct Trigger {
pub events: Vec<ChangeEvent>,
}
#[must_use]
pub fn spawn(
mut events_rx: mpsc::UnboundedReceiver<ChangeEvent>,
config: ScheduleConfig,
) -> mpsc::UnboundedReceiver<Trigger> {
let (tx, rx) = mpsc::unbounded_channel();
tokio::spawn(async move {
let mut buffered: Vec<ChangeEvent> = Vec::new();
let mut first_event_at: Option<time::Instant> = None;
loop {
let next_event = if buffered.is_empty() {
events_rx.recv().await
} else {
let deadline = compute_deadline(first_event_at, config);
let Ok(event) = time::timeout_at(deadline, events_rx.recv()).await else {
let trigger = Trigger {
events: std::mem::take(&mut buffered),
};
first_event_at = None;
if tx.send(trigger).is_err() {
return;
}
continue;
};
event
};
let Some(event) = next_event else { return };
if buffered.is_empty() {
first_event_at = Some(time::Instant::now());
}
buffered.push(event);
}
});
rx
}
fn compute_deadline(first: Option<time::Instant>, config: ScheduleConfig) -> time::Instant {
let now = time::Instant::now();
let quiet_deadline = now + config.quiet;
match first {
Some(start) => {
let cap_deadline = start + config.cap;
quiet_deadline.min(cap_deadline)
}
None => quiet_deadline,
}
}