A11 (C++ runtime)
Native C++ implementation of the A11 streaming action runtime
Loading...
Searching...
No Matches
future.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
27#ifndef A11_CONCURRENCY_FUTURE_H_
28#define A11_CONCURRENCY_FUTURE_H_
29
30#include <exception>
31#include <functional>
32#include <memory>
33#include <optional>
34#include <type_traits>
35#include <utility>
36#include <vector>
37
38#include <absl/base/nullability.h>
39#include <absl/functional/any_invocable.h>
40#include <absl/log/log.h>
41#include <absl/status/status.h>
42#include <absl/status/statusor.h>
43#include <absl/time/clock.h>
44#include <absl/time/time.h>
45
46#include "a11/exception_guard.h"
47#include "thread/boost_primitives.h"
48#include "thread/fiber.h"
49#include "thread/select.h"
50#include "thread/selectables.h"
51
52namespace a11 {
53
55struct Unit {
56 friend bool operator==(Unit, Unit) = default;
57};
58
59template <typename T>
60class Future;
61
62template <typename T>
63class Promise;
64
65namespace internal {
66
67template <typename T>
68struct FutureState {
69 mutable thread::Mutex mu;
70 thread::CondVar cv;
71 bool ready ABSL_GUARDED_BY(mu) = false;
72 std::optional<absl::StatusOr<T>> result ABSL_GUARDED_BY(mu);
73 std::function<void()> cancel ABSL_GUARDED_BY(mu);
74 thread::PermanentEvent event;
75 std::vector<absl::AnyInvocable<void(const absl::StatusOr<T>&)>> callbacks
76 ABSL_GUARDED_BY(mu);
77};
78
79template <typename T>
81 absl::AnyInvocable<void(const absl::StatusOr<T>&)> callback,
82 const absl::StatusOr<T>& result) {
83 // A completion callback belongs to whoever called OnReady, and so does this
84 // instantiation: if their callback can throw, their translation unit has
85 // exceptions and Attempt catches there.
86 const absl::Status raised = exception_guard::Attempt(
87 [&] { callback(result); }, "Future completion callback");
88 if (!raised.ok()) {
89 LOG(ERROR) << raised.message();
90 }
91}
92
93} // namespace internal
94
104template <typename T>
106 absl::AnyInvocable<absl::StatusOr<T>() &&> work,
107 std::function<void()> cancellation_hook,
108 thread::TreeOptions tree_options = {});
109
110template <typename T>
111Future<T> Submit(absl::AnyInvocable<absl::StatusOr<T>() &&> work,
112 thread::TreeOptions tree_options = {});
113
125template <typename T>
126class Future {
127 public:
128 Future() = default;
129
131 [[nodiscard]] bool valid() const { return state_ != nullptr; }
132
134 [[nodiscard]] bool IsReady() const {
135 if (state_ == nullptr) {
136 return false;
137 }
138 thread::MutexLock lock(&state_->mu);
139 return state_->ready;
140 }
141
147 absl::Status Cancel() const {
148 if (state_ == nullptr) {
149 return absl::FailedPreconditionError("Future is not valid");
150 }
151 std::function<void()> cancel;
152 {
153 thread::MutexLock lock(&state_->mu);
154 if (state_->ready) {
155 return absl::OkStatus();
156 }
157 cancel = state_->cancel;
158 }
159
160 if (cancel == nullptr) {
161 return absl::UnimplementedError(
162 "This Future does not have a cancellation source");
163 }
164
165 return exception_guard::Attempt([&] { cancel(); }, "Future cancellation");
166 }
167
174 absl::StatusOr<T> Await(absl::Time deadline = absl::InfiniteFuture()) const {
175 if (state_ == nullptr) {
176 return absl::FailedPreconditionError("Future is not valid");
177 }
178 {
179 thread::MutexLock lock(&state_->mu);
180 if (state_->ready) {
181 return *state_->result;
182 }
183 }
184
185 // A dynamic A11 fiber must yield its worker instead of blocking it. Plain
186 // external threads use a cv variable and do not need a fiber
187 // scheduler installed merely to wait for an A11 operation.
188 if (thread::GetPerThreadFiberPtr() != nullptr) {
189 const int selected = thread::SelectUntil(
190 deadline, {thread::OnCancel(), state_->event.OnEvent()});
191 if (selected == 0) {
192 return absl::CancelledError("Future wait cancelled");
193 }
194 if (selected < 0) {
195 return absl::DeadlineExceededError(
196 "Future was not ready before deadline");
197 }
198 } else {
199 thread::MutexLock lock(&state_->mu);
200 while (!state_->ready) {
201 if (state_->cv.WaitWithDeadline(&state_->mu, deadline) &&
202 !state_->ready) {
203 return absl::DeadlineExceededError(
204 "Future was not ready before deadline");
205 }
206 }
207 return *state_->result;
208 }
209
210 thread::MutexLock lock(&state_->mu);
211 if (!state_->ready) {
212 return absl::InternalError("Future wake-up did not publish a result");
213 }
214 return *state_->result;
215 }
216
224 absl::AnyInvocable<void(const absl::StatusOr<T>&)> callback) const {
225 if (callback == nullptr) {
226 return;
227 }
228 if (state_ == nullptr) {
229 const absl::StatusOr<T> invalid =
230 absl::FailedPreconditionError("Future is not valid");
231 internal::InvokeFutureCallback<T>(std::move(callback), invalid);
232 return;
233 }
234 const absl::StatusOr<T>* absl_nullable ready_result = nullptr;
235 {
236 thread::MutexLock lock(&state_->mu);
237 if (!state_->ready) {
238 state_->callbacks.push_back(std::move(callback));
239 return;
240 }
241 ready_result = &*state_->result;
242 }
243 internal::InvokeFutureCallback<T>(std::move(callback), *ready_result);
244 }
245
246 private:
247 void SetCancellationCallbackForExecutor(std::function<void()> cancel) {
248 if (state_ == nullptr) {
249 return;
250 }
251 thread::MutexLock lock(&state_->mu);
252 if (!state_->ready) {
253 state_->cancel = std::move(cancel);
254 }
255 }
256
257 explicit Future(std::shared_ptr<internal::FutureState<T>> state)
258 : state_(std::move(state)) {}
259
260 std::shared_ptr<internal::FutureState<T>> state_;
261
262 friend class Promise<T>;
263 template <typename U>
264 friend Future<U> Submit(absl::AnyInvocable<absl::StatusOr<U>() &&> work,
265 thread::TreeOptions tree_options);
266 template <typename U>
268 absl::AnyInvocable<absl::StatusOr<U>() &&> work,
269 std::function<void()> cancellation_hook,
270 thread::TreeOptions tree_options);
271};
272
281template <typename T>
282class Promise {
283 public:
284 Promise() : state_(std::make_shared<internal::FutureState<T>>()) {}
285
286 Promise(const Promise&) = delete;
287 Promise& operator=(const Promise&) = delete;
288
290 Promise(Promise&& other) noexcept : state_(std::move(other.state_)) {}
291
293 Promise& operator=(Promise&& other) noexcept {
294 if (this != &other) {
295 Abandon();
296 state_ = std::move(other.state_);
297 }
298 return *this;
299 }
300
301 ~Promise() { Abandon(); }
302
304 [[nodiscard]] Future<T> future() const { return Future<T>(state_); }
305
307 absl::Status SetCancellationCallback(std::function<void()> cancel) {
308 if (state_ == nullptr) {
309 return absl::FailedPreconditionError("Promise is not valid");
310 }
311 thread::MutexLock lock(&state_->mu);
312 if (state_->ready) {
313 return absl::FailedPreconditionError("Promise is already complete");
314 }
315 state_->cancel = std::move(cancel);
316 return absl::OkStatus();
317 }
318
320 absl::Status SetValue(T value) {
321 return SetResult(absl::StatusOr<T>(std::move(value)));
322 }
323
325 absl::Status SetStatus(absl::Status status) {
326 if (status.ok()) {
327 return absl::InvalidArgumentError(
328 "SetStatus requires a non-OK status; use SetValue for success");
329 }
330 absl::StatusOr<T> result;
331 result.AssignStatus(std::move(status));
332 return SetResult(std::move(result));
333 }
334
336 absl::Status SetResult(absl::StatusOr<T> result) {
337 if (state_ == nullptr) {
338 return absl::FailedPreconditionError("Promise is not valid");
339 }
340 std::vector<absl::AnyInvocable<void(const absl::StatusOr<T>&)>> callbacks;
341 const absl::StatusOr<T>* absl_nullable published = nullptr;
342 {
343 thread::MutexLock lock(&state_->mu);
344 if (state_->ready) {
345 return absl::AlreadyExistsError("Promise has already been completed");
346 }
347 state_->result.emplace(std::move(result));
348 state_->ready = true;
349 state_->cancel = {};
350 published = &*state_->result;
351 callbacks.swap(state_->callbacks);
352 }
353 state_->event.Notify();
354 state_->cv.SignalAll();
355 for (auto& callback : callbacks) {
356 internal::InvokeFutureCallback<T>(std::move(callback), *published);
357 }
358 return absl::OkStatus();
359 }
360
361 private:
362 void Abandon() {
363 if (state_ == nullptr) {
364 return;
365 }
366 bool ready = false;
367 {
368 thread::MutexLock lock(&state_->mu);
369 ready = state_->ready;
370 }
371 // Never release the last state owner while its embedded mutex is locked.
372 if (ready) {
373 state_.reset();
374 return;
375 }
376 SetStatus(absl::CancelledError("Promise was abandoned")).IgnoreError();
377 state_.reset();
378 }
379
380 std::shared_ptr<internal::FutureState<T>> state_;
381};
382
384template <typename T>
386 Promise<T> promise;
387 Future<T> future = promise.future();
388 promise.SetResult(std::move(value)).IgnoreError();
389 return future;
390}
391
393template <typename T>
394Future<T> CompletedFuture(absl::StatusOr<T> result) {
395 Promise<T> promise;
396 Future<T> future = promise.future();
397 promise.SetResult(std::move(result)).IgnoreError();
398 return future;
399}
400
402template <typename T>
403Future<T> FailedFuture(absl::Status status) {
404 Promise<T> promise;
405 Future<T> future = promise.future();
406 promise.SetStatus(std::move(status)).IgnoreError();
407 return future;
408}
409
412
414inline Task ReadyTask() {
415 return ReadyFuture(Unit{});
416}
417
419inline Task FailedTask(absl::Status status) {
420 return FailedFuture<Unit>(std::move(status));
421}
422
445template <typename T, typename Fn>
446auto Then(const Future<T>& future, Fn transform) -> Future<
447 typename std::invoke_result_t<Fn, const absl::StatusOr<T>&>::value_type> {
448 using Result = std::invoke_result_t<Fn, const absl::StatusOr<T>&>;
449 using U = typename Result::value_type;
450
451 if (future.IsReady()) {
452 return CompletedFuture<U>(transform(future.Await()));
453 }
454
455 Promise<U> promise;
456 Future<U> continued = promise.future();
457 promise.SetCancellationCallback([future]() mutable { (void)future.Cancel(); })
458 .IgnoreError();
459 future.OnReady(
460 [promise = std::move(promise), transform = std::move(transform)](
461 const absl::StatusOr<T>& result) mutable {
462 promise.SetResult(transform(result)).IgnoreError();
463 });
464 return continued;
465}
466
467} // namespace a11
468
469#endif // A11_CONCURRENCY_FUTURE_H_
Shared handle to one asynchronous result.
Definition future.h:126
absl::StatusOr< T > Await(absl::Time deadline=absl::InfiniteFuture()) const
Wait for and return the result up to an absolute deadline.
Definition future.h:174
bool IsReady() const
Whether the producer has published either a value or an error.
Definition future.h:134
Future()=default
friend Future< U > Submit(absl::AnyInvocable< absl::StatusOr< U >() && > work, thread::TreeOptions tree_options)
absl::Status Cancel() const
Request cancellation from the operation producing this result.
Definition future.h:147
friend Future< U > SubmitWithCancellationHook(absl::AnyInvocable< absl::StatusOr< U >() && > work, std::function< void()> cancellation_hook, thread::TreeOptions tree_options)
bool valid() const
Whether this handle refers to shared completion state.
Definition future.h:131
void OnReady(absl::AnyInvocable< void(const absl::StatusOr< T > &)> callback) const
Run callback once when the result becomes available.
Definition future.h:223
Move-only producer for a Future result.
Definition future.h:282
absl::Status SetValue(T value)
Complete successfully with value.
Definition future.h:320
absl::Status SetStatus(absl::Status status)
Complete with a non-OK status.
Definition future.h:325
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
Promise()
Definition future.h:284
Promise(Promise &&other) noexcept
Transfer responsibility for completing or abandoning the shared state.
Definition future.h:290
~Promise()
Definition future.h:301
absl::Status SetCancellationCallback(std::function< void()> cancel)
Install the operation invoked when a consumer calls Future::Cancel().
Definition future.h:307
Promise & operator=(const Promise &)=delete
Promise & operator=(Promise &&other) noexcept
Abandon this state, then take responsibility for other's state.
Definition future.h:293
Promise(const Promise &)=delete
std::string value
Definition discover.cc:114
Turns what a caller's callable throws into a Status, at the boundary.
thread::Mutex mu
Definition executor.cc:32
absl::Status Attempt(Callable &&callable, std::string_view what)
Runs callable, reporting anything it throws as a Status.
Definition exception_guard.h:114
void InvokeFutureCallback(absl::AnyInvocable< void(const absl::StatusOr< T > &)> callback, const absl::StatusOr< T > &result)
Definition future.h:80
Definition action.cc:63
Task FailedTask(absl::Status status)
Return an already-failed Task.
Definition future.h:419
auto Then(const Future< T > &future, Fn transform) -> Future< typename std::invoke_result_t< Fn, const absl::StatusOr< T > & >::value_type >
Continue with transform once future completes, without a fiber.
Definition future.h:446
Future< T > FailedFuture(absl::Status status)
Return an already-failed future containing status.
Definition future.h:403
Future< T > ReadyFuture(T value)
Return an already-successful future containing value.
Definition future.h:385
Task ReadyTask()
Return an already-successful Task.
Definition future.h:414
Future< T > CompletedFuture(absl::StatusOr< T > result)
Return an already-completed future containing result.
Definition future.h:394
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
Empty success value used by Future<Unit> operations that return no data.
Definition future.h:55
friend bool operator==(Unit, Unit)=default