use std::sync::atomic::{AtomicBool, Ordering}; use std::time::{Duration, Instant}; use parking_lot::{Condvar, Mutex}; /// tells hydrant's own threads to wind down. they sleep on it instead of /// `thread::sleep`, so a shutdown doesn't wait out a whole interval. #[derive(Default)] pub struct Stop { // the ingest shards check this per message, so it stays a plain load stopped: AtomicBool, lock: Mutex<()>, changed: Condvar, } impl Stop { pub fn trigger(&self) { // set under the lock, or a waiter that just checked could miss the notify let _lock = self.lock.lock(); self.stopped.store(true, Ordering::Release); self.changed.notify_all(); } pub fn is_set(&self) -> bool { self.stopped.load(Ordering::Acquire) } /// sleeps for `timeout`, or less if stopped meanwhile. true means stop. pub fn wait(&self, timeout: Duration) -> bool { let deadline = Instant::now() + timeout; let mut lock = self.lock.lock(); while !self.is_set() { if self.changed.wait_until(&mut lock, deadline).timed_out() { break; } } self.is_set() } } #[cfg(test)] mod tests { use super::*; use std::sync::Arc; #[test] fn wait_times_out_without_a_trigger() { let stop = Stop::default(); let started = Instant::now(); assert!(!stop.wait(Duration::from_millis(20))); assert!(started.elapsed() >= Duration::from_millis(20)); } #[test] fn trigger_wakes_a_waiting_thread_early() { let stop = Arc::new(Stop::default()); let waiter = std::thread::spawn({ let stop = stop.clone(); move || { let started = Instant::now(); (stop.wait(Duration::from_secs(60)), started.elapsed()) } }); std::thread::sleep(Duration::from_millis(20)); stop.trigger(); let (stopped, waited) = waiter.join().unwrap(); assert!(stopped); assert!(waited < Duration::from_secs(5)); // and stays set for anyone who checks later assert!(stop.wait(Duration::from_secs(60))); } }