Feature/topology refactor (#274)
* VersionFilterable for HashMap * Removed NymTopology trait in favour of concrete type * Removed providers from NymTopology * Made gateway conversion use reference, similarly to mixes * Using more concrete types in topology rather than b58 strings * Allowing gateways to have DNS-resolvable mix listener address * Error propagation for gateway key conversion
This commit is contained in:
committed by
GitHub
parent
e849e45b12
commit
a53d0a4aac
Generated
+2
-10
@@ -693,6 +693,7 @@ dependencies = [
|
||||
name = "directory-client-models"
|
||||
version = "0.1.0"
|
||||
dependencies = [
|
||||
"crypto",
|
||||
"serde",
|
||||
"topology",
|
||||
]
|
||||
@@ -1414,15 +1415,6 @@ dependencies = [
|
||||
"url 1.7.2",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "itertools"
|
||||
version = "0.8.2"
|
||||
source = "registry+https://github.com/rust-lang/crates.io-index"
|
||||
checksum = "f56a2d0bc861f9165be4eb3442afd3c236d8a98afd426f65d92324ae1091a484"
|
||||
dependencies = [
|
||||
"either",
|
||||
]
|
||||
|
||||
[[package]]
|
||||
name = "itoa"
|
||||
version = "0.4.5"
|
||||
@@ -3337,7 +3329,7 @@ name = "topology"
|
||||
version = "0.1.0"
|
||||
dependencies = [
|
||||
"bs58",
|
||||
"itertools",
|
||||
"crypto",
|
||||
"log 0.4.8",
|
||||
"nymsphinx-addressing",
|
||||
"nymsphinx-types",
|
||||
|
||||
@@ -27,12 +27,10 @@ use std::sync::Arc;
|
||||
use tokio::runtime::Handle;
|
||||
use tokio::task::JoinHandle;
|
||||
use tokio::time;
|
||||
use topology::NymTopology;
|
||||
|
||||
pub(crate) struct LoopCoverTrafficStream<R, T>
|
||||
pub(crate) struct LoopCoverTrafficStream<R>
|
||||
where
|
||||
R: CryptoRng + Rng,
|
||||
T: NymTopology,
|
||||
{
|
||||
/// Key used to encrypt and decrypt content of an ACK packet.
|
||||
ack_key: Arc<AckAes128Key>,
|
||||
@@ -61,13 +59,12 @@ where
|
||||
rng: R,
|
||||
|
||||
/// Accessor to the common instance of network topology.
|
||||
topology_access: TopologyAccessor<T>,
|
||||
topology_access: TopologyAccessor,
|
||||
}
|
||||
|
||||
impl<R, T> Stream for LoopCoverTrafficStream<R, T>
|
||||
impl<R> Stream for LoopCoverTrafficStream<R>
|
||||
where
|
||||
R: CryptoRng + Rng + Unpin,
|
||||
T: NymTopology, // this really confuses me, why T doesn't need to be Unpin?
|
||||
{
|
||||
// Item is only used to indicate we should create a new message rather than actual cover message
|
||||
// reason being to not introduce unnecessary complexity by having to keep state of topology
|
||||
@@ -98,7 +95,7 @@ where
|
||||
|
||||
// obviously when we finally make shared rng that is on 'higher' level, this should become
|
||||
// generic `R`
|
||||
impl<T: 'static + NymTopology> LoopCoverTrafficStream<OsRng, T> {
|
||||
impl LoopCoverTrafficStream<OsRng> {
|
||||
pub(crate) fn new(
|
||||
ack_key: Arc<AckAes128Key>,
|
||||
average_ack_delay: time::Duration,
|
||||
@@ -106,7 +103,7 @@ impl<T: 'static + NymTopology> LoopCoverTrafficStream<OsRng, T> {
|
||||
average_cover_message_sending_delay: time::Duration,
|
||||
mix_tx: MixMessageSender,
|
||||
our_full_destination: Recipient,
|
||||
topology_access: TopologyAccessor<T>,
|
||||
topology_access: TopologyAccessor,
|
||||
) -> Self {
|
||||
let rng = OsRng;
|
||||
|
||||
|
||||
@@ -38,7 +38,6 @@ use nymsphinx::NodeAddressBytes;
|
||||
use received_buffer::{ReceivedBufferMessage, ReconstructedMessagesReceiver};
|
||||
use std::sync::Arc;
|
||||
use tokio::runtime::Runtime;
|
||||
use topology::NymTopology;
|
||||
|
||||
mod cover_traffic_stream;
|
||||
pub(crate) mod inbound_messages;
|
||||
@@ -72,7 +71,9 @@ impl NymClient {
|
||||
|
||||
pub fn as_mix_recipient(&self) -> Recipient {
|
||||
Recipient::new(
|
||||
self.identity_keypair.public_key().derive_address(),
|
||||
self.identity_keypair
|
||||
.public_key()
|
||||
.derive_destination_address(),
|
||||
// TODO: below only works under assumption that gateway address == gateway id
|
||||
// (which currently is true)
|
||||
NodeAddressBytes::try_from_base58_string(self.config.get_gateway_id()).unwrap(),
|
||||
@@ -81,10 +82,10 @@ impl NymClient {
|
||||
|
||||
// future constantly pumping loop cover traffic at some specified average rate
|
||||
// the pumped traffic goes to the MixTrafficController
|
||||
fn start_cover_traffic_stream<T: 'static + NymTopology>(
|
||||
fn start_cover_traffic_stream(
|
||||
&self,
|
||||
ack_key: Arc<AckAes128Key>,
|
||||
topology_accessor: TopologyAccessor<T>,
|
||||
topology_accessor: TopologyAccessor,
|
||||
mix_tx: MixMessageSender,
|
||||
) {
|
||||
info!("Starting loop cover traffic stream...");
|
||||
@@ -108,9 +109,9 @@ impl NymClient {
|
||||
// TODO: I'm not a fan of this function signature, i.e. that it returns the ACK key, but then again
|
||||
// I think the ACK key should be generated by 'acknowledgement_control' rather than by
|
||||
// the client itself. However, it all might change once we have to implement key rotation.
|
||||
fn start_real_traffic_controller<T: 'static + NymTopology>(
|
||||
fn start_real_traffic_controller(
|
||||
&self,
|
||||
topology_accessor: TopologyAccessor<T>,
|
||||
topology_accessor: TopologyAccessor,
|
||||
ack_receiver: AcknowledgementReceiver,
|
||||
input_receiver: InputMessageReceiver,
|
||||
mix_sender: MixMessageSender,
|
||||
@@ -210,10 +211,7 @@ impl NymClient {
|
||||
|
||||
// future responsible for periodically polling directory server and updating
|
||||
// the current global view of topology
|
||||
fn start_topology_refresher(
|
||||
&mut self,
|
||||
topology_accessor: TopologyAccessor<directory_client::Topology>,
|
||||
) {
|
||||
fn start_topology_refresher(&mut self, topology_accessor: TopologyAccessor) {
|
||||
let topology_refresher_config = TopologyRefresherConfig::new(
|
||||
self.config.get_directory_server(),
|
||||
self.config.get_topology_refresh_rate(),
|
||||
@@ -256,9 +254,9 @@ impl NymClient {
|
||||
MixTrafficController::new(mix_rx, gateway_client).start(self.runtime.handle());
|
||||
}
|
||||
|
||||
fn start_websocket_listener<T: 'static + NymTopology>(
|
||||
fn start_websocket_listener(
|
||||
&self,
|
||||
topology_accessor: TopologyAccessor<T>,
|
||||
topology_accessor: TopologyAccessor,
|
||||
buffer_requester: ReceivedBufferRequestSender,
|
||||
msg_input: InputMessageSender,
|
||||
) {
|
||||
@@ -343,7 +341,7 @@ impl NymClient {
|
||||
|
||||
// channels responsible for controlling ack messages
|
||||
let (ack_sender, ack_receiver) = mpsc::unbounded();
|
||||
let shared_topology_accessor = TopologyAccessor::<directory_client::Topology>::new();
|
||||
let shared_topology_accessor = TopologyAccessor::new();
|
||||
|
||||
// the components are started in very specific order. Unless you know what you are doing,
|
||||
// do not change that.
|
||||
|
||||
+4
-7
@@ -25,13 +25,11 @@ use nymsphinx::{
|
||||
};
|
||||
use rand::{CryptoRng, Rng};
|
||||
use std::sync::Arc;
|
||||
use topology::NymTopology;
|
||||
|
||||
// responsible for splitting received message and initial sending attempt
|
||||
pub(super) struct InputMessageListener<R, T>
|
||||
pub(super) struct InputMessageListener<R>
|
||||
where
|
||||
R: CryptoRng + Rng,
|
||||
T: NymTopology,
|
||||
{
|
||||
ack_key: Arc<AckAes128Key>,
|
||||
ack_recipient: Recipient,
|
||||
@@ -39,13 +37,12 @@ where
|
||||
message_chunker: MessageChunker<R>,
|
||||
pending_acks: PendingAcksMap,
|
||||
real_message_sender: RealMessageSender,
|
||||
topology_access: TopologyAccessor<T>,
|
||||
topology_access: TopologyAccessor,
|
||||
}
|
||||
|
||||
impl<R, T> InputMessageListener<R, T>
|
||||
impl<R> InputMessageListener<R>
|
||||
where
|
||||
R: CryptoRng + Rng,
|
||||
T: NymTopology,
|
||||
{
|
||||
pub(super) fn new(
|
||||
ack_key: Arc<AckAes128Key>,
|
||||
@@ -54,7 +51,7 @@ where
|
||||
message_chunker: MessageChunker<R>,
|
||||
pending_acks: PendingAcksMap,
|
||||
real_message_sender: RealMessageSender,
|
||||
topology_access: TopologyAccessor<T>,
|
||||
topology_access: TopologyAccessor,
|
||||
) -> Self {
|
||||
InputMessageListener {
|
||||
ack_key,
|
||||
|
||||
@@ -38,7 +38,6 @@ use tokio::{
|
||||
sync::{Notify, RwLock},
|
||||
task::JoinHandle,
|
||||
};
|
||||
use topology::NymTopology;
|
||||
|
||||
mod acknowledgement_listener;
|
||||
mod input_message_listener;
|
||||
@@ -98,26 +97,24 @@ impl AcknowledgementControllerConnectors {
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) struct AcknowledgementController<R, T>
|
||||
pub(super) struct AcknowledgementController<R>
|
||||
where
|
||||
R: CryptoRng + Rng,
|
||||
T: NymTopology,
|
||||
{
|
||||
ack_key: Arc<AckAes128Key>,
|
||||
acknowledgement_listener: Option<AcknowledgementListener>,
|
||||
input_message_listener: Option<InputMessageListener<R, T>>,
|
||||
retransmission_request_listener: Option<RetransmissionRequestListener<R, T>>,
|
||||
input_message_listener: Option<InputMessageListener<R>>,
|
||||
retransmission_request_listener: Option<RetransmissionRequestListener<R>>,
|
||||
sent_notification_listener: Option<SentNotificationListener>,
|
||||
}
|
||||
|
||||
impl<R, T> AcknowledgementController<R, T>
|
||||
impl<R> AcknowledgementController<R>
|
||||
where
|
||||
R: 'static + CryptoRng + Rng + Clone + Send,
|
||||
T: 'static + NymTopology,
|
||||
{
|
||||
pub(super) fn new(
|
||||
mut rng: R,
|
||||
topology_access: TopologyAccessor<T>,
|
||||
topology_access: TopologyAccessor,
|
||||
ack_recipient: Recipient,
|
||||
average_packet_delay_duration: Duration,
|
||||
average_ack_delay_duration: Duration,
|
||||
|
||||
+4
-7
@@ -26,13 +26,11 @@ use nymsphinx::{
|
||||
};
|
||||
use rand::{CryptoRng, Rng};
|
||||
use std::sync::Arc;
|
||||
use topology::NymTopology;
|
||||
|
||||
// responsible for packet retransmission upon fired timer
|
||||
pub(super) struct RetransmissionRequestListener<R, T>
|
||||
pub(super) struct RetransmissionRequestListener<R>
|
||||
where
|
||||
R: CryptoRng + Rng,
|
||||
T: NymTopology,
|
||||
{
|
||||
ack_key: Arc<AckAes128Key>,
|
||||
ack_recipient: Recipient,
|
||||
@@ -40,13 +38,12 @@ where
|
||||
pending_acks: PendingAcksMap,
|
||||
real_message_sender: RealMessageSender,
|
||||
request_receiver: RetransmissionRequestReceiver,
|
||||
topology_access: TopologyAccessor<T>,
|
||||
topology_access: TopologyAccessor,
|
||||
}
|
||||
|
||||
impl<R, T> RetransmissionRequestListener<R, T>
|
||||
impl<R> RetransmissionRequestListener<R>
|
||||
where
|
||||
R: CryptoRng + Rng,
|
||||
T: NymTopology,
|
||||
{
|
||||
pub(super) fn new(
|
||||
ack_key: Arc<AckAes128Key>,
|
||||
@@ -55,7 +52,7 @@ where
|
||||
pending_acks: PendingAcksMap,
|
||||
real_message_sender: RealMessageSender,
|
||||
request_receiver: RetransmissionRequestReceiver,
|
||||
topology_access: TopologyAccessor<T>,
|
||||
topology_access: TopologyAccessor,
|
||||
) -> Self {
|
||||
RetransmissionRequestListener {
|
||||
ack_key,
|
||||
|
||||
@@ -34,7 +34,6 @@ use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use tokio::runtime::Handle;
|
||||
use tokio::task::JoinHandle;
|
||||
use topology::NymTopology;
|
||||
|
||||
mod acknowlegement_control;
|
||||
mod real_traffic_stream;
|
||||
@@ -68,24 +67,23 @@ impl Config {
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) struct RealMessagesController<R, T>
|
||||
pub(crate) struct RealMessagesController<R>
|
||||
where
|
||||
R: CryptoRng + Rng,
|
||||
T: NymTopology,
|
||||
{
|
||||
out_queue_control: Option<OutQueueControl<R, T>>,
|
||||
ack_control: Option<AcknowledgementController<R, T>>,
|
||||
out_queue_control: Option<OutQueueControl<R>>,
|
||||
ack_control: Option<AcknowledgementController<R>>,
|
||||
}
|
||||
|
||||
// obviously when we finally make shared rng that is on 'higher' level, this should become
|
||||
// generic `R`
|
||||
impl<T: 'static + NymTopology> RealMessagesController<OsRng, T> {
|
||||
impl RealMessagesController<OsRng> {
|
||||
pub(crate) fn new(
|
||||
config: Config,
|
||||
ack_receiver: AcknowledgementReceiver,
|
||||
input_receiver: InputMessageReceiver,
|
||||
mix_sender: MixMessageSender,
|
||||
topology_access: TopologyAccessor<T>,
|
||||
topology_access: TopologyAccessor,
|
||||
) -> Self {
|
||||
let rng = OsRng;
|
||||
|
||||
|
||||
@@ -31,12 +31,10 @@ use std::pin::Pin;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use tokio::time;
|
||||
use topology::NymTopology;
|
||||
|
||||
pub(crate) struct OutQueueControl<R, T>
|
||||
pub(crate) struct OutQueueControl<R>
|
||||
where
|
||||
R: CryptoRng + Rng,
|
||||
T: NymTopology,
|
||||
{
|
||||
/// Key used to encrypt and decrypt content of an ACK packet.
|
||||
ack_key: Arc<AckAes128Key>,
|
||||
@@ -72,7 +70,7 @@ where
|
||||
rng: R,
|
||||
|
||||
/// Accessor to the common instance of network topology.
|
||||
topology_access: TopologyAccessor<T>,
|
||||
topology_access: TopologyAccessor,
|
||||
}
|
||||
|
||||
pub(crate) struct RealMessage {
|
||||
@@ -105,10 +103,9 @@ pub(crate) enum StreamMessage {
|
||||
Real(RealMessage),
|
||||
}
|
||||
|
||||
impl<R, T> Stream for OutQueueControl<R, T>
|
||||
impl<R> Stream for OutQueueControl<R>
|
||||
where
|
||||
R: CryptoRng + Rng + Unpin,
|
||||
T: NymTopology, // this really confuses me, why T doesn't need to be Unpin?
|
||||
{
|
||||
type Item = StreamMessage;
|
||||
|
||||
@@ -144,10 +141,9 @@ where
|
||||
}
|
||||
}
|
||||
|
||||
impl<R, T> OutQueueControl<R, T>
|
||||
impl<R> OutQueueControl<R>
|
||||
where
|
||||
R: CryptoRng + Rng + Unpin,
|
||||
T: NymTopology,
|
||||
{
|
||||
pub(crate) fn new(
|
||||
ack_key: Arc<AckAes128Key>,
|
||||
@@ -159,7 +155,7 @@ where
|
||||
real_receiver: RealMessageReceiver,
|
||||
rng: R,
|
||||
our_full_destination: Recipient,
|
||||
topology_access: TopologyAccessor<T>,
|
||||
topology_access: TopologyAccessor,
|
||||
) -> Self {
|
||||
OutQueueControl {
|
||||
ack_key,
|
||||
|
||||
@@ -16,6 +16,8 @@ use crate::built_info;
|
||||
use directory_client::DirectoryClient;
|
||||
use log::*;
|
||||
use nymsphinx::addressing::clients::Recipient;
|
||||
use nymsphinx::params::DEFAULT_NUM_MIX_HOPS;
|
||||
use std::convert::TryInto;
|
||||
use std::ops::Deref;
|
||||
use std::sync::Arc;
|
||||
use std::time;
|
||||
@@ -27,66 +29,65 @@ use topology::{gateway, NymTopology};
|
||||
|
||||
// I'm extremely curious why compiler NEVER complained about lack of Debug here before
|
||||
#[derive(Debug)]
|
||||
struct TopologyAccessorInner<T: NymTopology>(Option<T>);
|
||||
pub(super) struct TopologyAccessorInner(Option<NymTopology>);
|
||||
|
||||
impl<T: NymTopology> Deref for TopologyAccessorInner<T> {
|
||||
type Target = Option<T>;
|
||||
|
||||
fn deref(&self) -> &Self::Target {
|
||||
impl AsRef<Option<NymTopology>> for TopologyAccessorInner {
|
||||
fn as_ref(&self) -> &Option<NymTopology> {
|
||||
&self.0
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: NymTopology> TopologyAccessorInner<T> {
|
||||
impl TopologyAccessorInner {
|
||||
fn new() -> Self {
|
||||
TopologyAccessorInner(None)
|
||||
}
|
||||
|
||||
fn update(&mut self, new: Option<T>) {
|
||||
fn update(&mut self, new: Option<NymTopology>) {
|
||||
self.0 = new;
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) struct TopologyReadPermit<'a, T: NymTopology> {
|
||||
permit: RwLockReadGuard<'a, TopologyAccessorInner<T>>,
|
||||
pub(super) struct TopologyReadPermit<'a> {
|
||||
permit: RwLockReadGuard<'a, TopologyAccessorInner>,
|
||||
}
|
||||
|
||||
impl<'a, T: NymTopology> TopologyReadPermit<'a, T> {
|
||||
impl<'a> Deref for TopologyReadPermit<'a> {
|
||||
type Target = TopologyAccessorInner;
|
||||
|
||||
fn deref(&self) -> &Self::Target {
|
||||
&self.permit
|
||||
}
|
||||
}
|
||||
|
||||
impl<'a> TopologyReadPermit<'a> {
|
||||
/// Using provided topology read permit, tries to get an immutable reference to the underlying
|
||||
/// topology. For obvious reasons the lifetime of the topology reference is bound to the permit.
|
||||
pub(super) fn try_get_valid_topology_ref(
|
||||
&'a self,
|
||||
ack_recipient: &Recipient,
|
||||
packet_recipient: &Recipient,
|
||||
) -> Option<&'a T> {
|
||||
// first we need to deref out of RwLockReadGuard
|
||||
// then we need to deref out of TopologyAccessorInner
|
||||
// then we must take ref of option, i.e. Option<&T>
|
||||
// and finally try to unwrap it to obtain &T
|
||||
let topology_ref_option = (*self.permit.deref()).as_ref();
|
||||
|
||||
if topology_ref_option.is_none() {
|
||||
return None;
|
||||
}
|
||||
|
||||
let topology_ref = topology_ref_option.unwrap();
|
||||
|
||||
// see if it's possible to route the packet to both gateways
|
||||
if !topology_ref.can_construct_path_through()
|
||||
|| !topology_ref.gateway_exists(&packet_recipient.gateway())
|
||||
|| !topology_ref.gateway_exists(&ack_recipient.gateway())
|
||||
{
|
||||
None
|
||||
} else {
|
||||
Some(topology_ref)
|
||||
) -> Option<&'a NymTopology> {
|
||||
// Note: implicit deref with Deref for TopologyReadPermit is happening here
|
||||
let topology_ref_option = self.permit.as_ref();
|
||||
match topology_ref_option {
|
||||
None => None,
|
||||
Some(topology_ref) => {
|
||||
// see if it's possible to route the packet to both gateways
|
||||
if !topology_ref.can_construct_path_through(DEFAULT_NUM_MIX_HOPS)
|
||||
|| !topology_ref.gateway_exists(&packet_recipient.gateway())
|
||||
|| !topology_ref.gateway_exists(&ack_recipient.gateway())
|
||||
{
|
||||
None
|
||||
} else {
|
||||
Some(topology_ref)
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl<'a, T: NymTopology> From<RwLockReadGuard<'a, TopologyAccessorInner<T>>>
|
||||
for TopologyReadPermit<'a, T>
|
||||
{
|
||||
fn from(read_permit: RwLockReadGuard<'a, TopologyAccessorInner<T>>) -> Self {
|
||||
impl<'a> From<RwLockReadGuard<'a, TopologyAccessorInner>> for TopologyReadPermit<'a> {
|
||||
fn from(read_permit: RwLockReadGuard<'a, TopologyAccessorInner>) -> Self {
|
||||
TopologyReadPermit {
|
||||
permit: read_permit,
|
||||
}
|
||||
@@ -94,26 +95,26 @@ impl<'a, T: NymTopology> From<RwLockReadGuard<'a, TopologyAccessorInner<T>>>
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
pub(crate) struct TopologyAccessor<T: NymTopology> {
|
||||
pub(crate) struct TopologyAccessor {
|
||||
// `RwLock` *seems to* be the better approach for this as write access is only requested every
|
||||
// few seconds, while reads are needed every single packet generated.
|
||||
// However, proper benchmarks will be needed to determine if `RwLock` is indeed a better
|
||||
// approach than a `Mutex`
|
||||
inner: Arc<RwLock<TopologyAccessorInner<T>>>,
|
||||
inner: Arc<RwLock<TopologyAccessorInner>>,
|
||||
}
|
||||
|
||||
impl<T: NymTopology> TopologyAccessor<T> {
|
||||
impl TopologyAccessor {
|
||||
pub(crate) fn new() -> Self {
|
||||
TopologyAccessor {
|
||||
inner: Arc::new(RwLock::new(TopologyAccessorInner::new())),
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) async fn get_read_permit(&self) -> TopologyReadPermit<'_, T> {
|
||||
pub(super) async fn get_read_permit(&self) -> TopologyReadPermit<'_> {
|
||||
self.inner.read().await.into()
|
||||
}
|
||||
|
||||
async fn update_global_topology(&mut self, new_topology: Option<T>) {
|
||||
async fn update_global_topology(&mut self, new_topology: Option<NymTopology>) {
|
||||
self.inner.write().await.update(new_topology);
|
||||
}
|
||||
|
||||
@@ -133,7 +134,7 @@ impl<T: NymTopology> TopologyAccessor<T> {
|
||||
pub(crate) async fn is_routable(&self) -> bool {
|
||||
match &self.inner.read().await.0 {
|
||||
None => false,
|
||||
Some(ref topology) => topology.can_construct_path_through(),
|
||||
Some(ref topology) => topology.can_construct_path_through(DEFAULT_NUM_MIX_HOPS),
|
||||
}
|
||||
}
|
||||
|
||||
@@ -167,17 +168,16 @@ impl TopologyRefresherConfig {
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) struct TopologyRefresher<T: NymTopology> {
|
||||
pub(crate) struct TopologyRefresher {
|
||||
directory_client: directory_client::Client,
|
||||
topology_accessor: TopologyAccessor<T>,
|
||||
topology_accessor: TopologyAccessor,
|
||||
refresh_rate: Duration,
|
||||
}
|
||||
|
||||
// TODO: consider (or maybe not) restoring generic TopologyRefresher<T>
|
||||
impl TopologyRefresher<directory_client::Topology> {
|
||||
impl TopologyRefresher {
|
||||
pub(crate) fn new_directory_client(
|
||||
cfg: TopologyRefresherConfig,
|
||||
topology_accessor: TopologyAccessor<directory_client::Topology>,
|
||||
topology_accessor: TopologyAccessor,
|
||||
) -> Self {
|
||||
let directory_client_config = directory_client::Config::new(cfg.directory_server);
|
||||
let directory_client = directory_client::Client::new(directory_client_config);
|
||||
@@ -189,13 +189,16 @@ impl TopologyRefresher<directory_client::Topology> {
|
||||
}
|
||||
}
|
||||
|
||||
async fn get_current_compatible_topology(&self) -> Option<directory_client::Topology> {
|
||||
async fn get_current_compatible_topology(&self) -> Option<NymTopology> {
|
||||
match self.directory_client.get_topology().await {
|
||||
Err(err) => {
|
||||
error!("failed to get network topology! - {:?}", err);
|
||||
None
|
||||
}
|
||||
Ok(topology) => Some(topology.filter_system_version(built_info::PKG_VERSION)),
|
||||
Ok(topology) => {
|
||||
let nym_topology: NymTopology = topology.try_into().ok()?;
|
||||
Some(nym_topology.filter_system_version(built_info::PKG_VERSION))
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
@@ -22,10 +22,10 @@ use directory_client::DirectoryClient;
|
||||
use gateway_client::GatewayClient;
|
||||
use gateway_requests::registration::handshake::SharedKey;
|
||||
use pemstore::pemstore::PemStore;
|
||||
use std::convert::TryInto;
|
||||
use std::sync::Arc;
|
||||
use std::time::Duration;
|
||||
use topology::gateway::Node;
|
||||
use topology::NymTopology;
|
||||
use topology::{gateway, NymTopology};
|
||||
|
||||
pub fn command_args<'a, 'b>() -> clap::App<'a, 'b> {
|
||||
App::new("init")
|
||||
@@ -64,26 +64,14 @@ pub fn command_args<'a, 'b>() -> clap::App<'a, 'b> {
|
||||
}
|
||||
|
||||
async fn try_gateway_registration(
|
||||
gateways: Vec<Node>,
|
||||
gateways: &Vec<gateway::Node>,
|
||||
our_identity: Arc<identity::KeyPair>,
|
||||
) -> Option<(String, String, SharedKey)> {
|
||||
let timeout = Duration::from_millis(1500);
|
||||
for gateway in gateways {
|
||||
let gateway_identity =
|
||||
match identity::PublicKey::from_base58_string(gateway.identity_key.clone()) {
|
||||
Ok(id) => id,
|
||||
Err(_) => {
|
||||
eprintln!(
|
||||
"gateway {} announces invalid identity!",
|
||||
gateway.identity_key
|
||||
);
|
||||
continue;
|
||||
}
|
||||
};
|
||||
|
||||
let mut gateway_client = GatewayClient::new_init(
|
||||
url::Url::parse(&gateway.client_listener).unwrap(),
|
||||
gateway_identity,
|
||||
gateway.identity_key.clone(),
|
||||
our_identity.clone(),
|
||||
timeout,
|
||||
);
|
||||
@@ -93,7 +81,11 @@ async fn try_gateway_registration(
|
||||
eprintln!("Error while closing connection to the gateway! - {:?}", err);
|
||||
continue;
|
||||
} else {
|
||||
return Some((gateway.identity_key, gateway.client_listener, shared_key));
|
||||
return Some((
|
||||
gateway.identity_key.to_base58_string(),
|
||||
gateway.client_listener.clone(),
|
||||
shared_key,
|
||||
));
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -108,8 +100,9 @@ async fn choose_gateway(
|
||||
let directory_client_config = directory_client::Config::new(directory_server.clone());
|
||||
let directory_client = directory_client::Client::new(directory_client_config);
|
||||
let topology = directory_client.get_topology().await.unwrap();
|
||||
let nym_topology: NymTopology = topology.try_into().expect("Invalid topology data!");
|
||||
|
||||
let version_filtered_topology = topology.filter_system_version(built_info::PKG_VERSION);
|
||||
let version_filtered_topology = nym_topology.filter_system_version(built_info::PKG_VERSION);
|
||||
// don't care about health of the networks as mixes can go up and down any time,
|
||||
// but DO care about gateways
|
||||
let gateways = version_filtered_topology.gateways();
|
||||
@@ -138,11 +131,14 @@ async fn get_gateway_listener(directory_server: String, gateway_identity: &str)
|
||||
let directory_client_config = directory_client::Config::new(directory_server);
|
||||
let directory_client = directory_client::Client::new(directory_client_config);
|
||||
let topology = directory_client.get_topology().await.unwrap();
|
||||
let gateways = topology.gateways();
|
||||
|
||||
// technically we don't need to do conversion here, but let's be consistent
|
||||
let nym_topology: NymTopology = topology.try_into().ok()?;
|
||||
let gateways = nym_topology.gateways();
|
||||
|
||||
for gateway in gateways {
|
||||
if gateway.identity_key == gateway_identity {
|
||||
return Some(gateway.client_listener);
|
||||
if gateway.identity_key.to_base58_string() == gateway_identity {
|
||||
return Some(gateway.client_listener.clone());
|
||||
}
|
||||
}
|
||||
None
|
||||
|
||||
@@ -31,7 +31,6 @@ use tokio_tungstenite::{
|
||||
tungstenite::{protocol::Message, Error as WsError},
|
||||
WebSocketStream,
|
||||
};
|
||||
use topology::NymTopology;
|
||||
|
||||
enum ReceivedResponseType {
|
||||
Binary,
|
||||
@@ -44,17 +43,17 @@ impl Default for ReceivedResponseType {
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) struct Handler<T: NymTopology> {
|
||||
pub(crate) struct Handler {
|
||||
msg_input: InputMessageSender,
|
||||
buffer_requester: ReceivedBufferRequestSender,
|
||||
self_full_address: Recipient,
|
||||
topology_accessor: TopologyAccessor<T>,
|
||||
topology_accessor: TopologyAccessor,
|
||||
socket: Option<WebSocketStream<TcpStream>>,
|
||||
received_response_type: ReceivedResponseType,
|
||||
}
|
||||
|
||||
// clone is used to use handler on a new connection, which initially is `None`
|
||||
impl<T: NymTopology> Clone for Handler<T> {
|
||||
impl Clone for Handler {
|
||||
fn clone(&self) -> Self {
|
||||
Handler {
|
||||
msg_input: self.msg_input.clone(),
|
||||
@@ -67,7 +66,7 @@ impl<T: NymTopology> Clone for Handler<T> {
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: NymTopology> Drop for Handler<T> {
|
||||
impl Drop for Handler {
|
||||
fn drop(&mut self) {
|
||||
self.buffer_requester
|
||||
.unbounded_send(ReceivedBufferMessage::ReceiverDisconnect)
|
||||
@@ -75,12 +74,12 @@ impl<T: NymTopology> Drop for Handler<T> {
|
||||
}
|
||||
}
|
||||
|
||||
impl<T: NymTopology> Handler<T> {
|
||||
impl Handler {
|
||||
pub(crate) fn new(
|
||||
msg_input: InputMessageSender,
|
||||
buffer_requester: ReceivedBufferRequestSender,
|
||||
self_full_address: Recipient,
|
||||
topology_accessor: TopologyAccessor<T>,
|
||||
topology_accessor: TopologyAccessor,
|
||||
) -> Self {
|
||||
Handler {
|
||||
msg_input,
|
||||
|
||||
@@ -20,7 +20,6 @@ use std::{
|
||||
};
|
||||
use tokio::runtime;
|
||||
use tokio::{sync::Notify, task::JoinHandle};
|
||||
use topology::NymTopology;
|
||||
|
||||
enum State {
|
||||
Connected,
|
||||
@@ -50,7 +49,7 @@ impl Listener {
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) async fn run<T: NymTopology + 'static>(&mut self, handler: Handler<T>) {
|
||||
pub(crate) async fn run(&mut self, handler: Handler) {
|
||||
let mut tcp_listener = tokio::net::TcpListener::bind(self.address)
|
||||
.await
|
||||
.expect("Failed to start websocket listener");
|
||||
@@ -100,11 +99,7 @@ impl Listener {
|
||||
}
|
||||
}
|
||||
|
||||
pub(crate) fn start<T: NymTopology + 'static>(
|
||||
mut self,
|
||||
rt_handle: &runtime::Handle,
|
||||
handler: Handler<T>,
|
||||
) -> JoinHandle<()> {
|
||||
pub(crate) fn start(mut self, rt_handle: &runtime::Handle, handler: Handler) -> JoinHandle<()> {
|
||||
info!("Running websocket on {:?}", self.address.to_string());
|
||||
|
||||
rt_handle.spawn(async move { self.run(handler).await })
|
||||
|
||||
@@ -21,7 +21,7 @@ serde = { version = "1.0", features = ["derive"] }
|
||||
serde_json = "1.0"
|
||||
slice_as_array = "1.1.0"
|
||||
wasm-bindgen = "0.2"
|
||||
rand = "0.7.2"
|
||||
rand = {version = "0.7.3", features = ["wasm-bindgen"]}
|
||||
|
||||
# internal
|
||||
crypto = { path = "../../common/crypto" }
|
||||
|
||||
@@ -11,24 +11,27 @@
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
use crypto::asymmetric::encryption;
|
||||
pub use models::keys::keygen;
|
||||
use models::topology::Topology;
|
||||
use nymsphinx::addressing::clients::Recipient;
|
||||
use nymsphinx::addressing::nodes::NymNodeRoutingAddress;
|
||||
use nymsphinx::params::DEFAULT_NUM_MIX_HOPS;
|
||||
use nymsphinx::Node as SphinxNode;
|
||||
use nymsphinx::{delays, Destination, NodeAddressBytes, SphinxPacket};
|
||||
use rand::rngs::OsRng;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::convert::TryFrom;
|
||||
use std::convert::TryInto;
|
||||
use std::net::SocketAddr;
|
||||
use std::time::Duration;
|
||||
use topology::NymTopology;
|
||||
use wasm_bindgen::prelude::*;
|
||||
|
||||
mod models;
|
||||
mod utils;
|
||||
|
||||
use crypto::asymmetric::encryption;
|
||||
pub use models::keys::keygen;
|
||||
use nymsphinx::addressing::clients::Recipient;
|
||||
use topology::NymTopology;
|
||||
const DEFAULT_RNG: OsRng = OsRng;
|
||||
|
||||
// When the `wee_alloc` feature is enabled, use `wee_alloc` as the global
|
||||
// allocator.
|
||||
@@ -63,7 +66,8 @@ pub fn create_sphinx_packet(topology_json: &str, msg: &str, recipient: &str) ->
|
||||
|
||||
let recipient = Recipient::try_from_string(recipient).unwrap();
|
||||
|
||||
let route = sphinx_route_to(topology_json, &recipient.gateway());
|
||||
let route =
|
||||
sphinx_route_to(topology_json, &recipient.gateway()).expect("todo: error handling...");
|
||||
let average_delay = Duration::from_secs_f64(0.1);
|
||||
let delays = delays::generate_from_average_duration(route.len(), average_delay);
|
||||
|
||||
@@ -104,13 +108,17 @@ fn payload(sphinx_packet: SphinxPacket, route: Vec<SphinxNode>) -> Vec<u8> {
|
||||
///
|
||||
/// This function panics if the supplied `raw_route` json string can't be
|
||||
/// extracted to a `JsonRoute`.
|
||||
fn sphinx_route_to(topology_json: &str, gateway_address: &NodeAddressBytes) -> Vec<SphinxNode> {
|
||||
fn sphinx_route_to(
|
||||
topology_json: &str,
|
||||
gateway_address: &NodeAddressBytes,
|
||||
) -> Option<Vec<SphinxNode>> {
|
||||
let topology = Topology::new(topology_json);
|
||||
let route = topology
|
||||
.random_route_to_gateway(gateway_address)
|
||||
let nym_topology: NymTopology = topology.try_into().ok()?;
|
||||
let route = nym_topology
|
||||
.random_route_to_gateway(&mut DEFAULT_RNG, DEFAULT_NUM_MIX_HOPS, gateway_address)
|
||||
.expect("invalid route produced");
|
||||
assert_eq!(4, route.len());
|
||||
route
|
||||
Some(route)
|
||||
}
|
||||
|
||||
impl TryFrom<NodeData> for SphinxNode {
|
||||
@@ -145,23 +153,6 @@ mod test_constructing_a_sphinx_packet {
|
||||
// );
|
||||
// assert_eq!(1404, packet.len());
|
||||
// }
|
||||
|
||||
#[test]
|
||||
#[cfg_attr(feature = "offline-test", ignore)]
|
||||
fn starts_with_a_mix_address() {
|
||||
let mut payload = create_sphinx_packet(
|
||||
topology_fixture(),
|
||||
"foomp",
|
||||
"5pgrc4gPHP2tBQgfezcdJ2ZAjipoAsy6evrqHdxBbVXq@CdqJCedY5d1geJNDjUqnEx8zF7mKjb6PCZ6k3T6xhxD",
|
||||
);
|
||||
// you don't really need 32 bytes here, but giving too much won't make it fail
|
||||
let mut address_buffer = [0; 32];
|
||||
let _ = payload.split_off(32);
|
||||
address_buffer.copy_from_slice(payload.as_slice());
|
||||
let address = NymNodeRoutingAddress::try_from_bytes(&address_buffer);
|
||||
|
||||
assert!(address.is_ok());
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
@@ -174,10 +165,11 @@ mod building_a_topology_from_json {
|
||||
sphinx_route_to(
|
||||
"",
|
||||
&NodeAddressBytes::try_from_base58_string(
|
||||
"CdqJCedY5d1geJNDjUqnEx8zF7mKjb6PCZ6k3T6xhxD",
|
||||
"FE7zC2sJZrhXgQWvzXXVH8GHi2xXRynX8UWK8rD8ikf3",
|
||||
)
|
||||
.unwrap(),
|
||||
);
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -186,10 +178,11 @@ mod building_a_topology_from_json {
|
||||
sphinx_route_to(
|
||||
"bad bad bad not json",
|
||||
&NodeAddressBytes::try_from_base58_string(
|
||||
"CdqJCedY5d1geJNDjUqnEx8zF7mKjb6PCZ6k3T6xhxD",
|
||||
"FE7zC2sJZrhXgQWvzXXVH8GHi2xXRynX8UWK8rD8ikf3",
|
||||
)
|
||||
.unwrap(),
|
||||
);
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -201,10 +194,11 @@ mod building_a_topology_from_json {
|
||||
sphinx_route_to(
|
||||
&json,
|
||||
&NodeAddressBytes::try_from_base58_string(
|
||||
"CdqJCedY5d1geJNDjUqnEx8zF7mKjb6PCZ6k3T6xhxD",
|
||||
"FE7zC2sJZrhXgQWvzXXVH8GHi2xXRynX8UWK8rD8ikf3",
|
||||
)
|
||||
.unwrap(),
|
||||
);
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
#[test]
|
||||
@@ -217,23 +211,24 @@ mod building_a_topology_from_json {
|
||||
sphinx_route_to(
|
||||
&json,
|
||||
&NodeAddressBytes::try_from_base58_string(
|
||||
"CdqJCedY5d1geJNDjUqnEx8zF7mKjb6PCZ6k3T6xhxD",
|
||||
"FE7zC2sJZrhXgQWvzXXVH8GHi2xXRynX8UWK8rD8ikf3",
|
||||
)
|
||||
.unwrap(),
|
||||
);
|
||||
)
|
||||
.unwrap();
|
||||
}
|
||||
|
||||
// JS: why is this an "offline-test" feature? It makes no network requests?
|
||||
#[test]
|
||||
#[cfg_attr(feature = "offline-test", ignore)]
|
||||
fn test_works_on_happy_json() {
|
||||
let route = sphinx_route_to(
|
||||
topology_fixture(),
|
||||
&NodeAddressBytes::try_from_base58_string(
|
||||
"CdqJCedY5d1geJNDjUqnEx8zF7mKjb6PCZ6k3T6xhxD",
|
||||
"FE7zC2sJZrhXgQWvzXXVH8GHi2xXRynX8UWK8rD8ikf3",
|
||||
)
|
||||
.unwrap(),
|
||||
);
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(4, route.len());
|
||||
}
|
||||
|
||||
@@ -244,14 +239,25 @@ mod building_a_topology_from_json {
|
||||
let route = sphinx_route_to(
|
||||
&json,
|
||||
&NodeAddressBytes::try_from_base58_string(
|
||||
"CdqJCedY5d1geJNDjUqnEx8zF7mKjb6PCZ6k3T6xhxD",
|
||||
"FE7zC2sJZrhXgQWvzXXVH8GHi2xXRynX8UWK8rD8ikf3",
|
||||
)
|
||||
.unwrap(),
|
||||
);
|
||||
)
|
||||
.unwrap();
|
||||
assert_eq!(4, route.len());
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod topology_fixture {
|
||||
use super::*;
|
||||
|
||||
#[test]
|
||||
fn is_valid() {
|
||||
let _nym_topology: NymTopology = Topology::new(topology_fixture()).try_into().unwrap();
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
fn topology_fixture() -> &'static str {
|
||||
r#"
|
||||
@@ -312,7 +318,7 @@ fn topology_fixture() -> &'static str {
|
||||
{
|
||||
"clientListener": "139.162.246.48:9000",
|
||||
"mixnetListener": "139.162.246.48:1789",
|
||||
"identityKey": "CdqJCedY5d1geJNDjUqnEx8zF7mKjb6PCZ6k3T6xhxD",
|
||||
"identityKey": "FE7zC2sJZrhXgQWvzXXVH8GHi2xXRynX8UWK8rD8ikf3",
|
||||
"sphinxKey": "BnLYqQjb8K6TmW5oFdNZrUTocGxa3rgzBvapQrf8XUbF",
|
||||
"version": "0.6.0",
|
||||
"location": "London, UK",
|
||||
@@ -326,7 +332,7 @@ fn topology_fixture() -> &'static str {
|
||||
{
|
||||
"clientListener": "127.0.0.1:9000",
|
||||
"mixnetListener": "127.0.0.1:1789",
|
||||
"identityKey": "B9xz9V6jpp1fEbDkeyR5f8miorw9bzXGKoMbKnaxkD41",
|
||||
"identityKey": "7hU4RNHtGC1FreLYLoBXXTH8AdaqU913NbqCv5Fu4z1r",
|
||||
"sphinxKey": "3KCpz1HCD8DqnQjemT1uuBZipmHFXM4V5btxLXwvM1gG",
|
||||
"version": "0.6.0",
|
||||
"location": "unknown",
|
||||
|
||||
@@ -43,7 +43,7 @@ impl TryInto<String> for GatewayIdentity {
|
||||
#[wasm_bindgen]
|
||||
pub fn keygen() -> String {
|
||||
let keypair = identity::KeyPair::new();
|
||||
let address = keypair.public_key().derive_address();
|
||||
let address = keypair.public_key().derive_destination_address();
|
||||
|
||||
GatewayIdentity {
|
||||
private_key: keypair.private_key().to_base58_string(),
|
||||
|
||||
@@ -13,7 +13,8 @@
|
||||
// limitations under the License.
|
||||
|
||||
use serde::Serializer;
|
||||
use topology::{coco, gateway, mix, provider, NymTopology};
|
||||
use std::convert::TryInto;
|
||||
use topology::NymTopology;
|
||||
|
||||
#[derive(Clone, Debug)]
|
||||
pub struct Topology {
|
||||
@@ -56,35 +57,10 @@ impl Topology {
|
||||
}
|
||||
}
|
||||
|
||||
impl NymTopology for Topology {
|
||||
fn new_from_nodes(
|
||||
mix_nodes: Vec<mix::Node>,
|
||||
mix_provider_nodes: Vec<provider::Node>,
|
||||
coco_nodes: Vec<coco::Node>,
|
||||
gateway_nodes: Vec<gateway::Node>,
|
||||
) -> Self {
|
||||
Topology {
|
||||
inner: directory_client_models::presence::Topology::new_from_nodes(
|
||||
mix_nodes,
|
||||
mix_provider_nodes,
|
||||
coco_nodes,
|
||||
gateway_nodes,
|
||||
),
|
||||
}
|
||||
}
|
||||
fn mix_nodes(&self) -> Vec<mix::Node> {
|
||||
self.inner.mix_nodes()
|
||||
}
|
||||
impl TryInto<NymTopology> for Topology {
|
||||
type Error = directory_client_models::presence::TopologyConversionError;
|
||||
|
||||
fn providers(&self) -> Vec<provider::Node> {
|
||||
self.inner.providers()
|
||||
}
|
||||
|
||||
fn gateways(&self) -> Vec<gateway::Node> {
|
||||
self.inner.gateways()
|
||||
}
|
||||
|
||||
fn coco_nodes(&self) -> Vec<topology::coco::Node> {
|
||||
self.inner.coco_nodes()
|
||||
fn try_into(self) -> Result<NymTopology, Self::Error> {
|
||||
self.inner.try_into()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -9,4 +9,5 @@ edition = "2018"
|
||||
[dependencies]
|
||||
serde = { version = "1.0.104", features = ["derive"] }
|
||||
|
||||
crypto = { path = "../../../crypto" }
|
||||
topology = { path = "../../../topology" }
|
||||
|
||||
@@ -12,8 +12,20 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use crypto::asymmetric::identity;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use topology::coco;
|
||||
use std::convert::TryInto;
|
||||
|
||||
#[derive(Debug)]
|
||||
pub enum ConversionError {
|
||||
InvalidKeyError,
|
||||
}
|
||||
|
||||
impl From<identity::SignatureError> for ConversionError {
|
||||
fn from(_: identity::SignatureError) -> Self {
|
||||
ConversionError::InvalidKeyError
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Deserialize, Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
@@ -25,26 +37,16 @@ pub struct CocoPresence {
|
||||
pub version: String,
|
||||
}
|
||||
|
||||
impl Into<topology::coco::Node> for CocoPresence {
|
||||
fn into(self) -> topology::coco::Node {
|
||||
topology::coco::Node {
|
||||
impl TryInto<topology::coco::Node> for CocoPresence {
|
||||
type Error = ConversionError;
|
||||
|
||||
fn try_into(self) -> Result<topology::coco::Node, Self::Error> {
|
||||
Ok(topology::coco::Node {
|
||||
location: self.location,
|
||||
host: self.host,
|
||||
pub_key: self.pub_key,
|
||||
pub_key: identity::PublicKey::from_base58_string(self.pub_key)?,
|
||||
last_seen: self.last_seen,
|
||||
version: self.version,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<topology::coco::Node> for CocoPresence {
|
||||
fn from(cn: coco::Node) -> Self {
|
||||
CocoPresence {
|
||||
location: cn.location,
|
||||
host: cn.host,
|
||||
pub_key: cn.pub_key,
|
||||
last_seen: cn.last_seen,
|
||||
version: cn.version,
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,8 +12,35 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use crypto::asymmetric::{encryption, identity};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use topology::gateway;
|
||||
use std::convert::TryInto;
|
||||
use std::io;
|
||||
use std::net::ToSocketAddrs;
|
||||
|
||||
#[derive(Debug)]
|
||||
pub enum ConversionError {
|
||||
InvalidKeyError,
|
||||
InvalidAddress(io::Error),
|
||||
}
|
||||
|
||||
impl From<identity::SignatureError> for ConversionError {
|
||||
fn from(_: identity::SignatureError) -> Self {
|
||||
ConversionError::InvalidKeyError
|
||||
}
|
||||
}
|
||||
|
||||
impl From<encryption::EncryptionKeyError> for ConversionError {
|
||||
fn from(_: encryption::EncryptionKeyError) -> Self {
|
||||
ConversionError::InvalidKeyError
|
||||
}
|
||||
}
|
||||
|
||||
impl From<io::Error> for ConversionError {
|
||||
fn from(err: io::Error) -> Self {
|
||||
ConversionError::InvalidAddress(err)
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Deserialize, Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
@@ -28,14 +55,24 @@ pub struct GatewayPresence {
|
||||
pub version: String,
|
||||
}
|
||||
|
||||
impl Into<topology::gateway::Node> for GatewayPresence {
|
||||
fn into(self) -> topology::gateway::Node {
|
||||
topology::gateway::Node {
|
||||
impl TryInto<topology::gateway::Node> for GatewayPresence {
|
||||
type Error = ConversionError;
|
||||
|
||||
fn try_into(self) -> Result<topology::gateway::Node, Self::Error> {
|
||||
let resolved_mix_hostname = self.mixnet_listener.to_socket_addrs()?.next();
|
||||
if resolved_mix_hostname.is_none() {
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::Other,
|
||||
"no valid socket address",
|
||||
))?;
|
||||
}
|
||||
|
||||
Ok(topology::gateway::Node {
|
||||
location: self.location,
|
||||
client_listener: self.client_listener.parse().unwrap(),
|
||||
mixnet_listener: self.mixnet_listener.parse().unwrap(),
|
||||
identity_key: self.identity_key,
|
||||
sphinx_key: self.sphinx_key,
|
||||
client_listener: self.client_listener,
|
||||
mixnet_listener: resolved_mix_hostname.unwrap(),
|
||||
identity_key: identity::PublicKey::from_base58_string(self.identity_key)?,
|
||||
sphinx_key: encryption::PublicKey::from_base58_string(self.sphinx_key)?,
|
||||
registered_clients: self
|
||||
.registered_clients
|
||||
.into_iter()
|
||||
@@ -43,26 +80,7 @@ impl Into<topology::gateway::Node> for GatewayPresence {
|
||||
.collect(),
|
||||
last_seen: self.last_seen,
|
||||
version: self.version,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<topology::gateway::Node> for GatewayPresence {
|
||||
fn from(mpn: gateway::Node) -> Self {
|
||||
GatewayPresence {
|
||||
location: mpn.location,
|
||||
client_listener: mpn.client_listener.to_string(),
|
||||
mixnet_listener: mpn.mixnet_listener.to_string(),
|
||||
identity_key: mpn.identity_key,
|
||||
sphinx_key: mpn.sphinx_key,
|
||||
registered_clients: mpn
|
||||
.registered_clients
|
||||
.into_iter()
|
||||
.map(|c| c.into())
|
||||
.collect(),
|
||||
last_seen: mpn.last_seen,
|
||||
version: mpn.version,
|
||||
}
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
@@ -79,11 +97,3 @@ impl Into<topology::gateway::Client> for GatewayClient {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<topology::gateway::Client> for GatewayClient {
|
||||
fn from(mpc: topology::gateway::Client) -> Self {
|
||||
GatewayClient {
|
||||
pub_key: mpc.pub_key,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -12,11 +12,29 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use crypto::asymmetric::encryption;
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::convert::TryInto;
|
||||
use std::io;
|
||||
use std::net::ToSocketAddrs;
|
||||
use topology::mix;
|
||||
|
||||
#[derive(Debug)]
|
||||
pub enum ConversionError {
|
||||
InvalidKeyError,
|
||||
InvalidAddress(io::Error),
|
||||
}
|
||||
|
||||
impl From<encryption::EncryptionKeyError> for ConversionError {
|
||||
fn from(_: encryption::EncryptionKeyError) -> Self {
|
||||
ConversionError::InvalidKeyError
|
||||
}
|
||||
}
|
||||
|
||||
impl From<io::Error> for ConversionError {
|
||||
fn from(err: io::Error) -> Self {
|
||||
ConversionError::InvalidAddress(err)
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Deserialize, Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
@@ -30,7 +48,7 @@ pub struct MixNodePresence {
|
||||
}
|
||||
|
||||
impl TryInto<topology::mix::Node> for MixNodePresence {
|
||||
type Error = io::Error;
|
||||
type Error = ConversionError;
|
||||
|
||||
fn try_into(self) -> Result<topology::mix::Node, Self::Error> {
|
||||
let resolved_hostname = self.host.to_socket_addrs()?.next();
|
||||
@@ -38,29 +56,16 @@ impl TryInto<topology::mix::Node> for MixNodePresence {
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::Other,
|
||||
"no valid socket address",
|
||||
));
|
||||
))?;
|
||||
}
|
||||
|
||||
Ok(topology::mix::Node {
|
||||
location: self.location,
|
||||
host: resolved_hostname.unwrap(),
|
||||
pub_key: self.pub_key,
|
||||
pub_key: encryption::PublicKey::from_base58_string(self.pub_key)?,
|
||||
layer: self.layer,
|
||||
last_seen: self.last_seen,
|
||||
version: self.version,
|
||||
})
|
||||
}
|
||||
}
|
||||
|
||||
impl From<topology::mix::Node> for MixNodePresence {
|
||||
fn from(mn: mix::Node) -> Self {
|
||||
MixNodePresence {
|
||||
location: mn.location,
|
||||
host: mn.host.to_string(),
|
||||
pub_key: mn.pub_key,
|
||||
layer: mn.layer,
|
||||
last_seen: mn.last_seen,
|
||||
version: mn.version,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -18,4 +18,4 @@ pub mod mixnodes;
|
||||
pub mod providers;
|
||||
pub mod topology;
|
||||
|
||||
pub use self::topology::Topology;
|
||||
pub use self::topology::{Topology, TopologyConversionError};
|
||||
|
||||
@@ -13,7 +13,6 @@
|
||||
// limitations under the License.
|
||||
|
||||
use serde::{Deserialize, Serialize};
|
||||
use topology::provider;
|
||||
|
||||
#[derive(Clone, Debug, Deserialize, Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
@@ -27,60 +26,8 @@ pub struct MixProviderPresence {
|
||||
pub version: String,
|
||||
}
|
||||
|
||||
impl Into<topology::provider::Node> for MixProviderPresence {
|
||||
fn into(self) -> topology::provider::Node {
|
||||
topology::provider::Node {
|
||||
location: self.location,
|
||||
client_listener: self.client_listener.parse().unwrap(),
|
||||
mixnet_listener: self.mixnet_listener.parse().unwrap(),
|
||||
pub_key: self.pub_key,
|
||||
registered_clients: self
|
||||
.registered_clients
|
||||
.into_iter()
|
||||
.map(|c| c.into())
|
||||
.collect(),
|
||||
last_seen: self.last_seen,
|
||||
version: self.version,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<topology::provider::Node> for MixProviderPresence {
|
||||
fn from(mpn: provider::Node) -> Self {
|
||||
MixProviderPresence {
|
||||
location: mpn.location,
|
||||
client_listener: mpn.client_listener.to_string(),
|
||||
mixnet_listener: mpn.mixnet_listener.to_string(),
|
||||
pub_key: mpn.pub_key,
|
||||
registered_clients: mpn
|
||||
.registered_clients
|
||||
.into_iter()
|
||||
.map(|c| c.into())
|
||||
.collect(),
|
||||
last_seen: mpn.last_seen,
|
||||
version: mpn.version,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Deserialize, Serialize)]
|
||||
#[serde(rename_all = "camelCase")]
|
||||
pub struct MixProviderClient {
|
||||
pub pub_key: String,
|
||||
}
|
||||
|
||||
impl Into<topology::provider::Client> for MixProviderClient {
|
||||
fn into(self) -> topology::provider::Client {
|
||||
topology::provider::Client {
|
||||
pub_key: self.pub_key,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
impl From<topology::provider::Client> for MixProviderClient {
|
||||
fn from(mpc: topology::provider::Client) -> Self {
|
||||
MixProviderClient {
|
||||
pub_key: mpc.pub_key,
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -15,7 +15,32 @@
|
||||
use super::{coconodes, gateways, mixnodes, providers};
|
||||
use serde::{Deserialize, Serialize};
|
||||
use std::convert::TryInto;
|
||||
use topology::{coco, gateway, mix, provider, NymTopology};
|
||||
use topology::{MixLayer, NymTopology};
|
||||
|
||||
#[derive(Debug)]
|
||||
pub enum TopologyConversionError {
|
||||
CocoError(self::coconodes::ConversionError),
|
||||
GatewayError(self::gateways::ConversionError),
|
||||
MixError(self::mixnodes::ConversionError),
|
||||
}
|
||||
|
||||
impl From<self::coconodes::ConversionError> for TopologyConversionError {
|
||||
fn from(err: self::coconodes::ConversionError) -> Self {
|
||||
TopologyConversionError::CocoError(err)
|
||||
}
|
||||
}
|
||||
|
||||
impl From<self::gateways::ConversionError> for TopologyConversionError {
|
||||
fn from(err: self::gateways::ConversionError) -> Self {
|
||||
TopologyConversionError::GatewayError(err)
|
||||
}
|
||||
}
|
||||
|
||||
impl From<self::mixnodes::ConversionError> for TopologyConversionError {
|
||||
fn from(err: self::mixnodes::ConversionError) -> Self {
|
||||
TopologyConversionError::MixError(err)
|
||||
}
|
||||
}
|
||||
|
||||
// Topology shows us the current state of the overall Nym network
|
||||
#[derive(Clone, Debug, Deserialize, Serialize)]
|
||||
@@ -27,47 +52,30 @@ pub struct Topology {
|
||||
pub gateway_nodes: Vec<gateways::GatewayPresence>,
|
||||
}
|
||||
|
||||
impl NymTopology for Topology {
|
||||
fn new_from_nodes(
|
||||
mix_nodes: Vec<mix::Node>,
|
||||
mix_provider_nodes: Vec<provider::Node>,
|
||||
coco_nodes: Vec<coco::Node>,
|
||||
gateway_nodes: Vec<gateway::Node>,
|
||||
) -> Self {
|
||||
Topology {
|
||||
coco_nodes: coco_nodes.into_iter().map(|node| node.into()).collect(),
|
||||
mix_nodes: mix_nodes.into_iter().map(|node| node.into()).collect(),
|
||||
mix_provider_nodes: mix_provider_nodes
|
||||
.into_iter()
|
||||
.map(|node| node.into())
|
||||
.collect(),
|
||||
gateway_nodes: gateway_nodes.into_iter().map(|node| node.into()).collect(),
|
||||
impl TryInto<NymTopology> for Topology {
|
||||
type Error = TopologyConversionError;
|
||||
|
||||
fn try_into(self) -> Result<NymTopology, TopologyConversionError> {
|
||||
use std::collections::HashMap;
|
||||
|
||||
let mut coco_nodes = Vec::with_capacity(self.coco_nodes.len());
|
||||
for coco in self.coco_nodes.into_iter() {
|
||||
coco_nodes.push(coco.try_into()?)
|
||||
}
|
||||
}
|
||||
|
||||
fn mix_nodes(&self) -> Vec<mix::Node> {
|
||||
self.mix_nodes
|
||||
.iter()
|
||||
.filter_map(|x| x.clone().try_into().ok())
|
||||
.collect()
|
||||
}
|
||||
let mut mixes = HashMap::new();
|
||||
for mix in self.mix_nodes.into_iter() {
|
||||
let layer = mix.layer as MixLayer;
|
||||
let layer_entry = mixes.entry(layer).or_insert(Vec::new());
|
||||
layer_entry.push(mix.try_into()?)
|
||||
}
|
||||
|
||||
fn providers(&self) -> Vec<provider::Node> {
|
||||
self.mix_provider_nodes
|
||||
.iter()
|
||||
.map(|x| x.clone().into())
|
||||
.collect()
|
||||
}
|
||||
let mut gateways = Vec::with_capacity(self.gateway_nodes.len());
|
||||
for gate in self.gateway_nodes.into_iter() {
|
||||
gateways.push(gate.try_into()?)
|
||||
}
|
||||
|
||||
fn gateways(&self) -> Vec<gateway::Node> {
|
||||
self.gateway_nodes
|
||||
.iter()
|
||||
.map(|x| x.clone().into())
|
||||
.collect()
|
||||
}
|
||||
|
||||
fn coco_nodes(&self) -> Vec<topology::coco::Node> {
|
||||
self.coco_nodes.iter().map(|x| x.clone().into()).collect()
|
||||
Ok(NymTopology::new(coco_nodes, mixes, gateways))
|
||||
}
|
||||
}
|
||||
|
||||
@@ -77,18 +85,20 @@ mod converting_mixnode_presence_into_topology_mixnode {
|
||||
|
||||
#[test]
|
||||
fn it_returns_error_on_unresolvable_hostname() {
|
||||
use topology::mix;
|
||||
|
||||
let unresolvable_hostname = "foomp.foomp.foomp:1234";
|
||||
|
||||
let mix_presence = mixnodes::MixNodePresence {
|
||||
location: "".to_string(),
|
||||
host: unresolvable_hostname.to_string(),
|
||||
pub_key: "".to_string(),
|
||||
pub_key: "BnLYqQjb8K6TmW5oFdNZrUTocGxa3rgzBvapQrf8XUbF".to_string(),
|
||||
layer: 0,
|
||||
last_seen: 0,
|
||||
version: "".to_string(),
|
||||
};
|
||||
|
||||
let result: Result<mix::Node, std::io::Error> = mix_presence.try_into();
|
||||
let result: Result<mix::Node, self::mixnodes::ConversionError> = mix_presence.try_into();
|
||||
assert!(result.is_err()) // This fails only for me. Why?
|
||||
// ¯\_(ツ)_/¯ - works on my machine (and travis)
|
||||
// Is it still broken?
|
||||
@@ -102,13 +112,15 @@ mod converting_mixnode_presence_into_topology_mixnode {
|
||||
let mix_presence = mixnodes::MixNodePresence {
|
||||
location: "".to_string(),
|
||||
host: resolvable_hostname.to_string(),
|
||||
pub_key: "".to_string(),
|
||||
pub_key: "BnLYqQjb8K6TmW5oFdNZrUTocGxa3rgzBvapQrf8XUbF".to_string(),
|
||||
layer: 0,
|
||||
last_seen: 0,
|
||||
version: "".to_string(),
|
||||
};
|
||||
|
||||
let result: Result<topology::mix::Node, std::io::Error> = mix_presence.try_into();
|
||||
assert!(result.is_ok())
|
||||
let result: Result<topology::mix::Node, self::mixnodes::ConversionError> =
|
||||
mix_presence.try_into();
|
||||
result.unwrap();
|
||||
// assert!(result.is_ok())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -376,7 +376,11 @@ impl<'a, R> GatewayClient<'static, R> {
|
||||
.as_ref()
|
||||
.unwrap_or_else(|| self.shared_key.as_ref().unwrap());
|
||||
let iv = AuthenticationIV::new_random(&mut DEFAULT_RNG);
|
||||
let self_address = self.local_identity.as_ref().public_key().derive_address();
|
||||
let self_address = self
|
||||
.local_identity
|
||||
.as_ref()
|
||||
.public_key()
|
||||
.derive_destination_address();
|
||||
let encrypted_address = EncryptedAddressBytes::new(&self_address, shared_key, &iv);
|
||||
|
||||
let msg =
|
||||
|
||||
@@ -201,6 +201,12 @@ impl Into<nymsphinx_types::PublicKey> for PublicKey {
|
||||
}
|
||||
}
|
||||
|
||||
impl<'a> Into<nymsphinx_types::PublicKey> for &'a PublicKey {
|
||||
fn into(self) -> nymsphinx_types::PublicKey {
|
||||
nymsphinx_types::PublicKey::from(self.to_bytes())
|
||||
}
|
||||
}
|
||||
|
||||
impl From<nymsphinx_types::PublicKey> for PublicKey {
|
||||
fn from(pub_key: nymsphinx_types::PublicKey) -> Self {
|
||||
Self(x25519_dalek::PublicKey::from(*pub_key.as_bytes()))
|
||||
|
||||
@@ -14,9 +14,11 @@
|
||||
|
||||
use crate::{PemStorableKey, PemStorableKeyPair};
|
||||
use bs58;
|
||||
use ed25519_dalek::SignatureError;
|
||||
pub use ed25519_dalek::SignatureError;
|
||||
pub use ed25519_dalek::{PUBLIC_KEY_LENGTH, SECRET_KEY_LENGTH, SIGNATURE_LENGTH};
|
||||
use nymsphinx_types::{DestinationAddressBytes, DESTINATION_ADDRESS_LENGTH};
|
||||
use nymsphinx_types::{
|
||||
DestinationAddressBytes, NodeAddressBytes, DESTINATION_ADDRESS_LENGTH, NODE_ADDRESS_LENGTH,
|
||||
};
|
||||
use rand::{rngs::OsRng, CryptoRng, RngCore};
|
||||
|
||||
/// Keypair for usage in ed25519 EdDSA.
|
||||
@@ -79,7 +81,7 @@ impl PemStorableKeyPair for KeyPair {
|
||||
pub struct PublicKey(ed25519_dalek::PublicKey);
|
||||
|
||||
impl PublicKey {
|
||||
pub fn derive_address(&self) -> DestinationAddressBytes {
|
||||
pub fn derive_destination_address(&self) -> DestinationAddressBytes {
|
||||
let mut temporary_address = [0u8; DESTINATION_ADDRESS_LENGTH];
|
||||
let public_key_bytes = self.to_bytes();
|
||||
|
||||
@@ -89,6 +91,16 @@ impl PublicKey {
|
||||
DestinationAddressBytes::from_bytes(temporary_address)
|
||||
}
|
||||
|
||||
pub fn derive_node_address(&self) -> NodeAddressBytes {
|
||||
let mut temporary_address = [0u8; NODE_ADDRESS_LENGTH];
|
||||
let public_key_bytes = self.to_bytes();
|
||||
|
||||
assert_eq!(NODE_ADDRESS_LENGTH, PUBLIC_KEY_LENGTH);
|
||||
|
||||
temporary_address.copy_from_slice(&public_key_bytes[..]);
|
||||
NodeAddressBytes::from_bytes(temporary_address)
|
||||
}
|
||||
|
||||
/// Convert this public key to a byte array.
|
||||
pub fn to_bytes(&self) -> [u8; PUBLIC_KEY_LENGTH] {
|
||||
self.0.to_bytes()
|
||||
|
||||
@@ -16,6 +16,7 @@ use crate::identifier::{prepare_identifier, AckAes128Key};
|
||||
use nymsphinx_addressing::clients::Recipient;
|
||||
use nymsphinx_addressing::nodes::{NymNodeRoutingAddress, MAX_NODE_ADDRESS_UNPADDED_LEN};
|
||||
use nymsphinx_params::packet_sizes::PacketSize;
|
||||
use nymsphinx_params::DEFAULT_NUM_MIX_HOPS;
|
||||
use nymsphinx_types::builder::SphinxPacketBuilder;
|
||||
use nymsphinx_types::{
|
||||
delays::{self, Delay},
|
||||
@@ -41,19 +42,19 @@ pub enum SURBAckRecoveryError {
|
||||
}
|
||||
|
||||
impl SURBAck {
|
||||
pub fn construct<R, T>(
|
||||
pub fn construct<R>(
|
||||
rng: &mut R,
|
||||
recipient: &Recipient,
|
||||
ack_key: &AckAes128Key,
|
||||
marshaled_fragment_id: [u8; 5],
|
||||
average_delay: time::Duration,
|
||||
topology: &T,
|
||||
topology: &NymTopology,
|
||||
) -> Result<Self, NymTopologyError>
|
||||
where
|
||||
R: RngCore + CryptoRng,
|
||||
T: NymTopology,
|
||||
{
|
||||
let route = topology.random_route_to_gateway(&recipient.gateway())?;
|
||||
let route =
|
||||
topology.random_route_to_gateway(rng, DEFAULT_NUM_MIX_HOPS, &recipient.gateway())?;
|
||||
let delays = delays::generate_from_average_duration(route.len(), average_delay);
|
||||
let destination = Destination::new(recipient.destination(), Default::default());
|
||||
|
||||
|
||||
@@ -22,6 +22,7 @@ use nymsphinx_acknowledgements::surb_ack::SURBAck;
|
||||
use nymsphinx_addressing::clients::Recipient;
|
||||
use nymsphinx_addressing::nodes::{NymNodeRoutingAddress, MAX_NODE_ADDRESS_UNPADDED_LEN};
|
||||
use nymsphinx_params::packet_sizes::PacketSize;
|
||||
use nymsphinx_params::DEFAULT_NUM_MIX_HOPS;
|
||||
use nymsphinx_types::builder::SphinxPacketBuilder;
|
||||
use nymsphinx_types::{delays, Delay, Destination, SphinxPacket};
|
||||
use rand::{rngs::OsRng, CryptoRng, Rng};
|
||||
@@ -134,7 +135,10 @@ impl MessageChunker<DefaultRng> {
|
||||
}
|
||||
}
|
||||
|
||||
impl<R: CryptoRng + Rng> MessageChunker<R> {
|
||||
impl<R> MessageChunker<R>
|
||||
where
|
||||
R: CryptoRng + Rng,
|
||||
{
|
||||
pub fn new_with_rng(
|
||||
rng: R,
|
||||
ack_recipient: Recipient,
|
||||
@@ -186,10 +190,10 @@ impl<R: CryptoRng + Rng> MessageChunker<R> {
|
||||
/// such that it contains required SURB-ACK.
|
||||
/// This method can fail if the provided network topology is invalid.
|
||||
/// It returns total expected delay as well as the `SphinxPacket` to be sent through the network.
|
||||
pub fn prepare_chunk_for_sending<T: NymTopology>(
|
||||
pub fn prepare_chunk_for_sending(
|
||||
&mut self,
|
||||
fragment: Fragment,
|
||||
topology: &T,
|
||||
topology: &NymTopology,
|
||||
ack_key: &AckAes128Key,
|
||||
packet_recipient: &Recipient,
|
||||
) -> Result<(Delay, (SocketAddr, SphinxPacket)), NymTopologyError> {
|
||||
@@ -203,7 +207,11 @@ impl<R: CryptoRng + Rng> MessageChunker<R> {
|
||||
.chain(fragment.into_bytes().into_iter())
|
||||
.collect();
|
||||
|
||||
let route = topology.random_route_to_gateway(&packet_recipient.gateway())?;
|
||||
let route = topology.random_route_to_gateway(
|
||||
&mut self.rng,
|
||||
DEFAULT_NUM_MIX_HOPS,
|
||||
&packet_recipient.gateway(),
|
||||
)?;
|
||||
let delays =
|
||||
delays::generate_from_average_duration(route.len(), self.average_packet_delay_duration);
|
||||
let destination = Destination::new(packet_recipient.destination(), Default::default());
|
||||
@@ -223,15 +231,12 @@ impl<R: CryptoRng + Rng> MessageChunker<R> {
|
||||
))
|
||||
}
|
||||
|
||||
fn generate_surb_ack<T>(
|
||||
fn generate_surb_ack(
|
||||
&mut self,
|
||||
fragment_id: &FragmentIdentifier,
|
||||
topology: &T,
|
||||
topology: &NymTopology,
|
||||
ack_key: &AckAes128Key,
|
||||
) -> Result<SURBAck, NymTopologyError>
|
||||
where
|
||||
T: NymTopology,
|
||||
{
|
||||
) -> Result<SURBAck, NymTopologyError> {
|
||||
SURBAck::construct(
|
||||
&mut self.rng,
|
||||
&self.ack_recipient,
|
||||
|
||||
@@ -18,6 +18,7 @@ use nymsphinx_addressing::clients::Recipient;
|
||||
use nymsphinx_addressing::nodes::{NymNodeRoutingAddress, NymNodeRoutingAddressError};
|
||||
use nymsphinx_chunking::fragment::COVER_FRAG_ID;
|
||||
use nymsphinx_params::packet_sizes::PacketSize;
|
||||
use nymsphinx_params::DEFAULT_NUM_MIX_HOPS;
|
||||
use nymsphinx_types::builder::SphinxPacketBuilder;
|
||||
use nymsphinx_types::{delays, Destination, Error as SphinxError, SphinxPacket};
|
||||
use rand::{CryptoRng, RngCore};
|
||||
@@ -55,16 +56,15 @@ impl From<NymTopologyError> for CoverMessageError {
|
||||
}
|
||||
}
|
||||
|
||||
pub fn generate_loop_cover_surb_ack<R, T>(
|
||||
pub fn generate_loop_cover_surb_ack<R>(
|
||||
rng: &mut R,
|
||||
topology: &T,
|
||||
topology: &NymTopology,
|
||||
ack_key: &AckAes128Key,
|
||||
full_address: &Recipient,
|
||||
average_ack_delay: time::Duration,
|
||||
) -> Result<SURBAck, CoverMessageError>
|
||||
where
|
||||
R: RngCore + CryptoRng,
|
||||
T: NymTopology,
|
||||
{
|
||||
Ok(SURBAck::construct(
|
||||
rng,
|
||||
@@ -76,9 +76,9 @@ where
|
||||
)?)
|
||||
}
|
||||
|
||||
pub fn generate_loop_cover_packet<R, T>(
|
||||
pub fn generate_loop_cover_packet<R>(
|
||||
rng: &mut R,
|
||||
topology: &T,
|
||||
topology: &NymTopology,
|
||||
ack_key: &AckAes128Key,
|
||||
full_address: &Recipient,
|
||||
average_ack_delay: time::Duration,
|
||||
@@ -86,7 +86,6 @@ pub fn generate_loop_cover_packet<R, T>(
|
||||
) -> Result<(SocketAddr, SphinxPacket), CoverMessageError>
|
||||
where
|
||||
R: RngCore + CryptoRng,
|
||||
T: NymTopology,
|
||||
{
|
||||
// we don't care about total ack delay - we will not be retransmitting it anyway
|
||||
let (_, ack_bytes) =
|
||||
@@ -105,7 +104,8 @@ where
|
||||
.take(plaintext_size)
|
||||
.collect();
|
||||
|
||||
let route = topology.random_route_to_gateway(&full_address.gateway())?;
|
||||
let route =
|
||||
topology.random_route_to_gateway(rng, DEFAULT_NUM_MIX_HOPS, &full_address.gateway())?;
|
||||
let delays = delays::generate_from_average_duration(route.len(), average_packet_delay);
|
||||
// in our design we don't care about SURB_ID
|
||||
let destination = Destination::new(full_address.destination(), Default::default());
|
||||
|
||||
@@ -13,3 +13,7 @@
|
||||
// limitations under the License.
|
||||
|
||||
pub mod packet_sizes;
|
||||
|
||||
// If somebody can provide an argument why it might be reasonable to have more than 255 mix hops,
|
||||
// I will change this to [`usize`]
|
||||
pub const DEFAULT_NUM_MIX_HOPS: u8 = 3;
|
||||
|
||||
@@ -8,13 +8,13 @@ edition = "2018"
|
||||
|
||||
[dependencies]
|
||||
bs58 = "0.3.0"
|
||||
itertools = "0.8.2"
|
||||
log = "0.4"
|
||||
pretty_env_logger = "0.3"
|
||||
rand = "0.7.2"
|
||||
serde = { version = "1.0.104", features = ["derive"] }
|
||||
|
||||
## internal
|
||||
nymsphinx-addressing = {path = "../nymsphinx/addressing"}
|
||||
nymsphinx-types = {path = "../nymsphinx/types"}
|
||||
version-checker = {path = "../version-checker" }
|
||||
crypto = { path = "../crypto" }
|
||||
nymsphinx-addressing = { path = "../nymsphinx/addressing" }
|
||||
nymsphinx-types = { path = "../nymsphinx/types" }
|
||||
version-checker = { path = "../version-checker" }
|
||||
|
||||
@@ -13,12 +13,13 @@
|
||||
// limitations under the License.
|
||||
|
||||
use crate::filter;
|
||||
use crypto::asymmetric::identity;
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct Node {
|
||||
pub location: String,
|
||||
pub host: String,
|
||||
pub pub_key: String,
|
||||
pub pub_key: identity::PublicKey,
|
||||
pub last_seen: u64,
|
||||
pub version: String,
|
||||
}
|
||||
|
||||
@@ -12,6 +12,9 @@
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use std::collections::HashMap;
|
||||
use std::hash::Hash;
|
||||
|
||||
pub trait Versioned: Clone {
|
||||
fn version(&self) -> String;
|
||||
}
|
||||
@@ -20,7 +23,10 @@ pub trait VersionFilterable<T> {
|
||||
fn filter_by_version(&self, expected_version: &str) -> Self;
|
||||
}
|
||||
|
||||
impl<T: Versioned> VersionFilterable<T> for Vec<T> {
|
||||
impl<T> VersionFilterable<T> for Vec<T>
|
||||
where
|
||||
T: Versioned,
|
||||
{
|
||||
fn filter_by_version(&self, expected_version: &str) -> Self {
|
||||
self.iter()
|
||||
.filter(|node| {
|
||||
@@ -30,3 +36,16 @@ impl<T: Versioned> VersionFilterable<T> for Vec<T> {
|
||||
.collect()
|
||||
}
|
||||
}
|
||||
|
||||
impl<T, K, V> VersionFilterable<T> for HashMap<K, V>
|
||||
where
|
||||
K: Eq + Hash + Clone,
|
||||
V: VersionFilterable<T>,
|
||||
T: Versioned,
|
||||
{
|
||||
fn filter_by_version(&self, expected_version: &str) -> Self {
|
||||
self.iter()
|
||||
.map(|(k, v)| (k.clone(), v.filter_by_version(expected_version)))
|
||||
.collect()
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,6 +13,7 @@
|
||||
// limitations under the License.
|
||||
|
||||
use crate::filter;
|
||||
use crypto::asymmetric::{encryption, identity};
|
||||
use nymsphinx_addressing::nodes::NymNodeRoutingAddress;
|
||||
use nymsphinx_types::Node as SphinxNode;
|
||||
use std::convert::TryInto;
|
||||
@@ -28,21 +29,14 @@ pub struct Node {
|
||||
pub location: String,
|
||||
pub client_listener: String,
|
||||
pub mixnet_listener: SocketAddr,
|
||||
// TODO: should we just import common/crypto and use 'proper' types for those directly?
|
||||
pub identity_key: String,
|
||||
pub sphinx_key: String,
|
||||
pub identity_key: identity::PublicKey,
|
||||
pub sphinx_key: encryption::PublicKey, // TODO: or nymsphinx::PublicKey? both are x25519
|
||||
pub registered_clients: Vec<Client>,
|
||||
pub last_seen: u64,
|
||||
pub version: String,
|
||||
}
|
||||
|
||||
impl Node {
|
||||
fn get_sphinx_key_bytes(&self) -> [u8; 32] {
|
||||
let mut key_bytes = [0; 32];
|
||||
bs58::decode(&self.sphinx_key).into(&mut key_bytes).unwrap();
|
||||
key_bytes
|
||||
}
|
||||
|
||||
pub fn has_client(&self, client_pub_key: String) -> bool {
|
||||
self.registered_clients
|
||||
.iter()
|
||||
@@ -57,14 +51,12 @@ impl filter::Versioned for Node {
|
||||
}
|
||||
}
|
||||
|
||||
impl Into<SphinxNode> for Node {
|
||||
impl<'a> Into<SphinxNode> for &'a Node {
|
||||
fn into(self) -> SphinxNode {
|
||||
let node_address_bytes = NymNodeRoutingAddress::from(self.mixnet_listener)
|
||||
.try_into()
|
||||
.unwrap();
|
||||
let key_bytes = self.get_sphinx_key_bytes();
|
||||
let key = nymsphinx_types::PublicKey::from(key_bytes);
|
||||
|
||||
SphinxNode::new(node_address_bytes, key)
|
||||
SphinxNode::new(node_address_bytes, (&self.sphinx_key).into())
|
||||
}
|
||||
}
|
||||
|
||||
+148
-145
@@ -13,162 +13,165 @@
|
||||
// limitations under the License.
|
||||
|
||||
use crate::filter::VersionFilterable;
|
||||
use itertools::Itertools;
|
||||
use nymsphinx_types::{Node as SphinxNode, NodeAddressBytes};
|
||||
use rand::seq::IteratorRandom;
|
||||
use std::cmp::max;
|
||||
use rand::Rng;
|
||||
use std::collections::HashMap;
|
||||
|
||||
pub mod coco;
|
||||
mod filter;
|
||||
pub mod gateway;
|
||||
pub mod mix;
|
||||
pub mod provider;
|
||||
|
||||
// TODO: Figure out why 'Clone' was required to have 'TopologyAccessor<T>' working
|
||||
// even though it only contains an Arc
|
||||
pub trait NymTopology: Sized + std::fmt::Debug + Send + Sync + Clone {
|
||||
fn new_from_nodes(
|
||||
mix_nodes: Vec<mix::Node>,
|
||||
mix_provider_nodes: Vec<provider::Node>,
|
||||
coco_nodes: Vec<coco::Node>,
|
||||
gateway_nodes: Vec<gateway::Node>,
|
||||
) -> Self;
|
||||
fn mix_nodes(&self) -> Vec<mix::Node>;
|
||||
fn providers(&self) -> Vec<provider::Node>;
|
||||
fn gateways(&self) -> Vec<gateway::Node>;
|
||||
fn coco_nodes(&self) -> Vec<coco::Node>;
|
||||
fn make_layered_topology(&self) -> Result<HashMap<u64, Vec<mix::Node>>, NymTopologyError> {
|
||||
let mut layered_topology: HashMap<u64, Vec<mix::Node>> = HashMap::new();
|
||||
let mut highest_layer = 0;
|
||||
for mix in self.mix_nodes() {
|
||||
// we need to have extra space for provider
|
||||
if mix.layer > nymsphinx_types::MAX_PATH_LENGTH as u64 {
|
||||
return Err(NymTopologyError::InvalidMixLayerError);
|
||||
}
|
||||
highest_layer = max(highest_layer, mix.layer);
|
||||
|
||||
let layer_nodes = layered_topology.entry(mix.layer).or_insert_with(Vec::new);
|
||||
layer_nodes.push(mix);
|
||||
}
|
||||
|
||||
// verify the topology - make sure there are no gaps and there is at least one node per layer
|
||||
let mut missing_layers = Vec::new();
|
||||
for layer in 1..=highest_layer {
|
||||
if !layered_topology.contains_key(&layer) {
|
||||
missing_layers.push(layer);
|
||||
}
|
||||
if layered_topology[&layer].is_empty() {
|
||||
missing_layers.push(layer);
|
||||
}
|
||||
}
|
||||
|
||||
if !missing_layers.is_empty() {
|
||||
return Err(NymTopologyError::MissingLayerError(missing_layers));
|
||||
}
|
||||
|
||||
Ok(layered_topology)
|
||||
}
|
||||
|
||||
// Tries to get a route through the mix network
|
||||
fn random_mix_route(&self) -> Result<Vec<SphinxNode>, NymTopologyError> {
|
||||
let mut layered_topology = self.make_layered_topology()?;
|
||||
let num_layers = layered_topology.len();
|
||||
let route = (1..=num_layers as u64)
|
||||
// unwrap is safe for 'remove' as it it failed, it implied the entry never existed
|
||||
// in the map in the first place which would contradict what we've just done
|
||||
.map(|layer| layered_topology.remove(&layer).unwrap()) // for each layer
|
||||
.map(|nodes| nodes.into_iter().choose(&mut rand::thread_rng()).unwrap()) // choose random node
|
||||
.map(|random_node| random_node.into()) // and convert it into sphinx specific node format
|
||||
.collect();
|
||||
|
||||
Ok(route)
|
||||
}
|
||||
|
||||
fn gateway_exists(&self, gateway_address: &NodeAddressBytes) -> bool {
|
||||
let b58_address = gateway_address.to_base58_string();
|
||||
self.gateways()
|
||||
.iter()
|
||||
.find(|&gateway| gateway.identity_key == b58_address)
|
||||
.is_some()
|
||||
}
|
||||
|
||||
fn random_route_to_gateway(
|
||||
&self,
|
||||
gateway_address: &NodeAddressBytes,
|
||||
) -> Result<Vec<SphinxNode>, NymTopologyError> {
|
||||
let b58_address = gateway_address.to_base58_string();
|
||||
|
||||
let gateway = self
|
||||
.gateways()
|
||||
.iter()
|
||||
.find(|&gateway| gateway.identity_key == b58_address)
|
||||
.ok_or_else(|| NymTopologyError::NonExistentGatewayError)?
|
||||
.clone();
|
||||
|
||||
Ok(self
|
||||
.random_mix_route()?
|
||||
.into_iter()
|
||||
.chain(std::iter::once(gateway.into()))
|
||||
.collect())
|
||||
}
|
||||
|
||||
fn all_paths(&self) -> Result<Vec<Vec<SphinxNode>>, NymTopologyError> {
|
||||
let mut layered_topology = self.make_layered_topology()?;
|
||||
let gateways = self.gateways();
|
||||
|
||||
let sorted_layers: Vec<Vec<SphinxNode>> = (1..=layered_topology.len() as u64)
|
||||
.map(|layer| layered_topology.remove(&layer).unwrap()) // get all nodes per layer
|
||||
.map(|layer_nodes| layer_nodes.into_iter().map(|node| node.into()).collect()) // convert them into 'proper' sphinx nodes
|
||||
.chain(std::iter::once(
|
||||
gateways.into_iter().map(|node| node.into()).collect(),
|
||||
)) // append all gateways to the end
|
||||
.collect();
|
||||
|
||||
let all_paths = sorted_layers
|
||||
.into_iter()
|
||||
.multi_cartesian_product() // create all possible paths through that
|
||||
.collect();
|
||||
|
||||
Ok(all_paths)
|
||||
}
|
||||
|
||||
fn filter_system_version(&self, expected_version: &str) -> Self {
|
||||
self.filter_node_versions(
|
||||
expected_version,
|
||||
expected_version,
|
||||
expected_version,
|
||||
expected_version,
|
||||
)
|
||||
}
|
||||
|
||||
fn filter_node_versions(
|
||||
&self,
|
||||
expected_mix_version: &str,
|
||||
expected_provider_version: &str,
|
||||
expected_gateway_version: &str,
|
||||
expected_coco_version: &str,
|
||||
) -> Self {
|
||||
let mixes = self.mix_nodes().filter_by_version(expected_mix_version);
|
||||
let providers = self
|
||||
.providers()
|
||||
.filter_by_version(expected_provider_version);
|
||||
let gateways = self.gateways().filter_by_version(expected_gateway_version);
|
||||
let cocos = self.coco_nodes().filter_by_version(expected_coco_version);
|
||||
|
||||
Self::new_from_nodes(mixes, providers, cocos, gateways)
|
||||
}
|
||||
|
||||
fn can_construct_path_through(&self) -> bool {
|
||||
!self.mix_nodes().is_empty()
|
||||
&& !self.gateways().is_empty()
|
||||
&& self.make_layered_topology().is_ok()
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Debug)]
|
||||
pub enum NymTopologyError {
|
||||
InvalidMixLayerError,
|
||||
MissingLayerError(Vec<u64>),
|
||||
NonExistentGatewayError,
|
||||
|
||||
InvalidNumberOfHopsError,
|
||||
NoMixesOnLayerAvailable(MixLayer),
|
||||
}
|
||||
|
||||
pub type MixLayer = u8;
|
||||
|
||||
#[derive(Debug)]
|
||||
pub struct NymTopology {
|
||||
coco_nodes: Vec<coco::Node>,
|
||||
mixes: HashMap<MixLayer, Vec<mix::Node>>,
|
||||
gateways: Vec<gateway::Node>,
|
||||
}
|
||||
|
||||
impl NymTopology {
|
||||
pub fn new(
|
||||
coco_nodes: Vec<coco::Node>,
|
||||
mixes: HashMap<MixLayer, Vec<mix::Node>>,
|
||||
gateways: Vec<gateway::Node>,
|
||||
) -> Self {
|
||||
NymTopology {
|
||||
coco_nodes,
|
||||
mixes,
|
||||
gateways,
|
||||
}
|
||||
}
|
||||
|
||||
pub fn coco_nodes(&self) -> &Vec<coco::Node> {
|
||||
&self.coco_nodes
|
||||
}
|
||||
|
||||
pub fn mixes(&self) -> &HashMap<MixLayer, Vec<mix::Node>> {
|
||||
&self.mixes
|
||||
}
|
||||
|
||||
pub fn gateways(&self) -> &Vec<gateway::Node> {
|
||||
&self.gateways
|
||||
}
|
||||
|
||||
fn get_gateway(&self, gateway_address: &NodeAddressBytes) -> Option<&gateway::Node> {
|
||||
self.gateways
|
||||
.iter()
|
||||
.find(|gateway| &gateway.identity_key.derive_node_address() == gateway_address)
|
||||
}
|
||||
|
||||
pub fn gateway_exists(&self, gateway_address: &NodeAddressBytes) -> bool {
|
||||
self.get_gateway(gateway_address).is_some()
|
||||
}
|
||||
|
||||
pub fn random_mix_route<R>(
|
||||
&self,
|
||||
rng: &mut R,
|
||||
num_mix_hops: u8,
|
||||
) -> Result<Vec<SphinxNode>, NymTopologyError>
|
||||
where
|
||||
// I don't think there's a need for this RNG to be crypto-secure
|
||||
R: Rng + ?Sized,
|
||||
{
|
||||
use rand::seq::SliceRandom;
|
||||
|
||||
if self.mixes.len() < num_mix_hops as usize {
|
||||
return Err(NymTopologyError::InvalidNumberOfHopsError);
|
||||
}
|
||||
let mut route = Vec::with_capacity(num_mix_hops as usize);
|
||||
|
||||
// there is no "layer 0"
|
||||
for layer in 1..=num_mix_hops {
|
||||
// get all mixes on particular layer
|
||||
let layer_mixes = match self.mixes.get(&layer) {
|
||||
Some(mixes) => mixes,
|
||||
None => return Err(NymTopologyError::NoMixesOnLayerAvailable(layer)),
|
||||
};
|
||||
|
||||
// choose a random mix from the above list
|
||||
// this can return a 'None' only if slice is empty
|
||||
let random_mix = match layer_mixes.choose(rng) {
|
||||
Some(random_mix) => random_mix,
|
||||
None => return Err(NymTopologyError::NoMixesOnLayerAvailable(layer)),
|
||||
};
|
||||
route.push(random_mix.into());
|
||||
}
|
||||
|
||||
Ok(route)
|
||||
}
|
||||
|
||||
pub fn random_route_to_gateway<R>(
|
||||
&self,
|
||||
rng: &mut R,
|
||||
num_mix_hops: u8,
|
||||
gateway_address: &NodeAddressBytes,
|
||||
) -> Result<Vec<SphinxNode>, NymTopologyError>
|
||||
where
|
||||
// I don't think there's a need for this RNG to be crypto-secure
|
||||
R: Rng + ?Sized,
|
||||
{
|
||||
let gateway = self
|
||||
.get_gateway(gateway_address)
|
||||
.ok_or_else(|| NymTopologyError::NonExistentGatewayError)?;
|
||||
|
||||
Ok(self
|
||||
.random_mix_route(rng, num_mix_hops)?
|
||||
.into_iter()
|
||||
.chain(std::iter::once(gateway.into()))
|
||||
.collect())
|
||||
}
|
||||
|
||||
pub fn can_construct_path_through(&self, num_mix_hops: u8) -> bool {
|
||||
// if there are no gateways present, we can't do anything
|
||||
if self.gateways.is_empty() {
|
||||
return false;
|
||||
}
|
||||
|
||||
// early termination
|
||||
if self.mixes.is_empty() {
|
||||
return false;
|
||||
}
|
||||
|
||||
// make sure there's at least one mix per layer
|
||||
for i in 1..=num_mix_hops {
|
||||
match self.mixes.get(&i) {
|
||||
None => return false,
|
||||
Some(layer_entry) => {
|
||||
if layer_entry.is_empty() {
|
||||
return false;
|
||||
}
|
||||
}
|
||||
}
|
||||
}
|
||||
true
|
||||
}
|
||||
|
||||
pub fn filter_system_version(&self, expected_version: &str) -> Self {
|
||||
self.filter_node_versions(expected_version, expected_version, expected_version)
|
||||
}
|
||||
|
||||
pub fn filter_node_versions(
|
||||
&self,
|
||||
expected_mix_version: &str,
|
||||
expected_gateway_version: &str,
|
||||
expected_coco_version: &str,
|
||||
) -> Self {
|
||||
NymTopology {
|
||||
mixes: self.mixes.filter_by_version(expected_mix_version),
|
||||
gateways: self.gateways.filter_by_version(expected_gateway_version),
|
||||
coco_nodes: self.coco_nodes.filter_by_version(expected_coco_version),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
@@ -13,6 +13,7 @@
|
||||
// limitations under the License.
|
||||
|
||||
use crate::filter;
|
||||
use crypto::asymmetric::encryption;
|
||||
use nymsphinx_addressing::nodes::NymNodeRoutingAddress;
|
||||
use nymsphinx_types::Node as SphinxNode;
|
||||
use std::convert::TryInto;
|
||||
@@ -22,32 +23,22 @@ use std::net::SocketAddr;
|
||||
pub struct Node {
|
||||
pub location: String,
|
||||
pub host: SocketAddr,
|
||||
pub pub_key: String,
|
||||
pub pub_key: encryption::PublicKey, // TODO: or nymsphinx::PublicKey? both are x25519
|
||||
pub layer: u64,
|
||||
pub last_seen: u64,
|
||||
pub version: String,
|
||||
}
|
||||
|
||||
impl Node {
|
||||
pub fn get_pub_key_bytes(&self) -> [u8; 32] {
|
||||
let mut key_bytes = [0; 32];
|
||||
bs58::decode(&self.pub_key).into(&mut key_bytes).unwrap();
|
||||
key_bytes
|
||||
}
|
||||
}
|
||||
|
||||
impl filter::Versioned for Node {
|
||||
fn version(&self) -> String {
|
||||
self.version.clone()
|
||||
}
|
||||
}
|
||||
|
||||
impl Into<SphinxNode> for Node {
|
||||
impl<'a> Into<SphinxNode> for &'a Node {
|
||||
fn into(self) -> SphinxNode {
|
||||
let node_address_bytes = NymNodeRoutingAddress::from(self.host).try_into().unwrap();
|
||||
let key_bytes = self.get_pub_key_bytes();
|
||||
let key = nymsphinx_types::PublicKey::from(key_bytes);
|
||||
|
||||
SphinxNode::new(node_address_bytes, key)
|
||||
SphinxNode::new(node_address_bytes, (&self.pub_key).into())
|
||||
}
|
||||
}
|
||||
|
||||
@@ -1,61 +0,0 @@
|
||||
// Copyright 2020 Nym Technologies SA
|
||||
//
|
||||
// Licensed under the Apache License, Version 2.0 (the "License");
|
||||
// you may not use this file except in compliance with the License.
|
||||
// You may obtain a copy of the License at
|
||||
//
|
||||
// http://www.apache.org/licenses/LICENSE-2.0
|
||||
//
|
||||
// Unless required by applicable law or agreed to in writing, software
|
||||
// distributed under the License is distributed on an "AS IS" BASIS,
|
||||
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
|
||||
// See the License for the specific language governing permissions and
|
||||
// limitations under the License.
|
||||
|
||||
use crate::filter;
|
||||
use nymsphinx_addressing::nodes::NymNodeRoutingAddress;
|
||||
use nymsphinx_types::Node as SphinxNode;
|
||||
use std::convert::TryInto;
|
||||
use std::net::SocketAddr;
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct Client {
|
||||
pub pub_key: String,
|
||||
}
|
||||
|
||||
#[derive(Debug, Clone)]
|
||||
pub struct Node {
|
||||
pub location: String,
|
||||
pub client_listener: SocketAddr,
|
||||
pub mixnet_listener: SocketAddr,
|
||||
pub pub_key: String,
|
||||
pub registered_clients: Vec<Client>,
|
||||
pub last_seen: u64,
|
||||
pub version: String,
|
||||
}
|
||||
|
||||
impl Node {
|
||||
pub fn get_pub_key_bytes(&self) -> [u8; 32] {
|
||||
let mut key_bytes = [0; 32];
|
||||
bs58::decode(&self.pub_key).into(&mut key_bytes).unwrap();
|
||||
key_bytes
|
||||
}
|
||||
}
|
||||
|
||||
impl filter::Versioned for Node {
|
||||
fn version(&self) -> String {
|
||||
self.version.clone()
|
||||
}
|
||||
}
|
||||
|
||||
impl Into<SphinxNode> for Node {
|
||||
fn into(self) -> SphinxNode {
|
||||
let node_address_bytes = NymNodeRoutingAddress::from(self.mixnet_listener)
|
||||
.try_into()
|
||||
.unwrap();
|
||||
let key_bytes = self.get_pub_key_bytes();
|
||||
let key = nymsphinx_types::PublicKey::from(key_bytes);
|
||||
|
||||
SphinxNode::new(node_address_bytes, key)
|
||||
}
|
||||
}
|
||||
@@ -296,7 +296,7 @@ impl<S> Handle<S> {
|
||||
Some(address) => address,
|
||||
None => return ServerResponse::new_error("malformed request"),
|
||||
};
|
||||
let remote_address = remote_identity.derive_address();
|
||||
let remote_address = remote_identity.derive_destination_address();
|
||||
|
||||
let derived_shared_key = match self.perform_registration_handshake(init_data).await {
|
||||
Ok(shared_key) => shared_key,
|
||||
|
||||
Reference in New Issue
Block a user