A11 (C++ runtime)
Native C++ implementation of the A11 action and streaming runtime
Loading...
Searching...
No Matches
a11::net::InProcessWireStream Class Referencefinal

A WireStream endpoint wired directly to a peer in the same process. More...

#include <cpp/a11/net/in_process_wire_stream.h>

Inheritance diagram for a11::net::InProcessWireStream:
[legend]

Classes

struct  State
 

Public Types

using Pair = std::pair< std::shared_ptr< InProcessWireStream >, std::shared_ptr< InProcessWireStream > >
 A connected pair of endpoints returned by CreatePair().
 

Public Member Functions

 ~InProcessWireStream () override=default
 
absl::Status Send (data::WireMessage message) override
 Enqueue a message for delivery to the peer.
 
a11::Task Start (OnMessage on_message, OnDone on_done) override
 Begin the stream as the initiating ("start") side.
 
a11::Task Accept (OnMessage on_message, OnDone on_done) override
 Begin the stream as the accepting ("accept") side.
 
absl::Status HalfClose (data::ByteMap trailers) override
 Signal that this side will send no more messages.
 
a11::Task DrainOutgoingMessages () override
 Await delivery of all buffered outbound messages.
 
absl::Status Abort (absl::Status status) override
 Abort the stream, discarding buffered work.
 
absl::Status SetDeadline (absl::Time deadline) override
 Set the absolute deadline after which the stream is aborted.
 
a11::Task Done () const
 
absl::Time deadline () const override
 
absl::Status GetStatus () const override
 
std::optional< data::ByteMapGetTrailers () const override
 
std::string GetId () const override
 
void *absl_nullable GetImpl () const override
 
 InProcessWireStream (ConstructorToken, std::shared_ptr< State > state)
 
absl::Status HalfClose ()
 Half-close with no trailers.
 
virtual absl::Status HalfClose (data::ByteMap trailers)=0
 Signal that this side will send no more messages.
 
absl::Status SetDeadline ()
 Clear any deadline (equivalent to an infinite deadline).
 
virtual absl::Status SetDeadline (absl::Time deadline)=0
 Set the absolute deadline after which the stream is aborted.
 
- Public Member Functions inherited from a11::net::WireStream
virtual ~WireStream ()=default
 
absl::Status HalfClose ()
 Half-close with no trailers.
 
absl::Status SetDeadline ()
 Clear any deadline (equivalent to an infinite deadline).
 

Static Public Member Functions

static absl::StatusOr< PairCreatePair (std::optional< WireStreamOptions > options=std::nullopt, std::optional< WireStreamOptions > first_options=std::nullopt, std::optional< WireStreamOptions > second_options=std::nullopt, std::string preassigned_id={})
 Creates a connected pair of in-process endpoints.
 

Detailed Description

A WireStream endpoint wired directly to a peer in the same process.

The two endpoints share in-memory queues, so a Send() on one is delivered to the OnMessage callback of the other with no network hop. One endpoint drives Start() and the other Accept(); either may HalfClose() or Abort(). Construct a connected pair with CreatePair() rather than instantiating directly.

Member Typedef Documentation

◆ Pair

using a11::net::InProcessWireStream::Pair = std::pair<std::shared_ptr<InProcessWireStream>, std::shared_ptr<InProcessWireStream> >

A connected pair of endpoints returned by CreatePair().

Constructor & Destructor Documentation

◆ ~InProcessWireStream()

a11::net::InProcessWireStream::~InProcessWireStream ( )
overridedefault

◆ InProcessWireStream()

a11::net::InProcessWireStream::InProcessWireStream ( ConstructorToken  ,
std::shared_ptr< State state 
)
inlineexplicit

Member Function Documentation

◆ Abort()

absl::Status a11::net::InProcessWireStream::Abort ( absl::Status  status)
overridevirtual

Abort the stream, discarding buffered work.

Parameters
statusA non-OK status reported to the peer as the abort reason.
Returns
OK if the abort was accepted.

Implements a11::net::WireStream.

◆ Accept()

a11::Task a11::net::InProcessWireStream::Accept ( OnMessage  on_message,
OnDone  on_done 
)
overridevirtual

Begin the stream as the accepting ("accept") side.

Parameters
on_messageInvoked for each inbound message (nullopt = peer half-closed).
on_doneInvoked once when the stream has finished.
Returns
An awaitable that resolves once the stream is accepted.

Implements a11::net::WireStream.

◆ CreatePair()

absl::StatusOr< InProcessWireStream::Pair > a11::net::InProcessWireStream::CreatePair ( std::optional< WireStreamOptions options = std::nullopt,
std::optional< WireStreamOptions first_options = std::nullopt,
std::optional< WireStreamOptions second_options = std::nullopt,
std::string  preassigned_id = {} 
)
static

Creates a connected pair of in-process endpoints.

