From e5655cc1bd02424b330578ae3657a7f83e8f9130 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sat, 26 Sep 2026 10:40:46 -0700 Subject: [PATCH 1/4] fix(proto): make BBR respond to classic ECN A CE-bearing ACK reached BBR as a lost-packet call with zero lost bytes, so Startup and ProbeUp, which skip the short-term loss response and see no lost bytes, kept accelerating into a marking bottleneck until it dropped. BBR now answers new CE feedback once per recovery episode: Startup exits and drains, as on high loss; Refill and ProbeUp stop the probe and bound inflight_longterm by the marked packet's tx_in_flight, as when loss shows inflight too high; every other state lowers the short-term model at the end of the round, as it does for loss. A mark adds no lost bytes or loss events, and an episode with CE is not undone as spurious. Co-Authored-By: Claude Opus 5.5 --- noq-proto/src/congestion.rs | 3 +- noq-proto/src/congestion/bbr3/mod.rs | 377 +++++++++++++++++++++++++-- noq-proto/src/tests/mod.rs | 153 +++++++++-- noq-proto/src/tests/util.rs | 97 +++++-- 4 files changed, 575 insertions(+), 55 deletions(-) diff --git a/noq-proto/src/congestion.rs b/noq-proto/src/congestion.rs index fa3361669..d374104c4 100644 --- a/noq-proto/src/congestion.rs +++ b/noq-proto/src/congestion.rs @@ -131,7 +131,8 @@ pub trait Controller: Send + Sync + std::fmt::Debug { /// lost. /// `lost_bytes` indicates how many bytes were lost. This value will be 0 for ECN triggers. /// `largest_lost` identifies the packet with the highest packet number in the congestion - /// event. + /// event, and `sent` is its send time. For ECN, that is the largest packet acknowledged by + /// the ACK whose CE count increased, and the event fires only on such an increase. /// /// Defaults to the deprecated [`Self::on_congestion_event`], dropping the space. fn on_congestion_event_space( diff --git a/noq-proto/src/congestion/bbr3/mod.rs b/noq-proto/src/congestion/bbr3/mod.rs index f8b336203..57bd3ca27 100644 --- a/noq-proto/src/congestion/bbr3/mod.rs +++ b/noq-proto/src/congestion/bbr3/mod.rs @@ -540,6 +540,9 @@ pub struct Bbr3 { loss_round_delivered: u64, /// equivalent to BBR.loss_in_round: flag set to true when loss occurs during the round loss_in_round: bool, + /// equivalent to Linux BBRv3's `ecn_in_round`: set when CE marks are reported during the + /// round, so the short-term model responds to them as it does to loss + ce_in_round: bool, /// equivalent to BBR.loss_events_in_round: count of discontiguous loss events /// observed in the current round trip, used by the STARTUP high-loss exit /// (BBRStartupFullLossCnt criterion). Reset at each loss-round boundary. @@ -563,6 +566,10 @@ pub struct Bbr3 { /// high-loss exit. Cleared once a packet sent after `recovery_start_time` is acknowledged. /// in_recovery: bool, + /// Whether BBR has responded to CE marks in the current recovery episode. It limits that + /// response to once per episode, and a marked episode's congestion is real, so its losses + /// are not undone as spurious. + ce_in_recovery: bool, /// equivalent to T_reno_bound: round-trip bound for the Reno-coexistence probe timer, /// re-picked from [`RENO_ROUNDS_BOUNDS`] each time the probe wait is randomized /// @@ -704,9 +711,11 @@ impl Bbr3 { recovery_start_time: None, recovery_start_round: 0, in_recovery: false, + ce_in_recovery: false, reno_rounds_bound: RENO_ROUNDS_BOUNDS[0], loss_round_delivered: 0, loss_in_round: false, + ce_in_round: false, probe_rtt_done_stamp: None, probe_rtt_round_done: false, prior_cwnd: 0, @@ -760,6 +769,7 @@ impl Bbr3 { } self.adapt_lower_bounds_from_congestion(); self.loss_in_round = false; + self.ce_in_round = false; } /// equivalent to BBRUpdateMaxBw @@ -799,7 +809,7 @@ impl Bbr3 { | BbrState::ProbeBw(ProbeBwSubstate::Up) | BbrState::Startup => {} _ => { - if self.loss_in_round { + if self.loss_in_round || self.ce_in_round { self.init_lower_bounds(); self.loss_lower_bounds(); } @@ -941,6 +951,7 @@ impl Bbr3 { self.recovery_start_time = Some(now); self.recovery_start_round = self.round_count; self.in_recovery = true; + self.ce_in_recovery = false; } /// Fast recovery ends when a packet sent after it began is acknowledged. @@ -1559,9 +1570,7 @@ impl Bbr3 { /// number is not the successor of the previous one. /// fn note_loss(&mut self, space: SpaceKind, packet_number: u64) { - if !self.loss_in_round { - self.loss_round_delivered = self.delivered; - } + self.start_congestion_round(); self.loss_in_round = true; let continues_range = self .last_lost_packet @@ -1572,6 +1581,58 @@ impl Bbr3 { self.last_lost_packet = Some((space, packet_number)); } + /// Opens the loss round on its first congestion signal, loss or CE, as draft-06's NoteLoss + /// does for loss. + fn start_congestion_round(&mut self) { + if !self.loss_in_round && !self.ce_in_round { + self.loss_round_delivered = self.delivered; + } + } + + /// Responds to newly reported CE marks as classic ECN: congestion that the bottleneck + /// signalled instead of dropping packets. + /// + /// Draft-06 section 3.7 requires treating CE as congestion without prescribing BBR's + /// response. As RFC 9002 reduces its window once per recovery period, this responds once per + /// recovery episode: Startup stops as it does on high loss, a bandwidth probe stops as it + /// does when loss shows inflight too high, and every other state lowers its short-term + /// model at the end of the round, as it does for loss. A mark loses no data, so it adds no + /// lost bytes or loss events, and the transport reports only an increased CE count, so old + /// marks never repeat the response. + /// + fn handle_ce(&mut self, now: Instant, sent: Instant, space: SpaceKind, packet_number: u64) { + self.enter_recovery(now, sent); + if std::mem::replace(&mut self.ce_in_recovery, true) { + return; + } + self.start_congestion_round(); + self.ce_in_round = true; + match self.state { + BbrState::Startup => { + // As Linux BBRv3's bbr_handle_queue_too_high_in_startup. + self.inflight_longterm = Ord::max(self.get_inflight(1.0), self.inflight_latest); + self.full_bw_reached = true; + self.full_bw_now = true; + self.enter_drain(); + } + BbrState::ProbeBw(ProbeBwSubstate::Refill | ProbeBwSubstate::Up) => { + // As handle_inflight_too_high, bounded by the inflight that drew the marks. + let packets = &self.packets[space as usize]; + if let Ok(i) = packets.binary_search_by_key(&packet_number, |p| p.packet_number) + && !packets[i].is_app_limited + { + self.inflight_longterm = Ord::max( + packets[i].tx_in_flight, + (self.target_inflight() as f64 * BETA) as u64, + ); + } + self.start_probe_bw_down(now); + } + _ => {} + } + self.update_control_parameters(0); + } + /// equivalent to BBRSaveStateUponLoss /// /// Runs once per recovery episode, so later losses in it cannot overwrite the pre-episode @@ -1654,6 +1715,7 @@ impl Bbr3 { /// equivalent to BBRResetCongestionSignals fn reset_congestion_signals(&mut self) { self.loss_in_round = false; + self.ce_in_round = false; self.bw_latest = 0.0; self.inflight_latest = 0; } @@ -1817,21 +1879,16 @@ impl Bbr3 { fn on_congestion_event( &mut self, now: Instant, - _sent: Instant, + sent: Instant, is_persistent_congestion: bool, is_ecn: bool, - lost_bytes: u64, - largest_lost_pn: u64, + _lost_bytes: u64, + largest_pn: u64, space: SpaceKind, ) { - // only process ecn here, regular packet loss is detected per packet in on_packet_lost. + // Loss is handled per packet in on_packet_lost, so only CE is handled here. if is_ecn { - self.lost += lost_bytes; - let p_index_result = self.packets[space as usize] - .binary_search_by_key(&largest_lost_pn, |p| p.packet_number); - if let Ok(p_index) = p_index_result { - self.process_lost_packet(p_index, space, now); - } + self.handle_ce(now, sent, space, largest_pn); } if is_persistent_congestion { self.cwnd = self.min_pipe_cwnd; @@ -1859,7 +1916,12 @@ impl Bbr3 { } /// equivalent to BBRHandleSpuriousLossDetection: + /// + /// CE marks in the episode confirm its congestion, so nothing is undone. fn on_spurious_congestion_event(&mut self) { + if self.ce_in_recovery { + return; + } self.restore_cwnd(); // No ACK processing may follow to re-bound the window, as when the ACK covers only // packets declared lost. ProbeRTT's exit restores the rest. @@ -2127,12 +2189,15 @@ mod test { pn: u64, send_ns: u64, ack_ns: u64, + /// whether the bottleneck marked the packet CE + ce: bool, } /// Single-bottleneck FIFO link simulator driving the real BBR /// `on_packet_sent`/`on_ack`/`on_end_acks` path against a constant bandwidth /// `bw`, constant propagation `rtt_ns`, and an infinite buffer (no loss). - /// Packets queue at the bottleneck and are served at `bw`. The sender always + /// Packets queue at the bottleneck and are served at `bw`, and are marked CE + /// past `ce_above_ns` of queueing delay if set. The sender always /// has data, paced at BBR's chosen rate, so it is cwnd-limited (never /// application-limited). Shared harness for the constant-link tests /// (A.1/A.3/A.5/A.8/A.9/A.10); tests needing loss, app-limiting, a mid-flight @@ -2153,6 +2218,11 @@ mod test { ret_ns: u64, // bottleneck serialization time for one MSS-sized packet btl_service_ns: u64, + /// queueing delay above which the bottleneck marks packets CE, as a classic AQM does; + /// `None` never marks + ce_above_ns: Option, + /// the longest queueing delay a packet has met at the bottleneck + max_queue_ns: u64, } impl Sim { @@ -2171,6 +2241,8 @@ mod test { fwd_ns: rtt_ns / 2, ret_ns: rtt_ns / 2, btl_service_ns: (mss as f64 / bw * 1e9).round() as u64, + ce_above_ns: None, + max_queue_ns: 0, } } @@ -2191,6 +2263,7 @@ mod test { pn: self.pn, send_ns: now_ns, ack_ns: now_ns, + ce: false, }); self.pn += 1; } @@ -2220,6 +2293,29 @@ mod test { .on_end_acks(now, self.inflight, app_limited, Some(0), SpaceKind::Data); } + /// Deliver one ACK frame as [`Self::ack`] does, whose ECN counts report new CE marks. The + /// transport then reports a congestion event for the largest acknowledged packet. + fn ack_ce(&mut self, now_ns: u64, pns: impl IntoIterator) { + let pns: Vec = pns.into_iter().collect(); + let largest = *pns.iter().max().unwrap(); + let sent_ns = self + .flight + .iter() + .find(|p| p.pn == largest) + .unwrap() + .send_ns; + self.ack(now_ns, pns, false); + self.bbr.on_congestion_event( + self.at(now_ns), + self.at(sent_ns), + false, + true, + 0, + largest, + SpaceKind::Data, + ); + } + /// Send `count` packets at `now_ns` and acknowledge them in one ACK `rtt_ns` later: one /// round when nothing else is in flight. fn round(&mut self, now_ns: u64, count: u64, rtt_ns: u64) { @@ -2228,6 +2324,13 @@ mod test { self.ack(now_ns + rtt_ns, first..self.pn, false); } + /// [`Self::round`], with the ACK reporting CE marks. + fn round_ce(&mut self, now_ns: u64, count: u64, rtt_ns: u64) { + let first = self.pn; + self.send(now_ns, count); + self.ack_ce(now_ns + rtt_ns, first..self.pn); + } + /// Report an empty transmit poll that nothing held back, as the transport does. fn starve(&mut self) { self.bbr.on_app_limited(self.inflight); @@ -2274,6 +2377,9 @@ mod test { let finish = service_start + self.btl_service_ns; self.btl_free_ns = finish; let ack_ns = finish + self.ret_ns; + let queue_ns = service_start - arrival; + self.max_queue_ns = self.max_queue_ns.max(queue_ns); + let ce = self.ce_above_ns.is_some_and(|limit| queue_ns > limit); self.bbr.on_packet_sent( self.base + Duration::from_nanos(send_ns), @@ -2286,6 +2392,7 @@ mod test { pn: self.pn, send_ns, ack_ns, + ce, }); // pace the next send at BBR's chosen pacing rate @@ -2316,6 +2423,18 @@ mod test { ); self.bbr .on_end_acks(now_at, self.inflight, false, Some(p.pn), SpaceKind::Data); + // Each ACK covers one packet, so a marked one raises the CE count. + if p.ce { + self.bbr.on_congestion_event( + now_at, + send_at, + false, + true, + 0, + p.pn, + SpaceKind::Data, + ); + } if on_ack(&mut self.bbr, self.now_ns, self.inflight, p.pn).is_break() { return; @@ -7815,10 +7934,14 @@ mod test { assert_eq!(bbr.last_lost_packet, Some((SpaceKind::Handshake, 0))); } - /// An ECN congestion event names its largest packet by space, like a loss. + /// An ECN congestion event names its largest packet by space: a probe it stops is bounded by + /// the inflight Handshake packet 0 was sent with, not Initial packet 0's. A mark loses + /// nothing, so every packet stays tracked. #[test] - fn ecn_congestion_marks_the_packet_from_its_own_space() { + fn ecn_congestion_reads_the_packet_from_its_own_space() { let mut bbr = Bbr3::new(Arc::new(Bbr3Config::default()), PACKET); + bbr.full_bw_reached = true; + bbr.start_probe_bw_up(); let t0 = Instant::now(); let at = |ms| t0 + Duration::from_millis(ms); let c: &mut dyn Controller = &mut bbr; @@ -7828,10 +7951,222 @@ mod test { c.on_packet_space_sent(at(2), PACKET, HANDSHAKE_0); c.on_congestion_event_space(at(10), at(2), false, true, 0, HANDSHAKE_0); - assert_eq!( - tracked_send_ms(&bbr, t0), - [0, 1], - "only Handshake 0, sent at 2ms, is gone" + assert_eq!(bbr.state, BbrState::ProbeBw(ProbeBwSubstate::Down)); + assert_eq!(bbr.inflight_longterm, 3 * PACKET as u64); + assert_eq!(tracked_send_ms(&bbr, t0), [0, 1, 2]); + } + + /// Rounds of 1, 2, 4, ..., 128 packets, 20ms apart, each acknowledged 10ms after it is sent, + /// with every ACK reporting CE marks or none. + fn doubling(ce: bool) -> Sim { + let mut sim = scripted(); + for i in 0..8 { + let (now, count) = (i * 20 * MS, 1 << i); + match ce { + true => sim.round_ce(now, count, 10 * MS), + false => sim.round(now, count, 10 * MS), + } + } + sim + } + + /// CE marks stop Startup and cut the sending load while the unmarked control keeps + /// doubling. Before, both stayed in Startup with the same window and pacing rate. + #[test] + fn startup_stops_on_ce() { + let control = doubling(false); + assert_eq!(control.bbr.state, BbrState::Startup); + assert_eq!(control.bbr.window(), 318_000); + assert!(control.bbr.recovery_start_time.is_none()); + + // The scripted sender ignores the window, so the model still sees the doubled flight. + let marked = doubling(true); + assert!(marked.bbr.full_bw_reached); + assert_eq!(marked.bbr.state, BbrState::ProbeBw(ProbeBwSubstate::Cruise)); + assert!(marked.bbr.window() * 2 < control.bbr.window()); + assert!(marked.bbr.pacing_rate * 2.0 < control.bbr.pacing_rate); + // A mark is not a loss. + assert_eq!(marked.bbr.lost, 0); + assert_eq!(marked.bbr.loss_events_in_round, 0); + } + + /// The first CE-bearing ACK in Startup drains at once, without waiting for the next ACK. + #[test] + fn startup_drains_on_first_ce() { + let mut sim = scripted(); + sim.round(0, 10, 10 * MS); + let pacing = sim.bbr.pacing_rate; + sim.round_ce(10 * MS, 20, 10 * MS); + assert_eq!(sim.bbr.state, BbrState::Drain); + assert!(sim.bbr.pacing_rate < pacing); + let bdp = sim.bbr.get_inflight(1.0); + assert_eq!(sim.bbr.inflight_longterm, bdp); + } + + /// CE ends a bandwidth probe as too much loss does, bounding inflight by the flight that drew + /// the marks; the unmarked control keeps probing. + #[test] + fn probe_up_stops_on_ce() { + let mut control = probing_up(); + control.ack(20 * MS, 10..20, false); + assert_eq!(control.bbr.state, BbrState::ProbeBw(ProbeBwSubstate::Up)); + + let mut marked = probing_up(); + marked.ack_ce(20 * MS, 10..20); + assert_eq!(marked.bbr.state, BbrState::ProbeBw(ProbeBwSubstate::Down)); + // Packet 19 was sent with ten packets in flight. + assert_eq!(marked.bbr.inflight_longterm, 12_000); + assert!(marked.bbr.window() <= 12_000); + assert!(marked.bbr.window() < control.bbr.window()); + assert!(marked.bbr.pacing_rate < control.bbr.pacing_rate); + } + + /// A scripted `Sim` cruising after a loss-free probe of 12 MB/s over a 10ms RTT. + fn cruising() -> Sim { + let mut sim = probed(); + for i in 0..5 { + sim.round(10 * MS + i * 20 * MS, 100, 20 * MS); + } + assert_eq!(sim.bbr.state, BbrState::ProbeBw(ProbeBwSubstate::Cruise)); + sim + } + + /// CE while cruising lowers the short-term model when the round ends, as loss does. + #[test] + fn cruise_lowers_short_term_model_on_ce() { + let mut control = cruising(); + control.round(110 * MS, 100, 20 * MS); + control.round(130 * MS, 100, 20 * MS); + assert_eq!(control.bbr.inflight_shortterm, u64::MAX); + + let mut marked = cruising(); + marked.round_ce(110 * MS, 100, 20 * MS); + marked.round(130 * MS, 100, 20 * MS); + assert_eq!(marked.bbr.state, BbrState::ProbeBw(ProbeBwSubstate::Cruise)); + assert!(marked.bbr.inflight_shortterm < u64::MAX); + assert!(marked.bbr.window() < control.bbr.window()); + } + + /// Several ACKs reporting CE in one recovery episode respond as one does, even once the + /// first has ended the probe; a mark on a packet sent after the episode began responds + /// again. + #[test] + fn ce_responds_once_per_recovery_episode() { + let mut once = probing_up(); + let mut several = probing_up(); + once.ack_ce(20 * MS, 10..15); + several.ack_ce(20 * MS, 10..15); + once.ack(21 * MS, 15..20, false); + several.ack_ce(21 * MS, 15..20); + for sim in [&mut once, &mut several] { + sim.round(21 * MS, 10, 10 * MS); + sim.round(31 * MS, 10, 10 * MS); + } + assert_eq!(several.bbr.inflight_shortterm, once.bbr.inflight_shortterm); + assert_eq!(several.bbr.window(), once.bbr.window()); + + // The next marked round is a new episode. + let cut = several.bbr.inflight_shortterm; + several.round_ce(41 * MS, 10, 10 * MS); + several.round(51 * MS, 10, 10 * MS); + assert!(several.bbr.inflight_shortterm < cut); + } + + /// CE during ProbeRTT neither raises its window nor delays its exit. + #[test] + fn probe_rtt_holds_on_ce() { + let run = |ce: bool| { + let mut sim = probing_rtt(); + let window = sim.bbr.window(); + match ce { + true => sim.round_ce(20 * MS, 4, 10 * MS), + false => sim.round(20 * MS, 4, 10 * MS), + } + assert!(sim.bbr.window() <= window); + let mut now = 30 * MS; + while sim.bbr.state == BbrState::ProbeRtt { + sim.round(now, 4, 10 * MS); + now += 10 * MS; + assert!(now < 1000 * MS, "ProbeRTT never ended"); + } + assert_eq!(sim.bbr.state, BbrState::ProbeBw(ProbeBwSubstate::Cruise)); + now + }; + assert_eq!(run(true), run(false)); + } + + /// CE on the ACK that ends a loss's round still stops Startup, though the loss began the + /// recovery episode, and declaring that loss spurious undoes neither. The mark adds no lost + /// bytes. + #[test] + fn ce_with_loss_survives_spurious_undo() { + let mut sim = scripted(); + sim.round(0, 10, 10 * MS); + let first = sim.pn; + sim.send(10 * MS, 20); + sim.lose(15 * MS, first); + sim.ack_ce(20 * MS, first + 1..sim.pn); + assert_ne!(sim.bbr.state, BbrState::Startup); + assert_eq!(sim.bbr.lost, PACKET as u64); + let (window, bound) = (sim.bbr.window(), sim.bbr.inflight_longterm); + + sim.bbr.on_spurious_congestion_event(); + assert_ne!(sim.bbr.state, BbrState::Startup); + assert!(sim.bbr.full_bw_reached); + assert_eq!(sim.bbr.window(), window); + assert_eq!(sim.bbr.inflight_longterm, bound); + } + + /// A 10 Mbit/s link with a 20ms RTT, a 25,000-byte BDP. + const LINK_BW: f64 = 1_250_000.0; + const LINK_RTT: u64 = 20 * MS; + + /// Runs `sim` until `until_ns`, returning the bytes delivered meanwhile. + fn run_until(sim: &mut Sim, until_ns: u64) -> u64 { + let start = sim.bbr.delivered; + sim.run( + 10_000_000, + |_| ControlFlow::Continue(()), + |_, now_ns, _, _| match now_ns >= until_ns { + true => ControlFlow::Break(()), + false => ControlFlow::Continue(()), + }, ); + sim.bbr.delivered - start + } + + /// A bottleneck that marks CE above 5ms of queue, a quarter of the RTT, holds a shorter + /// queue than the unmarked control at nearly the same goodput. Once marking stops, the flow + /// regains the control's rate. + #[test] + fn sustained_ce_bounds_the_queue_and_recovers() { + const SEC: u64 = 1_000_000_000; + let config = || Bbr3Config { + probe_rng_seed: Some([7; 16]), + ..Bbr3Config::default() + }; + let mut control = Sim::new(config(), 1200, LINK_BW, LINK_RTT); + let mut marked = Sim::new(config(), 1200, LINK_BW, LINK_RTT); + marked.ce_above_ns = Some(5 * MS); + + // Startup and its drain. + run_until(&mut control, SEC); + run_until(&mut marked, SEC); + control.max_queue_ns = 0; + marked.max_queue_ns = 0; + + let control_bytes = run_until(&mut control, 11 * SEC); + let marked_bytes = run_until(&mut marked, 11 * SEC); + // Measured at 19ms for the control and 11ms marked. + assert!(control.bbr.recovery_start_time.is_none()); + assert!(marked.bbr.recovery_start_time.is_some()); + assert!(marked.max_queue_ns * 4 < control.max_queue_ns * 3); + assert!(marked_bytes as f64 > 0.9 * control_bytes as f64); + + marked.ce_above_ns = None; + run_until(&mut control, 21 * SEC); + run_until(&mut marked, 21 * SEC); + assert!(marked.bbr.max_bw > 0.9 * control.bbr.max_bw); + assert_eq!(marked.bbr.inflight_shortterm, u64::MAX); } } diff --git a/noq-proto/src/tests/mod.rs b/noq-proto/src/tests/mod.rs index 13c47f071..10e673e37 100644 --- a/noq-proto/src/tests/mod.rs +++ b/noq-proto/src/tests/mod.rs @@ -38,7 +38,7 @@ use crate::{ StreamEvent, Transmit, TransportConfig, TransportErrorCode, VarInt, WriteError, cid_generator::{ConnectionIdGenerator, RandomConnectionIdGenerator}, coding::{Decodable, Encodable}, - congestion::{Controller, ControllerFactory, ControllerMetrics, PacketId, Space}, + congestion::{Bbr3Config, Controller, ControllerFactory, ControllerMetrics, PacketId, Space}, crypto::rustls::{QuicServerConfig, configured_provider}, frame::{self, Frame, FrameStruct}, packet::{FixedLengthConnectionIdParser, PartialDecode}, @@ -5065,6 +5065,64 @@ fn recorded_pair(factory: PacketRecorderFactory) -> (Pair, ConnectionHandle, Pac (pair, client_ch, log) } +/// Counts the congestion events one controller's log reports for CE marks. +fn ce_events(log: &PacketLog) -> usize { + log.lock() + .unwrap() + .iter() + .filter(|e| matches!(e, PacketEvent::Congestion { ecn: true, .. })) + .count() +} + +/// ACKs keep carrying the CE count after the marks stop, but only an increase is a congestion +/// event, so a controller never answers the same marks twice. +#[test] +fn old_ce_marks_report_no_new_congestion() { + let _guard = subscribe(); + let (mut pair, client_ch, log) = recorded_pair(Default::default()); + let s = pair.client_streams(client_ch).open(Dir::Uni).unwrap(); + pair.client_send(client_ch, s) + .write(&[42; 8 * 1024]) + .unwrap(); + pair.congestion_experienced = true; + pair.drive_client(); + pair.congestion_experienced = false; + pair.drive(); + let marked = ce_events(&log); + assert!(marked > 0); + + pair.client_send(client_ch, s) + .write(&[42; 64 * 1024]) + .unwrap(); + pair.drive(); + assert_eq!(ce_events(&log), marked); + assert!(pair.client_conn_mut(client_ch).using_ecn()); +} + +/// A path that starts bleaching the ECN field or re-marking it ECT(1) reports no CE, though the +/// bottleneck marks every packet, and re-marking fails validation, so the sender stops using +/// ECN. Bleaching is only caught once an ACK's count increase falls short of its ACK ranges, +/// which may not happen within this transfer. +#[test] +fn invalid_ecn_feedback_reports_no_congestion() { + let _guard = subscribe(); + for rewrite in [None, Some(EcnCodepoint::Ect1)] { + let (mut pair, client_ch, log) = recorded_pair(Default::default()); + assert!(pair.client_conn_mut(client_ch).using_ecn()); + pair.congestion_experienced = true; + pair.rewrite_ecn = Some(rewrite); + let s = pair.client_streams(client_ch).open(Dir::Uni).unwrap(); + pair.client_send(client_ch, s) + .write(&[42; 16 * 1024]) + .unwrap(); + pair.drive(); + assert_eq!(ce_events(&log), 0, "{rewrite:?}"); + if rewrite.is_some() { + assert!(!pair.client_conn_mut(client_ch).using_ecn()); + } + } +} + /// Sends with `send`, lets it deliver, then sends again after an idle gap. With nothing in /// flight no ACK arrives during the gap, so the empty poll that follows the last ACK must tell /// the controller of starvation before the resumed send. @@ -5375,12 +5433,87 @@ fn throughput() -> TestResult { BwLimitConfig { bytes_per_second: BPS_LIMIT, buffer_size: 50 * 1500, // buffer that fits ~50 full packets + marks_ce: true, latency: Duration::from_millis(3), }, )) .connect(); - let mut bytes_to_send = TOTAL_BYTES; + let time = upload(&mut pair, TOTAL_BYTES)?; + let bytes_per_second = TOTAL_BYTES as f64 / time.as_secs_f64(); + info!(?time, bytes_per_second); + + let expected_bps = BPS_LIMIT as f64; + // Less than 2% deviation from the BPS limit + assert!( + (bytes_per_second - expected_bps).abs() / expected_bps < 0.05, + "deviated too far from expected throughput limit" + ); + + Ok(()) +} + +/// A 1 MB/s bottleneck with a 20ms RTT and a one-BDP buffer carries a BBR upload, once marking +/// CE at half full and once only tail-dropping when full. BBR answers the marks, so the marking +/// run holds a shorter queue without drops at the dropping run's goodput. +#[test] +fn bbr_marking_versus_dropping() -> TestResult { + const TOTAL_BYTES: usize = 4_000_000; + const BPS_LIMIT: u64 = 1_000_000; + + let _guard = subscribe(); + let run = |marks_ce| -> TestResult<_> { + let mut transport = TransportConfig::default(); + transport.congestion_controller_factory(Arc::new(Bbr3Config::default())); + let mut pair = ConnPair::builder() + .with_transport_cfg(transport) + .with_routes(BwLimitedRouting::new( + Pair::CLIENT_ADDR, + Pair::SERVER_ADDR, + Instant::now(), + BwLimitConfig { + bytes_per_second: BPS_LIMIT, + buffer_size: 20_000, + marks_ce, + latency: Duration::from_millis(10), + }, + )) + .connect(); + let time = upload(&mut pair, TOTAL_BYTES)?; + let goodput = TOTAL_BYTES as f64 / time.as_secs_f64(); + let ecn = pair.conn(Client).using_ecn(); + let queue = pair.routes.as_bw_limited().client_to_server(); + info!( + marks_ce, + ecn, + goodput, + mean_delay = ?queue.mean_delay(), + max_delay = ?queue.max_delay, + dropped = queue.dropped, + marked = queue.marked, + congested = ?queue.congested, + ); + assert!(ecn); + Ok((goodput, queue.clone())) + }; + + let (marking_goodput, marking) = run(true)?; + let (dropping_goodput, dropping) = run(false)?; + // Measured at 3.2ms mean delay for marking and 7.4ms for dropping, 0.5s and 1.6s past the + // marking threshold, and 38 drops when dropping. The controller before classic ECN dropped + // 44 packets when marking, at 7.4ms and 1.6s. + assert!(marking.marked > 0); + assert_eq!(marking.dropped, 0); + assert!(dropping.dropped > 0); + assert!(marking.mean_delay() * 3 < dropping.mean_delay() * 2); + assert!(marking.congested * 2 < dropping.congested); + assert!(marking_goodput > 0.9 * dropping_goodput); + Ok(()) +} + +/// Uploads `total` bytes from client to server, returning how long it took. +fn upload(pair: &mut ConnPair, total: usize) -> TestResult { + let mut bytes_to_send = total; let mut bytes_received = 0; let start = pair.time; @@ -5407,20 +5540,8 @@ fn throughput() -> TestResult { } assert_eq!(bytes_to_send, 0); - assert_eq!(bytes_received, TOTAL_BYTES); - - let time = pair.time.saturating_duration_since(start); - let bytes_per_second = TOTAL_BYTES as f64 / time.as_secs_f64(); - info!(bytes_received, ?time, bytes_per_second); - - let expected_bps = BPS_LIMIT as f64; - // Less than 2% deviation from the BPS limit - assert!( - (bytes_per_second - expected_bps).abs() / expected_bps < 0.05, - "deviated too far from expected throughput limit" - ); - - Ok(()) + assert_eq!(bytes_received, total); + Ok(pair.time.saturating_duration_since(start)) } const ZEROES: [u8; 10_000] = [0u8; 10_000]; diff --git a/noq-proto/src/tests/util.rs b/noq-proto/src/tests/util.rs index 90af6ea8e..cec65ef13 100644 --- a/noq-proto/src/tests/util.rs +++ b/noq-proto/src/tests/util.rs @@ -49,6 +49,9 @@ pub(super) struct Pair { pub(super) mtu: usize, /// Simulates explicit congestion notification pub(super) congestion_experienced: bool, + /// Replaces the ECN codepoint of every delivered datagram after any marking, as a path that + /// bleaches the field (`Some(None)`) or re-marks it + pub(super) rewrite_ecn: Option>, /// Number of spin bit flips pub(super) spins: u64, /// The routing table used for resolving addresses observed for incoming packets @@ -106,6 +109,7 @@ impl Pair { spins: 0, last_spin: false, congestion_experienced: false, + rewrite_ecn: None, routes: BasicRouting { client_addr, server_addr, @@ -141,6 +145,7 @@ impl Pair { spins: 0, last_spin: false, congestion_experienced: false, + rewrite_ecn: None, routes: BasicRouting { client_addr: Self::CLIENT_ADDR, server_addr: Self::SERVER_ADDR, @@ -219,10 +224,12 @@ impl Pair { recv_time, congestion_experienced, } => { - let ecn = set_congestion_experienced( - packet.ecn, - self.congestion_experienced || congestion_experienced, - ); + let ecn = self.rewrite_ecn.unwrap_or_else(|| { + set_congestion_experienced( + packet.ecn, + self.congestion_experienced || congestion_experienced, + ) + }); self.server.inbound.push( recv_time, Inbound { @@ -257,10 +264,12 @@ impl Pair { recv_time, congestion_experienced, } => { - let ecn = set_congestion_experienced( - packet.ecn, - self.congestion_experienced || congestion_experienced, - ); + let ecn = self.rewrite_ecn.unwrap_or_else(|| { + set_congestion_experienced( + packet.ecn, + self.congestion_experienced || congestion_experienced, + ) + }); self.client.inbound.push( recv_time, Inbound { @@ -1690,7 +1699,7 @@ pub(super) enum Routing { Basic(BasicRouting), SimpleFirewall(SimpleFirewallRouting), ManyToMany(ManyToManyRouting), - BwLimited(BwLimitedRouting), + BwLimited(Box), } impl Routing { @@ -1739,6 +1748,13 @@ impl Routing { } } + pub(super) fn as_bw_limited(&self) -> &BwLimitedRouting { + match self { + Self::BwLimited(inner) => inner, + _ => panic!("cast to BwLimitedRouting failed, a different routing table is set"), + } + } + /// Sets the one-way latency between the two endpoints. pub(super) fn set_latency(&mut self, latency: Duration) { match self { @@ -2334,13 +2350,15 @@ pub(super) struct BwLimitConfig { /// /// Once the queue exceeds this, packets are tail-dropped. pub(super) buffer_size: u32, + /// Whether packets are marked CE once the queue is half full, as a classic AQM does. + pub(super) marks_ce: bool, /// The one-way latency of the simulated link, applied to both directions. pub(super) latency: Duration, } impl From for Routing { fn from(value: BwLimitedRouting) -> Self { - Self::BwLimited(value) + Self::BwLimited(Box::new(value)) } } @@ -2354,17 +2372,23 @@ impl BwLimitedRouting { let BwLimitConfig { bytes_per_second, buffer_size, + marks_ce, latency, } = config; Self { client_addr, server_addr, - limiter_client_to_server: Limiter::new(bytes_per_second, buffer_size, now), - limiter_server_to_client: Limiter::new(bytes_per_second, buffer_size, now), + limiter_client_to_server: Limiter::new(bytes_per_second, buffer_size, marks_ce, now), + limiter_server_to_client: Limiter::new(bytes_per_second, buffer_size, marks_ce, now), latency, } } + /// What the client-to-server queue has seen. + pub(super) fn client_to_server(&self) -> &QueueStats { + &self.limiter_client_to_server.stats + } + pub(super) fn set_latency(&mut self, latency: Duration) { self.latency = latency; } @@ -2416,15 +2440,42 @@ impl BwLimitedRouting { } } +/// What a [`BwLimitedRouting`] queue has seen. +#[derive(Debug, Default, Clone)] +pub(super) struct QueueStats { + /// Packets tail-dropped + pub(super) dropped: u64, + /// Packets marked CE + pub(super) marked: u64, + /// Packets queued for delivery + pub(super) queued: u32, + /// The queueing delay delivered packets met, summed + pub(super) total_delay: Duration, + /// The longest queueing delay a delivered packet met + pub(super) max_delay: Duration, + /// How long the link spent serving packets that queued past half full, the marking + /// threshold, so how long the sender took to answer congestion, summed over the run + pub(super) congested: Duration, +} + +impl QueueStats { + /// The mean queueing delay delivered packets met. + pub(super) fn mean_delay(&self) -> Duration { + self.total_delay / self.queued.max(1) + } +} + #[derive(Debug)] struct Limiter { time_per_byte: Duration, finished_sending_at: Instant, max_queue: Duration, + marks_ce: bool, + stats: QueueStats, } impl Limiter { - fn new(bytes_per_second: u64, buffer_size: u32, now: Instant) -> Self { + fn new(bytes_per_second: u64, buffer_size: u32, marks_ce: bool, now: Instant) -> Self { const ONE_SEC_IN_NANOS: u64 = 1_000_000_000; let time_per_byte = Duration::from_nanos(ONE_SEC_IN_NANOS / bytes_per_second); let max_queue = buffer_size * time_per_byte; @@ -2432,19 +2483,31 @@ impl Limiter { time_per_byte, max_queue, finished_sending_at: now, + marks_ce, + stats: QueueStats::default(), } } fn time_of_send(&mut self, now: Instant, num_bytes: usize) -> Option<(Instant, bool)> { - if self.finished_sending_at > now + self.max_queue { + let delay = self.finished_sending_at.saturating_duration_since(now); + if delay > self.max_queue { + self.stats.dropped += 1; return None; } + let service = num_bytes as u32 * self.time_per_byte; + self.stats.queued += 1; + self.stats.total_delay += delay; + self.stats.max_delay = self.stats.max_delay.max(delay); + if delay > self.max_queue / 2 { + self.stats.congested += service; + } - self.finished_sending_at = - cmp::max(now, self.finished_sending_at) + num_bytes as u32 * self.time_per_byte; + self.finished_sending_at = cmp::max(now, self.finished_sending_at) + service; // We mark the CE bit once the queue is 50% full - let experienced_congestion = self.finished_sending_at > now + self.max_queue / 2; + let experienced_congestion = + self.marks_ce && self.finished_sending_at > now + self.max_queue / 2; + self.stats.marked += experienced_congestion as u64; Some((self.finished_sending_at, experienced_congestion)) } From 039130b8ecc807c1cb73e316494938f107c1749d Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sat, 26 Sep 2026 11:06:31 -0700 Subject: [PATCH 2/4] fix(proto): bound a CE-stopped probe by the ACK's rate sample The CE event names the ACK's largest packet, which can be an ACK-only packet BBR never tracked, so the lookup found nothing and left the probe unbounded. Co-Authored-By: Claude Opus 5.5 --- noq-proto/src/congestion/bbr3/mod.rs | 44 +++++++++++++++++----------- 1 file changed, 27 insertions(+), 17 deletions(-) diff --git a/noq-proto/src/congestion/bbr3/mod.rs b/noq-proto/src/congestion/bbr3/mod.rs index 57bd3ca27..dcb008fb8 100644 --- a/noq-proto/src/congestion/bbr3/mod.rs +++ b/noq-proto/src/congestion/bbr3/mod.rs @@ -1600,7 +1600,7 @@ impl Bbr3 { /// lost bytes or loss events, and the transport reports only an increased CE count, so old /// marks never repeat the response. /// - fn handle_ce(&mut self, now: Instant, sent: Instant, space: SpaceKind, packet_number: u64) { + fn handle_ce(&mut self, now: Instant, sent: Instant) { self.enter_recovery(now, sent); if std::mem::replace(&mut self.ce_in_recovery, true) { return; @@ -1616,13 +1616,14 @@ impl Bbr3 { self.enter_drain(); } BbrState::ProbeBw(ProbeBwSubstate::Refill | ProbeBwSubstate::Up) => { - // As handle_inflight_too_high, bounded by the inflight that drew the marks. - let packets = &self.packets[space as usize]; - if let Ok(i) = packets.binary_search_by_key(&packet_number, |p| p.packet_number) - && !packets[i].is_app_limited + // As handle_inflight_too_high, bounded by the inflight that drew the marks: the + // marked ACK's rate sample. The event's own packet may be an untracked ACK-only + // packet. + if let Some(rs) = self.rs + && !rs.is_app_limited { self.inflight_longterm = Ord::max( - packets[i].tx_in_flight, + rs.tx_in_flight, (self.target_inflight() as f64 * BETA) as u64, ); } @@ -1883,12 +1884,12 @@ impl Bbr3 { is_persistent_congestion: bool, is_ecn: bool, _lost_bytes: u64, - largest_pn: u64, - space: SpaceKind, + _largest_lost_pn: u64, + _space: SpaceKind, ) { // Loss is handled per packet in on_packet_lost, so only CE is handled here. if is_ecn { - self.handle_ce(now, sent, space, largest_pn); + self.handle_ce(now, sent); } if is_persistent_congestion { self.cwnd = self.min_pipe_cwnd; @@ -7934,14 +7935,11 @@ mod test { assert_eq!(bbr.last_lost_packet, Some((SpaceKind::Handshake, 0))); } - /// An ECN congestion event names its largest packet by space: a probe it stops is bounded by - /// the inflight Handshake packet 0 was sent with, not Initial packet 0's. A mark loses - /// nothing, so every packet stays tracked. + /// A CE mark loses nothing, so the packet an ECN congestion event names stays tracked, + /// where it was once removed as lost. #[test] - fn ecn_congestion_reads_the_packet_from_its_own_space() { + fn ce_keeps_the_marked_packet_tracked() { let mut bbr = Bbr3::new(Arc::new(Bbr3Config::default()), PACKET); - bbr.full_bw_reached = true; - bbr.start_probe_bw_up(); let t0 = Instant::now(); let at = |ms| t0 + Duration::from_millis(ms); let c: &mut dyn Controller = &mut bbr; @@ -7951,9 +7949,8 @@ mod test { c.on_packet_space_sent(at(2), PACKET, HANDSHAKE_0); c.on_congestion_event_space(at(10), at(2), false, true, 0, HANDSHAKE_0); - assert_eq!(bbr.state, BbrState::ProbeBw(ProbeBwSubstate::Down)); - assert_eq!(bbr.inflight_longterm, 3 * PACKET as u64); assert_eq!(tracked_send_ms(&bbr, t0), [0, 1, 2]); + assert_eq!(bbr.lost, 0); } /// Rounds of 1, 2, 4, ..., 128 packets, 20ms apart, each acknowledged 10ms after it is sent, @@ -8021,6 +8018,19 @@ mod test { assert!(marked.bbr.pacing_rate < control.bbr.pacing_rate); } + /// The transport names the ACK's largest packet, which can be an ACK-only packet BBR never + /// tracked; the probe is still bounded, by the ACK's rate sample. + #[test] + fn probe_up_ce_on_an_untracked_packet_bounds_inflight() { + let mut sim = probing_up(); + sim.ack(20 * MS, 10..20, false); + let now = sim.at(20 * MS); + sim.bbr + .on_congestion_event(now, now, false, true, 0, 1000, SpaceKind::Data); + assert_eq!(sim.bbr.state, BbrState::ProbeBw(ProbeBwSubstate::Down)); + assert_eq!(sim.bbr.inflight_longterm, 12_000); + } + /// A scripted `Sim` cruising after a loss-free probe of 12 MB/s over a 10ms RTT. fn cruising() -> Sim { let mut sim = probed(); From 7664b06134b48b5cb5012a5af1fa89adfe1812b4 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sat, 26 Sep 2026 11:28:11 -0700 Subject: [PATCH 3/4] fix(proto): bound a CE-stopped probe only by the marked ACK's own sample CE follows loss's probe-feedback rule (bw_probe_samples), so marks on the ACK that itself ended the probe still bound it, and an ACK of only untracked packets no longer bounds it by an older ACK's sample. Co-Authored-By: Claude Opus 5.5 --- noq-proto/src/congestion/bbr3/mod.rs | 67 +++++++++++++++++++--------- 1 file changed, 47 insertions(+), 20 deletions(-) diff --git a/noq-proto/src/congestion/bbr3/mod.rs b/noq-proto/src/congestion/bbr3/mod.rs index dcb008fb8..94f773494 100644 --- a/noq-proto/src/congestion/bbr3/mod.rs +++ b/noq-proto/src/congestion/bbr3/mod.rs @@ -515,6 +515,9 @@ pub struct Bbr3 { /// equivalent to RS.has_data: true once the ACK being processed has delivered a tracked /// packet, so `rs` describes this ACK and is folded into the model when the ACK ends. rs_has_data: bool, + /// Whether the last ACK delivered a tracked packet, so `rs` describes it. A CE event is + /// reported after its ACK ends, and an ACK of only untracked packets leaves an older `rs`. + rs_from_last_ack: bool, /// equivalent to RS.newly_acked, accumulated over the ACK being processed. It is passed to /// the ACK's model steps rather than kept in `rs`, so no later `set_cwnd` counts it again. newly_acked: u64, @@ -696,6 +699,7 @@ impl Bbr3 { lost: 0, rs: None, rs_has_data: false, + rs_from_last_ack: false, newly_acked: 0, packets: Default::default(), rounds_since_bw_probe: 0, @@ -1556,7 +1560,7 @@ impl Bbr3 { if self.is_inflight_too_high() { rate_sample.tx_in_flight = self.inflight_at_loss(p.size as u64); self.rs = Some(rate_sample); - self.handle_inflight_too_high(now); + self.handle_inflight_too_high(now, self.rs); } } self.packets[space as usize].remove(packet_index); @@ -1594,9 +1598,10 @@ impl Bbr3 { /// /// Draft-06 section 3.7 requires treating CE as congestion without prescribing BBR's /// response. As RFC 9002 reduces its window once per recovery period, this responds once per - /// recovery episode: Startup stops as it does on high loss, a bandwidth probe stops as it - /// does when loss shows inflight too high, and every other state lowers its short-term - /// model at the end of the round, as it does for loss. A mark loses no data, so it adds no + /// recovery episode: Startup stops as it does on high loss, a bandwidth probe whose feedback + /// is arriving stops as it does when loss shows inflight too high, and every other state + /// lowers its short-term model at the end of the round, as it does for loss. A mark loses no + /// data, so it adds no /// lost bytes or loss events, and the transport reports only an increased CE count, so old /// marks never repeat the response. /// @@ -1615,19 +1620,12 @@ impl Bbr3 { self.full_bw_now = true; self.enter_drain(); } - BbrState::ProbeBw(ProbeBwSubstate::Refill | ProbeBwSubstate::Up) => { - // As handle_inflight_too_high, bounded by the inflight that drew the marks: the - // marked ACK's rate sample. The event's own packet may be an untracked ACK-only - // packet. - if let Some(rs) = self.rs - && !rs.is_app_limited - { - self.inflight_longterm = Ord::max( - rs.tx_in_flight, - (self.target_inflight() as f64 * BETA) as u64, - ); - } - self.start_probe_bw_down(now); + // The marked ACK's sample bounds the probe, since the event's own packet can be an + // untracked ACK-only one. The probe's feedback can still be arriving after the ACK + // itself ended the probe. + _ if self.bw_probe_samples => { + let rs = self.rs.filter(|_| self.rs_from_last_ack); + self.handle_inflight_too_high(now, rs); } _ => {} } @@ -1672,9 +1670,11 @@ impl Bbr3 { } /// equivalent to BBRHandleInflightTooHigh - fn handle_inflight_too_high(&mut self, now: Instant) { + /// + /// `rs` is the sample that showed it, if any. + fn handle_inflight_too_high(&mut self, now: Instant, rs: Option) { self.bw_probe_samples = false; - if let Some(rate_sample) = self.rs + if let Some(rate_sample) = rs && !rate_sample.is_app_limited { self.inflight_longterm = Ord::max( @@ -1855,7 +1855,8 @@ impl Bbr3 { let newly_acked = std::mem::take(&mut self.newly_acked); // An ACK that delivered no tracked packet has no sample, and a finished one is never // folded twice. - if !std::mem::take(&mut self.rs_has_data) { + self.rs_from_last_ack = std::mem::take(&mut self.rs_has_data); + if !self.rs_from_last_ack { return; } let Some(mut rs) = self.rs else { @@ -8031,6 +8032,32 @@ mod test { assert_eq!(sim.bbr.inflight_longterm, 12_000); } + /// CE on the ACK that itself ends the probe still bounds it, as the probe's feedback is + /// still arriving. + #[test] + fn ce_on_the_probe_ending_ack_bounds_inflight() { + let mut sim = probing_up(); + // This ACK finds the bandwidth plateau and leaves ProbeUp before its CE is reported. + sim.bbr.full_bw_now = true; + sim.ack_ce(20 * MS, 10..20); + assert_eq!(sim.bbr.state, BbrState::ProbeBw(ProbeBwSubstate::Down)); + assert_eq!(sim.bbr.inflight_longterm, 12_000); + } + + /// CE on an ACK of only untracked packets stops the probe without bounding it by an older + /// ACK's sample. + #[test] + fn ce_without_a_sample_keeps_the_bound() { + let mut sim = probing_up(); + let now = sim.at(20 * MS); + sim.bbr + .on_end_acks(now, sim.inflight, false, Some(1000), SpaceKind::Data); + sim.bbr + .on_congestion_event(now, now, false, true, 0, 1000, SpaceKind::Data); + assert_eq!(sim.bbr.state, BbrState::ProbeBw(ProbeBwSubstate::Down)); + assert_eq!(sim.bbr.inflight_longterm, 100_000); + } + /// A scripted `Sim` cruising after a loss-free probe of 12 MB/s over a 10ms RTT. fn cruising() -> Sim { let mut sim = probed(); From c85041b2d050a84616ca127807999373c3a377c0 Mon Sep 17 00:00:00 2001 From: Luke Curley Date: Sun, 27 Sep 2026 16:19:06 -0700 Subject: [PATCH 4/4] fix(proto): bound a CE-stopped probe by the ACK's own sample Loss detection runs between an ACK and its CE event and rewrites `rs` per lost packet, so CE could bound the probe by a lost packet's flight. Keep the completed ACK's sample apart for CE. Co-Authored-By: Claude Opus 5.5 --- noq-proto/src/congestion/bbr3/mod.rs | 36 +++++++++++++++++++++------- 1 file changed, 28 insertions(+), 8 deletions(-) diff --git a/noq-proto/src/congestion/bbr3/mod.rs b/noq-proto/src/congestion/bbr3/mod.rs index 94f773494..d7531b15c 100644 --- a/noq-proto/src/congestion/bbr3/mod.rs +++ b/noq-proto/src/congestion/bbr3/mod.rs @@ -515,9 +515,9 @@ pub struct Bbr3 { /// equivalent to RS.has_data: true once the ACK being processed has delivered a tracked /// packet, so `rs` describes this ACK and is folded into the model when the ACK ends. rs_has_data: bool, - /// Whether the last ACK delivered a tracked packet, so `rs` describes it. A CE event is - /// reported after its ACK ends, and an ACK of only untracked packets leaves an older `rs`. - rs_from_last_ack: bool, + /// The last ACK's own rate sample, if it delivered a tracked packet. A CE event is reported + /// after its ACK ends and after loss detection, which rewrites `rs` per lost packet. + ack_rs: Option, /// equivalent to RS.newly_acked, accumulated over the ACK being processed. It is passed to /// the ACK's model steps rather than kept in `rs`, so no later `set_cwnd` counts it again. newly_acked: u64, @@ -699,7 +699,7 @@ impl Bbr3 { lost: 0, rs: None, rs_has_data: false, - rs_from_last_ack: false, + ack_rs: None, newly_acked: 0, packets: Default::default(), rounds_since_bw_probe: 0, @@ -1624,8 +1624,7 @@ impl Bbr3 { // untracked ACK-only one. The probe's feedback can still be arriving after the ACK // itself ended the probe. _ if self.bw_probe_samples => { - let rs = self.rs.filter(|_| self.rs_from_last_ack); - self.handle_inflight_too_high(now, rs); + self.handle_inflight_too_high(now, self.ack_rs); } _ => {} } @@ -1855,8 +1854,8 @@ impl Bbr3 { let newly_acked = std::mem::take(&mut self.newly_acked); // An ACK that delivered no tracked packet has no sample, and a finished one is never // folded twice. - self.rs_from_last_ack = std::mem::take(&mut self.rs_has_data); - if !self.rs_from_last_ack { + self.ack_rs = None; + if !std::mem::take(&mut self.rs_has_data) { return; } let Some(mut rs) = self.rs else { @@ -1874,6 +1873,7 @@ impl Bbr3 { rs.delivery_rate = rs.delivered as f64 / rs.interval.as_secs_f64(); } self.rs = Some(rs); + self.ack_rs = Some(rs); self.update_model_and_state(rs.last_packet, newly_acked, now); self.update_control_parameters(newly_acked); } @@ -8058,6 +8058,26 @@ mod test { assert_eq!(sim.bbr.inflight_longterm, 100_000); } + /// The transport detects losses between an ACK and its CE event. A loss too small to show + /// inflight too high rewrites `rs` with the lost packet's flight, and CE still bounds the + /// probe by the marked ACK's own sample. + #[test] + fn ce_after_loss_detection_bounds_by_the_ack_sample() { + let mut sim = probing_up(); + // Packet 109 is sent with 100 packets in flight, so losing it alone stays under + // LOSS_THRESH. + sim.send(10 * MS, 90); + sim.ack(20 * MS, 10..20, false); + sim.lose(20 * MS, 109); + assert_eq!(sim.bbr.state, BbrState::ProbeBw(ProbeBwSubstate::Up)); + let now = sim.at(20 * MS); + sim.bbr + .on_congestion_event(now, sim.at(10 * MS), false, true, 0, 19, SpaceKind::Data); + assert_eq!(sim.bbr.state, BbrState::ProbeBw(ProbeBwSubstate::Down)); + // Packet 19 was sent with ten packets in flight. + assert_eq!(sim.bbr.inflight_longterm, 12_000); + } + /// A scripted `Sim` cruising after a loss-free probe of 12 MB/s over a 10ms RTT. fn cruising() -> Sim { let mut sim = probed();