Skip to content

Commit f564045

Browse files
committed
Dyn ByteStreamHandler
1 parent 08dd326 commit f564045

3 files changed

Lines changed: 57 additions & 0 deletions

File tree

grpc/Cargo.toml

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -68,6 +68,7 @@ flate2 = { version = "1.0", optional = true }
6868
zstd = { version = "0.13", optional = true }
6969
trait-variant = "0.1"
7070
send-future = "0.1"
71+
dynosaur = "0.3.0"
7172

7273
[dev-dependencies]
7374
async-stream = "0.3.6"

grpc/src/server/method_handler.rs

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -24,3 +24,11 @@ pub use codec_unary_method_handler::CodecUnaryMethodHandler;
2424

2525
mod message_allocator;
2626
pub use message_allocator::{HeapMessageAllocator, MessageAllocator, MessageHolder};
27+
28+
mod byte_stream_method_handler;
29+
pub use byte_stream_method_handler::{
30+
ByteStreamMethodHandler, DynByteStreamMethodHandler, GenericByteStreamAdapter,
31+
};
32+
33+
/// The default response body type produced by standard Codecs.
34+
pub type CodecRespB = bytes::Bytes;
Lines changed: 48 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,48 @@
1+
use crate::call::{HandlerCallOptions, Incoming, StreamingRequest, StreamingResponseWriter};
2+
use crate::server::method_handler::GenericByteStreamMethodHandler;
3+
use crate::stream::PushStreamProducer;
4+
use crate::Status;
5+
use bytes::Buf;
6+
7+
/// A method handler that processes raw bytes.
8+
#[trait_variant::make(Send)]
9+
#[dynosaur::dynosaur(pub DynByteStreamMethodHandler = dyn(box) ByteStreamMethodHandler)]
10+
pub trait ByteStreamMethodHandler<ReqB, W, P>: Send + Sync {
11+
type RespB: Buf + Send;
12+
13+
async fn call(
14+
&self,
15+
options: HandlerCallOptions,
16+
req: StreamingRequest<Incoming<ReqB>, P>,
17+
resp: W,
18+
) -> Result<(), Status>
19+
where
20+
ReqB: Buf + Send,
21+
P: PushStreamProducer<Item = Incoming<ReqB>> + Send,
22+
W: StreamingResponseWriter<Self::RespB>;
23+
}
24+
25+
/// Adapter to bridge `GenericByteStreamMethodHandler` to `ByteStreamMethodHandler`.
26+
pub struct GenericByteStreamAdapter<T>(pub T);
27+
28+
impl<T, ReqB, W, P> ByteStreamMethodHandler<ReqB, W, P> for GenericByteStreamAdapter<T>
29+
where
30+
T: GenericByteStreamMethodHandler,
31+
W: StreamingResponseWriter<T::RespB> + Send,
32+
{
33+
type RespB = T::RespB;
34+
35+
async fn call(
36+
&self,
37+
options: HandlerCallOptions,
38+
req: StreamingRequest<Incoming<ReqB>, P>,
39+
resp: W,
40+
) -> Result<(), Status>
41+
where
42+
ReqB: Buf + Send,
43+
W: StreamingResponseWriter<Self::RespB>,
44+
P: PushStreamProducer<Item = Incoming<ReqB>> + Send,
45+
{
46+
<T as GenericByteStreamMethodHandler>::call(&self.0, options, req, resp).await
47+
}
48+
}

0 commit comments

Comments
 (0)