#![allow(non_snake_case)]
use crate::{
dht::{dht_protocol::*, dht_trait::Dht},
engine::p2p_protocol::P2pProtocol,
gateway::P2pGateway,
transport::{
error::{TransportError, TransportResult},
protocol::{TransportCommand, TransportEvent},
transport_trait::Transport,
ConnectionId, ConnectionIdRef,
},
};
use lib3h_protocol::DidWork;
use rmp_serde::{Deserializer, Serializer};
use serde::{Deserialize, Serialize};
use url::Url;
fn transport_id_to_url(id: ConnectionId) -> Url {
Url::parse(id.as_str())
.expect("gateway_transport: transport_id_to_url: connection id is not a well formed url")
}
impl<T: Transport, D: Dht> Transport for P2pGateway<T, D> {
fn connect(&mut self, uri: &Url) -> TransportResult<ConnectionId> {
trace!("({}).connect() {}", self.identifier.clone(), uri);
let connection_id = self.inner_transport.borrow_mut().connect(&uri)?;
self.connection_map
.insert(uri.clone(), connection_id.clone());
Ok(connection_id)
}
fn close(&mut self, id: &ConnectionIdRef) -> TransportResult<()> {
self.inner_transport.borrow_mut().close(id)
}
fn close_all(&mut self) -> TransportResult<()> {
self.inner_transport.borrow_mut().close_all()
}
fn send(&mut self, dht_id_list: &[&ConnectionIdRef], payload: &[u8]) -> TransportResult<()> {
let dht_uri_list = self.connection_id_to_dht_uri_list(dht_id_list)?;
trace!(
"({}).send() {:?} -> {:?} | {}",
self.identifier.clone(),
dht_id_list,
dht_uri_list,
payload.len()
);
let mut net_uri_list = Vec::new();
for dht_uri in dht_uri_list {
let net_uri = self
.connection_map
.get(&dht_uri)
.expect("unknown dht_transport");
net_uri_list.push(net_uri);
trace!(
"({}).send() reversed mapped dht_uri {:?} to net_uri {:?}",
self.identifier.clone(),
dht_uri,
net_uri
)
}
let ref_list: Vec<&str> = net_uri_list.iter().map(|v| v.as_str()).collect();
self.inner_transport.borrow_mut().send(&ref_list, payload)
}
fn send_all(&mut self, payload: &[u8]) -> TransportResult<()> {
let connection_list = self.connection_id_list()?;
let dht_id_list: Vec<&str> = connection_list.iter().map(|v| &**v).collect();
trace!("({}) send_all() {:?}", self.identifier.clone(), dht_id_list);
self.send(&dht_id_list, payload)
}
fn bind(&mut self, url: &Url) -> TransportResult<Url> {
trace!("({}) bind() {}", self.identifier.clone(), url);
self.inner_transport.borrow_mut().bind(url)
}
fn post(&mut self, command: TransportCommand) -> TransportResult<()> {
self.transport_inbox.push_back(command);
Ok(())
}
fn process(&mut self) -> TransportResult<(DidWork, Vec<TransportEvent>)> {
let mut outbox = Vec::new();
let mut did_work = false;
loop {
let cmd = match self.transport_inbox.pop_front() {
None => break,
Some(msg) => msg,
};
let res = self.serve_TransportCommand(&cmd);
if let Ok(mut output) = res {
did_work = true;
outbox.append(&mut output);
}
}
trace!(
"({}).Transport.process() - output: {} {}",
self.identifier.clone(),
did_work,
outbox.len()
);
let (inner_did_work, mut event_list) = self.inner_transport.borrow_mut().process()?;
trace!(
"({}).Transport.inner_process() - output: {} {}",
self.identifier.clone(),
inner_did_work,
event_list.len()
);
if inner_did_work {
did_work = true;
outbox.append(&mut event_list);
}
for evt in outbox.clone() {
self.handle_TransportEvent(&evt)?;
}
Ok((did_work, outbox))
}
fn connection_id_list(&self) -> TransportResult<Vec<ConnectionId>> {
let peer_data_list = self.inner_dht.get_peer_list();
let mut id_list = Vec::new();
for peer_data in peer_data_list {
id_list.push(peer_data.peer_address);
}
Ok(id_list)
}
fn get_uri(&self, id: &ConnectionIdRef) -> Option<Url> {
self.inner_transport.borrow().get_uri(id)
}
}
impl<T: Transport, D: Dht> P2pGateway<T, D> {
pub(crate) fn connection_id_to_dht_uri_list(
&self,
id_list: &[&ConnectionIdRef],
) -> TransportResult<Vec<Url>> {
let mut uri_list = Vec::with_capacity(id_list.len());
for connectionId in id_list {
let maybe_peer = self.inner_dht.get_peer(connectionId);
match maybe_peer {
None => {
return Err(TransportError::new(format!(
"Unknown connectionId: {}",
connectionId
)));
}
Some(peer) => uri_list.push(peer.peer_uri),
}
}
Ok(uri_list)
}
pub(crate) fn handle_TransportEvent(&mut self, evt: &TransportEvent) -> TransportResult<()> {
debug!(
"<<< '({})' recv transport event: {:?}",
self.identifier.clone(),
evt
);
match evt {
TransportEvent::TransportError(id, e) => {
error!(
"(GatewayTransport) Connection Error for {}: {}\n Closing connection.",
id, e
);
self.inner_transport.borrow_mut().close(id)?;
}
TransportEvent::ConnectResult(id) => {
info!("({}) Connection opened id: {}", self.identifier.clone(), id);
if let Some(uri) = self.get_uri(id) {
trace!(
"(GatewayTransport).ConnectResult: mapping {} -> {}",
uri,
id
);
self.connection_map.insert(uri, id.clone());
let our_peer_address = P2pProtocol::PeerAddress(
self.identifier().to_string(),
self.this_peer().clone().peer_address,
);
let mut buf = Vec::new();
our_peer_address
.serialize(&mut Serializer::new(&mut buf))
.unwrap();
trace!(
"(GatewayTransport) P2pProtocol::PeerAddress: {:?} to {:?}",
our_peer_address,
id
);
self.inner_transport.borrow_mut().send(&[&id], &buf)?;
}
}
TransportEvent::Connection(_id) => {
unimplemented!();
}
TransportEvent::Closed(id) => {
warn!("Connection closed: {}", id);
self.inner_transport.borrow_mut().close(id)?;
}
TransportEvent::Received(id, payload) => {
debug!("Received message from: {}", id);
let mut de = Deserializer::new(&payload[..]);
let maybe_p2p_msg: Result<P2pProtocol, rmp_serde::decode::Error> =
Deserialize::deserialize(&mut de);
if let Ok(p2p_msg) = maybe_p2p_msg {
if let P2pProtocol::PeerAddress(gateway_id, peer_address) = p2p_msg {
debug!(
"Received PeerAddress: {} | {} ({})",
peer_address, gateway_id, self.identifier
);
if self.identifier == gateway_id {
let peer = PeerData {
peer_address: peer_address.clone(),
peer_uri: transport_id_to_url(id.clone()),
timestamp: 42, };
Dht::post(self, DhtCommand::HoldPeer(peer)).expect("FIXME");
Dht::process(self).expect("HACK");
}
}
}
}
};
Ok(())
}
#[allow(non_snake_case)]
fn serve_TransportCommand(
&mut self,
cmd: &TransportCommand,
) -> TransportResult<Vec<TransportEvent>> {
trace!(
"({}) serving transport cmd: {:?}",
self.identifier.clone(),
cmd
);
match cmd {
TransportCommand::Connect(url) => {
let id = self.connect(url)?;
let evt = TransportEvent::ConnectResult(id);
Ok(vec![evt])
}
TransportCommand::Send(id_list, payload) => {
let mut id_ref_list = Vec::with_capacity(id_list.len());
for id in id_list {
id_ref_list.push(id.as_str());
}
let _id = self.send(&id_ref_list, payload)?;
Ok(vec![])
}
TransportCommand::SendAll(payload) => {
let _id = self.send_all(payload)?;
Ok(vec![])
}
TransportCommand::Close(id) => {
self.close(id)?;
let evt = TransportEvent::Closed(id.to_string());
Ok(vec![evt])
}
TransportCommand::CloseAll => {
self.close_all()?;
let outbox = Vec::new();
Ok(outbox)
}
TransportCommand::Bind(url) => {
self.bind(url)?;
Ok(vec![])
}
}
}
}