17#ifndef A11_CONCURRENCY_EXECUTOR_H_
18#define A11_CONCURRENCY_EXECUTOR_H_
25#include <absl/functional/any_invocable.h>
26#include <absl/status/status.h>
27#include <absl/status/status_macros.h>
28#include <absl/status/statusor.h>
29#include <absl/time/time.h>
37void Schedule(absl::AnyInvocable<
void() &&> work,
38 thread::TreeOptions tree_options = {});
43 thread::TreeOptions tree_options = {});
60template <
typename T,
typename Fn>
62 ->
Future<
typename std::invoke_result_t<
63 Fn,
const absl::StatusOr<T>&>::value_type> {
64 using Result = std::invoke_result_t<Fn, const absl::StatusOr<T>&>;
65 using U =
typename Result::value_type;
67 if (future.IsReady()) {
68 return CompletedFuture<U>(transform(future.Await()));
71 [future = std::move(future), deadline,
72 transform = std::move(transform)]()
mutable -> absl::StatusOr<U> {
73 const absl::StatusOr<T> result = future.Await(deadline);
74 return transform(result);
80 absl::AnyInvocable<absl::StatusOr<T>() &&> work,
81 std::function<
void()> cancellation_hook, thread::TreeOptions tree_options) {
85 [promise = std::move(promise), work = std::move(work)]()
mutable {
86 absl::StatusOr<T> result;
87 if (thread::Cancelled()) {
88 result = absl::CancelledError(
"Task cancelled before it started");
94 [&] { result = std::move(work)(); },
"task");
99 const absl::Status completion = promise.
SetResult(std::move(result));
102 std::move(tree_options));
105 future.SetCancellationCallbackForExecutor(
106 [cancel = std::move(cancel),
107 cancellation_hook = std::move(cancellation_hook)]() {
108 if (cancellation_hook !=
nullptr) {
119 thread::TreeOptions tree_options) {
120 return SubmitWithCancellationHook<T>(std::move(work), {},
121 std::move(tree_options));
126 thread::TreeOptions tree_options = {}) {
128 [work = std::move(work)]()
mutable -> absl::StatusOr<Unit> {
129 ABSL_RETURN_IF_ERROR(std::move(work)());
Shared handle to one asynchronous result.
Definition future.h:126
Move-only producer for a Future result.
Definition future.h:282
absl::Status SetResult(absl::StatusOr< T > result)
Complete with either a value or an error, waking every observer.
Definition future.h:336
Future< T > future() const
Return a consumer handle sharing this promise's completion state.
Definition future.h:304
Completion values used by every asynchronous A11 operation.
absl::Status Attempt(Callable &&callable, std::string_view what)
Runs callable, reporting anything it throws as a Status.
Definition exception_guard.h:114
auto ThenAfterWaiting(Future< T > future, absl::Time deadline, Fn transform) -> Future< typename std::invoke_result_t< Fn, const absl::StatusOr< T > & >::value_type >
Continue with transform when future completes or the deadline expires.
Definition executor.h:61
std::function< void()> ScheduleCancelable(absl::AnyInvocable< void() && > work, thread::TreeOptions tree_options)
Schedule work and return an idempotent cooperative-cancellation function.
Definition executor.cc:59
Task SubmitTask(absl::AnyInvocable< absl::Status() && > work, thread::TreeOptions tree_options={})
Run an operation that returns only a completion status.
Definition executor.h:125
void Schedule(absl::AnyInvocable< void() && > work, thread::TreeOptions tree_options)
Schedule work on A11's fiber pool without returning a completion handle.
Definition executor.cc:47
Future< T > Submit(absl::AnyInvocable< absl::StatusOr< T >() && > work, thread::TreeOptions tree_options)
Run status-returning work on the fiber pool and expose its Future.
Definition executor.h:118
Future< T > SubmitWithCancellationHook(absl::AnyInvocable< absl::StatusOr< T >() && > work, std::function< void()> cancellation_hook, thread::TreeOptions tree_options)
Run work on A11's fiber pool with application-specific cancellation.
Definition executor.h:79