nxtquic-api 0.1.3

High-level async API for NxtQuic
Documentation
//! # NxtQuic Async API
//!
//! High-level async/await API for the NxtQuic QUIC transport library.
//! Provides ergonomic `Endpoint`, `Connection`, `SendStream`, and `RecvStream`
//! types built on top of the I/O-free `nxtquic-proto` core.
//!
//! This crate will be fully implemented in Phase 4 of the NxtQuic project.

#![forbid(unsafe_code)]
#![warn(missing_docs)]

pub mod connection;
pub mod endpoint;
pub mod runtime;
pub mod stream;

pub use connection::{Connection, ConnectionCloseInfo, ConnectionStats, HandshakeData};
pub use endpoint::{
    ClientConfig, ClientConfigBuilder, Connecting, Endpoint, EndpointConfig, EndpointStats,
    Incoming, ServerConfig, ServerConfigBuilder, TransportConfig,
};
pub use runtime::{Runtime, TokioRuntime};
pub use stream::{RecvStream, SendStream};

#[cfg(test)]
mod tests {
    use super::*;
    use std::net::SocketAddr;

    #[tokio::test]
    async fn bind_owns_a_real_udp_socket() {
        let server_addr: SocketAddr = "127.0.0.1:0".parse().unwrap();
        let server_endpoint = Endpoint::bind(server_addr).await.unwrap();
        let bound_addr = server_endpoint.local_addr().unwrap();

        let duplicate_bind = Endpoint::bind(bound_addr).await;
        assert!(duplicate_bind.is_err());
    }

    #[tokio::test]
    async fn endpoint_rejects_malformed_initial_headers() {
        let endpoint = Endpoint::bind("127.0.0.1:0".parse().unwrap())
            .await
            .unwrap();
        let sender = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
        let target = endpoint.local_addr().unwrap();

        sender.send_to(&[0x40, 0x00], target).await.unwrap();
        sender
            .send_to(&[0xc0, 0x00, 0x00, 0x00, 0x01], target)
            .await
            .unwrap();

        assert!(
            tokio::time::timeout(std::time::Duration::from_millis(100), endpoint.accept())
                .await
                .is_err()
        );
    }

    #[tokio::test]
    async fn stream_transfers_bytes_and_reports_eof_after_shutdown() {
        use tokio::io::{AsyncReadExt, AsyncWriteExt};

        let connection = Connection::new();
        let (mut send, mut recv) = connection.open_bi().await.unwrap();
        send.write_all(b"nxtquic").await.unwrap();
        send.shutdown().await.unwrap();

        let mut output = Vec::new();
        recv.read_to_end(&mut output).await.unwrap();
        assert_eq!(output, b"nxtquic");
    }

    #[tokio::test]
    async fn closed_connection_rejects_new_streams() {
        let connection = Connection::new();
        connection.close().await;
        assert!(connection.open_bi().await.is_err());
        assert!(connection.open_uni().await.is_err());
    }

    #[tokio::test]
    async fn network_connection_sends_datagrams_to_peer() {
        let peer = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
        let local = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
        let connection = Connection::from_udp(
            std::sync::Arc::new(local),
            peer.local_addr().unwrap(),
        );
        connection
            .send_datagram(bytes::Bytes::from_static(b"ping"))
            .await
            .unwrap();
        let mut buf = [0_u8; 16];
        let (len, _) = tokio::time::timeout(
            std::time::Duration::from_secs(1),
            peer.recv_from(&mut buf),
        )
        .await
        .unwrap()
        .unwrap();
        assert_eq!(&buf[..len], b"ping");
    }

    #[tokio::test]
    async fn network_connection_receives_peer_datagrams() {
        let peer = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
        let local = tokio::net::UdpSocket::bind("127.0.0.1:0").await.unwrap();
        let local_addr = local.local_addr().unwrap();
        let connection = Connection::from_udp(std::sync::Arc::new(local), peer.local_addr().unwrap());
        peer.send_to(b"pong", local_addr).await.unwrap();
        assert_eq!(connection.recv_datagram().await.unwrap(), bytes::Bytes::from_static(b"pong"));
    }

    #[tokio::test]
    async fn datagram_frames_are_exposed_by_connection_api() {
        use nxtquic_proto::frame::{DatagramFrame, Frame};
        let connection = Connection::new();
        let mut payload = Vec::new();
        Frame::Datagram(DatagramFrame { has_length: true, data: bytes::Bytes::from_static(b"frame") }).encode(&mut payload);
        connection.ingest_frames(&payload).await.unwrap();
        assert_eq!(connection.recv_datagram().await.unwrap(), bytes::Bytes::from_static(b"frame"));
    }

    #[tokio::test]
    async fn stream_frames_are_delivered_to_the_matching_queue() {
        use nxtquic_proto::{frame::{Frame, StreamFrame}, shared::StreamId, varint::VarInt};
        use tokio::io::AsyncReadExt;
        let connection = Connection::new();
        let mut payload = Vec::new();
        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);
        connection.ingest_frames(&payload).await.unwrap();
        let mut recv = connection.accept_bi().await.unwrap().1;
        let mut output = Vec::new();
        recv.read_to_end(&mut output).await.unwrap();
        assert_eq!(output, b"x");
    }

    #[tokio::test]
    async fn stream_frames_are_reassembled_into_accept_bi() {
        use nxtquic_proto::frame::{Frame, StreamFrame};
        use nxtquic_proto::{shared::StreamId, varint::VarInt};
        use tokio::io::AsyncReadExt;
        let connection = Connection::new();
        let mut payload = Vec::new();
        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);
        connection.ingest_frames(&payload).await.unwrap();
        let (_send, mut recv) = connection.accept_bi().await.unwrap();
        let mut out = Vec::new();
        recv.read_to_end(&mut out).await.unwrap();
        assert_eq!(out, b"hello");
    }

    #[tokio::test]
    async fn client_unidirectional_streams_are_delivered_to_accept_uni() {
        use nxtquic_proto::{frame::{Frame, StreamFrame}, shared::StreamId, varint::VarInt};
        use tokio::io::AsyncReadExt;

        let connection = Connection::new();
        let mut payload = Vec::new();
        Frame::Stream(StreamFrame {
            // Stream ID 2: client-initiated, unidirectional.
            stream_id: StreamId::from_u64(2),
            offset: VarInt::ZERO,
            length: Some(VarInt::from_u32(5)),
            fin: true,
            data: bytes::Bytes::from_static(b"hello"),
        })
        .encode(&mut payload);
        connection.ingest_frames(&payload).await.unwrap();

        let mut recv = connection.accept_uni().await.unwrap();
        let mut out = Vec::new();
        recv.read_to_end(&mut out).await.unwrap();
        assert_eq!(out, b"hello");
    }

    #[tokio::test]
    async fn closing_connection_wakes_pending_stream_accept() {
        let connection = std::sync::Arc::new(Connection::new());
        let waiter = {
            let connection = std::sync::Arc::clone(&connection);
            tokio::spawn(async move { connection.accept_uni().await })
        };
        tokio::task::yield_now().await;
        connection.close().await;
        let result = waiter.await.unwrap();
        assert!(result.is_err());
    }
}