File
Blob: src/workerd/api/kv.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 "kv.h" |
| 6 | |
| 7 | #include "system-streams.h" |
| 8 | #include "util.h" |
| 9 | |
| 10 | #include <workerd/io/features.h> |
| 11 | #include <workerd/io/io-context.h> |
| 12 | #include <workerd/io/limit-enforcer.h> |
| 13 | #include <workerd/util/http-util.h> |
| 14 | #include <workerd/util/mimetype.h> |
| 15 | |
| 16 | #include <kj/compat/http.h> |
| 17 | #include <kj/encoding.h> |
| 18 | |
| 19 | namespace workerd::api { |
| 20 | |
| 21 | // As documented in Cloudflare's Worker KV limits. |
| 22 | static constexpr size_t kMaxKeyLength = 512; |
| 23 | |
| 24 | static void checkForErrorStatus(kj::StringPtr method, const kj::HttpClient::Response& response) { |
| 25 | if (response.statusCode < 200 || response.statusCode >= 300) { |
| 26 | // Manually construct exception so that we can incorporate method and status into the text |
| 27 | // that JavaScript sees. |
| 28 | kj::throwFatalException(kj::Exception(kj::Exception::Type::FAILED, __FILE__, __LINE__, |
| 29 | kj::str(JSG_EXCEPTION(Error) ": KV ", method, " failed: ", response.statusCode, ' ', |
| 30 | response.statusText))); |
| 31 | } |
| 32 | } |
| 33 | |
| 34 | static void validateKeyName(kj::StringPtr method, kj::StringPtr name) { |
| 35 | JSG_REQUIRE(name != "", TypeError, "Key name cannot be empty."); |
| 36 | JSG_REQUIRE(name != ".", TypeError, "\".\" is not allowed as a key name."); |
| 37 | JSG_REQUIRE(name != "..", TypeError, "\"..\" is not allowed as a key name."); |
| 38 | JSG_REQUIRE(name.size() <= kMaxKeyLength, Error, "KV ", method, " failed: ", 414, |
| 39 | " UTF-8 encoded length of ", name.size(), " exceeds key length limit of ", kMaxKeyLength, |
| 40 | "."); |
| 41 | } |
| 42 | |
| 43 | static void parseListMetadata(TraceContext& traceContext, |
| 44 | jsg::Lock& js, |
| 45 | jsg::JsValue listResponse, |
| 46 | kj::Maybe<jsg::JsValue> cacheStatus) { |
| 47 | static constexpr auto METADATA = "metadata"_kjc; |
| 48 | static constexpr auto KEYS = "keys"_kjc; |
| 49 | static constexpr auto CURSOR = "cursor"_kjc; |
| 50 | static constexpr auto LIST_COMPLETE = "list_complete"_kjc; |
| 51 | static constexpr auto EXPIRATION = "expiration"_kjc; |
| 52 | |
| 53 | js.withinHandleScope([&] { |
| 54 | auto obj = KJ_ASSERT_NONNULL(listResponse.tryCast<jsg::JsObject>()); |
| 55 | |
| 56 | KJ_IF_SOME(boolVal, obj.get(js, LIST_COMPLETE).tryCast<jsg::JsBoolean>()) { |
| 57 | traceContext.setTag("cloudflare.kv.response.list_complete"_kjc, boolVal.value(js)); |
| 58 | } |
| 59 | |
| 60 | KJ_IF_SOME(cursor, obj.get(js, CURSOR).tryCast<jsg::JsString>()) { |
| 61 | traceContext.setTag("cloudflare.kv.response.cursor"_kjc, kj::str(cursor)); |
| 62 | } |
| 63 | |
| 64 | KJ_IF_SOME(expiration, obj.get(js, EXPIRATION).tryCast<jsg::JsNumber>()) { |
| 65 | KJ_IF_SOME(value, expiration.value(js)) { |
| 66 | traceContext.setTag("cloudflare.kv.response.expiration"_kjc, static_cast<int64_t>(value)); |
| 67 | } |
| 68 | } |
| 69 | |
| 70 | KJ_IF_SOME(keysArr, obj.get(js, KEYS).tryCast<jsg::JsArray>()) { |
| 71 | auto length = keysArr.size(); |
| 72 | traceContext.setTag("cloudflare.kv.response.returned_rows"_kjc, static_cast<int64_t>(length)); |
| 73 | for (int i = 0; i < length; i++) { |
| 74 | js.withinHandleScope([&] { |
| 75 | KJ_IF_SOME(key, keysArr.get(js, i).tryCast<jsg::JsObject>()) { |
| 76 | KJ_IF_SOME(str, key.get(js, METADATA).tryCast<jsg::JsString>()) { |
| 77 | key.set(js, METADATA, jsg::JsValue::fromJson(js, str)); |
| 78 | } |
| 79 | } |
| 80 | }); |
| 81 | } |
| 82 | } |
| 83 | |
| 84 | obj.set(js, "cacheStatus"_kjc, cacheStatus.orDefault(js.null())); |
| 85 | }); |
| 86 | } |
| 87 | |
| 88 | constexpr auto FLPROD_405_HEADER = "CF-KV-FLPROD-405"_kj; |
| 89 | |
| 90 | kj::Own<kj::HttpClient> KvNamespace::getHttpClient(IoContext& context, |
| 91 | kj::HttpHeaders& headers, |
| 92 | kj::OneOf<LimitEnforcer::KvOpType, kj::LiteralStringConst> opTypeOrName, |
| 93 | kj::StringPtr urlStr, |
| 94 | TraceContext& traceContext) { |
| 95 | |
| 96 | KJ_SWITCH_ONEOF(opTypeOrName) { |
| 97 | KJ_CASE_ONEOF(name, kj::LiteralStringConst) {} |
| 98 | KJ_CASE_ONEOF(opType, LimitEnforcer::KvOpType) { |
| 99 | // Check if we've hit KV usage limits. (This will throw if we have.) |
| 100 | context.getLimitEnforcer().newKvRequest(opType); |
| 101 | } |
| 102 | } |
| 103 | |
| 104 | auto client = context.getHttpClient(subrequestChannel, true, kj::none, traceContext); |
| 105 | |
| 106 | headers.addPtrPtr(FLPROD_405_HEADER, urlStr); |
| 107 | for (const auto& header: additionalHeaders) { |
| 108 | headers.addPtrPtr(header.name.asPtr(), header.value.asPtr()); |
| 109 | } |
| 110 | |
| 111 | return client; |
| 112 | } |
| 113 | |
| 114 | jsg::Promise<KvNamespace::GetResult> KvNamespace::getSingle(jsg::Lock& js, |
| 115 | IoContext& context, |
| 116 | TraceContext& traceContext, |
| 117 | kj::String name, |
| 118 | jsg::Optional<kj::OneOf<kj::String, GetOptions>> options) { |
| 119 | return js.evalNow([&] { |
| 120 | auto resp = getWithMetadataImpl( |
| 121 | js, context, traceContext, kj::mv(name), kj::mv(options), LimitEnforcer::KvOpType::GET); |
| 122 | return resp.then(js, |
| 123 | [](jsg::Lock&, KvNamespace::GetWithMetadataResult result) { return kj::mv(result.value); }); |
| 124 | }); |
| 125 | } |
| 126 | |
| 127 | jsg::Promise<jsg::JsRef<jsg::JsMap>> KvNamespace::getBulk(jsg::Lock& js, |
| 128 | IoContext& context, |
| 129 | TraceContext& traceContext, |
| 130 | kj::Array<kj::String> name, |
| 131 | jsg::Optional<kj::OneOf<kj::String, GetOptions>> options, |
| 132 | bool withMetadata) { |
| 133 | return js.evalNow([&] { |
| 134 | kj::Url url; |
| 135 | url.scheme = kj::str("https"); |
| 136 | url.host = kj::str("fake-host"); |
| 137 | url.path.add(kj::str("bulk")); |
| 138 | url.path.add(kj::str("get")); |
| 139 | |
| 140 | kj::String body = formBulkBodyString(js, name, withMetadata, options); |
| 141 | kj::Maybe<uint64_t> expectedBodySize = static_cast<uint64_t>(body.size()); |
| 142 | auto headers = kj::HttpHeaders(context.getHeaderTable()); |
| 143 | headers.set(kj::HttpHeaderId::CONTENT_TYPE, MimeType::JSON.toString()); |
| 144 | |
| 145 | auto urlStr = url.toString(kj::Url::Context::HTTP_PROXY_REQUEST); |
| 146 | |
| 147 | // This could be quite large, so let's limit the string length to 512 characters |
| 148 | auto keysStr = kj::strArray(name, ", "); |
| 149 | if (keysStr.size() > 512) { |
| 150 | keysStr = kj::str(keysStr.slice(0, 509), "..."); |
| 151 | } |
| 152 | traceContext.setTag("cloudflare.kv.query.keys"_kjc, kj::mv(keysStr)); |
| 153 | traceContext.setTag("cloudflare.kv.query.keys.count"_kjc, static_cast<int64_t>(name.size())); |
| 154 | |
| 155 | KJ_IF_SOME(_options, options) { |
| 156 | KJ_SWITCH_ONEOF(_options) { |
| 157 | KJ_CASE_ONEOF(type, kj::String) { |
| 158 | traceContext.setTag("cloudflare.kv.query.type"_kjc, kj::mv(type)); |
| 159 | } |
| 160 | KJ_CASE_ONEOF(o, GetOptions) { |
| 161 | KJ_IF_SOME(type, o.type) { |
| 162 | traceContext.setTag("cloudflare.kv.query.type"_kjc, kj::mv(type)); |
| 163 | } |
| 164 | KJ_IF_SOME(cacheTtl, o.cacheTtl) { |
| 165 | traceContext.setTag( |
| 166 | "cloudflare.kv.query.cache_ttl"_kjc, static_cast<int64_t>(cacheTtl)); |
| 167 | } |
| 168 | } |
| 169 | } |
| 170 | } |
| 171 | |
| 172 | auto client = |
| 173 | getHttpClient(context, headers, LimitEnforcer::KvOpType::GET_BULK, urlStr, traceContext); |
| 174 | |
| 175 | auto promise = context.waitForOutputLocks().then( |
| 176 | [client = kj::mv(client), urlStr = kj::mv(urlStr), headers = kj::mv(headers), |
| 177 | expectedBodySize, supportedBody = kj::mv(body)]() mutable { |
| 178 | auto innerReq = client->request(kj::HttpMethod::POST, urlStr, headers, expectedBodySize); |
| 179 | auto req = attachToRequest(kj::mv(innerReq), kj::refcountedWrapper(kj::mv(client))); |
| 180 | |
| 181 | kj::Promise<void> writePromise = nullptr; |
| 182 | writePromise = req.body->write(supportedBody.asBytes()).attach(kj::mv(supportedBody)); |
| 183 | |
| 184 | return writePromise.attach(kj::mv(req.body)).then([resp = kj::mv(req.response)]() mutable { |
| 185 | return resp.then([](kj::HttpClient::Response&& response) mutable { |
| 186 | checkForErrorStatus("GET_BULK", response); |
| 187 | return response.body->readAllText().attach(kj::mv(response.body)); |
| 188 | }); |
| 189 | }); |
| 190 | }); |
| 191 | |
| 192 | return context.awaitIo(js, kj::mv(promise), |
| 193 | [&, traceContext = kj::mv(traceContext)](jsg::Lock& js, kj::String text) mutable { |
| 194 | traceContext.setTag("cloudflare.kv.response.size"_kjc, static_cast<int64_t>(text.size())); |
| 195 | auto result = jsg::JsValue::fromJson(js, text); |
| 196 | auto map = js.map(); |
| 197 | KJ_IF_SOME(obj, result.tryCast<jsg::JsObject>()) { |
| 198 | auto values = obj.getPropertyNames(js, jsg::KeyCollectionFilter::OWN_ONLY, |
| 199 | jsg::PropertyFilter::SKIP_SYMBOLS, jsg::IndexFilter::SKIP_INDICES); |
| 200 | for (int i = 0; i < values.size(); i++) { |
| 201 | auto key = values.get(js, i); |
| 202 | map.set(js, kj::mv(key), obj.get(js, key)); |
| 203 | } |
| 204 | traceContext.setTag( |
| 205 | "cloudflare.kv.response.returned_rows"_kjc, static_cast<int64_t>(values.size())); |
| 206 | } |
| 207 | return jsg::JsRef(js, map); |
| 208 | }); |
| 209 | }); |
| 210 | } |
| 211 | |
| 212 | kj::String KvNamespace::formBulkBodyString(jsg::Lock& js, |
| 213 | kj::Array<kj::String>& names, |
| 214 | bool withMetadata, |
| 215 | jsg::Optional<kj::OneOf<kj::String, GetOptions>>& options) { |
| 216 | |
| 217 | kj::String type = kj::str(""); |
| 218 | kj::String cacheTtlStr = kj::str(""); |
| 219 | KJ_IF_SOME(oneOfOptions, options) { |
| 220 | KJ_SWITCH_ONEOF(oneOfOptions) { |
| 221 | KJ_CASE_ONEOF(t, kj::String) { |
| 222 | type = kj::str(t); |
| 223 | } |
| 224 | KJ_CASE_ONEOF(options, GetOptions) { |
| 225 | KJ_IF_SOME(t, options.type) { |
| 226 | type = kj::str(t); |
| 227 | } |
| 228 | KJ_IF_SOME(cacheTtl, options.cacheTtl) { |
| 229 | cacheTtlStr = kj::str(cacheTtl); |
| 230 | } |
| 231 | } |
| 232 | } |
| 233 | } |
| 234 | auto object = js.obj(); |
| 235 | |
| 236 | auto keysArray = |
| 237 | js.arr(names.asPtr(), [](jsg::Lock& js, const kj::String& val) { return js.str(val); }); |
| 238 | object.set(js, "keys", keysArray); |
| 239 | |
| 240 | if (type != kj::str("")) { |
| 241 | object.set(js, "type", js.str(type)); |
| 242 | } |
| 243 | if (withMetadata) { |
| 244 | object.set(js, "withMetadata", js.boolean(true)); |
| 245 | } |
| 246 | if (cacheTtlStr != kj::str("")) { |
| 247 | object.set(js, "cacheTtl", js.str(cacheTtlStr)); |
| 248 | } |
| 249 | return jsg::JsValue(object).toJson(js); |
| 250 | } |
| 251 | |
| 252 | kj::OneOf<jsg::Promise<KvNamespace::GetResult>, jsg::Promise<jsg::JsRef<jsg::JsMap>>> KvNamespace:: |
| 253 | get(jsg::Lock& js, |
| 254 | kj::OneOf<kj::String, kj::Array<kj::String>> name, |
| 255 | jsg::Optional<kj::OneOf<kj::String, GetOptions>> options) { |
| 256 | auto& context = IoContext::current(); |
| 257 | TraceContext traceContext = context.makeUserTraceSpan("kv_get"_kjc); |
| 258 | traceContext.setTag("db.system.name"_kjc, "cloudflare-kv"_kjc); |
| 259 | traceContext.setTag("db.operation.name"_kjc, "get"_kjc); |
| 260 | traceContext.setTag("cloudflare.binding.name"_kjc, bindingName.asPtr()); |
| 261 | traceContext.setTag("cloudflare.binding.type"_kjc, "KV"_kjc); |
| 262 | |
| 263 | KJ_SWITCH_ONEOF(name) { |
| 264 | KJ_CASE_ONEOF(arr, kj::Array<kj::String>) { |
| 265 | return context.attachSpans(js, |
| 266 | getBulk(js, context, traceContext, kj::mv(arr), kj::mv(options), false), |
| 267 | kj::mv(traceContext)); |
| 268 | } |
| 269 | KJ_CASE_ONEOF(str, kj::String) { |
| 270 | return context.attachSpans(js, |
| 271 | getSingle(js, context, traceContext, kj::mv(str), kj::mv(options)), kj::mv(traceContext)); |
| 272 | } |
| 273 | } |
| 274 | KJ_UNREACHABLE; |
| 275 | }; |
| 276 | |
| 277 | jsg::Promise<KvNamespace::GetWithMetadataResult> KvNamespace::getWithMetadataSingle(jsg::Lock& js, |
| 278 | IoContext& context, |
| 279 | TraceContext& traceContext, |
| 280 | kj::String name, |
| 281 | jsg::Optional<kj::OneOf<kj::String, GetOptions>> options) { |
| 282 | return getWithMetadataImpl( |
| 283 | js, context, traceContext, kj::mv(name), kj::mv(options), LimitEnforcer::KvOpType::GET_WITH); |
| 284 | } |
| 285 | |
| 286 | kj::OneOf<jsg::Promise<KvNamespace::GetWithMetadataResult>, jsg::Promise<jsg::JsRef<jsg::JsMap>>> |
| 287 | KvNamespace::getWithMetadata(jsg::Lock& js, |
| 288 | kj::OneOf<kj::Array<kj::String>, kj::String> name, |
| 289 | jsg::Optional<kj::OneOf<kj::String, GetOptions>> options) { |
| 290 | |
| 291 | auto& context = IoContext::current(); |
| 292 | TraceContext traceContext = context.makeUserTraceSpan("kv_getWithMetadata"_kjc); |
| 293 | traceContext.setTag("db.system.name"_kjc, "cloudflare-kv"_kjc); |
| 294 | traceContext.setTag("db.operation.name"_kjc, "get"_kjc); |
| 295 | traceContext.setTag("cloudflare.binding.name"_kjc, bindingName.asPtr()); |
| 296 | traceContext.setTag("cloudflare.binding.type"_kjc, "KV"_kjc); |
| 297 | KJ_SWITCH_ONEOF(name) { |
| 298 | KJ_CASE_ONEOF(arr, kj::Array<kj::String>) { |
| 299 | return context.attachSpans(js, |
| 300 | getBulk(js, context, traceContext, kj::mv(arr), kj::mv(options), true), |
| 301 | kj::mv(traceContext)); |
| 302 | } |
| 303 | KJ_CASE_ONEOF(str, kj::String) { |
| 304 | return context.attachSpans(js, |
| 305 | getWithMetadataSingle(js, context, traceContext, kj::mv(str), kj::mv(options)), |
| 306 | kj::mv(traceContext)); |
| 307 | } |
| 308 | } |
| 309 | KJ_UNREACHABLE; |
| 310 | } |
| 311 | |
| 312 | jsg::Promise<KvNamespace::GetWithMetadataResult> KvNamespace::getWithMetadataImpl(jsg::Lock& js, |
| 313 | IoContext& context, |
| 314 | TraceContext& traceContext, |
| 315 | kj::String name, |
| 316 | jsg::Optional<kj::OneOf<kj::String, GetOptions>> options, |
| 317 | LimitEnforcer::KvOpType op) { |
| 318 | validateKeyName("GET", name); |
| 319 | |
| 320 | traceContext.setTag("cloudflare.kv.query.keys"_kjc, name.asPtr()); |
| 321 | traceContext.setTag("cloudflare.kv.query.keys.count"_kjc, static_cast<int64_t>(1)); |
| 322 | |
| 323 | kj::Url url; |
| 324 | url.scheme = kj::str("https"); |
| 325 | url.host = kj::str("fake-host"); |
| 326 | url.path.add(kj::mv(name)); |
| 327 | url.query.add(kj::Url::QueryParam{kj::str("urlencoded"), kj::str("true")}); |
| 328 | |
| 329 | kj::Maybe<kj::String> type; |
| 330 | KJ_IF_SOME(oneOfOptions, options) { |
| 331 | KJ_SWITCH_ONEOF(oneOfOptions) { |
| 332 | KJ_CASE_ONEOF(t, kj::String) { |
| 333 | type = kj::str(t); |
| 334 | traceContext.setTag("cloudflare.kv.query.type"_kjc, kj::mv(t)); |
| 335 | } |
| 336 | KJ_CASE_ONEOF(options, GetOptions) { |
| 337 | KJ_IF_SOME(t, options.type) { |
| 338 | type = kj::str(t); |
| 339 | traceContext.setTag("cloudflare.kv.query.type"_kjc, kj::mv(t)); |
| 340 | } |
| 341 | KJ_IF_SOME(cacheTtl, options.cacheTtl) { |
| 342 | url.query.add(kj::Url::QueryParam{kj::str("cache_ttl"), kj::str(cacheTtl)}); |
| 343 | traceContext.setTag("cloudflare.kv.query.cache_ttl"_kjc, static_cast<int64_t>(cacheTtl)); |
| 344 | } |
| 345 | } |
| 346 | } |
| 347 | } |
| 348 | |
| 349 | auto urlStr = url.toString(kj::Url::Context::HTTP_PROXY_REQUEST); |
| 350 | |
| 351 | auto headers = kj::HttpHeaders(context.getHeaderTable()); |
| 352 | auto client = getHttpClient(context, headers, op, urlStr, traceContext); |
| 353 | |
| 354 | auto request = client->request(kj::HttpMethod::GET, urlStr, headers); |
| 355 | return context.awaitIo(js, kj::mv(request.response), |
| 356 | [type = kj::mv(type), &context, client = kj::mv(client), traceContext = kj::mv(traceContext)]( |
| 357 | jsg::Lock& js, kj::HttpClient::Response&& response) mutable |
| 358 | -> jsg::Promise<KvNamespace::GetWithMetadataResult> { |
| 359 | auto cacheStatus = |
| 360 | response.headers->get(context.getHeaderIds().cfCacheStatus).map([&](kj::StringPtr cs) { |
| 361 | traceContext.setTag("cloudflare.kv.response.cache_status"_kjc, cs); |
| 362 | return jsg::JsRef<jsg::JsValue>(js, js.strIntern(cs)); |
| 363 | }); |
| 364 | |
| 365 | if (response.statusCode == 404 || response.statusCode == 410) { |
| 366 | return js.resolvedPromise(KvNamespace::GetWithMetadataResult{ |
| 367 | .value = kj::none, |
| 368 | .metadata = kj::none, |
| 369 | .cacheStatus = kj::mv(cacheStatus), |
| 370 | }); |
| 371 | } |
| 372 | |
| 373 | checkForErrorStatus("GET", response); |
| 374 | |
| 375 | auto metaheader = response.headers->get(context.getHeaderIds().cfKvMetadata); |
| 376 | kj::Maybe<kj::String> maybeMeta; |
| 377 | KJ_IF_SOME(m, metaheader) { |
| 378 | traceContext.setTag("cloudflare.kv.response.metadata"_kjc, true); |
| 379 | maybeMeta = kj::str(m); |
| 380 | } |
| 381 | |
| 382 | auto typeName = |
| 383 | type.map([](const kj::String& s) -> kj::StringPtr { return s; }).orDefault("text"); |
| 384 | |
| 385 | auto& context = IoContext::current(); |
| 386 | auto stream = newSystemStream(response.body.attach(kj::mv(client)), |
| 387 | getContentEncoding( |
| 388 | context, *response.headers, Response::BodyEncoding::AUTO, FeatureFlags::get(js))); |
| 389 | |
| 390 | jsg::Promise<KvNamespace::GetResult> result = nullptr; |
| 391 | |
| 392 | KJ_IF_SOME(size, stream->tryGetLength(StreamEncoding::IDENTITY)) { |
| 393 | traceContext.setTag("cloudflare.kv.response.size"_kjc, static_cast<int64_t>(size)); |
| 394 | } |
| 395 | // This method always returns a single result, but this attribute should be consistent with getBulk |
| 396 | traceContext.setTag("cloudflare.kv.response.returned_rows"_kjc, static_cast<int64_t>(1)); |
| 397 | |
| 398 | if (typeName == "stream") { |
| 399 | result = js.resolvedPromise( |
| 400 | KvNamespace::GetResult(js.alloc<ReadableStream>(context, kj::mv(stream)))); |
| 401 | } else if (typeName == "text") { |
| 402 | // NOTE: In theory we should be using awaitIoLegacy() here since ReadableStreamSource is |
| 403 | // supposed to handle pending events on its own, but we also know that the HTTP client |
| 404 | // backing a KV namespace is never implemented in local JavaScript, so whatever. |
| 405 | result = context.awaitIo(js, |
| 406 | stream->readAllText(context.getLimitEnforcer().getBufferingLimit()) |
| 407 | .attach(kj::mv(stream)), |
| 408 | [](jsg::Lock&, kj::String text) { return KvNamespace::GetResult(kj::mv(text)); }); |
| 409 | } else if (typeName == "arrayBuffer") { |
| 410 | result = context.awaitIo(js, |
| 411 | stream->readAllBytes(context.getLimitEnforcer().getBufferingLimit()) |
| 412 | .attach(kj::mv(stream)), |
| 413 | [](jsg::Lock&, kj::Array<byte> text) { return KvNamespace::GetResult(kj::mv(text)); }); |
| 414 | } else if (typeName == "json") { |
| 415 | result = context.awaitIo(js, |
| 416 | stream->readAllText(context.getLimitEnforcer().getBufferingLimit()) |
| 417 | .attach(kj::mv(stream)), |
| 418 | [](jsg::Lock& js, kj::String text) { |
| 419 | auto ref = jsg::JsRef(js, jsg::JsValue::fromJson(js, text)); |
| 420 | return KvNamespace::GetResult(kj::mv(ref)); |
| 421 | }); |
| 422 | } else { |
| 423 | JSG_FAIL_REQUIRE(TypeError, |
| 424 | "Unknown response type. Possible types are \"text\", \"arrayBuffer\", " |
| 425 | "\"json\", and \"stream\"."); |
| 426 | } |
| 427 | return result.then(js, |
| 428 | [maybeMeta = kj::mv(maybeMeta), cacheStatus = kj::mv(cacheStatus)](jsg::Lock& js, |
| 429 | KvNamespace::GetResult result) mutable -> KvNamespace::GetWithMetadataResult { |
| 430 | kj::Maybe<jsg::JsRef<jsg::JsValue>> meta; |
| 431 | KJ_IF_SOME(metaStr, maybeMeta) { |
| 432 | meta = jsg::JsRef(js, jsg::JsValue::fromJson(js, metaStr)); |
| 433 | } |
| 434 | return KvNamespace::GetWithMetadataResult{ |
| 435 | kj::mv(result), |
| 436 | kj::mv(meta), |
| 437 | kj::mv(cacheStatus), |
| 438 | }; |
| 439 | }); |
| 440 | }); |
| 441 | } |
| 442 | |
| 443 | jsg::Promise<jsg::JsRef<jsg::JsValue>> KvNamespace::list( |
| 444 | jsg::Lock& js, jsg::Optional<ListOptions> options) { |
| 445 | return js.evalNow([&] { |
| 446 | auto& context = IoContext::current(); |
| 447 | TraceContext traceContext = context.makeUserTraceSpan("kv_list"_kjc); |
| 448 | |
| 449 | traceContext.setTag("db.system.name"_kjc, "cloudflare-kv"_kjc); |
| 450 | traceContext.setTag("db.operation.name"_kjc, "list"_kjc); |
| 451 | traceContext.setTag("cloudflare.binding.name"_kjc, bindingName.asPtr()); |
| 452 | traceContext.setTag("cloudflare.binding.type"_kjc, "KV"_kjc); |
| 453 | |
| 454 | kj::Url url; |
| 455 | url.scheme = kj::str("https"); |
| 456 | url.host = kj::str("fake-host"); |
| 457 | KJ_IF_SOME(o, options) { |
| 458 | KJ_IF_SOME(limit, o.limit) { |
| 459 | traceContext.setTag("cloudflare.kv.query.limit"_kjc, static_cast<int64_t>(limit)); |
| 460 | if (limit > 0) { |
| 461 | url.query.add(kj::Url::QueryParam{kj::str("key_count_limit"), kj::str(limit)}); |
| 462 | } |
| 463 | } |
| 464 | KJ_IF_SOME(maybePrefix, o.prefix) { |
| 465 | KJ_IF_SOME(prefix, maybePrefix) { |
| 466 | traceContext.setTag("cloudflare.kv.query.prefix"_kjc, prefix.asPtr()); |
| 467 | url.query.add(kj::Url::QueryParam{kj::str("prefix"), kj::str(prefix)}); |
| 468 | } |
| 469 | } |
| 470 | KJ_IF_SOME(maybeCursor, o.cursor) { |
| 471 | KJ_IF_SOME(cursor, maybeCursor) { |
| 472 | traceContext.setTag("cloudflare.kv.query.cursor"_kjc, cursor.asPtr()); |
| 473 | url.query.add(kj::Url::QueryParam{kj::str("cursor"), kj::str(cursor)}); |
| 474 | } |
| 475 | } |
| 476 | } |
| 477 | |
| 478 | auto urlStr = url.toString(kj::Url::Context::HTTP_PROXY_REQUEST); |
| 479 | |
| 480 | auto headers = kj::HttpHeaders(context.getHeaderTable()); |
| 481 | auto client = |
| 482 | getHttpClient(context, headers, LimitEnforcer::KvOpType::LIST, urlStr, traceContext); |
| 483 | |
| 484 | auto request = client->request(kj::HttpMethod::GET, urlStr, headers); |
| 485 | return context.attachSpans(js, |
| 486 | context.awaitIo(js, kj::mv(request.response), |
| 487 | [&context, client = kj::mv(client), traceContext = kj::mv(traceContext)]( |
| 488 | jsg::Lock& js, kj::HttpClient::Response&& response) mutable |
| 489 | -> jsg::Promise<jsg::JsRef<jsg::JsValue>> { |
| 490 | checkForErrorStatus("GET", response); |
| 491 | |
| 492 | kj::Maybe<jsg::JsRef<jsg::JsValue>> cacheStatus = |
| 493 | [&]() -> kj::Maybe<jsg::JsRef<jsg::JsValue>> { |
| 494 | KJ_IF_SOME(cs, response.headers->get(context.getHeaderIds().cfCacheStatus)) { |
| 495 | traceContext.setTag("cloudflare.kv.response.cache_status"_kjc, cs); |
| 496 | return jsg::JsRef<jsg::JsValue>(js, js.strIntern(cs)); |
| 497 | } |
| 498 | return kj::none; |
| 499 | }(); |
| 500 | |
| 501 | auto stream = newSystemStream(response.body.attach(kj::mv(client)), |
| 502 | getContentEncoding( |
| 503 | context, *response.headers, Response::BodyEncoding::AUTO, FeatureFlags::get(js))); |
| 504 | |
| 505 | KJ_IF_SOME(size, stream->tryGetLength(StreamEncoding::IDENTITY)) { |
| 506 | traceContext.setTag("cloudflare.kv.response.size"_kjc, static_cast<int64_t>(size)); |
| 507 | } |
| 508 | |
| 509 | return context.awaitIo(js, |
| 510 | stream->readAllText(context.getLimitEnforcer().getBufferingLimit()) |
| 511 | .attach(kj::mv(stream)), |
| 512 | [cacheStatus = kj::mv(cacheStatus), traceContext = kj::mv(traceContext)]( |
| 513 | jsg::Lock& js, kj::String text) mutable { |
| 514 | auto result = jsg::JsValue::fromJson(js, text); |
| 515 | parseListMetadata(traceContext, js, result, |
| 516 | cacheStatus.map( |
| 517 | [&](jsg::JsRef<jsg::JsValue>& cs) -> jsg::JsValue { return cs.getHandle(js); })); |
| 518 | return jsg::JsRef(js, result); |
| 519 | }); |
| 520 | }), |
| 521 | kj::mv(traceContext)); |
| 522 | }); |
| 523 | } |
| 524 | |
| 525 | jsg::Promise<void> KvNamespace::put(jsg::Lock& js, |
| 526 | kj::String name, |
| 527 | KvNamespace::PutBody body, |
| 528 | jsg::Optional<PutOptions> options, |
| 529 | const jsg::TypeHandler<KvNamespace::PutSupportedTypes>& putTypeHandler) { |
| 530 | return js.evalNow([&] { |
| 531 | validateKeyName("PUT", name); |
| 532 | |
| 533 | auto& context = IoContext::current(); |
| 534 | TraceContext traceContext = context.makeUserTraceSpan("kv_put"_kjc); |
| 535 | |
| 536 | traceContext.setTag("db.system.name"_kjc, "cloudflare-kv"_kjc); |
| 537 | traceContext.setTag("db.operation.name"_kjc, "put"_kjc); |
| 538 | traceContext.setTag("cloudflare.binding.name"_kjc, bindingName.asPtr()); |
| 539 | traceContext.setTag("cloudflare.binding.type"_kjc, "KV"_kjc); |
| 540 | traceContext.setTag("cloudflare.kv.query.keys"_kjc, name.asPtr()); |
| 541 | traceContext.setTag("cloudflare.kv.query.keys.count"_kjc, static_cast<int64_t>(1)); |
| 542 | |
| 543 | kj::Url url; |
| 544 | url.scheme = kj::str("https"); |
| 545 | url.host = kj::str("fake-host"); |
| 546 | url.path.add(kj::mv(name)); |
| 547 | url.query.add(kj::Url::QueryParam{kj::str("urlencoded"), kj::str("true")}); |
| 548 | |
| 549 | kj::HttpHeaders headers(context.getHeaderTable()); |
| 550 | |
| 551 | // If any optional parameters were specified by the client, append them to |
| 552 | // the URL's query parameters. |
| 553 | KJ_IF_SOME(o, options) { |
| 554 | KJ_IF_SOME(expiration, o.expiration) { |
| 555 | traceContext.setTag("cloudflare.kv.query.expiration"_kjc, static_cast<int64_t>(expiration)); |
| 556 | url.query.add(kj::Url::QueryParam{kj::str("expiration"), kj::str(expiration)}); |
| 557 | } |
| 558 | KJ_IF_SOME(expirationTtl, o.expirationTtl) { |
| 559 | traceContext.setTag( |
| 560 | "cloudflare.kv.query.expiration_ttl"_kjc, static_cast<int64_t>(expirationTtl)); |
| 561 | url.query.add(kj::Url::QueryParam{kj::str("expiration_ttl"), kj::str(expirationTtl)}); |
| 562 | } |
| 563 | KJ_IF_SOME(maybeMetadata, o.metadata) { |
| 564 | KJ_IF_SOME(metadata, maybeMetadata) { |
| 565 | kj::String json = metadata.getHandle(js).toJson(js); |
| 566 | headers.set(context.getHeaderIds().cfKvMetadata, kj::mv(json)); |
| 567 | traceContext.setTag("cloudflare.kv.query.metadata"_kjc, true); |
| 568 | } |
| 569 | } |
| 570 | } |
| 571 | |
| 572 | PutSupportedTypes supportedBody; |
| 573 | |
| 574 | KJ_SWITCH_ONEOF(body) { |
| 575 | KJ_CASE_ONEOF(text, kj::String) { |
| 576 | supportedBody = kj::mv(text); |
| 577 | } |
| 578 | KJ_CASE_ONEOF(object, jsg::JsObject) { |
| 579 | supportedBody = JSG_REQUIRE_NONNULL(putTypeHandler.tryUnwrap(js, object), TypeError, |
| 580 | "KV put() accepts only strings, ArrayBuffers, ArrayBufferViews, and " |
| 581 | "ReadableStreams as values."); |
| 582 | JSG_REQUIRE(!supportedBody.is<kj::String>(), TypeError, |
| 583 | "KV put() accepts only strings, ArrayBuffers, ArrayBufferViews, and " |
| 584 | "ReadableStreams as values."); |
| 585 | // TODO(someday): replace this with logic to do something smarter with Objects |
| 586 | } |
| 587 | } |
| 588 | |
| 589 | kj::Maybe<uint64_t> expectedBodySize; |
| 590 | |
| 591 | KJ_SWITCH_ONEOF(supportedBody) { |
| 592 | KJ_CASE_ONEOF(text, kj::String) { |
| 593 | headers.setPtr(kj::HttpHeaderId::CONTENT_TYPE, MimeType::PLAINTEXT_STRING); |
| 594 | expectedBodySize = static_cast<uint64_t>(text.size()); |
| 595 | traceContext.setTag("cloudflare.kv.query.value_type"_kjc, "text"_kjc); |
| 596 | } |
| 597 | KJ_CASE_ONEOF(data, kj::Array<byte>) { |
| 598 | expectedBodySize = static_cast<uint64_t>(data.size()); |
| 599 | traceContext.setTag("cloudflare.kv.query.value_type"_kjc, "ArrayBuffer"_kjc); |
| 600 | } |
| 601 | KJ_CASE_ONEOF(stream, jsg::Ref<ReadableStream>) { |
| 602 | expectedBodySize = stream->tryGetLength(StreamEncoding::IDENTITY); |
| 603 | traceContext.setTag("cloudflare.kv.query.value_type"_kjc, "ReadableStream"_kjc); |
| 604 | } |
| 605 | } |
| 606 | |
| 607 | KJ_IF_SOME(bodySize, expectedBodySize) { |
| 608 | traceContext.setTag("cloudflare.kv.query.payload.size"_kjc, static_cast<int64_t>(bodySize)); |
| 609 | } |
| 610 | |
| 611 | auto urlStr = url.toString(kj::Url::Context::HTTP_PROXY_REQUEST); |
| 612 | |
| 613 | auto client = |
| 614 | getHttpClient(context, headers, LimitEnforcer::KvOpType::PUT, urlStr, traceContext); |
| 615 | |
| 616 | auto promise = context.waitForOutputLocks().then( |
| 617 | [&context, client = kj::mv(client), urlStr = kj::mv(urlStr), headers = kj::mv(headers), |
| 618 | expectedBodySize, supportedBody = kj::mv(supportedBody)]() mutable { |
| 619 | auto innerReq = client->request(kj::HttpMethod::PUT, urlStr, headers, expectedBodySize); |
| 620 | // TODO(perf): More efficient to explicitly attach rcClient below? |
| 621 | auto req = attachToRequest(kj::mv(innerReq), kj::refcountedWrapper(kj::mv(client))); |
| 622 | |
| 623 | kj::Promise<void> writePromise = nullptr; |
| 624 | KJ_SWITCH_ONEOF(supportedBody) { |
| 625 | KJ_CASE_ONEOF(text, kj::String) { |
| 626 | writePromise = req.body->write(text.asBytes()).attach(kj::mv(text)); |
| 627 | } |
| 628 | KJ_CASE_ONEOF(data, kj::Array<byte>) { |
| 629 | writePromise = req.body->write(data).attach(kj::mv(data)); |
| 630 | } |
| 631 | KJ_CASE_ONEOF(stream, jsg::Ref<ReadableStream>) { |
| 632 | writePromise = context.run( |
| 633 | [dest = newSystemStream(kj::mv(req.body), StreamEncoding::IDENTITY, context), |
| 634 | stream = kj::mv(stream)](jsg::Lock& js) mutable { |
| 635 | return IoContext::current().waitForDeferredProxy( |
| 636 | stream->pumpTo(js, kj::mv(dest), true)); |
| 637 | }); |
| 638 | } |
| 639 | } |
| 640 | |
| 641 | return writePromise.attach(kj::mv(req.body)).then([resp = kj::mv(req.response)]() mutable { |
| 642 | return resp.then([](kj::HttpClient::Response&& response) mutable { |
| 643 | checkForErrorStatus("PUT", response); |
| 644 | |
| 645 | // Read and discard response body, otherwise we might burn the HTTP connection. |
| 646 | return response.body->readAllBytes().attach(kj::mv(response.body)).ignoreResult(); |
| 647 | }); |
| 648 | }); |
| 649 | }); |
| 650 | |
| 651 | return context.attachSpans(js, context.awaitIo(js, kj::mv(promise)), kj::mv(traceContext)); |
| 652 | }); |
| 653 | } |
| 654 | |
| 655 | jsg::Promise<void> KvNamespace::delete_(jsg::Lock& js, kj::String name) { |
| 656 | return js.evalNow([&] { |
| 657 | validateKeyName("DELETE", name); |
| 658 | |
| 659 | auto& context = IoContext::current(); |
| 660 | TraceContext traceContext = context.makeUserTraceSpan("kv_delete"_kjc); |
| 661 | |
| 662 | traceContext.setTag("db.system.name"_kjc, "cloudflare-kv"_kjc); |
| 663 | traceContext.setTag("db.operation.name"_kjc, "delete"_kjc); |
| 664 | traceContext.setTag("cloudflare.binding.name"_kjc, bindingName.asPtr()); |
| 665 | traceContext.setTag("cloudflare.binding.type"_kjc, "KV"_kjc); |
| 666 | traceContext.setTag("cloudflare.kv.query.keys"_kjc, name.asPtr()); |
| 667 | traceContext.setTag("cloudflare.kv.query.keys.count"_kjc, static_cast<int64_t>(1)); |
| 668 | |
| 669 | auto urlStr = kj::str("https://fake-host/", kj::encodeUriComponent(name), "?urlencoded=true"); |
| 670 | |
| 671 | kj::HttpHeaders headers(context.getHeaderTable()); |
| 672 | |
| 673 | auto client = |
| 674 | getHttpClient(context, headers, LimitEnforcer::KvOpType::DELETE, urlStr, traceContext); |
| 675 | |
| 676 | auto promise = context.waitForOutputLocks().then( |
| 677 | [headers = kj::mv(headers), client = kj::mv(client), urlStr = kj::mv(urlStr)]() mutable { |
| 678 | return client->request(kj::HttpMethod::DELETE, urlStr, headers, static_cast<uint64_t>(0)) |
| 679 | .response |
| 680 | .then([](kj::HttpClient::Response&& response) mutable { |
| 681 | checkForErrorStatus("DELETE", response); |
| 682 | }).attach(kj::mv(client)); |
| 683 | }); |
| 684 | |
| 685 | return context.attachSpans(js, context.awaitIo(js, kj::mv(promise)), kj::mv(traceContext)); |
| 686 | }); |
| 687 | } |
| 688 | |
| 689 | jsg::Ref<JsRpcPromise> KvNamespace::deleteBulk(const v8::FunctionCallbackInfo<v8::Value>& args) { |
| 690 | jsg::Lock& js = jsg::Lock::from(args.GetIsolate()); |
| 691 | auto fetcher = js.alloc<Fetcher>(subrequestChannel, Fetcher::RequiresHostAndProtocol::NO, true); |
| 692 | auto method = JSG_REQUIRE_NONNULL( |
| 693 | fetcher->getRpcMethodInternal(js, kj::str("delete"_kj)), Error, "missing delete method"); |
| 694 | return method->call(args); |
| 695 | } |
| 696 | |
| 697 | } // namespace workerd::api |