Add trait to mock client dependency in DelayForwarder (#1073)
* client-libs/mixnode: add trait for client sending * mixnode: instantiate and test DelayForwarder * mixnode: remove commented out line
This commit is contained in:
@@ -187,11 +187,15 @@ impl MixNode {
|
||||
) -> PacketDelayForwardSender {
|
||||
info!("Starting packet delay-forwarder...");
|
||||
|
||||
let mut packet_forwarder = DelayForwarder::new(
|
||||
let client_config = mixnet_client::Config::new(
|
||||
self.config.get_packet_forwarding_initial_backoff(),
|
||||
self.config.get_packet_forwarding_maximum_backoff(),
|
||||
self.config.get_initial_connection_timeout(),
|
||||
self.config.get_maximum_connection_buffer_size(),
|
||||
);
|
||||
|
||||
let mut packet_forwarder = DelayForwarder::new(
|
||||
mixnet_client::Client::new(client_config),
|
||||
node_stats_update_sender,
|
||||
);
|
||||
|
||||
|
||||
@@ -7,7 +7,7 @@ use futures::StreamExt;
|
||||
use nonexhaustive_delayqueue::{Expired, NonExhaustiveDelayQueue, TimerError};
|
||||
use nymsphinx::forwarding::packet::MixPacket;
|
||||
use std::io;
|
||||
use tokio::time::{Duration, Instant};
|
||||
use tokio::time::Instant;
|
||||
|
||||
// Delay + MixPacket vs Instant + MixPacket
|
||||
|
||||
@@ -17,34 +17,27 @@ pub(crate) type PacketDelayForwardSender = mpsc::UnboundedSender<(MixPacket, Opt
|
||||
type PacketDelayForwardReceiver = mpsc::UnboundedReceiver<(MixPacket, Option<Instant>)>;
|
||||
|
||||
/// Entity responsible for delaying received sphinx packet and forwarding it to next node.
|
||||
pub(crate) struct DelayForwarder {
|
||||
pub(crate) struct DelayForwarder<C>
|
||||
where
|
||||
C: mixnet_client::SendWithoutResponse,
|
||||
{
|
||||
delay_queue: NonExhaustiveDelayQueue<MixPacket>,
|
||||
mixnet_client: mixnet_client::Client,
|
||||
mixnet_client: C,
|
||||
packet_sender: PacketDelayForwardSender,
|
||||
packet_receiver: PacketDelayForwardReceiver,
|
||||
node_stats_update_sender: UpdateSender,
|
||||
}
|
||||
|
||||
impl DelayForwarder {
|
||||
pub(crate) fn new(
|
||||
initial_reconnection_backoff: Duration,
|
||||
maximum_reconnection_backoff: Duration,
|
||||
initial_connection_timeout: Duration,
|
||||
maximum_connection_buffer_size: usize,
|
||||
node_stats_update_sender: UpdateSender,
|
||||
) -> Self {
|
||||
let client_config = mixnet_client::Config::new(
|
||||
initial_reconnection_backoff,
|
||||
maximum_reconnection_backoff,
|
||||
initial_connection_timeout,
|
||||
maximum_connection_buffer_size,
|
||||
);
|
||||
|
||||
impl<C> DelayForwarder<C>
|
||||
where
|
||||
C: mixnet_client::SendWithoutResponse,
|
||||
{
|
||||
pub(crate) fn new(client: C, node_stats_update_sender: UpdateSender) -> DelayForwarder<C> {
|
||||
let (packet_sender, packet_receiver) = mpsc::unbounded();
|
||||
|
||||
DelayForwarder {
|
||||
DelayForwarder::<C> {
|
||||
delay_queue: NonExhaustiveDelayQueue::new(),
|
||||
mixnet_client: mixnet_client::Client::new(client_config),
|
||||
mixnet_client: client,
|
||||
packet_sender,
|
||||
packet_receiver,
|
||||
node_stats_update_sender,
|
||||
@@ -123,3 +116,115 @@ impl DelayForwarder {
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
|
||||
use std::net::{IpAddr, Ipv4Addr, SocketAddr};
|
||||
use std::sync::{Arc, Mutex};
|
||||
use std::time::Duration;
|
||||
|
||||
use nymsphinx::addressing::nodes::NymNodeRoutingAddress;
|
||||
use nymsphinx_params::packet_sizes::PacketSize;
|
||||
use nymsphinx_params::PacketMode;
|
||||
use nymsphinx_types::builder::SphinxPacketBuilder;
|
||||
use nymsphinx_types::{
|
||||
crypto, Delay as SphinxDelay, Destination, DestinationAddressBytes, Node, NodeAddressBytes,
|
||||
SphinxPacket, DESTINATION_ADDRESS_LENGTH, IDENTIFIER_LENGTH, NODE_ADDRESS_LENGTH,
|
||||
};
|
||||
|
||||
#[derive(Default)]
|
||||
struct TestClient {
|
||||
pub packets_sent: Arc<Mutex<Vec<(NymNodeRoutingAddress, SphinxPacket, PacketMode)>>>,
|
||||
}
|
||||
|
||||
impl mixnet_client::SendWithoutResponse for TestClient {
|
||||
fn send_without_response(
|
||||
&mut self,
|
||||
address: NymNodeRoutingAddress,
|
||||
packet: SphinxPacket,
|
||||
packet_mode: PacketMode,
|
||||
) -> io::Result<()> {
|
||||
self.packets_sent
|
||||
.lock()
|
||||
.unwrap()
|
||||
.push((address, packet, packet_mode));
|
||||
Ok(())
|
||||
}
|
||||
}
|
||||
|
||||
fn make_valid_sphinx_packet(size: PacketSize) -> SphinxPacket {
|
||||
let (_, node1_pk) = crypto::keygen();
|
||||
let node1 = Node::new(
|
||||
NodeAddressBytes::from_bytes([5u8; NODE_ADDRESS_LENGTH]),
|
||||
node1_pk,
|
||||
);
|
||||
let (_, node2_pk) = crypto::keygen();
|
||||
let node2 = Node::new(
|
||||
NodeAddressBytes::from_bytes([4u8; NODE_ADDRESS_LENGTH]),
|
||||
node2_pk,
|
||||
);
|
||||
let (_, node3_pk) = crypto::keygen();
|
||||
let node3 = Node::new(
|
||||
NodeAddressBytes::from_bytes([2u8; NODE_ADDRESS_LENGTH]),
|
||||
node3_pk,
|
||||
);
|
||||
|
||||
let route = [node1, node2, node3];
|
||||
let destination = Destination::new(
|
||||
DestinationAddressBytes::from_bytes([3u8; DESTINATION_ADDRESS_LENGTH]),
|
||||
[4u8; IDENTIFIER_LENGTH],
|
||||
);
|
||||
let delays = vec![
|
||||
SphinxDelay::new_from_nanos(42),
|
||||
SphinxDelay::new_from_nanos(42),
|
||||
SphinxDelay::new_from_nanos(42),
|
||||
];
|
||||
SphinxPacketBuilder::new()
|
||||
.with_payload_size(size.payload_size())
|
||||
.build_packet(b"foomp".to_vec(), &route, &destination, &delays)
|
||||
.unwrap()
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn packets_received_are_forwarded() {
|
||||
// Wire up the DelayForwarder
|
||||
let (stats_sender, _stats_receiver) = mpsc::unbounded();
|
||||
let node_stats_update_sender = UpdateSender::new(stats_sender);
|
||||
let client = TestClient::default();
|
||||
let client_packets_sent = client.packets_sent.clone();
|
||||
let mut delay_forwarder = DelayForwarder::new(client, node_stats_update_sender);
|
||||
let packet_sender = delay_forwarder.sender();
|
||||
|
||||
// Spawn the worker, listening on packet_sender channel
|
||||
tokio::spawn(async move { delay_forwarder.run().await });
|
||||
|
||||
// Send a `MixPacket` down the channel without any delay attached.
|
||||
let next_hop =
|
||||
NymNodeRoutingAddress::from(SocketAddr::new(IpAddr::V4(Ipv4Addr::new(1, 2, 3, 4)), 42));
|
||||
let mix_packet = MixPacket::new(
|
||||
next_hop,
|
||||
make_valid_sphinx_packet(PacketSize::default()),
|
||||
PacketMode::default(),
|
||||
);
|
||||
let forward_instant = None;
|
||||
packet_sender
|
||||
.unbounded_send((mix_packet, forward_instant))
|
||||
.unwrap();
|
||||
|
||||
// Give the the worker a chance to act
|
||||
tokio::time::sleep(Duration::from_millis(10)).await;
|
||||
|
||||
// The client should have forwarded the packet straight away
|
||||
assert_eq!(
|
||||
client_packets_sent
|
||||
.lock()
|
||||
.unwrap()
|
||||
.iter()
|
||||
.map(|(a, _, _)| *a)
|
||||
.collect::<Vec<_>>(),
|
||||
vec![next_hop]
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user