【问题标题】:impl AsyncRead for tonic::Streamingimpl AsyncRead for tonic::Streaming
【发布时间】:2021-03-04 20:42:03
【问题描述】:

我正在尝试服用补品routeguide tutorial,并将客户端变成rocket 服务器。我只是接受响应并将 gRPC 转换为字符串。

service RouteGuide {
    rpc GetFeature(Point) returns (Feature) {}
    rpc ListFeatures(Rectangle) returns (stream Feature) {}
}

这对于 GetFeature 来说已经足够好了。对于 ListFeatures 查询,就像 Tonic 允许客户端在响应中使用流一样,我想将其传递给 Rocket 客户端。我看到 Rocket 支持 streaming 响应,但我需要实现 AsyncRead 特征。

有没有办法做这样的事情?以下是关于我正在做的事情的精简版:

struct FeatureStream {
    stream: tonic::Streaming<Feature>,
}

impl AsyncRead for FeatureStream {
    fn poll_read(
        self: Pin<&mut Self>,
        cx: &mut Context<'_>,
        buf: &mut ReadBuf<'_>,
    ) -> Poll<std::io::Result<()>> {
        // Write out as utf8 any response messages.
        match Pin::new(&mut self.stream.message()).poll(cx) {
            Poll::Pending => Poll::Pending,
            Poll::Ready(feature) => Poll::Pending,
        }
    }
}

#[get("/list_features")]
async fn list_features(client: State<'_, RouteGuideClient<Channel>>) -> Stream<FeatureStream> {
    let rectangle = Rectangle {
        low: Some(Point {
            latitude: 400_000_000,
            longitude: -750_000_000,
        }),
        high: Some(Point {
            latitude: 420_000_000,
            longitude: -730_000_000,
        }),
    };
    let mut client = client.inner().clone();
    let stream = client
        .list_features(Request::new(rectangle))
        .await
        .unwrap()
        .into_inner();
    Stream::from(FeatureStream { stream })
}

#[rocket::launch]
async fn rocket() -> rocket::Rocket {
    rocket::ignite()
        .manage(
            create_route_guide_client("http://[::1]:10000")
                .await
                .unwrap(),
        )
        .mount("/", rocket::routes![list_features,])
}

出现错误:

error[E0277]: `from_generator::GenFuture<[static generator@Streaming<Feature>::message::{closure#0} for<'r, 's, 't0, 't1, 't2> {ResumeTy, &'r mut Streaming<Feature>, [closure@Streaming<Feature>::message::{closure#0}::{closure#0}], rocket::futures::future::PollFn<[closure@Streaming<Feature>::message::{closure#0}::{closure#0}]>, ()}]>` cannot be unpinned
   --> src/web_user.rs:34:15
    |
