use bytes::Bytes;
use std::io;
use std::io::ErrorKind;
use std::net::SocketAddr;
use std::sync::atomic::{AtomicI64, Ordering};
use std::sync::Arc;
use tokio::net::UdpSocket;
use tokio::sync::mpsc::{unbounded_channel, UnboundedReceiver, UnboundedSender};
pub type UDPPeer = Arc<UdpPeer>;
pub type UdpSender = UnboundedSender<io::Result<Bytes>>;
pub type UdpReader = UnboundedReceiver<io::Result<Bytes>>;
pub struct UdpPeer {
pub socket_id: usize,
pub udp_sock: Arc<UdpSocket>,
pub addr: SocketAddr,
sender: UdpSender,
last_read_time: AtomicI64,
}
impl Drop for UdpPeer {
fn drop(&mut self) {
log::trace!(
"udp_listen socket:{} udp peer:{} drop",
self.socket_id,
self.addr
)
}
}
impl UdpPeer {
#[inline]
pub fn new(
socket_id: usize,
udp_sock: Arc<UdpSocket>,
addr: SocketAddr,
) -> (UDPPeer, UdpReader) {
let (tx, rx) = unbounded_channel();
(
Arc::new(Self {
socket_id,
udp_sock,
addr,
sender: tx,
last_read_time: AtomicI64::new(timestamp_sec()),
}),
rx,
)
}
#[inline]
pub(crate) fn get_last_recv_sec(&self) -> i64 {
self.last_read_time.load(Ordering::Acquire)
}
#[inline]
pub(crate) fn push_data(&self, data: Bytes) -> io::Result<()> {
if let Err(err) = self.sender.send(Ok(data)) {
Err(io::Error::other(err))
} else {
Ok(())
}
}
#[inline]
pub(crate) async fn push_data_and_update_instant(&self, data: Bytes) -> io::Result<()> {
self.last_read_time
.store(timestamp_sec(), Ordering::Release);
self.push_data(data)
}
#[inline]
pub fn get_socket_id(&self) -> usize {
self.socket_id
}
#[inline]
pub fn get_addr(&self) -> SocketAddr {
self.addr
}
#[inline]
pub async fn send(&self, buf: &[u8]) -> io::Result<usize> {
self.udp_sock.send_to(buf, &self.addr).await
}
#[inline]
pub fn close(&self) {
if let Err(err) = self.sender.send(Err(io::Error::new(
ErrorKind::TimedOut,
"udp peer need close",
))) {
log::error!("send timeout to udp peer:{} error:{err}", self.get_addr());
}
}
}
#[inline]
fn timestamp_sec() -> i64 {
chrono::Utc::now().timestamp()
}