File
Blob: src/workerd/api/http.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 "http.h" |
| 6 | |
| 7 | #include "data-url.h" |
| 8 | #include "headers.h" |
| 9 | #include "queue.h" |
| 10 | #include "sockets.h" |
| 11 | #include "streams/readable-source.h" |
| 12 | #include "system-streams.h" |
| 13 | #include "util.h" |
| 14 | #include "worker-rpc.h" |
| 15 | #include "workerd/jsg/jsvalue.h" |
| 16 | |
| 17 | #include <workerd/io/features.h> |
| 18 | #include <workerd/io/io-context.h> |
| 19 | #include <workerd/jsg/ser.h> |
| 20 | #include <workerd/jsg/url.h> |
| 21 | #include <workerd/util/abortable.h> |
| 22 | #include <workerd/util/entropy.h> |
| 23 | #include <workerd/util/http-util.h> |
| 24 | #include <workerd/util/mimetype.h> |
| 25 | #include <workerd/util/own-util.h> |
| 26 | #include <workerd/util/stream-utils.h> |
| 27 | #include <workerd/util/strings.h> |
| 28 | #include <workerd/util/thread-scopes.h> |
| 29 | |
| 30 | #include <capnp/compat/http-over-capnp.capnp.h> |
| 31 | #include <kj/compat/url.h> |
| 32 | #include <kj/encoding.h> |
| 33 | #include <kj/memory.h> |
| 34 | #include <kj/parse/char.h> |
| 35 | |
| 36 | namespace workerd::api { |
| 37 | |
| 38 | namespace { |
| 39 | Request::CacheMode getCacheModeFromName(kj::StringPtr value) { |
| 40 | if (value == "no-store") return Request::CacheMode::NOSTORE; |
| 41 | if (value == "no-cache") return Request::CacheMode::NOCACHE; |
| 42 | if (value == "reload") return Request::CacheMode::RELOAD; |
| 43 | JSG_FAIL_REQUIRE(TypeError, kj::str("Unsupported cache mode: ", value)); |
| 44 | } |
| 45 | |
| 46 | jsg::Optional<kj::StringPtr> getCacheModeName(Request::CacheMode mode) { |
| 47 | switch (mode) { |
| 48 | case (Request::CacheMode::NONE): |
| 49 | return kj::none; |
| 50 | case (Request::CacheMode::NOCACHE): |
| 51 | return "no-cache"_kj; |
| 52 | case (Request::CacheMode::NOSTORE): |
| 53 | return "no-store"_kj; |
| 54 | case (Request::CacheMode::RELOAD): |
| 55 | return "reload"_kj; |
| 56 | } |
| 57 | KJ_UNREACHABLE; |
| 58 | } |
| 59 | |
| 60 | } // namespace |
| 61 | |
| 62 | // ----------------------------------------------------------------------------- |
| 63 | // serialization of headers |
| 64 | // |
| 65 | // http-over-capnp.capnp has a nice list of common header names, taken from the HTTP/2 standard. |
| 66 | // We'll use it as an optimization. |
| 67 | // |
| 68 | // Note that using numeric IDs for headers implies we lose the original capitalization. However, |
| 69 | // the JS Headers API doesn't actually give the application any way to observe the capitalization |
| 70 | // of header names -- it only becomes relevant when serializing over HTTP/1.1. And at that point, |
| 71 | // we are actually free to change the capitalization anyway, and we commonly do (KJ itself will |
| 72 | // normalize capitalization of all registered headers, and http-over-capnp also loses |
| 73 | // capitalization). So, it's certainly not worth it to try to keep the original capitalization |
| 74 | // across serialization. |
| 75 | |
| 76 | Body::Buffer Body::Buffer::clone(jsg::Lock& js) { |
| 77 | Buffer result; |
| 78 | result.view = view; |
| 79 | KJ_SWITCH_ONEOF(ownBytes) { |
| 80 | KJ_CASE_ONEOF(ref, jsg::JsRef<jsg::JsBufferSource>) { |
| 81 | result.ownBytes = ref.addRef(js); |
| 82 | } |
| 83 | KJ_CASE_ONEOF(refcounted, kj::Own<RefcountedBytes>) { |
| 84 | result.ownBytes = kj::addRef(*refcounted); |
| 85 | } |
| 86 | KJ_CASE_ONEOF(blob, jsg::Ref<Blob>) { |
| 87 | result.ownBytes = blob.addRef(); |
| 88 | } |
| 89 | } |
| 90 | return result; |
| 91 | } |
| 92 | |
| 93 | Body::ExtractedBody::ExtractedBody( |
| 94 | jsg::Ref<ReadableStream> stream, kj::Maybe<Buffer> buffer, kj::Maybe<kj::String> contentType) |
| 95 | : impl{kj::mv(stream), kj::mv(buffer)}, |
| 96 | contentType(kj::mv(contentType)) { |
| 97 | // This check is in the constructor rather than `extractBody()`, because we often construct |
| 98 | // ExtractedBodys from ReadableStreams directly. |
| 99 | JSG_REQUIRE(!impl.stream->isDisturbed(), TypeError, |
| 100 | "This ReadableStream is disturbed (has already been read from), and cannot " |
| 101 | "be used as a body."); |
| 102 | } |
| 103 | |
| 104 | Body::ExtractedBody Body::extractBody(jsg::Lock& js, Initializer init) { |
| 105 | Buffer buffer; |
| 106 | kj::Maybe<kj::String> contentType; |
| 107 | |
| 108 | KJ_SWITCH_ONEOF(init) { |
| 109 | KJ_CASE_ONEOF(stream, jsg::Ref<ReadableStream>) { |
| 110 | return kj::mv(stream); |
| 111 | } |
| 112 | KJ_CASE_ONEOF(gen, jsg::AsyncGeneratorIgnoringStrings<jsg::Value>) { |
| 113 | return ReadableStream::from(js, gen.release()); |
| 114 | } |
| 115 | KJ_CASE_ONEOF(text, kj::String) { |
| 116 | contentType = kj::str(MimeType::PLAINTEXT_STRING); |
| 117 | buffer = kj::mv(text); |
| 118 | } |
| 119 | KJ_CASE_ONEOF(bytesRef, jsg::JsRef<jsg::JsBufferSource>) { |
| 120 | // Per the Fetch spec we must copy the input buffer here. Beyond spec conformance, this |
| 121 | // fixes a UAF: the incoming data may alias a v8::BackingStore whose underlying memory can |
| 122 | // be freed if the original ArrayBuffer is detached and transferred (e.g. via structuredClone |
| 123 | // with a transfer list) and then garbage collected. This applies to both resizable and |
| 124 | // fixed-size buffers. Copying severs the dependency on the V8 backing store. |
| 125 | buffer = kj::heapArray(bytesRef.getHandle(js).asArrayPtr()); |
| 126 | } |
| 127 | KJ_CASE_ONEOF(blob, jsg::Ref<Blob>) { |
| 128 | // Blobs always have a type, but it defaults to an empty string. We should NOT set |
| 129 | // Content-Type when the blob type is empty. |
| 130 | kj::StringPtr blobType = blob->getType(); |
| 131 | if (blobType != nullptr) { |
| 132 | contentType = kj::str(blobType); |
| 133 | } |
| 134 | buffer = kj::mv(blob); |
| 135 | } |
| 136 | KJ_CASE_ONEOF(formData, jsg::Ref<FormData>) { |
| 137 | // Make an array of characters containing random hexadecimal digits. |
| 138 | // |
| 139 | // Note: Rather than use random hex digits, we could generate the hex digits by hashing the |
| 140 | // form-data content itself! This would give us pleasing assurance that our boundary string |
| 141 | // is not present in the content being divided. The downside is CPU usage if, say, a user |
| 142 | // uploads an enormous file. |
| 143 | kj::FixedArray<kj::byte, 16> boundaryBuffer; |
| 144 | workerd::getEntropy(boundaryBuffer); |
| 145 | auto boundary = kj::encodeHex(boundaryBuffer); |
| 146 | contentType = MimeType::formDataWithBoundary(boundary); |
| 147 | buffer = formData->serialize(boundary); |
| 148 | } |
| 149 | KJ_CASE_ONEOF(searchParams, jsg::Ref<URLSearchParams>) { |
| 150 | contentType = MimeType::formUrlEncodedWithCharset("UTF-8"_kj); |
| 151 | buffer = searchParams->toString(); |
| 152 | } |
| 153 | KJ_CASE_ONEOF(searchParams, jsg::Ref<url::URLSearchParams>) { |
| 154 | contentType = MimeType::formUrlEncodedWithCharset("UTF-8"_kj); |
| 155 | buffer = searchParams->toString(); |
| 156 | } |
| 157 | } |
| 158 | |
| 159 | auto buf = buffer.clone(js); |
| 160 | |
| 161 | // We use streams::newMemorySource() here rather than newSystemStream() wrapping a |
| 162 | // newMemoryInputStream() because we do NOT want deferred proxying for bodies with |
| 163 | // V8 heap provenance. Some buffer types (e.g. Blob data) may reference V8 heap memory |
| 164 | // and we must ensure the data is consumed and destroyed while under the isolate lock, |
| 165 | // which means deferred proxying is not allowed. |
| 166 | auto rs = streams::newMemorySource(buf.view, kj::heap(kj::mv(buf.ownBytes))); |
| 167 | |
| 168 | return {js.alloc<ReadableStream>(IoContext::current(), kj::mv(rs)), kj::mv(buffer), |
| 169 | kj::mv(contentType)}; |
| 170 | } |
| 171 | |
| 172 | Body::Body(jsg::Lock& js, kj::Maybe<ExtractedBody> init, Headers& headers) |
| 173 | : impl(kj::mv(init).map([&headers](auto i) -> Impl { |
| 174 | KJ_IF_SOME(ct, i.contentType) { |
| 175 | if (!headers.hasCommon(capnp::CommonHeaderName::CONTENT_TYPE)) { |
| 176 | // The spec allows the user to override the Content-Type, if they wish, so we only set |
| 177 | // the Content-Type if it doesn't already exist. |
| 178 | headers.setCommon(capnp::CommonHeaderName::CONTENT_TYPE, kj::mv(ct)); |
| 179 | } else KJ_IF_SOME(parsed, MimeType::tryParse(ct)) { |
| 180 | if (MimeType::FORM_DATA == parsed) { |
| 181 | // Custom content-type request/responses with FormData are broken since they require a |
| 182 | // boundary parameter only the FormData serializer can provide. Let's warn if a dev does this. |
| 183 | IoContext::current().logWarning( |
| 184 | "A FormData body was provided with a custom Content-Type header when constructing " |
| 185 | "a Request or Response object. This will prevent the recipient of the Request or " |
| 186 | "Response from being able to parse the body. Consider omitting the custom " |
| 187 | "Content-Type header."); |
| 188 | } |
| 189 | } |
| 190 | } |
| 191 | return kj::mv(i.impl); |
| 192 | })), |
| 193 | headersRef(headers) {} |
| 194 | |
| 195 | kj::Maybe<Body::Buffer> Body::getBodyBuffer(jsg::Lock& js) { |
| 196 | KJ_IF_SOME(i, impl) { |
| 197 | KJ_IF_SOME(b, i.buffer) { |
| 198 | return b.clone(js); |
| 199 | } |
| 200 | } |
| 201 | return kj::none; |
| 202 | } |
| 203 | |
| 204 | bool Body::canRewindBody() { |
| 205 | KJ_IF_SOME(i, impl) { |
| 206 | // We can only rewind buffer-backed bodies. |
| 207 | return i.buffer != kj::none; |
| 208 | } |
| 209 | // Null bodies are trivially "rewindable". |
| 210 | return true; |
| 211 | } |
| 212 | |
| 213 | void Body::rewindBody(jsg::Lock& js) { |
| 214 | KJ_DASSERT(canRewindBody()); |
| 215 | |
| 216 | KJ_IF_SOME(i, impl) { |
| 217 | auto bufferCopy = KJ_ASSERT_NONNULL(i.buffer).clone(js); |
| 218 | |
| 219 | // We use streams::newMemorySource() here rather than newSystemStream() wrapping a |
| 220 | // newMemoryInputStream() because we do NOT want deferred proxying for bodies with |
| 221 | // V8 heap provenance. Specifically, the bufferCopy.view here, while being a kj::ArrayPtr, |
| 222 | // will typically be wrapping a v8::BackingStore, and we must ensure that is is consumed |
| 223 | // and destroyed while under the isolate lock, which means deferred proxying is not allowed. |
| 224 | auto rs = streams::newMemorySource(bufferCopy.view, kj::heap(kj::mv(bufferCopy.ownBytes))); |
| 225 | i.stream = js.alloc<ReadableStream>(IoContext::current(), kj::mv(rs)); |
| 226 | } |
| 227 | } |
| 228 | |
| 229 | void Body::nullifyBody() { |
| 230 | impl = kj::none; |
| 231 | } |
| 232 | |
| 233 | kj::Maybe<jsg::Ref<ReadableStream>> Body::getBody() { |
| 234 | KJ_IF_SOME(i, impl) { |
| 235 | return i.stream.addRef(); |
| 236 | } |
| 237 | return kj::none; |
| 238 | } |
| 239 | bool Body::getBodyUsed() { |
| 240 | KJ_IF_SOME(i, impl) { |
| 241 | return i.stream->isDisturbed(); |
| 242 | } |
| 243 | return false; |
| 244 | } |
| 245 | jsg::Promise<jsg::BufferSource> Body::arrayBuffer(jsg::Lock& js) { |
| 246 | KJ_IF_SOME(i, impl) { |
| 247 | return js.evalNow([&] { |
| 248 | JSG_REQUIRE(!i.stream->isDisturbed(), TypeError, |
| 249 | "Body has already been used. " |
| 250 | "It can only be used once. Use tee() first if you need to read it twice."); |
| 251 | return i.stream->getController().readAllBytes( |
| 252 | js, IoContext::current().getLimitEnforcer().getBufferingLimit()); |
| 253 | }); |
| 254 | } |
| 255 | |
| 256 | // If there's no body, we just return an empty array. |
| 257 | // See https://fetch.spec.whatwg.org/#concept-body-consume-body |
| 258 | auto backing = jsg::BackingStore::alloc<v8::ArrayBuffer>(js, 0); |
| 259 | return js.resolvedPromise(jsg::BufferSource(js, kj::mv(backing))); |
| 260 | } |
| 261 | |
| 262 | jsg::Promise<jsg::BufferSource> Body::bytes(jsg::Lock& js) { |
| 263 | return arrayBuffer(js).then(js, |
| 264 | [](jsg::Lock& js, jsg::BufferSource data) { return data.getTypedView<v8::Uint8Array>(js); }); |
| 265 | } |
| 266 | |
| 267 | jsg::Promise<kj::String> Body::text(jsg::Lock& js) { |
| 268 | KJ_IF_SOME(i, impl) { |
| 269 | return js.evalNow([&] { |
| 270 | JSG_REQUIRE(!i.stream->isDisturbed(), TypeError, |
| 271 | "Body has already been used. " |
| 272 | "It can only be used once. Use tee() first if you need to read it twice."); |
| 273 | |
| 274 | // A common mistake is to call .text() on non-text content, e.g. because you're implementing a |
| 275 | // search-and-replace across your whole site and you forgot that it'll apply to images too. |
| 276 | // When running with a warning handler, let's warn the developer if they do this. |
| 277 | auto& context = IoContext::current(); |
| 278 | if (context.hasWarningHandler()) { |
| 279 | KJ_IF_SOME(type, headersRef.getCommon(js, capnp::CommonHeaderName::CONTENT_TYPE)) { |
| 280 | maybeWarnIfNotText(js, type); |
| 281 | } |
| 282 | } |
| 283 | |
| 284 | return i.stream->getController().readAllText( |
| 285 | js, context.getLimitEnforcer().getBufferingLimit()); |
| 286 | }); |
| 287 | } |
| 288 | |
| 289 | // If there's no body, we just return an empty string. |
| 290 | // See https://fetch.spec.whatwg.org/#concept-body-consume-body |
| 291 | return js.resolvedPromise(kj::String()); |
| 292 | } |
| 293 | |
| 294 | jsg::Promise<jsg::Ref<FormData>> Body::formData(jsg::Lock& js) { |
| 295 | auto formData = js.alloc<FormData>(); |
| 296 | |
| 297 | return js.evalNow([&] { |
| 298 | JSG_REQUIRE(!getBodyUsed(), TypeError, |
| 299 | "Body has already been used. " |
| 300 | "It can only be used once. Use tee() first if you need to read it twice."); |
| 301 | |
| 302 | auto contentType = |
| 303 | JSG_REQUIRE_NONNULL(headersRef.getCommon(js, capnp::CommonHeaderName::CONTENT_TYPE), |
| 304 | TypeError, "Parsing a Body as FormData requires a Content-Type header."); |
| 305 | |
| 306 | KJ_IF_SOME(i, impl) { |
| 307 | KJ_ASSERT(!i.stream->isDisturbed()); |
| 308 | auto& context = IoContext::current(); |
| 309 | return i.stream->getController() |
| 310 | .readAllText(js, context.getLimitEnforcer().getBufferingLimit()) |
| 311 | .then(js, |
| 312 | [contentType = kj::mv(contentType), formData = kj::mv(formData)]( |
| 313 | auto& js, kj::String rawText) mutable { |
| 314 | formData->parse(js, kj::mv(rawText), contentType, |
| 315 | !FeatureFlags::get(js).getFormDataParserSupportsFiles()); |
| 316 | return kj::mv(formData); |
| 317 | }); |
| 318 | } |
| 319 | |
| 320 | // Theoretically, we already know if this will throw: the empty string is a valid |
| 321 | // application/x-www-form-urlencoded body, but not multipart/form-data. However, best to let |
| 322 | // FormData::parse() make the decision, to keep the logic in one place. |
| 323 | formData->parse( |
| 324 | js, kj::String(), contentType, !FeatureFlags::get(js).getFormDataParserSupportsFiles()); |
| 325 | return js.resolvedPromise(kj::mv(formData)); |
| 326 | }); |
| 327 | } |
| 328 | |
| 329 | jsg::Promise<jsg::Value> Body::json(jsg::Lock& js) { |
| 330 | return text(js).then(js, [](jsg::Lock& js, kj::String text) { return js.parseJson(text); }); |
| 331 | } |
| 332 | |
| 333 | jsg::Promise<jsg::Ref<Blob>> Body::blob(jsg::Lock& js) { |
| 334 | return arrayBuffer(js).then(js, [this](jsg::Lock& js, jsg::BufferSource buffer) { |
| 335 | kj::String contentType = headersRef.getCommon(js, capnp::CommonHeaderName::CONTENT_TYPE) |
| 336 | .map([](auto&& b) -> kj::String { |
| 337 | return kj::mv(b); |
| 338 | }).orDefault(nullptr); |
| 339 | |
| 340 | if (FeatureFlags::get(js).getBlobStandardMimeType()) { |
| 341 | contentType = MimeType::extract(contentType) |
| 342 | .map([](MimeType&& mt) -> kj::String { |
| 343 | return mt.toString(); |
| 344 | }).orDefault(nullptr); |
| 345 | } |
| 346 | |
| 347 | return js.alloc<Blob>(js, buffer.getJsHandle(js), kj::mv(contentType)); |
| 348 | }); |
| 349 | } |
| 350 | |
| 351 | kj::Maybe<Body::ExtractedBody> Body::clone(jsg::Lock& js) { |
| 352 | KJ_IF_SOME(i, impl) { |
| 353 | auto branches = i.stream->tee(js); |
| 354 | |
| 355 | i.stream = kj::mv(branches[0]); |
| 356 | |
| 357 | return ExtractedBody{kj::mv(branches[1]), i.buffer.map([&](Buffer& b) { return b.clone(js); })}; |
| 358 | } |
| 359 | |
| 360 | return kj::none; |
| 361 | } |
| 362 | |
| 363 | // ======================================================================================= |
| 364 | |
| 365 | jsg::Ref<Request> Request::coerce( |
| 366 | jsg::Lock& js, Request::Info input, jsg::Optional<Request::Initializer> init) { |
| 367 | return input.is<jsg::Ref<Request>>() && init == kj::none |
| 368 | ? kj::mv(input.get<jsg::Ref<Request>>()) |
| 369 | : Request::constructor(js, kj::mv(input), kj::mv(init)); |
| 370 | } |
| 371 | |
| 372 | jsg::Optional<kj::StringPtr> Request::getCache(jsg::Lock& js) { |
| 373 | return getCacheModeName(cacheMode); |
| 374 | } |
| 375 | Request::CacheMode Request::getCacheMode() { |
| 376 | return cacheMode; |
| 377 | } |
| 378 | |
| 379 | jsg::Ref<Request> Request::constructor( |
| 380 | jsg::Lock& js, Request::Info input, jsg::Optional<Request::Initializer> init) { |
| 381 | kj::String url; |
| 382 | kj::HttpMethod method = kj::HttpMethod::GET; |
| 383 | kj::Maybe<jsg::Ref<Headers>> headers; |
| 384 | kj::Maybe<jsg::Ref<Fetcher>> fetcher; |
| 385 | kj::Maybe<jsg::Ref<AbortSignal>> signal; |
| 386 | CfProperty cf; |
| 387 | kj::Maybe<Body::ExtractedBody> body; |
| 388 | Redirect redirect = Redirect::FOLLOW; |
| 389 | CacheMode cacheMode = CacheMode::NONE; |
| 390 | Response_BodyEncoding responseBodyEncoding = Response_BodyEncoding::AUTO; |
| 391 | |
| 392 | KJ_SWITCH_ONEOF(input) { |
| 393 | KJ_CASE_ONEOF(u, kj::String) { |
| 394 | url = kj::mv(u); |
| 395 | |
| 396 | // TODO(later): This is rather unfortunate. The original implementation of |
| 397 | // this used non-standard URL parsing in violation of the spec. Unfortunately |
| 398 | // some users have come to depend on the non-standard behavior so we have to |
| 399 | // gate the standard behavior with a compat flag. Ideally we'd just be able to |
| 400 | // use the standard parsed URL throughout all of the code but in order to |
| 401 | // minimize the number of changes, we're going to ultimately end up double |
| 402 | // parsing (and serializing) the URL... here we parse it with the standard |
| 403 | // parser, reserialize it back into a string for the sake of not modifying |
| 404 | // the rest of the implementation. Fortunately the standard parser is fast |
| 405 | // but it would eventually be nice to eliminate the double parsing. |
| 406 | if (FeatureFlags::get(js).getFetchStandardUrl()) { |
| 407 | auto parsed = JSG_REQUIRE_NONNULL( |
| 408 | jsg::Url::tryParse(url.asPtr()), TypeError, kj::str("Invalid URL: ", url)); |
| 409 | url = kj::str(parsed.getHref()); |
| 410 | } |
| 411 | } |
| 412 | KJ_CASE_ONEOF(r, jsg::Ref<Request>) { |
| 413 | // Check to see if we're getting a new body from `init`. If so, we want to ignore `input`'s |
| 414 | // body. Note that this is technically non-conformant behavior, but the spec is broken: |
| 415 | // https://github.com/whatwg/fetch/issues/674 |
| 416 | // |
| 417 | // TODO(cleanup): The body extraction logic is getting difficult to follow with the current |
| 418 | // 2-pass initialization we perform (first `input`, then `init`). It'd be nice to defer |
| 419 | // checks like the one we're avoiding here until the very end, so the `init` pass has a |
| 420 | // chance to override `input`'s members *before* we check if the body we're extracting is |
| 421 | // disturbed. |
| 422 | bool ignoreInputBody = false; |
| 423 | KJ_IF_SOME(i, init) { |
| 424 | KJ_SWITCH_ONEOF(i) { |
| 425 | KJ_CASE_ONEOF(initDict, InitializerDict) { |
| 426 | if (initDict.body != kj::none) { |
| 427 | ignoreInputBody = true; |
| 428 | } |
| 429 | } |
| 430 | KJ_CASE_ONEOF(otherRequest, jsg::Ref<Request>) { |
| 431 | // If our initializer dictionary is another Request object, it will always have a `body` |
| 432 | // property. Even if it's null, we should treat it as an explicit body rewrite. |
| 433 | ignoreInputBody = true; |
| 434 | } |
| 435 | } |
| 436 | } |
| 437 | |
| 438 | jsg::Ref<Request> oldRequest = kj::mv(r); |
| 439 | url = kj::str(oldRequest->getUrl()); |
| 440 | method = oldRequest->method; |
| 441 | headers = js.alloc<Headers>(js, *oldRequest->headers); |
| 442 | cf = oldRequest->cf.deepClone(js); |
| 443 | if (!ignoreInputBody) { |
| 444 | JSG_REQUIRE(!oldRequest->getBodyUsed(), TypeError, |
| 445 | "Cannot reconstruct a Request with a used body."); |
| 446 | KJ_IF_SOME(oldJsBody, oldRequest->getBody()) { |
| 447 | // The stream spec says to "create a proxy" for the passed in readable, which it |
| 448 | // defines generically as creating a TransformStream and using pipeThrough to pass |
| 449 | // the input stream through, giving the TransformStream's readable to the extracted |
| 450 | // body below. We don't need to do that. Instead, we just create a new ReadableStream |
| 451 | // that takes over ownership of the internals of the given stream. The given stream |
| 452 | // is left in a locked/disturbed mode so that it can no longer be used. |
| 453 | body = Body::ExtractedBody((oldJsBody)->detach(js), oldRequest->getBodyBuffer(js)); |
| 454 | } |
| 455 | } |
| 456 | cacheMode = oldRequest->getCacheMode(); |
| 457 | redirect = oldRequest->getRedirectEnum(); |
| 458 | fetcher = oldRequest->getFetcher(); |
| 459 | signal = oldRequest->getSignal(); |
| 460 | } |
| 461 | } |
| 462 | |
| 463 | KJ_IF_SOME(i, init) { |
| 464 | KJ_SWITCH_ONEOF(i) { |
| 465 | KJ_CASE_ONEOF(initDict, InitializerDict) { |
| 466 | KJ_IF_SOME(integrity, initDict.integrity) { |
| 467 | JSG_REQUIRE(integrity.size() == 0, TypeError, |
| 468 | "Subrequest integrity checking is not implemented. " |
| 469 | "The integrity option must be either undefined or an empty string."); |
| 470 | } |
| 471 | |
| 472 | KJ_IF_SOME(m, initDict.method) { |
| 473 | KJ_IF_SOME(code, tryParseHttpMethod(m)) { |
| 474 | method = code; |
| 475 | } else KJ_IF_SOME(code, kj::tryParseHttpMethod(toUpper(m))) { |
| 476 | method = code; |
| 477 | if (!FeatureFlags::get(js).getUpperCaseAllHttpMethods()) { |
| 478 | // This is actually the spec defined behavior. We're expected to only |
| 479 | // upper case get, post, put, delete, head, and options per the spec. |
| 480 | // Other methods, even if they would be recognized if they were uppercased, |
| 481 | // are supposed to be rejected. |
| 482 | // Refs: https://fetch.spec.whatwg.org/#methods |
| 483 | switch (method) { |
| 484 | case kj::HttpMethod::GET: |
| 485 | case kj::HttpMethod::POST: |
| 486 | case kj::HttpMethod::PUT: |
| 487 | case kj::HttpMethod::DELETE: |
| 488 | case kj::HttpMethod::HEAD: |
| 489 | case kj::HttpMethod::OPTIONS: |
| 490 | break; |
| 491 | default: |
| 492 | JSG_FAIL_REQUIRE(TypeError, kj::str("Invalid HTTP method string: ", m)); |
| 493 | } |
| 494 | } |
| 495 | } else { |
| 496 | JSG_FAIL_REQUIRE(TypeError, kj::str("Invalid HTTP method string: ", m)); |
| 497 | } |
| 498 | } |
| 499 | |
| 500 | KJ_IF_SOME(h, initDict.headers) { |
| 501 | headers = Headers::constructor(js, kj::mv(h)); |
| 502 | } |
| 503 | |
| 504 | KJ_IF_SOME(p, initDict.fetcher) { |
| 505 | fetcher = kj::mv(p); |
| 506 | } |
| 507 | |
| 508 | KJ_IF_SOME(s, initDict.signal) { |
| 509 | // Note that since this is an optional-maybe, `s` is type Maybe<AbortSignal>. It could |
| 510 | // be null. But that seems like what we want. If someone doesn't specify `signal` at all, |
| 511 | // they want to inherit the `signal` property from the original request. But if they |
| 512 | // explicitly say `signal: null`, they must want to drop the signal that was on the |
| 513 | // original request. |
| 514 | signal = kj::mv(s); |
| 515 | initDict.signal = kj::none; |
| 516 | } |
| 517 | |
| 518 | KJ_IF_SOME(newCf, initDict.cf) { |
| 519 | // TODO(cleanup): When initDict.cf is updated to use jsg::JsRef instead |
| 520 | // of jsg::V8Ref, we can clean this up a bit further. |
| 521 | auto cloned = newCf.deepClone(js); |
| 522 | cf = CfProperty(js, jsg::JsObject(cloned.getHandle(js))); |
| 523 | } |
| 524 | |
| 525 | KJ_IF_SOME(b, kj::mv(initDict.body).orDefault(kj::none)) { |
| 526 | body = Body::extractBody(js, kj::mv(b)); |
| 527 | JSG_REQUIRE(method != kj::HttpMethod::GET && method != kj::HttpMethod::HEAD, TypeError, |
| 528 | "Request with a GET or HEAD method cannot have a body."); |
| 529 | } |
| 530 | |
| 531 | KJ_IF_SOME(r, initDict.redirect) { |
| 532 | redirect = JSG_REQUIRE_NONNULL(Request::tryParseRedirect(r), TypeError, |
| 533 | "Invalid redirect value, must be one of \"follow\" or \"manual\" (\"error\" won't be " |
| 534 | "implemented since it does not make sense at the edge; use \"manual\" and check the " |
| 535 | "response status code)."); |
| 536 | } |
| 537 | |
| 538 | KJ_IF_SOME(c, initDict.cache) { |
| 539 | cacheMode = getCacheModeFromName(c); |
| 540 | } |
| 541 | |
| 542 | KJ_IF_SOME(e, initDict.encodeResponseBody) { |
| 543 | if (e == "manual"_kj) { |
| 544 | responseBodyEncoding = Response_BodyEncoding::MANUAL; |
| 545 | } else if (e == "automatic"_kj) { |
| 546 | responseBodyEncoding = Response_BodyEncoding::AUTO; |
| 547 | } else { |
| 548 | JSG_FAIL_REQUIRE(TypeError, kj::str("encodeResponseBody: unexpected value: ", e)); |
| 549 | } |
| 550 | } |
| 551 | |
| 552 | if (initDict.method != kj::none || initDict.body != kj::none) { |
| 553 | // We modified at least one of the method or the body. In this case, we enforce the |
| 554 | // spec rule that GET/HEAD requests cannot have bodies. (On the other hand, if neither |
| 555 | // of these fields was modified, but the original Request object that we're rewriting |
| 556 | // already represented a GET/HEAD method with a body, we allow that to pass through. |
| 557 | // We support proxying such requests and rewriting their URL/headers/etc.) |
| 558 | JSG_REQUIRE( |
| 559 | (method != kj::HttpMethod::GET && method != kj::HttpMethod::HEAD) || body == kj::none, |
| 560 | TypeError, "Request with a GET or HEAD method cannot have a body."); |
| 561 | } |
| 562 | } |
| 563 | KJ_CASE_ONEOF(otherRequest, jsg::Ref<Request>) { |
| 564 | method = otherRequest->method; |
| 565 | redirect = otherRequest->redirect; |
| 566 | cacheMode = otherRequest->cacheMode; |
| 567 | responseBodyEncoding = otherRequest->responseBodyEncoding; |
| 568 | fetcher = otherRequest->getFetcher(); |
| 569 | signal = otherRequest->getSignal(); |
| 570 | headers = js.alloc<Headers>(js, *otherRequest->headers); |
| 571 | cf = otherRequest->cf.deepClone(js); |
| 572 | KJ_IF_SOME(b, otherRequest->getBody()) { |
| 573 | // Note that unlike when `input` (Request ctor's 1st parameter) is a Request object, here |
| 574 | // we're NOT stealing the other request's body, because we're supposed to pretend that the |
| 575 | // other request is just a dictionary. |
| 576 | body = Body::ExtractedBody(kj::mv(b)); |
| 577 | } |
| 578 | } |
| 579 | } |
| 580 | } |
| 581 | |
| 582 | if (headers == kj::none) { |
| 583 | headers = js.alloc<Headers>(); |
| 584 | } |
| 585 | |
| 586 | // TODO(conform): If `init` has a keepalive flag, pass it to the Body constructor. |
| 587 | return js.alloc<Request>(js, method, url, redirect, KJ_ASSERT_NONNULL(kj::mv(headers)), |
| 588 | kj::mv(fetcher), kj::mv(signal), kj::mv(cf), kj::mv(body), /* thisSignal */ kj::none, |
| 589 | cacheMode, responseBodyEncoding); |
| 590 | } |
| 591 | |
| 592 | jsg::Ref<Request> Request::clone(jsg::Lock& js) { |
| 593 | auto headersClone = headers->clone(js); |
| 594 | |
| 595 | auto cfClone = cf.deepClone(js); |
| 596 | auto bodyClone = Body::clone(js); |
| 597 | |
| 598 | return js.alloc<Request>(js, method, url, redirect, kj::mv(headersClone), getFetcher(), |
| 599 | /* signal */ getSignal(), kj::mv(cfClone), kj::mv(bodyClone), /* thisSignal */ kj::none, |
| 600 | cacheMode, responseBodyEncoding); |
| 601 | } |
| 602 | |
| 603 | kj::StringPtr Request::getMethod() { |
| 604 | return kj::toCharSequence(method); |
| 605 | } |
| 606 | kj::StringPtr Request::getUrl() { |
| 607 | return url; |
| 608 | } |
| 609 | jsg::Ref<Headers> Request::getHeaders(jsg::Lock& js) { |
| 610 | return headers.addRef(); |
| 611 | } |
| 612 | kj::StringPtr Request::getRedirect() { |
| 613 | // TODO(cleanup): Web IDL enum <-> JS string conversion boilerplate is a common need and could be |
| 614 | // factored out. |
| 615 | |
| 616 | switch (redirect) { |
| 617 | case Redirect::FOLLOW: |
| 618 | return "follow"; |
| 619 | case Redirect::MANUAL: |
| 620 | return "manual"; |
| 621 | } |
| 622 | |
| 623 | KJ_UNREACHABLE; |
| 624 | } |
| 625 | kj::Maybe<jsg::Ref<Fetcher>> Request::getFetcher() { |
| 626 | return fetcher.map([](jsg::Ref<Fetcher>& f) { return f.addRef(); }); |
| 627 | } |
| 628 | kj::Maybe<jsg::Ref<AbortSignal>> Request::getSignal() { |
| 629 | return signal.map([](jsg::Ref<AbortSignal>& s) { return s.addRef(); }); |
| 630 | } |
| 631 | |
| 632 | jsg::Optional<jsg::JsObject> Request::getCf(jsg::Lock& js) { |
| 633 | return cf.get(js); |
| 634 | } |
| 635 | |
| 636 | // If signal is given, getThisSignal returns a reference to it. |
| 637 | // Otherwise, we lazily create a new never-aborts AbortSignal that will not |
| 638 | // be used for anything because the spec wills it so. |
| 639 | // Note: To be pedantic, the spec actually calls for us to create a |
| 640 | // second AbortSignal in addition to the one being passed in, but |
| 641 | // that's a bit silly and unnecessary. |
| 642 | // The name "thisSignal" is derived from the fetch spec, which draws a |
| 643 | // distinction between the "signal" and "this' signal". |
| 644 | jsg::Ref<AbortSignal> Request::getThisSignal(jsg::Lock& js) { |
| 645 | KJ_IF_SOME(s, signal) { |
| 646 | return s.addRef(); |
| 647 | } |
| 648 | KJ_IF_SOME(s, thisSignal) { |
| 649 | return s.addRef(); |
| 650 | } |
| 651 | auto newSignal = js.alloc<AbortSignal>(kj::none, kj::none, AbortSignal::Flag::NEVER_ABORTS); |
| 652 | thisSignal = newSignal.addRef(); |
| 653 | return newSignal; |
| 654 | } |
| 655 | |
| 656 | void Request::clearSignalIfIgnoredForSubrequest(jsg::Lock& js) { |
| 657 | KJ_IF_SOME(s, signal) { |
| 658 | if (s->isIgnoredForSubrequests(js)) { |
| 659 | signal = kj::none; |
| 660 | } |
| 661 | } |
| 662 | } |
| 663 | |
| 664 | kj::Maybe<Request::Redirect> Request::tryParseRedirect(kj::StringPtr redirect) { |
| 665 | if (strcasecmp(redirect.cStr(), "follow") == 0) { |
| 666 | return Redirect::FOLLOW; |
| 667 | } |
| 668 | if (strcasecmp(redirect.cStr(), "manual") == 0) { |
| 669 | return Redirect::MANUAL; |
| 670 | } |
| 671 | return kj::none; |
| 672 | } |
| 673 | |
| 674 | void Request::shallowCopyHeadersTo(kj::HttpHeaders& out) { |
| 675 | headers->shallowCopyTo(out); |
| 676 | } |
| 677 | |
| 678 | kj::Maybe<kj::String> Request::serializeCfBlobJson(jsg::Lock& js) { |
| 679 | // We need to clone the cf object if we're going to modify it. We modify it when: |
| 680 | // 1. cacheMode != NONE (existing behavior: map cache option to cacheTtl/cacheLevel/etc.) |
| 681 | // 2. cacheControl is not explicitly set and we need to synthesize it from cacheTtl or cacheMode |
| 682 | // |
| 683 | // For backward compatibility during migration, we dual-write: keep cacheTtl as-is but also |
| 684 | // synthesize cacheControl so downstream services can start consuming the unified field. |
| 685 | // Once downstream fully migrates to cacheControl, cacheTtl can be removed. |
| 686 | |
| 687 | bool hasCacheMode = (cacheMode != CacheMode::NONE); |
| 688 | bool needsSynthesizedCacheControl = false; |
| 689 | |
| 690 | if (!hasCacheMode) { |
| 691 | // Check if cf has cacheTtl but no cacheControl โ we'll need to synthesize cacheControl. |
| 692 | KJ_IF_SOME(cfObj, cf.get(js)) { |
| 693 | if (!cfObj.has(js, "cacheControl") && cfObj.has(js, "cacheTtl")) { |
| 694 | auto ttlVal = cfObj.get(js, "cacheTtl"); |
| 695 | needsSynthesizedCacheControl = !ttlVal.isUndefined(); |
| 696 | } |
| 697 | } |
| 698 | } |
| 699 | |
| 700 | if (!hasCacheMode && !needsSynthesizedCacheControl) { |
| 701 | return cf.serialize(js); |
| 702 | } |
| 703 | |
| 704 | CfProperty clone; |
| 705 | KJ_IF_SOME(obj, cf.get(js)) { |
| 706 | (void)obj; |
| 707 | clone = cf.deepClone(js); |
| 708 | } else { |
| 709 | clone = CfProperty(js, js.obj()); |
| 710 | } |
| 711 | auto obj = KJ_ASSERT_NONNULL(clone.get(js)); |
| 712 | |
| 713 | constexpr int NOCACHE_TTL = -1; |
| 714 | if (hasCacheMode) { |
| 715 | switch (cacheMode) { |
| 716 | case CacheMode::NOSTORE: |
| 717 | if (obj.has(js, "cacheTtl")) { |
| 718 | jsg::JsValue oldTtl = obj.get(js, "cacheTtl"); |
| 719 | JSG_REQUIRE(oldTtl.strictEquals(js.num(NOCACHE_TTL)), TypeError, |
| 720 | kj::str("CacheTtl: ", oldTtl, ", is not compatible with cache: ", |
| 721 | getCacheModeName(cacheMode).orDefault("none"_kj), " header.")); |
| 722 | } else { |
| 723 | obj.set(js, "cacheTtl", js.num(NOCACHE_TTL)); |
| 724 | } |
| 725 | KJ_FALLTHROUGH; |
| 726 | case CacheMode::RELOAD: |
| 727 | obj.set(js, "cacheLevel", js.str("bypass"_kjc)); |
| 728 | break; |
| 729 | case CacheMode::NOCACHE: |
| 730 | obj.set(js, "cacheForceRevalidate", js.boolean(true)); |
| 731 | break; |
| 732 | case CacheMode::NONE: |
| 733 | KJ_UNREACHABLE; |
| 734 | } |
| 735 | } |
| 736 | |
| 737 | // Synthesize cacheControl from cacheTtl or cacheMode when cacheControl is not explicitly set. |
| 738 | // This dual-writes both fields so downstream can migrate to cacheControl incrementally. |
| 739 | if (!obj.has(js, "cacheControl")) { |
| 740 | if (hasCacheMode) { |
| 741 | // Synthesize from the cache request option. |
| 742 | switch (cacheMode) { |
| 743 | case CacheMode::NOSTORE: |
| 744 | obj.set(js, "cacheControl", js.str("no-store"_kjc)); |
| 745 | break; |
| 746 | case CacheMode::NOCACHE: |
| 747 | obj.set(js, "cacheControl", js.str("no-cache"_kjc)); |
| 748 | break; |
| 749 | case CacheMode::RELOAD: |
| 750 | break; |
| 751 | case CacheMode::NONE: |
| 752 | KJ_UNREACHABLE; |
| 753 | } |
| 754 | } else if (obj.has(js, "cacheTtl")) { |
| 755 | // Synthesize from cacheTtl value: positive/zero โ max-age=N, -1 โ no-store. |
| 756 | jsg::JsValue ttlVal = obj.get(js, "cacheTtl"); |
| 757 | if (ttlVal.strictEquals(js.num(NOCACHE_TTL))) { |
| 758 | obj.set(js, "cacheControl", js.str("no-store"_kjc)); |
| 759 | } else KJ_IF_SOME(ttlInt, ttlVal.tryCast<jsg::JsInt32>()) { |
| 760 | auto ttl = KJ_ASSERT_NONNULL(ttlInt.value(js)); |
| 761 | obj.set(js, "cacheControl", js.str(kj::str("max-age=", ttl))); |
| 762 | } |
| 763 | } |
| 764 | } |
| 765 | |
| 766 | return clone.serialize(js); |
| 767 | } |
| 768 | |
| 769 | void RequestInitializerDict::validate(jsg::Lock& js) { |
| 770 | KJ_IF_SOME(c, cache) { |
| 771 | // Check compatibility flag |
| 772 | JSG_REQUIRE(FeatureFlags::get(js).getCacheOptionEnabled(), Error, |
| 773 | kj::str("The 'cache' field on 'RequestInitializerDict' is not implemented.")); |
| 774 | |
| 775 | // Validate that the cache type is valid |
| 776 | auto cacheMode = getCacheModeFromName(c); |
| 777 | |
| 778 | bool invalidNoCache = |
| 779 | !FeatureFlags::get(js).getCacheNoCache() && (cacheMode == Request::CacheMode::NOCACHE); |
| 780 | bool invalidReload = |
| 781 | !FeatureFlags::get(js).getCacheReload() && (cacheMode == Request::CacheMode::RELOAD); |
| 782 | JSG_REQUIRE( |
| 783 | !invalidNoCache && !invalidReload, TypeError, kj::str("Unsupported cache mode: ", c)); |
| 784 | } |
| 785 | |
| 786 | // Validate mutual exclusion of cf.cacheControl with cf.cacheTtl and the cache request option. |
| 787 | // cacheControl provides explicit Cache-Control header override and cannot be combined with |
| 788 | // cacheTtl (which sets a simplified TTL) or the cache option (which maps to cacheTtl internally). |
| 789 | // cacheTtlByStatus is allowed alongside cacheControl since they serve different purposes. |
| 790 | KJ_IF_SOME(cfRef, cf) { |
| 791 | auto cfObj = jsg::JsObject(cfRef.getHandle(js)); |
| 792 | if (cfObj.has(js, "cacheControl")) { |
| 793 | auto cacheControlVal = cfObj.get(js, "cacheControl"); |
| 794 | if (!cacheControlVal.isUndefined()) { |
| 795 | // cacheControl + cacheTtl โ throw |
| 796 | if (cfObj.has(js, "cacheTtl")) { |
| 797 | auto cacheTtlVal = cfObj.get(js, "cacheTtl"); |
| 798 | JSG_REQUIRE(cacheTtlVal.isUndefined(), TypeError, |
| 799 | "The 'cacheControl' and 'cacheTtl' options on cf are mutually exclusive. " |
| 800 | "Use 'cacheControl' for explicit Cache-Control header directives, " |
| 801 | "or 'cacheTtl' for a simplified TTL, but not both."); |
| 802 | } |
| 803 | // cacheControl + cache option (no-store/no-cache) โ throw |
| 804 | // The cache request option maps to cacheTtl internally, so they conflict. |
| 805 | JSG_REQUIRE(cache == kj::none, TypeError, |
| 806 | "The 'cacheControl' option on cf cannot be used together with the 'cache' " |
| 807 | "request option. The 'cache' option ('no-store'/'no-cache') maps to cache TTL " |
| 808 | "behavior internally, which conflicts with explicit Cache-Control directives."); |
| 809 | } |
| 810 | } |
| 811 | } |
| 812 | |
| 813 | KJ_IF_SOME(e, encodeResponseBody) { |
| 814 | JSG_REQUIRE(e == "manual"_kj || e == "automatic"_kj, TypeError, |
| 815 | kj::str("encodeResponseBody: unexpected value: ", e)); |
| 816 | } |
| 817 | } |
| 818 | |
| 819 | void Request::serialize(jsg::Lock& js, |
| 820 | jsg::Serializer& serializer, |
| 821 | const jsg::TypeHandler<RequestInitializerDict>& initDictHandler) { |
| 822 | serializer.writeLengthDelimited(url); |
| 823 | |
| 824 | // Our strategy is to construct an initializer dict object and serialize that as a JS object. |
| 825 | // This makes the deserialization end really simple (just call the constructor), and it also |
| 826 | // gives us extensibility: we can add new fields without having to bump the serialization tag. |
| 827 | // clang-format off |
| 828 | serializer.write(js, jsg::JsValue(initDictHandler.wrap(js, RequestInitializerDict{ |
| 829 | // GET is the default, so only serialize the method if it's something else. |
| 830 | .method = method == kj::HttpMethod::GET ? jsg::Optional<kj::String>() : kj::str(method), |
| 831 | |
| 832 | .headers = headers.addRef(), |
| 833 | |
| 834 | .body = getBody().map([](jsg::Ref<ReadableStream> stream) -> Body::Initializer { |
| 835 | // jsg::Ref<ReadableStream> is one of the possible variants of Body::Initializer. |
| 836 | return kj::mv(stream); |
| 837 | }), |
| 838 | |
| 839 | // "manual" is the default for `redirect`, so only encode if it's not that. |
| 840 | .redirect = redirect == Redirect::MANUAL ? kj::str(getRedirect()) |
| 841 | : kj::Maybe<kj::String>(kj::none), |
| 842 | |
| 843 | // We have to ignore .fetcher for serialization. We can't simply fail if a fetcher is present |
| 844 | // because requests received by the top-level fetch handler actually have .fetcher set to |
| 845 | // the hidden "next" binding, which historically could be different from null (although in |
| 846 | // practice these days it is always the same). We obviously want to be able to serialize |
| 847 | // requests received by the top-level fetch handler so... we have to ignore this. This |
| 848 | // property should probably go away in any case. |
| 849 | |
| 850 | .cf = cf.getRef(js), |
| 851 | |
| 852 | .cache = getCacheModeName(cacheMode).map( |
| 853 | [](kj::StringPtr name) -> kj::String { return kj::str(name); }), |
| 854 | |
| 855 | // .mode is unimplemented |
| 856 | // .credentials is unimplemented |
| 857 | // .referrer is unimplemented |
| 858 | // .referrerPolicy is unimplemented |
| 859 | // .integrity is required to be empty |
| 860 | |
| 861 | // If an AbortSignal is present, we'll try to serialize it. As of this writing, AbortSignal |
| 862 | // is not serializable, but we could add support for sending it over RPC in the future. |
| 863 | // |
| 864 | // Note we have to double-Maybe this, so that if no signal is present, the property is absent |
| 865 | // instead of `null`. |
| 866 | .signal = |
| 867 | signal.map([&js](jsg::Ref<AbortSignal>& s) -> kj::Maybe<jsg::Ref<AbortSignal>> { |
| 868 | if (s->isIgnoredForSubrequests(js)) { |
| 869 | return kj::none; |
| 870 | } |
| 871 | |
| 872 | return s.addRef(); |
| 873 | }), |
| 874 | |
| 875 | // Only serialize responseBodyEncoding if it's not the default AUTO |
| 876 | .encodeResponseBody = responseBodyEncoding == Response_BodyEncoding::AUTO |
| 877 | ? jsg::Optional<kj::String>() |
| 878 | : kj::str("manual") |
| 879 | }))); |
| 880 | } |
| 881 | |
| 882 | jsg::Ref<Request> Request::deserialize(jsg::Lock& js, |
| 883 | rpc::SerializationTag tag, |
| 884 | jsg::Deserializer& deserializer, |
| 885 | const jsg::TypeHandler<RequestInitializerDict>& initDictHandler) { |
| 886 | auto url = deserializer.readLengthDelimitedString(); |
| 887 | auto init = KJ_UNWRAP_OR(initDictHandler.tryUnwrap(js, deserializer.readValue(js)), { |
| 888 | JSG_FAIL_REQUIRE(DOMDataCloneError, |
| 889 | "Deserialization failed: could not deserialize Request initializer"); |
| 890 | }); |
| 891 | return Request::constructor(js, kj::mv(url), kj::mv(init)); |
| 892 | } |
| 893 | |
| 894 | // ======================================================================================= |
| 895 | |
| 896 | namespace { |
| 897 | constexpr kj::StringPtr defaultStatusText(uint statusCode) { |
| 898 | // RFC 7231 recommendations, unless otherwise specified. |
| 899 | // https://tools.ietf.org/html/rfc7231#section-6.1 |
| 900 | #define STATUS(code, text) case code: return text##_kj |
| 901 | switch (statusCode) { |
| 902 | // Status code 0 is used exclusively with error responses |
| 903 | // created using Response.error() |
| 904 | STATUS(0, ""); |
| 905 | STATUS(100, "Continue"); |
| 906 | STATUS(101, "Switching Protocols"); |
| 907 | STATUS(102, "Processing"); // RFC 2518, WebDAV |
| 908 | STATUS(103, "Early Hints"); // RFC 8297 |
| 909 | STATUS(200, "OK"); |
| 910 | STATUS(201, "Created"); |
| 911 | STATUS(202, "Accepted"); |
| 912 | STATUS(203, "Non-Authoritative Information"); |
| 913 | STATUS(204, "No Content"); |
| 914 | STATUS(205, "Reset Content"); |
| 915 | STATUS(206, "Partial Content"); |
| 916 | STATUS(207, "Multi-Status"); // RFC 4918, WebDAV |
| 917 | STATUS(208, "Already Reported"); // RFC 5842, WebDAV |
| 918 | STATUS(226, "IM Used"); // RFC 3229 |
| 919 | STATUS(300, "Multiple Choices"); |
| 920 | STATUS(301, "Moved Permanently"); |
| 921 | STATUS(302, "Found"); |
| 922 | STATUS(303, "See Other"); |
| 923 | STATUS(304, "Not Modified"); |
| 924 | STATUS(305, "Use Proxy"); |
| 925 | |
| 926 | STATUS(307, "Temporary Redirect"); |
| 927 | STATUS(308, "Permanent Redirect"); // RFC 7538 |
| 928 | STATUS(400, "Bad Request"); |
| 929 | STATUS(401, "Unauthorized"); |
| 930 | STATUS(402, "Payment Required"); |
| 931 | STATUS(403, "Forbidden"); |
| 932 | STATUS(404, "Not Found"); |
| 933 | STATUS(405, "Method Not Allowed"); |
| 934 | STATUS(406, "Not Acceptable"); |
| 935 | STATUS(407, "Proxy Authentication Required"); |
| 936 | STATUS(408, "Request Timeout"); |
| 937 | STATUS(409, "Conflict"); |
| 938 | STATUS(410, "Gone"); |
| 939 | STATUS(411, "Length Required"); |
| 940 | STATUS(412, "Precondition Failed"); |
| 941 | STATUS(413, "Payload Too Large"); |
| 942 | STATUS(414, "URI Too Long"); |
| 943 | STATUS(415, "Unsupported Media Type"); |
| 944 | STATUS(416, "Range Not Satisfiable"); |
| 945 | STATUS(417, "Expectation Failed"); |
| 946 | STATUS(418, "I'm a teapot"); // RFC 2324 |
| 947 | STATUS(421, "Misdirected Request"); // RFC 7540 |
| 948 | STATUS(422, "Unprocessable Entity"); // RFC 4918, WebDAV |
| 949 | STATUS(423, "Locked"); // RFC 4918, WebDAV |
| 950 | STATUS(424, "Failed Dependency"); // RFC 4918, WebDAV |
| 951 | STATUS(426, "Upgrade Required"); |
| 952 | STATUS(428, "Precondition Required"); // RFC 6585 |
| 953 | STATUS(429, "Too Many Requests"); // RFC 6585 |
| 954 | STATUS(431, "Request Header Fields Too Large"); // RFC 6585 |
| 955 | STATUS(451, "Unavailable For Legal Reasons"); // RFC 7725 |
| 956 | STATUS(500, "Internal Server Error"); |
| 957 | STATUS(501, "Not Implemented"); |
| 958 | STATUS(502, "Bad Gateway"); |
| 959 | STATUS(503, "Service Unavailable"); |
| 960 | STATUS(504, "Gateway Timeout"); |
| 961 | STATUS(505, "HTTP Version Not Supported"); |
| 962 | STATUS(506, "Variant Also Negotiates"); // RFC 2295 |
| 963 | STATUS(507, "Insufficient Storage"); // RFC 4918, WebDAV |
| 964 | STATUS(508, "Loop Detected"); // RFC 5842, WebDAV |
| 965 | STATUS(510, "Not Extended"); // RFC 2774 |
| 966 | STATUS(511, "Network Authentication Required"); // RFC 6585 |
| 967 | default: |
| 968 | // If we don't recognize the status code, check which range it falls into and use the status |
| 969 | // code class defined by RFC 7231, section 6, as the status text. |
| 970 | if (statusCode >= 200 && statusCode < 300) { |
| 971 | return "Successful"_kj; |
| 972 | } else if (statusCode >= 300 && statusCode < 400) { |
| 973 | return "Redirection"_kj; |
| 974 | } else if (statusCode >= 400 && statusCode < 500) { |
| 975 | return "Client Error"_kj; |
| 976 | } else if (statusCode >= 500 && statusCode < 600) { |
| 977 | return "Server Error"_kj; |
| 978 | } else { |
| 979 | return ""_kj; |
| 980 | } |
| 981 | } |
| 982 | #undef STATUS |
| 983 | } |
| 984 | |
| 985 | constexpr bool isNullBodyStatusCode(uint statusCode) { |
| 986 | switch (statusCode) { |
| 987 | // Fetch spec section 2.2.3 defines these status codes as null body statuses: |
| 988 | // https://fetch.spec.whatwg.org/#null-body-status |
| 989 | case 101: |
| 990 | case 204: |
| 991 | case 205: |
| 992 | case 304: |
| 993 | return true; |
| 994 | default: |
| 995 | return false; |
| 996 | } |
| 997 | } |
| 998 | |
| 999 | constexpr bool isRedirectStatusCode(uint statusCode) { |
| 1000 | switch (statusCode) { |
| 1001 | // Fetch spec section 2.2.3 defines these status codes as redirect statuses: |
| 1002 | // https://fetch.spec.whatwg.org/#redirect-status |
| 1003 | case 301: |
| 1004 | case 302: |
| 1005 | case 303: |
| 1006 | case 307: |
| 1007 | case 308: |
| 1008 | return true; |
| 1009 | default: |
| 1010 | return false; |
| 1011 | } |
| 1012 | } |
| 1013 | } // namespace |
| 1014 | |
| 1015 | Response::Response(jsg::Lock& js, |
| 1016 | int statusCode, |
| 1017 | kj::Maybe<kj::String> statusText, |
| 1018 | jsg::Ref<Headers> headers, |
| 1019 | CfProperty&& cf, |
| 1020 | kj::Maybe<Body::ExtractedBody> body, |
| 1021 | kj::Array<kj::String> urlList, |
| 1022 | kj::Maybe<jsg::Ref<WebSocket>> webSocket, |
| 1023 | Response::BodyEncoding bodyEncoding) |
| 1024 | : Body(js, kj::mv(body), *headers), |
| 1025 | statusCode(statusCode), |
| 1026 | statusText(kj::mv(statusText)), |
| 1027 | headers(kj::mv(headers)), |
| 1028 | cf(kj::mv(cf)), |
| 1029 | urlList(kj::mv(urlList)), |
| 1030 | webSocket(kj::mv(webSocket)), |
| 1031 | bodyEncoding(bodyEncoding), |
| 1032 | asyncContext(jsg::AsyncContextFrame::currentRef(js)) {} |
| 1033 | |
| 1034 | jsg::Ref<Response> Response::error(jsg::Lock& js) { |
| 1035 | return js.alloc<Response>(js, 0, kj::none, js.alloc<Headers>(), CfProperty(), kj::none); |
| 1036 | }; |
| 1037 | |
| 1038 | jsg::Ref<Response> Response::constructor(jsg::Lock& js, |
| 1039 | jsg::Optional<kj::Maybe<Body::Initializer>> optionalBodyInit, |
| 1040 | jsg::Optional<Initializer> maybeInit) { |
| 1041 | auto bodyInit = kj::mv(optionalBodyInit).orDefault(kj::none); |
| 1042 | Initializer init = kj::mv(maybeInit).orDefault(InitializerDict()); |
| 1043 | |
| 1044 | int statusCode = 200; |
| 1045 | BodyEncoding bodyEncoding = Response::BodyEncoding::AUTO; |
| 1046 | |
| 1047 | kj::Maybe<kj::String> statusText; |
| 1048 | kj::Maybe<Body::ExtractedBody> body = kj::none; |
| 1049 | jsg::Ref<Headers> headers = nullptr; |
| 1050 | CfProperty cf; |
| 1051 | kj::Maybe<jsg::Ref<WebSocket>> webSocket = kj::none; |
| 1052 | |
| 1053 | KJ_SWITCH_ONEOF(init) { |
| 1054 | KJ_CASE_ONEOF(initDict, InitializerDict) { |
| 1055 | KJ_IF_SOME(status, initDict.status) { |
| 1056 | statusCode = status; |
| 1057 | } |
| 1058 | KJ_IF_SOME(t, initDict.statusText) { |
| 1059 | statusText = kj::mv(t); |
| 1060 | } |
| 1061 | KJ_IF_SOME(v, initDict.encodeBody) { |
| 1062 | if (v == "manual"_kj) { |
| 1063 | bodyEncoding = Response::BodyEncoding::MANUAL; |
| 1064 | } else if (v == "automatic"_kj) { |
| 1065 | bodyEncoding = Response::BodyEncoding::AUTO; |
| 1066 | } else { |
| 1067 | JSG_FAIL_REQUIRE(TypeError, kj::str("encodeBody: unexpected value: ", v)); |
| 1068 | } |
| 1069 | } |
| 1070 | |
| 1071 | KJ_IF_SOME(initHeaders, initDict.headers) { |
| 1072 | headers = Headers::constructor(js, kj::mv(initHeaders)); |
| 1073 | } else { |
| 1074 | headers = js.alloc<Headers>(); |
| 1075 | } |
| 1076 | |
| 1077 | KJ_IF_SOME(newCf, initDict.cf) { |
| 1078 | // TODO(cleanup): When initDict.cf is updated to use jsg::JsRef instead |
| 1079 | // of jsg::V8Ref, we can clean this up a bit further. |
| 1080 | auto cloned = newCf.deepClone(js); |
| 1081 | cf = CfProperty(js, jsg::JsObject(cloned.getHandle(js))); |
| 1082 | } |
| 1083 | |
| 1084 | KJ_IF_SOME(ws, initDict.webSocket) { |
| 1085 | KJ_IF_SOME(ws2, ws) { |
| 1086 | webSocket = ws2.addRef(); |
| 1087 | } |
| 1088 | } |
| 1089 | } |
| 1090 | KJ_CASE_ONEOF(otherResponse, jsg::Ref<Response>) { |
| 1091 | // Note that in a true Fetch-conformant implementation, this entire case is enabled by Web IDL |
| 1092 | // treating objects as dictionaries. However, some of our Response class's properties are |
| 1093 | // jsg::WontImplement, which prevent us from relying on that Web IDL behavior ourselves. |
| 1094 | |
| 1095 | statusCode = otherResponse->statusCode; |
| 1096 | bodyEncoding = otherResponse->bodyEncoding; |
| 1097 | kj::StringPtr otherStatusText = otherResponse->getStatusText(); |
| 1098 | if (otherStatusText != defaultStatusText(statusCode)) { |
| 1099 | statusText = kj::str(otherStatusText); |
| 1100 | } |
| 1101 | headers = js.alloc<Headers>(js, *otherResponse->headers); |
| 1102 | cf = otherResponse->cf.deepClone(js); |
| 1103 | KJ_IF_SOME(otherWs, otherResponse->webSocket) { |
| 1104 | webSocket = otherWs.addRef(); |
| 1105 | } |
| 1106 | } |
| 1107 | } |
| 1108 | |
| 1109 | if (webSocket == kj::none) { |
| 1110 | JSG_REQUIRE(statusCode >= 200 && statusCode <= 599, RangeError, |
| 1111 | "Responses may only be constructed with status codes in the range 200 to 599, inclusive."); |
| 1112 | } else { |
| 1113 | JSG_REQUIRE( |
| 1114 | statusCode == 101, RangeError, "Responses with a WebSocket must have status code 101."); |
| 1115 | } |
| 1116 | |
| 1117 | KJ_IF_SOME(s, statusText) { |
| 1118 | // Disallow control characters (especially \r and \n) in statusText since it could allow |
| 1119 | // header injection. |
| 1120 | // |
| 1121 | // TODO(cleanup): Once this is deployed, update open-source KJ HTTP to do this automatically. |
| 1122 | for (char c: s) { |
| 1123 | if (static_cast<byte>(c) < 0x20u) { |
| 1124 | JSG_FAIL_REQUIRE(TypeError, "Invalid statusText"); |
| 1125 | } |
| 1126 | } |
| 1127 | } |
| 1128 | |
| 1129 | KJ_IF_SOME(bi, bodyInit) { |
| 1130 | body = Body::extractBody(js, kj::mv(bi)); |
| 1131 | if (isNullBodyStatusCode(statusCode)) { |
| 1132 | // TODO(conform): We *should* fail unconditionally here, but during the Workers beta we |
| 1133 | // allowed Responses to have null body statuses with non-null, zero-length bodies. In order |
| 1134 | // not to break anything in production, for now we allow the author to construct a Response |
| 1135 | // with a zero-length buffer, but we give them a console warning. If we can ever verify that |
| 1136 | // no one relies on this behavior, we should remove this non-conformity. |
| 1137 | |
| 1138 | // Fail if the body is not backed by a buffer (i.e., it's an opaque ReadableStream). |
| 1139 | auto& buffer = JSG_REQUIRE_NONNULL(KJ_ASSERT_NONNULL(body).impl.buffer, TypeError, |
| 1140 | "Response with null body status (101, 204, 205, or 304) cannot have a body."); |
| 1141 | |
| 1142 | // Fail if the body is backed by a non-zero-length buffer. |
| 1143 | JSG_REQUIRE(buffer.view.size() == 0, TypeError, |
| 1144 | "Response with null body status (101, 204, 205, or 304) cannot have a body."); |
| 1145 | |
| 1146 | auto& context = IoContext::current(); |
| 1147 | if (context.hasWarningHandler()) { |
| 1148 | context.logWarning(kj::str("Constructing a Response with a null body status (", statusCode, |
| 1149 | ") and a non-null, " |
| 1150 | "zero-length body. This is technically incorrect, and we recommend you update your " |
| 1151 | "code to explicitly pass in a `null` body, e.g. `new Response(null, { status: ", |
| 1152 | statusCode, |
| 1153 | ", ... })`. (We continue to allow the zero-length body behavior because it " |
| 1154 | "was previously the only way to construct a Response with a null body status. This " |
| 1155 | "behavior may change in the future.)")); |
| 1156 | } |
| 1157 | |
| 1158 | // Treat the zero-length body as a null body. |
| 1159 | body = kj::none; |
| 1160 | } |
| 1161 | } |
| 1162 | |
| 1163 | return js.alloc<Response>(js, statusCode, kj::mv(statusText), kj::mv(headers), |
| 1164 | kj::mv(cf), kj::mv(body), nullptr, kj::mv(webSocket), bodyEncoding); |
| 1165 | } |
| 1166 | |
| 1167 | jsg::Ref<Response> Response::redirect(jsg::Lock& js, kj::String url, jsg::Optional<int> status) { |
| 1168 | auto statusCode = status.orDefault(302); |
| 1169 | if (!isRedirectStatusCode(statusCode)) { |
| 1170 | JSG_FAIL_REQUIRE(RangeError, |
| 1171 | kj::str(statusCode, |
| 1172 | " is not a redirect status code. " |
| 1173 | "It must be one of: 301, 302, 303, 307, or 308.")); |
| 1174 | } |
| 1175 | |
| 1176 | // TODO(conform): The URL is supposed to be parsed relative to the "current setting's object's API |
| 1177 | // base URL". |
| 1178 | kj::String parsedUrl = nullptr; |
| 1179 | if (FeatureFlags::get(js).getSpecCompliantResponseRedirect()) { |
| 1180 | auto parsed = JSG_REQUIRE_NONNULL( |
| 1181 | jsg::Url::tryParse(url.asPtr()), TypeError, "Unable to parse URL: ", url); |
| 1182 | parsedUrl = kj::str(parsed.getHref()); |
| 1183 | } else { |
| 1184 | auto urlOptions = kj::Url::Options{.percentDecode = false, .allowEmpty = true}; |
| 1185 | auto maybeParsedUrl = kj::Url::tryParse(url.asPtr(), kj::Url::REMOTE_HREF, urlOptions); |
| 1186 | if (maybeParsedUrl == kj::none) { |
| 1187 | JSG_FAIL_REQUIRE(TypeError, kj::str("Unable to parse URL: ", url)); |
| 1188 | } |
| 1189 | parsedUrl = KJ_ASSERT_NONNULL(kj::mv(maybeParsedUrl)).toString(); |
| 1190 | } |
| 1191 | |
| 1192 | if (!kj::HttpHeaders::isValidHeaderValue(parsedUrl)) { |
| 1193 | JSG_FAIL_REQUIRE( |
| 1194 | TypeError, kj::str("Redirect URL cannot contain '\\r', '\\n', or '\\0': ", url)); |
| 1195 | } |
| 1196 | |
| 1197 | // Build our headers object with `Location` set to the parsed URL. |
| 1198 | kj::HttpHeaders kjHeaders(IoContext::current().getHeaderTable()); |
| 1199 | kjHeaders.set(kj::HttpHeaderId::LOCATION, kj::mv(parsedUrl)); |
| 1200 | auto headers = js.alloc<Headers>(js, kjHeaders, Headers::Guard::IMMUTABLE); |
| 1201 | |
| 1202 | return js.alloc<Response>(js, statusCode, kj::none, kj::mv(headers), nullptr, kj::none); |
| 1203 | } |
| 1204 | |
| 1205 | jsg::Ref<Response> Response::json_( |
| 1206 | jsg::Lock& js, jsg::JsValue any, jsg::Optional<Initializer> maybeInit) { |
| 1207 | |
| 1208 | const auto maybeSetContentType = [](jsg::Lock& js, auto headers) { |
| 1209 | if (!headers->hasCommon(capnp::CommonHeaderName::CONTENT_TYPE)) { |
| 1210 | headers->setCommon(capnp::CommonHeaderName::CONTENT_TYPE, MimeType::JSON.toString()); |
| 1211 | } |
| 1212 | return kj::mv(headers); |
| 1213 | }; |
| 1214 | |
| 1215 | // While this all looks a bit complicated, all the following is doing is checking |
| 1216 | // to see if maybeInit contains a content-type header. If it does, the existing |
| 1217 | // value is left alone. If it does not, then we set the value of content-type |
| 1218 | // to the default content type for JSON payloads. The reason this all looks a bit |
| 1219 | // complicated is that maybeInit is an optional kj::OneOf that might be either |
| 1220 | // a dict or a jsg::Ref<Response>. If it is a dict, then the optional headers |
| 1221 | // field is also an optional kj::OneOf that can be either a dict or a jsg::Ref<Headers>. |
| 1222 | // We have to deal with all of the various possibilities here to set the content-type |
| 1223 | // appropriately. |
| 1224 | KJ_IF_SOME(init, maybeInit) { |
| 1225 | KJ_SWITCH_ONEOF(init) { |
| 1226 | KJ_CASE_ONEOF(dict, InitializerDict) { |
| 1227 | KJ_IF_SOME(headers, dict.headers) { |
| 1228 | dict.headers = maybeSetContentType(js, Headers::constructor(js, kj::mv(headers))); |
| 1229 | } else { |
| 1230 | dict.headers = maybeSetContentType(js, js.alloc<Headers>()); |
| 1231 | } |
| 1232 | } |
| 1233 | KJ_CASE_ONEOF(res, jsg::Ref<Response>) { |
| 1234 | auto otherStatusText = res->getStatusText(); |
| 1235 | auto newInit = InitializerDict{ |
| 1236 | .status = res->statusCode, |
| 1237 | .statusText = otherStatusText == nullptr || |
| 1238 | otherStatusText == defaultStatusText(res->statusCode) |
| 1239 | ? jsg::Optional<kj::String>() : kj::str(otherStatusText), |
| 1240 | .headers = maybeSetContentType(js, Headers::constructor(js, res->headers.addRef())), |
| 1241 | .cf = res->cf.getRef(js), |
| 1242 | .encodeBody = |
| 1243 | kj::str(res->bodyEncoding == Response::BodyEncoding::MANUAL ? "manual" : "automatic"), |
| 1244 | }; |
| 1245 | |
| 1246 | KJ_IF_SOME(otherWs, res->webSocket) { |
| 1247 | newInit.webSocket = otherWs.addRef(); |
| 1248 | } |
| 1249 | |
| 1250 | maybeInit = kj::mv(newInit); |
| 1251 | } |
| 1252 | } |
| 1253 | } else { |
| 1254 | maybeInit = InitializerDict{ |
| 1255 | .headers = maybeSetContentType(js, js.alloc<Headers>()), |
| 1256 | }; |
| 1257 | } |
| 1258 | |
| 1259 | return constructor(js, kj::Maybe(any.toJson(js)), kj::mv(maybeInit)); |
| 1260 | } |
| 1261 | |
| 1262 | jsg::Ref<Response> Response::clone(jsg::Lock& js) { |
| 1263 | JSG_REQUIRE( |
| 1264 | webSocket == kj::none, TypeError, "Cannot clone a response to a WebSocket handshake."); |
| 1265 | |
| 1266 | auto headersClone = headers->clone(js); |
| 1267 | auto cfClone = cf.deepClone(js); |
| 1268 | |
| 1269 | auto bodyClone = Body::clone(js); |
| 1270 | |
| 1271 | auto urlListClone = KJ_MAP(url, urlList) { return kj::str(url); }; |
| 1272 | |
| 1273 | return js.alloc<Response>(js, statusCode, |
| 1274 | mapCopyString(statusText), |
| 1275 | kj::mv(headersClone), kj::mv(cfClone), kj::mv(bodyClone), kj::mv(urlListClone)); |
| 1276 | } |
| 1277 | |
| 1278 | kj::Promise<DeferredProxy<void>> Response::send(jsg::Lock& js, |
| 1279 | kj::HttpService::Response& outer, |
| 1280 | SendOptions options, |
| 1281 | kj::Maybe<const kj::HttpHeaders&> maybeReqHeaders) { |
| 1282 | JSG_REQUIRE(!getBodyUsed(), TypeError, |
| 1283 | "Body has already been used. " |
| 1284 | "It can only be used once. Use tee() first if you need to read it twice."); |
| 1285 | |
| 1286 | // Careful: Keep in mind that the promise we return could be canceled in which case `outer` will |
| 1287 | // be destroyed. Additionally, the response body stream we get from calling send() must itself |
| 1288 | // be destroyed before `outer` is destroyed. So, it's important to make sure that only the |
| 1289 | // promise we return encapsulates any task that might write to the response body. We can't, for |
| 1290 | // example, put the response body into a JS heap object. That should all be fine as long as we |
| 1291 | // use a pumpTo() that can be canceled. |
| 1292 | |
| 1293 | auto& context = IoContext::current(); |
| 1294 | kj::HttpHeaders outHeaders(context.getHeaderTable()); |
| 1295 | headers->shallowCopyTo(outHeaders); |
| 1296 | |
| 1297 | KJ_IF_SOME(ws, webSocket) { |
| 1298 | // `Response::acceptWebSocket()` can throw if we did not ask for a WebSocket. This |
| 1299 | // would promote a js client error into an uncatchable server error. Thus, we throw early here |
| 1300 | // if we do not expect a WebSocket. This could also be a 426 status code response, but we think |
| 1301 | // that the majority of our users expect us to throw on a client-side fetch error instead of |
| 1302 | // returning a 4xx status code. A 426 status code error _might_ be more appropriate if the |
| 1303 | // request headers originate from outside the worker developer's control (e.g. a client |
| 1304 | // application by some other party). |
| 1305 | JSG_REQUIRE(options.allowWebSocket, TypeError, |
| 1306 | "Worker tried to return a WebSocket in a response to a request " |
| 1307 | "which did not contain the header \"Upgrade: websocket\"."); |
| 1308 | |
| 1309 | const bool hasEnabledWebSocketCompression = FeatureFlags::get(js).getWebSocketCompression(); |
| 1310 | |
| 1311 | if (hasEnabledWebSocketCompression && |
| 1312 | outHeaders.get(kj::HttpHeaderId::SEC_WEBSOCKET_EXTENSIONS) == kj::none) { |
| 1313 | // Since workerd uses `MANUAL_COMPRESSION` mode for websocket compression, we need to |
| 1314 | // pass the headers we want to support to `acceptWebSocket()`. |
| 1315 | KJ_IF_SOME(config, ws->getPreferredExtensions(kj::WebSocket::ExtensionsContext::RESPONSE)) { |
| 1316 | // We try to get extensions for use in a response (i.e. for a server side websocket). |
| 1317 | // This allows us to `optimizedPumpTo()` `webSocket`. |
| 1318 | outHeaders.set(kj::HttpHeaderId::SEC_WEBSOCKET_EXTENSIONS, kj::mv(config)); |
| 1319 | } else { |
| 1320 | // `webSocket` is not a WebSocketImpl, we want to support whatever valid config the client |
| 1321 | // requested, so we'll just use the client's requested headers. |
| 1322 | KJ_IF_SOME(reqHeaders, maybeReqHeaders) { |
| 1323 | KJ_IF_SOME(value, reqHeaders.get(kj::HttpHeaderId::SEC_WEBSOCKET_EXTENSIONS)) { |
| 1324 | outHeaders.setPtr(kj::HttpHeaderId::SEC_WEBSOCKET_EXTENSIONS, value); |
| 1325 | } |
| 1326 | } |
| 1327 | } |
| 1328 | } else if (!hasEnabledWebSocketCompression) { |
| 1329 | // While we guard against an origin server including `Sec-WebSocket-Extensions` in a Response |
| 1330 | // (we don't send the extension in an offer, and if the server includes it in a response we |
| 1331 | // will reject the connection), a Worker could still explicitly add the header to a Response. |
| 1332 | outHeaders.unset(kj::HttpHeaderId::SEC_WEBSOCKET_EXTENSIONS); |
| 1333 | } |
| 1334 | |
| 1335 | auto clientSocket = outer.acceptWebSocket(outHeaders); |
| 1336 | auto wsPromise = ws->couple(kj::mv(clientSocket), context.getMetrics()); |
| 1337 | |
| 1338 | KJ_IF_SOME(a, context.getActor()) { |
| 1339 | KJ_IF_SOME(hib, a.getHibernationManager()) { |
| 1340 | // We attach a reference to the deferred proxy task so the HibernationManager lives at least |
| 1341 | // as long as the websocket connection. |
| 1342 | // The actor still retains its reference to the manager, so any subsequent requests prior |
| 1343 | // to hibernation will not need to re-obtain a reference. |
| 1344 | wsPromise = wsPromise.attach(kj::addRef(hib)); |
| 1345 | } |
| 1346 | } |
| 1347 | return wsPromise; |
| 1348 | } else KJ_IF_SOME(jsBody, getBody()) { |
| 1349 | auto encoding = getContentEncoding(context, outHeaders, bodyEncoding, FeatureFlags::get(js)); |
| 1350 | auto maybeLength = jsBody->tryGetLength(encoding); |
| 1351 | auto stream = |
| 1352 | newSystemStream(outer.send(statusCode, getStatusText(), outHeaders, maybeLength), encoding); |
| 1353 | // We need to enter the AsyncContextFrame that was captured when the |
| 1354 | // Response was created before starting the loop. |
| 1355 | jsg::AsyncContextFrame::Scope scope(js, asyncContext); |
| 1356 | return jsBody->pumpTo(js, kj::mv(stream), true); |
| 1357 | } else { |
| 1358 | outer.send(statusCode, getStatusText(), outHeaders, static_cast<uint64_t>(0)); |
| 1359 | return addNoopDeferredProxy(kj::READY_NOW); |
| 1360 | } |
| 1361 | } |
| 1362 | |
| 1363 | int Response::getStatus() { |
| 1364 | return statusCode; |
| 1365 | } |
| 1366 | kj::StringPtr Response::getStatusText() { |
| 1367 | KJ_IF_SOME(text, statusText) { |
| 1368 | return text; |
| 1369 | } |
| 1370 | return defaultStatusText(statusCode); |
| 1371 | } |
| 1372 | jsg::Ref<Headers> Response::getHeaders(jsg::Lock& js) { |
| 1373 | return headers.addRef(); |
| 1374 | } |
| 1375 | |
| 1376 | bool Response::getOk() { |
| 1377 | return statusCode >= 200 && statusCode < 300; |
| 1378 | } |
| 1379 | bool Response::getRedirected() { |
| 1380 | return urlList.size() > 1; |
| 1381 | } |
| 1382 | kj::StringPtr Response::getUrl() { |
| 1383 | if (urlList.size() > 0) { |
| 1384 | // We're supposed to drop any fragment from the URL. Instead of doing it here, we rely on the |
| 1385 | // code that calls the Response constructor (e.g. makeHttpResponse()) to drop the fragments |
| 1386 | // before giving the stringified URL to us. |
| 1387 | return urlList.back(); |
| 1388 | } else { |
| 1389 | // Per spec, if the URL list is empty, we return an empty string. I dunno, man. |
| 1390 | return ""; |
| 1391 | } |
| 1392 | } |
| 1393 | |
| 1394 | kj::Maybe<jsg::Ref<WebSocket>> Response::getWebSocket(jsg::Lock& js) { |
| 1395 | return webSocket.map([&](jsg::Ref<WebSocket>& ptr) { return ptr.addRef(); }); |
| 1396 | } |
| 1397 | |
| 1398 | jsg::Optional<jsg::JsObject> Response::getCf(jsg::Lock& js) { |
| 1399 | return cf.get(js); |
| 1400 | } |
| 1401 | |
| 1402 | void Response::serialize(jsg::Lock& js, |
| 1403 | jsg::Serializer& serializer, |
| 1404 | const jsg::TypeHandler<InitializerDict>& initDictHandler, |
| 1405 | const jsg::TypeHandler<kj::Maybe<jsg::Ref<ReadableStream>>>& streamHandler) { |
| 1406 | serializer.write(js, jsg::JsValue(streamHandler.wrap(js, getBody()))); |
| 1407 | |
| 1408 | // As with Request, we serialize the initializer dict as a JS object. |
| 1409 | serializer.write(js, |
| 1410 | jsg::JsValue(initDictHandler.wrap(js, |
| 1411 | InitializerDict{ |
| 1412 | .status = statusCode == 200 ? jsg::Optional<int>() : statusCode, |
| 1413 | .statusText = statusText.map([](auto& txt) { return kj::str(txt); }), |
| 1414 | .headers = headers.addRef(), |
| 1415 | .cf = cf.getRef(js), |
| 1416 | |
| 1417 | // If a WebSocket is present, we'll try to serialize it. As of this writing, WebSocket |
| 1418 | // is not serializable, but we could add support for sending it over RPC in the future. |
| 1419 | // |
| 1420 | // Note we have to double-Maybe this, so that if no signal is present, the property is absent |
| 1421 | // instead of `null`. |
| 1422 | .webSocket = |
| 1423 | webSocket.map([](jsg::Ref<WebSocket>& s) -> kj::Maybe<jsg::Ref<WebSocket>> { |
| 1424 | return s.addRef(); |
| 1425 | }), |
| 1426 | |
| 1427 | .encodeBody = bodyEncoding == BodyEncoding::AUTO ? jsg::Optional<kj::String>() |
| 1428 | : kj::str("manual"), |
| 1429 | }))); |
| 1430 | } |
| 1431 | |
| 1432 | jsg::Ref<Response> Response::deserialize(jsg::Lock& js, |
| 1433 | rpc::SerializationTag tag, |
| 1434 | jsg::Deserializer& deserializer, |
| 1435 | const jsg::TypeHandler<InitializerDict>& initDictHandler, |
| 1436 | const jsg::TypeHandler<kj::Maybe<jsg::Ref<ReadableStream>>>& streamHandler) { |
| 1437 | auto body = KJ_UNWRAP_OR(streamHandler.tryUnwrap(js, deserializer.readValue(js)), { |
| 1438 | JSG_FAIL_REQUIRE(DOMDataCloneError, |
| 1439 | "Deserialization failed: could not deserialize Response body"); |
| 1440 | }); |
| 1441 | auto init = KJ_UNWRAP_OR(initDictHandler.tryUnwrap(js, deserializer.readValue(js)), { |
| 1442 | JSG_FAIL_REQUIRE(DOMDataCloneError, |
| 1443 | "Deserialization failed: could not deserialize Response initializer"); |
| 1444 | }); |
| 1445 | |
| 1446 | // If the status code is zero, then it was an error response. We cannot |
| 1447 | // use the Response::constructor. |
| 1448 | KJ_IF_SOME(status, init.status) { |
| 1449 | if (status == 0) { |
| 1450 | return Response::error(js); |
| 1451 | } |
| 1452 | } |
| 1453 | |
| 1454 | return Response::constructor(js, kj::mv(body), kj::mv(init)); |
| 1455 | } |
| 1456 | |
| 1457 | // ======================================================================================= |
| 1458 | |
| 1459 | jsg::Ref<Request> FetchEvent::getRequest() { |
| 1460 | return request.addRef(); |
| 1461 | } |
| 1462 | |
| 1463 | kj::Maybe<jsg::Promise<jsg::Ref<Response>>> FetchEvent::getResponsePromise(jsg::Lock& js) { |
| 1464 | KJ_SWITCH_ONEOF(state) { |
| 1465 | KJ_CASE_ONEOF(_, AwaitingRespondWith) { |
| 1466 | state = ResponseSent(); |
| 1467 | return kj::none; |
| 1468 | } |
| 1469 | KJ_CASE_ONEOF(called, RespondWithCalled) { |
| 1470 | auto result = kj::mv(called.promise); |
| 1471 | state = ResponseSent(); |
| 1472 | return kj::mv(result); |
| 1473 | } |
| 1474 | KJ_CASE_ONEOF(_, ResponseSent) { |
| 1475 | KJ_FAIL_REQUIRE("can only call getResponsePromise() once"); |
| 1476 | } |
| 1477 | } |
| 1478 | KJ_UNREACHABLE; |
| 1479 | } |
| 1480 | |
| 1481 | void FetchEvent::respondWith(jsg::Lock& js, jsg::Promise<jsg::Ref<Response>> promise) { |
| 1482 | preventDefault(); |
| 1483 | |
| 1484 | if (IoContext::current().hasOutputGate()) { |
| 1485 | // Once a Response is returned, we need to apply the output lock. |
| 1486 | promise = promise.then(js, [](jsg::Lock& js, jsg::Ref<Response>&& response) { |
| 1487 | auto& context = IoContext::current(); |
| 1488 | return context.awaitIo(js, context.waitForOutputLocks(), |
| 1489 | [response = kj::mv(response)](jsg::Lock&) mutable { return kj::mv(response); }); |
| 1490 | }); |
| 1491 | } |
| 1492 | |
| 1493 | KJ_SWITCH_ONEOF(state) { |
| 1494 | KJ_CASE_ONEOF(_, AwaitingRespondWith) { |
| 1495 | state = RespondWithCalled{kj::mv(promise)}; |
| 1496 | } |
| 1497 | KJ_CASE_ONEOF(called, RespondWithCalled) { |
| 1498 | JSG_FAIL_REQUIRE(DOMInvalidStateError, |
| 1499 | "FetchEvent.respondWith() has already been called; it can only be called once."); |
| 1500 | } |
| 1501 | KJ_CASE_ONEOF(_, ResponseSent) { |
| 1502 | JSG_FAIL_REQUIRE(DOMInvalidStateError, |
| 1503 | "Too late to call FetchEvent.respondWith(). It must be called synchronously in the " |
| 1504 | "event handler."); |
| 1505 | } |
| 1506 | } |
| 1507 | |
| 1508 | stopImmediatePropagation(); |
| 1509 | } |
| 1510 | |
| 1511 | void FetchEvent::passThroughOnException() { |
| 1512 | IoContext::current().setFailOpen(); |
| 1513 | } |
| 1514 | |
| 1515 | // ======================================================================================= |
| 1516 | |
| 1517 | namespace { |
| 1518 | |
| 1519 | // Fetch spec requires (suggests?) 20: https://fetch.spec.whatwg.org/#http-redirect-fetch |
| 1520 | constexpr auto MAX_REDIRECT_COUNT = 20; |
| 1521 | |
| 1522 | jsg::Promise<jsg::Ref<Response>> handleHttpResponse(jsg::Lock& js, |
| 1523 | jsg::Ref<Fetcher> fetcher, |
| 1524 | jsg::Ref<Request> jsRequest, |
| 1525 | kj::Vector<kj::Url> urlList, |
| 1526 | kj::HttpClient::Response&& response); |
| 1527 | jsg::Promise<jsg::Ref<Response>> handleHttpRedirectResponse(jsg::Lock& js, |
| 1528 | jsg::Ref<Fetcher> fetcher, |
| 1529 | jsg::Ref<Request> jsRequest, |
| 1530 | kj::Vector<kj::Url> urlList, |
| 1531 | uint status, |
| 1532 | kj::StringPtr location); |
| 1533 | |
| 1534 | jsg::Promise<jsg::Ref<Response>> fetchImplNoOutputLock(jsg::Lock& js, |
| 1535 | jsg::Ref<Fetcher> fetcher, |
| 1536 | jsg::Ref<Request> jsRequest, |
| 1537 | kj::Vector<kj::Url> urlList) { |
| 1538 | KJ_ASSERT(!urlList.empty()); |
| 1539 | |
| 1540 | auto& ioContext = IoContext::current(); |
| 1541 | |
| 1542 | auto signal = jsRequest->getSignal(); |
| 1543 | KJ_IF_SOME(s, signal) { |
| 1544 | // If the AbortSignal has already been triggered, then we need to stop here. |
| 1545 | if (s->getAborted(js)) { |
| 1546 | return js.rejectedPromise<jsg::Ref<Response>>(s->getReason(js)); |
| 1547 | } |
| 1548 | } |
| 1549 | |
| 1550 | // Get client and trace context (if needed) in one clean call |
| 1551 | auto clientWithTracing = fetcher->getClientWithTracing(ioContext, jsRequest->serializeCfBlobJson(js), "fetch"_kjc); |
| 1552 | auto traceContext = kj::mv(clientWithTracing.traceContext); |
| 1553 | |
| 1554 | // TODO(cleanup): Don't convert to HttpClient. Use the HttpService interface instead. This |
| 1555 | // requires a significant rewrite of the code below. It'll probably get simpler, though? |
| 1556 | kj::Own<kj::HttpClient> client = asHttpClient(kj::mv(clientWithTracing.client)); |
| 1557 | |
| 1558 | kj::HttpHeaders headers(ioContext.getHeaderTable()); |
| 1559 | jsRequest->shallowCopyHeadersTo(headers); |
| 1560 | |
| 1561 | // If the jsRequest has a CacheMode, we need to handle that here. |
| 1562 | // Currently, the only cache mode we support is undefined and no-store, no-cache, and reload |
| 1563 | auto headerIds = ioContext.getHeaderIds(); |
| 1564 | const auto cacheMode = jsRequest->getCacheMode(); |
| 1565 | switch (cacheMode) { |
| 1566 | case Request::CacheMode::RELOAD: |
| 1567 | KJ_FALLTHROUGH; |
| 1568 | case Request::CacheMode::NOSTORE: |
| 1569 | KJ_FALLTHROUGH; |
| 1570 | case Request::CacheMode::NOCACHE: |
| 1571 | if (headers.get(headerIds.cacheControl) == kj::none) { |
| 1572 | headers.setPtr(headerIds.cacheControl, "no-cache"); |
| 1573 | } |
| 1574 | if (headers.get(headerIds.pragma) == kj::none) { |
| 1575 | headers.setPtr(headerIds.pragma, "no-cache"); |
| 1576 | } |
| 1577 | KJ_FALLTHROUGH; |
| 1578 | case Request::CacheMode::NONE: |
| 1579 | break; |
| 1580 | default: |
| 1581 | KJ_UNREACHABLE; |
| 1582 | } |
| 1583 | |
| 1584 | KJ_IF_SOME(ctx, traceContext) { |
| 1585 | ctx.setTag("network.protocol.name"_kjc, "http"_kjc); |
| 1586 | ctx.setTag("network.protocol.version"_kjc, "HTTP/1.1"_kjc); |
| 1587 | ctx.setTag("http.request.method"_kjc, kj::str(jsRequest->getMethodEnum())); |
| 1588 | ctx.setTag("url.full"_kjc, jsRequest->getUrl()); |
| 1589 | |
| 1590 | KJ_IF_SOME(userAgent, headers.get(headerIds.userAgent)) { |
| 1591 | ctx.setTag("user_agent.original"_kjc, userAgent); |
| 1592 | } |
| 1593 | |
| 1594 | KJ_IF_SOME(contentType, headers.get(headerIds.contentType)) { |
| 1595 | ctx.setTag("http.request.header.content-type"_kjc, contentType); |
| 1596 | } |
| 1597 | |
| 1598 | KJ_IF_SOME(contentLength, headers.get(headerIds.contentLength)) { |
| 1599 | ctx.setTag("http.request.header.content-length"_kjc, contentLength); |
| 1600 | } |
| 1601 | |
| 1602 | KJ_IF_SOME(accept, headers.get(headerIds.accept)) { |
| 1603 | ctx.setTag("http.request.header.accept"_kjc, accept); |
| 1604 | } |
| 1605 | |
| 1606 | KJ_IF_SOME(acceptEncoding, headers.get(headerIds.acceptEncoding)) { |
| 1607 | ctx.setTag("http.request.header.accept-encoding"_kjc, acceptEncoding); |
| 1608 | } |
| 1609 | } |
| 1610 | |
| 1611 | kj::String url = |
| 1612 | uriEncodeControlChars(urlList.back().toString(kj::Url::HTTP_PROXY_REQUEST).asBytes()); |
| 1613 | |
| 1614 | if (headers.isWebSocket()) { |
| 1615 | if (!FeatureFlags::get(js).getWebSocketCompression()) { |
| 1616 | // If we haven't enabled the websocket compression compatibility flag, strip the header from the |
| 1617 | // subrequest. |
| 1618 | headers.unset(kj::HttpHeaderId::SEC_WEBSOCKET_EXTENSIONS); |
| 1619 | } |
| 1620 | auto webSocketResponse = client->openWebSocket(url, headers); |
| 1621 | return ioContext.awaitIo(js, |
| 1622 | AbortSignal::maybeCancelWrap(js, signal, kj::mv(webSocketResponse)), |
| 1623 | [fetcher = kj::mv(fetcher), jsRequest = kj::mv(jsRequest), urlList = kj::mv(urlList), |
| 1624 | client = kj::mv(client), signal = kj::mv(signal)]( |
| 1625 | jsg::Lock& js, kj::HttpClient::WebSocketResponse&& response) mutable |
| 1626 | -> jsg::Promise<jsg::Ref<Response>> { |
| 1627 | KJ_SWITCH_ONEOF(response.webSocketOrBody) { |
| 1628 | KJ_CASE_ONEOF(body, kj::Own<kj::AsyncInputStream>) { |
| 1629 | body = body.attach(kj::mv(client)); |
| 1630 | return handleHttpResponse(js, kj::mv(fetcher), kj::mv(jsRequest), kj::mv(urlList), |
| 1631 | {response.statusCode, response.statusText, response.headers, kj::mv(body)}); |
| 1632 | } |
| 1633 | KJ_CASE_ONEOF(webSocket, kj::Own<kj::WebSocket>) { |
| 1634 | KJ_ASSERT(response.statusCode == 101); |
| 1635 | webSocket = webSocket.attach(kj::mv(client)); |
| 1636 | KJ_IF_SOME(s, signal) { |
| 1637 | // If the AbortSignal has already been triggered, then we need to stop here. |
| 1638 | if (s->getAborted(js)) { |
| 1639 | return js.rejectedPromise<jsg::Ref<Response>>(s->getReason(js)); |
| 1640 | } |
| 1641 | webSocket = kj::refcounted<AbortableWebSocket>(kj::mv(webSocket), s->getCanceler()); |
| 1642 | } |
| 1643 | return js.resolvedPromise(makeHttpResponse(js, jsRequest->getMethodEnum(), |
| 1644 | kj::mv(urlList), response.statusCode, response.statusText, *response.headers, |
| 1645 | newNullInputStream(), js.alloc<WebSocket>(js, kj::mv(webSocket)), |
| 1646 | jsRequest->getResponseBodyEncoding(), kj::mv(signal))); |
| 1647 | } |
| 1648 | } |
| 1649 | KJ_UNREACHABLE; |
| 1650 | }); |
| 1651 | } else { |
| 1652 | kj::Maybe<kj::HttpClient::Request> nativeRequest; |
| 1653 | KJ_IF_SOME(jsBody, jsRequest->getBody()) { |
| 1654 | // Note that for requests, we do not automatically handle Content-Encoding, because the fetch() |
| 1655 | // standard does not say that we should. Hence, we always use StreamEncoding::IDENTITY. |
| 1656 | // https://github.com/whatwg/fetch/issues/589 |
| 1657 | auto maybeLength = jsBody->tryGetLength(StreamEncoding::IDENTITY); |
| 1658 | KJ_IF_SOME(ctx, traceContext) { |
| 1659 | KJ_IF_SOME(length, maybeLength) { |
| 1660 | ctx.setTag("http.request.body.size"_kjc, static_cast<int64_t>(length)); |
| 1661 | } |
| 1662 | } |
| 1663 | |
| 1664 | if (maybeLength.orDefault(1) == 0 && |
| 1665 | headers.get(kj::HttpHeaderId::CONTENT_LENGTH) == kj::none && |
| 1666 | headers.get(kj::HttpHeaderId::TRANSFER_ENCODING) == kj::none) { |
| 1667 | // Request has a non-null but explicitly empty body, and has neither a Content-Length nor |
| 1668 | // a Transfer-Encoding header. If we don't set one of those two, and the receiving end is |
| 1669 | // another worker (especially within a pipeline or reached via RPC, not real HTTP), then |
| 1670 | // the code in global-scope.c++ on the receiving end will decide the body should be null. |
| 1671 | // We'd like to avoid this weird discontinuity, so let's set Content-Length explicitly to |
| 1672 | // 0. |
| 1673 | headers.setPtr(kj::HttpHeaderId::CONTENT_LENGTH, "0"_kj); |
| 1674 | } |
| 1675 | |
| 1676 | nativeRequest = client->request(jsRequest->getMethodEnum(), url, headers, maybeLength); |
| 1677 | auto& nr = KJ_ASSERT_NONNULL(nativeRequest); |
| 1678 | auto stream = newSystemStream(kj::mv(nr.body), StreamEncoding::IDENTITY); |
| 1679 | |
| 1680 | // We want to support bidirectional streaming, so we actually don't want to wait for the |
| 1681 | // request to finish before we deliver the response to the app. |
| 1682 | |
| 1683 | // jsBody is not used directly within the function but is passed in so that |
| 1684 | // the coroutine frame keeps it alive. |
| 1685 | static constexpr auto handleCancelablePump = [](kj::Promise<void> promise, |
| 1686 | auto jsBody) -> kj::Promise<void> { |
| 1687 | try { |
| 1688 | co_await promise; |
| 1689 | } catch (...) { |
| 1690 | auto exception = kj::getCaughtExceptionAsKj(); |
| 1691 | if (exception.getType() != kj::Exception::Type::DISCONNECTED) { |
| 1692 | kj::throwFatalException(kj::mv(exception)); |
| 1693 | } |
| 1694 | // Ignore DISCONNECTED exceptions thrown by the writePromise, so that we always |
| 1695 | // return the server's response, which should identify if any issue occurred with |
| 1696 | // the body stream anyway. |
| 1697 | } |
| 1698 | }; |
| 1699 | |
| 1700 | // TODO(someday): Allow deferred proxying for bidirectional streaming. |
| 1701 | ioContext.addWaitUntil(handleCancelablePump( |
| 1702 | AbortSignal::maybeCancelWrap( |
| 1703 | js, signal, ioContext.waitForDeferredProxy(jsBody->pumpTo(js, kj::mv(stream), true))), |
| 1704 | jsBody.addRef())); |
| 1705 | } else { |
| 1706 | nativeRequest = client->request(jsRequest->getMethodEnum(), url, headers, static_cast<uint64_t>(0)); |
| 1707 | } |
| 1708 | return ioContext.awaitIo(js, |
| 1709 | AbortSignal::maybeCancelWrap(js, signal, kj::mv(KJ_ASSERT_NONNULL(nativeRequest).response)) |
| 1710 | .catch_([](kj::Exception&& exception) -> kj::Promise<kj::HttpClient::Response> { |
| 1711 | if (exception.getDescription().startsWith("invalid Content-Length header value")) { |
| 1712 | return JSG_KJ_EXCEPTION(FAILED, Error, exception.getDescription()); |
| 1713 | } else if (exception.getDescription().contains("NOSENTRY script not found")) { |
| 1714 | return JSG_KJ_EXCEPTION(FAILED, Error, "Worker not found."); |
| 1715 | } |
| 1716 | return kj::mv(exception); |
| 1717 | }), |
| 1718 | [fetcher = kj::mv(fetcher), jsRequest = kj::mv(jsRequest), urlList = kj::mv(urlList), |
| 1719 | client = kj::mv(client), traceContext = kj::mv(traceContext)](jsg::Lock& js, |
| 1720 | kj::HttpClient::Response&& response) mutable -> jsg::Promise<jsg::Ref<Response>> { |
| 1721 | response.body = response.body.attach(kj::mv(client)); |
| 1722 | KJ_IF_SOME(ctx, traceContext) { |
| 1723 | ctx.setTag("http.response.status_code"_kjc, static_cast<int64_t>(response.statusCode)); |
| 1724 | KJ_IF_SOME(length, response.body->tryGetLength()) { |
| 1725 | ctx.setTag("http.response.body.size"_kjc, static_cast<int64_t>(length)); |
| 1726 | } |
| 1727 | auto headerIds = IoContext::current().getHeaderIds(); |
| 1728 | KJ_IF_SOME(cfRay, response.headers->get(headerIds.cfRay)) { |
| 1729 | ctx.setTag("cloudflare.ray_id"_kjc, cfRay); |
| 1730 | } |
| 1731 | } |
| 1732 | return handleHttpResponse( |
| 1733 | js, kj::mv(fetcher), kj::mv(jsRequest), kj::mv(urlList), kj::mv(response)); |
| 1734 | }); |
| 1735 | } |
| 1736 | } |
| 1737 | |
| 1738 | jsg::Promise<jsg::Ref<Response>> fetchImpl(jsg::Lock& js, |
| 1739 | jsg::Ref<Fetcher> fetcher, |
| 1740 | jsg::Ref<Request> jsRequest, |
| 1741 | kj::Vector<kj::Url> urlList) { |
| 1742 | auto& context = IoContext::current(); |
| 1743 | // Optimization: For non-actors, which never have output locks, avoid the overhead of |
| 1744 | // awaitIo() and such by not going back to the event loop at all. |
| 1745 | KJ_IF_SOME(promise, context.waitForOutputLocksIfNecessary()) { |
| 1746 | return context.awaitIo(js, kj::mv(promise), |
| 1747 | [fetcher = kj::mv(fetcher), jsRequest = kj::mv(jsRequest), urlList = kj::mv(urlList)]( |
| 1748 | jsg::Lock& js) mutable { |
| 1749 | return fetchImplNoOutputLock(js, kj::mv(fetcher), kj::mv(jsRequest), kj::mv(urlList)); |
| 1750 | }); |
| 1751 | } else { |
| 1752 | return fetchImplNoOutputLock(js, kj::mv(fetcher), kj::mv(jsRequest), kj::mv(urlList)); |
| 1753 | } |
| 1754 | } |
| 1755 | |
| 1756 | jsg::Promise<jsg::Ref<Response>> handleHttpResponse(jsg::Lock& js, |
| 1757 | jsg::Ref<Fetcher> fetcher, |
| 1758 | jsg::Ref<Request> jsRequest, |
| 1759 | kj::Vector<kj::Url> urlList, |
| 1760 | kj::HttpClient::Response&& response) { |
| 1761 | auto signal = jsRequest->getSignal(); |
| 1762 | |
| 1763 | KJ_IF_SOME(s, signal) { |
| 1764 | // If the AbortSignal has already been triggered, then we need to stop here. |
| 1765 | if (s->getAborted(js)) { |
| 1766 | return js.rejectedPromise<jsg::Ref<Response>>(s->getReason(js)); |
| 1767 | } |
| 1768 | response.body = kj::refcounted<AbortableInputStream>(kj::mv(response.body), s->getCanceler()); |
| 1769 | } |
| 1770 | |
| 1771 | if (isRedirectStatusCode(response.statusCode) && |
| 1772 | jsRequest->getRedirectEnum() == Request::Redirect::FOLLOW) { |
| 1773 | KJ_IF_SOME(l, response.headers->get(kj::HttpHeaderId::LOCATION)) { |
| 1774 | |
| 1775 | // Pump the response body to a singleton null stream before following the redirect. |
| 1776 | auto& ioContext = IoContext::current(); |
| 1777 | return ioContext.awaitIo(js, |
| 1778 | response.body->pumpTo(getGlobalNullOutputStream()).ignoreResult().attach(kj::mv(response.body)), |
| 1779 | [fetcher = kj::mv(fetcher), jsRequest = kj::mv(jsRequest), urlList = kj::mv(urlList), |
| 1780 | status = response.statusCode, location = kj::str(l)](jsg::Lock& js) mutable { |
| 1781 | return handleHttpRedirectResponse( |
| 1782 | js, kj::mv(fetcher), kj::mv(jsRequest), kj::mv(urlList), status, kj::mv(location)); |
| 1783 | }); |
| 1784 | } else { |
| 1785 | // No Location header. That's OK, we just return the response as is. |
| 1786 | // See https://fetch.spec.whatwg.org/#http-redirect-fetch step 2. |
| 1787 | } |
| 1788 | } |
| 1789 | |
| 1790 | auto result = makeHttpResponse(js, jsRequest->getMethodEnum(), kj::mv(urlList), |
| 1791 | response.statusCode, response.statusText, *response.headers, kj::mv(response.body), kj::none, |
| 1792 | jsRequest->getResponseBodyEncoding(), kj::mv(signal)); |
| 1793 | |
| 1794 | return js.resolvedPromise(kj::mv(result)); |
| 1795 | } |
| 1796 | |
| 1797 | jsg::Promise<jsg::Ref<Response>> handleHttpRedirectResponse(jsg::Lock& js, |
| 1798 | jsg::Ref<Fetcher> fetcher, |
| 1799 | jsg::Ref<Request> jsRequest, |
| 1800 | kj::Vector<kj::Url> urlList, |
| 1801 | uint status, |
| 1802 | kj::StringPtr location) { |
| 1803 | // Reconstruct the request body stream for retransmission in the face of a redirect. Before |
| 1804 | // reconstructing the stream, however, this function: |
| 1805 | // |
| 1806 | // - Throws if `status` is non-303 and this request doesn't have a "rewindable" body. |
| 1807 | // - Translates POST requests that hit 301, 302, or 303 into GET requests with null bodies. |
| 1808 | // - Translates HEAD requests that hit 303 into HEAD requests with null bodies. |
| 1809 | // - Translates all other requests that hit 303 into GET requests with null bodies. |
| 1810 | |
| 1811 | auto redirectedLocation = ([&]() -> kj::Maybe<kj::Url> { |
| 1812 | // TODO(later): This is a bit unfortunate. Per the fetch spec, we're supposed to be |
| 1813 | // using standard WHATWG URL parsing to resolve the redirect URL. However, changing it |
| 1814 | // now requires a compat flag. In order to minimize changes to the rest of the impl |
| 1815 | // we end up double parsing the URL here, once with the standard parser to produce the |
| 1816 | // correct result, and again with kj::Url in order to produce something that works with |
| 1817 | // the existing code. Fortunately the standard parser is fast but it would be nice to |
| 1818 | // be able to avoid the double parse at some point. |
| 1819 | if (FeatureFlags::get(js).getFetchStandardUrl()) { |
| 1820 | auto base = urlList.back().toString(); |
| 1821 | KJ_IF_SOME(parsed, jsg::Url::tryParse(location, base.asPtr())) { |
| 1822 | auto str = kj::str(parsed.getHref()); |
| 1823 | return kj::Url::tryParse(str.asPtr(), kj::Url::Context::REMOTE_HREF, |
| 1824 | kj::Url::Options{ |
| 1825 | .percentDecode = false, |
| 1826 | .allowEmpty = true, |
| 1827 | }); |
| 1828 | } else { |
| 1829 | return kj::none; |
| 1830 | } |
| 1831 | } else { |
| 1832 | return urlList.back().tryParseRelative(location); |
| 1833 | } |
| 1834 | })(); |
| 1835 | |
| 1836 | if (redirectedLocation == kj::none) { |
| 1837 | auto exception = |
| 1838 | JSG_KJ_EXCEPTION(FAILED, TypeError, "Invalid Location header; unable to follow redirect."); |
| 1839 | return js.rejectedPromise<jsg::Ref<Response>>(kj::mv(exception)); |
| 1840 | } |
| 1841 | |
| 1842 | // Note: RFC7231 says we should propagate fragments from the current request URL to the |
| 1843 | // redirected URL. The Fetch spec seems to take the position that that's the navigator's |
| 1844 | // job -- i.e., that you should be using redirect manual mode and deciding what to do with |
| 1845 | // fragments in Location headers yourself. We follow the spec, and don't do any explicit |
| 1846 | // fragment propagation. |
| 1847 | |
| 1848 | if (urlList.size() - 1 >= MAX_REDIRECT_COUNT) { |
| 1849 | auto exception = JSG_KJ_EXCEPTION(FAILED, TypeError, "Too many redirects.", urlList); |
| 1850 | return js.rejectedPromise<jsg::Ref<Response>>(kj::mv(exception)); |
| 1851 | } |
| 1852 | |
| 1853 | if (FeatureFlags::get(js).getStripAuthorizationOnCrossOriginRedirect()) { |
| 1854 | auto base = urlList.back().toString(); |
| 1855 | |
| 1856 | auto currentUrl = KJ_UNWRAP_OR(jsg::Url::tryParse(base.asPtr()), { |
| 1857 | auto exception = |
| 1858 | JSG_KJ_EXCEPTION(FAILED, TypeError, "Invalid current URL; unable to follow redirect."); |
| 1859 | return js.rejectedPromise<jsg::Ref<Response>>(kj::mv(exception)); |
| 1860 | }); |
| 1861 | |
| 1862 | auto locationUrl = KJ_UNWRAP_OR(jsg::Url::tryParse(location, base.asPtr()), { |
| 1863 | auto exception = |
| 1864 | JSG_KJ_EXCEPTION(FAILED, TypeError, "Invalid Location header; unable to follow redirect."); |
| 1865 | return js.rejectedPromise<jsg::Ref<Response>>(kj::mv(exception)); |
| 1866 | }); |
| 1867 | |
| 1868 | if (currentUrl.getOrigin() != locationUrl.getOrigin()) { |
| 1869 | // If requestโs current URLโs origin is not same origin with locationURLโs origin, then |
| 1870 | // for each headerName of CORS non-wildcard request-header name, delete headerName from |
| 1871 | // requestโs header list. |
| 1872 | // -- Fetch spec s. 4.4.13 |
| 1873 | // <https://fetch.spec.whatwg.org/#http-redirect-fetch> |
| 1874 | // (NB: "CORS non-wildcard request-header name" consists solely of "Authorization") |
| 1875 | |
| 1876 | jsRequest->getHeaders(js)->deleteCommon(capnp::CommonHeaderName::AUTHORIZATION); |
| 1877 | } |
| 1878 | } |
| 1879 | |
| 1880 | urlList.add(kj::mv(KJ_ASSERT_NONNULL(redirectedLocation))); |
| 1881 | |
| 1882 | // "If actualResponseโs status is not 303, requestโs body is non-null, and |
| 1883 | // requestโs bodyโs source [buffer] is null, then return a network error." |
| 1884 | // https://fetch.spec.whatwg.org/#http-redirect-fetch step 9. |
| 1885 | // |
| 1886 | // TODO(conform): this check pedantically enforces the spec, even if a POST hits a 301 or |
| 1887 | // 302. In that case, we're going to null out the body, anyway, so it feels strange to |
| 1888 | // report an error. If we widen fetch()'s contract to allow POSTs with non-buffer-backed |
| 1889 | // bodies to survive 301/302 redirects, our logic would get simpler here. |
| 1890 | // |
| 1891 | // Follow up with the spec authors about this. |
| 1892 | if (status != 303 && !jsRequest->canRewindBody()) { |
| 1893 | auto exception = JSG_KJ_EXCEPTION(FAILED, TypeError, |
| 1894 | "A request with a one-time-use body (it was initialized from a stream, not a buffer) " |
| 1895 | "encountered a redirect requiring the body to be retransmitted. To avoid this error " |
| 1896 | "in the future, construct this request from a buffer-like body initializer."); |
| 1897 | return js.rejectedPromise<jsg::Ref<Response>>(kj::mv(exception)); |
| 1898 | } |
| 1899 | |
| 1900 | auto method = jsRequest->getMethodEnum(); |
| 1901 | |
| 1902 | // "If either actualResponseโs status is 301 or 302 and requestโs method is |
| 1903 | // `POST`, or actualResponseโs status is 303 and request's method is not `HEAD`, set requestโs |
| 1904 | // method to `GET` and requestโs body to null." |
| 1905 | // https://fetch.spec.whatwg.org/#http-redirect-fetch step 11. |
| 1906 | if (((status == 301 || status == 302) && method == kj::HttpMethod::POST) || |
| 1907 | (status == 303 && method != kj::HttpMethod::HEAD)) { |
| 1908 | // TODO(conform): When translating a request with a body to a GET request, should we |
| 1909 | // explicitly remove Content-* headers? See https://github.com/whatwg/fetch/issues/609 |
| 1910 | jsRequest->setMethodEnum(kj::HttpMethod::GET); |
| 1911 | jsRequest->nullifyBody(); |
| 1912 | } else { |
| 1913 | // Reconstruct the stream from our buffer. The spec does not specify that we should cancel the |
| 1914 | // current body transmission in HTTP/1.1, so I'm not neutering the stream. (For HTTP/2 it asks |
| 1915 | // us to send a RST_STREAM frame if possible.) |
| 1916 | // |
| 1917 | // We know `buffer` is non-null here because we checked `buffer`'s nullness when non-303, and |
| 1918 | // nulled out `impl` when 303. Combined, they guarantee that we have a backing buffer. |
| 1919 | jsRequest->rewindBody(js); |
| 1920 | } |
| 1921 | |
| 1922 | // No need to wait for output locks again when following a redirect, because we didn't interact |
| 1923 | // with the app state in any way. |
| 1924 | return fetchImplNoOutputLock(js, kj::mv(fetcher), kj::mv(jsRequest), kj::mv(urlList)); |
| 1925 | } |
| 1926 | |
| 1927 | } // namespace |
| 1928 | |
| 1929 | jsg::Ref<Response> makeHttpResponse(jsg::Lock& js, |
| 1930 | kj::HttpMethod method, |
| 1931 | kj::Vector<kj::Url> urlListParam, |
| 1932 | uint statusCode, |
| 1933 | kj::StringPtr statusText, |
| 1934 | const kj::HttpHeaders& headers, |
| 1935 | kj::Own<kj::AsyncInputStream> body, |
| 1936 | kj::Maybe<jsg::Ref<WebSocket>> webSocket, |
| 1937 | Response::BodyEncoding bodyEncoding, |
| 1938 | kj::Maybe<jsg::Ref<AbortSignal>> signal) { |
| 1939 | auto responseHeaders = js.alloc<Headers>(js, headers, Headers::Guard::RESPONSE); |
| 1940 | auto& context = IoContext::current(); |
| 1941 | |
| 1942 | // The Fetch spec defines responses to HEAD or CONNECT requests, or responses with null body |
| 1943 | // statuses, as having null bodies. |
| 1944 | // See https://fetch.spec.whatwg.org/#main-fetch step 21. |
| 1945 | // |
| 1946 | // Note that we don't handle the CONNECT case here because kj-http handles CONNECT specially, |
| 1947 | // and the Fetch spec doesn't allow users to create Requests with CONNECT methods. |
| 1948 | kj::Maybe<Body::ExtractedBody> responseBody = kj::none; |
| 1949 | if (method != kj::HttpMethod::HEAD && !isNullBodyStatusCode(statusCode)) { |
| 1950 | responseBody = Body::ExtractedBody(js.alloc<ReadableStream>(context, |
| 1951 | newSystemStream(kj::mv(body), |
| 1952 | getContentEncoding(context, headers, bodyEncoding, FeatureFlags::get(js))))); |
| 1953 | } |
| 1954 | |
| 1955 | // The Fetch spec defines "response URLs" as having no fragments. Since the last URL in the list |
| 1956 | // is the one reported by Response::getUrl(), we nullify its fragment before serialization. |
| 1957 | kj::Array<kj::String> urlList; |
| 1958 | if (!urlListParam.empty()) { |
| 1959 | urlListParam.back().fragment = kj::none; |
| 1960 | urlList = KJ_MAP(url, urlListParam) { return url.toString(); }; |
| 1961 | } |
| 1962 | |
| 1963 | // TODO(someday): Fill response CF blob from somewhere? |
| 1964 | kj::Maybe<kj::String> maybeStatusText = statusText == defaultStatusText(statusCode) |
| 1965 | ? kj::Maybe<kj::String>() |
| 1966 | : kj::str(statusText); |
| 1967 | return js.alloc<Response>(js, statusCode, kj::mv(maybeStatusText), kj::mv(responseHeaders), |
| 1968 | nullptr, kj::mv(responseBody), kj::mv(urlList), kj::mv(webSocket), bodyEncoding); |
| 1969 | } |
| 1970 | |
| 1971 | namespace { |
| 1972 | |
| 1973 | jsg::Promise<jsg::Ref<Response>> fetchImplNoOutputLock(jsg::Lock& js, |
| 1974 | kj::Maybe<jsg::Ref<Fetcher>> fetcher, |
| 1975 | Request::Info requestOrUrl, |
| 1976 | jsg::Optional<Request::Initializer> requestInit) { |
| 1977 | // This use of evalNow() is obsoleted by the capture_async_api_throws compatibility flag, but |
| 1978 | // we need to keep it here for people who don't have that flag set. |
| 1979 | return js.evalNow([&] { |
| 1980 | // The spec requires us to call Request's constructor here, so we do. This is unfortunate, but |
| 1981 | // important for a few reasons: |
| 1982 | // |
| 1983 | // 1. If Request's constructor would throw, we must throw here, too. |
| 1984 | // 2. If `requestOrUrl` is a Request object, we must disturb its body immediately and leave it |
| 1985 | // disturbed. The typical fetch() call will do this naturally, except those which encounter |
| 1986 | // 303 redirects: they become GET requests with null bodies, which could then be reused. |
| 1987 | // 3. Following from the previous point, we must not allow the original request's method to |
| 1988 | // mutate. |
| 1989 | // |
| 1990 | // We could emulate these behaviors with various hacks, but just reconstructing the request up |
| 1991 | // front is robust, and won't add significant overhead compared to the rest of fetch(). |
| 1992 | auto jsRequest = Request::constructor(js, kj::mv(requestOrUrl), kj::mv(requestInit)); |
| 1993 | |
| 1994 | // Clear the request's signal if the 'ignoreForSubrequests' flag is set. This happens when |
| 1995 | // a request from an incoming fetch is passed-through to another fetch. We want to avoid |
| 1996 | // aborting the subrequest in that case. |
| 1997 | jsRequest->clearSignalIfIgnoredForSubrequest(js); |
| 1998 | |
| 1999 | // This URL list keeps track of redirections and becomes a source for Response's URL list. The |
| 2000 | // first URL in the list is the Request's URL (visible to JS via Request::getUrl()). The last URL |
| 2001 | // in the list is the Request's "current" URL (eventually visible to JS via Response::getUrl()). |
| 2002 | auto urlList = kj::Vector<kj::Url>(1 + MAX_REDIRECT_COUNT); |
| 2003 | |
| 2004 | jsg::Ref<Fetcher> actualFetcher = nullptr; |
| 2005 | KJ_IF_SOME(f, fetcher) { |
| 2006 | actualFetcher = kj::mv(f); |
| 2007 | } else KJ_IF_SOME(f, jsRequest->getFetcher()) { |
| 2008 | actualFetcher = kj::mv(f); |
| 2009 | } else { |
| 2010 | actualFetcher = |
| 2011 | js.alloc<Fetcher>(IoContext::NULL_CLIENT_CHANNEL, Fetcher::RequiresHostAndProtocol::YES); |
| 2012 | } |
| 2013 | |
| 2014 | KJ_IF_SOME(dataUrl, DataUrl::tryParse(jsRequest->getUrl())) { |
| 2015 | // If the URL is a data URL, we need to handle it specially. |
| 2016 | kj::Maybe<jsg::Ref<ReadableStream>> maybeResponseBody; |
| 2017 | auto type = dataUrl.getMimeType().toString(); |
| 2018 | |
| 2019 | // The Fetch spec defines responses to HEAD or CONNECT requests, or responses with null body |
| 2020 | // statuses, as having null bodies. |
| 2021 | // See https://fetch.spec.whatwg.org/#main-fetch step 21. |
| 2022 | // |
| 2023 | // Note that we don't handle the CONNECT case here because kj-http handles CONNECT specially, |
| 2024 | // and the Fetch spec doesn't allow users to create Requests with CONNECT methods. |
| 2025 | if (jsRequest->getMethodEnum() == kj::HttpMethod::GET) { |
| 2026 | auto view = dataUrl.getData(); |
| 2027 | auto rs = streams::newMemorySource(view, kj::heap(kj::mv(dataUrl))); |
| 2028 | maybeResponseBody.emplace(js.alloc<ReadableStream>(IoContext::current(), kj::mv(rs))); |
| 2029 | } |
| 2030 | |
| 2031 | auto headers = js.alloc<Headers>(); |
| 2032 | headers->setCommon(capnp::CommonHeaderName::CONTENT_TYPE, kj::mv(type)); |
| 2033 | return js.resolvedPromise(Response::constructor(js, kj::mv(maybeResponseBody), |
| 2034 | Response::InitializerDict{ |
| 2035 | .status = 200, |
| 2036 | .headers = kj::mv(headers), |
| 2037 | })); |
| 2038 | } |
| 2039 | |
| 2040 | urlList.add(actualFetcher->parseUrl(js, jsRequest->getUrl())); |
| 2041 | return fetchImplNoOutputLock(js, kj::mv(actualFetcher), kj::mv(jsRequest), kj::mv(urlList)); |
| 2042 | }); |
| 2043 | } |
| 2044 | |
| 2045 | } // namespace |
| 2046 | |
| 2047 | jsg::Promise<jsg::Ref<Response>> fetchImpl(jsg::Lock& js, |
| 2048 | kj::Maybe<jsg::Ref<Fetcher>> fetcher, |
| 2049 | Request::Info requestOrUrl, |
| 2050 | jsg::Optional<Request::Initializer> requestInit) { |
| 2051 | auto& context = IoContext::current(); |
| 2052 | // Optimization: For non-actors, which never have output locks, avoid the overhead of |
| 2053 | // awaitIo() and such by not going back to the event loop at all. |
| 2054 | KJ_IF_SOME(promise, context.waitForOutputLocksIfNecessary()) { |
| 2055 | return context.awaitIo(js, kj::mv(promise), |
| 2056 | [fetcher = kj::mv(fetcher), requestOrUrl = kj::mv(requestOrUrl), |
| 2057 | requestInit = kj::mv(requestInit)](jsg::Lock& js) mutable { |
| 2058 | return fetchImplNoOutputLock(js, kj::mv(fetcher), kj::mv(requestOrUrl), kj::mv(requestInit)); |
| 2059 | }); |
| 2060 | } else { |
| 2061 | return fetchImplNoOutputLock(js, kj::mv(fetcher), kj::mv(requestOrUrl), kj::mv(requestInit)); |
| 2062 | } |
| 2063 | } |
| 2064 | |
| 2065 | jsg::Ref<Socket> Fetcher::connect( |
| 2066 | jsg::Lock& js, AnySocketAddress address, jsg::Optional<SocketOptions> options) { |
| 2067 | return connectImpl(js, JSG_THIS, kj::mv(address), kj::mv(options)); |
| 2068 | } |
| 2069 | |
| 2070 | jsg::Promise<jsg::Ref<Response>> Fetcher::fetch(jsg::Lock& js, |
| 2071 | kj::OneOf<jsg::Ref<Request>, kj::String> requestOrUrl, |
| 2072 | jsg::Optional<kj::OneOf<RequestInitializerDict, jsg::Ref<Request>>> requestInit) { |
| 2073 | return fetchImpl(js, JSG_THIS, kj::mv(requestOrUrl), kj::mv(requestInit)); |
| 2074 | } |
| 2075 | |
| 2076 | kj::Maybe<jsg::Ref<JsRpcProperty>> Fetcher::getRpcMethod(jsg::Lock& js, kj::String name) { |
| 2077 | // This is like JsRpcStub::getRpcMethod(), but we also initiate a whole new JS RPC session |
| 2078 | // each time the method is called (handled by `getClientForOneCall()`, below). |
| 2079 | |
| 2080 | auto flags = FeatureFlags::get(js); |
| 2081 | if (!flags.getFetcherRpc() && !flags.getWorkerdExperimental()) { |
| 2082 | // We need to pretend that we haven't implemented a wildcard property, as unfortunately it |
| 2083 | // breaks some workers in the wild. We would, however, like to warn users who are trying to use |
| 2084 | // RPC so they understand why it isn't working. |
| 2085 | |
| 2086 | if (name == "idFromName") { |
| 2087 | // HACK specifically for itty-durable: We will not write any warning here, since itty-durable |
| 2088 | // automatically checks for this property on all bindings in an effort to discover Durable |
| 2089 | // Object namespaces. The warning would be confusing. |
| 2090 | // |
| 2091 | // Reported here: https://github.com/kwhitley/itty-durable/issues/48 |
| 2092 | } else { |
| 2093 | IoContext::current().logWarningOnce(kj::str("WARNING: Tried to access method or property '", |
| 2094 | name, |
| 2095 | "' on a Service Binding or " |
| 2096 | "Durable Object stub. Are you trying to use RPC? If so, please enable the 'rpc' compat " |
| 2097 | "flag or update your compat date to 2024-04-03 or later (see " |
| 2098 | "https://developers.cloudflare.com/workers/configuration/compatibility-dates/ ). If you " |
| 2099 | "are not trying to use RPC, please note that in the future, this property (and all other " |
| 2100 | "property names) will appear to be present as an RPC method.")); |
| 2101 | } |
| 2102 | |
| 2103 | return kj::none; |
| 2104 | } |
| 2105 | |
| 2106 | return getRpcMethodInternal(js, kj::mv(name)); |
| 2107 | } |
| 2108 | |
| 2109 | kj::Maybe<jsg::Ref<JsRpcProperty>> Fetcher::getRpcMethodInternal(jsg::Lock& js, kj::String name) { |
| 2110 | // Same as getRpcMethod, but skips compatibility check to allow RPC to be used from bindings |
| 2111 | // attached to workers without rpc flag. |
| 2112 | |
| 2113 | // Do not return a method for `then`, otherwise JavaScript decides this is a thenable, i.e. a |
| 2114 | // custom Promise, which will mean a Promise that resolves to this object will attempt to chain |
| 2115 | // with it, which is not what you want! |
| 2116 | if (name == "then"_kj) return kj::none; |
| 2117 | |
| 2118 | return js.alloc<JsRpcProperty>(JSG_THIS, kj::mv(name)); |
| 2119 | } |
| 2120 | |
| 2121 | rpc::JsRpcTarget::Client Fetcher::getClientForOneCall( |
| 2122 | jsg::Lock& js, kj::Vector<kj::StringPtr>& path) { |
| 2123 | auto& ioContext = IoContext::current(); |
| 2124 | auto worker = getClient(ioContext, kj::none, "jsRpcSession"_kjc); |
| 2125 | auto event = kj::heap<api::JsRpcSessionCustomEvent>( |
| 2126 | JsRpcSessionCustomEvent::WORKER_RPC_EVENT_TYPE); |
| 2127 | |
| 2128 | auto result = event->getCap(); |
| 2129 | |
| 2130 | // Arrange to cancel the CustomEvent if our I/O context is destroyed. But otherwise, we don't |
| 2131 | // actually care about the result of the event. If it throws, the membrane will already have |
| 2132 | // propagated the exception to any RPC calls that we're waiting on, so we even ignore errors |
| 2133 | // here -- otherwise they'll end up logged as "uncaught exceptions" even if they were, in fact, |
| 2134 | // caught elsewhere. |
| 2135 | ioContext.addTask(worker->customEvent(kj::mv(event)).attach(kj::mv(worker)).then([](auto&&) { |
| 2136 | }, [](kj::Exception&&) {})); |
| 2137 | |
| 2138 | // (Don't extend `path` because we're the root.) |
| 2139 | |
| 2140 | return result; |
| 2141 | } |
| 2142 | |
| 2143 | void Fetcher::serialize(jsg::Lock& js, jsg::Serializer& serializer) { |
| 2144 | auto channel = getSubrequestChannel(IoContext::current()); |
| 2145 | channel->requireAllowsTransfer(); |
| 2146 | |
| 2147 | KJ_IF_SOME(handler, serializer.getExternalHandler()) { |
| 2148 | KJ_IF_SOME(frankenvalueHandler, kj::tryDowncast<Frankenvalue::CapTableBuilder>(handler)) { |
| 2149 | // Encoding a Frankenvalue (e.g. for dynamic loopback props or dynamic isolate env). |
| 2150 | serializer.writeRawUint32(frankenvalueHandler.add(kj::mv(channel))); |
| 2151 | return; |
| 2152 | } else KJ_IF_SOME(rpcHandler, kj::tryDowncast<RpcSerializerExternalHandler>(handler)) { |
| 2153 | JSG_REQUIRE(FeatureFlags::get(js).getWorkerdExperimental(), DOMDataCloneError, |
| 2154 | "ServiceStub serialization requires the 'experimental' compat flag."); |
| 2155 | |
| 2156 | auto token = channel->getToken(IoChannelFactory::ChannelTokenUsage::RPC); |
| 2157 | rpcHandler.write([token = kj::mv(token)](rpc::JsValue::External::Builder builder) { |
| 2158 | builder.setSubrequestChannelToken(token); |
| 2159 | }); |
| 2160 | return; |
| 2161 | } |
| 2162 | // TODO(someday): structuredClone() should have special handling that just reproduces the same |
| 2163 | // local object. At present we have no way to recognize structuredClone() here though. |
| 2164 | } |
| 2165 | |
| 2166 | // The allow_irrevocable_stub_storage flag allows us to just embed the token inline. This format |
| 2167 | // is temporary, anyone using this will lose their data later. |
| 2168 | JSG_REQUIRE(FeatureFlags::get(js).getAllowIrrevocableStubStorage(), DOMDataCloneError, |
| 2169 | "ServiceStub cannot be serialized in this context."); |
| 2170 | serializer.writeLengthDelimited(channel->getToken(IoChannelFactory::ChannelTokenUsage::STORAGE)); |
| 2171 | } |
| 2172 | |
| 2173 | jsg::Ref<Fetcher> Fetcher::deserialize(jsg::Lock& js, |
| 2174 | rpc::SerializationTag tag, jsg::Deserializer& deserializer) { |
| 2175 | KJ_IF_SOME(handler, deserializer.getExternalHandler()) { |
| 2176 | KJ_IF_SOME(frankenvalueHandler, kj::tryDowncast<Frankenvalue::CapTableReader>(handler)) { |
| 2177 | // Decoding a Frankenvalue (e.g. for dynamic loopback props or dynamic isolate env). |
| 2178 | auto& cap = KJ_REQUIRE_NONNULL(frankenvalueHandler.get(deserializer.readRawUint32()), |
| 2179 | "serialized ServiceStub had invalid cap table index"); |
| 2180 | |
| 2181 | KJ_IF_SOME(channel, kj::tryDowncast<IoChannelFactory::SubrequestChannel>(cap)) { |
| 2182 | // Probably decoding dynamic ctx.props. |
| 2183 | return js.alloc<Fetcher>(IoContext::current().addObject(kj::addRef(channel))); |
| 2184 | } else KJ_IF_SOME(channel, kj::tryDowncast<IoChannelCapTableEntry>(cap)) { |
| 2185 | // Probably decoding dynamic isolate env. |
| 2186 | return js.alloc<Fetcher>( |
| 2187 | channel.getChannelNumber(IoChannelCapTableEntry::Type::SUBREQUEST), |
| 2188 | RequiresHostAndProtocol::YES, /*isInHouse=*/false); |
| 2189 | } else { |
| 2190 | KJ_FAIL_REQUIRE("ServiceStub capability in Frankenvalue is not a SubrequestChannel?"); |
| 2191 | } |
| 2192 | } else KJ_IF_SOME(rpcHandler, kj::tryDowncast<RpcDeserializerExternalHandler>(handler)) { |
| 2193 | JSG_REQUIRE(FeatureFlags::get(js).getWorkerdExperimental(), DOMDataCloneError, |
| 2194 | "ServiceStub serialization requires the 'experimental' compat flag."); |
| 2195 | |
| 2196 | auto external = rpcHandler.read(); |
| 2197 | KJ_REQUIRE(external.isSubrequestChannelToken()); |
| 2198 | auto& ioctx = IoContext::current(); |
| 2199 | auto channel = ioctx.getIoChannelFactory().subrequestChannelFromToken( |
| 2200 | IoChannelFactory::ChannelTokenUsage::RPC, |
| 2201 | external.getSubrequestChannelToken()); |
| 2202 | return js.alloc<Fetcher>(ioctx.addObject(kj::mv(channel))); |
| 2203 | } |
| 2204 | } |
| 2205 | |
| 2206 | // The allow_irrevocable_stub_storage flag allows us to just embed the token inline. This format |
| 2207 | // is temporary, anyone using this will lose their data later. |
| 2208 | JSG_REQUIRE(FeatureFlags::get(js).getAllowIrrevocableStubStorage(), DOMDataCloneError, |
| 2209 | "ServiceStub cannot be deserialized in this context."); |
| 2210 | auto& ioctx = IoContext::current(); |
| 2211 | auto channel = ioctx.getIoChannelFactory().subrequestChannelFromToken( |
| 2212 | IoChannelFactory::ChannelTokenUsage::STORAGE, deserializer.readLengthDelimitedBytes()); |
| 2213 | return js.alloc<Fetcher>(ioctx.addObject(kj::mv(channel))); |
| 2214 | } |
| 2215 | |
| 2216 | static jsg::Promise<void> throwOnError( |
| 2217 | jsg::Lock& js, kj::StringPtr method, jsg::Promise<jsg::Ref<Response>> promise) { |
| 2218 | return promise.then(js, [method](jsg::Lock&, jsg::Ref<Response> response) { |
| 2219 | uint status = response->getStatus(); |
| 2220 | // TODO(someday): Would be nice to attach the response to the JavaScript error, maybe? Or |
| 2221 | // should people really use fetch() if they want to inspect error responses? |
| 2222 | JSG_REQUIRE(status >= 200 && status < 300, Error, |
| 2223 | kj::str("HTTP ", method, " request failed: ", response->getStatus(), " ", |
| 2224 | response->getStatusText())); |
| 2225 | }); |
| 2226 | } |
| 2227 | |
| 2228 | static jsg::Promise<Fetcher::GetResult> parseResponse( |
| 2229 | jsg::Lock& js, jsg::Ref<Response> response, jsg::Optional<kj::String> type) { |
| 2230 | auto typeName = |
| 2231 | type.map([](const kj::String& s) -> kj::StringPtr { return s; }).orDefault("text"); |
| 2232 | if (typeName == "stream") { |
| 2233 | KJ_IF_SOME(body, response->getBody()) { |
| 2234 | return js.resolvedPromise(Fetcher::GetResult(kj::mv(body))); |
| 2235 | } else { |
| 2236 | // Empty body. |
| 2237 | return js.resolvedPromise(Fetcher::GetResult(js.alloc<ReadableStream>( |
| 2238 | IoContext::current(), newSystemStream(newNullInputStream(), StreamEncoding::IDENTITY)))); |
| 2239 | } |
| 2240 | } |
| 2241 | |
| 2242 | if (typeName == "text") { |
| 2243 | return response->text(js).then(js, [response = kj::mv(response)](jsg::Lock&, auto x) { |
| 2244 | return Fetcher::GetResult(kj::mv(x)); |
| 2245 | }); |
| 2246 | } else if (typeName == "arrayBuffer") { |
| 2247 | return response->arrayBuffer(js).then(js, [response = kj::mv(response)](jsg::Lock&, auto x) { |
| 2248 | return Fetcher::GetResult(kj::mv(x)); |
| 2249 | }); |
| 2250 | } else if (typeName == "json") { |
| 2251 | return response->json(js).then(js, [response = kj::mv(response)](jsg::Lock&, auto x) { |
| 2252 | return Fetcher::GetResult(kj::mv(x)); |
| 2253 | }); |
| 2254 | } else { |
| 2255 | JSG_FAIL_REQUIRE(TypeError, |
| 2256 | "Unknown response type. Possible types are \"text\", \"arrayBuffer\", " |
| 2257 | "\"json\", and \"stream\"."); |
| 2258 | } |
| 2259 | } |
| 2260 | |
| 2261 | jsg::Promise<Fetcher::GetResult> Fetcher::get( |
| 2262 | jsg::Lock& js, kj::String url, jsg::Optional<kj::String> type) { |
| 2263 | RequestInitializerDict subInit; |
| 2264 | subInit.method = kj::str("GET"); |
| 2265 | |
| 2266 | return fetchImpl(js, JSG_THIS, kj::mv(url), kj::mv(subInit)) |
| 2267 | .then(js, |
| 2268 | [type = kj::mv(type)]( |
| 2269 | jsg::Lock& js, jsg::Ref<Response> response) mutable -> jsg::Promise<GetResult> { |
| 2270 | uint status = response->getStatus(); |
| 2271 | if (status == 404 || status == 410) { |
| 2272 | return js.resolvedPromise(GetResult(js.v8Ref(js.v8Null()))); |
| 2273 | } else if (!response->getOk()) { |
| 2274 | // Manually construct exception so that we can incorporate method and status into the text |
| 2275 | // that JavaScript sees. |
| 2276 | // TODO(someday): Would be nice to attach the response to the JavaScript error, maybe? Or |
| 2277 | // should people really use fetch() if they want to inspect error responses? |
| 2278 | JSG_FAIL_REQUIRE(Error, |
| 2279 | kj::str( |
| 2280 | "HTTP GET request failed: ", response->getStatus(), " ", response->getStatusText())); |
| 2281 | } else { |
| 2282 | return parseResponse(js, kj::mv(response), kj::mv(type)); |
| 2283 | } |
| 2284 | }); |
| 2285 | } |
| 2286 | |
| 2287 | jsg::Promise<void> Fetcher::put(jsg::Lock& js, |
| 2288 | kj::String url, |
| 2289 | Body::Initializer body, |
| 2290 | jsg::Optional<Fetcher::PutOptions> options) { |
| 2291 | // Note that this borrows liberally from fetchImpl(fetcher, request, init, isolate). |
| 2292 | // This use of evalNow() is obsoleted by the capture_async_api_throws compatibility flag, but |
| 2293 | // we need to keep it here for people who don't have that flag set. |
| 2294 | return throwOnError(js, "PUT", js.evalNow([&] { |
| 2295 | RequestInitializerDict subInit; |
| 2296 | subInit.method = kj::str("PUT"); |
| 2297 | subInit.body = kj::mv(body); |
| 2298 | auto jsRequest = Request::constructor(js, kj::mv(url), kj::mv(subInit)); |
| 2299 | auto urlList = kj::Vector<kj::Url>(1 + MAX_REDIRECT_COUNT); |
| 2300 | |
| 2301 | kj::Url parsedUrl = this->parseUrl(js, jsRequest->getUrl()); |
| 2302 | |
| 2303 | // If any optional parameters were specified by the client, append them to |
| 2304 | // the URL's query parameters. |
| 2305 | KJ_IF_SOME(o, options) { |
| 2306 | KJ_IF_SOME(expiration, o.expiration) { |
| 2307 | parsedUrl.query.add(kj::Url::QueryParam{kj::str("expiration"), kj::str(expiration)}); |
| 2308 | } |
| 2309 | KJ_IF_SOME(expirationTtl, o.expirationTtl) { |
| 2310 | parsedUrl.query.add(kj::Url::QueryParam{kj::str("expiration_ttl"), kj::str(expirationTtl)}); |
| 2311 | } |
| 2312 | } |
| 2313 | |
| 2314 | urlList.add(kj::mv(parsedUrl)); |
| 2315 | return fetchImpl(js, JSG_THIS, kj::mv(jsRequest), kj::mv(urlList)); |
| 2316 | })); |
| 2317 | } |
| 2318 | |
| 2319 | jsg::Promise<void> Fetcher::delete_(jsg::Lock& js, kj::String url) { |
| 2320 | RequestInitializerDict subInit; |
| 2321 | subInit.method = kj::str("DELETE"); |
| 2322 | return throwOnError(js, "DELETE", fetchImpl(js, JSG_THIS, kj::mv(url), kj::mv(subInit))); |
| 2323 | } |
| 2324 | |
| 2325 | jsg::Promise<Fetcher::QueueResult> Fetcher::queue(jsg::Lock& js, |
| 2326 | kj::String queueName, |
| 2327 | kj::Array<ServiceBindingQueueMessage> messages, |
| 2328 | jsg::Optional<MessageBatchMetadata> metadata) { |
| 2329 | auto& ioContext = IoContext::current(); |
| 2330 | |
| 2331 | auto encodedMessages = kj::heapArrayBuilder<IncomingQueueMessage>(messages.size()); |
| 2332 | for (auto& msg: messages) { |
| 2333 | KJ_IF_SOME(b, msg.body) { |
| 2334 | JSG_REQUIRE(msg.serializedBody == kj::none, TypeError, |
| 2335 | "Expected one of body or serializedBody for each message"); |
| 2336 | jsg::Serializer serializer(js, |
| 2337 | jsg::Serializer::Options{ |
| 2338 | .version = 15, |
| 2339 | .omitHeader = false, |
| 2340 | }); |
| 2341 | serializer.write(js, jsg::JsValue(b.getHandle(js))); |
| 2342 | encodedMessages.add(IncomingQueueMessage{.id = kj::mv(msg.id), |
| 2343 | .timestamp = msg.timestamp, |
| 2344 | .body = serializer.release().data, |
| 2345 | .attempts = msg.attempts}); |
| 2346 | } else KJ_IF_SOME(b, msg.serializedBody) { |
| 2347 | encodedMessages.add(IncomingQueueMessage{.id = kj::mv(msg.id), |
| 2348 | .timestamp = msg.timestamp, |
| 2349 | .body = kj::mv(b), |
| 2350 | .attempts = msg.attempts}); |
| 2351 | } else { |
| 2352 | JSG_FAIL_REQUIRE(TypeError, "Expected one of body or serializedBody for each message"); |
| 2353 | } |
| 2354 | } |
| 2355 | |
| 2356 | // Only create worker interface after the error checks above to reduce overhead in case of errors. |
| 2357 | auto worker = getClient(ioContext, kj::none, "queue"_kjc); |
| 2358 | auto event = kj::refcounted<api::QueueCustomEvent>(QueueEvent::Params{ |
| 2359 | .queueName = kj::mv(queueName), |
| 2360 | .messages = encodedMessages.finish(), |
| 2361 | .metadata = kj::mv(metadata).orDefault({}), |
| 2362 | }); |
| 2363 | |
| 2364 | auto eventRef = |
| 2365 | kj::addRef(*event); // attempt to work around windows-specific null pointer deref. |
| 2366 | return ioContext.awaitIo(js, worker->customEvent(kj::mv(eventRef)).attach(kj::mv(worker)), |
| 2367 | [event = kj::mv(event)](jsg::Lock& js, WorkerInterface::CustomEvent::Result result) { |
| 2368 | return Fetcher::QueueResult{ |
| 2369 | .outcome = kj::str(result.outcome), |
| 2370 | .ackAll = event->getAckAll(), |
| 2371 | .retryBatch = event->getRetryBatch(), |
| 2372 | .explicitAcks = event->getExplicitAcks(), |
| 2373 | .retryMessages = event->getRetryMessages(), |
| 2374 | }; |
| 2375 | }); |
| 2376 | } |
| 2377 | |
| 2378 | jsg::Promise<Fetcher::ScheduledResult> Fetcher::scheduled( |
| 2379 | jsg::Lock& js, jsg::Optional<ScheduledOptions> options) { |
| 2380 | auto& ioContext = IoContext::current(); |
| 2381 | auto worker = getClient(ioContext, kj::none, "scheduled"_kjc); |
| 2382 | |
| 2383 | auto scheduledTime = ioContext.now(); |
| 2384 | auto cron = kj::String(); |
| 2385 | KJ_IF_SOME(o, options) { |
| 2386 | KJ_IF_SOME(t, o.scheduledTime) { |
| 2387 | scheduledTime = t; |
| 2388 | } |
| 2389 | KJ_IF_SOME(c, o.cron) { |
| 2390 | cron = kj::mv(c); |
| 2391 | } |
| 2392 | } |
| 2393 | |
| 2394 | return ioContext.awaitIo(js, |
| 2395 | worker->runScheduled(scheduledTime, cron).attach(kj::mv(worker), kj::mv(cron)), |
| 2396 | [](jsg::Lock& js, WorkerInterface::ScheduledResult result) { |
| 2397 | return Fetcher::ScheduledResult{ |
| 2398 | .outcome = kj::str(result.outcome), |
| 2399 | .noRetry = !result.retry, |
| 2400 | }; |
| 2401 | }); |
| 2402 | } |
| 2403 | |
| 2404 | kj::Own<WorkerInterface> Fetcher::getClient( |
| 2405 | IoContext& ioContext, kj::Maybe<kj::String> cfStr, kj::ConstString operationName) { |
| 2406 | auto clientWithTracing = getClientWithTracing(ioContext, kj::mv(cfStr), kj::mv(operationName)); |
| 2407 | return clientWithTracing.client.attach(kj::mv(clientWithTracing.traceContext)); |
| 2408 | } |
| 2409 | |
| 2410 | Fetcher::ClientWithTracing Fetcher::getClientWithTracing( |
| 2411 | IoContext& ioContext, kj::Maybe<kj::String> cfStr, kj::ConstString operationName) { |
| 2412 | KJ_SWITCH_ONEOF(channelOrClientFactory) { |
| 2413 | KJ_CASE_ONEOF(channel, uint) { |
| 2414 | // For channels, create trace context |
| 2415 | auto traceContext = ioContext.makeUserTraceSpan(kj::mv(operationName)); |
| 2416 | auto client = ioContext.getSubrequestChannel(channel, isInHouse, kj::mv(cfStr), traceContext); |
| 2417 | return ClientWithTracing{kj::mv(client), kj::mv(traceContext)}; |
| 2418 | } |
| 2419 | KJ_CASE_ONEOF(channel, IoOwn<IoChannelFactory::SubrequestChannel>) { |
| 2420 | auto traceContext = ioContext.makeUserTraceSpan(kj::mv(operationName)); |
| 2421 | auto client = ioContext.getSubrequest( |
| 2422 | [&](TraceContext& tracing, IoChannelFactory& ioChannelFactory) { |
| 2423 | return channel->startRequest({.cfBlobJson = kj::mv(cfStr), |
| 2424 | .parentSpan = tracing.getInternalSpanParent(), |
| 2425 | .userSpanParent = tracing.getUserSpanParent()}); |
| 2426 | }, { |
| 2427 | .inHouse = isInHouse, |
| 2428 | .wrapMetrics = !isInHouse, |
| 2429 | .existingTraceContext = traceContext, |
| 2430 | }); |
| 2431 | return ClientWithTracing{kj::mv(client), kj::mv(traceContext)}; |
| 2432 | } |
| 2433 | KJ_CASE_ONEOF(outgoingFactory, IoOwn<OutgoingFactory>) { |
| 2434 | // Outgoing factories are responsible for routing through getSubrequestNoChecks() (or |
| 2435 | // getSubrequest()) internally if they create HTTP connections, to ensure external memory |
| 2436 | // adjustment and other subrequest accounting are applied. |
| 2437 | auto client = outgoingFactory->newSingleUseClient(kj::mv(cfStr)); |
| 2438 | return ClientWithTracing{kj::mv(client), kj::none}; |
| 2439 | } |
| 2440 | KJ_CASE_ONEOF(outgoingFactory, kj::Own<CrossContextOutgoingFactory>) { |
| 2441 | // Same as OutgoingFactory above -- the factory is responsible for routing through |
| 2442 | // getSubrequestNoChecks() internally. |
| 2443 | auto client = outgoingFactory->newSingleUseClient(ioContext, kj::mv(cfStr)); |
| 2444 | return ClientWithTracing{kj::mv(client), kj::none}; |
| 2445 | } |
| 2446 | } |
| 2447 | KJ_UNREACHABLE; |
| 2448 | } |
| 2449 | |
| 2450 | kj::Own<IoChannelFactory::SubrequestChannel> Fetcher::getSubrequestChannel(IoContext& ioContext) { |
| 2451 | KJ_SWITCH_ONEOF(channelOrClientFactory) { |
| 2452 | KJ_CASE_ONEOF(channel, uint) { |
| 2453 | return ioContext.getIoChannelFactory().getSubrequestChannel(channel); |
| 2454 | } |
| 2455 | KJ_CASE_ONEOF(channel, IoOwn<IoChannelFactory::SubrequestChannel>) { |
| 2456 | return kj::addRef(*channel); |
| 2457 | } |
| 2458 | KJ_CASE_ONEOF(outgoingFactory, IoOwn<OutgoingFactory>) { |
| 2459 | return outgoingFactory->getSubrequestChannel(); |
| 2460 | } |
| 2461 | KJ_CASE_ONEOF(outgoingFactory, kj::Own<CrossContextOutgoingFactory>) { |
| 2462 | return outgoingFactory->getSubrequestChannel(ioContext); |
| 2463 | } |
| 2464 | } |
| 2465 | KJ_UNREACHABLE; |
| 2466 | } |
| 2467 | |
| 2468 | kj::Url Fetcher::parseUrl(jsg::Lock& js, kj::StringPtr url) { |
| 2469 | // We need to prep the request's URL for transmission over HTTP. fetch() accepts URLs that have |
| 2470 | // "." and ".." components as well as fragments (stuff after '#'), all of which needs to be |
| 2471 | // removed/collapsed before the URL is HTTP-ready. Luckily our URL parser does all this if we |
| 2472 | // tell it the context is REMOTE_HREF. |
| 2473 | constexpr auto urlOptions = kj::Url::Options{.percentDecode = false, .allowEmpty = true}; |
| 2474 | kj::Maybe<kj::Url> maybeParsed; |
| 2475 | if (this->requiresHost == RequiresHostAndProtocol::YES) { |
| 2476 | maybeParsed = kj::Url::tryParse(url, kj::Url::REMOTE_HREF, urlOptions); |
| 2477 | } else { |
| 2478 | // We don't require a protocol nor hostname, but we accept them. The easiest way to implement |
| 2479 | // this is to parse relative to a dummy URL. |
| 2480 | static const kj::Url FAKE = |
| 2481 | kj::Url::parse("https://fake-host/", kj::Url::REMOTE_HREF, urlOptions); |
| 2482 | maybeParsed = FAKE.tryParseRelative(url); |
| 2483 | } |
| 2484 | |
| 2485 | KJ_IF_SOME(p, maybeParsed) { |
| 2486 | if (p.scheme != "http" && p.scheme != "https") { |
| 2487 | // A non-HTTP scheme was requested. We should probably throw an exception, but historically |
| 2488 | // we actually went ahead and passed `X-Forwarded-Proto: whatever` to FL, which it happily |
| 2489 | // ignored if the protocol specified was not "https". Whoops. Unfortunately, some workers |
| 2490 | // in production have grown dependent on the bug. We'll have to use a runtime versioning flag |
| 2491 | // to fix this. |
| 2492 | |
| 2493 | if (FeatureFlags::get(js).getFetchRefusesUnknownProtocols()) { |
| 2494 | // Backwards-compatibility flag not enabled, so just fail. |
| 2495 | JSG_FAIL_REQUIRE(TypeError, kj::str("Fetch API cannot load: ", url)); |
| 2496 | } |
| 2497 | |
| 2498 | if (p.scheme != nullptr && '0' <= p.scheme[0] && p.scheme[0] <= '9') { |
| 2499 | // First character of the scheme is a digit. This is a weird case: Normally the KJ URL |
| 2500 | // parser would treat a scheme starting with a digit as invalid. But, due to a bug, |
| 2501 | // `tryParseRelative()` does NOT treat it as invalid. So, we know we took the branch above |
| 2502 | // that used `tryParseRelative()` above. In any case, later stages of the runtime will |
| 2503 | // definitely try to parse this URL again and will reject it at that time, producing an |
| 2504 | // internal error. We might as well throw a transparent error here instead so that we don't |
| 2505 | // log a garbage sentry alert. |
| 2506 | JSG_FAIL_REQUIRE(TypeError, kj::str("Fetch API cannot load: ", url)); |
| 2507 | } |
| 2508 | |
| 2509 | // In preview, log a warning in hopes that people fix this. |
| 2510 | kj::StringPtr more = nullptr; |
| 2511 | if (p.scheme == "ws" || p.scheme == "wss") { |
| 2512 | // Include some extra text for ws:// and wss:// specifically, since this is the most common |
| 2513 | // mistake. |
| 2514 | more = " Note that fetch() treats WebSockets as a special kind of HTTP request, " |
| 2515 | "therefore WebSockets should use 'http:'/'https:', not 'ws:'/'wss:'."; |
| 2516 | } else if (p.scheme == "ftp") { |
| 2517 | // Include some extra text for ftp://, since we see this sometimes. |
| 2518 | more = " fetch() does not support the FTP protocol."; |
| 2519 | } |
| 2520 | IoContext::current().logWarning( |
| 2521 | kj::str("Worker passed an invalid URL to fetch(). URLs passed to fetch() must begin with " |
| 2522 | "either 'http:' or 'https:', not '", |
| 2523 | p.scheme, |
| 2524 | ":'. Due to a historical bug, any " |
| 2525 | "other protocol used here will be treated the same as 'http:'. We plan to correct " |
| 2526 | "this bug in the future, so please update your Worker to use 'http:' or 'https:' for " |
| 2527 | "all fetch() URLs.", |
| 2528 | more)); |
| 2529 | } |
| 2530 | |
| 2531 | return kj::mv(p); |
| 2532 | } else { |
| 2533 | JSG_FAIL_REQUIRE(TypeError, kj::str("Fetch API cannot load: ", url)); |
| 2534 | } |
| 2535 | } |
| 2536 | } // namespace workerd::api |