From e8a83c4e6c79ba019b74fcbbf7e77b93ef4f1033 Mon Sep 17 00:00:00 2001 From: Friedel Ziegelmayer Date: Wed, 28 May 2025 13:01:43 +0200 Subject: [PATCH] fix(quinn-proto): improve per path logic in poll_transmit --- quinn-proto/src/connection/mod.rs | 114 +++++++++++---------- quinn-proto/src/connection/transmit_buf.rs | 9 ++ 2 files changed, 69 insertions(+), 54 deletions(-) diff --git a/quinn-proto/src/connection/mod.rs b/quinn-proto/src/connection/mod.rs index bfd93996d..c918be5d8 100644 --- a/quinn-proto/src/connection/mod.rs +++ b/quinn-proto/src/connection/mod.rs @@ -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, }) diff --git a/quinn-proto/src/connection/transmit_buf.rs b/quinn-proto/src/connection/transmit_buf.rs index df59a89ba..5a5f1f4a2 100644 --- a/quinn-proto/src/connection/transmit_buf.rs +++ b/quinn-proto/src/connection/transmit_buf.rs @@ -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