use crate::{
dht::{DHTConfig, DHTKeyMaterial, DHTNode, DHTRecord},
records::NSRecord,
validator::Validator,
};
use anyhow::{anyhow, Result};
use futures::future::try_join_all;
use libp2p::Multiaddr;
use noosphere_core::authority::SUPPORTED_KEYS;
use noosphere_storage::{SphereDb, Storage};
use std::collections::HashMap;
use ucan::crypto::did::DidParser;
#[cfg(doc)]
use crate::NameSystemBuilder;
#[cfg(doc)]
use cid::Cid;
pub static BOOTSTRAP_PEERS_ADDRESSES: [&str; 1] =
["/ip4/134.122.20.28/tcp/6666/p2p/12D3KooWAKxaCWsSGauqhCZXyeYjDSQ6jna2SmVhLn4J7uQFdvot"];
lazy_static! {
pub static ref BOOTSTRAP_PEERS: [Multiaddr; 1] = BOOTSTRAP_PEERS_ADDRESSES.map(|addr| addr.parse().expect("parseable"));
}
pub struct NameSystem<S, K>
where
S: Storage + 'static,
K: DHTKeyMaterial,
{
pub(crate) bootstrap_peers: Option<Vec<Multiaddr>>,
pub(crate) dht: Option<DHTNode<Validator<S>>>,
pub(crate) dht_config: DHTConfig,
pub(crate) key_material: K,
pub(crate) store: SphereDb<S>,
hosted_records: HashMap<String, NSRecord>,
resolved_records: HashMap<String, NSRecord>,
did_parser: DidParser,
}
impl<S, K> NameSystem<S, K>
where
S: Storage,
K: DHTKeyMaterial,
{
pub(crate) fn new(
key_material: K,
store: SphereDb<S>,
bootstrap_peers: Option<Vec<Multiaddr>>,
dht_config: DHTConfig,
) -> Self {
NameSystem {
key_material,
store,
bootstrap_peers,
dht_config,
dht: None,
hosted_records: HashMap::new(),
resolved_records: HashMap::new(),
did_parser: DidParser::new(SUPPORTED_KEYS),
}
}
pub async fn connect(&mut self) -> Result<()> {
let mut dht = DHTNode::new(
&self.key_material,
self.bootstrap_peers.as_ref(),
Validator::new(&self.store),
&self.dht_config,
)?;
dht.run().map_err(|e| anyhow!(e.to_string()))?;
dht.bootstrap().await.map_err(|e| anyhow!(e.to_string()))?;
dht.wait_for_peers(1)
.await
.map_err(|e| anyhow!(e.to_string()))?;
self.dht = Some(dht);
Ok(())
}
pub fn disconnect(&mut self) -> Result<()> {
if let Some(mut dht) = self.dht.take() {
dht.terminate()?;
}
Ok(())
}
pub async fn propagate_records(&self) -> Result<()> {
let _ = self.require_dht()?;
if self.hosted_records.is_empty() {
return Ok(());
}
let pending_tasks: Vec<_> = self
.hosted_records
.iter()
.map(|(identity, record)| self.dht_put_record(identity, record))
.collect();
try_join_all(pending_tasks).await?;
Ok(())
}
pub async fn put_record(&mut self, record: NSRecord) -> Result<()> {
let _ = self.require_dht()?;
record.validate(&self.store, &mut self.did_parser).await?;
let identity = record.identity();
self.dht_put_record(identity, &record).await?;
self.hosted_records.insert(identity.to_owned(), record);
Ok(())
}
pub async fn get_record(&mut self, identity: &str) -> Result<Option<NSRecord>> {
if let Some(record) = self.resolved_records.get(identity) {
if !record.is_expired() {
return Ok(Some(record.clone()));
} else {
self.resolved_records.remove(identity);
}
}
match self.dht_get_record(identity).await? {
(_, Some(record)) => {
self.resolved_records
.insert(identity.to_owned(), record.clone());
Ok(Some(record))
}
(_, None) => Ok(None),
}
}
pub fn flush_records(&mut self) {
self.resolved_records.drain();
}
pub fn flush_records_for_identity(&mut self, identity: &String) -> bool {
self.resolved_records.remove(identity).is_some()
}
pub fn get_cache(&self) -> &HashMap<String, NSRecord> {
&self.resolved_records
}
pub fn get_cache_mut(&mut self) -> &mut HashMap<String, NSRecord> {
&mut self.resolved_records
}
pub fn p2p_address(&self) -> Option<&Multiaddr> {
if let Some(dht) = &self.dht {
dht.p2p_address()
} else {
None
}
}
async fn dht_get_record(&self, identity: &str) -> Result<(String, Option<NSRecord>)> {
let dht = self.require_dht()?;
match 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
.address()
.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: &str, record: &NSRecord) -> Result<()> {
let dht = self.require_dht()?;
let record: Vec<u8> = record.try_into()?;
match 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()))
}
}
}
fn require_dht(&self) -> Result<&DHTNode<Validator<S>>> {
self.dht.as_ref().ok_or_else(|| anyhow!("not connected"))
}
}
impl<S, K> Drop for NameSystem<S, K>
where
S: Storage,
K: DHTKeyMaterial,
{
fn drop(&mut self) {
if let Err(e) = self.disconnect() {
error!("{}", e.to_string());
}
}
}
#[cfg(test)]
mod test {
use super::*;
#[test]
fn bootstrap_peers_parseable() {
assert_eq!(BOOTSTRAP_PEERS.len(), 1);
}
}