diff --git a/src/stream/tcb.rs b/src/stream/tcb.rs index 9d322a2..208668d 100644 --- a/src/stream/tcb.rs +++ b/src/stream/tcb.rs @@ -453,10 +453,15 @@ impl Tcb { self.inflight_packets.values().collect::>() } + /// 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 } } diff --git a/src/stream/tcp.rs b/src/stream/tcp.rs index 28e06b2..4af9895 100644 --- a/src/stream/tcp.rs +++ b/src/stream/tcp.rs @@ -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(); @@ -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()); + } }