27#ifndef A11_CONCURRENCY_FUTURE_H_
28#define A11_CONCURRENCY_FUTURE_H_
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>
47#include "thread/boost_primitives.h"
48#include "thread/fiber.h"
49#include "thread/select.h"
50#include "thread/selectables.h"
69 mutable thread::Mutex
mu;
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
81 absl::AnyInvocable<
void(
const absl::StatusOr<T>&)> callback,
82 const absl::StatusOr<T>& result) {
87 [&] { callback(result); },
"Future completion callback");
89 LOG(ERROR) << raised.message();
106 absl::AnyInvocable<absl::StatusOr<T>() &&> work,
107 std::function<
void()> cancellation_hook,
108 thread::TreeOptions tree_options = {});
111Future<T>
Submit(absl::AnyInvocable<absl::StatusOr<T>() &&> work,
112 thread::TreeOptions tree_options = {});
131 [[nodiscard]]
bool valid()
const {
return state_ !=
nullptr; }
135 if (state_ ==
nullptr) {
138 thread::MutexLock lock(&state_->mu);
139 return state_->ready;
148 if (state_ ==
nullptr) {
149 return absl::FailedPreconditionError(
"Future is not valid");
151 std::function<void()> cancel;
153 thread::MutexLock lock(&state_->mu);
155 return absl::OkStatus();
157 cancel = state_->cancel;
160 if (cancel ==
nullptr) {
161 return absl::UnimplementedError(
162 "This Future does not have a cancellation source");
174 absl::StatusOr<T>
Await(absl::Time deadline = absl::InfiniteFuture())
const {
175 if (state_ ==
nullptr) {
176 return absl::FailedPreconditionError(
"Future is not valid");
179 thread::MutexLock lock(&state_->mu);
181 return *state_->result;
188 if (thread::GetPerThreadFiberPtr() !=
nullptr) {
189 const int selected = thread::SelectUntil(
190 deadline, {thread::OnCancel(), state_->event.OnEvent()});
192 return absl::CancelledError(
"Future wait cancelled");
195 return absl::DeadlineExceededError(
196 "Future was not ready before deadline");
199 thread::MutexLock lock(&state_->mu);
200 while (!state_->ready) {
201 if (state_->cv.WaitWithDeadline(&state_->mu, deadline) &&
203 return absl::DeadlineExceededError(
204 "Future was not ready before deadline");
207 return *state_->result;
210 thread::MutexLock lock(&state_->mu);
211 if (!state_->ready) {
212 return absl::InternalError(
"Future wake-up did not publish a result");
214 return *state_->result;
224 absl::AnyInvocable<
void(
const absl::StatusOr<T>&)> callback)
const {
225 if (callback ==
nullptr) {
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);
234 const absl::StatusOr<T>* absl_nullable ready_result =
nullptr;
236 thread::MutexLock lock(&state_->mu);
237 if (!state_->ready) {
238 state_->callbacks.push_back(std::move(callback));
241 ready_result = &*state_->result;
243 internal::InvokeFutureCallback<T>(std::move(callback), *ready_result);
247 void SetCancellationCallbackForExecutor(std::function<
void()> cancel) {
248 if (state_ ==
nullptr) {
251 thread::MutexLock lock(&state_->mu);
252 if (!state_->ready) {
253 state_->cancel = std::move(cancel);
257 explicit Future(std::shared_ptr<internal::FutureState<T>> state)
258 : state_(std::move(state)) {}
260 std::shared_ptr<internal::FutureState<T>> state_;
263 template <
typename U>
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);
284 Promise() : state_(std::make_shared<internal::FutureState<T>>()) {}
294 if (
this != &other) {
296 state_ = std::move(other.state_);
308 if (state_ ==
nullptr) {
309 return absl::FailedPreconditionError(
"Promise is not valid");
311 thread::MutexLock lock(&state_->mu);
313 return absl::FailedPreconditionError(
"Promise is already complete");
315 state_->cancel = std::move(cancel);
316 return absl::OkStatus();
327 return absl::InvalidArgumentError(
328 "SetStatus requires a non-OK status; use SetValue for success");
330 absl::StatusOr<T> result;
331 result.AssignStatus(std::move(status));
337 if (state_ ==
nullptr) {
338 return absl::FailedPreconditionError(
"Promise is not valid");
340 std::vector<absl::AnyInvocable<void(
const absl::StatusOr<T>&)>> callbacks;
341 const absl::StatusOr<T>* absl_nullable published =
nullptr;
343 thread::MutexLock lock(&state_->mu);
345 return absl::AlreadyExistsError(
"Promise has already been completed");
347 state_->result.emplace(std::move(result));
348 state_->ready =
true;
350 published = &*state_->result;
351 callbacks.swap(state_->callbacks);
353 state_->event.Notify();
354 state_->cv.SignalAll();
355 for (
auto& callback : callbacks) {
356 internal::InvokeFutureCallback<T>(std::move(callback), *published);
358 return absl::OkStatus();
363 if (state_ ==
nullptr) {
368 thread::MutexLock lock(&state_->mu);
369 ready = state_->ready;
376 SetStatus(absl::CancelledError(
"Promise was abandoned")).IgnoreError();
380 std::shared_ptr<internal::FutureState<T>> state_;
397 promise.
SetResult(std::move(result)).IgnoreError();
406 promise.
SetStatus(std::move(status)).IgnoreError();
420 return FailedFuture<Unit>(std::move(status));
445template <
typename T,
typename Fn>
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;
451 if (future.IsReady()) {
452 return CompletedFuture<U>(transform(future.Await()));
460 [promise = std::move(promise), transform = std::move(transform)](
461 const absl::StatusOr<T>& result)
mutable {
462 promise.
SetResult(transform(result)).IgnoreError();
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
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
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