Skip to main content

Workflow

Package github.com/digitallysavvy/go-ai/pkg/workflow is the Go port of the TypeScript @ai-sdk/workflow package. It provides an agent for durable execution, where one run spans several process executions, and the client and server pieces that stream a run to a chat UI. For a walkthrough of WorkflowAgent, see WorkflowAgent.

WorkflowAgent​

workflow.NewWorkflowAgent validates a WorkflowAgent value and returns it.

func NewWorkflowAgent(agent WorkflowAgent) (*WorkflowAgent, error)
FieldTypeDescription
Modelprovider.LanguageModel
Systemstring
Instructionsinterface{}Instructions is the TypeScript-compatible name for system instructions.
AllowSystemInMessagesboolAllowSystemInMessages permits system-role messages in Messages. By default WorkflowAgent rejects system messages; use Instructions/System for default system prompts.
Tools[]types.Tool
ToolSetmap[string]types.Tool
StopWhen[]ai.StopCondition
Outputinterface{}
Telemetry*ai.TelemetrySettings
IDstring
Promptstring
OnStartStartCallback
OnStepStartStepStartCallback
OnToolExecutionStartToolExecutionStartCallback
OnToolExecutionEndToolExecutionEndCallback
OnStepEndStepEndCallback
OnStepFinishStepFinishCallbackDeprecated: use OnStepEnd.
OnEndEndCallback
OnFinishFinishCallbackDeprecated: use OnEnd.
OnErrorErrorCallback
OnAbortAbortCallback
PrepareCallPrepareCallHook
PrepareStepPrepareStepHook
FilterActiveToolsFilterActiveToolsHook
CallOptionsSchemaschema.Schema
CallOptionsinterface{}
ActiveTools[]string
Temperature*float64
MaxTokens*int
TopP*float64
TopK*int
FrequencyPenalty*float64
PresencePenalty*float64
StopSequences[]string
Seed*int
Headersmap[string]string
Reasoning*types.ReasoningLevel
SendReasoning*bool
ProviderOptionsmap[string]interface{}
RuntimeContextinterface{}
ToolsContextmap[string]interface{}
ToolChoicetypes.ToolChoice
Include*ai.IncludeOptions
ExperimentalSandboxinterface{}
ExperimentalRefineToolInputmap[string]ai.ToolInputRefiner
RepairToolCallai.ToolCallRepairFunctionRepairToolCall attempts to repair tool calls that fail to parse.
ExperimentalRepairToolCallai.ToolCallRepairFunctionExperimentalRepairToolCall is a deprecated alias for RepairToolCall. Deprecated: use RepairToolCall.
ExperimentalToolApprovalSecret[]byteExperimentalToolApprovalSecret signs issued approval requests and verifies resumed approvals before approved tools execute.
MaxRetries*intMaxRetries controls transient provider call retries for each model call, forwarded to agent.AgentConfig.MaxRetries. Defaults to 2 when nil, matching TS WorkflowAgent's mergedGenerationSettings.maxRetries ?? 2.
Timeout*ai.TimeoutConfigTimeout provides granular timeout controls, forwarded to agent.AgentConfig.Timeout.
ExperimentalDownloadai.DownloadFunctionExperimentalDownload customizes remote file URL downloads before model calls, forwarded to agent.AgentConfig.ExperimentalDownload.
MethodDescription
Generate(ctx, prompt string, opts *agent.AgentGenerateOptions) (*WorkflowResult, error)Runs the agent and returns the final result.
GenerateWithOptions(ctx, WorkflowGenerateOptions) (*WorkflowResult, error)Runs with workflow-native call options.
Stream(ctx, prompt string, opts *agent.AgentStreamOptions) (*WorkflowStreamResult, error)Runs in streaming mode.
StreamWithOptions(ctx, WorkflowStreamOptions) (*WorkflowStreamResult, error)Streams with workflow-native stream options.

WorkflowGenerateOptions​

