#![allow(non_snake_case)]
use crate::{
dht::{dht_protocol::*, dht_trait::Dht},
engine::p2p_protocol::*,
gateway::P2pGateway,
transport::transport_trait::Transport,
};
use lib3h_protocol::{data_types::EntryData, AddressRef, DidWork, Lib3hResult};
use rmp_serde::Serializer;
use serde::Serialize;
impl<T: Transport, D: Dht> Dht for P2pGateway<T, D> {
fn get_peer(&self, peer_address: &str) -> Option<PeerData> {
self.inner_dht.get_peer(peer_address)
}
fn fetch_peer(&self, peer_address: &str) -> Option<PeerData> {
self.inner_dht.fetch_peer(peer_address)
}
fn get_entry(&self, entry_address: &AddressRef) -> Option<EntryData> {
self.inner_dht.get_entry(entry_address)
}
fn fetch_entry(&self, entry_address: &AddressRef) -> Option<EntryData> {
self.inner_dht.fetch_entry(entry_address)
}
fn post(&mut self, cmd: DhtCommand) -> Lib3hResult<()> {
self.inner_dht.post(cmd)
}
fn process(&mut self) -> Lib3hResult<(DidWork, Vec<DhtEvent>)> {
let (did_work, dht_event_list) = self.inner_dht.process()?;
println!(
"[t] ({}).Dht.process() - output: {} {}",
self.identifier.clone(),
did_work,
dht_event_list.len()
);
if did_work {
for evt in dht_event_list.clone() {
self.handle_DhtEvent(evt)?;
}
}
Ok((did_work, dht_event_list))
}
fn this_peer(&self) -> &PeerData {
self.inner_dht.this_peer()
}
fn get_peer_list(&self) -> Vec<PeerData> {
self.inner_dht.get_peer_list()
}
}
impl<T: Transport, D: Dht> P2pGateway<T, D> {
pub(crate) fn handle_DhtEvent(&mut self, evt: DhtEvent) -> Lib3hResult<()> {
println!(
"[t] ({}).handle_DhtEvent() {:?}",
self.identifier.clone(),
evt
);
match evt {
DhtEvent::GossipTo(data) => {
for to_peer_address in data.peer_address_list {
let me = &self.inner_dht.this_peer().peer_address;
if &to_peer_address == me {
continue;
}
let peer_transport = self
.inner_dht
.get_peer(&to_peer_address)
.expect("Should gossip to a known peer")
.transport;
println!(
"({}) GossipTo: {} {}",
self.identifier.clone(),
to_peer_address,
peer_transport
);
let p2p_gossip = P2pProtocol::Gossip(GossipData {
space_address: self.identifier().as_bytes().to_vec(),
to_peer_address: to_peer_address.as_bytes().to_vec(),
from_peer_address: self.this_peer().peer_address.as_bytes().to_vec(),
bundle: data.bundle.clone(),
});
let mut payload = Vec::new();
p2p_gossip
.serialize(&mut Serializer::new(&mut payload))
.unwrap();
self.inner_transport
.borrow_mut()
.send(&[peer_transport.as_str()], &payload)?;
}
}
DhtEvent::GossipUnreliablyTo(_data) => {
}
DhtEvent::HoldPeerRequested(_peer_address) => {
}
DhtEvent::PeerTimedOut(_data) => {
}
DhtEvent::HoldEntryRequested(_from, _data) => {
}
DhtEvent::FetchEntryResponse(_data) => {
}
DhtEvent::EntryPruned(_address) => {
}
}
Ok(())
}
}