Symbian platform (C++)
SDK native APIs, runtime and tooling
Loading...
Searching...
No Matches
worker_executor.h
Go to the documentation of this file.
1// Copyright 2026 The Symbian SDK Authors.
2// Licensed under the Apache License, Version 2.0.
3
4#ifndef SYMBIAN_CONCURRENCY_WORKER_EXECUTOR_H_
5#define SYMBIAN_CONCURRENCY_WORKER_EXECUTOR_H_
6
7#include <cstddef>
8#include <functional>
9#include <memory>
10#include <type_traits>
11#include <utility>
12
13#include "absl/functional/any_invocable.h"
14#include "absl/status/status.h"
16
17namespace symbian::concurrency {
18
28 private:
29 struct State;
30
31 public:
32 using Work = absl::AnyInvocable<void() &&>;
33
35 public:
37 absl::Status Post(Work work) const;
38
39 private:
40 friend class WorkerExecutor;
41
42 explicit DispatchHandle(std::weak_ptr<State> state)
43 : state_(std::move(state)) {}
44
45 std::weak_ptr<State> state_;
46 };
47
48 explicit WorkerExecutor(std::size_t max_outstanding = 128);
52
53 DispatchHandle handle() const { return DispatchHandle(state_); }
54
56 absl::Status Post(Work work) { return handle().Post(std::move(work)); }
57
65 Task PostFiber(Work work, std::size_t stack_bytes = 16 * 1024);
67 void Close();
69 Task Finish();
70
71 private:
72 std::shared_ptr<State> state_;
73};
74
75namespace internal {
76absl::StatusOr<WorkerExecutor*> WorkerFor(EventExecutor& executor);
77} // namespace internal
78
79// The transform is posted as a stackless callback. The completing thread
80// only enqueues a cheap Future handle; the worker copies the result there.
81// Queue rejection completes the returned Future with an error. Then()
82// retains its A11 inline semantics.
83template <typename T, typename Fn>
85 -> Future<typename std::invoke_result_t<
86 Fn, const absl::StatusOr<T>&>::value_type> {
87 using U =
88 typename std::invoke_result_t<Fn, const absl::StatusOr<T>&>::value_type;
89 auto promise = std::make_shared<Promise<U>>();
90 Future<U> continued = promise->future();
91 promise->SetCancellationCallback([future] { future.Cancel(); });
92 auto handle = worker.handle();
93 future.OnReady([promise = std::move(promise),
94 transform = std::move(transform), handle,
95 future](const absl::StatusOr<T>&) mutable {
96 absl::Status posted = handle.Post(
97 [promise, transform = std::move(transform), future]() mutable {
98 auto result = future.ResultIfReady();
99 if (!result) {
100 promise->SetError(absl::InternalError(
101 "Ready Future lost its result before worker dispatch"));
102 return;
103 }
104 promise->SetResult(transform(*result));
105 });
106 if (!posted.ok()) {
107 promise->SetError(std::move(posted));
108 }
109 });
110 return continued;
111}
112
113template <typename T>
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> {
118 using U =
119 typename std::invoke_result_t<Fn, const absl::StatusOr<T>&>::value_type;
120 auto worker = internal::WorkerFor(executor);
121 if (!worker.ok()) {
122 return FailedFuture<U>(worker.status());
123 }
124 return ThenOn(*this, **worker, std::move(transform));
125}
126
127} // namespace symbian::concurrency
128
129#endif // SYMBIAN_CONCURRENCY_WORKER_EXECUTOR_H_
Definition bounded_channel.h:38
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