use std::future::pending;
use tokio::sync::watch;
#[derive(Clone, Debug)]
pub struct GracefulShutdownSignal {
rx: watch::Receiver<()>,
}
impl GracefulShutdownSignal {
pub(crate) fn new(rx: watch::Receiver<()>) -> Self {
Self { rx }
}
pub async fn notified(&self) {
let mut rx = self.rx.clone();
if rx.changed().await.is_err() {
pending::<()>().await;
}
}
}
#[cfg(test)]
mod tests {
use std::time::Duration;
use actix_rt::time::timeout;
use super::GracefulShutdownSignal;
#[actix_rt::test]
async fn set_signal_notifies_listener() {
let (tx, rx) = tokio::sync::watch::channel(());
let signal = GracefulShutdownSignal::new(rx);
timeout(Duration::from_millis(100), signal.notified())
.await
.expect_err("signal notified listener before shutdown");
tx.send_replace(());
timeout(Duration::from_millis(100), signal.notified())
.await
.expect("set signal did not notify listener");
timeout(Duration::from_millis(100), signal.notified())
.await
.expect("set signal did not notify later listener");
}
#[actix_rt::test]
async fn closed_unset_signal_does_not_notify_listener() {
let (tx, rx) = tokio::sync::watch::channel(());
let signal = GracefulShutdownSignal::new(rx);
drop(tx);
assert!(timeout(Duration::from_millis(10), signal.notified())
.await
.is_err());
}
}