// Tracking of received packets and generating ACKs thereof.
use std::{
cmp::min,
collections::VecDeque,
fmt::{self, Display, Formatter},
time::{Duration, Instant},
};
use enum_map::{Enum, EnumMap}; use enumset::{EnumSet, EnumSetType}; use log::{Level, log_enabled}; use neqo_common::{Buffer, Ecn, MAX_VARINT, qdebug, qtrace, qwarn}; use nss::Epoch; use smallvec::SmallVec; use strum::{Display, EnumIter};
impl PacketRange { /// Make a single packet range. pubconstfn new(pn: packet::Number) -> Self { Self {
largest: pn,
smallest: pn,
ack_needed: true,
}
}
/// Get the number of acknowledged packets in the range. pubconstfn len(&self) -> u64 { self.largest - self.smallest + 1
}
/// Returns whether this needs to be sent. pubconstfn ack_needed(&self) -> bool { self.ack_needed
}
/// Return whether the given number is in the range. pubconstfn contains(&self, pn: packet::Number) -> bool {
(pn >= self.smallest) && (pn <= self.largest)
}
/// Maybe add a packet number to the range. Returns true if it was added /// at the small end (which indicates that this might need merging with a /// preceding range). pubfn add(&mutself, pn: packet::Number) -> InsertionResult {
assert!(!self.contains(pn)); // Only insert if this is adjacent the current range. if (self.largest + 1) == pn {
qtrace!("[{self}] Adding largest {pn}"); self.largest += 1; self.ack_needed = true;
InsertionResult::Largest
} elseifself.smallest == (pn + 1) {
qtrace!("[{self}] Adding smallest {pn}"); self.smallest -= 1; self.ack_needed = true;
InsertionResult::Smallest
} else {
InsertionResult::NotInserted
}
}
/// Maybe merge a higher-numbered range into this. fn merge_larger(&mutself, other: &Self) {
qdebug!("[{self}] Merging {other}"); // This only works if they are immediately adjacent.
assert_eq!(self.largest + 1, other.smallest);
/// When a packet containing the range `other` is acknowledged, /// clear the `ack_needed` attribute on this. /// Requires that other is equal to this, or a larger range. pubconstfn acknowledged(&mutself, other: &Self) { if (other.smallest <= self.smallest) && (other.largest >= self.largest) { self.ack_needed = false;
}
}
}
/// The default maximum ACK delay we use locally and advertise to the remote. pubconst DEFAULT_LOCAL_ACK_DELAY: Duration = Duration::from_millis(20); /// The default maximum ACK delay we assume the remote uses. /// /// > If this value is absent, a default of 25 milliseconds is assumed. /// /// <https://datatracker.ietf.org/doc/html/rfc9000#section-18.2> pubconst DEFAULT_REMOTE_ACK_DELAY: Duration = Duration::from_millis(25); /// The default number of in-order packets we will receive after /// largest acknowledged without sending an immediate acknowledgment. pubconst DEFAULT_ACK_PACKET_TOLERANCE: packet::Number = 1; const MAX_TRACKED_RANGES: usize = 32; const MAX_ACKS_PER_FRAME: usize = 32;
/// A structure that tracks what was included in an ACK. #[derive(Debug, Clone)] pubstruct AckToken {
space: PacketNumberSpace,
ranges: Box<[PacketRange]>,
}
impl AckToken { /// Get the space for this token. pubconstfn space(&self) -> PacketNumberSpace { self.space
}
}
/// A structure that tracks what packets have been received, /// and what needs acknowledgement for a packet number space. #[derive(Debug)] pubstruct RecvdPackets {
space: PacketNumberSpace,
ranges: VecDeque<PacketRange>, /// The packet number of the lowest number packet that we are tracking.
min_tracked: packet::Number, /// The time we got the largest acknowledged.
largest_pn_time: Option<Instant>, /// The time that we should be sending an ACK.
ack_time: Option<Instant>, /// The time we last sent an ACK.
last_ack_time: Option<Instant>, /// The current ACK frequency sequence number.
ack_frequency_seqno: u64, /// The time to delay after receiving the first packet that is /// not immediately acknowledged.
ack_delay: Duration, /// The number of ack-eliciting packets that have been received, but /// not acknowledged.
unacknowledged_count: packet::Number, /// The number of contiguous packets that can be received without /// acknowledging immediately.
unacknowledged_tolerance: packet::Number, /// Whether we are ignoring packets that arrive out of order /// for the purposes of generating immediate acknowledgment.
ignore_order: bool, // The counts of different ECN marks that have been received.
ecn_count: ecn::Count,
}
impl RecvdPackets { /// Make a new `RecvdPackets` for the indicated packet number space. pubfn new(space: PacketNumberSpace) -> Self { Self {
space,
ranges: VecDeque::new(),
min_tracked: 0,
largest_pn_time: None,
ack_time: None,
last_ack_time: None,
ack_frequency_seqno: 0,
ack_delay: DEFAULT_LOCAL_ACK_DELAY,
unacknowledged_count: 0,
unacknowledged_tolerance: if space == PacketNumberSpace::ApplicationData {
DEFAULT_ACK_PACKET_TOLERANCE
} else { // ACK more aggressively 0
},
ignore_order: false,
ecn_count: ecn::Count::default(),
}
}
/// Get the ECN counts. pubconstfn ecn_marks(&mutself) -> &mut ecn::Count {
&mutself.ecn_count
}
/// Get the time at which the next ACK should be sent. pubconstfn ack_time(&self) -> Option<Instant> { self.ack_time
}
/// Update acknowledgment delay parameters. pubconstfn ack_freq(
&mutself,
seqno: u64,
tolerance: packet::Number,
delay: Duration,
ignore_order: bool,
) { // Yes, this means that we will overwrite values if a sequence number is // reused, but that is better than using an `Option<packet::Number>` // when it will always be `Some`. if seqno >= self.ack_frequency_seqno { self.ack_frequency_seqno = seqno; self.unacknowledged_tolerance = tolerance; self.ack_delay = delay; self.ignore_order = ignore_order;
}
}
/// Returns true if an ACK frame should be sent now. fn ack_now(&self, now: Instant, rtt: Duration) -> bool { // If ack_time is Some, then we have something to acknowledge. // In that case, either ack because `now >= ack_time`, or // because it is more than an RTT since the last time we sent an ack. self.ack_time.is_some_and(|next| {
next <= now || self.last_ack_time.is_some_and(|last| last + rtt <= now)
})
}
// A simple addition of a packet number to the tracked set. // This doesn't do a binary search on the assumption that // new packets will generally be added to the start of the list. fn add(&mutself, pn: packet::Number) -> Res<()> { for i in0..self.ranges.len() { matchself.ranges[i].add(pn) {
InsertionResult::Largest => return Ok(()),
InsertionResult::Smallest => { // If this was the smallest, it might have filled a gap. let nxt = i + 1; if (nxt < self.ranges.len()) && (pn - 1 == self.ranges[nxt].largest) { let larger = self.ranges.remove(i).ok_or(Error::Internal)?; self.ranges[i].merge_larger(&larger);
} return Ok(());
}
InsertionResult::NotInserted => { ifself.ranges[i].largest < pn { self.ranges.insert(i, PacketRange::new(pn)); return Ok(());
}
}
}
} self.ranges.push_back(PacketRange::new(pn));
Ok(())
}
fn trim_ranges(&mutself, stats: &mut Stats) -> Res<()> { // Limit the number of ranges that are tracked to MAX_TRACKED_RANGES. ifself.ranges.len() > MAX_TRACKED_RANGES { let oldest = self.ranges.pop_back().ok_or(Error::Internal)?; if oldest.ack_needed {
qwarn!("[{self}] Dropping unacknowledged ACK range: {oldest}");
stats.unacked_range_dropped += 1;
} else {
qdebug!("[{self}] Drop ACK range: {oldest}");
} self.min_tracked = oldest.largest + 1;
}
Ok(())
}
/// Add the packet to the tracked set. /// Return true if the packet was the largest received so far. pubfn set_received(
&mutself,
now: Instant,
pn: packet::Number,
ack_eliciting: bool,
stats: &mut Stats,
) -> Res<bool> { let next_in_order_pn = self.ranges.front().map_or(0, |r| r.largest + 1);
qtrace!("[{self}] received {pn}, next: {next_in_order_pn}");
self.add(pn)?; self.trim_ranges(stats)?;
// The new addition was the largest, so update the time we use for calculating ACK delay. let largest = if pn >= next_in_order_pn { self.largest_pn_time = Some(now); true
} else { false
};
if ack_eliciting { self.unacknowledged_count += 1;
let ack_time = if immediate_ack {
now
} else { // Note that `ack_delay` can change and that won't take effect if // we are waiting on the previous delay timer. // If ACK delay increases, we might send an ACK a bit early; // if ACK delay decreases, we might send an ACK a bit later. // We could use min() here, but change is rare and the size // of the change is very small. self.ack_time.unwrap_or_else(|| now + self.ack_delay)
};
qdebug!("[{self}] Set ACK timer to {ack_time:?}"); self.ack_time = Some(ack_time);
}
Ok(largest)
}
/// If we just received a PING frame, we should immediately acknowledge. pubfn immediate_ack(&mutself, now: Instant) { self.ack_time = Some(now);
qdebug!("[{self}] immediate_ack at {now:?}");
}
/// Check if the packet is a duplicate. pubfn is_duplicate(&self, pn: packet::Number) -> bool { if pn < self.min_tracked { returntrue;
} self.ranges
.iter()
.take_while(|r| pn <= r.largest)
.any(|r| r.contains(pn))
}
/// Mark the given range as having been acknowledged. pubfn acknowledged(&mutself, acked: &[PacketRange]) { letmut range_iter = self.ranges.iter_mut(); letmut cur = range_iter.next().expect("should have at least one range"); for ack in acked { while cur.smallest > ack.largest { let Some(next) = range_iter.next() else { return;
};
cur = next;
}
cur.acknowledged(ack);
}
}
/// Length of the worst possible ACK frame, assuming only one range and ECN counts. /// Note that this assumes one byte for the type and count of extra ranges. pubconst USEFUL_ACK_LEN: usize = 1 + 8 + 8 + 1 + 8 + 3 * 8;
/// Generate an ACK frame for this packet number space. /// /// Unlike other frame generators this doesn't modify the underlying instance /// to track what has been sent. This only clears the delayed ACK timer. /// /// When sending ACKs, we want to always send the most recent ranges, /// even if they have been sent in other packets. /// /// We don't send ranges that have been acknowledged, but they still need /// to be tracked so that duplicates can be detected. fn write_frame<B: Buffer>(
&mutself,
now: Instant,
rtt: Duration,
builder: &mut packet::Builder<B>,
tokens: &mut recovery::Tokens,
stats: &mut FrameStats,
) { // Check that we aren't delaying ACKs. if !self.ack_now(now, rtt) { return;
}
// Drop extra ACK ranges to fit the available space. Do this based on // a worst-case estimate of frame size for simplicity. // // When congestion limited, ACK-only packets are 255 bytes at most // (`recovery::ACK_ONLY_SIZE_LIMIT - 1`). This results in limiting the // ranges to 13 here. let max_ranges = iflet Some(avail) = builder.remaining().checked_sub(Self::USEFUL_ACK_LEN)
{ // Apply a hard maximum to keep plenty of space for other stuff.
min(1 + (avail / 16), MAX_ACKS_PER_FRAME)
} else { return;
};
let ranges = self
.ranges
.iter()
.filter(|r| r.ack_needed())
.take(max_ranges)
.cloned()
.collect::<SmallVec<[_; MAX_TRACKED_RANGES]>>(); if ranges.is_empty() { return;
}
letmut iter = ranges.iter(); let Some(first) = iter.next() else { return };
stats.largest_acknowledged = first.largest;
stats.ack += 1;
let Some(largest_pn_time) = self.largest_pn_time else { return;
}; let elapsed = now.duration_since(largest_pn_time); // We use the default exponent, so delay is in multiples of 8 microseconds. let ack_delay = u64::try_from(elapsed.as_micros() / 8).unwrap_or(u64::MAX); let ack_delay = min(MAX_VARINT, ack_delay); let Ok(extra_ranges) = u64::try_from(ranges.len() - 1) else { return;
};
builder.encode_frame( ifself.ecn_count.is_some() {
FrameType::AckEcn
} else {
FrameType::Ack
},
|b| {
b.encode_varint(first.largest);
b.encode_varint(ack_delay);
b.encode_varint(extra_ranges); // extra ranges
b.encode_varint(first.len() - 1); // first range
letmut last = first.smallest; for r in iter { // The difference must be at least 2 because 0-length gaps, // (difference 1) are illegal.
b.encode_varint(last - r.largest - 2); // Gap
b.encode_varint(r.len() - 1); // Range
last = r.smallest;
}
/// Force an ACK to be generated immediately. pubfn immediate_ack(&mutself, space: PacketNumberSpace, now: Instant) { iflet Some(space) = self.get_mut(space) {
space.immediate_ack(now);
}
}
/// Determine the earliest time that an ACK might be needed. pubfn ack_time(&self, now: Instant) -> Option<Instant> { if log_enabled!(Level::Trace) { for (space, recvd) in &self.spaces { iflet Some(recvd) = recvd {
qtrace!("ack_time for {space} = {:?}", recvd.ack_time());
}
}
} ifself.spaces[PacketNumberSpace::Initial].is_none()
&& self.spaces[PacketNumberSpace::Handshake].is_none()
&& let Some(recvd) = &self.spaces[PacketNumberSpace::ApplicationData]
{ return recvd.ack_time();
}
// Ignore any time that is in the past relative to `now`. // That is something of a hack, but there are cases where we can't send ACK // frames for all spaces, which can mean that one space is stuck in the past. // That isn't a problem because we guarantee that earlier spaces will always // be able to send ACK frames. self.spaces
.values()
.flatten()
.filter_map(|recvd| recvd.ack_time().filter(|t| *t > now))
.min()
}
fn test_ack_range(pns: &[packet::Number], nranges: usize) { letmut rp = RecvdPackets::new(PacketNumberSpace::Initial); // Any space will do.
assert_eq!(rp.to_string(), "Recvd-in"); letmut packets = HashSet::new();
for pn in pns {
rp.set_received(now(), *pn, true, &mut Stats::default())
.unwrap();
packets.insert(*pn);
}
assert_eq!(rp.ranges.len(), nranges);
// Check that all these packets will be detected as duplicates. for pn in pns {
assert!(rp.is_duplicate(*pn));
}
// Check that the ranges decrease monotonically and don't overlap. letmut iter = rp.ranges.iter(); letmut last = iter.next().expect("should have at least one"); for n in iter {
assert!(n.largest + 1 < last.smallest);
last = n;
}
// Check that the ranges include the right values. letmut in_ranges = HashSet::new(); for range in &rp.ranges { for included in range.smallest..=range.largest {
in_ranges.insert(included);
}
}
assert_eq!(packets, in_ranges);
}
// Even though the range was dropped, we still consider it a duplicate.
assert!(rp.is_duplicate(0));
assert!(!rp.is_duplicate(1));
assert!(rp.is_duplicate(2));
}
#[test] fn ack_delay() { const COUNT: packet::Number = 9; const DELAY: Duration = Duration::from_millis(7); letmut stats = Stats::default(); // Only application data packets are delayed. letmut rp = RecvdPackets::new(PacketNumberSpace::ApplicationData);
assert!(rp.ack_time().is_none());
assert!(!rp.ack_now(now(), RTT));
rp.ack_freq(0, COUNT, DELAY, false);
// Some packets won't cause an ACK to be needed. for i in0..COUNT {
rp.set_received(now(), i, true, &mut stats).unwrap();
assert_eq!(Some(now() + DELAY), rp.ack_time());
assert!(!rp.ack_now(now(), RTT));
assert!(rp.ack_now(now() + DELAY, RTT));
}
// Exceeding COUNT will move the ACK time to now.
rp.set_received(now(), COUNT, true, &mut stats).unwrap();
assert_eq!(Some(now()), rp.ack_time());
assert!(rp.ack_now(now(), RTT));
}
#[test] fn no_ack_delay() { letmut stats = Stats::default(); for space in &[PacketNumberSpace::Initial, PacketNumberSpace::Handshake] { letmut rp = RecvdPackets::new(*space);
assert!(rp.ack_time().is_none());
assert!(!rp.ack_now(now(), RTT));
// Any packet in these spaces is acknowledged straight away.
rp.set_received(now(), 0, true, &mut stats).unwrap();
assert_eq!(Some(now()), rp.ack_time());
assert!(rp.ack_now(now(), RTT));
}
}
// Filling in behind the largest acknowledged causes immediate ACK.
rp.set_received(now(), 0, true, &mut stats).unwrap();
write_frame(&mut rp);
// A new packet ordinarily doesn't result in an ACK, but this time it does.
rp.set_received(now() + RTT, 2, true, &mut stats).unwrap();
write_frame_at(&mut rp, now() + RTT);
}
/// Test that an in-order packet that is not ack-eliciting doesn't /// increase the number of packets needed to cause an ACK. #[test] fn non_ack_eliciting_skip() { letmut stats = Stats::default(); letmut rp = RecvdPackets::new(PacketNumberSpace::ApplicationData);
rp.ack_freq(0, 1, Duration::from_millis(10), true);
// This should be ignored.
rp.set_received(now(), 0, false, &mut stats).unwrap();
assert_ne!(Some(now()), rp.ack_time()); // Skip 1 (it has no effect).
rp.set_received(now(), 2, true, &mut stats).unwrap();
assert_ne!(Some(now()), rp.ack_time());
rp.set_received(now(), 3, true, &mut stats).unwrap();
assert_eq!(Some(now()), rp.ack_time());
}
/// If a packet that is not ack-eliciting is reordered, that's fine too. #[test] fn non_ack_eliciting_reorder() { letmut stats = Stats::default(); letmut rp = RecvdPackets::new(PacketNumberSpace::ApplicationData);
rp.ack_freq(0, 1, Duration::from_millis(10), false);
// These are out of order, but they are not ack-eliciting.
rp.set_received(now(), 1, false, &mut stats).unwrap();
assert_ne!(Some(now()), rp.ack_time());
rp.set_received(now(), 0, false, &mut stats).unwrap();
assert_ne!(Some(now()), rp.ack_time());
// These are in order.
rp.set_received(now(), 2, true, &mut stats).unwrap();
assert_ne!(Some(now()), rp.ack_time());
rp.set_received(now(), 3, true, &mut stats).unwrap();
assert_eq!(Some(now()), rp.ack_time());
}
// This should be delayed.
tracker
.get_mut(PacketNumberSpace::ApplicationData)
.unwrap()
.set_received(now(), 0, true, &mut stats)
.unwrap();
assert_eq!(Some(now() + DELAY), tracker.ack_time(now()));
// This should move the time forward. let later = now() + (DELAY / 2);
tracker
.get_mut(PacketNumberSpace::Initial)
.unwrap()
.set_received(later, 0, true, &mut stats)
.unwrap();
assert_eq!(Some(later), tracker.ack_time(now()));
}
#[test] fn drop_spaces() { letmut stats = Stats::default(); letmut tracker = AckTracker::default(); letmut builder =
packet::Builder::short(Encoder::default(), false, None::<&[u8]>, packet::LIMIT);
tracker
.get_mut(PacketNumberSpace::Initial)
.unwrap()
.set_received(now(), 0, true, &mut stats)
.unwrap(); // The reference time for `ack_time` has to be in the past or we filter out the timer.
assert!(
tracker
.ack_time(now().checked_sub(Duration::from_millis(1)).unwrap())
.is_some()
);
// Mark another packet as received so we have cause to send another ACK in that space.
tracker
.get_mut(PacketNumberSpace::Initial)
.unwrap()
.set_received(now(), 1, true, &mut stats)
.unwrap();
assert!(
tracker
.ack_time(now().checked_sub(Duration::from_millis(1)).unwrap())
.is_some()
);
// Now drop that space.
tracker.drop_space(PacketNumberSpace::Initial);
assert!(tracker.get_mut(PacketNumberSpace::Initial).is_none());
assert!(
tracker
.ack_time(now().checked_sub(Duration::from_millis(1)).unwrap())
.is_none()
);
tracker.write_frame(
PacketNumberSpace::Initial,
now(),
RTT,
&mut builder,
&mut tokens,
&mut frame_stats,
);
assert_eq!(frame_stats.ack, 1); iflet recovery::Token::Ack(tok) = &tokens[0] {
tracker.acked(tok); // Should be a noop.
} else {
panic!("not an ACK token");
}
}
/// ACK delay encodes elapsed microseconds divided by 8 (the default ACK delay exponent). #[test] fn ack_delay_encoding() { let t = now(); // 16µs → ack_delay = 16/8 = 2. let elapsed = Duration::from_micros(16);
letmut builder =
packet::Builder::short(Encoder::default(), false, None::<&[u8]>, packet::LIMIT); letmut stats = FrameStats::default();
tracker.write_frame(
PacketNumberSpace::Initial,
t + elapsed,
RTT,
&mut builder,
&mut recovery::Tokens::new(),
&mut stats,
);
assert_eq!(stats.ack, 1, "ACK frame should have been written");
// Decode the ACK frame from the builder output and check ack_delay. let enc: Encoder = builder.into(); let bytes = Vec::from(enc); // Skip the 1-byte packet header; remainder is the ACK frame. letmut dec = Decoder::from(&bytes[1..]); let frame = Frame::decode(&mut dec).unwrap(); let Frame::Ack { ack_delay, .. } = frame else {
panic!("expected ACK frame, got {frame:?}");
};
assert_eq!(ack_delay, 2, "ack_delay must be 16\u{b5}s / 8 = 2");
}
letmut builder =
packet::Builder::short(Encoder::default(), false, None::<&[u8]>, packet::LIMIT); // The code pessimistically assumes that each range needs 16 bytes to express. // So this won't be enough for a second range.
builder.set_limit(RecvdPackets::USEFUL_ACK_LEN + 8);
// While we have multiple PN spaces, we ignore ACK timers from the past. // Send out of order to cause the delayed ack timer to be set to `now()`.
tracker
.get_mut(PacketNumberSpace::ApplicationData)
.unwrap()
.set_received(now(), 3, true, &mut Stats::default())
.unwrap();
assert!(tracker.ack_time(now() + Duration::from_millis(1)).is_none());
// When we are reduced to one space, that filter is off.
tracker.drop_space(PacketNumberSpace::Initial);
tracker.drop_space(PacketNumberSpace::Handshake);
assert_eq!(
tracker.ack_time(now() + Duration::from_millis(1)),
Some(now())
);
}
// Test partial overlap: other.smallest <= self.smallest but other.largest < self.largest. // Should NOT clear ack_needed. letmut r2 = PacketRange::new(5);
r2.add(6); // range is now 5..=6
assert_eq!(r2.to_string(), "6->5");
assert_eq!(r2.len(), 2);
r2.acknowledged(&PacketRange {
largest: 5,
smallest: 0,
ack_needed: false,
});
assert!(r2.ack_needed()); // Should still need ack
}
#[test] fn trim_ranges_increments_stat() { // Each dropped range increments the stat by 1. letmut rp = RecvdPackets::new(PacketNumberSpace::Initial); letmut stats = Stats::default(); // Fill with MAX_TRACKED_RANGES + 2 disjoint ack-eliciting packets. for i in0..=(MAX_TRACKED_RANGES + 1) {
rp.set_received(now(), (i * 2) as u64, true, &mut stats)
.unwrap();
} // Two ranges should have been dropped.
assert_eq!(stats.unacked_range_dropped, 2);
}
#[test] fn acknowledged_tracks_duplicates() { letmut rp = RecvdPackets::new(PacketNumberSpace::ApplicationData); letmut stats = Stats::default(); // Receive packets 0, 1, 2 — one contiguous range, ACK needed. for pn in0u64..3 {
rp.set_received(now(), pn, true, &mut stats).unwrap();
}
assert!(
rp.ack_time().is_some(), "ACK should be needed before acknowledging"
);
// Simulate peer acknowledging our ACK by calling acknowledged. let acked = [PacketRange::new(0)]; // covers packet 0
rp.acknowledged(&acked); // Ranges are still tracked for duplicate detection after acknowledgement.
assert!(rp.is_duplicate(0));
assert!(rp.is_duplicate(2));
}
}
Die Informationen auf dieser Webseite wurden
nach bestem Wissen sorgfältig zusammengestellt. Es wird jedoch weder Vollständigkeit, noch Richtigkeit,
noch Qualität der bereit gestellten Informationen zugesichert.
Bemerkung:
Die farbliche Syntaxdarstellung und die Messung sind noch experimentell.