diff --git a/.changeset/four-cobras-yawn.md b/.changeset/four-cobras-yawn.md new file mode 100644 index 000000000..b7ab74639 --- /dev/null +++ b/.changeset/four-cobras-yawn.md @@ -0,0 +1,25 @@ +--- +"@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. + +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/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/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 c18aad03b..28d0337a0 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,46 @@ export class WorkflowChain< return workflow.startAsync(input, options); } + /** + * 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, + ): 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 + * 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, + ): 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..a60b2557a 100644 --- a/packages/core/src/workflow/core.ts +++ b/packages/core/src/workflow/core.ts @@ -54,9 +54,11 @@ import type { WorkflowStepData, WorkflowStreamResult, WorkflowSuspensionMetadata, + WorkflowTimeTravelOptions, } 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); @@ -118,6 +120,68 @@ const isWorkflowStepStatus = (value: unknown): value is WorkflowStepData["status value === "cancelled" || value === "skipped"; +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"; +}; + +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; @@ -177,6 +241,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 => { @@ -997,9 +1078,44 @@ 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", { + replayExecutionId: options.replayFrom.executionId, + executionId, + }); + } + } catch (error) { + logger.warn("Failed to get source trace IDs for replay:", { + error, + replayExecutionId: options.replayFrom.executionId, + executionId, + }); + } + } else if (options?.resumeFrom?.executionId) { try { const workflowState = await executionMemory.getWorkflowState(executionId); // Look for trace IDs from the original execution @@ -1010,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, + }); } } @@ -1027,6 +1150,7 @@ export function createWorkflow< input: input, context: contextMap, resumedFrom, + replayedFrom, }); // Wrap entire execution in root span @@ -1062,7 +1186,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 @@ -2493,6 +2617,218 @@ 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 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}`); + } + + 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 = workflowSteps.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") ?? []; + 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 = 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] = { + input: checkpointSnapshot.input, + output: checkpointSnapshot.output, + status: toWorkflowStepStatus(checkpointSnapshot.status, logger), + error: serializeStepError(checkpointSnapshot.error), + }; + continue; + } + + 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 || + (isStepNameUnique && (event.from === stepName || event.name === stepName)) + ); + }); + + if (fallbackEvent) { + replayStepData[step.id] = { + input: fallbackEvent.input, + output: fallbackEvent.output, + status: toWorkflowStepStatus(fallbackEvent.status, logger), + 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 executionMemory.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 = 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, + userId: sourceState.userId, + conversationId: sourceState.conversationId, + context: sourceContext, + workflowState: effectiveWorkflowState, + memory: executionMemory, + metadata: lineageMetadata, + skipStateInit: true, + replayFrom: { + executionId: timeTravelOptions.executionId, + stepId: timeTravelOptions.stepId, + }, + 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 +2989,170 @@ 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(); + const replayExecutionMemory = timeTravelOptions.memory ?? defaultMemory; + + 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(); + }, + ); + + 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, + }, + memory: replayExecutionMemory, + suspendController: resumedSuspendController, + }; + + executeInternal(replayOriginalInput, 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 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, + 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; + return resumeSuspendedReplayStream(replayResult, input, opts); + }, + 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); }, @@ -2673,6 +3173,7 @@ export function createWorkflow< executionId, suspendController, }; + const streamExecutionMemory = executionOptions.memory ?? defaultMemory; // Save the original input for resume const originalInput = input; @@ -2714,6 +3215,92 @@ 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, + }, + memory: streamExecutionMemory, + 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, @@ -2729,117 +3316,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/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/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/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 new file mode 100644 index 000000000..019d8a2fe --- /dev/null +++ b/packages/core/src/workflow/time-travel.spec.ts @@ -0,0 +1,318 @@ +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.reset(); + }); + + 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 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() }); + + 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..d52aaaa16 100644 --- a/packages/core/src/workflow/types.ts +++ b/packages/core/src/workflow/types.ts @@ -224,6 +224,34 @@ 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; + /** + * 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 { /** * Number of retry attempts for a step when it throws an error @@ -287,6 +315,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) @@ -361,6 +394,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 @@ -733,6 +777,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/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..b1076e968 100644 --- a/packages/server-core/src/handlers/workflow.handlers.ts +++ b/packages/server-core/src/handlers/workflow.handlers.ts @@ -1,6 +1,13 @@ -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 { 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"; @@ -53,6 +60,11 @@ type ResumableStreamingWorkflowExecution = StreamingWorkflowExecution & { ) => Promise; }; +type WorkflowReplayRequestBody = z.infer; +type WorkflowTimeTravelRequest = Parameters< + NonNullable["timeTravel"]> +>[0]; + function parseReplaySequence(value: StreamQueryValue): number | undefined { if (value === undefined || value === null || value === "") { return undefined; @@ -824,6 +836,96 @@ 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: WorkflowReplayRequestBody | undefined, + 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 Partial< + Pick, "timeTravel"> + >; + + if (typeof workflowWithReplay.timeTravel !== "function") { + return { + success: false, + error: "Workflow does not support replay", + httpStatus: 400, + }; + } + + const replayOptions: WorkflowTimeTravelRequest = { + executionId, + stepId: stepId.trim(), + inputData, + resumeData, + workflowStateOverride, + }; + const result = await workflowWithReplay.timeTravel(replayOptions); + + 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 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") || + isReplayPreparationError + ? 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..75dbf0f7c 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().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 + .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); diff --git a/website/docs/api/endpoints/workflows.md b/website/docs/api/endpoints/workflows.md index 51ff72e5e..b2facbd05 100644 --- a/website/docs/api/endpoints/workflows.md +++ b/website/docs/api/endpoints/workflows.md @@ -510,6 +510,128 @@ 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 + } + } +} +``` + +**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) +- `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. @@ -525,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", 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..730de647d 100644 --- a/website/docs/workflows/streaming.md +++ b/website/docs/workflows/streaming.md @@ -42,11 +42,13 @@ Workflows emit these event types during execution: ### Consuming the Stream -VoltAgent provides three 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 // Method 1: Stream execution for real-time events @@ -84,6 +86,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 +804,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`: