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
11 changes: 11 additions & 0 deletions .changeset/fix-workflow-rpc-disposal.md
Original file line number Diff line number Diff line change
@@ -0,0 +1,11 @@
---
"@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. 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.
2 changes: 1 addition & 1 deletion packages/workflows-shared/src/binding.ts
Original file line number Diff line number Diff line change
Expand Up @@ -213,7 +213,7 @@ export class WorkflowBinding extends WorkerEntrypoint<Env> {
);

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);
Expand Down
71 changes: 51 additions & 20 deletions packages/workflows-shared/src/context.ts
Original file line number Diff line number Diff line change
Expand Up @@ -697,22 +697,26 @@ export class Context extends RpcTarget {
activeTimeoutTask?: Promise<never>
): Promise<unknown> => {
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 {
// 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);
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<Disposable> | null | undefined)?.[
Symbol.dispose
]?.();
Comment thread
samstarling marked this conversation as resolved.
}
}

streamResultSeen = true;
Expand Down Expand Up @@ -763,14 +767,41 @@ 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.
// 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 {
if (isReadableStreamLike(value)) {
await value.cancel(error);
}
} finally {
(value as Partial<Disposable> | null | undefined)?.[
Symbol.dispose
]?.();
}
})
// Late rejection or cleanup failure must not replace the
// attempt's original error or become an unhandled rejection.
.catch(() => {});
Comment thread
devin-ai-integration[bot] marked this conversation as resolved.
throw error;
}
}

// We store the value of `output` in an object with a `value` property. This allows us to store `undefined`,
Expand Down
16 changes: 15 additions & 1 deletion packages/workflows-shared/src/introspection.ts
Original file line number Diff line number Diff line change
Expand Up @@ -220,6 +220,7 @@ export class WorkflowIntrospectorHandle implements WorkflowIntrospector {
}

export class WorkflowInstanceIntrospectorHandle implements WorkflowInstanceIntrospector {
#disposed = false;
#instanceModifier: WorkflowInstanceModifier | undefined;
#instanceModifierPromise: Promise<WorkflowInstanceModifier> | undefined;

Expand Down Expand Up @@ -283,7 +284,20 @@ export class WorkflowInstanceIntrospectorHandle implements WorkflowInstanceIntro

/** Keep this bound; explicit resource management may call the disposer unbound. */
dispose = async (): Promise<void> => {
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.
// If it failed, there is no modifier to dispose; modify() still reports it.
const modifier =
this.#instanceModifier ??
(await this.#instanceModifierPromise?.catch(() => undefined));
(modifier as Partial<Disposable> | undefined)?.[Symbol.dispose]?.();
} finally {
await this.workflow.unsafeAbort(this.instanceId, "Instance dispose");
}
};

async [Symbol.asyncDispose](): Promise<void> {
Expand Down
Loading
Loading