5#ifndef SYMBIAN_CONCURRENCY_TIMER_PUMP_H_
6#define SYMBIAN_CONCURRENCY_TIMER_PUMP_H_
18#include "absl/time/time.h"
28 absl::Microseconds(std::numeric_limits<std::int32_t>::max());
30 return static_cast<int>(
31 absl::ToInt64Microseconds(absl::Ceil(
slice, absl::Microseconds(1))));
35 absl::Time deadline, absl::Time
now) {
36 if (deadline == absl::InfiniteFuture()) {
37 return absl::InfiniteDuration();
39 if (deadline <=
now) {
40 return absl::ZeroDuration();
42 const absl::Duration remaining = deadline -
now;
43 if (remaining == absl::InfiniteDuration()) {
44 return absl::OutOfRangeError(
"Absolute deadline is too distant");
69 if (closed_ || wake_) {
72 auto wake = std::make_shared<WakeTarget>();
73 const int result =
wake->Open();
75 wake_ = std::move(
wake);
83 return [
target = std::weak_ptr<WakeTarget>(wake_)] {
93 if (wall_now_ ==
nullptr) {
94 return FailedTask(absl::FailedPreconditionError(
"No wall clock"));
97 if (!remaining.ok()) {
104 if (closed_ || !wake_) {
105 return FailedTask(absl::FailedPreconditionError(
"Timer pump is closed"));
107 if (
delay < absl::ZeroDuration()) {
109 absl::InvalidArgumentError(
"Timer delay must be nonnegative"));
111 if (entries_.size() >= max_pending_) {
112 return FailedTask(absl::ResourceExhaustedError(
"Timer pump is full"));
114 auto entry = std::make_shared<Entry>();
115 entry->remaining =
delay;
116 if (
delay != absl::InfiniteDuration()) {
117 const int opened = entry->timer.Open();
123 Task task = entry->promise.future();
124 entry->promise.SetCancellationCallback(
125 [entry = std::weak_ptr<Entry>(entry),
126 wake = std::weak_ptr<WakeTarget>(wake_)] {
127 if (
auto owned = entry.lock()) {
128 owned->cancel_requested.store(1, std::memory_order_release);
129 if (auto target = wake.lock()) {
136 entries_.push_back(entry);
137 if (
delay != absl::InfiniteDuration()) {
138 const int started = ArmNext(*entry);
141 entry->timer.Close();
142 entry->promise.SetError(
152 if (closed_ || !wake_ ||
budget == 0) {
156 std::vector<std::pair<std::shared_ptr<Entry>,
int>> completed;
157 for (
auto it = entries_.begin();
it != entries_.end();) {
159 if (completed.size() ==
budget) {
163 if (entry->remaining == absl::InfiniteDuration()) {
164 if (entry->cancel_requested.load(std::memory_order_acquire) == 0) {
169 it = entries_.erase(
it);
172 if (!entry->cancel_submitted &&
173 entry->cancel_requested.load(std::memory_order_acquire) != 0) {
174 entry->timer.Cancel();
175 entry->cancel_submitted =
true;
177 if (!entry->timer.IsReady()) {
181 int code = entry->timer.Result();
184 std::chrono::duration_cast<std::chrono::nanoseconds>(
185 std::chrono::steady_clock::now() - entry->armed_at);
186 entry->remaining -= absl::Nanoseconds(
elapsed.count());
187 if (entry->remaining > absl::ZeroDuration()) {
188 if (entry->cancel_requested.load(std::memory_order_acquire) != 0) {
191 const int rearmed = ArmNext(*entry);
200 entry->timer.Close();
201 completed.emplace_back(entry,
code);
202 it = entries_.erase(
it);
204 for (
auto& [entry,
code] : completed) {
206 entry->promise.SetValue(
Unit{});
207 }
else if (
code == -3) {
208 entry->promise.SetError(absl::CancelledError(
"Timer cancelled"));
210 entry->promise.SetError(
214 return completed.size();
220 if (!closed_ && wake_) {
233 auto entries = std::move(entries_);
234 for (
auto& entry : entries) {
235 entry->timer.Close();
236 entry->promise.SetError(absl::CancelledError(
"Timer pump closed"));
247 if (native_ && !pending_) {
268 bool pending_ =
false;
273 Promise<Unit> promise;
274 std::atomic<int> cancel_requested{0};
275 bool cancel_submitted =
false;
276 absl::Duration remaining = absl::ZeroDuration();
277 std::chrono::steady_clock::time_point armed_at;
280 static int ArmNext(Entry& entry) {
281 entry.armed_at = std::chrono::steady_clock::now();
285 std::shared_ptr<WakeTarget> wake_;
286 std::vector<std::shared_ptr<Entry>> entries_;
287 std::size_t max_pending_;
289 bool closed_ =
false;
void SymbianRuntimeWaitForAnyRequest()
Definition sdk_timer.cc:77
void SymbianRuntimeWakeSignal(SymbianRuntimeWakeState *state)
Definition sdk_timer.cc:105
int SymbianRuntimeWakeCreate(SymbianRuntimeWakeState **state)
Definition sdk_timer.cc:85
void SymbianRuntimeWakeClose(SymbianRuntimeWakeState *state)
Definition sdk_timer.cc:111
Definition bounded_channel.h:38
Definition timer_pump.h:54
Task ScheduleAt(absl::Time deadline)
Definition timer_pump.h:92
std::size_t DispatchReady(std::size_t budget=64)
Definition timer_pump.h:151
~TimerPump()
Definition timer_pump.h:66
absl::Time(*)() WallClock
Definition timer_pump.h:57
TimerPump & operator=(const TimerPump &)=delete
static constexpr std::size_t kDefaultMaxPending
Definition timer_pump.h:56
TimerPump(std::size_t max_pending=kDefaultMaxPending, WallClock wall_now=&absl::Now)
Definition timer_pump.h:59
void Close()
Definition timer_pump.h:225
void Park() const
Definition timer_pump.h:219
TimerPump(const TimerPump &)=delete
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 boost_primitives.h:95
Definition boost_primitives.h:34
std::vector< uint32_t > code
Definition e32.cc:70
int TimerSliceMicroseconds(absl::Duration remaining)
Definition timer_pump.h:26
absl::StatusOr< absl::Duration > RemainingAtRegistration(absl::Time deadline, absl::Time now)
Definition timer_pump.h:34
Definition bounded_channel.h:32
Task FailedTask(absl::Status error)
Definition future.h:276
constexpr int kCancel
Definition native_status.h:27
absl::Status StatusFromNativeError(int native_code, absl::string_view operation)
Definition native_status.h:42
std::string target
Definition sis.cc:301
Definition sdk_timer.cc:81