File
Blob: src/workerd/io/io-context.h
| 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 | #pragma once |
| 6 | |
| 7 | #include "io-own.h" |
| 8 | #include "worker.h" |
| 9 | |
| 10 | #include <workerd/api/deferred-proxy.h> |
| 11 | #include <workerd/io/actor-id.h> |
| 12 | #include <workerd/io/external-pusher.h> |
| 13 | #include <workerd/io/io-channels.h> |
| 14 | #include <workerd/io/io-gate.h> |
| 15 | #include <workerd/io/io-thread-context.h> |
| 16 | #include <workerd/io/io-timers.h> |
| 17 | #include <workerd/io/limit-enforcer.h> |
| 18 | #include <workerd/io/trace.h> |
| 19 | #include <workerd/io/worker-fs.h> |
| 20 | #include <workerd/jsg/async-context.h> |
| 21 | #include <workerd/jsg/jsg.h> |
| 22 | #include <workerd/util/exception.h> |
| 23 | #include <workerd/util/uncaught-exception-source.h> |
| 24 | #include <workerd/util/weak-refs.h> |
| 25 | |
| 26 | #include <capnp/dynamic.h> |
| 27 | #include <kj/async-io.h> |
| 28 | #include <kj/compat/http.h> |
| 29 | #include <kj/function.h> |
| 30 | #include <kj/mutex.h> |
| 31 | |
| 32 | namespace workerd { |
| 33 | class WorkerTracer; |
| 34 | class BaseTracer; |
| 35 | } // namespace workerd |
| 36 | |
| 37 | namespace workerd { |
| 38 | class LimitEnforcer; |
| 39 | } |
| 40 | |
| 41 | namespace capnp { |
| 42 | class HttpOverCapnpFactory; |
| 43 | } |
| 44 | |
| 45 | namespace workerd { |
| 46 | |
| 47 | // This wishes it were IoContext::Runnable::Exceptional. |
| 48 | WD_STRONG_BOOL(IoContext_Runnable_Exceptional); |
| 49 | |
| 50 | [[noreturn]] void throwExceededMemoryLimit(bool isActor); |
| 51 | |
| 52 | class IoContext; |
| 53 | |
| 54 | // Represents one incoming request being handled by a IoContext. In non-actor scenarios, |
| 55 | // there is only ever one IncomingRequest per IoContext, but with actors there could be many. |
| 56 | // |
| 57 | // This should normally be referenced as IoContext::IncomingRequest, but it has been pulled |
| 58 | // out of the nested scope to allow forward-declaration. |
| 59 | // |
| 60 | // The purpose of tracking IncomingRequests at all is so that we can perform metrics, logging, |
| 61 | // and tracing on a "per-request basis", e.g. we can log that a particular incoming request |
| 62 | // generated N subrequests, and traces can trace through them. But this concept falls apart |
| 63 | // a bit when actors are in play, because we can't really say which incoming request "caused" |
| 64 | // any particular subrequest, especially when multiple incoming requests overlap. As a |
| 65 | // heuristic approximation, we attribute each subrequest (and all other forms of resource |
| 66 | // usage) to the "current" incoming request, which is defined as the newest request that hasn't |
| 67 | // already completed. |
| 68 | class IoContext_IncomingRequest final { |
| 69 | public: |
| 70 | IoContext_IncomingRequest(kj::Own<IoContext> context, |
| 71 | kj::Own<IoChannelFactory> ioChannelFactory, |
| 72 | kj::Own<RequestObserver> metrics, |
| 73 | kj::Maybe<kj::Own<BaseTracer>> workerTracer, |
| 74 | kj::Maybe<tracing::InvocationSpanContext> maybeTriggerInvocationSpan); |
| 75 | KJ_DISALLOW_COPY_AND_MOVE(IoContext_IncomingRequest); |
| 76 | ~IoContext_IncomingRequest() noexcept(false); |
| 77 | |
| 78 | IoContext& getContext() { |
| 79 | return *context; |
| 80 | } |
| 81 | |
| 82 | // Invoked when the request is actually delivered. |
| 83 | // |
| 84 | // If, for some reason, this is not invoked before the object is destroyed, this indicate that |
| 85 | // the event was canceled for some reason before delivery. No JavaScript was invoked. |
| 86 | // |
| 87 | // This method invokes metrics->delivered() and also makes this IncomingRequest "current" for |
| 88 | // the IoContext. |
| 89 | // |
| 90 | // If delivered() is never called, then drain() need not be called. |
| 91 | void delivered(kj::SourceLocation = kj::SourceLocation()); |
| 92 | |
| 93 | // Waits until the request is "done". For non-actor requests this means waiting until |
| 94 | // all "waitUntil" tasks finish, applying the "soft timeout" time limit from WorkerLimits. |
| 95 | // |
| 96 | // For actor requests, this means waiting until either all tasks have finished (not just |
| 97 | // waitUntil, all tasks), or a new incoming request has been received (which then takes over |
| 98 | // responsibility for waiting for tasks), or the actor is shut down. |
| 99 | kj::Promise<void> drain(); |
| 100 | |
| 101 | // Waits for all "waitUntil" tasks to finish, up to the time limit for scheduled events, as |
| 102 | // defined by `scheduledTimeoutMs` in `WorkerLimits`. Returns an enum indicating the event outcome |
| 103 | // based on whether the given tasks completed successfully, hit a timeout, or were aborted. |
| 104 | // |
| 105 | // Note that, while this is similar in some ways to `drain()`, `finishScheduled()` is intended |
| 106 | // to be called synchronously during request handling, i.e. where a client is waiting for the |
| 107 | // result, and the operation will be canceled if the client disconnects. `drain()` is intended |
| 108 | // to be called after the client has received a response or disconnected. |
| 109 | // |
| 110 | // This method is also used by some custom event handlers (see WorkerInterface::CustomEvent) that |
| 111 | // need similar behavior, as well as the test handler. TODO(cleanup): Rename to something more |
| 112 | // generic? |
| 113 | kj::Promise<EventOutcome> finishScheduled(); |
| 114 | |
| 115 | // Access the event loop's current time point. This will remain constant between ticks. This is |
| 116 | // used to implement IoContext::now(), which should be preferred so that time can be adjusted |
| 117 | // based on setTimeout() when needed. |
| 118 | kj::Date now(kj::Maybe<kj::Date> nextTimeout = kj::none); |
| 119 | |
| 120 | RequestObserver& getMetrics() { |
| 121 | return *metrics; |
| 122 | } |
| 123 | |
| 124 | kj::Maybe<BaseTracer&> getWorkerTracer() { |
| 125 | return workerTracer; |
| 126 | } |
| 127 | |
| 128 | // Returns a new reference to the root user trace span for this incoming request, or |
| 129 | // SpanParent(nullptr) if the request has no user-tracing root span. |
| 130 | SpanParent getRootUserTraceSpan() { |
| 131 | return rootUserTraceSpan.addRef(); |
| 132 | } |
| 133 | |
| 134 | // The invocation span context is a unique identifier for a specific |
| 135 | // worker invocation. |
| 136 | tracing::InvocationSpanContext& getInvocationSpanContext(); |
| 137 | |
| 138 | private: |
| 139 | kj::Own<IoContext> context; |
| 140 | kj::Own<RequestObserver> metrics; |
| 141 | kj::Maybe<kj::Own<BaseTracer>> workerTracer; |
| 142 | kj::Own<IoChannelFactory> ioChannelFactory; |
| 143 | |
| 144 | // Root user trace span for this request. Populated during delivered() via |
| 145 | // BaseTracer::makeUserRequestSpan(); otherwise a null SpanParent. The tracer it references |
| 146 | // is owned by workerTracer above; because user-tracing SpanSubmitters hold only a |
| 147 | // BaseTracer::WeakRef, stale SpanParent references (e.g. in AsyncContextFrame storage via |
| 148 | // IoOwn, kept alive past ~IncomingRequest by the IoContext's delete queue) cannot extend |
| 149 | // tracer lifetime. |
| 150 | SpanParent rootUserTraceSpan = SpanParent(nullptr); |
| 151 | |
| 152 | // The invocation span context identifies the trace id, invocation id, and root |
| 153 | // span for the current request. Every invocation of a worker function always |
| 154 | // has a root span, even if it is not explicitly traced. |
| 155 | kj::Maybe<tracing::InvocationSpanContext> maybeTriggerInvocationSpan; |
| 156 | kj::Maybe<tracing::InvocationSpanContext> invocationSpanContext; |
| 157 | |
| 158 | bool wasDelivered = false; |
| 159 | |
| 160 | // Used for debugging, tracks whether we properly called drain() or some other mechanism to |
| 161 | // wait for waitUntil tasks. |
| 162 | bool waitedForWaitUntil = false; |
| 163 | |
| 164 | // If drain() was already called, this is non-null and fulfilling it will cancel the drain. |
| 165 | // This is used in particular when a new IncomingRequest starts while the drain is being |
| 166 | // awaited. |
| 167 | kj::Maybe<kj::Own<kj::PromiseFulfiller<void>>> drainFulfiller; |
| 168 | |
| 169 | // Used by IoContext::incomingRequests. |
| 170 | kj::ListLink<IoContext_IncomingRequest> link; |
| 171 | |
| 172 | // Tracks the location where delivered() was called for debugging. |
| 173 | kj::Maybe<kj::SourceLocation> deliveredLocation; |
| 174 | |
| 175 | friend class IoContext; |
| 176 | }; |
| 177 | |
| 178 | // IoContext holds state associated with a single I/O context. For stateless requests, each |
| 179 | // incoming request runs in a unique I/O context. For actors, each actor runs in a unique I/O |
| 180 | // context (but all requests received by that actor run in the same context). |
| 181 | // |
| 182 | // The IoContext serves as a bridge between JavaScript objects and I/O objects. I/O |
| 183 | // objects are strongly tied to the KJ event loop, and thus must live on a single thread. The |
| 184 | // JS isolate, however, can move between threads, bringing all garbage-collected heap objects |
| 185 | // with it. So, when a GC'ed object holds a reference to I/O objects or tasks (KJ promises), it |
| 186 | // needs help from IoContext manage this. |
| 187 | // |
| 188 | // Whenever JavaScript is executing, the current IoContext can be obtained via |
| 189 | // `IoContext::current()`, and this can then be used to manage I/O, such as outgoing |
| 190 | // subrequests. When the IoContext is destroyed, all outstanding I/O objects and tasks |
| 191 | // created through it are destroyed immediately, even if objects on the JS heap still refer to |
| 192 | // them. Any attempt to access an I/O object from the wrong context will throw. |
| 193 | // |
| 194 | // This has an observable side-effect for workers: if a worker saves the request objects |
| 195 | // associated with one request into its global state and then attempts to access those objects |
| 196 | // within callbacks associated with some other request, an exception will be thrown. We actually |
| 197 | // like this. We don't want people leaking heavy objects or allowing simultaneous requests to |
| 198 | // interfere with each other. |
| 199 | class IoContext final: public kj::Refcounted, private kj::TaskSet::ErrorHandler { |
| 200 | public: |
| 201 | class TimeoutManagerImpl; |
| 202 | |
| 203 | // Construct a new IoContext. Before using it, you must also create an IncomingRequest. |
| 204 | IoContext(ThreadContext& thread, |
| 205 | kj::Own<const Worker> worker, |
| 206 | kj::Maybe<Worker::Actor&> actor, |
| 207 | kj::Own<LimitEnforcer> limitEnforcer); |
| 208 | |
| 209 | // On destruction, all outstanding tasks associated with this request are canceled. |
| 210 | ~IoContext() noexcept(false); |
| 211 | |
| 212 | using IncomingRequest = IoContext_IncomingRequest; |
| 213 | |
| 214 | const Worker& getWorker() { |
| 215 | return *worker; |
| 216 | } |
| 217 | Worker::Lock& getCurrentLock() { |
| 218 | return KJ_REQUIRE_NONNULL(currentLock); |
| 219 | } |
| 220 | |
| 221 | kj::Maybe<Worker::Actor&> getActor() { |
| 222 | return actor; |
| 223 | } |
| 224 | |
| 225 | // Gets the actor, throwing if there isn't one. |
| 226 | Worker::Actor& getActorOrThrow(); |
| 227 | |
| 228 | RequestObserver& getMetrics() { |
| 229 | return *getCurrentIncomingRequest().metrics; |
| 230 | } |
| 231 | |
| 232 | kj::Maybe<BaseTracer&> getWorkerTracer() { |
| 233 | if (incomingRequests.empty()) return kj::none; |
| 234 | return getCurrentIncomingRequest().getWorkerTracer(); |
| 235 | } |
| 236 | |
| 237 | // Returns the root user trace span for the current incoming request, if any. |
| 238 | SpanParent getRootUserTraceSpan() { |
| 239 | if (incomingRequests.empty()) return SpanParent(nullptr); |
| 240 | return getCurrentIncomingRequest().getRootUserTraceSpan(); |
| 241 | } |
| 242 | |
| 243 | LimitEnforcer& getLimitEnforcer() { |
| 244 | return *limitEnforcer; |
| 245 | } |
| 246 | |
| 247 | // Get the current input lock. Throws an exception if no input lock is held (e.g. because this is |
| 248 | // not an actor request). |
| 249 | InputGate::Lock getInputLock(); |
| 250 | |
| 251 | // Get the current CriticalSection, if there is one, or returns null if not. |
| 252 | kj::Maybe<kj::Own<InputGate::CriticalSection>> getCriticalSection(); |
| 253 | |
| 254 | // Runs `callback` within its own critical section, returning its final result. If `callback` |
| 255 | // throws, the input lock will break, resetting the actor. |
| 256 | // |
| 257 | // This can only be called when I/O gates are active, i.e. in an actor. |
| 258 | template <typename Func> |
| 259 | jsg::PromiseForResult<Func, void, true> blockConcurrencyWhile(jsg::Lock& js, Func&& callback); |
| 260 | |
| 261 | // Returns true if output lock gating is necessary. |
| 262 | // Can be used in optimizations to bypass wait* calls altogether. |
| 263 | bool hasOutputGate(); |
| 264 | |
| 265 | // Wait until all outstanding output locks have been unlocked. Does not wait for future output |
| 266 | // locks, even if they are created before past locks are unlocked. |
| 267 | // |
| 268 | // This is used in actors to block output while some storage writes are uncommitted. For |
| 269 | // non-actor requests, this always completes immediately. |
| 270 | kj::Promise<void> waitForOutputLocks(); |
| 271 | |
| 272 | // Like waitForOutputLocks() but, as an optimization, returns null in (some) cases where no |
| 273 | // wait is needed, such as when the request is not an actor request. |
| 274 | // |
| 275 | // Use the ...IoOwn() overload if you need to store this promise in a JS API object. |
| 276 | kj::Maybe<kj::Promise<void>> waitForOutputLocksIfNecessary(); |
| 277 | kj::Maybe<IoOwn<kj::Promise<void>>> waitForOutputLocksIfNecessaryIoOwn(); |
| 278 | |
| 279 | // Check if the output gate (only used by actors) is currently broken. This indicates that there |
| 280 | // was a problem with committing storage writes. |
| 281 | // |
| 282 | // For non-actor requests, this always returns false. |
| 283 | bool isOutputGateBroken(); |
| 284 | |
| 285 | // Lock output until the given promise completes. |
| 286 | // |
| 287 | // It is an error to call this outside of actors. |
| 288 | template <typename T> |
| 289 | kj::Promise<T> lockOutputWhile(kj::Promise<T> promise); |
| 290 | |
| 291 | bool isInspectorEnabled(); |
| 292 | |
| 293 | // Returns true if there is something listening for warnings โ the Chrome DevTools inspector, |
| 294 | // a streaming tail worker tracer, or --verbose stderr logging. Use this to guard expensive |
| 295 | // warning-message construction that should be skipped when nobody would see the result. |
| 296 | bool hasWarningHandler(); |
| 297 | |
| 298 | // Log a warning. Emits to the Chrome DevTools inspector (if connected), stderr, and to the |
| 299 | // streaming tail worker tracer (if active). |
| 300 | void logWarning(kj::StringPtr description); |
| 301 | |
| 302 | // Log a warning, deduplicating so that each unique message is only logged once for the lifetime |
| 303 | // of an isolate. Emits to the same destinations as logWarning(). |
| 304 | void logWarningOnce(kj::StringPtr description); |
| 305 | |
| 306 | // Log an internal error message. Deduplicates log messages such that a single unique message will |
| 307 | // only be logged once for the lifetime of an isolate. |
| 308 | void logErrorOnce(kj::StringPtr description); |
| 309 | |
| 310 | void logUncaughtException(kj::StringPtr description); |
| 311 | void logUncaughtException(UncaughtExceptionSource source, |
| 312 | const jsg::JsValue& exception, |
| 313 | const jsg::JsMessage& message = jsg::JsMessage()); |
| 314 | |
| 315 | // Log an uncaught exception from an asynchronous context, i.e. when the IoContext is not |
| 316 | // "current". |
| 317 | void logUncaughtExceptionAsync(UncaughtExceptionSource source, kj::Exception&& e); |
| 318 | |
| 319 | // Returns a promise that will reject with an exception if and when the request should be |
| 320 | // aborted, e.g. because its CPU time expired. This should be joined with any promises for |
| 321 | // incoming tasks. |
| 322 | kj::Promise<void> onAbort() { |
| 323 | return abortPromise.addBranch(); |
| 324 | } |
| 325 | |
| 326 | // Force context abort now. |
| 327 | // |
| 328 | // Note that abort() is safe to call while the IoContext is current. Becaues of this, it cannot |
| 329 | // cancel any tasks synchronously, as this might cancel the current promise, leading to a crash. |
| 330 | void abort(kj::Exception&& e); |
| 331 | |
| 332 | // Await the given promise and, if it throws, call `abort()` with the exception. The promise |
| 333 | // given here should just be a monitoring promise, it should not represent any sort of background |
| 334 | // work beyond monitoring. In particular, it must not be a task that attempts to enter the |
| 335 | // isolate by calling context.run(). |
| 336 | void abortWhen(kj::Promise<void> promise); |
| 337 | |
| 338 | // Has event.passThroughOnException() been called? |
| 339 | bool isFailOpen() { |
| 340 | return failOpen; |
| 341 | } |
| 342 | |
| 343 | // Called by event.passThroughOnException(). |
| 344 | void setFailOpen() { |
| 345 | failOpen = true; |
| 346 | } |
| 347 | |
| 348 | // ----------------------------------------------------------------- |
| 349 | // Tracking thread-local request |
| 350 | |
| 351 | // Asynchronously execute a callback inside the context. |
| 352 | // |
| 353 | // We don't use a "scope" class because this might actually switch to a larger stack for the |
| 354 | // duration of the callback. |
| 355 | // |
| 356 | // If `inputLock` is not provided, and this is an actor context, an input lock will be obtained |
| 357 | // before executing the callback. |
| 358 | template <typename Func> |
| 359 | kj::PromiseForResult<Func, Worker::Lock&> run( |
| 360 | Func&& func, kj::Maybe<InputGate::Lock> inputLock = kj::none) KJ_WARN_UNUSED_RESULT; |
| 361 | |
| 362 | // Like run() but executes within the given critical section, if it is non-null. If |
| 363 | // `criticalSection` is null, then this just forwards to the other run() (with null inputLock). |
| 364 | template <typename Func> |
| 365 | kj::PromiseForResult<Func, Worker::Lock&> run(Func&& func, |
| 366 | kj::Maybe<kj::Own<InputGate::CriticalSection>> criticalSection) KJ_WARN_UNUSED_RESULT; |
| 367 | |
| 368 | // Returns the current IoContext for the thread. |
| 369 | // Throws an exception if there is no current context (see hasCurrent() below). |
| 370 | static IoContext& current(); |
| 371 | |
| 372 | // Like current(), but returns kj::none if there is no current context. |
| 373 | static kj::Maybe<IoContext&> tryCurrent(); |
| 374 | |
| 375 | // True if there is a current IoContext for the thread (current() will not throw). |
| 376 | static bool hasCurrent(); |
| 377 | |
| 378 | // True if this is the IoContext for the current thread (same as `hasCurrent() && tcx == current()`). |
| 379 | bool isCurrent(); |
| 380 | |
| 381 | // Check if a current request is available. Used to provide better diagnostics when this is |
| 382 | // unexpectedly absent when reporting a user span. |
| 383 | // TODO(cleanup): This is a hack, remove after addressing the underlying issue. |
| 384 | bool hasCurrentIncomingRequest() { |
| 385 | return !incomingRequests.empty(); |
| 386 | } |
| 387 | |
| 388 | // Like requireCurrent() but throws a JS error if this IoContext is not the current. |
| 389 | void requireCurrentOrThrowJs(); |
| 390 | |
| 391 | // A WeakRef is a weak reference to a IoContext. Note that because IoContext is not |
| 392 | // itself ref-counted, we cannot follow the usual pattern of a weak reference that potentially |
| 393 | // converts to a strong reference. Instead, intended usage looks like so: |
| 394 | // ``` |
| 395 | // auto& context = IoContext::current(); |
| 396 | // return canOutliveContext().then([contextWeakRef = context.getWeakRef()]() mutable { |
| 397 | // auto hadContext = contextWeakRef.runIfAlive([&](IoContext& context){ |
| 398 | // useContextFinally(context); |
| 399 | // }); |
| 400 | // if (!hadContext) { |
| 401 | // doWhatMustBeDone(); |
| 402 | // } |
| 403 | // }); |
| 404 | // ``` |
| 405 | using WeakRef = workerd::WeakRef<IoContext>; |
| 406 | |
| 407 | kj::Own<WeakRef> getWeakRef() { |
| 408 | return kj::addRef(*selfRef); |
| 409 | } |
| 410 | |
| 411 | // If there is a current IoContext, return its WeakRef. |
| 412 | static kj::Maybe<kj::Own<WeakRef>> tryGetWeakRefForCurrent(); |
| 413 | |
| 414 | // Like requireCurrentOrThrowJs() but works on a WeakRef. |
| 415 | static void requireCurrentOrThrowJs(WeakRef& weak); |
| 416 | |
| 417 | // Just throw the error that requireCurrentOrThrowJs() would throw on failure. |
| 418 | [[noreturn]] static void throwNotCurrentJsError( |
| 419 | kj::Maybe<const std::type_info&> maybeType = kj::none); |
| 420 | |
| 421 | // ----------------------------------------------------------------- |
| 422 | // Task scheduling and object storage |
| 423 | |
| 424 | // Arrange for the given promise to execute as part of this request. It will be canceled if the |
| 425 | // request is canceled. |
| 426 | void addTask(kj::Promise<void> promise); |
| 427 | |
| 428 | template <typename T, typename Func> |
| 429 | jsg::PromiseForResult<Func, T, true> awaitIo(jsg::Lock& js, kj::Promise<T> promise, Func&& func); |
| 430 | |
| 431 | // Attach the objects to the promise by creating a continuation that holds them. |
| 432 | // This ensures the attachments stay alive until the promise resolves. |
| 433 | // This should ONLY be used with TraceContext or SpanBuilder objects. |
| 434 | template <typename T, typename... Attachments> |
| 435 | jsg::Promise<T> attachSpans(jsg::Lock& js, jsg::Promise<T> promise, Attachments&&... attachments) |
| 436 | requires(... && |
| 437 | (kj::isSameType<Attachments, SpanBuilder>() || kj::isSameType<Attachments, TraceContext>())) |
| 438 | { |
| 439 | return attachSpansInternalOnly(js, kj::mv(promise), kj::fwd<Attachments>(attachments)...); |
| 440 | } |
| 441 | |
| 442 | // public for tests |
| 443 | template <typename T, typename... Attachments> |
| 444 | jsg::Promise<T> attachSpansInternalOnly( |
| 445 | jsg::Lock& js, jsg::Promise<T> promise, Attachments&&... attachments) { |
| 446 | auto attachmentTuple = addObject(kj::heap(kj::tuple(kj::fwd<Attachments>(attachments)...))); |
| 447 | |
| 448 | if constexpr (kj::isSameType<T, void>()) { |
| 449 | return promise.then(js, [attachmentTuple = kj::mv(attachmentTuple)](jsg::Lock&) { |
| 450 | // The attachments are kept alive in this lambda's capture |
| 451 | }); |
| 452 | } else { |
| 453 | return promise.then(js, [attachmentTuple = kj::mv(attachmentTuple)](jsg::Lock&, T result) { |
| 454 | // The attachments are kept alive in this lambda's capture |
| 455 | return result; |
| 456 | }); |
| 457 | } |
| 458 | } |
| 459 | |
| 460 | // Waits for some background I/O to complete, then executes `func` on the result, returning a |
| 461 | // JavaScript promise for the result of that. If no `func` is provided, no transformation is |
| 462 | // applied. |
| 463 | // |
| 464 | // If the IoContext is canceled, the I/O promise will be canceled, `func` will be destroyed |
| 465 | // without being called, and the JS promise will never resolve. |
| 466 | // |
| 467 | // You might wonder why this function takes a continuation function as a parameter, rather than |
| 468 | // taking a single `kj::Promise<T>`, returning `jsg::Promise<T>`, and leaving it up to you to |
| 469 | // call `.then()` on the result. The answer is that `func` provides stronger guarantees about the |
| 470 | // context where it runs, which avoids the need for `IoOwn`s: |
| 471 | // - `func` itself can safely capture I/O objects without IoOwn, because the function itself |
| 472 | // is attached to the IoContext. (If the IoContext is canceled, `func` is destroyed.) |
| 473 | // - Similarly, the result of `promise` can be an I/O object without needing to be wrapped in |
| 474 | // IoOwn, because `func` is guaranteed to be called in this IoContext. |
| 475 | // |
| 476 | // Conversely, you might wonder why you wouldn't use `awaitIo(promise.then(func))` instead, which |
| 477 | // would also avoid the need for `IoOwn` since `func` would run as part of the KJ event loop. |
| 478 | // But, in this version, `func` cannot access any JavaScript objects, because it would not run |
| 479 | // with the isolate lock. |
| 480 | // |
| 481 | // Historically, we solved this with something called `capctx`. You would write something like: |
| 482 | // `awaitIo(promise.then(capctx(func)))`. This provided both properties: `func()` ran both in |
| 483 | // the KJ event loop and with the isolate lock held. However, this had the problem that it |
| 484 | // required returning to the KJ event loop between running func() and running whatever |
| 485 | // JavaScript code was waiting on it. This implies releasing the isolate lock just to |
| 486 | // immediately acquire it again, which was wasteful. Passing `func` as a parameter to `awaitIo()` |
| 487 | // allows it to run under the same isolate lock that then runs the awaiting JavaScript. |
| 488 | // |
| 489 | // Note that awaitIo() automatically implies registering a pending event while waiting for the |
| 490 | // promise (no need to call registerPendingEvent()). |
| 491 | template <typename T> |
| 492 | jsg::Promise<T> awaitIo(jsg::Lock& js, kj::Promise<T> promise); |
| 493 | |
| 494 | // Waits for the given I/O while holding the input lock, so that all other I/O is blocked from |
| 495 | // completing in the meantime (unless it is also holding the same input lock). |
| 496 | template <typename T> |
| 497 | jsg::Promise<T> awaitIoWithInputLock(jsg::Lock& js, kj::Promise<T> promise); |
| 498 | |
| 499 | template <typename T, typename Func> |
| 500 | jsg::PromiseForResult<Func, T, true> awaitIoWithInputLock( |
| 501 | jsg::Lock& js, kj::Promise<T> promise, Func&& func); |
| 502 | |
| 503 | // DEPRECATED: Like awaitIo() but: |
| 504 | // - Does not have a continuation function, so suffers from the problems described in |
| 505 | // `awaitIo()`'s doc comment. |
| 506 | // - Does not automatically register a pending event. |
| 507 | // |
| 508 | // This is used to implement the historical KJ-oriented PromiseWrapper behavior in terms of the |
| 509 | // new `awaitIo()` implementation. This should go away once all API implementations are |
| 510 | // refactored to use `awaitIo()`. |
| 511 | template <typename T> |
| 512 | jsg::Promise<T> awaitIoLegacy(jsg::Lock& js, kj::Promise<T> promise); |
| 513 | |
| 514 | // DEPRECATED: Like awaitIo() but: |
| 515 | // - Does not have a continuation function, so suffers from the problems described in |
| 516 | // `awaitIo()`'s doc comment. |
| 517 | // - Does not automatically register a pending event. |
| 518 | // |
| 519 | // This is used to implement the historical KJ-oriented PromiseWrapper behavior in terms of the |
| 520 | // new `awaitIo()` implementation. This should go away once all API implementations are |
| 521 | // refactored to use `awaitIo()`. |
| 522 | template <typename T> |
| 523 | jsg::Promise<T> awaitIoLegacyWithInputLock(jsg::Lock& js, kj::Promise<T> promise); |
| 524 | |
| 525 | // Returns a KJ promise that resolves when a particular JavaScript promise completes. |
| 526 | // |
| 527 | // The JS promise must complete within this IoContext. The KJ promise will reject |
| 528 | // immediately if any of these happen: |
| 529 | // - The JS promise is GC'ed without resolving. |
| 530 | // - The JS promise is resolved from the wrong context. |
| 531 | // - The system detects that no further progress will be made in this context (because there is no |
| 532 | // more JavaScript to run, and there is no outstanding I/O scheduled with awaitIo()). |
| 533 | // |
| 534 | // If `T` is `IoOwn<U>`, it will be unwrapped to just `U` in the result. If `U` is in turn |
| 535 | // `kj::Promise<V>`, then the promises will be chained as usual, so the final result is |
| 536 | // `kj::Promise<V>`. |
| 537 | template <typename T> |
| 538 | kj::_::ReducePromises<RemoveIoOwn<T>> awaitJs(jsg::Lock& js, jsg::Promise<T> promise); |
| 539 | |
| 540 | enum TopUpFlag { NO_TOP_UP, TOP_UP }; |
| 541 | |
| 542 | // Make a kj::Function which, when called, re-enters this IoContext to run some code. |
| 543 | // |
| 544 | // `func` is a function with a signature similar to: |
| 545 | // |
| 546 | // template <typename... Params, typename Result> |
| 547 | // jsg::Promise<Result> func(jsg::Lock& js, Params&&... params); |
| 548 | // |
| 549 | // (Optionally, the `jsg::Promise<Result>` can just be `Result` instead.) |
| 550 | // |
| 551 | // The returned lambda will a signature like: |
| 552 | // |
| 553 | // kj::Promise<Result> func(Params&&...); |
| 554 | // |
| 555 | // This function can be invoked without holding the isolate lock. |
| 556 | // |
| 557 | // You might think that all this does is set up a lambda that captures the IoContext and calls |
| 558 | // ctx.run(). But, it turns out getting this right is a lot more complicated. |
| 559 | // - What if the IoContext has been canceled / destroyed, or is destroyed during the callback? |
| 560 | // - What if it still exists, but it's an actor and there's no longer an IncomingRequest? |
| 561 | // - How do you prevent "the script will never generate a response" if the callback is the |
| 562 | // only thing being waited for? |
| 563 | // - What if the call was made within blockConcurrencyWhile()? The callback will be blocked until |
| 564 | // the critical section ends, which could lead to deadlock if the critical section code is |
| 565 | // waiting on it? |
| 566 | // |
| 567 | // This solves all that: |
| 568 | // - If the IoContext is destroyed, the callback throws an exception. |
| 569 | // - However, as long as the callback itself exists, it is treated as if a task were added using |
| 570 | // addTask(). In actors, this blocks hibernation and keeps the IncomingRequest live. |
| 571 | // - Additionally, the calback counts as a PendingEvent. |
| 572 | // - The callback is allowed to run within the critical section (blockConcurrencyWhile()) from |
| 573 | // which it was called. |
| 574 | // |
| 575 | // In short, you should almost never use ctx.run() to re-enter an existing context. You almost |
| 576 | // always want either awaitIo() (to re-enter the context after some KJ promise completes) or |
| 577 | // makeReentryCallback() (to re-enter the context on a callback). |
| 578 | // |
| 579 | // The returned function can be called multiple times. |
| 580 | // |
| 581 | // Note that when invoking the returned function, the function object itself must outlive the |
| 582 | // Promise it returns -- just like a coroutine lambda that has a capture. This should, of course, |
| 583 | // be assumed of all functions that return promises, but classically kj::Promise's own `.then()` |
| 584 | // does not keep its input continuation functions live in this way. If you want to pass the |
| 585 | // callback to `.then()`, you can wrap it in `kj::coCapture()`, but note that this means it can |
| 586 | // only be called once. |
| 587 | // |
| 588 | // Use `makeReentryCallback<IoContext::TOP_UP>(func)` to cause |
| 589 | // `ctx.getLimitEnforcer().topUpActor()` to be called each time the callback is invoked. This is |
| 590 | // useful because `topUpActor()` must be called before entering the isolate lock, so it can't be |
| 591 | // part of the body of the given callback function. |
| 592 | template <TopUpFlag topUp = NO_TOP_UP, typename Func> |
| 593 | auto makeReentryCallback(Func func); |
| 594 | |
| 595 | // Returns the number of times addTask() has been called (even if the tasks have completed). |
| 596 | uint taskCount() { |
| 597 | return addTaskCounter; |
| 598 | } |
| 599 | |
| 600 | // Indicates that the script has requested that it stay active until the given promise resolves. |
| 601 | // drain() waits until all such promises have completed. |
| 602 | void addWaitUntil(kj::Promise<void> promise); |
| 603 | |
| 604 | // Returns the status of waitUntil promises. If a promise fails, this sets the status to the |
| 605 | // one corresponding to the exception type. |
| 606 | EventOutcome waitUntilStatus() const { |
| 607 | return waitUntilStatusValue; |
| 608 | } |
| 609 | |
| 610 | // DO NOT USE, use `addWaitUntil()` instead. |
| 611 | kj::TaskSet& getWaitUntilTasks() { |
| 612 | // TODO(cleanup): This is only needed for use with RpcWorkerInterface, but we can eliminate |
| 613 | // that class's need for waitUntilTasks if we change the signature of sendTraces() to return |
| 614 | // a promise, I think. |
| 615 | return waitUntilTasks; |
| 616 | } |
| 617 | |
| 618 | // Wraps a reference in a wrapper which: |
| 619 | // 1. Will throw an exception if dereferenced while the IoContext is not current for the |
| 620 | // thread. |
| 621 | // 2. Can be safely destroyed from any thread. |
| 622 | // 3. Invalidates itself when the request ends (such that dereferencing throws). |
| 623 | template <typename T> |
| 624 | IoOwn<T> addObject(kj::Own<T> obj); |
| 625 | |
| 626 | // Wraps a reference in a wrapper which: |
| 627 | // 1. Will throw an exception if dereferenced while the IoContext is not current for the |
| 628 | // thread. |
| 629 | // 2. Can be safely destroyed from any thread. |
| 630 | // 3. Invalidates itself when the request ends (such that dereferencing throws). |
| 631 | template <typename T> |
| 632 | IoPtr<T> addObject(T& obj); |
| 633 | |
| 634 | // Like addObject() but takes a functor, returning a functor with the same signature but which |
| 635 | // holds the original functor under a `IoOwn`, and so will stop working if the IoContext |
| 636 | // is no longer valid. This is particularly useful for passing to `jsg::Promise::then()` when |
| 637 | // you need the continuation to run in the correct context. |
| 638 | template <typename Func> |
| 639 | auto addFunctor(Func&& func); |
| 640 | |
| 641 | // Attach an object to the IoContext such that it will be destroyed when either the returned |
| 642 | // reference is dropped OR the IoContext itself is destroyed. In the latter case, further |
| 643 | // attempts to access the returned reference will throw. The reference can only be used and |
| 644 | // destroyed within the same thread as the IoContext lives. |
| 645 | template <typename T> |
| 646 | ReverseIoOwn<T> addObjectReverse(kj::Own<T> obj); |
| 647 | |
| 648 | // Call this to indicate that the caller expects to call into JavaScript in this IoContext |
| 649 | // at some point in the future, in response to some *external* event that the caller is waiting |
| 650 | // for. Then, hold on to the returned handle until that time. This prevents finalizers from being |
| 651 | // called in the meantime. |
| 652 | kj::Own<void> registerPendingEvent(); |
| 653 | // TODO(cleanup): awaitIo() automatically applies this. Is the public method needed anymore? |
| 654 | |
| 655 | // When you want to perform a task that returns Promise<DeferredProxy<T>> and the application |
| 656 | // JavaScript is waiting for the result, use `context.waitForDeferredProxy(promise)` to turn it |
| 657 | // into a regular `Promise<T>`, including registering pending events as needed. |
| 658 | template <typename T> |
| 659 | kj::Promise<T> waitForDeferredProxy(kj::Promise<api::DeferredProxy<T>>&& promise) { |
| 660 | return promise.then([this](api::DeferredProxy<T> deferredProxy) { |
| 661 | return deferredProxy.proxyTask.attach(registerPendingEvent()); |
| 662 | }); |
| 663 | } |
| 664 | |
| 665 | // Like awaitIo(), but handles the specific case of Promise<DeferredProxy>. This is special |
| 666 | // because the convention is that the outer promise is NOT treated as a pending I/O event; it |
| 667 | // may actually be waiting for something to happen in JavaScript land. Once the outer promise |
| 668 | // resolves, the inner promise (the DeferredProxy<T>) is treated as external I/O. |
| 669 | template <typename T> |
| 670 | jsg::Promise<T> awaitDeferredProxy(jsg::Lock& js, kj::Promise<api::DeferredProxy<T>>&& promise) { |
| 671 | return awaitIoImpl( |
| 672 | js, waitForDeferredProxy(kj::mv(promise)), getCriticalSection(), IdentityFunc<T>()); |
| 673 | } |
| 674 | |
| 675 | // Called by ScheduledEvent |
| 676 | void setNoRetryScheduled() { |
| 677 | retryScheduled = false; |
| 678 | } |
| 679 | |
| 680 | // Called by ServiceWorkerGlobalScope::runScheduled |
| 681 | bool shouldRetryScheduled() { |
| 682 | return retryScheduled; |
| 683 | } |
| 684 | |
| 685 | // ----------------------------------------------------------------- |
| 686 | // Access to I/O |
| 687 | |
| 688 | // Used to implement setTimeout(). We don't expose the timer directly because the |
| 689 | // promises it returns need to live in this I/O context, anyway. |
| 690 | TimeoutId setTimeoutImpl( |
| 691 | TimeoutId::Generator& generator, bool repeat, jsg::Function<void()> function, double msDelay); |
| 692 | |
| 693 | // Used to implement clearTimeout(). We don't expose the timer directly because the |
| 694 | // promises it returns need to live in this I/O context, anyway. |
| 695 | void clearTimeoutImpl(TimeoutId key); |
| 696 | |
| 697 | size_t getTimeoutCount(); |
| 698 | |
| 699 | // Access the event loop's current time point. This will remain constant between ticks. |
| 700 | kj::Date now(IncomingRequest& incomingRequest); |
| 701 | |
| 702 | // Access the event loop's current time point. This will remain constant between ticks. |
| 703 | kj::Date now(); |
| 704 | |
| 705 | TmpDirStoreScope& getTmpDirStoreScope() { |
| 706 | KJ_IF_SOME(scope, tmpDirStoreScope) { |
| 707 | return *scope; |
| 708 | } |
| 709 | return *tmpDirStoreScope.emplace(TmpDirStoreScope::create()); |
| 710 | } |
| 711 | |
| 712 | // Returns a promise that resolves once `now() >= when`. |
| 713 | kj::Promise<void> atTime(kj::Date when) { |
| 714 | return getIoChannelFactory().getTimer().atTime(when); |
| 715 | } |
| 716 | |
| 717 | // Returns a promise that resolves after some time. This is intended to be used for implementing |
| 718 | // time limits on some sort of operation, not for implementing application-driven timing, as it |
| 719 | // does not maintain consistency with the clock as observed through Date.now(), e.g. when it |
| 720 | // comes to Spectre mitigations. |
| 721 | kj::Promise<void> afterLimitTimeout(kj::Duration t) { |
| 722 | return getIoChannelFactory().getTimer().afterLimitTimeout(t); |
| 723 | } |
| 724 | |
| 725 | // Provide access to the system CSPRNG. |
| 726 | kj::EntropySource& getEntropySource() { |
| 727 | return thread.getEntropySource(); |
| 728 | } |
| 729 | |
| 730 | capnp::HttpOverCapnpFactory& getHttpOverCapnpFactory() { |
| 731 | return thread.getHttpOverCapnpFactory(); |
| 732 | } |
| 733 | |
| 734 | capnp::ByteStreamFactory& getByteStreamFactory() { |
| 735 | return thread.getByteStreamFactory(); |
| 736 | } |
| 737 | |
| 738 | const kj::HttpHeaderTable& getHeaderTable() { |
| 739 | return thread.getHeaderTable(); |
| 740 | } |
| 741 | const ThreadContext::HeaderIdBundle& getHeaderIds() { |
| 742 | return thread.getHeaderIds(); |
| 743 | } |
| 744 | |
| 745 | kj::Rc<ExternalPusherImpl> getExternalPusher(); |
| 746 | |
| 747 | // Subrequest channel numbers for the two special channels. |
| 748 | // NULL = The channel used by global fetch() when the Request has no fetcher attached. |
| 749 | // NEXT = DEPRECATED: The fetcher attached to Requests delivered by a FetchEvent, so that we can |
| 750 | // detect when an incoming request is passed through to `fetch()` (perhaps with rewrites) |
| 751 | // and treat that case differently. In practice this has proven too confusing, so we don't |
| 752 | // plan to treat NEXT and NULL differently going forward. |
| 753 | static constexpr uint NULL_CLIENT_CHANNEL = 0; |
| 754 | static constexpr uint NEXT_CLIENT_CHANNEL = 1; |
| 755 | |
| 756 | // Number of subrequest channels that have special meaning (and so won't appear in any binding). |
| 757 | static constexpr uint SPECIAL_SUBREQUEST_CHANNEL_COUNT = 2; |
| 758 | |
| 759 | struct SubrequestOptions final { |
| 760 | // When inHouse is true, the subrequest is to an API provided internally. For example calls |
| 761 | // to KV. This primarily affects metrics and limits. |
| 762 | bool inHouse; |
| 763 | |
| 764 | // When true, the client is wrapped by metrics.wrapSubrequestClient() ensuring appropriate |
| 765 | // metrics collection. |
| 766 | bool wrapMetrics; |
| 767 | |
| 768 | // The name to use for the request's span if tracing is turned on. |
| 769 | kj::Maybe<kj::ConstString> operationName; |
| 770 | |
| 771 | // The tracing context to use for the subrequest if tracing is enabled. |
| 772 | kj::Maybe<TraceContext&> existingTraceContext; |
| 773 | }; |
| 774 | |
| 775 | // Wraps a WorkerInterface factory with subrequest accounting: tracing, optional metrics wrapping, |
| 776 | // and an external memory adjustment to pressure V8's GC. All code paths that create HTTP |
| 777 | // connections (including those built from capnp capabilities via getHttpOverCapnpFactory()) |
| 778 | // should route through this function or getSubrequest(). |
| 779 | kj::Own<WorkerInterface> getSubrequestNoChecks( |
| 780 | kj::FunctionParam<kj::Own<WorkerInterface>(TraceContext&, IoChannelFactory&)> func, |
| 781 | SubrequestOptions options); |
| 782 | |
| 783 | // If creating a new subrequest is permitted, calls the given factory function synchronously to |
| 784 | // create one. |
| 785 | // If operationName is specified within options and tracing is enabled, this will add a child span |
| 786 | // to the current trace span for both tracing formats. |
| 787 | // TODO(o11y): In the future we may need to change the interface to support having different span |
| 788 | // names and enforce that only documented spans can be emitted. |
| 789 | kj::Own<WorkerInterface> getSubrequest( |
| 790 | kj::FunctionParam<kj::Own<WorkerInterface>(TraceContext&, IoChannelFactory&)> func, |
| 791 | SubrequestOptions options); |
| 792 | |
| 793 | // Get WorkerInterface objects to use for subrequests. |
| 794 | // |
| 795 | // `channel` specifies which outgoing channel to use. The special channel 0 refers to the "null" |
| 796 | // binding (used for fetches where `request.fetcher` is not set), and channel 1 refers to the |
| 797 | // "next" binding (used when request.fetcher is carried over from the incoming request). |
| 798 | // Named bindings, e.g. Worker2Worker bindings, will have indices starting from 2. Fetcher |
| 799 | // bindings declared via Worker::Global::Fetcher have a corresponding `channel` property to refer |
| 800 | // to these outgoing bindings. |
| 801 | // |
| 802 | // `isInHouse` is true if this client represents an "in house" endpoint, i.e. some API provided |
| 803 | // by the Workers platform. For example, KV namespaces are in-house. This primarily affects |
| 804 | // metrics and limits: |
| 805 | // - In-house requests do not count as "subrequests" for metrics and logging purposes. |
| 806 | // - In-house requests are not subject to the same limits on the number of subrequests per |
| 807 | // request. |
| 808 | // - In preview, in-house requests do not show up in the network tab. |
| 809 | // |
| 810 | // `operationName` is the name to use for the request's span, if tracing is turned on. |
| 811 | kj::Own<WorkerInterface> getSubrequestChannel(uint channel, |
| 812 | bool isInHouse, |
| 813 | kj::Maybe<kj::String> cfBlobJson, |
| 814 | kj::ConstString operationName); |
| 815 | |
| 816 | // Get WorkerInterface objects to use for subrequests. |
| 817 | // |
| 818 | // `channel` specifies which outgoing channel to use. The special channel 0 refers to the "null" |
| 819 | // binding (used for fetches where `request.fetcher` is not set), and channel 1 refers to the |
| 820 | // "next" binding (used when request.fetcher is carried over from the incoming request). |
| 821 | // Named bindings, e.g. Worker2Worker bindings, will have indices starting from 2. Fetcher |
| 822 | // bindings declared via Worker::Global::Fetcher have a corresponding `channel` property to refer |
| 823 | // to these outgoing bindings. |
| 824 | // |
| 825 | // `isInHouse` is true if this client represents an "in house" endpoint, i.e. some API provided |
| 826 | // by the Workers platform. For example, KV namespaces are in-house. This primarily affects |
| 827 | // metrics and limits: |
| 828 | // - In-house requests do not count as "subrequests" for metrics and logging purposes. |
| 829 | // - In-house requests are not subject to the same limits on the number of subrequests per |
| 830 | // request. |
| 831 | // - In preview, in-house requests do not show up in the network tab. |
| 832 | // |
| 833 | // `traceContext` is the trace context to use for the subrequest, if tracing is turned on. |
| 834 | kj::Own<WorkerInterface> getSubrequestChannel( |
| 835 | uint channel, bool isInHouse, kj::Maybe<kj::String> cfBlobJson, TraceContext& traceContext); |
| 836 | |
| 837 | // Like getSubrequestChannel() but doesn't enforce limits. Use for trusted paths only. |
| 838 | kj::Own<WorkerInterface> getSubrequestChannelNoChecks(uint channel, |
| 839 | bool isInHouse, |
| 840 | kj::Maybe<kj::String> cfBlobJson, |
| 841 | kj::Maybe<kj::ConstString> operationName = kj::none); |
| 842 | |
| 843 | // Convenience methods that call getSubrequest*() and adapt the returned WorkerInterface objects |
| 844 | // to HttpClient. |
| 845 | kj::Own<kj::HttpClient> getHttpClient(uint channel, |
| 846 | bool isInHouse, |
| 847 | kj::Maybe<kj::String> cfBlobJson, |
| 848 | kj::ConstString operationName); |
| 849 | |
| 850 | kj::Own<kj::HttpClient> getHttpClient( |
| 851 | uint channel, bool isInHouse, kj::Maybe<kj::String> cfBlobJson, TraceContext& traceContext); |
| 852 | // TODO(cleanup): Make it the caller's job to call asHttpClient() on the result of |
| 853 | // getSubrequest*(). |
| 854 | |
| 855 | // Get a raw Cap'n Proto capability for the given channel. This is appropriate for pure RPC use |
| 856 | // cases (e.g. actor operations, email dispatch) that don't create HTTP connections. If you're |
| 857 | // converting the capability to an HTTP service via getHttpOverCapnpFactory(), use |
| 858 | // getSubrequestNoChecks() instead and call channelFactory.getCapability() from the callback, |
| 859 | // so that the external memory adjustment and other subrequest accounting are applied. |
| 860 | capnp::Capability::Client getCapnpChannel(uint channel) { |
| 861 | return getIoChannelFactory().getCapability(channel); |
| 862 | } |
| 863 | |
| 864 | kj::Own<IoChannelFactory::ActorChannel> getGlobalActorChannel(uint channel, |
| 865 | const ActorIdFactory::ActorId& id, |
| 866 | kj::Maybe<kj::String> locationHint, |
| 867 | ActorGetMode mode, |
| 868 | bool enableReplicaRouting, |
| 869 | ActorRoutingMode routingMode, |
| 870 | SpanParent parentSpan, |
| 871 | kj::Maybe<ActorVersion> version) { |
| 872 | return getIoChannelFactory().getGlobalActor(channel, id, kj::mv(locationHint), mode, |
| 873 | enableReplicaRouting, routingMode, kj::mv(parentSpan), kj::mv(version)); |
| 874 | } |
| 875 | kj::Own<IoChannelFactory::ActorChannel> getColoLocalActorChannel( |
| 876 | uint channel, kj::StringPtr id, SpanParent parentSpan) { |
| 877 | return getIoChannelFactory().getColoLocalActor(channel, id, kj::mv(parentSpan)); |
| 878 | } |
| 879 | |
| 880 | void abortAllActors(kj::Maybe<kj::Exception&> reason) { |
| 881 | getIoChannelFactory().abortAllActors(reason); |
| 882 | } |
| 883 | |
| 884 | void deleteAllActors(kj::Maybe<kj::Exception&> reason) { |
| 885 | getIoChannelFactory().deleteAllActors(reason); |
| 886 | } |
| 887 | |
| 888 | // Condemn and terminate JS isolate |
| 889 | void abortIsolate(kj::StringPtr reason = nullptr); |
| 890 | |
| 891 | // Get an HttpClient to use for Cache API subrequests. |
| 892 | kj::Own<CacheClient> getCacheClient(); |
| 893 | |
| 894 | // Returns an object that ensures an async JS operation started in the current scope captures the |
| 895 | // given trace span, or the current request's trace span, if no span is given. |
| 896 | jsg::AsyncContextFrame::StorageScope makeAsyncTraceScope( |
| 897 | Worker::Lock& lock, kj::Maybe<SpanParent> spanParent = kj::none) KJ_WARN_UNUSED_RESULT; |
| 898 | |
| 899 | // Returns an object that ensures an async JS operation started in the current scope captures |
| 900 | // the given user trace span, or the current incoming request's root user trace span if none is |
| 901 | // given. Storing the span in the AsyncContextFrame (which on actors outlives individual |
| 902 | // requests via the IoContext's delete queue) is safe because user-tracing SpanSubmitter |
| 903 | // implementations hold only a BaseTracer::WeakRef - stale references cannot extend tracer |
| 904 | // lifetime. |
| 905 | jsg::AsyncContextFrame::StorageScope makeUserAsyncTraceScope( |
| 906 | Worker::Lock& lock, kj::Maybe<SpanParent> userSpan = kj::none) KJ_WARN_UNUSED_RESULT; |
| 907 | |
| 908 | // Returns the current span being recorded. If called while the JS lock is held, uses the trace |
| 909 | // information from the current async context, if available. |
| 910 | SpanParent getCurrentTraceSpan(); |
| 911 | SpanParent getCurrentUserTraceSpan(); |
| 912 | |
| 913 | // Returns the invocation's traceId/invocationId paired with the currently-active user |
| 914 | // span's spanId (as pushed by `ctx.tracing.enterSpan`), falling back to the invocation |
| 915 | // root's spanId when no user span is active. |
| 916 | tracing::InvocationSpanContext getInvocationSpanContext() { |
| 917 | auto& base = getCurrentIncomingRequest().getInvocationSpanContext(); |
| 918 | tracing::SpanId sid = getCurrentUserTraceSpan().getSpanId(); |
| 919 | if (sid != tracing::SpanId::nullId) { |
| 920 | return tracing::InvocationSpanContext( |
| 921 | base.getTraceId(), base.getInvocationId(), sid, base.getTraceFlags()); |
| 922 | } |
| 923 | return base.clone(); |
| 924 | } |
| 925 | |
| 926 | // Returns a builder for recording tracing spans (or a no-op builder if tracing is inactive). |
| 927 | // If called while the JS lock is held, uses the trace information from the current async |
| 928 | // context, if available. |
| 929 | [[nodiscard]] SpanBuilder makeTraceSpan(kj::ConstString operationName); |
| 930 | // Returns both an internal and a user tracing span, this ensures that all user spans are |
| 931 | // available in internal tracing. |
| 932 | [[nodiscard]] TraceContext makeUserTraceSpan(kj::ConstString operationName); |
| 933 | |
| 934 | // Implement per-IoContext rate limiting for Cache.put(). Pass the body of a Cache API PUT |
| 935 | // request and get a possibly wrapped stream back. |
| 936 | // |
| 937 | // If the stream has an unknown length, you will get a wrapped stream back that is used to |
| 938 | // serialize PUT requests. |
| 939 | jsg::Promise<IoOwn<kj::AsyncInputStream>> makeCachePutStream( |
| 940 | jsg::Lock& js, kj::Own<kj::AsyncInputStream> stream); |
| 941 | // TODO(cleanup): Factor this into getCacheClient() somehow so it's not opt-in. |
| 942 | |
| 943 | // Gets a CapabilityServerSet representing the capnp capabilities hosted by this request or |
| 944 | // actor context. This allows us to implement the CapnpCapability::unwrap() method on |
| 945 | // capabilities which allows the application to get at the underlying server object, when the |
| 946 | // capability points to a local object. |
| 947 | capnp::CapabilityServerSet<capnp::DynamicCapability>& getLocalCapSet() { |
| 948 | return localCapSet; |
| 949 | } |
| 950 | |
| 951 | void writeLogfwdr(uint channel, kj::FunctionParam<void(capnp::AnyPointer::Builder)> buildMessage); |
| 952 | |
| 953 | jsg::JsObject getPromiseContextTag(jsg::Lock& js); |
| 954 | |
| 955 | // The IoChannelFactory must be accessed through the |
| 956 | // currentIncomingRequest because it has some tracing context built in. |
| 957 | // |
| 958 | // TODO(later): this is made public for Python Workers. It should be possible to make this private |
| 959 | // again later. |
| 960 | IoChannelFactory& getIoChannelFactory() { |
| 961 | return *getCurrentIncomingRequest().ioChannelFactory; |
| 962 | } |
| 963 | |
| 964 | void pumpMessageLoop(); |
| 965 | |
| 966 | private: |
| 967 | ThreadContext& thread; |
| 968 | |
| 969 | kj::Own<WeakRef> selfRef = kj::refcounted<WeakRef>(kj::Badge<IoContext>(), *this); |
| 970 | |
| 971 | kj::Maybe<kj::Own<TmpDirStoreScope>> tmpDirStoreScope; |
| 972 | |
| 973 | kj::Own<const Worker> worker; |
| 974 | kj::Maybe<Worker::Actor&> actor; |
| 975 | kj::Own<LimitEnforcer> limitEnforcer; |
| 976 | |
| 977 | // List of active IncomingRequests, ordered from most-recently-started to least-recently-started. |
| 978 | kj::List<IncomingRequest, &IncomingRequest::link> incomingRequests; |
| 979 | |
| 980 | kj::Maybe<kj::SourceLocation> lastDeliveredLocation; |
| 981 | |
| 982 | capnp::CapabilityServerSet<capnp::DynamicCapability> localCapSet; |
| 983 | |
| 984 | bool failOpen = false; |
| 985 | |
| 986 | // For debug checks. |
| 987 | void* threadId; |
| 988 | |
| 989 | // For scheduled workers noRetry calls |
| 990 | bool retryScheduled = true; |
| 991 | |
| 992 | kj::Maybe<Worker::Lock&> currentLock; |
| 993 | kj::Maybe<InputGate::Lock> currentInputLock; |
| 994 | |
| 995 | DeleteQueuePtr deleteQueue; |
| 996 | |
| 997 | kj::Maybe<kj::Exception> abortException; |
| 998 | kj::Own<kj::PromiseFulfiller<void>> abortFulfiller; |
| 999 | kj::ForkedPromise<void> abortPromise = nullptr; |
| 1000 | |
| 1001 | class PendingEvent; |
| 1002 | |
| 1003 | kj::Maybe<PendingEvent&> pendingEvent; |
| 1004 | kj::Maybe<kj::Promise<void>> abortFromHangTask; |
| 1005 | |
| 1006 | // Objects pointed to by IoOwn<T>s. |
| 1007 | // NOTE: This must live below `deleteQueue`, as some of these OwnedObjects may own attachctx()'ed |
| 1008 | // objects which reference `deleteQueue` in their destructors. |
| 1009 | OwnedObjectList ownedObjects; |
| 1010 | |
| 1011 | kj::Maybe<kj::Rc<ExternalPusherImpl>> externalPusher; |
| 1012 | |
| 1013 | // Implementation detail of makeCachePutStream(). |
| 1014 | |
| 1015 | // TODO: Used for Cache PUT serialization. |
| 1016 | kj::Promise<void> cachePutSerializer; |
| 1017 | |
| 1018 | // The timeout manager needs to live below `deleteQueue` because the promises may refer to |
| 1019 | // objects in the queue. |
| 1020 | // |
| 1021 | // ATTENTION: `timeoutManager` MUST be declared before both `waitUntilTasks` and `tasks` so it |
| 1022 | // outlives them. During TaskSet destruction, deferred callbacks (e.g. the one in Scheduler::wait |
| 1023 | // that clears the timer slot via clearTimeoutImpl) still need a live timeoutManager. C++ destroys |
| 1024 | // members in reverse declaration order, so declaring timeoutManager first ensures it is destroyed |
| 1025 | // last among these three. |
| 1026 | kj::Own<TimeoutManager> timeoutManager; |
| 1027 | |
| 1028 | kj::TaskSet waitUntilTasks; |
| 1029 | EventOutcome waitUntilStatusValue = EventOutcome::OK; |
| 1030 | |
| 1031 | void setTimeoutImpl(TimeoutId timeoutId, |
| 1032 | bool repeat, |
| 1033 | jsg::V8Ref<v8::Function> function, |
| 1034 | double msDelay, |
| 1035 | kj::Array<jsg::Value> args); |
| 1036 | |
| 1037 | uint addTaskCounter = 0; |
| 1038 | kj::TaskSet tasks; |
| 1039 | |
| 1040 | // This canceler will be canceled when the IoContext is destroyed. Use it to wrap promises that |
| 1041 | // need to be held externally but which should error if the IoContext is canceled. This is used |
| 1042 | // for `makeReentryCallback()` in particular. |
| 1043 | kj::Canceler canceler; |
| 1044 | |
| 1045 | kj::Own<WorkerInterface> getSubrequestChannelImpl(uint channel, |
| 1046 | bool isInHouse, |
| 1047 | kj::Maybe<kj::String> cfBlobJson, |
| 1048 | TraceContext& tracing, |
| 1049 | IoChannelFactory& channelFactory); |
| 1050 | |
| 1051 | friend class IoContext_IncomingRequest; |
| 1052 | template <typename T> |
| 1053 | friend class IoOwn; |
| 1054 | template <typename T> |
| 1055 | friend class IoPtr; |
| 1056 | |
| 1057 | void taskFailed(kj::Exception&& exception) override; |
| 1058 | void requireCurrent(); |
| 1059 | void checkFarGet(const DeleteQueue& expectedQueue, const std::type_info& type); |
| 1060 | |
| 1061 | kj::Maybe<jsg::JsRef<jsg::JsObject>> promiseContextTag; |
| 1062 | |
| 1063 | class Runnable { |
| 1064 | public: |
| 1065 | using Exceptional = IoContext_Runnable_Exceptional; |
| 1066 | virtual void run(Worker::Lock& lock) = 0; |
| 1067 | }; |
| 1068 | void runImpl(Runnable& runnable, |
| 1069 | Worker::LockType lockType, |
| 1070 | kj::Maybe<InputGate::Lock> inputLock, |
| 1071 | Runnable::Exceptional exceptional); |
| 1072 | |
| 1073 | void abortFromHang(Worker::AsyncLock& asyncLock); |
| 1074 | |
| 1075 | template <typename T> |
| 1076 | struct IdentityFunc { |
| 1077 | inline T operator()(jsg::Lock&, T&& value) const { |
| 1078 | return kj::mv(value); |
| 1079 | } |
| 1080 | }; |
| 1081 | template <> |
| 1082 | struct IdentityFunc<void> { |
| 1083 | inline void operator()(jsg::Lock&) const {} |
| 1084 | }; |
| 1085 | |
| 1086 | template <typename T> |
| 1087 | struct ExceptionOr_ { |
| 1088 | using Type = kj::OneOf<T, kj::Exception>; |
| 1089 | }; |
| 1090 | template <> |
| 1091 | struct ExceptionOr_<void> { |
| 1092 | using Type = kj::Maybe<kj::Exception>; |
| 1093 | }; |
| 1094 | template <typename T> |
| 1095 | using ExceptionOr = ExceptionOr_<T>::Type; |
| 1096 | |
| 1097 | template <typename T, typename InputLockOrMaybeCriticalSection, typename Func> |
| 1098 | jsg::PromiseForResult<Func, T, true> awaitIoImpl( |
| 1099 | jsg::Lock& js, kj::Promise<T> promise, InputLockOrMaybeCriticalSection ilOrCs, Func&& func); |
| 1100 | |
| 1101 | // The IncomingRequest that is currently considered "current". This is always the |
| 1102 | // latest-starting request that hasn't yet completed. |
| 1103 | // |
| 1104 | // For stateless requests, there is only ever one IncomingRequest per IoContext. For |
| 1105 | // actors, there is one IoContext per actor, and each incoming request to the actor |
| 1106 | // creates a new IncomingRequest. |
| 1107 | // |
| 1108 | // The current request is tracked for metrics, logging, and tracing purposes. Any resource |
| 1109 | // usage on the part of the actor, including outgoing subrequests, is attributed to the current |
| 1110 | // request for logging and tracing. This is a hack, we don't actually know which request |
| 1111 | // "caused" any particular resource usage, so this is merely our best guess. |
| 1112 | // |
| 1113 | // The IoChannelFactory must also be accessed through the currentIncomingRequest because it has |
| 1114 | // some tracing context built in. |
| 1115 | IncomingRequest& getCurrentIncomingRequest() { |
| 1116 | KJ_REQUIRE(!incomingRequests.empty(), "the IoContext has no current IncomingRequest", |
| 1117 | lastDeliveredLocation); |
| 1118 | return incomingRequests.front(); |
| 1119 | } |
| 1120 | |
| 1121 | // Run the given callback within the scope of this IoContext. This encapsulates the |
| 1122 | // setup of a number of scopes that must be entered prior to running within the |
| 1123 | // context, including entering the V8StackScope and acquiring the Worker::Lock. |
| 1124 | void runInContextScope(Worker::LockType lockType, |
| 1125 | kj::Maybe<InputGate::Lock> inputLock, |
| 1126 | kj::Function<void(Worker::Lock&)> func); |
| 1127 | |
| 1128 | kj::Promise<void> deleteQueueSignalTask; |
| 1129 | static kj::Promise<void> startDeleteQueueSignalTask(IoContext* context); |
| 1130 | |
| 1131 | friend class Finalizeable; |
| 1132 | friend class DeleteQueue; |
| 1133 | template <typename T> |
| 1134 | friend kj::Promise<ExceptionOr<T>> promiseForExceptionOrT(kj::Promise<T> promise); |
| 1135 | template <typename Result> |
| 1136 | friend Result throwOrReturnResult( |
| 1137 | jsg::Lock& js, IoContext::ExceptionOr<Result>&& exceptionOrResult); |
| 1138 | }; |
| 1139 | |
| 1140 | // The SuppressIoContextScope utility is used to temporarily suppress the active IoContext |
| 1141 | // on the current thread while it is in scope. |
| 1142 | struct SuppressIoContextScope { |
| 1143 | IoContext* cached; |
| 1144 | SuppressIoContextScope(); |
| 1145 | ~SuppressIoContextScope() noexcept(false); |
| 1146 | KJ_DISALLOW_COPY_AND_MOVE(SuppressIoContextScope); |
| 1147 | }; |
| 1148 | |
| 1149 | // ======================================================================================= |
| 1150 | // inline implementation details |
| 1151 | |
| 1152 | template <typename T> |
| 1153 | kj::Promise<T> IoContext::lockOutputWhile(kj::Promise<T> promise) { |
| 1154 | return getActorOrThrow().getOutputGate().lockWhile(kj::mv(promise), getCurrentTraceSpan()); |
| 1155 | } |
| 1156 | |
| 1157 | template <typename Func> |
| 1158 | kj::PromiseForResult<Func, Worker::Lock&> IoContext::run( |
| 1159 | Func&& func, kj::Maybe<kj::Own<InputGate::CriticalSection>> criticalSection) { |
| 1160 | KJ_IF_SOME(cs, criticalSection) { |
| 1161 | return cs.get() |
| 1162 | ->wait(getCurrentTraceSpan()) |
| 1163 | .then([this, func = kj::fwd<Func>(func)](InputGate::Lock&& inputLock) mutable { |
| 1164 | return run(kj::fwd<Func>(func), kj::mv(inputLock)); |
| 1165 | }); |
| 1166 | } else { |
| 1167 | return run(kj::fwd<Func>(func)); |
| 1168 | } |
| 1169 | } |
| 1170 | |
| 1171 | template <typename Func> |
| 1172 | kj::PromiseForResult<Func, Worker::Lock&> IoContext::run( |
| 1173 | Func&& func, kj::Maybe<InputGate::Lock> inputLock) { |
| 1174 | // Before we try running anything, let's make sure our IoContext hasn't been aborted. If it has |
| 1175 | // been aborted, there's likely not an active request so later operations will fail anyway. |
| 1176 | KJ_IF_SOME(ex, abortException) { |
| 1177 | return ex.clone(); |
| 1178 | } |
| 1179 | |
| 1180 | kj::Promise<Worker::AsyncLock> asyncLockPromise = nullptr; |
| 1181 | KJ_IF_SOME(a, actor) { |
| 1182 | if (inputLock == kj::none) { |
| 1183 | return a.getInputGate() |
| 1184 | .wait(getCurrentTraceSpan()) |
| 1185 | .then([this, func = kj::fwd<Func>(func)](InputGate::Lock&& inputLock) mutable { |
| 1186 | return run(kj::fwd<Func>(func), kj::mv(inputLock)); |
| 1187 | }); |
| 1188 | } |
| 1189 | |
| 1190 | asyncLockPromise = worker->takeAsyncLockWhenActorCacheReady(now(), a, getMetrics()); |
| 1191 | } else { |
| 1192 | asyncLockPromise = worker->takeAsyncLock(getMetrics()); |
| 1193 | } |
| 1194 | |
| 1195 | return asyncLockPromise.then([this, inputLock = kj::mv(inputLock), func = kj::fwd<Func>(func)]( |
| 1196 | Worker::AsyncLock lock) mutable { |
| 1197 | using Result = decltype(func(kj::instance<Worker::Lock&>())); |
| 1198 | |
| 1199 | if constexpr (kj::isSameType<Result, void>()) { |
| 1200 | struct RunnableImpl: public Runnable { |
| 1201 | Func func; |
| 1202 | |
| 1203 | RunnableImpl(Func&& func): func(kj::fwd<Func>(func)) {} |
| 1204 | void run(Worker::Lock& lock) override { |
| 1205 | func(lock); |
| 1206 | } |
| 1207 | }; |
| 1208 | |
| 1209 | RunnableImpl runnable(kj::fwd<Func>(func)); |
| 1210 | runImpl(runnable, lock, kj::mv(inputLock), Runnable::Exceptional(false)); |
| 1211 | } else { |
| 1212 | struct RunnableImpl: public Runnable { |
| 1213 | Func func; |
| 1214 | kj::Maybe<Result> result; |
| 1215 | |
| 1216 | RunnableImpl(Func&& func): func(kj::fwd<Func>(func)) {} |
| 1217 | void run(Worker::Lock& lock) override { |
| 1218 | result = func(lock); |
| 1219 | } |
| 1220 | }; |
| 1221 | |
| 1222 | RunnableImpl runnable{kj::fwd<Func>(func)}; |
| 1223 | runImpl(runnable, lock, kj::mv(inputLock), Runnable::Exceptional(false)); |
| 1224 | KJ_IF_SOME(r, runnable.result) { |
| 1225 | return kj::mv(r); |
| 1226 | } else { |
| 1227 | KJ_UNREACHABLE; |
| 1228 | } |
| 1229 | } |
| 1230 | }); |
| 1231 | } |
| 1232 | |
| 1233 | template <typename T, typename Func> |
| 1234 | jsg::PromiseForResult<Func, T, true> IoContext::awaitIo( |
| 1235 | jsg::Lock& js, kj::Promise<T> promise, Func&& func) { |
| 1236 | return awaitIoImpl( |
| 1237 | js, promise.attach(registerPendingEvent()), getCriticalSection(), kj::fwd<Func>(func)); |
| 1238 | } |
| 1239 | |
| 1240 | template <typename T> |
| 1241 | jsg::Promise<T> IoContext::awaitIo(jsg::Lock& js, kj::Promise<T> promise) { |
| 1242 | return awaitIoImpl( |
| 1243 | js, promise.attach(registerPendingEvent()), getCriticalSection(), IdentityFunc<T>()); |
| 1244 | } |
| 1245 | |
| 1246 | template <typename T, typename Func> |
| 1247 | jsg::PromiseForResult<Func, T, true> IoContext::awaitIoWithInputLock( |
| 1248 | jsg::Lock& js, kj::Promise<T> promise, Func&& func) { |
| 1249 | return awaitIoImpl( |
| 1250 | js, promise.attach(registerPendingEvent()), getInputLock(), kj::fwd<Func>(func)); |
| 1251 | } |
| 1252 | |
| 1253 | template <typename T> |
| 1254 | jsg::Promise<T> IoContext::awaitIoWithInputLock(jsg::Lock& js, kj::Promise<T> promise) { |
| 1255 | return awaitIoImpl(js, promise.attach(registerPendingEvent()), getInputLock(), IdentityFunc<T>()); |
| 1256 | } |
| 1257 | |
| 1258 | template <typename T> |
| 1259 | jsg::Promise<T> IoContext::awaitIoLegacy(jsg::Lock& js, kj::Promise<T> promise) { |
| 1260 | return awaitIoImpl(js, kj::mv(promise), getCriticalSection(), IdentityFunc<T>()); |
| 1261 | } |
| 1262 | |
| 1263 | template <typename T> |
| 1264 | jsg::Promise<T> IoContext::awaitIoLegacyWithInputLock(jsg::Lock& js, kj::Promise<T> promise) { |
| 1265 | return awaitIoImpl(js, kj::mv(promise), getInputLock(), IdentityFunc<T>()); |
| 1266 | } |
| 1267 | |
| 1268 | // To reduce the code size impact of awaitIoImpl, move promise continuation code out of |
| 1269 | // awaitIoImpl() where possible. This way, the then() parameters are only templated based on one |
| 1270 | // type each. |
| 1271 | template <typename T> |
| 1272 | kj::Promise<IoContext::ExceptionOr<T>> promiseForExceptionOrT(kj::Promise<T> promise) { |
| 1273 | if constexpr (jsg::isVoid<T>()) { |
| 1274 | return promise.then([]() -> IoContext::ExceptionOr<T> { return kj::none; }, |
| 1275 | [](kj::Exception&& exception) -> IoContext::ExceptionOr<T> { return kj::mv(exception); }); |
| 1276 | } else { |
| 1277 | return promise.then([](T&& result) -> IoContext::ExceptionOr<T> { return kj::mv(result); }, |
| 1278 | [](kj::Exception&& exception) -> IoContext::ExceptionOr<T> { return kj::mv(exception); }); |
| 1279 | } |
| 1280 | }; |
| 1281 | |
| 1282 | template <typename Result> |
| 1283 | Result throwOrReturnResult(jsg::Lock& js, IoContext::ExceptionOr<Result>&& exceptionOrResult) { |
| 1284 | if constexpr (jsg::isVoid<Result>()) { |
| 1285 | KJ_IF_SOME(e, exceptionOrResult) { |
| 1286 | // Now that we're in a promise continuation, we can convert the error and get a good stack |
| 1287 | // trace. |
| 1288 | js.throwException(kj::mv(e)); |
| 1289 | } |
| 1290 | } else { |
| 1291 | KJ_SWITCH_ONEOF(exceptionOrResult) { |
| 1292 | KJ_CASE_ONEOF(e, kj::Exception) { |
| 1293 | // Now that we're in a promise continuation, we can convert the error and get a good stack |
| 1294 | // trace. |
| 1295 | js.throwException(kj::mv(e)); |
| 1296 | } |
| 1297 | KJ_CASE_ONEOF(result, Result) { |
| 1298 | return kj::mv(result); |
| 1299 | } |
| 1300 | } |
| 1301 | KJ_UNREACHABLE; |
| 1302 | } |
| 1303 | }; |
| 1304 | |
| 1305 | template <typename T, typename InputLockOrMaybeCriticalSection, typename Func> |
| 1306 | jsg::PromiseForResult<Func, T, true> IoContext::awaitIoImpl( |
| 1307 | jsg::Lock& js, kj::Promise<T> promise, InputLockOrMaybeCriticalSection ilOrCs, Func&& func) { |
| 1308 | // WARNING: The fact that `promise` has been passed by value whereas `func` is by reference is |
| 1309 | // actually important, because this means that if we throw an exception here in the function |
| 1310 | // body, `promise` will be destroyed first, before `func`. That's important as often `func` |
| 1311 | // holds ownership of objects that `promise` depends on. |
| 1312 | |
| 1313 | requireCurrent(); |
| 1314 | |
| 1315 | // `T` is the type produced by the input promise. `Result` is the type of the final output |
| 1316 | // promise. `Func` transforms from `T` to `Result`. |
| 1317 | using Result = jsg::ReturnType<Func, T, true>; |
| 1318 | |
| 1319 | // It is necessary for us to grab a reference to the jsg::AsyncContextFrame here |
| 1320 | // and pass it into the then(). If the promise is rejected, and there is no rejection |
| 1321 | // handler attached to it, an unhandledrejection event will be scheduled, and scheduling |
| 1322 | // that event needs to be done within the appropriate frame to propagate the correct context. |
| 1323 | |
| 1324 | // We need to catch exceptions from KJ and merge them into the result, so that they can propagate |
| 1325 | // to JavaScript. |
| 1326 | kj::Promise<ExceptionOr<T>> promiseExceptionOrT = promiseForExceptionOrT(kj::mv(promise)); |
| 1327 | |
| 1328 | // Reminder: This can throw JsExceptionThrown if the execution context has been terminated. |
| 1329 | // it's important in that case that `promiseExceptionOrT` will be destroyed before `func`. |
| 1330 | auto [jsPromise, resolver] = js.newPromiseAndResolver<ExceptionOr<Result>>(); |
| 1331 | |
| 1332 | addTask(promiseExceptionOrT.then( |
| 1333 | [this, resolver = kj::mv(resolver), ilOrCs = kj::mv(ilOrCs), |
| 1334 | maybeAsyncContext = jsg::AsyncContextFrame::currentRef(js), |
| 1335 | // Reminder: It's important that `func` gets attached to the promise before the whole |
| 1336 | // thing is passed to `addTask()`, so that it's impossible for `func` to be destroyed |
| 1337 | // before the inner promise. |
| 1338 | func = kj::fwd<Func>(func)](ExceptionOr<T>&& exceptionOrT) mutable { |
| 1339 | struct FuncResultPair { |
| 1340 | // It's important that `exceptionOrT` is destroyed before `Func`. Lambda captures are |
| 1341 | // destroyed in unspecified order, so we wrap them in a struct to make it explicit. |
| 1342 | Func func; |
| 1343 | ExceptionOr<T> exceptionOrT; |
| 1344 | }; |
| 1345 | |
| 1346 | return run( |
| 1347 | [resolver = kj::mv(resolver), |
| 1348 | funcResultPair = FuncResultPair{kj::fwd<Func>(func), kj::mv(exceptionOrT)}, |
| 1349 | maybeAsyncContext = kj::mv(maybeAsyncContext)](Worker::Lock& lock) mutable { |
| 1350 | jsg::AsyncContextFrame::Scope asyncScope(lock, maybeAsyncContext); |
| 1351 | jsg::Lock& js = lock; |
| 1352 | |
| 1353 | if constexpr (jsg::isVoid<T>()) { |
| 1354 | KJ_IF_SOME(e, funcResultPair.exceptionOrT) { |
| 1355 | // We don't use `resolver.reject()` here because if we convert the kj::Exception into |
| 1356 | // a JS Error here, it won't have a useful stack trace. V8 can generate a good stack |
| 1357 | // trace as long as we construct the Error inside of a promise continuation, so we use |
| 1358 | // a `.then()` below that actually extracts the kj::Exception and turn it into a JS |
| 1359 | // Error. |
| 1360 | resolver.resolve(js, kj::mv(e)); |
| 1361 | } else { |
| 1362 | try { |
| 1363 | js.tryCatch([&]() { |
| 1364 | if constexpr (jsg::isVoid<Result>()) { |
| 1365 | funcResultPair.func(js); |
| 1366 | resolver.resolve(js, kj::none); |
| 1367 | } else { |
| 1368 | resolver.resolve(js, funcResultPair.func(js)); |
| 1369 | } |
| 1370 | }, [&](jsg::Value error) { |
| 1371 | // Here we can just `resolver.reject` because we already have a JS exception. |
| 1372 | resolver.reject(js, error.getHandle(js)); |
| 1373 | }); |
| 1374 | } catch (jsg::JsExceptionThrown&) { |
| 1375 | // An uncatchable JS exception -- presumably, the isolate has been terminated. We |
| 1376 | // can't convert this into a promise rejection, we need to just propagate it up. |
| 1377 | throw; |
| 1378 | } catch (...) { |
| 1379 | // Again, pass along the KJ exception so we can convert it later in the right context. |
| 1380 | resolver.resolve(js, kj::getCaughtExceptionAsKj()); |
| 1381 | } |
| 1382 | } |
| 1383 | } else { |
| 1384 | // T is not void. |
| 1385 | KJ_SWITCH_ONEOF(funcResultPair.exceptionOrT) { |
| 1386 | KJ_CASE_ONEOF(exception, kj::Exception) { |
| 1387 | // Again, pass along the KJ exception so we can convert it later in the right context. |
| 1388 | resolver.resolve(js, kj::mv(exception)); |
| 1389 | } |
| 1390 | KJ_CASE_ONEOF(result, T) { |
| 1391 | try { |
| 1392 | js.tryCatch([&]() { |
| 1393 | if constexpr (jsg::isVoid<Result>()) { |
| 1394 | funcResultPair.func(js, kj::mv(result)); |
| 1395 | resolver.resolve(js, kj::none); |
| 1396 | } else { |
| 1397 | // Here we can just `resolver.reject` because we already have a JS exception. |
| 1398 | resolver.resolve(js, funcResultPair.func(js, kj::mv(result))); |
| 1399 | } |
| 1400 | }, [&](jsg::Value error) { resolver.reject(js, error.getHandle(js)); }); |
| 1401 | } catch (jsg::JsExceptionThrown&) { |
| 1402 | // An uncatchable JS exception -- presumably, the isolate has been terminated. We |
| 1403 | // can't convert this into a promise rejection, we need to just propagate it up. |
| 1404 | throw; |
| 1405 | } catch (...) { |
| 1406 | // Again, pass along the KJ exception so we can convert it later in the right context. |
| 1407 | resolver.resolve(js, kj::getCaughtExceptionAsKj()); |
| 1408 | } |
| 1409 | } |
| 1410 | } |
| 1411 | } |
| 1412 | }, |
| 1413 | kj::mv(ilOrCs)); |
| 1414 | })); |
| 1415 | |
| 1416 | // Reminder: This can throw JsExceptionThrown if the execution context has been terminated. We |
| 1417 | // have already disowned `promise` and `func` by this point, though, so teardown order is no |
| 1418 | // longer our concern. |
| 1419 | return jsPromise.then(js, throwOrReturnResult<Result>); |
| 1420 | } |
| 1421 | |
| 1422 | template <typename T> |
| 1423 | kj::_::ReducePromises<RemoveIoOwn<T>> IoContext::awaitJs(jsg::Lock& js, jsg::Promise<T> jsPromise) { |
| 1424 | auto paf = kj::newPromiseAndFulfiller<RemoveIoOwn<T>>(); |
| 1425 | struct RefcountedFulfiller: public kj::Refcounted { |
| 1426 | kj::Own<kj::PromiseFulfiller<RemoveIoOwn<T>>> fulfiller; |
| 1427 | kj::Own<const AtomicWeakRef<Worker::Isolate>> maybeIsolate; |
| 1428 | bool isDone = false; |
| 1429 | |
| 1430 | RefcountedFulfiller(kj::Own<const AtomicWeakRef<Worker::Isolate>> maybeIsolate, |
| 1431 | kj::Own<kj::PromiseFulfiller<RemoveIoOwn<T>>> fulfiller) |
| 1432 | : fulfiller(kj::mv(fulfiller)), |
| 1433 | maybeIsolate(kj::mv(maybeIsolate)) {} |
| 1434 | |
| 1435 | ~RefcountedFulfiller() noexcept(false) { |
| 1436 | if (!isDone) { |
| 1437 | reject(); |
| 1438 | } |
| 1439 | } |
| 1440 | |
| 1441 | private: |
| 1442 | void reject() { |
| 1443 | // We use a weak isolate reference here in case the isolate gets dropped before this code |
| 1444 | // is executed. In that case we default to `false` as we cannot access the original isolate. |
| 1445 | auto hasExcessivelyExceededHeapLimit = maybeIsolate->tryAddStrongRef() |
| 1446 | .map([](kj::Own<const Worker::Isolate> isolate) { |
| 1447 | return isolate->getLimitEnforcer().hasExcessivelyExceededHeapLimit(); |
| 1448 | }).orDefault(false); |
| 1449 | if (hasExcessivelyExceededHeapLimit) { |
| 1450 | auto e = JSG_KJ_EXCEPTION(OVERLOADED, Error, "Worker has exceeded memory limit."); |
| 1451 | e.setDetail(MEMORY_LIMIT_DETAIL_ID, kj::heapArray<kj::byte>(0)); |
| 1452 | fulfiller->reject(kj::mv(e)); |
| 1453 | } else { |
| 1454 | // The JavaScript resolver was garbage collected, i.e. JavaScript will never resolve |
| 1455 | // this promise. |
| 1456 | fulfiller->reject(JSG_KJ_EXCEPTION(FAILED, Error, "Promise will never complete.")); |
| 1457 | } |
| 1458 | } |
| 1459 | }; |
| 1460 | auto& isolate = Worker::Isolate::from(js); |
| 1461 | auto fulfiller = kj::refcounted<RefcountedFulfiller>(isolate.getWeakRef(), kj::mv(paf.fulfiller)); |
| 1462 | |
| 1463 | auto errorHandler = [fulfiller = addObject(kj::addRef(*fulfiller))]( |
| 1464 | jsg::Lock& js, jsg::Value jsExceptionRef) mutable { |
| 1465 | // Note: `context` can possibly be different than the one that started the wait, if the |
| 1466 | // promise resolved from a different context. In that case the use of `fulfiller` will |
| 1467 | // throw later on. But it's OK to use the wrong context up until that point. |
| 1468 | auto& context = IoContext::current(); |
| 1469 | |
| 1470 | auto isolate = context.getCurrentLock().getIsolate(); |
| 1471 | auto jsException = jsExceptionRef.getHandle(js); |
| 1472 | |
| 1473 | // TODO(someday): We log an "uncaught exception" here whenever a promise returned from JS to |
| 1474 | // C++ rejects. However, the C++ code waiting on the promise may do its own logging (e.g. |
| 1475 | // event.respondWith() does), in which case this is redundant. But, it's difficult to be |
| 1476 | // sure that all C++ consumers log properly, and even if they do, the stack trace is lost |
| 1477 | // once the exception has been tunneled into a KJ exception, so the later logging won't be |
| 1478 | // as useful. We should improve the tunneling to include stack traces and ensure that all |
| 1479 | // consumers do in fact log exceptions, then we can remove this. |
| 1480 | context.logUncaughtException( |
| 1481 | UncaughtExceptionSource::INTERNAL_ASYNC, jsg::JsValue(jsException)); |
| 1482 | |
| 1483 | fulfiller->fulfiller->reject(jsg::createTunneledException(isolate, jsException)); |
| 1484 | fulfiller->isDone = true; |
| 1485 | }; |
| 1486 | |
| 1487 | if constexpr (jsg::isVoid<T>()) { |
| 1488 | jsPromise.then(js, [fulfiller = addObject(kj::mv(fulfiller))](jsg::Lock&) mutable { |
| 1489 | fulfiller->fulfiller->fulfill(); |
| 1490 | fulfiller->isDone = true; |
| 1491 | }, kj::mv(errorHandler)); |
| 1492 | } else { |
| 1493 | jsPromise.then(js, [fulfiller = addObject(kj::mv(fulfiller))](jsg::Lock&, T&& result) mutable { |
| 1494 | if constexpr (isIoOwn<T>()) { |
| 1495 | fulfiller->fulfiller->fulfill(kj::mv(*result)); |
| 1496 | } else { |
| 1497 | fulfiller->fulfiller->fulfill(kj::mv(result)); |
| 1498 | } |
| 1499 | fulfiller->isDone = true; |
| 1500 | }, kj::mv(errorHandler)); |
| 1501 | } |
| 1502 | |
| 1503 | return paf.promise.exclusiveJoin(onAbort().then([]() -> RemoveIoOwn<T> { KJ_UNREACHABLE; })); |
| 1504 | } |
| 1505 | |
| 1506 | template <IoContext::TopUpFlag topUp, typename Func> |
| 1507 | auto IoContext::makeReentryCallback(Func func) { |
| 1508 | // A reentry callback is meant for *re-*entry, so should only be created while already inside |
| 1509 | // the IoContext. Initial entry into the IoContext should just use run(). |
| 1510 | requireCurrent(); |
| 1511 | |
| 1512 | // We need to: |
| 1513 | // - Use addTask() to make sure that, if we're in an actor, the IncomingEvent stays alive while |
| 1514 | // the callback exists (and hibernation is blocked). |
| 1515 | // - Call registerPendingEvent() to make sure that, if we're NOT in an actor, we don't conclude |
| 1516 | // that there's nothing left to wait for while the callback exists. |
| 1517 | // TODO(perf): Probably both of these things could be done in simpler ways involving less |
| 1518 | // allocation, but it would require some refactoring. |
| 1519 | auto [promise, fulfiller] = kj::newPromiseAndFulfiller<void>(); |
| 1520 | addTask(kj::mv(promise)); |
| 1521 | auto releaseNotifier = |
| 1522 | kj::defer([fulfiller = kj::mv(fulfiller), pe = registerPendingEvent()]() mutable { |
| 1523 | fulfiller->fulfill(); |
| 1524 | }); |
| 1525 | |
| 1526 | auto ioFunc = addObjectReverse(kj::heap(kj::fwd<Func>(func))); |
| 1527 | |
| 1528 | return [self = getWeakRef(), cs = getCriticalSection(), releaseNotifier = kj::mv(releaseNotifier), |
| 1529 | ioFunc = kj::mv(ioFunc)](auto&&... params) mutable { |
| 1530 | auto& ctx = JSG_REQUIRE_NONNULL(self->tryGet(), Error, |
| 1531 | "The execution context which hosts this callback is no longer running."); |
| 1532 | |
| 1533 | if constexpr (topUp == TOP_UP) { |
| 1534 | ctx.getLimitEnforcer().topUpActor(); |
| 1535 | } |
| 1536 | |
| 1537 | return ctx.canceler.wrap(ctx.run( |
| 1538 | [&ctx, &ioFunc, ... params = kj::fwd<decltype(params)>(params)]( |
| 1539 | Worker::Lock& lock) mutable { |
| 1540 | using ResultType = kj::Decay<decltype(func(lock, kj::fwd<decltype(params)>(params)...))>; |
| 1541 | |
| 1542 | auto& func = *ioFunc; |
| 1543 | |
| 1544 | if constexpr (kj::isSameType<ResultType, void>()) { |
| 1545 | (void)ctx; |
| 1546 | func(lock, kj::fwd<decltype(params)>(params)...); |
| 1547 | } else if constexpr (jsg::isPromise<ResultType>()) { |
| 1548 | return ctx.awaitJs(lock, func(lock, kj::fwd<decltype(params)>(params)...)); |
| 1549 | } else { |
| 1550 | (void)ctx; |
| 1551 | return func(lock, kj::fwd<decltype(params)>(params)...); |
| 1552 | } |
| 1553 | }, |
| 1554 | kj::mv(cs))); |
| 1555 | }; |
| 1556 | } |
| 1557 | |
| 1558 | template <typename T> |
| 1559 | inline IoOwn<T> IoContext::addObject(kj::Own<T> obj) { |
| 1560 | requireCurrent(); |
| 1561 | return deleteQueue.queue->addObject(kj::mv(obj), ownedObjects); |
| 1562 | } |
| 1563 | |
| 1564 | template <typename T> |
| 1565 | inline IoPtr<T> IoContext::addObject(T& obj) { |
| 1566 | requireCurrent(); |
| 1567 | return IoPtr<T>(deleteQueue.queue.addRef(), &obj); |
| 1568 | } |
| 1569 | |
| 1570 | template <typename Func> |
| 1571 | auto IoContext::addFunctor(Func&& func) { |
| 1572 | if constexpr (kj::isReference<Func>()) { |
| 1573 | return [func = addObject(func)]( |
| 1574 | auto&&... params) mutable { return (*func)(kj::fwd<decltype(params)>(params)...); }; |
| 1575 | } else { |
| 1576 | return [func = addObject(kj::heap(kj::mv(func)))]( |
| 1577 | auto&&... params) mutable { return (*func)(kj::fwd<decltype(params)>(params)...); }; |
| 1578 | } |
| 1579 | } |
| 1580 | |
| 1581 | template <typename T> |
| 1582 | inline ReverseIoOwn<T> IoContext::addObjectReverse(kj::Own<T> obj) { |
| 1583 | // We intentionally don't requireCurrent() -- the only requirement is that the caller is in the |
| 1584 | // same thread. |
| 1585 | return deleteQueue.queue->addObjectReverse(getWeakRef(), kj::mv(obj), ownedObjects); |
| 1586 | } |
| 1587 | |
| 1588 | template <typename Func> |
| 1589 | jsg::PromiseForResult<Func, void, true> IoContext::blockConcurrencyWhile( |
| 1590 | jsg::Lock& js, Func&& callback) { |
| 1591 | auto lock = getInputLock(); |
| 1592 | auto cs = lock.startCriticalSection(); |
| 1593 | auto cs2 = kj::addRef(*cs); |
| 1594 | |
| 1595 | using T = jsg::RemovePromise<jsg::ReturnType<Func, void, true>>; |
| 1596 | auto [result, resolver] = js.newPromiseAndResolver<T>(); |
| 1597 | |
| 1598 | addTask( |
| 1599 | cs->wait(getCurrentTraceSpan()) |
| 1600 | .then([this, callback = kj::mv(callback), |
| 1601 | maybeAsyncContext = jsg::AsyncContextFrame::currentRef(js)]( |
| 1602 | InputGate::Lock inputLock) mutable { |
| 1603 | return run( |
| 1604 | [this, callback = kj::mv(callback), maybeAsyncContext = kj::mv(maybeAsyncContext)]( |
| 1605 | Worker::Lock& lock) mutable { |
| 1606 | jsg::AsyncContextFrame::Scope scope(lock, maybeAsyncContext); |
| 1607 | auto cb = kj::mv(callback); |
| 1608 | |
| 1609 | // Remember that this can throw synchronously, and it's important that we catch such throws |
| 1610 | // and call cs->failed(). |
| 1611 | auto promise = cb(lock); |
| 1612 | |
| 1613 | // Arrange to time out if the critical section runs more than 30 seconds, so that objects |
| 1614 | // won't be hung forever if they have a critical section that deadlocks. |
| 1615 | auto timeout = afterLimitTimeout(30 * kj::SECONDS).then([]() -> T { |
| 1616 | kj::throwFatalException(JSG_KJ_EXCEPTION(OVERLOADED, Error, |
| 1617 | "A call to blockConcurrencyWhile() in a Durable Object waited for " |
| 1618 | "too long. The call was canceled and the Durable Object was reset.")); |
| 1619 | }); |
| 1620 | |
| 1621 | return awaitJs(lock, kj::mv(promise)).exclusiveJoin(kj::mv(timeout)); |
| 1622 | }, |
| 1623 | kj::mv(inputLock)); |
| 1624 | }) |
| 1625 | .then( |
| 1626 | [this, cs = kj::mv(cs), resolver = kj::mv(resolver), |
| 1627 | maybeAsyncContext = jsg::AsyncContextFrame::currentRef(js)](T&& value) mutable { |
| 1628 | auto inputLock = cs->succeeded(); |
| 1629 | return run( |
| 1630 | [value = kj::mv(value), resolver = kj::mv(resolver), |
| 1631 | maybeAsyncContext = kj::mv(maybeAsyncContext)](Worker::Lock& lock) mutable { |
| 1632 | jsg::AsyncContextFrame::Scope scope(lock, maybeAsyncContext); |
| 1633 | resolver.resolve(lock, kj::mv(value)); |
| 1634 | }, |
| 1635 | kj::mv(inputLock)); |
| 1636 | }, |
| 1637 | [cs = kj::mv(cs2)](kj::Exception&& e) mutable { |
| 1638 | // Annotate as broken for periodic metrics. |
| 1639 | auto msg = e.getDescription(); |
| 1640 | if (!msg.startsWith("broken."_kj) && !msg.startsWith("remote.broken."_kj)) { |
| 1641 | // If we already set up a brokenness reason, we shouldn't override it. |
| 1642 | |
| 1643 | auto description = jsg::annotateBroken(msg, "broken.inputGateBroken"); |
| 1644 | e.setDescription(kj::mv(description)); |
| 1645 | } |
| 1646 | |
| 1647 | // Note that on failure, no further InputLocks will be obtainable and the actor will |
| 1648 | // shut down, so don't worry about holding a lock until we get back to application code -- |
| 1649 | // we won't! In fact, we don't even bother calling resolver.reject() because it's meaningless |
| 1650 | // at this point. |
| 1651 | cs->failed(e); |
| 1652 | |
| 1653 | kj::throwFatalException(kj::mv(e)); |
| 1654 | })); |
| 1655 | |
| 1656 | return kj::mv(result); |
| 1657 | } |
| 1658 | |
| 1659 | } // namespace workerd |