mod blob_apis;
mod commands;
mod data;
mod queries;
mod register_apis;
use crate::client::{
connections::{Session, SAFE_CLIENT_DIR},
errors::Error,
ClientConfig,
};
use crate::messaging::data::{CmdError, DataQuery, ServiceMsg};
use crate::types::{ChunkAddress, Keypair, PublicKey};
use crate::messaging::{ServiceAuth, WireMsg};
use crate::prefix_map::NetworkPrefixMap;
use crate::types::utils::read_prefix_map_from_disk;
use itertools::Itertools;
use rand::rngs::OsRng;
use std::collections::BTreeSet;
use std::net::SocketAddr;
use std::sync::Arc;
use tokio::{
sync::{mpsc::Receiver, RwLock},
time::Duration,
};
use tracing::{debug, info};
use xor_name::XorName;
#[derive(Clone, Debug)]
pub struct Client {
keypair: Keypair,
incoming_errors: Arc<RwLock<Receiver<CmdError>>>,
session: Session,
pub(crate) query_timeout: Duration,
}
impl Client {
#[instrument(skip_all, level = "debug", name = "New client")]
pub async fn new(
config: ClientConfig,
bootstrap_nodes: BTreeSet<SocketAddr>,
optional_keypair: Option<Keypair>,
) -> Result<Self, Error> {
Client::create_with(config, bootstrap_nodes, optional_keypair, true).await
}
pub(crate) async fn create_with(
config: ClientConfig,
bootstrap_nodes: BTreeSet<SocketAddr>,
optional_keypair: Option<Keypair>,
read_prefixmap: bool,
) -> Result<Self, Error> {
let mut rng = OsRng;
let keypair = match optional_keypair {
Some(id) => {
info!("Client started for specific pk: {:?}", id.public_key());
id
}
None => {
let keypair = Keypair::new_ed25519(&mut rng);
info!(
"Client started for new randomly created pk: {:?}",
keypair.public_key()
);
keypair
}
};
let home_dir = dirs_next::home_dir()
.ok_or_else(|| Error::Generic("Error opening home dir".to_string()))?;
let root_dir = &home_dir
.join(SAFE_CLIENT_DIR)
.join(format!("sn_client-{}", keypair.public_key()));
let prefix_map = if read_prefixmap {
match read_prefix_map_from_disk(&home_dir.join(".safe/prefix_map")).await {
Ok(prefix_map) => prefix_map,
Err(e) => {
warn!("Could not read PrefixMap at '.safe/prefix_map': {:?}", e);
info!("Checking Client's root dir");
match read_prefix_map_from_disk(root_dir).await {
Ok(map) => map,
Err(e) => {
warn!("Could not read PrefixMap at Client's root dir: {:?}", e);
info!("Defaulting to a fresh NetworkPrefixMap");
NetworkPrefixMap::new(config.genesis_key)
}
}
}
}
} else {
NetworkPrefixMap::new(config.genesis_key)
};
if config.genesis_key != prefix_map.genesis_key() {
return Err(Error::Generic(
"Genesis Key from the config and the PrefixMap mismatch".to_string(),
));
}
let (err_sender, err_receiver) = tokio::sync::mpsc::channel::<CmdError>(10);
let client_pk = keypair.public_key();
debug!(
"Creating new session with genesis key: {:?} ",
config.genesis_key
);
debug!(
"Creating new session with genesis key (in hex format): {} ",
hex::encode(config.genesis_key.to_bytes())
);
let session = Session::new(
client_pk,
config.genesis_key,
config.qp2p,
err_sender,
config.local_addr,
config.standard_wait,
prefix_map,
)
.await?;
let client = Self {
keypair,
session,
incoming_errors: Arc::new(RwLock::new(err_receiver)),
query_timeout: config.query_timeout,
};
let random_dst_addr = XorName::random();
let serialised_cmd = {
let msg = ServiceMsg::Query(DataQuery::GetChunk(ChunkAddress(random_dst_addr)));
WireMsg::serialize_msg_payload(&msg)?
};
let signature = client.keypair.sign(&serialised_cmd);
let auth = ServiceAuth {
public_key: client_pk,
signature,
};
let bootstrap_nodes = bootstrap_nodes.iter().copied().collect_vec();
client
.session
.make_contact_with_nodes(
bootstrap_nodes.clone(),
random_dst_addr,
auth,
serialised_cmd,
)
.await?;
Ok(client)
}
pub fn keypair(&self) -> Keypair {
self.keypair.clone()
}
pub fn public_key(&self) -> PublicKey {
self.keypair().public_key()
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::client::utils::test_utils::{
create_test_client, create_test_client_with, init_test_logger,
};
use crate::types::utils::random_bytes;
use crate::types::Scope;
use eyre::Result;
use std::{
collections::HashSet,
net::{IpAddr, Ipv4Addr, SocketAddr},
};
#[tokio::test(flavor = "multi_thread")]
async fn client_creation() -> Result<()> {
init_test_logger();
let _client = create_test_client().await?;
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
#[ignore]
async fn client_nonsense_bootstrap_fails() -> Result<()> {
init_test_logger();
let mut nonsense_bootstrap = HashSet::new();
let _ = nonsense_bootstrap.insert(SocketAddr::new(
IpAddr::V4(Ipv4Addr::new(127, 0, 0, 1)),
3033,
));
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn client_creation_with_existing_keypair() -> Result<()> {
init_test_logger();
let mut rng = OsRng;
let full_id = Keypair::new_ed25519(&mut rng);
let pk = full_id.public_key();
let client = create_test_client_with(Some(full_id), None, true).await?;
assert_eq!(pk, client.public_key());
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn long_lived_connection_survives() -> Result<()> {
init_test_logger();
let client = create_test_client().await?;
tokio::time::sleep(tokio::time::Duration::from_secs(40)).await;
let bytes = random_bytes(self_encryption::MIN_ENCRYPTABLE_BYTES / 2);
let _ = client.upload(bytes, Scope::Public).await?;
Ok(())
}
#[test]
fn client_is_send() {
init_test_logger();
fn require_send<T: Send>(_t: T) {}
require_send(create_test_client());
}
}