Symbian platform (C++)
SDK native APIs, runtime and tooling
Loading...
Searching...
No Matches
timer_pump.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// Native completion adapter for the A11-derived stackless Future profile.
4
5#ifndef SYMBIAN_CONCURRENCY_TIMER_PUMP_H_
6#define SYMBIAN_CONCURRENCY_TIMER_PUMP_H_
7
8#include <atomic>
9#include <chrono>
10#include <cstddef>
11#include <cstdint>
12#include <functional>
13#include <limits>
14#include <memory>
15#include <utility>
16#include <vector>
17
18#include "absl/time/time.h"
22
23namespace symbian::concurrency {
24
25namespace internal {
26inline int TimerSliceMicroseconds(absl::Duration remaining) {
27 const absl::Duration max_slice =
28 absl::Microseconds(std::numeric_limits<std::int32_t>::max());
29 const absl::Duration slice = remaining < max_slice ? remaining : max_slice;
30 return static_cast<int>(
31 absl::ToInt64Microseconds(absl::Ceil(slice, absl::Microseconds(1))));
32}
33
34inline absl::StatusOr<absl::Duration> RemainingAtRegistration(
35 absl::Time deadline, absl::Time now) {
36 if (deadline == absl::InfiniteFuture()) {
37 return absl::InfiniteDuration();
38 }
39 if (deadline <= now) {
40 return absl::ZeroDuration();
41 }
42 const absl::Duration remaining = deadline - now;
43 if (remaining == absl::InfiniteDuration()) {
44 return absl::OutOfRangeError("Absolute deadline is too distant");
45 }
46 return remaining;
47}
48} // namespace internal
49
50// One OS event thread owns this pump and every timer in it. The caller also
51// owns Window Server statuses: inspect and dispatch all ready statuses before
52// Park(), then inspect them again after Park() returns. Park is the only wait
53// in the combined loop; no timer worker or second request consumer is created.
54class TimerPump {
55 public:
56 static constexpr std::size_t kDefaultMaxPending = 64;
57 using WallClock = absl::Time (*)();
58
60 WallClock wall_now = &absl::Now)
61 : max_pending_(max_pending), wall_now_(wall_now) {}
62
63 TimerPump(const TimerPump&) = delete;
64 TimerPump& operator=(const TimerPump&) = delete;
65
67
68 int Open() {
69 if (closed_ || wake_) {
70 return -11; // KErrAlreadyExists.
71 }
72 auto wake = std::make_shared<WakeTarget>();
73 const int result = wake->Open();
74 if (result == 0) {
75 wake_ = std::move(wake);
76 }
77 return result;
78 }
79
80 // Another native request owner can share this event thread's one wait.
81 // The callback stays safe after Close and never consumes the semaphore.
82 std::function<void()> WakeCallback() const {
83 return [target = std::weak_ptr<WakeTarget>(wake_)] {
84 if (auto owned = target.lock()) {
85 owned->Notify();
86 }
87 };
88 }
89
90 // Absolute times are wall times. Convert once at acceptance; later device
91 // clock corrections do not move an already accepted timer.
92 Task ScheduleAt(absl::Time deadline) {
93 if (wall_now_ == nullptr) {
94 return FailedTask(absl::FailedPreconditionError("No wall clock"));
95 }
96 auto remaining = internal::RemainingAtRegistration(deadline, wall_now_());
97 if (!remaining.ok()) {
98 return FailedTask(remaining.status());
99 }
100 return ScheduleAfter(*remaining);
101 }
102
103 Task ScheduleAfter(absl::Duration delay) {
104 if (closed_ || !wake_) {
105 return FailedTask(absl::FailedPreconditionError("Timer pump is closed"));
106 }
107 if (delay < absl::ZeroDuration()) {
108 return FailedTask(
109 absl::InvalidArgumentError("Timer delay must be nonnegative"));
110 }
111 if (entries_.size() >= max_pending_) {
112 return FailedTask(absl::ResourceExhaustedError("Timer pump is full"));
113 }
114 auto entry = std::make_shared<Entry>();
115 entry->remaining = delay;
116 if (delay != absl::InfiniteDuration()) {
117 const int opened = entry->timer.Open();
118 if (opened != 0) {
119 return FailedTask(
120 symbian::StatusFromNativeError(opened, "RTimer create"));
121 }
122 }
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()) {
130 target->Notify();
131 }
132 }
133 });
134 // Ownership is visible to DispatchReady before an immediate native
135 // completion can occur. OnReady never runs until the entry is removed.
136 entries_.push_back(entry);
137 if (delay != absl::InfiniteDuration()) {
138 const int started = ArmNext(*entry);
139 if (started != 0) {
140 entries_.pop_back();
141 entry->timer.Close();
142 entry->promise.SetError(
144 }
145 }
146 return task;
147 }
148
149 // Only the event thread calls DispatchReady. A11 OnReady callbacks may run
150 // inline, so remove finished entries before publishing their results.
151 std::size_t DispatchReady(std::size_t budget = 64) {
152 if (closed_ || !wake_ || budget == 0) {
153 return 0;
154 }
155 wake_->Clear();
156 std::vector<std::pair<std::shared_ptr<Entry>, int>> completed;
157 for (auto it = entries_.begin(); it != entries_.end();) {
158 auto& entry = *it;
159 if (completed.size() == budget) {
160 wake_->Notify();
161 break;
162 }
163 if (entry->remaining == absl::InfiniteDuration()) {
164 if (entry->cancel_requested.load(std::memory_order_acquire) == 0) {
165 ++it;
166 continue;
167 }
168 completed.emplace_back(entry, symbian::native_error::kCancel);
169 it = entries_.erase(it);
170 continue;
171 }
172 if (!entry->cancel_submitted &&
173 entry->cancel_requested.load(std::memory_order_acquire) != 0) {
174 entry->timer.Cancel();
175 entry->cancel_submitted = true;
176 }
177 if (!entry->timer.IsReady()) {
178 ++it;
179 continue;
180 }
181 int code = entry->timer.Result();
182 if (code == 0) {
183 const auto elapsed =
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) {
190 } else {
191 const int rearmed = ArmNext(*entry);
192 if (rearmed == 0) {
193 ++it;
194 continue;
195 }
196 code = rearmed;
197 }
198 }
199 }
200 entry->timer.Close();
201 completed.emplace_back(entry, code);
202 it = entries_.erase(it);
203 }
204 for (auto& [entry, code] : completed) {
205 if (code == 0) {
206 entry->promise.SetValue(Unit{});
207 } else if (code == -3) { // KErrCancel.
208 entry->promise.SetError(absl::CancelledError("Timer cancelled"));
209 } else {
210 entry->promise.SetError(
211 symbian::StatusFromNativeError(code, "RTimer completion"));
212 }
213 }
214 return completed.size();
215 }
216
217 // The caller is the sole consumer of this OS thread's request semaphore.
218 // Check Window Server statuses and call DispatchReady before parking.
219 void Park() const {
220 if (!closed_ && wake_) {
222 }
223 }
224
225 void Close() {
226 if (closed_) {
227 return;
228 }
229 closed_ = true;
230 if (wake_) {
231 wake_->Close();
232 }
233 auto entries = std::move(entries_);
234 for (auto& entry : entries) {
235 entry->timer.Close();
236 entry->promise.SetError(absl::CancelledError("Timer pump closed"));
237 }
238 wake_.reset();
239 }
240
241 private:
242 struct WakeTarget {
243 int Open() { return SymbianRuntimeWakeCreate(&native_); }
244
245 void Notify() {
246 thread::MutexLock lock(&mu_);
247 if (native_ && !pending_) {
248 pending_ = true;
250 }
251 }
252
253 void Clear() {
254 thread::MutexLock lock(&mu_);
255 pending_ = false;
256 }
257
258 void Close() {
259 thread::MutexLock lock(&mu_);
260 if (native_) {
262 native_ = nullptr;
263 }
264 }
265
266 thread::Mutex mu_;
267 SymbianRuntimeWakeState* native_ = nullptr;
268 bool pending_ = false;
269 };
270
271 struct Entry {
272 NativeTimer timer;
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;
278 };
279
280 static int ArmNext(Entry& entry) {
281 entry.armed_at = std::chrono::steady_clock::now();
282 return entry.timer.Start(internal::TimerSliceMicroseconds(entry.remaining));
283 }
284
285 std::shared_ptr<WakeTarget> wake_;
286 std::vector<std::shared_ptr<Entry>> entries_;
287 std::size_t max_pending_;
288 WallClock wall_now_;
289 bool closed_ = false;
290};
291
292} // namespace symbian::concurrency
293
294#endif // SYMBIAN_CONCURRENCY_TIMER_PUMP_H_
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
Definition future.h:33