use tokio::sync::mpsc;
pub struct AsyncReceiver<T> {
rx: mpsc::Receiver<T>,
dropped_count: Option<std::sync::Arc<std::sync::atomic::AtomicU64>>,
}
impl<T> AsyncReceiver<T> {
pub(crate) fn new(rx: mpsc::Receiver<T>) -> Self {
Self {
rx,
dropped_count: None,
}
}
pub(crate) fn new_with_drop_counter(
rx: mpsc::Receiver<T>,
dropped_count: std::sync::Arc<std::sync::atomic::AtomicU64>,
) -> Self {
Self {
rx,
dropped_count: Some(dropped_count),
}
}
#[must_use]
pub fn dropped_message_count(&self) -> u64 {
self.dropped_count
.as_ref()
.map_or(0, |c| c.load(std::sync::atomic::Ordering::Relaxed))
}
#[must_use]
#[deprecated(
since = "0.2.0",
note = "单位为「条」而非「字节」,易与字节计数混用;请改用 dropped_message_count"
)]
pub fn dropped_count(&self) -> u64 {
self.dropped_message_count()
}
pub async fn recv(&mut self) -> Option<T> {
self.rx.recv().await
}
pub fn try_recv(&mut self) -> Result<T, AsyncRecvError> {
self.rx.try_recv().map_err(|e| match e {
mpsc::error::TryRecvError::Empty => AsyncRecvError::Empty,
mpsc::error::TryRecvError::Disconnected => AsyncRecvError::Disconnected,
})
}
pub fn close(&mut self) {
self.rx.close();
}
}
impl<T> std::fmt::Debug for AsyncReceiver<T> {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("AsyncReceiver").finish_non_exhaustive()
}
}
#[derive(thiserror::Error, Debug, Clone, Copy, PartialEq, Eq)]
pub enum AsyncRecvError {
#[error("异步通道当前无可用数据")]
Empty,
#[error("异步通道已关闭且剩余数据已耗尽")]
Disconnected,
}