Skip to content
Open
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 6 additions & 0 deletions packages/workflows/src/driver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -53,6 +53,12 @@ export interface EngineDriver {
*/
delete(key: Uint8Array): Promise<void>;

/**
* Batch delete multiple keys in a single operation.
* Should be atomic if possible.
*/
batchDelete(keys: Uint8Array[]): Promise<void>;

/**
* Delete all keys with a given prefix.
*/
Expand Down
20 changes: 20 additions & 0 deletions packages/workflows/src/rivetkit/driver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -139,6 +139,18 @@ class WorkflowStorage {
);
}

async batchDelete(keys: Uint8Array[]): Promise<void> {
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<void> {
const start = prefixWorkflowKey(prefix);
await this.#db.execute(
Expand Down Expand Up @@ -287,6 +299,10 @@ export class ActorWorkflowDriver implements EngineDriver {
await track(this.#runCtx, this.#storage.delete(key));
}

async batchDelete(keys: Uint8Array[]): Promise<void> {
await track(this.#runCtx, this.#storage.batchDelete(keys));
}

async deletePrefix(prefix: Uint8Array): Promise<void> {
await track(this.#runCtx, this.#storage.deletePrefix(prefix));
}
Expand Down Expand Up @@ -370,6 +386,10 @@ export class ActorWorkflowControlDriver implements EngineDriver {
await this.#storage.delete(key);
}

async batchDelete(keys: Uint8Array[]): Promise<void> {
await this.#storage.batchDelete(keys);
}

async deletePrefix(prefix: Uint8Array): Promise<void> {
await this.#storage.deletePrefix(prefix);
}
Expand Down
7 changes: 7 additions & 0 deletions packages/workflows/src/testing.ts
Original file line number Diff line number Diff line change
Expand Up @@ -177,6 +177,13 @@ export class InMemoryDriver implements EngineDriver {
this.kv.delete(keyToHex(key));
}

async batchDelete(keys: Uint8Array[]): Promise<void> {
await sleep(this.latency);
for (const key of keys) {
this.kv.delete(keyToHex(key));
}
}

async deletePrefix(prefix: Uint8Array): Promise<void> {
await sleep(this.latency);
for (const [hexKey, entry] of this.kv) {
Expand Down
6 changes: 6 additions & 0 deletions packages/workflows/tests/compat/fixture-driver.ts
Original file line number Diff line number Diff line change
Expand Up @@ -174,6 +174,12 @@ export class CompatibilityDriver {
this.#rows.delete(keyString(key));
}

async batchDelete(keys: Uint8Array[]): Promise<void> {
for (const key of keys) {
this.#rows.delete(keyString(key));
}
}

async deletePrefix(prefix: Uint8Array): Promise<void> {
for (const [encoded, row] of this.#rows) {
if (startsWith(row.key, prefix)) this.#rows.delete(encoded);
Expand Down
Loading