Skip to content
File

Blob: src/workerd/io/actor-sqlite.h

cpp330 lines
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 
13namespace workerd {
14 
15// An implementation of ActorCacheOps that is backed by SqliteKv.
16class 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