Skip to content

Commit e2f4bd0

Browse files
committed
feat(grpc): Remove trailers from SendStream and make Handle return Trailers .
1 parent 9cef196 commit e2f4bd0

4 files changed

Lines changed: 36 additions & 29 deletions

File tree

grpc/examples/inmemory.rs

Lines changed: 4 additions & 11 deletions
Original file line numberDiff line numberDiff line change
@@ -37,6 +37,7 @@ use grpc::core::RecvMessage;
3737
use grpc::core::RequestHeaders;
3838
use grpc::core::SendMessage;
3939
use grpc::core::ServerResponseStreamItem;
40+
use grpc::core::Trailers;
4041
use grpc::credentials::InsecureChannelCredentials;
4142
use grpc::inmemory;
4243
use grpc::server;
@@ -84,7 +85,7 @@ impl Handle for Handler {
8485
_options: CallOptions,
8586
tx: &mut impl server::SendStream,
8687
mut rx: impl server::RecvStream + 'static,
87-
) {
88+
) -> Trailers {
8889
let method = headers.method_name().clone();
8990
let id = self.id.clone();
9091
// Send headers
@@ -108,16 +109,8 @@ impl Handle for Handler {
108109
)
109110
.await;
110111
}
111-
// Send trailers
112-
let _ = tx
113-
.send(
114-
ServerResponseStreamItem::Trailers(grpc::core::Trailers::new(grpc::Status::new(
115-
grpc::StatusCode::Ok,
116-
"OK",
117-
))),
118-
server::SendOptions::default(),
119-
)
120-
.await;
112+
// Return trailers
113+
Trailers::new(grpc::Status::new(grpc::StatusCode::Ok, "OK"))
121114
}
122115
}
123116

grpc/src/core/mod.rs

Lines changed: 1 addition & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -101,8 +101,7 @@ impl dyn RecvMessage + '_ {
101101
}
102102
}
103103

104-
/// ResponseStreamItem represents an item in a response stream (either server
105-
/// sending or client receiving).
104+
/// ResponseStreamItem represents an item in a response stream from the client's view.
106105
///
107106
/// A response stream must always contain items exactly as follows:
108107
///
@@ -135,8 +134,6 @@ pub enum ServerResponseStreamItem<'a> {
135134
Headers(ResponseHeaders),
136135
/// Indicates a message on the stream.
137136
Message(&'a dyn SendMessage),
138-
/// Indicates trailers were received on the stream and includes the trailers.
139-
Trailers(Trailers),
140137
}
141138

142139
/// Contains all information transmitted in the response headers of an RPC.

grpc/src/inmemory/mod.rs

Lines changed: 18 additions & 6 deletions
Original file line numberDiff line numberDiff line change
@@ -85,6 +85,7 @@ struct InMemoryServerCall {
8585
headers: RequestHeaders,
8686
req_rx: mpsc::UnboundedReceiver<InMemoryRequestStreamItem>,
8787
resp_tx: mpsc::UnboundedSender<InMemoryResponseStreamItem>,
88+
trailer_tx: oneshot::Sender<Trailers>,
8889
}
8990

9091
enum InMemoryRequestStreamItem {
@@ -95,7 +96,6 @@ enum InMemoryRequestStreamItem {
9596
enum InMemoryResponseStreamItem {
9697
Headers(ResponseHeaders),
9798
Message(Box<dyn Buf + Send + Sync>),
98-
Trailers(Trailers),
9999
StreamClosed,
100100
}
101101

@@ -179,6 +179,7 @@ impl ServerListener for InMemoryListener {
179179
headers: call.headers,
180180
send: InMemoryServerSendStream { tx: call.resp_tx },
181181
recv: InMemoryServerRecvStream { rx: call.req_rx },
182+
trailers_tx: call.trailer_tx,
182183
})
183184
}
184185
_ = self.inner.close_notify.notified() => {
@@ -204,7 +205,6 @@ impl ServerSendStream for InMemoryServerSendStream {
204205
let buf = m.encode().map_err(|_| ())?;
205206
InMemoryResponseStreamItem::Message(buf)
206207
}
207-
ServerResponseStreamItem::Trailers(t) => InMemoryResponseStreamItem::Trailers(t),
208208
};
209209

