A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
in_process_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
31#ifndef A11_NET_IN_PROCESS_WIRE_STREAM_H_
32#define A11_NET_IN_PROCESS_WIRE_STREAM_H_
33
34#include <memory>
35#include <optional>
36#include <string>
37#include <utility>
38
39#include <absl/base/nullability.h>
40#include <absl/status/status.h>
41#include <absl/status/statusor.h>
42#include <absl/time/time.h>
43
45#include "a11/data/types.h"
46#include "a11/net/wire_stream.h"
47
48namespace a11::net {
49
58class InProcessWireStream final : public WireStream {
59 private:
60 struct State;
61
62 struct ConstructorToken {};
63
64 public:
66 using Pair = std::pair<std::shared_ptr<InProcessWireStream>,
67 std::shared_ptr<InProcessWireStream>>;
68
81 static absl::StatusOr<Pair> CreatePair(
82 std::optional<WireStreamOptions> options = std::nullopt,
83 std::optional<WireStreamOptions> first_options = std::nullopt,
84 std::optional<WireStreamOptions> second_options = std::nullopt,
85 std::string preassigned_id = {});
86
87 ~InProcessWireStream() override = default;
88
91
92 absl::Status Send(data::WireMessage message) override;
93 a11::Task Start(OnMessage on_message, OnDone on_done) override;
94 a11::Task Accept(OnMessage on_message, OnDone on_done) override;
95 absl::Status HalfClose(data::ByteMap trailers) override;
97 absl::Status Abort(absl::Status status) override;
98 absl::Status SetDeadline(absl::Time deadline) override;
99
102 [[nodiscard]] a11::Task Done() const;
103
104 [[nodiscard]] absl::Time deadline() const override;
105 [[nodiscard]] absl::Status GetStatus() const override;
106 [[nodiscard]] std::optional<data::ByteMap> GetTrailers() const override;
107 [[nodiscard]] std::string GetId() const override;
108 [[nodiscard]] void* absl_nullable GetImpl() const override;
109
110 explicit InProcessWireStream(ConstructorToken, std::shared_ptr<State> state)
111 : state_(std::move(state)) {}
112
113 private:
114 a11::Task StartEndpoint(OnMessage on_message, OnDone on_done);
115 static void Sender(const std::shared_ptr<State>& state);
118 static void DeliverClaimed(const std::shared_ptr<State>& state,
119 data::WireMessage message);
120 static void Receiver(const std::shared_ptr<State>& state);
121 static void WatchTiming(const std::shared_ptr<State>& state);
122 static void MarkActivity(const std::shared_ptr<State>& first,
123 const std::shared_ptr<State>& second);
124 static bool ForceAbort(const std::shared_ptr<State>& state,
125 absl::Status status);
126 static void MaybeFinish(const std::shared_ptr<State>& state);
127 static void Finish(const std::shared_ptr<State>& state);
128
129 std::shared_ptr<State> state_;
130};
131
132} // namespace a11::net
133
134#endif // A11_NET_IN_PROCESS_WIRE_STREAM_H_
A WireStream endpoint wired directly to a peer in the same process.
Definition in_process_wire_stream.h:58
~InProcessWireStream() override=default
absl::Status Send(data::WireMessage message) override
Enqueue a message for delivery to the peer.
Definition in_process_wire_stream.cc:167
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
absl::Status Abort(absl::Status status) override
Abort the stream, discarding buffered work.
Definition in_process_wire_stream.cc:369
std::string GetId() const override
Definition in_process_wire_stream.cc:442
a11::Task Done() const
Definition in_process_wire_stream.cc:413
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
static absl::StatusOr< Pair > CreatePair(std::optional< WireStreamOptions > options=std::nullopt, std::optional< WireStreamOptions > first_options=std::nullopt, std::optional< WireStreamOptions > second_options=std::nullopt, std::string preassigned_id={})
Creates a connected pair of in-process endpoints.
Definition in_process_wire_stream.cc:128
a11::Task Accept(OnMessage on_message, OnDone on_done) override
Begin the stream as the accepting ("accept") side.
Definition in_process_wire_stream.cc:287
a11::Task DrainOutgoingMessages() override
Await delivery of all buffered outbound messages.
Definition in_process_wire_stream.cc:353
std::optional< data::ByteMap > GetTrailers() const override
Definition in_process_wire_stream.cc:437
absl::Time deadline() const override
Definition in_process_wire_stream.cc:417
void *absl_nullable GetImpl() const override
Definition in_process_wire_stream.cc:446
a11::Task Start(OnMessage on_message, OnDone on_done) override
Begin the stream as the initiating ("start") side.
Definition in_process_wire_stream.cc:283
InProcessWireStream(ConstructorToken, std::shared_ptr< State > state)
Definition in_process_wire_stream.h:110
absl::Status GetStatus() const override
Definition in_process_wire_stream.cc:422
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
The top-level frame exchanged between two A11 endpoints.
Definition types.h:485
Definition in_process_wire_stream.cc:69
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...