From d41eab04db4bdcb4bb637199878a9bb7246c1da9 Mon Sep 17 00:00:00 2001 From: Sam Starling <42478+samstarling@users.noreply.github.com> Date: Fri, 18 Sep 2026 10:56:33 +0100 Subject: [PATCH 1/4] fix(workflows): dispose RPC results and introspection modifiers --- .changeset/fix-workflow-rpc-disposal.md | 9 ++ packages/workflows-shared/src/binding.ts | 2 +- packages/workflows-shared/src/context.ts | 35 +++---- .../workflows-shared/src/introspection.ts | 14 ++- .../tests/rpc-lifecycle.test.ts | 96 +++++++++++++++++++ 5 files changed, 138 insertions(+), 18 deletions(-) create mode 100644 .changeset/fix-workflow-rpc-disposal.md create mode 100644 packages/workflows-shared/tests/rpc-lifecycle.test.ts diff --git a/.changeset/fix-workflow-rpc-disposal.md b/.changeset/fix-workflow-rpc-disposal.md new file mode 100644 index 00000000000..a18800ee90f --- /dev/null +++ b/.changeset/fix-workflow-rpc-disposal.md @@ -0,0 +1,9 @@ +--- +"@cloudflare/workflows-shared": patch +"miniflare": patch +"@cloudflare/vitest-plugin": patch +--- + +Dispose Workflow step results and introspection modifiers after use + +Release RPC resources deterministically during local Workflow execution and introspection. This prevents undisposed RPC warnings and requests being cancelled after their execution context has ended. Live step results retain their original data shape, including typed-array offsets and backing buffers. diff --git a/packages/workflows-shared/src/binding.ts b/packages/workflows-shared/src/binding.ts index 20f5e86d009..87a5201ff31 100644 --- a/packages/workflows-shared/src/binding.ts +++ b/packages/workflows-shared/src/binding.ts @@ -213,7 +213,7 @@ export class WorkflowBinding extends WorkerEntrypoint { ); if (introspectionSession !== undefined) { - const modifier = stub.getInstanceModifier(); + using modifier = await stub.getInstanceModifier(); introspectionSession.instanceIds.push(id); for (const operation of introspectionSession.operations) { await applyWorkflowIntrospectionOperation(modifier, operation); diff --git a/packages/workflows-shared/src/context.ts b/packages/workflows-shared/src/context.ts index 76979d7540b..124c48d6014 100644 --- a/packages/workflows-shared/src/context.ts +++ b/packages/workflows-shared/src/context.ts @@ -697,22 +697,25 @@ export class Context extends RpcTarget { activeTimeoutTask?: Promise ): Promise => { if (!isReadableStreamLike(value)) { - // Typed-array views anywhere in the value tree are copied - // into a tight backing buffer so the full backing buffer - // does not ride along with each view (issue #14101). View - // types are preserved so cached replays observe the same - // constructor as the live execution path. The caller still - // receives the original `value` below — only the stored - // shape changes. - const stored = normalizeForStorage(value); - await this.#state.storage.put(valueKey, { value: stored }); - abortController.abort("step finished"); - // @ts-expect-error priorityQueue is initiated in init - this.#engine.priorityQueue.remove({ - hash: priorityQueueHash, - type: "timeout", - }); - return value; + try { + // Do not forward the callback's RPC disposer into the caller's + // execution context. Keep the live result's data shape intact. + const cloned = structuredClone(value); + // Compact typed-array backing buffers only for storage (#14101). + const stored = normalizeForStorage(cloned); + await this.#state.storage.put(valueKey, { value: stored }); + abortController.abort("step finished"); + // @ts-expect-error priorityQueue is initiated in init + this.#engine.priorityQueue.remove({ + hash: priorityQueueHash, + type: "timeout", + }); + return cloned; + } finally { + (value as Partial | null | undefined)?.[ + Symbol.dispose + ]?.(); + } } streamResultSeen = true; diff --git a/packages/workflows-shared/src/introspection.ts b/packages/workflows-shared/src/introspection.ts index ef337357c6a..f4f500e750f 100644 --- a/packages/workflows-shared/src/introspection.ts +++ b/packages/workflows-shared/src/introspection.ts @@ -220,6 +220,7 @@ export class WorkflowIntrospectorHandle implements WorkflowIntrospector { } export class WorkflowInstanceIntrospectorHandle implements WorkflowInstanceIntrospector { + #disposed = false; #instanceModifier: WorkflowInstanceModifier | undefined; #instanceModifierPromise: Promise | undefined; @@ -283,7 +284,18 @@ export class WorkflowInstanceIntrospectorHandle implements WorkflowInstanceIntro /** Keep this bound; explicit resource management may call the disposer unbound. */ dispose = async (): Promise => { - await this.workflow.unsafeAbort(this.instanceId, "Instance dispose"); + if (this.#disposed) { + return; + } + this.#disposed = true; + try { + // Acquisition starts in the constructor, even if modify() is never called. + const modifier = + this.#instanceModifier ?? (await this.#instanceModifierPromise); + (modifier as Partial | undefined)?.[Symbol.dispose]?.(); + } finally { + await this.workflow.unsafeAbort(this.instanceId, "Instance dispose"); + } }; async [Symbol.asyncDispose](): Promise { diff --git a/packages/workflows-shared/tests/rpc-lifecycle.test.ts b/packages/workflows-shared/tests/rpc-lifecycle.test.ts new file mode 100644 index 00000000000..78e1be5bbcd --- /dev/null +++ b/packages/workflows-shared/tests/rpc-lifecycle.test.ts @@ -0,0 +1,96 @@ +import { createExecutionContext } from "cloudflare:test"; +import { env } from "cloudflare:workers"; +import { afterEach, it, vi } from "vitest"; +import workerdUnsafe from "workerd:unsafe"; +import { WorkflowBinding } from "../src/binding"; +import { WorkflowInstanceIntrospectorHandle } from "../src/introspection"; +import { WorkflowInstanceModifier } from "../src/modifier"; +import { setTestWorkflowCallback } from "./test-entry"; +import { runWorkflowAndAwait } from "./utils"; +import type { WorkflowBinding as IntrospectionBinding } from "../src/types"; + +afterEach(async () => { + await workerdUnsafe.abortAllDurableObjects(); + Reflect.deleteProperty(WorkflowInstanceModifier.prototype, Symbol.dispose); +}); + +function createBinding(): WorkflowBinding { + return new WorkflowBinding(createExecutionContext(), { + ENGINE: env.ENGINE, + BINDING_NAME: "TEST_WORKFLOW", + WORKFLOW_NAME: "test-workflow", + }); +} + +it("releases the step callback's RPC result and preserves the live data shape", async ({ + expect, +}) => { + const dispose = vi.fn(); + let result: unknown; + await runWorkflowAndAwait(crypto.randomUUID(), async (_event, step) => { + result = await step.do("RPC result", async () => ({ + value: 42, + bytes: new Uint8Array(new ArrayBuffer(8), 2, 2), + [Symbol.dispose]: dispose, + })); + }); + + expect(result).toEqual({ value: 42, bytes: new Uint8Array(2) }); + expect(result).toHaveProperty("bytes.byteOffset", 2); + expect(result).toHaveProperty("bytes.buffer.byteLength", 8); + // The remote disposer runs asynchronously after the local stub is released. + await vi.waitFor(() => expect(dispose).toHaveBeenCalledOnce()); +}); + +it.for([false, true])( + "disposes an introspector's modifier (modified: %s)", + async (modified, { expect }) => { + const binding = createBinding(); + await using instance = new WorkflowInstanceIntrospectorHandle( + binding as unknown as IntrospectionBinding, + crypto.randomUUID() + ); + if (modified) { + await instance.modify((modifier) => modifier.disableSleeps()); + } + + await instance.dispose(); + await expect( + instance.modify((modifier) => modifier.disableSleeps()) + ).rejects.toThrow("RPC stub used after being disposed"); + // The await-using scope calls dispose() again, exercising idempotence. + } +); + +it.for([false, true])( + "releases the temporary modifier used by create (modified: %s)", + async (modified, { expect }) => { + const id = crypto.randomUUID(); + const binding = createBinding(); + const dispose = vi.fn(); + // RPC target disposers are looked up on the prototype. + Object.defineProperty(WorkflowInstanceModifier.prototype, Symbol.dispose, { + value: dispose, + configurable: true, + }); + + const sessionId = await binding.unsafeStartIntrospection(); + try { + if (modified) { + await binding.unsafeSetIntrospectionOperations(sessionId, [ + { type: "disableSleeps" }, + ]); + } + setTestWorkflowCallback(async () => undefined); + await binding.create({ id }); + const instance = await binding.get(id); + await vi.waitUntil( + async () => (await instance.status()).status === "complete" + ); + await vi.waitFor(() => expect(dispose).toHaveBeenCalledOnce()); + } finally { + await binding.unsafeStopIntrospection(sessionId); + await binding.unsafeAbort(id); + } + } +); From e14c104847b9150cbc891358033334bae0366763 Mon Sep 17 00:00:00 2001 From: Sam Starling <42478+samstarling@users.noreply.github.com> Date: Fri, 18 Sep 2026 13:42:36 +0100 Subject: [PATCH 2/4] fix(workflows): Preserve cleanup after modifier acquisition fails --- .changeset/fix-workflow-rpc-disposal.md | 2 + packages/workflows-shared/src/context.ts | 5 +- .../workflows-shared/src/introspection.ts | 4 +- .../tests/rpc-lifecycle.test.ts | 97 ++++++++++++++++++- 4 files changed, 104 insertions(+), 4 deletions(-) diff --git a/.changeset/fix-workflow-rpc-disposal.md b/.changeset/fix-workflow-rpc-disposal.md index a18800ee90f..8db27dcfe54 100644 --- a/.changeset/fix-workflow-rpc-disposal.md +++ b/.changeset/fix-workflow-rpc-disposal.md @@ -7,3 +7,5 @@ Dispose Workflow step results and introspection modifiers after use Release RPC resources deterministically during local Workflow execution and introspection. This prevents undisposed RPC warnings and requests being cancelled after their execution context has ended. Live step results retain their original data shape, including typed-array offsets and backing buffers. + +Disposing an introspector still aborts the instance if modifier acquisition failed, without introducing a teardown error. Non-stream results containing function-valued array properties are rejected according to the serialisable-output contract, instead of silently discarding those properties during storage normalisation. diff --git a/packages/workflows-shared/src/context.ts b/packages/workflows-shared/src/context.ts index 124c48d6014..e065ccc391c 100644 --- a/packages/workflows-shared/src/context.ts +++ b/packages/workflows-shared/src/context.ts @@ -698,8 +698,9 @@ export class Context extends RpcTarget { ): Promise => { if (!isReadableStreamLike(value)) { try { - // Do not forward the callback's RPC disposer into the caller's - // execution context. Keep the live result's data shape intact. + // Non-stream results must be structured-cloneable. Clone before + // normalisation, which can discard unsupported array properties, + // and avoid forwarding the callback's RPC disposer to the caller. const cloned = structuredClone(value); // Compact typed-array backing buffers only for storage (#14101). const stored = normalizeForStorage(cloned); diff --git a/packages/workflows-shared/src/introspection.ts b/packages/workflows-shared/src/introspection.ts index f4f500e750f..c439aa48121 100644 --- a/packages/workflows-shared/src/introspection.ts +++ b/packages/workflows-shared/src/introspection.ts @@ -290,8 +290,10 @@ export class WorkflowInstanceIntrospectorHandle implements WorkflowInstanceIntro this.#disposed = true; try { // Acquisition starts in the constructor, even if modify() is never called. + // If it failed, there is no modifier to dispose; modify() still reports it. const modifier = - this.#instanceModifier ?? (await this.#instanceModifierPromise); + this.#instanceModifier ?? + (await this.#instanceModifierPromise?.catch(() => undefined)); (modifier as Partial | undefined)?.[Symbol.dispose]?.(); } finally { await this.workflow.unsafeAbort(this.instanceId, "Instance dispose"); diff --git a/packages/workflows-shared/tests/rpc-lifecycle.test.ts b/packages/workflows-shared/tests/rpc-lifecycle.test.ts index 78e1be5bbcd..449b37afd1c 100644 --- a/packages/workflows-shared/tests/rpc-lifecycle.test.ts +++ b/packages/workflows-shared/tests/rpc-lifecycle.test.ts @@ -1,9 +1,11 @@ -import { createExecutionContext } from "cloudflare:test"; +import { createExecutionContext, runInDurableObject } from "cloudflare:test"; import { env } from "cloudflare:workers"; import { afterEach, it, vi } from "vitest"; import workerdUnsafe from "workerd:unsafe"; import { WorkflowBinding } from "../src/binding"; +import { InstanceEvent } from "../src/instance"; import { WorkflowInstanceIntrospectorHandle } from "../src/introspection"; +import { normalizeForStorage } from "../src/lib/serialization"; import { WorkflowInstanceModifier } from "../src/modifier"; import { setTestWorkflowCallback } from "./test-entry"; import { runWorkflowAndAwait } from "./utils"; @@ -94,3 +96,96 @@ it.for([false, true])( } } ); + +it.for([false, true])( + "aborts after modifier acquisition fails (modify attempted: %s)", + async (attemptModify, { expect }) => { + const binding = createBinding(); + const error = new Error("modifier acquisition failed"); + vi.spyOn(binding, "unsafeGetInstanceModifier").mockRejectedValue(error); + const abort = vi.spyOn(binding, "unsafeAbort"); + const id = crypto.randomUUID(); + { + await using instance = new WorkflowInstanceIntrospectorHandle( + binding as unknown as IntrospectionBinding, + id + ); + if (attemptModify) { + await expect( + instance.modify((modifier) => modifier.disableSleeps()) + ).rejects.toBe(error); + } + await expect(instance.dispose()).resolves.toBeUndefined(); + } + expect(abort).toHaveBeenCalledExactlyOnceWith(id, "Instance dispose"); + } +); + +it("still reports a modifier disposer failure after aborting", async ({ + expect, +}) => { + const binding = createBinding(); + const error = new Error("modifier disposal failed"); + vi.spyOn(binding, "unsafeGetInstanceModifier").mockResolvedValue({ + [Symbol.dispose]() { + throw error; + }, + }); + const abort = vi.spyOn(binding, "unsafeAbort"); + const id = crypto.randomUUID(); + const instance = new WorkflowInstanceIntrospectorHandle( + binding as unknown as IntrospectionBinding, + id + ); + await expect(instance.dispose()).rejects.toBe(error); + expect(abort).toHaveBeenCalledExactlyOnceWith(id, "Instance dispose"); +}); + +it.for([ + { name: "function-valued property", create: () => ({ helper: () => 42 }) }, + { + name: "function-valued array property", + create: () => Object.assign([42], { helper: () => 42 }), + }, +])( + "rejects a non-cloneable step result: $name", + async ({ create }, { expect }) => { + // The storage normaliser used to discard extra array properties. They must + // not let non-serialisable output bypass the step result contract. + const id = crypto.randomUUID(); + await runWorkflowAndAwait(id, async (_event, step) => { + await step.do("non-cloneable result", async () => create()); + }); + const engine = env.ENGINE.get(env.ENGINE.idFromName(id)); + const { logs } = await engine.readLogs(); + expect(logs).toContainEqual( + expect.objectContaining({ + event: InstanceEvent.WORKFLOW_FAILURE, + metadata: expect.objectContaining({ + error: expect.objectContaining({ + message: expect.stringContaining("not serialisable"), + }), + }), + }) + ); + } +); + +it("preserves data from class instances with prototype methods", async ({ + expect, +}) => { + class Result { + value = 42; + method(): number { + return this.value; + } + } + const stub = env.ENGINE.get(env.ENGINE.idFromName(crypto.randomUUID())); + await runInDurableObject(stub, async (_engine, state) => { + const value = new Result(); + // Both the old storage path and structuredClone ignore prototype methods. + await state.storage.put("result", { value: normalizeForStorage(value) }); + expect(await state.storage.get("result")).toEqual({ value: { value: 42 } }); + expect(structuredClone(value)).toEqual({ value: 42 }); + }); +}); From 1fa952f09eed6c2f5a50698ea64948f327a72448 Mon Sep 17 00:00:00 2001 From: Sam Starling <42478+samstarling@users.noreply.github.com> Date: Fri, 18 Sep 2026 14:06:17 +0100 Subject: [PATCH 3/4] fix(workflows): Dispose callback results that arrive after timeout --- .changeset/fix-workflow-rpc-disposal.md | 2 +- packages/workflows-shared/src/context.ts | 30 ++++- .../tests/rpc-lifecycle.test.ts | 107 ++++++++++++++++++ 3 files changed, 134 insertions(+), 5 deletions(-) diff --git a/.changeset/fix-workflow-rpc-disposal.md b/.changeset/fix-workflow-rpc-disposal.md index 8db27dcfe54..7a255bf6ce8 100644 --- a/.changeset/fix-workflow-rpc-disposal.md +++ b/.changeset/fix-workflow-rpc-disposal.md @@ -6,6 +6,6 @@ Dispose Workflow step results and introspection modifiers after use -Release RPC resources deterministically during local Workflow execution and introspection. This prevents undisposed RPC warnings and requests being cancelled after their execution context has ended. Live step results retain their original data shape, including typed-array offsets and backing buffers. +Release RPC resources deterministically during local Workflow execution and introspection. This prevents undisposed RPC warnings and requests being cancelled after their execution context has ended. Live step results retain their original data shape, including typed-array offsets and backing buffers. Callbacks that finish after a step times out also release their results, cancelling unused streams without delaying retries. Disposing an introspector still aborts the instance if modifier acquisition failed, without introducing a teardown error. Non-stream results containing function-valued array properties are rejected according to the serialisable-output contract, instead of silently discarding those properties during storage normalisation. diff --git a/packages/workflows-shared/src/context.ts b/packages/workflows-shared/src/context.ts index e065ccc391c..9ae354728ef 100644 --- a/packages/workflows-shared/src/context.ts +++ b/packages/workflows-shared/src/context.ts @@ -767,14 +767,36 @@ export class Context extends RpcTarget { } } else { timeoutTask = timeoutPromise(); - result = await Promise.race([ + const callbackTask = Promise.resolve( doWrapperClosure({ step: { name, count }, attempt: stepState.attemptedCount, config: toEngineStepConfig(config), - }), - timeoutTask, - ]); + }) + ); + try { + result = await Promise.race([callbackTask, timeoutTask]); + } catch (error) { + // A timeout does not cancel the callback RPC. Release any late + // result without waiting for it before retrying. Successful race + // winners remain owned by persistStepResult, including streams. + void callbackTask + .then(async (value) => { + try { + if (isReadableStreamLike(value)) { + await value.cancel(error); + } + } finally { + (value as Partial | null | undefined)?.[ + Symbol.dispose + ]?.(); + } + }) + // Late rejection or cleanup failure must not replace the + // attempt's original error or become an unhandled rejection. + .catch(() => {}); + throw error; + } } // We store the value of `output` in an object with a `value` property. This allows us to store `undefined`, diff --git a/packages/workflows-shared/tests/rpc-lifecycle.test.ts b/packages/workflows-shared/tests/rpc-lifecycle.test.ts index 449b37afd1c..0e7ae9afe69 100644 --- a/packages/workflows-shared/tests/rpc-lifecycle.test.ts +++ b/packages/workflows-shared/tests/rpc-lifecycle.test.ts @@ -13,6 +13,7 @@ import type { WorkflowBinding as IntrospectionBinding } from "../src/types"; afterEach(async () => { await workerdUnsafe.abortAllDurableObjects(); + vi.restoreAllMocks(); Reflect.deleteProperty(WorkflowInstanceModifier.prototype, Symbol.dispose); }); @@ -189,3 +190,109 @@ it("preserves data from class instances with prototype methods", async ({ expect(structuredClone(value)).toEqual({ value: 42 }); }); }); + +it.for(["object", "stream", "reject", "cancel-reject"] as const)( + "cleans up a late callback after a timeout: %s", + async (kind, { expect }) => { + const lateCleanup = vi.fn(); + const originalCancel = ReadableStream.prototype.cancel; + const cancel = vi.spyOn(ReadableStream.prototype, "cancel"); + if (kind === "cancel-reject") { + cancel.mockImplementation(async function (reason) { + await originalCancel.call(this, reason); + throw new Error("late stream cancellation failed"); + }); + } + const winningDispose = vi.fn(); + const winningCancel = vi.fn(); + const attempts: number[] = []; + const lateResult = Promise.withResolvers(); + let lateSettled = false; + let cleanupSeenByRetry = false; + let result: unknown; + await runWorkflowAndAwait(crypto.randomUUID(), async (_event, step) => { + const value = await step.do< + { attempt: number } | ReadableStream + >( + "late callback result", + { + timeout: "1 second", + retries: { limit: 1, delay: 1, backoff: "constant" }, + }, + async (ctx) => { + attempts.push(ctx.attempt); + if (ctx.attempt === 1) { + // Only the retry can release this callback: cleanup must not + // delay the retry until the abandoned callback finishes. + await lateResult.promise; + lateSettled = true; + if (kind === "reject") { + throw new Error("late callback failed"); + } + if (kind === "object") { + return { attempt: 1, [Symbol.dispose]: lateCleanup }; + } + // Observe cancellation on the receiving stream. Workerd does + // not reliably notify the source's cancel hook over RPC. + return new ReadableStream(); + } + lateResult.resolve(); + cleanupSeenByRetry = await vi + .waitUntil( + () => + kind === "reject" + ? lateSettled + : kind === "object" + ? lateCleanup.mock.calls.length === 1 + : cancel.mock.settledResults.length === 1, + { timeout: 500 } + ) + .then( + () => true, + () => false + ); + if (kind === "stream" || kind === "cancel-reject") { + return new ReadableStream({ + start(controller) { + controller.enqueue(new TextEncoder().encode("retry succeeded")); + controller.close(); + }, + cancel: winningCancel, + }); + } + return { attempt: 2, [Symbol.dispose]: winningDispose }; + } + ); + result = + value instanceof ReadableStream + ? await new Response(value).text() + : value; + }); + + expect(attempts).toEqual([1, 2]); + expect(lateSettled).toBe(true); + expect(cleanupSeenByRetry).toBe(true); + if (kind === "stream" || kind === "cancel-reject") { + expect(result).toBe("retry succeeded"); + expect(winningCancel).not.toHaveBeenCalled(); + } else { + expect(result).toEqual({ attempt: 2 }); + await vi.waitFor(() => expect(winningDispose).toHaveBeenCalledOnce()); + } + if (kind === "object") { + expect(lateCleanup).toHaveBeenCalledOnce(); + } + if (kind === "stream" || kind === "cancel-reject") { + expect(cancel).toHaveBeenCalledExactlyOnceWith( + expect.objectContaining({ + name: "WorkflowTimeoutError", + }) + ); + expect(cancel.mock.settledResults[0]?.type).toBe( + kind === "stream" ? "fulfilled" : "rejected" + ); + } else { + expect(cancel).not.toHaveBeenCalled(); + } + } +); From d68cb6fe6a19054fc7e0287911af700151f13aed Mon Sep 17 00:00:00 2001 From: Sam Starling <42478+samstarling@users.noreply.github.com> Date: Fri, 18 Sep 2026 14:57:18 +0100 Subject: [PATCH 4/4] fix(workflows): Cover late callback cleanup across a step boundary The cleanup for a callback abandoned by a step timeout runs from an untracked promise. The engine is a Durable Object, which stays active while there is ongoing work or pending I/O, so the cleanup is retained without a waitUntil (a documented no-op for Durable Objects). Add a regression test that releases the abandoned callback from a later step, and record the lifetime reasoning at the call site. Co-Authored-By: Claude Opus 5 (1M context) --- packages/workflows-shared/src/context.ts | 5 +++ .../tests/rpc-lifecycle.test.ts | 45 +++++++++++++++++++ 2 files changed, 50 insertions(+) diff --git a/packages/workflows-shared/src/context.ts b/packages/workflows-shared/src/context.ts index 9ae354728ef..ed3dbcbe7f0 100644 --- a/packages/workflows-shared/src/context.ts +++ b/packages/workflows-shared/src/context.ts @@ -780,6 +780,11 @@ export class Context extends RpcTarget { // A timeout does not cancel the callback RPC. Release any late // result without waiting for it before retrying. Successful race // winners remain owned by persistStepResult, including streams. + // The engine is a Durable Object, which stays active while there + // is ongoing work or pending I/O, so this untracked cleanup is + // retained without a waitUntil (a documented no-op for Durable + // Objects). Once the run() call that owns the callback stub ends, + // the stub is torn down with it and there is no result to release. void callbackTask .then(async (value) => { try { diff --git a/packages/workflows-shared/tests/rpc-lifecycle.test.ts b/packages/workflows-shared/tests/rpc-lifecycle.test.ts index 0e7ae9afe69..d983fb42f5c 100644 --- a/packages/workflows-shared/tests/rpc-lifecycle.test.ts +++ b/packages/workflows-shared/tests/rpc-lifecycle.test.ts @@ -296,3 +296,48 @@ it.for(["object", "stream", "reject", "cancel-reject"] as const)( } } ); + +it("cleans up a late callback that lands after its step has finished", async ({ + expect, +}) => { + const lateCleanup = vi.fn(); + const attempts: number[] = []; + const lateResult = Promise.withResolvers(); + let cleanupSeenByLaterStep = false; + + await runWorkflowAndAwait(crypto.randomUUID(), async (_event, step) => { + await step.do( + "late callback result", + { + timeout: "1 second", + retries: { limit: 1, delay: 1, backoff: "constant" }, + }, + async (ctx) => { + attempts.push(ctx.attempt); + if (ctx.attempt === 1) { + // Held open past the retry, so the abandoned callback's cleanup is + // left running from the engine with no step awaiting it. + await lateResult.promise; + return { attempt: 1, [Symbol.dispose]: lateCleanup }; + } + return { attempt: 2 }; + } + ); + + // The engine keeps the abandoned cleanup alive across the step boundary: + // releasing the first callback here still disposes its result. + await step.do("later step", async () => { + lateResult.resolve(); + cleanupSeenByLaterStep = await vi + .waitUntil(() => lateCleanup.mock.calls.length === 1, { timeout: 500 }) + .then( + () => true, + () => false + ); + }); + }); + + expect(attempts).toEqual([1, 2]); + expect(cleanupSeenByLaterStep).toBe(true); + expect(lateCleanup).toHaveBeenCalledOnce(); +});