pub struct ServerStream<B, RespView> { /* private fields */ }Expand description
Response from a server-streaming RPC.
Provides incremental access to response messages as they arrive from the server.
Messages are decoded one at a time from the HTTP response body using the
message() method, which returns Ok(None) for
a clean end and Err for a failed RPC — ? is the complete error
handling. Trailing metadata becomes available after the stream ends.
§Example
let mut stream = call_server_stream(&transport, &config, SVC_METHOD_SPEC.with_origin(SpecOrigin::Client), req, CallOptions::default()).await?;
println!("headers: {:?}", stream.headers());
while let Some(msg) = stream.message().await? {
println!("got message: {:?}", msg);
}
if let Some(trailers) = stream.trailers() {
println!("trailers: {:?}", trailers);
}Implementations§
Source§impl<B, RespView> ServerStream<B, RespView>where
B: Body<Data = Bytes> + Unpin,
B::Error: Display,
RespView: MessageView<'static> + Send,
RespView::Owned: Message + JsonDeserialize,
impl<B, RespView> ServerStream<B, RespView>where
B: Body<Data = Bytes> + Unpin,
B::Error: Display,
RespView: MessageView<'static> + Send,
RespView::Owned: Message + JsonDeserialize,
Sourcepub async fn message<M>(
&mut self,
) -> Result<Option<StreamMessage<M>>, ConnectError>where
RespView: MessageView<'static, Owned = M>,
M: HasMessageView<View<'static> = RespView>,
pub async fn message<M>(
&mut self,
) -> Result<Option<StreamMessage<M>>, ConnectError>where
RespView: MessageView<'static, Owned = M>,
M: HasMessageView<View<'static> = RespView>,
Fetch the next message from the stream.
Returns Ok(Some(msg)) for each message, Ok(None) when the stream
ends cleanly (gRPC status OK / error-free END_STREAM), or
Err(...) for everything else: protocol/decode/deadline errors and
a server error carried in the stream’s termination metadata (gRPC
trailers, gRPC-Web trailer frame, or Connect END_STREAM envelope).
Ok(None) means the RPC succeeded. Terminal errors arrive from
message() itself, as in tonic.
Every Err is terminal and sticky: the stream will never yield
another message, subsequent calls return the same Err (the same
policy as a failed stream construction; stronger than tonic, which
yields the error once and then reads as a clean end), and recovery
means making a new call — not re-polling this one. The terminal error also remains
inspectable via error(), and
trailers() is populated when termination metadata
was received — for both the Ok(None) and Err ends.
If a deadline was set on this call (via CallOptions::with_timeout
or ClientConfig::with_default_timeout), each message() poll is
bounded by it — gRPC deadline semantics are whole-call, so a hung
server won’t block indefinitely (matching grpc-java and connect-go).
§Errors
A response body that ends without its protocol’s termination
metadata is not a clean end and returns Err rather than
Ok(None): internal for a Connect stream missing its
END_STREAM envelope; for gRPC/gRPC-Web, internal when no
trailers arrived at all, unknown when trailers arrived without a
grpc-status, and unknown for a malformed grpc-status value —
matching grpc-go’s treatment of each case. A Trailers-Only response
carrying grpc-status: 0 in the headers (empty body) is a clean
end.
Whatever the cause, the returned error carries the response
metadata: ConnectError::response_headers() always, and
ConnectError::trailers() whenever termination metadata arrived.
Sourcepub fn trailers(&self) -> Option<&HeaderMap>
pub fn trailers(&self) -> Option<&HeaderMap>
Returns the trailing metadata, if available.
Only populated after message() reports the end of
the stream, and only when termination metadata was received — for
both the Ok(None) and Err ends.
Sourcepub fn error(&self) -> Option<&ConnectError>
pub fn error(&self) -> Option<&ConnectError>
Returns the terminal error that ended the stream, if any — a server error from the termination metadata (gRPC trailers / Connect END_STREAM), or a decode/transport/deadline failure.
message() already returns this same error, so most
callers never need this accessor; it exists for post-hoc inspection
alongside trailers().