File
Blob: src/workerd/api/sync-kv.c++
| 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 | |
| 11 | namespace workerd::api { |
| 12 | |
| 13 | jsg::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 | |
| 33 | jsg::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 | |
| 90 | kj::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 | |
| 103 | void 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 | |
| 116 | kj::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 |