Skip to content
File

Blob: src/workerd/api/cache.c++

31.3 KB
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#include "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 
14namespace workerd::api {
15 
16// =======================================================================================
17// Cache
18 
19namespace {
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".
38kj::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 
56Cache::Cache(kj::Maybe<kj::String> cacheName): cacheName(kj::mv(cacheName)) {}
57 
58jsg::Unimplemented Cache::add(Request::Info request) {
59 return {};
60}
61 
62jsg::Unimplemented Cache::addAll(kj::Array<Request::Info> requests) {
63 return {};
64}
65 
66jsg::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.
177jsg::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 
526jsg::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 
590kj::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 
614CacheStorage::CacheStorage(jsg::Lock& js): default_(js.alloc<Cache>(kj::none)) {}
615 
616jsg::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