如何将 tokio TcpStream 转换为可序列化/反序列化值的接收器/流?

The*_*vat 1 rust serde rust-tokio

我有一个东京TcpStream。我想T通过这个流传递一些类型。这种类型的T实现Serialize和Deserialize. 我怎样才能获得 aSink<T>和 a Stream<T>?

我找到了板条箱tokio_util和tokio_serde,但我不知道如何使用它们来做我想做的事情。

Ian*_* S. 8

我不知道您的代码结构或您计划使用的编解码器,但我已经弄清楚如何将所有内容粘合在一起形成一个可行的示例。

您的Sink<T>和Stream<Item=T>将会由Framed中的类型提供tokio-serde。该层处理通过 传递消息serde。此类型采用四个通用参数:Transport、Item(流项)、SinkItem和Codec。Codec是您要使用的特定序列化器和反序列化器的包装器。您可以在此处查看提供的选项。ItemandSinkItem将成为您必须实现Serializeand的消息类型Deserialize。Transport需要是 aSink<SinkItem>和/或Stream<Item=Item>本身,以便框架实现任何有用的特征。这就是tokio-util进来的地方。它提供了各种类型Framed*,允许您将实现AsyncRead/的AsyncWrite东西分别转换为流和接收器。为了构造这些帧,您需要指定一个编解码器来将帧与线路分隔开。为了简单起见,在我的示例中我只使用了LengthDelimitedCodec,但还提供了其他选项。

言归正传,这里有一个示例,说明如何将 atokio::net::TcpStream拆分为 anSink<T>和Stream<Item=T>。请注意,这T是流端的结果,因为如果消息格式错误,serde 层可能会失败。

use futures::{SinkExt, StreamExt};
use serde::{Deserialize, Serialize};
use tokio::net::{
    tcp::{OwnedReadHalf, OwnedWriteHalf},
    TcpListener,
    TcpStream,
};
use tokio_serde::{formats::Json, Framed};
use tokio_util::codec::{FramedRead, FramedWrite, LengthDelimitedCodec};

#[derive(Serialize, Deserialize, Debug)]
struct MyMessage {
    field: String,
}

type WrappedStream = FramedRead<OwnedReadHalf, LengthDelimitedCodec>;
type WrappedSink = FramedWrite<OwnedWriteHalf, LengthDelimitedCodec>;

// We use the unit type in place of the message types since we're
// only dealing with one half of the IO
type SerStream = Framed<WrappedStream, MyMessage, (), Json<MyMessage, ()>>;
type DeSink = Framed<WrappedSink, (), MyMessage, Json<(), MyMessage>>;

fn wrap_stream(stream: TcpStream) -> (SerStream, DeSink) {
    let (read, write) = stream.into_split();
    let stream = WrappedStream::new(read, LengthDelimitedCodec::new());
    let sink = WrappedSink::new(write, LengthDelimitedCodec::new());
    (
        SerStream::new(stream, Json::default()),
        DeSink::new(sink, Json::default()),
    )
}

#[tokio::main]
async fn main() {
    let listener = TcpListener::bind("0.0.0.0:8080")
        .await
        .expect("Failed to bind server to addr");

    tokio::task::spawn(async move {
        let (stream, _) = listener
            .accept()
            .await
            .expect("Failed to accept incoming connection");
        
        let (mut stream, mut sink) = wrap_stream(stream);

        println!(
            "Server received: {:?}",
            stream
                .next()
                .await
                .expect("No data in stream")
                .expect("Failed to parse ping")
        );

        sink.send(MyMessage {
            field: "pong".to_owned(),
        })
            .await
            .expect("Failed to send pong");
    });

    let stream = TcpStream::connect("127.0.0.1:8080")
        .await
        .expect("Failed to connect to server");

    let (mut stream, mut sink) = wrap_stream(stream);

    sink.send(MyMessage {
        field: "ping".to_owned(),
    })
        .await
        .expect("Failed to send ping to server");
        
    println!(
        "Client received: {:?}",
        stream
            .next()
            .await
            .expect("No data in stream")
            .expect("Failed to parse pong")
    );
}
Run Code Online (Sandbox Code Playgroud)

运行这个例子会产生:

Server received: MyMessage { field: "ping" }
Client received: MyMessage { field: "pong" }
Run Code Online (Sandbox Code Playgroud)

请注意,并不要求您拆分流。您可以改为tokio_util::codec::Framed从构造 a TcpStream,并tokio_serde::Framed用 a 构造 a tokio_serde::formats::SymmetricalJson<MyMessage>,然后相应地Framed实现Sink和Stream。此外,此示例中的许多功能都是功能门控的,因此请务必根据文档启用适当的功能。