File
Blob: src/workerd/api/global-scope.c++
| 1 | // Copyright (c) 2017-2022 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | |
| 5 | #include "global-scope.h" |
| 6 | |
| 7 | #include "simdutf.h" |
| 8 | |
| 9 | #include <workerd/api/cache.h> |
| 10 | #include <workerd/api/crypto/crypto.h> |
| 11 | #include <workerd/api/events.h> |
| 12 | #include <workerd/api/eventsource.h> |
| 13 | #ifdef WORKERD_FUZZILLI |
| 14 | #include <workerd/api/fuzzilli.h> |
| 15 | #endif |
| 16 | #include <workerd/api/hibernatable-web-socket.h> |
| 17 | #include <workerd/api/scheduled.h> |
| 18 | #include <workerd/api/sockets.h> |
| 19 | #include <workerd/api/system-streams.h> |
| 20 | #include <workerd/api/trace.h> |
| 21 | #include <workerd/api/tracing.h> |
| 22 | #include <workerd/api/util.h> |
| 23 | #include <workerd/api/worker-rpc.h> |
| 24 | #include <workerd/io/compatibility-date.h> |
| 25 | #include <workerd/io/features.h> |
| 26 | #include <workerd/io/io-context.h> |
| 27 | #include <workerd/io/tracer.h> |
| 28 | #include <workerd/jsg/async-context.h> |
| 29 | #include <workerd/jsg/ser.h> |
| 30 | #include <workerd/jsg/util.h> |
| 31 | #include <workerd/util/sentry.h> |
| 32 | #include <workerd/util/stream-utils.h> |
| 33 | #include <workerd/util/thread-scopes.h> |
| 34 | #include <workerd/util/uncaught-exception-source.h> |
| 35 | #include <workerd/util/use-perfetto-categories.h> |
| 36 | |
| 37 | #include <kj/encoding.h> |
| 38 | |
| 39 | namespace workerd::api { |
| 40 | |
| 41 | namespace { |
| 42 | |
| 43 | enum class NeuterReason { SENT_RESPONSE, THREW_EXCEPTION, CLIENT_DISCONNECTED }; |
| 44 | |
| 45 | kj::Exception makeNeuterException(NeuterReason reason) { |
| 46 | switch (reason) { |
| 47 | case NeuterReason::SENT_RESPONSE: |
| 48 | return JSG_KJ_EXCEPTION( |
| 49 | FAILED, TypeError, "Can't read from request stream after response has been sent."); |
| 50 | case NeuterReason::THREW_EXCEPTION: |
| 51 | return JSG_KJ_EXCEPTION( |
| 52 | FAILED, TypeError, "Can't read from request stream after responding with an exception."); |
| 53 | case NeuterReason::CLIENT_DISCONNECTED: |
| 54 | return JSG_KJ_EXCEPTION( |
| 55 | DISCONNECTED, TypeError, "Can't read from request stream because client disconnected."); |
| 56 | } |
| 57 | KJ_UNREACHABLE; |
| 58 | } |
| 59 | |
| 60 | } // namespace |
| 61 | |
| 62 | void ExecutionContext::waitUntil(kj::Promise<void> promise) { |
| 63 | IoContext::current().addWaitUntil(kj::mv(promise)); |
| 64 | } |
| 65 | |
| 66 | void ExecutionContext::passThroughOnException() { |
| 67 | IoContext::current().setFailOpen(); |
| 68 | } |
| 69 | |
| 70 | jsg::Optional<jsg::Ref<CacheContext>> ExecutionContext::getCache(jsg::Lock& js) { |
| 71 | // Hook for the embedding application to provide a CacheContext. |
| 72 | // The default Worker::Api implementation returns kj::none. |
| 73 | if (IoContext::hasCurrent()) { |
| 74 | return Worker::Isolate::from(js).getApi().getCtxCacheProperty(js); |
| 75 | } |
| 76 | return kj::none; |
| 77 | } |
| 78 | |
| 79 | jsg::Promise<CachePurgeResult> CacheContext::purge(jsg::Lock& js, |
| 80 | CachePurgeOptions options, |
| 81 | const jsg::TypeHandler<CachePurgeOptions>& optionsHandler, |
| 82 | const jsg::TypeHandler<CachePurgeResult>& resultHandler, |
| 83 | const jsg::TypeHandler<jsg::Ref<JsRpcProperty>>& rpcPropHandler) { |
| 84 | JSG_FAIL_REQUIRE(Error, "Cache purge is not available in this context."); |
| 85 | } |
| 86 | |
| 87 | jsg::Optional<jsg::Ref<Tracing>> ExecutionContext::getTracing(jsg::Lock& js) { |
| 88 | if (!FeatureFlags::get(js).getWorkerdExperimental()) { |
| 89 | return kj::none; |
| 90 | } |
| 91 | // A new Tracing handle is allocated on first access only - `JSG_LAZY_INSTANCE_PROPERTY` |
| 92 | // uses V8's SetLazyDataProperty, which caches the getter result on the instance after the |
| 93 | // first call. So `ctx.tracing === ctx.tracing` and only one allocation per |
| 94 | // ExecutionContext. Tracing itself is stateless (no fields); the per-request state lives |
| 95 | // on IoContext::getCurrentUserTraceSpan(), which enterSpan consults at call time. |
| 96 | return js.alloc<Tracing>(); |
| 97 | } |
| 98 | |
| 99 | void ExecutionContext::abort(jsg::Lock& js, jsg::Optional<jsg::Value> reason) { |
| 100 | KJ_IF_SOME(r, reason) { |
| 101 | IoContext::current().abort(js.exceptionToKj(kj::mv(r))); |
| 102 | } else { |
| 103 | auto e = |
| 104 | JSG_KJ_EXCEPTION(FAILED, Error, "Worker execution was aborted due to call to ctx.abort()."); |
| 105 | IoContext::current().abort(kj::mv(e)); |
| 106 | } |
| 107 | |
| 108 | js.terminateExecutionNow(); |
| 109 | } |
| 110 | |
| 111 | namespace { |
| 112 | template <typename T> |
| 113 | jsg::LenientOptional<T> mapAddRef(jsg::Lock& js, jsg::LenientOptional<T>& function) { |
| 114 | return function.map([&](T& a) { return a.addRef(js); }); |
| 115 | } |
| 116 | } // namespace |
| 117 | |
| 118 | ExportedHandler ExportedHandler::clone(jsg::Lock& js) { |
| 119 | return ExportedHandler{ |
| 120 | .fetch{mapAddRef(js, fetch)}, |
| 121 | .connect{mapAddRef(js, connect)}, |
| 122 | .tail{mapAddRef(js, tail)}, |
| 123 | .trace{mapAddRef(js, trace)}, |
| 124 | .tailStream{mapAddRef(js, tailStream)}, |
| 125 | .scheduled{mapAddRef(js, scheduled)}, |
| 126 | .alarm{mapAddRef(js, alarm)}, |
| 127 | .test{mapAddRef(js, test)}, |
| 128 | .webSocketMessage{mapAddRef(js, webSocketMessage)}, |
| 129 | .webSocketClose{mapAddRef(js, webSocketClose)}, |
| 130 | .webSocketError{mapAddRef(js, webSocketError)}, |
| 131 | .self{js.v8Isolate, self.getHandle(js.v8Isolate)}, |
| 132 | .env{env.addRef(js)}, |
| 133 | .ctx{getCtx()}, |
| 134 | .missingSuperclass = missingSuperclass, |
| 135 | }; |
| 136 | } |
| 137 | |
| 138 | ServiceWorkerGlobalScope::ServiceWorkerGlobalScope() |
| 139 | : unhandledRejections([this](jsg::Lock& js, |
| 140 | v8::PromiseRejectEvent event, |
| 141 | jsg::V8Ref<v8::Promise> promise, |
| 142 | jsg::Value value) { |
| 143 | // If async context tracking is enabled, then we need to ensure that we enter the frame |
| 144 | // associated with the promise before we invoke the unhandled rejection callback handling. |
| 145 | auto ev = js.alloc<PromiseRejectionEvent>(event, kj::mv(promise), kj::mv(value)); |
| 146 | dispatchEventImpl(js, kj::mv(ev)); |
| 147 | }) {} |
| 148 | |
| 149 | void ServiceWorkerGlobalScope::clear() { |
| 150 | removeAllHandlers(); |
| 151 | unhandledRejections.clear(); |
| 152 | } |
| 153 | |
| 154 | kj::Promise<void> ServiceWorkerGlobalScope::connect(kj::String host, |
| 155 | const kj::HttpHeaders& headers, |
| 156 | kj::AsyncIoStream& connection, |
| 157 | kj::HttpService::ConnectResponse& response, |
| 158 | Worker::Lock& lock, |
| 159 | kj::Maybe<ExportedHandler&> exportedHandler) { |
| 160 | ExportedHandler& eh = JSG_REQUIRE_NONNULL(exportedHandler, Error, |
| 161 | "Connect ingress is not currently supported with Service Workers syntax."); |
| 162 | KJ_REQUIRE(FeatureFlags::get(lock).getWorkerdExperimental(), |
| 163 | "connect handling requires the experimental flag."); |
| 164 | |
| 165 | KJ_IF_SOME(handler, eh.connect) { |
| 166 | // Has a connect handler! |
| 167 | response.accept(200, "OK", headers); |
| 168 | |
| 169 | // Using neuterable stream to manage lifetime of stream promises |
| 170 | auto ownConnection = newNeuterableIoStream(connection); |
| 171 | |
| 172 | auto& ioContext = IoContext::current(); |
| 173 | jsg::Lock& js = lock; |
| 174 | |
| 175 | // TLS support is not implemented so far. Note that setupSocket() expects the domain parameter |
| 176 | // to be set to the expected host name using startTLS, so that it can be provided to the TLS |
| 177 | // callback, so we'd need to change that or figure out a way to get the host domain. |
| 178 | auto nullTlsStarter = kj::heap<kj::TlsStarterCallback>(); |
| 179 | // We set isDefaultFetchPort to false here – sockets.c++ sets it for ports 443 and 8080 to |
| 180 | // provide a more descriptive error message for HTTP, but this is not relevant on the TCP server |
| 181 | // side. |
| 182 | jsg::Ref<Socket> jsSocket = |
| 183 | setupSocket(js, kj::mv(ownConnection), kj::none /* remoteAddress */, kj::mv(host), kj::none, |
| 184 | kj::mv(nullTlsStarter), SecureTransportKind::OFF, kj::none, false, kj::none); |
| 185 | // handleProxyStatus() is required to indicate that the socket was opened properly. Since the |
| 186 | // connection is already open at this point, exception handling is not required. |
| 187 | jsSocket->handleProxyStatus(js, kj::Promise<kj::Maybe<kj::Exception>>(kj::none)); |
| 188 | |
| 189 | kj::Maybe<SpanBuilder> span = ioContext.makeTraceSpan("connect_handler"_kjc); |
| 190 | auto promise = handler(js, kj::mv(jsSocket), eh.env.addRef(js), eh.getCtx()); |
| 191 | return ioContext.awaitJs(js, kj::mv(promise)).attach(kj::mv(span)); |
| 192 | } |
| 193 | lock.logWarningOnce("Received a connect event but we lack a handler. " |
| 194 | "Did you remember to export a connect() function?"); |
| 195 | JSG_FAIL_REQUIRE(Error, "Handler does not export a connect() function."); |
| 196 | } |
| 197 | |
| 198 | kj::Promise<DeferredProxy<void>> ServiceWorkerGlobalScope::request(kj::HttpMethod method, |
| 199 | kj::StringPtr url, |
| 200 | const kj::HttpHeaders& headers, |
| 201 | kj::AsyncInputStream& requestBody, |
| 202 | kj::HttpService::Response& response, |
| 203 | kj::Maybe<kj::StringPtr> cfBlobJson, |
| 204 | Worker::Lock& lock, |
| 205 | kj::Maybe<ExportedHandler&> exportedHandler, |
| 206 | kj::Maybe<jsg::Ref<AbortSignal>> abortSignal) { |
| 207 | TRACE_EVENT("workerd", "ServiceWorkerGlobalScope::request()"); |
| 208 | // To construct a ReadableStream object, we're supposed to pass in an Own<AsyncInputStream>, so |
| 209 | // that it can drop the reference whenever it gets GC'ed. But in this case the stream's lifetime |
| 210 | // is not under our control -- it's attached to the request. So, we wrap it in a |
| 211 | // NeuterableInputStream which allows us to disconnect the stream before it becomes invalid. |
| 212 | auto ownRequestBody = newNeuterableInputStream(requestBody); |
| 213 | auto deferredNeuter = kj::defer([ownRequestBody = kj::addRef(*ownRequestBody)]() mutable { |
| 214 | // Make sure to cancel the request body stream since the native stream is no longer valid once |
| 215 | // the returned promise completes. Note that the KJ HTTP library deals with the fact that we |
| 216 | // haven't consumed the entire request body. |
| 217 | ownRequestBody->neuter(makeNeuterException(NeuterReason::CLIENT_DISCONNECTED)); |
| 218 | }); |
| 219 | KJ_ON_SCOPE_FAILURE(ownRequestBody->neuter(makeNeuterException(NeuterReason::THREW_EXCEPTION))); |
| 220 | |
| 221 | auto& ioContext = IoContext::current(); |
| 222 | jsg::Lock& js = lock; |
| 223 | |
| 224 | CfProperty cf(cfBlobJson); |
| 225 | |
| 226 | // We only create the body stream if there is a body to read. |
| 227 | kj::Maybe<jsg::Ref<ReadableStream>> maybeJsStream = kj::none; |
| 228 | |
| 229 | // If the request has "no body", we want `request.body` to be null. But, this is not the same |
| 230 | // thing as the request having a body that happens to be empty. Unfortunately, KJ HTTP gives us |
| 231 | // a zero-length AsyncInputStream either way, so we can't just check the stream length. |
| 232 | // |
| 233 | // The HTTP spec says: "The presence of a message body in a request is signaled by a |
| 234 | // Content-Length or Transfer-Encoding header field." RFC 7230, section 3.3. |
| 235 | // https://tools.ietf.org/html/rfc7230#section-3.3 |
| 236 | // |
| 237 | // But, the request was not necessarily received over HTTP! It could be from another worker in |
| 238 | // a pipeline, or it could have been received over RPC. In either case, the headers don't |
| 239 | // necessarily mean anything; the calling worker can fill them in however it wants. |
| 240 | // |
| 241 | // So, we decide if the body is null if both headers are missing AND the stream is known to have |
| 242 | // zero length. And on the sending end (fetchImpl() in http.c++), if we're sending a request with |
| 243 | // a non-null body that is known to be empty, we explicitly set Content-Length: 0. This should |
| 244 | // mean that in all worker-to-worker interactions, if the sender provided a non-null body, the |
| 245 | // receiver will receive a non-null body, independent of anything else. |
| 246 | // |
| 247 | // TODO(cleanup): Should KJ HTTP interfaces explicitly communicate the difference between a |
| 248 | // missing body and an empty one? |
| 249 | kj::Maybe<Body::ExtractedBody> body; |
| 250 | if (headers.get(kj::HttpHeaderId::CONTENT_LENGTH) != kj::none || |
| 251 | headers.get(kj::HttpHeaderId::TRANSFER_ENCODING) != kj::none || |
| 252 | requestBody.tryGetLength().orDefault(1) > 0) { |
| 253 | // We do not automatically decode gzipped request bodies because the fetch() standard doesn't |
| 254 | // specify any automatic encoding of requests. https://github.com/whatwg/fetch/issues/589 |
| 255 | auto b = newSystemStream(kj::addRef(*ownRequestBody), StreamEncoding::IDENTITY); |
| 256 | auto jsStream = js.alloc<ReadableStream>(ioContext, kj::mv(b)); |
| 257 | body = Body::ExtractedBody(jsStream.addRef()); |
| 258 | maybeJsStream = kj::mv(jsStream); |
| 259 | } |
| 260 | |
| 261 | auto jsHeaders = js.alloc<Headers>(js, headers, Headers::Guard::REQUEST); |
| 262 | |
| 263 | // If the request doesn't specify "Content-Length" or "Transfer-Encoding", set "Content-Length" |
| 264 | // to the body length if it's known. This ensures handlers for worker-to-worker requests can |
| 265 | // access known body lengths if they're set, without buffering bodies. |
| 266 | if (body != kj::none && !jsHeaders->hasCommon(capnp::CommonHeaderName::CONTENT_LENGTH) && |
| 267 | !jsHeaders->hasCommon(capnp::CommonHeaderName::TRANSFER_ENCODING)) { |
| 268 | KJ_IF_SOME(l, requestBody.tryGetLength()) { |
| 269 | jsHeaders->setCommon(capnp::CommonHeaderName::CONTENT_LENGTH, kj::str(l)); |
| 270 | } else { |
| 271 | jsHeaders->setCommon(capnp::CommonHeaderName::TRANSFER_ENCODING, kj::str("chunked")); |
| 272 | } |
| 273 | } |
| 274 | |
| 275 | if (defaultFetcher == kj::none) { |
| 276 | // The default fetcher can be shared across requests since it does not persist any |
| 277 | // request-specific state. We can create one lazily here for the first request |
| 278 | // handled by this global scope and avoid the extra allocations on subsequence requests. |
| 279 | // We could also consider creating it lazily when first accessed, but Fetcher is very |
| 280 | // lightweight so this seems sufficient for now. |
| 281 | defaultFetcher = |
| 282 | js.alloc<Fetcher>(IoContext::NEXT_CLIENT_CHANNEL, Fetcher::RequiresHostAndProtocol::YES); |
| 283 | } |
| 284 | |
| 285 | auto jsRequest = js.alloc<Request>(js, method, url, Request::Redirect::MANUAL, kj::mv(jsHeaders), |
| 286 | KJ_ASSERT_NONNULL(defaultFetcher).addRef(), |
| 287 | /* signal */ kj::mv(abortSignal), kj::mv(cf), kj::mv(body), |
| 288 | /* thisSignal */ kj::none, Request::CacheMode::NONE); |
| 289 | |
| 290 | // signal vs thisSignal |
| 291 | // -------------------- |
| 292 | // The fetch spec definition of Request has a distinction between the |
| 293 | // "signal" (which is an optional AbortSignal passed in with the options), and "this' signal", |
| 294 | // which is an AbortSignal that is always available via the request.signal accessor. |
| 295 | // |
| 296 | // redirect |
| 297 | // -------- |
| 298 | // I set the redirect mode to manual here, so that by default scripts that just pass requests |
| 299 | // through to a fetch() call will behave the same as scripts which don't call .respondWith(): if |
| 300 | // the request results in a redirect, the visitor will see that redirect. |
| 301 | |
| 302 | auto event = js.alloc<FetchEvent>(kj::mv(jsRequest)); |
| 303 | |
| 304 | uint tasksBefore = ioContext.taskCount(); |
| 305 | |
| 306 | // We'll drop our span once the promise (fetch handler result) resolves. |
| 307 | kj::Maybe<SpanBuilder> span = ioContext.makeTraceSpan("fetch_handler"_kjc); |
| 308 | bool useDefaultHandling; |
| 309 | KJ_IF_SOME(h, exportedHandler) { |
| 310 | KJ_IF_SOME(f, h.fetch) { |
| 311 | auto promise = f(lock, event->getRequest(), h.env.addRef(js), h.getCtx()); |
| 312 | event->respondWith(lock, kj::mv(promise)); |
| 313 | useDefaultHandling = false; |
| 314 | } else { |
| 315 | // In modules mode we don't have a concept of "default handling". |
| 316 | lock.logWarningOnce("Received a FetchEvent but we lack a handler for FetchEvents. " |
| 317 | "Did you remember to export a fetch() function?"); |
| 318 | JSG_FAIL_REQUIRE(Error, "Handler does not export a fetch() function."); |
| 319 | } |
| 320 | } else { |
| 321 | // Fire off the handlers. |
| 322 | useDefaultHandling = dispatchEventImpl(lock, event.addRef()); |
| 323 | } |
| 324 | |
| 325 | if (useDefaultHandling) { |
| 326 | // No one called respondWith() or preventDefault(). Go directly to subrequest. |
| 327 | |
| 328 | if (ioContext.taskCount() > tasksBefore) { |
| 329 | lock.logWarning( |
| 330 | "FetchEvent handler did not call respondWith() before returning, but initiated some " |
| 331 | "asynchronous task. That task will be canceled and default handling will occur -- the " |
| 332 | "request will be sent unmodified to your origin. Remember that you must call " |
| 333 | "respondWith() *before* the event handler returns, if you don't want default handling. " |
| 334 | "You cannot call it asynchronously later on. If you need to wait for I/O (e.g. a " |
| 335 | "subrequest) before generating a Response, then call respondWith() with a Promise (for " |
| 336 | "the eventual Response) as the argument."); |
| 337 | } |
| 338 | |
| 339 | KJ_IF_SOME(jsStream, maybeJsStream) { |
| 340 | if (jsStream->isDisturbed()) { |
| 341 | lock.logUncaughtException( |
| 342 | "Script consumed request body but didn't call respondWith(). Can't forward request."); |
| 343 | return addNoopDeferredProxy( |
| 344 | response.sendError(500, "Internal Server Error", ioContext.getHeaderTable())); |
| 345 | } |
| 346 | maybeJsStream = kj::none; // Release our reference; we don't need it anymore. |
| 347 | } |
| 348 | |
| 349 | auto client = ioContext.getHttpClient( |
| 350 | IoContext::NEXT_CLIENT_CHANNEL, false, mapCopyString(cfBlobJson), "fetch_default"_kjc); |
| 351 | auto adapter = kj::newHttpService(*client); |
| 352 | auto promise = adapter->request(method, url, headers, requestBody, response); |
| 353 | // Default handling doesn't rely on the IoContext at all so we can return it as a |
| 354 | // deferred proxy task. |
| 355 | return DeferredProxy<void>{promise.attach(kj::mv(adapter), kj::mv(client))}; |
| 356 | } else KJ_IF_SOME(promise, event->getResponsePromise(lock)) { |
| 357 | maybeJsStream = kj::none; // Release reference to request body stream. We don't need it. |
| 358 | auto body2 = kj::addRef(*ownRequestBody); |
| 359 | |
| 360 | // HACK: If the client disconnects, the `response` reference is no longer valid. But our |
| 361 | // promise resolves in JavaScript space, so won't be canceled. So we need to track |
| 362 | // cancellation separately. We use a weird refcounted boolean. |
| 363 | // TODO(cleanup): Is there something less ugly we can do here? |
| 364 | auto canceled = kj::refcounted<kj::RefcountedWrapper<bool>>(false); |
| 365 | |
| 366 | return ioContext |
| 367 | .awaitJs(lock, |
| 368 | promise.then(kj::implicitCast<jsg::Lock&>(lock), |
| 369 | ioContext.addFunctor( |
| 370 | [&response, allowWebSocket = headers.isWebSocket(), |
| 371 | canceled = canceled->addWrappedRef(), &headers, span = kj::mv(span)]( |
| 372 | jsg::Lock& js, jsg::Ref<Response> innerResponse) mutable |
| 373 | -> IoOwn<kj::Promise<DeferredProxy<void>>> { |
| 374 | JSG_REQUIRE(innerResponse->getType() != "error"_kj, TypeError, |
| 375 | "Return value from serve handler must not be an error response (like Response.error())"); |
| 376 | |
| 377 | auto& context = IoContext::current(); |
| 378 | // Drop our fetch_handler span now that the promise has resolved. |
| 379 | span = kj::none; |
| 380 | if (*canceled) { |
| 381 | // Oops, the client disconnected before the response was ready to send. `response` is |
| 382 | // a dangling reference, let's not use it. |
| 383 | return context.addObject(kj::heap(addNoopDeferredProxy(kj::READY_NOW))); |
| 384 | } else { |
| 385 | return context.addObject(kj::heap( |
| 386 | innerResponse->send(js, response, {.allowWebSocket = allowWebSocket}, headers))); |
| 387 | } |
| 388 | }))) |
| 389 | .attach( |
| 390 | kj::defer([canceled = kj::mv(canceled)]() mutable { canceled->getWrapped() = true; })) |
| 391 | .then( |
| 392 | [ownRequestBody = kj::mv(ownRequestBody), deferredNeuter = kj::mv(deferredNeuter)]( |
| 393 | DeferredProxy<void> deferredProxy) mutable { |
| 394 | // In the case of bidirectional streaming, the request body stream needs to remain valid |
| 395 | // while proxying the response. So, arrange for neutering to happen only after the proxy |
| 396 | // task finishes. |
| 397 | deferredProxy.proxyTask = deferredProxy.proxyTask |
| 398 | .then([body = kj::addRef(*ownRequestBody)]() mutable { |
| 399 | body->neuter(makeNeuterException(NeuterReason::SENT_RESPONSE)); |
| 400 | }, [body = kj::addRef(*ownRequestBody)](kj::Exception&& e) mutable { |
| 401 | body->neuter(makeNeuterException(NeuterReason::THREW_EXCEPTION)); |
| 402 | kj::throwFatalException(kj::mv(e)); |
| 403 | }).attach(kj::mv(deferredNeuter)); |
| 404 | |
| 405 | return deferredProxy; |
| 406 | }, |
| 407 | [body = kj::mv(body2)](kj::Exception&& e) mutable -> DeferredProxy<void> { |
| 408 | // HACK: We depend on the fact that the success-case lambda above hasn't been destroyed yet |
| 409 | // so `deferredNeuter` hasn't been destroyed yet. |
| 410 | body->neuter(makeNeuterException(NeuterReason::THREW_EXCEPTION)); |
| 411 | kj::throwFatalException(kj::mv(e)); |
| 412 | }); |
| 413 | } else { |
| 414 | // The service worker API says that if default handling is prevented and respondWith() wasn't |
| 415 | // called, the request should result in "a network error". |
| 416 | return KJ_EXCEPTION(DISCONNECTED, "preventDefault() called but respondWith() not called"); |
| 417 | } |
| 418 | } |
| 419 | |
| 420 | void ServiceWorkerGlobalScope::sendTraces(kj::ArrayPtr<kj::Own<Trace>> traces, |
| 421 | Worker::Lock& lock, |
| 422 | kj::Maybe<ExportedHandler&> exportedHandler) { |
| 423 | auto isolate = lock.getIsolate(); |
| 424 | jsg::Lock& js = lock; |
| 425 | |
| 426 | KJ_IF_SOME(h, exportedHandler) { |
| 427 | KJ_IF_SOME(f, h.tail) { |
| 428 | auto tailEvent = js.alloc<TailEvent>(lock, "tail"_kjc, traces); |
| 429 | auto promise = f(lock, tailEvent->getEvents(), h.env.addRef(isolate), h.getCtx()); |
| 430 | tailEvent->waitUntil(kj::mv(promise)); |
| 431 | } else KJ_IF_SOME(f, h.trace) { |
| 432 | auto traceEvent = js.alloc<TailEvent>(lock, "trace"_kjc, traces); |
| 433 | auto promise = f(lock, traceEvent->getEvents(), h.env.addRef(isolate), h.getCtx()); |
| 434 | traceEvent->waitUntil(kj::mv(promise)); |
| 435 | } else { |
| 436 | lock.logWarningOnce("Attempted to send events but we lack a handler, " |
| 437 | "did you remember to export a tail() function?"); |
| 438 | JSG_FAIL_REQUIRE(Error, "Handler does not export a tail() function."); |
| 439 | } |
| 440 | } else { |
| 441 | // Fire off the handlers. |
| 442 | // We only create both events here. |
| 443 | auto tailEvent = js.alloc<TailEvent>(lock, "tail"_kjc, traces); |
| 444 | auto traceEvent = js.alloc<TailEvent>(lock, "trace"_kjc, traces); |
| 445 | dispatchEventImpl(lock, tailEvent.addRef()); |
| 446 | dispatchEventImpl(lock, traceEvent.addRef()); |
| 447 | |
| 448 | // We assume no action is necessary for "default" trace handling. |
| 449 | } |
| 450 | } |
| 451 | |
| 452 | void ServiceWorkerGlobalScope::startScheduled(kj::Date scheduledTime, |
| 453 | kj::StringPtr cron, |
| 454 | Worker::Lock& lock, |
| 455 | kj::Maybe<ExportedHandler&> exportedHandler) { |
| 456 | auto& context = IoContext::current(); |
| 457 | jsg::Lock& js = lock; |
| 458 | |
| 459 | double eventTime = (scheduledTime - kj::UNIX_EPOCH) / kj::MILLISECONDS; |
| 460 | |
| 461 | auto event = js.alloc<ScheduledEvent>(eventTime, cron); |
| 462 | |
| 463 | auto isolate = lock.getIsolate(); |
| 464 | |
| 465 | KJ_IF_SOME(h, exportedHandler) { |
| 466 | KJ_IF_SOME(f, h.scheduled) { |
| 467 | auto promise = |
| 468 | f(lock, js.alloc<ScheduledController>(event.addRef()), h.env.addRef(isolate), h.getCtx()) |
| 469 | .then([&context]() { |
| 470 | KJ_IF_SOME(t, context.getWorkerTracer()) { |
| 471 | t.setReturn(context.now()); |
| 472 | } |
| 473 | }); |
| 474 | event->waitUntil(kj::mv(promise)); |
| 475 | } else { |
| 476 | lock.logWarningOnce( |
| 477 | "Received a ScheduledEvent but we lack a handler for ScheduledEvents " |
| 478 | "(a.k.a. Cron Triggers). Did you remember to export a scheduled() function?"); |
| 479 | context.setNoRetryScheduled(); |
| 480 | JSG_FAIL_REQUIRE(Error, "Handler does not export a scheduled() function"); |
| 481 | } |
| 482 | } else { |
| 483 | // Fire off the handlers after confirming there is at least one. |
| 484 | if (getHandlerCount("scheduled") == 0) { |
| 485 | lock.logWarningOnce( |
| 486 | "Received a ScheduledEvent but we lack an event listener for scheduled events " |
| 487 | "(a.k.a. Cron Triggers). Did you remember to call addEventListener(\"scheduled\", ...)?"); |
| 488 | context.setNoRetryScheduled(); |
| 489 | JSG_FAIL_REQUIRE(Error, "No event listener registered for scheduled events."); |
| 490 | } |
| 491 | dispatchEventImpl(lock, event.addRef()); |
| 492 | } |
| 493 | } |
| 494 | |
| 495 | namespace { |
| 496 | // Returns true if an alarm failure should count against the user's retry limit. |
| 497 | // A failure is user-generated if any of: |
| 498 | // - The exception was explicitly tagged with EXCEPTION_IS_USER_ERROR at construction time |
| 499 | // (e.g. state.abort(), exceededCpu, exceededMemory, overload queue). |
| 500 | // - The exception originated from user code throwing inside blockConcurrencyWhile, which |
| 501 | // breaks the input gate as a secondary side-effect. |
| 502 | // - The exception is a plain jsg.* error without broken.* or jsg-internal.* prefixes, |
| 503 | // meaning the user's handler threw directly. |
| 504 | bool isAlarmFailureUserError(kj::StringPtr description, bool hasUserErrorDetail) { |
| 505 | if (hasUserErrorDetail) return true; |
| 506 | if (jsg::isExceptionFromInputGateBroken(description)) return true; |
| 507 | auto tunneled = jsg::tunneledErrorType(description); |
| 508 | return tunneled.isJsgError && !tunneled.isInternal && !tunneled.isDurableObjectReset; |
| 509 | } |
| 510 | } // namespace |
| 511 | |
| 512 | kj::Promise<WorkerInterface::AlarmResult> ServiceWorkerGlobalScope::runAlarm(kj::Date scheduledTime, |
| 513 | kj::Duration timeout, |
| 514 | uint32_t retryCount, |
| 515 | Worker::Lock& lock, |
| 516 | kj::Maybe<ExportedHandler&> exportedHandler) { |
| 517 | |
| 518 | auto& context = IoContext::current(); |
| 519 | auto& actor = KJ_ASSERT_NONNULL(context.getActor()); |
| 520 | auto& persistent = KJ_ASSERT_NONNULL(actor.getPersistent()); |
| 521 | |
| 522 | kj::String actorId; |
| 523 | KJ_SWITCH_ONEOF(actor.getId()) { |
| 524 | KJ_CASE_ONEOF(f, kj::Own<ActorIdFactory::ActorId>) { |
| 525 | actorId = f->toString(); |
| 526 | } |
| 527 | KJ_CASE_ONEOF(s, kj::String) { |
| 528 | actorId = kj::str(s); |
| 529 | } |
| 530 | } |
| 531 | |
| 532 | auto currentTime = context.now(); |
| 533 | KJ_SWITCH_ONEOF(persistent.armAlarmHandler( |
| 534 | scheduledTime, context.getCurrentTraceSpan(), currentTime, false, actorId)) { |
| 535 | KJ_CASE_ONEOF(armResult, ActorCacheInterface::RunAlarmHandler) { |
| 536 | auto& handler = KJ_REQUIRE_NONNULL(exportedHandler); |
| 537 | if (handler.alarm == kj::none) { |
| 538 | |
| 539 | lock.logWarningOnce("Attempted to run a scheduled alarm without a handler, " |
| 540 | "did you remember to export an alarm() function?"); |
| 541 | return WorkerInterface::AlarmResult{ |
| 542 | .retry = false, .outcome = EventOutcome::SCRIPT_NOT_FOUND}; |
| 543 | } |
| 544 | |
| 545 | auto& alarm = KJ_ASSERT_NONNULL(handler.alarm); |
| 546 | |
| 547 | return context |
| 548 | .run([exportedHandler, &context, timeout, retryCount, scheduledTime, &alarm, |
| 549 | maybeAsyncContext = jsg::AsyncContextFrame::currentRef(lock)]( |
| 550 | Worker::Lock& lock) mutable -> kj::Promise<WorkerInterface::AlarmResult> { |
| 551 | jsg::AsyncContextFrame::Scope asyncScope(lock, maybeAsyncContext); |
| 552 | // We want to limit alarm handler walltime to 15 minutes at most. If the timeout promise |
| 553 | // completes we want to cancel the alarm handler. If the alarm handler promise completes |
| 554 | // first timeout will be canceled. |
| 555 | jsg::Lock& js = lock; |
| 556 | auto timeoutPromise = context.afterLimitTimeout(timeout).then( |
| 557 | [&context]() -> kj::Promise<WorkerInterface::AlarmResult> { |
| 558 | // We don't want to delete the alarm since we have not successfully completed the alarm |
| 559 | // execution. |
| 560 | auto& actor = KJ_ASSERT_NONNULL(context.getActor()); |
| 561 | auto& persistent = KJ_ASSERT_NONNULL(actor.getPersistent()); |
| 562 | persistent.cancelDeferredAlarmDeletion(); |
| 563 | |
| 564 | LOG_NOSENTRY(WARNING, "Alarm exceeded its allowed execution time"); |
| 565 | // Report alarm handler failure and log it. |
| 566 | auto e = KJ_EXCEPTION(OVERLOADED, |
| 567 | "broken.dropped; worker_do_not_log; jsg.Error: Alarm exceeded its allowed execution time"); |
| 568 | e.setDetail(jsg::EXCEPTION_IS_USER_ERROR, kj::heapArray<kj::byte>(0)); |
| 569 | e.setDetail(CPU_LIMIT_DETAIL_ID, kj::heapArray<kj::byte>(0)); |
| 570 | context.getMetrics().reportFailure(e); |
| 571 | |
| 572 | // We don't want the handler to keep running after timeout. |
| 573 | context.abort(kj::mv(e)); |
| 574 | // We want timed out alarms to be treated as user errors. As such, we'll mark them as |
| 575 | // retriable, and we'll count the retries against the alarm retries limit. This will ensure |
| 576 | // that the handler will attempt to run for a number of times before giving up and deleting |
| 577 | // the alarm. |
| 578 | return WorkerInterface::AlarmResult{ |
| 579 | .retry = true, .retryCountsAgainstLimit = true, .outcome = EventOutcome::EXCEEDED_CPU}; |
| 580 | }); |
| 581 | |
| 582 | return alarm(lock, js.alloc<AlarmInvocationInfo>(scheduledTime, retryCount)) |
| 583 | .then([]() -> kj::Promise<WorkerInterface::AlarmResult> { |
| 584 | return WorkerInterface::AlarmResult{.retry = false, .outcome = EventOutcome::OK}; |
| 585 | }).exclusiveJoin(kj::mv(timeoutPromise)); |
| 586 | }) |
| 587 | .catch_([&context, deferredDelete = kj::mv(armResult.deferredDelete)]( |
| 588 | kj::Exception&& e) mutable { |
| 589 | auto& actor = KJ_ASSERT_NONNULL(context.getActor()); |
| 590 | auto& persistent = KJ_ASSERT_NONNULL(actor.getPersistent()); |
| 591 | persistent.cancelDeferredAlarmDeletion(); |
| 592 | |
| 593 | context.getMetrics().reportFailure(e); |
| 594 | |
| 595 | auto description = kj::str(e.getDescription()); // because e is moved before this is used |
| 596 | auto log = !jsg::isTunneledException(description) && !jsg::isDoNotLogException(description); |
| 597 | auto isUserError = e.getDetail(jsg::EXCEPTION_IS_USER_ERROR) != kj::none; |
| 598 | |
| 599 | // This will include the error in inspector/tracers and log to syslog if internal. |
| 600 | context.logUncaughtExceptionAsync(UncaughtExceptionSource::ALARM_HANDLER, kj::mv(e)); |
| 601 | |
| 602 | EventOutcome outcome = EventOutcome::EXCEPTION; |
| 603 | KJ_IF_SOME(status, context.getLimitEnforcer().getLimitsExceeded()) { |
| 604 | outcome = status; |
| 605 | } |
| 606 | |
| 607 | kj::String actorId; |
| 608 | KJ_SWITCH_ONEOF(actor.getId()) { |
| 609 | KJ_CASE_ONEOF(f, kj::Own<ActorIdFactory::ActorId>) { |
| 610 | actorId = f->toString(); |
| 611 | } |
| 612 | KJ_CASE_ONEOF(s, kj::String) { |
| 613 | actorId = kj::str(s); |
| 614 | } |
| 615 | } |
| 616 | |
| 617 | auto isUserGeneratedError = isAlarmFailureUserError(description, isUserError); |
| 618 | auto shouldRetryCountsAgainstLimits = !context.isOutputGateBroken() || isUserGeneratedError; |
| 619 | |
| 620 | // We want to alert if we aren't going to count this alarm retry against limits. |
| 621 | // Skip logging when the output gate broke as a secondary effect of a user-generated error: |
| 622 | // that is expected behaviour and already counted as a user error. |
| 623 | if (!isUserGeneratedError && log && context.isOutputGateBroken()) { |
| 624 | LOG_NOSENTRY(ERROR, "output lock broke during alarm execution", actorId, description); |
| 625 | } else if (!isUserGeneratedError && context.isOutputGateBroken()) { |
| 626 | // Tunneled or do-not-log non-user error with a broken output gate. Log for diagnostics |
| 627 | // so we can investigate stuck alarms. |
| 628 | LOG_NOSENTRY(ERROR, |
| 629 | "output lock broke during alarm execution without an interesting error description", |
| 630 | actorId, description, shouldRetryCountsAgainstLimits); |
| 631 | } |
| 632 | return WorkerInterface::AlarmResult{.retry = true, |
| 633 | .retryCountsAgainstLimit = shouldRetryCountsAgainstLimits, |
| 634 | .outcome = outcome, |
| 635 | .errorDescription = kj::str(description)}; |
| 636 | }) |
| 637 | .then([&context](WorkerInterface::AlarmResult result) |
| 638 | -> kj::Promise<WorkerInterface::AlarmResult> { |
| 639 | return context.waitForOutputLocks().then([result = kj::mv(result)]() mutable { |
| 640 | return kj::mv(result); |
| 641 | }, [&context](kj::Exception&& e) { |
| 642 | auto& actor = KJ_ASSERT_NONNULL(context.getActor()); |
| 643 | kj::String actorId; |
| 644 | KJ_SWITCH_ONEOF(actor.getId()) { |
| 645 | KJ_CASE_ONEOF(f, kj::Own<ActorIdFactory::ActorId>) { |
| 646 | actorId = f->toString(); |
| 647 | } |
| 648 | KJ_CASE_ONEOF(s, kj::String) { |
| 649 | actorId = kj::str(s); |
| 650 | } |
| 651 | } |
| 652 | auto isUserGeneratedError = isAlarmFailureUserError( |
| 653 | e.getDescription(), e.getDetail(jsg::EXCEPTION_IS_USER_ERROR) != kj::none); |
| 654 | auto shouldRetryCountsAgainstLimits = isUserGeneratedError; |
| 655 | if (auto desc = e.getDescription(); |
| 656 | !jsg::isTunneledException(desc) && !jsg::isDoNotLogException(desc)) { |
| 657 | if (!isUserGeneratedError) { |
| 658 | if (isInterestingException(e)) { |
| 659 | LOG_EXCEPTION("alarmOutputLock"_kj, e); |
| 660 | } else { |
| 661 | LOG_NOSENTRY(ERROR, "output lock broke after executing alarm", actorId, e); |
| 662 | } |
| 663 | } |
| 664 | } else if (!isUserGeneratedError) { |
| 665 | // Tunneled or do-not-log exception that is not a user error. Forward as-is without |
| 666 | // counting against retry limits. |
| 667 | LOG_NOSENTRY(ERROR, |
| 668 | "output lock broke after executing alarm with tunneled non-user error", actorId, |
| 669 | e.getDescription()); |
| 670 | } |
| 671 | return WorkerInterface::AlarmResult{.retry = true, |
| 672 | .retryCountsAgainstLimit = shouldRetryCountsAgainstLimits, |
| 673 | .outcome = EventOutcome::EXCEPTION, |
| 674 | .errorDescription = kj::str(e.getDescription())}; |
| 675 | }); |
| 676 | }); |
| 677 | } |
| 678 | KJ_CASE_ONEOF(armResult, ActorCacheInterface::CancelAlarmHandler) { |
| 679 | return armResult.waitBeforeCancel.then([]() { |
| 680 | return WorkerInterface::AlarmResult{.retry = false, .outcome = EventOutcome::CANCELED}; |
| 681 | }); |
| 682 | } |
| 683 | } |
| 684 | KJ_UNREACHABLE; |
| 685 | } |
| 686 | |
| 687 | jsg::Promise<void> ServiceWorkerGlobalScope::test( |
| 688 | Worker::Lock& lock, kj::Maybe<ExportedHandler&> exportedHandler) { |
| 689 | // TODO(someday): For Service Workers syntax, do we want addEventListener("test")? Not supporting |
| 690 | // it for now. |
| 691 | ExportedHandler& eh = JSG_REQUIRE_NONNULL( |
| 692 | exportedHandler, Error, "Tests are not currently supported with Service Workers syntax."); |
| 693 | |
| 694 | auto& testHandler = |
| 695 | JSG_REQUIRE_NONNULL(eh.test, Error, "Entrypoint does not export a test() function."); |
| 696 | |
| 697 | jsg::Lock& js = lock; |
| 698 | return testHandler(lock, js.alloc<TestController>(), eh.env.addRef(lock), eh.getCtx()); |
| 699 | } |
| 700 | |
| 701 | // This promise is used to set the timeout for hibernatable websocket events. It's expected to be |
| 702 | // dropped in most cases, as long as the hibernatable websocket event promise completes before it. |
| 703 | kj::Promise<void> ServiceWorkerGlobalScope::eventTimeoutPromise(uint32_t timeoutMs) { |
| 704 | auto& actor = KJ_ASSERT_NONNULL(IoContext::current().getActor()); |
| 705 | co_await IoContext::current().afterLimitTimeout(timeoutMs * kj::MILLISECONDS); |
| 706 | // This is the ActorFlushReason for eviction in Cloudflare's internal implementation. |
| 707 | auto evictionCode = 2; |
| 708 | auto e = KJ_EXCEPTION(DISCONNECTED, |
| 709 | "broken.dropped; jsg.Error: Actor exceeded event execution time and was disconnected."); |
| 710 | e.setDetail(jsg::EXCEPTION_IS_USER_ERROR, kj::heapArray<kj::byte>(0)); |
| 711 | actor.shutdown(evictionCode, kj::mv(e)); |
| 712 | } |
| 713 | |
| 714 | kj::Promise<void> ServiceWorkerGlobalScope::setHibernatableEventTimeout( |
| 715 | kj::Promise<void> event, kj::Maybe<uint32_t> eventTimeoutMs) { |
| 716 | // If we have a maximum event duration timeout set, we should prevent the actor from running |
| 717 | // for more than the user selected duration. |
| 718 | auto timeoutMs = eventTimeoutMs.orDefault(static_cast<uint32_t>(0)); |
| 719 | if (timeoutMs > 0) { |
| 720 | return event.exclusiveJoin(eventTimeoutPromise(timeoutMs)); |
| 721 | } |
| 722 | return event; |
| 723 | } |
| 724 | |
| 725 | // TODO(cleanup): the hibernatable websocket handler functions here are largely identical – consider |
| 726 | // folding them. |
| 727 | // |
| 728 | // Note: The hibernatable WebSocket message path passes kj::OneOf<kj::String, kj::Array<byte>> |
| 729 | // directly to the webSocketMessage() handler, so binary data is always delivered as ArrayBuffer. |
| 730 | // The WebSocket binaryType property (and the websocket_standard_binary_type compat flag) has no |
| 731 | // effect here — this is by design for the Durable Object handler API, which bypasses the |
| 732 | // normal WebSocket read loop and its Blob/ArrayBuffer dispatch logic. |
| 733 | void ServiceWorkerGlobalScope::sendHibernatableWebSocketMessage(IoContext& context, |
| 734 | kj::OneOf<kj::String, kj::Array<byte>> message, |
| 735 | kj::Maybe<uint32_t> eventTimeoutMs, |
| 736 | kj::String websocketId, |
| 737 | Worker::Lock& lock, |
| 738 | kj::Maybe<ExportedHandler&> exportedHandler) { |
| 739 | jsg::Lock& js = lock; |
| 740 | auto event = js.alloc<HibernatableWebSocketEvent>(); |
| 741 | // Even if no handler is exported, we need to claim the websocket so it's removed from the map. |
| 742 | auto websocket = event->claimWebSocket(lock, websocketId); |
| 743 | |
| 744 | KJ_IF_SOME(h, exportedHandler) { |
| 745 | KJ_IF_SOME(handler, h.webSocketMessage) { |
| 746 | event->waitUntil(setHibernatableEventTimeout( |
| 747 | handler(lock, kj::mv(websocket), kj::mv(message)), eventTimeoutMs) |
| 748 | .then([&context]() { |
| 749 | KJ_IF_SOME(t, context.getWorkerTracer()) { |
| 750 | t.setReturn(context.now()); |
| 751 | } |
| 752 | })); |
| 753 | } |
| 754 | // We want to deliver a message, but if no webSocketMessage handler is exported, we shouldn't fail |
| 755 | } |
| 756 | } |
| 757 | |
| 758 | void ServiceWorkerGlobalScope::sendHibernatableWebSocketClose(IoContext& context, |
| 759 | HibernatableSocketParams::Close close, |
| 760 | kj::Maybe<uint32_t> eventTimeoutMs, |
| 761 | kj::String websocketId, |
| 762 | Worker::Lock& lock, |
| 763 | kj::Maybe<ExportedHandler&> exportedHandler) { |
| 764 | jsg::Lock& js = lock; |
| 765 | auto event = js.alloc<HibernatableWebSocketEvent>(); |
| 766 | |
| 767 | // Even if no handler is exported, we need to claim the websocket so it's removed from the map. |
| 768 | // |
| 769 | // We won't be dispatching any further events because we've received a close, so we return the |
| 770 | // owned websocket back to the api::WebSocket. |
| 771 | auto releasePackage = event->prepareForRelease(lock, websocketId); |
| 772 | auto websocket = kj::mv(releasePackage.webSocketRef); |
| 773 | websocket->initiateHibernatableRelease(lock, kj::mv(releasePackage.ownedWebSocket), |
| 774 | kj::mv(releasePackage.tags), api::WebSocket::HibernatableReleaseState::CLOSE); |
| 775 | KJ_IF_SOME(h, exportedHandler) { |
| 776 | KJ_IF_SOME(handler, h.webSocketClose) { |
| 777 | event->waitUntil(setHibernatableEventTimeout( |
| 778 | handler(lock, kj::mv(websocket), close.code, kj::mv(close.reason), close.wasClean), |
| 779 | eventTimeoutMs) |
| 780 | .then([&context]() { |
| 781 | KJ_IF_SOME(t, context.getWorkerTracer()) { |
| 782 | t.setReturn(context.now()); |
| 783 | } |
| 784 | })); |
| 785 | } |
| 786 | // We want to deliver close, but if no webSocketClose handler is exported, we shouldn't fail |
| 787 | } |
| 788 | } |
| 789 | |
| 790 | void ServiceWorkerGlobalScope::sendHibernatableWebSocketError(IoContext& context, |
| 791 | kj::Exception e, |
| 792 | kj::Maybe<uint32_t> eventTimeoutMs, |
| 793 | kj::String websocketId, |
| 794 | Worker::Lock& lock, |
| 795 | kj::Maybe<ExportedHandler&> exportedHandler) { |
| 796 | jsg::Lock& js = lock; |
| 797 | auto event = js.alloc<HibernatableWebSocketEvent>(); |
| 798 | |
| 799 | // Even if no handler is exported, we need to claim the websocket so it's removed from the map. |
| 800 | // |
| 801 | // We won't be dispatching any further events because we've encountered an error, so we return |
| 802 | // the owned websocket back to the api::WebSocket. |
| 803 | auto releasePackage = event->prepareForRelease(lock, websocketId); |
| 804 | auto& websocket = releasePackage.webSocketRef; |
| 805 | websocket->initiateHibernatableRelease(lock, kj::mv(releasePackage.ownedWebSocket), |
| 806 | kj::mv(releasePackage.tags), WebSocket::HibernatableReleaseState::ERROR); |
| 807 | |
| 808 | KJ_IF_SOME(h, exportedHandler) { |
| 809 | KJ_IF_SOME(handler, h.webSocketError) { |
| 810 | event->waitUntil(setHibernatableEventTimeout( |
| 811 | handler(js, kj::mv(websocket), js.exceptionToJs(kj::mv(e))), eventTimeoutMs) |
| 812 | .then([&context]() { |
| 813 | KJ_IF_SOME(t, context.getWorkerTracer()) { |
| 814 | t.setReturn(context.now()); |
| 815 | } |
| 816 | })); |
| 817 | } |
| 818 | // We want to deliver an error, but if no webSocketError handler is exported, we shouldn't fail |
| 819 | } |
| 820 | } |
| 821 | |
| 822 | void ServiceWorkerGlobalScope::emitPromiseRejection(jsg::Lock& js, |
| 823 | v8::PromiseRejectEvent event, |
| 824 | jsg::V8Ref<v8::Promise> promise, |
| 825 | jsg::Value value) { |
| 826 | |
| 827 | const auto hasHandlers = [this] { |
| 828 | return getHandlerCount("unhandledrejection"_kj) + getHandlerCount("rejectionhandled"_kj); |
| 829 | }; |
| 830 | |
| 831 | const auto hasInspector = [] { |
| 832 | KJ_IF_SOME(ioContext, IoContext::tryCurrent()) { |
| 833 | return ioContext.isInspectorEnabled(); |
| 834 | } else { |
| 835 | return false; |
| 836 | } |
| 837 | }; |
| 838 | |
| 839 | if (hasHandlers() || hasInspector()) { |
| 840 | unhandledRejections.setUseMicrotasksCompletedCallback( |
| 841 | FeatureFlags::get(js).getUnhandledRejectionAfterMicrotaskCheckpoint()); |
| 842 | unhandledRejections.report(js, event, kj::mv(promise), kj::mv(value)); |
| 843 | } |
| 844 | } |
| 845 | |
| 846 | void ServiceWorkerGlobalScope::setConnectOverride(kj::String networkAddress, ConnectFn connectFn) { |
| 847 | connectOverrides.upsert(kj::mv(networkAddress), kj::mv(connectFn)); |
| 848 | } |
| 849 | |
| 850 | kj::Maybe<ServiceWorkerGlobalScope::ConnectFn&> ServiceWorkerGlobalScope::getConnectOverride( |
| 851 | kj::StringPtr networkAddress) { |
| 852 | return connectOverrides.find(networkAddress); |
| 853 | } |
| 854 | |
| 855 | jsg::JsString ServiceWorkerGlobalScope::btoa(jsg::Lock& js, jsg::JsString str) { |
| 856 | // We could implement btoa() by accepting a kj::String, but then we'd have to check that it |
| 857 | // doesn't have any multibyte code points. Easier to perform that test using v8::String's |
| 858 | // ContainsOnlyOneByte() function. |
| 859 | JSG_REQUIRE(str.containsOnlyOneByte(), DOMInvalidCharacterError, |
| 860 | "btoa() can only operate on characters in the Latin1 (ISO/IEC 8859-1) range."); |
| 861 | auto strArray = str.toArray<kj::byte>(js); |
| 862 | auto expected_length = simdutf::base64_length_from_binary(strArray.size()); |
| 863 | KJ_STACK_ARRAY(kj::byte, result, expected_length, 1024, 1024); |
| 864 | auto written = simdutf::binary_to_base64( |
| 865 | strArray.asChars().begin(), strArray.size(), result.asChars().begin()); |
| 866 | return js.str(result.first(written)); |
| 867 | } |
| 868 | |
| 869 | #ifdef WORKERD_FUZZILLI |
| 870 | void ServiceWorkerGlobalScope::fuzzilli(jsg::Lock& js, jsg::Arguments<jsg::Value> args) { |
| 871 | // Delegate to the fuzzilli handler in fuzzilli.c++ |
| 872 | fuzzilli_handler(js, args); |
| 873 | } |
| 874 | #endif |
| 875 | |
| 876 | jsg::JsString ServiceWorkerGlobalScope::atob(jsg::Lock& js, kj::String data) { |
| 877 | auto decoded = kj::decodeBase64(data.asArray()); |
| 878 | |
| 879 | JSG_REQUIRE(!decoded.hadErrors, DOMInvalidCharacterError, |
| 880 | "atob() called with invalid base64-encoded data. (Only whitespace, '+', '/', alphanumeric " |
| 881 | "ASCII, and up to two terminal '=' signs when the input data length is divisible by 4 are " |
| 882 | "allowed.)"); |
| 883 | |
| 884 | // Similar to btoa() taking a v8::Value, we return a v8::String directly, as this allows us to |
| 885 | // construct a string from the non-nul-terminated array returned from decodeBase64(). This avoids |
| 886 | // making a copy purely to append a nul byte. |
| 887 | return js.str(decoded.asBytes()); |
| 888 | } |
| 889 | |
| 890 | void ServiceWorkerGlobalScope::queueMicrotask(jsg::Lock& js, jsg::Function<void()> task) { |
| 891 | auto fn = js.wrapSimpleFunction(js.v8Context(), |
| 892 | JSG_VISITABLE_LAMBDA((this, fn = kj::mv(task)), (fn), |
| 893 | (jsg::Lock& js, const v8::FunctionCallbackInfo<v8::Value>& args) { |
| 894 | js.tryCatch([&] { |
| 895 | // The function won't be called with any arguments, so we can |
| 896 | // safely ignore anything passed in to args. |
| 897 | fn(js); |
| 898 | }, [&](jsg::Value exception) { |
| 899 | // The reportError call itself can potentially throw errors. Let's catch |
| 900 | // and report them as well. |
| 901 | js.tryCatch([&] { reportError(js, jsg::JsValue(exception.getHandle(js))); }, |
| 902 | [&](jsg::Value exception) { |
| 903 | // An error was thrown by the 'error' event handler. That's unfortunate. |
| 904 | // Let's log the error and just continue. It won't be possible to actually |
| 905 | // catch or handle this error so logging is really the only way to notify |
| 906 | // folks about it. |
| 907 | auto val = jsg::JsValue(exception.getHandle(js)); |
| 908 | // If the value is an object that has a stack property, log that so we get |
| 909 | // the stack trace if it is an exception. |
| 910 | KJ_IF_SOME(obj, val.tryCast<jsg::JsObject>()) { |
| 911 | auto stack = obj.get(js, "stack"_kj); |
| 912 | if (!stack.isUndefined()) { |
| 913 | js.reportError(stack); |
| 914 | return; |
| 915 | } |
| 916 | } else { |
| 917 | } // Here to avoid a compile warning |
| 918 | // Otherwise just log the stringified value generically. |
| 919 | js.reportError(val); |
| 920 | }); |
| 921 | }); |
| 922 | })); |
| 923 | |
| 924 | js.v8Isolate->EnqueueMicrotask(fn); |
| 925 | } |
| 926 | |
| 927 | jsg::JsValue ServiceWorkerGlobalScope::structuredClone( |
| 928 | jsg::Lock& js, jsg::JsValue value, jsg::Optional<StructuredCloneOptions> maybeOptions) { |
| 929 | KJ_IF_SOME(options, maybeOptions) { |
| 930 | KJ_IF_SOME(transfer, options.transfer) { |
| 931 | auto transfers = KJ_MAP(i, transfer) { return i.getHandle(js); }; |
| 932 | return value.structuredClone(js, kj::mv(transfers)); |
| 933 | } |
| 934 | } |
| 935 | return value.structuredClone(js); |
| 936 | } |
| 937 | |
| 938 | TimeoutId::NumberType ServiceWorkerGlobalScope::setTimeoutInternal( |
| 939 | jsg::Function<void()> function, double msDelay) { |
| 940 | auto timeoutId = IoContext::current().setTimeoutImpl(timeoutIdGenerator, |
| 941 | /* repeat */ false, kj::mv(function), msDelay); |
| 942 | return timeoutId.toNumber(); |
| 943 | } |
| 944 | |
| 945 | TimeoutId::NumberType ServiceWorkerGlobalScope::setTimeout(jsg::Lock& js, |
| 946 | jsg::Function<void(jsg::Arguments<jsg::Value>)> function, |
| 947 | jsg::Optional<double> msDelay, |
| 948 | jsg::Arguments<jsg::Value> args) { |
| 949 | function.setReceiver(js.v8Ref<v8::Value>(js.v8Context()->Global())); |
| 950 | auto fn = [function = kj::mv(function), args = kj::mv(args), |
| 951 | context = jsg::AsyncContextFrame::currentRef(js)](jsg::Lock& js) mutable { |
| 952 | jsg::AsyncContextFrame::Scope scope(js, context); |
| 953 | function(js, kj::mv(args)); |
| 954 | }; |
| 955 | auto timeoutId = IoContext::current().setTimeoutImpl(timeoutIdGenerator, |
| 956 | /* repeat */ false, [function = kj::mv(fn)](jsg::Lock& js) mutable { function(js); }, |
| 957 | msDelay.orDefault(0)); |
| 958 | return timeoutId.toNumber(); |
| 959 | } |
| 960 | |
| 961 | void ServiceWorkerGlobalScope::clearTimeout(jsg::Lock& js, kj::Maybe<jsg::JsNumber> timeoutId) { |
| 962 | KJ_IF_SOME(rawId, timeoutId) { |
| 963 | // Browsers does not throw an error when "unsafe" integers are passed to the clearTimeout method. |
| 964 | // Let's make sure we ignore those values, just like browsers and other runtimes. |
| 965 | KJ_IF_SOME(id, rawId.toSafeInteger(js)) { |
| 966 | IoContext::current().clearTimeoutImpl(TimeoutId::fromNumber(id)); |
| 967 | } |
| 968 | } |
| 969 | } |
| 970 | |
| 971 | TimeoutId::NumberType ServiceWorkerGlobalScope::setInterval(jsg::Lock& js, |
| 972 | jsg::Function<void(jsg::Arguments<jsg::Value>)> function, |
| 973 | jsg::Optional<double> msDelay, |
| 974 | jsg::Arguments<jsg::Value> args) { |
| 975 | function.setReceiver(js.v8Ref<v8::Value>(js.v8Context()->Global())); |
| 976 | auto fn = [function = kj::mv(function), args = kj::mv(args), |
| 977 | context = jsg::AsyncContextFrame::currentRef(js)](jsg::Lock& js) mutable { |
| 978 | jsg::AsyncContextFrame::Scope scope(js, context); |
| 979 | // Because the fn is called multiple times, we will clone the args on each call. |
| 980 | auto argv = KJ_MAP(i, args) { return i.addRef(js); }; |
| 981 | function(js, jsg::Arguments(kj::mv(argv))); |
| 982 | }; |
| 983 | auto timeoutId = IoContext::current().setTimeoutImpl(timeoutIdGenerator, |
| 984 | /* repeat */ true, [function = kj::mv(fn)](jsg::Lock& js) mutable { function(js); }, |
| 985 | msDelay.orDefault(0)); |
| 986 | return timeoutId.toNumber(); |
| 987 | } |
| 988 | |
| 989 | void ServiceWorkerGlobalScope::clearInterval(jsg::Lock& js, kj::Maybe<jsg::JsNumber> timeoutId) { |
| 990 | clearTimeout(js, kj::mv(timeoutId)); |
| 991 | } |
| 992 | |
| 993 | jsg::Ref<Crypto> ServiceWorkerGlobalScope::getCrypto(jsg::Lock& js) { |
| 994 | return js.alloc<Crypto>(js); |
| 995 | } |
| 996 | |
| 997 | jsg::Ref<CacheStorage> ServiceWorkerGlobalScope::getCaches(jsg::Lock& js) { |
| 998 | return js.alloc<CacheStorage>(js); |
| 999 | } |
| 1000 | |
| 1001 | jsg::Promise<jsg::Ref<Response>> ServiceWorkerGlobalScope::fetch(jsg::Lock& js, |
| 1002 | kj::OneOf<jsg::Ref<Request>, kj::String> requestOrUrl, |
| 1003 | jsg::Optional<Request::Initializer> requestInit) { |
| 1004 | return fetchImpl(js, kj::none, kj::mv(requestOrUrl), kj::mv(requestInit)); |
| 1005 | } |
| 1006 | |
| 1007 | void ServiceWorkerGlobalScope::reportError(jsg::Lock& js, jsg::JsValue error) { |
| 1008 | // Per the spec, we are going to first emit an error event on the global object. |
| 1009 | // If that event is not prevented, we will log the error to the console. Note |
| 1010 | // that we do not throw the error at all. |
| 1011 | auto message = v8::Exception::CreateMessage(js.v8Isolate, error); |
| 1012 | auto event = js.alloc<ErrorEvent>(ErrorEvent::ErrorEventInit{.message = kj::str(message->Get()), |
| 1013 | .filename = kj::str(message->GetScriptResourceName()), |
| 1014 | .lineno = jsg::check(message->GetLineNumber(js.v8Context())), |
| 1015 | .colno = jsg::check(message->GetStartColumn(js.v8Context())), |
| 1016 | .error = jsg::JsRef(js, error)}); |
| 1017 | if (dispatchEventImpl(js, kj::mv(event))) { |
| 1018 | // If the value is an object that has a stack property, log that so we get |
| 1019 | // the stack trace if it is an exception. |
| 1020 | KJ_IF_SOME(obj, error.tryCast<jsg::JsObject>()) { |
| 1021 | auto stack = obj.get(js, "stack"_kj); |
| 1022 | if (!stack.isUndefined()) { |
| 1023 | js.reportError(stack); |
| 1024 | return; |
| 1025 | } |
| 1026 | } |
| 1027 | // Otherwise just log the stringified value generically. |
| 1028 | js.reportError(error); |
| 1029 | } |
| 1030 | } |
| 1031 | |
| 1032 | jsg::JsValue ServiceWorkerGlobalScope::getBuffer(jsg::Lock& js) { |
| 1033 | KJ_IF_SOME(p, bufferValue) { |
| 1034 | return p.getHandle(js); |
| 1035 | } |
| 1036 | constexpr auto kSpecifier = "node:buffer"_kj; |
| 1037 | KJ_IF_SOME(module, js.resolveModule(kSpecifier)) { |
| 1038 | KJ_IF_SOME(p, bufferValue) { |
| 1039 | // There's a chance that resolving the module caused side-effects |
| 1040 | // that set the bufferValue, we let's check again. |
| 1041 | return p.getHandle(js); |
| 1042 | } |
| 1043 | // When requireReturnsDefaultExport flag is enabled, resolveModule returns the |
| 1044 | // default export directly. Otherwise it returns the module namespace. |
| 1045 | auto obj = module.has(js, "default"_kj) |
| 1046 | ? KJ_ASSERT_NONNULL(module.get(js, "default"_kj).tryCast<jsg::JsObject>()) |
| 1047 | : module; |
| 1048 | auto buffer = obj.get(js, "Buffer"_kj); |
| 1049 | JSG_REQUIRE(buffer.isFunction(), TypeError, "Invalid node:buffer implementation"); |
| 1050 | bufferValue = jsg::JsRef(js, buffer); |
| 1051 | return buffer; |
| 1052 | } else { |
| 1053 | // If we are unable to resolve the node:buffer module here, it likely |
| 1054 | // means that we don't actually have a module registry installed. Just |
| 1055 | // return undefined in this case. |
| 1056 | bufferValue = jsg::JsRef(js, js.undefined()); |
| 1057 | return js.undefined(); |
| 1058 | } |
| 1059 | } |
| 1060 | |
| 1061 | void ServiceWorkerGlobalScope::setProcess(jsg::Lock& js, jsg::JsValue newProcess) { |
| 1062 | processValue = jsg::JsRef(js, newProcess); |
| 1063 | } |
| 1064 | |
| 1065 | void ServiceWorkerGlobalScope::setBuffer(jsg::Lock& js, jsg::JsValue newBuffer) { |
| 1066 | bufferValue = jsg::JsRef(js, newBuffer); |
| 1067 | } |
| 1068 | |
| 1069 | jsg::JsValue ServiceWorkerGlobalScope::getProcess(jsg::Lock& js) { |
| 1070 | KJ_IF_SOME(p, processValue) { |
| 1071 | return p.getHandle(js); |
| 1072 | } |
| 1073 | // Handle process module redirection based on enable_nodejs_process_v2 flag |
| 1074 | auto specifier = ([&]() -> kj::StringPtr { |
| 1075 | if (FeatureFlags::get(js).getEnableNodeJsProcessV2()) { |
| 1076 | return "node-internal:public_process"_kj; |
| 1077 | } else { |
| 1078 | return "node-internal:legacy_process"_kj; |
| 1079 | } |
| 1080 | })(); |
| 1081 | |
| 1082 | KJ_IF_SOME(module, js.resolveInternalModule(specifier)) { |
| 1083 | KJ_IF_SOME(p, processValue) { |
| 1084 | // There's a chance that resolving the module caused side-effects |
| 1085 | // that set the processValue, we let's check again. |
| 1086 | return p.getHandle(js); |
| 1087 | } |
| 1088 | // When requireReturnsDefaultExport flag is enabled, resolveInternalModule returns the |
| 1089 | // default export directly. Otherwise it returns the module namespace. |
| 1090 | auto def = module.has(js, "default"_kj) ? module.get(js, "default"_kj) : jsg::JsValue(module); |
| 1091 | JSG_REQUIRE(def.isObject(), TypeError, "Invalid node:process implementation"); |
| 1092 | processValue = jsg::JsRef(js, def); |
| 1093 | return def; |
| 1094 | } else { |
| 1095 | // If we are unable to resolve the internal module here, it likely |
| 1096 | // means that we don't actually have a module registry installed. Just |
| 1097 | // return undefined in this case. |
| 1098 | processValue = jsg::JsRef(js, js.undefined()); |
| 1099 | return js.undefined(); |
| 1100 | } |
| 1101 | } |
| 1102 | |
| 1103 | jsg::Ref<StorageManager> Navigator::getStorage(jsg::Lock& js) { |
| 1104 | return js.alloc<StorageManager>(); |
| 1105 | } |
| 1106 | |
| 1107 | bool Navigator::sendBeacon(jsg::Lock& js, kj::String url, jsg::Optional<Body::Initializer> body) { |
| 1108 | KJ_IF_SOME(context, IoContext::tryCurrent()) { |
| 1109 | auto v8Context = js.v8Context(); |
| 1110 | auto& global = |
| 1111 | jsg::extractInternalPointer<ServiceWorkerGlobalScope, true>(v8Context, v8Context->Global()); |
| 1112 | auto promise = global.fetch(js, kj::mv(url), |
| 1113 | Request::InitializerDict{ |
| 1114 | .method = kj::str("POST"), |
| 1115 | .body = kj::mv(body), |
| 1116 | }); |
| 1117 | |
| 1118 | context.addWaitUntil(context.awaitJs(js, kj::mv(promise)).ignoreResult()); |
| 1119 | return true; |
| 1120 | } |
| 1121 | |
| 1122 | // We cannot schedule a beacon to be sent outside of a request context. |
| 1123 | return false; |
| 1124 | } |
| 1125 | |
| 1126 | // ====================================================================================== |
| 1127 | |
| 1128 | Immediate::Immediate(IoContext& context, TimeoutId timeoutId) |
| 1129 | : ioContext(context.addObject(context)), |
| 1130 | timeoutId(timeoutId) {} |
| 1131 | |
| 1132 | void Immediate::dispose() { |
| 1133 | // IoPtr will throw if the IoContext is no longer valid, which is fine - accessing the Immediate |
| 1134 | // from the wrong IoContext should throw an error, just as accessing the Immediate when it's |
| 1135 | // owning IoContext is gone. |
| 1136 | ioContext->clearTimeoutImpl(timeoutId); |
| 1137 | } |
| 1138 | |
| 1139 | jsg::Ref<Immediate> ServiceWorkerGlobalScope::setImmediate(jsg::Lock& js, |
| 1140 | jsg::Function<void(jsg::Arguments<jsg::Value>)> function, |
| 1141 | jsg::Arguments<jsg::Value> args) { |
| 1142 | |
| 1143 | // This is an approximation of the Node.js setImmediate global API. |
| 1144 | // We implement it in terms of setting a 0 ms timeout. This is not |
| 1145 | // how Node.js does it so there will be some edge cases where the |
| 1146 | // timing of the callback will differ relative to the equivalent |
| 1147 | // operations in Node.js. For the vast majority of cases, users |
| 1148 | // really shouldn't be able to tell a difference. It would likely |
| 1149 | // only be somewhat pathological edge cases that could be affected |
| 1150 | // by the differences. Unfortunately, changing this later to match |
| 1151 | // Node.js would likely be a breaking change for some users that |
| 1152 | // would require a compat flag... but that's OK for now? |
| 1153 | |
| 1154 | auto& context = IoContext::current(); |
| 1155 | auto fn = [function = kj::mv(function), args = kj::mv(args), |
| 1156 | context = jsg::AsyncContextFrame::currentRef(js)](jsg::Lock& js) mutable { |
| 1157 | jsg::AsyncContextFrame::Scope scope(js, context); |
| 1158 | function(js, kj::mv(args)); |
| 1159 | }; |
| 1160 | auto timeoutId = context.setTimeoutImpl(timeoutIdGenerator, |
| 1161 | /* repeat */ false, [function = kj::mv(fn)](jsg::Lock& js) mutable { function(js); }, 0); |
| 1162 | return js.alloc<Immediate>(context, timeoutId); |
| 1163 | } |
| 1164 | |
| 1165 | void ServiceWorkerGlobalScope::clearImmediate(kj::Maybe<jsg::Ref<Immediate>> maybeImmediate) { |
| 1166 | KJ_IF_SOME(immediate, maybeImmediate) { |
| 1167 | immediate->dispose(); |
| 1168 | } |
| 1169 | } |
| 1170 | |
| 1171 | jsg::JsObject Cloudflare::getCompatibilityFlags(jsg::Lock& js) { |
| 1172 | auto flags = FeatureFlags::get(js); |
| 1173 | auto obj = js.objNoProto(); |
| 1174 | auto dynamic = capnp::toDynamic(flags); |
| 1175 | auto schema = dynamic.getSchema(); |
| 1176 | |
| 1177 | bool skipExperimental = !flags.getWorkerdExperimental(); |
| 1178 | |
| 1179 | for (auto field: schema.getFields()) { |
| 1180 | // If this is an experimental flag, we expose it only if the experimental mode |
| 1181 | // is enabled. |
| 1182 | auto annotations = field.getProto().getAnnotations(); |
| 1183 | bool skip = false; |
| 1184 | if (skipExperimental) { |
| 1185 | for (auto annotation: annotations) { |
| 1186 | if (annotation.getId() == EXPERIMENTAl_ANNOTATION_ID) { |
| 1187 | skip = true; |
| 1188 | break; |
| 1189 | } |
| 1190 | } |
| 1191 | } |
| 1192 | if (skip) continue; |
| 1193 | |
| 1194 | // Note that disable flags are not exposed. |
| 1195 | for (auto annotation: annotations) { |
| 1196 | if (annotation.getId() == COMPAT_ENABLE_FLAG_ANNOTATION_ID) { |
| 1197 | obj.setReadOnly( |
| 1198 | js, annotation.getValue().getText(), js.boolean(dynamic.get(field).as<bool>())); |
| 1199 | } |
| 1200 | } |
| 1201 | } |
| 1202 | |
| 1203 | obj.seal(js); |
| 1204 | return obj; |
| 1205 | } |
| 1206 | |
| 1207 | } // namespace workerd::api |