Skip to content
File

Blob: src/workerd/server/alarm-scheduler.c++

10.4 KB
1// Copyright (c) 2023 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#include "alarm-scheduler.h"
6 
7#include <kj/debug.h>
8 
9#include <cmath>
10 
11namespace workerd::server {
12 
13int AlarmScheduler::maxJitterMsForDelay(kj::Duration delay) {
14 double delayMs = delay / kj::MILLISECONDS;
15 return std::floor(RETRY_JITTER_FACTOR * delayMs);
16}
17 
18namespace {
19 
20std::default_random_engine makeSeededRandomEngine() {
21 // Using the time as a seed here is fine, we just want to have some randomness for retry jitter
22 auto time = kj::systemPreciseMonotonicClock().now();
23 auto seed = (time - kj::origin<kj::TimePoint>()) / kj::NANOSECONDS;
24 
25 std::default_random_engine engine(seed);
26 return engine;
27}
28 
29} // namespace
30 
31AlarmScheduler::AlarmScheduler(const kj::Clock& clock,
32 kj::Timer& timer,
33 const SqliteDatabase::Vfs& vfs,
34 kj::Path path,
35 GetActorFn getActor)
36 : clock(clock),
37 timer(timer),
38 random(makeSeededRandomEngine()),
39 getActor(kj::mv(getActor)),
40 db([&] {
41 auto db = kj::heap<SqliteDatabase>(vfs, kj::mv(path),
42 kj::WriteMode::CREATE | kj::WriteMode::MODIFY | kj::WriteMode::CREATE_PARENT);
43 ensureInitialized(*db);
44 return kj::mv(db);
45 }()),
46 tasks(*this) {
47 loadAlarmsFromDb();
48}
49 
50void AlarmScheduler::ensureInitialized(SqliteDatabase& db) {
51 // TODO(sqlite): Do this automatically at a lower layer?
52 db.run("PRAGMA journal_mode=WAL;");
53 
54 db.run(R"(
55 CREATE TABLE IF NOT EXISTS _cf_ALARM (
56 actor_id TEXT PRIMARY KEY,
57 scheduled_time INTEGER
58 ) WITHOUT ROWID;
59 )");
60}
61 
62void AlarmScheduler::loadAlarmsFromDb() {
63 auto now = clock.now();
64 
65 // TODO(someday): don't maintain the entire alarm set in memory -- right now for the usecase of
66 // local development, doing so is sufficient.
67 auto query = db->run(R"(
68 SELECT actor_id, scheduled_time FROM _cf_ALARM;
69 )");
70 
71 while (!query.isDone()) {
72 auto date = kj::UNIX_EPOCH + (kj::NANOSECONDS * query.getInt64(1));
73 
74 auto ownActorId = kj::str(query.getText(0));
75 auto actor = kj::attachVal(ActorKey{.actorId = ownActorId}, kj::mv(ownActorId));
76 auto& actorRef = *actor;
77 
78 alarms.insert(actorRef, scheduleAlarm(now, kj::mv(actor), date));
79 
80 query.nextRow();
81 }
82}
83 
84kj::Maybe<kj::Date> AlarmScheduler::getAlarm(ActorKey actor) {
85 // TODO(someday): Might be able to simplify AlarmScheduler somewhat, now that ActorSqlite no
86 // longer relies on it for getAlarm()?
87 KJ_IF_SOME(alarm, alarms.find(actor)) {
88 if (alarm.status == AlarmStatus::STARTED) {
89 // getAlarm() when the alarm handler is running should return null,
90 // unless an alarm is queued;
91 return alarm.queuedAlarm;
92 } else {
93 return alarm.scheduledTime;
94 }
95 } else {
96 // We currently retain the entire set of queued alarms in memory, no need to hit sqlite
97 return kj::none;
98 }
99}
100 
101bool AlarmScheduler::setAlarm(ActorKey actor, kj::Date scheduledTime) {
102 int64_t scheduledTimeNs = (scheduledTime - kj::UNIX_EPOCH) / kj::NANOSECONDS;
103 auto query = stmtSetAlarm.run(actor.actorId, scheduledTimeNs);
104 
105 bool existing = true;
106 auto& entry = alarms.findOrCreate(actor, [&]() {
107 existing = false;
108 
109 auto ownActorId = kj::str(actor.actorId);
110 auto ownActor = kj::attachVal(ActorKey{.actorId = ownActorId}, kj::mv(ownActorId));
111 
112 return decltype(alarms)::Entry{
113 *ownActor, scheduleAlarm(clock.now(), kj::mv(ownActor), scheduledTime)};
114 });
115 
116 if (existing) {
117 if (entry.status != AlarmStatus::WAITING) {
118 // We queue any new alarm after the existing alarm even if the new alarm has the same scheduled
119 // time, as receiving a notification directly maps to a write for that time in the actor.
120 entry.queuedAlarm = scheduledTime;
121 } else {
122 entry = scheduleAlarm(clock.now(), kj::mv(entry.actor), scheduledTime);
123 }
124 }
125 
126 return query.changeCount() > 0;
127}
128 
129void AlarmScheduler::deleteAll() {
130 // Cancel all in-memory alarm tasks.
131 alarms.clear();
132 // Wipe the persistent store.
133 db->run("DELETE FROM _cf_ALARM;");
134}
135 
136bool AlarmScheduler::deleteAlarm(ActorKey actor) {
137 auto query = stmtDeleteAlarm.run(actor.actorId);
138 
139 KJ_IF_SOME(entry, alarms.findEntry(actor)) {
140 KJ_IF_SOME(queued, entry.value.queuedAlarm) {
141 if (entry.value.status == AlarmStatus::STARTED) {
142 // If we are currently running an alarm, we want to delete the queued instead of current.
143 entry.value.queuedAlarm = kj::none;
144 } else {
145 entry.value = scheduleAlarm(clock.now(), kj::mv(entry.value.actor), queued);
146 }
147 } else {
148 if (entry.value.status != AlarmStatus::STARTED) {
149 // We can't remove running alarms.
150 alarms.erase(entry);
151 }
152 }
153 }
154 
155 return query.changeCount() > 0;
156}
157 
158kj::Promise<AlarmScheduler::RetryInfo> AlarmScheduler::runAlarm(
159 const ActorKey& actor, kj::Date scheduledTime, uint32_t retryCount) {
160 auto result = co_await getActor(kj::str(actor.actorId))->runAlarm(scheduledTime, retryCount);
161 
162 co_return RetryInfo{.retry = result.outcome != EventOutcome::OK && result.retry,
163 .retryCountsAgainstLimit = result.retryCountsAgainstLimit};
164}
165 
166AlarmScheduler::ScheduledAlarm AlarmScheduler::scheduleAlarm(
167 kj::Date now, kj::Own<ActorKey> actor, kj::Date scheduledTime) {
168 auto task = makeAlarmTask(scheduledTime - now, *actor, scheduledTime);
169 
170 return ScheduledAlarm{kj::mv(actor), scheduledTime, kj::mv(task)};
171}
172 
173kj::Promise<void> AlarmScheduler::checkTimestamp(kj::Duration delay, kj::Date scheduledTime) {
174 co_await timer.afterDelay(delay);
175 
176 // Since we are waiting on timer.afterDelay, it's possible that timer.now() was behind
177 // the real time by a few ms, leading to premature alarm() execution. This checks it the current
178 // time is >= than scheduledTime to ensure we run alarms only on or after their scheduled time.
179 auto now = clock.now();
180 if (now < scheduledTime) {
181 // If it's not yet time to trigger the alarm, we shall wait a while longer until we can
182 // trigger it. This repeats until it's time for the alarm to run.
183 co_await checkTimestamp(scheduledTime - now, scheduledTime);
184 }
185}
186 
187kj::Promise<void> AlarmScheduler::makeAlarmTask(
188 kj::Duration delay, const ActorKey& actorRef, kj::Date scheduledTime) {
189 co_await checkTimestamp(delay, scheduledTime);
190 uint32_t retryCount = 0;
191 {
192 auto& entry = KJ_ASSERT_NONNULL(alarms.findEntry(actorRef));
193 entry.value.status = AlarmStatus::STARTED;
194 retryCount = entry.value.countedRetry;
195 }
196 
197 auto retryInfo = co_await ([&]() -> kj::Promise<RetryInfo> {
198 try {
199 co_return co_await runAlarm(actorRef, scheduledTime, retryCount);
200 } catch (...) {
201 auto exception = kj::getCaughtExceptionAsKj();
202 KJ_LOG(WARNING, exception);
203 co_return RetryInfo{.retry = true,
204 
205 // An exception here is "weird", they should normally
206 // be turned into AlarmResult statuses in the sandbox
207 // for any user-caused error. Let's not count this
208 // retry attempt against the limit.
209 .retryCountsAgainstLimit = false};
210 }
211 })();
212 
213 try {
214 auto& entry = KJ_ASSERT_NONNULL(alarms.findEntry(actorRef));
215 
216 // We can't overwrite our entry before moving ourselves out of it, as a promise cannot
217 // delete itself.
218 tasks.add(kj::mv(entry.value.task));
219 
220 // If an alarm is queued, there's no point in retrying the current one -- proceed
221 // to running the queued alarm instead.
222 KJ_IF_SOME(a, entry.value.queuedAlarm) {
223 // creating a new alarm and overwriting the old one will reset
224 // `status` to WAITING and `queuedAlarm` to null
225 entry.value = scheduleAlarm(clock.now(), kj::mv(entry.value.actor), a);
226 co_return;
227 }
228 
229 // When we reach this block of code and alarm has either succeeded or failed and may (or may
230 // not) retry. Setting the status of an alarm as FINISHED here, will allow deletion of alarms
231 // between retries. If there's a retry, `makeAlarmTask` is called, setting status as RUNNING
232 // again.
233 entry.value.status = AlarmStatus::FINISHED;
234 
235 if (retryInfo.retry) {
236 // recreate the task, running after a delay determined using the retry factor
237 if (entry.value.countedRetry >= AlarmScheduler::RETRY_MAX_TRIES) {
238 // Notify the actor to clear its in-memory alarm state so getAlarm() reflects the
239 // deletion. We ignore the returned remaining time — the workerd-local alarm scheduler
240 // already has visibility into the actor's alarm state via its SQLite hooks.
241 // If the notification fails, we keep the alarm in the scheduler so it is not silently
242 // lost.
243 try {
244 co_await getActor(kj::str(actorRef.actorId))->abandonAlarm(scheduledTime).ignoreResult();
245 } catch (...) {
246 auto exception = kj::getCaughtExceptionAsKj();
247 KJ_LOG(
248 WARNING, "abandonAlarm notification failed, keeping alarm in scheduler", exception);
249 co_return;
250 }
251 deleteAlarm(*entry.value.actor);
252 co_return;
253 }
254 if (retryInfo.retryCountsAgainstLimit) {
255 entry.value.countedRetry++;
256 
257 if (!entry.value.previousRetryCountedAgainstLimit) {
258 // The last retry didn't count against the limit, indicating it was due to some internal
259 // error. However, this retry does, meaning it's due to an error in user code,
260 // most likely a different error. We should reset the retry counter used for
261 // calculating backoff, so user-caused retries don't have an unnecessarily high backoff
262 // time if they come after internal-caused retries.
263 
264 entry.value.backoff = 0;
265 }
266 }
267 entry.value.previousRetryCountedAgainstLimit = retryInfo.retryCountsAgainstLimit;
268 
269 entry.value.backoff = kj::min(AlarmScheduler::RETRY_BACKOFF_MAX, entry.value.backoff);
270 auto delay = (AlarmScheduler::RETRY_START_SECONDS << entry.value.backoff) * kj::SECONDS;
271 
272 std::uniform_int_distribution<> distribution(0, maxJitterMsForDelay(delay));
273 delay += distribution(random) * kj::MILLISECONDS;
274 
275 entry.value.backoff++;
276 entry.value.retry++;
277 
278 entry.value.task = makeAlarmTask(delay, actorRef, scheduledTime);
279 } else {
280 KJ_ASSERT(entry.value.queuedAlarm == kj::none);
281 deleteAlarm(actorRef);
282 }
283 } catch (...) {
284 auto exception = kj::getCaughtExceptionAsKj();
285 KJ_LOG(ERROR, "Failed to run alarm and was unable to schedule a retry", exception);
286 }
287}
288 
289void AlarmScheduler::taskFailed(kj::Exception&& e) {
290 KJ_LOG(WARNING, e);
291}
292 
293} // namespace workerd::server