A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
session.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
33#ifndef A11_SERVICE_SESSION_H_
34#define A11_SERVICE_SESSION_H_
35
36#include <cstddef>
37#include <functional>
38#include <memory>
39#include <optional>
40#include <string>
41#include <string_view>
42#include <utility>
43#include <vector>
44
45#include <absl/status/status.h>
46#include <absl/status/statusor.h>
47#include <absl/time/time.h>
48
49#include "a11/actions/action.h"
52#include "a11/data/types.h"
53#include "a11/net/wire_stream.h"
54#include "a11/nodes/node_map.h"
55
56namespace a11::service {
57
59inline constexpr std::string_view kSessionStatusHeader = "x-a11-session-status";
61inline constexpr size_t kMaxSingleMessageSize = 32 * 1024 * 1024;
62
84 size_t max_buffered_bytes_total = 32 * 1024 * 1024;
86 size_t max_buffered_bytes_per_stream = 4 * 1024 * 1024;
88 absl::Duration no_stream_timeout = absl::Seconds(30);
90 absl::Time deadline = absl::InfiniteFuture();
91
96 absl::Status Validate() const;
97};
98
99class Session;
100
103using OnSessionStreamMessage = std::function<a11::Task(
104 std::optional<data::WireMessage>, std::shared_ptr<net::WireStream>,
105 std::shared_ptr<Session>)>;
107using OnSessionStreamDone = std::function<a11::Task(
108 std::shared_ptr<net::WireStream>, std::shared_ptr<Session>)>;
109
112enum class StreamMode { kStart, kAccept };
113
126class Session : public std::enable_shared_from_this<Session> {
127 public:
143 static absl::StatusOr<std::shared_ptr<Session>> Create(
144 std::string session_id = {},
145 OnSessionStreamMessage on_stream_message = {},
146 OnSessionStreamDone on_stream_done = {}, data::ByteMap headers = {},
147 SessionOptions options = {},
148 std::shared_ptr<nodes::NodeMap> node_map = nullptr,
149 std::shared_ptr<actions::ActionRegistry> action_registry = nullptr);
150
151 virtual ~Session() = default;
152
158 absl::StatusOr<
159 std::vector<std::pair<std::string, std::shared_ptr<net::WireStream>>>>
160 Streams() const;
161
168 absl::StatusOr<std::shared_ptr<net::WireStream>> GetStream(
169 std::string_view stream_id) const;
170
176 [[nodiscard]] std::string GetId() const;
177
182 [[nodiscard]] std::shared_ptr<nodes::NodeMap> GetNodeMap() const;
183
194 absl::Status SetNodeMap(const std::shared_ptr<nodes::NodeMap>& node_map);
195
200 [[nodiscard]] std::shared_ptr<actions::ActionRegistry> GetActionRegistry()
201 const;
202
213 absl::Status SetActionRegistry(
214 const std::shared_ptr<actions::ActionRegistry>& registry);
215
221 [[nodiscard]] std::vector<
222 std::pair<std::string, std::shared_ptr<actions::Action>>>
223 Actions() const;
224
230 absl::StatusOr<std::shared_ptr<actions::Action>> GetAction(
231 std::string_view action_id) const;
232
239 absl::Status CancelAction(std::string_view action_id);
240
246 absl::Status CancelAllActions();
247
254 a11::Task AwaitAllActions(absl::Duration timeout = absl::InfiniteDuration());
255
262
270 data::ActionMessage message,
271 std::shared_ptr<net::WireStream> origin_stream = nullptr);
272
278 a11::Task DispatchAction(const std::shared_ptr<actions::Action>& action);
279
288 data::WireMessage message,
289 std::shared_ptr<net::WireStream> origin_stream = nullptr);
290
295 [[nodiscard]] bool IsClosed() const;
296
302 [[nodiscard]] bool IsDone() const;
303
309 [[nodiscard]] a11::Task Done() const;
310
315 [[nodiscard]] absl::Status GetStatus() const;
316
324 absl::StatusOr<a11::Task> AddStream(std::shared_ptr<net::WireStream> stream,
326
332 absl::Status HalfClose();
333
339 virtual absl::Status Abort(absl::Status status);
340
348 absl::Status Send(data::WireMessage message, std::string_view stream_id = {});
349
354 [[nodiscard]] absl::Time deadline() const;
355
361 absl::Status SetDeadline(absl::Time deadline = absl::InfiniteFuture());
362
363 protected:
364 Session() = default;
365 absl::Status Initialize(
366 const std::shared_ptr<Session>& self, std::string session_id,
367 OnSessionStreamMessage on_stream_message,
368 OnSessionStreamDone on_stream_done, data::ByteMap headers,
369 SessionOptions options, std::shared_ptr<nodes::NodeMap> node_map,
370 std::shared_ptr<actions::ActionRegistry> action_registry);
371
372 private:
373 struct State;
374 struct StreamState;
375 std::shared_ptr<State> state_;
376
377 a11::Task HandleStreamMessage(
378 const std::shared_ptr<StreamState>& stream_state,
379 std::optional<data::WireMessage> message);
380 void ProcessStreamMessages(const std::shared_ptr<StreamState>& stream_state);
381 a11::Task HandleStreamDone(const std::shared_ptr<StreamState>& stream_state);
382 void RemoveStream(const std::shared_ptr<StreamState>& stream_state);
383 void FinishIfPossible();
384 void NotifyStateChanged();
385
386 absl::Status TrackAction(const std::shared_ptr<actions::Action>& action);
387 void UntrackAction(const std::shared_ptr<actions::Action>& action);
388 std::shared_ptr<actions::ActionLimiter> GetActionLimiter(bool nested) const;
389
390 friend class actions::Action;
391};
392
398
408class SessionWithRecv final : public Session {
409 public:
421 static absl::StatusOr<std::shared_ptr<SessionWithRecv>> Create(
422 std::string session_id = {}, data::ByteMap headers = {},
423 SessionOptions options = {},
424 std::shared_ptr<nodes::NodeMap> node_map = nullptr,
425 std::shared_ptr<actions::ActionRegistry> action_registry = nullptr);
426
434 absl::Time deadline = absl::InfiniteFuture());
435
443 absl::Time deadline = absl::InfiniteFuture());
444
450 absl::Status Abort(absl::Status status) override;
451
452 private:
453 struct ReceiveState;
454 std::shared_ptr<ReceiveState> receive_state_;
455
456 a11::Task OnMessage(std::optional<data::WireMessage> message,
457 std::shared_ptr<net::WireStream> stream);
458 a11::Task OnDone(std::shared_ptr<net::WireStream> stream);
459 void SignalReceiveError(absl::Status status);
460};
461
467absl::StatusOr<data::ByteMap> NormalizeSessionHeaders(data::ByteMap headers);
468
469} // namespace a11::service
470
471#endif // A11_SERVICE_SESSION_H_
A11's unit of work: the Action and its supporting types.
std::string stream_id
Definition authorization.cc:328
std::string session_id
Definition authorization.cc:327
A11's unit of work: a schema-described, asynchronously run operation.
Definition action.h:136
A Session variant that buffers inbound messages for pull-style reception.
Definition session.h:408
a11::Future< std::optional< data::WireMessage > > Receive(absl::Time deadline=absl::InfiniteFuture())
Await the next inbound wire message.
Definition session.cc:1862
static absl::StatusOr< std::shared_ptr< SessionWithRecv > > Create(std::string session_id={}, data::ByteMap headers={}, SessionOptions options={}, std::shared_ptr< nodes::NodeMap > node_map=nullptr, std::shared_ptr< actions::ActionRegistry > action_registry=nullptr)
Create a pull-style session.
Definition session.cc:1690
absl::Status Abort(absl::Status status) override
Abort the session, also failing any pending receivers.
Definition session.cc:1906
a11::Future< std::optional< ReceivedSessionMessage > > ReceiveWithStreamId(absl::Time deadline=absl::InfiniteFuture())
Await the next inbound message together with its stream id.
Definition session.cc:1809
A connection-scoped runtime that multiplexes wire streams and runs actions.
Definition session.h:126
static absl::StatusOr< std::shared_ptr< Session > > Create(std::string session_id={}, OnSessionStreamMessage on_stream_message={}, OnSessionStreamDone on_stream_done={}, data::ByteMap headers={}, SessionOptions options={}, std::shared_ptr< nodes::NodeMap > node_map=nullptr, std::shared_ptr< actions::ActionRegistry > action_registry=nullptr)
Create a session.
Definition session.cc:241
virtual absl::Status Abort(absl::Status status)
Abort the session immediately, cancelling streams and actions.
Definition session.cc:1556
absl::Status HalfClose()
Signal that this side will send no more messages.
Definition session.cc:1515
bool IsDone() const
Report whether the session has fully finished.
Definition session.cc:1098
absl::StatusOr< std::vector< std::pair< std::string, std::shared_ptr< net::WireStream > > > > Streams() const
Return the streams currently attached to the session.
Definition session.cc:383
std::shared_ptr< nodes::NodeMap > GetNodeMap() const
Return the NodeMap backing this session's node state.
Definition session.cc:412
a11::Task Done() const
Await the session's full completion.
Definition session.cc:1103
absl::Status Initialize(const std::shared_ptr< Session > &self, std::string session_id, OnSessionStreamMessage on_stream_message, OnSessionStreamDone on_stream_done, data::ByteMap headers, SessionOptions options, std::shared_ptr< nodes::NodeMap > node_map, std::shared_ptr< actions::ActionRegistry > action_registry)
Definition session.cc:256
a11::Task DispatchAction(const std::shared_ptr< actions::Action > &action)
Run an already-constructed action within the session.
Definition session.cc:947
std::vector< std::pair< std::string, std::shared_ptr< actions::Action > > > Actions() const
Return the actions currently running in the session.
Definition session.cc:462
absl::Status GetStatus() const
Return the session's terminal status.
Definition session.cc:1108
absl::Status Send(data::WireMessage message, std::string_view stream_id={})
Enqueue a wire message for delivery.
Definition session.cc:1624
absl::Status CancelAction(std::string_view action_id)
Request cancellation of a running action.
Definition session.cc:478
std::string GetId() const
Return the session's unique identifier.
Definition session.cc:407
a11::Task AwaitAllActions(absl::Duration timeout=absl::InfiniteDuration())
Wait for all in-flight actions to finish.
Definition session.cc:500
absl::StatusOr< a11::Task > AddStream(std::shared_ptr< net::WireStream > stream, StreamMode mode=StreamMode::kStart)
Attach a wire stream and begin pumping its messages.
Definition session.cc:1122
a11::Task DispatchWireMessage(data::WireMessage message, std::shared_ptr< net::WireStream > origin_stream=nullptr)
Route a wire message through the session as if it arrived on a stream.
Definition session.cc:968
absl::Time deadline() const
Return the absolute deadline after which the session is aborted.
Definition session.cc:1660
absl::Status SetNodeMap(const std::shared_ptr< nodes::NodeMap > &node_map)
Replace the NodeMap backing this session's node state.
Definition session.cc:417
bool IsClosed() const
Report whether the session has been closed.
Definition session.cc:1093
std::shared_ptr< actions::ActionRegistry > GetActionRegistry() const
Return the registry used to resolve incoming action messages.
Definition session.cc:438
virtual ~Session()=default
absl::Status SetDeadline(absl::Time deadline=absl::InfiniteFuture())
Set the absolute deadline after which the session is aborted.
Definition session.cc:1665
absl::Status CancelAllActions()
Request cancellation of every running action.
Definition session.cc:484
a11::Future< std::uint32_t > DispatchNodeFragment(data::NodeFragment fragment)
Dispatch a node fragment into the session's NodeMap.
Definition session.cc:599
absl::StatusOr< std::shared_ptr< actions::Action > > GetAction(std::string_view action_id) const
Look up a running action by id.
Definition session.cc:467
absl::StatusOr< std::shared_ptr< net::WireStream > > GetStream(std::string_view stream_id) const
Look up an attached stream by id.
Definition session.cc:396
a11::Task DispatchActionMessage(data::ActionMessage message, std::shared_ptr< net::WireStream > origin_stream=nullptr)
Resolve an action message against the registry and run it.
Definition session.cc:824
absl::Status SetActionRegistry(const std::shared_ptr< actions::ActionRegistry > &registry)
Replace the registry used to resolve incoming action messages.
Definition session.cc:443
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:74
absl::StatusOr< data::ByteMap > NormalizeSessionHeaders(data::ByteMap headers)
Validate and case-normalize session headers.
Definition session.cc:232
StreamMode
Whether this side starts (kStart) or accepts (kAccept) a stream during its startup handshake.
Definition session.h:112
constexpr size_t kMaxSingleMessageSize
Hard upper bound for any one WireMessage admitted by a Session.
Definition session.h:61
constexpr std::string_view kSessionStatusHeader
Trailer/header used to communicate the session's structured terminal status.
Definition session.h:59
std::function< a11::Task(std::shared_ptr< net::WireStream >, std::shared_ptr< Session >)> OnSessionStreamDone
Callback invoked once a session stream has finished; may be a coroutine.
Definition session.h:108
std::function< a11::Task(std::optional< data::WireMessage >, std::shared_ptr< net::WireStream >, std::shared_ptr< Session >)> OnSessionStreamMessage
Callback invoked for each message received on a session stream (nullopt signals end-of-stream); may b...
Definition session.h:105
Future< Unit > Task
Asynchronous operation whose only successful result is completion itself.
Definition future.h:411
A registry of AsyncNodes keyed by node id.
The action registry: a catalogue of named schema/handler pairs.
The wire description of an action invocation.
Definition types.h:435
One piece of a node's stream: an inline chunk or a node reference.
Definition types.h:335
The top-level frame exchanged between two A11 endpoints.
Definition types.h:485
An inbound wire message paired with the id of the stream it arrived on.
Definition session.h:394
std::string stream_id
Id of the transport that delivered it.
Definition session.h:396
data::WireMessage message
Received application message.
Definition session.h:395
Limits and timeouts governing a Session's buffering, concurrency, and lifetime.
Definition session.h:72
size_t max_buffered_messages_per_stream
Maximum number of messages buffered per stream.
Definition session.h:76
size_t max_concurrent_nested_actions
Maximum number of concurrently running nested (child) actions.
Definition session.h:80
size_t max_concurrent_root_actions
Maximum number of concurrently running root (top-level) actions.
Definition session.h:78
absl::Time deadline
Absolute time after which the session is aborted.
Definition session.h:90
absl::Duration no_stream_timeout
How long the session waits with no active stream before finishing.
Definition session.h:88
size_t max_buffered_messages_total
Maximum number of messages buffered across all streams.
Definition session.h:74
size_t max_buffered_bytes_total
Maximum total bytes buffered across all streams.
Definition session.h:84
size_t max_buffered_bytes_per_stream
Maximum bytes buffered per stream.
Definition session.h:86
size_t max_single_message_size
Maximum size in bytes of a single wire message.
Definition session.h:82
absl::Status Validate() const
Validate the option values.
Definition session.cc:204
Definition session.cc:138
Definition session.cc:113
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...