FieldTypeDescription
Promptstring
Messages[]types.Message
Systemstring
Instructionsinterface{}
AllowSystemInMessagesbool
Tools[]types.Tool
ToolSetmap[string]types.Tool
StopWhen[]ai.StopCondition
Telemetry*ai.TelemetrySettings
RuntimeContextinterface{}
ToolsContextmap[string]interface{}
Include*ai.IncludeOptions
ExperimentalSandboxinterface{}
ExperimentalRefineToolInputmap[string]ai.ToolInputRefiner
RepairToolCallai.ToolCallRepairFunctionRepairToolCall attempts to repair tool calls that fail to parse.
ExperimentalRepairToolCallai.ToolCallRepairFunctionExperimentalRepairToolCall is a deprecated alias for RepairToolCall. Deprecated: use RepairToolCall.
ExperimentalToolApprovalSecret[]byteExperimentalToolApprovalSecret overrides the agent's approval secret.
MaxRetries*intMaxRetries overrides the agent's MaxRetries for this call.
Timeout*ai.TimeoutConfigTimeout overrides the agent's Timeout for this call.
ExperimentalDownloadai.DownloadFunctionExperimentalDownload overrides the agent's ExperimentalDownload for this call.
OnStartStartCallback
OnStepStartStepStartCallback
OnToolExecutionStartToolExecutionStartCallback
OnToolExecutionEndToolExecutionEndCallback
OnStepEndStepEndCallback
OnStepFinishStepFinishCallbackDeprecated: use OnStepEnd.
OnEndEndCallback
OnFinishFinishCallbackDeprecated: use OnEnd.
OnErrorErrorCallback
OnAbortAbortCallback

WorkflowStreamOptions​

FieldTypeDescription
Promptstring
Messages[]types.Message
Systemstring
Instructionsinterface{}
AllowSystemInMessagesbool
Tools[]types.Tool
ToolSetmap[string]types.Tool
StopWhen[]ai.StopCondition
Telemetry*ai.TelemetrySettings
ActiveTools[]string
RuntimeContextinterface{}
ToolsContextmap[string]interface{}
Include*ai.IncludeOptions
ExperimentalSandboxinterface{}
ExperimentalRefineToolInputmap[string]ai.ToolInputRefiner
RepairToolCallai.ToolCallRepairFunctionRepairToolCall attempts to repair tool calls that fail to parse.
ExperimentalRepairToolCallai.ToolCallRepairFunctionExperimentalRepairToolCall is a deprecated alias for RepairToolCall. Deprecated: use RepairToolCall.
ExperimentalToolApprovalSecret[]byteExperimentalToolApprovalSecret overrides the agent's approval secret.
MaxRetries*intMaxRetries overrides the agent's MaxRetries for this call.
Timeout*ai.TimeoutConfigTimeout overrides the agent's Timeout for this call.
ExperimentalDownloadai.DownloadFunctionExperimentalDownload overrides the agent's ExperimentalDownload for this call.
ExperimentalTransform[]ai.StreamTransformFuncExperimentalTransform is an ordered list of transforms applied to raw model stream chunks before they reach OnChunk or the returned WorkflowStreamResult, matching TS WorkflowAgent's experimental_transform stream option (hash 165455d). Previously declared-but-unused in TS; this forwards it through to the underlying ai.StreamText call so it actually runs.
OnChunkfunc(chunk provider.StreamChunk)
OnStartStartCallback
OnStepStartStepStartCallback
OnToolExecutionStartToolExecutionStartCallback
OnToolExecutionEndToolExecutionEndCallback
OnStepEndStepEndCallback
OnStepFinishStepFinishCallbackDeprecated: use OnStepEnd.
OnEndEndCallback
OnFinishFinishCallbackDeprecated: use OnEnd.
OnErrorErrorCallback
OnAbortAbortCallback

Results​

workflow.WorkflowResult is the final non-streaming result. IsLoopFinished() reports whether the run ended naturally.

FieldTypeDescription
(embedded)*agent.AgentResultEmbedded.

workflow.WorkflowStreamResult wraps the streaming result.

FieldTypeDescription
(embedded)*ai.StreamTextResultEmbedded.

Hooks and callbacks​

The hook and callback function types mirror the ToolLoopAgent ones.

