From c187e2421a187e06f1b66d7ef3cd10a8466233d5 Mon Sep 17 00:00:00 2001 From: Omer Aplak Date: Sat, 21 Feb 2026 18:38:46 -0800 Subject: [PATCH 1/7] feat(core): add workflow time-travel replay APIs --- .changeset/four-cobras-yawn.md | 23 ++ packages/core/src/memory/types.ts | 4 + packages/core/src/workflow/chain.ts | 35 ++ packages/core/src/workflow/core.ts | 379 ++++++++++++++++++ packages/core/src/workflow/index.ts | 1 + .../core/src/workflow/time-travel.spec.ts | 250 ++++++++++++ packages/core/src/workflow/types.ts | 37 ++ website/docs/workflows/overview.md | 47 +++ website/docs/workflows/streaming.md | 34 +- website/docs/workflows/suspend-resume.md | 70 ++++ 10 files changed, 879 insertions(+), 1 deletion(-) create mode 100644 .changeset/four-cobras-yawn.md create mode 100644 packages/core/src/workflow/time-travel.spec.ts diff --git a/.changeset/four-cobras-yawn.md b/.changeset/four-cobras-yawn.md new file mode 100644 index 000000000..7640fc94c --- /dev/null +++ b/.changeset/four-cobras-yawn.md @@ -0,0 +1,23 @@ +--- +"@voltagent/core": minor +--- + +Add workflow time-travel and deterministic replay APIs. + +New APIs: + +- `workflow.timeTravel(options)` +- `workflow.timeTravelStream(options)` +- `workflowChain.timeTravel(options)` +- `workflowChain.timeTravelStream(options)` + +`timeTravel` replays a historical execution from a selected step with a new execution ID, preserving the original execution history. Replay runs can optionally override selected-step input (`inputData`), resume payload (`resumeData`), and shared workflow state (`workflowStateOverride`). + +Replay lineage metadata is now persisted on workflow state records: + +- `replayedFromExecutionId` +- `replayFromStepId` + +New public type exports from `@voltagent/core` include `WorkflowTimeTravelOptions`. + +Also adds workflow documentation and usage examples for deterministic replay in overview, suspend/resume, and streaming docs. diff --git a/packages/core/src/memory/types.ts b/packages/core/src/memory/types.ts index bc502de3b..d73d09bce 100644 --- a/packages/core/src/memory/types.ts +++ b/packages/core/src/memory/types.ts @@ -189,6 +189,10 @@ export interface WorkflowStateEntry { userId?: string; /** Conversation ID if applicable */ conversationId?: string; + /** Source execution ID if this run is a replay */ + replayedFromExecutionId?: string; + /** Source step ID used when this run was replayed */ + replayFromStepId?: string; /** Additional metadata */ metadata?: Record; /** Timestamps */ diff --git a/packages/core/src/workflow/chain.ts b/packages/core/src/workflow/chain.ts index c18aad03b..5fd27a5b2 100644 --- a/packages/core/src/workflow/chain.ts +++ b/packages/core/src/workflow/chain.ts @@ -58,6 +58,7 @@ import type { WorkflowStepState, WorkflowStreamResult, WorkflowStreamWriter, + WorkflowTimeTravelOptions, } from "./types"; export type { AgentConfig } from "./steps/and-agent"; @@ -987,6 +988,40 @@ export class WorkflowChain< return workflow.startAsync(input, options); } + /** + * Replay a historical execution from the selected step + */ + async timeTravel( + options: WorkflowTimeTravelOptions, + ): Promise> { + const workflow = createWorkflow( + this.config, + // @ts-expect-error - upstream types work and this is nature of how the createWorkflow function is typed using variadic args + ...this.steps, + ); + return (await workflow.timeTravel(options)) as unknown as WorkflowExecutionResult< + RESULT_SCHEMA, + RESUME_SCHEMA + >; + } + + /** + * Stream a historical replay from the selected step + */ + timeTravelStream( + options: WorkflowTimeTravelOptions, + ): WorkflowStreamResult { + const workflow = createWorkflow( + this.config, + // @ts-expect-error - upstream types work and this is nature of how the createWorkflow function is typed using variadic args + ...this.steps, + ); + return workflow.timeTravelStream(options) as unknown as WorkflowStreamResult< + RESULT_SCHEMA, + RESUME_SCHEMA + >; + } + /** * Restart an interrupted execution from persisted checkpoint state * This recreates a workflow instance via `createWorkflow(...)` on each call. diff --git a/packages/core/src/workflow/core.ts b/packages/core/src/workflow/core.ts index 511cdc107..1c023be59 100644 --- a/packages/core/src/workflow/core.ts +++ b/packages/core/src/workflow/core.ts @@ -54,6 +54,7 @@ import type { WorkflowStepData, WorkflowStreamResult, WorkflowSuspensionMetadata, + WorkflowTimeTravelOptions, } from "./types"; export const VOLTAGENT_RESTART_CHECKPOINT_KEY = "__voltagent_restart_checkpoint"; @@ -118,6 +119,13 @@ const isWorkflowStepStatus = (value: unknown): value is WorkflowStepData["status value === "cancelled" || value === "skipped"; +const toWorkflowStepStatus = (value: unknown): WorkflowStepData["status"] => { + if (isWorkflowStepStatus(value)) { + return value; + } + return "success"; +}; + const deserializeCheckpointStepData = (value: unknown): WorkflowStepData | undefined => { if (!isObjectRecord(value) || !isWorkflowStepStatus(value.status)) { return undefined; @@ -177,6 +185,23 @@ const toValidContextMap = (context: unknown): Map | un return new Map(entries); }; +const toContextEntries = ( + context: Map | undefined, +): Array<[string | symbol, unknown]> | undefined => + context ? Array.from(context.entries()) : undefined; + +const withoutRestartCheckpointMetadata = ( + metadata?: Record, +): Record | undefined => { + if (!metadata) { + return undefined; + } + + const nextMetadata = { ...metadata }; + delete nextMetadata[VOLTAGENT_RESTART_CHECKPOINT_KEY]; + return nextMetadata; +}; + const getRestartCheckpointFromMetadata = ( metadata: Record | undefined, ): WorkflowRestartCheckpoint | undefined => { @@ -2493,6 +2518,193 @@ export function createWorkflow< }; }; + type PreparedTimeTravelExecution = { + executionId: string; + startAt: Date; + workflowInput: WorkflowInput; + executionOptions: WorkflowRunOptions; + }; + + const prepareTimeTravelExecution = async ( + timeTravelOptions: WorkflowTimeTravelOptions, + replayExecutionId: string = randomUUID(), + replayStartAt: Date = new Date(), + ): Promise => { + const sourceState = await defaultMemory.getWorkflowState(timeTravelOptions.executionId); + if (!sourceState) { + throw new Error(`Workflow state not found: ${timeTravelOptions.executionId}`); + } + + if (sourceState.workflowId !== id) { + throw new Error( + `Execution ${timeTravelOptions.executionId} belongs to workflow '${sourceState.workflowId}', expected '${id}'`, + ); + } + + if (sourceState.status === "running") { + throw new Error( + `Execution ${timeTravelOptions.executionId} is still running. Use restart() for crash recovery or wait for completion before time travel.`, + ); + } + + const targetStepIndex = (steps as BaseStep[]).findIndex( + (step) => step.id === timeTravelOptions.stepId, + ); + if (targetStepIndex === -1) { + throw new Error(`Step '${timeTravelOptions.stepId}' not found in workflow '${id}'`); + } + + const workflowStartEventInput = sourceState.events?.find( + (event) => event.type === "workflow-start", + )?.input; + const sourceWorkflowInput = sourceState.input ?? workflowStartEventInput; + if (sourceWorkflowInput === undefined) { + throw new Error( + `Cannot time travel execution ${timeTravelOptions.executionId}: missing persisted workflow input`, + ); + } + + const sourceCheckpoint = getRestartCheckpointFromMetadata(sourceState.metadata); + const sourceStepData = sourceCheckpoint?.stepData ?? {}; + const replayStepData: Record = {}; + + const sourceStepCompleteEvents = + sourceState.events?.filter((event) => event.type === "step-complete") ?? []; + + for (let index = 0; index < targetStepIndex; index += 1) { + const step = (steps as BaseStep[])[index]; + const checkpointSnapshot = sourceStepData[step.id]; + if (checkpointSnapshot) { + replayStepData[step.id] = { + input: checkpointSnapshot.input, + output: checkpointSnapshot.output, + status: toWorkflowStepStatus(checkpointSnapshot.status), + error: serializeStepError(checkpointSnapshot.error), + }; + continue; + } + + const fallbackEvent = sourceStepCompleteEvents.find( + (event) => + event.from === step.id || + event.name === step.id || + event.from === step.name || + event.name === step.name, + ); + + if (fallbackEvent) { + replayStepData[step.id] = { + input: fallbackEvent.input, + output: fallbackEvent.output, + status: toWorkflowStepStatus(fallbackEvent.status), + error: null, + }; + } + } + + const missingHistoricalSteps = (steps as BaseStep[]) + .slice(0, targetStepIndex) + .map((step) => step.id) + .filter((stepId) => replayStepData[stepId] === undefined); + if (missingHistoricalSteps.length > 0) { + throw new Error( + `Cannot time travel from step '${timeTravelOptions.stepId}': missing historical snapshots for steps ${missingHistoricalSteps.join(", ")}`, + ); + } + + const previousStepOutput = + targetStepIndex > 0 + ? replayStepData[(steps as BaseStep[])[targetStepIndex - 1]?.id]?.output + : undefined; + const sourceTargetStepInput = sourceStepData[timeTravelOptions.stepId]?.input; + const checkpointInputFallback = + sourceCheckpoint?.resumeStepIndex === targetStepIndex + ? sourceCheckpoint.stepExecutionState + : undefined; + + const replayStepInput = + timeTravelOptions.inputData ?? + sourceTargetStepInput ?? + previousStepOutput ?? + (targetStepIndex === 0 ? sourceWorkflowInput : checkpointInputFallback); + + if (replayStepInput === undefined) { + throw new Error( + `Cannot time travel from step '${timeTravelOptions.stepId}': missing historical input data (provide inputData override).`, + ); + } + + const effectiveWorkflowState = + timeTravelOptions.workflowStateOverride ?? + sourceCheckpoint?.workflowState ?? + sourceState.workflowState ?? + {}; + + const sourceContext = toValidContextMap(sourceState.context); + const lineageMetadata = { + ...(withoutRestartCheckpointMetadata(sourceState.metadata) ?? {}), + replayedFromExecutionId: timeTravelOptions.executionId, + replayFromStepId: timeTravelOptions.stepId, + replayedAt: replayStartAt.toISOString(), + }; + + await defaultMemory.setWorkflowState(replayExecutionId, { + id: replayExecutionId, + workflowId: id, + workflowName: name, + status: "running", + input: sourceWorkflowInput, + context: toContextEntries(sourceContext), + workflowState: effectiveWorkflowState, + userId: sourceState.userId, + conversationId: sourceState.conversationId, + replayedFromExecutionId: timeTravelOptions.executionId, + replayFromStepId: timeTravelOptions.stepId, + metadata: lineageMetadata, + createdAt: replayStartAt, + updatedAt: replayStartAt, + }); + + const completedStepsData = (steps as BaseStep[]) + .slice(0, targetStepIndex) + .map((step, stepIndex) => ({ + stepId: step.id, + stepName: step.name ?? step.id, + stepIndex, + output: replayStepData[step.id]?.output, + status: replayStepData[step.id]?.status, + })); + + const executionOptions: WorkflowRunOptions = { + executionId: replayExecutionId, + userId: sourceState.userId, + conversationId: sourceState.conversationId, + context: sourceContext, + workflowState: effectiveWorkflowState, + metadata: lineageMetadata, + resumeFrom: { + executionId: replayExecutionId, + resumeStepIndex: targetStepIndex, + lastEventSequence: sourceCheckpoint?.eventSequence, + resumeData: timeTravelOptions.resumeData, + checkpoint: { + stepExecutionState: replayStepInput, + completedStepsData, + workflowState: effectiveWorkflowState, + stepData: replayStepData, + usage: sourceCheckpoint?.usage, + }, + }, + }; + + return { + executionId: replayExecutionId, + startAt: replayStartAt, + workflowInput: sourceWorkflowInput as WorkflowInput, + executionOptions, + }; + }; + const workflow: Workflow & { __setDefaultMemory?: (memory: MemoryV2) => void; } = { @@ -2653,6 +2865,173 @@ export function createWorkflow< startAt, }; }, + timeTravel: async ( + timeTravelOptions: WorkflowTimeTravelOptions, + ): Promise> => { + const preparedReplay = await prepareTimeTravelExecution(timeTravelOptions); + return executeInternal(preparedReplay.workflowInput, preparedReplay.executionOptions); + }, + timeTravelStream: (timeTravelOptions: WorkflowTimeTravelOptions) => { + const streamController = new WorkflowStreamController(); + const executionId = randomUUID(); + const startAt = new Date(); + const suspendController = createDefaultSuspendController(); + + let replayOriginalInput: WorkflowInput | undefined; + + const replayPromise = (async () => { + const preparedReplay = await prepareTimeTravelExecution( + timeTravelOptions, + executionId, + startAt, + ); + replayOriginalInput = preparedReplay.workflowInput; + const replayExecutionOptions: WorkflowRunOptions = { + ...preparedReplay.executionOptions, + suspendController, + }; + return executeInternal( + preparedReplay.workflowInput, + replayExecutionOptions, + streamController, + ); + })(); + + replayPromise + .then( + (result) => { + if (result.status !== "suspended") { + streamController.close(); + } + }, + () => { + streamController.close(); + }, + ) + .catch(() => { + // Error is surfaced through promise-backed fields on stream result. + }); + + const streamResult: WorkflowStreamResult = { + executionId, + workflowId: id, + startAt, + endAt: replayPromise.then((r) => r.endAt), + status: replayPromise.then((r) => r.status), + result: replayPromise.then((r) => r.result), + suspension: replayPromise.then((r) => r.suspension), + cancellation: replayPromise.then((r) => r.cancellation), + error: replayPromise.then((r) => r.error), + usage: replayPromise.then((r) => r.usage), + toUIMessageStreamResponse: eventToUIMessageStreamResponse(streamController), + resume: async (input: z.infer, opts?: { stepId?: string }) => { + const replayResult = await replayPromise; + if (replayResult.status !== "suspended") { + throw new Error(`Cannot resume workflow in ${replayResult.status} state`); + } + + if (!replayResult.suspension) { + throw new Error("No suspension metadata found"); + } + + if (!replayOriginalInput) { + throw new Error("Missing replay input for resume"); + } + + let resumeStepIndex = replayResult.suspension.suspendedStepIndex; + if (opts?.stepId) { + const overrideIndex = (steps as BaseStep[]).findIndex( + (step) => step.id === opts.stepId, + ); + if (overrideIndex === -1) { + throw new Error(`Step '${opts.stepId}' not found in workflow '${id}'`); + } + resumeStepIndex = overrideIndex; + } + + let resumedResolve: ( + value: WorkflowExecutionResult, + ) => void; + let resumedReject: (error: any) => void; + const resumedPromise = new Promise>( + (resolve, reject) => { + resumedResolve = resolve; + resumedReject = reject; + }, + ); + + const resumedSuspendController = createDefaultSuspendController(); + const resumeOptions: WorkflowRunOptions = { + executionId: replayResult.executionId, + resumeFrom: { + executionId: replayResult.executionId, + checkpoint: replayResult.suspension.checkpoint, + resumeStepIndex, + resumeData: input, + }, + suspendController: resumedSuspendController, + }; + + executeInternal(replayOriginalInput, resumeOptions, streamController) + .then( + (result) => { + if (result.status !== "suspended") { + streamController.close(); + } + resumedResolve(result); + }, + (error) => { + streamController.close(); + resumedReject(error); + }, + ) + .catch(() => {}); + + const resumedStreamResult: WorkflowStreamResult = { + executionId: replayResult.executionId, + workflowId: replayResult.workflowId, + startAt: replayResult.startAt, + endAt: resumedPromise.then((r) => r.endAt), + status: resumedPromise.then((r) => r.status), + result: resumedPromise.then((r) => r.result), + suspension: resumedPromise.then((r) => r.suspension), + cancellation: resumedPromise.then((r) => r.cancellation), + error: resumedPromise.then((r) => r.error), + usage: resumedPromise.then((r) => r.usage), + resume: async (input2: z.infer, opts2?: { stepId?: string }) => { + const nextResult = await resumedPromise; + if (nextResult.status !== "suspended") { + throw new Error(`Cannot resume workflow in ${nextResult.status} state`); + } + return streamResult.resume(input2, opts2); + }, + suspend: (reason?: string) => { + resumedSuspendController.suspend(reason); + }, + cancel: (reason?: string) => { + resumedSuspendController.cancel(reason); + }, + abort: () => streamController.abort(), + toUIMessageStreamResponse: eventToUIMessageStreamResponse(streamController), + [Symbol.asyncIterator]: () => streamController.getStream(), + }; + + return resumedStreamResult; + }, + suspend: (reason?: string) => { + suspendController.suspend(reason); + }, + cancel: (reason?: string) => { + suspendController.cancel(reason); + }, + abort: () => { + streamController.abort(); + }, + [Symbol.asyncIterator]: () => streamController.getStream(), + }; + + return streamResult; + }, restart: (executionId: string, options?: WorkflowRunOptions) => { return restartExecution(executionId, options); }, diff --git a/packages/core/src/workflow/index.ts b/packages/core/src/workflow/index.ts index 4077146d9..0bd0cd9d3 100644 --- a/packages/core/src/workflow/index.ts +++ b/packages/core/src/workflow/index.ts @@ -33,6 +33,7 @@ export type { WorkflowRunOptions, WorkflowResumeOptions, WorkflowStartAsyncResult, + WorkflowTimeTravelOptions, WorkflowSuspensionMetadata, WorkflowSuspendController, WorkflowStateStore, diff --git a/packages/core/src/workflow/time-travel.spec.ts b/packages/core/src/workflow/time-travel.spec.ts new file mode 100644 index 000000000..caad64cb9 --- /dev/null +++ b/packages/core/src/workflow/time-travel.spec.ts @@ -0,0 +1,250 @@ +import { beforeEach, describe, expect, it } from "vitest"; +import { z } from "zod"; +import { Memory } from "../memory"; +import { InMemoryStorageAdapter } from "../memory/adapters/storage/in-memory"; +import { createWorkflow } from "./core"; +import { WorkflowRegistry } from "./registry"; +import { andThen } from "./steps"; + +describe.sequential("workflow.timeTravel", () => { + beforeEach(() => { + const registry = WorkflowRegistry.getInstance(); + (registry as any).workflows.clear(); + }); + + it("should replay from middle step with a new execution id", async () => { + const memory = new Memory({ storage: new InMemoryStorageAdapter() }); + const counters = { step1: 0, step2: 0, step3: 0 }; + + const workflow = createWorkflow( + { + id: "time-travel-middle-step", + name: "Time Travel Middle Step", + input: z.object({ value: z.number() }), + result: z.object({ result: z.number() }), + memory, + }, + andThen({ + id: "step-1", + execute: async ({ data }) => { + counters.step1 += 1; + return { value: data.value + 1 }; + }, + }), + andThen({ + id: "step-2", + execute: async ({ data }) => { + counters.step2 += 1; + return { value: data.value * 2 }; + }, + }), + andThen({ + id: "step-3", + execute: async ({ data }) => { + counters.step3 += 1; + return { result: data.value + 3 }; + }, + }), + ); + + const registry = WorkflowRegistry.getInstance(); + registry.registerWorkflow(workflow); + + const original = await workflow.run({ value: 1 }); + expect(original.status).toBe("completed"); + expect(original.result).toEqual({ result: 7 }); + expect(counters).toEqual({ step1: 1, step2: 1, step3: 1 }); + + const replay = await workflow.timeTravel({ + executionId: original.executionId, + stepId: "step-2", + }); + + expect(replay.status).toBe("completed"); + expect(replay.result).toEqual({ result: 7 }); + expect(replay.executionId).not.toBe(original.executionId); + expect(counters).toEqual({ step1: 1, step2: 2, step3: 2 }); + + const originalState = await memory.getWorkflowState(original.executionId); + const replayState = await memory.getWorkflowState(replay.executionId); + + expect(originalState?.status).toBe("completed"); + expect(originalState?.output).toEqual({ result: 7 }); + expect(originalState?.replayedFromExecutionId).toBeUndefined(); + expect(originalState?.replayFromStepId).toBeUndefined(); + + expect(replayState?.status).toBe("completed"); + expect(replayState?.replayedFromExecutionId).toBe(original.executionId); + expect(replayState?.replayFromStepId).toBe("step-2"); + }); + + it("should apply inputData override only to replay run", async () => { + const memory = new Memory({ storage: new InMemoryStorageAdapter() }); + + const workflow = createWorkflow( + { + id: "time-travel-input-override", + name: "Time Travel Input Override", + input: z.object({ value: z.number() }), + result: z.object({ result: z.number() }), + memory, + }, + andThen({ + id: "step-1", + execute: async ({ data }) => ({ value: data.value + 1 }), + }), + andThen({ + id: "step-2", + execute: async ({ data }) => ({ value: data.value * 2 }), + }), + andThen({ + id: "step-3", + execute: async ({ data }) => ({ result: data.value + 3 }), + }), + ); + + const registry = WorkflowRegistry.getInstance(); + registry.registerWorkflow(workflow); + + const original = await workflow.run({ value: 1 }); + expect(original.result).toEqual({ result: 7 }); + + const replay = await workflow.timeTravel({ + executionId: original.executionId, + stepId: "step-2", + inputData: { value: 100 }, + }); + + expect(replay.status).toBe("completed"); + expect(replay.result).toEqual({ result: 203 }); + + const originalState = await memory.getWorkflowState(original.executionId); + expect(originalState?.output).toEqual({ result: 7 }); + }); + + it("should fail with actionable error when step does not exist", async () => { + const memory = new Memory({ storage: new InMemoryStorageAdapter() }); + + const workflow = createWorkflow( + { + id: "time-travel-invalid-step", + name: "Time Travel Invalid Step", + input: z.object({ value: z.number() }), + result: z.object({ value: z.number() }), + memory, + }, + andThen({ + id: "step-1", + execute: async ({ data }) => data, + }), + ); + + const registry = WorkflowRegistry.getInstance(); + registry.registerWorkflow(workflow); + + const original = await workflow.run({ value: 1 }); + await expect( + workflow.timeTravel({ + executionId: original.executionId, + stepId: "missing-step", + }), + ).rejects.toThrow("Step 'missing-step' not found"); + }); + + it("should keep original execution history unchanged after replay", async () => { + const memory = new Memory({ storage: new InMemoryStorageAdapter() }); + + const workflow = createWorkflow( + { + id: "time-travel-history-integrity", + name: "Time Travel History Integrity", + input: z.object({ value: z.number() }), + result: z.object({ result: z.number() }), + memory, + }, + andThen({ + id: "step-1", + execute: async ({ data }) => ({ value: data.value + 1 }), + }), + andThen({ + id: "step-2", + execute: async ({ data }) => ({ result: data.value + 5 }), + }), + ); + + const registry = WorkflowRegistry.getInstance(); + registry.registerWorkflow(workflow); + + const original = await workflow.run({ value: 3 }); + const originalBeforeReplay = await memory.getWorkflowState(original.executionId); + + await workflow.timeTravel({ + executionId: original.executionId, + stepId: "step-2", + }); + + const originalAfterReplay = await memory.getWorkflowState(original.executionId); + expect(originalAfterReplay?.status).toBe("completed"); + expect(originalAfterReplay?.output).toEqual(originalBeforeReplay?.output); + expect(originalAfterReplay?.updatedAt).toEqual(originalBeforeReplay?.updatedAt); + expect(originalAfterReplay?.replayedFromExecutionId).toBeUndefined(); + expect(originalAfterReplay?.replayFromStepId).toBeUndefined(); + }); + + it("should emit workflow stream events during replay streaming", async () => { + const memory = new Memory({ storage: new InMemoryStorageAdapter() }); + + const workflow = createWorkflow( + { + id: "time-travel-stream", + name: "Time Travel Stream", + input: z.object({ value: z.number() }), + result: z.object({ result: z.number() }), + memory, + }, + andThen({ + id: "step-1", + execute: async ({ data }) => ({ value: data.value + 1 }), + }), + andThen({ + id: "step-2", + execute: async ({ data }) => ({ value: data.value * 2 }), + }), + andThen({ + id: "step-3", + execute: async ({ data }) => ({ result: data.value + 3 }), + }), + ); + + const registry = WorkflowRegistry.getInstance(); + registry.registerWorkflow(workflow); + + const original = await workflow.run({ value: 2 }); + + const stream = workflow.timeTravelStream({ + executionId: original.executionId, + stepId: "step-2", + }); + + const events: Array<{ type: string; from?: string }> = []; + for await (const event of stream) { + events.push({ type: event.type, from: event.from }); + } + + const replayResult = await stream.result; + expect(replayResult).toEqual({ result: 9 }); + + const eventTypes = events.map((event) => event.type); + expect(eventTypes).toContain("workflow-start"); + expect(eventTypes).toContain("step-start"); + expect(eventTypes).toContain("step-complete"); + expect(eventTypes).toContain("workflow-complete"); + + const startedSteps = events + .filter((event) => event.type === "step-start") + .map((event) => event.from); + expect(startedSteps).toContain("step-2"); + expect(startedSteps).toContain("step-3"); + expect(startedSteps).not.toContain("step-1"); + }); +}); diff --git a/packages/core/src/workflow/types.ts b/packages/core/src/workflow/types.ts index 1f5b4991a..b26c9a151 100644 --- a/packages/core/src/workflow/types.ts +++ b/packages/core/src/workflow/types.ts @@ -224,6 +224,29 @@ export interface WorkflowStartAsyncResult { startAt: Date; } +export interface WorkflowTimeTravelOptions { + /** + * Source execution ID to replay from + */ + executionId: string; + /** + * Step ID to restart execution from + */ + stepId: string; + /** + * Optional override for the selected step input/state data + */ + inputData?: DangerouslyAllowAny; + /** + * Optional resume payload passed as `resumeData` to the selected step + */ + resumeData?: DangerouslyAllowAny; + /** + * Optional override for shared workflow state during replay + */ + workflowStateOverride?: WorkflowStateStore; +} + export interface WorkflowRetryConfig { /** * Number of retry attempts for a step when it throws an error @@ -733,6 +756,20 @@ export type Workflow< input: WorkflowInput, options?: WorkflowRunOptions, ) => Promise; + /** + * Replay an existing execution from a selected historical step. + * A new execution ID is created and linked to the source run for audit safety. + */ + timeTravel: ( + options: WorkflowTimeTravelOptions, + ) => Promise>; + /** + * Stream replay execution from a selected historical step. + * A new execution ID is created and linked to the source run for audit safety. + */ + timeTravelStream: ( + options: WorkflowTimeTravelOptions, + ) => WorkflowStreamResult; /** * Restart an interrupted execution from persisted checkpoint state * @param executionId - Execution ID to restart diff --git a/website/docs/workflows/overview.md b/website/docs/workflows/overview.md index f10f138e3..a242b8172 100644 --- a/website/docs/workflows/overview.md +++ b/website/docs/workflows/overview.md @@ -609,6 +609,53 @@ These are equivalent to `workflow.toWorkflow().restart(...)` and `workflow.toWor For cross-workflow recovery, use `WorkflowRegistry.getInstance().restartAllActiveWorkflowRuns()`. +### Time Travel (Deterministic Replay) + +Use time travel when you need to replay a historical execution from a specific step for debugging or operational recovery. + +Unlike `restart()`, time travel creates a new execution ID and keeps the original run immutable. + +```typescript +const runnableWorkflow = workflow.toWorkflow(); + +const original = await runnableWorkflow.run({ value: 1 }); + +const replay = await runnableWorkflow.timeTravel({ + executionId: original.executionId, + stepId: "step-2", + // Optional overrides: + // inputData: { value: 100 }, + // resumeData: { approved: true }, + // workflowStateOverride: { replayReason: "manual-debug" }, +}); + +console.log(replay.executionId); // New replay execution ID +console.log(replay.result); // Final replay output + +const replayState = await runnableWorkflow.memory.getWorkflowState(replay.executionId); +console.log(replayState?.replayedFromExecutionId); // original.executionId +console.log(replayState?.replayFromStepId); // "step-2" +``` + +`WorkflowChain` exposes the same APIs: + +```typescript +await workflow.timeTravel({ + executionId: original.executionId, + stepId: "step-2", +}); + +const replayStream = workflow.timeTravelStream({ + executionId: original.executionId, + stepId: "step-2", +}); +``` + +Notes: + +- Time travel only works for non-running executions. +- If the source execution is still `running`, use `restart(...)` for crash recovery or wait until it reaches a terminal status. + ### Executing Workflows via REST API Once your workflows are registered with VoltAgent, they can also be executed through the REST API. This is useful for triggering workflows from web applications, mobile apps, or any external system. diff --git a/website/docs/workflows/streaming.md b/website/docs/workflows/streaming.md index e6eaa0a33..355877da4 100644 --- a/website/docs/workflows/streaming.md +++ b/website/docs/workflows/streaming.md @@ -42,11 +42,12 @@ Workflows emit these event types during execution: ### Consuming the Stream -VoltAgent provides three methods for workflow execution: +VoltAgent provides four methods for workflow execution: - `.stream()` - Real-time execution with event streaming - `.run()` - Standard execution without streaming - `.startAsync()` - Fire-and-forget execution (returns immediately) +- `.timeTravelStream()` - Real-time streaming replay from a historical execution step ```typescript // Method 1: Stream execution for real-time events @@ -84,6 +85,21 @@ console.log("Started execution:", started.executionId); // Later, inspect status/output from workflow memory const state = await workflow.memory.getWorkflowState(started.executionId); console.log("Current status:", state?.status); + +// Method 4: Stream a deterministic replay from a historical run +const sourceExecution = await workflow.run(input); + +const replayStream = workflow.timeTravelStream({ + executionId: sourceExecution.executionId, + stepId: "step-2", // Replay starts from this step +}); + +for await (const event of replayStream) { + console.log("Replay event:", event.type, event.from); +} + +const replayResult = await replayStream.result; +console.log("Replay result:", replayResult); ``` ## Writer API @@ -787,6 +803,22 @@ interface WorkflowStartAsyncResult { } ``` +### WorkflowTimeTravelOptions + +Used by `.timeTravel()` and `.timeTravelStream()` to replay a historical execution: + +```typescript +interface WorkflowTimeTravelOptions { + executionId: string; // Source execution ID + stepId: string; // Step to replay from + inputData?: unknown; // Optional selected-step input override + resumeData?: unknown; // Optional resume payload override + workflowStateOverride?: Record; // Optional shared workflow state override +} +``` + +`.timeTravelStream()` returns `WorkflowStreamResult`, just like `.stream()`, but uses historical state as its starting point. + ### Key Differences | Feature | `.run()` | `.startAsync()` | `.stream()` | diff --git a/website/docs/workflows/suspend-resume.md b/website/docs/workflows/suspend-resume.md index f01f54b25..ff38b13e5 100644 --- a/website/docs/workflows/suspend-resume.md +++ b/website/docs/workflows/suspend-resume.md @@ -446,6 +446,76 @@ console.log(summary.restarted.length, summary.failed.length); - VoltAgent restores checkpointed workflow data, shared workflow state, context, and usage before continuing. - Steps should be idempotent where possible, because external side effects may have already occurred before a crash. +## Time Travel & Deterministic Replay + +Use time travel to replay a completed/suspended/cancelled/error execution from a specific historical step. + +Time travel differs from restart: + +- `restart(executionId)` continues a `running` execution after crash/interruption. +- `timeTravel({ executionId, stepId })` creates a new execution from historical state. + +### Replay from a Specific Step + +```typescript +const workflow = myWorkflowChain.toWorkflow(); + +const original = await workflow.run({ value: 1 }); + +const replay = await workflow.timeTravel({ + executionId: original.executionId, + stepId: "step-2", +}); + +console.log(replay.executionId); // New execution ID +console.log(replay.result); // Replay result +``` + +### Replay with Overrides + +You can override the selected step input, resume payload, or shared workflow state: + +```typescript +const replay = await workflow.timeTravel({ + executionId: original.executionId, + stepId: "approval-step", + inputData: { amount: 2500 }, + resumeData: { approved: true, approvedBy: "ops-user-1" }, + workflowStateOverride: { replayReason: "incident-1234" }, +}); +``` + +### Streaming Replay + +For real-time replay events, use `timeTravelStream`: + +```typescript +const stream = workflow.timeTravelStream({ + executionId: original.executionId, + stepId: "step-2", +}); + +for await (const event of stream) { + console.log(event.type, event.from); +} + +const replayResult = await stream.result; +console.log(replayResult); +``` + +### Replay Lineage Metadata + +Replay executions persist lineage fields so you can trace origin: + +- `replayedFromExecutionId` +- `replayFromStepId` + +```typescript +const replayState = await workflow.memory.getWorkflowState(replay.executionId); +console.log(replayState?.replayedFromExecutionId); +console.log(replayState?.replayFromStepId); +``` + ## External Suspension You can also pause workflows from outside using `createSuspendController`: From c8d2d931b55b061d24d65df00def7db45032e9bd Mon Sep 17 00:00:00 2001 From: Omer Aplak Date: Sat, 21 Feb 2026 21:34:44 -0800 Subject: [PATCH 2/7] docs: add workflow replay REST endpoint docs and examples --- .changeset/four-cobras-yawn.md | 2 + website/docs/api/endpoints/workflows.md | 102 ++++++++++++++++++++++++ 2 files changed, 104 insertions(+) diff --git a/.changeset/four-cobras-yawn.md b/.changeset/four-cobras-yawn.md index 7640fc94c..b7ab74639 100644 --- a/.changeset/four-cobras-yawn.md +++ b/.changeset/four-cobras-yawn.md @@ -21,3 +21,5 @@ Replay lineage metadata is now persisted on workflow state records: New public type exports from `@voltagent/core` include `WorkflowTimeTravelOptions`. Also adds workflow documentation and usage examples for deterministic replay in overview, suspend/resume, and streaming docs. + +Adds REST API documentation for replay endpoint `POST /workflows/:id/executions/:executionId/replay`, including request/response details and both cURL and JavaScript (`fetch`) code examples for default replay and replay with overrides (`inputData`, `resumeData`, `workflowStateOverride`). diff --git a/website/docs/api/endpoints/workflows.md b/website/docs/api/endpoints/workflows.md index 51ff72e5e..8dd45a8b8 100644 --- a/website/docs/api/endpoints/workflows.md +++ b/website/docs/api/endpoints/workflows.md @@ -510,6 +510,108 @@ curl -X POST http://localhost:3141/workflows/order-approval/executions/exec_123/ }' ``` +## Replay Workflow + +Create a deterministic replay execution from a historical run and selected step. + +**Endpoint:** `POST /workflows/:id/executions/:executionId/replay` + +**Request Body:** + +```json +{ + "stepId": "approval-step", + "inputData": { + "amount": 2500 + }, + "resumeData": { + "approved": true, + "approvedBy": "ops-user-1" + }, + "workflowStateOverride": { + "replayReason": "incident-1234" + } +} +``` + +**Parameters:** + +| Field | Type | Description | +| ----------------------- | ------ | --------------------------------------- | +| `stepId` | string | Historical step ID to replay from | +| `inputData` | any | Optional selected-step input override | +| `resumeData` | any | Optional resume payload override | +| `workflowStateOverride` | object | Optional shared workflow state override | + +**Response:** + +```json +{ + "success": true, + "data": { + "executionId": "exec_replay_123", + "startAt": "2024-01-15T11:00:00.000Z", + "endAt": "2024-01-15T11:00:02.250Z", + "status": "completed", + "result": { + "approved": true, + "finalAmount": 2500 + } + } +} +``` + +**Error Cases:** + +- `400` - Invalid replay parameters (for example invalid `stepId` or source execution still running) +- `404` - Workflow or source execution not found +- `500` - Replay failed due to server error + +**cURL Example (Default Replay):** + +```bash +curl -X POST http://localhost:3141/workflows/order-approval/executions/exec_123/replay \ + -H "Content-Type: application/json" \ + -d '{ + "stepId": "approval-step" + }' +``` + +**cURL Example (Replay With Overrides):** + +```bash +curl -X POST http://localhost:3141/workflows/order-approval/executions/exec_123/replay \ + -H "Content-Type: application/json" \ + -d '{ + "stepId": "approval-step", + "inputData": { "amount": 2500 }, + "resumeData": { "approved": true, "approvedBy": "ops-user-1" }, + "workflowStateOverride": { "replayReason": "incident-1234" } + }' +``` + +**JavaScript Example:** + +```javascript +const response = await fetch( + "http://localhost:3141/workflows/order-approval/executions/exec_123/replay", + { + method: "POST", + headers: { "Content-Type": "application/json" }, + body: JSON.stringify({ + stepId: "approval-step", + inputData: { amount: 2500 }, + resumeData: { approved: true, approvedBy: "ops-user-1" }, + workflowStateOverride: { replayReason: "incident-1234" }, + }), + } +); + +const replay = await response.json(); +console.log("Replay execution ID:", replay.data.executionId); +console.log("Replay status:", replay.data.status); +``` + ## Get Workflow State Retrieve the current state of a workflow execution. From d196943a8772bbe59285613408f1666521169413 Mon Sep 17 00:00:00 2001 From: Omer Aplak Date: Sat, 21 Feb 2026 21:57:28 -0800 Subject: [PATCH 3/7] feat: add endpoints --- packages/core/src/observability/types.ts | 5 + packages/core/src/workflow/chain.ts | 6 ++ packages/core/src/workflow/core.ts | 59 ++++++++++-- .../open-telemetry/trace-context.spec.ts | 49 +++++++++- .../workflow/open-telemetry/trace-context.ts | 49 +++++++++- .../core/src/workflow/time-travel.spec.ts | 68 +++++++++++++ packages/core/src/workflow/types.ts | 16 ++++ packages/server-core/src/auth/defaults.ts | 1 + .../src/handlers/workflow.handlers.ts | 95 +++++++++++++++++++ .../server-core/src/routes/definitions.ts | 27 ++++++ .../server-core/src/schemas/agent.schemas.ts | 23 +++++ .../src/routes/workflow.routes.ts | 39 ++++++++ packages/server-elysia/src/schemas.ts | 4 + .../server-hono/src/routes/agent.routes.ts | 67 +++++++++++++ packages/server-hono/src/routes/index.ts | 18 ++++ packages/serverless-hono/src/routes.ts | 12 +++ 16 files changed, 529 insertions(+), 9 deletions(-) diff --git a/packages/core/src/observability/types.ts b/packages/core/src/observability/types.ts index 0f6ea9d35..beea33eab 100644 --- a/packages/core/src/observability/types.ts +++ b/packages/core/src/observability/types.ts @@ -162,6 +162,11 @@ export interface SpanAttributes { "workflow.step.index"?: number; "workflow.step.type"?: string; "workflow.step.name"?: string; + "workflow.replayed"?: boolean; + "workflow.replay.source_trace_id"?: string; + "workflow.replay.source_span_id"?: string; + "workflow.replay.source_execution_id"?: string; + "workflow.replay.source_step_id"?: string; // Tool-specific attributes "tool.name"?: string; diff --git a/packages/core/src/workflow/chain.ts b/packages/core/src/workflow/chain.ts index 5fd27a5b2..28d0337a0 100644 --- a/packages/core/src/workflow/chain.ts +++ b/packages/core/src/workflow/chain.ts @@ -990,6 +990,9 @@ export class WorkflowChain< /** * Replay a historical execution from the selected step + * This recreates a workflow instance via `createWorkflow(...)` on each call. + * Use persistent/shared memory (or register the workflow) so source execution state is discoverable. + * For ephemeral setup patterns, prefer `chain.toWorkflow().timeTravel(...)` and reuse that instance. */ async timeTravel( options: WorkflowTimeTravelOptions, @@ -1007,6 +1010,9 @@ export class WorkflowChain< /** * Stream a historical replay from the selected step + * This recreates a workflow instance via `createWorkflow(...)` on each call. + * Use persistent/shared memory (or register the workflow) so source execution state is discoverable. + * For ephemeral setup patterns, prefer `chain.toWorkflow().timeTravelStream(...)` and reuse that instance. */ timeTravelStream( options: WorkflowTimeTravelOptions, diff --git a/packages/core/src/workflow/core.ts b/packages/core/src/workflow/core.ts index 1c023be59..ab68bbd10 100644 --- a/packages/core/src/workflow/core.ts +++ b/packages/core/src/workflow/core.ts @@ -58,6 +58,7 @@ import type { } from "./types"; export const VOLTAGENT_RESTART_CHECKPOINT_KEY = "__voltagent_restart_checkpoint"; +const workflowReplayLogger = new LoggerProxy({ component: "workflow-core-replay" }); const isObjectRecord = (value: unknown): value is Record => typeof value === "object" && value !== null && !Array.isArray(value); @@ -119,10 +120,22 @@ const isWorkflowStepStatus = (value: unknown): value is WorkflowStepData["status value === "cancelled" || value === "skipped"; -const toWorkflowStepStatus = (value: unknown): WorkflowStepData["status"] => { +const toWorkflowStepStatus = ( + value: unknown, + logger?: Pick, +): WorkflowStepData["status"] => { if (isWorkflowStepStatus(value)) { return value; } + + const targetLogger = logger ?? workflowReplayLogger; + targetLogger.warn( + "Unexpected workflow step status in replay checkpoint; defaulting to 'success'", + { + rawStatusValue: value, + }, + ); + return "success"; }; @@ -1022,9 +1035,37 @@ export function createWorkflow< : undefined; const workflowStateStore = options?.workflowState ?? {}; - // Get previous trace IDs if resuming + // Resolve trace lineage for resume/replay links let resumedFrom: { traceId: string; spanId: string } | undefined; - if (options?.resumeFrom?.executionId) { + let replayedFrom: + | { + traceId: string; + spanId: string; + executionId: string; + stepId: string; + } + | undefined; + + if (options?.replayFrom?.executionId) { + try { + const workflowState = await executionMemory.getWorkflowState( + options.replayFrom.executionId, + ); + if (workflowState?.metadata?.traceId && workflowState?.metadata?.spanId) { + replayedFrom = { + traceId: workflowState.metadata.traceId as string, + spanId: workflowState.metadata.spanId as string, + executionId: options.replayFrom.executionId, + stepId: options.replayFrom.stepId, + }; + logger.debug("Found source trace IDs for replay:", replayedFrom); + } else { + logger.warn("No source trace IDs found in replay workflow state metadata"); + } + } catch (error) { + logger.warn("Failed to get source trace IDs for replay:", { error }); + } + } else if (options?.resumeFrom?.executionId) { try { const workflowState = await executionMemory.getWorkflowState(executionId); // Look for trace IDs from the original execution @@ -1052,6 +1093,7 @@ export function createWorkflow< input: input, context: contextMap, resumedFrom, + replayedFrom, }); // Wrap entire execution in root span @@ -1087,7 +1129,7 @@ export function createWorkflow< }); // Check if resuming an existing execution - if (options?.resumeFrom?.executionId) { + if (options?.resumeFrom?.executionId && !options?.replayFrom) { runLogger.debug(`Resuming execution ${executionId} for workflow ${id}`); // Record resume in trace @@ -2578,7 +2620,7 @@ export function createWorkflow< replayStepData[step.id] = { input: checkpointSnapshot.input, output: checkpointSnapshot.output, - status: toWorkflowStepStatus(checkpointSnapshot.status), + status: toWorkflowStepStatus(checkpointSnapshot.status, logger), error: serializeStepError(checkpointSnapshot.error), }; continue; @@ -2596,7 +2638,7 @@ export function createWorkflow< replayStepData[step.id] = { input: fallbackEvent.input, output: fallbackEvent.output, - status: toWorkflowStepStatus(fallbackEvent.status), + status: toWorkflowStepStatus(fallbackEvent.status, logger), error: null, }; } @@ -2682,6 +2724,11 @@ export function createWorkflow< context: sourceContext, workflowState: effectiveWorkflowState, metadata: lineageMetadata, + skipStateInit: true, + replayFrom: { + executionId: timeTravelOptions.executionId, + stepId: timeTravelOptions.stepId, + }, resumeFrom: { executionId: replayExecutionId, resumeStepIndex: targetStepIndex, diff --git a/packages/core/src/workflow/open-telemetry/trace-context.spec.ts b/packages/core/src/workflow/open-telemetry/trace-context.spec.ts index dd852e32a..84657e3e9 100644 --- a/packages/core/src/workflow/open-telemetry/trace-context.spec.ts +++ b/packages/core/src/workflow/open-telemetry/trace-context.spec.ts @@ -8,14 +8,14 @@ async function waitForSpan( traceId: string, spanId: string, ) { - for (let attempt = 0; attempt < 10; attempt++) { + for (let attempt = 0; attempt < 30; attempt++) { await observability.forceFlush(); const traceSpans = await observability.getTraceFromStorage(traceId); const span = traceSpans.find((candidate: any) => candidate.spanId === spanId); if (span) { return span; } - await new Promise((resolve) => setTimeout(resolve, 10)); + await new Promise((resolve) => setTimeout(resolve, 20)); } return undefined; @@ -131,4 +131,49 @@ describe.sequential("WorkflowTraceContext", () => { ]), ); }); + + it("adds replay lineage links and attributes for time-travel traces", async () => { + const tracer = observability.getTracer(); + const sourceSpan = tracer.startSpan("workflow.source"); + const sourceContext = sourceSpan.spanContext(); + sourceSpan.end(); + + const traceContext = new WorkflowTraceContext(observability, "workflow.replay", { + workflowId: "wf-3", + workflowName: "workflow-replay", + executionId: "exec-3", + replayedFrom: { + traceId: sourceContext.traceId, + spanId: sourceContext.spanId, + executionId: "exec-source", + stepId: "step-2", + }, + }); + const rootSpan = traceContext.getRootSpan() as any; + const rootLinks = rootSpan.links; + const rootAttributes = rootSpan.attributes as Record | undefined; + traceContext.end("completed"); + + expect(rootLinks).toEqual( + expect.arrayContaining([ + expect.objectContaining({ + context: expect.objectContaining({ + traceId: sourceContext.traceId, + spanId: sourceContext.spanId, + }), + attributes: expect.objectContaining({ + "link.type": "replay", + "workflow.replayed": true, + "workflow.replay.source_execution_id": "exec-source", + "workflow.replay.source_step_id": "step-2", + }), + }), + ]), + ); + + expect(rootAttributes?.["workflow.replayed"]).toBe(true); + expect(rootAttributes?.["workflow.replay.source_trace_id"]).toBe(sourceContext.traceId); + expect(rootAttributes?.["workflow.replay.source_step_id"]).toBe("step-2"); + expect(rootAttributes?.["workflow.resumed"]).not.toBe(true); + }); }); diff --git a/packages/core/src/workflow/open-telemetry/trace-context.ts b/packages/core/src/workflow/open-telemetry/trace-context.ts index 02632d5e8..188e95d3b 100644 --- a/packages/core/src/workflow/open-telemetry/trace-context.ts +++ b/packages/core/src/workflow/open-telemetry/trace-context.ts @@ -54,6 +54,12 @@ export interface WorkflowTraceContextOptions { traceId: string; spanId: string; }; + replayedFrom?: { + traceId: string; + spanId: string; + executionId: string; + stepId: string; + }; } export class WorkflowTraceContext { @@ -138,6 +144,24 @@ export class WorkflowTraceContext { }, ] : []), + ...(options.replayedFrom + ? [ + { + context: { + traceId: options.replayedFrom.traceId, + spanId: options.replayedFrom.spanId, + traceFlags: 1, // Sampled + traceState: undefined, + }, + attributes: { + "link.type": "replay", + "workflow.replayed": true, + "workflow.replay.source_execution_id": options.replayedFrom.executionId, + "workflow.replay.source_step_id": options.replayedFrom.stepId, + }, + }, + ] + : []), ]; this.rootSpan = this.tracer.startSpan( @@ -152,12 +176,28 @@ export class WorkflowTraceContext { "workflow.previous_trace_id": options.resumedFrom.traceId, "workflow.previous_span_id": options.resumedFrom.spanId, }), + ...(options.replayedFrom && { + "workflow.replayed": true, + "workflow.replay.source_trace_id": options.replayedFrom.traceId, + "workflow.replay.source_span_id": options.replayedFrom.spanId, + "workflow.replay.source_execution_id": options.replayedFrom.executionId, + "workflow.replay.source_step_id": options.replayedFrom.stepId, + }), }, links: links.length > 0 ? links : undefined, }, parentContext, ); + if (options.replayedFrom) { + this.rootSpan.addEvent("workflow.replayed", { + "replay.source_trace_id": options.replayedFrom.traceId, + "replay.source_span_id": options.replayedFrom.spanId, + "replay.source_execution_id": options.replayedFrom.executionId, + "replay.source_step_id": options.replayedFrom.stepId, + }); + } + // Set active context with root span this.activeContext = trace.setSpan(context.active(), this.rootSpan); pushActiveSpan(this.rootSpan); @@ -584,8 +624,15 @@ export function addWorkflowAttributesToSpan( } } + // Replay lineage + if (options?.replayFrom) { + span.setAttribute("workflow.replayed", true); + span.setAttribute("workflow.replay.source_execution_id", options.replayFrom.executionId); + span.setAttribute("workflow.replay.source_step_id", options.replayFrom.stepId); + } + // Resume information - if (options?.resumeFrom) { + if (options?.resumeFrom && !options?.replayFrom) { span.setAttribute("workflow.resumed", true); span.setAttribute("workflow.resume.execution_id", options.resumeFrom.executionId); span.setAttribute("workflow.resume.step_index", options.resumeFrom.resumeStepIndex); diff --git a/packages/core/src/workflow/time-travel.spec.ts b/packages/core/src/workflow/time-travel.spec.ts index caad64cb9..77917e9f5 100644 --- a/packages/core/src/workflow/time-travel.spec.ts +++ b/packages/core/src/workflow/time-travel.spec.ts @@ -151,6 +151,74 @@ describe.sequential("workflow.timeTravel", () => { ).rejects.toThrow("Step 'missing-step' not found"); }); + it("should reject time travel when source execution is still running", async () => { + const memory = new Memory({ storage: new InMemoryStorageAdapter() }); + const runningExecutionId = "exec-running-time-travel"; + + const workflow = createWorkflow( + { + id: "time-travel-running-source", + name: "Time Travel Running Source", + input: z.object({ value: z.number() }), + result: z.object({ value: z.number() }), + memory, + }, + andThen({ + id: "step-1", + execute: async ({ data }) => data, + }), + ); + + const registry = WorkflowRegistry.getInstance(); + registry.registerWorkflow(workflow); + + const now = new Date(); + await memory.setWorkflowState(runningExecutionId, { + id: runningExecutionId, + workflowId: workflow.id, + workflowName: workflow.name, + status: "running", + input: { value: 1 }, + createdAt: now, + updatedAt: now, + }); + + await expect( + workflow.timeTravel({ + executionId: runningExecutionId, + stepId: "step-1", + }), + ).rejects.toThrow("running"); + }); + + it("should fail with actionable error when execution does not exist", async () => { + const memory = new Memory({ storage: new InMemoryStorageAdapter() }); + + const workflow = createWorkflow( + { + id: "time-travel-missing-execution", + name: "Time Travel Missing Execution", + input: z.object({ value: z.number() }), + result: z.object({ value: z.number() }), + memory, + }, + andThen({ + id: "step-1", + execute: async ({ data }) => data, + }), + ); + + const registry = WorkflowRegistry.getInstance(); + registry.registerWorkflow(workflow); + + await expect( + workflow.timeTravel({ + executionId: "exec-missing-time-travel", + stepId: "step-1", + }), + ).rejects.toThrow("Workflow state not found"); + }); + it("should keep original execution history unchanged after replay", async () => { const memory = new Memory({ storage: new InMemoryStorageAdapter() }); diff --git a/packages/core/src/workflow/types.ts b/packages/core/src/workflow/types.ts index b26c9a151..2a8fc2f41 100644 --- a/packages/core/src/workflow/types.ts +++ b/packages/core/src/workflow/types.ts @@ -310,6 +310,11 @@ export interface WorkflowRunOptions { * Options for resuming a suspended workflow */ resumeFrom?: WorkflowResumeOptions; + /** + * Internal replay lineage context for deterministic time-travel executions + * @internal + */ + replayFrom?: WorkflowReplayOptions; /** * Suspension mode: * - 'graceful': Wait for current step to complete before suspending (default) @@ -384,6 +389,17 @@ export interface WorkflowResumeOptions { resumeData?: DangerouslyAllowAny; } +export interface WorkflowReplayOptions { + /** + * Source execution ID used for replay lineage + */ + executionId: string; + /** + * Source step ID where replay starts + */ + stepId: string; +} + export interface WorkflowRestartCheckpoint { /** * Zero-based step index where execution should continue diff --git a/packages/server-core/src/auth/defaults.ts b/packages/server-core/src/auth/defaults.ts index faf69fe9b..3dc005d58 100644 --- a/packages/server-core/src/auth/defaults.ts +++ b/packages/server-core/src/auth/defaults.ts @@ -119,6 +119,7 @@ export const PROTECTED_ROUTES = [ // ======================================== "POST /workflows/:id/executions/:executionId/suspend", // Suspend execution "POST /workflows/:id/executions/:executionId/resume", // Resume execution + "POST /workflows/:id/executions/:executionId/replay", // Replay execution from historical step "POST /workflows/:id/executions/:executionId/cancel", // Cancel execution // ======================================== diff --git a/packages/server-core/src/handlers/workflow.handlers.ts b/packages/server-core/src/handlers/workflow.handlers.ts index 16a41f341..f5b9513ef 100644 --- a/packages/server-core/src/handlers/workflow.handlers.ts +++ b/packages/server-core/src/handlers/workflow.handlers.ts @@ -824,6 +824,101 @@ export async function handleResumeWorkflow( } } +/** + * Handler for replaying a workflow execution from a historical step + * Returns replay result + */ +export async function handleReplayWorkflow( + workflowId: string, + executionId: string, + body: any, + deps: ServerProviderDeps, + logger: Logger, +): Promise { + try { + const { stepId, inputData, resumeData, workflowStateOverride } = body || {}; + + if (typeof stepId !== "string" || stepId.trim().length === 0) { + return { + success: false, + error: "stepId is required", + httpStatus: 400, + }; + } + + const registeredWorkflow = deps.workflowRegistry.getWorkflow(workflowId); + if (!registeredWorkflow) { + return { + success: false, + error: "Workflow not found", + httpStatus: 404, + }; + } + + const workflowWithReplay = registeredWorkflow.workflow as typeof registeredWorkflow.workflow & { + timeTravel?: (options: { + executionId: string; + stepId: string; + inputData?: unknown; + resumeData?: unknown; + workflowStateOverride?: Record; + }) => Promise<{ + executionId: string; + startAt: Date | string; + endAt: Date | string; + status: string; + result: unknown; + }>; + }; + + if (typeof workflowWithReplay.timeTravel !== "function") { + return { + success: false, + error: "Workflow does not support replay", + httpStatus: 400, + }; + } + + const result = await workflowWithReplay.timeTravel({ + executionId, + stepId: stepId.trim(), + inputData, + resumeData, + workflowStateOverride, + }); + + return { + success: true, + data: { + executionId: result.executionId, + startAt: result.startAt instanceof Date ? result.startAt.toISOString() : result.startAt, + endAt: result.endAt instanceof Date ? result.endAt.toISOString() : result.endAt, + status: result.status, + result: result.result, + }, + }; + } catch (error) { + logger.error("Failed to replay workflow", { error, workflowId, executionId }); + + const message = error instanceof Error ? error.message : "Failed to replay workflow"; + const normalizedMessage = message.toLowerCase(); + const httpStatus = normalizedMessage.includes("not found") + ? 404 + : normalizedMessage.includes("cannot time travel") || + normalizedMessage.includes("still running") || + normalizedMessage.includes("belongs to workflow") || + normalizedMessage.includes("step") + ? 400 + : 500; + + return { + success: false, + error: message, + httpStatus, + }; + } +} + function formatWorkflowState(workflowState: WorkflowStateEntry) { return { ...workflowState, diff --git a/packages/server-core/src/routes/definitions.ts b/packages/server-core/src/routes/definitions.ts index 0766e1b57..749b98c0b 100644 --- a/packages/server-core/src/routes/definitions.ts +++ b/packages/server-core/src/routes/definitions.ts @@ -617,6 +617,33 @@ export const WORKFLOW_ROUTES = { }, }, }, + replayWorkflow: { + method: "post" as const, + path: "/workflows/:id/executions/:executionId/replay", + summary: "Replay workflow execution from a step", + description: + "Create a deterministic replay execution from a historical workflow run and selected step. Replay creates a new execution ID and preserves the original run history.", + tags: ["Workflow Management"], + operationId: "replayWorkflow", + responses: { + 200: { + description: "Successfully replayed workflow execution", + contentType: "application/json", + }, + 400: { + description: "Invalid replay parameters", + contentType: "application/json", + }, + 404: { + description: "Workflow or source execution not found", + contentType: "application/json", + }, + 500: { + description: "Failed to replay workflow due to server error", + contentType: "application/json", + }, + }, + }, getWorkflowState: { method: "get" as const, path: "/workflows/:id/executions/:executionId/state", diff --git a/packages/server-core/src/schemas/agent.schemas.ts b/packages/server-core/src/schemas/agent.schemas.ts index 488681f78..b4dcc5e1d 100644 --- a/packages/server-core/src/schemas/agent.schemas.ts +++ b/packages/server-core/src/schemas/agent.schemas.ts @@ -455,3 +455,26 @@ export const WorkflowResumeResponseSchema = z.object({ }) .describe("Workflow resume result"), }); + +export const WorkflowReplayRequestSchema = z.object({ + stepId: z.string().describe("Step ID to replay from"), + inputData: z.any().optional().describe("Optional input override for the selected step"), + resumeData: z.any().optional().describe("Optional resume payload override"), + workflowStateOverride: z + .record(z.string(), z.any()) + .optional() + .describe("Optional workflow state override for replay run"), +}); + +export const WorkflowReplayResponseSchema = z.object({ + success: z.literal(true), + data: z + .object({ + executionId: z.string(), + startAt: z.string(), + endAt: z.string().optional(), + status: z.string(), + result: z.any(), + }) + .describe("Workflow replay result"), +}); diff --git a/packages/server-elysia/src/routes/workflow.routes.ts b/packages/server-elysia/src/routes/workflow.routes.ts index 160c05cf8..b2aa3fafa 100644 --- a/packages/server-elysia/src/routes/workflow.routes.ts +++ b/packages/server-elysia/src/routes/workflow.routes.ts @@ -8,6 +8,7 @@ import { handleGetWorkflowState, handleGetWorkflows, handleListWorkflowRuns, + handleReplayWorkflow, handleResumeWorkflow, handleStreamWorkflow, handleSuspendWorkflow, @@ -22,6 +23,8 @@ import { WorkflowExecutionRequestSchema, WorkflowExecutionResponseSchema, WorkflowListSchema, + WorkflowReplayRequestSchema, + WorkflowReplayResponseSchema, WorkflowResponseSchema, WorkflowResumeRequestSchema, WorkflowResumeResponseSchema, @@ -379,6 +382,42 @@ export function registerWorkflowRoutes( }, ); + // POST /workflows/:id/executions/:executionId/replay - Replay workflow execution from a historical step + app.post( + "/workflows/:id/executions/:executionId/replay", + async ({ params, body, set }) => { + const response = await handleReplayWorkflow( + params.id, + params.executionId, + body, + deps, + logger, + ); + if (!response.success) { + set.status = response.httpStatus || 500; + return response; + } + set.status = 200; + return response; + }, + { + params: WorkflowExecutionParams, + body: WorkflowReplayRequestSchema, + response: { + 200: WorkflowReplayResponseSchema, + 400: ErrorSchema, + 404: ErrorSchema, + 500: ErrorSchema, + }, + detail: { + summary: "Replay workflow execution from step", + description: + "Creates a deterministic replay execution from a source run and the selected step ID", + tags: ["Workflows"], + }, + }, + ); + // GET /workflows/:id/executions/:executionId/state - Get workflow execution state app.get( "/workflows/:id/executions/:executionId/state", diff --git a/packages/server-elysia/src/schemas.ts b/packages/server-elysia/src/schemas.ts index c5a96a6d1..8ce5cae8e 100644 --- a/packages/server-elysia/src/schemas.ts +++ b/packages/server-elysia/src/schemas.ts @@ -12,6 +12,8 @@ import { WorkflowExecutionRequestSchema as ZodWorkflowExecutionRequestSchema, WorkflowExecutionResponseSchema as ZodWorkflowExecutionResponseSchema, WorkflowListSchema as ZodWorkflowListSchema, + WorkflowReplayRequestSchema as ZodWorkflowReplayRequestSchema, + WorkflowReplayResponseSchema as ZodWorkflowReplayResponseSchema, WorkflowResponseSchema as ZodWorkflowResponseSchema, WorkflowResumeRequestSchema as ZodWorkflowResumeRequestSchema, WorkflowResumeResponseSchema as ZodWorkflowResumeResponseSchema, @@ -62,6 +64,8 @@ export const WorkflowCancelRequestSchema = zodToTypeBox(ZodWorkflowCancelRequest export const WorkflowCancelResponseSchema = zodToTypeBox(ZodWorkflowCancelResponseSchema); export const WorkflowResumeRequestSchema = zodToTypeBox(ZodWorkflowResumeRequestSchema); export const WorkflowResumeResponseSchema = zodToTypeBox(ZodWorkflowResumeResponseSchema); +export const WorkflowReplayRequestSchema = zodToTypeBox(ZodWorkflowReplayRequestSchema); +export const WorkflowReplayResponseSchema = zodToTypeBox(ZodWorkflowReplayResponseSchema); // Update schemas export const UpdateCheckResponseSchema = t.Object({ diff --git a/packages/server-hono/src/routes/agent.routes.ts b/packages/server-hono/src/routes/agent.routes.ts index aa326f9de..f36e5b882 100644 --- a/packages/server-hono/src/routes/agent.routes.ts +++ b/packages/server-hono/src/routes/agent.routes.ts @@ -14,6 +14,8 @@ import { WorkflowExecutionRequestSchema, WorkflowExecutionResponseSchema, WorkflowListSchema, + WorkflowReplayRequestSchema, + WorkflowReplayResponseSchema, WorkflowResumeRequestSchema, WorkflowResumeResponseSchema, WorkflowStreamEventSchema, @@ -59,6 +61,8 @@ export { WorkflowSuspendResponseSchema, WorkflowCancelRequestSchema, WorkflowCancelResponseSchema, + WorkflowReplayRequestSchema, + WorkflowReplayResponseSchema, WorkflowResumeRequestSchema, WorkflowResumeResponseSchema, } from "@voltagent/server-core"; @@ -1091,3 +1095,66 @@ export const resumeWorkflowRoute = createRoute({ summary: WORKFLOW_ROUTES.resumeWorkflow.summary, description: WORKFLOW_ROUTES.resumeWorkflow.description, }); + +// Replay workflow route +export const replayWorkflowRoute = createRoute({ + method: WORKFLOW_ROUTES.replayWorkflow.method, + path: WORKFLOW_ROUTES.replayWorkflow.path + .replace(":id", "{id}") + .replace(":executionId", "{executionId}"), + request: { + params: z.object({ + id: workflowIdParam(), + executionId: executionIdParam(), + }), + body: { + content: { + "application/json": { + schema: WorkflowReplayRequestSchema, + }, + }, + }, + }, + responses: { + 200: { + content: { + "application/json": { + schema: WorkflowReplayResponseSchema, + }, + }, + description: + WORKFLOW_ROUTES.replayWorkflow.responses?.[200]?.description || + "Successful workflow replay", + }, + 400: { + content: { + "application/json": { + schema: ErrorSchema, + }, + }, + description: + WORKFLOW_ROUTES.replayWorkflow.responses?.[400]?.description || "Invalid replay request", + }, + 404: { + content: { + "application/json": { + schema: ErrorSchema, + }, + }, + description: + WORKFLOW_ROUTES.replayWorkflow.responses?.[404]?.description || + "Workflow or execution not found", + }, + 500: { + content: { + "application/json": { + schema: ErrorSchema, + }, + }, + description: WORKFLOW_ROUTES.replayWorkflow.responses?.[500]?.description || "Server error", + }, + }, + tags: [...WORKFLOW_ROUTES.replayWorkflow.tags], + summary: WORKFLOW_ROUTES.replayWorkflow.summary, + description: WORKFLOW_ROUTES.replayWorkflow.description, +}); diff --git a/packages/server-hono/src/routes/index.ts b/packages/server-hono/src/routes/index.ts index bbcab3173..2e4f75c7b 100644 --- a/packages/server-hono/src/routes/index.ts +++ b/packages/server-hono/src/routes/index.ts @@ -23,6 +23,7 @@ import { handleListAgentWorkspaceSkills, handleListWorkflowRuns, handleReadAgentWorkspaceFile, + handleReplayWorkflow, handleResumeChatStream, handleResumeWorkflow, handleStreamObject, @@ -41,6 +42,7 @@ import { getAgentsRoute, getWorkflowsRoute, objectRoute, + replayWorkflowRoute, resumeChatStreamRoute, resumeWorkflowRoute, streamObjectRoute, @@ -432,6 +434,22 @@ export function registerWorkflowRoutes( return c.json(response, 200); }); + // Replay workflow execution from a historical step + app.openapi(replayWorkflowRoute, async (c) => { + const workflowId = c.req.param("id"); + const executionId = c.req.param("executionId"); + if (!workflowId || !executionId) { + throw new Error("Missing workflow or execution id parameter"); + } + const body = await c.req.json(); + const response = await handleReplayWorkflow(workflowId, executionId, body, deps, logger); + if (!response.success) { + const status = response.httpStatus || 500; + return c.json(response, status); + } + return c.json(response, 200); + }); + app.get("/workflows/executions", async (c) => { const query = c.req.query(); const response = await handleListWorkflowRuns(undefined, query, deps, logger); diff --git a/packages/serverless-hono/src/routes.ts b/packages/serverless-hono/src/routes.ts index 481edc82c..72e73dc51 100644 --- a/packages/serverless-hono/src/routes.ts +++ b/packages/serverless-hono/src/routes.ts @@ -68,6 +68,7 @@ import { handleListTools, handleListWorkflowRuns, handleReadAgentWorkspaceFile, + handleReplayWorkflow, handleResumeChatStream, handleResumeWorkflow, handleSaveMemoryMessages, @@ -565,6 +566,17 @@ export function registerWorkflowRoutes(app: Hono, deps: ServerProviderDeps, logg return c.json(response, response.success ? 200 : 500); }); + app.post(WORKFLOW_ROUTES.replayWorkflow.path, async (c) => { + const workflowId = c.req.param("id"); + const executionId = c.req.param("executionId"); + const body = await readJsonBody(c, logger); + if (!body) { + return c.json({ success: false, error: "Invalid JSON body" }, 400); + } + const response = await handleReplayWorkflow(workflowId, executionId, body, deps, logger); + return c.json(response, response.success ? 200 : (response.httpStatus ?? 500)); + }); + app.get(WORKFLOW_ROUTES.listWorkflowRuns.path, async (c) => { const query = c.req.query(); const response = await handleListWorkflowRuns(undefined, query, deps, logger); From 62efd42326792d2ae3d372bbaf749161350d1e95 Mon Sep 17 00:00:00 2001 From: Omer Aplak Date: Sat, 21 Feb 2026 21:57:39 -0800 Subject: [PATCH 4/7] chore: update docs --- website/docs/workflows/streaming.md | 3 ++- 1 file changed, 2 insertions(+), 1 deletion(-) diff --git a/website/docs/workflows/streaming.md b/website/docs/workflows/streaming.md index 355877da4..730de647d 100644 --- a/website/docs/workflows/streaming.md +++ b/website/docs/workflows/streaming.md @@ -42,11 +42,12 @@ Workflows emit these event types during execution: ### Consuming the Stream -VoltAgent provides four methods for workflow execution: +VoltAgent provides five methods for workflow execution: - `.stream()` - Real-time execution with event streaming - `.run()` - Standard execution without streaming - `.startAsync()` - Fire-and-forget execution (returns immediately) +- `.timeTravel()` - Deterministic replay from a historical execution step - `.timeTravelStream()` - Real-time streaming replay from a historical execution step ```typescript From 7a6c2607ea711ed9e80b3f258e6253e512be6571 Mon Sep 17 00:00:00 2001 From: Omer Aplak Date: Sat, 21 Feb 2026 22:32:46 -0800 Subject: [PATCH 5/7] chore: fix code reviews --- packages/core/src/workflow/core.ts | 473 ++++++++++-------- .../core/src/workflow/time-travel.spec.ts | 1 + .../src/handlers/workflow.handlers.ts | 37 +- .../server-core/src/schemas/agent.schemas.ts | 2 +- website/docs/api/endpoints/workflows.md | 24 +- 5 files changed, 297 insertions(+), 240 deletions(-) diff --git a/packages/core/src/workflow/core.ts b/packages/core/src/workflow/core.ts index ab68bbd10..868cfb64e 100644 --- a/packages/core/src/workflow/core.ts +++ b/packages/core/src/workflow/core.ts @@ -139,6 +139,49 @@ const toWorkflowStepStatus = ( return "success"; }; +const getEventStepIndex = (event: { stepIndex?: unknown; metadata?: unknown }): + | number + | undefined => { + if (typeof event.stepIndex === "number" && Number.isInteger(event.stepIndex)) { + return event.stepIndex; + } + + if (!isObjectRecord(event.metadata)) { + return undefined; + } + + const metadataStepIndex = event.metadata.stepIndex; + if (typeof metadataStepIndex === "number" && Number.isInteger(metadataStepIndex)) { + return metadataStepIndex; + } + + if (typeof metadataStepIndex === "string") { + const parsed = Number(metadataStepIndex); + if (Number.isInteger(parsed)) { + return parsed; + } + } + + return undefined; +}; + +const getEventStepId = (event: { stepId?: unknown; metadata?: unknown }): string | undefined => { + if (typeof event.stepId === "string" && event.stepId.length > 0) { + return event.stepId; + } + + if (!isObjectRecord(event.metadata)) { + return undefined; + } + + const metadataStepId = event.metadata.stepId; + if (typeof metadataStepId === "string" && metadataStepId.length > 0) { + return metadataStepId; + } + + return undefined; +}; + const deserializeCheckpointStepData = (value: unknown): WorkflowStepData | undefined => { if (!isObjectRecord(value) || !isWorkflowStepStatus(value.status)) { return undefined; @@ -1060,10 +1103,17 @@ export function createWorkflow< }; logger.debug("Found source trace IDs for replay:", replayedFrom); } else { - logger.warn("No source trace IDs found in replay workflow state metadata"); + logger.warn("No source trace IDs found in replay workflow state metadata", { + replayExecutionId: options.replayFrom.executionId, + executionId, + }); } } catch (error) { - logger.warn("Failed to get source trace IDs for replay:", { error }); + logger.warn("Failed to get source trace IDs for replay:", { + error, + replayExecutionId: options.replayFrom.executionId, + executionId, + }); } } else if (options?.resumeFrom?.executionId) { try { @@ -1076,10 +1126,17 @@ export function createWorkflow< }; logger.debug("Found previous trace IDs for resume:", resumedFrom); } else { - logger.warn("No suspended trace IDs found in workflow state metadata"); + logger.warn("No suspended trace IDs found in workflow state metadata", { + resumeExecutionId: options.resumeFrom.executionId, + executionId, + }); } } catch (error) { - logger.warn("Failed to get previous trace IDs for resume:", { error }); + logger.warn("Failed to get previous trace IDs for resume:", { + error, + resumeExecutionId: options.resumeFrom.executionId, + executionId, + }); } } @@ -2626,13 +2683,23 @@ export function createWorkflow< continue; } - const fallbackEvent = sourceStepCompleteEvents.find( - (event) => + const fallbackEvent = sourceStepCompleteEvents.find((event) => { + const eventStepIndex = getEventStepIndex(event); + if (eventStepIndex !== undefined) { + return eventStepIndex === index; + } + + const eventStepId = getEventStepId(event); + if (eventStepId !== undefined) { + return eventStepId === step.id; + } + + return ( event.from === step.id || event.name === step.id || - event.from === step.name || - event.name === step.name, - ); + (step.name !== undefined && (event.from === step.name || event.name === step.name)) + ); + }); if (fallbackEvent) { replayStepData[step.id] = { @@ -2944,20 +3011,105 @@ export function createWorkflow< ); })(); - replayPromise - .then( + replayPromise.then( + (result) => { + if (result.status !== "suspended") { + streamController.close(); + } + }, + () => { + streamController.close(); + }, + ); + + const resumeSuspendedReplayStream = async ( + suspendedResult: WorkflowExecutionResult, + resumeInput: z.infer, + opts?: { stepId?: string }, + ): Promise> => { + if (suspendedResult.status !== "suspended") { + throw new Error(`Cannot resume workflow in ${suspendedResult.status} state`); + } + + if (!suspendedResult.suspension) { + throw new Error("No suspension metadata found"); + } + + if (!replayOriginalInput) { + throw new Error("Missing replay input for resume"); + } + + let resumeStepIndex = suspendedResult.suspension.suspendedStepIndex; + if (opts?.stepId) { + const overrideIndex = (steps as BaseStep[]).findIndex((step) => step.id === opts.stepId); + if (overrideIndex === -1) { + throw new Error(`Step '${opts.stepId}' not found in workflow '${id}'`); + } + resumeStepIndex = overrideIndex; + } + + let resumedResolve: (value: WorkflowExecutionResult) => void; + let resumedReject: (error: any) => void; + const resumedPromise = new Promise>( + (resolve, reject) => { + resumedResolve = resolve; + resumedReject = reject; + }, + ); + + const resumedSuspendController = createDefaultSuspendController(); + const resumeOptions: WorkflowRunOptions = { + executionId: suspendedResult.executionId, + resumeFrom: { + executionId: suspendedResult.executionId, + checkpoint: suspendedResult.suspension.checkpoint, + resumeStepIndex, + resumeData: resumeInput, + }, + suspendController: resumedSuspendController, + }; + + executeInternal(replayOriginalInput, resumeOptions, streamController).then( (result) => { if (result.status !== "suspended") { streamController.close(); } + resumedResolve(result); }, - () => { + (error) => { streamController.close(); + resumedReject(error); }, - ) - .catch(() => { - // Error is surfaced through promise-backed fields on stream result. - }); + ); + + const resumedStreamResult: WorkflowStreamResult = { + executionId: suspendedResult.executionId, + workflowId: suspendedResult.workflowId, + startAt: suspendedResult.startAt, + endAt: resumedPromise.then((r) => r.endAt), + status: resumedPromise.then((r) => r.status), + result: resumedPromise.then((r) => r.result), + suspension: resumedPromise.then((r) => r.suspension), + cancellation: resumedPromise.then((r) => r.cancellation), + error: resumedPromise.then((r) => r.error), + usage: resumedPromise.then((r) => r.usage), + resume: async (nextInput: z.infer, nextOpts?: { stepId?: string }) => { + const nextResult = await resumedPromise; + return resumeSuspendedReplayStream(nextResult, nextInput, nextOpts); + }, + suspend: (reason?: string) => { + resumedSuspendController.suspend(reason); + }, + cancel: (reason?: string) => { + resumedSuspendController.cancel(reason); + }, + abort: () => streamController.abort(), + toUIMessageStreamResponse: eventToUIMessageStreamResponse(streamController), + [Symbol.asyncIterator]: () => streamController.getStream(), + }; + + return resumedStreamResult; + }; const streamResult: WorkflowStreamResult = { executionId, @@ -2973,97 +3125,7 @@ export function createWorkflow< toUIMessageStreamResponse: eventToUIMessageStreamResponse(streamController), resume: async (input: z.infer, opts?: { stepId?: string }) => { const replayResult = await replayPromise; - if (replayResult.status !== "suspended") { - throw new Error(`Cannot resume workflow in ${replayResult.status} state`); - } - - if (!replayResult.suspension) { - throw new Error("No suspension metadata found"); - } - - if (!replayOriginalInput) { - throw new Error("Missing replay input for resume"); - } - - let resumeStepIndex = replayResult.suspension.suspendedStepIndex; - if (opts?.stepId) { - const overrideIndex = (steps as BaseStep[]).findIndex( - (step) => step.id === opts.stepId, - ); - if (overrideIndex === -1) { - throw new Error(`Step '${opts.stepId}' not found in workflow '${id}'`); - } - resumeStepIndex = overrideIndex; - } - - let resumedResolve: ( - value: WorkflowExecutionResult, - ) => void; - let resumedReject: (error: any) => void; - const resumedPromise = new Promise>( - (resolve, reject) => { - resumedResolve = resolve; - resumedReject = reject; - }, - ); - - const resumedSuspendController = createDefaultSuspendController(); - const resumeOptions: WorkflowRunOptions = { - executionId: replayResult.executionId, - resumeFrom: { - executionId: replayResult.executionId, - checkpoint: replayResult.suspension.checkpoint, - resumeStepIndex, - resumeData: input, - }, - suspendController: resumedSuspendController, - }; - - executeInternal(replayOriginalInput, resumeOptions, streamController) - .then( - (result) => { - if (result.status !== "suspended") { - streamController.close(); - } - resumedResolve(result); - }, - (error) => { - streamController.close(); - resumedReject(error); - }, - ) - .catch(() => {}); - - const resumedStreamResult: WorkflowStreamResult = { - executionId: replayResult.executionId, - workflowId: replayResult.workflowId, - startAt: replayResult.startAt, - endAt: resumedPromise.then((r) => r.endAt), - status: resumedPromise.then((r) => r.status), - result: resumedPromise.then((r) => r.result), - suspension: resumedPromise.then((r) => r.suspension), - cancellation: resumedPromise.then((r) => r.cancellation), - error: resumedPromise.then((r) => r.error), - usage: resumedPromise.then((r) => r.usage), - resume: async (input2: z.infer, opts2?: { stepId?: string }) => { - const nextResult = await resumedPromise; - if (nextResult.status !== "suspended") { - throw new Error(`Cannot resume workflow in ${nextResult.status} state`); - } - return streamResult.resume(input2, opts2); - }, - suspend: (reason?: string) => { - resumedSuspendController.suspend(reason); - }, - cancel: (reason?: string) => { - resumedSuspendController.cancel(reason); - }, - abort: () => streamController.abort(), - toUIMessageStreamResponse: eventToUIMessageStreamResponse(streamController), - [Symbol.asyncIterator]: () => streamController.getStream(), - }; - - return resumedStreamResult; + return resumeSuspendedReplayStream(replayResult, input, opts); }, suspend: (reason?: string) => { suspendController.suspend(reason); @@ -3140,6 +3202,91 @@ export function createWorkflow< }); // Return stream result immediately + const resumeSuspendedStream = async ( + suspendedResult: WorkflowExecutionResult, + resumeInput: z.infer, + opts?: { stepId?: string }, + ): Promise> => { + if (suspendedResult.status !== "suspended") { + throw new Error(`Cannot resume workflow in ${suspendedResult.status} state`); + } + + if (!suspendedResult.suspension) { + throw new Error("No suspension metadata found"); + } + + let resumeStepIndex = suspendedResult.suspension.suspendedStepIndex; + if (opts?.stepId) { + const overrideIndex = (steps as BaseStep[]).findIndex((step) => step.id === opts.stepId); + if (overrideIndex === -1) { + throw new Error(`Step '${opts.stepId}' not found in workflow '${id}'`); + } + resumeStepIndex = overrideIndex; + } + + let resumedResolve: (value: WorkflowExecutionResult) => void; + let resumedReject: (error: any) => void; + const resumedPromise = new Promise>( + (resolve, reject) => { + resumedResolve = resolve; + resumedReject = reject; + }, + ); + + const resumedSuspendController = createDefaultSuspendController(); + const resumeOptions: WorkflowRunOptions = { + executionId: suspendedResult.executionId, + resumeFrom: { + executionId: suspendedResult.executionId, + checkpoint: suspendedResult.suspension.checkpoint, + resumeStepIndex, + resumeData: resumeInput, + }, + suspendController: resumedSuspendController, + }; + + executeInternal(originalInput, resumeOptions, streamController).then( + (result) => { + if (result.status !== "suspended") { + streamController?.close(); + } + resumedResolve(result); + }, + (error) => { + streamController?.close(); + resumedReject(error); + }, + ); + + const resumedStreamResult: WorkflowStreamResult = { + executionId: suspendedResult.executionId, + workflowId: suspendedResult.workflowId, + startAt: suspendedResult.startAt, + endAt: resumedPromise.then((r) => r.endAt), + status: resumedPromise.then((r) => r.status), + result: resumedPromise.then((r) => r.result), + suspension: resumedPromise.then((r) => r.suspension), + cancellation: resumedPromise.then((r) => r.cancellation), + error: resumedPromise.then((r) => r.error), + usage: resumedPromise.then((r) => r.usage), + resume: async (nextInput: z.infer, nextOpts?: { stepId?: string }) => { + const nextResult = await resumedPromise; + return resumeSuspendedStream(nextResult, nextInput, nextOpts); + }, + suspend: (reason?: string) => { + resumedSuspendController.suspend(reason); + }, + cancel: (reason?: string) => { + resumedSuspendController.cancel(reason); + }, + abort: () => streamController.abort(), + toUIMessageStreamResponse: eventToUIMessageStreamResponse(streamController), + [Symbol.asyncIterator]: () => streamController.getStream(), + }; + + return resumedStreamResult; + }; + const streamResult: WorkflowStreamResult = { executionId, workflowId: id, @@ -3155,117 +3302,7 @@ export function createWorkflow< resume: async (input: z.infer, opts?: { stepId?: string }) => { const execResult = await resultPromise; - if (execResult.status !== "suspended") { - throw new Error(`Cannot resume workflow in ${execResult.status} state`); - } - - // Continue with the same stream controller - don't create a new one - // Create new promise for the resumed execution - let resumedResolve: ( - value: WorkflowExecutionResult, - ) => void; - let resumedReject: (error: any) => void; - const resumedPromise = new Promise>( - (resolve, reject) => { - resumedResolve = resolve; - resumedReject = reject; - }, - ); - - // Use a fresh controller for this resumed run. - const resumedSuspendController = createDefaultSuspendController(); - - // Execute the resume by calling stream again with resume options - const executeResume = async () => { - // Get the suspension metadata - if (!execResult.suspension) { - throw new Error("No suspension metadata found"); - } - - let resumeStepIndex = execResult.suspension.suspendedStepIndex; - if (opts?.stepId) { - const overrideIndex = (steps as BaseStep[]).findIndex( - (step) => step.id === opts.stepId, - ); - if (overrideIndex === -1) { - throw new Error(`Step '${opts.stepId}' not found in workflow '${id}'`); - } - resumeStepIndex = overrideIndex; - } - - // Create resume options to continue from where we left off - const resumeOptions: WorkflowRunOptions = { - executionId: execResult.executionId, - resumeFrom: { - executionId: execResult.executionId, - checkpoint: execResult.suspension.checkpoint, - resumeStepIndex, - resumeData: input, - }, - suspendController: resumedSuspendController, - }; - - // Re-execute with streaming from the suspension point - // This will emit events to the same stream controller - const resumed = await executeInternal( - originalInput, // Use the original input saved in closure - resumeOptions, - streamController, - ); - return resumed; - }; - - // Start resume execution and emit events to the same stream - executeResume() - .then( - (result) => { - // Only close stream if workflow completed or errored (not suspended again) - if (result.status !== "suspended") { - streamController?.close(); - } - resumedResolve(result); - }, - (error) => { - streamController?.close(); - resumedReject(error); - }, - ) - .catch(() => {}); - - // Return a stream result that continues using the same stream - const resumedStreamResult: WorkflowStreamResult = { - executionId: execResult.executionId, // Keep same execution ID - workflowId: execResult.workflowId, - startAt: execResult.startAt, - endAt: resumedPromise.then((r) => r.endAt), - status: resumedPromise.then((r) => r.status), - result: resumedPromise.then((r) => r.result), - suspension: resumedPromise.then((r) => r.suspension), - cancellation: resumedPromise.then((r) => r.cancellation), - error: resumedPromise.then((r) => r.error), - usage: resumedPromise.then((r) => r.usage), - resume: async (input2: z.infer, opts?: { stepId?: string }) => { - // Resume again using the same stream - const nextResult = await resumedPromise; - if (nextResult.status !== "suspended") { - throw new Error(`Cannot resume workflow in ${nextResult.status} state`); - } - // Recursively call resume on the stream result (which will use the same stream controller) - return streamResult.resume(input2, opts); - }, - suspend: (reason?: string) => { - resumedSuspendController.suspend(reason); - }, - cancel: (reason?: string) => { - resumedSuspendController.cancel(reason); - }, - abort: () => streamController.abort(), - toUIMessageStreamResponse: eventToUIMessageStreamResponse(streamController), - // Continue using the same stream iterator - [Symbol.asyncIterator]: () => streamController.getStream(), - }; - - return resumedStreamResult; + return resumeSuspendedStream(execResult, input, opts); }, suspend: (reason?: string) => { suspendController.suspend(reason); diff --git a/packages/core/src/workflow/time-travel.spec.ts b/packages/core/src/workflow/time-travel.spec.ts index 77917e9f5..64a7d2fb4 100644 --- a/packages/core/src/workflow/time-travel.spec.ts +++ b/packages/core/src/workflow/time-travel.spec.ts @@ -10,6 +10,7 @@ describe.sequential("workflow.timeTravel", () => { beforeEach(() => { const registry = WorkflowRegistry.getInstance(); (registry as any).workflows.clear(); + (registry as any).activeExecutions.clear(); }); it("should replay from middle step with a new execution id", async () => { diff --git a/packages/server-core/src/handlers/workflow.handlers.ts b/packages/server-core/src/handlers/workflow.handlers.ts index f5b9513ef..062a3703c 100644 --- a/packages/server-core/src/handlers/workflow.handlers.ts +++ b/packages/server-core/src/handlers/workflow.handlers.ts @@ -1,4 +1,9 @@ -import type { ServerProviderDeps, WorkflowRunQuery, WorkflowStateEntry } from "@voltagent/core"; +import type { + ServerProviderDeps, + Workflow, + WorkflowRunQuery, + WorkflowStateEntry, +} from "@voltagent/core"; import { zodSchemaToJsonUI } from "@voltagent/core"; import type { Logger } from "@voltagent/internal"; import type { ApiResponse, ErrorResponse } from "../types"; @@ -53,6 +58,10 @@ type ResumableStreamingWorkflowExecution = StreamingWorkflowExecution & { ) => Promise; }; +type WorkflowTimeTravelRequest = Parameters< + NonNullable["timeTravel"]> +>[0]; + function parseReplaySequence(value: StreamQueryValue): number | undefined { if (value === undefined || value === null || value === "") { return undefined; @@ -855,21 +864,9 @@ export async function handleReplayWorkflow( }; } - const workflowWithReplay = registeredWorkflow.workflow as typeof registeredWorkflow.workflow & { - timeTravel?: (options: { - executionId: string; - stepId: string; - inputData?: unknown; - resumeData?: unknown; - workflowStateOverride?: Record; - }) => Promise<{ - executionId: string; - startAt: Date | string; - endAt: Date | string; - status: string; - result: unknown; - }>; - }; + const workflowWithReplay = registeredWorkflow.workflow as Partial< + Pick, "timeTravel"> + >; if (typeof workflowWithReplay.timeTravel !== "function") { return { @@ -879,13 +876,14 @@ export async function handleReplayWorkflow( }; } - const result = await workflowWithReplay.timeTravel({ + const replayOptions: WorkflowTimeTravelRequest = { executionId, stepId: stepId.trim(), inputData, resumeData, workflowStateOverride, - }); + }; + const result = await workflowWithReplay.timeTravel(replayOptions); return { success: true, @@ -906,8 +904,7 @@ export async function handleReplayWorkflow( ? 404 : normalizedMessage.includes("cannot time travel") || normalizedMessage.includes("still running") || - normalizedMessage.includes("belongs to workflow") || - normalizedMessage.includes("step") + normalizedMessage.includes("belongs to workflow") ? 400 : 500; diff --git a/packages/server-core/src/schemas/agent.schemas.ts b/packages/server-core/src/schemas/agent.schemas.ts index b4dcc5e1d..75dbf0f7c 100644 --- a/packages/server-core/src/schemas/agent.schemas.ts +++ b/packages/server-core/src/schemas/agent.schemas.ts @@ -457,7 +457,7 @@ export const WorkflowResumeResponseSchema = z.object({ }); export const WorkflowReplayRequestSchema = z.object({ - stepId: z.string().describe("Step ID to replay from"), + stepId: z.string().min(1).describe("Step ID to replay from"), inputData: z.any().optional().describe("Optional input override for the selected step"), resumeData: z.any().optional().describe("Optional resume payload override"), workflowStateOverride: z diff --git a/website/docs/api/endpoints/workflows.md b/website/docs/api/endpoints/workflows.md index 8dd45a8b8..b2facbd05 100644 --- a/website/docs/api/endpoints/workflows.md +++ b/website/docs/api/endpoints/workflows.md @@ -561,9 +561,29 @@ Create a deterministic replay execution from a historical run and selected step. } ``` +**Response (Suspended):** + +```json +{ + "success": true, + "data": { + "executionId": "exec_replay_123", + "startAt": "2024-01-15T11:00:00.000Z", + "endAt": null, + "status": "suspended", + "result": null, + "suspension": { + "suspendedAt": "2024-01-15T11:00:01.250Z", + "reason": "Awaiting manual approval", + "suspendedStepIndex": 2 + } + } +} +``` + **Error Cases:** -- `400` - Invalid replay parameters (for example invalid `stepId` or source execution still running) +- `400` - Invalid replay parameters (for example, invalid `stepId` or source execution still running) - `404` - Workflow or source execution not found - `500` - Replay failed due to server error @@ -627,6 +647,8 @@ Retrieve the current state of a workflow execution. "executionId": "exec_1234567890_abc123", "workflowId": "order-approval", "status": "suspended", + "replayedFromExecutionId": "exec_0987654321_replay", + "replayFromStepId": "step-approval-1", "startAt": "2024-01-15T10:00:00.000Z", "suspension": { "suspendedAt": "2024-01-15T10:00:02.500Z", From 90f92a759b04971fb538cdc0fe8cb2b70bb1292c Mon Sep 17 00:00:00 2001 From: Omer Aplak Date: Sat, 21 Feb 2026 22:50:17 -0800 Subject: [PATCH 6/7] fix: address review feedback for time-travel replay and typing --- packages/core/src/workflow/core.ts | 42 ++++++++++++------- packages/core/src/workflow/registry.ts | 8 ++++ .../core/src/workflow/time-travel.spec.ts | 3 +- packages/core/src/workflow/types.ts | 5 +++ .../src/handlers/workflow.handlers.ts | 5 ++- 5 files changed, 44 insertions(+), 19 deletions(-) diff --git a/packages/core/src/workflow/core.ts b/packages/core/src/workflow/core.ts index 868cfb64e..c7eb9005f 100644 --- a/packages/core/src/workflow/core.ts +++ b/packages/core/src/workflow/core.ts @@ -2629,7 +2629,10 @@ export function createWorkflow< replayExecutionId: string = randomUUID(), replayStartAt: Date = new Date(), ): Promise => { - const sourceState = await defaultMemory.getWorkflowState(timeTravelOptions.executionId); + const executionMemory = timeTravelOptions.memory ?? defaultMemory; + const workflowSteps = steps as BaseStep[]; + + const sourceState = await executionMemory.getWorkflowState(timeTravelOptions.executionId); if (!sourceState) { throw new Error(`Workflow state not found: ${timeTravelOptions.executionId}`); } @@ -2646,9 +2649,7 @@ export function createWorkflow< ); } - const targetStepIndex = (steps as BaseStep[]).findIndex( - (step) => step.id === timeTravelOptions.stepId, - ); + const targetStepIndex = workflowSteps.findIndex((step) => step.id === timeTravelOptions.stepId); if (targetStepIndex === -1) { throw new Error(`Step '${timeTravelOptions.stepId}' not found in workflow '${id}'`); } @@ -2669,9 +2670,19 @@ export function createWorkflow< const sourceStepCompleteEvents = sourceState.events?.filter((event) => event.type === "step-complete") ?? []; + const stepNameCounts = new Map(); + for (const step of workflowSteps) { + if (typeof step.name !== "string" || step.name.length === 0) { + continue; + } + stepNameCounts.set(step.name, (stepNameCounts.get(step.name) ?? 0) + 1); + } for (let index = 0; index < targetStepIndex; index += 1) { - const step = (steps as BaseStep[])[index]; + const step = workflowSteps[index]; + const stepName = step.name; + const isStepNameUnique = + typeof stepName === "string" && stepName.length > 0 && stepNameCounts.get(stepName) === 1; const checkpointSnapshot = sourceStepData[step.id]; if (checkpointSnapshot) { replayStepData[step.id] = { @@ -2697,7 +2708,7 @@ export function createWorkflow< return ( event.from === step.id || event.name === step.id || - (step.name !== undefined && (event.from === step.name || event.name === step.name)) + (isStepNameUnique && (event.from === stepName || event.name === stepName)) ); }); @@ -2757,7 +2768,7 @@ export function createWorkflow< replayedAt: replayStartAt.toISOString(), }; - await defaultMemory.setWorkflowState(replayExecutionId, { + await executionMemory.setWorkflowState(replayExecutionId, { id: replayExecutionId, workflowId: id, workflowName: name, @@ -2774,15 +2785,13 @@ export function createWorkflow< updatedAt: replayStartAt, }); - const completedStepsData = (steps as BaseStep[]) - .slice(0, targetStepIndex) - .map((step, stepIndex) => ({ - stepId: step.id, - stepName: step.name ?? step.id, - stepIndex, - output: replayStepData[step.id]?.output, - status: replayStepData[step.id]?.status, - })); + const completedStepsData = workflowSteps.slice(0, targetStepIndex).map((step, stepIndex) => ({ + stepId: step.id, + stepName: step.name ?? step.id, + stepIndex, + output: replayStepData[step.id]?.output, + status: replayStepData[step.id]?.status, + })); const executionOptions: WorkflowRunOptions = { executionId: replayExecutionId, @@ -2790,6 +2799,7 @@ export function createWorkflow< conversationId: sourceState.conversationId, context: sourceContext, workflowState: effectiveWorkflowState, + memory: executionMemory, metadata: lineageMetadata, skipStateInit: true, replayFrom: { diff --git a/packages/core/src/workflow/registry.ts b/packages/core/src/workflow/registry.ts index c1f3e0be9..c039b0bce 100644 --- a/packages/core/src/workflow/registry.ts +++ b/packages/core/src/workflow/registry.ts @@ -62,6 +62,14 @@ export class WorkflowRegistry extends SimpleEventEmitter { return globalThis.___voltagent_workflow_registry; } + /** + * Clears registry state. Primarily used by tests for deterministic isolation. + */ + public reset(): void { + this.workflows.clear(); + this.activeExecutions.clear(); + } + /** * Register a workflow with the registry */ diff --git a/packages/core/src/workflow/time-travel.spec.ts b/packages/core/src/workflow/time-travel.spec.ts index 64a7d2fb4..019d8a2fe 100644 --- a/packages/core/src/workflow/time-travel.spec.ts +++ b/packages/core/src/workflow/time-travel.spec.ts @@ -9,8 +9,7 @@ import { andThen } from "./steps"; describe.sequential("workflow.timeTravel", () => { beforeEach(() => { const registry = WorkflowRegistry.getInstance(); - (registry as any).workflows.clear(); - (registry as any).activeExecutions.clear(); + registry.reset(); }); it("should replay from middle step with a new execution id", async () => { diff --git a/packages/core/src/workflow/types.ts b/packages/core/src/workflow/types.ts index 2a8fc2f41..d52aaaa16 100644 --- a/packages/core/src/workflow/types.ts +++ b/packages/core/src/workflow/types.ts @@ -245,6 +245,11 @@ export interface WorkflowTimeTravelOptions { * Optional override for shared workflow state during replay */ workflowStateOverride?: WorkflowStateStore; + /** + * Optional memory adapter to read source execution and persist replay execution state. + * Falls back to workflow default memory when omitted. + */ + memory?: Memory; } export interface WorkflowRetryConfig { diff --git a/packages/server-core/src/handlers/workflow.handlers.ts b/packages/server-core/src/handlers/workflow.handlers.ts index 062a3703c..57c4244f0 100644 --- a/packages/server-core/src/handlers/workflow.handlers.ts +++ b/packages/server-core/src/handlers/workflow.handlers.ts @@ -6,6 +6,8 @@ import type { } from "@voltagent/core"; import { zodSchemaToJsonUI } from "@voltagent/core"; import type { Logger } from "@voltagent/internal"; +import type { z } from "zod"; +import type { WorkflowReplayRequestSchema } from "../schemas/agent.schemas"; import type { ApiResponse, ErrorResponse } from "../types"; import { processWorkflowOptions } from "../utils/options"; import { formatSSE } from "../utils/sse"; @@ -58,6 +60,7 @@ type ResumableStreamingWorkflowExecution = StreamingWorkflowExecution & { ) => Promise; }; +type WorkflowReplayRequestBody = z.infer; type WorkflowTimeTravelRequest = Parameters< NonNullable["timeTravel"]> >[0]; @@ -840,7 +843,7 @@ export async function handleResumeWorkflow( export async function handleReplayWorkflow( workflowId: string, executionId: string, - body: any, + body: WorkflowReplayRequestBody | undefined, deps: ServerProviderDeps, logger: Logger, ): Promise { From 682907220cf3684a74d65148a47abe960e6862b1 Mon Sep 17 00:00:00 2001 From: Omer Aplak Date: Sun, 22 Feb 2026 06:43:43 -0800 Subject: [PATCH 7/7] fix: preserve memory adapter on stream resumes and map replay prep errors --- packages/core/src/workflow/core.ts | 4 ++++ packages/server-core/src/handlers/workflow.handlers.ts | 9 ++++++++- 2 files changed, 12 insertions(+), 1 deletion(-) diff --git a/packages/core/src/workflow/core.ts b/packages/core/src/workflow/core.ts index c7eb9005f..a60b2557a 100644 --- a/packages/core/src/workflow/core.ts +++ b/packages/core/src/workflow/core.ts @@ -3000,6 +3000,7 @@ export function createWorkflow< const executionId = randomUUID(); const startAt = new Date(); const suspendController = createDefaultSuspendController(); + const replayExecutionMemory = timeTravelOptions.memory ?? defaultMemory; let replayOriginalInput: WorkflowInput | undefined; @@ -3076,6 +3077,7 @@ export function createWorkflow< resumeStepIndex, resumeData: resumeInput, }, + memory: replayExecutionMemory, suspendController: resumedSuspendController, }; @@ -3171,6 +3173,7 @@ export function createWorkflow< executionId, suspendController, }; + const streamExecutionMemory = executionOptions.memory ?? defaultMemory; // Save the original input for resume const originalInput = input; @@ -3252,6 +3255,7 @@ export function createWorkflow< resumeStepIndex, resumeData: resumeInput, }, + memory: streamExecutionMemory, suspendController: resumedSuspendController, }; diff --git a/packages/server-core/src/handlers/workflow.handlers.ts b/packages/server-core/src/handlers/workflow.handlers.ts index 57c4244f0..b1076e968 100644 --- a/packages/server-core/src/handlers/workflow.handlers.ts +++ b/packages/server-core/src/handlers/workflow.handlers.ts @@ -903,11 +903,18 @@ export async function handleReplayWorkflow( const message = error instanceof Error ? error.message : "Failed to replay workflow"; const normalizedMessage = message.toLowerCase(); + const isReplayPreparationError = + normalizedMessage.includes("missing historical snapshots") || + normalizedMessage.includes("missing snapshot") || + normalizedMessage.includes("missing input") || + normalizedMessage.includes("no historical") || + normalizedMessage.includes("missing history"); const httpStatus = normalizedMessage.includes("not found") ? 404 : normalizedMessage.includes("cannot time travel") || normalizedMessage.includes("still running") || - normalizedMessage.includes("belongs to workflow") + normalizedMessage.includes("belongs to workflow") || + isReplayPreparationError ? 400 : 500;