diff --git a/packages/workflows/src/driver.ts b/packages/workflows/src/driver.ts index 1df9977..e1b5f8c 100644 --- a/packages/workflows/src/driver.ts +++ b/packages/workflows/src/driver.ts @@ -53,6 +53,12 @@ export interface EngineDriver { */ delete(key: Uint8Array): Promise; + /** + * Batch delete multiple keys in a single operation. + * Should be atomic if possible. + */ + batchDelete(keys: Uint8Array[]): Promise; + /** * Delete all keys with a given prefix. */ diff --git a/packages/workflows/src/rivetkit/driver.ts b/packages/workflows/src/rivetkit/driver.ts index a174fd4..3da4330 100644 --- a/packages/workflows/src/rivetkit/driver.ts +++ b/packages/workflows/src/rivetkit/driver.ts @@ -139,6 +139,18 @@ class WorkflowStorage { ); } + async batchDelete(keys: Uint8Array[]): Promise { + if (keys.length === 0) return; + await this.#db.transaction(async (tx) => { + for (const key of keys) { + await tx.execute( + "DELETE FROM _rivet_wf_kv WHERE key = ?", + prefixWorkflowKey(key), + ); + } + }); + } + async deletePrefix(prefix: Uint8Array): Promise { const start = prefixWorkflowKey(prefix); await this.#db.execute( @@ -287,6 +299,10 @@ export class ActorWorkflowDriver implements EngineDriver { await track(this.#runCtx, this.#storage.delete(key)); } + async batchDelete(keys: Uint8Array[]): Promise { + await track(this.#runCtx, this.#storage.batchDelete(keys)); + } + async deletePrefix(prefix: Uint8Array): Promise { await track(this.#runCtx, this.#storage.deletePrefix(prefix)); } @@ -370,6 +386,10 @@ export class ActorWorkflowControlDriver implements EngineDriver { await this.#storage.delete(key); } + async batchDelete(keys: Uint8Array[]): Promise { + await this.#storage.batchDelete(keys); + } + async deletePrefix(prefix: Uint8Array): Promise { await this.#storage.deletePrefix(prefix); } diff --git a/packages/workflows/src/testing.ts b/packages/workflows/src/testing.ts index 8dbc696..d6b51ef 100644 --- a/packages/workflows/src/testing.ts +++ b/packages/workflows/src/testing.ts @@ -177,6 +177,13 @@ export class InMemoryDriver implements EngineDriver { this.kv.delete(keyToHex(key)); } + async batchDelete(keys: Uint8Array[]): Promise { + await sleep(this.latency); + for (const key of keys) { + this.kv.delete(keyToHex(key)); + } + } + async deletePrefix(prefix: Uint8Array): Promise { await sleep(this.latency); for (const [hexKey, entry] of this.kv) { diff --git a/packages/workflows/tests/compat/fixture-driver.ts b/packages/workflows/tests/compat/fixture-driver.ts index 2882488..35cf1bc 100644 --- a/packages/workflows/tests/compat/fixture-driver.ts +++ b/packages/workflows/tests/compat/fixture-driver.ts @@ -174,6 +174,12 @@ export class CompatibilityDriver { this.#rows.delete(keyString(key)); } + async batchDelete(keys: Uint8Array[]): Promise { + for (const key of keys) { + this.#rows.delete(keyString(key)); + } + } + async deletePrefix(prefix: Uint8Array): Promise { for (const [encoded, row] of this.#rows) { if (startsWith(row.key, prefix)) this.#rows.delete(encoded);