Skip to content
File

Blob: src/workerd/io/actor-sqlite-test.c++

118.6 KB
1// Copyright (c) 2024 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#include "actor-sqlite.h"
6#include "io-gate.h"
7 
8#include <workerd/util/capnp-mock.h>
9#include <workerd/util/test.h>
10 
11#include <sqlite3.h>
12 
13#include <kj/debug.h>
14#include <kj/test.h>
15 
16namespace workerd {
17namespace {
18 
19static constexpr kj::Date oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH;
20static constexpr kj::Date twoMs = 2 * kj::MILLISECONDS + kj::UNIX_EPOCH;
21static constexpr kj::Date threeMs = 3 * kj::MILLISECONDS + kj::UNIX_EPOCH;
22static constexpr kj::Date fourMs = 4 * kj::MILLISECONDS + kj::UNIX_EPOCH;
23static constexpr kj::Date fiveMs = 5 * kj::MILLISECONDS + kj::UNIX_EPOCH;
24static constexpr kj::Date sixMs = 6 * kj::MILLISECONDS + kj::UNIX_EPOCH;
25static constexpr kj::Date tenMs = 10 * kj::MILLISECONDS + kj::UNIX_EPOCH;
26// Used as the "current time" parameter for armAlarmHandler in tests.
27// Set to epoch (before all test alarm times) so existing tests aren't affected by
28// the overdue alarm check.
29static constexpr kj::Date testCurrentTime = kj::UNIX_EPOCH;
30 
31template <typename T>
32kj::Promise<T> eagerlyReportExceptions(kj::Promise<T> promise, kj::SourceLocation location = {}) {
33 return promise.eagerlyEvaluate([location](kj::Exception&& e) -> T {
34 KJ_LOG_AT(ERROR, location, e);
35 kj::throwFatalException(kj::mv(e));
36 });
37}
38 
39// Expect that a synchronous result is returned.
40template <typename T>
41T expectSync(kj::OneOf<T, kj::Promise<T>> result, kj::SourceLocation location = {}) {
42 KJ_SWITCH_ONEOF(result) {
43 KJ_CASE_ONEOF(promise, kj::Promise<T>) {
44 KJ_FAIL_ASSERT_AT(location, "result was unexpectedly asynchronous");
45 }
46 KJ_CASE_ONEOF(value, T) {
47 return kj::mv(value);
48 }
49 }
50 KJ_UNREACHABLE;
51}
52 
53struct ActorSqliteTestOptions final {
54 bool monitorOutputGate = true;
55};
56 
57struct ActorSqliteTest final {
58 kj::EventLoop loop;
59 kj::WaitScope ws;
60 
61 OutputGate gate;
62 kj::Own<const kj::Directory> vfsDir;
63 SqliteDatabase::Vfs vfs;
64 SqliteDatabase db;
65 
66 struct Call final {
67 kj::String desc;
68 kj::Own<kj::PromiseFulfiller<void>> fulfiller;
69 };
70 kj::Vector<Call> calls;
71 
72 struct ActorSqliteTestHooks final: public ActorSqlite::Hooks {
73 public:
74 explicit ActorSqliteTestHooks(ActorSqliteTest& parent): parent(parent) {}
75 
76 kj::Promise<void> scheduleRun(
77 kj::Maybe<kj::Date> newAlarmTime, kj::Promise<void> priorTask) override {
78 KJ_IF_SOME(h, parent.scheduleRunWithPriorHandler) {
79 return h(newAlarmTime, kj::mv(priorTask));
80 }
81 KJ_IF_SOME(h, parent.scheduleRunHandler) {
82 return h(newAlarmTime);
83 }
84 auto desc = newAlarmTime.map([](auto& t) {
85 return kj::str("scheduleRun(", t, ")");
86 }).orDefault(kj::str("scheduleRun(none)"));
87 auto [promise, fulfiller] = kj::newPromiseAndFulfiller<void>();
88 parent.calls.add(Call{kj::mv(desc), kj::mv(fulfiller)});
89 return kj::mv(promise);
90 }
91 
92 ActorSqliteTest& parent;
93 };
94 kj::Maybe<kj::Function<kj::Promise<void>(kj::Maybe<kj::Date>)>> scheduleRunHandler;
95 kj::Maybe<kj::Function<kj::Promise<void>(kj::Maybe<kj::Date>, kj::Promise<void>)>>
96 scheduleRunWithPriorHandler;
97 ActorSqliteTestHooks hooks = ActorSqliteTestHooks(*this);
98 
99 ActorSqlite actor;
100 
101 kj::Promise<void> gateBrokenPromise;
102 kj::UnwindDetector unwindDetector;
103 
104 explicit ActorSqliteTest(ActorSqliteTestOptions options = {})
105 : ws(loop),
106 vfsDir(kj::newInMemoryDirectory(kj::nullClock())),
107 vfs(*vfsDir),
108 db(vfs, kj::Path({"foo"}), kj::WriteMode::CREATE | kj::WriteMode::MODIFY),
109 actor(kj::attachRef(db), gate, KJ_BIND_METHOD(*this, commitCallback), hooks),
110 gateBrokenPromise(options.monitorOutputGate ? eagerlyReportExceptions(gate.onBroken())
111 : kj::Promise<void>(kj::READY_NOW)) {
112 db.afterReset([this](SqliteDatabase& cbDb) {
113 KJ_ASSERT(&db == &cbDb);
114 KJ_ASSERT(actor.isCommitScheduled(),
115 "actor should have a commit scheduled during afterReset() callback");
116 });
117 }
118 
119 ~ActorSqliteTest() noexcept(false) {
120 if (!unwindDetector.isUnwinding()) {
121 // Make sure if the output gate has been broken, the exception was reported. This is
122 // important to report errors thrown inside flush(), since those won't otherwise propagate
123 // into the test body.
124 gateBrokenPromise.poll(ws);
125 
126 // Make sure there's no outstanding async work we haven't considered:
127 pollAndExpectCalls({}, "unexpected calls at end of test");
128 }
129 }
130 
131 kj::Promise<void> commitCallback(SpanParent) {
132 auto [promise, fulfiller] = kj::newPromiseAndFulfiller<void>();
133 calls.add(Call{kj::str("commit"), kj::mv(fulfiller)});
134 return kj::mv(promise);
135 }
136 
137 // Polls the event loop, then asserts that the description of calls up to this point match the
138 // expectation and returns their fulfillers. Also clears the call log.
139 //
140 // TODO(cleanup): Is there a better way to do mocks? capnp-mock looks nice, but seems a bit
141 // heavyweight for this test.
142 kj::Vector<kj::Own<kj::PromiseFulfiller<void>>> pollAndExpectCalls(
143 std::initializer_list<kj::StringPtr> expCallDescs,
144 kj::StringPtr message = ""_kj,
145 kj::SourceLocation location = {}) {
146 ws.poll();
147 auto callDescs = KJ_MAP(c, calls) { return kj::str(c.desc); };
148 KJ_ASSERT_AT(callDescs == heapArray(expCallDescs), location, kj::str(message));
149 auto fulfillers = KJ_MAP(c, calls) { return kj::mv(c.fulfiller); };
150 calls.clear();
151 return kj::mv(fulfillers);
152 }
153 
154 // A few driver methods for convenience.
155 auto get(kj::StringPtr key, ActorCache::ReadOptions options = {}) {
156 return actor.get(kj::str(key), options);
157 }
158 auto getAlarm(ActorCache::ReadOptions options = {}) {
159 return actor.getAlarm(options);
160 }
161 auto put(kj::StringPtr key, kj::StringPtr value, ActorCache::WriteOptions options = {}) {
162 return actor.put(kj::str(key), kj::heapArray(value.asBytes()), options, nullptr);
163 }
164 auto putMultiple(
165 kj::Array<ActorCache::KeyValuePair> pairs, ActorCache::WriteOptions options = {}) {
166 return actor.put(kj::mv(pairs), options, nullptr);
167 }
168 auto putMultipleExplicitTxn(
169 kj::Array<ActorCache::KeyValuePair> pairs, ActorCache::WriteOptions options = {}) {
170 auto txn = actor.startTransaction();
171 txn->put(kj::mv(pairs), options, nullptr);
172 return txn->commit();
173 }
174 auto deleteMultiple(kj::Array<kj::String> keys, ActorCache::WriteOptions options = {}) {
175 return actor.delete_(kj::mv(keys), options, nullptr);
176 }
177 auto setAlarm(kj::Maybe<kj::Date> newTime, ActorCache::WriteOptions options = {}) {
178 return actor.setAlarm(newTime, options, nullptr);
179 }
180 auto sync() {
181 return actor.onNoPendingFlush(nullptr);
182 }
183};
184 
185KJ_TEST("initial alarm value is unset") {
186 ActorSqliteTest test;
187 
188 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
189}
190 
191KJ_TEST("can set and get alarm") {
192 ActorSqliteTest test;
193 
194 test.setAlarm(oneMs);
195 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
196 test.pollAndExpectCalls({"commit"})[0]->fulfill();
197 
198 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
199}
200 
201KJ_TEST("check put multiple wraps operations in a transaction") {
202 ActorSqliteTest test;
203 
204 kj::Vector<ActorCache::KeyValuePair> putKVs;
205 putKVs.add(ActorCache::KeyValuePair{kj::str("foo"), kj::heapArray(kj::str("bar").asBytes())});
206 
207 // NoTxn test
208 {
209 // Check that we're in a NoTxn
210 KJ_ASSERT(!test.actor.isCommitScheduled());
211 test.putMultiple(putKVs.releaseAsArray());
212 // During write, all NoTxn operations are wrapped in an ImplicitTxn.
213 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
214 commitFulfiller->fulfill();
215 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
216 }
217 
218 // ExplicitTxn test
219 {
220 putKVs.add(ActorCache::KeyValuePair{kj::str("foo2"), kj::heapArray(kj::str("bar2").asBytes())});
221 KJ_ASSERT(!test.actor.isCommitScheduled());
222 // Similar to the previous putMultiple, but wrapped in a transactionSync (ExplicitTxn)
223 test.putMultipleExplicitTxn(putKVs.releaseAsArray());
224 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
225 commitFulfiller->fulfill();
226 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo2"))) == kj::str("bar2").asBytes());
227 }
228 
229 // ImplicitTxn test
230 {
231 // A single put will create an ImplicitTxn that we can use to wrap our putMultiple into.
232 KJ_ASSERT(!test.actor.isCommitScheduled());
233 test.put("baz", "bat");
234 
235 // By now, we should check there's a commit scheduled in a ImplicitTxn.
236 KJ_ASSERT(test.actor.isCommitScheduled());
237 putKVs.add(ActorCache::KeyValuePair{kj::str("foo3"), kj::heapArray(kj::str("bar3").asBytes())});
238 test.putMultiple(putKVs.releaseAsArray());
239 
240 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
241 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("bat").asBytes());
242 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo3"))) == kj::str("bar3").asBytes());
243 commitFulfiller->fulfill();
244 }
245}
246 
247KJ_TEST("check put multiple wraps operations in a transaction") {
248 ActorSqliteTest test;
249 
250 kj::Vector<ActorCache::KeyValuePair> putKVs;
251 putKVs.add(ActorCache::KeyValuePair{kj::str("foo"), kj::heapArray(kj::str("bar").asBytes())});
252 
253 // NoTxn test
254 {
255 // Check that we're in a NoTxn
256 KJ_ASSERT(!test.actor.isCommitScheduled());
257 test.putMultiple(putKVs.releaseAsArray());
258 // During write, all NoTxn operations are wrapped in an ImplicitTxn.
259 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
260 commitFulfiller->fulfill();
261 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
262 }
263 
264 // ExplicitTxn test
265 {
266 putKVs.add(ActorCache::KeyValuePair{kj::str("foo2"), kj::heapArray(kj::str("bar2").asBytes())});
267 KJ_ASSERT(!test.actor.isCommitScheduled());
268 // Similar to the previous putMultiple, but wrapped in a transactionSync (ExplicitTxn)
269 test.putMultipleExplicitTxn(putKVs.releaseAsArray());
270 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
271 commitFulfiller->fulfill();
272 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo2"))) == kj::str("bar2").asBytes());
273 }
274 
275 // ImplicitTxn test
276 {
277 // A single put will create an ImplicitTxn that we can use to wrap our putMultiple into.
278 KJ_ASSERT(!test.actor.isCommitScheduled());
279 test.put("baz", "bat");
280 
281 // By now, we should check there's a commit scheduled in a ImplicitTxn.
282 KJ_ASSERT(test.actor.isCommitScheduled());
283 putKVs.add(ActorCache::KeyValuePair{kj::str("foo3"), kj::heapArray(kj::str("bar3").asBytes())});
284 test.putMultiple(putKVs.releaseAsArray());
285 
286 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
287 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("bat").asBytes());
288 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo3"))) == kj::str("bar3").asBytes());
289 commitFulfiller->fulfill();
290 }
291}
292 
293KJ_TEST("check put multiple wraps operations in a transaction and rollback on error") {
294 ActorSqliteTest test;
295 
296 // We expect that putMultiple is all or nothing, rolling back if a single put fails.
297 
298 kj::Vector<ActorCache::KeyValuePair> putKVs;
299 
300 // Add some regular key-value pairs that we know are supported
301 putKVs.add(ActorCache::KeyValuePair{kj::str("foo"), kj::heapArray(kj::str("bar").asBytes())});
302 putKVs.add(ActorCache::KeyValuePair{kj::str("foo2"), kj::heapArray(kj::str("bar2").asBytes())});
303 putKVs.add(ActorCache::KeyValuePair{kj::str("foo3"), kj::heapArray(kj::str("bar3").asBytes())});
304 
305 // Now create a key that's too large. Should fail with string or blob too big: SQLITE_TOOBIG
306 auto tooLongKey = kj::heapString(2200000);
307 tooLongKey.asArray().fill('a');
308 // Add it to our KV array
309 putKVs.add(
310 ActorCache::KeyValuePair{kj::str(tooLongKey), kj::heapArray(kj::str("bar").asBytes())});
311 
312 // NoTxn test
313 {
314 // Check that we're in a NoTxn
315 KJ_ASSERT(!test.actor.isCommitScheduled());
316 try {
317 test.putMultiple(putKVs.releaseAsArray());
318 // During write, all NoTxn operations are wrapped in an ImplicitTxn.
319 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
320 commitFulfiller->fulfill();
321 KJ_UNREACHABLE;
322 } catch (kj::Exception& e) {
323 KJ_ASSERT(
324 e.getDescription() == "expected false; jsg.Error: string or blob too big: SQLITE_TOOBIG");
325 }
326 KJ_ASSERT(expectSync(test.get(kj::str("foo"))) == nullptr);
327 KJ_ASSERT(expectSync(test.get(kj::str("foo2"))) == nullptr);
328 KJ_ASSERT(expectSync(test.get(kj::str("foo3"))) == nullptr);
329 }
330 
331 // Reset the transaction state by going async, which will cause the ImplicitTxn to commit.
332 {
333 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
334 commitFulfiller->fulfill();
335 }
336 
337 // ExplicitTxn test
338 {
339 KJ_ASSERT(!test.actor.isCommitScheduled());
340 // Similar to the previous putMultiple, but wrapped in a transactionSync (ExplicitTxn)
341 test.putMultipleExplicitTxn(putKVs.releaseAsArray());
342 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
343 commitFulfiller->fulfill();
344 KJ_ASSERT(expectSync(test.get(kj::str("foo"))) == nullptr);
345 KJ_ASSERT(expectSync(test.get(kj::str("foo2"))) == nullptr);
346 KJ_ASSERT(expectSync(test.get(kj::str("foo3"))) == nullptr);
347 }
348 
349 // ImplicitTxn test
350 {
351 // A single put will create an ImplicitTxn that we can use to wrap our putMultiple into.
352 KJ_ASSERT(!test.actor.isCommitScheduled());
353 test.put("baz", "bat");
354 
355 // By now, we should check there's a commit scheduled in a ImplicitTxn.
356 KJ_ASSERT(test.actor.isCommitScheduled());
357 test.putMultiple(putKVs.releaseAsArray());
358 
359 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
360 // The single put succeeded, but the putMultiple did not.
361 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("bat").asBytes());
362 KJ_ASSERT(expectSync(test.get(kj::str("foo"))) == nullptr);
363 KJ_ASSERT(expectSync(test.get(kj::str("foo2"))) == nullptr);
364 KJ_ASSERT(expectSync(test.get(kj::str("foo3"))) == nullptr);
365 commitFulfiller->fulfill();
366 }
367}
368 
369KJ_TEST("alarm write happens transactionally with storage ops") {
370 ActorSqliteTest test;
371 
372 test.setAlarm(oneMs);
373 test.put("foo", "bar");
374 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
375 test.pollAndExpectCalls({"commit"})[0]->fulfill();
376 
377 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
378 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
379}
380 
381KJ_TEST("storage op without alarm change does not wait on scheduler") {
382 ActorSqliteTest test;
383 
384 test.put("foo", "bar");
385 test.pollAndExpectCalls({"commit"})[0]->fulfill();
386 
387 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
388 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
389}
390 
391KJ_TEST("alarm scheduling starts synchronously before implicit local db commit") {
392 ActorSqliteTest test;
393 
394 // In workerd (unlike edgeworker), there is no remote storage, so there is no work done in
395 // commitCallback(); the local db is considered durably stored after the synchronous sqlite
396 // commit() call returns. If a commit includes an alarm state change that requires scheduling
397 // before the commit call, it needs to happen synchronously. Since workerd synchronously
398 // schedules alarms, we just need to ensure that the database is in a pre-commit state when
399 // scheduleRun() is called.
400 
401 // Initialize alarm state to 2ms.
402 test.setAlarm(twoMs);
403 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
404 test.pollAndExpectCalls({"commit"})[0]->fulfill();
405 test.pollAndExpectCalls({});
406 
407 bool startedScheduleRun = false;
408 test.scheduleRunHandler = [&](kj::Maybe<kj::Date>) -> kj::Promise<void> {
409 startedScheduleRun = true;
410 
411 KJ_EXPECT_THROW_MESSAGE(
412 "cannot start a transaction within a transaction", test.db.run("BEGIN TRANSACTION"));
413 
414 return kj::READY_NOW;
415 };
416 
417 test.setAlarm(oneMs);
418 KJ_ASSERT(!startedScheduleRun);
419 test.ws.poll();
420 KJ_ASSERT(startedScheduleRun);
421 
422 test.pollAndExpectCalls({"commit"})[0]->fulfill();
423 
424 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
425}
426 
427KJ_TEST("alarm scheduling starts synchronously before explicit local db commit") {
428 ActorSqliteTest test;
429 
430 // Initialize alarm state to 2ms.
431 test.setAlarm(twoMs);
432 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
433 test.pollAndExpectCalls({"commit"})[0]->fulfill();
434 test.pollAndExpectCalls({});
435 
436 bool startedScheduleRun = false;
437 test.scheduleRunHandler = [&](kj::Maybe<kj::Date>) -> kj::Promise<void> {
438 startedScheduleRun = true;
439 
440 // Not sure if there is a good way to detect savepoint presence without mutating the db state,
441 // but this is sufficient to verify the test properties:
442 
443 // Verify that we are not within a nested savepoint.
444 KJ_EXPECT_THROW_MESSAGE(
445 "no such savepoint: _cf_savepoint_1", test.db.run("RELEASE _cf_savepoint_1"));
446 
447 // Verify that we are within the root savepoint.
448 test.db.run("RELEASE _cf_savepoint_0");
449 KJ_EXPECT_THROW_MESSAGE(
450 "no such savepoint: _cf_savepoint_0", test.db.run("RELEASE _cf_savepoint_0"));
451 
452 // We don't actually care what happens in the test after this point, but it's slightly simpler
453 // to re-add the savepoint to allow the test to complete cleanly:
454 test.db.run("SAVEPOINT _cf_savepoint_0");
455 
456 return kj::READY_NOW;
457 };
458 
459 {
460 auto txn = test.actor.startTransaction();
461 txn->setAlarm(oneMs, {}, nullptr);
462 
463 KJ_ASSERT(!startedScheduleRun);
464 txn->commit();
465 KJ_ASSERT(startedScheduleRun);
466 
467 test.pollAndExpectCalls({"commit"})[0]->fulfill();
468 }
469 
470 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
471}
472 
473KJ_TEST("alarm scheduling does not start synchronously before nested explicit local db commit") {
474 ActorSqliteTest test;
475 
476 // Initialize alarm state to 2ms.
477 test.setAlarm(twoMs);
478 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
479 test.pollAndExpectCalls({"commit"})[0]->fulfill();
480 test.pollAndExpectCalls({});
481 
482 bool startedScheduleRun = false;
483 test.scheduleRunHandler = [&](kj::Maybe<kj::Date>) -> kj::Promise<void> {
484 startedScheduleRun = true;
485 return kj::READY_NOW;
486 };
487 
488 {
489 auto txn1 = test.actor.startTransaction();
490 
491 {
492 auto txn2 = test.actor.startTransaction();
493 txn2->setAlarm(oneMs, {}, nullptr);
494 
495 txn2->commit();
496 KJ_ASSERT(!startedScheduleRun);
497 }
498 
499 txn1->commit();
500 KJ_ASSERT(startedScheduleRun);
501 
502 test.pollAndExpectCalls({"commit"})[0]->fulfill();
503 }
504 
505 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
506}
507 
508KJ_TEST("synchronous alarm scheduling failure causes local db commit to throw synchronously") {
509 ActorSqliteTest test({.monitorOutputGate = false});
510 auto promise = test.gate.onBroken();
511 
512 auto getLocalAlarm = [&]() -> kj::Maybe<kj::Date> {
513 auto query = test.db.run("SELECT value FROM _cf_METADATA WHERE key = 1");
514 if (query.isDone() || query.isNull(0)) {
515 return kj::none;
516 } else {
517 return kj::UNIX_EPOCH + query.getInt64(0) * kj::NANOSECONDS;
518 }
519 };
520 
521 // Initialize alarm state to 2ms.
522 test.setAlarm(twoMs);
523 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
524 test.pollAndExpectCalls({"commit"})[0]->fulfill();
525 test.pollAndExpectCalls({});
526 
527 // Override scheduleRun handler with one that throws synchronously.
528 bool startedScheduleRun = false;
529 test.scheduleRunHandler = [&](kj::Maybe<kj::Date>) -> kj::Promise<void> {
530 startedScheduleRun = true;
531 // Must throw synchronously; returning an exception is insufficient.
532 kj::throwFatalException(KJ_EXCEPTION(FAILED, "a_sync_fail"));
533 };
534 
535 KJ_ASSERT(!promise.poll(test.ws));
536 test.setAlarm(oneMs);
537 
538 // Expect that polling will attempt to commit the implicit transaction, which should
539 // synchronously fail when attempting to call scheduleRun() before the db commit, and roll back the
540 // local db state to the 2ms alarm.
541 KJ_ASSERT(!startedScheduleRun);
542 KJ_ASSERT(KJ_REQUIRE_NONNULL(getLocalAlarm()) == oneMs);
543 test.ws.poll();
544 KJ_ASSERT(startedScheduleRun);
545 KJ_ASSERT(KJ_REQUIRE_NONNULL(getLocalAlarm()) == twoMs);
546 
547 KJ_ASSERT(promise.poll(test.ws));
548 KJ_EXPECT_THROW_MESSAGE("a_sync_fail", promise.wait(test.ws));
549}
550 
551KJ_TEST("can clear alarm") {
552 ActorSqliteTest test;
553 
554 // Initialize alarm state to 1ms.
555 test.setAlarm(oneMs);
556 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
557 test.pollAndExpectCalls({"commit"})[0]->fulfill();
558 test.pollAndExpectCalls({});
559 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
560 
561 test.setAlarm(kj::none);
562 test.pollAndExpectCalls({"commit"})[0]->fulfill();
563 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
564 
565 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
566}
567 
568KJ_TEST("can set alarm twice") {
569 ActorSqliteTest test;
570 
571 test.setAlarm(oneMs);
572 test.setAlarm(twoMs);
573 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
574 test.pollAndExpectCalls({"commit"})[0]->fulfill();
575 
576 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
577}
578 
579KJ_TEST("setting duplicate alarm is no-op") {
580 ActorSqliteTest test;
581 
582 test.setAlarm(kj::none);
583 test.pollAndExpectCalls({});
584 
585 test.setAlarm(oneMs);
586 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
587 test.pollAndExpectCalls({"commit"})[0]->fulfill();
588 
589 test.setAlarm(oneMs);
590 test.pollAndExpectCalls({});
591}
592 
593KJ_TEST("tells alarm handler to cancel when committed alarm is empty") {
594 ActorSqliteTest test;
595 
596 {
597 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
598 // We expect armAlarmHandler() to tell us to cancel the alarm.
599 KJ_ASSERT(armResult.is<ActorCache::CancelAlarmHandler>());
600 auto waitPromise = kj::mv(armResult.get<ActorCache::CancelAlarmHandler>().waitBeforeCancel);
601 
602 // We also expect the alarm cancellation to contain a scheduling request to delete the alarm,
603 // to handle cases where alarm deletion was durably committed to the database, but a failure
604 // occurred before the alarm deletion was conveyed to the alarm scheduler.
605 KJ_ASSERT(!waitPromise.poll(test.ws));
606 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
607 KJ_ASSERT(waitPromise.poll(test.ws));
608 waitPromise.wait(test.ws);
609 }
610}
611 
612KJ_TEST("tells alarm handler to reschedule when handler alarm is later than committed alarm") {
613 ActorSqliteTest test;
614 
615 // Initialize alarm state to 1ms.
616 test.setAlarm(oneMs);
617 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
618 test.pollAndExpectCalls({"commit"})[0]->fulfill();
619 test.pollAndExpectCalls({});
620 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
621 
622 // Request handler run at 2ms. Expect cancellation with rescheduling.
623 auto armResult = test.actor.armAlarmHandler(twoMs, nullptr, testCurrentTime);
624 KJ_ASSERT(armResult.is<ActorSqlite::CancelAlarmHandler>());
625 auto cancelResult = kj::mv(armResult.get<ActorSqlite::CancelAlarmHandler>());
626 
627 // Expect rescheduling was requested and that returned promise resolves after fulfillment.
628 auto waitBeforeCancel = kj::mv(cancelResult.waitBeforeCancel);
629 auto rescheduleFulfiller = kj::mv(test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]);
630 KJ_ASSERT(!waitBeforeCancel.poll(test.ws));
631 rescheduleFulfiller->fulfill();
632 KJ_ASSERT(waitBeforeCancel.poll(test.ws));
633 waitBeforeCancel.wait(test.ws);
634}
635 
636KJ_TEST("tells alarm handler to reschedule when handler alarm is earlier than committed alarm") {
637 ActorSqliteTest test;
638 
639 // Initialize alarm state to 2ms.
640 test.setAlarm(twoMs);
641 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
642 test.pollAndExpectCalls({"commit"})[0]->fulfill();
643 test.pollAndExpectCalls({});
644 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
645 
646 // Expect that armAlarmHandler() tells caller to cancel after rescheduling completes.
647 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
648 KJ_ASSERT(armResult.is<ActorSqlite::CancelAlarmHandler>());
649 auto cancelResult = kj::mv(armResult.get<ActorSqlite::CancelAlarmHandler>());
650 
651 // Expect rescheduling was requested and that returned promise resolves after fulfillment.
652 auto waitBeforeCancel = kj::mv(cancelResult.waitBeforeCancel);
653 auto rescheduleFulfiller = kj::mv(test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]);
654 KJ_ASSERT(!waitBeforeCancel.poll(test.ws));
655 rescheduleFulfiller->fulfill();
656 KJ_ASSERT(waitBeforeCancel.poll(test.ws));
657 waitBeforeCancel.wait(test.ws);
658}
659 
660KJ_TEST("runs overdue alarm immediately when local alarm time is in the past") {
661 ActorSqliteTest test;
662 
663 // Initialize alarm state to 2ms.
664 test.setAlarm(twoMs);
665 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
666 test.pollAndExpectCalls({"commit"})[0]->fulfill();
667 test.pollAndExpectCalls({});
668 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
669 
670 // The local state says the alarm is due to fire at 2ms, but we're saying the AlarmManager has 1ms,
671 // usually this would result in a rescheduling of the alarm, but since our currentTime is 5ms, we
672 // will just run the alarm now since it's already overdue.
673 {
674 auto overdueCurrentTime = fiveMs;
675 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, overdueCurrentTime);
676 
677 // Should run the handler immediately instead of canceling/rescheduling.
678 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
679 }
680 
681 // commit and delete the alarm after we drop the alarm handler (this is a deferred delete).
682 test.pollAndExpectCalls({"commit"})[0]->fulfill();
683 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
684}
685 
686KJ_TEST("does not cancel handler when local db alarm state is later than scheduled alarm") {
687 ActorSqliteTest test;
688 
689 // Initialize alarm state to 1ms.
690 test.setAlarm(oneMs);
691 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
692 test.pollAndExpectCalls({"commit"})[0]->fulfill();
693 test.pollAndExpectCalls({});
694 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
695 
696 test.setAlarm(twoMs);
697 {
698 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
699 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
700 }
701 test.pollAndExpectCalls({"commit"})[0]->fulfill();
702 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
703}
704 
705KJ_TEST("does not cancel handler when local db alarm state is earlier than scheduled alarm") {
706 ActorSqliteTest test;
707 
708 // Initialize alarm state to 2ms.
709 test.setAlarm(twoMs);
710 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
711 test.pollAndExpectCalls({"commit"})[0]->fulfill();
712 test.pollAndExpectCalls({});
713 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
714 
715 test.setAlarm(oneMs);
716 {
717 auto armResult = test.actor.armAlarmHandler(twoMs, nullptr, testCurrentTime);
718 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
719 }
720 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
721 test.pollAndExpectCalls({"commit"})[0]->fulfill();
722}
723 
724KJ_TEST("getAlarm() returns null during handler") {
725 ActorSqliteTest test;
726 
727 // Initialize alarm state to 1ms.
728 test.setAlarm(oneMs);
729 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
730 test.pollAndExpectCalls({"commit"})[0]->fulfill();
731 test.pollAndExpectCalls({});
732 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
733 
734 {
735 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
736 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
737 test.pollAndExpectCalls({});
738 
739 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
740 }
741 test.pollAndExpectCalls({"commit"})[0]->fulfill();
742 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
743}
744 
745KJ_TEST("alarm handler handle clears alarm when dropped with no writes") {
746 ActorSqliteTest test;
747 
748 // Initialize alarm state to 1ms.
749 test.setAlarm(oneMs);
750 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
751 test.pollAndExpectCalls({"commit"})[0]->fulfill();
752 test.pollAndExpectCalls({});
753 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
754 
755 {
756 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
757 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
758 }
759 test.pollAndExpectCalls({"commit"})[0]->fulfill();
760 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
761 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
762}
763 
764KJ_TEST("alarm deleter does not clear alarm when dropped with writes") {
765 ActorSqliteTest test;
766 
767 // Initialize alarm state to 1ms.
768 test.setAlarm(oneMs);
769 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
770 test.pollAndExpectCalls({"commit"})[0]->fulfill();
771 test.pollAndExpectCalls({});
772 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
773 
774 {
775 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
776 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
777 test.setAlarm(twoMs);
778 }
779 test.pollAndExpectCalls({"commit"})[0]->fulfill();
780 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
781 
782 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
783}
784 
785KJ_TEST("can cancel deferred alarm deletion during handler") {
786 ActorSqliteTest test;
787 
788 // Initialize alarm state to 1ms.
789 test.setAlarm(oneMs);
790 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
791 test.pollAndExpectCalls({"commit"})[0]->fulfill();
792 test.pollAndExpectCalls({});
793 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
794 
795 {
796 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
797 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
798 test.actor.cancelDeferredAlarmDeletion();
799 }
800 
801 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
802}
803 
804KJ_TEST("canceling deferred alarm deletion outside handler has no effect") {
805 ActorSqliteTest test;
806 
807 // Initialize alarm state to 1ms.
808 test.setAlarm(oneMs);
809 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
810 test.pollAndExpectCalls({"commit"})[0]->fulfill();
811 test.pollAndExpectCalls({});
812 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
813 
814 {
815 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
816 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
817 }
818 test.pollAndExpectCalls({"commit"})[0]->fulfill();
819 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
820 
821 test.actor.cancelDeferredAlarmDeletion();
822 
823 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
824}
825 
826KJ_TEST("canceling deferred alarm deletion outside handler edge case") {
827 // Presumably harmless to cancel deletion if the client requests it after the handler ends but
828 // before the event loop runs the commit code? Trying to cancel deletion outside the handler is
829 // a bit of a contract violation anyway -- maybe we should just assert against it?
830 ActorSqliteTest test;
831 
832 // Initialize alarm state to 1ms.
833 test.setAlarm(oneMs);
834 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
835 test.pollAndExpectCalls({"commit"})[0]->fulfill();
836 test.pollAndExpectCalls({});
837 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
838 
839 {
840 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
841 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
842 }
843 test.actor.cancelDeferredAlarmDeletion();
844 test.pollAndExpectCalls({"commit"})[0]->fulfill();
845 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
846 
847 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
848}
849 
850KJ_TEST("canceling deferred alarm deletion is idempotent") {
851 // Not sure if important, but matches ActorCache behavior.
852 ActorSqliteTest test;
853 
854 // Initialize alarm state to 1ms.
855 test.setAlarm(oneMs);
856 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
857 test.pollAndExpectCalls({"commit"})[0]->fulfill();
858 test.pollAndExpectCalls({});
859 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
860 
861 {
862 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
863 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
864 test.actor.cancelDeferredAlarmDeletion();
865 test.actor.cancelDeferredAlarmDeletion();
866 }
867 
868 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
869}
870 
871KJ_TEST("alarm handler cleanup succeeds when output gate is broken") {
872 auto runWithSetup = [](auto testFunc) {
873 ActorSqliteTest test({.monitorOutputGate = false});
874 auto promise = test.gate.onBroken();
875 
876 // Initialize alarm state to 1ms.
877 test.setAlarm(oneMs);
878 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
879 test.pollAndExpectCalls({"commit"})[0]->fulfill();
880 test.pollAndExpectCalls({});
881 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
882 
883 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
884 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
885 auto deferredDelete = kj::mv(armResult.get<ActorSqlite::RunAlarmHandler>().deferredDelete);
886 
887 // Break gate
888 test.put("foo", "bar");
889 test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "a_rejected_commit"));
890 KJ_EXPECT_THROW_MESSAGE("a_rejected_commit", promise.wait(test.ws));
891 // Ensure taskFailed handler runs and notices brokenness:
892 test.ws.poll();
893 
894 testFunc(test, kj::mv(deferredDelete));
895 };
896 
897 // Here, we test that the deferred deleter destructor doesn't throw, both in the case when the
898 // caller cancels deletion and when it does not cancel it:
899 
900 runWithSetup([](ActorSqliteTest& test, kj::Own<void> deferredDelete) {
901 // In the case where the handler fails, we assume the caller will explicitly cancel deferred
902 // alarm deletion:
903 test.actor.cancelDeferredAlarmDeletion();
904 
905 // Dropping the DeferredAlarmDeleter should succeed -- that is, it should not throw here:
906 { auto drop = kj::mv(deferredDelete); }
907 });
908 
909 runWithSetup([](ActorSqliteTest& test, kj::Own<void> deferredDelete) {
910 // In the case where the handler succeeds, the caller will not cancel deferred deletion before
911 // dropping the DeferredAlarmDeleter. Dropping the DeferredAlarmDeleter should still succeed,
912 // even if the output gate happens to already be broken:
913 { auto drop = kj::mv(deferredDelete); }
914 });
915}
916 
917KJ_TEST("handler alarm is not deleted when commit fails") {
918 ActorSqliteTest test({.monitorOutputGate = false});
919 
920 auto promise = test.gate.onBroken();
921 
922 // Initialize alarm state to 1ms.
923 test.setAlarm(oneMs);
924 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
925 test.pollAndExpectCalls({"commit"})[0]->fulfill();
926 test.pollAndExpectCalls({});
927 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
928 
929 {
930 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
931 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
932 
933 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
934 }
935 test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "a_rejected_commit"));
936 
937 KJ_EXPECT_THROW_MESSAGE("a_rejected_commit", promise.wait(test.ws));
938}
939 
940KJ_TEST("setting earlier alarm persists alarm scheduling before db") {
941 ActorSqliteTest test;
942 
943 // Initialize alarm state to 2ms.
944 test.setAlarm(twoMs);
945 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
946 test.pollAndExpectCalls({"commit"})[0]->fulfill();
947 test.pollAndExpectCalls({});
948 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
949 
950 // Update alarm to be earlier. We expect the alarm scheduling to be persisted before the db.
951 test.setAlarm(oneMs);
952 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
953 test.pollAndExpectCalls({"commit"})[0]->fulfill();
954 
955 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
956}
957 
958KJ_TEST("setting later alarm persists db before alarm scheduling") {
959 ActorSqliteTest test;
960 
961 // Initialize alarm state to 1ms.
962 test.setAlarm(oneMs);
963 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
964 test.pollAndExpectCalls({"commit"})[0]->fulfill();
965 test.pollAndExpectCalls({});
966 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
967 
968 // Update alarm to be later. We expect the db to be persisted before the alarm scheduling.
969 test.setAlarm(twoMs);
970 test.pollAndExpectCalls({"commit"})[0]->fulfill();
971 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
972 
973 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
974}
975 
976KJ_TEST("multiple set-earlier in-flight alarms wait for earliest before committing db") {
977 ActorSqliteTest test;
978 
979 // Initialize alarm state to 5ms.
980 test.setAlarm(fiveMs);
981 test.pollAndExpectCalls({"scheduleRun(5ms)"})[0]->fulfill();
982 test.pollAndExpectCalls({"commit"})[0]->fulfill();
983 test.pollAndExpectCalls({});
984 KJ_ASSERT(expectSync(test.getAlarm()) == fiveMs);
985 
986 // Gate is not blocked.
987 auto gateWaitBefore = test.gate.wait(nullptr);
988 KJ_ASSERT(gateWaitBefore.poll(test.ws));
989 
990 // Update alarm to be earlier (4ms). We expect the alarm scheduling to start.
991 test.setAlarm(fourMs);
992 auto fulfiller4Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(4ms)"})[0]);
993 test.pollAndExpectCalls({});
994 KJ_ASSERT(expectSync(test.getAlarm()) == fourMs);
995 
996 // Gate as-of 4ms update is blocked.
997 auto gateWait4ms = test.gate.wait(nullptr);
998 KJ_ASSERT(!gateWait4ms.poll(test.ws));
999 
1000 // While 4ms scheduling request is in-flight, update alarm to be even earlier (3ms). We expect
1001 // the 4ms request to block the 3ms scheduling request.
1002 test.setAlarm(threeMs);
1003 test.pollAndExpectCalls({});
1004 KJ_ASSERT(expectSync(test.getAlarm()) == threeMs);
1005 
1006 // Gate as-of 3ms update is blocked.
1007 auto gateWait3ms = test.gate.wait(nullptr);
1008 KJ_ASSERT(!gateWait3ms.poll(test.ws));
1009 
1010 // Update alarm to be even earlier (2ms). We expect scheduling requests to still be blocked.
1011 test.setAlarm(twoMs);
1012 test.pollAndExpectCalls({});
1013 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
1014 
1015 // Gate as-of 2ms update is blocked.
1016 auto gateWait2ms = test.gate.wait(nullptr);
1017 KJ_ASSERT(!gateWait2ms.poll(test.ws));
1018 
1019 // Fulfill the 4ms request. We expect the 2ms scheduling to start, because that is the current
1020 // alarm value.
1021 fulfiller4Ms->fulfill();
1022 auto fulfiller2Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]);
1023 test.pollAndExpectCalls({});
1024 
1025 // While waiting for 2ms request, update alarm time to be 1ms. Expect scheduling to be blocked.
1026 test.setAlarm(oneMs);
1027 test.pollAndExpectCalls({});
1028 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1029 
1030 // Gate as-of 1ms update is blocked.
1031 auto gateWait1ms = test.gate.wait(nullptr);
1032 KJ_ASSERT(!gateWait1ms.poll(test.ws));
1033 
1034 // Fulfill the 2ms request. We expect the 1ms scheduling to start.
1035 fulfiller2Ms->fulfill();
1036 auto fulfiller1Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]);
1037 test.pollAndExpectCalls({});
1038 
1039 // Fulfill the 1ms request. We expect a single db commit to start (coalescing all previous db
1040 // commits together).
1041 fulfiller1Ms->fulfill();
1042 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
1043 test.pollAndExpectCalls({});
1044 
1045 // We expect all earlier gates to be blocked until commit completes.
1046 KJ_ASSERT(!gateWait4ms.poll(test.ws));
1047 KJ_ASSERT(!gateWait3ms.poll(test.ws));
1048 KJ_ASSERT(!gateWait2ms.poll(test.ws));
1049 KJ_ASSERT(!gateWait1ms.poll(test.ws));
1050 commitFulfiller->fulfill();
1051 KJ_ASSERT(gateWait4ms.poll(test.ws));
1052 KJ_ASSERT(gateWait3ms.poll(test.ws));
1053 KJ_ASSERT(gateWait2ms.poll(test.ws));
1054 KJ_ASSERT(gateWait1ms.poll(test.ws));
1055 
1056 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1057}
1058 
1059KJ_TEST("setting later alarm times does scheduling after db commit") {
1060 ActorSqliteTest test;
1061 
1062 // Initialize alarm state to 1ms.
1063 test.setAlarm(oneMs);
1064 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
1065 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1066 test.pollAndExpectCalls({});
1067 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1068 
1069 // Gate is not blocked.
1070 auto gateWaitBefore = test.gate.wait(nullptr);
1071 KJ_ASSERT(gateWaitBefore.poll(test.ws));
1072 
1073 // Set alarm to 2ms. Expect 2ms db commit to start.
1074 test.setAlarm(twoMs);
1075 auto commit2MsFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
1076 test.pollAndExpectCalls({});
1077 
1078 // Gate as-of 2ms update is blocked.
1079 auto gateWait2Ms = test.gate.wait(nullptr);
1080 KJ_ASSERT(!gateWait2Ms.poll(test.ws));
1081 
1082 // Set alarm to 3ms. Expect 3ms db commit to start. The 2ms scheduleRun will never happen now
1083 // that we've overwritten it while it was persisting to SQLite (before we send an update to the
1084 // alarm manager).
1085 test.setAlarm(threeMs);
1086 auto commit3MsFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
1087 test.pollAndExpectCalls({});
1088 
1089 // Gate as-of 3ms update is blocked.
1090 auto gateWait3Ms = test.gate.wait(nullptr);
1091 KJ_ASSERT(!gateWait3Ms.poll(test.ws));
1092 
1093 // Expect 2ms gate to be unblocked once the commit finishes, but don't expect a scheduleRun(2ms).
1094 KJ_ASSERT(!gateWait2Ms.poll(test.ws));
1095 commit2MsFulfiller->fulfill();
1096 KJ_ASSERT(gateWait2Ms.poll(test.ws));
1097 test.pollAndExpectCalls({});
1098 
1099 // Fulfill 3ms db commit. Expect 3ms alarm to be scheduled and 3ms gate to be unblocked.
1100 KJ_ASSERT(!gateWait3Ms.poll(test.ws));
1101 
1102 commit3MsFulfiller->fulfill();
1103 KJ_ASSERT(gateWait3Ms.poll(test.ws));
1104 
1105 auto fulfiller3Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(3ms)"})[0]);
1106 test.pollAndExpectCalls({});
1107 
1108 fulfiller3Ms->fulfill();
1109}
1110 
1111KJ_TEST("rejected move-earlier alarm scheduling request breaks gate") {
1112 ActorSqliteTest test({.monitorOutputGate = false});
1113 
1114 auto promise = test.gate.onBroken();
1115 
1116 test.setAlarm(oneMs);
1117 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->reject(
1118 KJ_EXCEPTION(FAILED, "a_rejected_scheduleRun"));
1119 
1120 KJ_EXPECT_THROW_MESSAGE("a_rejected_scheduleRun", promise.wait(test.ws));
1121}
1122 
1123KJ_TEST("rejected move-later alarm scheduling request does not break gate") {
1124 ActorSqliteTest test;
1125 
1126 // Initialize alarm state to 1ms.
1127 test.setAlarm(oneMs);
1128 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
1129 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1130 test.pollAndExpectCalls({});
1131 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1132 
1133 // Update alarm to be later. We expect the db to be persisted before the alarm scheduling.
1134 // We simulate a failure during the alarm rescheduling, but expect it to not break the output
1135 // gate.
1136 test.setAlarm(twoMs);
1137 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1138 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->reject(
1139 KJ_EXCEPTION(FAILED, "a_rejected_scheduleRun"));
1140 
1141 // Subsequent kv put succeeds. In an earlier version of the code, this failed, due to capturing
1142 // the scheduling failure as if it had broke the output gate, without actually breaking the
1143 // output gate.
1144 test.put("foo", "bar");
1145 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1146}
1147 
1148KJ_TEST("rapid move-later alarm changes coalesce into bounded scheduleRun calls") {
1149 // When many commits each move the alarm time later while a scheduleRun is already in-flight,
1150 // the scheduleLaterAlarm mechanism should coalesce them into at most one pending request,
1151 // rather than chaining N promises (one per commit).
1152 ActorSqliteTest test;
1153 
1154 // Initialize alarm state to 1ms.
1155 test.setAlarm(oneMs);
1156 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
1157 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1158 test.pollAndExpectCalls({});
1159 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1160 
1161 // Move alarm to 2ms. The db commit completes, triggering a post-commit scheduleRun(2ms)
1162 // since the alarm moved later.
1163 test.setAlarm(twoMs);
1164 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1165 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
1166 // The first move-later scheduleRun starts.
1167 auto fulfiller2Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]);
1168 
1169 // While 2ms scheduleRun is in-flight, move alarm to 3ms, 4ms, 5ms in rapid succession.
1170 // Each commit completes immediately but the scheduleRun for 2ms is still pending.
1171 // Only the final value (5ms) should be scheduled after the 2ms scheduleRun completes.
1172 test.setAlarm(threeMs);
1173 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1174 test.pollAndExpectCalls({}); // No new scheduleRun -- coalesced into pending.
1175 KJ_ASSERT(expectSync(test.getAlarm()) == threeMs);
1176 
1177 test.setAlarm(fourMs);
1178 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1179 test.pollAndExpectCalls({}); // No new scheduleRun -- coalesced into pending.
1180 KJ_ASSERT(expectSync(test.getAlarm()) == fourMs);
1181 
1182 test.setAlarm(fiveMs);
1183 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1184 test.pollAndExpectCalls({}); // No new scheduleRun -- coalesced into pending.
1185 KJ_ASSERT(expectSync(test.getAlarm()) == fiveMs);
1186 
1187 // Now fulfill the 2ms scheduleRun. The coalesced pending time (5ms) should be scheduled next.
1188 fulfiller2Ms->fulfill();
1189 auto fulfiller5Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(5ms)"})[0]);
1190 // Importantly, there is exactly one scheduleRun(5ms), not three separate calls for 3ms, 4ms, 5ms.
1191 
1192 fulfiller5Ms->fulfill();
1193 test.pollAndExpectCalls({});
1194 
1195 KJ_ASSERT(expectSync(test.getAlarm()) == fiveMs);
1196}
1197 
1198KJ_TEST("armAlarmHandler with coalesced pending alarms schedules reschedule exactly once") {
1199 // Verifies two properties:
1200 // 1. No duplicate scheduleRun(6ms): armAlarmHandler clears pendingLaterAlarmTime so the
1201 // FORK_A completion handler does not re-issue it.
1202 // 2. Future commits (10ms) that arrive after armAlarmHandler fires are correctly handled:
1203 // they queue in pendingLaterAlarmTime, get picked up by FORK_A's completion handler,
1204 // and chain off FORK_B (armAlarmHandler's fork) so the order is 3ms -> 6ms -> 10ms.
1205 ActorSqliteTest test;
1206 
1207 // Initialize alarm to 1ms and fully commit it so lastConfirmedAlarmDbState = 1ms.
1208 test.setAlarm(oneMs);
1209 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
1210 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1211 test.pollAndExpectCalls({});
1212 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1213 
1214 // Move alarm to 3ms -- scheduleRun(3ms) goes in-flight via scheduleLaterAlarm.
1215 // alarmLaterIsInFlight=true, alarmLaterInFlight=FORK_A.
1216 test.setAlarm(threeMs);
1217 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1218 auto fulfiller3Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(3ms)"})[0]);
1219 
1220 // While 3ms scheduleRun is in-flight, rapidly move to 4ms then 6ms.
1221 // Both coalesce into pendingLaterAlarmTime=6ms; no new scheduleRun issued.
1222 test.setAlarm(fourMs);
1223 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1224 test.pollAndExpectCalls({});
1225 
1226 test.setAlarm(sixMs);
1227 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1228 test.pollAndExpectCalls({});
1229 KJ_ASSERT(expectSync(test.getAlarm()) == sixMs);
1230 
1231 // The 1ms alarm fires. armAlarmHandler sees scheduledTime=1ms, localAlarmState=6ms.
1232 // willFireEarlier(1ms, 6ms) => reschedule-later path:
1233 // requestScheduledAlarm(6ms, FORK_A.addBranch()) called synchronously -> FORK_B
1234 // pendingLaterAlarmTime cleared to kj::none
1235 // alarmLaterInFlight = FORK_B
1236 // alarmLaterIsInFlight unchanged (still true, owned by FORK_A lifecycle)
1237 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
1238 KJ_ASSERT(armResult.is<ActorSqlite::CancelAlarmHandler>());
1239 auto& cancelResult = armResult.get<ActorSqlite::CancelAlarmHandler>();
1240 
1241 // scheduleRun(6ms) issued exactly once -- synchronously inside armAlarmHandler.
1242 auto fulfiller6Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(6ms)"})[0]);
1243 
1244 // Commit for 10ms arrives while scheduleRun(3ms) is still in-flight.
1245 // alarmLaterIsInFlight=true (FORK_A lifecycle still active) so 10ms is correctly
1246 // queued: pendingLaterAlarmTime=Some(10ms). FORK_B is referenced by alarmLaterInFlight.
1247 test.setAlarm(tenMs);
1248 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1249 test.pollAndExpectCalls({}); // No scheduleRun yet -- coalesced into pending.
1250 
1251 // Fulfill scheduleRun(3ms). FORK_A resolves. FORK_A completion handler fires:
1252 // alarmLaterIsInFlight=false
1253 // pendingLaterAlarmTime=Some(10ms) -> scheduleLaterAlarm(10ms)
1254 // requestScheduledAlarm(10ms, FORK_B.addBranch()) -> FORK_C chains off FORK_B
1255 // scheduleRun(10ms) issued synchronously
1256 // Importantly: scheduleRun(6ms) is NOT issued again here -- pendingLaterAlarmTime
1257 // held 10ms (not 6ms), because armAlarmHandler had already cleared the 6ms.
1258 fulfiller3Ms->fulfill();
1259 auto fulfiller10Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(10ms)"})[0]);
1260 
1261 // Fulfill scheduleRun(6ms). FORK_B resolves cleanly -- it has no completion handler,
1262 // so overwriting alarmLaterInFlight with FORK_C is safe: FORK_C captured a branch of
1263 // FORK_B as priorTask before the field was overwritten, keeping FORK_B alive. FORK_B
1264 // resolving propagates into FORK_C's priorTask silently. No new scheduleRun here.
1265 fulfiller6Ms->fulfill();
1266 KJ_ASSERT(cancelResult.waitBeforeCancel.poll(test.ws));
1267 test.pollAndExpectCalls({});
1268 
1269 // Fulfill scheduleRun(10ms). FORK_C resolves, its completion handler fires with no
1270 // pending times. Done.
1271 fulfiller10Ms->fulfill();
1272 test.pollAndExpectCalls({});
1273 KJ_ASSERT(expectSync(test.getAlarm()) == tenMs);
1274}
1275 
1276KJ_TEST("coalesced move-later followed by move-earlier does not race") {
1277 // Regression test for a race condition where a coalesced pendingLaterAlarmTime could
1278 // be drained concurrently with a move-earlier scheduleRun. The fix is that
1279 // startPrecommitAlarmScheduling() clears pendingLaterAlarmTime when setting up a
1280 // move-earlier, so the completion handler finds nothing to drain.
1281 //
1282 // Scenario: alarm at 1ms -> move to 5ms (later, in-flight) -> move to 10ms (coalesced)
1283 // -> move to 2ms (earlier). Without the fix, after the 5ms RPC completes, both
1284 // scheduleRun(10ms) and scheduleRun(2ms) would fire concurrently. With the fix,
1285 // only scheduleRun(2ms) fires because the coalesced 10ms was cleared.
1286 ActorSqliteTest test;
1287 
1288 uint activeRpcs = 0;
1289 uint maxConcurrentRpcs = 0;
1290 
1291 // Custom handler that respects priorTask ordering like the real alarm manager.
1292 // The real alarm manager awaits priorTask before sending its RPC; we replicate
1293 // that here and track concurrent calls.
1294 test.scheduleRunWithPriorHandler = [&](kj::Maybe<kj::Date> newAlarmTime,
1295 kj::Promise<void> priorTask) -> kj::Promise<void> {
1296 return priorTask.then([&, newAlarmTime]() mutable -> kj::Promise<void> {
1297 activeRpcs++;
1298 maxConcurrentRpcs = kj::max(maxConcurrentRpcs, activeRpcs);
1299 auto desc = newAlarmTime.map([](auto& t) {
1300 return kj::str("scheduleRun(", t, ")");
1301 }).orDefault(kj::str("scheduleRun(none)"));
1302 auto [promise, fulfiller] = kj::newPromiseAndFulfiller<void>();
1303 test.calls.add(ActorSqliteTest::Call{kj::mv(desc), kj::mv(fulfiller)});
1304 return promise.then([&]() { activeRpcs--; });
1305 });
1306 };
1307 
1308 // Poll event loop until at least `count` calls accumulate. With the
1309 // priorTask-respecting handler, calls take extra event loop turns to appear.
1310 auto drainCalls = [&](std::initializer_list<kj::StringPtr> expected, kj::StringPtr msg = ""_kj) {
1311 size_t need = expected.size();
1312 for (int i = 0; i < 100; i++) {
1313 test.ws.poll();
1314 if (need == 0 && i >= 10) break;
1315 if (need > 0 && test.calls.size() >= need) break;
1316 }
1317 auto callDescs = KJ_MAP(c, test.calls) { return kj::str(c.desc); };
1318 KJ_ASSERT(callDescs == kj::heapArray(expected), msg);
1319 auto fulfillers = KJ_MAP(c, test.calls) { return kj::mv(c.fulfiller); };
1320 test.calls.clear();
1321 return kj::mv(fulfillers);
1322 };
1323 
1324 // 1. Initialize alarm state to 1ms.
1325 test.setAlarm(oneMs);
1326 drainCalls({"scheduleRun(1ms)"})[0]->fulfill();
1327 drainCalls({"commit"})[0]->fulfill();
1328 drainCalls({});
1329 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1330 
1331 // 2. Move alarm to 5ms (later). The db commit completes, then scheduleRun(5ms)
1332 // fires post-commit via scheduleLaterAlarm. Hold the fulfiller to keep it in-flight.
1333 test.setAlarm(fiveMs);
1334 drainCalls({"commit"})[0]->fulfill();
1335 auto fulfiller5Ms = kj::mv(drainCalls({"scheduleRun(5ms)"})[0]);
1336 
1337 // 3. While scheduleRun(5ms) is in-flight, move alarm to 10ms (later).
1338 // Since alarmLaterIsInFlight is true, 10ms is coalesced into pendingLaterAlarmTime.
1339 test.setAlarm(tenMs);
1340 drainCalls({"commit"})[0]->fulfill();
1341 drainCalls({}); // No scheduleRun -- coalesced into pending.
1342 
1343 // 4. Move alarm earlier to 2ms. startPrecommitAlarmScheduling() clears
1344 // pendingLaterAlarmTime and calls requestScheduledAlarm(2ms, FORK_5.addBranch()).
1345 // The priorTask-respecting handler blocks until FORK_5 resolves.
1346 test.setAlarm(twoMs);
1347 drainCalls({}); // scheduleRun(2ms) blocked on priorTask.
1348 
1349 // 5. Fulfill scheduleRun(5ms). FORK_5 resolves:
1350 // - Completion handler: pendingLaterAlarmTime was cleared -> no drain, no-op.
1351 // - Move-earlier priorTask resolves -> scheduleRun(2ms) fires.
1352 fulfiller5Ms->fulfill();
1353 auto fulfiller2Ms = kj::mv(drainCalls(
1354 {"scheduleRun(2ms)"}, "expected only scheduleRun(2ms), no concurrent scheduleRun(10ms)")[0]);
1355 
1356 // Verify no concurrent RPCs occurred.
1357 KJ_ASSERT(maxConcurrentRpcs <= 1,
1358 "scheduleRun RPCs were sent concurrently -- "
1359 "the coalesced move-later raced with the move-earlier");
1360 
1361 // 6. Complete the move-earlier and its commit.
1362 fulfiller2Ms->fulfill();
1363 drainCalls({"commit"})[0]->fulfill();
1364 
1365 // Let any remaining completion handlers settle.
1366 for (int i = 0; i < 20; i++) test.ws.poll();
1367 
1368 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
1369}
1370 
1371KJ_TEST("an exception thrown during merged commits does not hang") {
1372 ActorSqliteTest test({.monitorOutputGate = false});
1373 
1374 auto promise = test.gate.onBroken();
1375 
1376 // Initialize alarm state to 5ms.
1377 test.setAlarm(fiveMs);
1378 test.pollAndExpectCalls({"scheduleRun(5ms)"})[0]->fulfill();
1379 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1380 test.pollAndExpectCalls({});
1381 KJ_ASSERT(expectSync(test.getAlarm()) == fiveMs);
1382 
1383 // Update alarm to be earlier (4ms). We expect the alarm scheduling to start.
1384 test.setAlarm(fourMs);
1385 auto fulfiller4Ms = kj::mv(test.pollAndExpectCalls({"scheduleRun(4ms)"})[0]);
1386 auto gateWait4ms = test.gate.wait(nullptr);
1387 
1388 // While 4ms scheduling request is in-flight, update alarm to be earlier (3ms). We expect
1389 // the two commit requests to merge and be blocked on the alarm scheduling request.
1390 test.setAlarm(threeMs);
1391 test.pollAndExpectCalls({});
1392 auto gateWait3ms = test.gate.wait(nullptr);
1393 
1394 // Reject the 4ms request. We expect both gate waiting promises to unblock with exceptions.
1395 KJ_ASSERT(!gateWait4ms.poll(test.ws));
1396 KJ_ASSERT(!gateWait3ms.poll(test.ws));
1397 fulfiller4Ms->reject(KJ_EXCEPTION(FAILED, "a_rejected_scheduleRun"));
1398 KJ_ASSERT(gateWait4ms.poll(test.ws));
1399 KJ_ASSERT(gateWait3ms.poll(test.ws));
1400 
1401 KJ_EXPECT_THROW_MESSAGE("a_rejected_scheduleRun", gateWait4ms.wait(test.ws));
1402 KJ_EXPECT_THROW_MESSAGE("a_rejected_scheduleRun", gateWait3ms.wait(test.ws));
1403 KJ_EXPECT_THROW_MESSAGE("a_rejected_scheduleRun", promise.wait(test.ws));
1404}
1405 
1406KJ_TEST("getAlarm/setAlarm check for brokenness") {
1407 ActorSqliteTest test({.monitorOutputGate = false});
1408 
1409 auto promise = test.gate.onBroken();
1410 
1411 // Break gate
1412 test.put("foo", "bar");
1413 test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "a_rejected_commit"));
1414 
1415 KJ_EXPECT_THROW_MESSAGE("a_rejected_commit", promise.wait(test.ws));
1416 
1417 // Apparently we don't actually set brokenness until the taskFailed handler runs, but presumably
1418 // this is OK?
1419 test.getAlarm();
1420 
1421 // Ensure taskFailed handler runs and notices brokenness:
1422 test.ws.poll();
1423 
1424 KJ_EXPECT_THROW_MESSAGE("a_rejected_commit", test.getAlarm());
1425 KJ_EXPECT_THROW_MESSAGE("a_rejected_commit", test.setAlarm(kj::none));
1426 test.pollAndExpectCalls({});
1427}
1428 
1429KJ_TEST("calling deleteAll() preserves alarm state if alarm is set") {
1430 ActorSqliteTest test;
1431 
1432 // Initialize alarm state to 1ms.
1433 test.setAlarm(oneMs);
1434 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
1435 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1436 test.pollAndExpectCalls({});
1437 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1438 
1439 {
1440 KJ_ASSERT(!test.actor.isCommitScheduled());
1441 ActorCache::DeleteAllResults results = test.actor.deleteAll({}, nullptr);
1442 KJ_ASSERT(test.actor.isCommitScheduled());
1443 KJ_ASSERT(results.backpressure == kj::none);
1444 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1445 
1446 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
1447 KJ_ASSERT(results.count.wait(test.ws) == 0);
1448 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1449 
1450 commitFulfiller->fulfill();
1451 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1452 
1453 test.pollAndExpectCalls({});
1454 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1455 }
1456 
1457 {
1458 // Should be fine to call deleteAll() a few times in succession, too:
1459 KJ_ASSERT(!test.actor.isCommitScheduled());
1460 ActorCache::DeleteAllResults results1 = test.actor.deleteAll({}, nullptr);
1461 ActorCache::DeleteAllResults results2 = test.actor.deleteAll({}, nullptr);
1462 KJ_ASSERT(test.actor.isCommitScheduled());
1463 KJ_ASSERT(results1.backpressure == kj::none);
1464 KJ_ASSERT(results2.backpressure == kj::none);
1465 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1466 
1467 // Presumably fine to be performing the alarm state restoration after each db reset:
1468 auto commitFulfillers = test.pollAndExpectCalls({"commit", "commit"});
1469 KJ_ASSERT(results1.count.wait(test.ws) == 0);
1470 KJ_ASSERT(results2.count.wait(test.ws) == 0);
1471 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1472 
1473 commitFulfillers[0]->fulfill();
1474 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1475 commitFulfillers[1]->fulfill();
1476 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1477 
1478 test.pollAndExpectCalls({});
1479 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1480 }
1481}
1482 
1483KJ_TEST("calling deleteAll() preserves alarm state if alarm is not set") {
1484 ActorSqliteTest test;
1485 
1486 // Initialize alarm state to empty value in metadata table.
1487 test.setAlarm(oneMs);
1488 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
1489 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1490 test.pollAndExpectCalls({});
1491 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1492 test.setAlarm(kj::none);
1493 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1494 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
1495 test.pollAndExpectCalls({});
1496 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1497 
1498 {
1499 KJ_ASSERT(!test.actor.isCommitScheduled());
1500 ActorCache::DeleteAllResults results = test.actor.deleteAll({}, nullptr);
1501 KJ_ASSERT(test.actor.isCommitScheduled());
1502 KJ_ASSERT(results.backpressure == kj::none);
1503 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1504 
1505 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
1506 KJ_ASSERT(results.count.wait(test.ws) == 0);
1507 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1508 
1509 commitFulfiller->fulfill();
1510 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1511 
1512 // We can also assert that we leave the database empty, in case that turns out to be useful later:
1513 auto q =
1514 test.db.run("SELECT name FROM sqlite_master WHERE type='table' AND name='_cf_METADATA'");
1515 KJ_ASSERT(q.isDone());
1516 }
1517 
1518 {
1519 // Should be fine to call deleteAll() a few times in succession, too:
1520 KJ_ASSERT(!test.actor.isCommitScheduled());
1521 ActorCache::DeleteAllResults results1 = test.actor.deleteAll({}, nullptr);
1522 ActorCache::DeleteAllResults results2 = test.actor.deleteAll({}, nullptr);
1523 KJ_ASSERT(test.actor.isCommitScheduled());
1524 KJ_ASSERT(results1.backpressure == kj::none);
1525 KJ_ASSERT(results2.backpressure == kj::none);
1526 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1527 
1528 // Presumably fine that the deletion commits coalesce:
1529 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
1530 KJ_ASSERT(results1.count.wait(test.ws) == 0);
1531 KJ_ASSERT(results2.count.wait(test.ws) == 0);
1532 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1533 
1534 commitFulfiller->fulfill();
1535 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1536 
1537 auto q =
1538 test.db.run("SELECT name FROM sqlite_master WHERE type='table' AND name='_cf_METADATA'");
1539 KJ_ASSERT(q.isDone());
1540 }
1541}
1542 
1543KJ_TEST("calling deleteAll() during an implicit transaction preserves alarm state") {
1544 ActorSqliteTest test;
1545 
1546 KJ_ASSERT(!test.actor.isCommitScheduled());
1547 
1548 // Initialize alarm state to 1ms.
1549 test.setAlarm(oneMs);
1550 
1551 ActorCache::DeleteAllResults results = test.actor.deleteAll({}, nullptr);
1552 KJ_ASSERT(test.actor.isCommitScheduled());
1553 KJ_ASSERT(results.backpressure == kj::none);
1554 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1555 
1556 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
1557 
1558 auto commitFulfiller = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
1559 KJ_ASSERT(results.count.wait(test.ws) == 0);
1560 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1561 
1562 commitFulfiller->fulfill();
1563 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1564 
1565 test.pollAndExpectCalls({});
1566 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1567}
1568 
1569KJ_TEST("deleteAll with deleteAlarm option deletes alarm") {
1570 // Tests that deleteAll() with deleteAlarm=true deletes the alarm along with KV data,
1571 // instead of preserving the alarm as it does by default.
1572 ActorSqliteTest test;
1573 
1574 // Initialize alarm state to 1ms.
1575 test.setAlarm(oneMs);
1576 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
1577 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1578 test.pollAndExpectCalls({});
1579 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1580 
1581 // Call deleteAll() with deleteAlarm=true.
1582 ActorCache::DeleteAllResults results = test.actor.deleteAll({}, nullptr, {.deleteAlarm = true});
1583 
1584 // The alarm should now be deleted.
1585 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1586 
1587 // Commit should include scheduling the alarm cancellation.
1588 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1589 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
1590 test.pollAndExpectCalls({});
1591 
1592 KJ_ASSERT(results.count.wait(test.ws) == 0);
1593 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1594}
1595 
1596KJ_TEST("deleteAll without deleteAlarm option preserves alarm") {
1597 // Tests that deleteAll() without deleteAlarm (the default) preserves the alarm,
1598 // which is the existing behavior.
1599 ActorSqliteTest test;
1600 
1601 // Initialize alarm state to 1ms.
1602 test.setAlarm(oneMs);
1603 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
1604 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1605 test.pollAndExpectCalls({});
1606 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1607 
1608 // Call deleteAll() without deleteAlarm (default behavior).
1609 ActorCache::DeleteAllResults results = test.actor.deleteAll({}, nullptr);
1610 
1611 // The alarm should be preserved.
1612 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1613 
1614 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1615 test.pollAndExpectCalls({});
1616 
1617 KJ_ASSERT(results.count.wait(test.ws) == 0);
1618 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1619}
1620 
1621KJ_TEST("deleteAll with deleteAlarm during alarm handler cancels deferred delete") {
1622 // Tests that calling deleteAll() with deleteAlarm=true while an alarm handler is running
1623 // correctly deletes the alarm and cancels the deferred alarm deletion (haveDeferredDelete).
1624 // When the handler's DeferredAlarmDeleter is dropped, it should NOT write a null alarm row
1625 // since deleteAll already handled the deletion.
1626 ActorSqliteTest test;
1627 
1628 // Initialize alarm state to 1ms.
1629 test.setAlarm(oneMs);
1630 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
1631 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1632 test.pollAndExpectCalls({});
1633 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1634 
1635 {
1636 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
1637 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
1638 
1639 // During the handler, getAlarm() should return none (deferred delete is active).
1640 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1641 
1642 // Call deleteAll() with deleteAlarm=true while the handler is running.
1643 auto results = test.actor.deleteAll({}, nullptr, {.deleteAlarm = true});
1644 
1645 // getAlarm() should still return none.
1646 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1647 
1648 // Drop the DeferredAlarmDeleter (simulating handler success). Since deleteAll already
1649 // cleared haveDeferredDelete, this should NOT write to the metadata table.
1650 }
1651 
1652 // The deleteAll commit should go through commitImpl(), which detects the alarm moved to none
1653 // and schedules the cancellation.
1654 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1655 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
1656 test.pollAndExpectCalls({});
1657 
1658 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1659}
1660 
1661KJ_TEST("deleteAll without deleteAlarm during alarm handler still has deferred delete") {
1662 // Tests that calling deleteAll() without deleteAlarm while an alarm handler is running
1663 // restores the alarm in metadata but leaves haveDeferredDelete active. When the handler
1664 // finishes, the deferred deletion deletes the restored alarm.
1665 ActorSqliteTest test;
1666 
1667 // Initialize alarm state to 1ms.
1668 test.setAlarm(oneMs);
1669 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
1670 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1671 test.pollAndExpectCalls({});
1672 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1673 
1674 {
1675 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
1676 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
1677 
1678 // During the handler, getAlarm() should return none (deferred delete is active).
1679 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1680 
1681 // Call deleteAll() without deleteAlarm while the handler is running.
1682 // This restores the alarm in metadata, but haveDeferredDelete is still true.
1683 test.actor.deleteAll({}, nullptr);
1684 
1685 // getAlarm() still returns none because haveDeferredDelete is still active.
1686 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1687 
1688 // Drop the DeferredAlarmDeleter (simulating handler success). This triggers
1689 // maybeDeleteDeferredAlarm() which deletes the restored alarm.
1690 }
1691 
1692 // The deleteAll commit goes first, then the deferred alarm deletion triggers its own commit
1693 // with alarm scheduling.
1694 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1695 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
1696 test.pollAndExpectCalls({});
1697 
1698 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1699}
1700 
1701KJ_TEST("deleteAll deleteAlarm does not schedule alarm cancellation if setAlarm interleaves") {
1702 ActorSqliteTest test;
1703 
1704 // Initialize alarm state to 1ms.
1705 test.setAlarm(oneMs);
1706 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
1707 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1708 test.pollAndExpectCalls({});
1709 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1710 
1711 // Start deleteAll with deleteAlarm=true and hold the commit.
1712 test.actor.deleteAll({}, nullptr, {.deleteAlarm = true});
1713 auto deleteAllCommit = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
1714 
1715 // While deleteAll commit is in-flight, set a later alarm.
1716 test.setAlarm(twoMs);
1717 auto setAlarmCommit = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
1718 test.pollAndExpectCalls({});
1719 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
1720 
1721 // Completing the deleteAll commit should NOT schedule a cancel because setAlarm interleaved.
1722 deleteAllCommit->fulfill();
1723 test.pollAndExpectCalls({});
1724 
1725 // Completing the setAlarm commit should schedule the new alarm time.
1726 setAlarmCommit->fulfill();
1727 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
1728 test.pollAndExpectCalls({});
1729 
1730 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
1731}
1732 
1733KJ_TEST("rolling back transaction leaves alarm in expected state") {
1734 ActorSqliteTest test;
1735 
1736 // Initialize alarm state to 2ms.
1737 test.setAlarm(twoMs);
1738 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
1739 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1740 test.pollAndExpectCalls({});
1741 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
1742 
1743 {
1744 auto txn = test.actor.startTransaction();
1745 KJ_ASSERT(expectSync(txn->getAlarm({})) == twoMs);
1746 txn->setAlarm(oneMs, {}, nullptr);
1747 KJ_ASSERT(expectSync(txn->getAlarm({})) == oneMs);
1748 // Dropping transaction without committing; should roll back.
1749 }
1750 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
1751}
1752 
1753KJ_TEST("rolling back transaction leaves deferred alarm deletion in expected state") {
1754 ActorSqliteTest test;
1755 
1756 // Initialize alarm state to 2ms.
1757 test.setAlarm(twoMs);
1758 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
1759 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1760 test.pollAndExpectCalls({});
1761 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
1762 
1763 {
1764 auto armResult = test.actor.armAlarmHandler(twoMs, nullptr, testCurrentTime);
1765 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
1766 
1767 auto txn = test.actor.startTransaction();
1768 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1769 test.setAlarm(oneMs);
1770 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1771 txn->rollback().wait(test.ws);
1772 
1773 // After rollback, getAlarm() still returns the deferred deletion result.
1774 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1775 
1776 // After rollback, no changes committed, no change in scheduled alarm.
1777 test.pollAndExpectCalls({});
1778 }
1779 
1780 // After handler, 2ms alarm is deleted.
1781 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1782 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
1783 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1784}
1785 
1786KJ_TEST("committing transaction leaves deferred alarm deletion in expected state") {
1787 ActorSqliteTest test;
1788 
1789 // Initialize alarm state to 2ms.
1790 test.setAlarm(twoMs);
1791 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
1792 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1793 test.pollAndExpectCalls({});
1794 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
1795 
1796 {
1797 auto armResult = test.actor.armAlarmHandler(twoMs, nullptr, testCurrentTime);
1798 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
1799 
1800 auto txn = test.actor.startTransaction();
1801 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1802 test.setAlarm(oneMs);
1803 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1804 txn->commit();
1805 
1806 // After commit, getAlarm() returns the committed value.
1807 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1808 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
1809 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1810 test.pollAndExpectCalls({});
1811 }
1812 
1813 // Alarm not deleted
1814 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1815}
1816 
1817KJ_TEST("rolling back nested transaction leaves deferred alarm deletion in expected state") {
1818 ActorSqliteTest test;
1819 
1820 // Initialize alarm state to 2ms.
1821 test.setAlarm(twoMs);
1822 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
1823 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1824 test.pollAndExpectCalls({});
1825 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
1826 
1827 {
1828 auto armResult = test.actor.armAlarmHandler(twoMs, nullptr, testCurrentTime);
1829 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
1830 
1831 auto txn1 = test.actor.startTransaction();
1832 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1833 {
1834 // Rolling back nested transaction change leaves deferred deletion in place.
1835 auto txn2 = test.actor.startTransaction();
1836 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1837 test.setAlarm(oneMs);
1838 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1839 txn2->rollback().wait(test.ws);
1840 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1841 }
1842 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1843 {
1844 // Committing nested transaction changes parent transaction state to dirty.
1845 auto txn3 = test.actor.startTransaction();
1846 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1847 test.setAlarm(oneMs);
1848 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1849 txn3->commit();
1850 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1851 }
1852 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1853 {
1854 // Nested transaction of dirty transaction is dirty, rollback has no effect.
1855 auto txn4 = test.actor.startTransaction();
1856 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1857 txn4->rollback().wait(test.ws);
1858 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1859 }
1860 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
1861 txn1->rollback().wait(test.ws);
1862 
1863 // After root transaction rollback, getAlarm() still returns the deferred deletion result.
1864 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1865 
1866 // After rollback, no changes committed, no change in scheduled alarm.
1867 test.pollAndExpectCalls({});
1868 }
1869 
1870 // After handler, 2ms alarm is deleted.
1871 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1872 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
1873 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
1874}
1875 
1876KJ_TEST("database write operations check for brokenness") {
1877 ActorSqliteTest test({.monitorOutputGate = false});
1878 
1879 auto promise = test.gate.onBroken();
1880 
1881 // Break gate
1882 test.put("foo", "bar");
1883 test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "a_rejected_commit"));
1884 
1885 KJ_EXPECT_THROW_MESSAGE("a_rejected_commit", promise.wait(test.ws));
1886 
1887 // We don't actually set ActorSqlite's brokenness until the taskFailed handler runs...
1888 test.ws.poll();
1889 
1890 // Try making a write operation to the database, expecting it to throw the broken message via
1891 // the onWrite handler:
1892 KJ_EXPECT_THROW_MESSAGE(
1893 "a_rejected_commit", test.db.run("CREATE TABLE IF NOT EXISTS counter (count INTEGER)"));
1894 test.pollAndExpectCalls({});
1895}
1896 
1897KJ_TEST("allowUnconfirmed put does not block output gate") {
1898 ActorSqliteTest test;
1899 
1900 // Gate is currently not blocked.
1901 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
1902 
1903 // Do an unconfirmed put
1904 test.put("foo", "bar", {.allowUnconfirmed = true});
1905 
1906 // Gate still isn't blocked, because we set `allowUnconfirmed`.
1907 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
1908 
1909 // Complete the transaction.
1910 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1911 
1912 // Gate should still not be blocked after commit completes
1913 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
1914 
1915 // Verify data was written
1916 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
1917}
1918 
1919KJ_TEST("confirmed put blocks output gate") {
1920 ActorSqliteTest test;
1921 
1922 // Gate is currently not blocked.
1923 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
1924 
1925 // Do a confirmed put (default behavior)
1926 test.put("foo", "bar", {.allowUnconfirmed = false});
1927 
1928 // Now it should be blocked.
1929 KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws));
1930 
1931 // Complete the transaction.
1932 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1933 
1934 // Gate should unblock after commit completes
1935 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
1936 
1937 // Verify data was written
1938 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
1939}
1940 
1941KJ_TEST("mixed confirmed and unconfirmed writes in same transaction use output gate") {
1942 ActorSqliteTest test;
1943 
1944 // Gate is currently not blocked.
1945 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
1946 
1947 // Do an unconfirmed put followed by a confirmed put in the same transaction batch
1948 test.put("foo", "bar", {.allowUnconfirmed = true});
1949 test.put("baz", "quux", {.allowUnconfirmed = false});
1950 
1951 // Since any write in the batch needs confirmation, the entire batch should use output gate
1952 KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws));
1953 
1954 // Complete the transaction.
1955 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1956 
1957 // Gate should unblock after commit completes
1958 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
1959 
1960 // Both writes should be committed
1961 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
1962 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("quux").asBytes());
1963}
1964 
1965KJ_TEST("allowUnconfirmed delete does not block output gate") {
1966 ActorSqliteTest test;
1967 
1968 // First set up some data
1969 test.put("foo", "bar");
1970 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1971 
1972 // Gate should be unblocked after setup
1973 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
1974 
1975 // Perform an unconfirmed delete - need to add delete helper method or use actor directly
1976 expectSync(test.actor.delete_(kj::str("foo"), {.allowUnconfirmed = true}, nullptr));
1977 
1978 // Gate still isn't blocked, because we set `allowUnconfirmed`.
1979 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
1980 
1981 // Complete the transaction.
1982 test.pollAndExpectCalls({"commit"})[0]->fulfill();
1983 
1984 // Gate should still not be blocked after commit completes
1985 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
1986 
1987 // Data should be deleted
1988 KJ_ASSERT(expectSync(test.get("foo")) == kj::none);
1989}
1990 
1991KJ_TEST("allowUnconfirmed putMultiple does not block output gate") {
1992 ActorSqliteTest test;
1993 
1994 // Gate should be unblocked at start
1995 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
1996 
1997 // Create multiple key-value pairs for the test
1998 kj::Vector<ActorCache::KeyValuePair> putKVs;
1999 putKVs.add(ActorCache::KeyValuePair{kj::str("foo"), kj::heapArray(kj::str("bar").asBytes())});
2000 putKVs.add(ActorCache::KeyValuePair{kj::str("baz"), kj::heapArray(kj::str("qux").asBytes())});
2001 putKVs.add(ActorCache::KeyValuePair{kj::str("key3"), kj::heapArray(kj::str("value3").asBytes())});
2002 
2003 // Perform an unconfirmed putMultiple within the implicit transaction
2004 test.putMultiple(putKVs.releaseAsArray(), {.allowUnconfirmed = true});
2005 
2006 // Gate still isn't blocked, because we set `allowUnconfirmed`.
2007 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2008 
2009 // Complete the transaction.
2010 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2011 
2012 // Gate should still not be blocked after commit completes
2013 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2014 
2015 // Verify all data was written correctly
2016 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
2017 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("qux").asBytes());
2018 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("key3"))) == kj::str("value3").asBytes());
2019}
2020 
2021KJ_TEST("allowUnconfirmed deleteMultiple does not block output gate") {
2022 ActorSqliteTest test;
2023 
2024 // First set up some data
2025 kj::Vector<ActorCache::KeyValuePair> putKVs;
2026 putKVs.add(ActorCache::KeyValuePair{kj::str("foo"), kj::heapArray(kj::str("bar").asBytes())});
2027 putKVs.add(ActorCache::KeyValuePair{kj::str("baz"), kj::heapArray(kj::str("qux").asBytes())});
2028 putKVs.add(ActorCache::KeyValuePair{kj::str("key3"), kj::heapArray(kj::str("value3").asBytes())});
2029 
2030 test.putMultiple(putKVs.releaseAsArray());
2031 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2032 
2033 // Gate should be unblocked after setup
2034 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2035 
2036 // Create array of keys to delete
2037 kj::Vector<kj::String> deleteKeys;
2038 deleteKeys.add(kj::str("foo"));
2039 deleteKeys.add(kj::str("baz"));
2040 deleteKeys.add(kj::str("key3"));
2041 
2042 // Perform an unconfirmed deleteMultiple
2043 KJ_EXPECT(expectSync(
2044 test.deleteMultiple(deleteKeys.releaseAsArray(), {.allowUnconfirmed = true})) == 3);
2045 
2046 // Gate still isn't blocked, because we set `allowUnconfirmed`.
2047 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2048 
2049 // Complete the transaction.
2050 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2051 
2052 // Gate should still not be blocked after commit completes
2053 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2054 
2055 // Verify all data was deleted
2056 KJ_ASSERT(expectSync(test.get("foo")) == kj::none);
2057 KJ_ASSERT(expectSync(test.get("baz")) == kj::none);
2058 KJ_ASSERT(expectSync(test.get("key3")) == kj::none);
2059}
2060 
2061KJ_TEST("unconfirmed write failure still breaks output gate") {
2062 ActorSqliteTest test({.monitorOutputGate = false});
2063 
2064 auto promise = test.gate.onBroken();
2065 
2066 // Do an unconfirmed put
2067 test.put("foo", "bar", {.allowUnconfirmed = true});
2068 
2069 // The output gate is not applied initially.
2070 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2071 KJ_ASSERT(!promise.poll(test.ws));
2072 
2073 // Reject the commit to simulate failure
2074 test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "flush failed hard"));
2075 
2076 // Gate should be broken due to commit failure
2077 KJ_EXPECT_THROW_MESSAGE("flush failed hard", promise.wait(test.ws));
2078}
2079 
2080KJ_TEST("Direct SQL queries are confirmed writes") {
2081 ActorSqliteTest test;
2082 
2083 // Gate is currently not blocked.
2084 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2085 
2086 auto& db = KJ_ASSERT_NONNULL(test.actor.getSqliteDatabase());
2087 
2088 db.run("CREATE TABLE myTable (i INTEGER PRIMARY KEY, s TEXT)");
2089 db.run("INSERT INTO myTable VALUES (1, \"a\")");
2090 
2091 // Now the gate should be blocked.
2092 KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws));
2093 
2094 // Complete the transaction.
2095 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2096 
2097 // Gate should unblock after commit completes
2098 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2099 
2100 // Make sure that the write actually succeeded.
2101 {
2102 auto query = db.run("SELECT * FROM myTable");
2103 KJ_ASSERT(!query.isDone());
2104 KJ_EXPECT(query.getInt64(0) == 1);
2105 KJ_EXPECT(query.getText(1) == "a");
2106 query.nextRow();
2107 KJ_ASSERT(query.isDone());
2108 }
2109}
2110 
2111KJ_TEST("An unconfirmed put followed by a direct SQL queries requires the output gate") {
2112 ActorSqliteTest test;
2113 
2114 // Gate is currently not blocked.
2115 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2116 
2117 test.put("foo", "bar", {.allowUnconfirmed = true});
2118 auto& db = KJ_ASSERT_NONNULL(test.actor.getSqliteDatabase());
2119 db.run("CREATE TABLE myTable (i INTEGER PRIMARY KEY, s TEXT)");
2120 db.run("INSERT INTO myTable VALUES (1, \"a\")");
2121 
2122 // Now the gate should be blocked.
2123 KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws));
2124 
2125 // Complete the transaction.
2126 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2127 
2128 // Gate should unblock after commit completes
2129 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2130 
2131 // Make sure that the write actually succeeded.
2132 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
2133 {
2134 auto query = db.run("SELECT * FROM myTable");
2135 KJ_ASSERT(!query.isDone());
2136 KJ_EXPECT(query.getInt64(0) == 1);
2137 KJ_EXPECT(query.getText(1) == "a");
2138 query.nextRow();
2139 KJ_ASSERT(query.isDone());
2140 }
2141}
2142 
2143KJ_TEST("sync() returns immediately when no writes are pending") {
2144 ActorSqliteTest test;
2145 
2146 // When there are no pending writes, sync() should return a resolved promise
2147 auto syncResult = test.sync();
2148 auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult));
2149 KJ_ASSERT(syncPromise.poll(test.ws));
2150}
2151 
2152KJ_TEST("sync() waits for confirmed writes to complete") {
2153 ActorSqliteTest test;
2154 
2155 // Do a confirmed write (default behavior)
2156 test.put("foo", "bar", {.allowUnconfirmed = false});
2157 
2158 // sync() should return a promise that blocks until the commit completes
2159 auto syncResult = test.sync();
2160 KJ_ASSERT(syncResult != kj::none);
2161 
2162 auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult));
2163 
2164 // The sync promise should not be ready yet
2165 KJ_ASSERT(!syncPromise.poll(test.ws));
2166 
2167 // Complete the commit
2168 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2169 
2170 // Now the sync promise should be ready
2171 KJ_ASSERT(syncPromise.poll(test.ws));
2172 syncPromise.wait(test.ws);
2173 
2174 // Verify data was written
2175 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
2176}
2177 
2178KJ_TEST("sync() waits for unconfirmed writes to complete") {
2179 ActorSqliteTest test;
2180 
2181 // Do an unconfirmed write
2182 test.put("foo", "bar", {.allowUnconfirmed = true});
2183 
2184 // sync() should still return a promise that blocks until the commit completes
2185 auto syncResult = test.sync();
2186 KJ_ASSERT(syncResult != kj::none);
2187 
2188 auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult));
2189 
2190 // The sync promise should not be ready yet
2191 KJ_ASSERT(!syncPromise.poll(test.ws));
2192 
2193 // Complete the commit
2194 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2195 
2196 // Now the sync promise should be ready
2197 KJ_ASSERT(syncPromise.poll(test.ws));
2198 syncPromise.wait(test.ws);
2199 
2200 // Verify data was written
2201 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
2202}
2203 
2204KJ_TEST("sync() waits for multiple unconfirmed writes in a row") {
2205 ActorSqliteTest test;
2206 
2207 // Do multiple unconfirmed writes - they should batch into a single transaction
2208 test.put("foo", "bar", {.allowUnconfirmed = true});
2209 test.put("baz", "qux", {.allowUnconfirmed = true});
2210 test.put("key3", "value3", {.allowUnconfirmed = true});
2211 
2212 // sync() should wait for the batched commit
2213 auto syncResult = test.sync();
2214 KJ_ASSERT(syncResult != kj::none);
2215 auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult));
2216 
2217 // The sync promise should not be ready yet
2218 KJ_ASSERT(!syncPromise.poll(test.ws));
2219 
2220 // Complete the single batched commit
2221 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2222 
2223 // Now the sync promise should be ready
2224 KJ_ASSERT(syncPromise.poll(test.ws));
2225 syncPromise.wait(test.ws);
2226 
2227 // Verify all writes were committed
2228 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
2229 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("qux").asBytes());
2230 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("key3"))) == kj::str("value3").asBytes());
2231}
2232 
2233KJ_TEST("sync() only waits for writes before it was called") {
2234 ActorSqliteTest test;
2235 
2236 // First write
2237 test.put("foo", "bar", {.allowUnconfirmed = true});
2238 
2239 // Call sync for the first write
2240 auto syncResult = test.sync();
2241 KJ_ASSERT(syncResult != kj::none);
2242 auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult));
2243 
2244 // Complete first commit
2245 auto firstCommit = kj::mv(test.pollAndExpectCalls({"commit"})[0]);
2246 firstCommit->fulfill();
2247 
2248 // First sync should complete
2249 KJ_ASSERT(syncPromise.poll(test.ws));
2250 syncPromise.wait(test.ws);
2251 
2252 // Second write after sync was called
2253 test.put("baz", "qux", {.allowUnconfirmed = true});
2254 
2255 // The original sync should still be complete (doesn't wait for new write)
2256 // To verify this, let's get a new sync that should wait for the second write
2257 auto syncResult2 = test.sync();
2258 KJ_ASSERT(syncResult2 != kj::none);
2259 auto syncPromise2 = kj::mv(KJ_ASSERT_NONNULL(syncResult2));
2260 
2261 // Second sync should not be ready
2262 KJ_ASSERT(!syncPromise2.poll(test.ws));
2263 
2264 // Complete second commit
2265 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2266 
2267 // Now second sync should complete
2268 KJ_ASSERT(syncPromise2.poll(test.ws));
2269 syncPromise2.wait(test.ws);
2270}
2271 
2272KJ_TEST("sync() propagates commit errors") {
2273 ActorSqliteTest test({.monitorOutputGate = false});
2274 
2275 auto promise = test.gate.onBroken();
2276 
2277 // Do an unconfirmed write
2278 test.put("foo", "bar", {.allowUnconfirmed = true});
2279 
2280 // Call sync
2281 auto syncResult = test.sync();
2282 KJ_ASSERT(syncResult != kj::none);
2283 auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult));
2284 
2285 // The sync promise should not be ready yet
2286 KJ_ASSERT(!syncPromise.poll(test.ws));
2287 
2288 // Reject the commit to simulate failure
2289 test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "commit failed"));
2290 
2291 // sync promise should become ready with an exception
2292 KJ_ASSERT(syncPromise.poll(test.ws));
2293 KJ_EXPECT_THROW_MESSAGE("commit failed", syncPromise.wait(test.ws));
2294 
2295 // Gate should also be broken
2296 KJ_EXPECT_THROW_MESSAGE("commit failed", promise.wait(test.ws));
2297}
2298 
2299KJ_TEST("sync() with mixed confirmed and unconfirmed writes") {
2300 ActorSqliteTest test;
2301 
2302 // Do an unconfirmed write followed by a confirmed write
2303 test.put("foo", "bar", {.allowUnconfirmed = true});
2304 test.put("baz", "qux", {.allowUnconfirmed = false});
2305 
2306 // sync() should wait for both writes
2307 auto syncResult = test.sync();
2308 KJ_ASSERT(syncResult != kj::none);
2309 auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult));
2310 
2311 // The sync promise should not be ready yet
2312 KJ_ASSERT(!syncPromise.poll(test.ws));
2313 
2314 // Complete the commit (both writes are in the same transaction)
2315 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2316 
2317 // Now the sync promise should be ready
2318 KJ_ASSERT(syncPromise.poll(test.ws));
2319 syncPromise.wait(test.ws);
2320 
2321 // Both writes should be committed
2322 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
2323 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("qux").asBytes());
2324}
2325 
2326KJ_TEST("multiple sync() calls for same commit") {
2327 ActorSqliteTest test;
2328 
2329 // Do a write
2330 test.put("foo", "bar", {.allowUnconfirmed = true});
2331 
2332 // Call sync multiple times - they should all wait for the same commit
2333 auto syncResult1 = test.sync();
2334 auto syncResult2 = test.sync();
2335 auto syncResult3 = test.sync();
2336 
2337 KJ_ASSERT(syncResult1 != kj::none);
2338 KJ_ASSERT(syncResult2 != kj::none);
2339 KJ_ASSERT(syncResult3 != kj::none);
2340 
2341 auto syncPromise1 = kj::mv(KJ_ASSERT_NONNULL(syncResult1));
2342 auto syncPromise2 = kj::mv(KJ_ASSERT_NONNULL(syncResult2));
2343 auto syncPromise3 = kj::mv(KJ_ASSERT_NONNULL(syncResult3));
2344 
2345 // None should be ready yet
2346 KJ_ASSERT(!syncPromise1.poll(test.ws));
2347 KJ_ASSERT(!syncPromise2.poll(test.ws));
2348 KJ_ASSERT(!syncPromise3.poll(test.ws));
2349 
2350 // Complete the commit
2351 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2352 
2353 // All sync promises should become ready
2354 KJ_ASSERT(syncPromise1.poll(test.ws));
2355 KJ_ASSERT(syncPromise2.poll(test.ws));
2356 KJ_ASSERT(syncPromise3.poll(test.ws));
2357 
2358 syncPromise1.wait(test.ws);
2359 syncPromise2.wait(test.ws);
2360 syncPromise3.wait(test.ws);
2361}
2362 
2363KJ_TEST("allowUnconfirmed setAlarm does not block output gate") {
2364 ActorSqliteTest test;
2365 
2366 // Gate is currently not blocked.
2367 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2368 
2369 // Do an unconfirmed setAlarm
2370 test.setAlarm(oneMs, {.allowUnconfirmed = true});
2371 
2372 // Gate still isn't blocked, because we set `allowUnconfirmed`.
2373 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2374 
2375 // Complete the transaction - alarm scheduling happens before commit
2376 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
2377 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2378 
2379 // Gate should still not be blocked after commit completes
2380 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2381 
2382 // Verify alarm was set
2383 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
2384}
2385 
2386KJ_TEST("confirmed setAlarm blocks output gate") {
2387 ActorSqliteTest test;
2388 
2389 // Gate is currently not blocked.
2390 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2391 
2392 // Do a confirmed setAlarm (default behavior)
2393 test.setAlarm(oneMs, {.allowUnconfirmed = false});
2394 
2395 // Gate should be blocked after scheduling starts
2396 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
2397 KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws));
2398 
2399 // Complete the transaction.
2400 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2401 
2402 // Gate should unblock after commit completes
2403 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2404 
2405 // Verify alarm was set
2406 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
2407}
2408 
2409KJ_TEST("allowUnconfirmed setAlarm then confirmed put uses output gate") {
2410 ActorSqliteTest test;
2411 
2412 // Gate is currently not blocked.
2413 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2414 
2415 // Do an unconfirmed setAlarm followed by a confirmed put in the same transaction batch
2416 test.setAlarm(oneMs, {.allowUnconfirmed = true});
2417 test.put("foo", "bar", {.allowUnconfirmed = false});
2418 
2419 // Since any write in the batch needs confirmation, the entire batch should use output gate
2420 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
2421 KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws));
2422 
2423 // Complete the transaction.
2424 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2425 
2426 // Gate should unblock after commit completes
2427 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2428 
2429 // Both operations should be committed
2430 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
2431 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
2432}
2433 
2434KJ_TEST("allowUnconfirmed setAlarm with storage ops") {
2435 ActorSqliteTest test;
2436 
2437 // Gate is currently not blocked.
2438 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2439 
2440 // Do unconfirmed setAlarm with unconfirmed storage writes
2441 test.setAlarm(oneMs, {.allowUnconfirmed = true});
2442 test.put("foo", "bar", {.allowUnconfirmed = true});
2443 
2444 // Gate still isn't blocked since both operations are unconfirmed
2445 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2446 
2447 // Complete the transaction
2448 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
2449 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2450 
2451 // Gate should still not be blocked
2452 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2453 
2454 // Verify both alarm and storage writes committed
2455 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
2456 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
2457}
2458 
2459KJ_TEST("allowUnconfirmed setAlarm updating existing alarm") {
2460 ActorSqliteTest test;
2461 
2462 // Initialize alarm state to 2ms.
2463 test.setAlarm(twoMs);
2464 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
2465 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2466 test.pollAndExpectCalls({});
2467 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
2468 
2469 // Gate should be unblocked after setup
2470 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2471 
2472 // Update alarm to earlier time with allowUnconfirmed
2473 test.setAlarm(oneMs, {.allowUnconfirmed = true});
2474 
2475 // Gate still isn't blocked
2476 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2477 
2478 // Complete the transaction - when moving alarm earlier, schedule happens first
2479 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
2480 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2481 
2482 // Gate should still not be blocked
2483 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2484 
2485 // Verify alarm was updated
2486 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
2487}
2488 
2489KJ_TEST("allowUnconfirmed setAlarm to later time") {
2490 ActorSqliteTest test;
2491 
2492 // Initialize alarm state to 1ms.
2493 test.setAlarm(oneMs);
2494 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
2495 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2496 test.pollAndExpectCalls({});
2497 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
2498 
2499 // Gate should be unblocked after setup
2500 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2501 
2502 // Update alarm to later time with allowUnconfirmed
2503 test.setAlarm(twoMs, {.allowUnconfirmed = true});
2504 
2505 // Gate still isn't blocked
2506 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2507 
2508 // Complete the transaction - when moving alarm later, commit happens first
2509 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2510 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
2511 
2512 // Gate should still not be blocked
2513 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2514 
2515 // Verify alarm was updated
2516 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
2517}
2518 
2519KJ_TEST("allowUnconfirmed setAlarm to clear alarm") {
2520 ActorSqliteTest test;
2521 
2522 // Initialize alarm state to 1ms.
2523 test.setAlarm(oneMs);
2524 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
2525 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2526 test.pollAndExpectCalls({});
2527 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
2528 
2529 // Gate should be unblocked after setup
2530 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2531 
2532 // Clear alarm with allowUnconfirmed
2533 test.setAlarm(kj::none, {.allowUnconfirmed = true});
2534 
2535 // Gate still isn't blocked
2536 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2537 
2538 // Complete the transaction
2539 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2540 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
2541 
2542 // Gate should still not be blocked
2543 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2544 
2545 // Verify alarm was cleared
2546 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
2547}
2548 
2549KJ_TEST("unconfirmed setAlarm failure still breaks output gate") {
2550 ActorSqliteTest test({.monitorOutputGate = false});
2551 
2552 auto promise = test.gate.onBroken();
2553 
2554 // Do an unconfirmed setAlarm
2555 test.setAlarm(oneMs, {.allowUnconfirmed = true});
2556 
2557 // The output gate is not applied initially.
2558 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2559 KJ_ASSERT(!promise.poll(test.ws));
2560 
2561 // Fulfill scheduleRun but reject the commit to simulate failure
2562 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
2563 test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "alarm commit failed"));
2564 
2565 // Gate should be broken due to commit failure
2566 KJ_EXPECT_THROW_MESSAGE("alarm commit failed", promise.wait(test.ws));
2567}
2568 
2569KJ_TEST("sync() throws after critical error in explicit transaction") {
2570 ActorSqliteTest test({.monitorOutputGate = false});
2571 auto heapLimit = [&]() {
2572 auto row = test.db.run("PRAGMA hard_heap_limit");
2573 return row.getInt(0);
2574 }();
2575 KJ_DBG(heapLimit);
2576 KJ_DEFER(sqlite3_hard_heap_limit64(heapLimit););
2577 
2578 // Start an explicit transaction
2579 auto txn = test.actor.startTransaction();
2580 
2581 // Do a write within the transaction
2582 txn->put(kj::str("foo"), kj::heapArray(kj::str("bar").asBytes()), {}, nullptr);
2583 
2584 // Trigger a critical error using SQLITE_NOMEM by setting a very low heap limit
2585 // and then trying to insert a large value.
2586 try {
2587 // Set SQLite's memory limit very low to trigger SQLITE_NOMEM
2588 test.db.run("PRAGMA hard_heap_limit=8192"); // 8KB limit
2589 
2590 // Create data that will exceed the memory limit
2591 auto largeData = kj::heapArray<byte>(50000, 'X'); // 50KB
2592 
2593 // This should trigger SQLITE_NOMEM, causing SQLite to auto-rollback the transaction,
2594 // which will trigger the critical error handler.
2595 //
2596 // We have to copy it again in order to convert to a Array<const byte> from an Array<byte>.
2597 txn->put(kj::str("large_key"), kj::heapArray<const byte>(largeData.asBytes()), /*options=*/{},
2598 nullptr);
2599 KJ_FAIL_ASSERT("Query should have failed with SQLITE_NOMEM");
2600 } catch (kj::Exception& e) {
2601 // Expected: out of memory error. We catch and ignore this to continue the test.
2602 KJ_ASSERT(e.getDescription().contains("out of memory"));
2603 }
2604 
2605 // sync() should also throw an exception because the storage is now broken
2606 auto syncResult = test.sync();
2607 KJ_ASSERT(syncResult != kj::none);
2608 auto syncPromise = kj::mv(KJ_ASSERT_NONNULL(syncResult));
2609 KJ_EXPECT_THROW_MESSAGE("broken", syncPromise.wait(test.ws));
2610 
2611 // The transaction is now in a broken state due to the critical error.
2612 // Attempting to commit should fail.
2613 KJ_EXPECT_THROW_MESSAGE("broken", txn->commit());
2614}
2615 
2616KJ_TEST("allowUnconfirmed put in explicit transaction does not block output gate") {
2617 ActorSqliteTest test;
2618 
2619 // Gate is currently not blocked.
2620 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2621 
2622 // Start an explicit transaction
2623 auto txn = test.actor.startTransaction();
2624 
2625 // Do an unconfirmed put within the transaction
2626 txn->put(
2627 kj::str("foo"), kj::heapArray(kj::str("bar").asBytes()), {.allowUnconfirmed = true}, nullptr);
2628 
2629 // Gate still isn't blocked during the transaction, because we set `allowUnconfirmed`.
2630 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2631 
2632 // Commit the transaction
2633 txn->commit();
2634 
2635 // Gate should still not be blocked during commit because all writes were unconfirmed
2636 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2637 
2638 // Complete the commit
2639 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2640 
2641 // Gate should still not be blocked after commit completes
2642 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2643 
2644 // Verify data was written
2645 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
2646}
2647 
2648KJ_TEST("confirmed put in explicit transaction blocks output gate on commit") {
2649 ActorSqliteTest test;
2650 
2651 // Gate is currently not blocked.
2652 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2653 
2654 // Start an explicit transaction
2655 auto txn = test.actor.startTransaction();
2656 
2657 // Do a confirmed put (default behavior)
2658 txn->put(kj::str("foo"), kj::heapArray(kj::str("bar").asBytes()), {.allowUnconfirmed = false},
2659 nullptr);
2660 
2661 // Gate should still not be blocked during the transaction - explicit txns only lock on commit
2662 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2663 
2664 // Commit the transaction
2665 txn->commit();
2666 
2667 // Now the gate should be blocked because we're committing a confirmed write
2668 KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws));
2669 
2670 // Complete the commit
2671 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2672 
2673 // Gate should unblock after commit completes
2674 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2675 
2676 // Verify data was written
2677 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
2678}
2679 
2680KJ_TEST("mixed confirmed and unconfirmed puts in explicit transaction use output gate") {
2681 ActorSqliteTest test;
2682 
2683 // Gate is currently not blocked.
2684 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2685 
2686 // Start an explicit transaction
2687 auto txn = test.actor.startTransaction();
2688 
2689 // Do an unconfirmed put followed by a confirmed put
2690 txn->put(
2691 kj::str("foo"), kj::heapArray(kj::str("bar").asBytes()), {.allowUnconfirmed = true}, nullptr);
2692 txn->put(kj::str("baz"), kj::heapArray(kj::str("quux").asBytes()), {.allowUnconfirmed = false},
2693 nullptr);
2694 
2695 // Gate should still not be blocked during the transaction
2696 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2697 
2698 // Commit the transaction
2699 txn->commit();
2700 
2701 // Since any write in the transaction needs confirmation, commit should use output gate
2702 KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws));
2703 
2704 // Complete the commit
2705 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2706 
2707 // Gate should unblock after commit completes
2708 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2709 
2710 // Both writes should be committed
2711 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
2712 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("quux").asBytes());
2713}
2714 
2715KJ_TEST("allowUnconfirmed delete in explicit transaction does not block output gate") {
2716 ActorSqliteTest test;
2717 
2718 // First set up some data
2719 test.put("foo", "bar");
2720 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2721 
2722 // Gate should be unblocked after setup
2723 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2724 
2725 // Start an explicit transaction
2726 auto txn = test.actor.startTransaction();
2727 
2728 // Perform an unconfirmed delete
2729 expectSync(txn->delete_(kj::str("foo"), {.allowUnconfirmed = true}, nullptr));
2730 
2731 // Gate still isn't blocked during the transaction
2732 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2733 
2734 // Commit the transaction
2735 txn->commit();
2736 
2737 // Gate should still not be blocked during commit because the delete was unconfirmed
2738 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2739 
2740 // Complete the commit
2741 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2742 
2743 // Gate should still not be blocked after commit completes
2744 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2745 
2746 // Verify data was deleted
2747 KJ_ASSERT(expectSync(test.get("foo")) == kj::none);
2748}
2749 
2750KJ_TEST("allowUnconfirmed putMultiple in explicit transaction does not block output gate") {
2751 ActorSqliteTest test;
2752 
2753 // Gate is currently not blocked.
2754 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2755 
2756 // Start an explicit transaction
2757 auto txn = test.actor.startTransaction();
2758 
2759 // Do an unconfirmed putMultiple
2760 auto pairs = kj::heapArrayBuilder<ActorCacheOps::KeyValuePair>(2);
2761 pairs.add(ActorCacheOps::KeyValuePair{kj::str("foo"), kj::heapArray(kj::str("bar").asBytes())});
2762 pairs.add(ActorCacheOps::KeyValuePair{kj::str("baz"), kj::heapArray(kj::str("quux").asBytes())});
2763 txn->put(pairs.finish(), {.allowUnconfirmed = true}, nullptr);
2764 
2765 // Gate still isn't blocked during the transaction
2766 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2767 
2768 // Commit the transaction
2769 txn->commit();
2770 
2771 // Gate should still not be blocked during commit
2772 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2773 
2774 // Complete the commit
2775 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2776 
2777 // Gate should still not be blocked after commit completes
2778 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2779 
2780 // Verify data was written
2781 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("foo"))) == kj::str("bar").asBytes());
2782 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("baz"))) == kj::str("quux").asBytes());
2783}
2784 
2785KJ_TEST("allowUnconfirmed deleteMultiple in explicit transaction does not block output gate") {
2786 ActorSqliteTest test;
2787 
2788 // First set up some data
2789 test.put("foo", "bar");
2790 test.put("baz", "quux");
2791 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2792 
2793 // Gate should be unblocked after setup
2794 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2795 
2796 // Start an explicit transaction
2797 auto txn = test.actor.startTransaction();
2798 
2799 // Perform an unconfirmed deleteMultiple
2800 auto keys = kj::heapArrayBuilder<ActorCacheOps::Key>(2);
2801 keys.add(kj::str("foo"));
2802 keys.add(kj::str("baz"));
2803 expectSync(txn->delete_(keys.finish(), {.allowUnconfirmed = true}, nullptr));
2804 
2805 // Gate still isn't blocked during the transaction
2806 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2807 
2808 // Commit the transaction
2809 txn->commit();
2810 
2811 // Gate should still not be blocked during commit
2812 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2813 
2814 // Complete the commit
2815 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2816 
2817 // Gate should still not be blocked after commit completes
2818 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2819 
2820 // Verify data was deleted
2821 KJ_ASSERT(expectSync(test.get("foo")) == kj::none);
2822 KJ_ASSERT(expectSync(test.get("baz")) == kj::none);
2823}
2824 
2825KJ_TEST("allowUnconfirmed setAlarm in explicit transaction does not block output gate") {
2826 ActorSqliteTest test;
2827 
2828 // Gate is currently not blocked.
2829 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2830 
2831 // Start an explicit transaction
2832 auto txn = test.actor.startTransaction();
2833 
2834 // Set an alarm with allowUnconfirmed
2835 txn->setAlarm(oneMs, {.allowUnconfirmed = true}, nullptr);
2836 
2837 // Gate still isn't blocked during the transaction
2838 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2839 
2840 // Commit the transaction
2841 txn->commit();
2842 
2843 // Gate should still not be blocked during commit
2844 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2845 
2846 // Complete the scheduleRun and commit
2847 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
2848 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2849 
2850 // Gate should still not be blocked after commit completes
2851 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2852 
2853 // Verify alarm was set
2854 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.getAlarm())) == oneMs);
2855}
2856 
2857KJ_TEST("nested transaction: unconfirmed child commit does not block output gate") {
2858 ActorSqliteTest test;
2859 
2860 // Gate is currently not blocked.
2861 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2862 
2863 // Start a parent transaction
2864 auto parentTxn = test.actor.startTransaction();
2865 
2866 // Do an unconfirmed put in the parent
2867 parentTxn->put(kj::str("parent"), kj::heapArray(kj::str("data").asBytes()),
2868 {.allowUnconfirmed = true}, nullptr);
2869 
2870 {
2871 // Start a nested child transaction
2872 auto childTxn = test.actor.startTransaction();
2873 
2874 // Do an unconfirmed put in the child
2875 childTxn->put(kj::str("child"), kj::heapArray(kj::str("data").asBytes()),
2876 {.allowUnconfirmed = true}, nullptr);
2877 
2878 // Gate still isn't blocked
2879 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2880 
2881 // Commit the child transaction
2882 childTxn->commit();
2883 }
2884 
2885 // Gate should still not be blocked after child commit
2886 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2887 
2888 // Commit the parent transaction
2889 parentTxn->commit();
2890 
2891 // Gate should still not be blocked during parent commit because all writes were unconfirmed
2892 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2893 
2894 // Complete the commit
2895 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2896 
2897 // Gate should still not be blocked after commit completes
2898 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2899 
2900 // Verify both writes were committed
2901 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("parent"))) == kj::str("data").asBytes());
2902 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("child"))) == kj::str("data").asBytes());
2903}
2904 
2905KJ_TEST("nested transaction: confirmed child propagates to parent commit") {
2906 ActorSqliteTest test;
2907 
2908 // Gate is currently not blocked.
2909 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2910 
2911 // Start a parent transaction with unconfirmed write
2912 auto parentTxn = test.actor.startTransaction();
2913 parentTxn->put(kj::str("parent"), kj::heapArray(kj::str("data").asBytes()),
2914 {.allowUnconfirmed = true}, nullptr);
2915 
2916 {
2917 // Start a nested child transaction
2918 auto childTxn = test.actor.startTransaction();
2919 
2920 // Do a confirmed put in the child
2921 childTxn->put(kj::str("child"), kj::heapArray(kj::str("data").asBytes()),
2922 {.allowUnconfirmed = false}, nullptr);
2923 
2924 // Gate still isn't blocked during the transaction
2925 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2926 
2927 // Commit the child transaction - this should propagate someWriteConfirmed to parent
2928 childTxn->commit();
2929 }
2930 
2931 // Gate should still not be blocked after child commit (no real commit yet)
2932 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2933 
2934 // Commit the parent transaction
2935 parentTxn->commit();
2936 
2937 // Now the gate should be blocked because the child had a confirmed write
2938 KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws));
2939 
2940 // Complete the commit
2941 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2942 
2943 // Gate should unblock after commit completes
2944 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2945 
2946 // Verify both writes were committed
2947 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("parent"))) == kj::str("data").asBytes());
2948 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("child"))) == kj::str("data").asBytes());
2949}
2950 
2951KJ_TEST("nested transaction: confirmed parent with unconfirmed child blocks output gate") {
2952 ActorSqliteTest test;
2953 
2954 // Gate is currently not blocked.
2955 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2956 
2957 // Start a parent transaction with confirmed write
2958 auto parentTxn = test.actor.startTransaction();
2959 parentTxn->put(kj::str("parent"), kj::heapArray(kj::str("data").asBytes()),
2960 {.allowUnconfirmed = false}, nullptr);
2961 
2962 {
2963 // Start a nested child transaction
2964 auto childTxn = test.actor.startTransaction();
2965 
2966 // Do an unconfirmed put in the child
2967 childTxn->put(kj::str("child"), kj::heapArray(kj::str("data").asBytes()),
2968 {.allowUnconfirmed = true}, nullptr);
2969 
2970 // Gate still isn't blocked during the transaction
2971 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2972 
2973 // Commit the child transaction
2974 childTxn->commit();
2975 }
2976 
2977 // Gate should still not be blocked after child commit
2978 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2979 
2980 // Commit the parent transaction
2981 parentTxn->commit();
2982 
2983 // Now the gate should be blocked because the parent had a confirmed write
2984 KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws));
2985 
2986 // Complete the commit
2987 test.pollAndExpectCalls({"commit"})[0]->fulfill();
2988 
2989 // Gate should unblock after commit completes
2990 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
2991 
2992 // Verify both writes were committed
2993 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("parent"))) == kj::str("data").asBytes());
2994 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("child"))) == kj::str("data").asBytes());
2995}
2996 
2997KJ_TEST("nested transaction: deeply nested confirmed write propagates to root") {
2998 ActorSqliteTest test;
2999 
3000 // Gate is currently not blocked.
3001 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
3002 
3003 // Start a parent transaction with unconfirmed write
3004 auto txn1 = test.actor.startTransaction();
3005 txn1->put(kj::str("level1"), kj::heapArray(kj::str("data").asBytes()), {.allowUnconfirmed = true},
3006 nullptr);
3007 
3008 {
3009 // Start a second level nested transaction with unconfirmed write
3010 auto txn2 = test.actor.startTransaction();
3011 txn2->put(kj::str("level2"), kj::heapArray(kj::str("data").asBytes()),
3012 {.allowUnconfirmed = true}, nullptr);
3013 
3014 {
3015 // Start a third level nested transaction with confirmed write
3016 auto txn3 = test.actor.startTransaction();
3017 txn3->put(kj::str("level3"), kj::heapArray(kj::str("data").asBytes()),
3018 {.allowUnconfirmed = false}, nullptr);
3019 
3020 // Gate still isn't blocked during the transaction
3021 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
3022 
3023 // Commit level 3 - should propagate someWriteConfirmed to level 2
3024 txn3->commit();
3025 }
3026 
3027 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
3028 
3029 // Commit level 2 - should propagate someWriteConfirmed to level 1
3030 txn2->commit();
3031 }
3032 
3033 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
3034 
3035 // Commit level 1 (root transaction)
3036 txn1->commit();
3037 
3038 // Now the gate should be blocked because level 3 had a confirmed write
3039 KJ_ASSERT(!test.gate.wait(nullptr).poll(test.ws));
3040 
3041 // Complete the commit
3042 test.pollAndExpectCalls({"commit"})[0]->fulfill();
3043 
3044 // Gate should unblock after commit completes
3045 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
3046 
3047 // Verify all writes were committed
3048 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("level1"))) == kj::str("data").asBytes());
3049 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("level2"))) == kj::str("data").asBytes());
3050 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("level3"))) == kj::str("data").asBytes());
3051}
3052 
3053KJ_TEST("nested transaction: rollback resets someWriteConfirmed flag") {
3054 ActorSqliteTest test;
3055 
3056 // Gate is currently not blocked.
3057 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
3058 
3059 // Start a parent transaction with unconfirmed write
3060 auto parentTxn = test.actor.startTransaction();
3061 parentTxn->put(kj::str("parent"), kj::heapArray(kj::str("data").asBytes()),
3062 {.allowUnconfirmed = true}, nullptr);
3063 
3064 {
3065 // Start a nested child transaction
3066 auto childTxn = test.actor.startTransaction();
3067 
3068 // Do a confirmed put in the child
3069 childTxn->put(kj::str("child"), kj::heapArray(kj::str("data").asBytes()),
3070 {.allowUnconfirmed = false}, nullptr);
3071 
3072 // Gate still isn't blocked during the transaction
3073 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
3074 
3075 // Rollback the child transaction instead of committing
3076 childTxn->rollback().wait(test.ws);
3077 }
3078 
3079 // Gate should still not be blocked
3080 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
3081 
3082 // Commit the parent transaction
3083 parentTxn->commit();
3084 
3085 // Gate should still not be blocked because child was rolled back
3086 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
3087 
3088 // Complete the commit
3089 test.pollAndExpectCalls({"commit"})[0]->fulfill();
3090 
3091 // Gate should still not be blocked after commit completes
3092 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
3093 
3094 // Verify only parent write was committed
3095 KJ_ASSERT(KJ_ASSERT_NONNULL(expectSync(test.get("parent"))) == kj::str("data").asBytes());
3096 KJ_ASSERT(expectSync(test.get("child")) == kj::none);
3097}
3098 
3099KJ_TEST("explicit transaction: commit failure breaks output gate even for unconfirmed writes") {
3100 ActorSqliteTest test({.monitorOutputGate = false});
3101 
3102 auto promise = test.gate.onBroken();
3103 
3104 // Gate is currently not blocked.
3105 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
3106 
3107 // Start an explicit transaction
3108 auto txn = test.actor.startTransaction();
3109 
3110 // Do an unconfirmed put
3111 txn->put(
3112 kj::str("foo"), kj::heapArray(kj::str("bar").asBytes()), {.allowUnconfirmed = true}, nullptr);
3113 
3114 // Commit the transaction
3115 txn->commit();
3116 
3117 // Gate should not be blocked yet because write was unconfirmed
3118 KJ_ASSERT(test.gate.wait(nullptr).poll(test.ws));
3119 
3120 // Reject the commit to simulate failure
3121 test.pollAndExpectCalls({"commit"})[0]->reject(KJ_EXCEPTION(FAILED, "commit failed"));
3122 
3123 // Gate should now be broken due to commit failure, even though write was unconfirmed
3124 KJ_EXPECT_THROW_MESSAGE("commit failed", promise.wait(test.ws));
3125}
3126 
3127KJ_TEST("ActorSqlite alarm cleared by abandonAlarm") {
3128 
3129 ActorSqliteTest test;
3130 
3131 test.setAlarm(oneMs);
3132 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
3133 test.pollAndExpectCalls({"commit"})[0]->fulfill();
3134 test.pollAndExpectCalls({});
3135 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
3136 
3137 // abandonAlarm() clears the alarm from SQLite:
3138 // setAlarm(null) -> commit -> scheduleRun(none) (move-later path).
3139 auto result = test.actor.abandonAlarm(oneMs).wait(test.ws);
3140 test.pollAndExpectCalls({"commit"})[0]->fulfill();
3141 test.pollAndExpectCalls({"scheduleRun(none)"})[0]->fulfill();
3142 test.pollAndExpectCalls({});
3143 
3144 // Returns kj::none: alarm was cleared, AlarmManager should not re-register.
3145 KJ_ASSERT(result == kj::none);
3146 
3147 // getAlarm() now returns null (alarm deleted from SQLite).
3148 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
3149}
3150 
3151KJ_TEST("ActorSqlite alarm preserved after ALARM_RETRY_MAX_TRIES uncounted (internal) failures") {
3152 // When all ALARM_RETRY_MAX_TRIES failures are uncounted (retryCountsAgainstLimit=false,
3153 // i.e. infrastructure errors), the alarm scheduler's countedRetry never reaches the limit and
3154 // abandonAlarm is NEVER called. The alarm must remain set in SQLite throughout so that
3155 // the scheduler can keep retrying indefinitely.
3156 
3157 ActorSqliteTest test;
3158 
3159 test.setAlarm(oneMs);
3160 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
3161 test.pollAndExpectCalls({"commit"})[0]->fulfill();
3162 test.pollAndExpectCalls({});
3163 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
3164 
3165 // Simulate uncounted failures well past ALARM_RETRY_MAX_TRIES (= 6).
3166 // countedRetry stays at 0; AlarmManager never gives up; abandonAlarm is never called.
3167 // We've seen alarms fail hundreds of times due to infrastructure errors in production,
3168 // so we check both at the boundary (6) and well beyond it (100).
3169 for (auto i = 0; i < 100; i++) {
3170 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
3171 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
3172 test.actor.cancelDeferredAlarmDeletion();
3173 test.pollAndExpectCalls({});
3174 
3175 // Check at the ALARM_RETRY_MAX_TRIES boundary and at the end.
3176 if (i == 5 || i == 99) {
3177 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
3178 }
3179 }
3180}
3181 
3182KJ_TEST("ActorSqlite abandonAlarm is a no-op when a newer alarm has replaced the abandoned one") {
3183 // If the user sets a new alarm between the last retry failure and the abandonAlarm() call,
3184 // and it has already committed to SQLite, abandonAlarm() must compare the time and leave
3185 // the new alarm untouched.
3186 
3187 ActorSqliteTest test;
3188 
3189 // Set the original alarm and commit it.
3190 test.setAlarm(oneMs);
3191 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
3192 test.pollAndExpectCalls({"commit"})[0]->fulfill();
3193 test.pollAndExpectCalls({});
3194 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
3195 
3196 // User sets a new alarm (twoMs).
3197 // The commit fires first; then the post-commit "move-later" logic fires scheduleRun(2ms)
3198 // because alarmScheduledNoLaterThan (oneMs) is earlier than the newly committed twoMs.
3199 test.setAlarm(twoMs);
3200 test.pollAndExpectCalls({"commit"})[0]->fulfill();
3201 test.pollAndExpectCalls({"scheduleRun(2ms)"})[0]->fulfill();
3202 test.pollAndExpectCalls({});
3203 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
3204 
3205 // abandonAlarm() for the original oneMs alarm must be a no-op: storedTime (twoMs) !=
3206 // scheduledTime (oneMs), so the time check prevents clearing the new alarm.
3207 // Returns twoMs so AlarmManager can re-register the actor's real alarm.
3208 auto result = test.actor.abandonAlarm(oneMs).wait(test.ws);
3209 test.pollAndExpectCalls({}); // No commit or scheduleRun -- correct no-op.
3210 
3211 KJ_ASSERT(KJ_ASSERT_NONNULL(result) == twoMs);
3212 
3213 // getAlarm() must still return twoMs.
3214 KJ_ASSERT(expectSync(test.getAlarm()) == twoMs);
3215}
3216 
3217KJ_TEST("ActorSqlite abandonAlarm returns kj::none when no alarm is stored") {
3218 ActorSqliteTest test;
3219 
3220 // No alarm ever set. abandonAlarm should be a pure no-op, returning kj::none.
3221 auto result = test.actor.abandonAlarm(oneMs).wait(test.ws);
3222 test.pollAndExpectCalls({});
3223 
3224 KJ_ASSERT(result == kj::none);
3225}
3226 
3227KJ_TEST("ActorSqlite abandonAlarm returns kj::none when inAlarmHandler") {
3228 ActorSqliteTest test;
3229 
3230 test.setAlarm(oneMs);
3231 test.pollAndExpectCalls({"scheduleRun(1ms)"})[0]->fulfill();
3232 test.pollAndExpectCalls({"commit"})[0]->fulfill();
3233 test.pollAndExpectCalls({});
3234 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
3235 
3236 // Arm the handler — inAlarmHandler is now true.
3237 auto armResult = test.actor.armAlarmHandler(oneMs, nullptr, testCurrentTime);
3238 KJ_ASSERT(armResult.is<ActorSqlite::RunAlarmHandler>());
3239 
3240 // abandonAlarm while handler is running: returns kj::none (handler owns the alarm).
3241 auto result = test.actor.abandonAlarm(oneMs).wait(test.ws);
3242 test.pollAndExpectCalls({});
3243 
3244 KJ_ASSERT(result == kj::none);
3245 
3246 // kj::none because haveDeferredDelete hides the alarm during the handler.
3247 KJ_ASSERT(expectSync(test.getAlarm()) == kj::none);
3248 
3249 // Cancel the deferred delete so cleanup doesn't trigger a commit.
3250 test.actor.cancelDeferredAlarmDeletion();
3251 
3252 // After cancellation, getAlarm() reads SQLite again -- oneMs is still there since
3253 // abandonAlarm was a no-op while the handler was running.
3254 KJ_ASSERT(expectSync(test.getAlarm()) == oneMs);
3255}
3256 
3257} // namespace
3258} // namespace workerd