File
Blob: src/workerd/api/actor-state.h
| 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 | #pragma once |
| 6 | // APIs that an Actor (Durable Object) uses to access its own state. |
| 7 | // |
| 8 | // See actor.h for APIs used by other Workers to talk to Actors. |
| 9 | |
| 10 | #include <workerd/api/actor.h> |
| 11 | #include <workerd/api/container.h> |
| 12 | #include <workerd/io/actor-cache.h> |
| 13 | #include <workerd/io/actor-id.h> |
| 14 | #include <workerd/io/compatibility-date.capnp.h> |
| 15 | #include <workerd/io/io-own.h> |
| 16 | #include <workerd/io/worker.h> |
| 17 | #include <workerd/jsg/jsg.h> |
| 18 | |
| 19 | #include <kj/async.h> |
| 20 | |
| 21 | namespace workerd::api { |
| 22 | class SqlStorage; |
| 23 | class SyncKvStorage; |
| 24 | |
| 25 | // Forward-declared to avoid dependency cycle (actor.h -> http.h -> basics.h -> actor-state.h) |
| 26 | class DurableObject; |
| 27 | class DurableObjectId; |
| 28 | class WebSocket; |
| 29 | class DurableObjectClass; |
| 30 | class LoopbackDurableObjectNamespace; |
| 31 | class LoopbackColoLocalActorNamespace; |
| 32 | |
| 33 | kj::Array<kj::byte> serializeV8Value(jsg::Lock& js, const jsg::JsValue& value); |
| 34 | |
| 35 | jsg::JsValue deserializeV8Value( |
| 36 | jsg::Lock& js, kj::ArrayPtr<const char> key, kj::ArrayPtr<const kj::byte> buf); |
| 37 | |
| 38 | // Common implementation of DurableObjectStorage and DurableObjectTransaction. This class is |
| 39 | // designed to be used as a mixin. |
| 40 | class DurableObjectStorageOperations { |
| 41 | public: |
| 42 | struct GetOptions { |
| 43 | jsg::Optional<bool> allowConcurrency; |
| 44 | jsg::Optional<bool> noCache; |
| 45 | |
| 46 | inline operator ActorCacheOps::ReadOptions() const { |
| 47 | return {.noCache = noCache.orDefault(false)}; |
| 48 | } |
| 49 | |
| 50 | JSG_STRUCT(allowConcurrency, noCache); |
| 51 | JSG_STRUCT_TS_OVERRIDE(DurableObjectGetOptions); // Rename from DurableObjectStorageOperationsGetOptions |
| 52 | }; |
| 53 | |
| 54 | jsg::Promise<jsg::JsRef<jsg::JsValue>> get(jsg::Lock& js, |
| 55 | kj::OneOf<kj::String, kj::Array<kj::String>> keys, |
| 56 | jsg::Optional<GetOptions> options); |
| 57 | |
| 58 | struct GetAlarmOptions { |
| 59 | jsg::Optional<bool> allowConcurrency; |
| 60 | |
| 61 | JSG_STRUCT(allowConcurrency); |
| 62 | JSG_STRUCT_TS_OVERRIDE(DurableObjectGetAlarmOptions); // Rename from DurableObjectStorageOperationsGetAlarmOptions |
| 63 | }; |
| 64 | |
| 65 | jsg::Promise<kj::Maybe<double>> getAlarm(jsg::Lock& js, jsg::Optional<GetAlarmOptions> options); |
| 66 | |
| 67 | struct ListOptions { |
| 68 | jsg::Optional<kj::String> start; |
| 69 | jsg::Optional<kj::String> startAfter; |
| 70 | jsg::Optional<kj::String> end; |
| 71 | jsg::Optional<kj::String> prefix; |
| 72 | jsg::Optional<bool> reverse; |
| 73 | jsg::Optional<int> limit; |
| 74 | |
| 75 | jsg::Optional<bool> allowConcurrency; |
| 76 | jsg::Optional<bool> noCache; |
| 77 | |
| 78 | inline operator ActorCacheOps::ReadOptions() const { |
| 79 | return {.noCache = noCache.orDefault(false)}; |
| 80 | } |
| 81 | |
| 82 | JSG_STRUCT(start, startAfter, end, prefix, reverse, limit, allowConcurrency, noCache); |
| 83 | JSG_STRUCT_TS_OVERRIDE(DurableObjectListOptions); // Rename from DurableObjectStorageOperationsListOptions |
| 84 | }; |
| 85 | |
| 86 | // A more convenient form of `ListOptions` for actually implementing the operation -- but less |
| 87 | // convenient for specifying it. |
| 88 | struct CompiledListOptions { |
| 89 | kj::String start; |
| 90 | kj::Maybe<kj::String> end; |
| 91 | bool reverse; |
| 92 | kj::Maybe<uint> limit; |
| 93 | }; |
| 94 | |
| 95 | // Compile `ListOptions` into `CompiledListOptions`. Returns null if the list operation would |
| 96 | // provably return no results (e.g. the end key is before the start key). This may (or may not) |
| 97 | // move some of the strings from the input to the output. |
| 98 | // |
| 99 | // This is public so that SyncKvStorage can reuse it. |
| 100 | static kj::Maybe<CompiledListOptions> compileListOptions(kj::Maybe<ListOptions>& maybeOptions); |
| 101 | |
| 102 | jsg::Promise<jsg::JsRef<jsg::JsValue>> list(jsg::Lock& js, jsg::Optional<ListOptions> options); |
| 103 | |
| 104 | struct PutOptions { |
| 105 | jsg::Optional<bool> allowConcurrency; |
| 106 | jsg::Optional<bool> allowUnconfirmed; |
| 107 | jsg::Optional<bool> noCache; |
| 108 | |
| 109 | inline operator ActorCacheOps::WriteOptions() const { |
| 110 | return { |
| 111 | .allowUnconfirmed = allowUnconfirmed.orDefault(false), .noCache = noCache.orDefault(false)}; |
| 112 | } |
| 113 | |
| 114 | JSG_STRUCT(allowConcurrency, allowUnconfirmed, noCache); |
| 115 | JSG_STRUCT_TS_OVERRIDE(DurableObjectPutOptions); // Rename from DurableObjectStorageOperationsPutOptions |
| 116 | }; |
| 117 | |
| 118 | jsg::Promise<void> put(jsg::Lock& js, |
| 119 | kj::OneOf<kj::String, jsg::Dict<jsg::JsValue>> keyOrEntries, |
| 120 | jsg::Optional<jsg::JsValue> value, |
| 121 | jsg::Optional<PutOptions> options, |
| 122 | const jsg::TypeHandler<PutOptions>& optionsTypeHandler); |
| 123 | |
| 124 | kj::OneOf<jsg::Promise<bool>, jsg::Promise<int>> delete_(jsg::Lock& js, |
| 125 | kj::OneOf<kj::String, kj::Array<kj::String>> keys, |
| 126 | jsg::Optional<PutOptions> options); |
| 127 | |
| 128 | struct SetAlarmOptions { |
| 129 | jsg::Optional<bool> allowConcurrency; |
| 130 | jsg::Optional<bool> allowUnconfirmed; |
| 131 | // We don't allow noCache for alarm puts. |
| 132 | |
| 133 | inline operator ActorCacheOps::WriteOptions() const { |
| 134 | return { |
| 135 | .allowUnconfirmed = allowUnconfirmed.orDefault(false), |
| 136 | }; |
| 137 | } |
| 138 | |
| 139 | JSG_STRUCT(allowConcurrency, allowUnconfirmed); |
| 140 | JSG_STRUCT_TS_OVERRIDE(DurableObjectSetAlarmOptions); // Rename from DurableObjectStorageOperationsSetAlarmOptions |
| 141 | }; |
| 142 | |
| 143 | jsg::Promise<void> setAlarm( |
| 144 | jsg::Lock& js, kj::Date scheduledTime, jsg::Optional<SetAlarmOptions> options); |
| 145 | jsg::Promise<void> deleteAlarm(jsg::Lock& js, jsg::Optional<SetAlarmOptions> options); |
| 146 | |
| 147 | protected: |
| 148 | using OpName = kj::StringPtr; |
| 149 | static constexpr OpName OP_GET = "get()"_kj; |
| 150 | static constexpr OpName OP_GET_ALARM = "getAlarm()"_kj; |
| 151 | static constexpr OpName OP_LIST = "list()"_kj; |
| 152 | static constexpr OpName OP_PUT = "put()"_kj; |
| 153 | static constexpr OpName OP_PUT_ALARM = "setAlarm()"_kj; |
| 154 | static constexpr OpName OP_DELETE = "delete()"_kj; |
| 155 | static constexpr OpName OP_DELETE_ALARM = "deleteAlarm()"_kj; |
| 156 | static constexpr OpName OP_RENAME = "rename()"_kj; |
| 157 | static constexpr OpName OP_ROLLBACK = "rollback()"_kj; |
| 158 | |
| 159 | static bool readOnlyOp(OpName op) { |
| 160 | return op == OP_GET || op == OP_LIST || op == OP_ROLLBACK; |
| 161 | } |
| 162 | |
| 163 | virtual ActorCacheOps& getCache(OpName op) = 0; |
| 164 | |
| 165 | // Whether to skip caching and allow concurrency on all operations. |
| 166 | virtual bool useDirectIo() = 0; |
| 167 | |
| 168 | // Method that should be called at the start of each storage operation to override any of the |
| 169 | // options as appropriate. |
| 170 | template <typename T> |
| 171 | T configureOptions(T&& options) { |
| 172 | if (useDirectIo()) { |
| 173 | options.allowConcurrency = true; |
| 174 | options.noCache = true; |
| 175 | } |
| 176 | return kj::mv(options); |
| 177 | } |
| 178 | |
| 179 | private: |
| 180 | jsg::Promise<jsg::JsRef<jsg::JsValue>> getOne( |
| 181 | jsg::Lock& js, kj::String key, const GetOptions& options); |
| 182 | jsg::Promise<jsg::JsRef<jsg::JsValue>> getMultiple( |
| 183 | jsg::Lock& js, kj::Array<kj::String> keys, const GetOptions& options); |
| 184 | |
| 185 | jsg::Promise<void> putOne( |
| 186 | jsg::Lock& js, kj::String key, jsg::JsValue value, const PutOptions& options); |
| 187 | jsg::Promise<void> putMultiple( |
| 188 | jsg::Lock& js, jsg::Dict<jsg::JsValue> entries, const PutOptions& options); |
| 189 | |
| 190 | jsg::Promise<bool> deleteOne(jsg::Lock& js, kj::String key, const PutOptions& options); |
| 191 | jsg::Promise<int> deleteMultiple( |
| 192 | jsg::Lock& js, kj::Array<kj::String> keys, const PutOptions& options); |
| 193 | }; |
| 194 | |
| 195 | class DurableObjectTransaction; |
| 196 | |
| 197 | class DurableObjectStorage: public jsg::Object, public DurableObjectStorageOperations { |
| 198 | public: |
| 199 | DurableObjectStorage(jsg::Lock&, IoPtr<ActorCacheInterface> cache, bool enableSql) |
| 200 | : cache(kj::mv(cache)), |
| 201 | enableSql(enableSql) {} |
| 202 | |
| 203 | // This constructor is only used when we're setting up the `DurableObjectStorage` for a replica |
| 204 | // Durable Object instance. Replicas need to retain a reference to their primary so they can |
| 205 | // forward write requests, and since we already have a reference to the primary prior to |
| 206 | // constructing the `DurableObjectStorage`, we can just pass in the information we need to build |
| 207 | // a stub. The stub is then stored in `maybePrimary`. |
| 208 | DurableObjectStorage(jsg::Lock& js, |
| 209 | IoPtr<ActorCacheInterface> cache, |
| 210 | bool enableSql, |
| 211 | kj::Own<IoChannelFactory::ActorChannel> primaryActorChannel, |
| 212 | kj::Own<ActorIdFactory::ActorId> primaryActorId); |
| 213 | |
| 214 | ActorCacheInterface& getActorCacheInterface() { |
| 215 | return *cache; |
| 216 | } |
| 217 | |
| 218 | // Throws if not SQLite-backed. |
| 219 | SqliteDatabase& getSqliteDb(jsg::Lock& js); |
| 220 | SqliteKv& getSqliteKv(jsg::Lock& js); |
| 221 | |
| 222 | struct TransactionOptions { |
| 223 | jsg::Optional<kj::Date> asOfTime; |
| 224 | jsg::Optional<bool> lowPriority; |
| 225 | |
| 226 | JSG_STRUCT(asOfTime, lowPriority); |
| 227 | JSG_STRUCT_TS_OVERRIDE(type TransactionOptions = never); |
| 228 | // Omit from definitions |
| 229 | }; |
| 230 | |
| 231 | jsg::Promise<jsg::JsRef<jsg::JsValue>> transaction(jsg::Lock& js, |
| 232 | jsg::Function<jsg::Promise<jsg::JsRef<jsg::JsValue>>(jsg::Ref<DurableObjectTransaction>)> |
| 233 | closure, |
| 234 | jsg::Optional<TransactionOptions> options); |
| 235 | |
| 236 | jsg::JsRef<jsg::JsValue> transactionSync( |
| 237 | jsg::Lock& js, jsg::Function<jsg::JsRef<jsg::JsValue>()> callback); |
| 238 | |
| 239 | jsg::Promise<void> deleteAll(jsg::Lock& js, jsg::Optional<PutOptions> options); |
| 240 | |
| 241 | jsg::Promise<void> sync(jsg::Lock& js); |
| 242 | |
| 243 | jsg::Ref<SqlStorage> getSql(jsg::Lock& js); |
| 244 | |
| 245 | jsg::Ref<SyncKvStorage> getKv(jsg::Lock& js); |
| 246 | |
| 247 | // Get a bookmark for the current state of the database. Note that since this is async, the |
| 248 | // bookmark will include any writes in the current atomic batch, including writes that are |
| 249 | // performed after this call begins. It could also include concurrent writes that haven't happened |
| 250 | // yet, unless blockConcurrencyWhile() is used to prevent them. |
| 251 | kj::Promise<kj::String> getCurrentBookmark(); |
| 252 | |
| 253 | // Get a bookmark representing approximately the given timestamp, which is a time up to 30 days |
| 254 | // in the past (or whatever the backup retention period is). |
| 255 | kj::Promise<kj::String> getBookmarkForTime(kj::Date timestamp); |
| 256 | |
| 257 | // Arrange that the next time the Durable Object restarts, the database will be restored to |
| 258 | // the state represented by the given bookmark. This returns a bookmark string which represents |
| 259 | // the state immediately before the restoration takes place, and thus can be used to undo the |
| 260 | // restore. (This bookmark technically refers to a *future* state -- it specifies the state the |
| 261 | // object will have at the end of the current session.) |
| 262 | // |
| 263 | // It is up to the caller to force a restart in order to complete the restoration, for instance |
| 264 | // by calling state.abort() or by throwing from a blockConcurrencyWhile() callback. |
| 265 | kj::Promise<kj::String> onNextSessionRestoreBookmark(kj::String bookmark); |
| 266 | |
| 267 | // Wait until the database has been updated to the state represented by `bookmark`. |
| 268 | // |
| 269 | // `waitForBookmark` is useful synchronizing requests across replicas of the same database. On |
| 270 | // primary databases, `waitForBookmark` will resolve immediately. On replica databases, |
| 271 | // `waitForBookmark` will resolve when the replica has been updated to a point at or after |
| 272 | // `bookmark`. |
| 273 | kj::Promise<void> waitForBookmark(kj::String bookmark); |
| 274 | |
| 275 | // Arrange to create replicas for this Durable Object. |
| 276 | // |
| 277 | // Once a Durable Object instance calls `ensureReplicas`, all subsequent calls will be no-ops, |
| 278 | // making it idempotent, unless `disableReplicas` has been called between `ensureReplicas` calls. |
| 279 | // |
| 280 | // Deprecated: See DurableObjectState::configureReadReplication. |
| 281 | void ensureReplicas(); |
| 282 | |
| 283 | // Arrange to disable replicas for this Durable Object. |
| 284 | // |
| 285 | // If replicas have never been created, this is a no-op. Similar to `ensureReplicas`, repeated |
| 286 | // calls are no-ops unless `ensureReplicas` re-enabled the replicas. |
| 287 | // |
| 288 | // Deprecated: See DurableObjectState::configureReadReplication. |
| 289 | void disableReplicas(); |
| 290 | |
| 291 | // getPrimary returns a new jsg::Ref to the primary stub, if there is one. |
| 292 | // If you only need to find out if we're a replica, isReplica does so without additional |
| 293 | // refcounting requirements. |
| 294 | jsg::Optional<jsg::Ref<DurableObject>> getPrimary(jsg::Lock& js); |
| 295 | |
| 296 | // isReplica returns whether we are a replica. |
| 297 | bool isReplica(); |
| 298 | |
| 299 | JSG_RESOURCE_TYPE(DurableObjectStorage, CompatibilityFlags::Reader flags) { |
| 300 | JSG_METHOD(get); |
| 301 | JSG_METHOD(list); |
| 302 | JSG_METHOD(put); |
| 303 | JSG_METHOD_NAMED(delete, delete_); |
| 304 | JSG_METHOD(deleteAll); |
| 305 | JSG_METHOD(transaction); |
| 306 | JSG_METHOD(getAlarm); |
| 307 | JSG_METHOD(setAlarm); |
| 308 | JSG_METHOD(deleteAlarm); |
| 309 | JSG_METHOD(sync); |
| 310 | |
| 311 | JSG_LAZY_INSTANCE_PROPERTY(sql, getSql); |
| 312 | JSG_LAZY_INSTANCE_PROPERTY(kv, getKv); |
| 313 | JSG_METHOD(transactionSync); |
| 314 | |
| 315 | JSG_METHOD(getCurrentBookmark); |
| 316 | JSG_METHOD(getBookmarkForTime); |
| 317 | JSG_METHOD(onNextSessionRestoreBookmark); |
| 318 | |
| 319 | if (flags.getWorkerdExperimental()) { |
| 320 | JSG_METHOD(waitForBookmark); |
| 321 | JSG_READONLY_INSTANCE_PROPERTY(primary, getPrimary); |
| 322 | } |
| 323 | |
| 324 | if (flags.getReplicaRouting()) { |
| 325 | JSG_METHOD(ensureReplicas); |
| 326 | JSG_METHOD(disableReplicas); |
| 327 | } |
| 328 | |
| 329 | JSG_TS_OVERRIDE({ |
| 330 | get<T = unknown>(key: string, options?: DurableObjectGetOptions): Promise<T | undefined>; |
| 331 | get<T = unknown>(keys: string[], options?: DurableObjectGetOptions): Promise<Map<string, T>>; |
| 332 | |
| 333 | list<T = unknown>(options?: DurableObjectListOptions): Promise<Map<string, T>>; |
| 334 | |
| 335 | put<T>(key: string, value: T, options?: DurableObjectPutOptions): Promise<void>; |
| 336 | put<T>(entries: Record<string, T>, options?: DurableObjectPutOptions): Promise<void>; |
| 337 | |
| 338 | delete(key: string, options?: DurableObjectPutOptions): Promise<boolean>; |
| 339 | delete(keys: string[], options?: DurableObjectPutOptions): Promise<number>; |
| 340 | |
| 341 | transaction<T>(closure: (txn: DurableObjectTransaction) => Promise<T>): Promise<T>; |
| 342 | transactionSync<T>(closure: () => T): T; |
| 343 | }); |
| 344 | } |
| 345 | |
| 346 | protected: |
| 347 | ActorCacheOps& getCache(kj::StringPtr op) override; |
| 348 | |
| 349 | bool useDirectIo() override { |
| 350 | return false; |
| 351 | } |
| 352 | |
| 353 | private: |
| 354 | IoPtr<ActorCacheInterface> cache; |
| 355 | bool enableSql; |
| 356 | uint transactionSyncDepth = 0; |
| 357 | |
| 358 | // Set if this is a replica Durable Object. |
| 359 | kj::Maybe<jsg::Ref<DurableObject>> maybePrimary; |
| 360 | }; |
| 361 | |
| 362 | class DurableObjectTransaction final: public jsg::Object, public DurableObjectStorageOperations { |
| 363 | public: |
| 364 | DurableObjectTransaction(IoOwn<ActorCacheInterface::Transaction> cacheTxn) |
| 365 | : cacheTxn(kj::mv(cacheTxn)) {} |
| 366 | |
| 367 | // Called from C++, not JS, after the transaction callback has completed (successfully or not). |
| 368 | // These methods do nothing if the transaction is already committed / rolled back. |
| 369 | kj::Promise<void> maybeCommit(); |
| 370 | |
| 371 | // Called from C++, not JS, after the transaction callback has completed (successfully or not). |
| 372 | // These methods do nothing if the transaction is already committed / rolled back. |
| 373 | void maybeRollback(); |
| 374 | |
| 375 | void rollback(); // called from JS |
| 376 | |
| 377 | // Just throws an exception saying this isn't supported. |
| 378 | void deleteAll(); |
| 379 | |
| 380 | JSG_RESOURCE_TYPE(DurableObjectTransaction) { |
| 381 | JSG_METHOD(get); |
| 382 | JSG_METHOD(list); |
| 383 | JSG_METHOD(put); |
| 384 | JSG_METHOD_NAMED(delete, delete_); |
| 385 | JSG_METHOD(deleteAll); |
| 386 | JSG_METHOD(rollback); |
| 387 | JSG_METHOD(getAlarm); |
| 388 | JSG_METHOD(setAlarm); |
| 389 | JSG_METHOD(deleteAlarm); |
| 390 | |
| 391 | JSG_TS_OVERRIDE({ |
| 392 | get<T = unknown>(key: string, options?: DurableObjectGetOptions): Promise<T | undefined>; |
| 393 | get<T = unknown>(keys: string[], options?: DurableObjectGetOptions): Promise<Map<string, T>>; |
| 394 | |
| 395 | list<T = unknown>(options?: DurableObjectListOptions): Promise<Map<string, T>>; |
| 396 | |
| 397 | put<T>(key: string, value: T, options?: DurableObjectPutOptions): Promise<void>; |
| 398 | put<T>(entries: Record<string, T>, options?: DurableObjectPutOptions): Promise<void>; |
| 399 | |
| 400 | delete(key: string, options?: DurableObjectPutOptions): Promise<boolean>; |
| 401 | delete(keys: string[], options?: DurableObjectPutOptions): Promise<number>; |
| 402 | |
| 403 | deleteAll: never; |
| 404 | }); |
| 405 | } |
| 406 | |
| 407 | protected: |
| 408 | ActorCacheOps& getCache(kj::StringPtr op) override; |
| 409 | |
| 410 | bool useDirectIo() override { |
| 411 | return false; |
| 412 | } |
| 413 | |
| 414 | private: |
| 415 | // Becomes null when committed or rolled back. |
| 416 | kj::Maybe<IoOwn<ActorCacheInterface::Transaction>> cacheTxn; |
| 417 | |
| 418 | bool rolledBack = false; |
| 419 | |
| 420 | friend DurableObjectStorage; |
| 421 | }; |
| 422 | |
| 423 | class DurableObjectFacets: public jsg::Object { |
| 424 | public: |
| 425 | DurableObjectFacets(kj::Maybe<IoPtr<Worker::Actor::FacetManager>> facetManager) |
| 426 | : facetManager(kj::mv(facetManager)) {} |
| 427 | |
| 428 | // Describes how to run a facet. The app provides this when first accessing a facet that isn't |
| 429 | // already running. |
| 430 | struct StartupOptions { |
| 431 | // The actor class to use to implement the facet. |
| 432 | // |
| 433 | // Note that the $ is needed only because `class` is a keyword in C++. JSG removes the $ from |
| 434 | // the name in the JS API. C++ does not officially recognize the existence of a $ symbol but |
| 435 | // all major compilers support using it as if it were a letter. |
| 436 | kj::OneOf<jsg::Ref<DurableObjectClass>, |
| 437 | jsg::Ref<LoopbackDurableObjectNamespace>, |
| 438 | jsg::Ref<LoopbackColoLocalActorNamespace>> |
| 439 | $class; |
| 440 | |
| 441 | // Value to expose as `ctx.id` in the facet. |
| 442 | jsg::Optional<kj::OneOf<jsg::Ref<DurableObjectId>, kj::String>> id; |
| 443 | |
| 444 | JSG_STRUCT($class, id); |
| 445 | |
| 446 | JSG_STRUCT_TS_OVERRIDE(FacetStartupOptions< |
| 447 | T extends Rpc.DurableObjectBranded | undefined = undefined> { |
| 448 | class: DurableObjectClass<T>; |
| 449 | id?: DurableObjectId | string; |
| 450 | |
| 451 | $class: never; // work around generate-types bug |
| 452 | }); |
| 453 | }; |
| 454 | |
| 455 | // Get a facet by name, starting it if it isn't already running. `getStartupOptions` is invoked |
| 456 | // only if the facet wasn't already running, to get information needed to start the facet. |
| 457 | // |
| 458 | // Returns a `Fetcher` instead of a `DurableObject` becasue the returend stub does not have the |
| 459 | // `id` or `name` methods that a DO stub normally has. |
| 460 | jsg::Ref<Fetcher> get(jsg::Lock& js, |
| 461 | kj::String name, |
| 462 | jsg::Function<jsg::Promise<StartupOptions>()> getStartupOptions); |
| 463 | |
| 464 | void abort(jsg::Lock& js, kj::String name, jsg::JsValue reason); |
| 465 | void delete_(jsg::Lock& js, kj::String name); |
| 466 | |
| 467 | JSG_RESOURCE_TYPE(DurableObjectFacets) { |
| 468 | JSG_METHOD(get); |
| 469 | JSG_METHOD(abort); |
| 470 | JSG_METHOD_NAMED(delete, delete_); |
| 471 | |
| 472 | JSG_TS_OVERRIDE({ |
| 473 | get<T extends Rpc.DurableObjectBranded | undefined = undefined>( |
| 474 | name: string, |
| 475 | getStartupOptions: () => FacetStartupOptions<T> | Promise<FacetStartupOptions<T>>) |
| 476 | : Fetcher<T>; |
| 477 | }); |
| 478 | } |
| 479 | |
| 480 | private: |
| 481 | kj::Maybe<IoPtr<Worker::Actor::FacetManager>> facetManager; |
| 482 | |
| 483 | Worker::Actor::FacetManager& getFacetManager() { |
| 484 | return *JSG_REQUIRE_NONNULL( |
| 485 | facetManager, Error, "This Durable Object does not support creating facets."); |
| 486 | } |
| 487 | }; |
| 488 | |
| 489 | // The type placed in event.actorState (pre-modules API). |
| 490 | // NOTE: It hasn't been renamed under the assumption that it will only be |
| 491 | // used for colo-local namespaces. |
| 492 | class ActorState: public jsg::Object { |
| 493 | // TODO(cleanup): Remove getPersistent method that isn't supported for colo-local actors anymore. |
| 494 | public: |
| 495 | ActorState(Worker::Actor::Id actorId, |
| 496 | kj::Maybe<jsg::JsRef<jsg::JsValue>> transient, |
| 497 | kj::Maybe<jsg::Ref<DurableObjectStorage>> persistent); |
| 498 | |
| 499 | kj::OneOf<jsg::Ref<DurableObjectId>, kj::StringPtr> getId(jsg::Lock& js); |
| 500 | |
| 501 | jsg::Optional<jsg::JsValue> getTransient(jsg::Lock& js) { |
| 502 | return transient.map([&](jsg::JsRef<jsg::JsValue>& v) { return v.getHandle(js); }); |
| 503 | } |
| 504 | |
| 505 | jsg::Optional<jsg::Ref<DurableObjectStorage>> getPersistent() { |
| 506 | return persistent.map([&](jsg::Ref<DurableObjectStorage>& p) { return p.addRef(); }); |
| 507 | } |
| 508 | |
| 509 | JSG_RESOURCE_TYPE(ActorState) { |
| 510 | JSG_READONLY_INSTANCE_PROPERTY(id, getId); |
| 511 | JSG_READONLY_INSTANCE_PROPERTY(transient, getTransient); |
| 512 | JSG_READONLY_INSTANCE_PROPERTY(persistent, getPersistent); |
| 513 | |
| 514 | JSG_TS_OVERRIDE(type ActorState = never); |
| 515 | } |
| 516 | |
| 517 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 518 | KJ_SWITCH_ONEOF(id) { |
| 519 | KJ_CASE_ONEOF(str, kj::String) { |
| 520 | tracker.trackField("id", str); |
| 521 | } |
| 522 | KJ_CASE_ONEOF(id, kj::Own<ActorIdFactory::ActorId>) { |
| 523 | // TODO(later): This only yields the shallow size of the ActorId and not the |
| 524 | // size of the actual value. Should probably make ActorID a MemoryRetainer. |
| 525 | tracker.trackFieldWithSize("id", sizeof(ActorIdFactory::ActorId)); |
| 526 | } |
| 527 | } |
| 528 | tracker.trackField("transient", transient); |
| 529 | tracker.trackField("persistent", persistent); |
| 530 | } |
| 531 | |
| 532 | private: |
| 533 | Worker::Actor::Id id; |
| 534 | kj::Maybe<jsg::JsRef<jsg::JsValue>> transient; |
| 535 | kj::Maybe<jsg::Ref<DurableObjectStorage>> persistent; |
| 536 | }; |
| 537 | |
| 538 | class WebSocketRequestResponsePair: public jsg::Object { |
| 539 | public: |
| 540 | WebSocketRequestResponsePair(kj::String request, kj::String response) |
| 541 | : request(kj::mv(request)), |
| 542 | response(kj::mv(response)) {}; |
| 543 | |
| 544 | static jsg::Ref<WebSocketRequestResponsePair> constructor( |
| 545 | jsg::Lock& js, kj::String request, kj::String response) { |
| 546 | return js.alloc<WebSocketRequestResponsePair>(kj::mv(request), kj::mv(response)); |
| 547 | }; |
| 548 | |
| 549 | kj::StringPtr getRequest() { |
| 550 | return request.asPtr(); |
| 551 | } |
| 552 | kj::StringPtr getResponse() { |
| 553 | return response.asPtr(); |
| 554 | } |
| 555 | |
| 556 | JSG_RESOURCE_TYPE(WebSocketRequestResponsePair) { |
| 557 | JSG_READONLY_PROTOTYPE_PROPERTY(request, getRequest); |
| 558 | JSG_READONLY_PROTOTYPE_PROPERTY(response, getResponse); |
| 559 | } |
| 560 | |
| 561 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 562 | tracker.trackField("request", request); |
| 563 | tracker.trackField("response", response); |
| 564 | } |
| 565 | |
| 566 | private: |
| 567 | kj::String request; |
| 568 | kj::String response; |
| 569 | }; |
| 570 | |
| 571 | // The type passed as the first parameter to durable object class's constructor. |
| 572 | class DurableObjectState: public jsg::Object { |
| 573 | public: |
| 574 | DurableObjectState(jsg::Lock& js, |
| 575 | Worker::Actor::Id actorId, |
| 576 | jsg::JsValue exports, |
| 577 | jsg::JsValue props, |
| 578 | kj::Maybe<jsg::Ref<DurableObjectStorage>> storage, |
| 579 | kj::Maybe<rpc::Container::Client> container, |
| 580 | bool containerRunning, |
| 581 | kj::Maybe<Worker::Actor::FacetManager&> facetManager, |
| 582 | kj::Maybe<ActorVersion> version = kj::none); |
| 583 | |
| 584 | void waitUntil(kj::Promise<void> promise); |
| 585 | |
| 586 | jsg::JsValue getExports(jsg::Lock& js) { |
| 587 | return exports.getHandle(js); |
| 588 | } |
| 589 | |
| 590 | jsg::JsValue getProps(jsg::Lock& js) { |
| 591 | return props.getHandle(js); |
| 592 | } |
| 593 | |
| 594 | kj::OneOf<jsg::Ref<DurableObjectId>, kj::StringPtr> getId(jsg::Lock& js); |
| 595 | |
| 596 | jsg::Optional<jsg::Ref<DurableObjectStorage>> getStorage() { |
| 597 | return storage.map([&](jsg::Ref<DurableObjectStorage>& p) { return p.addRef(); }); |
| 598 | } |
| 599 | |
| 600 | struct Version { |
| 601 | jsg::Optional<kj::StringPtr> cohort; |
| 602 | JSG_STRUCT(cohort); |
| 603 | }; |
| 604 | jsg::Optional<Version> getVersion() { |
| 605 | return version.map([](ActorVersion& v) -> Version { |
| 606 | return Version{.cohort = v.cohort.map([](kj::String& s) -> kj::StringPtr { return s; })}; |
| 607 | }); |
| 608 | } |
| 609 | jsg::Optional<jsg::Ref<Container>> getContainer() { |
| 610 | return container.map([](jsg::Ref<Container>& c) { return c.addRef(); }); |
| 611 | } |
| 612 | |
| 613 | jsg::Ref<DurableObjectFacets> getFacets(jsg::Lock& js) { |
| 614 | return js.alloc<DurableObjectFacets>(facetManager); |
| 615 | } |
| 616 | |
| 617 | jsg::Promise<jsg::JsRef<jsg::JsValue>> blockConcurrencyWhile( |
| 618 | jsg::Lock& js, jsg::Function<jsg::Promise<jsg::JsRef<jsg::JsValue>>()> callback); |
| 619 | |
| 620 | // Reset the object, including breaking the output gate and canceling any writes that haven't |
| 621 | // been committed yet. |
| 622 | void abort(jsg::Lock& js, jsg::Optional<kj::String> reason); |
| 623 | |
| 624 | // Sets and returns a new hibernation manager in an actor if there's none or returns the existing. |
| 625 | Worker::Actor::HibernationManager& maybeInitHibernationManager(Worker::Actor& actor); |
| 626 | |
| 627 | // Adds a WebSocket to the set attached to this object. |
| 628 | // `ws.accept()` must NOT have been called separately. |
| 629 | // Once called, any incoming messages will be delivered |
| 630 | // by calling the Durable Object's webSocketMessage() |
| 631 | // handler, and webSocketClose() will be invoked upon |
| 632 | // disconnect. |
| 633 | // |
| 634 | // After calling this, the WebSocket is accepted, so |
| 635 | // its send() and close() methods can be used to send |
| 636 | // messages. It should be noted that calling addEventListener() |
| 637 | // on the websocket does nothing, since inbound events will |
| 638 | // automatically be delivered to one of the webSocketMessage()/ |
| 639 | // webSocketClose()/webSocketError() handlers. No inbound events |
| 640 | // to a WebSocket accepted via acceptWebSocket() will ever be |
| 641 | // delivered to addEventListener(), so there is no reason to call it. |
| 642 | // |
| 643 | // `tags` are string tags which can be used to look up |
| 644 | // the WebSocket with getWebSockets(). |
| 645 | void acceptWebSocket(jsg::Ref<WebSocket> ws, jsg::Optional<kj::Array<kj::String>> tags); |
| 646 | |
| 647 | // Gets an array of accepted WebSockets matching the given tag. |
| 648 | // If no tag is provided, an array of all accepted WebSockets is returned. |
| 649 | // Disconnected WebSockets are automatically removed from the list. |
| 650 | kj::Array<jsg::Ref<api::WebSocket>> getWebSockets(jsg::Lock& js, jsg::Optional<kj::String> tag); |
| 651 | |
| 652 | // Sets an object-wide websocket auto response message for a specific |
| 653 | // request string. All websockets belonging to the same object must |
| 654 | // reply to the request with the matching response, then store the timestamp at which |
| 655 | // the request was received. |
| 656 | // If maybeReqResp is not set, we consider it as unset and remove any set request response pair. |
| 657 | void setWebSocketAutoResponse( |
| 658 | jsg::Optional<jsg::Ref<api::WebSocketRequestResponsePair>> maybeReqResp); |
| 659 | |
| 660 | // Gets the currently set object-wide websocket auto response. |
| 661 | kj::Maybe<jsg::Ref<api::WebSocketRequestResponsePair>> getWebSocketAutoResponse(jsg::Lock& js); |
| 662 | |
| 663 | // Get the last auto response timestamp or null |
| 664 | kj::Maybe<kj::Date> getWebSocketAutoResponseTimestamp(jsg::Ref<WebSocket> ws); |
| 665 | |
| 666 | // Sets or unsets the timeout for hibernatable websocket events, preventing the execution of |
| 667 | // the event from taking longer than the specified timeout, if set. |
| 668 | void setHibernatableWebSocketEventTimeout(jsg::Optional<uint32_t> timeoutMs); |
| 669 | |
| 670 | // Get the currently set hibernatable websocket event timeout if set, or kj::none if not. |
| 671 | kj::Maybe<uint32_t> getHibernatableWebSocketEventTimeout(); |
| 672 | |
| 673 | // Gets an array of tags that this websocket was accepted with. If the given websocket is not |
| 674 | // hibernatable, we'll throw an error because regular websockets do not have tags. |
| 675 | kj::Array<kj::StringPtr> getTags(jsg::Lock& js, jsg::Ref<api::WebSocket> ws); |
| 676 | |
| 677 | // Returns a stub for the primary if there is one. |
| 678 | jsg::Optional<jsg::Ref<DurableObject>> getPrimaryStub(jsg::Lock& js); |
| 679 | |
| 680 | struct ReadReplicationOptions { |
| 681 | kj::String mode; |
| 682 | |
| 683 | JSG_STRUCT(mode); |
| 684 | JSG_STRUCT_TS_OVERRIDE(DurableObjectReadReplicationOptions { mode: "auto" | "disabled"; }); |
| 685 | }; |
| 686 | |
| 687 | // Change replica settings for this Durable Object. |
| 688 | // |
| 689 | // Must be called with a mode of "auto" or "disabled". Repeat calls that set the same settings are |
| 690 | // idempotent. |
| 691 | jsg::Promise<void> configureReadReplication(jsg::Lock& js, ReadReplicationOptions options); |
| 692 | |
| 693 | JSG_RESOURCE_TYPE(DurableObjectState, CompatibilityFlags::Reader flags) { |
| 694 | JSG_METHOD(waitUntil); |
| 695 | if (flags.getEnableCtxExports()) { |
| 696 | JSG_LAZY_INSTANCE_PROPERTY(exports, getExports); |
| 697 | } |
| 698 | JSG_LAZY_INSTANCE_PROPERTY(props, getProps); |
| 699 | JSG_LAZY_INSTANCE_PROPERTY(id, getId); |
| 700 | JSG_LAZY_INSTANCE_PROPERTY(storage, getStorage); |
| 701 | JSG_LAZY_INSTANCE_PROPERTY(container, getContainer); |
| 702 | JSG_LAZY_INSTANCE_PROPERTY(facets, getFacets); |
| 703 | if (flags.getEnableVersionApi()) { |
| 704 | JSG_LAZY_INSTANCE_PROPERTY(version, getVersion); |
| 705 | } |
| 706 | |
| 707 | if (flags.getWorkerdExperimental()) { |
| 708 | JSG_LAZY_READONLY_INSTANCE_PROPERTY(primaryStub, getPrimaryStub); |
| 709 | } |
| 710 | |
| 711 | JSG_METHOD(blockConcurrencyWhile); |
| 712 | JSG_METHOD(acceptWebSocket); |
| 713 | JSG_METHOD(getWebSockets); |
| 714 | JSG_METHOD(setWebSocketAutoResponse); |
| 715 | JSG_METHOD(getWebSocketAutoResponse); |
| 716 | JSG_METHOD(getWebSocketAutoResponseTimestamp); |
| 717 | JSG_METHOD(setHibernatableWebSocketEventTimeout); |
| 718 | JSG_METHOD(getHibernatableWebSocketEventTimeout); |
| 719 | JSG_METHOD(getTags); |
| 720 | |
| 721 | JSG_METHOD(abort); |
| 722 | |
| 723 | if (flags.getReplicaRouting()) { |
| 724 | JSG_METHOD(configureReadReplication); |
| 725 | } |
| 726 | |
| 727 | JSG_TS_ROOT(); |
| 728 | |
| 729 | // Type overrides: |
| 730 | // * Define Props/Exports type parameters. |
| 731 | // * Make `storage` non-optional |
| 732 | // * Make `id` strictly `DurableObjectId` (it's only a string for colo-local actors which are |
| 733 | // not available publicly). |
| 734 | if (flags.getEnableCtxExports()) { |
| 735 | JSG_TS_OVERRIDE(<Props = unknown> { |
| 736 | readonly props: Props; |
| 737 | readonly exports: Cloudflare.Exports; |
| 738 | readonly id: DurableObjectId; |
| 739 | readonly storage: DurableObjectStorage; |
| 740 | blockConcurrencyWhile<T>(callback: () => Promise<T>): Promise<T>; |
| 741 | }); |
| 742 | } else { |
| 743 | // No ctx.exports yet. |
| 744 | JSG_TS_OVERRIDE(<Props = unknown> { |
| 745 | readonly props: Props; |
| 746 | readonly id: DurableObjectId; |
| 747 | readonly storage: DurableObjectStorage; |
| 748 | blockConcurrencyWhile<T>(callback: () => Promise<T>): Promise<T>; |
| 749 | }); |
| 750 | } |
| 751 | } |
| 752 | |
| 753 | void visitForMemoryInfo(jsg::MemoryTracker& tracker) const { |
| 754 | KJ_SWITCH_ONEOF(id) { |
| 755 | KJ_CASE_ONEOF(str, kj::String) { |
| 756 | tracker.trackField("id", str); |
| 757 | } |
| 758 | KJ_CASE_ONEOF(id, kj::Own<ActorIdFactory::ActorId>) { |
| 759 | // TODO(later): This only yields the shallow size of the ActorId and not the |
| 760 | // size of the actual value. Should probably make ActorID a MemoryRetainer. |
| 761 | tracker.trackFieldWithSize("id", sizeof(ActorIdFactory::ActorId)); |
| 762 | } |
| 763 | } |
| 764 | tracker.trackField("storage", storage); |
| 765 | } |
| 766 | |
| 767 | private: |
| 768 | Worker::Actor::Id id; |
| 769 | jsg::JsRef<jsg::JsValue> exports; |
| 770 | jsg::JsRef<jsg::JsValue> props; |
| 771 | kj::Maybe<jsg::Ref<DurableObjectStorage>> storage; |
| 772 | kj::Maybe<jsg::Ref<Container>> container; |
| 773 | kj::Maybe<IoPtr<Worker::Actor::FacetManager>> facetManager; |
| 774 | kj::Maybe<ActorVersion> version; |
| 775 | |
| 776 | // Limits for Hibernatable WebSocket tags. |
| 777 | |
| 778 | const size_t MAX_TAGS_PER_CONNECTION = 10; |
| 779 | const size_t MAX_TAG_LENGTH = 256; |
| 780 | }; |
| 781 | |
| 782 | #define EW_ACTOR_STATE_ISOLATE_TYPES \ |
| 783 | api::ActorState, api::DurableObjectState, api::DurableObjectTransaction, \ |
| 784 | api::DurableObjectStorage, api::DurableObjectState::ReadReplicationOptions, \ |
| 785 | api::DurableObjectStorage::TransactionOptions, \ |
| 786 | api::DurableObjectStorageOperations::ListOptions, \ |
| 787 | api::DurableObjectStorageOperations::GetOptions, \ |
| 788 | api::DurableObjectStorageOperations::GetAlarmOptions, \ |
| 789 | api::DurableObjectStorageOperations::PutOptions, \ |
| 790 | api::DurableObjectStorageOperations::SetAlarmOptions, api::WebSocketRequestResponsePair, \ |
| 791 | api::DurableObjectFacets, api::DurableObjectFacets::StartupOptions, \ |
| 792 | api::DurableObjectState::Version |
| 793 | |
| 794 | } // namespace workerd::api |