|
@@ -31,6 +31,7 @@ use futures::{
|
|
|
stream::{FuturesUnordered, StreamExt},
|
|
stream::{FuturesUnordered, StreamExt},
|
|
|
};
|
|
};
|
|
|
use smol::Timer;
|
|
use smol::Timer;
|
|
|
|
|
+use tracing::warn;
|
|
|
use url::Url;
|
|
use url::Url;
|
|
|
|
|
|
|
|
#[cfg(feature = "upnp-igd")]
|
|
#[cfg(feature = "upnp-igd")]
|
|
@@ -58,6 +59,33 @@ use crate::{
|
|
|
/// Atomic pointer to Acceptor
|
|
/// Atomic pointer to Acceptor
|
|
|
pub type AcceptorPtr = Arc<Acceptor>;
|
|
pub type AcceptorPtr = Arc<Acceptor>;
|
|
|
|
|
|
|
|
|
|
+const ACCEPT_RETRY_MIN: Duration = Duration::from_millis(100);
|
|
|
|
|
+const ACCEPT_RETRY_MAX: Duration = Duration::from_secs(5);
|
|
|
|
|
+
|
|
|
|
|
+struct ResourceExhaustionBackoff {
|
|
|
|
|
+ next: Duration,
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+impl ResourceExhaustionBackoff {
|
|
|
|
|
+ fn new() -> Self {
|
|
|
|
|
+ Self { next: ACCEPT_RETRY_MIN }
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ fn next_delay(&mut self) -> Duration {
|
|
|
|
|
+ let delay = self.next;
|
|
|
|
|
+ self.next = self.next.saturating_mul(2).min(ACCEPT_RETRY_MAX);
|
|
|
|
|
+ delay
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ fn reset(&mut self) {
|
|
|
|
|
+ self.next = ACCEPT_RETRY_MIN;
|
|
|
|
|
+ }
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
|
|
+fn is_descriptor_exhaustion(error: &std::io::Error) -> bool {
|
|
|
|
|
+ error.raw_os_error().is_some_and(|code| code == libc::EMFILE || code == libc::ENFILE)
|
|
|
|
|
+}
|
|
|
|
|
+
|
|
|
fn with_handshake_timeout(negotiation: PtNegotiation, timeout: Duration) -> PtNegotiation {
|
|
fn with_handshake_timeout(negotiation: PtNegotiation, timeout: Duration) -> PtNegotiation {
|
|
|
Box::pin(async move {
|
|
Box::pin(async move {
|
|
|
let timer = Timer::after(timeout);
|
|
let timer = Timer::after(timeout);
|
|
@@ -209,6 +237,8 @@ impl Acceptor {
|
|
|
let hosts = self.session.upgrade().unwrap().p2p().hosts();
|
|
let hosts = self.session.upgrade().unwrap().p2p().hosts();
|
|
|
let mut negotiations = FuturesUnordered::<PtNegotiation>::new();
|
|
let mut negotiations = FuturesUnordered::<PtNegotiation>::new();
|
|
|
let mut accepting = None;
|
|
let mut accepting = None;
|
|
|
|
|
+ let mut accept_retry = None;
|
|
|
|
|
+ let mut resource_backoff = ResourceExhaustionBackoff::new();
|
|
|
|
|
|
|
|
loop {
|
|
loop {
|
|
|
// Reserve capacity for established channels, transport negotiations,
|
|
// Reserve capacity for established channels, transport negotiations,
|
|
@@ -219,7 +249,7 @@ impl Acceptor {
|
|
|
negotiations.len() +
|
|
negotiations.len() +
|
|
|
usize::from(accepting.is_some());
|
|
usize::from(accepting.is_some());
|
|
|
|
|
|
|
|
- if reserved < limit && accepting.is_none() {
|
|
|
|
|
|
|
+ if reserved < limit && accepting.is_none() && accept_retry.is_none() {
|
|
|
accepting = Some(listener.next());
|
|
accepting = Some(listener.next());
|
|
|
}
|
|
}
|
|
|
|
|
|
|
@@ -231,6 +261,7 @@ impl Acceptor {
|
|
|
if negotiations.is_empty() {
|
|
if negotiations.is_empty() {
|
|
|
match accept.await {
|
|
match accept.await {
|
|
|
Ok(negotiation) => {
|
|
Ok(negotiation) => {
|
|
|
|
|
+ resource_backoff.reset();
|
|
|
negotiations
|
|
negotiations
|
|
|
.push(with_handshake_timeout(negotiation, handshake_timeout));
|
|
.push(with_handshake_timeout(negotiation, handshake_timeout));
|
|
|
continue
|
|
continue
|
|
@@ -243,6 +274,7 @@ impl Acceptor {
|
|
|
|
|
|
|
|
match select(accept, negotiation).await {
|
|
match select(accept, negotiation).await {
|
|
|
Either::Left((Ok(negotiation), _)) => {
|
|
Either::Left((Ok(negotiation), _)) => {
|
|
|
|
|
+ resource_backoff.reset();
|
|
|
negotiations
|
|
negotiations
|
|
|
.push(with_handshake_timeout(negotiation, handshake_timeout));
|
|
.push(with_handshake_timeout(negotiation, handshake_timeout));
|
|
|
continue
|
|
continue
|
|
@@ -258,6 +290,26 @@ impl Acceptor {
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
+ } else if let Some(retry) = accept_retry.take() {
|
|
|
|
|
+ if negotiations.is_empty() {
|
|
|
|
|
+ retry.await;
|
|
|
|
|
+ continue
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ let negotiation = negotiations.next();
|
|
|
|
|
+ pin_mut!(negotiation);
|
|
|
|
|
+
|
|
|
|
|
+ match select(retry, negotiation).await {
|
|
|
|
|
+ Either::Left((_, _)) => continue,
|
|
|
|
|
+ Either::Right((Some(result), retry)) => {
|
|
|
|
|
+ accept_retry = Some(retry);
|
|
|
|
|
+ result
|
|
|
|
|
+ }
|
|
|
|
|
+ Either::Right((None, retry)) => {
|
|
|
|
|
+ accept_retry = Some(retry);
|
|
|
|
|
+ continue
|
|
|
|
|
+ }
|
|
|
|
|
+ }
|
|
|
} else if let Some(result) = negotiations.next().await {
|
|
} else if let Some(result) = negotiations.next().await {
|
|
|
result
|
|
result
|
|
|
} else {
|
|
} else {
|
|
@@ -302,6 +354,17 @@ impl Acceptor {
|
|
|
self.channel_publisher.notify(Ok(channel)).await;
|
|
self.channel_publisher.notify(Ok(channel)).await;
|
|
|
}
|
|
}
|
|
|
|
|
|
|
|
|
|
+ Err(e) if is_descriptor_exhaustion(&e) => {
|
|
|
|
|
+ let delay = resource_backoff.next_delay();
|
|
|
|
|
+ warn!(
|
|
|
|
|
+ target: "net::acceptor::run_accept_loop",
|
|
|
|
|
+ "[P2P] Listener descriptor exhaustion: {e}; retrying accepts in {} ms",
|
|
|
|
|
+ delay.as_millis(),
|
|
|
|
|
+ );
|
|
|
|
|
+ accept_retry = Some(Timer::after(delay));
|
|
|
|
|
+ continue
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
// As per accept(2) recommendation:
|
|
// As per accept(2) recommendation:
|
|
|
Err(e) if e.raw_os_error().is_some() => match e.raw_os_error().unwrap() {
|
|
Err(e) if e.raw_os_error().is_some() => match e.raw_os_error().unwrap() {
|
|
|
libc::EAGAIN | libc::ECONNABORTED | libc::EPROTO | libc::EINTR => continue,
|
|
libc::EAGAIN | libc::ECONNABORTED | libc::EPROTO | libc::EINTR => continue,
|
|
@@ -386,3 +449,29 @@ impl Acceptor {
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
}
|
|
}
|
|
|
|
|
+
|
|
|
|
|
+#[cfg(test)]
|
|
|
|
|
+mod tests {
|
|
|
|
|
+ use super::*;
|
|
|
|
|
+
|
|
|
|
|
+ #[test]
|
|
|
|
|
+ fn descriptor_exhaustion_classification() {
|
|
|
|
|
+ assert!(is_descriptor_exhaustion(&std::io::Error::from_raw_os_error(libc::EMFILE)));
|
|
|
|
|
+ assert!(is_descriptor_exhaustion(&std::io::Error::from_raw_os_error(libc::ENFILE)));
|
|
|
|
|
+ assert!(!is_descriptor_exhaustion(&std::io::Error::from_raw_os_error(libc::EAGAIN)));
|
|
|
|
|
+ assert!(!is_descriptor_exhaustion(&std::io::Error::other("listener error")));
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ #[test]
|
|
|
|
|
+ fn descriptor_exhaustion_backoff_is_capped_and_resettable() {
|
|
|
|
|
+ let mut backoff = ResourceExhaustionBackoff::new();
|
|
|
|
|
+ let expected = [100, 200, 400, 800, 1600, 3200, 5000, 5000];
|
|
|
|
|
+
|
|
|
|
|
+ for millis in expected {
|
|
|
|
|
+ assert_eq!(backoff.next_delay(), Duration::from_millis(millis));
|
|
|
|
|
+ }
|
|
|
|
|
+
|
|
|
|
|
+ backoff.reset();
|
|
|
|
|
+ assert_eq!(backoff.next_delay(), ACCEPT_RETRY_MIN);
|
|
|
|
|
+ }
|
|
|
|
|
+}
|