34  |         match Pin::new(&mut self.stream.message()).poll(cx) {
    |               ^^^^^^^^ within `impl std::future::Future`, the trait `Unpin` is not implemented for `from_generator::GenFuture<[static generator@Streaming<Feature>::message::{closure#0} for<'r, 's, 't0, 't1, 't2> {ResumeTy, &'r mut Streaming<Feature>, [closure@Streaming<Feature>::message::{closure#0}::{closure#0}], rocket::futures::future::PollFn<[closure@Streaming<Feature>::message::{closure#0}::{closure#0}]>, ()}]>`
    | 
   ::: /home/matan/.cargo/registry/src/github.com-1ecc6299db9ec823/tonic-0.4.0/src/codec/decode.rs:106:40
    |
106 |     pub async fn message(&mut self) -> Result<Option<T>, Status> {
    |                                        ------------------------- within this `impl std::future::Future`
    |
    = note: required because it appears within the type `impl std::future::Future`
    = note: required because it appears within the type `impl std::future::Future`
    = note: required by `Pin::<P>::new`

【问题讨论】:

  • 看起来 message() 函数是从 tonic Streaming 中选择下一条消息的助手。你不需要message() 函数来为AsyncRead 选择下一条消息你已经有Stream,你可以自己选择下一条消息,这里是代码play.rust-lang.org/…(它返回Pending 对于所有情况正如您的代码所做的那样,您可以根据需要更改它)

标签: rust async-await rust-tokio rust-rocket


【解决方案1】:

问题是从tonic::Streaming&lt;Feature&gt;::message() 生成的Future 没有实现Unpin,因为它是一个async 函数。让我们将此类型标记为MessageFuture,您不能安全地固定&amp;mut MessageFuture指针,因为取消引用的类型MessageFuture没有实现Unpin

为什么不安全?

来自reference,实现Unpin带来:

固定后可以安全移动的类型。

这意味着如果T:!UnpinPin&lt;&amp;mut T&gt; 不可移动,这很重要,因为由异步块创建的Futures 没有Unpin 实现,因为它可能持有来自自身的成员引用,并且如果您移动T 这个引用的指针也将被移动,但引用仍将指向相同的地址,为防止这种情况,它不应该是可移动的。请阅读"Pinning" section from async-book 了解原因。

注意:T:!Unpin 表示T 是没有Unpin 实现的类型。

解决方案

message() 函数是从tonic::Streaming&lt;T&gt; 中选择下一条消息的助手。您不需要特别调用 message() 从流中选择下一个元素,您的结构中已经有了实际的流。

struct FeatureStream {stream: tonic::Streaming<Feature>}

您可以等待AsyncRead 的下一条消息,例如:

impl AsyncRead for FeatureStream {
    fn poll_read(
        mut self: Pin<&mut Self>,
        cx: &mut Context<'_>,
        buf: &mut ReadBuf<'_>,
    ) -> Poll<std::io::Result<()>> {

       //it returns Pending for all cases as your code does, you can change it as you want
        match self.stream.poll_next_unpin(cx) {
            Poll::Ready(Some(Ok(m))) => Poll::Pending,
            Poll::Ready(Some(Err(e))) => Poll::Pending,
            Poll::Ready(None) => Poll::Pending,
            Poll::Pending => Poll::Pending
        }
    }
}

请注意tonic::Streaming&lt;T&gt; 实现了Unpin(reference)

【讨论】:

    【解决方案2】:

    感谢 Omer Erden 回答这个问题。因此,它归结为基于 tonic::Streaming 实现的 futures::Stream 特征来实现 AsyncRead。这是我实际使用的代码。

    impl AsyncRead for FeatureStream {
        fn poll_read(
            mut self: Pin<&mut Self>,
            cx: &mut Context<'_>,
            buf: &mut ReadBuf<'_>,
        ) -> Poll<std::io::Result<()>> {
            use futures::stream::StreamExt;
            use std::io::{Error, ErrorKind};
    
            match self.stream.poll_next_unpin(cx) {
                Poll::Ready(Some(Ok(m))) => {
                    buf.put_slice(format!("{:?}\n", m).as_bytes());
                    Poll::Ready(Ok(()))
                }
                Poll::Ready(Some(Err(e))) => {
                    Poll::Ready(Err(Error::new(ErrorKind::Other, format!("{:?}", e))))
                }
                Poll::Ready(None) => {
                    // None from a stream means the stream terminated. To indicate
                    // that from AsyncRead we return Ok and leave buf unchanged.
                    Poll::Ready(Ok(()))
                }
                Poll::Pending => Poll::Pending,
            }
        }
    }
    

    与此同时,我的解决方法是创建一个 TcpStream(它实现 AsyncRead)的两端并返回它的一端,同时生成一个单独的任务来实际写出结果。

    #[get("/list_features")]
    async fn list_features(
        client: State<'_, RouteGuideClient<Channel>>,
        tasks: State<'_, Mutex<Vec<tokio::task::JoinHandle<()>>>>,
    ) -> Result<Stream<TcpStream>, Debug<std::io::Error>> {
        let mut client = client.inner().clone();
        let mut feature_stream = client
            .list_features(Request::new(Rectangle {
                low: Some(Point {
                    latitude: 400000000,
                    longitude: -750000000,
                }),
                high: Some(Point {
                    latitude: 420000000,
                    longitude: -730000000,
                }),
            }))
            .await
            .unwrap()
            .into_inner();
    
        // Port 0 tells to operating system to choose an unused port.
        let tcp_listener = TcpListener::bind(("127.0.0.1", 0)).await?;
        let socket_addr = tcp_listener.local_addr().unwrap();
        tasks.lock().unwrap().push(tokio::spawn(async move {
            let mut tcp_stream = TcpStream::connect(socket_addr).await.unwrap();
    
            while let Some(feature) = feature_stream.message().await.unwrap() {
                match tcp_stream
                    .write_all(format!("{:?}\n", feature).as_bytes())
                    .await
                {
                    Ok(()) => (),
                    Err(e) => panic!(e),
                }
            }
            println!("End task");
            ()
        }));
        Ok(Stream::from(tcp_listener.accept().await?.0))
    }
    

    【讨论】:

    • 您不应该在轮询时运行阻塞代码,而是可以使用 tokio 计时器来创建适当的延迟,更多信息:stackoverflow.com/questions/48735952/…
    • 我假设你指的是我的睡眠?那是来自我正在运行的测试,在此处发布时忘记将其取出,感谢您指出并粘贴该链接。
    猜你喜欢
    • 2021-04-11
    • 1970-01-01
    • 2020-10-15
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    • 1970-01-01
    相关资源
    最近更新 更多