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