32#ifndef A11_ACTIONS_ACTION_H_
33#define A11_ACTIONS_ACTION_H_
45#include <absl/base/thread_annotations.h>
46#include <absl/container/flat_hash_map.h>
47#include <absl/container/flat_hash_set.h>
48#include <absl/status/status.h>
49#include <absl/status/statusor.h>
50#include <absl/strings/str_format.h>
51#include <absl/time/time.h>
59#include "thread/boost_primitives.h"
82struct VerifiedAuthorization;
109 static absl::StatusOr<std::shared_ptr<ActionLimiter>>
Create(
size_t maximum);
121 const size_t maximum_;
122 size_t active_ ABSL_GUARDED_BY(mu_) = 0;
123 std::shared_ptr<thread::PermanentEvent> changed_ ABSL_GUARDED_BY(mu_);
136class Action :
public std::enable_shared_from_this<Action> {
150 static absl::StatusOr<std::shared_ptr<Action>>
Create(
153 std::shared_ptr<nodes::NodeMap> node_map =
nullptr,
154 std::shared_ptr<net::WireStream> stream =
nullptr,
155 std::shared_ptr<service::Session>
session =
nullptr,
156 std::shared_ptr<ActionRegistry> registry =
nullptr,
157 size_t max_concurrent_nested_actions =
163 static absl::StatusOr<std::string>
MakeNodeId(std::string_view action_id,
164 std::string_view node_name);
167 [[nodiscard]] std::string
GetId()
const;
169 absl::Status
SetId(std::string action_id);
195 absl::Status
BindNodeMap(std::shared_ptr<nodes::NodeMap> node_map);
197 [[nodiscard]] std::shared_ptr<nodes::NodeMap>
GetNodeMap()
const;
199 absl::Status
BindStream(std::shared_ptr<net::WireStream> stream);
201 [[nodiscard]] std::shared_ptr<net::WireStream>
GetStream()
const;
203 absl::Status
BindRegistry(std::shared_ptr<ActionRegistry> registry);
205 [[nodiscard]] std::shared_ptr<ActionRegistry>
GetRegistry()
const;
209 [[nodiscard]] std::shared_ptr<service::Session>
GetSession()
const;
212 std::shared_ptr<const VerifiedAuthorization> authorization);
214 [[nodiscard]] std::shared_ptr<const VerifiedAuthorization>
218 absl::StatusOr<std::shared_ptr<nodes::AsyncNode>>
GetNode(
219 std::string node_id);
226 absl::StatusOr<std::shared_ptr<nodes::AsyncNode>>
GetInput(
227 const std::string&
name, std::optional<bool> bind_stream = std::nullopt);
234 absl::StatusOr<std::shared_ptr<nodes::AsyncNode>>
GetOutput(
235 const std::string&
name, std::optional<bool> bind_stream = std::nullopt);
237 absl::StatusOr<std::shared_ptr<nodes::AsyncNode>>
GetPort(
238 const std::string&
name);
271 template <
typename T>
292 template <
typename... Args>
293 absl::Status
Logf(
const absl::FormatSpec<Args...>&
format,
294 const Args&... args);
305 template <
typename... Args>
307 const absl::FormatSpec<Args...>&
format,
308 const Args&... args);
320 absl::StatusOr<std::shared_ptr<nodes::AsyncNode>>
GetLogNode();
332 absl::StatusOr<std::optional<data::Bytes>>
GetHeader(
333 std::string_view
name)
const;
342 std::string_view
name)
const;
345 const std::shared_ptr<Action>&
target,
362 absl::StatusOr<std::shared_ptr<Action>>
MakeNested(
364 bool forward_headers =
true);
376 absl::StatusOr<std::shared_ptr<Action>>
MakeNested(
377 std::string_view action_name,
bool propagate_io =
true,
378 bool forward_headers =
true);
389 absl::StatusOr<std::shared_ptr<Action>>
Run();
402 absl::Duration
timeout = absl::InfiniteDuration());
410 absl::Duration
timeout = absl::InfiniteDuration());
422 [[nodiscard]] std::string
TraceId()
const;
424 [[nodiscard]] std::string
SpanId()
const;
451 [[nodiscard]]
bool IsDone()
const;
459 [[nodiscard]] absl::Status
GetStatus()
const;
466 enum class Mode { kNone, kRun, kCall, kCancelled };
469 std::shared_ptr<nodes::NodeMap> node_map,
470 std::shared_ptr<net::WireStream> stream,
471 const std::shared_ptr<service::Session>&
session,
472 std::shared_ptr<ActionRegistry> registry,
473 std::shared_ptr<ActionLimiter> nested_limiter);
475 absl::Status Begin(Mode mode);
476 absl::Status RemapDefaultPorts() ABSL_EXCLUSIVE_LOCKS_REQUIRED(mu_);
477 absl::Status AttachStreamIfRequested(
478 const std::shared_ptr<nodes::AsyncNode>&
node,
bool bind);
479 absl::Status ValidateMessagePorts(
480 const std::vector<data::Port>& ports,
481 const absl::flat_hash_map<std::string, ActionPortSchema>& schema_ports,
482 std::string_view
kind)
const;
483 void RunHandler(
const std::shared_ptr<ActionLimiter>& limiter);
484 absl::Status ApplyInputAutofills();
485 [[nodiscard]] std::vector<data::NodeFragment> CollectAutofillFragments()
487 void StartFinish(absl::Status status);
488 absl::Status FinishRun(absl::Status status);
491 [[nodiscard]]
bool HasInputAutofills()
const;
494 absl::StatusOr<std::shared_ptr<Action>> RunHandlerWithoutFiber(
495 const std::shared_ptr<Action>& self,
496 std::shared_ptr<ActionLimiter> acquired_limiter =
nullptr);
497 void ReleaseAcquiredLimiter();
498 absl::Status FinishOutputNodes(
const absl::Status& status);
502 static absl::Status CloseUnwrittenOutput(
503 const std::shared_ptr<nodes::AsyncNode>&
node,
504 const absl::Status& status);
509 absl::Status CommunicateStatus(
const absl::Status& status);
510 absl::Status AbortInputs(
const absl::Status& status);
511 absl::Status SendNodeAbortStatuses(
512 const absl::flat_hash_set<std::string>& node_ids,
513 const absl::Status& status);
514 absl::Status ReleaseNodesAfterRun();
515 absl::Status DetachBoundStreamNodes();
516 absl::Status SendRemoteCancel();
517 void CompleteCall(
const absl::Status& status,
bool remove_from_session);
518 void AbortLocalCallOutputs(
const absl::Status& status);
519 absl::Status TrackInSession(
const std::shared_ptr<service::Session>&
session);
520 void UntrackFromSession();
521 void SetDispatchStatus(absl::Status status);
522 void SetCompletionStatus(
const absl::Status& status);
525 absl::Status StartActionSpan(Mode mode);
526 void EndActionSpan(
const absl::Status& status);
528 void RecordActionCallEvent(std::string_view
name, std::string_view
id);
530 mutable thread::Mutex mu_;
533 std::string id_ ABSL_GUARDED_BY(mu_);
536 std::shared_ptr<nodes::NodeMap> node_map_ ABSL_GUARDED_BY(mu_);
537 std::shared_ptr<net::WireStream> stream_ ABSL_GUARDED_BY(mu_);
538 std::weak_ptr<service::Session> session_ ABSL_GUARDED_BY(mu_);
539 std::weak_ptr<service::Session> tracked_session_ ABSL_GUARDED_BY(mu_);
540 std::shared_ptr<ActionRegistry> registry_ ABSL_GUARDED_BY(mu_);
541 std::shared_ptr<const VerifiedAuthorization> verified_authorization_
542 ABSL_GUARDED_BY(mu_);
543 absl::flat_hash_map<std::string, std::string> input_ids_ ABSL_GUARDED_BY(mu_);
544 absl::flat_hash_map<std::string, std::string> output_ids_
545 ABSL_GUARDED_BY(mu_);
546 absl::flat_hash_set<std::shared_ptr<nodes::AsyncNode>> input_nodes_
547 ABSL_GUARDED_BY(mu_);
548 absl::flat_hash_set<std::shared_ptr<nodes::AsyncNode>> output_nodes_
549 ABSL_GUARDED_BY(mu_);
550 absl::flat_hash_set<std::shared_ptr<nodes::AsyncNode>> stream_bound_nodes_
551 ABSL_GUARDED_BY(mu_);
552 Mode mode_ ABSL_GUARDED_BY(mu_) = Mode::kNone;
553 bool input_autofills_applied_ ABSL_GUARDED_BY(mu_) =
false;
555 std::shared_ptr<ActionLimiter> acquired_limiter_ ABSL_GUARDED_BY(mu_);
557 bool span_status_set_by_user_ ABSL_GUARDED_BY(mu_) =
false;
558 bool cancel_requested_ ABSL_GUARDED_BY(mu_) =
false;
559 bool finishing_ ABSL_GUARDED_BY(mu_) =
false;
560 std::optional<absl::Status> completion_status_ ABSL_GUARDED_BY(mu_);
564 bool outputs_finished_ ABSL_GUARDED_BY(mu_) =
false;
565 absl::Status outputs_final_status_ ABSL_GUARDED_BY(mu_);
568 bool log_claimed_ ABSL_GUARDED_BY(mu_) =
false;
569 std::optional<absl::Status> dispatch_status_ ABSL_GUARDED_BY(mu_);
570 std::shared_ptr<a11::Promise<a11::Unit>> done_promise_ ABSL_GUARDED_BY(mu_);
571 a11::Task done_future_ ABSL_GUARDED_BY(mu_);
572 std::shared_ptr<a11::Promise<a11::Unit>> dispatch_promise_
573 ABSL_GUARDED_BY(mu_);
574 a11::Task dispatch_future_ ABSL_GUARDED_BY(mu_);
575 std::vector<OnActionCancelled> cancel_callbacks_ ABSL_GUARDED_BY(mu_);
576 std::weak_ptr<Action> parent_ ABSL_GUARDED_BY(mu_);
577 absl::flat_hash_set<std::shared_ptr<Action>> children_ ABSL_GUARDED_BY(mu_);
578 std::shared_ptr<ActionLimiter> nested_limiter_ ABSL_GUARDED_BY(mu_);
597inline constexpr bool kLogsAsText =
598 std::is_convertible_v<const T&, std::string_view>;
606 absl::StatusOr<data::Chunk> chunk;
607 if constexpr (internal::kLogsAsText<T>) {
610 const std::string
text{std::string_view(
value)};
618 return chunk.status();
620 return WriteLog(*std::move(chunk), options);
623template <
typename... Args>
625 const Args&... args) {
629template <
typename... Args>
631 const absl::FormatSpec<Args...>&
format,
632 const Args&... args) {
633 return Log(absl::StrFormat(
format, args...), options);
Schemas describing an action's typed interface and its settings.
std::shared_ptr< service::Session > session
Definition authorization.cc:326
Shared handle to one asynchronous result.
Definition future.h:126
Cancellation-aware counting semaphore for nested-action concurrency.
Definition action.h:106
absl::Status Acquire()
Acquires a slot, blocking until one is free or cancelled.
Definition action.cc:161
bool TryAcquire()
Acquires a slot if one is immediately available.
Definition action.cc:152
void Release()
Releases a previously acquired slot.
Definition action.cc:181
static absl::StatusOr< std::shared_ptr< ActionLimiter > > Create(size_t maximum)
Creates a limiter admitting at most maximum holders.
Definition action.cc:138
A thread-safe catalogue mapping action names to schema and handler.
Definition registry.h:59
A11's unit of work: a schema-described, asynchronously run operation.
Definition action.h:136
std::shared_ptr< const VerifiedAuthorization > GetVerifiedAuthorization() const
Returns the verified authorization bound to this action.
Definition action.cc:480
absl::Status ClearInputsAfterRun(bool clear=true)
Sets whether inputs are released after each run.
Definition action.cc:360
void SetSpanStatus(obs::SpanStatus status, std::string_view description={})
Sets the span status explicitly.
Definition action.cc:1360
absl::Status BindStream(std::shared_ptr< net::WireStream > stream)
Binds the wire stream used to dispatch the action remotely.
Definition action.cc:383
absl::Status BindSession(const std::shared_ptr< service::Session > &session)
Binds the owning session.
Definition action.cc:440
std::optional< absl::Status > GetDispatchStatus() const
The remote dispatch status, or nullopt when not (yet) dispatched.
Definition action.cc:1241
absl::Status SetHeader(const std::string &name, data::Bytes value)
Sets header name to value.
Definition action.cc:803
bool HasBeenRun() const
Whether the action has been started with Run.
Definition action.cc:1219
absl::Status SetId(std::string action_id)
Sets this action's instance id.
Definition action.cc:279
absl::StatusOr< std::optional< data::Bytes > > GetHeader(std::string_view name) const
Returns header name, or nullopt when absent.
Definition action.cc:783
data::ByteMap Headers() const
Returns a copy of all headers.
Definition action.cc:778
absl::Status SetOnCancelled(OnActionCancelled callback)
Registers a callback invoked when the action is cancelled.
Definition action.cc:1205
std::string GetId() const
Returns this action's instance id.
Definition action.cc:274
absl::Status RemoveHeader(std::string_view name)
Removes header name.
Definition action.cc:810
std::string TraceId() const
This action's trace id as lowercase hex.
Definition action.cc:1325
bool HasBeenCalled() const
Whether the action has been dispatched with Call.
Definition action.cc:1224
ActionSchema GetSchema() const
Returns this action's schema.
Definition action.cc:298
std::shared_ptr< nodes::NodeMap > GetNodeMap() const
Returns the bound node map.
Definition action.cc:378
absl::StatusOr< std::shared_ptr< nodes::AsyncNode > > GetInput(const std::string &name, std::optional< bool > bind_stream=std::nullopt)
Returns the input port node named name.
Definition action.cc:498
std::shared_ptr< service::Session > GetSession() const
Returns the owning session, if any.
Definition action.cc:464
ActionSettings GetSettings() const
Returns the action's current settings.
Definition action.cc:337
absl::Status BindVerifiedAuthorization(std::shared_ptr< const VerifiedAuthorization > authorization)
Binds identity and authority already verified for this action.
Definition action.cc:469
absl::Status LogfWith(const LogOptions &options, const absl::FormatSpec< Args... > &format, const Args &... args)
Logs a formatted string with explicit options.
Definition action.h:630
void SetSpanName(std::string_view name)
Overrides the display name of this action's span.
Definition action.cc:1355
absl::Status SetSchema(ActionSchema schema)
Replaces this action's schema (validated).
Definition action.cc:303
bool ContainsPort(std::string_view name) const
Whether the schema declares a port named name.
Definition action.cc:598
absl::Status ForwardHeadersWithPrefix(const std::shared_ptr< Action > &target, std::string_view prefix=kActionHeaderPrefix) const
Copies all headers starting with prefix onto target.
Definition action.cc:828
absl::Status Log(const T &value, const LogOptions &options={})
Logs value on the reserved kActionLogOutput port.
Definition action.h:603
absl::Status ClearOutputsAfterRun(bool clear=true)
Sets whether outputs are released after each run.
Definition action.cc:366
absl::Status BindRegistry(std::shared_ptr< ActionRegistry > registry)
Binds the registry used to resolve nested actions by name.
Definition action.cc:429
static absl::StatusOr< std::shared_ptr< Action > > Create(ActionSchema schema, std::string action_id={}, ActionHandler handler={}, std::shared_ptr< nodes::NodeMap > node_map=nullptr, std::shared_ptr< net::WireStream > stream=nullptr, std::shared_ptr< service::Session > session=nullptr, std::shared_ptr< ActionRegistry > registry=nullptr, size_t max_concurrent_nested_actions=kDefaultMaxConcurrentNestedActions)
Creates an action.
Definition action.cc:195
absl::Status MapPortsFromMessage(const data::ActionMessage &message)
Binds this action's ports to the nodes named in message.
Definition action.cc:758
absl::Status BindStreamsOnInputsByDefault(bool bind)
Sets whether input port streams are bound by default.
Definition action.cc:348
bool Cancelled() const
Whether cancellation has been requested/applied.
Definition action.cc:1229
a11::Future< std::shared_ptr< Action > > Wait(absl::Duration timeout=absl::InfiniteDuration())
Awaits completion of the action.
Definition action.cc:1124
a11::Future< std::shared_ptr< Action > > Call(data::ByteMap wire_headers={})
Dispatches the action to a peer over the bound wire stream.
Definition action.cc:1027
absl::Status GetStatus() const
The action's completion status (OK while still running).
Definition action.cc:1236
absl::Status Logf(const absl::FormatSpec< Args... > &format, const Args &... args)
Logs a formatted string, as ::absl::StrFormat would format it.
Definition action.h:624
absl::Status BindNodeMap(std::shared_ptr< nodes::NodeMap > node_map)
Binds the node map that backs this action's ports.
Definition action.cc:372
absl::Status BindHandler(ActionHandler handler)
Binds the handler invoked when the action runs.
Definition action.cc:314
absl::StatusOr< std::shared_ptr< nodes::AsyncNode > > GetPort(const std::string &name)
Returns the input or output port node named name.
Definition action.cc:576
data::ActionMessage GetActionMessage() const
Returns the wire a11::data::ActionMessage describing this action.
Definition action.cc:741
std::string SpanId() const
This action's span id as lowercase hex, or empty when untraced.
Definition action.cc:1330
absl::Status Cancel()
Requests cancellation of the action (local or remote).
Definition action.cc:1150
bool HasHandler() const
Whether a handler is bound.
Definition action.cc:332
bool IsDone() const
Whether the action has finished (successfully or not).
Definition action.cc:1214
std::shared_ptr< ActionRegistry > GetRegistry() const
Returns the bound registry.
Definition action.cc:435
absl::StatusOr< std::shared_ptr< nodes::AsyncNode > > GetLogNode()
Returns the log port's node, claiming it for this consumer.
Definition action.cc:612
bool HasHeader(std::string_view name) const
Whether header name is set.
Definition action.cc:795
void SetSpanAttribute(std::string_view key, std::string_view value)
Sets a string attribute on this action's span.
Definition action.cc:1335
a11::Future< absl::Status > WaitForDispatch(absl::Duration timeout=absl::InfiniteDuration())
Awaits acceptance of a remote dispatch.
Definition action.cc:1081
static absl::StatusOr< std::string > MakeNodeId(std::string_view action_id, std::string_view node_name)
Derives the node id for port node_name of action action_id.
Definition action.cc:265
absl::Status BindStreamsOnOutputsByDefault(bool bind)
Sets whether output port streams are bound by default.
Definition action.cc:354
absl::StatusOr< std::shared_ptr< Action > > Run()
Runs the action's handler locally.
Definition action.cc:902
std::shared_ptr< net::WireStream > GetStream() const
Returns the bound wire stream.
Definition action.cc:424
absl::StatusOr< std::shared_ptr< nodes::AsyncNode > > GetNode(std::string node_id)
Returns the port node with raw id node_id.
Definition action.cc:485
absl::Status ForwardHeader(const std::shared_ptr< Action > &target, std::string_view name) const
Copies header name from this action onto target.
Definition action.cc:817
absl::Status SetSettings(ActionSettings settings)
Replaces the action's settings.
Definition action.cc:342
ActionHandler GetHandler() const
Returns the currently bound handler.
Definition action.cc:327
absl::StatusOr< std::shared_ptr< nodes::AsyncNode > > GetOutput(const std::string &name, std::optional< bool > bind_stream=std::nullopt)
Returns the output port node named name.
Definition action.cc:525
absl::StatusOr< std::shared_ptr< Action > > MakeNested(const ActionSchema &schema, bool propagate_io=true, bool forward_headers=true)
Creates a nested action from a schema, parented to this action.
Definition action.cc:843
Definition serialization.h:161
absl::StatusOr< Chunk > ToChunk(const T &value, std::string_view mimetype={}) const
Serializes value into a Chunk.
Definition serialization.h:291
A bidirectional, message-oriented channel between two A11 endpoints.
Definition wire_stream.h:103
An asynchronous, ordered stream of chunks read from and written to A11.
Definition async_node.h:111
A thread-safe registry of AsyncNodes keyed by id.
Definition node_map.h:66
Move-only RAII handle to one action, session, or transport span.
Definition span.h:51
A connection-scoped runtime that multiplexes wire streams and runs actions.
Definition session.h:126
std::string value
Definition discover.cc:114
std::string key
Definition discover.cc:782
std::optional< std::string > text
The text, where the value is a string literal or a name bound to one.
Definition discover.cc:612
Completion values used by every asynchronous A11 operation.
std::optional< absl::Duration > timeout
Definition main.cc:144
Log levels, log-chunk metadata and the process-wide action log sink.
ActionHandler MakeAsyncActionHandler(SyncActionHandler handler)
Adapts a synchronous handler into an asynchronous ActionHandler.
Definition action.cc:128
std::function< absl::Status(std::shared_ptr< Action >)> SyncActionHandler
Synchronous action handler returning a completion status.
Definition action.h:92
constexpr std::string_view kActionHeaderPrefix
Prefix reserved for A11's framework headers.
Definition schema.h:70
constexpr size_t kDefaultMaxConcurrentNestedActions
Default cap on concurrently running nested actions.
Definition action.h:85
std::function< absl::Status(std::shared_ptr< Action >)> OnActionCancelled
Callback invoked when an action is cancelled.
Definition action.h:94
std::function< a11::Task(std::shared_ptr< Action >)> ActionHandler
Asynchronous action handler: runs the action, returns an awaitable.
Definition action.h:90
std::string Bytes
Raw byte payload; an alias of std::string.
Definition types.h:56
constexpr std::string_view kTextMimetype
Media type of UTF-8 text carried as itself.
Definition serialization.h:73
SerializationRegistry & GlobalSerializationRegistry()
Returns the process-wide registry (defaults pre-installed).
Definition serialization.cc:524
absl::flat_hash_map< std::string, Bytes > ByteMap
String-keyed map of byte values (headers, attributes, etc.).
Definition types.h:58
SpanStatus
OpenTelemetry outcome states, mirrored to keep OTel out of this API.
Definition span.h:38
SpanKind
OpenTelemetry span relationships, mirrored to keep OTel out of this API.
Definition span.h:35
Future< Unit > Task
Asynchronous operation whose only successful result is completion itself.
Definition future.h:411
graph::RefId node
The graph ref this answer corresponds to, when a graph is being built.
Definition resolve.cc:698
Type-and-mimetype indexed serialization of values to/from Chunks.
The full typed interface of an action.
Definition schema.h:138
Per-action runtime settings for stream binding and cleanup.
Definition schema.h:172
Everything about a log other than the object being logged.
Definition log.h:115
std::string_view mimetype
Media type hint for the serializer.
Definition log.h:124
The wire description of an action invocation.
Definition types.h:435
A unit of data: bytes plus optional descriptive metadata.
Definition types.h:185
Output format
Definition main.cc:64
std::string target
Which editor definition syntax is about; every one when empty, which is the default.
Definition main.cc:74
A11's core wire value types: chunks, node fragments and messages.