File
Blob: src/workerd/io/actor-sqlite.c++
| 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 | #include "actor-sqlite.h" |
| 6 | |
| 7 | #include "io-gate.h" |
| 8 | |
| 9 | #include <workerd/jsg/exception.h> |
| 10 | #include <workerd/util/sentry.h> |
| 11 | |
| 12 | #include <kj/exception.h> |
| 13 | #include <kj/function.h> |
| 14 | |
| 15 | #include <algorithm> |
| 16 | |
| 17 | namespace workerd { |
| 18 | |
| 19 | namespace { |
| 20 | |
| 21 | // Returns true if a given (set or unset) alarm will fire earlier than another. |
| 22 | static bool willFireEarlier(kj::Maybe<kj::Date> alarm1, kj::Maybe<kj::Date> alarm2) { |
| 23 | // Intuitively, an unset alarm is effectively indistinguishable from an alarm set at infinity. |
| 24 | return alarm1.orDefault(kj::maxValue) < alarm2.orDefault(kj::maxValue); |
| 25 | } |
| 26 | |
| 27 | // Helper to make kj::Maybe<kj::Date> loggable - returns the date or kj::maxValue for logging |
| 28 | static kj::Date logDate(kj::Maybe<kj::Date> maybeDate) { |
| 29 | return maybeDate.orDefault(kj::maxValue); |
| 30 | } |
| 31 | |
| 32 | // Set options.allowUnconfirmed to false and log a reason why. |
| 33 | void disableAllowUnconfirmed(ActorCacheOps::WriteOptions& options, kj::StringPtr reason) { |
| 34 | if (options.allowUnconfirmed) { |
| 35 | KJ_LOG(WARNING, "NOSENTRY allowUnconfirmed disabled", reason); |
| 36 | options.allowUnconfirmed = false; |
| 37 | } |
| 38 | } |
| 39 | |
| 40 | } // namespace |
| 41 | |
| 42 | ActorSqlite::ActorSqlite(kj::Own<SqliteDatabase> dbParam, |
| 43 | OutputGate& outputGate, |
| 44 | kj::Function<kj::Promise<void>(SpanParent)> commitCallback, |
| 45 | Hooks& hooks, |
| 46 | bool debugAlarmSyncParam) |
| 47 | : db(kj::mv(dbParam)), |
| 48 | outputGate(outputGate), |
| 49 | commitCallback(kj::mv(commitCallback)), |
| 50 | hooks(hooks), |
| 51 | kv(*db), |
| 52 | metadata(*db), |
| 53 | commitTasks(*this), |
| 54 | debugAlarmSync(debugAlarmSyncParam) { |
| 55 | db->onWrite(KJ_BIND_METHOD(*this, onWrite)); |
| 56 | db->onCriticalError(KJ_BIND_METHOD(*this, onCriticalError)); |
| 57 | lastConfirmedAlarmDbState = metadata.getAlarm(); |
| 58 | |
| 59 | // Because we preserve an invariant that scheduled alarms are always at or earlier than |
| 60 | // persisted db alarm state, it should be OK to populate our idea of the latest scheduled alarm |
| 61 | // using the current db alarm state. At worst, it may perform one unnecessary scheduling |
| 62 | // request in cases where a previous alarm-state-altering transaction failed. |
| 63 | alarmScheduledNoLaterThan = metadata.getAlarm(); |
| 64 | } |
| 65 | |
| 66 | ActorSqlite::ImplicitTxn::ImplicitTxn(ActorSqlite& parent): parent(parent) { |
| 67 | KJ_REQUIRE(parent.currentTxn.is<NoTxn>()); |
| 68 | parent.beginTxn.run(); |
| 69 | parent.currentTxn = this; |
| 70 | } |
| 71 | ActorSqlite::ImplicitTxn::~ImplicitTxn() noexcept(false) { |
| 72 | KJ_IF_SOME(c, parent.currentTxn.tryGet<ImplicitTxn*>()) { |
| 73 | if (c == this) { |
| 74 | parent.currentTxn.init<NoTxn>(); |
| 75 | } |
| 76 | } |
| 77 | if (!committed && parent.broken == kj::none) { |
| 78 | // Failed to commit, so roll back. |
| 79 | // |
| 80 | // This should only happen in cases of catastrophic error. Since this is rarely actually |
| 81 | // executed, we don't prepare a statement for it. |
| 82 | parent.db->run("ROLLBACK TRANSACTION"); |
| 83 | } |
| 84 | } |
| 85 | |
| 86 | void ActorSqlite::ImplicitTxn::commit() { |
| 87 | // Ignore redundant commit()s. |
| 88 | if (!committed) { |
| 89 | parent.commitTxn.run(); |
| 90 | committed = true; |
| 91 | } |
| 92 | } |
| 93 | |
| 94 | void ActorSqlite::ImplicitTxn::rollback() { |
| 95 | // Ignore redundant commit()s. |
| 96 | if (!committed) { |
| 97 | // As of this writing, rollback() is only called when the database is about to be reset. |
| 98 | // Preparing a statement for it would be a waste since that statement would never be executed |
| 99 | // more than once, since resetting requires repreparing all statements anyway. So we don't |
| 100 | // bother. |
| 101 | parent.db->run("ROLLBACK TRANSACTION"); |
| 102 | committed = true; |
| 103 | } |
| 104 | } |
| 105 | |
| 106 | void ActorSqlite::ImplicitTxn::setSomeWriteConfirmed(bool someWriteConfirmed) { |
| 107 | this->someWriteConfirmed = someWriteConfirmed; |
| 108 | } |
| 109 | |
| 110 | bool ActorSqlite::ImplicitTxn::isSomeWriteConfirmed() const { |
| 111 | return someWriteConfirmed; |
| 112 | } |
| 113 | |
| 114 | ActorSqlite::ExplicitTxn::ExplicitTxn(ActorSqlite& actorSqlite): actorSqlite(actorSqlite) { |
| 115 | KJ_SWITCH_ONEOF(actorSqlite.currentTxn) { |
| 116 | KJ_CASE_ONEOF(_, NoTxn) {} |
| 117 | KJ_CASE_ONEOF(implicit, ImplicitTxn*) { |
| 118 | // An implicit transaction is open, commit it now because it would be weird if writes |
| 119 | // performed before the explicit transaction started were postponed until the transaction |
| 120 | // completes. Note that this isn't violating any atomicity guarantees because the transaction |
| 121 | // API is async, and atomicity is only guaranteed over synchronous code. |
| 122 | implicit->commit(); |
| 123 | } |
| 124 | KJ_CASE_ONEOF(exp, ExplicitTxn*) { |
| 125 | KJ_REQUIRE(!exp->hasChild, |
| 126 | "critical section should have blocked creation of more than one child at a time"); |
| 127 | parent = kj::addRef(*exp); |
| 128 | exp->hasChild = true; |
| 129 | depth = exp->depth + 1; |
| 130 | alarmDirty = exp->alarmDirty; |
| 131 | someWriteConfirmed = exp->someWriteConfirmed; |
| 132 | } |
| 133 | } |
| 134 | actorSqlite.currentTxn = this; |
| 135 | |
| 136 | // To support nested transactions, we assign each savepoint a name based on its nesting depth. |
| 137 | // Unfortunately this means we cannot prepare the statement, unless we prepare a series of |
| 138 | // statements for each depth. (Actually, it could be reasonable to prepare statements for |
| 139 | // depth 0 specifically, but I'm not going to try it for now.) |
| 140 | actorSqlite.db->run( |
| 141 | {.regulator = SqliteDatabase::TRUSTED}, kj::str("SAVEPOINT _cf_savepoint_", depth)); |
| 142 | } |
| 143 | ActorSqlite::ExplicitTxn::~ExplicitTxn() noexcept(false) { |
| 144 | KJ_DEFER([&]() noexcept { |
| 145 | // We'd better crash if any of this state update fails, otherwise dangling pointers. |
| 146 | |
| 147 | // This is wrapped in a KJ_DEFER because we want it to run no matter what *after* the rollback |
| 148 | // fix and KJ_DEFER seemed like the cleanest way to do this. |
| 149 | |
| 150 | KJ_ASSERT(!hasChild); |
| 151 | auto cur = KJ_ASSERT_NONNULL(actorSqlite.currentTxn.tryGet<ExplicitTxn*>()); |
| 152 | KJ_ASSERT(cur == this); |
| 153 | KJ_IF_SOME(p, parent) { |
| 154 | p.get()->hasChild = false; |
| 155 | actorSqlite.currentTxn = p.get(); |
| 156 | } else { |
| 157 | actorSqlite.currentTxn.init<NoTxn>(); |
| 158 | } |
| 159 | }();); |
| 160 | |
| 161 | if (!committed && actorSqlite.broken == kj::none) { |
| 162 | // Assume rollback if not committed. |
| 163 | rollbackImpl(); |
| 164 | } |
| 165 | } |
| 166 | |
| 167 | bool ActorSqlite::ExplicitTxn::getAlarmDirty() { |
| 168 | return alarmDirty; |
| 169 | } |
| 170 | |
| 171 | void ActorSqlite::ExplicitTxn::setAlarmDirty() { |
| 172 | alarmDirty = true; |
| 173 | } |
| 174 | |
| 175 | void ActorSqlite::ExplicitTxn::setSomeWriteConfirmed(bool someWriteConfirmed) { |
| 176 | this->someWriteConfirmed = someWriteConfirmed; |
| 177 | } |
| 178 | |
| 179 | bool ActorSqlite::ExplicitTxn::isSomeWriteConfirmed() const { |
| 180 | return someWriteConfirmed; |
| 181 | } |
| 182 | |
| 183 | kj::Maybe<kj::Promise<void>> ActorSqlite::ExplicitTxn::commit() { |
| 184 | actorSqlite.requireNotBroken(); |
| 185 | KJ_REQUIRE(!hasChild, |
| 186 | "critical sections should have prevented committing transaction while " |
| 187 | "nested txn is outstanding"); |
| 188 | |
| 189 | // Start the schedule request before root transaction commit(), for correctness in workerd. |
| 190 | kj::Maybe<PrecommitAlarmState> precommitAlarmState; |
| 191 | if (parent == kj::none) { |
| 192 | precommitAlarmState = actorSqlite.startPrecommitAlarmScheduling(); |
| 193 | } |
| 194 | |
| 195 | actorSqlite.db->run( |
| 196 | {.regulator = SqliteDatabase::TRUSTED}, kj::str("RELEASE _cf_savepoint_", depth)); |
| 197 | committed = true; |
| 198 | |
| 199 | KJ_IF_SOME(p, parent) { |
| 200 | if (alarmDirty) { |
| 201 | p->alarmDirty = true; |
| 202 | } |
| 203 | if (someWriteConfirmed) { |
| 204 | p->someWriteConfirmed = true; |
| 205 | } |
| 206 | } else { |
| 207 | if (alarmDirty) { |
| 208 | actorSqlite.haveDeferredDelete = false; |
| 209 | } |
| 210 | |
| 211 | // We committed the root transaction, so it's time to signal any replication layer and lock |
| 212 | // the output gate in the meantime. |
| 213 | |
| 214 | // Unlike ImplicitTxn, which locks the output gate at the start of the first write that requires |
| 215 | // confirmation, ExplicitTxn only locks when we're going to confirm the commit. I think this |
| 216 | // makes since given the explicit commit call. |
| 217 | auto commitPromise = kj::evalNow([this, &precommitAlarmState]() { |
| 218 | return actorSqlite.commitImpl( |
| 219 | kj::mv(KJ_ASSERT_NONNULL(precommitAlarmState)), actorSqlite.currentCommitSpan.addRef()); |
| 220 | }) |
| 221 | .catch_([outputGate = &actorSqlite.outputGate, |
| 222 | spanParent = actorSqlite.currentCommitSpan.addRef()]( |
| 223 | kj::Exception&& e) mutable { |
| 224 | // Unconditionally break the output gate if commit threw an error, no matter whether the |
| 225 | // commit was confirmed or unconfirmed. |
| 226 | return outputGate->lockWhile(kj::Promise<void>(kj::mv(e)), kj::mv(spanParent)); |
| 227 | }); |
| 228 | if (someWriteConfirmed) { |
| 229 | commitPromise = actorSqlite.outputGate.lockWhile( |
| 230 | kj::mv(commitPromise), actorSqlite.currentCommitSpan.addRef()); |
| 231 | } |
| 232 | auto forkedPromise = commitPromise.fork(); |
| 233 | actorSqlite.commitTasks.add(forkedPromise.addBranch()); |
| 234 | actorSqlite.lastCommit = kj::mv(forkedPromise); |
| 235 | } |
| 236 | |
| 237 | // No backpressure for SQLite. |
| 238 | return kj::none; |
| 239 | } |
| 240 | |
| 241 | kj::Promise<void> ActorSqlite::ExplicitTxn::rollback() { |
| 242 | actorSqlite.requireNotBroken(); |
| 243 | JSG_REQUIRE(!hasChild, Error, |
| 244 | "Cannot roll back an outer transaction while a nested transaction is still running."); |
| 245 | if (!committed) { |
| 246 | rollbackImpl(); |
| 247 | committed = true; |
| 248 | } |
| 249 | return kj::READY_NOW; |
| 250 | } |
| 251 | |
| 252 | void ActorSqlite::ExplicitTxn::rollbackImpl() noexcept(false) { |
| 253 | actorSqlite.db->run( |
| 254 | {.regulator = SqliteDatabase::TRUSTED}, kj::str("ROLLBACK TO _cf_savepoint_", depth)); |
| 255 | actorSqlite.db->run( |
| 256 | {.regulator = SqliteDatabase::TRUSTED}, kj::str("RELEASE _cf_savepoint_", depth)); |
| 257 | KJ_IF_SOME(p, parent) { |
| 258 | alarmDirty = p->alarmDirty; |
| 259 | someWriteConfirmed = p->someWriteConfirmed; |
| 260 | } else { |
| 261 | alarmDirty = false; |
| 262 | someWriteConfirmed = false; |
| 263 | } |
| 264 | } |
| 265 | |
| 266 | void ActorSqlite::onCriticalError( |
| 267 | kj::StringPtr errorMessage, kj::Maybe<kj::Exception> maybeException) { |
| 268 | // If we have already experienced a terminal exception, no need to replace it |
| 269 | if (broken == kj::none) { |
| 270 | kj::Exception exception = kj::mv(maybeException).orDefault([&]() { |
| 271 | return JSG_KJ_EXCEPTION(FAILED, Error, errorMessage); |
| 272 | }); |
| 273 | exception.setDescription(kj::str("broken.outputGateBroken; ", exception.getDescription())); |
| 274 | broken.emplace(exception.clone()); |
| 275 | |
| 276 | // Also ensure output gate is explicitly broken. |
| 277 | commitTasks.add( |
| 278 | outputGate.lockWhile(kj::Promise<void>(kj::mv(exception)), currentCommitSpan.addRef())); |
| 279 | } |
| 280 | } |
| 281 | |
| 282 | void ActorSqlite::startImplicitTxn() { |
| 283 | auto txn = kj::heap<ImplicitTxn>(*this); |
| 284 | |
| 285 | // We implement the magic of accumulating all of the writes between JavaScript awaits in one |
| 286 | // transaction by evaluating by wrapping the commit function with kj::evalLater, which runs the |
| 287 | // function on the next turn of the event loop |
| 288 | auto commitPromise = |
| 289 | kj::evalLater([this, txn = kj::mv(txn)]() mutable -> kj::Promise<void> { |
| 290 | // Don't commit if shutdown() has been called. |
| 291 | requireNotBroken(); |
| 292 | |
| 293 | // Start the schedule request before commit(), for correctness in workerd. |
| 294 | auto precommitAlarmState = startPrecommitAlarmScheduling(); |
| 295 | |
| 296 | try { |
| 297 | txn->commit(); |
| 298 | } catch (...) { |
| 299 | // HACK: If we became broken during `COMMIT TRANSACTION` then throw the broken exception |
| 300 | // instead of whatever SQLite threw. |
| 301 | requireNotBroken(); |
| 302 | |
| 303 | // No, we're not broken, so propagate the exception as-is. |
| 304 | throw; |
| 305 | } |
| 306 | |
| 307 | // The callback is only expected to commit writes up until this point. Any new writes that |
| 308 | // occur while the callback is in progress are NOT included, therefore require a new commit |
| 309 | // to be scheduled. So, we should drop `txn` to cause `currentTxn` to become NoTxn now, |
| 310 | // rather than after the callback. |
| 311 | { auto drop = kj::mv(txn); } |
| 312 | |
| 313 | // Move the commit span out immediately so new writes can capture a fresh span. |
| 314 | return commitImpl(kj::mv(precommitAlarmState), kj::mv(currentCommitSpan)); |
| 315 | }) |
| 316 | // Unconditionally break the output gate if commit threw an error, no matter whether the |
| 317 | // commit was confirmed or unconfirmed. |
| 318 | .catch_([this](kj::Exception&& e) { |
| 319 | return outputGate.lockWhile(kj::Promise<void>(kj::mv(e)), nullptr); |
| 320 | }) |
| 321 | // We need to wait for this in commitTasks and in lastCommit. |
| 322 | .fork(); |
| 323 | |
| 324 | commitTasks.add(commitPromise.addBranch()); |
| 325 | |
| 326 | // Commits must be executed in order, so we only have to track the most recent commit promise. |
| 327 | lastCommit = kj::mv(commitPromise); |
| 328 | } |
| 329 | |
| 330 | void ActorSqlite::onWrite(bool allowUnconfirmed) { |
| 331 | requireNotBroken(); |
| 332 | if (currentTxn.is<NoTxn>()) { |
| 333 | startImplicitTxn(); |
| 334 | } |
| 335 | |
| 336 | // Update the status of the current transaction. |
| 337 | KJ_SWITCH_ONEOF(currentTxn) { |
| 338 | KJ_CASE_ONEOF(_, NoTxn) { |
| 339 | KJ_FAIL_REQUIRE("we must have a transaction at this point"); |
| 340 | } |
| 341 | KJ_CASE_ONEOF(implicitTxn, ImplicitTxn*) { |
| 342 | if (!implicitTxn->isSomeWriteConfirmed() && !allowUnconfirmed) { |
| 343 | // This is adding a must-confirm write to the transaction, so we must ensure the outputGate |
| 344 | // locks for remainder of this transaction. |
| 345 | implicitTxn->setSomeWriteConfirmed(!allowUnconfirmed); |
| 346 | commitTasks.add(outputGate.lockWhile(lastCommit.addBranch(), currentCommitSpan.addRef())); |
| 347 | } |
| 348 | } |
| 349 | KJ_CASE_ONEOF(explicitTxn, ExplicitTxn*) { |
| 350 | if (!explicitTxn->isSomeWriteConfirmed() && !allowUnconfirmed) { |
| 351 | // This is adding a must-confirm write to the transaction, so we must ensure the outputGate |
| 352 | // locks for remainder of this transaction. |
| 353 | explicitTxn->setSomeWriteConfirmed(!allowUnconfirmed); |
| 354 | // ExplicitTxns don't have a pending commit and don't lock the output gate during the |
| 355 | // transaction, so there's nothing to do here. |
| 356 | } |
| 357 | } |
| 358 | } |
| 359 | } |
| 360 | |
| 361 | kj::Promise<void> ActorSqlite::requestScheduledAlarm( |
| 362 | kj::Maybe<kj::Date> requestedTime, kj::Promise<void> priorTask) { |
| 363 | // Not using coroutines here, because it's important for correctness in workerd that a |
| 364 | // synchronously thrown exception in scheduleRun() can escape synchronously to the caller. |
| 365 | |
| 366 | bool movingAlarmLater = willFireEarlier(alarmScheduledNoLaterThan, requestedTime); |
| 367 | if (movingAlarmLater) { |
| 368 | // Since we are setting the alarm to be later, we can update alarmScheduledNoLaterThan |
| 369 | // immediately and still preserve the invariant that the scheduled alarm time is equal to or |
| 370 | // earlier than the persisted db alarm value. Doing the immediate update ensures that |
| 371 | // subsequent invocations of commitImpl() will compare against the correct value in their |
| 372 | // precommit alarm checks, even if other later-setting requests are still in-flight, without |
| 373 | // needing to wait for them to complete. |
| 374 | alarmScheduledNoLaterThan = requestedTime; |
| 375 | } |
| 376 | |
| 377 | return hooks.scheduleRun(requestedTime, kj::mv(priorTask)) |
| 378 | .then([this, movingAlarmLater, requestedTime]() { |
| 379 | if (!movingAlarmLater) { |
| 380 | alarmScheduledNoLaterThan = requestedTime; |
| 381 | } |
| 382 | }); |
| 383 | } |
| 384 | |
| 385 | void ActorSqlite::scheduleLaterAlarm(kj::Maybe<kj::Date> newAlarmTime, SpanParent parentSpan) { |
| 386 | if (alarmLaterIsInFlight) { |
| 387 | // There's already a move-later request in-flight. Just store the desired time; the in-flight |
| 388 | // request's completion handler will pick it up and start a new request. This overwrites any |
| 389 | // previously pending time, which is fine -- only the latest value matters. |
| 390 | pendingLaterAlarmTime = newAlarmTime; |
| 391 | return; |
| 392 | } |
| 393 | |
| 394 | alarmLaterIsInFlight = true; |
| 395 | alarmLaterInFlight = requestScheduledAlarm(newAlarmTime, alarmLaterInFlight.addBranch()) |
| 396 | .attach(parentSpan.newChild("actor_sqlite_alarm_sync"_kjc)) |
| 397 | .catch_([](kj::Exception&& e) { |
| 398 | // If an exception occurs when scheduling the alarm later, it's OK -- the alarm will |
| 399 | // eventually fire at the earlier time, and the rescheduling will be retried. |
| 400 | // We catch here to prevent the chain from breaking on errors. |
| 401 | LOG_WARNING_PERIODICALLY("NOSENTRY SQLite reschedule later alarm failed", e); |
| 402 | }).fork(); |
| 403 | |
| 404 | commitTasks.add(alarmLaterInFlight.addBranch() |
| 405 | .then([this]() { |
| 406 | alarmLaterIsInFlight = false; |
| 407 | KJ_IF_SOME(nextTime, kj::mv(pendingLaterAlarmTime)) { |
| 408 | scheduleLaterAlarm(nextTime, nullptr); |
| 409 | } |
| 410 | }).catch_([](kj::Exception&& e) { |
| 411 | // Move-later alarm failures are non-fatal; catch here to prevent taskFailed() from |
| 412 | // breaking the output gate. |
| 413 | LOG_WARNING_PERIODICALLY("NOSENTRY SQLite reschedule later alarm drain failed", e); |
| 414 | })); |
| 415 | } |
| 416 | |
| 417 | ActorSqlite::PrecommitAlarmState ActorSqlite::startPrecommitAlarmScheduling() { |
| 418 | PrecommitAlarmState state; |
| 419 | if (pendingCommit == kj::none && |
| 420 | willFireEarlier(metadata.getAlarm(), alarmScheduledNoLaterThan)) { |
| 421 | // We must wait on the `alarmLaterInFlight` promise here, otherwise, if there is an in-flight |
| 422 | // "move later" alarm task and it fails, our "move earlier" alarm might interleave, succeed, |
| 423 | // and be followed by a retry of the in-flight "move later" alarm. This happens because "move later" |
| 424 | // alarms complete after we commit to local SQLite. |
| 425 | // |
| 426 | // By waiting on any in-flight "move later" alarm, we correctly serialize our `scheduleRun()` |
| 427 | // calls to the alarm manager. |
| 428 | // Clear any pending move-later alarm time. Since we are about to move the alarm |
| 429 | // earlier, any coalesced later time is now obsolete. This also prevents the |
| 430 | // scheduleLaterAlarm completion handler from starting a concurrent scheduleRun |
| 431 | // when it drains pendingLaterAlarmTime after the current in-flight request resolves. |
| 432 | pendingLaterAlarmTime = kj::none; |
| 433 | state.schedulingPromise = |
| 434 | requestScheduledAlarm(metadata.getAlarm(), alarmLaterInFlight.addBranch()); |
| 435 | } |
| 436 | return kj::mv(state); |
| 437 | } |
| 438 | |
| 439 | kj::Promise<void> ActorSqlite::commitImpl( |
| 440 | ActorSqlite::PrecommitAlarmState precommitAlarmState, SpanParent parentSpan) { |
| 441 | auto commitSpan = parentSpan.newChild("actor_sqlite_commit"_kjc); |
| 442 | |
| 443 | // We assume that exceptions thrown during commit will propagate to the caller, such that they |
| 444 | // will ensure cancelDeferredAlarmDeletion() is called, if necessary. |
| 445 | |
| 446 | bool haveAlarmForDebug = false; |
| 447 | |
| 448 | if (debugAlarmSync) { |
| 449 | KJ_LOG(WARNING, "NOSENTRY DEBUG_ALARM: commitImpl entered", (pendingCommit != kj::none), |
| 450 | alarmVersion, logDate(metadata.getAlarm())); |
| 451 | } |
| 452 | |
| 453 | KJ_IF_SOME(pending, pendingCommit) { |
| 454 | // If an earlier commitImpl() invocation is already in the process of updating precommit |
| 455 | // alarms but has not yet made the commitCallback() call, it should be OK to wait on it to |
| 456 | // perform the precommit alarm update and db commit for this invocation, too. |
| 457 | kj::Maybe<kj::Date> alarmBeforeMerge; |
| 458 | if (debugAlarmSync) { |
| 459 | alarmBeforeMerge = metadata.getAlarm(); |
| 460 | KJ_LOG(WARNING, "NOSENTRY DEBUG_ALARM: Commit merge waiting", logDate(alarmBeforeMerge), |
| 461 | alarmVersion); |
| 462 | } |
| 463 | commitSpan.setTag("merged_with_pending_commit"_kjc, true); |
| 464 | co_await pending.addBranch(); |
| 465 | if (debugAlarmSync) { |
| 466 | auto alarmAfterMerge = metadata.getAlarm(); |
| 467 | KJ_LOG(WARNING, "NOSENTRY DEBUG_ALARM: Commit merge resumed", logDate(alarmBeforeMerge), |
| 468 | logDate(alarmAfterMerge), alarmVersion); |
| 469 | } |
| 470 | co_return; |
| 471 | } |
| 472 | |
| 473 | // There are no pending commits in-flight, so we set up a forked promise that other callers can |
| 474 | // wait on, to perform the alarm scheduling and database persistence work for all of them. Note |
| 475 | // that the fulfiller is owned by this coroutine context, so if an exception is thrown below, |
| 476 | // the fulfiller's destructor will detect that the stack is unwinding and will automatically |
| 477 | // propagate the thrown exception to the other waiters. |
| 478 | auto [promise, fulfiller] = kj::newPromiseAndFulfiller<void>(); |
| 479 | pendingCommit = promise.fork(); |
| 480 | |
| 481 | // Wait for the first precommit alarm scheduling request to complete, if any. This was set up |
| 482 | // in startPrecommitAlarmScheduling() and is essentially the first iteration of the below |
| 483 | // while() loop, but needed to be initiated synchronously before the local database commit to |
| 484 | // ensure correctness in workerd. |
| 485 | KJ_IF_SOME(p, precommitAlarmState.schedulingPromise) { |
| 486 | auto alarmSpan = commitSpan.newChild("actor_sqlite_alarm_sync"_kjc); |
| 487 | haveAlarmForDebug = true; |
| 488 | co_await p; |
| 489 | } |
| 490 | |
| 491 | // While the local db state requires an earlier alarm than is known might be scheduled, issue an |
| 492 | // alarm update request for the earlier time and wait for it to complete. This helps ensure |
| 493 | // that the successfully scheduled alarm time is always earlier or equal to the alarm state in |
| 494 | // the successfully persisted db. |
| 495 | int syncIterations = 0; |
| 496 | auto startAlarmState = metadata.getAlarm(); |
| 497 | while (willFireEarlier(metadata.getAlarm(), alarmScheduledNoLaterThan)) { |
| 498 | auto alarmSpan = commitSpan.newChild("actor_sqlite_alarm_sync"_kjc); |
| 499 | if (debugAlarmSync) { |
| 500 | haveAlarmForDebug = true; |
| 501 | auto currentAlarmState = metadata.getAlarm(); |
| 502 | KJ_LOG(WARNING, "NOSENTRY DEBUG_ALARM: Move earlier loop iteration", syncIterations, |
| 503 | logDate(currentAlarmState), logDate(alarmScheduledNoLaterThan), alarmVersion); |
| 504 | } |
| 505 | // Note that we do not pass alarmLaterInFlight here. We don't need to for the following |
| 506 | // reasons: |
| 507 | // |
| 508 | // 1. We already waited for it in the precommitAlarmState promise above. |
| 509 | // 2. We set the `pendingCommit` prior to yielding to the event loop earlier, so any subsequent |
| 510 | // commits have to wait for us to fulfill the pendingCommit promise. In short, no one could |
| 511 | // have started another "move-later" alarm, not until we finish. |
| 512 | // |
| 513 | // While we *could* pass the alarmLaterInFlight promise (it wouldn't be incorrect), when |
| 514 | // calling addBranch() on a resolved ForkedPromise, the continuation would be evaluated on a |
| 515 | // future turn of the event loop. That means we're going to suspend, even if the promise is |
| 516 | // ready, which means we'd take a performance hit. |
| 517 | co_await requestScheduledAlarm(metadata.getAlarm(), kj::READY_NOW); |
| 518 | syncIterations++; |
| 519 | } |
| 520 | if (debugAlarmSync && syncIterations > 0) { |
| 521 | KJ_LOG(WARNING, "NOSENTRY DEBUG_ALARM: Move earlier loop complete", logDate(startAlarmState), |
| 522 | "ended_with", logDate(metadata.getAlarm()), "iterations", syncIterations, alarmVersion); |
| 523 | } |
| 524 | |
| 525 | // Issue the commitCallback() request to persist the db state, then synchronously clear the |
| 526 | // pending commit so that the next commitImpl() invocation starts its own set of precommit alarm |
| 527 | // updates and db commit. |
| 528 | auto alarmStateForCommit = metadata.getAlarm(); |
| 529 | |
| 530 | // Capture the alarm version before going async to detect concurrent alarm changes. If the |
| 531 | // alarmVersion changes while we are in-flight, we should skip attempting any move-later alarm |
| 532 | // update. |
| 533 | auto alarmVersionBeforeAsync = alarmVersion; |
| 534 | |
| 535 | if (debugAlarmSync) { |
| 536 | KJ_LOG(WARNING, "NOSENTRY DEBUG_ALARM: Captured state before persisting to SQLite async", |
| 537 | logDate(alarmStateForCommit), alarmVersionBeforeAsync); |
| 538 | } |
| 539 | |
| 540 | auto commitCallbackPromise = commitCallback(SpanParent(commitSpan)); |
| 541 | pendingCommit = kj::none; |
| 542 | |
| 543 | // Wait for the db to persist. |
| 544 | co_await commitCallbackPromise; |
| 545 | lastConfirmedAlarmDbState = alarmStateForCommit; |
| 546 | |
| 547 | if (debugAlarmSync && haveAlarmForDebug) { |
| 548 | KJ_LOG(WARNING, "NOSENTRY DEBUG_ALARM: Persisted in SQLite", "sqlite_has", |
| 549 | logDate(alarmStateForCommit), "alarmScheduledNoLaterThan", |
| 550 | logDate(alarmScheduledNoLaterThan), alarmVersion); |
| 551 | } |
| 552 | |
| 553 | // Notify any merged commitImpl() requests that the db persistence completed. |
| 554 | fulfiller->fulfill(); |
| 555 | |
| 556 | if (debugAlarmSync) { |
| 557 | KJ_LOG(WARNING, "NOSENTRY DEBUG_ALARM: Version check", alarmVersionBeforeAsync, alarmVersion, |
| 558 | "match", (alarmVersion == alarmVersionBeforeAsync)); |
| 559 | } |
| 560 | // If another commit modified the alarm while we were async, skip post-commit alarm sync. |
| 561 | // |
| 562 | // We do this for a few reasons: |
| 563 | // 1. The other commit will handle its own alarm sync |
| 564 | // 2. Post-commit syncs are inherently optional (the alarm will self-correct) |
| 565 | // 3. This coalesces redundant alarm updates for better performance |
| 566 | // 4. This avoids race conditions where a later commit moved the alarm earlier, requiring a |
| 567 | // pre-commit alarm update, and this update may have already been made before we get here. |
| 568 | if (alarmVersion == alarmVersionBeforeAsync) { |
| 569 | // No intervening alarm changes, it is safe to schedule a move-later alarm update if needed. |
| 570 | if (willFireEarlier(alarmScheduledNoLaterThan, alarmStateForCommit)) { |
| 571 | if (debugAlarmSync) { |
| 572 | KJ_LOG(WARNING, "NOSENTRY DEBUG_ALARM: Moving alarm later", "sqlite_has", |
| 573 | logDate(alarmStateForCommit), logDate(alarmScheduledNoLaterThan), alarmVersion); |
| 574 | } |
| 575 | scheduleLaterAlarm(alarmStateForCommit, commitSpan); |
| 576 | } |
| 577 | } |
| 578 | } |
| 579 | |
| 580 | void ActorSqlite::taskFailed(kj::Exception&& exception) { |
| 581 | // The output gate should already have been broken since it wraps all commit tasks that can |
| 582 | // throw. So, we don't have to report anything here, the exception will already propagate |
| 583 | // elsewhere. We should block further operations, though. |
| 584 | if (broken == kj::none) { |
| 585 | broken = kj::mv(exception); |
| 586 | if (!outputGate.isBroken()) { |
| 587 | LOG_PERIODICALLY( |
| 588 | ERROR, "SQLite actor recorded broken exception without breaking output gate"); |
| 589 | } |
| 590 | } |
| 591 | } |
| 592 | |
| 593 | void ActorSqlite::requireNotBroken() { |
| 594 | KJ_IF_SOME(e, broken) { |
| 595 | kj::throwFatalException(e.clone()); |
| 596 | } |
| 597 | } |
| 598 | |
| 599 | void ActorSqlite::maybeDeleteDeferredAlarm() { |
| 600 | if (!inAlarmHandler) { |
| 601 | // Pretty sure this can't happen. |
| 602 | LOG_WARNING_ONCE("expected to be in alarm handler when trying to delete alarm"); |
| 603 | } |
| 604 | inAlarmHandler = false; |
| 605 | |
| 606 | if (haveDeferredDelete) { |
| 607 | // If we have reached this point, the client is destroying its DeferredAlarmDeleter at the end |
| 608 | // of an alarm handler run, and deletion hasn't been cancelled, indicating that the handler |
| 609 | // returned success. |
| 610 | // |
| 611 | // If the output gate has somehow broken in the interim, attempting to write the deletion here |
| 612 | // will cause the DeferredAlarmDeleter destructor to throw, which the caller probably isn't |
| 613 | // expecting. So we'll skip the deletion attempt, and let the caller detect the gate |
| 614 | // brokenness through other means. |
| 615 | if (broken == kj::none) { |
| 616 | // Use the span captured in armAlarmHandler() for this internal write, since |
| 617 | // metadata.setAlarm() doesn't go through the regular write path with a traceSpan parameter. |
| 618 | currentCommitSpan = kj::mv(deferredAlarmSpan); |
| 619 | // the safe thing to do is to require confirmation. |
| 620 | if (metadata.setAlarm(kj::none, /*allowUnconfirmed=*/false)) { |
| 621 | ++alarmVersion; |
| 622 | if (debugAlarmSync) { |
| 623 | KJ_LOG(WARNING, "NOSENTRY DEBUG_ALARM: maybeDeleteDeferredAlarm cleared alarm", |
| 624 | alarmVersion); |
| 625 | } |
| 626 | } |
| 627 | } |
| 628 | haveDeferredDelete = false; |
| 629 | deferredAlarmSpan = nullptr; |
| 630 | } |
| 631 | } |
| 632 | |
| 633 | // ======================================================================================= |
| 634 | // ActorCacheInterface implementation |
| 635 | |
| 636 | kj::OneOf<kj::Maybe<ActorCacheOps::Value>, kj::Promise<kj::Maybe<ActorCacheOps::Value>>> |
| 637 | ActorSqlite::get(Key key, ReadOptions options) { |
| 638 | requireNotBroken(); |
| 639 | |
| 640 | kj::Maybe<ActorCacheOps::Value> result; |
| 641 | kv.get(key, [&](ValuePtr value) { result = kj::heapArray(value); }); |
| 642 | return result; |
| 643 | } |
| 644 | |
| 645 | kj::OneOf<ActorCacheOps::GetResultList, kj::Promise<ActorCacheOps::GetResultList>> ActorSqlite::get( |
| 646 | kj::Array<Key> keys, ReadOptions options) { |
| 647 | requireNotBroken(); |
| 648 | |
| 649 | kj::Vector<KeyValuePair> results; |
| 650 | for (auto& key: keys) { |
| 651 | kv.get( |
| 652 | key, [&](ValuePtr value) { results.add(KeyValuePair{kj::mv(key), kj::heapArray(value)}); }); |
| 653 | } |
| 654 | std::sort(results.begin(), results.end(), [](auto& a, auto& b) { return a.key < b.key; }); |
| 655 | return GetResultList(kj::mv(results)); |
| 656 | } |
| 657 | |
| 658 | kj::OneOf<kj::Maybe<kj::Date>, kj::Promise<kj::Maybe<kj::Date>>> ActorSqlite::getAlarm( |
| 659 | ReadOptions options) { |
| 660 | requireNotBroken(); |
| 661 | |
| 662 | bool transactionAlarmDirty = false; |
| 663 | KJ_IF_SOME(exp, currentTxn.tryGet<ExplicitTxn*>()) { |
| 664 | transactionAlarmDirty = exp->getAlarmDirty(); |
| 665 | } |
| 666 | |
| 667 | if (haveDeferredDelete && !transactionAlarmDirty) { |
| 668 | // If an alarm handler is currently running, and a new alarm time has not been set yet, We |
| 669 | // need to return that there is no alarm. |
| 670 | return kj::Maybe<kj::Date>(kj::none); |
| 671 | } else { |
| 672 | return metadata.getAlarm(); |
| 673 | } |
| 674 | KJ_UNREACHABLE; |
| 675 | } |
| 676 | |
| 677 | kj::OneOf<ActorCacheOps::GetResultList, kj::Promise<ActorCacheOps::GetResultList>> ActorSqlite:: |
| 678 | list(Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) { |
| 679 | requireNotBroken(); |
| 680 | |
| 681 | kj::Vector<KeyValuePair> results; |
| 682 | kv.list(begin, end, limit, SqliteKv::FORWARD, [&](KeyPtr key, ValuePtr value) { |
| 683 | results.add(KeyValuePair{kj::str(key), kj::heapArray(value)}); |
| 684 | }); |
| 685 | |
| 686 | // Already guaranteed sorted. |
| 687 | return GetResultList(kj::mv(results)); |
| 688 | } |
| 689 | |
| 690 | kj::OneOf<ActorCacheOps::GetResultList, kj::Promise<ActorCacheOps::GetResultList>> ActorSqlite:: |
| 691 | listReverse(Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) { |
| 692 | requireNotBroken(); |
| 693 | |
| 694 | kj::Vector<KeyValuePair> results; |
| 695 | kv.list(begin, end, limit, SqliteKv::REVERSE, [&](KeyPtr key, ValuePtr value) { |
| 696 | results.add(KeyValuePair{kj::str(key), kj::heapArray(value)}); |
| 697 | }); |
| 698 | |
| 699 | // Already guaranteed sorted (reversed). |
| 700 | return GetResultList(kj::mv(results)); |
| 701 | } |
| 702 | |
| 703 | kj::Maybe<kj::Promise<void>> ActorSqlite::put( |
| 704 | Key key, Value value, WriteOptions options, SpanParent traceSpan) { |
| 705 | requireNotBroken(); |
| 706 | // Capture trace span for the output gate lock hold trace. |
| 707 | currentCommitSpan = kj::mv(traceSpan); |
| 708 | kv.put(key, value, {.allowUnconfirmed = options.allowUnconfirmed}); |
| 709 | return kj::none; |
| 710 | } |
| 711 | |
| 712 | kj::Maybe<kj::Promise<void>> ActorSqlite::put( |
| 713 | kj::Array<KeyValuePair> pairs, WriteOptions options, SpanParent traceSpan) { |
| 714 | requireNotBroken(); |
| 715 | // Capture trace span for the output gate lock hold trace. |
| 716 | currentCommitSpan = kj::mv(traceSpan); |
| 717 | if (currentTxn.is<NoTxn>()) { |
| 718 | // If we are not in a transaction, start an ImplicitTxn since that's what would happen on the |
| 719 | // first write anyway. |
| 720 | startImplicitTxn(); |
| 721 | } |
| 722 | |
| 723 | KJ_ASSERT(!currentTxn.is<NoTxn>()); |
| 724 | |
| 725 | kv.put(pairs, {.allowUnconfirmed = options.allowUnconfirmed}); |
| 726 | |
| 727 | return kj::none; |
| 728 | } |
| 729 | |
| 730 | kj::OneOf<bool, kj::Promise<bool>> ActorSqlite::delete_( |
| 731 | Key key, WriteOptions options, SpanParent traceSpan) { |
| 732 | requireNotBroken(); |
| 733 | // Capture trace span for the output gate lock hold trace. |
| 734 | currentCommitSpan = kj::mv(traceSpan); |
| 735 | |
| 736 | return kv.delete_(key, {.allowUnconfirmed = options.allowUnconfirmed}); |
| 737 | } |
| 738 | |
| 739 | kj::OneOf<uint, kj::Promise<uint>> ActorSqlite::delete_( |
| 740 | kj::Array<Key> keys, WriteOptions options, SpanParent traceSpan) { |
| 741 | requireNotBroken(); |
| 742 | // Capture trace span for the output gate lock hold trace. |
| 743 | currentCommitSpan = kj::mv(traceSpan); |
| 744 | |
| 745 | uint count = 0; |
| 746 | for (auto& key: keys) { |
| 747 | count += kv.delete_(key, {.allowUnconfirmed = options.allowUnconfirmed}); |
| 748 | } |
| 749 | return count; |
| 750 | } |
| 751 | |
| 752 | kj::Maybe<kj::Promise<void>> ActorSqlite::setAlarm( |
| 753 | kj::Maybe<kj::Date> newAlarmTime, WriteOptions options, SpanParent traceSpan) { |
| 754 | requireNotBroken(); |
| 755 | // Capture trace span for the output gate lock hold trace. |
| 756 | currentCommitSpan = kj::mv(traceSpan); |
| 757 | |
| 758 | // TODO(someday): When deleting alarm data in an otherwise empty database, clear the database to |
| 759 | // free up resources? |
| 760 | |
| 761 | // Only increment version counter if the alarm value actually changed. This is important because |
| 762 | // if the value didn't change, no SQLite write occurs, so no implicit transaction is started, |
| 763 | // and we don't want to invalidate in-flight commits without a replacement commit. |
| 764 | if (metadata.setAlarm(newAlarmTime, options.allowUnconfirmed)) { |
| 765 | ++alarmVersion; |
| 766 | if (debugAlarmSync) { |
| 767 | KJ_LOG(WARNING, "NOSENTRY DEBUG_ALARM: setAlarm called", logDate(newAlarmTime), alarmVersion); |
| 768 | } |
| 769 | } |
| 770 | |
| 771 | KJ_IF_SOME(exp, currentTxn.tryGet<ExplicitTxn*>()) { |
| 772 | exp->setAlarmDirty(); |
| 773 | } else { |
| 774 | haveDeferredDelete = false; |
| 775 | } |
| 776 | |
| 777 | return kj::none; |
| 778 | } |
| 779 | |
| 780 | kj::Own<ActorCacheInterface::Transaction> ActorSqlite::startTransaction() { |
| 781 | requireNotBroken(); |
| 782 | |
| 783 | return kj::refcounted<ExplicitTxn>(*this); |
| 784 | } |
| 785 | |
| 786 | ActorCacheInterface::DeleteAllResults ActorSqlite::deleteAll( |
| 787 | WriteOptions options, SpanParent traceSpan, DeleteAllOptions deleteAllOptions) { |
| 788 | requireNotBroken(); |
| 789 | disableAllowUnconfirmed(options, "deleteAll is not supported"); |
| 790 | |
| 791 | // Capture trace span for the output gate lock (deleteAll always requires confirmation). |
| 792 | currentCommitSpan = kj::mv(traceSpan); |
| 793 | |
| 794 | // kv.deleteAll() clears the database, so we need to save and possibly restore alarm state in |
| 795 | // the metadata table to maintain behavior from before the deleteAllDeletesAlarm compat flag. |
| 796 | auto localAlarmState = metadata.getAlarm(); |
| 797 | |
| 798 | // deleteAll() cannot be part of a transaction because it deletes the database altogether. So, |
| 799 | // we have to close our transactions or fail. |
| 800 | KJ_SWITCH_ONEOF(currentTxn) { |
| 801 | KJ_CASE_ONEOF(_, NoTxn) { |
| 802 | // good |
| 803 | } |
| 804 | KJ_CASE_ONEOF(implicit, ImplicitTxn*) { |
| 805 | // Whatever the implicit transaction did, it's about to be blown away anyway. Roll it back |
| 806 | // so we don't waste time flushing these writes anywhere. |
| 807 | implicit->rollback(); |
| 808 | currentTxn = NoTxn(); |
| 809 | } |
| 810 | KJ_CASE_ONEOF(exp, ExplicitTxn*) { |
| 811 | // Keep in mind: |
| 812 | // |
| 813 | // ctx.storage.transaction(txn => { |
| 814 | // txn.deleteAll(); // calls `DurableObjectTransaction::deleteAll()` |
| 815 | // ctx.storage.deleteAll(); // calls this method, `ActorSqlite::deleteAll()` |
| 816 | // }); |
| 817 | // |
| 818 | // `DurableObjectTransaction::deleteAll()` throws this exception, since `deleteAll()` is not |
| 819 | // supported inside a transaction. Under the new SQLite-backed storage system, directly |
| 820 | // calling `cxt.storage` inside a transaction (as opposed to using the `txn` object) should |
| 821 | // still be treated as part of the transaction, and so should throw the same thing. |
| 822 | JSG_FAIL_REQUIRE(Error, "Cannot call deleteAll() within a transaction"); |
| 823 | } |
| 824 | } |
| 825 | |
| 826 | if (!deleteAllCommitScheduled) { |
| 827 | // Make sure a commit callback is queued for the deleteAll(). |
| 828 | commitTasks.add(outputGate.lockWhile(kj::evalLater([this]() mutable -> kj::Promise<void> { |
| 829 | // Don't commit if shutdown() has been called. |
| 830 | requireNotBroken(); |
| 831 | |
| 832 | deleteAllCommitScheduled = false; |
| 833 | if (currentTxn.is<ImplicitTxn*>()) { |
| 834 | // An implicit transaction is already scheduled, so we'll count on it to perform a commit when it's |
| 835 | // done. This is particularly important for the case where deleteAll() was called while an alarm |
| 836 | // is outstanding; resetting the alarm state (below) starts an implicit transaction. |
| 837 | // We don't want to commit the deletion without that transaction. |
| 838 | return kj::READY_NOW; |
| 839 | } else { |
| 840 | // Use commitImpl() rather than commitCallback() so that alarm scheduling is handled. |
| 841 | // This is important when deleteAll() deletes an alarm: commitImpl() detects that |
| 842 | // metadata.getAlarm() moved to kj::none and notifies the scheduler via |
| 843 | // requestScheduledAlarm(kj::none, ...). |
| 844 | auto precommitAlarmState = startPrecommitAlarmScheduling(); |
| 845 | return commitImpl(kj::mv(precommitAlarmState), currentCommitSpan.addRef()); |
| 846 | } |
| 847 | }), |
| 848 | currentCommitSpan.addRef())); |
| 849 | deleteAllCommitScheduled = true; |
| 850 | } |
| 851 | |
| 852 | uint count = kv.deleteAll(); |
| 853 | |
| 854 | // Reset alarm state, if necessary. If no alarm is set, OK to just leave metadata table |
| 855 | // uninitialized. |
| 856 | if (localAlarmState != kj::none) { |
| 857 | if (deleteAllOptions.deleteAlarm) { |
| 858 | // The caller wants the alarm deleted along with KV data. Since kv.deleteAll() already |
| 859 | // wiped the database (including the alarm metadata), metadata.getAlarm() will naturally |
| 860 | // return kj::none without creating any tables or rows. Increment alarmVersion so in‑flight |
| 861 | // commits don’t perform stale post‑commit alarm scheduling, and the deleteAll commit can sync |
| 862 | // cancellation. |
| 863 | ++alarmVersion; |
| 864 | haveDeferredDelete = false; |
| 865 | } else { |
| 866 | // TODO(correctness): Since workerd doesn't have a separate durability step, in the unlikely |
| 867 | // event of a failure here, between deleteAll() and setAlarm(), we could theoretically lose the |
| 868 | // current alarm state when running under workerd. Not sure if there's a practical way to avoid |
| 869 | // this. |
| 870 | if (metadata.setAlarm(localAlarmState, options.allowUnconfirmed)) { |
| 871 | ++alarmVersion; |
| 872 | if (debugAlarmSync) { |
| 873 | KJ_LOG(WARNING, "NOSENTRY DEBUG_ALARM: deleteAll restored alarm", |
| 874 | logDate(localAlarmState), alarmVersion); |
| 875 | } |
| 876 | } |
| 877 | } |
| 878 | } |
| 879 | |
| 880 | return { |
| 881 | .backpressure = kj::none, |
| 882 | .count = count, |
| 883 | }; |
| 884 | } |
| 885 | |
| 886 | kj::Maybe<kj::Promise<void>> ActorSqlite::evictStale(kj::Date now) { |
| 887 | // This implementation never needs to apply backpressure. |
| 888 | return kj::none; |
| 889 | } |
| 890 | |
| 891 | void ActorSqlite::shutdown(kj::Maybe<const kj::Exception&> maybeException) { |
| 892 | // TODO(cleanup): Logic copied from ActorCache::shutdown(). Should they share somehow? |
| 893 | |
| 894 | if (broken == kj::none) { |
| 895 | auto exception = [&]() { |
| 896 | KJ_IF_SOME(e, maybeException) { |
| 897 | // We were given an exception, use it. |
| 898 | return e.clone(); |
| 899 | } |
| 900 | |
| 901 | // Use the direct constructor so that we can reuse the constexpr message variable for testing. |
| 902 | auto exception = kj::Exception(kj::Exception::Type::DISCONNECTED, __FILE__, __LINE__, |
| 903 | kj::heapString(ActorCache::SHUTDOWN_ERROR_MESSAGE)); |
| 904 | |
| 905 | // Add trace info sufficient to tell us which operation caused the failure. |
| 906 | exception.addTraceHere(); |
| 907 | exception.addTrace(__builtin_return_address(0)); |
| 908 | return exception; |
| 909 | }(); |
| 910 | |
| 911 | // Any scheduled flushes will fail once `flushImpl()` is invoked and notices that |
| 912 | // `maybeTerminalException` has a value. Any in-flight flushes will continue to run in the |
| 913 | // background. Remember that these in-flight flushes may or may not be awaited by the worker, |
| 914 | // but they still hold the output lock as long as `allowUnconfirmed` wasn't used. |
| 915 | broken.emplace(kj::mv(exception)); |
| 916 | |
| 917 | // We explicitly do not schedule a flush to break the output gate. This means that if a request |
| 918 | // is ongoing after the actor cache is shutting down, the output gate is only broken if they |
| 919 | // had to send a flush after shutdown, either from a scheduled flush or a retry after failure. |
| 920 | } else { |
| 921 | // We've already experienced a terminal exception either from shutdown or OOM, there should |
| 922 | // already be a flush scheduled that will break the output gate. |
| 923 | } |
| 924 | } |
| 925 | |
| 926 | kj::OneOf<ActorSqlite::CancelAlarmHandler, ActorSqlite::RunAlarmHandler> ActorSqlite:: |
| 927 | armAlarmHandler(kj::Date scheduledTime, |
| 928 | SpanParent parentSpan, |
| 929 | kj::Date currentTime, |
| 930 | bool noCache, |
| 931 | kj::StringPtr actorId) { |
| 932 | KJ_ASSERT(!inAlarmHandler); |
| 933 | |
| 934 | if (haveDeferredDelete) { |
| 935 | // Unlikely to happen, unless caller is starting new alarm handler before previous alarm |
| 936 | // handler cleanup has completed. |
| 937 | LOG_WARNING_ONCE("expected previous alarm handler to be cleaned up"); |
| 938 | } |
| 939 | |
| 940 | auto localAlarmState = metadata.getAlarm(); |
| 941 | if (localAlarmState != scheduledTime) { |
| 942 | if (localAlarmState == lastConfirmedAlarmDbState) { |
| 943 | // If the local alarm time is already in the past, just run the handler now. This avoids |
| 944 | // blocking alarm execution on the AlarmManager sync when storage is overloaded. The alarm |
| 945 | // will either delete itself on success or reschedule on failure. |
| 946 | if ((willFireEarlier(localAlarmState, currentTime))) { |
| 947 | auto localAlarmTime = KJ_ASSERT_NONNULL(localAlarmState); |
| 948 | LOG_WARNING_PERIODICALLY( |
| 949 | "NOSENTRY SQLite alarm overdue, running despite AlarmManager mismatch", scheduledTime, |
| 950 | localAlarmTime, currentTime, actorId); |
| 951 | haveDeferredDelete = true; |
| 952 | inAlarmHandler = true; |
| 953 | deferredAlarmSpan = kj::mv(parentSpan); |
| 954 | static const DeferredAlarmDeleter disposer; |
| 955 | return RunAlarmHandler{.deferredDelete = kj::Own<void>(this, disposer)}; |
| 956 | } |
| 957 | |
| 958 | // If there's a clean db time that differs from the requested handler's scheduled time, this |
| 959 | // run should be canceled. |
| 960 | if (willFireEarlier(scheduledTime, localAlarmState)) { |
| 961 | // If the handler's scheduled time is earlier than the clean scheduled time, we may be |
| 962 | // recovering from a failed db commit or scheduling request, so we need to request that |
| 963 | // the alarm be rescheduled for the current db time, and tell the caller to wait for |
| 964 | // successful rescheduling before cancelling the current handler invocation. |
| 965 | // |
| 966 | // TODO(perf): If we already have such a rescheduling request in-flight, might want to |
| 967 | // coalesce with the existing request? |
| 968 | LOG_WARNING_PERIODICALLY( |
| 969 | "NOSENTRY SQLite alarm handler canceled with requestScheduledAlarm.", scheduledTime, |
| 970 | localAlarmState.orDefault(kj::UNIX_EPOCH), actorId); |
| 971 | |
| 972 | // Since we're requesting to move the alarm time to later, we need to update the |
| 973 | // alarmLaterInFlight promise. We issue a single requestScheduledAlarm call, fork it, |
| 974 | // and use one branch for alarmLaterInFlight (with error catching so the fork remains |
| 975 | // usable) and the other for the CancelAlarmHandler return value (which propagates |
| 976 | // errors to the caller). We directly update alarmLaterInFlight here rather than using |
| 977 | // scheduleLaterAlarm(), because we need that separate un-caught branch. |
| 978 | auto schedulingPromise = |
| 979 | requestScheduledAlarm(localAlarmState, alarmLaterInFlight.addBranch()).fork(); |
| 980 | // Clear any stale pending time so that when the existing completion handler |
| 981 | // fires it does not start a redundant scheduleLaterAlarm for the same time that |
| 982 | // armAlarmHandler is already scheduling. |
| 983 | pendingLaterAlarmTime = kj::none; |
| 984 | alarmLaterInFlight = schedulingPromise.addBranch() |
| 985 | .catch_([](kj::Exception&& e) { |
| 986 | // If an exception occurs when scheduling the alarm later, it's OK -- the alarm will |
| 987 | // eventually fire at the earlier time, and the rescheduling will be retried. |
| 988 | // We catch here to prevent the chain from breaking on errors. |
| 989 | LOG_WARNING_PERIODICALLY("NOSENTRY SQLite reschedule later alarm failed", e); |
| 990 | }).fork(); |
| 991 | return CancelAlarmHandler{.waitBeforeCancel = schedulingPromise.addBranch()}; |
| 992 | } else { |
| 993 | // We have a clean local alarm time that is earlier than the handler's scheduled time, |
| 994 | // which suggests that either the alarm manager is working with stale data or that local |
| 995 | // alarm time has somehow gotten out of sync with the scheduled alarm time. |
| 996 | |
| 997 | // We know localAlarmState has a value here because we're in the branch where it's earlier |
| 998 | // than scheduledTime (not equal, and not later). |
| 999 | auto localTime = KJ_ASSERT_NONNULL(localAlarmState); |
| 1000 | |
| 1001 | // Only log if the alarm manager is significantly late (>10 seconds behind SQLite) |
| 1002 | if (scheduledTime - localTime > 10 * kj::SECONDS) { |
| 1003 | LOG_WARNING_PERIODICALLY( |
| 1004 | "NOSENTRY SQLite alarm handler canceled.", scheduledTime, actorId, localTime); |
| 1005 | } |
| 1006 | |
| 1007 | // Tell the caller to wait for successful rescheduling before cancelling the current |
| 1008 | // handler invocation. |
| 1009 | // |
| 1010 | // We pass kj::READY_NOW because being in this branch (SQLite is ahead of the alarm manager) |
| 1011 | // means there's no recent move-later operation to wait for, so no need for alarmLaterInFlight. |
| 1012 | return CancelAlarmHandler{ |
| 1013 | .waitBeforeCancel = requestScheduledAlarm(localAlarmState, kj::READY_NOW)}; |
| 1014 | } |
| 1015 | } else { |
| 1016 | // There's a alarm write that hasn't been set yet pending for a time different than ours -- |
| 1017 | // We won't cancel the alarm because it hasn't been confirmed, but we shouldn't delete |
| 1018 | // the pending write. |
| 1019 | haveDeferredDelete = false; |
| 1020 | } |
| 1021 | } else { |
| 1022 | haveDeferredDelete = true; |
| 1023 | deferredAlarmSpan = kj::mv(parentSpan); |
| 1024 | } |
| 1025 | inAlarmHandler = true; |
| 1026 | |
| 1027 | static const DeferredAlarmDeleter disposer; |
| 1028 | return RunAlarmHandler{.deferredDelete = kj::Own<void>(this, disposer)}; |
| 1029 | } |
| 1030 | |
| 1031 | void ActorSqlite::cancelDeferredAlarmDeletion() { |
| 1032 | if (!inAlarmHandler) { |
| 1033 | // Pretty sure this can't happen. |
| 1034 | LOG_WARNING_ONCE("expected to be in alarm handler when trying to cancel deleted alarm"); |
| 1035 | } |
| 1036 | haveDeferredDelete = false; |
| 1037 | } |
| 1038 | |
| 1039 | kj::Promise<kj::Maybe<kj::Date>> ActorSqlite::abandonAlarm(kj::Date scheduledTime) { |
| 1040 | // Called when AlarmManager has given up retrying an alarm after too many counted failures. |
| 1041 | // Clear the alarm from SQLite so getAlarm() returns null instead of a stale time. |
| 1042 | // Only clear if SQLite currently has the exact alarm being abandoned and we're not mid-handler. |
| 1043 | // The time check guards against the race where the user set a new alarm (which always has a |
| 1044 | // time >= now() > scheduledTime due to past-time clamping in setAlarm) before this call arrived. |
| 1045 | if (inAlarmHandler) { |
| 1046 | // Shouldn't happen -- AlarmManager shouldn't call abandonAlarm while a handler is running. |
| 1047 | LOG_WARNING_ONCE("abandonAlarm called while alarm handler is still running"); |
| 1048 | return kj::Maybe<kj::Date>(kj::none); |
| 1049 | } |
| 1050 | KJ_IF_SOME(storedTime, metadata.getAlarm()) { |
| 1051 | if (storedTime == scheduledTime) { |
| 1052 | setAlarm(kj::none, {}, nullptr); |
| 1053 | return kj::Maybe<kj::Date>(kj::none); |
| 1054 | } else { |
| 1055 | // The user set a different alarm. Return it so AlarmManager can re-register. |
| 1056 | return kj::Maybe<kj::Date>(storedTime); |
| 1057 | } |
| 1058 | } |
| 1059 | return kj::Maybe<kj::Date>(kj::none); |
| 1060 | } |
| 1061 | |
| 1062 | kj::Maybe<kj::Promise<void>> ActorSqlite::onNoPendingFlush(SpanParent parentSpan) { |
| 1063 | // This implements sync(). |
| 1064 | // |
| 1065 | // sync() should wait for ALL writes (both confirmed and unconfirmed) that are outstanding at the |
| 1066 | // time sync() is called. We use lastCommit which keeps track of the most recent commit to be |
| 1067 | // formed. We join with the outputGate because there are a lot of edge cases where we break the |
| 1068 | // output gate and it's easiest to catch all of those instances here rather than updating |
| 1069 | // everything to also break lastCommit. |
| 1070 | return kj::joinPromisesFailFast( |
| 1071 | kj::arr(lastCommit.addBranch(), outputGate.wait(kj::mv(parentSpan)))); |
| 1072 | } |
| 1073 | |
| 1074 | kj::Promise<kj::String> ActorSqlite::getCurrentBookmark(SpanParent parentSpan) { |
| 1075 | // This is an ersatz implementation that's good enough for local dev with D1's Session API. |
| 1076 | // |
| 1077 | // The returned bookmark satisfies the properties that D1 cares about: |
| 1078 | // |
| 1079 | // * Later bookmarks sort after earlier bookmarks. We implement this by incrementing the bookmark |
| 1080 | // * whenever getCurrentBookmark() is called. |
| 1081 | // |
| 1082 | // * Bookmarks from the current workerd session sort after bookmarks from previous sessions. We |
| 1083 | // implement this by saving an ersatz bookmark in the SqliteMetadata table. |
| 1084 | |
| 1085 | requireNotBroken(); |
| 1086 | uint64_t bookmark = 0; |
| 1087 | KJ_IF_SOME(b, metadata.getLocalDevelopmentBookmark()) { |
| 1088 | bookmark = b + 1; |
| 1089 | } |
| 1090 | metadata.setLocalDevelopmentBookmark(bookmark); |
| 1091 | |
| 1092 | // TODO(cleanup): Left-padded number stringification should maybe be in KJ? |
| 1093 | auto paddedHex = [](uint32_t n) { |
| 1094 | kj::FixedArray<char, 8> result; |
| 1095 | for (auto i = 0; i < result.size(); i++) { |
| 1096 | char digit = n % 16; |
| 1097 | n /= 16; |
| 1098 | digit += digit < 10 ? '0' : ('a' - 10); |
| 1099 | result[result.size() - 1 - i] = digit; |
| 1100 | } |
| 1101 | return result; |
| 1102 | }; |
| 1103 | |
| 1104 | // Turn the bookmark into a format matching what Cloudflare's production returns. |
| 1105 | constexpr uint32_t uint32_max = kj::maxValue; |
| 1106 | kj::FixedArray<char, 32> pad; |
| 1107 | pad.fill('0'); |
| 1108 | return kj::str(paddedHex(bookmark / uint32_max), '-', paddedHex(bookmark % uint32_max), '-', |
| 1109 | paddedHex(0), '-', pad); |
| 1110 | } |
| 1111 | |
| 1112 | kj::Promise<void> ActorSqlite::waitForBookmark(kj::StringPtr bookmark, SpanParent parentSpan) { |
| 1113 | // This is an ersatz implementation that's good enough for local dev with D1's Session API. |
| 1114 | requireNotBroken(); |
| 1115 | return kj::READY_NOW; |
| 1116 | } |
| 1117 | |
| 1118 | void ActorSqlite::TxnCommitRegulator::onError( |
| 1119 | kj::Maybe<int> sqliteErrorCode, kj::StringPtr message) const { |
| 1120 | KJ_IF_SOME(c, sqliteErrorCode) { |
| 1121 | // We cannot `#include <sqlite3.h>` in the same compilation unit as `#include |
| 1122 | // <workerd/io/trace.h>` because the latter includes v8 and v8 seems to conflict with sqlite. |
| 1123 | // So we copy the value of SQLITE_CONSTRAINT from sqlite3.h |
| 1124 | constexpr int SQLITE_CONSTRAINT = 19; |
| 1125 | if (c == SQLITE_CONSTRAINT) { |
| 1126 | JSG_ASSERT(false, Error, |
| 1127 | "Durable Object was reset and rolled back to its last known good state because the " |
| 1128 | "application left the database in a state where constraints were violated: ", |
| 1129 | message); |
| 1130 | } |
| 1131 | } |
| 1132 | |
| 1133 | // For any other type of error, fall back to the default behavior (throwing a non-JSG exception) |
| 1134 | // as we don't know for sure that the problem is the application's fault. |
| 1135 | } |
| 1136 | |
| 1137 | const ActorSqlite::Hooks ActorSqlite::Hooks::DEFAULT = ActorSqlite::Hooks{}; |
| 1138 | |
| 1139 | kj::Promise<void> ActorSqlite::Hooks::scheduleRun( |
| 1140 | kj::Maybe<kj::Date> newAlarmTime, kj::Promise<void> priorTask) { |
| 1141 | JSG_FAIL_REQUIRE(Error, "alarms are not yet implemented for SQLite-backed Durable Objects"); |
| 1142 | } |
| 1143 | |
| 1144 | kj::OneOf<kj::Maybe<ActorCacheOps::Value>, kj::Promise<kj::Maybe<ActorCacheOps::Value>>> |
| 1145 | ActorSqlite::ExplicitTxn::get(Key key, ReadOptions options) { |
| 1146 | return actorSqlite.get(kj::mv(key), options); |
| 1147 | } |
| 1148 | kj::OneOf<ActorCacheOps::GetResultList, kj::Promise<ActorCacheOps::GetResultList>> ActorSqlite:: |
| 1149 | ExplicitTxn::get(kj::Array<Key> keys, ReadOptions options) { |
| 1150 | return actorSqlite.get(kj::mv(keys), options); |
| 1151 | } |
| 1152 | kj::OneOf<kj::Maybe<kj::Date>, kj::Promise<kj::Maybe<kj::Date>>> ActorSqlite::ExplicitTxn::getAlarm( |
| 1153 | ReadOptions options) { |
| 1154 | return actorSqlite.getAlarm(options); |
| 1155 | } |
| 1156 | kj::OneOf<ActorCacheOps::GetResultList, kj::Promise<ActorCacheOps::GetResultList>> ActorSqlite:: |
| 1157 | ExplicitTxn::list(Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) { |
| 1158 | return actorSqlite.list(kj::mv(begin), kj::mv(end), limit, options); |
| 1159 | } |
| 1160 | kj::OneOf<ActorCacheOps::GetResultList, kj::Promise<ActorCacheOps::GetResultList>> ActorSqlite:: |
| 1161 | ExplicitTxn::listReverse( |
| 1162 | Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) { |
| 1163 | return actorSqlite.listReverse(kj::mv(begin), kj::mv(end), limit, options); |
| 1164 | } |
| 1165 | kj::Maybe<kj::Promise<void>> ActorSqlite::ExplicitTxn::put( |
| 1166 | Key key, Value value, WriteOptions options, SpanParent traceSpan) { |
| 1167 | return actorSqlite.put(kj::mv(key), kj::mv(value), options, kj::mv(traceSpan)); |
| 1168 | } |
| 1169 | kj::Maybe<kj::Promise<void>> ActorSqlite::ExplicitTxn::put( |
| 1170 | kj::Array<KeyValuePair> pairs, WriteOptions options, SpanParent traceSpan) { |
| 1171 | return actorSqlite.put(kj::mv(pairs), options, kj::mv(traceSpan)); |
| 1172 | } |
| 1173 | kj::OneOf<bool, kj::Promise<bool>> ActorSqlite::ExplicitTxn::delete_( |
| 1174 | Key key, WriteOptions options, SpanParent traceSpan) { |
| 1175 | return actorSqlite.delete_(kj::mv(key), options, kj::mv(traceSpan)); |
| 1176 | } |
| 1177 | kj::OneOf<uint, kj::Promise<uint>> ActorSqlite::ExplicitTxn::delete_( |
| 1178 | kj::Array<Key> keys, WriteOptions options, SpanParent traceSpan) { |
| 1179 | return actorSqlite.delete_(kj::mv(keys), options, kj::mv(traceSpan)); |
| 1180 | } |
| 1181 | kj::Maybe<kj::Promise<void>> ActorSqlite::ExplicitTxn::setAlarm( |
| 1182 | kj::Maybe<kj::Date> newAlarmTime, WriteOptions options, SpanParent traceSpan) { |
| 1183 | return actorSqlite.setAlarm(newAlarmTime, options, kj::mv(traceSpan)); |
| 1184 | } |
| 1185 | |
| 1186 | } // namespace workerd |