use crate::message::{Message, MessageHeader, MessageType};
use log::error;
use tokio::io::AsyncReadExt;
const BUFFER_SIZE: usize = 4096;
pub async fn recv<F>(
client_id: u32,
user_id: u32,
mut con: tokio::net::tcp::OwnedReadHalf,
send_queue: tokio::sync::mpsc::UnboundedSender<Message>,
close_user: F,
) where
F: std::future::Future<Output = ()>,
{
loop {
let mut buf = vec![0; BUFFER_SIZE];
match con.read(&mut buf).await {
Ok(0) => {
break;
}
Ok(n) => {
let header = MessageHeader::new(user_id, MessageType::Data, n as u64);
let msg = Message::new(header, buf);
match send_queue.send(msg) {
Ok(_) => {}
Err(e) => {
error!(
"[{}][{}] Forwarding message to client: {}",
client_id, user_id, e
);
break;
}
};
}
Err(ref e) if e.kind() == std::io::ErrorKind::WouldBlock => {
continue;
}
Err(e) => {
error!("[{}][{}] Reading from User-Con: {}", client_id, user_id, e);
break;
}
}
}
close_user.await;
}