Skip to main content

ServerStream

Struct ServerStream 

Source
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,

Source

pub fn headers(&self) -> &HeaderMap

Returns the response headers.

Source

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.

Source

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.

Source

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().

Trait Implementations§

Source§

impl<B, RespView> Debug for ServerStream<B, RespView>

Source§

fn fmt(&self, f: &mut Formatter<'_>) -> Result

Formats the value using the given formatter. Read more

Auto Trait Implementations§

§

impl<B, RespView> !RefUnwindSafe for ServerStream<B, RespView>

§

impl<B, RespView> !UnwindSafe for ServerStream<B, RespView>

§

impl<B, RespView> Freeze for ServerStream<B, RespView>
where B: Freeze, PhantomData<RespView>: Freeze,

§

impl<B, RespView> Send for ServerStream<B, RespView>
where B: Send, PhantomData<RespView>: Send,

§

impl<B, RespView> Sync for ServerStream<B, RespView>
where B: Sync, PhantomData<RespView>: Sync,

§

impl<B, RespView> Unpin for ServerStream<B, RespView>
where B: Unpin, PhantomData<RespView>: Unpin,

§

impl<B, RespView> UnsafeUnpin for ServerStream<B, RespView>
where B: UnsafeUnpin, PhantomData<RespView>: UnsafeUnpin,

Blanket Implementations§

Source§

impl<T> Any for T
where T: 'static + ?Sized,

Source§

fn type_id(&self) -> TypeId

Gets the TypeId of self. Read more
Source§

impl<T> Borrow<T> for T
where T: ?Sized,

Source§

fn borrow(&self) -> &T

Immutably borrows from an owned value. Read more
Source§

impl<T> BorrowMut<T> for T
where T: ?Sized,

Source§

fn borrow_mut(&mut self) -> &mut T

Mutably borrows from an owned value. Read more
Source§

impl<T> From<T> for T

Source§

fn from(t: T) -> T

Returns the argument unchanged.

Source§

impl<T> Instrument for T

Source§

fn instrument(self, span: Span) -> Instrumented<Self> ⓘ

Instruments this type with the provided Span, returning an Instrumented wrapper. Read more
Source§

fn in_current_span(self) -> Instrumented<Self> ⓘ

Instruments this type with the current Span, returning an Instrumented wrapper. Read more
Source§

impl<T, U> Into<U> for T
where U: From<T>,

Source§

fn into(self) -> U

Calls U::from(self).

That is, this conversion is whatever the implementation of From<T> for U chooses to do.

Source§

impl<T, U> TryFrom<U> for T
where U: Into<T>,

Source§

type Error = !

The type returned in the event of a conversion error.
Source§

fn try_from(value: U) -> Result<T, !>

Performs the conversion.
Source§

impl<T, U> TryInto<U> for T
where U: TryFrom<T>,

Source§

type Error = <U as TryFrom<T>>::Error

The type returned in the event of a conversion error.
Source§

fn try_into(self) -> Result<U, <U as TryFrom<T>>::Error>

Performs the conversion.
Source§

impl<T> WithSubscriber for T

Source§

fn with_subscriber<S>(self, subscriber: S) -> WithDispatch<Self> ⓘ
where S: Into<Dispatch>,

Attaches the provided Subscriber to this type, returning a WithDispatch wrapper. Read more
Source§

fn with_current_subscriber(self) -> WithDispatch<Self> ⓘ

Attaches the current default Subscriber to this type, returning a WithDispatch wrapper. Read more