A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
wire_stream_with_recv.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_WIRE_STREAM_WITH_RECV_H_
18#define A11_NET_WIRE_STREAM_WITH_RECV_H_
19
20#include <memory>
21#include <optional>
22#include <string>
23
24#include <absl/base/nullability.h>
25#include <absl/status/status.h>
26#include <absl/status/statusor.h>
27#include <absl/time/time.h>
28
30#include "a11/data/types.h"
31#include "a11/net/wire_stream.h"
32
33namespace a11::net {
34
44 : public WireStream,
45 public std::enable_shared_from_this<WireStreamWithRecv> {
46 private:
47 struct State;
48
49 struct ConstructorToken {};
50
51 public:
53 static absl::StatusOr<std::shared_ptr<WireStreamWithRecv>> Create(
54 std::shared_ptr<WireStream> stream);
55
58
59 absl::Status Send(data::WireMessage message) override;
60 a11::Task Start(OnMessage on_message, OnDone on_done) override;
61 a11::Task Accept(OnMessage on_message, OnDone on_done) override;
66 absl::Status HalfClose(data::ByteMap trailers) override;
68 absl::Status Abort(absl::Status status) override;
69 absl::Status SetDeadline(absl::Time deadline) override;
70
71 [[nodiscard]] absl::Time deadline() const override;
72 [[nodiscard]] absl::Status GetStatus() const override;
73 [[nodiscard]] std::optional<data::ByteMap> GetTrailers() const override;
74 [[nodiscard]] std::string GetId() const override;
75 [[nodiscard]] void* absl_nullable GetImpl() const override;
76
82 absl::Duration timeout = absl::InfiniteDuration());
83
85 [[nodiscard]] std::shared_ptr<WireStream> wrapped_stream() const {
86 return stream_;
87 }
88
89 WireStreamWithRecv(ConstructorToken, std::shared_ptr<WireStream> stream,
90 std::string id, std::shared_ptr<State> state)
91 : stream_(std::move(stream)),
92 id_(std::move(id)),
93 state_(std::move(state)) {}
94
95 private:
96 a11::Task StartImpl(bool accept, OnMessage on_message, OnDone on_done);
97 a11::Task HandleMessage(OnMessage observer,
98 std::optional<data::WireMessage> message);
99 a11::Task HandleDone(OnDone observer);
100 void RecordCurrentStatus() const;
101 void SignalError(absl::Status status) const;
102
103 std::shared_ptr<WireStream> stream_;
104 std::string id_;
105 std::shared_ptr<State> state_;
106};
107
108} // namespace a11::net
109
110#endif // A11_NET_WIRE_STREAM_WITH_RECV_H_
Pull-oriented adapter for a callback-driven WireStream.
Definition wire_stream_with_recv.h:45
absl::Status SetDeadline()
Clear any deadline (equivalent to an infinite deadline).
Definition wire_stream.h:168
a11::Task DrainOutgoingMessages() override
Await delivery of all buffered outbound messages.
Definition wire_stream_with_recv.cc:116
absl::Status HalfClose()
Half-close with no trailers.
Definition wire_stream.h:139
std::shared_ptr< WireStream > wrapped_stream() const
Return the underlying callback-oriented stream.
Definition wire_stream_with_recv.h:85
absl::Status GetStatus() const override
Definition wire_stream_with_recv.cc:138
absl::Status Abort(absl::Status status) override
Abort the stream, discarding buffered work.
Definition wire_stream_with_recv.cc:120
a11::Task Accept()
Accept the wrapped stream and route inbound messages to Receive().
Definition wire_stream_with_recv.cc:87
std::optional< data::ByteMap > GetTrailers() const override
Definition wire_stream_with_recv.cc:142
void *absl_nullable GetImpl() const override
Definition wire_stream_with_recv.cc:150
a11::Task Start()
Start the wrapped stream and route inbound messages to Receive().
Definition wire_stream_with_recv.cc:83
absl::Time deadline() const override
Definition wire_stream_with_recv.cc:134
WireStreamWithRecv(ConstructorToken, std::shared_ptr< WireStream > stream, std::string id, std::shared_ptr< State > state)
Definition wire_stream_with_recv.h:89
std::string GetId() const override
Definition wire_stream_with_recv.cc:146
absl::Status Send(data::WireMessage message) override
Enqueue a message for delivery to the peer.
Definition wire_stream_with_recv.cc:71
static absl::StatusOr< std::shared_ptr< WireStreamWithRecv > > Create(std::shared_ptr< WireStream > stream)
Wrap stream, preserving its id, lifecycle, and transport handle.
Definition wire_stream_with_recv.cc:59
a11::Future< std::optional< data::WireMessage > > Receive(absl::Duration timeout=absl::InfiniteDuration())
Await one inbound message, or nullopt after the peer half-closes.
Definition wire_stream_with_recv.cc:154
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.
std::optional< absl::Duration > timeout
Definition main.cc:144
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 wire_stream_with_recv.cc:48
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...