// Copyright (c) 2017-2022 Cloudflare, Inc. // Licensed under the Apache 2.0 license found in the LICENSE file or at: // https://opensource.org/licenses/Apache-2.0 #include "kv.h" #include "system-streams.h" #include "util.h" #include #include #include #include #include #include #include namespace workerd::api { // As documented in Cloudflare's Worker KV limits. static constexpr size_t kMaxKeyLength = 512; static void checkForErrorStatus(kj::StringPtr method, const kj::HttpClient::Response& response) { if (response.statusCode < 200 || response.statusCode >= 300) { // Manually construct exception so that we can incorporate method and status into the text // that JavaScript sees. kj::throwFatalException(kj::Exception(kj::Exception::Type::FAILED, __FILE__, __LINE__, kj::str(JSG_EXCEPTION(Error) ": KV ", method, " failed: ", response.statusCode, ' ', response.statusText))); } } static void validateKeyName(kj::StringPtr method, kj::StringPtr name) { JSG_REQUIRE(name != "", TypeError, "Key name cannot be empty."); JSG_REQUIRE(name != ".", TypeError, "\".\" is not allowed as a key name."); JSG_REQUIRE(name != "..", TypeError, "\"..\" is not allowed as a key name."); JSG_REQUIRE(name.size() <= kMaxKeyLength, Error, "KV ", method, " failed: ", 414, " UTF-8 encoded length of ", name.size(), " exceeds key length limit of ", kMaxKeyLength, "."); } static void parseListMetadata(TraceContext& traceContext, jsg::Lock& js, jsg::JsValue listResponse, kj::Maybe cacheStatus) { static constexpr auto METADATA = "metadata"_kjc; static constexpr auto KEYS = "keys"_kjc; static constexpr auto CURSOR = "cursor"_kjc; static constexpr auto LIST_COMPLETE = "list_complete"_kjc; static constexpr auto EXPIRATION = "expiration"_kjc; js.withinHandleScope([&] { auto obj = KJ_ASSERT_NONNULL(listResponse.tryCast()); KJ_IF_SOME(boolVal, obj.get(js, LIST_COMPLETE).tryCast()) { traceContext.setTag("cloudflare.kv.response.list_complete"_kjc, boolVal.value(js)); } KJ_IF_SOME(cursor, obj.get(js, CURSOR).tryCast()) { traceContext.setTag("cloudflare.kv.response.cursor"_kjc, kj::str(cursor)); } KJ_IF_SOME(expiration, obj.get(js, EXPIRATION).tryCast()) { KJ_IF_SOME(value, expiration.value(js)) { traceContext.setTag("cloudflare.kv.response.expiration"_kjc, static_cast(value)); } } KJ_IF_SOME(keysArr, obj.get(js, KEYS).tryCast()) { auto length = keysArr.size(); traceContext.setTag("cloudflare.kv.response.returned_rows"_kjc, static_cast(length)); for (int i = 0; i < length; i++) { js.withinHandleScope([&] { KJ_IF_SOME(key, keysArr.get(js, i).tryCast()) { KJ_IF_SOME(str, key.get(js, METADATA).tryCast()) { key.set(js, METADATA, jsg::JsValue::fromJson(js, str)); } } }); } } obj.set(js, "cacheStatus"_kjc, cacheStatus.orDefault(js.null())); }); } constexpr auto FLPROD_405_HEADER = "CF-KV-FLPROD-405"_kj; kj::Own KvNamespace::getHttpClient(IoContext& context, kj::HttpHeaders& headers, kj::OneOf opTypeOrName, kj::StringPtr urlStr, TraceContext& traceContext) { KJ_SWITCH_ONEOF(opTypeOrName) { KJ_CASE_ONEOF(name, kj::LiteralStringConst) {} KJ_CASE_ONEOF(opType, LimitEnforcer::KvOpType) { // Check if we've hit KV usage limits. (This will throw if we have.) context.getLimitEnforcer().newKvRequest(opType); } } auto client = context.getHttpClient(subrequestChannel, true, kj::none, traceContext); headers.addPtrPtr(FLPROD_405_HEADER, urlStr); for (const auto& header: additionalHeaders) { headers.addPtrPtr(header.name.asPtr(), header.value.asPtr()); } return client; } jsg::Promise KvNamespace::getSingle(jsg::Lock& js, IoContext& context, TraceContext& traceContext, kj::String name, jsg::Optional> options) { return js.evalNow([&] { auto resp = getWithMetadataImpl( js, context, traceContext, kj::mv(name), kj::mv(options), LimitEnforcer::KvOpType::GET); return resp.then(js, [](jsg::Lock&, KvNamespace::GetWithMetadataResult result) { return kj::mv(result.value); }); }); } jsg::Promise> KvNamespace::getBulk(jsg::Lock& js, IoContext& context, TraceContext& traceContext, kj::Array name, jsg::Optional> options, bool withMetadata) { return js.evalNow([&] { kj::Url url; url.scheme = kj::str("https"); url.host = kj::str("fake-host"); url.path.add(kj::str("bulk")); url.path.add(kj::str("get")); kj::String body = formBulkBodyString(js, name, withMetadata, options); kj::Maybe expectedBodySize = static_cast(body.size()); auto headers = kj::HttpHeaders(context.getHeaderTable()); headers.set(kj::HttpHeaderId::CONTENT_TYPE, MimeType::JSON.toString()); auto urlStr = url.toString(kj::Url::Context::HTTP_PROXY_REQUEST); // This could be quite large, so let's limit the string length to 512 characters auto keysStr = kj::strArray(name, ", "); if (keysStr.size() > 512) { keysStr = kj::str(keysStr.slice(0, 509), "..."); } traceContext.setTag("cloudflare.kv.query.keys"_kjc, kj::mv(keysStr)); traceContext.setTag("cloudflare.kv.query.keys.count"_kjc, static_cast(name.size())); KJ_IF_SOME(_options, options) { KJ_SWITCH_ONEOF(_options) { KJ_CASE_ONEOF(type, kj::String) { traceContext.setTag("cloudflare.kv.query.type"_kjc, kj::mv(type)); } KJ_CASE_ONEOF(o, GetOptions) { KJ_IF_SOME(type, o.type) { traceContext.setTag("cloudflare.kv.query.type"_kjc, kj::mv(type)); } KJ_IF_SOME(cacheTtl, o.cacheTtl) { traceContext.setTag( "cloudflare.kv.query.cache_ttl"_kjc, static_cast(cacheTtl)); } } } } auto client = getHttpClient(context, headers, LimitEnforcer::KvOpType::GET_BULK, urlStr, traceContext); auto promise = context.waitForOutputLocks().then( [client = kj::mv(client), urlStr = kj::mv(urlStr), headers = kj::mv(headers), expectedBodySize, supportedBody = kj::mv(body)]() mutable { auto innerReq = client->request(kj::HttpMethod::POST, urlStr, headers, expectedBodySize); auto req = attachToRequest(kj::mv(innerReq), kj::refcountedWrapper(kj::mv(client))); kj::Promise writePromise = nullptr; writePromise = req.body->write(supportedBody.asBytes()).attach(kj::mv(supportedBody)); return writePromise.attach(kj::mv(req.body)).then([resp = kj::mv(req.response)]() mutable { return resp.then([](kj::HttpClient::Response&& response) mutable { checkForErrorStatus("GET_BULK", response); return response.body->readAllText().attach(kj::mv(response.body)); }); }); }); return context.awaitIo(js, kj::mv(promise), [&, traceContext = kj::mv(traceContext)](jsg::Lock& js, kj::String text) mutable { traceContext.setTag("cloudflare.kv.response.size"_kjc, static_cast(text.size())); auto result = jsg::JsValue::fromJson(js, text); auto map = js.map(); KJ_IF_SOME(obj, result.tryCast()) { auto values = obj.getPropertyNames(js, jsg::KeyCollectionFilter::OWN_ONLY, jsg::PropertyFilter::SKIP_SYMBOLS, jsg::IndexFilter::SKIP_INDICES); for (int i = 0; i < values.size(); i++) { auto key = values.get(js, i); map.set(js, kj::mv(key), obj.get(js, key)); } traceContext.setTag( "cloudflare.kv.response.returned_rows"_kjc, static_cast(values.size())); } return jsg::JsRef(js, map); }); }); } kj::String KvNamespace::formBulkBodyString(jsg::Lock& js, kj::Array& names, bool withMetadata, jsg::Optional>& options) { kj::String type = kj::str(""); kj::String cacheTtlStr = kj::str(""); KJ_IF_SOME(oneOfOptions, options) { KJ_SWITCH_ONEOF(oneOfOptions) { KJ_CASE_ONEOF(t, kj::String) { type = kj::str(t); } KJ_CASE_ONEOF(options, GetOptions) { KJ_IF_SOME(t, options.type) { type = kj::str(t); } KJ_IF_SOME(cacheTtl, options.cacheTtl) { cacheTtlStr = kj::str(cacheTtl); } } } } auto object = js.obj(); auto keysArray = js.arr(names.asPtr(), [](jsg::Lock& js, const kj::String& val) { return js.str(val); }); object.set(js, "keys", keysArray); if (type != kj::str("")) { object.set(js, "type", js.str(type)); } if (withMetadata) { object.set(js, "withMetadata", js.boolean(true)); } if (cacheTtlStr != kj::str("")) { object.set(js, "cacheTtl", js.str(cacheTtlStr)); } return jsg::JsValue(object).toJson(js); } kj::OneOf, jsg::Promise>> KvNamespace:: get(jsg::Lock& js, kj::OneOf> name, jsg::Optional> options) { auto& context = IoContext::current(); TraceContext traceContext = context.makeUserTraceSpan("kv_get"_kjc); traceContext.setTag("db.system.name"_kjc, "cloudflare-kv"_kjc); traceContext.setTag("db.operation.name"_kjc, "get"_kjc); traceContext.setTag("cloudflare.binding.name"_kjc, bindingName.asPtr()); traceContext.setTag("cloudflare.binding.type"_kjc, "KV"_kjc); KJ_SWITCH_ONEOF(name) { KJ_CASE_ONEOF(arr, kj::Array) { return context.attachSpans(js, getBulk(js, context, traceContext, kj::mv(arr), kj::mv(options), false), kj::mv(traceContext)); } KJ_CASE_ONEOF(str, kj::String) { return context.attachSpans(js, getSingle(js, context, traceContext, kj::mv(str), kj::mv(options)), kj::mv(traceContext)); } } KJ_UNREACHABLE; }; jsg::Promise KvNamespace::getWithMetadataSingle(jsg::Lock& js, IoContext& context, TraceContext& traceContext, kj::String name, jsg::Optional> options) { return getWithMetadataImpl( js, context, traceContext, kj::mv(name), kj::mv(options), LimitEnforcer::KvOpType::GET_WITH); } kj::OneOf, jsg::Promise>> KvNamespace::getWithMetadata(jsg::Lock& js, kj::OneOf, kj::String> name, jsg::Optional> options) { auto& context = IoContext::current(); TraceContext traceContext = context.makeUserTraceSpan("kv_getWithMetadata"_kjc); traceContext.setTag("db.system.name"_kjc, "cloudflare-kv"_kjc); traceContext.setTag("db.operation.name"_kjc, "get"_kjc); traceContext.setTag("cloudflare.binding.name"_kjc, bindingName.asPtr()); traceContext.setTag("cloudflare.binding.type"_kjc, "KV"_kjc); KJ_SWITCH_ONEOF(name) { KJ_CASE_ONEOF(arr, kj::Array) { return context.attachSpans(js, getBulk(js, context, traceContext, kj::mv(arr), kj::mv(options), true), kj::mv(traceContext)); } KJ_CASE_ONEOF(str, kj::String) { return context.attachSpans(js, getWithMetadataSingle(js, context, traceContext, kj::mv(str), kj::mv(options)), kj::mv(traceContext)); } } KJ_UNREACHABLE; } jsg::Promise KvNamespace::getWithMetadataImpl(jsg::Lock& js, IoContext& context, TraceContext& traceContext, kj::String name, jsg::Optional> options, LimitEnforcer::KvOpType op) { validateKeyName("GET", name); traceContext.setTag("cloudflare.kv.query.keys"_kjc, name.asPtr()); traceContext.setTag("cloudflare.kv.query.keys.count"_kjc, static_cast(1)); kj::Url url; url.scheme = kj::str("https"); url.host = kj::str("fake-host"); url.path.add(kj::mv(name)); url.query.add(kj::Url::QueryParam{kj::str("urlencoded"), kj::str("true")}); kj::Maybe type; KJ_IF_SOME(oneOfOptions, options) { KJ_SWITCH_ONEOF(oneOfOptions) { KJ_CASE_ONEOF(t, kj::String) { type = kj::str(t); traceContext.setTag("cloudflare.kv.query.type"_kjc, kj::mv(t)); } KJ_CASE_ONEOF(options, GetOptions) { KJ_IF_SOME(t, options.type) { type = kj::str(t); traceContext.setTag("cloudflare.kv.query.type"_kjc, kj::mv(t)); } KJ_IF_SOME(cacheTtl, options.cacheTtl) { url.query.add(kj::Url::QueryParam{kj::str("cache_ttl"), kj::str(cacheTtl)}); traceContext.setTag("cloudflare.kv.query.cache_ttl"_kjc, static_cast(cacheTtl)); } } } } auto urlStr = url.toString(kj::Url::Context::HTTP_PROXY_REQUEST); auto headers = kj::HttpHeaders(context.getHeaderTable()); auto client = getHttpClient(context, headers, op, urlStr, traceContext); auto request = client->request(kj::HttpMethod::GET, urlStr, headers); return context.awaitIo(js, kj::mv(request.response), [type = kj::mv(type), &context, client = kj::mv(client), traceContext = kj::mv(traceContext)]( jsg::Lock& js, kj::HttpClient::Response&& response) mutable -> jsg::Promise { auto cacheStatus = response.headers->get(context.getHeaderIds().cfCacheStatus).map([&](kj::StringPtr cs) { traceContext.setTag("cloudflare.kv.response.cache_status"_kjc, cs); return jsg::JsRef(js, js.strIntern(cs)); }); if (response.statusCode == 404 || response.statusCode == 410) { return js.resolvedPromise(KvNamespace::GetWithMetadataResult{ .value = kj::none, .metadata = kj::none, .cacheStatus = kj::mv(cacheStatus), }); } checkForErrorStatus("GET", response); auto metaheader = response.headers->get(context.getHeaderIds().cfKvMetadata); kj::Maybe maybeMeta; KJ_IF_SOME(m, metaheader) { traceContext.setTag("cloudflare.kv.response.metadata"_kjc, true); maybeMeta = kj::str(m); } auto typeName = type.map([](const kj::String& s) -> kj::StringPtr { return s; }).orDefault("text"); auto& context = IoContext::current(); auto stream = newSystemStream(response.body.attach(kj::mv(client)), getContentEncoding( context, *response.headers, Response::BodyEncoding::AUTO, FeatureFlags::get(js))); jsg::Promise result = nullptr; KJ_IF_SOME(size, stream->tryGetLength(StreamEncoding::IDENTITY)) { traceContext.setTag("cloudflare.kv.response.size"_kjc, static_cast(size)); } // This method always returns a single result, but this attribute should be consistent with getBulk traceContext.setTag("cloudflare.kv.response.returned_rows"_kjc, static_cast(1)); if (typeName == "stream") { result = js.resolvedPromise( KvNamespace::GetResult(js.alloc(context, kj::mv(stream)))); } else if (typeName == "text") { // NOTE: In theory we should be using awaitIoLegacy() here since ReadableStreamSource is // supposed to handle pending events on its own, but we also know that the HTTP client // backing a KV namespace is never implemented in local JavaScript, so whatever. result = context.awaitIo(js, stream->readAllText(context.getLimitEnforcer().getBufferingLimit()) .attach(kj::mv(stream)), [](jsg::Lock&, kj::String text) { return KvNamespace::GetResult(kj::mv(text)); }); } else if (typeName == "arrayBuffer") { result = context.awaitIo(js, stream->readAllBytes(context.getLimitEnforcer().getBufferingLimit()) .attach(kj::mv(stream)), [](jsg::Lock&, kj::Array text) { return KvNamespace::GetResult(kj::mv(text)); }); } else if (typeName == "json") { result = context.awaitIo(js, stream->readAllText(context.getLimitEnforcer().getBufferingLimit()) .attach(kj::mv(stream)), [](jsg::Lock& js, kj::String text) { auto ref = jsg::JsRef(js, jsg::JsValue::fromJson(js, text)); return KvNamespace::GetResult(kj::mv(ref)); }); } else { JSG_FAIL_REQUIRE(TypeError, "Unknown response type. Possible types are \"text\", \"arrayBuffer\", " "\"json\", and \"stream\"."); } return result.then(js, [maybeMeta = kj::mv(maybeMeta), cacheStatus = kj::mv(cacheStatus)](jsg::Lock& js, KvNamespace::GetResult result) mutable -> KvNamespace::GetWithMetadataResult { kj::Maybe> meta; KJ_IF_SOME(metaStr, maybeMeta) { meta = jsg::JsRef(js, jsg::JsValue::fromJson(js, metaStr)); } return KvNamespace::GetWithMetadataResult{ kj::mv(result), kj::mv(meta), kj::mv(cacheStatus), }; }); }); } jsg::Promise> KvNamespace::list( jsg::Lock& js, jsg::Optional options) { return js.evalNow([&] { auto& context = IoContext::current(); TraceContext traceContext = context.makeUserTraceSpan("kv_list"_kjc); traceContext.setTag("db.system.name"_kjc, "cloudflare-kv"_kjc); traceContext.setTag("db.operation.name"_kjc, "list"_kjc); traceContext.setTag("cloudflare.binding.name"_kjc, bindingName.asPtr()); traceContext.setTag("cloudflare.binding.type"_kjc, "KV"_kjc); kj::Url url; url.scheme = kj::str("https"); url.host = kj::str("fake-host"); KJ_IF_SOME(o, options) { KJ_IF_SOME(limit, o.limit) { traceContext.setTag("cloudflare.kv.query.limit"_kjc, static_cast(limit)); if (limit > 0) { url.query.add(kj::Url::QueryParam{kj::str("key_count_limit"), kj::str(limit)}); } } KJ_IF_SOME(maybePrefix, o.prefix) { KJ_IF_SOME(prefix, maybePrefix) { traceContext.setTag("cloudflare.kv.query.prefix"_kjc, prefix.asPtr()); url.query.add(kj::Url::QueryParam{kj::str("prefix"), kj::str(prefix)}); } } KJ_IF_SOME(maybeCursor, o.cursor) { KJ_IF_SOME(cursor, maybeCursor) { traceContext.setTag("cloudflare.kv.query.cursor"_kjc, cursor.asPtr()); url.query.add(kj::Url::QueryParam{kj::str("cursor"), kj::str(cursor)}); } } } auto urlStr = url.toString(kj::Url::Context::HTTP_PROXY_REQUEST); auto headers = kj::HttpHeaders(context.getHeaderTable()); auto client = getHttpClient(context, headers, LimitEnforcer::KvOpType::LIST, urlStr, traceContext); auto request = client->request(kj::HttpMethod::GET, urlStr, headers); return context.attachSpans(js, context.awaitIo(js, kj::mv(request.response), [&context, client = kj::mv(client), traceContext = kj::mv(traceContext)]( jsg::Lock& js, kj::HttpClient::Response&& response) mutable -> jsg::Promise> { checkForErrorStatus("GET", response); kj::Maybe> cacheStatus = [&]() -> kj::Maybe> { KJ_IF_SOME(cs, response.headers->get(context.getHeaderIds().cfCacheStatus)) { traceContext.setTag("cloudflare.kv.response.cache_status"_kjc, cs); return jsg::JsRef(js, js.strIntern(cs)); } return kj::none; }(); auto stream = newSystemStream(response.body.attach(kj::mv(client)), getContentEncoding( context, *response.headers, Response::BodyEncoding::AUTO, FeatureFlags::get(js))); KJ_IF_SOME(size, stream->tryGetLength(StreamEncoding::IDENTITY)) { traceContext.setTag("cloudflare.kv.response.size"_kjc, static_cast(size)); } return context.awaitIo(js, stream->readAllText(context.getLimitEnforcer().getBufferingLimit()) .attach(kj::mv(stream)), [cacheStatus = kj::mv(cacheStatus), traceContext = kj::mv(traceContext)]( jsg::Lock& js, kj::String text) mutable { auto result = jsg::JsValue::fromJson(js, text); parseListMetadata(traceContext, js, result, cacheStatus.map( [&](jsg::JsRef& cs) -> jsg::JsValue { return cs.getHandle(js); })); return jsg::JsRef(js, result); }); }), kj::mv(traceContext)); }); } jsg::Promise KvNamespace::put(jsg::Lock& js, kj::String name, KvNamespace::PutBody body, jsg::Optional options, const jsg::TypeHandler& putTypeHandler) { return js.evalNow([&] { validateKeyName("PUT", name); auto& context = IoContext::current(); TraceContext traceContext = context.makeUserTraceSpan("kv_put"_kjc); traceContext.setTag("db.system.name"_kjc, "cloudflare-kv"_kjc); traceContext.setTag("db.operation.name"_kjc, "put"_kjc); traceContext.setTag("cloudflare.binding.name"_kjc, bindingName.asPtr()); traceContext.setTag("cloudflare.binding.type"_kjc, "KV"_kjc); traceContext.setTag("cloudflare.kv.query.keys"_kjc, name.asPtr()); traceContext.setTag("cloudflare.kv.query.keys.count"_kjc, static_cast(1)); kj::Url url; url.scheme = kj::str("https"); url.host = kj::str("fake-host"); url.path.add(kj::mv(name)); url.query.add(kj::Url::QueryParam{kj::str("urlencoded"), kj::str("true")}); kj::HttpHeaders headers(context.getHeaderTable()); // If any optional parameters were specified by the client, append them to // the URL's query parameters. KJ_IF_SOME(o, options) { KJ_IF_SOME(expiration, o.expiration) { traceContext.setTag("cloudflare.kv.query.expiration"_kjc, static_cast(expiration)); url.query.add(kj::Url::QueryParam{kj::str("expiration"), kj::str(expiration)}); } KJ_IF_SOME(expirationTtl, o.expirationTtl) { traceContext.setTag( "cloudflare.kv.query.expiration_ttl"_kjc, static_cast(expirationTtl)); url.query.add(kj::Url::QueryParam{kj::str("expiration_ttl"), kj::str(expirationTtl)}); } KJ_IF_SOME(maybeMetadata, o.metadata) { KJ_IF_SOME(metadata, maybeMetadata) { kj::String json = metadata.getHandle(js).toJson(js); headers.set(context.getHeaderIds().cfKvMetadata, kj::mv(json)); traceContext.setTag("cloudflare.kv.query.metadata"_kjc, true); } } } PutSupportedTypes supportedBody; KJ_SWITCH_ONEOF(body) { KJ_CASE_ONEOF(text, kj::String) { supportedBody = kj::mv(text); } KJ_CASE_ONEOF(object, jsg::JsObject) { supportedBody = JSG_REQUIRE_NONNULL(putTypeHandler.tryUnwrap(js, object), TypeError, "KV put() accepts only strings, ArrayBuffers, ArrayBufferViews, and " "ReadableStreams as values."); JSG_REQUIRE(!supportedBody.is(), TypeError, "KV put() accepts only strings, ArrayBuffers, ArrayBufferViews, and " "ReadableStreams as values."); // TODO(someday): replace this with logic to do something smarter with Objects } } kj::Maybe expectedBodySize; KJ_SWITCH_ONEOF(supportedBody) { KJ_CASE_ONEOF(text, kj::String) { headers.setPtr(kj::HttpHeaderId::CONTENT_TYPE, MimeType::PLAINTEXT_STRING); expectedBodySize = static_cast(text.size()); traceContext.setTag("cloudflare.kv.query.value_type"_kjc, "text"_kjc); } KJ_CASE_ONEOF(data, kj::Array) { expectedBodySize = static_cast(data.size()); traceContext.setTag("cloudflare.kv.query.value_type"_kjc, "ArrayBuffer"_kjc); } KJ_CASE_ONEOF(stream, jsg::Ref) { expectedBodySize = stream->tryGetLength(StreamEncoding::IDENTITY); traceContext.setTag("cloudflare.kv.query.value_type"_kjc, "ReadableStream"_kjc); } } KJ_IF_SOME(bodySize, expectedBodySize) { traceContext.setTag("cloudflare.kv.query.payload.size"_kjc, static_cast(bodySize)); } auto urlStr = url.toString(kj::Url::Context::HTTP_PROXY_REQUEST); auto client = getHttpClient(context, headers, LimitEnforcer::KvOpType::PUT, urlStr, traceContext); auto promise = context.waitForOutputLocks().then( [&context, client = kj::mv(client), urlStr = kj::mv(urlStr), headers = kj::mv(headers), expectedBodySize, supportedBody = kj::mv(supportedBody)]() mutable { auto innerReq = client->request(kj::HttpMethod::PUT, urlStr, headers, expectedBodySize); // TODO(perf): More efficient to explicitly attach rcClient below? auto req = attachToRequest(kj::mv(innerReq), kj::refcountedWrapper(kj::mv(client))); kj::Promise writePromise = nullptr; KJ_SWITCH_ONEOF(supportedBody) { KJ_CASE_ONEOF(text, kj::String) { writePromise = req.body->write(text.asBytes()).attach(kj::mv(text)); } KJ_CASE_ONEOF(data, kj::Array) { writePromise = req.body->write(data).attach(kj::mv(data)); } KJ_CASE_ONEOF(stream, jsg::Ref) { writePromise = context.run( [dest = newSystemStream(kj::mv(req.body), StreamEncoding::IDENTITY, context), stream = kj::mv(stream)](jsg::Lock& js) mutable { return IoContext::current().waitForDeferredProxy( stream->pumpTo(js, kj::mv(dest), true)); }); } } return writePromise.attach(kj::mv(req.body)).then([resp = kj::mv(req.response)]() mutable { return resp.then([](kj::HttpClient::Response&& response) mutable { checkForErrorStatus("PUT", response); // Read and discard response body, otherwise we might burn the HTTP connection. return response.body->readAllBytes().attach(kj::mv(response.body)).ignoreResult(); }); }); }); return context.attachSpans(js, context.awaitIo(js, kj::mv(promise)), kj::mv(traceContext)); }); } jsg::Promise KvNamespace::delete_(jsg::Lock& js, kj::String name) { return js.evalNow([&] { validateKeyName("DELETE", name); auto& context = IoContext::current(); TraceContext traceContext = context.makeUserTraceSpan("kv_delete"_kjc); traceContext.setTag("db.system.name"_kjc, "cloudflare-kv"_kjc); traceContext.setTag("db.operation.name"_kjc, "delete"_kjc); traceContext.setTag("cloudflare.binding.name"_kjc, bindingName.asPtr()); traceContext.setTag("cloudflare.binding.type"_kjc, "KV"_kjc); traceContext.setTag("cloudflare.kv.query.keys"_kjc, name.asPtr()); traceContext.setTag("cloudflare.kv.query.keys.count"_kjc, static_cast(1)); auto urlStr = kj::str("https://fake-host/", kj::encodeUriComponent(name), "?urlencoded=true"); kj::HttpHeaders headers(context.getHeaderTable()); auto client = getHttpClient(context, headers, LimitEnforcer::KvOpType::DELETE, urlStr, traceContext); auto promise = context.waitForOutputLocks().then( [headers = kj::mv(headers), client = kj::mv(client), urlStr = kj::mv(urlStr)]() mutable { return client->request(kj::HttpMethod::DELETE, urlStr, headers, static_cast(0)) .response .then([](kj::HttpClient::Response&& response) mutable { checkForErrorStatus("DELETE", response); }).attach(kj::mv(client)); }); return context.attachSpans(js, context.awaitIo(js, kj::mv(promise)), kj::mv(traceContext)); }); } jsg::Ref KvNamespace::deleteBulk(const v8::FunctionCallbackInfo& args) { jsg::Lock& js = jsg::Lock::from(args.GetIsolate()); auto fetcher = js.alloc(subrequestChannel, Fetcher::RequiresHostAndProtocol::NO, true); auto method = JSG_REQUIRE_NONNULL( fetcher->getRpcMethodInternal(js, kj::str("delete"_kj)), Error, "missing delete method"); return method->call(args); } } // namespace workerd::api