streaming
This commit is contained in:
+11
-2
@@ -1,6 +1,6 @@
|
|||||||
use crate::{Code, Status};
|
use crate::{Code, Status};
|
||||||
use bytes::{Bytes, IntoBuf};
|
use bytes::{Bytes, IntoBuf};
|
||||||
use futures_core::TryStream;
|
use futures_core::Stream;
|
||||||
use futures_util::{ready, TryStreamExt};
|
use futures_util::{ready, TryStreamExt};
|
||||||
use http::HeaderMap;
|
use http::HeaderMap;
|
||||||
use http_body::Body;
|
use http_body::Body;
|
||||||
@@ -13,9 +13,18 @@ pub struct AsyncBody<S> {
|
|||||||
error: Option<Status>,
|
error: Option<Status>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
impl<S> AsyncBody<S>
|
||||||
|
where
|
||||||
|
S: Stream<Item = Result<crate::body::BytesBuf, Status>> + Unpin,
|
||||||
|
{
|
||||||
|
pub fn new(inner: S) -> Self {
|
||||||
|
Self { inner, error: None }
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
impl<S> Body for AsyncBody<S>
|
impl<S> Body for AsyncBody<S>
|
||||||
where
|
where
|
||||||
S: TryStream<Ok = BytesBuf, Error = Status> + Unpin,
|
S: Stream<Item = Result<crate::body::BytesBuf, Status>> + Unpin,
|
||||||
{
|
{
|
||||||
type Data = BytesBuf;
|
type Data = BytesBuf;
|
||||||
type Error = Status;
|
type Error = Status;
|
||||||
|
|||||||
@@ -18,6 +18,8 @@ pub use response::Response;
|
|||||||
pub use status::{Code, Status};
|
pub use status::{Code, Status};
|
||||||
pub use tonic_macros::server;
|
pub use tonic_macros::server;
|
||||||
|
|
||||||
|
pub(crate) use error::Error;
|
||||||
|
|
||||||
use std::future::Future;
|
use std::future::Future;
|
||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
|
|
||||||
|
|||||||
+133
-14
@@ -1,14 +1,15 @@
|
|||||||
use crate::{Request, Response, Status};
|
#![allow(dead_code)]
|
||||||
|
|
||||||
|
use crate::{Code, Request, Response, Status};
|
||||||
use async_stream::stream;
|
use async_stream::stream;
|
||||||
use bytes::{Bytes, BytesMut, IntoBuf};
|
use bytes::{Buf, BufMut, Bytes, BytesMut, IntoBuf};
|
||||||
use futures_core::TryStream;
|
use futures_core::{Stream, TryStream};
|
||||||
use futures_util::{stream, StreamExt, TryStreamExt};
|
use futures_util::{future, stream, StreamExt, TryStreamExt};
|
||||||
|
use http_body::Body;
|
||||||
use std::future::Future;
|
use std::future::Future;
|
||||||
use tokio_codec::{Decoder, Encoder};
|
use tokio_codec::{Decoder, Encoder};
|
||||||
use tower_service::Service;
|
use tower_service::Service;
|
||||||
|
use tracing::{debug, trace};
|
||||||
#[allow(dead_code)]
|
|
||||||
type Result<T> = std::result::Result<Response<T>, Status>;
|
|
||||||
|
|
||||||
pub trait Codec {
|
pub trait Codec {
|
||||||
type Encode;
|
type Encode;
|
||||||
@@ -16,7 +17,7 @@ pub trait Codec {
|
|||||||
}
|
}
|
||||||
|
|
||||||
pub struct Encode<T, U> {
|
pub struct Encode<T, U> {
|
||||||
inner: T,
|
encoder: T,
|
||||||
source: U,
|
source: U,
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -25,19 +26,19 @@ where
|
|||||||
T: Encoder,
|
T: Encoder,
|
||||||
U: TryStream<Ok = T::Item, Error = Status> + Unpin,
|
U: TryStream<Ok = T::Item, Error = Status> + Unpin,
|
||||||
{
|
{
|
||||||
pub fn new(inner: T, source: U) -> Self {
|
pub fn new(encoder: T, source: U) -> Self {
|
||||||
Encode { inner, source }
|
Encode { encoder, source }
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn encode<'a>(
|
pub fn encode<'a>(
|
||||||
&'a mut self,
|
&'a mut self,
|
||||||
buf: &'a mut BytesMut,
|
buf: &'a mut BytesMut,
|
||||||
) -> impl TryStream<Ok = crate::body::BytesBuf, Error = Status> + 'a {
|
) -> impl Stream<Item = Result<crate::body::BytesBuf, Status>> + 'a {
|
||||||
stream! {
|
stream! {
|
||||||
loop {
|
loop {
|
||||||
match self.source.try_next().await {
|
match self.source.try_next().await {
|
||||||
Ok(Some(item)) => {
|
Ok(Some(item)) => {
|
||||||
self.inner.encode(item, buf).map_err(drop).unwrap();
|
self.encoder.encode(item, buf).map_err(drop).unwrap();
|
||||||
let len = buf.len();
|
let len = buf.len();
|
||||||
yield Ok(buf.split_to(len).freeze().into_buf());
|
yield Ok(buf.split_to(len).freeze().into_buf());
|
||||||
},
|
},
|
||||||
@@ -49,17 +50,135 @@ where
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
pub struct Streaming<T> {
|
||||||
|
decoder: T,
|
||||||
|
buf: BytesMut,
|
||||||
|
state: State,
|
||||||
|
}
|
||||||
|
|
||||||
|
#[derive(Debug)]
|
||||||
|
enum State {
|
||||||
|
ReadHeader,
|
||||||
|
ReadBody { compression: bool, len: usize },
|
||||||
|
Done,
|
||||||
|
}
|
||||||
|
|
||||||
|
impl<T> Streaming<T>
|
||||||
|
where
|
||||||
|
T: Decoder,
|
||||||
|
T::Item: Unpin + 'static,
|
||||||
|
{
|
||||||
|
pub fn decode<'a, B>(
|
||||||
|
&'a mut self,
|
||||||
|
source: &'a mut B,
|
||||||
|
) -> impl Stream<Item = Result<T::Item, Status>> + 'a
|
||||||
|
where
|
||||||
|
B: Body,
|
||||||
|
B::Error: Into<crate::Error>,
|
||||||
|
{
|
||||||
|
stream! {
|
||||||
|
loop {
|
||||||
|
// TODO: use try_stream! and ?
|
||||||
|
if let Some(item) = self.decode_chunk().unwrap() {
|
||||||
|
yield Ok(item);
|
||||||
|
}
|
||||||
|
|
||||||
|
let chunk = match future::poll_fn(|cx| source.poll_data(cx)).await {
|
||||||
|
Some(Ok(d)) => Some(d),
|
||||||
|
Some(Err(e)) => {
|
||||||
|
let err = e.into();
|
||||||
|
debug!("decoder inner stream error: {:?}", err);
|
||||||
|
let status = Status::from_error(&*err);
|
||||||
|
yield Err(status);
|
||||||
|
break;
|
||||||
|
},
|
||||||
|
None => None,
|
||||||
|
};
|
||||||
|
|
||||||
|
if let Some(data)= chunk {
|
||||||
|
self.buf.put(data);
|
||||||
|
} else {
|
||||||
|
if self.buf.has_remaining_mut() {
|
||||||
|
trace!("unexpected EOF decoding stream");
|
||||||
|
yield Err(Status::new(
|
||||||
|
Code::Internal,
|
||||||
|
"Unexpected EOF decoding stream.".to_string(),
|
||||||
|
));
|
||||||
|
} else {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
fn decode_chunk(&mut self) -> Result<Option<T::Item>, Status> {
|
||||||
|
let buf = (&self.buf).into_buf();
|
||||||
|
|
||||||
|
if let State::ReadHeader = self.state {
|
||||||
|
if buf.remaining() < 5 {
|
||||||
|
return Ok(None);
|
||||||
|
}
|
||||||
|
|
||||||
|
let is_compressed = match buf.get_u8() {
|
||||||
|
0 => false,
|
||||||
|
1 => {
|
||||||
|
trace!("message compressed, compression not supported yet");
|
||||||
|
return Err(crate::Status::new(
|
||||||
|
crate::Code::Unimplemented,
|
||||||
|
"Message compressed, compression not supported yet.".to_string(),
|
||||||
|
));
|
||||||
|
}
|
||||||
|
f => {
|
||||||
|
trace!("unexpected compression flag");
|
||||||
|
return Err(crate::Status::new(
|
||||||
|
crate::Code::Internal,
|
||||||
|
format!("Unexpected compression flag: {}", f),
|
||||||
|
));
|
||||||
|
}
|
||||||
|
};
|
||||||
|
let len = (&self.buf[..]).into_buf().get_u32_be() as usize;
|
||||||
|
|
||||||
|
self.state = State::ReadBody {
|
||||||
|
compression: is_compressed,
|
||||||
|
len,
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
if let State::ReadBody { len, .. } = self.state {
|
||||||
|
if buf.remaining() < len {
|
||||||
|
return Ok(None);
|
||||||
|
}
|
||||||
|
|
||||||
|
match self.decoder.decode(&mut self.buf) {
|
||||||
|
Ok(Some(msg)) => {
|
||||||
|
self.state = State::ReadHeader;
|
||||||
|
return Ok(Some(msg));
|
||||||
|
}
|
||||||
|
Err(e) => {
|
||||||
|
return Err(e);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
Ok(None)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
#[cfg(test)]
|
#[cfg(test)]
|
||||||
mod tests {
|
mod tests {
|
||||||
use crate::body::AsyncBody;
|
use crate::body::AsyncBody;
|
||||||
use crate::server::Encode;
|
use crate::server::Encode;
|
||||||
use bytes::Bytes;
|
use bytes::{Bytes, BytesMut};
|
||||||
use tokio_codec::BytesCodec;
|
use tokio_codec::BytesCodec;
|
||||||
|
|
||||||
#[test]
|
#[test]
|
||||||
fn body() {
|
fn body() {
|
||||||
let stream = futures_util::stream::iter(vec![Ok(Bytes::new())]);
|
let stream = futures_util::stream::iter(vec![Ok(Bytes::new())]);
|
||||||
let encode = Encode::new(BytesCodec::new(), stream);
|
let mut encode = Encode::new(BytesCodec::new(), stream);
|
||||||
|
|
||||||
|
let mut buf = BytesMut::with_capacity(1024);
|
||||||
|
AsyncBody::new(Box::pin(encode.encode(&mut buf)));
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user