Parameters
optionsShared WireStreamOptions applied to both endpoints.
first_optionsOptional overrides for the first endpoint.
second_optionsOptional overrides for the second endpoint.
preassigned_idWhen non-empty, fixes the shared stream id (and hence the tracing trace id) instead of generating one. This is the implementation-level hook for preassigning a stream's trace id without widening the WireStream interface.
Returns
The two connected endpoints, or an error status.

◆ deadline()

absl::Time a11::net::InProcessWireStream::deadline ( ) const
overridevirtual
Returns
The stream's current deadline.

Implements a11::net::WireStream.

◆ Done()

a11::Task a11::net::InProcessWireStream::Done ( ) const
Returns
An awaitable that resolves when this endpoint has fully finished (its peer half-closed or the stream aborted).

◆ DrainOutgoingMessages()

a11::Task a11::net::InProcessWireStream::DrainOutgoingMessages ( )
overridevirtual

Await delivery of all buffered outbound messages.

Returns
An awaitable that resolves once this endpoint's queued messages have been flushed to the transport (requires a prior HalfClose).

Implements a11::net::WireStream.

◆ GetId()

std::string a11::net::InProcessWireStream::GetId ( ) const
overridevirtual
Returns
This stream's transport-assigned identifier.

Implements a11::net::WireStream.

◆ GetImpl()

void *absl_nullable a11::net::InProcessWireStream::GetImpl ( ) const
overridevirtual
Returns
An opaque handle to the underlying transport object (advanced; may be null).

Implements a11::net::WireStream.

◆ GetStatus()

absl::Status a11::net::InProcessWireStream::GetStatus ( ) const
overridevirtual
Returns
The stream's terminal status (OK unless it failed or was aborted).

Implements a11::net::WireStream.

◆ GetTrailers()

std::optional< data::ByteMap > a11::net::InProcessWireStream::GetTrailers ( ) const
overridevirtual
Returns
The peer's closing trailers, if the stream has received them.

Implements a11::net::WireStream.

◆ HalfClose() [1/3]

absl::Status a11::net::WireStream::HalfClose ( )
inline

Half-close with no trailers.

See also
HalfClose(data::ByteMap)

◆ HalfClose() [2/3]

absl::Status a11::net::InProcessWireStream::HalfClose ( data::ByteMap  trailers)
overridevirtual

Signal that this side will send no more messages.

The terminal marker is queued after buffered outbound messages and, once the peer also half-closes, the stream completes. Inbound messages continue to be delivered until then. Call DrainOutgoingMessages() to await local transport delivery.

Parameters
trailersOptional closing metadata delivered to the peer.
Returns
OK, or a non-OK status if the stream cannot be half-closed.

Implements a11::net::WireStream.

◆ HalfClose() [3/3]

virtual absl::Status a11::net::WireStream::HalfClose ( data::ByteMap  trailers)
virtual

Signal that this side will send no more messages.

The terminal marker is queued after buffered outbound messages and, once the peer also half-closes, the stream completes. Inbound messages continue to be delivered until then. Call DrainOutgoingMessages() to await local transport delivery.

Parameters
trailersOptional closing metadata delivered to the peer.
Returns
OK, or a non-OK status if the stream cannot be half-closed.

Implements a11::net::WireStream.

◆ Send()

absl::Status a11::net::InProcessWireStream::Send ( data::WireMessage  message)
overridevirtual

Enqueue a message for delivery to the peer.

Non-blocking.

The message is admitted to this endpoint's outbound queue and delivered asynchronously by the transport task, which is also where backpressure is applied; there is no delivery-order guarantee across messages (see the class comment).

Parameters
messageThe message to send.
Returns
OK once the message is queued, or a non-OK status if the stream is not writable (e.g. already half-closed or aborted).

Implements a11::net::WireStream.

◆ SetDeadline() [1/3]

absl::Status a11::net::WireStream::SetDeadline ( )
inline

Clear any deadline (equivalent to an infinite deadline).

◆ SetDeadline() [2/3]

absl::Status a11::net::InProcessWireStream::SetDeadline ( absl::Time  deadline)
overridevirtual

Set the absolute deadline after which the stream is aborted.

Parameters
deadlineThe deadline; absl::InfiniteFuture() disables it.

Implements a11::net::WireStream.

◆ SetDeadline() [3/3]

virtual absl::Status a11::net::WireStream::SetDeadline ( absl::Time  deadline)
virtual

Set the absolute deadline after which the stream is aborted.

Parameters
deadlineThe deadline; absl::InfiniteFuture() disables it.

Implements a11::net::WireStream.

◆ Start()

a11::Task a11::net::InProcessWireStream::Start ( OnMessage  on_message,
OnDone  on_done 
)
overridevirtual

Begin the stream as the initiating ("start") side.

Parameters
on_messageInvoked for each inbound message (nullopt = peer half-closed).
on_doneInvoked once when the stream has finished.
Returns
An awaitable that resolves once the startup handshake completes.

Implements a11::net::WireStream.


The documentation for this class was generated from the following files: