#![allow(non_snake_case)]
use crate::{
dht::{dht_protocol::*, dht_trait::Dht},
engine::{p2p_protocol::P2pProtocol, RealEngine},
transport::{protocol::*, transport_trait::Transport, ConnectionIdRef},
};
use lib3h_crypto_api::{Buffer, CryptoSystem};
use lib3h_protocol::{data_types::*, protocol_server::Lib3hServerProtocol, DidWork, Lib3hResult};
use rmp_serde::{Deserializer, Serializer};
use serde::{Deserialize, Serialize};
impl<T: Transport, D: Dht, SecBuf: Buffer, Crypto: CryptoSystem> RealEngine<T, D, SecBuf, Crypto> {
pub(crate) fn process_network_gateway(
&mut self,
) -> Lib3hResult<(DidWork, Vec<Lib3hServerProtocol>)> {
let mut outbox = Vec::new();
let (tranport_did_work, event_list) =
Transport::process(&mut *self.network_gateway.borrow_mut())?;
debug!(
"{} - network_gateway Transport.process(): {} {}",
self.name.clone(),
tranport_did_work,
event_list.len()
);
if tranport_did_work {
for evt in event_list {
let mut output = self.handle_netTransportEvent(&evt)?;
outbox.append(&mut output);
}
}
let (dht_did_work, event_list) = Dht::process(&mut *self.network_gateway.borrow_mut())?;
if dht_did_work {
for evt in event_list {
let mut output = self.handle_netDhtEvent(evt)?;
outbox.append(&mut output);
}
}
Ok((tranport_did_work || dht_did_work, outbox))
}
fn handle_netDhtEvent(&mut self, cmd: DhtEvent) -> Lib3hResult<Vec<Lib3hServerProtocol>> {
debug!("{} << handle_netDhtEvent: {:?}", self.name.clone(), cmd);
let outbox = Vec::new();
match cmd {
DhtEvent::GossipTo(_data) => {
}
DhtEvent::GossipUnreliablyTo(_data) => {
}
DhtEvent::HoldPeerRequested(_peer_address) => {
}
DhtEvent::PeerTimedOut(_data) => {
}
DhtEvent::HoldEntryRequested(_from, _entry) => {
}
DhtEvent::FetchEntryResponse(_data) => {
}
DhtEvent::EntryPruned(_address) => {
}
}
Ok(outbox)
}
fn handle_netTransportEvent(
&mut self,
evt: &TransportEvent,
) -> Lib3hResult<Vec<Lib3hServerProtocol>> {
debug!(
"{} << handle_netTransportEvent: {:?}",
self.name.clone(),
evt
);
let mut outbox = Vec::new();
match evt {
TransportEvent::TransportError(id, e) => {
self.network_connections.remove(id);
error!("{} Network error from {} : {:?}", self.name.clone(), id, e);
if self.network_connections.is_empty() {
let data = DisconnectedData {
network_id: "FIXME".to_string(),
};
outbox.push(Lib3hServerProtocol::Disconnected(data));
}
}
TransportEvent::ConnectResult(id) => {
let mut network_gateway = self.network_gateway.borrow_mut();
if let Some(uri) = network_gateway.get_uri(id) {
info!("Network Connection opened: {} ({})", id, uri);
let space_list = self.get_all_spaces();
let our_joined_space_list = P2pProtocol::AllJoinedSpaceList(space_list);
let mut buf = Vec::new();
our_joined_space_list
.serialize(&mut Serializer::new(&mut buf))
.unwrap();
trace!(
"(GatewayTransport) P2pProtocol::AllJoinedSpaceList: {:?} to {:?}",
our_joined_space_list,
id
);
let peer_list = network_gateway.get_peer_list();
trace!(
"(GatewayTransport) P2pProtocol::AllJoinedSpaceList: get_peer_list = {:?}",
peer_list
);
let maybe_peer_data = peer_list.iter().find(|pd| pd.peer_uri == uri);
if let Some(peer_data) = maybe_peer_data {
trace!(
"(GatewayTransport) P2pProtocol::AllJoinedSpaceList ; sending back to {:?}",
peer_data,
);
network_gateway.send(&[&peer_data.peer_address], &buf)?;
}
if self.network_connections.is_empty() {
let data = ConnectedData {
request_id: "FIXME".to_string(),
uri,
};
outbox.push(Lib3hServerProtocol::Connected(data));
}
let _ = self.network_connections.insert(id.to_owned());
}
}
TransportEvent::Connection(_id) => {
unimplemented!();
}
TransportEvent::Closed(id) => {
self.network_connections.remove(id);
if self.network_connections.is_empty() {
let data = DisconnectedData {
network_id: "FIXME".to_string(),
};
outbox.push(Lib3hServerProtocol::Disconnected(data));
}
}
TransportEvent::Received(id, payload) => {
debug!("Received message from: {} | {}", id, payload.len());
let mut de = Deserializer::new(&payload[..]);
let maybe_msg: Result<P2pProtocol, rmp_serde::decode::Error> =
Deserialize::deserialize(&mut de);
if let Err(e) = maybe_msg {
return Err(format_err!("Failed deserializing msg: {:?}", e));
}
let p2p_msg = maybe_msg.unwrap();
let mut output = self.serve_P2pProtocol(id, &p2p_msg)?;
outbox.append(&mut output);
}
};
Ok(outbox)
}
fn serve_P2pProtocol(
&mut self,
_from_id: &ConnectionIdRef,
p2p_msg: &P2pProtocol,
) -> Lib3hResult<Vec<Lib3hServerProtocol>> {
let mut outbox = Vec::new();
match p2p_msg {
P2pProtocol::Gossip(msg) => {
let space_gateway = self
.space_gateway_map
.get_mut(&(msg.space_address.to_owned(), msg.to_peer_address.to_owned()))
.ok_or_else(|| format_err!("space_gateway not found"))?;
let from_peer_address =
std::string::String::from_utf8_lossy(&msg.from_peer_address).into_owned();
let cmd = DhtCommand::HandleGossip(RemoteGossipBundleData {
from_peer_address,
bundle: msg.bundle.clone(),
});
space_gateway.post_dht(cmd)?;
}
P2pProtocol::DirectMessage(data) => {
let lib3_msg = Lib3hServerProtocol::HandleSendDirectMessage(data.clone());
outbox.push(lib3_msg);
}
P2pProtocol::DirectMessageResult(data) => {
let lib3_msg = Lib3hServerProtocol::SendDirectMessageResult(data.clone());
outbox.push(lib3_msg);
}
P2pProtocol::FetchData => {
}
P2pProtocol::FetchDataResponse => {
}
P2pProtocol::PeerAddress(_, _) => {
}
P2pProtocol::BroadcastJoinSpace(gateway_id, peer_data) => {
debug!("Received JoinSpace: {} {:?}", gateway_id, peer_data);
for (_, space_gateway) in self.space_gateway_map.iter_mut() {
space_gateway.post_dht(DhtCommand::HoldPeer(peer_data.clone()))?;
}
}
P2pProtocol::AllJoinedSpaceList(join_list) => {
debug!("Received AllJoinedSpaceList: {:?}", join_list);
for (space_address, peer_data) in join_list {
let maybe_space_gateway = self.get_first_space_mut(space_address);
if let Some(space_gateway) = maybe_space_gateway {
space_gateway.post_dht(DhtCommand::HoldPeer(peer_data.clone()))?;
}
}
}
};
Ok(outbox)
}
}