dope 0.2.3

Thin io_uring adaptor with "Manifolds"
Documentation
use dope_futures::AsyncRead;
use dope_futures::listener::Listener;
use dope_futures::stream::Stream;
use dope_runtime::{DefaultBackend, block_on, spawn};
use dope_transport::Tcp;

#[test]
fn echo_basic() {
    let mut listener = Listener::<Tcp, DefaultBackend>::bind("127.0.0.1:0").unwrap();
    let addr = listener.local_addr().unwrap();

    block_on(async move {
        let server = spawn(async move {
            let (mut tcp, _) = listener.accept().await.unwrap();
            let mut buf = [0u8; 12];
            tcp.read_exact(&mut buf).await.unwrap();
            tcp.write_all(&buf).await.unwrap();
        });

        let mut tcp = Stream::<Tcp, DefaultBackend>::connect_with(addr)
            .await
            .unwrap();
        let msg = b"hello-dope!!";
        tcp.write_all(msg).await.unwrap();
        let mut out = [0u8; 12];
        tcp.read_exact(&mut out).await.unwrap();
        assert_eq!(&out, msg);

        server.await;
    });
}

#[test]
fn echo_large_payload() {
    const SIZE: usize = 128 * 1024;

    let mut listener = Listener::<Tcp, DefaultBackend>::bind("127.0.0.1:0").unwrap();
    let addr = listener.local_addr().unwrap();

    block_on(async move {
        let server = spawn(async move {
            let (mut tcp, _) = listener.accept().await.unwrap();
            let mut received = Vec::with_capacity(SIZE);
            while received.len() < SIZE {
                let want = (SIZE - received.len()).min(16384);
                let (res, buf) = tcp.read(vec![0u8; want]).await;
                let n = res.unwrap();
                if n == 0 {
                    break;
                }
                received.extend_from_slice(&buf[..n]);
            }
            received
        });

        let mut tcp = Stream::<Tcp, DefaultBackend>::connect_with(addr)
            .await
            .unwrap();
        let payload: Vec<u8> = (0..SIZE).map(|i| (i % 251) as u8).collect();
        let (res, _) = tcp.write_all_owned(payload.clone()).await;
        res.unwrap();

        let received = server.await;
        assert_eq!(received.len(), SIZE);
        assert_eq!(received, payload);
    });
}

#[test]
fn echo_concurrent_connections() {
    const CONNS: usize = 16;
    const MSG_LEN: usize = 32;

    let mut listener = Listener::<Tcp, DefaultBackend>::bind("127.0.0.1:0").unwrap();
    let addr = listener.local_addr().unwrap();

    block_on(async move {
        let server = spawn(async move {
            for _ in 0..CONNS {
                let (mut tcp, _) = listener.accept().await.unwrap();
                spawn(async move {
                    let mut buf = [0u8; MSG_LEN];
                    tcp.read_exact(&mut buf).await.unwrap();
                    tcp.write_all(&buf).await.unwrap();
                });
            }
        });

        for i in 0..CONNS {
            let mut tcp = Stream::<Tcp, DefaultBackend>::connect_with(addr)
                .await
                .unwrap();
            let mut msg = [0u8; MSG_LEN];
            for (j, b) in msg.iter_mut().enumerate() {
                *b = ((i * 17 + j) % 251) as u8;
            }
            tcp.write_all(&msg).await.unwrap();
            let mut out = [0u8; MSG_LEN];
            tcp.read_exact(&mut out).await.unwrap();
            assert_eq!(out, msg);
        }

        server.await;
    });
}