From ca66cd5a3be807f4e33d06a884090cafc4f6413a Mon Sep 17 00:00:00 2001 From: vaish Date: Thu, 23 Jul 2026 15:33:10 +0530 Subject: [PATCH] feat: support individual and batch workflow instance deletion --- .changeset/workflows-instances-delete.md | 11 + .../scripts/openapi-filter-config.ts | 101 ++++++++ packages/miniflare/src/index.ts | 202 +++++++++++++-- .../miniflare/src/plugins/workflows/index.ts | 2 + .../workers/local-explorer/explorer.worker.ts | 13 + .../workers/local-explorer/generated/index.ts | 5 + .../local-explorer/generated/types.gen.ts | 42 +++ .../local-explorer/generated/zod.gen.ts | 44 ++++ .../workers/local-explorer/openapi.local.json | 109 ++++++++ .../local-explorer/resources/workflows.ts | 54 ++++ .../src/workers/local-explorer/route-names.ts | 4 + .../workflows/wrapped-binding.worker.ts | 13 +- .../test/plugins/workflows/index.spec.ts | 237 ++++++++++++++++- packages/workflows-shared/src/binding.ts | 169 ++++++++++++ packages/workflows-shared/src/engine.ts | 32 +++ packages/workflows-shared/src/lib/errors.ts | 5 + packages/workflows-shared/src/types.ts | 9 + .../workflows-shared/tests/binding.test.ts | 240 ++++++++++++++++++ .../wrangler/src/__tests__/workflows.test.ts | 148 +++++++++++ packages/wrangler/src/index.ts | 5 + .../workflows/commands/instances/delete.ts | 125 +++++++++ 21 files changed, 1543 insertions(+), 27 deletions(-) create mode 100644 .changeset/workflows-instances-delete.md create mode 100644 packages/wrangler/src/workflows/commands/instances/delete.ts diff --git a/.changeset/workflows-instances-delete.md b/.changeset/workflows-instances-delete.md new file mode 100644 index 00000000000..90ab8cf6ec8 --- /dev/null +++ b/.changeset/workflows-instances-delete.md @@ -0,0 +1,11 @@ +--- +"@cloudflare/workflows-shared": minor +"wrangler": minor +"miniflare": minor +--- + +Add individual and batch Workflow instance deletion to the runtime and SDK. + +- `WorkflowInstance.delete()` deletes one instance. Self-deletion stops the current execution. +- `env.MY_WORKFLOW.deleteBatch(instanceIds)` deletes up to 100 instances and returns `{ deleted, errors }` per input position. +- `wrangler workflows instances delete [id..]` deletes instances remotely or with `--local`; IDs can also come from a JSON array passed with `--filename`, with a combined limit of 100. diff --git a/packages/miniflare/scripts/openapi-filter-config.ts b/packages/miniflare/scripts/openapi-filter-config.ts index d0968087c33..16dcd17278e 100644 --- a/packages/miniflare/scripts/openapi-filter-config.ts +++ b/packages/miniflare/scripts/openapi-filter-config.ts @@ -985,6 +985,107 @@ const config = { tags: ["Workflows"], }, }, + "/workflows/{workflow_name}/instances/batch/delete": { + post: { + description: "Deletes multiple workflow instances.", + operationId: "workflows-batch-delete-instances", + parameters: [ + { + in: "path", + name: "workflow_name", + required: true, + schema: { + $ref: "#/components/schemas/workflows_workflow-name", + }, + }, + ], + requestBody: { + required: true, + content: { + "application/json": { + schema: { + type: "object", + properties: { + instances: { + type: "array", + minItems: 1, + maxItems: 100, + items: { + type: "string", + minLength: 1, + maxLength: 100, + pattern: "^[a-zA-Z0-9_][a-zA-Z0-9-_]*$", + }, + }, + }, + required: ["instances"], + }, + }, + }, + }, + responses: { + "200": { + content: { + "application/json": { + schema: { + allOf: [ + { + $ref: "#/components/schemas/workers_api-response-common", + }, + { + type: "object", + properties: { + result: { + type: "object", + properties: { + deleted: { + type: "array", + items: { + type: "object", + properties: { + id: { type: "string" }, + }, + required: ["id"], + }, + }, + errors: { + type: "array", + items: { + type: "object", + properties: { + id: { type: "string" }, + code: { type: "number" }, + message: { type: "string" }, + }, + required: ["id", "code", "message"], + }, + }, + }, + required: ["deleted", "errors"], + }, + }, + }, + ], + }, + }, + }, + description: "Batch delete Workflow Instances response.", + }, + "4XX": { + content: { + "application/json": { + schema: { + $ref: "#/components/schemas/workers_api-response-common-failure", + }, + }, + }, + description: "Batch delete Workflow Instances response failure.", + }, + }, + summary: "Batch Delete Workflow Instances", + tags: ["Workflows"], + }, + }, "/workflows/{workflow_name}/instances/{instance_id}": { get: { description: "Returns the status details of a workflow instance.", diff --git a/packages/miniflare/src/index.ts b/packages/miniflare/src/index.ts index 87314a6b275..5408bf16246 100644 --- a/packages/miniflare/src/index.ts +++ b/packages/miniflare/src/index.ts @@ -7,6 +7,7 @@ import net from "node:net"; import os from "node:os"; import path from "node:path"; import { ReadableStream } from "node:stream/web"; +import { setTimeout as wait } from "node:timers/promises"; import util from "node:util"; import zlib from "node:zlib"; import { checkMacOSVersion } from "@cloudflare/cli-shared-helpers"; @@ -928,6 +929,15 @@ export function _initialiseInstanceRegistry() { return (maybeInstanceRegistry = new Map()); } +type PendingWorkflowStorageDelete = { + promise: Promise; + failed: boolean; + deleted: boolean; +}; + +const WORKFLOW_STORAGE_EXTENSIONS = [".sqlite", ".sqlite-shm", ".sqlite-wal"]; +const WORKFLOW_STORAGE_DELETE_ATTEMPTS = 41; + export class Miniflare { #previousSharedOpts?: PluginSharedOptions; #previousWorkerOpts?: PluginWorkerOptions[]; @@ -946,6 +956,10 @@ export class Miniflare { string, { browserProcess: Process; wsEndpoint: string } > = new Map(); + #pendingWorkflowStorageDeletes = new Map< + string, + PendingWorkflowStorageDelete + >(); readonly #runtime?: Runtime; readonly #removeExitHook?: () => void; @@ -1385,6 +1399,131 @@ export class Miniflare { } } + async #deleteWorkflowStorageFiles( + instancePath: string, + pendingDelete: PendingWorkflowStorageDelete + ): Promise { + let firstError: unknown; + let failed = false; + for (const ext of WORKFLOW_STORAGE_EXTENSIONS) { + const filePath = `${instancePath}${ext}`; + for ( + let attempt = 0; + attempt < WORKFLOW_STORAGE_DELETE_ATTEMPTS; + attempt++ + ) { + try { + await fs.promises.unlink(filePath); + if (ext === ".sqlite") { + pendingDelete.deleted = true; + } + break; + } catch (error) { + if (isFileNotFoundError(error)) { + break; + } + const code = + typeof error === "object" && error !== null && "code" in error + ? error.code + : undefined; + if ( + (code !== "EBUSY" && code !== "EPERM") || + attempt === WORKFLOW_STORAGE_DELETE_ATTEMPTS - 1 + ) { + if (!failed) { + firstError = error; + failed = true; + } + break; + } + await wait(50); + } + } + } + if (failed) { + throw firstError; + } + } + + async #runWorkflowStorageDelete( + instancePath: string, + defer: boolean, + pendingDelete: PendingWorkflowStorageDelete, + previousDelete?: PendingWorkflowStorageDelete + ): Promise { + await previousDelete?.promise; + pendingDelete.deleted = previousDelete?.deleted ?? false; + if (defer) { + await wait(100); + } + try { + await this.#deleteWorkflowStorageFiles(instancePath, pendingDelete); + } catch (error) { + pendingDelete.failed = true; + this.#log.error( + error instanceof Error ? error : new Error(String(error)) + ); + } + if ( + !pendingDelete.failed && + this.#pendingWorkflowStorageDeletes.get(instancePath) === pendingDelete + ) { + this.#pendingWorkflowStorageDeletes.delete(instancePath); + } + } + + #queueWorkflowStorageDelete( + instancePath: string, + defer: boolean + ): PendingWorkflowStorageDelete { + const previousDelete = + this.#pendingWorkflowStorageDeletes.get(instancePath); + const pendingDelete: PendingWorkflowStorageDelete = { + deleted: false, + failed: false, + promise: Promise.resolve(), + }; + this.#pendingWorkflowStorageDeletes.set(instancePath, pendingDelete); + pendingDelete.promise = this.#runWorkflowStorageDelete( + instancePath, + defer, + pendingDelete, + previousDelete + ); + return pendingDelete; + } + + async #waitForWorkflowStorageDelete( + instancePath: string, + retried = false + ): Promise { + const pendingDelete = this.#pendingWorkflowStorageDeletes.get(instancePath); + if (pendingDelete === undefined) { + return new Response(null, { status: 204 }); + } + await pendingDelete.promise; + + const latestDelete = this.#pendingWorkflowStorageDeletes.get(instancePath); + if (latestDelete === undefined) { + return new Response(null, { status: 204 }); + } + if (latestDelete !== pendingDelete) { + return this.#waitForWorkflowStorageDelete(instancePath, retried); + } + if (!pendingDelete.failed) { + this.#pendingWorkflowStorageDeletes.delete(instancePath); + return new Response(null, { status: 204 }); + } + if (retried || this.#disposeController.signal.aborted) { + return new Response("Failed to delete workflow instance", { + status: 500, + }); + } + + this.#queueWorkflowStorageDelete(instancePath, false); + return this.#waitForWorkflowStorageDelete(instancePath, true); + } + /** * Deletes a Workflow Engine DO instance by removing its .sqlite file * (and any associated -shm/-wal files) from the persistence directory. @@ -1405,9 +1544,14 @@ export class Miniflare { const hexId = slashIndex === -1 ? null - : decodeURIComponent(pathAfterPrefix.slice(slashIndex + 1)); + : decodeURIComponent( + pathAfterPrefix.slice(slashIndex + 1) + ).toLowerCase(); assert(workflowName, "Workflow name is required"); + if (url.searchParams.has("waitForPendingDelete") && !hexId) { + return new Response("Instance ID is required", { status: 400 }); + } const coreSharedOpts = this.#sharedOpts.core; const workflowsPersistPath = getPersistPath( @@ -1426,28 +1570,30 @@ export class Miniflare { return new Response("Invalid workflow name", { status: 400 }); } - const extensions = [".sqlite", ".sqlite-shm", ".sqlite-wal"]; - if (hexId) { - // Delete a single instance - let deleted = false; - for (const ext of extensions) { - const filePath = path.join(namespacePath, `${hexId}${ext}`); - if (!filePath.startsWith(namespacePath + path.sep)) { - return new Response("Invalid instance ID", { status: 400 }); - } - try { - await fs.promises.unlink(filePath); - if (ext === ".sqlite") { - deleted = true; - } - } catch (e) { - if (!isFileNotFoundError(e)) { - throw e; - } - } + const instancePath = path.join(namespacePath, hexId); + if (!instancePath.startsWith(namespacePath + path.sep)) { + return new Response("Invalid instance ID", { status: 400 }); } - if (!deleted) { + + if (url.searchParams.has("waitForPendingDelete")) { + return this.#waitForWorkflowStorageDelete(instancePath); + } + + const pendingDelete = this.#queueWorkflowStorageDelete( + instancePath, + url.searchParams.has("defer") + ); + if (url.searchParams.has("defer")) { + return new Response("Accepted", { status: 202 }); + } + await pendingDelete.promise; + if (pendingDelete.failed) { + return new Response("Failed to delete workflow instance", { + status: 500, + }); + } + if (!pendingDelete.deleted) { return new Response("Not Found", { status: 404 }); } } else { @@ -1456,7 +1602,9 @@ export class Miniflare { const dirEntries = await fs.promises.readdir(namespacePath); await Promise.all( dirEntries - .filter((name) => extensions.some((ext) => name.endsWith(ext))) + .filter((name) => + WORKFLOW_STORAGE_EXTENSIONS.some((ext) => name.endsWith(ext)) + ) .map((name) => fs.promises.unlink(path.join(namespacePath, name)).catch(() => {}) ) @@ -1639,7 +1787,10 @@ export class Miniflare { } else if (url.pathname.startsWith("/core/do-storage/")) { response = await this.#handleLoopbackDOStorageRequest(url); } else if (url.pathname.startsWith("/core/workflow-storage/")) { - if (request.method === "DELETE") { + if ( + request.method === "DELETE" || + url.searchParams.has("waitForPendingDelete") + ) { response = await this.#handleLoopbackWorkflowStorageDeleteRequest(url); } else { @@ -3348,6 +3499,11 @@ export class Miniflare { // Cleanup as much as possible even if `#init()` threw await this.#proxyClient?.dispose(); await this.#runtime?.dispose(); + await Promise.all( + [...this.#pendingWorkflowStorageDeletes.values()].map( + ({ promise }) => promise + ) + ); // Close the undici Pool used for dispatching fetch requests to the // runtime. This must happen after the runtime is disposed, so that // in-flight connections are broken and close immediately. Without this, diff --git a/packages/miniflare/src/plugins/workflows/index.ts b/packages/miniflare/src/plugins/workflows/index.ts index e6c8b0108bd..25d50278e16 100644 --- a/packages/miniflare/src/plugins/workflows/index.ts +++ b/packages/miniflare/src/plugins/workflows/index.ts @@ -12,6 +12,7 @@ import { getUserBindingServiceName, ProxyNodeBinding, SERVICE_DEV_REGISTRY_PROXY, + WORKER_BINDING_SERVICE_LOOPBACK, } from "../shared"; import type { Service } from "../../runtime"; import type { Plugin, RemoteProxyConnectionString } from "../shared"; @@ -217,6 +218,7 @@ export const WORKFLOWS_PLUGIN: Plugin< name: "WORKFLOW_NAME", json: JSON.stringify(workflow.name), }, + WORKER_BINDING_SERVICE_LOOPBACK, ...(workflow.stepLimit !== undefined ? [ { diff --git a/packages/miniflare/src/workers/local-explorer/explorer.worker.ts b/packages/miniflare/src/workers/local-explorer/explorer.worker.ts index e9b4b127867..3258ee6f6c5 100644 --- a/packages/miniflare/src/workers/local-explorer/explorer.worker.ts +++ b/packages/miniflare/src/workers/local-explorer/explorer.worker.ts @@ -18,6 +18,7 @@ import { zWorkersKvNamespaceListANamespaceSKeysData, zWorkersKvNamespaceListNamespacesData, zObservabilityQueryData, + zWorkflowsBatchDeleteInstancesData, zWorkflowsChangeInstanceStatusData, zWorkflowsListInstancesData, } from "./generated/zod.gen"; @@ -45,6 +46,7 @@ import { createWorkflowInstance, deleteWorkflow, deleteWorkflowInstance, + deleteWorkflowInstances, getWorkflowDetails, getWorkflowInstanceDetails, listWorkflowInstances, @@ -315,6 +317,17 @@ app.post("/api/workflows/:workflow_name/instances", (c) => createWorkflowInstance(c, c.req.param("workflow_name")) ); +app.post( + "/api/workflows/:workflow_name/instances/batch/delete", + validateRequestBody(zWorkflowsBatchDeleteInstancesData.shape.body), + (c) => + deleteWorkflowInstances( + c, + c.req.param("workflow_name"), + c.req.valid("json") + ) +); + app.get("/api/workflows/:workflow_name/instances/:instance_id", (c) => getWorkflowInstanceDetails( c, diff --git a/packages/miniflare/src/workers/local-explorer/generated/index.ts b/packages/miniflare/src/workers/local-explorer/generated/index.ts index 4533d808860..92472cf8f41 100644 --- a/packages/miniflare/src/workers/local-explorer/generated/index.ts +++ b/packages/miniflare/src/workers/local-explorer/generated/index.ts @@ -174,6 +174,11 @@ export type { WorkersNamespaceWritable, WorkersObject, WorkersSchemasId, + WorkflowsBatchDeleteInstancesData, + WorkflowsBatchDeleteInstancesError, + WorkflowsBatchDeleteInstancesErrors, + WorkflowsBatchDeleteInstancesResponse, + WorkflowsBatchDeleteInstancesResponses, WorkflowsChangeInstanceStatusData, WorkflowsChangeInstanceStatusError, WorkflowsChangeInstanceStatusErrors, diff --git a/packages/miniflare/src/workers/local-explorer/generated/types.gen.ts b/packages/miniflare/src/workers/local-explorer/generated/types.gen.ts index 3bf15682d9f..5d017920996 100644 --- a/packages/miniflare/src/workers/local-explorer/generated/types.gen.ts +++ b/packages/miniflare/src/workers/local-explorer/generated/types.gen.ts @@ -1666,6 +1666,48 @@ export type WorkflowsCreateInstanceResponses = { export type WorkflowsCreateInstanceResponse = WorkflowsCreateInstanceResponses[keyof WorkflowsCreateInstanceResponses]; +export type WorkflowsBatchDeleteInstancesData = { + body: { + instances: Array; + }; + path: { + workflow_name: WorkflowsWorkflowName; + }; + query?: never; + url: "/workflows/{workflow_name}/instances/batch/delete"; +}; + +export type WorkflowsBatchDeleteInstancesErrors = { + /** + * Batch delete Workflow Instances response failure. + */ + "4XX": WorkersApiResponseCommonFailure; +}; + +export type WorkflowsBatchDeleteInstancesError = + WorkflowsBatchDeleteInstancesErrors[keyof WorkflowsBatchDeleteInstancesErrors]; + +export type WorkflowsBatchDeleteInstancesResponses = { + /** + * Batch delete Workflow Instances response. + */ + 200: WorkersApiResponseCommon & { + result?: { + deleted: Array<{ + id: string; + }>; + errors: Array<{ + id: string; + code: number; + message: string; + }>; + }; + }; +}; + +export type WorkflowsBatchDeleteInstancesResponse = + WorkflowsBatchDeleteInstancesResponses[keyof WorkflowsBatchDeleteInstancesResponses]; + export type WorkflowsDeleteInstanceData = { body?: never; path: { diff --git a/packages/miniflare/src/workers/local-explorer/generated/zod.gen.ts b/packages/miniflare/src/workers/local-explorer/generated/zod.gen.ts index 998064d0c3d..e4f0e1ed95d 100644 --- a/packages/miniflare/src/workers/local-explorer/generated/zod.gen.ts +++ b/packages/miniflare/src/workers/local-explorer/generated/zod.gen.ts @@ -1056,6 +1056,50 @@ export const zWorkflowsCreateInstanceResponse = zWorkersApiResponseCommon.and( }) ); +export const zWorkflowsBatchDeleteInstancesData = z.object({ + body: z.object({ + instances: z + .array( + z + .string() + .min(1) + .max(100) + .regex(/^[a-zA-Z0-9_][a-zA-Z0-9-_]*$/) + ) + .min(1) + .max(100), + }), + path: z.object({ + workflow_name: zWorkflowsWorkflowName, + }), + query: z.never().optional(), +}); + +/** + * Batch delete Workflow Instances response. + */ +export const zWorkflowsBatchDeleteInstancesResponse = + zWorkersApiResponseCommon.and( + z.object({ + result: z + .object({ + deleted: z.array( + z.object({ + id: z.string(), + }) + ), + errors: z.array( + z.object({ + id: z.string(), + code: z.number(), + message: z.string(), + }) + ), + }) + .optional(), + }) + ); + export const zWorkflowsDeleteInstanceData = z.object({ body: z.never().optional(), path: z.object({ diff --git a/packages/miniflare/src/workers/local-explorer/openapi.local.json b/packages/miniflare/src/workers/local-explorer/openapi.local.json index bdff2fa1e49..2cedc8e0a05 100644 --- a/packages/miniflare/src/workers/local-explorer/openapi.local.json +++ b/packages/miniflare/src/workers/local-explorer/openapi.local.json @@ -1632,6 +1632,115 @@ "tags": ["Workflows"] } }, + "/workflows/{workflow_name}/instances/batch/delete": { + "post": { + "description": "Deletes multiple workflow instances.", + "operationId": "workflows-batch-delete-instances", + "parameters": [ + { + "in": "path", + "name": "workflow_name", + "required": true, + "schema": { + "$ref": "#/components/schemas/workflows_workflow-name" + } + } + ], + "requestBody": { + "required": true, + "content": { + "application/json": { + "schema": { + "type": "object", + "properties": { + "instances": { + "type": "array", + "minItems": 1, + "maxItems": 100, + "items": { + "type": "string", + "minLength": 1, + "maxLength": 100, + "pattern": "^[a-zA-Z0-9_][a-zA-Z0-9-_]*$" + } + } + }, + "required": ["instances"] + } + } + } + }, + "responses": { + "200": { + "content": { + "application/json": { + "schema": { + "allOf": [ + { + "$ref": "#/components/schemas/workers_api-response-common" + }, + { + "type": "object", + "properties": { + "result": { + "type": "object", + "properties": { + "deleted": { + "type": "array", + "items": { + "type": "object", + "properties": { + "id": { + "type": "string" + } + }, + "required": ["id"] + } + }, + "errors": { + "type": "array", + "items": { + "type": "object", + "properties": { + "id": { + "type": "string" + }, + "code": { + "type": "number" + }, + "message": { + "type": "string" + } + }, + "required": ["id", "code", "message"] + } + } + }, + "required": ["deleted", "errors"] + } + } + } + ] + } + } + }, + "description": "Batch delete Workflow Instances response." + }, + "4XX": { + "content": { + "application/json": { + "schema": { + "$ref": "#/components/schemas/workers_api-response-common-failure" + } + } + }, + "description": "Batch delete Workflow Instances response failure." + } + }, + "summary": "Batch Delete Workflow Instances", + "tags": ["Workflows"] + } + }, "/workflows/{workflow_name}/instances/{instance_id}": { "get": { "description": "Returns the status details of a workflow instance.", diff --git a/packages/miniflare/src/workers/local-explorer/resources/workflows.ts b/packages/miniflare/src/workers/local-explorer/resources/workflows.ts index f626a48d5a8..dec098b84bb 100644 --- a/packages/miniflare/src/workers/local-explorer/resources/workflows.ts +++ b/packages/miniflare/src/workers/local-explorer/resources/workflows.ts @@ -7,6 +7,7 @@ import { errorResponse, wrapResponse } from "../common"; import type { AppContext } from "../common"; import type { Env } from "../explorer.worker"; import type { + WorkflowsBatchDeleteInstancesData, WorkflowsChangeInstanceStatusData, WorkflowsWorkflow, } from "../generated"; @@ -15,6 +16,7 @@ import type { RestartFromStep, WorkflowInstanceTerminateOptions, } from "@cloudflare/workflows-shared/src/binding"; +import type { WorkflowBatchDeleteResult } from "@cloudflare/workflows-shared/src/types"; import type { z } from "zod"; // ============================================================================ @@ -33,6 +35,10 @@ interface DirectoryEntry { birthtimeMs: number; } +interface WorkflowWithBatchDelete { + deleteBatch(instanceIds: string[]): Promise; +} + /** Methods on a WorkflowInstance handle (from workflow.get()). */ interface WorkflowHandle { pause(): Promise; @@ -945,6 +951,54 @@ export async function changeWorkflowInstanceStatus( } } +/** + * Delete multiple workflow instances through the local Workflow binding. + */ +export async function deleteWorkflowInstances( + c: AppContext, + workflowName: string, + body: WorkflowsBatchDeleteInstancesData["body"] +): Promise { + const workflow = getWorkflowBinding(c.env, workflowName); + + if (!workflow) { + const ownerMiniflare = await findWorkflowOwner(c, workflowName); + if (ownerMiniflare) { + const response = await fetchFromPeer( + ownerMiniflare, + `/workflows/${encodeURIComponent(workflowName)}/instances/batch/delete`, + { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify(body), + } + ); + if (response) { + return response; + } + } + + return errorResponse( + 404, + WORKFLOW_ERROR_NOT_FOUND, + `Workflow '${workflowName}' not found.` + ); + } + + try { + // TODO(vaish): remove cast once @cloudflare/workers-types ships deleteBatch + const result = await ( + workflow as unknown as WorkflowWithBatchDelete + ).deleteBatch(body.instances); + statusCountsCache.delete(workflowName); + return c.json(wrapResponse(result)); + } catch (error) { + const message = + error instanceof Error ? error.message : "Failed to delete instances"; + return errorResponse(500, 10001, message); + } +} + /** * Delete a workflow instance by removing its .sqlite persistence files. * diff --git a/packages/miniflare/src/workers/local-explorer/route-names.ts b/packages/miniflare/src/workers/local-explorer/route-names.ts index 1aa3c0464ff..5e1e2be3b32 100644 --- a/packages/miniflare/src/workers/local-explorer/route-names.ts +++ b/packages/miniflare/src/workers/local-explorer/route-names.ts @@ -16,6 +16,10 @@ const ROUTE_PATTERNS: [RegExp, string][] = [ [/^\/r2\/buckets\/[^/]+\/objects$/, "r2.objects"], [/^\/r2\/buckets\/[^/]+$/, "r2.bucket"], [/^\/r2\/buckets$/, "r2.buckets"], + [ + /^\/workflows\/[^/]+\/instances\/batch\/delete$/, + "workflows.instances.batch_delete", + ], [ /^\/workflows\/[^/]+\/instances\/[^/]+\/events\/[^/]+$/, "workflows.instance.event", diff --git a/packages/miniflare/src/workers/workflows/wrapped-binding.worker.ts b/packages/miniflare/src/workers/workflows/wrapped-binding.worker.ts index fdaa9521068..2e8d98c3de5 100644 --- a/packages/miniflare/src/workers/workflows/wrapped-binding.worker.ts +++ b/packages/miniflare/src/workers/workflows/wrapped-binding.worker.ts @@ -3,7 +3,10 @@ import type { WorkflowInstanceRestartOptions, WorkflowInstanceTerminateOptions, } from "@cloudflare/workflows-shared/src/binding"; -import type { WorkflowIntrospectionOperation } from "@cloudflare/workflows-shared/src/types"; +import type { + WorkflowBatchDeleteResult, + WorkflowIntrospectionOperation, +} from "@cloudflare/workflows-shared/src/types"; class WorkflowImpl implements Workflow { constructor(private binding: WorkflowBinding) {} @@ -34,6 +37,10 @@ class WorkflowImpl implements Workflow { }); } + async deleteBatch(instanceIds: string[]): Promise { + return this.binding.deleteBatch({ instances: instanceIds }); + } + async unsafeGetBindingName(): Promise { return this.binding.unsafeGetBindingName(); } @@ -124,6 +131,10 @@ class InstanceImpl implements WorkflowInstance { await instance.restart(options); } + public async delete(): Promise { + await this.binding.deleteInstance(this.id); + } + public async status(): Promise { using instance = await this.getInstance(); using res = (await instance.status()) as InstanceStatus & Disposable; diff --git a/packages/miniflare/test/plugins/workflows/index.spec.ts b/packages/miniflare/test/plugins/workflows/index.spec.ts index 9ab62d25b6d..239da69ec72 100644 --- a/packages/miniflare/test/plugins/workflows/index.spec.ts +++ b/packages/miniflare/test/plugins/workflows/index.spec.ts @@ -2,7 +2,8 @@ import * as fs from "node:fs/promises"; import path from "node:path"; import { scheduler } from "node:timers/promises"; import { Miniflare, WORKFLOWS_PLUGIN_NAME } from "miniflare"; -import { describe, test } from "vitest"; +import { describe, test, vi } from "vitest"; +import { CorePaths } from "../../../src/workers/core/constants"; import { useDispose, useTmp } from "../../test-shared"; import type { MiniflareOptions } from "miniflare"; @@ -118,6 +119,13 @@ const LIFECYCLE_WORKFLOW_SCRIPT = () => ` import { WorkflowEntrypoint } from "cloudflare:workers"; export class LifecycleWorkflow extends WorkflowEntrypoint { async run(event, step) { + if (event.payload?.selfDelete) { + await step.waitForEvent("self-delete", { type: "self-delete" }); + const instance = await this.env.LIFECYCLE_WORKFLOW.get(event.instanceId); + await instance.delete(); + throw new Error("continued after self-delete"); + } + const first = await step.do("first step", async () => "step-1-done"); await step.do("long step", async () => { @@ -136,7 +144,10 @@ export default { const id = url.searchParams.get("id") || "lifecycle-test"; if (url.pathname === "/create") { - const instance = await env.LIFECYCLE_WORKFLOW.create({ id }); + const instance = await env.LIFECYCLE_WORKFLOW.create({ + id, + params: { selfDelete: url.searchParams.has("selfDelete") }, + }); const status = await instance.status(); return Response.json({ id: instance.id, status }); } @@ -170,9 +181,24 @@ export default { return Response.json(await instance.status()); } + if (url.pathname === "/delete") { + const instance = await env.LIFECYCLE_WORKFLOW.get(id); + await instance.delete(); + return Response.json({ ok: true }); + } + + if (url.pathname === "/deleteBatch") { + return Response.json( + await env.LIFECYCLE_WORKFLOW.deleteBatch(url.searchParams.getAll("id")) + ); + } + if (url.pathname === "/sendEvent") { const instance = await env.LIFECYCLE_WORKFLOW.get(id); - await instance.sendEvent({ type: "continue", payload: { sent: true } }); + await instance.sendEvent({ + type: url.searchParams.get("type") || "continue", + payload: { sent: true }, + }); return Response.json({ ok: true }); } @@ -196,6 +222,31 @@ function lifecycleMiniflareOpts(tmp: string): MiniflareOptions { }; } +async function getPersistedInstanceFiles(tmp: string): Promise { + try { + const files = await fs.readdir( + path.join( + tmp, + WORKFLOWS_PLUGIN_NAME, + "miniflare-workflows-LIFECYCLE_WORKFLOW" + ) + ); + return files.filter( + (file) => file.endsWith(".sqlite") && file !== "metadata.sqlite" + ); + } catch (error) { + if ( + typeof error === "object" && + error !== null && + "code" in error && + error.code === "ENOENT" + ) { + return []; + } + throw error; + } +} + async function waitForStatus( mf: Miniflare, id: string, @@ -306,6 +357,186 @@ describe("workflow instance lifecycle methods", () => { await waitForStatus(mf, "terminate-test", "terminated"); }); + test("delete a workflow", async ({ expect }) => { + const tmp = await useTmp(); + const mf = new Miniflare(lifecycleMiniflareOpts(tmp)); + useDispose(mf); + + const createResponse = await mf.dispatchFetch( + "http://localhost/create?id=delete-one" + ); + await createResponse.text(); + + expect(await getPersistedInstanceFiles(tmp)).toHaveLength(1); + const deleteResponse = await mf.dispatchFetch( + "http://localhost/delete?id=delete-one" + ); + expect(await deleteResponse.json()).toEqual({ ok: true }); + expect(await getPersistedInstanceFiles(tmp)).toHaveLength(0); + + const statusResponse = await mf.dispatchFetch( + "http://localhost/status?id=delete-one" + ); + expect(statusResponse.status).toBe(500); + expect(await statusResponse.text()).toContain("instance.not_found"); + }); + + test("reports overlapping storage deletion as successful", async ({ + expect, + }) => { + const tmp = await useTmp(); + const mf = new Miniflare({ + ...lifecycleMiniflareOpts(tmp), + unsafeLocalExplorer: true, + }); + useDispose(mf); + + const createResponse = await mf.dispatchFetch( + "http://localhost/create?id=overlapping-delete" + ); + await createResponse.text(); + + const bindingDelete = mf + .dispatchFetch("http://localhost/delete?id=overlapping-delete") + .then((response) => response.json()); + await scheduler.wait(25); + const explorerDelete = await mf.dispatchFetch( + `http://localhost${CorePaths.EXPLORER}/api/workflows/LIFECYCLE_WORKFLOW/instances/overlapping-delete`, + { method: "DELETE" } + ); + expect(explorerDelete.status).toBe(200); + await explorerDelete.text(); + expect(await bindingDelete).toEqual({ ok: true }); + }); + + test("continues deleting storage files after an unlink error", async ({ + expect, + }) => { + const tmp = await useTmp(); + const mf = new Miniflare({ + ...lifecycleMiniflareOpts(tmp), + unsafeLocalExplorer: true, + }); + useDispose(mf); + + const hexId = "a".repeat(64); + const instancePath = path.join( + tmp, + WORKFLOWS_PLUGIN_NAME, + "miniflare-workflows-LIFECYCLE_WORKFLOW", + hexId + ); + await fs.mkdir(path.dirname(instancePath), { recursive: true }); + await fs.writeFile(`${instancePath}.sqlite`, ""); + await fs.mkdir(`${instancePath}.sqlite-shm`); + await fs.writeFile(`${instancePath}.sqlite-wal`, ""); + + const response = await mf.dispatchFetch( + `http://localhost${CorePaths.EXPLORER}/api/workflows/LIFECYCLE_WORKFLOW/instances/${hexId}`, + { method: "DELETE" } + ); + const body = await response.text(); + expect(response.status, body).toBe(500); + await expect(fs.stat(`${instancePath}.sqlite`)).rejects.toMatchObject({ + code: "ENOENT", + }); + await expect(fs.stat(`${instancePath}.sqlite-wal`)).rejects.toMatchObject({ + code: "ENOENT", + }); + expect((await fs.stat(`${instancePath}.sqlite-shm`)).isDirectory()).toBe( + true + ); + }); + + test("recreates a workflow immediately after deletion", async ({ + expect, + }) => { + const tmp = await useTmp(); + const mf = new Miniflare(lifecycleMiniflareOpts(tmp)); + useDispose(mf); + + let response = await mf.dispatchFetch( + "http://localhost/create?id=delete-recreate" + ); + await response.text(); + await waitForStatus(mf, "delete-recreate", "complete"); + response = await mf.dispatchFetch( + "http://localhost/delete?id=delete-recreate" + ); + await response.text(); + response = await mf.dispatchFetch( + "http://localhost/create?id=delete-recreate" + ); + await response.text(); + + await waitForStatus(mf, "delete-recreate", "complete"); + expect(await getPersistedInstanceFiles(tmp)).toHaveLength(1); + }); + + test("delete a workflow from its own execution", async ({ expect }) => { + const tmp = await useTmp(); + const mf = new Miniflare(lifecycleMiniflareOpts(tmp)); + useDispose(mf); + + const createResponse = await mf.dispatchFetch( + "http://localhost/create?id=self-delete&selfDelete" + ); + await createResponse.text(); + expect(await getPersistedInstanceFiles(tmp)).toHaveLength(1); + + const eventResponse = await mf.dispatchFetch( + "http://localhost/sendEvent?id=self-delete&type=self-delete" + ); + await eventResponse.text(); + await vi.waitUntil( + async () => (await getPersistedInstanceFiles(tmp)).length === 0, + { timeout: 5000 } + ); + + const statusResponse = await mf.dispatchFetch( + "http://localhost/status?id=self-delete" + ); + expect(statusResponse.status).toBe(500); + expect(await statusResponse.text()).toContain("instance.not_found"); + }); + + test("delete multiple workflows", async ({ expect }) => { + const tmp = await useTmp(); + const mf = new Miniflare(lifecycleMiniflareOpts(tmp)); + useDispose(mf); + + for (const id of ["delete-1", "delete-2"]) { + const response = await mf.dispatchFetch( + `http://localhost/create?id=${id}` + ); + await response.text(); + } + + expect(await getPersistedInstanceFiles(tmp)).toHaveLength(2); + const response = await mf.dispatchFetch( + "http://localhost/deleteBatch?id=delete-1&id=missing&id=delete-2&id=delete-1" + ); + expect(await response.json()).toEqual({ + deleted: [{ id: "delete-1" }, { id: "delete-2" }, { id: "delete-1" }], + errors: [ + { + id: "missing", + code: 10400, + message: "workflows.api.error.instance.not_found", + }, + ], + }); + expect(await getPersistedInstanceFiles(tmp)).toHaveLength(0); + + for (const id of ["delete-1", "delete-2"]) { + const statusResponse = await mf.dispatchFetch( + `http://localhost/status?id=${id}` + ); + expect(statusResponse.status).toBe(500); + await statusResponse.text(); + } + }); + test("restart a running workflow", async ({ expect }) => { const tmp = await useTmp(); const mf = new Miniflare(lifecycleMiniflareOpts(tmp)); diff --git a/packages/workflows-shared/src/binding.ts b/packages/workflows-shared/src/binding.ts index 9574eeb4c48..e5ce8cb4bbc 100644 --- a/packages/workflows-shared/src/binding.ts +++ b/packages/workflows-shared/src/binding.ts @@ -1,8 +1,10 @@ import { RpcTarget, WorkerEntrypoint } from "cloudflare:workers"; import { InstanceEvent, instanceStatusName } from "./instance"; import { + isUserTriggeredDelete, isUserTriggeredPause, isUserTriggeredRestart, + createWorkflowError, isUserTriggeredTerminate, WorkflowError, } from "./lib/errors"; @@ -16,6 +18,7 @@ import type { } from "./engine"; import type { InstanceStatus as EngineInstanceStatus } from "./instance"; import type { + WorkflowBatchDeleteResult, WorkflowInstanceModifier, WorkflowIntrospectionOperation, WorkflowIntrospectionStreamResult, @@ -25,8 +28,45 @@ type Env = { ENGINE: DurableObjectNamespace; BINDING_NAME: string; WORKFLOW_NAME: string; + MINIFLARE_LOOPBACK?: Fetcher; }; +async function waitForPersistedInstanceDelete( + env: Env, + id: string | undefined +): Promise { + if (id === undefined || env.MINIFLARE_LOOPBACK === undefined) { + return; + } + + const hexId = env.ENGINE.idFromName(id).toString(); + const response = await env.MINIFLARE_LOOPBACK.fetch( + `http://localhost/core/workflow-storage/${encodeURIComponent(env.WORKFLOW_NAME)}/${hexId}?waitForPendingDelete=1` + ); + if (!response.ok) { + throw new Error( + `Failed to wait for persisted workflow instance '${id}' deletion` + ); + } +} + +async function deletePersistedInstance(env: Env, id: string): Promise { + if (env.MINIFLARE_LOOPBACK === undefined) { + return; + } + + const stub = env.ENGINE.get(env.ENGINE.idFromName(id)); + await stub.unsafeAbort(); + + const response = await env.MINIFLARE_LOOPBACK.fetch( + `http://localhost/core/workflow-storage/${encodeURIComponent(env.WORKFLOW_NAME)}/${stub.id.toString()}`, + { method: "DELETE" } + ); + if (!response.ok && response.status !== 404) { + throw new Error(`Failed to delete persisted workflow instance '${id}'`); + } +} + type WorkflowIntrospectionSession = { id: string; operations: WorkflowIntrospectionOperation[]; @@ -151,6 +191,7 @@ export class WorkflowBinding extends WorkerEntrypoint { throw new WorkflowError("Workflow instance has invalid id"); } + await waitForPersistedInstanceDelete(this.env, id); const stubId = this.env.ENGINE.idFromName(id); const stub = this.env.ENGINE.get(stubId); const introspectionSession = workflowIntrospectionSessions.get( @@ -240,6 +281,123 @@ export class WorkflowBinding extends WorkerEntrypoint { ); } + // deleteInstance, not delete: Fetcher.delete(url) shadows a same-named JSRPC method. + public async deleteInstance(id: string): Promise { + if (!isValidWorkflowInstanceId(id)) { + throw createWorkflowError( + "Instance ID is invalid", + "instance.invalid_id" + ); + } + + const stub = this.env.ENGINE.get(this.env.ENGINE.idFromName(id)); + try { + await stub.deleteInstance(); + } catch (error) { + // delete aborts the instance + if (!isUserTriggeredDelete(error)) { + throw error; + } + } + await waitForPersistedInstanceDelete(this.env, id); + } + + public async deleteBatch(options: { + instances: string[]; + }): Promise { + const instanceIds = options?.instances; + if (!Array.isArray(instanceIds)) { + throw createWorkflowError("Provided argument is invalid", "body"); + } + if (instanceIds.length > 100) { + throw createWorkflowError( + "batchDeleteInstances only supports 100 instances at a time", + "body" + ); + } + if (instanceIds.length === 0) { + throw createWorkflowError( + "batchDeleteInstances should have at least 1 instance", + "body" + ); + } + if (!instanceIds.every(isValidWorkflowInstanceId)) { + throw createWorkflowError( + "Instance ID is invalid", + "instance.invalid_id" + ); + } + + const uniqueIds = [...new Set(instanceIds)]; + const settled = await Promise.allSettled( + uniqueIds.map((id) => + this.env.ENGINE.get(this.env.ENGINE.idFromName(id)).deleteInstance() + ) + ); + const resultsById = new Map( + uniqueIds.map((id, index) => [id, settled[index]]) + ); + const result: WorkflowBatchDeleteResult = { deleted: [], errors: [] }; + for (const id of instanceIds) { + const deletion = resultsById.get(id); + if (deletion === undefined) { + throw new Error("Missing batch deletion result"); + } + if ( + deletion.status === "fulfilled" || + isUserTriggeredDelete(deletion.reason) + ) { + result.deleted.push({ id }); + continue; + } + + const isNotFound = + deletion.reason instanceof Error && + deletion.reason.message.includes("(instance.not_found)"); + result.errors.push({ + id, + code: isNotFound ? 10400 : 10001, + message: isNotFound + ? "workflows.api.error.instance.not_found" + : "workflows.api.error.internal_server", + }); + } + + const missingIds = new Set( + result.errors.filter(({ code }) => code === 10400).map(({ id }) => id) + ); + const cleanupIds = [ + ...new Set([...result.deleted.map(({ id }) => id), ...missingIds]), + ]; + const cleanups = await Promise.allSettled( + cleanupIds.map((id) => + missingIds.has(id) + ? deletePersistedInstance(this.env, id) + : waitForPersistedInstanceDelete(this.env, id) + ) + ); + const failedCleanupIds = new Set( + cleanupIds.filter((_, index) => cleanups[index]?.status === "rejected") + ); + if (failedCleanupIds.size === 0) { + return result; + } + + const errorsById = new Map(result.errors.map((error) => [error.id, error])); + return { + deleted: result.deleted.filter(({ id }) => !failedCleanupIds.has(id)), + errors: instanceIds.flatMap((id) => { + if (failedCleanupIds.has(id)) { + return [ + { id, code: 10001, message: "workflows.api.error.internal_server" }, + ]; + } + const error = errorsById.get(id); + return error === undefined ? [] : [error]; + }), + }; + } + public async unsafeGetBindingName(): Promise { // async because of rpc return this.env.BINDING_NAME; @@ -377,6 +535,17 @@ export class WorkflowHandle extends RpcTarget implements WorkflowInstance { } } + public async delete(): Promise { + try { + await this.stub.deleteInstance(); + } catch (e) { + // delete aborts the instance + if (!isUserTriggeredDelete(e)) { + throw e; + } + } + } + public async restart( options?: WorkflowInstanceRestartOptions ): Promise { diff --git a/packages/workflows-shared/src/engine.ts b/packages/workflows-shared/src/engine.ts index a6e96855e1f..054d56dcf2b 100644 --- a/packages/workflows-shared/src/engine.ts +++ b/packages/workflows-shared/src/engine.ts @@ -64,6 +64,8 @@ import type { interface Env { ENGINE: DurableObjectNamespace; USER_WORKFLOW: WorkflowEntrypoint; + MINIFLARE_LOOPBACK?: Fetcher; + WORKFLOW_NAME?: string; STEP_LIMIT?: string; // JSON-encoded number from miniflare binding } @@ -1037,6 +1039,36 @@ export class Engine extends DurableObject { await this.abort(ABORT_REASONS.USER_TERMINATE); } + async deleteInstance(): Promise { + if ((await this.ctx.storage.get(INSTANCE_METADATA)) === undefined) { + throw createWorkflowError( + "Instance does not exist", + "instance.not_found" + ); + } + + await this.ctx.storage.deleteAll(); + + if ( + this.env.MINIFLARE_LOOPBACK !== undefined && + this.env.WORKFLOW_NAME !== undefined + ) { + try { + const response = await this.env.MINIFLARE_LOOPBACK.fetch( + `http://localhost/core/workflow-storage/${encodeURIComponent(this.env.WORKFLOW_NAME)}/${this.ctx.id.toString()}?defer=1`, + { method: "DELETE" } + ); + if (!response.ok && response.status !== 404) { + console.error("Failed to delete persisted workflow instance"); + } + } catch (error) { + console.error("Failed to delete persisted workflow instance", error); + } + } + + await this.abort(ABORT_REASONS.USER_DELETE); + } + async userTriggeredPause() { const status = await this.getStatus(); diff --git a/packages/workflows-shared/src/lib/errors.ts b/packages/workflows-shared/src/lib/errors.ts index 25614b3cd46..af9632da56c 100644 --- a/packages/workflows-shared/src/lib/errors.ts +++ b/packages/workflows-shared/src/lib/errors.ts @@ -70,6 +70,7 @@ export const ABORT_REASONS = { USER_PAUSE: `${ABORT_PREFIX} User called pause`, USER_RESTART: `${ABORT_PREFIX} User called restart`, USER_TERMINATE: `${ABORT_PREFIX} User called terminate`, + USER_DELETE: `${ABORT_PREFIX} User called delete`, NON_RETRYABLE_ERROR: `${ABORT_PREFIX} A step threw a NonRetryableError`, NOT_SERIALISABLE: `${ABORT_PREFIX} Value is not serialisable`, STORAGE_LIMIT_EXCEEDED: `${ABORT_PREFIX} Storage limit exceeded`, @@ -110,6 +111,10 @@ export function isUserTriggeredTerminate(e: unknown): boolean { return getErrorMessage(e) === ABORT_REASONS.USER_TERMINATE; } +export function isUserTriggeredDelete(e: unknown): boolean { + return getErrorMessage(e) === ABORT_REASONS.USER_DELETE; +} + function getCompatFlag(name: string): boolean { // eslint-disable-next-line @typescript-eslint/no-explicit-any -- safe globalThis access for environments where cloudflare global may not exist return (globalThis as any).Cloudflare?.compatibilityFlags?.[name] ?? false; diff --git a/packages/workflows-shared/src/types.ts b/packages/workflows-shared/src/types.ts index ae8358ac358..d5ae9b869e0 100644 --- a/packages/workflows-shared/src/types.ts +++ b/packages/workflows-shared/src/types.ts @@ -1,3 +1,12 @@ +export interface WorkflowBatchDeleteResult { + deleted: { id: string }[]; + errors: Array<{ + id: string; + code: number; + message: string; + }>; +} + export type WorkflowStepSelector = { name: string; index?: number; diff --git a/packages/workflows-shared/tests/binding.test.ts b/packages/workflows-shared/tests/binding.test.ts index 03417019660..86bd237f787 100644 --- a/packages/workflows-shared/tests/binding.test.ts +++ b/packages/workflows-shared/tests/binding.test.ts @@ -108,6 +108,23 @@ describe("WorkflowBinding", () => { "Workflow instance has invalid id" ); }); + + it("should block creation when pending persistence deletion fails", async ({ + expect, + }) => { + const binding = new WorkflowBinding(createExecutionContext(), { + ENGINE: env.ENGINE, + BINDING_NAME: "TEST_WORKFLOW", + WORKFLOW_NAME: "test-workflow", + MINIFLARE_LOOPBACK: { + fetch: () => Promise.resolve(new Response(null, { status: 500 })), + } as unknown as Fetcher, + }); + + await expect(binding.create({ id: "cleanup-failed" })).rejects.toThrow( + "Failed to wait for persisted workflow instance 'cleanup-failed' deletion" + ); + }); }); describe("get()", () => { @@ -130,6 +147,7 @@ describe("WorkflowBinding", () => { resume: expect.any(Function), terminate: expect.any(Function), restart: expect.any(Function), + delete: expect.any(Function), }); // Wait for the workflow to complete before the test ends so @@ -144,6 +162,228 @@ describe("WorkflowBinding", () => { }); }); + describe("instance deletion", () => { + it("deleteInstance should delete an instance and wipe its stored state", async ({ + expect, + }) => { + const id = uniqueId(); + const binding = createBinding(); + + setTestWorkflowCallback(async () => "done"); + await binding.create({ id }); + + const instance = await binding.get(id); + await vi.waitUntil( + async () => { + const status = await instance.status(); + return status.status === "complete"; + }, + { timeout: 5000 } + ); + + await expect(binding.deleteInstance(id)).resolves.toBeUndefined(); + await expect(binding.get(id)).rejects.toThrow("instance.not_found"); + }); + + it("should reject an invalid instance ID", async ({ expect }) => { + await expect(createBinding().deleteInstance("#invalid!")).rejects.toThrow( + "(instance.invalid_id) Instance ID is invalid" + ); + }); + + it("should let a running instance delete itself and stop execution", async ({ + expect, + }) => { + const id = uniqueId(); + const binding = createBinding(); + let deleteStarted = false; + let continuedAfterDelete = false; + + setTestWorkflowCallback(async () => { + deleteStarted = true; + const instance = await binding.get(id); + await (instance as unknown as { delete(): Promise }).delete(); + continuedAfterDelete = true; + }); + await binding.create({ id }); + await vi.waitUntil(() => deleteStarted, { timeout: 5000 }); + + await vi.waitUntil( + async () => { + try { + await binding.get(id); + return false; + } catch { + return true; + } + }, + { timeout: 5000 } + ); + + await scheduler.wait(50); + expect(continuedAfterDelete).toBe(false); + await expect(binding.get(id)).rejects.toThrow("instance.not_found"); + }); + }); + + describe("deleteBatch()", () => { + it("should delete instances and wipe their stored state", async ({ + expect, + }) => { + const ids = [uniqueId(), uniqueId()]; + const binding = createBinding(); + + setTestWorkflowCallback(async () => "done"); + await binding.createBatch(ids.map((id) => ({ id }))); + + for (const id of ids) { + const instance = await binding.get(id); + await vi.waitUntil( + async () => { + const status = await instance.status(); + return status.status === "complete"; + }, + { timeout: 5000 } + ); + } + + await expect(binding.deleteBatch({ instances: ids })).resolves.toEqual({ + deleted: ids.map((id) => ({ id })), + errors: [], + }); + for (const id of ids) { + await expect(binding.get(id)).rejects.toThrow("instance.not_found"); + } + }); + + it("should report each duplicate missing ID", async ({ expect }) => { + const binding = createBinding(); + await expect( + binding.deleteBatch({ + instances: ["batch", "batch"], + }) + ).resolves.toEqual({ + deleted: [], + errors: [ + { + id: "batch", + code: 10400, + message: "workflows.api.error.instance.not_found", + }, + { + id: "batch", + code: 10400, + message: "workflows.api.error.instance.not_found", + }, + ], + }); + }); + + it("should normalize unexpected deletion errors", async ({ expect }) => { + const deleteInstance = vi + .fn() + .mockRejectedValue(new Error("sensitive failure")); + const binding = new WorkflowBinding(createExecutionContext(), { + ENGINE: { + idFromName: (id: string) => id, + get: () => ({ deleteInstance }), + } as unknown as DurableObjectNamespace, + BINDING_NAME: "TEST_WORKFLOW", + WORKFLOW_NAME: "test-workflow", + }); + + await expect( + binding.deleteBatch({ instances: ["broken-instance"] }) + ).resolves.toEqual({ + deleted: [], + errors: [ + { + id: "broken-instance", + code: 10001, + message: "workflows.api.error.internal_server", + }, + ], + }); + expect(deleteInstance).toHaveBeenCalledOnce(); + }); + + it("should report persistence cleanup failures per instance", async ({ + expect, + }) => { + const abort = vi.fn(() => Promise.resolve()); + const binding = new WorkflowBinding(createExecutionContext(), { + ENGINE: { + idFromName: (id: string) => ({ toString: () => id }), + get: (id: { toString(): string }) => ({ + id, + deleteInstance: () => { + if (id.toString() === "missing") { + return Promise.reject( + new Error("(instance.not_found) Instance does not exist") + ); + } + return Promise.resolve(); + }, + unsafeAbort: abort, + }), + } as unknown as DurableObjectNamespace, + BINDING_NAME: "TEST_WORKFLOW", + WORKFLOW_NAME: "test-workflow", + MINIFLARE_LOOPBACK: { + fetch: (url: string) => + Promise.resolve( + new Response(null, { + status: url.includes("/cleanup-failed") ? 500 : 204, + }) + ), + } as unknown as Fetcher, + }); + + await expect( + binding.deleteBatch({ + instances: ["cleanup-failed", "missing", "deleted", "cleanup-failed"], + }) + ).resolves.toEqual({ + deleted: [{ id: "deleted" }], + errors: [ + { + id: "cleanup-failed", + code: 10001, + message: "workflows.api.error.internal_server", + }, + { + id: "missing", + code: 10400, + message: "workflows.api.error.instance.not_found", + }, + { + id: "cleanup-failed", + code: 10001, + message: "workflows.api.error.internal_server", + }, + ], + }); + expect(abort).toHaveBeenCalledOnce(); + }); + + it("should reject invalid batches", async ({ expect }) => { + const binding = createBinding(); + await expect(binding.deleteBatch({ instances: [] })).rejects.toThrow( + "(body) batchDeleteInstances should have at least 1 instance" + ); + await expect( + binding.deleteBatch({ + instances: Array.from({ length: 101 }, (_, i) => `instance-${i}`), + }) + ).rejects.toThrow( + "(body) batchDeleteInstances only supports 100 instances at a time" + ); + await expect( + binding.deleteBatch({ instances: ["#invalid!"] }) + ).rejects.toThrow("(instance.invalid_id) Instance ID is invalid"); + }); + }); + describe("createBatch()", () => { it("should create multiple instances in a batch", async ({ expect }) => { const binding = createBinding(); diff --git a/packages/wrangler/src/__tests__/workflows.test.ts b/packages/wrangler/src/__tests__/workflows.test.ts index fa5b27b085f..2c7ba1aee21 100644 --- a/packages/wrangler/src/__tests__/workflows.test.ts +++ b/packages/wrangler/src/__tests__/workflows.test.ts @@ -19,6 +19,7 @@ import { runWrangler } from "./helpers/run-wrangler"; import { writeWorkerSource } from "./helpers/write-worker-source"; import type { Instance, Workflow } from "../workflows/types"; import type { RawConfig } from "@cloudflare/workers-utils"; +import type { WorkflowBatchDeleteResult } from "@cloudflare/workflows-shared/src/types"; import type { ExpectStatic } from "vitest"; describe("wrangler workflows", () => { @@ -199,6 +200,7 @@ describe("wrangler workflows", () => { wrangler workflows instances restart Restart a workflow instance wrangler workflows instances pause Pause a workflow instance wrangler workflows instances resume Resume a workflow instance + wrangler workflows instances delete [id..] Delete workflow instances GLOBAL FLAGS -c, --config Path to Wrangler configuration file [string] @@ -687,6 +689,103 @@ describe("wrangler workflows", () => { }); }); + describe("instances delete", () => { + const mockDeleteInstances = ( + expect: ExpectStatic, + expectedIds: string[], + result: WorkflowBatchDeleteResult = { + deleted: expectedIds.map((id) => ({ id })), + errors: [], + } + ) => { + msw.use( + http.post( + `*/accounts/:accountId/workflows/:workflowName/instances/batch/delete`, + async ({ request }) => { + expect(await request.json()).toEqual({ instances: expectedIds }); + return HttpResponse.json({ + success: true, + errors: [], + messages: [], + result, + }); + }, + { once: true } + ) + ); + }; + + it("should delete multiple instances", async ({ expect }) => { + writeWranglerConfig(); + mockDeleteInstances(expect, ["foo", "bar"]); + + await runWrangler(`workflows instances delete some-workflow foo bar`); + expect(std.info).toMatchInlineSnapshot( + `"🗑️ Deleted workflow instances from "some-workflow": "foo", "bar""` + ); + }); + + it("should report per-instance errors after logging deletions", async ({ + expect, + }) => { + writeWranglerConfig(); + mockDeleteInstances(expect, ["foo", "bar"], { + deleted: [{ id: "foo" }], + errors: [{ id: "bar", code: 500, message: "delete failed" }], + }); + + await expect( + runWrangler(`workflows instances delete some-workflow foo bar`) + ).rejects.toThrow( + "Failed to delete 1 workflow instance(s):\n - bar: delete failed" + ); + expect(std.info).toContain('"foo"'); + }); + + it("should read instance IDs from a JSON file", async ({ expect }) => { + writeWranglerConfig(); + fs.writeFileSync("instance-ids.json", JSON.stringify(["bar"])); + mockDeleteInstances(expect, ["foo", "bar"]); + + await runWrangler( + "workflows instances delete some-workflow foo --filename instance-ids.json" + ); + expect(std.info).toContain('"foo", "bar"'); + }); + + it("should require at least one instance ID", async ({ expect }) => { + writeWranglerConfig(); + await expect( + runWrangler("workflows instances delete some-workflow") + ).rejects.toThrow("Provide at least one workflow instance ID"); + }); + + it("should reject an invalid IDs file", async ({ expect }) => { + writeWranglerConfig(); + fs.writeFileSync("instance-ids.json", JSON.stringify(["foo", 1])); + await expect( + runWrangler( + "workflows instances delete some-workflow --filename instance-ids.json" + ) + ).rejects.toThrow( + 'Unexpected JSON input from "instance-ids.json". Expected an array of strings.' + ); + }); + + it("should reject more than 100 combined instances", async ({ expect }) => { + writeWranglerConfig(); + const ids = Array.from({ length: 100 }, (_, i) => `instance-${i}`); + fs.writeFileSync("instance-ids.json", JSON.stringify(["overflow"])); + await expect( + runWrangler( + `workflows instances delete some-workflow ${ids.join(" ")} --filename instance-ids.json` + ) + ).rejects.toThrow( + "You can delete at most 100 workflow instances at a time" + ); + }); + }); + describe("instances restart", () => { const mockInstances: Instance[] = [ { @@ -1808,6 +1907,55 @@ describe("wrangler workflows", () => { }); }); + describe("workflows instances delete --local", () => { + it("should resolve latest once before local deletion", async ({ + expect, + }) => { + writeWranglerConfig(); + let listRequests = 0; + const ids = ["newest-instance", "newest-instance", "explicit-instance"]; + msw.use( + http.get(`${LOCAL_BASE}/workflows/:workflowName/instances`, () => { + listRequests++; + return HttpResponse.json({ + success: true, + errors: [], + messages: [], + result: [ + { + id: "newest-instance", + created_on: "2024-06-01T00:00:00Z", + }, + ], + }); + }), + http.post( + `${LOCAL_BASE}/workflows/:workflowName/instances/batch/delete`, + async ({ request }) => { + expect(await request.json()).toEqual({ instances: ids }); + return HttpResponse.json({ + success: true, + errors: [], + messages: [], + result: { + deleted: ids.map((id) => ({ id })), + errors: [], + }, + }); + } + ) + ); + + await runWrangler( + "workflows instances delete my-workflow latest latest explicit-instance --local" + ); + expect(listRequests).toBe(1); + expect(std.info).toContain( + '"newest-instance", "newest-instance", "explicit-instance"' + ); + }); + }); + describe("workflows instances restart --local", () => { it("should restart an instance in local dev session", async ({ expect, diff --git a/packages/wrangler/src/index.ts b/packages/wrangler/src/index.ts index 2501578b7e7..8199c2f65d0 100644 --- a/packages/wrangler/src/index.ts +++ b/packages/wrangler/src/index.ts @@ -536,6 +536,7 @@ import { websearchSearchCommand } from "./websearch/search"; import { workflowsInstanceNamespace, workflowsNamespace } from "./workflows"; import { workflowsDeleteCommand } from "./workflows/commands/delete"; import { workflowsDescribeCommand } from "./workflows/commands/describe"; +import { workflowsInstancesDeleteCommand } from "./workflows/commands/instances/delete"; import { workflowsInstancesDescribeCommand } from "./workflows/commands/instances/describe"; import { workflowsInstancesListCommand } from "./workflows/commands/instances/list"; import { workflowsInstancesPauseCommand } from "./workflows/commands/instances/pause"; @@ -2269,6 +2270,10 @@ export function createCLIParser(argv: string[]) { command: "wrangler workflows instances resume", definition: workflowsInstancesResumeCommand, }, + { + command: "wrangler workflows instances delete", + definition: workflowsInstancesDeleteCommand, + }, ]); registry.registerNamespace("workflows"); diff --git a/packages/wrangler/src/workflows/commands/instances/delete.ts b/packages/wrangler/src/workflows/commands/instances/delete.ts new file mode 100644 index 00000000000..b5d4fa00b1b --- /dev/null +++ b/packages/wrangler/src/workflows/commands/instances/delete.ts @@ -0,0 +1,125 @@ +import { parseJSON, readFileSync, UserError } from "@cloudflare/workers-utils"; +import { fetchResult } from "../../../cfetch"; +import { createCommand } from "../../../core/create-command"; +import { logger } from "../../../logger"; +import { requireAuth } from "../../../user"; +import { + fetchLocalResult, + getLocalInstanceIdFromArgs, + localWorkflowArgs, +} from "../../local"; +import { getInstanceIdFromArgs } from "../../utils"; +import type { WorkflowBatchDeleteResult } from "@cloudflare/workflows-shared/src/types"; + +export const workflowsInstancesDeleteCommand = createCommand({ + metadata: { + description: "Delete workflow instances", + owner: "Product: Workflows", + status: "stable", + }, + positionalArgs: ["name", "id"], + args: { + ...localWorkflowArgs, + name: { + describe: "Name of the workflow", + type: "string", + demandOption: true, + }, + id: { + describe: + "IDs of the instances - you can type 'latest' to get the latest instance and delete it", + type: "string", + array: true, + }, + filename: { + describe: "Path to a JSON file containing an array of instance IDs", + type: "string", + }, + }, + + async handler(args, { config }) { + let fileIds: string[] = []; + if (args.filename) { + const parsed = parseJSON(readFileSync(args.filename), args.filename); + if ( + !Array.isArray(parsed) || + !parsed.every((id) => typeof id === "string") + ) { + throw new UserError( + `Unexpected JSON input from "${args.filename}". Expected an array of strings.`, + { telemetryMessage: "workflows batch delete invalid filename" } + ); + } + fileIds = parsed; + } + + const requestedIds = [...(args.id ?? []), ...fileIds]; + if (requestedIds.length === 0) { + throw new UserError("Provide at least one workflow instance ID", { + telemetryMessage: "workflows batch delete no ids", + }); + } + if (requestedIds.length > 100) { + throw new UserError( + "You can delete at most 100 workflow instances at a time", + { + telemetryMessage: "workflows batch delete too large", + } + ); + } + + let ids = requestedIds; + let result: WorkflowBatchDeleteResult; + + if (args.local) { + if (ids.includes("latest")) { + const latestId = await getLocalInstanceIdFromArgs(args.port, { + id: "latest", + name: args.name, + }); + ids = ids.map((id) => (id === "latest" ? latestId : id)); + } + result = await fetchLocalResult( + args.port, + `/workflows/${encodeURIComponent(args.name)}/instances/batch/delete`, + { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ instances: ids }), + } + ); + } else { + const accountId = await requireAuth(config); + if (ids.includes("latest")) { + const latestId = await getInstanceIdFromArgs( + accountId, + { id: "latest", name: args.name }, + config + ); + ids = ids.map((id) => (id === "latest" ? latestId : id)); + } + result = await fetchResult( + config, + `/accounts/${accountId}/workflows/${args.name}/instances/batch/delete`, + { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ instances: ids }), + } + ); + } + + if (result.deleted.length > 0) { + logger.info( + `🗑️ Deleted workflow instances from "${args.name}": ${result.deleted.map(({ id }) => `"${id}"`).join(", ")}` + ); + } + + if (result.errors.length > 0) { + throw new UserError( + `Failed to delete ${result.errors.length} workflow instance(s):\n${result.errors.map(({ id, message }) => ` - ${id}: ${message}`).join("\n")}`, + { telemetryMessage: "workflows batch delete partial failure" } + ); + } + }, +});