Add docs to codec and server
This commit is contained in:
@@ -96,7 +96,6 @@ impl HttpBody for BoxBody {
|
|||||||
type Error = Status;
|
type Error = Status;
|
||||||
|
|
||||||
fn is_end_stream(&self) -> bool {
|
fn is_end_stream(&self) -> bool {
|
||||||
// Body::is_end_stream(&self.inner)
|
|
||||||
self.inner.is_end_stream()
|
self.inner.is_end_stream()
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -62,10 +62,8 @@ impl<T> Grpc<T> {
|
|||||||
<T::ResponseBody as HttpBody>::Error: Into<crate::Error>,
|
<T::ResponseBody as HttpBody>::Error: Into<crate::Error>,
|
||||||
<T::ResponseBody as HttpBody>::Data: Into<Bytes>,
|
<T::ResponseBody as HttpBody>::Data: Into<Bytes>,
|
||||||
C: Codec<Encode = M1, Decode = M2>,
|
C: Codec<Encode = M1, Decode = M2>,
|
||||||
C::Encoder: Send + 'static,
|
|
||||||
C::Decoder: Send + 'static,
|
|
||||||
M1: Send + 'static,
|
M1: Send + 'static,
|
||||||
{
|
M2: Send + 'static {
|
||||||
let request = request.map(|m| stream::once(future::ok(m)));
|
let request = request.map(|m| stream::once(future::ok(m)));
|
||||||
self.client_streaming(request, path, codec).await
|
self.client_streaming(request, path, codec).await
|
||||||
}
|
}
|
||||||
@@ -84,9 +82,8 @@ impl<T> Grpc<T> {
|
|||||||
<T::ResponseBody as HttpBody>::Data: Into<Bytes>,
|
<T::ResponseBody as HttpBody>::Data: Into<Bytes>,
|
||||||
S: Stream<Item = Result<M1, Status>> + Send + 'static,
|
S: Stream<Item = Result<M1, Status>> + Send + 'static,
|
||||||
C: Codec<Encode = M1, Decode = M2>,
|
C: Codec<Encode = M1, Decode = M2>,
|
||||||
C::Encoder: Send + 'static,
|
M1: Send + 'static,
|
||||||
C::Decoder: Send + 'static,
|
M2: Send + 'static
|
||||||
M1: Send,
|
|
||||||
{
|
{
|
||||||
let (mut parts, body) = self.streaming(request, path, codec).await?.into_parts();
|
let (mut parts, body) = self.streaming(request, path, codec).await?.into_parts();
|
||||||
|
|
||||||
@@ -117,9 +114,8 @@ impl<T> Grpc<T> {
|
|||||||
<T::ResponseBody as HttpBody>::Error: Into<crate::Error>,
|
<T::ResponseBody as HttpBody>::Error: Into<crate::Error>,
|
||||||
<T::ResponseBody as HttpBody>::Data: Into<Bytes>,
|
<T::ResponseBody as HttpBody>::Data: Into<Bytes>,
|
||||||
C: Codec<Encode = M1, Decode = M2>,
|
C: Codec<Encode = M1, Decode = M2>,
|
||||||
C::Encoder: Send + 'static,
|
|
||||||
C::Decoder: Send + 'static,
|
|
||||||
M1: Send + 'static,
|
M1: Send + 'static,
|
||||||
|
M2: Send + 'static
|
||||||
{
|
{
|
||||||
let request = request.map(|m| stream::once(future::ok(m)));
|
let request = request.map(|m| stream::once(future::ok(m)));
|
||||||
self.streaming(request, path, codec).await
|
self.streaming(request, path, codec).await
|
||||||
@@ -139,9 +135,8 @@ impl<T> Grpc<T> {
|
|||||||
<T::ResponseBody as HttpBody>::Error: Into<crate::Error>,
|
<T::ResponseBody as HttpBody>::Error: Into<crate::Error>,
|
||||||
S: Stream<Item = Result<M1, Status>> + Send + 'static,
|
S: Stream<Item = Result<M1, Status>> + Send + 'static,
|
||||||
C: Codec<Encode = M1, Decode = M2>,
|
C: Codec<Encode = M1, Decode = M2>,
|
||||||
C::Encoder: Send + 'static,
|
M1: Send + 'static,
|
||||||
C::Decoder: Send + 'static,
|
M2: Send + 'static
|
||||||
M1: Send,
|
|
||||||
{
|
{
|
||||||
let mut parts = Parts::default();
|
let mut parts = Parts::default();
|
||||||
parts.path_and_query = Some(path);
|
parts.path_and_query = Some(path);
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
//! gRPC over HTTP2 client implementation.
|
//! gRPC client implementation.
|
||||||
//!
|
//!
|
||||||
//! This module contains the low level components to build a gRPC client. It
|
//! This module contains the low level components to build a gRPC client. It
|
||||||
//! provides a codec agnostic gRPC client dispatcher and a decorated tower
|
//! provides a codec agnostic gRPC client dispatcher and a decorated tower
|
||||||
|
|||||||
+14
-16
@@ -1,18 +1,17 @@
|
|||||||
use super::Decoder;
|
use super::Decoder;
|
||||||
use crate::body::BoxBody;
|
use crate::{Code, Status, BoxBody, metadata::MetadataMap};
|
||||||
use crate::metadata::MetadataMap;
|
|
||||||
use crate::{Code, Status};
|
|
||||||
use bytes::{Buf, BufMut, Bytes, BytesMut, IntoBuf};
|
use bytes::{Buf, BufMut, Bytes, BytesMut, IntoBuf};
|
||||||
use futures_core::Stream;
|
use futures_core::Stream;
|
||||||
use futures_util::{future, ready};
|
use futures_util::{future, ready};
|
||||||
use http::StatusCode;
|
use http::StatusCode;
|
||||||
use http_body::Body;
|
use http_body::Body;
|
||||||
use std::fmt;
|
use std::{fmt, pin::Pin, task::{Context, Poll}};
|
||||||
use std::pin::Pin;
|
|
||||||
use std::task::{Context, Poll};
|
|
||||||
use tracing::{debug, trace};
|
use tracing::{debug, trace};
|
||||||
|
|
||||||
// #[derive(Debug)]
|
/// Streaming requests and responses.
|
||||||
|
///
|
||||||
|
/// This will wrap some inner [`Body`] and [`Decoder`] and provide an interface
|
||||||
|
/// to fetch the message stream and trailing metadata
|
||||||
pub struct Streaming<T> {
|
pub struct Streaming<T> {
|
||||||
decoder: Box<dyn Decoder<Item = T, Error = Status> + Send + 'static>,
|
decoder: Box<dyn Decoder<Item = T, Error = Status> + Send + 'static>,
|
||||||
body: BoxBody,
|
body: BoxBody,
|
||||||
@@ -38,7 +37,7 @@ enum Direction {
|
|||||||
}
|
}
|
||||||
|
|
||||||
impl<T> Streaming<T> {
|
impl<T> Streaming<T> {
|
||||||
pub fn new_response<B, D>(decoder: D, body: B, status_code: StatusCode) -> Self
|
pub(crate) fn new_response<B, D>(decoder: D, body: B, status_code: StatusCode) -> Self
|
||||||
where
|
where
|
||||||
B: Body + Send + 'static,
|
B: Body + Send + 'static,
|
||||||
B::Data: Into<Bytes>,
|
B::Data: Into<Bytes>,
|
||||||
@@ -48,7 +47,7 @@ impl<T> Streaming<T> {
|
|||||||
Self::new(decoder, body, Direction::Response(status_code))
|
Self::new(decoder, body, Direction::Response(status_code))
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn new_empty<B, D>(decoder: D, body: B) -> Self
|
pub(crate) fn new_empty<B, D>(decoder: D, body: B) -> Self
|
||||||
where
|
where
|
||||||
B: Body + Send + 'static,
|
B: Body + Send + 'static,
|
||||||
B::Data: Into<Bytes>,
|
B::Data: Into<Bytes>,
|
||||||
@@ -58,7 +57,7 @@ impl<T> Streaming<T> {
|
|||||||
Self::new(decoder, body, Direction::EmptyResponse)
|
Self::new(decoder, body, Direction::EmptyResponse)
|
||||||
}
|
}
|
||||||
|
|
||||||
pub fn new_request<B, D>(decoder: D, body: B) -> Self
|
pub(crate) fn new_request<B, D>(decoder: D, body: B) -> Self
|
||||||
where
|
where
|
||||||
B: Body + Send + 'static,
|
B: Body + Send + 'static,
|
||||||
B::Data: Into<Bytes>,
|
B::Data: Into<Bytes>,
|
||||||
@@ -112,7 +111,7 @@ impl<T> Streaming<T> {
|
|||||||
// Trailers were not caught during poll_next and thus lets poll for
|
// Trailers were not caught during poll_next and thus lets poll for
|
||||||
// them manually.
|
// them manually.
|
||||||
let map =
|
let map =
|
||||||
future::poll_fn(|cx| unsafe { Pin::new_unchecked(&mut self.body) }.poll_trailers(cx))
|
future::poll_fn(|cx| Pin::new(&mut self.body).poll_trailers(cx))
|
||||||
.await
|
.await
|
||||||
.map_err(|e| Status::from_error(&e))?;
|
.map_err(|e| Status::from_error(&e))?;
|
||||||
|
|
||||||
@@ -189,8 +188,7 @@ impl<T> Stream for Streaming<T> {
|
|||||||
None => (),
|
None => (),
|
||||||
}
|
}
|
||||||
|
|
||||||
// FIXME: Figure out how to verify that this is safe
|
let chunk = match ready!(Pin::new(&mut self.body).poll_data(cx)) {
|
||||||
let chunk = match ready!(unsafe { Pin::new_unchecked(&mut self.body) }.poll_data(cx)) {
|
|
||||||
Some(Ok(d)) => Some(d),
|
Some(Ok(d)) => Some(d),
|
||||||
Some(Err(e)) => {
|
Some(Err(e)) => {
|
||||||
let err: crate::Error = e.into();
|
let err: crate::Error = e.into();
|
||||||
@@ -205,7 +203,7 @@ impl<T> Stream for Streaming<T> {
|
|||||||
if let Some(data) = chunk {
|
if let Some(data) = chunk {
|
||||||
self.buf.put(data);
|
self.buf.put(data);
|
||||||
} else {
|
} else {
|
||||||
// FIXME: get BytesMut to impl `Buf` directlty?
|
// TODO: get BytesMut to impl `Buf` directlty?
|
||||||
let buf1 = (&self.buf[..]).into_buf();
|
let buf1 = (&self.buf[..]).into_buf();
|
||||||
if buf1.has_remaining() {
|
if buf1.has_remaining() {
|
||||||
trace!("unexpected EOF decoding stream");
|
trace!("unexpected EOF decoding stream");
|
||||||
@@ -220,7 +218,7 @@ impl<T> Stream for Streaming<T> {
|
|||||||
}
|
}
|
||||||
|
|
||||||
if let Direction::Response(status) = self.direction {
|
if let Direction::Response(status) = self.direction {
|
||||||
match ready!(unsafe { Pin::new_unchecked(&mut self.body) }.poll_trailers(cx)) {
|
match ready!(Pin::new(&mut self.body).poll_trailers(cx)) {
|
||||||
Ok(trailer) => {
|
Ok(trailer) => {
|
||||||
if let Err(e) = crate::status::infer_grpc_status(trailer.as_ref(), status) {
|
if let Err(e) = crate::status::infer_grpc_status(trailer.as_ref(), status) {
|
||||||
return Some(Err(e)).into();
|
return Some(Err(e)).into();
|
||||||
@@ -243,6 +241,6 @@ impl<T> Stream for Streaming<T> {
|
|||||||
|
|
||||||
impl<T> fmt::Debug for Streaming<T> {
|
impl<T> fmt::Debug for Streaming<T> {
|
||||||
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
|
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
|
||||||
write!(f, "Streaming")
|
f.debug_struct("Streaming").finish()
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|||||||
+23
-6
@@ -1,23 +1,40 @@
|
|||||||
|
//! gRPC encoding and decoding.
|
||||||
|
//!
|
||||||
|
//! This module contains the generic `Codec` trait and a protobuf codec
|
||||||
|
//! based on prost.
|
||||||
|
|
||||||
mod decode;
|
mod decode;
|
||||||
mod encode;
|
mod encode;
|
||||||
mod prost;
|
mod prost;
|
||||||
|
|
||||||
pub use self::decode::Streaming;
|
pub use self::decode::Streaming;
|
||||||
pub use self::encode::{encode_client, encode_server, EncodeBody};
|
pub(crate) use self::encode::{encode_client, encode_server};
|
||||||
pub use self::prost::ProstCodec;
|
pub use self::prost::ProstCodec;
|
||||||
|
pub use tokio_codec::{Decoder, Encoder};
|
||||||
|
|
||||||
use crate::Status;
|
use crate::Status;
|
||||||
use tokio_codec::{Decoder, Encoder};
|
|
||||||
|
|
||||||
|
/// Triat that knows how to encode and decode gRPC messages.
|
||||||
pub trait Codec {
|
pub trait Codec {
|
||||||
type Encode;
|
/// The encodable message.
|
||||||
type Decode;
|
type Encode: Send + 'static;
|
||||||
|
/// The decodable message.
|
||||||
|
type Decode: Send + 'static;
|
||||||
|
|
||||||
type Encoder: Encoder<Item = Self::Encode, Error = Status>;
|
/// The encoder that can encode a message.
|
||||||
type Decoder: Decoder<Item = Self::Decode, Error = Status>;
|
type Encoder: Encoder<Item = Self::Encode, Error = Status> + Send + 'static;
|
||||||
|
/// The encoder that can decode a message.
|
||||||
|
type Decoder: Decoder<Item = Self::Decode, Error = Status> + Send + 'static;
|
||||||
|
|
||||||
|
/// The content type of this codec.
|
||||||
|
///
|
||||||
|
/// This should follow the `Content-Type` definition [here].
|
||||||
|
///
|
||||||
|
/// [here]: https://github.com/grpc/grpc/blob/master/doc/PROTOCOL-HTTP2.md#requests
|
||||||
const CONTENT_TYPE: &'static str;
|
const CONTENT_TYPE: &'static str;
|
||||||
|
|
||||||
|
/// Fetch the encoder.
|
||||||
fn encoder(&mut self) -> Self::Encoder;
|
fn encoder(&mut self) -> Self::Encoder;
|
||||||
|
/// Fetch the decoder.
|
||||||
fn decoder(&mut self) -> Self::Decoder;
|
fn decoder(&mut self) -> Self::Decoder;
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -1,16 +1,17 @@
|
|||||||
use super::Codec;
|
use super::{Codec, Decoder, Encoder};
|
||||||
use crate::{Code, Status};
|
use crate::{Code, Status};
|
||||||
use bytes::{BufMut, BytesMut};
|
use bytes::{BufMut, BytesMut};
|
||||||
use prost::Message;
|
use prost::Message;
|
||||||
use std::marker::PhantomData;
|
use std::marker::PhantomData;
|
||||||
use tokio_codec::{Decoder, Encoder};
|
|
||||||
|
|
||||||
|
/// A [`Codec`] that implements `application/grpc+proto` via the prost library..
|
||||||
#[derive(Debug, Clone)]
|
#[derive(Debug, Clone)]
|
||||||
pub struct ProstCodec<T, U> {
|
pub struct ProstCodec<T, U> {
|
||||||
_pd: PhantomData<(T, U)>,
|
_pd: PhantomData<(T, U)>,
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<T, U> ProstCodec<T, U> {
|
impl<T, U> ProstCodec<T, U> {
|
||||||
|
/// Create a new codec that knows how to encode `T` and decode `U`.
|
||||||
pub fn new() -> Self {
|
pub fn new() -> Self {
|
||||||
Self { _pd: PhantomData }
|
Self { _pd: PhantomData }
|
||||||
}
|
}
|
||||||
@@ -18,8 +19,8 @@ impl<T, U> ProstCodec<T, U> {
|
|||||||
|
|
||||||
impl<T, U> Codec for ProstCodec<T, U>
|
impl<T, U> Codec for ProstCodec<T, U>
|
||||||
where
|
where
|
||||||
T: Message,
|
T: Message + Send + 'static,
|
||||||
U: Message + Default,
|
U: Message + Default + Send + 'static,
|
||||||
{
|
{
|
||||||
type Encode = T;
|
type Encode = T;
|
||||||
type Decode = U;
|
type Decode = U;
|
||||||
@@ -38,6 +39,7 @@ where
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A [`Encoder`] that knows how to encode `T`.
|
||||||
pub struct ProstEncoder<T>(PhantomData<T>);
|
pub struct ProstEncoder<T>(PhantomData<T>);
|
||||||
|
|
||||||
impl<T: Message> Encoder for ProstEncoder<T> {
|
impl<T: Message> Encoder for ProstEncoder<T> {
|
||||||
@@ -56,6 +58,7 @@ impl<T: Message> Encoder for ProstEncoder<T> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
/// A [`Decoder`] that knows how to decode `U`.
|
||||||
pub struct ProstDecoder<U>(PhantomData<U>);
|
pub struct ProstDecoder<U>(PhantomData<U>);
|
||||||
|
|
||||||
impl<U: Message + Default> Decoder for ProstDecoder<U> {
|
impl<U: Message + Default> Decoder for ProstDecoder<U> {
|
||||||
|
|||||||
@@ -7,6 +7,7 @@ use bytes::Bytes;
|
|||||||
use futures_core::TryStream;
|
use futures_core::TryStream;
|
||||||
use futures_util::{future, stream, TryStreamExt};
|
use futures_util::{future, stream, TryStreamExt};
|
||||||
use http_body::Body;
|
use http_body::Body;
|
||||||
|
use std::fmt;
|
||||||
|
|
||||||
/// A gRPC Server handler.
|
/// A gRPC Server handler.
|
||||||
///
|
///
|
||||||
@@ -24,10 +25,6 @@ pub struct Grpc<T> {
|
|||||||
impl<T> Grpc<T>
|
impl<T> Grpc<T>
|
||||||
where
|
where
|
||||||
T: Codec,
|
T: Codec,
|
||||||
T::Decoder: Send + 'static,
|
|
||||||
T::Decode: Send + Unpin + 'static,
|
|
||||||
T::Encoder: Send + 'static,
|
|
||||||
T::Encode: Send + Unpin + 'static,
|
|
||||||
{
|
{
|
||||||
/// Creates a new gRPC client with the provided [`Codec`].
|
/// Creates a new gRPC client with the provided [`Codec`].
|
||||||
pub fn new(codec: T) -> Self {
|
pub fn new(codec: T) -> Self {
|
||||||
@@ -97,8 +94,6 @@ where
|
|||||||
) -> http::Response<BoxBody>
|
) -> http::Response<BoxBody>
|
||||||
where
|
where
|
||||||
S: ClientStreamingService<Streaming<T::Decode>, Response = T::Encode>,
|
S: ClientStreamingService<Streaming<T::Decode>, Response = T::Encode>,
|
||||||
T::Decode: Send + 'static,
|
|
||||||
T::Decoder: Send + 'static,
|
|
||||||
B: Body + Send + 'static,
|
B: Body + Send + 'static,
|
||||||
B::Data: Into<Bytes> + Send + 'static,
|
B::Data: Into<Bytes> + Send + 'static,
|
||||||
B::Error: Into<crate::Error> + Send + 'static,
|
B::Error: Into<crate::Error> + Send + 'static,
|
||||||
@@ -205,3 +200,9 @@ where
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
impl<T> fmt::Debug for Grpc<T> {
|
||||||
|
fn fmt(&self, f: &mut fmt::Formatter) -> fmt::Result {
|
||||||
|
f.debug_struct("Grpc").finish()
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|||||||
@@ -1,4 +1,4 @@
|
|||||||
//! gRPC over HTTP2 server implementation.
|
//! gRPC server implementation.
|
||||||
//!
|
//!
|
||||||
//! This module contains the low level components to build a gRPC server. It
|
//! This module contains the low level components to build a gRPC server. It
|
||||||
//! provides a codec agnostic gRPC server handler.
|
//! provides a codec agnostic gRPC server handler.
|
||||||
|
|||||||
Reference in New Issue
Block a user