rifts 0.2.0

Rift Realtime Protocol / 1.0 — server-side implementation
Documentation
//! Ntex WebSocket adapter.
//!
//! Wraps ntex WebSocket (`ntex::ws::WsSink` + message stream) as a
//! Rift `TransportConnection` via a channel bridge.

use std::net::SocketAddr;

use crate::transport::TransportConnection;
use crate::transport::bridge::spawn_bridge_local;

/// Wrap an ntex WebSocket pair into a `TransportConnection`.
///
/// `stream` is the message stream from `ntex::web::ws::start()`.
/// `sink` is the `ntex::ws::WsSink` for outgoing messages.
///
/// ```ignore
/// use ntex::web;
/// use ntex::ws;
///
/// async fn handler(req: web::HttpRequest, stream: web::Payload)
///     -> Result<web::HttpResponse, web::Error>
/// {
///     ws::start(req, stream, |msg_stream, sink| async move {
///         let conn = rift::transport::ntex::into_connection(sink, msg_stream, None);
///         tokio::spawn(async move {
///             rift_server.accept_and_spawn(conn);
///         });
///     })
/// }
/// ```
pub fn into_connection<S, E>(
    sink: ntex::ws::WsSink,
    mut stream: S,
    peer: Option<SocketAddr>,
) -> Box<dyn TransportConnection>
where
    S: futures_util::Stream<Item = Result<ntex::ws::Message, E>> + Unpin + 'static,
    E: std::fmt::Debug,
{
    spawn_bridge_local(
        peer,
        256,
        // Reader: pull from ntex message stream → tokio channel.
        move |tx| {
            ntex::rt::spawn(async move {
                use futures_util::StreamExt;
                while let Some(msg) = stream.next().await {
                    let raw = match msg {
                        Ok(ntex::ws::Message::Binary(bin)) => {
                            let mut v = Vec::with_capacity(1 + bin.len());
                            v.push(b'B');
                            v.extend_from_slice(&bin);
                            v
                        }
                        Ok(ntex::ws::Message::Text(text)) => {
                            let mut v = Vec::with_capacity(1 + text.len());
                            v.push(b'T');
                            v.extend_from_slice(text.as_bytes());
                            v
                        }
                        Ok(ntex::ws::Message::Close(_)) => vec![b'C'],
                        Ok(_) => continue,
                        Err(_) => break,
                    };
                    if tx.send(raw).await.is_err() {
                        break;
                    }
                }
            });
        },
        // Writer: receive from tokio channel → ntex WsSink.
        move |mut rx| {
            ntex::rt::spawn(async move {
                while let Some(raw) = rx.recv().await {
                    if raw.first() == Some(&b'C') {
                        let _ = sink.send(ntex::ws::Message::Close(None)).await;
                        break;
                    }
                    if raw.len() > 1 {
                        let _ = sink
                            .send(ntex::ws::Message::Binary(raw[1..].to_vec().into()))
                            .await;
                    }
                }
            });
        },
    )
}