use crate::{
client::NameSystemClient,
dht::{
DhtConfig, DhtError, DhtKeyMaterial, DhtNode, DhtRecord, NetworkInfo, Peer, RecordValidator,
},
records::NsRecord,
utils::make_p2p_address,
PeerId,
};
use anyhow::{anyhow, Result};
use async_trait::async_trait;
#[cfg(doc)]
use cid::Cid;
use futures::future::try_join_all;
use libp2p::Multiaddr;
use noosphere_core::data::Did;
use std::collections::HashMap;
use tokio::sync::{Mutex, MutexGuard};
pub static BOOTSTRAP_PEERS_ADDRESSES: [&str; 1] =
["/ip4/134.122.20.28/tcp/6666/p2p/12D3KooWPyjAB3XWUboGmLLPkR53fTyj4GaNi65RvQ61BVwqV4HG"];
lazy_static! {
pub static ref BOOTSTRAP_PEERS: [Multiaddr; 1] = BOOTSTRAP_PEERS_ADDRESSES.map(|addr| addr.parse().expect("parseable"));
}
pub struct NameSystem {
pub(crate) dht: DhtNode,
hosted_records: Mutex<HashMap<Did, NsRecord>>,
resolved_records: Mutex<HashMap<Did, NsRecord>>,
#[cfg(feature = "api_server")]
api_server: Option<APIServer>,
}
impl NameSystem {
pub fn new<K: DhtKeyMaterial, V: RecordValidator + 'static>(
key_material: &K,
dht_config: DhtConfig,
validator: Option<V>,
) -> Result<Self> {
Ok(NameSystem {
dht: DhtNode::new(key_material, dht_config, validator)?,
hosted_records: Mutex::new(HashMap::new()),
resolved_records: Mutex::new(HashMap::new()),
})
}
pub async fn propagate_records(&self) -> Result<()> {
let hosted_records = self.hosted_records.lock().await;
if hosted_records.is_empty() {
return Ok(());
}
let pending_tasks: Vec<_> = hosted_records
.iter()
.map(|(identity, record)| self.dht_put_record(identity, record))
.collect();
try_join_all(pending_tasks).await?;
Ok(())
}
pub async fn flush_records(&self) {
let mut resolved_records = self.resolved_records.lock().await;
resolved_records.drain();
}
pub async fn flush_records_for_identity(&self, identity: &Did) -> bool {
let mut resolved_records = self.resolved_records.lock().await;
resolved_records.remove(identity).is_some()
}
pub async fn get_cache(&self) -> MutexGuard<HashMap<Did, NsRecord>> {
self.resolved_records.lock().await
}
async fn dht_get_record(&self, identity: &Did) -> Result<(Did, Option<NsRecord>)> {
match self.dht.get_record(identity.as_bytes()).await {
Ok(DhtRecord { key: _, value }) => match value {
Some(value) => {
let record = NsRecord::try_from(value)?;
info!(
"NameSystem: GetRecord: {} {}",
identity,
record
.link()
.map_or_else(|| String::from("None"), |cid| cid.to_string())
);
Ok((identity.to_owned(), Some(record)))
}
None => {
warn!("NameSystem: GetRecord: No record found for {}.", identity);
Ok((identity.to_owned(), None))
}
},
Err(e) => {
warn!("NameSystem: GetRecord: Failure for {} {:?}.", identity, e);
Err(anyhow!(e.to_string()))
}
}
}
async fn dht_put_record(&self, identity: &Did, record: &NsRecord) -> Result<()> {
let record: Vec<u8> = record.try_into()?;
match self.dht.put_record(identity.as_bytes(), &record).await {
Ok(_) => {
info!("NameSystem: PutRecord: {}", identity);
Ok(())
}
Err(e) => {
warn!("NameSystem: PutRecord: Failure for {} {:?}.", identity, e);
Err(anyhow!(e.to_string()))
}
}
}
}
#[async_trait]
impl NameSystemClient for NameSystem {
async fn network_info(&self) -> Result<NetworkInfo> {
self.dht.network_info().await.map_err(|e| e.into())
}
fn peer_id(&self) -> &PeerId {
self.dht.peer_id()
}
async fn add_peers(&self, peers: Vec<Multiaddr>) -> Result<()> {
self.dht.add_peers(peers).await.map_err(|e| e.into())
}
async fn peers(&self) -> Result<Vec<Peer>> {
self.dht.peers().await.map_err(|e| e.into())
}
async fn listen(&self, listening_address: Multiaddr) -> Result<Multiaddr> {
self.dht
.listen(listening_address)
.await
.map_err(|e| e.into())
}
async fn stop_listening(&self) -> Result<()> {
self.dht.stop_listening().await.map_err(|e| e.into())
}
async fn bootstrap(&self) -> Result<()> {
self.dht.bootstrap().await.map_err(|e| e.into())
}
async fn address(&self) -> Result<Option<Multiaddr>> {
let mut addresses = self
.dht
.addresses()
.await
.map_err(<DhtError as Into<anyhow::Error>>::into)?;
if !addresses.is_empty() {
let peer_id = self.peer_id().to_owned();
let address = make_p2p_address(addresses.swap_remove(0), peer_id);
Ok(Some(address))
} else {
Ok(None)
}
}
async fn put_record(&self, record: NsRecord) -> Result<()> {
let identity = Did::from(record.identity());
self.dht_put_record(&identity, &record).await?;
self.hosted_records.lock().await.insert(identity, record);
Ok(())
}
async fn get_record(&self, identity: &Did) -> Result<Option<NsRecord>> {
{
let mut resolved_records = self.resolved_records.lock().await;
if let Some(record) = resolved_records.get(identity) {
if !record.is_expired() {
return Ok(Some(record.clone()));
} else {
resolved_records.remove(identity);
}
}
};
match self.dht_get_record(identity).await? {
(_, Some(record)) => {
let mut resolved_records = self.resolved_records.lock().await;
resolved_records.insert(identity.to_owned(), record.clone());
Ok(Some(record))
}
(_, None) => Ok(None),
}
}
}
#[cfg(test)]
mod test {
use super::*;
#[test]
fn bootstrap_peers_parseable() {
assert_eq!(BOOTSTRAP_PEERS.len(), 1);
}
use crate::{ns_client_tests, Validator};
use crate::{utils::wait_for_peers, NameSystemBuilder, NameSystemClient};
use noosphere_core::authority::generate_ed25519_key;
use noosphere_storage::{MemoryStorage, SphereDb};
use std::sync::Arc;
use tokio::sync::Mutex;
struct DataPlaceholder {
_bootstrap: NameSystem,
_ns: Arc<Mutex<NameSystem>>,
}
async fn before_each() -> Result<(DataPlaceholder, Arc<Mutex<NameSystem>>)> {
let (bootstrap, bootstrap_address) = {
let key_material = generate_ed25519_key();
let store = SphereDb::new(&MemoryStorage::default()).await.unwrap();
let ns = NameSystemBuilder::default()
.validator(Validator::new(store.clone()))
.key_material(&key_material)
.listening_port(0)
.use_test_config()
.build()
.await
.unwrap();
ns.bootstrap().await.unwrap();
let address = ns.address().await?.unwrap();
(ns, address)
};
let ns = {
let key_material = generate_ed25519_key();
let store = SphereDb::new(&MemoryStorage::default()).await.unwrap();
let ns = NameSystemBuilder::default()
.validator(Validator::new(store.clone()))
.key_material(&key_material)
.bootstrap_peers(&[bootstrap_address.clone()])
.use_test_config()
.build()
.await
.unwrap();
ns.bootstrap().await.unwrap();
wait_for_peers::<NameSystem>(&ns, 1).await?;
ns
};
let client = Arc::new(Mutex::new(ns));
let reference = client.clone();
let data = DataPlaceholder {
_ns: reference,
_bootstrap: bootstrap,
};
Ok((data, client))
}
ns_client_tests!(NameSystem, before_each, DataPlaceholder);
}