use crate::dht::{
dht_protocol::*,
dht_trait::{Dht, DhtConfig},
};
use lib3h_protocol::{data_types::EntryData, Address, DidWork, Lib3hResult};
use std::collections::{HashMap, HashSet, VecDeque};
use rmp_serde::{Deserializer, Serializer};
use serde::{Deserialize, Serialize};
use url::Url;
#[derive(Debug, Clone, PartialEq, Deserialize, Serialize)]
enum MirrorGossip {
Entry(EntryData),
Peer(PeerData),
}
pub struct MirrorDht {
inbox: VecDeque<DhtCommand>,
entry_list: HashMap<Address, HashSet<Address>>,
peer_list: HashMap<String, PeerData>,
this_peer: PeerData,
pending_fetch_request_list: HashSet<String>,
}
impl MirrorDht {
pub fn new(peer_address: &str, peer_transport: &Url) -> Self {
MirrorDht {
inbox: VecDeque::new(),
peer_list: HashMap::new(),
entry_list: HashMap::new(),
this_peer: PeerData {
peer_address: peer_address.to_string(),
peer_uri: peer_transport.clone(),
timestamp: 0, },
pending_fetch_request_list: HashSet::new(),
}
}
pub fn new_with_config(config: &DhtConfig) -> Lib3hResult<Self> {
Ok(Self::new(&config.this_peer_address, &config.this_peer_uri))
}
}
impl Dht for MirrorDht {
fn get_peer_list(&self) -> Vec<PeerData> {
self.peer_list.values().map(|v| v.clone()).collect()
}
fn get_peer(&self, peer_address: &str) -> Option<PeerData> {
let res = self.peer_list.get(peer_address);
if let Some(pd) = res {
return Some(pd.clone());
}
None
}
fn this_peer(&self) -> &PeerData {
&self.this_peer
}
fn get_entry_address_list(&self) -> Vec<&Address> {
self.entry_list.iter().map(|kv| kv.0).collect()
}
fn get_aspects_of(&self, entry_address: &Address) -> Option<Vec<Address>> {
match self.entry_list.get(entry_address) {
None => None,
Some(set) => {
let vec = set.iter().map(|addr| addr.clone()).collect();
Some(vec)
}
}
}
fn post(&mut self, cmd: DhtCommand) -> Lib3hResult<()> {
self.inbox.push_back(cmd);
Ok(())
}
fn process(&mut self) -> Lib3hResult<(DidWork, Vec<DhtEvent>)> {
let mut outbox = Vec::new();
let mut did_work = false;
loop {
let cmd = match self.inbox.pop_front() {
None => break,
Some(msg) => msg,
};
let res = self.serve_DhtCommand(&cmd);
if let Ok(mut output) = res {
did_work = true;
outbox.append(&mut output);
} else {
error!("serve_DhtCommand() failed: {:?}", res);
}
}
Ok((did_work, outbox))
}
}
impl MirrorDht {
fn add_peer(&mut self, peer_info: &PeerData) -> bool {
trace!("@MirrorDht@ Adding peer: {:?}", peer_info);
let maybe_peer = self.peer_list.get_mut(&peer_info.peer_address);
match maybe_peer {
None => {
trace!("@MirrorDht@ Adding peer - OK NEW");
self.peer_list
.insert(peer_info.peer_address.clone(), peer_info.clone());
true
}
Some(mut peer) => {
if peer_info.timestamp <= peer.timestamp {
trace!("@MirrorDht@ Adding peer - BAD");
return false;
}
trace!("@MirrorDht@ Adding peer - OK UPDATED");
peer.timestamp = peer_info.timestamp;
true
}
}
}
fn diff_aspects(&self, entry: &EntryData) -> HashSet<Address> {
let aspect_address_set: HashSet<_> = entry
.aspect_list
.iter()
.map(|aspect| aspect.aspect_address.clone())
.collect();
let maybe_aspects = self.entry_list.get(&entry.entry_address);
if maybe_aspects.is_none() {
return aspect_address_set;
}
let held_aspects: HashSet<_> = maybe_aspects.unwrap().clone();
let diff: HashSet<_> = aspect_address_set
.difference(&held_aspects)
.map(|item| item.clone())
.collect();
diff
}
fn add_entry_aspects(&mut self, entry: &EntryData) -> bool {
let diff: HashSet<_> = self.diff_aspects(&entry);
if diff.len() == 0 {
return false;
}
let maybe_known_aspects = self.entry_list.get(&entry.entry_address);
let new_aspects: HashSet<_> = match maybe_known_aspects {
None => diff,
Some(known_aspects) => known_aspects
.union(&diff)
.map(|item| item.clone())
.collect(),
};
self.entry_list
.insert(entry.entry_address.clone(), new_aspects);
true
}
fn gossip_entry(&self, entry: &EntryData) -> DhtEvent {
let entry_gossip = MirrorGossip::Entry(entry.clone());
let mut buf = Vec::new();
entry_gossip
.serialize(&mut Serializer::new(&mut buf))
.unwrap();
let gossip_evt = GossipToData {
peer_address_list: self
.get_peer_list()
.iter()
.map(|pi| pi.peer_address.clone())
.collect(),
bundle: buf,
};
DhtEvent::GossipTo(gossip_evt)
}
#[allow(non_snake_case)]
fn serve_DhtCommand(&mut self, cmd: &DhtCommand) -> Lib3hResult<Vec<DhtEvent>> {
debug!("@MirrorDht@ serving cmd: {:?}", cmd);
match cmd {
DhtCommand::HandleGossip(msg) => {
trace!("Deserializer msg.bundle: {:?}", msg.bundle);
let mut de = Deserializer::new(&msg.bundle[..]);
let maybe_gossip: Result<MirrorGossip, rmp_serde::decode::Error> =
Deserialize::deserialize(&mut de);
if let Err(e) = maybe_gossip {
return Err(format_err!("Failed deserializing gossip: {:?}", e));
}
match maybe_gossip.unwrap() {
MirrorGossip::Entry(entry) => {
let diff = self.diff_aspects(&entry);
if diff.len() > 0 {
return Ok(vec![DhtEvent::HoldEntryRequested(
self.this_peer.peer_address.clone(),
entry,
)]);
}
return Ok(vec![]);
}
MirrorGossip::Peer(peer) => {
let is_new = self.add_peer(&peer);
if is_new {
return Ok(vec![DhtEvent::HoldPeerRequested(peer)]);
}
return Ok(vec![]);
}
}
}
DhtCommand::FetchEntry(fetch_entry) => {
self.pending_fetch_request_list
.insert(fetch_entry.msg_id.clone());
return Ok(vec![DhtEvent::EntryDataRequested(fetch_entry.clone())]);
}
DhtCommand::HoldPeer(msg) => {
let peer_address_list: Vec<String> = self
.get_peer_list()
.iter()
.map(|pi| pi.peer_address.clone())
.collect();
let received_new_content = self.add_peer(msg);
if !received_new_content {
return Ok(vec![]);
}
let mut event_list = Vec::new();
let peer = self
.peer_list
.get(&msg.peer_address)
.expect("Should have peer by now");
let peer_gossip = MirrorGossip::Peer(peer.clone());
let mut buf = Vec::new();
peer_gossip
.serialize(&mut Serializer::new(&mut buf))
.unwrap();
trace!("@MirrorDht@ gossiping peer: {:?}", peer);
let gossip_evt = GossipToData {
peer_address_list,
bundle: buf,
};
event_list.push(DhtEvent::GossipTo(gossip_evt));
let peer = self.this_peer();
if msg.peer_address != peer.peer_address {
let peer_gossip = MirrorGossip::Peer(peer.clone());
let mut buf = Vec::new();
peer_gossip
.serialize(&mut Serializer::new(&mut buf))
.unwrap();
trace!(
"@MirrorDht@ gossiping peer back: {:?} | to: {}",
peer,
msg.peer_address
);
let gossip_evt = GossipToData {
peer_address_list: vec![msg.peer_address.clone()],
bundle: buf,
};
event_list.push(DhtEvent::GossipTo(gossip_evt));
}
Ok(event_list)
}
DhtCommand::HoldEntryAspectAddress(entry) => {
let received_new_content = self.add_entry_aspects(&entry);
if !received_new_content {
return Ok(vec![]);
}
let address_str = entry.entry_address.clone();
self.pending_fetch_request_list
.insert(address_str.to_string());
let fetch_entry = FetchDhtEntryData {
msg_id: address_str.to_string(),
entry_address: entry.entry_address.to_owned(),
};
Ok(vec![DhtEvent::EntryDataRequested(fetch_entry)])
}
DhtCommand::BroadcastEntry(entry) => {
let received_new_content = self.add_entry_aspects(&entry);
if !received_new_content {
return Ok(vec![]);
}
let gossip_evt = self.gossip_entry(entry);
Ok(vec![gossip_evt])
}
DhtCommand::DropEntryAddress(_) => Ok(vec![]),
DhtCommand::EntryDataResponse(response) => {
if !self.pending_fetch_request_list.remove(&response.msg_id) {
return Err(format_err!("Received response for an unknown request"));
}
let address_str: String = (&response.entry.entry_address).clone().into();
if address_str == response.msg_id {
let gossip_evt = self.gossip_entry(&response.entry);
return Ok(vec![gossip_evt]);
}
Ok(vec![DhtEvent::FetchEntryResponse(response.clone())])
}
}
}
}