Symbian platform (C++)
SDK native APIs, runtime and tooling
Loading...
Searching...
No Matches
event_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_EVENT_EXECUTOR_H_
5#define SYMBIAN_CONCURRENCY_EVENT_EXECUTOR_H_
6
7#include <chrono>
8#include <cstddef>
9#include <functional>
10#include <memory>
11#include <mutex>
12#include <utility>
13
14#include "absl/status/status.h"
15#include "absl/status/statusor.h"
20#include "thread/fiber.h"
21
22namespace symbian::concurrency {
23
24// The event thread's sole native request-semaphore consumer. Window Server
25// statuses remain owned by the caller and must be inspected on every turn.
26// Native adapters remove completed requests before invoking inline callbacks.
28 public:
29 EventExecutor() : scheduler_(&policy_) {}
30
31 EventExecutor(const EventExecutor&) = delete;
33
35
36 absl::Status Open() {
37 if (opened_ || closed_) {
38 return absl::FailedPreconditionError("Event executor already opened");
39 }
40 const int result = timers_.Open();
41 if (result != 0) {
42 return symbian::StatusFromNativeError(result, "Event executor wake");
43 }
44 auto wake = timers_.WakeCallback();
45 policy_.SetWake(wake);
46 auto mailbox = std::make_shared<EventMailbox>(std::move(wake));
47 {
48 std::lock_guard lock(dispatch_mu_);
49 mailbox_ = std::move(mailbox);
50 opened_ = true;
51 }
52 return absl::OkStatus();
53 }
54
55 // Explicit event affinity. A11 Post/PostAt retain their worker-pool meaning.
56 absl::Status DispatchToEvent(std::function<void()> callback) {
57 std::shared_ptr<EventMailbox> mailbox;
58 {
59 std::lock_guard lock(dispatch_mu_);
60 if (!opened_ || closed_) {
61 return absl::FailedPreconditionError("Event executor is closed");
62 }
63 mailbox = mailbox_;
64 }
65 return mailbox->Enqueue(std::move(callback));
66 }
67
68 Task ScheduleAfter(absl::Duration delay) {
69 return timers_.ScheduleAfter(delay);
70 }
71
72 Task ScheduleAt(absl::Time deadline) { return timers_.ScheduleAt(deadline); }
73
74 absl::Status OpenProperty(int category, unsigned int key) {
75 if (!opened_ || closed_ || property_) {
76 return absl::FailedPreconditionError("Property owner unavailable");
77 }
78 auto property = std::make_unique<PropertyWatch>(timers_.WakeCallback());
79 absl::Status status = property->Open(category, key);
80 if (status.ok()) {
81 property_ = std::move(property);
82 }
83 return status;
84 }
85
87 if (!property_ || closed_) {
88 return FailedFuture<int>(
89 absl::FailedPreconditionError("Property owner is closed"));
90 }
91 return property_->Next();
92 }
93
94 absl::Status SetProperty(int value) {
95 if (!property_ || closed_) {
96 return absl::FailedPreconditionError("Property owner is closed");
97 }
98 return property_->Set(value);
99 }
100
101 thread::Scheduler& fibers() { return scheduler_; }
102
103 // Lazily create one shared worker for explicit compute placement. The
104 // event thread calls this during setup; ordinary DispatchReady turns never
105 // create an OS thread or offload callbacks implicitly.
106 absl::StatusOr<WorkerExecutor*> workers() {
107 std::lock_guard lock(dispatch_mu_);
108 if (!opened_ || closed_) {
109 return absl::FailedPreconditionError("Event executor is closed");
110 }
111 if (!worker_) {
112 worker_ = std::make_unique<WorkerExecutor>();
113 }
114 return worker_.get();
115 }
116
117 // Bound each source independently. Call again when HasReady() is true;
118 // native and mailbox adapters resignal when a turn leaves work behind.
119 absl::Status DispatchReady(std::size_t budget = 64) {
120 if (!opened_ || closed_) {
121 return absl::FailedPreconditionError("Event executor is closed");
122 }
123 if (budget == 0) {
124 return absl::OkStatus();
125 }
126 if (property_) {
127 property_->DispatchReady();
128 }
129 timers_.DispatchReady(budget);
130 mailbox_->DispatchReady(budget);
131 absl::Status status = scheduler_.RunReady(budget);
132 if (!status.ok()) {
133 return status;
134 }
135 return RearmFiberDeadline();
136 }
137
138 bool HasReady() const {
139 return opened_ && !closed_ &&
140 (scheduler_.HasReady() || mailbox_->Pending() != 0 ||
141 (fiber_alarm_.valid() && fiber_alarm_.IsReady()));
142 }
143
144 // The caller checks Window Server statuses before this call. The request
145 // semaphore preserves completion signals arriving after that check.
146 void Park() const {
147 if (opened_ && !closed_ && !HasReady()) {
148 timers_.Park();
149 }
150 }
151
152 void Close() {
153 {
154 std::lock_guard lock(dispatch_mu_);
155 if (closed_) {
156 return;
157 }
158 closed_ = true;
159 }
160 if (property_) {
161 property_->Close();
162 property_.reset();
163 }
164 if (worker_) {
165 worker_->Close();
166 }
167 if (fiber_alarm_.valid()) {
168 fiber_alarm_.Cancel();
169 }
170 timers_.Close();
171 if (mailbox_) {
172 mailbox_->Close();
173 mailbox_.reset();
174 }
175 policy_.SetWake({});
176 }
177
178 private:
179 class WakePolicy final : public thread::SchedulerPolicy {
180 public:
181 std::size_t PickNext(std::span<thread::Fiber* const>) override { return 0; }
182
183 void NotifyReady() noexcept override {
184 std::function<void()> callback;
185 {
186 std::lock_guard lock(mu_);
187 callback = wake_;
188 }
189 if (callback) {
190 callback();
191 }
192 }
193
194 void SetWake(std::function<void()> wake) {
195 std::lock_guard lock(mu_);
196 wake_ = std::move(wake);
197 }
198
199 private:
200 std::mutex mu_;
201 std::function<void()> wake_;
202 };
203
204 absl::Status RearmFiberDeadline() {
205 const auto next = scheduler_.NextDeadline();
206 if (next == fiber_deadline_ &&
207 (!fiber_alarm_.valid() || !fiber_alarm_.IsReady())) {
208 return absl::OkStatus();
209 }
210 if (fiber_alarm_.valid() && !fiber_alarm_.IsReady()) {
211 fiber_alarm_.Cancel();
212 }
213 fiber_alarm_ = {};
214 fiber_deadline_ = next;
215 if (next == std::chrono::steady_clock::time_point::max()) {
216 return absl::OkStatus();
217 }
218 const auto now = std::chrono::steady_clock::now();
219 const auto delay =
220 next <= now ? std::chrono::nanoseconds::zero() : next - now;
221 fiber_alarm_ = timers_.ScheduleAfter(absl::Nanoseconds(delay.count()));
222 if (auto result = fiber_alarm_.ResultIfReady(); result && !result->ok()) {
223 return result->status();
224 }
225 return absl::OkStatus();
226 }
227
228 TimerPump timers_;
229 WakePolicy policy_;
230 thread::Scheduler scheduler_;
231 std::mutex dispatch_mu_;
232 std::shared_ptr<EventMailbox> mailbox_;
233 std::unique_ptr<PropertyWatch> property_;
234 std::unique_ptr<WorkerExecutor> worker_;
235 Task fiber_alarm_;
236 std::chrono::steady_clock::time_point fiber_deadline_ =
237 std::chrono::steady_clock::time_point::max();
238 bool opened_ = false;
239 bool closed_ = false;
240};
241
242} // namespace symbian::concurrency
243
244#endif // SYMBIAN_CONCURRENCY_EVENT_EXECUTOR_H_
Definition bounded_channel.h:38
Definition event_executor.h:27
EventExecutor(const EventExecutor &)=delete
absl::StatusOr< WorkerExecutor * > workers()
Definition event_executor.h:106
bool HasReady() const
Definition event_executor.h:138
absl::Status SetProperty(int value)
Definition event_executor.h:94
void Park() const
Definition event_executor.h:146
EventExecutor & operator=(const EventExecutor &)=delete
absl::Status Open()
Definition event_executor.h:36
Task ScheduleAfter(absl::Duration delay)
Definition event_executor.h:68
absl::Status DispatchToEvent(std::function< void()> callback)
Definition event_executor.h:56
EventExecutor()
Definition event_executor.h:29
~EventExecutor()
Definition event_executor.h:34
absl::Status DispatchReady(std::size_t budget=64)
Definition event_executor.h:119
void Close()
Definition event_executor.h:152
Task ScheduleAt(absl::Time deadline)
Definition event_executor.h:72
Future< int > NextProperty()
Definition event_executor.h:86
thread::Scheduler & fibers()
Definition event_executor.h:101
absl::Status OpenProperty(int category, unsigned int key)
Definition event_executor.h:74
bool IsReady() const
Definition future.h:136
bool Cancel() const
Definition future.h:146
std::optional< absl::StatusOr< T > > ResultIfReady() const
Definition future.h:166
bool valid() const
Definition future.h:134
Task ScheduleAt(absl::Time deadline)
Definition timer_pump.h:92
std::size_t DispatchReady(std::size_t budget=64)
Definition timer_pump.h:151
void Close()
Definition timer_pump.h:225
void Park() const
Definition timer_pump.h:219
std::function< void()> WakeCallback() const
Definition timer_pump.h:82
Task ScheduleAfter(absl::Duration delay)
Definition timer_pump.h:103
int Open()
Definition timer_pump.h:68
Definition fiber.h:36
Definition fiber.h:46
bool HasReady() const
Definition fiber.cc:142
std::chrono::steady_clock::time_point NextDeadline() const
Definition fiber.cc:156
absl::Status RunReady(std::size_t max_turns)
Definition fiber.cc:167
Definition bounded_channel.h:32
Future< Unit > Task
Definition future.h:270
absl::Status StatusFromNativeError(int native_code, absl::string_view operation)
Definition native_status.h:42
libusb_context * value
Definition usb.cc:35