TypeDescription
workflow.PrepareStepHookMutates per-step options using the current history.
workflow.PrepareCallHookMutates per-step call options before model invocation.
workflow.FilterActiveToolsHookReduces the available tools for a step.
workflow.LanguageModelCallOptionsThe per-step call configuration, as in ToolLoopAgent.
workflow.StartCallbackFires once before the first step.
workflow.StepStartCallbackFires before each model step.
workflow.StepEndCallback, workflow.StepFinishCallbackFire after each completed step.
workflow.ToolExecutionStartCallback, workflow.ToolExecutionEndCallbackFire around local tool execution.
workflow.EndCallback, workflow.FinishCallbackFire once after the run completes.
workflow.ErrorCallbackFires when execution returns an error.
workflow.AbortCallbackFires when the context cancels the run.
FieldTypeDescription
StepNumberint
Systemstring
AllowSystemInMessagesbool
Messages[]types.Message
Tools[]types.Tool
ToolChoicetypes.ToolChoice
CallOptionsinterface{}
Temperature*float64
MaxTokens*int
TopP*float64
TopK*int
FrequencyPenalty*float64
PresencePenalty*float64
StopSequences[]string
Seed*int
Headersmap[string]string
Reasoning*types.ReasoningLevel
SendReasoning*bool
ProviderOptionsmap[string]interface{}
RuntimeContextinterface{}
ToolsContextmap[string]interface{}
ExperimentalSandboxinterface{}
PreviousSteps[]types.StepResult
AccumulatedUsagetypes.Usage
CustomDatainterface{}
StopWhen[]ai.StopConditionStopWhen, ActiveTools and ExperimentalDownload mirror TS prepareCall's per-call stopWhen/activeTools/download settings parity (d56638a): a PrepareCall/PrepareStep hook can read the current effective value here and mutate it to override the call.
ActiveTools[]string
ExperimentalDownloadai.DownloadFunction
MaxRetries*intMaxRetries and Timeout mirror TS's maxRetries/abortSignal prepareCall parity (419adc7). Timeout is the Go stand-in for TS's AbortSignal.
Timeout*ai.TimeoutConfig
InitialInstructionsstringInitialInstructions and InitialMessages are the original (unmutated) instructions/messages the call was invoked with, matching TS prepareCall's initial-inputs parity (b666f57). They are read-only: a hook should mutate System/Messages to change what is sent, not these.
InitialMessages[]types.Message

Iterating steps​

workflow.StreamTextIterator iterates the step results of a run. Next(ctx) returns the next *types.StepResult and io.EOF when the steps are exhausted. Close() marks the iterator closed.

ConstructorDescription
workflow.NewStreamTextIterator(steps []types.StepResult)Iterates a slice you already have.
workflow.NewStreamTextIteratorFromChannel(ch <-chan types.StepResult)Iterates live steps from a channel.
workflow.NewStreamTextIteratorFromResult(result *WorkflowResult)Iterates the steps of a result.

Serializable tools​

A durable run cannot persist a Go function. These helpers carry the tool definitions across an execution boundary, and you attach Execute functions again by name afterward.

FunctionDescription
workflow.SerializeToolSet(tools []types.Tool) (map[string]SerializableToolDef, error)Converts tools to a JSON-safe map keyed by name.
workflow.ResolveSerializableTools(defs map[string]SerializableToolDef) []types.ToolRebuilds non-executable tool descriptors.
workflow.ValidateSerializableToolInput(def SerializableToolDef, input interface{}) (interface{}, error)Validates input against the serialized schema. It applies schema defaults first, and returns the defaulted value. Pass that value to the tool, not the raw input.
workflow.MarshalSchema(schema interface{}) (map[string]interface{}, error)Converts a JSON-serializable schema to a map.
workflow.UnmarshalSchema(data map[string]interface{}) (interface{}, error)Rebuilds a schema as a generic JSON value.
FieldTypeDescription
Namestring
Descriptionstring
Titlestring
Parametersmap[string]interface{}
InputSchemamap[string]interface{}
Typestring
IDstring
Argsmap[string]interface{}
Strict*bool
ProviderExecutedbool
IsProviderExecutedbool
SupportsDeferredResultsbool
InputExamples[]types.ToolInputExample
ProviderMetadatamap[string]interface{}

