// Copyright (c) 2017-2022 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #include "actor-cache.h" #include #include #include #include // for api::StreamEncoding #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #include #if _WIN32 #include #include #include #include #else #include #include #endif namespace workerd { namespace { constexpr kj::StringPtr logLevelToString(LogLevel level) { switch (level) { case LogLevel::DEBUG_: return "debug"; case LogLevel::INFO: return "info"; case LogLevel::LOG: return "log"; case LogLevel::WARN: return "warn"; case LogLevel::ERROR: return "error"; default: return "log"; } } void headersToCDP(const kj::HttpHeaders& in, capnp::JsonValue::Builder out) { std::map> inMap; in.forEach([&](kj::StringPtr name, kj::StringPtr value) { inMap.try_emplace(name, 1).first->second.add(value); }); auto outObj = out.initObject(inMap.size()); auto headersPos = 0; for (auto& entry: inMap) { auto field = outObj[headersPos++]; field.setName(entry.first); // CDP uses strange header representation where headers with multiple // values are merged into one newline-delimited string field.initValue().setString(kj::strArray(entry.second, "\n")); } } void stackTraceToCDP(jsg::Lock& js, cdp::Runtime::StackTrace::Builder builder) { // TODO(cleanup): Maybe use V8Inspector::captureStackTrace() which does this for us. However, it // produces protocol objects in its own format which want to handle their whole serialization // to JSON. Also, those protocol objects are defined in generated code which we currently don't // include in our cached V8 build artifacts; we'd need to fix that. But maybe we should really // be using the V8-generated protocol objects rather than our parallel capnp versions! auto stackTrace = v8::StackTrace::CurrentStackTrace(js.v8Isolate, 10); auto frameCount = stackTrace->GetFrameCount(); auto callFrames = builder.initCallFrames(frameCount); for (int i = 0; i < frameCount; i++) { auto src = stackTrace->GetFrame(js.v8Isolate, i); auto dest = callFrames[i]; auto url = src->GetScriptNameOrSourceURL(); if (!url.IsEmpty()) { dest.setUrl(kj::str(url)); } else { dest.setUrl(""_kj); } dest.setScriptId(kj::str(src->GetScriptId())); auto func = src->GetFunctionName(); if (!func.IsEmpty()) { dest.setFunctionName(kj::str(func)); } else { dest.setFunctionName(""_kj); } // V8 locations are 1-based, but CDP locations are 0-based... oh, well dest.setLineNumber(src->GetLineNumber() - 1); dest.setColumnNumber(src->GetColumn() - 1); } } kj::Own makeCdpJsonCodec() { auto codec = kj::heap(); codec->handleByAnnotation(); codec->handleByAnnotation(); return codec; } const capnp::JsonCodec& getCdpJsonCodec() { static const kj::Own codec = makeCdpJsonCodec(); return *codec; } } // namespace // ======================================================================================= namespace { using ExceptionOrDuration = kj::OneOf; // Inform the inspector of an exception thrown. // // Passes `source` as the exception's short message. Reconstructs `message` from `exception` if // `message` is empty. void sendExceptionToInspector(jsg::Lock& js, v8_inspector::V8Inspector& inspector, UncaughtExceptionSource source, const jsg::JsValue& exception, jsg::JsMessage message) { jsg::sendExceptionToInspector(js, inspector, kj::str(source), exception, message); } void addExceptionToTrace(jsg::Lock& js, IoContext& ioContext, BaseTracer& tracer, UncaughtExceptionSource source, const jsg::JsValue& exception, const jsg::TypeHandler& errorTypeHandler) { if (source == UncaughtExceptionSource::INTERNAL || source == UncaughtExceptionSource::INTERNAL_ASYNC) { // Skip redundant intermediate JS->C++ exception reporting. See: IoContext::runImpl(), // PromiseWrapper::tryUnwrap() // // TODO(someday): Arguably it could make sense to store these exceptions off to the side and // report them only if they don't end up being duplicates of a later exception that has a more // specific context. This would cover cases where the C++ code that eventually received the // exception never ended up reporting it. return; } auto timestamp = ioContext.now(); Worker::Api::ErrorInterface error; if (exception.isObject()) { error = KJ_REQUIRE_NONNULL(errorTypeHandler.tryUnwrap(js, exception), "Should always be possible to unwrap error interface from an object."); } kj::String name; KJ_IF_SOME(n, error.name) { name = kj::str(n); } else { name = kj::str("Error"); } kj::String message; KJ_IF_SOME(m, error.message) { message = kj::str(m); } else { // This doesn't appear to be an Error object. Fall back to stringifying the whole value as // the message. if (!js.v8Isolate->IsExecutionTerminating()) { v8::TryCatch tryCatch(js.v8Isolate); try { message = exception.toString(js); } catch (jsg::JsExceptionThrown&) { // Failed to stringify. // // Note that we're intentionally not checking tryCatch.CanContinue() here, because we still // want to continue even if the isolate has been terminated. } } } kj::Maybe stack; KJ_IF_SOME(s, error.stack) { kj::StringPtr slice = s; // Normally `error.stack` repeats the error type and message first. We don't want send two // copies of that to the trace so we'll strip it off. if (slice.startsWith(name)) { slice = slice.slice(name.size()); if (slice.startsWith(": "_kj)) { slice = slice.slice(2); } } if (slice.startsWith(message)) { slice = slice.slice(message.size()); if (slice.startsWith("\n")) { slice = slice.slice(1); } } if (slice.size() > 0) { stack = kj::str(slice); } } tracer.addException(ioContext.getInvocationSpanContext(), timestamp, kj::mv(name), kj::mv(message), kj::mv(stack)); } void reportStartupError(kj::StringPtr id, jsg::Lock& js, const kj::Maybe>& inspector, const IsolateLimitEnforcer& limitEnforcer, ExceptionOrDuration limitErrorOrTime, v8::TryCatch& catcher, kj::Maybe errorReporter, kj::Maybe& permanentException, SpanParent parentSpan, bool isDynamicWorker) { v8::TryCatch catcher2(js.v8Isolate); ExceptionOrDuration limitErrorOrTime2 = 0 * kj::NANOSECONDS; try { KJ_SWITCH_ONEOF(limitErrorOrTime) { KJ_CASE_ONEOF(limitError, kj::Exception) { auto description = jsg::extractTunneledExceptionDescription(limitError.getDescription()); auto& ex = permanentException.emplace(kj::mv(limitError)); KJ_IF_SOME(e, errorReporter) { e.addError(kj::heapString(description)); } else KJ_IF_SOME(i, inspector) { // We want to extend just enough CPU time as is necessary to report the exception // to the inspector here. 10 milliseconds should be more than enough. auto limitScope = limitEnforcer.enterLoggingJs(js, limitErrorOrTime2); jsg::sendExceptionToInspector(js, *i.get(), description); // When the inspector is active, we don't want to throw here because then the inspector // won't be able to connect and the developer will never know what happened. } else { // We should never get here in production if we've validated scripts before deployment. KJ_LOG(WARNING, "script startup exceeded resource limits", id, ex); kj::throwFatalException(ex.clone()); } } KJ_CASE_ONEOF_DEFAULT { if (catcher.HasCaught()) { js.withinHandleScope([&] { auto exception = catcher.Exception(); permanentException = js.exceptionToKj(js.v8Ref(exception)); KJ_IF_SOME(e, errorReporter) { auto limitScope = limitEnforcer.enterLoggingJs(js, limitErrorOrTime2); kj::Vector lines; lines.add(kj::str("Uncaught ", jsg::extractTunneledExceptionDescription( KJ_ASSERT_NONNULL(permanentException).getDescription()))); jsg::JsMessage message(catcher.Message()); message.addJsStackTrace(js, lines); e.addError(kj::strArray(lines, "\n")); } else KJ_IF_SOME(i, inspector) { auto limitScope = limitEnforcer.enterLoggingJs(js, limitErrorOrTime2); sendExceptionToInspector(js, *i.get(), UncaughtExceptionSource::INTERNAL, jsg::JsValue(exception), jsg::JsMessage(catcher.Message())); // When the inspector is active, we don't want to throw here because then the inspector // won't be able to connect and the developer will never know what happened. } else { // We should never get here in production if we've validated scripts before deployment. // (unless this is a dynamic worker) kj::Vector lines; jsg::JsMessage message(catcher.Message()); message.addJsStackTrace(js, lines); auto trace = kj::strArray(lines, "; "); auto description = KJ_ASSERT_NONNULL(permanentException).getDescription(); auto span = parentSpan.newChild("script_startup_exception"_kjc); span.setTag("error"_kjc, true); span.addLog(kj::systemPreciseCalendarClock().now(), "exception"_kjc, kj::ConstString( kj::str("script startup threw exception", id, description, trace))); if (isDynamicWorker) { // Rethrow the tunneled JSG exception so it converts back to a JS Error. kj::throwFatalException(KJ_ASSERT_NONNULL(permanentException).clone()); } else { KJ_LOG(ERROR, "script startup threw exception", id, description, trace); KJ_FAIL_REQUIRE("script startup threw exception"); } } }); } else { kj::throwFatalException(permanentException .emplace(KJ_EXCEPTION(FAILED, "returned empty handle but didn't throw exception?", id)) .clone()); } } } } catch (const jsg::JsExceptionThrown&) { #define LOG_AND_SET_PERM_EXCEPTION(...) \ KJ_LOG(ERROR, __VA_ARGS__); \ if (permanentException == kj::none) { \ permanentException = KJ_EXCEPTION(FAILED, __VA_ARGS__); \ } KJ_SWITCH_ONEOF(limitErrorOrTime2) { KJ_CASE_ONEOF(limitError2, kj::Exception) { // TODO(cleanup): If we see this error show up in production, stop logging it, because I // guess it's not necessarily an error? The other two cases below are more worrying though. KJ_LOG(ERROR, limitError2); if (permanentException == kj::none) { permanentException = kj::mv(limitError2); } } KJ_CASE_ONEOF_DEFAULT { if (catcher2.HasTerminated()) { LOG_AND_SET_PERM_EXCEPTION( "script startup threw exception; during our attempt to stringify the exception, " "the script apparently was terminated for non-resource-limit reasons.", id); } else { LOG_AND_SET_PERM_EXCEPTION( "script startup threw exception; furthermore, an attempt to stringify the exception " "threw another exception, which shouldn't be possible?", id); } } } #undef LOG_AND_SET_PERM_EXCEPTION } } uint64_t getCurrentThreadId() { #if __linux__ return syscall(SYS_gettid); #elif _WIN32 return GetCurrentThreadId(); #else // Assume MacOS or BSD uint64_t tid; pthread_threadid_np(nullptr, &tid); return tid; #endif } } // namespace // Represents a thread's attempt to take an async lock. Each Isolate has a linked list of // `AsyncWaiter`s. A particular thread only ever owns one `AsyncWaiter` at a time. class Worker::AsyncWaiter: public kj::Refcounted { public: AsyncWaiter(kj::Own isolate); ~AsyncWaiter() noexcept; KJ_DISALLOW_COPY_AND_MOVE(AsyncWaiter); private: // Executor for this waiter's thread. const kj::Executor& executor; // The isolate for which this waiter is currently waiting. kj::Own isolate; // Promise/fulfiller to fire when the waiter reaches the front of the list for the corresponding // isolate. kj::ForkedPromise readyPromise = nullptr; kj::Own> readyFulfiller; // Promise/fulfiller to fire when the AsyncLock is finally released. This is used when a thread // tries to take locks on multiple different isolates concurrently, in order to serialize the // locks so only one is taken at a time. This is NOT a cross-thread fulfiller; it can only be // fulfilled by the thread that owns the waiter. kj::ForkedPromise releasePromise = nullptr; kj::Own> releaseFulfiller; // Protected by the lock on `Isolate::asyncWaiters` for the isolate identified by // `currentIsolate`. Must be null if `currentIsolate` is null. (All other members of `Waiter` // can only be accessed by the thread that created the `Waiter`.) kj::Maybe next; kj::Maybe* prev; static const kj::EventLoopLocal threadCurrentWaiter; friend class Worker::Isolate; friend class Worker::AsyncLock; }; class Worker::InspectorClient: public v8_inspector::V8InspectorClient { public: // Wall time in milliseconds with millisecond precision. console.time() and friends rely on this // function to implement timers. double currentTimeMS() override { auto timePoint = kj::UNIX_EPOCH; KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { // We're on a request-serving thread. timePoint = ioContext.now(); } else { auto lockedState = state.lockExclusive(); KJ_IF_SOME(info, lockedState->inspectorTimerInfo) { if (info.threadId == getCurrentThreadId()) { // We're on an inspector-serving thread. timePoint = info.timer.now() + info.timerOffset - kj::origin() + kj::UNIX_EPOCH; } } // We're at script startup time -- just return the Epoch. } return (timePoint - kj::UNIX_EPOCH) / kj::MILLISECONDS; } void setInspectorTimerInfo(kj::Timer& timer, kj::Duration timerOffset) { auto lockedState = state.lockExclusive(); lockedState->inspectorTimerInfo = InspectorTimerInfo{timer, timerOffset, getCurrentThreadId()}; } void setChannel(Worker::Isolate::InspectorChannelImpl& channel) { auto lockedState = state.lockExclusive(); // There is only one active inspector channel at a time in workerd. The teardown of any // previous channel should have invalidated `lockedState->channel`. KJ_REQUIRE(lockedState->channel == kj::none); lockedState->channel = channel; } void resetChannel() { auto lockedState = state.lockExclusive(); lockedState->channel = kj::none; } // This method is called by v8 when a breakpoint or debugger statement is hit. This method // processes debugger messages until `Debugger.resume()` is called, when v8 then calls // `quitMessageLoopOnPause()`. // // This method is ultimately called from the `InspectorChannelImpl` and the isolate lock is // held when this method is called. void runMessageLoopOnPause(int contextGroupId) override { auto lockedState = state.lockExclusive(); KJ_IF_SOME(channel, lockedState->channel) { runMessageLoop = true; do { if (!dispatchOneMessageDuringPause(channel)) { break; } } while (runMessageLoop); } } // This method is called by v8 to resume execution after a breakpoint is hit. void quitMessageLoopOnPause() override { runMessageLoop = false; } private: static bool dispatchOneMessageDuringPause(Worker::Isolate::InspectorChannelImpl& channel); struct InspectorTimerInfo { kj::Timer& timer; kj::Duration timerOffset; uint64_t threadId; }; bool runMessageLoop; // State that may be set on a thread other than the isolate thread. // These are typically set in attachInspector when an inspector connection is // made. struct State { // Inspector channel to use to pump messages. kj::Maybe channel; // The timer and offset for the inspector-serving thread. kj::Maybe inspectorTimerInfo; }; kj::MutexGuarded state; }; static thread_local const Worker::Api* currentApi = nullptr; const Worker::Api& Worker::Api::current() { KJ_REQUIRE(currentApi != nullptr, "not running JavaScript"); return *currentApi; } kj::Maybe Worker::Api::tryCurrent() { if (currentApi != nullptr) { return *currentApi; } return kj::none; } jsg::Optional> Worker::Api::getCtxCacheProperty(jsg::Lock& js) const { return kj::none; } struct Worker::Impl { kj::Maybe> context; // The environment blob to pass to handlers. kj::Maybe env; kj::Maybe ctxExports; // Note: The default export is given the string name "default", because that's what V8 tells us, // and so it's easiest to go with it. I guess that means that you can't actually name an export // "default"? kj::HashMap namedHandlers; kj::HashMap actorClasses; kj::HashMap statelessClasses; kj::HashMap workflowClasses; // If set, then any attempt to use this worker shall throw this exception. kj::Maybe permanentException; }; // Note that Isolate mutable state is protected by locking the JsgWorkerIsolate unless otherwise // noted. struct Worker::Isolate::Impl { IsolateObserver& metrics; kj::Own inspectorClient; kj::Maybe> inspector; InspectorPolicy inspectorPolicy; kj::Maybe> profiler; ActorCache::SharedLru actorCacheLru; // Used by JSG/Rust integration. ::rust::Box<::workerd::rust::jsg::Realm> realm; // UUID for this isolate, initialized first time getUuid() is called. kj::Lazy uuid; // Notification messages to deliver to the next inspector client when it connects. kj::Vector queuedNotifications; // Set of warning log lines that should not be logged to the inspector again. kj::HashSet warningOnceDescriptions; // Set of error log lines that should not be logged again. kj::HashSet errorOnceDescriptions; // Instantaneous count of how many threads are trying to or have successfully obtained an // AsyncLock on this isolate, used to implement getCurrentLoad(). mutable uint lockAttemptGauge = 0; // Atomically incremented upon every successful lock. The ThreadProgressCounter in Impl::Lock // registers a reference to `lockSuccessCounter` as the thread's progress counter during a lock // attempt. This allows watchdogs to see evidence of forward progress in other threads, even if // their own thread has blocked waiting for the lock for a long time. mutable uint64_t lockSuccessCount = 0; // Wrapper around JsgWorkerIsolate::Lock and various RAII objects which help us report metrics, // measure instantaneous load, avoid spurious watchdog kills, and defer context destruction. // // Always use this wrapper in code which may face lock contention (that's mostly everywhere). class Lock { public: explicit Lock( const Worker::Isolate& isolate, Worker::LockType lockType, jsg::V8StackScope& stackScope) : impl(*isolate.impl), metrics([&isolate, &lockType]() -> kj::Maybe> { KJ_SWITCH_ONEOF(lockType.origin) { KJ_CASE_ONEOF(sync, Worker::Lock::TakeSynchronously) { // TODO(perf): We could do some tracking here to discover overly harmful synchronous // locks. return isolate.getMetrics().tryCreateLockTiming(sync.getRequest()); } KJ_CASE_ONEOF(async, AsyncLock*) { KJ_REQUIRE(async->waiter->isolate.get() == &isolate, "async lock was taken against a different isolate than the synchronous lock"); return kj::mv(async->lockTiming); } } KJ_UNREACHABLE; }()), progressCounter(impl.lockSuccessCount), oldCurrentApi(currentApi), limitEnforcer(isolate.getLimitEnforcer()), loggingOptions(isolate.loggingOptions), lock(isolate.api->lock(stackScope)) { WarnAboutIsolateLockScope::maybeWarn(); // Increment the success count to expose forward progress to all threads. __atomic_add_fetch(&impl.lockSuccessCount, 1, __ATOMIC_RELAXED); metrics.locked(); // We record the current lock so our GC prologue/epilogue callbacks can report GC time via // Jaeger tracing. KJ_DASSERT(impl.currentLock == kj::none, "Isolate lock taken recursively"); impl.currentLock = *this; // Now's a good time to destroy any workers queued up for destruction. auto workersToDestroy = impl.workerDestructionQueue.lockExclusive()->pop(); for (auto& workerImpl: workersToDestroy.asArrayPtr()) { KJ_IF_SOME(c, workerImpl->context) { disposeContext(kj::mv(c)); } workerImpl = nullptr; } currentApi = isolate.api.get(); } ~Lock() noexcept(false) { currentApi = oldCurrentApi; #ifdef KJ_DEBUG // We lack a KJ_DASSERT_NONNULL because it would have to look a lot like KJ_IF_SOME, thus // we use a pragma around KJ_DEBUG here. auto& implCurrentLock = KJ_ASSERT_NONNULL(impl.currentLock, "Isolate lock released twice"); KJ_ASSERT(&implCurrentLock == this, "Isolate lock released recursively"); #endif if (shouldReportIsolateMetrics) { // The isolate asked this lock to report the stats when it released. Let's do it. limitEnforcer.reportMetrics(impl.metrics); } impl.currentLock = kj::none; } KJ_DISALLOW_COPY_AND_MOVE(Lock); void setupContext(v8::Local context) { // The V8Inspector implements the `console` object. KJ_IF_SOME(i, impl.inspector) { i.get()->contextCreated( v8_inspector::V8ContextInfo(context, 1, jsg::toInspectorStringView("Worker"))); } Worker::setupContext(*lock, context, loggingOptions); } void disposeContext(jsg::JsContext context) { lock->withinHandleScope([&] { auto v8Context = context.getHandle(*lock); context->clear(); KJ_IF_SOME(i, impl.inspector) { i.get()->contextDestroyed(v8Context); } { auto drop = kj::mv(context); } lock->v8Isolate->ContextDisposedNotification(v8::ContextDependants::kNoDependants); }); } void gcPrologue() { metrics.gcPrologue(); // Filter out tracked WASM instance entries where the instance has been // garbage-collected (weak instanceRef is empty), allowing the linear memory // to be reclaimed. limitEnforcer.getTrackedWasmInstances().filter(*lock); } void gcEpilogue() { metrics.gcEpilogue(); } // Call limitEnforcer.exitJs(), and also schedule to call limitEnforcer.reportMetrics() // later. Returns true if condemned. We take a mutable reference to it to make sure the caller // believes it has exclusive access. bool checkInWithLimitEnforcer(Worker::Isolate& isolate); private: const Impl& impl; IsolateObserver::LockRecord metrics; ThreadProgressCounter progressCounter; bool shouldReportIsolateMetrics = false; const Api* oldCurrentApi; const IsolateLimitEnforcer& limitEnforcer; // only so we can call getIsolateStats() // When structuredLogging is YES AND consoleMode is STDOUT js logs will be emitted to STDOUT // as newline separated json objects LoggingOptions loggingOptions; public: kj::Own lock; }; // Protected by v8::Locker -- if v8::Locker::IsLocked(isolate) is true, then it is safe to access // this variable. mutable kj::Maybe currentLock; static constexpr auto WORKER_DESTRUCTION_QUEUE_INITIAL_SIZE = 8; static constexpr auto WORKER_DESTRUCTION_QUEUE_MAX_CAPACITY = 100; // Similar in spirit to the deferred destruction queue in jsg::IsolateBase. When a Worker is // destroyed, it puts its Impl, which contains objects that need to be destroyed under the isolate // lock, into this queue. Our own Isolate::Impl::Lock implementation then clears this queue the // next time the isolate is locked, whether that be by a connection thread, or the Worker's own // destructor if it owns the last `kj::Own` reference. // // Fairly obviously, this member is protected by its own mutex, not the isolate lock. const kj::MutexGuarded>> workerDestructionQueue{ WORKER_DESTRUCTION_QUEUE_INITIAL_SIZE, WORKER_DESTRUCTION_QUEUE_MAX_CAPACITY}; // TODO(cleanup): The only reason this exists and we can't just rely on the isolate's regular // deferred destruction queue to lazily destroy the various V8 objects in Worker::Impl is // because our GlobalScope object needs to have a function called on it, and any attached // inspector needs to be notified. JSG doesn't know about these things. struct IsolateState { kj::Own inspectorClient; kj::Maybe> inspector; ::rust::Box<::workerd::rust::jsg::Realm> realm; }; static IsolateState initIsolate( const Api& api, IsolateLimitEnforcer& limitEnforcer, InspectorPolicy inspectorPolicy) { auto inspectorClient = kj::heap(); // Default constructor of ::rust::Box is deleted, so we use a Maybe to delay initialization. kj::Maybe<::rust::Box<::workerd::rust::jsg::Realm>> realm; kj::Maybe> inspector; jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) { auto lock = api.lock(stackScope); auto featureFlagsWords = capnp::canonicalize(api.getFeatureFlags()); realm = ::workerd::rust::jsg::realm_create( lock->v8Isolate, featureFlagsWords.asBytes().as()); lock->v8Isolate->SetData( ::workerd::jsg::SetDataIndex::SET_DATA_RUST_REALM, &*KJ_REQUIRE_NONNULL(realm)); limitEnforcer.customizeIsolate(lock->v8Isolate); if (inspectorPolicy != InspectorPolicy::DISALLOW) { // We just created our isolate, so we don't need to use Isolate::Impl::Lock. KJ_ASSERT(!isMultiTenantProcess(), "inspector is not safe in multi-tenant processes"); inspector = v8_inspector::V8Inspector::create(lock->v8Isolate, inspectorClient.get()); } }); return {kj::mv(inspectorClient), kj::mv(inspector), kj::mv(KJ_REQUIRE_NONNULL(realm))}; } Impl(IsolateObserver& metrics, IsolateLimitEnforcer& limitEnforcer, InspectorPolicy inspectorPolicy, IsolateState state) : metrics(metrics), inspectorClient(kj::mv(state.inspectorClient)), inspector(kj::mv(state.inspector)), inspectorPolicy(inspectorPolicy), actorCacheLru(limitEnforcer.getActorCacheLruOptions()), realm(kj::mv(state.realm)) {} Impl(const Api& api, IsolateObserver& metrics, IsolateLimitEnforcer& limitEnforcer, InspectorPolicy inspectorPolicy) : Impl(metrics, limitEnforcer, inspectorPolicy, initIsolate(api, limitEnforcer, inspectorPolicy)) {} }; namespace { class CpuProfilerDisposer final: public kj::Disposer { public: virtual void disposeImpl(void* pointer) const override { reinterpret_cast(pointer)->Dispose(); } static const CpuProfilerDisposer instance; }; const CpuProfilerDisposer CpuProfilerDisposer::instance{}; static constexpr kj::StringPtr PROFILE_NAME = "Default Profile"_kj; static void setSamplingInterval(v8::CpuProfiler& profiler, int interval) { profiler.SetSamplingInterval(interval); } static void startProfiling(jsg::Lock& js, v8::CpuProfiler& profiler) { js.withinHandleScope([&] { v8::CpuProfilingOptions options( v8::kLeafNodeLineNumbers, v8::CpuProfilingOptions::kNoSampleLimit); profiler.StartProfiling(jsg::v8StrIntern(js.v8Isolate, PROFILE_NAME), kj::mv(options)); }); } static void stopProfiling(jsg::Lock& js, v8::CpuProfiler& profiler, cdp::Command::Builder& cmd) { js.withinHandleScope([&] { auto cpuProfile = profiler.StopProfiling(jsg::v8StrIntern(js.v8Isolate, PROFILE_NAME)); if (cpuProfile == nullptr) return; // profiling never started kj::Vector allNodes; kj::Vector unvisited; unvisited.add(cpuProfile->GetTopDownRoot()); while (!unvisited.empty()) { auto next = unvisited.back(); allNodes.add(next); unvisited.removeLast(); for (int i = 0; i < next->GetChildrenCount(); i++) { unvisited.add(next->GetChild(i)); } } auto res = cmd.getProfilerStop().initResult(); auto profile = res.initProfile(); profile.setStartTime(cpuProfile->GetStartTime()); profile.setEndTime(cpuProfile->GetEndTime()); auto nodes = profile.initNodes(allNodes.size()); for (auto i: kj::indices(allNodes)) { auto nodeBuilder = nodes[i]; nodeBuilder.setId(allNodes[i]->GetNodeId()); auto callFrame = nodeBuilder.initCallFrame(); callFrame.setFunctionName(allNodes[i]->GetFunctionNameStr()); callFrame.setScriptId(kj::str(allNodes[i]->GetScriptId())); callFrame.setUrl(allNodes[i]->GetScriptResourceNameStr()); // V8 locations are 1-based, but CDP locations are 0-based... callFrame.setLineNumber(allNodes[i]->GetLineNumber() - 1); callFrame.setColumnNumber(allNodes[i]->GetColumnNumber() - 1); nodeBuilder.setHitCount(allNodes[i]->GetHitCount()); auto children = nodeBuilder.initChildren(allNodes[i]->GetChildrenCount()); for (int j = 0; j < allNodes[i]->GetChildrenCount(); j++) { children.set(j, allNodes[i]->GetChild(j)->GetNodeId()); } auto hitLineCount = allNodes[i]->GetHitLineCount(); auto lineBuffer = kj::heapArray(hitLineCount); allNodes[i]->GetLineTicks(lineBuffer.begin(), lineBuffer.size()); auto positionTicks = nodeBuilder.initPositionTicks(hitLineCount); for (uint j = 0; j < hitLineCount; j++) { auto positionTick = positionTicks[j]; positionTick.setLine(lineBuffer[j].line); positionTick.setTicks(lineBuffer[j].hit_count); } } auto sampleCount = cpuProfile->GetSamplesCount(); auto samples = profile.initSamples(sampleCount); auto timeDeltas = profile.initTimeDeltas(sampleCount); auto lastTimestamp = cpuProfile->GetStartTime(); for (int i = 0; i < sampleCount; i++) { samples.set(i, cpuProfile->GetSample(i)->GetNodeId()); auto sampleTime = cpuProfile->GetSampleTimestamp(i); timeDeltas.set(i, sampleTime - lastTimestamp); lastTimestamp = sampleTime; } }); } } // anonymous namespace struct Worker::Script::Impl { kj::Own vfs; kj::Maybe> maybeNewModuleRegistry; // When using the new module registry, the module registry itself holds the // SchemaLoader, so we don't need to hold it here. When using the original // module registry, however, we need a schema loader to instantiate capnp // modules and bindings. kj::Maybe> maybeSchemaLoader; kj::OneOf unboundScriptOrMainModule; kj::Array globals; kj::Maybe> moduleContext; // If set, then any attempt to use this script shall throw this exception. kj::Maybe permanentException; Impl(kj::Own vfs, kj::Maybe> maybeNewModuleRegistry) : vfs(kj::mv(vfs)), maybeNewModuleRegistry(kj::mv(maybeNewModuleRegistry)) { if (this->maybeNewModuleRegistry == kj::none) { maybeSchemaLoader = kj::heap(); } } struct DynamicImportResult { jsg::Value value; bool isException = false; DynamicImportResult(jsg::Value value, bool isException = false) : value(kj::mv(value)), isException(isException) {} }; using DynamicImportHandler = kj::Function; void configureDynamicImports(jsg::Lock& js, jsg::ModuleRegistry& modules) { // This is only used with the original module registry implementation. KJ_ASSERT(!FeatureFlags::get(js).getNewModuleRegistry(), "legacy dynamic imports must not be used with the new module registry"); static auto constexpr handleDynamicImport = [](kj::Own worker, DynamicImportHandler handler, kj::Maybe> asyncContext) -> kj::Promise { co_await kj::yield(); auto asyncLock = co_await worker->takeAsyncLockWithoutRequest(nullptr); co_return worker->runInLockScope(asyncLock, [&](Worker::Lock& lock) { TmpDirStoreScope tmpDirStoreScope; return JSG_WITHIN_CONTEXT_SCOPE(lock, lock.getContext(), [&](jsg::Lock& js) { jsg::AsyncContextFrame::Scope asyncContextScope(js, asyncContext); // We have to wrap the call to handler in a try catch here because // we have to tunnel any jsg::JsExceptionThrown instances back. v8::TryCatch tryCatch(js.v8Isolate); ExceptionOrDuration limitErrorOrTime = 0 * kj::NANOSECONDS; try { auto limitScope = worker->getIsolate().getLimitEnforcer().enterDynamicImportJs( lock, limitErrorOrTime); return DynamicImportResult(handler()); } catch (jsg::JsExceptionThrown&) { // Handled below... } catch (kj::Exception& ex) { kj::throwFatalException(kj::mv(ex)); } KJ_ASSERT(tryCatch.HasCaught()); if (!tryCatch.CanContinue() || tryCatch.Exception().IsEmpty()) { // There's nothing else we can do here but throw a generic fatal exception. KJ_SWITCH_ONEOF(limitErrorOrTime) { KJ_CASE_ONEOF(limitError, kj::Exception) { kj::throwFatalException(kj::mv(limitError)); } KJ_CASE_ONEOF_DEFAULT { kj::throwFatalException( JSG_KJ_EXCEPTION(FAILED, Error, "Failed to load dynamic module.")); } } } return DynamicImportResult(js.v8Ref(tryCatch.Exception()), true); }); }); }; modules.setDynamicImportCallback([](jsg::Lock& js, DynamicImportHandler handler) mutable { KJ_IF_SOME(context, IoContext::tryCurrent()) { // If we are within the scope of a IoContext, then we are going to pop // out of it to perform the actual module instantiation. return context.awaitIo(js, handleDynamicImport(kj::atomicAddRef(context.getWorker()), kj::mv(handler), jsg::AsyncContextFrame::currentRef(js)), [](jsg::Lock& js, DynamicImportResult result) { if (result.isException) { return js.rejectedPromise(kj::mv(result.value)); } return js.resolvedPromise(kj::mv(result.value)); }); } // If we got here, there is no current IoContext. We're going to perform the // module resolution synchronously and we do not have to worry about blocking any // i/o. We get here, for instance, when dynamic import is used at the top level of // a script (which is weird, but allowed). // // We do not need to use limitEnforcer.enterDynamicImportJs() here because this should // already be covered by the startup resource limiter. return js.resolvedPromise(handler()); }); } kj::Maybe getNewModuleRegistry() const { return maybeNewModuleRegistry.map( [](auto& r) -> const workerd::jsg::modules::ModuleRegistry& { return *r.get(); }); } }; namespace { // Given an array of strings, return a valid serialized JSON string like: // {"flags":["minimal_subrequests",...]} // // Return null if the array is empty. kj::Maybe makeCompatJson(kj::ArrayPtr enableFlags) { if (enableFlags.size() == 0) { return kj::none; } // Calculate the size of the string we're going to generate. constexpr auto PREFIX = "{\"flags\":["_kj; constexpr auto SUFFIX = "]}"_kj; uint size = std::accumulate(enableFlags.begin(), enableFlags.end(), // We need two quotes and one comma for each enable-flag past the first, plus a NUL char. PREFIX.size() + SUFFIX.size() + 3 * enableFlags.size(), [](uint z, kj::StringPtr s) { return z + s.size(); }); kj::Vector json(size); json.addAll(PREFIX); bool first = true; for (auto flag: enableFlags) { if (first) { first = false; } else { json.add(','); } json.add('"'); for (auto& c: flag.asArray()) { // TODO(cleanup): Copied from simpleJsonStringCheck(). Hopefully this will // go away forever soon. KJ_REQUIRE(c != '\"'); KJ_REQUIRE(c != '\\'); KJ_REQUIRE(c >= 0x20); } json.addAll(flag); json.add('"'); } json.addAll(SUFFIX); json.add('\0'); return kj::String(json.releaseAsArray()); } // When a promise is created in a different IoContext, we need to use a // kj::CrossThreadFulfiller in order to wait on it. The Waiter instance will // be held on the Promise itself, and will be fulfilled/rejected when the // promise is resolved or rejected. This will signal all of the waiters // from other IoContexts. jsg::Promise addCrossThreadPromiseWaiter(jsg::Lock& js, v8::Local& promise) { auto waiter = kj::newPromiseAndCrossThreadFulfiller(); struct Waiter: public kj::Refcounted { kj::Maybe>> fulfiller; void done() { KJ_IF_SOME(f, fulfiller) { // Done this way so that the fulfiller is released as soon as possible // when done as the JS promise may not clean up reactions right away. f->fulfill(); fulfiller = kj::none; } } Waiter(kj::Own> fulfiller) : fulfiller(kj::mv(fulfiller)) {} }; auto fulfiller = kj::refcounted(kj::mv(waiter.fulfiller)); auto onSuccess = [waiter = kj::addRef(*fulfiller)]( jsg::Lock& js, jsg::Value value) mutable { waiter->done(); }; auto onFailure = [waiter = kj::mv(fulfiller)]( jsg::Lock& js, jsg::Value exception) mutable { waiter->done(); }; js.toPromise(promise).then(js, kj::mv(onSuccess), kj::mv(onFailure)); return IoContext::current().awaitIo(js, kj::mv(waiter.promise)); } struct HeapSnapshotDeleter: public kj::Disposer { static const HeapSnapshotDeleter INSTANCE; void disposeImpl(void* ptr) const override { auto snapshot = const_cast(static_cast(ptr)); snapshot->Delete(); } }; const HeapSnapshotDeleter HeapSnapshotDeleter::INSTANCE; } // namespace Worker::Isolate::Isolate(kj::Own apiParam, kj::Own metricsParam, kj::StringPtr id, kj::Own limitEnforcerParam, InspectorPolicy inspectorPolicy, LoggingOptions loggingOptions) : metrics(kj::mv(metricsParam)), id(kj::str(id)), limitEnforcer(kj::mv(limitEnforcerParam)), cpuLimitNearlyExceededCallback( kj::MutexGuarded>>(kj::none)), api(kj::mv(apiParam)), loggingOptions(loggingOptions), featureFlagsForFl(makeCompatJson(decompileCompatibilityFlagsForFl(api->getFeatureFlags()))), impl(kj::heap(*api, *metrics, *limitEnforcer, inspectorPolicy)), weakIsolateRef(WeakIsolateRef::wrap(this)), traceAsyncContextKey(kj::refcounted()), userTraceAsyncContextKey(kj::refcounted()) { api->setIsolateObserver(*metrics); metrics->created(); // We just created our isolate, so we don't need to use Isolate::Impl::Lock (nor an async lock). jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) { auto lock = api->lock(stackScope); auto features = api->getFeatureFlags(); KJ_DASSERT(lock->v8Isolate->GetNumberOfDataSlots() >= jsg::SET_DATA_SLOTS_IN_USE); KJ_DASSERT(lock->v8Isolate->GetData(jsg::SET_DATA_ISOLATE) == nullptr); lock->v8Isolate->SetData(jsg::SET_DATA_ISOLATE, this); lock->setCaptureThrowsAsRejections(features.getCaptureThrowsAsRejections()); // TODO(cleanup): Now that this list has grown significantly, we should probably // refactor to pass all of the options in a single call instead of one by one. if (features.getSetToStringTag()) { lock->setToStringTag(); } if (features.getShouldSetImmutablePrototype() || features.getPythonWorkers()) { lock->setImmutablePrototype(); } if (features.getSpecCompliantPropertyAttributes()) { lock->setSpecCompliantPropertyAttributes(); } if (features.getNodeJsCompatV2()) { lock->setNodeJsCompatEnabled(); } if (features.getEnableNodeJsProcessV2()) { lock->setNodeJsProcessV2Enabled(); } if (features.getRequireReturnsDefaultExport()) { lock->setRequireReturnsDefaultExportEnabled(); } if (features.getThrowOnUnrecognizedImportAssertion()) { lock->setThrowOnUnrecognizedImportAssertion(); } if (features.getNoTopLevelAwaitInRequire()) { lock->disableTopLevelAwait(); } if (features.getEnhancedErrorSerialization()) { lock->setUsingEnhancedErrorSerialization(); } if (features.getFastJsgStruct()) { lock->setUsingFastJsgStruct(); } if (impl->inspector != kj::none || ::kj::_::Debug::shouldLog(::kj::LogSeverity::INFO)) { lock->setLoggerCallback([this](jsg::Lock& js, kj::StringPtr message) { if (impl->inspector != kj::none) { logMessage(js, static_cast(cdp::LogType::WARNING), message); } KJ_LOG(INFO, "console warning", message); }); lock->setErrorReporterCallback([this](jsg::Lock& js, kj::String desc, const jsg::JsValue& error, const jsg::JsMessage& message) { // Only add exception to trace when running within an I/O context with a tracer. KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { KJ_IF_SOME(tracer, ioContext.getWorkerTracer()) { addExceptionToTrace(js, ioContext, tracer, UncaughtExceptionSource::REQUEST_HANDLER, error, api->getErrorInterfaceTypeHandler(js)); } } KJ_IF_SOME(i, impl->inspector) { jsg::sendExceptionToInspector(js, *i.get(), kj::str(desc), error, message); } // Run with --verbose to log JS exceptions to stderr. Useful when running tests. KJ_LOG(INFO, "uncaught exception", desc); }); } // By default, V8's memory pressure level is "none". This tells V8 that no one else on the // machine is competing for memory so it might as well use all it wants and be lazy about GC. // // In our production environment, however, we can safely assume that there is always memory // pressure, because every machine is handling thousands of tenants all the time. So we might // as well just throw the switch to "moderate" right away. lock->v8Isolate->MemoryPressureNotification(v8::MemoryPressureLevel::kModerate); // Register GC prologue and epilogue callbacks so that we can report GC CPU time via the // "request_context" Jaeger span. lock->v8Isolate->AddGCPrologueCallback( [](v8::Isolate* isolate, v8::GCType type, v8::GCCallbackFlags flags, void* data) noexcept { // We assume that a v8::Locker is alive during GC. KJ_DASSERT(v8::Locker::IsLocked(isolate)); auto& self = *reinterpret_cast(data); // However, currentLock might not be available, if (like in our Worker::Isolate constructor) we // don't use a Worker::Isolate::Impl::Lock. KJ_IF_SOME(currentLock, self.impl->currentLock) { currentLock.gcPrologue(); } }, this); lock->v8Isolate->AddGCEpilogueCallback( [](v8::Isolate* isolate, v8::GCType type, v8::GCCallbackFlags flags, void* data) noexcept { // We make similar assumptions about v8::Locker and currentLock as in the prologue callback. KJ_DASSERT(v8::Locker::IsLocked(isolate)); auto& self = *reinterpret_cast(data); KJ_IF_SOME(currentLock, self.impl->currentLock) { currentLock.gcEpilogue(); } }, this); lock->v8Isolate->SetPromiseRejectCallback([](v8::PromiseRejectMessage message) { // TODO(cleanup): IoContext doesn't really need to be involved here. We are trying to call // a method of ServiceWorkerGlobalScope, which is the context object. So we should be able to // do something like unwrap(lock, isolate->GetCurrentContext()).emitPromiseRejection(). // However, JSG doesn't currently provide an easy way to do this. KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { try { ioContext.getCurrentLock().reportPromiseRejectEvent(message); } catch (jsg::JsExceptionThrown&) { // V8 expects us to just return. return; } } }); // The PromiseCrossContextCallback is used to allow cross-IoContext promise following. // When the IoContext scope is entered, we set the "promise context tag" associated // with the IoContext on the Isolate that is locked. Any Promise that is created within // that scope will be tagged with the same promise context tag. When an attempt to // follow a promise occurs (e.g. either using Promise.prototype.then() or await, etc) // our patched v8 logic will check to see if the followed promise's tag matches the // current Isolate tag. If they do not, then v8 will invoke this callback. The promise // here is the promise that belongs to a different IoContext. lock->v8Isolate->SetPromiseCrossContextCallback( [](v8::Local context, v8::Local promise, v8::Local tag) -> v8::MaybeLocal { auto& js = jsg::Lock::current(); try { // Generally this condition is only going to happen when using dynamic imports. // It should not be common. JSG_REQUIRE(IoContext::hasCurrent(), Error, "Unable to wait on a promise created within a request when not running within a " "request."); return js.wrapSimplePromise( addCrossThreadPromiseWaiter(js, promise) .then(js, [promise = js.v8Ref(promise.As())](auto& js) mutable { // Once the waiter has been resolved, return the now settled promise. // Since the promise has been settled, it is now safe to access from // other requests. Note that the resolved value of the promise still // might not be safe to access! (e.g. if it contains any IoOwns attached // to the other request IoContext). return kj::mv(promise); })); } catch (jsg::JsExceptionThrown&) { // Exceptions here are generally unexpected but possible because the jsg::Promise // then can fail if the isolate is in the process of being torn down. Let's just // return control back to V8 which should handle the case. return v8::MaybeLocal(); } catch (...) { auto ex = kj::getCaughtExceptionAsKj(); KJ_LOG(ERROR, "Setting promise cross context follower failed unexpectedly", ex); jsg::throwInternalError(js.v8Isolate, kj::mv(ex)); return v8::MaybeLocal(); } }); // The PromiseCrossContextResolveCallback is used to ensure that promise reactions // are only scheduled on the microtask queue from the appropriate IoContext for the // promise. Huh? Yeah, that's not super clear... let me explain a bit more. // Every request runs in its own IoContext. // Some I/O objects are bound to the IoContext when they are created. // If these objects are accessed from the wrong IoContext, things blow up. // If I create a promise in one request and pass the resolve/reject functions // off to a different request, bad things can happen because the IoContext can // actually change in the promise continuation. Take the following case for example: // // In request one: // // const ab = AbortSignal.abort(); // AbortSignal is bound to the IoContext // const { promise, resolve } = Promise.withResolvers(); // globalThis.resolve = resolve; // await promise; // console.log(ab.aborted); // // In request two: // // globalThis.resolve(); // // What previously would happen is that the `console.log(ab.aborted) after the // `await promise` in request one would fail with an error because the current // IoContext would change! (it would be the IoContext from request two!). // // That's bad. // // So this callback is added to ensure that the promise reactions for the promise // being resolved are not scheduled until we are back in the correct IoContext for // the promise. // // This happens by (ab)using the DeleteQueue that is specific to the owning // IoContext. When the IoContext is entered, the isolate is updated with a // current "promise tag". Whenever a promise is created, it is associated with // the isolate's current tag. Whenever a promise is followed (calling .then, etc), // we check the tag and arrange for a cross-thread resolve. When the promise is // resolved or rejected, we check the tag also. If the promise tag and the current // isolate tag do not match, the function below is called. if (features.getHandleCrossRequestPromiseResolution()) { lock->v8Isolate->SetPromiseCrossContextResolveCallback( [](v8::Isolate* isolate, v8::Local tag, v8::Local reactions, v8::Local argument, std::function reactions, v8::Local argument)> callback) -> v8::Maybe { try { auto& js = jsg::Lock::from(isolate); // The promise tag is generally opaque except for right here. The tag // wraps an instanceof kj::Own, which wraps an atomically // refcounted pointer to the DeleteQueue for the correct isolate. // We simply pass the given callback, reactions, and argument to // a function that will be added to the queue inside DeleteQueue. // The next time the relevant IoContext is entered, this queue will // be drained and the actions will be run. Adding the task to the // delete queue will also signal the IoContext that it should wake // up and drain the queue. Simple, eh? // // A word of warning tho! It is possible for the IoContext to be // destroyed before the promise is resolved. Any actions that have // already been added to the queue would end up being dropped silently // on the floor. Actions that are added to the queue now will be run // immediately in the wrong IoContext. auto& ref = jsg::unwrapOpaqueRef>(isolate, tag); ref->execute(js, [reactions = jsg::Data(isolate, reactions), argument = jsg::V8Ref(isolate, argument), callback = kj::mv(callback)](jsg::Lock& js) mutable { callback(js.v8Isolate, reactions.getHandle(js), argument.getHandle(js)); }); return v8::JustVoid(); } catch (jsg::JsExceptionThrown&) { // Exceptions here are generally unexpected but possible because the jsg::Promise // then can fail if the isolate is in the process of being torn down. Let's just // return control back to V8 which should handle the case. // Note that errors thrown here and below should cause the resolve() or reject() // function calls to throw, which is unusual. Just important to keep that in mind. // Most likely errors thrown here are fatal so that should be OK. return v8::Nothing(); } catch (...) { jsg::throwInternalError(isolate, kj::getCaughtExceptionAsKj()); return v8::Nothing(); } }); } }); } Worker::Script::Script(kj::Own isolateParam, kj::StringPtr id, const Script::Source& source, IsolateObserver::StartType startType, bool logNewScript, kj::Maybe errorReporter, kj::Maybe> artifacts, SpanParent parentSpan, kj::Own vfs, kj::Maybe> maybeNewModuleRegistry) : isolate(kj::mv(isolateParam)), id(kj::str(id)), modular(source.variant.is()), python(modular && source.variant.get().isPython), impl(kj::heap(kj::mv(vfs), kj::mv(maybeNewModuleRegistry))), dynamicEnvBuilder(source.dynamicEnvBuilder.map( [](const auto& inst) -> kj::Arc { return inst.addRef(); })) { auto parseMetrics = isolate->metrics->parse(startType); // TODO(perf): It could make sense to take an async lock when constructing a script if we // co-locate multiple scripts in the same isolate. As of this writing, we do not, except in // previews, where it doesn't matter. If we ever do co-locate multiple scripts in the same // isolate, we may wish to make the RequestObserver object available here, in order to // attribute lock timing to that request. jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) { Isolate::Impl::Lock recordedLock( *isolate, Worker::Lock::TakeSynchronously(kj::none), stackScope); auto& lock = *recordedLock.lock; // If we throw an exception, it's important that `impl` is destroyed under lock. KJ_ON_SCOPE_FAILURE({ auto implToDestroy = kj::mv(impl); KJ_IF_SOME(c, implToDestroy->moduleContext) { recordedLock.disposeContext(kj::mv(c)); } else { // Else block to avoid dangling else clang warning. } }); lock.withinHandleScope([&] { if (isolate->impl->inspector != kj::none || errorReporter != kj::none) { lock.v8Isolate->SetCaptureStackTraceForUncaughtExceptions(true); } v8::Local context; if (modular) { // Modules can't be compiled for multiple contexts. We need to create the real context now. auto& mContext = impl->moduleContext.emplace(isolate->getApi().newContext(lock, { .newModuleRegistry = impl->getNewModuleRegistry(), .schemaLoader = getSchemaLoader(), })); mContext->enableWarningOnSpecialEvents(); context = mContext.getHandle(lock); recordedLock.setupContext(context); } else { // Although we're going to compile a script independent of context, V8 requires that // there be an active context, otherwise it will segfault, I guess. So we create a // dummy context. (Undocumented, as usual.) context = v8::Context::New(lock.v8Isolate, nullptr, v8::ObjectTemplate::New(lock.v8Isolate)); // We need to set the highest used index in every context we create to be a nullptr // This is because we might later on call GetAlignedPointerFromEmbedderData which fails with // a fatal error if the array is smaller than the given index. jsg::setAlignedPointerInEmbedderData( context, jsg::ContextPointerSlot::MAX_POINTER_SLOT, nullptr); } JSG_WITHIN_CONTEXT_SCOPE(lock, context, [&](jsg::Lock& js) { // const_cast OK because we hold the isolate lock. Worker::Isolate& lockedWorkerIsolate = const_cast(*isolate); if (logNewScript) { // HACK: Log a message indicating that a new script was loaded. This is used only when the // inspector is enabled. We want to do this immediately after the context is created, // before the user gets a chance to modify the behavior of the console, which if they // did, we'd then need to be more careful to apply time limits and such. lockedWorkerIsolate.logMessage(lock, static_cast(cdp::LogType::WARNING), "Script modified; context reset."); } // We need to register this context with the inspector, otherwise errors won't be // reported. But we want it to be un-registered as soon as the script has been // compiled, otherwise the inspector will end up with multiple contexts active which // is very confusing for the user (since they'll have to select from the drop-down // which context to use). // // (For modules, the context was already registered by `setupContext()`, above. KJ_IF_SOME(i, isolate->impl->inspector) { if (!modular) { i.get()->contextCreated( v8_inspector::V8ContextInfo(context, 1, jsg::toInspectorStringView("Compiler"))); } } else { } // Here to squash a compiler warning KJ_DEFER({ if (!modular) { KJ_IF_SOME(i, isolate->impl->inspector) { i.get()->contextDestroyed(context); } else { } // Here to squash a compiler warning } }); v8::TryCatch catcher(lock.v8Isolate); ExceptionOrDuration limitErrorOrTime = 0 * kj::NANOSECONDS; try { try { KJ_SWITCH_ONEOF(source.variant) { KJ_CASE_ONEOF(script, ScriptSource) { // This path is used for the older, service worker syntax workers. if (script.capnpSchemas.size() > 0) { // const_cast OK because we hold the isolate lock. auto& schemaLoader = const_cast(getSchemaLoader()); for (auto node: script.capnpSchemas) { schemaLoader.load(node); } } impl->globals = isolate->getApi().compileServiceWorkerGlobals(lock, script, *isolate); { // It's unclear to me if CompileUnboundScript() can get trapped in any // infinite loops or excessively-expensive computation requiring a time // limit. We'll go ahead and apply a time limit just to be safe. Don't // add it to the rollover bank, though. auto limitScope = isolate->getLimitEnforcer().enterStartupJs(lock, limitErrorOrTime); impl->unboundScriptOrMainModule = jsg::NonModuleScript::compile(lock, script.mainScript, script.mainScriptName); } } KJ_CASE_ONEOF(modulesSource, ModulesSource) { // This path is used for the new ESM worker syntax. if (modulesSource.capnpSchemas.size() > 0) { // const_cast OK because we hold the isolate lock. auto& schemaLoader = const_cast(getSchemaLoader()); for (auto node: modulesSource.capnpSchemas) { schemaLoader.load(node); } } if (!isolate->getApi().getFeatureFlags().getNewModuleRegistry()) { kj::Own limitScope; if (modulesSource.isPython) { limitScope = isolate->getLimitEnforcer().enterStartupPython(js, limitErrorOrTime); } else { limitScope = isolate->getLimitEnforcer().enterStartupJs(js, limitErrorOrTime); } impl->configureDynamicImports(lock, *jsg::ModuleRegistry::from(lock)); isolate->getApi().compileModules( lock, modulesSource, *isolate, kj::mv(artifacts), parentSpan.addRef()); } impl->unboundScriptOrMainModule = kj::Path::parse(modulesSource.mainModule); } } parseMetrics->done(); } catch (const kj::Exception& e) { lock.throwException(e.clone()); // lock.throwException() here will throw a jsg::JsExceptionThrown which we catch // in the outer try/catch. } } catch (const jsg::JsExceptionThrown&) { reportStartupError(id, lock, isolate->impl->inspector, isolate->getLimitEnforcer(), kj::mv(limitErrorOrTime), catcher, errorReporter, impl->permanentException, parentSpan.addRef(), dynamicEnvBuilder != kj::none); } }); }); }); } void Worker::Script::installVirtualFileSystemOnContext(v8::Local context) const { jsg::setAlignedPointerInEmbedderData(context, jsg::ContextPointerSlot::VIRTUAL_FILE_SYSTEM, const_cast(impl->vfs.get())); } const capnp::SchemaLoader& Worker::Script::getSchemaLoader() const { KJ_IF_SOME(moduleRegistry, impl->maybeNewModuleRegistry) { return moduleRegistry->getSchemaLoader(); } else { return *KJ_ASSERT_NONNULL(impl->maybeSchemaLoader); } } kj::Own Worker::Isolate::getWeakRef() const { return weakIsolateRef->addRef(); } kj::StringPtr Worker::Isolate::getUuid() const { // As of this writing, getUuid() is only used by actors, for metrics. We don't want to bother // generating it if not used. The call site does not have nor want an isolate lock, so we use a // kj::Lazy to make initialization thread-safe. return impl->uuid.get( [](kj::SpaceFor& space) { return space.construct(randomUUID(kj::none)); }); } Worker::Isolate::~Isolate() noexcept(false) { metrics->teardownStarted(); // Update the isolate stats one last time to make sure we're accurate for cleanup in // `evicted()`. limitEnforcer->reportMetrics(*metrics); metrics->evicted(); weakIsolateRef->invalidate(); // The cpuLimitNearlyExceededCallback may hold references to objects owned by the isolate and // their destructors need the isolate to still exist. So destroy them before we destroy the // isolate. *cpuLimitNearlyExceededCallback.lockExclusive() = kj::none; // Make sure to destroy things under lock. This lock should never be contended since the isolate // is about to be destroyed, but we have to take the lock in order to enter the isolate. // It's also important that we lock one last time, in order to destroy any remaining workers in // worker destruction queue. jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) { Isolate::Impl::Lock recordedLock(*this, Worker::Lock::TakeSynchronously(kj::none), stackScope); metrics->teardownLockAcquired(); auto inspector = kj::mv(impl->inspector); auto dropTraceAsyncContextKey = kj::mv(traceAsyncContextKey); auto dropUserTraceAsyncContextKey = kj::mv(userTraceAsyncContextKey); // The Rust Realm must be dropped under lock since Realm::drop() accesses V8 globals // and calls drop functions that may interact with V8. auto dropRealm = kj::mv(impl->realm); // Release all tracked WASM instance entries while V8 is still alive. Each entry holds a // shared_ptr whose destructor who needs the isolate to still be alive. // This is analogous to the cpuTimeLimitNearlyExceededCallback detaching above ^^^ limitEnforcer->getTrackedWasmInstances().clear(*recordedLock.lock); }); } Worker::Script::~Script() noexcept(false) { // Make sure to destroy things under lock. // TODO(perf): It could make sense to try to obtain an async lock before destroying a script if // multiple scripts are co-located in the same isolate. As of this writing, that doesn't happen // except in preview. In any case, Scripts are destroyed in the GC thread, where we don't care // too much about lock latency. jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) { Isolate::Impl::Lock recordedLock( *isolate, Worker::Lock::TakeSynchronously(kj::none), stackScope); KJ_IF_SOME(c, impl->moduleContext) { recordedLock.disposeContext(kj::mv(c)); } impl = nullptr; }); } const Worker::Isolate& Worker::Isolate::from(jsg::Lock& js) { auto ptr = js.v8Isolate->GetData(jsg::SET_DATA_ISOLATE); KJ_ASSERT(ptr != nullptr); return *static_cast(ptr); } bool Worker::Isolate::Impl::Lock::checkInWithLimitEnforcer(Worker::Isolate& isolate) { shouldReportIsolateMetrics = true; return limitEnforcer.exitJs(*lock); } kj::Maybe> Worker::Isolate::getCpuLimitNearlyExceededCallback() const { auto lock = cpuLimitNearlyExceededCallback.lockExclusive(); KJ_IF_SOME(cb, *lock) { return cb.reference(); } return kj::none; } void Worker::Isolate::setCpuLimitNearlyExceededCallback(kj::Function cb) const { auto lock = cpuLimitNearlyExceededCallback.lockExclusive(); // Make sure we don't reassign the callback so we don't invalidate references we've passed out. if (*lock == kj::none) { *lock = kj::mv(cb); return; } kj::throwRecoverableException(KJ_EXCEPTION( FAILED, "Python Workers Internal Error: CpuLimitNearlyExceededCallback already set")); } void Worker::Isolate::registerTrackedWasmInstance(jsg::Lock& js, v8::Local instance, kj::Array memory, kj::Maybe signalOffset, kj::Maybe terminatedOffset) const { // Register the WASM module for receiving shutdown signals. The signal handler will // iterate the list unconditionally when CPU time is nearly exhausted. KJ_IF_SOME(entry, limitEnforcer->getTrackedWasmInstances().registerSignal( js, kj::mv(memory), signalOffset, terminatedOffset)) { // Set up a weak reference to the instance. When V8 collects it, the handle becomes // empty and the GC prologue filter removes the entry, releasing the strong memory ref. entry.instanceRef.Reset(js.v8Isolate, instance); entry.instanceRef.SetWeak(); } } // EW-1319: Set WebAssembly.Module @@HasInstance // // The instanceof operator can be changed by setting the @@HasInstance method // on the object, https://tc39.es/ecma262/#sec-instanceofoperator. void setWebAssemblyModuleHasInstance(jsg::Lock& lock, v8::Local context) { JSG_WITHIN_CONTEXT_SCOPE(lock, context, [&](jsg::Lock& lock) { auto instanceof = [](const v8::FunctionCallbackInfo& info) { jsg::Lock::from(info.GetIsolate()).withinHandleScope([&] { info.GetReturnValue().Set(info[0]->IsWasmModuleObject()); }); }; v8::Local function = jsg::check(v8::Function::New(context, instanceof)); auto webAssembly = KJ_ASSERT_NONNULL(lock.global().get(lock, "WebAssembly").tryCast()); auto module = KJ_ASSERT_NONNULL(webAssembly.get(lock, "Module").tryCast()); jsg::check(v8::Local(module)->DefineOwnProperty( context, v8::Symbol::GetHasInstance(lock.v8Isolate), function)); }); } // Installs a shim around WebAssembly.instantiate and WebAssembly.Instance that hooks into the // shutdown signal if it exists void shimWebAssemblyInstantiate(jsg::Lock& lock, v8::Local context) { // We need to enter the context because this function compiles and executes JavaScript via // v8::Script::Compile/Run. setupContext() is called before JSG_WITHIN_CONTEXT_SCOPE, so the // context is not yet entered at this point. v8::Context::Scope contextScope(context); // Create a C++ callback that the JS shims call to register a {instance, memory, signalOffset, // terminatedOffset} tuple. // __registerTrackedWasmInstance(instance: WebAssembly.Instance, // memory: WebAssembly.Memory, signalOffset: number, // terminatedOffset: number) // signalOffset or terminatedOffset may be -1, indicating the corresponding export is absent. auto registerCb = [](const v8::FunctionCallbackInfo& info) { auto& js = jsg::Lock::from(info.GetIsolate()); js.withinHandleScope([&] { if (info.Length() < 4 || !info[0]->IsObject() || !info[1]->IsWasmMemoryObject() || !info[2]->IsNumber() || !info[3]->IsNumber()) { js.v8Isolate->ThrowException( js.str("registerTrackedWasmInstance: expected " "(WebAssembly.Instance, WebAssembly.Memory, number, number)"_kj)); return; } auto instance = info[0].As(); auto memory = info[1].As(); // signalOffset is -1 when __instance_signal was not exported. auto signalRaw = info[2].As()->Value(); kj::Maybe signalOffset; if (signalRaw >= 0) { signalOffset = static_cast(signalRaw); } // terminatedOffset is -1 when __instance_terminated was not exported. auto terminatedRaw = info[3].As()->Value(); kj::Maybe terminatedOffset; if (terminatedRaw >= 0) { terminatedOffset = static_cast(terminatedRaw); } auto backingStore = memory->Buffer()->GetBackingStore(); auto wasmMemory = kj::arrayPtr(static_cast(backingStore->Data()), backingStore->ByteLength()) .attach(kj::mv(backingStore)); KJ_IF_SOME(e, kj::runCatchingExceptions([&] { Worker::Isolate::from(js).registerTrackedWasmInstance( js, instance, kj::mv(wasmMemory), signalOffset, terminatedOffset); })) { js.v8Isolate->ThrowException(js.exceptionToJs(kj::mv(e)).getHandle(js)); } }); }; auto registerFn = jsg::check(v8::Function::New(context, registerCb)); // Build the shim in JavaScript. It wraps both WebAssembly.instantiate (async) and // WebAssembly.Instance (sync constructor). auto shimScript = jsg::NonModuleScript::compile(lock, WASM_INSTANTIATE_SHIM, "wasm-instantiate-shim.js"_kj); auto shimFn = KJ_ASSERT_NONNULL(shimScript.runAndReturn(lock).tryCast()); // Call the factory — it mutates `WebAssembly` in place. shimFn.call(lock, lock.global(), jsg::JsFunction(registerFn)); } void Worker::setupContext( jsg::Lock& lock, v8::Local context, const LoggingOptions& loggingOptions) { // Set WebAssembly.Module @@HasInstance setWebAssemblyModuleHasInstance(lock, context); // Shim WebAssembly.instantiate to detect modules exporting "__instance_signal". if (util::Autogate::isEnabled(util::AutogateKey::WASM_SHUTDOWN_SIGNAL_SHIM)) { shimWebAssemblyInstantiate(lock, context); } // We replace the default V8 console.log(), etc. methods, to give the worker access to // logged content, and log formatted values to stdout/stderr locally. auto global = context->Global(); auto consoleStr = jsg::v8StrIntern(lock.v8Isolate, "console"); auto console = jsg::check(global->Get(context, consoleStr)).As(); auto setHandler = [&](const char* method, LogLevel level) { auto methodStr = jsg::v8StrIntern(lock.v8Isolate, method); v8::Global original( lock.v8Isolate, jsg::check(console->Get(context, methodStr)).As()); auto f = lock.wrapSimpleFunction(context, [loggingOptions, level, original = kj::mv(original)]( jsg::Lock& js, const v8::FunctionCallbackInfo& info) { handleLog(js, loggingOptions, level, original, info); }); jsg::check(console->Set(context, methodStr, f)); }; setHandler("debug", LogLevel::DEBUG_); setHandler("error", LogLevel::ERROR); setHandler("info", LogLevel::INFO); setHandler("log", LogLevel::LOG); setHandler("warn", LogLevel::WARN); } // ======================================================================================= namespace { kj::Maybe tryResolveMainModule(jsg::Lock& js, const kj::Path& mainModule, jsg::JsContext& jsContext, const Worker::Script& script, ExceptionOrDuration& limitErrorOrTime) { kj::Own limitScope; if (script.isPython()) { limitScope = script.getIsolate().getLimitEnforcer().enterStartupPython(js, limitErrorOrTime); } else { limitScope = script.getIsolate().getLimitEnforcer().enterStartupJs(js, limitErrorOrTime); } KJ_DEFER({ if (limitErrorOrTime.is()) { // If we hit the limit in PerformMicrotaskCheckpoint() we may not have actually // thrown an exception. throw jsg::JsExceptionThrown(); } }); // Before resolving the main module, if both nodejs_compat_v2 and the new // module registry are enabled, let's pre-resolve the process and buffer modules. // Why? Great question! Resolving these modules synchronously causes the microtask // queue to be pumped, which we don't actually want to do while resolving the main // module until we are ready. Both process and buffer are exposed via globalThis // when the nodejs_compat_v2 flag is used, and if the top-level scope is accessing // either globalThis.process or globalThis.buffer, then we need to make sure that // the modules are already resolved so we don't pump the microtask queue while // synchronously accessing those globals. Resolving them here ensures that they are // ready to go before we begin evaluating the main module. auto featureFlags = FeatureFlags::get(js); if (featureFlags.getNodeJsCompatV2() && featureFlags.getNewModuleRegistry()) { JSG_REQUIRE_NONNULL(js.resolveModule("node:process", jsg::RequireEsm::YES), Error, "Failed to initialize node:process module"); JSG_REQUIRE_NONNULL(js.resolveModule("node:buffer", jsg::RequireEsm::YES), Error, "Failed to initialize node:buffer module"); } // When enable_nodejs_global_timers is enabled, load the module that makes all 6 timer // functions (setTimeout, setInterval, clearTimeout, clearInterval, setImmediate, // clearImmediate) available on globalThis as Node.js-compatible versions from node:timers. if (featureFlags.getEnableNodejsGlobalTimers()) { JSG_REQUIRE_NONNULL(js.resolveInternalModule("node-internal:internal_timers_global_override"), Error, "Failed to initialize node-internal:internal_timers_global_override module"); } return js.resolveModule(mainModule.toString(false), jsg::RequireEsm::YES); } } // anonymous namespace Worker::Worker(kj::Own scriptParam, kj::Own metricsParam, kj::FunctionParam target, v8::Local ctxExports)> compileBindings, IsolateObserver::StartType startType, SpanParent parentSpan, LockType lockType, kj::Maybe errorReporter, kj::Maybe startupTime) : script(kj::mv(scriptParam)), metrics(kj::mv(metricsParam)), impl(kj::heap()) { // Enter/lock isolate. jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) { Isolate::Impl::Lock recordedLock(*script->isolate, lockType, stackScope); auto& lock = *recordedLock.lock; // If we throw an exception, it's important that `impl` is destroyed under lock. KJ_ON_SCOPE_FAILURE({ auto implToDestroy = kj::mv(impl); KJ_IF_SOME(c, implToDestroy->context) { recordedLock.disposeContext(kj::mv(c)); } else { // Else block to avoid dangling else clang warning. } }); auto maybeMakeSpan = [&](auto operationName) -> SpanBuilder { auto span = parentSpan.newChild(kj::mv(operationName)); if (span.isObserved()) { span.setTag("truncated_script_id"_kjc, truncateScriptId(script->getId())); } return span; }; auto currentSpan = maybeMakeSpan("lw:new_startup_metrics"_kjc); auto startupMetrics = metrics->startup(startType); currentSpan = maybeMakeSpan("lw:new_context"_kjc); // Create a stack-allocated handle scope. lock.withinHandleScope([&] { jsg::JsContext* jsContext; KJ_IF_SOME(c, script->impl->moduleContext) { // Use the shared context from the script. // const_cast OK because guarded by `lock`. jsContext = const_cast*>(&c); currentSpan.setTag("module_context"_kjc, true); } else { // Create a new context. jsContext = &this->impl->context.emplace(script->isolate->getApi().newContext(lock, { .newModuleRegistry = script->impl->getNewModuleRegistry(), .schemaLoader = script->getSchemaLoader(), })); } v8::Local context = KJ_REQUIRE_NONNULL(jsContext).getHandle(lock); // Install the virtual file system on the context. Keep in mind that for service // worker style workers, the Script may be shared between multiple Workers, even // across different accounts. Currently, the internal state of the VFS does not // contain any account-specific or worker-specific state so this is OK for now. // The VFS would contain the script files only and any temporary files created // within the context of a worker are always stored in temporary space attached // to the IoContext or the current execution context. If we extend these capabilities // in the future, we may need to revisit this. For modular workers, this is not // an issue since each Worker gets its own Script instance. script->installVirtualFileSystemOnContext(context); if (!script->modular) { recordedLock.setupContext(context); } if (script->impl->unboundScriptOrMainModule == nullptr) { // Script failed to parse. Act as if the script was empty -- i.e. do nothing. impl->permanentException = script->impl->permanentException.map([](auto& e) { return e.clone(); }); return; } // Enter the context for compiling and running the script. JSG_WITHIN_CONTEXT_SCOPE(lock, context, [&](jsg::Lock& js) { v8::TryCatch catcher(lock.v8Isolate); ExceptionOrDuration limitErrorOrTime = 0 * kj::NANOSECONDS; try { try { currentSpan = maybeMakeSpan("lw:globals_instantiation"_kjc); v8::Local bindingsScope; if (script->isModular()) { // Use `env` variable. bindingsScope = v8::Object::New(lock.v8Isolate); if (!FeatureFlags::get(js).getDisableImportableEnv()) { lock.setWorkerEnv(lock.v8Ref(bindingsScope)); } } else { // Use global-scope bindings. bindingsScope = context->Global(); } // Load globals. // const_cast OK because we hold the lock. for (auto& global: const_cast(*script).impl->globals) { lock.v8Set(bindingsScope, global.name, global.value); } v8::Local ctxExports = v8::Object::New(lock.v8Isolate); compileBindings(lock, script->isolate->getApi(), bindingsScope, ctxExports); // Execute script. currentSpan = maybeMakeSpan("lw:top_level_execution"_kjc); // Ensure that our worker top-level bootstrap has a temporary directory // storage scope. This is used to store temporary files created within // the top-level evaluation of the worker. With this instantiated on // the stack, temporary files will be cleaned up when the scope is // destroyed, which means any temporary files created in the top-level // evaluation will *not* be available to the worker after the top-level // evaluation is complete. TmpDirStoreScope tmpDirStoreScope; // We allow eval and new Function() during startup, becaues startup time is entirely // deterministic, so we can easily reproduce the input to eval() by just running the // worker again. We do not allow eval() at runtime because we need to have a record of // all code that executes in production for forensic purposes, and at runtime the input // to eval() could have come from a remote source on which we don't have a record. js.setAllowEval(FeatureFlags::get(js).getAllowEvalDuringStartup()); KJ_DEFER(js.setAllowEval(false)); KJ_SWITCH_ONEOF(script->impl->unboundScriptOrMainModule) { KJ_CASE_ONEOF(unboundScript, jsg::NonModuleScript) { auto limitScope = script->isolate->getLimitEnforcer().enterStartupJs(lock, limitErrorOrTime); unboundScript.run(lock); // Flush microtasks enqueued during top-level script evaluation. // Without this flush, microtasks (e.g. promise continuations from async // initialization) remain on the per-isolate microtask queue and can leak across // V8 contexts when multiple Workers share an isolate (same script, different // zones). The leaked microtasks then execute under the wrong IoContext, making // things go boom. lock.runMicrotasks(); } KJ_CASE_ONEOF(mainModule, kj::Path) { KJ_IF_SOME(ns, tryResolveMainModule(lock, mainModule, *jsContext, *script, limitErrorOrTime)) { impl->env = lock.v8Ref(bindingsScope.As()); impl->ctxExports = lock.v8Ref(ctxExports.As()); if (!FeatureFlags::get(js).getDisableImportableEnv()) { lock.setWorkerExports(lock.v8Ref(ctxExports)); } auto& api = script->isolate->getApi(); auto handlers = api.unwrapExports(lock, ns); auto entrypointClasses = api.getEntrypointClasses(lock); for (auto& handler: handlers.fields) { KJ_SWITCH_ONEOF(handler.value) { KJ_CASE_ONEOF(obj, api::ExportedHandler) { obj.env = lock.v8Ref(bindingsScope.As()); // Historically, non-class-based handlers reused the same ctx object for all requests. // This was an accident, but some Workers depend on it. // Newer worker with the unique_ctx_per_invocation will allocate a new ctx for every request. obj.ctx = js.alloc(lock, jsg::JsValue(ctxExports)); // Python Workers append all durable objects, worker entrypoint and workflow // entrypoint classes in the pythonEntrypoints named export. bool isPythonWorker = FeatureFlags::get(js).getPythonWorkers(); if (handler.name == "pythonEntrypoints" && isPythonWorker) { auto handle = obj.self.getHandle(js); auto dict = js.toDict(handle); for (auto& field: dict.fields) { auto unwrapped = api.unwrapExport(lock, field.value); KJ_SWITCH_ONEOF(unwrapped) { KJ_CASE_ONEOF(cls, EntrypointClass) { processEntrypointClass( js, kj::mv(cls), entrypointClasses, kj::mv(field.name)); } KJ_CASE_ONEOF(obj, api::ExportedHandler) { KJ_FAIL_ASSERT("Expected EntrypointClass"); } } } } else { impl->namedHandlers.insert(kj::mv(handler.name), kj::mv(obj)); } } KJ_CASE_ONEOF(cls, EntrypointClass) { processEntrypointClass( js, kj::mv(cls), entrypointClasses, kj::mv(handler.name)); } } } } else { JSG_FAIL_REQUIRE(TypeError, "Main module name is not present in bundle."); } } } KJ_IF_SOME(s, startupTime) { KJ_SWITCH_ONEOF(limitErrorOrTime) { KJ_CASE_ONEOF(startupTimeElapsed, kj::Duration) { s = startupTimeElapsed; } KJ_CASE_ONEOF(limitError, kj::Exception) {} } } else { } startupMetrics->done(); } catch (const kj::Exception& e) { lock.throwException(e.clone()); // lock.throwException() here will throw a jsg::JsExceptionThrown which we catch // in the outer try/catch. } } catch (const jsg::JsExceptionThrown&) { reportStartupError(script->id, lock, script->isolate->impl->inspector, script->isolate->getLimitEnforcer(), kj::mv(limitErrorOrTime), catcher, errorReporter, impl->permanentException, currentSpan, script->getDynamicEnvBuilder() != kj::none); } }); // Reset this back to its default after startup execution // Leaving it on comes at the expense of collecting stack traces for all thrown exceptions // Ref: https://github.com/cloudflare/workerd/issues/5332 if (script->isolate->impl->inspector == kj::none) { lock.v8Isolate->SetCaptureStackTraceForUncaughtExceptions(false); } }); }); } Worker::~Worker() noexcept(false) { metrics->teardownStarted(); auto& isolateImpl = *script->getIsolate().impl; auto lock = isolateImpl.workerDestructionQueue.lockExclusive(); // Previously, this metric meant the isolate lock. We might as well make it mean the worker // destruction queue lock now to verify it is much less-contended than the isolate lock. metrics->teardownLockAcquired(); // Defer destruction of our V8 objects, in particular our jsg::Context, which requires some // finalization. lock->push(kj::mv(impl)); } void Worker::processEntrypointClass(jsg::Lock& js, EntrypointClass cls, EntrypointClasses entrypointClasses, kj::String handlerName) { js.withinHandleScope([&]() { jsg::JsObject handle(KJ_ASSERT_NONNULL(cls.tryGetHandle(js.v8Isolate))); for (;;) { if (handle == entrypointClasses.durableObject) { impl->actorClasses.insert(kj::mv(handlerName), ActorClassInfo{ .cls = kj::mv(cls), .missingSuperclass = false, }); return; } else if (handle == entrypointClasses.workerEntrypoint) { impl->statelessClasses.insert(kj::mv(handlerName), kj::mv(cls)); return; } else if (handle == entrypointClasses.workflowEntrypoint) { impl->workflowClasses.insert(kj::mv(handlerName), kj::mv(cls)); return; } handle = KJ_UNWRAP_OR(handle.getPrototype(js).tryCast(), { // Reached end of prototype chain. // For historical reasons, we assume a class is a Durable Object // class if it doesn't inherit anything. // TODO(someday): Log a warning suggesting extending DurableObject. // TODO(someday): Introduce a compat flag that makes this required. impl->actorClasses.insert(kj::mv(handlerName), ActorClassInfo{ .cls = kj::mv(cls), .missingSuperclass = true, }); return; }); } }); } void Worker::handleLog(jsg::Lock& js, const LoggingOptions& loggingOptions, LogLevel level, const v8::Global& original, const v8::FunctionCallbackInfo& info) { // Call original V8 implementation so messages sent to connected inspector if any auto context = js.v8Context(); int length = info.Length(); // to pass additional arguments from this function to js' `formatLog` we add arguments to the end // of the arguments vector, then in formatLog we `pop` these from the vector. // 3 is just the number of args we currently pass. v8::LocalVector args(js.v8Isolate, length + 3); for (auto i: kj::zeroTo(length)) args[i] = info[i]; jsg::check(original.Get(js.v8Isolate)->Call(context, info.This(), length, args.data())); // The TryCatch is initialized here to catch cases where the v8 isolate's execution is // terminating, usually as a result of an infinite loop. We need to perform the initialization // here because `message` is called multiple times. v8::TryCatch tryCatch(js.v8Isolate); auto message = [&]() { int length = info.Length(); kj::Vector stringified(length); for (auto i: kj::zeroTo(length)) { auto arg = info[i]; // serializeJson and v8::Value::ToString can throw JS exceptions // (e.g. for recursive objects) so we eat them here, to ensure logging and non-logging code // have the same exception behavior. if (!tryCatch.CanContinue()) { stringified.add(kj::str("{}")); break; } // The following code checks the `arg` to see if it should be serialised to JSON. // // We use the following criteria: if arg is null, a number, a boolean, an array, a string, an // object or it defines a `toJSON` property that is a function, then the arg gets serialised // to JSON. // // Otherwise we stringify the argument. js.withinHandleScope([&] { auto context = js.v8Context(); bool shouldSerialiseToJson = false; if (arg->IsNull() || arg->IsNumber() || arg->IsArray() || arg->IsBoolean() || arg->IsString() || arg->IsUndefined()) { // This is special cased for backwards compatibility. shouldSerialiseToJson = true; } if (arg->IsObject()) { v8::Local obj = arg.As(); v8::Local freshObj = v8::Object::New(js.v8Isolate); // Determine whether `obj` is constructed using `{}` or `new Object()`. This ensures // we don't serialise values like Promises to JSON. #if V8_MAJOR_VERSION >= 15 || (V8_MAJOR_VERSION == 14 && V8_MINOR_VERSION >= 7) if (obj->GetPrototype()->SameValue(freshObj->GetPrototype()) || obj->GetPrototype()->IsNull()) { #else // TODO(cleanup): Remove when unnecessary. if (obj->GetPrototypeV2()->SameValue(freshObj->GetPrototypeV2()) || obj->GetPrototypeV2()->IsNull()) { #endif shouldSerialiseToJson = true; } // Check if arg has a `toJSON` property which is a function. auto toJSONStr = jsg::v8StrIntern(js.v8Isolate, "toJSON"_kj); v8::MaybeLocal toJSON = obj->GetRealNamedProperty(context, toJSONStr); if (!toJSON.IsEmpty()) { if (jsg::check(toJSON)->IsFunction()) { shouldSerialiseToJson = true; } } } if (kj::runCatchingExceptions([&]() { // On the off chance the the arg is the request.cf object, let's make // sure we do not log proxied fields here. if (shouldSerialiseToJson) { auto s = js.serializeJson(arg); // serializeJson returns the string "undefined" for some values (undefined, // Symbols, functions). We remap these values to null to ensure valid JSON output. if (s == "undefined"_kj) { stringified.add(kj::str("null")); } else { stringified.add(kj::mv(s)); } } else { stringified.add(js.serializeJson(jsg::check(arg->ToString(context)))); } }) != kj::none) { stringified.add(kj::str("{}")); }; }); } return kj::str("[", kj::delimited(stringified, ", "_kj), "]"); }; // Only check tracing if console.log() was not invoked at the top level. KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { KJ_IF_SOME(tracer, ioContext.getWorkerTracer()) { auto timestamp = ioContext.now(); tracer.addLog(ioContext.getInvocationSpanContext(), timestamp, level, message()); } } if (loggingOptions.consoleMode == Worker::ConsoleMode::INSPECTOR_ONLY) { // Lets us dump console.log()s to stdout when running test-runner with --verbose flag, to make // it easier to debug tests. Note that when --verbose is not passed, KJ_LOG(INFO, ...) will // not even evaluate its arguments, so `message()` will not be called at all. KJ_LOG(INFO, "console.log()", message()); } else { // Write to stdio if allowed by console mode. This is making use of our internal // built-in implementation of the node:util inspect API. static const ColorMode COLOR_MODE = permitsColor(); #if _WIN32 static bool STDOUT_TTY = _isatty(_fileno(stdout)); static bool STDERR_TTY = _isatty(_fileno(stderr)); #else static bool STDOUT_TTY = isatty(STDOUT_FILENO); static bool STDERR_TTY = isatty(STDERR_FILENO); #endif // Log warnings and errors to stderr // Always log to stdout when structuredLogging is enabled. auto useStderr = level >= LogLevel::WARN && !loggingOptions.structuredLogging; auto fd = useStderr ? stderr : stdout; auto tty = useStderr ? STDERR_TTY : STDOUT_TTY; auto colors = COLOR_MODE == ColorMode::ENABLED || (COLOR_MODE == ColorMode::ENABLED_IF_TTY && tty); constexpr auto kSpecifier = "node-internal:internal_inspect"_kj; auto inspectModule = KJ_ASSERT_NONNULL(js.resolveInternalModule(kSpecifier)); v8::Local formatLogVal = inspectModule.get(js, "formatLog"_kj); KJ_ASSERT(formatLogVal->IsFunction()); auto formatLog = formatLogVal.As(); auto levelStr = logLevelToString(level); args[length] = js.boolean(colors); args[length + 1] = js.boolean(loggingOptions.structuredLogging.toBool()); args[length + 2] = js.strIntern(levelStr); auto formatted = js.toString( jsg::check(formatLog->Call(context, js.v8Undefined(), length + 3, args.data()))); fprintf(fd, "%s\n", formatted.cStr()); fflush(fd); } } Worker::Lock::TakeSynchronously::TakeSynchronously(kj::Maybe requestParam) { KJ_IF_SOME(r, requestParam) { request = &r; } } kj::Maybe Worker::Lock::TakeSynchronously::getRequest() { if (request != nullptr) { return *request; } return kj::none; } struct Worker::Lock::Impl { Isolate::Impl::Lock recordedLock; jsg::Lock& inner; Impl(const Worker& worker, LockType lockType, jsg::V8StackScope& stackScope) : recordedLock(worker.getIsolate(), lockType, stackScope), inner(*recordedLock.lock) {} }; Worker::Lock::Lock(const Worker& constWorker, LockType lockType, jsg::V8StackScope& stackScope) : // const_cast OK because we took out a lock. worker(const_cast(constWorker)), impl(kj::heap(worker, lockType, stackScope)) { kj::requireOnStack(this, "Worker::Lock MUST be allocated on the stack."); } Worker::Lock::~Lock() noexcept(false) { // const_cast OK because we hold -- nay, we *are* -- a lock on the script. auto& isolate = const_cast(worker.getIsolate()); if (impl->recordedLock.checkInWithLimitEnforcer(isolate)) { isolate.disconnectInspector(); } } void Worker::Lock::requireNoPermanentException() { KJ_IF_SOME(e, worker.impl->permanentException) { // Block taking lock when worker failed to start up. kj::throwFatalException(e.clone()); } } Worker::Lock::operator jsg::Lock&() { return impl->inner; } v8::Isolate* Worker::Lock::getIsolate() { return impl->inner.v8Isolate; } v8::Local Worker::Lock::getContext() { KJ_IF_SOME(c, worker.impl->context) { return c.getHandle(impl->inner); } else KJ_IF_SOME(c, const_cast(*worker.script).impl->moduleContext) { return c.getHandle(impl->inner); } else { KJ_UNREACHABLE; } } template static inline kj::Own fakeOwn(T& ref) { return kj::Own(&ref, kj::NullDisposer::instance); } kj::Maybe> Worker::Lock::getExportedHandler( kj::Maybe name, kj::Maybe versionInfo, Frankenvalue props, kj::Maybe actor, bool isDynamicDispatch) { KJ_IF_SOME(a, actor) { KJ_IF_SOME(h, a.getHandler()) { return fakeOwn(h); } } kj::StringPtr n = name.orDefault("default"_kj); auto getHandlerFromEntrypointClass = [&](EntrypointClass& cls) -> kj::Maybe> { jsg::Lock& js = *this; auto handler = kj::heap(cls(js, js.alloc(js, jsg::JsValue(KJ_ASSERT_NONNULL(worker.impl->ctxExports).getHandle(js)), props.toJs(js), kj::mv(versionInfo)), KJ_ASSERT_NONNULL(worker.impl->env).addRef(js))); // HACK: We set handler.env and handler.ctx to undefined because we already passed the real // env and ctx into the constructor, and we want the handler methods to act like they take // just one parameter. handler->env = js.v8Ref(js.v8Undefined()); handler->ctx = kj::none; return handler; }; KJ_IF_SOME(h, worker.impl->namedHandlers.find(n)) { jsg::Lock& js = *this; if (!FeatureFlags::get(js).getReuseCtxAcrossNonclassEvents()) { api::ExportedHandler constructedHandler = h.clone(js); constructedHandler.ctx = js.alloc(js, jsg::JsValue(KJ_ASSERT_NONNULL(worker.impl->ctxExports).getHandle(js)), props.toJs(js), kj::mv(versionInfo)); return kj::heap(kj::mv(constructedHandler)); } return fakeOwn(h); } else KJ_IF_SOME(cls, worker.impl->statelessClasses.find(n)) { return getHandlerFromEntrypointClass(cls); } else KJ_IF_SOME(cls, worker.impl->workflowClasses.find(n)) { return getHandlerFromEntrypointClass(cls); } else if (name == kj::none) { // If the default export was requested, and we didn't find a handler for it, we'll fall back // to addEventListener(). // // Note: The original intention was that we only use addEventListener() for // service-worker-syntax scripts, but apparently the code has long allowed it for // modules-based script too, if they lacked an `export default`. Yikes! Sadly, there are // Workers in production relying on this so we are stuck with it. return kj::none; } else { if (worker.impl->actorClasses.find(n) != kj::none) { if (isDynamicDispatch) { JSG_FAIL_REQUIRE(TypeError, "The entrypoint name ", n, " refers to a Durable Object class, but the incoming request is trying to invoke it as" " a stateless worker."); } else { LOG_ERROR_PERIODICALLY("worker is not an actor but class name was requested", n); } } else if (isDynamicDispatch) { JSG_FAIL_REQUIRE(TypeError, "The entrypoint name ", n, " was not found in this worker. Ensure the worker exports an entrypoint with that name."); } else { LOG_ERROR_PERIODICALLY("worker has no such named entrypoint", n); } KJ_FAIL_ASSERT("worker_do_not_log; Unable to get exported handler"); }; } api::ServiceWorkerGlobalScope& Worker::Lock::getGlobalScope() { return KJ_ASSERT_NONNULL(jsg::getAlignedPointerFromEmbedderData( getContext(), jsg::ContextPointerSlot::GLOBAL_WRAPPER)); } TimeoutId::Generator& Worker::Lock::getTimeoutIdGenerator() { return getGlobalScope().timeoutIdGenerator; } jsg::AsyncContextFrame::StorageKey& Worker::Lock::getTraceAsyncContextKey() { // const_cast OK because we are a lock on this isolate. auto& isolate = const_cast(worker.getIsolate()); return *(isolate.traceAsyncContextKey); } jsg::AsyncContextFrame::StorageKey& Worker::Lock::getUserTraceAsyncContextKey() { // const_cast OK because we are a lock on this isolate. auto& isolate = const_cast(worker.getIsolate()); return *(isolate.userTraceAsyncContextKey); } bool Worker::Lock::isInspectorEnabled() { return worker.script->isolate->impl->inspector != kj::none; } void Worker::Lock::logWarning(kj::StringPtr description) { // const_cast OK because we are a lock on this isolate. const_cast(worker.getIsolate()).logWarning(description, *this); } void Worker::Lock::logWarningOnce(kj::StringPtr description) { // const_cast OK because we are a lock on this isolate. const_cast(worker.getIsolate()).logWarningOnce(description, *this); } void Worker::Lock::logErrorOnce(kj::StringPtr description) { // const_cast OK because we are a lock on this isolate. const_cast(worker.getIsolate()).logErrorOnce(description); } void Worker::Lock::logUncaughtException(kj::StringPtr description) { // We don't add the exception to traces here, since it turns out that this path only gets hit by // intermediate exception handling. KJ_IF_SOME(i, worker.script->isolate->impl->inspector) { JSG_WITHIN_CONTEXT_SCOPE(*this, getContext(), [&](jsg::Lock& js) { jsg::sendExceptionToInspector(js, *i.get(), description); }); } // Run with --verbose to log JS exceptions to stderr. Useful when running tests. KJ_LOG(INFO, "uncaught exception", description); } void Worker::Lock::logUncaughtException( UncaughtExceptionSource source, const jsg::JsValue& exception, const jsg::JsMessage& message) { // Only add exception to trace when running within an I/O context with a tracer. KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { KJ_IF_SOME(tracer, ioContext.getWorkerTracer()) { JSG_WITHIN_CONTEXT_SCOPE(*this, getContext(), [&](jsg::Lock& js) { addExceptionToTrace(impl->inner, ioContext, tracer, source, exception, worker.getIsolate().getApi().getErrorInterfaceTypeHandler(*this)); }); } } KJ_IF_SOME(i, worker.script->isolate->impl->inspector) { JSG_WITHIN_CONTEXT_SCOPE(*this, getContext(), [&](jsg::Lock& js) { sendExceptionToInspector(js, *i.get(), source, exception, message); }); } // Run with --verbose to log JS exceptions to stderr. Useful when running tests. if (kj::_::Debug::shouldLog(::kj::LogSeverity::INFO)) { JSG_WITHIN_CONTEXT_SCOPE(*this, getContext(), [&](jsg::Lock& js) { // Try to log `error.stack` if it exists. KJ_IF_SOME(obj, exception.tryCast()) { auto stack = obj.get(js, "stack"); if (!stack.isUndefined()) { KJ_LOG(INFO, "uncaught exception", source, stack); return; } } else { // Compiler gives a spurious warning if this `else` isn't here. } KJ_LOG(INFO, "uncaught exception", source, exception); }); } } void Worker::Lock::logUncaughtException(UncaughtExceptionSource source, kj::Exception&& exception) { jsg::Lock& js = *this; try { auto jsError = js.exceptionToJsValue(kj::mv(exception), { .trusted = true, }); logUncaughtException(source, jsError.getHandle(js)); } catch (const jsg::JsExceptionThrown&) { // An exception occurred while trying to convert the exception to a JS value. // With exceptionToJs, this should only happen if the isolate is terminating // because of a fatal error when trying to deserialize a tunneled exception // detail. In this case, we will want to log the original exception instead, // so let's try exceptionToJs again but this time ignoring the detail, and // if it throws again, we'll give up and propagate that exception to the // caller. auto jsError = js.exceptionToJsValue(exception.clone(), {.ignoreDetail = true}); logUncaughtException(source, jsError.getHandle(js)); } } void Worker::Lock::reportPromiseRejectEvent(v8::PromiseRejectMessage& message) { getGlobalScope().emitPromiseRejection(*this, message.GetEvent(), jsg::V8Ref(getIsolate(), message.GetPromise()), jsg::V8Ref(getIsolate(), message.GetValue())); } void Worker::Lock::validateHandlers(ValidationErrorReporter& errorReporter) { JSG_WITHIN_CONTEXT_SCOPE(*this, getContext(), [&](jsg::Lock& js) { kj::HashSet ignoredHandlers; ignoredHandlers.insert("alarm"_kj); ignoredHandlers.insert("unhandledrejection"_kj); ignoredHandlers.insert("rejectionhandled"_kj); // Helper function to collect methods from a prototype chain auto collectMethodsFromPrototypeChain = [&](jsg::JsValue startProto, kj::HashSet& seenNames) { // Find the prototype for `Object` by creating one. auto obj = js.obj(); jsg::JsValue prototypeOfObject = obj.getPrototype(js); // Walk the prototype chain. jsg::JsValue proto = startProto; for (;;) { auto protoObj = KJ_UNWRAP_OR(proto.tryCast(), { errorReporter.addError( kj::str("Exported value's prototype chain does not end in Object.")); return; }); if (protoObj == prototypeOfObject) { // Reached the prototype for `Object`. Stop here. break; } // Awkwardly, the prototype's members are not typically enumerable, so we have to // enumerate them rather directly. jsg::JsArray properties = protoObj.getPropertyNames(js, jsg::KeyCollectionFilter::OWN_ONLY, jsg::PropertyFilter::SKIP_SYMBOLS, jsg::IndexFilter::SKIP_INDICES); for (auto i: kj::zeroTo(properties.size())) { auto name = properties.get(js, i).toString(js); if (name == "constructor"_kj) { // Don't treat special method `constructor` as an exported handler. continue; } if (!ignoredHandlers.contains(name)) { // Only report each method name once, even if it overrides a method in a superclass. seenNames.upsert(kj::mv(name), [&](auto&, auto&&) {}); } } proto = protoObj.getPrototype(js); } }; KJ_IF_SOME(c, worker.impl->context) { // Service workers syntax. auto handlerNames = c->getHandlerNames(); kj::Vector handlers; for (auto& name: handlerNames) { if (!ignoredHandlers.contains(name)) { handlers.add(kj::str(name)); } } if (handlers.empty()) { errorReporter.addError( kj::str("No event handlers were registered. This script does nothing.")); } errorReporter.addEntrypoint(kj::none, handlers.releaseAsArray()); } else { auto report = [&](kj::Maybe name, api::ExportedHandler& exported) { auto handle = exported.self.getHandle(js); if (handle->IsArray()) { // HACK: toDict() will throw a TypeError if given an array, because jsg::DictWrapper is // designed to treat arrays as not matching when a dict is expected. However, // StructWrapper has no such restriction, and therefore an exported array will // successfully produce an ExportedHandler (presumably with no handler functions), and // hence we will see it here. Rather than try to correct this inconsistency between // struct and dict handling (which could have unintended consequences), let's just // work around by ignoring arrays here. errorReporter.addEntrypoint(name, kj::Array()); } else { // Use a HashSet to avoid duplicates when methods exist both as own properties // and in the prototype chain kj::HashSet methodSet; // First, check for own properties (like a plain object literal) auto dict = js.toDict(handle); for (auto& field: dict.fields) { if (!ignoredHandlers.contains(field.name)) { methodSet.upsert(kj::mv(field.name), [&](auto&, auto&&) {}); } } // Then, check for methods in the prototype chain (like a class instance) js.withinHandleScope([&]() { collectMethodsFromPrototypeChain(jsg::JsObject(handle).getPrototype(js), methodSet); }); // Convert HashSet to Array for reporting errorReporter.addEntrypoint(name, KJ_MAP(n, methodSet) { return kj::mv(n); }); } }; auto getEntrypointName = [&](kj::StringPtr key) -> kj::Maybe { if (key == "default"_kj) { return kj::none; } else { return key; } }; for (auto& entry: worker.impl->namedHandlers) { report(getEntrypointName(entry.key), entry.value); } for (auto& entry: worker.impl->actorClasses) { KJ_IF_SOME(entrypointName, getEntrypointName(entry.key)) { errorReporter.addActorClass(entrypointName); } else { // Hmm, it appears someone tried to export a Durable Object class as a default // entrypoint. This doesn't actually work: the runtime will not allow this DO class // to be used, either for actors or as an entrypoint. // // TODO(someday): Make this a hard error. I'm hesitant to do it in my current change // for fear that it'll break someone somewhere forcing a rollback. For now we log. LOG_PERIODICALLY(ERROR, "Exported actor class as default entrypoint. This doesn't work, but historically " "did not produce a startup-time error."); } } for (auto& entry: worker.impl->statelessClasses) { // We want to report all of the stateless class's members. To do this, we examine its // prototype, and its prototype's prototype, and so on, until we get to Object's // prototype, which we ignore. auto entrypointName = getEntrypointName(entry.key); kj::HashSet seenNames; js.withinHandleScope([&]() { // For stateless classes, we need to get the class's prototype property jsg::JsObject ctor(KJ_ASSERT_NONNULL(entry.value.tryGetHandle(js.v8Isolate))); jsg::JsValue proto = ctor.get(js, "prototype"); collectMethodsFromPrototypeChain(proto, seenNames); }); errorReporter.addEntrypoint(entrypointName, KJ_MAP(n, seenNames) { return kj::mv(n); }); } for (auto& entry: worker.impl->workflowClasses) { KJ_IF_SOME(entrypointName, getEntrypointName(entry.key)) { kj::HashSet seenNames; js.withinHandleScope([&]() { // For stateless classes, we need to get the class's prototype property jsg::JsObject ctor(KJ_ASSERT_NONNULL(entry.value.tryGetHandle(js.v8Isolate))); jsg::JsValue proto = ctor.get(js, "prototype"); collectMethodsFromPrototypeChain(proto, seenNames); }); errorReporter.addWorkflowClass(entrypointName, KJ_MAP(n, seenNames) { return kj::mv(n); }); } else { } } } }); } // ======================================================================================= // AsyncLock implementation const kj::EventLoopLocal Worker::AsyncWaiter::threadCurrentWaiter; Worker::Isolate::AsyncWaiterList::~AsyncWaiterList() noexcept { // It should be impossible for this list to be non-empty since each member of the list holds a // strong reference back to us. But if the list is non-empty, we'd better crash here, to avoid // dangling pointers. KJ_ASSERT(head == kj::none, "destroying non-empty waiter list?"); KJ_ASSERT(tail == &head, "tail pointer corrupted?"); } kj::Promise Worker::Isolate::takeAsyncLockWithoutRequest( SpanParent parentSpan) const { auto lockTiming = getMetrics().tryCreateLockTiming(kj::mv(parentSpan)); return takeAsyncLockImpl(kj::mv(lockTiming)); } kj::Promise Worker::Isolate::takeAsyncLock(RequestObserver& request) const { auto lockTiming = getMetrics().tryCreateLockTiming(kj::Maybe(request)); return takeAsyncLockImpl(kj::mv(lockTiming)); } kj::Promise Worker::Isolate::takeAsyncLockImpl( kj::Maybe> lockTiming) const { kj::Maybe currentLoad; if (lockTiming != kj::none) { currentLoad = getCurrentLoad(); } for (uint threadWaitingDifferentLockCount = 0;; ++threadWaitingDifferentLockCount) { AsyncWaiter* waiter = *AsyncWaiter::threadCurrentWaiter; if (waiter == nullptr) { // Thread is not currently waiting on a lock. KJ_IF_SOME(lt, lockTiming) { lt.get()->reportAsyncInfo(KJ_ASSERT_NONNULL(currentLoad), false /* threadWaitingSameLock */, threadWaitingDifferentLockCount); } auto newWaiter = kj::refcounted(kj::atomicAddRef(*this)); co_await newWaiter->readyPromise; co_return AsyncLock(kj::mv(newWaiter), kj::mv(lockTiming)); } else if (waiter->isolate == this) { // Thread is waiting on a lock already, and it's for the same isolate. We can coalesce the // locks. KJ_IF_SOME(lt, lockTiming) { lt.get()->reportAsyncInfo(KJ_ASSERT_NONNULL(currentLoad), true /* threadWaitingSameLock */, threadWaitingDifferentLockCount); } auto newWaiterRef = kj::addRef(*waiter); co_await newWaiterRef->readyPromise; co_return AsyncLock(kj::mv(newWaiterRef), kj::mv(lockTiming)); } else { // Thread is already waiting for or holding a different isolate lock. Wait for that one to // be released before we try to lock a different isolate. // TODO(perf): Use of ForkedPromise leads to thundering herd here. Should be minor in practice, // but we could consider creating another linked list instead... KJ_IF_SOME(lt, lockTiming) { lt.get()->waitingForOtherIsolate(waiter->isolate->getId()); } co_await waiter->releasePromise; } } } kj::Promise Worker::takeAsyncLockWithoutRequest(SpanParent parentSpan) const { return script->getIsolate().takeAsyncLockWithoutRequest(kj::mv(parentSpan)); } kj::Promise Worker::takeAsyncLock(RequestObserver& request) const { return script->getIsolate().takeAsyncLock(request); } Worker::AsyncWaiter::AsyncWaiter(kj::Own isolateParam) : executor(kj::getCurrentThreadExecutor()), isolate(kj::mv(isolateParam)) { // Init `releasePromise` / `releaseFulfiller`. { auto paf = kj::newPromiseAndFulfiller(); releasePromise = paf.promise.fork(); releaseFulfiller = kj::mv(paf.fulfiller); } // Add ourselves to the wait queue for this isolate. auto lock = isolate->asyncWaiters.lockExclusive(); if (lock->tail == &lock->head) { // Looks like the queue is empty, so we immediately get the lock. readyPromise = kj::Promise(kj::READY_NOW).fork(); // We can leave `readyFulfiller` null as no one will ever invoke it anyway. } else { // Arrange to get notified later. auto paf = kj::newPromiseAndCrossThreadFulfiller(); readyPromise = paf.promise.fork(); readyFulfiller = kj::mv(paf.fulfiller); } next = kj::none; prev = lock->tail; *lock->tail = this; lock->tail = &next; *threadCurrentWaiter = this; __atomic_add_fetch(&isolate->impl->lockAttemptGauge, 1, __ATOMIC_RELAXED); } Worker::AsyncWaiter::~AsyncWaiter() noexcept { // This destructor is `noexcept` because an exception here probably leaves the process in a bad // state. __atomic_sub_fetch(&isolate->impl->lockAttemptGauge, 1, __ATOMIC_RELAXED); auto lock = isolate->asyncWaiters.lockExclusive(); releaseFulfiller->fulfill(); // Remove ourselves from the list. *prev = next; KJ_IF_SOME(n, next) { n.prev = prev; } else { lock->tail = prev; } if (prev == &lock->head) { // We held the lock before now. Alert the next waiter that they are now at the front of the // line. KJ_IF_SOME(n, next) { n.readyFulfiller->fulfill(); } } auto& w = *threadCurrentWaiter; KJ_ASSERT(w == this); w = nullptr; } kj::Promise Worker::AsyncLock::whenThreadIdle() { AsyncWaiter*& currentWaiter = *AsyncWaiter::threadCurrentWaiter; for (;;) { if (currentWaiter != nullptr) { co_await currentWaiter->releasePromise; continue; } co_await kj::yieldUntilQueueEmpty(); if (currentWaiter == nullptr) { co_return; } // Whoops, a new lock attempt appeared, loop. } } // ======================================================================================= // A proxy for OutputStream that internally buffers data as long as it's beyond a given limit. // Also, it counts size of all the data it has seen (whether it has hit the limit or not). // // We use this in the Network tab to report response stats and preview [decompressed] bodies, // but we don't want to keep buffering extremely large ones, so just discard buffered data // upon hitting a limit and don't return any body to the devtools frontend afterwards. class Worker::Isolate::LimitedBodyWrapper: public kj::OutputStream { public: LimitedBodyWrapper(size_t limit = 1 * 1024 * 1024): limit(limit) { if (limit > 0) { inner.emplace(); } } KJ_DISALLOW_COPY_AND_MOVE(LimitedBodyWrapper); void reset() { this->inner = kj::none; } void write(kj::ArrayPtr data) override { this->size += data.size(); KJ_IF_SOME(inner, this->inner) { if (this->size <= this->limit) { inner.write(data); } else { reset(); } } } size_t getWrittenSize() { return this->size; } kj::Maybe> getArray() { KJ_IF_SOME(inner, this->inner) { return inner.getArray(); } else { return kj::none; } } private: size_t size = 0; size_t limit = 0; kj::Maybe inner; }; struct MessageQueue { kj::Vector messages; size_t head; enum class Status { ACTIVE, CLOSED } status; }; class Worker::Isolate::InspectorChannelImpl final: public v8_inspector::V8Inspector::Channel { public: InspectorChannelImpl(kj::Own isolateParam, kj::Own isolateThreadExecutor, kj::WebSocket& webSocket) : ioHandler(kj::mv(isolateThreadExecutor), webSocket), state(kj::heap(this, kj::mv(isolateParam))) { ioHandler.connect(*this); } // In preview sessions, synchronous locks are not an issue. We declare an alternate spelling of // the type so that all the individual locks below don't turn up in a search for synchronous // locks. using InspectorLock = Worker::Lock::TakeSynchronously; ~InspectorChannelImpl() noexcept try { // Stop message pump. ioHandler.disconnect(); // Delete session under lock. auto state = this->state.lockExclusive(); jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) { Isolate::Impl::Lock recordedLock(*state->get()->isolate, InspectorLock(kj::none), stackScope); if (state->get()->isolate->currentInspectorSession != kj::none) { const_cast(*state->get()->isolate).disconnectInspector(); } state->get()->teardownUnderLock(); }); } catch (...) { // Unfortunately since we're inheriting from Channel which declares a virtual destructor with // default exception constraints, we have to catch all exceptions here and log them. // But different kinds of exceptions call for different ways to stringify the exception. // kj::runCatchingExceptions() normally does this for us, but there's no way to use it while // wrapping the whole destructor (including destructors of members). So... we do a native // catch(...) and then we rethrow the exception inside a kj::runCatchingExceptions and then log // that. Yeah. // // TODO(cleanup): Maybe we could add a kj::stringifyCurrentException() or // kj::logUncaughtException() or something? KJ_IF_SOME(exception, kj::runCatchingExceptions([&]() { throw; })) { KJ_LOG(ERROR, "uncaught exception in ~Script() and the C++ standard is broken", exception); } } void disconnect() { // Fake like the client requested close. This will cause outgoingLoop() to exit and everything // will be cleaned up. ioHandler.disconnect(); } void dispatchProtocolMessage(kj::String message, v8_inspector::V8InspectorSession& session, Isolate& isolate, jsg::V8StackScope& stackScope, Isolate::Impl::Lock& recordedLock) { capnp::MallocMessageBuilder messageBuilder; auto cmd = messageBuilder.initRoot(); getCdpJsonCodec().decode(message, cmd); switch (cmd.which()) { case cdp::Command::UNKNOWN: { break; } case cdp::Command::NETWORK_ENABLE: { setNetworkEnabled(true); cmd.getNetworkEnable().initResult(); break; } case cdp::Command::NETWORK_DISABLE: { setNetworkEnabled(false); cmd.getNetworkDisable().initResult(); break; } case cdp::Command::NETWORK_GET_RESPONSE_BODY: { auto err = cmd.getNetworkGetResponseBody().initError(); err.setCode(-32600); err.setMessage("Network.getResponseBody is not supported in this fork"); break; } case cdp::Command::PROFILER_STOP: { KJ_IF_SOME(p, isolate.impl->profiler) { auto& lock = recordedLock.lock; stopProfiling(*lock, *p, cmd); } break; } case cdp::Command::PROFILER_START: { KJ_IF_SOME(p, isolate.impl->profiler) { auto& lock = recordedLock.lock; startProfiling(*lock, *p); } break; } case cdp::Command::PROFILER_SET_SAMPLING_INTERVAL: { KJ_IF_SOME(p, isolate.impl->profiler) { auto interval = cmd.getProfilerSetSamplingInterval().getParams().getInterval(); setSamplingInterval(*p, interval); } break; } case cdp::Command::PROFILER_ENABLE: { auto& lock = recordedLock.lock; isolate.impl->profiler = kj::Own( v8::CpuProfiler::New(lock->v8Isolate, v8::kDebugNaming, v8::kLazyLogging), CpuProfilerDisposer::instance); break; } case cdp::Command::TAKE_HEAP_SNAPSHOT: { auto& lock = recordedLock.lock; takeHeapSnapshot(*lock, cmd.getTakeHeapSnapshot().getParams()); break; } } if (!cmd.isUnknown()) { sendNotification(cmd); return; } auto& lock = recordedLock.lock; // We have at times observed V8 bugs where the inspector queues a background task and // then synchronously waits for it to complete, which would deadlock if background // threads are disallowed. Since the inspector is in a process sandbox anyway, it's not // a big deal to just permit those background threads. AllowV8BackgroundThreadsScope allowBackgroundThreads; ExceptionOrDuration limitErrorOrTime = 0 * kj::NANOSECONDS; { auto limitScope = isolate.getLimitEnforcer().enterInspectorJs(*lock, limitErrorOrTime); session.dispatchProtocolMessage(jsg::toInspectorStringView(message)); } // Run microtasks in case the user made an async call. if (!limitErrorOrTime.is()) { auto limitScope = isolate.getLimitEnforcer().enterInspectorJs(*lock, limitErrorOrTime); lock->runMicrotasks(); } else { // Oops, we already exceeded the limit, so force the microtask queue to be thrown away. lock->terminateNextExecution(); lock->runMicrotasks(); } KJ_SWITCH_ONEOF(limitErrorOrTime) { KJ_CASE_ONEOF(limitError, kj::Exception) { lock->withinHandleScope([&] { // HACK: We want to print the error, but we need a context to do that. // We don't know which contexts exist in this isolate, so I guess we have to // create one. Ugh. auto dummyContext = v8::Context::New(lock->v8Isolate); // We need to set the highest used index in every context we create to be a nullptr // This is because we might later on call GetAlignedPointerFromEmbedderData which fails with // a fatal error if the array is smaller than the given index. jsg::setAlignedPointerInEmbedderData( dummyContext, jsg::ContextPointerSlot::MAX_POINTER_SLOT, nullptr); auto& inspector = *KJ_ASSERT_NONNULL(isolate.impl->inspector); inspector.contextCreated(v8_inspector::V8ContextInfo(dummyContext, 1, v8_inspector::StringView(reinterpret_cast("Worker"), 6))); JSG_WITHIN_CONTEXT_SCOPE(*lock, dummyContext, [&](jsg::Lock& js) { jsg::sendExceptionToInspector(js, inspector, jsg::extractTunneledExceptionDescription(limitError.getDescription())); }); inspector.contextDestroyed(dummyContext); }); } KJ_CASE_ONEOF(startupTimeElapsed, kj::Duration) {} } if (recordedLock.checkInWithLimitEnforcer(isolate)) { disconnect(); } } kj::Promise messagePump() { return ioHandler.messagePump(); } void handleDispatchProtocolMessage( Worker::AsyncLock& asyncLock, kj::MutexGuarded& incomingQueue) { auto lockedState = state.lockExclusive(); v8_inspector::V8InspectorSession& session = *lockedState->get()->session; Isolate& isolate = const_cast(*lockedState->get()->isolate); jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) { Isolate::Impl::Lock recordedLock(isolate, asyncLock, stackScope); auto lockedQueue = incomingQueue.lockExclusive(); if (lockedQueue->status != MessageQueue::Status::ACTIVE) { return; } auto messages = lockedQueue->messages.slice(lockedQueue->head, lockedQueue->messages.size()); for (auto& message: messages) { dispatchProtocolMessage(kj::mv(message), session, isolate, stackScope, recordedLock); } lockedQueue->messages.clear(); lockedQueue->head = 0; }); } kj::Promise dispatchProtocolMessages(kj::MutexGuarded& incomingQueue) { // This method is called on the I/O thread, which also adds messages to the `incomingQueue`. // So long as this method does not yield/resume mid-way, there is no concern about how // long the queue lock is held for whilst dispatching messages. auto i = kj::atomicAddRef(*this->state.lockExclusive()->get()->isolate); auto asyncLock = co_await i->takeAsyncLockWithoutRequest(nullptr); handleDispatchProtocolMessage(asyncLock, incomingQueue); } // --------------------------------------------------------------------------- // implements Channel // // Keep in mind that these methods will be called from various threads! void sendResponse(int callId, std::unique_ptr message) override { // callId is encoded in the message, too. Unsure why this method even exists. sendNotification(kj::mv(message)); } bool isNetworkEnabled() { return __atomic_load_n(&networkEnabled, __ATOMIC_RELAXED); } void setNetworkEnabled(bool enable) { __atomic_store_n(&networkEnabled, enable, __ATOMIC_RELAXED); } void sendNotification(kj::String message) { ioHandler.send(kj::mv(message)); } template void sendNotification(T&& message) { sendNotification(getCdpJsonCodec().encode(message)); } void sendNotification(std::unique_ptr message) override { sendNotification(kj::str(message->string())); } void flushProtocolNotifications() override { // Are we supposed to do anything here? There's no documentation, so who knows? Maybe we could // delay signaling the outgoing loop until this call? } // Dispatches one message whilst automatic CDP messages on the I/O worker thread is paused, called // on the thread executing the isolate whilst execution is suspended due to a breakpoint or // debugger statement. bool dispatchOneMessageDuringPause(); private: // Class that manages the I/O for devtools connections. I/O is performed on the // thread associated with the InspectorService (the thread that calls attachInspector). // Most of the public API is intended for code running on the isolate thread, such as // the InspectorChannelImpl and the InspectorClient. class WebSocketIoHandler final { public: WebSocketIoHandler(kj::Own isolateThreadExecutor, kj::WebSocket& webSocket) : isolateThreadExecutor(kj::mv(isolateThreadExecutor)), webSocket(webSocket) { // Assume we are being instantiated on the InspectorService thread, the thread that will do // I/O for CDP messages. Messages are delivered to the InspectorChannelImpl on the Isolate thread. outgoingQueueNotifier = XThreadNotifier::create(); } // Sets the channel that messages are delivered to. void connect(InspectorChannelImpl& inspectorChannel) { channel = inspectorChannel; } void disconnect() { channel = kj::none; shutdown(); } // Blocked the current thread until a message arrives. This is intended // for use in the InspectorClient when breakpoints are hit. The InspectorClient // has to remain in runMessageLoopOnPause() but still receive CDP messages // (e.g. resume). kj::Maybe waitForMessage() { return incomingQueue.when([](const MessageQueue& incomingQueue) { return (incomingQueue.head < incomingQueue.messages.size() || incomingQueue.status == MessageQueue::Status::CLOSED); }, [](MessageQueue& incomingQueue) -> kj::Maybe { if (incomingQueue.status == MessageQueue::Status::CLOSED) return {}; return pollMessage(incomingQueue); }); } // Message pumping promise that should be evaluated on the InspectorService // thread. kj::Promise messagePump() { // Although inspector I/O must happen on the InspectorService thread (to make sure breakpoints // don't block inspector I/O), inspector messages must be actually dispatched on the Isolate // thread. So, we run the dispatch loop on the Isolate thread. // // Note that the above comment is only really accurate in vanilla workerd. In the case of the // internal Cloudflare Workers runtime, `isolateThreadExecutor` may actually refer to the // current thread's `kj::Executor`. That's fine; calling `executeAsync()` on the current // thread's executor just posts the task to the event loop, and everything works as expected. // Since the dispatch loop and the receive loop communicate over a XThreadNotifier, and // XThreadNotifiers must be created on the thread which will call their `awaitNotification()` // function, we awkwardly perform two `executeAsync()`s here, one to create the // XThreadNotifier, then another to spawn the dispatch loop. // // We create a new XThreadNotifier for each `messagePump()` call, rather than try to re-use // one long-term, because XThreadNotifiers' `awaitNotification()` function is not cancel-safe. // That is, once its promise is cancelled, the notifier is broken. auto incomingQueueNotifier = co_await isolateThreadExecutor->executeAsync([]() { return XThreadNotifier::create(); }); auto dispatchLoopPromise = isolateThreadExecutor->executeAsync( [this, notifier = kj::atomicAddRef(*incomingQueueNotifier)]() mutable { return dispatchLoop(kj::mv(notifier)); }); co_return co_await receiveLoop(kj::mv(incomingQueueNotifier)) .exclusiveJoin(kj::mv(dispatchLoopPromise)) .exclusiveJoin(transmitLoop()); } void send(kj::String message) { auto lockedOutgoingQueue = outgoingQueue.lockExclusive(); if (lockedOutgoingQueue->status == MessageQueue::Status::CLOSED) return; lockedOutgoingQueue->messages.add(kj::mv(message)); outgoingQueueNotifier->notify(); } private: static kj::Maybe pollMessage(MessageQueue& messageQueue) { if (messageQueue.head < messageQueue.messages.size()) { kj::String message = kj::mv(messageQueue.messages[messageQueue.head++]); if (messageQueue.head == messageQueue.messages.size()) { messageQueue.head = 0; messageQueue.messages.clear(); } return kj::mv(message); } return {}; } void shutdown() { // Drain incoming queue, the isolate thread may be waiting on it // on will notice it is closed if woken without any messages to // deliver in WebSocketIoWorker::waitForMessage(). { auto lockedIncomingQueue = incomingQueue.lockExclusive(); lockedIncomingQueue->head = 0; lockedIncomingQueue->messages.clear(); lockedIncomingQueue->status = MessageQueue::Status::CLOSED; } { auto lockedOutgoingQueue = outgoingQueue.lockExclusive(); lockedOutgoingQueue->status = MessageQueue::Status::CLOSED; } // Wake any waiters since queue status fields have been updated. outgoingQueueNotifier->notify(); } // Must be called on the InspectorService thread. kj::Promise receiveLoop(kj::Own incomingQueueNotifier) { for (;;) { auto message = co_await webSocket.receive(MAX_MESSAGE_SIZE); KJ_SWITCH_ONEOF(message) { KJ_CASE_ONEOF(text, kj::String) { incomingQueue.lockExclusive()->messages.add(kj::mv(text)); incomingQueueNotifier->notify(); } KJ_CASE_ONEOF(blob, kj::Array) { // Ignore. } KJ_CASE_ONEOF(close, kj::WebSocket::Close) { shutdown(); // Pause here to give transmitLoop() the chance to finish and send a reply close. // When `transmitLoop()` ends, `messagePump()` as a whole will end, canceling // `receiveLoop()`. co_await kj::Promise(kj::NEVER_DONE); } } } } // Must be called on the Isolate thread. kj::Promise dispatchLoop(kj::Own incomingQueueNotifier) { for (;;) { co_await incomingQueueNotifier->awaitNotification(); KJ_IF_SOME(c, channel) { co_await c.dispatchProtocolMessages(this->incomingQueue); } } } // Must be called on the InspectorService thread. kj::Promise transmitLoop() { for (;;) { co_await outgoingQueueNotifier->awaitNotification(); try { auto lockedOutgoingQueue = outgoingQueue.lockExclusive(); auto messages = kj::mv(lockedOutgoingQueue->messages); bool receivedClose = lockedOutgoingQueue->status == MessageQueue::Status::CLOSED; lockedOutgoingQueue.release(); co_await sendToWebSocket(kj::mv(messages)); if (receivedClose) { co_await webSocket.close(1000, "client closed connection"); co_return; } } catch (kj::Exception& e) { shutdown(); throw; } } } kj::Promise sendToWebSocket(kj::Vector messages) { for (auto& message: messages) { co_await webSocket.send(message); } } // We need access to the Isolate thread's kj::Executor to run the inspector dispatch loop. This // doesn't actually have to be an Own, because the Isolate thread will destroy the Isolate // before it exits, but it doesn't hurt. kj::Own isolateThreadExecutor; kj::MutexGuarded incomingQueue; // The notifier for `incomingQueue`, `incomingQueueNotifier`, is created once per // `messagePump()` call, and never re-used, so it doesn't live here. kj::MutexGuarded outgoingQueue; // This XThreadNotifier must be created on the InspectorService thread. kj::Own outgoingQueueNotifier; kj::WebSocket& webSocket; // only accessed on the InspectorService thread. std::atomic_bool receivedClose; // accessed on any thread (only transitions false -> true). kj::Maybe channel; // only accessed on the isolate thread. // Sometimes the inspector protocol sends large messages. KJ defaults to a 1MB size limit // for WebSocket messages, which makes sense for production use cases, but for debug we should // be OK to go larger. So, we'll accept 128MB. static constexpr size_t MAX_MESSAGE_SIZE = 128u << 20; }; WebSocketIoHandler ioHandler; void takeHeapSnapshot( jsg::Lock& js, cdp::HeapProfiler::Command::TakeHeapSnapshot::Params::Reader params) { struct Activity: public v8::ActivityControl { InspectorChannelImpl& channel; Activity(InspectorChannelImpl& channel): channel(channel) {} ControlOption ReportProgressValue(uint32_t done, uint32_t total) override { capnp::MallocMessageBuilder message; auto event = message.initRoot(); auto progressParams = event.initReportHeapSnapshotProgress(); progressParams.setDone(done); progressParams.setTotal(total); if (done == total) { progressParams.setFinished(true); } auto notification = getCdpJsonCodec().encode(event); channel.sendNotification(kj::mv(notification)); return ControlOption::kContinue; } }; struct Writer: public v8::OutputStream { InspectorChannelImpl& channel; Writer(InspectorChannelImpl& channel): channel(channel) {} void EndOfStream() override {} int GetChunkSize() override { return 65536; // big chunks == faster // The chunk size here will determine the actual number of individual // messages that are sent. The default is... rather small. Experience // node and node-heapdump shows that this can be bumped up // much higher to get better performance. Here we use the value // that Node.js uses (see Node.js' FileOutputStream impl). } v8::OutputStream::WriteResult WriteAsciiChunk(char* data, int size) override { capnp::MallocMessageBuilder message; auto event = message.initRoot(); auto params = event.initAddHeapSnapshotChunk(); params.setChunk(kj::heapString(data, size)); auto notification = getCdpJsonCodec().encode(event); channel.sendNotification(kj::mv(notification)); return v8::OutputStream::WriteResult::kContinue; } }; Activity activity(*this); Writer writer(*this); v8::HeapProfiler::HeapSnapshotOptions options{}; if (params.getReportProgress()) { options.control = &activity; } if (params.getExposeInternals()) { options.snapshot_mode = v8::HeapProfiler::HeapSnapshotMode::kExposeInternals; } if (params.getCaptureNumericValue()) { options.numerics_mode = v8::HeapProfiler::NumericsMode::kExposeNumericValues; } auto profiler = js.v8Isolate->GetHeapProfiler(); auto snapshot = kj::Own( profiler->TakeHeapSnapshot(options), HeapSnapshotDeleter::INSTANCE); snapshot->Serialize(&writer); } struct State { kj::Own isolate; std::unique_ptr session; State(InspectorChannelImpl* self, kj::Own isolateParam) : isolate(kj::mv(isolateParam)), session(KJ_ASSERT_NONNULL(isolate->impl->inspector) ->connect(1, self, v8_inspector::StringView(), isolate->impl->inspectorPolicy == InspectorPolicy::ALLOW_UNTRUSTED ? v8_inspector::V8Inspector::kUntrusted : v8_inspector::V8Inspector::kFullyTrusted)) {} ~State() noexcept(false) { if (session != nullptr) { KJ_LOG(ERROR, "Deleting InspectorChannelImpl::State without having called " "teardownUnderLock()", kj::getStackTrace()); // Isolate locks are recursive so it should be safe to lock here. jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) { Isolate::Impl::Lock recordedLock(*isolate, InspectorLock(kj::none), stackScope); session = nullptr; }); } } // Must be called with the worker isolate locked. Should be called immediately before // destruction. void teardownUnderLock() { session = nullptr; } KJ_DISALLOW_COPY_AND_MOVE(State); }; // Mutex ordering: You must lock this *before* locking the isolate. kj::MutexGuarded> state; // Not under `state` lock due to lock ordering complications. volatile bool networkEnabled = false; }; bool Worker::Isolate::InspectorChannelImpl::dispatchOneMessageDuringPause() { auto maybeMessage = ioHandler.waitForMessage(); // We can be paused by either hitting a debugger statement in a script or from hitting // a breakpoint or someone hit break. KJ_IF_SOME(message, maybeMessage) { auto lockedState = this->state.lockExclusive(); // Received a message whilst script is running, probably in a breakpoint. v8_inspector::V8InspectorSession& session = *lockedState->get()->session; // const_cast OK because the IoContext has the lock. Isolate& isolate = const_cast(*lockedState->get()->isolate); Worker::Lock& workerLock = IoContext::current().getCurrentLock(); Isolate::Impl::Lock& recordedLock = workerLock.impl->recordedLock; jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) { dispatchProtocolMessage(kj::mv(message), session, isolate, stackScope, recordedLock); }); return true; } else { // No message from waitForMessage() implies the connection is broken. return false; } } bool Worker::InspectorClient::dispatchOneMessageDuringPause( Worker::Isolate::InspectorChannelImpl& channel) { return channel.dispatchOneMessageDuringPause(); } kj::Promise Worker::Isolate::attachInspector(kj::Timer& timer, kj::Duration timerOffset, kj::HttpService::Response& response, const kj::HttpHeaderTable& headerTable, kj::HttpHeaderId controlHeaderId) const { KJ_REQUIRE(impl->inspector != kj::none); kj::HttpHeaders headers(headerTable); headers.setPtr(controlHeaderId, "{\"ewLog\":{\"status\":\"ok\"}}"); auto webSocket = response.acceptWebSocket(headers); // This `attachInspector()` overload is used by the internal Cloudflare Workers runtime, which has // no concept of a single Isolate thread. Instead, it's OK for all inspector messages to be // dispatched on the calling thread. auto executor = kj::getCurrentThreadExecutor().addRef(); return attachInspector(kj::mv(executor), timer, timerOffset, *webSocket) .attach(kj::mv(webSocket)); } kj::Promise Worker::Isolate::attachInspector( kj::Own isolateThreadExecutor, kj::Timer& timer, kj::Duration timerOffset, kj::WebSocket& webSocket) const { KJ_REQUIRE(impl->inspector != kj::none); return jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) { Isolate::Impl::Lock recordedLock( *this, InspectorChannelImpl::InspectorLock(kj::none), stackScope); auto& lock = *recordedLock.lock; auto& lockedSelf = const_cast(*this); // If another inspector was already connected, boot it, on the assumption that that connection // is dead and this is why the user reconnected. While we could actually allow both inspector // sessions to stay open (V8 supports this!), we'd then need to store a set of all connected // inspectors in order to be able to disconnect all of them in case of an isolate purge... let's // just not. lockedSelf.disconnectInspector(); lockedSelf.impl->inspectorClient->setInspectorTimerInfo(timer, timerOffset); auto channel = kj::heap( kj::atomicAddRef(*this), kj::mv(isolateThreadExecutor), webSocket); lockedSelf.currentInspectorSession = *channel; lockedSelf.impl->inspectorClient->setChannel(*channel); // Send any queued notifications. lock.withinHandleScope([&] { for (auto& notification: lockedSelf.impl->queuedNotifications) { channel->sendNotification(kj::mv(notification)); } lockedSelf.impl->queuedNotifications.clear(); }); return channel->messagePump().attach(kj::mv(channel)); }); } void Worker::Isolate::disconnectInspector() { // If an inspector session is connected, proactively drop it, so as to force it to drop its // reference on the script, so that the script can be deleted. KJ_IF_SOME(current, currentInspectorSession) { current.disconnect(); currentInspectorSession = kj::none; } impl->inspectorClient->resetChannel(); } void Worker::Isolate::logWarning(kj::StringPtr description, Lock& lock) { if (impl->inspector != kj::none) { JSG_WITHIN_CONTEXT_SCOPE(lock, lock.getContext(), [&](jsg::Lock& js) { logMessage(js, static_cast(cdp::LogType::WARNING), description); }); } if (loggingOptions.consoleMode == Worker::ConsoleMode::INSPECTOR_ONLY) { // Run with --verbose to log JS exceptions to stderr. Useful when running tests. KJ_LOG(INFO, "console warning", description); } else { fprintf(stderr, "%s\n", description.cStr()); fflush(stderr); } KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { KJ_IF_SOME(tracer, ioContext.getWorkerTracer()) { // json encoding is required over simply wrapping it in quotes to correctly escape the string. capnp::JsonCodec json; auto jsonDescription = kj::str("[", json.encode(capnp::Text::Reader(description)), "]"); auto timestamp = ioContext.now(); tracer.addLog( ioContext.getInvocationSpanContext(), timestamp, LogLevel::WARN, kj::mv(jsonDescription)); } } } void Worker::Isolate::logWarningOnce(kj::StringPtr description, Lock& lock) { impl->warningOnceDescriptions.findOrCreate(description, [&] { logWarning(description, lock); return kj::str(description); }); } void Worker::Isolate::logErrorOnce(kj::StringPtr description) { impl->errorOnceDescriptions.findOrCreate(description, [&] { KJ_LOG(ERROR, description); return kj::str(description); }); } void Worker::Isolate::logMessage(jsg::Lock& js, uint16_t type, kj::StringPtr description) { if (impl->inspector != kj::none) { // We want to log a warning to the devtools console, as if `console.warn()` were called. // However, the only public interface to call the real `console.warn()` is via JavaScript, // where it could have been monkey-patched by the guest. We'd like to avoid having to worry // about that blowing up in our face. So instead we arrange to send the proper devtools // protocol messages ourselves. // // TODO(cleanup): It would be better if we could directly add the message to the inspector's // console log (without calling through JavaScript). What we're doing here has some problems. // In particular, if no client is connected yet, we attempt to queue up the messages to send // later, much like the real inspector does. This is kind of complicated, and doesn't quite // work right: // - The messages won't necessarily be in the right order with normal console logs made at // the same time (with identical timestamps). // - In theory we should queue *all* logged warnings and deliver them to every future client, // not just the next client to connect. But if we do that, we also need to respect the // protocol command to clear the history when requested. This was further than I cared to // go. // To fix these problems, maybe we should just patch V8 with a direct interface into the // inspector's own log. (Also, how does Chrome handle this?) js.withinHandleScope([&] { capnp::MallocMessageBuilder message; auto event = message.initRoot(); auto params = event.initRuntimeConsoleApiCalled(); params.setType(static_cast(type)); params.initArgs(1)[0].initString().setValue(description); params.setExecutionContextId(v8_inspector::V8ContextInfo::executionContextId(js.v8Context())); params.setTimestamp(impl->inspectorClient->currentTimeMS()); stackTraceToCDP(js, params.initStackTrace()); auto notification = getCdpJsonCodec().encode(event); KJ_IF_SOME(i, currentInspectorSession) { i.sendNotification(kj::mv(notification)); } else { impl->queuedNotifications.add(kj::mv(notification)); } }); } } // ======================================================================================= struct Worker::Actor::Impl { Actor::Id actorId; Frankenvalue props; MakeStorageFunc makeStorage; kj::Own metrics; // When a boolean, indicates whether a `transient` should exist. If true, it will be initialized // on the first `ensureConstructed()`. kj::OneOf> transient; kj::Maybe> actorCache; kj::Maybe> ctxObject; kj::Maybe container; kj::Maybe facetManager; kj::Maybe version; struct NoClass {}; struct Initializing {}; // If the actor is backed by a class, this field tracks the instance through its stages. The // instance is constructed as part of the first request to be delivered. kj::OneOf classInstance; class HooksImpl: public InputGate::Hooks, public OutputGate::Hooks, public ActorCache::Hooks { public: HooksImpl(kj::Own loopback, TimerChannel& timerChannel, ActorObserver& metrics) : loopback(kj::mv(loopback)), timerChannel(timerChannel), metrics(metrics) {} void inputGateLocked() override { metrics.inputGateLocked(); } void inputGateReleased() override { metrics.inputGateReleased(); } void inputGateWaiterAdded() override { metrics.inputGateWaiterAdded(); } void inputGateWaiterRemoved() override { metrics.inputGateWaiterRemoved(); } // Implements InputGate::Hooks. kj::Promise makeTimeoutPromise() override { // This really only protects against total hangs. Lowering the timeout drastically is risky, // since low timeouts can spuriously fire when under heavy CPU load, failing requests that // would otherwise succeed. auto timeout = 30 * kj::SECONDS; co_await timerChannel.afterLimitTimeout(timeout); kj::throwFatalException(KJ_EXCEPTION(OVERLOADED, "broken.outputGateBroken; jsg.Error: Durable Object storage operation exceeded " "timeout which caused object to be reset.")); } // Implements OutputGate::Hooks. void outputGateLocked() override { metrics.outputGateLocked(); } void outputGateReleased() override { metrics.outputGateReleased(); } void outputGateWaiterAdded() override { metrics.outputGateWaiterAdded(); } void outputGateWaiterRemoved() override { metrics.outputGateWaiterRemoved(); } // Implements ActorCache::Hooks void updateAlarmInMemory(kj::Maybe newAlarmTime) override; void storageReadCompleted(kj::Duration latency) override { metrics.storageReadCompleted(latency); } void storageWriteCompleted(kj::Duration latency) override { metrics.storageWriteCompleted(latency); } private: kj::Own loopback; // only for updateAlarmInMemory() TimerChannel& timerChannel; // only for afterLimitTimeout() and updateAlarmInMemory() ActorObserver& metrics; kj::Maybe> maybeAlarmPreviewTask; }; HooksImpl hooks; // Handles both input locks and request locks. InputGate inputGate; // Handles output locks. OutputGate outputGate; // `ioContext` is initialized upon delivery of the first request. kj::Maybe> ioContext; // If onBroken() is called while `ioContext` is still null, this is initialized. When // `ioContext` is constructed, this will be fulfilled with `ioContext.onAbort()`. kj::Maybe>>> abortFulfiller; // Task which periodically flushes metrics. Initialized after `ioContext` is initialized. kj::Maybe> metricsFlushLoopTask; // Allows sending requests back into this actor, recreating it as necessary. Safe to hold longer // than the Worker::Actor is alive. kj::Own loopback; TimerChannel& timerChannel; kj::ForkedPromise shutdownPromise; kj::Own> shutdownFulfiller; // If this Actor has a HibernationManager, it means the Actor has recently accepted a Hibernatable // websocket. We eventually move the HibernationManager into the DeferredProxy task // (since it's long lived), but can still refer to the HibernationManager by passing a reference // in each CustomEvent. kj::Maybe> hibernationManager; kj::Maybe hibernationEventType; struct ScheduledAlarm { ScheduledAlarm( kj::Date scheduledTime, kj::PromiseFulfillerPair pf) : scheduledTime(scheduledTime), resultFulfiller(kj::mv(pf.fulfiller)), resultPromise(pf.promise.fork()) {} KJ_DISALLOW_COPY(ScheduledAlarm); ScheduledAlarm(ScheduledAlarm&&) = default; ~ScheduledAlarm() noexcept(false) {} kj::Date scheduledTime; WorkerInterface::AlarmFulfiller resultFulfiller; kj::ForkedPromise resultPromise; kj::Promise cleanupPromise = resultPromise.addBranch().then( [](WorkerInterface::AlarmOutcome&&) {}, [](kj::Exception&&) {}); // The first thing we do after we get a result should be to remove the running alarm (if we got // that far). So we grab the first branch now and ignore any results, before anyone else has a // chance to do so. }; struct RunningAlarm { kj::Date scheduledTime; kj::ForkedPromise resultPromise; }; // If valid, we have an alarm invocation that has not yet received an `AlarmFulfiller` and thus // is either waiting for a running alarm or its scheduled time. kj::Maybe maybeScheduledAlarm; // If valid, we have an alarm invocation that has received an `AlarmFulfiller` and is currently // considered running. This alarm is no longer cancelable. kj::Maybe maybeRunningAlarm; // This is a forked promise so that we can schedule and then cancel multiple alarms while an alarm // is running. kj::ForkedPromise runningAlarmTask = kj::Promise(kj::READY_NOW).fork(); Impl(Worker::Actor& self, Actor::Id actorId, bool hasTransient, MakeActorCacheFunc makeActorCache, Frankenvalue props, MakeStorageFunc makeStorage, kj::Own loopback, TimerChannel& timerChannel, kj::Own metricsParam, kj::Maybe> manager, kj::Maybe& hibernationEventType, kj::Maybe container, kj::Maybe facetManager, kj::PromiseFulfillerPair paf = kj::newPromiseAndFulfiller()) : actorId(kj::mv(actorId)), props(kj::mv(props)), makeStorage(kj::mv(makeStorage)), metrics(kj::mv(metricsParam)), transient(hasTransient), container(kj::mv(container)), facetManager(facetManager), hooks(loopback->addRef(), timerChannel, *metrics), inputGate(hooks), outputGate(hooks), loopback(kj::mv(loopback)), timerChannel(timerChannel), shutdownPromise(paf.promise.fork()), shutdownFulfiller(kj::mv(paf.fulfiller)), hibernationManager(kj::mv(manager)), hibernationEventType(kj::mv(hibernationEventType)) { actorCache = makeActorCache(self.worker->getIsolate().impl->actorCacheLru, outputGate, hooks, *metrics); } }; kj::Promise Worker::takeAsyncLockWhenActorCacheReady( kj::Date now, Actor& actor, RequestObserver& request) const { auto lockTiming = getIsolate().getMetrics().tryCreateLockTiming(kj::Maybe(request)); KJ_IF_SOME(c, actor.impl->actorCache) { KJ_IF_SOME(p, c.get()->evictStale(now)) { // Got backpressure, wait for it. // TODO(someday): Count this time period differently in lock timing data? co_await p; } } co_return co_await getIsolate().takeAsyncLockImpl(kj::mv(lockTiming)); } Worker::Actor::Actor(const Worker& worker, kj::Maybe tracker, Actor::Id actorId, bool hasTransient, MakeActorCacheFunc makeActorCache, kj::Maybe className, Frankenvalue props, MakeStorageFunc makeStorage, kj::Own loopback, TimerChannel& timerChannel, kj::Own metrics, kj::Maybe> manager, kj::Maybe hibernationEventType, kj::Maybe container, kj::Maybe facetManager, kj::Maybe version) : worker(kj::atomicAddRef(worker)), tracker(tracker.map([](RequestTracker& tracker) { return tracker.addRef(); })) { impl = kj::heap(*this, kj::mv(actorId), hasTransient, kj::mv(makeActorCache), kj::mv(props), kj::mv(makeStorage), kj::mv(loopback), timerChannel, kj::mv(metrics), kj::mv(manager), hibernationEventType, kj::mv(container), facetManager); impl->version = kj::mv(version); KJ_IF_SOME(c, className) { KJ_IF_SOME(cls, worker.impl->actorClasses.find(c)) { // const_cast OK because we're just storing the pointer and will only use this under lock. impl->classInstance = const_cast(&cls); } else { auto e = KJ_EXCEPTION(FAILED, "broken.ignored; no such actor class", c); e.setDetail(jsg::EXCEPTION_IS_USER_ERROR, kj::heapArray(0)); kj::throwFatalException(kj::mv(e)); } } else { impl->classInstance = Impl::NoClass(); } } void Worker::Actor::ensureConstructed(IoContext& context) { KJ_IF_SOME(info, impl->classInstance.tryGet()) { // IMPORTANT: We need to set the state to "Initializing" synchronously, before // ensureConstructedImpl() actually executes and acquires the input lock. // This prevents multiple concurrent initialization attempts if multiple calls to // ensureConstructed() arrive back-to-back. // // This doesn't create a race condition with getHandler() because InputGate::wait() // synchronously adds the caller to the wait queue, even though it completes // asynchronously. Any call to getHandler() that arrives after this point will // have to wait for the input lock, which is only acquired and released by // ensureConstructedImpl() when it completes initialization. // // So the "actor still initializing" error in getHandler() should be impossible // unless a code path is bypassing the input lock mechanism. context.addWaitUntil(ensureConstructedImpl(context, *info)); impl->classInstance = Impl::Initializing(); } } kj::Promise Worker::Actor::ensureConstructedImpl(IoContext& context, ActorClassInfo& info) { InputGate::Lock inputLock = co_await impl->inputGate.wait(context.getCurrentTraceSpan()); try { bool containerRunning = false; KJ_IF_SOME(c, impl->container) { // We need to do an RPC to check if the container is running. // TODO(perf): It would be nice if we could have started this RPC earlier, e.g. in parallel // with starting the script, and also if we could save the status across hibernations. But // that would require some refactoring, and this RPC should (eventally) be local, so it's // not a huge deal. auto status = co_await c.statusRequest(capnp::MessageSize{4, 0}).send(); containerRunning = status.getRunning(); } co_await context.run([this, &info, containerRunning](Worker::Lock& lock) { jsg::Lock& js = lock; kj::Maybe> storage; KJ_IF_SOME(c, impl->actorCache) { storage = impl->makeStorage(lock, worker->getIsolate().getApi(), *c); } auto ctx = js.alloc(js, cloneId(), jsg::JsValue(KJ_ASSERT_NONNULL(lock.getWorker().impl->ctxExports).getHandle(js)), impl->props.toJs(js), kj::mv(storage), kj::mv(impl->container), containerRunning, impl->facetManager, impl->version.map([](ActorVersion& v) { return ActorVersion{.cohort = v.cohort.map([](kj::String& s) { return kj::str(s); })}; })); auto handler = info.cls(lock, ctx.addRef(), KJ_ASSERT_NONNULL(lock.getWorker().impl->env).addRef(js)); // Since we JUST passed `ctx` into the class constructor, it definitely has a handle // attached. Let's grab it and stash it to implement getCtx(). auto ctxHandle = jsg::JsObject(KJ_ASSERT_NONNULL(ctx.tryGetHandle(js))); impl->ctxObject = jsg::JsRef(js, ctxHandle); // HACK: We set handler.env to undefined because we already passed the real env into the // constructor, and we want the handler methods to act like they take just one parameter. // We do the same for handler.ctx, as ExecutionContext related tasks are performed // on the actor's state field instead. handler.env = js.v8Ref(js.v8Undefined()); handler.ctx = kj::none; handler.missingSuperclass = info.missingSuperclass; impl->classInstance = kj::mv(handler); }, inputLock.addRef(context.getCurrentTraceSpan())); // We addRef() the inputLock above rather than kj::mv() it so that the lock remains held // through the catch block below, if an exception is thrown. This is important since we // MUST update `impl->classInstance` to something other than `Initializing` before we // release the lock. } catch (...) { // Get the KJ exception auto e = kj::getCaughtExceptionAsKj(); auto msg = e.getDescription(); if (!msg.startsWith("broken."_kj) && !msg.startsWith("remote.broken."_kj)) { // If we already set up a brokenness reason, we shouldn't override it. auto description = jsg::annotateBroken(msg, "broken.constructorFailed"); e.setDescription(kj::mv(description)); } context.abort(e.clone()); impl->classInstance = kj::mv(e); } } Worker::Actor::~Actor() noexcept(false) { // Note: We do not need an isolate lock to destroy the actor impl. Everything in it is specific // to our thread, or is a handle that can be dropped outside of the lock. } void Worker::Actor::shutdown(uint16_t reasonCode, kj::Maybe error) { // We're officially canceling all background work and we're going to destruct the Actor as soon // as all IoContexts that reference it go out of scope. We might still log additional // periodic messages, and that's good because we might care about that information. That said, // we're officially "broken" from this point because we cannot service background work and our // capability server should have triggered this (potentially indirectly) via its destructor. KJ_IF_SOME(r, impl->ioContext) { impl->metrics->shutdown(reasonCode, r.get()->getLimitEnforcer()); } else { // The actor was shut down before the IoContext was even constructed, so no metrics are // written. } shutdownActorCache(error); impl->shutdownFulfiller->fulfill(); } void Worker::Actor::shutdownActorCache(kj::Maybe error) { KJ_IF_SOME(ac, impl->actorCache) { ac.get()->shutdown(error); } else { // The actor was aborted before the actor cache was constructed, nothing to do. } } kj::Promise Worker::Actor::onShutdown() { return impl->shutdownPromise.addBranch(); } kj::Promise Worker::Actor::onBroken() { // TODO(soon): Detect and report other cases of brokenness, as described in worker.capnp. kj::Promise abortPromise = nullptr; KJ_IF_SOME(rc, impl->ioContext) { abortPromise = rc.get()->onAbort(); } else { auto paf = kj::newPromiseAndFulfiller>(); abortPromise = kj::mv(paf.promise); impl->abortFulfiller = kj::mv(paf.fulfiller); } return abortPromise; } const Worker::Actor::Id& Worker::Actor::getId() { return impl->actorId; } bool Worker::Actor::idsEqual(const Id& a, const Id& b) { if (a.which() != b.which()) return false; KJ_SWITCH_ONEOF(a) { KJ_CASE_ONEOF(actorId, kj::Own) { return actorId->equals(*b.get>()); } KJ_CASE_ONEOF(str, kj::String) { return str == b.get(); } } KJ_UNREACHABLE; } Worker::Actor::Id Worker::Actor::cloneId(Worker::Actor::Id& id) { KJ_SWITCH_ONEOF(id) { KJ_CASE_ONEOF(coloLocalId, kj::String) { return kj::str(coloLocalId); } KJ_CASE_ONEOF(globalId, kj::Own) { return globalId->clone(); } } KJ_UNREACHABLE; } Worker::Actor::Id Worker::Actor::cloneId() { return cloneId(impl->actorId); } kj::Maybe> Worker::Actor::getTransient(Worker::Lock& lock) { KJ_REQUIRE(&lock.getWorker() == worker.get()); if (impl->transient.tryGet().orDefault(false)) { // First call and `hasTransient` was true. Initialize it now, since we have the lock. jsg::Lock& js = lock; impl->transient.init>(js, js.obj()); } return impl->transient.tryGet>().map( [&](jsg::JsRef& val) { return val.addRef(lock); }); } kj::Maybe Worker::Actor::getPersistent() { return impl->actorCache; } kj::Own Worker::Actor::getLoopback() { return impl->loopback->addRef(); } kj::Maybe> Worker::Actor::makeStorageForSwSyntax( Worker::Lock& lock) { return impl->actorCache.map([&](kj::Own& cache) { return impl->makeStorage(lock, worker->getIsolate().getApi(), *cache); }); } void Worker::Actor::assertCanSetAlarm() { KJ_SWITCH_ONEOF(impl->classInstance) { KJ_CASE_ONEOF(_, Impl::NoClass) { // Once upon a time, we allowed actors without classes. Let's make a nicer message if we // we somehow see a classless actor attempt to run an alarm in the wild. JSG_FAIL_REQUIRE( TypeError, "Your Durable Object must be class-based in order to call setAlarm()"); } KJ_CASE_ONEOF(_, Worker::ActorClassInfo*) { KJ_FAIL_ASSERT("setAlarm() invoked before Durable Object ctor"); } KJ_CASE_ONEOF(_, Impl::Initializing) { // We don't explicitly know if we have an alarm handler or not, so just let it happen. We'll // handle it when we go to run the alarm. return; } KJ_CASE_ONEOF(handler, api::ExportedHandler) { JSG_REQUIRE(handler.alarm != kj::none, TypeError, "Your Durable Object class must have an alarm() handler in order to call setAlarm()"); return; } KJ_CASE_ONEOF(exception, kj::Exception) { // We've failed in the ctor, might as well just throw that exception for now. kj::throwFatalException(exception.clone()); } } KJ_UNREACHABLE; } void Worker::Actor::Impl::HooksImpl::updateAlarmInMemory(kj::Maybe newTime) { if (newTime == kj::none) { maybeAlarmPreviewTask = kj::none; return; } auto scheduledTime = KJ_ASSERT_NONNULL(newTime); auto retry = kj::coCapture([this, originalTime = scheduledTime]() -> kj::Promise { kj::Date scheduledTime = originalTime; for (auto i: kj::zeroTo(WorkerInterface::ALARM_RETRY_MAX_TRIES)) { co_await timerChannel.atTime(scheduledTime); auto result = co_await loopback->getWorker(IoChannelFactory::SubrequestMetadata{}) ->runAlarm(originalTime, i); if (result.outcome == EventOutcome::OK || !result.retry) { break; } auto delay = (WorkerInterface::ALARM_RETRY_START_SECONDS << i++) * kj::SECONDS; scheduledTime = timerChannel.now() + delay; } }); maybeAlarmPreviewTask = retry(); } kj::Maybe> Worker::Actor::getAlarm( kj::Date scheduledTime) { KJ_IF_SOME(runningAlarm, impl->maybeRunningAlarm) { if (runningAlarm.scheduledTime == scheduledTime) { // The running alarm has the same time, we can just wait for it. return runningAlarm.resultPromise.addBranch(); } } KJ_IF_SOME(scheduledAlarm, impl->maybeScheduledAlarm) { if (scheduledAlarm.scheduledTime == scheduledTime) { // The scheduled alarm has the same time, we can just wait for it. return scheduledAlarm.resultPromise.addBranch(); } } return kj::none; } kj::Promise Worker::Actor::scheduleAlarm( kj::Date scheduledTime) { KJ_IF_SOME(runningAlarm, impl->maybeRunningAlarm) { if (runningAlarm.scheduledTime == scheduledTime) { // The running alarm has the same time, we can just wait for it. auto result = co_await runningAlarm.resultPromise; co_return result; } } KJ_IF_SOME(scheduledAlarm, impl->maybeScheduledAlarm) { // We had a previously scheduled alarm, let's cancel it. scheduledAlarm.resultFulfiller.cancel(); impl->maybeScheduledAlarm = kj::none; } KJ_IASSERT(impl->maybeScheduledAlarm == kj::none); auto& scheduledAlarm = impl->maybeScheduledAlarm.emplace( scheduledTime, kj::newPromiseAndFulfiller()); // Probably don't need to use kj::coCapture for this but doing so just to be on the // safe side... auto whenCanceled = (kj::coCapture([&scheduledAlarm]() -> kj::Promise { // We've been cancelled, so return that result. Note that we cannot be resolved any other // way until we return an AlarmFulfiller below. co_return co_await scheduledAlarm.resultPromise; }))(); // Date.now() < scheduledTime when the alarm comes in, since we subtract elapsed CPU time from // the time of last I/O in the implementation of Date.now(). This difference could be used to // implement a Spectre timer, so we have to wait a little longer until // `Date.now() == scheduledTime`. Note that this also means that we could invoke ahead of its // `scheduledTime` and we'll delay until appropriate, this may be useful in cases of clock skew. co_return co_await handleAlarm(scheduledTime).exclusiveJoin(kj::mv(whenCanceled)); } kj::Promise Worker::Actor::handleAlarm( kj::Date scheduledTime) { // Let's wait for any running alarm to cleanup before we even delay. co_await impl->runningAlarmTask; co_await KJ_ASSERT_NONNULL(impl->ioContext)->atTime(scheduledTime); // It's time to run! Let's tear apart the scheduled alarm and make a running alarm. // `maybeScheduledAlarm` should have the same value we emplaced above. If another call to // `scheduleAlarm()` emplaced a new value, then `whenCanceled` should have resolved which // cancels this this promise chain. auto scheduledAlarm = KJ_ASSERT_NONNULL(kj::mv(impl->maybeScheduledAlarm)); impl->maybeScheduledAlarm = kj::none; impl->maybeRunningAlarm.emplace(Impl::RunningAlarm{ .scheduledTime = scheduledAlarm.scheduledTime, .resultPromise = kj::mv(scheduledAlarm.resultPromise), }); impl->runningAlarmTask = scheduledAlarm.cleanupPromise .attach(kj::defer([&impl = *impl]() { // As soon as we get fulfilled or rejected, let's unset this alarm as the running alarm. // // NOTE: We could get here during `Actor`'s destructor, which in turn calls `Actor::Impl`'s // destructor, which destroys `runningAlarmTask`, which is us. But in this case, `actor.impl` // is already nulled out (the pointer gets nulled before the destructor runs). This is why we // captured `impl` by reference above, rather than capturing `this`. impl.maybeRunningAlarm = kj::none; })).eagerlyEvaluate([](kj::Exception&& e) { LOG_EXCEPTION("actorAlarmCleanup", e); }).fork(); co_return kj::mv(scheduledAlarm.resultFulfiller); } kj::Maybe Worker::Actor::getHandler() { KJ_SWITCH_ONEOF(impl->classInstance) { KJ_CASE_ONEOF(_, Impl::NoClass) { return kj::none; } KJ_CASE_ONEOF(_, Worker::ActorClassInfo*) { KJ_FAIL_ASSERT("ensureConstructed() wasn't called"); } KJ_CASE_ONEOF(_, Impl::Initializing) { // This shouldn't be possible because ensureConstructed() would have initiated the // construction task which would have taken an input lock as well as the isolate lock, // which should have prevented any other code from executing on the actor until they // were released. KJ_FAIL_ASSERT("actor still initializing when getHandler() called"); } KJ_CASE_ONEOF(handler, api::ExportedHandler) { return handler; } KJ_CASE_ONEOF(exception, kj::Exception) { kj::throwFatalException(exception.clone()); } } KJ_UNREACHABLE; } ActorObserver& Worker::Actor::getMetrics() { return *impl->metrics; } InputGate& Worker::Actor::getInputGate() { return impl->inputGate; } OutputGate& Worker::Actor::getOutputGate() { return impl->outputGate; } kj::Maybe Worker::Actor::getIoContext() { return impl->ioContext.map([](kj::Own& rc) -> IoContext& { return *rc; }); } void Worker::Actor::setIoContext(kj::Own context) { KJ_REQUIRE(impl->ioContext == kj::none); KJ_IF_SOME(f, impl->abortFulfiller) { f.get()->fulfill(context->onAbort()); impl->abortFulfiller = kj::none; } auto& limitEnforcer = context->getLimitEnforcer(); impl->ioContext = kj::mv(context); impl->metricsFlushLoopTask = impl->metrics->flushLoop(impl->timerChannel, limitEnforcer) .eagerlyEvaluate([](kj::Exception&& e) { LOG_EXCEPTION("actorMetricsFlushLoop", e); }); } jsg::JsObject Worker::Actor::getCtx(jsg::Lock& js) { return KJ_REQUIRE_NONNULL(impl->ctxObject).getHandle(js); } jsg::JsValue Worker::Actor::getEnv(jsg::Lock& js) { return jsg::JsValue(KJ_REQUIRE_NONNULL(worker->impl->env).getHandle(js)); } kj::Maybe Worker::Actor::getHibernationManager() { return impl->hibernationManager.map( [](kj::Own& hib) -> HibernationManager& { return *hib; }); } void Worker::Actor::setHibernationManager(kj::Own hib) { KJ_REQUIRE(impl->hibernationManager == kj::none); hib->setTimerChannel(impl->timerChannel); // Not the cleanest way to provide hibernation manager with a timer channel reference, but // where HibernationManager is constructed (actor-state), we don't have a timer channel ref. impl->hibernationManager = kj::mv(hib); } kj::Maybe Worker::Actor::getHibernationEventType() { return impl->hibernationEventType; } kj::Own Worker::Actor::addRef() { KJ_IF_SOME(t, tracker) { // We can attachToThisReference() here, attached object's lifetime being tied to refcounted // instance is deliberate. return kj::addRef(*this).attachToThisReference(t.get()->startRequest()); } else { return kj::addRef(*this); } } // ======================================================================================= uint Worker::Isolate::getCurrentLoad() const { return __atomic_load_n(&impl->lockAttemptGauge, __ATOMIC_RELAXED); } uint Worker::Isolate::getLockSuccessCount() const { return __atomic_load_n(&impl->lockSuccessCount, __ATOMIC_RELAXED); } kj::Own Worker::Isolate::newScript(kj::StringPtr scriptId, const Script::Source& source, IsolateObserver::StartType startType, SpanParent parentSpan, kj::Own vfs, bool logNewScript, kj::Maybe errorReporter, kj::Maybe> artifacts, kj::Maybe> maybeNewModuleRegistry) const { // Script doesn't already exist, so compile it. return kj::atomicRefcounted