Skip to content
File

Blob: src/workerd/api/global-scope.c++

53.2 KB
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#include "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 
39namespace workerd::api {
40 
41namespace {
42 
43enum class NeuterReason { SENT_RESPONSE, THREW_EXCEPTION, CLIENT_DISCONNECTED };
44 
45kj::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 
62void ExecutionContext::waitUntil(kj::Promise<void> promise) {
63 IoContext::current().addWaitUntil(kj::mv(promise));
64}
65 
66void ExecutionContext::passThroughOnException() {
67 IoContext::current().setFailOpen();
68}
69 
70jsg::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 
79jsg::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 
87jsg::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 
99void 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 
111namespace {
112template <typename T>
113jsg::LenientOptional<T> mapAddRef(jsg::Lock& js, jsg::LenientOptional<T>& function) {
114 return function.map([&](T& a) { return a.addRef(js); });
115}
116} // namespace
117 
118ExportedHandler 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 
138ServiceWorkerGlobalScope::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 
149void ServiceWorkerGlobalScope::clear() {
150 removeAllHandlers();
151 unhandledRejections.clear();
152}
153 
154kj::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 
198kj::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 
420void 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 
452void 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 
495namespace {
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.
504bool 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 
512kj::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 
687jsg::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.
703kj::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 
714kj::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.
733void 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 
758void 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 
790void 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 
822void 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 
846void ServiceWorkerGlobalScope::setConnectOverride(kj::String networkAddress, ConnectFn connectFn) {
847 connectOverrides.upsert(kj::mv(networkAddress), kj::mv(connectFn));
848}
849 
850kj::Maybe<ServiceWorkerGlobalScope::ConnectFn&> ServiceWorkerGlobalScope::getConnectOverride(
851 kj::StringPtr networkAddress) {
852 return connectOverrides.find(networkAddress);
853}
854 
855jsg::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
870void 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 
876jsg::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 
890void 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 
927jsg::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 
938TimeoutId::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 
945TimeoutId::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 
961void 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 
971TimeoutId::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 
989void ServiceWorkerGlobalScope::clearInterval(jsg::Lock& js, kj::Maybe<jsg::JsNumber> timeoutId) {
990 clearTimeout(js, kj::mv(timeoutId));
991}
992 
993jsg::Ref<Crypto> ServiceWorkerGlobalScope::getCrypto(jsg::Lock& js) {
994 return js.alloc<Crypto>(js);
995}
996 
997jsg::Ref<CacheStorage> ServiceWorkerGlobalScope::getCaches(jsg::Lock& js) {
998 return js.alloc<CacheStorage>(js);
999}
1000 
1001jsg::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 
1007void 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 
1032jsg::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 
1061void ServiceWorkerGlobalScope::setProcess(jsg::Lock& js, jsg::JsValue newProcess) {
1062 processValue = jsg::JsRef(js, newProcess);
1063}
1064 
1065void ServiceWorkerGlobalScope::setBuffer(jsg::Lock& js, jsg::JsValue newBuffer) {
1066 bufferValue = jsg::JsRef(js, newBuffer);
1067}
1068 
1069jsg::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 
1103jsg::Ref<StorageManager> Navigator::getStorage(jsg::Lock& js) {
1104 return js.alloc<StorageManager>();
1105}
1106 
1107bool 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 
1128Immediate::Immediate(IoContext& context, TimeoutId timeoutId)
1129 : ioContext(context.addObject(context)),
1130 timeoutId(timeoutId) {}
1131 
1132void 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 
1139jsg::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 
1165void ServiceWorkerGlobalScope::clearImmediate(kj::Maybe<jsg::Ref<Immediate>> maybeImmediate) {
1166 KJ_IF_SOME(immediate, maybeImmediate) {
1167 immediate->dispose();
1168 }
1169}
1170 
1171jsg::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