Skip to content
File

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

cpp1169 lines
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#pragma once
6 
7#include <workerd/io/actor-storage.capnp.h>
8#include <workerd/io/trace.h>
9#include <workerd/jsg/exception.h>
10#include <workerd/util/strong-bool.h>
11 
12#include <kj/async.h>
13#include <kj/debug.h>
14#include <kj/list.h>
15#include <kj/map.h>
16#include <kj/mutex.h>
17#include <kj/one-of.h>
18#include <kj/time.h>
19 
20#include <atomic>
21 
22namespace workerd {
23 
24using kj::byte;
25using kj::uint;
26class OutputGate;
27class SqliteDatabase;
28class SqliteKv;
29 
30WD_STRONG_BOOL(ReadReplicationIsEnabled);
31 
32struct ActorCacheReadOptions {
33 // If the entry is not already in cache and has to be read from disk, don't store the result in
34 // cache, only return it to the caller.
35 //
36 // If there is already a matching entry in cache, that value will be returned as normal. Hence,
37 // `noCache` does not affect consistency, only performance.
38 bool noCache = false;
39};
40 
41struct ActorCacheWriteOptions {
42 // Instructs that the output gate should not wait for this write to be confirmed on disk. Write
43 // failures will still break the output gate -- but the application could potentially return a
44 // result before the failure is observed, leading to a prematurely confirmed write.
45 bool allowUnconfirmed = false;
46 
47 // Once the value has been confirmed written to disk, immediately evict it from the cache.
48 //
49 // Until the value is safely on disk, the dirty value will be used to fulfill reads for the same
50 // key. Hence, `noCache` does not affect consistency, only performance.
51 bool noCache = false;
52};
53 
54struct DeleteAllOptions {
55 // When true, deleteAll() will also delete any scheduled alarm. The alarm deletion is
56 // guaranteed to take effect only after the deleteAll() itself succeeds, so that we never
57 // end up in a state where the alarm is deleted but KV data remains.
58 bool deleteAlarm = false;
59};
60 
61// Common interface between ActorCache and ActorCache::Transaction.
62class ActorCacheOps {
63 public:
64 using Key = kj::String;
65 using KeyPtr = kj::StringPtr;
66 // Keys are text for now, but we could also change this to `Array<const byte>`.
67 static inline Key cloneKey(KeyPtr ptr) {
68 return kj::str(ptr);
69 }
70 
71 // Values are raw bytes.
72 using Value = kj::Array<const byte>;
73 using ValuePtr = kj::ArrayPtr<const byte>;
74 
75 struct KeyValuePair {
76 Key key;
77 Value value;
78 };
79 struct KeyValuePtrPair {
80 KeyPtr key;
81 ValuePtr value;
82 
83 KeyValuePtrPair(const KeyValuePair& other): key(other.key), value(other.value) {}
84 KeyValuePtrPair(KeyPtr key, ValuePtr value): key(key), value(value) {}
85 };
86 
87 struct KeyRename {
88 Key oldKey;
89 Key newKey;
90 };
91 
92 enum class CacheStatus { CACHED, UNCACHED };
93 
94 struct KeyValuePtrPairWithCache: public KeyValuePtrPair {
95 CacheStatus status;
96 
97 KeyValuePtrPairWithCache(const KeyValuePtrPair& other, CacheStatus status)
98 : KeyValuePtrPair(other),
99 status(status) {}
100 KeyValuePtrPairWithCache(const KeyValuePtrPairWithCache& other)
101 : KeyValuePtrPair(other.key, other.value),
102 status(other.status) {}
103 KeyValuePtrPairWithCache(KeyPtr key, ValuePtr value, CacheStatus status)
104 : KeyValuePtrPair(key, value),
105 status(status) {}
106 };
107 
108 // An iterable type where each element is a KeyValuePtrPair.
109 class GetResultList;
110 
111 struct CleanAlarm {};
112 
113 struct DirtyAlarm {
114 kj::Maybe<kj::Date> newTime;
115 };
116 
117 using MaybeAlarmChange = kj::OneOf<CleanAlarm, DirtyAlarm>;
118 
119 struct DirtyAlarmWithOptions: public DirtyAlarm {
120 ActorCacheWriteOptions options;
121 };
122 
123 using ReadOptions = ActorCacheReadOptions;
124 
125 // Get the values for some key, keys, or range of keys.
126 //
127 // Returns a Maybe<Value> or GetResultList if the result is immediately available from cache,
128 // otherwise returns a Promise and fetches the results from storage.
129 //
130 // `listReverse()` lists in reverse, which turns out to require a subtly different implementation
131 // of pretty much the entire algorithm.
132 //
133 // Passing a null key for `end` means list through the last key in the actor. For `begin`, you
134 // can pass an empty string to list from the first key in the actor, since the empty string is
135 // the first possible key.
136 virtual kj::OneOf<kj::Maybe<Value>, kj::Promise<kj::Maybe<Value>>> get(
137 Key key, ReadOptions options) = 0;
138 virtual kj::OneOf<GetResultList, kj::Promise<GetResultList>> get(
139 kj::Array<Key> keys, ReadOptions options) = 0;
140 virtual kj::OneOf<kj::Maybe<kj::Date>, kj::Promise<kj::Maybe<kj::Date>>> getAlarm(
141 ReadOptions options) = 0;
142 virtual kj::OneOf<GetResultList, kj::Promise<GetResultList>> list(
143 Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) = 0;
144 virtual kj::OneOf<GetResultList, kj::Promise<GetResultList>> listReverse(
145 Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) = 0;
146 
147 using WriteOptions = ActorCacheWriteOptions;
148 
149 // Writes a key/value into cache and schedules it to be flushed to disk later.
150 //
151 // The cache will automatically arrange to flush changes to disk, adding the flush to the
152 // `OutputGate` passed to the constructor.
153 //
154 // This returns a promise for backpressure. If a promise is returned, the application should
155 // delay further puts until the promise resolves. This happens when too much data is pinned in
156 // cache because writes haven't been flushed to disk yet. Dropping this promise will not cancel
157 // the put.
158 //
159 // The traceSpan parameter is used for output gate lock hold tracing. It is captured only by
160 // the first write that starts a new flush batch.
161 virtual kj::Maybe<kj::Promise<void>> put(
162 Key key, Value value, WriteOptions options, SpanParent traceSpan) = 0;
163 virtual kj::Maybe<kj::Promise<void>> put(
164 kj::Array<KeyValuePair> pairs, WriteOptions options, SpanParent traceSpan) = 0;
165 
166 // Writes a new alarm time into cache and schedules it to be flushed to disk later, same as put().
167 virtual kj::Maybe<kj::Promise<void>> setAlarm(
168 kj::Maybe<kj::Date> newTime, WriteOptions options, SpanParent traceSpan) = 0;
169 
170 // Delete the given keys.
171 //
172 // Returns a `bool` or `uint` if it can be immediately determined from cache how many keys were
173 // present before the call. Otherwise, returns a promise which resolves after getting a response
174 // from underlying storage. The promise also applies backpressure if needed, as with put().
175 //
176 // The traceSpan parameter is used for output gate lock hold tracing.
177 virtual kj::OneOf<bool, kj::Promise<bool>> delete_(
178 Key key, WriteOptions options, SpanParent traceSpan) = 0;
179 virtual kj::OneOf<uint, kj::Promise<uint>> delete_(
180 kj::Array<Key> keys, WriteOptions options, SpanParent traceSpan) = 0;
181};
182 
183// Abstract interface that is implemented by ActorCache as well as ActorSqlite.
184//
185// This extends ActorCacheOps and adds some methods that don't make sense as part of
186// ActorCache::Transaction.
187class ActorCacheInterface: public ActorCacheOps {
188 public:
189 // If the actor's storage is backed by SQLite, return the underlying database.
190 virtual kj::Maybe<SqliteDatabase&> getSqliteDatabase() = 0;
191 
192 // If the actor's storage is backed by SQLite, return the SqliteKv object which provides a
193 // synchronous interface to KV storage. This is only available for SQLite-basked DOs because
194 // old-style DOs have asyncronous storage.
195 virtual kj::Maybe<SqliteKv&> getSqliteKv() = 0;
196 
197 class Transaction: public ActorCacheOps {
198 public:
199 // Write all changes to the underlying ActorCache.
200 //
201 // If commit() is not called before the Transaction is destroyed, nothing is written.
202 //
203 // Returns a promise if backpressure needs to be applied (like ActorCache::put()).
204 //
205 // This will NOT detect conflicts, it will always just write blindly, because conflicts
206 // inherently cannot happen.
207 virtual kj::Maybe<kj::Promise<void>> commit() = 0;
208 
209 virtual kj::Promise<void> rollback() = 0;
210 };
211 
212 virtual kj::Own<Transaction> startTransaction() = 0;
213 
214 // We split these up so client code that doesn't need the count doesn't have to
215 // wait for it just to account for backpressure
216 struct DeleteAllResults {
217 kj::Maybe<kj::Promise<void>> backpressure;
218 kj::Promise<uint> count;
219 };
220 
221 // Delete everything in the actor's storage. This is not part of ActorCacheOps because it
222 // is not supported as part of a transaction.
223 //
224 // The returned count only includes keys that were actually deleted from storage, not keys in
225 // cache -- we only use the returned deleteAll count for billing, and not counting deletes of
226 // entries that are only in cache is no problem for billing, those deletes don't cost us anything.
227 virtual DeleteAllResults deleteAll(
228 WriteOptions options, SpanParent traceSpan, DeleteAllOptions deleteAllOptions = {}) = 0;
229 
230 // Call each time the isolate lock is taken to evict stale entries. If this returns a promise,
231 // then the caller must hold off on JavaScript execution until the promise resolves -- this
232 // creates back pressure when the write queue is too deep.
233 //
234 // (This takes a Date rather than a TimePoint because it is based on Date.now(), to avoid
235 // bypassing Spectre mitigations.)
236 virtual kj::Maybe<kj::Promise<void>> evictStale(kj::Date now) = 0;
237 
238 virtual void shutdown(kj::Maybe<const kj::Exception&> maybeException) = 0;
239 
240 // Possible armAlarmHandler() return values:
241 //
242 // Alarm should be canceled without retry (because alarm state has changed such that the
243 // requested alarm time is no longer valid).
244 struct CancelAlarmHandler {
245 // Caller should wait for this promise to complete before canceling.
246 kj::Promise<void> waitBeforeCancel;
247 };
248 // Alarm should be run.
249 struct RunAlarmHandler {
250 // RAII object to delete the alarm, if object is destroyed before setAlarm() or
251 // cancelDeferredAlarmDeletion() are called. Caller should attach it to a promise
252 // representing the alarm handler's execution.
253 kj::Own<void> deferredDelete;
254 };
255 
256 // Call when entering the alarm handler.
257 //
258 // `currentTime` is used to determine if an overdue alarm should run immediately even when
259 // the local alarm state differs from the scheduled time (to avoid blocking on storage sync).
260 virtual kj::OneOf<CancelAlarmHandler, RunAlarmHandler> armAlarmHandler(kj::Date scheduledTime,
261 SpanParent parentSpan,
262 kj::Date currentTime,
263 bool noCache = false,
264 kj::StringPtr actorId = "") = 0;
265 
266 virtual void cancelDeferredAlarmDeletion() = 0;
267 
268 // Called by AlarmManager when it has given up retrying an alarm after too many counted failures.
269 // Implementations should clear the alarm from their local state so getAlarm() reflects the
270 // deletion. Returns the stored alarm time if it differs from scheduledTime (the user set a new
271 // alarm), or kj::none if the alarm was cleared or no alarm was stored.
272 virtual kj::Promise<kj::Maybe<kj::Date>> abandonAlarm(kj::Date scheduledTime) {
273 return kj::Maybe<kj::Date>(kj::none);
274 }
275 
276 virtual kj::Maybe<kj::Promise<void>> onNoPendingFlush(SpanParent parentSpan) = 0;
277 
278 // Implements the respective PITR API calls. The default implementations throw JSG errors saying
279 // PITR is not implemented. These methods are meant to be implemented internally.
280 virtual kj::Promise<kj::String> getCurrentBookmark(SpanParent parentSpan) {
281 JSG_FAIL_REQUIRE(
282 Error, "This Durable Object's storage back-end does not implement point-in-time recovery.");
283 }
284 
285 virtual kj::Promise<kj::String> getBookmarkForTime(kj::Date timestamp) {
286 JSG_FAIL_REQUIRE(
287 Error, "This Durable Object's storage back-end does not implement point-in-time recovery.");
288 }
289 
290 virtual kj::Promise<kj::String> onNextSessionRestoreBookmark(kj::StringPtr bookmark) {
291 JSG_FAIL_REQUIRE(
292 Error, "This Durable Object's storage back-end does not implement point-in-time recovery.");
293 }
294 
295 virtual kj::Promise<void> waitForBookmark(kj::StringPtr bookmark, SpanParent parentSpan) {
296 JSG_FAIL_REQUIRE(
297 Error, "This Durable Object's storage back-end does not implement point-in-time recovery.");
298 }
299 
300 virtual void ensureReplicas() {
301 JSG_FAIL_REQUIRE(Error, "This Durable Object's storage back-end does not support replication.");
302 }
303 
304 virtual void disableReplicas() {
305 JSG_FAIL_REQUIRE(Error, "This Durable Object's storage back-end does not support replication.");
306 }
307 
308 virtual kj::Promise<void> configureReadReplication(ReadReplicationIsEnabled) {
309 JSG_FAIL_REQUIRE(Error, "This Durable Object's storage back-end does not support replication.");
310 }
311};
312 
313// An in-memory caching layer on top of ActorStorage.Stage RPC interface.
314//
315// This cache assumes that it is the only client of the underlying storage -- which is, of
316// course, true for actors.
317//
318// Writes complete "instantly" -- but the OutputGate is told to block output until the write is
319// confirmed durable.
320//
321// Ordering is carefully preserved. A read will always return results consistent with the time
322// when it was called, never reflecting later writes -- even writes that are performed before
323// the read actually completes. Writes are never committed out-of-order (this is accomplished by
324// brute force -- the cache always performs a transaction committing all dirty keys at once).
325//
326// The cache implements LRU eviction triggered by both time and memory pressure. Memory usage is
327// accounted across many actors (typically, all actors in the same isolate), so that the cache
328// size limit can be set based on the per-isolate memory limit.
329class ActorCache final: public ActorCacheInterface {
330 public:
331 // Shared LRU for a whole isolate.
332 class SharedLru;
333 
334 // Hooks that can be used to customize ActorCache behavior or report statistics.
335 class Hooks {
336 public:
337 // Called when the alarm time is dirty when neverFlush is set and ensureFlushScheduled is called.
338 virtual void updateAlarmInMemory(kj::Maybe<kj::Date> newAlarmTime) {};
339 
340 // Used to track metrics of read and write operation latencies from the isolate's perspective.
341 virtual void storageReadCompleted(kj::Duration latency) {}
342 virtual void storageWriteCompleted(kj::Duration latency) {}
343 
344 static const Hooks DEFAULT;
345 };
346 
347 static constexpr auto SHUTDOWN_ERROR_MESSAGE =
348 "broken.ignored; jsg.Error: "
349 "Durable Object storage is no longer accessible."_kj;
350 
351 ActorCache(rpc::ActorStorage::Stage::Client storage,
352 const SharedLru& lru,
353 OutputGate& gate,
354 // Hooks has no member variables, so const_cast is acceptable.
355 Hooks& hooks = const_cast<Hooks&>(Hooks::DEFAULT));
356 ~ActorCache() noexcept(false);
357 
358 kj::Maybe<SqliteDatabase&> getSqliteDatabase() override {
359 return kj::none;
360 }
361 kj::Maybe<SqliteKv&> getSqliteKv() override {
362 return kj::none;
363 }
364 kj::OneOf<kj::Maybe<Value>, kj::Promise<kj::Maybe<Value>>> get(
365 Key key, ReadOptions options) override;
366 kj::OneOf<GetResultList, kj::Promise<GetResultList>> get(
367 kj::Array<Key> keys, ReadOptions options) override;
368 kj::OneOf<kj::Maybe<kj::Date>, kj::Promise<kj::Maybe<kj::Date>>> getAlarm(
369 ReadOptions options) override;
370 kj::OneOf<GetResultList, kj::Promise<GetResultList>> list(
371 Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) override;
372 kj::OneOf<GetResultList, kj::Promise<GetResultList>> listReverse(
373 Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) override;
374 kj::Maybe<kj::Promise<void>> put(
375 Key key, Value value, WriteOptions options, SpanParent traceSpan) override;
376 kj::Maybe<kj::Promise<void>> put(
377 kj::Array<KeyValuePair> pairs, WriteOptions options, SpanParent traceSpan) override;
378 kj::OneOf<bool, kj::Promise<bool>> delete_(
379 Key key, WriteOptions options, SpanParent traceSpan) override;
380 kj::OneOf<uint, kj::Promise<uint>> delete_(
381 kj::Array<Key> keys, WriteOptions options, SpanParent traceSpan) override;
382 kj::Maybe<kj::Promise<void>> setAlarm(
383 kj::Maybe<kj::Date> newAlarmTime, WriteOptions options, SpanParent traceSpan) override;
384 // See ActorCacheOps.
385 
386 kj::Own<ActorCacheInterface::Transaction> startTransaction() override;
387 DeleteAllResults deleteAll(
388 WriteOptions options, SpanParent traceSpan, DeleteAllOptions deleteAllOptions = {}) override;
389 kj::Maybe<kj::Promise<void>> evictStale(kj::Date now) override;
390 void shutdown(kj::Maybe<const kj::Exception&> maybeException) override;
391 
392 kj::OneOf<CancelAlarmHandler, RunAlarmHandler> armAlarmHandler(kj::Date scheduledTime,
393 SpanParent parentSpan,
394 kj::Date currentTime,
395 bool noCache = false,
396 kj::StringPtr actorId = "") override;
397 void cancelDeferredAlarmDeletion() override;
398 kj::Promise<kj::Maybe<kj::Date>> abandonAlarm(kj::Date scheduledTime) override;
399 kj::Maybe<kj::Promise<void>> onNoPendingFlush(SpanParent parentSpan) override;
400 // See ActorCacheInterface
401 
402 class Transaction;
403 // Check for inconsistencies in the cache, e.g. redundant entries.
404 void verifyConsistencyForTest();
405 
406 private:
407 // Backs the `kj::Own<void>` returned by `armAlarmHandler()`.
408 class DeferredAlarmDeleter: public kj::Disposer {
409 public:
410 // The `Own<void>` returned by `armAlarmHandler()` is actually set up to point to the
411 // `ActorCache` itself, but with an alternate disposer that deletes the alarm rather than
412 // the whole object.
413 void disposeImpl(void* pointer) const override {
414 auto p = reinterpret_cast<ActorCache*>(pointer);
415 KJ_IF_SOME(d, p->currentAlarmTime.tryGet<DeferredAlarmDelete>()) {
416 d.status = DeferredAlarmDelete::Status::READY;
417 p->ensureFlushScheduled(WriteOptions{.noCache = d.noCache}, kj::mv(d.traceSpan));
418 }
419 }
420 };
421 
422 enum class EntrySyncStatus : int8_t {
423 // The value was set by the app via put() or delete(), and we have not yet initiated a write
424 // to disk. The entry is appended to `dirtyList` whenever entering this state.
425 //
426 // Next state: CLEAN (if the flush succeeds) or NOT_IN_CACHE (if a new put()/delete()
427 // overwrites this entry first).
428 DIRTY,
429 
430 // The entry matches what is currently on disk. The entry is currently present in the LRU
431 // queue.
432 //
433 // Next state: NOT_IN_CACHE (if a new put()/delete() overwrites the entry), or deleted (if
434 // evicted due to memory pressure).
435 CLEAN,
436 
437 // This entry is not currently in the cache -- it is an orphaned object. This happens e.g. if
438 // a put() or delete() overwrote the entry, in which case the `Entry` object is removed from
439 // the map and replaced with a new object. The old object may continue to exist if it is still
440 // the subject of an outstanding get() which was initiated before the entry was overwritten.
441 // This is not the only use of NOT_IN_CACHE, but in general, any Entry which is not in the
442 // cache's `currentValues` map must have this state.
443 //
444 // Next state: deleted (Entry will be destroyed when the refcount reaches zero), or any
445 // other state if the entry is inserted into the cache.
446 NOT_IN_CACHE
447 };
448 
449 enum class EntryValueStatus : uint8_t {
450 // This entry has a known value. (Note that while there is nothing wrong per say with a value
451 // size of zero, v8 serialized data will always have a greater size.)
452 PRESENT,
453 
454 // This entry is known to be absent.
455 ABSENT,
456 
457 // This entry has not been fetched into cache yet, but the previous entry has `gapIsKnownEmpty =
458 // true`. Such entries are created as a result of list() operations, to mark the endpoint of the
459 // list range. List ranges are exclusive of their endpoint, hence the value associated with this
460 // key is commonly unknown.
461 UNKNOWN,
462 };
463 
464 struct CountedDelete;
465 
466 struct Entry: public kj::AtomicRefcounted {
467 // A cache entry.
468 //
469 // Entries are refcounted so that an operation which cares about a particular entry can keep
470 // it live even after it has been evicted or overwritten. In particular, because read
471 // operations are consistent with the time when read() was called, they may need to hold
472 // strong references to the entries they are reading, so that if the entries are overwritten,
473 // the read operation still has the original value from when it was called.
474 //
475 // The mutable content of an `Entry` is protected by the same mutex that protects
476 // `lru.cleanList`. `key` and `value` are declared `const` so that they can safely be used
477 // without a lock.
478 
479 Entry(ActorCache& cache, Key key, Value value);
480 Entry(ActorCache& cache, Key key, EntryValueStatus status);
481 Entry(Key key, Value value);
482 Entry(Key key, EntryValueStatus status);
483 ~Entry() noexcept(false);
484 KJ_DISALLOW_COPY_AND_MOVE(Entry);
485 
486 kj::Maybe<ActorCache&> maybeCache;
487 const Key key;
488 
489 private:
490 // The value associated with this key. If our `valueStatus` below is `ABSENT` or `UNKNOWN`, it
491 // will have size 0.
492 //
493 // `value` cannot change after the `Entry` is constructed. When a key is overwritten, the
494 // existing `Entry` is removed from the map and replaced with a new one, so that `value` does
495 // not need to be modified. This allows us to avoid copying `Entry` objects by refcounting
496 // them instead, especially in the case of a read operation which is only partially fulfilled
497 // from cache and needs to remember the original cached values even if they are overwritten
498 // before the read completes.
499 const Value value;
500 EntryValueStatus valueStatus;
501 
502 // This enum indicates how synchronized this entry is with storage.
503 EntrySyncStatus syncStatus = EntrySyncStatus::NOT_IN_CACHE;
504 
505 public:
506 EntryValueStatus getValueStatus() const {
507 return valueStatus;
508 }
509 
510 inline EntrySyncStatus getSyncStatus() const {
511 return syncStatus;
512 }
513 
514 kj::Maybe<ValuePtr> getValuePtr() const {
515 if (valueStatus == EntryValueStatus::PRESENT) {
516 return value.asPtr();
517 } else {
518 return kj::none;
519 }
520 }
521 kj::Maybe<Value> getValue() const {
522 KJ_IF_SOME(ptr, getValuePtr()) {
523 return ptr.attach(kj::atomicAddRef(*this));
524 } else {
525 return kj::none;
526 }
527 }
528 
529 void setNotInCache() {
530 syncStatus = EntrySyncStatus::NOT_IN_CACHE;
531 }
532 
533 // Avoid using setClean() and setDirty() directly. If you want to set the status to CLEAN or
534 // DIRTY, consider using the addToCleanList() and addToDirtyList() methods. This helps us keep
535 // the state transitions manageable.
536 void setClean() {
537 syncStatus = EntrySyncStatus::CLEAN;
538 }
539 
540 void setDirty() {
541 syncStatus = EntrySyncStatus::DIRTY;
542 }
543 
544 bool isDirty() const {
545 switch (getSyncStatus()) {
546 case EntrySyncStatus::DIRTY: {
547 return true;
548 }
549 case EntrySyncStatus::CLEAN: {
550 return false;
551 }
552 case EntrySyncStatus::NOT_IN_CACHE: {
553 KJ_FAIL_ASSERT("NOT_IN_CACHE entries should not be in the map or flushing");
554 }
555 }
556 }
557 
558 bool isStale = false;
559 bool flushStarted = false;
560 
561 // If true, then a past list() operation covered the space between this entry and the following
562 // entry, meaning that we know for sure that there are no other keys on disk between them.
563 bool gapIsKnownEmpty = false;
564 
565 // If true, then this entry should be evicted from cache immediately when it becomes CLEAN.
566 // The entry still needs to reside in cache while DIRTY since we need to store it
567 // somewhere, and so we might as well serve cache hits based on it in the meantime.
568 bool noCache = false;
569 
570 // In the DIRTY state, if this entry was originally created as the result of a
571 // `delete()` call, and as such the caller needs to receive a count of deletions, then this
572 // tracks that need. Note that only one caller could ever be waiting on this, because
573 // subsequent delete() calls can be counted based on the cache content. This can be false
574 // if no delete operations need a count from this entry.
575 bool isCountedDelete = false;
576 
577 // This Entry is part of a CountedDelete, but has since been overwritten via a put().
578 // This is really only useful in determining if we need to retry the deletion of this entry from
579 // storage, since we're interested in the number of deleted records. If we already got the count,
580 // we won't include this entry as part of our retried delete.
581 bool overwritingCountedDelete = false;
582 
583 // If CLEAN, the entry will be in the SharedLru's `cleanList`.
584 //
585 // If DIRTY, the entry will be in `dirtyList`.
586 kj::ListLink<Entry> link;
587 
588 size_t size() const {
589 return sizeof(*this) + key.size() + value.size();
590 }
591 };
592 
593 // Callbacks for a kj::TreeIndex for a kj::Table<kj::Own<Entry>>.
594 class EntryTableCallbacks {
595 public:
596 inline KeyPtr keyForRow(const kj::Own<Entry>& row) const {
597 return row->key;
598 }
599 
600 inline bool isBefore(const kj::Own<Entry>& row, KeyPtr key) const {
601 return row->key < key;
602 }
603 inline bool isBefore(const kj::Own<Entry>& a, const kj::Own<Entry>& b) const {
604 return a->key < b->key;
605 }
606 
607 inline bool matches(const kj::Own<Entry>& row, KeyPtr key) const {
608 return row->key == key;
609 }
610 };
611 
612 // When delete() is called with one or more keys that aren't in cache, we will need to get
613 // feedback from the database in order to report a count of deletions back to the application.
614 // Entries that were originally added to the cache as part of such a `delete()` will reference
615 // a `CountedDelete`.
616 //
617 // This object can only be manipulated in the thread that owns the specific actor that made
618 // the request. That works out fine since CountedDelete only ever exists for dirty entries,
619 // which won't be touched cross-thread by the LRU.
620 struct CountedDelete final: public kj::Refcounted {
621 CountedDelete() = default;
622 KJ_DISALLOW_COPY_AND_MOVE(CountedDelete);
623 
624 kj::Promise<void> forgiveIfFinished(kj::Promise<void> promise) {
625 try {
626 co_await promise;
627 } catch (...) {
628 if (isFinished) {
629 // We already flushed, so it's OK that the promise threw.
630 co_return;
631 } else {
632 throw;
633 }
634 }
635 }
636 
637 // Running count of entries that existed before the delete.
638 uint countDeleted = 0;
639 
640 // Did this particular counted delete succeed within a transaction? In other words, did we
641 // already get the count? Even if we got the count, we may need to retry if the transaction
642 // itself failed, though we won't need to get the count again.
643 bool completedInTransaction = false;
644 
645 // Did this particular counted delete succeed? Note that this can be true even if the flush
646 // failed on a different batch of operations.
647 bool isFinished = false;
648 
649 // The entries are associated with this counted delete.
650 kj::Vector<kj::Own<Entry>> entries;
651 };
652 
653 class CountedDeleteWaiter {
654 public:
655 explicit CountedDeleteWaiter(ActorCache& cache, kj::Own<CountedDelete> state)
656 : cache(cache),
657 state(kj::mv(state)) {
658 // Register this operation so that we can batch it properly during flush.
659 cache.countedDeletes.insert(this->state.get());
660 }
661 KJ_DISALLOW_COPY_AND_MOVE(CountedDeleteWaiter);
662 ~CountedDeleteWaiter() noexcept(false) {
663 for (auto& entry: state->entries) {
664 // Let each entry associated with this counted delete know that we aren't waiting anymore.
665 entry->isCountedDelete = false;
666 }
667 
668 // Since the count of deleted pairs is no longer required, we don't need to batch the ops.
669 // Note that we're doing eraseMatch since the pointer is a temporary literal.
670 cache.countedDeletes.eraseMatch(state.get());
671 }
672 
673 const CountedDelete& getCountedDelete() const {
674 return *state;
675 }
676 
677 private:
678 ActorCache& cache;
679 kj::Own<CountedDelete> state;
680 };
681 
682 kj::HashSet<CountedDelete*> countedDeletes;
683 
684 rpc::ActorStorage::Stage::Client storage;
685 const SharedLru& lru;
686 OutputGate& gate;
687 Hooks& hooks;
688 const kj::MonotonicClock& clock;
689 
690 // Wrapper around kj::List that keeps track of the total size of all elements.
691 class DirtyList {
692 public:
693 void add(Entry& entry) {
694 inner.add(entry);
695 innerSize += entry.size();
696 }
697 
698 void remove(Entry& entry) {
699 inner.remove(entry);
700 innerSize -= entry.size();
701 }
702 
703 size_t sizeInBytes() {
704 return innerSize;
705 }
706 
707 auto begin() {
708 return inner.begin();
709 }
710 auto end() {
711 return inner.end();
712 }
713 
714 private:
715 kj::List<Entry, &Entry::link> inner;
716 size_t innerSize = 0;
717 };
718 
719 // List of entries in DIRTY state. New dirty entries are added to the end. If any
720 // flushing entries are present, they always appear strictly before non-flushing entries.
721 DirtyList dirtyList;
722 
723 // Map of current known values for keys. Searchable by key, including ordered iteration.
724 //
725 // This map is protected by the same lock as lru.cleanList. ExternalMutexGuarded helps enforce
726 // this.
727 kj::ExternalMutexGuarded<kj::Table<kj::Own<Entry>, kj::TreeIndex<EntryTableCallbacks>>>
728 currentValues;
729 
730 struct UnknownAlarmTime {};
731 struct KnownAlarmTime {
732 enum class Status { CLEAN, DIRTY, FLUSHING } status;
733 kj::Maybe<kj::Date> time;
734 bool noCache = false;
735 };
736 
737 // Used by armAlarmHandler to know if a write needs to happen after the handler finishes
738 // to clear the alarm time.
739 struct DeferredAlarmDelete {
740 enum class Status { WAITING, READY, FLUSHING } status;
741 
742 // Set to a time to pass as `timeToDelete` when making the delete call.
743 kj::Date timeToDelete;
744 
745 // When the delete finishes, set to whether or not it succeeded.
746 kj::Maybe<bool> wasDeleted;
747 
748 bool noCache = false;
749 
750 // Trace span for the alarm handler, used when scheduling the deferred alarm deletion flush.
751 SpanParent traceSpan = nullptr;
752 };
753 
754 kj::OneOf<UnknownAlarmTime, KnownAlarmTime, DeferredAlarmDelete> currentAlarmTime =
755 UnknownAlarmTime{};
756 
757 struct ReadCompletionChain: public kj::Refcounted {
758 kj::Maybe<kj::Own<ReadCompletionChain>> next;
759 kj::Maybe<kj::Own<kj::PromiseFulfiller<void>>> fulfiller;
760 ReadCompletionChain() = default;
761 ~ReadCompletionChain() noexcept(false);
762 KJ_DISALLOW_COPY_AND_MOVE(ReadCompletionChain);
763 };
764 // Used to implement waitForPastReads(). See that function to understand how it works...
765 kj::Own<ReadCompletionChain> readCompletionChain = kj::refcounted<ReadCompletionChain>();
766 
767 // True if ensureFlushScheduled() has been called but the flush has not started yet.
768 bool flushScheduled = false;
769 
770 // When flushScheduled is true, indicates whether the output gate is already waiting on said
771 // flush. The first write that does *not* set `allowUnconfirmed` causes the output gate to be
772 // applied.
773 bool flushScheduledWithOutputGate = false;
774 
775 // Trace span for the current flush operation, captured from the first write that triggers
776 // a flush batch. Used for the output gate lock hold trace.
777 SpanParent currentFlushSpan = nullptr;
778 
779 // The count of the number of flushes that have been queued without yet resolving.
780 size_t flushesEnqueued = 0;
781 
782 struct DeleteAllState {
783 // If deleteAll() was called since the last flush, these are all the dirty entries that existed
784 // in the cache immediately before the deleteAll(). Since deleteAll() cannot be part of a
785 // transaction, in order to maintain ordering guarantees, we'll need to flush these entries
786 // first, then perform the deleteAll(), then flush any entries that were dirtied after the
787 // deleteAll().
788 kj::Vector<kj::Own<Entry>> deletedDirty;
789 kj::Own<kj::PromiseFulfiller<uint>> countFulfiller;
790 
791 // If true, the alarm should also be deleted after the deleteAll() RPC succeeds.
792 bool deleteAlarm = false;
793 };
794 
795 kj::Maybe<DeleteAllState> requestedDeleteAll;
796 
797 // Promise for the completion of the previous flush. We can only execute one flushImpl() at a time
798 // because we can't allow out-of-order writes.
799 kj::ForkedPromise<void> lastFlush = kj::Promise<void>(kj::READY_NOW).fork();
800 // TODO(perf): If we could rely on e-order on the ActorStorage API, we could pipeline additional
801 // writes and not have to worry about this. However, at present, ActorStorage has automatic
802 // reconnect behavior at the supervisor layer which violates e-order.
803 
804 // Did we hit a problem that makes the ActorCache unusable? If so this is the exception that
805 // describes the problem.
806 kj::Maybe<kj::Exception> maybeTerminalException;
807 
808 // Will be canceled if and when `oomException` becomes non-null.
809 kj::Canceler oomCanceler;
810 
811 // Type of a lock on `SharedLru::cleanList`. We use the same lock to protect `currentValues`.
812 using Lock = kj::Locked<kj::List<Entry, &Entry::link>>;
813 
814 // Add this entry to the clean list and set its status to CLEAN.
815 // This doesn't do much, but it makes it easier to track what's going on.
816 void addToCleanList(Lock& listLock, Entry& entryRef) {
817 entryRef.setClean();
818 listLock->add(entryRef);
819 }
820 
821 // Add this entry to the dirty list and set its status to DIRTY.
822 // This doesn't do much, but it makes it easier to track what's going on.
823 void addToDirtyList(Entry& entryRef) {
824 entryRef.setDirty();
825 dirtyList.add(entryRef);
826 }
827 
828 // Indicate that an entry was observed by a read operation and so should be moved to the end of
829 // the LRU queue.
830 void touchEntry(Lock& lock, Entry& entry);
831 
832 // TODO(soon) This function mostly belongs on the SharedLru, not the ActorCache. Notably,
833 // `removeEntry()` has to do with the shared clean list but `evictEntry()` has to do with
834 // the non-shared map. It is like this for now because generalizing the SharedLru into an
835 // IsolateCache is bigger work.
836 void removeEntry(Lock& lock, Entry& entry);
837 
838 // Look for a key in cache, returning a strong reference on the matching entry.
839 //
840 // Note that the returned entry could have `EntryValueStatus::UNKNOWN` which means we do not know
841 // if it is in storage or `EntryValueStatus::ABSENT` which means we know it is not in storage.
842 kj::Own<Entry> findInCache(Lock& lock, KeyPtr key, const ReadOptions& options);
843 
844 // Add an entry to the cache, where the entry was the result of reading from storage. If another
845 // entry with the same key has been inserted in the meantime, then the new entry will not be
846 // inserted and will instead immediately have state NOT_IN_CACHE.
847 //
848 // Either way, a strong reference to the entry is returned.
849 kj::Own<Entry> addReadResultToCache(
850 Lock& lock, Key key, kj::Maybe<capnp::Data::Reader> value, const ReadOptions& readOptions);
851 
852 // Mark all gaps empty between the begin and end key.
853 void markGapsEmpty(Lock& lock, KeyPtr begin, kj::Maybe<KeyPtr> end, const ReadOptions& options);
854 
855 // Implements put() or delete(). Multi-key variants call this for each key.
856 void putImpl(Lock& lock,
857 kj::Own<Entry> newEntry,
858 const WriteOptions& options,
859 kj::Maybe<CountedDelete&> counted,
860 SpanParent traceSpan);
861 
862 kj::Promise<kj::Maybe<Value>> getImpl(kj::Own<Entry> entry, ReadOptions options);
863 
864 // Ensure that we will flush dirty entries soon.
865 // The traceSpan is captured only on the first call that starts a new flush batch.
866 void ensureFlushScheduled(const WriteOptions& options, SpanParent traceSpan);
867 
868 // Schedule a read RPC. The given function will be invoked and provided with an
869 // ActorStorage::Operations::Client on which the read operation should be performed. The function
870 // might be called multiple times. The first call may be synchronous.
871 //
872 // This method has two purposes:
873 // - Retry operations that fail due to disconnects.
874 // - Ensure that reads cannot be re-ordered after writes that were originally scheduled later.
875 //
876 // Note that `function()` must return a plain `Promise`, not a `capnp::RemotePromise`, because
877 // it is necessary to `.attach()` something to it. Use `.dropPipeline()` to convert a
878 // `RemotePromise` to a plain `Promise`.
879 template <typename Func>
880 kj::PromiseForResult<Func, rpc::ActorStorage::Operations::Client> scheduleStorageRead(
881 Func&& function);
882 
883 // Wait until all read operations that are currently in-flight have completed or failed
884 // (including exhausting all retries). Does not propagate the read exception, if any. This is
885 // used for ordering, to make sure a write is not committed too early such that it interferes
886 // with a previous read.
887 kj::Promise<void> waitForPastReads();
888 
889 kj::Promise<void> flushImpl(uint retryCount = 0);
890 kj::Promise<void> flushImplDeleteAll(uint retryCount = 0);
891 
892 struct FlushBatch {
893 size_t pairCount = 0;
894 size_t wordCount = 0;
895 };
896 struct PutFlush {
897 kj::Vector<kj::Own<Entry>> entries;
898 kj::Vector<FlushBatch> batches;
899 };
900 struct MutedDeleteFlush {
901 kj::Vector<kj::Own<Entry>> entries;
902 kj::Vector<FlushBatch> batches;
903 };
904 struct CountedDeleteFlush {
905 kj::Own<CountedDelete> countedDelete;
906 kj::Vector<FlushBatch> batches;
907 };
908 using CountedDeleteFlushes = kj::Array<CountedDeleteFlush>;
909 kj::Promise<void> startFlushTransaction();
910 kj::Promise<void> flushImplUsingSinglePut(PutFlush putFlush);
911 kj::Promise<void> flushImplUsingSingleMutedDelete(MutedDeleteFlush mutedFlush);
912 kj::Promise<void> flushImplUsingSingleCountedDelete(CountedDeleteFlush countedFlush);
913 kj::Promise<void> flushImplAlarmOnly(DirtyAlarm dirty);
914 kj::Promise<void> flushImplUsingTxn(PutFlush putFlush,
915 MutedDeleteFlush mutedDeleteFlush,
916 CountedDeleteFlushes countedDeleteFlushes,
917 MaybeAlarmChange maybeAlarmChange);
918 
919 // Carefully remove a clean entry from `currentValues`, making sure to update gaps.
920 void evictEntry(Lock& lock, Entry& entry);
921 
922 // Drop the entire cache. Called during destructor and on OOM.
923 void clear(Lock& lock);
924 
925 // Throws OOM exception if `oom` is true.
926 void requireNotTerminal(SpanParent traceSpan);
927 
928 // Evict cache entries as needed to reach the target memory usage. If the cache has exceeded the
929 // hard limit, trigger an OOM, canceling all RPCs and breaking the output gate.
930 void evictOrOomIfNeeded(Lock& lock);
931 
932 // If the LRU is currently over the soft limit, returns a promise that resolves when it is
933 // back under the limit.
934 kj::Maybe<kj::Promise<void>> getBackpressure();
935 
936 class GetMultiStreamImpl;
937 class ForwardListStreamImpl;
938 class ReverseListStreamImpl;
939 friend class ActorCacheOps::GetResultList;
940};
941 
942class ActorCacheOps::GetResultList {
943 using Entry = ActorCache::Entry;
944 
945 public:
946 class Iterator {
947 public:
948 KeyValuePtrPairWithCache operator*() {
949 KJ_IREQUIRE(ptr->get()->getValueStatus() == ActorCache::EntryValueStatus::PRESENT);
950 return {ptr->get()->key, ptr->get()->getValuePtr().orDefault({}), *statusPtr};
951 }
952 Iterator& operator++() {
953 ++ptr;
954 ++statusPtr;
955 return *this;
956 }
957 Iterator operator++(int) {
958 auto copy = *this;
959 ++ptr;
960 ++statusPtr;
961 return copy;
962 }
963 bool operator==(const Iterator& other) const {
964 return ptr == other.ptr && statusPtr == other.statusPtr;
965 }
966 
967 private:
968 const kj::Own<Entry>* ptr;
969 const CacheStatus* statusPtr;
970 
971 explicit Iterator(const kj::Own<Entry>* ptr, const CacheStatus* statusPtr)
972 : ptr(ptr),
973 statusPtr(statusPtr) {}
974 friend class GetResultList;
975 };
976 
977 Iterator begin() const {
978 return Iterator(entries.begin(), cacheStatuses.begin());
979 }
980 Iterator end() const {
981 return Iterator(entries.end(), cacheStatuses.end());
982 }
983 size_t size() const {
984 return entries.size();
985 }
986 
987 // Construct a simple GetResultList from key-value pairs.
988 explicit GetResultList(kj::Vector<KeyValuePair> contents);
989 
990 private:
991 kj::Vector<kj::Own<Entry>> entries;
992 kj::Vector<CacheStatus> cacheStatuses;
993 
994 enum Order { FORWARD, REVERSE };
995 
996 // Merges `cachedEntries` and `fetchedEntries`, which should each already be sorted in the
997 // given order. If a key exists in both, `cachedEntries` is preferred.
998 //
999 // After merging, if an entry's value is null, it is dropped.
1000 //
1001 // The final result is truncated to `limit`, if any.
1002 //
1003 // The idea is that `cachedEntries` is the set of entries that were loaded from cache while
1004 // `fetchedEntries` is the set read from storage.
1005 explicit GetResultList(kj::Vector<kj::Own<Entry>> cachedEntries,
1006 kj::Vector<kj::Own<Entry>> fetchedEntries,
1007 Order order,
1008 kj::Maybe<uint> limit = kj::none);
1009 
1010 friend class ActorCache;
1011};
1012 
1013// Options to ActorCache::SharedLru's constructor. Declared at top level so that it can be
1014// forward-declared elsewhere.
1015struct ActorCacheSharedLruOptions {
1016 // Memory usage that the LRU will try to stay under by evicting clean values.
1017 size_t softLimit;
1018 
1019 // Memory usage at which operations should start failing and actors should be killed for
1020 // exceeding memory limits.
1021 size_t hardLimit;
1022 
1023 // Time period after which a value that hasn't been accessed at all should be evicted even if
1024 // the total cache size is below `softLimit`.
1025 kj::Duration staleTimeout;
1026 
1027 // How many bytes in a particular ActorCache can be dirty before backpressure is applied on the
1028 // app.
1029 size_t dirtyListByteLimit;
1030 
1031 // Maximum number of keys in a single RPC message during a flush. If a message would be larger
1032 // than this, it'll be split into multiple calls.
1033 //
1034 // This should typically be set to ActorStorageClientImpl::MAX_KEYS from
1035 // supervisor/actor-storage.h.
1036 size_t maxKeysPerRpc;
1037 
1038 // If true, assume `noCache` for all operations.
1039 bool noCache = false;
1040 
1041 // If true, don't actually flush anything. This is used in preview sessions, since they keep
1042 // state strictly in memory.
1043 bool neverFlush = false;
1044};
1045 
1046class ActorCache::SharedLru {
1047 public:
1048 using Options = ActorCacheSharedLruOptions;
1049 
1050 explicit SharedLru(Options options);
1051 
1052 ~SharedLru() noexcept(false);
1053 KJ_DISALLOW_COPY_AND_MOVE(SharedLru);
1054 
1055 // Mostly for testing.
1056 size_t currentSize() const {
1057 return size.load(std::memory_order_relaxed);
1058 }
1059 
1060 private:
1061 const Options options;
1062 
1063 // List of clean values, across all caches, ordered from least-recently-used to
1064 // most-recently-used.
1065 kj::MutexGuarded<kj::List<Entry, &Entry::link>> cleanList;
1066 
1067 // Total byte size of everything that is cached, including dirty values that aren't in `cleanList`.
1068 mutable std::atomic<size_t> size = 0;
1069 
1070 // TimePoint when we should next evict stale entries. Represented as an int64_t of nanoseconds
1071 // instead of kj::TimePoint to allow for atomic operations.
1072 mutable std::atomic<int64_t> nextStaleCheckNs = 0;
1073 
1074 // Evict cache entries as needed according to the cache limits. Returns true if the hard limit
1075 // is exceeded and nothing can be evicted, in which case the caller should fail out in the
1076 // appropriate way for the kind of operation being performed.
1077 bool evictIfNeeded(Lock& lock) const KJ_WARN_UNUSED_RESULT;
1078 
1079 friend class ActorCache;
1080};
1081 
1082// A transaction represents a set of writes that haven't been committed. The transaction can be
1083// discarded without committing.
1084//
1085// ActorCache::Transaction intentionally does NOT detect conflicts with concurrent transactions.
1086// It is up to a higher layer to make sure that only one transaction occurs at a time, perhaps
1087// using critical sections.
1088class ActorCache::Transaction final: public ActorCacheInterface::Transaction {
1089 public:
1090 Transaction(ActorCache& cache);
1091 ~Transaction() noexcept(false);
1092 
1093 kj::OneOf<kj::Maybe<Value>, kj::Promise<kj::Maybe<Value>>> get(
1094 Key key, ReadOptions options) override;
1095 kj::OneOf<GetResultList, kj::Promise<GetResultList>> get(
1096 kj::Array<Key> keys, ReadOptions options) override;
1097 kj::OneOf<kj::Maybe<kj::Date>, kj::Promise<kj::Maybe<kj::Date>>> getAlarm(
1098 ReadOptions options) override;
1099 kj::OneOf<GetResultList, kj::Promise<GetResultList>> list(
1100 Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) override;
1101 kj::OneOf<GetResultList, kj::Promise<GetResultList>> listReverse(
1102 Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) override;
1103 kj::Maybe<kj::Promise<void>> put(
1104 Key key, Value value, WriteOptions options, SpanParent traceSpan) override;
1105 kj::Maybe<kj::Promise<void>> put(
1106 kj::Array<KeyValuePair> pairs, WriteOptions options, SpanParent traceSpan) override;
1107 kj::OneOf<bool, kj::Promise<bool>> delete_(
1108 Key key, WriteOptions options, SpanParent traceSpan) override;
1109 kj::OneOf<uint, kj::Promise<uint>> delete_(
1110 kj::Array<Key> keys, WriteOptions options, SpanParent traceSpan) override;
1111 kj::Maybe<kj::Promise<void>> setAlarm(
1112 kj::Maybe<kj::Date> newAlarmTime, WriteOptions options, SpanParent traceSpan) override;
1113 // Same interface as ActorCache.
1114 //
1115 // Read ops will reflect the previous writes made to the transaction even though they aren't
1116 // committed yet.
1117 
1118 kj::Maybe<kj::Promise<void>> commit() override;
1119 kj::Promise<void> rollback() override;
1120 // Implements ActorCacheInterface::Transaction.
1121 
1122 private:
1123 ActorCache& cache;
1124 
1125 struct Change {
1126 kj::Own<Entry> entry;
1127 WriteOptions options;
1128 };
1129 
1130 // Callbacks for a kj::TreeIndex for a kj::Table<Change>.
1131 class ChangeTableCallbacks {
1132 public:
1133 inline KeyPtr keyForRow(const Change& row) const {
1134 return row.entry->key;
1135 }
1136 
1137 inline bool isBefore(const Change& row, KeyPtr key) const {
1138 return row.entry->key < key;
1139 }
1140 inline bool matches(const Change& row, KeyPtr key) const {
1141 return row.entry->key == key;
1142 }
1143 };
1144 
1145 kj::Table<Change, kj::TreeIndex<ChangeTableCallbacks>> entriesToWrite;
1146 
1147 kj::Maybe<DirtyAlarmWithOptions> alarmChange;
1148 
1149 // Trace span captured from each write, to be used when commit() flushes changes.
1150 SpanParent commitSpan = nullptr;
1151 
1152 // Merge the changes in the transaction with the results from reading from the underlying
1153 // ActorCache.
1154 kj::OneOf<GetResultList, kj::Promise<GetResultList>> merge(
1155 kj::Vector<kj::Own<Entry>> changedEntries,
1156 kj::OneOf<GetResultList, kj::Promise<GetResultList>> cacheRead,
1157 GetResultList::Order order);
1158 
1159 // Adds the given key/value pair to `changes`. If an existing entry is replaced, *count is
1160 // incremented if it was a positive entry. If no existing entry is replaced, then the key
1161 // is returned, indicating that if a count is needed, we'll need to inspect cache/disk.
1162 kj::Maybe<KeyPtr> putImpl(Lock& lock,
1163 kj::Own<Entry> entry,
1164 const WriteOptions& options,
1165 kj::Maybe<uint&> count = kj::none);
1166};
1167 
1168} // namespace workerd