File
Blob: src/workerd/server/alarm-scheduler.h
| 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 | |
| 18 | namespace workerd::server { |
| 19 | |
| 20 | struct 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 | |
| 34 | inline 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. |
| 40 | class 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 |