File
Blob: src/workerd/io/worker-entrypoint.c++
| 1 | // Copyright (c) 2017-2022 Cloudflare, Inc. |
| 2 | // Licensed under the Apache 2.0 license found in the LICENSE file or at: |
| 3 | // https://opensource.org/licenses/Apache-2.0 |
| 4 | |
| 5 | #include "worker-entrypoint.h" |
| 6 | |
| 7 | #include <workerd/api/basics.h> |
| 8 | #include <workerd/api/global-scope.h> |
| 9 | #include <workerd/api/util.h> |
| 10 | #include <workerd/io/features.h> |
| 11 | #include <workerd/io/io-context.h> |
| 12 | #include <workerd/io/limit-enforcer.h> |
| 13 | #include <workerd/io/tracer.h> |
| 14 | #include <workerd/jsg/jsg.h> |
| 15 | #include <workerd/util/http-util.h> |
| 16 | #include <workerd/util/sentry.h> |
| 17 | #include <workerd/util/strings.h> |
| 18 | #include <workerd/util/thread-scopes.h> |
| 19 | #include <workerd/util/uncaught-exception-source.h> |
| 20 | #include <workerd/util/use-perfetto-categories.h> |
| 21 | |
| 22 | #include <capnp/message.h> |
| 23 | #include <kj/compat/http.h> |
| 24 | |
| 25 | namespace workerd { |
| 26 | |
| 27 | namespace { |
| 28 | // Wrapper around a Worker that handles receiving a new event from the outside. In particular, |
| 29 | // this handles: |
| 30 | // - Creating a IoContext and making it current. |
| 31 | // - Executing the worker under lock. |
| 32 | // - Catching exceptions and converting them to HTTP error responses. |
| 33 | // - Or, falling back to proxying if passThroughOnException() was used. |
| 34 | // - Finish waitUntil() tasks. |
| 35 | class WorkerEntrypoint final: public WorkerInterface { |
| 36 | public: |
| 37 | // Call this instead of the constructor. It actually adds a wrapper object around the |
| 38 | // `WorkerEntrypoint`, but the wrapper still implements `WorkerInterface`. |
| 39 | // |
| 40 | // WorkerEntrypoint will create a IoContext, and that IoContext may outlive the |
| 41 | // WorkerEntrypoint by means of a waitUntil() task. Any object(s) which must be kept alive to |
| 42 | // support the worker for the lifetime of the IoContext (e.g., subsequent pipeline stages) |
| 43 | // must be passed in via `ioContextDependency`. |
| 44 | // |
| 45 | // If this is NOT a zone worker, then `zoneDefaultWorkerLimits` should be a default instance of |
| 46 | // WorkerLimits::Reader. Hence this is not necessarily the same as |
| 47 | // topLevelRequest.getZoneDefaultWorkerLimits(), since the top level request may be shared between |
| 48 | // zone and non-zone workers. |
| 49 | static kj::Own<WorkerInterface> construct(ThreadContext& threadContext, |
| 50 | kj::Own<const Worker> worker, |
| 51 | kj::Maybe<kj::StringPtr> entrypointName, |
| 52 | Frankenvalue props, |
| 53 | kj::Maybe<kj::Own<Worker::Actor>> actor, |
| 54 | kj::Own<LimitEnforcer> limitEnforcer, |
| 55 | kj::Own<void> ioContextDependency, |
| 56 | kj::Own<IoChannelFactory> ioChannelFactory, |
| 57 | kj::Own<RequestObserver> metrics, |
| 58 | kj::TaskSet& waitUntilTasks, |
| 59 | bool tunnelExceptions, |
| 60 | kj::Maybe<kj::Own<BaseTracer>> workerTracer, |
| 61 | kj::Maybe<kj::String> cfBlobJson, |
| 62 | kj::Maybe<Worker::VersionInfo> versionInfo, |
| 63 | kj::Maybe<tracing::InvocationSpanContext> maybeTriggerInvocationSpan, |
| 64 | bool isDynamicDispatch); |
| 65 | |
| 66 | kj::Promise<void> request(kj::HttpMethod method, |
| 67 | kj::StringPtr url, |
| 68 | const kj::HttpHeaders& headers, |
| 69 | kj::AsyncInputStream& requestBody, |
| 70 | Response& response) override; |
| 71 | kj::Promise<void> connect(kj::StringPtr host, |
| 72 | const kj::HttpHeaders& headers, |
| 73 | kj::AsyncIoStream& connection, |
| 74 | ConnectResponse& response, |
| 75 | kj::HttpConnectSettings settings) override; |
| 76 | kj::Promise<void> prewarm(kj::StringPtr url) override; |
| 77 | kj::Promise<ScheduledResult> runScheduled(kj::Date scheduledTime, kj::StringPtr cron) override; |
| 78 | kj::Promise<AlarmResult> runAlarm(kj::Date scheduledTime, uint32_t retryCount) override; |
| 79 | kj::Promise<kj::Maybe<kj::Date>> abandonAlarm(kj::Date scheduledTime) override; |
| 80 | kj::Promise<bool> test() override; |
| 81 | kj::Promise<CustomEvent::Result> customEvent(kj::Own<CustomEvent> event) override; |
| 82 | |
| 83 | private: |
| 84 | class ResponseSentTracker; |
| 85 | |
| 86 | // Members initialized at startup. |
| 87 | |
| 88 | ThreadContext& threadContext; |
| 89 | kj::TaskSet& waitUntilTasks; |
| 90 | kj::Maybe<kj::Own<IoContext::IncomingRequest>> incomingRequest; |
| 91 | bool tunnelExceptions; |
| 92 | bool isDynamicDispatch; |
| 93 | kj::Maybe<kj::StringPtr> entrypointName; |
| 94 | Frankenvalue props; |
| 95 | kj::Maybe<kj::String> cfBlobJson; |
| 96 | kj::Maybe<Worker::VersionInfo> versionInfo; |
| 97 | |
| 98 | // Hacky members used to hold some temporary state while processing a request. |
| 99 | // See gory details in WorkerEntrypoint::request(). |
| 100 | |
| 101 | kj::Maybe<kj::Promise<void>> proxyTask; |
| 102 | kj::Maybe<kj::Own<WorkerInterface>> failOpenService; |
| 103 | bool loggedExceptionEarlier = false; |
| 104 | kj::Maybe<jsg::Ref<api::AbortController>> abortController; |
| 105 | |
| 106 | void init(kj::Own<const Worker> worker, |
| 107 | kj::Maybe<kj::Own<Worker::Actor>> actor, |
| 108 | kj::Own<LimitEnforcer> limitEnforcer, |
| 109 | kj::Own<void> ioContextDependency, |
| 110 | kj::Own<IoChannelFactory> ioChannelFactory, |
| 111 | kj::Own<RequestObserver> metrics, |
| 112 | kj::Maybe<kj::Own<BaseTracer>> workerTracer, |
| 113 | kj::Maybe<tracing::InvocationSpanContext> maybeTriggerInvocationSpan); |
| 114 | |
| 115 | template <typename T> |
| 116 | kj::Promise<T> maybeAddGcPassForTest(IoContext& context, kj::Promise<T> promise); |
| 117 | |
| 118 | kj::Promise<WorkerEntrypoint::AlarmResult> runAlarmImpl( |
| 119 | kj::Own<IoContext::IncomingRequest> incomingRequest, |
| 120 | kj::Date scheduledTime, |
| 121 | uint32_t retryCount); |
| 122 | |
| 123 | public: // For kj::heap() only; pretend this is private. |
| 124 | WorkerEntrypoint(kj::Badge<WorkerEntrypoint> badge, |
| 125 | ThreadContext& threadContext, |
| 126 | kj::TaskSet& waitUntilTasks, |
| 127 | bool tunnelExceptions, |
| 128 | bool isDynamicDispatch, |
| 129 | kj::Maybe<kj::StringPtr> entrypointName, |
| 130 | Frankenvalue props, |
| 131 | kj::Maybe<kj::String> cfBlobJson, |
| 132 | kj::Maybe<Worker::VersionInfo> versionInfo); |
| 133 | }; |
| 134 | |
| 135 | // Simple wrapper around `HttpService::Response` to let us know if the response was sent |
| 136 | // already. |
| 137 | class WorkerEntrypoint::ResponseSentTracker final: public kj::HttpService::Response { |
| 138 | public: |
| 139 | ResponseSentTracker(kj::HttpService::Response& inner): inner(inner) {} |
| 140 | KJ_DISALLOW_COPY_AND_MOVE(ResponseSentTracker); |
| 141 | |
| 142 | bool isSent() const { |
| 143 | return sent; |
| 144 | } |
| 145 | uint getHttpResponseStatus() const { |
| 146 | return httpResponseStatus; |
| 147 | } |
| 148 | |
| 149 | kj::Own<kj::AsyncOutputStream> send(uint statusCode, |
| 150 | kj::StringPtr statusText, |
| 151 | const kj::HttpHeaders& headers, |
| 152 | kj::Maybe<uint64_t> expectedBodySize = kj::none) override { |
| 153 | TRACE_EVENT( |
| 154 | "workerd", "WorkerEntrypoint::ResponseSentTracker::send()", "statusCode", statusCode); |
| 155 | sent = true; |
| 156 | httpResponseStatus = statusCode; |
| 157 | return inner.send(statusCode, statusText, headers, expectedBodySize); |
| 158 | } |
| 159 | |
| 160 | kj::Own<kj::WebSocket> acceptWebSocket(const kj::HttpHeaders& headers) override { |
| 161 | TRACE_EVENT("workerd", "WorkerEntrypoint::ResponseSentTracker::acceptWebSocket()"); |
| 162 | sent = true; |
| 163 | return inner.acceptWebSocket(headers); |
| 164 | } |
| 165 | |
| 166 | private: |
| 167 | uint httpResponseStatus = 0; |
| 168 | kj::HttpService::Response& inner; |
| 169 | bool sent = false; |
| 170 | }; |
| 171 | |
| 172 | kj::Own<WorkerInterface> WorkerEntrypoint::construct(ThreadContext& threadContext, |
| 173 | kj::Own<const Worker> worker, |
| 174 | kj::Maybe<kj::StringPtr> entrypointName, |
| 175 | Frankenvalue props, |
| 176 | kj::Maybe<kj::Own<Worker::Actor>> actor, |
| 177 | kj::Own<LimitEnforcer> limitEnforcer, |
| 178 | kj::Own<void> ioContextDependency, |
| 179 | kj::Own<IoChannelFactory> ioChannelFactory, |
| 180 | kj::Own<RequestObserver> metrics, |
| 181 | kj::TaskSet& waitUntilTasks, |
| 182 | bool tunnelExceptions, |
| 183 | kj::Maybe<kj::Own<BaseTracer>> workerTracer, |
| 184 | kj::Maybe<kj::String> cfBlobJson, |
| 185 | kj::Maybe<Worker::VersionInfo> versionInfo, |
| 186 | kj::Maybe<tracing::InvocationSpanContext> maybeTriggerInvocationSpan, |
| 187 | bool isDynamicDispatch) { |
| 188 | TRACE_EVENT("workerd", "WorkerEntrypoint::construct()"); |
| 189 | |
| 190 | auto obj = kj::heap<WorkerEntrypoint>(kj::Badge<WorkerEntrypoint>(), threadContext, |
| 191 | waitUntilTasks, tunnelExceptions, isDynamicDispatch, entrypointName, kj::mv(props), |
| 192 | kj::mv(cfBlobJson), kj::mv(versionInfo)); |
| 193 | obj->init(kj::mv(worker), kj::mv(actor), kj::mv(limitEnforcer), kj::mv(ioContextDependency), |
| 194 | kj::mv(ioChannelFactory), kj::addRef(*metrics), kj::mv(workerTracer), |
| 195 | kj::mv(maybeTriggerInvocationSpan)); |
| 196 | auto& wrapper = metrics->wrapWorkerInterface(*obj); |
| 197 | return kj::attachRef(wrapper, kj::mv(obj), kj::mv(metrics)); |
| 198 | } |
| 199 | |
| 200 | WorkerEntrypoint::WorkerEntrypoint(kj::Badge<WorkerEntrypoint> badge, |
| 201 | ThreadContext& threadContext, |
| 202 | kj::TaskSet& waitUntilTasks, |
| 203 | bool tunnelExceptions, |
| 204 | bool isDynamicDispatch, |
| 205 | kj::Maybe<kj::StringPtr> entrypointName, |
| 206 | Frankenvalue props, |
| 207 | kj::Maybe<kj::String> cfBlobJson, |
| 208 | kj::Maybe<Worker::VersionInfo> versionInfo) |
| 209 | : threadContext(threadContext), |
| 210 | waitUntilTasks(waitUntilTasks), |
| 211 | tunnelExceptions(tunnelExceptions), |
| 212 | isDynamicDispatch(isDynamicDispatch), |
| 213 | entrypointName(entrypointName), |
| 214 | props(kj::mv(props)), |
| 215 | cfBlobJson(kj::mv(cfBlobJson)), |
| 216 | versionInfo(kj::mv(versionInfo)) {} |
| 217 | |
| 218 | void WorkerEntrypoint::init(kj::Own<const Worker> worker, |
| 219 | kj::Maybe<kj::Own<Worker::Actor>> actor, |
| 220 | kj::Own<LimitEnforcer> limitEnforcer, |
| 221 | kj::Own<void> ioContextDependency, |
| 222 | kj::Own<IoChannelFactory> ioChannelFactory, |
| 223 | kj::Own<RequestObserver> metrics, |
| 224 | kj::Maybe<kj::Own<BaseTracer>> workerTracer, |
| 225 | kj::Maybe<tracing::InvocationSpanContext> maybeTriggerInvocationSpan) { |
| 226 | TRACE_EVENT("workerd", "WorkerEntrypoint::init()"); |
| 227 | // We need to construct the IoContext -- unless this is an actor and it already has a |
| 228 | // IoContext, in which case we reuse it. |
| 229 | |
| 230 | auto newContext = [&]() { |
| 231 | TRACE_EVENT("workerd", "WorkerEntrypoint::init() create new IoContext"); |
| 232 | auto actorRef = actor.map([](kj::Own<Worker::Actor>& ptr) -> Worker::Actor& { return *ptr; }); |
| 233 | |
| 234 | // Attaching to refcount instance is safe here since this instance stays alive for the lifetime |
| 235 | // of the associated WorkerInterface, other references may be created below for actors requests |
| 236 | // in separate init() calls but this ioContextDependency does not need to live as long as those |
| 237 | // instances. |
| 238 | return kj::refcounted<IoContext>(threadContext, kj::mv(worker), actorRef, kj::mv(limitEnforcer)) |
| 239 | .attachToThisReference(kj::mv(ioContextDependency)); |
| 240 | }; |
| 241 | |
| 242 | kj::Own<IoContext> context; |
| 243 | KJ_IF_SOME(a, actor) { |
| 244 | KJ_IF_SOME(rc, a.get()->getIoContext()) { |
| 245 | context = kj::addRef(rc); |
| 246 | } else { |
| 247 | context = newContext(); |
| 248 | a.get()->setIoContext(kj::addRef(*context)); |
| 249 | } |
| 250 | } else { |
| 251 | context = newContext(); |
| 252 | } |
| 253 | |
| 254 | incomingRequest = kj::heap<IoContext::IncomingRequest>(kj::mv(context), kj::mv(ioChannelFactory), |
| 255 | kj::mv(metrics), kj::mv(workerTracer), kj::mv(maybeTriggerInvocationSpan)) |
| 256 | .attach(kj::mv(actor)); |
| 257 | } |
| 258 | |
| 259 | kj::Exception exceptionToPropagate(bool isInternalException, kj::Exception&& exception) { |
| 260 | if (isInternalException) { |
| 261 | // We've already logged it here, the only thing that matters to the client is that we failed |
| 262 | // due to an internal error. Note that this does not need to be labeled "remote." since jsg |
| 263 | // will sanitize it as an internal error. Note that we use `setDescription()` to preserve |
| 264 | // the exception type for `jsg::exceptionToJs(...)` downstream. |
| 265 | exception.setDescription(kj::str("worker_do_not_log; Request failed due to internal error")); |
| 266 | return kj::mv(exception); |
| 267 | } else { |
| 268 | // We do not care how many remote capnp servers this went through since we are returning |
| 269 | // it to the worker via jsg. |
| 270 | // TODO(someday) We also do this stripping when making the tunneled exception for |
| 271 | // `jsg::isTunneledException(...)`. It would be lovely if we could simply store some type |
| 272 | // instead of `loggedExceptionEarlier`. It would save use some work. |
| 273 | auto description = jsg::stripRemoteExceptionPrefix(exception.getDescription()); |
| 274 | if (!description.startsWith("remote.")) { |
| 275 | // If we already were annotated as remote from some other worker entrypoint, no point |
| 276 | // adding an additional prefix. |
| 277 | exception.setDescription(kj::str("remote.", description)); |
| 278 | } |
| 279 | return kj::mv(exception); |
| 280 | } |
| 281 | } |
| 282 | |
| 283 | kj::Promise<void> WorkerEntrypoint::request(kj::HttpMethod method, |
| 284 | kj::StringPtr url, |
| 285 | const kj::HttpHeaders& headers, |
| 286 | kj::AsyncInputStream& requestBody, |
| 287 | Response& response) { |
| 288 | TRACE_EVENT("workerd", "WorkerEntrypoint::request()", "url", url.cStr(), |
| 289 | PERFETTO_FLOW_FROM_POINTER(this)); |
| 290 | auto incomingRequest = |
| 291 | kj::mv(KJ_REQUIRE_NONNULL(this->incomingRequest, "request() can only be called once")); |
| 292 | this->incomingRequest = kj::none; |
| 293 | auto& context = incomingRequest->getContext(); |
| 294 | |
| 295 | auto wrappedResponse = kj::heap<ResponseSentTracker>(response); |
| 296 | |
| 297 | bool isActor = context.getActor() != kj::none; |
| 298 | // HACK: Capture workerTracer directly, it's unclear how to acquire the right tracer from context |
| 299 | // when we need it (for DOs, IoContext may point to a different WorkerTracer by the time we use |
| 300 | // it). The tracer lives as long or longer than the IoContext (based on being co-owned |
| 301 | // by IncomingRequest and PipelineTracer) so long enough. |
| 302 | kj::Maybe<BaseTracer&> workerTracer; |
| 303 | |
| 304 | KJ_IF_SOME(t, incomingRequest->getWorkerTracer()) { |
| 305 | kj::String cfJson; |
| 306 | KJ_IF_SOME(c, cfBlobJson) { |
| 307 | cfJson = kj::str(c); |
| 308 | } |
| 309 | |
| 310 | // To match our historical behavior (when we used to pull the headers from the JavaScript |
| 311 | // object later on), we need to canonicalize the headers, including: |
| 312 | // - Lower-case the header name. |
| 313 | // - Combine multiple headers with the same name into a comma-delimited list. (This explicitly |
| 314 | // breaks the Set-Cookie header, incidentally, but should be equivalent for all other |
| 315 | // headers.) |
| 316 | kj::TreeMap<kj::String, kj::Vector<kj::StringPtr>> traceHeaders; |
| 317 | headers.forEach([&](kj::StringPtr name, kj::StringPtr value) { |
| 318 | kj::String lower = toLower(name); |
| 319 | auto& slot = traceHeaders.findOrCreate( |
| 320 | lower, [&]() { return decltype(traceHeaders)::Entry{kj::mv(lower), {}}; }); |
| 321 | slot.add(value); |
| 322 | }); |
| 323 | auto traceHeadersArray = KJ_MAP(entry, traceHeaders) { |
| 324 | return tracing::FetchEventInfo::Header(kj::mv(entry.key), kj::strArray(entry.value, ", ")); |
| 325 | }; |
| 326 | |
| 327 | t.setEventInfo(*incomingRequest, |
| 328 | tracing::FetchEventInfo(method, kj::str(url), kj::mv(cfJson), kj::mv(traceHeadersArray))); |
| 329 | workerTracer = t; |
| 330 | } |
| 331 | |
| 332 | incomingRequest->delivered(); |
| 333 | |
| 334 | auto metricsForCatch = kj::addRef(incomingRequest->getMetrics()); |
| 335 | auto metricsForProxyTask = kj::addRef(incomingRequest->getMetrics()); |
| 336 | |
| 337 | TRACE_EVENT_BEGIN("workerd", "WorkerEntrypoint::request() waiting on context", |
| 338 | PERFETTO_TRACK_FROM_POINTER(&context), PERFETTO_FLOW_FROM_POINTER(this)); |
| 339 | |
| 340 | return context |
| 341 | .run([this, &context, method, url, &headers, &requestBody, |
| 342 | &metrics = incomingRequest->getMetrics(), &wrappedResponse = *wrappedResponse, |
| 343 | entrypointName = entrypointName](Worker::Lock& lock) mutable { |
| 344 | TRACE_EVENT_END("workerd", PERFETTO_TRACK_FROM_POINTER(&context)); |
| 345 | TRACE_EVENT("workerd", "WorkerEntrypoint::request() run", PERFETTO_FLOW_FROM_POINTER(this)); |
| 346 | jsg::AsyncContextFrame::StorageScope traceScope = context.makeAsyncTraceScope(lock); |
| 347 | jsg::AsyncContextFrame::StorageScope userTraceScope = context.makeUserAsyncTraceScope(lock); |
| 348 | auto featureFlags = FeatureFlags::get(lock); |
| 349 | |
| 350 | kj::Maybe<jsg::Ref<api::AbortSignal>> signal; |
| 351 | |
| 352 | if (featureFlags.getEnableRequestSignal()) { |
| 353 | auto abortSignalFlag = featureFlags.getRequestSignalPassthrough() |
| 354 | ? api::AbortSignal::Flag::NONE |
| 355 | : api::AbortSignal::Flag::IGNORE_FOR_SUBREQUESTS; |
| 356 | jsg::Lock& js = lock; |
| 357 | signal.emplace(abortController.emplace(js.alloc<api::AbortController>(js, abortSignalFlag)) |
| 358 | ->getSignal()); |
| 359 | } |
| 360 | |
| 361 | return lock.getGlobalScope().request(method, url, headers, requestBody, wrappedResponse, |
| 362 | cfBlobJson, lock, |
| 363 | lock.getExportedHandler(entrypointName, kj::mv(versionInfo), kj::mv(props), |
| 364 | context.getActor(), isDynamicDispatch), |
| 365 | kj::mv(signal)); |
| 366 | }) |
| 367 | .then([this, &context, &wrappedResponse = *wrappedResponse, workerTracer]( |
| 368 | api::DeferredProxy<void> deferredProxy) { |
| 369 | TRACE_EVENT("workerd", "WorkerEntrypoint::request() deferred proxy step", |
| 370 | PERFETTO_FLOW_FROM_POINTER(this)); |
| 371 | proxyTask = kj::mv(deferredProxy.proxyTask); |
| 372 | KJ_IF_SOME(t, workerTracer) { |
| 373 | auto httpResponseStatus = wrappedResponse.getHttpResponseStatus(); |
| 374 | if (httpResponseStatus != 0) { |
| 375 | t.setReturn(context.now(), tracing::FetchResponseInfo(httpResponseStatus)); |
| 376 | } else { |
| 377 | t.setReturn(context.now()); |
| 378 | } |
| 379 | } |
| 380 | }) |
| 381 | .catch_([this, &context](kj::Exception&& exception) mutable -> kj::Promise<void> { |
| 382 | TRACE_EVENT("workerd", "WorkerEntrypoint::request() catch", PERFETTO_FLOW_FROM_POINTER(this)); |
| 383 | // Log JS exceptions to the JS console, if inspector is attached. This also has the effect of |
| 384 | // logging internal errors to syslog. |
| 385 | loggedExceptionEarlier = true; |
| 386 | context.logUncaughtExceptionAsync(UncaughtExceptionSource::REQUEST_HANDLER, exception.clone()); |
| 387 | |
| 388 | // Do not allow the exception to escape the isolate without waiting for the output gate to |
| 389 | // open. Note that in the success path, this is taken care of in `FetchEvent::respondWith()`. |
| 390 | return context.waitForOutputLocks().then( |
| 391 | #ifdef WORKERD_USE_PERFETTO |
| 392 | [exception = kj::mv(exception), |
| 393 | flow = PERFETTO_TERMINATING_FLOW_FROM_POINTER(this)]() mutable -> kj::Promise<void> { |
| 394 | TRACE_EVENT("workerd", "WorkerEntrypoint::request() after output lock wait", flow); |
| 395 | return kj::mv(exception); |
| 396 | }); |
| 397 | #else |
| 398 | [exception = kj::mv(exception)]() mutable -> kj::Promise<void> { |
| 399 | return kj::mv(exception); |
| 400 | }); |
| 401 | #endif // defined(WORKERD_USE_PERFETTO) |
| 402 | }) |
| 403 | .attach(kj::defer([this, incomingRequest = kj::mv(incomingRequest), &context]() mutable { |
| 404 | // The request has been canceled, but allow it to continue executing in the background. |
| 405 | if (context.isFailOpen()) { |
| 406 | // Fail-open behavior has been chosen, we'd better save an interface that we can use for |
| 407 | // that purpose later. |
| 408 | failOpenService = context.getSubrequestChannelNoChecks( |
| 409 | IoContext::NEXT_CLIENT_CHANNEL, false, kj::mv(cfBlobJson)); |
| 410 | } |
| 411 | |
| 412 | if (proxyTask == kj::none && !loggedExceptionEarlier) { |
| 413 | // When the client disconnects, trigger an abort on request.signal, unless the request has |
| 414 | // already completed normally, or failed with an exception. |
| 415 | |
| 416 | // TODO(perf): Don't add a task to trigger the abort unless we know it has at least one |
| 417 | // listener. |
| 418 | KJ_IF_SOME(ctrl, abortController) { |
| 419 | context.addWaitUntil(context.run([ctrl = ctrl.addRef()](Worker::Lock& lock) mutable { |
| 420 | ctrl->getSignal()->triggerAbort( |
| 421 | lock, JSG_KJ_EXCEPTION(DISCONNECTED, DOMAbortError, "The client has disconnected")); |
| 422 | })); |
| 423 | } |
| 424 | } |
| 425 | |
| 426 | // Release reference to the AbortController. |
| 427 | // Either the waitUntilTask holds a reference to it, or it will never be triggered at all. |
| 428 | abortController = kj::none; |
| 429 | |
| 430 | auto promise = incomingRequest->drain().attach(kj::mv(incomingRequest)); |
| 431 | waitUntilTasks.add(maybeAddGcPassForTest(context, kj::mv(promise))); |
| 432 | })) |
| 433 | .then([this, metrics = kj::mv(metricsForProxyTask)]() mutable -> kj::Promise<void> { |
| 434 | TRACE_EVENT("workerd", "WorkerEntrypoint::request() finish proxying", |
| 435 | PERFETTO_TERMINATING_FLOW_FROM_POINTER(this)); |
| 436 | // Now that the IoContext is dropped (unless it had waitUntil()s), we can finish proxying |
| 437 | // without pinning it or the isolate into memory. |
| 438 | KJ_IF_SOME(p, proxyTask) { |
| 439 | return p.catch_([metrics = kj::mv(metrics)](kj::Exception&& e) mutable -> kj::Promise<void> { |
| 440 | metrics->reportFailure(e, RequestObserver::FailureSource::DEFERRED_PROXY); |
| 441 | return kj::mv(e); |
| 442 | }); |
| 443 | } else { |
| 444 | return kj::READY_NOW; |
| 445 | } |
| 446 | }) |
| 447 | .attach(kj::defer([this]() mutable { |
| 448 | // If we're being cancelled, we need to make sure `proxyTask` gets canceled. |
| 449 | proxyTask = kj::none; |
| 450 | })) |
| 451 | .catch_([this, wrappedResponse = kj::mv(wrappedResponse), isActor, method, url, &headers, |
| 452 | &requestBody, metrics = kj::mv(metricsForCatch), |
| 453 | workerTracer](kj::Exception&& exception) mutable -> kj::Promise<void> { |
| 454 | // Don't return errors to end user. |
| 455 | TRACE_EVENT("workerd", "WorkerEntrypoint::request() exception", |
| 456 | PERFETTO_TERMINATING_FLOW_FROM_POINTER(this)); |
| 457 | |
| 458 | auto isInternalException = !jsg::isTunneledException(exception.getDescription()) && |
| 459 | !jsg::isDoNotLogException(exception.getDescription()); |
| 460 | if (!loggedExceptionEarlier) { |
| 461 | // This exception seems to have originated during the deferred proxy task, so it was not |
| 462 | // logged to the IoContext earlier. |
| 463 | if (exception.getType() != kj::Exception::Type::DISCONNECTED && isInternalException) { |
| 464 | LOG_EXCEPTION("workerEntrypoint", exception); |
| 465 | } else { |
| 466 | KJ_LOG(INFO, exception); // Run with --verbose to see exception logs. |
| 467 | } |
| 468 | } |
| 469 | |
| 470 | if (wrappedResponse->isSent()) { |
| 471 | // We can't fail open if the response was already sent, so set `failOpenService` null so that |
| 472 | // that branch isn't taken below. |
| 473 | failOpenService = kj::none; |
| 474 | } |
| 475 | |
| 476 | if (isActor) { |
| 477 | // We want to tunnel exceptions from actors back to the caller. |
| 478 | // TODO(cleanup): We'd really like to tunnel exceptions any time a worker is calling another |
| 479 | // worker, not just for actors (and W2W below), but getting that right will require cleaning |
| 480 | // up error handling more generally. |
| 481 | return exceptionToPropagate(isInternalException, kj::mv(exception)); |
| 482 | } else KJ_IF_SOME(service, failOpenService) { |
| 483 | // Fall back to origin. |
| 484 | |
| 485 | // We're catching the exception, but metrics should still indicate an exception. |
| 486 | metrics->reportFailure(exception); |
| 487 | |
| 488 | auto promise = kj::evalNow([&] { |
| 489 | auto promise = service.get()->request(method, url, headers, requestBody, *wrappedResponse); |
| 490 | metrics->setFailedOpen(true); |
| 491 | return promise.attach(kj::mv(service)); |
| 492 | }); |
| 493 | return promise.catch_([this, wrappedResponse = kj::mv(wrappedResponse), workerTracer, |
| 494 | metrics = kj::mv(metrics)](kj::Exception&& e) mutable { |
| 495 | metrics->setFailedOpen(false); |
| 496 | if (e.getType() != kj::Exception::Type::DISCONNECTED && |
| 497 | // Avoid logging recognized external errors here, such as invalid headers returned from |
| 498 | // the server. |
| 499 | !jsg::isTunneledException(e.getDescription()) && |
| 500 | !jsg::isDoNotLogException(e.getDescription())) { |
| 501 | LOG_EXCEPTION("failOpenFallback", e); |
| 502 | } |
| 503 | if (!wrappedResponse->isSent()) { |
| 504 | kj::HttpHeaders headers(threadContext.getHeaderTable()); |
| 505 | wrappedResponse->send(500, "Internal Server Error", headers, static_cast<uint64_t>(0)); |
| 506 | KJ_IF_SOME(t, workerTracer) { |
| 507 | t.setReturn(kj::none, tracing::FetchResponseInfo(500)); |
| 508 | } |
| 509 | } |
| 510 | }); |
| 511 | } else if (tunnelExceptions) { |
| 512 | // Like with the isActor check, we want to return exceptions back to the caller. |
| 513 | // We don't want to handle this case the same as the isActor case though, since we want |
| 514 | // fail-open to operate normally, which means this case must happen after fail-open handling. |
| 515 | return exceptionToPropagate(isInternalException, kj::mv(exception)); |
| 516 | } else { |
| 517 | // Return error. |
| 518 | |
| 519 | // We're catching the exception and replacing it with 5xx, but metrics should still indicate |
| 520 | // an exception. |
| 521 | metrics->reportFailure(exception); |
| 522 | |
| 523 | // We can't send an error response if a response was already started; we can only drop the |
| 524 | // connection in that case. |
| 525 | if (!wrappedResponse->isSent()) { |
| 526 | kj::HttpHeaders headers(threadContext.getHeaderTable()); |
| 527 | if (exception.getType() == kj::Exception::Type::OVERLOADED) { |
| 528 | wrappedResponse->send(503, "Service Unavailable", headers, static_cast<uint64_t>(0)); |
| 529 | } else { |
| 530 | wrappedResponse->send(500, "Internal Server Error", headers, static_cast<uint64_t>(0)); |
| 531 | } |
| 532 | KJ_IF_SOME(t, workerTracer) { |
| 533 | t.setReturn( |
| 534 | kj::none, tracing::FetchResponseInfo(wrappedResponse->getHttpResponseStatus())); |
| 535 | } |
| 536 | } |
| 537 | |
| 538 | return kj::READY_NOW; |
| 539 | } |
| 540 | }); |
| 541 | } |
| 542 | |
| 543 | kj::Promise<void> WorkerEntrypoint::connect(kj::StringPtr host, |
| 544 | const kj::HttpHeaders& headers, |
| 545 | kj::AsyncIoStream& connection, |
| 546 | ConnectResponse& response, |
| 547 | kj::HttpConnectSettings settings) { |
| 548 | TRACE_EVENT("workerd", "WorkerEntrypoint::connect()"); |
| 549 | auto incomingRequest = |
| 550 | kj::mv(KJ_REQUIRE_NONNULL(this->incomingRequest, "connect() can only be called once")); |
| 551 | this->incomingRequest = kj::none; |
| 552 | auto& context = incomingRequest->getContext(); |
| 553 | auto featureFlags = context.getWorker().getIsolate().getApi().getFeatureFlags(); |
| 554 | |
| 555 | if (featureFlags.getConnectPassThrough()) { |
| 556 | incomingRequest->delivered(); |
| 557 | |
| 558 | KJ_DEFER({ |
| 559 | // Since we called incomingRequest->delivered, we are obliged to call `drain()`. |
| 560 | auto promise = incomingRequest->drain().attach(kj::mv(incomingRequest)); |
| 561 | waitUntilTasks.add(maybeAddGcPassForTest(context, kj::mv(promise))); |
| 562 | }); |
| 563 | // connect_pass_through feature flag means we should just forward the connect request on to |
| 564 | // the global outbound. |
| 565 | |
| 566 | auto next = context.getSubrequestChannelNoChecks( |
| 567 | IoContext::NEXT_CLIENT_CHANNEL, false, kj::mv(cfBlobJson)); |
| 568 | |
| 569 | // Note: Intentionally return without co_await so that the `incomingRequest` is destroyed, |
| 570 | // because we don't have any need to keep the context around. |
| 571 | return next->connect(host, headers, connection, response, settings); |
| 572 | } else if (!featureFlags.getWorkerdExperimental()) { |
| 573 | JSG_FAIL_REQUIRE(TypeError, "Incoming CONNECT on a worker not supported"); |
| 574 | } |
| 575 | |
| 576 | // TODO(soon): Implement basic TLS support for connect handler. |
| 577 | JSG_REQUIRE(!settings.useTls, Error, "Incoming CONNECT with TLS not supported"); |
| 578 | // Capture workerTracer, see request() for rationale. |
| 579 | kj::Maybe<BaseTracer&> workerTracer; |
| 580 | |
| 581 | bool isActor = context.getActor() != kj::none; |
| 582 | |
| 583 | KJ_IF_SOME(t, incomingRequest->getWorkerTracer()) { |
| 584 | t.setEventInfo(*incomingRequest, tracing::ConnectEventInfo()); |
| 585 | workerTracer = t; |
| 586 | } |
| 587 | incomingRequest->delivered(); |
| 588 | |
| 589 | auto metricsForCatch = kj::addRef(incomingRequest->getMetrics()); |
| 590 | |
| 591 | return context |
| 592 | .run( |
| 593 | [this, &headers, &context, &connection, &response, entrypointName = entrypointName, |
| 594 | versionInfo = kj::mv(versionInfo), host = kj::str(host)](Worker::Lock& lock) mutable { |
| 595 | jsg::AsyncContextFrame::StorageScope traceScope = context.makeAsyncTraceScope(lock); |
| 596 | jsg::AsyncContextFrame::StorageScope userTraceScope = context.makeUserAsyncTraceScope(lock); |
| 597 | |
| 598 | return lock.getGlobalScope().connect(kj::mv(host), headers, connection, response, lock, |
| 599 | lock.getExportedHandler(entrypointName, kj::mv(versionInfo), kj::mv(props), |
| 600 | context.getActor(), isDynamicDispatch)); |
| 601 | }) |
| 602 | .then([&context, workerTracer]() { |
| 603 | KJ_IF_SOME(t, workerTracer) { |
| 604 | t.setReturn(context.now()); |
| 605 | } |
| 606 | }) |
| 607 | .catch_([this, &context](kj::Exception&& exception) mutable -> kj::Promise<void> { |
| 608 | // Log JS exceptions to the JS console, if inspector is attached. This also has the effect of |
| 609 | // logging internal errors to syslog. |
| 610 | loggedExceptionEarlier = true; |
| 611 | context.logUncaughtExceptionAsync(UncaughtExceptionSource::REQUEST_HANDLER, exception.clone()); |
| 612 | |
| 613 | // Do not allow the exception to escape the isolate without waiting for the output gate to |
| 614 | // open. Note that in the success path, this is taken care of in `FetchEvent::respondWith()`. |
| 615 | return context.waitForOutputLocks().then( |
| 616 | [exception = kj::mv(exception)]() mutable -> kj::Promise<void> { |
| 617 | return kj::mv(exception); |
| 618 | }); |
| 619 | }) |
| 620 | .attach(kj::defer([this, incomingRequest = kj::mv(incomingRequest), &context]() mutable { |
| 621 | // The request has been canceled, but allow it to continue executing in the background. |
| 622 | auto promise = incomingRequest->drain().attach(kj::mv(incomingRequest)); |
| 623 | waitUntilTasks.add(maybeAddGcPassForTest(context, kj::mv(promise))); |
| 624 | })) |
| 625 | .catch_([this, isActor, &response, metrics = kj::mv(metricsForCatch), workerTracer]( |
| 626 | kj::Exception&& exception) mutable -> kj::Promise<void> { |
| 627 | // Don't return errors to end user. |
| 628 | auto isInternalException = !jsg::isTunneledException(exception.getDescription()) && |
| 629 | !jsg::isDoNotLogException(exception.getDescription()); |
| 630 | if (!loggedExceptionEarlier) { |
| 631 | // This exception seems to have originated during the deferred proxy task, so it was not |
| 632 | // logged to the IoContext earlier. |
| 633 | if (exception.getType() != kj::Exception::Type::DISCONNECTED && isInternalException) { |
| 634 | LOG_EXCEPTION("workerEntrypoint", exception); |
| 635 | } else { |
| 636 | KJ_LOG(INFO, exception); // Run with --verbose to see exception logs. |
| 637 | } |
| 638 | } |
| 639 | |
| 640 | if (isActor || tunnelExceptions) { |
| 641 | // We want to tunnel exceptions from actors back to the caller. |
| 642 | // TODO(cleanup): We'd really like to tunnel exceptions any time a worker is calling another |
| 643 | // worker, not just for actors (and W2W below), but getting that right will require cleaning |
| 644 | // up error handling more generally. |
| 645 | return exceptionToPropagate(isInternalException, kj::mv(exception)); |
| 646 | } else { |
| 647 | // Return error. |
| 648 | |
| 649 | // We're catching the exception and replacing it with 5xx, but metrics should still indicate |
| 650 | // an exception. |
| 651 | metrics->reportFailure(exception); |
| 652 | |
| 653 | kj::HttpHeaders headers(threadContext.getHeaderTable()); |
| 654 | if (exception.getType() == kj::Exception::Type::OVERLOADED) { |
| 655 | response.reject(503, "Service Unavailable", headers, static_cast<uint64_t>(0)); |
| 656 | } else { |
| 657 | response.reject(500, "Internal Server Error", headers, static_cast<uint64_t>(0)); |
| 658 | } |
| 659 | // TODO(o11y): Should we also indicate a return response code for TCP? |
| 660 | KJ_IF_SOME(t, workerTracer) { |
| 661 | t.setReturn(kj::none); |
| 662 | } |
| 663 | |
| 664 | return kj::READY_NOW; |
| 665 | } |
| 666 | }); |
| 667 | } |
| 668 | |
| 669 | kj::Promise<void> WorkerEntrypoint::prewarm(kj::StringPtr url) { |
| 670 | // Nothing to do, the worker is already loaded. |
| 671 | TRACE_EVENT("workerd", "WorkerEntrypoint::prewarm()", "url", url.cStr()); |
| 672 | auto incomingRequest = |
| 673 | kj::mv(KJ_REQUIRE_NONNULL(this->incomingRequest, "prewarm() can only be called once")); |
| 674 | incomingRequest->getMetrics().setIsPrewarm(); |
| 675 | |
| 676 | // Intentionally don't call incomingRequest->delivered() for prewarm requests and do not create |
| 677 | // an Onset event, prewarm is not being traced. |
| 678 | |
| 679 | // TODO(someday): Ideally, middleware workers would forward prewarm() to the next stage. At |
| 680 | // present we don't have a good way to decide what stage that is, especially given that we'll |
| 681 | // be switching to `next` being a binding in the future. |
| 682 | return kj::READY_NOW; |
| 683 | } |
| 684 | |
| 685 | kj::Promise<WorkerInterface::ScheduledResult> WorkerEntrypoint::runScheduled( |
| 686 | kj::Date scheduledTime, kj::StringPtr cron) { |
| 687 | TRACE_EVENT("workerd", "WorkerEntrypoint::runScheduled()"); |
| 688 | auto incomingRequest = |
| 689 | kj::mv(KJ_REQUIRE_NONNULL(this->incomingRequest, "runScheduled() can only be called once")); |
| 690 | this->incomingRequest = kj::none; |
| 691 | auto& context = incomingRequest->getContext(); |
| 692 | |
| 693 | KJ_ASSERT(context.getActor() == kj::none); |
| 694 | // This code currently doesn't work with actors because cancellations occur immediately, without |
| 695 | // calling context->drain(). We don't ever send scheduled events to actors. If we do, we'll have |
| 696 | // to think more about this. |
| 697 | |
| 698 | double eventTime = (scheduledTime - kj::UNIX_EPOCH) / kj::MILLISECONDS; |
| 699 | |
| 700 | KJ_IF_SOME(t, incomingRequest->getWorkerTracer()) { |
| 701 | t.setEventInfo(*incomingRequest, tracing::ScheduledEventInfo(eventTime, kj::str(cron))); |
| 702 | } |
| 703 | |
| 704 | incomingRequest->delivered(); |
| 705 | |
| 706 | // Scheduled handlers run entirely in waitUntil() tasks. |
| 707 | context.addWaitUntil( |
| 708 | context.run([scheduledTime, cron, entrypointName = entrypointName, |
| 709 | versionInfo = kj::mv(versionInfo), props = kj::mv(props), &context, |
| 710 | &metrics = incomingRequest->getMetrics()](Worker::Lock& lock) mutable { |
| 711 | TRACE_EVENT("workerd", "WorkerEntrypoint::runScheduled() run"); |
| 712 | jsg::AsyncContextFrame::StorageScope traceScope = context.makeAsyncTraceScope(lock); |
| 713 | jsg::AsyncContextFrame::StorageScope userTraceScope = context.makeUserAsyncTraceScope(lock); |
| 714 | |
| 715 | lock.getGlobalScope().startScheduled(scheduledTime, cron, lock, |
| 716 | lock.getExportedHandler( |
| 717 | entrypointName, kj::mv(versionInfo), kj::mv(props), context.getActor())); |
| 718 | })); |
| 719 | |
| 720 | static auto constexpr waitForFinished = [](IoContext& context, |
| 721 | kj::Own<IoContext::IncomingRequest> request) |
| 722 | -> kj::Promise<WorkerInterface::ScheduledResult> { |
| 723 | TRACE_EVENT("workerd", "WorkerEntrypoint::runScheduled() waitForFinished()"); |
| 724 | auto scheduledResult = co_await request->finishScheduled(); |
| 725 | bool completed = scheduledResult == EventOutcome::OK; |
| 726 | co_return WorkerInterface::ScheduledResult{.retry = context.shouldRetryScheduled(), |
| 727 | .outcome = completed ? context.waitUntilStatus() : scheduledResult}; |
| 728 | }; |
| 729 | |
| 730 | auto promise = waitForFinished(context, kj::mv(incomingRequest)); |
| 731 | |
| 732 | return maybeAddGcPassForTest(context, kj::mv(promise)); |
| 733 | } |
| 734 | |
| 735 | kj::Promise<WorkerInterface::AlarmResult> WorkerEntrypoint::runAlarmImpl( |
| 736 | kj::Own<IoContext::IncomingRequest> incomingRequest, |
| 737 | kj::Date scheduledTime, |
| 738 | uint32_t retryCount) { |
| 739 | // We want to de-duplicate alarm requests as follows: |
| 740 | // - An alarm must not be canceled once it is running, UNLESS the whole actor is shut down. |
| 741 | // - If multiple alarm invocations arrive with the same scheduled time, we only run one. |
| 742 | // - If we are asked to schedule an alarm while one is running, we wait for the running alarm to |
| 743 | // finish. |
| 744 | // - However, we schedule no more than one alarm. If another one (with yet another different |
| 745 | // scheduled time) arrives while we still have one running and one scheduled, we discard the |
| 746 | // previous scheduled alarm. |
| 747 | |
| 748 | TRACE_EVENT("workerd", "WorkerEntrypoint::runAlarmImpl()"); |
| 749 | |
| 750 | auto& context = incomingRequest->getContext(); |
| 751 | auto& actor = KJ_REQUIRE_NONNULL(context.getActor(), "alarm() should only work with actors"); |
| 752 | |
| 753 | KJ_IF_SOME(promise, actor.getAlarm(scheduledTime)) { |
| 754 | // There is a pre-existing alarm for `scheduledTime`, we can just wait for its result. |
| 755 | // TODO(someday) If the request responsible for fulfilling this alarm were to be cancelled, then |
| 756 | // we could probably take over and try to fulfill it ourselves. Maybe we'd want to loop on |
| 757 | // `actor.getAlarm()`? We'd have to distinguish between rescheduling and request cancellation. |
| 758 | auto outcome = co_await promise; |
| 759 | co_return AlarmResult{.retry = outcome.retry, |
| 760 | .retryCountsAgainstLimit = outcome.retryCountsAgainstLimit, |
| 761 | .outcome = outcome.outcome}; |
| 762 | } |
| 763 | |
| 764 | // There isn't a pre-existing alarm, we can set event info and call `delivered()` (which emits |
| 765 | // metrics events). |
| 766 | KJ_IF_SOME(t, incomingRequest->getWorkerTracer()) { |
| 767 | t.setEventInfo(*incomingRequest, tracing::AlarmEventInfo(scheduledTime)); |
| 768 | } |
| 769 | |
| 770 | incomingRequest->delivered(); |
| 771 | |
| 772 | auto scheduleAlarmResult = co_await actor.scheduleAlarm(scheduledTime); |
| 773 | KJ_SWITCH_ONEOF(scheduleAlarmResult) { |
| 774 | KJ_CASE_ONEOF(af, WorkerInterface::AlarmFulfiller) { |
| 775 | // We're now in charge of running this alarm! |
| 776 | auto cancellationGuard = kj::defer([&af]() { |
| 777 | // Our promise chain was cancelled, let's cancel our fulfiller for any other requests |
| 778 | // that were waiting on us. |
| 779 | af.cancel(); |
| 780 | }); |
| 781 | |
| 782 | KJ_DEFER({ |
| 783 | // The alarm has finished but allow the request to continue executing in the background. |
| 784 | waitUntilTasks.add(incomingRequest->drain().attach(kj::mv(incomingRequest))); |
| 785 | }); |
| 786 | |
| 787 | try { |
| 788 | auto result = |
| 789 | co_await context.run([scheduledTime, retryCount, entrypointName = entrypointName, |
| 790 | versionInfo = kj::mv(versionInfo), props = kj::mv(props), |
| 791 | &context](Worker::Lock& lock) mutable { |
| 792 | jsg::AsyncContextFrame::StorageScope traceScope = context.makeAsyncTraceScope(lock); |
| 793 | jsg::AsyncContextFrame::StorageScope userTraceScope = |
| 794 | context.makeUserAsyncTraceScope(lock); |
| 795 | |
| 796 | // If we have an invalid timeout, set it to the default value of 15 minutes. |
| 797 | auto timeout = context.getLimitEnforcer().getAlarmLimit(); |
| 798 | if (timeout == 0 * kj::MILLISECONDS) { |
| 799 | LOG_NOSENTRY(WARNING, "Invalid alarm timeout value. Using 15 minutes", timeout); |
| 800 | timeout = 15 * kj::MINUTES; |
| 801 | } |
| 802 | |
| 803 | auto handler = lock.getExportedHandler( |
| 804 | entrypointName, kj::mv(versionInfo), kj::mv(props), context.getActor()); |
| 805 | return lock.getGlobalScope().runAlarm(scheduledTime, timeout, retryCount, lock, handler); |
| 806 | }); |
| 807 | |
| 808 | // The alarm handler was successfully complete. We must guarantee this same alarm does not |
| 809 | // run again. |
| 810 | if (result.outcome == EventOutcome::OK) { |
| 811 | // When an alarm handler completes its execution, the alarm is marked ready for deletion in |
| 812 | // actor-cache. This alarm change will only be reflected in the alarmsXX table, once cache |
| 813 | // flushes and changes are written to storage. |
| 814 | // If there are any pending flushes, they are locked with the actor output gate until |
| 815 | // they complete. We should wait until the output gate locks are released. |
| 816 | // If we don't wait, it's possible for alarm manager to pull the wrong alarm value (the |
| 817 | // same alarm that just completed) from storage before these changes are actually made, |
| 818 | // rerunning it, when it shouldn't. |
| 819 | co_await actor.getOutputGate().wait(context.getCurrentTraceSpan()); |
| 820 | } |
| 821 | |
| 822 | // We succeeded, inform any other entrypoints that may be waiting upon us. |
| 823 | af.fulfill(result.asOutcome()); |
| 824 | cancellationGuard.cancel(); |
| 825 | co_return kj::mv(result); |
| 826 | } catch (const kj::Exception& e) { |
| 827 | // We failed, inform any other entrypoints that may be waiting upon us. |
| 828 | af.reject(e); |
| 829 | cancellationGuard.cancel(); |
| 830 | throw; |
| 831 | } |
| 832 | } |
| 833 | KJ_CASE_ONEOF(outcome, WorkerInterface::AlarmOutcome) { |
| 834 | // The alarm was cancelled while we were waiting to run, go ahead and return the result. |
| 835 | co_return AlarmResult{.retry = outcome.retry, |
| 836 | .retryCountsAgainstLimit = outcome.retryCountsAgainstLimit, |
| 837 | .outcome = outcome.outcome}; |
| 838 | } |
| 839 | } |
| 840 | |
| 841 | KJ_UNREACHABLE; |
| 842 | } |
| 843 | |
| 844 | kj::Promise<WorkerInterface::AlarmResult> WorkerEntrypoint::runAlarm( |
| 845 | kj::Date scheduledTime, uint32_t retryCount) { |
| 846 | TRACE_EVENT("workerd", "WorkerEntrypoint::runAlarm()"); |
| 847 | auto incomingRequest = |
| 848 | kj::mv(KJ_REQUIRE_NONNULL(this->incomingRequest, "runAlarm() can only be called once")); |
| 849 | this->incomingRequest = kj::none; |
| 850 | |
| 851 | auto& context = incomingRequest->getContext(); |
| 852 | auto promise = runAlarmImpl(kj::mv(incomingRequest), scheduledTime, retryCount); |
| 853 | auto result = co_await maybeAddGcPassForTest(context, kj::mv(promise)); |
| 854 | KJ_IF_SOME(t, context.getWorkerTracer()) { |
| 855 | t.setReturn(context.now()); |
| 856 | } |
| 857 | co_return result; |
| 858 | } |
| 859 | |
| 860 | kj::Promise<kj::Maybe<kj::Date>> WorkerEntrypoint::abandonAlarm(kj::Date scheduledTime) { |
| 861 | TRACE_EVENT("workerd", "WorkerEntrypoint::abandonAlarm()"); |
| 862 | // This does not require running the user's alarm handler -- it's a pure actor-state cleanup. |
| 863 | // Access the actor directly from the IoContext without going through the JS dispatch machinery. |
| 864 | auto& req = |
| 865 | KJ_REQUIRE_NONNULL(incomingRequest, "abandonAlarm() called without an incoming request"); |
| 866 | auto& actor = KJ_REQUIRE_NONNULL( |
| 867 | req->getContext().getActor(), "abandonAlarm() should only work with actors"); |
| 868 | auto& persistent = KJ_REQUIRE_NONNULL( |
| 869 | actor.getPersistent(), "abandonAlarm() requires actor with persistent storage"); |
| 870 | return persistent.abandonAlarm(scheduledTime); |
| 871 | } |
| 872 | |
| 873 | kj::Promise<bool> WorkerEntrypoint::test() { |
| 874 | TRACE_EVENT("workerd", "WorkerEntrypoint::test()"); |
| 875 | auto incomingRequest = |
| 876 | kj::mv(KJ_REQUIRE_NONNULL(this->incomingRequest, "test() can only be called once")); |
| 877 | this->incomingRequest = kj::none; |
| 878 | auto& context = incomingRequest->getContext(); |
| 879 | KJ_IF_SOME(t, incomingRequest->getWorkerTracer()) { |
| 880 | t.setEventInfo(*incomingRequest, tracing::CustomEventInfo()); |
| 881 | } |
| 882 | |
| 883 | incomingRequest->delivered(); |
| 884 | |
| 885 | context.addWaitUntil( |
| 886 | context.run([entrypointName = entrypointName, versionInfo = kj::mv(versionInfo), |
| 887 | props = kj::mv(props), &context, &metrics = incomingRequest->getMetrics()]( |
| 888 | Worker::Lock& lock) mutable -> kj::Promise<void> { |
| 889 | TRACE_EVENT("workerd", "WorkerEntrypoint::test() run"); |
| 890 | jsg::AsyncContextFrame::StorageScope traceScope = context.makeAsyncTraceScope(lock); |
| 891 | jsg::AsyncContextFrame::StorageScope userTraceScope = context.makeUserAsyncTraceScope(lock); |
| 892 | |
| 893 | return context.awaitJs(lock, |
| 894 | lock.getGlobalScope().test(lock, |
| 895 | lock.getExportedHandler( |
| 896 | entrypointName, kj::mv(versionInfo), kj::mv(props), context.getActor()))); |
| 897 | })); |
| 898 | |
| 899 | static auto constexpr waitForFinished = |
| 900 | [](IoContext& context, kj::Own<IoContext::IncomingRequest> request) -> kj::Promise<bool> { |
| 901 | TRACE_EVENT("workerd", "WorkerEntrypoint::test() waitForFinished()"); |
| 902 | auto scheduledResult = co_await request->finishScheduled(); |
| 903 | |
| 904 | if (scheduledResult == EventOutcome::EXCEPTION) { |
| 905 | // If the test handler throws an exception (without aborting - just a regular exception), |
| 906 | // then `outcome` ends up being EventOutcome::EXCEPTION, which causes us to return false. |
| 907 | // But in that case we are separately relying on the exception being logged as an uncaught |
| 908 | // exception, rather than throwing it. |
| 909 | // This is why we don't rethrow the exception but rather log it as an uncaught exception. |
| 910 | try { |
| 911 | co_await context.onAbort(); |
| 912 | } catch (...) { |
| 913 | auto exception = kj::getCaughtExceptionAsKj(); |
| 914 | KJ_LOG(ERROR, exception); |
| 915 | } |
| 916 | } |
| 917 | |
| 918 | // Not adding a return event here – we only provide rudimentary tracing support for test events |
| 919 | // (enough so that we can get logs/spans from them in wd-tests), so this is not needed in |
| 920 | // practice. |
| 921 | |
| 922 | bool completed = scheduledResult == EventOutcome::OK; |
| 923 | auto outcome = completed ? context.waitUntilStatus() : scheduledResult; |
| 924 | co_return outcome == EventOutcome::OK; |
| 925 | }; |
| 926 | |
| 927 | return maybeAddGcPassForTest(context, waitForFinished(context, kj::mv(incomingRequest))); |
| 928 | } |
| 929 | |
| 930 | kj::Promise<WorkerInterface::CustomEvent::Result> WorkerEntrypoint::customEvent( |
| 931 | kj::Own<CustomEvent> event) { |
| 932 | TRACE_EVENT("workerd", "WorkerEntrypoint::customEvent()", "type", event->getType()); |
| 933 | auto incomingRequest = |
| 934 | kj::mv(KJ_REQUIRE_NONNULL(this->incomingRequest, "customEvent() can only be called once")); |
| 935 | this->incomingRequest = kj::none; |
| 936 | |
| 937 | auto& context = incomingRequest->getContext(); |
| 938 | |
| 939 | // Set event info BEFORE calling run() to ensure onset event is reported before |
| 940 | // any user code executes (particularly important for actors whose constructors may run |
| 941 | // during delivered()). |
| 942 | KJ_IF_SOME(t, incomingRequest->getWorkerTracer()) { |
| 943 | t.setEventInfo(*incomingRequest, event->getEventInfo()); |
| 944 | } |
| 945 | |
| 946 | auto promise = event |
| 947 | ->run(kj::mv(incomingRequest), entrypointName, kj::mv(versionInfo), |
| 948 | kj::mv(props), waitUntilTasks, isDynamicDispatch) |
| 949 | .attach(kj::mv(event)); |
| 950 | |
| 951 | // TODO(cleanup): In theory `context` may have been destroyed by now if `event->run()` dropped |
| 952 | // the `incomingRequest` synchronously. No current implementation does that, and |
| 953 | // maybeAddGcPassForTest() is a no-op outside of tests, so I'm ignoring the theoretical problem |
| 954 | // for now. Otherwise we will need to `atomicAddRef()` the `Worker` at some point earlier on |
| 955 | // but I'd like to avoid that in the non-test case. |
| 956 | return maybeAddGcPassForTest(context, kj::mv(promise)); |
| 957 | } |
| 958 | |
| 959 | #ifdef KJ_DEBUG |
| 960 | void requestGc(const Worker& worker) { |
| 961 | TRACE_EVENT("workerd", "Debug: requestGc()"); |
| 962 | jsg::runInV8Stack([&](jsg::V8StackScope& stackScope) { |
| 963 | auto& isolate = worker.getIsolate(); |
| 964 | auto lock = isolate.getApi().lock(stackScope); |
| 965 | lock->requestGcForTesting(); |
| 966 | }); |
| 967 | } |
| 968 | |
| 969 | template <typename T> |
| 970 | kj::Promise<T> addGcPassForTest(IoContext& context, kj::Promise<T> promise) { |
| 971 | TRACE_EVENT("workerd", "Debug: addGcPassForTest"); |
| 972 | auto worker = kj::atomicAddRef(context.getWorker()); |
| 973 | if constexpr (kj::isSameType<T, void>()) { |
| 974 | co_await promise; |
| 975 | requestGc(*worker); |
| 976 | } else { |
| 977 | auto ret = co_await promise; |
| 978 | requestGc(*worker); |
| 979 | co_return kj::mv(ret); |
| 980 | } |
| 981 | } |
| 982 | #endif |
| 983 | |
| 984 | template <typename T> |
| 985 | kj::Promise<T> WorkerEntrypoint::maybeAddGcPassForTest(IoContext& context, kj::Promise<T> promise) { |
| 986 | #ifdef KJ_DEBUG |
| 987 | if (isPredictableModeForTest()) { |
| 988 | return addGcPassForTest(context, kj::mv(promise)); |
| 989 | } |
| 990 | #endif |
| 991 | return kj::mv(promise); |
| 992 | } |
| 993 | |
| 994 | } // namespace |
| 995 | |
| 996 | kj::Own<WorkerInterface> newWorkerEntrypoint(ThreadContext& threadContext, |
| 997 | kj::Own<const Worker> worker, |
| 998 | kj::Maybe<kj::StringPtr> entrypointName, |
| 999 | Frankenvalue props, |
| 1000 | kj::Maybe<kj::Own<Worker::Actor>> actor, |
| 1001 | kj::Own<LimitEnforcer> limitEnforcer, |
| 1002 | kj::Own<void> ioContextDependency, |
| 1003 | kj::Own<IoChannelFactory> ioChannelFactory, |
| 1004 | kj::Own<RequestObserver> metrics, |
| 1005 | kj::TaskSet& waitUntilTasks, |
| 1006 | bool tunnelExceptions, |
| 1007 | kj::Maybe<kj::Own<BaseTracer>> workerTracer, |
| 1008 | kj::Maybe<kj::String> cfBlobJson, |
| 1009 | kj::Maybe<Worker::VersionInfo> versionInfo, |
| 1010 | kj::Maybe<tracing::InvocationSpanContext> maybeTriggerInvocationSpan, |
| 1011 | bool isDynamicDispatch) { |
| 1012 | return WorkerEntrypoint::construct(threadContext, kj::mv(worker), kj::mv(entrypointName), |
| 1013 | kj::mv(props), kj::mv(actor), kj::mv(limitEnforcer), kj::mv(ioContextDependency), |
| 1014 | kj::mv(ioChannelFactory), kj::mv(metrics), waitUntilTasks, tunnelExceptions, |
| 1015 | kj::mv(workerTracer), kj::mv(cfBlobJson), kj::mv(versionInfo), |
| 1016 | kj::mv(maybeTriggerInvocationSpan), isDynamicDispatch); |
| 1017 | } |
| 1018 | |
| 1019 | } // namespace workerd |