use crossbeam_channel::{bounded, select, tick, Sender};
use std::thread::JoinHandle;
use std::time::Duration;
pub struct PeriodicClosure {
join_handle: Option<JoinHandle<()>>,
canceller: Sender<()>,
}
impl PeriodicClosure {
pub fn start<F: FnMut() + Send + Sync + 'static>(
name: String,
period: Duration,
mut func: F,
) -> Self {
let ticker = tick(period);
let (canceller, cancelled) = bounded(1);
let join_handle = std::thread::Builder::new()
.name(name)
.spawn(move || {
loop {
if cancelled.try_recv().is_ok() {
break;
}
select! {
recv(cancelled) -> _ => break,
recv(ticker) -> _ => func(),
}
}
})
.unwrap();
Self {
join_handle: Some(join_handle),
canceller,
}
}
}
impl Drop for PeriodicClosure {
fn drop(&mut self) {
if let Some(join_handle) = self.join_handle.take() {
self.canceller
.send(())
.expect("The receiver must exists in detached thread.");
join_handle.join().unwrap();
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::sync::Arc;
#[test]
fn test() {
let i = Arc::new(std::sync::atomic::AtomicI64::new(0));
let j = Arc::clone(&i);
let p = PeriodicClosure::start(
"test thread".to_string(),
Duration::from_millis(100),
move || {
j.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
},
);
std::thread::sleep(Duration::from_millis(1050));
std::mem::drop(p);
assert_eq!(i.load(std::sync::atomic::Ordering::Relaxed), 10);
}
#[test]
fn test_collapsed() {
let i = Arc::new(std::sync::atomic::AtomicI64::new(0));
let j = Arc::clone(&i);
let p = PeriodicClosure::start(
"test thread".to_string(),
Duration::from_millis(1),
move || {
j.fetch_add(1, std::sync::atomic::Ordering::Relaxed);
std::thread::sleep(Duration::from_millis(100));
},
);
std::thread::sleep(Duration::from_millis(950));
std::mem::drop(p);
assert_eq!(i.load(std::sync::atomic::Ordering::Relaxed), 10);
}
}