Skip to main content

WriteBlockHandle

Struct WriteBlockHandle 

Source
pub struct WriteBlockHandle {
    pub request_tx: Option<Sender<WriteRequest>>,
    /* private fields */
}
Expand description

Handle for an in-progress WriteBlock bidirectional streaming RPC.

The gRPC call runs in a background tokio task. The caller sends data through request_tx and receives responses via recv_response(). When done, call close() to drop the request channel and wait for the server to finalize.

Fields§

§request_tx: Option<Sender<WriteRequest>>

Sender for client → server WriteRequest messages (data chunks, flush commands).

Wrapped in Option so that close() can take() the sender (closing the client→server half of the stream) without violating the move semantics imposed by this type’s Drop impl. None after close() has run; senders attempting to use it should treat that as “stream already closed”.

Implementations§

Source§

impl WriteBlockHandle

Source

pub async fn recv_response(&mut self) -> Result<Option<WriteResponse>>

Receive the next WriteResponse from the server (e.g., flush ack).

Returns None if the server has closed the response stream.

Source

pub async fn close(self) -> Result<()>

Close the write stream by dropping the request sender and wait for any final response from the server.

Source

pub async fn cancel(self)

Cancel the write stream without waiting for server finalization.

Drops the request sender and response receiver immediately and aborts the background gRPC task so its resources are released promptly (rather than relying on the implicit “task exits because channels were dropped” behaviour, which leaves the JoinHandle detached on drop). Matches Java’s GrpcBlockingStream.cancel().

Trait Implementations§

Source§

impl Drop for WriteBlockHandle

Safety net: aborts the background gRPC task if the handle is dropped without going through close() / cancel().

Without this, an early ? return on the error path leaves a detached tokio task that can hang indefinitely on stream.message().await (e.g. on a half-open server connection that never sends a final response), keeping the underlying tonic Channel alive and leaking resources.

cancel() and close() already take() the task_handle, so on the happy path task_handle is None here and abort() is a no-op — matching the doc-comment on task_handle above.

Source§

fn drop(&mut self)

Executes the destructor for this type. Read more
Source§

fn pin_drop(self: Pin<&mut Self>)

🔬This is a nightly-only experimental API. (pin_ergonomics)
Execute the destructor for this type, but different to Drop::drop, it requires self to be pinned. Read more

Auto Trait Implementations§

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<ST, DT> CastableFrom<ST, Initialized, Initialized> for DT
where ST: ?Sized, DT: ?Sized,

Source§

impl<ST, DT> CastableFrom<ST, Uninit, Uninit> for DT
where ST: ?Sized, DT: ?Sized,

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> IntoRequest<T> for T

Source§

fn into_request(self) -> Request<T>

Wrap the input message T in a tonic::Request
Source§

impl<L> LayerExt<L> for L

Source§

fn named_layer<S>(&self, service: S) -> Layered<<L as Layer<S>>::Service, S>
where L: Layer<S>,

Applies the layer to a service and wraps it in Layered.
Source§

impl<T> Read<Exclusive, BecauseExclusive> for T
where T: ?Sized,

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<V, T> VZip<V> for T
where V: MultiLane<T>,

Source§

fn vzip(self) -> V

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