Adapt runtime stream retransmit pacing
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
ci / test (push) Canceled after 0s
ci / fuzz-smoke (push) Canceled after 0s
ci / windows-client (push) Canceled after 0s
ci / package-release (linux-x86_64, ubuntu-latest) (push) Canceled after 0s
ci / package-release (macos-aarch64, macos-14) (push) Canceled after 0s
ci / package-release (macos-x86_64, macos-13) (push) Canceled after 0s
ci / package-release (windows-x86_64, windows-latest) (push) Canceled after 0s
ci / remote-bench (push) Canceled after 0s
ci / publish-gitea-release (push) Canceled after 0s
This commit is contained in:
+28
-20
@@ -29,8 +29,8 @@ use tokio::net::UdpSocket;
|
||||
|
||||
pub const DEFAULT_INITIAL_WINDOW: usize = 1024 * 1024;
|
||||
pub const DEFAULT_RETRANSMIT_AFTER: Duration = Duration::from_millis(200);
|
||||
const ADAPTIVE_RETRANSMIT_PAD: Duration = Duration::from_millis(10);
|
||||
const ADAPTIVE_RETRANSMIT_MIN: Duration = Duration::from_millis(10);
|
||||
pub const ADAPTIVE_RETRANSMIT_PAD: Duration = Duration::from_millis(10);
|
||||
pub const ADAPTIVE_RETRANSMIT_MIN: Duration = Duration::from_millis(10);
|
||||
pub const DEFAULT_KEEPALIVE_AFTER: Duration = Duration::from_secs(2);
|
||||
pub const DEFAULT_RETIRED_STREAM_TOMBSTONES: usize = 16 * 1024;
|
||||
pub const SERVICE_TARGET_PREFIX: &str = "@dosh-";
|
||||
@@ -553,17 +553,7 @@ impl StreamMux {
|
||||
}
|
||||
|
||||
pub fn effective_retransmit_after(&self) -> Duration {
|
||||
let Some(srtt) = self.srtt else {
|
||||
return self.config.retransmit_after;
|
||||
};
|
||||
let adaptive = srtt
|
||||
.saturating_add(ADAPTIVE_RETRANSMIT_PAD)
|
||||
.max(ADAPTIVE_RETRANSMIT_MIN);
|
||||
if self.config.retransmit_after < ADAPTIVE_RETRANSMIT_MIN {
|
||||
self.config.retransmit_after
|
||||
} else {
|
||||
adaptive.min(self.config.retransmit_after)
|
||||
}
|
||||
adaptive_retransmit_after(self.config.retransmit_after, self.srtt)
|
||||
}
|
||||
|
||||
pub fn is_open(&self, stream_id: u64) -> bool {
|
||||
@@ -709,13 +699,7 @@ impl StreamMux {
|
||||
}
|
||||
|
||||
fn observe_rtt_sample(&mut self, sample: Duration) {
|
||||
self.srtt = Some(match self.srtt {
|
||||
None => sample,
|
||||
Some(srtt) => {
|
||||
let smoothed = ((srtt.as_micros() * 7) + sample.as_micros()) / 8;
|
||||
Duration::from_micros(smoothed.min(u64::MAX as u128) as u64)
|
||||
}
|
||||
});
|
||||
observe_retransmit_rtt(&mut self.srtt, sample);
|
||||
}
|
||||
|
||||
fn add_credit(&mut self, stream_id: u64, bytes: usize) {
|
||||
@@ -761,6 +745,30 @@ impl StreamMux {
|
||||
}
|
||||
}
|
||||
|
||||
pub fn adaptive_retransmit_after(configured: Duration, srtt: Option<Duration>) -> Duration {
|
||||
let Some(srtt) = srtt else {
|
||||
return configured;
|
||||
};
|
||||
let adaptive = srtt
|
||||
.saturating_add(ADAPTIVE_RETRANSMIT_PAD)
|
||||
.max(ADAPTIVE_RETRANSMIT_MIN);
|
||||
if configured < ADAPTIVE_RETRANSMIT_MIN {
|
||||
configured
|
||||
} else {
|
||||
adaptive.min(configured)
|
||||
}
|
||||
}
|
||||
|
||||
pub fn observe_retransmit_rtt(srtt: &mut Option<Duration>, sample: Duration) {
|
||||
*srtt = Some(match *srtt {
|
||||
None => sample,
|
||||
Some(srtt) => {
|
||||
let smoothed = ((srtt.as_micros() * 7) + sample.as_micros()) / 8;
|
||||
Duration::from_micros(smoothed.min(u64::MAX as u128) as u64)
|
||||
}
|
||||
});
|
||||
}
|
||||
|
||||
pub struct DoshTransport {
|
||||
socket: Arc<UdpSocket>,
|
||||
role: SessionRole,
|
||||
|
||||
Reference in New Issue
Block a user