Moved listening loop to provider struct impl

This commit is contained in:
Jedrzej Stuczynski
2019-12-09 12:47:38 +00:00
parent 86861ccb70
commit 836cff226a
2 changed files with 57 additions and 37 deletions
+4 -2
View File
@@ -1,3 +1,4 @@
use crate::provider::ServiceProvider;
use clap::{App, Arg, ArgMatches, SubCommand};
use curve25519_dalek::scalar::Scalar;
use std::net::ToSocketAddrs;
@@ -48,7 +49,8 @@ fn run(matches: &ArgMatches) {
// make sure our socket_address is equal to our predefined-hardcoded value
assert_eq!("127.0.0.1:8081", socket_address.to_string());
provider::listening_loop().unwrap();
let provider = ServiceProvider::new(socket_address, secret_key);
provider.start_listening().unwrap()
}
fn main() {
@@ -58,7 +60,7 @@ fn main() {
.about("Implementation of the Loopix-based Service Provider")
.subcommand(
SubCommand::with_name("run")
.about("Starts the mixnode")
.about("Starts the service provider")
.arg(
Arg::with_name("host")
.short("h")
+53 -35
View File
@@ -1,47 +1,65 @@
use sphinx::{ProcessedPacket, SphinxPacket};
use tokio::prelude::*;
use tokio::runtime::Runtime;
use curve25519_dalek::scalar::Scalar;
use std::net::SocketAddr;
pub fn listening_loop() -> Result<(), Box<dyn std::error::Error>> {
// Create the runtime, probably later move it to Provider struct itself?
let mut rt = Runtime::new()?;
pub struct ServiceProvider {
network_address: SocketAddr,
secret_key: Scalar,
}
// Spawn the root task
rt.block_on(async {
let my_address = "127.0.0.1:8081";
let mut listener = tokio::net::TcpListener::bind(my_address).await?;
impl ServiceProvider {
pub fn new(network_address: SocketAddr, secret_key: Scalar) -> Self {
ServiceProvider {
network_address,
secret_key,
}
}
println!("Starting Nym store-and-forward Provider on address {:?}", my_address);
println!("Waiting for input...");
pub fn start_listening(&self) -> Result<(), Box<dyn std::error::Error>> {
// Create the runtime, probably later move it to Provider struct itself?
let mut rt = Runtime::new()?;
loop {
let (mut inbound, _) = listener.accept().await?;
// Spawn the root task
rt.block_on(async {
let my_address = "127.0.0.1:8081";
let mut listener = tokio::net::TcpListener::bind(my_address).await?;
tokio::spawn(async move {
let mut buf = [0; sphinx::PACKET_SIZE];
println!("Starting Nym store-and-forward Provider on address {:?}", my_address);
println!("Waiting for input...");
loop {
match inbound.read(&mut buf).await {
Ok(length) if length == 0 =>
{
println!("Remote connection closed.");
loop {
let (mut inbound, _) = listener.accept().await?;
tokio::spawn(async move {
let mut buf = [0; sphinx::PACKET_SIZE];
loop {
match inbound.read(&mut buf).await {
Ok(length) if length == 0 =>
{
println!("Remote connection closed.");
return;
}
Ok(_) => {
let packet = SphinxPacket::from_bytes(buf.to_vec()).unwrap();
let payload = match packet.process(Default::default()) {
ProcessedPacket::ProcessedPacketFinalHop(_, _, payload) => Some(payload),
_ => None,
}.unwrap();
let message = payload.get_content();
}
Err(e) => {
println!("failed to read from socket; err = {:?}", e);
return;
}
Ok(_) => {
let packet = SphinxPacket::from_bytes(buf.to_vec()).unwrap();
let payload = match packet.process(Default::default()) {
ProcessedPacket::ProcessedPacketFinalHop(_, _, payload) => Some(payload),
_ => None,
}.unwrap();
let message = payload.get_content();
}
Err(e) => {
println!("failed to read from socket; err = {:?}", e);
return;
}
};
}
});
}
})
};
}
});
}
})
}
}