Skip to content
File

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

5.2 KB
1// Copyright (c) 2017-2025 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 "sync-kv.h"
6 
7#include "actor-state.h"
8 
9#include <workerd/util/sqlite-kv.h>
10 
11namespace workerd::api {
12 
13jsg::JsValue SyncKvStorage::get(jsg::Lock& js, kj::String key) {
14 TraceContext traceContext =
15 IoContext::current().makeUserTraceSpan("durable_object_storage_kv_get"_kjc);
16 
17 SqliteKv& sqliteKv = getSqliteKv(js);
18 
19 traceContext.setTag("db.system.name"_kjc, "cloudflare-durable-object-sql"_kjc);
20 traceContext.setTag("db.operation.name"_kjc, "get"_kjc);
21 traceContext.setTag("cloudflare.durable_object.kv.query.keys"_kjc, key.asPtr());
22 traceContext.setTag("cloudflare.durable_object.kv.query.keys.count"_kjc, static_cast<int64_t>(1));
23 
24 kj::Maybe<jsg::JsValue> result;
25 if (sqliteKv.get(key,
26 [&](kj::ArrayPtr<const byte> value) { result = deserializeV8Value(js, key, value); })) {
27 return KJ_ASSERT_NONNULL(result);
28 } else {
29 return js.undefined();
30 }
31}
32 
33jsg::Ref<SyncKvStorage::ListIterator> SyncKvStorage::list(
34 jsg::Lock& js, jsg::Optional<ListOptions> maybeOptions) {
35 TraceContext traceContext =
36 IoContext::current().makeUserTraceSpan("durable_object_storage_kv_list"_kjc);
37 SqliteKv& sqliteKv = getSqliteKv(js);
38 
39 traceContext.setTag("db.system.name"_kjc, "cloudflare-durable-object-sql"_kjc);
40 traceContext.setTag("db.operation.name"_kjc, "list"_kjc);
41 
42 KJ_IF_SOME(o, maybeOptions) {
43 KJ_IF_SOME(start, o.start) {
44 traceContext.setTag("cloudflare.durable_object.kv.query.start"_kjc, start.asPtr());
45 }
46 KJ_IF_SOME(startAfter, o.startAfter) {
47 traceContext.setTag("cloudflare.durable_object.kv.query.startAfter"_kjc, startAfter.asPtr());
48 }
49 KJ_IF_SOME(end, o.end) {
50 traceContext.setTag("cloudflare.durable_object.kv.query.end"_kjc, end.asPtr());
51 }
52 KJ_IF_SOME(prefix, o.prefix) {
53 traceContext.setTag("cloudflare.durable_object.kv.query.prefix"_kjc, prefix.asPtr());
54 }
55 KJ_IF_SOME(reverse, o.reverse) {
56 traceContext.setTag("cloudflare.durable_object.kv.query.reverse"_kjc, reverse);
57 }
58 KJ_IF_SOME(limit, o.limit) {
59 traceContext.setTag(
60 "cloudflare.durable_object.kv.query.limit"_kjc, static_cast<int64_t>(limit));
61 }
62 }
63 
64 // Convert our options to DurableObjectStorageOperations::ListOptions (which also have the
65 // `allowConcurrency` and `noCache` options, which are irrelevant in the sync interface).
66 auto asyncOptions = kj::mv(maybeOptions).map([&](ListOptions&& options) {
67 return DurableObjectStorageOperations::ListOptions{
68 .start = kj::mv(options.start),
69 .startAfter = kj::mv(options.startAfter),
70 .end = kj::mv(options.end),
71 .prefix = kj::mv(options).prefix,
72 .reverse = options.reverse,
73 .limit = options.limit,
74 };
75 });
76 
77 auto [start, end, reverse, limit] =
78 KJ_UNWRAP_OR(DurableObjectStorageOperations::compileListOptions(asyncOptions), {
79 // Key range is empty. Return empty map.
80 return js.alloc<SyncKvStorage::ListIterator>(
81 IoContext::current().addObject(kj::heap<SqliteKv::ListCursor>(nullptr)));
82 });
83 
84 auto cursor = sqliteKv.list(start, end, limit, reverse ? SqliteKv::REVERSE : SqliteKv::FORWARD)
85 .attach(kj::mv(start), kj::mv(end));
86 
87 return js.alloc<SyncKvStorage::ListIterator>(IoContext::current().addObject(kj::mv(cursor)));
88}
89 
90kj::Maybe<jsg::JsArray> SyncKvStorage::listNext(jsg::Lock& js, IoOwn<SqliteKv::ListCursor>& state) {
91 auto& stateRef = *state;
92 KJ_IF_SOME(pair, stateRef.next()) {
93 return js.arr(js.str(pair.key), deserializeV8Value(js, pair.key, pair.value));
94 } else if (stateRef.wasCanceled()) {
95 JSG_FAIL_REQUIRE(Error,
96 "kv.list() iterator was invalidated because a new call to kv.list() was started. Only one "
97 "kv.list() iterator can exist at a time.");
98 } else {
99 return kj::none;
100 }
101}
102 
103void SyncKvStorage::put(jsg::Lock& js, kj::String key, jsg::JsValue value) {
104 TraceContext traceContext =
105 IoContext::current().makeUserTraceSpan("durable_object_storage_kv_put"_kjc);
106 SqliteKv& sqliteKv = getSqliteKv(js);
107 
108 traceContext.setTag("db.system.name"_kjc, "cloudflare-durable-object-sql"_kjc);
109 traceContext.setTag("db.operation.name"_kjc, "put"_kjc);
110 traceContext.setTag("cloudflare.durable_object.kv.query.keys"_kjc, key.asPtr());
111 traceContext.setTag("cloudflare.durable_object.kv.query.keys.count"_kjc, static_cast<int64_t>(1));
112 
113 sqliteKv.put(key, serializeV8Value(js, value));
114}
115 
116kj::OneOf<bool, int> SyncKvStorage::delete_(jsg::Lock& js, kj::String key) {
117 TraceContext traceContext =
118 IoContext::current().makeUserTraceSpan("durable_object_storage_kv_delete"_kjc);
119 SqliteKv& sqliteKv = getSqliteKv(js);
120 
121 traceContext.setTag("db.system.name"_kjc, "cloudflare-durable-object-sql"_kjc);
122 traceContext.setTag("db.operation.name"_kjc, "delete"_kjc);
123 traceContext.setTag("cloudflare.durable_object.kv.query.keys"_kjc, key.asPtr());
124 traceContext.setTag("cloudflare.durable_object.kv.query.keys.count"_kjc, static_cast<int64_t>(1));
125 
126 auto deleted = sqliteKv.delete_(key);
127 
128 traceContext.setTag("cloudflare.durable_object.kv.response.deleted_count"_kjc,
129 static_cast<int64_t>(deleted ? 1 : 0));
130 
131 return deleted;
132}
133 
134} // namespace workerd::api