Skip to content

Commit 1bdcc41

Browse files
committed
This change introduces StreamingMethodHandler as a generic trait for handling all types of gRPC methods on the server side, along with adapters for specific method types.
Key changes: - Introduce `StreamingMethodHandler` trait to provide a common interface for method execution. - Add `UnaryMethodAdapter`, `ServerStreamingAdapter`, `ClientStreamingAdapter`, and `BidiStreamingAdapter` to bridge specific method traits to the unified handler. - This generic API has intentionally been kept private due to the pitfalls below. Known pitfalls: - Lazy and ResponseHolder are not very great abstractions to workaround the fact that we may not always have ownership of the underlying message. - Asymmetry between the API: Request is lazy while response is in a holder. - The API is very generic heavy, but this should be acceptable at least until we talk about interceptors. C++ has a similar setup with Generic handler until being type erased in the codec layer.
1 parent b787609 commit 1bdcc41

8 files changed

Lines changed: 1332 additions & 0 deletions

File tree

grpc/src/server/method_handler.rs

Lines changed: 20 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,20 @@
1+
mod bidi_streaming_adapter;
2+
mod client_streaming_adapter;
3+
mod server_streaming_adapter;
4+
5+
pub use bidi_streaming_adapter::BidiStreamingAdapter;
6+
pub use client_streaming_adapter::ClientStreamingAdapter;
7+
8+
pub use server_streaming_adapter::ServerStreamingAdapter;
9+
10+
mod unary_adapter;
11+
pub use unary_adapter::UnaryMethodAdapter;
12+
13+
mod message_stream_handler;
14+
pub use message_stream_handler::MessageStreamHandler;
15+
16+
mod message_allocator;
17+
pub use message_allocator::{
18+
HeapMessageAllocator, HeapMessageHolder, HeapRequestHolder, HeapResponseHolder,
19+
RpcMessageAllocator, RpcMessageHolder, RpcRequestHolder, RpcResponseHolder,
20+
};
Lines changed: 248 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -0,0 +1,248 @@
1+
use crate::send_future::SendFuture;
2+
3+
use crate::server::call::HandlerCallOptions;
4+
use crate::server::call::{
5+
metadata_writer::TrailingMetadataWriter, Metadata, Outgoing, StreamingRequest,
6+
StreamingResponseWriter,
7+
};
8+
use crate::server::method_handler::MessageStreamHandler;
9+
use crate::server::stream::{
10+
PushStreamConsumer, PushStreamExt, PushStreamProducer, PushStreamWriter,
11+
};
12+
use crate::server::BidiStreamingMethod;
13+
use crate::Status;
14+
15+
use crate::server::call::Lazy;
16+
use crate::server::message::AsMut;
17+
use crate::server::method_handler::message_allocator::HeapResponseHolder;
18+
19+
/// Helper struct to adapt `PushStreamConsumer<Item=Outgoing<...>>` to `PushStreamConsumer<Item=Resp>`.
20+
struct ResponseConsumerAdapter<W, Resp> {
21+
writer: W,
22+
_phantom: std::marker::PhantomData<Resp>,
23+
}
24+
25+
impl<W, Resp> ResponseConsumerAdapter<W, Resp>
26+
where
27+
W: PushStreamConsumer<Outgoing<HeapResponseHolder<Resp>>> + Send,
28+
Resp: Send + 'static,
29+
{
30+
fn new(writer: W) -> Self {
31+
Self {
32+
writer,
33+
_phantom: std::marker::PhantomData,
34+
}
35+
}
36+
}
37+
38+
impl<W, Resp> PushStreamConsumer<Resp> for ResponseConsumerAdapter<W, Resp>
39+
where
40+
W: PushStreamConsumer<Outgoing<HeapResponseHolder<Resp>>> + Send,
41+
Resp: Send,
42+
{
43+
async fn write(&mut self, item: Resp) -> Result<(), Status> {
44+
let holder = HeapResponseHolder::new(item);
45+
self.writer.write(Outgoing::new(holder)).await
46+
}
47+
}
48+
49+
/// Adapter for `BidiStreamingMethod`.
50+
pub struct BidiStreamingAdapter<T>(pub T);
51+
52+
impl<T, Req, Resp> MessageStreamHandler for BidiStreamingAdapter<T>
53+
where
54+
T: BidiStreamingMethod<Req = Req, Resp = Resp> + Sync,
55+
Req: AsMut + Default + Send + 'static,
56+
Resp: AsMut + Default + Send + 'static,
57+
{
58+
type Req = Req;
59+
type Resp = Resp;
60+
61+
type ResponseHolder = HeapResponseHolder<Resp>;
62+
63+
async fn call<P, W, L>(
64+
&self,
65+
_options: HandlerCallOptions,
66+
req: StreamingRequest<P>,
67+
writer: W,
68+
) -> Result<(), Status>
69+
where
70+
P: PushStreamProducer<Item = L> + Send + 'static,
71+
W: StreamingResponseWriter<Outgoing<Self::ResponseHolder>> + Send,
72+
L: Lazy<Req>,
73+
<W as StreamingResponseWriter<Outgoing<Self::ResponseHolder>>>::MessageWriter: 'static,
74+
{
75+
// 1. Send Initial Metadata
76+
let (msg_writer, trailer_writer) =
77+
writer.send_initial_metadata(Metadata::default()).await?;
78+
79+
// 2. Adapt input stream (L -> Req)
80+
let (_, stream) = req.into_parts();
81+
let req_stream = stream.then(|lazy_req| async move {
82+
let mut req = Req::default();
83+
lazy_req.resolve(req.as_mut()).make_send().await?;
84+
Ok(req)
85+
});
86+
87+
// 3. Call method
88+
{
89+
let adapter = ResponseConsumerAdapter::new(msg_writer);
90+
let stream_writer = PushStreamWriter::new(adapter);
91+
self.0
92+
.bidi_streaming(req_stream, stream_writer)
93+
.await
94+
.map_err(|s| s.into_status())?;
95+
}
96+
97+
// 4. Send Trailers
98+
trailer_writer
99+
.send_trailing_metadata(Metadata::default())
100+
.await
101+
}
102+
}
103+
104+
#[cfg(test)]
105+
mod tests {
106+
use super::*;
107+
use crate::server::call::metadata_writer::{InitialMetadataWriter, TrailingMetadataWriter};
108+
use crate::server::call::test_util::StreamingResponseImpl;
109+
use crate::server::stream::{
110+
PushStream, PushStreamConsumer, PushStreamProducer, PushStreamWriter,
111+
};
112+
use crate::server::BidiStreamingMethod;
113+
use crate::Status;
114+
115+
struct MockBidiStreamingMethod {
116+
resp_to_return: i32,
117+
}
118+
119+
impl BidiStreamingMethod for MockBidiStreamingMethod {
120+
type Req = i32;
121+
type Resp = i32;
122+
async fn bidi_streaming<P, C>(
123+
&self,
124+
_req: PushStream<P>,
125+
writer: PushStreamWriter<C>,
126+
) -> Result<(), crate::ServerStatus>
127+
where
128+
P: PushStreamProducer<Item = i32> + Send,
129+
C: PushStreamConsumer<i32> + Send,
130+
{
131+
Ok(())
132+
}
133+
}
134+
135+
struct MockMetadataWriter {
136+
sent_initial: bool,
137+
sent_trailing: bool,
138+
}
139+
140+
impl InitialMetadataWriter for MockMetadataWriter {
141+
async fn send_initial_metadata(mut self, _metadata: Metadata) -> Result<(), Status> {
142+
self.sent_initial = true;
143+
Ok(())
144+
}
145+
}
146+
147+
impl TrailingMetadataWriter for MockMetadataWriter {
148+
async fn send_trailing_metadata(mut self, _metadata: Metadata) -> Result<(), Status> {
149+
self.sent_trailing = true;
150+
Ok(())
151+
}
152+
}
153+
154+
struct MockProducer;
155+
impl PushStreamProducer for MockProducer {
156+
type Item = i32;
157+
async fn produce(
158+
self,
159+
_writer: PushStreamWriter<impl crate::server::stream::PushStreamConsumer<Self::Item>,
160+
>,
161+
) -> Result<(), Status> {
162+
Ok(())
163+
}
164+
}
165+
166+
#[tokio::test]
167+
async fn test_bidi_streaming_adapter_v2_success() {
168+
use protobuf_well_known_types::Timestamp;
169+
170+
struct MockBidiStreamingMethodV2;
171+
impl BidiStreamingMethod for MockBidiStreamingMethodV2 {
172+
type Req = Timestamp;
173+
type Resp = Timestamp;
174+
async fn bidi_streaming<P, C>(
175+
&self,
176+
_req: PushStream<P>,
177+
mut writer: PushStreamWriter<C>,
178+
) -> Result<(), crate::ServerStatus>
179+
where
180+
P: PushStreamProducer<Item = Timestamp> + Send,
181+
C: PushStreamConsumer<Timestamp> + Send,
182+
{
183+
// Write one item
184+
let mut msg = Timestamp::new();
185+
msg.set_seconds(100);
186+
writer.write(msg).await.unwrap();
187+
Ok(())
188+
}
189+
}
190+
191+
let method = MockBidiStreamingMethodV2;
192+
let adapter = BidiStreamingAdapter(method);
193+
194+
// V2 expects Lazy<Req>
195+
struct MockLazy(Timestamp);
196+
impl Lazy<Timestamp> for MockLazy {
197+
async fn resolve(self, mut dest: <Timestamp as AsMut>::Mut<'_>) -> Result<(), Status> {
198+
dest.set_seconds(self.0.seconds());
199+
Ok(())
200+
}
201+
}
202+
203+
struct MockLazyProducer;
204+
impl PushStreamProducer for MockLazyProducer {
205+
type Item = MockLazy;
206+
async fn produce(
207+
self,
208+
writer: PushStreamWriter<impl PushStreamConsumer<Self::Item>>,
209+
) -> Result<(), Status> {
210+
// Just close
211+
Ok(())
212+
}
213+
}
214+
215+
let producer = MockLazyProducer;
216+
let stream = PushStream::new(producer);
217+
let req = StreamingRequest::new(stream, Metadata::default());
218+
219+
// Consumer for Outgoing<HeapResponseHolder<Timestamp>>
220+
struct MockV2Consumer;
221+
impl PushStreamConsumer<Outgoing<HeapResponseHolder<Timestamp>>> for MockV2Consumer {
222+
async fn write(
223+
&mut self,
224+
_item: Outgoing<HeapResponseHolder<Timestamp>>,
225+
) -> Result<(), Status> {
226+
Ok(())
227+
}
228+
}
229+
230+
let stream_writer = PushStreamWriter::new(MockV2Consumer);
231+
232+
let initial_writer = MockMetadataWriter {
233+
sent_initial: false,
234+
sent_trailing: false,
235+
};
236+
let trailing_writer = MockMetadataWriter {
237+
sent_initial: false,
238+
sent_trailing: false,
239+
};
240+
let writer = StreamingResponseImpl::new(stream_writer, initial_writer, trailing_writer);
241+
242+
let result = adapter
243+
.call(HandlerCallOptions::default(), req, writer)
244+
.await;
245+
246+
assert!(result.is_ok());
247+
}
248+
}

0 commit comments

Comments
 (0)