A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
http_sse_wire_stream.h
Go to the documentation of this file.
1/*
2 * Copyright 2026 The A11 Authors
3 *
4 * Licensed under the Apache License, Version 2.0 (the "License");
5 * you may not use this file except in compliance with the License.
6 * You may obtain a copy of the License at
7 *
8 * http://www.apache.org/licenses/LICENSE-2.0
9 *
10 * Unless required by applicable law or agreed to in writing, software
11 * distributed under the License is distributed on an "AS IS" BASIS,
12 * WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
13 * See the License for the specific language governing permissions and
14 * limitations under the License.
15 */
16
42#ifndef A11_NET_HTTP_SSE_WIRE_STREAM_H_
43#define A11_NET_HTTP_SSE_WIRE_STREAM_H_
44
45#include <cstdint>
46#include <functional>
47#include <memory>
48#include <optional>
49#include <string>
50
51#include <absl/base/nullability.h>
52#include <absl/status/status.h>
53#include <absl/status/statusor.h>
54#include <absl/time/time.h>
55
57#include "a11/data/types.h"
59#include "a11/net/http2.h"
62#include "a11/net/wire_stream.h"
63
64namespace a11::net {
65
67inline constexpr std::string_view kSseStreamIdHeader = "x-a11-stream-id";
69inline constexpr std::string_view kSseHttpHeaderPrefix = "x-a11-http-";
76inline constexpr std::string_view kSseOutboundModesHeader = "x-a11-outbound";
78inline constexpr std::string_view kSseOutboundPostToken = "post";
80inline constexpr std::string_view kSseOutboundStreamToken = "stream";
92inline constexpr std::string_view kSseWireStreamContentType =
93 "application/vnd.a11.wire-stream+json";
95inline constexpr std::string_view kDefaultSseConnectEndpoint = "/connect";
97inline constexpr std::string_view kDefaultSseMessageEndpoint =
98 "/streams/{id}/message";
99
118 kPost,
130 kStream,
131};
132
144 std::string connect_endpoint =
145 std::string(kDefaultSseConnectEndpoint);
156 std::string message_endpoint =
157 std::string(kDefaultSseMessageEndpoint);
163 // Server-side: whether a streamed outbound request body is accepted, and so
164 // advertised in kSseOutboundModesHeader.
175 // Server-side response-header policy: the `Server` header, cross-origin
176 // access, and the hints that tell a client and its intermediaries how to
177 // treat a reply.
183
185 absl::Status Validate() const;
186};
187
197 : public WireStream,
198 public std::enable_shared_from_this<HttpSseWireStream> {
199 public:
200 ~HttpSseWireStream() override = default;
201
204
205 absl::Status Send(data::WireMessage message) override;
206 a11::Task Start(OnMessage on_message, OnDone on_done) override;
207 a11::Task Accept(OnMessage on_message, OnDone on_done) override;
208 absl::Status HalfClose(data::ByteMap trailers) override;
210 absl::Status Abort(absl::Status status) override;
211 absl::Status SetDeadline(absl::Time deadline) override;
212
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;
218
220 [[nodiscard]] HttpHeaders GetHttpRequestHeaders() const;
221
229 [[nodiscard]] std::string GetRequestPath() const;
232 [[nodiscard]] std::optional<HttpHeaders> GetHttpResponseHeaders() const;
234 absl::Status SetHttpRequestHeaders(HttpHeaders headers);
236 absl::Status SetHttpResponseHeaders(HttpHeaders headers);
240
241 protected:
242 enum class Role { kClient, kServer };
243 struct State;
244
245 HttpSseWireStream(Role role, const std::string& id,
246 const HttpSseOptions& options,
248 std::shared_ptr<State> state);
249
250 a11::Task StartEndpoint(bool accept, OnMessage on_message, OnDone on_done);
252 a11::Task HandleBridgeMessage(std::optional<data::WireMessage> message);
255 void FailTransport(absl::Status status);
256 void MarkHttpHeadersReady(HttpHeaders headers);
257 void SetId(std::string id);
258 [[nodiscard]] HttpSseOptions options() const;
259 [[nodiscard]] std::shared_ptr<InProcessWireStream> bridge() const;
260
262 virtual absl::Status Transmit(data::WireMessage message) = 0;
263 virtual void* absl_nullable TransportImpl() const = 0;
264
265 virtual void TransportDone() {}
266
267 private:
268 Role role_;
269 std::shared_ptr<InProcessWireStream> application_;
270 std::shared_ptr<InProcessWireStream> bridge_;
271 std::shared_ptr<State> state_;
272
273 friend class HttpSseServer;
274};
275
283 private:
284 struct ClientState;
285
286 struct ConstructorToken {};
287
288 public:
299 static absl::StatusOr<std::shared_ptr<HttpSseClientWireStream>> Create(
300 std::string url, HttpSseOptions options = {},
301 std::shared_ptr<Http2Client> client = nullptr,
302 HttpHeaders request_headers = {});
303
306 [[nodiscard]] std::shared_ptr<Http2Client> client() const;
307
315 [[nodiscard]] SseOutboundDelivery outbound_delivery() const;
316
317 HttpSseClientWireStream(ConstructorToken, const std::string& url,
318 const HttpSseOptions& options,
320 std::shared_ptr<State> state,
321 std::shared_ptr<ClientState> client_state);
322
323 protected:
324 a11::Task OpenTransport() override;
325 absl::Status Transmit(data::WireMessage message) override;
326 void* absl_nullable TransportImpl() const override;
327 void TransportDone() override;
328
329 private:
330 static void ReceiveSseLoop(
331 const std::shared_ptr<HttpSseClientWireStream>& self);
332
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();
346
347 std::shared_ptr<ClientState> client_state_;
348};
349
350class HttpSseServer;
351
359 private:
360 struct ServerStreamState;
361
362 struct ConstructorToken {};
363
364 public:
367 a11::Task Accepted() const;
368
369 HttpSseServerWireStream(ConstructorToken, const std::string& id,
370 const HttpSseOptions& options,
372 std::shared_ptr<State> state,
373 std::shared_ptr<ServerStreamState> server_state);
374
375 protected:
376 a11::Task OpenTransport() override;
377 absl::Status Transmit(data::WireMessage message) override;
378 void* absl_nullable TransportImpl() const override;
379 void TransportDone() override;
380
381 private:
382 std::shared_ptr<ServerStreamState> server_state_;
383
384 friend class HttpSseServer;
385};
386
389 std::function<a11::Task(std::shared_ptr<HttpSseServerWireStream>)>;
390
397class HttpSseServer : public std::enable_shared_from_this<HttpSseServer> {
398 public:
408 static absl::StatusOr<std::shared_ptr<HttpSseServer>> Create(
409 std::string bind_address, std::uint16_t port,
410 OnHttpSseConnect on_connect = {}, HttpSseOptions options = {});
411
413
417 absl::Status Stop();
419 [[nodiscard]] std::uint16_t port() const;
421 [[nodiscard]] bool running() const;
423 [[nodiscard]] std::shared_ptr<Http2Server> http2_server() const;
424
425 private:
426 struct State;
427
428 explicit HttpSseServer(std::shared_ptr<State> state)
429 : state_(std::move(state)) {}
430
431 static a11::Task HandleRequest(const std::shared_ptr<State>& state,
432 HttpRequest request,
433 std::shared_ptr<Http2ResponseWriter> response);
434 static a11::Task HandleConnect(const std::shared_ptr<State>& state,
435 HttpRequest request,
436 std::shared_ptr<Http2ResponseWriter> response);
437 static a11::Task HandleMessage(const std::shared_ptr<State>& state,
438 std::string stream_id, HttpRequest request,
439 std::shared_ptr<Http2ResponseWriter> response);
441 static a11::Task HandleMessageStream(
442 const std::shared_ptr<State>& state, std::string stream_id,
443 HttpRequest request, std::shared_ptr<Http2ResponseWriter> response);
444
445 std::shared_ptr<State> state_;
446
448};
449
452
453} // namespace a11::net
454
455#endif // A11_NET_HTTP_SSE_WIRE_STREAM_H_
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
Definition action.h:65
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 response headers every A11 HTTP server sends, in one place.
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
The response-header policy of one A11 HTTP surface.
Definition server_headers.h:112
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...