Skip to content
File

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

28.0 KB
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#include "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 
19namespace workerd::api {
20 
21// As documented in Cloudflare's Worker KV limits.
22static constexpr size_t kMaxKeyLength = 512;
23 
24static 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 
34static 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 
43static 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 
88constexpr auto FLPROD_405_HEADER = "CF-KV-FLPROD-405"_kj;
89 
90kj::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 
114jsg::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 
127jsg::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 
212kj::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 
252kj::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 
277jsg::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 
286kj::OneOf<jsg::Promise<KvNamespace::GetWithMetadataResult>, jsg::Promise<jsg::JsRef<jsg::JsMap>>>
287KvNamespace::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 
312jsg::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 
443jsg::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 
525jsg::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 
655jsg::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 
689jsg::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