33#ifndef A11_SERVICE_SESSION_H_
34#define A11_SERVICE_SESSION_H_
45#include <absl/status/status.h>
46#include <absl/status/statusor.h>
47#include <absl/time/time.h>
104 std::optional<data::WireMessage>, std::shared_ptr<net::WireStream>,
105 std::shared_ptr<Session>)>;
108 std::shared_ptr<net::WireStream>, std::shared_ptr<Session>)>;
126class Session :
public std::enable_shared_from_this<Session> {
143 static absl::StatusOr<std::shared_ptr<Session>>
Create(
148 std::shared_ptr<nodes::NodeMap> node_map =
nullptr,
149 std::shared_ptr<actions::ActionRegistry> action_registry =
nullptr);
159 std::vector<std::pair<std::string, std::shared_ptr<net::WireStream>>>>
168 absl::StatusOr<std::shared_ptr<net::WireStream>>
GetStream(
176 [[nodiscard]] std::string
GetId()
const;
182 [[nodiscard]] std::shared_ptr<nodes::NodeMap>
GetNodeMap()
const;
194 absl::Status
SetNodeMap(
const std::shared_ptr<nodes::NodeMap>& node_map);
214 const std::shared_ptr<actions::ActionRegistry>& registry);
221 [[nodiscard]] std::vector<
222 std::pair<std::string, std::shared_ptr<actions::Action>>>
230 absl::StatusOr<std::shared_ptr<actions::Action>>
GetAction(
231 std::string_view action_id)
const;
271 std::shared_ptr<net::WireStream> origin_stream =
nullptr);
289 std::shared_ptr<net::WireStream> origin_stream =
nullptr);
295 [[nodiscard]]
bool IsClosed()
const;
302 [[nodiscard]]
bool IsDone()
const;
315 [[nodiscard]] absl::Status
GetStatus()
const;
324 absl::StatusOr<a11::Task>
AddStream(std::shared_ptr<net::WireStream> stream,
339 virtual absl::Status
Abort(absl::Status status);
354 [[nodiscard]] absl::Time
deadline()
const;
366 const std::shared_ptr<Session>& self, std::string
session_id,
370 std::shared_ptr<actions::ActionRegistry> action_registry);
375 std::shared_ptr<State> state_;
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();
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;
421 static absl::StatusOr<std::shared_ptr<SessionWithRecv>>
Create(
424 std::shared_ptr<nodes::NodeMap> node_map =
nullptr,
425 std::shared_ptr<actions::ActionRegistry> action_registry =
nullptr);
434 absl::Time
deadline = absl::InfiniteFuture());
443 absl::Time
deadline = absl::InfiniteFuture());
450 absl::Status
Abort(absl::Status status)
override;
454 std::shared_ptr<ReceiveState> receive_state_;
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);
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 > ®istry)
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
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:1678
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...