A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
channel_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
17#ifndef A11_NET_CHANNEL_WIRE_STREAM_H_
18#define A11_NET_CHANNEL_WIRE_STREAM_H_
19
20#include <cstddef>
21#include <functional>
22#include <memory>
23#include <optional>
24#include <string>
25
26#include <absl/base/nullability.h>
27#include <absl/status/status.h>
28#include <absl/status/statusor.h>
29#include <absl/time/time.h>
30
32#include "a11/data/types.h"
34#include "a11/net/wire_stream.h"
35
36namespace a11::net {
37
38namespace internal {
39class BinaryChannel;
40} // namespace internal
41
44
47 // Serialized byte packets never exceed split_size. Complete messages carry
48 // a small suffix; larger messages use length-suffixed first chunks followed
49 // by ordinary chunks in A11's byte-chunking protocol.
50 size_t split_size = 64 * 1024;
52 size_t max_pending_bytes = 64 * 1024 * 1024;
53
55 absl::Status Validate() const;
56};
57
68 : public WireStream,
69 public std::enable_shared_from_this<ChannelWireStream> {
70 public:
72 using OpenOperation = std::function<absl::Status()>;
73
76
77 ~ChannelWireStream() override;
78
79 absl::Status Send(data::WireMessage message) override;
80 a11::Task Start(OnMessage on_message, OnDone on_done) override;
81 a11::Task Accept(OnMessage on_message, OnDone on_done) override;
82 absl::Status HalfClose(data::ByteMap trailers) override;
84 absl::Status Abort(absl::Status status) override;
85 absl::Status SetDeadline(absl::Time deadline) override;
86
87 [[nodiscard]] absl::Time deadline() const override;
88 [[nodiscard]] absl::Status GetStatus() const override;
89 [[nodiscard]] std::optional<data::ByteMap> GetTrailers() const override;
90 [[nodiscard]] std::string GetId() const override;
91 [[nodiscard]] void* absl_nullable GetImpl() const override;
92
93 protected:
94 struct State;
95
96 explicit ChannelWireStream(std::shared_ptr<State> state)
97 : state_(std::move(state)) {}
98
99 static absl::StatusOr<std::shared_ptr<State>> MakeState(
100 std::shared_ptr<internal::BinaryChannel> channel, std::string id,
101 ChannelEndpointRole role, OpenOperation open_operation,
102 WireStreamOptions options, ChannelFramingOptions framing = {});
103
104 private:
105 a11::Task StartEndpoint(bool accept, OnMessage on_message, OnDone on_done);
106 static void Sender(const std::shared_ptr<State>& state);
109 static void DeliverClaimed(const std::shared_ptr<State>& state,
110 const data::WireMessage& message,
111 std::uint64_t message_id);
112 static void Receiver(const std::shared_ptr<State>& state);
113 static void WatchTiming(const std::shared_ptr<State>& state);
114 static void MarkActivity(const std::shared_ptr<State>& state);
115 static void ForceAbort(const std::shared_ptr<State>& state,
116 absl::Status status, bool can_communicate = true);
117 static void MaybeFinish(const std::shared_ptr<State>& state);
118 static void Finish(const std::shared_ptr<State>& state,
119 std::optional<absl::Status> terminal_error = std::nullopt);
120 static void Notify(const std::shared_ptr<State>& state);
121
122 std::shared_ptr<State> state_;
123};
124
125} // namespace a11::net
126
127#endif // A11_NET_CHANNEL_WIRE_STREAM_H_
Bounded packetisation and reassembly for binary channel transports.
Shared WireStream lifecycle and framing for binary channels.
Definition channel_wire_stream.h:69
std::optional< data::ByteMap > GetTrailers() const override
Definition channel_wire_stream.cc:547
ChannelWireStream(std::shared_ptr< State > state)
Definition channel_wire_stream.h:96
absl::Status GetStatus() const override
Definition channel_wire_stream.cc:532
absl::Status Abort(absl::Status status) override
Abort the stream, discarding buffered work.
Definition channel_wire_stream.cc:483
a11::Task Start(OnMessage on_message, OnDone on_done) override
Begin the stream as the initiating ("start") side.
Definition channel_wire_stream.cc:263
std::function< absl::Status()> OpenOperation
Transport-specific operation that initiates or accepts the channel.
Definition channel_wire_stream.h:72
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
~ChannelWireStream() override
Definition channel_wire_stream.cc:171
a11::Task DrainOutgoingMessages() override
Await delivery of all buffered outbound messages.
Definition channel_wire_stream.cc:471
std::string GetId() const override
Definition channel_wire_stream.cc:552
absl::Status Send(data::WireMessage message) override
Enqueue a message for delivery to the peer.
Definition channel_wire_stream.cc:189
void *absl_nullable GetImpl() const override
Definition channel_wire_stream.cc:557
a11::Task Accept(OnMessage on_message, OnDone on_done) override
Begin the stream as the accepting ("accept") side.
Definition channel_wire_stream.cc:267
absl::Time deadline() const override
Definition channel_wire_stream.cc:527
static absl::StatusOr< std::shared_ptr< State > > MakeState(std::shared_ptr< internal::BinaryChannel > channel, std::string id, ChannelEndpointRole role, OpenOperation open_operation, WireStreamOptions options, ChannelFramingOptions framing={})
Definition channel_wire_stream.cc:148
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
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
ChannelEndpointRole
Handshake role a binary channel endpoint is permitted to assume.
Definition channel_wire_stream.h:43
The top-level frame exchanged between two A11 endpoints.
Definition types.h:485
Packetisation and incomplete-message bounds for a binary WireStream.
Definition channel_wire_stream.h:46
size_t max_pending_bytes
Pending inbound byte limit.
Definition channel_wire_stream.h:52
size_t max_pending_messages
Incomplete inbound message limit.
Definition channel_wire_stream.h:51
size_t split_size
Maximum encoded channel packet size.
Definition channel_wire_stream.h:50
absl::Status Validate() const
Validate the framing limits before a transport starts.
Definition channel_wire_stream.cc:135
Definition channel_wire_stream.cc:57
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...