mirror of
https://github.com/n0-computer/noq.git
synced 2026-09-23 19:48:19 +00:00
fix(quinn-proto): improve per path logic in poll_transmit
This commit is contained in:
committed by
GitHub
parent
3507beef71
commit
e8a83c4e6c
@@ -539,20 +539,6 @@ impl Connection {
|
||||
// What about PATH_CHALLENGE or PATH_RESPONSE? We need to check if we need to send
|
||||
// any of those.
|
||||
|
||||
// TODO(flub): We only have PathId(0) for now. For multipath we need to figure
|
||||
// out which path we want to send the packet on before we start building it.
|
||||
let path_id = PathId(0);
|
||||
|
||||
let mut buf = TransmitBuf::new(
|
||||
buf,
|
||||
max_datagrams,
|
||||
self.path_data(path_id).current_mtu().into(),
|
||||
);
|
||||
|
||||
if let Some(challenge) = self.send_prev_path_challenge(now, &mut buf, path_id) {
|
||||
return Some(challenge);
|
||||
}
|
||||
|
||||
// Check whether we need to send a close message
|
||||
let close = match self.state {
|
||||
State::Drained => {
|
||||
@@ -602,10 +588,6 @@ impl Connection {
|
||||
let mut last_packet_number = None;
|
||||
|
||||
let mut path_id = *self.paths.first_key_value().expect("one path must exist").0;
|
||||
let mut space_id = match path_id {
|
||||
PathId(0) => SpaceId::Initial,
|
||||
_ => SpaceId::Data,
|
||||
};
|
||||
|
||||
// If there is any available path we only want to send frames to a backup path that
|
||||
// must be sent on that path.
|
||||
@@ -614,15 +596,29 @@ impl Connection {
|
||||
.values()
|
||||
.any(|path| path.data.status == PathStatus::Available);
|
||||
|
||||
// Setup for the first path_id
|
||||
let mut transmit = TransmitBuf::new(
|
||||
buf,
|
||||
max_datagrams,
|
||||
self.path_data(path_id).current_mtu().into(),
|
||||
);
|
||||
if let Some(challenge) = self.send_prev_path_challenge(now, &mut transmit, path_id) {
|
||||
return Some(challenge);
|
||||
}
|
||||
let mut space_id = match path_id {
|
||||
PathId(0) => SpaceId::Initial,
|
||||
_ => SpaceId::Data,
|
||||
};
|
||||
|
||||
loop {
|
||||
// Determine if anything can be sent in this packet number space (SpaceId +
|
||||
// PathId).
|
||||
let max_packet_size = if buf.datagram_remaining_mut() > 0 {
|
||||
let max_packet_size = if transmit.datagram_remaining_mut() > 0 {
|
||||
// We are trying to coalesce another packet into this datagram.
|
||||
buf.datagram_remaining_mut()
|
||||
transmit.datagram_remaining_mut()
|
||||
} else {
|
||||
// A new datagram needs to be started.
|
||||
buf.segment_size()
|
||||
transmit.segment_size()
|
||||
};
|
||||
let can_send = self.space_can_send(space_id, path_id, max_packet_size, close);
|
||||
let path_should_send = {
|
||||
@@ -644,9 +640,9 @@ impl Connection {
|
||||
continue;
|
||||
}
|
||||
|
||||
let send_blocked = if path_should_send && buf.datagram_remaining_mut() == 0 {
|
||||
let send_blocked = if path_should_send && transmit.datagram_remaining_mut() == 0 {
|
||||
// Only check congestion control if a new datagram is needed.
|
||||
self.path_congestion_check(space_id, path_id, &buf, &can_send, now)
|
||||
self.path_congestion_check(space_id, path_id, &transmit, &can_send, now)
|
||||
} else {
|
||||
PathBlocked::No
|
||||
};
|
||||
@@ -665,7 +661,7 @@ impl Connection {
|
||||
|
||||
// If there are any datagrams in the transmit, packets for another path can
|
||||
// not be built.
|
||||
if buf.num_datagrams() > 0 {
|
||||
if transmit.num_datagrams() > 0 {
|
||||
break;
|
||||
}
|
||||
|
||||
@@ -675,6 +671,15 @@ impl Connection {
|
||||
trace!(?space_id, ?path_id, ?next_path_id, "trying next path");
|
||||
path_id = *next_path_id;
|
||||
space_id = SpaceId::Data;
|
||||
|
||||
// update per path state
|
||||
transmit.set_segment_size(self.path_data(path_id).current_mtu().into());
|
||||
if let Some(challenge) =
|
||||
self.send_prev_path_challenge(now, &mut transmit, path_id)
|
||||
{
|
||||
return Some(challenge);
|
||||
}
|
||||
|
||||
continue;
|
||||
}
|
||||
None => {
|
||||
@@ -686,14 +691,14 @@ impl Connection {
|
||||
}
|
||||
|
||||
// If the datagram is full, we need to start a new one.
|
||||
if buf.datagram_remaining_mut() == 0 {
|
||||
if buf.num_datagrams() >= buf.max_datagrams() {
|
||||
if transmit.datagram_remaining_mut() == 0 {
|
||||
if transmit.num_datagrams() >= transmit.max_datagrams() {
|
||||
// No more datagrams allowed
|
||||
break;
|
||||
}
|
||||
|
||||
match self.spaces[space_id].for_path(path_id).loss_probes {
|
||||
0 => buf.start_new_datagram(),
|
||||
0 => transmit.start_new_datagram(),
|
||||
_ => {
|
||||
// We need something to send for a tail-loss probe.
|
||||
let request_immediate_ack =
|
||||
@@ -709,21 +714,21 @@ impl Connection {
|
||||
// Clamp the datagram to at most the minimum MTU to ensure that loss
|
||||
// probes can get through and enable recovery even if the path MTU
|
||||
// has shrank unexpectedly.
|
||||
buf.start_new_datagram_with_size(std::cmp::min(
|
||||
transmit.start_new_datagram_with_size(std::cmp::min(
|
||||
usize::from(INITIAL_MTU),
|
||||
buf.segment_size(),
|
||||
transmit.segment_size(),
|
||||
));
|
||||
}
|
||||
}
|
||||
trace!(count = buf.num_datagrams(), "new datagram started");
|
||||
trace!(count = transmit.num_datagrams(), "new datagram started");
|
||||
coalesce = true;
|
||||
pad_datagram = false;
|
||||
}
|
||||
|
||||
// If coalescing another packet into the existing datagram, there should
|
||||
// still be enough space for a whole packet.
|
||||
if buf.datagram_start_offset() < buf.len() {
|
||||
debug_assert!(buf.datagram_remaining_mut() >= MIN_PACKET_SPACE);
|
||||
if transmit.datagram_start_offset() < transmit.len() {
|
||||
debug_assert!(transmit.datagram_remaining_mut() >= MIN_PACKET_SPACE);
|
||||
}
|
||||
|
||||
//
|
||||
@@ -753,7 +758,7 @@ impl Connection {
|
||||
space_id,
|
||||
path_id,
|
||||
self.rem_cids.get(&path_id).unwrap().active(),
|
||||
&mut buf,
|
||||
&mut transmit,
|
||||
can_send.other,
|
||||
self,
|
||||
)?;
|
||||
@@ -855,10 +860,10 @@ impl Connection {
|
||||
},
|
||||
false,
|
||||
);
|
||||
self.stats.udp_tx.on_sent(1, buf.len());
|
||||
self.stats.udp_tx.on_sent(1, transmit.len());
|
||||
return Some(Transmit {
|
||||
destination: remote,
|
||||
size: buf.len(),
|
||||
size: transmit.len(),
|
||||
ecn: None,
|
||||
segment_size: None,
|
||||
src_ip: self.local_ip,
|
||||
@@ -969,8 +974,8 @@ impl Connection {
|
||||
|
||||
builder.finish_and_track(now, self, path_id, sent_frames, pad_datagram);
|
||||
|
||||
if buf.num_datagrams() == 1 {
|
||||
buf.clip_datagram_size();
|
||||
if transmit.num_datagrams() == 1 {
|
||||
transmit.clip_datagram_size();
|
||||
}
|
||||
}
|
||||
}
|
||||
@@ -980,15 +985,15 @@ impl Connection {
|
||||
// be the one from the highest packet space.
|
||||
self.path_data_mut(path_id).congestion.on_sent(
|
||||
now,
|
||||
buf.len() as u64,
|
||||
transmit.len() as u64,
|
||||
last_packet_number,
|
||||
);
|
||||
}
|
||||
|
||||
self.app_limited = buf.is_empty() && !congestion_blocked;
|
||||
self.app_limited = transmit.is_empty() && !congestion_blocked;
|
||||
|
||||
// Send MTU probe if necessary
|
||||
if buf.is_empty() && self.state.is_established() {
|
||||
if transmit.is_empty() && self.state.is_established() {
|
||||
let space_id = SpaceId::Data;
|
||||
let next_pn = self.spaces[space_id].for_path(path_id).peek_tx_number();
|
||||
let probe_size = self
|
||||
@@ -996,10 +1001,10 @@ impl Connection {
|
||||
.mtud
|
||||
.poll_transmit(now, next_pn)?;
|
||||
|
||||
debug_assert_eq!(buf.num_datagrams(), 0);
|
||||
buf.start_new_datagram_with_size(probe_size as usize);
|
||||
debug_assert_eq!(transmit.num_datagrams(), 0);
|
||||
transmit.start_new_datagram_with_size(probe_size as usize);
|
||||
|
||||
debug_assert_eq!(buf.datagram_start_offset(), 0);
|
||||
debug_assert_eq!(transmit.datagram_start_offset(), 0);
|
||||
// TODO(flub): I'm not particularly happy about this unwrap. But let's leave it
|
||||
// for now until more stuff is settled. We probably should check earlier on
|
||||
// in poll_transmit that we have a valid CID to use.
|
||||
@@ -1008,7 +1013,7 @@ impl Connection {
|
||||
space_id,
|
||||
path_id,
|
||||
self.rem_cids.get(&path_id).unwrap().active(),
|
||||
&mut buf,
|
||||
&mut transmit,
|
||||
true,
|
||||
self,
|
||||
)?;
|
||||
@@ -1037,34 +1042,35 @@ impl Connection {
|
||||
trace!(?probe_size, "writing MTUD probe");
|
||||
}
|
||||
|
||||
if buf.is_empty() {
|
||||
if transmit.is_empty() {
|
||||
return None;
|
||||
}
|
||||
|
||||
trace!(
|
||||
segment_size = buf.segment_size(),
|
||||
last_datagram_len = buf.len() % buf.segment_size(),
|
||||
segment_size = transmit.segment_size(),
|
||||
last_datagram_len = transmit.len() % transmit.segment_size(),
|
||||
"sending {} bytes in {} datagrams",
|
||||
buf.len(),
|
||||
buf.num_datagrams()
|
||||
transmit.len(),
|
||||
transmit.num_datagrams()
|
||||
);
|
||||
self.path_data_mut(path_id).inc_total_sent(buf.len() as u64);
|
||||
self.path_data_mut(path_id)
|
||||
.inc_total_sent(transmit.len() as u64);
|
||||
|
||||
self.stats
|
||||
.udp_tx
|
||||
.on_sent(buf.num_datagrams() as u64, buf.len());
|
||||
.on_sent(transmit.num_datagrams() as u64, transmit.len());
|
||||
|
||||
Some(Transmit {
|
||||
destination: self.path_data(path_id).remote,
|
||||
size: buf.len(),
|
||||
size: transmit.len(),
|
||||
ecn: if self.path_data(path_id).sending_ecn {
|
||||
Some(EcnCodepoint::Ect0)
|
||||
} else {
|
||||
None
|
||||
},
|
||||
segment_size: match buf.num_datagrams() {
|
||||
segment_size: match transmit.num_datagrams() {
|
||||
1 => None,
|
||||
_ => Some(buf.segment_size()),
|
||||
_ => Some(transmit.segment_size()),
|
||||
},
|
||||
src_ip: self.local_ip,
|
||||
})
|
||||
|
||||
@@ -68,6 +68,15 @@ impl<'a> TransmitBuf<'a> {
|
||||
}
|
||||
}
|
||||
|
||||
pub(super) fn set_segment_size(&mut self, mtu: usize) {
|
||||
debug_assert!(
|
||||
self.datagram_start == 0 && self.buf_capacity == 0 && self.num_datagrams == 0,
|
||||
"can only change the segment size if nothing has been written yet"
|
||||
);
|
||||
|
||||
self.segment_size = mtu;
|
||||
}
|
||||
|
||||
/// Starts a datagram with a custom datagram size
|
||||
///
|
||||
/// This is a specialized version of [`TransmitBuf::start_new_datagram`] which sets the
|
||||
|
||||
Reference in New Issue
Block a user