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(n) => {
let message_type = if n > 0 {
MessageType::Data
} else {
MessageType::EOF
};
let header = MessageHeader::new(user_id, message_type, 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;
}
};
if n == 0 {
break;
}
}
Err(e) => {
error!("[{}][{}] Reading from User-Con: {}", client_id, user_id, e);
break;
}
}
}
close_user.await;
}