210210
self.tx.send(inmemory_item).map_err(|_| ())
@@ -246,18 +246,23 @@ impl Invoke for InMemoryConnection {
246246
) -> (Self::SendStream, Self::RecvStream) {
247247
let (req_tx, req_rx) = mpsc::unbounded_channel::<InMemoryRequestStreamItem>();
248248
let (resp_tx, resp_rx) = mpsc::unbounded_channel::<InMemoryResponseStreamItem>();
249+
let (trailer_tx, trailer_rx) = oneshot::channel();
249250

250251
let call = InMemoryServerCall {
251252
headers,
252253
req_rx,
253254
resp_tx,
255+
trailer_tx,
254256
};
255257

256-
let _ = self.s.try_send(call);
258+
let _ = self.s.send(call).await;
257259

258260
(
259261
Box::new(InMemoryClientSendStream { tx: Some(req_tx) }),
260-
Box::new(InMemoryClientRecvStream { rx: resp_rx }),
262+
Box::new(InMemoryClientRecvStream {
263+
rx: resp_rx,
264+
trailer_rx: Some(trailer_rx),
265+
}),
261266
)
262267
}
263268
}
@@ -299,6 +304,7 @@ impl Drop for InMemoryClientSendStream {
299304

300305
pub struct InMemoryClientRecvStream {
301306
rx: mpsc::UnboundedReceiver<InMemoryResponseStreamItem>,
307+
trailer_rx: Option<oneshot::Receiver<Trailers>>,
302308
}
303309

304310
impl ClientRecvStream for InMemoryClientRecvStream {
@@ -309,8 +315,14 @@ impl ClientRecvStream for InMemoryClientRecvStream {
309315
msg.decode(&mut buf).unwrap();
310316
ClientResponseStreamItem::Message(())
311317
}
312-
Some(InMemoryResponseStreamItem::Trailers(t)) => ClientResponseStreamItem::Trailers(t),
313-
_ => ClientResponseStreamItem::StreamClosed,
318+
_ => {
319+
if let Some(trailer_rx) = self.trailer_rx.take() {
320+
if let Ok(trailers) = trailer_rx.await {
321+
return ClientResponseStreamItem::Trailers(trailers);
322+
}
323+
}
324+
ClientResponseStreamItem::StreamClosed
325+
}
314326
}
315327
}
316328
}

grpc/src/server/mod.rs

Lines changed: 13 additions & 8 deletions
Original file line numberDiff line numberDiff line change
@@ -29,6 +29,8 @@ use crate::client::CallOptions;
2929
use crate::core::RecvMessage;
3030
use crate::core::RequestHeaders;
3131
use crate::core::ServerResponseStreamItem;
32+
use crate::core::Trailers;
33+
use tokio::sync::oneshot;
3234

3335
pub struct Server {
3436
handler: Option<Arc<dyn DynHandle>>,
@@ -38,6 +40,7 @@ pub struct Call<SS, RS> {
3840
pub headers: RequestHeaders,
3941
pub send: SS,
4042
pub recv: RS,
43+
pub trailers_tx: oneshot::Sender<Trailers>,
4144
}
4245

4346
#[trait_variant::make(Send)]
@@ -64,11 +67,14 @@ impl Server {
6467
let mut send: Box<dyn DynSendStream> = Box::new(call.send);
6568
let recv = BoxedRecvStream(Box::new(call.recv));
6669
let options = CallOptions::default();
67-
self.handler
70+
let trailers_tx = call.trailers_tx;
71+
let trailers = self
72+
.handler
6873
.as_ref()
6974
.unwrap()
7075
.dyn_handle(call.headers, options, &mut *send, recv)
7176
.await;
77+
let _ = trailers_tx.send(trailers);
7278
}
7379
}
7480
}
@@ -92,7 +98,7 @@ pub trait Handle: Send + Sync {
9298
options: CallOptions,
9399
tx: &mut impl SendStream,
94100
rx: impl RecvStream + 'static,
95-
);
101+
) -> Trailers;
96102
}
97103

98104
#[async_trait]
@@ -103,7 +109,7 @@ trait DynHandle: Send + Sync {
103109
options: CallOptions,
104110
tx: &mut dyn DynSendStream,
105111
rx: BoxedRecvStream,
106-
);
112+
) -> Trailers;
107113
}
108114

109115
#[async_trait]
@@ -114,7 +120,7 @@ impl<T: Handle> DynHandle for T {
114120
options: CallOptions,
115121
mut tx: &mut dyn DynSendStream,
116122
rx: BoxedRecvStream,
117-
) {
123+
) -> Trailers {
118124
self.handle(headers, options, &mut tx, rx).await
119125
}
120126
}
@@ -146,10 +152,9 @@ impl RecvStream for BoxedRecvStream {
146152
pub trait SendStream {
147153
/// Sends the next item on the stream. Returns `Ok(())` on success, or
148154
/// `Err(())` on failure. `Err(())` is a terminal state.
149-
/// Sending is also considered complete after successfully sending
150-
/// `ServerResponseStreamItem::Trailers`.
151-
/// Calling this method after an error or after sending trailers should be
152-
/// avoided and is unspecified.
155+
/// Sending is considered complete when the `handle` method returns
156+
/// `Trailers`.
157+
/// Calling this method after an error should be avoided and is unspecified.
153158
///
154159
/// # Cancel safety
155160
///

0 commit comments

Comments
 (0)