WorkflowAgent
pkg/workflow provides a workflow-oriented agent surface for durable execution scenarios.
It mirrors the TypeScript @ai-sdk/workflow package with Go option structs instead of
JavaScript overloads.
Key Types
workflow.WorkflowAgentworkflow.WorkflowResultworkflow.WorkflowStreamResultworkflow.WorkflowRunMultiplexerworkflow.WorkflowChatTransportworkflow.StreamTextIteratorworkflow.SerializableToolDef
Constructor
agent, err := workflow.NewWorkflowAgent(workflow.WorkflowAgent{
ID: "support-agent",
Model: model,
Instructions: "You are a support assistant.",
Tools: []types.Tool{...},
StopWhen: []ai.StopCondition{ai.IsStepCount(20)},
ActiveTools: []string{"search"},
Telemetry: &ai.TelemetrySettings{FunctionID: "support-agent"},
})
System is still accepted for compatibility, but Instructions is the preferred name.
Instructions accepts a string, a system types.Message, or a slice of system messages.
Tools can be supplied as either a []types.Tool or a TypeScript-style keyed ToolSet
using map[string]types.Tool; missing tool names are filled from the map key.
MaxSteps is intentionally not part of WorkflowAgent; use StopWhen with
ai.IsStepCount(n).
Generate
result, err := agent.GenerateWithOptions(ctx, workflow.WorkflowGenerateOptions{
Prompt: "Help me debug this issue",
OnError: func(ctx context.Context, err error) {
log.Printf("workflow error: %v", err)
},
})
if err != nil {
return err
}
_ = result.IsLoopFinished()
_ = result.Output
When Output is an ai.Output specification, GenerateWithOptions forwards its
response format to the model and parses the final text into WorkflowResult.Output.
Stream
streamResult, err := agent.StreamWithOptions(ctx, workflow.WorkflowStreamOptions{
Prompt: "Walk me through this fix",
ActiveTools: []string{"search", "final_answer"},
OnAbort: func(ctx context.Context, steps []types.StepResult) {
log.Printf("aborted after %d steps", len(steps))
},
})
if err != nil {
return err
}
ExperimentalTransform applies an ordered list of ai.StreamTransformFunc to raw model stream chunks before they reach OnChunk or the returned stream, letting you rewrite or drop chunks (e.g. smoothing text, or replacing an oversized provider tool result):
redactLargeResults := func(ctx context.Context, chunk provider.StreamChunk) []provider.StreamChunk {
if chunk.Type == provider.ChunkTypeToolResult && chunk.ToolResult != nil {
if s, ok := chunk.ToolResult.Result.(string); ok && len(s) > 2048 {
replaced := *chunk.ToolResult
replaced.Result = map[string]interface{}{"reason": "too_large_for_stream"}
chunk.ToolResult = &replaced
}
}
return []provider.StreamChunk{chunk}
}
streamResult, err := agent.StreamWithOptions(ctx, workflow.WorkflowStreamOptions{
Prompt: "Walk me through this fix",
ExperimentalTransform: []ai.StreamTransformFunc{redactLargeResults},
})
if err != nil {
return err
}
Chat Transport
WorkflowRunMultiplexer is the server-side http.Handler that emits run-scoped SSE and
supports reattaching to a run:
transport := workflow.NewWorkflowRunMultiplexer()
http.Handle("/api/workflow", transport)
http.Handle("/api/workflow/resume", transport.Resume(runID))
Resume supports a startIndex query parameter (including negative tail-relative offsets).
SendMessages/ReconnectToStream on WorkflowRunMultiplexer return ([]*streaming.SSEEvent, error)
and take workflow.SendMessagesOptions/workflow.ReconnectToStreamOptions respectively:
events, err := transport.SendMessages(ctx, workflow.SendMessagesOptions{
ChatID: "chat-1",
Trigger: "submit-message",
Messages: messages,
})
events, err = transport.ReconnectToStream(ctx, workflow.ReconnectToStreamOptions{
RunID: "run-1",
ChatID: "chat-1",
StartIndex: -10,
})
WorkflowChatTransport is the separate Go client-side equivalent of the TypeScript
WorkflowChatTransport: it implements ai.ChatTransport, POSTing a chat's message history to a
workflow endpoint and streaming the response as ai.UIMessageChunk values over channels:
transport := workflow.NewWorkflowChatTransport(workflow.WorkflowChatTransportOptions{
API: "https://example.com/api/chat",
InitialStartIndex: -10, // tail-relative offset used on reconnect
})
chunks, errs := transport.SendMessages(ctx, ai.ChatTransportSendMessagesRequest{
ChatID: "chat-1",
Trigger: "submit-message",
Messages: messages,
})
chunks, errs = transport.ReconnectToStream(ctx, ai.ChatTransportReconnectToStreamRequest{
ChatID: "chat-1",
})
Serializable Schemas And Tools
Use workflow.MarshalSchema and workflow.UnmarshalSchema for JSON Schema round trips.
Use workflow.SerializeToolSet to strip function fields from tool definitions before
crossing a workflow boundary, then workflow.ResolveSerializableTools to restore
non-executable tool descriptors in the resumed step.
Use workflow.ValidateSerializableToolInput to validate reconstructed tool inputs
against the serialized JSON Schema before executing an attached Go function.
Provider Serialization
providerutils.SerializeModel returns a JSON-safe map containing provider, modelId,
and config. Provider auth fields and function-valued fields are omitted. Use
providerutils.DeserializeModel(providerID, modelID, config) after importing the target
provider package so its deserializer is registered.
providerutils.SerializeModel/DeserializeModel are convenience wrappers around
provider.SerializeModel(provider.LanguageModel) (provider.SerializedModel, error) and
provider.DeserializeModel(provider.SerializedModel) (provider.LanguageModel, error) for
the language model case specifically. Every other model kind that a provider can
serialize has the same pair of functions and a matching Register*Deserializer, one set
per kind — for every provider model TS marks WORKFLOW_SERIALIZE:
provider.SerializeImageModel / DeserializeImageModel / RegisterImageModelDeserializer
provider.SerializeVideoModel / DeserializeVideoModel / RegisterVideoModelDeserializer
provider.SerializeSpeechModel / DeserializeSpeechModel / RegisterSpeechModelDeserializer
provider.SerializeTranscriptionModel / DeserializeTranscriptionModel / RegisterTranscriptionModelDeserializer
provider.SerializeEmbeddingModel / DeserializeEmbeddingModel / RegisterEmbeddingModelDeserializer
provider.SerializeEvaluationModel / DeserializeEvaluationModel / RegisterEvaluationModelDeserializer
None of this is wired into the workflow runtime automatically — there is no
implicit "serialize every model field on a workflow struct" step. Your
application calls Serialize*Model explicitly before persisting workflow
state or crossing a boundary, and calls Deserialize*Model (after importing
the target provider package, so its deserializer registers itself via
init()) explicitly when resuming. A model that implements
provider.SerializableModelStrict instead of the plain SerializableModel
interface can refuse to serialize — for example, an Open Responses model with
registered extension codecs, since the codecs are functions and cannot be
reconstructed from JSON.