use std::time::{Duration, Instant};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
pub enum Directive {
ArmQuietDeadline(Instant),
StartBuild(BuildGeneration),
DiscardStale(BuildGeneration),
PromoteCandidate(BuildGeneration),
None,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq, PartialOrd, Ord)]
pub struct BuildGeneration(pub u64);
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum State {
Idle,
Burst { deadline: Instant },
Building { generation: BuildGeneration },
BuildingDirty {
generation: BuildGeneration,
deadline: Instant,
ripe: bool,
},
}
#[derive(Debug)]
pub struct Coalescer {
quiet: Duration,
state: State,
next_generation: u64,
}
impl Coalescer {
#[must_use]
pub fn new(quiet: Duration) -> Self {
Self {
quiet,
state: State::Idle,
next_generation: 1,
}
}
pub fn relevant_event(&mut self, now: Instant) -> Directive {
let deadline = now + self.quiet;
match self.state {
State::Idle | State::Burst { .. } => {
self.state = State::Burst { deadline };
Directive::ArmQuietDeadline(deadline)
}
State::Building { generation } | State::BuildingDirty { generation, .. } => {
self.state = State::BuildingDirty {
generation,
deadline,
ripe: false,
};
Directive::ArmQuietDeadline(deadline)
}
}
}
pub fn quiet_deadline_elapsed(&mut self, now: Instant) -> Directive {
match self.state {
State::Idle | State::Building { .. } => Directive::None,
State::Burst { deadline } => {
if now < deadline {
return Directive::None;
}
let generation = self.mint_generation();
self.state = State::Building { generation };
Directive::StartBuild(generation)
}
State::BuildingDirty {
generation,
deadline,
ripe,
} => {
if now < deadline {
return Directive::None;
}
let _ = ripe;
self.state = State::BuildingDirty {
generation,
deadline,
ripe: true,
};
Directive::None
}
}
}
pub fn build_completed(&mut self, completed: BuildGeneration, now: Instant) -> [Directive; 2] {
match self.state {
State::Building { generation } if generation == completed => {
self.state = State::Idle;
[Directive::PromoteCandidate(completed), Directive::None]
}
State::BuildingDirty {
generation,
deadline,
ripe,
} if generation == completed => {
if ripe || now >= deadline {
let follow_up = self.mint_generation();
self.state = State::Building {
generation: follow_up,
};
[
Directive::DiscardStale(completed),
Directive::StartBuild(follow_up),
]
} else {
self.state = State::Burst { deadline };
[Directive::DiscardStale(completed), Directive::None]
}
}
State::Idle
| State::Burst { .. }
| State::Building { .. }
| State::BuildingDirty { .. } => [Directive::DiscardStale(completed), Directive::None],
}
}
pub fn watcher_desynchronized(&mut self, now: Instant) -> Directive {
match self.state {
State::Idle | State::Burst { .. } => {
let generation = self.mint_generation();
self.state = State::Building { generation };
Directive::StartBuild(generation)
}
State::Building { generation } | State::BuildingDirty { generation, .. } => {
self.state = State::BuildingDirty {
generation,
deadline: now,
ripe: true,
};
Directive::None
}
}
}
fn mint_generation(&mut self) -> BuildGeneration {
let generation = BuildGeneration(self.next_generation);
self.next_generation += 1;
generation
}
}
#[cfg(test)]
mod tests {
use super::{BuildGeneration, Coalescer, Directive};
use std::time::{Duration, Instant};
const QUIET: Duration = Duration::from_millis(100);
fn machine() -> (Coalescer, Instant) {
(Coalescer::new(QUIET), Instant::now())
}
#[test]
fn one_event_one_deadline_one_build() {
let (mut m, t0) = machine();
assert_eq!(
m.relevant_event(t0),
Directive::ArmQuietDeadline(t0 + QUIET)
);
assert_eq!(
m.quiet_deadline_elapsed(t0 + QUIET),
Directive::StartBuild(BuildGeneration(1))
);
assert_eq!(
m.build_completed(BuildGeneration(1), t0 + QUIET * 2),
[
Directive::PromoteCandidate(BuildGeneration(1)),
Directive::None
]
);
}
#[test]
fn rapid_burst_coalesces_to_one_build() {
let (mut m, t0) = machine();
for i in 0..10 {
let at = t0 + Duration::from_millis(i * 10);
assert_eq!(
m.relevant_event(at),
Directive::ArmQuietDeadline(at + QUIET),
"every event re-arms the same single deadline"
);
}
let last_event = t0 + Duration::from_millis(90);
assert_eq!(m.quiet_deadline_elapsed(t0 + QUIET), Directive::None);
assert_eq!(
m.quiet_deadline_elapsed(last_event + QUIET),
Directive::StartBuild(BuildGeneration(1))
);
assert_eq!(
m.build_completed(BuildGeneration(1), last_event + QUIET * 2),
[
Directive::PromoteCandidate(BuildGeneration(1)),
Directive::None
]
);
}
#[test]
fn edit_during_build_discards_stale_and_rebuilds_once() {
let (mut m, t0) = machine();
m.relevant_event(t0);
assert_eq!(
m.quiet_deadline_elapsed(t0 + QUIET),
Directive::StartBuild(BuildGeneration(1))
);
let edit = t0 + QUIET + Duration::from_millis(20);
assert_eq!(
m.relevant_event(edit),
Directive::ArmQuietDeadline(edit + QUIET)
);
assert_eq!(m.quiet_deadline_elapsed(edit + QUIET), Directive::None);
assert_eq!(
m.build_completed(BuildGeneration(1), edit + QUIET * 2),
[
Directive::DiscardStale(BuildGeneration(1)),
Directive::StartBuild(BuildGeneration(2))
]
);
assert_eq!(
m.build_completed(BuildGeneration(2), edit + QUIET * 3),
[
Directive::PromoteCandidate(BuildGeneration(2)),
Directive::None
]
);
}
#[test]
fn open_burst_at_completion_waits_for_quiet() {
let (mut m, t0) = machine();
m.relevant_event(t0);
m.quiet_deadline_elapsed(t0 + QUIET);
let edit = t0 + QUIET + Duration::from_millis(20);
m.relevant_event(edit);
let completion = edit + Duration::from_millis(10);
assert_eq!(
m.build_completed(BuildGeneration(1), completion),
[Directive::DiscardStale(BuildGeneration(1)), Directive::None]
);
assert_eq!(
m.quiet_deadline_elapsed(edit + QUIET),
Directive::StartBuild(BuildGeneration(2))
);
}
#[test]
fn idle_arms_nothing() {
let (mut m, t0) = machine();
assert_eq!(m.quiet_deadline_elapsed(t0 + QUIET), Directive::None);
assert_eq!(m.quiet_deadline_elapsed(t0 + QUIET * 100), Directive::None);
}
#[test]
fn same_sequence_same_generations() {
let script = |m: &mut Coalescer, t0: Instant| {
let mut log = vec![
m.relevant_event(t0),
m.quiet_deadline_elapsed(t0 + QUIET),
m.relevant_event(t0 + QUIET + Duration::from_millis(5)),
m.quiet_deadline_elapsed(t0 + QUIET * 2 + Duration::from_millis(5)),
];
for d in m.build_completed(BuildGeneration(1), t0 + QUIET * 3) {
log.push(d);
}
for d in m.build_completed(BuildGeneration(2), t0 + QUIET * 4) {
log.push(d);
}
log
};
let t0 = Instant::now();
let (mut a, mut b) = (Coalescer::new(QUIET), Coalescer::new(QUIET));
assert_eq!(script(&mut a, t0), script(&mut b, t0));
}
#[test]
fn desync_rebuilds_immediately_when_not_building() {
let (mut m, t0) = machine();
assert_eq!(
m.watcher_desynchronized(t0),
Directive::StartBuild(BuildGeneration(1))
);
}
#[test]
fn desync_during_build_rebuilds_at_completion() {
let (mut m, t0) = machine();
m.relevant_event(t0);
m.quiet_deadline_elapsed(t0 + QUIET);
assert_eq!(
m.watcher_desynchronized(t0 + QUIET + Duration::from_millis(1)),
Directive::None
);
assert_eq!(
m.build_completed(BuildGeneration(1), t0 + QUIET * 2),
[
Directive::DiscardStale(BuildGeneration(1)),
Directive::StartBuild(BuildGeneration(2))
]
);
}
#[test]
fn unknown_generation_completion_never_promotes() {
let (mut m, t0) = machine();
assert_eq!(
m.build_completed(BuildGeneration(7), t0),
[Directive::DiscardStale(BuildGeneration(7)), Directive::None]
);
}
}