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;
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 2: client-initiated, unidirectional.
162            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}