A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
executor.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_CONCURRENCY_EXECUTOR_H_
18#define A11_CONCURRENCY_EXECUTOR_H_
19
20#include <exception>
21#include <functional>
22#include <type_traits>
23#include <utility>
24
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>
30
32
33namespace a11 {
34
37void Schedule(absl::AnyInvocable<void() &&> work,
38 thread::TreeOptions tree_options = {});
39
42std::function<void()> ScheduleCancelable(absl::AnyInvocable<void() &&> work,
43 thread::TreeOptions tree_options = {});
44
60template <typename T, typename Fn>
61auto ThenAfterWaiting(Future<T> future, absl::Time deadline, Fn transform)
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;
66
67 if (future.IsReady()) {
68 return CompletedFuture<U>(transform(future.Await()));
69 }
70 return Submit<U>(
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);
75 });
76}
77
78template <typename T>
80 absl::AnyInvocable<absl::StatusOr<T>() &&> work,
81 std::function<void()> cancellation_hook, thread::TreeOptions tree_options) {
82 Promise<T> promise;
83 Future<T> future = promise.future();
84 std::function<void()> cancel = ScheduleCancelable(
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");
89 } else {
90 // The work is the caller's, and so is this instantiation; see
91 // a11/exception_guard.h for why the guard belongs here rather than at
92 // the call.
93 const absl::Status raised = exception_guard::Attempt(
94 [&] { result = std::move(work)(); }, "task");
95 if (!raised.ok()) {
96 result = raised;
97 }
98 }
99 const absl::Status completion = promise.SetResult(std::move(result));
100 (void)completion;
101 },
102 std::move(tree_options));
103 // The promise has moved into the task, but both handles share its state.
104 // Install cancellation through a temporary handle recovered from Future.
105 future.SetCancellationCallbackForExecutor(
106 [cancel = std::move(cancel),
107 cancellation_hook = std::move(cancellation_hook)]() {
108 if (cancellation_hook != nullptr) {
109 cancellation_hook();
110 }
111 cancel();
112 });
113 return future;
114}
115
117template <typename T>
118Future<T> Submit(absl::AnyInvocable<absl::StatusOr<T>() &&> work,
119 thread::TreeOptions tree_options) {
120 return SubmitWithCancellationHook<T>(std::move(work), {},
121 std::move(tree_options));
122}
123
125inline Task SubmitTask(absl::AnyInvocable<absl::Status() &&> work,
126 thread::TreeOptions tree_options = {}) {
127 return Submit<Unit>(
128 [work = std::move(work)]() mutable -> absl::StatusOr<Unit> {
129 ABSL_RETURN_IF_ERROR(std::move(work)());
130 return Unit{};
131 },
132 tree_options);
133}
134
135} // namespace a11
136
137#endif // A11_CONCURRENCY_EXECUTOR_H_
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
Definition action.cc:63
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