13#ifndef A11_CONCURRENCY_FUTURE_H_
14#define A11_CONCURRENCY_FUTURE_H_
24#include <absl/base/nullability.h>
25#include <absl/functional/any_invocable.h>
26#include <absl/log/log.h>
27#include <absl/status/status.h>
28#include <absl/status/statusor.h>
29#include <absl/time/clock.h>
30#include <absl/time/time.h>
32#include "thread/boost_primitives.h"
33#include "thread/fiber.h"
34#include "thread/select.h"
35#include "thread/selectables.h"
54 mutable thread::Mutex
mu;
56 bool ready ABSL_GUARDED_BY(mu) =
false;
57 std::optional<absl::StatusOr<T>> result ABSL_GUARDED_BY(mu);
58 std::function<
void()>
cancel ABSL_GUARDED_BY(mu);
59 thread::PermanentEvent event;
60 std::vector<absl::AnyInvocable<
void(
const absl::StatusOr<T>&)>>
callbacks
66 absl::AnyInvocable<
void(
const absl::StatusOr<T>&)>
callback,
67 const absl::StatusOr<T>& result) {
70 }
catch (
const std::exception& error) {
71 LOG(
ERROR) <<
"Future completion callback raised: " << error.what();
73 LOG(
ERROR) <<
"Future completion callback raised a non-standard exception";
90 absl::AnyInvocable<absl::StatusOr<T>() &&>
work,
119 if (state_ ==
nullptr) {
122 thread::MutexLock
lock(&state_->mu);
123 return state_->ready;
132 if (state_ ==
nullptr) {
133 return absl::FailedPreconditionError(
"Future is not valid");
137 thread::MutexLock
lock(&state_->mu);
139 return absl::OkStatus();
145 return absl::UnimplementedError(
146 "This Future does not have a cancellation source");
151 return absl::OkStatus();
152 }
catch (
const std::exception& error) {
153 return absl::UnknownError(error.what());
155 return absl::UnknownError(
156 "Future cancellation raised a non-standard exception");
166 absl::StatusOr<T>
Await(absl::Time deadline = absl::InfiniteFuture())
const {
167 if (state_ ==
nullptr) {
168 return absl::FailedPreconditionError(
"Future is not valid");
171 thread::MutexLock
lock(&state_->mu);
173 return *state_->result;
180 if (thread::GetPerThreadFiberPtr() !=
nullptr) {
181 const int selected = thread::SelectUntil(
182 deadline, {thread::OnCancel(), state_->event.OnEvent()});
184 return absl::CancelledError(
"Future wait cancelled");
187 return absl::DeadlineExceededError(
188 "Future was not ready before deadline");
191 thread::MutexLock
lock(&state_->mu);
192 while (!state_->ready) {
193 if (state_->cv.WaitWithDeadline(&state_->mu, deadline) &&
195 return absl::DeadlineExceededError(
196 "Future was not ready before deadline");
199 return *state_->result;
202 thread::MutexLock
lock(&state_->mu);
203 if (!state_->ready) {
204 return absl::InternalError(
"Future wake-up did not publish a result");
206 return *state_->result;
216 absl::AnyInvocable<
void(
const absl::StatusOr<T>&)>
callback)
const {
220 if (state_ ==
nullptr) {
221 const absl::StatusOr<T>
invalid =
222 absl::FailedPreconditionError(
"Future is not valid");
228 thread::MutexLock
lock(&state_->mu);
229 if (!state_->ready) {
230 state_->callbacks.push_back(std::move(
callback));
239 void SetCancellationCallbackForExecutor(std::function<
void()>
cancel) {
240 if (state_ ==
nullptr) {
243 thread::MutexLock
lock(&state_->mu);
244 if (!state_->ready) {
245 state_->cancel = std::move(
cancel);
249 explicit Future(std::shared_ptr<internal::FutureState<T>> state)
250 : state_(std::
move(state)) {}
252 std::shared_ptr<internal::FutureState<T>> state_;
255 template <
typename U>
258 template <
typename U>
260 absl::AnyInvocable<absl::StatusOr<U>() &&>
work,
286 if (
this != &
other) {
288 state_ = std::move(
other.state_);
300 if (state_ ==
nullptr) {
301 return absl::FailedPreconditionError(
"Promise is not valid");
303 thread::MutexLock
lock(&state_->mu);
305 return absl::FailedPreconditionError(
"Promise is already complete");
307 state_->cancel = std::move(
cancel);
308 return absl::OkStatus();
313 return SetResult(absl::StatusOr<T>(std::move(value)));
319 return absl::InvalidArgumentError(
320 "SetStatus requires a non-OK status; use SetValue for success");
322 absl::StatusOr<T> result;
323 result.AssignStatus(std::move(status));
329 if (state_ ==
nullptr) {
330 return absl::FailedPreconditionError(
"Promise is not valid");
332 std::vector<absl::AnyInvocable<
void(
const absl::StatusOr<T>&)>>
callbacks;
335 thread::MutexLock
lock(&state_->mu);
337 return absl::AlreadyExistsError(
"Promise has already been completed");
339 state_->result.emplace(std::move(result));
340 state_->ready =
true;
345 state_->event.Notify();
346 state_->cv.SignalAll();
350 return absl::OkStatus();
355 if (state_ ==
nullptr) {
360 thread::MutexLock
lock(&state_->mu);
361 ready = state_->ready;
368 SetStatus(absl::CancelledError(
"Promise was abandoned")).IgnoreError();
372 std::shared_ptr<internal::FutureState<T>> state_;
380 promise.
SetResult(std::move(value)).IgnoreError();
389 promise.
SetResult(std::move(result)).IgnoreError();
398 promise.
SetStatus(std::move(status)).IgnoreError();
Shared handle to one asynchronous result.
Definition future.h:110
absl::StatusOr< T > Await(absl::Time deadline=absl::InfiniteFuture()) const
Wait for and return the result up to an absolute deadline.
Definition future.h:166
bool IsReady() const
Whether the producer has published either a value or an error.
Definition future.h:118
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:131
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:115
void OnReady(absl::AnyInvocable< void(const absl::StatusOr< T > &)> callback) const
Run callback once when the result becomes available.
Definition future.h:215
Move-only producer for a Future result.
Definition future.h:274
absl::Status SetValue(T value)
Complete successfully with value.
Definition future.h:312
absl::Status SetStatus(absl::Status status)
Complete with a non-OK status.
Definition future.h:317
absl::Status SetResult(absl::StatusOr< T > result)
Complete with either a value or an error, waking every observer.
Definition future.h:328
Future< T > future() const
Return a consumer handle sharing this promise's completion state.
Definition future.h:296
Promise()
Definition future.h:276
Promise(Promise &&other) noexcept
Transfer responsibility for completing or abandoning the shared state.
Definition future.h:282
~Promise()
Definition future.h:293
absl::Status SetCancellationCallback(std::function< void()> cancel)
Install the operation invoked when a consumer calls Future::Cancel().
Definition future.h:299
Promise & operator=(const Promise &)=delete
Promise & operator=(Promise &&other) noexcept
Abandon this state, then take responsibility for other's state.
Definition future.h:285
Promise(const Promise &)=delete
thread::Mutex mu
Definition executor.cc:20
void InvokeFutureCallback(absl::AnyInvocable< void(const absl::StatusOr< T > &)> callback, const absl::StatusOr< T > &result)
Definition future.h:65
Task FailedTask(absl::Status status)
Return an already-failed Task.
Definition future.h:411
Future< T > FailedFuture(absl::Status status)
Return an already-failed future containing status.
Definition future.h:395
Future< T > ReadyFuture(T value)
Return an already-successful future containing value.
Definition future.h:377
Task ReadyTask()
Return an already-successful Task.
Definition future.h:406
Future< T > CompletedFuture(absl::StatusOr< T > result)
Return an already-completed future containing result.
Definition future.h:386
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:66
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:30
Empty success value used by Future<Unit> operations that return no data.
Definition future.h:40
friend bool operator==(Unit, Unit)=default