diff --git a/toyos-net-shard/tcp/src/conn.rs b/toyos-net-shard/tcp/src/conn.rs index b0086dfff11..b2cadac0b0a 100644 --- a/toyos-net-shard/tcp/src/conn.rs +++ b/toyos-net-shard/tcp/src/conn.rs @@ -22,6 +22,11 @@ const TS_RECENT_VALID: Duration = Duration::from_secs(24 * 24 * 3600); const HEADERS: u16 = 40; const TS_OPTION: u32 = 12; const MIN_MTU: u16 = 576; +/// RFC 8985 §7.2's WCDelAckT: the longest a peer may delay the ACK of one segment. +const WORST_DELAYED_ACK: Duration = Duration::from_millis(200); +/// What a probe waits past two round trips with more than one segment out: Linux's +/// `TCP_TIMEOUT_MIN`, so an ACK on time is never raced. +const PROBE_SLACK: Duration = Duration::from_millis(2); /// A parsed segment as TCP reads it. #[derive(Clone, Copy, Debug)] @@ -334,6 +339,12 @@ pub struct Sync { rto_pending: bool, timeout_rtx: bool, timing: Option<(Seq, Instant)>, + /// RFC 8985 §7: when the loss probe is due, whether it is, the probe in flight (SND.NXT once + /// it left and whether it was a retransmission), and whether an RTT sample came since the last. + probe_at: Option, + probe_due: bool, + probe: Option<(Seq, bool)>, + sampled: bool, persist: Option, sws: Option, sws_fired: bool, @@ -380,7 +391,7 @@ impl Sync { let mut sync = Self { phase, tx, - rx: Rx::new(p.rcv_next, p.receive_buffer, n.rcv_shift, p.offered), + rx: Rx::new(p.rcv_next, p.receive_buffer, n.rcv_shift, p.offered, now, (path_mss, ts_bytes)), rtt, cc: Cc::new(smss, p.handshake_timeouts), peer_mss: u32::from(n.peer_mss), @@ -399,6 +410,10 @@ impl Sync { rto_pending: false, timeout_rtx: false, timing: None, + probe_at: None, + probe_due: false, + probe: None, + sampled: true, persist: None, sws: None, sws_fired: false, @@ -533,6 +548,10 @@ impl Sync { if placed == Placed::RangeLimit { ctx.log.count(Counter::OooRangeLimit); } + if matches!(placed, Placed::InOrder | Placed::Filled) { + let echo = self.ts.zip(seg.options.timestamps()).map(|(ts, t)| (t.echo, ts.echo_rtt(t.echo, now))); + self.rx.sample_rtt(us32(text.len()), echo, now); + } self.rx.owe_for_text(placed, now); } if fin && !peer_closed { @@ -608,17 +627,36 @@ impl Sync { self.unsolicited(ctx); return false; } - let newly = if self.sack_ok { + let (newly, dsack) = if self.sack_ok { self.tx.read_sack(ack, &seg.options, ctx.log) } else { if seg.options.sack_blocks().len() > 0 { ctx.log.count(Counter::SackUnnegotiated); } - false + (false, None) }; let flight = self.tx.flight(); let acked = ack.since(self.tx.una); let window = u32::from(seg.window).checked_shl(u32::from(self.tx.shift)).unwrap_or(u32::MAX); + let same = acked == 0 && seg.payload.is_empty() && !seg.syn() && !seg.fin() && window == self.tx.wnd; + // A duplicate ACK (RFC 5681 §2) or new SACK information leaves loss to recovery. + if newly || (same && flight > 0) { + self.probe_at = None; + } + // RFC 8985 §7.4.2: at or past the probe's end, a probe of new data, a D-SACK of the probe or + // a duplicate without SACK ends the episode with nothing lost; only an ACK past the end + // without either says a resent probe repaired a loss. + if let Some((end, resent)) = self.probe.filter(|&(end, _)| ack.at_or_after(end) && acked <= flight) { + if !resent || dsack == Some(end) || (same && seg.options.sack_blocks().len() == 0) { + self.probe = None; + } else if ack.after(end) { + self.probe = None; + self.cc.on_loss(flight); + self.cc.cwnd = self.cc.cwnd.min(self.cc.ssthresh); + self.cc.end_recovery(); + ctx.log.count(Counter::LossProbeRecovery); + } + } if acked > 0 && acked <= flight { self.new_ack(seg, ack, acked, flight, ctx); if newly { @@ -698,6 +736,24 @@ impl Sync { self.rtx_timer = None; self.arm(now); } + self.probe_due = false; + self.schedule_probe(now); + } + + /// RFC 8985 §7.2, on new data sent and on an ACK that moves SND.UNA: with SACK, outside fast + /// and RTO recovery, with nothing SACKed, no probe in flight and an RTT sample since the last + /// (§7.3), a probe is due two round trips on, plus WCDelAckT when one segment is out or the + /// slack when more are, and never after the RTO, which it then stands in for. + fn schedule_probe(&mut self, now: Instant) { + self.probe_at = None; + let Some(srtt) = self.rtt.srtt() else { return }; + let Some(rto) = self.rtx_timer else { return }; + let recovering = self.recovery != Recovery::None || (self.episode && self.tx.una.at_or_before(self.recover)); + if !self.sack_ok || recovering || !self.tx.sacked().is_empty() || self.probe.is_some() || !self.sampled || self.persist.is_some() { + return; + } + let delayed = if self.tx.flight() <= self.smss() { WORST_DELAYED_ACK } else { PROBE_SLACK }; + self.probe_at = Some(now.after(srtt.saturating_mul(2).saturating_add(delayed)).min(rto)); } /// RFC 6298 (5.1): the timer runs while sequence space that left is outstanding, and never @@ -717,6 +773,7 @@ impl Sync { Some(rtt) => { let expected = flight.div_ceil(self.smss().saturating_mul(2).max(1)).max(1); self.rtt.sample(rtt, expected); + self.sampled = true; } None => ctx.log.count(Counter::TsEcrInvalid), } @@ -725,6 +782,7 @@ impl Sync { if let Some((_, at)) = self.timing.filter(|&(end, _)| ack.at_or_after(end)) { self.timing = None; self.rtt.sample(now.since(at), 1); + self.sampled = true; } } } @@ -756,6 +814,9 @@ impl Sync { } return; } + // RFC 8985 §7.1: fast recovery starts the loss probe's state afresh. + (self.probe_at, self.probe_due, self.probe) = (None, false, None); + self.arm(ctx.now); let flight = self.tx.flight(); self.cc.on_loss(flight.saturating_sub(self.lt_bytes)); self.recover = self.tx.nxt.sub(1); @@ -785,7 +846,7 @@ impl Sync { if self.persist.is_none() { let rto = self.rtt.rto(); self.persist = Some(Persist { at: now.after(rto), interval: rto, from: self.tx.nxt, due: false, unanswered: false }); - self.rtx_timer = None; + (self.rtx_timer, self.probe_at, self.probe_due) = (None, None, false); } } else if let Some(persist) = self.persist.take() { if self.tx.nxt.after(persist.from) { @@ -851,10 +912,10 @@ impl Sync { Ok(n) } - pub fn recv(&mut self, out: &mut [u8]) -> Result { + pub fn recv(&mut self, out: &mut [u8], now: Instant) -> Result { let n = self.rx.read(out); if n > 0 { - self.rx.after_read(self.smss()); + self.rx.after_read(n, now, self.smss()); Ok(Received::Data(n)) } else if self.rx.closed { Ok(Received::End) @@ -866,13 +927,13 @@ impl Sync { } /// [`Self::recv`] for a reader that may take less than it is shown. - pub fn recv_with(&mut self, take: impl FnOnce(&[u8]) -> usize) -> Result { + pub fn recv_with(&mut self, take: impl FnOnce(&[u8]) -> usize, now: Instant) -> Result { if self.rx.unread() == 0 { return if self.rx.closed { Ok(Received::End) } else { Err(Error::WouldBlock) }; } let n = self.rx.read_with(take); if n > 0 { - self.rx.after_read(self.smss()); + self.rx.after_read(n, now, self.smss()); } Ok(Received::Data(n)) } @@ -935,6 +996,7 @@ impl Sync { pub fn deadline(&self, ctx: &Ctx<'_>) -> Option { [ self.rtx_timer, + self.probe_at, self.persist.filter(|p| !p.due).map(|p| p.at), self.keepalive_at(ctx), self.sws.filter(|_| !self.sws_fired), @@ -957,6 +1019,10 @@ impl Sync { ctx.log.count(Counter::OrphanIdleAbort); return Tick::Orphan; } + // RFC 8985 §7.3: the RTO runs again from the probe's leaving. + if self.probe_at.is_some_and(|at| at <= now) { + (self.probe_at, self.probe_due, self.rtx_timer) = (None, true, None); + } if self.rtx_timer.is_some_and(|at| at <= now) { self.expire(ctx); } @@ -991,6 +1057,7 @@ impl Sync { self.rto_pending = true; self.rtt.back_off(); self.timing = None; + (self.probe_at, self.probe_due, self.probe) = (None, false, None); let repeat = core::mem::replace(&mut self.timeout_rtx, true); self.cc.on_timeout(self.tx.flight(), repeat); self.recovery = Recovery::None; @@ -1024,6 +1091,9 @@ impl Sync { self.rx.dup_owed = self.rx.dup_owed.saturating_sub(1); return Ok(true); } + if self.probe_due && self.loss_probe(ctx, exit)? { + return Ok(true); + } if self.retransmission(ctx, exit)? || self.new_data(ctx, exit)? || self.probe(ctx, exit)? { return Ok(true); } @@ -1106,6 +1176,9 @@ impl Sync { self.timing = Some((end, now)); } self.arm(now); + if len > 0 && !resent { + self.schedule_probe(now); + } self.refresh(now); self.acknowledged(); } @@ -1286,6 +1359,33 @@ impl Sync { Ok(true) } + /// RFC 8985 §7.3: a segment of new data where the peer's window takes one whole, outside cwnd; + /// else the last segment sent again. Its ACK carries what the lost one would have: the window + /// reopened, or the tail's acknowledgment. + fn loss_probe(&mut self, ctx: &mut Ctx<'_>, exit: &mut dyn Exit) -> Result { + let blocks = self.blocks(); + let room = self.room(&blocks); + let whole = self.tx.unsent().min(room); + let usable = u32::try_from(self.tx.usable().max(0)).unwrap_or(u32::MAX); + let (start, len) = if whole > 0 && usable >= whole { + (self.tx.nxt, whole) + } else { + let end = self.tx.nxt.earlier(self.tx.data_end()); + let start = if end.since(self.tx.una) > room { end.sub(room) } else { self.tx.una }; + (start, if start.before(end) { end.since(start) } else { 0 }) + }; + if len == 0 { + self.probe_due = false; + self.arm(ctx.now); + return Ok(false); + } + let resent = start.before(self.tx.nxt); + self.emit(ctx, exit, start, len, false, blocks)?; + (self.probe_at, self.probe_due, self.probe, self.sampled) = (None, false, Some((self.tx.nxt, resent)), false); + ctx.log.count(Counter::LossProbe); + Ok(true) + } + /// A due persist probe (RFC 9293 §3.8.6.1) or keepalive (§3.8.4). fn probe(&mut self, ctx: &mut Ctx<'_>, exit: &mut dyn Exit) -> Result { let now = ctx.now; diff --git a/toyos-net-shard/tcp/src/counters.rs b/toyos-net-shard/tcp/src/counters.rs index 5a6a941a545..d3dda025058 100644 --- a/toyos-net-shard/tcp/src/counters.rs +++ b/toyos-net-shard/tcp/src/counters.rs @@ -65,6 +65,8 @@ toyos_net_wire::counters! { SackRecovery = "tcp.sack-recovery"; LimitedTransmit = "tcp.limited-transmit"; PersistProbe = "tcp.persist-probe"; + LossProbe = "tcp.loss-probe"; + LossProbeRecovery = "tcp.loss-probe-recovery"; KeepaliveProbe = "tcp.keepalive-probe"; EventOverflow = "tcp.event-overflow"; } diff --git a/toyos-net-shard/tcp/src/lib.rs b/toyos-net-shard/tcp/src/lib.rs index 1c9e560a264..0cbb12006d7 100644 --- a/toyos-net-shard/tcp/src/lib.rs +++ b/toyos-net-shard/tcp/src/lib.rs @@ -104,7 +104,9 @@ pub struct Secrets { pub struct Config { /// The outgoing interface's IP MTU. pub mtu: u16, + /// The most a connection's receive buffer grows to; it sets the window scale offered. pub receive_buffer: u32, + /// The most a connection's send buffer holds; storage is taken as it fills. pub send_buffer: u32, pub secrets: Secrets, } @@ -150,6 +152,9 @@ pub mod limits { pub const EVENTS: usize = 1_024; /// The largest window scaling can offer (RFC 7323 §2.3): no receive buffer is larger. pub const RECEIVE_BUFFER_MAX: u32 = 65_535 << 14; + /// A connection's receive capacity until its reader grows it (`rx`'s module): the most a + /// SYN's unscaled window offers, and the most a connection nobody reads ever holds. + pub const RECEIVE_BUFFER_INITIAL: u32 = 65_535; } /// A [`Config`] that [`Tcp::new`] refuses. diff --git a/toyos-net-shard/tcp/src/props.rs b/toyos-net-shard/tcp/src/props.rs index fd9faa06d53..65b77b78886 100644 --- a/toyos-net-shard/tcp/src/props.rs +++ b/toyos-net-shard/tcp/src/props.rs @@ -114,7 +114,8 @@ impl Checker { } } - /// PROP-01, 04, 05, 06 and 07 on every synchronized connection, and every refusal counted. + /// PROP-01, 04, 05, 06 and 07 and `rx`'s capacity invariant on every synchronized connection, + /// and every refusal counted. fn state(&mut self, node: usize, tcp: &mut Tcp, now: Instant) { for (tuple, sync) in tcp.each_sync() { let expiries = tcp.counters().get(Counter::Rto); @@ -134,6 +135,8 @@ impl Checker { assert!(tx.nxt.at_or_before(left), "SND.NXT {:?} past what left, {left:?}", tx.nxt); } let edge = sync.rx.edge(); + let promised = sync.rx.unread().saturating_add(usize::try_from(edge.since(sync.rx.next)).unwrap_or(usize::MAX)); + assert!(promised <= sync.rx.capacity(), "rx: unread + window {promised} past the capacity {}", sync.rx.capacity()); if let Some(&(acked, offered)) = self.told.get(&(node, tuple)) { let sent = sync.rx.last_ack_sent; assert!(sent.at_or_before(acked), "acknowledged to {sent:?}, past what left, {acked:?}"); @@ -250,7 +253,8 @@ impl Run { judged: 0, sacked: [Vec::new(), Vec::new()], })); - let mut net = Net::new(10); + // Odd seeds run at netstack's buffer, which grows and is scaled. + let mut net = Net::buffered(10, if seed.is_multiple_of(2) { 65_535 } else { 4 << 20 }); net.keep_streams = true; net.impair = link(Rng::new(seed ^ 0xaaaa), s); let shared = Rc::clone(&checker); diff --git a/toyos-net-shard/tcp/src/ring.rs b/toyos-net-shard/tcp/src/ring.rs index e91f30384fb..aaf45daab7f 100644 --- a/toyos-net-shard/tcp/src/ring.rs +++ b/toyos-net-shard/tcp/src/ring.rs @@ -1,5 +1,6 @@ -//! A byte ring of fixed capacity whose storage is taken at the first write, so a connection that -//! carries no data holds no buffer. Offsets count from the oldest byte held. +//! A byte ring whose storage grows, by doubling, with the furthest offset written and never past +//! the capacity, so a connection holds memory for what it holds and not for what it may. The +//! capacity only rises. Offsets count from the oldest byte held. use alloc::vec::Vec; @@ -37,18 +38,35 @@ impl Ring { self.capacity.saturating_sub(self.len) } - /// The physical index of `offset`, which is below the capacity. + /// Raises the capacity to `capacity`; a lower one is no change. + pub fn grow(&mut self, capacity: usize) { + self.capacity = self.capacity.max(capacity); + } + + /// The physical index of `offset`, which is below the storage's length. fn at(&self, offset: usize) -> usize { let index = self.head.saturating_add(offset); - index.checked_sub(self.capacity).unwrap_or(index) + index.checked_sub(self.bytes.len()).unwrap_or(index) + } + + /// Storage for every offset below `end`, at most the capacity: each byte stored keeps its + /// offset, laid out again from the head. + fn reserve(&mut self, end: usize) { + if end > self.bytes.len() { + self.bytes.rotate_left(self.head); + self.head = 0; + let len = end.max(self.bytes.len().saturating_mul(2)).min(self.capacity); + self.bytes.resize(len, 0); + } } /// Stores `data` at `offset` without counting it held; what lies past the capacity is not stored. pub fn write_at(&mut self, offset: usize, data: &[u8]) -> usize { - if self.bytes.is_empty() { - self.bytes.resize(self.capacity, 0); - } let fits = self.capacity.saturating_sub(offset).min(data.len()); + if fits == 0 { + return 0; + } + self.reserve(offset.saturating_add(fits)); let start = self.at(offset); let (data, _) = data.split_at(fits); let first = self.bytes.get_mut(start..).map_or(0, |to| copy(to, data)); @@ -81,7 +99,7 @@ impl Ring { let offset = offset.min(self.len); let len = len.min(self.len.saturating_sub(offset)); let start = self.at(offset); - let first_len = len.min(self.capacity.saturating_sub(start)); + let first_len = len.min(self.bytes.len().saturating_sub(start)); let first = self.bytes.get(start..start.saturating_add(first_len)).unwrap_or_default(); let second = self.bytes.get(..len.saturating_sub(first_len)).unwrap_or_default(); (first, second) diff --git a/toyos-net-shard/tcp/src/rx.rs b/toyos-net-shard/tcp/src/rx.rs index 8a5900d80f0..ec4873b08b8 100644 --- a/toyos-net-shard/tcp/src/rx.rs +++ b/toyos-net-shard/tcp/src/rx.rs @@ -4,6 +4,23 @@ //! `unread + (edge − next) ≤ capacity` always holds, so every byte the peer may send has room. The //! edge never retreats (RFC 7323 §2.4), bytes already stored are never overwritten, and bytes once //! reported in a SACK block are kept until delivered: this receiver never reneges. +//! +//! **The capacity grows only by what the user reads.** It starts at +//! [`limits::RECEIVE_BUFFER_INITIAL`](crate::limits::RECEIVE_BUFFER_INITIAL) and, once a round +//! trip has passed, rises to twice what the user read per round trip, up to the configured +//! buffer: enough that a reader keeping up never waits on the window while the sender's rate +//! doubles. Text nobody reads grows nothing, so a peer alone cannot make a connection hold more +//! than the initial capacity, and a connection whose round trip is unknown keeps it. +//! +//! **The round trip is the receiver's own** (Linux's `tcp_rcv_rtt_measure_ts` and +//! `tcp_rcv_rtt_measure`): a downloader sends no data, so the sender's SRTT keeps the handshake's +//! sample while queueing lengthens the path. With timestamps, a full-sized segment carrying a +//! TSecr not yet seen samples the age of that echo, at least one tick and at most twice the +//! estimate, averaged with gain 1/8: an echo aged by the peer's own silence moves the estimate +//! by an eighth at most. Full-sized is the largest segment the peer sends, learned from what +//! arrives (Linux's `tcp_measure_rcv_mss`), never this end's send MSS: a peer whose segments are +//! shorter would never be sampled. Without timestamps, the time the peer took to fill the window +//! offered is at least one round trip, and the least such time is kept. use alloc::vec::Vec; use core::time::Duration; @@ -16,6 +33,8 @@ use crate::Instant; pub const OOO_RANGES: usize = 32; pub const DELAYED_ACK: Duration = Duration::from_millis(40); +/// One TSval tick (RFC 7323 §5.4): an echo's age is never less. +const TICK: Duration = Duration::from_millis(1); /// Distinct duplicate ACKs held for a starved transmit: enough for the peer's fast retransmit. pub const DUP_ACKS_OWED: u8 = 3; @@ -65,16 +84,43 @@ pub struct Rx { /// The window, in bytes, the last segment sent offered. pub last_window: u32, pub last_ack_sent: Seq, + /// The most the capacity grows to. + max: usize, + /// When the current measurement began, and what the user has read since. + round: (Instant, usize), + /// The receiver's round-trip estimate. + rtt: Option, + /// The last TSecr sampled, or without timestamps the edge whose filling is timed, and since when. + sampler: Sampler, + /// The largest segment the peer sends, the length of the last shorter one, and the bounds + /// of both: what our SYN offered less the timestamp option, and the floor MSS less it. + rcv_mss: u32, + short: u32, + mss_bounds: (u32, u32), +} + +#[derive(Clone, Copy, Debug)] +enum Sampler { + Echo(Option), + Fill(Option<(Seq, Instant)>), } impl Rx { - /// `next` is IRS + 1; the SYN or SYN-ACK offered `window`, unscaled. - pub fn new(next: Seq, capacity: usize, shift: u8, window: u32) -> Self { + /// `next` is IRS + 1; the SYN or SYN-ACK offered `window`, unscaled; the capacity grows to + /// `max`, or to the most a window field at `shift` offers if that is less: an edge past it + /// would be one the peer was never told of. `(mss, options)`: the MSS our SYN offered and the + /// timestamp option's length on every segment, 0 without timestamps, which picks how the + /// round trip is sampled. + pub fn new(next: Seq, max: usize, shift: u8, window: u32, now: Instant, (mss, options): (u32, u32)) -> Self { + let initial = usize::try_from(crate::limits::RECEIVE_BUFFER_INITIAL).unwrap_or(usize::MAX); + let offerable = usize::from(u16::MAX).checked_shl(u32::from(shift)).unwrap_or(usize::MAX); + let max = max.min(offerable); + let floor = u32::from(crate::limits::MSS_FLOOR).saturating_sub(options); Self { next, edge: next.add(window), shift, - buf: Ring::new(capacity), + buf: Ring::new(max.min(initial)), ranges: Vec::new(), stamp: 0, fin: None, @@ -88,6 +134,59 @@ impl Rx { trigger: None, last_window: window, last_ack_sent: next, + max, + round: (now, 0), + rtt: None, + sampler: if options > 0 { Sampler::Echo(None) } else { Sampler::Fill(None) }, + rcv_mss: floor, + short: 0, + mss_bounds: (floor, mss.saturating_sub(options)), + } + } + + /// The most the peer may have unread and in flight to us: the receive buffer's present size. + #[cfg(test)] + pub fn capacity(&self) -> usize { + self.buf.capacity() + } + + /// After in-order text of `len` bytes arrived at `now`, with its TSecr and that echo's age + /// where it is one this end sent. + pub fn sample_rtt(&mut self, len: u32, echo: Option<(u32, Option)>, now: Instant) { + self.measure_mss(len); + match &mut self.sampler { + Sampler::Echo(last) => { + let Some((echo, age)) = echo.filter(|&(echo, _)| *last != Some(echo)) else { return }; + *last = Some(echo); + let Some(age) = age.filter(|_| len >= self.rcv_mss) else { return }; + let age = age.max(TICK); + self.rtt = Some(self.rtt.map_or(age, |rtt| rtt.saturating_mul(7).saturating_add(age.min(rtt.saturating_mul(2))).checked_div(8).unwrap_or(rtt))); + } + Sampler::Fill(mark) => { + if let Some((edge, since)) = *mark { + if self.next.before(edge) { + return; + } + let took = now.since(since).max(Duration::from_micros(1)); + self.rtt = Some(self.rtt.map_or(took, |rtt| rtt.min(took))); + } + *mark = Some((self.edge, now)); + } + } + } + + /// Linux's `tcp_measure_rcv_mss`: a segment as long as the largest so far raises it, up to what + /// our SYN offered; two in a row of one shorter length, not under the floor, lower it to that. + fn measure_mss(&mut self, len: u32) { + let (floor, offered) = self.mss_bounds; + let short = core::mem::replace(&mut self.short, 0); + if len >= self.rcv_mss { + self.rcv_mss = len.min(offered); + } else if len >= floor { + self.short = len; + if len == short { + self.rcv_mss = len; + } } } @@ -252,6 +351,11 @@ impl Rx { /// The right edge and window field a segment built now carries; [`Self::advertise`] commits them. pub fn offer(&self, mss: u32) -> (Seq, u16) { let edge = self.candidate(mss).unwrap_or(self.edge); + // A window the shift would round to zero is offered as one unit where the buffer has the + // room, as Linux does: else a sender owing the text of a hole waits on a window not shut. + let unit = 1u32.checked_shl(u32::from(self.shift)).unwrap_or(u32::MAX); + let window = edge.since(self.next); + let edge = if window > 0 && window < unit && unit <= self.free() { self.next.add(unit) } else { edge }; let field = (edge.since(self.next) >> self.shift).min(u32::from(u16::MAX)); (edge, u16::try_from(field).unwrap_or(u16::MAX)) } @@ -266,13 +370,17 @@ impl Rx { if !self.ranges.is_empty() { return None; } - let free = u32::try_from(self.buf.capacity().saturating_sub(self.buf.len())).unwrap_or(u32::MAX); let mask = u32::MAX.checked_shl(u32::from(self.shift)).unwrap_or(0); - let candidate = self.next.add(free & mask); + let candidate = self.next.add(self.free() & mask); let threshold = self.threshold(mss); (candidate.after(self.edge) && candidate.since(self.edge) >= threshold).then_some(candidate) } + /// Room for text the user has not read. + fn free(&self) -> u32 { + u32::try_from(self.buf.capacity().saturating_sub(self.buf.len())).unwrap_or(u32::MAX) + } + fn threshold(&self, mss: u32) -> u32 { (u32::try_from(self.buf.capacity()).unwrap_or(u32::MAX) / 2).min(mss) } @@ -354,11 +462,90 @@ impl Rx { self.buf.consume(n); } - /// After the user read: a window update is owed when the window offered was below the - /// threshold and the edge would now move by at least it (RFC 9293 §3.8.6.2.2). - pub fn after_read(&mut self, mss: u32) { - if self.candidate(mss).is_some() && self.last_window < self.threshold(mss) { - self.ack_now = true; + /// After the user read `n` bytes: the capacity grows by the module's rule, and a window update + /// is owed when the window offered was below the threshold and the edge would now move by at + /// least it (RFC 9293 §3.8.6.2.2). + pub fn after_read(&mut self, n: usize, now: Instant, mss: u32) { + let (began, read) = self.round; + let read = read.saturating_add(n); + self.round = (began, read); + if let Some(rtt) = self.rtt.filter(|&rtt| now.since(began) >= rtt) { + let read = u128::try_from(read).unwrap_or(u128::MAX); + let per_rtt = read.saturating_mul(rtt.as_nanos()).checked_div(now.since(began).as_nanos()).unwrap_or(0); + let want = usize::try_from(per_rtt.saturating_mul(2)).unwrap_or(usize::MAX); + self.buf.grow(want.min(self.max)); + self.round = (now, 0); } + if let Some(candidate) = self.candidate(mss) { + if self.last_window < self.threshold(mss) || candidate.since(self.next) >= self.window().saturating_mul(2) { + self.ack_now = true; + } + } + } +} + +#[cfg(test)] +mod tests { + use super::*; + + fn ms(n: u64) -> Duration { + Duration::from_millis(n) + } + + /// A timestamped connection's receive side at t = 0: our SYN offered an MSS of 1,460. + fn stamped() -> Rx { + Rx::new(Seq::new(1), 4 << 20, 7, 65_535, Instant::from_nanos(0), (1460, 12)) + } + + fn echo(rx: &mut Rx, len: u32, tsecr: u32, age: Duration) { + rx.sample_rtt(len, Some((tsecr, Some(age))), Instant::from_nanos(u64::from(tsecr) * 1_000_000)); + } + + /// Full-sized is the peer's largest segment: 1,380-byte segments are sampled though our own + /// send MSS is larger. After a 1,448-byte one, a single 1,380-byte segment is not; the second + /// in a row lowers full-sized to 1,380 and is. + #[test] + fn full_sized_is_learned_from_what_arrives() { + let mut rx = stamped(); + echo(&mut rx, 1380, 1, ms(15)); + assert_eq!(rx.rtt, Some(ms(15))); + echo(&mut rx, 1448, 2, ms(15)); + echo(&mut rx, 1380, 3, ms(23)); + assert_eq!(rx.rtt, Some(ms(15))); + echo(&mut rx, 1380, 4, ms(23)); + assert_eq!(rx.rtt, Some(ms(16))); + } + + /// The echo of the last ACK before the peer fell silent for 10 s ages by the silence: it moves + /// the estimate by an eighth, no more. An echo of this tick's TSval ages one tick. + #[test] + fn an_echo_aged_by_the_peers_silence_moves_the_estimate_an_eighth() { + let mut rx = stamped(); + for tsecr in 1..40 { + echo(&mut rx, 1448, tsecr, ms(16)); + } + echo(&mut rx, 1448, 10_040, ms(10_016)); + assert_eq!(rx.rtt, Some(ms(18))); + echo(&mut rx, 1448, 10_041, Duration::ZERO); + assert_eq!(rx.rtt, Some(Duration::from_micros(15_875))); + } + + /// Without timestamps the time to fill a window offered is at least a round trip, so the + /// least is kept: the window offered at 0 ms fills at 80 ms, the next at 200 ms. + #[test] + fn without_timestamps_the_least_fill_time_is_kept() { + let mut rx = Rx::new(Seq::new(1), 65_535, 0, 65_535, Instant::from_nanos(0), (1460, 0)); + let arrive = |rx: &mut Rx, len: usize, at_ms: u64| { + let next = rx.next; + assert_eq!(rx.place(next, &vec![0; len]), Placed::InOrder); + rx.sample_rtt(u32::try_from(len).unwrap(), None, Instant::from_nanos(at_ms * 1_000_000)); + rx.read(&mut vec![0; len]); + rx.advertise(1460); + }; + arrive(&mut rx, 1460, 0); + arrive(&mut rx, 64_075, 80); + assert_eq!(rx.rtt, Some(ms(80))); + arrive(&mut rx, 65_535, 200); + assert_eq!(rx.rtt, Some(ms(80))); } } diff --git a/toyos-net-shard/tcp/src/stack.rs b/toyos-net-shard/tcp/src/stack.rs index b0c7a839123..acb35480fef 100644 --- a/toyos-net-shard/tcp/src/stack.rs +++ b/toyos-net-shard/tcp/src/stack.rs @@ -414,7 +414,7 @@ impl Tcp { Local { mss: mtu.saturating_sub(40), shift: self.shift, - window: u16::try_from(self.config.receive_buffer.min(u32::from(u16::MAX))).unwrap_or(u16::MAX), + window: u16::try_from(self.config.receive_buffer.min(limits::RECEIVE_BUFFER_INITIAL)).unwrap_or(u16::MAX), receive_buffer: usize::try_from(self.config.receive_buffer).unwrap_or(0), send_buffer: usize::try_from(self.config.send_buffer).unwrap_or(0), mtu, @@ -1009,7 +1009,7 @@ impl Tcp { let conn = self.conn(id)?; let result = match &mut conn.state { Tcb::SynSent(_) | Tcb::SynRcvd(_) => Err(Error::WouldBlock), - Tcb::Sync(sync) => sync.recv(out), + Tcb::Sync(sync) => sync.recv(out, now), Tcb::Ended(Ended { failure: Some(failure), .. }) => Err(Error::Failed(*failure)), Tcb::Ended(Ended { failure: None, rx }) => match rx.as_mut().map(|rx| rx.read(out)) { Some(n) if n > 0 => Ok(Received::Data(n)), @@ -1028,7 +1028,7 @@ impl Tcp { let conn = self.conn(id)?; let result = match &mut conn.state { Tcb::SynSent(_) | Tcb::SynRcvd(_) => Err(Error::WouldBlock), - Tcb::Sync(sync) => sync.recv_with(take), + Tcb::Sync(sync) => sync.recv_with(take, now), Tcb::Ended(Ended { failure: Some(failure), .. }) => Err(Error::Failed(*failure)), Tcb::Ended(Ended { failure: None, rx }) => match rx.as_mut() { Some(rx) if rx.unread() > 0 => Ok(Received::Data(rx.read_with(take))), diff --git a/toyos-net-shard/tcp/src/tx.rs b/toyos-net-shard/tcp/src/tx.rs index 203fc7354c2..d56a94ba09f 100644 --- a/toyos-net-shard/tcp/src/tx.rs +++ b/toyos-net-shard/tcp/src/tx.rs @@ -130,16 +130,18 @@ impl Tx { }); } - /// Reads an ACK's SACK blocks (RFC 2018 §3, RFC 2883 §4). `true` when a block covered bytes - /// the scoreboard did not hold: RFC 6675's duplicate acknowledgment. - pub fn read_sack(&mut self, ack: Seq, options: &TcpOptions<'_>, log: &mut Log) -> bool { + /// Reads an ACK's SACK blocks (RFC 2018 §3, RFC 2883 §4): whether a block covered bytes the + /// scoreboard did not hold, RFC 6675's duplicate acknowledgment, and a D-SACK's right edge. + pub fn read_sack(&mut self, ack: Seq, options: &TcpOptions<'_>, log: &mut Log) -> (bool, Option) { let second = options.sack_blocks().nth(1); let mut newly = false; + let mut dsack = None; for (i, block) in options.sack_blocks().enumerate() { let (left, right) = (Seq::from(block.left), Seq::from(block.right)); let inside_second = second.is_some_and(|s| Seq::from(s.left).at_or_before(left) && right.at_or_before(Seq::from(s.right))); if i == 0 && (right.at_or_before(ack) || inside_second) { log.count(Counter::DsackRcvd); + dsack = Some(right); continue; } if !(ack.before(left) && left.before(right) && right.at_or_before(self.nxt)) { @@ -152,7 +154,7 @@ impl Tx { newly |= self.insert(left, right); } } - newly + (newly, dsack) } fn insert(&mut self, left: Seq, right: Seq) -> bool { diff --git a/toyos-net-shard/tcp/tests/common/net.rs b/toyos-net-shard/tcp/tests/common/net.rs index 66227e1c1f9..20b10075b58 100644 --- a/toyos-net-shard/tcp/tests/common/net.rs +++ b/toyos-net-shard/tcp/tests/common/net.rs @@ -152,11 +152,16 @@ pub fn ns(ms: u64) -> u64 { impl Net { /// Two nodes, a one-way delay of `rtt / 2`, clocks starting an hour in. pub fn new(rtt_ms: u64) -> Self { + Self::buffered(rtt_ms, 65_535) + } + + /// [`Self::new`] with each node's receive and send buffers at `buffer`. + pub fn buffered(rtt_ms: u64, buffer: u32) -> Self { let node = |addr, seed: u8| Node { tcp: Tcp::new(toyos_net_tcp::Config { mtu: 1500, - receive_buffer: 65_535, - send_buffer: 65_535, + receive_buffer: buffer, + send_buffer: buffer, secrets: toyos_net_tcp::Secrets { isn: key(seed), timestamp: key(seed.wrapping_add(0x10)), diff --git a/toyos-net-shard/tcp/tests/net.rs b/toyos-net-shard/tcp/tests/net.rs index e60980f98f5..91b753e36f9 100644 --- a/toyos-net-shard/tcp/tests/net.rs +++ b/toyos-net-shard/tcp/tests/net.rs @@ -68,22 +68,32 @@ fn s_net_002_newreno_when_one_side_has_no_sack() { } } +/// Every second pure ACK is lost, either parity. At 65,535 the window, not cwnd, bounds the sender, +/// so a round's window reopens on the update its reader's read sends, and that update is sometimes +/// the tail's only ACK; at netstack's buffer cwnd bounds it. No loss is ever inferred: no RTO +/// fires, and what is sent twice is only a loss probe's segment, which the receiver reports as a +/// duplicate. #[test] fn s_net_005_lost_acks() { - let mut net = Net::new(10); - net.connections(1, 80, [MIB, 0]); - let mut acks = 0; - net.impair = Box::new(move |_, o| { - if o.payload.is_empty() && o.flags & (SYN | FIN | RST) == 0 { - acks += 1; - if acks % 2 == 0 { - return Fate::Drop; + for (buffer, parity) in [(65_535, 0), (65_535, 1), (BUFFER, 0), (BUFFER, 1)] { + let mut net = Net::buffered(10, buffer); + net.connections(1, 80, [MIB, 0]); + let mut acks = 0; + net.impair = Box::new(move |_, o| { + if o.payload.is_empty() && o.flags & (SYN | FIN | RST) == 0 { + acks += 1; + if acks % 2 == parity { + return Fate::Drop; + } } - } - Fate::Pass - }); - finish(&mut net); - assert_eq!(net.count(0, Counter::RetransmitBytes), 0); + Fate::Pass + }); + finish(&mut net); + let count = |c| net.count(0, c); + let (probes, arm) = (count(Counter::LossProbe), format!("buffer {buffer}, parity {parity}")); + assert_eq!((count(Counter::Rto), count(Counter::LossProbeRecovery)), (0, 0), "{arm}"); + assert!(count(Counter::RetransmitBytes) <= probes * 1448 && count(Counter::DsackRcvd) >= probes, "{arm}: {probes} probes"); + } } #[test] @@ -303,3 +313,229 @@ fn s_ls_012_a_syn_flood() { h.input(64_010, seg(5001).ack(1001)); assert!(h.tcp.accept(listener).unwrap().is_some()); } + +/// The buffer netstack gives each connection each way: 4 MiB, window scale 7. +const BUFFER: u32 = 4 << 20; + +/// What one node saw of its connection over a whole transfer. +#[derive(Debug, Default)] +struct Peaks { + /// SND.NXT − SND.UNA. + flight: u32, + /// RCV.WND as offered: the right edge less RCV.NXT. + window: u32, + unread: usize, + /// (snd_shift, rcv_shift), once established. + shifts: Option<(u8, u8)>, +} + +/// Steps `net` a millisecond at a time for at most `limit_ms` or until it finishes, reading each +/// node's connection after every step. +fn watch(net: &mut Net, limit_ms: u64) -> [Peaks; 2] { + let mut peaks = [Peaks::default(), Peaks::default()]; + let end = net.now + ns(limit_ms); + while !net.finished() && net.now < end { + net.advance(1); + for i in 0..net.apps.len() { + let (node, id) = (net.apps[i].node, net.apps[i].id); + let Some(info) = net.nodes[node].tcp.info(id) else { continue }; + let p = &mut peaks[node]; + p.flight = p.flight.max(info.snd_nxt.since(info.snd_una)); + p.window = p.window.max(info.rcv_edge.since(info.rcv_nxt)); + p.unread = p.unread.max(info.unread); + assert_eq!(*p.shifts.get_or_insert((info.snd_shift, info.rcv_shift)), (info.snd_shift, info.rcv_shift)); + } + } + peaks +} + +/// RFC 7323 §2: both SYNs carry Window Scale, so each end offers its buffer at shift 7, and over +/// a gigabit link with a 20 ms round trip a reader that keeps up grows its window to the whole +/// 4 MiB, and each way the sender keeps more than 65,535 bytes in flight; every byte arrives. +/// With one data segment in 3,000 dropped, out-of-order text is held across the buffer's growth +/// and the flight still passes 65,535. +#[test] +fn rfc_7323_2_a_scaled_window_carries_more_than_64_kib_in_flight_each_way() { + for drop in [None, Some(3000)] { + let mut net = Net::buffered(20, BUFFER); + for node in &mut net.nodes { + node.credit_per_ms = Some(125_000); + } + net.connections(1, 80, [32 * MIB, 32 * MIB]); + if let Some(n) = drop { + let (mut impair, _) = every_nth(n); + net.impair = Box::new(move |from, o| impair(from, o)); + } + let peaks = watch(&mut net, 600_000); + assert!(net.finished(), "the transfer did not finish:\n{}", net.dump()); + net.assert_exact(); + for (node, p) in peaks.iter().enumerate() { + assert_eq!(p.shifts, Some((7, 7)), "node {node}"); + assert!(p.flight > 65_535, "node {node}, drop {drop:?}: {p:?}"); + assert!(drop.is_some() || p.window == BUFFER, "node {node}: {p:?}"); + } + } +} + +/// RFC 7323 §2.2: a peer whose SYN carries no Window Scale gets an unscaled window, however large +/// the buffer behind it, and none of its windows is read scaled: neither end has more than 65,535 +/// bytes in flight, and every byte arrives. +#[test] +fn rfc_7323_2_2_no_scaling_unless_both_syns_offer_it() { + let mut net = Net::buffered(20, BUFFER); + net.rewrite = Some(Box::new(bare_syn)); + net.connections(1, 80, [4 * MIB, 4 * MIB]); + let peaks = watch(&mut net, 600_000); + assert!(net.finished(), "the transfer did not finish:\n{}", net.dump()); + net.assert_exact(); + for (node, p) in peaks.iter().enumerate() { + assert_eq!(p.shifts, Some((0, 0)), "node {node}"); + assert!(p.flight <= 65_535 && p.window <= 65_535, "node {node}: {p:?}"); + } +} + +/// A receive buffer grows only by what its reader takes: a peer that sends into a connection +/// nobody reads for ten seconds finds 65,535 bytes of room and no more, and once the reader reads +/// the buffer grows and every byte arrives. +#[test] +fn a_receive_buffer_nobody_reads_never_grows() { + let mut net = Net::buffered(20, BUFFER); + net.connections(1, 80, [16 * MIB, 0]); + net.run(100, |n| n.apps.len() == 2); + let reader = net.apps.iter().position(|a| a.node == 1).unwrap(); + net.apps[reader].reading = false; + net.apps[reader].read_limit = Some(0); + let held = watch(&mut net, 10_000); + assert!(held[1].unread <= 65_535 && held[1].window <= 65_535, "{:?}", held[1]); + assert!(held[0].flight <= 65_535, "{:?}", held[0]); + net.apps[reader].read_limit = None; + net.apps[reader].reading = true; + let read = watch(&mut net, 600_000); + assert!(net.finished(), "the transfer did not finish:\n{}", net.dump()); + net.assert_exact(); + assert!(read[1].window > 65_535, "{:?}", read[1]); +} + +/// Node 0's SYN with its timestamps taken out: both ends scale and SACK, and neither stamps. +fn syn_without_timestamps(from: usize, o: &O) -> Option> { + (from == 0 && o.flags & SYN != 0).then(|| { + let s = seg(o.seq).syn().wnd(o.wnd).mss(o.mss.unwrap()).sackok().ws(o.ws.unwrap()); + s.bytes(o.src, o.dst, 0, 0) + }) +} + +/// The receive buffer grows to twice what its reader takes per round trip, by the receiver's own +/// estimate of the round trip, with timestamps and without. The path's round trip quadruples once +/// the handshake is done, so a downloader's SRTT, which only its own data samples, still says +/// 20 ms. The reader takes a burst every millisecond twice per 20 ms: it is never offered more +/// than twice its rate over an estimate at most half again the path's 80 ms, plus a burst and a +/// segment, and once the window has grown it is never kept waiting, which an estimate under the +/// path's would. +#[test] +fn the_receive_window_grows_by_what_is_read_per_round_trip() { + const BURST: u64 = 40_000; + const PERIOD: u64 = 20; + const RATE: u64 = 2 * BURST / PERIOD; + const ESTIMATE: u64 = 120; + for timestamps in [true, false] { + let mut net = Net::buffered(20, BUFFER); + net.auto_shut = false; + if !timestamps { + net.rewrite = Some(Box::new(syn_without_timestamps)); + } + net.connections(1, 80, [64 * MIB, 0]); + assert!(net.run(100, |n| n.apps.len() == 2)); + net.delay = ns(40); + let reader = net.apps.iter().position(|a| a.node == 1).unwrap(); + let handshake = net.nodes[1].tcp.info(net.apps[reader].id).unwrap().srtt; + assert!(handshake < Some(std::time::Duration::from_millis(30))); + let (mut window, mut settled) = (0, 0); + for ms in 0..4_000 { + if ms == 2_000 { + settled = net.apps[reader].received.1; + } + let burst = ms % PERIOD < 2; + (net.apps[reader].reading, net.apps[reader].read_limit) = (burst, burst.then_some(BURST as usize)); + net.advance(1); + let info = net.nodes[1].tcp.info(net.apps[reader].id).unwrap(); + assert_eq!(info.ts_recent.is_some(), timestamps); + window = window.max(u64::from(info.rcv_edge.since(info.rcv_nxt))); + assert_eq!(info.srtt, handshake, "a downloader's SRTT keeps the handshake's"); + } + let arm = format!("timestamps {timestamps}, window {window}"); + let bound = 2 * (RATE * ESTIMATE + BURST) + 1448; + assert!(window <= bound, "{arm}: past {bound}"); + assert_eq!(net.apps[reader].received.1 - settled, 2_000 * RATE, "{arm}"); + } +} + +/// Node 0's SYN offering an MSS of 1,392: node 1 sends 1,380-byte segments, 68 less than node 0's +/// own send MSS of 1,448. +fn syn_offering_1392(from: usize, o: &O) -> Option> { + (from == 0 && o.flags & SYN != 0).then(|| { + let (value, echo) = o.ts.unwrap(); + let s = seg(o.seq).syn().wnd(o.wnd).mss(1392).sackok().ws(o.ws.unwrap()).ts(value, echo); + s.bytes(o.src, o.dst, 0, 0) + }) +} + +/// A download shaped like the T14's from a CDN: gigabit, a 15 ms round trip, timestamps and SACK, +/// and a peer whose segments are shorter than the downloader's own send MSS. Full-sized is what +/// the peer sends, so its echoes are sampled and the window grows to the whole buffer; measured +/// against the send MSS, none ever was, and the window stayed at 65,535. +#[test] +fn a_peer_sending_segments_shorter_than_our_send_mss_still_grows_the_window() { + let mut net = Net::buffered(15, BUFFER); + net.nodes[1].credit_per_ms = Some(125_000); + net.rewrite = Some(Box::new(syn_offering_1392)); + net.keep_wire = true; + net.connections(1, 443, [0, 32 * MIB]); + let peaks = watch(&mut net, 600_000); + assert!(net.finished(), "the transfer did not finish:\n{}", net.dump()); + net.assert_exact(); + let longest = net.wire.iter().filter(|(from, _)| *from == 1).map(|(_, o)| o.payload.len()).max(); + assert_eq!((longest, peaks[0].shifts), (Some(1380), Some((7, 7)))); + assert_eq!(peaks[0].window, BUFFER, "{:?}", peaks[0]); +} + +/// RFC 8985 §7.3: no probe leaves unless an RTT sample came since the last. Without timestamps a +/// retransmission voids the round's timing, so once the path's round trip passes twice the SRTT +/// a probe every round would leave before every ACK and SRTT would never move again. Here, after +/// a 20 ms handshake, the round trip is 100 ms, and node 0 writes ten whole segments a round, +/// 150 ms after the last was acknowledged, by when the D-SACK of the last probe has ended its +/// episode: the rounds without a probe take samples, SRTT climbs to the path's, and the probes +/// stop. +#[test] +fn rfc_8985_7_3_no_probe_without_an_rtt_sample_since_the_last() { + const TAIL: usize = 10 * 1460; + let mut net = Net::buffered(20, BUFFER); + net.auto_shut = false; + net.rewrite = Some(Box::new(syn_without_timestamps)); + net.connections(1, 80, [0, 0]); + assert!(net.run(100, |n| n.apps.len() == 2)); + net.delay = ns(50); + let (sender, receiver) = (net.apps.iter().position(|a| a.node == 0).unwrap(), net.apps.iter().position(|a| a.node == 1).unwrap()); + let mut probes = vec![0]; + for round in 0..40u8 { + let (now, id) = (net.instant(0), net.apps[sender].id); + assert_eq!(net.nodes[0].tcp.send(now, id, &[round; TAIL]), Ok(TAIL)); + net.apps[sender].sent.feed(&[round; TAIL]); + let want = net.apps[sender].sent.1; + for _ in 0..10_000 { + let info = net.nodes[0].tcp.info(id).unwrap(); + if net.apps[receiver].received.1 == want && info.snd_una == info.snd_nxt { + break; + } + net.advance(1); + } + assert_eq!(net.apps[receiver].received.1, want, "round {round}"); + net.advance(150); + probes.push(net.count(0, Counter::LossProbe)); + } + let srtt = net.nodes[0].tcp.info(net.apps[sender].id).unwrap().srtt.unwrap(); + let report = format!("probes after each round {probes:?}, SRTT {srtt:?}"); + assert!((90..=110).contains(&srtt.as_millis()), "{report}"); + assert_eq!(probes[40], probes[20], "{report}"); + assert!(probes.windows(3).all(|w| w[2] - w[0] <= 1), "two rounds in a row probed: {report}"); + assert_eq!((net.count(0, Counter::Rto), net.count(0, Counter::LossProbeRecovery)), (0, 0), "{report}"); +} diff --git a/toyos-net-shard/tcp/tests/receive.rs b/toyos-net-shard/tcp/tests/receive.rs index f2c3fdbea0d..b291672a377 100644 --- a/toyos-net-shard/tcp/tests/receive.rs +++ b/toyos-net-shard/tcp/tests/receive.rs @@ -281,3 +281,18 @@ fn s_rx_027_half_close_keeps_receiving() { check(&last, "ACK=15001"); assert_eq!(h.info().state, State::FinWait2); } + +/// At shift 7 a window below one unit reads as zero, and is offered as one unit only where the +/// buffer has that unit free: with 100 bytes free and a hole open, a whole unit would let the peer +/// send past the capacity, so the edge holds and the field stays 0. +#[test] +fn a_sub_unit_window_rounds_up_only_into_free_room() { + let mut h = client(4 << 20, seg(5000).ack(1001).syn().wnd(65_535).mss(1460).sackok().ws(7)); + assert_eq!(h.info().rcv_shift, 7); + h.input(1, seg(5001).ack(1001).len(65_435)); + let edge = h.info().rcv_edge; + assert_eq!(edge.since(h.info().rcv_nxt), 100); + expect(&h.input(2, seg(70_486).ack(1001).len(50)), &["ACK=70436 WND=0"]); + let info = h.info(); + assert_eq!((info.rcv_edge, info.ooo_ranges, info.unread), (edge, 1, 65_435)); +} diff --git a/toyos-net-shard/tcp/tests/recovery.rs b/toyos-net-shard/tcp/tests/recovery.rs index 7ec0a9b9e80..3832d41f5d6 100644 --- a/toyos-net-shard/tcp/tests/recovery.rs +++ b/toyos-net-shard/tcp/tests/recovery.rs @@ -268,12 +268,21 @@ fn s_lr_017_a_second_timeout() { assert_eq!((info.ssthresh, info.cwnd, info.rto), (10_220, 1460, ms(800))); } +/// EF's ten segments meet silence: two round trips and the slack on, a loss probe sends the last +/// one again (RFC 8985 §7.3), and the RTO runs from it. +fn probed_then_expired(h: &mut H) { + ten_out(h); + expect(&h.at(22), &["SEQ=14033 LEN=1448"]); + assert_eq!(h.count(Counter::LossProbe), 1); + nothing(&h.at(221)); + expect(&h.at(222), &["SEQ=1001"]); +} + #[test] fn s_lr_018_no_sack_recovery_before_the_timeout_point() { let mut h = fixture_ef(); - ten_out(&mut h); - h.at(200); - for (t, end) in [(210, 3897), (211, 5345), (212, 6793)] { + probed_then_expired(&mut h); + for (t, end) in [(232, 3897), (233, 5345), (234, 6793)] { sack_dup(&mut h, t, &[(2449, end)]); } assert!(!h.info().in_recovery); @@ -283,9 +292,8 @@ fn s_lr_018_no_sack_recovery_before_the_timeout_point() { #[test] fn s_lr_019_sacks_after_a_timeout_skip_held_ranges() { let mut h = fixture_ef(); - ten_out(&mut h); - expect(&h.at(200), &["SEQ=1001"]); - let outs = h.input_full(210, seg(5001).ack(2449).sack(&[(1001, 2449), (3897, 15_481)])); + probed_then_expired(&mut h); + let outs = h.input_full(232, seg(5001).ack(2449).sack(&[(1001, 2449), (3897, 15_481)])); expect(&outs, &["SEQ=2449 LEN=1448"]); assert_eq!(h.info().cwnd, 2896); assert_eq!(h.count(Counter::DsackRcvd), 1); @@ -331,3 +339,176 @@ fn s_lr_023_a_lost_fast_retransmission() { assert_eq!(info.ssthresh, flight * 7 / 10); assert_eq!(info.cwnd, 1460); } + +/// RFC 8985 §7.3: with data queued past cwnd and the peer's window open, the probe is a segment +/// of new data, outside cwnd; its ACK infers no loss. +#[test] +fn rfc_8985_7_3_the_probe_is_new_data_where_the_window_takes_a_segment() { + let mut h = fixture_ef(); + assert_eq!(h.send(0, 30_000).len(), 10); + nothing(&h.at(21)); + expect(&h.at(22), &["SEQ=15481 LEN=1448"]); + let cwnd = h.info().cwnd; + h.input_full(30, seg(5001).ack(16_929)); + assert_eq!((h.count(Counter::LossProbe), h.count(Counter::LossProbeRecovery), h.count(Counter::RetransmitBytes)), (1, 0, 0)); + assert!(h.info().cwnd >= cwnd); +} + +/// RFC 8985 §7.4.2: the probe sent the last segment again. The ACK that reaches its end without +/// a D-SACK leaves the episode open, since the original's ACK reads the same; the ACK past the +/// end with none says one copy was lost: cwnd is reduced as for a loss, once. +#[test] +fn rfc_8985_7_4_a_resent_probe_acknowledged_past_its_end_without_a_dsack_repaired_a_loss() { + let mut h = fixture_ef(); + ten_out(&mut h); + expect(&h.at(22), &["SEQ=14033 LEN=1448"]); + nothing(&h.send(25, 1448)); + expect(&h.input_full(30, seg(5001).ack(15_481)), &["SEQ=15481 LEN=1448"]); + assert_eq!(h.count(Counter::LossProbeRecovery), 0); + h.input_full(40, seg(5001).ack(16_929)); + let info = h.info(); + assert_eq!((info.ssthresh, info.cwnd, info.in_recovery), (2896, 2896, false)); + assert_eq!(h.count(Counter::LossProbeRecovery), 1); +} + +/// RFC 8985 §7.4.2: the ACK that reaches the resent probe's end carries no D-SACK, and the next +/// reports the probe as a duplicate: both copies arrived, nothing was lost, and cwnd stands through +/// the ACK of what is sent next. +#[test] +fn rfc_8985_7_4_a_dsack_after_the_ack_at_the_probes_end_infers_no_loss() { + let mut h = fixture_ef(); + ten_out(&mut h); + expect(&h.at(22), &["SEQ=14033 LEN=1448"]); + h.input_full(30, seg(5001).ack(15_481)); + let cwnd = h.info().cwnd; + h.input_full(31, seg(5001).ack(15_481).sack(&[(14_033, 15_481)])); + nothing(&h.send(32, 1448).into_iter().filter(|o| o.payload.is_empty()).collect::>()); + h.input_full(40, seg(5001).ack(16_929)); + assert_eq!((h.info().cwnd >= cwnd, h.count(Counter::LossProbeRecovery), h.count(Counter::DsackRcvd)), (true, 0, 1)); +} + +/// RFC 8985 §7.4: the same, but the ACK reports the probe's segment as a duplicate: nothing was +/// lost, and cwnd stands. +#[test] +fn rfc_8985_7_4_a_resent_probe_reported_as_a_duplicate_infers_no_loss() { + let mut h = fixture_ef(); + ten_out(&mut h); + expect(&h.at(22), &["SEQ=14033 LEN=1448"]); + let cwnd = h.info().cwnd; + h.input_full(30, seg(5001).ack(15_481).sack(&[(14_033, 15_481)])); + assert_eq!((h.info().cwnd >= cwnd, h.count(Counter::LossProbeRecovery), h.count(Counter::DsackRcvd)), (true, 0, 1)); +} + +/// RFC 8985 §7.2: with one segment out, the probe waits WCDelAckT past two round trips, which is +/// past the RTO here, so it goes at the RTO's time in the RTO's place, and the RTO runs from it. +#[test] +fn rfc_8985_7_2_with_one_segment_out_the_probe_stands_in_for_the_first_rto() { + let mut h = fixture_ef(); + h.send(0, 1448); + let rto = h.info().rto.as_millis() as i64; + nothing(&h.at(rto - 1)); + expect(&h.at(rto), &["SEQ=1001 LEN=1448"]); + assert_eq!((h.count(Counter::LossProbe), h.count(Counter::Rto)), (1, 0)); + nothing(&h.at(2 * rto - 1)); + expect(&h.at(2 * rto), &["SEQ=1001 LEN=1448"]); + assert_eq!(h.count(Counter::Rto), 1); +} + +/// RFC 8985 §7.2: no probe is scheduled in RTO recovery. After the probe at 22 ms and the RTO at +/// 222 ms, a cumulative ACK without SACK moves SND.UNA while go-back-N has the rest still to send: +/// no second probe follows it, however long the next ACK takes. +#[test] +fn rfc_8985_7_2_no_probe_in_rto_recovery() { + let mut h = fixture_ef(); + probed_then_expired(&mut h); + h.input_full(232, seg(5001).ack(2449)); + let pto = 2 * h.info().srtt.unwrap().as_millis() as i64 + 2; + let rto = h.info().rto.as_millis() as i64; + h.at(232 + pto + 1); + h.at(232 + rto - 1); + assert_eq!((h.count(Counter::LossProbe), h.count(Counter::Rto)), (1, 1)); +} + +/// RFC 8985 §7.1: entering fast recovery ends the probe's episode. The resent probe is still +/// outstanding when SACK recovery begins, and that recovery's own reduction is the only one: the +/// ACK at the probe's end and the one past it reduce nothing more. +#[test] +fn rfc_8985_7_1_fast_recovery_ends_the_probes_episode() { + let mut h = fixture_ef(); + ten_out(&mut h); + expect(&h.at(22), &["SEQ=14033 LEN=1448"]); + nothing(&sack_dup(&mut h, 30, &[(2449, 3897)])); + nothing(&sack_dup(&mut h, 31, &[(2449, 5345)])); + expect(&sack_dup(&mut h, 32, &[(2449, 6793)]), &["SEQ=1001 LEN=1448"]); + let ssthresh = h.info().ssthresh; + h.input_full(40, seg(5001).ack(15_481)); + assert!(!h.info().in_recovery); + expect(&h.send(41, 1448), &["SEQ=15481 LEN=1448"]); + h.input_full(50, seg(5001).ack(16_929)); + assert_eq!((h.info().ssthresh, h.count(Counter::SackRecovery), h.count(Counter::LossProbeRecovery)), (ssthresh, 1, 0)); +} + +/// A probe that came due while the next hop was not ready has not left when an ACK shuts the +/// window: persist takes over, and the last segment is not sent again into the shut window. +#[test] +fn a_probe_due_when_the_window_shuts_gives_way_to_persist() { + let mut h = fixture_ef(); + ten_out(&mut h); + h.hop = Box::new(|t, _| if t < 100 { toyos_net_tcp::Hop::Pending } else { toyos_net_tcp::Hop::Ready(()) }); + nothing(&h.at(22)); + let latest = h.log.iter().rev().find_map(|o| o.ts.map(|(v, _)| v)).unwrap(); + nothing(&h.at(30)); + h.deliver(seg(5001).ack(1001).wnd(0).ts(50_030, latest)); + assert_eq!(h.info().snd_wnd, 0); + h.tcp.wake(B); + nothing(&h.at(100)); + assert_eq!(h.count(Counter::LossProbe), 0); +} + +/// RFC 8985 §7.4.2: only a D-SACK matching the probe's end says the probe was a duplicate. The ACK +/// past the end reports an older segment as a duplicate, not the probe: one copy of the probe's +/// segment was still lost, and cwnd is reduced. +#[test] +fn rfc_8985_7_4_a_dsack_of_another_segment_still_infers_the_loss() { + let mut h = fixture_ef(); + ten_out(&mut h); + expect(&h.at(22), &["SEQ=14033 LEN=1448"]); + nothing(&h.send(25, 1448)); + expect(&h.input_full(30, seg(5001).ack(15_481)), &["SEQ=15481 LEN=1448"]); + h.input_full(40, seg(5001).ack(16_929).sack(&[(1001, 2449)])); + assert_eq!((h.count(Counter::DsackRcvd), h.count(Counter::LossProbeRecovery), h.info().ssthresh), (1, 1, 2896)); +} + +/// RFC 8985 §7.4.2 and §7.1 on one ACK: it passes the resent probe's end and SACKs three +/// segments above a hole, so it both infers the probe's loss and enters SACK recovery, and cwnd is +/// cut twice: the probe's cut, then recovery's from the flight after the ACK. Linux v6.12 does the +/// same: `tcp_process_tlp_ack` reduces and leaves CWR through `tcp_try_keep_open`, so +/// `tcp_enter_recovery` finds no reduction in progress and reduces again. +#[test] +fn rfc_8985_7_4_an_ack_that_infers_the_probes_loss_and_enters_recovery_cuts_twice() { + let mut h = fixture_ef(); + ten_out(&mut h); + expect(&h.at(22), &["SEQ=14033 LEN=1448"]); + nothing(&h.send(25, 5 * 1448)); + assert_eq!(h.input_full(30, seg(5001).ack(15_481)).len(), 5); + h.input_full(40, seg(5001).ack(16_929).sack(&[(18_377, 22_721)])); + let info = h.info(); + assert!(info.in_recovery); + assert_eq!((h.count(Counter::LossProbeRecovery), h.count(Counter::SackRecovery)), (1, 1)); + assert_eq!((info.ssthresh, info.cwnd), (4 * 1448 * 7 / 10, 4 * 1448 * 7 / 10)); +} + +/// RFC 8985 §7.4.2, Case 2: after the ACK at the resent probe's end, a duplicate ACK without SACK +/// says both copies arrived; the ACK past the end that follows infers nothing. +#[test] +fn rfc_8985_7_4_a_duplicate_without_sack_at_the_probes_end_infers_no_loss() { + let mut h = fixture_ef(); + ten_out(&mut h); + expect(&h.at(22), &["SEQ=14033 LEN=1448"]); + nothing(&h.send(25, 1448)); + expect(&h.input_full(30, seg(5001).ack(15_481)), &["SEQ=15481 LEN=1448"]); + let cwnd = h.info().cwnd; + nothing(&h.input_full(31, seg(5001).ack(15_481))); + h.input_full(40, seg(5001).ack(16_929)); + assert_eq!((h.info().cwnd >= cwnd, h.count(Counter::LossProbeRecovery)), (true, 0)); +} diff --git a/userland/netstack/src/main.rs b/userland/netstack/src/main.rs index 619aa93f2a1..568d89f0fae 100644 --- a/userland/netstack/src/main.rs +++ b/userland/netstack/src/main.rs @@ -81,26 +81,28 @@ pub const HOSTNAME: &str = "toyos-t14"; /// boot and a lease that lands later is applied like any other. const LEASE_BOUND: Duration = Duration::from_millis(toyos_tco::LEASE_BOUND_MS); -/// Payload bytes each direction of a connection buffers in the stack, before -/// the window closes on the peer or the client's pipe stops being read. -const TCP_BUFFER: u32 = 65_535; +/// The most payload each direction of a connection buffers in the stack, and +/// so the most a window offers: gigabit for a round trip up to 33 ms, at +/// window scale 7. Storage is taken as text is held, and the receive side +/// grows to this only as its reader keeps up (`toyos-net-tcp`'s `rx`). +const TCP_BUFFER: u32 = 4 << 20; /// A kernel pipe is one 2 MiB page (`kernel/src/pipe.rs`). The client /// allocates it, and netstack holding its far end is what keeps it alive. const PIPE_BYTES: u64 = 2 * 1024 * 1024; /// What one place can make this machine hold, at the largest of the three -/// things a place is (`toyos-net-node`'s `places`): a listener, whose peers -/// fill its queue of `LISTEN_READY` finished connections with a receive -/// buffer of text each, and its wake pipe. A stream is its two pipes and two -/// buffers, a datagram socket its two pipes and two queues, and both are -/// less. +/// things a place is (`toyos-net-node`'s `places`): a stream, its two pipes +/// and two buffers; a listener, whose peers fill its queue of `LISTEN_READY` +/// finished connections with an initial receive buffer of text each, nobody +/// reading them to grow one, and its wake pipe; a datagram socket, its two +/// pipes and two queues. const PLACE_BYTES: u64 = { - let listener = toyos_net_tcp::limits::LISTEN_READY as u64 * TCP_BUFFER as u64 + PIPE_BYTES; let stream = 2 * PIPE_BYTES + 2 * TCP_BUFFER as u64; + let listener = toyos_net_tcp::limits::LISTEN_READY as u64 * toyos_net_tcp::limits::RECEIVE_BUFFER_INITIAL as u64 + PIPE_BYTES; let datagram = 2 * PIPE_BYTES + (toyos_net_udp::limits::RX_BYTES + toyos_net_udp::limits::TX_BYTES) as u64; - assert!(listener >= stream && listener >= datagram); - listener + assert!(stream >= listener && stream >= datagram); + stream }; /// Share of physical memory netstack lets its clients' places tie up.