Chat transports​

WorkflowChatTransport​

workflow.WorkflowChatTransport is a client-side ai.ChatTransport. It POSTs the message history to a workflow chat endpoint and reads the response as a stream of UI message chunks. When the response ends without a finish chunk, for example after a network drop or a function timeout, it reconnects with a GET to <api>/<runId>/stream?startIndex=N. It also repairs UI message stream framing and, on a resume with a negative start index, drops deltas and ends whose start fell outside the resumed window.

func NewWorkflowChatTransport(opts WorkflowChatTransportOptions) *WorkflowChatTransport

SendMessages(ctx, ai.ChatTransportSendMessagesRequest) and ReconnectToStream(ctx, ai.ChatTransportReconnectToStreamRequest) implement ai.ChatTransport. See Stream transport helpers.

FieldTypeDescription
APIstringAPI is the chat endpoint. Defaults to "/api/chat".
HTTPClient*http.ClientHTTPClient is used for both the initial POST and any reconnect GETs. Defaults to http.DefaultClient.
OnChatSendMessagefunc(resp *http.Response, req ai.ChatTransportSendMessagesRequest) errorOnChatSendMessage is invoked after the initial POST completes, useful for inspecting response headers or tracking chat history.
OnChatEndfunc(WorkflowChatTransportEndEvent) errorOnChatEnd is invoked once, after a "finish" chunk is observed.
MaxConsecutiveErrorsintMaxConsecutiveErrors bounds reconnect attempts. Defaults to 3.
InitialStartIndexintInitialStartIndex is the default startIndex used by ReconnectToStream when it is called directly (not as part of a SendMessages recovery). Negative values are tail-relative (e.g. -10 reads the last 10 chunks), useful for resuming a chat UI after a page refresh without replaying the whole conversation. Defaults to 0 (replay from the beginning).
PrepareSendMessagesRequestfunc(ai.ChatTransportSendMessagesRequest) (PreparedChatRequest, error)PrepareSendMessagesRequest customizes the API endpoint, body, and headers used for the initial POST.
PrepareReconnectToStreamRequestfunc(WorkflowChatTransportReconnectContext) (PreparedChatRequest, error)PrepareReconnectToStreamRequest customizes the API endpoint and headers used for reconnect GETs.
FieldTypeDescription
ChatIDstring
MessageIDstring
Triggerstring
Messagesinterface{}
Bodymap[string]interface{}
Headersmap[string]string
FieldTypeDescription
RunIDstring
ChatIDstring
StartIndexint
Headersmap[string]string

workflow.WorkflowChatTransportReconnectContext is passed to PrepareReconnectToStreamRequest. workflow.WorkflowChatTransportEndEvent is passed to OnChatEnd. workflow.ChatEndEvent is the equivalent event for the multiplexer's OnChatEnd. It carries ChatID, ChunkIndex, RunID and the SSE Events.

FieldTypeDescription
ChatIDstring
ChunkIndexint

WorkflowRunMultiplexer​

workflow.WorkflowRunMultiplexer is a Go-only server and client helper that predates WorkflowChatTransport. It serves run-scoped SSE and resume handlers. Its client helpers return raw SSE events, not UI message chunks, so it does not implement ai.ChatTransport. For new client code, use WorkflowChatTransport.

func NewWorkflowRunMultiplexer(opts ...WorkflowRunMultiplexerOptions) *WorkflowRunMultiplexer
MethodDescription
ServeHTTP(w, r)Starts a new run stream and emits lifecycle events.
Resume(runID string) http.HandlerReturns a handler that attaches an SSE reader to a run in progress.
SendMessages(ctx, SendMessagesOptions) ([]*streaming.SSEEvent, error)POSTs chat messages and returns parsed SSE events.
ReconnectToStream(ctx, ReconnectToStreamOptions) ([]*streaming.SSEEvent, error)GETs the API to resume a run-scoped SSE stream.
FieldTypeDescription
APIstring
HTTPClient*http.Client
MaxConsecutiveErrorsint
InitialStartIndexint
OnChatSendMessagefunc(*http.Response, SendMessagesOptions) error
OnChatEndfunc(ChatEndEvent) error
PrepareSendMessagesRequestfunc(SendMessagesOptions) (PreparedChatRequest, error)
PrepareReconnectToStreamRequestfunc(ReconnectToStreamOptions) (PreparedChatRequest, error)
FieldTypeDescription
APIstring
Bodymap[string]interface{}
Headersmap[string]string

