File
Blob: src/workerd/api/r2-rpc.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 "r2-rpc.h" |
| 6 | |
| 7 | #include <workerd/api/r2-api.capnp.h> |
| 8 | #include <workerd/api/system-streams.h> |
| 9 | #include <workerd/api/util.h> |
| 10 | #include <workerd/util/http-util.h> |
| 11 | // This is imported for the error type and that's shared between internal and public beta. |
| 12 | |
| 13 | #include <capnp/compat/json.h> |
| 14 | #include <capnp/message.h> |
| 15 | #include <kj/compat/http.h> |
| 16 | |
| 17 | namespace workerd::api { |
| 18 | static kj::Own<R2Error> toError(uint statusCode, kj::StringPtr responseBody) { |
| 19 | capnp::JsonCodec json; |
| 20 | json.handleByAnnotation<public_beta::R2ErrorResponse>(); |
| 21 | capnp::MallocMessageBuilder errorMessageArena; |
| 22 | auto errorMessage = errorMessageArena.initRoot<public_beta::R2ErrorResponse>(); |
| 23 | json.decode(responseBody, errorMessage); |
| 24 | |
| 25 | return kj::refcounted<R2Error>(errorMessage.getV4code(), kj::str(errorMessage.getMessage())); |
| 26 | } |
| 27 | |
| 28 | jsg::JsValue R2Error::getStack(jsg::Lock& js) { |
| 29 | return jsg::JsObject(KJ_ASSERT_NONNULL(errorForStack).Get(js.v8Isolate)).get(js, "stack"_kj); |
| 30 | } |
| 31 | |
| 32 | kj::Maybe<uint> R2Result::v4ErrorCode() { |
| 33 | KJ_IF_SOME(e, toThrow) { |
| 34 | return e->v4Code; |
| 35 | } |
| 36 | return kj::none; |
| 37 | } |
| 38 | |
| 39 | kj::Maybe<kj::String> R2Result::getR2ErrorMessage() { |
| 40 | KJ_IF_SOME(e, toThrow) { |
| 41 | return kj::str(e->getMessage()); |
| 42 | } |
| 43 | return kj::none; |
| 44 | } |
| 45 | |
| 46 | void R2Result::throwIfError( |
| 47 | kj::StringPtr action, const jsg::TypeHandler<jsg::Ref<R2Error>>& errorType) { |
| 48 | KJ_IF_SOME(e, toThrow) { |
| 49 | // TODO(soon): Once jsg::JsPromise exists, switch to using that to tunnel out the exception. As |
| 50 | // it stands today, unfortunately, all we can send back to the user is a message. R2Error isn't |
| 51 | // a registered type in the runtime. When reenabling, make sure to update overrides/r2.d.ts to |
| 52 | // reenable the type |
| 53 | #if 0 |
| 54 | auto isolate = IoContext::current().getCurrentLock().getIsolate(); |
| 55 | (*e)->action = kj::str(action); |
| 56 | (*e)->errorForStack = v8::Global<v8::Object>( |
| 57 | isolate, v8::Exception::Error(v8::String::Empty(isolate)).As<v8::Object>()); |
| 58 | isolate->ThrowException(errorType.wrapRef(kj::mv(*e))); |
| 59 | throw jsg::JsExceptionThrown(); |
| 60 | #else |
| 61 | JSG_FAIL_REQUIRE(Error, kj::str(action, ": ", e.get()->getMessage(), " (", e->v4Code, ')')); |
| 62 | #endif |
| 63 | } |
| 64 | } |
| 65 | |
| 66 | namespace { |
| 67 | kj::String getFakeUrl(kj::ArrayPtr<kj::StringPtr> path) { |
| 68 | kj::Url url; |
| 69 | url.scheme = kj::str("https"); |
| 70 | url.host = kj::str("fake-host"); |
| 71 | for (const auto& p: path) { |
| 72 | url.path.add(kj::str(p)); |
| 73 | } |
| 74 | return url.toString(kj::Url::Context::HTTP_PROXY_REQUEST); |
| 75 | } |
| 76 | } // namespace |
| 77 | |
| 78 | kj::Promise<R2Result> doR2HTTPGetRequest(kj::Own<kj::HttpClient> client, |
| 79 | kj::String metadataPayload, |
| 80 | kj::ArrayPtr<kj::StringPtr> path, |
| 81 | kj::Maybe<kj::StringPtr> jwt, |
| 82 | CompatibilityFlags::Reader flags) { |
| 83 | auto& context = IoContext::current(); |
| 84 | auto url = getFakeUrl(path); |
| 85 | |
| 86 | auto& headerIds = context.getHeaderIds(); |
| 87 | |
| 88 | auto requestHeaders = kj::HttpHeaders(context.getHeaderTable()); |
| 89 | requestHeaders.set(headerIds.cfBlobRequest, kj::mv(metadataPayload)); |
| 90 | KJ_IF_SOME(j, jwt) { |
| 91 | requestHeaders.set(headerIds.authorization, kj::str("Bearer ", j)); |
| 92 | } |
| 93 | |
| 94 | static auto constexpr processStream = |
| 95 | [](kj::StringPtr metadata, kj::HttpClient::Response& response, kj::Own<kj::HttpClient> client, |
| 96 | CompatibilityFlags::Reader flags, IoContext& context) -> kj::Promise<R2Result> { |
| 97 | auto stream = newSystemStream(response.body.attach(kj::mv(client)), |
| 98 | getContentEncoding(context, *response.headers, Response::BodyEncoding::AUTO, flags), |
| 99 | context); |
| 100 | auto metadataSize = atoi((metadata).cStr()); |
| 101 | // R2 itself will try to stick to a cap of 256 KiB of response here. However for listing |
| 102 | // sometimes our heuristics have corner cases. This way we're more lenient in case someone |
| 103 | // finds a corner case for the heuristic so that we don't fail the GET with an opaque |
| 104 | // internal error. |
| 105 | KJ_REQUIRE(metadataSize <= 1024 * 1024, "R2 metadata size seems way too large"); |
| 106 | KJ_REQUIRE(metadataSize >= 0, "R2 metadata size parsed as negative"); |
| 107 | |
| 108 | auto metadataBuffer = kj::heapArray<char>(metadataSize); |
| 109 | auto metadataReadLength = |
| 110 | co_await stream->tryRead(metadataBuffer.begin(), metadataSize, metadataSize); |
| 111 | |
| 112 | KJ_ASSERT( |
| 113 | metadataReadLength == metadataBuffer.size(), "R2 metadata buffer not read fully/overflow?"); |
| 114 | |
| 115 | co_return R2Result{.httpStatus = response.statusCode, |
| 116 | .metadataPayload = kj::mv(metadataBuffer), |
| 117 | .stream = kj::mv(stream)}; |
| 118 | }; |
| 119 | |
| 120 | auto request = |
| 121 | client->request(kj::HttpMethod::GET, url, requestHeaders, static_cast<uint64_t>(0)); |
| 122 | |
| 123 | auto response = co_await request.response; |
| 124 | |
| 125 | if (response.statusCode >= 400) { |
| 126 | // Error responses should have a cfR2ErrorHeader but don't always. If there |
| 127 | // isn't one, we'll use a generic error. |
| 128 | if (response.headers->get(headerIds.cfR2ErrorHeader) == kj::none) { |
| 129 | LOG_WARNING_ONCE( |
| 130 | "R2 error response does not contain the CF-R2-Error header.", response.statusCode); |
| 131 | } |
| 132 | auto error = |
| 133 | response.headers->get(headerIds.cfR2ErrorHeader) |
| 134 | .orDefault("{\"version\":0,\"v4code\":0,\"message\":\"Unspecified error\"}"_kj); |
| 135 | |
| 136 | R2Result result = { |
| 137 | .httpStatus = response.statusCode, |
| 138 | .toThrow = toError(response.statusCode, error), |
| 139 | }; |
| 140 | |
| 141 | KJ_IF_SOME(m, response.headers->get(headerIds.cfBlobMetadataSize)) { |
| 142 | auto processed = co_await processStream(m, response, kj::mv(client), flags, context); |
| 143 | result.metadataPayload = kj::mv(processed.metadataPayload); |
| 144 | result.stream = kj::mv(processed.stream); |
| 145 | } |
| 146 | |
| 147 | co_return kj::mv(result); |
| 148 | } |
| 149 | |
| 150 | KJ_IF_SOME(m, response.headers->get(headerIds.cfBlobMetadataSize)) { |
| 151 | co_return co_await processStream(m, response, kj::mv(client), flags, context); |
| 152 | } else { |
| 153 | co_return R2Result{.httpStatus = response.statusCode}; |
| 154 | } |
| 155 | } |
| 156 | |
| 157 | kj::Promise<R2Result> doR2HTTPPutRequest(kj::Own<kj::HttpClient> client, |
| 158 | kj::Maybe<R2PutValue> supportedBody, |
| 159 | kj::Maybe<uint64_t> streamSize, |
| 160 | kj::String metadataPayload, |
| 161 | kj::ArrayPtr<kj::StringPtr> path, |
| 162 | kj::Maybe<kj::StringPtr> jwt) { |
| 163 | // NOTE: A lot of code here is duplicated with kv.c++. Maybe it can be refactored to be more |
| 164 | // reusable? |
| 165 | auto& context = IoContext::current(); |
| 166 | auto headers = kj::HttpHeaders(context.getHeaderTable()); |
| 167 | auto url = getFakeUrl(path); |
| 168 | |
| 169 | kj::Maybe<uint64_t> expectedBodySize; |
| 170 | |
| 171 | KJ_IF_SOME(b, supportedBody) { |
| 172 | KJ_SWITCH_ONEOF(b) { |
| 173 | KJ_CASE_ONEOF(stream, jsg::Ref<ReadableStream>) { |
| 174 | expectedBodySize = stream->tryGetLength(StreamEncoding::IDENTITY); |
| 175 | if (expectedBodySize == kj::none) { |
| 176 | expectedBodySize = streamSize; |
| 177 | } |
| 178 | JSG_REQUIRE(expectedBodySize != kj::none, TypeError, |
| 179 | "Provided readable stream must have a known length (request/response body or readable " |
| 180 | "half of FixedLengthStream)"); |
| 181 | JSG_REQUIRE(streamSize.orDefault(KJ_ASSERT_NONNULL(expectedBodySize)) == expectedBodySize, |
| 182 | RangeError, "Provided stream length (", streamSize.orDefault(-1), |
| 183 | ") doesn't match what " |
| 184 | "the stream reports (", |
| 185 | KJ_ASSERT_NONNULL(expectedBodySize), ")"); |
| 186 | } |
| 187 | KJ_CASE_ONEOF(text, jsg::NonCoercible<kj::String>) { |
| 188 | expectedBodySize = text.value.size(); |
| 189 | KJ_REQUIRE(streamSize == kj::none); |
| 190 | } |
| 191 | KJ_CASE_ONEOF(data, kj::Array<kj::byte>) { |
| 192 | expectedBodySize = data.size(); |
| 193 | KJ_REQUIRE(streamSize == kj::none); |
| 194 | } |
| 195 | KJ_CASE_ONEOF(data, jsg::Ref<Blob>) { |
| 196 | expectedBodySize = data->getSize(); |
| 197 | KJ_REQUIRE(streamSize == kj::none); |
| 198 | } |
| 199 | } |
| 200 | } else { |
| 201 | expectedBodySize = static_cast<uint64_t>(0); |
| 202 | KJ_REQUIRE(streamSize == kj::none); |
| 203 | } |
| 204 | |
| 205 | headers.set(context.getHeaderIds().cfBlobMetadataSize, kj::str(metadataPayload.size())); |
| 206 | KJ_IF_SOME(j, jwt) { |
| 207 | headers.set(context.getHeaderIds().authorization, kj::str("Bearer ", j)); |
| 208 | } |
| 209 | |
| 210 | uint64_t combinedSize = metadataPayload.size() + KJ_ASSERT_NONNULL(expectedBodySize); |
| 211 | |
| 212 | co_await context.waitForOutputLocks(); |
| 213 | |
| 214 | auto request = client->request(kj::HttpMethod::PUT, url, headers, combinedSize); |
| 215 | |
| 216 | co_await request.body->write(metadataPayload.asBytes()); |
| 217 | |
| 218 | KJ_IF_SOME(b, supportedBody) { |
| 219 | KJ_SWITCH_ONEOF(b) { |
| 220 | KJ_CASE_ONEOF(text, jsg::NonCoercible<kj::String>) { |
| 221 | co_await request.body->write(text.value.asBytes()); |
| 222 | } |
| 223 | KJ_CASE_ONEOF(data, kj::Array<byte>) { |
| 224 | co_await request.body->write(data); |
| 225 | } |
| 226 | KJ_CASE_ONEOF(blob, jsg::Ref<Blob>) { |
| 227 | auto data = blob->getData(); |
| 228 | co_await request.body->write(data); |
| 229 | } |
| 230 | KJ_CASE_ONEOF(stream, jsg::Ref<ReadableStream>) { |
| 231 | // Because the ReadableStream might be a fully JavaScript-backed stream, we must |
| 232 | // start running the pump within the IoContext/isolate lock. |
| 233 | co_await context.run( |
| 234 | [dest = newSystemStream(kj::mv(request.body), StreamEncoding::IDENTITY, context), |
| 235 | stream = kj::mv(stream)](jsg::Lock& js) mutable { |
| 236 | return IoContext::current().waitForDeferredProxy(stream->pumpTo(js, kj::mv(dest), true)); |
| 237 | }); |
| 238 | } |
| 239 | } |
| 240 | } |
| 241 | |
| 242 | auto response = co_await request.response; |
| 243 | |
| 244 | if (response.statusCode >= 400) { |
| 245 | // Error responses should have a cfR2ErrorHeader but don't always. If there |
| 246 | // isn't one, we'll use a generic error. |
| 247 | auto& headerIds = context.getHeaderIds(); |
| 248 | if (response.headers->get(headerIds.cfR2ErrorHeader) == kj::none) { |
| 249 | LOG_WARNING_ONCE( |
| 250 | "R2 error response does not contain the CF-R2-Error header.", response.statusCode); |
| 251 | } |
| 252 | auto error = |
| 253 | response.headers->get(headerIds.cfR2ErrorHeader) |
| 254 | .orDefault("{\"version\":0,\"v4code\":0,\"message\":\"Unspecified error\"}"_kj); |
| 255 | |
| 256 | co_return R2Result{ |
| 257 | .httpStatus = response.statusCode, |
| 258 | .toThrow = toError(response.statusCode, error), |
| 259 | }; |
| 260 | } |
| 261 | |
| 262 | auto responseBody = co_await response.body->readAllText(); |
| 263 | |
| 264 | co_return R2Result{ |
| 265 | .httpStatus = response.statusCode, |
| 266 | .metadataPayload = responseBody.releaseArray(), |
| 267 | }; |
| 268 | } |
| 269 | } // namespace workerd::api |