A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
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
29#ifndef A11_NET_WIRE_STREAM_H_
30#define A11_NET_WIRE_STREAM_H_
31
32#include <cstddef>
33#include <functional>
34#include <optional>
35#include <string>
36
37#include <absl/base/nullability.h>
38#include <absl/status/status.h>
39#include <absl/time/time.h>
40
42#include "a11/data/types.h"
43
44namespace a11::net {
45
47inline constexpr std::string_view kAbortStatusHeader = "x-a11-abort-status";
49inline constexpr size_t kMaxSingleMessageSize = 32 * 1024 * 1024;
50
64 size_t max_buffered_incoming_bytes = 32 * 1024 * 1024;
67 absl::Duration message_timeout = absl::InfiniteDuration();
69 absl::Time deadline = absl::InfiniteFuture();
70
72 absl::Status Validate() const;
73};
74
78using OnMessage =
79 std::function<a11::Task(std::optional<data::WireMessage> message)>;
81using OnDone = std::function<a11::Task()>;
82
104 public:
105 virtual ~WireStream() = default;
106
118 virtual absl::Status Send(data::WireMessage message) = 0;
119
127 virtual a11::Task Start(OnMessage on_message, OnDone on_done) = 0;
128
136 virtual a11::Task Accept(OnMessage on_message, OnDone on_done) = 0;
137
139 absl::Status HalfClose() { return HalfClose(data::ByteMap{}); }
140
151 virtual absl::Status HalfClose(data::ByteMap trailers) = 0;
152
159
165 virtual absl::Status Abort(absl::Status status) = 0;
166
168 absl::Status SetDeadline() { return SetDeadline(absl::InfiniteFuture()); }
169
174 virtual absl::Status SetDeadline(absl::Time deadline) = 0;
175
177 [[nodiscard]] virtual absl::Time deadline() const = 0;
179 [[nodiscard]] virtual absl::Status GetStatus() const = 0;
181 [[nodiscard]] virtual std::optional<data::ByteMap> GetTrailers() const = 0;
183 [[nodiscard]] virtual std::string GetId() const = 0;
186 [[nodiscard]] virtual void* absl_nullable GetImpl() const = 0;
187};
188
191absl::StatusOr<data::ByteMap> NormalizeWireHeaders(data::ByteMap headers);
192
193} // namespace a11::net
194
195#endif // A11_NET_WIRE_STREAM_H_
A bidirectional, message-oriented channel between two A11 endpoints.
Definition wire_stream.h:103
virtual std::optional< data::ByteMap > GetTrailers() const =0
virtual ~WireStream()=default
virtual a11::Task Accept(OnMessage on_message, OnDone on_done)=0
Begin the stream as the accepting ("accept") side.
virtual void *absl_nullable GetImpl() const =0
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
virtual a11::Task DrainOutgoingMessages()=0
Await delivery of all buffered outbound messages.
virtual absl::Status Send(data::WireMessage message)=0
Enqueue a message for delivery to the peer.
virtual absl::Status SetDeadline(absl::Time deadline)=0
Set the absolute deadline after which the stream is aborted.
virtual absl::Status HalfClose(data::ByteMap trailers)=0
Signal that this side will send no more messages.
virtual absl::Time deadline() const =0
virtual a11::Task Start(OnMessage on_message, OnDone on_done)=0
Begin the stream as the initiating ("start") side.
virtual absl::Status GetStatus() const =0
virtual absl::Status Abort(absl::Status status)=0
Abort the stream, discarding buffered work.
virtual std::string GetId() const =0
Completion values used by every asynchronous A11 operation.
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
constexpr size_t kMaxSingleMessageSize
Hard ceiling on the size of a single reassembled inbound WireMessage.
Definition wire_stream.h:49
absl::StatusOr< data::ByteMap > NormalizeWireHeaders(data::ByteMap headers)
Validate and case-normalise a wire header map, returning the normalised copy or a non-OK status if a ...
Definition wire_stream.cc:53
constexpr std::string_view kAbortStatusHeader
Trailer key under which an aborting endpoint reports its terminal status.
Definition wire_stream.h:47
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
Buffering, sizing, and deadline limits for a WireStream endpoint.
Definition wire_stream.h:58
size_t max_single_message_size
Reject any single reassembled message larger than this.
Definition wire_stream.h:62
size_t max_buffered_incoming_messages
Maximum inbound messages buffered before backpressure is applied.
Definition wire_stream.h:60
absl::Duration message_timeout
Fail a message that cannot be delivered within this duration (infinite = no per-message timeout).
Definition wire_stream.h:67
absl::Status Validate() const
Validate the option values, returning a non-OK status if inconsistent.
Definition wire_stream.cc:29
size_t max_buffered_incoming_bytes
Maximum total bytes of buffered inbound messages.
Definition wire_stream.h:64
absl::Time deadline
Absolute deadline after which the stream is aborted.
Definition wire_stream.h:69
A11's core wire value types: chunks, node fragments and messages.