File
Blob: src/workerd/api/cache.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 "cache.h" |
| 6 | |
| 7 | #include "util.h" |
| 8 | |
| 9 | #include <workerd/io/io-context.h> |
| 10 | #include <workerd/util/own-util.h> |
| 11 | |
| 12 | #include <kj/encoding.h> |
| 13 | |
| 14 | namespace workerd::api { |
| 15 | |
| 16 | // ======================================================================================= |
| 17 | // Cache |
| 18 | |
| 19 | namespace { |
| 20 | |
| 21 | #define LOG_CACHE_ERROR_ONCE(TEXT, RESPONSE) |
| 22 | // TODO(someday): Fix Cache API bugs. We logged them for two years as a reminder, but... they |
| 23 | // never got fixed. The logging is making it hard to see other problems. So we're ending it. |
| 24 | // If someone decides to take this on again, you can restore this macro's implementation as |
| 25 | // follows: |
| 26 | // |
| 27 | // #define LOG_CACHE_ERROR_ONCE(TEXT, RESPONSE) ({ \ |
| 28 | // static bool seen = false; \ |
| 29 | // if (!seen) { \ |
| 30 | // seen = true; \ |
| 31 | // KJ_LOG(ERROR, TEXT, RESPONSE.statusCode, RESPONSE.statusText); \ |
| 32 | // } \ |
| 33 | // }) |
| 34 | |
| 35 | // Throw an application-visible exception if the URL won't be parsed correctly at a lower |
| 36 | // layer. If the URL is valid then just return it. The purpose of this function is to avoid |
| 37 | // throwing an "internal error". |
| 38 | kj::StringPtr validateUrl(kj::StringPtr url) { |
| 39 | // TODO(bug): We should parse and process URLs the same way we would URLs passed to fetch(). |
| 40 | // But, that might mean e.g. discarding fragments ("hashes", stuff after a '#'), which would |
| 41 | // be a change in behavior that could subtly affect production workers... |
| 42 | |
| 43 | static constexpr auto urlOptions = kj::Url::Options{ |
| 44 | .percentDecode = false, |
| 45 | .allowEmpty = true, |
| 46 | }; |
| 47 | |
| 48 | JSG_REQUIRE(kj::Url::tryParse(url, kj::Url::HTTP_PROXY_REQUEST, urlOptions) != kj::none, |
| 49 | TypeError, "Invalid URL. Cache API keys must be fully-qualified, valid URLs."); |
| 50 | |
| 51 | return url; |
| 52 | } |
| 53 | |
| 54 | } // namespace |
| 55 | |
| 56 | Cache::Cache(kj::Maybe<kj::String> cacheName): cacheName(kj::mv(cacheName)) {} |
| 57 | |
| 58 | jsg::Unimplemented Cache::add(Request::Info request) { |
| 59 | return {}; |
| 60 | } |
| 61 | |
| 62 | jsg::Unimplemented Cache::addAll(kj::Array<Request::Info> requests) { |
| 63 | return {}; |
| 64 | } |
| 65 | |
| 66 | jsg::Promise<jsg::Optional<jsg::Ref<Response>>> Cache::match(jsg::Lock& js, |
| 67 | Request::Info requestOrUrl, |
| 68 | jsg::Optional<CacheQueryOptions> options, |
| 69 | CompatibilityFlags::Reader flags) { |
| 70 | auto& context = IoContext::current(); |
| 71 | TraceContext traceContext = context.makeUserTraceSpan("cache_match"_kjc); |
| 72 | |
| 73 | KJ_IF_SOME(o, options) { |
| 74 | KJ_IF_SOME(ignoreMethod, o.ignoreMethod) { |
| 75 | traceContext.setTag("cache.request.ignore_method"_kjc, ignoreMethod); |
| 76 | } |
| 77 | } |
| 78 | |
| 79 | // This use of evalNow() is obsoleted by the capture_async_api_throws compatibility flag, but |
| 80 | // we need to keep it here for people who don't have that flag set. |
| 81 | return js.evalNow([&]() -> jsg::Promise<jsg::Optional<jsg::Ref<Response>>> { |
| 82 | auto jsRequest = Request::coerce(js, kj::mv(requestOrUrl), kj::none); |
| 83 | |
| 84 | traceContext.setTag("cache.request.url"_kjc, jsRequest->getUrl()); |
| 85 | traceContext.setTag("cache.request.method"_kjc, kj::str(jsRequest->getMethodEnum())); |
| 86 | |
| 87 | if (!options.orDefault({}).ignoreMethod.orDefault(false) && |
| 88 | jsRequest->getMethodEnum() != kj::HttpMethod::GET) { |
| 89 | return js.resolvedPromise(jsg::Optional<jsg::Ref<Response>>()); |
| 90 | } |
| 91 | |
| 92 | auto httpClient = getHttpClient( |
| 93 | context, jsRequest->serializeCfBlobJson(js), traceContext, flags.getCacheApiCompatFlags()); |
| 94 | auto requestHeaders = kj::HttpHeaders(context.getHeaderTable()); |
| 95 | jsRequest->shallowCopyHeadersTo(requestHeaders); |
| 96 | |
| 97 | auto headerIds = context.getHeaderIds(); |
| 98 | // parse each of the request headers to add info to span |
| 99 | KJ_IF_SOME(range, requestHeaders.get(headerIds.range)) { |
| 100 | traceContext.setTag("cache.request.header.range"_kjc, range); |
| 101 | } |
| 102 | KJ_IF_SOME(ifModifiedSince, requestHeaders.get(headerIds.ifModifiedSince)) { |
| 103 | traceContext.setTag("cache.request.header.if_modified_since"_kjc, ifModifiedSince); |
| 104 | } |
| 105 | KJ_IF_SOME(ifNoneMatch, requestHeaders.get(headerIds.ifNoneMatch)) { |
| 106 | traceContext.setTag("cache.request.header.if_none_match"_kjc, ifNoneMatch); |
| 107 | } |
| 108 | |
| 109 | requestHeaders.setPtr(context.getHeaderIds().cacheControl, "only-if-cached"); |
| 110 | auto nativeRequest = httpClient->request(kj::HttpMethod::GET, validateUrl(jsRequest->getUrl()), |
| 111 | requestHeaders, static_cast<uint64_t>(0)); |
| 112 | |
| 113 | return context.awaitIo(js, kj::mv(nativeRequest.response), |
| 114 | [httpClient = kj::mv(httpClient), &context, traceContext = kj::mv(traceContext)]( |
| 115 | jsg::Lock& js, |
| 116 | kj::HttpClient::Response&& response) mutable -> jsg::Optional<jsg::Ref<Response>> { |
| 117 | response.body = response.body.attach(kj::mv(httpClient)); |
| 118 | |
| 119 | traceContext.setTag( |
| 120 | "cache.response.status_code"_kjc, static_cast<int64_t>(response.statusCode)); |
| 121 | KJ_IF_SOME(length, response.body->tryGetLength()) { |
| 122 | traceContext.setTag("cache.response.body.size"_kjc, static_cast<int64_t>(length)); |
| 123 | } |
| 124 | |
| 125 | kj::StringPtr cacheStatus; |
| 126 | KJ_IF_SOME(cs, response.headers->get(context.getHeaderIds().cfCacheStatus)) { |
| 127 | cacheStatus = cs; |
| 128 | traceContext.setTag("cache.response.cache_status"_kjc, cacheStatus); |
| 129 | } else { |
| 130 | // This is an internal error representing a violation of the contract between us and |
| 131 | // the cache. Since it is always conformant to return undefined from Cache::match() |
| 132 | // (because we are allowed to evict any asset at any time), we don't really need to make the |
| 133 | // script fail. However, it might be indicative of a larger problem, and should be |
| 134 | // investigated. |
| 135 | LOG_CACHE_ERROR_ONCE("Response to Cache API GET has no CF-Cache-Status: ", response); |
| 136 | traceContext.setTag("cache.response.success"_kjc, false); |
| 137 | return kj::none; |
| 138 | } |
| 139 | |
| 140 | // The status code should be a 504 on cache miss, but we need to rely on CF-Cache-Status |
| 141 | // because someone might cache a 504. |
| 142 | // See https://httpwg.org/specs/rfc7234.html#cache-request-directive.only-if-cached |
| 143 | // |
| 144 | // TODO(cleanup): CACHE-5949 We should never receive EXPIRED or UPDATING responses, but we do. |
| 145 | // We treat them the same as a MISS mostly to keep from blowing up our Sentry reports. |
| 146 | // TODO(someday): If the cache status is EXPIRED and we return undefined here, does a PURGE on |
| 147 | // this URL result in a 200, causing us to return true from Cache::delete_()? If so, that's |
| 148 | // a small inconsistency: we shouldn't have a match failure but a delete success. |
| 149 | if (cacheStatus == "MISS" || cacheStatus == "EXPIRED" || cacheStatus == "UPDATING") { |
| 150 | traceContext.setTag("cache.response.success"_kjc, false); |
| 151 | return kj::none; |
| 152 | } else if (cacheStatus != "HIT") { |
| 153 | // Another internal error. See above comment where we retrieve the CF-Cache-Status header. |
| 154 | LOG_CACHE_ERROR_ONCE("Response to Cache API GET has invalid CF-Cache-Status: ", response); |
| 155 | traceContext.setTag("cache.response.success"_kjc, false); |
| 156 | return kj::none; |
| 157 | } |
| 158 | |
| 159 | traceContext.setTag("cache.response.success"_kjc, true); |
| 160 | return makeHttpResponse(js, kj::HttpMethod::GET, {}, response.statusCode, response.statusText, |
| 161 | *response.headers, kj::mv(response.body), kj::none); |
| 162 | }); |
| 163 | }); |
| 164 | } |
| 165 | |
| 166 | // Send a PUT request to the cache whose URL is the original request URL and whose body is the |
| 167 | // HTTP response we'd like to cache for that request. |
| 168 | // |
| 169 | // The HTTP response in the PUT request body (the "PUT payload") must itself be an HTTP message, |
| 170 | // except that it MUST NOT have chunked encoding applied to it, even if it has a |
| 171 | // Transfer-Encoding: chunked header. To be clear, the PUT request itself may be chunked, but it |
| 172 | // must not have any nested chunked encoding. |
| 173 | // |
| 174 | // In order to extract the response's data to serialize it, we'll need to call |
| 175 | // `jsResponse->send()`, which will properly encode the response's body if a Content-Encoding |
| 176 | // header is present. This means we'll need to create an instance of kj::HttpService::Response. |
| 177 | jsg::Promise<void> Cache::put(jsg::Lock& js, |
| 178 | Request::Info requestOrUrl, |
| 179 | jsg::Ref<Response> jsResponse, |
| 180 | CompatibilityFlags::Reader flags) { |
| 181 | |
| 182 | JSG_REQUIRE( |
| 183 | jsResponse->getType() != "error"_kj, TypeError, "Cache is unable to store an error response"); |
| 184 | |
| 185 | // Fake kj::HttpService::Response implementation that allows us to reuse jsResponse->send() to |
| 186 | // serialize the response (headers + body) in the format needed to serve as the payload of |
| 187 | // our cache PUT request. |
| 188 | class ResponseSerializer final: public kj::HttpService::Response { |
| 189 | public: |
| 190 | struct Payload { |
| 191 | // The serialized form of the response to be cached. This stream itself contains a full |
| 192 | // HTTP response, with headers and body, representing the content of jsResponse to be written |
| 193 | // to the cache. |
| 194 | kj::Own<kj::AsyncInputStream> stream; |
| 195 | |
| 196 | // A promise which resolves once the payload's headers have been written. Normally, this |
| 197 | // couldn't possibly resolve until the body has been written, and jsRepsonse->send() won't |
| 198 | // complete until then -- except if the body is empty, in which case jsResponse->send() may |
| 199 | // return immediately. |
| 200 | kj::Promise<void> writeHeadersPromise; |
| 201 | }; |
| 202 | |
| 203 | Payload getPayload() { |
| 204 | return KJ_ASSERT_NONNULL(kj::mv(payload)); |
| 205 | } |
| 206 | |
| 207 | private: |
| 208 | kj::Own<kj::AsyncOutputStream> send(uint statusCode, |
| 209 | kj::StringPtr statusText, |
| 210 | const kj::HttpHeaders& headers, |
| 211 | kj::Maybe<uint64_t> expectedBodySize) override { |
| 212 | kj::String contentLength; |
| 213 | |
| 214 | kj::StringPtr connectionHeaders[kj::HttpHeaders::CONNECTION_HEADERS_COUNT]; |
| 215 | KJ_IF_SOME(ebs, expectedBodySize) { |
| 216 | contentLength = kj::str(ebs); |
| 217 | connectionHeaders[kj::HttpHeaders::BuiltinIndices::CONTENT_LENGTH] = contentLength; |
| 218 | } else { |
| 219 | connectionHeaders[kj::HttpHeaders::BuiltinIndices::TRANSFER_ENCODING] = "chunked"; |
| 220 | } |
| 221 | |
| 222 | auto serializedHeaders = headers.serializeResponse(statusCode, statusText, connectionHeaders); |
| 223 | |
| 224 | auto expectedPayloadSize = |
| 225 | expectedBodySize.map([&](uint64_t size) { return size + serializedHeaders.size(); }); |
| 226 | |
| 227 | // We want to create an AsyncInputStream that represents the payload, including both headers |
| 228 | // and body. To do this, we'll create a one-way pipe, using the input end of the pipe as |
| 229 | // said stream. This means we have to write the headers, followed by the body, to the output |
| 230 | // end of the pipe. |
| 231 | // |
| 232 | // send() needs to return a stream to which the caller can write the body. Since we need to |
| 233 | // make sure the headers are written first, we'll return a kj::newPromisedStream(), using a |
| 234 | // promise that resolves to the pipe output as soon as the headers are written. |
| 235 | // |
| 236 | // There's a catch: Unfortunately, if the caller doesn't intend to write any body, then they |
| 237 | // will probably drop the return stream immediately. This could prematurely cancel our header |
| 238 | // write. To avoid that, we split the promise and keep a branch in `writeHeadersPromise`, |
| 239 | // which will have to be awaited separately. |
| 240 | auto payloadPipe = kj::newOneWayPipe(expectedPayloadSize); |
| 241 | |
| 242 | static auto constexpr handleHeaders = [](kj::Own<kj::AsyncOutputStream> out, |
| 243 | kj::String serializedHeaders) |
| 244 | -> kj::Promise<kj::Tuple<kj::Own<kj::AsyncOutputStream>, bool>> { |
| 245 | co_await out->write(serializedHeaders.asBytes()); |
| 246 | co_return kj::tuple(kj::mv(out), false); |
| 247 | }; |
| 248 | |
| 249 | auto headersPromises = |
| 250 | handleHeaders(kj::mv(payloadPipe.out), kj::mv(serializedHeaders)).split(); |
| 251 | |
| 252 | payload = Payload{.stream = kj::mv(payloadPipe.in), |
| 253 | .writeHeadersPromise = kj::get<1>(headersPromises).ignoreResult()}; |
| 254 | |
| 255 | return kj::newPromisedStream(kj::mv(kj::get<0>(headersPromises))); |
| 256 | } |
| 257 | |
| 258 | kj::Own<kj::WebSocket> acceptWebSocket(const kj::HttpHeaders&) override { |
| 259 | JSG_FAIL_REQUIRE(TypeError, "Cannot cache WebSocket upgrade response."); |
| 260 | } |
| 261 | |
| 262 | kj::Maybe<Payload> payload; |
| 263 | }; |
| 264 | |
| 265 | // This use of evalNow() is obsoleted by the capture_async_api_throws compatibility flag, but |
| 266 | // we need to keep it here for people who don't have that flag set. |
| 267 | return js.evalNow([&] { |
| 268 | auto jsRequest = Request::coerce(js, kj::mv(requestOrUrl), kj::none); |
| 269 | |
| 270 | auto& context = IoContext::current(); |
| 271 | TraceContext traceContext = context.makeUserTraceSpan("cache_put"_kjc); |
| 272 | |
| 273 | traceContext.setTag("cache.request.url"_kjc, jsRequest->getUrl()); |
| 274 | traceContext.setTag("cache.request.method"_kjc, kj::str(jsRequest->getMethodEnum())); |
| 275 | traceContext.setTag( |
| 276 | "cache.request.payload.status_code"_kjc, static_cast<int64_t>(jsResponse->getStatus())); |
| 277 | |
| 278 | // TODO(conform): Require that jsRequest's url has an http or https scheme. This is only |
| 279 | // important if api::Request is changed to parse its URL eagerly (as required by spec), rather |
| 280 | // than at fetch()-time. |
| 281 | |
| 282 | JSG_REQUIRE(jsRequest->getMethodEnum() == kj::HttpMethod::GET, TypeError, |
| 283 | "Cannot cache response to non-GET request."); |
| 284 | |
| 285 | JSG_REQUIRE(jsResponse->getStatus() != 206, TypeError, |
| 286 | "Cannot cache response to a range request (206 Partial Content)."); |
| 287 | |
| 288 | auto responseHeadersRef = jsResponse->getHeaders(js); |
| 289 | auto cacheControl = responseHeadersRef->getCommon(js, capnp::CommonHeaderName::CACHE_CONTROL); |
| 290 | |
| 291 | KJ_IF_SOME(vary, responseHeadersRef->getCommon(js, capnp::CommonHeaderName::VARY)) { |
| 292 | JSG_REQUIRE(vary.findFirst('*') == kj::none, TypeError, |
| 293 | "Cannot cache response with 'Vary: *' header."); |
| 294 | } |
| 295 | |
| 296 | KJ_IF_SOME(cacheControl, |
| 297 | responseHeadersRef->getCommon(js, capnp::CommonHeaderName::CACHE_CONTROL)) { |
| 298 | traceContext.setTag("cache.request.payload.header.cache_control"_kjc, cacheControl.asPtr()); |
| 299 | } |
| 300 | KJ_IF_SOME(cacheTag, responseHeadersRef->getPtr(js, "cache-tag"_kj)) { |
| 301 | traceContext.setTag("cache.request.payload.header.cache_tag"_kjc, cacheTag.asPtr()); |
| 302 | } |
| 303 | KJ_IF_SOME(etag, responseHeadersRef->getCommon(js, capnp::CommonHeaderName::ETAG)) { |
| 304 | traceContext.setTag("cache.request.payload.header.etag"_kjc, etag.asPtr()); |
| 305 | } |
| 306 | KJ_IF_SOME(expires, responseHeadersRef->getCommon(js, capnp::CommonHeaderName::EXPIRES)) { |
| 307 | traceContext.setTag("cache.request.payload.header.expires"_kjc, expires.asPtr()); |
| 308 | } |
| 309 | KJ_IF_SOME(lastModified, |
| 310 | responseHeadersRef->getCommon(js, capnp::CommonHeaderName::LAST_MODIFIED)) { |
| 311 | traceContext.setTag("cache.request.payload.header.last_modified"_kjc, lastModified.asPtr()); |
| 312 | } |
| 313 | |
| 314 | if (jsResponse->getStatus() == 304) { |
| 315 | // Silently discard 304 status responses to conditional requests. Caching 304s could be a |
| 316 | // source of bugs in a worker, since a worker which blindly stuffs responses from `fetch()` |
| 317 | // into cache could end up caching one, then later respond to non-conditional requests with |
| 318 | // the cached 304. |
| 319 | // |
| 320 | // Unlike the 206 response status check above, we don't throw here because we used to allow |
| 321 | // this behavior. Silently discarding 304s maintains backwards compatibility and is actually |
| 322 | // still spec-conformant. |
| 323 | |
| 324 | if (context.hasWarningHandler()) { |
| 325 | context.logWarning( |
| 326 | "Ignoring attempt to Cache.put() a 304 status response. 304 responses " |
| 327 | "are not meaningful to cache, and a potential source of bugs. Consider validating that " |
| 328 | "the response status is meaningful to cache before calling Cache.put()."); |
| 329 | } |
| 330 | |
| 331 | return js.resolvedPromise(); |
| 332 | } |
| 333 | |
| 334 | ResponseSerializer serializer; |
| 335 | // We need to send the response to our serializer immediately in order to fulfill Cache.put()'s |
| 336 | // contract: the caller should be able to observe that the response body is disturbed as soon |
| 337 | // as put() returns. |
| 338 | auto serializePromise = jsResponse->send(js, serializer, {}, kj::none); |
| 339 | auto payload = serializer.getPayload(); |
| 340 | |
| 341 | KJ_IF_SOME(length, payload.stream->tryGetLength()) { |
| 342 | traceContext.setTag("cache.request.payload.size"_kjc, static_cast<int64_t>(length)); |
| 343 | } |
| 344 | |
| 345 | // Wait for output locks and cache put quota, trying to avoid returning to the KJ event loop |
| 346 | // in the common case where no waits are needed. |
| 347 | // TODO(later): With Cache streams no longer having a size limit enforced by the runtime, |
| 348 | // explore if we can clean up stream serialization too. |
| 349 | jsg::Promise<IoOwn<kj::AsyncInputStream>> startStreamPromise = nullptr; |
| 350 | auto makeCachePutStream = [&context, stream = kj::mv(payload.stream)](jsg::Lock& js) mutable { |
| 351 | return context.makeCachePutStream(js, kj::mv(stream)); |
| 352 | }; |
| 353 | KJ_IF_SOME(p, context.waitForOutputLocksIfNecessary()) { |
| 354 | startStreamPromise = context.awaitIo(js, kj::mv(p), kj::mv(makeCachePutStream)); |
| 355 | } else { |
| 356 | startStreamPromise = makeCachePutStream(js); |
| 357 | } |
| 358 | |
| 359 | return startStreamPromise.then(js, |
| 360 | context.addFunctor( |
| 361 | [this, &context, jsRequest = kj::mv(jsRequest), cacheControl = kj::mv(cacheControl), |
| 362 | serializePromise = kj::mv(serializePromise), |
| 363 | writePayloadHeadersPromise = kj::mv(payload.writeHeadersPromise), |
| 364 | enableCompatFlags = flags.getCacheApiCompatFlags(), |
| 365 | traceContext = kj::mv(traceContext)](jsg::Lock& js, |
| 366 | IoOwn<kj::AsyncInputStream> payloadStream) mutable -> jsg::Promise<void> { |
| 367 | // Make the PUT request to cache. |
| 368 | auto httpClient = getHttpClient( |
| 369 | context, jsRequest->serializeCfBlobJson(js), traceContext, enableCompatFlags); |
| 370 | auto requestHeaders = kj::HttpHeaders(context.getHeaderTable()); |
| 371 | jsRequest->shallowCopyHeadersTo(requestHeaders); |
| 372 | auto nativeRequest = httpClient->request(kj::HttpMethod::PUT, |
| 373 | validateUrl(jsRequest->getUrl()), requestHeaders, payloadStream->tryGetLength()); |
| 374 | |
| 375 | auto pumpRequestBodyPromise = payloadStream->pumpTo(*nativeRequest.body).ignoreResult(); |
| 376 | // NOTE: We don't attach nativeRequest.body here because we want to control its |
| 377 | // destruction timing in the event of an error; see below. |
| 378 | |
| 379 | // The next step is a bit complicated as it occurs in two separate async flows. |
| 380 | // First, we await the serialization promise, then enter "deferred proxying" by issuing |
| 381 | // `KJ_CO_MAGIC BEGIN_DEFERRED_PROXYING` from our coroutine. Everything after that |
| 382 | // `KJ_CO_MAGIC` constitutes the second async flow that actually handles the request and |
| 383 | // response. |
| 384 | // |
| 385 | // Weird: It's important that these objects be torn down in the right order and that the |
| 386 | // DeferredProxy promise is handled separately from the inner promise. |
| 387 | // |
| 388 | // Moreover, there is an interesting property: In the event that `httpClient` is destroyed |
| 389 | // immediately after `bodyStream` (i.e. without returning to the KJ event loop in between), |
| 390 | // and the body is chunked, then the connection will be closed before the terminating chunk |
| 391 | // can be written. This is actually convenient as it allows us to make sure that when we |
| 392 | // bail out due to an error, the cache is able to see that the request was incomplete and |
| 393 | // should therefore not commit the cache entry. |
| 394 | // |
| 395 | // This is a bit of an accident. It would be much better if KJ's AsyncOutputStream had an |
| 396 | // explicit `end()` method to indicate all data had been written successfully, rather than |
| 397 | // just assume so in the destructor. But, that's a major refactor, and it's immediately |
| 398 | // important to us that we don't write incomplete cache entries, so we rely on this hack for |
| 399 | // now. See EW-812 for the broader problem. |
| 400 | // |
| 401 | // A little funky: The process of "serializing" the cache entry payload means reading all the |
| 402 | // data from the payload body stream and writing it to cache. But, the payload body might |
| 403 | // originate from the app's own JavaScript, rather than being the response to some remote |
| 404 | // request. If the stream is JS-backed, then we want to be careful to track "pending events". |
| 405 | // Specifically, if the stream hasn't reported EOF yet, but JavaScript stops executing and |
| 406 | // there is no external I/O that we're waiting for, then we know that the stream will never |
| 407 | // end, and we want to cancel out the IoContext proactively. |
| 408 | // |
| 409 | // If we were to use `context.awaitIo(serializePromise)` here, we'd lose this property, |
| 410 | // because the context would believe that waiting for the stream itself constituted I/O, even |
| 411 | // if the stream is backed by JS. |
| 412 | // |
| 413 | // On the other hand, once the serialization step completes, we need to wait for the cache |
| 414 | // backend to respond. At that point, we *are* awaiting I/O, and want to record that |
| 415 | // correctly. |
| 416 | // |
| 417 | // So basically, we have an asynchronous promise we need to wait for, and for the first part |
| 418 | // of that wait, we don't want to count it as pending I/O, but for the second part, we do. |
| 419 | // How do we accomplish this? |
| 420 | // |
| 421 | // Well, it just so happens that `serializePromise` is a special kind of promise that might |
| 422 | // help us -- it's kj::Promise<DeferredProxy<void>>, a deferred proxy stream promise. This |
| 423 | // is a promise-for-a-promise, with an interesting property: the outer promise is used to |
| 424 | // wait for JavaScript-backed stream events, while the inner promise represents pure external |
| 425 | // I/O. The method context.awaitDeferredProxy() awaits this special kind of promise, and it |
| 426 | // already only counts the inner promise as being external pending I/O. |
| 427 | // |
| 428 | // However, we have some additional work we want to do *after* serializePromise (both parts) |
| 429 | // completes -- additional work that is also external I/O. So how do we handle that? Well... |
| 430 | // we can actually append it to `serializePromise`'s inner promise! Then awaitDeferredProxy() |
| 431 | // will properly treat it as pending I/O, but only *after* the outer promise completes. This |
| 432 | // gets us everything we want. |
| 433 | // |
| 434 | // Hence, what you see here: we first await the serializePromise, then enter deferred proxying |
| 435 | // with our magic `KJ_CO_MAGIC`, then perform all our additional work. Then we |
| 436 | // `awaitDeferredProxy()` the whole thing. |
| 437 | |
| 438 | // Here we handle the promise for the DeferredProxy itself. |
| 439 | static auto constexpr handleSerialize = |
| 440 | [](kj::Promise<DeferredProxy<void>> serialize, kj::Own<kj::HttpClient> httpClient, |
| 441 | kj::Promise<kj::HttpClient::Response> responsePromise, |
| 442 | kj::Own<kj::AsyncOutputStream> bodyStream, kj::Promise<void> pumpRequestBodyPromise, |
| 443 | kj::Promise<void> writePayloadHeadersPromise, |
| 444 | kj::Own<kj::AsyncInputStream> payloadStream, |
| 445 | TraceContext traceContext) -> kj::Promise<DeferredProxy<void>> { |
| 446 | // This is extremely odd and a bit annoying but we have to make sure |
| 447 | // these are destroyed in a particular order due to cross-dependencies |
| 448 | // for each. If the kj::Promise returned by handleSerialize is dropped |
| 449 | // before the co_await serialize completes, then these won't ever be |
| 450 | // moved away into the handleResponse method (which ensures proper |
| 451 | // cleanup order). In such a case, we explicitly layout a cleanup order |
| 452 | // here to make it clear. |
| 453 | // Note: we could do this by ordering the arguments in a particular way, |
| 454 | // or by doing what a previous iteration of this code did and put everything |
| 455 | // into a struct in a particular order but pulling things out like this makes |
| 456 | // what is going on here much more intentional and explicit. |
| 457 | // |
| 458 | // If these are not cleaned up in the right order, there can be subtle |
| 459 | // use-after-free issues reported by asan and certain flows can end up |
| 460 | // hanging. |
| 461 | KJ_DEFER({ |
| 462 | pumpRequestBodyPromise = nullptr; |
| 463 | payloadStream = nullptr; |
| 464 | bodyStream = nullptr; |
| 465 | responsePromise = nullptr; |
| 466 | writePayloadHeadersPromise = nullptr; |
| 467 | httpClient = nullptr; |
| 468 | }); |
| 469 | try { |
| 470 | auto deferred = co_await serialize; |
| 471 | |
| 472 | // With our `serialize` promise having resolved to a DeferredProxy, we can now enter |
| 473 | // deferred proxying ourselves. |
| 474 | KJ_CO_MAGIC BEGIN_DEFERRED_PROXYING; |
| 475 | |
| 476 | co_await deferred.proxyTask; |
| 477 | // Make sure headers get written even if the body was empty -- see comments earlier. |
| 478 | co_await writePayloadHeadersPromise; |
| 479 | // Make sure the request body is done being pumped and had no errors. If serialization |
| 480 | // completed successfully, then this should also complete immediately thereafter. |
| 481 | co_await pumpRequestBodyPromise; |
| 482 | // It is important to destroy the bodyStream before actually waiting on the |
| 483 | // responsePromise to ensure that the terminal chunk is written since the bodyStream |
| 484 | // may only write the terminal chunk in the streams destructor. |
| 485 | bodyStream = nullptr; |
| 486 | payloadStream = nullptr; |
| 487 | auto response = co_await responsePromise; |
| 488 | // We expect to see either 204 (success) or 413 (failure). Any other status code is a |
| 489 | // violation of the contract between us and the cache, and is an internal |
| 490 | // error, which we log. However, there's no need to throw, since the Cache API is an |
| 491 | // ephemeral K/V store, and we never guaranteed the script we'd actually cache anything. |
| 492 | if (response.statusCode != 204 && response.statusCode != 413) { |
| 493 | LOG_CACHE_ERROR_ONCE("Response to Cache API PUT was neither 204 nor 413: ", response); |
| 494 | } else if (response.statusCode == 204) { |
| 495 | traceContext.setTag("cache.response.success"_kjc, true); |
| 496 | } else if (response.statusCode == 413) { |
| 497 | traceContext.setTag("cache.response.success"_kjc, false); |
| 498 | } |
| 499 | } catch (...) { |
| 500 | traceContext.setTag("cache.response.success"_kjc, false); |
| 501 | auto exception = kj::getCaughtExceptionAsKj(); |
| 502 | if (exception.getType() != kj::Exception::Type::DISCONNECTED) { |
| 503 | kj::throwFatalException(kj::mv(exception)); |
| 504 | } |
| 505 | // If the origin or the cache disconnected, we don't treat this as an error, as put() |
| 506 | // doesn't guarantee that it stores anything anyway. |
| 507 | // |
| 508 | // TODO(someday): I (Kenton) don't understand why we'd explicitly want to hide this |
| 509 | // error, even though hiding it is technically not a violation of the contract. To me |
| 510 | // this seems undesirable, especially when it was the origin that failed. The caller |
| 511 | // can always choose to ignore errors if they want (and many do, by passing to |
| 512 | // waitUntil()). However, there is at least one test which depends on this behavior, |
| 513 | // and probably production Workers in the wild, so I'm not changing it for now. |
| 514 | } |
| 515 | }; |
| 516 | |
| 517 | return context.awaitDeferredProxy(js, |
| 518 | handleSerialize(kj::mv(serializePromise), kj::mv(httpClient), |
| 519 | kj::mv(nativeRequest.response), kj::mv(nativeRequest.body), |
| 520 | kj::mv(pumpRequestBodyPromise), kj::mv(writePayloadHeadersPromise), |
| 521 | kj::mv(payloadStream), kj::mv(traceContext))); |
| 522 | })); |
| 523 | }); |
| 524 | } |
| 525 | |
| 526 | jsg::Promise<bool> Cache::delete_(jsg::Lock& js, |
| 527 | Request::Info requestOrUrl, |
| 528 | jsg::Optional<CacheQueryOptions> options, |
| 529 | CompatibilityFlags::Reader flags) { |
| 530 | auto& context = IoContext::current(); |
| 531 | TraceContext traceContext = context.makeUserTraceSpan("cache_delete"_kjc); |
| 532 | |
| 533 | KJ_IF_SOME(o, options) { |
| 534 | KJ_IF_SOME(ignoreMethod, o.ignoreMethod) { |
| 535 | traceContext.setTag("cache.request.ignore_method"_kjc, ignoreMethod); |
| 536 | } |
| 537 | } |
| 538 | |
| 539 | // This use of evalNow() is obsoleted by the capture_async_api_throws compatibility flag, but |
| 540 | // we need to keep it here for people who don't have that flag set. |
| 541 | return js.evalNow([&]() -> jsg::Promise<bool> { |
| 542 | auto jsRequest = Request::coerce(js, kj::mv(requestOrUrl), kj::none); |
| 543 | |
| 544 | traceContext.setTag("cache.request.url"_kjc, jsRequest->getUrl()); |
| 545 | traceContext.setTag("cache.request.method"_kjc, kj::str(jsRequest->getMethodEnum())); |
| 546 | if (!options.orDefault({}).ignoreMethod.orDefault(false) && |
| 547 | jsRequest->getMethodEnum() != kj::HttpMethod::GET) { |
| 548 | return js.resolvedPromise(false); |
| 549 | } |
| 550 | |
| 551 | // Make the PURGE request to cache. |
| 552 | |
| 553 | auto httpClient = getHttpClient( |
| 554 | context, jsRequest->serializeCfBlobJson(js), traceContext, flags.getCacheApiCompatFlags()); |
| 555 | auto requestHeaders = kj::HttpHeaders(context.getHeaderTable()); |
| 556 | jsRequest->shallowCopyHeadersTo(requestHeaders); |
| 557 | // HACK: The cache doesn't permit PURGE requests from the outside world. It does this by |
| 558 | // filtering on X-Real-IP, which can't be set from the outside world. X-Real-IP can, however, |
| 559 | // be set by a Worker when making requests to its own origin, as "spoofing" client IPs to |
| 560 | // your own origin isn't a security flaw. Also, a Worker sending PURGE requests to its own |
| 561 | // origin's cache is not a security flaw (that's what this very API is implementing after |
| 562 | // all) so it all lines up nicely. |
| 563 | requestHeaders.addPtrPtr("X-Real-IP"_kj, "127.0.0.1"_kj); |
| 564 | auto nativeRequest = httpClient->request(kj::HttpMethod::PURGE, |
| 565 | validateUrl(jsRequest->getUrl()), requestHeaders, static_cast<uint64_t>(0)); |
| 566 | |
| 567 | return context.awaitIo(js, kj::mv(nativeRequest.response), |
| 568 | [httpClient = kj::mv(httpClient), traceContext = kj::mv(traceContext)]( |
| 569 | jsg::Lock&, kj::HttpClient::Response&& response) mutable -> bool { |
| 570 | traceContext.setTag( |
| 571 | "cache.response.status_code"_kjc, static_cast<int64_t>(response.statusCode)); |
| 572 | if (response.statusCode == 200) { |
| 573 | traceContext.setTag("cache.response.success"_kjc, true); |
| 574 | return true; |
| 575 | } else if (response.statusCode == 404) { |
| 576 | traceContext.setTag("cache.response.success"_kjc, false); |
| 577 | return false; |
| 578 | } else if (response.statusCode == 429) { |
| 579 | traceContext.setTag("cache.response.success"_kjc, false); |
| 580 | // Throw, but do not log the response to Sentry, as rate-limited subrequests are normal |
| 581 | JSG_FAIL_REQUIRE( |
| 582 | Error, "Unable to delete cached response. Subrequests are being rate-limited."); |
| 583 | } |
| 584 | LOG_CACHE_ERROR_ONCE("Response to Cache API PURGE was neither 200 nor 404: ", response); |
| 585 | JSG_FAIL_REQUIRE(Error, "Unable to delete cached response."); |
| 586 | }); |
| 587 | }); |
| 588 | } |
| 589 | |
| 590 | kj::Own<kj::HttpClient> Cache::getHttpClient(IoContext& context, |
| 591 | kj::Maybe<kj::String> cfBlobJson, |
| 592 | TraceContext& traceContext, |
| 593 | bool enableCompatFlags) { |
| 594 | auto cacheClient = context.getCacheClient(); |
| 595 | auto metadata = CacheClient::SubrequestMetadata{ |
| 596 | .cfBlobJson = kj::mv(cfBlobJson), |
| 597 | .parentSpan = traceContext.getInternalSpanParent(), |
| 598 | .featureFlagsForFl = kj::none, |
| 599 | }; |
| 600 | if (enableCompatFlags) { |
| 601 | metadata.featureFlagsForFl = |
| 602 | mapCopyString(context.getWorker().getIsolate().getFeatureFlagsForFl()); |
| 603 | } |
| 604 | auto httpClient = |
| 605 | cacheName.map([&](kj::String& n) { |
| 606 | return cacheClient->getNamespace(n, kj::mv(metadata)); |
| 607 | }).orDefault([&]() { return cacheClient->getDefault(kj::mv(metadata)); }); |
| 608 | return httpClient.attach(kj::mv(cacheClient)); |
| 609 | } |
| 610 | |
| 611 | // ======================================================================================= |
| 612 | // CacheStorage |
| 613 | |
| 614 | CacheStorage::CacheStorage(jsg::Lock& js): default_(js.alloc<Cache>(kj::none)) {} |
| 615 | |
| 616 | jsg::Promise<jsg::Ref<Cache>> CacheStorage::open(jsg::Lock& js, kj::String cacheName) { |
| 617 | // Set some reasonable limit to prevent scripts from blowing up our control header size. |
| 618 | static constexpr auto MAX_CACHE_NAME_LENGTH = 1024; |
| 619 | JSG_REQUIRE(cacheName.size() < MAX_CACHE_NAME_LENGTH, TypeError, |
| 620 | "Cache name is too long."); // Mah spoon is toooo big. |
| 621 | |
| 622 | return js.resolvedPromise(js.alloc<Cache>(kj::mv(cacheName))); |
| 623 | } |
| 624 | |
| 625 | } // namespace workerd::api |