4#ifndef SYMBIAN_CONCURRENCY_WORKER_EXECUTOR_H_
5#define SYMBIAN_CONCURRENCY_WORKER_EXECUTOR_H_
13#include "absl/functional/any_invocable.h"
14#include "absl/status/status.h"
45 std::weak_ptr<State> state_;
72 std::shared_ptr<State> state_;
76absl::StatusOr<WorkerExecutor*> WorkerFor(EventExecutor& executor);
83template <
typename T,
typename Fn>
85 ->
Future<
typename std::invoke_result_t<
86 Fn,
const absl::StatusOr<T>&>::value_type> {
88 typename std::invoke_result_t<Fn, const absl::StatusOr<T>&>::value_type;
89 auto promise = std::make_shared<Promise<U>>();
91 promise->SetCancellationCallback([future] { future.Cancel(); });
93 future.OnReady([promise = std::move(promise),
95 future](
const absl::StatusOr<T>&)
mutable {
98 auto result = future.ResultIfReady();
100 promise->SetError(absl::InternalError(
101 "Ready Future lost its result before worker dispatch"));
107 promise->SetError(std::move(
posted));
114template <
typename Fn>
115auto Future<T>::ThenOnWorker(EventExecutor& executor, Fn transform)
const
116 -> Future<
typename std::invoke_result_t<
117 Fn,
const absl::StatusOr<T>&>::value_type> {
119 typename std::invoke_result_t<Fn, const absl::StatusOr<T>&>::value_type;
120 auto worker = internal::WorkerFor(executor);
122 return FailedFuture<U>(worker.status());
124 return ThenOn(*
this, **worker, std::move(transform));
Definition bounded_channel.h:38
Definition worker_executor.h:34
absl::Status Post(Work work) const
Enqueue stackless work or return a full/closed status.
Definition worker_executor.cc:171
Bounded work destination on one guest OS thread.
Definition worker_executor.h:27
Task Finish()
Close and return a task that completes after all jobs drain.
Definition worker_executor.cc:206
absl::Status Post(Work work)
Enqueue stackless work on the single worker thread.
Definition worker_executor.h:56
WorkerExecutor & operator=(const WorkerExecutor &)=delete
WorkerExecutor(const WorkerExecutor &)=delete
DispatchHandle handle() const
Definition worker_executor.h:53
Task PostFiber(Work work, std::size_t stack_bytes=16 *1024)
Run one job on a bounded guest fiber hosted by this worker.
Definition worker_executor.cc:179
~WorkerExecutor()
Definition worker_executor.cc:167
void Close()
Stop accepting work while queued jobs continue to drain.
Definition worker_executor.cc:194
absl::AnyInvocable< void() && > Work
Definition worker_executor.h:32
Definition exception_count.cc:8
Definition bounded_channel.h:32
auto ThenOn(const Future< T > &future, WorkerExecutor &worker, Fn transform) -> Future< typename std::invoke_result_t< Fn, const absl::StatusOr< T > & >::value_type >
Definition worker_executor.h:84
Definition worker_executor.cc:30
libusb_device_handle * handle
Definition usb.cc:79