File
Blob: src/workerd/io/actor-sqlite.h
| 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 | #pragma once |
| 6 | |
| 7 | #include "actor-cache.h" |
| 8 | |
| 9 | #include <workerd/io/trace.h> |
| 10 | #include <workerd/util/sqlite-kv.h> |
| 11 | #include <workerd/util/sqlite-metadata.h> |
| 12 | |
| 13 | namespace workerd { |
| 14 | |
| 15 | // An implementation of ActorCacheOps that is backed by SqliteKv. |
| 16 | class ActorSqlite final: public ActorCacheInterface, private kj::TaskSet::ErrorHandler { |
| 17 | // TODO(perf): This interface is not designed ideally for wrapping SqliteKv. In particular, we |
| 18 | // end up allocating extra copies of all the results. It would be nicer if we could actually |
| 19 | // parse the V8-serialized values directly from the blob pointers that SQLite spits out. |
| 20 | // However, that probably requires rewriting `DurableObjectStorageOperations`. For now, hooking |
| 21 | // here is easier and not too costly. |
| 22 | |
| 23 | public: |
| 24 | // Hooks to configure ActorSqlite behavior, right now only used to allow plugging in a backend |
| 25 | // for alarm operations. |
| 26 | class Hooks { |
| 27 | public: |
| 28 | // Makes a request to the alarm manager to run the alarm handler at the given time, returning |
| 29 | // a promise that resolves when the scheduling has succeeded. `priorTask` is any work we must |
| 30 | // wait on prior to scheduling the new request, as of this writing, this would be the |
| 31 | // alarmLaterInFlight promise, which tracks any in-flight request to move the alarm "later" |
| 32 | // than is currently set. |
| 33 | virtual kj::Promise<void> scheduleRun( |
| 34 | kj::Maybe<kj::Date> newAlarmTime, kj::Promise<void> priorTask); |
| 35 | |
| 36 | static const Hooks DEFAULT; |
| 37 | |
| 38 | static constexpr inline Hooks& getDefaultHooks() { |
| 39 | // Hooks has no member variables, so const_cast is acceptable. |
| 40 | return const_cast<Hooks&>(Hooks::DEFAULT); |
| 41 | } |
| 42 | }; |
| 43 | |
| 44 | // Constructs ActorSqlite, arranging to honor the output gate, that is, any writes to the |
| 45 | // database which occur without any `await`s in between will automatically be combined into a |
| 46 | // single atomic write. This is accomplished using transactions. In addition to ensuring |
| 47 | // atomicity, this tends to improve performance, as SQLite is able to coalesce writes across |
| 48 | // statements that modify the same page. |
| 49 | // |
| 50 | // `commitCallback` will be invoked after committing a transaction. The output gate will block on |
| 51 | // the returned promise. This can be used e.g. when the database needs to be replicated to other |
| 52 | // machines before being considered durable. |
| 53 | explicit ActorSqlite(kj::Own<SqliteDatabase> dbParam, |
| 54 | OutputGate& outputGate, |
| 55 | kj::Function<kj::Promise<void>(SpanParent)> commitCallback, |
| 56 | Hooks& hooks = Hooks::getDefaultHooks(), |
| 57 | bool debugAlarmSync = false); |
| 58 | |
| 59 | bool isCommitScheduled() { |
| 60 | return !currentTxn.is<NoTxn>() || deleteAllCommitScheduled; |
| 61 | } |
| 62 | |
| 63 | kj::Maybe<SqliteDatabase&> getSqliteDatabase() override { |
| 64 | return *db; |
| 65 | } |
| 66 | |
| 67 | kj::Maybe<SqliteKv&> getSqliteKv() override { |
| 68 | requireNotBroken(); |
| 69 | return kv; |
| 70 | } |
| 71 | |
| 72 | kj::OneOf<kj::Maybe<Value>, kj::Promise<kj::Maybe<Value>>> get( |
| 73 | Key key, ReadOptions options) override; |
| 74 | kj::OneOf<GetResultList, kj::Promise<GetResultList>> get( |
| 75 | kj::Array<Key> keys, ReadOptions options) override; |
| 76 | kj::OneOf<kj::Maybe<kj::Date>, kj::Promise<kj::Maybe<kj::Date>>> getAlarm( |
| 77 | ReadOptions options) override; |
| 78 | kj::OneOf<GetResultList, kj::Promise<GetResultList>> list( |
| 79 | Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) override; |
| 80 | kj::OneOf<GetResultList, kj::Promise<GetResultList>> listReverse( |
| 81 | Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) override; |
| 82 | kj::Maybe<kj::Promise<void>> put( |
| 83 | Key key, Value value, WriteOptions options, SpanParent traceSpan) override; |
| 84 | kj::Maybe<kj::Promise<void>> put( |
| 85 | kj::Array<KeyValuePair> pairs, WriteOptions options, SpanParent traceSpan) override; |
| 86 | kj::OneOf<bool, kj::Promise<bool>> delete_( |
| 87 | Key key, WriteOptions options, SpanParent traceSpan) override; |
| 88 | kj::OneOf<uint, kj::Promise<uint>> delete_( |
| 89 | kj::Array<Key> keys, WriteOptions options, SpanParent traceSpan) override; |
| 90 | kj::Maybe<kj::Promise<void>> setAlarm( |
| 91 | kj::Maybe<kj::Date> newAlarmTime, WriteOptions options, SpanParent traceSpan) override; |
| 92 | // See ActorCacheOps. |
| 93 | |
| 94 | kj::Own<ActorCacheInterface::Transaction> startTransaction() override; |
| 95 | DeleteAllResults deleteAll( |
| 96 | WriteOptions options, SpanParent traceSpan, DeleteAllOptions deleteAllOptions = {}) override; |
| 97 | kj::Maybe<kj::Promise<void>> evictStale(kj::Date now) override; |
| 98 | void shutdown(kj::Maybe<const kj::Exception&> maybeException) override; |
| 99 | kj::OneOf<CancelAlarmHandler, RunAlarmHandler> armAlarmHandler(kj::Date scheduledTime, |
| 100 | SpanParent parentSpan, |
| 101 | kj::Date currentTime, |
| 102 | bool noCache = false, |
| 103 | kj::StringPtr actorId = "") override; |
| 104 | void cancelDeferredAlarmDeletion() override; |
| 105 | kj::Promise<kj::Maybe<kj::Date>> abandonAlarm(kj::Date scheduledTime) override; |
| 106 | kj::Maybe<kj::Promise<void>> onNoPendingFlush(SpanParent parentSpan) override; |
| 107 | kj::Promise<kj::String> getCurrentBookmark(SpanParent parentSpan) override; |
| 108 | kj::Promise<void> waitForBookmark(kj::StringPtr bookmark, SpanParent parentSpan) override; |
| 109 | // See ActorCacheInterface |
| 110 | |
| 111 | private: |
| 112 | kj::Own<SqliteDatabase> db; |
| 113 | OutputGate& outputGate; |
| 114 | kj::Function<kj::Promise<void>(SpanParent)> commitCallback; |
| 115 | Hooks& hooks; |
| 116 | SqliteKv kv; |
| 117 | SqliteMetadata metadata; |
| 118 | |
| 119 | // Define a SqliteDatabase::Regulator that is similar to TRUSTED but turns certain SQLite errors |
| 120 | // into application errors as appropriate when committing an implicit transaction. |
| 121 | class TxnCommitRegulator: public SqliteDatabase::Regulator { |
| 122 | public: |
| 123 | void onError(kj::Maybe<int> sqliteErrorCode, kj::StringPtr message) const override; |
| 124 | }; |
| 125 | static constexpr TxnCommitRegulator TRUSTED_TXN_COMMIT; |
| 126 | |
| 127 | SqliteDatabase::Statement beginTxn = db->prepare("BEGIN TRANSACTION"); |
| 128 | SqliteDatabase::Statement commitTxn = db->prepare(TRUSTED_TXN_COMMIT, "COMMIT TRANSACTION"); |
| 129 | |
| 130 | kj::Maybe<kj::Exception> broken; |
| 131 | |
| 132 | struct NoTxn {}; |
| 133 | |
| 134 | class ImplicitTxn { |
| 135 | public: |
| 136 | explicit ImplicitTxn(ActorSqlite& parent); |
| 137 | ~ImplicitTxn() noexcept(false); |
| 138 | KJ_DISALLOW_COPY_AND_MOVE(ImplicitTxn); |
| 139 | |
| 140 | void commit(); |
| 141 | void rollback(); |
| 142 | |
| 143 | void setSomeWriteConfirmed(bool someWriteConfirmed); |
| 144 | bool isSomeWriteConfirmed() const; |
| 145 | |
| 146 | private: |
| 147 | ActorSqlite& parent; |
| 148 | |
| 149 | bool committed = false; |
| 150 | |
| 151 | // True if any of the writes in this commit are confirmed writes. |
| 152 | bool someWriteConfirmed = false; |
| 153 | }; |
| 154 | |
| 155 | class ExplicitTxn: public ActorCacheInterface::Transaction, public kj::Refcounted { |
| 156 | public: |
| 157 | ExplicitTxn(ActorSqlite& actorSqlite); |
| 158 | ~ExplicitTxn() noexcept(false); |
| 159 | KJ_DISALLOW_COPY_AND_MOVE(ExplicitTxn); |
| 160 | |
| 161 | bool getAlarmDirty(); |
| 162 | void setAlarmDirty(); |
| 163 | |
| 164 | void setSomeWriteConfirmed(bool someWriteConfirmed); |
| 165 | bool isSomeWriteConfirmed() const; |
| 166 | |
| 167 | kj::Maybe<kj::Promise<void>> commit() override; |
| 168 | kj::Promise<void> rollback() override; |
| 169 | // Implements ActorCacheInterface::Transaction. |
| 170 | |
| 171 | kj::OneOf<kj::Maybe<Value>, kj::Promise<kj::Maybe<Value>>> get( |
| 172 | Key key, ReadOptions options) override; |
| 173 | kj::OneOf<GetResultList, kj::Promise<GetResultList>> get( |
| 174 | kj::Array<Key> keys, ReadOptions options) override; |
| 175 | kj::OneOf<kj::Maybe<kj::Date>, kj::Promise<kj::Maybe<kj::Date>>> getAlarm( |
| 176 | ReadOptions options) override; |
| 177 | kj::OneOf<GetResultList, kj::Promise<GetResultList>> list( |
| 178 | Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) override; |
| 179 | kj::OneOf<GetResultList, kj::Promise<GetResultList>> listReverse( |
| 180 | Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) override; |
| 181 | kj::Maybe<kj::Promise<void>> put( |
| 182 | Key key, Value value, WriteOptions options, SpanParent traceSpan) override; |
| 183 | kj::Maybe<kj::Promise<void>> put( |
| 184 | kj::Array<KeyValuePair> pairs, WriteOptions options, SpanParent traceSpan) override; |
| 185 | kj::OneOf<bool, kj::Promise<bool>> delete_( |
| 186 | Key key, WriteOptions options, SpanParent traceSpan) override; |
| 187 | kj::OneOf<uint, kj::Promise<uint>> delete_( |
| 188 | kj::Array<Key> keys, WriteOptions options, SpanParent traceSpan) override; |
| 189 | kj::Maybe<kj::Promise<void>> setAlarm( |
| 190 | kj::Maybe<kj::Date> newAlarmTime, WriteOptions options, SpanParent traceSpan) override; |
| 191 | // Implements ActorCacheOps. These will all forward to the ActorSqlite instance. |
| 192 | |
| 193 | private: |
| 194 | ActorSqlite& actorSqlite; |
| 195 | kj::Maybe<kj::Own<ExplicitTxn>> parent; |
| 196 | uint depth = 0; |
| 197 | bool hasChild = false; |
| 198 | bool committed = false; |
| 199 | bool alarmDirty = false; |
| 200 | // True if any of the writes in this commit are confirmed writes. |
| 201 | bool someWriteConfirmed = false; |
| 202 | |
| 203 | void rollbackImpl(); |
| 204 | }; |
| 205 | |
| 206 | // When set to NoTxn, there is no transaction outstanding. |
| 207 | // |
| 208 | // When set to `ImplicitTxn*`, an implicit transaction is currently open, owned by `commitTasks`. |
| 209 | // If there is a need to commit this early, e.g. to start an explicit transaction, that can be |
| 210 | // done through this reference. |
| 211 | // |
| 212 | // When set to `ExplicitTxn*`, an explicit transaction is currently open, so no implicit |
| 213 | // transactions should be used in the meantime. |
| 214 | kj::OneOf<NoTxn, ImplicitTxn*, ExplicitTxn*> currentTxn = NoTxn(); |
| 215 | |
| 216 | // If true, then a commit is scheduled as a result of deleteAll() having been called. |
| 217 | bool deleteAllCommitScheduled = false; |
| 218 | |
| 219 | // State for tracking completion of all commits (both confirmed and unconfirmed) for implementing |
| 220 | // sync() in onNoPendingFlush. |
| 221 | kj::ForkedPromise<void> lastCommit = kj::Promise<void>(kj::READY_NOW).fork(); |
| 222 | |
| 223 | // Backs the `kj::Own<void>` returned by `armAlarmHandler()`. |
| 224 | class DeferredAlarmDeleter: public kj::Disposer { |
| 225 | public: |
| 226 | // The `Own<void>` returned by `armAlarmHandler()` is actually set up to point to the |
| 227 | // `ActorSqlite` itself, but with an alternate disposer that deletes the alarm rather than |
| 228 | // the whole object. |
| 229 | void disposeImpl(void* pointer) const override { |
| 230 | reinterpret_cast<ActorSqlite*>(pointer)->maybeDeleteDeferredAlarm(); |
| 231 | } |
| 232 | }; |
| 233 | |
| 234 | // We need to track some additional alarm state to guarantee at-least-once alarm delivery: |
| 235 | // Within an alarm handler, we want the observable alarm state to look like the running alarm |
| 236 | // was deleted at the start of the handler (when armAlarmHandler() is called), but we don't |
| 237 | // actually want to persist that deletion until after the handler has successfully completed. |
| 238 | bool haveDeferredDelete = false; |
| 239 | |
| 240 | // Trace span for the deferred alarm deletion, captured from armAlarmHandler and used when |
| 241 | // the alarm is actually deleted. This is separate from currentCommitSpan because the alarm |
| 242 | // deletion is an internal write (via metadata.setAlarm) that doesn't go through the regular |
| 243 | // write methods with a traceSpan parameter. If the alarm handler does no other writes, |
| 244 | // currentCommitSpan would be null, so we need this saved span for the output gate lock trace. |
| 245 | SpanParent deferredAlarmSpan = nullptr; |
| 246 | |
| 247 | // Some state only used for tracking calling invariants. |
| 248 | bool inAlarmHandler = false; |
| 249 | |
| 250 | // The alarm state for which we last received confirmation that the db was durably stored. |
| 251 | kj::Maybe<kj::Date> lastConfirmedAlarmDbState; |
| 252 | |
| 253 | // The latest time we'd expect a scheduled alarm to fire, given the current set of in-flight |
| 254 | // scheduling requests, without yet knowing if any of them succeeded or failed. We use this |
| 255 | // value to maintain the invariant that the scheduled alarm is always equal to or earlier than |
| 256 | // the alarm value in the persisted database state. |
| 257 | kj::Maybe<kj::Date> alarmScheduledNoLaterThan; |
| 258 | |
| 259 | // A promise for an in-progress alarm notification update and database commit. |
| 260 | kj::Maybe<kj::ForkedPromise<void>> pendingCommit; |
| 261 | |
| 262 | kj::TaskSet commitTasks; |
| 263 | |
| 264 | // Trace span for the current commit operation. Captured from each write and used |
| 265 | // for the output gate lock hold trace when a non-allowUnconfirmed write occurs. |
| 266 | SpanParent currentCommitSpan = nullptr; |
| 267 | |
| 268 | // Promise for the currently in-flight "move alarm later" operation, if any. |
| 269 | // Used to serialize move-earlier operations against any pending move-later operation. |
| 270 | kj::ForkedPromise<void> alarmLaterInFlight = kj::Promise<void>(kj::READY_NOW).fork(); |
| 271 | |
| 272 | // True when a "move alarm later" request is currently in-flight via scheduleLaterAlarm(). |
| 273 | bool alarmLaterIsInFlight = false; |
| 274 | |
| 275 | // When a "move alarm later" request is already in-flight and we need to schedule |
| 276 | // another one, we store the desired alarm time here. When the in-flight request |
| 277 | // completes, it checks this variable and starts a new request if needed. The outer |
| 278 | // Maybe indicates whether there is a pending time at all; the inner Maybe<Date> is |
| 279 | // the alarm time to set (where kj::none means "clear the alarm"). |
| 280 | kj::Maybe<kj::Maybe<kj::Date>> pendingLaterAlarmTime; |
| 281 | |
| 282 | // Version counter that increments on every alarm change. Used to detect if another commit |
| 283 | // modified the alarm while we were async, allowing us to skip redundant post-commit alarm |
| 284 | // syncs. This provides automatic coalescing of rapid alarm changes. |
| 285 | uint64_t alarmVersion = 0; |
| 286 | |
| 287 | // Debug flag for tracing alarm synchronization issues for specific namespaces |
| 288 | bool debugAlarmSync = false; |
| 289 | |
| 290 | void startImplicitTxn(); |
| 291 | |
| 292 | void onWrite(bool allowUnconfirmed); |
| 293 | |
| 294 | void onCriticalError(kj::StringPtr errorMessage, kj::Maybe<kj::Exception> maybeException); |
| 295 | |
| 296 | // Issues a request to the alarm scheduler for the given time, returning a promise that resolves |
| 297 | // when the request is confirmed. |
| 298 | kj::Promise<void> requestScheduledAlarm( |
| 299 | kj::Maybe<kj::Date> requestedTime, kj::Promise<void> priorTask); |
| 300 | |
| 301 | // Schedules a "move alarm later" operation. If no move-later is currently in-flight, starts one |
| 302 | // immediately. If one is already in-flight, stores the desired time in `pendingLaterAlarmTime` |
| 303 | // so it will be picked up when the current in-flight operation completes. |
| 304 | void scheduleLaterAlarm(kj::Maybe<kj::Date> newAlarmTime, SpanParent parentSpan); |
| 305 | |
| 306 | struct PrecommitAlarmState { |
| 307 | // Promise for the completion of precommit alarm scheduling |
| 308 | kj::Maybe<kj::Promise<void>> schedulingPromise; |
| 309 | }; |
| 310 | |
| 311 | // To be called just before committing the local sqlite db, to synchronously start any necessary |
| 312 | // alarm scheduling: |
| 313 | PrecommitAlarmState startPrecommitAlarmScheduling(); |
| 314 | |
| 315 | // Performs the rest of the asynchronous commit, to be waited on after committing the local |
| 316 | // sqlite db. Should be called in the same turn of the event loop as |
| 317 | // startPrecommitAlarmScheduling() and passed the state that it returned. |
| 318 | kj::Promise<void> commitImpl(PrecommitAlarmState precommitAlarmState, SpanParent parentSpan); |
| 319 | |
| 320 | void taskFailed(kj::Exception&& exception) override; |
| 321 | |
| 322 | void requireNotBroken(); |
| 323 | |
| 324 | // Called when DeferredAlarmDeleter is destroyed, to delete alarm if not reset or cancelled |
| 325 | // during handler. |
| 326 | void maybeDeleteDeferredAlarm(); |
| 327 | }; |
| 328 | |
| 329 | } // namespace workerd |