Skip to content
File

Blob: src/workerd/server/alarm-scheduler.h

cpp134 lines
1// Copyright (c) 2017-2022 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#pragma once
6 
7#include <workerd/io/worker-interface.h>
8#include <workerd/util/sqlite.h>
9 
10#include <kj/async.h>
11#include <kj/common.h>
12#include <kj/map.h>
13#include <kj/time.h>
14#include <kj/timer.h>
15 
16#include <random>
17 
18namespace workerd::server {
19 
20struct ActorKey {
21 kj::StringPtr actorId;
22 
23 bool operator==(const ActorKey& other) const {
24 return actorId == other.actorId;
25 }
26 
27 kj::Own<ActorKey> clone() const {
28 auto ownActorId = kj::str(actorId);
29 
30 return kj::attachVal(ActorKey{.actorId = ownActorId}, kj::mv(ownActorId));
31 }
32};
33 
34inline uint KJ_HASHCODE(const ActorKey& k) {
35 return kj::hashCode(k.actorId);
36}
37 
38// Allows scheduling alarm executions at specific times, returning a promise representing
39// the completion of the alarm event.
40class AlarmScheduler final: kj::TaskSet::ErrorHandler {
41 public:
42 static constexpr auto RETRY_START_SECONDS = WorkerInterface::ALARM_RETRY_START_SECONDS;
43 
44 // Max number of "valid" retry attempts, i.e the worker returned an error
45 static constexpr auto RETRY_MAX_TRIES = WorkerInterface::ALARM_RETRY_MAX_TRIES;
46 
47 // Bound for exponential backoff when RETRY_MAX_TRIES is exceeded due to internal errors.
48 // 2 << 9 is 1024 seconds, about 17 minutes. Total time spent in retries once the backoff limit
49 // is reached is over 30 minutes.
50 static constexpr auto RETRY_BACKOFF_MAX = 9;
51 
52 // How much jitter should be applied to retry times to avoid bundled retries overloading
53 // some common dependency between a set of failed alarms
54 static constexpr auto RETRY_JITTER_FACTOR = 0.25;
55 
56 using GetActorFn = kj::Function<kj::Own<WorkerInterface>(kj::String)>;
57 
58 AlarmScheduler(const kj::Clock& clock,
59 kj::Timer& timer,
60 const SqliteDatabase::Vfs& vfs,
61 kj::Path path,
62 GetActorFn getActor);
63 
64 kj::Maybe<kj::Date> getAlarm(ActorKey actor);
65 bool setAlarm(ActorKey actor, kj::Date scheduledTime);
66 bool deleteAlarm(ActorKey actor);
67 
68 // Cancels all pending alarms and removes them from persistent storage.
69 void deleteAll();
70 
71 private:
72 enum class AlarmStatus { WAITING, STARTED, FINISHED };
73 const kj::Clock& clock;
74 kj::Timer& timer;
75 std::default_random_engine random;
76 GetActorFn getActor;
77 kj::Own<SqliteDatabase> db;
78 kj::TaskSet tasks;
79 
80 struct ScheduledAlarm {
81 kj::Own<ActorKey> actor;
82 kj::Date scheduledTime;
83 kj::Promise<void> task;
84 kj::Maybe<kj::Date> queuedAlarm = kj::none;
85 // Once started, an alarm can have a single alarm queued behind it.
86 AlarmStatus status = AlarmStatus::WAITING;
87 
88 bool previousRetryCountedAgainstLimit = false;
89 
90 // Counter for calculating backoff -- separate from retry, so we can reset backoff without losing
91 // the total count of retry attempts
92 uint32_t backoff = 0;
93 
94 // Counter for retry attempts, whether or not they apply to the limit
95 uint32_t retry = 0;
96 
97 // Counter for retry attempts that apply to the retry limit.
98 uint32_t countedRetry = 0;
99 };
100 
101 kj::HashMap<ActorKey, ScheduledAlarm> alarms;
102 
103 struct RetryInfo {
104 bool retry;
105 bool retryCountsAgainstLimit;
106 };
107 kj::Promise<RetryInfo> runAlarm(
108 const ActorKey& actor, kj::Date scheduledTime, uint32_t retryCount);
109 
110 ScheduledAlarm scheduleAlarm(kj::Date now, kj::Own<ActorKey> actor, kj::Date scheduledTime);
111 
112 kj::Promise<void> makeAlarmTask(
113 kj::Duration delay, const ActorKey& actor, kj::Date scheduledTime);
114 
115 kj::Promise<void> checkTimestamp(kj::Duration delay, kj::Date scheduledTime);
116 
117 SqliteDatabase::Statement stmtSetAlarm = db->prepare(R"(
118 INSERT INTO _cf_ALARM VALUES(?, ?)
119 ON CONFLICT DO UPDATE SET scheduled_time = excluded.scheduled_time;
120 )");
121 SqliteDatabase::Statement stmtDeleteAlarm = db->prepare(R"(
122 DELETE FROM _cf_ALARM WHERE actor_id = ?
123 )");
124 
125 void taskFailed(kj::Exception&& exception) override;
126 
127 int maxJitterMsForDelay(kj::Duration delay);
128 
129 static void ensureInitialized(SqliteDatabase& db);
130 void loadAlarmsFromDb();
131};
132 
133} // namespace workerd::server