4#ifndef SYMBIAN_CONCURRENCY_EVENT_EXECUTOR_H_
5#define SYMBIAN_CONCURRENCY_EVENT_EXECUTOR_H_
14#include "absl/status/status.h"
15#include "absl/status/statusor.h"
20#include "thread/fiber.h"
37 if (opened_ || closed_) {
38 return absl::FailedPreconditionError(
"Event executor already opened");
40 const int result = timers_.
Open();
45 policy_.SetWake(
wake);
46 auto mailbox = std::make_shared<EventMailbox>(std::move(
wake));
48 std::lock_guard lock(dispatch_mu_);
52 return absl::OkStatus();
57 std::shared_ptr<EventMailbox>
mailbox;
59 std::lock_guard lock(dispatch_mu_);
60 if (!opened_ || closed_) {
61 return absl::FailedPreconditionError(
"Event executor is closed");
65 return mailbox->Enqueue(std::move(callback));
75 if (!opened_ || closed_ || property_) {
76 return absl::FailedPreconditionError(
"Property owner unavailable");
78 auto property = std::make_unique<PropertyWatch>(timers_.
WakeCallback());
79 absl::Status status =
property->Open(category, key);
81 property_ = std::move(property);
87 if (!property_ || closed_) {
89 absl::FailedPreconditionError(
"Property owner is closed"));
91 return property_->Next();
95 if (!property_ || closed_) {
96 return absl::FailedPreconditionError(
"Property owner is closed");
98 return property_->Set(
value);
107 std::lock_guard lock(dispatch_mu_);
108 if (!opened_ || closed_) {
109 return absl::FailedPreconditionError(
"Event executor is closed");
112 worker_ = std::make_unique<WorkerExecutor>();
114 return worker_.get();
120 if (!opened_ || closed_) {
121 return absl::FailedPreconditionError(
"Event executor is closed");
124 return absl::OkStatus();
127 property_->DispatchReady();
130 mailbox_->DispatchReady(
budget);
135 return RearmFiberDeadline();
139 return opened_ && !closed_ &&
140 (scheduler_.
HasReady() || mailbox_->Pending() != 0 ||
147 if (opened_ && !closed_ && !
HasReady()) {
154 std::lock_guard lock(dispatch_mu_);
167 if (fiber_alarm_.
valid()) {
181 std::size_t PickNext(std::span<thread::Fiber* const>)
override {
return 0; }
183 void NotifyReady() noexcept
override {
184 std::function<void()> callback;
186 std::lock_guard lock(mu_);
194 void SetWake(std::function<
void()> wake) {
195 std::lock_guard lock(mu_);
196 wake_ = std::move(wake);
201 std::function<void()> wake_;
204 absl::Status RearmFiberDeadline() {
206 if (next == fiber_deadline_ &&
208 return absl::OkStatus();
214 fiber_deadline_ = next;
215 if (next == std::chrono::steady_clock::time_point::max()) {
216 return absl::OkStatus();
218 const auto now = std::chrono::steady_clock::now();
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();
225 return absl::OkStatus();
231 std::mutex dispatch_mu_;
232 std::shared_ptr<EventMailbox> mailbox_;
233 std::unique_ptr<PropertyWatch> property_;
234 std::unique_ptr<WorkerExecutor> worker_;
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;
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
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