Skip to content
File

Blob: src/workerd/io/io-context.c++

60.9 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 "io-context.h"
6 
7#include <workerd/io/io-gate.h>
8#include <workerd/io/tracer.h>
9#include <workerd/io/worker.h>
10#include <workerd/jsg/jsg.h>
11#include <workerd/jsg/setup.h>
12#include <workerd/util/autogate.h>
13#include <workerd/util/own-util.h>
14#include <workerd/util/sentry.h>
15#include <workerd/util/thread-scopes.h>
16#include <workerd/util/uncaught-exception-source.h>
17 
18#include <kj/debug.h>
19 
20#include <cmath>
21#include <map>
22 
23namespace workerd {
24 
25static thread_local IoContext* threadLocalRequest = nullptr;
26 
27SuppressIoContextScope::SuppressIoContextScope(): cached(threadLocalRequest) {
28 threadLocalRequest = nullptr;
29}
30 
31SuppressIoContextScope::~SuppressIoContextScope() noexcept(false) {
32 threadLocalRequest = cached;
33}
34 
35static const kj::EventLoopLocal<int> threadId;
36 
37static void* getThreadId() {
38 return threadId.get();
39}
40 
41class IoContext::TimeoutManagerImpl final: public TimeoutManager {
42 public:
43 class TimeoutState;
44 using Map = std::map<TimeoutId, TimeoutState>;
45 using Iterator = Map::iterator;
46 
47 TimeoutManagerImpl() = default;
48 KJ_DISALLOW_COPY_AND_MOVE(TimeoutManagerImpl);
49 
50 TimeoutId setTimeout(
51 IoContext& context, TimeoutId::Generator& generator, TimeoutParameters params) override {
52 // Verify the generator is from the correct ServiceWorkerGlobalScope. If we have been passed a
53 // different `timeoutIdGenerator`, then that means this IoContext is active at a time when
54 // JavaScript in a different V8 context is executing. This _should_ be impossible, but we're
55 // occasionally seeing timeout ID collision assertion failures in `addState()`, and one possible
56 // explanation is that an IoContext is somehow current for a different V8 context.
57 //
58 // TODO(cleanup): Find a more general way to assert that the JS API surface is being used under
59 // the correct IoContext, get rid of this function's `generator` parameter, and instead rely
60 // on the IoContext to provide the generator.
61 KJ_ASSERT(&generator == &context.getCurrentLock().getTimeoutIdGenerator(),
62 "TimeoutId Generator mismatch - using a generator from wrong ServiceWorkerGlobalScope");
63 
64 auto [id, it] = addState(generator, kj::mv(params));
65 setTimeoutImpl(context, it);
66 return id;
67 }
68 
69 void clearTimeout(IoContext&, TimeoutId id) override;
70 
71 size_t getTimeoutCount() const override {
72 return timeoutsStarted - timeoutsFinished;
73 }
74 
75 kj::Maybe<kj::Date> getNextTimeout() const override {
76 if (timeoutTimes.size() == 0) {
77 return kj::none;
78 } else {
79 return timeoutTimes.begin()->key.when;
80 }
81 }
82 
83 void cancelAll() override {
84 timerTask = nullptr;
85 timeouts.clear();
86 timeoutTimes.clear();
87 }
88 
89 private:
90 struct IdAndIterator {
91 TimeoutId id;
92 Iterator it;
93 };
94 IdAndIterator addState(TimeoutId::Generator& generator, TimeoutParameters params);
95 
96 void setTimeoutImpl(IoContext& context, Iterator it);
97 
98 // A pair of a Date and a numeric ID, used as entry in timeoutTimes set, below.
99 struct TimeoutTime {
100 kj::Date when;
101 uint tiebreaker; // Unique number, in case two timeouts target same time.
102 
103 inline bool operator<(const TimeoutTime& other) const {
104 if (when < other.when) return true;
105 if (when > other.when) return false;
106 return tiebreaker < other.tiebreaker;
107 }
108 inline bool operator==(const TimeoutTime& other) const {
109 return when == other.when && tiebreaker == other.tiebreaker;
110 }
111 };
112 
113 // Tracks registered timeouts sorted by the next time the timeout is expected to fire.
114 //
115 // The associated fulfiller should be fulfilled when the time has been reached AND all previous
116 // timeouts have completed.
117 kj::TreeMap<TimeoutTime, kj::Own<kj::PromiseFulfiller<void>>> timeoutTimes;
118 uint timeoutTimesTiebreakerCounter = 0;
119 
120 uint timeoutsStarted = 0;
121 uint timeoutsFinished = 0;
122 Map timeouts;
123 
124 // Promise that is waiting for the closest timeout, and will fulfill its fulfiller. We only ever
125 // actually wait on the next timeout in `timeoutTasks`, so that we can't fulfill timer callbacks
126 // out-of-order. This task gets replaced each time the lead timeout changes.
127 kj::Promise<void> timerTask = nullptr;
128 
129 // Must be called any time timeoutTimes.begin() changes.
130 void resetTimerTask(TimerChannel& timerChannel);
131};
132 
133class IoContext::TimeoutManagerImpl::TimeoutState {
134 public:
135 TimeoutState(TimeoutManagerImpl& manager, TimeoutParameters params);
136 ~TimeoutState();
137 
138 void trigger(Worker::Lock& lock);
139 void cancel();
140 
141 TimeoutManagerImpl& manager;
142 TimeoutParameters params;
143 
144 bool isCanceled = false;
145 bool isRunning = false;
146 
147 kj::Maybe<kj::Promise<void>> maybePromise;
148};
149 
150IoContext::IoContext(ThreadContext& thread,
151 kj::Own<const Worker> workerParam,
152 kj::Maybe<Worker::Actor&> actorParam,
153 kj::Own<LimitEnforcer> limitEnforcerParam)
154 : thread(thread),
155 worker(kj::mv(workerParam)),
156 actor(actorParam),
157 limitEnforcer(kj::mv(limitEnforcerParam)),
158 threadId(getThreadId()),
159 deleteQueue(kj::arc<DeleteQueue>()),
160 cachePutSerializer(kj::READY_NOW),
161 timeoutManager(kj::heap<TimeoutManagerImpl>()),
162 waitUntilTasks(*this),
163 tasks(*this),
164 deleteQueueSignalTask(startDeleteQueueSignalTask(this)) {
165 kj::PromiseFulfillerPair<void> paf = kj::newPromiseAndFulfiller<void>();
166 abortFulfiller = kj::mv(paf.fulfiller);
167 abortPromise = paf.promise.fork();
168 
169 // Arrange to complain if execution resource limits (CPU/memory) are exceeded.
170 auto makeLimitsPromise = [this]() {
171 auto promise = limitEnforcer->onLimitsExceeded();
172 if (isInspectorEnabled()) {
173 // Arrange to report the problem to the inspector in addition to aborting.
174 // TODO(cleanup): This is weird. Should it go somewhere else?
175 promise = (kj::coCapture([this, promise = kj::mv(promise)]() mutable -> kj::Promise<void> {
176 kj::Maybe<kj::Exception> maybeException;
177 try {
178 co_await promise;
179 } catch (...) {
180 // Just capture the exception here, we'll handle is below since we cannot have
181 // a co_await in the body of a catch clause.
182 maybeException = kj::getCaughtExceptionAsKj();
183 }
184 
185 KJ_IF_SOME(exception, maybeException) {
186 Worker::AsyncLock asyncLock = co_await worker->takeAsyncLockWithoutRequest(nullptr);
187 worker->runInLockScope(asyncLock, [&](Worker::Lock& lock) {
188 lock.logUncaughtException(
189 jsg::extractTunneledExceptionDescription(exception.getDescription()));
190 kj::throwFatalException(kj::mv(exception));
191 });
192 }
193 }))();
194 }
195 
196 return promise;
197 };
198 KJ_IF_SOME(cb, this->worker->getIsolate().getCpuLimitNearlyExceededCallback()) {
199 limitEnforcer->setCpuLimitNearlyExceededCallback(kj::mv(cb));
200 }
201 
202 // Arrange to abort when limits expire.
203 abortWhen(makeLimitsPromise());
204 
205 KJ_IF_SOME(a, actor) {
206 // Arrange to complain if the input gate is broken, which indicates a critical section failed
207 // and the actor can no longer be used.
208 abortWhen(a.getInputGate().onBroken());
209 
210 // Also complain if the output gate is broken, which indicates a critical storage failure that
211 // means we cannot continue execution. (In fact, we need to retroactively pretend that previous
212 // execution didn't happen, but that is taken care of elsewhere.)
213 abortWhen(a.getOutputGate().onBroken());
214 }
215}
216 
217IoContext::IncomingRequest::IoContext_IncomingRequest(kj::Own<IoContext> contextParam,
218 kj::Own<IoChannelFactory> ioChannelFactoryParam,
219 kj::Own<RequestObserver> metricsParam,
220 kj::Maybe<kj::Own<BaseTracer>> workerTracer,
221 kj::Maybe<tracing::InvocationSpanContext> maybeTriggerInvocationSpan)
222 : context(kj::mv(contextParam)),
223 metrics(kj::mv(metricsParam)),
224 workerTracer(kj::mv(workerTracer)),
225 ioChannelFactory(kj::mv(ioChannelFactoryParam)),
226 maybeTriggerInvocationSpan(kj::mv(maybeTriggerInvocationSpan)) {}
227 
228tracing::InvocationSpanContext& IoContext::IncomingRequest::getInvocationSpanContext() {
229 // Creating a new InvocationSpanContext can be a bit expensive since it needs to
230 // generate random IDs, so we only create it lazily when requested, which should
231 // only be when tracing is enabled and we need to record spans.
232 KJ_IF_SOME(ctx, invocationSpanContext) {
233 return ctx;
234 }
235 
236 invocationSpanContext = tracing::InvocationSpanContext::newForInvocation(
237 maybeTriggerInvocationSpan.map(
238 [](auto& trigger) -> tracing::InvocationSpanContext& { return trigger; }),
239 context->getEntropySource());
240 return KJ_ASSERT_NONNULL(invocationSpanContext);
241}
242 
243// A call to delivered() implies a promise to call drain() later (or one of the other methods
244// that sets waitedForWaitUntil). So, we can now safely add the request to
245// context->incomingRequests, which implies taking responsibility for draining on the way out.
246void IoContext::IncomingRequest::delivered(kj::SourceLocation location) {
247 KJ_REQUIRE(!wasDelivered, "delivered() can only be called once");
248 if (!context->incomingRequests.empty()) {
249 // There is already an IncomingRequest running in this context, and we're going to make it no
250 // longer current. Make sure to attribute accumulated CPU time to it.
251 auto& oldFront = context->incomingRequests.front();
252 context->limitEnforcer->reportMetrics(*oldFront.metrics);
253 
254 KJ_IF_SOME(f, oldFront.drainFulfiller) {
255 // Allow the previous current IncomingRequest to finish draining, because the new request
256 // will take over responsibility for completing any tasks that aren't done yet.
257 f.get()->fulfill();
258 }
259 }
260 
261 context->incomingRequests.addFront(*this);
262 wasDelivered = true;
263 deliveredLocation = location;
264 metrics->delivered();
265 
266 // Create the root user trace span once per request. Stale references to the span (e.g. from
267 // AsyncContextFrame storage via IoOwn, which for actors can outlive this request via the
268 // IoContext's delete queue) are safe: user-tracing SpanSubmitters hold only a
269 // BaseTracer::WeakRef, so they cannot extend tracer lifetime.
270 KJ_IF_SOME(workerTracer, workerTracer) {
271 if (util::Autogate::isEnabled(util::AutogateKey::USER_SPAN_CONTEXT_PROPAGATION)) {
272 auto& invCtx = getInvocationSpanContext();
273 rootUserTraceSpan =
274 workerTracer->makeUserRequestSpan(invCtx.getTraceId(), invCtx.getTraceFlags());
275 } else {
276 rootUserTraceSpan = workerTracer->makeUserRequestSpan(tracing::TraceId(nullptr), kj::none);
277 }
278 }
279 
280 KJ_IF_SOME(a, context->actor) {
281 // Re-synchronize the timer and top up limits for every new incoming request to an actor.
282 ioChannelFactory->getTimer().syncTime();
283 context->limitEnforcer->topUpActor();
284 
285 // Run the Actor's constructor if it hasn't been run already.
286 a.ensureConstructed(*context);
287 
288 // Record a new incoming request to actor metrics.
289 a.getMetrics().startRequest();
290 }
291}
292 
293kj::Date IoContext::IncomingRequest::now(kj::Maybe<kj::Date> nextTimeout) {
294 metrics->clockRead();
295 return ioChannelFactory->getTimer().now(kj::mv(nextTimeout));
296}
297 
298IoContext::IncomingRequest::~IoContext_IncomingRequest() noexcept(false) {
299 if (!wasDelivered) {
300 KJ_IF_SOME(w, workerTracer) {
301 w->markUnused();
302 }
303 // Request was never added to context->incomingRequests in the first place.
304 return;
305 }
306 
307 // Hack: We need to report an accurate time stamps for the STW outcome event, but the timer may
308 // not be available when the outcome event gets reported. Define the outcome event time as the
309 // time when the incoming request shuts down.
310 KJ_IF_SOME(w, workerTracer) {
311 w->recordTimestamp(now());
312 }
313 
314 if (&context->incomingRequests.front() == this) {
315 // We're the current request, make sure to consume CPU time attribution.
316 context->limitEnforcer->reportMetrics(*metrics);
317 context->lastDeliveredLocation = deliveredLocation;
318 
319 if (!waitedForWaitUntil && !context->waitUntilTasks.isEmpty()) {
320 KJ_LOG(WARNING, "failed to invoke drain() on IncomingRequest before destroying it",
321 kj::getStackTrace());
322 }
323 }
324 
325 KJ_IF_SOME(a, context->actor) {
326 a.getMetrics().endRequest();
327 }
328 context->worker->getIsolate().completedRequest();
329 metrics->jsDone();
330 
331 if (context->isShared()) {
332 // This context is not about to be destroyed when we drop it, but if it was aborted, we would
333 // prefer for it to get cleaned up promptly.
334 
335 KJ_IF_SOME(e, context->abortException) {
336 // The context was aborted. It's possible that the event ended with background work still
337 // scheduled, because `drain()` ends early on abort. We should cancel that background work
338 // now.
339 //
340 // We couldn't do this in abort() because it can be called from inside a task that could
341 // be canceled, and a self-cancellation would lead to a crash.
342 
343 if (!context->canceler.isEmpty()) {
344 context->canceler.cancel(e);
345 }
346 context->timeoutManager->cancelAll();
347 context->tasks.clear();
348 context->waitUntilTasks.clear();
349 }
350 }
351 
352 // Remove incoming request after canceling waitUntil tasks, which may have spans attached that
353 // require accessing a timer from the active request.
354 context->incomingRequests.remove(*this);
355}
356 
357InputGate::Lock IoContext::getInputLock() {
358 return KJ_ASSERT_NONNULL(currentInputLock, "no input lock available in this context")
359 .addRef(getCurrentTraceSpan());
360}
361 
362kj::Maybe<kj::Own<InputGate::CriticalSection>> IoContext::getCriticalSection() {
363 KJ_IF_SOME(l, currentInputLock) {
364 return l.getCriticalSection().map(
365 [](InputGate::CriticalSection& cs) { return kj::addRef(cs); });
366 } else {
367 return kj::none;
368 }
369}
370 
371kj::Promise<void> IoContext::waitForOutputLocks() {
372 KJ_IF_SOME(p, waitForOutputLocksIfNecessary()) {
373 return kj::mv(p);
374 } else {
375 return kj::READY_NOW;
376 }
377}
378 
379bool IoContext::hasOutputGate() {
380 return actor != kj::none;
381}
382 
383kj::Maybe<kj::Promise<void>> IoContext::waitForOutputLocksIfNecessary() {
384 return actor.map(
385 [this](Worker::Actor& actor) { return actor.getOutputGate().wait(getCurrentTraceSpan()); });
386}
387 
388kj::Maybe<IoOwn<kj::Promise<void>>> IoContext::waitForOutputLocksIfNecessaryIoOwn() {
389 return waitForOutputLocksIfNecessary().map(
390 [this](kj::Promise<void> promise) { return addObject(kj::heap(kj::mv(promise))); });
391}
392 
393bool IoContext::isOutputGateBroken() {
394 KJ_IF_SOME(a, actor) {
395 return a.getOutputGate().isBroken();
396 } else {
397 return false;
398 }
399}
400 
401bool IoContext::isInspectorEnabled() {
402 return worker->getIsolate().isInspectorEnabled();
403}
404 
405bool IoContext::hasWarningHandler() {
406 return isInspectorEnabled() || getWorkerTracer() != kj::none ||
407 ::kj::_::Debug::shouldLog(::kj::LogSeverity::INFO);
408}
409 
410void IoContext::logWarning(kj::StringPtr description) {
411 KJ_REQUIRE_NONNULL(currentLock).logWarning(description);
412}
413 
414void IoContext::logWarningOnce(kj::StringPtr description) {
415 KJ_REQUIRE_NONNULL(currentLock).logWarningOnce(description);
416}
417 
418void IoContext::logErrorOnce(kj::StringPtr description) {
419 KJ_REQUIRE_NONNULL(currentLock).logErrorOnce(description);
420}
421 
422void IoContext::logUncaughtException(kj::StringPtr description) {
423 KJ_REQUIRE_NONNULL(currentLock).logUncaughtException(description);
424}
425 
426void IoContext::logUncaughtException(
427 UncaughtExceptionSource source, const jsg::JsValue& exception, const jsg::JsMessage& message) {
428 KJ_REQUIRE_NONNULL(currentLock).logUncaughtException(source, exception, message);
429}
430 
431void IoContext::logUncaughtExceptionAsync(
432 UncaughtExceptionSource source, kj::Exception&& exception) {
433 if (getWorkerTracer() == kj::none && !worker->getIsolate().isInspectorEnabled()) {
434 // We don't need to take the isolate lock as neither inspecting nor tracing is enabled. We
435 // do still want to syslog if relevant, but we can do that without a lock.
436 if (!jsg::isTunneledException(exception.getDescription()) &&
437 !jsg::isDoNotLogException(exception.getDescription()) &&
438 // TODO(soon): Figure out why client disconnects are getting logged here if we don't
439 // ignore DISCONNECTED. If we fix that, do we still want to filter these?
440 exception.getType() != kj::Exception::Type::DISCONNECTED) {
441 LOG_EXCEPTION("jsgInternalError", exception);
442 } else {
443 KJ_LOG(INFO, "uncaught exception", exception); // Run with --verbose to see exception logs.
444 }
445 return;
446 }
447 
448 struct RunnableImpl: public Runnable {
449 UncaughtExceptionSource source;
450 kj::Exception exception;
451 
452 RunnableImpl(UncaughtExceptionSource source, kj::Exception&& exception)
453 : source(source),
454 exception(kj::mv(exception)) {}
455 void run(Worker::Lock& lock) override {
456 // TODO(soon): Add logUncaughtException to jsg::Lock.
457 lock.logUncaughtException(source, kj::mv(exception));
458 }
459 };
460 
461 // Make sure this is logged even if another exception occurs trying to log it to the devtools inspector,
462 // e.g. if `runImpl` throws before calling logUncaughtException.
463 // This is useful for tests (and in fact only affects tests, since it's logged at an INFO level).
464 KJ_ON_SCOPE_FAILURE({ KJ_LOG(INFO, "uncaught exception", source, exception); });
465 RunnableImpl runnable(source, kj::mv(exception));
466 // TODO(perf): Is it worth using an async lock here? The only case where it really matters is
467 // when a trace worker is active, but maybe they'll be more common in the future. To take an
468 // async lock here, we'll probably have to update all the call sites of this method... ick.
469 kj::Maybe<RequestObserver&> metrics;
470 if (!incomingRequests.empty()) metrics = getMetrics();
471 runImpl(
472 runnable, Worker::Lock::TakeSynchronously(metrics), kj::none, Runnable::Exceptional(true));
473}
474 
475void IoContext::abort(kj::Exception&& e) {
476 if (abortException != kj::none) {
477 return;
478 }
479 abortException = e.clone();
480 KJ_IF_SOME(a, actor) {
481 // Stop the ActorCache from flushing any scheduled write operations to prevent any unnecessary
482 // or unintentional async work
483 a.shutdownActorCache(e.clone());
484 }
485 abortFulfiller->reject(kj::mv(e));
486}
487 
488void IoContext::abortIsolate(kj::StringPtr reason) {
489 getIoChannelFactory().abortIsolate(reason);
490}
491 
492void IoContext::abortWhen(kj::Promise<void> promise) {
493 // Unlike addTask(), abortWhen() always uses `tasks`, even in actors, because we do not want
494 // these tasks to block hibernation.
495 if (abortException == kj::none) {
496 tasks.add(promise.catch_([this](kj::Exception&& e) { abort(kj::mv(e)); }));
497 }
498}
499 
500void IoContext::addTask(kj::Promise<void> promise) {
501 ++addTaskCounter;
502 
503 // In Actors, we treat all tasks as wait-until tasks, because it's perfectly legit to start a
504 // task under one request and then expect some other request to handle it later.
505 if (actor != kj::none) {
506 addWaitUntil(kj::mv(promise));
507 return;
508 }
509 
510 if (actor == kj::none) {
511 // This metric won't work correctly in actors since it's being tracked per-request, but tasks
512 // are not tied to requests in actors. So we just skip it in actors. (Actually this code path
513 // is not even executed in the actor case but I'm leaving the check in just in case that ever
514 // changes.)
515 auto& metrics = getMetrics();
516 if (metrics.getSpan().isObserved()) {
517 promise = promise.attach(metrics.addedContextTask());
518 }
519 }
520 
521 tasks.add(kj::mv(promise));
522}
523 
524void IoContext::addWaitUntil(kj::Promise<void> promise) {
525 if (actor == kj::none) {
526 // This metric won't work correctly in actors since it's being tracked per-request, but tasks
527 // are not tied to requests in actors. So we just skip it in actors.
528 auto& metrics = getMetrics();
529 if (metrics.getSpan().isObserved()) {
530 promise = promise.attach(metrics.addedWaitUntilTask());
531 }
532 }
533 
534 if (incomingRequests.empty()) {
535 DEBUG_FATAL_RELEASE_LOG(WARNING, "Adding task to IoContext with no current IncomingRequest",
536 lastDeliveredLocation, kj::getStackTrace());
537 }
538 
539 waitUntilTasks.add(kj::mv(promise));
540}
541 
542// Mark ourselves so we know that we made a best effort attempt to wait for waitUntilTasks.
543kj::Promise<void> IoContext::IncomingRequest::drain() {
544 waitedForWaitUntil = true;
545 
546 if (&context->incomingRequests.front() != this) {
547 // A newer request was received, so draining isn't our job.
548 return kj::READY_NOW;
549 }
550 
551 kj::Promise<void> timeoutPromise = nullptr;
552 KJ_IF_SOME(a, context->actor) {
553 // For actors, all promises are canceled on actor shutdown, not on a fixed timeout,
554 // because work doesn't necessarily happen on a per-request basis in actors and we don't want
555 // work being unexpectedly canceled based on which request initiated it.
556 timeoutPromise = a.onShutdown();
557 
558 // Also arrange to cancel the drain if a new request arrives, since it will take over
559 // responsibility for background tasks.
560 auto drainPaf = kj::newPromiseAndFulfiller<void>();
561 drainFulfiller = kj::mv(drainPaf.fulfiller);
562 timeoutPromise = timeoutPromise.exclusiveJoin(kj::mv(drainPaf.promise));
563 } else {
564 // For non-actor requests, apply the configured soft timeout, typically 30 seconds.
565 auto timeoutLogPromise = [this]() -> kj::Promise<void> {
566 return context->run([this](Worker::Lock&) {
567 context->logWarning(
568 "waitUntil() tasks did not complete within the allowed time after invocation end and have been cancelled. "
569 "See: https://developers.cloudflare.com/workers/runtime-apis/context/#waituntil");
570 });
571 };
572 timeoutPromise = context->limitEnforcer->limitDrain().then(kj::mv(timeoutLogPromise));
573 }
574 return context->waitUntilTasks.onEmpty()
575 .exclusiveJoin(kj::mv(timeoutPromise))
576 .exclusiveJoin(context->onAbort().catch_([](kj::Exception&&) {}));
577}
578 
579kj::Promise<EventOutcome> IoContext::IncomingRequest::finishScheduled() {
580 // TODO(someday): In principle we should be able to support delivering the "scheduled" event type
581 // to an actor, and this may be important if we open up the whole of WorkerInterface to be
582 // callable from any stub. However, the logic around async tasks would have to be different. We
583 // cannot assume that just because an async task fails while the scheduled event is running,
584 // that the scheduled event itself failed -- the failure could have been a task initiated by
585 // an unrelated concurrent event.
586 KJ_ASSERT(context->actor == kj::none,
587 "this code isn't designed to allow scheduled events to be delivered to actors");
588 
589 // Mark ourselves so we know that we made a best effort attempt to wait for waitUntilTasks.
590 KJ_ASSERT(context->incomingRequests.size() == 1);
591 context->incomingRequests.front().waitedForWaitUntil = true;
592 
593 auto timeoutPromise = context->limitEnforcer->limitScheduled().then([] {
594 // TODO(soon): The limit being hit here is a wall time limit. Can we report an
595 // "exceededWallTime" outcome instead?
596 return EventOutcome::EXCEEDED_CPU;
597 });
598 return context->waitUntilTasks.onEmpty()
599 .then([]() { return EventOutcome::OK; })
600 .exclusiveJoin(kj::mv(timeoutPromise))
601 .exclusiveJoin(context->onAbort().then([] {
602 // abortFulfiller should only ever be rejected instead of being fulfilled, return an
603 // internalError outcome if it does happen
604 return EventOutcome::INTERNAL_ERROR;
605 }, [](kj::Exception&& e) { return RequestObserver::outcomeFromException(e); }));
606}
607 
608class IoContext::PendingEvent: public kj::Refcounted {
609 public:
610 explicit PendingEvent(IoContext& context): maybeContext(context) {}
611 ~PendingEvent() noexcept(false);
612 KJ_DISALLOW_COPY_AND_MOVE(PendingEvent);
613 
614 kj::Maybe<IoContext&> maybeContext;
615};
616 
617IoContext::~IoContext() noexcept(false) {
618 if (!canceler.isEmpty()) {
619 KJ_IF_SOME(e, abortException) {
620 // Assume the abort exception is why we are canceling.
621 canceler.cancel(e);
622 } else {
623 canceler.cancel(JSG_KJ_EXCEPTION(
624 FAILED, Error, "The execution context responding to this call was canceled."));
625 }
626 }
627 
628 // Detach the PendingEvent if it still exists.
629 KJ_IF_SOME(pe, pendingEvent) {
630 pe.maybeContext = kj::none;
631 }
632 
633 // Kill the sentinel so that no weak references can refer to this IoContext anymore.
634 selfRef->invalidate();
635}
636 
637IoContext::PendingEvent::~PendingEvent() noexcept(false) {
638 IoContext& context = KJ_UNWRAP_OR(maybeContext, {
639 // IoContext must have been destroyed before the PendingEvent was.
640 return;
641 });
642 
643 context.pendingEvent = kj::none;
644 
645 // We can't abort just yet. We need to run the event loop to see if any queued
646 // events come back into JavaScript. If registerPendingEvent() is called in the meantime, this
647 // will be canceled.
648 context.abortFromHangTask = Worker::AsyncLock::whenThreadIdle()
649 .then([&context = context]() noexcept {
650 // We have nothing left to do and no PendingEvent has been registered. Abort now.
651 return context.worker->takeAsyncLock(context.getMetrics())
652 .then([&context](Worker::AsyncLock asyncLock) { context.abortFromHang(asyncLock); });
653 }).eagerlyEvaluate(nullptr);
654}
655 
656kj::Own<void> IoContext::registerPendingEvent() {
657 if (actor != kj::none) {
658 // Actors don't use the pending event system, because different requests to the same Actor are
659 // explicitly allowed to resolve each other's promises.
660 return {};
661 }
662 
663 KJ_IF_SOME(pe, pendingEvent) {
664 return kj::addRef(pe);
665 } else {
666 KJ_IF_SOME(e, abortException) {
667 kj::throwFatalException(e.clone());
668 }
669 
670 // Cancel any already-scheduled finalization.
671 abortFromHangTask = kj::none;
672 
673 auto result = kj::refcounted<PendingEvent>(*this);
674 pendingEvent = *result;
675 return result;
676 }
677}
678 
679IoContext::TimeoutManagerImpl::TimeoutState::TimeoutState(
680 TimeoutManagerImpl& manager, TimeoutParameters params)
681 : manager(manager),
682 params(kj::mv(params)) {
683 ++manager.timeoutsStarted;
684}
685 
686IoContext::TimeoutManagerImpl::TimeoutState::~TimeoutState() {
687 KJ_ASSERT(!isRunning);
688 if (!isCanceled) {
689 ++manager.timeoutsFinished;
690 }
691}
692 
693void IoContext::TimeoutManagerImpl::TimeoutState::trigger(Worker::Lock& lock) {
694 isRunning = true;
695 auto cleanupGuard = kj::defer([&] { isRunning = false; });
696 
697 // Now it's safe to call the user's callback.
698 KJ_IF_SOME(function, params.function) {
699 (function)(lock);
700 }
701}
702 
703void IoContext::TimeoutManagerImpl::TimeoutState::cancel() {
704 if (isCanceled) {
705 return;
706 }
707 
708 auto wasCanceled = isCanceled;
709 isCanceled = true;
710 
711 if (!isRunning && !wasCanceled) {
712 params.function = kj::none;
713 maybePromise = kj::none;
714 }
715 
716 ++manager.timeoutsFinished;
717}
718 
719auto IoContext::TimeoutManagerImpl::addState(
720 TimeoutId::Generator& generator, TimeoutParameters params) -> IdAndIterator {
721 JSG_REQUIRE(getTimeoutCount() < MAX_TIMEOUTS, DOMQuotaExceededError,
722 "You have exceeded the number of active timeouts you may set.",
723 " max active timeouts: ", MAX_TIMEOUTS, ", current active timeouts: ", getTimeoutCount(),
724 ", finished timeouts: ", timeoutsFinished);
725 
726 auto id = generator.getNext();
727 auto [it, wasEmplaced] = timeouts.try_emplace(id, *this, kj::mv(params));
728 if (!wasEmplaced) {
729 // We shouldn't have reached here because the `TimeoutId::Generator` throws if it reaches
730 // Number.MAX_SAFE_INTEGER, much less wraps around the uint64_t number space. Let's throw with
731 // as many details as possible.
732 auto& state = it->second;
733 auto delay = state.params.msDelay;
734 auto repeat = state.params.repeat;
735 KJ_FAIL_ASSERT("Saw a timeout id collision", getTimeoutCount(), timeoutsStarted, id.toNumber(),
736 delay, repeat);
737 }
738 
739 return {id, it};
740}
741 
742void IoContext::TimeoutManagerImpl::setTimeoutImpl(IoContext& context, Iterator it) {
743 auto& state = it->second;
744 
745 auto stateGuard = kj::defer([&]() {
746 if (state.maybePromise == kj::none) {
747 // Something threw, erase the state.
748 timeouts.erase(it);
749 }
750 });
751 
752 auto paf = kj::newPromiseAndFulfiller<void>();
753 
754 // Schedule relative to Date.now() so the delay appears exact to the application.
755 auto when = context.now() + state.params.msDelay * kj::MILLISECONDS;
756 // TODO(cleanup): The manual use of run() here (including carrying over the critical section) is
757 // kind of ugly, but using awaitIo() doesn't work here because we need the ability to cancel
758 // the timer, so we don't want to addTask() it, which awaitIo() does implicitly.
759 auto promise =
760 paf.promise.then([this, &context, it, cs = context.getCriticalSection()]() mutable {
761 return context.run([this, &context, it](Worker::Lock& lock) mutable {
762 auto& state = it->second;
763 
764 auto stateGuard = kj::defer([&] {
765 if (state.maybePromise == kj::none) {
766 // At the end of this block, there was no new timeout, so we should remove the state.
767 // Note that this can happen from cancelTimeout or a non-repeating timeout.
768 timeouts.erase(it);
769 }
770 });
771 
772 if (state.isCanceled) {
773 // We've been canceled before running. Nothing more to do.
774 KJ_ASSERT(state.maybePromise == kj::none);
775 return;
776 }
777 
778 KJ_IF_SOME(promise, state.maybePromise) {
779 // We could KJ_ASSERT_NONNULL(iter->second) instead if we are sure clearTimeout() couldn't
780 // race us. However, I'm not sure about that.
781 
782 // First, move our timeout promise to the task set so it's safe to call clearInterval()
783 // inside the user's callback. We don't yet null out the Maybe<Promise>, because we need to
784 // be able to detect whether the user does call clearInterval(). We leave the actual map
785 // entry in place because this aids in reporting cross-request-context timeout cancellation
786 // errors to the user.
787 context.addTask(kj::mv(promise));
788 
789 // Because Promise has an underspecified move ctor, we need to explicitly nullify the Maybe
790 // to indicate that we've consumed the promise.
791 state.maybePromise = kj::none;
792 
793 // The user's callback might throw, but we need to at least attempt to reschedule interval
794 // callbacks even if they throw. This deferred action takes care of that. Note that we don't
795 // run the user's callback directly in this->run(), because that function throws a fatal
796 // exception if a JS exception is thrown, which complicates our logic here.
797 //
798 // TODO(perf): If we can guarantee that `timeout->second = nullptr` will never throw, it
799 // might be worthwhile having an early-out path for non-interval timeouts.
800 kj::UnwindDetector unwindDetector;
801 KJ_DEFER(unwindDetector.catchExceptionsIfUnwinding([&] {
802 if (state.isCanceled) {
803 // The user's callback has called clearInterval(), nothing more to do.
804 KJ_ASSERT(state.maybePromise == kj::none);
805 return;
806 }
807 
808 // If this is an interval task and the script has CPU time left, reschedule the task;
809 // otherwise leave the dead map entry in place.
810 if (state.params.repeat && context.limitEnforcer->getLimitsExceeded() == kj::none) {
811 setTimeoutImpl(context, it);
812 }
813 }););
814 
815 state.trigger(lock);
816 }
817 }, kj::mv(cs));
818 }, [](kj::Exception&&) {});
819 
820 promise = promise.attach(context.registerPendingEvent());
821 
822 // Add an entry to the timeoutTimes map, to track when the nearest timeout is. Arrange for it
823 // to be removed when the promise completes.
824 TimeoutTime timeoutTimesKey{when, timeoutTimesTiebreakerCounter++};
825 timeoutTimes.insert(timeoutTimesKey, kj::mv(paf.fulfiller));
826 auto deferredTimeoutTimeRemoval = kj::defer([this, &context, timeoutTimesKey]() {
827 // If the promise is being destroyed due to IoContext teardown then IoChannelFactory may
828 // no longer be available, but we can just skip starting a new timer in that case as it'd be
829 // canceled anyway. Similarly we should skip rescheduling if the context has been aborted since
830 // there's no way the events can run anyway (and we'll cause trouble if `cancelAll()` is being
831 // called in ~IoContext_IncomingRequest).
832 if (context.selfRef->isValid() && context.abortException == kj::none) {
833 bool isNext = timeoutTimes.begin()->key == timeoutTimesKey;
834 timeoutTimes.erase(timeoutTimesKey);
835 if (isNext) resetTimerTask(context.getIoChannelFactory().getTimer());
836 }
837 });
838 
839 if (timeoutTimes.begin()->key == timeoutTimesKey) {
840 resetTimerTask(context.getIoChannelFactory().getTimer());
841 }
842 promise = promise.attach(kj::mv(deferredTimeoutTimeRemoval));
843 
844 if (context.actor != kj::none) {
845 // Add a wait-until task which resolves when this timer completes. This ensures that
846 // `IncomingRequest::drain()` waits until all timers finish.
847 auto paf = kj::newPromiseAndFulfiller<void>();
848 promise = promise.attach(
849 kj::defer([fulfiller = kj::mv(paf.fulfiller)]() mutable { fulfiller->fulfill(); }));
850 context.addWaitUntil(kj::mv(paf.promise));
851 }
852 
853 state.maybePromise = promise.eagerlyEvaluate(nullptr);
854}
855 
856void IoContext::TimeoutManagerImpl::resetTimerTask(TimerChannel& timerChannel) {
857 if (timeoutTimes.size() == 0) {
858 // Not waiting for any timer, clear the existing timer task.
859 timerTask = nullptr;
860 } else {
861 // Wait for the first timer.
862 auto& entry = *timeoutTimes.begin();
863 timerTask = timerChannel.atTime(entry.key.when)
864 .then([this, key = entry.key]() {
865 auto& newEntry = *timeoutTimes.begin();
866 KJ_ASSERT(newEntry.key == key,
867 "front of timeoutTimes changed without calling resetTimerTask(), we probably missed "
868 "a timeout!");
869 newEntry.value->fulfill();
870 }).eagerlyEvaluate([](kj::Exception&& e) { KJ_LOG(ERROR, e); });
871 }
872}
873 
874void IoContext::TimeoutManagerImpl::clearTimeout(IoContext& context, TimeoutId timeoutId) {
875 auto timeout = timeouts.find(timeoutId);
876 if (timeout == timeouts.end()) {
877 // We can't find this timeout, thus we act as if it was already canceled.
878 return;
879 }
880 
881 // Cancel the timeout.
882 timeout->second.cancel();
883}
884 
885TimeoutId IoContext::setTimeoutImpl(
886 TimeoutId::Generator& generator, bool repeat, jsg::Function<void()> function, double msDelay) {
887 static constexpr int64_t max = 3153600000000; // Milliseconds in 100 years
888 // Clamp the range on timers to [0, 3153600000000] (inclusive). The specs
889 // do not indicate a clear maximum range for setTimeout/setInterval so the
890 // limit here is fairly arbitrary. 100 years max should be plenty safe.
891 int64_t delay = msDelay <= 0 || std::isnan(msDelay) ? 0
892 : msDelay >= static_cast<double>(max) ? max
893 : static_cast<int64_t>(msDelay);
894 auto params = TimeoutManager::TimeoutParameters(repeat, delay, kj::mv(function));
895 return timeoutManager->setTimeout(*this, generator, kj::mv(params));
896}
897 
898void IoContext::clearTimeoutImpl(TimeoutId id) {
899 timeoutManager->clearTimeout(*this, id);
900}
901 
902size_t IoContext::getTimeoutCount() {
903 return timeoutManager->getTimeoutCount();
904}
905 
906kj::Date IoContext::now(IncomingRequest& incomingRequest) {
907 if (getWorker().getScript().getIsolate().getApi().getFeatureFlags().getPreciseTimers()) {
908 auto now = kj::systemPreciseCalendarClock().now();
909 // Round to 3ms granularity
910 int64_t ms = (now - kj::UNIX_EPOCH) / kj::MILLISECONDS;
911 int64_t roundedMs = (ms / 3) * 3;
912 return kj::UNIX_EPOCH + roundedMs * kj::MILLISECONDS;
913 }
914 
915 // Let TimerChannel decide whether to clamp to the next timeout time. This is how Spectre
916 // mitigations ensure Date.now() inside a callback returns exactly the scheduled time.
917 return incomingRequest.now(timeoutManager->getNextTimeout());
918}
919 
920kj::Date IoContext::now() {
921 return now(getCurrentIncomingRequest());
922}
923 
924kj::Rc<ExternalPusherImpl> IoContext::getExternalPusher() {
925 KJ_IF_SOME(ep, externalPusher) {
926 return ep.addRef();
927 } else {
928 return externalPusher.emplace(kj::rc<ExternalPusherImpl>(getByteStreamFactory())).addRef();
929 }
930}
931 
932kj::Own<WorkerInterface> IoContext::getSubrequestNoChecks(
933 kj::FunctionParam<kj::Own<WorkerInterface>(TraceContext&, IoChannelFactory&)> func,
934 SubrequestOptions options) {
935 TraceContext tracing;
936 KJ_IF_SOME(n, options.operationName) {
937 tracing = makeUserTraceSpan(n.clone());
938 }
939 
940 kj::Own<WorkerInterface> ret;
941 KJ_IF_SOME(existing, options.existingTraceContext) {
942 ret = func(existing, getIoChannelFactory());
943 } else {
944 ret = func(tracing, getIoChannelFactory());
945 }
946 
947 if (options.wrapMetrics) {
948 auto& metrics = getMetrics();
949 ret = metrics.wrapSubrequestClient(kj::mv(ret));
950 ret = worker->getIsolate().wrapSubrequestClient(
951 kj::mv(ret), getHeaderIds().contentEncoding, metrics);
952 }
953 
954 if (tracing.isObserved()) {
955 auto ioOwnedSpan = addObject(kj::heap(kj::mv(tracing)));
956 ret = ret.attach(kj::mv(ioOwnedSpan));
957 }
958 
959 // Subrequests use a lot of unaccounted C++ memory, so we adjust V8's external memory counter to
960 // pressure the GC and protect against OOMs. We apply this adjustment to ALL subrequests (not
961 // just fetch). We only apply this when the JS lock is held (i.e., when JS code initiated the
962 // subrequest); infrastructure paths that bypass JS don't need it.
963 KJ_IF_SOME(lock, currentLock) {
964 jsg::Lock& js = lock;
965 ret = ret.attach(js.getExternalMemoryAdjustment(8 * 1024));
966 }
967 
968 return kj::mv(ret);
969}
970 
971kj::Own<WorkerInterface> IoContext::getSubrequest(
972 kj::FunctionParam<kj::Own<WorkerInterface>(TraceContext&, IoChannelFactory&)> func,
973 SubrequestOptions options) {
974 limitEnforcer->newSubrequest(options.inHouse);
975 return getSubrequestNoChecks(kj::mv(func), kj::mv(options));
976}
977 
978kj::Own<WorkerInterface> IoContext::getSubrequestChannel(
979 uint channel, bool isInHouse, kj::Maybe<kj::String> cfBlobJson, kj::ConstString operationName) {
980 return getSubrequest(
981 [&](TraceContext& tracing, IoChannelFactory& channelFactory) {
982 return getSubrequestChannelImpl(
983 channel, isInHouse, kj::mv(cfBlobJson), tracing, channelFactory);
984 },
985 SubrequestOptions{
986 .inHouse = isInHouse,
987 .wrapMetrics = !isInHouse,
988 .operationName = kj::mv(operationName),
989 });
990}
991 
992kj::Own<WorkerInterface> IoContext::getSubrequestChannel(
993 uint channel, bool isInHouse, kj::Maybe<kj::String> cfBlobJson, TraceContext& traceContext) {
994 return getSubrequest(
995 [&](TraceContext& tracing, IoChannelFactory& channelFactory) {
996 return getSubrequestChannelImpl(
997 channel, isInHouse, kj::mv(cfBlobJson), tracing, channelFactory);
998 },
999 SubrequestOptions{
1000 .inHouse = isInHouse,
1001 .wrapMetrics = !isInHouse,
1002 .existingTraceContext = traceContext,
1003 });
1004}
1005 
1006kj::Own<WorkerInterface> IoContext::getSubrequestChannelNoChecks(uint channel,
1007 bool isInHouse,
1008 kj::Maybe<kj::String> cfBlobJson,
1009 kj::Maybe<kj::ConstString> operationName) {
1010 return getSubrequestNoChecks(
1011 [&](TraceContext& tracing, IoChannelFactory& channelFactory) {
1012 return getSubrequestChannelImpl(
1013 channel, isInHouse, kj::mv(cfBlobJson), tracing, channelFactory);
1014 },
1015 SubrequestOptions{
1016 .inHouse = isInHouse,
1017 .wrapMetrics = !isInHouse,
1018 .operationName = kj::mv(operationName),
1019 });
1020}
1021 
1022kj::Own<WorkerInterface> IoContext::getSubrequestChannelImpl(uint channel,
1023 bool isInHouse,
1024 kj::Maybe<kj::String> cfBlobJson,
1025 TraceContext& tracing,
1026 IoChannelFactory& channelFactory) {
1027 IoChannelFactory::SubrequestMetadata metadata{
1028 .cfBlobJson = kj::mv(cfBlobJson),
1029 .parentSpan = tracing.getInternalSpanParent(),
1030 .userSpanParent = tracing.getUserSpanParent(),
1031 .featureFlagsForFl = mapCopyString(worker->getIsolate().getFeatureFlagsForFl()),
1032 };
1033 
1034 auto client = channelFactory.startSubrequest(channel, kj::mv(metadata));
1035 
1036 return client;
1037}
1038 
1039kj::Own<kj::HttpClient> IoContext::getHttpClient(
1040 uint channel, bool isInHouse, kj::Maybe<kj::String> cfBlobJson, kj::ConstString operationName) {
1041 return asHttpClient(
1042 getSubrequestChannel(channel, isInHouse, kj::mv(cfBlobJson), kj::mv(operationName)));
1043}
1044 
1045kj::Own<kj::HttpClient> IoContext::getHttpClient(
1046 uint channel, bool isInHouse, kj::Maybe<kj::String> cfBlobJson, TraceContext& traceContext) {
1047 return asHttpClient(getSubrequestChannel(channel, isInHouse, kj::mv(cfBlobJson), traceContext));
1048}
1049 
1050kj::Own<CacheClient> IoContext::getCacheClient() {
1051 // TODO(someday): Should Cache API requests be considered in-house? They are already not counted
1052 // as subrequests in metrics and logs (like in-house requests aren't), but historically the
1053 // subrequest limit still applied. Since I can't currently think of a use case for more than 50
1054 // cache API requests per request, I'm leaving it as-is for now.
1055 limitEnforcer->newSubrequest(false);
1056 auto ret = getIoChannelFactory().getCache();
1057 
1058 // Apply external memory adjustment for Cache API subrequests (same as other subrequests in
1059 // getSubrequestNoChecks).
1060 KJ_IF_SOME(lock, currentLock) {
1061 jsg::Lock& js = lock;
1062 ret = ret.attach(js.getExternalMemoryAdjustment(8 * 1024));
1063 }
1064 
1065 return kj::mv(ret);
1066}
1067 
1068jsg::AsyncContextFrame::StorageScope IoContext::makeAsyncTraceScope(
1069 Worker::Lock& lock, kj::Maybe<SpanParent> spanParentOverride) {
1070 static const SpanParent dummySpanParent = nullptr;
1071 
1072 jsg::Lock& js = lock;
1073 kj::Own<SpanParent> spanParent;
1074 KJ_IF_SOME(spo, kj::mv(spanParentOverride)) {
1075 spanParent = kj::heap(kj::mv(spo));
1076 } else {
1077 // TODO(cleanup): Can we also elide the other memory allocations for the (unused) storage
1078 // scope if tracing is disabled?
1079 SpanParent metricsSpan = getMetrics().getSpan();
1080 if (!metricsSpan.isObserved()) {
1081 // const_cast is ok: There's no state that could be changed in a non-observed span parent.
1082 spanParent = kj::Own<SpanParent>(
1083 &const_cast<SpanParent&>(dummySpanParent), kj::NullDisposer::instance);
1084 } else {
1085 spanParent = kj::heap(kj::mv(metricsSpan));
1086 }
1087 }
1088 auto ioOwnSpanParent = IoContext::current().addObject(kj::mv(spanParent));
1089 auto spanHandle = jsg::wrapOpaque(js.v8Context(), kj::mv(ioOwnSpanParent));
1090 return jsg::AsyncContextFrame::StorageScope(
1091 js, lock.getTraceAsyncContextKey(), js.v8Ref(spanHandle));
1092}
1093 
1094jsg::AsyncContextFrame::StorageScope IoContext::makeUserAsyncTraceScope(
1095 Worker::Lock& lock, kj::Maybe<SpanParent> userSpanOverride) {
1096 jsg::Lock& js = lock;
1097 kj::Own<SpanParent> userSpan;
1098 KJ_IF_SOME(sp, kj::mv(userSpanOverride)) {
1099 userSpan = kj::heap(kj::mv(sp));
1100 } else {
1101 userSpan = kj::heap(getRootUserTraceSpan());
1102 }
1103 auto ioOwnSpan = IoContext::current().addObject(kj::mv(userSpan));
1104 auto spanHandle = jsg::wrapOpaque(js.v8Context(), kj::mv(ioOwnSpan));
1105 return jsg::AsyncContextFrame::StorageScope(
1106 js, lock.getUserTraceAsyncContextKey(), js.v8Ref(spanHandle));
1107}
1108 
1109SpanParent IoContext::getCurrentTraceSpan() {
1110 // If called while lock is held, try to use the trace info stored in the async context.
1111 KJ_IF_SOME(lock, currentLock) {
1112 KJ_IF_SOME(frame, jsg::AsyncContextFrame::current(lock)) {
1113 KJ_IF_SOME(value, frame.get(lock.getTraceAsyncContextKey())) {
1114 auto handle = value.getHandle(lock);
1115 jsg::Lock& js = lock;
1116 auto& spanParent = jsg::unwrapOpaqueRef<IoOwn<SpanParent>>(js.v8Isolate, handle);
1117 return spanParent->addRef();
1118 }
1119 }
1120 }
1121 
1122 // If async context is unavailable (unset, or JS lock is not held), fall back to heuristic of
1123 // using the trace info from the most recent active request.
1124 return getMetrics().getSpan();
1125}
1126 
1127SpanParent IoContext::getCurrentUserTraceSpan() {
1128 // Skip the AsyncContextFrame probe when user tracing isn't wired up: an unobserved
1129 // root means enterSpan can't have pushed anything (see Tracing::enterSpan).
1130 if (incomingRequests.empty()) {
1131 return SpanParent(nullptr);
1132 }
1133 SpanParent root = getCurrentIncomingRequest().getRootUserTraceSpan();
1134 if (!root.isObserved()) {
1135 return kj::mv(root);
1136 }
1137 
1138 // If called while lock is held, try to use the trace info stored in the async context.
1139 KJ_IF_SOME(lock, currentLock) {
1140 KJ_IF_SOME(frame, jsg::AsyncContextFrame::current(lock)) {
1141 KJ_IF_SOME(value, frame.get(lock.getUserTraceAsyncContextKey())) {
1142 auto handle = value.getHandle(lock);
1143 jsg::Lock& js = lock;
1144 auto& userSpan = jsg::unwrapOpaqueRef<IoOwn<SpanParent>>(js.v8Isolate, handle);
1145 return userSpan->addRef();
1146 }
1147 }
1148 }
1149 return kj::mv(root);
1150}
1151 
1152SpanBuilder IoContext::makeTraceSpan(kj::ConstString operationName) {
1153 return getCurrentTraceSpan().newChild(kj::mv(operationName));
1154}
1155 
1156TraceContext IoContext::makeUserTraceSpan(kj::ConstString operationName) {
1157 auto span = makeTraceSpan(operationName.clone());
1158 auto userSpan = getCurrentUserTraceSpan().newChild(kj::mv(operationName));
1159 return TraceContext(kj::mv(span), kj::mv(userSpan));
1160}
1161 
1162void IoContext::taskFailed(kj::Exception&& exception) {
1163 if (waitUntilStatusValue == EventOutcome::OK) {
1164 KJ_IF_SOME(status, limitEnforcer->getLimitsExceeded()) {
1165 waitUntilStatusValue = status;
1166 } else {
1167 waitUntilStatusValue = RequestObserver::outcomeFromException(exception);
1168 }
1169 }
1170 
1171 // If `taskFailed()` throws the whole event loop blows up... let's be careful not to let that
1172 // happen.
1173 KJ_IF_SOME(e, kj::runCatchingExceptions([&]() {
1174 logUncaughtExceptionAsync(UncaughtExceptionSource::ASYNC_TASK, kj::mv(exception));
1175 })) {
1176 KJ_LOG(ERROR, "logUncaughtExceptionAsync() threw an exception?", e);
1177 }
1178}
1179 
1180void IoContext::requireCurrent() {
1181 KJ_REQUIRE(threadLocalRequest == this, "request is not current in this thread");
1182}
1183 
1184void IoContext::checkFarGet(const DeleteQueue& expectedQueue, const std::type_info& type) {
1185 requireCurrent();
1186 
1187 if (&expectedQueue == deleteQueue.queue.get()) {
1188 // same request or same actor, success
1189 } else {
1190 throwNotCurrentJsError(type);
1191 }
1192}
1193 
1194Worker::Actor& IoContext::getActorOrThrow() {
1195 return KJ_ASSERT_NONNULL(actor, "not an actor request");
1196}
1197 
1198void IoContext::runInContextScope(Worker::LockType lockType,
1199 kj::Maybe<InputGate::Lock> inputLock,
1200 kj::Function<void(Worker::Lock&)> func) {
1201 // The previously-current context, before we entered this scope. We have to allow opening
1202 // multiple nested scopes especially to support destructors: destroying objects related to a
1203 // subrequest in one worker could transitively destroy resources belonging to the next worker in
1204 // the pipeline. We can't delay destruction to a future turn of the event loop because it's
1205 // common for child objects to contain pointers back to stuff owned by the parent that could
1206 // then be dangling.
1207 KJ_REQUIRE(threadId == getThreadId(), "IoContext cannot switch threads");
1208 SuppressIoContextScope previousRequest;
1209 threadLocalRequest = this;
1210 
1211 worker->runInLockScope(lockType, [&](Worker::Lock& lock) {
1212 KJ_REQUIRE(currentInputLock == kj::none);
1213 KJ_REQUIRE(currentLock == kj::none);
1214 KJ_DEFER(currentLock = kj::none; currentInputLock = kj::none);
1215 currentInputLock = kj::mv(inputLock);
1216 currentLock = lock;
1217 
1218 JSG_WITHIN_CONTEXT_SCOPE(lock, lock.getContext(), [&](jsg::Lock& js) {
1219 v8::Isolate::PromiseContextScope promiseContextScope(
1220 lock.getIsolate(), getPromiseContextTag(lock));
1221 
1222 {
1223 // Handle any pending deletions that arrived while the worker was processing a different
1224 // request.
1225 auto l = deleteQueue.queue->crossThreadDeleteQueue.lockExclusive();
1226 auto& state = KJ_ASSERT_NONNULL(*l);
1227 for (auto& object: state.queue) {
1228 OwnedObjectList::unlink(*object);
1229 }
1230 state.queue.clear();
1231 }
1232 
1233 func(lock);
1234 });
1235 });
1236}
1237 
1238void IoContext::runImpl(Runnable& runnable,
1239 Worker::LockType lockType,
1240 kj::Maybe<InputGate::Lock> inputLock,
1241 Runnable::Exceptional exceptional) {
1242 KJ_IF_SOME(l, inputLock) {
1243 KJ_REQUIRE(l.isFor(KJ_ASSERT_NONNULL(actor).getInputGate()));
1244 }
1245 
1246 getIoChannelFactory().getTimer().syncTime();
1247 
1248 runInContextScope(lockType, kj::mv(inputLock), [&](Worker::Lock& workerLock) {
1249 kj::Own<void> event;
1250 if (!exceptional) {
1251 workerLock.requireNoPermanentException();
1252 // Prevent prematurely detecting a hang while we're still executing JavaScript.
1253 // TODO(cleanup): Is this actually still needed or is this vestigial? Seems like it should
1254 // not be necessary.
1255 event = registerPendingEvent();
1256 }
1257 
1258 auto limiterScope = limitEnforcer->enterJs(workerLock, *this);
1259 
1260 bool gotTermination = false;
1261 
1262 KJ_DEFER({
1263 // Always clear out all pending V8 events before leaving the scope. This ensures that
1264 // there's never any unfinished work waiting to run when we return to the event loop.
1265 //
1266 // Alternatively, we could use kj::evalLater() to queue a callback which runs the microtasks.
1267 // This would perhaps prevent a microtask loop from blocking incoming I/O events. However,
1268 // in practice this seems like a dubious scenario. A script that does while(1) will always
1269 // block I/O, so why should a script in a promise loop not? If scripts want to use 100% of
1270 // CPU but also receive I/O as it arrives, we should offer some API to explicitly request
1271 // polling for I/O.
1272 jsg::Lock& js = workerLock;
1273 
1274 if (gotTermination) {
1275 // We already consumed the termination pseudo-exception, so if we call RunMicrotasks() now,
1276 // they will run with no limit. But if we call terminateNextExecution() again now, it will
1277 // conveniently cause RunMicrotasks() to terminate _right after_ dequeuing the contents of
1278 // the task queue, which is perfect, because it effectively cancels them all.
1279 js.terminateNextExecution();
1280 }
1281 
1282 // Run microtask checkpoint with an active IoContext
1283 {
1284 // Running the microtask queue can itself trigger a pending exception in the isolate.
1285 v8::TryCatch tryCatch(workerLock.getIsolate());
1286 
1287 js.runMicrotasks();
1288 
1289 if (tryCatch.HasCaught()) {
1290 // It really shouldn't be possible for microtasks to throw regular exceptions.
1291 // so if we got here it should be a terminal condition.
1292 KJ_ASSERT(tryCatch.HasTerminated());
1293 // If we do not reset here we end up with a dangling exception in the isolate that
1294 // leads to an assert in v8 when the Lock is destroyed.
1295 tryCatch.Reset();
1296 // Ensure we don't pump the message loop in this case
1297 gotTermination = true;
1298 }
1299 }
1300 
1301 // With --gc-stress, force a full GC after microtasks run. This catches objects that
1302 // became unreachable during JS execution / microtask processing (e.g., a
1303 // ReadableStreamDefaultReader with no JS variable binding whose closed promise is
1304 // still pending).
1305 if (isGcStressModeForTest()) {
1306 workerLock.getIsolate()->RequestGarbageCollectionForTesting(
1307 v8::Isolate::kFullGarbageCollection);
1308 }
1309 
1310 // Run FinalizationRegistry cleanup tasks without an IoContext
1311 {
1312 SuppressIoContextScope noIoCtxt;
1313 while (!gotTermination && js.pumpMsgLoop()) {
1314 // Check if FinalizationRegistry cleanup callbacks have not breached our limits
1315 if (limitEnforcer->getLimitsExceeded() != kj::none) {
1316 // We can potentially log this, but due to a lack of IoContext we cannot notify
1317 // the worker
1318 break;
1319 }
1320 
1321 // It is possible that a microtask got enqueued during pumpMsgLoop execution
1322 // Microtasks enqueued by FinalizationRegistry cleanup tasks should also run
1323 // without an active IoContext
1324 v8::TryCatch tryCatch(workerLock.getIsolate());
1325 
1326 js.runMicrotasks();
1327 
1328 if (tryCatch.HasCaught()) {
1329 // It really shouldn't be possible for microtasks to throw regular exceptions.
1330 // so if we got here it should be a terminal condition.
1331 KJ_ASSERT(tryCatch.HasTerminated());
1332 // If we do not reset here we end up with a dangling exception in the isolate that
1333 // leads to an assert in v8 when the Lock is destroyed.
1334 tryCatch.Reset();
1335 // Ensure we don't pump the message loop in this case
1336 gotTermination = true;
1337 }
1338 }
1339 }
1340 });
1341 
1342 // With --gc-stress, force a full GC before each awaitIo continuation. This helps detect
1343 // KJ async objects (promises, streams, etc.) stored on the JS heap without IoOwn
1344 // wrapping. Such objects crash under DISALLOW_KJ_IO_DESTRUCTORS_SCOPE when collected
1345 // by GC, but normally the timing window is too brief to hit. Forcing GC at every
1346 // continuation makes these bugs deterministic.
1347 if (isGcStressModeForTest()) {
1348 workerLock.getIsolate()->RequestGarbageCollectionForTesting(
1349 v8::Isolate::kFullGarbageCollection);
1350 }
1351 
1352 v8::TryCatch tryCatch(workerLock.getIsolate());
1353 try {
1354 runnable.run(workerLock);
1355 } catch (const jsg::JsExceptionThrown&) {
1356 if (tryCatch.HasTerminated()) {
1357 gotTermination = true;
1358 limiterScope = nullptr;
1359 
1360 // Check if we hit a limit.
1361 limitEnforcer->requireLimitsNotExceeded();
1362 
1363 // Check if we were aborted. TerminateExecution() may be called after abort() in order
1364 // to prevent any more JavaScript from executing.
1365 KJ_IF_SOME(e, abortException) {
1366 kj::throwFatalException(e.clone());
1367 }
1368 
1369 // That should have thrown, so we shouldn't get here.
1370 KJ_FAIL_ASSERT("script terminated for unknown reasons");
1371 } else {
1372 if (tryCatch.Message().IsEmpty()) {
1373 // Should never happen, but check for it because otherwise V8 will crash.
1374 KJ_LOG(ERROR, "tryCatch.Message() was empty even when not HasTerminated()??",
1375 kj::getStackTrace());
1376 JSG_FAIL_REQUIRE(Error, "(JavaScript exception with no message)");
1377 } else {
1378 auto jsException = tryCatch.Exception();
1379 
1380 // TODO(someday): We log "uncaught exception" here whenever throwing from JS to C++.
1381 // However, the C++ code calling us may still catch the exception and do its own logging,
1382 // or may even tunnel it back to JavaScript, making this log line redundant or maybe even
1383 // wrong (if the exception is in fact caught later). But, it's difficult to be sure that
1384 // all C++ consumers log properly, and even if they do, the stack trace is lost once the
1385 // exception has been tunneled into a KJ exception, so the later logging won't be as
1386 // useful. We should improve the tunneling to include stack traces and ensure that all
1387 // consumers do in fact log exceptions, then we can remove this.
1388 workerLock.logUncaughtException(UncaughtExceptionSource::INTERNAL,
1389 jsg::JsValue(jsException), jsg::JsMessage(tryCatch.Message()));
1390 
1391 jsg::throwTunneledException(workerLock.getIsolate(), jsException);
1392 }
1393 }
1394 }
1395 });
1396}
1397 
1398static constexpr auto kAsyncIoErrorMessage =
1399 "Disallowed operation called within global scope. Asynchronous I/O "
1400 "(ex: fetch() or connect()), setting a timeout, and generating random "
1401 "values are not allowed within global scope. To fix this error, perform this "
1402 "operation within a handler. "
1403 "https://developers.cloudflare.com/workers/runtime-apis/handlers/";
1404 
1405IoContext& IoContext::current() {
1406 if (threadLocalRequest == nullptr) {
1407 v8::Isolate* isolate = v8::Isolate::TryGetCurrent();
1408 KJ_REQUIRE(isolate != nullptr, "there is no current request on this thread");
1409 isolate->ThrowError(jsg::v8StrIntern(isolate, kAsyncIoErrorMessage));
1410 throw jsg::JsExceptionThrown();
1411 } else {
1412 return *threadLocalRequest;
1413 }
1414}
1415 
1416kj::Maybe<IoContext&> IoContext::tryCurrent() {
1417 if (threadLocalRequest == nullptr) {
1418 return kj::none;
1419 } else {
1420 return *threadLocalRequest;
1421 }
1422}
1423 
1424bool IoContext::hasCurrent() {
1425 return threadLocalRequest != nullptr;
1426}
1427 
1428bool IoContext::isCurrent() {
1429 return this == threadLocalRequest;
1430}
1431 
1432auto IoContext::tryGetWeakRefForCurrent() -> kj::Maybe<kj::Own<WeakRef>> {
1433 KJ_IF_SOME(ioContext, tryCurrent()) {
1434 return ioContext.getWeakRef();
1435 } else {
1436 return kj::none;
1437 }
1438}
1439 
1440void IoContext::abortFromHang(Worker::AsyncLock& asyncLock) {
1441 KJ_ASSERT(actor == kj::none); // we don't perform hang detection on actor requests
1442 
1443 // Don't bother aborting if limits were exceeded because in that case the abort promise will be
1444 // fulfilled shortly anyway.
1445 if (limitEnforcer->getLimitsExceeded() == kj::none) {
1446 abort(JSG_KJ_EXCEPTION(FAILED, Error,
1447 "The Workers runtime canceled this request because it detected that your Worker's code "
1448 "had hung and would never generate a response. Refer to: "
1449 "https://developers.cloudflare.com/workers/observability/errors/"));
1450 }
1451}
1452 
1453namespace {
1454 
1455class CacheSerializedInputStream final: public kj::AsyncInputStream {
1456 public:
1457 CacheSerializedInputStream(
1458 kj::Own<kj::AsyncInputStream> inner, kj::Own<kj::PromiseFulfiller<void>> fulfiller)
1459 : inner(kj::mv(inner)),
1460 fulfiller(kj::mv(fulfiller)) {}
1461 
1462 ~CacheSerializedInputStream() noexcept(false) {
1463 fulfiller->fulfill();
1464 }
1465 
1466 kj::Promise<size_t> tryRead(void* buffer, size_t minBytes, size_t maxBytes) override {
1467 return inner->tryRead(buffer, minBytes, maxBytes);
1468 }
1469 
1470 kj::Maybe<uint64_t> tryGetLength() override {
1471 return inner->tryGetLength();
1472 }
1473 
1474 kj::Promise<uint64_t> pumpTo(kj::AsyncOutputStream& output, uint64_t amount) override {
1475 return inner->pumpTo(output, amount);
1476 }
1477 
1478 private:
1479 kj::Own<kj::AsyncInputStream> inner;
1480 kj::Own<kj::PromiseFulfiller<void>> fulfiller;
1481};
1482 
1483} // namespace
1484 
1485jsg::Promise<IoOwn<kj::AsyncInputStream>> IoContext::makeCachePutStream(
1486 jsg::Lock& js, kj::Own<kj::AsyncInputStream> stream) {
1487 auto paf = kj::newPromiseAndFulfiller<void>();
1488 
1489 KJ_DEFER(cachePutSerializer = kj::mv(paf.promise));
1490 
1491 return awaitIo(js,
1492 cachePutSerializer.then(
1493 [fulfiller = kj::mv(paf.fulfiller),
1494 stream = kj::mv(stream)]() mutable -> kj::Own<kj::AsyncInputStream> {
1495 if (stream->tryGetLength() != kj::none) {
1496 // PUT with Content-Length. We can just return immediately, allowing the next PUT to start.
1497 KJ_DEFER(fulfiller->fulfill());
1498 return kj::mv(stream);
1499 } else {
1500 // TODO(later): With Cache streams no longer having a size limit enforced by the runtime,
1501 // explore if we can clean up stream serialization too.
1502 // PUT with Transfer-Encoding: chunked. We have no idea how big this request body is going to
1503 // be, so wrap the stream that only unblocks the next PUT after this one is complete.
1504 return kj::heap<CacheSerializedInputStream>(kj::mv(stream), kj::mv(fulfiller));
1505 }
1506 }),
1507 [this](
1508 jsg::Lock&, kj::Own<kj::AsyncInputStream> result) { return addObject(kj::mv(result)); });
1509}
1510 
1511void IoContext::writeLogfwdr(
1512 uint channel, kj::FunctionParam<void(capnp::AnyPointer::Builder)> buildMessage) {
1513 addWaitUntil(getIoChannelFactory()
1514 .writeLogfwdr(channel, kj::mv(buildMessage))
1515 .attach(registerPendingEvent()));
1516}
1517 
1518void IoContext::requireCurrentOrThrowJs() {
1519 if (!isCurrent()) {
1520 throwNotCurrentJsError();
1521 }
1522}
1523 
1524void IoContext::requireCurrentOrThrowJs(WeakRef& weak) {
1525 KJ_IF_SOME(ctx, weak.tryGet()) {
1526 if (ctx.isCurrent()) {
1527 return;
1528 }
1529 }
1530 throwNotCurrentJsError();
1531}
1532 
1533void IoContext::throwNotCurrentJsError(kj::Maybe<const std::type_info&> maybeType) {
1534 auto type = maybeType
1535 .map([](const std::type_info& type) {
1536 return kj::str(" (I/O type: ", jsg::typeName(type), ")");
1537 }).orDefault(kj::String());
1538 
1539 if (threadLocalRequest != nullptr && threadLocalRequest->actor != kj::none) {
1540 JSG_FAIL_REQUIRE(Error,
1541 kj::str(
1542 "Cannot perform I/O on behalf of a different Durable Object. I/O objects "
1543 "(such as streams, request/response bodies, and others) created in the context of one "
1544 "Durable Object cannot be accessed from a different Durable Object in the same isolate. "
1545 "This is a limitation of Cloudflare Workers which allows us to improve overall "
1546 "performance.",
1547 type));
1548 } else {
1549 JSG_FAIL_REQUIRE(Error,
1550 kj::str(
1551 "Cannot perform I/O on behalf of a different request. I/O objects (such as "
1552 "streams, request/response bodies, and others) created in the context of one request "
1553 "handler cannot be accessed from a different request's handler. This is a limitation "
1554 "of Cloudflare Workers which allows us to improve overall performance.",
1555 type));
1556 }
1557}
1558 
1559jsg::JsObject IoContext::getPromiseContextTag(jsg::Lock& js) {
1560 if (promiseContextTag == kj::none) {
1561 auto deferral = kj::heap<IoCrossContextExecutor>(deleteQueue.queue.addRef());
1562 promiseContextTag = jsg::JsRef(js, js.opaque(kj::mv(deferral)));
1563 }
1564 return KJ_REQUIRE_NONNULL(promiseContextTag).getHandle(js);
1565}
1566 
1567kj::Promise<void> IoContext::startDeleteQueueSignalTask(IoContext* context) {
1568 // The promise that is returned is held by the IoContext itself, so when the
1569 // IoContext is destroyed, the promise will be canceled and the loop will
1570 // end. On each iteration of the loop we want to reset the cross thread
1571 // signal in the delete queue, then wait on the promise. Once the promise
1572 // is fulfilled, we will run an empty task to prompt the IoContext to drain
1573 // the DeleteQueue.
1574 try {
1575 for (;;) {
1576 co_await context->deleteQueue.queue->resetCrossThreadSignal();
1577 co_await context->run([](auto& lock) {
1578 auto& context = IoContext::current();
1579 auto l = context.deleteQueue.queue->crossThreadDeleteQueue.lockExclusive();
1580 auto& state = KJ_ASSERT_NONNULL(*l);
1581 for (auto& action: state.actions) {
1582 action(lock);
1583 }
1584 state.actions.clear();
1585 });
1586 }
1587 } catch (...) {
1588 context->abort(kj::getCaughtExceptionAsKj());
1589 }
1590}
1591} // namespace workerd