Durable harness runner​

These functions run a harness agent across process executions. Each execution returns a serializable workflow.HarnessWorkflowState. Your durable-workflow runtime stores it as the checkpoint, and you pass it to the next execution.

FunctionDescription
workflow.CreateHarnessWorkflowState(HarnessWorkflowInput) HarnessWorkflowStateBuilds the initial state for one user turn.
workflow.RunHarnessAgent(ctx, RunHarnessAgentOptions) (HarnessWorkflowState, error)Runs one durable execution. It resumes or starts the session and streams the turn's chunks to Writable. When TimeSliceSeconds is positive, it races the turn against that budget and suspends the turn when the slice ends.
workflow.RunHarnessAgentTimeSlice(ctx, RunHarnessAgentTimeSliceOptions) (HarnessWorkflowState, error)Runs one time-boxed slice. When the slice ends before the turn, the status is ready_for_next_step.
workflow.RunHarnessAgentStep(ctx, RunHarnessAgentStepOptions) (HarnessWorkflowState, error)Runs until the next step boundary. Set StopWhen on the agent, for example ai.IsStepCount(1).
workflow.RunHarnessAgentSlice(ctx, RunHarnessAgentSliceOptions)Deprecated alias of RunHarnessAgentTimeSlice that reports timed_out where the new function reports ready_for_next_step.
workflow.FinalizeHarnessWorkflow(state) (HarnessWorkflowFinalResult, error)Collapses a terminal state into its result. It returns an error if the run failed.

workflow.DefaultTimeSliceSeconds is 750. A Vercel Fluid Compute instance is recycled at about 800 seconds, so a 750-second slice leaves time for the next execution to reattach to the sandbox that is still running.

Status​

ConstantValueDescription
HarnessWorkflowStatusNotStarted"not_started"HarnessWorkflowStatusNotStarted is the fresh state before any execution has run.
HarnessWorkflowStatusReadyForNextStep"ready_for_next_step"HarnessWorkflowStatusReadyForNextStep means the turn remains unfinished and ContinueFrom carries the cursor for the next execution.
HarnessWorkflowStatusAwaitingToolApproval"awaiting_tool_approval"HarnessWorkflowStatusAwaitingToolApproval means the turn emitted one or more tool approval/result requests and ContinueFrom carries the suspended turn.
HarnessWorkflowStatusFinished"finished"HarnessWorkflowStatusFinished means the agent turn completed on its own; FinalResult is set.
HarnessWorkflowStatusFailed"failed"HarnessWorkflowStatusFailed means the turn errored; Error is set.
HarnessWorkflowStatusTimedOut"timed_out"HarnessWorkflowStatusTimedOut is returned by the deprecated RunHarnessAgentSlice in place of HarnessWorkflowStatusReadyForNextStep. Deprecated: use HarnessWorkflowStatusReadyForNextStep.

Types​

workflow.HarnessWorkflowAgent is the subset of *harness.Agent the runner drives: HasOutput, CreateSession, Stream and ContinueStream. *harness.Agent satisfies it.

