1#![forbid(unsafe_code)]
10#![warn(missing_docs)]
11
12pub mod connection;
13pub mod endpoint;
14pub mod runtime;
15pub mod stream;
16
17pub use connection::Connection;
18pub use endpoint::{
19 ClientConfig, Connecting, Endpoint, EndpointConfig, Incoming, ServerConfig, TransportConfig,
20};
21pub use stream::{RecvStream, SendStream};
22
23#[cfg(test)]
24mod tests {
25 use super::*;
26 use std::net::SocketAddr;
27
28 #[tokio::test]
29 async fn bind_owns_a_real_udp_socket() {
30 let server_addr: SocketAddr = "127.0.0.1:0".parse().unwrap();
31 let server_endpoint = Endpoint::bind(server_addr).await.unwrap();
32 let bound_addr = server_endpoint.local_addr().unwrap();
33
34 let duplicate_bind = Endpoint::bind(bound_addr).await;
35 assert!(duplicate_bind.is_err());
36 }
37
38 #[tokio::test]
39 async fn endpoint_rejects_malformed_initial_headers() {
40 let endpoint = Endpoint::bind("127.0.0.1:0".parse().unwrap())
41 .await
42 .unwrap();
43 let sender = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
44 let target = endpoint.local_addr().unwrap();
45
46 sender.send_to(&[0x40, 0x00], target).await.unwrap();
47 sender
48 .send_to(&[0xc0, 0x00, 0x00, 0x00, 0x01], target)
49 .await
50 .unwrap();
51
52 assert!(
53 tokio::time::timeout(std::time::Duration::from_millis(100), endpoint.accept())
54 .await
55 .is_err()
56 );
57 }
58
59 #[tokio::test]
60 async fn stream_transfers_bytes_and_reports_eof_after_shutdown() {
61 use tokio::io::{AsyncReadExt, AsyncWriteExt};
62
63 let connection = Connection::new();
64 let (mut send, mut recv) = connection.open_bi().await.unwrap();
65 send.write_all(b"nxtquic").await.unwrap();
66 send.shutdown().await.unwrap();
67
68 let mut output = Vec::new();
69 recv.read_to_end(&mut output).await.unwrap();
70 assert_eq!(output, b"nxtquic");
71 }
72
73 #[tokio::test]
74 async fn closed_connection_rejects_new_streams() {
75 let connection = Connection::new();
76 connection.close().await;
77 assert!(connection.open_bi().await.is_err());
78 assert!(connection.open_uni().await.is_err());
79 }
80
81 #[tokio::test]
82 async fn network_connection_sends_datagrams_to_peer() {
83 let peer = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
84 let local = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
85 let connection = Connection::from_udp(
86 std::sync::Arc::new(local),
87 peer.local_addr().unwrap(),
88 );
89 connection
90 .send_datagram(bytes::Bytes::from_static(b"ping"))
91 .await
92 .unwrap();
93 let mut buf = [0_u8; 16];
94 let (len, _) = tokio::time::timeout(
95 std::time::Duration::from_secs(1),
96 peer.recv_from(&mut buf),
97 )
98 .await
99 .unwrap()
100 .unwrap();
101 assert_eq!(&buf[..len], b"ping");
102 }
103
104 #[tokio::test]
105 async fn network_connection_receives_peer_datagrams() {
106 let peer = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
107 let local = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
108 let local_addr = local.local_addr().unwrap();
109 let connection = Connection::from_udp(std::sync::Arc::new(local), peer.local_addr().unwrap());
110 peer.send_to(b"pong", local_addr).await.unwrap();
111 assert_eq!(connection.recv_datagram().await.unwrap(), bytes::Bytes::from_static(b"pong"));
112 }
113
114 #[tokio::test]
115 async fn datagram_frames_are_exposed_by_connection_api() {
116 use nxtquic_proto::frame::{DatagramFrame, Frame};
117 let connection = Connection::new();
118 let mut payload = Vec::new();
119 Frame::Datagram(DatagramFrame { has_length: true, data: bytes::Bytes::from_static(b"frame") }).encode(&mut payload);
120 connection.ingest_frames(&payload).await.unwrap();
121 assert_eq!(connection.recv_datagram().await.unwrap(), bytes::Bytes::from_static(b"frame"));
122 }
123
124 #[tokio::test]
125 async fn stream_frames_are_delivered_to_the_matching_queue() {
126 use nxtquic_proto::{frame::{Frame, StreamFrame}, shared::StreamId, varint::VarInt};
127 use tokio::io::AsyncReadExt;
128 let connection = Connection::new();
129 let mut payload = Vec::new();
130 Frame::Stream(StreamFrame { stream_id: StreamId::from_u64(0), offset: VarInt::ZERO, length: Some(VarInt::from_u32(1)), fin: true, data: bytes::Bytes::from_static(b"x") }).encode(&mut payload);
131 connection.ingest_frames(&payload).await.unwrap();
132 let mut recv = connection.accept_bi().await.unwrap().1;
133 let mut output = Vec::new();
134 recv.read_to_end(&mut output).await.unwrap();
135 assert_eq!(output, b"x");
136 }
137
138 #[tokio::test]
139 async fn stream_frames_are_reassembled_into_accept_bi() {
140 use nxtquic_proto::frame::{Frame, StreamFrame};
141 use nxtquic_proto::{shared::StreamId, varint::VarInt};
142 use tokio::io::AsyncReadExt;
143 let connection = Connection::new();
144 let mut payload = Vec::new();
145 Frame::Stream(StreamFrame { stream_id: StreamId::from_u64(0), offset: VarInt::ZERO, length: Some(VarInt::from_u32(5)), fin: true, data: bytes::Bytes::from_static(b"hello") }).encode(&mut payload);
146 connection.ingest_frames(&payload).await.unwrap();
147 let (_send, mut recv) = connection.accept_bi().await.unwrap();
148 let mut out = Vec::new();
149 recv.read_to_end(&mut out).await.unwrap();
150 assert_eq!(out, b"hello");
151 }
152
153 #[tokio::test]
154 async fn client_unidirectional_streams_are_delivered_to_accept_uni() {
155 use nxtquic_proto::{frame::{Frame, StreamFrame}, shared::StreamId, varint::VarInt};
156 use tokio::io::AsyncReadExt;
157
158 let connection = Connection::new();
159 let mut payload = Vec::new();
160 Frame::Stream(StreamFrame {
161 stream_id: StreamId::from_u64(2),
163 offset: VarInt::ZERO,
164 length: Some(VarInt::from_u32(5)),
165 fin: true,
166 data: bytes::Bytes::from_static(b"hello"),
167 })
168 .encode(&mut payload);
169 connection.ingest_frames(&payload).await.unwrap();
170
171 let mut recv = connection.accept_uni().await.unwrap();
172 let mut out = Vec::new();
173 recv.read_to_end(&mut out).await.unwrap();
174 assert_eq!(out, b"hello");
175 }
176
177 #[tokio::test]
178 async fn closing_connection_wakes_pending_stream_accept() {
179 let connection = std::sync::Arc::new(Connection::new());
180 let waiter = {
181 let connection = std::sync::Arc::clone(&connection);
182 tokio::spawn(async move { connection.accept_uni().await })
183 };
184 tokio::task::yield_now().await;
185 connection.close().await;
186 let result = waiter.await.unwrap();
187 assert!(result.is_err());
188 }
189}