feat: expose tcp_nodelay for clients and servers (#145)
* expose tcp_nodelay * do not depend on difference versions of rand
This commit is contained in:
committed by
Lucio Franco
parent
6b43f63578
commit
0eb9991b9f
@@ -80,7 +80,7 @@ tower = { git = "https://github.com/tower-rs/tower" }
|
|||||||
# Required for routeguide
|
# Required for routeguide
|
||||||
serde = { version = "1.0", features = ["derive"] }
|
serde = { version = "1.0", features = ["derive"] }
|
||||||
serde_json = "1.0"
|
serde_json = "1.0"
|
||||||
rand = "0.7.2"
|
rand = "0.6"
|
||||||
|
|
||||||
# Required for wellknown types
|
# Required for wellknown types
|
||||||
prost-types = "0.5"
|
prost-types = "0.5"
|
||||||
|
|||||||
+1
-1
@@ -76,7 +76,7 @@ rustls-native-certs = { version = "0.1", optional = true }
|
|||||||
[dev-dependencies]
|
[dev-dependencies]
|
||||||
tokio = { version = "0.2", features = ["rt-core", "macros"] }
|
tokio = { version = "0.2", features = ["rt-core", "macros"] }
|
||||||
static_assertions = "1.0"
|
static_assertions = "1.0"
|
||||||
rand = "0.7.2"
|
rand = "0.6"
|
||||||
criterion = "0.3"
|
criterion = "0.3"
|
||||||
|
|
||||||
[package.metadata.docs.rs]
|
[package.metadata.docs.rs]
|
||||||
|
|||||||
@@ -30,6 +30,7 @@ pub struct Endpoint {
|
|||||||
pub(super) init_stream_window_size: Option<u32>,
|
pub(super) init_stream_window_size: Option<u32>,
|
||||||
pub(super) init_connection_window_size: Option<u32>,
|
pub(super) init_connection_window_size: Option<u32>,
|
||||||
pub(super) tcp_keepalive: Option<Duration>,
|
pub(super) tcp_keepalive: Option<Duration>,
|
||||||
|
pub(super) tcp_nodelay: bool,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Endpoint {
|
impl Endpoint {
|
||||||
@@ -171,6 +172,14 @@ impl Endpoint {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Set the value of `TCP_NODELAY` option for accepted connections. Enabled by default.
|
||||||
|
pub fn tcp_nodelay(self, enabled: bool) -> Self {
|
||||||
|
Endpoint {
|
||||||
|
tcp_nodelay: enabled,
|
||||||
|
..self
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Create a channel from this config.
|
/// Create a channel from this config.
|
||||||
pub async fn connect(&self) -> Result<Channel, super::Error> {
|
pub async fn connect(&self) -> Result<Channel, super::Error> {
|
||||||
Channel::connect(self.clone()).await
|
Channel::connect(self.clone()).await
|
||||||
@@ -191,6 +200,7 @@ impl From<Uri> for Endpoint {
|
|||||||
init_stream_window_size: None,
|
init_stream_window_size: None,
|
||||||
init_connection_window_size: None,
|
init_connection_window_size: None,
|
||||||
tcp_keepalive: None,
|
tcp_keepalive: None,
|
||||||
|
tcp_nodelay: true,
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -56,6 +56,7 @@ pub struct Server {
|
|||||||
init_connection_window_size: Option<u32>,
|
init_connection_window_size: Option<u32>,
|
||||||
max_concurrent_streams: Option<u32>,
|
max_concurrent_streams: Option<u32>,
|
||||||
tcp_keepalive: Option<Duration>,
|
tcp_keepalive: Option<Duration>,
|
||||||
|
tcp_nodelay: bool,
|
||||||
}
|
}
|
||||||
|
|
||||||
/// A stack based `Service` router.
|
/// A stack based `Service` router.
|
||||||
@@ -77,7 +78,10 @@ pub trait ServiceName {
|
|||||||
impl Server {
|
impl Server {
|
||||||
/// Create a new server builder that can configure a [`Server`].
|
/// Create a new server builder that can configure a [`Server`].
|
||||||
pub fn builder() -> Self {
|
pub fn builder() -> Self {
|
||||||
Default::default()
|
Server {
|
||||||
|
tcp_nodelay: true,
|
||||||
|
..Default::default()
|
||||||
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -164,6 +168,14 @@ impl Server {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// Set the value of `TCP_NODELAY` option for accepted connections. Enabled by default.
|
||||||
|
pub fn tcp_nodelay(self, enabled: bool) -> Self {
|
||||||
|
Server {
|
||||||
|
tcp_nodelay: enabled,
|
||||||
|
..self
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
/// Intercept the execution of gRPC methods.
|
/// Intercept the execution of gRPC methods.
|
||||||
///
|
///
|
||||||
/// ```
|
/// ```
|
||||||
@@ -221,12 +233,13 @@ impl Server {
|
|||||||
let init_connection_window_size = self.init_connection_window_size;
|
let init_connection_window_size = self.init_connection_window_size;
|
||||||
let init_stream_window_size = self.init_stream_window_size;
|
let init_stream_window_size = self.init_stream_window_size;
|
||||||
let max_concurrent_streams = self.max_concurrent_streams;
|
let max_concurrent_streams = self.max_concurrent_streams;
|
||||||
let tcp_keepalive = self.tcp_keepalive;
|
|
||||||
// let timeout = self.timeout.clone();
|
// let timeout = self.timeout.clone();
|
||||||
|
|
||||||
let incoming = hyper::server::accept::from_stream::<_, _, crate::Error>(
|
let incoming = hyper::server::accept::from_stream::<_, _, crate::Error>(
|
||||||
async_stream::try_stream! {
|
async_stream::try_stream! {
|
||||||
let mut tcp = TcpIncoming::bind(addr, tcp_keepalive)?;
|
let mut tcp = TcpIncoming::bind(addr)?
|
||||||
|
.set_nodelay(self.tcp_nodelay)
|
||||||
|
.set_keepalive(self.tcp_keepalive);
|
||||||
|
|
||||||
while let Some(stream) = tcp.try_next().await? {
|
while let Some(stream) = tcp.try_next().await? {
|
||||||
#[cfg(feature = "tls")]
|
#[cfg(feature = "tls")]
|
||||||
@@ -418,13 +431,20 @@ struct TcpIncoming {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl TcpIncoming {
|
impl TcpIncoming {
|
||||||
fn bind(addr: SocketAddr, tcp_keepalive: Option<Duration>) -> Result<Self, crate::Error> {
|
fn bind(addr: SocketAddr) -> Result<Self, crate::Error> {
|
||||||
let mut inner = conn::AddrIncoming::bind(&addr).map_err(Box::new)?;
|
let inner = conn::AddrIncoming::bind(&addr).map_err(Box::new)?;
|
||||||
inner.set_nodelay(true);
|
|
||||||
inner.set_keepalive(tcp_keepalive);
|
|
||||||
|
|
||||||
Ok(Self { inner })
|
Ok(Self { inner })
|
||||||
}
|
}
|
||||||
|
|
||||||
|
fn set_nodelay(mut self, enabled: bool) -> Self {
|
||||||
|
self.inner.set_nodelay(enabled);
|
||||||
|
self
|
||||||
|
}
|
||||||
|
|
||||||
|
fn set_keepalive(mut self, tcp_keepalive: Option<Duration>) -> Self {
|
||||||
|
self.inner.set_keepalive(tcp_keepalive);
|
||||||
|
self
|
||||||
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl Stream for TcpIncoming {
|
impl Stream for TcpIncoming {
|
||||||
|
|||||||
@@ -28,10 +28,14 @@ pub(crate) struct Connection {
|
|||||||
impl Connection {
|
impl Connection {
|
||||||
pub(crate) async fn new(endpoint: Endpoint) -> Result<Self, crate::Error> {
|
pub(crate) async fn new(endpoint: Endpoint) -> Result<Self, crate::Error> {
|
||||||
#[cfg(feature = "tls")]
|
#[cfg(feature = "tls")]
|
||||||
let connector = connector(endpoint.tls.clone(), endpoint.tcp_keepalive);
|
let connector = connector(endpoint.tls.clone())
|
||||||
|
.set_keepalive(endpoint.tcp_keepalive)
|
||||||
|
.set_nodelay(endpoint.tcp_nodelay);
|
||||||
|
|
||||||
#[cfg(not(feature = "tls"))]
|
#[cfg(not(feature = "tls"))]
|
||||||
let connector = connector(endpoint.tcp_keepalive);
|
let connector = connector()
|
||||||
|
.set_keepalive(endpoint.tcp_keepalive)
|
||||||
|
.set_nodelay(endpoint.tcp_nodelay);
|
||||||
|
|
||||||
let settings = Builder::new()
|
let settings = Builder::new()
|
||||||
.http2_initial_stream_window_size(endpoint.init_stream_window_size)
|
.http2_initial_stream_window_size(endpoint.init_stream_window_size)
|
||||||
|
|||||||
@@ -11,17 +11,13 @@ use tower_make::MakeConnection;
|
|||||||
use tower_service::Service;
|
use tower_service::Service;
|
||||||
|
|
||||||
#[cfg(not(feature = "tls"))]
|
#[cfg(not(feature = "tls"))]
|
||||||
pub(crate) fn connector(tcp_keepalive: Option<Duration>) -> HttpConnector {
|
pub(crate) fn connector() -> Connector {
|
||||||
let mut http = HttpConnector::new();
|
Connector::new()
|
||||||
http.enforce_http(false);
|
|
||||||
http.set_nodelay(true);
|
|
||||||
http.set_keepalive(tcp_keepalive);
|
|
||||||
http
|
|
||||||
}
|
}
|
||||||
|
|
||||||
#[cfg(feature = "tls")]
|
#[cfg(feature = "tls")]
|
||||||
pub(crate) fn connector(tls: Option<TlsConnector>, tcp_keepalive: Option<Duration>) -> Connector {
|
pub(crate) fn connector(tls: Option<TlsConnector>) -> Connector {
|
||||||
Connector::new(tls, tcp_keepalive)
|
Connector::new(tls)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub(crate) struct Connector {
|
pub(crate) struct Connector {
|
||||||
@@ -31,13 +27,35 @@ pub(crate) struct Connector {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl Connector {
|
impl Connector {
|
||||||
|
#[cfg(not(feature = "tls"))]
|
||||||
|
pub(crate) fn new() -> Self {
|
||||||
|
Self {
|
||||||
|
http: Self::http_connector(),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(feature = "tls")]
|
#[cfg(feature = "tls")]
|
||||||
pub(crate) fn new(tls: Option<TlsConnector>, tcp_keepalive: Option<Duration>) -> Self {
|
fn new(tls: Option<TlsConnector>) -> Self {
|
||||||
|
Self {
|
||||||
|
http: Self::http_connector(),
|
||||||
|
tls,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn set_nodelay(mut self, enabled: bool) -> Self {
|
||||||
|
self.http.set_nodelay(enabled);
|
||||||
|
self
|
||||||
|
}
|
||||||
|
|
||||||
|
pub(crate) fn set_keepalive(mut self, duration: Option<Duration>) -> Self {
|
||||||
|
self.http.set_keepalive(duration);
|
||||||
|
self
|
||||||
|
}
|
||||||
|
|
||||||
|
fn http_connector() -> HttpConnector {
|
||||||
let mut http = HttpConnector::new();
|
let mut http = HttpConnector::new();
|
||||||
http.enforce_http(false);
|
http.enforce_http(false);
|
||||||
http.set_nodelay(true);
|
http
|
||||||
http.set_keepalive(tcp_keepalive);
|
|
||||||
Self { http, tls }
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user