A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
byte_chunking.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
25#ifndef A11_NET_BYTE_CHUNKING_H_
26#define A11_NET_BYTE_CHUNKING_H_
27
28#include <cstddef>
29#include <cstdint>
30#include <optional>
31#include <string>
32#include <string_view>
33#include <vector>
34
35#include <absl/base/nullability.h>
36#include <absl/status/status.h>
37#include <absl/status/statusor.h>
38
39namespace a11::net {
40
42enum class BytePacketType : std::uint8_t {
43 kCompleteBytes = 0x00,
44 kByteChunk = 0x01,
46};
47
49struct BytePacket {
51 std::string payload;
52 std::uint64_t transient_id = 0;
53 std::uint32_t sequence = 0;
54 std::uint32_t packet_count =
55 0;
56
57 template <typename Sink>
58 friend void AbslStringify(Sink& sink, const BytePacket& packet) {
59 sink.Append("BytePacket{type=");
60 sink.Append(packet.type == BytePacketType::kCompleteBytes ? "complete"
61 : packet.type == BytePacketType::kByteChunk ? "chunk"
62 : "first_chunk");
63 sink.Append(", id=");
64 sink.Append(packet.transient_id);
65 sink.Append(", sequence=");
66 sink.Append(packet.sequence);
67 sink.Append(", packet_count=");
68 sink.Append(packet.packet_count);
69 sink.Append(", payload_size=");
70 sink.Append(packet.payload.size());
71 sink.Append("}");
72 }
73};
74
77 size_t packet_size = 64 * 1024;
78 size_t max_message_size = 32 * 1024 * 1024;
81 64 * 1024 * 1024;
82
84 absl::Status Validate() const;
85};
86
88absl::StatusOr<std::vector<std::string>> SplitBytesIntoPackets(
89 std::string_view bytes, std::uint64_t transient_id, size_t packet_size);
96absl::StatusOr<std::vector<std::string>> SplitOwnedBytesIntoPackets(
97 std::string bytes, std::uint64_t transient_id, size_t packet_size);
99absl::StatusOr<BytePacket> ParseBytePacket(std::string_view packet);
105absl::StatusOr<BytePacket> ParseOwnedBytePacket(std::string packet);
106
116 public:
118 explicit ByteReassembler(ByteChunkingOptions options);
120
123
125 absl::StatusOr<std::optional<std::string>> Feed(std::string packet);
127 void Clear();
128
130 [[nodiscard]] size_t pending_message_count() const;
132 [[nodiscard]] size_t pending_byte_count() const;
133
134 private:
135 struct Impl;
136 static constexpr size_t kImplSize = 256;
137 static constexpr size_t kImplAlignment = alignof(std::max_align_t);
138
139 Impl* absl_nonnull GetImpl();
140 [[nodiscard]] const Impl* absl_nonnull GetImpl() const;
141
142 alignas(kImplAlignment) std::byte impl_[kImplSize];
143};
144
145} // namespace a11::net
146
147#endif // A11_NET_BYTE_CHUNKING_H_
Bounded, thread-safe reassembly for interleaved binary messages.
Definition byte_chunking.h:115
size_t pending_message_count() const
Number of message ids currently awaiting more packets.
Definition byte_chunking.cc:385
ByteReassembler(const ByteReassembler &)=delete
ByteReassembler & operator=(const ByteReassembler &)=delete
absl::StatusOr< std::optional< std::string > > Feed(std::string packet)
Admit one packet and return a complete message when this finishes one.
Definition byte_chunking.cc:274
size_t pending_byte_count() const
Aggregate payload bytes retained by incomplete messages.
Definition byte_chunking.cc:391
~ByteReassembler()
Definition byte_chunking.cc:262
void Clear()
Discard every incomplete message, for example when a channel aborts.
Definition byte_chunking.cc:378
Definition action.h:65
absl::StatusOr< BytePacket > ParseOwnedBytePacket(std::string packet)
Parse one packet, reusing its buffer as the payload.
Definition byte_chunking.cc:232
absl::StatusOr< std::vector< std::string > > SplitOwnedBytesIntoPackets(std::string bytes, std::uint64_t transient_id, size_t packet_size)
Split bytes the caller owns, reusing the buffer when it fits a packet.
Definition byte_chunking.cc:207
absl::StatusOr< BytePacket > ParseBytePacket(std::string_view packet)
Parse and validate one packet without retaining the input view.
Definition byte_chunking.cc:224
BytePacketType
Packet shapes in the A11 byte-chunking wire format.
Definition byte_chunking.h:42
absl::StatusOr< std::vector< std::string > > SplitBytesIntoPackets(std::string_view bytes, std::uint64_t transient_id, size_t packet_size)
Split bytes into A11 packets with fixed little-endian suffixes.
Definition byte_chunking.cc:115
Bounds packet size and incomplete-message memory during reassembly.
Definition byte_chunking.h:76
size_t packet_size
Maximum encoded packet size.
Definition byte_chunking.h:77
size_t max_pending_messages
Simultaneous incomplete messages.
Definition byte_chunking.h:79
size_t max_message_size
Reassembled message limit.
Definition byte_chunking.h:78
size_t max_pending_bytes
Aggregate pending payload limit.
Definition byte_chunking.h:80
absl::Status Validate() const
Validate that all limits can represent at least one useful packet.
Definition byte_chunking.cc:98
Parsed packet metadata plus the owned piece of application payload.
Definition byte_chunking.h:49
std::uint32_t sequence
Zero-based position in the message.
Definition byte_chunking.h:53
BytePacketType type
Packet shape.
Definition byte_chunking.h:50
std::uint32_t packet_count
Total count, when supplied by the first packet.
Definition byte_chunking.h:54
std::string payload
Bytes contributed by this packet.
Definition byte_chunking.h:51
std::uint64_t transient_id
Definition byte_chunking.h:52
friend void AbslStringify(Sink &sink, const BytePacket &packet)
Definition byte_chunking.h:58
Definition byte_chunking.cc:241