Skip to content
File

Blob: src/workerd/api/actor-state.c++

55.8 KB
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#include "actor-state.h"
6 
7#include "actor.h"
8#include "export-loopback.h"
9#include "sql.h"
10#include "sync-kv.h"
11#include "util.h"
12 
13#include <workerd/api/web-socket.h>
14#include <workerd/io/actor-cache.h>
15#include <workerd/io/actor-id.h>
16#include <workerd/io/actor-sqlite.h>
17#include <workerd/io/features.h>
18#include <workerd/io/hibernation-manager.h>
19#include <workerd/jsg/jsg.h>
20#include <workerd/jsg/ser.h>
21#include <workerd/jsg/util.h>
22 
23#include <v8.h>
24 
25namespace workerd::api {
26 
27namespace {
28 
29constexpr size_t BILLING_UNIT = 4096;
30 
31enum class BillAtLeastOne { NO, YES };
32 
33uint32_t billingUnits(size_t bytes, BillAtLeastOne billAtLeastOne = BillAtLeastOne::YES) {
34 if (billAtLeastOne == BillAtLeastOne::YES && bytes == 0) {
35 return 1; // always bill for at least 1 billing unit
36 }
37 return bytes / BILLING_UNIT + (bytes % BILLING_UNIT != 0);
38}
39 
40jsg::JsValue deserializeMaybeV8Value(
41 jsg::Lock& js, kj::ArrayPtr<const char> key, kj::Maybe<kj::ArrayPtr<const kj::byte>> buf) {
42 KJ_IF_SOME(b, buf) {
43 return deserializeV8Value(js, key, b);
44 } else {
45 return js.undefined();
46 }
47}
48 
49template <typename T, typename Options, typename Func>
50auto transformCacheResult(jsg::Lock& js,
51 kj::OneOf<T, kj::Promise<T>> input,
52 const Options& options,
53 Func&& func) -> jsg::Promise<decltype(func(js, kj::instance<T>()))> {
54 KJ_SWITCH_ONEOF(input) {
55 KJ_CASE_ONEOF(value, T) {
56 return js.resolvedPromise(func(js, kj::mv(value)));
57 }
58 KJ_CASE_ONEOF(promise, kj::Promise<T>) {
59 auto& context = IoContext::current();
60 if (options.allowConcurrency.orDefault(false)) {
61 return context.awaitIo(
62 js, kj::mv(promise), [func = kj::fwd<Func>(func)](jsg::Lock& js, T&& value) mutable {
63 return func(js, kj::mv(value));
64 });
65 } else {
66 return context.awaitIoWithInputLock(
67 js, kj::mv(promise), [func = kj::fwd<Func>(func)](jsg::Lock& js, T&& value) mutable {
68 return func(js, kj::mv(value));
69 });
70 }
71 }
72 }
73 KJ_UNREACHABLE;
74}
75 
76template <typename T, typename Options, typename Func>
77auto transformCacheResultWithCacheStatus(jsg::Lock& js,
78 kj::OneOf<T, kj::Promise<T>> input,
79 const Options& options,
80 Func&& func) -> jsg::Promise<decltype(func(js, kj::instance<T>(), kj::instance<bool>()))> {
81 KJ_SWITCH_ONEOF(input) {
82 KJ_CASE_ONEOF(value, T) {
83 return js.resolvedPromise(func(js, kj::mv(value), true));
84 }
85 KJ_CASE_ONEOF(promise, kj::Promise<T>) {
86 auto& context = IoContext::current();
87 if (options.allowConcurrency.orDefault(false)) {
88 return context.awaitIo(
89 js, kj::mv(promise), [func = kj::fwd<Func>(func)](jsg::Lock& js, T&& value) mutable {
90 return func(js, kj::mv(value), false);
91 });
92 } else {
93 return context.awaitIoWithInputLock(
94 js, kj::mv(promise), [func = kj::fwd<Func>(func)](jsg::Lock& js, T&& value) mutable {
95 return func(js, kj::mv(value), false);
96 });
97 }
98 }
99 }
100 KJ_UNREACHABLE;
101}
102 
103template <typename Options>
104jsg::Promise<void> transformMaybeBackpressure(
105 jsg::Lock& js, const Options& options, kj::Maybe<kj::Promise<void>> maybeBackpressure) {
106 KJ_IF_SOME(backpressure, maybeBackpressure) {
107 // Note: In practice `allowConcurrency` will have no effect on a backpressure promise since
108 // backpressure blocks everything anyway, but we pass the option through for consistency in
109 // case of future changes.
110 auto& context = IoContext::current();
111 if (options.allowConcurrency.orDefault(false)) {
112 return context.awaitIo(js, kj::mv(backpressure));
113 } else {
114 return context.awaitIoWithInputLock(js, kj::mv(backpressure), [](jsg::Lock&) {});
115 }
116 } else {
117 return js.resolvedPromise();
118 }
119}
120 
121ActorObserver& currentActorMetrics() {
122 return IoContext::current().getActorOrThrow().getMetrics();
123}
124 
125jsg::JsRef<jsg::JsValue> listResultsToMap(
126 jsg::Lock& js, ActorCacheOps::GetResultList value, bool completelyCached) {
127 return js
128 .withinHandleScope([&] {
129 auto map = js.map();
130 size_t cachedReadBytes = 0;
131 size_t uncachedReadBytes = 0;
132 for (const auto& entry: value) {
133 auto& bytesRef =
134 entry.status == ActorCacheOps::CacheStatus::CACHED ? cachedReadBytes : uncachedReadBytes;
135 bytesRef += entry.key.size() + entry.value.size();
136 map.set(js, entry.key, deserializeV8Value(js, entry.key, entry.value));
137 }
138 auto& actorMetrics = currentActorMetrics();
139 if (cachedReadBytes || uncachedReadBytes) {
140 size_t totalReadBytes = cachedReadBytes + uncachedReadBytes;
141 uint32_t totalUnits = billingUnits(totalReadBytes);
142 
143 // If we went to disk, we want to ensure we bill at least 1 uncached unit.
144 // Otherwise, we disable this behavior, to ensure a fully cached list will have
145 // uncachedUnits == 0.
146 auto billAtLeastOne = completelyCached ? BillAtLeastOne::NO : BillAtLeastOne::YES;
147 uint32_t uncachedUnits = billingUnits(uncachedReadBytes, billAtLeastOne);
148 uint32_t cachedUnits = totalUnits - uncachedUnits;
149 
150 actorMetrics.addUncachedStorageReadUnits(uncachedUnits);
151 actorMetrics.addCachedStorageReadUnits(cachedUnits);
152 } else {
153 // We bill 1 uncached read unit if there was no results from the list.
154 actorMetrics.addUncachedStorageReadUnits(1);
155 }
156 
157 return jsg::JsValue(map);
158 }).addRef(js);
159}
160 
161kj::Function<jsg::JsRef<jsg::JsValue>(jsg::Lock&, ActorCacheOps::GetResultList)>
162getMultipleResultsToMap(size_t numInputKeys) {
163 return [numInputKeys](jsg::Lock& js, ActorCacheOps::GetResultList value) mutable {
164 return js
165 .withinHandleScope([&] {
166 auto map = js.map();
167 uint32_t cachedUnits = 0;
168 uint32_t uncachedUnits = 0;
169 for (const auto& entry: value) {
170 auto& unitsRef =
171 entry.status == ActorCacheOps::CacheStatus::CACHED ? cachedUnits : uncachedUnits;
172 unitsRef += billingUnits(entry.key.size() + entry.value.size());
173 map.set(js, entry.key, deserializeV8Value(js, entry.key, entry.value));
174 }
175 auto& actorMetrics = currentActorMetrics();
176 actorMetrics.addCachedStorageReadUnits(cachedUnits);
177 
178 size_t leftoverKeys = 0;
179 if (numInputKeys >= value.size()) {
180 leftoverKeys = numInputKeys - value.size();
181 } else {
182 KJ_LOG(ERROR, "More returned pairs than provided input keys in getMultipleResultsToMap",
183 numInputKeys, value.size());
184 }
185 
186 // leftover keys weren't in the result set, but potentially still
187 // had to be queried for existence.
188 //
189 // TODO(someday): This isn't quite accurate -- we do cache negative entries.
190 // Billing will still be correct today, but if we do ever start billing
191 // only for uncached reads, we'll need to address this.
192 actorMetrics.addUncachedStorageReadUnits(leftoverKeys + uncachedUnits);
193 
194 return jsg::JsValue(map);
195 }).addRef(js);
196 };
197}
198 
199kj::Promise<void> updateStorageWriteUnit(
200 IoContext& context, ActorObserver& metrics, uint32_t units) {
201 // The ActorObserver& reference here is guaranteed to outlive this task, so
202 // accessing it after the co_await here is safe.
203 co_await context.waitForOutputLocks();
204 metrics.addStorageWriteUnits(units);
205}
206 
207kj::Promise<void> updateStorageDeletes(
208 IoContext& context, ActorObserver& metrics, kj::Promise<uint> promise) {
209 // The ActorObserver& reference here is guaranteed to outlive this task, so
210 // accessing it after the co_await here is safe.
211 auto deleted = co_await promise;
212 if (deleted == 0) deleted = 1;
213 metrics.addStorageDeletes(deleted);
214};
215 
216// Return the id of the current actor (or the empty string if there is no current actor).
217kj::Maybe<kj::String> getCurrentActorId() {
218 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
219 KJ_IF_SOME(actor, ioContext.getActor()) {
220 KJ_SWITCH_ONEOF(actor.getId()) {
221 KJ_CASE_ONEOF(s, kj::String) {
222 return kj::heapString(s);
223 }
224 KJ_CASE_ONEOF(actorId, kj::Own<ActorIdFactory::ActorId>) {
225 return actorId->toString();
226 }
227 }
228 KJ_UNREACHABLE;
229 }
230 }
231 return kj::none;
232}
233 
234} // namespace
235 
236DurableObjectStorage::DurableObjectStorage(jsg::Lock& js,
237 IoPtr<ActorCacheInterface> cache,
238 bool enableSql,
239 kj::Own<IoChannelFactory::ActorChannel> primaryActorChannel,
240 kj::Own<ActorIdFactory::ActorId> primaryActorId)
241 : cache(kj::mv(cache)),
242 enableSql(enableSql) {
243 
244 auto replicaFactory = kj::heap<ReplicaActorOutgoingFactory>(
245 kj::mv(primaryActorChannel), primaryActorId->toString());
246 auto outgoingFactory =
247 IoContext::current().addObject<Fetcher::OutgoingFactory>(kj::mv(replicaFactory));
248 auto requiresHost = FeatureFlags::get(IoContext::current().getCurrentLock())
249 .getDurableObjectFetchRequiresSchemeAuthority()
250 ? Fetcher::RequiresHostAndProtocol::YES
251 : Fetcher::RequiresHostAndProtocol::NO;
252 
253 this->maybePrimary = js.alloc<DurableObject>(
254 js.alloc<DurableObjectId>(kj::mv(primaryActorId)), kj::mv(outgoingFactory), requiresHost);
255}
256 
257jsg::Promise<jsg::JsRef<jsg::JsValue>> DurableObjectStorageOperations::get(jsg::Lock& js,
258 kj::OneOf<kj::String, kj::Array<kj::String>> keys,
259 jsg::Optional<GetOptions> maybeOptions) {
260 auto& context = IoContext::current();
261 auto traceContext = context.makeUserTraceSpan("durable_object_storage_get"_kjc);
262 auto options = configureOptions(kj::mv(maybeOptions).orDefault(GetOptions{}));
263 KJ_SWITCH_ONEOF(keys) {
264 KJ_CASE_ONEOF(s, kj::String) {
265 return context.attachSpans(js, getOne(js, kj::mv(s), options), kj::mv(traceContext));
266 }
267 KJ_CASE_ONEOF(a, kj::Array<kj::String>) {
268 return context.attachSpans(js, getMultiple(js, kj::mv(a), options), kj::mv(traceContext));
269 }
270 }
271 KJ_UNREACHABLE
272}
273 
274jsg::Promise<jsg::JsRef<jsg::JsValue>> DurableObjectStorageOperations::getOne(
275 jsg::Lock& js, kj::String key, const GetOptions& options) {
276 auto result = getCache(OP_GET).get(kj::str(key), options);
277 return transformCacheResultWithCacheStatus(js, kj::mv(result), options,
278 [key = kj::mv(key)](jsg::Lock& js, kj::Maybe<ActorCacheOps::Value> value, bool cached) {
279 uint32_t units = 1;
280 KJ_IF_SOME(v, value) {
281 units = billingUnits(v.size());
282 }
283 auto& actorMetrics = currentActorMetrics();
284 if (cached) {
285 actorMetrics.addCachedStorageReadUnits(units);
286 } else {
287 actorMetrics.addUncachedStorageReadUnits(units);
288 }
289 return deserializeMaybeV8Value(js, key, value).addRef(js);
290 });
291}
292 
293jsg::Promise<kj::Maybe<double>> DurableObjectStorageOperations::getAlarm(
294 jsg::Lock& js, jsg::Optional<GetAlarmOptions> maybeOptions) {
295 auto& context = IoContext::current();
296 auto traceContext = context.makeUserTraceSpan("durable_object_storage_getAlarm"_kjc);
297 // Even if we do not have an alarm handler, we might once have had one. It's fine to return
298 // whatever a previous alarm setting or a falsy result.
299 auto options = configureOptions(maybeOptions
300 .map([](auto& o) {
301 return GetOptions{.allowConcurrency = o.allowConcurrency, .noCache = false};
302 }).orDefault(GetOptions{}));
303 auto result = getCache(OP_GET_ALARM).getAlarm(options);
304 
305 return context.attachSpans(js,
306 transformCacheResult(js, kj::mv(result), options,
307 [](jsg::Lock&, kj::Maybe<kj::Date> date) {
308 return date.map(
309 [](auto& date) { return static_cast<double>((date - kj::UNIX_EPOCH) / kj::MILLISECONDS); });
310 }),
311 kj::mv(traceContext));
312}
313 
314kj::Maybe<DurableObjectStorageOperations::CompiledListOptions> DurableObjectStorageOperations::
315 compileListOptions(kj::Maybe<ListOptions>& maybeOptions) {
316 kj::String start;
317 kj::Maybe<kj::String> end;
318 bool reverse = false;
319 kj::Maybe<uint> limit;
320 
321 KJ_IF_SOME(o, maybeOptions) {
322 KJ_IF_SOME(s, o.start) {
323 if (o.startAfter != kj::none) {
324 KJ_FAIL_REQUIRE(
325 "jsg.TypeError: list() cannot be called with both start and startAfter values.");
326 }
327 start = kj::mv(s);
328 }
329 KJ_IF_SOME(sks, o.startAfter) {
330 // Convert an exclusive startAfter into an inclusive start key here so that the implementation
331 // doesn't need to handle both. This can be done simply by adding two NULL bytes. One to the end of
332 // the startAfter and another to set the start key after startAfter.
333 auto startAfterKey = kj::heapArray<char>(sks.size() + 2);
334 
335 // Copy over the original string.
336 memcpy(startAfterKey.begin(), sks.begin(), sks.size());
337 // Add one additional null byte to set the new start as the key immediately
338 // after startAfter. This looks a little sketchy to be doing with strings rather
339 // than arrays, but kj::String explicitly allows for NULL bytes inside of strings.
340 startAfterKey[startAfterKey.size() - 2] = '\0';
341 // kj::String automatically reads the last NULL as string termination, so we need to add it twice
342 // to make it stick in the final string.
343 startAfterKey[startAfterKey.size() - 1] = '\0';
344 start = kj::String(kj::mv(startAfterKey));
345 }
346 KJ_IF_SOME(e, o.end) {
347 end = kj::mv(e);
348 }
349 KJ_IF_SOME(r, o.reverse) {
350 reverse = r;
351 }
352 KJ_IF_SOME(l, o.limit) {
353 JSG_REQUIRE(l > 0, TypeError, "List limit must be positive.");
354 limit = l;
355 }
356 KJ_IF_SOME(prefix, o.prefix) {
357 // Let's clamp `start` and `end` to include only keys with the given prefix.
358 if (prefix.size() > 0) {
359 if (start < prefix) {
360 // `start` is before `prefix`, so listing should actually start at `prefix`.
361 start = kj::str(prefix);
362 } else if (start.startsWith(prefix)) {
363 // `start` is within the prefix, so need not be modified.
364 } else {
365 // `start` comes after the last value with the prefix, so there's no overlap.
366 return kj::none;
367 }
368 
369 // Calculate the first key that sorts after all keys with the given prefix.
370 kj::Vector<char> keyAfterPrefix(prefix.size());
371 keyAfterPrefix.addAll(prefix);
372 while (!keyAfterPrefix.empty() && static_cast<byte>(keyAfterPrefix.back()) == 0xff) {
373 keyAfterPrefix.removeLast();
374 }
375 if (keyAfterPrefix.empty()) {
376 // The prefix is a string of some number of 0xff bytes, so includes the entire key space
377 // up through the last possible key. Hence, there is no end. (But if an end was specified
378 // earlier, that's still valid.)
379 } else {
380 keyAfterPrefix.back()++;
381 keyAfterPrefix.add('\0');
382 auto keyAfterPrefixStr = kj::String(keyAfterPrefix.releaseAsArray());
383 
384 KJ_IF_SOME(e, end) {
385 if (e <= prefix) {
386 // No keys could possibly match both the end and the prefix.
387 return kj::none;
388 } else if (e.startsWith(prefix)) {
389 // `end` is within the prefix, so need not be modified.
390 } else {
391 // `end` comes after all keys with the prefix, so we should stop at the end of the
392 // prefix.
393 end = kj::mv(keyAfterPrefixStr);
394 }
395 } else {
396 // We didn't have any end set, so use the end of the prefix range.
397 end = kj::mv(keyAfterPrefixStr);
398 }
399 }
400 }
401 }
402 }
403 
404 KJ_IF_SOME(e, end) {
405 if (e <= start) {
406 // Key range is empty.
407 return kj::none;
408 }
409 }
410 
411 return CompiledListOptions{
412 .start = kj::mv(start),
413 .end = kj::mv(end),
414 .reverse = reverse,
415 .limit = limit,
416 };
417}
418 
419jsg::Promise<jsg::JsRef<jsg::JsValue>> DurableObjectStorageOperations::list(
420 jsg::Lock& js, jsg::Optional<ListOptions> maybeOptions) {
421 auto& context = IoContext::current();
422 auto traceContext = context.makeUserTraceSpan("durable_object_storage_list"_kjc);
423 auto [start, end, reverse, limit] = KJ_UNWRAP_OR(compileListOptions(maybeOptions),
424 { return js.resolvedPromise(jsg::JsValue(js.map()).addRef(js)); });
425 
426 auto options = configureOptions(kj::mv(maybeOptions).orDefault(ListOptions{}));
427 ActorCacheOps::ReadOptions readOptions = options;
428 
429 auto result = reverse
430 ? getCache(OP_LIST).listReverse(kj::mv(start), kj::mv(end), limit, readOptions)
431 : getCache(OP_LIST).list(kj::mv(start), kj::mv(end), limit, readOptions);
432 return context.attachSpans(js,
433 transformCacheResultWithCacheStatus(js, kj::mv(result), options, &listResultsToMap),
434 kj::mv(traceContext));
435}
436 
437jsg::Promise<void> DurableObjectStorageOperations::put(jsg::Lock& js,
438 kj::OneOf<kj::String, jsg::Dict<jsg::JsValue>> keyOrEntries,
439 jsg::Optional<jsg::JsValue> value,
440 jsg::Optional<PutOptions> maybeOptions,
441 const jsg::TypeHandler<PutOptions>& optionsTypeHandler) {
442 auto& context = IoContext::current();
443 auto traceContext = context.makeUserTraceSpan("durable_object_storage_put"_kjc);
444 // TODO(soon): Add tests of data generated at current versions to ensure we'll
445 // know before releasing any backwards-incompatible serializer changes,
446 // potentially checking the header in addition to the value.
447 auto options = configureOptions(kj::mv(maybeOptions).orDefault(PutOptions{}));
448 KJ_SWITCH_ONEOF(keyOrEntries) {
449 KJ_CASE_ONEOF(k, kj::String) {
450 KJ_IF_SOME(v, value) {
451 return context.attachSpans(js, putOne(js, kj::mv(k), v, options), kj::mv(traceContext));
452 } else {
453 JSG_FAIL_REQUIRE(TypeError, "put() called with undefined value.");
454 }
455 }
456 KJ_CASE_ONEOF(o, jsg::Dict<jsg::JsValue>) {
457 KJ_IF_SOME(v, value) {
458 KJ_IF_SOME(opt, optionsTypeHandler.tryUnwrap(js, v)) {
459 // return putMultiple(js, kj::mv(o), configureOptions(kj::mv(opt)));
460 return context.attachSpans(
461 js, putMultiple(js, kj::mv(o), configureOptions(kj::mv(opt))), kj::mv(traceContext));
462 } else {
463 JSG_FAIL_REQUIRE(TypeError,
464 "put() may only be called with a single key-value pair and optional options as put(key, value, options) or with multiple key-value pairs and optional options as put(entries, options)");
465 }
466 } else {
467 // return putMultiple(js, kj::mv(o), options);
468 return context.attachSpans(js, putMultiple(js, kj::mv(o), options), kj::mv(traceContext));
469 }
470 }
471 }
472 KJ_UNREACHABLE;
473}
474 
475jsg::Promise<void> DurableObjectStorageOperations::setAlarm(
476 jsg::Lock& js, kj::Date scheduledTime, jsg::Optional<SetAlarmOptions> maybeOptions) {
477 JSG_REQUIRE(scheduledTime > kj::origin<kj::Date>(), TypeError,
478 "setAlarm() cannot be called with an alarm time <= 0");
479 
480 auto& context = IoContext::current();
481 auto traceContext = context.makeUserTraceSpan("durable_object_storage_setAlarm"_kjc);
482 // This doesn't check if we have an alarm handler per say. It checks if we have an initialized
483 // (post-ctor) JS durable object with an alarm handler. Notably, this means this won't throw if
484 // `setAlarm` is invoked in the DO ctor even if the DO class does not have an alarm handler. This
485 // is better than throwing even if we do have an alarm handler.
486 context.getActorOrThrow().assertCanSetAlarm();
487 
488 auto options = configureOptions(maybeOptions
489 .map([](auto& o) {
490 return PutOptions{.allowConcurrency = o.allowConcurrency,
491 .allowUnconfirmed = o.allowUnconfirmed,
492 .noCache = false};
493 }).orDefault(PutOptions{}));
494 
495 // We fudge times set in the past to Date.now() to ensure that any one user can't DDOS the alarm
496 // polling system by putting dates far in the past and therefore getting sorted earlier by the index.
497 // This also ensures uniqueness of alarm times (which is required for correctness),
498 // in the situation where customers use a constant date in the past to indicate
499 // they want immediate execution.
500 kj::Date dateNowKjDate = static_cast<int64_t>(dateNow()) * kj::MILLISECONDS + kj::UNIX_EPOCH;
501 
502 auto maybeBackpressure = transformMaybeBackpressure(js, options,
503 getCache(OP_PUT_ALARM)
504 .setAlarm(kj::max(scheduledTime, dateNowKjDate), options, context.getCurrentTraceSpan()));
505 
506 // setAlarm() is billed as a single write unit.
507 context.addTask(updateStorageWriteUnit(context, currentActorMetrics(), 1));
508 
509 return context.attachSpans(js, kj::mv(maybeBackpressure), kj::mv(traceContext));
510}
511 
512jsg::Promise<void> DurableObjectStorageOperations::putOne(
513 jsg::Lock& js, kj::String key, jsg::JsValue value, const PutOptions& options) {
514 
515 kj::Array<byte> buffer = serializeV8Value(js, value);
516 
517 auto units = billingUnits(key.size() + buffer.size());
518 
519 auto& context = IoContext::current();
520 
521 jsg::Promise<void> maybeBackpressure = transformMaybeBackpressure(js, options,
522 getCache(OP_PUT).put(kj::mv(key), kj::mv(buffer), options, context.getCurrentTraceSpan()));
523 
524 context.addTask(updateStorageWriteUnit(context, currentActorMetrics(), units));
525 return maybeBackpressure;
526}
527 
528kj::OneOf<jsg::Promise<bool>, jsg::Promise<int>> DurableObjectStorageOperations::delete_(
529 jsg::Lock& js,
530 kj::OneOf<kj::String, kj::Array<kj::String>> keys,
531 jsg::Optional<PutOptions> maybeOptions) {
532 auto& context = IoContext::current();
533 auto traceContext = context.makeUserTraceSpan("durable_object_storage_delete"_kjc);
534 auto options = configureOptions(kj::mv(maybeOptions).orDefault(PutOptions{}));
535 KJ_SWITCH_ONEOF(keys) {
536 KJ_CASE_ONEOF(s, kj::String) {
537 return context.attachSpans(js, deleteOne(js, kj::mv(s), options), kj::mv(traceContext));
538 }
539 KJ_CASE_ONEOF(a, kj::Array<kj::String>) {
540 return context.attachSpans(js, deleteMultiple(js, kj::mv(a), options), kj::mv(traceContext));
541 }
542 }
543 KJ_UNREACHABLE
544}
545 
546jsg::Promise<void> DurableObjectStorageOperations::deleteAlarm(
547 jsg::Lock& js, jsg::Optional<SetAlarmOptions> maybeOptions) {
548 auto& context = IoContext::current();
549 auto traceContext = context.makeUserTraceSpan("durable_object_storage_deleteAlarm"_kjc);
550 // Even if we do not have an alarm handler, we might once have had one. It's fine to remove that
551 // alarm or noop on the absence of one.
552 auto options = configureOptions(maybeOptions
553 .map([](auto& o) {
554 return PutOptions{.allowConcurrency = o.allowConcurrency,
555 .allowUnconfirmed = o.allowUnconfirmed,
556 .noCache = false};
557 }).orDefault(PutOptions{}));
558 
559 return context.attachSpans(js,
560 transformMaybeBackpressure(js, options,
561 getCache(OP_DELETE_ALARM).setAlarm(kj::none, options, context.getCurrentTraceSpan())),
562 kj::mv(traceContext));
563}
564 
565jsg::Promise<void> DurableObjectStorage::deleteAll(
566 jsg::Lock& js, jsg::Optional<PutOptions> maybeOptions) {
567 auto& context = IoContext::current();
568 auto traceContext = context.makeUserTraceSpan("durable_object_storage_deleteAll"_kjc);
569 auto options = configureOptions(kj::mv(maybeOptions).orDefault(PutOptions{}));
570 
571 DeleteAllOptions deleteAllOptions{
572 .deleteAlarm = FeatureFlags::get(js).getDeleteAllDeletesAlarm(),
573 };
574 auto deleteAll = cache->deleteAll(options, context.getCurrentTraceSpan(), deleteAllOptions);
575 
576 context.addTask(updateStorageDeletes(context, currentActorMetrics(), kj::mv(deleteAll.count)));
577 
578 return context.attachSpans(js,
579 transformMaybeBackpressure(js, options, kj::mv(deleteAll.backpressure)),
580 kj::mv(traceContext));
581}
582 
583void DurableObjectTransaction::deleteAll() {
584 JSG_FAIL_REQUIRE(Error, "Cannot call deleteAll() within a transaction");
585}
586 
587jsg::Promise<bool> DurableObjectStorageOperations::deleteOne(
588 jsg::Lock& js, kj::String key, const PutOptions& options) {
589 auto& context = IoContext::current();
590 
591 return transformCacheResult(js,
592 getCache(OP_DELETE).delete_(kj::mv(key), options, context.getCurrentTraceSpan()), options,
593 [](jsg::Lock&, bool value) {
594 currentActorMetrics().addStorageDeletes(1);
595 return value;
596 });
597}
598 
599jsg::Promise<jsg::JsRef<jsg::JsValue>> DurableObjectStorageOperations::getMultiple(
600 jsg::Lock& js, kj::Array<kj::String> keys, const GetOptions& options) {
601 auto numKeys = keys.size();
602 
603 return transformCacheResult(
604 js, getCache(OP_GET).get(kj::mv(keys), options), options, getMultipleResultsToMap(numKeys));
605}
606 
607jsg::Promise<void> DurableObjectStorageOperations::putMultiple(
608 jsg::Lock& js, jsg::Dict<jsg::JsValue> entries, const PutOptions& options) {
609 kj::Vector<ActorCacheOps::KeyValuePair> kvs(entries.fields.size());
610 
611 uint32_t units = 0;
612 for (auto& field: entries.fields) {
613 if (field.value.isUndefined()) continue;
614 // We silently drop fields with value=undefined in putMultiple. There aren't many good options here, as
615 // deleting an undefined field is confusing, throwing could break otherwise working code, and
616 // a stray undefined here or there is probably closer to what the user desires.
617 
618 kj::Array<byte> buffer = serializeV8Value(js, field.value);
619 
620 units += billingUnits(field.name.size() + buffer.size());
621 
622 kvs.add(ActorCacheOps::KeyValuePair{kj::mv(field.name), kj::mv(buffer)});
623 }
624 
625 auto& context = IoContext::current();
626 
627 jsg::Promise<void> maybeBackpressure = transformMaybeBackpressure(js, options,
628 getCache(OP_PUT).put(kvs.releaseAsArray(), options, context.getCurrentTraceSpan()));
629 
630 context.addTask(updateStorageWriteUnit(context, currentActorMetrics(), units));
631 
632 return maybeBackpressure;
633}
634 
635jsg::Promise<int> DurableObjectStorageOperations::deleteMultiple(
636 jsg::Lock& js, kj::Array<kj::String> keys, const PutOptions& options) {
637 auto numKeys = keys.size();
638 
639 auto& context = IoContext::current();
640 
641 return transformCacheResult(js,
642 getCache(OP_DELETE).delete_(kj::mv(keys), options, context.getCurrentTraceSpan()), options,
643 [numKeys](jsg::Lock&, uint count) -> int {
644 currentActorMetrics().addStorageDeletes(numKeys);
645 return count;
646 });
647}
648 
649ActorCacheOps& DurableObjectStorage::getCache(OpName op) {
650 return *cache;
651}
652 
653jsg::Promise<jsg::JsRef<jsg::JsValue>> DurableObjectStorage::transaction(jsg::Lock& js,
654 jsg::Function<jsg::Promise<jsg::JsRef<jsg::JsValue>>(jsg::Ref<DurableObjectTransaction>)>
655 callback,
656 jsg::Optional<TransactionOptions> options) {
657 auto& context = IoContext::current();
658 auto traceContext = context.makeUserTraceSpan("durable_object_storage_transaction"_kjc);
659 
660 struct TxnResult {
661 jsg::JsRef<jsg::JsValue> value;
662 bool isError;
663 };
664 
665 return context.attachSpans(js,
666 context
667 .blockConcurrencyWhile(js,
668 [callback = kj::mv(callback), &context, &cache = *cache](
669 jsg::Lock& js) mutable -> jsg::Promise<TxnResult> {
670 // Note that the call to `startTransaction()` is when the SQLite-backed implementation will
671 // actually invoke `BEGIN TRANSACTION`, so it's important that we're inside the
672 // blockConcurrencyWhile block before that point so we don't accidentally catch some other
673 // asynchronous event in our transaction.
674 //
675 // For the ActorCache-based implementation, it doesn't matter when we call `startTransaction()`
676 // as the method merely allocates an object and returns it with no side effects.
677 auto txn = js.alloc<DurableObjectTransaction>(context.addObject(cache.startTransaction()));
678 
679 return js.resolvedPromise(txn.addRef())
680 .then(js, kj::mv(callback))
681 .then(js, [txn = txn.addRef()](jsg::Lock& js, jsg::JsRef<jsg::JsValue> value) mutable {
682 // In correct usage, `context` should not have changed here, particularly because we're in
683 // a critical section so it should have been impossible for any other context to receive
684 // control. However, depending on all that is a bit precarious. jsg::Promise::then() itself
685 // does NOT guarantee it runs in the same context (the application could have returned a
686 // custom Promise and then resolved in from some other context). So let's be safe and grab
687 // IoContext::current() again here, rather than capture it in the lambda.
688 auto& context = IoContext::current();
689 return context.awaitIoWithInputLock(js, txn->maybeCommit(),
690 [value = kj::mv(value)](jsg::Lock&) mutable { return TxnResult{kj::mv(value), false}; });
691 }, [txn = txn.addRef()](jsg::Lock& js, jsg::Value exception) mutable {
692 // The transaction callback threw an exception. We don't actually want to reset the object,
693 // we only want to roll back the transaction and propagate the exception. So, we carefully
694 // pack the exception away into a value.
695 txn->maybeRollback();
696 return js.resolvedPromise(TxnResult{
697 // TODO(cleanup): Simplify this once exception is passed using jsg::JsRef instead
698 // of jsg::V8Ref
699 jsg::JsValue(exception.getHandle(js)).addRef(js), true});
700 });
701 })
702 .then(js,
703 [](jsg::Lock& js, TxnResult result) -> jsg::JsRef<jsg::JsValue> {
704 if (result.isError) {
705 js.throwException(result.value.getHandle(js));
706 } else {
707 return kj::mv(result.value);
708 }
709 }),
710 kj::mv(traceContext));
711}
712 
713jsg::JsRef<jsg::JsValue> DurableObjectStorage::transactionSync(
714 jsg::Lock& js, jsg::Function<jsg::JsRef<jsg::JsValue>()> callback) {
715 KJ_IF_SOME(sqlite, cache->getSqliteDatabase()) {
716 // SAVEPOINT is a readonly statement, but we need to trigger an outer TRANSACTION
717 sqlite.notifyWrite();
718 
719 uint depth = transactionSyncDepth++;
720 KJ_DEFER(--transactionSyncDepth);
721 
722 // TODO(perf): SQLite actually allows multiple savepoints with the same name. The name refers
723 // to the most-recent of these savepoints. This means we don't actually have to append the
724 // depth to each savepoint name like I originally thought. We should refactor this -- and use
725 // prepared statements.
726 
727 sqlite.run(
728 {.regulator = SqliteDatabase::TRUSTED}, kj::str("SAVEPOINT _cf_sync_savepoint_", depth));
729 return js.tryCatch([&]() {
730 auto result = callback(js);
731 
732 // If a critical error forced an automatic rollback, we throw an exception to convey failure
733 // to the caller of transactionSync(), even if the callback did not throw.
734 JSG_REQUIRE(!sqlite.observedCriticalError(), Error,
735 "Cannot commit transaction due to an earlier SQL critical error");
736 
737 sqlite.run(
738 {.regulator = SqliteDatabase::TRUSTED}, kj::str("RELEASE _cf_sync_savepoint_", depth));
739 return kj::mv(result);
740 }, [&](jsg::Value exception) -> jsg::JsRef<jsg::JsValue> {
741 // If a critical error forced an automatic rollback, we skip the rollback and release
742 // attempt, because savepoints should already be released.
743 if (!sqlite.observedCriticalError()) {
744 sqlite.run({.regulator = SqliteDatabase::TRUSTED},
745 kj::str("ROLLBACK TO _cf_sync_savepoint_", depth));
746 sqlite.run(
747 {.regulator = SqliteDatabase::TRUSTED}, kj::str("RELEASE _cf_sync_savepoint_", depth));
748 }
749 js.throwException(kj::mv(exception));
750 });
751 } else {
752 JSG_FAIL_REQUIRE(Error, "Durable Object is not backed by SQL.");
753 }
754}
755 
756jsg::Promise<void> DurableObjectStorage::sync(jsg::Lock& js) {
757 auto& context = IoContext::current();
758 auto traceContext = context.makeUserTraceSpan("durable_object_storage_sync"_kjc);
759 KJ_IF_SOME(p, cache->onNoPendingFlush(traceContext.getInternalSpanParent())) {
760 // Note that we're not actually flushing since that will happen anyway once we go async. We're
761 // merely checking if we have any pending or in-flight operations, and providing a promise that
762 // resolves when they succeed. This promise only covers operations that were scheduled before
763 // this method was invoked. If the cache has to flush again later from future operations, this
764 // promise will resolve before they complete. If this promise were to reject, then the actor's
765 // output gate will be broken first and the isolate will not resume synchronous execution.
766 return context.attachSpans(js, context.awaitIo(js, kj::mv(p)), kj::mv(traceContext));
767 } else {
768 return js.resolvedPromise();
769 }
770}
771 
772SqliteDatabase& DurableObjectStorage::getSqliteDb(jsg::Lock& js) {
773 KJ_IF_SOME(db, cache->getSqliteDatabase()) {
774 // Actor is SQLite-backed but let's make sure SQL is configured to be enabled.
775 if (enableSql) {
776 return db;
777 } else if (FeatureFlags::get(js).getWorkerdExperimental()) {
778 // For backwards-compatibility, if the `experimental` compat flag is on, enable SQL. This is
779 // deprecated, though, so warn in this case.
780 
781 // TODO(soon): Uncomment this warning after the D1 simulator has been updated to use
782 // `enableSql`. Otherwise, people doing local dev against D1 may see the warning
783 // spuriously.
784 
785 // IoContext::current().logWarningOnce(
786 // "Enabling SQL API based on the 'experimental' flag, but this will stop working soon. "
787 // "Instead, please set `enableSql = true` in your workerd config for the DO namespace. "
788 // "If using wrangler, under `[[migrations]]` in wrangler.toml, change `new_classes` to "
789 // "`new_sqlite_classes`.");
790 
791 return db;
792 } else {
793 // We're presumably running local workerd, which always uses SQLite for DO storage, but we're
794 // trying to simulate a non-SQLite DO namespace for testing purposes.
795 JSG_FAIL_REQUIRE(Error,
796 "SQL is not enabled for this Durable Object class. To enable it, change "
797 "`new_classes` to `new_sqlite_classes` within the 'migrations' field in "
798 "your wrangler.jsonc or wrangler.toml file. If using workerd directly,"
799 "set `enableSql = true` in your workerd config for the class. Note "
800 "that this change cannot be made after the class is "
801 "already deployed to production.");
802 }
803 } else {
804 // We're in production (not local workerd) and this DO namespace is not backed by SQLite.
805 JSG_FAIL_REQUIRE(Error,
806 "This Durable Object is not backed by SQLite storage, so the SQL API is not available. "
807 "SQL can be enabled on a new Durable Object class by using the `new_sqlite_classes` "
808 "instead of `new_classes` under `[[migrations]]` in your wrangler.toml, but an "
809 "already-deployed class cannot be converted to SQLite (except by deleting the existing "
810 "data).");
811 }
812}
813 
814SqliteKv& DurableObjectStorage::getSqliteKv(jsg::Lock& js) {
815 KJ_IF_SOME(kv, cache->getSqliteKv()) {
816 // Actor is SQLite-backed but let's make sure SQL is configured to be enabled.
817 if (enableSql) {
818 return kv;
819 } else {
820 // We're presumably running local workerd, which always uses SQLite for DO storage, but we're
821 // trying to simulate a non-SQLite DO namespace for testing purposes.
822 JSG_FAIL_REQUIRE(Error,
823 "The storage.kv (synchronous KV) API is only available for SQLite-backed Durable "
824 "Objects, but this object's namespace is not declared to use SQLite. You can use "
825 "the older, asyncronous interface via methods of `storage` itself (e.g. "
826 "`storage.get()`). Alternatively, to enable SQLite, change `new_classes` to "
827 "`new_sqlite_classes` within the 'migrations' field in your wrangler.jsonc or "
828 "wrangler.toml file. If using workerd directly, set `enableSql = true` in your workerd "
829 "config for the class. Note that this change cannot be made after the class is "
830 "already deployed to production.");
831 }
832 } else {
833 // We're in production (not local workerd) and this DO namespace is not backed by SQLite.
834 JSG_FAIL_REQUIRE(Error,
835 "The storage.kv (synchronous KV) API is only available for SQLite-backed Durable "
836 "Objects, but this object's namespace is not declared to use SQLite. You can use "
837 "the older, asyncronous interface via methods of `storage` itself (e.g. "
838 "`storage.get()`). SQLite can be enabled on a new Durable Object class by using the "
839 "`new_sqlite_classes` instead of `new_classes` under `migrations` in your "
840 "wrangler.jsonc or wrangler.toml, but an already-deployed class cannot be converted "
841 "to SQLite (except by deleting the existing data).");
842 }
843}
844 
845jsg::Ref<SqlStorage> DurableObjectStorage::getSql(jsg::Lock& js) {
846 return js.alloc<SqlStorage>(JSG_THIS);
847}
848 
849jsg::Ref<SyncKvStorage> DurableObjectStorage::getKv(jsg::Lock& js) {
850 return js.alloc<SyncKvStorage>(JSG_THIS);
851}
852 
853kj::Promise<kj::String> DurableObjectStorage::getCurrentBookmark() {
854 auto& context = IoContext::current();
855 auto traceContext = context.makeUserTraceSpan("durable_object_storage_getCurrentBookmark"_kjc);
856 
857 return cache->getCurrentBookmark(traceContext.getInternalSpanParent())
858 .attach(kj::mv(traceContext));
859}
860 
861kj::Promise<kj::String> DurableObjectStorage::getBookmarkForTime(kj::Date timestamp) {
862 return cache->getBookmarkForTime(timestamp);
863}
864 
865kj::Promise<kj::String> DurableObjectStorage::onNextSessionRestoreBookmark(kj::String bookmark) {
866 return cache->onNextSessionRestoreBookmark(bookmark);
867}
868 
869kj::Promise<void> DurableObjectStorage::waitForBookmark(kj::String bookmark) {
870 auto& context = IoContext::current();
871 auto traceContext = context.makeUserTraceSpan("durable_object_storage_waitForBookmark"_kjc);
872 
873 return cache->waitForBookmark(bookmark, traceContext.getInternalSpanParent())
874 .attach(kj::mv(traceContext));
875}
876 
877void DurableObjectStorage::ensureReplicas() {
878 if (maybePrimary != kj::none) {
879 KJ_FAIL_ASSERT("Replica Durable Objects cannot call ensureReplicas().");
880 }
881 return cache->ensureReplicas();
882}
883 
884void DurableObjectStorage::disableReplicas() {
885 if (maybePrimary != kj::none) {
886 KJ_FAIL_ASSERT("Replica Durable Objects cannot call disableReplicas().");
887 }
888 return cache->disableReplicas();
889}
890 
891jsg::Optional<jsg::Ref<DurableObject>> DurableObjectStorage::getPrimary(jsg::Lock& js) {
892 // TODO(cleanup): the primary stub should live on DurableObjectState instead of DurableObjectStorage.
893 KJ_IF_SOME(primary, maybePrimary) {
894 return primary.addRef();
895 }
896 return kj::none;
897}
898 
899bool DurableObjectStorage::isReplica() {
900 // TODO(cleanup): the primary stub should live on DurableObjectState instead of DurableObjectStorage.
901 return maybePrimary != kj::none;
902}
903 
904ActorCacheOps& DurableObjectTransaction::getCache(OpName op) {
905 JSG_REQUIRE(!rolledBack, Error, kj::str("Cannot ", op, " on rolled back transaction"));
906 auto& result = *JSG_REQUIRE_NONNULL(cacheTxn, Error,
907 kj::str("Cannot call ", op,
908 " on transaction that has already committed: did you move `txn` outside of the closure?"));
909 return result;
910}
911 
912void DurableObjectTransaction::rollback() {
913 if (rolledBack) return; // allow multiple calls to rollback()
914 getCache(OP_ROLLBACK); // just for the checks
915 KJ_IF_SOME(t, cacheTxn) {
916 auto prom = t->rollback();
917 IoContext::current().addWaitUntil(kj::mv(prom).attach(kj::mv(cacheTxn)));
918 cacheTxn = kj::none;
919 }
920 rolledBack = true;
921}
922 
923kj::Promise<void> DurableObjectTransaction::maybeCommit() {
924 // cacheTxn is null if rollback() was called, in which case we don't want to commit anything.
925 KJ_IF_SOME(t, cacheTxn) {
926 auto maybePromise = t->commit();
927 cacheTxn = kj::none;
928 KJ_IF_SOME(promise, maybePromise) {
929 return kj::mv(promise);
930 }
931 }
932 return kj::READY_NOW;
933}
934 
935void DurableObjectTransaction::maybeRollback() {
936 cacheTxn = kj::none;
937 rolledBack = true;
938}
939 
940namespace {
941 
942// Maximum length of a facet name, in characters.
943constexpr size_t MAX_FACET_NAME_LENGTH = 256;
944 
945// Maximum depth of the facet tree, including the root Durable Object. Root is at depth 0, so
946// the deepest allowed facet is at depth MAX_FACET_TREE_DEPTH - 1.
947constexpr uint MAX_FACET_TREE_DEPTH = 4;
948 
949inline void requireValidFacetName(kj::StringPtr name) {
950 JSG_REQUIRE(name.size() <= MAX_FACET_NAME_LENGTH, TypeError, "Facet name is too long (max ",
951 MAX_FACET_NAME_LENGTH, " characters).");
952}
953 
954} // namespace
955 
956class FacetOutgoingFactory final: public Fetcher::OutgoingFactory {
957 public:
958 FacetOutgoingFactory(Worker::Actor::FacetManager& facetManager,
959 kj::String name,
960 kj::Function<kj::Promise<Worker::Actor::FacetManager::StartInfo>()> getStartInfo)
961 : facetManager(facetManager),
962 name(kj::mv(name)),
963 getStartInfo(kj::mv(getStartInfo)) {}
964 
965 kj::Own<WorkerInterface> newSingleUseClient(kj::Maybe<kj::String> cfStr) override {
966 auto& context = IoContext::current();
967 
968 return context.getMetrics().wrapActorSubrequestClient(context.getSubrequest(
969 [&](TraceContext& tracing, IoChannelFactory& ioChannelFactory) {
970 tracing.setTag("facet_name"_kjc, name.asPtr());
971 
972 // Lazily initialize actorChannel
973 if (actorChannel == kj::none) {
974 actorChannel = facetManager.getFacet(name, kj::mv(getStartInfo));
975 }
976 
977 return KJ_REQUIRE_NONNULL(actorChannel)
978 ->startRequest({.cfBlobJson = kj::mv(cfStr),
979 .parentSpan = tracing.getInternalSpanParent(),
980 .userSpanParent = tracing.getUserSpanParent()});
981 },
982 {.inHouse = true,
983 .wrapMetrics = true,
984 .operationName = kj::ConstString("facet_subrequest"_kjc)}));
985 }
986 
987 private:
988 Worker::Actor::FacetManager& facetManager;
989 kj::String name;
990 
991 // This is moved away when `actorChannel` is initialized.
992 kj::Function<kj::Promise<Worker::Actor::FacetManager::StartInfo>()> getStartInfo;
993 
994 kj::Maybe<kj::Own<IoChannelFactory::ActorChannel>> actorChannel;
995};
996 
997jsg::Ref<Fetcher> DurableObjectFacets::get(jsg::Lock& js,
998 kj::String name,
999 jsg::Function<jsg::Promise<StartupOptions>()> getStartupOptions) {
1000 requireValidFacetName(name);
1001 
1002 auto& fm = getFacetManager();
1003 
1004 JSG_REQUIRE(fm.getDepth() + 1 < MAX_FACET_TREE_DEPTH, Error,
1005 "Facet nesting depth limit exceeded. The maximum depth including the root Durable Object is ",
1006 MAX_FACET_TREE_DEPTH, ".");
1007 
1008 auto& ioCtx = IoContext::current();
1009 
1010 kj::Function<kj::Promise<Worker::Actor::FacetManager::StartInfo>()> getStartInfo =
1011 ioCtx.makeReentryCallback(
1012 [&ioCtx, getStartupOptions = kj::mv(getStartupOptions)](jsg::Lock& js) mutable {
1013 return getStartupOptions(js).then(js, [&ioCtx](jsg::Lock& js, StartupOptions options) {
1014 Worker::Actor::Id id;
1015 KJ_IF_SOME(i, options.id) {
1016 KJ_SWITCH_ONEOF(i) {
1017 KJ_CASE_ONEOF(doId, jsg::Ref<DurableObjectId>) {
1018 id = doId->getInner().clone();
1019 }
1020 KJ_CASE_ONEOF(strId, kj::String) {
1021 id = kj::mv(strId);
1022 }
1023 }
1024 } else {
1025 // Child inherits parent ID.
1026 id = ioCtx.getActorOrThrow().cloneId();
1027 }
1028 
1029 DurableObjectClass& actorClass = [&]() -> DurableObjectClass& {
1030 KJ_SWITCH_ONEOF(options.$class) {
1031 KJ_CASE_ONEOF(bare, jsg::Ref<DurableObjectClass>) {
1032 return *bare.get();
1033 }
1034 KJ_CASE_ONEOF(loopback, jsg::Ref<LoopbackDurableObjectNamespace>) {
1035 return loopback->getClass();
1036 }
1037 KJ_CASE_ONEOF(loopback, jsg::Ref<LoopbackColoLocalActorNamespace>) {
1038 return loopback->getClass();
1039 }
1040 }
1041 KJ_UNREACHABLE;
1042 }();
1043 
1044 return Worker::Actor::FacetManager::StartInfo{
1045 .actorClass = actorClass.getChannel(ioCtx),
1046 .id = kj::mv(id),
1047 };
1048 });
1049 });
1050 
1051 kj::Own<Fetcher::OutgoingFactory> factory =
1052 kj::heap<FacetOutgoingFactory>(fm, kj::mv(name), kj::mv(getStartInfo));
1053 
1054 auto requiresHost = FeatureFlags::get(js).getDurableObjectFetchRequiresSchemeAuthority()
1055 ? Fetcher::RequiresHostAndProtocol::YES
1056 : Fetcher::RequiresHostAndProtocol::NO;
1057 
1058 // We return a plain Fetcher, not a DurableObject, because we don't want the stub to have
1059 // `name` or `id` properties.
1060 return js.alloc<Fetcher>(ioCtx.addObject(kj::mv(factory)), requiresHost, true /* isInHouse */);
1061}
1062 
1063void DurableObjectFacets::abort(jsg::Lock& js, kj::String name, jsg::JsValue reason) {
1064 requireValidFacetName(name);
1065 getFacetManager().abortFacet(name, js.exceptionToKj(reason));
1066}
1067 
1068void DurableObjectFacets::delete_(jsg::Lock& js, kj::String name) {
1069 requireValidFacetName(name);
1070 getFacetManager().deleteFacet(name);
1071}
1072 
1073ActorState::ActorState(Worker::Actor::Id actorId,
1074 kj::Maybe<jsg::JsRef<jsg::JsValue>> transient,
1075 kj::Maybe<jsg::Ref<DurableObjectStorage>> persistent)
1076 : id(kj::mv(actorId)),
1077 transient(kj::mv(transient)),
1078 persistent(kj::mv(persistent)) {}
1079 
1080kj::OneOf<jsg::Ref<DurableObjectId>, kj::StringPtr> ActorState::getId(jsg::Lock& js) {
1081 KJ_SWITCH_ONEOF(id) {
1082 KJ_CASE_ONEOF(coloLocalId, kj::String) {
1083 return coloLocalId.asPtr();
1084 }
1085 KJ_CASE_ONEOF(globalId, kj::Own<ActorIdFactory::ActorId>) {
1086 return js.alloc<DurableObjectId>(globalId->clone());
1087 }
1088 }
1089 KJ_UNREACHABLE;
1090}
1091 
1092DurableObjectState::DurableObjectState(jsg::Lock& js,
1093 Worker::Actor::Id actorId,
1094 jsg::JsValue exports,
1095 jsg::JsValue props,
1096 kj::Maybe<jsg::Ref<DurableObjectStorage>> storage,
1097 kj::Maybe<rpc::Container::Client> container,
1098 bool containerRunning,
1099 kj::Maybe<Worker::Actor::FacetManager&> facetManager,
1100 kj::Maybe<ActorVersion> version)
1101 : id(kj::mv(actorId)),
1102 exports(js, exports),
1103 props(js, props),
1104 storage(kj::mv(storage)),
1105 container(container.map([&](rpc::Container::Client& cap) {
1106 return js.alloc<Container>(kj::mv(cap), containerRunning);
1107 })),
1108 facetManager(facetManager.map(
1109 [](Worker::Actor::FacetManager& ref) { return IoContext::current().addObject(ref); })),
1110 version(kj::mv(version)) {}
1111 
1112void DurableObjectState::waitUntil(kj::Promise<void> promise) {
1113 IoContext::current().addWaitUntil(kj::mv(promise));
1114}
1115 
1116kj::OneOf<jsg::Ref<DurableObjectId>, kj::StringPtr> DurableObjectState::getId(jsg::Lock& js) {
1117 KJ_SWITCH_ONEOF(id) {
1118 KJ_CASE_ONEOF(coloLocalId, kj::String) {
1119 return coloLocalId.asPtr();
1120 }
1121 KJ_CASE_ONEOF(globalId, kj::Own<ActorIdFactory::ActorId>) {
1122 return js.alloc<DurableObjectId>(globalId->clone());
1123 }
1124 }
1125 KJ_UNREACHABLE;
1126}
1127 
1128jsg::Promise<jsg::JsRef<jsg::JsValue>> DurableObjectState::blockConcurrencyWhile(
1129 jsg::Lock& js, jsg::Function<jsg::Promise<jsg::JsRef<jsg::JsValue>>()> callback) {
1130 return IoContext::current().blockConcurrencyWhile(js, kj::mv(callback));
1131}
1132 
1133void DurableObjectState::abort(jsg::Lock& js, jsg::Optional<kj::String> reason) {
1134 kj::String description = kj::mv(reason)
1135 .map([](kj::String&& text) {
1136 return kj::str("broken.outputGateBroken; jsg.Error: ", text);
1137 }).orDefault([]() {
1138 return kj::str("broken.outputGateBroken; jsg.Error: Application called abort() to reset "
1139 "Durable Object.");
1140 });
1141 
1142 kj::Exception error(kj::Exception::Type::FAILED, __FILE__, __LINE__, kj::mv(description));
1143 error.setDetail(jsg::EXCEPTION_IS_USER_ERROR, kj::heapArray<byte>(0));
1144 
1145 KJ_IF_SOME(s, storage) {
1146 // Make sure we _synchronously_ break storage so that there's no chance our promise fulfilling
1147 // will race against the output gate, possibly allowing writes to complete before being
1148 // canceled.
1149 s.get()->getActorCacheInterface().shutdown(error);
1150 }
1151 
1152 IoContext::current().abort(kj::mv(error));
1153 js.terminateExecutionNow();
1154}
1155 
1156Worker::Actor::HibernationManager& DurableObjectState::maybeInitHibernationManager(
1157 Worker::Actor& actor) {
1158 if (actor.getHibernationManager() == kj::none) {
1159 // If there's no hibernation manager created yet, we should create one.
1160 actor.setHibernationManager(kj::refcounted<HibernationManagerImpl>(
1161 actor.getLoopback(), KJ_REQUIRE_NONNULL(actor.getHibernationEventType())));
1162 }
1163 return KJ_REQUIRE_NONNULL(actor.getHibernationManager());
1164}
1165 
1166void DurableObjectState::acceptWebSocket(
1167 jsg::Ref<WebSocket> ws, jsg::Optional<kj::Array<kj::String>> tags) {
1168 JSG_ASSERT(!ws->isAccepted(), Error,
1169 "Cannot call `acceptWebSocket()` if the WebSocket was already accepted via `accept()`");
1170 JSG_ASSERT(ws->peerIsAwaitingCoupling(), Error,
1171 "Cannot call `acceptWebSocket()` on this WebSocket because its pair has already been "
1172 "accepted or used in a Response.");
1173 
1174 // We need to get a HibernationManager to give the websocket to.
1175 auto& a = KJ_REQUIRE_NONNULL(IoContext::current().getActor());
1176 // HibernationManager's acceptWebSocket() will throw if the websocket is in an incompatible state.
1177 // Note that not providing a tag is equivalent to providing an empty tag array.
1178 // Any duplicate tags will be ignored.
1179 kj::Array<kj::String> distinctTags = [&]() -> kj::Array<kj::String> {
1180 KJ_IF_SOME(t, tags) {
1181 kj::HashSet<kj::String> seen;
1182 size_t distinctTagCount = 0;
1183 for (auto tag = t.begin(); tag < t.end(); tag++) {
1184 JSG_REQUIRE(distinctTagCount < MAX_TAGS_PER_CONNECTION, Error,
1185 "a Hibernatable WebSocket cannot have more than ", MAX_TAGS_PER_CONNECTION, " tags");
1186 JSG_REQUIRE(tag->size() <= MAX_TAG_LENGTH, Error, "\"", *tag, "\" ",
1187 "is longer than the max tag length (", MAX_TAG_LENGTH, " characters).");
1188 if (!seen.contains(*tag)) {
1189 seen.insert(kj::mv(*tag));
1190 distinctTagCount++;
1191 }
1192 }
1193 
1194 return KJ_MAP(tag, seen) { return kj::mv(tag); };
1195 }
1196 return kj::Array<kj::String>();
1197 }();
1198 maybeInitHibernationManager(a).acceptWebSocket(kj::mv(ws), distinctTags);
1199}
1200 
1201kj::Array<jsg::Ref<api::WebSocket>> DurableObjectState::getWebSockets(
1202 jsg::Lock& js, jsg::Optional<kj::String> tag) {
1203 auto& a = KJ_REQUIRE_NONNULL(IoContext::current().getActor());
1204 KJ_IF_SOME(manager, a.getHibernationManager()) {
1205 return manager.getWebSockets(js, tag.map([](kj::StringPtr t) { return t; })).releaseAsArray();
1206 }
1207 return kj::Array<jsg::Ref<api::WebSocket>>();
1208}
1209 
1210void DurableObjectState::setWebSocketAutoResponse(
1211 jsg::Optional<jsg::Ref<WebSocketRequestResponsePair>> maybeReqResp) {
1212 auto& a = KJ_REQUIRE_NONNULL(IoContext::current().getActor());
1213 
1214 if (maybeReqResp == kj::none) {
1215 // If there's no request/response pair, we unset any current set auto response configuration.
1216 KJ_IF_SOME(manager, a.getHibernationManager()) {
1217 // If there's no hibernation manager created yet, there's nothing to do here.
1218 manager.setWebSocketAutoResponse(kj::none, kj::none);
1219 }
1220 return;
1221 }
1222 
1223 auto reqResp = KJ_REQUIRE_NONNULL(kj::mv(maybeReqResp));
1224 auto maxRequestOrResponseSize = 2048;
1225 
1226 JSG_REQUIRE(reqResp->getRequest().size() <= maxRequestOrResponseSize, RangeError,
1227 kj::str("Request cannot be larger than ", maxRequestOrResponseSize, " bytes. ",
1228 "A request of size ", reqResp->getRequest().size(), " was provided."));
1229 
1230 JSG_REQUIRE(reqResp->getResponse().size() <= maxRequestOrResponseSize, RangeError,
1231 kj::str("Response cannot be larger than ", maxRequestOrResponseSize, " bytes. ",
1232 "A response of size ", reqResp->getResponse().size(), " was provided."));
1233 
1234 maybeInitHibernationManager(a).setWebSocketAutoResponse(
1235 reqResp->getRequest(), reqResp->getResponse());
1236}
1237 
1238kj::Maybe<jsg::Ref<api::WebSocketRequestResponsePair>> DurableObjectState::getWebSocketAutoResponse(
1239 jsg::Lock& js) {
1240 auto& a = KJ_REQUIRE_NONNULL(IoContext::current().getActor());
1241 KJ_IF_SOME(manager, a.getHibernationManager()) {
1242 // If there's no hibernation manager created yet, there's nothing to do here.
1243 return manager.getWebSocketAutoResponse(js);
1244 }
1245 return kj::none;
1246}
1247 
1248kj::Maybe<kj::Date> DurableObjectState::getWebSocketAutoResponseTimestamp(jsg::Ref<WebSocket> ws) {
1249 return ws->getAutoResponseTimestamp();
1250}
1251 
1252void DurableObjectState::setHibernatableWebSocketEventTimeout(jsg::Optional<uint32_t> timeoutMs) {
1253 auto& a = KJ_REQUIRE_NONNULL(IoContext::current().getActor());
1254 
1255 // Setting a timeout = 0ms or an empty value will unset any currently set event timeout.
1256 // If there's no hibernation manager instantiated, we can skip the event timeout unsetting.
1257 if (timeoutMs == kj::none || KJ_REQUIRE_NONNULL(timeoutMs) == 0) {
1258 KJ_IF_SOME(hibernationManager, a.getHibernationManager()) {
1259 hibernationManager.setEventTimeout(kj::none);
1260 }
1261 return;
1262 }
1263 
1264 auto t = timeoutMs.orDefault(static_cast<uint32_t>(0));
1265 
1266 // We want to limit the duration of an event to a maximum of 7 days (604800 * 1000 millis).
1267 JSG_REQUIRE(t <= 604800 * 1000, Error, "Event timeout should not exceed 604800000 ms.");
1268 
1269 maybeInitHibernationManager(a).setEventTimeout(t);
1270}
1271 
1272kj::Maybe<uint32_t> DurableObjectState::getHibernatableWebSocketEventTimeout() {
1273 KJ_IF_SOME(a, IoContext::current().getActor()) {
1274 KJ_IF_SOME(manager, a.getHibernationManager()) {
1275 return manager.getEventTimeout();
1276 }
1277 }
1278 return kj::none;
1279}
1280 
1281kj::Array<kj::StringPtr> DurableObjectState::getTags(jsg::Lock& js, jsg::Ref<api::WebSocket> ws) {
1282 return ws->getHibernatableTags();
1283}
1284 
1285jsg::Optional<jsg::Ref<DurableObject>> DurableObjectState::getPrimaryStub(jsg::Lock& js) {
1286 KJ_IF_SOME(s, storage) {
1287 return s->getPrimary(js);
1288 }
1289 return kj::none;
1290}
1291 
1292jsg::Promise<void> DurableObjectState::configureReadReplication(
1293 jsg::Lock& js, DurableObjectState::ReadReplicationOptions options) {
1294 
1295 auto& context = IoContext::current();
1296 auto traceContext =
1297 context.makeUserTraceSpan("durable_object_state_configureReadReplication"_kjc);
1298 
1299 auto& s =
1300 JSG_REQUIRE_NONNULL(storage, TypeError, "This actor does not support read replication.");
1301 
1302 if (s->isReplica()) {
1303 JSG_FAIL_REQUIRE(Error, "Replica Durable Objects cannot call configureReadReplication().");
1304 }
1305 
1306 bool enabled = [&]() {
1307 if (options.mode == "auto"_kj) {
1308 return true;
1309 } else if (options.mode == "disabled"_kj) {
1310 return false;
1311 }
1312 JSG_FAIL_REQUIRE(TypeError,
1313 "configureReadReplication() called with unknown mode setting: ", options.mode, ".");
1314 }();
1315 
1316 auto promise =
1317 s->getActorCacheInterface().configureReadReplication(ReadReplicationIsEnabled(enabled));
1318 
1319 return context.attachSpans(js, context.awaitIo(js, kj::mv(promise)), kj::mv(traceContext));
1320}
1321 
1322kj::Array<kj::byte> serializeV8Value(jsg::Lock& js, const jsg::JsValue& value) {
1323 jsg::Serializer serializer(js,
1324 jsg::Serializer::Options{
1325 .version = 15,
1326 .omitHeader = false,
1327 });
1328 serializer.write(js, value);
1329 auto released = serializer.release();
1330 return kj::mv(released.data);
1331}
1332 
1333jsg::JsValue deserializeV8Value(
1334 jsg::Lock& js, kj::ArrayPtr<const char> key, kj::ArrayPtr<const kj::byte> buf) {
1335 
1336 KJ_ASSERT(buf.size() > 0, "unexpectedly empty value buffer", key);
1337 try {
1338 // The js.tryCatch will handle the normal exception path. We wrap this in an
1339 // additional try/catch in case the js.tryCatch hits an exception that is
1340 // terminal for the isolate, causing exception to be rethrown, in which case
1341 // we throw a kj::Exception wrapping a jsg.Error.
1342 return js.tryCatch([&]() -> jsg::JsValue {
1343 jsg::Deserializer::Options options{};
1344 if (buf[0] != 0xFF) {
1345 // When Durable Objects was first released, it did not properly write headers when serializing
1346 // to storage. If we find that the header is missing (as indicated by the first byte not being
1347 // 0xFF), it's safe to assume that the data was written at the only serialization version we
1348 // used during that early time period, so we explicitly set that version here.
1349 options.version = 13;
1350 options.readHeader = false;
1351 }
1352 
1353 jsg::Deserializer deserializer(js, buf, kj::none, kj::none, options);
1354 
1355 return deserializer.readValue(js);
1356 }, [&](jsg::Value&& exception) mutable -> jsg::JsValue {
1357 // If we do hit a deserialization error, we log information that will be helpful in
1358 // understanding the problem but that won't leak too much about the customer's data. We
1359 // include the key (to help find the data in the database if it hasn't been deleted), the
1360 // length of the value, and the first three bytes of the value (which is just the v8-internal
1361 // version header and the tag that indicates the type of the value, but not its contents).
1362 kj::String actorId = getCurrentActorId().orDefault([]() { return kj::String(); });
1363 KJ_FAIL_ASSERT("actor storage deserialization failed", "failed to deserialize stored value",
1364 actorId, exception.getHandle(js), key, buf.size(),
1365 buf.first(std::min(static_cast<size_t>(3), buf.size())));
1366 });
1367 } catch (jsg::JsExceptionThrown&) {
1368 // We can occasionally hit an isolate termination here -- we prefix the error with jsg to avoid
1369 // counting it against our internal storage error metrics but also throw a KJ exception rather
1370 // than a jsExceptionThrown error to avoid confusing the normal termination handling code.
1371 // We don't expect users to ever actually see this error.
1372 JSG_FAIL_REQUIRE(Error,
1373 "isolate terminated while deserializing value from Durable Object "
1374 "storage; contact us if you're wondering why you're seeing this");
1375 }
1376}
1377 
1378} // namespace workerd::api