use tokio::sync::watch;
#[derive(Clone)]
pub(crate) struct CancelHandle(watch::Sender<bool>);
#[derive(Clone)]
pub(crate) struct CancelSignal(watch::Receiver<bool>);
impl CancelHandle {
pub(crate) fn new() -> (Self, CancelSignal) {
let (sender, receiver) = watch::channel(false);
(Self(sender), CancelSignal(receiver))
}
pub(crate) fn cancel(&self) {
self.0.send_replace(true);
}
}
impl CancelSignal {
pub(crate) async fn cancelled(&mut self) {
if *self.0.borrow() {
return;
}
while self.0.changed().await.is_ok() {
if *self.0.borrow() {
return;
}
}
}
pub(crate) fn is_cancelled(&self) -> bool {
*self.0.borrow()
}
}
#[cfg(test)]
mod tests {
use super::CancelHandle;
#[test]
fn cancellation_is_reusable() {
let runtime = tokio::runtime::Builder::new_current_thread()
.enable_all()
.build()
.unwrap();
runtime.block_on(async {
let (handle, mut signal) = CancelHandle::new();
let mut other = signal.clone();
let waiter = tokio::spawn(async move { other.cancelled().await });
handle.cancel();
waiter.await.unwrap();
signal.cancelled().await;
assert!(signal.is_cancelled());
});
}
}