// Copyright (c) 2024 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #include "actor-sqlite.h" #include "io-gate.h" #include #include #include #include #include namespace workerd { namespace { static constexpr kj::Date oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH; static constexpr kj::Date twoMs = 2 * kj::MILLISECONDS + kj::UNIX_EPOCH; static constexpr kj::Date threeMs = 3 * kj::MILLISECONDS + kj::UNIX_EPOCH; static constexpr kj::Date fourMs = 4 * kj::MILLISECONDS + kj::UNIX_EPOCH; static constexpr kj::Date fiveMs = 5 * kj::MILLISECONDS + kj::UNIX_EPOCH; static constexpr kj::Date sixMs = 6 * kj::MILLISECONDS + kj::UNIX_EPOCH; static constexpr kj::Date tenMs = 10 * kj::MILLISECONDS + kj::UNIX_EPOCH; // Used as the "current time" parameter for armAlarmHandler in tests. // Set to epoch (before all test alarm times) so existing tests aren't affected by // the overdue alarm check. static constexpr kj::Date testCurrentTime = kj::UNIX_EPOCH; template kj::Promise eagerlyReportExceptions(kj::Promise promise, kj::SourceLocation location = {}) { return promise.eagerlyEvaluate([location](kj::Exception&& e) -> T { KJ_LOG_AT(ERROR, location, e); kj::throwFatalException(kj::mv(e)); }); } // Expect that a synchronous result is returned. template T expectSync(kj::OneOf> result, kj::SourceLocation location = {}) { KJ_SWITCH_ONEOF(result) { KJ_CASE_ONEOF(promise, kj::Promise) { KJ_FAIL_ASSERT_AT(location, "result was unexpectedly asynchronous"); } KJ_CASE_ONEOF(value, T) { return kj::mv(value); } } KJ_UNREACHABLE; } struct ActorSqliteTestOptions final { bool monitorOutputGate = true; }; struct ActorSqliteTest final { kj::EventLoop loop; kj::WaitScope ws; OutputGate gate; kj::Own vfsDir; SqliteDatabase::Vfs vfs; SqliteDatabase db; struct Call final { kj::String desc; kj::Own> fulfiller; }; kj::Vector calls; struct ActorSqliteTestHooks final: public ActorSqlite::Hooks { public: explicit ActorSqliteTestHooks(ActorSqliteTest& parent): parent(parent) {} kj::Promise scheduleRun( kj::Maybe newAlarmTime, kj::Promise priorTask) override { KJ_IF_SOME(h, parent.scheduleRunWithPriorHandler) { return h(newAlarmTime, kj::mv(priorTask)); } KJ_IF_SOME(h, parent.scheduleRunHandler) { return h(newAlarmTime); } auto desc = newAlarmTime.map([](auto& t) { return kj::str("scheduleRun(", t, ")"); }).orDefault(kj::str("scheduleRun(none)")); auto [promise, fulfiller] = kj::newPromiseAndFulfiller(); parent.calls.add(Call{kj::mv(desc), kj::mv(fulfiller)}); return kj::mv(promise); } ActorSqliteTest& parent; }; kj::Maybe(kj::Maybe)>> scheduleRunHandler; kj::Maybe(kj::Maybe, kj::Promise)>> scheduleRunWithPriorHandler; ActorSqliteTestHooks hooks = ActorSqliteTestHooks(*this); ActorSqlite actor; kj::Promise gateBrokenPromise; kj::UnwindDetector unwindDetector; explicit ActorSqliteTest(ActorSqliteTestOptions options = {}) : ws(loop), vfsDir(kj::newInMemoryDirectory(kj::nullClock())), vfs(*vfsDir), db(vfs, kj::Path({"foo"}), kj::WriteMode::CREATE | kj::WriteMode::MODIFY), actor(kj::attachRef(db), gate, KJ_BIND_METHOD(*this, commitCallback), hooks), gateBrokenPromise(options.monitorOutputGate ? eagerlyReportExceptions(gate.onBroken()) : kj::Promise(kj::READY_NOW)) { db.afterReset([this](SqliteDatabase& cbDb) { KJ_ASSERT(&db == &cbDb); KJ_ASSERT(actor.isCommitScheduled(), "actor should have a commit scheduled during afterReset() callback"); }); } ~ActorSqliteTest() noexcept(false) { if (!unwindDetector.isUnwinding()) { // Make sure if the output gate has been broken, the exception was reported. This is // important to report errors thrown inside flush(), since those won't otherwise propagate // into the test body. gateBrokenPromise.poll(ws); // Make sure there's no outstanding async work we haven't considered: pollAndExpectCalls({}, "unexpected calls at end of test"); } } kj::Promise commitCallback(SpanParent) { auto [promise, fulfiller] = kj::newPromiseAndFulfiller(); calls.add(Call{kj::str("commit"), kj::mv(fulfiller)}); return kj::mv(promise); } // Polls the event loop, then asserts that the description of calls up to this point match the // expectation and returns their fulfillers. Also clears the call log. // // TODO(cleanup): Is there a better way to do mocks? capnp-mock looks nice, but seems a bit // heavyweight for this test. kj::Vector>> pollAndExpectCalls( std::initializer_list expCallDescs, kj::StringPtr message = ""_kj, kj::SourceLocation location = {}) { ws.poll(); auto callDescs = KJ_MAP(c, calls) { return kj::str(c.desc); }; KJ_ASSERT_AT(callDescs == heapArray(expCallDescs), location, kj::str(message)); auto fulfillers = KJ_MAP(c, calls) { return kj::mv(c.fulfiller); }; calls.clear(); return kj::mv(fulfillers); } // A few driver methods for convenience. auto get(kj::StringPtr key, ActorCache::ReadOptions options = {}) { return actor.get(kj::str(key), options); } auto getAlarm(ActorCache::ReadOptions options = {}) { return actor.getAlarm(options); } auto put(kj::StringPtr key, kj::StringPtr value, ActorCache::WriteOptions options = {}) { return actor.put(kj::str(key), kj::heapArray(value.asBytes()), options, nullptr); } auto putMultiple( kj::Array pairs, ActorCache::WriteOptions options = {}) { return actor.put(kj::mv(pairs), options, nullptr); } auto putMultipleExplicitTxn( kj::Array pairs, ActorCache::WriteOptions options = {}) { auto txn = actor.startTransaction(); txn->put(kj::mv(pairs), options, nullptr); return txn->commit(); } auto deleteMultiple(kj::Array keys, ActorCache::WriteOptions options = {}) { return actor.delete_(kj::mv(keys), options, nullptr); } auto setAlarm(kj::Maybe newTime, ActorCache::WriteOptions options = {}) { return actor.setAlarm(newTime, options, nullptr); } auto sync() { return actor.onNoPendingFlush(nullptr); } }; KJ_TEST("initial alarm value is unset") { ActorSqliteTest test; KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } KJ_TEST("can set and get alarm") { ActorSqliteTest test; test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_TEST("check put multiple wraps operations in a transaction") { ActorSqliteTest test; kj::Vector putKVs; putKVs.add(ActorCache::KeyValuePair{kj::str("foo"), kj::heapArray(kj::str("bar").asBytes())}); // NoTxn test { // Check that we're in a NoTxn KJ_ASSERT(!test.actor.isCommitScheduled()); test.putMultiple(putKVs.releaseAsArray()); // During write, all NoTxn operations are wrapped in an ImplicitTxn. auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); commitFulfiller->fulfill(); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); } // ExplicitTxn test { putKVs.add(ActorCache::KeyValuePair{kj::str("foo2"), kj::heapArray(kj::str("bar2").asBytes())}); KJ_ASSERT(!test.actor.isCommitScheduled()); // Similar to the previous putMultiple, but wrapped in a transactionSync (ExplicitTxn) test.putMultipleExplicitTxn(putKVs.releaseAsArray()); auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); commitFulfiller->fulfill(); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo2"))) == kj::str("bar2").asBytes()); } // ImplicitTxn test { // A single put will create an ImplicitTxn that we can use to wrap our putMultiple into. KJ_ASSERT(!test.actor.isCommitScheduled()); test.put("baz", "bat"); // By now, we should check there's a commit scheduled in a ImplicitTxn. KJ_ASSERT(test.actor.isCommitScheduled()); putKVs.add(ActorCache::KeyValuePair{kj::str("foo3"), kj::heapArray(kj::str("bar3").asBytes())}); test.putMultiple(putKVs.releaseAsArray()); auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("bat").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo3"))) == kj::str("bar3").asBytes()); commitFulfiller->fulfill(); } } KJ_TEST("check put multiple wraps operations in a transaction") { ActorSqliteTest test; kj::Vector putKVs; putKVs.add(ActorCache::KeyValuePair{kj::str("foo"), kj::heapArray(kj::str("bar").asBytes())}); // NoTxn test { // Check that we're in a NoTxn KJ_ASSERT(!test.actor.isCommitScheduled()); test.putMultiple(putKVs.releaseAsArray()); // During write, all NoTxn operations are wrapped in an ImplicitTxn. auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); commitFulfiller->fulfill(); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); } // ExplicitTxn test { putKVs.add(ActorCache::KeyValuePair{kj::str("foo2"), kj::heapArray(kj::str("bar2").asBytes())}); KJ_ASSERT(!test.actor.isCommitScheduled()); // Similar to the previous putMultiple, but wrapped in a transactionSync (ExplicitTxn) test.putMultipleExplicitTxn(putKVs.releaseAsArray()); auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); commitFulfiller->fulfill(); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo2"))) == kj::str("bar2").asBytes()); } // ImplicitTxn test { // A single put will create an ImplicitTxn that we can use to wrap our putMultiple into. KJ_ASSERT(!test.actor.isCommitScheduled()); test.put("baz", "bat"); // By now, we should check there's a commit scheduled in a ImplicitTxn. KJ_ASSERT(test.actor.isCommitScheduled()); putKVs.add(ActorCache::KeyValuePair{kj::str("foo3"), kj::heapArray(kj::str("bar3").asBytes())}); test.putMultiple(putKVs.releaseAsArray()); auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("bat").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo3"))) == kj::str("bar3").asBytes()); commitFulfiller->fulfill(); } } KJ_TEST("check put multiple wraps operations in a transaction and rollback on error") { ActorSqliteTest test; // We expect that putMultiple is all or nothing, rolling back if a single put fails. kj::Vector putKVs; // Add some regular key-value pairs that we know are supported putKVs.add(ActorCache::KeyValuePair{kj::str("foo"), kj::heapArray(kj::str("bar").asBytes())}); putKVs.add(ActorCache::KeyValuePair{kj::str("foo2"), kj::heapArray(kj::str("bar2").asBytes())}); putKVs.add(ActorCache::KeyValuePair{kj::str("foo3"), kj::heapArray(kj::str("bar3").asBytes())}); // Now create a key that's too large. Should fail with string or blob too big: SQLITE_TOOBIG auto tooLongKey = kj::heapString(2200000); tooLongKey.asArray().fill('a'); // Add it to our KV array putKVs.add( ActorCache::KeyValuePair{kj::str(tooLongKey), kj::heapArray(kj::str("bar").asBytes())}); // NoTxn test { // Check that we're in a NoTxn KJ_ASSERT(!test.actor.isCommitScheduled()); try { test.putMultiple(putKVs.releaseAsArray()); // During write, all NoTxn operations are wrapped in an ImplicitTxn. auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); commitFulfiller->fulfill(); KJ_UNREACHABLE; } catch (kj::Exception& e) { KJ_ASSERT( e.getDescription() == "expected false; jsg.Error: string or blob too big: SQLITE_TOOBIG"); } KJ_ASSERT(expectSync(test.get(kj::str("foo"))) == nullptr); KJ_ASSERT(expectSync(test.get(kj::str("foo2"))) == nullptr); KJ_ASSERT(expectSync(test.get(kj::str("foo3"))) == nullptr); } // Reset the transaction state by going async, which will cause the ImplicitTxn to commit. { auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); commitFulfiller->fulfill(); } // ExplicitTxn test { KJ_ASSERT(!test.actor.isCommitScheduled()); // Similar to the previous putMultiple, but wrapped in a transactionSync (ExplicitTxn) test.putMultipleExplicitTxn(putKVs.releaseAsArray()); auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); commitFulfiller->fulfill(); KJ_ASSERT(expectSync(test.get(kj::str("foo"))) == nullptr); KJ_ASSERT(expectSync(test.get(kj::str("foo2"))) == nullptr); KJ_ASSERT(expectSync(test.get(kj::str("foo3"))) == nullptr); } // ImplicitTxn test { // A single put will create an ImplicitTxn that we can use to wrap our putMultiple into. KJ_ASSERT(!test.actor.isCommitScheduled()); test.put("baz", "bat"); // By now, we should check there's a commit scheduled in a ImplicitTxn. KJ_ASSERT(test.actor.isCommitScheduled()); test.putMultiple(putKVs.releaseAsArray()); auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); // The single put succeeded, but the putMultiple did not. KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("bat").asBytes()); KJ_ASSERT(expectSync(test.get(kj::str("foo"))) == nullptr); KJ_ASSERT(expectSync(test.get(kj::str("foo2"))) == nullptr); KJ_ASSERT(expectSync(test.get(kj::str("foo3"))) == nullptr); commitFulfiller->fulfill(); } } KJ_TEST("alarm write happens transactionally with storage ops") { ActorSqliteTest test; test.setAlarm(oneMs); test.put("foo", "bar"); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); } KJ_TEST("storage op without alarm change does not wait on scheduler") { ActorSqliteTest test; test.put("foo", "bar"); test.pollAndExpectCalls({"commit"})[0]->fulfill(); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } KJ_TEST("alarm scheduling starts synchronously before implicit local db commit") { ActorSqliteTest test; // In workerd (unlike edgeworker), there is no remote storage, so there is no work done in // commitCallback(); the local db is considered durably stored after the synchronous sqlite // commit() call returns. If a commit includes an alarm state change that requires scheduling // before the commit call, it needs to happen synchronously. Since workerd synchronously // schedules alarms, we just need to ensure that the database is in a pre-commit state when // scheduleRun() is called. // Initialize alarm state to 2ms. test.setAlarm(twoMs); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); bool startedScheduleRun = false; test.scheduleRunHandler = [&](kj::Maybe) -> kj::Promise { startedScheduleRun = true; KJ_EXPECT_THROW_MESSAGE( "cannot start a transaction within a transaction", test.db.run("BEGIN TRANSACTION")); return kj::READY_NOW; }; test.setAlarm(oneMs); KJ_ASSERT(!startedScheduleRun); test.ws.poll(); KJ_ASSERT(startedScheduleRun); test.pollAndExpectCalls({"commit"})[0]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_TEST("alarm scheduling starts synchronously before explicit local db commit") { ActorSqliteTest test; // Initialize alarm state to 2ms. test.setAlarm(twoMs); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); bool startedScheduleRun = false; test.scheduleRunHandler = [&](kj::Maybe) -> kj::Promise { startedScheduleRun = true; // Not sure if there is a good way to detect savepoint presence without mutating the db state, // but this is sufficient to verify the test properties: // Verify that we are not within a nested savepoint. KJ_EXPECT_THROW_MESSAGE( "no such savepoint: _cf_savepoint_1", test.db.run("RELEASE _cf_savepoint_1")); // Verify that we are within the root savepoint. test.db.run("RELEASE _cf_savepoint_0"); KJ_EXPECT_THROW_MESSAGE( "no such savepoint: _cf_savepoint_0", test.db.run("RELEASE _cf_savepoint_0")); // We don't actually care what happens in the test after this point, but it's slightly simpler // to re-add the savepoint to allow the test to complete cleanly: test.db.run("SAVEPOINT _cf_savepoint_0"); return kj::READY_NOW; }; { auto txn = test.actor.startTransaction(); txn->setAlarm(oneMs, {}, nullptr); KJ_ASSERT(!startedScheduleRun); txn->commit(); KJ_ASSERT(startedScheduleRun); test.pollAndExpectCalls({"commit"})[0]->fulfill(); } KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_TEST("alarm scheduling does not start synchronously before nested explicit local db commit") { ActorSqliteTest test; // Initialize alarm state to 2ms. test.setAlarm(twoMs); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); bool startedScheduleRun = false; test.scheduleRunHandler = [&](kj::Maybe) -> kj::Promise { startedScheduleRun = true; return kj::READY_NOW; }; { auto txn1 = test.actor.startTransaction(); { auto txn2 = test.actor.startTransaction(); txn2->setAlarm(oneMs, {}, nullptr); txn2->commit(); KJ_ASSERT(!startedScheduleRun); } txn1->commit(); KJ_ASSERT(startedScheduleRun); test.pollAndExpectCalls({"commit"})[0]->fulfill(); } KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_TEST("synchronous alarm scheduling failure causes local db commit to throw synchronously") { ActorSqliteTest test({.monitorOutputGate = false}); auto promise = test.gate.onBroken(); auto getLocalAlarm = [&]() -> kj::Maybe { auto query = test.db.run("SELECT value FROM _cf_METADATA WHERE key = 1"); if (query.isDone() || query.isNull(0)) { return kj::none; } else { return kj::UNIX_EPOCH + query.getInt64(0) * kj::NANOSECONDS; } }; // Initialize alarm state to 2ms. test.setAlarm(twoMs); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); // Override scheduleRun handler with one that throws synchronously. bool startedScheduleRun = false; test.scheduleRunHandler = [&](kj::Maybe) -> kj::Promise { startedScheduleRun = true; // Must throw synchronously; returning an exception is insufficient. kj::throwFatalException(KJ_EXCEPTION(FAILED, "a_sync_fail")); }; KJ_ASSERT(!promise.poll(test.ws)); test.setAlarm(oneMs); // Expect that polling will attempt to commit the implicit transaction, which should // synchronously fail when attempting to call scheduleRun() before the db commit, and roll back the // local db state to the 2ms alarm. KJ_ASSERT(!startedScheduleRun); KJ_ASSERT(KJ_REQUIRE_NONNULL(getLocalAlarm()) == oneMs); test.ws.poll(); KJ_ASSERT(startedScheduleRun); KJ_ASSERT(KJ_REQUIRE_NONNULL(getLocalAlarm()) == twoMs); KJ_ASSERT(promise.poll(test.ws)); KJ_EXPECT_THROW_MESSAGE("a_sync_fail", promise.wait(test.ws)); } KJ_TEST("can clear alarm") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); test.setAlarm(kj::none); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } KJ_TEST("can set alarm twice") { ActorSqliteTest test; test.setAlarm(oneMs); test.setAlarm(twoMs); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); } KJ_TEST("setting duplicate alarm is no-op") { ActorSqliteTest test; test.setAlarm(kj::none); test.pollAndExpectCalls({}); test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.setAlarm(oneMs); test.pollAndExpectCalls({}); } KJ_TEST("tells alarm handler to cancel when committed alarm is empty") { ActorSqliteTest test; { auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); // We expect armAlarmHandler() to tell us to cancel the alarm. KJ_ASSERT(armResult.is()); auto waitPromise = kj::mv(armResult.get().waitBeforeCancel); // We also expect the alarm cancellation to contain a scheduling request to delete the alarm, // to handle cases where alarm deletion was durably committed to the database, but a failure // occurred before the alarm deletion was conveyed to the alarm scheduler. KJ_ASSERT(!waitPromise.poll(test.ws)); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); KJ_ASSERT(waitPromise.poll(test.ws)); waitPromise.wait(test.ws); } } KJ_TEST("tells alarm handler to reschedule when handler alarm is later than committed alarm") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Request handler run at 2ms. Expect cancellation with rescheduling. auto armResult = test.actor.armAlarmHandler(twoMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); auto cancelResult = kj::mv(armResult.get()); // Expect rescheduling was requested and that returned promise resolves after fulfillment. auto waitBeforeCancel = kj::mv(cancelResult.waitBeforeCancel); auto rescheduleFulfiller = kj::mv(test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]); KJ_ASSERT(!waitBeforeCancel.poll(test.ws)); rescheduleFulfiller->fulfill(); KJ_ASSERT(waitBeforeCancel.poll(test.ws)); waitBeforeCancel.wait(test.ws); } KJ_TEST("tells alarm handler to reschedule when handler alarm is earlier than committed alarm") { ActorSqliteTest test; // Initialize alarm state to 2ms. test.setAlarm(twoMs); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); // Expect that armAlarmHandler() tells caller to cancel after rescheduling completes. auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); auto cancelResult = kj::mv(armResult.get()); // Expect rescheduling was requested and that returned promise resolves after fulfillment. auto waitBeforeCancel = kj::mv(cancelResult.waitBeforeCancel); auto rescheduleFulfiller = kj::mv(test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]); KJ_ASSERT(!waitBeforeCancel.poll(test.ws)); rescheduleFulfiller->fulfill(); KJ_ASSERT(waitBeforeCancel.poll(test.ws)); waitBeforeCancel.wait(test.ws); } KJ_TEST("runs overdue alarm immediately when local alarm time is in the past") { ActorSqliteTest test; // Initialize alarm state to 2ms. test.setAlarm(twoMs); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); // The local state says the alarm is due to fire at 2ms, but we're saying the AlarmManager has 1ms, // usually this would result in a rescheduling of the alarm, but since our currentTime is 5ms, we // will just run the alarm now since it's already overdue. { auto overdueCurrentTime = fiveMs; auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, overdueCurrentTime); // Should run the handler immediately instead of canceling/rescheduling. KJ_ASSERT(armResult.is()); } // commit and delete the alarm after we drop the alarm handler (this is a deferred delete). test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); } KJ_TEST("does not cancel handler when local db alarm state is later than scheduled alarm") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); test.setAlarm(twoMs); { auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); } test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); } KJ_TEST("does not cancel handler when local db alarm state is earlier than scheduled alarm") { ActorSqliteTest test; // Initialize alarm state to 2ms. test.setAlarm(twoMs); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); test.setAlarm(oneMs); { auto armResult = test.actor.armAlarmHandler(twoMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); } test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); } KJ_TEST("getAlarm() returns null during handler") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); { auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); } KJ_TEST("alarm handler handle clears alarm when dropped with no writes") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); { auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); } test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } KJ_TEST("alarm deleter does not clear alarm when dropped with writes") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); { auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); test.setAlarm(twoMs); } test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); } KJ_TEST("can cancel deferred alarm deletion during handler") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); { auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); test.actor.cancelDeferredAlarmDeletion(); } KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_TEST("canceling deferred alarm deletion outside handler has no effect") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); { auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); } test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); test.actor.cancelDeferredAlarmDeletion(); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } KJ_TEST("canceling deferred alarm deletion outside handler edge case") { // Presumably harmless to cancel deletion if the client requests it after the handler ends but // before the event loop runs the commit code? Trying to cancel deletion outside the handler is // a bit of a contract violation anyway -- maybe we should just assert against it? ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); { auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); } test.actor.cancelDeferredAlarmDeletion(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } KJ_TEST("canceling deferred alarm deletion is idempotent") { // Not sure if important, but matches ActorCache behavior. ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); { auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); test.actor.cancelDeferredAlarmDeletion(); test.actor.cancelDeferredAlarmDeletion(); } KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_TEST("alarm handler cleanup succeeds when output gate is broken") { auto runWithSetup = [](auto testFunc) { ActorSqliteTest test({.monitorOutputGate = false}); auto promise = test.gate.onBroken(); // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); auto deferredDelete = kj::mv(armResult.get().deferredDelete); // Break gate test.put("foo", "bar"); test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "a_rejected_commit")); KJ_EXPECT_THROW_MESSAGE("a_rejected_commit", promise.wait(test.ws)); // Ensure taskFailed handler runs and notices brokenness: test.ws.poll(); testFunc(test, kj::mv(deferredDelete)); }; // Here, we test that the deferred deleter destructor doesn't throw, both in the case when the // caller cancels deletion and when it does not cancel it: runWithSetup([](ActorSqliteTest& test, kj::Own deferredDelete) { // In the case where the handler fails, we assume the caller will explicitly cancel deferred // alarm deletion: test.actor.cancelDeferredAlarmDeletion(); // Dropping the DeferredAlarmDeleter should succeed -- that is, it should not throw here: { auto drop = kj::mv(deferredDelete); } }); runWithSetup([](ActorSqliteTest& test, kj::Own deferredDelete) { // In the case where the handler succeeds, the caller will not cancel deferred deletion before // dropping the DeferredAlarmDeleter. Dropping the DeferredAlarmDeleter should still succeed, // even if the output gate happens to already be broken: { auto drop = kj::mv(deferredDelete); } }); } KJ_TEST("handler alarm is not deleted when commit fails") { ActorSqliteTest test({.monitorOutputGate = false}); auto promise = test.gate.onBroken(); // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); { auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "a_rejected_commit")); KJ_EXPECT_THROW_MESSAGE("a_rejected_commit", promise.wait(test.ws)); } KJ_TEST("setting earlier alarm persists alarm scheduling before db") { ActorSqliteTest test; // Initialize alarm state to 2ms. test.setAlarm(twoMs); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); // Update alarm to be earlier. We expect the alarm scheduling to be persisted before the db. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_TEST("setting later alarm persists db before alarm scheduling") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Update alarm to be later. We expect the db to be persisted before the alarm scheduling. test.setAlarm(twoMs); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); } KJ_TEST("multiple set-earlier in-flight alarms wait for earliest before committing db") { ActorSqliteTest test; // Initialize alarm state to 5ms. test.setAlarm(fiveMs); test.pollAndExpectCalls({"scheduleRun(5ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == fiveMs); // Gate is not blocked. auto gateWaitBefore = test.gate.wait(nullptr); KJ_ASSERT(gateWaitBefore.poll(test.ws)); // Update alarm to be earlier (4ms). We expect the alarm scheduling to start. test.setAlarm(fourMs); auto fulfiller4Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(4ms)"})[0]); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == fourMs); // Gate as-of 4ms update is blocked. auto gateWait4ms = test.gate.wait(nullptr); KJ_ASSERT(!gateWait4ms.poll(test.ws)); // While 4ms scheduling request is in-flight, update alarm to be even earlier (3ms). We expect // the 4ms request to block the 3ms scheduling request. test.setAlarm(threeMs); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == threeMs); // Gate as-of 3ms update is blocked. auto gateWait3ms = test.gate.wait(nullptr); KJ_ASSERT(!gateWait3ms.poll(test.ws)); // Update alarm to be even earlier (2ms). We expect scheduling requests to still be blocked. test.setAlarm(twoMs); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); // Gate as-of 2ms update is blocked. auto gateWait2ms = test.gate.wait(nullptr); KJ_ASSERT(!gateWait2ms.poll(test.ws)); // Fulfill the 4ms request. We expect the 2ms scheduling to start, because that is the current // alarm value. fulfiller4Ms->fulfill(); auto fulfiller2Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]); test.pollAndExpectCalls({}); // While waiting for 2ms request, update alarm time to be 1ms. Expect scheduling to be blocked. test.setAlarm(oneMs); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Gate as-of 1ms update is blocked. auto gateWait1ms = test.gate.wait(nullptr); KJ_ASSERT(!gateWait1ms.poll(test.ws)); // Fulfill the 2ms request. We expect the 1ms scheduling to start. fulfiller2Ms->fulfill(); auto fulfiller1Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]); test.pollAndExpectCalls({}); // Fulfill the 1ms request. We expect a single db commit to start (coalescing all previous db // commits together). fulfiller1Ms->fulfill(); auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); test.pollAndExpectCalls({}); // We expect all earlier gates to be blocked until commit completes. KJ_ASSERT(!gateWait4ms.poll(test.ws)); KJ_ASSERT(!gateWait3ms.poll(test.ws)); KJ_ASSERT(!gateWait2ms.poll(test.ws)); KJ_ASSERT(!gateWait1ms.poll(test.ws)); commitFulfiller->fulfill(); KJ_ASSERT(gateWait4ms.poll(test.ws)); KJ_ASSERT(gateWait3ms.poll(test.ws)); KJ_ASSERT(gateWait2ms.poll(test.ws)); KJ_ASSERT(gateWait1ms.poll(test.ws)); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_TEST("setting later alarm times does scheduling after db commit") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Gate is not blocked. auto gateWaitBefore = test.gate.wait(nullptr); KJ_ASSERT(gateWaitBefore.poll(test.ws)); // Set alarm to 2ms. Expect 2ms db commit to start. test.setAlarm(twoMs); auto commit2MsFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); test.pollAndExpectCalls({}); // Gate as-of 2ms update is blocked. auto gateWait2Ms = test.gate.wait(nullptr); KJ_ASSERT(!gateWait2Ms.poll(test.ws)); // Set alarm to 3ms. Expect 3ms db commit to start. The 2ms scheduleRun will never happen now // that we've overwritten it while it was persisting to SQLite (before we send an update to the // alarm manager). test.setAlarm(threeMs); auto commit3MsFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); test.pollAndExpectCalls({}); // Gate as-of 3ms update is blocked. auto gateWait3Ms = test.gate.wait(nullptr); KJ_ASSERT(!gateWait3Ms.poll(test.ws)); // Expect 2ms gate to be unblocked once the commit finishes, but don't expect a scheduleRun(2ms). KJ_ASSERT(!gateWait2Ms.poll(test.ws)); commit2MsFulfiller->fulfill(); KJ_ASSERT(gateWait2Ms.poll(test.ws)); test.pollAndExpectCalls({}); // Fulfill 3ms db commit. Expect 3ms alarm to be scheduled and 3ms gate to be unblocked. KJ_ASSERT(!gateWait3Ms.poll(test.ws)); commit3MsFulfiller->fulfill(); KJ_ASSERT(gateWait3Ms.poll(test.ws)); auto fulfiller3Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(3ms)"})[0]); test.pollAndExpectCalls({}); fulfiller3Ms->fulfill(); } KJ_TEST("rejected move-earlier alarm scheduling request breaks gate") { ActorSqliteTest test({.monitorOutputGate = false}); auto promise = test.gate.onBroken(); test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->reject( KJ_EXCEPTION(FAILED, "a_rejected_scheduleRun")); KJ_EXPECT_THROW_MESSAGE("a_rejected_scheduleRun", promise.wait(test.ws)); } KJ_TEST("rejected move-later alarm scheduling request does not break gate") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Update alarm to be later. We expect the db to be persisted before the alarm scheduling. // We simulate a failure during the alarm rescheduling, but expect it to not break the output // gate. test.setAlarm(twoMs); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->reject( KJ_EXCEPTION(FAILED, "a_rejected_scheduleRun")); // Subsequent kv put succeeds. In an earlier version of the code, this failed, due to capturing // the scheduling failure as if it had broke the output gate, without actually breaking the // output gate. test.put("foo", "bar"); test.pollAndExpectCalls({"commit"})[0]->fulfill(); } KJ_TEST("rapid move-later alarm changes coalesce into bounded scheduleRun calls") { // When many commits each move the alarm time later while a scheduleRun is already in-flight, // the scheduleLaterAlarm mechanism should coalesce them into at most one pending request, // rather than chaining N promises (one per commit). ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Move alarm to 2ms. The db commit completes, triggering a post-commit scheduleRun(2ms) // since the alarm moved later. test.setAlarm(twoMs); test.pollAndExpectCalls({"commit"})[0]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); // The first move-later scheduleRun starts. auto fulfiller2Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]); // While 2ms scheduleRun is in-flight, move alarm to 3ms, 4ms, 5ms in rapid succession. // Each commit completes immediately but the scheduleRun for 2ms is still pending. // Only the final value (5ms) should be scheduled after the 2ms scheduleRun completes. test.setAlarm(threeMs); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); // No new scheduleRun -- coalesced into pending. KJ_ASSERT(expectSync(test.getAlarm()) == threeMs); test.setAlarm(fourMs); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); // No new scheduleRun -- coalesced into pending. KJ_ASSERT(expectSync(test.getAlarm()) == fourMs); test.setAlarm(fiveMs); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); // No new scheduleRun -- coalesced into pending. KJ_ASSERT(expectSync(test.getAlarm()) == fiveMs); // Now fulfill the 2ms scheduleRun. The coalesced pending time (5ms) should be scheduled next. fulfiller2Ms->fulfill(); auto fulfiller5Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(5ms)"})[0]); // Importantly, there is exactly one scheduleRun(5ms), not three separate calls for 3ms, 4ms, 5ms. fulfiller5Ms->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == fiveMs); } KJ_TEST("armAlarmHandler with coalesced pending alarms schedules reschedule exactly once") { // Verifies two properties: // 1. No duplicate scheduleRun(6ms): armAlarmHandler clears pendingLaterAlarmTime so the // FORK_A completion handler does not re-issue it. // 2. Future commits (10ms) that arrive after armAlarmHandler fires are correctly handled: // they queue in pendingLaterAlarmTime, get picked up by FORK_A's completion handler, // and chain off FORK_B (armAlarmHandler's fork) so the order is 3ms -> 6ms -> 10ms. ActorSqliteTest test; // Initialize alarm to 1ms and fully commit it so lastConfirmedAlarmDbState = 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Move alarm to 3ms -- scheduleRun(3ms) goes in-flight via scheduleLaterAlarm. // alarmLaterIsInFlight=true, alarmLaterInFlight=FORK_A. test.setAlarm(threeMs); test.pollAndExpectCalls({"commit"})[0]->fulfill(); auto fulfiller3Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(3ms)"})[0]); // While 3ms scheduleRun is in-flight, rapidly move to 4ms then 6ms. // Both coalesce into pendingLaterAlarmTime=6ms; no new scheduleRun issued. test.setAlarm(fourMs); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); test.setAlarm(sixMs); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == sixMs); // The 1ms alarm fires. armAlarmHandler sees scheduledTime=1ms, localAlarmState=6ms. // willFireEarlier(1ms, 6ms) => reschedule-later path: // requestScheduledAlarm(6ms, FORK_A.addBranch()) called synchronously -> FORK_B // pendingLaterAlarmTime cleared to kj::none // alarmLaterInFlight = FORK_B // alarmLaterIsInFlight unchanged (still true, owned by FORK_A lifecycle) auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); auto& cancelResult = armResult.get(); // scheduleRun(6ms) issued exactly once -- synchronously inside armAlarmHandler. auto fulfiller6Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(6ms)"})[0]); // Commit for 10ms arrives while scheduleRun(3ms) is still in-flight. // alarmLaterIsInFlight=true (FORK_A lifecycle still active) so 10ms is correctly // queued: pendingLaterAlarmTime=Some(10ms). FORK_B is referenced by alarmLaterInFlight. test.setAlarm(tenMs); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); // No scheduleRun yet -- coalesced into pending. // Fulfill scheduleRun(3ms). FORK_A resolves. FORK_A completion handler fires: // alarmLaterIsInFlight=false // pendingLaterAlarmTime=Some(10ms) -> scheduleLaterAlarm(10ms) // requestScheduledAlarm(10ms, FORK_B.addBranch()) -> FORK_C chains off FORK_B // scheduleRun(10ms) issued synchronously // Importantly: scheduleRun(6ms) is NOT issued again here -- pendingLaterAlarmTime // held 10ms (not 6ms), because armAlarmHandler had already cleared the 6ms. fulfiller3Ms->fulfill(); auto fulfiller10Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(10ms)"})[0]); // Fulfill scheduleRun(6ms). FORK_B resolves cleanly -- it has no completion handler, // so overwriting alarmLaterInFlight with FORK_C is safe: FORK_C captured a branch of // FORK_B as priorTask before the field was overwritten, keeping FORK_B alive. FORK_B // resolving propagates into FORK_C's priorTask silently. No new scheduleRun here. fulfiller6Ms->fulfill(); KJ_ASSERT(cancelResult.waitBeforeCancel.poll(test.ws)); test.pollAndExpectCalls({}); // Fulfill scheduleRun(10ms). FORK_C resolves, its completion handler fires with no // pending times. Done. fulfiller10Ms->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == tenMs); } KJ_TEST("coalesced move-later followed by move-earlier does not race") { // Regression test for a race condition where a coalesced pendingLaterAlarmTime could // be drained concurrently with a move-earlier scheduleRun. The fix is that // startPrecommitAlarmScheduling() clears pendingLaterAlarmTime when setting up a // move-earlier, so the completion handler finds nothing to drain. // // Scenario: alarm at 1ms -> move to 5ms (later, in-flight) -> move to 10ms (coalesced) // -> move to 2ms (earlier). Without the fix, after the 5ms RPC completes, both // scheduleRun(10ms) and scheduleRun(2ms) would fire concurrently. With the fix, // only scheduleRun(2ms) fires because the coalesced 10ms was cleared. ActorSqliteTest test; uint activeRpcs = 0; uint maxConcurrentRpcs = 0; // Custom handler that respects priorTask ordering like the real alarm manager. // The real alarm manager awaits priorTask before sending its RPC; we replicate // that here and track concurrent calls. test.scheduleRunWithPriorHandler = [&](kj::Maybe newAlarmTime, kj::Promise priorTask) -> kj::Promise { return priorTask.then([&, newAlarmTime]() mutable -> kj::Promise { activeRpcs++; maxConcurrentRpcs = kj::max(maxConcurrentRpcs, activeRpcs); auto desc = newAlarmTime.map([](auto& t) { return kj::str("scheduleRun(", t, ")"); }).orDefault(kj::str("scheduleRun(none)")); auto [promise, fulfiller] = kj::newPromiseAndFulfiller(); test.calls.add(ActorSqliteTest::Call{kj::mv(desc), kj::mv(fulfiller)}); return promise.then([&]() { activeRpcs--; }); }); }; // Poll event loop until at least `count` calls accumulate. With the // priorTask-respecting handler, calls take extra event loop turns to appear. auto drainCalls = [&](std::initializer_list expected, kj::StringPtr msg = ""_kj) { size_t need = expected.size(); for (int i = 0; i < 100; i++) { test.ws.poll(); if (need == 0 && i >= 10) break; if (need > 0 && test.calls.size() >= need) break; } auto callDescs = KJ_MAP(c, test.calls) { return kj::str(c.desc); }; KJ_ASSERT(callDescs == kj::heapArray(expected), msg); auto fulfillers = KJ_MAP(c, test.calls) { return kj::mv(c.fulfiller); }; test.calls.clear(); return kj::mv(fulfillers); }; // 1. Initialize alarm state to 1ms. test.setAlarm(oneMs); drainCalls({"scheduleRun(1ms)"})[0]->fulfill(); drainCalls({"commit"})[0]->fulfill(); drainCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // 2. Move alarm to 5ms (later). The db commit completes, then scheduleRun(5ms) // fires post-commit via scheduleLaterAlarm. Hold the fulfiller to keep it in-flight. test.setAlarm(fiveMs); drainCalls({"commit"})[0]->fulfill(); auto fulfiller5Ms = kj::mv(drainCalls({"scheduleRun(5ms)"})[0]); // 3. While scheduleRun(5ms) is in-flight, move alarm to 10ms (later). // Since alarmLaterIsInFlight is true, 10ms is coalesced into pendingLaterAlarmTime. test.setAlarm(tenMs); drainCalls({"commit"})[0]->fulfill(); drainCalls({}); // No scheduleRun -- coalesced into pending. // 4. Move alarm earlier to 2ms. startPrecommitAlarmScheduling() clears // pendingLaterAlarmTime and calls requestScheduledAlarm(2ms, FORK_5.addBranch()). // The priorTask-respecting handler blocks until FORK_5 resolves. test.setAlarm(twoMs); drainCalls({}); // scheduleRun(2ms) blocked on priorTask. // 5. Fulfill scheduleRun(5ms). FORK_5 resolves: // - Completion handler: pendingLaterAlarmTime was cleared -> no drain, no-op. // - Move-earlier priorTask resolves -> scheduleRun(2ms) fires. fulfiller5Ms->fulfill(); auto fulfiller2Ms = kj::mv(drainCalls( {"scheduleRun(2ms)"}, "expected only scheduleRun(2ms), no concurrent scheduleRun(10ms)")[0]); // Verify no concurrent RPCs occurred. KJ_ASSERT(maxConcurrentRpcs <= 1, "scheduleRun RPCs were sent concurrently -- " "the coalesced move-later raced with the move-earlier"); // 6. Complete the move-earlier and its commit. fulfiller2Ms->fulfill(); drainCalls({"commit"})[0]->fulfill(); // Let any remaining completion handlers settle. for (int i = 0; i < 20; i++) test.ws.poll(); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); } KJ_TEST("an exception thrown during merged commits does not hang") { ActorSqliteTest test({.monitorOutputGate = false}); auto promise = test.gate.onBroken(); // Initialize alarm state to 5ms. test.setAlarm(fiveMs); test.pollAndExpectCalls({"scheduleRun(5ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == fiveMs); // Update alarm to be earlier (4ms). We expect the alarm scheduling to start. test.setAlarm(fourMs); auto fulfiller4Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(4ms)"})[0]); auto gateWait4ms = test.gate.wait(nullptr); // While 4ms scheduling request is in-flight, update alarm to be earlier (3ms). We expect // the two commit requests to merge and be blocked on the alarm scheduling request. test.setAlarm(threeMs); test.pollAndExpectCalls({}); auto gateWait3ms = test.gate.wait(nullptr); // Reject the 4ms request. We expect both gate waiting promises to unblock with exceptions. KJ_ASSERT(!gateWait4ms.poll(test.ws)); KJ_ASSERT(!gateWait3ms.poll(test.ws)); fulfiller4Ms->reject(KJ_EXCEPTION(FAILED, "a_rejected_scheduleRun")); KJ_ASSERT(gateWait4ms.poll(test.ws)); KJ_ASSERT(gateWait3ms.poll(test.ws)); KJ_EXPECT_THROW_MESSAGE("a_rejected_scheduleRun", gateWait4ms.wait(test.ws)); KJ_EXPECT_THROW_MESSAGE("a_rejected_scheduleRun", gateWait3ms.wait(test.ws)); KJ_EXPECT_THROW_MESSAGE("a_rejected_scheduleRun", promise.wait(test.ws)); } KJ_TEST("getAlarm/setAlarm check for brokenness") { ActorSqliteTest test({.monitorOutputGate = false}); auto promise = test.gate.onBroken(); // Break gate test.put("foo", "bar"); test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "a_rejected_commit")); KJ_EXPECT_THROW_MESSAGE("a_rejected_commit", promise.wait(test.ws)); // Apparently we don't actually set brokenness until the taskFailed handler runs, but presumably // this is OK? test.getAlarm(); // Ensure taskFailed handler runs and notices brokenness: test.ws.poll(); KJ_EXPECT_THROW_MESSAGE("a_rejected_commit", test.getAlarm()); KJ_EXPECT_THROW_MESSAGE("a_rejected_commit", test.setAlarm(kj::none)); test.pollAndExpectCalls({}); } KJ_TEST("calling deleteAll() preserves alarm state if alarm is set") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); { KJ_ASSERT(!test.actor.isCommitScheduled()); ActorCache::DeleteAllResults results = test.actor.deleteAll({}, nullptr); KJ_ASSERT(test.actor.isCommitScheduled()); KJ_ASSERT(results.backpressure == kj::none); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); KJ_ASSERT(results.count.wait(test.ws) == 0); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); commitFulfiller->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } { // Should be fine to call deleteAll() a few times in succession, too: KJ_ASSERT(!test.actor.isCommitScheduled()); ActorCache::DeleteAllResults results1 = test.actor.deleteAll({}, nullptr); ActorCache::DeleteAllResults results2 = test.actor.deleteAll({}, nullptr); KJ_ASSERT(test.actor.isCommitScheduled()); KJ_ASSERT(results1.backpressure == kj::none); KJ_ASSERT(results2.backpressure == kj::none); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Presumably fine to be performing the alarm state restoration after each db reset: auto commitFulfillers = test.pollAndExpectCalls({"commit", "commit"}); KJ_ASSERT(results1.count.wait(test.ws) == 0); KJ_ASSERT(results2.count.wait(test.ws) == 0); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); commitFulfillers[0]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); commitFulfillers[1]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } } KJ_TEST("calling deleteAll() preserves alarm state if alarm is not set") { ActorSqliteTest test; // Initialize alarm state to empty value in metadata table. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); test.setAlarm(kj::none); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); { KJ_ASSERT(!test.actor.isCommitScheduled()); ActorCache::DeleteAllResults results = test.actor.deleteAll({}, nullptr); KJ_ASSERT(test.actor.isCommitScheduled()); KJ_ASSERT(results.backpressure == kj::none); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); KJ_ASSERT(results.count.wait(test.ws) == 0); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); commitFulfiller->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); // We can also assert that we leave the database empty, in case that turns out to be useful later: auto q = test.db.run("SELECT name FROM sqlite_master WHERE type='table' AND name='_cf_METADATA'"); KJ_ASSERT(q.isDone()); } { // Should be fine to call deleteAll() a few times in succession, too: KJ_ASSERT(!test.actor.isCommitScheduled()); ActorCache::DeleteAllResults results1 = test.actor.deleteAll({}, nullptr); ActorCache::DeleteAllResults results2 = test.actor.deleteAll({}, nullptr); KJ_ASSERT(test.actor.isCommitScheduled()); KJ_ASSERT(results1.backpressure == kj::none); KJ_ASSERT(results2.backpressure == kj::none); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); // Presumably fine that the deletion commits coalesce: auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); KJ_ASSERT(results1.count.wait(test.ws) == 0); KJ_ASSERT(results2.count.wait(test.ws) == 0); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); commitFulfiller->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); auto q = test.db.run("SELECT name FROM sqlite_master WHERE type='table' AND name='_cf_METADATA'"); KJ_ASSERT(q.isDone()); } } KJ_TEST("calling deleteAll() during an implicit transaction preserves alarm state") { ActorSqliteTest test; KJ_ASSERT(!test.actor.isCommitScheduled()); // Initialize alarm state to 1ms. test.setAlarm(oneMs); ActorCache::DeleteAllResults results = test.actor.deleteAll({}, nullptr); KJ_ASSERT(test.actor.isCommitScheduled()); KJ_ASSERT(results.backpressure == kj::none); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]); KJ_ASSERT(results.count.wait(test.ws) == 0); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); commitFulfiller->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_TEST("deleteAll with deleteAlarm option deletes alarm") { // Tests that deleteAll() with deleteAlarm=true deletes the alarm along with KV data, // instead of preserving the alarm as it does by default. ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Call deleteAll() with deleteAlarm=true. ActorCache::DeleteAllResults results = test.actor.deleteAll({}, nullptr, {.deleteAlarm = true}); // The alarm should now be deleted. KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); // Commit should include scheduling the alarm cancellation. test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(results.count.wait(test.ws) == 0); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } KJ_TEST("deleteAll without deleteAlarm option preserves alarm") { // Tests that deleteAll() without deleteAlarm (the default) preserves the alarm, // which is the existing behavior. ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Call deleteAll() without deleteAlarm (default behavior). ActorCache::DeleteAllResults results = test.actor.deleteAll({}, nullptr); // The alarm should be preserved. KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(results.count.wait(test.ws) == 0); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_TEST("deleteAll with deleteAlarm during alarm handler cancels deferred delete") { // Tests that calling deleteAll() with deleteAlarm=true while an alarm handler is running // correctly deletes the alarm and cancels the deferred alarm deletion (haveDeferredDelete). // When the handler's DeferredAlarmDeleter is dropped, it should NOT write a null alarm row // since deleteAll already handled the deletion. ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); { auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); // During the handler, getAlarm() should return none (deferred delete is active). KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); // Call deleteAll() with deleteAlarm=true while the handler is running. auto results = test.actor.deleteAll({}, nullptr, {.deleteAlarm = true}); // getAlarm() should still return none. KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); // Drop the DeferredAlarmDeleter (simulating handler success). Since deleteAll already // cleared haveDeferredDelete, this should NOT write to the metadata table. } // The deleteAll commit should go through commitImpl(), which detects the alarm moved to none // and schedules the cancellation. test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } KJ_TEST("deleteAll without deleteAlarm during alarm handler still has deferred delete") { // Tests that calling deleteAll() without deleteAlarm while an alarm handler is running // restores the alarm in metadata but leaves haveDeferredDelete active. When the handler // finishes, the deferred deletion deletes the restored alarm. ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); { auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); // During the handler, getAlarm() should return none (deferred delete is active). KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); // Call deleteAll() without deleteAlarm while the handler is running. // This restores the alarm in metadata, but haveDeferredDelete is still true. test.actor.deleteAll({}, nullptr); // getAlarm() still returns none because haveDeferredDelete is still active. KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); // Drop the DeferredAlarmDeleter (simulating handler success). This triggers // maybeDeleteDeferredAlarm() which deletes the restored alarm. } // The deleteAll commit goes first, then the deferred alarm deletion triggers its own commit // with alarm scheduling. test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } KJ_TEST("deleteAll deleteAlarm does not schedule alarm cancellation if setAlarm interleaves") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Start deleteAll with deleteAlarm=true and hold the commit. test.actor.deleteAll({}, nullptr, {.deleteAlarm = true}); auto deleteAllCommit = kj::mv(test.pollAndExpectCalls({"commit"})[0]); // While deleteAll commit is in-flight, set a later alarm. test.setAlarm(twoMs); auto setAlarmCommit = kj::mv(test.pollAndExpectCalls({"commit"})[0]); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); // Completing the deleteAll commit should NOT schedule a cancel because setAlarm interleaved. deleteAllCommit->fulfill(); test.pollAndExpectCalls({}); // Completing the setAlarm commit should schedule the new alarm time. setAlarmCommit->fulfill(); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); } KJ_TEST("rolling back transaction leaves alarm in expected state") { ActorSqliteTest test; // Initialize alarm state to 2ms. test.setAlarm(twoMs); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); { auto txn = test.actor.startTransaction(); KJ_ASSERT(expectSync(txn->getAlarm({})) == twoMs); txn->setAlarm(oneMs, {}, nullptr); KJ_ASSERT(expectSync(txn->getAlarm({})) == oneMs); // Dropping transaction without committing; should roll back. } KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); } KJ_TEST("rolling back transaction leaves deferred alarm deletion in expected state") { ActorSqliteTest test; // Initialize alarm state to 2ms. test.setAlarm(twoMs); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); { auto armResult = test.actor.armAlarmHandler(twoMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); auto txn = test.actor.startTransaction(); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); test.setAlarm(oneMs); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); txn->rollback().wait(test.ws); // After rollback, getAlarm() still returns the deferred deletion result. KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); // After rollback, no changes committed, no change in scheduled alarm. test.pollAndExpectCalls({}); } // After handler, 2ms alarm is deleted. test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } KJ_TEST("committing transaction leaves deferred alarm deletion in expected state") { ActorSqliteTest test; // Initialize alarm state to 2ms. test.setAlarm(twoMs); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); { auto armResult = test.actor.armAlarmHandler(twoMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); auto txn = test.actor.startTransaction(); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); test.setAlarm(oneMs); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); txn->commit(); // After commit, getAlarm() returns the committed value. KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); } // Alarm not deleted KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_TEST("rolling back nested transaction leaves deferred alarm deletion in expected state") { ActorSqliteTest test; // Initialize alarm state to 2ms. test.setAlarm(twoMs); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); { auto armResult = test.actor.armAlarmHandler(twoMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); auto txn1 = test.actor.startTransaction(); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); { // Rolling back nested transaction change leaves deferred deletion in place. auto txn2 = test.actor.startTransaction(); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); test.setAlarm(oneMs); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); txn2->rollback().wait(test.ws); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); { // Committing nested transaction changes parent transaction state to dirty. auto txn3 = test.actor.startTransaction(); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); test.setAlarm(oneMs); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); txn3->commit(); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); { // Nested transaction of dirty transaction is dirty, rollback has no effect. auto txn4 = test.actor.startTransaction(); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); txn4->rollback().wait(test.ws); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); txn1->rollback().wait(test.ws); // After root transaction rollback, getAlarm() still returns the deferred deletion result. KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); // After rollback, no changes committed, no change in scheduled alarm. test.pollAndExpectCalls({}); } // After handler, 2ms alarm is deleted. test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } KJ_TEST("database write operations check for brokenness") { ActorSqliteTest test({.monitorOutputGate = false}); auto promise = test.gate.onBroken(); // Break gate test.put("foo", "bar"); test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "a_rejected_commit")); KJ_EXPECT_THROW_MESSAGE("a_rejected_commit", promise.wait(test.ws)); // We don't actually set ActorSqlite's brokenness until the taskFailed handler runs... test.ws.poll(); // Try making a write operation to the database, expecting it to throw the broken message via // the onWrite handler: KJ_EXPECT_THROW_MESSAGE( "a_rejected_commit", test.db.run("CREATE TABLE IF NOT EXISTS counter (count INTEGER)")); test.pollAndExpectCalls({}); } KJ_TEST("allowUnconfirmed put does not block output gate") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Do an unconfirmed put test.put("foo", "bar", {.allowUnconfirmed = true}); // Gate still isn't blocked, because we set `allowUnconfirmed`. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction. test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should still not be blocked after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify data was written KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); } KJ_TEST("confirmed put blocks output gate") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Do a confirmed put (default behavior) test.put("foo", "bar", {.allowUnconfirmed = false}); // Now it should be blocked. KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction. test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should unblock after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify data was written KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); } KJ_TEST("mixed confirmed and unconfirmed writes in same transaction use output gate") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Do an unconfirmed put followed by a confirmed put in the same transaction batch test.put("foo", "bar", {.allowUnconfirmed = true}); test.put("baz", "quux", {.allowUnconfirmed = false}); // Since any write in the batch needs confirmation, the entire batch should use output gate KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction. test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should unblock after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Both writes should be committed KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("quux").asBytes()); } KJ_TEST("allowUnconfirmed delete does not block output gate") { ActorSqliteTest test; // First set up some data test.put("foo", "bar"); test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should be unblocked after setup KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Perform an unconfirmed delete - need to add delete helper method or use actor directly expectSync(test.actor.delete_(kj::str("foo"), {.allowUnconfirmed = true}, nullptr)); // Gate still isn't blocked, because we set `allowUnconfirmed`. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction. test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should still not be blocked after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Data should be deleted KJ_ASSERT(expectSync(test.get("foo")) == kj::none); } KJ_TEST("allowUnconfirmed putMultiple does not block output gate") { ActorSqliteTest test; // Gate should be unblocked at start KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Create multiple key-value pairs for the test kj::Vector putKVs; putKVs.add(ActorCache::KeyValuePair{kj::str("foo"), kj::heapArray(kj::str("bar").asBytes())}); putKVs.add(ActorCache::KeyValuePair{kj::str("baz"), kj::heapArray(kj::str("qux").asBytes())}); putKVs.add(ActorCache::KeyValuePair{kj::str("key3"), kj::heapArray(kj::str("value3").asBytes())}); // Perform an unconfirmed putMultiple within the implicit transaction test.putMultiple(putKVs.releaseAsArray(), {.allowUnconfirmed = true}); // Gate still isn't blocked, because we set `allowUnconfirmed`. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction. test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should still not be blocked after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify all data was written correctly KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("qux").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("key3"))) == kj::str("value3").asBytes()); } KJ_TEST("allowUnconfirmed deleteMultiple does not block output gate") { ActorSqliteTest test; // First set up some data kj::Vector putKVs; putKVs.add(ActorCache::KeyValuePair{kj::str("foo"), kj::heapArray(kj::str("bar").asBytes())}); putKVs.add(ActorCache::KeyValuePair{kj::str("baz"), kj::heapArray(kj::str("qux").asBytes())}); putKVs.add(ActorCache::KeyValuePair{kj::str("key3"), kj::heapArray(kj::str("value3").asBytes())}); test.putMultiple(putKVs.releaseAsArray()); test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should be unblocked after setup KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Create array of keys to delete kj::Vector deleteKeys; deleteKeys.add(kj::str("foo")); deleteKeys.add(kj::str("baz")); deleteKeys.add(kj::str("key3")); // Perform an unconfirmed deleteMultiple KJ_EXPECT(expectSync( test.deleteMultiple(deleteKeys.releaseAsArray(), {.allowUnconfirmed = true})) == 3); // Gate still isn't blocked, because we set `allowUnconfirmed`. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction. test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should still not be blocked after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify all data was deleted KJ_ASSERT(expectSync(test.get("foo")) == kj::none); KJ_ASSERT(expectSync(test.get("baz")) == kj::none); KJ_ASSERT(expectSync(test.get("key3")) == kj::none); } KJ_TEST("unconfirmed write failure still breaks output gate") { ActorSqliteTest test({.monitorOutputGate = false}); auto promise = test.gate.onBroken(); // Do an unconfirmed put test.put("foo", "bar", {.allowUnconfirmed = true}); // The output gate is not applied initially. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); KJ_ASSERT(!promise.poll(test.ws)); // Reject the commit to simulate failure test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "flush failed hard")); // Gate should be broken due to commit failure KJ_EXPECT_THROW_MESSAGE("flush failed hard", promise.wait(test.ws)); } KJ_TEST("Direct SQL queries are confirmed writes") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); auto& db = KJ_ASSERT_NONNULL(test.actor.getSqliteDatabase()); db.run("CREATE TABLE myTable (i INTEGER PRIMARY KEY, s TEXT)"); db.run("INSERT INTO myTable VALUES (1, \"a\")"); // Now the gate should be blocked. KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction. test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should unblock after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Make sure that the write actually succeeded. { auto query = db.run("SELECT * FROM myTable"); KJ_ASSERT(!query.isDone()); KJ_EXPECT(query.getInt64(0) == 1); KJ_EXPECT(query.getText(1) == "a"); query.nextRow(); KJ_ASSERT(query.isDone()); } } KJ_TEST("An unconfirmed put followed by a direct SQL queries requires the output gate") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); test.put("foo", "bar", {.allowUnconfirmed = true}); auto& db = KJ_ASSERT_NONNULL(test.actor.getSqliteDatabase()); db.run("CREATE TABLE myTable (i INTEGER PRIMARY KEY, s TEXT)"); db.run("INSERT INTO myTable VALUES (1, \"a\")"); // Now the gate should be blocked. KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction. test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should unblock after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Make sure that the write actually succeeded. KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); { auto query = db.run("SELECT * FROM myTable"); KJ_ASSERT(!query.isDone()); KJ_EXPECT(query.getInt64(0) == 1); KJ_EXPECT(query.getText(1) == "a"); query.nextRow(); KJ_ASSERT(query.isDone()); } } KJ_TEST("sync() returns immediately when no writes are pending") { ActorSqliteTest test; // When there are no pending writes, sync() should return a resolved promise auto syncResult = test.sync(); auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult)); KJ_ASSERT(syncPromise.poll(test.ws)); } KJ_TEST("sync() waits for confirmed writes to complete") { ActorSqliteTest test; // Do a confirmed write (default behavior) test.put("foo", "bar", {.allowUnconfirmed = false}); // sync() should return a promise that blocks until the commit completes auto syncResult = test.sync(); KJ_ASSERT(syncResult != kj::none); auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult)); // The sync promise should not be ready yet KJ_ASSERT(!syncPromise.poll(test.ws)); // Complete the commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Now the sync promise should be ready KJ_ASSERT(syncPromise.poll(test.ws)); syncPromise.wait(test.ws); // Verify data was written KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); } KJ_TEST("sync() waits for unconfirmed writes to complete") { ActorSqliteTest test; // Do an unconfirmed write test.put("foo", "bar", {.allowUnconfirmed = true}); // sync() should still return a promise that blocks until the commit completes auto syncResult = test.sync(); KJ_ASSERT(syncResult != kj::none); auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult)); // The sync promise should not be ready yet KJ_ASSERT(!syncPromise.poll(test.ws)); // Complete the commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Now the sync promise should be ready KJ_ASSERT(syncPromise.poll(test.ws)); syncPromise.wait(test.ws); // Verify data was written KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); } KJ_TEST("sync() waits for multiple unconfirmed writes in a row") { ActorSqliteTest test; // Do multiple unconfirmed writes - they should batch into a single transaction test.put("foo", "bar", {.allowUnconfirmed = true}); test.put("baz", "qux", {.allowUnconfirmed = true}); test.put("key3", "value3", {.allowUnconfirmed = true}); // sync() should wait for the batched commit auto syncResult = test.sync(); KJ_ASSERT(syncResult != kj::none); auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult)); // The sync promise should not be ready yet KJ_ASSERT(!syncPromise.poll(test.ws)); // Complete the single batched commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Now the sync promise should be ready KJ_ASSERT(syncPromise.poll(test.ws)); syncPromise.wait(test.ws); // Verify all writes were committed KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("qux").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("key3"))) == kj::str("value3").asBytes()); } KJ_TEST("sync() only waits for writes before it was called") { ActorSqliteTest test; // First write test.put("foo", "bar", {.allowUnconfirmed = true}); // Call sync for the first write auto syncResult = test.sync(); KJ_ASSERT(syncResult != kj::none); auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult)); // Complete first commit auto firstCommit = kj::mv(test.pollAndExpectCalls({"commit"})[0]); firstCommit->fulfill(); // First sync should complete KJ_ASSERT(syncPromise.poll(test.ws)); syncPromise.wait(test.ws); // Second write after sync was called test.put("baz", "qux", {.allowUnconfirmed = true}); // The original sync should still be complete (doesn't wait for new write) // To verify this, let's get a new sync that should wait for the second write auto syncResult2 = test.sync(); KJ_ASSERT(syncResult2 != kj::none); auto syncPromise2 = kj::mv(KJ_ASSERT_NONNULL(syncResult2)); // Second sync should not be ready KJ_ASSERT(!syncPromise2.poll(test.ws)); // Complete second commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Now second sync should complete KJ_ASSERT(syncPromise2.poll(test.ws)); syncPromise2.wait(test.ws); } KJ_TEST("sync() propagates commit errors") { ActorSqliteTest test({.monitorOutputGate = false}); auto promise = test.gate.onBroken(); // Do an unconfirmed write test.put("foo", "bar", {.allowUnconfirmed = true}); // Call sync auto syncResult = test.sync(); KJ_ASSERT(syncResult != kj::none); auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult)); // The sync promise should not be ready yet KJ_ASSERT(!syncPromise.poll(test.ws)); // Reject the commit to simulate failure test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "commit failed")); // sync promise should become ready with an exception KJ_ASSERT(syncPromise.poll(test.ws)); KJ_EXPECT_THROW_MESSAGE("commit failed", syncPromise.wait(test.ws)); // Gate should also be broken KJ_EXPECT_THROW_MESSAGE("commit failed", promise.wait(test.ws)); } KJ_TEST("sync() with mixed confirmed and unconfirmed writes") { ActorSqliteTest test; // Do an unconfirmed write followed by a confirmed write test.put("foo", "bar", {.allowUnconfirmed = true}); test.put("baz", "qux", {.allowUnconfirmed = false}); // sync() should wait for both writes auto syncResult = test.sync(); KJ_ASSERT(syncResult != kj::none); auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult)); // The sync promise should not be ready yet KJ_ASSERT(!syncPromise.poll(test.ws)); // Complete the commit (both writes are in the same transaction) test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Now the sync promise should be ready KJ_ASSERT(syncPromise.poll(test.ws)); syncPromise.wait(test.ws); // Both writes should be committed KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("qux").asBytes()); } KJ_TEST("multiple sync() calls for same commit") { ActorSqliteTest test; // Do a write test.put("foo", "bar", {.allowUnconfirmed = true}); // Call sync multiple times - they should all wait for the same commit auto syncResult1 = test.sync(); auto syncResult2 = test.sync(); auto syncResult3 = test.sync(); KJ_ASSERT(syncResult1 != kj::none); KJ_ASSERT(syncResult2 != kj::none); KJ_ASSERT(syncResult3 != kj::none); auto syncPromise1 = kj::mv(KJ_ASSERT_NONNULL(syncResult1)); auto syncPromise2 = kj::mv(KJ_ASSERT_NONNULL(syncResult2)); auto syncPromise3 = kj::mv(KJ_ASSERT_NONNULL(syncResult3)); // None should be ready yet KJ_ASSERT(!syncPromise1.poll(test.ws)); KJ_ASSERT(!syncPromise2.poll(test.ws)); KJ_ASSERT(!syncPromise3.poll(test.ws)); // Complete the commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // All sync promises should become ready KJ_ASSERT(syncPromise1.poll(test.ws)); KJ_ASSERT(syncPromise2.poll(test.ws)); KJ_ASSERT(syncPromise3.poll(test.ws)); syncPromise1.wait(test.ws); syncPromise2.wait(test.ws); syncPromise3.wait(test.ws); } KJ_TEST("allowUnconfirmed setAlarm does not block output gate") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Do an unconfirmed setAlarm test.setAlarm(oneMs, {.allowUnconfirmed = true}); // Gate still isn't blocked, because we set `allowUnconfirmed`. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction - alarm scheduling happens before commit test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should still not be blocked after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify alarm was set KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_TEST("confirmed setAlarm blocks output gate") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Do a confirmed setAlarm (default behavior) test.setAlarm(oneMs, {.allowUnconfirmed = false}); // Gate should be blocked after scheduling starts test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction. test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should unblock after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify alarm was set KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_TEST("allowUnconfirmed setAlarm then confirmed put uses output gate") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Do an unconfirmed setAlarm followed by a confirmed put in the same transaction batch test.setAlarm(oneMs, {.allowUnconfirmed = true}); test.put("foo", "bar", {.allowUnconfirmed = false}); // Since any write in the batch needs confirmation, the entire batch should use output gate test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction. test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should unblock after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Both operations should be committed KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); } KJ_TEST("allowUnconfirmed setAlarm with storage ops") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Do unconfirmed setAlarm with unconfirmed storage writes test.setAlarm(oneMs, {.allowUnconfirmed = true}); test.put("foo", "bar", {.allowUnconfirmed = true}); // Gate still isn't blocked since both operations are unconfirmed KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should still not be blocked KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify both alarm and storage writes committed KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); } KJ_TEST("allowUnconfirmed setAlarm updating existing alarm") { ActorSqliteTest test; // Initialize alarm state to 2ms. test.setAlarm(twoMs); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); // Gate should be unblocked after setup KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Update alarm to earlier time with allowUnconfirmed test.setAlarm(oneMs, {.allowUnconfirmed = true}); // Gate still isn't blocked KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction - when moving alarm earlier, schedule happens first test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should still not be blocked KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify alarm was updated KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } KJ_TEST("allowUnconfirmed setAlarm to later time") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Gate should be unblocked after setup KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Update alarm to later time with allowUnconfirmed test.setAlarm(twoMs, {.allowUnconfirmed = true}); // Gate still isn't blocked KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction - when moving alarm later, commit happens first test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); // Gate should still not be blocked KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify alarm was updated KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); } KJ_TEST("allowUnconfirmed setAlarm to clear alarm") { ActorSqliteTest test; // Initialize alarm state to 1ms. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Gate should be unblocked after setup KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Clear alarm with allowUnconfirmed test.setAlarm(kj::none, {.allowUnconfirmed = true}); // Gate still isn't blocked KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the transaction test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); // Gate should still not be blocked KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify alarm was cleared KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } KJ_TEST("unconfirmed setAlarm failure still breaks output gate") { ActorSqliteTest test({.monitorOutputGate = false}); auto promise = test.gate.onBroken(); // Do an unconfirmed setAlarm test.setAlarm(oneMs, {.allowUnconfirmed = true}); // The output gate is not applied initially. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); KJ_ASSERT(!promise.poll(test.ws)); // Fulfill scheduleRun but reject the commit to simulate failure test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "alarm commit failed")); // Gate should be broken due to commit failure KJ_EXPECT_THROW_MESSAGE("alarm commit failed", promise.wait(test.ws)); } KJ_TEST("sync() throws after critical error in explicit transaction") { ActorSqliteTest test({.monitorOutputGate = false}); auto heapLimit = [&]() { auto row = test.db.run("PRAGMA hard_heap_limit"); return row.getInt(0); }(); KJ_DBG(heapLimit); KJ_DEFER(sqlite3_hard_heap_limit64(heapLimit);); // Start an explicit transaction auto txn = test.actor.startTransaction(); // Do a write within the transaction txn->put(kj::str("foo"), kj::heapArray(kj::str("bar").asBytes()), {}, nullptr); // Trigger a critical error using SQLITE_NOMEM by setting a very low heap limit // and then trying to insert a large value. try { // Set SQLite's memory limit very low to trigger SQLITE_NOMEM test.db.run("PRAGMA hard_heap_limit=8192"); // 8KB limit // Create data that will exceed the memory limit auto largeData = kj::heapArray(50000, 'X'); // 50KB // This should trigger SQLITE_NOMEM, causing SQLite to auto-rollback the transaction, // which will trigger the critical error handler. // // We have to copy it again in order to convert to a Array from an Array. txn->put(kj::str("large_key"), kj::heapArray(largeData.asBytes()), /*options=*/{}, nullptr); KJ_FAIL_ASSERT("Query should have failed with SQLITE_NOMEM"); } catch (kj::Exception& e) { // Expected: out of memory error. We catch and ignore this to continue the test. KJ_ASSERT(e.getDescription().contains("out of memory")); } // sync() should also throw an exception because the storage is now broken auto syncResult = test.sync(); KJ_ASSERT(syncResult != kj::none); auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult)); KJ_EXPECT_THROW_MESSAGE("broken", syncPromise.wait(test.ws)); // The transaction is now in a broken state due to the critical error. // Attempting to commit should fail. KJ_EXPECT_THROW_MESSAGE("broken", txn->commit()); } KJ_TEST("allowUnconfirmed put in explicit transaction does not block output gate") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Start an explicit transaction auto txn = test.actor.startTransaction(); // Do an unconfirmed put within the transaction txn->put( kj::str("foo"), kj::heapArray(kj::str("bar").asBytes()), {.allowUnconfirmed = true}, nullptr); // Gate still isn't blocked during the transaction, because we set `allowUnconfirmed`. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit the transaction txn->commit(); // Gate should still not be blocked during commit because all writes were unconfirmed KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should still not be blocked after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify data was written KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); } KJ_TEST("confirmed put in explicit transaction blocks output gate on commit") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Start an explicit transaction auto txn = test.actor.startTransaction(); // Do a confirmed put (default behavior) txn->put(kj::str("foo"), kj::heapArray(kj::str("bar").asBytes()), {.allowUnconfirmed = false}, nullptr); // Gate should still not be blocked during the transaction - explicit txns only lock on commit KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit the transaction txn->commit(); // Now the gate should be blocked because we're committing a confirmed write KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws)); // Complete the commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should unblock after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify data was written KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); } KJ_TEST("mixed confirmed and unconfirmed puts in explicit transaction use output gate") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Start an explicit transaction auto txn = test.actor.startTransaction(); // Do an unconfirmed put followed by a confirmed put txn->put( kj::str("foo"), kj::heapArray(kj::str("bar").asBytes()), {.allowUnconfirmed = true}, nullptr); txn->put(kj::str("baz"), kj::heapArray(kj::str("quux").asBytes()), {.allowUnconfirmed = false}, nullptr); // Gate should still not be blocked during the transaction KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit the transaction txn->commit(); // Since any write in the transaction needs confirmation, commit should use output gate KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws)); // Complete the commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should unblock after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Both writes should be committed KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("quux").asBytes()); } KJ_TEST("allowUnconfirmed delete in explicit transaction does not block output gate") { ActorSqliteTest test; // First set up some data test.put("foo", "bar"); test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should be unblocked after setup KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Start an explicit transaction auto txn = test.actor.startTransaction(); // Perform an unconfirmed delete expectSync(txn->delete_(kj::str("foo"), {.allowUnconfirmed = true}, nullptr)); // Gate still isn't blocked during the transaction KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit the transaction txn->commit(); // Gate should still not be blocked during commit because the delete was unconfirmed KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should still not be blocked after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify data was deleted KJ_ASSERT(expectSync(test.get("foo")) == kj::none); } KJ_TEST("allowUnconfirmed putMultiple in explicit transaction does not block output gate") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Start an explicit transaction auto txn = test.actor.startTransaction(); // Do an unconfirmed putMultiple auto pairs = kj::heapArrayBuilder(2); pairs.add(ActorCacheOps::KeyValuePair{kj::str("foo"), kj::heapArray(kj::str("bar").asBytes())}); pairs.add(ActorCacheOps::KeyValuePair{kj::str("baz"), kj::heapArray(kj::str("quux").asBytes())}); txn->put(pairs.finish(), {.allowUnconfirmed = true}, nullptr); // Gate still isn't blocked during the transaction KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit the transaction txn->commit(); // Gate should still not be blocked during commit KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should still not be blocked after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify data was written KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("quux").asBytes()); } KJ_TEST("allowUnconfirmed deleteMultiple in explicit transaction does not block output gate") { ActorSqliteTest test; // First set up some data test.put("foo", "bar"); test.put("baz", "quux"); test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should be unblocked after setup KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Start an explicit transaction auto txn = test.actor.startTransaction(); // Perform an unconfirmed deleteMultiple auto keys = kj::heapArrayBuilder(2); keys.add(kj::str("foo")); keys.add(kj::str("baz")); expectSync(txn->delete_(keys.finish(), {.allowUnconfirmed = true}, nullptr)); // Gate still isn't blocked during the transaction KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit the transaction txn->commit(); // Gate should still not be blocked during commit KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should still not be blocked after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify data was deleted KJ_ASSERT(expectSync(test.get("foo")) == kj::none); KJ_ASSERT(expectSync(test.get("baz")) == kj::none); } KJ_TEST("allowUnconfirmed setAlarm in explicit transaction does not block output gate") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Start an explicit transaction auto txn = test.actor.startTransaction(); // Set an alarm with allowUnconfirmed txn->setAlarm(oneMs, {.allowUnconfirmed = true}, nullptr); // Gate still isn't blocked during the transaction KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit the transaction txn->commit(); // Gate should still not be blocked during commit KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the scheduleRun and commit test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should still not be blocked after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify alarm was set KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.getAlarm())) == oneMs); } KJ_TEST("nested transaction: unconfirmed child commit does not block output gate") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Start a parent transaction auto parentTxn = test.actor.startTransaction(); // Do an unconfirmed put in the parent parentTxn->put(kj::str("parent"), kj::heapArray(kj::str("data").asBytes()), {.allowUnconfirmed = true}, nullptr); { // Start a nested child transaction auto childTxn = test.actor.startTransaction(); // Do an unconfirmed put in the child childTxn->put(kj::str("child"), kj::heapArray(kj::str("data").asBytes()), {.allowUnconfirmed = true}, nullptr); // Gate still isn't blocked KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit the child transaction childTxn->commit(); } // Gate should still not be blocked after child commit KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit the parent transaction parentTxn->commit(); // Gate should still not be blocked during parent commit because all writes were unconfirmed KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should still not be blocked after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify both writes were committed KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("parent"))) == kj::str("data").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("child"))) == kj::str("data").asBytes()); } KJ_TEST("nested transaction: confirmed child propagates to parent commit") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Start a parent transaction with unconfirmed write auto parentTxn = test.actor.startTransaction(); parentTxn->put(kj::str("parent"), kj::heapArray(kj::str("data").asBytes()), {.allowUnconfirmed = true}, nullptr); { // Start a nested child transaction auto childTxn = test.actor.startTransaction(); // Do a confirmed put in the child childTxn->put(kj::str("child"), kj::heapArray(kj::str("data").asBytes()), {.allowUnconfirmed = false}, nullptr); // Gate still isn't blocked during the transaction KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit the child transaction - this should propagate someWriteConfirmed to parent childTxn->commit(); } // Gate should still not be blocked after child commit (no real commit yet) KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit the parent transaction parentTxn->commit(); // Now the gate should be blocked because the child had a confirmed write KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws)); // Complete the commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should unblock after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify both writes were committed KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("parent"))) == kj::str("data").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("child"))) == kj::str("data").asBytes()); } KJ_TEST("nested transaction: confirmed parent with unconfirmed child blocks output gate") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Start a parent transaction with confirmed write auto parentTxn = test.actor.startTransaction(); parentTxn->put(kj::str("parent"), kj::heapArray(kj::str("data").asBytes()), {.allowUnconfirmed = false}, nullptr); { // Start a nested child transaction auto childTxn = test.actor.startTransaction(); // Do an unconfirmed put in the child childTxn->put(kj::str("child"), kj::heapArray(kj::str("data").asBytes()), {.allowUnconfirmed = true}, nullptr); // Gate still isn't blocked during the transaction KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit the child transaction childTxn->commit(); } // Gate should still not be blocked after child commit KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit the parent transaction parentTxn->commit(); // Now the gate should be blocked because the parent had a confirmed write KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws)); // Complete the commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should unblock after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify both writes were committed KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("parent"))) == kj::str("data").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("child"))) == kj::str("data").asBytes()); } KJ_TEST("nested transaction: deeply nested confirmed write propagates to root") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Start a parent transaction with unconfirmed write auto txn1 = test.actor.startTransaction(); txn1->put(kj::str("level1"), kj::heapArray(kj::str("data").asBytes()), {.allowUnconfirmed = true}, nullptr); { // Start a second level nested transaction with unconfirmed write auto txn2 = test.actor.startTransaction(); txn2->put(kj::str("level2"), kj::heapArray(kj::str("data").asBytes()), {.allowUnconfirmed = true}, nullptr); { // Start a third level nested transaction with confirmed write auto txn3 = test.actor.startTransaction(); txn3->put(kj::str("level3"), kj::heapArray(kj::str("data").asBytes()), {.allowUnconfirmed = false}, nullptr); // Gate still isn't blocked during the transaction KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit level 3 - should propagate someWriteConfirmed to level 2 txn3->commit(); } KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit level 2 - should propagate someWriteConfirmed to level 1 txn2->commit(); } KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit level 1 (root transaction) txn1->commit(); // Now the gate should be blocked because level 3 had a confirmed write KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws)); // Complete the commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should unblock after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify all writes were committed KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("level1"))) == kj::str("data").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("level2"))) == kj::str("data").asBytes()); KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("level3"))) == kj::str("data").asBytes()); } KJ_TEST("nested transaction: rollback resets someWriteConfirmed flag") { ActorSqliteTest test; // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Start a parent transaction with unconfirmed write auto parentTxn = test.actor.startTransaction(); parentTxn->put(kj::str("parent"), kj::heapArray(kj::str("data").asBytes()), {.allowUnconfirmed = true}, nullptr); { // Start a nested child transaction auto childTxn = test.actor.startTransaction(); // Do a confirmed put in the child childTxn->put(kj::str("child"), kj::heapArray(kj::str("data").asBytes()), {.allowUnconfirmed = false}, nullptr); // Gate still isn't blocked during the transaction KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Rollback the child transaction instead of committing childTxn->rollback().wait(test.ws); } // Gate should still not be blocked KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Commit the parent transaction parentTxn->commit(); // Gate should still not be blocked because child was rolled back KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Complete the commit test.pollAndExpectCalls({"commit"})[0]->fulfill(); // Gate should still not be blocked after commit completes KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Verify only parent write was committed KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("parent"))) == kj::str("data").asBytes()); KJ_ASSERT(expectSync(test.get("child")) == kj::none); } KJ_TEST("explicit transaction: commit failure breaks output gate even for unconfirmed writes") { ActorSqliteTest test({.monitorOutputGate = false}); auto promise = test.gate.onBroken(); // Gate is currently not blocked. KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Start an explicit transaction auto txn = test.actor.startTransaction(); // Do an unconfirmed put txn->put( kj::str("foo"), kj::heapArray(kj::str("bar").asBytes()), {.allowUnconfirmed = true}, nullptr); // Commit the transaction txn->commit(); // Gate should not be blocked yet because write was unconfirmed KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws)); // Reject the commit to simulate failure test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "commit failed")); // Gate should now be broken due to commit failure, even though write was unconfirmed KJ_EXPECT_THROW_MESSAGE("commit failed", promise.wait(test.ws)); } KJ_TEST("ActorSqlite alarm cleared by abandonAlarm") { ActorSqliteTest test; test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // abandonAlarm() clears the alarm from SQLite: // setAlarm(null) -> commit -> scheduleRun(none) (move-later path). auto result = test.actor.abandonAlarm(oneMs).wait(test.ws); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill(); test.pollAndExpectCalls({}); // Returns kj::none: alarm was cleared, AlarmManager should not re-register. KJ_ASSERT(result == kj::none); // getAlarm() now returns null (alarm deleted from SQLite). KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); } KJ_TEST("ActorSqlite alarm preserved after ALARM_RETRY_MAX_TRIES uncounted (internal) failures") { // When all ALARM_RETRY_MAX_TRIES failures are uncounted (retryCountsAgainstLimit=false, // i.e. infrastructure errors), the alarm scheduler's countedRetry never reaches the limit and // abandonAlarm is NEVER called. The alarm must remain set in SQLite throughout so that // the scheduler can keep retrying indefinitely. ActorSqliteTest test; test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Simulate uncounted failures well past ALARM_RETRY_MAX_TRIES (= 6). // countedRetry stays at 0; AlarmManager never gives up; abandonAlarm is never called. // We've seen alarms fail hundreds of times due to infrastructure errors in production, // so we check both at the boundary (6) and well beyond it (100). for (auto i = 0; i < 100; i++) { auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); test.actor.cancelDeferredAlarmDeletion(); test.pollAndExpectCalls({}); // Check at the ALARM_RETRY_MAX_TRIES boundary and at the end. if (i == 5 || i == 99) { KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } } } KJ_TEST("ActorSqlite abandonAlarm is a no-op when a newer alarm has replaced the abandoned one") { // If the user sets a new alarm between the last retry failure and the abandonAlarm() call, // and it has already committed to SQLite, abandonAlarm() must compare the time and leave // the new alarm untouched. ActorSqliteTest test; // Set the original alarm and commit it. test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // User sets a new alarm (twoMs). // The commit fires first; then the post-commit "move-later" logic fires scheduleRun(2ms) // because alarmScheduledNoLaterThan (oneMs) is earlier than the newly committed twoMs. test.setAlarm(twoMs); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); // abandonAlarm() for the original oneMs alarm must be a no-op: storedTime (twoMs) != // scheduledTime (oneMs), so the time check prevents clearing the new alarm. // Returns twoMs so AlarmManager can re-register the actor's real alarm. auto result = test.actor.abandonAlarm(oneMs).wait(test.ws); test.pollAndExpectCalls({}); // No commit or scheduleRun -- correct no-op. KJ_ASSERT(KJ_ASSERT_NONNULL(result) == twoMs); // getAlarm() must still return twoMs. KJ_ASSERT(expectSync(test.getAlarm()) == twoMs); } KJ_TEST("ActorSqlite abandonAlarm returns kj::none when no alarm is stored") { ActorSqliteTest test; // No alarm ever set. abandonAlarm should be a pure no-op, returning kj::none. auto result = test.actor.abandonAlarm(oneMs).wait(test.ws); test.pollAndExpectCalls({}); KJ_ASSERT(result == kj::none); } KJ_TEST("ActorSqlite abandonAlarm returns kj::none when inAlarmHandler") { ActorSqliteTest test; test.setAlarm(oneMs); test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill(); test.pollAndExpectCalls({"commit"})[0]->fulfill(); test.pollAndExpectCalls({}); KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); // Arm the handler — inAlarmHandler is now true. auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime); KJ_ASSERT(armResult.is()); // abandonAlarm while handler is running: returns kj::none (handler owns the alarm). auto result = test.actor.abandonAlarm(oneMs).wait(test.ws); test.pollAndExpectCalls({}); KJ_ASSERT(result == kj::none); // kj::none because haveDeferredDelete hides the alarm during the handler. KJ_ASSERT(expectSync(test.getAlarm()) == kj::none); // Cancel the deferred delete so cleanup doesn't trigger a commit. test.actor.cancelDeferredAlarmDeletion(); // After cancellation, getAlarm() reads SQLite again -- oneMs is still there since // abandonAlarm was a no-op while the handler was running. KJ_ASSERT(expectSync(test.getAlarm()) == oneMs); } } // namespace } // namespace workerd