Skip to main content

nxtquic_api/
lib.rs

1//! # NxtQuic Async API
2//!
3//! High-level async/await API for the NxtQuic QUIC transport library.
4//! Provides ergonomic `Endpoint`, `Connection`, `SendStream`, and `RecvStream`
5//! types built on top of the I/O-free `nxtquic-proto` core.
6//!
7//! This crate will be fully implemented in Phase 4 of the NxtQuic project.
8
9#![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 2: client-initiated, unidirectional.
164            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}