#![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: 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());
}
}