use std::sync::{
Mutex, PoisonError,
atomic::{AtomicBool, AtomicU64, Ordering},
};
pub struct RedrawGate<T> {
slot: Mutex<Option<T>>,
pending: AtomicBool,
generation: AtomicU64,
}
impl<T> RedrawGate<T> {
pub const fn new() -> Self {
Self {
slot: Mutex::new(None),
pending: AtomicBool::new(false),
generation: AtomicU64::new(0),
}
}
pub fn submit(&self, value: T) -> Option<u64> {
*self.slot.lock().unwrap_or_else(PoisonError::into_inner) = Some(value);
let generation = self.generation.fetch_add(1, Ordering::AcqRel) + 1;
self.pending
.compare_exchange(false, true, Ordering::AcqRel, Ordering::Acquire)
.is_ok()
.then_some(generation)
}
pub fn take(&self) -> Option<T> {
self.pending.store(false, Ordering::Release);
self.slot
.lock()
.unwrap_or_else(PoisonError::into_inner)
.take()
}
pub fn peek<R>(&self, look: impl FnOnce(&T) -> R) -> Option<R> {
self.slot
.lock()
.unwrap_or_else(PoisonError::into_inner)
.as_ref()
.map(look)
}
#[cfg(test)]
pub fn current_generation(&self) -> u64 {
self.generation.load(Ordering::Acquire)
}
#[cfg(test)]
pub fn is_pending(&self) -> bool {
self.pending.load(Ordering::Acquire)
}
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn first_submit_must_enqueue() {
let gate: RedrawGate<u32> = RedrawGate::new();
assert_eq!(gate.submit(1), Some(1));
assert!(gate.is_pending());
}
#[test]
fn a_second_submit_before_take_must_not_enqueue_again() {
let gate: RedrawGate<u32> = RedrawGate::new();
assert!(gate.submit(1).is_some());
assert_eq!(
gate.submit(2),
None,
"a closure is already pending -- must not enqueue a second one"
);
assert_eq!(
gate.submit(3),
None,
"still pending -- must not enqueue a third one either"
);
}
#[test]
fn every_submit_bumps_generation_even_when_coalesced() {
let gate: RedrawGate<u32> = RedrawGate::new();
gate.submit(1);
gate.submit(2);
gate.submit(3);
assert_eq!(
gate.current_generation(),
3,
"generation must advance for every submit, not just the ones that enqueue"
);
}
#[test]
fn take_returns_the_latest_value_not_the_first_one() {
let gate: RedrawGate<u32> = RedrawGate::new();
assert!(gate.submit(1).is_some());
assert_eq!(gate.submit(2), None, "coalesced");
assert_eq!(gate.submit(3), None, "coalesced");
assert_eq!(
gate.take(),
Some(3),
"the one closure that does run must see the LATEST payload, not the one \
that originally caused it to be enqueued"
);
}
#[test]
fn take_clears_pending_so_a_later_submit_can_enqueue_again() {
let gate: RedrawGate<u32> = RedrawGate::new();
assert!(gate.submit(1).is_some());
assert!(gate.is_pending());
gate.take();
assert!(!gate.is_pending());
assert!(
gate.submit(2).is_some(),
"pending was cleared, a fresh closure may be queued"
);
}
#[test]
fn take_on_an_empty_gate_returns_none() {
let gate: RedrawGate<u32> = RedrawGate::new();
assert_eq!(gate.take(), None);
}
#[test]
fn peek_shows_the_queued_payload_and_leaves_it_there() {
let gate: RedrawGate<u32> = RedrawGate::new();
assert_eq!(gate.peek(|queued| *queued), None, "nothing queued yet");
gate.submit(1);
gate.submit(2);
assert_eq!(gate.peek(|queued| *queued), Some(2), "the latest one");
assert_eq!(
gate.peek(|queued| *queued),
Some(2),
"peeking takes nothing"
);
assert_eq!(gate.take(), Some(2));
assert_eq!(gate.peek(|queued| *queued), None, "taken");
}
#[test]
fn a_burst_of_n_submits_coalesces_to_exactly_one_pending_closure() {
let gate: RedrawGate<u32> = RedrawGate::new();
let mut enqueued = 0;
for i in 0..50 {
if gate.submit(i).is_some() {
enqueued += 1;
}
}
assert_eq!(enqueued, 1, "only the first of a burst may enqueue");
assert_eq!(
gate.take(),
Some(49),
"and it must see the last value submitted"
);
}
#[test]
fn a_second_burst_after_take_enqueues_exactly_once_more() {
let gate: RedrawGate<u32> = RedrawGate::new();
for i in 0..10 {
gate.submit(i);
}
gate.take();
let mut enqueued = 0;
for i in 10..20 {
if gate.submit(i).is_some() {
enqueued += 1;
}
}
assert_eq!(enqueued, 1);
assert_eq!(gate.take(), Some(19));
}
#[test]
fn works_as_a_pure_signal_with_a_unit_payload() {
let gate: RedrawGate<()> = RedrawGate::new();
assert!(gate.submit(()).is_some());
assert_eq!(gate.submit(()), None, "coalesced");
assert_eq!(gate.take(), Some(()));
assert!(gate.submit(()).is_some(), "pending was cleared");
}
}