FieldTypeDescription
Promptharness.Prompt
Messages[]types.Message
SessionIDstring
ResumeFrom*harness.ResumeSessionState
ContinueFrom*harness.ContinueTurnState
FieldTypeDescription
SessionIDstringSessionID is the stable harness session id; doubles as the sandbox name across processes.
Promptharness.PromptPrompt is the new user turn for this run. Sent once, on the execution that starts the turn.
Messages[]types.MessageMessages carries full model messages for continuing a suspended approval turn (e.g. tool-approval-response content). When non-nil, the next execution sends these instead of Prompt/ContinueFrom.
StatusHarnessWorkflowStatus
ResumeFrom*harness.ResumeSessionStateResumeFrom carries resume coordinates for the next user turn.
ContinueFrom*harness.ContinueTurnStateContinueFrom carries continuation coordinates for this run's current suspended turn.
StreamContext*HarnessWorkflowStreamContext
FinalResult*HarnessWorkflowFinalResult
Errorstring
FieldTypeDescription
SessionIDstring
FinishReasonstring
Usage*HarnessWorkflowUsageSummary
OutputanyOutput is the agent's parsed and schema-validated output when the agent has an output specification (HarnessWorkflowAgent.HasOutput).
FieldTypeDescription
AgentHarnessWorkflowAgent
StateHarnessWorkflowState
SandboxSessionproviderutils.SandboxSessionSandboxSession, when set, is forwarded to HarnessWorkflowAgent. CreateSession as a caller-owned sandbox (harness.CreateSessionOptions. SandboxSession) instead of letting the agent's own SandboxProvider create/resume one. Mirrors TS RunHarnessAgentOptions.sandboxSession.
TimeSliceSecondsfloat64TimeSliceSeconds is the wall-clock budget for this execution. Zero means no time slice: the run continues until the harness's own turn (or, with a StopWhen-configured agent, step) boundary.
DestroyOnFinishboolDestroyOnFinish controls whether to destroy the sandbox when the run finishes or fails. Defaults to false: the session is parked and a fresh resume state is returned in ResumeFrom, so the next user turn reattaches to the same conversation (multi-turn chat). Set true for a one-shot run that should release the sandbox when the run ends.
WritableHarnessWorkflowWriterWritable is where to write the turn's UI-message chunks. Required — TS's default resolveWorkflowWritable() (a Workflow DevKit getWritable()) has no Go equivalent.
FieldTypeDescription
AgentHarnessWorkflowAgent
StateHarnessWorkflowState
SandboxSessionproviderutils.SandboxSession
TimeSliceSecondsfloat64TimeSliceSeconds defaults to DefaultTimeSliceSeconds when zero.
DestroyOnFinishbool
WritableHarnessWorkflowWriter
FieldTypeDescription
AgentHarnessWorkflowAgent
StateHarnessWorkflowState
SandboxSessionproviderutils.SandboxSession
DestroyOnFinishbool
WritableHarnessWorkflowWriter
FieldTypeDescription
AgentHarnessWorkflowAgent
StateHarnessWorkflowState
SandboxSessionproviderutils.SandboxSession
TimeSliceSecondsfloat64
SliceTimeoutSecondsfloat64SliceTimeoutSeconds is used when TimeSliceSeconds is zero. Deprecated: use TimeSliceSeconds.
DestroyOnFinishbool
WritableHarnessWorkflowWriter

workflow.HarnessWorkflowStreamContext is the serializable subset of in-flight UI message state carried across an execution boundary. workflow.HarnessWorkflowActiveToolInput is a tool input that was partly streamed when the execution ended. workflow.HarnessWorkflowUsageSummary is a minimal token usage summary.

FieldTypeDescription
ActiveTextPartsmap[string]ai.UIMessageChunk
ActiveReasoningPartsmap[string]ai.UIMessageChunk
ActiveToolInputsmap[string]HarnessWorkflowActiveToolInput
PendingToolInputsmap[string]ai.UIMessageChunk
FieldTypeDescription
InputTokens*int64
OutputTokens*int64

Writer​

workflow.HarnessWorkflowWriter receives the UI message chunks of one execution:

type HarnessWorkflowWriter interface {
Write(chunk ai.UIMessageChunk) error
Close() error
}

Write is called once per chunk, in order. Close is called only when the run reaches finished. The other statuses leave the writer open because a later execution keeps writing, or the failure propagates. workflow.NewChanHarnessWorkflowWriter(ch) returns a *workflow.ChanHarnessWorkflowWriter that forwards every chunk to ch and closes ch on Close.