File
Blob: src/workerd/io/actor-cache.c++
| 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 | |
| 7 | #include <workerd/io/actor-storage.h> |
| 8 | #include <workerd/io/io-gate.h> |
| 9 | #include <workerd/jsg/exception.h> |
| 10 | #include <workerd/util/duration-exceeded-logger.h> |
| 11 | #include <workerd/util/exception.h> |
| 12 | #include <workerd/util/sentry.h> |
| 13 | |
| 14 | #include <kj/debug.h> |
| 15 | |
| 16 | #include <algorithm> |
| 17 | |
| 18 | namespace workerd { |
| 19 | |
| 20 | // Max size, in words, of a storage RPC request. Set to 16MiB because our storage backend has a |
| 21 | // hard limit of 16MiB per operation. |
| 22 | // |
| 23 | // (Also, at 64MiB we'd hit the Cap'n Proto message size limit.) |
| 24 | // |
| 25 | // Note that in practice, the key size limit (options.maxKeysPerRpc) will kick in long before we |
| 26 | // hit this limit, so this is just a sanity check. |
| 27 | static constexpr size_t MAX_ACTOR_STORAGE_RPC_WORDS = (16u << 20) / sizeof(capnp::word); |
| 28 | |
| 29 | const ActorCache::Hooks ActorCache::Hooks::DEFAULT; |
| 30 | |
| 31 | namespace { |
| 32 | |
| 33 | // Utility functions for recording latency metrics via a one-liner in the callers below. |
| 34 | auto recordStorageRead(ActorCache::Hooks& hooks, const kj::MonotonicClock& clock) { |
| 35 | auto start = clock.now(); |
| 36 | return kj::defer([start, &hooks, &clock]() { hooks.storageReadCompleted(clock.now() - start); }); |
| 37 | } |
| 38 | auto recordStorageWrite(ActorCache::Hooks& hooks, const kj::MonotonicClock& clock) { |
| 39 | auto start = clock.now(); |
| 40 | return kj::defer([start, &hooks, &clock]() { hooks.storageWriteCompleted(clock.now() - start); }); |
| 41 | } |
| 42 | |
| 43 | } // namespace |
| 44 | |
| 45 | ActorCache::ActorCache( |
| 46 | rpc::ActorStorage::Stage::Client storage, const SharedLru& lru, OutputGate& gate, Hooks& hooks) |
| 47 | : storage(kj::mv(storage)), |
| 48 | lru(lru), |
| 49 | gate(gate), |
| 50 | hooks(hooks), |
| 51 | clock(kj::systemPreciseMonotonicClock()), |
| 52 | currentValues(lru.cleanList.lockExclusive()) {} |
| 53 | |
| 54 | ActorCache::~ActorCache() noexcept(false) { |
| 55 | // Need to remove all entries from any lists they might be in. |
| 56 | auto lock = lru.cleanList.lockExclusive(); |
| 57 | clear(lock); |
| 58 | } |
| 59 | |
| 60 | void ActorCache::clear(Lock& lock) { |
| 61 | for (auto& entry: currentValues.get(lock)) { |
| 62 | removeEntry(lock, *entry); |
| 63 | } |
| 64 | currentValues.get(lock).clear(); |
| 65 | } |
| 66 | |
| 67 | ActorCache::Entry::Entry(ActorCache& cache, Key key, Value value) |
| 68 | : maybeCache(cache), |
| 69 | key(kj::mv(key)), |
| 70 | value(kj::mv(value)), |
| 71 | valueStatus(EntryValueStatus::PRESENT) { |
| 72 | KJ_IF_SOME(c, maybeCache) { |
| 73 | c.lru.size.fetch_add(size(), std::memory_order_relaxed); |
| 74 | } |
| 75 | } |
| 76 | |
| 77 | ActorCache::Entry::Entry(ActorCache& cache, Key key, EntryValueStatus valueStatus) |
| 78 | : maybeCache(cache), |
| 79 | key(kj::mv(key)), |
| 80 | valueStatus(valueStatus) { |
| 81 | KJ_IASSERT(valueStatus != EntryValueStatus::PRESENT, |
| 82 | "Pass a serialized empty v8 value if you want a present but empty entry!"); |
| 83 | KJ_IF_SOME(c, maybeCache) { |
| 84 | c.lru.size.fetch_add(size(), std::memory_order_relaxed); |
| 85 | } |
| 86 | } |
| 87 | |
| 88 | ActorCache::Entry::Entry(Key key, Value value) |
| 89 | : key(kj::mv(key)), |
| 90 | value(kj::mv(value)), |
| 91 | valueStatus(EntryValueStatus::PRESENT) {} |
| 92 | ActorCache::Entry::Entry(Key key, EntryValueStatus valueStatus) |
| 93 | : key(kj::mv(key)), |
| 94 | valueStatus(valueStatus) {} |
| 95 | |
| 96 | ActorCache::Entry::~Entry() noexcept(false) { |
| 97 | KJ_IF_SOME(c, maybeCache) { |
| 98 | size_t size = this->size(); |
| 99 | |
| 100 | size_t before = c.lru.size.fetch_sub(size, std::memory_order_relaxed); |
| 101 | |
| 102 | if (KJ_UNLIKELY(before < size)) { |
| 103 | // underflow -- shouldn't happen, but just in case, let's fix |
| 104 | KJ_LOG(ERROR, "SharedLru size tracking inconsistency detected", before, size, |
| 105 | kj::getStackTrace()); |
| 106 | c.lru.size.store(0, std::memory_order_relaxed); |
| 107 | } |
| 108 | |
| 109 | if (link.isLinked()) { |
| 110 | switch (getSyncStatus()) { |
| 111 | case EntrySyncStatus::CLEAN: { |
| 112 | KJ_LOG(WARNING, "Entry destructed while still in the clean list"); |
| 113 | break; |
| 114 | } |
| 115 | case EntrySyncStatus::DIRTY: { |
| 116 | // Ah, we don't need a lock so we can just unlink ourselves. This is safe because we will |
| 117 | // only destruct a DIRTY entry on the actor's event loop. (We can destruct a CLEAN entry |
| 118 | // as part of evicting entries from the shared lru on a different event loop.) |
| 119 | c.dirtyList.remove(*this); |
| 120 | break; |
| 121 | } |
| 122 | case EntrySyncStatus::NOT_IN_CACHE: { |
| 123 | KJ_LOG(WARNING, "Entry with sync status NOT_IN_CACHE still in a list"); |
| 124 | break; |
| 125 | } |
| 126 | } |
| 127 | } |
| 128 | } |
| 129 | } |
| 130 | |
| 131 | ActorCache::SharedLru::SharedLru(Options options): options(options) {} |
| 132 | |
| 133 | ActorCache::SharedLru::~SharedLru() noexcept(false) { |
| 134 | KJ_REQUIRE(cleanList.getWithoutLock().empty(), |
| 135 | "ActorCache::SharedLru destroyed while an ActorCache still exists?"); |
| 136 | if (size.load(std::memory_order_relaxed) != 0) { |
| 137 | KJ_LOG(ERROR, |
| 138 | "SharedLru destroyed while cache entries still exist, " |
| 139 | "this will lead to use-after-free"); |
| 140 | } |
| 141 | } |
| 142 | |
| 143 | kj::Maybe<kj::Promise<void>> ActorCache::evictStale(kj::Date now) { |
| 144 | int64_t nowNs = (now - kj::UNIX_EPOCH) / kj::NANOSECONDS; |
| 145 | int64_t oldValue = lru.nextStaleCheckNs.load(std::memory_order_relaxed); |
| 146 | |
| 147 | if (nowNs >= oldValue) { |
| 148 | int64_t newValue = nowNs + lru.options.staleTimeout / kj::NANOSECONDS; |
| 149 | if (lru.nextStaleCheckNs.compare_exchange_strong(oldValue, newValue)) { |
| 150 | auto lock = lru.cleanList.lockExclusive(); |
| 151 | for (auto& entry: *lock) { |
| 152 | if (entry.isStale) { |
| 153 | auto& cache = KJ_ASSERT_NONNULL(entry.maybeCache); |
| 154 | cache.removeEntry(lock, entry); |
| 155 | cache.evictEntry(lock, entry); |
| 156 | } else { |
| 157 | entry.isStale = true; |
| 158 | } |
| 159 | } |
| 160 | } |
| 161 | } |
| 162 | |
| 163 | // Apply backpressure if we're over the soft limit. |
| 164 | return getBackpressure(); |
| 165 | } |
| 166 | |
| 167 | kj::OneOf<ActorCache::CancelAlarmHandler, ActorCache::RunAlarmHandler> ActorCache::armAlarmHandler( |
| 168 | kj::Date scheduledTime, |
| 169 | SpanParent parentSpan, |
| 170 | kj::Date currentTime KJ_UNUSED, |
| 171 | bool noCache, |
| 172 | kj::StringPtr actorId) { |
| 173 | noCache = noCache || lru.options.noCache; |
| 174 | |
| 175 | KJ_ASSERT(!currentAlarmTime.is<DeferredAlarmDelete>()); |
| 176 | bool alarmDeleteNeeded = true; |
| 177 | KJ_IF_SOME(t, currentAlarmTime.tryGet<KnownAlarmTime>()) { |
| 178 | if (t.time != scheduledTime) { |
| 179 | if (t.status == KnownAlarmTime::Status::CLEAN) { |
| 180 | // If there's a clean scheduledTime that is different from ours, this run should be |
| 181 | // canceled. |
| 182 | LOG_WARNING_PERIODICALLY("NOSENTRY actor-cache alarm handler canceled.", scheduledTime, |
| 183 | t.time.orDefault(kj::UNIX_EPOCH), actorId); |
| 184 | return CancelAlarmHandler{.waitBeforeCancel = kj::READY_NOW}; |
| 185 | } else { |
| 186 | // There's a alarm write that hasn't been set yet pending for a time different than ours -- |
| 187 | // We won't cancel the alarm because it hasn't been confirmed, but we shouldn't delete |
| 188 | // the pending write. |
| 189 | alarmDeleteNeeded = false; |
| 190 | } |
| 191 | } |
| 192 | } |
| 193 | |
| 194 | if (alarmDeleteNeeded) { |
| 195 | currentAlarmTime = DeferredAlarmDelete{ |
| 196 | .status = DeferredAlarmDelete::Status::WAITING, |
| 197 | .timeToDelete = scheduledTime, |
| 198 | .noCache = noCache, |
| 199 | .traceSpan = kj::mv(parentSpan), |
| 200 | }; |
| 201 | } |
| 202 | static const DeferredAlarmDeleter disposer; |
| 203 | return RunAlarmHandler{.deferredDelete = kj::Own<void>(this, disposer)}; |
| 204 | } |
| 205 | |
| 206 | void ActorCache::cancelDeferredAlarmDeletion() { |
| 207 | KJ_IF_SOME(deferredDelete, currentAlarmTime.tryGet<DeferredAlarmDelete>()) { |
| 208 | currentAlarmTime = KnownAlarmTime{.status = KnownAlarmTime::Status::CLEAN, |
| 209 | .time = deferredDelete.timeToDelete, |
| 210 | .noCache = deferredDelete.noCache}; |
| 211 | } |
| 212 | } |
| 213 | |
| 214 | kj::Promise<kj::Maybe<kj::Date>> ActorCache::abandonAlarm(kj::Date scheduledTime) { |
| 215 | // Called when AlarmManager has given up retrying an alarm after too many counted failures. |
| 216 | // Reset the in-memory alarm state to unknown so the next getAlarm() refetches from storage |
| 217 | // rather than serving a potentially stale cached value. |
| 218 | // Only act if we still have a clean KnownAlarmTime whose time matches the abandoned alarm. |
| 219 | KJ_IF_SOME(t, currentAlarmTime.tryGet<KnownAlarmTime>()) { |
| 220 | KJ_IF_SOME(storedTime, t.time) { |
| 221 | if (t.status == KnownAlarmTime::Status::CLEAN) { |
| 222 | if (storedTime == scheduledTime) { |
| 223 | currentAlarmTime = UnknownAlarmTime{}; |
| 224 | return kj::Maybe<kj::Date>(kj::none); |
| 225 | } else { |
| 226 | // The user set a different alarm. Return it so AlarmManager can re-register. |
| 227 | return kj::Maybe<kj::Date>(storedTime); |
| 228 | } |
| 229 | } |
| 230 | } |
| 231 | } |
| 232 | return kj::Maybe<kj::Date>(kj::none); |
| 233 | } |
| 234 | |
| 235 | kj::Maybe<kj::Promise<void>> ActorCache::getBackpressure() { |
| 236 | if (dirtyList.sizeInBytes() > lru.options.dirtyListByteLimit && !lru.options.neverFlush) { |
| 237 | // Wait for dirty entries to be flushed. |
| 238 | return lastFlush.addBranch().then([this]() -> kj::Promise<void> { |
| 239 | KJ_IF_SOME(p, getBackpressure()) { |
| 240 | return kj::mv(p); |
| 241 | } else { |
| 242 | return kj::READY_NOW; |
| 243 | } |
| 244 | }); |
| 245 | } |
| 246 | |
| 247 | // At one point, we tried applying backpressure if the total cache size was greater than |
| 248 | // `softLimit`. This turned out to be a bad idea. If the cache is over the limit due to dirty |
| 249 | // entries waiting to be flushed, then `dirtyListByteLimit` will actually kick in first (since |
| 250 | // it's by default 8MB of data). So if the cache is over the soft limit (which is typically more |
| 251 | // like 16MB), it could only be because a very large read operation has loaded a bunch of entries |
| 252 | // into memory but hasn't delivered them to the app yet. In this case, if we apply backpressure, |
| 253 | // then the app cannot make progress and therefore cannot receive the result of these reads! So it |
| 254 | // will just deadlock. |
| 255 | // |
| 256 | // Hence, it only makes sense to wait for dirty entries to be flushed, not to wait for overall |
| 257 | // size to go down. |
| 258 | return kj::none; |
| 259 | } |
| 260 | |
| 261 | void ActorCache::requireNotTerminal(SpanParent traceSpan) { |
| 262 | KJ_IF_SOME(e, maybeTerminalException) { |
| 263 | if (!gate.isBroken()) { |
| 264 | // We've tried to use storage after shutdown, break the output gate via `flushImpl()` so that |
| 265 | // we don't let the worker return stale state. This isn't strictly necessary but it does |
| 266 | // mirror previous behavior wherein we would use disabled storage via `flushImpl()` and break |
| 267 | // the output gate. |
| 268 | ensureFlushScheduled({}, kj::mv(traceSpan)); |
| 269 | } |
| 270 | |
| 271 | kj::throwFatalException(e.clone()); |
| 272 | } |
| 273 | } |
| 274 | |
| 275 | void ActorCache::evictOrOomIfNeeded(Lock& lock) { |
| 276 | if (lru.evictIfNeeded(lock)) { |
| 277 | auto exception = KJ_EXCEPTION(OVERLOADED, |
| 278 | "broken.exceededMemory; jsg.Error: Durable Object's isolate exceeded its memory limit due to overflowing the " |
| 279 | "storage cache. This could be due to writing too many values to storage without stopping " |
| 280 | "to wait for writes to complete, or due to reading too many values in a single operation " |
| 281 | "(e.g. a large list()). All objects in the isolate were reset."); |
| 282 | |
| 283 | // Add trace info sufficient to tell us which operation caused the failure. |
| 284 | exception.addTraceHere(); |
| 285 | exception.addTrace(__builtin_return_address(0)); |
| 286 | // We know this exception happens due to user error. Let's add an exception detail so we can |
| 287 | // parse it later. |
| 288 | exception.setDetail(jsg::EXCEPTION_IS_USER_ERROR, kj::heapArray<byte>(0)); |
| 289 | exception.setDetail(MEMORY_LIMIT_DETAIL_ID, kj::heapArray<byte>(0)); |
| 290 | |
| 291 | if (maybeTerminalException == kj::none) { |
| 292 | maybeTerminalException.emplace(exception.clone()); |
| 293 | } else { |
| 294 | // We've already experienced a terminal exception either from shutdown or OOM. Note that we |
| 295 | // still schedule the flush since shutdown does not. |
| 296 | } |
| 297 | |
| 298 | clear(lock); |
| 299 | oomCanceler.cancel(exception); |
| 300 | |
| 301 | if (!gate.isBroken()) { |
| 302 | // We want to break the OutputGate. We can't quite just do `gate.lockWhile(exception)` because |
| 303 | // that returns a promise which we'd then have to put somewhere so that we don't immediately |
| 304 | // cancel it. Instead, we can ensure that a flush has been scheduled. `flushImpl()`, when |
| 305 | // called, will throw an exception which breaks the gate. |
| 306 | ensureFlushScheduled(WriteOptions(), nullptr); |
| 307 | } |
| 308 | |
| 309 | kj::throwFatalException(kj::mv(exception)); |
| 310 | } |
| 311 | } |
| 312 | |
| 313 | bool ActorCache::SharedLru::evictIfNeeded(Lock& lock) const { |
| 314 | for (;;) { |
| 315 | size_t current = size.load(std::memory_order_relaxed); |
| 316 | if (current <= options.softLimit) { |
| 317 | // All good. |
| 318 | return false; |
| 319 | } |
| 320 | |
| 321 | // We're over the limit, let's evict stuff. |
| 322 | if (lock->empty()) { |
| 323 | // Nothing to evict. |
| 324 | return current > options.hardLimit; |
| 325 | } |
| 326 | |
| 327 | Entry& entry = lock->front(); |
| 328 | auto& cache = KJ_ASSERT_NONNULL(entry.maybeCache); |
| 329 | cache.removeEntry(lock, entry); |
| 330 | cache.evictEntry(lock, entry); |
| 331 | } |
| 332 | } |
| 333 | |
| 334 | void ActorCache::touchEntry(Lock& lock, Entry& entry) { |
| 335 | if (entry.getSyncStatus() == EntrySyncStatus::CLEAN) { |
| 336 | entry.isStale = false; |
| 337 | lock->remove(entry); |
| 338 | addToCleanList(lock, entry); |
| 339 | } |
| 340 | |
| 341 | // We only call `touchEntry` when the operation or the LRU has !noCache, so we want to cache this. |
| 342 | // |
| 343 | // If this is a dirty entry previously marked no-cache, remove that mark. This results in the |
| 344 | // same end state as if the entry had been flushed and evicted before the read -- it would have |
| 345 | // been read back, and then into cache. |
| 346 | entry.noCache = false; |
| 347 | } |
| 348 | |
| 349 | void ActorCache::removeEntry(Lock& lock, Entry& entry) { |
| 350 | switch (entry.getSyncStatus()) { |
| 351 | case EntrySyncStatus::DIRTY: { |
| 352 | dirtyList.remove(entry); |
| 353 | break; |
| 354 | } |
| 355 | case EntrySyncStatus::CLEAN: { |
| 356 | lock->remove(entry); |
| 357 | break; |
| 358 | } |
| 359 | case EntrySyncStatus::NOT_IN_CACHE: { |
| 360 | // Nothing to do! |
| 361 | break; |
| 362 | } |
| 363 | } |
| 364 | |
| 365 | entry.setNotInCache(); |
| 366 | } |
| 367 | |
| 368 | void ActorCache::evictEntry(Lock& lock, Entry& entry) { |
| 369 | auto& map = currentValues.get(lock); |
| 370 | auto ordered = map.ordered(); |
| 371 | auto iter = map.seek(entry.key); |
| 372 | |
| 373 | KJ_ASSERT(iter != ordered.end() && iter->get() == &entry); |
| 374 | |
| 375 | // If the previous entry has gapIsKnownEmpty, we need to set that false, because when we delete |
| 376 | // this entry, the previous entry's "gap" will now extend to the *next* entry. We definitely know |
| 377 | // that that the new gap is non-empty because we're evicting an entry inside that very gap. |
| 378 | // |
| 379 | // TODO(perf): Maybe we should instead replace the evicted entry with an UNKNOWN entry in this |
| 380 | // case? The problem is, when the app accesses a key in the gap, the LRU time of the previous |
| 381 | // entry gets bumped, but the _next_ entry does not get bumped. Hence these accesses won't |
| 382 | // prevent the next entry from being evicted, and when it is, the gap effectively gets evicted |
| 383 | // too, leading to a cache miss on a key that had been recently accessed. This is a pretty |
| 384 | // obscure scenario, though, and after one cache miss the key would then be in cache again. |
| 385 | if (iter != ordered.begin()) { |
| 386 | auto prev = iter; |
| 387 | --prev; |
| 388 | prev->get()->gapIsKnownEmpty = false; |
| 389 | } |
| 390 | |
| 391 | map.erase(*iter); |
| 392 | } |
| 393 | |
| 394 | void ActorCache::verifyConsistencyForTest() { |
| 395 | auto lock = lru.cleanList.lockExclusive(); |
| 396 | currentValues.get(lock).verify(); // verify the table's BTreeIndex |
| 397 | bool prevGapIsKnownEmpty = false; |
| 398 | kj::Maybe<kj::StringPtr> prevKey = kj::none; |
| 399 | for (auto& entry: currentValues.get(lock).ordered()) { |
| 400 | KJ_IF_SOME(p, prevKey) { |
| 401 | KJ_ASSERT(entry->key > p, "keys out of order?", p, entry->key); |
| 402 | } |
| 403 | prevKey = entry->key; |
| 404 | auto& key = entry->key; |
| 405 | switch (entry->getValueStatus()) { |
| 406 | case EntryValueStatus::ABSENT: { |
| 407 | KJ_ASSERT(!prevGapIsKnownEmpty || !entry->gapIsKnownEmpty, |
| 408 | "clean negative entry in the middle of a known-empty gap is redundant", key); |
| 409 | break; |
| 410 | } |
| 411 | case EntryValueStatus::PRESENT: { |
| 412 | // Nothing to do for PRESENT! |
| 413 | break; |
| 414 | } |
| 415 | case EntryValueStatus::UNKNOWN: { |
| 416 | KJ_ASSERT(!entry->gapIsKnownEmpty, "entry can't be followed by known-empty gap", key); |
| 417 | break; |
| 418 | } |
| 419 | } |
| 420 | |
| 421 | KJ_ASSERT(entry->getSyncStatus() != EntrySyncStatus::NOT_IN_CACHE, |
| 422 | "entry should not appear in map", entry->key); |
| 423 | KJ_ASSERT(entry->link.isLinked()); |
| 424 | |
| 425 | prevGapIsKnownEmpty = entry->gapIsKnownEmpty; |
| 426 | } |
| 427 | } |
| 428 | |
| 429 | // ======================================================================================= |
| 430 | // read operations |
| 431 | |
| 432 | kj::OneOf<kj::Maybe<ActorCache::Value>, kj::Promise<kj::Maybe<ActorCache::Value>>> ActorCache::get( |
| 433 | Key key, ReadOptions options) { |
| 434 | ActorStorageLimits::checkMaxKeySize(key); |
| 435 | |
| 436 | options.noCache = options.noCache || lru.options.noCache; |
| 437 | requireNotTerminal(nullptr); |
| 438 | |
| 439 | auto lock = lru.cleanList.lockExclusive(); |
| 440 | auto entry = findInCache(lock, kj::mv(key), options); |
| 441 | switch (entry->getValueStatus()) { |
| 442 | case EntryValueStatus::PRESENT: |
| 443 | case EntryValueStatus::ABSENT: { |
| 444 | return entry->getValue(); |
| 445 | } |
| 446 | case EntryValueStatus::UNKNOWN: { |
| 447 | return getImpl(kj::mv(entry), options); |
| 448 | } |
| 449 | } |
| 450 | } |
| 451 | |
| 452 | auto ActorCache::getImpl( |
| 453 | kj::Own<Entry> entry, ReadOptions options) -> kj::Promise<kj::Maybe<Value>> { |
| 454 | auto response = co_await scheduleStorageRead( |
| 455 | [key = entry->key.asBytes()](rpc::ActorStorage::Operations::Client client) { |
| 456 | auto req = client.getRequest(capnp::MessageSize{4 + key.size() / sizeof(capnp::word), 0}); |
| 457 | req.setKey(key); |
| 458 | return req.send().dropPipeline(); |
| 459 | }); |
| 460 | |
| 461 | kj::Maybe<capnp::Data::Reader> value; |
| 462 | if (response.hasValue()) { |
| 463 | value = response.getValue(); |
| 464 | } |
| 465 | auto lock = lru.cleanList.lockExclusive(); |
| 466 | auto newEntry = addReadResultToCache(lock, cloneKey(entry->key), value, options); |
| 467 | evictOrOomIfNeeded(lock); |
| 468 | co_return newEntry->getValue(); |
| 469 | } |
| 470 | |
| 471 | class ActorCache::GetMultiStreamImpl final: public rpc::ActorStorage::ListStream::Server { |
| 472 | public: |
| 473 | GetMultiStreamImpl(ActorCache& cache, |
| 474 | kj::Vector<kj::Own<Entry>> cachedEntries, |
| 475 | kj::Vector<Key> keysToFetchParam, |
| 476 | kj::Own<kj::PromiseFulfiller<GetResultList>> fulfiller, |
| 477 | const ReadOptions& options) |
| 478 | : cache(cache), |
| 479 | cachedEntries(kj::mv(cachedEntries)), |
| 480 | keysToFetch(kj::mv(keysToFetchParam)), |
| 481 | nextExpectedKey(keysToFetch.begin()), |
| 482 | fulfiller(kj::mv(fulfiller)), |
| 483 | options(options) {} |
| 484 | |
| 485 | kj::Promise<void> values(ValuesContext context) override { |
| 486 | if (!fulfiller->isWaiting()) { |
| 487 | // The original caller stopped listening. Try to cancel the stream by throwing. |
| 488 | return KJ_EXCEPTION(DISCONNECTED, "canceled"); |
| 489 | } |
| 490 | |
| 491 | auto lock = cache.lru.cleanList.lockExclusive(); |
| 492 | auto params = context.getParams(); |
| 493 | kj::String prevKey; |
| 494 | for (auto kv: params.getList()) { |
| 495 | KJ_ASSERT(kv.hasValue()); // values that don't exist aren't listed! |
| 496 | KJ_ASSERT(nextExpectedKey != keysToFetch.end()); |
| 497 | |
| 498 | // TODO(perf): This copy of the key is not really needed, we use the key from `keysToFetch` |
| 499 | // instead. But the capnp representation is a byte array which isn't null-terminated |
| 500 | // which would make the code difficult below. |
| 501 | auto key = kj::str(kv.getKey().asChars()); |
| 502 | |
| 503 | KJ_ASSERT(key >= prevKey, "storage returned keys in non-sorted order?"); |
| 504 | |
| 505 | // Find matching key in keysToFetch, possibly marking missing keys as absent. |
| 506 | for (;;) { |
| 507 | if (nextExpectedKey == keysToFetch.end() || key < *nextExpectedKey) { |
| 508 | // This may be a duplicate due to a retry. Ignore it. |
| 509 | break; |
| 510 | } else if (key == *nextExpectedKey) { |
| 511 | fetchedEntries.add( |
| 512 | cache.addReadResultToCache(lock, kj::mv(*nextExpectedKey), kv.getValue(), options)); |
| 513 | ++nextExpectedKey; |
| 514 | break; |
| 515 | } |
| 516 | |
| 517 | // It seems the list results have moved past `nextExpectedKey`, meaning it wasn't present |
| 518 | // on disk. Write a negative cache entry. |
| 519 | cache.addReadResultToCache(lock, kj::mv(*nextExpectedKey), kj::none, options); |
| 520 | ++nextExpectedKey; |
| 521 | } |
| 522 | |
| 523 | if (nextExpectedKey == keysToFetch.end()) { |
| 524 | fulfill(); |
| 525 | } |
| 526 | |
| 527 | prevKey = kj::mv(key); |
| 528 | } |
| 529 | cache.evictOrOomIfNeeded(lock); |
| 530 | return kj::READY_NOW; |
| 531 | } |
| 532 | |
| 533 | kj::Promise<void> end(EndContext context) override { |
| 534 | if (!fulfiller->isWaiting()) { |
| 535 | // Just ignore end() if we've already stopped waiting. |
| 536 | return kj::READY_NOW; |
| 537 | } |
| 538 | |
| 539 | if (nextExpectedKey < keysToFetch.end()) { |
| 540 | // Some trailing keys weren't seen, better mark them as not present. |
| 541 | auto lock = cache.lru.cleanList.lockExclusive(); |
| 542 | while (nextExpectedKey < keysToFetch.end()) { |
| 543 | cache.addReadResultToCache(lock, kj::mv(*nextExpectedKey++), kj::none, options); |
| 544 | } |
| 545 | cache.evictOrOomIfNeeded(lock); |
| 546 | } |
| 547 | |
| 548 | fulfill(); |
| 549 | |
| 550 | return kj::READY_NOW; |
| 551 | } |
| 552 | |
| 553 | void fulfill() { |
| 554 | // We return results in sorted order. You might argue that it could make sense to return |
| 555 | // results in the same order as the keys were originally specified. Even though we return |
| 556 | // a `Map` in JavaScript, the iteration order of a `Map` is defined to be the order of |
| 557 | // insertion, therefore the order in which we return results here is actually observable by |
| 558 | // the application. Trying to match the input order, however, almost certainly wouldn't be |
| 559 | // useful to apps. The only plausible way it could be useful is if the app could do e.g. |
| 560 | // `[...map.values()]` and end up with an array of values that exactly corresponds to the |
| 561 | // input array of keys. However, it won't exactly correspond for two reasons: |
| 562 | // - Keys that weren't present on disk aren't listed at all. To meaningfully change this, |
| 563 | // we would need to say that the Map object returned to JavaScript would contain entries |
| 564 | // even for missing keys, where the value is explicitly set to `undefined`. However, |
| 565 | // changing that would be a breaking change. |
| 566 | // - Keys that were listed twice in the input list won't be reported twice. This is an |
| 567 | // inherent limitation of the fact that we return a `Map`. |
| 568 | // |
| 569 | // Hence, applications that tried to depend on this ordering would be shooting themselves |
| 570 | // in the foot. We do, however, want to produce a consistent ordering for reproducibility's |
| 571 | // sake, but any consistent ordering will due. Sorted order is as good as anything else, and |
| 572 | // happens to be nice and easy for us. |
| 573 | fulfiller->fulfill( |
| 574 | GetResultList(kj::mv(cachedEntries), kj::mv(fetchedEntries), GetResultList::FORWARD)); |
| 575 | } |
| 576 | |
| 577 | // Indicates that the operation is being canceled. Proactively drops all entries. This |
| 578 | // is important because the destructor of an `Entry` updates the cache's accounting of memory |
| 579 | // usage, so it's important that an `Entry` cannot be held beyond the lifetime of the cache |
| 580 | // itself. |
| 581 | void cancel() { |
| 582 | KJ_ASSERT(!fulfiller->isWaiting()); // proves further RPCs will be ignored |
| 583 | cachedEntries.clear(); |
| 584 | fetchedEntries.clear(); |
| 585 | } |
| 586 | |
| 587 | ActorCache& cache; |
| 588 | kj::Vector<kj::Own<Entry>> cachedEntries; |
| 589 | kj::Vector<kj::Own<Entry>> fetchedEntries; |
| 590 | kj::Vector<Key> keysToFetch; |
| 591 | Key* nextExpectedKey; |
| 592 | kj::Own<kj::PromiseFulfiller<GetResultList>> fulfiller; |
| 593 | ReadOptions options; |
| 594 | }; |
| 595 | |
| 596 | kj::OneOf<ActorCache::GetResultList, kj::Promise<ActorCache::GetResultList>> ActorCache::get( |
| 597 | kj::Array<Key> keys, ReadOptions options) { |
| 598 | ActorStorageLimits::checkMaxPairsCount(keys.size()); |
| 599 | |
| 600 | options.noCache = options.noCache || lru.options.noCache; |
| 601 | requireNotTerminal(nullptr); |
| 602 | |
| 603 | std::sort(keys.begin(), keys.end()); |
| 604 | |
| 605 | kj::Vector<kj::Own<Entry>> cachedEntries(keys.size()); |
| 606 | // Entries satisfying the requested keys. |
| 607 | |
| 608 | kj::Vector<Key> keysToFetch(keys.size()); |
| 609 | // Keys that were not satisfied from cache. |
| 610 | |
| 611 | capnp::MessageSize sizeHint{4, 1}; |
| 612 | |
| 613 | { |
| 614 | auto lock = lru.cleanList.lockExclusive(); |
| 615 | for (auto& key: keys) { |
| 616 | auto entry = findInCache(lock, key, options); |
| 617 | switch (entry->getValueStatus()) { |
| 618 | case EntryValueStatus::PRESENT: |
| 619 | case EntryValueStatus::ABSENT: { |
| 620 | cachedEntries.add(kj::mv(entry)); |
| 621 | break; |
| 622 | } |
| 623 | case EntryValueStatus::UNKNOWN: { |
| 624 | // +1 word for padding, +1 word for the pointer in the key list. |
| 625 | sizeHint.wordCount += key.size() / sizeof(capnp::word) + 2; |
| 626 | keysToFetch.add(kj::mv(key)); |
| 627 | } |
| 628 | } |
| 629 | } |
| 630 | } |
| 631 | |
| 632 | if (keysToFetch.empty()) { |
| 633 | // All satisfied, return early. |
| 634 | return GetResultList(kj::mv(cachedEntries), {}, GetResultList::FORWARD); |
| 635 | } |
| 636 | |
| 637 | auto paf = kj::newPromiseAndFulfiller<GetResultList>(); |
| 638 | auto streamServer = kj::heap<GetMultiStreamImpl>( |
| 639 | *this, kj::mv(cachedEntries), kj::mv(keysToFetch), kj::mv(paf.fulfiller), options); |
| 640 | auto& streamServerRef = *streamServer; |
| 641 | |
| 642 | rpc::ActorStorage::ListStream::Client streamClient = kj::mv(streamServer); |
| 643 | |
| 644 | auto sendPromise = scheduleStorageRead( |
| 645 | [sizeHint, streamClient, &streamServerRef]( |
| 646 | rpc::ActorStorage::Operations::Client client) mutable -> kj::Promise<void> { |
| 647 | if (streamServerRef.nextExpectedKey == streamServerRef.keysToFetch.end()) { |
| 648 | // No more keys expected, must have finished listing on a previous try. |
| 649 | return kj::READY_NOW; |
| 650 | } |
| 651 | auto req = client.getMultipleRequest(sizeHint); |
| 652 | auto keysToFetch = |
| 653 | kj::arrayPtr(streamServerRef.nextExpectedKey, streamServerRef.keysToFetch.end()); |
| 654 | auto list = req.initKeys(keysToFetch.size()); |
| 655 | for (auto i: kj::indices(keysToFetch)) { |
| 656 | list.set(i, keysToFetch[i].asBytes()); |
| 657 | } |
| 658 | req.setStream(streamClient); |
| 659 | return req.sendIgnoringResult(); |
| 660 | }); |
| 661 | |
| 662 | // Wait on the RPC only until stream.end() is called, then report the results. We prevent |
| 663 | // `stream` from being destroyed until we have a result so that if the RPC throws an exception, |
| 664 | // we don't accidentally report "PromiseFulfiller not fulfilled" instead of the exception. |
| 665 | auto promise = sendPromise.then([&streamServerRef]() -> kj::Promise<ActorCache::GetResultList> { |
| 666 | if (streamServerRef.fulfiller->isWaiting()) { |
| 667 | return KJ_EXCEPTION(FAILED, "getMultiple() never called stream.end()"); |
| 668 | } else { |
| 669 | // We'll be canceled momentarily... |
| 670 | return kj::NEVER_DONE; |
| 671 | } |
| 672 | }); |
| 673 | return paf.promise.exclusiveJoin(kj::mv(promise)) |
| 674 | .attach(kj::defer( |
| 675 | [client = kj::mv(streamClient), &streamServerRef]() { streamServerRef.cancel(); })); |
| 676 | } |
| 677 | |
| 678 | kj::OneOf<kj::Maybe<kj::Date>, kj::Promise<kj::Maybe<kj::Date>>> ActorCache::getAlarm( |
| 679 | ReadOptions options) { |
| 680 | options.noCache = options.noCache || lru.options.noCache; |
| 681 | |
| 682 | // If in cache return time |
| 683 | // Else schedule alarm read |
| 684 | KJ_SWITCH_ONEOF(currentAlarmTime) { |
| 685 | KJ_CASE_ONEOF(entry, ActorCache::DeferredAlarmDelete) { |
| 686 | // An alarm handler is currently running, and a new alarm time has not been set yet. |
| 687 | // We need to return that there is no alarm. |
| 688 | return kj::Maybe<kj::Date>(kj::none); |
| 689 | } |
| 690 | KJ_CASE_ONEOF(entry, ActorCache::KnownAlarmTime) { |
| 691 | return entry.time; |
| 692 | } |
| 693 | KJ_CASE_ONEOF(_, ActorCache::UnknownAlarmTime) { |
| 694 | return scheduleStorageRead([](rpc::ActorStorage::Operations::Client client) { |
| 695 | auto req = client.getAlarmRequest(); |
| 696 | return req.send().dropPipeline(); |
| 697 | }) |
| 698 | .then([this, options](capnp::Response<rpc::ActorStorage::Operations::GetAlarmResults> |
| 699 | response) mutable -> kj::Maybe<kj::Date> { |
| 700 | auto scheduledTimeMs = response.getScheduledTimeMs(); |
| 701 | auto result = [&]() -> kj::Maybe<kj::Date> { |
| 702 | if (scheduledTimeMs == 0) { |
| 703 | return kj::none; |
| 704 | } else { |
| 705 | return scheduledTimeMs * kj::MILLISECONDS + kj::UNIX_EPOCH; |
| 706 | } |
| 707 | }(); |
| 708 | |
| 709 | if (!options.noCache && currentAlarmTime.is<UnknownAlarmTime>()) { |
| 710 | // If we don't end up in this branch, the time that's already in currentAlarmTime must |
| 711 | // be at least as fresh as the one we just read. |
| 712 | // |
| 713 | // If it was created by a setAlarm(), then it is actually fresher. If it was created |
| 714 | // by a concurrent getAlarm(), then it should be exactly the same time. |
| 715 | |
| 716 | currentAlarmTime = |
| 717 | ActorCache::KnownAlarmTime{ActorCache::KnownAlarmTime::Status::CLEAN, result}; |
| 718 | } |
| 719 | |
| 720 | return result; |
| 721 | }); |
| 722 | } |
| 723 | } |
| 724 | |
| 725 | KJ_UNREACHABLE; |
| 726 | } |
| 727 | |
| 728 | // ----------------------------------------------------------------------------- |
| 729 | |
| 730 | namespace { |
| 731 | // To simplify the handling of Maybe<Key> representing the end point of a list range, we define |
| 732 | // these operators to allow comparison between a Key and a Maybe<Key>, where a null Maybe<Key> |
| 733 | // sorts after all other keys. |
| 734 | |
| 735 | inline bool operator==(const ActorCache::Key& a, const kj::Maybe<ActorCache::Key>& b) { |
| 736 | KJ_IF_SOME(bb, b) { |
| 737 | return a == bb; |
| 738 | } else { |
| 739 | return false; |
| 740 | } |
| 741 | } |
| 742 | inline bool operator<(const ActorCache::Key& a, const kj::Maybe<ActorCache::Key>& b) { |
| 743 | KJ_IF_SOME(bb, b) { |
| 744 | return a < bb; |
| 745 | } else { |
| 746 | return true; |
| 747 | } |
| 748 | } |
| 749 | inline bool operator>=(const ActorCache::Key& a, const kj::Maybe<ActorCache::Key>& b) { |
| 750 | KJ_IF_SOME(bb, b) { |
| 751 | return a >= bb; |
| 752 | } else { |
| 753 | return false; |
| 754 | } |
| 755 | } |
| 756 | inline bool operator>(const ActorCache::Key& a, const kj::Maybe<ActorCache::KeyPtr>& b) { |
| 757 | KJ_IF_SOME(bb, b) { |
| 758 | return a > bb; |
| 759 | } else { |
| 760 | return false; |
| 761 | } |
| 762 | } |
| 763 | |
| 764 | inline auto seekOrEnd(auto& map, kj::Maybe<ActorCache::KeyPtr> key) { |
| 765 | KJ_IF_SOME(k, key) { |
| 766 | return map.seek(k); |
| 767 | } else { |
| 768 | return map.ordered().end(); |
| 769 | } |
| 770 | } |
| 771 | |
| 772 | } // namespace |
| 773 | |
| 774 | class ActorCache::ForwardListStreamImpl final: public rpc::ActorStorage::ListStream::Server { |
| 775 | public: |
| 776 | ForwardListStreamImpl(ActorCache& cache, |
| 777 | Key beginKey, |
| 778 | kj::Maybe<Key> endKey, |
| 779 | kj::Vector<kj::Own<Entry>> cachedEntries, |
| 780 | kj::Own<kj::PromiseFulfiller<GetResultList>> fulfiller, |
| 781 | kj::Maybe<uint> originalLimit, |
| 782 | kj::Maybe<uint> adjustedLimit, |
| 783 | bool beginKeyIsKnown, |
| 784 | const ReadOptions& options) |
| 785 | : cache(cache), |
| 786 | beginKey(kj::mv(beginKey)), |
| 787 | endKey(kj::mv(endKey)), |
| 788 | cachedEntries(kj::mv(cachedEntries)), |
| 789 | fulfiller(kj::mv(fulfiller)), |
| 790 | originalLimit(originalLimit), |
| 791 | adjustedLimit(adjustedLimit), |
| 792 | beginKeyIsKnown(beginKeyIsKnown), |
| 793 | options(options) {} |
| 794 | |
| 795 | kj::Promise<void> values(ValuesContext context) override { |
| 796 | if (!fulfiller->isWaiting()) { |
| 797 | // The original caller stopped listening. Try to cancel the stream by throwing. |
| 798 | return KJ_EXCEPTION(DISCONNECTED, "canceled"); |
| 799 | } |
| 800 | |
| 801 | { |
| 802 | auto lock = cache.lru.cleanList.lockExclusive(); |
| 803 | auto list = context.getParams().getList(); |
| 804 | |
| 805 | bool insertedAny = false; |
| 806 | |
| 807 | for (auto kv: list) { |
| 808 | Key key = kj::str(kv.getKey().asChars()); |
| 809 | |
| 810 | if (!beginKeyIsKnown) { |
| 811 | if (key != beginKey) { |
| 812 | // This is the first set of results we've received, and it does not include the start |
| 813 | // point of the list. Therefore, we should insert an entry with a null value, to make |
| 814 | // sure the whole range can be marked as empty. We'll end up marking this entry as |
| 815 | // part of markGapsEmpty(), later. |
| 816 | markBeginAsEmpty(lock); |
| 817 | } |
| 818 | } else { |
| 819 | if (key <= beginKey) { |
| 820 | // Out-of-order result. This is probably the result of restarting the list operation |
| 821 | // due to a disconnect. We assume this is actually a duplicate of a result we |
| 822 | // received earlier. Ignore it. |
| 823 | continue; |
| 824 | } |
| 825 | } |
| 826 | |
| 827 | KJ_ASSERT(kv.hasValue()); // values that don't exist aren't listed! |
| 828 | auto entry = cache.addReadResultToCache(lock, kj::mv(key), kv.getValue(), options); |
| 829 | fetchedEntries.add(kj::mv(entry)); |
| 830 | insertedAny = true; |
| 831 | } |
| 832 | |
| 833 | if (insertedAny) { |
| 834 | // Update `gapIsKnownEmpty` on the whole range. |
| 835 | cache.markGapsEmpty(lock, beginKey, fetchedEntries.back()->key.asPtr(), options); |
| 836 | beginKey = cloneKey(fetchedEntries.back()->key); |
| 837 | beginKeyIsKnown = true; |
| 838 | } |
| 839 | |
| 840 | cache.evictOrOomIfNeeded(lock); |
| 841 | } |
| 842 | |
| 843 | if (fetchedEntries.size() >= adjustedLimit.orDefault(kj::maxValue)) { |
| 844 | // Oh we're already done. |
| 845 | fulfill(); |
| 846 | } |
| 847 | return kj::READY_NOW; |
| 848 | } |
| 849 | |
| 850 | kj::Promise<void> end(EndContext context) override { |
| 851 | if (!fulfiller->isWaiting()) { |
| 852 | // Just ignore end() if we've already stopped waiting. In particular this happens in |
| 853 | // limit requests that reach the limit -- the last call to values() will have already |
| 854 | // fulfilled the fulfiller. |
| 855 | return kj::READY_NOW; |
| 856 | } |
| 857 | |
| 858 | // Mark the rest of the range as empty. |
| 859 | { |
| 860 | auto lock = cache.lru.cleanList.lockExclusive(); |
| 861 | |
| 862 | if (!beginKeyIsKnown) { |
| 863 | // We received no results at all, so the start of the list is definitely not in storage. |
| 864 | markBeginAsEmpty(lock); |
| 865 | } |
| 866 | |
| 867 | if (fetchedEntries.size() < adjustedLimit.orDefault(kj::maxValue)) { |
| 868 | // We didn't reach the limit, so the rest of the range must be empty. |
| 869 | cache.markGapsEmpty(lock, beginKey, endKey, options); |
| 870 | } |
| 871 | |
| 872 | cache.evictOrOomIfNeeded(lock); |
| 873 | } |
| 874 | |
| 875 | fulfill(); |
| 876 | |
| 877 | return kj::READY_NOW; |
| 878 | } |
| 879 | |
| 880 | void fulfill() { |
| 881 | fulfiller->fulfill(GetResultList( |
| 882 | kj::mv(cachedEntries), kj::mv(fetchedEntries), GetResultList::FORWARD, originalLimit)); |
| 883 | }; |
| 884 | |
| 885 | // Mark the start of the list operation will a null entry, because we did not see it listed. |
| 886 | // |
| 887 | // Note that this insertion attempt will be ignored in two cases: |
| 888 | // 1. An entry already exists with this key, perhaps as the result of a put(). This is |
| 889 | // fine, because the existing entry means we have something to mark. |
| 890 | // 2. The entry doesn't exist, but the previous entry has `gapIsKnownEmpty = true`, and |
| 891 | // so the insertion of a new null entry is ignored for being redundant. This case is |
| 892 | // fine too, as the gap is already marked. Our markGapsEmpty() call will start with the |
| 893 | // following entry. |
| 894 | void markBeginAsEmpty(Lock& lock) { |
| 895 | cache.addReadResultToCache(lock, cloneKey(beginKey), kj::none, options); |
| 896 | } |
| 897 | |
| 898 | // Indicates that the operation is being canceled. Proactively drops all entries. This |
| 899 | // is important because the destructor of an `Entry` updates the cache's accounting of memory |
| 900 | // usage, so it's important that an `Entry` cannot be held beyond the lifetime of the cache |
| 901 | // itself. |
| 902 | void cancel() { |
| 903 | KJ_ASSERT(!fulfiller->isWaiting()); // proves further RPCs will be ignored |
| 904 | cachedEntries.clear(); |
| 905 | fetchedEntries.clear(); |
| 906 | } |
| 907 | |
| 908 | ActorCache& cache; |
| 909 | |
| 910 | // Either: |
| 911 | // - No prefix of the list is known yet, and `beginKey` is the original begin point passed to |
| 912 | // list(). |
| 913 | // - Some prefix is already satisfied, either from cache or from a previous batch of results |
| 914 | // streamed from storage, and `beginKey` is the key of the last known entry in this prefix. |
| 915 | Key beginKey; |
| 916 | |
| 917 | // The end of the list range, as originally passed to list(). |
| 918 | kj::Maybe<Key> endKey; |
| 919 | |
| 920 | // Entries we gathered from cache. |
| 921 | kj::Vector<kj::Own<Entry>> cachedEntries; |
| 922 | |
| 923 | // Entries that have streamed in from disk. |
| 924 | kj::Vector<kj::Own<Entry>> fetchedEntries; |
| 925 | |
| 926 | // Fulfiller for the final results. |
| 927 | kj::Own<kj::PromiseFulfiller<GetResultList>> fulfiller; |
| 928 | |
| 929 | // The original requested limit, if any. |
| 930 | kj::Maybe<uint> originalLimit; |
| 931 | |
| 932 | // The limit we sent to storage. |
| 933 | kj::Maybe<uint> adjustedLimit; |
| 934 | |
| 935 | // Does `beginKey` point to a key where we already know the associated value? This is |
| 936 | // especially true when `beginKey` points to the last entry of a previous batch received via |
| 937 | // a call to `values()`. |
| 938 | bool beginKeyIsKnown; |
| 939 | |
| 940 | ReadOptions options; |
| 941 | }; |
| 942 | |
| 943 | kj::OneOf<ActorCache::GetResultList, kj::Promise<ActorCache::GetResultList>> ActorCache::list( |
| 944 | Key beginKey, kj::Maybe<Key> endKey, kj::Maybe<uint> limit, ReadOptions options) { |
| 945 | options.noCache = options.noCache || lru.options.noCache; |
| 946 | requireNotTerminal(nullptr); |
| 947 | |
| 948 | // We start by scanning the cache for entries satisfying the list range. If we can fully satisfy |
| 949 | // the list using these, then we're done! Otherwise, we make a storage request to get the rest. |
| 950 | // When the storage request produces results, we must discard any that conflict with what was |
| 951 | // in cache before hand, since what's in cache could have come from a put() that wasn't flushed |
| 952 | // yet. However, we need to be careful NOT to use any entries that were put() *after* the list() |
| 953 | // operation started. |
| 954 | |
| 955 | kj::Vector<kj::Own<Entry>> cachedEntries; |
| 956 | size_t positiveCount = 0; // number of positive entries in `cachedEntries` |
| 957 | if (limit.orDefault(kj::maxValue) == 0 || beginKey >= endKey) { |
| 958 | // No results in these cases, just return. |
| 959 | return ActorCache::GetResultList(kj::mv(cachedEntries), {}, GetResultList::FORWARD); |
| 960 | } |
| 961 | |
| 962 | uint limitAdjustment = 0; |
| 963 | // When requesting to storage, we need to adjust the limit to increase it by the number of cached |
| 964 | // negative entries in the range, since each of those negative entries could potentially negate a |
| 965 | // positive entry read from disk. |
| 966 | |
| 967 | auto lock = lru.cleanList.lockExclusive(); |
| 968 | auto& map = currentValues.get(lock); |
| 969 | auto ordered = map.ordered(); |
| 970 | |
| 971 | kj::Maybe<KeyPtr> storageListStart; |
| 972 | // If we must do a storage operation, what key shall it start at? |
| 973 | // |
| 974 | // Note that we never do more than one storage operation, even if we have a patchwork of cache |
| 975 | // entries matching different subsets of the list. Trying to split the operation into multiple |
| 976 | // smaller list operations to avoid re-listing things we already know seems like too much work to |
| 977 | // be worth it. So, we only track the first key which we know needs to be listed, and then we |
| 978 | // list the rest of the space from there. |
| 979 | |
| 980 | bool storageListStartIsKnown = false; |
| 981 | // Does `storageListStart` point to a key for which we already know the value? If so we can |
| 982 | // avoid listing that key specifically. |
| 983 | |
| 984 | uint knownPrefixSize = 0; |
| 985 | // How many keys were matched from cache before (and not including) `storageListStart`? We will |
| 986 | // use this to reduce the `limit` we pass in the storage op (if there is one). |
| 987 | |
| 988 | // Let's iterate over the cache starting from `beginKey`. |
| 989 | auto iter = map.seek(beginKey); |
| 990 | |
| 991 | // We need some special logic to handle the starting point with regard to gaps. |
| 992 | if (iter != ordered.end() && iter->get()->key == beginKey) { |
| 993 | // There is an entry specifically for `beginKey`, so we'll start there. |
| 994 | } else { |
| 995 | // `beginKey` does not match an entry, but we can check if it is in a known-empty gap. |
| 996 | if (iter == ordered.begin()) { |
| 997 | // No, because there is no previous entry. Oh well. We will have to start the storage list |
| 998 | // from `beginKey`. |
| 999 | storageListStart = beginKey; |
| 1000 | storageListStartIsKnown = false; |
| 1001 | } else { |
| 1002 | // There is a previous key in cache, let's take a look. |
| 1003 | auto prev = iter; |
| 1004 | --prev; |
| 1005 | if (prev->get()->gapIsKnownEmpty) { |
| 1006 | // `beginKey` is in a known-empty gap, so we know that this key simply doesn't exist in |
| 1007 | // storage. |
| 1008 | } else { |
| 1009 | // We don't know if `beginKey` exists in storage so we'll have to start the storage list |
| 1010 | // there. |
| 1011 | storageListStart = beginKey; |
| 1012 | storageListStartIsKnown = false; |
| 1013 | } |
| 1014 | } |
| 1015 | } |
| 1016 | |
| 1017 | // Now we can start scanning normally. We need to scan entries within the list range to build |
| 1018 | // a list of possible results, as well as to determine whether we need to do a storage request. |
| 1019 | // Even if we end up having to go to disk to find more data, we don't need to scan more than |
| 1020 | // `limit` entries from cache because any entries beyond that couldn't possibly end up in the |
| 1021 | // final results anyway. |
| 1022 | // |
| 1023 | // Note that we must keep scanning the cache *even if* we've seen an empty gap and |
| 1024 | // `storageListStart` is non-null. This is because our results must include recent put()s, which |
| 1025 | // may still be DIRTY so won't be returned when we list the database. Later on we'll merge the |
| 1026 | // entries we find in cache with those we get from disk. |
| 1027 | for (; iter != ordered.end() && iter->get()->key < endKey && |
| 1028 | positiveCount < limit.orDefault(kj::maxValue); |
| 1029 | ++iter) { |
| 1030 | Entry& entry = **iter; |
| 1031 | |
| 1032 | if (!options.noCache) { |
| 1033 | touchEntry(lock, entry); |
| 1034 | } |
| 1035 | |
| 1036 | switch (entry.getValueStatus()) { |
| 1037 | case EntryValueStatus::ABSENT: { |
| 1038 | cachedEntries.add(kj::atomicAddRef(entry)); |
| 1039 | if (storageListStart != kj::none && entry.isDirty()) { |
| 1040 | // This negative entry could negate something read from storage later, so we need to |
| 1041 | // increase the storage list limit. |
| 1042 | ++limitAdjustment; |
| 1043 | } |
| 1044 | break; |
| 1045 | } |
| 1046 | case EntryValueStatus::PRESENT: { |
| 1047 | cachedEntries.add(kj::atomicAddRef(entry)); |
| 1048 | ++positiveCount; |
| 1049 | if (storageListStart == kj::none) { |
| 1050 | ++knownPrefixSize; |
| 1051 | } |
| 1052 | break; |
| 1053 | } |
| 1054 | case EntryValueStatus::UNKNOWN: { |
| 1055 | // Ignore entry that exists only to mark a previous list range. |
| 1056 | break; |
| 1057 | } |
| 1058 | } |
| 1059 | |
| 1060 | if (storageListStart == kj::none && !entry.gapIsKnownEmpty) { |
| 1061 | // The gap after this entry is not cached so we'll have to start our list operation here. |
| 1062 | storageListStart = entry.key; |
| 1063 | storageListStartIsKnown = entry.getValueStatus() != EntryValueStatus::UNKNOWN; |
| 1064 | } |
| 1065 | } |
| 1066 | |
| 1067 | if (iter != ordered.end() && iter->get()->key == endKey) { |
| 1068 | // We have an entry exactly at our end, it might even be a previously inserted UNKNOWN. Let's |
| 1069 | // touch it for freshness. |
| 1070 | if (!options.noCache) { |
| 1071 | touchEntry(lock, **iter); |
| 1072 | } |
| 1073 | } |
| 1074 | |
| 1075 | if (storageListStart == kj::none || knownPrefixSize >= limit.orDefault(kj::maxValue)) { |
| 1076 | // We fully satisfied the list operation from cache. |
| 1077 | return GetResultList(kj::mv(cachedEntries), {}, GetResultList::FORWARD, limit); |
| 1078 | } |
| 1079 | |
| 1080 | auto adjustedLimit = |
| 1081 | limit.map([&](uint orig) { return orig + limitAdjustment - knownPrefixSize; }); |
| 1082 | |
| 1083 | auto paf = kj::newPromiseAndFulfiller<GetResultList>(); |
| 1084 | auto streamServer = kj::heap<ForwardListStreamImpl>(*this, |
| 1085 | cloneKey(KJ_ASSERT_NONNULL(storageListStart)), kj::mv(endKey), kj::mv(cachedEntries), |
| 1086 | kj::mv(paf.fulfiller), limit, adjustedLimit, storageListStartIsKnown, options); |
| 1087 | auto& streamServerRef = *streamServer; |
| 1088 | |
| 1089 | rpc::ActorStorage::ListStream::Client streamClient = kj::mv(streamServer); |
| 1090 | |
| 1091 | auto sendPromise = scheduleStorageRead( |
| 1092 | [&streamServerRef, streamClient]( |
| 1093 | rpc::ActorStorage::Operations::Client client) mutable -> kj::Promise<void> { |
| 1094 | auto req = client.listRequest( |
| 1095 | capnp::MessageSize{8 + streamServerRef.beginKey.size() / sizeof(capnp::word) + |
| 1096 | streamServerRef.endKey.map([](KeyPtr k) { |
| 1097 | return k.size() / sizeof(capnp::word); |
| 1098 | }).orDefault(0), |
| 1099 | 1}); |
| 1100 | |
| 1101 | if (streamServerRef.beginKeyIsKnown) { |
| 1102 | // `streamServerRef.beginKey` points to a key for which we already know the value, either |
| 1103 | // because it was already in cache when we started, or because we are retrying and a previous |
| 1104 | // call to `values()` produced this key. Querying it again would be redundant. But, list |
| 1105 | // operations are inclusive of the start key. So, we compute the successor of the start key, |
| 1106 | // which is the key with a zero byte appended. |
| 1107 | auto buffer = req.initStart(streamServerRef.beginKey.size() + 1); |
| 1108 | memcpy(buffer.begin(), streamServerRef.beginKey.begin(), buffer.size() - 1); |
| 1109 | // Technically capnp is zero-initialized so this is redundant, but just for safety and |
| 1110 | // clarity... |
| 1111 | buffer[buffer.size() - 1] = 0; |
| 1112 | } else { |
| 1113 | if (streamServerRef.beginKey.size() > 0) { |
| 1114 | req.setStart(streamServerRef.beginKey.asBytes()); |
| 1115 | } |
| 1116 | } |
| 1117 | |
| 1118 | KJ_IF_SOME(e, streamServerRef.endKey) { |
| 1119 | req.setEnd(e.asBytes()); |
| 1120 | } |
| 1121 | |
| 1122 | KJ_IF_SOME(l, streamServerRef.adjustedLimit) { |
| 1123 | if (streamServerRef.fetchedEntries.size() >= l) { |
| 1124 | // Oh it turns out we actually satisfied the limit already so we don't actually have to |
| 1125 | // retry. The fulfiller would have already been fulfilled. |
| 1126 | return kj::READY_NOW; |
| 1127 | } |
| 1128 | req.setLimit(l - streamServerRef.fetchedEntries.size()); |
| 1129 | } |
| 1130 | |
| 1131 | req.setStream(streamClient); |
| 1132 | return req.sendIgnoringResult(); |
| 1133 | }); |
| 1134 | |
| 1135 | // Wait on the RPC only until stream.end() is called, then report the results. We prevent |
| 1136 | // `stream` from being destroyed until we have a result so that if the RPC throws an exception, |
| 1137 | // we don't accidentally report "PromiseFulfiller not fulfilled" instead of the exception. |
| 1138 | auto promise = sendPromise.then([&streamServerRef]() -> kj::Promise<ActorCache::GetResultList> { |
| 1139 | if (streamServerRef.fulfiller->isWaiting()) { |
| 1140 | return KJ_EXCEPTION(FAILED, "list() never called stream.end()"); |
| 1141 | } else { |
| 1142 | // We'll be canceled momentarily... |
| 1143 | return kj::NEVER_DONE; |
| 1144 | } |
| 1145 | }); |
| 1146 | |
| 1147 | return paf.promise.exclusiveJoin(kj::mv(promise)) |
| 1148 | .attach(kj::defer( |
| 1149 | [client = kj::mv(streamClient), &streamServerRef]() { streamServerRef.cancel(); })); |
| 1150 | } |
| 1151 | |
| 1152 | // ----------------------------------------------------------------------------- |
| 1153 | |
| 1154 | class ActorCache::ReverseListStreamImpl final: public rpc::ActorStorage::ListStream::Server { |
| 1155 | public: |
| 1156 | ReverseListStreamImpl(ActorCache& cache, |
| 1157 | Key beginKey, |
| 1158 | kj::Maybe<Key> endKey, |
| 1159 | kj::Vector<kj::Own<Entry>> cachedEntries, |
| 1160 | kj::Own<kj::PromiseFulfiller<GetResultList>> fulfiller, |
| 1161 | kj::Maybe<uint> originalLimit, |
| 1162 | kj::Maybe<uint> adjustedLimit, |
| 1163 | ReadOptions options) |
| 1164 | : cache(cache), |
| 1165 | beginKey(kj::mv(beginKey)), |
| 1166 | endKey(kj::mv(endKey)), |
| 1167 | cachedEntries(kj::mv(cachedEntries)), |
| 1168 | fulfiller(kj::mv(fulfiller)), |
| 1169 | originalLimit(originalLimit), |
| 1170 | adjustedLimit(adjustedLimit), |
| 1171 | options(options) {} |
| 1172 | |
| 1173 | kj::Promise<void> values(ValuesContext context) override { |
| 1174 | if (!fulfiller->isWaiting()) { |
| 1175 | // The original caller stopped listening. Try to cancel the stream by throwing. |
| 1176 | return KJ_EXCEPTION(DISCONNECTED, "canceled"); |
| 1177 | } |
| 1178 | |
| 1179 | { |
| 1180 | auto lock = cache.lru.cleanList.lockExclusive(); |
| 1181 | auto list = context.getParams().getList(); |
| 1182 | |
| 1183 | bool insertedAny = false; |
| 1184 | |
| 1185 | for (auto kv: list) { |
| 1186 | Key key = kj::str(kv.getKey().asChars()); |
| 1187 | |
| 1188 | if (key >= endKey) { |
| 1189 | // Out-of-order result. This is probably the result of restarting the list operation |
| 1190 | // due to a disconnect. We assume this is actually a duplicate of a result we |
| 1191 | // received earlier. Ignore it. |
| 1192 | continue; |
| 1193 | } |
| 1194 | |
| 1195 | KJ_ASSERT(kv.hasValue()); // values that don't exist aren't listed! |
| 1196 | auto entry = cache.addReadResultToCache(lock, kj::mv(key), kv.getValue(), options); |
| 1197 | fetchedEntries.add(kj::mv(entry)); |
| 1198 | insertedAny = true; |
| 1199 | } |
| 1200 | |
| 1201 | if (insertedAny) { |
| 1202 | // Update `gapIsKnownEmpty` on the whole range. |
| 1203 | cache.markGapsEmpty(lock, fetchedEntries.back()->key, endKey, options); |
| 1204 | endKey = cloneKey(fetchedEntries.back()->key); |
| 1205 | } |
| 1206 | |
| 1207 | cache.evictOrOomIfNeeded(lock); |
| 1208 | } |
| 1209 | |
| 1210 | if (fetchedEntries.size() >= adjustedLimit.orDefault(kj::maxValue) || beginKey == endKey) { |
| 1211 | // Oh we're already done. |
| 1212 | fulfill(); |
| 1213 | } |
| 1214 | return kj::READY_NOW; |
| 1215 | } |
| 1216 | |
| 1217 | kj::Promise<void> end(EndContext context) override { |
| 1218 | if (!fulfiller->isWaiting()) { |
| 1219 | // Just ignore end() if we've already stopped waiting. In particular this happens in |
| 1220 | // limit requests that reach the limit, or when we see an entry matching the beginning |
| 1221 | // key of the list range -- in both cases, the last call to values() will have already |
| 1222 | // fulfilled the fulfiller. |
| 1223 | return kj::READY_NOW; |
| 1224 | } |
| 1225 | |
| 1226 | // Mark the rest of the range as empty. |
| 1227 | { |
| 1228 | auto lock = cache.lru.cleanList.lockExclusive(); |
| 1229 | |
| 1230 | if (fetchedEntries.size() < adjustedLimit.orDefault(kj::maxValue)) { |
| 1231 | // We didn't reach the limit, so the rest of the range must be empty. |
| 1232 | |
| 1233 | // We may need to insert a negative entry at the beginning of the list range, since we |
| 1234 | // didn't see it, implying it's not present on disk. addResultToCache() will conveniently |
| 1235 | // avoid adding anything if it turns out this is already in a known-empty gap. |
| 1236 | auto beginEntry = cache.addReadResultToCache(lock, cloneKey(beginKey), kj::none, options); |
| 1237 | |
| 1238 | // And we need to mark gaps empty from there to the final entry we actually saw. |
| 1239 | cache.markGapsEmpty(lock, beginEntry->key, endKey, options); |
| 1240 | } |
| 1241 | |
| 1242 | cache.evictOrOomIfNeeded(lock); |
| 1243 | } |
| 1244 | |
| 1245 | fulfill(); |
| 1246 | |
| 1247 | return kj::READY_NOW; |
| 1248 | } |
| 1249 | |
| 1250 | void fulfill() { |
| 1251 | fulfiller->fulfill(GetResultList( |
| 1252 | kj::mv(cachedEntries), kj::mv(fetchedEntries), GetResultList::REVERSE, originalLimit)); |
| 1253 | } |
| 1254 | |
| 1255 | // Indicates that the operation is being canceled. Proactively drops all entries. This |
| 1256 | // is important because the destructor of an `Entry` updates the cache's accounting of memory |
| 1257 | // usage, so it's important that an `Entry` cannot be held beyond the lifetime of the cache |
| 1258 | // itself. |
| 1259 | void cancel() { |
| 1260 | KJ_ASSERT(!fulfiller->isWaiting()); // proves further RPCs will be ignored |
| 1261 | cachedEntries.clear(); |
| 1262 | fetchedEntries.clear(); |
| 1263 | } |
| 1264 | |
| 1265 | ActorCache& cache; |
| 1266 | |
| 1267 | // The beginning of the list range, as originally passed to list(). |
| 1268 | Key beginKey; |
| 1269 | |
| 1270 | // Either: |
| 1271 | // - No suffix of the list is known yet, and `endKey` is the original end point passed to |
| 1272 | // list(). |
| 1273 | // - Some suffix is already satisfied , either from cache or from a previous batch of results |
| 1274 | // streamed from storage, and `endKey` is the key of the first known entry in this suffix. |
| 1275 | kj::Maybe<Key> endKey; |
| 1276 | |
| 1277 | // Entries we gathered from cache. |
| 1278 | kj::Vector<kj::Own<Entry>> cachedEntries; |
| 1279 | |
| 1280 | // Entries that have streamed in from disk. |
| 1281 | kj::Vector<kj::Own<Entry>> fetchedEntries; |
| 1282 | |
| 1283 | // Fulfiller for the final results. |
| 1284 | kj::Own<kj::PromiseFulfiller<GetResultList>> fulfiller; |
| 1285 | |
| 1286 | // The original requested limit, if any. |
| 1287 | kj::Maybe<uint> originalLimit; |
| 1288 | |
| 1289 | // The limit we sent to storage. |
| 1290 | kj::Maybe<uint> adjustedLimit; |
| 1291 | |
| 1292 | ReadOptions options; |
| 1293 | }; |
| 1294 | |
| 1295 | kj::OneOf<ActorCache::GetResultList, kj::Promise<ActorCache::GetResultList>> ActorCache:: |
| 1296 | listReverse(Key beginKey, kj::Maybe<Key> endKey, kj::Maybe<uint> limit, ReadOptions options) { |
| 1297 | options.noCache = options.noCache || lru.options.noCache; |
| 1298 | requireNotTerminal(nullptr); |
| 1299 | |
| 1300 | // Alas, everything needs to be done slightly differently when listing in reverse. This function |
| 1301 | // is an adjusted version of the previous function. |
| 1302 | |
| 1303 | kj::Vector<kj::Own<Entry>> cachedEntries; |
| 1304 | size_t positiveCount = 0; // number of positive entries in `cachedEntries` |
| 1305 | if (limit.orDefault(kj::maxValue) == 0 || beginKey >= endKey) { |
| 1306 | // No results in these cases, just return. |
| 1307 | return ActorCache::GetResultList(kj::mv(cachedEntries), {}, GetResultList::REVERSE); |
| 1308 | } |
| 1309 | |
| 1310 | uint limitAdjustment = 0; |
| 1311 | // When requesting to storage, we need to adjust the limit to increase it by the number of cached |
| 1312 | // negative entries in the range, since each of those negative entries could potentially negate a |
| 1313 | // positive entry read from disk. |
| 1314 | |
| 1315 | auto lock = lru.cleanList.lockExclusive(); |
| 1316 | auto& map = currentValues.get(lock); |
| 1317 | auto ordered = map.ordered(); |
| 1318 | |
| 1319 | kj::Maybe<KeyPtr> storageListEnd; |
| 1320 | // If we must do a storage operation, what key shall it end at? |
| 1321 | // |
| 1322 | // As an extra hack, if the Maybe is non-null but the KeyPtr is null, this indicates there is |
| 1323 | // no end. It's impossible for storageListEnd to point at a null key and intend this to mean that |
| 1324 | // the end should be the empty-string key because this would suggest an empty list range. |
| 1325 | |
| 1326 | uint knownSuffixSize = 0; |
| 1327 | // How many keys were matched from cache after (and including) `storageListEnd`? We will |
| 1328 | // use this to reduce the `limit` we pass in the storage op (if there is one). |
| 1329 | |
| 1330 | // Let's iterate backwards over the cache starting from `endKey`. Iterating backwards is a |
| 1331 | // bit mind-bendy. |
| 1332 | // |
| 1333 | // Note that we must keep scanning the cache *even if* we've seen an empty gap and |
| 1334 | // `storageListEnd` is non-null. This is because our results must include recent put()s, which |
| 1335 | // may still be DIRTY so won't be returned when we list the database. Later on we'll merge the |
| 1336 | // entries we find in cache with those we get from disk. |
| 1337 | KeyPtr nextKey = endKey.orDefault({}); // "the last key we saw in backwards order" |
| 1338 | auto iter = seekOrEnd(map, endKey); |
| 1339 | if (iter != map.ordered().end() && iter->get()->key == endKey) { |
| 1340 | // We have an entry exactly at our end, it might even be a previously inserted UNKNOWN. Let's |
| 1341 | // touch it for freshness. |
| 1342 | if (!options.noCache) { |
| 1343 | touchEntry(lock, **iter); |
| 1344 | } |
| 1345 | } |
| 1346 | for (; positiveCount < limit.orDefault(kj::maxValue);) { |
| 1347 | if (iter == ordered.begin()) { |
| 1348 | // No earlier entries, treat same as if previous entry were before beginKey and had |
| 1349 | // gapIsKnownEmpty = false. |
| 1350 | if (storageListEnd == kj::none) storageListEnd = nextKey; |
| 1351 | break; |
| 1352 | } |
| 1353 | |
| 1354 | // Step backwards. |
| 1355 | --iter; |
| 1356 | auto& entry = **iter; |
| 1357 | |
| 1358 | // If the gap after this entry is not known empty, then we've exhausted our known-suffix and |
| 1359 | // will need to cover this gap using a storage RPC. |
| 1360 | if (storageListEnd == kj::none && !entry.gapIsKnownEmpty) { |
| 1361 | storageListEnd = nextKey; |
| 1362 | } |
| 1363 | |
| 1364 | if (entry.key < beginKey) { |
| 1365 | // We've traversed past the beginning of our range so exit the loop here. |
| 1366 | break; |
| 1367 | } |
| 1368 | |
| 1369 | if (!options.noCache) { |
| 1370 | touchEntry(lock, entry); |
| 1371 | } |
| 1372 | |
| 1373 | // Note that we need to add even negative entries to `cachedEntries` so that they override |
| 1374 | // whatever we read from storage later. However, they should not count against the limit. |
| 1375 | switch (entry.getValueStatus()) { |
| 1376 | case EntryValueStatus::ABSENT: { |
| 1377 | cachedEntries.add(kj::atomicAddRef(entry)); |
| 1378 | if (storageListEnd != kj::none && entry.isDirty()) { |
| 1379 | // This negative entry could negate something read from storage later, so we need to |
| 1380 | // increase the storage list limit. |
| 1381 | ++limitAdjustment; |
| 1382 | } |
| 1383 | break; |
| 1384 | } |
| 1385 | case EntryValueStatus::PRESENT: { |
| 1386 | cachedEntries.add(kj::atomicAddRef(entry)); |
| 1387 | ++positiveCount; |
| 1388 | if (storageListEnd == kj::none) { |
| 1389 | ++knownSuffixSize; |
| 1390 | } |
| 1391 | break; |
| 1392 | } |
| 1393 | case EntryValueStatus::UNKNOWN: { |
| 1394 | // Ignore entry that exists only to mark a previous list range. |
| 1395 | break; |
| 1396 | } |
| 1397 | } |
| 1398 | |
| 1399 | if (entry.key == beginKey) { |
| 1400 | // We've traversed through the beginning of our range so exit the loop here. |
| 1401 | break; |
| 1402 | } |
| 1403 | |
| 1404 | nextKey = entry.key; |
| 1405 | } |
| 1406 | |
| 1407 | if (storageListEnd == kj::none || knownSuffixSize >= limit.orDefault(kj::maxValue)) { |
| 1408 | // We fully satisfied the list operation from cache. |
| 1409 | return GetResultList(kj::mv(cachedEntries), {}, GetResultList::REVERSE, limit); |
| 1410 | } |
| 1411 | |
| 1412 | { |
| 1413 | KeyPtr k = KJ_ASSERT_NONNULL(storageListEnd); |
| 1414 | if (k.size() == 0) { |
| 1415 | // Empty string inside non-null storageListEnd means that our endpoint is the end of the |
| 1416 | // keyspace. (It couldn't possibly mean that our endpoint is the *beginning* of the keyspace, |
| 1417 | // because that would mean that we're listing a zero-sized range, in which case we would have |
| 1418 | // returned earlier.) |
| 1419 | endKey = kj::none; |
| 1420 | } else { |
| 1421 | endKey = cloneKey(k); |
| 1422 | } |
| 1423 | } |
| 1424 | |
| 1425 | auto adjustedLimit = |
| 1426 | limit.map([&](uint orig) { return orig + limitAdjustment - knownSuffixSize; }); |
| 1427 | |
| 1428 | auto paf = kj::newPromiseAndFulfiller<GetResultList>(); |
| 1429 | auto streamServer = kj::heap<ReverseListStreamImpl>(*this, kj::mv(beginKey), kj::mv(endKey), |
| 1430 | kj::mv(cachedEntries), kj::mv(paf.fulfiller), limit, adjustedLimit, options); |
| 1431 | auto& streamServerRef = *streamServer; |
| 1432 | |
| 1433 | rpc::ActorStorage::ListStream::Client streamClient = kj::mv(streamServer); |
| 1434 | |
| 1435 | auto sendPromise = scheduleStorageRead( |
| 1436 | [&streamServerRef, streamClient]( |
| 1437 | rpc::ActorStorage::Operations::Client client) mutable -> kj::Promise<void> { |
| 1438 | auto req = client.listRequest( |
| 1439 | capnp::MessageSize{8 + streamServerRef.beginKey.size() / sizeof(capnp::word) + |
| 1440 | streamServerRef.endKey.map([](KeyPtr k) { |
| 1441 | return k.size() / sizeof(capnp::word); |
| 1442 | }).orDefault(0), |
| 1443 | 1}); |
| 1444 | if (streamServerRef.beginKey.size() > 0) { |
| 1445 | req.setStart(streamServerRef.beginKey.asBytes()); |
| 1446 | } |
| 1447 | KJ_IF_SOME(e, streamServerRef.endKey) { |
| 1448 | req.setEnd(e.asBytes()); |
| 1449 | } |
| 1450 | req.setReverse(true); |
| 1451 | KJ_IF_SOME(l, streamServerRef.adjustedLimit) { |
| 1452 | if (streamServerRef.fetchedEntries.size() >= l) { |
| 1453 | // Oh it turns out we actually satisfied the limit already so we don't actually have to |
| 1454 | // retry. The fulfiller would have already been fulfilled. |
| 1455 | return kj::READY_NOW; |
| 1456 | } |
| 1457 | req.setLimit(l - streamServerRef.fetchedEntries.size()); |
| 1458 | } |
| 1459 | req.setStream(streamClient); |
| 1460 | return req.sendIgnoringResult(); |
| 1461 | }); |
| 1462 | |
| 1463 | // Wait on the RPC only until stream.end() is called, then report the results. We prevent |
| 1464 | // `stream` from being destroyed until we have a result so that if the RPC throws an exception, |
| 1465 | // we don't accidentally report "PromiseFulfiller not fulfilled" instead of the exception. |
| 1466 | auto promise = sendPromise.then([&streamServerRef]() -> kj::Promise<ActorCache::GetResultList> { |
| 1467 | if (streamServerRef.fulfiller->isWaiting()) { |
| 1468 | return KJ_EXCEPTION(FAILED, "list() never called stream.end()"); |
| 1469 | } else { |
| 1470 | // We'll be canceled momentarily... |
| 1471 | return kj::NEVER_DONE; |
| 1472 | } |
| 1473 | }); |
| 1474 | |
| 1475 | return paf.promise.exclusiveJoin(kj::mv(promise)) |
| 1476 | .attach(kj::defer( |
| 1477 | [client = kj::mv(streamClient), &streamServerRef]() { streamServerRef.cancel(); })); |
| 1478 | } |
| 1479 | |
| 1480 | // ----------------------------------------------------------------------------- |
| 1481 | // Helpers for read operations |
| 1482 | |
| 1483 | kj::Own<ActorCache::Entry> ActorCache::findInCache( |
| 1484 | Lock& lock, KeyPtr key, const ReadOptions& options) { |
| 1485 | auto& map = currentValues.get(lock); |
| 1486 | auto iter = map.seek(key); |
| 1487 | auto ordered = map.ordered(); |
| 1488 | |
| 1489 | if (iter != ordered.end() && iter->get()->key == key) { |
| 1490 | // Found exact matching entry. |
| 1491 | Entry& entry = **iter; |
| 1492 | if (!options.noCache) { |
| 1493 | touchEntry(lock, entry); |
| 1494 | } |
| 1495 | return kj::atomicAddRef(entry); |
| 1496 | } else { |
| 1497 | // Key is not in the map, but we have to check for outstanding list() operations by checking |
| 1498 | // the previous entry's gapState. |
| 1499 | |
| 1500 | if (iter != ordered.begin()) { |
| 1501 | Entry& prev = **--iter; |
| 1502 | if (prev.gapIsKnownEmpty) { |
| 1503 | // A previous list() operation covered this section of the key space and did not find this |
| 1504 | // key, so we know it's not present. Return a dummy entry saying this. |
| 1505 | return kj::atomicRefcounted<ActorCache::Entry>(cloneKey(key), EntryValueStatus::ABSENT); |
| 1506 | } |
| 1507 | } |
| 1508 | |
| 1509 | // We don't know whether this key exists in storage. |
| 1510 | return kj::atomicRefcounted<ActorCache::Entry>(cloneKey(key), EntryValueStatus::UNKNOWN); |
| 1511 | } |
| 1512 | } |
| 1513 | |
| 1514 | kj::Own<ActorCache::Entry> ActorCache::addReadResultToCache( |
| 1515 | Lock& lock, Key key, kj::Maybe<capnp::Data::Reader> maybeReader, const ReadOptions& options) { |
| 1516 | if (options.noCache) { |
| 1517 | // We don't actually want to add this to the cache, just return the entry. |
| 1518 | KJ_IF_SOME(reader, maybeReader) { |
| 1519 | return kj::atomicRefcounted<Entry>(kj::mv(key), kj::heapArray(reader)); |
| 1520 | } else { |
| 1521 | return kj::atomicRefcounted<Entry>(kj::mv(key), EntryValueStatus::ABSENT); |
| 1522 | } |
| 1523 | } |
| 1524 | |
| 1525 | auto& map = currentValues.get(lock); |
| 1526 | |
| 1527 | kj::Own<Entry> entry; |
| 1528 | KJ_IF_SOME(reader, maybeReader) { |
| 1529 | entry = kj::atomicRefcounted<Entry>(*this, kj::mv(key), kj::heapArray(reader)); |
| 1530 | } else { |
| 1531 | // Inserting a negative entry. Let's check if the new insertion is redundant due to the |
| 1532 | // previous entry having `gapIsKnownEmpty`. |
| 1533 | auto iter = map.seek(key); |
| 1534 | auto ordered = map.ordered(); |
| 1535 | if ((iter == ordered.end() || iter->get()->key != key) && iter != ordered.begin()) { |
| 1536 | // We did not find an exact match for the key, so we got an iterator pointing to the next |
| 1537 | // entry after the key. It's not the first entry, so we can back it up one to get the |
| 1538 | // entry before the key. |
| 1539 | --iter; |
| 1540 | |
| 1541 | if (iter->get()->gapIsKnownEmpty) { |
| 1542 | // This entry is redundant, so we won't insert it. |
| 1543 | return kj::atomicRefcounted<Entry>(kj::mv(key), EntryValueStatus::ABSENT); |
| 1544 | ; |
| 1545 | } |
| 1546 | } |
| 1547 | |
| 1548 | entry = kj::atomicRefcounted<Entry>(*this, kj::mv(key), EntryValueStatus::ABSENT); |
| 1549 | // TODO(perf): It's a little sad that we are going to do a findOrCreate() below that is going |
| 1550 | // to repeat the same lookup that produced `iter`. Maybe we could extend kj::Table with a |
| 1551 | // way to provide an existing iterator as a hint when inserting? |
| 1552 | } |
| 1553 | |
| 1554 | // At this point, we know we definitely want there to exist an entry matching this key. So now |
| 1555 | // try to insert it. |
| 1556 | auto& slot = map.findOrCreate(entry->key, [&]() { |
| 1557 | // No existing entry has this key, so insert our new entry. |
| 1558 | // |
| 1559 | // Note that it's definitely guaranteed that the entry *before* the one we're inserting cannot |
| 1560 | // possibly have `gapIsKnownEmpty = true`, because: |
| 1561 | // 1. If our new entry has a null value, then we could have returned early above in this case. |
| 1562 | // 2. If our new entry has a non-null value, then it would be inconsistent for a previous |
| 1563 | // entry to claim that the gap is empty -- this new entry proves it was not! Remember that |
| 1564 | // we are inserting an entry that was the result of reading from disk, so it *must* be |
| 1565 | // consistent with any existing knowledge about the state of disk -- unless we have a bug in |
| 1566 | // the caching logic. |
| 1567 | // |
| 1568 | // Because of this, we know it is correct to leave `gapIsKnownEmpty = false` on our new entry. |
| 1569 | addToCleanList(lock, *entry); |
| 1570 | return kj::atomicAddRef(*entry); |
| 1571 | }); |
| 1572 | |
| 1573 | if (slot.get() != entry.get()) { |
| 1574 | // There was a pre-existing entry with the key, so ours wasn't inserted. |
| 1575 | switch (slot->getValueStatus()) { |
| 1576 | case EntryValueStatus::UNKNOWN: { |
| 1577 | // Oh, it's just a marker for the end of a list range. Go ahead and insert our new entry |
| 1578 | // into the same slot. |
| 1579 | KJ_ASSERT(!slot->gapIsKnownEmpty); // UNKNOWN entry should never have gapIsKnownEmpty. |
| 1580 | removeEntry(lock, *slot); |
| 1581 | |
| 1582 | addToCleanList(lock, *entry); |
| 1583 | slot = kj::atomicAddRef(*entry); |
| 1584 | break; |
| 1585 | } |
| 1586 | case EntryValueStatus::PRESENT: |
| 1587 | case EntryValueStatus::ABSENT: { |
| 1588 | // The entry that's already in the map must be at least as fresh as the one we just created. |
| 1589 | // If it was created by a put() or delete(), then it is actually fresher. If it was created |
| 1590 | // by a concurrent get() or list() that fetched the same key, then it should be exactly the |
| 1591 | // same value. So, either way, our new entry isn't needed. We mark it NOT_IN_CACHE since it |
| 1592 | // won't be placed in the map. |
| 1593 | // |
| 1594 | // NOTE: You might be tempted to say that if the existing entry is DIRTY, but its value |
| 1595 | // matches the value that we just read off disk, then we can cancel the write, because |
| 1596 | // we've discovered it is redundant. Unfortunately, this is NOT true, because it's possible |
| 1597 | // something else has been written in between. Specifically, we could currently be in the |
| 1598 | // process of building a transaction that wrote some other value to this specific key, but |
| 1599 | // hasn't been committed yet, probably because it is waiting for this read operation to |
| 1600 | // complete. Meanwhile, another put() or delete() could have just been performed |
| 1601 | // momentarily ago that changed the flushing entry back to DIRTY and changed its value to |
| 1602 | // one that coincidentally matches what we pulled off disk. However, the open transaction |
| 1603 | // is still going to be committed, writing the intermediate value, so we still need to plan |
| 1604 | // to write this value again in the next transaction. |
| 1605 | touchEntry(lock, *slot); |
| 1606 | break; |
| 1607 | } |
| 1608 | } |
| 1609 | } |
| 1610 | |
| 1611 | return kj::mv(entry); |
| 1612 | } |
| 1613 | |
| 1614 | // Set `gapIsKnownEmpty` across the range covered by a new batch of entries arriving from |
| 1615 | // storage via a list() operation. Since we just listed this range, we know that all the gaps |
| 1616 | // between entries in this range can now be marked as empty. |
| 1617 | // |
| 1618 | // You might ask: "But what if an entry was evicted from the cache between when list() was |
| 1619 | // called and now, creating a gap?" |
| 1620 | // |
| 1621 | // There are two possibilities: |
| 1622 | // 1. The evicted entry was clean at the time list() was called. In this case, the list() |
| 1623 | // operation will have returned it, so it would have been re-added to the cache just |
| 1624 | // before this method call. |
| 1625 | // 2. The evicted entry was dirty at the time list() was called. This can't cause a problem |
| 1626 | // because we ensure that any flush is ordered after all previous read operations, so such |
| 1627 | // entries could not possibly be marked clean until after the list operation completes. |
| 1628 | // And, they cannot be evicted until they are marked clean. So these entries could not |
| 1629 | // have been evicted yet. |
| 1630 | void ActorCache::markGapsEmpty( |
| 1631 | Lock& lock, KeyPtr beginKey, kj::Maybe<KeyPtr> endKey, const ReadOptions& options) { |
| 1632 | if (options.noCache) { |
| 1633 | // Oops, never mind. We're not caching the list() results, so we can't mark anything |
| 1634 | // known-empty. |
| 1635 | return; |
| 1636 | } |
| 1637 | |
| 1638 | auto& map = currentValues.get(lock); |
| 1639 | |
| 1640 | auto endIter = seekOrEnd(map, endKey); |
| 1641 | { |
| 1642 | auto ordered = map.ordered(); |
| 1643 | if (endIter == ordered.end() || endIter->get()->key > endKey) { |
| 1644 | // The key that we're marking up *to* is not in the map. |
| 1645 | if (endIter == ordered.begin()) { |
| 1646 | // Whoops, it appears we don't actually have any entries in the marking range. This could |
| 1647 | // happen during a forward list() due to entries from previous values() calls having |
| 1648 | // already been evicted before end() was called. In this case, nothing would actually be |
| 1649 | // marked below. But then our UNKNOWN entry would be inconsistent, so we'd better not |
| 1650 | // insert it at all. |
| 1651 | // |
| 1652 | // Note that this does NOT happen as a result of a list() returning no results, because |
| 1653 | // in that case the list operation would have inserted a negative entry at the beginning |
| 1654 | // of the range. The only reason why we wouldn't have found that negative entry here is |
| 1655 | // because it has since been evicted. |
| 1656 | return; |
| 1657 | } |
| 1658 | |
| 1659 | --endIter; |
| 1660 | if (endIter->get()->key < beginKey) { |
| 1661 | // Same as above, it appears we have no suitable entries to mark, so we can't insert an |
| 1662 | // UNKNOWN. |
| 1663 | return; |
| 1664 | } |
| 1665 | |
| 1666 | if (endIter->get()->gapIsKnownEmpty) { |
| 1667 | // The end key is in an already-known-empty gap, so there's no need to insert an UNKNOWN. |
| 1668 | // We intentionally leave `endIter` pointing to the start of the gap even though it's not |
| 1669 | // the end of our list range, because we know the stuff from there to the end of the range |
| 1670 | // is already marked. |
| 1671 | } else { |
| 1672 | // We must insert an UNKNOWN entry to cap our range. |
| 1673 | KJ_IF_SOME(k, endKey) { |
| 1674 | auto entry = kj::atomicRefcounted<Entry>(*this, cloneKey(k), EntryValueStatus::UNKNOWN); |
| 1675 | addToCleanList(lock, *entry); |
| 1676 | map.insert(kj::mv(entry)); |
| 1677 | } else { |
| 1678 | // No UNKNOWN needed since the end is actually the end of the key space. |
| 1679 | } |
| 1680 | |
| 1681 | // Oops, that invalidated our iterator, so find it again. |
| 1682 | endIter = seekOrEnd(map, endKey); |
| 1683 | } |
| 1684 | } |
| 1685 | } |
| 1686 | |
| 1687 | kj::Vector<KeyPtr> keysToErase; |
| 1688 | auto beginIter = map.seek(beginKey); |
| 1689 | auto mapEnd = map.ordered().end(); |
| 1690 | for (auto iter = beginIter; iter != mapEnd; ++iter) { |
| 1691 | auto& entry = **iter; |
| 1692 | |
| 1693 | if (entry.getValueStatus() != EntryValueStatus::PRESENT && !entry.isDirty() && |
| 1694 | (iter != endIter || iter->get()->gapIsKnownEmpty)) { |
| 1695 | // Either: |
| 1696 | // (a) This is an UNKNOWN entry. |
| 1697 | // (b) This is a clean negative entry. |
| 1698 | // |
| 1699 | // And either: |
| 1700 | // (a) This is not the last entry, so we're about to set `gapIsKnownEmpty` on it. |
| 1701 | // (b) It is the last entry, and it is already `gapIsKnownEmpty`. |
| 1702 | // |
| 1703 | // Either way, if the *previous* entry also has `gapIsKnownEmpty`, then *this* entry |
| 1704 | // becomes redundant. In that case we need to delete it instead. |
| 1705 | // |
| 1706 | // Note that a negative entry that is DIRTY is not necessarily redundant, because it |
| 1707 | // could be that a different value was written to that entry and then deleted between |
| 1708 | // when the list() was initiated and the current state of the cache. A negative DIRTY |
| 1709 | // entry will become redundant once it becomes CLEAN, so we'll have to deal with it then. |
| 1710 | |
| 1711 | bool prevGapIsEmpty; |
| 1712 | if (iter == beginIter) { |
| 1713 | // This is the first entry in the range, so we have to check if the previous entry |
| 1714 | // was marked. |
| 1715 | if (iter == map.ordered().begin()) { |
| 1716 | prevGapIsEmpty = false; |
| 1717 | } else { |
| 1718 | auto prev = iter; |
| 1719 | --prev; |
| 1720 | prevGapIsEmpty = prev->get()->gapIsKnownEmpty; |
| 1721 | } |
| 1722 | } else { |
| 1723 | // This isn't the first entry we've iterated over so we must have marked the previous |
| 1724 | // one with gapIsKnownEmpty. |
| 1725 | prevGapIsEmpty = true; |
| 1726 | } |
| 1727 | |
| 1728 | if (prevGapIsEmpty) { |
| 1729 | // Unfortunately erasing from the map will invalidate our iterator, so we need to make |
| 1730 | // a second pass to erase, below. |
| 1731 | keysToErase.add(entry.key); |
| 1732 | } |
| 1733 | } |
| 1734 | |
| 1735 | if (iter == endIter) { |
| 1736 | // We didn't check for `iter == endIter` earlier because the conditional above -- which |
| 1737 | // potentially deletes redundant entries -- can actually apply to the end of the range, even |
| 1738 | // though that entry itself isn't considered part of the range. Marking the range could cause |
| 1739 | // the entry immediately after the end to become redundant. |
| 1740 | // |
| 1741 | // We do want to break here, though, because we do not want to mark an entry that is past |
| 1742 | // the end of the range. |
| 1743 | break; |
| 1744 | } |
| 1745 | |
| 1746 | entry.gapIsKnownEmpty = true; |
| 1747 | } |
| 1748 | |
| 1749 | for (auto& key: keysToErase) { |
| 1750 | auto& entry = KJ_ASSERT_NONNULL(map.find(key)); |
| 1751 | removeEntry(lock, *entry); |
| 1752 | map.erase(entry); |
| 1753 | } |
| 1754 | } |
| 1755 | |
| 1756 | ActorCache::GetResultList::GetResultList(kj::Vector<KeyValuePair> contents) |
| 1757 | : entries(contents.size()), |
| 1758 | cacheStatuses(contents.size()) { |
| 1759 | // TODO(perf): Allocating an `Entry` object for every key/value pair is lame but to avoid it |
| 1760 | // we'd have to make the common case worse... |
| 1761 | for (auto& kv: contents) { |
| 1762 | entries.add(kj::atomicRefcounted<Entry>(kj::mv(kv.key), kj::mv(kv.value))); |
| 1763 | cacheStatuses.add(CacheStatus::UNCACHED); |
| 1764 | } |
| 1765 | } |
| 1766 | |
| 1767 | // Merges `cachedEntries` and `fetchedEntries`, which should each already be sorted in the |
| 1768 | // given order. If a key exists in both, `cachedEntries` is preferred. |
| 1769 | // |
| 1770 | // After merging, if an entry's value is null, it is dropped. |
| 1771 | // |
| 1772 | // The final result is truncated to `limit`, if any. |
| 1773 | // |
| 1774 | // The idea is that `cachedEntries` is the set of entries that were loaded from cache while |
| 1775 | // `fetchedEntries` is the set read from storage. |
| 1776 | ActorCache::GetResultList::GetResultList(kj::Vector<kj::Own<Entry>> cachedEntries, |
| 1777 | kj::Vector<kj::Own<Entry>> fetchedEntries, |
| 1778 | Order order, |
| 1779 | kj::Maybe<uint> maybeLimit) { |
| 1780 | uint limit = maybeLimit.orDefault(kj::maxValue); |
| 1781 | entries.reserve(kj::min(cachedEntries.size() + fetchedEntries.size(), limit)); |
| 1782 | |
| 1783 | auto cachedIter = cachedEntries.begin(); |
| 1784 | auto fetchedIter = fetchedEntries.begin(); |
| 1785 | |
| 1786 | auto add = [&](kj::Own<ActorCache::Entry>&& entry, CacheStatus status) { |
| 1787 | // Remove null values. |
| 1788 | if (entry->getValueStatus() == ActorCache::EntryValueStatus::PRESENT) { |
| 1789 | entries.add(kj::mv(entry)); |
| 1790 | cacheStatuses.add(status); |
| 1791 | } |
| 1792 | }; |
| 1793 | |
| 1794 | while ((cachedIter != cachedEntries.end() || fetchedIter != fetchedEntries.end()) && |
| 1795 | entries.size() < limit) { |
| 1796 | if (cachedIter == cachedEntries.end()) { |
| 1797 | add(kj::mv(*fetchedIter++), CacheStatus::UNCACHED); |
| 1798 | } else if (fetchedIter == fetchedEntries.end()) { |
| 1799 | add(kj::mv(*cachedIter++), CacheStatus::CACHED); |
| 1800 | } else if (order == REVERSE ? cachedIter->get()->key > fetchedIter->get()->key |
| 1801 | : cachedIter->get()->key < fetchedIter->get()->key) { |
| 1802 | add(kj::mv(*cachedIter++), CacheStatus::CACHED); |
| 1803 | } else if (cachedIter->get()->key == fetchedIter->get()->key) { |
| 1804 | // Same key in both. Prefer the cached entry because it will reflect the state as of when the |
| 1805 | // operation began. |
| 1806 | // Uncached status because we still fetched from disk. |
| 1807 | add(kj::mv(*cachedIter++), CacheStatus::UNCACHED); |
| 1808 | ++fetchedIter; |
| 1809 | } else { |
| 1810 | add(kj::mv(*fetchedIter++), CacheStatus::UNCACHED); |
| 1811 | } |
| 1812 | } |
| 1813 | |
| 1814 | #ifdef KJ_DEBUG |
| 1815 | // Verify sort. |
| 1816 | kj::Maybe<KeyPtr> prev; |
| 1817 | for (auto& entry: entries) { |
| 1818 | KJ_IF_SOME(p, prev) { |
| 1819 | if (order == REVERSE) { |
| 1820 | KJ_ASSERT(entry->key < p); |
| 1821 | } else { |
| 1822 | KJ_ASSERT(entry->key > p); |
| 1823 | } |
| 1824 | } |
| 1825 | prev = entry->key; |
| 1826 | } |
| 1827 | #endif |
| 1828 | } |
| 1829 | |
| 1830 | template <typename Func> |
| 1831 | kj::PromiseForResult<Func, rpc::ActorStorage::Operations::Client> ActorCache::scheduleStorageRead( |
| 1832 | Func&& function) { |
| 1833 | // This is basically kj::retryOnDisconnect() except that we make the first call synchronously. |
| 1834 | // For our use case, this is safe, and I wanted to make sure reads get sent concurrently with |
| 1835 | // further JavaScript execution if possible. |
| 1836 | auto promise = kj::evalNow( |
| 1837 | [&]() mutable { return function(storage).attach(recordStorageRead(hooks, clock)); }); |
| 1838 | return oomCanceler.wrap( |
| 1839 | promise |
| 1840 | .catch_([this, function = kj::mv(function)](kj::Exception&& e) mutable |
| 1841 | -> kj::PromiseForResult<Func, rpc::ActorStorage::Operations::Client> { |
| 1842 | if (e.getType() == kj::Exception::Type::DISCONNECTED) { |
| 1843 | return function(storage).attach(recordStorageRead(hooks, clock)); |
| 1844 | } else { |
| 1845 | return kj::mv(e); |
| 1846 | } |
| 1847 | }).attach(kj::addRef(*readCompletionChain))); |
| 1848 | } |
| 1849 | |
| 1850 | kj::Promise<void> ActorCache::waitForPastReads() { |
| 1851 | if (!readCompletionChain->isShared()) { |
| 1852 | // No reads are in flight right now. |
| 1853 | return kj::READY_NOW; |
| 1854 | } |
| 1855 | |
| 1856 | // Create a new chain link. |
| 1857 | auto next = kj::refcounted<ReadCompletionChain>(); |
| 1858 | |
| 1859 | // Update previous chain so that when it is destroyed, it'll fulfill us and also drop its |
| 1860 | // reference on the next link. |
| 1861 | auto paf = kj::newPromiseAndFulfiller<void>(); |
| 1862 | readCompletionChain->fulfiller = kj::mv(paf.fulfiller); |
| 1863 | readCompletionChain->next = kj::addRef(*next); |
| 1864 | |
| 1865 | // Make `next` the current link. |
| 1866 | readCompletionChain = kj::mv(next); |
| 1867 | |
| 1868 | return kj::mv(paf.promise); |
| 1869 | } |
| 1870 | |
| 1871 | ActorCache::ReadCompletionChain::~ReadCompletionChain() noexcept(false) { |
| 1872 | KJ_IF_SOME(f, fulfiller) { |
| 1873 | f->fulfill(); |
| 1874 | } |
| 1875 | } |
| 1876 | |
| 1877 | // ======================================================================================= |
| 1878 | // write operations |
| 1879 | |
| 1880 | kj::Maybe<kj::Promise<void>> ActorCache::put( |
| 1881 | Key key, Value value, WriteOptions options, SpanParent traceSpan) { |
| 1882 | ActorStorageLimits::checkMaxKeySize(key); |
| 1883 | ActorStorageLimits::checkMaxValueSize(key, value); |
| 1884 | |
| 1885 | options.noCache = options.noCache || lru.options.noCache; |
| 1886 | requireNotTerminal(traceSpan.addRef()); |
| 1887 | { |
| 1888 | auto lock = lru.cleanList.lockExclusive(); |
| 1889 | kj::Maybe<CountedDelete> maybeCountedDelete; |
| 1890 | auto entry = kj::atomicRefcounted<Entry>(*this, kj::mv(key), kj::mv(value)); |
| 1891 | putImpl(lock, kj::mv(entry), options, maybeCountedDelete, kj::mv(traceSpan)); |
| 1892 | evictOrOomIfNeeded(lock); |
| 1893 | } |
| 1894 | return getBackpressure(); |
| 1895 | } |
| 1896 | |
| 1897 | kj::Maybe<kj::Promise<void>> ActorCache::put( |
| 1898 | kj::Array<KeyValuePair> pairs, WriteOptions options, SpanParent traceSpan) { |
| 1899 | for (auto& pair: pairs) { |
| 1900 | // We check limits in a separate loop to fail the whole operation when any pair fails a check |
| 1901 | ActorStorageLimits::checkMaxKeySize(pair.key); |
| 1902 | ActorStorageLimits::checkMaxValueSize(pair.key, pair.value); |
| 1903 | } |
| 1904 | |
| 1905 | options.noCache = options.noCache || lru.options.noCache; |
| 1906 | requireNotTerminal(traceSpan.addRef()); |
| 1907 | { |
| 1908 | auto lock = lru.cleanList.lockExclusive(); |
| 1909 | for (auto& pair: pairs) { |
| 1910 | kj::Maybe<CountedDelete> maybeCountedDelete; |
| 1911 | auto entry = kj::atomicRefcounted<Entry>(*this, kj::mv(pair.key), kj::mv(pair.value)); |
| 1912 | putImpl(lock, kj::mv(entry), options, maybeCountedDelete, traceSpan.addRef()); |
| 1913 | } |
| 1914 | evictOrOomIfNeeded(lock); |
| 1915 | } |
| 1916 | return getBackpressure(); |
| 1917 | } |
| 1918 | |
| 1919 | kj::Maybe<kj::Promise<void>> ActorCache::setAlarm( |
| 1920 | kj::Maybe<kj::Date> newAlarmTime, WriteOptions options, SpanParent traceSpan) { |
| 1921 | options.noCache = options.noCache || lru.options.noCache; |
| 1922 | KJ_IF_SOME(time, currentAlarmTime.tryGet<KnownAlarmTime>()) { |
| 1923 | // If we're in the alarm handler and haven't set the time yet, |
| 1924 | // we can't perform this optimization as currentAlarmTime will be equal |
| 1925 | // to the currently running time but we indicate to the actor in getAlarm() that there |
| 1926 | // is no alarm set, therefore we need to act like that in setAlarm(). |
| 1927 | // |
| 1928 | // After the first write in the handler occurs, which would set KnownAlarmTime, |
| 1929 | // the logic here is correct again as currentAlarmTime would match what we are reporting |
| 1930 | // to the user from getAlarm(). |
| 1931 | // |
| 1932 | // So, we only apply this for KnownAlarmTime. |
| 1933 | |
| 1934 | if (time.time == newAlarmTime) { |
| 1935 | // No change! May as well skip the storage operation. |
| 1936 | return kj::none; |
| 1937 | } |
| 1938 | } |
| 1939 | |
| 1940 | currentAlarmTime = ActorCache::KnownAlarmTime{ |
| 1941 | ActorCache::KnownAlarmTime::Status::DIRTY, newAlarmTime, options.noCache}; |
| 1942 | |
| 1943 | ensureFlushScheduled(options, kj::mv(traceSpan)); |
| 1944 | |
| 1945 | return getBackpressure(); |
| 1946 | } |
| 1947 | |
| 1948 | namespace { |
| 1949 | template <typename F> |
| 1950 | kj::OneOf<std::invoke_result_t<F>, kj::PromiseForResult<F, void>> mapPromise( |
| 1951 | kj::Maybe<kj::Promise<void>> maybePromise, F&& f) { |
| 1952 | KJ_IF_SOME(promise, maybePromise) { |
| 1953 | return promise.then(kj::fwd<F>(f)); |
| 1954 | } else { |
| 1955 | return kj::fwd<F>(f)(); |
| 1956 | } |
| 1957 | } |
| 1958 | } // namespace |
| 1959 | |
| 1960 | kj::OneOf<bool, kj::Promise<bool>> ActorCache::delete_( |
| 1961 | Key key, WriteOptions options, SpanParent traceSpan) { |
| 1962 | ActorStorageLimits::checkMaxKeySize(key); |
| 1963 | |
| 1964 | options.noCache = options.noCache || lru.options.noCache; |
| 1965 | requireNotTerminal(traceSpan.addRef()); |
| 1966 | |
| 1967 | auto countedDelete = kj::refcounted<CountedDelete>(); |
| 1968 | { |
| 1969 | auto lock = lru.cleanList.lockExclusive(); |
| 1970 | auto entry = kj::atomicRefcounted<Entry>(*this, kj::mv(key), EntryValueStatus::ABSENT); |
| 1971 | putImpl(lock, kj::mv(entry), options, *countedDelete, kj::mv(traceSpan)); |
| 1972 | evictOrOomIfNeeded(lock); |
| 1973 | } |
| 1974 | |
| 1975 | auto waiter = kj::heap<CountedDeleteWaiter>(*this, kj::addRef(*countedDelete)); |
| 1976 | kj::Maybe<kj::Promise<void>> maybePromise; |
| 1977 | KJ_IF_SOME(p, getBackpressure()) { |
| 1978 | // This might be more than one flush but that's OK as long as our state gets taken care of. |
| 1979 | maybePromise = countedDelete->forgiveIfFinished(kj::mv(p)); |
| 1980 | } else if (!countedDelete->entries.empty()) { |
| 1981 | maybePromise = countedDelete->forgiveIfFinished(lastFlush.addBranch()); |
| 1982 | } |
| 1983 | return mapPromise(kj::mv(maybePromise), |
| 1984 | [waiter = kj::mv(waiter)]() { return waiter->getCountedDelete().countDeleted > 0; }); |
| 1985 | } |
| 1986 | |
| 1987 | kj::OneOf<uint, kj::Promise<uint>> ActorCache::delete_( |
| 1988 | kj::Array<Key> keys, WriteOptions options, SpanParent traceSpan) { |
| 1989 | for (auto& key: keys) { |
| 1990 | ActorStorageLimits::checkMaxKeySize(key); |
| 1991 | } |
| 1992 | |
| 1993 | options.noCache = options.noCache || lru.options.noCache; |
| 1994 | requireNotTerminal(traceSpan.addRef()); |
| 1995 | |
| 1996 | auto countedDelete = kj::refcounted<CountedDelete>(); |
| 1997 | { |
| 1998 | auto lock = lru.cleanList.lockExclusive(); |
| 1999 | for (auto& key: keys) { |
| 2000 | auto entry = kj::atomicRefcounted<Entry>(*this, kj::mv(key), EntryValueStatus::ABSENT); |
| 2001 | putImpl(lock, kj::mv(entry), options, *countedDelete, traceSpan.addRef()); |
| 2002 | } |
| 2003 | evictOrOomIfNeeded(lock); |
| 2004 | } |
| 2005 | |
| 2006 | auto waiter = kj::heap<CountedDeleteWaiter>(*this, kj::addRef(*countedDelete)); |
| 2007 | kj::Maybe<kj::Promise<void>> maybePromise; |
| 2008 | KJ_IF_SOME(p, getBackpressure()) { |
| 2009 | // This might be more than one flush but that's OK as long as our state gets taken care of. |
| 2010 | maybePromise = countedDelete->forgiveIfFinished(kj::mv(p)); |
| 2011 | } else if (!countedDelete->entries.empty()) { |
| 2012 | maybePromise = countedDelete->forgiveIfFinished(lastFlush.addBranch()); |
| 2013 | } |
| 2014 | return mapPromise(kj::mv(maybePromise), |
| 2015 | [waiter = kj::mv(waiter)]() { return waiter->getCountedDelete().countDeleted; }); |
| 2016 | } |
| 2017 | |
| 2018 | kj::Own<ActorCacheInterface::Transaction> ActorCache::startTransaction() { |
| 2019 | return kj::heap<Transaction>(*this); |
| 2020 | } |
| 2021 | |
| 2022 | ActorCache::DeleteAllResults ActorCache::deleteAll( |
| 2023 | WriteOptions options, SpanParent traceSpan, DeleteAllOptions deleteAllOptions) { |
| 2024 | // Since deleteAll() cannot be performed as part of another transaction, in order to maintain |
| 2025 | // our ordering guarantees, we will have to complete all writes that occurred prior to the |
| 2026 | // deleteAll(), then submit the deleteAll(), then do any writes afterwards. Conveniently, though, |
| 2027 | // a deleteAll() invalidates the whole map. So, we can take all the dirty entries out and place |
| 2028 | // them off to the side for the moment, so that overwrites won't affect them. (Otherwise, an |
| 2029 | // overwritten entry would be moved to the end of the dirty list, which might mean it is |
| 2030 | // committed in the wrong order with respect to the deleteAll().) |
| 2031 | |
| 2032 | options.noCache = options.noCache || lru.options.noCache; |
| 2033 | requireNotTerminal(traceSpan.addRef()); |
| 2034 | |
| 2035 | kj::Promise<uint> result{static_cast<uint>(0)}; |
| 2036 | |
| 2037 | { |
| 2038 | auto lock = lru.cleanList.lockExclusive(); |
| 2039 | auto& map = currentValues.get(lock); |
| 2040 | |
| 2041 | kj::Vector<kj::Own<Entry>> deletedDirty; |
| 2042 | for (auto& entry: dirtyList) { |
| 2043 | // We will be removing all entries from their respective lists soon, so let's preserve the |
| 2044 | // dirty list so we can run it before our delete all. |
| 2045 | deletedDirty.add(kj::atomicAddRef(entry)); |
| 2046 | }; |
| 2047 | |
| 2048 | // Clear out the entire map. |
| 2049 | for (auto& entry: map) { |
| 2050 | removeEntry(lock, *entry); |
| 2051 | } |
| 2052 | map.clear(); |
| 2053 | |
| 2054 | // Insert a dummy entry with an ABSENT key and gapIsKnownEmpty = true to indicate that |
| 2055 | // everything is empty. |
| 2056 | map.findOrCreate(Key{}, [&]() { |
| 2057 | Key key; |
| 2058 | auto entry = kj::atomicRefcounted<Entry>(*this, kj::mv(key), EntryValueStatus::ABSENT); |
| 2059 | addToCleanList(lock, *entry); |
| 2060 | entry->gapIsKnownEmpty = true; |
| 2061 | return entry; |
| 2062 | }); |
| 2063 | |
| 2064 | KJ_IF_SOME(existing, requestedDeleteAll) { |
| 2065 | // A previous deleteAll() was scheduled and hasn't been committed yet. This means that we |
| 2066 | // can actually coalesce the two, and there's no need to commit any writes that happened |
| 2067 | // between them. So we can throw away `deletedDirty`. |
| 2068 | // We also don't want to double-bill for a coalesced deleteAll, so we don't update |
| 2069 | // result in this branch. We do ensure the alarm is deleted if requested, though. |
| 2070 | existing.deleteAlarm = existing.deleteAlarm || deleteAllOptions.deleteAlarm; |
| 2071 | } else { |
| 2072 | // If no previous deleteAll() was scheduled, then schedule this one. |
| 2073 | auto paf = kj::newPromiseAndFulfiller<uint>(); |
| 2074 | result = kj::mv(paf.promise); |
| 2075 | requestedDeleteAll = DeleteAllState{ |
| 2076 | .deletedDirty = kj::mv(deletedDirty), |
| 2077 | .countFulfiller = kj::mv(paf.fulfiller), |
| 2078 | .deleteAlarm = deleteAllOptions.deleteAlarm, |
| 2079 | }; |
| 2080 | ensureFlushScheduled(options, kj::mv(traceSpan)); |
| 2081 | } |
| 2082 | |
| 2083 | if (deleteAllOptions.deleteAlarm) { |
| 2084 | // Update the in-memory alarm state immediately so that getAlarm() returns null right away. |
| 2085 | // The actual alarm deletion from storage is deferred to flushImplDeleteAll(), which sets |
| 2086 | // currentAlarmTime to DIRTY after the deleteAll RPC succeeds, causing the post-deleteAll |
| 2087 | // flush to send the deleteAlarm RPC. |
| 2088 | currentAlarmTime = KnownAlarmTime{ |
| 2089 | .status = KnownAlarmTime::Status::CLEAN, |
| 2090 | .time = kj::none, |
| 2091 | }; |
| 2092 | } |
| 2093 | |
| 2094 | // This is called for consistency, but deleteAll() strictly reduces cache usage, so it's not |
| 2095 | // entirely necessary. |
| 2096 | evictOrOomIfNeeded(lock); |
| 2097 | } |
| 2098 | |
| 2099 | return DeleteAllResults{.backpressure = getBackpressure(), .count = kj::mv(result)}; |
| 2100 | } |
| 2101 | |
| 2102 | void ActorCache::putImpl(Lock& lock, |
| 2103 | kj::Own<Entry> newEntry, |
| 2104 | const WriteOptions& options, |
| 2105 | kj::Maybe<CountedDelete&> maybeCountedDelete, |
| 2106 | SpanParent traceSpan) { |
| 2107 | auto& map = currentValues.get(lock); |
| 2108 | auto ordered = map.ordered(); |
| 2109 | |
| 2110 | // This gets a little complicated because we want to avoid redundant insertions. |
| 2111 | |
| 2112 | newEntry->noCache = options.noCache; |
| 2113 | |
| 2114 | auto iter = map.seek(newEntry->key); |
| 2115 | if (iter != ordered.end() && iter->get()->key == newEntry->key) { |
| 2116 | // Exact same entry already exists. |
| 2117 | auto& slot = *iter; |
| 2118 | |
| 2119 | switch (slot->getValueStatus()) { |
| 2120 | case EntryValueStatus::PRESENT: { |
| 2121 | if (slot->getValuePtr() == newEntry->getValuePtr()) { |
| 2122 | // No change! The entry already had this value. Might as well skip the whole storage |
| 2123 | // operation. |
| 2124 | return; |
| 2125 | } |
| 2126 | |
| 2127 | KJ_IF_SOME(c, maybeCountedDelete) { |
| 2128 | // Overwrote an entry that was in cache, so we can count it now. Note that because we |
| 2129 | // are PRESENT, we will not be added to the CountedDelete's `entries`, since we only |
| 2130 | // do this for UNKNOWN entries! Instead, we'll be part of a regular delete. |
| 2131 | ++c.countDeleted; |
| 2132 | } |
| 2133 | break; |
| 2134 | } |
| 2135 | case EntryValueStatus::ABSENT: { |
| 2136 | if (slot->getValuePtr() == newEntry->getValuePtr()) { |
| 2137 | // No change! The entry already had this value. Might as well skip the whole storage |
| 2138 | // operation. |
| 2139 | return; |
| 2140 | } |
| 2141 | |
| 2142 | if (slot->isCountedDelete) { |
| 2143 | // We are overwriting an entry that is slated for a counted delete operation. |
| 2144 | // There may be a situation where all the entries associated with a counted delete are |
| 2145 | // actually successfully deleted (and we get the count), but the transaction the deletes |
| 2146 | // execute within fails. |
| 2147 | // |
| 2148 | // Since we are currently overwriting the Entry, we might as well inform the |
| 2149 | // `CountedDelete` that this Entry has since been overwritten. Then, if we hit the case |
| 2150 | // described above, we won't need to include this Entry in a subsequent counted delete |
| 2151 | // retry, since we already have the count AND the Entry has been overwritten. |
| 2152 | // |
| 2153 | // For more details, see how we filter the entries to be deleted for a CountedDeleteFlush |
| 2154 | // as part of a flush. |
| 2155 | slot->overwritingCountedDelete = true; |
| 2156 | } |
| 2157 | // We don't have to worry about the counted delete since we were already deleted. |
| 2158 | break; |
| 2159 | } |
| 2160 | case EntryValueStatus::UNKNOWN: { |
| 2161 | // This was a list end marker, we should just overwrite it. |
| 2162 | |
| 2163 | KJ_IF_SOME(c, maybeCountedDelete) { |
| 2164 | // Despite an entry being present, we don't know if the key exists, because it's just an |
| 2165 | // UNKNOWN entry. So we will still have to arrange to count the delete later. |
| 2166 | newEntry->isCountedDelete = true; |
| 2167 | c.entries.add(kj::atomicAddRef(*newEntry)); |
| 2168 | } |
| 2169 | break; |
| 2170 | } |
| 2171 | } |
| 2172 | |
| 2173 | KJ_DASSERT(slot->key == newEntry->key); |
| 2174 | |
| 2175 | // Inherit gap state. |
| 2176 | newEntry->gapIsKnownEmpty = slot->gapIsKnownEmpty; |
| 2177 | |
| 2178 | // Swap in the new entry. |
| 2179 | removeEntry(lock, *slot); |
| 2180 | |
| 2181 | slot = kj::mv(newEntry); |
| 2182 | addToDirtyList(*slot); |
| 2183 | } else { |
| 2184 | // No exact matching entry exists, insert a new one. |
| 2185 | |
| 2186 | // Does the previous entry have a known-empty gap? |
| 2187 | bool previousGapKnownEmpty = false; |
| 2188 | if (iter != ordered.begin()) { |
| 2189 | --iter; |
| 2190 | previousGapKnownEmpty = iter->get()->gapIsKnownEmpty; |
| 2191 | } |
| 2192 | if (previousGapKnownEmpty && newEntry->getValueStatus() == EntryValueStatus::ABSENT) { |
| 2193 | // No change! The entry is already known not to exist, and we're trying to delete it. Might |
| 2194 | // as well skip the whole storage operation. |
| 2195 | return; |
| 2196 | } |
| 2197 | |
| 2198 | // Create the new entry. |
| 2199 | // TODO(perf): Extend kj::TreeIndex to allow supplying the existing iterator as a hint when |
| 2200 | // inserting a new entry, to avoid repeating the lookup. |
| 2201 | auto& slot = map.insert(kj::mv(newEntry)); |
| 2202 | slot->gapIsKnownEmpty = previousGapKnownEmpty; |
| 2203 | KJ_IF_SOME(c, maybeCountedDelete) { |
| 2204 | slot->isCountedDelete = true; |
| 2205 | c.entries.add(kj::atomicAddRef(*slot)); |
| 2206 | } |
| 2207 | addToDirtyList(*slot); |
| 2208 | } |
| 2209 | |
| 2210 | ensureFlushScheduled(options, kj::mv(traceSpan)); |
| 2211 | } |
| 2212 | |
| 2213 | void ActorCache::ensureFlushScheduled(const WriteOptions& options, SpanParent traceSpan) { |
| 2214 | if (lru.options.neverFlush) { |
| 2215 | // Skip all flushes. Used for preview sessions where data is strictly kept in memory. |
| 2216 | |
| 2217 | // Handle deleteAll state that would normally be processed during flushImplDeleteAll(). |
| 2218 | KJ_IF_SOME(deleteAllState, requestedDeleteAll) { |
| 2219 | deleteAllState.countFulfiller->fulfill(0); |
| 2220 | if (deleteAllState.deleteAlarm) { |
| 2221 | currentAlarmTime = KnownAlarmTime{ |
| 2222 | .status = KnownAlarmTime::Status::DIRTY, |
| 2223 | .time = kj::none, |
| 2224 | }; |
| 2225 | } |
| 2226 | requestedDeleteAll = kj::none; |
| 2227 | } |
| 2228 | |
| 2229 | // Also, we need to handle scheduling or canceling any alarm changes locally. |
| 2230 | KJ_SWITCH_ONEOF(currentAlarmTime) { |
| 2231 | KJ_CASE_ONEOF(knownAlarmTime, ActorCache::KnownAlarmTime) { |
| 2232 | if (knownAlarmTime.status == KnownAlarmTime::Status::DIRTY) { |
| 2233 | knownAlarmTime.status = KnownAlarmTime::Status::CLEAN; |
| 2234 | hooks.updateAlarmInMemory(knownAlarmTime.time); |
| 2235 | } |
| 2236 | } |
| 2237 | KJ_CASE_ONEOF(deferredDelete, ActorCache::DeferredAlarmDelete) { |
| 2238 | if (deferredDelete.status == DeferredAlarmDelete::Status::READY) { |
| 2239 | currentAlarmTime = KnownAlarmTime{ |
| 2240 | .status = KnownAlarmTime::Status::CLEAN, |
| 2241 | .time = kj::none, |
| 2242 | }; |
| 2243 | hooks.updateAlarmInMemory(kj::none); |
| 2244 | } |
| 2245 | } |
| 2246 | KJ_CASE_ONEOF(_, UnknownAlarmTime) {} |
| 2247 | } |
| 2248 | |
| 2249 | return; |
| 2250 | } |
| 2251 | |
| 2252 | if (!flushScheduled) { |
| 2253 | flushScheduled = true; |
| 2254 | // Capture the trace span from the first write in this flush batch. |
| 2255 | currentFlushSpan = kj::mv(traceSpan); |
| 2256 | |
| 2257 | auto flushPromise = lastFlush.addBranch() |
| 2258 | .attach(kj::defer([this]() { |
| 2259 | flushScheduled = false; |
| 2260 | flushScheduledWithOutputGate = false; |
| 2261 | // Reset the flush span for the next batch |
| 2262 | currentFlushSpan = nullptr; |
| 2263 | })).then([this]() { |
| 2264 | ++flushesEnqueued; |
| 2265 | return kj::evalNow([this]() { |
| 2266 | // `flushImpl()` can throw, so we need to wrap it in `evalNow()` to observe all pathways. |
| 2267 | return flushImpl(); |
| 2268 | }).attach(kj::defer([this]() { --flushesEnqueued; })); |
| 2269 | }); |
| 2270 | |
| 2271 | if (options.allowUnconfirmed) { |
| 2272 | // Don't apply output gate. But, if an exception is thrown, we still want to break the gate, |
| 2273 | // so arrange for that. |
| 2274 | flushPromise = flushPromise.catch_([this](kj::Exception&& e) { |
| 2275 | return gate.lockWhile(kj::Promise<void>(kj::mv(e)), nullptr); |
| 2276 | }); |
| 2277 | } else { |
| 2278 | flushPromise = gate.lockWhile(kj::mv(flushPromise), currentFlushSpan.addRef()); |
| 2279 | flushScheduledWithOutputGate = true; |
| 2280 | } |
| 2281 | |
| 2282 | lastFlush = flushPromise.fork(); |
| 2283 | } else if (!flushScheduledWithOutputGate && !options.allowUnconfirmed) { |
| 2284 | // The flush has already been scheduled without the output gate, but we want to upgrade it to |
| 2285 | // use the output gate now. The span was already captured when the flush was first scheduled. |
| 2286 | lastFlush = gate.lockWhile(lastFlush.addBranch(), currentFlushSpan.addRef()).fork(); |
| 2287 | flushScheduledWithOutputGate = true; |
| 2288 | } |
| 2289 | } |
| 2290 | |
| 2291 | // This function returns a Maybe<Promise> because a falsy maybe allows the jsg interface to make |
| 2292 | // a resolved jsg::Promise. This is meaningfully different from a ready kj::Promise because it |
| 2293 | // allows the next continuation to run immediately on the microtask queue instead of returning to |
| 2294 | // the kj event loop and fulfilling a resolver that enqueues the continuation. |
| 2295 | kj::Maybe<kj::Promise<void>> ActorCache::onNoPendingFlush(SpanParent parentSpan) { |
| 2296 | if (lru.options.neverFlush) { |
| 2297 | // We won't ever flush (usually because we're a preview session), so return a falsy maybe. |
| 2298 | return kj::none; |
| 2299 | } |
| 2300 | |
| 2301 | if (flushScheduled) { |
| 2302 | // There is a flush that is currently scheduled but not yet running, we need to wait for that |
| 2303 | // flush to complete before resolving the jsg::Promise. |
| 2304 | return lastFlush.addBranch(); |
| 2305 | } |
| 2306 | |
| 2307 | if (flushesEnqueued > 0) { |
| 2308 | // There is no flush that is scheduled but there is one running, we need to wait for that flush |
| 2309 | // to complete before resolving the jsg::Promise. |
| 2310 | return lastFlush.addBranch(); |
| 2311 | } |
| 2312 | |
| 2313 | // There are no scheduled or in-flight flushes (and there may never have been any), we can return |
| 2314 | // a false Maybe. |
| 2315 | return kj::none; |
| 2316 | } |
| 2317 | |
| 2318 | void ActorCache::shutdown(kj::Maybe<const kj::Exception&> maybeException) { |
| 2319 | if (maybeTerminalException == kj::none) { |
| 2320 | auto exception = [&]() { |
| 2321 | KJ_IF_SOME(e, maybeException) { |
| 2322 | // We were given an exception, use it. |
| 2323 | return e.clone(); |
| 2324 | } |
| 2325 | |
| 2326 | // Use the direct constructor so that we can reuse the constexpr message variable for testing. |
| 2327 | auto exception = kj::Exception(kj::Exception::Type::DISCONNECTED, __FILE__, __LINE__, |
| 2328 | kj::heapString(SHUTDOWN_ERROR_MESSAGE)); |
| 2329 | |
| 2330 | // Add trace info sufficient to tell us which operation caused the failure. |
| 2331 | exception.addTraceHere(); |
| 2332 | exception.addTrace(__builtin_return_address(0)); |
| 2333 | return exception; |
| 2334 | }(); |
| 2335 | |
| 2336 | // Any scheduled flushes will fail once `flushImpl()` is invoked and notices that |
| 2337 | // `maybeTerminalException` has a value. Any in-flight flushes will continue to run in the |
| 2338 | // background. Remember that these in-flight flushes may or may not be awaited by the worker, |
| 2339 | // but they still hold the output lock as long as `allowUnconfirmed` wasn't used. |
| 2340 | maybeTerminalException.emplace(kj::mv(exception)); |
| 2341 | |
| 2342 | // We explicitly do not schedule a flush to break the output gate. This means that if a request |
| 2343 | // is ongoing after the actor cache is shutting down, the output gate is only broken if they |
| 2344 | // had to send a flush after shutdown, either from a scheduled flush or a retry after failure. |
| 2345 | } else { |
| 2346 | // We've already experienced a terminal exception either from shutdown or OOM, there should |
| 2347 | // already be a flush scheduled that will break the output gate. |
| 2348 | } |
| 2349 | } |
| 2350 | |
| 2351 | constexpr size_t bytesToWordsRoundUp(size_t bytes) { |
| 2352 | return (bytes + sizeof(capnp::word) - 1) / sizeof(capnp::word); |
| 2353 | } |
| 2354 | |
| 2355 | namespace { |
| 2356 | using RpcPutRequest = capnp::Request<rpc::ActorStorage::Operations::PutParams, |
| 2357 | rpc::ActorStorage::Operations::PutResults>; |
| 2358 | |
| 2359 | using RpcDeleteRequest = capnp::Request<rpc::ActorStorage::Operations::DeleteParams, |
| 2360 | rpc::ActorStorage::Operations::DeleteResults>; |
| 2361 | } // namespace |
| 2362 | |
| 2363 | kj::Promise<void> ActorCache::startFlushTransaction() { |
| 2364 | // Whenever we flush, we MUST write ALL dirty entries in a single transaction. This is necessary |
| 2365 | // because our cache design doesn't necessarily remember the order in which writes were |
| 2366 | // originally initiated, and thus it's not possible to choose a consistent prefix of writes |
| 2367 | // to transact at once. In particular, when two writes occur on the same key with no (successful) |
| 2368 | // flush in between, the first value is thrown away and never written at all. If we then wanted |
| 2369 | // to perform a partial write that brings storage up-to-date with some point in time between the |
| 2370 | // first and second puts, we wouldn't be able to, because we don't have the old value. |
| 2371 | // |
| 2372 | // Perhaps this would be possible to fix by adding more complex logic. But, it doesn't seem |
| 2373 | // like a big deal to require all flushes to be complete flushes. |
| 2374 | |
| 2375 | // We don't take a lock on `lru.cleanList` here, because we don't need it. We only access |
| 2376 | // `dirtyList`, which is only ever accessed within the actor's thread, so it's safe. We know |
| 2377 | // that `SharedLru` will only ever mess with CLEAN entries, which we don't look at here. |
| 2378 | |
| 2379 | // We have three kinds of writes: Puts, counted deletes, and muted deletes. Counted deletes are |
| 2380 | // delete operations for which the application still wants to know exactly how many keys are |
| 2381 | // actually deleted. We must make a separate RPC call for each counted delete, in order to get |
| 2382 | // the counts back. But we may also have deletes where the application doesn't need to know the |
| 2383 | // count, either because it discarded the promise already, or because we were able to determine |
| 2384 | // the count based on cache. We call these "muted" deletes, and we can batch them all together. |
| 2385 | // We can also batch all the puts together, because applications don't expect puts to return |
| 2386 | // anything. |
| 2387 | // |
| 2388 | // There's another wrinkle, which is that we don't want to send more than 128 keys per batch. |
| 2389 | // This per-batch limit is historically enforced by our storage back-end |
| 2390 | // (supervisor/actor-storage.c++). Truth be told, the limit is artificial and the original |
| 2391 | // motivations for it don't apply anymore. However, splitting huge batches into smaller ones is |
| 2392 | // beneficial to avoid writing overly large capnp messages and other reasons. So, for puts and |
| 2393 | // muted deletes, we go ahead and construct batches of no more than 128 keys. They all end up |
| 2394 | // being part of the same transaction in the end, though. |
| 2395 | // |
| 2396 | // TODO(perf): Currently we send all the batches at the same time. If the batches are large, |
| 2397 | // it could be worth spacing them out a bit so we don't saturate the connection. However, we |
| 2398 | // still need to make sure that the whole transaction represents a consistent snapshot in time, |
| 2399 | // so getting this right, without making a copy of everything upfront, could get complicated. |
| 2400 | // Punting for now. |
| 2401 | |
| 2402 | PutFlush putFlush; |
| 2403 | MutedDeleteFlush mutedDeleteFlush; |
| 2404 | |
| 2405 | auto includeInCurrentBatch = [this](kj::Vector<FlushBatch>& batches, size_t words) { |
| 2406 | KJ_ASSERT(words < MAX_ACTOR_STORAGE_RPC_WORDS); |
| 2407 | |
| 2408 | if (batches.empty()) { |
| 2409 | // This is the first one, let's just set up a current batch. |
| 2410 | batches.add(FlushBatch{}); |
| 2411 | } else if (auto& tailBatch = batches.back(); tailBatch.pairCount >= lru.options.maxKeysPerRpc || |
| 2412 | ((tailBatch.wordCount + words) > MAX_ACTOR_STORAGE_RPC_WORDS)) { |
| 2413 | // We've filled this batch, add a new one. |
| 2414 | batches.add(FlushBatch{}); |
| 2415 | } |
| 2416 | |
| 2417 | auto& batch = batches.back(); |
| 2418 | ++batch.pairCount; |
| 2419 | batch.wordCount += words; |
| 2420 | }; |
| 2421 | |
| 2422 | kj::Vector<CountedDeleteFlush> countedDeleteFlushes(countedDeletes.size()); |
| 2423 | for (auto countedDelete: countedDeletes) { |
| 2424 | if (countedDelete->isFinished) { |
| 2425 | // This countedDelete has already be executed, but we haven't delivered the final count to |
| 2426 | // the waiter yet. We'll skip it here since the destructor of CountedDeleteWaiter should |
| 2427 | // eventually remove this entry from `countedDeletes`. |
| 2428 | continue; |
| 2429 | } |
| 2430 | |
| 2431 | // We might have successfully deleted these entries, but had the broader transaction fail. |
| 2432 | // In that case, we might have entries that have since been overwritten, and which no longer |
| 2433 | // need to be scheduled for deletion. |
| 2434 | kj::Vector<kj::Own<Entry>> entriesToDelete(countedDelete->entries.size()); |
| 2435 | for (auto& entry: countedDelete->entries) { |
| 2436 | if (entry->overwritingCountedDelete && countedDelete->completedInTransaction) { |
| 2437 | // Not only is this a retry, but we have since modified the entry with a put(). |
| 2438 | // Since we already have the delete count, we don't need to delete this entry again. |
| 2439 | continue; |
| 2440 | } |
| 2441 | entriesToDelete.add(kj::mv(entry)); |
| 2442 | } |
| 2443 | |
| 2444 | // We will skip this CountedDelete if there are no entries that need to be deleted. |
| 2445 | // It will be removed from `countedDeletes` by the next flush. |
| 2446 | if (entriesToDelete.empty()) { |
| 2447 | continue; |
| 2448 | } |
| 2449 | |
| 2450 | auto& countedDeleteFlush = countedDeleteFlushes.add(CountedDeleteFlush{ |
| 2451 | .countedDelete = kj::addRef(*countedDelete), |
| 2452 | }); |
| 2453 | // Now that we've filtered our entries down to only those that need to be deleted, |
| 2454 | // we need to overwrite the CountedDelete's `entries`. |
| 2455 | countedDelete->entries = kj::mv(entriesToDelete); |
| 2456 | for (auto& entry: countedDelete->entries) { |
| 2457 | // A delete() call on this key is waiting to find out if the key existed in storage. Since |
| 2458 | // each delete() call needs to return the count of keys deleted, we must issue |
| 2459 | // corresponding delete calls to storage with the same batching, so that storage returns |
| 2460 | // the right counts to us. We can't batch all the deletes into a single delete operation |
| 2461 | // since then we'd only get a single count back and we wouldn't know how to split that up |
| 2462 | // to satisfy all the callers. |
| 2463 | // |
| 2464 | // Note that a subsequent put() call could have set entry.value to non-null, but we still |
| 2465 | // have to perform the delete first in order to determine the count that the delete() call |
| 2466 | // should return. |
| 2467 | // |
| 2468 | // There is a minor quirk here because the counted delete set does not distinguish between |
| 2469 | // before and after a delete all. That's actually okay because we should be able to |
| 2470 | // immediately resolve counted deletes requested after a delete all (either the values are |
| 2471 | // absent or they have a dirty put). This might also be an issue if we respected noCache for |
| 2472 | // delete all's dummy value, but we do not. |
| 2473 | entry->flushStarted = true; |
| 2474 | |
| 2475 | auto keySizeInWords = bytesToWordsRoundUp(entry->key.size()); |
| 2476 | auto words = keySizeInWords + 1; |
| 2477 | includeInCurrentBatch(countedDeleteFlush.batches, words); |
| 2478 | } |
| 2479 | } |
| 2480 | |
| 2481 | auto countEntry = [&](Entry& entry) { |
| 2482 | // Counts up the number of operations and RPC message sizes we'll need to cover this entry. |
| 2483 | |
| 2484 | if (entry.isCountedDelete) { |
| 2485 | // We should have already put this entry into a batch, so just skip it. |
| 2486 | KJ_ASSERT(entry.flushStarted); |
| 2487 | return; |
| 2488 | } |
| 2489 | |
| 2490 | entry.flushStarted = true; |
| 2491 | |
| 2492 | auto keySizeInWords = bytesToWordsRoundUp(entry.key.size()); |
| 2493 | |
| 2494 | KJ_IF_SOME(v, entry.getValuePtr()) { |
| 2495 | auto words = keySizeInWords + bytesToWordsRoundUp(v.size()) + |
| 2496 | capnp::sizeInWords<rpc::ActorStorage::KeyValue>(); |
| 2497 | includeInCurrentBatch(putFlush.batches, words); |
| 2498 | putFlush.entries.add(kj::atomicAddRef(entry)); |
| 2499 | } else { |
| 2500 | auto words = keySizeInWords + 1; |
| 2501 | includeInCurrentBatch(mutedDeleteFlush.batches, words); |
| 2502 | mutedDeleteFlush.entries.add(kj::atomicAddRef(entry)); |
| 2503 | } |
| 2504 | }; |
| 2505 | |
| 2506 | MaybeAlarmChange maybeAlarmChange = CleanAlarm{}; |
| 2507 | KJ_SWITCH_ONEOF(currentAlarmTime) { |
| 2508 | KJ_CASE_ONEOF(knownAlarmTime, ActorCache::KnownAlarmTime) { |
| 2509 | if (knownAlarmTime.status == KnownAlarmTime::Status::DIRTY || |
| 2510 | knownAlarmTime.status == KnownAlarmTime::Status::FLUSHING) { |
| 2511 | knownAlarmTime.status = KnownAlarmTime::Status::FLUSHING; |
| 2512 | maybeAlarmChange = DirtyAlarm{knownAlarmTime.time}; |
| 2513 | } |
| 2514 | } |
| 2515 | KJ_CASE_ONEOF(deferredDelete, ActorCache::DeferredAlarmDelete) { |
| 2516 | if (deferredDelete.status == DeferredAlarmDelete::Status::READY || |
| 2517 | deferredDelete.status == DeferredAlarmDelete::Status::FLUSHING) { |
| 2518 | deferredDelete.status = DeferredAlarmDelete::Status::FLUSHING; |
| 2519 | maybeAlarmChange = DirtyAlarm{kj::none}; |
| 2520 | } |
| 2521 | } |
| 2522 | KJ_CASE_ONEOF(_, UnknownAlarmTime) {} |
| 2523 | } |
| 2524 | |
| 2525 | // We have to remember _before_ waiting for the flush whether or not it was a pre-deleteAll() |
| 2526 | // flush. Otherwise, if it wasn't, but someone calls deleteAll() while we're flushing, then |
| 2527 | // `requestedDeleteAll` might be non-null afterwards, but that would not indicate that we were |
| 2528 | // ready to issue the delete-all. |
| 2529 | KJ_IF_SOME(r, requestedDeleteAll) { |
| 2530 | for (auto& entry: r.deletedDirty) { |
| 2531 | countEntry(*entry); |
| 2532 | } |
| 2533 | } else { |
| 2534 | for (auto& entry: dirtyList) { |
| 2535 | countEntry(entry); |
| 2536 | } |
| 2537 | } |
| 2538 | |
| 2539 | // We don't want to write anything until we know that any past reads have completed, because one |
| 2540 | // or more of those reads could have been on the previous value of a key that was then overwritten |
| 2541 | // by a put() that we're about to flush, and we don't want it to be possible for that read to end |
| 2542 | // up receiving a value that was written later (especially if the read retries due to a |
| 2543 | // disconnect). |
| 2544 | // |
| 2545 | // In practice, most code probably will not have any reads in flight when a flush occurs. |
| 2546 | // |
| 2547 | // Note that we have cached strong references to all entries we intend to mutate above. This means |
| 2548 | // that we can be confident that flushing the cached set will not conflict with future reads |
| 2549 | // because: |
| 2550 | // - All our cached entries are dirty. |
| 2551 | // - Dirty entries can only be removed from the cache map if replaced by a new dirty entry. |
| 2552 | // - Thus all new read requests for our cached entries keys will be served from cache. |
| 2553 | co_await waitForPastReads(); |
| 2554 | |
| 2555 | // Actually flush out the changes. |
| 2556 | auto useTransactionToFlush = [&]() { |
| 2557 | return flushImplUsingTxn(kj::mv(putFlush), kj::mv(mutedDeleteFlush), |
| 2558 | countedDeleteFlushes.releaseAsArray(), kj::mv(maybeAlarmChange)); |
| 2559 | }; |
| 2560 | |
| 2561 | uint typesOfDataToFlush = 0; |
| 2562 | if (!putFlush.batches.empty()) { |
| 2563 | ++typesOfDataToFlush; |
| 2564 | } |
| 2565 | if (!mutedDeleteFlush.batches.empty()) { |
| 2566 | ++typesOfDataToFlush; |
| 2567 | } |
| 2568 | if (!countedDeleteFlushes.empty()) { |
| 2569 | ++typesOfDataToFlush; |
| 2570 | } |
| 2571 | if (maybeAlarmChange.is<DirtyAlarm>()) { |
| 2572 | ++typesOfDataToFlush; |
| 2573 | } |
| 2574 | |
| 2575 | if (typesOfDataToFlush == 0) { |
| 2576 | // Oh, nothing to do. |
| 2577 | } else if (typesOfDataToFlush > 1) { |
| 2578 | // We have multiple types of operations, so we have to use a transaction. |
| 2579 | co_await useTransactionToFlush(); |
| 2580 | } else if (maybeAlarmChange.is<DirtyAlarm>()) { |
| 2581 | // We only had an alarm, we can skip the transaction. |
| 2582 | co_await flushImplAlarmOnly(maybeAlarmChange.get<DirtyAlarm>()); |
| 2583 | } else if (putFlush.batches.size() == 1) { |
| 2584 | // As an optimization for the common case where there are only puts and they all fit in a |
| 2585 | // single batch, just send a simple put rather than complicating things with a transaction. |
| 2586 | co_await flushImplUsingSinglePut(kj::mv(putFlush)); |
| 2587 | } else if (mutedDeleteFlush.batches.size() == 1) { |
| 2588 | // Same as for puts, but for muted deletes. |
| 2589 | co_await flushImplUsingSingleMutedDelete(kj::mv(mutedDeleteFlush)); |
| 2590 | } else if (countedDeleteFlushes.size() == 1 && countedDeleteFlushes[0].batches.size() == 1) { |
| 2591 | // Same as for puts, but for muted deletes. |
| 2592 | co_await flushImplUsingSingleCountedDelete(kj::mv(countedDeleteFlushes[0])); |
| 2593 | } else { |
| 2594 | // None of the special cases above triggered. Default to using a transaction in all other cases, |
| 2595 | // such as when there are so many keys to be flushed that they don't fit into a single batch. |
| 2596 | co_await useTransactionToFlush(); |
| 2597 | } |
| 2598 | } |
| 2599 | |
| 2600 | kj::Promise<void> ActorCache::flushImpl(uint retryCount) { |
| 2601 | KJ_IF_SOME(e, maybeTerminalException) { |
| 2602 | // If we have a terminal exception, throw here to break the output gate and prevent any calls |
| 2603 | // to storage. This does not use `requireNotTerminal()` so that we don't recursively schedule |
| 2604 | // flushes. |
| 2605 | kj::throwFatalException(e.clone()); |
| 2606 | } |
| 2607 | |
| 2608 | auto flushProm = startFlushTransaction(); |
| 2609 | |
| 2610 | bool flushingBeforeDeleteAll = requestedDeleteAll != kj::none; |
| 2611 | return oomCanceler.wrap(kj::mv(flushProm)) |
| 2612 | .then([this, flushingBeforeDeleteAll]() -> kj::Promise<void> { |
| 2613 | // We need to process the alarm result before we (potentially) start the delete all because if |
| 2614 | // we did not our alarm state can't know if it need to flush a new time or not after the delete |
| 2615 | // all. This might be another reason why delete all should not be considered truly deleting the |
| 2616 | // durable object: alarms are not cleared by a delete all. |
| 2617 | KJ_SWITCH_ONEOF(currentAlarmTime) { |
| 2618 | KJ_CASE_ONEOF(knownAlarmTime, ActorCache::KnownAlarmTime) { |
| 2619 | if (knownAlarmTime.status == KnownAlarmTime::Status::FLUSHING) { |
| 2620 | if (knownAlarmTime.noCache) { |
| 2621 | currentAlarmTime = UnknownAlarmTime{}; |
| 2622 | } else { |
| 2623 | knownAlarmTime.status = KnownAlarmTime::Status::CLEAN; |
| 2624 | } |
| 2625 | } |
| 2626 | } |
| 2627 | KJ_CASE_ONEOF(deferredDelete, ActorCache::DeferredAlarmDelete) { |
| 2628 | if (deferredDelete.status == DeferredAlarmDelete::Status::FLUSHING) { |
| 2629 | bool wasDeleted = KJ_ASSERT_NONNULL(deferredDelete.wasDeleted); |
| 2630 | if (deferredDelete.noCache || !wasDeleted) { |
| 2631 | currentAlarmTime = UnknownAlarmTime{}; |
| 2632 | } else { |
| 2633 | currentAlarmTime = KnownAlarmTime{.status = KnownAlarmTime::Status::CLEAN, |
| 2634 | .time = kj::none, |
| 2635 | .noCache = deferredDelete.noCache}; |
| 2636 | } |
| 2637 | } |
| 2638 | } |
| 2639 | KJ_CASE_ONEOF(_, ActorCache::UnknownAlarmTime) {} |
| 2640 | } |
| 2641 | if (flushingBeforeDeleteAll) { |
| 2642 | // The writes we flushed were writes that had occurred before a deleteAll. Now that they are |
| 2643 | // written, we must perform the deleteAll() itself. |
| 2644 | return flushImplDeleteAll(); |
| 2645 | } |
| 2646 | |
| 2647 | auto lock = lru.cleanList.lockExclusive(); |
| 2648 | |
| 2649 | KJ_IF_SOME(r, requestedDeleteAll) { |
| 2650 | // It would appear that all dirty entries were moved into `requestedDeleteAll` during the |
| 2651 | // time that we were waiting for the flushImpl(). We want to remove the flushing entries |
| 2652 | // from that vector now. |
| 2653 | // TODO(cleanup): kj::Vector<T>::filter() would be nice to have here. |
| 2654 | auto dst = r.deletedDirty.begin(); |
| 2655 | for (auto src = r.deletedDirty.begin(); src != r.deletedDirty.end(); ++src) { |
| 2656 | if (!src->get()->flushStarted) { |
| 2657 | if (dst != src) *dst = kj::mv(*src); |
| 2658 | ++dst; |
| 2659 | } |
| 2660 | } |
| 2661 | r.deletedDirty.resize(dst - r.deletedDirty.begin()); |
| 2662 | } else { |
| 2663 | // Mark all flushing entries as `CLEAN`. Note that we know that all flushing entries must |
| 2664 | // form a prefix of `dirtyList` since any new entries would have been added to the end. |
| 2665 | for (auto& entry: dirtyList) { |
| 2666 | if (!entry.flushStarted) { |
| 2667 | // Completed all flushing entries. |
| 2668 | break; |
| 2669 | } |
| 2670 | |
| 2671 | KJ_ASSERT(entry.flushStarted); |
| 2672 | |
| 2673 | // We know all `countedDelete` operations were satisfied so we can remove this if it's |
| 2674 | // present. The `CountedDeleteWaiters` will resolve once the flush is finished, and will |
| 2675 | // remove the `CountedDelete`s from `countedDeletes`. Even if it doesn't happen by the |
| 2676 | // next flush, each `CountedDelete` should have `isFinished` set so even if we encounter it |
| 2677 | // next flush we won't attempt to delete again. |
| 2678 | entry.isCountedDelete = false; |
| 2679 | |
| 2680 | dirtyList.remove(entry); |
| 2681 | if (entry.noCache) { |
| 2682 | entry.setNotInCache(); |
| 2683 | evictEntry(lock, entry); |
| 2684 | } else { |
| 2685 | if (entry.gapIsKnownEmpty && entry.getValueStatus() == EntryValueStatus::ABSENT) { |
| 2686 | // This is a negative entry, and is followed by a known-empty gap. If the previous entry |
| 2687 | // also has `gapIsKnownEmpty`, then this entry is entirely redundant. |
| 2688 | auto& map = currentValues.get(lock); |
| 2689 | auto entryIter = map.seek(entry.key); |
| 2690 | KJ_ASSERT(entryIter->get() == &entry); |
| 2691 | |
| 2692 | if (entryIter != map.ordered().begin()) { |
| 2693 | auto prevIter = entryIter; |
| 2694 | --prevIter; |
| 2695 | if (prevIter->get()->gapIsKnownEmpty) { |
| 2696 | // Yep! |
| 2697 | entry.setNotInCache(); |
| 2698 | map.erase(*entryIter); |
| 2699 | // WARNING: We might have just deleted `entry`. |
| 2700 | continue; |
| 2701 | } |
| 2702 | } |
| 2703 | } |
| 2704 | |
| 2705 | addToCleanList(lock, entry); |
| 2706 | } |
| 2707 | } |
| 2708 | } |
| 2709 | |
| 2710 | evictOrOomIfNeeded(lock); |
| 2711 | |
| 2712 | return kj::READY_NOW; |
| 2713 | }, [this, retryCount](kj::Exception&& e) -> kj::Promise<void> { |
| 2714 | static const size_t MAX_RETRIES = 4; |
| 2715 | if (e.getType() == kj::Exception::Type::DISCONNECTED && retryCount < MAX_RETRIES) { |
| 2716 | return flushImpl(retryCount + 1); |
| 2717 | } else if (jsg::isTunneledException(e.getDescription()) || |
| 2718 | jsg::isDoNotLogException(e.getDescription())) { |
| 2719 | // Before passing along the exception, give it the proper brokenness reason. |
| 2720 | // We were overriding any exception that came through here by ioGateBroken (now outputGateBroken). |
| 2721 | // without checking for previous brokenness reasons we would be unable to throw |
| 2722 | // exceededConcurrentStorageOps at all. |
| 2723 | auto msg = jsg::stripRemoteExceptionPrefix(e.getDescription()); |
| 2724 | if (!(msg.startsWith("broken."))) { |
| 2725 | e.setDescription(kj::str("broken.outputGateBroken; ", msg)); |
| 2726 | } |
| 2727 | return kj::mv(e); |
| 2728 | } else { |
| 2729 | if (isInterestingException(e)) { |
| 2730 | LOG_EXCEPTION("actorCacheFlush", e); |
| 2731 | } else { |
| 2732 | LOG_NOSENTRY(ERROR, "actor cache flush failed", e); |
| 2733 | } |
| 2734 | // Pass through exception type to convey appropriate retry behavior. |
| 2735 | return kj::Exception(e.getType(), __FILE__, __LINE__, |
| 2736 | kj::str("broken.outputGateBroken; jsg.Error: Internal error in Durable " |
| 2737 | "Object storage write caused object to be reset.")); |
| 2738 | } |
| 2739 | }); |
| 2740 | } |
| 2741 | |
| 2742 | kj::Promise<void> ActorCache::flushImplUsingSinglePut(PutFlush putFlush) { |
| 2743 | KJ_ASSERT(putFlush.batches.size() == 1); |
| 2744 | auto& batch = putFlush.batches[0]; |
| 2745 | |
| 2746 | KJ_ASSERT(batch.wordCount < MAX_ACTOR_STORAGE_RPC_WORDS); |
| 2747 | KJ_ASSERT(batch.pairCount == putFlush.entries.size()); |
| 2748 | |
| 2749 | auto request = storage.putRequest(capnp::MessageSize{4 + batch.wordCount, 0}); |
| 2750 | auto list = request.initEntries(batch.pairCount); |
| 2751 | auto entryIt = putFlush.entries.begin(); |
| 2752 | for (auto kv: list) { |
| 2753 | auto& entry = **(entryIt++); |
| 2754 | auto v = KJ_ASSERT_NONNULL(entry.getValuePtr()); |
| 2755 | kv.setKey(entry.key.asBytes()); |
| 2756 | kv.setValue(v); |
| 2757 | } |
| 2758 | |
| 2759 | // We're done with the batching instructions, free them before we go async. |
| 2760 | putFlush.entries.clear(); |
| 2761 | putFlush.batches.clear(); |
| 2762 | { |
| 2763 | auto writeObserver = recordStorageWrite(hooks, clock); |
| 2764 | util::DurationExceededLogger logger( |
| 2765 | clock, 1 * kj::SECONDS, "storage operation took longer than expected: single put"); |
| 2766 | co_await request.sendIgnoringResult(); |
| 2767 | } |
| 2768 | } |
| 2769 | |
| 2770 | kj::Promise<void> ActorCache::flushImplUsingSingleMutedDelete(MutedDeleteFlush mutedFlush) { |
| 2771 | KJ_ASSERT(mutedFlush.batches.size() == 1); |
| 2772 | auto& batch = mutedFlush.batches[0]; |
| 2773 | |
| 2774 | KJ_ASSERT(batch.wordCount < MAX_ACTOR_STORAGE_RPC_WORDS); |
| 2775 | KJ_ASSERT(batch.pairCount == mutedFlush.entries.size()); |
| 2776 | |
| 2777 | auto request = storage.deleteRequest(capnp::MessageSize{4 + batch.wordCount, 0}); |
| 2778 | auto listBuilder = request.initKeys(batch.pairCount); |
| 2779 | auto entryIt = mutedFlush.entries.begin(); |
| 2780 | for (size_t i = 0; i < batch.pairCount; ++i) { |
| 2781 | auto& entry = **(entryIt++); |
| 2782 | listBuilder.set(i, entry.key.asBytes()); |
| 2783 | } |
| 2784 | |
| 2785 | // We're done with the batching instructions, free them before we go async. |
| 2786 | mutedFlush.entries.clear(); |
| 2787 | mutedFlush.batches.clear(); |
| 2788 | |
| 2789 | { |
| 2790 | auto writeObserver = recordStorageWrite(hooks, clock); |
| 2791 | util::DurationExceededLogger logger( |
| 2792 | clock, 1 * kj::SECONDS, "storage operation took longer than expected: muted delete"); |
| 2793 | co_await request.sendIgnoringResult(); |
| 2794 | } |
| 2795 | } |
| 2796 | |
| 2797 | kj::Promise<void> ActorCache::flushImplUsingSingleCountedDelete(CountedDeleteFlush countedFlush) { |
| 2798 | KJ_ASSERT(countedFlush.batches.size() == 1); |
| 2799 | auto& batch = countedFlush.batches[0]; |
| 2800 | |
| 2801 | auto countedDelete = kj::mv(countedFlush.countedDelete); |
| 2802 | KJ_ASSERT(batch.wordCount < MAX_ACTOR_STORAGE_RPC_WORDS); |
| 2803 | KJ_ASSERT(batch.pairCount == countedDelete->entries.size()); |
| 2804 | |
| 2805 | auto request = storage.deleteRequest(capnp::MessageSize{4 + batch.wordCount, 0}); |
| 2806 | auto listBuilder = request.initKeys(batch.pairCount); |
| 2807 | auto entryIt = countedDelete->entries.begin(); |
| 2808 | for (size_t i = 0; i < batch.pairCount; ++i) { |
| 2809 | auto& entry = **(entryIt++); |
| 2810 | listBuilder.set(i, entry.key.asBytes()); |
| 2811 | } |
| 2812 | |
| 2813 | // We're done with the batching instructions, free them before we go async. |
| 2814 | countedFlush.batches.clear(); |
| 2815 | |
| 2816 | auto writeObserver = recordStorageWrite(hooks, clock); |
| 2817 | util::DurationExceededLogger logger( |
| 2818 | clock, 1 * kj::SECONDS, "storage operation took longer than expected: counted delete"); |
| 2819 | auto response = co_await request.send(); |
| 2820 | countedDelete->countDeleted += response.getNumDeleted(); |
| 2821 | countedDelete->isFinished = true; |
| 2822 | } |
| 2823 | |
| 2824 | kj::Promise<void> ActorCache::flushImplAlarmOnly(DirtyAlarm dirty) { |
| 2825 | auto writeObserver = recordStorageWrite(hooks, clock); |
| 2826 | util::DurationExceededLogger logger( |
| 2827 | clock, 1 * kj::SECONDS, "storage operation took longer than expected: set/delete alarm"); |
| 2828 | |
| 2829 | // TODO(someday) This could be templated to reuse the same code for this and the transaction case. |
| 2830 | // Handle alarm writes first, since they're simplest. |
| 2831 | KJ_IF_SOME(newTime, dirty.newTime) { |
| 2832 | auto req = storage.setAlarmRequest(); |
| 2833 | req.setScheduledTimeMs((newTime - kj::UNIX_EPOCH) / kj::MILLISECONDS); |
| 2834 | co_await req.sendIgnoringResult(); |
| 2835 | co_return; |
| 2836 | } else { |
| 2837 | // Alarm deletes are a bit trickier because we have to take DeferredAlarmDeletes into account. |
| 2838 | auto req = storage.deleteAlarmRequest(); |
| 2839 | KJ_IF_SOME(deferredDelete, currentAlarmTime.tryGet<DeferredAlarmDelete>()) { |
| 2840 | if (deferredDelete.status == DeferredAlarmDelete::Status::FLUSHING) { |
| 2841 | req.setTimeToDeleteMs((deferredDelete.timeToDelete - kj::UNIX_EPOCH) / kj::MILLISECONDS); |
| 2842 | auto response = co_await req.send(); |
| 2843 | KJ_IF_SOME(deferredDelete, currentAlarmTime.tryGet<DeferredAlarmDelete>()) { |
| 2844 | if (deferredDelete.status == DeferredAlarmDelete::Status::FLUSHING) { |
| 2845 | // We always update wasDeleted regardless of whether or not it is true |
| 2846 | // because this continuation can succeed even if the greater transaction fails, |
| 2847 | // and so we want to make sure we end up with the correct value if the first |
| 2848 | // attempt succeeds to delete, the txn fails, and the retry fails to delete. |
| 2849 | // The early update is OK because we don't actually use the incorrect state until |
| 2850 | // the transaction succeeds in the .then() below. |
| 2851 | deferredDelete.wasDeleted = response.getDeleted(); |
| 2852 | } |
| 2853 | } |
| 2854 | } else { |
| 2855 | // Not sending a delete request for WAITING or READY is intentional. The WAITING |
| 2856 | // state refers to when the alarm run has started but has not completed successfully, |
| 2857 | // and READY is set when the run completes -- only FLUSHING indicates we actually |
| 2858 | // need to send a request. |
| 2859 | } |
| 2860 | } else { |
| 2861 | co_await req.send(); |
| 2862 | } |
| 2863 | } |
| 2864 | } |
| 2865 | |
| 2866 | kj::Promise<void> ActorCache::flushImplUsingTxn(PutFlush putFlush, |
| 2867 | MutedDeleteFlush mutedDeleteFlush, |
| 2868 | CountedDeleteFlushes countedDeleteFlushes, |
| 2869 | MaybeAlarmChange maybeAlarmChange) { |
| 2870 | auto txnProm = storage.txnRequest(capnp::MessageSize{4, 0}).send(); |
| 2871 | auto txn = txnProm.getTransaction(); |
| 2872 | |
| 2873 | struct RpcCountedDelete { |
| 2874 | kj::Own<CountedDelete> countedDelete; |
| 2875 | kj::Array<RpcDeleteRequest> rpcDeletes; |
| 2876 | }; |
| 2877 | auto rpcCountedDeletes = kj::heapArrayBuilder<RpcCountedDelete>(countedDeleteFlushes.size()); |
| 2878 | auto rpcMutedDeletes = kj::heapArrayBuilder<RpcDeleteRequest>(mutedDeleteFlush.batches.size()); |
| 2879 | auto rpcPuts = kj::heapArrayBuilder<RpcPutRequest>(putFlush.batches.size()); |
| 2880 | |
| 2881 | for (auto& flush: countedDeleteFlushes) { |
| 2882 | auto countedDelete = kj::mv(flush.countedDelete); |
| 2883 | auto entryIt = countedDelete->entries.begin(); |
| 2884 | kj::Vector<RpcDeleteRequest> rpcDeletes; |
| 2885 | for (auto& batch: flush.batches) { |
| 2886 | KJ_ASSERT(batch.wordCount < MAX_ACTOR_STORAGE_RPC_WORDS); |
| 2887 | |
| 2888 | auto request = txn.deleteRequest(capnp::MessageSize{4 + batch.wordCount, 0}); |
| 2889 | auto listBuilder = request.initKeys(batch.pairCount); |
| 2890 | for (size_t i = 0; i < batch.pairCount; ++i) { |
| 2891 | KJ_ASSERT(entryIt != countedDelete->entries.end()); |
| 2892 | auto& entry = **(entryIt++); |
| 2893 | listBuilder.set(i, entry.key.asBytes()); |
| 2894 | } |
| 2895 | |
| 2896 | rpcDeletes.add(kj::mv(request)); |
| 2897 | } |
| 2898 | KJ_ASSERT(entryIt == countedDelete->entries.end()); |
| 2899 | rpcCountedDeletes.add(RpcCountedDelete{ |
| 2900 | .countedDelete = kj::mv(countedDelete), |
| 2901 | .rpcDeletes = rpcDeletes.releaseAsArray(), |
| 2902 | }); |
| 2903 | } |
| 2904 | countedDeleteFlushes = nullptr; |
| 2905 | |
| 2906 | { |
| 2907 | auto entryIt = mutedDeleteFlush.entries.begin(); |
| 2908 | for (auto& batch: mutedDeleteFlush.batches) { |
| 2909 | KJ_ASSERT(batch.wordCount < MAX_ACTOR_STORAGE_RPC_WORDS); |
| 2910 | |
| 2911 | auto request = txn.deleteRequest(capnp::MessageSize{4 + batch.wordCount, 0}); |
| 2912 | auto listBuilder = request.initKeys(batch.pairCount); |
| 2913 | for (size_t i = 0; i < batch.pairCount; ++i) { |
| 2914 | KJ_ASSERT(entryIt != mutedDeleteFlush.entries.end()); |
| 2915 | auto& entry = **(entryIt++); |
| 2916 | listBuilder.set(i, entry.key.asBytes()); |
| 2917 | } |
| 2918 | rpcMutedDeletes.add(kj::mv(request)); |
| 2919 | } |
| 2920 | KJ_ASSERT(entryIt == mutedDeleteFlush.entries.end()); |
| 2921 | } |
| 2922 | mutedDeleteFlush.entries.clear(); |
| 2923 | mutedDeleteFlush.batches.clear(); |
| 2924 | |
| 2925 | { |
| 2926 | auto entryIt = putFlush.entries.begin(); |
| 2927 | for (auto& batch: putFlush.batches) { |
| 2928 | KJ_ASSERT(batch.wordCount < MAX_ACTOR_STORAGE_RPC_WORDS); |
| 2929 | |
| 2930 | auto request = txn.putRequest(capnp::MessageSize{4 + batch.wordCount, 0}); |
| 2931 | auto listBuilder = request.initEntries(batch.pairCount); |
| 2932 | for (auto kv: listBuilder) { |
| 2933 | KJ_ASSERT(entryIt != putFlush.entries.end()); |
| 2934 | auto& entry = **(entryIt++); |
| 2935 | auto v = KJ_ASSERT_NONNULL(entry.getValuePtr()); |
| 2936 | kv.setKey(entry.key.asBytes()); |
| 2937 | kv.setValue(v); |
| 2938 | } |
| 2939 | rpcPuts.add(kj::mv(request)); |
| 2940 | } |
| 2941 | KJ_ASSERT(entryIt == putFlush.entries.end()); |
| 2942 | } |
| 2943 | putFlush.entries.clear(); |
| 2944 | putFlush.batches.clear(); |
| 2945 | |
| 2946 | // Send all the RPCs. It's important that counted deletes are sent first since they can overlap |
| 2947 | // with puts. Specifically this can happen if someone does a delete() immediately followed by a |
| 2948 | // put() on the same key. These two writes may have been coalesced into a single flush. |
| 2949 | // Unfortunately, we can't just skip the delete because we still need to count it. So we issue |
| 2950 | // a delete, followed by a put, in the same transaction. |
| 2951 | // The constant extra 2 promises are those added outside of the rpc batches, currently one |
| 2952 | // to work around a bug in capnp::autoreconnect, and one to actually commit the flush txn |
| 2953 | // A 3rd promise may be added to write the alarm time if necessary. |
| 2954 | auto promises = kj::heapArrayBuilder<kj::Promise<void>>(rpcPuts.size() + rpcMutedDeletes.size() + |
| 2955 | rpcCountedDeletes.size() + 2 + !maybeAlarmChange.is<CleanAlarm>()); |
| 2956 | |
| 2957 | auto joinCountedDelete = [](RpcCountedDelete& rpcCountedDelete) -> kj::Promise<void> { |
| 2958 | auto promises = KJ_MAP(request, rpcCountedDelete.rpcDeletes) { |
| 2959 | return request.send().then( |
| 2960 | [](capnp::Response<rpc::ActorStorage::Operations::DeleteResults>&& response) mutable |
| 2961 | -> uint { return response.getNumDeleted(); }); |
| 2962 | }; |
| 2963 | |
| 2964 | size_t recordsDeleted = 0; |
| 2965 | for (auto& promise: promises) { |
| 2966 | recordsDeleted += co_await promise; |
| 2967 | } |
| 2968 | |
| 2969 | // This may be a retry following a successful counted delete within a failed transaction. |
| 2970 | // In that case, we don't want to update the count again, since we've already considered it. |
| 2971 | if (!rpcCountedDelete.countedDelete->completedInTransaction) { |
| 2972 | // We only increment our `countDeleted` if *ALL* the delete batches succeeded. |
| 2973 | rpcCountedDelete.countedDelete->countDeleted += recordsDeleted; |
| 2974 | } |
| 2975 | |
| 2976 | // This delete succeeded, but we may need to retry it in some cases, ex. if the transaction fails. |
| 2977 | // If we *do* retry after a successful counted delete, we won't want to update our |
| 2978 | // `countDeleted` since we already got it. |
| 2979 | rpcCountedDelete.countedDelete->completedInTransaction = true; |
| 2980 | }; |
| 2981 | |
| 2982 | for (auto& rpcCountedDelete: rpcCountedDeletes) { |
| 2983 | promises.add(joinCountedDelete(rpcCountedDelete)); |
| 2984 | } |
| 2985 | |
| 2986 | for (auto& request: rpcMutedDeletes) { |
| 2987 | promises.add(request.sendIgnoringResult()); |
| 2988 | } |
| 2989 | |
| 2990 | for (auto& request: rpcPuts) { |
| 2991 | promises.add(request.sendIgnoringResult()); |
| 2992 | } |
| 2993 | |
| 2994 | KJ_SWITCH_ONEOF(maybeAlarmChange) { |
| 2995 | KJ_CASE_ONEOF(dirty, DirtyAlarm) { |
| 2996 | KJ_IF_SOME(newTime, dirty.newTime) { |
| 2997 | auto req = txn.setAlarmRequest(); |
| 2998 | req.setScheduledTimeMs((newTime - kj::UNIX_EPOCH) / kj::MILLISECONDS); |
| 2999 | promises.add(req.sendIgnoringResult()); |
| 3000 | } else { |
| 3001 | auto req = txn.deleteAlarmRequest(); |
| 3002 | KJ_IF_SOME(deferredDelete, currentAlarmTime.tryGet<DeferredAlarmDelete>()) { |
| 3003 | if (deferredDelete.status == DeferredAlarmDelete::Status::FLUSHING) { |
| 3004 | req.setTimeToDeleteMs( |
| 3005 | (deferredDelete.timeToDelete - kj::UNIX_EPOCH) / kj::MILLISECONDS); |
| 3006 | auto prom = req.send().then([this](auto response) { |
| 3007 | KJ_IF_SOME(deferredDelete, currentAlarmTime.tryGet<DeferredAlarmDelete>()) { |
| 3008 | if (deferredDelete.status == DeferredAlarmDelete::Status::FLUSHING) { |
| 3009 | // We always update wasDeleted regardless of whether or not it is true |
| 3010 | // because this continuation can succeed even if the greater transaction fails, |
| 3011 | // and so we want to make sure we end up with the correct value if the first |
| 3012 | // attempt succeeds to delete, the txn fails, and the retry fails to delete. |
| 3013 | // The early update is OK because we don't actually use the incorrect state until |
| 3014 | // the transaction succeeds in the .then() below. |
| 3015 | deferredDelete.wasDeleted = response.getDeleted(); |
| 3016 | } |
| 3017 | } |
| 3018 | }); |
| 3019 | promises.add(kj::mv(prom)); |
| 3020 | } |
| 3021 | // Not sending a delete request for WAITING or READY is intentional. The WAITING |
| 3022 | // state refers to when the alarm run has started but has not completed successfully, |
| 3023 | // and READY is set when the run completes -- only FLUSHING indicates we actually |
| 3024 | // need to send a request. |
| 3025 | } else { |
| 3026 | promises.add(req.sendIgnoringResult()); |
| 3027 | } |
| 3028 | } |
| 3029 | } |
| 3030 | KJ_CASE_ONEOF(_, CleanAlarm) {} |
| 3031 | } |
| 3032 | |
| 3033 | // We have to wait on the transaction promise so we don't cancel the catch_ branch that triggers |
| 3034 | // our autoReconnect logic on storage failures. |
| 3035 | // TODO(cleanup): We should probably fix ReconnectHook so the catch_ doesn't get canceled |
| 3036 | // if the promise is dropped but the pipeline stays alive. |
| 3037 | promises.add(txnProm.ignoreResult()); |
| 3038 | |
| 3039 | { |
| 3040 | auto writeObserver = recordStorageWrite(hooks, clock); |
| 3041 | util::DurationExceededLogger logger(clock, 1 * kj::SECONDS, |
| 3042 | "storage operation took longer than expected: commit flush transaction"); |
| 3043 | promises.add(txn.commitRequest(capnp::MessageSize{4, 0}).sendIgnoringResult()); |
| 3044 | |
| 3045 | co_await kj::joinPromises(promises.finish()); |
| 3046 | for (auto& rpcCountedDelete: rpcCountedDeletes) { |
| 3047 | // Now that the transaction has successfully completed, we can mark all our CountedDeletes |
| 3048 | // as having completed as well. |
| 3049 | rpcCountedDelete.countedDelete->isFinished = true; |
| 3050 | } |
| 3051 | } |
| 3052 | } |
| 3053 | |
| 3054 | kj::Promise<void> ActorCache::flushImplDeleteAll(uint retryCount) { |
| 3055 | // By this point, we've completed any writes that had originally been performed before |
| 3056 | // deleteAll() was called, and we're ready to perform the deleteAll() itself. |
| 3057 | // |
| 3058 | // Note that we intentionally don't time deleteAll() with hooks.startStorageWrite() because it's |
| 3059 | // expected to be much slower than all other storage operations, taking linear time with respect |
| 3060 | // to how much data is stored in the actor. |
| 3061 | |
| 3062 | KJ_ASSERT(requestedDeleteAll != kj::none); |
| 3063 | |
| 3064 | return storage.deleteAllRequest(capnp::MessageSize{2, 0}) |
| 3065 | .send() |
| 3066 | .then( |
| 3067 | [this](capnp::Response<rpc::ActorStorage::Operations::DeleteAllResults> results) |
| 3068 | -> kj::Promise<void> { |
| 3069 | auto& deleteAllState = KJ_ASSERT_NONNULL(requestedDeleteAll); |
| 3070 | deleteAllState.countFulfiller->fulfill(results.getNumDeleted()); |
| 3071 | bool shouldDeleteAlarm = deleteAllState.deleteAlarm; |
| 3072 | |
| 3073 | // Success! We can now null out `requestedDeleteAll`. Note that we don't have to worry about |
| 3074 | // `requestedDeleteAll` having changed since we flushed it earlier, because it can't change |
| 3075 | // until it is first nulled out. If deleteAll() is called multiple times before the first one |
| 3076 | // finishes, subsequent ones see `requestedDeleteAll` is already non-null and they don't change |
| 3077 | // it. Instead, the writes that occurred between the deleteAll()s are simply discarded, as if |
| 3078 | // the two deleteAll()s had been coalesced into a single one. |
| 3079 | requestedDeleteAll = kj::none; |
| 3080 | |
| 3081 | if (shouldDeleteAlarm) { |
| 3082 | // The deleteAll() was requested with deleteAlarm=true. Now that the deleteAll RPC has |
| 3083 | // succeeded, we dirty the alarm to kj::none so that it will be flushed in the post-deleteAll |
| 3084 | // flushImpl() call below. This ordering ensures the alarm is only deleted after KV data has |
| 3085 | // been successfully deleted. |
| 3086 | // |
| 3087 | // However, if currentAlarmTime is already DIRTY, that means a subsequent setAlarm() call |
| 3088 | // was made after the deleteAll() was initiated, and that newer alarm value should take |
| 3089 | // precedence over the deletion. |
| 3090 | bool alreadyDirty = false; |
| 3091 | KJ_IF_SOME(known, currentAlarmTime.tryGet<KnownAlarmTime>()) { |
| 3092 | alreadyDirty = known.status == KnownAlarmTime::Status::DIRTY; |
| 3093 | } |
| 3094 | if (!alreadyDirty) { |
| 3095 | currentAlarmTime = ActorCache::KnownAlarmTime{ |
| 3096 | .status = ActorCache::KnownAlarmTime::Status::DIRTY, |
| 3097 | .time = kj::none, |
| 3098 | }; |
| 3099 | } |
| 3100 | } |
| 3101 | |
| 3102 | { |
| 3103 | auto lock = lru.cleanList.lockExclusive(); |
| 3104 | evictOrOomIfNeeded(lock); |
| 3105 | } |
| 3106 | |
| 3107 | // Now we must flush any writes that happened after the deleteAll(). (If there are none, this |
| 3108 | // will complete quickly.) |
| 3109 | // TODO(soon) This will use the write options for the deleteAll() even if the options for future |
| 3110 | // operations differ. This can mean that we will not wait for the output gate when we were asked |
| 3111 | // to do so. We should fix this. |
| 3112 | return flushImpl(); |
| 3113 | }, |
| 3114 | [this, retryCount](kj::Exception&& e) -> kj::Promise<void> { |
| 3115 | static const size_t MAX_RETRIES = 4; |
| 3116 | if (e.getType() == kj::Exception::Type::DISCONNECTED && retryCount < MAX_RETRIES) { |
| 3117 | return flushImplDeleteAll(retryCount + 1); |
| 3118 | } else if (jsg::isTunneledException(e.getDescription()) || |
| 3119 | jsg::isDoNotLogException(e.getDescription())) { |
| 3120 | // Before passing along the exception, give it the proper brokenness reason. |
| 3121 | auto msg = jsg::stripRemoteExceptionPrefix(e.getDescription()); |
| 3122 | e.setDescription(kj::str("broken.outputGateBroken; ", msg)); |
| 3123 | return kj::mv(e); |
| 3124 | } else { |
| 3125 | if (isInterestingException(e)) { |
| 3126 | LOG_EXCEPTION("actorCacheDeleteAll", e); |
| 3127 | } else { |
| 3128 | LOG_NOSENTRY(ERROR, "actorCacheDeleteAll failed", e); |
| 3129 | } |
| 3130 | // Pass through exception type to convey appropriate retry behavior. |
| 3131 | return kj::Exception(e.getType(), __FILE__, __LINE__, |
| 3132 | kj::str( |
| 3133 | "broken.outputGateBroken; jsg.Error: Internal error in Durable Object storage deleteAll() caused object to be reset.")); |
| 3134 | } |
| 3135 | }); |
| 3136 | } |
| 3137 | |
| 3138 | // ======================================================================================= |
| 3139 | // ActorCache::Transaction |
| 3140 | |
| 3141 | ActorCache::Transaction::Transaction(ActorCache& cache): cache(cache) {} |
| 3142 | ActorCache::Transaction::~Transaction() noexcept(false) { |
| 3143 | // If not commit()ed... we don't have to do anything in particular here, just drop the changes. |
| 3144 | } |
| 3145 | |
| 3146 | kj::Maybe<kj::Promise<void>> ActorCache::Transaction::commit() { |
| 3147 | { |
| 3148 | auto lock = cache.lru.cleanList.lockExclusive(); |
| 3149 | for (auto& change: entriesToWrite) { |
| 3150 | cache.putImpl(lock, kj::mv(change.entry), change.options, kj::none, commitSpan.addRef()); |
| 3151 | } |
| 3152 | entriesToWrite.clear(); |
| 3153 | cache.evictOrOomIfNeeded(lock); |
| 3154 | } |
| 3155 | |
| 3156 | KJ_IF_SOME(change, alarmChange) { |
| 3157 | cache.setAlarm(change.newTime, change.options, commitSpan.addRef()); |
| 3158 | } |
| 3159 | alarmChange = kj::none; |
| 3160 | |
| 3161 | return cache.getBackpressure(); |
| 3162 | } |
| 3163 | |
| 3164 | kj::Promise<void> ActorCache::Transaction::rollback() { |
| 3165 | entriesToWrite.clear(); |
| 3166 | alarmChange = kj::none; |
| 3167 | return kj::READY_NOW; |
| 3168 | } |
| 3169 | |
| 3170 | // ----------------------------------------------------------------------------- |
| 3171 | // transaction reads |
| 3172 | |
| 3173 | kj::OneOf<kj::Maybe<ActorCache::Value>, kj::Promise<kj::Maybe<ActorCache::Value>>> ActorCache:: |
| 3174 | Transaction::get(Key key, ReadOptions options) { |
| 3175 | options.noCache = options.noCache || cache.lru.options.noCache; |
| 3176 | KJ_IF_SOME(change, entriesToWrite.find(key)) { |
| 3177 | return change.entry->getValue(); |
| 3178 | } else { |
| 3179 | return cache.get(kj::mv(key), options); |
| 3180 | } |
| 3181 | } |
| 3182 | |
| 3183 | kj::OneOf<ActorCache::GetResultList, kj::Promise<ActorCache::GetResultList>> ActorCache:: |
| 3184 | Transaction::get(kj::Array<Key> keys, ReadOptions options) { |
| 3185 | options.noCache = options.noCache || cache.lru.options.noCache; |
| 3186 | |
| 3187 | kj::Vector<kj::Own<Entry>> changedEntries; |
| 3188 | kj::Vector<Key> keysToFetch; |
| 3189 | |
| 3190 | for (auto& key: keys) { |
| 3191 | KJ_IF_SOME(change, entriesToWrite.find(key)) { |
| 3192 | changedEntries.add(kj::atomicAddRef(*change.entry)); |
| 3193 | } else { |
| 3194 | keysToFetch.add(kj::mv(key)); |
| 3195 | } |
| 3196 | } |
| 3197 | |
| 3198 | std::sort(changedEntries.begin(), changedEntries.end(), |
| 3199 | [](auto& a, auto& b) { return a.get()->key < b.get()->key; }); |
| 3200 | |
| 3201 | return merge(kj::mv(changedEntries), cache.get(keysToFetch.releaseAsArray(), options), |
| 3202 | GetResultList::FORWARD); |
| 3203 | } |
| 3204 | |
| 3205 | kj::OneOf<kj::Maybe<kj::Date>, kj::Promise<kj::Maybe<kj::Date>>> ActorCache::Transaction::getAlarm( |
| 3206 | ReadOptions options) { |
| 3207 | options.noCache = options.noCache || cache.lru.options.noCache; |
| 3208 | KJ_IF_SOME(a, alarmChange) { |
| 3209 | return a.newTime; |
| 3210 | } else { |
| 3211 | return cache.getAlarm(options); |
| 3212 | } |
| 3213 | } |
| 3214 | |
| 3215 | kj::OneOf<ActorCache::GetResultList, kj::Promise<ActorCache::GetResultList>> ActorCache:: |
| 3216 | Transaction::list(Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) { |
| 3217 | options.noCache = options.noCache || cache.lru.options.noCache; |
| 3218 | kj::Vector<kj::Own<Entry>> changedEntries; |
| 3219 | if (limit.orDefault(kj::maxValue) == 0 || begin >= end) { |
| 3220 | // No results in these cases, just return. |
| 3221 | return ActorCache::GetResultList(kj::mv(changedEntries), {}, GetResultList::REVERSE); |
| 3222 | } |
| 3223 | auto beginIter = entriesToWrite.seek(begin); |
| 3224 | auto endIter = seekOrEnd(entriesToWrite, end); |
| 3225 | uint positiveCount = 0; |
| 3226 | // TODO(cleanup): Add `iterRange()` to KJ's public interface. |
| 3227 | for (auto& change: kj::_::iterRange(beginIter, endIter)) { |
| 3228 | changedEntries.add(kj::atomicAddRef(*change.entry)); |
| 3229 | if (change.entry->getValueStatus() == EntryValueStatus::PRESENT) { |
| 3230 | ++positiveCount; |
| 3231 | } |
| 3232 | if (positiveCount == limit.orDefault(kj::maxValue)) break; |
| 3233 | } |
| 3234 | |
| 3235 | // Increase limit to make sure it can't be underrun by negative entries negating it. |
| 3236 | limit = limit.map([&](uint n) { return n + (changedEntries.size() - positiveCount); }); |
| 3237 | |
| 3238 | return merge(kj::mv(changedEntries), cache.list(kj::mv(begin), kj::mv(end), limit, options), |
| 3239 | GetResultList::FORWARD); |
| 3240 | } |
| 3241 | |
| 3242 | kj::OneOf<ActorCache::GetResultList, kj::Promise<ActorCache::GetResultList>> ActorCache:: |
| 3243 | Transaction::listReverse( |
| 3244 | Key begin, kj::Maybe<Key> end, kj::Maybe<uint> limit, ReadOptions options) { |
| 3245 | options.noCache = options.noCache || cache.lru.options.noCache; |
| 3246 | kj::Vector<kj::Own<Entry>> changedEntries; |
| 3247 | if (limit.orDefault(kj::maxValue) == 0 || begin >= end) { |
| 3248 | // No results in these cases, just return. |
| 3249 | return ActorCache::GetResultList(kj::mv(changedEntries), {}, GetResultList::REVERSE); |
| 3250 | } |
| 3251 | auto beginIter = entriesToWrite.seek(begin); |
| 3252 | auto endIter = seekOrEnd(entriesToWrite, end); |
| 3253 | uint positiveCount = 0; |
| 3254 | for (auto iter = endIter; iter != beginIter;) { |
| 3255 | --iter; |
| 3256 | changedEntries.add(kj::atomicAddRef(*iter->entry)); |
| 3257 | if (iter->entry->getValueStatus() == EntryValueStatus::PRESENT) { |
| 3258 | ++positiveCount; |
| 3259 | } |
| 3260 | if (positiveCount == limit.orDefault(kj::maxValue)) break; |
| 3261 | } |
| 3262 | |
| 3263 | // Increase limit to make sure it can't be underrun by negative entries negating it. |
| 3264 | limit = limit.map([&](uint n) { return n + (changedEntries.size() - positiveCount); }); |
| 3265 | |
| 3266 | return merge(kj::mv(changedEntries), |
| 3267 | cache.listReverse(kj::mv(begin), kj::mv(end), limit, options), GetResultList::REVERSE); |
| 3268 | } |
| 3269 | |
| 3270 | kj::OneOf<ActorCache::GetResultList, kj::Promise<ActorCache::GetResultList>> ActorCache:: |
| 3271 | Transaction::merge(kj::Vector<kj::Own<Entry>> changedEntries, |
| 3272 | kj::OneOf<GetResultList, kj::Promise<GetResultList>> cacheRead, |
| 3273 | GetResultList::Order order) { |
| 3274 | KJ_SWITCH_ONEOF(cacheRead) { |
| 3275 | KJ_CASE_ONEOF(results, GetResultList) { |
| 3276 | return GetResultList(kj::mv(changedEntries), kj::mv(results.entries), order); |
| 3277 | } |
| 3278 | KJ_CASE_ONEOF(promise, kj::Promise<GetResultList>) { |
| 3279 | return promise.then( |
| 3280 | [changedEntries = kj::mv(changedEntries), order](GetResultList results) mutable { |
| 3281 | return GetResultList(kj::mv(changedEntries), kj::mv(results.entries), order); |
| 3282 | }); |
| 3283 | } |
| 3284 | } |
| 3285 | KJ_UNREACHABLE; |
| 3286 | } |
| 3287 | |
| 3288 | // ----------------------------------------------------------------------------- |
| 3289 | // transaction writes |
| 3290 | |
| 3291 | kj::Maybe<kj::Promise<void>> ActorCache::Transaction::put( |
| 3292 | Key key, Value value, WriteOptions options, SpanParent traceSpan) { |
| 3293 | options.noCache = options.noCache || cache.lru.options.noCache; |
| 3294 | auto lock = cache.lru.cleanList.lockExclusive(); |
| 3295 | auto entry = kj::atomicRefcounted<Entry>(cache, kj::mv(key), kj::mv(value)); |
| 3296 | putImpl(lock, kj::mv(entry), options); |
| 3297 | |
| 3298 | // Capture span for use at commit time |
| 3299 | commitSpan = kj::mv(traceSpan); |
| 3300 | |
| 3301 | // Don't apply backpressure because transactions can't be flushed anyway. |
| 3302 | return kj::none; |
| 3303 | } |
| 3304 | |
| 3305 | kj::Maybe<kj::Promise<void>> ActorCache::Transaction::put( |
| 3306 | kj::Array<KeyValuePair> pairs, WriteOptions options, SpanParent traceSpan) { |
| 3307 | options.noCache = options.noCache || cache.lru.options.noCache; |
| 3308 | auto lock = cache.lru.cleanList.lockExclusive(); |
| 3309 | |
| 3310 | for (auto& pair: pairs) { |
| 3311 | auto entry = kj::atomicRefcounted<Entry>(cache, kj::mv(pair.key), kj::mv(pair.value)); |
| 3312 | putImpl(lock, kj::mv(entry), options); |
| 3313 | } |
| 3314 | |
| 3315 | // Capture span for use at commit time |
| 3316 | commitSpan = kj::mv(traceSpan); |
| 3317 | |
| 3318 | // Don't apply backpressure because transactions can't be flushed anyway. |
| 3319 | return kj::none; |
| 3320 | } |
| 3321 | |
| 3322 | kj::Maybe<kj::Promise<void>> ActorCache::Transaction::setAlarm( |
| 3323 | kj::Maybe<kj::Date> newTime, WriteOptions options, SpanParent traceSpan) { |
| 3324 | options.noCache = options.noCache || cache.lru.options.noCache; |
| 3325 | alarmChange = DirtyAlarmWithOptions{DirtyAlarm{newTime}, options}; |
| 3326 | |
| 3327 | // Capture span for use at commit time |
| 3328 | commitSpan = kj::mv(traceSpan); |
| 3329 | |
| 3330 | return kj::none; |
| 3331 | } |
| 3332 | |
| 3333 | kj::OneOf<bool, kj::Promise<bool>> ActorCache::Transaction::delete_( |
| 3334 | Key key, WriteOptions options, SpanParent traceSpan) { |
| 3335 | options.noCache = options.noCache || cache.lru.options.noCache; |
| 3336 | |
| 3337 | uint count = 0; |
| 3338 | kj::Maybe<KeyPtr> keyToCount; |
| 3339 | |
| 3340 | { |
| 3341 | auto lock = cache.lru.cleanList.lockExclusive(); |
| 3342 | auto entry = kj::atomicRefcounted<Entry>(cache, kj::mv(key), EntryValueStatus::ABSENT); |
| 3343 | keyToCount = putImpl(lock, kj::mv(entry), options, count); |
| 3344 | } |
| 3345 | |
| 3346 | // Capture span for use at commit time |
| 3347 | commitSpan = kj::mv(traceSpan); |
| 3348 | |
| 3349 | KJ_IF_SOME(k, keyToCount) { |
| 3350 | // Unfortunately, to find out the count, we have to do a read. |
| 3351 | KJ_SWITCH_ONEOF(cache.get(cloneKey(k), {})) { |
| 3352 | KJ_CASE_ONEOF(value, kj::Maybe<ActorCache::Value>) { |
| 3353 | return value != kj::none; |
| 3354 | } |
| 3355 | KJ_CASE_ONEOF(promise, kj::Promise<kj::Maybe<ActorCache::Value>>) { |
| 3356 | return promise.then([](kj::Maybe<ActorCache::Value> value) { return value != kj::none; }); |
| 3357 | } |
| 3358 | } |
| 3359 | KJ_UNREACHABLE; |
| 3360 | } else { |
| 3361 | return count > 0; |
| 3362 | } |
| 3363 | } |
| 3364 | |
| 3365 | kj::OneOf<uint, kj::Promise<uint>> ActorCache::Transaction::delete_( |
| 3366 | kj::Array<Key> keys, WriteOptions options, SpanParent traceSpan) { |
| 3367 | options.noCache = options.noCache || cache.lru.options.noCache; |
| 3368 | |
| 3369 | if (keys.size() == 0) { |
| 3370 | return 0u; |
| 3371 | } |
| 3372 | |
| 3373 | uint count = 0; |
| 3374 | kj::Vector<kj::Vector<Key>> keysToCount; |
| 3375 | auto startNewBatch = [&]() { return &keysToCount.add(); }; |
| 3376 | auto currentBatch = startNewBatch(); |
| 3377 | |
| 3378 | { |
| 3379 | auto lock = cache.lru.cleanList.lockExclusive(); |
| 3380 | for (auto& key: keys) { |
| 3381 | auto entry = kj::atomicRefcounted<Entry>(cache, kj::mv(key), EntryValueStatus::ABSENT); |
| 3382 | KJ_IF_SOME(keyToCount, putImpl(lock, kj::mv(entry), options, count)) { |
| 3383 | if (currentBatch->size() >= cache.lru.options.maxKeysPerRpc) { |
| 3384 | currentBatch = startNewBatch(); |
| 3385 | } |
| 3386 | currentBatch->add(cloneKey(keyToCount)); |
| 3387 | } |
| 3388 | } |
| 3389 | } |
| 3390 | |
| 3391 | // Capture span for use at commit time |
| 3392 | commitSpan = kj::mv(traceSpan); |
| 3393 | |
| 3394 | if (keysToCount.empty()) { |
| 3395 | return count; |
| 3396 | } else { |
| 3397 | // HACK: Since we allow deletes of larger than our maxKeysPerRpc but these deletes can provoke |
| 3398 | // gets, we need to batch said gets. This all would be much simpler if our default get behavior |
| 3399 | // did batching/sync. |
| 3400 | kj::Maybe<kj::Promise<uint>> maybeTotalPromise; |
| 3401 | for (auto& batch: keysToCount) { |
| 3402 | // Unfortunately, to find out the count, we have to do a read. Note that even returning this |
| 3403 | // value separate from a committed transaction means that non-transaction storage ops can make |
| 3404 | // the value incorrect. |
| 3405 | KJ_SWITCH_ONEOF(cache.get(batch.releaseAsArray(), {})) { |
| 3406 | KJ_CASE_ONEOF(results, GetResultList) { |
| 3407 | count = count + results.size(); |
| 3408 | } |
| 3409 | KJ_CASE_ONEOF(promise, kj::Promise<GetResultList>) { |
| 3410 | if (maybeTotalPromise == kj::none) { |
| 3411 | // We had to do a remote get, start a promise |
| 3412 | maybeTotalPromise.emplace(0); |
| 3413 | } |
| 3414 | maybeTotalPromise = KJ_ASSERT_NONNULL(maybeTotalPromise) |
| 3415 | .then([promise = kj::mv(promise)]( |
| 3416 | uint previousResult) mutable -> kj::Promise<uint> { |
| 3417 | return promise.then([previousResult](GetResultList results) mutable -> uint { |
| 3418 | return previousResult + kj::implicitCast<uint>(results.size()); |
| 3419 | }); |
| 3420 | }); |
| 3421 | } |
| 3422 | } |
| 3423 | } |
| 3424 | |
| 3425 | KJ_IF_SOME(totalPromise, maybeTotalPromise) { |
| 3426 | return totalPromise.then([count](uint result) { return count + result; }); |
| 3427 | } else { |
| 3428 | return count; |
| 3429 | } |
| 3430 | } |
| 3431 | } |
| 3432 | |
| 3433 | kj::Maybe<ActorCache::KeyPtr> ActorCache::Transaction::putImpl( |
| 3434 | Lock& lock, kj::Own<Entry> entry, const WriteOptions& options, kj::Maybe<uint&> count) { |
| 3435 | Change change{ |
| 3436 | .entry = kj::mv(entry), |
| 3437 | .options = options, |
| 3438 | }; |
| 3439 | bool replaced = false; |
| 3440 | auto& slot = entriesToWrite.upsert(kj::mv(change), [&](auto& existing, auto&& replacement) { |
| 3441 | replaced = true; |
| 3442 | KJ_IF_SOME(c, count) { |
| 3443 | c += existing.entry->getValueStatus() == EntryValueStatus::PRESENT; |
| 3444 | } |
| 3445 | existing = kj::mv(replacement); |
| 3446 | }); |
| 3447 | if (replaced) { |
| 3448 | // Already counted. |
| 3449 | return kj::none; |
| 3450 | } else { |
| 3451 | return KeyPtr(slot.entry->key); |
| 3452 | } |
| 3453 | } |
| 3454 | } // namespace workerd |