use std::path::PathBuf;
use std::sync::Arc;
use std::time::Duration;
use nym_bandwidth_controller::config::BandwidthControllerConfig;
use nym_bandwidth_controller::requests::BandwidthControllerRequestSender;
use nym_bandwidth_controller::{BandwidthController, BandwidthTicketProvider, TicketType};
use nym_bandwidth_fetcher::NyxdCredentialFetcher;
use nym_credentials_interface::BandwidthCredential;
use nym_crypto::asymmetric::{ed25519, x25519};
use nym_lp::peer::{DHKeyPair, LpRemotePeer};
use nym_network_defaults::NymNetworkDetails;
use nym_registration_client::{LpRegistrationClient, NestedLpSession};
use nym_registration_common::WireguardConfiguration;
use nym_task::ShutdownToken;
use nym_validator_client::nym_api::NymApiClientExt;
use nym_validator_client::DirectSigningHttpRpcNyxdClient;
use rand09::SeedableRng;
use sha2::{Digest, Sha256};
use time::OffsetDateTime;
use tokio::net::TcpStream;
use tokio::task::JoinHandle;
use tokio_util::sync::CancellationToken;
use url::Url;
use zeroize::Zeroizing;
use nym_api_requests::models::described::v2::NymNodeDescriptionV2;
use crate::config::{RestockPolicy, SessionConfig};
use crate::dvpn::{DvpnDirectory, QuicBridge};
use crate::error::SessionError;
use crate::fetcher::TimeoutFetcher;
use crate::gateway::{self, GatewayInfo, GatewaySpec, SelectedGateway, WgRole};
use crate::registration_cache::RegistrationCache;
const TICKETS_TO_SPEND: u32 = 1;
const API_TIMEOUT: Duration = Duration::from_secs(30);
const PROVISIONING_TIMEOUT: Duration = Duration::from_secs(5 * 60);
fn wireguard_ticket_types() -> Vec<TicketType> {
vec![TicketType::V1WireguardEntry, TicketType::V1WireguardExit]
}
fn needed_ticket_types(two_hop: bool) -> Vec<TicketType> {
let mut types = vec![TicketType::V1WireguardEntry];
if two_hop {
types.push(TicketType::V1WireguardExit);
}
types
}
pub struct HopConfig {
pub wg_config: WireguardConfiguration,
pub client_private_key: x25519::PrivateKey,
pub gateway_identity: ed25519::PublicKey,
pub gateway: GatewayInfo,
pub bridge: Option<QuicBridge>,
}
pub struct Registration {
pub entry: HopConfig,
pub exit: Option<HopConfig>,
}
struct OwnedController {
sender: BandwidthControllerRequestSender,
task: JoinHandle<()>,
}
pub struct Session {
api: nym_http_api_client::Client,
provider: Arc<dyn BandwidthTicketProvider>,
owned: Option<OwnedController>,
cancel: CancellationToken,
directory: Option<DvpnDirectory>,
reg_cache: Option<std::sync::Mutex<RegistrationCache>>,
}
impl Session {
pub async fn new(
config: SessionConfig,
cancel: CancellationToken,
) -> Result<Self, SessionError> {
let SessionConfig {
mnemonic,
network,
credential_store_path,
data_path,
dvpn_directory_url,
automatic_topups,
bandwidth_provider,
reuse_registrations,
} = config;
let reg_cache = reuse_registrations.then(|| {
std::sync::Mutex::new(RegistrationCache::load(
&data_path,
network.network_name.clone(),
))
});
let api_url_str = network
.endpoints
.iter()
.find_map(|e| e.api_url.clone())
.ok_or(SessionError::MissingEndpoint { which: "nym-api" })?;
let api_url = Url::parse(&api_url_str).map_err(|source| SessionError::InvalidUrl {
which: "nym-api",
url: api_url_str.clone(),
source,
})?;
let api = nym_http_api_client::Client::new(api_url, Some(API_TIMEOUT));
let directory = match dvpn_directory_url {
Some(url) => match DvpnDirectory::fetch(&url).await {
Ok(dir) => Some(dir),
Err(e) => {
tracing::warn!("failed to fetch dVPN directory at {url}: {e}");
Some(DvpnDirectory::default())
}
},
None => None,
};
let (provider, owned) = match bandwidth_provider {
Some(external) => (external, None),
None => {
let (provider, owned) = Self::spawn_controller(
mnemonic,
network,
credential_store_path,
data_path,
automatic_topups,
cancel.clone(),
)
.await?;
(provider, Some(owned))
}
};
Ok(Self {
api,
provider,
owned,
cancel,
directory,
reg_cache,
})
}
async fn spawn_controller(
mnemonic: bip39::Mnemonic,
network: NymNetworkDetails,
credential_store_path: Option<PathBuf>,
data_path: PathBuf,
automatic_topups: Option<RestockPolicy>,
cancel: CancellationToken,
) -> Result<(Arc<dyn BandwidthTicketProvider>, OwnedController), SessionError> {
let nyxd_url = network
.endpoints
.first()
.map(|e| e.nyxd_url.clone())
.ok_or(SessionError::MissingEndpoint { which: "nyxd" })?;
let client_id = derive_client_id(&mnemonic);
let nyxd = DirectSigningHttpRpcNyxdClient::connect_with_mnemonic_and_network_details(
nyxd_url.as_str(),
network,
mnemonic,
)?;
let nyxd = Arc::new(nyxd);
let store_path = credential_store_path.unwrap_or_else(|| data_path.join("credentials.db"));
if let Some(parent) = store_path.parent() {
std::fs::create_dir_all(parent).map_err(|e| SessionError::Storage(e.to_string()))?;
}
let storage = nym_credential_storage::initialise_persistent_storage(&store_path).await;
let fetcher_db = data_path.join("fetcher-requests.db");
let fetcher = NyxdCredentialFetcher::new(nyxd, &fetcher_db, client_id)
.await
.map_err(|e| SessionError::Issuance(e.to_string()))?;
let fetcher = TimeoutFetcher::new(fetcher);
let config = match automatic_topups {
Some(policy) => {
let mut config: BandwidthControllerConfig = policy.into();
config.managed_ticket_types = wireguard_ticket_types();
config
}
None => BandwidthControllerConfig {
managed_ticket_types: Vec::new(),
..Default::default()
},
};
let controller = BandwidthController::new(storage)
.with_config(config)
.with_credential_fetcher(fetcher);
let sender = controller.get_request_sender();
let shutdown = ShutdownToken::new_from_tokio_token(cancel.clone());
let task = tokio::spawn(async move { controller.run(shutdown).await });
let provider: Arc<dyn BandwidthTicketProvider> = Arc::new(sender.clone());
Ok((provider, OwnedController { sender, task }))
}
pub async fn ensure_ticketbooks(&self, two_hop: bool) -> Result<(), SessionError> {
self.ensure_ticket_types(needed_ticket_types(two_hop)).await
}
async fn ensure_ticket_types(&self, types: Vec<TicketType>) -> Result<(), SessionError> {
if types.is_empty() {
return Ok(());
}
let Some(owned) = &self.owned else {
return Ok(());
};
tokio::select! {
biased;
_ = self.cancel.cancelled() => Err(SessionError::Cancelled),
res = tokio::time::timeout(PROVISIONING_TIMEOUT, async {
owned
.sender
.restock_ticketbooks(types.clone())
.await
.map_err(|e| SessionError::Issuance(e.to_string()))?;
owned
.sender
.wait_for_ticketbooks(types)
.await
.map_err(|e| SessionError::Issuance(e.to_string()))
}) => res.unwrap_or(Err(SessionError::ProvisioningTimeout {
after: PROVISIONING_TIMEOUT,
})),
}
}
pub async fn obtain_wireguard_credential(
&self,
gateway_id: ed25519::PublicKey,
role: WgRole,
) -> Result<BandwidthCredential, SessionError> {
let ticket_type = match role {
WgRole::Entry => TicketType::V1WireguardEntry,
WgRole::Exit => TicketType::V1WireguardExit,
};
let prepared = self
.provider
.get_ecash_ticket(
ticket_type,
gateway_id,
TICKETS_TO_SPEND,
OffsetDateTime::now_utc(),
)
.await
.map_err(|e| SessionError::Issuance(e.to_string()))?
.ok_or_else(|| {
SessionError::Issuance("no stored ticket available for top-up".into())
})?;
Ok(BandwidthCredential::from(prepared.data))
}
pub fn bandwidth_provider(&self) -> Arc<dyn BandwidthTicketProvider> {
self.provider.clone()
}
async fn fetch_topology(&self) -> Result<Vec<NymNodeDescriptionV2>, SessionError> {
Ok(self.api.get_all_described_nodes_v2().await?)
}
async fn fetch_topology_cancellable(&self) -> Result<Vec<NymNodeDescriptionV2>, SessionError> {
tokio::select! {
biased;
_ = self.cancel.cancelled() => Err(SessionError::Cancelled),
res = self.fetch_topology() => res,
}
}
pub async fn select_gateway(
&self,
spec: &GatewaySpec,
role: WgRole,
) -> Result<SelectedGateway, SessionError> {
tokio::select! {
biased;
_ = self.cancel.cancelled() => Err(SessionError::Cancelled),
res = async {
let nodes = self.fetch_topology().await?;
gateway::select(&nodes, spec, role, self.directory.as_ref(), false, None)
} => res,
}
}
pub async fn register_single_hop(
&self,
gateway: &GatewaySpec,
) -> Result<Registration, SessionError> {
self.register_single_inner(gateway).await
}
fn with_cache<R>(&self, f: impl FnOnce(&mut RegistrationCache) -> R) -> Option<R> {
self.reg_cache.as_ref().map(|cache| {
f(&mut cache
.lock()
.unwrap_or_else(|poisoned| poisoned.into_inner()))
})
}
fn cached_hop(
&self,
identity: &ed25519::PublicKey,
gateway: GatewayInfo,
role: WgRole,
) -> Option<HopConfig> {
let cached = self.with_cache(|cache| cache.lookup(identity, role))??;
tracing::info!(
"reusing cached registration for {} ({role:?}) — no ticket spent",
identity.to_base58_string()
);
Some(HopConfig {
wg_config: cached.wg_config,
client_private_key: cached.client_private_key,
gateway_identity: *identity,
gateway,
bridge: None,
})
}
fn finalize_hop(
&self,
identity: &ed25519::PublicKey,
gateway: GatewayInfo,
role: WgRole,
client_private_key: x25519::PrivateKey,
wg_config: WireguardConfiguration,
) -> HopConfig {
self.with_cache(|cache| cache.insert(identity, role, &client_private_key, &wg_config));
HopConfig {
wg_config,
client_private_key,
gateway_identity: *identity,
gateway,
bridge: None,
}
}
pub fn invalidate_registration(&self, gateway: &ed25519::PublicKey, role: WgRole) {
self.with_cache(|cache| cache.remove(gateway, role));
}
async fn register_single_inner(
&self,
gateway: &GatewaySpec,
) -> Result<Registration, SessionError> {
let nodes = self.fetch_topology_cancellable().await?;
let selected = gateway::select(
&nodes,
gateway,
WgRole::Entry,
self.directory.as_ref(),
false,
None,
)?;
if let Some(hop) = self.cached_hop(&selected.identity, selected.info(), WgRole::Entry) {
return Ok(Registration {
entry: hop,
exit: None,
});
}
self.ensure_ticket_types(vec![TicketType::V1WireguardEntry])
.await?;
let hop = self
.register_hop(&selected, TicketType::V1WireguardEntry)
.await?;
Ok(Registration {
entry: hop,
exit: None,
})
}
pub async fn register_two_hop(
&self,
entry: &GatewaySpec,
exit: &GatewaySpec,
) -> Result<Registration, SessionError> {
self.register_two_hop_inner(entry, exit, false).await
}
pub async fn register_two_hop_quic(
&self,
entry: &GatewaySpec,
exit: &GatewaySpec,
) -> Result<Registration, SessionError> {
self.register_two_hop_inner(entry, exit, true).await
}
async fn register_two_hop_inner(
&self,
entry: &GatewaySpec,
exit: &GatewaySpec,
entry_quic: bool,
) -> Result<Registration, SessionError> {
let nodes = self.fetch_topology_cancellable().await?;
let entry_gw = gateway::select(
&nodes,
entry,
WgRole::Entry,
self.directory.as_ref(),
entry_quic,
None,
)?;
let exit_gw = gateway::select(
&nodes,
exit,
WgRole::Exit,
self.directory.as_ref(),
false,
Some(&entry_gw.identity),
)?;
let entry_bridge = if entry_quic {
entry_gw.quic.clone()
} else {
None
};
let (cached_entry, cached_exit) = match (
self.cached_hop(&entry_gw.identity, entry_gw.info(), WgRole::Entry),
self.cached_hop(&exit_gw.identity, exit_gw.info(), WgRole::Exit),
) {
(Some(mut entry_hop), Some(exit_hop)) => {
entry_hop.bridge = entry_bridge;
return Ok(Registration {
entry: entry_hop,
exit: Some(exit_hop),
});
}
partial => partial,
};
let mut needed = Vec::new();
if cached_entry.is_none() {
needed.push(TicketType::V1WireguardEntry);
}
if cached_exit.is_none() {
needed.push(TicketType::V1WireguardExit);
}
self.ensure_ticket_types(needed).await?;
let entry_lp = lp_info(&entry_gw)?;
let exit_lp = lp_info(&exit_gw)?;
let entry_keypair = Arc::new(DHKeyPair::new(&mut rand09::rng()));
let entry_peer =
LpRemotePeer::new(entry_lp.x25519).with_key_digests(entry_lp.expected_kem_key_hashes);
let mut entry_client = LpRegistrationClient::<TcpStream>::new_with_default_config(
entry_keypair,
entry_peer,
entry_lp.address,
entry_lp.ciphersuite,
entry_lp.lp_protocol_version,
);
tokio::select! {
biased;
_ = self.cancel.cancelled() => return Err(SessionError::Cancelled),
r = entry_client.perform_handshake() => r.map_err(|source| SessionError::Registration {
address: entry_lp.address,
source,
})?,
}
let mut rng = rand09::rngs::StdRng::from_os_rng();
let exit_hop = match cached_exit {
Some(hop) => hop,
None => {
let exit_keypair = Arc::new(DHKeyPair::new(&mut rand09::rng()));
let exit_peer = LpRemotePeer::new(exit_lp.x25519)
.with_key_digests(exit_lp.expected_kem_key_hashes);
let mut nested = NestedLpSession::new(
exit_lp.address,
exit_keypair,
exit_peer,
exit_lp.ciphersuite,
exit_lp.lp_protocol_version,
);
let exit_wg = x25519::KeyPair::new(&mut rand::thread_rng());
let exit_cfg = nested
.handshake_and_register_dvpn::<TcpStream, _>(
&mut entry_client,
&mut rng,
&exit_wg,
&exit_gw.identity,
self.provider.as_ref(),
None,
TicketType::V1WireguardExit,
)
.await
.map_err(|source| SessionError::Registration {
address: exit_lp.address,
source,
})?;
self.finalize_hop(
&exit_gw.identity,
exit_gw.info(),
WgRole::Exit,
x25519::PrivateKey::from_secret(exit_wg.private_key().to_bytes()),
exit_cfg,
)
}
};
let mut entry_hop = match cached_entry {
Some(hop) => hop,
None => {
let entry_wg = x25519::KeyPair::new(&mut rand::thread_rng());
let entry_cfg = entry_client
.register_dvpn(
&mut rng,
&entry_wg,
&entry_gw.identity,
self.provider.as_ref(),
None,
TicketType::V1WireguardEntry,
)
.await
.map_err(|source| SessionError::Registration {
address: entry_lp.address,
source,
})?;
self.finalize_hop(
&entry_gw.identity,
entry_gw.info(),
WgRole::Entry,
x25519::PrivateKey::from_secret(entry_wg.private_key().to_bytes()),
entry_cfg,
)
}
};
entry_hop.bridge = entry_bridge;
Ok(Registration {
entry: entry_hop,
exit: Some(exit_hop),
})
}
async fn register_hop(
&self,
selected: &SelectedGateway,
ticket_type: TicketType,
) -> Result<HopConfig, SessionError> {
let lp = lp_info(selected)?;
let keypair = Arc::new(DHKeyPair::new(&mut rand09::rng()));
let peer = LpRemotePeer::new(lp.x25519).with_key_digests(lp.expected_kem_key_hashes);
let mut client = LpRegistrationClient::<TcpStream>::new_with_default_config(
keypair,
peer,
lp.address,
lp.ciphersuite,
lp.lp_protocol_version,
);
tokio::select! {
biased;
_ = self.cancel.cancelled() => return Err(SessionError::Cancelled),
r = client.perform_handshake() => r.map_err(|source| SessionError::Registration {
address: lp.address,
source,
})?,
}
let mut rng = rand09::rngs::StdRng::from_os_rng();
let wg = x25519::KeyPair::new(&mut rand::thread_rng());
let cfg = client
.register_dvpn(
&mut rng,
&wg,
&selected.identity,
self.provider.as_ref(),
None,
ticket_type,
)
.await
.map_err(|source| SessionError::Registration {
address: lp.address,
source,
})?;
let role = match ticket_type {
TicketType::V1WireguardExit => WgRole::Exit,
_ => WgRole::Entry,
};
Ok(self.finalize_hop(
&selected.identity,
selected.info(),
role,
x25519::PrivateKey::from_secret(wg.private_key().to_bytes()),
cfg,
))
}
pub async fn shutdown(mut self) {
self.cancel.cancel();
if let Some(owned) = self.owned.take() {
let _ = owned.task.await;
}
}
}
impl Drop for Session {
fn drop(&mut self) {
self.cancel.cancel();
}
}
fn derive_client_id(mnemonic: &bip39::Mnemonic) -> Zeroizing<Vec<u8>> {
let entropy = Zeroizing::new(mnemonic.to_entropy());
let mut hasher = Sha256::new();
hasher.update(b"nym-sdk-session::client-id::v1");
hasher.update(entropy.as_slice());
Zeroizing::new(hasher.finalize().to_vec())
}
fn lp_info(
selected: &SelectedGateway,
) -> Result<nym_registration_common::NymNodeLPInformation, SessionError> {
selected
.node
.node
.lp_data
.clone()
.ok_or_else(|| SessionError::MalformedGateway {
identity: selected.identity.to_base58_string(),
reason: "gateway advertises no LP data".to_string(),
})
}
#[cfg(test)]
#[path = "session_tests.rs"]
mod tests;