// Copyright (c) 2017-2022 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #pragma once #include #include #include #include #include #include #include #include namespace workerd::server { struct ActorKey { kj::StringPtr actorId; bool operator==(const ActorKey& other) const { return actorId == other.actorId; } kj::Own clone() const { auto ownActorId = kj::str(actorId); return kj::attachVal(ActorKey{.actorId = ownActorId}, kj::mv(ownActorId)); } }; inline uint KJ_HASHCODE(const ActorKey& k) { return kj::hashCode(k.actorId); } // Allows scheduling alarm executions at specific times, returning a promise representing // the completion of the alarm event. class AlarmScheduler final: kj::TaskSet::ErrorHandler { public: static constexpr auto RETRY_START_SECONDS = WorkerInterface::ALARM_RETRY_START_SECONDS; // Max number of "valid" retry attempts, i.e the worker returned an error static constexpr auto RETRY_MAX_TRIES = WorkerInterface::ALARM_RETRY_MAX_TRIES; // Bound for exponential backoff when RETRY_MAX_TRIES is exceeded due to internal errors. // 2 << 9 is 1024 seconds, about 17 minutes. Total time spent in retries once the backoff limit // is reached is over 30 minutes. static constexpr auto RETRY_BACKOFF_MAX = 9; // How much jitter should be applied to retry times to avoid bundled retries overloading // some common dependency between a set of failed alarms static constexpr auto RETRY_JITTER_FACTOR = 0.25; using GetActorFn = kj::Function(kj::String)>; AlarmScheduler(const kj::Clock& clock, kj::Timer& timer, const SqliteDatabase::Vfs& vfs, kj::Path path, GetActorFn getActor); kj::Maybe getAlarm(ActorKey actor); bool setAlarm(ActorKey actor, kj::Date scheduledTime); bool deleteAlarm(ActorKey actor); // Cancels all pending alarms and removes them from persistent storage. void deleteAll(); private: enum class AlarmStatus { WAITING, STARTED, FINISHED }; const kj::Clock& clock; kj::Timer& timer; std::default_random_engine random; GetActorFn getActor; kj::Own db; kj::TaskSet tasks; struct ScheduledAlarm { kj::Own actor; kj::Date scheduledTime; kj::Promise task; kj::Maybe queuedAlarm = kj::none; // Once started, an alarm can have a single alarm queued behind it. AlarmStatus status = AlarmStatus::WAITING; bool previousRetryCountedAgainstLimit = false; // Counter for calculating backoff -- separate from retry, so we can reset backoff without losing // the total count of retry attempts uint32_t backoff = 0; // Counter for retry attempts, whether or not they apply to the limit uint32_t retry = 0; // Counter for retry attempts that apply to the retry limit. uint32_t countedRetry = 0; }; kj::HashMap alarms; struct RetryInfo { bool retry; bool retryCountsAgainstLimit; }; kj::Promise runAlarm( const ActorKey& actor, kj::Date scheduledTime, uint32_t retryCount); ScheduledAlarm scheduleAlarm(kj::Date now, kj::Own actor, kj::Date scheduledTime); kj::Promise makeAlarmTask( kj::Duration delay, const ActorKey& actor, kj::Date scheduledTime); kj::Promise checkTimestamp(kj::Duration delay, kj::Date scheduledTime); SqliteDatabase::Statement stmtSetAlarm = db->prepare(R"( INSERT INTO _cf_ALARM VALUES(?, ?) ON CONFLICT DO UPDATE SET scheduled_time = excluded.scheduled_time; )"); SqliteDatabase::Statement stmtDeleteAlarm = db->prepare(R"( DELETE FROM _cf_ALARM WHERE actor_id = ? )"); void taskFailed(kj::Exception&& exception) override; int maxJitterMsForDelay(kj::Duration delay); static void ensureInitialized(SqliteDatabase& db); void loadAlarmsFromDb(); }; } // namespace workerd::server