Skip to content
File

Blob: src/workerd/api/memory-cache.c++

27.0 KB
1#include "memory-cache.h"
2 
3#include <workerd/api/util.h>
4#include <workerd/io/io-context.h>
5#include <workerd/io/io-util.h>
6#include <workerd/jsg/jsg.h>
7#include <workerd/jsg/ser.h>
8#include <workerd/util/weak-refs.h>
9 
10namespace workerd::api {
11 
12static constexpr size_t MAX_KEY_SIZE = 2 * 1024;
13 
14// Returns the current calendar time as a double, just like Date.now() would,
15// except without the safeguards that exist within an I/O context. This
16// function is used only when a worker is being created or destroyed.
17static double getCurrentTimeOutsideIoContext() {
18 KJ_ASSERT(!IoContext::hasCurrent());
19 auto now = kj::systemCoarseCalendarClock().now();
20 return (now - kj::UNIX_EPOCH) / kj::MILLISECONDS;
21}
22 
23// Returns true if the given expiration time exists and has passed. If this is
24// called in an I/O context, the I/O context's timer is used. Otherwise,
25// if allowOutsideIoContext is true, the system clock is used (see above).
26// Lastly, if this function is called from outside of an I/O context and if
27// allowOutsideIoContext is false, this function returns false regardless
28// of whether the expiration time has passed.
29static bool hasExpired(const kj::Maybe<double>& expiration, bool allowOutsideIoContext = false) {
30 KJ_IF_SOME(e, expiration) {
31 double now = (allowOutsideIoContext && !IoContext::hasCurrent())
32 ? getCurrentTimeOutsideIoContext()
33 : dateNow();
34 return e < now;
35 }
36 return false;
37}
38 
39SharedMemoryCache::SharedMemoryCache(kj::Maybe<const MemoryCacheProvider&> provider,
40 kj::StringPtr id,
41 kj::Maybe<AdditionalResizeMemoryLimitHandler&> additionalResizeMemoryLimitHandler,
42 const kj::MonotonicClock& timer)
43 : provider(provider),
44 id(kj::str(id)),
45 additionalResizeMemoryLimitHandler(additionalResizeMemoryLimitHandler),
46 timer(timer) {}
47 
48SharedMemoryCache::~SharedMemoryCache() noexcept(false) {
49 KJ_IF_SOME(p, provider) {
50 p.removeInstance(*this);
51 }
52}
53 
54void SharedMemoryCache::suggest(const Limits& limits) const {
55 auto data = this->data.lockExclusive();
56 bool isKnownLimit = data->suggestedLimits.contains(limits);
57 data->suggestedLimits.insert(limits);
58 if (!isKnownLimit) {
59 resize(*data);
60 }
61}
62 
63void SharedMemoryCache::unsuggest(const Limits& limits) const {
64 auto data = this->data.lockExclusive();
65 auto loc = data->suggestedLimits.find(limits);
66 KJ_ASSERT(loc != data->suggestedLimits.end());
67 data->suggestedLimits.erase(loc);
68 resize(*data);
69}
70 
71void SharedMemoryCache::resize(ThreadUnsafeData& data) const {
72 data.effectiveLimits = Limits::min();
73 for (const auto& limits: data.suggestedLimits) {
74 data.effectiveLimits = Limits::max(data.effectiveLimits, limits.normalize());
75 }
76 
77 KJ_IF_SOME(handler, additionalResizeMemoryLimitHandler) {
78 // Allow the embedder to adjust the effective limits.
79 handler(data);
80 }
81 
82 // Fast path for clearing the cache.
83 if (data.effectiveLimits.maxKeys == 0) {
84 data.totalValueSize = 0;
85 data.cache.clear();
86 return;
87 }
88 
89 // First, remove any values that might be too large.
90 while (data.cache.size() != 0) {
91 MemoryCacheEntry& largestEntry = *data.cache.ordered<2>().begin();
92 if (largestEntry.size() <= data.effectiveLimits.maxValueSize) {
93 break;
94 }
95 data.totalValueSize -= largestEntry.size();
96 data.cache.erase(largestEntry);
97 }
98 
99 // Now just keep keep evicting until we are within limits.
100 while (data.totalValueSize > data.effectiveLimits.maxTotalValueSize ||
101 data.cache.size() > data.effectiveLimits.maxKeys) {
102 evictNextWhileLocked(data, true);
103 }
104}
105 
106kj::Maybe<kj::Own<CacheValue>> SharedMemoryCache::getWhileLocked(
107 ThreadUnsafeData& data, const kj::String& key) const {
108 KJ_IF_SOME(existingCacheEntry, data.cache.find(key)) {
109 if (hasExpired(existingCacheEntry.expiration)) {
110 // The cache entry has an associated expiration time and that time has
111 // passed (according to the calling IoContext's timer).
112 data.totalValueSize -= existingCacheEntry.size();
113 data.cache.erase(existingCacheEntry);
114 return kj::none;
115 }
116 
117 // Obtain a reference to the cache value before we kj::mv the cache entry.
118 auto cacheValue = kj::atomicAddRef(*existingCacheEntry.value);
119 
120 // Update the liveliness.
121 MemoryCacheEntry entry = data.cache.release(existingCacheEntry);
122 entry.liveliness = data.stepLiveliness();
123 data.cache.insert(kj::mv(entry));
124 
125 return kj::mv(cacheValue);
126 } else {
127 return kj::none;
128 }
129}
130 
131void SharedMemoryCache::putWhileLocked(ThreadUnsafeData& data,
132 const kj::String& key,
133 kj::Own<CacheValue>&& value,
134 kj::Maybe<double> expiration) const {
135 size_t valueSize = value->bytes.size();
136 
137 auto writeSpan = IoContext::current().makeTraceSpan("memory_cache_write"_kjc);
138 writeSpan.setTag("key"_kjc, key.asPtr());
139 writeSpan.setTag("value_size"_kjc, static_cast<double>(valueSize));
140 writeSpan.setTag("has_expiration"_kjc, expiration != kj::none);
141 
142 if (valueSize > data.effectiveLimits.maxValueSize) {
143 // Silently drop the value. For consistency, also drop the previous value,
144 // if one exists, such that a subsequent read() will not return an outdated
145 // value. Note that removeIfExistsWhileLocked(key) will update the
146 // totalValueSize if necessary, so we don't need to do that here.
147 writeSpan.setTag("write_rejected"_kjc, true);
148 writeSpan.setTag("rejection_reason"_kjc, "value_too_large"_kjc);
149 writeSpan.setTag("max_value_size"_kjc, static_cast<double>(data.effectiveLimits.maxValueSize));
150 removeIfExistsWhileLocked(data, key);
151 return;
152 }
153 
154 if (hasExpired(expiration)) {
155 writeSpan.setTag("write_rejected"_kjc, true);
156 writeSpan.setTag("rejection_reason"_kjc, "already_expired"_kjc);
157 removeIfExistsWhileLocked(data, key);
158 return;
159 }
160 
161 kj::Maybe<MemoryCacheEntry&> existingEntry = data.cache.find(key.asPtr());
162 bool isUpdate = existingEntry != kj::none;
163 size_t evictionCount = 0;
164 
165 KJ_IF_SOME(entry, existingEntry) {
166 size_t oldValueSize = entry.size();
167 KJ_ASSERT(data.totalValueSize >= oldValueSize);
168 MemoryCacheEntry updatedEntry = data.cache.release(entry);
169 data.totalValueSize -= oldValueSize;
170 while (data.totalValueSize + valueSize > data.effectiveLimits.maxTotalValueSize) {
171 // We have already released the existing entry for our key, so there is no
172 // risk of evicting it.
173 evictNextWhileLocked(data);
174 evictionCount++;
175 }
176 updatedEntry.liveliness = data.stepLiveliness();
177 updatedEntry.value = kj::mv(value);
178 updatedEntry.expiration = expiration;
179 data.cache.insert(kj::mv(updatedEntry));
180 data.totalValueSize += valueSize;
181 } else {
182 // Ensure that adding a new key won't push us over the limit.
183 if (data.cache.size() >= data.effectiveLimits.maxKeys) {
184 evictNextWhileLocked(data);
185 evictionCount++;
186 }
187 // Ensure that the size of the new value won't push us over the limit.
188 while (data.totalValueSize + valueSize > data.effectiveLimits.maxTotalValueSize) {
189 evictNextWhileLocked(data);
190 evictionCount++;
191 }
192 MemoryCacheEntry newEntry = {
193 kj::str(key),
194 data.stepLiveliness(),
195 kj::mv(value),
196 expiration,
197 };
198 data.cache.insert(kj::mv(newEntry));
199 data.totalValueSize += valueSize;
200 }
201 
202 writeSpan.setTag("write_success"_kjc, true);
203 writeSpan.setTag("is_update"_kjc, isUpdate);
204 writeSpan.setTag("evictions_triggered"_kjc, static_cast<double>(evictionCount));
205 writeSpan.setTag("cache_total_size_after"_kjc, static_cast<double>(data.totalValueSize));
206 writeSpan.setTag("cache_entry_count_after"_kjc, static_cast<double>(data.cache.size()));
207}
208 
209void SharedMemoryCache::evictNextWhileLocked(
210 ThreadUnsafeData& data, bool allowOutsideIoContext) const {
211 // The caller is responsible for ensuring that the cache is not empty already.
212 KJ_REQUIRE(data.cache.size() > 0);
213 
214 // Create eviction span - only called from IO context
215 auto evictionSpan = IoContext::current().makeTraceSpan("memory_cache_eviction"_kjc);
216 
217 // If there is an entry that has expired already, evict that one.
218 MemoryCacheEntry& maybeExpired = *data.cache.ordered<3>().begin();
219 KJ_ASSERT(data.totalValueSize >= maybeExpired.size());
220 if (hasExpired(maybeExpired.expiration, allowOutsideIoContext)) {
221 evictionSpan.setTag("eviction_reason"_kjc, "expiration"_kjc);
222 evictionSpan.setTag("evicted_key"_kjc, maybeExpired.key.asPtr());
223 evictionSpan.setTag("evicted_size"_kjc, static_cast<double>(maybeExpired.size()));
224 evictionSpan.setTag("cache_size_before"_kjc, static_cast<double>(data.totalValueSize));
225 evictionSpan.setTag("cache_entries_before"_kjc, static_cast<double>(data.cache.size()));
226 data.totalValueSize -= maybeExpired.size();
227 data.cache.erase(maybeExpired);
228 return;
229 }
230 
231 // Otherwise, if no entry has expired, evict the least recently used entry.
232 MemoryCacheEntry& leastRecentlyUsed = *data.cache.ordered<1>().begin();
233 evictionSpan.setTag("eviction_reason"_kjc, "lru"_kjc);
234 evictionSpan.setTag("evicted_key"_kjc, leastRecentlyUsed.key.asPtr());
235 evictionSpan.setTag("evicted_size"_kjc, static_cast<double>(leastRecentlyUsed.size()));
236 evictionSpan.setTag("cache_size_before"_kjc, static_cast<double>(data.totalValueSize));
237 evictionSpan.setTag("cache_entries_before"_kjc, static_cast<double>(data.cache.size()));
238 KJ_ASSERT(data.totalValueSize >= leastRecentlyUsed.size());
239 data.totalValueSize -= leastRecentlyUsed.size();
240 data.cache.erase(leastRecentlyUsed);
241}
242 
243void SharedMemoryCache::removeIfExistsWhileLocked(
244 ThreadUnsafeData& data, const kj::String& key) const {
245 KJ_IF_SOME(entry, data.cache.find(key)) {
246 // This DOES NOT count as an eviction because it might happen while
247 // replacing the existing cache entry with a new one, when the new one is
248 // being evicted immediately. It is up to the caller to count that.
249 size_t valueSize = entry.size();
250 KJ_ASSERT(valueSize <= data.totalValueSize);
251 data.totalValueSize -= valueSize;
252 data.cache.erase(entry);
253 }
254}
255 
256kj::Own<const SharedMemoryCache> SharedMemoryCache::create(
257 kj::Maybe<const MemoryCacheProvider&> provider,
258 kj::StringPtr id,
259 kj::Maybe<AdditionalResizeMemoryLimitHandler&> handler,
260 const kj::MonotonicClock& timer) {
261 return kj::atomicRefcounted<const SharedMemoryCache>(provider, id, handler, timer);
262}
263 
264SharedMemoryCache::Use::Use(kj::Own<const SharedMemoryCache> cache, const Limits& limits)
265 : cache(kj::mv(cache)),
266 limits(limits) {
267 this->cache->suggest(limits);
268}
269 
270SharedMemoryCache::Use::Use(Use&& other): cache(kj::mv(other.cache)), limits(other.limits) {
271 this->cache->suggest(limits);
272}
273 
274SharedMemoryCache::Use::~Use() noexcept(false) {
275 if (cache.get() != nullptr) {
276 cache->unsuggest(limits);
277 }
278}
279 
280kj::Maybe<kj::Own<CacheValue>> SharedMemoryCache::Use::getWithoutFallback(
281 const kj::String& key, SpanBuilder& readSpan) const {
282 kj::Locked<ThreadUnsafeData> data = [&] {
283 auto memoryCacheLockRecord =
284 ScopedDurationTagger(readSpan, memoryCachekLockWaitTimeTag, cache->timer);
285 return cache->data.lockExclusive();
286 }();
287 auto result = cache->getWhileLocked(*data, key);
288 
289 // Track cache hit/miss
290 readSpan.setTag("cache_hit"_kjc, result != kj::none);
291 KJ_IF_SOME(value, result) {
292 readSpan.setTag("entry_size"_kjc, static_cast<double>(value->bytes.size()));
293 }
294 readSpan.setTag("cache_total_size"_kjc, static_cast<double>(data->totalValueSize));
295 readSpan.setTag("cache_entry_count"_kjc, static_cast<double>(data->cache.size()));
296 
297 return result;
298}
299 
300kj::OneOf<kj::Own<CacheValue>, kj::Promise<SharedMemoryCache::Use::GetWithFallbackOutcome>>
301SharedMemoryCache::Use::getWithFallback(const kj::String& key, SpanBuilder& readSpan) const {
302 kj::Locked<ThreadUnsafeData> data = [&] {
303 auto memoryCacheLockRecord =
304 ScopedDurationTagger(readSpan, memoryCachekLockWaitTimeTag, cache->timer);
305 return cache->data.lockExclusive();
306 }();
307 KJ_IF_SOME(existingValue, cache->getWhileLocked(*data, key)) {
308 // Cache hit
309 readSpan.setTag("cache_hit"_kjc, true);
310 readSpan.setTag("entry_size"_kjc, static_cast<double>(existingValue->bytes.size()));
311 readSpan.setTag("cache_total_size"_kjc, static_cast<double>(data->totalValueSize));
312 readSpan.setTag("cache_entry_count"_kjc, static_cast<double>(data->cache.size()));
313 return kj::mv(existingValue);
314 } else KJ_IF_SOME(existingInProgress, data->inProgress.find(key)) {
315 // Cache miss - but another request is already fetching this key
316 readSpan.setTag("cache_hit"_kjc, false);
317 readSpan.setTag("coalesced_request"_kjc, true);
318 readSpan.setTag("waiting_on_inflight"_kjc, true);
319 readSpan.setTag(
320 "inflight_waiters_count"_kjc, static_cast<double>(existingInProgress->waiting.size() + 1));
321 
322 // Create a span to track how long we wait for the inflight request
323 auto waitSpan = readSpan.newChild("memory_cache_coalesce_wait"_kjc);
324 waitSpan.setTag("key"_kjc, key.asPtr());
325 waitSpan.setTag("waiters_ahead"_kjc, static_cast<double>(existingInProgress->waiting.size()));
326 
327 // We return a Promise, but we keep the fulfiller. We might fulfill it
328 // from a different thread, so we need a cross-thread fulfiller here.
329 auto pair = kj::newPromiseAndCrossThreadFulfiller<GetWithFallbackOutcome>();
330 existingInProgress->waiting.emplace(kj::mv(pair.fulfiller));
331 // We have to register a pending event with the I/O context so that the
332 // runtime does not detect a hanging promise. Another fallback is in
333 // progress and once it settles, we will fulfill the promise that we return
334 // here, either with the produced value or with another fallback task.
335 return pair.promise.attach(IoContext::current().registerPendingEvent(), kj::mv(waitSpan));
336 } else {
337 // Cache miss - this request will fetch from upstream
338 readSpan.setTag("cache_hit"_kjc, false);
339 readSpan.setTag("coalesced_request"_kjc, false);
340 readSpan.setTag("initiating_fallback"_kjc, true);
341 readSpan.setTag("cache_total_size"_kjc, static_cast<double>(data->totalValueSize));
342 readSpan.setTag("cache_entry_count"_kjc, static_cast<double>(data->cache.size()));
343 
344 auto& newEntry = data->inProgress.insert(kj::heap<InProgress>(kj::str(key)));
345 auto inProgress = newEntry.get();
346 return kj::Promise<GetWithFallbackOutcome>(prepareFallback(*inProgress));
347 }
348}
349 
350SharedMemoryCache::Use::FallbackDoneCallback SharedMemoryCache::Use::prepareFallback(
351 InProgress& inProgress) const {
352 // We need to detect if the Promise that we are about to create ever settles,
353 // as opposed to being destroyed without either being resolved or rejecting.
354 struct FallbackStatus {
355 bool hasSettled = false;
356 };
357 auto status = kj::heap<FallbackStatus>();
358 auto& statusRef = *status;
359 
360 auto deferredCancel = kj::defer([this, status = kj::mv(status), &inProgress]() {
361 // If the callback was destroyed without having run (for example, because
362 // it was added to an I/O context that has since been canceled), we treat
363 // it as if the promise had failed.
364 if (!status->hasSettled) {
365 handleFallbackFailure(inProgress);
366 }
367 });
368 
369 return [this, &inProgress, &status = statusRef, deferredCancel = kj::mv(deferredCancel)](
370 kj::Maybe<FallbackResult> maybeResult, SpanBuilder& fallbackSpan) mutable {
371 KJ_IF_SOME(result, maybeResult) {
372 // The fallback succeeded. Store the value in the cache and propagate it to
373 // all waiting requests, even if it has expired already.
374 status.hasSettled = true;
375 
376 auto data = cache->data.lockExclusive();
377 size_t waiterCount = inProgress.waiting.size();
378 
379 cache->putWhileLocked(
380 *data, kj::str(inProgress.key), kj::atomicAddRef(*result.value), result.expiration);
381 
382 inProgress.waiting.drainTo(
383 [&](auto&& waiter) { waiter.fulfiller->fulfill(kj::atomicAddRef(*result.value)); });
384 data->inProgress.eraseMatch(inProgress.key);
385 
386 // Track the completion of fallback and distribution to waiters
387 fallbackSpan.setTag("waiters_notified"_kjc, static_cast<double>(waiterCount));
388 } else {
389 // The fallback failed for some reason. We do not care much about why it
390 // failed. If there are other queued fallbacks, handelFallbackFailure will
391 // schedule the next one.
392 status.hasSettled = true;
393 handleFallbackFailure(inProgress);
394 }
395 };
396}
397 
398void SharedMemoryCache::Use::handleFallbackFailure(InProgress& inProgress) const {
399 kj::Own<kj::CrossThreadPromiseFulfiller<GetWithFallbackOutcome>> nextFulfiller;
400 
401 // If there is another queued fallback, retrieve it and remove it from the
402 // queue. Otherwise, just delete the queue entirely.
403 {
404 auto data = cache->data.lockExclusive();
405 
406 KJ_IF_SOME(next, inProgress.waiting.pop()) {
407 nextFulfiller = kj::mv(next.fulfiller);
408 } else {
409 data->inProgress.eraseMatch(inProgress.key);
410 }
411 }
412 
413 // fulfill() might destroy the Promise returned by prepareFallback(). In
414 // particular, that will happen if the I/O context that the fulfiller was
415 // created for has been canceled or destroyed, in which case the promise
416 // associated with the fulfiller has been destroyed. When the promise returned
417 // by prepareFallback() is destroyed without having settled, it will recover
418 // from that, but it will lock the cache while doing so. That is why it is
419 // important that the cache is not already locked when we call fulfill().
420 if (nextFulfiller) {
421 nextFulfiller->fulfill(prepareFallback(inProgress));
422 }
423}
424 
425void SharedMemoryCache::Use::delete_(const kj::String& key) const {
426 auto data = cache->data.lockExclusive();
427 cache->removeIfExistsWhileLocked(*data, key);
428}
429 
430// Attempts to serialize a JavaScript value. If that fails, this function throws
431// a tunneled exception, see jsg::createTunneledException().
432static kj::Own<CacheValue> hackySerialize(jsg::Lock& js, jsg::JsRef<jsg::JsValue>& value) {
433 JSG_TRY(js) {
434 jsg::Serializer serializer(js);
435 serializer.write(js, value.getHandle(js));
436 return kj::atomicRefcounted<CacheValue>(serializer.release().data);
437 }
438 JSG_CATCH(exception) {
439 // We run into big problems with tunneled exceptions here. When
440 // the toString() function of the JavaScript error is not marked
441 // as side effect free, tunneling the exception fails entirely
442 // because kj::str() returns an empty string for the error. As a
443 // workaround, we drop the error object in that case and return
444 // a generic error that only includes the type of the value.
445 // TODO(later): remove this workaround
446 if (kj::str(exception.getHandle(js)).size() == 0) {
447 throw JSG_KJ_EXCEPTION(
448 FAILED, DOMDataCloneError, "failed to serialize ", value.getHandle(js).typeOf(js));
449 }
450 
451 // This is still pretty bad. We lose the original error stack.
452 // TODO(later): remove string-based error tunneling
453 throw js.exceptionToKj(kj::mv(exception));
454 }
455}
456 
457jsg::Promise<jsg::JsRef<jsg::JsValue>> MemoryCache::read(jsg::Lock& js,
458 jsg::NonCoercible<kj::String> key,
459 jsg::Optional<FallbackFunction> optionalFallback) {
460 if (key.value.size() > MAX_KEY_SIZE) {
461 return js.rejectedPromise<jsg::JsRef<jsg::JsValue>>(js.rangeError("Key too large."_kj));
462 }
463 
464 auto readSpan = IoContext::current().makeTraceSpan("memory_cache_read"_kjc);
465 auto userReadSpan = IoContext::current().makeUserTraceSpan("memory_cache_read"_kjc);
466 
467 KJ_IF_SOME(fallback, optionalFallback) {
468 KJ_SWITCH_ONEOF(cacheUse.getWithFallback(key.value, readSpan)) {
469 KJ_CASE_ONEOF(result, kj::Own<CacheValue>) {
470 // Optimization: Don't even release the isolate lock if the value is already in cache.
471 jsg::Deserializer deserializer(js, result->bytes.asPtr());
472 auto value = jsg::JsRef(js, deserializer.readValue(js));
473 
474 return js.resolvedPromise(kj::mv(value));
475 }
476 KJ_CASE_ONEOF(promise, kj::Promise<SharedMemoryCache::Use::GetWithFallbackOutcome>) {
477 return IoContext::current().awaitIo(js, kj::mv(promise),
478 [fallback = kj::mv(fallback), key = kj::str(key.value), readSpan = kj::mv(readSpan),
479 userSpan = kj::mv(userReadSpan), self = JSG_THIS](
480 jsg::Lock& js, SharedMemoryCache::Use::GetWithFallbackOutcome cacheResult) mutable
481 -> jsg::Promise<jsg::JsRef<jsg::JsValue>> {
482 KJ_SWITCH_ONEOF(cacheResult) {
483 KJ_CASE_ONEOF(serialized, kj::Own<CacheValue>) {
484 readSpan.setTag("fallback_cache_hit"_kjc, true);
485 readSpan.setTag("entry_size"_kjc, static_cast<double>(serialized->bytes.size()));
486 
487 jsg::Deserializer deserializer(js, serialized->bytes.asPtr());
488 return js.resolvedPromise(jsg::JsRef(js, deserializer.readValue(js)));
489 }
490 KJ_CASE_ONEOF(callback, SharedMemoryCache::Use::FallbackDoneCallback) {
491 auto& context = IoContext::current();
492 auto heapCallback = kj::heap(kj::mv(callback));
493 
494 // Create a span for the fallback execution
495 auto fallbackSpan = readSpan.newChild("memory_cache_fallback"_kjc);
496 fallbackSpan.setTag("key"_kjc, key.asPtr());
497 
498 // Wrap the spans in RefcountedWrapper so they can be shared between then/catch
499 auto fallbackSpanRc = kj::refcountedWrapper<SpanBuilder>(kj::mv(fallbackSpan));
500 auto readSpanRc = kj::refcountedWrapper<SpanBuilder>(kj::mv(readSpan));
501 
502 return js.evalNow([&]() { return fallback(js, kj::mv(key)); })
503 .then(js,
504 [callback = context.addObject(*heapCallback),
505 fallbackSpan = fallbackSpanRc->addWrappedRef(),
506 readSpan = readSpanRc->addWrappedRef()](jsg::Lock& js,
507 CacheValueProduceResult result) mutable -> jsg::JsRef<jsg::JsValue> {
508 // NOTE: `callback` is IoPtr, not IoOwn. The catch block gets the IoOwn, which
509 // ensures the object still exists at this point.
510 fallbackSpan->setTag("fallback_success"_kjc, true);
511 
512 auto serialized = hackySerialize(js, result.value);
513 fallbackSpan->setTag(
514 "fallback_result_size"_kjc, static_cast<double>(serialized->bytes.size()));
515 
516 KJ_IF_SOME(expiration, result.expiration) {
517 JSG_REQUIRE(
518 !kj::isNaN(expiration), TypeError, "Expiration time must not be NaN.");
519 fallbackSpan->setTag("has_expiration"_kjc, true);
520 } else {
521 fallbackSpan->setTag("has_expiration"_kjc, false);
522 }
523 (*callback)(
524 SharedMemoryCache::Use::FallbackResult{kj::mv(serialized), result.expiration},
525 *fallbackSpan);
526 return kj::mv(result.value);
527 })
528 .catch_(js,
529 JSG_VISITABLE_LAMBDA(
530 (self = kj::mv(self), callback = context.addObject(kj::mv(heapCallback)),
531 fallbackSpan = fallbackSpanRc->addWrappedRef(),
532 readSpan = readSpanRc->addWrappedRef()),
533 (self),
534 (jsg::Lock & js, jsg::Value&& exception) mutable->jsg::JsRef<jsg::JsValue> {
535 fallbackSpan->setTag("fallback_success"_kjc, false);
536 fallbackSpan->setTag(
537 "fallback_error"_kjc, kj::str(exception.getHandle(js)));
538 (*callback)(kj::none, *fallbackSpan);
539 js.throwException(kj::mv(exception));
540 }));
541 }
542 }
543 KJ_UNREACHABLE;
544 });
545 }
546 }
547 KJ_UNREACHABLE;
548 } else {
549 KJ_IF_SOME(cacheValue, cacheUse.getWithoutFallback(key.value, readSpan)) {
550 jsg::Deserializer deserializer(js, cacheValue->bytes.asPtr());
551 return js.resolvedPromise(jsg::JsRef(js, deserializer.readValue(js)));
552 }
553 return js.resolvedPromise(jsg::JsRef(js, js.undefined()));
554 }
555}
556 
557void MemoryCache::delete_(jsg::Lock& js, jsg::NonCoercible<kj::String> key) {
558 // Ignore operations on keys exceeding key max size.
559 if (key.value.size() > MAX_KEY_SIZE) {
560 js.throwException(js.rangeError("Key too large."_kj));
561 return;
562 }
563 
564 auto deleteSpan = IoContext::current().makeTraceSpan("memory_cache_delete"_kjc);
565 deleteSpan.setTag("key"_kjc, key.value.asPtr());
566 
567 cacheUse.delete_(key.value);
568 
569 deleteSpan.setTag("delete_completed"_kjc, true);
570}
571 
572// ======================================================================================
573 
574MemoryCacheProvider::MemoryCacheProvider(const kj::MonotonicClock& timer,
575 kj::Maybe<SharedMemoryCache::AdditionalResizeMemoryLimitHandler>
576 additionalResizeMemoryLimitHandler)
577 : additionalResizeMemoryLimitHandler(kj::mv(additionalResizeMemoryLimitHandler)),
578 timer(timer) {}
579 
580MemoryCacheProvider::~MemoryCacheProvider() noexcept(false) {
581 // TODO(cleanup): Later, assuming progress is made on kj::Ptr<T>, we ought to be able
582 // to remove this. For now we just need to make sure that the MemoryCacheProvider instance
583 // outlives any SharedMemoryCache instances that are referencing it.
584 KJ_REQUIRE(caches.lockShared()->size() == 0,
585 "There are still active SharedMemoryCache instances. Use-after-free errors are likely.");
586}
587 
588kj::Own<const SharedMemoryCache> MemoryCacheProvider::getInstance(
589 kj::Maybe<kj::StringPtr> cacheId) const {
590 
591 const auto makeCache = [this](kj::Maybe<const MemoryCacheProvider&> provider, kj::StringPtr id) {
592 // The cache doesn't exist in the map. Let's create it.
593 auto handler = additionalResizeMemoryLimitHandler.map(
594 [](const SharedMemoryCache::AdditionalResizeMemoryLimitHandler& handler)
595 -> SharedMemoryCache::AdditionalResizeMemoryLimitHandler& {
596 return const_cast<SharedMemoryCache::AdditionalResizeMemoryLimitHandler&>(handler);
597 });
598 return SharedMemoryCache::create(provider, id, handler, timer);
599 };
600 
601 KJ_IF_SOME(cid, cacheId) {
602 auto lock = caches.lockExclusive();
603 
604 // First, let's see if the cache already exists. If it does, we'll just return
605 // a strong reference to it.
606 KJ_IF_SOME(found, lock->find(cid)) {
607 KJ_IF_SOME(ref, kj::atomicAddRefWeak(*found)) {
608 return kj::mv(ref);
609 } else {
610 // We found an entry in the map, but atomicAddRefWeak failed. Doh. We have
611 // to replace the map entry with a new cache instance.
612 auto cache = makeCache(kj::Maybe<const MemoryCacheProvider&>(*this), cid);
613 lock->upsert(kj::str(cid), cache.get());
614 return kj::mv(cache);
615 }
616 }
617 
618 // The cache doesn't exist, let's create it and add it to the map
619 auto cache = makeCache(kj::Maybe<const MemoryCacheProvider&>(*this), cid);
620 lock->insert(kj::str(cid), cache.get());
621 return kj::mv(cache);
622 }
623 
624 // Since we don't have a cache id, we'll just create a new cache and return it.
625 return makeCache(kj::none, nullptr);
626}
627 
628void MemoryCacheProvider::removeInstance(const SharedMemoryCache& instance) const {
629 // This is fun. We have to make sure that the instance to be removed is actually
630 // what we expect it to be.
631 auto lock = caches.lockExclusive();
632 KJ_IF_SOME(found, lock->findEntry(instance.getId())) {
633 if (found.value != &instance) {
634 // Not the instance we expected it to be. Cache instance was likely replaced
635 // by a new instance with the same id. Do nothing.
636 return;
637 }
638 lock->erase(found);
639 }
640}
641 
642} // namespace workerd::api