Skip to content
File

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

208.5 KB
1// Copyright (c) 2017-2022 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-cache.h"
6#include "io-gate.h"
7 
8#include <workerd/util/capnp-mock.h>
9#include <workerd/util/test.h>
10 
11#include <capnp/dynamic.h>
12#include <kj/debug.h>
13#include <kj/list.h>
14#include <kj/source-location.h>
15#include <kj/test.h>
16#include <kj/thread.h>
17 
18namespace workerd {
19namespace {
20 
21// =======================================================================================
22// Test helpers specific to ActorCache test.
23 
24template <typename T>
25kj::Promise<T> eagerlyReportExceptions(kj::Promise<T> promise, kj::SourceLocation location = {}) {
26 // TODO(cleanup): Move to KJ somewhere?
27 return promise.eagerlyEvaluate([location](kj::Exception&& e) -> T {
28 KJ_LOG_AT(ERROR, location, e);
29 kj::throwFatalException(kj::mv(e));
30 });
31}
32 
33template <typename T>
34kj::Promise<T> expectUncached(
35 kj::OneOf<T, kj::Promise<T>> result, kj::SourceLocation location = {}) {
36 // Expect that a result returned by get()/list()/delete() was not served entirely from cache,
37 // and return the promise.
38 KJ_SWITCH_ONEOF(result) {
39 KJ_CASE_ONEOF(promise, kj::Promise<T>) {
40 return eagerlyReportExceptions(kj::mv(promise), location);
41 }
42 KJ_CASE_ONEOF(value, T) {
43 KJ_FAIL_ASSERT_AT(location, "result was unexpectedly cached");
44 }
45 }
46 KJ_UNREACHABLE;
47}
48 
49template <typename T>
50T expectCached(kj::OneOf<T, kj::Promise<T>> result, kj::SourceLocation location = {}) {
51 // Expect that a result returned by get()/list()/delete() was served entirely from cache, and
52 // return that.
53 KJ_SWITCH_ONEOF(result) {
54 KJ_CASE_ONEOF(promise, kj::Promise<T>) {
55 KJ_FAIL_ASSERT_AT(location, "result was unexpectedly uncached");
56 }
57 KJ_CASE_ONEOF(value, T) {
58 return kj::mv(value);
59 }
60 }
61 KJ_UNREACHABLE;
62}
63 
64struct KeyValuePtr {
65 kj::StringPtr key;
66 kj::StringPtr value;
67 
68 inline bool operator==(const KeyValuePtr& other) const {
69 return key == other.key && value == other.value;
70 }
71 inline kj::String toString() const {
72 return kj::str(key, ": ", value);
73 }
74};
75 
76struct KeyValue {
77 kj::String key;
78 kj::String value;
79 
80 inline bool operator==(const KeyValuePtr& other) const {
81 return key == other.key && value == other.value;
82 }
83 inline bool operator==(const KeyValue& other) const {
84 return key == other.key && value == other.value;
85 }
86 inline kj::String toString() const {
87 return kj::str(key, ": ", value);
88 }
89};
90 
91kj::ArrayPtr<const KeyValuePtr> kvs(kj::ArrayPtr<const KeyValuePtr> a) {
92 return a;
93}
94// We want to be able to write checks like:
95//
96// KJ_ASSERT(results == {{"bar", "456"}, {"foo", "123"}});
97//
98// Unfortunately, the compiler is not smart enough to figure out how to interpret the braced list
99// the right of the comparison. So, we give it a little help by wrapping it in `kvs()`, like:
100//
101// KJ_ASSERT(results == kvs({{"bar", "456"}, {"foo", "123"}}));
102 
103// stringifyValues() is a convenience function that turns byte-array values returned by ActorCache
104// into strings, for a variety of different return types.
105 
106kj::String stringifyValues(ActorCache::ValuePtr value) {
107 return kj::str(value.asChars());
108}
109KeyValue stringifyValues(ActorCache::KeyValuePtrPair kv) {
110 return {kj::str(kv.key), stringifyValues(kv.value)};
111}
112kj::Array<KeyValue> stringifyValues(const ActorCache::GetResultList& list) {
113 return KJ_MAP(e, list) { return stringifyValues(e); };
114}
115 
116template <typename T>
117auto stringifyValues(const kj::Maybe<T>& values) {
118 return values.map([](auto& v) { return stringifyValues(v); });
119}
120 
121template <typename T>
122kj::OneOf<decltype(stringifyValues(kj::instance<T>())),
123 kj::Promise<decltype(stringifyValues(kj::instance<T>()))>>
124stringifyValues(kj::OneOf<T, kj::Promise<T>> result) {
125 KJ_SWITCH_ONEOF(result) {
126 KJ_CASE_ONEOF(promise, kj::Promise<T>) {
127 return promise.then([](T result) { return stringifyValues(kj::mv(result)); });
128 }
129 KJ_CASE_ONEOF(value, T) {
130 return stringifyValues(kj::mv(value));
131 }
132 }
133 KJ_UNREACHABLE;
134}
135 
136struct ActorCacheConvenienceWrappers {
137 // Convenience methods to make test more concise by handling value conversions to/from strings,
138 // and allowing parameters to be string literals instead of owned strings.
139 //
140 // This is formulated as a mixin inherited by ActorCacheTest, below, so that it can be reused
141 // for transactions as well.
142 
143 ActorCacheConvenienceWrappers(ActorCacheOps& target): target(target) {}
144 
145 auto get(kj::StringPtr key, ActorCache::ReadOptions options = {}) {
146 return stringifyValues(target.get(kj::str(key), options));
147 }
148 auto get(kj::ArrayPtr<const kj::StringPtr> keys, ActorCache::ReadOptions options = {}) {
149 return stringifyValues(target.get(KJ_MAP(k, keys) { return kj::str(k); }, options));
150 }
151 auto getAlarm(ActorCache::ReadOptions options = {}) {
152 return target.getAlarm(options);
153 }
154 
155 auto list(kj::StringPtr begin,
156 kj::StringPtr end,
157 kj::Maybe<uint> limit = kj::none,
158 ActorCache::ReadOptions options = {}) {
159 return stringifyValues(target.list(kj::str(begin), kj::str(end), limit, options));
160 }
161 auto listReverse(kj::StringPtr begin,
162 kj::StringPtr end,
163 kj::Maybe<uint> limit = kj::none,
164 ActorCache::ReadOptions options = {}) {
165 return stringifyValues(target.listReverse(kj::str(begin), kj::str(end), limit, options));
166 }
167 
168 auto list(kj::StringPtr begin,
169 decltype(nullptr),
170 kj::Maybe<uint> limit = kj::none,
171 ActorCache::ReadOptions options = {}) {
172 return stringifyValues(target.list(kj::str(begin), kj::none, limit, options));
173 }
174 auto listReverse(kj::StringPtr begin,
175 decltype(nullptr),
176 kj::Maybe<uint> limit = kj::none,
177 ActorCache::ReadOptions options = {}) {
178 return stringifyValues(target.listReverse(kj::str(begin), kj::none, limit, options));
179 }
180 
181 auto put(kj::StringPtr key, kj::StringPtr value, ActorCache::WriteOptions options = {}) {
182 return target.put(kj::str(key), kj::heapArray(value.asBytes()), options, nullptr);
183 }
184 auto put(kj::ArrayPtr<const KeyValuePtr> kvs, ActorCache::WriteOptions options = {}) {
185 return target.put(KJ_MAP(kv, kvs) {
186 return ActorCache::KeyValuePair{kj::str(kv.key), kj::heapArray(kv.value.asBytes())};
187 }, options, nullptr);
188 }
189 auto setAlarm(kj::Maybe<kj::Date> newTime, ActorCache::WriteOptions options = {}) {
190 return target.setAlarm(newTime, options, nullptr);
191 }
192 
193 auto delete_(kj::StringPtr key, ActorCache::WriteOptions options = {}) {
194 return target.delete_(kj::str(key), options, nullptr);
195 }
196 auto delete_(kj::ArrayPtr<const kj::StringPtr> keys, ActorCache::WriteOptions options = {}) {
197 return target.delete_(KJ_MAP(k, keys) { return kj::str(k); }, options, nullptr);
198 }
199 
200 private:
201 ActorCacheOps& target;
202};
203 
204struct ActorCacheTestOptions {
205 bool monitorOutputGate = true;
206 size_t softLimit = 512 * 1024;
207 size_t hardLimit = 1024 * 1024;
208 kj::Duration staleTimeout = 1 * kj::SECONDS;
209 size_t dirtyListByteLimit = 64 * 1024;
210 size_t maxKeysPerRpc = 128;
211 bool noCache = false;
212 bool neverFlush = false;
213};
214 
215struct ActorCacheTest: public ActorCacheConvenienceWrappers {
216 // Common test setup code and helpers used in many test cases.
217 
218 kj::EventLoop loop;
219 kj::WaitScope ws;
220 kj::Own<MockServer> mockStorage;
221 
222 ActorCache::SharedLru lru;
223 OutputGate gate;
224 ActorCache cache;
225 
226 kj::Promise<void> gateBrokenPromise;
227 
228 kj::UnwindDetector unwindDetector;
229 
230 ActorCacheTest(ActorCacheTestOptions options = {},
231 MockServer::Pair<rpc::ActorStorage::Stage> mockPair =
232 MockServer::make<rpc::ActorStorage::Stage>())
233 : ActorCacheConvenienceWrappers(cache),
234 ws(loop),
235 mockStorage(kj::mv(mockPair.mock)),
236 lru({options.softLimit, options.hardLimit, options.staleTimeout, options.dirtyListByteLimit,
237 options.maxKeysPerRpc, options.noCache, options.neverFlush}),
238 cache(kj::mv(mockPair.client), lru, gate),
239 gateBrokenPromise(options.monitorOutputGate ? eagerlyReportExceptions(gate.onBroken())
240 : kj::Promise<void>(kj::READY_NOW)) {}
241 
242 // Simulates `count` counted alarm handler failures for `alarmTime`, leaving the cache
243 // in KnownAlarmTime{CLEAN, alarmTime} as AlarmManager would after each retry.
244 ~ActorCacheTest() noexcept(false) {
245 // Make sure if the output gate has been broken, the exception was reported. This is important
246 // to report errors thrown inside flush(), since those won't otherwise propagate into the test
247 // body.
248 gateBrokenPromise.poll(ws);
249 
250 if (!unwindDetector.isUnwinding()) {
251 // On successful test completion, also check that there were no extra calls to the mock.
252 mockStorage->expectNoActivity(ws);
253 
254 cache.verifyConsistencyForTest();
255 }
256 }
257};
258 
259KJ_TEST("ActorCache single-key basics") {
260 ActorCacheTest test;
261 auto& ws = test.ws;
262 auto& mockStorage = test.mockStorage;
263 
264 // Get value that is present on disk.
265 {
266 auto promise = expectUncached(test.get("foo"));
267 
268 mockStorage->expectCall("get", ws)
269 .withParams(CAPNP(key = "foo"))
270 .thenReturn(CAPNP(value = "bar"));
271 
272 auto result = KJ_ASSERT_NONNULL(promise.wait(ws));
273 KJ_EXPECT(result == "bar");
274 }
275 
276 // Get value that is absent on disk.
277 {
278 auto promise = expectUncached(test.get("bar"));
279 
280 mockStorage->expectCall("get", ws).withParams(CAPNP(key = "bar")).thenReturn(CAPNP());
281 
282 auto result = promise.wait(ws);
283 KJ_EXPECT(result == kj::none);
284 }
285 
286 // Get cached.
287 {
288 auto result = KJ_ASSERT_NONNULL(expectCached(test.get("foo")));
289 KJ_EXPECT(result == "bar");
290 }
291 {
292 auto result = expectCached(test.get("bar"));
293 KJ_EXPECT(result == kj::none);
294 }
295 
296 // Overwrite with a put().
297 {
298 test.put("foo", "baz");
299 
300 mockStorage->expectCall("put", ws)
301 .withParams(CAPNP(entries = [(key = "foo", value = "baz")]))
302 .thenReturn(CAPNP());
303 }
304 
305 {
306 auto result = KJ_ASSERT_NONNULL(expectCached(test.get("foo")));
307 KJ_EXPECT(result == "baz");
308 }
309 
310 {
311 KJ_ASSERT(expectCached(test.delete_("foo")));
312 
313 mockStorage->expectCall("delete", ws)
314 .withParams(CAPNP(keys = ["foo"]))
315 .thenReturn(CAPNP(numDeleted = 1));
316 }
317 
318 {
319 auto result = expectCached(test.get("foo"));
320 KJ_EXPECT(result == kj::none);
321 }
322}
323 
324KJ_TEST("ActorCache multi-key basics") {
325 ActorCacheTest test;
326 auto& ws = test.ws;
327 auto& mockStorage = test.mockStorage;
328 
329 {
330 // Request four keys, but only return two. The others should be marked empty. Note we
331 // intentionally make sure that, in alphabetical order, the keys alternate between present
332 // and absent, with the last one being absent, for maximum code coverage.
333 auto promise = expectUncached(test.get({"foo"_kj, "bar"_kj, "baz"_kj, "qux"_kj}));
334 
335 mockStorage->expectCall("getMultiple", ws)
336 .withParams(CAPNP(keys = [ "bar", "baz", "foo", "qux" ]), "stream"_kj)
337 .useCallback("stream", [&](MockClient stream) {
338 stream
339 .call("values",
340 CAPNP(list =
341 [
342 (key = "bar", value = "456"),
343 // baz absent
344 (key = "foo", value = "123"),
345 // qux absent
346 ]))
347 .expectReturns(CAPNP(), ws);
348 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
349 }).expectCanceled();
350 
351 auto results = promise.wait(ws);
352 KJ_ASSERT(results == kvs({{"bar", "456"}, {"foo", "123"}}));
353 }
354 
355 {
356 auto results = expectCached(test.get({"foo"_kj, "bar"_kj, "baz"_kj, "qux"_kj}));
357 KJ_ASSERT(results == kvs({{"bar", "456"}, {"foo", "123"}}));
358 }
359 
360 {
361 test.put({{"foo", "321"}, {"bar", "654"}});
362 
363 mockStorage->expectCall("put", ws)
364 .withParams(CAPNP(entries = [ (key = "foo", value = "321"), (key = "bar", value = "654") ]))
365 .thenReturn(CAPNP());
366 }
367 
368 {
369 auto results = expectCached(test.get({"foo"_kj, "bar"_kj}));
370 KJ_ASSERT(results == kvs({{"bar", "654"}, {"foo", "321"}}));
371 }
372 
373 {
374 KJ_ASSERT(expectCached(test.delete_({"foo"_kj, "bar"_kj, "baz"_kj, "qux"_kj})) == 2);
375 
376 mockStorage->expectCall("delete", ws)
377 .withParams(CAPNP(keys = [ "foo", "bar" ]))
378 .thenReturn(CAPNP(numDeleted = 2));
379 }
380 
381 {
382 auto results = expectCached(test.get({"foo"_kj, "bar"_kj}));
383 KJ_ASSERT(results == kvs({}));
384 }
385}
386 
387// =======================================================================================
388 
389KJ_TEST("ActorCache more puts") {
390 ActorCacheTest test;
391 auto& ws = test.ws;
392 auto& mockStorage = test.mockStorage;
393 
394 {
395 test.put("foo", "bar");
396 
397 // Value is immediately in cache.
398 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "bar");
399 
400 auto inProgressFlush = mockStorage->expectCall("put", ws).withParams(
401 CAPNP(entries = [(key = "foo", value = "bar")]));
402 
403 // Still in cache during flush.
404 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "bar");
405 
406 kj::mv(inProgressFlush).thenReturn(CAPNP());
407 }
408 
409 // Still in cache after transaction completion.
410 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "bar");
411 
412 // Putting the exact same value is redundant, so doesn't do an RPC.
413 {
414 test.put("foo", "bar");
415 mockStorage->expectNoActivity(ws);
416 }
417 
418 // Putting a different value is not redundant.
419 {
420 test.put("foo", "baz");
421 
422 mockStorage->expectCall("put", ws)
423 .withParams(CAPNP(entries = [(key = "foo", value = "baz")]))
424 .thenReturn(CAPNP());
425 }
426 
427 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "baz");
428}
429 
430KJ_TEST("ActorCache more deletes") {
431 ActorCacheTest test;
432 auto& ws = test.ws;
433 auto& mockStorage = test.mockStorage;
434 
435 {
436 auto promise = expectUncached(test.delete_("foo"));
437 
438 // Value is immediately in cache.
439 KJ_ASSERT(expectCached(test.get("foo")) == kj::none);
440 
441 auto mockDelete = mockStorage->expectCall("delete", ws).withParams(CAPNP(keys = ["foo"]));
442 
443 // Still in cache during flush.
444 KJ_ASSERT(!promise.poll(ws));
445 KJ_ASSERT(expectCached(test.get("foo")) == kj::none);
446 
447 kj::mv(mockDelete).thenReturn(CAPNP(numDeleted = 1));
448 
449 // Delete call returned true due to numDeleted = 1.
450 KJ_ASSERT(promise.wait(ws));
451 }
452 
453 // Still in cache after transaction completion.
454 KJ_ASSERT(expectCached(test.get("foo")) == kj::none);
455 
456 // Try a case where the key isn't on disk.
457 {
458 auto promise = expectUncached(test.delete_("bar"));
459 
460 mockStorage->expectCall("delete", ws)
461 .withParams(CAPNP(keys = ["bar"]))
462 .thenReturn(CAPNP(numDeleted = 0));
463 
464 // Delete call returned false due to numDeleted = 0.
465 KJ_ASSERT(!promise.wait(ws));
466 }
467 
468 // Deleting an already-deleted key is redundant, so doesn't do an RPC.
469 {
470 KJ_ASSERT(!expectCached(test.delete_("foo")));
471 KJ_ASSERT(!expectCached(test.delete_("bar")));
472 
473 mockStorage->expectNoActivity(ws);
474 }
475 
476 // Putting over the deleted key is not redundant.
477 {
478 test.put("foo", "baz");
479 
480 mockStorage->expectCall("put", ws)
481 .withParams(CAPNP(entries = [(key = "foo", value = "baz")]))
482 .thenReturn(CAPNP());
483 }
484 
485 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "baz");
486 
487 // Deleting it again is not redundant.
488 {
489 KJ_ASSERT(expectCached(test.delete_("foo")));
490 
491 mockStorage->expectCall("delete", ws)
492 .withParams(CAPNP(keys = ["foo"]))
493 .thenReturn(CAPNP(numDeleted = 1));
494 }
495 
496 KJ_ASSERT(expectCached(test.get("foo")) == kj::none);
497}
498 
499KJ_TEST("ActorCache more multi-puts") {
500 ActorCacheTest test;
501 auto& ws = test.ws;
502 auto& mockStorage = test.mockStorage;
503 
504 // Create a scenario where we have several cached and uncached keys.
505 // foo, bar = cached with values
506 // baz, qux = cached as absent
507 // corge, grault = not cached
508 {
509 auto promise = expectUncached(test.get({"foo", "bar", "baz", "qux"}));
510 
511 mockStorage->expectCall("getMultiple", ws)
512 .withParams(CAPNP(keys = [ "bar", "baz", "foo", "qux" ]), "stream"_kj)
513 .useCallback("stream", [&](MockClient stream) {
514 stream
515 .call("values",
516 CAPNP(list = [ (key = "bar", value = "456"), (key = "foo", value = "123") ]))
517 .expectReturns(CAPNP(), ws);
518 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
519 }).expectCanceled();
520 
521 auto results = promise.wait(ws);
522 KJ_ASSERT(results == kvs({{"bar", "456"}, {"foo", "123"}}));
523 }
524 
525 {
526 test.put({{"foo", "321"}, {"bar", "456"}, {"baz", "654"}, {"corge", "987"}});
527 
528 // Values are immediately in cache.
529 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "321");
530 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "456");
531 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "654");
532 KJ_ASSERT(expectCached(test.get("qux")) == nullptr);
533 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("corge"))) == "987");
534 
535 mockStorage->expectCall("put", ws)
536 .withParams(CAPNP(entries =
537 [
538 (key = "foo", value = "321"),
539 // bar omitted because it was redundant
540 (key = "baz", value = "654"), (key = "corge", value = "987")
541 ]))
542 .thenReturn(CAPNP());
543 }
544 
545 // Fetch everything again for good measure.
546 {
547 auto promise =
548 expectUncached(test.get({"foo"_kj, "bar"_kj, "baz"_kj, "qux"_kj, "corge"_kj, "grault"_kj}));
549 
550 // Only "grault" is not cached.
551 mockStorage->expectCall("getMultiple", ws)
552 .withParams(CAPNP(keys = ["grault"]), "stream"_kj)
553 .useCallback("stream", [&](MockClient stream) {
554 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
555 }).expectCanceled();
556 
557 auto results = promise.wait(ws);
558 KJ_ASSERT(results ==
559 kvs({
560 {"bar", "456"},
561 {"baz", "654"},
562 {"corge", "987"},
563 {"foo", "321"},
564 }));
565 }
566}
567 
568KJ_TEST("ActorCache more multi-deletes") {
569 ActorCacheTest test;
570 auto& ws = test.ws;
571 auto& mockStorage = test.mockStorage;
572 
573 // Create a scenario where we have several cached and uncached keys.
574 // foo, bar = cached with values
575 // baz, qux = cached as absent
576 // corge, grault = not cached
577 {
578 auto promise = expectUncached(test.get({"foo", "bar", "baz", "qux"}));
579 
580 mockStorage->expectCall("getMultiple", ws)
581 .withParams(CAPNP(keys = [ "bar", "baz", "foo", "qux" ]), "stream"_kj)
582 .useCallback("stream", [&](MockClient stream) {
583 stream
584 .call("values",
585 CAPNP(list = [ (key = "bar", value = "456"), (key = "foo", value = "123") ]))
586 .expectReturns(CAPNP(), ws);
587 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
588 }).expectCanceled();
589 
590 auto results = promise.wait(ws);
591 KJ_ASSERT(results == kvs({{"bar", "456"}, {"foo", "123"}}));
592 }
593 
594 {
595 auto promise = expectUncached(test.delete_({"bar"_kj, "qux"_kj, "corge"_kj, "grault"_kj}));
596 
597 // Values are immediately in cache.
598 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
599 KJ_ASSERT(expectCached(test.get("bar")) == nullptr);
600 KJ_ASSERT(expectCached(test.get("baz")) == nullptr);
601 KJ_ASSERT(expectCached(test.get("qux")) == nullptr);
602 KJ_ASSERT(expectCached(test.get("corge")) == nullptr);
603 KJ_ASSERT(expectCached(test.get("grault")) == nullptr);
604 
605 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
606 mockTxn->expectCall("delete", ws)
607 .withParams(CAPNP(keys = [ "corge", "grault" ]))
608 .thenReturn(CAPNP(numDeleted = 1));
609 mockTxn->expectCall("delete", ws)
610 .withParams(CAPNP(keys = ["bar"]))
611 .thenReturn(CAPNP(numDeleted = 65382)); // count is ignored
612 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
613 mockTxn->expectDropped(ws);
614 
615 KJ_ASSERT(promise.wait(ws) == 2);
616 }
617 
618 // Fetch everything again for good measure.
619 {
620 auto promise = expectUncached(
621 test.get({"foo"_kj, "bar"_kj, "baz"_kj, "qux"_kj, "corge"_kj, "grault"_kj, "garply"_kj}));
622 
623 // Only "garply" is not cached.
624 mockStorage->expectCall("getMultiple", ws)
625 .withParams(CAPNP(keys = ["garply"]), "stream"_kj)
626 .useCallback("stream", [&](MockClient stream) {
627 stream.call("values", CAPNP(list = [(key = "garply", value = "abcd")]))
628 .expectReturns(CAPNP(), ws);
629 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
630 }).expectCanceled();
631 
632 auto results = promise.wait(ws);
633 KJ_ASSERT(results == kvs({{"foo", "123"}, {"garply", "abcd"}}));
634 }
635}
636 
637KJ_TEST("ActorCache batching due to maxKeysPerRpc") {
638 ActorCacheTest test({.maxKeysPerRpc = 2});
639 auto& ws = test.ws;
640 auto& mockStorage = test.mockStorage;
641 
642 // Do 5 puts and 3 deletes and expect a transaction that is batched accordingly given we set the
643 // batch size to 2.
644 test.put({{"foo", "123"}, {"bar", "456"}, {"baz", "789"}});
645 test.put("qux", "555");
646 test.put("corge", "999");
647 
648 // Note that because we drop the returned promises from these deletes, they end up as "muted"
649 // deletes, so the resulting batches don't have to match the original calls.
650 test.delete_("grault");
651 test.delete_({"garply"_kj, "waldo"_kj});
652 
653 // We keep these promises, so they should not be "muted". Specifically, "count4" should be its own
654 // batch despite fitting in a batch with "count3" because it's a separate delete.
655 auto deleteProm1 = expectUncached(test.delete_({"count1"_kj, "count2"_kj, "count3"_kj}));
656 auto deleteProm2 = expectUncached(test.delete_({"count4"_kj}));
657 auto deleteProm3 = expectUncached(test.delete_({"count5"_kj, "count6"_kj}));
658 
659 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
660 mockTxn->expectCall("delete", ws)
661 .withParams(CAPNP(keys = [ "count1", "count2" ]))
662 .thenReturn(CAPNP(numDeleted = 1)); // Treat one of this batch as present, 2 total.
663 mockTxn->expectCall("delete", ws)
664 .withParams(CAPNP(keys = ["count3"]))
665 .thenReturn(CAPNP(numDeleted = 1)); // Treat one of this batch as present, 2 total.
666 mockTxn->expectCall("delete", ws)
667 .withParams(CAPNP(keys = ["count4"]))
668 .thenReturn(CAPNP(numDeleted = 0)); // Treat this batch as absent.
669 mockTxn->expectCall("delete", ws)
670 .withParams(CAPNP(keys = [ "count5", "count6" ]))
671 .thenReturn(CAPNP(numDeleted = 2)); // Treat all of this batch as present.
672 mockTxn->expectCall("delete", ws)
673 .withParams(CAPNP(keys = [ "grault", "garply" ]))
674 .thenReturn(CAPNP(numDeleted = 1));
675 mockTxn->expectCall("delete", ws)
676 .withParams(CAPNP(keys = ["waldo"]))
677 .thenReturn(CAPNP(numDeleted = 1));
678 mockTxn->expectCall("put", ws)
679 .withParams(CAPNP(entries = [ (key = "foo", value = "123"), (key = "bar", value = "456") ]))
680 .thenReturn(CAPNP());
681 mockTxn->expectCall("put", ws)
682 .withParams(CAPNP(entries = [ (key = "baz", value = "789"), (key = "qux", value = "555") ]))
683 .thenReturn(CAPNP());
684 mockTxn->expectCall("put", ws)
685 .withParams(CAPNP(entries = [(key = "corge", value = "999")]))
686 .thenReturn(CAPNP());
687 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
688 mockTxn->expectDropped(ws);
689 
690 KJ_EXPECT(deleteProm1.wait(ws) == 2);
691 KJ_EXPECT(deleteProm2.wait(ws) == 0);
692 KJ_EXPECT(deleteProm3.wait(ws) == 2);
693}
694 
695KJ_TEST("ActorCache batching due to max storage RPC words") {
696 ActorCacheTest test({.hardLimit = 128 * 1024 * 1024});
697 auto& ws = test.ws;
698 auto& mockStorage = test.mockStorage;
699 
700 // Doing 128 puts with 128 KiB values should exceed the 16 MiB limit enforced on storage RPCs.
701 auto bigVal = kj::heapArray<const byte>(128 * 1024);
702 for (int i = 0; i < 128; ++i) {
703 test.cache.put(kj::str(i),
704 kj::Array<const byte>(bigVal.begin(), bigVal.size(), kj::NullArrayDisposer::instance),
705 ActorCache::WriteOptions(), nullptr);
706 }
707 
708 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
709 mockTxn->expectCall("put", ws).thenReturn(CAPNP());
710 mockTxn->expectCall("put", ws).thenReturn(CAPNP());
711 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
712 mockTxn->expectDropped(ws);
713}
714 
715KJ_TEST("ActorCache deleteAll()") {
716 ActorCacheTest test;
717 auto& ws = test.ws;
718 auto& mockStorage = test.mockStorage;
719 
720 // Populate the cache with some stuff.
721 {
722 auto promise = expectUncached(test.get({"qux"_kj, "corge"_kj}));
723 
724 mockStorage->expectCall("getMultiple", ws)
725 .withParams(CAPNP(keys = [ "corge", "qux" ]), "stream"_kj)
726 .useCallback("stream", [&](MockClient stream) {
727 stream.call("values", CAPNP(list = [(key = "corge", value = "555")]))
728 .expectReturns(CAPNP(), ws);
729 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
730 }).expectCanceled();
731 
732 auto results = promise.wait(ws);
733 KJ_ASSERT(results == kvs({{"corge", "555"}}));
734 }
735 
736 test.put("foo", "123"); // plain put
737 auto deletePromise = expectUncached(test.delete_({"bar"_kj, "baz"_kj, "grault"_kj}));
738 test.put("baz", "789"); // overwrites a counted delete
739 test.delete_("garply"); // uncounted delete
740 
741 auto deleteAll = test.cache.deleteAll({}, nullptr);
742 
743 // Post-deleteAll writes.
744 test.put("grault", "12345");
745 test.put("garply", "54321");
746 test.put("waldo", "99999");
747 
748 // Alarms are not affected by deleteAll, so this alarm set should actually end up in
749 // the pre-deleteAll flush.
750 test.setAlarm(12345 * kj::MILLISECONDS + kj::UNIX_EPOCH);
751 
752 KJ_ASSERT(expectCached(test.get("foo")) == nullptr);
753 KJ_ASSERT(expectCached(test.get("baz")) == nullptr);
754 KJ_ASSERT(expectCached(test.get("corge")) == nullptr);
755 KJ_ASSERT(expectCached(test.get("a")) == nullptr);
756 KJ_ASSERT(expectCached(test.get("z")) == nullptr);
757 KJ_ASSERT(expectCached(test.get("")) == nullptr);
758 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("grault"))) == "12345");
759 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("garply"))) == "54321");
760 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("waldo"))) == "99999");
761 
762 {
763 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
764 mockTxn->expectCall("delete", ws)
765 .withParams(CAPNP(keys = [ "bar", "baz", "grault" ]))
766 .thenReturn(CAPNP(numDeleted = 2));
767 mockTxn->expectCall("delete", ws)
768 .withParams(CAPNP(keys = ["garply"]))
769 .thenReturn(CAPNP(numDeleted = 2));
770 mockTxn->expectCall("put", ws)
771 .withParams(CAPNP(entries = [ (key = "foo", value = "123"), (key = "baz", value = "789") ]))
772 .thenReturn(CAPNP());
773 mockTxn->expectCall("setAlarm", ws)
774 .withParams(CAPNP(scheduledTimeMs = 12345))
775 .thenReturn(CAPNP());
776 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
777 mockTxn->expectDropped(ws);
778 }
779 
780 mockStorage->expectCall("deleteAll", ws).thenReturn(CAPNP(numDeleted = 2));
781 
782 KJ_ASSERT(deleteAll.count.wait(ws) == 2);
783 
784 // Post-deleteAll writes in a new flush.
785 {
786 mockStorage->expectCall("put", ws)
787 .withParams(CAPNP(entries =
788 [
789 (key = "grault", value = "12345"),
790 (key = "garply", value = "54321"), (key = "waldo", value = "99999")
791 ]))
792 .thenReturn(CAPNP());
793 }
794 
795 KJ_ASSERT(deletePromise.wait(ws) == 2);
796 
797 KJ_ASSERT(expectCached(test.get("foo")) == nullptr);
798 KJ_ASSERT(expectCached(test.get("baz")) == nullptr);
799 KJ_ASSERT(expectCached(test.get("corge")) == nullptr);
800 KJ_ASSERT(expectCached(test.get("a")) == nullptr);
801 KJ_ASSERT(expectCached(test.get("z")) == nullptr);
802 KJ_ASSERT(expectCached(test.get("")) == nullptr);
803 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("grault"))) == "12345");
804 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("garply"))) == "54321");
805 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("waldo"))) == "99999");
806}
807 
808KJ_TEST("ActorCache deleteAll() during transaction commit") {
809 // This tests a race condition that existed previously in the code.
810 
811 ActorCacheTest test;
812 auto& ws = test.ws;
813 auto& mockStorage = test.mockStorage;
814 
815 // Get a transaction going, and then issue a deleteAll() in the middle of it.
816 test.put("foo", "123");
817 
818 {
819 auto inProgressFlush = mockStorage->expectCall("put", ws).withParams(
820 CAPNP(entries = [(key = "foo", value = "123")]));
821 
822 // Issue a put and a deleteAll() here!
823 test.put("bar", "456");
824 test.cache.deleteAll({}, nullptr);
825 
826 kj::mv(inProgressFlush).thenReturn(CAPNP());
827 }
828 
829 // We should see a new flush happen for the pre-deleteAll() write.
830 {
831 mockStorage->expectCall("put", ws)
832 .withParams(CAPNP(entries = [(key = "bar", value = "456")]))
833 .thenReturn(CAPNP());
834 }
835 
836 // Now the deleteAll() actually happens.
837 mockStorage->expectCall("deleteAll", ws).thenReturn(CAPNP());
838}
839 
840KJ_TEST("ActorCache deleteAll() again when previous one isn't done yet") {
841 ActorCacheTest test;
842 auto& ws = test.ws;
843 auto& mockStorage = test.mockStorage;
844 
845 // Populate the cache with some stuff.
846 {
847 auto promise = expectUncached(test.get({"qux"_kj, "corge"_kj}));
848 
849 mockStorage->expectCall("getMultiple", ws)
850 .withParams(CAPNP(keys = [ "corge", "qux" ]), "stream"_kj)
851 .useCallback("stream", [&](MockClient stream) {
852 stream.call("values", CAPNP(list = [(key = "corge", value = "555")]))
853 .expectReturns(CAPNP(), ws);
854 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
855 }).expectCanceled();
856 
857 auto results = promise.wait(ws);
858 KJ_ASSERT(results == kvs({{"corge", "555"}}));
859 }
860 
861 test.put("foo", "123"); // plain put
862 auto deletePromise = expectUncached(test.delete_({"bar"_kj, "baz"_kj, "grault"_kj}));
863 test.put("baz", "789"); // overwrites a counted delete
864 test.delete_("garply"); // uncounted delete
865 
866 auto deleteAllA = test.cache.deleteAll({}, nullptr);
867 
868 // Post-deleteAll writes.
869 test.put("grault", "12345");
870 test.put("garply", "54321");
871 test.put("waldo", "99999");
872 
873 {
874 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
875 mockTxn->expectCall("delete", ws)
876 .withParams(CAPNP(keys = [ "bar", "baz", "grault" ]))
877 .thenReturn(CAPNP(numDeleted = 2));
878 mockTxn->expectCall("delete", ws)
879 .withParams(CAPNP(keys = ["garply"]))
880 .thenReturn(CAPNP(numDeleted = 2));
881 mockTxn->expectCall("put", ws)
882 .withParams(CAPNP(entries = [ (key = "foo", value = "123"), (key = "baz", value = "789") ]))
883 .thenReturn(CAPNP());
884 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
885 mockTxn->expectDropped(ws);
886 }
887 
888 // Do another deleteAll() before the first one is done.
889 auto deleteAllB = test.cache.deleteAll({}, nullptr);
890 
891 // And a write after that.
892 test.put("fred", "2323");
893 
894 // Now finish it.
895 mockStorage->expectCall("deleteAll", ws).thenReturn(CAPNP(numDeleted = 2));
896 KJ_ASSERT(deleteAllA.count.wait(ws) == 2);
897 KJ_ASSERT(deleteAllB.count.wait(ws) == 0);
898 
899 // The deleteAll()s were coalesced, so only the final write is committed.
900 {
901 mockStorage->expectCall("put", ws)
902 .withParams(CAPNP(entries = [(key = "fred", value = "2323")]))
903 .thenReturn(CAPNP());
904 }
905 KJ_ASSERT(deletePromise.wait(ws) == 2);
906}
907 
908KJ_TEST("ActorCache coalescing") {
909 ActorCacheTest test;
910 auto& ws = test.ws;
911 auto& mockStorage = test.mockStorage;
912 
913 // Create a scenario where we have several cached and uncached keys.
914 // foo, bar = cached with values
915 // baz, qux = cached as absent
916 // corge, grault, others = not cached
917 {
918 auto promise = expectUncached(test.get({"foo", "bar", "baz", "qux"}));
919 
920 mockStorage->expectCall("getMultiple", ws)
921 .withParams(CAPNP(keys = [ "bar", "baz", "foo", "qux" ]), "stream"_kj)
922 .useCallback("stream", [&](MockClient stream) {
923 stream
924 .call("values",
925 CAPNP(list = [ (key = "bar", value = "456"), (key = "foo", value = "123") ]))
926 .expectReturns(CAPNP(), ws);
927 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
928 }).expectCanceled();
929 
930 auto results = promise.wait(ws);
931 KJ_ASSERT(results == kvs({{"bar", "456"}, {"foo", "123"}}));
932 }
933 
934 // Now do several puts and deletes that overwrite each other, and make sure they coalesce
935 // properly.
936 {
937 test.put({{"bar", "654"}, {"qux", "555"}, {"corge", "789"}});
938 test.put("corge", "987");
939 auto promise1 = expectUncached(test.delete_({"bar"_kj, "grault"_kj}));
940 KJ_ASSERT(expectCached(test.delete_("foo")));
941 auto promise2 = expectUncached(test.delete_({"garply"_kj, "waldo"_kj, "fred"_kj}));
942 
943 // Note this final put undoes a delete. However, the delete was of a key not in cache, so it
944 // still has to be performed in order to produce the deletion count.
945 test.put("waldo", "odlaw");
946 
947 auto values = expectCached(test.get({"foo"_kj, "bar"_kj, "baz"_kj, "qux"_kj, "corge"_kj,
948 "grault"_kj, "garply"_kj, "waldo"_kj, "fred"_kj}));
949 KJ_EXPECT(values ==
950 kvs({
951 {"corge", "987"},
952 {"qux", "555"},
953 {"waldo", "odlaw"},
954 }));
955 
956 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
957 mockTxn->expectCall("delete", ws)
958 .withParams(CAPNP(keys = ["grault"]))
959 .thenReturn(CAPNP(numDeleted = 0));
960 mockTxn->expectCall("delete", ws)
961 .withParams(CAPNP(keys = [ "garply", "waldo", "fred" ]))
962 .thenReturn(CAPNP(numDeleted = 2));
963 mockTxn->expectCall("delete", ws)
964 .withParams(CAPNP(keys = [ "bar", "foo" ]))
965 .thenReturn(CAPNP(numDeleted = 65382)); // count is ignored
966 mockTxn->expectCall("put", ws)
967 .withParams(CAPNP(entries =
968 [
969 (key = "qux", value = "555"), (key = "corge", value = "987"),
970 (key = "waldo", value = "odlaw")
971 ]))
972 .thenReturn(CAPNP());
973 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
974 mockTxn->expectDropped(ws);
975 
976 KJ_ASSERT(promise1.wait(ws) == 1);
977 KJ_ASSERT(promise2.wait(ws) == 2);
978 }
979}
980 
981KJ_TEST("ActorCache canceled deletes are coalesced") {
982 ActorCacheTest test;
983 auto& ws = test.ws;
984 auto& mockStorage = test.mockStorage;
985 
986 // A bunch of deletes where we immediately drop the returned promises.
987 (void)expectUncached(test.delete_("foo"));
988 (void)expectUncached(test.delete_({"bar"_kj, "baz"_kj}));
989 (void)expectUncached(test.delete_("qux"));
990 
991 // Keep one promise.
992 auto promise = expectUncached(test.delete_("corge"));
993 
994 // Overwrite one of them.
995 test.put("qux", "blah");
996 
997 // The deletes where the caller stopped listening will be coalesced into one, or dropped entirely
998 // if overwritten by a later put().
999 {
1000 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
1001 mockTxn->expectCall("delete", ws)
1002 .withParams(CAPNP(keys = ["corge"]))
1003 .thenReturn(CAPNP(numDeleted = 0));
1004 mockTxn->expectCall("delete", ws)
1005 .withParams(CAPNP(keys = [ "foo", "bar", "baz" ]))
1006 .thenReturn(CAPNP(numDeleted = 1234)); // count ignored
1007 mockTxn->expectCall("put", ws)
1008 .withParams(CAPNP(entries = [(key = "qux", value = "blah")]))
1009 .thenReturn(CAPNP());
1010 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
1011 mockTxn->expectDropped(ws);
1012 }
1013 
1014 KJ_ASSERT(!promise.wait(ws));
1015}
1016 
1017KJ_TEST("ActorCache get-put ordering") {
1018 ActorCacheTest test;
1019 auto& ws = test.ws;
1020 auto& mockStorage = test.mockStorage;
1021 
1022 // Initiate a get, followed by a put and a delete that affect the same keys. Since the get()
1023 // started first, its final results later should not reflect the put and delete.
1024 auto promise1 = expectUncached(test.get({"foo"_kj, "bar"_kj, "baz"_kj}));
1025 test.put("foo", "123");
1026 auto deletePromise = expectUncached(test.delete_("bar"));
1027 
1028 // Verify cache content.
1029 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
1030 KJ_ASSERT(expectCached(test.get("bar")) == nullptr);
1031 
1032 // Start another get. This time, "foo" and "bar" will be served from cache, but "baz" is still
1033 // on disk. This means this get won't complete immediately. We'll then overwrite the value of
1034 // "bar", but hope that the get() has already picked up the cached value for consistency.
1035 auto promise2 = expectUncached(test.get({"foo"_kj, "bar"_kj, "baz"_kj}));
1036 test.put("bar", "456");
1037 
1038 // Verify cache content.
1039 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
1040 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "456");
1041 
1042 // Expect to receive the storage gets. But, don't return from them yet!
1043 KJ_ASSERT(!promise1.poll(ws));
1044 auto mockGet1 = mockStorage->expectCall("getMultiple", ws)
1045 .withParams(CAPNP(keys = [ "bar", "baz", "foo" ]), "stream"_kj);
1046 
1047 KJ_ASSERT(!promise2.poll(ws));
1048 auto mockGet2 =
1049 mockStorage->expectCall("getMultiple", ws).withParams(CAPNP(keys = ["baz"]), "stream"_kj);
1050 
1051 // No writes will be done until our reads finish!
1052 mockStorage->expectNoActivity(ws);
1053 
1054 // Let's have the second read complete first.
1055 kj::mv(mockGet2)
1056 .useCallback("stream", [&](MockClient stream) {
1057 stream.call("values", CAPNP(list = [(key = "baz", value = "987")])).expectReturns(CAPNP(), ws);
1058 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
1059 }).expectCanceled();
1060 
1061 // The completed read returns cached results as of when it was called, merged with what it
1062 // read from disk.
1063 KJ_ASSERT(promise2.wait(ws) == kvs({{"baz", "987"}, {"foo", "123"}}));
1064 
1065 // The completed read brought "baz" into cache but didn't change "foo" or "bar".
1066 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
1067 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "456");
1068 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "987");
1069 
1070 // Still no writes because the first read is still outstanding.
1071 mockStorage->expectNoActivity(ws);
1072 
1073 // Finally, have the first read complete.
1074 kj::mv(mockGet1)
1075 .useCallback("stream", [&](MockClient stream) {
1076 stream
1077 .call("values",
1078 CAPNP(list =
1079 [
1080 (key = "bar", value = "654"), (key = "baz", value = "987"),
1081 (key = "foo", value = "321")
1082 ]))
1083 .expectReturns(CAPNP(), ws);
1084 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
1085 }).expectCanceled();
1086 
1087 // This returns exactly what came off disk, not reflecting any later writes.
1088 KJ_ASSERT(promise1.wait(ws) == kvs({{"bar", "654"}, {"baz", "987"}, {"foo", "321"}}));
1089 
1090 // The completed read didn't mess with the cache.
1091 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
1092 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "456");
1093 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "987");
1094 
1095 // Next up, the flush transaction proceeds.
1096 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
1097 mockTxn->expectCall("delete", ws)
1098 .withParams(CAPNP(keys = ["bar"]))
1099 .thenReturn(CAPNP(numDeleted = 1));
1100 mockTxn->expectCall("put", ws)
1101 .withParams(CAPNP(entries = [ (key = "foo", value = "123"), (key = "bar", value = "456") ]))
1102 .thenReturn(CAPNP());
1103 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
1104 mockTxn->expectDropped(ws);
1105 
1106 // Our delete finally finished.
1107 KJ_ASSERT(deletePromise.wait(ws) == 1);
1108 
1109 // Cache is still good.
1110 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
1111 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "456");
1112 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "987");
1113}
1114 
1115KJ_TEST("ActorCache put during flush") {
1116 ActorCacheTest test;
1117 auto& ws = test.ws;
1118 auto& mockStorage = test.mockStorage;
1119 
1120 test.put({{"foo", "123"}, {"bar", "456"}});
1121 
1122 {
1123 auto inProgressFlush = mockStorage->expectCall("put", ws).withParams(
1124 CAPNP(entries = [ (key = "foo", value = "123"), (key = "bar", value = "456") ]));
1125 
1126 // We're in the middle of flushing... do a put. Should be fine.
1127 test.put("bar", "654");
1128 
1129 kj::mv(inProgressFlush).thenReturn(CAPNP());
1130 }
1131 
1132 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
1133 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "654");
1134 
1135 {
1136 mockStorage->expectCall("put", ws)
1137 .withParams(CAPNP(entries = [(key = "bar", value = "654")]))
1138 .thenReturn(CAPNP());
1139 }
1140 
1141 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
1142 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "654");
1143}
1144 
1145KJ_TEST("ActorCache flush retry") {
1146 ActorCacheTest test;
1147 auto& ws = test.ws;
1148 auto& mockStorage = test.mockStorage;
1149 
1150 test.put({{"foo", "123"}, {"bar", "456"}, {"baz", "789"}});
1151 auto promise1 = expectUncached(test.delete_({"qux"_kj, "quux"_kj}));
1152 auto promise2 = expectUncached(test.delete_({"corge"_kj, "grault"_kj}));
1153 
1154 {
1155 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
1156 // One delete succeeds, the other throws (later).
1157 mockTxn->expectCall("delete", ws)
1158 .withParams(CAPNP(keys = [ "qux", "quux" ]))
1159 .thenReturn(CAPNP(numDeleted = 1));
1160 auto mockDelete =
1161 mockTxn->expectCall("delete", ws).withParams(CAPNP(keys = [ "corge", "grault" ]));
1162 mockTxn->expectCall("put", ws)
1163 .withParams(CAPNP(entries =
1164 [
1165 (key = "foo", value = "123"), (key = "bar", value = "456"),
1166 (key = "baz", value = "789")
1167 ]))
1168 .thenReturn(CAPNP());
1169 
1170 // While the transaction is outstanding, some more puts and deletes mess with things...
1171 test.put("bar", "654");
1172 KJ_ASSERT(expectCached(test.delete_("baz")));
1173 test.put("qux", "987");
1174 test.put("corge", "555");
1175 
1176 kj::mv(mockDelete).thenThrow(KJ_EXCEPTION(DISCONNECTED, "delete failed"));
1177 mockTxn->expectCall("commit", ws).thenThrow(KJ_EXCEPTION(DISCONNECTED, "flush failed"));
1178 mockTxn->expectDropped(ws);
1179 }
1180 
1181 // Verify cache.
1182 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
1183 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "654");
1184 KJ_ASSERT(expectCached(test.get("baz")) == nullptr);
1185 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("qux"))) == "987");
1186 KJ_ASSERT(expectCached(test.get("quux")) == nullptr);
1187 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("corge"))) == "555");
1188 KJ_ASSERT(expectCached(test.get("grault")) == nullptr);
1189 
1190 // Although the counted delete succeeded, the promise will not resolve until our flush succeeds!
1191 KJ_ASSERT(!promise1.poll(ws));
1192 // The second delete failed and is also still outstanding until a flush succeeds.
1193 KJ_ASSERT(!promise2.poll(ws));
1194 
1195 // The transaction will be retried, with the updated puts and deletes.
1196 {
1197 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
1198 // Note that "corge" is still the subject of a delete, even though it has since been
1199 // overwritten by a put, because we still need to count the delete. "qux", on the other
1200 // hand, no longer needs counting, and has also been overwritten by a put(), so it doesn't
1201 // need to be deleted anymore. "quux" is still deleted, even though the count was returned
1202 // last time, because it hasn't been further overwritten, and that delete from last time
1203 // wasn't actually committed.
1204 mockTxn->expectCall("delete", ws)
1205 .withParams(CAPNP(keys = ["quux"]))
1206 .thenReturn(CAPNP()); // count ignored because we got it on the first try!
1207 mockTxn->expectCall("delete", ws)
1208 .withParams(CAPNP(keys = [ "corge", "grault" ]))
1209 .thenReturn(CAPNP(numDeleted = 2));
1210 mockTxn->expectCall("delete", ws)
1211 .withParams(CAPNP(keys = ["baz"]))
1212 .thenReturn(CAPNP(numDeleted = 1234)); // count ignored
1213 mockTxn->expectCall("put", ws)
1214 .withParams(CAPNP(entries =
1215 [
1216 (key = "foo", value = "123"), (key = "bar", value = "654"),
1217 (key = "qux", value = "987"), (key = "corge", value = "555")
1218 ]))
1219 .thenReturn(CAPNP());
1220 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
1221 mockTxn->expectDropped(ws);
1222 }
1223 
1224 // Verify cache.
1225 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
1226 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "654");
1227 KJ_ASSERT(expectCached(test.get("baz")) == nullptr);
1228 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("qux"))) == "987");
1229 KJ_ASSERT(expectCached(test.get("quux")) == nullptr);
1230 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("corge"))) == "555");
1231 KJ_ASSERT(expectCached(test.get("grault")) == nullptr);
1232 
1233 // Second delete finished this time.
1234 KJ_ASSERT(promise2.wait(ws) == 2);
1235 
1236 // Our flush has succeeded and we've obtained our count!
1237 KJ_ASSERT(promise1.wait(ws) == 1);
1238}
1239 
1240KJ_TEST("ActorCache output gate blocked during flush") {
1241 ActorCacheTest test({.monitorOutputGate = false});
1242 auto& ws = test.ws;
1243 auto& mockStorage = test.mockStorage;
1244 
1245 // Gate is currently not blocked.
1246 KJ_ASSERT(test.gate.wait(nullptr).poll(ws));
1247 
1248 // Do a put.
1249 test.put("foo", "123");
1250 test.delete_("bar");
1251 
1252 // Now it is blocked.
1253 auto gatePromise = test.gate.wait(nullptr);
1254 KJ_ASSERT(!gatePromise.poll(ws));
1255 
1256 // Complete the transaction.
1257 {
1258 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
1259 mockTxn->expectCall("delete", ws)
1260 .withParams(CAPNP(keys = ["bar"]))
1261 .thenReturn(CAPNP(numDeleted = 0));
1262 mockTxn->expectCall("put", ws)
1263 .withParams(CAPNP(entries = [(key = "foo", value = "123")]))
1264 .thenReturn(CAPNP());
1265 auto commitCall = mockTxn->expectCall("commit", ws);
1266 
1267 // Still blocked until commit completes.
1268 KJ_ASSERT(!gatePromise.poll(ws));
1269 
1270 kj::mv(commitCall).thenReturn(CAPNP());
1271 mockTxn->expectDropped(ws);
1272 }
1273 
1274 KJ_ASSERT(gatePromise.poll(ws));
1275 gatePromise.wait(ws);
1276}
1277 
1278KJ_TEST("ActorCache output gate bypass") {
1279 ActorCacheTest test({.monitorOutputGate = false});
1280 auto& ws = test.ws;
1281 auto& mockStorage = test.mockStorage;
1282 
1283 // Gate is currently not blocked.
1284 test.gate.wait(nullptr).wait(ws);
1285 
1286 // Do a put.
1287 test.put("foo", "123", {.allowUnconfirmed = true});
1288 
1289 // Gate still isn't blocked, because we set `allowUnconfirmed`.
1290 test.gate.wait(nullptr).wait(ws);
1291 
1292 // Complete the transaction.
1293 {
1294 mockStorage->expectCall("put", ws)
1295 .withParams(CAPNP(entries = [(key = "foo", value = "123")]))
1296 .thenReturn(CAPNP());
1297 }
1298 
1299 test.gate.wait(nullptr).wait(ws);
1300}
1301 
1302KJ_TEST("ActorCache output gate bypass on one put but not the next") {
1303 ActorCacheTest test({.monitorOutputGate = false});
1304 auto& ws = test.ws;
1305 auto& mockStorage = test.mockStorage;
1306 
1307 // Gate is currently not blocked.
1308 test.gate.wait(nullptr).wait(ws);
1309 
1310 // Do two puts, only bypassing on the first. The net result should be that the output gate is
1311 // in effect.
1312 test.put("foo", "123", {.allowUnconfirmed = true});
1313 test.put("bar", "456");
1314 
1315 // Now it is blocked.
1316 auto gatePromise = test.gate.wait(nullptr);
1317 KJ_ASSERT(!gatePromise.poll(ws));
1318 
1319 // Complete the transaction.
1320 {
1321 auto inProgressFlush = mockStorage->expectCall("put", ws).withParams(
1322 CAPNP(entries = [ (key = "foo", value = "123"), (key = "bar", value = "456") ]));
1323 
1324 // Still blocked until the flush completes.
1325 KJ_ASSERT(!gatePromise.poll(ws));
1326 
1327 kj::mv(inProgressFlush).thenReturn(CAPNP());
1328 }
1329 
1330 KJ_ASSERT(gatePromise.poll(ws));
1331 gatePromise.wait(ws);
1332}
1333 
1334KJ_TEST("ActorCache flush hard failure") {
1335 ActorCacheTest test({.monitorOutputGate = false});
1336 auto& ws = test.ws;
1337 auto& mockStorage = test.mockStorage;
1338 
1339 auto promise = test.gate.onBroken();
1340 
1341 test.put("foo", "123");
1342 
1343 KJ_ASSERT(!promise.poll(ws));
1344 
1345 {
1346 mockStorage->expectCall("put", ws)
1347 .withParams(CAPNP(entries = [(key = "foo", value = "123")]))
1348 .thenThrow(KJ_EXCEPTION(FAILED, "jsg.Error: flush failed hard"));
1349 }
1350 
1351 KJ_EXPECT_THROW_MESSAGE(
1352 "broken.outputGateBroken; jsg.Error: flush failed hard", promise.wait(ws));
1353 
1354 // Further writes won't even try to start any new transactions because the failure killed them all.
1355 test.put("bar", "456");
1356}
1357 
1358KJ_TEST("ActorCache flush hard failure with output gate bypass") {
1359 ActorCacheTest test({.monitorOutputGate = false});
1360 auto& ws = test.ws;
1361 auto& mockStorage = test.mockStorage;
1362 
1363 auto promise = test.gate.onBroken();
1364 
1365 test.put("foo", "123", {.allowUnconfirmed = true});
1366 
1367 // The output gate is not applied.
1368 test.gate.wait(nullptr).wait(ws);
1369 KJ_ASSERT(!promise.poll(ws));
1370 
1371 {
1372 mockStorage->expectCall("put", ws)
1373 .withParams(CAPNP(entries = [(key = "foo", value = "123")]))
1374 .thenThrow(KJ_EXCEPTION(FAILED, "jsg.Error: flush failed hard"));
1375 }
1376 
1377 // The failure was still propagated to the output gate.
1378 KJ_EXPECT_THROW_MESSAGE("flush failed hard", promise.wait(ws));
1379 KJ_EXPECT_THROW_MESSAGE("flush failed hard", test.gate.wait(nullptr).wait(ws));
1380 
1381 // Further writes won't even try to start any new transactions because the failure killed them all.
1382 test.put("bar", "456");
1383}
1384 
1385KJ_TEST("ActorCache read retry") {
1386 ActorCacheTest test;
1387 auto& ws = test.ws;
1388 auto& mockStorage = test.mockStorage;
1389 
1390 auto promise = expectUncached(test.get("foo"));
1391 test.put("bar", "456");
1392 test.delete_("baz");
1393 
1394 // Expect the get, but don't resolve yet.
1395 auto mockGet = mockStorage->expectCall("get", ws).withParams(CAPNP(key = "foo"));
1396 
1397 // No activity because reads are outstanding.
1398 mockStorage->expectNoActivity(ws);
1399 
1400 // Fail out the read with a disconnect.
1401 kj::mv(mockGet).thenThrow(KJ_EXCEPTION(DISCONNECTED, "read failed"));
1402 
1403 // It will be retried.
1404 auto mockGet2 = mockStorage->expectCall("get", ws).withParams(CAPNP(key = "foo"));
1405 
1406 // Still no activity because of the read.
1407 mockStorage->expectNoActivity(ws);
1408 
1409 // Finish it.
1410 kj::mv(mockGet2).thenReturn(CAPNP(value = "123"));
1411 
1412 // Now the transaction starts actually writing (and completes).
1413 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
1414 mockTxn->expectCall("delete", ws)
1415 .withParams(CAPNP(keys = ["baz"]))
1416 .thenReturn(CAPNP(numDeleted = 0));
1417 mockTxn->expectCall("put", ws)
1418 .withParams(CAPNP(entries = [(key = "bar", value = "456")]))
1419 .thenReturn(CAPNP());
1420 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
1421 mockTxn->expectDropped(ws);
1422 
1423 // And the read finishes.
1424 KJ_ASSERT(KJ_ASSERT_NONNULL(promise.wait(ws)) == "123");
1425}
1426 
1427KJ_TEST("ActorCache read retry on flush containing only puts") {
1428 ActorCacheTest test;
1429 auto& ws = test.ws;
1430 auto& mockStorage = test.mockStorage;
1431 
1432 auto promise = expectUncached(test.get("foo"));
1433 test.put("bar", "456");
1434 
1435 // Expect the get, but don't resolve yet.
1436 auto mockGet = mockStorage->expectCall("get", ws).withParams(CAPNP(key = "foo"));
1437 
1438 // No activity on the flush yet (not even starting a txn), because reads are outstanding.
1439 mockStorage->expectNoActivity(ws);
1440 
1441 // Fail out the read with a disconnect.
1442 kj::mv(mockGet).thenThrow(KJ_EXCEPTION(DISCONNECTED, "read failed"));
1443 
1444 // It will be retried.
1445 auto mockGet2 = mockStorage->expectCall("get", ws).withParams(CAPNP(key = "foo"));
1446 
1447 // Still no transaction activity.
1448 mockStorage->expectNoActivity(ws);
1449 
1450 // Finish it.
1451 kj::mv(mockGet2).thenReturn(CAPNP(value = "123"));
1452 
1453 // Now the transaction starts actually writing (and completes).
1454 mockStorage->expectCall("put", ws)
1455 .withParams(CAPNP(entries = [(key = "bar", value = "456")]))
1456 .thenReturn(CAPNP());
1457 
1458 // And the read finishes.
1459 KJ_ASSERT(KJ_ASSERT_NONNULL(promise.wait(ws)) == "123");
1460}
1461 
1462KJ_TEST("ActorCache read hard fail") {
1463 ActorCacheTest test;
1464 auto& ws = test.ws;
1465 auto& mockStorage = test.mockStorage;
1466 
1467 // Don't use expectCached() this time because we don't want eagerlyReportExceptions(), because
1468 // we actually expect an exception.
1469 auto promise = test.get("foo").get<kj::Promise<kj::Maybe<kj::String>>>();
1470 test.put("bar", "456");
1471 test.delete_("baz");
1472 
1473 // Expect the get, but don't resolve yet.
1474 auto mockGet = mockStorage->expectCall("get", ws).withParams(CAPNP(key = "foo"));
1475 
1476 // We won't write anything until the read completes.
1477 mockStorage->expectNoActivity(ws);
1478 
1479 // Fail out the read with non-disconnect.
1480 kj::mv(mockGet).thenThrow(KJ_EXCEPTION(FAILED, "read failed"));
1481 
1482 // The read propagates the error.
1483 KJ_EXPECT_THROW_MESSAGE("read failed", promise.wait(ws));
1484 
1485 // The read is NOT retried, so expect the transaction to run now.
1486 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
1487 mockTxn->expectCall("delete", ws)
1488 .withParams(CAPNP(keys = ["baz"]))
1489 .thenReturn(CAPNP(numDeleted = 0));
1490 mockTxn->expectCall("put", ws)
1491 .withParams(CAPNP(entries = [(key = "bar", value = "456")]))
1492 .thenReturn(CAPNP());
1493 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
1494 mockTxn->expectDropped(ws);
1495 
1496 // The read is NOT retried.
1497 mockStorage->expectNoActivity(ws);
1498}
1499 
1500KJ_TEST("ActorCache read cancel") {
1501 ActorCacheTest test;
1502 auto& ws = test.ws;
1503 auto& mockStorage = test.mockStorage;
1504 
1505 auto promise = expectUncached(test.get("foo"));
1506 test.put("bar", "456");
1507 test.delete_("baz");
1508 
1509 // Expect the get, but intentionally don't resolve it.
1510 auto mockGet = mockStorage->expectCall("get", ws).withParams(CAPNP(key = "foo"));
1511 
1512 // We won't write anything until the read completes.
1513 mockStorage->expectNoActivity(ws);
1514 
1515 // Cancel the read.
1516 promise = nullptr;
1517 mockGet.expectCanceled();
1518 
1519 // The transaction proceeds.
1520 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
1521 mockTxn->expectCall("delete", ws)
1522 .withParams(CAPNP(keys = ["baz"]))
1523 .thenReturn(CAPNP(numDeleted = 0));
1524 mockTxn->expectCall("put", ws)
1525 .withParams(CAPNP(entries = [(key = "bar", value = "456")]))
1526 .thenReturn(CAPNP());
1527 
1528 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
1529 mockTxn->expectDropped(ws);
1530}
1531 
1532KJ_TEST("ActorCache get-multiple multiple blocks") {
1533 ActorCacheTest test;
1534 auto& ws = test.ws;
1535 auto& mockStorage = test.mockStorage;
1536 
1537 auto promise = expectUncached(test.get({"foo"_kj, "bar"_kj, "baz"_kj, "qux"_kj, "corge"_kj}));
1538 
1539 mockStorage->expectCall("getMultiple", ws)
1540 .withParams(CAPNP(keys = [ "bar", "baz", "corge", "foo", "qux" ]), "stream"_kj)
1541 .useCallback("stream", [&](MockClient stream) {
1542 stream.call("values", CAPNP(list = [(key = "baz", value = "456")])).expectReturns(CAPNP(), ws);
1543 
1544 // At this point, "bar" and "baz" are considered cached.
1545 KJ_ASSERT(expectCached(test.get("bar")) == nullptr);
1546 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "456");
1547 (void)expectUncached(test.get("corge"));
1548 (void)expectUncached(test.get("foo"));
1549 (void)expectUncached(test.get("qux"));
1550 
1551 stream.call("values", CAPNP(list = [(key = "foo", value = "789")])).expectReturns(CAPNP(), ws);
1552 
1553 // At this point, everything except "qux" is cached.
1554 KJ_ASSERT(expectCached(test.get("bar")) == nullptr);
1555 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "456");
1556 KJ_ASSERT(expectCached(test.get("corge")) == nullptr);
1557 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "789");
1558 (void)expectUncached(test.get("qux"));
1559 
1560 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
1561 
1562 // Now it's all cached.
1563 KJ_ASSERT(expectCached(test.get("bar")) == nullptr);
1564 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "456");
1565 KJ_ASSERT(expectCached(test.get("corge")) == nullptr);
1566 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "789");
1567 KJ_ASSERT(expectCached(test.get("qux")) == nullptr);
1568 }).expectCanceled();
1569 
1570 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "456"}, {"foo", "789"}}));
1571}
1572 
1573KJ_TEST("ActorCache get-multiple partial retry") {
1574 ActorCacheTest test;
1575 auto& ws = test.ws;
1576 auto& mockStorage = test.mockStorage;
1577 
1578 auto promise = expectUncached(test.get({"foo"_kj, "bar"_kj, "baz"_kj, "qux"_kj}));
1579 
1580 mockStorage->expectCall("getMultiple", ws)
1581 .withParams(CAPNP(keys = [ "bar", "baz", "foo", "qux" ]), "stream"_kj)
1582 .useCallback("stream", [&](MockClient stream) {
1583 stream.call("values", CAPNP(list = [(key = "baz", value = "456")])).expectReturns(CAPNP(), ws);
1584 }).thenThrow(KJ_EXCEPTION(DISCONNECTED, "read failed"));
1585 
1586 mockStorage
1587 ->expectCall("getMultiple", ws)
1588 // Since "baz" was received, the caller knows that it only has to retry keys after that.
1589 .withParams(CAPNP(keys = [ "foo", "qux" ]), "stream"_kj)
1590 .useCallback("stream", [&](MockClient stream) {
1591 stream.call("values", CAPNP(list = [(key = "qux", value = "789")])).expectReturns(CAPNP(), ws);
1592 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
1593 }).expectCanceled();
1594 
1595 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "456"}, {"qux", "789"}}));
1596}
1597 
1598// =======================================================================================
1599// OK... time for hard mode. Let's test list().
1600 
1601KJ_TEST("ActorCache list()") {
1602 ActorCacheTest test;
1603 auto& ws = test.ws;
1604 auto& mockStorage = test.mockStorage;
1605 
1606 {
1607 auto promise = expectUncached(test.list("bar", "qux"));
1608 
1609 mockStorage->expectCall("list", ws)
1610 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
1611 .useCallback("stream", [&](MockClient stream) {
1612 stream
1613 .call("values",
1614 CAPNP(list =
1615 [
1616 (key = "bar", value = "456"), (key = "baz", value = "789"),
1617 (key = "foo", value = "123")
1618 ]))
1619 .expectReturns(CAPNP(), ws);
1620 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
1621 }).expectCanceled();
1622 
1623 KJ_ASSERT(promise.wait(ws) == kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
1624 }
1625 
1626 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
1627 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "456");
1628 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "789");
1629 
1630 // Stuff in range that wasn't reported is cached as absent.
1631 KJ_ASSERT(expectCached(test.get("bara")) == nullptr);
1632 KJ_ASSERT(expectCached(test.get("corge")) == nullptr);
1633 KJ_ASSERT(expectCached(test.get("quw")) == nullptr);
1634 
1635 // Listing the same range again is fully cached.
1636 KJ_ASSERT(expectCached(test.list("bar", "qux")) ==
1637 kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
1638 
1639 // Limits can be applied to the cached results.
1640 KJ_ASSERT(expectCached(test.list("bar", "qux", 0u)) == kvs({}));
1641 KJ_ASSERT(expectCached(test.list("bar", "qux", 1)) == kvs({{"bar", "456"}}));
1642 KJ_ASSERT(expectCached(test.list("bar", "qux", 2)) == kvs({{"bar", "456"}, {"baz", "789"}}));
1643 KJ_ASSERT(expectCached(test.list("bar", "qux", 3)) ==
1644 kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
1645 KJ_ASSERT(expectCached(test.list("bar", "qux", 4)) ==
1646 kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
1647 KJ_ASSERT(expectCached(test.list("bar", "qux", 1000)) ==
1648 kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
1649 
1650 // The endpoint of the list is not cached.
1651 {
1652 auto promise = expectUncached(test.get("qux"));
1653 
1654 mockStorage->expectCall("get", ws)
1655 .withParams(CAPNP(key = "qux"))
1656 .thenReturn(CAPNP(value = "555"));
1657 
1658 auto result = KJ_ASSERT_NONNULL(promise.wait(ws));
1659 KJ_EXPECT(result == "555");
1660 }
1661}
1662 
1663KJ_TEST("ActorCache list() all") {
1664 ActorCacheTest test;
1665 auto& ws = test.ws;
1666 auto& mockStorage = test.mockStorage;
1667 
1668 {
1669 auto promise = expectUncached(test.list(nullptr, nullptr));
1670 
1671 mockStorage->expectCall("list", ws)
1672 .withParams(CAPNP(), "stream"_kj)
1673 .useCallback("stream", [&](MockClient stream) {
1674 stream
1675 .call("values",
1676 CAPNP(list =
1677 [
1678 (key = "bar", value = "456"), (key = "baz", value = "789"),
1679 (key = "foo", value = "123")
1680 ]))
1681 .expectReturns(CAPNP(), ws);
1682 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
1683 }).expectCanceled();
1684 
1685 KJ_ASSERT(promise.wait(ws) == kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
1686 }
1687 
1688 KJ_ASSERT(expectCached(test.get("")) == nullptr);
1689 KJ_ASSERT(expectCached(test.list(nullptr, nullptr)) ==
1690 kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
1691 KJ_ASSERT(expectCached(test.list("bar", "qux")) ==
1692 kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
1693 KJ_ASSERT(expectCached(test.list(nullptr, nullptr)) ==
1694 kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
1695 KJ_ASSERT(expectCached(test.list("baz", nullptr)) == kvs({{"baz", "789"}, {"foo", "123"}}));
1696 KJ_ASSERT(expectCached(test.list(nullptr, "foo")) == kvs({{"bar", "456"}, {"baz", "789"}}));
1697}
1698 
1699KJ_TEST("ActorCache list() with limit") {
1700 ActorCacheTest test;
1701 auto& ws = test.ws;
1702 auto& mockStorage = test.mockStorage;
1703 
1704 {
1705 auto promise = expectUncached(test.list("bar", "qux", 3));
1706 
1707 mockStorage->expectCall("list", ws)
1708 .withParams(CAPNP(start = "bar", end = "qux", limit = 3), "stream"_kj)
1709 .useCallback("stream", [&](MockClient stream) {
1710 stream
1711 .call("values",
1712 CAPNP(list =
1713 [
1714 (key = "bar", value = "456"), (key = "baz", value = "789"),
1715 (key = "foo", value = "123")
1716 ]))
1717 .expectReturns(CAPNP(), ws);
1718 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
1719 }).expectCanceled();
1720 
1721 KJ_ASSERT(promise.wait(ws) == kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
1722 }
1723 
1724 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
1725 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "456");
1726 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "789");
1727 
1728 // Stuff in range that wasn't reported is cached as absent -- but not past the last reported
1729 // value, which was "foo".
1730 KJ_ASSERT(expectCached(test.get("bara")) == nullptr);
1731 KJ_ASSERT(expectCached(test.get("corge")) == nullptr);
1732 KJ_ASSERT(expectCached(test.get("fon")) == nullptr);
1733 
1734 // Stuff after the last key is not in cache.
1735 (void)expectUncached(test.get("fooa"));
1736 
1737 // Listing the same range again, with the same limit or lower, is fully cached.
1738 KJ_ASSERT(expectCached(test.list("bar", "qux", 3)) ==
1739 kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
1740 KJ_ASSERT(expectCached(test.list("bar", "qux", 2)) == kvs({{"bar", "456"}, {"baz", "789"}}));
1741 KJ_ASSERT(expectCached(test.list("bar", "qux", 1)) == kvs({{"bar", "456"}}));
1742 KJ_ASSERT(expectCached(test.list("bar", "qux", 0u)) == kvs({}));
1743 
1744 // But a larger limit won't be cached.
1745 {
1746 auto promise = expectUncached(test.list("bar", "qux", 4));
1747 
1748 // The new list will start at "foo\0" with a limit of 1, so that it won't redundantly list foo
1749 // itself and will only get the one remaining key that it needs.
1750 mockStorage->expectCall("list", ws)
1751 .withParams(CAPNP(start = "foo\0", end = "qux", limit = 1), "stream"_kj)
1752 .useCallback("stream", [&](MockClient stream) {
1753 stream.call("values", CAPNP(list = [(key = "garply", value = "54321")]))
1754 .expectReturns(CAPNP(), ws);
1755 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
1756 }).expectCanceled();
1757 
1758 KJ_ASSERT(promise.wait(ws) ==
1759 kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}, {"garply", "54321"}}));
1760 }
1761 
1762 // Cached if we try it again though.
1763 KJ_ASSERT(expectCached(test.list("bar", "qux", 4)) ==
1764 kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}, {"garply", "54321"}}));
1765}
1766 
1767KJ_TEST("ActorCache list() with limit around negative entries") {
1768 // This checks for a bug where the initial scan through cache for list() applies the limit to
1769 // the total number of entries seen (positive or negative), when it really needs to apply only
1770 // to positive entries.
1771 
1772 ActorCacheTest test;
1773 auto& ws = test.ws;
1774 auto& mockStorage = test.mockStorage;
1775 
1776 // Set up a bunch of negative entries and a positive one after them.
1777 test.delete_({"bar1"_kj, "bar2"_kj, "bar3"_kj, "bar4"_kj});
1778 test.put("baz", "789");
1779 
1780 // Now do a list through them. It should see the positive entry in cache.
1781 {
1782 auto promise = expectUncached(test.list("bar", "qux", 3));
1783 
1784 mockStorage->expectCall("list", ws)
1785 .withParams(CAPNP(start = "bar", end = "qux", limit = 7), "stream"_kj)
1786 .useCallback("stream", [&](MockClient stream) {
1787 stream
1788 .call("values",
1789 CAPNP(list =
1790 [
1791 (key = "bar", value = "456"), (key = "bar1", value = "xxx"),
1792 (key = "bar3", value = "yyy"), (key = "foo", value = "123")
1793 ]))
1794 .expectReturns(CAPNP(), ws);
1795 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
1796 }).expectCanceled();
1797 
1798 KJ_ASSERT(promise.wait(ws) == kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
1799 }
1800 
1801 KJ_ASSERT(expectCached(test.list("bar", "qux", 4)) ==
1802 kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
1803 
1804 // Acknowledge the transaction.
1805 {
1806 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
1807 mockTxn->expectCall("delete", ws).thenReturn(CAPNP());
1808 mockTxn->expectCall("put", ws).thenReturn(CAPNP());
1809 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
1810 mockTxn->expectDropped(ws);
1811 }
1812}
1813 
1814KJ_TEST("ActorCache list() start point is not present") {
1815 ActorCacheTest test;
1816 auto& ws = test.ws;
1817 auto& mockStorage = test.mockStorage;
1818 
1819 {
1820 auto promise = expectUncached(test.list("bar", "qux"));
1821 
1822 mockStorage->expectCall("list", ws)
1823 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
1824 .useCallback("stream", [&](MockClient stream) {
1825 stream
1826 .call("values",
1827 CAPNP(list = [ (key = "baz", value = "789"), (key = "foo", value = "123") ]))
1828 .expectReturns(CAPNP(), ws);
1829 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
1830 }).expectCanceled();
1831 
1832 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "789"}, {"foo", "123"}}));
1833 }
1834 
1835 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
1836 
1837 KJ_ASSERT(expectCached(test.get("bar")) == nullptr);
1838 KJ_ASSERT(expectCached(test.get("bara")) == nullptr);
1839 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "789");
1840 KJ_ASSERT(expectCached(test.get("baza")) == nullptr);
1841 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
1842 KJ_ASSERT(expectCached(test.get("fooa")) == nullptr);
1843}
1844 
1845KJ_TEST("ActorCache list() multiple ranges") {
1846 ActorCacheTest test;
1847 auto& ws = test.ws;
1848 auto& mockStorage = test.mockStorage;
1849 
1850 {
1851 auto promise = expectUncached(test.list("a", "c"));
1852 
1853 mockStorage->expectCall("list", ws)
1854 .withParams(CAPNP(start = "a", end = "c"), "stream"_kj)
1855 .useCallback("stream", [&](MockClient stream) {
1856 stream.call("values", CAPNP(list = [ (key = "a", value = "1"), (key = "b", value = "2") ]))
1857 .expectReturns(CAPNP(), ws);
1858 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
1859 }).expectCanceled();
1860 
1861 KJ_ASSERT(promise.wait(ws) == kvs({{"a", "1"}, {"b", "2"}}));
1862 }
1863 
1864 KJ_ASSERT(expectCached(test.list("a", "c")) == kvs({{"a", "1"}, {"b", "2"}}));
1865 
1866 {
1867 auto promise = expectUncached(test.list("x", "z"));
1868 
1869 mockStorage->expectCall("list", ws)
1870 .withParams(CAPNP(start = "x", end = "z"), "stream"_kj)
1871 .useCallback("stream", [&](MockClient stream) {
1872 stream.call("values", CAPNP(list = [(key = "y", value = "9")])).expectReturns(CAPNP(), ws);
1873 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
1874 }).expectCanceled();
1875 
1876 KJ_ASSERT(promise.wait(ws) == kvs({{"y", "9"}}));
1877 }
1878 
1879 KJ_ASSERT(expectCached(test.list("a", "c")) == kvs({{"a", "1"}, {"b", "2"}}));
1880 KJ_ASSERT(expectCached(test.list("x", "z")) == kvs({{"y", "9"}}));
1881 
1882 (void)expectUncached(test.get("w"));
1883 (void)expectUncached(test.get("d"));
1884 (void)expectUncached(test.get("c"));
1885}
1886 
1887KJ_TEST("ActorCache list() with some already-cached keys in range") {
1888 ActorCacheTest test;
1889 auto& ws = test.ws;
1890 auto& mockStorage = test.mockStorage;
1891 
1892 // Initialize cache with some clean entries, both positive and negative.
1893 {
1894 auto promise1 = expectUncached(test.get("bbb"));
1895 auto promise2 = expectUncached(test.get("ccc"));
1896 
1897 mockStorage->expectCall("get", ws).withParams(CAPNP(key = "bbb")).thenReturn(CAPNP());
1898 mockStorage->expectCall("get", ws)
1899 .withParams(CAPNP(key = "ccc"))
1900 .thenReturn(CAPNP(value = "cval"));
1901 
1902 KJ_ASSERT(promise1.wait(ws) == nullptr);
1903 KJ_ASSERT(KJ_ASSERT_NONNULL(promise2.wait(ws)) == "cval");
1904 }
1905 
1906 // Also some newly-written entries, positive and negative.
1907 test.put("ddd", "dval");
1908 auto deletePromise = expectUncached(test.delete_("eee"));
1909 
1910 // Now list the range. Explicitly produce results that contradict the recent writes.
1911 {
1912 auto promise = expectUncached(test.list("aaa", "fff"));
1913 
1914 mockStorage->expectCall("list", ws)
1915 .withParams(CAPNP(start = "aaa", end = "fff"), "stream"_kj)
1916 .useCallback("stream", [&](MockClient stream) {
1917 stream
1918 .call("values",
1919 CAPNP(list = [ (key = "ccc", value = "cval"), (key = "eee", value = "eval") ]))
1920 .expectReturns(CAPNP(), ws);
1921 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
1922 }).expectCanceled();
1923 
1924 KJ_ASSERT(promise.wait(ws) == kvs({{"ccc", "cval"}, {"ddd", "dval"}}));
1925 }
1926 
1927 {
1928 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
1929 mockTxn->expectCall("delete", ws)
1930 .withParams(CAPNP(keys = ["eee"]))
1931 .thenReturn(CAPNP(numDeleted = 1));
1932 mockTxn->expectCall("put", ws)
1933 .withParams(CAPNP(entries = [(key = "ddd", value = "dval")]))
1934 .thenReturn(CAPNP());
1935 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
1936 mockTxn->expectDropped(ws);
1937 }
1938 
1939 KJ_ASSERT(deletePromise.wait(ws) == 1);
1940}
1941 
1942KJ_TEST("ActorCache list() with seemingly-redundant dirty entries") {
1943 ActorCacheTest test;
1944 auto& ws = test.ws;
1945 auto& mockStorage = test.mockStorage;
1946 
1947 // Write some stuff.
1948 auto deletePromise = expectUncached(test.delete_("bbb"));
1949 test.put("ccc", "cval");
1950 
1951 // Initiate a list operation, but don't complete it yet.
1952 auto listPromise = expectUncached(test.list("aaa", "fff"));
1953 auto listCall = mockStorage->expectCall("list", ws)
1954 .withParams(CAPNP(start = "aaa", end = "fff"), "stream"_kj);
1955 
1956 // The delete won't do any work until the list completes.
1957 mockStorage->expectNoActivity(ws);
1958 
1959 // Now write some contradictory values.
1960 test.put("bbb", "bval");
1961 KJ_ASSERT(expectCached(test.delete_("ccc")) == 1);
1962 
1963 // Now let the list complete in a way that matches what was just written.
1964 kj::mv(listCall)
1965 .useCallback("stream", [&](MockClient stream) {
1966 stream.call("values", CAPNP(list = [(key = "bbb", value = "bval")])).expectReturns(CAPNP(), ws);
1967 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
1968 }).expectCanceled();
1969 
1970 // The list produces results consistent with when it started.
1971 KJ_ASSERT(listPromise.wait(ws) == kvs({{"ccc", "cval"}}));
1972 
1973 // But the later writes are still there in cache.
1974 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bbb"))) == "bval");
1975 KJ_ASSERT(expectCached(test.get("ccc")) == nullptr);
1976 
1977 // Now the transaction runs, notably containing only the original writes, not the later writes,
1978 // despite our flush being delayed by the reads.
1979 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
1980 mockTxn->expectCall("delete", ws)
1981 .withParams(CAPNP(keys = ["bbb"]))
1982 .thenReturn(CAPNP(numDeleted = 1));
1983 mockTxn->expectCall("put", ws)
1984 .withParams(CAPNP(entries = [(key = "ccc", value = "cval")]))
1985 .thenReturn(CAPNP());
1986 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
1987 mockTxn->expectDropped(ws);
1988 KJ_ASSERT(deletePromise.wait(ws) == 1);
1989 
1990 // And then there's a new transaction to write things back to the original values.
1991 // This is NOT REDUNDANT, even though the list results seemed to match the current cached values!
1992 // (I wrote this test to prove to myself that a DIRTY entry can't be marked CLEAN just because a
1993 // read result from disk came back with the same value.)
1994 mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
1995 mockTxn->expectCall("delete", ws)
1996 .withParams(CAPNP(keys = ["ccc"]))
1997 .thenReturn(CAPNP(numDeleted = 1));
1998 mockTxn->expectCall("put", ws)
1999 .withParams(CAPNP(entries = [(key = "bbb", value = "bval")]))
2000 .thenReturn(CAPNP());
2001 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
2002 mockTxn->expectDropped(ws);
2003 
2004 // For good measure, verify list result can be served from cache.
2005 KJ_ASSERT(expectCached(test.list("aaa", "fff")) == kvs({{"bbb", "bval"}}));
2006}
2007 
2008KJ_TEST("ActorCache list() starting from known value") {
2009 ActorCacheTest test;
2010 auto& ws = test.ws;
2011 auto& mockStorage = test.mockStorage;
2012 
2013 {
2014 test.put("bar", "123");
2015 
2016 mockStorage->expectCall("put", ws)
2017 .withParams(CAPNP(entries = [(key = "bar", value = "123")]))
2018 .thenReturn(CAPNP());
2019 }
2020 
2021 {
2022 auto promise = expectUncached(test.list("bar", "qux"));
2023 
2024 mockStorage->expectCall("list", ws)
2025 .withParams(CAPNP(start = "bar\0", end = "qux"), "stream"_kj)
2026 .useCallback("stream", [&](MockClient stream) {
2027 stream.call("values", CAPNP(list = [(key = "baz", value = "456")]))
2028 .expectReturns(CAPNP(), ws);
2029 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2030 }).expectCanceled();
2031 
2032 KJ_ASSERT(promise.wait(ws) == kvs({{"bar", "123"}, {"baz", "456"}}));
2033 }
2034 
2035 KJ_ASSERT(expectCached(test.list("bar", "qux")) == kvs({{"bar", "123"}, {"baz", "456"}}));
2036}
2037 
2038KJ_TEST("ActorCache list() starting from unknown value") {
2039 ActorCacheTest test;
2040 auto& ws = test.ws;
2041 auto& mockStorage = test.mockStorage;
2042 
2043 {
2044 test.put("baz", "456");
2045 
2046 mockStorage->expectCall("put", ws)
2047 .withParams(CAPNP(entries = [(key = "baz", value = "456")]))
2048 .thenReturn(CAPNP());
2049 }
2050 
2051 {
2052 auto promise = expectUncached(test.list("bar", "qux"));
2053 
2054 mockStorage->expectCall("list", ws)
2055 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
2056 .useCallback("stream", [&](MockClient stream) {
2057 stream
2058 .call("values",
2059 CAPNP(list = [ (key = "baz", value = "456"), (key = "foo", value = "123") ]))
2060 .expectReturns(CAPNP(), ws);
2061 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2062 }).expectCanceled();
2063 
2064 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "456"}, {"foo", "123"}}));
2065 }
2066 
2067 KJ_ASSERT(expectCached(test.list("bar", "qux")) == kvs({{"baz", "456"}, {"foo", "123"}}));
2068}
2069 
2070KJ_TEST("ActorCache list() consecutively, absent midpoint") {
2071 ActorCacheTest test;
2072 auto& ws = test.ws;
2073 auto& mockStorage = test.mockStorage;
2074 
2075 {
2076 auto promise = expectUncached(test.list("bar", "corge"));
2077 
2078 mockStorage->expectCall("list", ws)
2079 .withParams(CAPNP(start = "bar", end = "corge"), "stream"_kj)
2080 .useCallback("stream", [&](MockClient stream) {
2081 stream.call("values", CAPNP(list = [(key = "baz", value = "456")]))
2082 .expectReturns(CAPNP(), ws);
2083 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2084 }).expectCanceled();
2085 
2086 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "456"}}));
2087 }
2088 
2089 {
2090 auto promise = expectUncached(test.list("corge", "qux"));
2091 
2092 mockStorage->expectCall("list", ws)
2093 .withParams(CAPNP(start = "corge", end = "qux"), "stream"_kj)
2094 .useCallback("stream", [&](MockClient stream) {
2095 stream.call("values", CAPNP(list = [(key = "foo", value = "123")]))
2096 .expectReturns(CAPNP(), ws);
2097 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2098 }).expectCanceled();
2099 
2100 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}}));
2101 }
2102 
2103 KJ_ASSERT(expectCached(test.list("bar", "qux")) == kvs({{"baz", "456"}, {"foo", "123"}}));
2104}
2105 
2106KJ_TEST("ActorCache list() consecutively reverse, absent midpoint") {
2107 ActorCacheTest test;
2108 auto& ws = test.ws;
2109 auto& mockStorage = test.mockStorage;
2110 
2111 {
2112 auto promise = expectUncached(test.list("corge", "qux"));
2113 
2114 mockStorage->expectCall("list", ws)
2115 .withParams(CAPNP(start = "corge", end = "qux"), "stream"_kj)
2116 .useCallback("stream", [&](MockClient stream) {
2117 stream.call("values", CAPNP(list = [(key = "foo", value = "123")]))
2118 .expectReturns(CAPNP(), ws);
2119 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2120 }).expectCanceled();
2121 
2122 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}}));
2123 }
2124 
2125 {
2126 auto promise = expectUncached(test.list("bar", "corge"));
2127 
2128 mockStorage->expectCall("list", ws)
2129 .withParams(CAPNP(start = "bar", end = "corge"), "stream"_kj)
2130 .useCallback("stream", [&](MockClient stream) {
2131 stream.call("values", CAPNP(list = [(key = "baz", value = "456")]))
2132 .expectReturns(CAPNP(), ws);
2133 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2134 }).expectCanceled();
2135 
2136 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "456"}}));
2137 }
2138 
2139 KJ_ASSERT(expectCached(test.list("bar", "qux")) == kvs({{"baz", "456"}, {"foo", "123"}}));
2140}
2141 
2142KJ_TEST("ActorCache list() consecutively, present midpoint") {
2143 ActorCacheTest test;
2144 auto& ws = test.ws;
2145 auto& mockStorage = test.mockStorage;
2146 
2147 {
2148 auto promise = expectUncached(test.list("bar", "corge"));
2149 
2150 mockStorage->expectCall("list", ws)
2151 .withParams(CAPNP(start = "bar", end = "corge"), "stream"_kj)
2152 .useCallback("stream", [&](MockClient stream) {
2153 stream.call("values", CAPNP(list = [(key = "baz", value = "456")]))
2154 .expectReturns(CAPNP(), ws);
2155 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2156 }).expectCanceled();
2157 
2158 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "456"}}));
2159 }
2160 
2161 {
2162 auto promise = expectUncached(test.list("corge", "qux"));
2163 
2164 mockStorage->expectCall("list", ws)
2165 .withParams(CAPNP(start = "corge", end = "qux"), "stream"_kj)
2166 .useCallback("stream", [&](MockClient stream) {
2167 stream
2168 .call("values",
2169 CAPNP(list = [ (key = "corge", value = "789"), (key = "foo", value = "123") ]))
2170 .expectReturns(CAPNP(), ws);
2171 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2172 }).expectCanceled();
2173 
2174 KJ_ASSERT(promise.wait(ws) == kvs({{"corge", "789"}, {"foo", "123"}}));
2175 }
2176 
2177 KJ_ASSERT(expectCached(test.list("bar", "qux")) ==
2178 kvs({{"baz", "456"}, {"corge", "789"}, {"foo", "123"}}));
2179}
2180 
2181KJ_TEST("ActorCache list() consecutively reverse, present midpoint") {
2182 ActorCacheTest test;
2183 auto& ws = test.ws;
2184 auto& mockStorage = test.mockStorage;
2185 
2186 {
2187 auto promise = expectUncached(test.list("corge", "qux"));
2188 
2189 mockStorage->expectCall("list", ws)
2190 .withParams(CAPNP(start = "corge", end = "qux"), "stream"_kj)
2191 .useCallback("stream", [&](MockClient stream) {
2192 stream
2193 .call("values",
2194 CAPNP(list = [ (key = "corge", value = "789"), (key = "foo", value = "123") ]))
2195 .expectReturns(CAPNP(), ws);
2196 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2197 }).expectCanceled();
2198 
2199 KJ_ASSERT(promise.wait(ws) == kvs({{"corge", "789"}, {"foo", "123"}}));
2200 }
2201 
2202 {
2203 auto promise = expectUncached(test.list("bar", "corge"));
2204 
2205 mockStorage->expectCall("list", ws)
2206 .withParams(CAPNP(start = "bar", end = "corge"), "stream"_kj)
2207 .useCallback("stream", [&](MockClient stream) {
2208 stream.call("values", CAPNP(list = [(key = "baz", value = "456")]))
2209 .expectReturns(CAPNP(), ws);
2210 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2211 }).expectCanceled();
2212 
2213 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "456"}}));
2214 }
2215 
2216 KJ_ASSERT(expectCached(test.list("bar", "qux")) ==
2217 kvs({{"baz", "456"}, {"corge", "789"}, {"foo", "123"}}));
2218}
2219 
2220KJ_TEST("ActorCache list() starting in known-empty gap") {
2221 ActorCacheTest test;
2222 auto& ws = test.ws;
2223 auto& mockStorage = test.mockStorage;
2224 
2225 // Create a known-empty gap between "bar" and "corge".
2226 {
2227 auto promise = expectUncached(test.list("bar", "corge"));
2228 
2229 mockStorage->expectCall("list", ws)
2230 .withParams(CAPNP(start = "bar", end = "corge"), "stream"_kj)
2231 .useCallback("stream", [&](MockClient stream) {
2232 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2233 }).expectCanceled();
2234 
2235 KJ_ASSERT(promise.wait(ws) == kvs({}));
2236 }
2237 
2238 // Now list from "baz" to "qux", which starts in the gap.
2239 {
2240 auto promise = expectUncached(test.list("baz", "qux"));
2241 
2242 mockStorage->expectCall("list", ws)
2243 .withParams(CAPNP(start = "corge", end = "qux"), "stream"_kj)
2244 .useCallback("stream", [&](MockClient stream) {
2245 stream.call("values", CAPNP(list = [(key = "foo", value = "123")]))
2246 .expectReturns(CAPNP(), ws);
2247 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2248 }).expectCanceled();
2249 
2250 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}}));
2251 }
2252 
2253 KJ_ASSERT(expectCached(test.list("bar", "qux")) == kvs({{"foo", "123"}}));
2254}
2255 
2256KJ_TEST("ActorCache list() ending in known-empty gap") {
2257 ActorCacheTest test;
2258 auto& ws = test.ws;
2259 auto& mockStorage = test.mockStorage;
2260 
2261 // Create a known-empty gap between "corge" and "qux".
2262 {
2263 auto promise = expectUncached(test.list("corge", "qux"));
2264 
2265 mockStorage->expectCall("list", ws)
2266 .withParams(CAPNP(start = "corge", end = "qux"), "stream"_kj)
2267 .useCallback("stream", [&](MockClient stream) {
2268 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2269 }).expectCanceled();
2270 
2271 KJ_ASSERT(promise.wait(ws) == kvs({}));
2272 }
2273 
2274 // Now list from "bar" to "foo", which ends in the gap.
2275 {
2276 auto promise = expectUncached(test.list("bar", "foo"));
2277 
2278 // Note that the implementation of `list()` only looks for a prefix that it can skip, not a
2279 // suffix. Hence, the underlying list() call will go all the way to "foo", even though the
2280 // range from "qux" to "foo" is entirely in cache and hence in theory could be skipped. This
2281 // optimization is missing because the code is complex enough already and it doesn't seem like
2282 // it would be a win that often.
2283 mockStorage->expectCall("list", ws)
2284 .withParams(CAPNP(start = "bar", end = "foo"), "stream"_kj)
2285 .useCallback("stream", [&](MockClient stream) {
2286 stream.call("values", CAPNP(list = [(key = "baz", value = "123")]))
2287 .expectReturns(CAPNP(), ws);
2288 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2289 }).expectCanceled();
2290 
2291 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "123"}}));
2292 }
2293 
2294 KJ_ASSERT(expectCached(test.list("bar", "qux")) == kvs({{"baz", "123"}}));
2295}
2296 
2297KJ_TEST("ActorCache list() with limit and dirty puts that end up past the limit") {
2298 ActorCacheTest test;
2299 auto& ws = test.ws;
2300 auto& mockStorage = test.mockStorage;
2301 
2302 test.put("corge", "123");
2303 test.put("grault", "321");
2304 
2305 {
2306 auto promise = expectUncached(test.list("bar", "qux", 3));
2307 
2308 mockStorage->expectCall("list", ws)
2309 .withParams(CAPNP(start = "bar", end = "qux", limit = 3), "stream"_kj)
2310 .useCallback("stream", [&](MockClient stream) {
2311 stream
2312 .call("values",
2313 CAPNP(list =
2314 [
2315 (key = "bar", value = "456"), (key = "baz", value = "654"),
2316 (key = "foo", value = "789")
2317 ]))
2318 .expectReturns(CAPNP(), ws);
2319 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2320 }).expectCanceled();
2321 
2322 KJ_ASSERT(promise.wait(ws) == kvs({{"bar", "456"}, {"baz", "654"}, {"corge", "123"}}));
2323 }
2324 
2325 // Although we only requested 3 results above, we actually listed through "foo" at least, so
2326 // now we can list 4 results and they'll all come from cache.
2327 KJ_ASSERT(expectCached(test.list("bar", "qux", 4)) ==
2328 kvs({{"bar", "456"}, {"baz", "654"}, {"corge", "123"}, {"foo", "789"}}));
2329 
2330 // Acknowledge the transaction.
2331 mockStorage->expectCall("put", ws).thenReturn(CAPNP());
2332}
2333 
2334KJ_TEST("ActorCache list() overwrite endpoint") {
2335 ActorCacheTest test;
2336 auto& ws = test.ws;
2337 auto& mockStorage = test.mockStorage;
2338 
2339 {
2340 auto promise = expectUncached(test.list("corge", "qux"));
2341 
2342 mockStorage->expectCall("list", ws)
2343 .withParams(CAPNP(start = "corge", end = "qux"), "stream"_kj)
2344 .useCallback("stream", [&](MockClient stream) {
2345 stream
2346 .call("values",
2347 CAPNP(list = [ (key = "corge", value = "789"), (key = "foo", value = "123") ]))
2348 .expectReturns(CAPNP(), ws);
2349 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2350 }).expectCanceled();
2351 
2352 KJ_ASSERT(promise.wait(ws) == kvs({{"corge", "789"}, {"foo", "123"}}));
2353 }
2354 
2355 test.put("qux", "456");
2356 
2357 KJ_ASSERT(expectCached(test.list("corge", "xyzzy", 3)) ==
2358 kvs({{"corge", "789"}, {"foo", "123"}, {"qux", "456"}}));
2359 
2360 // Acknowledge the transaction.
2361 mockStorage->expectCall("put", ws).thenReturn(CAPNP());
2362}
2363 
2364KJ_TEST("ActorCache list() delete endpoint") {
2365 ActorCacheTest test;
2366 auto& ws = test.ws;
2367 auto& mockStorage = test.mockStorage;
2368 
2369 {
2370 auto promise = expectUncached(test.list("corge", "qux"));
2371 
2372 mockStorage->expectCall("list", ws)
2373 .withParams(CAPNP(start = "corge", end = "qux"), "stream"_kj)
2374 .useCallback("stream", [&](MockClient stream) {
2375 stream
2376 .call("values",
2377 CAPNP(list = [ (key = "corge", value = "789"), (key = "foo", value = "123") ]))
2378 .expectReturns(CAPNP(), ws);
2379 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2380 }).expectCanceled();
2381 
2382 KJ_ASSERT(promise.wait(ws) == kvs({{"corge", "789"}, {"foo", "123"}}));
2383 }
2384 
2385 auto deletePromise = expectUncached(test.delete_("qux"));
2386 
2387 // Acknowledge the delete transaction.
2388 {
2389 mockStorage->expectCall("delete", ws)
2390 .withParams(CAPNP(keys = ["qux"]))
2391 .thenReturn(CAPNP(numDeleted = 1));
2392 }
2393 
2394 KJ_ASSERT(deletePromise.wait(ws) == 1);
2395 
2396 // Do another list() through the deleted entry to make sure it didn't cause confusion. We apply
2397 // a limit to this list to check for a bug where negative entries in the fully-cached prefix
2398 // were incorrectly counted against the limit; only positive entries should be.
2399 {
2400 auto promise = expectUncached(test.list("corge", "xyzzy", 4));
2401 
2402 mockStorage->expectCall("list", ws)
2403 .withParams(CAPNP(start = "qux\0", end = "xyzzy", limit = 2), "stream"_kj)
2404 .useCallback("stream", [&](MockClient stream) {
2405 stream.call("values", CAPNP(list = [(key = "waldo", value = "555")]))
2406 .expectReturns(CAPNP(), ws);
2407 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2408 }).expectCanceled();
2409 
2410 KJ_ASSERT(promise.wait(ws) == kvs({{"corge", "789"}, {"foo", "123"}, {"waldo", "555"}}));
2411 }
2412}
2413 
2414KJ_TEST("ActorCache list() delete endpoint empty range") {
2415 // Same as last test except the listed range is totally empty.
2416 ActorCacheTest test;
2417 auto& ws = test.ws;
2418 auto& mockStorage = test.mockStorage;
2419 
2420 {
2421 auto promise = expectUncached(test.list("corge", "qux"));
2422 
2423 mockStorage->expectCall("list", ws)
2424 .withParams(CAPNP(start = "corge", end = "qux"), "stream"_kj)
2425 .useCallback("stream", [&](MockClient stream) {
2426 stream.call("values", CAPNP(list = [])).expectReturns(CAPNP(), ws);
2427 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2428 }).expectCanceled();
2429 
2430 KJ_ASSERT(promise.wait(ws) == kvs({}));
2431 }
2432 
2433 auto deletePromise = expectUncached(test.delete_("qux"));
2434 
2435 // Acknowledge the delete transaction.
2436 {
2437 mockStorage->expectCall("delete", ws)
2438 .withParams(CAPNP(keys = ["qux"]))
2439 .thenReturn(CAPNP(numDeleted = 1));
2440 }
2441 
2442 KJ_ASSERT(deletePromise.wait(ws) == 1);
2443 
2444 // Do another list() through the deleted entry to make sure it didn't cause confusion. We apply
2445 // a limit to this list to check for a bug where negative entries in the fully-cached prefix
2446 // were incorrectly counted against the limit; only positive entries should be.
2447 {
2448 auto promise = expectUncached(test.list("corge", "xyzzy", 4));
2449 
2450 mockStorage->expectCall("list", ws)
2451 .withParams(CAPNP(start = "qux\0", end = "xyzzy", limit = 4), "stream"_kj)
2452 .useCallback("stream", [&](MockClient stream) {
2453 stream.call("values", CAPNP(list = [(key = "qux", value = "555")]))
2454 .expectReturns(CAPNP(), ws);
2455 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2456 }).expectCanceled();
2457 
2458 KJ_ASSERT(promise.wait(ws) == kvs({}));
2459 }
2460}
2461 
2462KJ_TEST("ActorCache list() interleave streaming with other ops") {
2463 ActorCacheTest test;
2464 auto& ws = test.ws;
2465 auto& mockStorage = test.mockStorage;
2466 
2467 auto promise = expectUncached(test.list("bar", "qux"));
2468 
2469 mockStorage->expectCall("list", ws)
2470 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
2471 .useCallback("stream", [&](MockClient stream) {
2472 stream
2473 .call("values",
2474 CAPNP(list = [ (key = "bar", value = "123"), (key = "corge", value = "456") ]))
2475 .expectReturns(CAPNP(), ws);
2476 
2477 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "123");
2478 KJ_ASSERT(expectCached(test.get("baz")) == nullptr);
2479 auto promise2 = expectUncached(test.get("grault"));
2480 mockStorage->expectCall("get", ws).withParams(CAPNP(key = "grault")).thenReturn(CAPNP());
2481 KJ_ASSERT(promise2.wait(ws) == nullptr);
2482 
2483 test.put("foo", "987");
2484 
2485 stream
2486 .call("values",
2487 CAPNP(list = [ (key = "foo", value = "789"), (key = "garply", value = "555") ]))
2488 .expectReturns(CAPNP(), ws);
2489 
2490 KJ_ASSERT(expectCached(test.delete_("garply")) == 1);
2491 
2492 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2493 }).expectCanceled();
2494 
2495 KJ_ASSERT(promise.wait(ws) ==
2496 kvs({{"bar", "123"}, {"corge", "456"}, {"foo", "789"}, {"garply", "555"}}));
2497 
2498 KJ_ASSERT(expectCached(test.list("bar", "qux")) ==
2499 kvs({{"bar", "123"}, {"corge", "456"}, {"foo", "987"}}));
2500 
2501 // There will be two flushes waiting since the put of "foo" will have started before the
2502 // delete of "garply"
2503 {
2504 mockStorage->expectCall("put", ws)
2505 .withParams(CAPNP(entries = [(key = "foo", value = "987")]))
2506 .thenReturn(CAPNP());
2507 }
2508 {
2509 mockStorage->expectCall("delete", ws).withParams(CAPNP(keys = ["garply"])).thenReturn(CAPNP());
2510 }
2511}
2512 
2513KJ_TEST("ActorCache list() end of first block deleted at inopportune time") {
2514 ActorCacheTest test;
2515 auto& ws = test.ws;
2516 auto& mockStorage = test.mockStorage;
2517 
2518 // Do a delete, wait for the commit... and then hold it open.
2519 auto deletePromise = expectUncached(test.delete_("corge"));
2520 
2521 auto mockDelete = mockStorage->expectCall("delete", ws).withParams(CAPNP(keys = ["corge"]));
2522 
2523 // Now do a list.
2524 auto promise = expectUncached(test.list("bar", "qux"));
2525 
2526 mockStorage->expectCall("list", ws)
2527 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
2528 .useCallback("stream", [&](MockClient stream) {
2529 // First block ends at the deleted entry.
2530 stream
2531 .call("values",
2532 CAPNP(list = [ (key = "bar", value = "123"), (key = "corge", value = "456") ]))
2533 .expectReturns(CAPNP(), ws);
2534 
2535 // Let the delete finish. So now the last key in the first block is cached as a negative
2536 // clean entry.
2537 kj::mv(mockDelete).thenReturn(CAPNP());
2538 
2539 // Continue on.
2540 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2541 }).expectCanceled();
2542 
2543 KJ_ASSERT(promise.wait(ws) == kvs({{"bar", "123"}}));
2544 
2545 KJ_ASSERT(expectCached(test.list("bar", "qux")) == kvs({{"bar", "123"}}));
2546}
2547 
2548KJ_TEST("ActorCache list() retry on failure") {
2549 ActorCacheTest test;
2550 auto& ws = test.ws;
2551 auto& mockStorage = test.mockStorage;
2552 
2553 {
2554 auto promise = expectUncached(test.list("bar", "qux", 4));
2555 
2556 mockStorage->expectCall("list", ws)
2557 .withParams(CAPNP(start = "bar", end = "qux", limit = 4), "stream"_kj)
2558 .useCallback("stream", [&](MockClient stream) {
2559 stream
2560 .call("values",
2561 CAPNP(list = [ (key = "bar", value = "456"), (key = "baz", value = "789") ]))
2562 .expectReturns(CAPNP(), ws);
2563 }).thenThrow(KJ_EXCEPTION(DISCONNECTED, "oops"));
2564 
2565 // Retry starts from `baz`.
2566 mockStorage->expectCall("list", ws)
2567 .withParams(CAPNP(start = "baz\0", end = "qux", limit = 2), "stream"_kj)
2568 .useCallback("stream", [&](MockClient stream) {
2569 // Duplicates of earlier keys will be ignored.
2570 stream
2571 .call("values",
2572 CAPNP(list =
2573 [
2574 (key = "bar", value = "IGNORE"), (key = "baz", value = "IGNORE"),
2575 (key = "foo", value = "123")
2576 ]))
2577 .expectReturns(CAPNP(), ws);
2578 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2579 }).expectCanceled();
2580 
2581 KJ_ASSERT(promise.wait(ws) == kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
2582 }
2583 
2584 KJ_ASSERT(expectCached(test.list("bar", "qux")) ==
2585 kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
2586}
2587 
2588KJ_TEST("ActorCache get() of endpoint of previous list() returning negative is cached correctly") {
2589 // This tests for a bug that once existed in ActorCache::addReadResultToCache() where we compared
2590 // against a moved-away value.
2591 ActorCacheTest test;
2592 auto& ws = test.ws;
2593 auto& mockStorage = test.mockStorage;
2594 
2595 {
2596 auto promise = expectUncached(test.list("bar", "qux", 4));
2597 
2598 mockStorage->expectCall("list", ws)
2599 .withParams(CAPNP(start = "bar", end = "qux", limit = 4), "stream"_kj)
2600 .useCallback("stream", [&](MockClient stream) {
2601 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2602 }).expectCanceled();
2603 
2604 KJ_ASSERT(promise.wait(ws) == kvs({}));
2605 }
2606 
2607 {
2608 auto promise = expectUncached(test.get("qux"));
2609 mockStorage->expectCall("get", ws).withParams(CAPNP(key = "qux")).thenReturn(CAPNP());
2610 KJ_ASSERT(promise.wait(ws) == nullptr);
2611 }
2612 
2613 KJ_ASSERT(expectCached(test.get("qux")) == nullptr);
2614}
2615 
2616// =======================================================================================
2617// And now... listReverse()... needs all its own tests...
2618 
2619KJ_TEST("ActorCache listReverse()") {
2620 ActorCacheTest test;
2621 auto& ws = test.ws;
2622 auto& mockStorage = test.mockStorage;
2623 
2624 {
2625 auto promise = expectUncached(test.listReverse("bar", "qux"));
2626 
2627 mockStorage->expectCall("list", ws)
2628 .withParams(CAPNP(start = "bar", end = "qux", reverse = true), "stream"_kj)
2629 .useCallback("stream", [&](MockClient stream) {
2630 stream
2631 .call("values",
2632 CAPNP(list =
2633 [
2634 (key = "foo", value = "123"), (key = "baz", value = "789"),
2635 (key = "bar", value = "456")
2636 ]))
2637 .expectReturns(CAPNP(), ws);
2638 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2639 }).expectCanceled();
2640 
2641 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}}));
2642 }
2643 
2644 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
2645 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "456");
2646 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "789");
2647 
2648 // Stuff in range that wasn't reported is cached as absent.
2649 KJ_ASSERT(expectCached(test.get("bara")) == nullptr);
2650 KJ_ASSERT(expectCached(test.get("corge")) == nullptr);
2651 KJ_ASSERT(expectCached(test.get("quw")) == nullptr);
2652 
2653 // Listing the same range again is fully cached.
2654 KJ_ASSERT(expectCached(test.listReverse("bar", "qux")) ==
2655 kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}}));
2656 
2657 // Limits can be applied to the cached results.
2658 KJ_ASSERT(expectCached(test.listReverse("bar", "qux", 0u)) == kvs({}));
2659 KJ_ASSERT(expectCached(test.listReverse("bar", "qux", 1)) == kvs({{"foo", "123"}}));
2660 KJ_ASSERT(
2661 expectCached(test.listReverse("bar", "qux", 2)) == kvs({{"foo", "123"}, {"baz", "789"}}));
2662 KJ_ASSERT(expectCached(test.listReverse("bar", "qux", 3)) ==
2663 kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}}));
2664 KJ_ASSERT(expectCached(test.listReverse("bar", "qux", 4)) ==
2665 kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}}));
2666 KJ_ASSERT(expectCached(test.listReverse("bar", "qux", 1000)) ==
2667 kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}}));
2668 
2669 // The endpoint of the list is not cached.
2670 {
2671 auto promise = expectUncached(test.get("qux"));
2672 
2673 mockStorage->expectCall("get", ws)
2674 .withParams(CAPNP(key = "qux"))
2675 .thenReturn(CAPNP(value = "555"));
2676 
2677 auto result = KJ_ASSERT_NONNULL(promise.wait(ws));
2678 KJ_EXPECT(result == "555");
2679 }
2680}
2681 
2682KJ_TEST("ActorCache listReverse() all") {
2683 ActorCacheTest test;
2684 auto& ws = test.ws;
2685 auto& mockStorage = test.mockStorage;
2686 
2687 {
2688 auto promise = expectUncached(test.listReverse(nullptr, nullptr));
2689 
2690 mockStorage->expectCall("list", ws)
2691 .withParams(CAPNP(reverse = true), "stream"_kj)
2692 .useCallback("stream", [&](MockClient stream) {
2693 stream
2694 .call("values",
2695 CAPNP(list =
2696 [
2697 (key = "foo", value = "123"), (key = "baz", value = "789"),
2698 (key = "bar", value = "456")
2699 ]))
2700 .expectReturns(CAPNP(), ws);
2701 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2702 }).expectCanceled();
2703 
2704 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}}));
2705 }
2706 
2707 KJ_ASSERT(expectCached(test.get("")) == nullptr);
2708 KJ_ASSERT(expectCached(test.list(nullptr, nullptr)) ==
2709 kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
2710 KJ_ASSERT(expectCached(test.listReverse("bar", "qux")) ==
2711 kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}}));
2712 KJ_ASSERT(expectCached(test.listReverse(nullptr, nullptr)) ==
2713 kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}}));
2714 KJ_ASSERT(
2715 expectCached(test.listReverse("baz", nullptr)) == kvs({{"foo", "123"}, {"baz", "789"}}));
2716 KJ_ASSERT(
2717 expectCached(test.listReverse(nullptr, "foo")) == kvs({{"baz", "789"}, {"bar", "456"}}));
2718}
2719 
2720KJ_TEST("ActorCache listReverse() with limit") {
2721 ActorCacheTest test;
2722 auto& ws = test.ws;
2723 auto& mockStorage = test.mockStorage;
2724 
2725 {
2726 auto promise = expectUncached(test.listReverse("abc", "qux", 3));
2727 
2728 mockStorage->expectCall("list", ws)
2729 .withParams(CAPNP(start = "abc", end = "qux", limit = 3, reverse = true), "stream"_kj)
2730 .useCallback("stream", [&](MockClient stream) {
2731 stream
2732 .call("values",
2733 CAPNP(list =
2734 [
2735 (key = "foo", value = "123"), (key = "baz", value = "789"),
2736 (key = "bar", value = "456")
2737 ]))
2738 .expectReturns(CAPNP(), ws);
2739 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2740 }).expectCanceled();
2741 
2742 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}}));
2743 }
2744 
2745 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
2746 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "456");
2747 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "789");
2748 
2749 // Stuff in range that wasn't reported is cached as absent -- but not past the last reported
2750 // value, which was "foo".
2751 KJ_ASSERT(expectCached(test.get("bara")) == nullptr);
2752 KJ_ASSERT(expectCached(test.get("corge")) == nullptr);
2753 KJ_ASSERT(expectCached(test.get("fon")) == nullptr);
2754 
2755 // Stuff before the first key is not in cache.
2756 (void)expectUncached(test.get("baq"));
2757 
2758 // Listing the same range again, with the same limit or lower, is fully cached.
2759 KJ_ASSERT(expectCached(test.listReverse("bar", "qux", 3)) ==
2760 kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}}));
2761 KJ_ASSERT(
2762 expectCached(test.listReverse("bar", "qux", 2)) == kvs({{"foo", "123"}, {"baz", "789"}}));
2763 KJ_ASSERT(expectCached(test.listReverse("bar", "qux", 1)) == kvs({{"foo", "123"}}));
2764 KJ_ASSERT(expectCached(test.listReverse("bar", "qux", 0u)) == kvs({}));
2765 
2766 // But a larger limit won't be cached.
2767 {
2768 auto promise = expectUncached(test.listReverse("abc", "qux", 4));
2769 
2770 // The new list will end at "bar" with a limit of 1, so that it will only get the one remaining
2771 // key that it needs.
2772 mockStorage->expectCall("list", ws)
2773 .withParams(CAPNP(start = "abc", end = "bar", limit = 1, reverse = true), "stream"_kj)
2774 .useCallback("stream", [&](MockClient stream) {
2775 stream.call("values", CAPNP(list = [(key = "baa", value = "xyz")]))
2776 .expectReturns(CAPNP(), ws);
2777 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2778 }).expectCanceled();
2779 
2780 KJ_ASSERT(
2781 promise.wait(ws) == kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}, {"baa", "xyz"}}));
2782 }
2783 
2784 // Cached if we try it again though.
2785 KJ_ASSERT(expectCached(test.listReverse("abc", "qux", 4)) ==
2786 kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}, {"baa", "xyz"}}));
2787}
2788 
2789KJ_TEST("ActorCache listReverse() with limit around negative entries") {
2790 // This checks for a bug where the initial scan through cache for list() applies the limit to
2791 // the total number of entries seen (positive or negative), when it really needs to apply only
2792 // to positive entries.
2793 
2794 ActorCacheTest test;
2795 auto& ws = test.ws;
2796 auto& mockStorage = test.mockStorage;
2797 
2798 // Set up a bunch of negative entries and a positive one after them.
2799 test.delete_({"bar1"_kj, "bar2"_kj, "bar3"_kj, "bar4"_kj});
2800 test.put("bar", "456");
2801 
2802 // Now do a list through them. It should see the positive entry in cache.
2803 {
2804 auto promise = expectUncached(test.listReverse("bar", "qux", 3));
2805 
2806 mockStorage->expectCall("list", ws)
2807 .withParams(CAPNP(start = "bar", end = "qux", limit = 7, reverse = true), "stream"_kj)
2808 .useCallback("stream", [&](MockClient stream) {
2809 stream
2810 .call("values",
2811 CAPNP(list =
2812 [
2813 (key = "foo", value = "123"), (key = "baz", value = "789"),
2814 (key = "bar3", value = "yyy"), (key = "bar1", value = "xxx")
2815 ]))
2816 .expectReturns(CAPNP(), ws);
2817 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2818 }).expectCanceled();
2819 
2820 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}}));
2821 }
2822 
2823 KJ_ASSERT(expectCached(test.listReverse("bar", "qux", 3)) ==
2824 kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}}));
2825 
2826 // Acknowledge the transaction.
2827 {
2828 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
2829 mockTxn->expectCall("delete", ws).thenReturn(CAPNP());
2830 mockTxn->expectCall("put", ws).thenReturn(CAPNP());
2831 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
2832 mockTxn->expectDropped(ws);
2833 }
2834}
2835 
2836KJ_TEST("ActorCache listReverse() start point is not present") {
2837 ActorCacheTest test;
2838 auto& ws = test.ws;
2839 auto& mockStorage = test.mockStorage;
2840 
2841 {
2842 auto promise = expectUncached(test.listReverse("bar", "qux"));
2843 
2844 mockStorage->expectCall("list", ws)
2845 .withParams(CAPNP(start = "bar", end = "qux", reverse = true), "stream"_kj)
2846 .useCallback("stream", [&](MockClient stream) {
2847 stream
2848 .call("values",
2849 CAPNP(list = [ (key = "foo", value = "123"), (key = "baz", value = "789") ]))
2850 .expectReturns(CAPNP(), ws);
2851 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2852 }).expectCanceled();
2853 
2854 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}, {"baz", "789"}}));
2855 }
2856 
2857 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
2858 
2859 KJ_ASSERT(expectCached(test.get("bar")) == nullptr);
2860 KJ_ASSERT(expectCached(test.get("bara")) == nullptr);
2861 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "789");
2862 KJ_ASSERT(expectCached(test.get("baza")) == nullptr);
2863 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
2864 KJ_ASSERT(expectCached(test.get("fooa")) == nullptr);
2865}
2866 
2867KJ_TEST("ActorCache listReverse() multiple ranges") {
2868 ActorCacheTest test;
2869 auto& ws = test.ws;
2870 auto& mockStorage = test.mockStorage;
2871 
2872 {
2873 auto promise = expectUncached(test.listReverse("a", "c"));
2874 
2875 mockStorage->expectCall("list", ws)
2876 .withParams(CAPNP(start = "a", end = "c", reverse = true), "stream"_kj)
2877 .useCallback("stream", [&](MockClient stream) {
2878 stream.call("values", CAPNP(list = [ (key = "b", value = "2"), (key = "a", value = "1") ]))
2879 .expectReturns(CAPNP(), ws);
2880 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2881 }).expectCanceled();
2882 
2883 KJ_ASSERT(promise.wait(ws) == kvs({{"b", "2"}, {"a", "1"}}));
2884 }
2885 
2886 KJ_ASSERT(expectCached(test.listReverse("a", "c")) == kvs({{"b", "2"}, {"a", "1"}}));
2887 
2888 {
2889 auto promise = expectUncached(test.listReverse("x", "z"));
2890 
2891 mockStorage->expectCall("list", ws)
2892 .withParams(CAPNP(start = "x", end = "z", reverse = true), "stream"_kj)
2893 .useCallback("stream", [&](MockClient stream) {
2894 stream.call("values", CAPNP(list = [(key = "y", value = "9")])).expectReturns(CAPNP(), ws);
2895 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2896 }).expectCanceled();
2897 
2898 KJ_ASSERT(promise.wait(ws) == kvs({{"y", "9"}}));
2899 }
2900 
2901 KJ_ASSERT(expectCached(test.listReverse("a", "c")) == kvs({{"b", "2"}, {"a", "1"}}));
2902 KJ_ASSERT(expectCached(test.listReverse("x", "z")) == kvs({{"y", "9"}}));
2903 
2904 (void)expectUncached(test.get("w"));
2905 (void)expectUncached(test.get("d"));
2906 (void)expectUncached(test.get("c"));
2907}
2908 
2909KJ_TEST("ActorCache listReverse() with some already-cached keys in range") {
2910 ActorCacheTest test;
2911 auto& ws = test.ws;
2912 auto& mockStorage = test.mockStorage;
2913 
2914 // Initialize cache with some clean entries, both positive and negative.
2915 {
2916 auto promise1 = expectUncached(test.get("bbb"));
2917 auto promise2 = expectUncached(test.get("ccc"));
2918 
2919 mockStorage->expectCall("get", ws).withParams(CAPNP(key = "bbb")).thenReturn(CAPNP());
2920 mockStorage->expectCall("get", ws)
2921 .withParams(CAPNP(key = "ccc"))
2922 .thenReturn(CAPNP(value = "cval"));
2923 
2924 KJ_ASSERT(promise1.wait(ws) == nullptr);
2925 KJ_ASSERT(KJ_ASSERT_NONNULL(promise2.wait(ws)) == "cval");
2926 }
2927 
2928 // Also some newly-written entries, positive and negative.
2929 test.put("ddd", "dval");
2930 auto deletePromise = expectUncached(test.delete_("eee"));
2931 
2932 // Now list the range. Explicitly produce results that contradict the recent writes.
2933 {
2934 auto promise = expectUncached(test.listReverse("aaa", "fff"));
2935 
2936 mockStorage->expectCall("list", ws)
2937 .withParams(CAPNP(start = "aaa", end = "fff", reverse = true), "stream"_kj)
2938 .useCallback("stream", [&](MockClient stream) {
2939 stream
2940 .call("values",
2941 CAPNP(list = [ (key = "eee", value = "eval"), (key = "ccc", value = "cval") ]))
2942 .expectReturns(CAPNP(), ws);
2943 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2944 }).expectCanceled();
2945 
2946 KJ_ASSERT(promise.wait(ws) == kvs({{"ddd", "dval"}, {"ccc", "cval"}}));
2947 }
2948 
2949 {
2950 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
2951 mockTxn->expectCall("delete", ws)
2952 .withParams(CAPNP(keys = ["eee"]))
2953 .thenReturn(CAPNP(numDeleted = 1));
2954 mockTxn->expectCall("put", ws)
2955 .withParams(CAPNP(entries = [(key = "ddd", value = "dval")]))
2956 .thenReturn(CAPNP());
2957 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
2958 mockTxn->expectDropped(ws);
2959 }
2960 
2961 KJ_ASSERT(deletePromise.wait(ws) == 1);
2962}
2963 
2964KJ_TEST("ActorCache listReverse() with seemingly-redundant dirty entries") {
2965 ActorCacheTest test;
2966 auto& ws = test.ws;
2967 auto& mockStorage = test.mockStorage;
2968 
2969 // Write some stuff.
2970 auto deletePromise = expectUncached(test.delete_("bbb"));
2971 test.put("ccc", "cval");
2972 
2973 // Initiate a list operation, but don't complete it yet.
2974 auto listPromise = expectUncached(test.listReverse("aaa", "fff"));
2975 auto listCall = mockStorage->expectCall("list", ws)
2976 .withParams(CAPNP(start = "aaa", end = "fff", reverse = true), "stream"_kj);
2977 
2978 // We won't do any work while the list is outstanding.
2979 mockStorage->expectNoActivity(ws);
2980 
2981 // Now write some contradictory values.
2982 test.put("bbb", "bval");
2983 KJ_ASSERT(expectCached(test.delete_("ccc")) == 1);
2984 
2985 // Now let the list complete in a way that matches what was just written.
2986 kj::mv(listCall)
2987 .useCallback("stream", [&](MockClient stream) {
2988 stream.call("values", CAPNP(list = [(key = "bbb", value = "bval")])).expectReturns(CAPNP(), ws);
2989 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
2990 }).expectCanceled();
2991 
2992 // The list produces results consistent with when it started.
2993 KJ_ASSERT(listPromise.wait(ws) == kvs({{"ccc", "cval"}}));
2994 
2995 // But the later writes are still there in cache.
2996 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bbb"))) == "bval");
2997 KJ_ASSERT(expectCached(test.get("ccc")) == nullptr);
2998 
2999 // The transaction completes now.
3000 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
3001 mockTxn->expectCall("delete", ws)
3002 .withParams(CAPNP(keys = ["bbb"]))
3003 .thenReturn(CAPNP(numDeleted = 1));
3004 mockTxn->expectCall("put", ws)
3005 .withParams(CAPNP(entries = [(key = "ccc", value = "cval")]))
3006 .thenReturn(CAPNP());
3007 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
3008 mockTxn->expectDropped(ws);
3009 KJ_ASSERT(deletePromise.wait(ws) == 1);
3010 
3011 // And then there's a new transaction to write things back to the original values.
3012 // This is NOT REDUNDANT, even though the list results seemed to match the current cached values!
3013 // (I wrote this test to prove to myself that a DIRTY entry can't be marked CLEAN just because a
3014 // read result from disk came back with the same value.)
3015 mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
3016 mockTxn->expectCall("delete", ws)
3017 .withParams(CAPNP(keys = ["ccc"]))
3018 .thenReturn(CAPNP(numDeleted = 1));
3019 mockTxn->expectCall("put", ws)
3020 .withParams(CAPNP(entries = [(key = "bbb", value = "bval")]))
3021 .thenReturn(CAPNP());
3022 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
3023 mockTxn->expectDropped(ws);
3024 
3025 // For good measure, verify list result can be served from cache.
3026 KJ_ASSERT(expectCached(test.listReverse("aaa", "fff")) == kvs({{"bbb", "bval"}}));
3027}
3028 
3029KJ_TEST("ActorCache listReverse() starting from known value") {
3030 ActorCacheTest test;
3031 auto& ws = test.ws;
3032 auto& mockStorage = test.mockStorage;
3033 
3034 {
3035 test.put("bar", "123");
3036 
3037 mockStorage->expectCall("put", ws).thenReturn(CAPNP());
3038 }
3039 
3040 {
3041 auto promise = expectUncached(test.listReverse("bar", "qux"));
3042 
3043 mockStorage->expectCall("list", ws)
3044 .withParams(CAPNP(start = "bar", end = "qux", reverse = true), "stream"_kj)
3045 .useCallback("stream", [&](MockClient stream) {
3046 stream.call("values", CAPNP(list = [(key = "baz", value = "456")]))
3047 .expectReturns(CAPNP(), ws);
3048 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3049 }).expectCanceled();
3050 
3051 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "456"}, {"bar", "123"}}));
3052 }
3053 
3054 KJ_ASSERT(expectCached(test.listReverse("bar", "qux")) == kvs({{"baz", "456"}, {"bar", "123"}}));
3055}
3056 
3057KJ_TEST("ActorCache listReverse() starting from unknown value") {
3058 ActorCacheTest test;
3059 auto& ws = test.ws;
3060 auto& mockStorage = test.mockStorage;
3061 
3062 {
3063 test.put("baz", "456");
3064 
3065 mockStorage->expectCall("put", ws).thenReturn(CAPNP());
3066 }
3067 
3068 {
3069 auto promise = expectUncached(test.listReverse("bar", "qux"));
3070 
3071 mockStorage->expectCall("list", ws)
3072 .withParams(CAPNP(start = "bar", end = "qux", reverse = true), "stream"_kj)
3073 .useCallback("stream", [&](MockClient stream) {
3074 stream
3075 .call("values",
3076 CAPNP(list = [ (key = "foo", value = "123"), (key = "baz", value = "456") ]))
3077 .expectReturns(CAPNP(), ws);
3078 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3079 }).expectCanceled();
3080 
3081 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}, {"baz", "456"}}));
3082 }
3083 
3084 KJ_ASSERT(expectCached(test.listReverse("bar", "qux")) == kvs({{"foo", "123"}, {"baz", "456"}}));
3085}
3086 
3087KJ_TEST("ActorCache listReverse() consecutively, absent midpoint") {
3088 ActorCacheTest test;
3089 auto& ws = test.ws;
3090 auto& mockStorage = test.mockStorage;
3091 
3092 {
3093 auto promise = expectUncached(test.listReverse("bar", "corge"));
3094 
3095 mockStorage->expectCall("list", ws)
3096 .withParams(CAPNP(start = "bar", end = "corge", reverse = true), "stream"_kj)
3097 .useCallback("stream", [&](MockClient stream) {
3098 stream.call("values", CAPNP(list = [(key = "baz", value = "456")]))
3099 .expectReturns(CAPNP(), ws);
3100 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3101 }).expectCanceled();
3102 
3103 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "456"}}));
3104 }
3105 
3106 {
3107 auto promise = expectUncached(test.listReverse("corge", "qux"));
3108 
3109 mockStorage->expectCall("list", ws)
3110 .withParams(CAPNP(start = "corge", end = "qux", reverse = true), "stream"_kj)
3111 .useCallback("stream", [&](MockClient stream) {
3112 stream.call("values", CAPNP(list = [(key = "foo", value = "123")]))
3113 .expectReturns(CAPNP(), ws);
3114 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3115 }).expectCanceled();
3116 
3117 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}}));
3118 }
3119 
3120 KJ_ASSERT(expectCached(test.listReverse("bar", "qux")) == kvs({{"foo", "123"}, {"baz", "456"}}));
3121}
3122 
3123KJ_TEST("ActorCache listReverse() consecutively reverse, absent midpoint") {
3124 ActorCacheTest test;
3125 auto& ws = test.ws;
3126 auto& mockStorage = test.mockStorage;
3127 
3128 {
3129 auto promise = expectUncached(test.listReverse("corge", "qux"));
3130 
3131 mockStorage->expectCall("list", ws)
3132 .withParams(CAPNP(start = "corge", end = "qux", reverse = true), "stream"_kj)
3133 .useCallback("stream", [&](MockClient stream) {
3134 stream.call("values", CAPNP(list = [(key = "foo", value = "123")]))
3135 .expectReturns(CAPNP(), ws);
3136 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3137 }).expectCanceled();
3138 
3139 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}}));
3140 }
3141 
3142 {
3143 auto promise = expectUncached(test.listReverse("bar", "corge"));
3144 
3145 mockStorage->expectCall("list", ws)
3146 .withParams(CAPNP(start = "bar", end = "corge", reverse = true), "stream"_kj)
3147 .useCallback("stream", [&](MockClient stream) {
3148 stream.call("values", CAPNP(list = [(key = "baz", value = "456")]))
3149 .expectReturns(CAPNP(), ws);
3150 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3151 }).expectCanceled();
3152 
3153 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "456"}}));
3154 }
3155 
3156 KJ_ASSERT(expectCached(test.listReverse("bar", "qux")) == kvs({{"foo", "123"}, {"baz", "456"}}));
3157}
3158 
3159KJ_TEST("ActorCache listReverse() consecutively, present midpoint") {
3160 ActorCacheTest test;
3161 auto& ws = test.ws;
3162 auto& mockStorage = test.mockStorage;
3163 
3164 {
3165 auto promise = expectUncached(test.listReverse("bar", "corge"));
3166 
3167 mockStorage->expectCall("list", ws)
3168 .withParams(CAPNP(start = "bar", end = "corge", reverse = true), "stream"_kj)
3169 .useCallback("stream", [&](MockClient stream) {
3170 stream.call("values", CAPNP(list = [(key = "baz", value = "456")]))
3171 .expectReturns(CAPNP(), ws);
3172 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3173 }).expectCanceled();
3174 
3175 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "456"}}));
3176 }
3177 
3178 {
3179 auto promise = expectUncached(test.listReverse("corge", "qux"));
3180 
3181 mockStorage->expectCall("list", ws)
3182 .withParams(CAPNP(start = "corge", end = "qux", reverse = true), "stream"_kj)
3183 .useCallback("stream", [&](MockClient stream) {
3184 stream
3185 .call("values",
3186 CAPNP(list = [ (key = "foo", value = "123"), (key = "corge", value = "789") ]))
3187 .expectReturns(CAPNP(), ws);
3188 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3189 }).expectCanceled();
3190 
3191 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}, {"corge", "789"}}));
3192 }
3193 
3194 KJ_ASSERT(expectCached(test.listReverse("bar", "qux")) ==
3195 kvs({{"foo", "123"}, {"corge", "789"}, {"baz", "456"}}));
3196}
3197 
3198KJ_TEST("ActorCache listReverse() consecutively reverse, present midpoint") {
3199 ActorCacheTest test;
3200 auto& ws = test.ws;
3201 auto& mockStorage = test.mockStorage;
3202 
3203 {
3204 auto promise = expectUncached(test.listReverse("corge", "qux"));
3205 
3206 mockStorage->expectCall("list", ws)
3207 .withParams(CAPNP(start = "corge", end = "qux", reverse = true), "stream"_kj)
3208 .useCallback("stream", [&](MockClient stream) {
3209 stream
3210 .call("values",
3211 CAPNP(list = [ (key = "foo", value = "123"), (key = "corge", value = "789") ]))
3212 .expectReturns(CAPNP(), ws);
3213 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3214 }).expectCanceled();
3215 
3216 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}, {"corge", "789"}}));
3217 }
3218 
3219 {
3220 auto promise = expectUncached(test.listReverse("bar", "corge"));
3221 
3222 mockStorage->expectCall("list", ws)
3223 .withParams(CAPNP(start = "bar", end = "corge", reverse = true), "stream"_kj)
3224 .useCallback("stream", [&](MockClient stream) {
3225 stream.call("values", CAPNP(list = [(key = "baz", value = "456")]))
3226 .expectReturns(CAPNP(), ws);
3227 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3228 }).expectCanceled();
3229 
3230 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "456"}}));
3231 }
3232 
3233 KJ_ASSERT(expectCached(test.listReverse("bar", "qux")) ==
3234 kvs({{"foo", "123"}, {"corge", "789"}, {"baz", "456"}}));
3235}
3236 
3237KJ_TEST("ActorCache listReverse() starting in known-empty gap") {
3238 ActorCacheTest test;
3239 auto& ws = test.ws;
3240 auto& mockStorage = test.mockStorage;
3241 
3242 // Create a known-empty gap between "bar" and "corge".
3243 {
3244 auto promise = expectUncached(test.list("bar", "corge"));
3245 
3246 mockStorage->expectCall("list", ws)
3247 .withParams(CAPNP(start = "bar", end = "corge"), "stream"_kj)
3248 .useCallback("stream", [&](MockClient stream) {
3249 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3250 }).expectCanceled();
3251 
3252 KJ_ASSERT(promise.wait(ws) == kvs({}));
3253 }
3254 
3255 // Now list from "baz" to "qux", which starts in the gap.
3256 {
3257 auto promise = expectUncached(test.listReverse("baz", "qux"));
3258 
3259 mockStorage->expectCall("list", ws)
3260 .withParams(CAPNP(start = "baz", end = "qux", reverse = true), "stream"_kj)
3261 .useCallback("stream", [&](MockClient stream) {
3262 stream.call("values", CAPNP(list = [(key = "foo", value = "123")]))
3263 .expectReturns(CAPNP(), ws);
3264 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3265 }).expectCanceled();
3266 
3267 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}}));
3268 }
3269 
3270 KJ_ASSERT(expectCached(test.list("bar", "qux")) == kvs({{"foo", "123"}}));
3271}
3272 
3273KJ_TEST("ActorCache listReverse() ending in known-empty gap") {
3274 ActorCacheTest test;
3275 auto& ws = test.ws;
3276 auto& mockStorage = test.mockStorage;
3277 
3278 // Create a known-empty gap between "corge" and "qux".
3279 {
3280 auto promise = expectUncached(test.list("corge", "qux"));
3281 
3282 mockStorage->expectCall("list", ws)
3283 .withParams(CAPNP(start = "corge", end = "qux"), "stream"_kj)
3284 .useCallback("stream", [&](MockClient stream) {
3285 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3286 }).expectCanceled();
3287 
3288 KJ_ASSERT(promise.wait(ws) == kvs({}));
3289 }
3290 
3291 // Now list from "bar" to "foo", which ends in the gap.
3292 {
3293 auto promise = expectUncached(test.listReverse("bar", "foo"));
3294 
3295 mockStorage->expectCall("list", ws)
3296 .withParams(CAPNP(start = "bar", end = "corge", reverse = true), "stream"_kj)
3297 .useCallback("stream", [&](MockClient stream) {
3298 stream.call("values", CAPNP(list = [(key = "baz", value = "123")]))
3299 .expectReturns(CAPNP(), ws);
3300 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3301 }).expectCanceled();
3302 
3303 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "123"}}));
3304 }
3305 
3306 KJ_ASSERT(expectCached(test.listReverse("bar", "qux")) == kvs({{"baz", "123"}}));
3307}
3308 
3309KJ_TEST("ActorCache listReverse() with limit and dirty puts that end up past the limit") {
3310 ActorCacheTest test;
3311 auto& ws = test.ws;
3312 auto& mockStorage = test.mockStorage;
3313 
3314 test.put("corge", "123");
3315 test.put("bar", "321");
3316 
3317 {
3318 auto promise = expectUncached(test.listReverse("bar", "qux", 3));
3319 
3320 mockStorage->expectCall("list", ws)
3321 .withParams(CAPNP(start = "bar", end = "qux", limit = 3, reverse = true), "stream"_kj)
3322 .useCallback("stream", [&](MockClient stream) {
3323 stream
3324 .call("values",
3325 CAPNP(list =
3326 [
3327 (key = "grault", value = "456"), (key = "foo", value = "654"),
3328 (key = "baz", value = "789")
3329 ]))
3330 .expectReturns(CAPNP(), ws);
3331 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3332 }).expectCanceled();
3333 
3334 KJ_ASSERT(promise.wait(ws) == kvs({{"grault", "456"}, {"foo", "654"}, {"corge", "123"}}));
3335 }
3336 
3337 // Although we only requested 3 results above, we actually listed through "baz" at least, so
3338 // now we can list 4 results and they'll all come from cache.
3339 KJ_ASSERT(expectCached(test.listReverse("bar", "qux", 4)) ==
3340 kvs({{"grault", "456"}, {"foo", "654"}, {"corge", "123"}, {"baz", "789"}}));
3341 
3342 // Acknowledge the transaction.
3343 mockStorage->expectCall("put", ws).thenReturn(CAPNP());
3344}
3345 
3346KJ_TEST("ActorCache listReverse() overwrite endpoint") {
3347 ActorCacheTest test;
3348 auto& ws = test.ws;
3349 auto& mockStorage = test.mockStorage;
3350 
3351 {
3352 auto promise = expectUncached(test.listReverse("corge", "qux"));
3353 
3354 mockStorage->expectCall("list", ws)
3355 .withParams(CAPNP(start = "corge", end = "qux", reverse = true), "stream"_kj)
3356 .useCallback("stream", [&](MockClient stream) {
3357 stream
3358 .call("values",
3359 CAPNP(list = [ (key = "foo", value = "123"), (key = "corge", value = "789") ]))
3360 .expectReturns(CAPNP(), ws);
3361 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3362 }).expectCanceled();
3363 
3364 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}, {"corge", "789"}}));
3365 }
3366 
3367 test.put("qux", "456");
3368 
3369 KJ_ASSERT(expectCached(test.list("corge", "xyzzy", 3)) ==
3370 kvs({{"corge", "789"}, {"foo", "123"}, {"qux", "456"}}));
3371 
3372 // Acknowledge the transaction.
3373 mockStorage->expectCall("put", ws).thenReturn(CAPNP());
3374}
3375 
3376KJ_TEST("ActorCache listReverse() delete endpoint") {
3377 ActorCacheTest test;
3378 auto& ws = test.ws;
3379 auto& mockStorage = test.mockStorage;
3380 
3381 {
3382 auto promise = expectUncached(test.listReverse("corge", "qux"));
3383 
3384 mockStorage->expectCall("list", ws)
3385 .withParams(CAPNP(start = "corge", end = "qux", reverse = true), "stream"_kj)
3386 .useCallback("stream", [&](MockClient stream) {
3387 stream
3388 .call("values",
3389 CAPNP(list = [ (key = "foo", value = "123"), (key = "corge", value = "789") ]))
3390 .expectReturns(CAPNP(), ws);
3391 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3392 }).expectCanceled();
3393 
3394 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}, {"corge", "789"}}));
3395 }
3396 
3397 KJ_ASSERT(expectCached(test.delete_("corge")) == 1);
3398 
3399 // Acknowledge the delete transaction.
3400 {
3401 mockStorage->expectCall("delete", ws)
3402 .withParams(CAPNP(keys = ["corge"]))
3403 .thenReturn(CAPNP(numDeleted = 1));
3404 }
3405 
3406 {
3407 auto promise = expectUncached(test.listReverse("bar", "qux", 4));
3408 
3409 mockStorage->expectCall("list", ws)
3410 .withParams(CAPNP(start = "bar", end = "corge", limit = 3, reverse = true), "stream"_kj)
3411 .useCallback("stream", [&](MockClient stream) {
3412 stream.call("values", CAPNP(list = [(key = "baz", value = "555")]))
3413 .expectReturns(CAPNP(), ws);
3414 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3415 }).expectCanceled();
3416 
3417 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}, {"baz", "555"}}));
3418 }
3419}
3420 
3421KJ_TEST("ActorCache listReverse() interleave streaming with other ops") {
3422 ActorCacheTest test;
3423 auto& ws = test.ws;
3424 auto& mockStorage = test.mockStorage;
3425 
3426 auto promise = expectUncached(test.listReverse("baa", "qux"));
3427 
3428 mockStorage->expectCall("list", ws)
3429 .withParams(CAPNP(start = "baa", end = "qux", reverse = true), "stream"_kj)
3430 .useCallback("stream", [&](MockClient stream) {
3431 stream
3432 .call("values",
3433 CAPNP(list = [ (key = "garply", value = "555"), (key = "foo", value = "789") ]))
3434 .expectReturns(CAPNP(), ws);
3435 
3436 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("garply"))) == "555");
3437 KJ_ASSERT(expectCached(test.get("grault")) == nullptr);
3438 KJ_ASSERT(expectCached(test.get("gah")) == nullptr);
3439 auto promise2 = expectUncached(test.get("baz"));
3440 mockStorage->expectCall("get", ws).withParams(CAPNP(key = "baz")).thenReturn(CAPNP());
3441 KJ_ASSERT(promise2.wait(ws) == nullptr);
3442 
3443 test.put("corge", "987");
3444 
3445 stream
3446 .call("values",
3447 CAPNP(list = [ (key = "corge", value = "456"), (key = "bar", value = "123") ]))
3448 .expectReturns(CAPNP(), ws);
3449 
3450 KJ_ASSERT(expectCached(test.delete_("bar")) == 1);
3451 
3452 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3453 }).expectCanceled();
3454 
3455 KJ_ASSERT(promise.wait(ws) ==
3456 kvs({{"garply", "555"}, {"foo", "789"}, {"corge", "456"}, {"bar", "123"}}));
3457 
3458 KJ_ASSERT(expectCached(test.listReverse("bar", "qux")) ==
3459 kvs({{"garply", "555"}, {"foo", "789"}, {"corge", "987"}}));
3460 
3461 // There will be two flushes waiting since the put of "foo" will have started before the
3462 // delete of "garply"
3463 {
3464 mockStorage->expectCall("put", ws)
3465 .withParams(CAPNP(entries = [(key = "corge", value = "987")]))
3466 .thenReturn(CAPNP());
3467 }
3468 { mockStorage->expectCall("delete", ws).withParams(CAPNP(keys = ["bar"])).thenReturn(CAPNP()); }
3469}
3470 
3471KJ_TEST("ActorCache listReverse() end of first block deleted at inopportune time") {
3472 ActorCacheTest test;
3473 auto& ws = test.ws;
3474 auto& mockStorage = test.mockStorage;
3475 
3476 // Do a delete, wait for the commit... and then hold it open.
3477 auto deletePromise = expectUncached(test.delete_("corge"));
3478 
3479 auto mockDelete = mockStorage->expectCall("delete", ws).withParams(CAPNP(keys = ["corge"]));
3480 
3481 // Now do a list.
3482 auto promise = expectUncached(test.listReverse("bar", "qux"));
3483 
3484 mockStorage->expectCall("list", ws)
3485 .withParams(CAPNP(start = "bar", end = "qux", reverse = true), "stream"_kj)
3486 .useCallback("stream", [&](MockClient stream) {
3487 // First block ends at the deleted entry.
3488 stream
3489 .call("values",
3490 CAPNP(list = [ (key = "foo", value = "456"), (key = "corge", value = "123") ]))
3491 .expectReturns(CAPNP(), ws);
3492 
3493 // Let the delete finish. So now the last key in the first block is cached as a negative
3494 // clean entry.
3495 kj::mv(mockDelete).thenReturn(CAPNP());
3496 
3497 // Continue on.
3498 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3499 }).expectCanceled();
3500 
3501 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "456"}}));
3502 
3503 KJ_ASSERT(expectCached(test.listReverse("bar", "qux")) == kvs({{"foo", "456"}}));
3504}
3505 
3506KJ_TEST("ActorCache listReverse() retry on failure") {
3507 ActorCacheTest test;
3508 auto& ws = test.ws;
3509 auto& mockStorage = test.mockStorage;
3510 
3511 {
3512 auto promise = expectUncached(test.listReverse("bar", "qux", 4));
3513 
3514 mockStorage->expectCall("list", ws)
3515 .withParams(CAPNP(start = "bar", end = "qux", limit = 4, reverse = true), "stream"_kj)
3516 .useCallback("stream", [&](MockClient stream) {
3517 stream
3518 .call("values",
3519 CAPNP(list = [ (key = "foo", value = "123"), (key = "baz", value = "789") ]))
3520 .expectReturns(CAPNP(), ws);
3521 }).thenThrow(KJ_EXCEPTION(DISCONNECTED, "oops"));
3522 
3523 // Retry starts from `baz`.
3524 mockStorage->expectCall("list", ws)
3525 .withParams(CAPNP(start = "bar", end = "baz", limit = 2, reverse = true), "stream"_kj)
3526 .useCallback("stream", [&](MockClient stream) {
3527 // Duplicates of earlier keys will be ignored.
3528 stream
3529 .call("values",
3530 CAPNP(list =
3531 [
3532 (key = "foo", value = "IGNORE"), (key = "baz", value = "IGNORE"),
3533 (key = "bar", value = "456")
3534 ]))
3535 .expectReturns(CAPNP(), ws);
3536 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3537 }).expectCanceled();
3538 
3539 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}}));
3540 }
3541 
3542 KJ_ASSERT(expectCached(test.listReverse("bar", "qux")) ==
3543 kvs({{"foo", "123"}, {"baz", "789"}, {"bar", "456"}}));
3544}
3545 
3546// =======================================================================================
3547// LRU purge
3548 
3549constexpr size_t ENTRY_SIZE = 120;
3550KJ_TEST("ActorCache LRU purge") {
3551 ActorCacheTest test({.softLimit = 1 * ENTRY_SIZE});
3552 auto& ws = test.ws;
3553 auto& mockStorage = test.mockStorage;
3554 
3555 auto promise = expectUncached(test.get("foo"));
3556 mockStorage->expectCall("get", ws)
3557 .withParams(CAPNP(key = "foo"))
3558 .thenReturn(CAPNP(value = "123"));
3559 
3560 KJ_ASSERT(KJ_ASSERT_NONNULL(promise.wait(ws)) == "123");
3561 
3562 // Still cached.
3563 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
3564 
3565 promise = expectUncached(test.get("bar"));
3566 mockStorage->expectCall("get", ws)
3567 .withParams(CAPNP(key = "bar"))
3568 .thenReturn(CAPNP(value = "456"));
3569 
3570 KJ_ASSERT(KJ_ASSERT_NONNULL(promise.wait(ws)) == "456");
3571 
3572 // Still cached.
3573 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == "456");
3574 
3575 // But foo was evicted.
3576 (void)expectUncached(test.get("foo"));
3577}
3578 
3579KJ_TEST("ActorCache LRU purge ordering") {
3580 ActorCacheTest test({.softLimit = 4 * ENTRY_SIZE});
3581 auto& ws = test.ws;
3582 auto& mockStorage = test.mockStorage;
3583 
3584 test.put("foo", "123");
3585 test.put("bar", "456");
3586 test.put("baz", "789");
3587 test.put("qux", "555");
3588 
3589 // Let the flush of the puts complete.
3590 mockStorage->expectCall("put", ws).thenReturn(CAPNP());
3591 
3592 // Ensure the flush actually completes (marking dirty entries as clean) before continuing.
3593 test.gate.wait(nullptr).wait(ws);
3594 
3595 // Touch foo.
3596 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
3597 
3598 // Write two new values to push things out.
3599 test.put("xxx", "aaa");
3600 test.put("yyy", "bbb");
3601 
3602 // More puts flushing.
3603 mockStorage->expectCall("put", ws).thenReturn(CAPNP());
3604 
3605 // Foo and qux live, bar and baz evicted.
3606 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
3607 (void)expectUncached(test.get("bar"));
3608 (void)expectUncached(test.get("baz"));
3609 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("qux"))) == "555");
3610 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("xxx"))) == "aaa");
3611 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("yyy"))) == "bbb");
3612}
3613 
3614KJ_TEST("ActorCache LRU purge larger") {
3615 ActorCacheTest test({.softLimit = 32 * ENTRY_SIZE});
3616 auto& ws = test.ws;
3617 auto& mockStorage = test.mockStorage;
3618 
3619 auto kilobyte = kj::str(kj::repeat('x', 1024));
3620 
3621 auto promise = expectUncached(test.get("foo"));
3622 mockStorage->expectCall("get", ws)
3623 .withParams(CAPNP(key = "foo"))
3624 .thenReturn(CAPNP(value = "123"));
3625 
3626 KJ_ASSERT(KJ_ASSERT_NONNULL(promise.wait(ws)) == "123");
3627 
3628 test.put("bar", kilobyte);
3629 test.put("baz", kilobyte);
3630 test.put("qux", kilobyte);
3631 
3632 // Still cached.
3633 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
3634 
3635 test.put("corge", kilobyte);
3636 
3637 // Dropped from cache, because the puts are in-flight and so cannot be dropped. This read gets
3638 // sent off before the puts above because the event loop hasn't been yielded yet.
3639 // TODO(cleanup): We hold onto the promise here (even though in theory it'd be fine to drop)
3640 // because the capnp-mock framework doesn't handle dropped client promises well (capnp destructs
3641 // the ReceivedCall before waitForEvent resolves and hands control back to expectCall, leaving
3642 // receivedPromises empty in expectCall).
3643 promise = expectUncached(test.get("foo"));
3644 
3645 test.put("grault", kilobyte);
3646 test.put("garply", kilobyte);
3647 
3648 // Everything dirty is still in cache despite exceeding cache bounds.
3649 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("bar"))) == kilobyte);
3650 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == kilobyte);
3651 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("qux"))) == kilobyte);
3652 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("corge"))) == kilobyte);
3653 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("grault"))) == kilobyte);
3654 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("garply"))) == kilobyte);
3655 
3656 {
3657 // We have to wait for the get before the flush since capnp-mock doesn't continue waiting
3658 // after receiving the first call, and in this case the first call received will be the get.
3659 mockStorage->expectCall("get", ws)
3660 .withParams(CAPNP(key = "foo"))
3661 .thenReturn(CAPNP(value = "123"));
3662 mockStorage->expectCall("put", ws).thenReturn(CAPNP());
3663 // Ensure the flush actually completes (marking dirty entries as clean) before continuing.
3664 test.gate.wait(nullptr).wait(ws);
3665 }
3666 
3667 (void)expectUncached(test.get("bar"));
3668 (void)expectUncached(test.get("baz"));
3669 (void)expectUncached(test.get("qux"));
3670 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("corge"))) == kilobyte);
3671 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("grault"))) == kilobyte);
3672 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("garply"))) == kilobyte);
3673}
3674 
3675KJ_TEST("ActorCache LRU purge") {
3676 ActorCacheTest test({.softLimit = 1});
3677 auto& ws = test.ws;
3678 auto& mockStorage = test.mockStorage;
3679 
3680 auto promise = expectUncached(test.get({"foo"_kj, "bar"_kj, "baz"_kj}));
3681 mockStorage->expectCall("getMultiple", ws)
3682 .withParams(CAPNP(keys = [ "bar", "baz", "foo" ]), "stream"_kj)
3683 .useCallback("stream", [&](MockClient stream) {
3684 stream
3685 .call("values",
3686 CAPNP(list =
3687 [
3688 (key = "bar", value = "456"), (key = "baz", value = "789"),
3689 (key = "foo", value = "123")
3690 ]))
3691 .expectReturns(CAPNP(), ws);
3692 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3693 }).expectCanceled();
3694 
3695 KJ_ASSERT(promise.wait(ws) == kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "123"}}));
3696 
3697 // Nothing was cached, because nothing fit in the LRU.
3698 KJ_ASSERT(test.lru.currentSize() == 0);
3699 (void)expectUncached(test.get("foo"));
3700 (void)expectUncached(test.get("bar"));
3701 (void)expectUncached(test.get("baz"));
3702}
3703 
3704KJ_TEST("ActorCache evict on timeout") {
3705 ActorCacheTest test;
3706 auto& ws = test.ws;
3707 auto& mockStorage = test.mockStorage;
3708 
3709 auto timePoint = kj::UNIX_EPOCH;
3710 KJ_ASSERT(test.cache.evictStale(timePoint) == kj::none);
3711 
3712 auto ackFlush = [&]() {
3713 mockStorage->expectCall("put", ws).thenReturn(CAPNP());
3714 // Ensure the flush actually completes (marking dirty entries as clean) before continuing.
3715 test.gate.wait(nullptr).wait(ws);
3716 };
3717 
3718 test.put("foo", "123");
3719 test.put("bar", "456");
3720 ackFlush();
3721 
3722 KJ_ASSERT(test.cache.evictStale(timePoint + 100 * kj::MILLISECONDS) == kj::none);
3723 KJ_ASSERT(test.cache.evictStale(timePoint + 200 * kj::MILLISECONDS) == kj::none);
3724 KJ_ASSERT(test.cache.evictStale(timePoint + 500 * kj::MILLISECONDS) == kj::none);
3725 
3726 expectCached(test.get("foo"));
3727 expectCached(test.get("bar"));
3728 
3729 KJ_ASSERT(test.cache.evictStale(timePoint + 1000 * kj::MILLISECONDS) == nullptr);
3730 // foo and bar are now stale
3731 
3732 // add baz
3733 test.put("baz", "789");
3734 ackFlush();
3735 
3736 // don't check foo because we want it to be evicted, but touch bar
3737 expectCached(test.get("bar"));
3738 
3739 KJ_ASSERT(test.cache.evictStale(timePoint + 2000 * kj::MILLISECONDS) == nullptr);
3740 // Now foo should be evicted and bar and baz stale.
3741 
3742 // Verify foo is evicted.
3743 (void)expectUncached(test.get("foo"));
3744 
3745 // Touch bar.
3746 expectCached(test.get("bar"));
3747 
3748 KJ_ASSERT(test.cache.evictStale(timePoint + 3000 * kj::MILLISECONDS) == nullptr);
3749 // Now baz should have been evicted, but bar is still here because we keep touching it.
3750 
3751 expectCached(test.get("bar"));
3752 (void)expectUncached(test.get("baz"));
3753}
3754 
3755KJ_TEST("ActorCache backpressure due to dirtyPressureThreshold") {
3756 ActorCacheTest test({.dirtyListByteLimit = 2 * ENTRY_SIZE});
3757 auto& ws = test.ws;
3758 auto& mockStorage = test.mockStorage;
3759 
3760 auto timePoint = kj::UNIX_EPOCH;
3761 KJ_ASSERT(test.cache.evictStale(timePoint) == nullptr);
3762 
3763 KJ_ASSERT(test.put("foo", "123") == nullptr);
3764 KJ_ASSERT(test.put("bar", "456") == nullptr);
3765 auto promise1 = KJ_ASSERT_NONNULL(test.put("baz", "789"));
3766 auto promise2 = KJ_ASSERT_NONNULL(test.put("qux", "555"));
3767 
3768 // These deletes are actually cached, BUT backpressure will apply to make them return a promise.
3769 auto promise3 = expectUncached(test.delete_("baz"));
3770 auto promise4 = expectUncached(test.delete_({"qux"_kj}));
3771 
3772 // A delete of an unknown keys will also apply backpressure, of course.
3773 auto promise5 = expectUncached(test.delete_("corge"));
3774 auto promise6 = expectUncached(test.delete_({"grault"_kj}));
3775 
3776 // Let the write transaction complete.
3777 {
3778 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
3779 mockTxn->expectCall("delete", ws).thenReturn(CAPNP(numDeleted = 0));
3780 mockTxn->expectCall("delete", ws).thenReturn(CAPNP(numDeleted = 0));
3781 mockTxn->expectCall("delete", ws).thenReturn(CAPNP(numDeleted = 0));
3782 mockTxn->expectCall("put", ws).thenReturn(CAPNP());
3783 
3784 // Test for bogus `KJ_ASSERT(flushScheduled)` in `ActorCache::getBackpressure()`.
3785 auto promise7 = KJ_ASSERT_NONNULL(test.cache.evictStale(timePoint));
3786 
3787 KJ_ASSERT(!promise1.poll(ws));
3788 KJ_ASSERT(!promise2.poll(ws));
3789 KJ_ASSERT(!promise3.poll(ws));
3790 KJ_ASSERT(!promise4.poll(ws));
3791 KJ_ASSERT(!promise5.poll(ws));
3792 KJ_ASSERT(!promise6.poll(ws));
3793 KJ_ASSERT(!promise7.poll(ws));
3794 
3795 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
3796 
3797 promise1.wait(ws);
3798 promise2.wait(ws);
3799 KJ_ASSERT(promise3.wait(ws));
3800 KJ_ASSERT(promise4.wait(ws) == 1);
3801 KJ_ASSERT(!promise5.wait(ws));
3802 KJ_ASSERT(promise6.wait(ws) == 0);
3803 promise7.wait(ws);
3804 
3805 mockTxn->expectDropped(ws);
3806 }
3807}
3808 
3809KJ_TEST("ActorCache lru evict entry with known-empty gaps") {
3810 ActorCacheTest test({.softLimit = 5 * ENTRY_SIZE});
3811 auto& ws = test.ws;
3812 auto& mockStorage = test.mockStorage;
3813 
3814 // Populate cache.
3815 {
3816 auto promise = expectUncached(test.list("bar", "qux"));
3817 
3818 mockStorage->expectCall("list", ws)
3819 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
3820 .useCallback("stream", [&](MockClient stream) {
3821 stream
3822 .call("values",
3823 CAPNP(list =
3824 [
3825 (key = "bar", value = "456"), (key = "baz", value = "789"),
3826 (key = "corge", value = "555"), (key = "foo", value = "123")
3827 ]))
3828 .expectReturns(CAPNP(), ws);
3829 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3830 }).expectCanceled();
3831 
3832 KJ_ASSERT(promise.wait(ws) ==
3833 kvs({{"bar", "456"}, {"baz", "789"}, {"corge", "555"}, {"foo", "123"}}));
3834 }
3835 
3836 KJ_ASSERT(expectCached(test.list("bar", "qux")) ==
3837 kvs({{"bar", "456"}, {"baz", "789"}, {"corge", "555"}, {"foo", "123"}}));
3838 
3839 // touch some stuff so that "corge" is the oldest entry.
3840 expectCached(test.list("foo", "qux"));
3841 expectCached(test.get("bar"));
3842 expectCached(test.get("baz"));
3843 
3844 // do a put() to force an eviction.
3845 {
3846 test.put("xyzzy", "x");
3847 
3848 mockStorage->expectCall("put", ws)
3849 .withParams(CAPNP(entries = [(key = "xyzzy", value = "x")]))
3850 .thenReturn(CAPNP());
3851 }
3852 
3853 // The ranges before and after "corge" are missing, but everything else is still in cache.
3854 KJ_ASSERT(expectCached(test.list("bar", "baz")) == kvs({{"bar", "456"}}));
3855 KJ_ASSERT(expectCached(test.get("bay")) == nullptr);
3856 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "789");
3857 KJ_ASSERT(expectCached(test.list("foo", "qux")) == kvs({{"foo", "123"}}));
3858 KJ_ASSERT(expectCached(test.get("fooa")) == nullptr);
3859 
3860 (void)expectUncached(test.get("baza"));
3861 (void)expectUncached(test.get("corge"));
3862 (void)expectUncached(test.get("fo"));
3863}
3864 
3865KJ_TEST("ActorCache lru evict gap entry with known-empty gaps") {
3866 ActorCacheTest test({.softLimit = 5 * ENTRY_SIZE});
3867 auto& ws = test.ws;
3868 auto& mockStorage = test.mockStorage;
3869 
3870 // Populate cache.
3871 {
3872 auto promise = expectUncached(test.list("bar", "qux"));
3873 
3874 mockStorage->expectCall("list", ws)
3875 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
3876 .useCallback("stream", [&](MockClient stream) {
3877 stream
3878 .call("values",
3879 CAPNP(list =
3880 [
3881 (key = "bar", value = "456"), (key = "baz", value = "789"),
3882 (key = "corge", value = "555"), (key = "foo", value = "123")
3883 ]))
3884 .expectReturns(CAPNP(), ws);
3885 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3886 }).expectCanceled();
3887 
3888 KJ_ASSERT(promise.wait(ws) ==
3889 kvs({{"bar", "456"}, {"baz", "789"}, {"corge", "555"}, {"foo", "123"}}));
3890 }
3891 
3892 KJ_ASSERT(expectCached(test.list("bar", "qux")) ==
3893 kvs({{"bar", "456"}, {"baz", "789"}, {"corge", "555"}, {"foo", "123"}}));
3894 
3895 // touch some stuff so that "qux" is the oldest entry.
3896 expectCached(test.get("bar"));
3897 expectCached(test.get("baz"));
3898 expectCached(test.get("corge"));
3899 expectCached(test.get("foo"));
3900 
3901 // We still have a cached gap between "foo" and "qux".
3902 KJ_ASSERT(expectCached(test.get("foo+1")) == nullptr);
3903 
3904 // do a put() to force an eviction.
3905 {
3906 test.put("xyzzy", "x");
3907 
3908 mockStorage->expectCall("put", ws)
3909 .withParams(CAPNP(entries = [(key = "xyzzy", value = "x")]))
3910 .thenReturn(CAPNP());
3911 }
3912 
3913 // Okay, that gap is gone now.
3914 (void)expectUncached(test.get("foo+1"));
3915}
3916 
3917KJ_TEST("ActorCache lru evict entry with trailing known-empty gap (followed by END_GAP)") {
3918 ActorCacheTest test({.softLimit = 5 * ENTRY_SIZE});
3919 auto& ws = test.ws;
3920 auto& mockStorage = test.mockStorage;
3921 
3922 // Populate cache.
3923 {
3924 auto promise = expectUncached(test.list("bar", "qux"));
3925 
3926 mockStorage->expectCall("list", ws)
3927 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
3928 .useCallback("stream", [&](MockClient stream) {
3929 stream
3930 .call("values",
3931 CAPNP(list =
3932 [
3933 (key = "bar", value = "456"), (key = "baz", value = "789"),
3934 (key = "corge", value = "555"), (key = "foo", value = "123")
3935 ]))
3936 .expectReturns(CAPNP(), ws);
3937 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3938 }).expectCanceled();
3939 
3940 KJ_ASSERT(promise.wait(ws) ==
3941 kvs({{"bar", "456"}, {"baz", "789"}, {"corge", "555"}, {"foo", "123"}}));
3942 }
3943 
3944 KJ_ASSERT(expectCached(test.list("bar", "qux")) ==
3945 kvs({{"bar", "456"}, {"baz", "789"}, {"corge", "555"}, {"foo", "123"}}));
3946 
3947 // touch some stuff so that "foo" is the oldest entry.
3948 expectCached(test.get("bar"));
3949 expectCached(test.get("baz"));
3950 expectCached(test.get("corge"));
3951 
3952 // do a put() to force an eviction.
3953 {
3954 test.put("xyzzy", "x");
3955 
3956 mockStorage->expectCall("put", ws)
3957 .withParams(CAPNP(entries = [(key = "xyzzy", value = "x")]))
3958 .thenReturn(CAPNP());
3959 }
3960 
3961 // The range after "foo" is missing, but everything else is still in cache.
3962 KJ_ASSERT(expectCached(test.list("bar", "corge")) == kvs({{"bar", "456"}, {"baz", "789"}}));
3963 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("corge"))) == "555");
3964 
3965 (void)expectUncached(test.get("corgf"));
3966 (void)expectUncached(test.get("foo"));
3967 (void)expectUncached(test.get("quw"));
3968 (void)expectUncached(test.get("qux"));
3969 (void)expectUncached(test.get("quy"));
3970}
3971 
3972KJ_TEST("ActorCache timeout entry with known-empty gaps") {
3973 ActorCacheTest test({.softLimit = 5 * ENTRY_SIZE});
3974 auto& ws = test.ws;
3975 auto& mockStorage = test.mockStorage;
3976 
3977 auto startTime = kj::UNIX_EPOCH;
3978 test.cache.evictStale(startTime);
3979 
3980 // Populate cache.
3981 {
3982 auto promise = expectUncached(test.list("bar", "qux"));
3983 
3984 mockStorage->expectCall("list", ws)
3985 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
3986 .useCallback("stream", [&](MockClient stream) {
3987 stream
3988 .call("values",
3989 CAPNP(list =
3990 [
3991 (key = "bar", value = "456"), (key = "baz", value = "789"),
3992 (key = "corge", value = "555"), (key = "foo", value = "123")
3993 ]))
3994 .expectReturns(CAPNP(), ws);
3995 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
3996 }).expectCanceled();
3997 
3998 KJ_ASSERT(promise.wait(ws) ==
3999 kvs({{"bar", "456"}, {"baz", "789"}, {"corge", "555"}, {"foo", "123"}}));
4000 }
4001 
4002 KJ_ASSERT(expectCached(test.list("bar", "qux")) ==
4003 kvs({{"bar", "456"}, {"baz", "789"}, {"corge", "555"}, {"foo", "123"}}));
4004 
4005 // Make all entries STALE.
4006 test.cache.evictStale(startTime + 1 * kj::SECONDS);
4007 
4008 // touch some stuff so that "corge" is the only STALE entry.
4009 expectCached(test.list("foo", "qux"));
4010 expectCached(test.get("bar"));
4011 expectCached(test.get("baz"));
4012 
4013 // Time out "corge".
4014 test.cache.evictStale(startTime + 2 * kj::SECONDS);
4015 
4016 // The ranges before and after "corge" are missing, but everything else is still in cache.
4017 KJ_ASSERT(expectCached(test.list("bar", "baz")) == kvs({{"bar", "456"}}));
4018 KJ_ASSERT(expectCached(test.get("bay")) == nullptr);
4019 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("baz"))) == "789");
4020 KJ_ASSERT(expectCached(test.list("foo", "qux")) == kvs({{"foo", "123"}}));
4021 KJ_ASSERT(expectCached(test.get("fooa")) == nullptr);
4022 
4023 (void)expectUncached(test.get("baza"));
4024 (void)expectUncached(test.get("corge"));
4025 (void)expectUncached(test.get("fo"));
4026}
4027 
4028KJ_TEST("ActorCache evictStale entire list with end marker") {
4029 ActorCacheTest test;
4030 auto& ws = test.ws;
4031 auto& mockStorage = test.mockStorage;
4032 
4033 auto timePoint = kj::UNIX_EPOCH;
4034 KJ_ASSERT(test.cache.evictStale(timePoint) == nullptr);
4035 
4036 {
4037 // Populate a decent list.
4038 auto promise = expectUncached(test.list("bar", "qux"));
4039 
4040 mockStorage->expectCall("list", ws)
4041 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
4042 .useCallback("stream", [&](MockClient stream) {
4043 stream
4044 .call("values",
4045 CAPNP(list =
4046 [
4047 (key = "bar", value = "456"), (key = "baz", value = "789"),
4048 (key = "corge", value = "555"), (key = "foo", value = "123")
4049 ]))
4050 .expectReturns(CAPNP(), ws);
4051 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
4052 }).expectCanceled();
4053 
4054 KJ_ASSERT(promise.wait(ws) ==
4055 kvs({{"bar", "456"}, {"baz", "789"}, {"corge", "555"}, {"foo", "123"}}));
4056 }
4057 
4058 KJ_EXPECT(test.lru.currentSize() > 0);
4059 
4060 // First mark the entire cache as stale.
4061 timePoint += 1 * kj::SECONDS;
4062 KJ_ASSERT(test.cache.evictStale(timePoint) == nullptr);
4063 KJ_EXPECT(test.lru.currentSize() > 0);
4064 
4065 // Evict the entire cache.
4066 timePoint += 1 * kj::SECONDS;
4067 KJ_ASSERT(test.cache.evictStale(timePoint) == nullptr);
4068 KJ_EXPECT(test.lru.currentSize() == 0);
4069}
4070 
4071KJ_TEST("ActorCache purge everything while listing") {
4072 ActorCacheTest test({.softLimit = 1}); // evict everything immediately
4073 auto& ws = test.ws;
4074 auto& mockStorage = test.mockStorage;
4075 
4076 {
4077 auto promise = expectUncached(test.list("bar", "qux"));
4078 
4079 mockStorage->expectCall("list", ws)
4080 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
4081 .useCallback("stream", [&](MockClient stream) {
4082 stream
4083 .call("values",
4084 CAPNP(list = [ (key = "bar", value = "456"), (key = "baz", value = "789") ]))
4085 .expectReturns(CAPNP(), ws);
4086 stream
4087 .call("values",
4088 CAPNP(list = [ (key = "corge", value = "555"), (key = "foo", value = "123") ]))
4089 .expectReturns(CAPNP(), ws);
4090 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
4091 }).expectCanceled();
4092 
4093 KJ_ASSERT(promise.wait(ws) ==
4094 kvs({{"bar", "456"}, {"baz", "789"}, {"corge", "555"}, {"foo", "123"}}));
4095 }
4096 
4097 (void)expectUncached(test.get("bar"));
4098 (void)expectUncached(test.get("baz"));
4099 (void)expectUncached(test.get("corge"));
4100 (void)expectUncached(test.get("foo"));
4101}
4102 
4103KJ_TEST("ActorCache purge everything while listing; has previous entry") {
4104 ActorCacheTest test({.softLimit = 1}); // evict everything immediately
4105 auto& ws = test.ws;
4106 auto& mockStorage = test.mockStorage;
4107 
4108 // This is the same as the previous test, except we put an entry into cache first that appears
4109 // before the list range. This exercises a slightly different code path in markGapsEmpty().
4110 test.put("a", "x");
4111 
4112 {
4113 auto promise = expectUncached(test.list("bar", "qux"));
4114 
4115 mockStorage->expectCall("list", ws)
4116 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
4117 .useCallback("stream", [&](MockClient stream) {
4118 stream
4119 .call("values",
4120 CAPNP(list = [ (key = "bar", value = "456"), (key = "baz", value = "789") ]))
4121 .expectReturns(CAPNP(), ws);
4122 stream
4123 .call("values",
4124 CAPNP(list = [ (key = "corge", value = "555"), (key = "foo", value = "123") ]))
4125 .expectReturns(CAPNP(), ws);
4126 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
4127 }).expectCanceled();
4128 
4129 KJ_ASSERT(promise.wait(ws) ==
4130 kvs({{"bar", "456"}, {"baz", "789"}, {"corge", "555"}, {"foo", "123"}}));
4131 }
4132 
4133 // Acknowledge the flush.
4134 mockStorage->expectCall("put", ws).thenReturn(CAPNP());
4135}
4136 
4137KJ_TEST("ActorCache exceed hard limit on read") {
4138 ActorCacheTest test(
4139 {.monitorOutputGate = false, .softLimit = 2 * ENTRY_SIZE, .hardLimit = 2 * ENTRY_SIZE});
4140 auto& ws = test.ws;
4141 auto& mockStorage = test.mockStorage;
4142 
4143 auto brokenPromise = test.gate.onBroken();
4144 
4145 {
4146 // Don't use expectUncached() since it will log exceptions as test failures.
4147 auto promise = test.list("bar"_kj, "qux"_kj).get<kj::Promise<kj::Array<KeyValue>>>();
4148 
4149 mockStorage->expectCall("list", ws)
4150 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
4151 .useCallback("stream", [&](MockClient stream) {
4152 stream
4153 .call("values",
4154 CAPNP(list = [ (key = "bar", value = "456"), (key = "baz", value = "789") ]))
4155 .expectReturns(CAPNP(), ws);
4156 KJ_ASSERT(!brokenPromise.poll(ws));
4157 
4158 // The next value delivered overflows the cache.
4159 stream.call("values", CAPNP(list = [(key = "corge", value = "555")]))
4160 .expectThrows(kj::Exception::Type::OVERLOADED,
4161 "exceeded its memory limit due to overflowing the storage cache", ws);
4162 
4163 KJ_ASSERT(brokenPromise.poll(ws));
4164 
4165 // The exception propagates to further calls do to capnp streaming semantics.
4166 stream.call("values", CAPNP(list = [(key = "foo", value = "123")]))
4167 .expectThrows(kj::Exception::Type::OVERLOADED,
4168 "exceeded its memory limit due to overflowing the storage cache", ws);
4169 stream.call("end", CAPNP())
4170 .expectThrows(kj::Exception::Type::OVERLOADED,
4171 "exceeded its memory limit due to overflowing the storage cache", ws);
4172 
4173 // The call will actually have been canceled when the first call failed.
4174 }).expectCanceled();
4175 
4176 KJ_EXPECT_THROW_MESSAGE(
4177 "exceeded its memory limit due to overflowing the storage cache", promise.wait(ws));
4178 }
4179 
4180 KJ_EXPECT_THROW_MESSAGE(
4181 "exceeded its memory limit due to overflowing the storage cache", brokenPromise.wait(ws));
4182}
4183 
4184KJ_TEST("ActorCache exceed hard limit on write") {
4185 ActorCacheTest test(
4186 {.monitorOutputGate = false, .softLimit = 2 * ENTRY_SIZE, .hardLimit = 2 * ENTRY_SIZE});
4187 auto& ws = test.ws;
4188 
4189 auto brokenPromise = test.gate.onBroken();
4190 
4191 test.put("foo", "123");
4192 test.put("bar", "456");
4193 KJ_EXPECT_THROW_MESSAGE(
4194 "exceeded its memory limit due to overflowing the storage cache", test.put("baz", "789"));
4195 
4196 KJ_ASSERT(brokenPromise.poll(ws));
4197 KJ_EXPECT_THROW_MESSAGE(
4198 "exceeded its memory limit due to overflowing the storage cache", brokenPromise.wait(ws));
4199}
4200 
4201// =======================================================================================
4202 
4203KJ_TEST("ActorCache skip cache") {
4204 ActorCacheTest test;
4205 auto& ws = test.ws;
4206 auto& mockStorage = test.mockStorage;
4207 
4208 // Read a value.
4209 {
4210 auto promise = expectUncached(test.get("foo", {.noCache = true}));
4211 
4212 mockStorage->expectCall("get", ws)
4213 .withParams(CAPNP(key = "foo"))
4214 .thenReturn(CAPNP(value = "bar"));
4215 
4216 auto result = KJ_ASSERT_NONNULL(promise.wait(ws));
4217 KJ_EXPECT(result == "bar");
4218 }
4219 
4220 // Read it again -- not in cache!
4221 {
4222 auto promise = expectUncached(test.get("foo", {.noCache = true}));
4223 
4224 mockStorage->expectCall("get", ws)
4225 .withParams(CAPNP(key = "foo"))
4226 .thenReturn(CAPNP(value = "baz"));
4227 
4228 auto result = KJ_ASSERT_NONNULL(promise.wait(ws));
4229 KJ_EXPECT(result == "baz");
4230 }
4231 
4232 // Put a value.
4233 {
4234 test.put("foo", "qux", {.noCache = true});
4235 
4236 // If we read it right now, it's in cache.
4237 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo", {.noCache = true}))) == "qux");
4238 
4239 mockStorage->expectCall("put", ws)
4240 .withParams(CAPNP(entries = [(key = "foo", value = "qux")]))
4241 .thenReturn(CAPNP());
4242 }
4243 
4244 // Wait on the output gate to make sure the flush is actually done.
4245 test.gate.wait(nullptr).wait(test.ws);
4246 
4247 // After the put completes, it's not in cache anymore.
4248 {
4249 auto promise = expectUncached(test.get("foo", {.noCache = true}));
4250 
4251 mockStorage->expectCall("get", ws)
4252 .withParams(CAPNP(key = "foo"))
4253 .thenReturn(CAPNP(value = "baz"));
4254 
4255 auto result = KJ_ASSERT_NONNULL(promise.wait(ws));
4256 KJ_EXPECT(result == "baz");
4257 }
4258 
4259 // Do it again. This time, though, the read that happens while dirty doesn't have .noCache.
4260 {
4261 test.put("foo", "qux", {.noCache = true});
4262 
4263 // If we read it right now, it's in cache.
4264 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "qux");
4265 
4266 mockStorage->expectCall("put", ws)
4267 .withParams(CAPNP(entries = [(key = "foo", value = "qux")]))
4268 .thenReturn(CAPNP());
4269 }
4270 
4271 // Wait on the output gate to make sure the flush is actually done.
4272 test.gate.wait(nullptr).wait(test.ws);
4273 
4274 // This time it stayed in cache!
4275 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "qux");
4276 
4277 // Do an uncached list.
4278 {
4279 auto promise = expectUncached(test.list("bar", "qux", kj::none, {.noCache = true}));
4280 
4281 mockStorage->expectCall("list", ws)
4282 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
4283 .useCallback("stream", [&](MockClient stream) {
4284 stream
4285 .call("values",
4286 CAPNP(list =
4287 [
4288 (key = "bar", value = "456"), (key = "baz", value = "789"),
4289 (key = "foo", value = "123")
4290 ]))
4291 .expectReturns(CAPNP(), ws);
4292 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
4293 }).expectCanceled();
4294 
4295 KJ_ASSERT(promise.wait(ws) == kvs({{"bar", "456"}, {"baz", "789"}, {"foo", "qux"}}));
4296 }
4297 
4298 // `foo` is still cached.
4299 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "qux");
4300 
4301 // The other things that were returned weren't cached.
4302 (void)expectUncached(test.get("bar"));
4303 (void)expectUncached(test.get("baz"));
4304 
4305 // No gaps were cached as empty either.
4306 (void)expectUncached(test.get("corge"));
4307 (void)expectUncached(test.get("grault"));
4308 
4309 // Again, but reverse list.
4310 {
4311 auto promise = expectUncached(test.listReverse("bar", "qux", kj::none, {.noCache = true}));
4312 
4313 mockStorage->expectCall("list", ws)
4314 .withParams(CAPNP(start = "bar", end = "qux", reverse = true), "stream"_kj)
4315 .useCallback("stream", [&](MockClient stream) {
4316 stream
4317 .call("values",
4318 CAPNP(list =
4319 [
4320 (key = "foo", value = "123"), (key = "baz", value = "789"),
4321 (key = "bar", value = "456")
4322 ]))
4323 .expectReturns(CAPNP(), ws);
4324 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
4325 }).expectCanceled();
4326 
4327 KJ_ASSERT(promise.wait(ws) == kvs({{"foo", "qux"}, {"baz", "789"}, {"bar", "456"}}));
4328 }
4329 
4330 // `foo` is still cached.
4331 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "qux");
4332 
4333 // The other things that were returned weren't cached.
4334 (void)expectUncached(test.get("bar"));
4335 (void)expectUncached(test.get("baz"));
4336 
4337 // No gaps were cached as empty either.
4338 (void)expectUncached(test.get("corge"));
4339 (void)expectUncached(test.get("grault"));
4340}
4341 
4342// =======================================================================================
4343 
4344KJ_TEST("ActorCache transaction read-through") {
4345 ActorCacheTest test;
4346 auto& ws = test.ws;
4347 auto& mockStorage = test.mockStorage;
4348 
4349 ActorCache::Transaction txn(test.cache);
4350 ActorCacheConvenienceWrappers eztxn(txn);
4351 
4352 {
4353 auto promise = expectUncached(eztxn.get("foo"));
4354 mockStorage->expectCall("get", ws)
4355 .withParams(CAPNP(key = "foo"))
4356 .thenReturn(CAPNP(value = "123"));
4357 KJ_ASSERT(KJ_ASSERT_NONNULL(promise.wait(ws)) == "123");
4358 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(eztxn.get("foo"))) == "123");
4359 }
4360 
4361 {
4362 auto promise = expectUncached(eztxn.get("bar"));
4363 mockStorage->expectCall("get", ws).withParams(CAPNP(key = "bar")).thenReturn(CAPNP());
4364 KJ_ASSERT(promise.wait(ws) == nullptr);
4365 KJ_ASSERT(expectCached(eztxn.get("bar")) == nullptr);
4366 }
4367 
4368 {
4369 auto promise = expectUncached(eztxn.get({"baz"_kj, "qux"_kj, "corge"_kj}));
4370 
4371 mockStorage->expectCall("getMultiple", ws)
4372 .withParams(CAPNP(keys = [ "baz", "corge", "qux" ]), "stream"_kj)
4373 .useCallback("stream", [&](MockClient stream) {
4374 stream
4375 .call("values",
4376 CAPNP(list = [ (key = "baz", value = "456"), (key = "qux", value = "789") ]))
4377 .expectReturns(CAPNP(), ws);
4378 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
4379 }).expectCanceled();
4380 
4381 KJ_ASSERT(promise.wait(ws) == kvs({{"baz", "456"}, {"qux", "789"}}));
4382 
4383 KJ_ASSERT(expectCached(eztxn.get({"foo"_kj, "bar"_kj, "baz"_kj, "qux"_kj, "corge"_kj})) ==
4384 kvs({{"baz", "456"}, {"foo", "123"}, {"qux", "789"}}));
4385 }
4386 
4387 {
4388 auto promise = expectUncached(eztxn.list("a", "z", 10));
4389 
4390 mockStorage->expectCall("list", ws)
4391 .withParams(CAPNP(start = "a", end = "z", limit = 10), "stream"_kj)
4392 .useCallback("stream", [&](MockClient stream) {
4393 stream
4394 .call("values",
4395 CAPNP(list =
4396 [
4397 (key = "baz", value = "456"), (key = "foo", value = "123"),
4398 (key = "grault", value = "555"), (key = "qux", value = "789")
4399 ]))
4400 .expectReturns(CAPNP(), ws);
4401 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
4402 }).expectCanceled();
4403 
4404 KJ_ASSERT(promise.wait(ws) ==
4405 kvs({{"baz", "456"}, {"foo", "123"}, {"grault", "555"}, {"qux", "789"}}));
4406 
4407 KJ_ASSERT(expectCached(eztxn.list("a", "z", 10)) ==
4408 kvs({{"baz", "456"}, {"foo", "123"}, {"grault", "555"}, {"qux", "789"}}));
4409 
4410 KJ_ASSERT(expectCached(eztxn.listReverse("a", "z")) ==
4411 kvs({{"qux", "789"}, {"grault", "555"}, {"foo", "123"}, {"baz", "456"}}));
4412 }
4413}
4414 
4415KJ_TEST("ActorCache transaction overlay changes") {
4416 ActorCacheTest test;
4417 auto& ws = test.ws;
4418 auto& mockStorage = test.mockStorage;
4419 
4420 ActorCache::Transaction txn(test.cache);
4421 ActorCacheConvenienceWrappers eztxn(txn);
4422 
4423 eztxn.put("foo", "321");
4424 eztxn.put({{"bar", "654"}, {"qux", "987"}});
4425 auto deletePromise1 = expectUncached(eztxn.delete_("grault"));
4426 auto deletePromise2 = expectUncached(eztxn.delete_({"baz"_kj, "garply"_kj}));
4427 
4428 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(eztxn.get("foo"))) == "321");
4429 KJ_ASSERT(expectCached(eztxn.get("baz")) == nullptr);
4430 KJ_ASSERT(expectCached(eztxn.get({"bar"_kj, "baz"_kj, "qux"_kj})) ==
4431 kvs({{"bar", "654"}, {"qux", "987"}}));
4432 
4433 // The deletes will force reads in order to compute counts.
4434 mockStorage->expectCall("get", ws)
4435 .withParams(CAPNP(key = "grault"))
4436 .thenReturn(CAPNP(value = "555"));
4437 mockStorage->expectCall("getMultiple", ws)
4438 .withParams(CAPNP(keys = [ "baz", "garply" ]), "stream"_kj)
4439 .useCallback("stream", [&](MockClient stream) {
4440 stream.call("values", CAPNP(list = [(key = "baz", value = "456")])).expectReturns(CAPNP(), ws);
4441 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
4442 }).expectCanceled();
4443 
4444 KJ_ASSERT(deletePromise1.wait(ws));
4445 KJ_ASSERT(deletePromise2.wait(ws) == 1);
4446 
4447 {
4448 auto promise = expectUncached(eztxn.get({"baz"_kj, "qux"_kj, "corge"_kj}));
4449 
4450 mockStorage->expectCall("getMultiple", ws)
4451 .withParams(CAPNP(keys = ["corge"]), "stream"_kj)
4452 .useCallback("stream", [&](MockClient stream) {
4453 stream.call("values", CAPNP(list = [])).expectReturns(CAPNP(), ws);
4454 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
4455 }).expectCanceled();
4456 
4457 KJ_ASSERT(promise.wait(ws) == kvs({{"qux", "987"}}));
4458 
4459 KJ_ASSERT(expectCached(eztxn.get({"foo"_kj, "bar"_kj, "baz"_kj, "qux"_kj, "corge"_kj})) ==
4460 kvs({{"bar", "654"}, {"foo", "321"}, {"qux", "987"}}));
4461 }
4462 
4463 {
4464 auto promise = expectUncached(eztxn.list("a", "z", 10));
4465 
4466 mockStorage
4467 ->expectCall("list", ws)
4468 // limit is adjusted by 3 because it could return values that have already been deleted
4469 // in the transaction.
4470 .withParams(CAPNP(start = "a", end = "z", limit = 13), "stream"_kj)
4471 .useCallback("stream", [&](MockClient stream) {
4472 stream
4473 .call("values",
4474 CAPNP(list =
4475 [
4476 (key = "baz", value = "456"), (key = "foo", value = "123"),
4477 (key = "grault", value = "555"), (key = "qux", value = "789")
4478 ]))
4479 .expectReturns(CAPNP(), ws);
4480 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
4481 }).expectCanceled();
4482 
4483 KJ_ASSERT(promise.wait(ws) == kvs({{"bar", "654"}, {"foo", "321"}, {"qux", "987"}}));
4484 
4485 KJ_ASSERT(expectCached(eztxn.list("a", "z", 10)) ==
4486 kvs({{"bar", "654"}, {"foo", "321"}, {"qux", "987"}}));
4487 
4488 KJ_ASSERT(expectCached(eztxn.listReverse("a", "z")) ==
4489 kvs({{"qux", "987"}, {"foo", "321"}, {"bar", "654"}}));
4490 }
4491 
4492 mockStorage->expectNoActivity(ws);
4493 
4494 txn.commit();
4495 
4496 {
4497 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
4498 mockTxn->expectCall("delete", ws)
4499 .withParams(CAPNP(keys = [ "grault", "baz" ]))
4500 .thenReturn(CAPNP(numDeleted = 2));
4501 mockTxn->expectCall("put", ws)
4502 .withParams(CAPNP(entries =
4503 [
4504 (key = "foo", value = "321"), (key = "bar", value = "654"),
4505 (key = "qux", value = "987")
4506 ]))
4507 .thenReturn(CAPNP());
4508 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
4509 mockTxn->expectDropped(ws);
4510 }
4511}
4512 
4513KJ_TEST("ActorCache transaction overlay changes precached") {
4514 // Like previous test, but have the range cached in the underlying cache before the transaction
4515 // touches it.
4516 
4517 ActorCacheTest test;
4518 auto& ws = test.ws;
4519 auto& mockStorage = test.mockStorage;
4520 
4521 {
4522 auto promise = expectUncached(test.listReverse("a", "z", 10));
4523 
4524 mockStorage->expectCall("list", ws)
4525 .withParams(CAPNP(start = "a", end = "z", limit = 10, reverse = true), "stream"_kj)
4526 .useCallback("stream", [&](MockClient stream) {
4527 stream
4528 .call("values",
4529 CAPNP(list =
4530 [
4531 (key = "qux", value = "789"), (key = "grault", value = "555"),
4532 (key = "foo", value = "123"), (key = "baz", value = "456")
4533 ]))
4534 .expectReturns(CAPNP(), ws);
4535 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
4536 }).expectCanceled();
4537 
4538 KJ_ASSERT(promise.wait(ws) ==
4539 kvs({{"qux", "789"}, {"grault", "555"}, {"foo", "123"}, {"baz", "456"}}));
4540 }
4541 
4542 ActorCache::Transaction txn(test.cache);
4543 ActorCacheConvenienceWrappers eztxn(txn);
4544 
4545 eztxn.put("foo", "321");
4546 eztxn.put({{"bar", "654"}, {"qux", "987"}});
4547 KJ_ASSERT(expectCached(eztxn.delete_("grault")));
4548 KJ_ASSERT(expectCached(eztxn.delete_({"baz"_kj, "garply"_kj})) == 1);
4549 
4550 KJ_ASSERT(KJ_ASSERT_NONNULL(expectCached(eztxn.get("foo"))) == "321");
4551 KJ_ASSERT(expectCached(eztxn.get("baz")) == nullptr);
4552 KJ_ASSERT(expectCached(eztxn.get({"bar"_kj, "baz"_kj, "qux"_kj})) ==
4553 kvs({{"bar", "654"}, {"qux", "987"}}));
4554 
4555 KJ_ASSERT(expectCached(eztxn.get({"foo"_kj, "bar"_kj, "baz"_kj, "qux"_kj, "corge"_kj})) ==
4556 kvs({{"bar", "654"}, {"foo", "321"}, {"qux", "987"}}));
4557 KJ_ASSERT(expectCached(eztxn.list("a", "z", 10)) ==
4558 kvs({{"bar", "654"}, {"foo", "321"}, {"qux", "987"}}));
4559 KJ_ASSERT(expectCached(eztxn.listReverse("a", "z")) ==
4560 kvs({{"qux", "987"}, {"foo", "321"}, {"bar", "654"}}));
4561 
4562 mockStorage->expectNoActivity(ws);
4563 
4564 txn.commit();
4565 
4566 {
4567 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
4568 mockTxn->expectCall("delete", ws)
4569 .withParams(CAPNP(keys = [ "grault", "baz" ]))
4570 .thenReturn(CAPNP(numDeleted = 2));
4571 mockTxn->expectCall("put", ws)
4572 .withParams(CAPNP(entries =
4573 [
4574 (key = "foo", value = "321"), (key = "bar", value = "654"),
4575 (key = "qux", value = "987")
4576 ]))
4577 .thenReturn(CAPNP());
4578 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
4579 mockTxn->expectDropped(ws);
4580 }
4581}
4582 
4583KJ_TEST("ActorCache transaction output gate blocked during flush") {
4584 ActorCacheTest test({.monitorOutputGate = false});
4585 auto& ws = test.ws;
4586 auto& mockStorage = test.mockStorage;
4587 
4588 // Gate is currently not blocked.
4589 test.gate.wait(nullptr).wait(ws);
4590 
4591 // Do a transaction with a put.
4592 ActorCache::Transaction txn(test.cache);
4593 ActorCacheConvenienceWrappers eztxn(txn);
4594 eztxn.put("foo", "123");
4595 txn.commit();
4596 
4597 // Now it is blocked.
4598 auto gatePromise = test.gate.wait(nullptr);
4599 KJ_ASSERT(!gatePromise.poll(ws));
4600 
4601 // Complete the transaction.
4602 {
4603 auto inProgressFlush = mockStorage->expectCall("put", ws).withParams(
4604 CAPNP(entries = [(key = "foo", value = "123")]));
4605 
4606 // Still blocked until the flush completes.
4607 KJ_ASSERT(!gatePromise.poll(ws));
4608 
4609 kj::mv(inProgressFlush).thenReturn(CAPNP());
4610 }
4611 
4612 KJ_ASSERT(gatePromise.poll(ws));
4613 gatePromise.wait(ws);
4614}
4615 
4616KJ_TEST("ActorCache transaction output gate bypass") {
4617 ActorCacheTest test({.monitorOutputGate = false});
4618 auto& ws = test.ws;
4619 auto& mockStorage = test.mockStorage;
4620 
4621 // Gate is currently not blocked.
4622 test.gate.wait(nullptr).wait(ws);
4623 
4624 // Do a transaction with a put.
4625 ActorCache::Transaction txn(test.cache);
4626 ActorCacheConvenienceWrappers eztxn(txn);
4627 eztxn.put("foo", "123", {.allowUnconfirmed = true});
4628 txn.commit();
4629 
4630 // Gate still isn't blocked, because we set `allowUnconfirmed`.
4631 test.gate.wait(nullptr).wait(ws);
4632 
4633 // Complete the transaction with a flush.
4634 mockStorage->expectCall("put", ws)
4635 .withParams(CAPNP(entries = [(key = "foo", value = "123")]))
4636 .thenReturn(CAPNP());
4637 
4638 test.gate.wait(nullptr).wait(ws);
4639}
4640 
4641KJ_TEST("ActorCache transaction output gate bypass on one put but not the next") {
4642 ActorCacheTest test({.monitorOutputGate = false});
4643 auto& ws = test.ws;
4644 auto& mockStorage = test.mockStorage;
4645 
4646 // Gate is currently not blocked.
4647 test.gate.wait(nullptr).wait(ws);
4648 
4649 // Do a transaction with two puts, only bypassing on the first. The net result should be that
4650 // the output gate is in effect.
4651 ActorCache::Transaction txn(test.cache);
4652 ActorCacheConvenienceWrappers eztxn(txn);
4653 eztxn.put("foo", "123", {.allowUnconfirmed = true});
4654 eztxn.put("bar", "456");
4655 txn.commit();
4656 
4657 // Now it is blocked.
4658 auto gatePromise = test.gate.wait(nullptr);
4659 KJ_ASSERT(!gatePromise.poll(ws));
4660 
4661 // Complete the transaction.
4662 {
4663 auto inProgressFlush = mockStorage->expectCall("put", ws).withParams(
4664 CAPNP(entries = [ (key = "foo", value = "123"), (key = "bar", value = "456") ]));
4665 
4666 // Still blocked until the flush completes.
4667 KJ_ASSERT(!gatePromise.poll(ws));
4668 
4669 kj::mv(inProgressFlush).thenReturn(CAPNP());
4670 }
4671 
4672 KJ_ASSERT(gatePromise.poll(ws));
4673 gatePromise.wait(ws);
4674}
4675 
4676KJ_TEST("ActorCache transaction multiple put batches") {
4677 ActorCacheTest test({.maxKeysPerRpc = 2});
4678 auto& ws = test.ws;
4679 auto& mockStorage = test.mockStorage;
4680 
4681 // Do a transaction with enough puts to batch.
4682 ActorCache::Transaction txn(test.cache);
4683 ActorCacheConvenienceWrappers eztxn(txn);
4684 eztxn.put({{"foo", "123"}, {"bar", "456"}, {"baz", "789"}});
4685 
4686 // Poll the wait scope to make sure we haven't slipped through to the cache already.
4687 ws.poll();
4688 
4689 eztxn.put({{"qux", "555"}, {"corge", "999"}});
4690 txn.commit();
4691 
4692 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
4693 mockTxn->expectCall("put", ws)
4694 .withParams(CAPNP(entries = [ (key = "foo", value = "123"), (key = "bar", value = "456") ]))
4695 .thenReturn(CAPNP());
4696 mockTxn->expectCall("put", ws)
4697 .withParams(CAPNP(entries = [ (key = "baz", value = "789"), (key = "qux", value = "555") ]))
4698 .thenReturn(CAPNP());
4699 mockTxn->expectCall("put", ws)
4700 .withParams(CAPNP(entries = [(key = "corge", value = "999")]))
4701 .thenReturn(CAPNP());
4702 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
4703 mockTxn->expectDropped(ws);
4704}
4705 
4706KJ_TEST("ActorCache transaction multiple counted delete batches") {
4707 // Do a transaction with a big counted delete. The rpc getMultiple and delete should batch
4708 // according to maxKeysPerRpc.
4709 
4710 ActorCacheTest test({.maxKeysPerRpc = 2});
4711 auto& ws = test.ws;
4712 auto& mockStorage = test.mockStorage;
4713 
4714 ActorCache::Transaction txn(test.cache);
4715 ActorCacheConvenienceWrappers eztxn(txn);
4716 
4717 {
4718 // Load one of our values to delete into the cache itself which will avoid rpc deletes for
4719 // counting.
4720 test.put("count2", "2");
4721 mockStorage->expectCall("put", ws)
4722 .withParams(CAPNP(entries = [(key = "count2", value = "2")]))
4723 .thenReturn(CAPNP());
4724 }
4725 
4726 {
4727 // Load one of our values to delete into the transaction which will avoid even talking to the
4728 // cache.
4729 eztxn.put("count3", "3");
4730 }
4731 
4732 auto deletePromise =
4733 eztxn.delete_({"count1"_kj, "count2"_kj, "count3"_kj, "count4"_kj, "count5"_kj})
4734 .get<kj::Promise<uint>>();
4735 
4736 mockStorage
4737 ->expectCall("getMultiple", ws)
4738 // Note that this batch is smaller because "count2" was known to the actor cache.
4739 .withParams(CAPNP(keys = ["count1"]), "stream"_kj)
4740 .useCallback("stream", [&](MockClient stream) {
4741 // Pretend that "count1" already exists but was not in the cache.
4742 stream.call("values", CAPNP(list = [(key = "count1", value = "1")])).expectReturns(CAPNP(), ws);
4743 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
4744 }).expectCanceled();
4745 mockStorage->expectCall("getMultiple", ws)
4746 .withParams(CAPNP(keys = [ "count4", "count5" ]), "stream"_kj)
4747 .useCallback("stream", [&](MockClient stream) {
4748 stream.call("end", CAPNP()).expectReturns(CAPNP(), ws);
4749 }).expectCanceled();
4750 
4751 // For hacky reasons, we are able to observe the counted delete before we submit the
4752 // transaction.
4753 KJ_EXPECT(deletePromise.wait(ws) == 3);
4754 
4755 txn.commit();
4756 
4757 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
4758 mockTxn
4759 ->expectCall("delete", ws)
4760 // "count3" comes first because it entered the transaction first.
4761 .withParams(CAPNP(keys = [ "count3", "count1" ]))
4762 .thenReturn(CAPNP(numDeleted = 1));
4763 mockTxn
4764 ->expectCall("delete", ws)
4765 // Neither "count4" or "count5" are deleted because we observed them in the get.
4766 .withParams(CAPNP(keys = ["count2"]))
4767 .thenReturn(CAPNP(numDeleted = 1));
4768 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
4769 mockTxn->expectDropped(ws);
4770}
4771 
4772KJ_TEST("ActorCache transaction negative list range returns nothing") {
4773 ActorCacheTest test({.monitorOutputGate = false});
4774 
4775 ActorCache::Transaction txn(test.cache);
4776 ActorCacheConvenienceWrappers eztxn(txn);
4777 
4778 eztxn.put("foo", "123");
4779 
4780 KJ_ASSERT(expectCached(eztxn.list("qux", "bar")) == kvs({}));
4781 KJ_ASSERT(expectCached(eztxn.listReverse("qux", "bar")) == kvs({}));
4782}
4783 
4784// =======================================================================================
4785 
4786KJ_TEST("ActorCache list stream cancellation") {
4787 // Test for cases where implementations of ListStream might stay alive longer than expected, due
4788 // to capabilities being held remotely.
4789 //
4790 // We can't use ActorCacheTest in this test as we need to manage allocation and destruction to
4791 // set up the problematic circumstances.
4792 
4793 kj::EventLoop loop;
4794 kj::WaitScope ws(loop);
4795 
4796 auto mockPair = MockServer::make<rpc::ActorStorage::Stage>();
4797 kj::Own<MockServer> mockStorage = kj::mv(mockPair.mock);
4798 auto mockClient = kj::mv(mockPair.client);
4799 
4800 ActorCacheTestOptions options;
4801 
4802 kj::Maybe<MockServer::ExpectedCall> call;
4803 kj::Maybe<MockClient> listClient;
4804 
4805 // Try get-multiple.
4806 {
4807 // We allocate `lru` on the heap to assist Valgrind in being able to detect when it is used
4808 // after free.
4809 auto lru = kj::heap<ActorCache::SharedLru>(ActorCache::SharedLru::Options{options.softLimit,
4810 options.hardLimit, options.staleTimeout, options.dirtyListByteLimit, options.maxKeysPerRpc});
4811 OutputGate gate;
4812 ActorCache cache(mockClient, *lru, gate);
4813 
4814 ActorCacheConvenienceWrappers ezCache(cache);
4815 
4816 auto promise = expectUncached(ezCache.get({"foo"_kj, "bar"_kj}));
4817 
4818 call = mockStorage->expectCall("getMultiple", ws)
4819 .withParams(CAPNP(keys = [ "bar", "foo" ]), "stream"_kj)
4820 .useCallback("stream", [&](MockClient stream) {
4821 stream.call("values", CAPNP(list = [(key = "bar", value = "123")]))
4822 .expectReturns(CAPNP(), ws);
4823 listClient = kj::mv(stream);
4824 });
4825 
4826 // Now we're going to cancel the promise and destroy the cache while the call is still
4827 // outstanding, with unreported entries in it!
4828 }
4829 
4830 KJ_ASSERT_NONNULL(listClient)
4831 .call("values", CAPNP(list = [(key = "foo", value = "456")]))
4832 .expectThrows(kj::Exception::Type::DISCONNECTED, "canceled", ws);
4833 KJ_ASSERT_NONNULL(kj::mv(call)).expectCanceled();
4834 
4835 // Try list().
4836 {
4837 // We allocate `lru` on the heap to assist Valgrind in being able to detect when it is used
4838 // after free.
4839 auto lru = kj::heap<ActorCache::SharedLru>(ActorCache::SharedLru::Options{options.softLimit,
4840 options.hardLimit, options.staleTimeout, options.dirtyListByteLimit, options.maxKeysPerRpc});
4841 OutputGate gate;
4842 ActorCache cache(mockClient, *lru, gate);
4843 
4844 ActorCacheConvenienceWrappers ezCache(cache);
4845 
4846 auto promise = expectUncached(ezCache.list("bar", "qux"));
4847 
4848 call = mockStorage->expectCall("list", ws)
4849 .withParams(CAPNP(start = "bar", end = "qux"), "stream"_kj)
4850 .useCallback("stream", [&](MockClient stream) {
4851 stream.call("values", CAPNP(list = [(key = "bar", value = "123")]))
4852 .expectReturns(CAPNP(), ws);
4853 listClient = kj::mv(stream);
4854 });
4855 
4856 // Now we're going to cancel the promise and destroy the cache while the call is still
4857 // outstanding, with unreported entries in it!
4858 }
4859 
4860 KJ_ASSERT_NONNULL(listClient)
4861 .call("values", CAPNP(list = [(key = "foo", value = "456")]))
4862 .expectThrows(kj::Exception::Type::DISCONNECTED, "canceled", ws);
4863 KJ_ASSERT_NONNULL(kj::mv(call)).expectCanceled();
4864 
4865 // Try listReverse().
4866 {
4867 // We allocate `lru` on the heap to assist Valgrind in being able to detect when it is used
4868 // after free.
4869 auto lru = kj::heap<ActorCache::SharedLru>(ActorCache::SharedLru::Options{options.softLimit,
4870 options.hardLimit, options.staleTimeout, options.dirtyListByteLimit, options.maxKeysPerRpc});
4871 OutputGate gate;
4872 ActorCache cache(mockClient, *lru, gate);
4873 
4874 ActorCacheConvenienceWrappers ezCache(cache);
4875 
4876 auto promise = expectUncached(ezCache.listReverse("bar", "qux"));
4877 
4878 call = mockStorage->expectCall("list", ws)
4879 .withParams(CAPNP(start = "bar", end = "qux", reverse = true), "stream"_kj)
4880 .useCallback("stream", [&](MockClient stream) {
4881 stream.call("values", CAPNP(list = [(key = "foo", value = "123")]))
4882 .expectReturns(CAPNP(), ws);
4883 listClient = kj::mv(stream);
4884 });
4885 
4886 // Now we're going to cancel the promise and destroy the cache while the call is still
4887 // outstanding, with unreported entries in it!
4888 }
4889 
4890 KJ_ASSERT_NONNULL(listClient)
4891 .call("values", CAPNP(list = [(key = "bar", value = "456")]))
4892 .expectThrows(kj::Exception::Type::DISCONNECTED, "canceled", ws);
4893 KJ_ASSERT_NONNULL(kj::mv(call)).expectCanceled();
4894}
4895 
4896KJ_TEST("ActorCache never-flush") {
4897 ActorCacheTest test({.neverFlush = true});
4898 auto& ws = test.ws;
4899 auto& mockStorage = test.mockStorage;
4900 
4901 // Puts don't start a transaction.
4902 KJ_EXPECT(test.put("foo", "123") == nullptr);
4903 KJ_EXPECT(test.cache.onNoPendingFlush(nullptr) == nullptr);
4904 mockStorage->expectNoActivity(ws);
4905 
4906 // Gets still see the put() value.
4907 KJ_EXPECT(KJ_ASSERT_NONNULL(expectCached(test.get("foo"))) == "123");
4908 
4909 // Uncached reads work normally.
4910 {
4911 auto promise = expectUncached(test.get("bar"));
4912 
4913 mockStorage->expectCall("get", ws)
4914 .withParams(CAPNP(key = "bar"))
4915 .thenReturn(CAPNP(value = "456"));
4916 
4917 auto result = KJ_ASSERT_NONNULL(promise.wait(ws));
4918 KJ_EXPECT(result == "456");
4919 }
4920}
4921 
4922KJ_TEST("ActorCache alarm get/put") {
4923 ActorCacheTest test;
4924 auto& ws = test.ws;
4925 auto& mockStorage = test.mockStorage;
4926 
4927 {
4928 auto time = expectUncached(test.getAlarm());
4929 
4930 mockStorage->expectCall("getAlarm", ws).thenReturn(CAPNP(scheduledTimeMs = 0));
4931 
4932 KJ_ASSERT(time.wait(ws) == kj::none);
4933 }
4934 
4935 {
4936 auto time = expectCached(test.getAlarm());
4937 
4938 KJ_ASSERT(time == kj::none);
4939 }
4940 
4941 auto oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH;
4942 auto twoMs = 2 * kj::MILLISECONDS + kj::UNIX_EPOCH;
4943 // Used as the "current time" parameter for armAlarmHandler in tests.
4944 auto testCurrentTime = kj::UNIX_EPOCH;
4945 {
4946 // Test alarm writes happen transactionally with storage ops
4947 test.setAlarm(oneMs);
4948 test.put("foo", "bar");
4949 
4950 auto mockTxn = mockStorage->expectCall("txn", ws).returnMock("transaction");
4951 mockTxn->expectCall("put", ws)
4952 .withParams(CAPNP(entries = [(key = "foo", value = "bar")]))
4953 .thenReturn(CAPNP());
4954 mockTxn->expectCall("setAlarm", ws).withParams(CAPNP(scheduledTimeMs = 1)).thenReturn(CAPNP());
4955 mockTxn->expectCall("commit", ws).thenReturn(CAPNP());
4956 mockTxn->expectDropped(ws);
4957 }
4958 
4959 {
4960 auto time = expectCached(test.getAlarm());
4961 
4962 KJ_ASSERT(time == oneMs);
4963 }
4964 
4965 {
4966 // Test clearing alarm
4967 test.setAlarm(kj::none);
4968 
4969 // When there are no other storage operations to be flushed, alarm modifications can be flushed
4970 // without a wrapping txn.
4971 mockStorage->expectCall("deleteAlarm", ws)
4972 .withParams(CAPNP(timeToDeleteMs = 0))
4973 .thenReturn(CAPNP(deleted = true));
4974 // Wait on the output gate to make sure the flush is actually done before checking the cache.
4975 test.gate.wait(nullptr).wait(test.ws);
4976 }
4977 
4978 {
4979 auto time = expectCached(test.getAlarm());
4980 
4981 KJ_ASSERT(time == kj::none);
4982 }
4983 
4984 {
4985 // we have a cached time == nullptr, so we should not attempt to run an alarm
4986 auto armResult =
4987 test.cache.armAlarmHandler(10 * kj::SECONDS + kj::UNIX_EPOCH, nullptr, testCurrentTime);
4988 KJ_ASSERT(armResult.is<ActorCache::CancelAlarmHandler>());
4989 auto cancelResult = kj::mv(armResult.get<ActorCache::CancelAlarmHandler>());
4990 KJ_ASSERT(cancelResult.waitBeforeCancel.poll(ws));
4991 cancelResult.waitBeforeCancel.wait(ws);
4992 }
4993 
4994 {
4995 test.setAlarm(oneMs);
4996 
4997 mockStorage->expectCall("setAlarm", ws)
4998 .withParams(CAPNP(scheduledTimeMs = 1))
4999 .thenReturn(CAPNP());
5000 }
5001 
5002 {
5003 // Test that alarm handler handle clears alarm when dropped with no writes
5004 {
5005 auto armResult = test.cache.armAlarmHandler(oneMs, nullptr, testCurrentTime);
5006 KJ_ASSERT(armResult.is<ActorCache::RunAlarmHandler>());
5007 }
5008 mockStorage->expectCall("deleteAlarm", ws)
5009 .withParams(CAPNP(timeToDeleteMs = 1))
5010 .thenReturn(CAPNP(deleted = true));
5011 }
5012 
5013 {
5014 test.setAlarm(oneMs);
5015 
5016 // Test that alarm handler handle does not clear alarm when dropped with writes
5017 {
5018 auto armResult = test.cache.armAlarmHandler(oneMs, nullptr, testCurrentTime);
5019 KJ_ASSERT(armResult.is<ActorCache::RunAlarmHandler>());
5020 test.setAlarm(twoMs);
5021 }
5022 mockStorage->expectCall("setAlarm", ws)
5023 .withParams(CAPNP(scheduledTimeMs = 2))
5024 .thenReturn(CAPNP());
5025 }
5026 
5027 {
5028 test.setAlarm(oneMs);
5029 
5030 // Test that alarm handler handle does not cache delete when it fails
5031 {
5032 auto armResult = test.cache.armAlarmHandler(oneMs, nullptr, testCurrentTime);
5033 KJ_ASSERT(armResult.is<ActorCache::RunAlarmHandler>());
5034 }
5035 mockStorage->expectCall("deleteAlarm", ws)
5036 .withParams(CAPNP(timeToDeleteMs = 1))
5037 .thenReturn(CAPNP(deleted = false));
5038 test.gate.wait(nullptr).wait(test.ws);
5039 }
5040 
5041 {
5042 // Test that alarm handler handle does not cache alarm delete when noCache == true
5043 {
5044 auto armResult = test.cache.armAlarmHandler(twoMs, nullptr, testCurrentTime, true);
5045 KJ_ASSERT(armResult.is<ActorCache::RunAlarmHandler>());
5046 }
5047 mockStorage->expectCall("deleteAlarm", ws)
5048 .withParams(CAPNP(timeToDeleteMs = 2))
5049 .thenReturn(CAPNP(deleted = true));
5050 test.gate.wait(nullptr).wait(test.ws);
5051 }
5052 
5053 {
5054 auto time = expectUncached(test.getAlarm());
5055 
5056 mockStorage->expectCall("getAlarm", ws).thenReturn(CAPNP(scheduledTimeMs = 0));
5057 
5058 KJ_ASSERT(time.wait(ws) == kj::none);
5059 }
5060}
5061 
5062KJ_TEST("ActorCache uncached nonnull alarm get") {
5063 ActorCacheTest test;
5064 auto& ws = test.ws;
5065 auto& mockStorage = test.mockStorage;
5066 
5067 auto time = expectUncached(test.getAlarm());
5068 auto oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5069 
5070 mockStorage->expectCall("getAlarm", ws).thenReturn(CAPNP(scheduledTimeMs = 1));
5071 
5072 KJ_ASSERT(time.wait(ws) == oneMs);
5073}
5074 
5075KJ_TEST("ActorCache alarm delete when flush fails") {
5076 ActorCacheTest test;
5077 auto& ws = test.ws;
5078 auto& mockStorage = test.mockStorage;
5079 
5080 auto oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5081 auto testCurrentTime = kj::UNIX_EPOCH;
5082 
5083 {
5084 auto time = expectUncached(test.getAlarm());
5085 
5086 mockStorage->expectCall("getAlarm", ws).thenReturn(CAPNP(scheduledTimeMs = 1));
5087 
5088 KJ_ASSERT(time.wait(ws) == oneMs);
5089 }
5090 
5091 {
5092 auto time = KJ_ASSERT_NONNULL(expectCached(test.getAlarm()));
5093 KJ_ASSERT(time == oneMs);
5094 }
5095 
5096 // we want to test that even if a flush is retried
5097 // that the post-delete actions for a checked delete happen.
5098 {
5099 auto handle = test.cache.armAlarmHandler(oneMs, nullptr, testCurrentTime);
5100 
5101 auto time = expectCached(test.getAlarm());
5102 KJ_ASSERT(time == kj::none);
5103 }
5104 
5105 for (auto i = 0; i < 2; i++) {
5106 mockStorage->expectCall("deleteAlarm", ws)
5107 .withParams(CAPNP(timeToDeleteMs = 1))
5108 .thenThrow(KJ_EXCEPTION(DISCONNECTED, "foo"));
5109 }
5110 
5111 {
5112 mockStorage->expectCall("deleteAlarm", ws)
5113 .withParams(CAPNP(timeToDeleteMs = 1))
5114 .thenReturn(CAPNP(deleted = false));
5115 // Wait on the output gate to make sure the flush is actually done.
5116 test.gate.wait(nullptr).wait(test.ws);
5117 }
5118 
5119 {
5120 auto time = expectUncached(test.getAlarm());
5121 
5122 mockStorage->expectCall("getAlarm", ws).thenReturn(CAPNP(scheduledTimeMs = 10));
5123 
5124 KJ_ASSERT(time.wait(ws) == 10 * kj::MILLISECONDS + kj::UNIX_EPOCH);
5125 }
5126}
5127 
5128KJ_TEST("ActorCache deleteAll() with deleteAlarm deletes alarm after deleteAll succeeds") {
5129 // Tests that deleteAll() with deleteAlarm=true deletes the alarm in the post-deleteAll
5130 // flush, ensuring the alarm is only deleted after the deleteAll RPC succeeds.
5131 ActorCacheTest test;
5132 auto& ws = test.ws;
5133 auto& mockStorage = test.mockStorage;
5134 
5135 auto oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5136 
5137 // First, set an alarm so we have something to delete.
5138 test.setAlarm(oneMs);
5139 
5140 mockStorage->expectCall("setAlarm", ws)
5141 .withParams(CAPNP(scheduledTimeMs = 1))
5142 .thenReturn(CAPNP());
5143 
5144 {
5145 auto time = expectCached(test.getAlarm());
5146 KJ_ASSERT(time == oneMs);
5147 }
5148 
5149 // Call deleteAll() with deleteAlarm=true.
5150 auto deleteAll = test.cache.deleteAll({}, nullptr, {.deleteAlarm = true});
5151 
5152 // The deleteAll RPC goes through first.
5153 mockStorage->expectCall("deleteAll", ws).thenReturn(CAPNP(numDeleted = 0));
5154 
5155 KJ_ASSERT(deleteAll.count.wait(ws) == 0);
5156 
5157 // After deleteAll succeeds, the alarm deletion is flushed in the post-deleteAll flush.
5158 mockStorage->expectCall("deleteAlarm", ws)
5159 .withParams(CAPNP(timeToDeleteMs = 0))
5160 .thenReturn(CAPNP(deleted = true));
5161 test.gate.wait(nullptr).wait(test.ws);
5162 
5163 {
5164 auto time = expectCached(test.getAlarm());
5165 KJ_ASSERT(time == kj::none);
5166 }
5167}
5168 
5169KJ_TEST("ActorCache deleteAll() without deleteAlarm preserves alarm") {
5170 // Tests that deleteAll() without deleteAlarm (the default) does not delete the alarm.
5171 ActorCacheTest test;
5172 auto& ws = test.ws;
5173 auto& mockStorage = test.mockStorage;
5174 
5175 auto oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5176 
5177 // First, set an alarm.
5178 test.setAlarm(oneMs);
5179 
5180 mockStorage->expectCall("setAlarm", ws)
5181 .withParams(CAPNP(scheduledTimeMs = 1))
5182 .thenReturn(CAPNP());
5183 
5184 {
5185 auto time = expectCached(test.getAlarm());
5186 KJ_ASSERT(time == oneMs);
5187 }
5188 
5189 // Call deleteAll() without deleteAlarm (default).
5190 auto deleteAll = test.cache.deleteAll({}, nullptr);
5191 
5192 // The deleteAll RPC goes through.
5193 mockStorage->expectCall("deleteAll", ws).thenReturn(CAPNP(numDeleted = 0));
5194 
5195 KJ_ASSERT(deleteAll.count.wait(ws) == 0);
5196 
5197 // Wait for the output gate to complete.
5198 test.gate.wait(nullptr).wait(test.ws);
5199 
5200 // The alarm should still be set.
5201 {
5202 auto time = expectCached(test.getAlarm());
5203 KJ_ASSERT(time == oneMs);
5204 }
5205}
5206 
5207KJ_TEST("ActorCache deleteAll() with deleteAlarm does not overwrite subsequent setAlarm") {
5208 // Tests that if setAlarm() is called after deleteAll({.deleteAlarm = true}) but before the
5209 // deleteAll RPC completes, the new alarm is preserved rather than being clobbered by the
5210 // post-deleteAll alarm deletion.
5211 ActorCacheTest test;
5212 auto& ws = test.ws;
5213 auto& mockStorage = test.mockStorage;
5214 
5215 auto oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5216 auto twoMs = 2 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5217 
5218 // First, set an alarm so we have something to delete.
5219 test.setAlarm(oneMs);
5220 
5221 mockStorage->expectCall("setAlarm", ws)
5222 .withParams(CAPNP(scheduledTimeMs = 1))
5223 .thenReturn(CAPNP());
5224 
5225 {
5226 auto time = expectCached(test.getAlarm());
5227 KJ_ASSERT(time == oneMs);
5228 }
5229 
5230 // Call deleteAll() with deleteAlarm=true.
5231 auto deleteAll = test.cache.deleteAll({}, nullptr, {.deleteAlarm = true});
5232 
5233 // Hold the deleteAll RPC in flight.
5234 auto mockDeleteAll = mockStorage->expectCall("deleteAll", ws);
5235 
5236 // While the deleteAll RPC is in flight, set a new alarm. This simulates the user's JS code
5237 // calling setAlarm() after deleteAll() but before the deleteAll RPC has completed.
5238 test.setAlarm(twoMs);
5239 
5240 // Now complete the deleteAll RPC.
5241 kj::mv(mockDeleteAll).thenReturn(CAPNP(numDeleted = 0));
5242 
5243 KJ_ASSERT(deleteAll.count.wait(ws) == 0);
5244 
5245 // The post-deleteAll flush should send setAlarm for the new time, not deleteAlarm.
5246 mockStorage->expectCall("setAlarm", ws)
5247 .withParams(CAPNP(scheduledTimeMs = 2))
5248 .thenReturn(CAPNP());
5249 test.gate.wait(nullptr).wait(test.ws);
5250 
5251 // Verify the alarm is set to the new time.
5252 {
5253 auto time = expectCached(test.getAlarm());
5254 KJ_ASSERT(time == twoMs);
5255 }
5256}
5257 
5258KJ_TEST("ActorCache deleteAll() with deleteAlarm during alarm handler cancels deferred delete") {
5259 // Tests that calling deleteAll() with deleteAlarm=true while an alarm handler is running
5260 // overwrites the DeferredAlarmDelete state, so when the handler finishes and the
5261 // DeferredAlarmDeleter fires, it does nothing (since currentAlarmTime is no longer
5262 // DeferredAlarmDelete). The alarm deletion instead happens via the normal post-deleteAll
5263 // flush path.
5264 ActorCacheTest test;
5265 auto& ws = test.ws;
5266 auto& mockStorage = test.mockStorage;
5267 
5268 auto oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5269 auto testCurrentTime = kj::UNIX_EPOCH;
5270 
5271 // Set an alarm so we have something to delete.
5272 test.setAlarm(oneMs);
5273 
5274 mockStorage->expectCall("setAlarm", ws)
5275 .withParams(CAPNP(scheduledTimeMs = 1))
5276 .thenReturn(CAPNP());
5277 
5278 {
5279 auto time = expectCached(test.getAlarm());
5280 KJ_ASSERT(time == oneMs);
5281 }
5282 
5283 {
5284 // Arm the alarm handler, putting us in DeferredAlarmDelete state.
5285 auto armResult = test.cache.armAlarmHandler(oneMs, nullptr, testCurrentTime);
5286 KJ_ASSERT(armResult.is<ActorCache::RunAlarmHandler>());
5287 
5288 // During the handler, getAlarm() should return none (deferred delete is active).
5289 auto time = expectCached(test.getAlarm());
5290 KJ_ASSERT(time == kj::none);
5291 
5292 // Call deleteAll() with deleteAlarm=true while the handler is running. This overwrites
5293 // currentAlarmTime from DeferredAlarmDelete to KnownAlarmTime{CLEAN, none}.
5294 auto deleteAll = test.cache.deleteAll({}, nullptr, {.deleteAlarm = true});
5295 
5296 // getAlarm() should still return none.
5297 {
5298 auto time = expectCached(test.getAlarm());
5299 KJ_ASSERT(time == kj::none);
5300 }
5301 
5302 // Drop the RunAlarmHandler. Since currentAlarmTime is no longer DeferredAlarmDelete,
5303 // the DeferredAlarmDeleter does nothing.
5304 }
5305 
5306 // The deleteAll RPC goes through.
5307 mockStorage->expectCall("deleteAll", ws).thenReturn(CAPNP(numDeleted = 0));
5308 
5309 // After deleteAll succeeds, the alarm deletion is flushed in the post-deleteAll flush.
5310 mockStorage->expectCall("deleteAlarm", ws)
5311 .withParams(CAPNP(timeToDeleteMs = 0))
5312 .thenReturn(CAPNP(deleted = true));
5313 test.gate.wait(nullptr).wait(test.ws);
5314 
5315 {
5316 auto time = expectCached(test.getAlarm());
5317 KJ_ASSERT(time == kj::none);
5318 }
5319}
5320 
5321KJ_TEST(
5322 "ActorCache deleteAll() without deleteAlarm during alarm handler preserves deferred delete") {
5323 // Tests that calling deleteAll() without deleteAlarm while an alarm handler is running
5324 // leaves the DeferredAlarmDelete state intact. When the handler finishes, the deferred
5325 // deletion fires normally and the alarm is deleted via the pre-deleteAll flush.
5326 ActorCacheTest test;
5327 auto& ws = test.ws;
5328 auto& mockStorage = test.mockStorage;
5329 
5330 auto oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5331 auto testCurrentTime = kj::UNIX_EPOCH;
5332 
5333 // Set an alarm so we have something to delete.
5334 test.setAlarm(oneMs);
5335 
5336 mockStorage->expectCall("setAlarm", ws)
5337 .withParams(CAPNP(scheduledTimeMs = 1))
5338 .thenReturn(CAPNP());
5339 
5340 {
5341 auto time = expectCached(test.getAlarm());
5342 KJ_ASSERT(time == oneMs);
5343 }
5344 
5345 {
5346 // Arm the alarm handler, putting us in DeferredAlarmDelete state.
5347 auto armResult = test.cache.armAlarmHandler(oneMs, nullptr, testCurrentTime);
5348 KJ_ASSERT(armResult.is<ActorCache::RunAlarmHandler>());
5349 
5350 // During the handler, getAlarm() should return none (deferred delete is active).
5351 auto time = expectCached(test.getAlarm());
5352 KJ_ASSERT(time == kj::none);
5353 
5354 // Call deleteAll() without deleteAlarm while the handler is running.
5355 // This does NOT touch currentAlarmTime, so DeferredAlarmDelete is preserved.
5356 auto deleteAll = test.cache.deleteAll({}, nullptr);
5357 
5358 // getAlarm() still returns none because DeferredAlarmDelete is still active.
5359 {
5360 auto time = expectCached(test.getAlarm());
5361 KJ_ASSERT(time == kj::none);
5362 }
5363 
5364 // Drop the RunAlarmHandler. DeferredAlarmDeleter fires, sets status to READY,
5365 // and calls ensureFlushScheduled().
5366 }
5367 
5368 // The deferred alarm delete is flushed in the pre-deleteAll flush, using the original
5369 // alarm time as the timeToDeleteMs.
5370 mockStorage->expectCall("deleteAlarm", ws)
5371 .withParams(CAPNP(timeToDeleteMs = 1))
5372 .thenReturn(CAPNP(deleted = true));
5373 
5374 // Then the deleteAll RPC goes through.
5375 mockStorage->expectCall("deleteAll", ws).thenReturn(CAPNP(numDeleted = 0));
5376 
5377 test.gate.wait(nullptr).wait(test.ws);
5378 
5379 {
5380 auto time = expectCached(test.getAlarm());
5381 KJ_ASSERT(time == kj::none);
5382 }
5383}
5384 
5385KJ_TEST("ActorCache deleteAll() failure with deleteAlarm does not delete alarm") {
5386 // Tests that if the deleteAll RPC fails, the alarm deletion RPC is never sent.
5387 ActorCacheTest test({.monitorOutputGate = false});
5388 auto& ws = test.ws;
5389 auto& mockStorage = test.mockStorage;
5390 
5391 auto brokenPromise = test.gate.onBroken();
5392 
5393 auto oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5394 
5395 // First, set an alarm so we have something to delete.
5396 test.setAlarm(oneMs);
5397 
5398 mockStorage->expectCall("setAlarm", ws)
5399 .withParams(CAPNP(scheduledTimeMs = 1))
5400 .thenReturn(CAPNP());
5401 
5402 {
5403 auto time = expectCached(test.getAlarm());
5404 KJ_ASSERT(time == oneMs);
5405 }
5406 
5407 // Call deleteAll() with deleteAlarm=true.
5408 auto deleteAll = test.cache.deleteAll({}, nullptr, {.deleteAlarm = true});
5409 
5410 // The deleteAll RPC fails.
5411 mockStorage->expectCall("deleteAll", ws)
5412 .thenThrow(KJ_EXCEPTION(FAILED, "jsg.Error: deleteAll failed"));
5413 
5414 // The output gate should be broken due to the failure.
5415 KJ_EXPECT_THROW_MESSAGE("deleteAll failed", brokenPromise.wait(ws));
5416 
5417 // No deleteAlarm RPC should have been sent since the deleteAll failed.
5418 mockStorage->expectNoActivity(ws);
5419}
5420 
5421KJ_TEST("ActorCache can wait for flush") {
5422 // This test confirms that `onNoPendingFlush()` will return a promise that resolves when any
5423 // scheduled or in-flight flush completes.
5424 
5425 ActorCacheTest test;
5426 auto& ws = test.ws;
5427 auto& mockStorage = test.mockStorage;
5428 
5429 struct InFlightRequest {
5430 workerd::MockServer::ExpectedCall op;
5431 kj::Maybe<kj::Own<workerd::MockServer>> maybeTxn;
5432 };
5433 
5434 // There is no pending flush since nothing has been done!
5435 KJ_ASSERT(test.cache.onNoPendingFlush(nullptr) == nullptr);
5436 
5437 struct VerifyOptions {
5438 bool skipSecondOperation;
5439 };
5440 size_t secondaryPutIndex = 0;
5441 auto verify = [&](auto receiveRequest, auto sendResponse, VerifyOptions options) {
5442 // We haven't sent our request yet, but we should have a promise now.
5443 auto scheduledPromise = KJ_ASSERT_NONNULL(test.cache.onNoPendingFlush(nullptr));
5444 
5445 // We have sent our request, but it hasn't responded yet. We should still have a promise.
5446 auto req = receiveRequest();
5447 auto inFlightPromise = KJ_ASSERT_NONNULL(test.cache.onNoPendingFlush(nullptr));
5448 
5449 // Do an additional put to make a separate flush.
5450 struct SecondOperation {
5451 kj::String key;
5452 kj::Promise<void> scheduledPromise;
5453 };
5454 kj::Maybe<SecondOperation> maybeSecondOperation;
5455 if (!options.skipSecondOperation) {
5456 auto key = kj::str("foo-", secondaryPutIndex++);
5457 test.put(key, "bar");
5458 auto secondPromise = KJ_ASSERT_NONNULL(test.cache.onNoPendingFlush(nullptr));
5459 KJ_ASSERT(!secondPromise.poll(ws));
5460 maybeSecondOperation.emplace(SecondOperation{
5461 .key = kj::mv(key),
5462 .scheduledPromise = kj::mv(secondPromise),
5463 });
5464 }
5465 
5466 // No promise should have resolved yet.
5467 KJ_ASSERT(!scheduledPromise.poll(ws) && !inFlightPromise.poll(ws));
5468 
5469 // Resolve the operations and confirm that the promises resolve.
5470 sendResponse(kj::mv(req));
5471 scheduledPromise.wait(ws);
5472 inFlightPromise.wait(ws);
5473 
5474 KJ_IF_SOME(secondOperation, maybeSecondOperation) {
5475 // This promise is for a later flush, so it should not have resolved yet.
5476 KJ_ASSERT(!secondOperation.scheduledPromise.poll(ws));
5477 
5478 // Finish our secondary put and observe the second flush resolving.
5479 auto params =
5480 kj::str(R"((entries = [(key = ")", secondOperation.key, R"(", value = "bar")]))");
5481 mockStorage->expectCall("put", ws).withParams(params).thenReturn(CAPNP());
5482 
5483 secondOperation.scheduledPromise.wait(ws);
5484 }
5485 
5486 // We finished our flush, nothing left to do.
5487 KJ_ASSERT(test.cache.onNoPendingFlush(nullptr) == nullptr);
5488 };
5489 
5490 {
5491 // Join in on a simple put.
5492 test.put("foo", "bar");
5493 
5494 verify(
5495 [&]() {
5496 return InFlightRequest{
5497 .op = mockStorage->expectCall("put", ws).withParams(
5498 CAPNP(entries = [(key = "foo", value = "bar")])),
5499 };
5500 }, [&](auto req) { kj::mv(req.op).thenReturn(CAPNP()); },
5501 {
5502 .skipSecondOperation = false,
5503 });
5504 }
5505 
5506 {
5507 // Join in on a delete.
5508 test.delete_("foo");
5509 
5510 verify(
5511 [&]() {
5512 return InFlightRequest{
5513 .op = mockStorage->expectCall("delete", ws).withParams(CAPNP(keys = ["foo"])),
5514 };
5515 }, [&](auto req) { kj::mv(req.op).thenReturn(CAPNP(numDeleted = 1)); },
5516 {
5517 .skipSecondOperation = false,
5518 });
5519 }
5520 
5521 {
5522 // Join in on a simple put with allowUnconfirmed.
5523 test.put("foo", "baz", ActorCacheWriteOptions{.allowUnconfirmed = true});
5524 
5525 verify(
5526 [&]() {
5527 return InFlightRequest{
5528 .op = mockStorage->expectCall("put", ws).withParams(
5529 CAPNP(entries = [(key = "foo", value = "baz")])),
5530 };
5531 }, [&](auto req) { kj::mv(req.op).thenReturn(CAPNP()); },
5532 {
5533 .skipSecondOperation = false,
5534 });
5535 }
5536 
5537 {
5538 // Join in on a delete with allowUnconfirmed.
5539 test.delete_("foo", ActorCacheWriteOptions{.allowUnconfirmed = true});
5540 
5541 verify(
5542 [&]() {
5543 return InFlightRequest{
5544 .op = mockStorage->expectCall("delete", ws).withParams(CAPNP(keys = ["foo"])),
5545 };
5546 }, [&](auto req) { kj::mv(req.op).thenReturn(CAPNP(numDeleted = 1)); },
5547 {
5548 .skipSecondOperation = false,
5549 });
5550 }
5551 
5552 {
5553 // Join in on a scheduled setAlarm.
5554 test.setAlarm(1 * kj::MILLISECONDS + kj::UNIX_EPOCH);
5555 
5556 verify(
5557 [&]() {
5558 return InFlightRequest{
5559 .op = mockStorage->expectCall("setAlarm", ws).withParams(CAPNP(scheduledTimeMs = 1)),
5560 };
5561 }, [&](auto req) { kj::mv(req.op).thenReturn(CAPNP()); },
5562 {
5563 .skipSecondOperation = false,
5564 });
5565 }
5566 
5567 {
5568 // Join in on a scheduled setAlarm with allowUnconfirmed.
5569 test.setAlarm(
5570 2 * kj::MILLISECONDS + kj::UNIX_EPOCH, ActorCacheWriteOptions{.allowUnconfirmed = true});
5571 
5572 verify(
5573 [&]() {
5574 return InFlightRequest{
5575 .op = mockStorage->expectCall("setAlarm", ws).withParams(CAPNP(scheduledTimeMs = 2)),
5576 };
5577 }, [&](auto req) { kj::mv(req.op).thenReturn(CAPNP()); },
5578 {
5579 .skipSecondOperation = false,
5580 });
5581 }
5582 
5583 {
5584 // Join in on a scheduled deleteAll.
5585 test.cache.deleteAll(ActorCacheWriteOptions{.allowUnconfirmed = false}, nullptr);
5586 
5587 verify(
5588 [&]() {
5589 return InFlightRequest{
5590 .op = mockStorage->expectCall("deleteAll", ws).withParams(CAPNP()),
5591 };
5592 }, [&](auto req) { kj::mv(req.op).thenReturn(CAPNP()); },
5593 {
5594 // We can't test the second operation because deleteAll immediately follows up with any puts
5595 // that happened while it was in flight. This means that we invoke the mock twice in the same
5596 // promise chain without being able to set up exceptions in time.
5597 .skipSecondOperation = true,
5598 });
5599 }
5600 
5601 {
5602 // Join in on a scheduled deleteAll with allowUnconfirmed.
5603 test.cache.deleteAll(ActorCacheWriteOptions{.allowUnconfirmed = true}, nullptr);
5604 
5605 verify(
5606 [&]() {
5607 return InFlightRequest{
5608 .op = mockStorage->expectCall("deleteAll", ws).withParams(CAPNP()),
5609 };
5610 }, [&](auto req) { kj::mv(req.op).thenReturn(CAPNP()); },
5611 {
5612 .skipSecondOperation = true,
5613 });
5614 }
5615}
5616 
5617KJ_TEST("ActorCache can shutdown") {
5618 // This test confirms that `shutdown()` stops scheduled flushes but does not stop in-flight
5619 // flushes. It also confirms that `shutdown()` prevents future operations.
5620 
5621 struct InFlightRequest {
5622 MockServer::ExpectedCall op;
5623 kj::Promise<void> promise;
5624 };
5625 
5626 struct BeforeShutdownResult {
5627 kj::Maybe<InFlightRequest> maybeReq;
5628 bool shouldBreakOutputGate;
5629 };
5630 
5631 struct VerifyOptions {
5632 kj::Maybe<const kj::Exception&> maybeError;
5633 };
5634 auto verifyWithOptions = [&](auto&& beforeShutdown, auto&& afterShutdown, VerifyOptions options) {
5635 auto test = ActorCacheTest({.monitorOutputGate = false});
5636 auto& ws = test.ws;
5637 
5638 BeforeShutdownResult res = beforeShutdown(test);
5639 
5640 // Shutdown and observe the pending flush to break the io gate.
5641 test.cache.shutdown(options.maybeError);
5642 auto maybeShutdownPromise = test.cache.onNoPendingFlush(nullptr);
5643 
5644 afterShutdown(test, kj::mv(res.maybeReq));
5645 
5646 auto error = options.maybeError.map([](const kj::Exception& e) {
5647 return e.clone();
5648 }).orDefault(KJ_EXCEPTION(DISCONNECTED, kj::str(ActorCache::SHUTDOWN_ERROR_MESSAGE)));
5649 
5650 if (res.shouldBreakOutputGate) {
5651 // We expected the output gate to break async after shutdown.
5652 auto& shutdownPromise = KJ_REQUIRE_NONNULL(maybeShutdownPromise);
5653 WD_EXPECT_THROW(error, shutdownPromise.wait(ws));
5654 KJ_EXPECT(test.cache.onNoPendingFlush(nullptr) == nullptr);
5655 WD_EXPECT_THROW(error, test.gate.wait(nullptr).wait(ws));
5656 } else KJ_IF_SOME(promise, maybeShutdownPromise) {
5657 // The in-flight flush should resolve cleanly without any follow on or breaking the output
5658 // gate.
5659 promise.wait(ws);
5660 KJ_EXPECT(test.cache.onNoPendingFlush(nullptr) == nullptr);
5661 test.gate.wait(nullptr).wait(ws);
5662 }
5663 
5664 // Puts and deletes, even with allowedUnconfirmed, should throw.
5665 WD_EXPECT_THROW(error, test.put("foo", "baz"));
5666 WD_EXPECT_THROW(error, test.put("foo", "bat", {.allowUnconfirmed = true}));
5667 WD_EXPECT_THROW(error, test.delete_("foo"));
5668 WD_EXPECT_THROW(error, test.delete_("foo", {.allowUnconfirmed = true}));
5669 
5670 if (!res.shouldBreakOutputGate) {
5671 // We tried to use storage after shutdown, we should now be breaking the output gate.
5672 auto afterShutdownPromise = KJ_ASSERT_NONNULL(test.cache.onNoPendingFlush(nullptr));
5673 WD_EXPECT_THROW(error, afterShutdownPromise.wait(ws));
5674 KJ_EXPECT(test.cache.onNoPendingFlush(nullptr) == nullptr);
5675 WD_EXPECT_THROW(error, test.gate.wait(nullptr).wait(ws));
5676 }
5677 };
5678 
5679 auto verify = [&](auto&& beforeShutdown, auto&& afterShutdown) {
5680 verifyWithOptions(beforeShutdown, afterShutdown, {.maybeError = kj::none});
5681 verifyWithOptions(beforeShutdown, afterShutdown, {.maybeError = KJ_EXCEPTION(FAILED, "Nope.")});
5682 };
5683 
5684 verify([](ActorCacheTest& test) {
5685 // Do nothing and expect nothing!
5686 return BeforeShutdownResult{
5687 .maybeReq = kj::none,
5688 .shouldBreakOutputGate = false,
5689 };
5690 }, [](ActorCacheTest& test, kj::Maybe<InFlightRequest>) {
5691 // Nothing should have made it to storage.
5692 test.mockStorage->expectNoActivity(test.ws);
5693 });
5694 
5695 verify([](ActorCacheTest& test) {
5696 // Do a confirmed put (which schedules a flush).
5697 test.put("foo", "bar", {.allowUnconfirmed = false});
5698 
5699 // Expect the put to be cancelled and break the gate.
5700 return BeforeShutdownResult{
5701 .maybeReq = kj::none,
5702 .shouldBreakOutputGate = true,
5703 };
5704 }, [](ActorCacheTest& test, kj::Maybe<InFlightRequest>) {
5705 // Nothing should have made it to storage.
5706 test.mockStorage->expectNoActivity(test.ws);
5707 });
5708 
5709 verify([](ActorCacheTest& test) {
5710 // Do an unconfirmed put (which schedules a flush).
5711 test.put("foo", "bar", {.allowUnconfirmed = true});
5712 
5713 // Expect the put to be cancelled and break the gate.
5714 return BeforeShutdownResult{
5715 .maybeReq = kj::none,
5716 .shouldBreakOutputGate = true,
5717 };
5718 }, [](ActorCacheTest& test, kj::Maybe<InFlightRequest>) {
5719 // Nothing should have made it to storage.
5720 test.mockStorage->expectNoActivity(test.ws);
5721 });
5722 
5723 verify([](ActorCacheTest& test) {
5724 // Do a confirmed put and wait for it to be in-flight.
5725 test.put("foo", "bar", {.allowUnconfirmed = false});
5726 
5727 auto op = test.mockStorage->expectCall("put", test.ws)
5728 .withParams(CAPNP(entries = [(key = "foo", value = "bar")]));
5729 auto promise = KJ_REQUIRE_NONNULL(test.cache.onNoPendingFlush(nullptr));
5730 KJ_EXPECT(!promise.poll(test.ws));
5731 
5732 return BeforeShutdownResult{
5733 .maybeReq =
5734 InFlightRequest{
5735 .op = kj::mv(op),
5736 .promise = kj::mv(promise),
5737 },
5738 .shouldBreakOutputGate = false,
5739 };
5740 }, [](ActorCacheTest& test, kj::Maybe<InFlightRequest> maybeReq) {
5741 // Finish the storage response and wait to see our pre-shutdown in-flight flush finish.
5742 auto req = KJ_ASSERT_NONNULL(kj::mv(maybeReq));
5743 kj::mv(req.op).thenReturn(CAPNP());
5744 req.promise.wait(test.ws);
5745 
5746 // Nothing else should have made it to storage.
5747 test.mockStorage->expectNoActivity(test.ws);
5748 });
5749 
5750 verify([](ActorCacheTest& test) {
5751 // Do an unconfirmed put and wait for it to be in-flight.
5752 test.put("foo", "bar", {.allowUnconfirmed = true});
5753 
5754 auto op = test.mockStorage->expectCall("put", test.ws)
5755 .withParams(CAPNP(entries = [(key = "foo", value = "bar")]));
5756 auto promise = KJ_REQUIRE_NONNULL(test.cache.onNoPendingFlush(nullptr));
5757 KJ_EXPECT(!promise.poll(test.ws));
5758 
5759 return BeforeShutdownResult{
5760 .maybeReq =
5761 InFlightRequest{
5762 .op = kj::mv(op),
5763 .promise = kj::mv(promise),
5764 },
5765 .shouldBreakOutputGate = false,
5766 };
5767 }, [](ActorCacheTest& test, kj::Maybe<InFlightRequest> maybeReq) {
5768 // Finish the storage response and wait to see our pre-shutdown in-flight flush finish.
5769 auto req = KJ_ASSERT_NONNULL(kj::mv(maybeReq));
5770 kj::mv(req.op).thenReturn(CAPNP());
5771 req.promise.wait(test.ws);
5772 
5773 // Nothing else should have made it to storage.
5774 test.mockStorage->expectNoActivity(test.ws);
5775 });
5776}
5777 
5778KJ_TEST("ActorCache alarm cleared by abandonAlarm") {
5779 // After the alarm scheduler calls abandonAlarm(), the cache correctly forgets the alarm.
5780 
5781 ActorCacheTest test;
5782 auto& ws = test.ws;
5783 auto& mockStorage = test.mockStorage;
5784 
5785 auto oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5786 
5787 test.setAlarm(oneMs);
5788 mockStorage->expectCall("setAlarm", ws)
5789 .withParams(CAPNP(scheduledTimeMs = 1))
5790 .thenReturn(CAPNP());
5791 // Drive the event loop so the flush completes and transitions DIRTY -> CLEAN.
5792 ws.poll();
5793 
5794 // abandonAlarm() resets currentAlarmTime to UnknownAlarmTime so the next getAlarm()
5795 // refetches from storage rather than serving a stale cached value.
5796 auto result = test.cache.abandonAlarm(oneMs).wait(ws);
5797 
5798 // Returns kj::none: alarm was cleared, AlarmManager should not re-register.
5799 KJ_ASSERT(result == kj::none);
5800 
5801 // getAlarm() now triggers a storage read (cache was reset to unknown).
5802 auto promise = expectUncached(test.getAlarm());
5803 mockStorage->expectCall("getAlarm", ws).thenReturn(CAPNP(scheduledTimeMs = 0));
5804 auto time = promise.wait(ws);
5805 KJ_ASSERT(time == kj::none);
5806}
5807 
5808KJ_TEST("ActorCache alarm preserved after ALARM_RETRY_MAX_TRIES uncounted (internal) failures") {
5809 // When all ALARM_RETRY_MAX_TRIES failures are uncounted (retryCountsAgainstLimit=false,
5810 // i.e. infrastructure errors), the alarm scheduler's countedRetry never reaches the limit and
5811 // abandonAlarm is NEVER called. The alarm must remain set throughout so that the scheduler
5812 // can keep retrying indefinitely until the infrastructure issue resolves.
5813 
5814 ActorCacheTest test;
5815 auto& ws = test.ws;
5816 auto& mockStorage = test.mockStorage;
5817 
5818 auto oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5819 auto testCurrentTime = kj::UNIX_EPOCH;
5820 
5821 test.setAlarm(oneMs);
5822 mockStorage->expectCall("setAlarm", ws)
5823 .withParams(CAPNP(scheduledTimeMs = 1))
5824 .thenReturn(CAPNP());
5825 
5826 // Simulate uncounted failures well past ALARM_RETRY_MAX_TRIES (= 6).
5827 // countedRetry stays at 0; AlarmManager never gives up; abandonAlarm is never called.
5828 // We've seen alarms fail hundreds of times due to infrastructure errors in production,
5829 // so we check both at the boundary (6) and well beyond it (100).
5830 for (auto i = 0; i < 100; i++) {
5831 auto armResult = test.cache.armAlarmHandler(oneMs, nullptr, testCurrentTime);
5832 KJ_ASSERT(armResult.is<ActorCache::RunAlarmHandler>());
5833 test.cache.cancelDeferredAlarmDeletion();
5834 
5835 // Check at the ALARM_RETRY_MAX_TRIES boundary and at the end.
5836 if (i == 5 || i == 99) {
5837 auto time = expectCached(test.getAlarm());
5838 KJ_ASSERT(time == oneMs);
5839 }
5840 }
5841}
5842 
5843KJ_TEST("ActorCache abandonAlarm is a no-op when a newer alarm has replaced the abandoned one") {
5844 // If the user sets a new alarm between the last retry failure and the abandonAlarm() call,
5845 // and that new alarm has already flushed to CLEAN, abandonAlarm() must compare the time
5846 // and leave the new alarm untouched.
5847 
5848 ActorCacheTest test;
5849 auto& ws = test.ws;
5850 auto& mockStorage = test.mockStorage;
5851 
5852 auto oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5853 auto twoMs = 2 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5854 
5855 // Set the original alarm and flush it to storage.
5856 test.setAlarm(oneMs);
5857 mockStorage->expectCall("setAlarm", ws)
5858 .withParams(CAPNP(scheduledTimeMs = 1))
5859 .thenReturn(CAPNP());
5860 
5861 // User sets a new alarm (twoMs). It flushes to CLEAN, leaving KnownAlarmTime{CLEAN, twoMs}.
5862 test.setAlarm(twoMs);
5863 mockStorage->expectCall("setAlarm", ws)
5864 .withParams(CAPNP(scheduledTimeMs = 2))
5865 .thenReturn(CAPNP());
5866 // Advance the event loop to process the storage response and complete the FLUSHING→CLEAN
5867 // transition. Without this poll, the state is still FLUSHING when abandonAlarm runs, and
5868 // the existing status check would protect it by accident, hiding the time-check regression.
5869 ws.poll();
5870 
5871 // abandonAlarm() for the original oneMs alarm must be a no-op: storedTime (twoMs) !=
5872 // scheduledTime (oneMs), so the time check prevents clearing the new alarm.
5873 // Returns twoMs so AlarmManager can re-register the actor's real alarm.
5874 auto result = test.cache.abandonAlarm(oneMs).wait(ws);
5875 
5876 KJ_ASSERT(KJ_ASSERT_NONNULL(result) == twoMs);
5877 
5878 // getAlarm() must still return twoMs -- the new alarm was NOT incorrectly cleared.
5879 auto time = expectCached(test.getAlarm());
5880 KJ_ASSERT(time == twoMs);
5881}
5882 
5883KJ_TEST("ActorCache abandonAlarm returns kj::none when no alarm is stored") {
5884 ActorCacheTest test;
5885 auto& ws = test.ws;
5886 
5887 // No alarm ever set. abandonAlarm should return kj::none.
5888 auto result = test.cache.abandonAlarm(1 * kj::MILLISECONDS + kj::UNIX_EPOCH).wait(ws);
5889 KJ_ASSERT(result == kj::none);
5890}
5891 
5892KJ_TEST("ActorCache abandonAlarm returns kj::none when alarm is DIRTY, not the uncommitted time") {
5893 // If the user set a new alarm T2 that hasn't flushed to storage yet (DIRTY state),
5894 // abandonAlarm must NOT return T2 to AlarmManager for re-registration.
5895 // Returning an uncommitted time would let AlarmManager fire an alarm that was never
5896 // acknowledged by storage, potentially racing with the flush. Returns kj::none instead.
5897 
5898 ActorCacheTest test;
5899 auto& ws = test.ws;
5900 auto& mockStorage = test.mockStorage;
5901 
5902 auto oneMs = 1 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5903 auto twoMs = 2 * kj::MILLISECONDS + kj::UNIX_EPOCH;
5904 
5905 // Flush T0 to storage.
5906 test.setAlarm(oneMs);
5907 mockStorage->expectCall("setAlarm", ws)
5908 .withParams(CAPNP(scheduledTimeMs = 1))
5909 .thenReturn(CAPNP());
5910 
5911 // User sets T2. State becomes DIRTY (flush task queued but not yet delivered to storage).
5912 test.setAlarm(twoMs);
5913 
5914 // Set up expectation for the T2 flush before driving the event loop.
5915 mockStorage->expectCall("setAlarm", ws)
5916 .withParams(CAPNP(scheduledTimeMs = 2))
5917 .thenReturn(CAPNP());
5918 
5919 // abandonAlarm(T0) arrives while T2 is DIRTY. The CLEAN status guard prevents returning
5920 // the uncommitted T2 time -- AlarmManager must not re-register based on it.
5921 auto result = test.cache.abandonAlarm(oneMs).wait(ws);
5922 KJ_ASSERT(result == kj::none);
5923 
5924 // Let the T2 flush complete normally.
5925 ws.poll();
5926 
5927 // T2 is still in cache -- the DIRTY alarm was preserved, not cleared.
5928 auto time = expectCached(test.getAlarm());
5929 KJ_ASSERT(time == twoMs);
5930}
5931 
5932} // namespace
5933} // namespace workerd