use std::fmt;
use std::fmt::Debug;
use std::hash::Hash;
use std::hash::Hasher;
use futures::sync::mpsc::*;
use network_messages::Message;
use utils::unique_id::UniqueId;
use crate::connection::close_type::CloseType;
use crate::websocket::Message as WebSocketMessage;
use crate::connection::network_connection::ClosedFlag;
#[derive(Clone)]
pub struct PeerSink {
sink: UnboundedSender<WebSocketMessage>,
unique_id: UniqueId,
closed_flag: ClosedFlag,
}
impl PeerSink {
pub fn new(channel: UnboundedSender<WebSocketMessage>, unique_id: UniqueId, closed_flag: ClosedFlag) -> Self {
PeerSink {
sink: channel,
unique_id,
closed_flag,
}
}
pub fn send(&self, msg: Message) -> Result<(), SendError<WebSocketMessage>> {
if self.closed_flag.is_closed() {
return Ok(());
}
self.sink.unbounded_send(WebSocketMessage::Message(msg))
}
pub fn close(&self, ty: CloseType, reason: Option<String>) {
if self.closed_flag.set_closed(true) {
return;
}
self.closed_flag.set_close_type(ty);
debug!("Closing connection, reason: {:?} ({:?})", ty, reason);
if let Err(error) = self.sink.unbounded_send(WebSocketMessage::Close(None)) {
debug!("Error closing connection: {}", error);
}
}
}
impl Debug for PeerSink {
fn fmt(&self, f: &mut fmt::Formatter) -> Result<(), fmt::Error> {
write!(f, "PeerSink {{}}")
}
}
impl PartialEq for PeerSink {
fn eq(&self, other: &PeerSink) -> bool {
self.unique_id == other.unique_id
}
}
impl Eq for PeerSink {}
impl Hash for PeerSink {
fn hash<H: Hasher>(&self, state: &mut H) {
self.unique_id.hash(state);
}
}