use async_channel::Receiver;
use tracing::debug;
#[derive(Debug)]
pub struct Terminator<T> {
receiver: Receiver<T>,
}
impl<T> Terminator<T> {
pub fn new(receiver: Receiver<T>) -> Self {
Self { receiver }
}
pub async fn terminate(&self) {
debug!("terminator has started.");
while self.receiver.recv().await.is_ok() {}
debug!("terminator has been completed.");
}
}
#[cfg(test)]
mod tests {
use super::*;
#[tokio::test]
async fn terminates_when_channel_closed() {
let (sender, receiver) = async_channel::bounded::<String>(10);
sender.send("a".to_string()).await.unwrap();
sender.send("b".to_string()).await.unwrap();
sender.close();
let terminator = Terminator::new(receiver);
terminator.terminate().await;
}
#[tokio::test]
async fn terminates_empty_channel() {
let (sender, receiver) = async_channel::bounded::<u32>(10);
sender.close();
let terminator = Terminator::new(receiver);
terminator.terminate().await;
}
#[tokio::test]
async fn terminates_after_sender_dropped() {
let (sender, receiver) = async_channel::bounded::<i32>(10);
sender.send(42).await.unwrap();
drop(sender);
let terminator = Terminator::new(receiver);
terminator.terminate().await;
}
}