Skip to content
File

Blob: src/workerd/io/worker.c++

195.0 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-cache.h"
6 
7#include <workerd/api/actor-state.h>
8#include <workerd/api/global-scope.h>
9#include <workerd/api/sockets.h>
10#include <workerd/api/streams/common.h> // for api::StreamEncoding
11#include <workerd/io/cdp.capnp.h>
12#include <workerd/io/compatibility-date.h>
13#include <workerd/io/features.h>
14#include <workerd/io/frankenvalue.h>
15#include <workerd/io/tracer.h>
16#include <workerd/io/wasm-instantiate-shim.embed.h>
17#include <workerd/io/worker.h>
18#include <workerd/jsg/async-context.h>
19#include <workerd/jsg/inspector.h>
20#include <workerd/jsg/jsg.h>
21#include <workerd/jsg/modules-new.h>
22#include <workerd/jsg/script.h>
23#include <workerd/jsg/setup.h>
24#include <workerd/jsg/util.h>
25#include <workerd/rust/jsg/lib.rs.h>
26#include <workerd/rust/jsg/v8.rs.h>
27#include <workerd/util/autogate.h>
28#include <workerd/util/batch-queue.h>
29#include <workerd/util/color-util.h>
30#include <workerd/util/mimetype.h>
31#include <workerd/util/stream-utils.h>
32#include <workerd/util/thread-scopes.h>
33#include <workerd/util/uuid.h>
34#include <workerd/util/xthreadnotifier.h>
35 
36#include <rust/jsg/ffi.h>
37#include <v8-inspector.h>
38#include <v8-profiler.h>
39#include <v8-wasm.h>
40 
41#include <capnp/compat/json.h>
42#include <capnp/message.h>
43#include <kj/compat/brotli.h>
44#include <kj/compat/gzip.h>
45#include <kj/encoding.h>
46#include <kj/filesystem.h>
47#include <kj/map.h>
48 
49#include <cstdint>
50#include <ctime>
51#include <map>
52#include <numeric>
53 
54#if _WIN32
55#include <io.h>
56#include <windows.h>
57 
58#include <kj/win32-api-version.h>
59#include <kj/windows-sanity.h>
60#else
61#include <sys/syscall.h>
62#include <unistd.h>
63#endif
64 
65namespace workerd {
66 
67namespace {
68 
69constexpr kj::StringPtr logLevelToString(LogLevel level) {
70 switch (level) {
71 case LogLevel::DEBUG_:
72 return "debug";
73 case LogLevel::INFO:
74 return "info";
75 case LogLevel::LOG:
76 return "log";
77 case LogLevel::WARN:
78 return "warn";
79 case LogLevel::ERROR:
80 return "error";
81 default:
82 return "log";
83 }
84}
85 
86void headersToCDP(const kj::HttpHeaders& in, capnp::JsonValue::Builder out) {
87 std::map<kj::StringPtr, kj::Vector<kj::StringPtr>> inMap;
88 in.forEach([&](kj::StringPtr name, kj::StringPtr value) {
89 inMap.try_emplace(name, 1).first->second.add(value);
90 });
91 
92 auto outObj = out.initObject(inMap.size());
93 auto headersPos = 0;
94 for (auto& entry: inMap) {
95 auto field = outObj[headersPos++];
96 field.setName(entry.first);
97 
98 // CDP uses strange header representation where headers with multiple
99 // values are merged into one newline-delimited string
100 field.initValue().setString(kj::strArray(entry.second, "\n"));
101 }
102}
103 
104void stackTraceToCDP(jsg::Lock& js, cdp::Runtime::StackTrace::Builder builder) {
105 // TODO(cleanup): Maybe use V8Inspector::captureStackTrace() which does this for us. However, it
106 // produces protocol objects in its own format which want to handle their whole serialization
107 // to JSON. Also, those protocol objects are defined in generated code which we currently don't
108 // include in our cached V8 build artifacts; we'd need to fix that. But maybe we should really
109 // be using the V8-generated protocol objects rather than our parallel capnp versions!
110 
111 auto stackTrace = v8::StackTrace::CurrentStackTrace(js.v8Isolate, 10);
112 auto frameCount = stackTrace->GetFrameCount();
113 auto callFrames = builder.initCallFrames(frameCount);
114 for (int i = 0; i < frameCount; i++) {
115 auto src = stackTrace->GetFrame(js.v8Isolate, i);
116 auto dest = callFrames[i];
117 auto url = src->GetScriptNameOrSourceURL();
118 if (!url.IsEmpty()) {
119 dest.setUrl(kj::str(url));
120 } else {
121 dest.setUrl(""_kj);
122 }
123 dest.setScriptId(kj::str(src->GetScriptId()));
124 auto func = src->GetFunctionName();
125 if (!func.IsEmpty()) {
126 dest.setFunctionName(kj::str(func));
127 } else {
128 dest.setFunctionName(""_kj);
129 }
130 // V8 locations are 1-based, but CDP locations are 0-based... oh, well
131 dest.setLineNumber(src->GetLineNumber() - 1);
132 dest.setColumnNumber(src->GetColumn() - 1);
133 }
134}
135 
136kj::Own<capnp::JsonCodec> makeCdpJsonCodec() {
137 auto codec = kj::heap<capnp::JsonCodec>();
138 codec->handleByAnnotation<cdp::Command>();
139 codec->handleByAnnotation<cdp::Event>();
140 return codec;
141}
142const capnp::JsonCodec& getCdpJsonCodec() {
143 static const kj::Own<capnp::JsonCodec> codec = makeCdpJsonCodec();
144 return *codec;
145}
146 
147} // namespace
148 
149// =======================================================================================
150 
151namespace {
152 
153using ExceptionOrDuration = kj::OneOf<kj::Exception, kj::Duration>;
154 
155// Inform the inspector of an exception thrown.
156//
157// Passes `source` as the exception's short message. Reconstructs `message` from `exception` if
158// `message` is empty.
159void sendExceptionToInspector(jsg::Lock& js,
160 v8_inspector::V8Inspector& inspector,
161 UncaughtExceptionSource source,
162 const jsg::JsValue& exception,
163 jsg::JsMessage message) {
164 jsg::sendExceptionToInspector(js, inspector, kj::str(source), exception, message);
165}
166 
167void addExceptionToTrace(jsg::Lock& js,
168 IoContext& ioContext,
169 BaseTracer& tracer,
170 UncaughtExceptionSource source,
171 const jsg::JsValue& exception,
172 const jsg::TypeHandler<Worker::Api::ErrorInterface>& errorTypeHandler) {
173 if (source == UncaughtExceptionSource::INTERNAL ||
174 source == UncaughtExceptionSource::INTERNAL_ASYNC) {
175 // Skip redundant intermediate JS->C++ exception reporting. See: IoContext::runImpl(),
176 // PromiseWrapper::tryUnwrap()
177 //
178 // TODO(someday): Arguably it could make sense to store these exceptions off to the side and
179 // report them only if they don't end up being duplicates of a later exception that has a more
180 // specific context. This would cover cases where the C++ code that eventually received the
181 // exception never ended up reporting it.
182 return;
183 }
184 
185 auto timestamp = ioContext.now();
186 Worker::Api::ErrorInterface error;
187 
188 if (exception.isObject()) {
189 error = KJ_REQUIRE_NONNULL(errorTypeHandler.tryUnwrap(js, exception),
190 "Should always be possible to unwrap error interface from an object.");
191 }
192 
193 kj::String name;
194 KJ_IF_SOME(n, error.name) {
195 name = kj::str(n);
196 } else {
197 name = kj::str("Error");
198 }
199 kj::String message;
200 KJ_IF_SOME(m, error.message) {
201 message = kj::str(m);
202 } else {
203 // This doesn't appear to be an Error object. Fall back to stringifying the whole value as
204 // the message.
205 if (!js.v8Isolate->IsExecutionTerminating()) {
206 v8::TryCatch tryCatch(js.v8Isolate);
207 try {
208 message = exception.toString(js);
209 } catch (jsg::JsExceptionThrown&) {
210 // Failed to stringify.
211 //
212 // Note that we're intentionally not checking tryCatch.CanContinue() here, because we still
213 // want to continue even if the isolate has been terminated.
214 }
215 }
216 }
217 
218 kj::Maybe<kj::String> stack;
219 KJ_IF_SOME(s, error.stack) {
220 kj::StringPtr slice = s;
221 
222 // Normally `error.stack` repeats the error type and message first. We don't want send two
223 // copies of that to the trace so we'll strip it off.
224 if (slice.startsWith(name)) {
225 slice = slice.slice(name.size());
226 if (slice.startsWith(": "_kj)) {
227 slice = slice.slice(2);
228 }
229 }
230 
231 if (slice.startsWith(message)) {
232 slice = slice.slice(message.size());
233 if (slice.startsWith("\n")) {
234 slice = slice.slice(1);
235 }
236 }
237 
238 if (slice.size() > 0) {
239 stack = kj::str(slice);
240 }
241 }
242 
243 tracer.addException(ioContext.getInvocationSpanContext(), timestamp, kj::mv(name),
244 kj::mv(message), kj::mv(stack));
245}
246 
247void reportStartupError(kj::StringPtr id,
248 jsg::Lock& js,
249 const kj::Maybe<std::unique_ptr<v8_inspector::V8Inspector>>& inspector,
250 const IsolateLimitEnforcer& limitEnforcer,
251 ExceptionOrDuration limitErrorOrTime,
252 v8::TryCatch& catcher,
253 kj::Maybe<Worker::ValidationErrorReporter&> errorReporter,
254 kj::Maybe<kj::Exception>& permanentException,
255 SpanParent parentSpan,
256 bool isDynamicWorker) {
257 v8::TryCatch catcher2(js.v8Isolate);
258 ExceptionOrDuration limitErrorOrTime2 = 0 * kj::NANOSECONDS;
259 try {
260 KJ_SWITCH_ONEOF(limitErrorOrTime) {
261 KJ_CASE_ONEOF(limitError, kj::Exception) {
262 auto description = jsg::extractTunneledExceptionDescription(limitError.getDescription());
263 
264 auto& ex = permanentException.emplace(kj::mv(limitError));
265 KJ_IF_SOME(e, errorReporter) {
266 e.addError(kj::heapString(description));
267 } else KJ_IF_SOME(i, inspector) {
268 // We want to extend just enough CPU time as is necessary to report the exception
269 // to the inspector here. 10 milliseconds should be more than enough.
270 auto limitScope = limitEnforcer.enterLoggingJs(js, limitErrorOrTime2);
271 jsg::sendExceptionToInspector(js, *i.get(), description);
272 // When the inspector is active, we don't want to throw here because then the inspector
273 // won't be able to connect and the developer will never know what happened.
274 } else {
275 // We should never get here in production if we've validated scripts before deployment.
276 KJ_LOG(WARNING, "script startup exceeded resource limits", id, ex);
277 kj::throwFatalException(ex.clone());
278 }
279 }
280 KJ_CASE_ONEOF_DEFAULT {
281 if (catcher.HasCaught()) {
282 js.withinHandleScope([&] {
283 auto exception = catcher.Exception();
284 
285 permanentException = js.exceptionToKj(js.v8Ref(exception));
286 
287 KJ_IF_SOME(e, errorReporter) {
288 auto limitScope = limitEnforcer.enterLoggingJs(js, limitErrorOrTime2);
289 
290 kj::Vector<kj::String> lines;
291 lines.add(kj::str("Uncaught ",
292 jsg::extractTunneledExceptionDescription(
293 KJ_ASSERT_NONNULL(permanentException).getDescription())));
294 jsg::JsMessage message(catcher.Message());
295 message.addJsStackTrace(js, lines);
296 e.addError(kj::strArray(lines, "\n"));
297 
298 } else KJ_IF_SOME(i, inspector) {
299 auto limitScope = limitEnforcer.enterLoggingJs(js, limitErrorOrTime2);
300 sendExceptionToInspector(js, *i.get(), UncaughtExceptionSource::INTERNAL,
301 jsg::JsValue(exception), jsg::JsMessage(catcher.Message()));
302 // When the inspector is active, we don't want to throw here because then the inspector
303 // won't be able to connect and the developer will never know what happened.
304 } else {
305 // We should never get here in production if we've validated scripts before deployment.
306 // (unless this is a dynamic worker)
307 kj::Vector<kj::String> lines;
308 jsg::JsMessage message(catcher.Message());
309 message.addJsStackTrace(js, lines);
310 auto trace = kj::strArray(lines, "; ");
311 auto description = KJ_ASSERT_NONNULL(permanentException).getDescription();
312 auto span = parentSpan.newChild("script_startup_exception"_kjc);
313 span.setTag("error"_kjc, true);
314 span.addLog(kj::systemPreciseCalendarClock().now(), "exception"_kjc,
315 kj::ConstString(
316 kj::str("script startup threw exception", id, description, trace)));
317 if (isDynamicWorker) {
318 // Rethrow the tunneled JSG exception so it converts back to a JS Error.
319 kj::throwFatalException(KJ_ASSERT_NONNULL(permanentException).clone());
320 } else {
321 KJ_LOG(ERROR, "script startup threw exception", id, description, trace);
322 KJ_FAIL_REQUIRE("script startup threw exception");
323 }
324 }
325 });
326 } else {
327 kj::throwFatalException(permanentException
328 .emplace(KJ_EXCEPTION(FAILED,
329 "returned empty handle but didn't throw exception?", id))
330 .clone());
331 }
332 }
333 }
334 } catch (const jsg::JsExceptionThrown&) {
335#define LOG_AND_SET_PERM_EXCEPTION(...) \
336 KJ_LOG(ERROR, __VA_ARGS__); \
337 if (permanentException == kj::none) { \
338 permanentException = KJ_EXCEPTION(FAILED, __VA_ARGS__); \
339 }
340 
341 KJ_SWITCH_ONEOF(limitErrorOrTime2) {
342 KJ_CASE_ONEOF(limitError2, kj::Exception) {
343 // TODO(cleanup): If we see this error show up in production, stop logging it, because I
344 // guess it's not necessarily an error? The other two cases below are more worrying though.
345 KJ_LOG(ERROR, limitError2);
346 if (permanentException == kj::none) {
347 permanentException = kj::mv(limitError2);
348 }
349 }
350 KJ_CASE_ONEOF_DEFAULT {
351 if (catcher2.HasTerminated()) {
352 LOG_AND_SET_PERM_EXCEPTION(
353 "script startup threw exception; during our attempt to stringify the exception, "
354 "the script apparently was terminated for non-resource-limit reasons.",
355 id);
356 } else {
357 LOG_AND_SET_PERM_EXCEPTION(
358 "script startup threw exception; furthermore, an attempt to stringify the exception "
359 "threw another exception, which shouldn't be possible?",
360 id);
361 }
362 }
363 }
364#undef LOG_AND_SET_PERM_EXCEPTION
365 }
366}
367 
368uint64_t getCurrentThreadId() {
369#if __linux__
370 return syscall(SYS_gettid);
371#elif _WIN32
372 return GetCurrentThreadId();
373#else
374 // Assume MacOS or BSD
375 uint64_t tid;
376 pthread_threadid_np(nullptr, &tid);
377 return tid;
378#endif
379}
380 
381} // namespace
382 
383// Represents a thread's attempt to take an async lock. Each Isolate has a linked list of
384// `AsyncWaiter`s. A particular thread only ever owns one `AsyncWaiter` at a time.
385class Worker::AsyncWaiter: public kj::Refcounted {
386 public:
387 AsyncWaiter(kj::Own<const Isolate> isolate);
388 ~AsyncWaiter() noexcept;
389 KJ_DISALLOW_COPY_AND_MOVE(AsyncWaiter);
390 
391 private:
392 // Executor for this waiter's thread.
393 const kj::Executor& executor;
394 
395 // The isolate for which this waiter is currently waiting.
396 kj::Own<const Isolate> isolate;
397 
398 // Promise/fulfiller to fire when the waiter reaches the front of the list for the corresponding
399 // isolate.
400 kj::ForkedPromise<void> readyPromise = nullptr;
401 kj::Own<kj::CrossThreadPromiseFulfiller<void>> readyFulfiller;
402 
403 // Promise/fulfiller to fire when the AsyncLock is finally released. This is used when a thread
404 // tries to take locks on multiple different isolates concurrently, in order to serialize the
405 // locks so only one is taken at a time. This is NOT a cross-thread fulfiller; it can only be
406 // fulfilled by the thread that owns the waiter.
407 kj::ForkedPromise<void> releasePromise = nullptr;
408 kj::Own<kj::PromiseFulfiller<void>> releaseFulfiller;
409 
410 // Protected by the lock on `Isolate::asyncWaiters` for the isolate identified by
411 // `currentIsolate`. Must be null if `currentIsolate` is null. (All other members of `Waiter`
412 // can only be accessed by the thread that created the `Waiter`.)
413 kj::Maybe<AsyncWaiter&> next;
414 kj::Maybe<AsyncWaiter&>* prev;
415 
416 static const kj::EventLoopLocal<AsyncWaiter*> threadCurrentWaiter;
417 
418 friend class Worker::Isolate;
419 friend class Worker::AsyncLock;
420};
421 
422class Worker::InspectorClient: public v8_inspector::V8InspectorClient {
423 public:
424 // Wall time in milliseconds with millisecond precision. console.time() and friends rely on this
425 // function to implement timers.
426 double currentTimeMS() override {
427 auto timePoint = kj::UNIX_EPOCH;
428 
429 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
430 // We're on a request-serving thread.
431 timePoint = ioContext.now();
432 } else {
433 auto lockedState = state.lockExclusive();
434 KJ_IF_SOME(info, lockedState->inspectorTimerInfo) {
435 if (info.threadId == getCurrentThreadId()) {
436 // We're on an inspector-serving thread.
437 timePoint =
438 info.timer.now() + info.timerOffset - kj::origin<kj::TimePoint>() + kj::UNIX_EPOCH;
439 }
440 }
441 // We're at script startup time -- just return the Epoch.
442 }
443 return (timePoint - kj::UNIX_EPOCH) / kj::MILLISECONDS;
444 }
445 
446 void setInspectorTimerInfo(kj::Timer& timer, kj::Duration timerOffset) {
447 auto lockedState = state.lockExclusive();
448 lockedState->inspectorTimerInfo = InspectorTimerInfo{timer, timerOffset, getCurrentThreadId()};
449 }
450 
451 void setChannel(Worker::Isolate::InspectorChannelImpl& channel) {
452 auto lockedState = state.lockExclusive();
453 // There is only one active inspector channel at a time in workerd. The teardown of any
454 // previous channel should have invalidated `lockedState->channel`.
455 KJ_REQUIRE(lockedState->channel == kj::none);
456 lockedState->channel = channel;
457 }
458 
459 void resetChannel() {
460 auto lockedState = state.lockExclusive();
461 lockedState->channel = kj::none;
462 }
463 
464 // This method is called by v8 when a breakpoint or debugger statement is hit. This method
465 // processes debugger messages until `Debugger.resume()` is called, when v8 then calls
466 // `quitMessageLoopOnPause()`.
467 //
468 // This method is ultimately called from the `InspectorChannelImpl` and the isolate lock is
469 // held when this method is called.
470 void runMessageLoopOnPause(int contextGroupId) override {
471 auto lockedState = state.lockExclusive();
472 KJ_IF_SOME(channel, lockedState->channel) {
473 runMessageLoop = true;
474 do {
475 if (!dispatchOneMessageDuringPause(channel)) {
476 break;
477 }
478 } while (runMessageLoop);
479 }
480 }
481 
482 // This method is called by v8 to resume execution after a breakpoint is hit.
483 void quitMessageLoopOnPause() override {
484 runMessageLoop = false;
485 }
486 
487 private:
488 static bool dispatchOneMessageDuringPause(Worker::Isolate::InspectorChannelImpl& channel);
489 
490 struct InspectorTimerInfo {
491 kj::Timer& timer;
492 kj::Duration timerOffset;
493 uint64_t threadId;
494 };
495 
496 bool runMessageLoop;
497 
498 // State that may be set on a thread other than the isolate thread.
499 // These are typically set in attachInspector when an inspector connection is
500 // made.
501 struct State {
502 // Inspector channel to use to pump messages.
503 kj::Maybe<Worker::Isolate::InspectorChannelImpl&> channel;
504 
505 // The timer and offset for the inspector-serving thread.
506 kj::Maybe<InspectorTimerInfo> inspectorTimerInfo;
507 };
508 kj::MutexGuarded<State> state;
509};
510 
511static thread_local const Worker::Api* currentApi = nullptr;
512 
513const Worker::Api& Worker::Api::current() {
514 KJ_REQUIRE(currentApi != nullptr, "not running JavaScript");
515 return *currentApi;
516}
517 
518kj::Maybe<const Worker::Api&> Worker::Api::tryCurrent() {
519 if (currentApi != nullptr) {
520 return *currentApi;
521 }
522 return kj::none;
523}
524 
525jsg::Optional<jsg::Ref<api::CacheContext>> Worker::Api::getCtxCacheProperty(jsg::Lock& js) const {
526 return kj::none;
527}
528 
529struct Worker::Impl {
530 kj::Maybe<jsg::JsContext<api::ServiceWorkerGlobalScope>> context;
531 
532 // The environment blob to pass to handlers.
533 kj::Maybe<jsg::Value> env;
534 kj::Maybe<jsg::Value> ctxExports;
535 
536 // Note: The default export is given the string name "default", because that's what V8 tells us,
537 // and so it's easiest to go with it. I guess that means that you can't actually name an export
538 // "default"?
539 kj::HashMap<kj::String, api::ExportedHandler> namedHandlers;
540 kj::HashMap<kj::String, ActorClassInfo> actorClasses;
541 kj::HashMap<kj::String, EntrypointClass> statelessClasses;
542 kj::HashMap<kj::String, EntrypointClass> workflowClasses;
543 
544 // If set, then any attempt to use this worker shall throw this exception.
545 kj::Maybe<kj::Exception> permanentException;
546};
547 
548// Note that Isolate mutable state is protected by locking the JsgWorkerIsolate unless otherwise
549// noted.
550struct Worker::Isolate::Impl {
551 IsolateObserver& metrics;
552 kj::Own<InspectorClient> inspectorClient;
553 kj::Maybe<std::unique_ptr<v8_inspector::V8Inspector>> inspector;
554 InspectorPolicy inspectorPolicy;
555 kj::Maybe<kj::Own<v8::CpuProfiler>> profiler;
556 ActorCache::SharedLru actorCacheLru;
557 
558 // Used by JSG/Rust integration.
559 ::rust::Box<::workerd::rust::jsg::Realm> realm;
560 
561 // UUID for this isolate, initialized first time getUuid() is called.
562 kj::Lazy<kj::String> uuid;
563 
564 // Notification messages to deliver to the next inspector client when it connects.
565 kj::Vector<kj::String> queuedNotifications;
566 
567 // Set of warning log lines that should not be logged to the inspector again.
568 kj::HashSet<kj::String> warningOnceDescriptions;
569 
570 // Set of error log lines that should not be logged again.
571 kj::HashSet<kj::String> errorOnceDescriptions;
572 
573 // Instantaneous count of how many threads are trying to or have successfully obtained an
574 // AsyncLock on this isolate, used to implement getCurrentLoad().
575 mutable uint lockAttemptGauge = 0;
576 
577 // Atomically incremented upon every successful lock. The ThreadProgressCounter in Impl::Lock
578 // registers a reference to `lockSuccessCounter` as the thread's progress counter during a lock
579 // attempt. This allows watchdogs to see evidence of forward progress in other threads, even if
580 // their own thread has blocked waiting for the lock for a long time.
581 mutable uint64_t lockSuccessCount = 0;
582 
583 // Wrapper around JsgWorkerIsolate::Lock and various RAII objects which help us report metrics,
584 // measure instantaneous load, avoid spurious watchdog kills, and defer context destruction.
585 //
586 // Always use this wrapper in code which may face lock contention (that's mostly everywhere).
587 class Lock {
588 
589 public:
590 explicit Lock(
591 const Worker::Isolate& isolate, Worker::LockType lockType, jsg::V8StackScope& stackScope)
592 : impl(*isolate.impl),
593 metrics([&isolate, &lockType]() -> kj::Maybe<kj::Own<IsolateObserver::LockTiming>> {
594 KJ_SWITCH_ONEOF(lockType.origin) {
595 KJ_CASE_ONEOF(sync, Worker::Lock::TakeSynchronously) {
596 // TODO(perf): We could do some tracking here to discover overly harmful synchronous
597 // locks.
598 return isolate.getMetrics().tryCreateLockTiming(sync.getRequest());
599 }
600 KJ_CASE_ONEOF(async, AsyncLock*) {
601 KJ_REQUIRE(async->waiter->isolate.get() == &isolate,
602 "async lock was taken against a different isolate than the synchronous lock");
603 return kj::mv(async->lockTiming);
604 }
605 }
606 KJ_UNREACHABLE;
607 }()),
608 progressCounter(impl.lockSuccessCount),
609 oldCurrentApi(currentApi),
610 limitEnforcer(isolate.getLimitEnforcer()),
611 loggingOptions(isolate.loggingOptions),
612 lock(isolate.api->lock(stackScope)) {
613 WarnAboutIsolateLockScope::maybeWarn();
614 
615 // Increment the success count to expose forward progress to all threads.
616 __atomic_add_fetch(&impl.lockSuccessCount, 1, __ATOMIC_RELAXED);
617 metrics.locked();
618 
619 // We record the current lock so our GC prologue/epilogue callbacks can report GC time via
620 // Jaeger tracing.
621 KJ_DASSERT(impl.currentLock == kj::none, "Isolate lock taken recursively");
622 impl.currentLock = *this;
623 
624 // Now's a good time to destroy any workers queued up for destruction.
625 auto workersToDestroy = impl.workerDestructionQueue.lockExclusive()->pop();
626 for (auto& workerImpl: workersToDestroy.asArrayPtr()) {
627 KJ_IF_SOME(c, workerImpl->context) {
628 disposeContext(kj::mv(c));
629 }
630 workerImpl = nullptr;
631 }
632 
633 currentApi = isolate.api.get();
634 }
635 ~Lock() noexcept(false) {
636 currentApi = oldCurrentApi;
637 
638#ifdef KJ_DEBUG
639 // We lack a KJ_DASSERT_NONNULL because it would have to look a lot like KJ_IF_SOME, thus
640 // we use a pragma around KJ_DEBUG here.
641 auto& implCurrentLock = KJ_ASSERT_NONNULL(impl.currentLock, "Isolate lock released twice");
642 KJ_ASSERT(&implCurrentLock == this, "Isolate lock released recursively");
643#endif
644 
645 if (shouldReportIsolateMetrics) {
646 // The isolate asked this lock to report the stats when it released. Let's do it.
647 limitEnforcer.reportMetrics(impl.metrics);
648 }
649 impl.currentLock = kj::none;
650 }
651 KJ_DISALLOW_COPY_AND_MOVE(Lock);
652 
653 void setupContext(v8::Local<v8::Context> context) {
654 // The V8Inspector implements the `console` object.
655 KJ_IF_SOME(i, impl.inspector) {
656 i.get()->contextCreated(
657 v8_inspector::V8ContextInfo(context, 1, jsg::toInspectorStringView("Worker")));
658 }
659 Worker::setupContext(*lock, context, loggingOptions);
660 }
661 
662 void disposeContext(jsg::JsContext<api::ServiceWorkerGlobalScope> context) {
663 lock->withinHandleScope([&] {
664 auto v8Context = context.getHandle(*lock);
665 context->clear();
666 KJ_IF_SOME(i, impl.inspector) {
667 i.get()->contextDestroyed(v8Context);
668 }
669 { auto drop = kj::mv(context); }
670 lock->v8Isolate->ContextDisposedNotification(v8::ContextDependants::kNoDependants);
671 });
672 }
673 
674 void gcPrologue() {
675 metrics.gcPrologue();
676 // Filter out tracked WASM instance entries where the instance has been
677 // garbage-collected (weak instanceRef is empty), allowing the linear memory
678 // to be reclaimed.
679 limitEnforcer.getTrackedWasmInstances().filter(*lock);
680 }
681 void gcEpilogue() {
682 metrics.gcEpilogue();
683 }
684 
685 // Call limitEnforcer.exitJs(), and also schedule to call limitEnforcer.reportMetrics()
686 // later. Returns true if condemned. We take a mutable reference to it to make sure the caller
687 // believes it has exclusive access.
688 bool checkInWithLimitEnforcer(Worker::Isolate& isolate);
689 
690 private:
691 const Impl& impl;
692 IsolateObserver::LockRecord metrics;
693 ThreadProgressCounter progressCounter;
694 bool shouldReportIsolateMetrics = false;
695 const Api* oldCurrentApi;
696 
697 const IsolateLimitEnforcer& limitEnforcer; // only so we can call getIsolateStats()
698 
699 // When structuredLogging is YES AND consoleMode is STDOUT js logs will be emitted to STDOUT
700 // as newline separated json objects
701 LoggingOptions loggingOptions;
702 
703 public:
704 kj::Own<jsg::Lock> lock;
705 };
706 
707 // Protected by v8::Locker -- if v8::Locker::IsLocked(isolate) is true, then it is safe to access
708 // this variable.
709 mutable kj::Maybe<Lock&> currentLock;
710 
711 static constexpr auto WORKER_DESTRUCTION_QUEUE_INITIAL_SIZE = 8;
712 static constexpr auto WORKER_DESTRUCTION_QUEUE_MAX_CAPACITY = 100;
713 
714 // Similar in spirit to the deferred destruction queue in jsg::IsolateBase. When a Worker is
715 // destroyed, it puts its Impl, which contains objects that need to be destroyed under the isolate
716 // lock, into this queue. Our own Isolate::Impl::Lock implementation then clears this queue the
717 // next time the isolate is locked, whether that be by a connection thread, or the Worker's own
718 // destructor if it owns the last `kj::Own<const Script>` reference.
719 //
720 // Fairly obviously, this member is protected by its own mutex, not the isolate lock.
721 const kj::MutexGuarded<BatchQueue<kj::Own<Worker::Impl>>> workerDestructionQueue{
722 WORKER_DESTRUCTION_QUEUE_INITIAL_SIZE, WORKER_DESTRUCTION_QUEUE_MAX_CAPACITY};
723 // TODO(cleanup): The only reason this exists and we can't just rely on the isolate's regular
724 // deferred destruction queue to lazily destroy the various V8 objects in Worker::Impl is
725 // because our GlobalScope object needs to have a function called on it, and any attached
726 // inspector needs to be notified. JSG doesn't know about these things.
727 
728 struct IsolateState {
729 kj::Own<InspectorClient> inspectorClient;
730 kj::Maybe<std::unique_ptr<v8_inspector::V8Inspector>> inspector;
731 ::rust::Box<::workerd::rust::jsg::Realm> realm;
732 };
733 
734 static IsolateState initIsolate(
735 const Api& api, IsolateLimitEnforcer& limitEnforcer, InspectorPolicy inspectorPolicy) {
736 auto inspectorClient = kj::heap<InspectorClient>();
737 // Default constructor of ::rust::Box is deleted, so we use a Maybe to delay initialization.
738 kj::Maybe<::rust::Box<::workerd::rust::jsg::Realm>> realm;
739 kj::Maybe<std::unique_ptr<v8_inspector::V8Inspector>> inspector;
740 jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) {
741 auto lock = api.lock(stackScope);
742 auto featureFlagsWords = capnp::canonicalize(api.getFeatureFlags());
743 realm = ::workerd::rust::jsg::realm_create(
744 lock->v8Isolate, featureFlagsWords.asBytes().as<kj_rs::Rust>());
745 lock->v8Isolate->SetData(
746 ::workerd::jsg::SetDataIndex::SET_DATA_RUST_REALM, &*KJ_REQUIRE_NONNULL(realm));
747 
748 limitEnforcer.customizeIsolate(lock->v8Isolate);
749 if (inspectorPolicy != InspectorPolicy::DISALLOW) {
750 // We just created our isolate, so we don't need to use Isolate::Impl::Lock.
751 KJ_ASSERT(!isMultiTenantProcess(), "inspector is not safe in multi-tenant processes");
752 inspector = v8_inspector::V8Inspector::create(lock->v8Isolate, inspectorClient.get());
753 }
754 });
755 return {kj::mv(inspectorClient), kj::mv(inspector), kj::mv(KJ_REQUIRE_NONNULL(realm))};
756 }
757 
758 Impl(IsolateObserver& metrics,
759 IsolateLimitEnforcer& limitEnforcer,
760 InspectorPolicy inspectorPolicy,
761 IsolateState state)
762 : metrics(metrics),
763 inspectorClient(kj::mv(state.inspectorClient)),
764 inspector(kj::mv(state.inspector)),
765 inspectorPolicy(inspectorPolicy),
766 actorCacheLru(limitEnforcer.getActorCacheLruOptions()),
767 realm(kj::mv(state.realm)) {}
768 
769 Impl(const Api& api,
770 IsolateObserver& metrics,
771 IsolateLimitEnforcer& limitEnforcer,
772 InspectorPolicy inspectorPolicy)
773 : Impl(metrics,
774 limitEnforcer,
775 inspectorPolicy,
776 initIsolate(api, limitEnforcer, inspectorPolicy)) {}
777};
778 
779namespace {
780 
781class CpuProfilerDisposer final: public kj::Disposer {
782 public:
783 virtual void disposeImpl(void* pointer) const override {
784 reinterpret_cast<v8::CpuProfiler*>(pointer)->Dispose();
785 }
786 
787 static const CpuProfilerDisposer instance;
788};
789 
790const CpuProfilerDisposer CpuProfilerDisposer::instance{};
791 
792static constexpr kj::StringPtr PROFILE_NAME = "Default Profile"_kj;
793 
794static void setSamplingInterval(v8::CpuProfiler& profiler, int interval) {
795 profiler.SetSamplingInterval(interval);
796}
797 
798static void startProfiling(jsg::Lock& js, v8::CpuProfiler& profiler) {
799 js.withinHandleScope([&] {
800 v8::CpuProfilingOptions options(
801 v8::kLeafNodeLineNumbers, v8::CpuProfilingOptions::kNoSampleLimit);
802 profiler.StartProfiling(jsg::v8StrIntern(js.v8Isolate, PROFILE_NAME), kj::mv(options));
803 });
804}
805 
806static void stopProfiling(jsg::Lock& js, v8::CpuProfiler& profiler, cdp::Command::Builder& cmd) {
807 js.withinHandleScope([&] {
808 auto cpuProfile = profiler.StopProfiling(jsg::v8StrIntern(js.v8Isolate, PROFILE_NAME));
809 if (cpuProfile == nullptr) return; // profiling never started
810 
811 kj::Vector<const v8::CpuProfileNode*> allNodes;
812 kj::Vector<const v8::CpuProfileNode*> unvisited;
813 
814 unvisited.add(cpuProfile->GetTopDownRoot());
815 while (!unvisited.empty()) {
816 auto next = unvisited.back();
817 allNodes.add(next);
818 unvisited.removeLast();
819 for (int i = 0; i < next->GetChildrenCount(); i++) {
820 unvisited.add(next->GetChild(i));
821 }
822 }
823 
824 auto res = cmd.getProfilerStop().initResult();
825 auto profile = res.initProfile();
826 profile.setStartTime(cpuProfile->GetStartTime());
827 profile.setEndTime(cpuProfile->GetEndTime());
828 
829 auto nodes = profile.initNodes(allNodes.size());
830 for (auto i: kj::indices(allNodes)) {
831 auto nodeBuilder = nodes[i];
832 nodeBuilder.setId(allNodes[i]->GetNodeId());
833 
834 auto callFrame = nodeBuilder.initCallFrame();
835 callFrame.setFunctionName(allNodes[i]->GetFunctionNameStr());
836 callFrame.setScriptId(kj::str(allNodes[i]->GetScriptId()));
837 callFrame.setUrl(allNodes[i]->GetScriptResourceNameStr());
838 // V8 locations are 1-based, but CDP locations are 0-based...
839 callFrame.setLineNumber(allNodes[i]->GetLineNumber() - 1);
840 callFrame.setColumnNumber(allNodes[i]->GetColumnNumber() - 1);
841 
842 nodeBuilder.setHitCount(allNodes[i]->GetHitCount());
843 
844 auto children = nodeBuilder.initChildren(allNodes[i]->GetChildrenCount());
845 for (int j = 0; j < allNodes[i]->GetChildrenCount(); j++) {
846 children.set(j, allNodes[i]->GetChild(j)->GetNodeId());
847 }
848 
849 auto hitLineCount = allNodes[i]->GetHitLineCount();
850 auto lineBuffer = kj::heapArray<v8::CpuProfileNode::LineTick>(hitLineCount);
851 allNodes[i]->GetLineTicks(lineBuffer.begin(), lineBuffer.size());
852 
853 auto positionTicks = nodeBuilder.initPositionTicks(hitLineCount);
854 for (uint j = 0; j < hitLineCount; j++) {
855 auto positionTick = positionTicks[j];
856 positionTick.setLine(lineBuffer[j].line);
857 positionTick.setTicks(lineBuffer[j].hit_count);
858 }
859 }
860 
861 auto sampleCount = cpuProfile->GetSamplesCount();
862 auto samples = profile.initSamples(sampleCount);
863 auto timeDeltas = profile.initTimeDeltas(sampleCount);
864 auto lastTimestamp = cpuProfile->GetStartTime();
865 for (int i = 0; i < sampleCount; i++) {
866 samples.set(i, cpuProfile->GetSample(i)->GetNodeId());
867 auto sampleTime = cpuProfile->GetSampleTimestamp(i);
868 timeDeltas.set(i, sampleTime - lastTimestamp);
869 lastTimestamp = sampleTime;
870 }
871 });
872}
873 
874} // anonymous namespace
875 
876struct Worker::Script::Impl {
877 kj::Own<workerd::VirtualFileSystem> vfs;
878 kj::Maybe<kj::Arc<workerd::jsg::modules::ModuleRegistry>> maybeNewModuleRegistry;
879 // When using the new module registry, the module registry itself holds the
880 // SchemaLoader, so we don't need to hold it here. When using the original
881 // module registry, however, we need a schema loader to instantiate capnp
882 // modules and bindings.
883 kj::Maybe<kj::Own<capnp::SchemaLoader>> maybeSchemaLoader;
884 
885 kj::OneOf<jsg::NonModuleScript, kj::Path> unboundScriptOrMainModule;
886 
887 kj::Array<CompiledGlobal> globals;
888 
889 kj::Maybe<jsg::JsContext<api::ServiceWorkerGlobalScope>> moduleContext;
890 
891 // If set, then any attempt to use this script shall throw this exception.
892 kj::Maybe<kj::Exception> permanentException;
893 
894 Impl(kj::Own<workerd::VirtualFileSystem> vfs,
895 kj::Maybe<kj::Arc<workerd::jsg::modules::ModuleRegistry>> maybeNewModuleRegistry)
896 : vfs(kj::mv(vfs)),
897 maybeNewModuleRegistry(kj::mv(maybeNewModuleRegistry)) {
898 if (this->maybeNewModuleRegistry == kj::none) {
899 maybeSchemaLoader = kj::heap<capnp::SchemaLoader>();
900 }
901 }
902 
903 struct DynamicImportResult {
904 jsg::Value value;
905 bool isException = false;
906 DynamicImportResult(jsg::Value value, bool isException = false)
907 : value(kj::mv(value)),
908 isException(isException) {}
909 };
910 using DynamicImportHandler = kj::Function<jsg::Value()>;
911 
912 void configureDynamicImports(jsg::Lock& js, jsg::ModuleRegistry& modules) {
913 // This is only used with the original module registry implementation.
914 KJ_ASSERT(!FeatureFlags::get(js).getNewModuleRegistry(),
915 "legacy dynamic imports must not be used with the new module registry");
916 static auto constexpr handleDynamicImport =
917 [](kj::Own<const Worker> worker, DynamicImportHandler handler,
918 kj::Maybe<jsg::Ref<jsg::AsyncContextFrame>> asyncContext)
919 -> kj::Promise<DynamicImportResult> {
920 co_await kj::yield();
921 auto asyncLock = co_await worker->takeAsyncLockWithoutRequest(nullptr);
922 
923 co_return worker->runInLockScope(asyncLock, [&](Worker::Lock& lock) {
924 TmpDirStoreScope tmpDirStoreScope;
925 return JSG_WITHIN_CONTEXT_SCOPE(lock, lock.getContext(), [&](jsg::Lock& js) {
926 jsg::AsyncContextFrame::Scope asyncContextScope(js, asyncContext);
927 
928 // We have to wrap the call to handler in a try catch here because
929 // we have to tunnel any jsg::JsExceptionThrown instances back.
930 v8::TryCatch tryCatch(js.v8Isolate);
931 ExceptionOrDuration limitErrorOrTime = 0 * kj::NANOSECONDS;
932 try {
933 auto limitScope = worker->getIsolate().getLimitEnforcer().enterDynamicImportJs(
934 lock, limitErrorOrTime);
935 return DynamicImportResult(handler());
936 } catch (jsg::JsExceptionThrown&) {
937 // Handled below...
938 } catch (kj::Exception& ex) {
939 kj::throwFatalException(kj::mv(ex));
940 }
941 
942 KJ_ASSERT(tryCatch.HasCaught());
943 if (!tryCatch.CanContinue() || tryCatch.Exception().IsEmpty()) {
944 // There's nothing else we can do here but throw a generic fatal exception.
945 KJ_SWITCH_ONEOF(limitErrorOrTime) {
946 KJ_CASE_ONEOF(limitError, kj::Exception) {
947 kj::throwFatalException(kj::mv(limitError));
948 }
949 KJ_CASE_ONEOF_DEFAULT {
950 kj::throwFatalException(
951 JSG_KJ_EXCEPTION(FAILED, Error, "Failed to load dynamic module."));
952 }
953 }
954 }
955 return DynamicImportResult(js.v8Ref(tryCatch.Exception()), true);
956 });
957 });
958 };
959 
960 modules.setDynamicImportCallback([](jsg::Lock& js, DynamicImportHandler handler) mutable {
961 KJ_IF_SOME(context, IoContext::tryCurrent()) {
962 // If we are within the scope of a IoContext, then we are going to pop
963 // out of it to perform the actual module instantiation.
964 
965 return context.awaitIo(js,
966 handleDynamicImport(kj::atomicAddRef(context.getWorker()), kj::mv(handler),
967 jsg::AsyncContextFrame::currentRef(js)),
968 [](jsg::Lock& js, DynamicImportResult result) {
969 if (result.isException) {
970 return js.rejectedPromise<jsg::Value>(kj::mv(result.value));
971 }
972 return js.resolvedPromise(kj::mv(result.value));
973 });
974 }
975 
976 // If we got here, there is no current IoContext. We're going to perform the
977 // module resolution synchronously and we do not have to worry about blocking any
978 // i/o. We get here, for instance, when dynamic import is used at the top level of
979 // a script (which is weird, but allowed).
980 //
981 // We do not need to use limitEnforcer.enterDynamicImportJs() here because this should
982 // already be covered by the startup resource limiter.
983 return js.resolvedPromise(handler());
984 });
985 }
986 
987 kj::Maybe<const workerd::jsg::modules::ModuleRegistry&> getNewModuleRegistry() const {
988 return maybeNewModuleRegistry.map(
989 [](auto& r) -> const workerd::jsg::modules::ModuleRegistry& { return *r.get(); });
990 }
991};
992 
993namespace {
994 
995// Given an array of strings, return a valid serialized JSON string like:
996// {"flags":["minimal_subrequests",...]}
997//
998// Return null if the array is empty.
999kj::Maybe<kj::String> makeCompatJson(kj::ArrayPtr<kj::StringPtr> enableFlags) {
1000 if (enableFlags.size() == 0) {
1001 return kj::none;
1002 }
1003 
1004 // Calculate the size of the string we're going to generate.
1005 constexpr auto PREFIX = "{\"flags\":["_kj;
1006 constexpr auto SUFFIX = "]}"_kj;
1007 uint size = std::accumulate(enableFlags.begin(), enableFlags.end(),
1008 // We need two quotes and one comma for each enable-flag past the first, plus a NUL char.
1009 PREFIX.size() + SUFFIX.size() + 3 * enableFlags.size(),
1010 [](uint z, kj::StringPtr s) { return z + s.size(); });
1011 
1012 kj::Vector<char> json(size);
1013 
1014 json.addAll(PREFIX);
1015 
1016 bool first = true;
1017 for (auto flag: enableFlags) {
1018 if (first) {
1019 first = false;
1020 } else {
1021 json.add(',');
1022 }
1023 
1024 json.add('"');
1025 
1026 for (auto& c: flag.asArray()) {
1027 // TODO(cleanup): Copied from simpleJsonStringCheck(). Hopefully this will
1028 // go away forever soon.
1029 KJ_REQUIRE(c != '\"');
1030 KJ_REQUIRE(c != '\\');
1031 KJ_REQUIRE(c >= 0x20);
1032 }
1033 json.addAll(flag);
1034 
1035 json.add('"');
1036 }
1037 
1038 json.addAll(SUFFIX);
1039 json.add('\0');
1040 
1041 return kj::String(json.releaseAsArray());
1042}
1043 
1044// When a promise is created in a different IoContext, we need to use a
1045// kj::CrossThreadFulfiller in order to wait on it. The Waiter instance will
1046// be held on the Promise itself, and will be fulfilled/rejected when the
1047// promise is resolved or rejected. This will signal all of the waiters
1048// from other IoContexts.
1049jsg::Promise<void> addCrossThreadPromiseWaiter(jsg::Lock& js, v8::Local<v8::Promise>& promise) {
1050 auto waiter = kj::newPromiseAndCrossThreadFulfiller<void>();
1051 
1052 struct Waiter: public kj::Refcounted {
1053 kj::Maybe<kj::Own<kj::CrossThreadPromiseFulfiller<void>>> fulfiller;
1054 void done() {
1055 KJ_IF_SOME(f, fulfiller) {
1056 // Done this way so that the fulfiller is released as soon as possible
1057 // when done as the JS promise may not clean up reactions right away.
1058 f->fulfill();
1059 fulfiller = kj::none;
1060 }
1061 }
1062 Waiter(kj::Own<kj::CrossThreadPromiseFulfiller<void>> fulfiller)
1063 : fulfiller(kj::mv(fulfiller)) {}
1064 };
1065 
1066 auto fulfiller = kj::refcounted<Waiter>(kj::mv(waiter.fulfiller));
1067 
1068 auto onSuccess = [waiter = kj::addRef(*fulfiller)](
1069 jsg::Lock& js, jsg::Value value) mutable { waiter->done(); };
1070 
1071 auto onFailure = [waiter = kj::mv(fulfiller)](
1072 jsg::Lock& js, jsg::Value exception) mutable { waiter->done(); };
1073 
1074 js.toPromise(promise).then(js, kj::mv(onSuccess), kj::mv(onFailure));
1075 
1076 return IoContext::current().awaitIo(js, kj::mv(waiter.promise));
1077}
1078 
1079struct HeapSnapshotDeleter: public kj::Disposer {
1080 static const HeapSnapshotDeleter INSTANCE;
1081 void disposeImpl(void* ptr) const override {
1082 auto snapshot = const_cast<v8::HeapSnapshot*>(static_cast<const v8::HeapSnapshot*>(ptr));
1083 snapshot->Delete();
1084 }
1085};
1086const HeapSnapshotDeleter HeapSnapshotDeleter::INSTANCE;
1087 
1088} // namespace
1089 
1090Worker::Isolate::Isolate(kj::Own<Api> apiParam,
1091 kj::Own<IsolateObserver> metricsParam,
1092 kj::StringPtr id,
1093 kj::Own<IsolateLimitEnforcer> limitEnforcerParam,
1094 InspectorPolicy inspectorPolicy,
1095 LoggingOptions loggingOptions)
1096 : metrics(kj::mv(metricsParam)),
1097 id(kj::str(id)),
1098 limitEnforcer(kj::mv(limitEnforcerParam)),
1099 cpuLimitNearlyExceededCallback(
1100 kj::MutexGuarded<kj::Maybe<kj::Function<void(void)>>>(kj::none)),
1101 api(kj::mv(apiParam)),
1102 loggingOptions(loggingOptions),
1103 featureFlagsForFl(makeCompatJson(decompileCompatibilityFlagsForFl(api->getFeatureFlags()))),
1104 impl(kj::heap<Impl>(*api, *metrics, *limitEnforcer, inspectorPolicy)),
1105 weakIsolateRef(WeakIsolateRef::wrap(this)),
1106 traceAsyncContextKey(kj::refcounted<jsg::AsyncContextFrame::StorageKey>()),
1107 userTraceAsyncContextKey(kj::refcounted<jsg::AsyncContextFrame::StorageKey>()) {
1108 api->setIsolateObserver(*metrics);
1109 metrics->created();
1110 // We just created our isolate, so we don't need to use Isolate::Impl::Lock (nor an async lock).
1111 jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) {
1112 auto lock = api->lock(stackScope);
1113 auto features = api->getFeatureFlags();
1114 
1115 KJ_DASSERT(lock->v8Isolate->GetNumberOfDataSlots() >= jsg::SET_DATA_SLOTS_IN_USE);
1116 KJ_DASSERT(lock->v8Isolate->GetData(jsg::SET_DATA_ISOLATE) == nullptr);
1117 lock->v8Isolate->SetData(jsg::SET_DATA_ISOLATE, this);
1118 
1119 lock->setCaptureThrowsAsRejections(features.getCaptureThrowsAsRejections());
1120 // TODO(cleanup): Now that this list has grown significantly, we should probably
1121 // refactor to pass all of the options in a single call instead of one by one.
1122 if (features.getSetToStringTag()) {
1123 lock->setToStringTag();
1124 }
1125 if (features.getShouldSetImmutablePrototype() || features.getPythonWorkers()) {
1126 lock->setImmutablePrototype();
1127 }
1128 if (features.getSpecCompliantPropertyAttributes()) {
1129 lock->setSpecCompliantPropertyAttributes();
1130 }
1131 if (features.getNodeJsCompatV2()) {
1132 lock->setNodeJsCompatEnabled();
1133 }
1134 if (features.getEnableNodeJsProcessV2()) {
1135 lock->setNodeJsProcessV2Enabled();
1136 }
1137 if (features.getRequireReturnsDefaultExport()) {
1138 lock->setRequireReturnsDefaultExportEnabled();
1139 }
1140 if (features.getThrowOnUnrecognizedImportAssertion()) {
1141 lock->setThrowOnUnrecognizedImportAssertion();
1142 }
1143 if (features.getNoTopLevelAwaitInRequire()) {
1144 lock->disableTopLevelAwait();
1145 }
1146 if (features.getEnhancedErrorSerialization()) {
1147 lock->setUsingEnhancedErrorSerialization();
1148 }
1149 if (features.getFastJsgStruct()) {
1150 lock->setUsingFastJsgStruct();
1151 }
1152 
1153 if (impl->inspector != kj::none || ::kj::_::Debug::shouldLog(::kj::LogSeverity::INFO)) {
1154 lock->setLoggerCallback([this](jsg::Lock& js, kj::StringPtr message) {
1155 if (impl->inspector != kj::none) {
1156 logMessage(js, static_cast<uint16_t>(cdp::LogType::WARNING), message);
1157 }
1158 KJ_LOG(INFO, "console warning", message);
1159 });
1160 lock->setErrorReporterCallback([this](jsg::Lock& js, kj::String desc,
1161 const jsg::JsValue& error, const jsg::JsMessage& message) {
1162 // Only add exception to trace when running within an I/O context with a tracer.
1163 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
1164 KJ_IF_SOME(tracer, ioContext.getWorkerTracer()) {
1165 addExceptionToTrace(js, ioContext, tracer, UncaughtExceptionSource::REQUEST_HANDLER,
1166 error, api->getErrorInterfaceTypeHandler(js));
1167 }
1168 }
1169 
1170 KJ_IF_SOME(i, impl->inspector) {
1171 jsg::sendExceptionToInspector(js, *i.get(), kj::str(desc), error, message);
1172 }
1173 
1174 // Run with --verbose to log JS exceptions to stderr. Useful when running tests.
1175 KJ_LOG(INFO, "uncaught exception", desc);
1176 });
1177 }
1178 
1179 // By default, V8's memory pressure level is "none". This tells V8 that no one else on the
1180 // machine is competing for memory so it might as well use all it wants and be lazy about GC.
1181 //
1182 // In our production environment, however, we can safely assume that there is always memory
1183 // pressure, because every machine is handling thousands of tenants all the time. So we might
1184 // as well just throw the switch to "moderate" right away.
1185 lock->v8Isolate->MemoryPressureNotification(v8::MemoryPressureLevel::kModerate);
1186 
1187 // Register GC prologue and epilogue callbacks so that we can report GC CPU time via the
1188 // "request_context" Jaeger span.
1189 lock->v8Isolate->AddGCPrologueCallback(
1190 [](v8::Isolate* isolate, v8::GCType type, v8::GCCallbackFlags flags, void* data) noexcept {
1191 // We assume that a v8::Locker is alive during GC.
1192 KJ_DASSERT(v8::Locker::IsLocked(isolate));
1193 auto& self = *reinterpret_cast<Isolate*>(data);
1194 // However, currentLock might not be available, if (like in our Worker::Isolate constructor) we
1195 // don't use a Worker::Isolate::Impl::Lock.
1196 KJ_IF_SOME(currentLock, self.impl->currentLock) {
1197 currentLock.gcPrologue();
1198 }
1199 }, this);
1200 lock->v8Isolate->AddGCEpilogueCallback(
1201 [](v8::Isolate* isolate, v8::GCType type, v8::GCCallbackFlags flags, void* data) noexcept {
1202 // We make similar assumptions about v8::Locker and currentLock as in the prologue callback.
1203 KJ_DASSERT(v8::Locker::IsLocked(isolate));
1204 auto& self = *reinterpret_cast<Isolate*>(data);
1205 KJ_IF_SOME(currentLock, self.impl->currentLock) {
1206 currentLock.gcEpilogue();
1207 }
1208 }, this);
1209 lock->v8Isolate->SetPromiseRejectCallback([](v8::PromiseRejectMessage message) {
1210 // TODO(cleanup): IoContext doesn't really need to be involved here. We are trying to call
1211 // a method of ServiceWorkerGlobalScope, which is the context object. So we should be able to
1212 // do something like unwrap(lock, isolate->GetCurrentContext()).emitPromiseRejection().
1213 // However, JSG doesn't currently provide an easy way to do this.
1214 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
1215 try {
1216 ioContext.getCurrentLock().reportPromiseRejectEvent(message);
1217 } catch (jsg::JsExceptionThrown&) {
1218 // V8 expects us to just return.
1219 return;
1220 }
1221 }
1222 });
1223 
1224 // The PromiseCrossContextCallback is used to allow cross-IoContext promise following.
1225 // When the IoContext scope is entered, we set the "promise context tag" associated
1226 // with the IoContext on the Isolate that is locked. Any Promise that is created within
1227 // that scope will be tagged with the same promise context tag. When an attempt to
1228 // follow a promise occurs (e.g. either using Promise.prototype.then() or await, etc)
1229 // our patched v8 logic will check to see if the followed promise's tag matches the
1230 // current Isolate tag. If they do not, then v8 will invoke this callback. The promise
1231 // here is the promise that belongs to a different IoContext.
1232 lock->v8Isolate->SetPromiseCrossContextCallback(
1233 [](v8::Local<v8::Context> context, v8::Local<v8::Promise> promise,
1234 v8::Local<v8::Object> tag) -> v8::MaybeLocal<v8::Promise> {
1235 auto& js = jsg::Lock::current();
1236 try {
1237 // Generally this condition is only going to happen when using dynamic imports.
1238 // It should not be common.
1239 JSG_REQUIRE(IoContext::hasCurrent(), Error,
1240 "Unable to wait on a promise created within a request when not running within a "
1241 "request.");
1242 
1243 return js.wrapSimplePromise(
1244 addCrossThreadPromiseWaiter(js, promise)
1245 .then(js, [promise = js.v8Ref(promise.As<v8::Value>())](auto& js) mutable {
1246 // Once the waiter has been resolved, return the now settled promise.
1247 // Since the promise has been settled, it is now safe to access from
1248 // other requests. Note that the resolved value of the promise still
1249 // might not be safe to access! (e.g. if it contains any IoOwns attached
1250 // to the other request IoContext).
1251 return kj::mv(promise);
1252 }));
1253 } catch (jsg::JsExceptionThrown&) {
1254 // Exceptions here are generally unexpected but possible because the jsg::Promise
1255 // then can fail if the isolate is in the process of being torn down. Let's just
1256 // return control back to V8 which should handle the case.
1257 return v8::MaybeLocal<v8::Promise>();
1258 } catch (...) {
1259 auto ex = kj::getCaughtExceptionAsKj();
1260 KJ_LOG(ERROR, "Setting promise cross context follower failed unexpectedly", ex);
1261 jsg::throwInternalError(js.v8Isolate, kj::mv(ex));
1262 return v8::MaybeLocal<v8::Promise>();
1263 }
1264 });
1265 
1266 // The PromiseCrossContextResolveCallback is used to ensure that promise reactions
1267 // are only scheduled on the microtask queue from the appropriate IoContext for the
1268 // promise. Huh? Yeah, that's not super clear... let me explain a bit more.
1269 // Every request runs in its own IoContext.
1270 // Some I/O objects are bound to the IoContext when they are created.
1271 // If these objects are accessed from the wrong IoContext, things blow up.
1272 // If I create a promise in one request and pass the resolve/reject functions
1273 // off to a different request, bad things can happen because the IoContext can
1274 // actually change in the promise continuation. Take the following case for example:
1275 //
1276 // In request one:
1277 //
1278 // const ab = AbortSignal.abort(); // AbortSignal is bound to the IoContext
1279 // const { promise, resolve } = Promise.withResolvers();
1280 // globalThis.resolve = resolve;
1281 // await promise;
1282 // console.log(ab.aborted);
1283 //
1284 // In request two:
1285 //
1286 // globalThis.resolve();
1287 //
1288 // What previously would happen is that the `console.log(ab.aborted) after the
1289 // `await promise` in request one would fail with an error because the current
1290 // IoContext would change! (it would be the IoContext from request two!).
1291 //
1292 // That's bad.
1293 //
1294 // So this callback is added to ensure that the promise reactions for the promise
1295 // being resolved are not scheduled until we are back in the correct IoContext for
1296 // the promise.
1297 //
1298 // This happens by (ab)using the DeleteQueue that is specific to the owning
1299 // IoContext. When the IoContext is entered, the isolate is updated with a
1300 // current "promise tag". Whenever a promise is created, it is associated with
1301 // the isolate's current tag. Whenever a promise is followed (calling .then, etc),
1302 // we check the tag and arrange for a cross-thread resolve. When the promise is
1303 // resolved or rejected, we check the tag also. If the promise tag and the current
1304 // isolate tag do not match, the function below is called.
1305 if (features.getHandleCrossRequestPromiseResolution()) {
1306 lock->v8Isolate->SetPromiseCrossContextResolveCallback(
1307 [](v8::Isolate* isolate, v8::Local<v8::Value> tag, v8::Local<v8::Data> reactions,
1308 v8::Local<v8::Value> argument,
1309 std::function<void(v8::Isolate * isolate, v8::Local<v8::Data> reactions,
1310 v8::Local<v8::Value> argument)> callback) -> v8::Maybe<void> {
1311 try {
1312 auto& js = jsg::Lock::from(isolate);
1313 
1314 // The promise tag is generally opaque except for right here. The tag
1315 // wraps an instanceof kj::Own<IoCrossContextExecutor>, which wraps an atomically
1316 // refcounted pointer to the DeleteQueue for the correct isolate.
1317 // We simply pass the given callback, reactions, and argument to
1318 // a function that will be added to the queue inside DeleteQueue.
1319 // The next time the relevant IoContext is entered, this queue will
1320 // be drained and the actions will be run. Adding the task to the
1321 // delete queue will also signal the IoContext that it should wake
1322 // up and drain the queue. Simple, eh?
1323 //
1324 // A word of warning tho! It is possible for the IoContext to be
1325 // destroyed before the promise is resolved. Any actions that have
1326 // already been added to the queue would end up being dropped silently
1327 // on the floor. Actions that are added to the queue now will be run
1328 // immediately in the wrong IoContext.
1329 auto& ref = jsg::unwrapOpaqueRef<kj::Own<IoCrossContextExecutor>>(isolate, tag);
1330 ref->execute(js,
1331 [reactions = jsg::Data(isolate, reactions), argument = jsg::V8Ref(isolate, argument),
1332 callback = kj::mv(callback)](jsg::Lock& js) mutable {
1333 callback(js.v8Isolate, reactions.getHandle(js), argument.getHandle(js));
1334 });
1335 return v8::JustVoid();
1336 } catch (jsg::JsExceptionThrown&) {
1337 // Exceptions here are generally unexpected but possible because the jsg::Promise
1338 // then can fail if the isolate is in the process of being torn down. Let's just
1339 // return control back to V8 which should handle the case.
1340 // Note that errors thrown here and below should cause the resolve() or reject()
1341 // function calls to throw, which is unusual. Just important to keep that in mind.
1342 // Most likely errors thrown here are fatal so that should be OK.
1343 return v8::Nothing<void>();
1344 } catch (...) {
1345 jsg::throwInternalError(isolate, kj::getCaughtExceptionAsKj());
1346 return v8::Nothing<void>();
1347 }
1348 });
1349 }
1350 });
1351}
1352 
1353Worker::Script::Script(kj::Own<const Isolate> isolateParam,
1354 kj::StringPtr id,
1355 const Script::Source& source,
1356 IsolateObserver::StartType startType,
1357 bool logNewScript,
1358 kj::Maybe<ValidationErrorReporter&> errorReporter,
1359 kj::Maybe<kj::Own<api::pyodide::ArtifactBundler_State>> artifacts,
1360 SpanParent parentSpan,
1361 kj::Own<workerd::VirtualFileSystem> vfs,
1362 kj::Maybe<kj::Arc<workerd::jsg::modules::ModuleRegistry>> maybeNewModuleRegistry)
1363 : isolate(kj::mv(isolateParam)),
1364 id(kj::str(id)),
1365 modular(source.variant.is<ModulesSource>()),
1366 python(modular && source.variant.get<ModulesSource>().isPython),
1367 impl(kj::heap<Impl>(kj::mv(vfs), kj::mv(maybeNewModuleRegistry))),
1368 dynamicEnvBuilder(source.dynamicEnvBuilder.map(
1369 [](const auto& inst) -> kj::Arc<DynamicEnvBuilder> { return inst.addRef(); })) {
1370 auto parseMetrics = isolate->metrics->parse(startType);
1371 // TODO(perf): It could make sense to take an async lock when constructing a script if we
1372 // co-locate multiple scripts in the same isolate. As of this writing, we do not, except in
1373 // previews, where it doesn't matter. If we ever do co-locate multiple scripts in the same
1374 // isolate, we may wish to make the RequestObserver object available here, in order to
1375 // attribute lock timing to that request.
1376 jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) {
1377 Isolate::Impl::Lock recordedLock(
1378 *isolate, Worker::Lock::TakeSynchronously(kj::none), stackScope);
1379 auto& lock = *recordedLock.lock;
1380 
1381 // If we throw an exception, it's important that `impl` is destroyed under lock.
1382 KJ_ON_SCOPE_FAILURE({
1383 auto implToDestroy = kj::mv(impl);
1384 KJ_IF_SOME(c, implToDestroy->moduleContext) {
1385 recordedLock.disposeContext(kj::mv(c));
1386 } else {
1387 // Else block to avoid dangling else clang warning.
1388 }
1389 });
1390 
1391 lock.withinHandleScope([&] {
1392 if (isolate->impl->inspector != kj::none || errorReporter != kj::none) {
1393 lock.v8Isolate->SetCaptureStackTraceForUncaughtExceptions(true);
1394 }
1395 
1396 v8::Local<v8::Context> context;
1397 if (modular) {
1398 // Modules can't be compiled for multiple contexts. We need to create the real context now.
1399 auto& mContext = impl->moduleContext.emplace(isolate->getApi().newContext(lock,
1400 {
1401 .newModuleRegistry = impl->getNewModuleRegistry(),
1402 .schemaLoader = getSchemaLoader(),
1403 }));
1404 mContext->enableWarningOnSpecialEvents();
1405 context = mContext.getHandle(lock);
1406 recordedLock.setupContext(context);
1407 } else {
1408 // Although we're going to compile a script independent of context, V8 requires that
1409 // there be an active context, otherwise it will segfault, I guess. So we create a
1410 // dummy context. (Undocumented, as usual.)
1411 context =
1412 v8::Context::New(lock.v8Isolate, nullptr, v8::ObjectTemplate::New(lock.v8Isolate));
1413 // We need to set the highest used index in every context we create to be a nullptr
1414 // This is because we might later on call GetAlignedPointerFromEmbedderData which fails with
1415 // a fatal error if the array is smaller than the given index.
1416 jsg::setAlignedPointerInEmbedderData(
1417 context, jsg::ContextPointerSlot::MAX_POINTER_SLOT, nullptr);
1418 }
1419 
1420 JSG_WITHIN_CONTEXT_SCOPE(lock, context, [&](jsg::Lock& js) {
1421 // const_cast OK because we hold the isolate lock.
1422 Worker::Isolate& lockedWorkerIsolate = const_cast<Isolate&>(*isolate);
1423 
1424 if (logNewScript) {
1425 // HACK: Log a message indicating that a new script was loaded. This is used only when the
1426 // inspector is enabled. We want to do this immediately after the context is created,
1427 // before the user gets a chance to modify the behavior of the console, which if they
1428 // did, we'd then need to be more careful to apply time limits and such.
1429 lockedWorkerIsolate.logMessage(lock, static_cast<uint16_t>(cdp::LogType::WARNING),
1430 "Script modified; context reset.");
1431 }
1432 
1433 // We need to register this context with the inspector, otherwise errors won't be
1434 // reported. But we want it to be un-registered as soon as the script has been
1435 // compiled, otherwise the inspector will end up with multiple contexts active which
1436 // is very confusing for the user (since they'll have to select from the drop-down
1437 // which context to use).
1438 //
1439 // (For modules, the context was already registered by `setupContext()`, above.
1440 KJ_IF_SOME(i, isolate->impl->inspector) {
1441 if (!modular) {
1442 i.get()->contextCreated(
1443 v8_inspector::V8ContextInfo(context, 1, jsg::toInspectorStringView("Compiler")));
1444 }
1445 } else {
1446 } // Here to squash a compiler warning
1447 KJ_DEFER({
1448 if (!modular) {
1449 KJ_IF_SOME(i, isolate->impl->inspector) {
1450 i.get()->contextDestroyed(context);
1451 } else {
1452 } // Here to squash a compiler warning
1453 }
1454 });
1455 
1456 v8::TryCatch catcher(lock.v8Isolate);
1457 ExceptionOrDuration limitErrorOrTime = 0 * kj::NANOSECONDS;
1458 
1459 try {
1460 try {
1461 KJ_SWITCH_ONEOF(source.variant) {
1462 KJ_CASE_ONEOF(script, ScriptSource) {
1463 // This path is used for the older, service worker syntax workers.
1464 
1465 if (script.capnpSchemas.size() > 0) {
1466 // const_cast OK because we hold the isolate lock.
1467 auto& schemaLoader = const_cast<capnp::SchemaLoader&>(getSchemaLoader());
1468 for (auto node: script.capnpSchemas) {
1469 schemaLoader.load(node);
1470 }
1471 }
1472 
1473 impl->globals =
1474 isolate->getApi().compileServiceWorkerGlobals(lock, script, *isolate);
1475 
1476 {
1477 // It's unclear to me if CompileUnboundScript() can get trapped in any
1478 // infinite loops or excessively-expensive computation requiring a time
1479 // limit. We'll go ahead and apply a time limit just to be safe. Don't
1480 // add it to the rollover bank, though.
1481 auto limitScope =
1482 isolate->getLimitEnforcer().enterStartupJs(lock, limitErrorOrTime);
1483 impl->unboundScriptOrMainModule =
1484 jsg::NonModuleScript::compile(lock, script.mainScript, script.mainScriptName);
1485 }
1486 }
1487 
1488 KJ_CASE_ONEOF(modulesSource, ModulesSource) {
1489 // This path is used for the new ESM worker syntax.
1490 
1491 if (modulesSource.capnpSchemas.size() > 0) {
1492 // const_cast OK because we hold the isolate lock.
1493 auto& schemaLoader = const_cast<capnp::SchemaLoader&>(getSchemaLoader());
1494 for (auto node: modulesSource.capnpSchemas) {
1495 schemaLoader.load(node);
1496 }
1497 }
1498 
1499 if (!isolate->getApi().getFeatureFlags().getNewModuleRegistry()) {
1500 kj::Own<void> limitScope;
1501 if (modulesSource.isPython) {
1502 limitScope =
1503 isolate->getLimitEnforcer().enterStartupPython(js, limitErrorOrTime);
1504 } else {
1505 limitScope = isolate->getLimitEnforcer().enterStartupJs(js, limitErrorOrTime);
1506 }
1507 impl->configureDynamicImports(lock, *jsg::ModuleRegistry::from(lock));
1508 isolate->getApi().compileModules(
1509 lock, modulesSource, *isolate, kj::mv(artifacts), parentSpan.addRef());
1510 }
1511 impl->unboundScriptOrMainModule = kj::Path::parse(modulesSource.mainModule);
1512 }
1513 }
1514 
1515 parseMetrics->done();
1516 } catch (const kj::Exception& e) {
1517 lock.throwException(e.clone());
1518 // lock.throwException() here will throw a jsg::JsExceptionThrown which we catch
1519 // in the outer try/catch.
1520 }
1521 } catch (const jsg::JsExceptionThrown&) {
1522 reportStartupError(id, lock, isolate->impl->inspector, isolate->getLimitEnforcer(),
1523 kj::mv(limitErrorOrTime), catcher, errorReporter, impl->permanentException,
1524 parentSpan.addRef(), dynamicEnvBuilder != kj::none);
1525 }
1526 });
1527 });
1528 });
1529}
1530 
1531void Worker::Script::installVirtualFileSystemOnContext(v8::Local<v8::Context> context) const {
1532 jsg::setAlignedPointerInEmbedderData(context, jsg::ContextPointerSlot::VIRTUAL_FILE_SYSTEM,
1533 const_cast<VirtualFileSystem*>(impl->vfs.get()));
1534}
1535 
1536const capnp::SchemaLoader& Worker::Script::getSchemaLoader() const {
1537 KJ_IF_SOME(moduleRegistry, impl->maybeNewModuleRegistry) {
1538 return moduleRegistry->getSchemaLoader();
1539 } else {
1540 return *KJ_ASSERT_NONNULL(impl->maybeSchemaLoader);
1541 }
1542}
1543 
1544kj::Own<const Worker::Isolate::WeakIsolateRef> Worker::Isolate::getWeakRef() const {
1545 return weakIsolateRef->addRef();
1546}
1547 
1548kj::StringPtr Worker::Isolate::getUuid() const {
1549 // As of this writing, getUuid() is only used by actors, for metrics. We don't want to bother
1550 // generating it if not used. The call site does not have nor want an isolate lock, so we use a
1551 // kj::Lazy to make initialization thread-safe.
1552 return impl->uuid.get(
1553 [](kj::SpaceFor<kj::String>& space) { return space.construct(randomUUID(kj::none)); });
1554}
1555 
1556Worker::Isolate::~Isolate() noexcept(false) {
1557 metrics->teardownStarted();
1558 
1559 // Update the isolate stats one last time to make sure we're accurate for cleanup in
1560 // `evicted()`.
1561 limitEnforcer->reportMetrics(*metrics);
1562 
1563 metrics->evicted();
1564 weakIsolateRef->invalidate();
1565 // The cpuLimitNearlyExceededCallback may hold references to objects owned by the isolate and
1566 // their destructors need the isolate to still exist. So destroy them before we destroy the
1567 // isolate.
1568 *cpuLimitNearlyExceededCallback.lockExclusive() = kj::none;
1569 
1570 // Make sure to destroy things under lock. This lock should never be contended since the isolate
1571 // is about to be destroyed, but we have to take the lock in order to enter the isolate.
1572 // It's also important that we lock one last time, in order to destroy any remaining workers in
1573 // worker destruction queue.
1574 jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) {
1575 Isolate::Impl::Lock recordedLock(*this, Worker::Lock::TakeSynchronously(kj::none), stackScope);
1576 metrics->teardownLockAcquired();
1577 auto inspector = kj::mv(impl->inspector);
1578 auto dropTraceAsyncContextKey = kj::mv(traceAsyncContextKey);
1579 auto dropUserTraceAsyncContextKey = kj::mv(userTraceAsyncContextKey);
1580 // The Rust Realm must be dropped under lock since Realm::drop() accesses V8 globals
1581 // and calls drop functions that may interact with V8.
1582 auto dropRealm = kj::mv(impl->realm);
1583 
1584 // Release all tracked WASM instance entries while V8 is still alive. Each entry holds a
1585 // shared_ptr<v8::BackingStore> whose destructor who needs the isolate to still be alive.
1586 // This is analogous to the cpuTimeLimitNearlyExceededCallback detaching above ^^^
1587 limitEnforcer->getTrackedWasmInstances().clear(*recordedLock.lock);
1588 });
1589}
1590 
1591Worker::Script::~Script() noexcept(false) {
1592 // Make sure to destroy things under lock.
1593 // TODO(perf): It could make sense to try to obtain an async lock before destroying a script if
1594 // multiple scripts are co-located in the same isolate. As of this writing, that doesn't happen
1595 // except in preview. In any case, Scripts are destroyed in the GC thread, where we don't care
1596 // too much about lock latency.
1597 jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) {
1598 Isolate::Impl::Lock recordedLock(
1599 *isolate, Worker::Lock::TakeSynchronously(kj::none), stackScope);
1600 KJ_IF_SOME(c, impl->moduleContext) {
1601 recordedLock.disposeContext(kj::mv(c));
1602 }
1603 impl = nullptr;
1604 });
1605}
1606 
1607const Worker::Isolate& Worker::Isolate::from(jsg::Lock& js) {
1608 auto ptr = js.v8Isolate->GetData(jsg::SET_DATA_ISOLATE);
1609 KJ_ASSERT(ptr != nullptr);
1610 return *static_cast<const Worker::Isolate*>(ptr);
1611}
1612 
1613bool Worker::Isolate::Impl::Lock::checkInWithLimitEnforcer(Worker::Isolate& isolate) {
1614 shouldReportIsolateMetrics = true;
1615 return limitEnforcer.exitJs(*lock);
1616}
1617 
1618kj::Maybe<kj::Function<void(void)>> Worker::Isolate::getCpuLimitNearlyExceededCallback() const {
1619 auto lock = cpuLimitNearlyExceededCallback.lockExclusive();
1620 KJ_IF_SOME(cb, *lock) {
1621 return cb.reference();
1622 }
1623 return kj::none;
1624}
1625 
1626void Worker::Isolate::setCpuLimitNearlyExceededCallback(kj::Function<void(void)> cb) const {
1627 auto lock = cpuLimitNearlyExceededCallback.lockExclusive();
1628 // Make sure we don't reassign the callback so we don't invalidate references we've passed out.
1629 if (*lock == kj::none) {
1630 *lock = kj::mv(cb);
1631 return;
1632 }
1633 kj::throwRecoverableException(KJ_EXCEPTION(
1634 FAILED, "Python Workers Internal Error: CpuLimitNearlyExceededCallback already set"));
1635}
1636 
1637void Worker::Isolate::registerTrackedWasmInstance(jsg::Lock& js,
1638 v8::Local<v8::Object> instance,
1639 kj::Array<kj::byte> memory,
1640 kj::Maybe<uint32_t> signalOffset,
1641 kj::Maybe<uint32_t> terminatedOffset) const {
1642 // Register the WASM module for receiving shutdown signals. The signal handler will
1643 // iterate the list unconditionally when CPU time is nearly exhausted.
1644 KJ_IF_SOME(entry,
1645 limitEnforcer->getTrackedWasmInstances().registerSignal(
1646 js, kj::mv(memory), signalOffset, terminatedOffset)) {
1647 // Set up a weak reference to the instance. When V8 collects it, the handle becomes
1648 // empty and the GC prologue filter removes the entry, releasing the strong memory ref.
1649 entry.instanceRef.Reset(js.v8Isolate, instance);
1650 entry.instanceRef.SetWeak();
1651 }
1652}
1653 
1654// EW-1319: Set WebAssembly.Module @@HasInstance
1655//
1656// The instanceof operator can be changed by setting the @@HasInstance method
1657// on the object, https://tc39.es/ecma262/#sec-instanceofoperator.
1658void setWebAssemblyModuleHasInstance(jsg::Lock& lock, v8::Local<v8::Context> context) {
1659 JSG_WITHIN_CONTEXT_SCOPE(lock, context, [&](jsg::Lock& lock) {
1660 auto instanceof = [](const v8::FunctionCallbackInfo<v8::Value>& info) {
1661 jsg::Lock::from(info.GetIsolate()).withinHandleScope([&] {
1662 info.GetReturnValue().Set(info[0]->IsWasmModuleObject());
1663 });
1664 };
1665 v8::Local<v8::Function> function = jsg::check(v8::Function::New(context, instanceof));
1666 
1667 auto webAssembly =
1668 KJ_ASSERT_NONNULL(lock.global().get(lock, "WebAssembly").tryCast<jsg::JsObject>());
1669 auto module = KJ_ASSERT_NONNULL(webAssembly.get(lock, "Module").tryCast<jsg::JsObject>());
1670 
1671 jsg::check(v8::Local<v8::Object>(module)->DefineOwnProperty(
1672 context, v8::Symbol::GetHasInstance(lock.v8Isolate), function));
1673 });
1674}
1675 
1676// Installs a shim around WebAssembly.instantiate and WebAssembly.Instance that hooks into the
1677// shutdown signal if it exists
1678void shimWebAssemblyInstantiate(jsg::Lock& lock, v8::Local<v8::Context> context) {
1679 // We need to enter the context because this function compiles and executes JavaScript via
1680 // v8::Script::Compile/Run. setupContext() is called before JSG_WITHIN_CONTEXT_SCOPE, so the
1681 // context is not yet entered at this point.
1682 v8::Context::Scope contextScope(context);
1683 
1684 // Create a C++ callback that the JS shims call to register a {instance, memory, signalOffset,
1685 // terminatedOffset} tuple.
1686 // __registerTrackedWasmInstance(instance: WebAssembly.Instance,
1687 // memory: WebAssembly.Memory, signalOffset: number,
1688 // terminatedOffset: number)
1689 // signalOffset or terminatedOffset may be -1, indicating the corresponding export is absent.
1690 auto registerCb = [](const v8::FunctionCallbackInfo<v8::Value>& info) {
1691 auto& js = jsg::Lock::from(info.GetIsolate());
1692 js.withinHandleScope([&] {
1693 if (info.Length() < 4 || !info[0]->IsObject() || !info[1]->IsWasmMemoryObject() ||
1694 !info[2]->IsNumber() || !info[3]->IsNumber()) {
1695 js.v8Isolate->ThrowException(
1696 js.str("registerTrackedWasmInstance: expected "
1697 "(WebAssembly.Instance, WebAssembly.Memory, number, number)"_kj));
1698 return;
1699 }
1700 auto instance = info[0].As<v8::Object>();
1701 auto memory = info[1].As<v8::WasmMemoryObject>();
1702 // signalOffset is -1 when __instance_signal was not exported.
1703 auto signalRaw = info[2].As<v8::Number>()->Value();
1704 kj::Maybe<uint32_t> signalOffset;
1705 if (signalRaw >= 0) {
1706 signalOffset = static_cast<uint32_t>(signalRaw);
1707 }
1708 // terminatedOffset is -1 when __instance_terminated was not exported.
1709 auto terminatedRaw = info[3].As<v8::Number>()->Value();
1710 kj::Maybe<uint32_t> terminatedOffset;
1711 if (terminatedRaw >= 0) {
1712 terminatedOffset = static_cast<uint32_t>(terminatedRaw);
1713 }
1714 auto backingStore = memory->Buffer()->GetBackingStore();
1715 auto wasmMemory =
1716 kj::arrayPtr(static_cast<kj::byte*>(backingStore->Data()), backingStore->ByteLength())
1717 .attach(kj::mv(backingStore));
1718 KJ_IF_SOME(e, kj::runCatchingExceptions([&] {
1719 Worker::Isolate::from(js).registerTrackedWasmInstance(
1720 js, instance, kj::mv(wasmMemory), signalOffset, terminatedOffset);
1721 })) {
1722 js.v8Isolate->ThrowException(js.exceptionToJs(kj::mv(e)).getHandle(js));
1723 }
1724 });
1725 };
1726 auto registerFn = jsg::check(v8::Function::New(context, registerCb));
1727 
1728 // Build the shim in JavaScript. It wraps both WebAssembly.instantiate (async) and
1729 // WebAssembly.Instance (sync constructor).
1730 auto shimScript =
1731 jsg::NonModuleScript::compile(lock, WASM_INSTANTIATE_SHIM, "wasm-instantiate-shim.js"_kj);
1732 auto shimFn = KJ_ASSERT_NONNULL(shimScript.runAndReturn(lock).tryCast<jsg::JsFunction>());
1733 
1734 // Call the factory โ€” it mutates `WebAssembly` in place.
1735 shimFn.call(lock, lock.global(), jsg::JsFunction(registerFn));
1736}
1737 
1738void Worker::setupContext(
1739 jsg::Lock& lock, v8::Local<v8::Context> context, const LoggingOptions& loggingOptions) {
1740 // Set WebAssembly.Module @@HasInstance
1741 setWebAssemblyModuleHasInstance(lock, context);
1742 
1743 // Shim WebAssembly.instantiate to detect modules exporting "__instance_signal".
1744 if (util::Autogate::isEnabled(util::AutogateKey::WASM_SHUTDOWN_SIGNAL_SHIM)) {
1745 shimWebAssemblyInstantiate(lock, context);
1746 }
1747 
1748 // We replace the default V8 console.log(), etc. methods, to give the worker access to
1749 // logged content, and log formatted values to stdout/stderr locally.
1750 auto global = context->Global();
1751 auto consoleStr = jsg::v8StrIntern(lock.v8Isolate, "console");
1752 auto console = jsg::check(global->Get(context, consoleStr)).As<v8::Object>();
1753 
1754 auto setHandler = [&](const char* method, LogLevel level) {
1755 auto methodStr = jsg::v8StrIntern(lock.v8Isolate, method);
1756 v8::Global<v8::Function> original(
1757 lock.v8Isolate, jsg::check(console->Get(context, methodStr)).As<v8::Function>());
1758 
1759 auto f = lock.wrapSimpleFunction(context,
1760 [loggingOptions, level, original = kj::mv(original)](
1761 jsg::Lock& js, const v8::FunctionCallbackInfo<v8::Value>& info) {
1762 handleLog(js, loggingOptions, level, original, info);
1763 });
1764 jsg::check(console->Set(context, methodStr, f));
1765 };
1766 
1767 setHandler("debug", LogLevel::DEBUG_);
1768 setHandler("error", LogLevel::ERROR);
1769 setHandler("info", LogLevel::INFO);
1770 setHandler("log", LogLevel::LOG);
1771 setHandler("warn", LogLevel::WARN);
1772}
1773// =======================================================================================
1774 
1775namespace {
1776kj::Maybe<jsg::JsObject> tryResolveMainModule(jsg::Lock& js,
1777 const kj::Path& mainModule,
1778 jsg::JsContext<api::ServiceWorkerGlobalScope>& jsContext,
1779 const Worker::Script& script,
1780 ExceptionOrDuration& limitErrorOrTime) {
1781 kj::Own<void> limitScope;
1782 if (script.isPython()) {
1783 limitScope = script.getIsolate().getLimitEnforcer().enterStartupPython(js, limitErrorOrTime);
1784 } else {
1785 limitScope = script.getIsolate().getLimitEnforcer().enterStartupJs(js, limitErrorOrTime);
1786 }
1787 
1788 KJ_DEFER({
1789 if (limitErrorOrTime.is<kj::Exception>()) {
1790 // If we hit the limit in PerformMicrotaskCheckpoint() we may not have actually
1791 // thrown an exception.
1792 throw jsg::JsExceptionThrown();
1793 }
1794 });
1795 
1796 // Before resolving the main module, if both nodejs_compat_v2 and the new
1797 // module registry are enabled, let's pre-resolve the process and buffer modules.
1798 // Why? Great question! Resolving these modules synchronously causes the microtask
1799 // queue to be pumped, which we don't actually want to do while resolving the main
1800 // module until we are ready. Both process and buffer are exposed via globalThis
1801 // when the nodejs_compat_v2 flag is used, and if the top-level scope is accessing
1802 // either globalThis.process or globalThis.buffer, then we need to make sure that
1803 // the modules are already resolved so we don't pump the microtask queue while
1804 // synchronously accessing those globals. Resolving them here ensures that they are
1805 // ready to go before we begin evaluating the main module.
1806 auto featureFlags = FeatureFlags::get(js);
1807 if (featureFlags.getNodeJsCompatV2() && featureFlags.getNewModuleRegistry()) {
1808 JSG_REQUIRE_NONNULL(js.resolveModule("node:process", jsg::RequireEsm::YES), Error,
1809 "Failed to initialize node:process module");
1810 JSG_REQUIRE_NONNULL(js.resolveModule("node:buffer", jsg::RequireEsm::YES), Error,
1811 "Failed to initialize node:buffer module");
1812 }
1813 
1814 // When enable_nodejs_global_timers is enabled, load the module that makes all 6 timer
1815 // functions (setTimeout, setInterval, clearTimeout, clearInterval, setImmediate,
1816 // clearImmediate) available on globalThis as Node.js-compatible versions from node:timers.
1817 if (featureFlags.getEnableNodejsGlobalTimers()) {
1818 JSG_REQUIRE_NONNULL(js.resolveInternalModule("node-internal:internal_timers_global_override"),
1819 Error, "Failed to initialize node-internal:internal_timers_global_override module");
1820 }
1821 
1822 return js.resolveModule(mainModule.toString(false), jsg::RequireEsm::YES);
1823}
1824} // anonymous namespace
1825 
1826Worker::Worker(kj::Own<const Script> scriptParam,
1827 kj::Own<WorkerObserver> metricsParam,
1828 kj::FunctionParam<void(jsg::Lock& lock,
1829 const Api& api,
1830 v8::Local<v8::Object> target,
1831 v8::Local<v8::Object> ctxExports)> compileBindings,
1832 IsolateObserver::StartType startType,
1833 SpanParent parentSpan,
1834 LockType lockType,
1835 kj::Maybe<ValidationErrorReporter&> errorReporter,
1836 kj::Maybe<kj::Duration&> startupTime)
1837 : script(kj::mv(scriptParam)),
1838 metrics(kj::mv(metricsParam)),
1839 impl(kj::heap<Impl>()) {
1840 // Enter/lock isolate.
1841 jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) {
1842 Isolate::Impl::Lock recordedLock(*script->isolate, lockType, stackScope);
1843 auto& lock = *recordedLock.lock;
1844 
1845 // If we throw an exception, it's important that `impl` is destroyed under lock.
1846 KJ_ON_SCOPE_FAILURE({
1847 auto implToDestroy = kj::mv(impl);
1848 KJ_IF_SOME(c, implToDestroy->context) {
1849 recordedLock.disposeContext(kj::mv(c));
1850 } else {
1851 // Else block to avoid dangling else clang warning.
1852 }
1853 });
1854 
1855 auto maybeMakeSpan = [&](auto operationName) -> SpanBuilder {
1856 auto span = parentSpan.newChild(kj::mv(operationName));
1857 if (span.isObserved()) {
1858 span.setTag("truncated_script_id"_kjc, truncateScriptId(script->getId()));
1859 }
1860 return span;
1861 };
1862 
1863 auto currentSpan = maybeMakeSpan("lw:new_startup_metrics"_kjc);
1864 
1865 auto startupMetrics = metrics->startup(startType);
1866 
1867 currentSpan = maybeMakeSpan("lw:new_context"_kjc);
1868 
1869 // Create a stack-allocated handle scope.
1870 lock.withinHandleScope([&] {
1871 jsg::JsContext<api::ServiceWorkerGlobalScope>* jsContext;
1872 
1873 KJ_IF_SOME(c, script->impl->moduleContext) {
1874 // Use the shared context from the script.
1875 // const_cast OK because guarded by `lock`.
1876 jsContext = const_cast<jsg::JsContext<api::ServiceWorkerGlobalScope>*>(&c);
1877 currentSpan.setTag("module_context"_kjc, true);
1878 } else {
1879 // Create a new context.
1880 jsContext = &this->impl->context.emplace(script->isolate->getApi().newContext(lock,
1881 {
1882 .newModuleRegistry = script->impl->getNewModuleRegistry(),
1883 .schemaLoader = script->getSchemaLoader(),
1884 }));
1885 }
1886 
1887 v8::Local<v8::Context> context = KJ_REQUIRE_NONNULL(jsContext).getHandle(lock);
1888 
1889 // Install the virtual file system on the context. Keep in mind that for service
1890 // worker style workers, the Script may be shared between multiple Workers, even
1891 // across different accounts. Currently, the internal state of the VFS does not
1892 // contain any account-specific or worker-specific state so this is OK for now.
1893 // The VFS would contain the script files only and any temporary files created
1894 // within the context of a worker are always stored in temporary space attached
1895 // to the IoContext or the current execution context. If we extend these capabilities
1896 // in the future, we may need to revisit this. For modular workers, this is not
1897 // an issue since each Worker gets its own Script instance.
1898 script->installVirtualFileSystemOnContext(context);
1899 
1900 if (!script->modular) {
1901 recordedLock.setupContext(context);
1902 }
1903 
1904 if (script->impl->unboundScriptOrMainModule == nullptr) {
1905 // Script failed to parse. Act as if the script was empty -- i.e. do nothing.
1906 impl->permanentException =
1907 script->impl->permanentException.map([](auto& e) { return e.clone(); });
1908 return;
1909 }
1910 
1911 // Enter the context for compiling and running the script.
1912 JSG_WITHIN_CONTEXT_SCOPE(lock, context, [&](jsg::Lock& js) {
1913 v8::TryCatch catcher(lock.v8Isolate);
1914 ExceptionOrDuration limitErrorOrTime = 0 * kj::NANOSECONDS;
1915 
1916 try {
1917 try {
1918 currentSpan = maybeMakeSpan("lw:globals_instantiation"_kjc);
1919 
1920 v8::Local<v8::Object> bindingsScope;
1921 if (script->isModular()) {
1922 // Use `env` variable.
1923 bindingsScope = v8::Object::New(lock.v8Isolate);
1924 if (!FeatureFlags::get(js).getDisableImportableEnv()) {
1925 lock.setWorkerEnv(lock.v8Ref(bindingsScope));
1926 }
1927 } else {
1928 // Use global-scope bindings.
1929 bindingsScope = context->Global();
1930 }
1931 
1932 // Load globals.
1933 // const_cast OK because we hold the lock.
1934 for (auto& global: const_cast<Script&>(*script).impl->globals) {
1935 lock.v8Set(bindingsScope, global.name, global.value);
1936 }
1937 
1938 v8::Local<v8::Object> ctxExports = v8::Object::New(lock.v8Isolate);
1939 
1940 compileBindings(lock, script->isolate->getApi(), bindingsScope, ctxExports);
1941 
1942 // Execute script.
1943 currentSpan = maybeMakeSpan("lw:top_level_execution"_kjc);
1944 
1945 // Ensure that our worker top-level bootstrap has a temporary directory
1946 // storage scope. This is used to store temporary files created within
1947 // the top-level evaluation of the worker. With this instantiated on
1948 // the stack, temporary files will be cleaned up when the scope is
1949 // destroyed, which means any temporary files created in the top-level
1950 // evaluation will *not* be available to the worker after the top-level
1951 // evaluation is complete.
1952 TmpDirStoreScope tmpDirStoreScope;
1953 
1954 // We allow eval and new Function() during startup, becaues startup time is entirely
1955 // deterministic, so we can easily reproduce the input to eval() by just running the
1956 // worker again. We do not allow eval() at runtime because we need to have a record of
1957 // all code that executes in production for forensic purposes, and at runtime the input
1958 // to eval() could have come from a remote source on which we don't have a record.
1959 js.setAllowEval(FeatureFlags::get(js).getAllowEvalDuringStartup());
1960 KJ_DEFER(js.setAllowEval(false));
1961 
1962 KJ_SWITCH_ONEOF(script->impl->unboundScriptOrMainModule) {
1963 KJ_CASE_ONEOF(unboundScript, jsg::NonModuleScript) {
1964 auto limitScope =
1965 script->isolate->getLimitEnforcer().enterStartupJs(lock, limitErrorOrTime);
1966 unboundScript.run(lock);
1967 // Flush microtasks enqueued during top-level script evaluation.
1968 // Without this flush, microtasks (e.g. promise continuations from async
1969 // initialization) remain on the per-isolate microtask queue and can leak across
1970 // V8 contexts when multiple Workers share an isolate (same script, different
1971 // zones). The leaked microtasks then execute under the wrong IoContext, making
1972 // things go boom.
1973 lock.runMicrotasks();
1974 }
1975 KJ_CASE_ONEOF(mainModule, kj::Path) {
1976 KJ_IF_SOME(ns,
1977 tryResolveMainModule(lock, mainModule, *jsContext, *script, limitErrorOrTime)) {
1978 impl->env = lock.v8Ref(bindingsScope.As<v8::Value>());
1979 impl->ctxExports = lock.v8Ref(ctxExports.As<v8::Value>());
1980 
1981 if (!FeatureFlags::get(js).getDisableImportableEnv()) {
1982 lock.setWorkerExports(lock.v8Ref(ctxExports));
1983 }
1984 
1985 auto& api = script->isolate->getApi();
1986 auto handlers = api.unwrapExports(lock, ns);
1987 auto entrypointClasses = api.getEntrypointClasses(lock);
1988 
1989 for (auto& handler: handlers.fields) {
1990 KJ_SWITCH_ONEOF(handler.value) {
1991 KJ_CASE_ONEOF(obj, api::ExportedHandler) {
1992 obj.env = lock.v8Ref(bindingsScope.As<v8::Value>());
1993 // Historically, non-class-based handlers reused the same ctx object for all requests.
1994 // This was an accident, but some Workers depend on it.
1995 // Newer worker with the unique_ctx_per_invocation will allocate a new ctx for every request.
1996 obj.ctx = js.alloc<api::ExecutionContext>(lock, jsg::JsValue(ctxExports));
1997 
1998 // Python Workers append all durable objects, worker entrypoint and workflow
1999 // entrypoint classes in the pythonEntrypoints named export.
2000 bool isPythonWorker = FeatureFlags::get(js).getPythonWorkers();
2001 if (handler.name == "pythonEntrypoints" && isPythonWorker) {
2002 auto handle = obj.self.getHandle(js);
2003 auto dict = js.toDict(handle);
2004 for (auto& field: dict.fields) {
2005 auto unwrapped = api.unwrapExport(lock, field.value);
2006 KJ_SWITCH_ONEOF(unwrapped) {
2007 KJ_CASE_ONEOF(cls, EntrypointClass) {
2008 processEntrypointClass(
2009 js, kj::mv(cls), entrypointClasses, kj::mv(field.name));
2010 }
2011 KJ_CASE_ONEOF(obj, api::ExportedHandler) {
2012 KJ_FAIL_ASSERT("Expected EntrypointClass");
2013 }
2014 }
2015 }
2016 } else {
2017 impl->namedHandlers.insert(kj::mv(handler.name), kj::mv(obj));
2018 }
2019 }
2020 KJ_CASE_ONEOF(cls, EntrypointClass) {
2021 processEntrypointClass(
2022 js, kj::mv(cls), entrypointClasses, kj::mv(handler.name));
2023 }
2024 }
2025 }
2026 } else {
2027 JSG_FAIL_REQUIRE(TypeError, "Main module name is not present in bundle.");
2028 }
2029 }
2030 }
2031 
2032 KJ_IF_SOME(s, startupTime) {
2033 KJ_SWITCH_ONEOF(limitErrorOrTime) {
2034 KJ_CASE_ONEOF(startupTimeElapsed, kj::Duration) {
2035 s = startupTimeElapsed;
2036 }
2037 KJ_CASE_ONEOF(limitError, kj::Exception) {}
2038 }
2039 } else {
2040 }
2041 startupMetrics->done();
2042 } catch (const kj::Exception& e) {
2043 lock.throwException(e.clone());
2044 // lock.throwException() here will throw a jsg::JsExceptionThrown which we catch
2045 // in the outer try/catch.
2046 }
2047 } catch (const jsg::JsExceptionThrown&) {
2048 reportStartupError(script->id, lock, script->isolate->impl->inspector,
2049 script->isolate->getLimitEnforcer(), kj::mv(limitErrorOrTime), catcher, errorReporter,
2050 impl->permanentException, currentSpan, script->getDynamicEnvBuilder() != kj::none);
2051 }
2052 });
2053 
2054 // Reset this back to its default after startup execution
2055 // Leaving it on comes at the expense of collecting stack traces for all thrown exceptions
2056 // Ref: https://github.com/cloudflare/workerd/issues/5332
2057 if (script->isolate->impl->inspector == kj::none) {
2058 lock.v8Isolate->SetCaptureStackTraceForUncaughtExceptions(false);
2059 }
2060 });
2061 });
2062}
2063 
2064Worker::~Worker() noexcept(false) {
2065 metrics->teardownStarted();
2066 
2067 auto& isolateImpl = *script->getIsolate().impl;
2068 auto lock = isolateImpl.workerDestructionQueue.lockExclusive();
2069 
2070 // Previously, this metric meant the isolate lock. We might as well make it mean the worker
2071 // destruction queue lock now to verify it is much less-contended than the isolate lock.
2072 metrics->teardownLockAcquired();
2073 
2074 // Defer destruction of our V8 objects, in particular our jsg::Context, which requires some
2075 // finalization.
2076 lock->push(kj::mv(impl));
2077}
2078 
2079void Worker::processEntrypointClass(jsg::Lock& js,
2080 EntrypointClass cls,
2081 EntrypointClasses entrypointClasses,
2082 kj::String handlerName) {
2083 js.withinHandleScope([&]() {
2084 jsg::JsObject handle(KJ_ASSERT_NONNULL(cls.tryGetHandle(js.v8Isolate)));
2085 
2086 for (;;) {
2087 if (handle == entrypointClasses.durableObject) {
2088 impl->actorClasses.insert(kj::mv(handlerName),
2089 ActorClassInfo{
2090 .cls = kj::mv(cls),
2091 .missingSuperclass = false,
2092 });
2093 return;
2094 } else if (handle == entrypointClasses.workerEntrypoint) {
2095 impl->statelessClasses.insert(kj::mv(handlerName), kj::mv(cls));
2096 return;
2097 } else if (handle == entrypointClasses.workflowEntrypoint) {
2098 impl->workflowClasses.insert(kj::mv(handlerName), kj::mv(cls));
2099 return;
2100 }
2101 
2102 handle = KJ_UNWRAP_OR(handle.getPrototype(js).tryCast<jsg::JsObject>(), {
2103 // Reached end of prototype chain.
2104 
2105 // For historical reasons, we assume a class is a Durable Object
2106 // class if it doesn't inherit anything.
2107 // TODO(someday): Log a warning suggesting extending DurableObject.
2108 // TODO(someday): Introduce a compat flag that makes this required.
2109 impl->actorClasses.insert(kj::mv(handlerName),
2110 ActorClassInfo{
2111 .cls = kj::mv(cls),
2112 .missingSuperclass = true,
2113 });
2114 return;
2115 });
2116 }
2117 });
2118}
2119 
2120void Worker::handleLog(jsg::Lock& js,
2121 const LoggingOptions& loggingOptions,
2122 LogLevel level,
2123 const v8::Global<v8::Function>& original,
2124 const v8::FunctionCallbackInfo<v8::Value>& info) {
2125 // Call original V8 implementation so messages sent to connected inspector if any
2126 auto context = js.v8Context();
2127 int length = info.Length();
2128 // to pass additional arguments from this function to js' `formatLog` we add arguments to the end
2129 // of the arguments vector, then in formatLog we `pop` these from the vector.
2130 // 3 is just the number of args we currently pass.
2131 v8::LocalVector<v8::Value> args(js.v8Isolate, length + 3);
2132 for (auto i: kj::zeroTo(length)) args[i] = info[i];
2133 jsg::check(original.Get(js.v8Isolate)->Call(context, info.This(), length, args.data()));
2134 
2135 // The TryCatch is initialized here to catch cases where the v8 isolate's execution is
2136 // terminating, usually as a result of an infinite loop. We need to perform the initialization
2137 // here because `message` is called multiple times.
2138 v8::TryCatch tryCatch(js.v8Isolate);
2139 auto message = [&]() {
2140 int length = info.Length();
2141 kj::Vector<kj::String> stringified(length);
2142 for (auto i: kj::zeroTo(length)) {
2143 auto arg = info[i];
2144 // serializeJson and v8::Value::ToString can throw JS exceptions
2145 // (e.g. for recursive objects) so we eat them here, to ensure logging and non-logging code
2146 // have the same exception behavior.
2147 if (!tryCatch.CanContinue()) {
2148 stringified.add(kj::str("{}"));
2149 break;
2150 }
2151 // The following code checks the `arg` to see if it should be serialised to JSON.
2152 //
2153 // We use the following criteria: if arg is null, a number, a boolean, an array, a string, an
2154 // object or it defines a `toJSON` property that is a function, then the arg gets serialised
2155 // to JSON.
2156 //
2157 // Otherwise we stringify the argument.
2158 js.withinHandleScope([&] {
2159 auto context = js.v8Context();
2160 bool shouldSerialiseToJson = false;
2161 if (arg->IsNull() || arg->IsNumber() || arg->IsArray() || arg->IsBoolean() ||
2162 arg->IsString() ||
2163 arg->IsUndefined()) { // This is special cased for backwards compatibility.
2164 shouldSerialiseToJson = true;
2165 }
2166 if (arg->IsObject()) {
2167 v8::Local<v8::Object> obj = arg.As<v8::Object>();
2168 v8::Local<v8::Object> freshObj = v8::Object::New(js.v8Isolate);
2169 
2170 // Determine whether `obj` is constructed using `{}` or `new Object()`. This ensures
2171 // we don't serialise values like Promises to JSON.
2172#if V8_MAJOR_VERSION >= 15 || (V8_MAJOR_VERSION == 14 && V8_MINOR_VERSION >= 7)
2173 if (obj->GetPrototype()->SameValue(freshObj->GetPrototype()) ||
2174 obj->GetPrototype()->IsNull()) {
2175#else
2176 // TODO(cleanup): Remove when unnecessary.
2177 if (obj->GetPrototypeV2()->SameValue(freshObj->GetPrototypeV2()) ||
2178 obj->GetPrototypeV2()->IsNull()) {
2179#endif
2180 shouldSerialiseToJson = true;
2181 }
2182 
2183 // Check if arg has a `toJSON` property which is a function.
2184 auto toJSONStr = jsg::v8StrIntern(js.v8Isolate, "toJSON"_kj);
2185 v8::MaybeLocal<v8::Value> toJSON = obj->GetRealNamedProperty(context, toJSONStr);
2186 if (!toJSON.IsEmpty()) {
2187 if (jsg::check(toJSON)->IsFunction()) {
2188 shouldSerialiseToJson = true;
2189 }
2190 }
2191 }
2192 
2193 if (kj::runCatchingExceptions([&]() {
2194 // On the off chance the the arg is the request.cf object, let's make
2195 // sure we do not log proxied fields here.
2196 if (shouldSerialiseToJson) {
2197 auto s = js.serializeJson(arg);
2198 // serializeJson returns the string "undefined" for some values (undefined,
2199 // Symbols, functions). We remap these values to null to ensure valid JSON output.
2200 if (s == "undefined"_kj) {
2201 stringified.add(kj::str("null"));
2202 } else {
2203 stringified.add(kj::mv(s));
2204 }
2205 } else {
2206 stringified.add(js.serializeJson(jsg::check(arg->ToString(context))));
2207 }
2208 }) != kj::none) {
2209 stringified.add(kj::str("{}"));
2210 };
2211 });
2212 }
2213 return kj::str("[", kj::delimited(stringified, ", "_kj), "]");
2214 };
2215 
2216 // Only check tracing if console.log() was not invoked at the top level.
2217 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
2218 KJ_IF_SOME(tracer, ioContext.getWorkerTracer()) {
2219 auto timestamp = ioContext.now();
2220 tracer.addLog(ioContext.getInvocationSpanContext(), timestamp, level, message());
2221 }
2222 }
2223 
2224 if (loggingOptions.consoleMode == Worker::ConsoleMode::INSPECTOR_ONLY) {
2225 // Lets us dump console.log()s to stdout when running test-runner with --verbose flag, to make
2226 // it easier to debug tests. Note that when --verbose is not passed, KJ_LOG(INFO, ...) will
2227 // not even evaluate its arguments, so `message()` will not be called at all.
2228 KJ_LOG(INFO, "console.log()", message());
2229 } else {
2230 // Write to stdio if allowed by console mode. This is making use of our internal
2231 // built-in implementation of the node:util inspect API.
2232 static const ColorMode COLOR_MODE = permitsColor();
2233#if _WIN32
2234 static bool STDOUT_TTY = _isatty(_fileno(stdout));
2235 static bool STDERR_TTY = _isatty(_fileno(stderr));
2236#else
2237 static bool STDOUT_TTY = isatty(STDOUT_FILENO);
2238 static bool STDERR_TTY = isatty(STDERR_FILENO);
2239#endif
2240 
2241 // Log warnings and errors to stderr
2242 // Always log to stdout when structuredLogging is enabled.
2243 auto useStderr = level >= LogLevel::WARN && !loggingOptions.structuredLogging;
2244 auto fd = useStderr ? stderr : stdout;
2245 auto tty = useStderr ? STDERR_TTY : STDOUT_TTY;
2246 auto colors =
2247 COLOR_MODE == ColorMode::ENABLED || (COLOR_MODE == ColorMode::ENABLED_IF_TTY && tty);
2248 
2249 constexpr auto kSpecifier = "node-internal:internal_inspect"_kj;
2250 auto inspectModule = KJ_ASSERT_NONNULL(js.resolveInternalModule(kSpecifier));
2251 v8::Local<v8::Value> formatLogVal = inspectModule.get(js, "formatLog"_kj);
2252 KJ_ASSERT(formatLogVal->IsFunction());
2253 auto formatLog = formatLogVal.As<v8::Function>();
2254 
2255 auto levelStr = logLevelToString(level);
2256 args[length] = js.boolean(colors);
2257 args[length + 1] = js.boolean(loggingOptions.structuredLogging.toBool());
2258 args[length + 2] = js.strIntern(levelStr);
2259 auto formatted = js.toString(
2260 jsg::check(formatLog->Call(context, js.v8Undefined(), length + 3, args.data())));
2261 fprintf(fd, "%s\n", formatted.cStr());
2262 fflush(fd);
2263 }
2264}
2265 
2266Worker::Lock::TakeSynchronously::TakeSynchronously(kj::Maybe<RequestObserver&> requestParam) {
2267 KJ_IF_SOME(r, requestParam) {
2268 request = &r;
2269 }
2270}
2271 
2272kj::Maybe<RequestObserver&> Worker::Lock::TakeSynchronously::getRequest() {
2273 if (request != nullptr) {
2274 return *request;
2275 }
2276 return kj::none;
2277}
2278 
2279struct Worker::Lock::Impl {
2280 Isolate::Impl::Lock recordedLock;
2281 jsg::Lock& inner;
2282 
2283 Impl(const Worker& worker, LockType lockType, jsg::V8StackScope& stackScope)
2284 : recordedLock(worker.getIsolate(), lockType, stackScope),
2285 inner(*recordedLock.lock) {}
2286};
2287 
2288Worker::Lock::Lock(const Worker& constWorker, LockType lockType, jsg::V8StackScope& stackScope)
2289 : // const_cast OK because we took out a lock.
2290 worker(const_cast<Worker&>(constWorker)),
2291 impl(kj::heap<Impl>(worker, lockType, stackScope)) {
2292 kj::requireOnStack(this, "Worker::Lock MUST be allocated on the stack.");
2293}
2294 
2295Worker::Lock::~Lock() noexcept(false) {
2296 // const_cast OK because we hold -- nay, we *are* -- a lock on the script.
2297 auto& isolate = const_cast<Isolate&>(worker.getIsolate());
2298 if (impl->recordedLock.checkInWithLimitEnforcer(isolate)) {
2299 isolate.disconnectInspector();
2300 }
2301}
2302 
2303void Worker::Lock::requireNoPermanentException() {
2304 KJ_IF_SOME(e, worker.impl->permanentException) {
2305 // Block taking lock when worker failed to start up.
2306 kj::throwFatalException(e.clone());
2307 }
2308}
2309 
2310Worker::Lock::operator jsg::Lock&() {
2311 return impl->inner;
2312}
2313 
2314v8::Isolate* Worker::Lock::getIsolate() {
2315 return impl->inner.v8Isolate;
2316}
2317 
2318v8::Local<v8::Context> Worker::Lock::getContext() {
2319 KJ_IF_SOME(c, worker.impl->context) {
2320 return c.getHandle(impl->inner);
2321 } else KJ_IF_SOME(c, const_cast<Script&>(*worker.script).impl->moduleContext) {
2322 return c.getHandle(impl->inner);
2323 } else {
2324 KJ_UNREACHABLE;
2325 }
2326}
2327 
2328template <typename T>
2329static inline kj::Own<T> fakeOwn(T& ref) {
2330 return kj::Own<T>(&ref, kj::NullDisposer::instance);
2331}
2332 
2333kj::Maybe<kj::Own<api::ExportedHandler>> Worker::Lock::getExportedHandler(
2334 kj::Maybe<kj::StringPtr> name,
2335 kj::Maybe<VersionInfo> versionInfo,
2336 Frankenvalue props,
2337 kj::Maybe<Worker::Actor&> actor,
2338 bool isDynamicDispatch) {
2339 KJ_IF_SOME(a, actor) {
2340 KJ_IF_SOME(h, a.getHandler()) {
2341 return fakeOwn(h);
2342 }
2343 }
2344 
2345 kj::StringPtr n = name.orDefault("default"_kj);
2346 
2347 auto getHandlerFromEntrypointClass =
2348 [&](EntrypointClass& cls) -> kj::Maybe<kj::Own<api::ExportedHandler>> {
2349 jsg::Lock& js = *this;
2350 auto handler = kj::heap(cls(js,
2351 js.alloc<api::ExecutionContext>(js,
2352 jsg::JsValue(KJ_ASSERT_NONNULL(worker.impl->ctxExports).getHandle(js)), props.toJs(js),
2353 kj::mv(versionInfo)),
2354 KJ_ASSERT_NONNULL(worker.impl->env).addRef(js)));
2355 
2356 // HACK: We set handler.env and handler.ctx to undefined because we already passed the real
2357 // env and ctx into the constructor, and we want the handler methods to act like they take
2358 // just one parameter.
2359 handler->env = js.v8Ref(js.v8Undefined());
2360 handler->ctx = kj::none;
2361 
2362 return handler;
2363 };
2364 
2365 KJ_IF_SOME(h, worker.impl->namedHandlers.find(n)) {
2366 jsg::Lock& js = *this;
2367 if (!FeatureFlags::get(js).getReuseCtxAcrossNonclassEvents()) {
2368 api::ExportedHandler constructedHandler = h.clone(js);
2369 constructedHandler.ctx = js.alloc<api::ExecutionContext>(js,
2370 jsg::JsValue(KJ_ASSERT_NONNULL(worker.impl->ctxExports).getHandle(js)), props.toJs(js),
2371 kj::mv(versionInfo));
2372 return kj::heap(kj::mv(constructedHandler));
2373 }
2374 return fakeOwn(h);
2375 } else KJ_IF_SOME(cls, worker.impl->statelessClasses.find(n)) {
2376 return getHandlerFromEntrypointClass(cls);
2377 } else KJ_IF_SOME(cls, worker.impl->workflowClasses.find(n)) {
2378 return getHandlerFromEntrypointClass(cls);
2379 } else if (name == kj::none) {
2380 // If the default export was requested, and we didn't find a handler for it, we'll fall back
2381 // to addEventListener().
2382 //
2383 // Note: The original intention was that we only use addEventListener() for
2384 // service-worker-syntax scripts, but apparently the code has long allowed it for
2385 // modules-based script too, if they lacked an `export default`. Yikes! Sadly, there are
2386 // Workers in production relying on this so we are stuck with it.
2387 return kj::none;
2388 } else {
2389 if (worker.impl->actorClasses.find(n) != kj::none) {
2390 if (isDynamicDispatch) {
2391 JSG_FAIL_REQUIRE(TypeError, "The entrypoint name ", n,
2392 " refers to a Durable Object class, but the incoming request is trying to invoke it as"
2393 " a stateless worker.");
2394 } else {
2395 LOG_ERROR_PERIODICALLY("worker is not an actor but class name was requested", n);
2396 }
2397 } else if (isDynamicDispatch) {
2398 JSG_FAIL_REQUIRE(TypeError, "The entrypoint name ", n,
2399 " was not found in this worker. Ensure the worker exports an entrypoint with that name.");
2400 } else {
2401 LOG_ERROR_PERIODICALLY("worker has no such named entrypoint", n);
2402 }
2403 
2404 KJ_FAIL_ASSERT("worker_do_not_log; Unable to get exported handler");
2405 };
2406}
2407 
2408api::ServiceWorkerGlobalScope& Worker::Lock::getGlobalScope() {
2409 return KJ_ASSERT_NONNULL(jsg::getAlignedPointerFromEmbedderData<api::ServiceWorkerGlobalScope>(
2410 getContext(), jsg::ContextPointerSlot::GLOBAL_WRAPPER));
2411}
2412 
2413TimeoutId::Generator& Worker::Lock::getTimeoutIdGenerator() {
2414 return getGlobalScope().timeoutIdGenerator;
2415}
2416 
2417jsg::AsyncContextFrame::StorageKey& Worker::Lock::getTraceAsyncContextKey() {
2418 // const_cast OK because we are a lock on this isolate.
2419 auto& isolate = const_cast<Isolate&>(worker.getIsolate());
2420 return *(isolate.traceAsyncContextKey);
2421}
2422 
2423jsg::AsyncContextFrame::StorageKey& Worker::Lock::getUserTraceAsyncContextKey() {
2424 // const_cast OK because we are a lock on this isolate.
2425 auto& isolate = const_cast<Isolate&>(worker.getIsolate());
2426 return *(isolate.userTraceAsyncContextKey);
2427}
2428 
2429bool Worker::Lock::isInspectorEnabled() {
2430 return worker.script->isolate->impl->inspector != kj::none;
2431}
2432 
2433void Worker::Lock::logWarning(kj::StringPtr description) {
2434 // const_cast OK because we are a lock on this isolate.
2435 const_cast<Isolate&>(worker.getIsolate()).logWarning(description, *this);
2436}
2437 
2438void Worker::Lock::logWarningOnce(kj::StringPtr description) {
2439 // const_cast OK because we are a lock on this isolate.
2440 const_cast<Isolate&>(worker.getIsolate()).logWarningOnce(description, *this);
2441}
2442 
2443void Worker::Lock::logErrorOnce(kj::StringPtr description) {
2444 // const_cast OK because we are a lock on this isolate.
2445 const_cast<Isolate&>(worker.getIsolate()).logErrorOnce(description);
2446}
2447 
2448void Worker::Lock::logUncaughtException(kj::StringPtr description) {
2449 // We don't add the exception to traces here, since it turns out that this path only gets hit by
2450 // intermediate exception handling.
2451 KJ_IF_SOME(i, worker.script->isolate->impl->inspector) {
2452 JSG_WITHIN_CONTEXT_SCOPE(*this, getContext(),
2453 [&](jsg::Lock& js) { jsg::sendExceptionToInspector(js, *i.get(), description); });
2454 }
2455 
2456 // Run with --verbose to log JS exceptions to stderr. Useful when running tests.
2457 KJ_LOG(INFO, "uncaught exception", description);
2458}
2459 
2460void Worker::Lock::logUncaughtException(
2461 UncaughtExceptionSource source, const jsg::JsValue& exception, const jsg::JsMessage& message) {
2462 // Only add exception to trace when running within an I/O context with a tracer.
2463 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
2464 KJ_IF_SOME(tracer, ioContext.getWorkerTracer()) {
2465 JSG_WITHIN_CONTEXT_SCOPE(*this, getContext(), [&](jsg::Lock& js) {
2466 addExceptionToTrace(impl->inner, ioContext, tracer, source, exception,
2467 worker.getIsolate().getApi().getErrorInterfaceTypeHandler(*this));
2468 });
2469 }
2470 }
2471 
2472 KJ_IF_SOME(i, worker.script->isolate->impl->inspector) {
2473 JSG_WITHIN_CONTEXT_SCOPE(*this, getContext(),
2474 [&](jsg::Lock& js) { sendExceptionToInspector(js, *i.get(), source, exception, message); });
2475 }
2476 
2477 // Run with --verbose to log JS exceptions to stderr. Useful when running tests.
2478 if (kj::_::Debug::shouldLog(::kj::LogSeverity::INFO)) {
2479 JSG_WITHIN_CONTEXT_SCOPE(*this, getContext(), [&](jsg::Lock& js) {
2480 // Try to log `error.stack` if it exists.
2481 KJ_IF_SOME(obj, exception.tryCast<jsg::JsObject>()) {
2482 auto stack = obj.get(js, "stack");
2483 if (!stack.isUndefined()) {
2484 KJ_LOG(INFO, "uncaught exception", source, stack);
2485 return;
2486 }
2487 } else {
2488 // Compiler gives a spurious warning if this `else` isn't here.
2489 }
2490 
2491 KJ_LOG(INFO, "uncaught exception", source, exception);
2492 });
2493 }
2494}
2495 
2496void Worker::Lock::logUncaughtException(UncaughtExceptionSource source, kj::Exception&& exception) {
2497 jsg::Lock& js = *this;
2498 try {
2499 auto jsError = js.exceptionToJsValue(kj::mv(exception),
2500 {
2501 .trusted = true,
2502 });
2503 logUncaughtException(source, jsError.getHandle(js));
2504 } catch (const jsg::JsExceptionThrown&) {
2505 // An exception occurred while trying to convert the exception to a JS value.
2506 // With exceptionToJs, this should only happen if the isolate is terminating
2507 // because of a fatal error when trying to deserialize a tunneled exception
2508 // detail. In this case, we will want to log the original exception instead,
2509 // so let's try exceptionToJs again but this time ignoring the detail, and
2510 // if it throws again, we'll give up and propagate that exception to the
2511 // caller.
2512 auto jsError = js.exceptionToJsValue(exception.clone(), {.ignoreDetail = true});
2513 logUncaughtException(source, jsError.getHandle(js));
2514 }
2515}
2516 
2517void Worker::Lock::reportPromiseRejectEvent(v8::PromiseRejectMessage& message) {
2518 getGlobalScope().emitPromiseRejection(*this, message.GetEvent(),
2519 jsg::V8Ref<v8::Promise>(getIsolate(), message.GetPromise()),
2520 jsg::V8Ref<v8::Value>(getIsolate(), message.GetValue()));
2521}
2522 
2523void Worker::Lock::validateHandlers(ValidationErrorReporter& errorReporter) {
2524 JSG_WITHIN_CONTEXT_SCOPE(*this, getContext(), [&](jsg::Lock& js) {
2525 kj::HashSet<kj::StringPtr> ignoredHandlers;
2526 ignoredHandlers.insert("alarm"_kj);
2527 ignoredHandlers.insert("unhandledrejection"_kj);
2528 ignoredHandlers.insert("rejectionhandled"_kj);
2529 
2530 // Helper function to collect methods from a prototype chain
2531 auto collectMethodsFromPrototypeChain = [&](jsg::JsValue startProto,
2532 kj::HashSet<kj::String>& seenNames) {
2533 // Find the prototype for `Object` by creating one.
2534 auto obj = js.obj();
2535 jsg::JsValue prototypeOfObject = obj.getPrototype(js);
2536 
2537 // Walk the prototype chain.
2538 jsg::JsValue proto = startProto;
2539 for (;;) {
2540 auto protoObj = KJ_UNWRAP_OR(proto.tryCast<jsg::JsObject>(), {
2541 errorReporter.addError(
2542 kj::str("Exported value's prototype chain does not end in Object."));
2543 return;
2544 });
2545 if (protoObj == prototypeOfObject) {
2546 // Reached the prototype for `Object`. Stop here.
2547 break;
2548 }
2549 
2550 // Awkwardly, the prototype's members are not typically enumerable, so we have to
2551 // enumerate them rather directly.
2552 jsg::JsArray properties = protoObj.getPropertyNames(js, jsg::KeyCollectionFilter::OWN_ONLY,
2553 jsg::PropertyFilter::SKIP_SYMBOLS, jsg::IndexFilter::SKIP_INDICES);
2554 for (auto i: kj::zeroTo(properties.size())) {
2555 auto name = properties.get(js, i).toString(js);
2556 if (name == "constructor"_kj) {
2557 // Don't treat special method `constructor` as an exported handler.
2558 continue;
2559 }
2560 
2561 if (!ignoredHandlers.contains(name)) {
2562 // Only report each method name once, even if it overrides a method in a superclass.
2563 seenNames.upsert(kj::mv(name), [&](auto&, auto&&) {});
2564 }
2565 }
2566 
2567 proto = protoObj.getPrototype(js);
2568 }
2569 };
2570 
2571 KJ_IF_SOME(c, worker.impl->context) {
2572 // Service workers syntax.
2573 auto handlerNames = c->getHandlerNames();
2574 kj::Vector<kj::String> handlers;
2575 for (auto& name: handlerNames) {
2576 if (!ignoredHandlers.contains(name)) {
2577 handlers.add(kj::str(name));
2578 }
2579 }
2580 if (handlers.empty()) {
2581 errorReporter.addError(
2582 kj::str("No event handlers were registered. This script does nothing."));
2583 }
2584 errorReporter.addEntrypoint(kj::none, handlers.releaseAsArray());
2585 } else {
2586 auto report = [&](kj::Maybe<kj::StringPtr> name, api::ExportedHandler& exported) {
2587 auto handle = exported.self.getHandle(js);
2588 if (handle->IsArray()) {
2589 // HACK: toDict() will throw a TypeError if given an array, because jsg::DictWrapper is
2590 // designed to treat arrays as not matching when a dict is expected. However,
2591 // StructWrapper has no such restriction, and therefore an exported array will
2592 // successfully produce an ExportedHandler (presumably with no handler functions), and
2593 // hence we will see it here. Rather than try to correct this inconsistency between
2594 // struct and dict handling (which could have unintended consequences), let's just
2595 // work around by ignoring arrays here.
2596 errorReporter.addEntrypoint(name, kj::Array<kj::String>());
2597 } else {
2598 // Use a HashSet to avoid duplicates when methods exist both as own properties
2599 // and in the prototype chain
2600 kj::HashSet<kj::String> methodSet;
2601 
2602 // First, check for own properties (like a plain object literal)
2603 auto dict = js.toDict(handle);
2604 for (auto& field: dict.fields) {
2605 if (!ignoredHandlers.contains(field.name)) {
2606 methodSet.upsert(kj::mv(field.name), [&](auto&, auto&&) {});
2607 }
2608 }
2609 
2610 // Then, check for methods in the prototype chain (like a class instance)
2611 js.withinHandleScope([&]() {
2612 collectMethodsFromPrototypeChain(jsg::JsObject(handle).getPrototype(js), methodSet);
2613 });
2614 
2615 // Convert HashSet to Array for reporting
2616 errorReporter.addEntrypoint(name, KJ_MAP(n, methodSet) { return kj::mv(n); });
2617 }
2618 };
2619 
2620 auto getEntrypointName = [&](kj::StringPtr key) -> kj::Maybe<kj::StringPtr> {
2621 if (key == "default"_kj) {
2622 return kj::none;
2623 } else {
2624 return key;
2625 }
2626 };
2627 
2628 for (auto& entry: worker.impl->namedHandlers) {
2629 report(getEntrypointName(entry.key), entry.value);
2630 }
2631 for (auto& entry: worker.impl->actorClasses) {
2632 KJ_IF_SOME(entrypointName, getEntrypointName(entry.key)) {
2633 errorReporter.addActorClass(entrypointName);
2634 } else {
2635 // Hmm, it appears someone tried to export a Durable Object class as a default
2636 // entrypoint. This doesn't actually work: the runtime will not allow this DO class
2637 // to be used, either for actors or as an entrypoint.
2638 //
2639 // TODO(someday): Make this a hard error. I'm hesitant to do it in my current change
2640 // for fear that it'll break someone somewhere forcing a rollback. For now we log.
2641 LOG_PERIODICALLY(ERROR,
2642 "Exported actor class as default entrypoint. This doesn't work, but historically "
2643 "did not produce a startup-time error.");
2644 }
2645 }
2646 for (auto& entry: worker.impl->statelessClasses) {
2647 // We want to report all of the stateless class's members. To do this, we examine its
2648 // prototype, and its prototype's prototype, and so on, until we get to Object's
2649 // prototype, which we ignore.
2650 auto entrypointName = getEntrypointName(entry.key);
2651 kj::HashSet<kj::String> seenNames;
2652 
2653 js.withinHandleScope([&]() {
2654 // For stateless classes, we need to get the class's prototype property
2655 jsg::JsObject ctor(KJ_ASSERT_NONNULL(entry.value.tryGetHandle(js.v8Isolate)));
2656 jsg::JsValue proto = ctor.get(js, "prototype");
2657 collectMethodsFromPrototypeChain(proto, seenNames);
2658 });
2659 
2660 errorReporter.addEntrypoint(entrypointName, KJ_MAP(n, seenNames) { return kj::mv(n); });
2661 }
2662 
2663 for (auto& entry: worker.impl->workflowClasses) {
2664 KJ_IF_SOME(entrypointName, getEntrypointName(entry.key)) {
2665 kj::HashSet<kj::String> seenNames;
2666 
2667 js.withinHandleScope([&]() {
2668 // For stateless classes, we need to get the class's prototype property
2669 jsg::JsObject ctor(KJ_ASSERT_NONNULL(entry.value.tryGetHandle(js.v8Isolate)));
2670 jsg::JsValue proto = ctor.get(js, "prototype");
2671 collectMethodsFromPrototypeChain(proto, seenNames);
2672 });
2673 
2674 errorReporter.addWorkflowClass(entrypointName, KJ_MAP(n, seenNames) { return kj::mv(n); });
2675 } else {
2676 }
2677 }
2678 }
2679 });
2680}
2681 
2682// =======================================================================================
2683// AsyncLock implementation
2684 
2685const kj::EventLoopLocal<Worker::AsyncWaiter*> Worker::AsyncWaiter::threadCurrentWaiter;
2686 
2687Worker::Isolate::AsyncWaiterList::~AsyncWaiterList() noexcept {
2688 // It should be impossible for this list to be non-empty since each member of the list holds a
2689 // strong reference back to us. But if the list is non-empty, we'd better crash here, to avoid
2690 // dangling pointers.
2691 KJ_ASSERT(head == kj::none, "destroying non-empty waiter list?");
2692 KJ_ASSERT(tail == &head, "tail pointer corrupted?");
2693}
2694 
2695kj::Promise<Worker::AsyncLock> Worker::Isolate::takeAsyncLockWithoutRequest(
2696 SpanParent parentSpan) const {
2697 auto lockTiming = getMetrics().tryCreateLockTiming(kj::mv(parentSpan));
2698 return takeAsyncLockImpl(kj::mv(lockTiming));
2699}
2700 
2701kj::Promise<Worker::AsyncLock> Worker::Isolate::takeAsyncLock(RequestObserver& request) const {
2702 auto lockTiming = getMetrics().tryCreateLockTiming(kj::Maybe<RequestObserver&>(request));
2703 return takeAsyncLockImpl(kj::mv(lockTiming));
2704}
2705 
2706kj::Promise<Worker::AsyncLock> Worker::Isolate::takeAsyncLockImpl(
2707 kj::Maybe<kj::Own<IsolateObserver::LockTiming>> lockTiming) const {
2708 kj::Maybe<uint> currentLoad;
2709 if (lockTiming != kj::none) {
2710 currentLoad = getCurrentLoad();
2711 }
2712 
2713 for (uint threadWaitingDifferentLockCount = 0;; ++threadWaitingDifferentLockCount) {
2714 AsyncWaiter* waiter = *AsyncWaiter::threadCurrentWaiter;
2715 
2716 if (waiter == nullptr) {
2717 // Thread is not currently waiting on a lock.
2718 KJ_IF_SOME(lt, lockTiming) {
2719 lt.get()->reportAsyncInfo(KJ_ASSERT_NONNULL(currentLoad), false /* threadWaitingSameLock */,
2720 threadWaitingDifferentLockCount);
2721 }
2722 auto newWaiter = kj::refcounted<AsyncWaiter>(kj::atomicAddRef(*this));
2723 co_await newWaiter->readyPromise;
2724 co_return AsyncLock(kj::mv(newWaiter), kj::mv(lockTiming));
2725 } else if (waiter->isolate == this) {
2726 // Thread is waiting on a lock already, and it's for the same isolate. We can coalesce the
2727 // locks.
2728 KJ_IF_SOME(lt, lockTiming) {
2729 lt.get()->reportAsyncInfo(KJ_ASSERT_NONNULL(currentLoad), true /* threadWaitingSameLock */,
2730 threadWaitingDifferentLockCount);
2731 }
2732 auto newWaiterRef = kj::addRef(*waiter);
2733 co_await newWaiterRef->readyPromise;
2734 co_return AsyncLock(kj::mv(newWaiterRef), kj::mv(lockTiming));
2735 } else {
2736 // Thread is already waiting for or holding a different isolate lock. Wait for that one to
2737 // be released before we try to lock a different isolate.
2738 // TODO(perf): Use of ForkedPromise leads to thundering herd here. Should be minor in practice,
2739 // but we could consider creating another linked list instead...
2740 KJ_IF_SOME(lt, lockTiming) {
2741 lt.get()->waitingForOtherIsolate(waiter->isolate->getId());
2742 }
2743 co_await waiter->releasePromise;
2744 }
2745 }
2746}
2747 
2748kj::Promise<Worker::AsyncLock> Worker::takeAsyncLockWithoutRequest(SpanParent parentSpan) const {
2749 return script->getIsolate().takeAsyncLockWithoutRequest(kj::mv(parentSpan));
2750}
2751 
2752kj::Promise<Worker::AsyncLock> Worker::takeAsyncLock(RequestObserver& request) const {
2753 return script->getIsolate().takeAsyncLock(request);
2754}
2755 
2756Worker::AsyncWaiter::AsyncWaiter(kj::Own<const Isolate> isolateParam)
2757 : executor(kj::getCurrentThreadExecutor()),
2758 isolate(kj::mv(isolateParam)) {
2759 // Init `releasePromise` / `releaseFulfiller`.
2760 {
2761 auto paf = kj::newPromiseAndFulfiller<void>();
2762 releasePromise = paf.promise.fork();
2763 releaseFulfiller = kj::mv(paf.fulfiller);
2764 }
2765 
2766 // Add ourselves to the wait queue for this isolate.
2767 auto lock = isolate->asyncWaiters.lockExclusive();
2768 if (lock->tail == &lock->head) {
2769 // Looks like the queue is empty, so we immediately get the lock.
2770 readyPromise = kj::Promise<void>(kj::READY_NOW).fork();
2771 // We can leave `readyFulfiller` null as no one will ever invoke it anyway.
2772 } else {
2773 // Arrange to get notified later.
2774 auto paf = kj::newPromiseAndCrossThreadFulfiller<void>();
2775 readyPromise = paf.promise.fork();
2776 readyFulfiller = kj::mv(paf.fulfiller);
2777 }
2778 
2779 next = kj::none;
2780 prev = lock->tail;
2781 *lock->tail = this;
2782 lock->tail = &next;
2783 
2784 *threadCurrentWaiter = this;
2785 
2786 __atomic_add_fetch(&isolate->impl->lockAttemptGauge, 1, __ATOMIC_RELAXED);
2787}
2788 
2789Worker::AsyncWaiter::~AsyncWaiter() noexcept {
2790 // This destructor is `noexcept` because an exception here probably leaves the process in a bad
2791 // state.
2792 
2793 __atomic_sub_fetch(&isolate->impl->lockAttemptGauge, 1, __ATOMIC_RELAXED);
2794 
2795 auto lock = isolate->asyncWaiters.lockExclusive();
2796 
2797 releaseFulfiller->fulfill();
2798 
2799 // Remove ourselves from the list.
2800 *prev = next;
2801 KJ_IF_SOME(n, next) {
2802 n.prev = prev;
2803 } else {
2804 lock->tail = prev;
2805 }
2806 
2807 if (prev == &lock->head) {
2808 // We held the lock before now. Alert the next waiter that they are now at the front of the
2809 // line.
2810 KJ_IF_SOME(n, next) {
2811 n.readyFulfiller->fulfill();
2812 }
2813 }
2814 
2815 auto& w = *threadCurrentWaiter;
2816 KJ_ASSERT(w == this);
2817 w = nullptr;
2818}
2819 
2820kj::Promise<void> Worker::AsyncLock::whenThreadIdle() {
2821 AsyncWaiter*& currentWaiter = *AsyncWaiter::threadCurrentWaiter;
2822 for (;;) {
2823 if (currentWaiter != nullptr) {
2824 co_await currentWaiter->releasePromise;
2825 continue;
2826 }
2827 
2828 co_await kj::yieldUntilQueueEmpty();
2829 
2830 if (currentWaiter == nullptr) {
2831 co_return;
2832 }
2833 // Whoops, a new lock attempt appeared, loop.
2834 }
2835}
2836 
2837// =======================================================================================
2838 
2839// A proxy for OutputStream that internally buffers data as long as it's beyond a given limit.
2840// Also, it counts size of all the data it has seen (whether it has hit the limit or not).
2841//
2842// We use this in the Network tab to report response stats and preview [decompressed] bodies,
2843// but we don't want to keep buffering extremely large ones, so just discard buffered data
2844// upon hitting a limit and don't return any body to the devtools frontend afterwards.
2845class Worker::Isolate::LimitedBodyWrapper: public kj::OutputStream {
2846 public:
2847 LimitedBodyWrapper(size_t limit = 1 * 1024 * 1024): limit(limit) {
2848 if (limit > 0) {
2849 inner.emplace();
2850 }
2851 }
2852 
2853 KJ_DISALLOW_COPY_AND_MOVE(LimitedBodyWrapper);
2854 
2855 void reset() {
2856 this->inner = kj::none;
2857 }
2858 
2859 void write(kj::ArrayPtr<const byte> data) override {
2860 this->size += data.size();
2861 KJ_IF_SOME(inner, this->inner) {
2862 if (this->size <= this->limit) {
2863 inner.write(data);
2864 } else {
2865 reset();
2866 }
2867 }
2868 }
2869 
2870 size_t getWrittenSize() {
2871 return this->size;
2872 }
2873 
2874 kj::Maybe<kj::ArrayPtr<byte>> getArray() {
2875 KJ_IF_SOME(inner, this->inner) {
2876 return inner.getArray();
2877 } else {
2878 return kj::none;
2879 }
2880 }
2881 
2882 private:
2883 size_t size = 0;
2884 size_t limit = 0;
2885 kj::Maybe<kj::VectorOutputStream> inner;
2886};
2887 
2888struct MessageQueue {
2889 kj::Vector<kj::String> messages;
2890 size_t head;
2891 enum class Status { ACTIVE, CLOSED } status;
2892};
2893 
2894class Worker::Isolate::InspectorChannelImpl final: public v8_inspector::V8Inspector::Channel {
2895 public:
2896 InspectorChannelImpl(kj::Own<const Worker::Isolate> isolateParam,
2897 kj::Own<const kj::Executor> isolateThreadExecutor,
2898 kj::WebSocket& webSocket)
2899 : ioHandler(kj::mv(isolateThreadExecutor), webSocket),
2900 state(kj::heap<State>(this, kj::mv(isolateParam))) {
2901 ioHandler.connect(*this);
2902 }
2903 
2904 // In preview sessions, synchronous locks are not an issue. We declare an alternate spelling of
2905 // the type so that all the individual locks below don't turn up in a search for synchronous
2906 // locks.
2907 using InspectorLock = Worker::Lock::TakeSynchronously;
2908 
2909 ~InspectorChannelImpl() noexcept try {
2910 // Stop message pump.
2911 ioHandler.disconnect();
2912 
2913 // Delete session under lock.
2914 auto state = this->state.lockExclusive();
2915 
2916 jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) {
2917 Isolate::Impl::Lock recordedLock(*state->get()->isolate, InspectorLock(kj::none), stackScope);
2918 if (state->get()->isolate->currentInspectorSession != kj::none) {
2919 const_cast<Isolate&>(*state->get()->isolate).disconnectInspector();
2920 }
2921 state->get()->teardownUnderLock();
2922 });
2923 } catch (...) {
2924 // Unfortunately since we're inheriting from Channel which declares a virtual destructor with
2925 // default exception constraints, we have to catch all exceptions here and log them.
2926 // But different kinds of exceptions call for different ways to stringify the exception.
2927 // kj::runCatchingExceptions() normally does this for us, but there's no way to use it while
2928 // wrapping the whole destructor (including destructors of members). So... we do a native
2929 // catch(...) and then we rethrow the exception inside a kj::runCatchingExceptions and then log
2930 // that. Yeah.
2931 //
2932 // TODO(cleanup): Maybe we could add a kj::stringifyCurrentException() or
2933 // kj::logUncaughtException() or something?
2934 KJ_IF_SOME(exception, kj::runCatchingExceptions([&]() { throw; })) {
2935 KJ_LOG(ERROR, "uncaught exception in ~Script() and the C++ standard is broken", exception);
2936 }
2937 }
2938 
2939 void disconnect() {
2940 // Fake like the client requested close. This will cause outgoingLoop() to exit and everything
2941 // will be cleaned up.
2942 ioHandler.disconnect();
2943 }
2944 
2945 void dispatchProtocolMessage(kj::String message,
2946 v8_inspector::V8InspectorSession& session,
2947 Isolate& isolate,
2948 jsg::V8StackScope& stackScope,
2949 Isolate::Impl::Lock& recordedLock) {
2950 capnp::MallocMessageBuilder messageBuilder;
2951 auto cmd = messageBuilder.initRoot<cdp::Command>();
2952 getCdpJsonCodec().decode(message, cmd);
2953 
2954 switch (cmd.which()) {
2955 case cdp::Command::UNKNOWN: {
2956 break;
2957 }
2958 case cdp::Command::NETWORK_ENABLE: {
2959 setNetworkEnabled(true);
2960 cmd.getNetworkEnable().initResult();
2961 break;
2962 }
2963 case cdp::Command::NETWORK_DISABLE: {
2964 setNetworkEnabled(false);
2965 cmd.getNetworkDisable().initResult();
2966 break;
2967 }
2968 case cdp::Command::NETWORK_GET_RESPONSE_BODY: {
2969 auto err = cmd.getNetworkGetResponseBody().initError();
2970 err.setCode(-32600);
2971 err.setMessage("Network.getResponseBody is not supported in this fork");
2972 break;
2973 }
2974 case cdp::Command::PROFILER_STOP: {
2975 KJ_IF_SOME(p, isolate.impl->profiler) {
2976 auto& lock = recordedLock.lock;
2977 stopProfiling(*lock, *p, cmd);
2978 }
2979 break;
2980 }
2981 case cdp::Command::PROFILER_START: {
2982 KJ_IF_SOME(p, isolate.impl->profiler) {
2983 auto& lock = recordedLock.lock;
2984 startProfiling(*lock, *p);
2985 }
2986 break;
2987 }
2988 case cdp::Command::PROFILER_SET_SAMPLING_INTERVAL: {
2989 KJ_IF_SOME(p, isolate.impl->profiler) {
2990 auto interval = cmd.getProfilerSetSamplingInterval().getParams().getInterval();
2991 setSamplingInterval(*p, interval);
2992 }
2993 break;
2994 }
2995 case cdp::Command::PROFILER_ENABLE: {
2996 auto& lock = recordedLock.lock;
2997 isolate.impl->profiler = kj::Own<v8::CpuProfiler>(
2998 v8::CpuProfiler::New(lock->v8Isolate, v8::kDebugNaming, v8::kLazyLogging),
2999 CpuProfilerDisposer::instance);
3000 break;
3001 }
3002 case cdp::Command::TAKE_HEAP_SNAPSHOT: {
3003 auto& lock = recordedLock.lock;
3004 takeHeapSnapshot(*lock, cmd.getTakeHeapSnapshot().getParams());
3005 break;
3006 }
3007 }
3008 
3009 if (!cmd.isUnknown()) {
3010 sendNotification(cmd);
3011 return;
3012 }
3013 
3014 auto& lock = recordedLock.lock;
3015 
3016 // We have at times observed V8 bugs where the inspector queues a background task and
3017 // then synchronously waits for it to complete, which would deadlock if background
3018 // threads are disallowed. Since the inspector is in a process sandbox anyway, it's not
3019 // a big deal to just permit those background threads.
3020 AllowV8BackgroundThreadsScope allowBackgroundThreads;
3021 
3022 ExceptionOrDuration limitErrorOrTime = 0 * kj::NANOSECONDS;
3023 {
3024 auto limitScope = isolate.getLimitEnforcer().enterInspectorJs(*lock, limitErrorOrTime);
3025 session.dispatchProtocolMessage(jsg::toInspectorStringView(message));
3026 }
3027 
3028 // Run microtasks in case the user made an async call.
3029 if (!limitErrorOrTime.is<kj::Exception>()) {
3030 auto limitScope = isolate.getLimitEnforcer().enterInspectorJs(*lock, limitErrorOrTime);
3031 lock->runMicrotasks();
3032 } else {
3033 // Oops, we already exceeded the limit, so force the microtask queue to be thrown away.
3034 lock->terminateNextExecution();
3035 lock->runMicrotasks();
3036 }
3037 
3038 KJ_SWITCH_ONEOF(limitErrorOrTime) {
3039 KJ_CASE_ONEOF(limitError, kj::Exception) {
3040 lock->withinHandleScope([&] {
3041 // HACK: We want to print the error, but we need a context to do that.
3042 // We don't know which contexts exist in this isolate, so I guess we have to
3043 // create one. Ugh.
3044 auto dummyContext = v8::Context::New(lock->v8Isolate);
3045 // We need to set the highest used index in every context we create to be a nullptr
3046 // This is because we might later on call GetAlignedPointerFromEmbedderData which fails with
3047 // a fatal error if the array is smaller than the given index.
3048 jsg::setAlignedPointerInEmbedderData(
3049 dummyContext, jsg::ContextPointerSlot::MAX_POINTER_SLOT, nullptr);
3050 auto& inspector = *KJ_ASSERT_NONNULL(isolate.impl->inspector);
3051 inspector.contextCreated(v8_inspector::V8ContextInfo(dummyContext, 1,
3052 v8_inspector::StringView(reinterpret_cast<const uint8_t*>("Worker"), 6)));
3053 JSG_WITHIN_CONTEXT_SCOPE(*lock, dummyContext, [&](jsg::Lock& js) {
3054 jsg::sendExceptionToInspector(js, inspector,
3055 jsg::extractTunneledExceptionDescription(limitError.getDescription()));
3056 });
3057 inspector.contextDestroyed(dummyContext);
3058 });
3059 }
3060 KJ_CASE_ONEOF(startupTimeElapsed, kj::Duration) {}
3061 }
3062 
3063 if (recordedLock.checkInWithLimitEnforcer(isolate)) {
3064 disconnect();
3065 }
3066 }
3067 
3068 kj::Promise<void> messagePump() {
3069 return ioHandler.messagePump();
3070 }
3071 
3072 void handleDispatchProtocolMessage(
3073 Worker::AsyncLock& asyncLock, kj::MutexGuarded<MessageQueue>& incomingQueue) {
3074 auto lockedState = state.lockExclusive();
3075 v8_inspector::V8InspectorSession& session = *lockedState->get()->session;
3076 Isolate& isolate = const_cast<Isolate&>(*lockedState->get()->isolate);
3077 jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) {
3078 Isolate::Impl::Lock recordedLock(isolate, asyncLock, stackScope);
3079 
3080 auto lockedQueue = incomingQueue.lockExclusive();
3081 if (lockedQueue->status != MessageQueue::Status::ACTIVE) {
3082 return;
3083 }
3084 
3085 auto messages = lockedQueue->messages.slice(lockedQueue->head, lockedQueue->messages.size());
3086 for (auto& message: messages) {
3087 dispatchProtocolMessage(kj::mv(message), session, isolate, stackScope, recordedLock);
3088 }
3089 lockedQueue->messages.clear();
3090 lockedQueue->head = 0;
3091 });
3092 }
3093 
3094 kj::Promise<void> dispatchProtocolMessages(kj::MutexGuarded<MessageQueue>& incomingQueue) {
3095 // This method is called on the I/O thread, which also adds messages to the `incomingQueue`.
3096 // So long as this method does not yield/resume mid-way, there is no concern about how
3097 // long the queue lock is held for whilst dispatching messages.
3098 auto i = kj::atomicAddRef(*this->state.lockExclusive()->get()->isolate);
3099 auto asyncLock = co_await i->takeAsyncLockWithoutRequest(nullptr);
3100 handleDispatchProtocolMessage(asyncLock, incomingQueue);
3101 }
3102 
3103 // ---------------------------------------------------------------------------
3104 // implements Channel
3105 //
3106 // Keep in mind that these methods will be called from various threads!
3107 
3108 void sendResponse(int callId, std::unique_ptr<v8_inspector::StringBuffer> message) override {
3109 // callId is encoded in the message, too. Unsure why this method even exists.
3110 sendNotification(kj::mv(message));
3111 }
3112 
3113 bool isNetworkEnabled() {
3114 return __atomic_load_n(&networkEnabled, __ATOMIC_RELAXED);
3115 }
3116 
3117 void setNetworkEnabled(bool enable) {
3118 __atomic_store_n(&networkEnabled, enable, __ATOMIC_RELAXED);
3119 }
3120 
3121 void sendNotification(kj::String message) {
3122 ioHandler.send(kj::mv(message));
3123 }
3124 
3125 template <typename T>
3126 void sendNotification(T&& message) {
3127 sendNotification(getCdpJsonCodec().encode(message));
3128 }
3129 
3130 void sendNotification(std::unique_ptr<v8_inspector::StringBuffer> message) override {
3131 sendNotification(kj::str(message->string()));
3132 }
3133 
3134 void flushProtocolNotifications() override {
3135 // Are we supposed to do anything here? There's no documentation, so who knows? Maybe we could
3136 // delay signaling the outgoing loop until this call?
3137 }
3138 
3139 // Dispatches one message whilst automatic CDP messages on the I/O worker thread is paused, called
3140 // on the thread executing the isolate whilst execution is suspended due to a breakpoint or
3141 // debugger statement.
3142 bool dispatchOneMessageDuringPause();
3143 
3144 private:
3145 // Class that manages the I/O for devtools connections. I/O is performed on the
3146 // thread associated with the InspectorService (the thread that calls attachInspector).
3147 // Most of the public API is intended for code running on the isolate thread, such as
3148 // the InspectorChannelImpl and the InspectorClient.
3149 class WebSocketIoHandler final {
3150 public:
3151 WebSocketIoHandler(kj::Own<const kj::Executor> isolateThreadExecutor, kj::WebSocket& webSocket)
3152 : isolateThreadExecutor(kj::mv(isolateThreadExecutor)),
3153 webSocket(webSocket) {
3154 // Assume we are being instantiated on the InspectorService thread, the thread that will do
3155 // I/O for CDP messages. Messages are delivered to the InspectorChannelImpl on the Isolate thread.
3156 outgoingQueueNotifier = XThreadNotifier::create();
3157 }
3158 
3159 // Sets the channel that messages are delivered to.
3160 void connect(InspectorChannelImpl& inspectorChannel) {
3161 channel = inspectorChannel;
3162 }
3163 
3164 void disconnect() {
3165 channel = kj::none;
3166 shutdown();
3167 }
3168 
3169 // Blocked the current thread until a message arrives. This is intended
3170 // for use in the InspectorClient when breakpoints are hit. The InspectorClient
3171 // has to remain in runMessageLoopOnPause() but still receive CDP messages
3172 // (e.g. resume).
3173 kj::Maybe<kj::String> waitForMessage() {
3174 return incomingQueue.when([](const MessageQueue& incomingQueue) {
3175 return (incomingQueue.head < incomingQueue.messages.size() ||
3176 incomingQueue.status == MessageQueue::Status::CLOSED);
3177 }, [](MessageQueue& incomingQueue) -> kj::Maybe<kj::String> {
3178 if (incomingQueue.status == MessageQueue::Status::CLOSED) return {};
3179 return pollMessage(incomingQueue);
3180 });
3181 }
3182 
3183 // Message pumping promise that should be evaluated on the InspectorService
3184 // thread.
3185 kj::Promise<void> messagePump() {
3186 // Although inspector I/O must happen on the InspectorService thread (to make sure breakpoints
3187 // don't block inspector I/O), inspector messages must be actually dispatched on the Isolate
3188 // thread. So, we run the dispatch loop on the Isolate thread.
3189 //
3190 // Note that the above comment is only really accurate in vanilla workerd. In the case of the
3191 // internal Cloudflare Workers runtime, `isolateThreadExecutor` may actually refer to the
3192 // current thread's `kj::Executor`. That's fine; calling `executeAsync()` on the current
3193 // thread's executor just posts the task to the event loop, and everything works as expected.
3194 
3195 // Since the dispatch loop and the receive loop communicate over a XThreadNotifier, and
3196 // XThreadNotifiers must be created on the thread which will call their `awaitNotification()`
3197 // function, we awkwardly perform two `executeAsync()`s here, one to create the
3198 // XThreadNotifier, then another to spawn the dispatch loop.
3199 //
3200 // We create a new XThreadNotifier for each `messagePump()` call, rather than try to re-use
3201 // one long-term, because XThreadNotifiers' `awaitNotification()` function is not cancel-safe.
3202 // That is, once its promise is cancelled, the notifier is broken.
3203 auto incomingQueueNotifier =
3204 co_await isolateThreadExecutor->executeAsync([]() { return XThreadNotifier::create(); });
3205 
3206 auto dispatchLoopPromise = isolateThreadExecutor->executeAsync(
3207 [this, notifier = kj::atomicAddRef(*incomingQueueNotifier)]() mutable {
3208 return dispatchLoop(kj::mv(notifier));
3209 });
3210 
3211 co_return co_await receiveLoop(kj::mv(incomingQueueNotifier))
3212 .exclusiveJoin(kj::mv(dispatchLoopPromise))
3213 .exclusiveJoin(transmitLoop());
3214 }
3215 
3216 void send(kj::String message) {
3217 auto lockedOutgoingQueue = outgoingQueue.lockExclusive();
3218 if (lockedOutgoingQueue->status == MessageQueue::Status::CLOSED) return;
3219 lockedOutgoingQueue->messages.add(kj::mv(message));
3220 outgoingQueueNotifier->notify();
3221 }
3222 
3223 private:
3224 static kj::Maybe<kj::String> pollMessage(MessageQueue& messageQueue) {
3225 if (messageQueue.head < messageQueue.messages.size()) {
3226 kj::String message = kj::mv(messageQueue.messages[messageQueue.head++]);
3227 if (messageQueue.head == messageQueue.messages.size()) {
3228 messageQueue.head = 0;
3229 messageQueue.messages.clear();
3230 }
3231 return kj::mv(message);
3232 }
3233 return {};
3234 }
3235 
3236 void shutdown() {
3237 // Drain incoming queue, the isolate thread may be waiting on it
3238 // on will notice it is closed if woken without any messages to
3239 // deliver in WebSocketIoWorker::waitForMessage().
3240 {
3241 auto lockedIncomingQueue = incomingQueue.lockExclusive();
3242 lockedIncomingQueue->head = 0;
3243 lockedIncomingQueue->messages.clear();
3244 lockedIncomingQueue->status = MessageQueue::Status::CLOSED;
3245 }
3246 {
3247 auto lockedOutgoingQueue = outgoingQueue.lockExclusive();
3248 lockedOutgoingQueue->status = MessageQueue::Status::CLOSED;
3249 }
3250 // Wake any waiters since queue status fields have been updated.
3251 outgoingQueueNotifier->notify();
3252 }
3253 
3254 // Must be called on the InspectorService thread.
3255 kj::Promise<void> receiveLoop(kj::Own<XThreadNotifier> incomingQueueNotifier) {
3256 for (;;) {
3257 auto message = co_await webSocket.receive(MAX_MESSAGE_SIZE);
3258 KJ_SWITCH_ONEOF(message) {
3259 KJ_CASE_ONEOF(text, kj::String) {
3260 incomingQueue.lockExclusive()->messages.add(kj::mv(text));
3261 incomingQueueNotifier->notify();
3262 }
3263 KJ_CASE_ONEOF(blob, kj::Array<byte>) {
3264 // Ignore.
3265 }
3266 KJ_CASE_ONEOF(close, kj::WebSocket::Close) {
3267 shutdown();
3268 // Pause here to give transmitLoop() the chance to finish and send a reply close.
3269 // When `transmitLoop()` ends, `messagePump()` as a whole will end, canceling
3270 // `receiveLoop()`.
3271 co_await kj::Promise<void>(kj::NEVER_DONE);
3272 }
3273 }
3274 }
3275 }
3276 
3277 // Must be called on the Isolate thread.
3278 kj::Promise<void> dispatchLoop(kj::Own<XThreadNotifier> incomingQueueNotifier) {
3279 for (;;) {
3280 co_await incomingQueueNotifier->awaitNotification();
3281 KJ_IF_SOME(c, channel) {
3282 co_await c.dispatchProtocolMessages(this->incomingQueue);
3283 }
3284 }
3285 }
3286 
3287 // Must be called on the InspectorService thread.
3288 kj::Promise<void> transmitLoop() {
3289 for (;;) {
3290 co_await outgoingQueueNotifier->awaitNotification();
3291 try {
3292 auto lockedOutgoingQueue = outgoingQueue.lockExclusive();
3293 auto messages = kj::mv(lockedOutgoingQueue->messages);
3294 bool receivedClose = lockedOutgoingQueue->status == MessageQueue::Status::CLOSED;
3295 lockedOutgoingQueue.release();
3296 co_await sendToWebSocket(kj::mv(messages));
3297 if (receivedClose) {
3298 co_await webSocket.close(1000, "client closed connection");
3299 co_return;
3300 }
3301 } catch (kj::Exception& e) {
3302 shutdown();
3303 throw;
3304 }
3305 }
3306 }
3307 
3308 kj::Promise<void> sendToWebSocket(kj::Vector<kj::String> messages) {
3309 for (auto& message: messages) {
3310 co_await webSocket.send(message);
3311 }
3312 }
3313 
3314 // We need access to the Isolate thread's kj::Executor to run the inspector dispatch loop. This
3315 // doesn't actually have to be an Own, because the Isolate thread will destroy the Isolate
3316 // before it exits, but it doesn't hurt.
3317 kj::Own<const kj::Executor> isolateThreadExecutor;
3318 
3319 kj::MutexGuarded<MessageQueue> incomingQueue;
3320 // The notifier for `incomingQueue`, `incomingQueueNotifier`, is created once per
3321 // `messagePump()` call, and never re-used, so it doesn't live here.
3322 
3323 kj::MutexGuarded<MessageQueue> outgoingQueue;
3324 // This XThreadNotifier must be created on the InspectorService thread.
3325 kj::Own<XThreadNotifier> outgoingQueueNotifier;
3326 
3327 kj::WebSocket& webSocket; // only accessed on the InspectorService thread.
3328 std::atomic_bool receivedClose; // accessed on any thread (only transitions false -> true).
3329 kj::Maybe<InspectorChannelImpl&> channel; // only accessed on the isolate thread.
3330 
3331 // Sometimes the inspector protocol sends large messages. KJ defaults to a 1MB size limit
3332 // for WebSocket messages, which makes sense for production use cases, but for debug we should
3333 // be OK to go larger. So, we'll accept 128MB.
3334 static constexpr size_t MAX_MESSAGE_SIZE = 128u << 20;
3335 };
3336 
3337 WebSocketIoHandler ioHandler;
3338 
3339 void takeHeapSnapshot(
3340 jsg::Lock& js, cdp::HeapProfiler::Command::TakeHeapSnapshot::Params::Reader params) {
3341 struct Activity: public v8::ActivityControl {
3342 InspectorChannelImpl& channel;
3343 Activity(InspectorChannelImpl& channel): channel(channel) {}
3344 
3345 ControlOption ReportProgressValue(uint32_t done, uint32_t total) override {
3346 capnp::MallocMessageBuilder message;
3347 auto event = message.initRoot<cdp::Event>();
3348 auto progressParams = event.initReportHeapSnapshotProgress();
3349 progressParams.setDone(done);
3350 progressParams.setTotal(total);
3351 if (done == total) {
3352 progressParams.setFinished(true);
3353 }
3354 auto notification = getCdpJsonCodec().encode(event);
3355 channel.sendNotification(kj::mv(notification));
3356 return ControlOption::kContinue;
3357 }
3358 };
3359 
3360 struct Writer: public v8::OutputStream {
3361 InspectorChannelImpl& channel;
3362 
3363 Writer(InspectorChannelImpl& channel): channel(channel) {}
3364 void EndOfStream() override {}
3365 
3366 int GetChunkSize() override {
3367 return 65536; // big chunks == faster
3368 // The chunk size here will determine the actual number of individual
3369 // messages that are sent. The default is... rather small. Experience
3370 // node and node-heapdump shows that this can be bumped up
3371 // much higher to get better performance. Here we use the value
3372 // that Node.js uses (see Node.js' FileOutputStream impl).
3373 }
3374 
3375 v8::OutputStream::WriteResult WriteAsciiChunk(char* data, int size) override {
3376 capnp::MallocMessageBuilder message;
3377 auto event = message.initRoot<cdp::Event>();
3378 
3379 auto params = event.initAddHeapSnapshotChunk();
3380 params.setChunk(kj::heapString(data, size));
3381 auto notification = getCdpJsonCodec().encode(event);
3382 channel.sendNotification(kj::mv(notification));
3383 
3384 return v8::OutputStream::WriteResult::kContinue;
3385 }
3386 };
3387 
3388 Activity activity(*this);
3389 Writer writer(*this);
3390 
3391 v8::HeapProfiler::HeapSnapshotOptions options{};
3392 if (params.getReportProgress()) {
3393 options.control = &activity;
3394 }
3395 if (params.getExposeInternals()) {
3396 options.snapshot_mode = v8::HeapProfiler::HeapSnapshotMode::kExposeInternals;
3397 }
3398 if (params.getCaptureNumericValue()) {
3399 options.numerics_mode = v8::HeapProfiler::NumericsMode::kExposeNumericValues;
3400 }
3401 
3402 auto profiler = js.v8Isolate->GetHeapProfiler();
3403 auto snapshot = kj::Own<const v8::HeapSnapshot>(
3404 profiler->TakeHeapSnapshot(options), HeapSnapshotDeleter::INSTANCE);
3405 snapshot->Serialize(&writer);
3406 }
3407 
3408 struct State {
3409 kj::Own<const Worker::Isolate> isolate;
3410 std::unique_ptr<v8_inspector::V8InspectorSession> session;
3411 
3412 State(InspectorChannelImpl* self, kj::Own<const Worker::Isolate> isolateParam)
3413 : isolate(kj::mv(isolateParam)),
3414 session(KJ_ASSERT_NONNULL(isolate->impl->inspector)
3415 ->connect(1,
3416 self,
3417 v8_inspector::StringView(),
3418 isolate->impl->inspectorPolicy == InspectorPolicy::ALLOW_UNTRUSTED
3419 ? v8_inspector::V8Inspector::kUntrusted
3420 : v8_inspector::V8Inspector::kFullyTrusted)) {}
3421 ~State() noexcept(false) {
3422 if (session != nullptr) {
3423 KJ_LOG(ERROR,
3424 "Deleting InspectorChannelImpl::State without having called "
3425 "teardownUnderLock()",
3426 kj::getStackTrace());
3427 
3428 // Isolate locks are recursive so it should be safe to lock here.
3429 jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) {
3430 Isolate::Impl::Lock recordedLock(*isolate, InspectorLock(kj::none), stackScope);
3431 session = nullptr;
3432 });
3433 }
3434 }
3435 
3436 // Must be called with the worker isolate locked. Should be called immediately before
3437 // destruction.
3438 void teardownUnderLock() {
3439 session = nullptr;
3440 }
3441 
3442 KJ_DISALLOW_COPY_AND_MOVE(State);
3443 };
3444 // Mutex ordering: You must lock this *before* locking the isolate.
3445 kj::MutexGuarded<kj::Own<State>> state;
3446 
3447 // Not under `state` lock due to lock ordering complications.
3448 volatile bool networkEnabled = false;
3449};
3450 
3451bool Worker::Isolate::InspectorChannelImpl::dispatchOneMessageDuringPause() {
3452 auto maybeMessage = ioHandler.waitForMessage();
3453 // We can be paused by either hitting a debugger statement in a script or from hitting
3454 // a breakpoint or someone hit break.
3455 KJ_IF_SOME(message, maybeMessage) {
3456 auto lockedState = this->state.lockExclusive();
3457 // Received a message whilst script is running, probably in a breakpoint.
3458 v8_inspector::V8InspectorSession& session = *lockedState->get()->session;
3459 // const_cast OK because the IoContext has the lock.
3460 Isolate& isolate = const_cast<Isolate&>(*lockedState->get()->isolate);
3461 Worker::Lock& workerLock = IoContext::current().getCurrentLock();
3462 Isolate::Impl::Lock& recordedLock = workerLock.impl->recordedLock;
3463 jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) {
3464 dispatchProtocolMessage(kj::mv(message), session, isolate, stackScope, recordedLock);
3465 });
3466 return true;
3467 } else {
3468 // No message from waitForMessage() implies the connection is broken.
3469 return false;
3470 }
3471}
3472 
3473bool Worker::InspectorClient::dispatchOneMessageDuringPause(
3474 Worker::Isolate::InspectorChannelImpl& channel) {
3475 return channel.dispatchOneMessageDuringPause();
3476}
3477 
3478kj::Promise<void> Worker::Isolate::attachInspector(kj::Timer& timer,
3479 kj::Duration timerOffset,
3480 kj::HttpService::Response& response,
3481 const kj::HttpHeaderTable& headerTable,
3482 kj::HttpHeaderId controlHeaderId) const {
3483 KJ_REQUIRE(impl->inspector != kj::none);
3484 
3485 kj::HttpHeaders headers(headerTable);
3486 headers.setPtr(controlHeaderId, "{\"ewLog\":{\"status\":\"ok\"}}");
3487 auto webSocket = response.acceptWebSocket(headers);
3488 
3489 // This `attachInspector()` overload is used by the internal Cloudflare Workers runtime, which has
3490 // no concept of a single Isolate thread. Instead, it's OK for all inspector messages to be
3491 // dispatched on the calling thread.
3492 auto executor = kj::getCurrentThreadExecutor().addRef();
3493 
3494 return attachInspector(kj::mv(executor), timer, timerOffset, *webSocket)
3495 .attach(kj::mv(webSocket));
3496}
3497 
3498kj::Promise<void> Worker::Isolate::attachInspector(
3499 kj::Own<const kj::Executor> isolateThreadExecutor,
3500 kj::Timer& timer,
3501 kj::Duration timerOffset,
3502 kj::WebSocket& webSocket) const {
3503 KJ_REQUIRE(impl->inspector != kj::none);
3504 
3505 return jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) {
3506 Isolate::Impl::Lock recordedLock(
3507 *this, InspectorChannelImpl::InspectorLock(kj::none), stackScope);
3508 auto& lock = *recordedLock.lock;
3509 auto& lockedSelf = const_cast<Worker::Isolate&>(*this);
3510 
3511 // If another inspector was already connected, boot it, on the assumption that that connection
3512 // is dead and this is why the user reconnected. While we could actually allow both inspector
3513 // sessions to stay open (V8 supports this!), we'd then need to store a set of all connected
3514 // inspectors in order to be able to disconnect all of them in case of an isolate purge... let's
3515 // just not.
3516 lockedSelf.disconnectInspector();
3517 
3518 lockedSelf.impl->inspectorClient->setInspectorTimerInfo(timer, timerOffset);
3519 
3520 auto channel = kj::heap<Worker::Isolate::InspectorChannelImpl>(
3521 kj::atomicAddRef(*this), kj::mv(isolateThreadExecutor), webSocket);
3522 lockedSelf.currentInspectorSession = *channel;
3523 lockedSelf.impl->inspectorClient->setChannel(*channel);
3524 
3525 // Send any queued notifications.
3526 lock.withinHandleScope([&] {
3527 for (auto& notification: lockedSelf.impl->queuedNotifications) {
3528 channel->sendNotification(kj::mv(notification));
3529 }
3530 lockedSelf.impl->queuedNotifications.clear();
3531 });
3532 
3533 return channel->messagePump().attach(kj::mv(channel));
3534 });
3535}
3536 
3537void Worker::Isolate::disconnectInspector() {
3538 // If an inspector session is connected, proactively drop it, so as to force it to drop its
3539 // reference on the script, so that the script can be deleted.
3540 KJ_IF_SOME(current, currentInspectorSession) {
3541 current.disconnect();
3542 currentInspectorSession = kj::none;
3543 }
3544 impl->inspectorClient->resetChannel();
3545}
3546 
3547void Worker::Isolate::logWarning(kj::StringPtr description, Lock& lock) {
3548 if (impl->inspector != kj::none) {
3549 JSG_WITHIN_CONTEXT_SCOPE(lock, lock.getContext(), [&](jsg::Lock& js) {
3550 logMessage(js, static_cast<uint16_t>(cdp::LogType::WARNING), description);
3551 });
3552 }
3553 
3554 if (loggingOptions.consoleMode == Worker::ConsoleMode::INSPECTOR_ONLY) {
3555 // Run with --verbose to log JS exceptions to stderr. Useful when running tests.
3556 KJ_LOG(INFO, "console warning", description);
3557 } else {
3558 fprintf(stderr, "%s\n", description.cStr());
3559 fflush(stderr);
3560 }
3561 
3562 KJ_IF_SOME(ioContext, IoContext::tryCurrent()) {
3563 KJ_IF_SOME(tracer, ioContext.getWorkerTracer()) {
3564 // json encoding is required over simply wrapping it in quotes to correctly escape the string.
3565 capnp::JsonCodec json;
3566 auto jsonDescription = kj::str("[", json.encode(capnp::Text::Reader(description)), "]");
3567 
3568 auto timestamp = ioContext.now();
3569 tracer.addLog(
3570 ioContext.getInvocationSpanContext(), timestamp, LogLevel::WARN, kj::mv(jsonDescription));
3571 }
3572 }
3573}
3574 
3575void Worker::Isolate::logWarningOnce(kj::StringPtr description, Lock& lock) {
3576 impl->warningOnceDescriptions.findOrCreate(description, [&] {
3577 logWarning(description, lock);
3578 return kj::str(description);
3579 });
3580}
3581 
3582void Worker::Isolate::logErrorOnce(kj::StringPtr description) {
3583 impl->errorOnceDescriptions.findOrCreate(description, [&] {
3584 KJ_LOG(ERROR, description);
3585 return kj::str(description);
3586 });
3587}
3588 
3589void Worker::Isolate::logMessage(jsg::Lock& js, uint16_t type, kj::StringPtr description) {
3590 if (impl->inspector != kj::none) {
3591 // We want to log a warning to the devtools console, as if `console.warn()` were called.
3592 // However, the only public interface to call the real `console.warn()` is via JavaScript,
3593 // where it could have been monkey-patched by the guest. We'd like to avoid having to worry
3594 // about that blowing up in our face. So instead we arrange to send the proper devtools
3595 // protocol messages ourselves.
3596 //
3597 // TODO(cleanup): It would be better if we could directly add the message to the inspector's
3598 // console log (without calling through JavaScript). What we're doing here has some problems.
3599 // In particular, if no client is connected yet, we attempt to queue up the messages to send
3600 // later, much like the real inspector does. This is kind of complicated, and doesn't quite
3601 // work right:
3602 // - The messages won't necessarily be in the right order with normal console logs made at
3603 // the same time (with identical timestamps).
3604 // - In theory we should queue *all* logged warnings and deliver them to every future client,
3605 // not just the next client to connect. But if we do that, we also need to respect the
3606 // protocol command to clear the history when requested. This was further than I cared to
3607 // go.
3608 // To fix these problems, maybe we should just patch V8 with a direct interface into the
3609 // inspector's own log. (Also, how does Chrome handle this?)
3610 
3611 js.withinHandleScope([&] {
3612 capnp::MallocMessageBuilder message;
3613 auto event = message.initRoot<cdp::Event>();
3614 
3615 auto params = event.initRuntimeConsoleApiCalled();
3616 params.setType(static_cast<cdp::LogType>(type));
3617 params.initArgs(1)[0].initString().setValue(description);
3618 params.setExecutionContextId(v8_inspector::V8ContextInfo::executionContextId(js.v8Context()));
3619 params.setTimestamp(impl->inspectorClient->currentTimeMS());
3620 stackTraceToCDP(js, params.initStackTrace());
3621 
3622 auto notification = getCdpJsonCodec().encode(event);
3623 KJ_IF_SOME(i, currentInspectorSession) {
3624 i.sendNotification(kj::mv(notification));
3625 } else {
3626 impl->queuedNotifications.add(kj::mv(notification));
3627 }
3628 });
3629 }
3630}
3631 
3632// =======================================================================================
3633 
3634struct Worker::Actor::Impl {
3635 Actor::Id actorId;
3636 Frankenvalue props;
3637 MakeStorageFunc makeStorage;
3638 
3639 kj::Own<ActorObserver> metrics;
3640 
3641 // When a boolean, indicates whether a `transient` should exist. If true, it will be initialized
3642 // on the first `ensureConstructed()`.
3643 kj::OneOf<bool, jsg::JsRef<jsg::JsValue>> transient;
3644 
3645 kj::Maybe<kj::Own<ActorCacheInterface>> actorCache;
3646 
3647 kj::Maybe<jsg::JsRef<jsg::JsObject>> ctxObject;
3648 
3649 kj::Maybe<rpc::Container::Client> container;
3650 kj::Maybe<FacetManager&> facetManager;
3651 kj::Maybe<ActorVersion> version;
3652 
3653 struct NoClass {};
3654 struct Initializing {};
3655 
3656 // If the actor is backed by a class, this field tracks the instance through its stages. The
3657 // instance is constructed as part of the first request to be delivered.
3658 kj::OneOf<NoClass, // not class-based
3659 Worker::ActorClassInfo*, // constructor not run yet
3660 Initializing, // constructor currently running
3661 api::ExportedHandler, // fully constructed
3662 kj::Exception // constructor threw
3663 >
3664 classInstance;
3665 
3666 class HooksImpl: public InputGate::Hooks, public OutputGate::Hooks, public ActorCache::Hooks {
3667 public:
3668 HooksImpl(kj::Own<Loopback> loopback, TimerChannel& timerChannel, ActorObserver& metrics)
3669 : loopback(kj::mv(loopback)),
3670 timerChannel(timerChannel),
3671 metrics(metrics) {}
3672 
3673 void inputGateLocked() override {
3674 metrics.inputGateLocked();
3675 }
3676 void inputGateReleased() override {
3677 metrics.inputGateReleased();
3678 }
3679 void inputGateWaiterAdded() override {
3680 metrics.inputGateWaiterAdded();
3681 }
3682 void inputGateWaiterRemoved() override {
3683 metrics.inputGateWaiterRemoved();
3684 }
3685 // Implements InputGate::Hooks.
3686 
3687 kj::Promise<void> makeTimeoutPromise() override {
3688 // This really only protects against total hangs. Lowering the timeout drastically is risky,
3689 // since low timeouts can spuriously fire when under heavy CPU load, failing requests that
3690 // would otherwise succeed.
3691 auto timeout = 30 * kj::SECONDS;
3692 co_await timerChannel.afterLimitTimeout(timeout);
3693 
3694 kj::throwFatalException(KJ_EXCEPTION(OVERLOADED,
3695 "broken.outputGateBroken; jsg.Error: Durable Object storage operation exceeded "
3696 "timeout which caused object to be reset."));
3697 }
3698 
3699 // Implements OutputGate::Hooks.
3700 
3701 void outputGateLocked() override {
3702 metrics.outputGateLocked();
3703 }
3704 void outputGateReleased() override {
3705 metrics.outputGateReleased();
3706 }
3707 void outputGateWaiterAdded() override {
3708 metrics.outputGateWaiterAdded();
3709 }
3710 void outputGateWaiterRemoved() override {
3711 metrics.outputGateWaiterRemoved();
3712 }
3713 
3714 // Implements ActorCache::Hooks
3715 
3716 void updateAlarmInMemory(kj::Maybe<kj::Date> newAlarmTime) override;
3717 void storageReadCompleted(kj::Duration latency) override {
3718 metrics.storageReadCompleted(latency);
3719 }
3720 void storageWriteCompleted(kj::Duration latency) override {
3721 metrics.storageWriteCompleted(latency);
3722 }
3723 
3724 private:
3725 kj::Own<Loopback> loopback; // only for updateAlarmInMemory()
3726 TimerChannel& timerChannel; // only for afterLimitTimeout() and updateAlarmInMemory()
3727 ActorObserver& metrics;
3728 
3729 kj::Maybe<kj::Promise<void>> maybeAlarmPreviewTask;
3730 };
3731 
3732 HooksImpl hooks;
3733 
3734 // Handles both input locks and request locks.
3735 InputGate inputGate;
3736 
3737 // Handles output locks.
3738 OutputGate outputGate;
3739 
3740 // `ioContext` is initialized upon delivery of the first request.
3741 kj::Maybe<kj::Own<IoContext>> ioContext;
3742 
3743 // If onBroken() is called while `ioContext` is still null, this is initialized. When
3744 // `ioContext` is constructed, this will be fulfilled with `ioContext.onAbort()`.
3745 kj::Maybe<kj::Own<kj::PromiseFulfiller<kj::Promise<void>>>> abortFulfiller;
3746 
3747 // Task which periodically flushes metrics. Initialized after `ioContext` is initialized.
3748 kj::Maybe<kj::Promise<void>> metricsFlushLoopTask;
3749 
3750 // Allows sending requests back into this actor, recreating it as necessary. Safe to hold longer
3751 // than the Worker::Actor is alive.
3752 kj::Own<Loopback> loopback;
3753 
3754 TimerChannel& timerChannel;
3755 
3756 kj::ForkedPromise<void> shutdownPromise;
3757 kj::Own<kj::PromiseFulfiller<void>> shutdownFulfiller;
3758 
3759 // If this Actor has a HibernationManager, it means the Actor has recently accepted a Hibernatable
3760 // websocket. We eventually move the HibernationManager into the DeferredProxy task
3761 // (since it's long lived), but can still refer to the HibernationManager by passing a reference
3762 // in each CustomEvent.
3763 kj::Maybe<kj::Own<HibernationManager>> hibernationManager;
3764 kj::Maybe<uint16_t> hibernationEventType;
3765 
3766 struct ScheduledAlarm {
3767 ScheduledAlarm(
3768 kj::Date scheduledTime, kj::PromiseFulfillerPair<WorkerInterface::AlarmOutcome> pf)
3769 : scheduledTime(scheduledTime),
3770 resultFulfiller(kj::mv(pf.fulfiller)),
3771 resultPromise(pf.promise.fork()) {}
3772 KJ_DISALLOW_COPY(ScheduledAlarm);
3773 ScheduledAlarm(ScheduledAlarm&&) = default;
3774 ~ScheduledAlarm() noexcept(false) {}
3775 
3776 kj::Date scheduledTime;
3777 WorkerInterface::AlarmFulfiller resultFulfiller;
3778 kj::ForkedPromise<WorkerInterface::AlarmOutcome> resultPromise;
3779 kj::Promise<void> cleanupPromise = resultPromise.addBranch().then(
3780 [](WorkerInterface::AlarmOutcome&&) {}, [](kj::Exception&&) {});
3781 // The first thing we do after we get a result should be to remove the running alarm (if we got
3782 // that far). So we grab the first branch now and ignore any results, before anyone else has a
3783 // chance to do so.
3784 };
3785 struct RunningAlarm {
3786 kj::Date scheduledTime;
3787 kj::ForkedPromise<WorkerInterface::AlarmOutcome> resultPromise;
3788 };
3789 // If valid, we have an alarm invocation that has not yet received an `AlarmFulfiller` and thus
3790 // is either waiting for a running alarm or its scheduled time.
3791 kj::Maybe<ScheduledAlarm> maybeScheduledAlarm;
3792 
3793 // If valid, we have an alarm invocation that has received an `AlarmFulfiller` and is currently
3794 // considered running. This alarm is no longer cancelable.
3795 kj::Maybe<RunningAlarm> maybeRunningAlarm;
3796 
3797 // This is a forked promise so that we can schedule and then cancel multiple alarms while an alarm
3798 // is running.
3799 kj::ForkedPromise<void> runningAlarmTask = kj::Promise<void>(kj::READY_NOW).fork();
3800 
3801 Impl(Worker::Actor& self,
3802 Actor::Id actorId,
3803 bool hasTransient,
3804 MakeActorCacheFunc makeActorCache,
3805 Frankenvalue props,
3806 MakeStorageFunc makeStorage,
3807 kj::Own<Loopback> loopback,
3808 TimerChannel& timerChannel,
3809 kj::Own<ActorObserver> metricsParam,
3810 kj::Maybe<kj::Own<HibernationManager>> manager,
3811 kj::Maybe<uint16_t>& hibernationEventType,
3812 kj::Maybe<rpc::Container::Client> container,
3813 kj::Maybe<FacetManager&> facetManager,
3814 kj::PromiseFulfillerPair<void> paf = kj::newPromiseAndFulfiller<void>())
3815 : actorId(kj::mv(actorId)),
3816 props(kj::mv(props)),
3817 makeStorage(kj::mv(makeStorage)),
3818 metrics(kj::mv(metricsParam)),
3819 transient(hasTransient),
3820 container(kj::mv(container)),
3821 facetManager(facetManager),
3822 hooks(loopback->addRef(), timerChannel, *metrics),
3823 inputGate(hooks),
3824 outputGate(hooks),
3825 loopback(kj::mv(loopback)),
3826 timerChannel(timerChannel),
3827 shutdownPromise(paf.promise.fork()),
3828 shutdownFulfiller(kj::mv(paf.fulfiller)),
3829 hibernationManager(kj::mv(manager)),
3830 hibernationEventType(kj::mv(hibernationEventType)) {
3831 actorCache =
3832 makeActorCache(self.worker->getIsolate().impl->actorCacheLru, outputGate, hooks, *metrics);
3833 }
3834};
3835 
3836kj::Promise<Worker::AsyncLock> Worker::takeAsyncLockWhenActorCacheReady(
3837 kj::Date now, Actor& actor, RequestObserver& request) const {
3838 auto lockTiming =
3839 getIsolate().getMetrics().tryCreateLockTiming(kj::Maybe<RequestObserver&>(request));
3840 
3841 KJ_IF_SOME(c, actor.impl->actorCache) {
3842 KJ_IF_SOME(p, c.get()->evictStale(now)) {
3843 // Got backpressure, wait for it.
3844 // TODO(someday): Count this time period differently in lock timing data?
3845 co_await p;
3846 }
3847 }
3848 
3849 co_return co_await getIsolate().takeAsyncLockImpl(kj::mv(lockTiming));
3850}
3851 
3852Worker::Actor::Actor(const Worker& worker,
3853 kj::Maybe<RequestTracker&> tracker,
3854 Actor::Id actorId,
3855 bool hasTransient,
3856 MakeActorCacheFunc makeActorCache,
3857 kj::Maybe<kj::StringPtr> className,
3858 Frankenvalue props,
3859 MakeStorageFunc makeStorage,
3860 kj::Own<Loopback> loopback,
3861 TimerChannel& timerChannel,
3862 kj::Own<ActorObserver> metrics,
3863 kj::Maybe<kj::Own<HibernationManager>> manager,
3864 kj::Maybe<uint16_t> hibernationEventType,
3865 kj::Maybe<rpc::Container::Client> container,
3866 kj::Maybe<FacetManager&> facetManager,
3867 kj::Maybe<ActorVersion> version)
3868 : worker(kj::atomicAddRef(worker)),
3869 tracker(tracker.map([](RequestTracker& tracker) { return tracker.addRef(); })) {
3870 impl = kj::heap<Impl>(*this, kj::mv(actorId), hasTransient, kj::mv(makeActorCache), kj::mv(props),
3871 kj::mv(makeStorage), kj::mv(loopback), timerChannel, kj::mv(metrics), kj::mv(manager),
3872 hibernationEventType, kj::mv(container), facetManager);
3873 impl->version = kj::mv(version);
3874 
3875 KJ_IF_SOME(c, className) {
3876 KJ_IF_SOME(cls, worker.impl->actorClasses.find(c)) {
3877 // const_cast OK because we're just storing the pointer and will only use this under lock.
3878 impl->classInstance = const_cast<ActorClassInfo*>(&cls);
3879 } else {
3880 auto e = KJ_EXCEPTION(FAILED, "broken.ignored; no such actor class", c);
3881 e.setDetail(jsg::EXCEPTION_IS_USER_ERROR, kj::heapArray<kj::byte>(0));
3882 kj::throwFatalException(kj::mv(e));
3883 }
3884 } else {
3885 impl->classInstance = Impl::NoClass();
3886 }
3887}
3888 
3889void Worker::Actor::ensureConstructed(IoContext& context) {
3890 KJ_IF_SOME(info, impl->classInstance.tryGet<ActorClassInfo*>()) {
3891 // IMPORTANT: We need to set the state to "Initializing" synchronously, before
3892 // ensureConstructedImpl() actually executes and acquires the input lock.
3893 // This prevents multiple concurrent initialization attempts if multiple calls to
3894 // ensureConstructed() arrive back-to-back.
3895 //
3896 // This doesn't create a race condition with getHandler() because InputGate::wait()
3897 // synchronously adds the caller to the wait queue, even though it completes
3898 // asynchronously. Any call to getHandler() that arrives after this point will
3899 // have to wait for the input lock, which is only acquired and released by
3900 // ensureConstructedImpl() when it completes initialization.
3901 //
3902 // So the "actor still initializing" error in getHandler() should be impossible
3903 // unless a code path is bypassing the input lock mechanism.
3904 context.addWaitUntil(ensureConstructedImpl(context, *info));
3905 impl->classInstance = Impl::Initializing();
3906 }
3907}
3908 
3909kj::Promise<void> Worker::Actor::ensureConstructedImpl(IoContext& context, ActorClassInfo& info) {
3910 InputGate::Lock inputLock = co_await impl->inputGate.wait(context.getCurrentTraceSpan());
3911 
3912 try {
3913 bool containerRunning = false;
3914 KJ_IF_SOME(c, impl->container) {
3915 // We need to do an RPC to check if the container is running.
3916 // TODO(perf): It would be nice if we could have started this RPC earlier, e.g. in parallel
3917 // with starting the script, and also if we could save the status across hibernations. But
3918 // that would require some refactoring, and this RPC should (eventally) be local, so it's
3919 // not a huge deal.
3920 auto status = co_await c.statusRequest(capnp::MessageSize{4, 0}).send();
3921 containerRunning = status.getRunning();
3922 }
3923 
3924 co_await context.run([this, &info, containerRunning](Worker::Lock& lock) {
3925 jsg::Lock& js = lock;
3926 
3927 kj::Maybe<jsg::Ref<api::DurableObjectStorage>> storage;
3928 KJ_IF_SOME(c, impl->actorCache) {
3929 storage = impl->makeStorage(lock, worker->getIsolate().getApi(), *c);
3930 }
3931 
3932 auto ctx = js.alloc<api::DurableObjectState>(js, cloneId(),
3933 jsg::JsValue(KJ_ASSERT_NONNULL(lock.getWorker().impl->ctxExports).getHandle(js)),
3934 impl->props.toJs(js), kj::mv(storage), kj::mv(impl->container), containerRunning,
3935 impl->facetManager, impl->version.map([](ActorVersion& v) {
3936 return ActorVersion{.cohort = v.cohort.map([](kj::String& s) { return kj::str(s); })};
3937 }));
3938 
3939 auto handler =
3940 info.cls(lock, ctx.addRef(), KJ_ASSERT_NONNULL(lock.getWorker().impl->env).addRef(js));
3941 
3942 // Since we JUST passed `ctx` into the class constructor, it definitely has a handle
3943 // attached. Let's grab it and stash it to implement getCtx().
3944 auto ctxHandle = jsg::JsObject(KJ_ASSERT_NONNULL(ctx.tryGetHandle(js)));
3945 impl->ctxObject = jsg::JsRef<jsg::JsObject>(js, ctxHandle);
3946 
3947 // HACK: We set handler.env to undefined because we already passed the real env into the
3948 // constructor, and we want the handler methods to act like they take just one parameter.
3949 // We do the same for handler.ctx, as ExecutionContext related tasks are performed
3950 // on the actor's state field instead.
3951 handler.env = js.v8Ref(js.v8Undefined());
3952 handler.ctx = kj::none;
3953 handler.missingSuperclass = info.missingSuperclass;
3954 
3955 impl->classInstance = kj::mv(handler);
3956 }, inputLock.addRef(context.getCurrentTraceSpan()));
3957 // We addRef() the inputLock above rather than kj::mv() it so that the lock remains held
3958 // through the catch block below, if an exception is thrown. This is important since we
3959 // MUST update `impl->classInstance` to something other than `Initializing` before we
3960 // release the lock.
3961 } catch (...) {
3962 // Get the KJ exception
3963 auto e = kj::getCaughtExceptionAsKj();
3964 
3965 auto msg = e.getDescription();
3966 if (!msg.startsWith("broken."_kj) && !msg.startsWith("remote.broken."_kj)) {
3967 // If we already set up a brokenness reason, we shouldn't override it.
3968 auto description = jsg::annotateBroken(msg, "broken.constructorFailed");
3969 e.setDescription(kj::mv(description));
3970 }
3971 
3972 context.abort(e.clone());
3973 impl->classInstance = kj::mv(e);
3974 }
3975}
3976 
3977Worker::Actor::~Actor() noexcept(false) {
3978 // Note: We do not need an isolate lock to destroy the actor impl. Everything in it is specific
3979 // to our thread, or is a handle that can be dropped outside of the lock.
3980}
3981 
3982void Worker::Actor::shutdown(uint16_t reasonCode, kj::Maybe<const kj::Exception&> error) {
3983 // We're officially canceling all background work and we're going to destruct the Actor as soon
3984 // as all IoContexts that reference it go out of scope. We might still log additional
3985 // periodic messages, and that's good because we might care about that information. That said,
3986 // we're officially "broken" from this point because we cannot service background work and our
3987 // capability server should have triggered this (potentially indirectly) via its destructor.
3988 KJ_IF_SOME(r, impl->ioContext) {
3989 impl->metrics->shutdown(reasonCode, r.get()->getLimitEnforcer());
3990 } else {
3991 // The actor was shut down before the IoContext was even constructed, so no metrics are
3992 // written.
3993 }
3994 
3995 shutdownActorCache(error);
3996 
3997 impl->shutdownFulfiller->fulfill();
3998}
3999 
4000void Worker::Actor::shutdownActorCache(kj::Maybe<const kj::Exception&> error) {
4001 KJ_IF_SOME(ac, impl->actorCache) {
4002 ac.get()->shutdown(error);
4003 } else {
4004 // The actor was aborted before the actor cache was constructed, nothing to do.
4005 }
4006}
4007 
4008kj::Promise<void> Worker::Actor::onShutdown() {
4009 return impl->shutdownPromise.addBranch();
4010}
4011 
4012kj::Promise<void> Worker::Actor::onBroken() {
4013 // TODO(soon): Detect and report other cases of brokenness, as described in worker.capnp.
4014 
4015 kj::Promise<void> abortPromise = nullptr;
4016 
4017 KJ_IF_SOME(rc, impl->ioContext) {
4018 abortPromise = rc.get()->onAbort();
4019 } else {
4020 auto paf = kj::newPromiseAndFulfiller<kj::Promise<void>>();
4021 abortPromise = kj::mv(paf.promise);
4022 impl->abortFulfiller = kj::mv(paf.fulfiller);
4023 }
4024 
4025 return abortPromise;
4026}
4027 
4028const Worker::Actor::Id& Worker::Actor::getId() {
4029 return impl->actorId;
4030}
4031 
4032bool Worker::Actor::idsEqual(const Id& a, const Id& b) {
4033 if (a.which() != b.which()) return false;
4034 
4035 KJ_SWITCH_ONEOF(a) {
4036 KJ_CASE_ONEOF(actorId, kj::Own<ActorIdFactory::ActorId>) {
4037 return actorId->equals(*b.get<kj::Own<ActorIdFactory::ActorId>>());
4038 }
4039 KJ_CASE_ONEOF(str, kj::String) {
4040 return str == b.get<kj::String>();
4041 }
4042 }
4043 KJ_UNREACHABLE;
4044}
4045 
4046Worker::Actor::Id Worker::Actor::cloneId(Worker::Actor::Id& id) {
4047 KJ_SWITCH_ONEOF(id) {
4048 KJ_CASE_ONEOF(coloLocalId, kj::String) {
4049 return kj::str(coloLocalId);
4050 }
4051 KJ_CASE_ONEOF(globalId, kj::Own<ActorIdFactory::ActorId>) {
4052 return globalId->clone();
4053 }
4054 }
4055 KJ_UNREACHABLE;
4056}
4057 
4058Worker::Actor::Id Worker::Actor::cloneId() {
4059 return cloneId(impl->actorId);
4060}
4061 
4062kj::Maybe<jsg::JsRef<jsg::JsValue>> Worker::Actor::getTransient(Worker::Lock& lock) {
4063 KJ_REQUIRE(&lock.getWorker() == worker.get());
4064 
4065 if (impl->transient.tryGet<bool>().orDefault(false)) {
4066 // First call and `hasTransient` was true. Initialize it now, since we have the lock.
4067 jsg::Lock& js = lock;
4068 impl->transient.init<jsg::JsRef<jsg::JsValue>>(js, js.obj());
4069 }
4070 
4071 return impl->transient.tryGet<jsg::JsRef<jsg::JsValue>>().map(
4072 [&](jsg::JsRef<jsg::JsValue>& val) { return val.addRef(lock); });
4073}
4074 
4075kj::Maybe<ActorCacheInterface&> Worker::Actor::getPersistent() {
4076 return impl->actorCache;
4077}
4078 
4079kj::Own<Worker::Actor::Loopback> Worker::Actor::getLoopback() {
4080 return impl->loopback->addRef();
4081}
4082 
4083kj::Maybe<jsg::Ref<api::DurableObjectStorage>> Worker::Actor::makeStorageForSwSyntax(
4084 Worker::Lock& lock) {
4085 return impl->actorCache.map([&](kj::Own<ActorCacheInterface>& cache) {
4086 return impl->makeStorage(lock, worker->getIsolate().getApi(), *cache);
4087 });
4088}
4089 
4090void Worker::Actor::assertCanSetAlarm() {
4091 KJ_SWITCH_ONEOF(impl->classInstance) {
4092 KJ_CASE_ONEOF(_, Impl::NoClass) {
4093 // Once upon a time, we allowed actors without classes. Let's make a nicer message if we
4094 // we somehow see a classless actor attempt to run an alarm in the wild.
4095 JSG_FAIL_REQUIRE(
4096 TypeError, "Your Durable Object must be class-based in order to call setAlarm()");
4097 }
4098 KJ_CASE_ONEOF(_, Worker::ActorClassInfo*) {
4099 KJ_FAIL_ASSERT("setAlarm() invoked before Durable Object ctor");
4100 }
4101 KJ_CASE_ONEOF(_, Impl::Initializing) {
4102 // We don't explicitly know if we have an alarm handler or not, so just let it happen. We'll
4103 // handle it when we go to run the alarm.
4104 return;
4105 }
4106 KJ_CASE_ONEOF(handler, api::ExportedHandler) {
4107 JSG_REQUIRE(handler.alarm != kj::none, TypeError,
4108 "Your Durable Object class must have an alarm() handler in order to call setAlarm()");
4109 return;
4110 }
4111 KJ_CASE_ONEOF(exception, kj::Exception) {
4112 // We've failed in the ctor, might as well just throw that exception for now.
4113 kj::throwFatalException(exception.clone());
4114 }
4115 }
4116 KJ_UNREACHABLE;
4117}
4118 
4119void Worker::Actor::Impl::HooksImpl::updateAlarmInMemory(kj::Maybe<kj::Date> newTime) {
4120 if (newTime == kj::none) {
4121 maybeAlarmPreviewTask = kj::none;
4122 return;
4123 }
4124 
4125 auto scheduledTime = KJ_ASSERT_NONNULL(newTime);
4126 
4127 auto retry = kj::coCapture([this, originalTime = scheduledTime]() -> kj::Promise<void> {
4128 kj::Date scheduledTime = originalTime;
4129 
4130 for (auto i: kj::zeroTo(WorkerInterface::ALARM_RETRY_MAX_TRIES)) {
4131 co_await timerChannel.atTime(scheduledTime);
4132 auto result = co_await loopback->getWorker(IoChannelFactory::SubrequestMetadata{})
4133 ->runAlarm(originalTime, i);
4134 
4135 if (result.outcome == EventOutcome::OK || !result.retry) {
4136 break;
4137 }
4138 
4139 auto delay = (WorkerInterface::ALARM_RETRY_START_SECONDS << i++) * kj::SECONDS;
4140 scheduledTime = timerChannel.now() + delay;
4141 }
4142 });
4143 
4144 maybeAlarmPreviewTask = retry();
4145}
4146 
4147kj::Maybe<kj::Promise<WorkerInterface::AlarmOutcome>> Worker::Actor::getAlarm(
4148 kj::Date scheduledTime) {
4149 KJ_IF_SOME(runningAlarm, impl->maybeRunningAlarm) {
4150 if (runningAlarm.scheduledTime == scheduledTime) {
4151 // The running alarm has the same time, we can just wait for it.
4152 return runningAlarm.resultPromise.addBranch();
4153 }
4154 }
4155 
4156 KJ_IF_SOME(scheduledAlarm, impl->maybeScheduledAlarm) {
4157 if (scheduledAlarm.scheduledTime == scheduledTime) {
4158 // The scheduled alarm has the same time, we can just wait for it.
4159 return scheduledAlarm.resultPromise.addBranch();
4160 }
4161 }
4162 
4163 return kj::none;
4164}
4165 
4166kj::Promise<WorkerInterface::ScheduleAlarmResult> Worker::Actor::scheduleAlarm(
4167 kj::Date scheduledTime) {
4168 KJ_IF_SOME(runningAlarm, impl->maybeRunningAlarm) {
4169 if (runningAlarm.scheduledTime == scheduledTime) {
4170 // The running alarm has the same time, we can just wait for it.
4171 auto result = co_await runningAlarm.resultPromise;
4172 co_return result;
4173 }
4174 }
4175 
4176 KJ_IF_SOME(scheduledAlarm, impl->maybeScheduledAlarm) {
4177 // We had a previously scheduled alarm, let's cancel it.
4178 scheduledAlarm.resultFulfiller.cancel();
4179 impl->maybeScheduledAlarm = kj::none;
4180 }
4181 
4182 KJ_IASSERT(impl->maybeScheduledAlarm == kj::none);
4183 auto& scheduledAlarm = impl->maybeScheduledAlarm.emplace(
4184 scheduledTime, kj::newPromiseAndFulfiller<WorkerInterface::AlarmOutcome>());
4185 
4186 // Probably don't need to use kj::coCapture for this but doing so just to be on the
4187 // safe side...
4188 auto whenCanceled =
4189 (kj::coCapture([&scheduledAlarm]() -> kj::Promise<WorkerInterface::ScheduleAlarmResult> {
4190 // We've been cancelled, so return that result. Note that we cannot be resolved any other
4191 // way until we return an AlarmFulfiller below.
4192 co_return co_await scheduledAlarm.resultPromise;
4193 }))();
4194 
4195 // Date.now() < scheduledTime when the alarm comes in, since we subtract elapsed CPU time from
4196 // the time of last I/O in the implementation of Date.now(). This difference could be used to
4197 // implement a Spectre timer, so we have to wait a little longer until
4198 // `Date.now() == scheduledTime`. Note that this also means that we could invoke ahead of its
4199 // `scheduledTime` and we'll delay until appropriate, this may be useful in cases of clock skew.
4200 
4201 co_return co_await handleAlarm(scheduledTime).exclusiveJoin(kj::mv(whenCanceled));
4202}
4203 
4204kj::Promise<WorkerInterface::ScheduleAlarmResult> Worker::Actor::handleAlarm(
4205 kj::Date scheduledTime) {
4206 // Let's wait for any running alarm to cleanup before we even delay.
4207 co_await impl->runningAlarmTask;
4208 
4209 co_await KJ_ASSERT_NONNULL(impl->ioContext)->atTime(scheduledTime);
4210 // It's time to run! Let's tear apart the scheduled alarm and make a running alarm.
4211 
4212 // `maybeScheduledAlarm` should have the same value we emplaced above. If another call to
4213 // `scheduleAlarm()` emplaced a new value, then `whenCanceled` should have resolved which
4214 // cancels this this promise chain.
4215 auto scheduledAlarm = KJ_ASSERT_NONNULL(kj::mv(impl->maybeScheduledAlarm));
4216 impl->maybeScheduledAlarm = kj::none;
4217 
4218 impl->maybeRunningAlarm.emplace(Impl::RunningAlarm{
4219 .scheduledTime = scheduledAlarm.scheduledTime,
4220 .resultPromise = kj::mv(scheduledAlarm.resultPromise),
4221 });
4222 impl->runningAlarmTask = scheduledAlarm.cleanupPromise
4223 .attach(kj::defer([&impl = *impl]() {
4224 // As soon as we get fulfilled or rejected, let's unset this alarm as the running alarm.
4225 //
4226 // NOTE: We could get here during `Actor`'s destructor, which in turn calls `Actor::Impl`'s
4227 // destructor, which destroys `runningAlarmTask`, which is us. But in this case, `actor.impl`
4228 // is already nulled out (the pointer gets nulled before the destructor runs). This is why we
4229 // captured `impl` by reference above, rather than capturing `this`.
4230 impl.maybeRunningAlarm = kj::none;
4231 })).eagerlyEvaluate([](kj::Exception&& e) {
4232 LOG_EXCEPTION("actorAlarmCleanup", e);
4233 }).fork();
4234 co_return kj::mv(scheduledAlarm.resultFulfiller);
4235}
4236 
4237kj::Maybe<api::ExportedHandler&> Worker::Actor::getHandler() {
4238 KJ_SWITCH_ONEOF(impl->classInstance) {
4239 KJ_CASE_ONEOF(_, Impl::NoClass) {
4240 return kj::none;
4241 }
4242 KJ_CASE_ONEOF(_, Worker::ActorClassInfo*) {
4243 KJ_FAIL_ASSERT("ensureConstructed() wasn't called");
4244 }
4245 KJ_CASE_ONEOF(_, Impl::Initializing) {
4246 // This shouldn't be possible because ensureConstructed() would have initiated the
4247 // construction task which would have taken an input lock as well as the isolate lock,
4248 // which should have prevented any other code from executing on the actor until they
4249 // were released.
4250 KJ_FAIL_ASSERT("actor still initializing when getHandler() called");
4251 }
4252 KJ_CASE_ONEOF(handler, api::ExportedHandler) {
4253 return handler;
4254 }
4255 KJ_CASE_ONEOF(exception, kj::Exception) {
4256 kj::throwFatalException(exception.clone());
4257 }
4258 }
4259 KJ_UNREACHABLE;
4260}
4261 
4262ActorObserver& Worker::Actor::getMetrics() {
4263 return *impl->metrics;
4264}
4265 
4266InputGate& Worker::Actor::getInputGate() {
4267 return impl->inputGate;
4268}
4269 
4270OutputGate& Worker::Actor::getOutputGate() {
4271 return impl->outputGate;
4272}
4273 
4274kj::Maybe<IoContext&> Worker::Actor::getIoContext() {
4275 return impl->ioContext.map([](kj::Own<IoContext>& rc) -> IoContext& { return *rc; });
4276}
4277 
4278void Worker::Actor::setIoContext(kj::Own<IoContext> context) {
4279 KJ_REQUIRE(impl->ioContext == kj::none);
4280 KJ_IF_SOME(f, impl->abortFulfiller) {
4281 f.get()->fulfill(context->onAbort());
4282 impl->abortFulfiller = kj::none;
4283 }
4284 auto& limitEnforcer = context->getLimitEnforcer();
4285 impl->ioContext = kj::mv(context);
4286 impl->metricsFlushLoopTask =
4287 impl->metrics->flushLoop(impl->timerChannel, limitEnforcer)
4288 .eagerlyEvaluate([](kj::Exception&& e) { LOG_EXCEPTION("actorMetricsFlushLoop", e); });
4289}
4290 
4291jsg::JsObject Worker::Actor::getCtx(jsg::Lock& js) {
4292 return KJ_REQUIRE_NONNULL(impl->ctxObject).getHandle(js);
4293}
4294 
4295jsg::JsValue Worker::Actor::getEnv(jsg::Lock& js) {
4296 return jsg::JsValue(KJ_REQUIRE_NONNULL(worker->impl->env).getHandle(js));
4297}
4298 
4299kj::Maybe<Worker::Actor::HibernationManager&> Worker::Actor::getHibernationManager() {
4300 return impl->hibernationManager.map(
4301 [](kj::Own<HibernationManager>& hib) -> HibernationManager& { return *hib; });
4302}
4303 
4304void Worker::Actor::setHibernationManager(kj::Own<HibernationManager> hib) {
4305 KJ_REQUIRE(impl->hibernationManager == kj::none);
4306 hib->setTimerChannel(impl->timerChannel);
4307 // Not the cleanest way to provide hibernation manager with a timer channel reference, but
4308 // where HibernationManager is constructed (actor-state), we don't have a timer channel ref.
4309 impl->hibernationManager = kj::mv(hib);
4310}
4311 
4312kj::Maybe<uint16_t> Worker::Actor::getHibernationEventType() {
4313 return impl->hibernationEventType;
4314}
4315 
4316kj::Own<Worker::Actor> Worker::Actor::addRef() {
4317 KJ_IF_SOME(t, tracker) {
4318 // We can attachToThisReference() here, attached object's lifetime being tied to refcounted
4319 // instance is deliberate.
4320 return kj::addRef(*this).attachToThisReference(t.get()->startRequest());
4321 } else {
4322 return kj::addRef(*this);
4323 }
4324}
4325 
4326// =======================================================================================
4327 
4328uint Worker::Isolate::getCurrentLoad() const {
4329 return __atomic_load_n(&impl->lockAttemptGauge, __ATOMIC_RELAXED);
4330}
4331 
4332uint Worker::Isolate::getLockSuccessCount() const {
4333 return __atomic_load_n(&impl->lockSuccessCount, __ATOMIC_RELAXED);
4334}
4335 
4336kj::Own<const Worker::Script> Worker::Isolate::newScript(kj::StringPtr scriptId,
4337 const Script::Source& source,
4338 IsolateObserver::StartType startType,
4339 SpanParent parentSpan,
4340 kj::Own<workerd::VirtualFileSystem> vfs,
4341 bool logNewScript,
4342 kj::Maybe<ValidationErrorReporter&> errorReporter,
4343 kj::Maybe<kj::Own<api::pyodide::ArtifactBundler_State>> artifacts,
4344 kj::Maybe<kj::Arc<workerd::jsg::modules::ModuleRegistry>> maybeNewModuleRegistry) const {
4345 // Script doesn't already exist, so compile it.
4346 return kj::atomicRefcounted<Script>(kj::atomicAddRef(*this), scriptId, source, startType,
4347 logNewScript, errorReporter, kj::mv(artifacts), kj::mv(parentSpan), kj::mv(vfs),
4348 kj::mv(maybeNewModuleRegistry));
4349}
4350 
4351void Worker::Isolate::completedRequest() const {
4352 limitEnforcer->completedRequest(id);
4353}
4354 
4355bool Worker::Isolate::isInspectorEnabled() const {
4356 return impl->inspector != kj::none;
4357}
4358 
4359namespace {
4360 
4361// We only run the inspector within process sandboxes. There, it is safe to query the real clock
4362// for some things, and we do so because we may not have a IoContext available to get
4363// Spectre-safe time.
4364 
4365// Monotonic time in seconds with millisecond precision.
4366double getMonotonicTimeForProcessSandboxOnly() {
4367 KJ_REQUIRE(!isMultiTenantProcess(), "precise timing not safe in multi-tenant processes");
4368 auto timePoint = kj::systemPreciseMonotonicClock().now();
4369 return (timePoint - kj::origin<kj::TimePoint>()) / kj::MILLISECONDS / 1e3;
4370}
4371 
4372// Wall time in seconds with millisecond precision.
4373double getWallTimeForProcessSandboxOnly() {
4374 KJ_REQUIRE(!isMultiTenantProcess(), "precise timing not safe in multi-tenant processes");
4375 auto timePoint = kj::systemPreciseCalendarClock().now();
4376 return (timePoint - kj::UNIX_EPOCH) / kj::MILLISECONDS / 1e3;
4377}
4378} // namespace
4379 
4380class Worker::Isolate::ResponseStreamWrapper final: public kj::AsyncOutputStream {
4381 public:
4382 ResponseStreamWrapper(kj::Own<const Isolate> isolate,
4383 kj::String requestId,
4384 kj::Own<kj::AsyncOutputStream> inner,
4385 api::StreamEncoding encoding,
4386 RequestObserver& requestMetrics)
4387 : constIsolate(kj::mv(isolate)),
4388 requestId(kj::mv(requestId)),
4389 inner(kj::mv(inner)),
4390 requestMetrics(requestMetrics) {
4391 if (encoding == api::StreamEncoding::GZIP) {
4392 compStream.emplace().init<kj::GzipOutputStream>(decodedBuf, kj::GzipOutputStream::DECOMPRESS);
4393 } else if (encoding == api::StreamEncoding::BROTLI) {
4394 compStream.emplace().init<kj::BrotliOutputStream>(
4395 decodedBuf, kj::BrotliOutputStream::DECOMPRESS);
4396 }
4397 }
4398 
4399 ~ResponseStreamWrapper() noexcept(false) {
4400 // It's possible that we already have an isolate lock, in which case we
4401 // don't want to grab another one. Here, we can determine if we have a
4402 // lock by checking if there is a current IoContext, if we do then we
4403 // definitely have a current lock.
4404 // While it is possible for us to have an isolate lock without a current
4405 // IoContext, it is quite unlikely that we'd be cleaning up a
4406 // ResponseStreamWrapper in that situation, so checking for the current
4407 // IoContext should work fine.
4408 if (IoContext::hasCurrent()) {
4409 reportToInspector();
4410 } else {
4411 // In this case we assume we don't have a lock and need to grab one.
4412 // If we continue to get warnings that we're taking the isolate lock
4413 // recursively here, that means we're cleaning these outside of the
4414 // IoContext but still have the isolate lock. In that case, we would
4415 // likely need to add an API to jsg::Lock to get the current lock
4416 // rather than relying on the IoContext.
4417 jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) {
4418 Isolate::Impl::Lock recordedLock(*constIsolate, InspectorLock(requestMetrics), stackScope);
4419 reportToInspector();
4420 });
4421 }
4422 }
4423 
4424 kj::Promise<void> write(kj::ArrayPtr<const byte> buffer) override {
4425 reportBytes(buffer);
4426 return inner->write(buffer);
4427 }
4428 kj::Promise<void> write(kj::ArrayPtr<const kj::ArrayPtr<const byte>> pieces) override {
4429 for (auto& piece: pieces) {
4430 reportBytes(piece);
4431 }
4432 return inner->write(pieces);
4433 }
4434 void reportBytes(kj::ArrayPtr<const byte> buffer) {
4435 if (buffer.size() == 0) {
4436 return;
4437 }
4438 
4439 rawSize += buffer.size();
4440 
4441 auto prevDecodedSize = decodedBuf.getWrittenSize();
4442 KJ_IF_SOME(comp, compStream) {
4443 KJ_SWITCH_ONEOF(comp) {
4444 KJ_CASE_ONEOF(gzip, kj::GzipOutputStream) {
4445 // On invalid gzip discard the previously decoded body and rethrow to stop the stream.
4446 // This way we will report sizes up to this point but won't read any more invalid data.
4447 KJ_ON_SCOPE_FAILURE(decodedBuf.reset());
4448 
4449 gzip.write(buffer);
4450 gzip.flush();
4451 }
4452 KJ_CASE_ONEOF(brotli, kj::BrotliOutputStream) {
4453 KJ_ON_SCOPE_FAILURE(decodedBuf.reset());
4454 
4455 brotli.write(buffer);
4456 brotli.flush();
4457 }
4458 }
4459 } else {
4460 decodedBuf.write(buffer);
4461 }
4462 auto decodedChunkSize = decodedBuf.getWrittenSize() - prevDecodedSize;
4463 
4464 jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) {
4465 Isolate::Impl::Lock recordedLock(*constIsolate, InspectorLock(requestMetrics), stackScope);
4466 auto& isolate = const_cast<Isolate&>(*constIsolate);
4467 
4468 KJ_IF_SOME(i, isolate.currentInspectorSession) {
4469 capnp::MallocMessageBuilder message;
4470 
4471 auto event = message.initRoot<cdp::Event>();
4472 
4473 auto params = event.initNetworkDataReceived();
4474 params.setRequestId(requestId);
4475 params.setEncodedDataLength(buffer.size());
4476 params.setDataLength(decodedChunkSize);
4477 params.setTimestamp(getMonotonicTimeForProcessSandboxOnly());
4478 
4479 i.sendNotification(event);
4480 }
4481 });
4482 }
4483 
4484 // Intentionally not wrapping `tryPumpFrom` to force consumer to use `write` in a loop which,
4485 // in turn, will report each chunk to the inspector to show progress of a slow response.
4486 
4487 kj::Promise<void> whenWriteDisconnected() override {
4488 return inner->whenWriteDisconnected();
4489 }
4490 
4491 private:
4492 using InspectorLock = InspectorChannelImpl::InspectorLock;
4493 
4494 kj::Own<const Isolate> constIsolate;
4495 kj::String requestId;
4496 kj::Own<kj::AsyncOutputStream> inner;
4497 size_t rawSize = 0;
4498 LimitedBodyWrapper decodedBuf;
4499 kj::Maybe<kj::OneOf<kj::GzipOutputStream, kj::BrotliOutputStream>> compStream;
4500 RequestObserver& requestMetrics;
4501 
4502 // Called when the wrapper is destroyed.
4503 // This should only ever be called when we are holding the isolate lock.
4504 void reportToInspector() {
4505 auto& isolate = const_cast<Isolate&>(*constIsolate);
4506 
4507 KJ_IF_SOME(i, isolate.currentInspectorSession) {
4508 capnp::MallocMessageBuilder message;
4509 
4510 auto event = message.initRoot<cdp::Event>();
4511 
4512 auto params = event.initNetworkLoadingFinished();
4513 params.setRequestId(requestId);
4514 params.setEncodedDataLength(rawSize);
4515 params.setTimestamp(getMonotonicTimeForProcessSandboxOnly());
4516 auto response = params.initCfResponse();
4517 KJ_IF_SOME(body, decodedBuf.getArray()) {
4518 response.setBase64Encoded(true);
4519 response.setBody(kj::encodeBase64(body));
4520 }
4521 
4522 i.sendNotification(event);
4523 }
4524 }
4525};
4526 
4527class Worker::Isolate::SubrequestClient final: public WorkerInterface {
4528 public:
4529 explicit SubrequestClient(kj::Own<const Isolate> isolate,
4530 kj::Own<WorkerInterface> inner,
4531 kj::HttpHeaderId contentEncodingHeaderId,
4532 RequestObserver& requestMetrics)
4533 : constIsolate(kj::mv(isolate)),
4534 inner(kj::mv(inner)),
4535 contentEncodingHeaderId(contentEncodingHeaderId),
4536 requestMetrics(kj::addRef(requestMetrics)) {}
4537 KJ_DISALLOW_COPY_AND_MOVE(SubrequestClient);
4538 kj::Promise<void> request(kj::HttpMethod method,
4539 kj::StringPtr url,
4540 const kj::HttpHeaders& headers,
4541 kj::AsyncInputStream& requestBody,
4542 kj::HttpService::Response& response) override;
4543 kj::Promise<void> connect(kj::StringPtr host,
4544 const kj::HttpHeaders& headers,
4545 kj::AsyncIoStream& connection,
4546 kj::HttpService::ConnectResponse& tunnel,
4547 kj::HttpConnectSettings settings) override;
4548 kj::Promise<void> prewarm(kj::StringPtr url) override;
4549 kj::Promise<ScheduledResult> runScheduled(kj::Date scheduledTime, kj::StringPtr cron) override;
4550 kj::Promise<AlarmResult> runAlarm(kj::Date scheduledTime, uint32_t retryCount) override;
4551 kj::Promise<CustomEvent::Result> customEvent(kj::Own<CustomEvent> event) override;
4552 
4553 private:
4554 kj::Own<const Isolate> constIsolate;
4555 kj::Own<WorkerInterface> inner;
4556 kj::HttpHeaderId contentEncodingHeaderId;
4557 kj::Own<RequestObserver> requestMetrics;
4558};
4559 
4560kj::Promise<void> Worker::Isolate::SubrequestClient::request(kj::HttpMethod method,
4561 kj::StringPtr url,
4562 const kj::HttpHeaders& headers,
4563 kj::AsyncInputStream& requestBody,
4564 kj::HttpService::Response& response) {
4565 using InspectorLock = InspectorChannelImpl::InspectorLock;
4566 
4567 auto signalRequest = [this, method, urlCopy = kj::str(url),
4568 headersCopy = headers.clone()]() -> kj::Maybe<kj::String> {
4569 return jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) -> kj::Maybe<kj::String> {
4570 Isolate::Impl::Lock recordedLock(*constIsolate, InspectorLock(*requestMetrics), stackScope);
4571 auto& lock = *recordedLock.lock;
4572 auto& isolate = const_cast<Isolate&>(*constIsolate);
4573 
4574 if (isolate.currentInspectorSession == kj::none) {
4575 return kj::none;
4576 }
4577 
4578 auto& i = KJ_ASSERT_NONNULL(isolate.currentInspectorSession);
4579 if (!i.isNetworkEnabled()) {
4580 return kj::none;
4581 }
4582 
4583 return lock.withinHandleScope([&] {
4584 auto requestId = kj::str(isolate.nextRequestId++);
4585 
4586 capnp::MallocMessageBuilder message;
4587 
4588 auto event = message.initRoot<cdp::Event>();
4589 
4590 auto params = event.initNetworkRequestWillBeSent();
4591 params.setRequestId(requestId);
4592 params.setLoaderId("");
4593 params.setTimestamp(getMonotonicTimeForProcessSandboxOnly());
4594 params.setWallTime(getWallTimeForProcessSandboxOnly());
4595 params.setType(cdp::Page::ResourceType::FETCH);
4596 
4597 auto initiator = params.initInitiator();
4598 initiator.setType(cdp::Network::Initiator::Type::SCRIPT);
4599 stackTraceToCDP(lock, initiator.initStack());
4600 
4601 auto request = params.initRequest();
4602 request.setUrl(urlCopy);
4603 request.setMethod(kj::str(method));
4604 
4605 headersToCDP(headersCopy, request.initHeaders());
4606 
4607 i.sendNotification(event);
4608 return kj::mv(requestId);
4609 });
4610 });
4611 };
4612 
4613 auto signalResponse =
4614 [this](kj::String requestId, uint statusCode, kj::StringPtr statusText,
4615 const kj::HttpHeaders& headers,
4616 kj::Own<kj::AsyncOutputStream> responseBody) -> kj::Own<kj::AsyncOutputStream> {
4617 // Note that we cannot take the isolate lock here, because if this is a worker-to-worker
4618 // subrequest, the destination isolate's lock may already be held, and we can't take multiple
4619 // isolate locks at once as this could lead to deadlock if the lock orders aren't consistent.
4620 //
4621 // Meanwhile, though, `statusText` and `headers` may point to things that will go away
4622 // immediately after we return. So, let's construct our message now, so that we don't have to
4623 // make redundant copies.
4624 //
4625 // Note that signalResponse() is only called at all if signalRequest() determined that network
4626 // inspection is enabled.
4627 
4628 auto message = kj::heap<capnp::MallocMessageBuilder>();
4629 
4630 auto event = message->initRoot<cdp::Event>();
4631 
4632 auto params = event.initNetworkResponseReceived();
4633 params.setRequestId(requestId);
4634 params.setTimestamp(getMonotonicTimeForProcessSandboxOnly());
4635 params.setType(cdp::Page::ResourceType::OTHER);
4636 
4637 auto response = params.initResponse();
4638 response.setStatus(statusCode);
4639 response.setStatusText(statusText);
4640 response.setProtocol("http/1.1");
4641 KJ_IF_SOME(type, headers.get(kj::HttpHeaderId::CONTENT_TYPE)) {
4642 KJ_IF_SOME(parsed, MimeType::tryParse(type, MimeType::IGNORE_PARAMS)) {
4643 response.setMimeType(parsed.toString());
4644 
4645 // Normally Chrome would know what it's loading based on an element or API used for
4646 // the request. We don't have that privilege, but still want network filters to work,
4647 // so we do our best-effort guess of the resource type based on its mime type.
4648 if (MimeType::HTML == parsed || MimeType::XHTML == parsed) {
4649 params.setType(cdp::Page::ResourceType::DOCUMENT);
4650 } else if (MimeType::CSS == parsed) {
4651 params.setType(cdp::Page::ResourceType::STYLESHEET);
4652 } else if (MimeType::isJavascript(parsed)) {
4653 params.setType(cdp::Page::ResourceType::SCRIPT);
4654 } else if (MimeType::isImage(parsed)) {
4655 params.setType(cdp::Page::ResourceType::IMAGE);
4656 } else if (MimeType::isAudio(parsed) || MimeType::isVideo(parsed)) {
4657 params.setType(cdp::Page::ResourceType::MEDIA);
4658 } else if (MimeType::isFont(parsed)) {
4659 params.setType(cdp::Page::ResourceType::FONT);
4660 } else if (MimeType::MANIFEST_JSON == parsed) {
4661 params.setType(cdp::Page::ResourceType::MANIFEST);
4662 } else if (MimeType::VTT == parsed) {
4663 params.setType(cdp::Page::ResourceType::TEXT_TRACK);
4664 } else if (MimeType::EVENT_STREAM == parsed) {
4665 params.setType(cdp::Page::ResourceType::EVENT_SOURCE);
4666 } else if (MimeType::isXml(parsed) || MimeType::isJson(parsed)) {
4667 params.setType(cdp::Page::ResourceType::XHR);
4668 }
4669 
4670 } else {
4671 response.setMimeType(MimeType::PLAINTEXT_STRING);
4672 }
4673 } else {
4674 response.setMimeType(MimeType::PLAINTEXT_STRING);
4675 }
4676 headersToCDP(headers, response.initHeaders());
4677 
4678 auto encoding = api::StreamEncoding::IDENTITY;
4679 KJ_IF_SOME(encodingStr, headers.get(contentEncodingHeaderId)) {
4680 if (encodingStr == "gzip") {
4681 encoding = api::StreamEncoding::GZIP;
4682 } else if (encodingStr == "br") {
4683 encoding = api::StreamEncoding::BROTLI;
4684 }
4685 }
4686 
4687 // Defer to a later turn of the event loop so that it's safe to take a lock.
4688 return kj::newPromisedStream(kj::evalLater(
4689 [this, responseBody = kj::mv(responseBody), message = kj::mv(message), event, encoding,
4690 requestId = kj::mv(requestId)]() mutable -> kj::Own<kj::AsyncOutputStream> {
4691 // Now we know we can lock...
4692 return jsg::runInV8Stack(
4693 [&](jsg::V8StackScope& stackScope) mutable -> kj::Own<kj::AsyncOutputStream> {
4694 Isolate::Impl::Lock recordedLock(*constIsolate, InspectorLock(*requestMetrics), stackScope);
4695 auto& isolate = const_cast<Isolate&>(*constIsolate);
4696 
4697 // We shouldn't even get here if network inspection isn't active since signalRequest() would
4698 // have returned null... but double-check anyway.
4699 if (isolate.currentInspectorSession == kj::none) {
4700 return kj::mv(responseBody);
4701 }
4702 
4703 auto& i = KJ_ASSERT_NONNULL(isolate.currentInspectorSession);
4704 if (!i.isNetworkEnabled()) {
4705 return kj::mv(responseBody);
4706 }
4707 
4708 i.sendNotification(event);
4709 
4710 return kj::heap<ResponseStreamWrapper>(kj::atomicAddRef(*constIsolate), kj::mv(requestId),
4711 kj::mv(responseBody), encoding, *requestMetrics);
4712 });
4713 }));
4714 };
4715 using SignalResponse = decltype(signalResponse);
4716 
4717 class ResponseWrapper final: public kj::HttpService::Response {
4718 public:
4719 ResponseWrapper(
4720 kj::HttpService::Response& inner, kj::String requestId, SignalResponse signalResponse)
4721 : inner(inner),
4722 requestId(kj::mv(requestId)),
4723 signalResponse(kj::mv(signalResponse)) {}
4724 
4725 kj::Own<kj::AsyncOutputStream> send(uint statusCode,
4726 kj::StringPtr statusText,
4727 const kj::HttpHeaders& headers,
4728 kj::Maybe<uint64_t> expectedBodySize = kj::none) override {
4729 auto body = inner.send(statusCode, statusText, headers, expectedBodySize);
4730 return signalResponse(kj::mv(requestId), statusCode, statusText, headers, kj::mv(body));
4731 }
4732 
4733 kj::Own<kj::WebSocket> acceptWebSocket(const kj::HttpHeaders& headers) override {
4734 auto webSocket = inner.acceptWebSocket(headers);
4735 // TODO(someday): Support sending WebSocket frames over CDP. For now we fake an empty
4736 // response.
4737 signalResponse(kj::mv(requestId), 101, "Switching Protocols", headers, newNullOutputStream());
4738 return kj::mv(webSocket);
4739 }
4740 
4741 private:
4742 kj::HttpService::Response& inner;
4743 kj::String requestId;
4744 SignalResponse signalResponse;
4745 };
4746 
4747 // For accurate lock metrics, we want to avoid taking a recursive isolate lock, so we postpone
4748 // the request until a later turn of the event loop.
4749 auto maybeRequestId = co_await kj::evalLater(kj::mv(signalRequest));
4750 
4751 // While we checked above that the headers are valid, let's check again
4752 // after the co_await...
4753 KJ_IF_SOME(rid, maybeRequestId) {
4754 ResponseWrapper wrapper(response, kj::mv(rid), kj::mv(signalResponse));
4755 co_await inner->request(method, url, headers, requestBody, wrapper);
4756 } else {
4757 co_await inner->request(method, url, headers, requestBody, response);
4758 }
4759}
4760 
4761kj::Promise<void> Worker::Isolate::SubrequestClient::connect(kj::StringPtr host,
4762 const kj::HttpHeaders& headers,
4763 kj::AsyncIoStream& connection,
4764 kj::HttpService::ConnectResponse& tunnel,
4765 kj::HttpConnectSettings settings) {
4766 // TODO(someday): EW-7116 Figure out how to represent TCP connections in the devtools network tab.
4767 return inner->connect(host, headers, connection, tunnel, kj::mv(settings));
4768}
4769 
4770// TODO(someday): Log other kinds of subrequests?
4771kj::Promise<void> Worker::Isolate::SubrequestClient::prewarm(kj::StringPtr url) {
4772 return inner->prewarm(url);
4773}
4774kj::Promise<WorkerInterface::ScheduledResult> Worker::Isolate::SubrequestClient::runScheduled(
4775 kj::Date scheduledTime, kj::StringPtr cron) {
4776 return inner->runScheduled(scheduledTime, cron);
4777}
4778kj::Promise<WorkerInterface::AlarmResult> Worker::Isolate::SubrequestClient::runAlarm(
4779 kj::Date scheduledTime, uint32_t retryCount) {
4780 return inner->runAlarm(scheduledTime, retryCount);
4781}
4782kj::Promise<WorkerInterface::CustomEvent::Result> Worker::Isolate::SubrequestClient::customEvent(
4783 kj::Own<CustomEvent> event) {
4784 return inner->customEvent(kj::mv(event));
4785}
4786 
4787kj::Own<WorkerInterface> Worker::Isolate::wrapSubrequestClient(kj::Own<WorkerInterface> client,
4788 kj::HttpHeaderId contentEncodingHeaderId,
4789 RequestObserver& requestMetrics) const {
4790 if (impl->inspector != kj::none) {
4791 client = kj::heap<SubrequestClient>(
4792 kj::atomicAddRef(*this), kj::mv(client), contentEncodingHeaderId, requestMetrics);
4793 }
4794 
4795 return client;
4796}
4797 
4798} // namespace workerd