42#ifndef A11_NET_HTTP_SSE_WIRE_STREAM_H_
43#define A11_NET_HTTP_SSE_WIRE_STREAM_H_
51#include <absl/base/nullability.h>
52#include <absl/status/status.h>
53#include <absl/status/statusor.h>
54#include <absl/time/time.h>
93 "application/vnd.a11.wire-stream+json";
98 "/streams/{id}/message";
198 public std::enable_shared_from_this<HttpSseWireStream> {
210 absl::Status
Abort(absl::Status status)
override;
213 [[nodiscard]] absl::Time
deadline()
const override;
214 [[nodiscard]] absl::Status
GetStatus()
const override;
215 [[nodiscard]] std::optional<data::ByteMap>
GetTrailers()
const override;
216 [[nodiscard]] std::string
GetId()
const override;
217 [[nodiscard]]
void* absl_nullable
GetImpl()
const override;
248 std::shared_ptr<State> state);
257 void SetId(std::string
id);
259 [[nodiscard]] std::shared_ptr<InProcessWireStream>
bridge()
const;
269 std::shared_ptr<InProcessWireStream> application_;
270 std::shared_ptr<InProcessWireStream> bridge_;
271 std::shared_ptr<State> state_;
286 struct ConstructorToken {};
299 static absl::StatusOr<std::shared_ptr<HttpSseClientWireStream>>
Create(
301 std::shared_ptr<Http2Client>
client =
nullptr,
306 [[nodiscard]] std::shared_ptr<Http2Client>
client()
const;
320 std::shared_ptr<State> state,
321 std::shared_ptr<ClientState> client_state);
330 static void ReceiveSseLoop(
331 const std::shared_ptr<HttpSseClientWireStream>& self);
335 absl::Status OpenOutboundStream(
const HttpHeaders& response_headers);
337 absl::Status TransmitOnStream(
const std::string& payload,
bool terminal);
339 absl::Status TransmitAsPost(std::string payload);
341 absl::Status ClaimPostSlot();
343 void ReleasePostSlot();
345 absl::Status AwaitPostsDelivered();
347 std::shared_ptr<ClientState> client_state_;
362 struct ConstructorToken {};
372 std::shared_ptr<State> state,
373 std::shared_ptr<ServerStreamState> server_state);
382 std::shared_ptr<ServerStreamState> server_state_;
389 std::function<
a11::Task(std::shared_ptr<HttpSseServerWireStream>)>;
408 static absl::StatusOr<std::shared_ptr<HttpSseServer>>
Create(
409 std::string bind_address, std::uint16_t
port,
419 [[nodiscard]] std::uint16_t
port()
const;
421 [[nodiscard]]
bool running()
const;
423 [[nodiscard]] std::shared_ptr<Http2Server>
http2_server()
const;
429 : state_(std::move(state)) {}
431 static a11::Task HandleRequest(
const std::shared_ptr<State>& state,
433 std::shared_ptr<Http2ResponseWriter> response);
434 static a11::Task HandleConnect(
const std::shared_ptr<State>& state,
436 std::shared_ptr<Http2ResponseWriter> response);
437 static a11::Task HandleMessage(
const std::shared_ptr<State>& state,
439 std::shared_ptr<Http2ResponseWriter> response);
442 const std::shared_ptr<State>& state, std::string
stream_id,
443 HttpRequest request, std::shared_ptr<Http2ResponseWriter> response);
445 std::shared_ptr<State> state_;
std::string stream_id
Definition authorization.cc:328
The client-side HTTP SSE wire stream, dialing out to a server URL.
Definition http_sse_wire_stream.h:282
std::shared_ptr< Http2Client > client() const
Definition http_sse_wire_stream.cc:691
void TransportDone() override
Definition http_sse_wire_stream.cc:1146
SseOutboundDelivery outbound_delivery() const
The outbound delivery method actually in use.
Definition http_sse_wire_stream.cc:786
void *absl_nullable TransportImpl() const override
Definition http_sse_wire_stream.cc:1163
a11::Task OpenTransport() override
Definition http_sse_wire_stream.cc:696
static absl::StatusOr< std::shared_ptr< HttpSseClientWireStream > > Create(std::string url, HttpSseOptions options={}, std::shared_ptr< Http2Client > client=nullptr, HttpHeaders request_headers={})
Creates a client SSE wire stream connecting to url.
Definition http_sse_wire_stream.cc:635
absl::Status Transmit(data::WireMessage message) override
Hands one outbound message to the transport.
Definition http_sse_wire_stream.cc:894
The server-side HTTP SSE wire stream, accepted from a client.
Definition http_sse_wire_stream.h:358
a11::Task Accepted() const
Definition http_sse_wire_stream.cc:1192
void *absl_nullable TransportImpl() const override
Definition http_sse_wire_stream.cc:1253
void TransportDone() override
Definition http_sse_wire_stream.cc:1257
absl::Status Transmit(data::WireMessage message) override
Definition http_sse_wire_stream.cc:1237
a11::Task OpenTransport() override
Definition http_sse_wire_stream.cc:1196
Hosts HTTP SSE wire streams on top of an HTTP/2 server.
Definition http_sse_wire_stream.h:397
std::uint16_t port() const
Definition http_sse_wire_stream.cc:1777
absl::Status Stop()
Stops the server and releases its resources.
Definition http_sse_wire_stream.cc:1746
~HttpSseServer()
Definition http_sse_wire_stream.cc:1402
std::shared_ptr< Http2Server > http2_server() const
Definition http_sse_wire_stream.cc:1788
bool running() const
Definition http_sse_wire_stream.cc:1782
static absl::StatusOr< std::shared_ptr< HttpSseServer > > Create(std::string bind_address, std::uint16_t port, OnHttpSseConnect on_connect={}, HttpSseOptions options={})
Creates and starts an SSE server accepting A11 wire streams.
Definition http_sse_wire_stream.cc:1343
a11::Future< std::shared_ptr< HttpSseServerWireStream > > WaitForStream()
Definition http_sse_wire_stream.cc:1726
Common base for the client and server HTTP SSE wire streams.
Definition http_sse_wire_stream.h:198
Role
Definition http_sse_wire_stream.h:242
absl::Status Abort(absl::Status status) override
Abort the stream, discarding buffered work.
Definition http_sse_wire_stream.cc:475
std::string GetRequestPath() const
The path a server stream was accepted on, query string included.
Definition http_sse_wire_stream.cc:522
a11::Task Start(OnMessage on_message, OnDone on_done) override
Begin the stream as the initiating ("start") side.
Definition http_sse_wire_stream.cc:291
void FailTransport(absl::Status status)
Definition http_sse_wire_stream.cc:449
a11::Task ReceiveTransportMessage(data::WireMessage message)
Definition http_sse_wire_stream.cc:425
absl::Status SetHttpRequestHeaders(HttpHeaders headers)
Sets HTTP headers to send on the SSE request; call before connecting.
Definition http_sse_wire_stream.cc:532
virtual absl::Status Transmit(data::WireMessage message)=0
a11::Task Accept(OnMessage on_message, OnDone on_done) override
Begin the stream as the accepting ("accept") side.
Definition http_sse_wire_stream.cc:295
absl::Status SetDeadline()
Clear any deadline (equivalent to an infinite deadline).
Definition wire_stream.h:168
absl::Status GetStatus() const override
Definition http_sse_wire_stream.cc:493
absl::Status Send(data::WireMessage message) override
Enqueue a message for delivery to the peer.
Definition http_sse_wire_stream.cc:287
absl::Status HalfClose()
Half-close with no trailers.
Definition wire_stream.h:139
void *absl_nullable GetImpl() const override
Definition http_sse_wire_stream.cc:513
void MarkHttpHeadersReady(HttpHeaders headers)
Definition http_sse_wire_stream.cc:568
a11::Task StartEndpoint(bool accept, OnMessage on_message, OnDone on_done)
Definition http_sse_wire_stream.cc:299
HttpHeaders GetHttpRequestHeaders() const
Definition http_sse_wire_stream.cc:517
virtual void TransportDone()
Definition http_sse_wire_stream.h:265
std::string GetId() const override
Definition http_sse_wire_stream.cc:508
std::optional< HttpHeaders > GetHttpResponseHeaders() const
Definition http_sse_wire_stream.cc:527
a11::Task StartInternalBridge()
Definition http_sse_wire_stream.cc:346
absl::Status SetHttpResponseHeaders(HttpHeaders headers)
Sets HTTP headers to send on the SSE response (server side).
Definition http_sse_wire_stream.cc:548
a11::Task WaitForHttpHeaders() const
Definition http_sse_wire_stream.cc:564
absl::Time deadline() const override
Definition http_sse_wire_stream.cc:489
virtual void *absl_nullable TransportImpl() const =0
std::shared_ptr< InProcessWireStream > bridge() const
Definition http_sse_wire_stream.cc:593
void SetId(std::string id)
Definition http_sse_wire_stream.cc:583
virtual a11::Task OpenTransport()=0
std::optional< data::ByteMap > GetTrailers() const override
Definition http_sse_wire_stream.cc:504
a11::Task HandleBridgeMessage(std::optional< data::WireMessage > message)
Definition http_sse_wire_stream.cc:368
HttpSseOptions options() const
Definition http_sse_wire_stream.cc:588
a11::Task HandleBridgeDone()
Definition http_sse_wire_stream.cc:389
a11::Task DrainOutgoingMessages() override
Await delivery of all buffered outbound messages.
Definition http_sse_wire_stream.cc:471
~HttpSseWireStream() override=default
std::pair< std::shared_ptr< InProcessWireStream >, std::shared_ptr< InProcessWireStream > > Pair
A connected pair of endpoints returned by CreatePair().
Definition in_process_wire_stream.h:67
A bidirectional, message-oriented channel between two A11 endpoints.
Definition wire_stream.h:103
absl::Status SetDeadline()
Clear any deadline (equivalent to an infinite deadline).
Definition wire_stream.h:168
absl::Status HalfClose()
Half-close with no trailers.
Definition wire_stream.h:139
GET /actions, on whichever server happens to hold the port.
Completion values used by every asynchronous A11 operation.
A11's nghttp2 HTTP/2 client, server, and streaming primitives.
In-memory paired WireStream that connects two endpoints without a network.
absl::flat_hash_map< std::string, Bytes > ByteMap
String-keyed map of byte values (headers, attributes, etc.).
Definition types.h:58
std::function< a11::Task(std::optional< data::WireMessage > message)> OnMessage
Called for each inbound message; std::nullopt signals the peer half-closed.
Definition wire_stream.h:79
std::function< a11::Task()> OnDone
Called once, when the stream has fully finished (cleanly or via abort).
Definition wire_stream.h:81
std::vector< std::pair< std::string, std::string > > HttpHeaders
An ordered list of (name, value) HTTP/2 header fields.
Definition http2.h:70
constexpr std::string_view kSseOutboundPostToken
Token for one-POST-per-message outbound delivery.
Definition http_sse_wire_stream.h:78
constexpr std::string_view kSseOutboundStreamToken
Token for long-lived-request-body outbound delivery.
Definition http_sse_wire_stream.h:80
constexpr std::string_view kSseStreamIdHeader
Response header naming the stream id assigned to an SSE connection.
Definition http_sse_wire_stream.h:67
SseOutboundDelivery
How an SSE client hands its outbound WireMessages to the server.
Definition http_sse_wire_stream.h:109
@ kPost
One HTTP POST per message.
@ kStream
One long-lived request body carrying every outbound message.
constexpr std::string_view kSseOutboundModesHeader
Response header on the connect response listing the outbound delivery modes the server accepts,...
Definition http_sse_wire_stream.h:76
constexpr std::string_view kDefaultSseConnectEndpoint
Default path on which a client opens the SSE event stream.
Definition http_sse_wire_stream.h:95
constexpr std::string_view kSseHttpHeaderPrefix
Prefix under which application HTTP headers are tunneled over SSE.
Definition http_sse_wire_stream.h:69
constexpr std::string_view kDefaultSseMessageEndpoint
Default message-post path template ({id} is the stream id).
Definition http_sse_wire_stream.h:97
std::function< a11::Task(std::shared_ptr< HttpSseServerWireStream >)> OnHttpSseConnect
Callback invoked with each accepted server-side SSE wire stream.
Definition http_sse_wire_stream.h:389
constexpr std::string_view kSseWireStreamContentType
Content type of a streamed outbound request body.
Definition http_sse_wire_stream.h:92
Future< Unit > Task
Asynchronous operation whose only successful result is completion itself.
Definition future.h:411
The top-level frame exchanged between two A11 endpoints.
Definition types.h:485
Where and whether a server answers discovery over plain HTTP.
Definition describe_endpoint.h:72
Body-size limits, buffering thresholds, deadline, TLS, and HTTP protocol negotiation for an HTTP clie...
Definition http2.h:150
A parsed HTTP/2 request: pseudo-headers, fields, and body.
Definition http2.h:90
Definition http_sse_wire_stream.cc:597
Endpoint paths and transport tuning for an HTTP SSE wire stream.
Definition http_sse_wire_stream.h:141
WireStreamOptions stream_options
Per-stream buffering and deadline.
Definition http_sse_wire_stream.h:142
ServerHeaderOptions headers
Server-side response-header policy: the Server header, cross-origin access, and the hints that tell a...
Definition http_sse_wire_stream.h:182
std::string message_endpoint
Outbound endpoint template.
Definition http_sse_wire_stream.h:156
absl::Status Validate() const
Definition http_sse_wire_stream.cc:249
DescribeEndpointOptions describe
Server-side: GET /actions, when something above filled in the handler.
Definition http_sse_wire_stream.h:162
std::string connect_endpoint_prefix
Server-side: also accept a connect POST anywhere under this prefix.
Definition http_sse_wire_stream.h:155
SseOutboundDelivery outbound
Client-side outbound delivery method; servers accept either.
Definition http_sse_wire_stream.h:159
std::string connect_endpoint
SSE connect POST path.
Definition http_sse_wire_stream.h:144
Http2Options http2_options
Shared HTTP/2 transport and TLS policy.
Definition http_sse_wire_stream.h:143
size_t max_concurrent_posts
Outbound POSTs a client keeps in flight at once (kPost only).
Definition http_sse_wire_stream.h:174
bool accept_streamed_outbound
Server-side: whether a streamed outbound request body is accepted, and so advertised in kSseOutboundM...
Definition http_sse_wire_stream.h:170
Definition http_sse_wire_stream.cc:1171
Definition http_sse_wire_stream.cc:1269
Buffering, sizing, and deadline limits for a WireStream endpoint.
Definition wire_stream.h:58
A11's core wire value types: chunks, node fragments and messages.
A11's transport abstraction: the bidirectional WireStream channel and the options/callbacks that driv...