Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
11 changes: 8 additions & 3 deletions src/stream/tcb.rs
Original file line number Diff line number Diff line change
Expand Up @@ -453,10 +453,15 @@ impl Tcb {
self.inflight_packets.values().collect::<Vec<_>>()
}

/// Bytes that can still be sent before reaching the right edge of the peer window or the cap on
/// unacknowledged bytes, whichever is nearer. Simplified version: min(cwnd, rwnd) - in flight.
pub(super) fn get_usable_send_window(&self) -> u32 {
let in_flight = self.seq.distance(self.get_last_received_ack());
self.max_unacked_bytes.min(self.get_send_window()).saturating_sub(in_flight)
}

pub fn is_send_buffer_full(&self) -> bool {
// To respect the receiver's window (remote_window) size and avoid sending too many unacknowledged packets, which may cause packet loss
// Simplified version: min(cwnd, rwnd)
self.seq.distance(self.get_last_received_ack()) >= self.max_unacked_bytes.min(self.get_send_window())
self.get_usable_send_window() == 0
}
}

Expand Down
26 changes: 26 additions & 0 deletions src/stream/tcp.rs
Original file line number Diff line number Diff line change
Expand Up @@ -391,6 +391,8 @@ impl AsyncWrite for IpStackTcpStream {
return Poll::Pending;
}

// A segment past the right edge of the peer window is trimmed by the peer (RFC 9293 §3.10.7.4).
let buf = &buf[..buf.len().min(tcb.get_usable_send_window() as usize)];
let sender = &self.up_packet_sender;
let payload_len = write_packet_to_device(sender, nt, &tcb, None, ACK | PSH, None, Some(buf.to_vec()))?;
let was_empty = tcb.is_inflight_queue_empty();
Expand Down Expand Up @@ -1470,4 +1472,28 @@ mod tests {
assert!(matches!(Pin::new(&mut stream).poll_write(&mut cx, b"sent"), Poll::Ready(Ok(4))));
assert_eq!(up_rx.recv().await.unwrap().payload.as_deref(), Some(&b"sent"[..]));
}

/// A write that would pass the right edge of the peer window is cut at the edge, and the next
/// one waits for an ACK instead of sending past it.
#[tokio::test]
async fn writer_stops_at_the_right_edge_of_the_peer_window() {
let (mut stream, syn_ack, mut up_rx) = open(2000, &[]).await;
let peer_seq = SeqNum(syn_ack.acknowledgment_number);
feed(&stream, &syn_ack, peer_seq, 2000, &[1; 4]);
up_rx.recv().await.unwrap();

let (_, waker) = recording_waker();
let mut cx = Context::from_waker(&waker);
let chunk = [7; 1460];
assert!(matches!(Pin::new(&mut stream).poll_write(&mut cx, &chunk), Poll::Ready(Ok(1460))));
assert_eq!(up_rx.recv().await.unwrap().payload.map(|p| p.len()), Some(1460));

assert!(
matches!(Pin::new(&mut stream).poll_write(&mut cx, &chunk), Poll::Ready(Ok(540))),
"the second write is not cut at the right edge of the 2000-byte window"
);
assert_eq!(up_rx.recv().await.unwrap().payload.map(|p| p.len()), Some(540));

assert!(Pin::new(&mut stream).poll_write(&mut cx, &chunk).is_pending());
}
}
Loading