Skip to content
Merged
116 changes: 108 additions & 8 deletions toyos-net-shard/tcp/src/conn.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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)]
Expand Down Expand Up @@ -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<Instant>,
probe_due: bool,
probe: Option<(Seq, bool)>,
sampled: bool,
persist: Option<Persist>,
sws: Option<Instant>,
sws_fired: bool,
Expand Down Expand Up @@ -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),
Expand All @@ -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,
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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 {
Expand Down Expand Up @@ -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
Expand All @@ -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),
}
Expand All @@ -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;
}
}
}
Expand Down Expand Up @@ -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);
Expand Down Expand Up @@ -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) {
Expand Down Expand Up @@ -851,10 +912,10 @@ impl Sync {
Ok(n)
}

pub fn recv(&mut self, out: &mut [u8]) -> Result<Received, Error> {
pub fn recv(&mut self, out: &mut [u8], now: Instant) -> Result<Received, Error> {
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)
Expand All @@ -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<Received, Error> {
pub fn recv_with(&mut self, take: impl FnOnce(&[u8]) -> usize, now: Instant) -> Result<Received, Error> {
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))
}
Expand Down Expand Up @@ -935,6 +996,7 @@ impl Sync {
pub fn deadline(&self, ctx: &Ctx<'_>) -> Option<Instant> {
[
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),
Expand All @@ -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);
}
Expand Down Expand Up @@ -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;
Expand Down Expand Up @@ -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);
}
Expand Down Expand Up @@ -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();
}
Expand Down Expand Up @@ -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<T>(&mut self, ctx: &mut Ctx<'_>, exit: &mut dyn Exit<T>) -> Result<bool, NotReady> {
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<T>(&mut self, ctx: &mut Ctx<'_>, exit: &mut dyn Exit<T>) -> Result<bool, NotReady> {
let now = ctx.now;
Expand Down
2 changes: 2 additions & 0 deletions toyos-net-shard/tcp/src/counters.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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";
}
Expand Down
5 changes: 5 additions & 0 deletions toyos-net-shard/tcp/src/lib.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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,
}
Expand Down Expand Up @@ -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.
Expand Down
8 changes: 6 additions & 2 deletions toyos-net-shard/tcp/src/props.rs
Original file line number Diff line number Diff line change
Expand Up @@ -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);
Expand All @@ -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:?}");
Expand Down Expand Up @@ -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);
Expand Down
34 changes: 26 additions & 8 deletions toyos-net-shard/tcp/src/ring.rs
Original file line number Diff line number Diff line change
@@ -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;

Expand Down Expand Up @@ -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));
Expand Down Expand Up @@ -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)
Expand Down
Loading
Loading