use crate::connections::Connections;
use crate::message::{Message, MessageHeader, MessageType};
use crate::streams::mpsc;
use log::{debug, error};
pub struct Sender {
id: u32,
tx: tokio::sync::mpsc::UnboundedSender<Message>,
total_client_cons: std::sync::Arc<Connections<mpsc::StreamWriter<Message>>>,
}
impl Sender {
pub fn new(
id: u32,
tx: tokio::sync::mpsc::UnboundedSender<Message>,
cons: std::sync::Arc<Connections<mpsc::StreamWriter<Message>>>,
) -> Self {
Self {
id,
tx,
total_client_cons: cons,
}
}
pub fn send(&self, data: Vec<u8>, length: u64) -> bool {
let header = MessageHeader::new(self.id, MessageType::Data, length);
let msg = Message::new(header, data);
self.tx.send(msg).is_ok()
}
}
impl Drop for Sender {
fn drop(&mut self) {
self.total_client_cons.remove(self.id);
debug!("[Sender][{}] Removed Connection", self.id);
let close_msg = Message::new(MessageHeader::new(self.id, MessageType::Close, 0), vec![]);
match self.tx.send(close_msg) {
Ok(_) => {
debug!("[Sender][{}] Sent Close", self.id);
}
Err(e) => {
error!("Sending Close-Message for {}: {}", self.id, e);
}
}
}
}