Skip to content
File

Blob: src/workerd/api/memory-cache.h

cpp541 lines
1#pragma once
2 
3#include <workerd/io/compatibility-date.capnp.h>
4#include <workerd/jsg/jsg.h>
5#include <workerd/util/checked-queue.h>
6 
7#include <kj/hash.h>
8#include <kj/map.h>
9#include <kj/mutex.h>
10#include <kj/table.h>
11#include <kj/time.h>
12 
13#include <set>
14 
15namespace workerd {
16class SpanBuilder;
17}
18 
19namespace workerd::api {
20 
21// The MemoryCache mechanism is an in-process, memory-resident data cache that
22// can be configured for workers. A single cache instance can be unique to an
23// individual worker or shared across multiple workers / isolates.
24//
25// Instances are configured as bindings on the worker (set up in the workers
26// configuration) and accessible via the environment bindings passed into the
27// worker handler functions:
28//
29// async fetch(req, env) {
30// await env.MY_CACHE.read('key', () => {
31// // Called if the 'key' does not exist in the cache
32// return 'new value';
33// });
34// }
35//
36// The cache is only capable of storing values that are v8 serializable (so
37// JS primitives other than Symbol, ordinary JavaScript objects but not class
38// instances, etc). Objects that represent i/o (like streams or promises are
39// explicitly not supported.
40 
41struct CacheValue: kj::AtomicRefcounted {
42 CacheValue(kj::Array<kj::byte>&& bytes): bytes(kj::mv(bytes)) {}
43 
44 kj::Array<kj::byte> bytes;
45};
46 
47struct MemoryCacheEntry {
48 // The key that this entry is associated with.
49 kj::String key;
50 
51 // Whenever an entry is created, updated, or retrieved, its liveliness is
52 // set to the value of a monotonically increasing counter.
53 uint64_t liveliness;
54 // TODO(cleanup): The liveliness index accomplishes the same thing as
55 // kj::InsertionOrderIndex.
56 //
57 // TODO(perf): Updating a cache entry's liveliness requires a re-insertion,
58 // which means that cache reads require an exclusive lock. This may be
59 // suboptimal for a read-heavy workload. WorkerSet avoids this by atomically
60 // updating a `lastUsed` timestamp. The tradeoff is that LRU-eviction
61 // becomes O(n) instead of O(1), since we can no longer use kj::Table's
62 // index to find the LRU entry.
63 
64 // The stored JavaScript value, serialized by V8. It is atomicRefcounted to
65 // allow threads to deserialize the value without having to lock the cache,
66 // so the value can even be deserialized while the cache entry is being
67 // evicted.
68 kj::Own<CacheValue> value;
69 
70 inline size_t size() const {
71 return value->bytes.size();
72 }
73 
74 // The expiration timestamp of this cache entry, usually the time at which the
75 // entry was created plus some TTL. This is measured in milliseconds and
76 // stored as a double so that it is compatible with api::dateNow() and
77 // EdgeWorkerPlatform::CurrentClockTimeMillis().
78 kj::Maybe<double> expiration;
79};
80 
81struct CacheValueProduceResult {
82 jsg::JsRef<jsg::JsValue> value;
83 jsg::Optional<double> expiration;
84 JSG_STRUCT(value, expiration);
85};
86 
87class MemoryCacheProvider;
88 
89// An in-memory cache that can be accessed by any number of workers/isolates
90// within the same process.
91// TODO(soon): We plan to explore replacing this implementation with a memcached-based
92// implementation in the near future. The memcached-based impl would likely be
93// fairly different from this implementation so quite a few of the details here
94// are expected to change.
95class SharedMemoryCache: public kj::AtomicRefcounted {
96 private:
97 struct InProgress;
98 
99 public:
100 struct ThreadUnsafeData;
101 
102 struct Limits {
103 // The maximum number of keys that may exist within the cache at the same
104 // time. The cache size grows at least linearly in the number of entries.
105 uint32_t maxKeys;
106 
107 // The maximum size of each individual value, when serialized.
108 uint32_t maxValueSize;
109 
110 // The maximum sum of all stored values. This is essentially the cache size,
111 // except that it only includes the sizes of the values and does not account
112 // for keys and the overhead of the data structures themselves.
113 uint64_t maxTotalValueSize;
114 
115 bool operator<(const Limits& b) const {
116 if (maxTotalValueSize != b.maxTotalValueSize) {
117 return maxTotalValueSize < b.maxTotalValueSize;
118 }
119 if (maxKeys != b.maxKeys) {
120 return maxKeys < b.maxKeys;
121 }
122 return maxTotalValueSize < b.maxTotalValueSize;
123 }
124 
125 Limits normalize() const KJ_WARN_UNUSED_RESULT {
126 // Avoid surprises due to misconfigured bindings that set one or more limits to 0.
127 if (maxKeys == 0 || maxValueSize == 0 || maxTotalValueSize == 0) {
128 return min();
129 }
130 
131 // If a binding specifies a maxValueSize that exceeds the maxTotalValueSize, remedy
132 // that by reducing the maxValueSize.
133 return Limits{
134 .maxKeys = maxKeys,
135 .maxValueSize = static_cast<uint32_t>(kj::min(maxValueSize, maxTotalValueSize)),
136 .maxTotalValueSize = maxTotalValueSize,
137 };
138 }
139 
140 static constexpr Limits min() {
141 return {0, 0, 0};
142 }
143 
144 static Limits max(const Limits& a, const Limits& b) {
145 return Limits{
146 std::max(a.maxKeys, b.maxKeys),
147 std::max(a.maxValueSize, b.maxValueSize),
148 std::max(a.maxTotalValueSize, b.maxTotalValueSize),
149 };
150 }
151 };
152 
153 KJ_DISALLOW_COPY_AND_MOVE(SharedMemoryCache);
154 
155 using AdditionalResizeMemoryLimitHandler = kj::Function<void(ThreadUnsafeData&)>;
156 
157 SharedMemoryCache(kj::Maybe<const MemoryCacheProvider&> provider,
158 kj::StringPtr id,
159 kj::Maybe<AdditionalResizeMemoryLimitHandler&> additionalResizeMemoryLimitHandler,
160 const kj::MonotonicClock& timer);
161 
162 ~SharedMemoryCache() noexcept(false);
163 
164 kj::StringPtr getId() const {
165 return id;
166 }
167 
168 static kj::Own<const SharedMemoryCache> create(kj::Maybe<const MemoryCacheProvider&> provider,
169 kj::StringPtr id,
170 kj::Maybe<AdditionalResizeMemoryLimitHandler&> additionalResizeMemoryLimitHandler,
171 const kj::MonotonicClock& timer);
172 
173 // RAII class that attaches itself to a cache, suggests cache limits to the
174 // cache it is attached to, and allows interacting with the cache.
175 class Use {
176 public:
177 KJ_DISALLOW_COPY(Use);
178 
179 Use(kj::Own<const SharedMemoryCache> cache, const Limits& limits);
180 Use(Use&& other);
181 ~Use() noexcept(false);
182 
183 // Returns a cached value for the given key if one exists (and has not
184 // expired). If no such value exists, nothing is returned, regardless of any
185 // in-progress fallbacks trying to produce such a value.
186 kj::Maybe<kj::Own<CacheValue>> getWithoutFallback(
187 const kj::String& key, SpanBuilder& readSpan) const;
188 
189 struct FallbackResult {
190 kj::Own<CacheValue> value;
191 kj::Maybe<double> expiration;
192 };
193 using FallbackDoneCallback = kj::Function<void(kj::Maybe<FallbackResult>, SpanBuilder&)>;
194 using GetWithFallbackOutcome = kj::OneOf<kj::Own<CacheValue>, FallbackDoneCallback>;
195 
196 // Returns either:
197 // 1. The immediate value, if already in cache.
198 // 2. A Promise that will eventually resolve either to the cached value
199 // or to a FallbackDoneCallback. In the latter case, the caller should
200 // invoke the fallback function.
201 kj::OneOf<kj::Own<CacheValue>, kj::Promise<GetWithFallbackOutcome>> getWithFallback(
202 const kj::String& key, SpanBuilder& readSpan) const;
203 
204 void delete_(const kj::String& key) const;
205 
206 private:
207 // Creates a new FallbackDoneCallback associated with the given
208 // InProgress struct. This is called whenever getWithFallback() wants to
209 // invoke a fallback but it does not call the fallback directly. The caller
210 // is responsible for passing the returned task and fulfiller to the
211 // respective I/O context in which the fallback will run.
212 FallbackDoneCallback prepareFallback(InProgress& inProgress) const;
213 
214 // Called whenever a fallback has failed. The fallback might have thrown an
215 // error or it might have returned a Promise that rejected, or the I/O
216 // context in which the fallback should have been invoked has already been
217 // destroyed. If other concurrent read operations have queued fallbacks,
218 // this schedules the next fallback. Otherwise, the InProgress struct is
219 // erased.
220 void handleFallbackFailure(InProgress& inProgress) const;
221 
222 kj::Own<const SharedMemoryCache> cache;
223 static constexpr auto memoryCachekLockWaitTimeTag = "memory_cache_lock_wait_time_ns"_kjc;
224 Limits limits;
225 };
226 
227 private:
228 struct InProgress {
229 const kj::String key;
230 
231 struct Waiter {
232 kj::Own<kj::CrossThreadPromiseFulfiller<Use::GetWithFallbackOutcome>> fulfiller;
233 };
234 workerd::util::Queue<Waiter> waiting;
235 
236 InProgress(kj::String&& key): key(kj::mv(key)) {}
237 
238 // Callbacks for a HashIndex that allow locating an InProgress struct
239 // based on the cache key.
240 class KeyCallbacks {
241 public:
242 inline const kj::String& keyForRow(const kj::Own<InProgress>& entry) const {
243 return entry->key;
244 }
245 
246 template <typename KeyLike>
247 inline bool matches(const kj::Own<InProgress>& e, KeyLike&& key) const {
248 return e->key == key;
249 }
250 
251 template <typename KeyLike>
252 inline auto hashCode(KeyLike&& key) const {
253 return kj::hashCode(key);
254 }
255 };
256 };
257 
258 // Called when initializing globals (i.e., bindings) for an isolate. Each
259 // cache binding holds one SharedMemoryCache::Use, which automatically calls
260 // this function when created. This call will never reduce the effective cache
261 // limits, but might increase them.
262 void suggest(const Limits& limits) const;
263 
264 // Called when a cache global and its associated SharedMemoryCache::Use is
265 // destroyed. This call might reduce the effective cache limits. If all uses
266 // have been destroyed, the effective limits will be reset to Limits::min(),
267 // effectively clearing the cache.
268 void unsuggest(const Limits& limits) const;
269 
270 // Used internally by suggest() and unsuggest() to dynamically resize the
271 // cache as appropriate. This function also recomputed the effective cache
272 // limits and thus must be called even when the cache size is increased (which
273 // does not change the cache contents).
274 void resize(ThreadUnsafeData& data) const;
275 
276 // Returns a cached value while the cache's data is already locked by the
277 // calling thread. If such a cache entry exists, it will be marked as the
278 // most recently used entry.
279 kj::Maybe<kj::Own<CacheValue>> getWhileLocked(
280 ThreadUnsafeData& data, const kj::String& key) const;
281 
282 // Stores a value in the cache, with an optional expiration timestamp. It is
283 // marked as the most recently used entry.
284 void putWhileLocked(ThreadUnsafeData& data,
285 const kj::String& key,
286 kj::Own<CacheValue>&& value,
287 kj::Maybe<double> expiration) const;
288 
289 // Evicts at least one cache entry. The cache's data must already be locked by
290 // the calling thread, and the cache must not be empty. Expiration timestamps
291 // are only considered if called from within an I/O context or if
292 // allowOutsideIoContext is true.
293 void evictNextWhileLocked(ThreadUnsafeData& data, bool allowOutsideIoContext = false) const;
294 
295 // Removes the cache entry with the given key, if it exists.
296 void removeIfExistsWhileLocked(ThreadUnsafeData& data, const kj::String& key) const;
297 
298 // Callbacks for a HashIndex that allow locating cache entries based on the
299 // cache key, which is a string. This is used for all key-based cache
300 // operations.
301 class KeyCallbacks {
302 public:
303 inline const kj::String& keyForRow(const MemoryCacheEntry& entry) const {
304 return entry.key;
305 }
306 
307 template <typename KeyLike>
308 inline bool matches(const MemoryCacheEntry& e, KeyLike&& key) const {
309 return e.key == key;
310 }
311 
312 template <typename KeyLike>
313 inline auto hashCode(KeyLike&& key) const {
314 return kj::hashCode(key);
315 }
316 };
317 
318 // Callbacks for a TreeIndex that allow sorting cache entries by their
319 // liveliness. This is used to evict the least recently used entry.
320 class LivelinessCallbacks {
321 public:
322 inline const uint64_t& keyForRow(const MemoryCacheEntry& entry) const {
323 return entry.liveliness;
324 }
325 
326 template <typename KeyLike>
327 inline bool matches(const MemoryCacheEntry& e, KeyLike&& key) const {
328 return e.liveliness == key;
329 }
330 
331 template <typename KeyLike>
332 inline bool isBefore(const MemoryCacheEntry& e, KeyLike&& key) const {
333 return e.liveliness < key;
334 }
335 };
336 
337 // Callbacks for a TreeIndex that allow sorting cache entries by the sizes
338 // of the serialized values. The entries are sorted in reverse order, i.e.,
339 // the first entry contains the largest value. This is used to quickly evict
340 // the largest cache values when the maximum value size is reduced, e.g.,
341 // when a new version of a worker is deployed.
342 class ValueSizeCallbacks {
343 public:
344 inline const MemoryCacheEntry& keyForRow(const MemoryCacheEntry& entry KJ_LIFETIMEBOUND) const {
345 return entry;
346 }
347 
348 template <typename KeyLike>
349 inline bool matches(const MemoryCacheEntry& e, KeyLike&& key) const {
350 return e.size() == key.size() && e.key == key.key;
351 }
352 
353 template <typename KeyLike>
354 inline bool isBefore(const MemoryCacheEntry& e, KeyLike&& key) const {
355 size_t szl = e.size(), szr = key.size();
356 if (szl != szr) return szl > szr;
357 return e.key < key.key;
358 }
359 };
360 
361 // Callbacks for a TreeIndex that allow sorting cache entries by their
362 // expiration times. This is used to quickly evict expired entries even when
363 // they are not least recently used. Values with no expiration timestamp are
364 // at the very end, ordered by their cache keys.
365 class ExpirationCallbacks {
366 public:
367 inline const MemoryCacheEntry& keyForRow(const MemoryCacheEntry& entry KJ_LIFETIMEBOUND) const {
368 return entry;
369 }
370 
371 template <typename KeyLike>
372 inline bool matches(const MemoryCacheEntry& e, KeyLike&& key) const {
373 return e.expiration == key.expiration && e.key == key.key;
374 }
375 
376 template <typename KeyLike>
377 inline bool isBefore(const MemoryCacheEntry& e, KeyLike&& key) const {
378 const kj::Maybe<double>&expl = e.expiration, expr = key.expiration;
379 if (expl != expr) return isBefore(expl, expr);
380 return e.key < key.key;
381 }
382 
383 private:
384 inline bool isBefore(const kj::Maybe<double>& a, const kj::Maybe<double>& b) const {
385 KJ_IF_SOME(da, a) {
386 KJ_IF_SOME(db, b) {
387 return da < db;
388 } else {
389 return false;
390 }
391 } else {
392 KJ_ASSERT(b != kj::none);
393 return true;
394 }
395 }
396 };
397 
398 public:
399 struct ThreadUnsafeData {
400 KJ_DISALLOW_COPY_AND_MOVE(ThreadUnsafeData);
401 
402 ThreadUnsafeData() {}
403 
404 // All limits that have been suggested by isolates that are currently using
405 // this cache.
406 std::multiset<Limits> suggestedLimits;
407 
408 // The computed effective limits. These are updated whenever new isolates
409 // are attached to this cache.
410 Limits effectiveLimits = Limits::min();
411 
412 // Returns the next liveliness and increments it so that the next call to
413 // this function will return a different value.
414 inline uint64_t stepLiveliness() {
415 return nextLiveliness++;
416 }
417 
418 // We do not handle integer overflow, but a 64-bit counter should never wrap
419 // around, at least not in the foreseeable future. (Even at a billion cache
420 // operations per second, it would take almost 600 years.)
421 uint64_t nextLiveliness = 0;
422 
423 // The sum of the sizes of all values that are currently stored in the cache.
424 // This is technically redundant information, but more efficient than
425 // iterating over all cache entries every time we need this information.
426 size_t totalValueSize = 0;
427 
428 // The actual cache contents.
429 kj::Table<MemoryCacheEntry, // row type
430 kj::HashIndex<KeyCallbacks>, // index over keys
431 kj::TreeIndex<LivelinessCallbacks>, // index over liveliness
432 kj::TreeIndex<ValueSizeCallbacks>, // index over value sizes
433 kj::TreeIndex<ExpirationCallbacks> // index over expiration
434 >
435 cache;
436 
437 // Whenever a fallback is active for a particular key, this table will
438 // contain one corresponding row. Other concurrent read operations can add
439 // themselves to the InProgress struct to be notified once the fallback
440 // completes. When a fallback succeeds, this immediately notifies all
441 // waiting read operations, but when it fails, this behaves like a queue and
442 // and invokes the next available fallback only.
443 kj::Table<kj::Own<InProgress>, kj::HashIndex<InProgress::KeyCallbacks>> inProgress;
444 };
445 
446 private:
447 // To ensure thread-safety, all mutable data is guarded by a mutex. Each cache
448 // operation requires an exclusive lock. Even read-only operations need to
449 // update the liveliness of cache entries, which currently requires a lock.
450 kj::MutexGuarded<ThreadUnsafeData> data;
451 
452 // The MemoryCacheProvider instance needs to be guaranteed to outlive the SharedMemoryCache
453 // instance. When the SharedMemoryCache is destroyed, it will remove itself from the provider.
454 // TODO(cleanup): Eventually, assuming/once the kj::Ptr<T> work progresses, it would be safer
455 // to replace this with a kj::Ptr<MemoryCacheProvider>
456 kj::Maybe<const MemoryCacheProvider&> provider;
457 
458 // It's a bit unfortunate that we need to keep a copy of the id here as well as in the map
459 // in the MemoryCacheProvider, however, it's entirely possible (at least theoretically) that
460 // the map entry in the MemoryCacheProvider could be removed before the SharedMemoryCache is
461 // fully destroyed, leaving a dangling reference. This be safe and keep a copy.
462 kj::String id;
463 
464 // Same as above, the MemoryCacheProvider owns the actual handler here. Since that is guaranteed
465 // to outlive this SharedMemoryCache instance, so is the handler.
466 kj::Maybe<AdditionalResizeMemoryLimitHandler&> additionalResizeMemoryLimitHandler;
467 
468 const kj::MonotonicClock& timer;
469};
470 
471// JavaScript class that allows accessing an in-memory cache.
472// Each instance of this class holds a SharedMemoryCache::Use object and
473// all calls from JavaScript are essentially forwarded to that object, which
474// manages interaction with the shared cache in a thread-safe manner.
475class MemoryCache: public jsg::Object {
476 public:
477 MemoryCache(SharedMemoryCache::Use&& use): cacheUse(kj::mv(use)) {}
478 
479 using FallbackFunction = jsg::Function<jsg::Promise<CacheValueProduceResult>(kj::String)>;
480 
481 // Reads a value from the cache or invokes a fallback function to obtain the
482 // value, if a fallback function was given.
483 jsg::Promise<jsg::JsRef<jsg::JsValue>> read(jsg::Lock& js,
484 jsg::NonCoercible<kj::String> key,
485 jsg::Optional<FallbackFunction> optionalFallback);
486 
487 // Delete a value from the cache.
488 void delete_(jsg::Lock& js, jsg::NonCoercible<kj::String> key);
489 
490 JSG_RESOURCE_TYPE(MemoryCache, CompatibilityFlags::Reader flags) {
491 JSG_METHOD(read);
492 if (flags.getMemoryCacheDelete()) {
493 JSG_METHOD_NAMED(delete, delete_);
494 }
495 }
496 
497 private:
498 SharedMemoryCache::Use cacheUse;
499};
500 
501// The MemoryCacheProvider provides the internal implementation of the MemoryCache mechanism.
502// It is responsible for owning the SharedMemoryCache instances and providing them to the
503// bindings as needed. The default implementation (created and returned by createDefault())
504// uses a simple in-memory map to store the SharedMemoryCache instances.
505// TODO(later): It may be worth considering some kind of metrics observer for the provider
506// that can be passed along to the individual cache instances so we can monitor just how much
507// the in memory cache is being used.
508class MemoryCacheProvider {
509 public:
510 MemoryCacheProvider(const kj::MonotonicClock& timer,
511 kj::Maybe<SharedMemoryCache::AdditionalResizeMemoryLimitHandler>
512 additionalResizeMemoryLimitHandler = kj::none);
513 KJ_DISALLOW_COPY_AND_MOVE(MemoryCacheProvider);
514 ~MemoryCacheProvider() noexcept(false);
515 
516 kj::Own<const SharedMemoryCache> getInstance(kj::Maybe<kj::StringPtr> cacheId = kj::none) const;
517 
518 void removeInstance(const SharedMemoryCache& instance) const;
519 
520 private:
521 kj::Maybe<SharedMemoryCache::AdditionalResizeMemoryLimitHandler>
522 additionalResizeMemoryLimitHandler;
523 
524 // All existing in-memory *shared* caches. This table will not include caches created
525 // that do not have an id (and therefore cannot be shared).
526 // TODO(cleanup): Later, assuming progress is made on kj::Ptr<T>, it would be nice
527 // to avoid the use of the bare pointer to SharedMemoryCache* here. When the SharedMemoryCache
528 // is destroyed, it will remove itself from this cache by calling removeInstance.
529 kj::MutexGuarded<kj::HashMap<kj::String, const SharedMemoryCache*>> caches;
530 
531 const kj::MonotonicClock& timer;
532};
533 
534// clang-format off
535#define EW_MEMORY_CACHE_ISOLATE_TYPES \
536 api::MemoryCache, \
537 api::CacheValueProduceResult
538// clang-format on
539 
540} // namespace workerd::api