runsync-transfer 2026.1.0

High-throughput P2P file transfer engine: adaptive compression, end-to-end AEAD, parallel chunked pipeline over QUIC or any async transport.
Documentation
//! In-process transport over `tokio::io::duplex`.
//!
//! Exists so the protocol, codec, and I/O pipeline can be tested at full speed
//! without a socket, a certificate, or a port. The duplex buffers are bounded,
//! so backpressure behaves the same way it does over a real stream — which is
//! the property the large-file tests actually depend on.

use super::{BoxRecv, BoxSend, Transport};
use crate::error::{Error, Result};
use tokio::io::DuplexStream;
use tokio::sync::{mpsc, Mutex};

type BiPair = (DuplexStream, DuplexStream);

pub struct MemTransport {
    label: String,
    buf: usize,
    uni_out: mpsc::UnboundedSender<DuplexStream>,
    uni_in: Mutex<mpsc::UnboundedReceiver<DuplexStream>>,
    bi_out: mpsc::UnboundedSender<BiPair>,
    bi_in: Mutex<mpsc::UnboundedReceiver<BiPair>>,
}

/// Create two connected endpoints. `buf` is the per-stream buffer in bytes and
/// sets how much data may sit in flight before the writer blocks.
pub fn pair(buf: usize) -> (MemTransport, MemTransport) {
    let (a_uni_tx, a_uni_rx) = mpsc::unbounded_channel();
    let (b_uni_tx, b_uni_rx) = mpsc::unbounded_channel();
    let (a_bi_tx, a_bi_rx) = mpsc::unbounded_channel();
    let (b_bi_tx, b_bi_rx) = mpsc::unbounded_channel();

    let a = MemTransport {
        label: "mem:a".into(),
        buf,
        uni_out: b_uni_tx,
        uni_in: Mutex::new(a_uni_rx),
        bi_out: b_bi_tx,
        bi_in: Mutex::new(a_bi_rx),
    };
    let b = MemTransport {
        label: "mem:b".into(),
        buf,
        uni_out: a_uni_tx,
        uni_in: Mutex::new(b_uni_rx),
        bi_out: a_bi_tx,
        bi_in: Mutex::new(b_bi_rx),
    };
    (a, b)
}

#[async_trait::async_trait]
impl Transport for MemTransport {
    async fn open_uni(&self) -> Result<BoxSend> {
        let (mine, theirs) = tokio::io::duplex(self.buf);
        self.uni_out
            .send(theirs)
            .map_err(|_| Error::Closed("peer endpoint dropped".into()))?;
        Ok(Box::new(mine))
    }

    async fn accept_uni(&self) -> Result<BoxRecv> {
        let mut rx = self.uni_in.lock().await;
        match rx.recv().await {
            Some(s) => Ok(Box::new(s)),
            None => Err(Error::Closed("peer endpoint dropped".into())),
        }
    }

    async fn open_bi(&self) -> Result<(BoxSend, BoxRecv)> {
        // Two independent duplexes so each direction has its own buffer and a
        // stalled reader cannot deadlock the writer in the other direction.
        let (my_w, their_r) = tokio::io::duplex(self.buf);
        let (their_w, my_r) = tokio::io::duplex(self.buf);
        self.bi_out
            .send((their_w, their_r))
            .map_err(|_| Error::Closed("peer endpoint dropped".into()))?;
        Ok((Box::new(my_w), Box::new(my_r)))
    }

    async fn accept_bi(&self) -> Result<(BoxSend, BoxRecv)> {
        let mut rx = self.bi_in.lock().await;
        match rx.recv().await {
            Some((w, r)) => Ok((Box::new(w), Box::new(r))),
            None => Err(Error::Closed("peer endpoint dropped".into())),
        }
    }

    fn close(&self, _code: u32, _reason: &[u8]) {}

    fn peer_label(&self) -> String {
        self.label.clone()
    }
}