add aquatic_common_tcp crate, move common functionality there

This commit is contained in:
Joakim Frostegård 2020-07-02 16:34:36 +02:00
parent 1eaf2a0351
commit 2e53a2adc1
14 changed files with 163 additions and 127 deletions

11
Cargo.lock generated
View file

@ -54,6 +54,16 @@ dependencies = [
"rand", "rand",
] ]
[[package]]
name = "aquatic_common_tcp"
version = "0.1.0"
dependencies = [
"anyhow",
"aquatic_common",
"mio",
"native-tls",
]
[[package]] [[package]]
name = "aquatic_http" name = "aquatic_http"
version = "0.1.0" version = "0.1.0"
@ -61,6 +71,7 @@ dependencies = [
"anyhow", "anyhow",
"aquatic_cli_helpers", "aquatic_cli_helpers",
"aquatic_common", "aquatic_common",
"aquatic_common_tcp",
"bendy", "bendy",
"either", "either",
"flume", "flume",

View file

@ -3,6 +3,7 @@
members = [ members = [
"aquatic_cli_helpers", "aquatic_cli_helpers",
"aquatic_common", "aquatic_common",
"aquatic_common_tcp",
"aquatic_http", "aquatic_http",
"aquatic_udp", "aquatic_udp",
"aquatic_udp_bench", "aquatic_udp_bench",

View file

@ -6,6 +6,7 @@
and maybe run scripts should be adjusted and maybe run scripts should be adjusted
## aquatic_http ## aquatic_http
* move stuff to common crate with ws: what about Request/InMessage etc?
* handshake stuff * handshake stuff
* fix overcomplicated and probably incorrect implementation * fix overcomplicated and probably incorrect implementation
* support TLS and plain at the same time?? * support TLS and plain at the same time??
@ -13,7 +14,6 @@
* fixed size buffer is probably bad * fixed size buffer is probably bad
* compact peer representation in announce response: is implementation correct? * compact peer representation in announce response: is implementation correct?
* scrape info hash parsing: multiple ought to be accepted * scrape info hash parsing: multiple ought to be accepted
* move stuff to common crate with ws: what about Request/InMessage etc?
* info hashes, peer ids: check that whole deserialization and url decoding * info hashes, peer ids: check that whole deserialization and url decoding
works as it should. There are suspicously many `\u{fffd}` works as it should. There are suspicously many `\u{fffd}`

View file

@ -0,0 +1,15 @@
[package]
name = "aquatic_common_tcp"
version = "0.1.0"
authors = ["Joakim Frostegård <joakim.frostegard@gmail.com>"]
edition = "2018"
license = "Apache-2.0"
[lib]
name = "aquatic_common_tcp"
[dependencies]
anyhow = "1"
aquatic_common = { path = "../aquatic_common" }
mio = { version = "0.7", features = ["tcp", "os-poll", "os-util"] }
native-tls = "0.2"

View file

View file

@ -0,0 +1,2 @@
pub mod config;
pub mod network;

View file

@ -0,0 +1 @@
pub mod stream;

View file

@ -0,0 +1,76 @@
use std::net::{SocketAddr};
use std::io::{Read, Write};
use mio::net::TcpStream;
use native_tls::TlsStream;
pub enum Stream {
TcpStream(TcpStream),
TlsStream(TlsStream<TcpStream>),
}
impl Stream {
#[inline]
pub fn get_peer_addr(&self) -> SocketAddr {
match self {
Self::TcpStream(stream) => stream.peer_addr().unwrap(),
Self::TlsStream(stream) => stream.get_ref().peer_addr().unwrap(),
}
}
}
impl Read for Stream {
#[inline]
fn read(&mut self, buf: &mut [u8]) -> Result<usize, ::std::io::Error> {
match self {
Self::TcpStream(stream) => stream.read(buf),
Self::TlsStream(stream) => stream.read(buf),
}
}
/// Not used but provided for completeness
#[inline]
fn read_vectored(
&mut self,
bufs: &mut [::std::io::IoSliceMut<'_>]
) -> ::std::io::Result<usize> {
match self {
Self::TcpStream(stream) => stream.read_vectored(bufs),
Self::TlsStream(stream) => stream.read_vectored(bufs),
}
}
}
impl Write for Stream {
#[inline]
fn write(&mut self, buf: &[u8]) -> ::std::io::Result<usize> {
match self {
Self::TcpStream(stream) => stream.write(buf),
Self::TlsStream(stream) => stream.write(buf),
}
}
/// Not used but provided for completeness
#[inline]
fn write_vectored(
&mut self,
bufs: &[::std::io::IoSlice<'_>]
) -> ::std::io::Result<usize> {
match self {
Self::TcpStream(stream) => stream.write_vectored(bufs),
Self::TlsStream(stream) => stream.write_vectored(bufs),
}
}
#[inline]
fn flush(&mut self) -> ::std::io::Result<()> {
match self {
Self::TcpStream(stream) => stream.flush(),
Self::TlsStream(stream) => stream.flush(),
}
}
}

View file

@ -17,6 +17,7 @@ path = "src/bin/main.rs"
anyhow = "1" anyhow = "1"
aquatic_cli_helpers = { path = "../aquatic_cli_helpers" } aquatic_cli_helpers = { path = "../aquatic_cli_helpers" }
aquatic_common = { path = "../aquatic_common" } aquatic_common = { path = "../aquatic_common" }
aquatic_common_tcp = { path = "../aquatic_common_tcp" }
bendy = { version = "0.3", features = ["std", "serde"] } bendy = { version = "0.3", features = ["std", "serde"] }
either = "1" either = "1"
flume = "0.7" flume = "0.7"

View file

@ -1,4 +1,4 @@
use std::net::{SocketAddr, IpAddr}; use std::net::SocketAddr;
use std::sync::Arc; use std::sync::Arc;
use flume::{Sender, Receiver}; use flume::{Sender, Receiver};
@ -34,7 +34,7 @@ pub enum PeerStatus {
} }
// identical to ws version - FIXME only if bytes left is optional // almost identical to ws version - FIXME only if bytes left is optional
impl PeerStatus { impl PeerStatus {
/// Determine peer status from announce event and number of bytes left. /// Determine peer status from announce event and number of bytes left.
/// ///

View file

@ -1,6 +1,6 @@
use std::net::{SocketAddr}; use std::net::{SocketAddr};
use std::io::{Read, Write};
use std::io::ErrorKind; use std::io::ErrorKind;
use std::io::{Read, Write};
use either::Either; use either::Either;
use hashbrown::HashMap; use hashbrown::HashMap;
@ -9,81 +9,12 @@ use mio::Token;
use mio::net::TcpStream; use mio::net::TcpStream;
use native_tls::{TlsAcceptor, TlsStream, MidHandshakeTlsStream}; use native_tls::{TlsAcceptor, TlsStream, MidHandshakeTlsStream};
use aquatic_common_tcp::network::stream::Stream;
use crate::common::*; use crate::common::*;
use crate::protocol::{Request, Response}; use crate::protocol::{Request, Response};
pub enum Stream {
TcpStream(TcpStream),
TlsStream(TlsStream<TcpStream>),
}
impl Stream {
#[inline]
pub fn get_peer_addr(&self) -> SocketAddr {
match self {
Self::TcpStream(stream) => stream.peer_addr().unwrap(),
Self::TlsStream(stream) => stream.get_ref().peer_addr().unwrap(),
}
}
}
impl Read for Stream {
#[inline]
fn read(&mut self, buf: &mut [u8]) -> Result<usize, ::std::io::Error> {
match self {
Self::TcpStream(stream) => stream.read(buf),
Self::TlsStream(stream) => stream.read(buf),
}
}
/// Not used but provided for completeness
#[inline]
fn read_vectored(
&mut self,
bufs: &mut [::std::io::IoSliceMut<'_>]
) -> ::std::io::Result<usize> {
match self {
Self::TcpStream(stream) => stream.read_vectored(bufs),
Self::TlsStream(stream) => stream.read_vectored(bufs),
}
}
}
impl Write for Stream {
#[inline]
fn write(&mut self, buf: &[u8]) -> ::std::io::Result<usize> {
match self {
Self::TcpStream(stream) => stream.write(buf),
Self::TlsStream(stream) => stream.write(buf),
}
}
/// Not used but provided for completeness
#[inline]
fn write_vectored(
&mut self,
bufs: &[::std::io::IoSlice<'_>]
) -> ::std::io::Result<usize> {
match self {
Self::TcpStream(stream) => stream.write_vectored(bufs),
Self::TlsStream(stream) => stream.write_vectored(bufs),
}
}
#[inline]
fn flush(&mut self) -> ::std::io::Result<()> {
match self {
Self::TcpStream(stream) => stream.flush(),
Self::TlsStream(stream) => stream.flush(),
}
}
}
#[derive(Debug)] #[derive(Debug)]
pub enum RequestParseError { pub enum RequestParseError {
NeedMoreData, NeedMoreData,
@ -296,4 +227,4 @@ impl Connection {
} }
pub type ConnectionMap<'a> = HashMap<Token, Connection>; pub type ConnectionMap = HashMap<Token, Connection>;

View file

@ -1,3 +1,6 @@
pub mod connection;
pub mod utils;
use std::time::Duration; use std::time::Duration;
use std::io::ErrorKind; use std::io::ErrorKind;
@ -11,13 +14,50 @@ use crate::common::*;
use crate::config::Config; use crate::config::Config;
use crate::protocol::*; use crate::protocol::*;
pub mod connection;
pub mod utils;
use connection::*; use connection::*;
use utils::*; use utils::*;
fn accept_new_streams(
listener: &mut TcpListener,
poll: &mut Poll,
connections: &mut ConnectionMap,
valid_until: ValidUntil,
poll_token_counter: &mut Token,
){
loop {
match listener.accept(){
Ok((mut stream, _)) => {
poll_token_counter.0 = poll_token_counter.0.wrapping_add(1);
if poll_token_counter.0 == 0 {
poll_token_counter.0 = 1;
}
let token = *poll_token_counter;
remove_connection_if_exists(connections, token);
poll.registry()
.register(&mut stream, token, Interest::READABLE)
.unwrap();
let connection = Connection::new(valid_until, stream);
connections.insert(token, connection);
},
Err(err) => {
if err.kind() == ErrorKind::WouldBlock {
break
}
info!("error while accepting streams: {}", err);
}
}
}
}
// will be almost identical to ws version // will be almost identical to ws version
pub fn run_socket_worker( pub fn run_socket_worker(
config: Config, config: Config,
@ -119,54 +159,14 @@ pub fn run_poll_loop(
} }
// will be identical to ws version
fn accept_new_streams(
listener: &mut TcpListener,
poll: &mut Poll,
connections: &mut ConnectionMap,
valid_until: ValidUntil,
poll_token_counter: &mut Token,
){
loop {
match listener.accept(){
Ok((mut stream, _)) => {
poll_token_counter.0 = poll_token_counter.0.wrapping_add(1);
if poll_token_counter.0 == 0 {
poll_token_counter.0 = 1;
}
let token = *poll_token_counter;
remove_connection_if_exists(connections, token);
poll.registry()
.register(&mut stream, token, Interest::READABLE)
.unwrap();
let connection = Connection::new(valid_until, stream);
connections.insert(token, connection);
},
Err(err) => {
if err.kind() == ErrorKind::WouldBlock {
break
}
info!("error while accepting streams: {}", err);
}
}
}
}
/// On the stream given by poll_token, get TLS (if requested) and tungstenite /// On the stream given by poll_token, get TLS (if requested) and tungstenite
/// up and running, then read messages and pass on through channel. /// up and running, then read messages and pass on through channel.
pub fn run_handshake_and_read_requests<'a>( pub fn run_handshake_and_read_requests(
socket_worker_index: usize, socket_worker_index: usize,
request_channel_sender: &RequestChannelSender, request_channel_sender: &RequestChannelSender,
opt_tls_acceptor: &'a Option<TlsAcceptor>, // If set, run TLS opt_tls_acceptor: &Option<TlsAcceptor>, // If set, run TLS
connections: &'a mut ConnectionMap<'a>, connections: &mut ConnectionMap,
poll_token: Token, poll_token: Token,
valid_until: ValidUntil, valid_until: ValidUntil,
){ ){

View file

@ -6,10 +6,10 @@ use socket2::{Socket, Domain, Type, Protocol};
use crate::config::Config; use crate::config::Config;
use super::connection::*; use super::*;
// will be identical to ws version // will be almost identical to ws version
pub fn create_listener( pub fn create_listener(
config: &Config config: &Config
) -> ::anyhow::Result<::std::net::TcpListener> { ) -> ::anyhow::Result<::std::net::TcpListener> {
@ -38,7 +38,6 @@ pub fn create_listener(
} }
// will be identical to ws version
/// Don't bother with deregistering from Poll. In my understanding, this is /// Don't bother with deregistering from Poll. In my understanding, this is
/// done automatically when the stream is dropped, as long as there are no /// done automatically when the stream is dropped, as long as there are no
/// other references to the file descriptor, such as when it is accessed /// other references to the file descriptor, such as when it is accessed
@ -53,7 +52,6 @@ pub fn remove_connection_if_exists(
} }
// will be identical to ws version
// Close and remove inactive connections // Close and remove inactive connections
pub fn remove_inactive_connections( pub fn remove_inactive_connections(
connections: &mut ConnectionMap, connections: &mut ConnectionMap,

View file

@ -1,6 +1,6 @@
use std::net::IpAddr; use std::net::IpAddr;
use hashbrown::HashMap; use hashbrown::HashMap;
use serde::{Serialize, Deserialize, Serializer}; use serde::{Serialize, Deserialize};
use crate::common::Peer; use crate::common::Peer;